UN-4044 [FEAT] Agent-KV extraction engine — public API, schema compiler and hardened code sandbox - #2309
UN-4044 [FEAT] Agent-KV extraction engine — public API, schema compiler and hardened code sandbox#2309vishnuszipstack wants to merge 86 commits into
Conversation
Design for productizing the unstract-agentic-table KV extractor as a metered async API: OSS scaffold (app, keys, jobs, dispatch, metering seams) + cloud executor plugin (engine port, codegen sandbox worker). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Corrects the spec to the platform's real mechanisms: executor-RPC dispatch_with_callback (single UUID task_id) instead of workflow transport rules; backend capability-probe gating instead of executor registry introspection; schema compiler carved into OSS as single source of truth; real usage-reporting paths (v1/usage/batch + cloud extras attribution); periodic-job scheduling homes; precise cancel semantics; webhook SSRF controls; rate-limiter scoping; pdfplumber for page counting. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Schema-compiler package, filesystem type, app+models with terminal write guard, key management, public mount+auth, submit serializer, limiters, dispatch glue, job endpoints, validate, internal APIs, callbacks, SSRF-guarded webhooks, sweeps, docs. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Two failure-path defects found in review, both inherited from the task-8 brief's own provisional code: 1. Concurrency slot leak: stage_input()/job.save() ran unguarded after check_and_acquire() succeeded, so an object-store or DB error left the slot stuck for up to the 6h TTL and returned an unhandled 500. Now wrapped in try/except: always release the slot, mark_terminal(FAILED) only when the row was actually persisted, and return a safe "Job could not be accepted; nothing was billed." 500. 2. dispatch_job's ExecutionContext construction and the platform-key lookup ran outside its try block, so a raw (non-DispatchError) exception there skipped mark_terminal, slot release, and the safe 500 entirely. Widened dispatch.py's try to cover the full fallible body (frozen executor_params/callback contract unchanged), and added a belt-and-braces except Exception in execution_views.py with identical cleanup to except DispatchError, via a shared _fail_job_response helper. The 300s sync-timeout polling loop (a third review finding) is intentionally untouched per controller ruling. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Add /internal/v1/agent-kv/jobs/<job_id>/stage/ and .../finalize/, mounted
via a new agent_kv/internal_urls.py included in internal_base_urls.py.
Auth is ambient (InternalAPIAuthMiddleware); both views take empty
auth/permission classes and require org_id in the body (400 without it).
StageReportView merges one stage's progress into job.stages, flipping
PENDING/DISPATCHED -> RUNNING on first report only. It is the sole write
gate for job.stages: only {status, seconds?, ...counters} is persisted,
so unexpected top-level body keys never leak into a stored entry (carried
constraint from task 9's status-endpoint review). A late report against
an already-terminal job is a 200 no-op via the TERMINAL-excluding update
queryset.
FinalizeView writes the result then calls mark_terminal on success, or
calls mark_terminal(FAILED, error=...) on failure; the concurrency slot
is always released in a finally. It pre-checks the job's terminal state
so a duplicate/late finalize call returns finalized:false without ever
rewriting an already-written result. Response is
{finalized, webhook_url, status} -- built out now so task 12 doesn't
need to reopen this file.
79 previous agent_kv tests + 13 new all pass (80 total).
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
…body 400s
Controller-authorized fixes from task-11 review (3 Important issues):
1. Stage-merge lost-update race: replace the Python read-modify-write
(dict(job.stages or {}) -> merge -> full-column .update(stages=...))
with a DB-side jsonb `||` concatenation expression
(_stage_merge_expression), so .update() only ever merges the reported
stage's single-entry dict at the database level. Two concurrent
reports for different stages can no longer clobber each other. The
earlier read is kept only for the RUNNING-flip decision and the
terminal no-op check, both re-guarded at update time by job_qs's
existing TERMINAL exclusion.
2. counters bypassing the write gate: add _sanitize_counters(), which
drops any counter key colliding with the reserved {"status",
"seconds"} fields and any non-scalar (dict/list) value, before
merging into the stage entry. Closes the task-9-review write-gate
constraint against counters clobbering reserved fields or admitting
nested structures.
3. Malformed bodies: StageReportView now 400s on missing `stage` or a
`status` that isn't exactly "running"/"done"; FinalizeView now 400s
unless `success` is a strict bool (isinstance check, not truthy/
falsy) -- a malformed finalize call no longer silently persists a
FAILED job with an empty error. Both 400 paths return before the
concurrency slot could be touched, so no release() call happens on
them (nothing was finalized).
test_internal_views.py: 13 -> 23 tests (10 new: expression-builder unit
tests, counters-sanitization tests, and the three new 400 cases).
90 passed in agent_kv/tests/ (67 pre-existing outside this file + 23
here).
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
…k call Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
…ow_http paths Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Recording is built and wired (LLM usage_records -> v1/usage/batch, page usage -> platform-service subscription hook, per-job usage_summary); the fail-closed pre-dispatch admission/quota reserve and Stripe invoicing are deferred to the billing sub-project. Correct the now-stale branch state. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
…allowlist
The architecture review flagged that §8a overstated the accepted residual.
Confirmed against workers/sandbox/gate.py: the import allowlist is exactly
{json, math, statistics, decimal, datetime, re, collections, itertools,
functools, sys}; pathlib, os, io, socket and urllib are all outside it.
The residual is real (open() is a builtin and cannot be denylisted without
breaking the runner stub contract) but narrower than documented: a single
read of a known path, not directory traversal and not exfiltration.
Corrected in the sandbox design §8a, the architecture overview, and the API
reference §12.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
# Conflicts: # docker/docker-compose.yaml # tests/compose/docker-compose.test.yaml # tests/groups.yaml
… PG transport
UN-4046 made the Postgres task transport unconditional and disabled the
Celery executor. Two Agent-KV queues were left wired to the broker-era
config, and both fail SILENTLY: a PG queue with no consumer still accepts
work (enqueue succeeds, rows sit in pg_queue_message), so jobs hang in
DISPATCHED with no error at the producer.
Finding 1 - celery_executor_agentic_kv had no consumer. Three sites:
* docker-compose.yaml: add it to worker-pg-executor's
WORKER_PG_QUEUE_CONSUMER_QUEUE.
* run-worker.sh: add it to the pg-executor role's queue list, for
host-run fleets.
* tests/compose: the e2e override set CELERY_QUEUES_EXECUTOR, which is
read only by workers/executor/worker.py -- the disabled Celery worker.
Re-pointed at WORKER_PG_QUEUE_CONSUMER_QUEUE, the variable the running
pg-queue-consumer actually reads.
The cloud chart already wired workerPgExecutor correctly (2abc96c4); no
change needed there.
Finding 4 - the sandbox consumed sandbox_codegen over the retired broker.
Added a pg-sandbox PG_CONSUMER_ROLES role ("sandbox;sandbox_codegen") and
switched the compose service to pg-queue-consumer, dropping the rabbitmq
dependency. It loads the same workers/sandbox/tasks.py via
WORKER_PG_QUEUE_CONSUMER_WORKER_TYPE, so every spec §6.3 layer is enforced
identically -- only the transport changed. All hardening is preserved:
read-only rootfs, cap_drop ALL, no-new-privileges, tmpfs /tmp, no published
ports, no host_gateway.
Adds workers/tests/test_queue_consumer_wiring.py: the OSS analogue of the
chart's validate-pg-worker-fleet guard. It asserts both queues have a
consumer in the variable the live consumer reads, that the dead Celery
variable is not reused, and that the sandbox's pod hardening survives --
_pg-worker.tpl cannot express any of it, so a future move onto the shared
template would otherwise drop it unnoticed. Each assertion was verified to
fail against the pre-fix state.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
…ments are
A submit dispatches paid work (LLM tokens + OCR pages) but had no billing
gate: the pre-dispatch guards were capability probe (501), per-key rate
(429), schema caps and concurrency (429). An organization whose trial had
expired, or whose subscription was inactive, could keep submitting.
Cloud gates API deployments with SubscriptionMiddleware, which resolves the
org from the URL (/deployment/api/{org_name}/...). Agent-KV cannot be gated
by it as-is: its URL carries no org segment -- the org lives in the Bearer
AgentKVKey, which only the view validates -- so the middleware's
UserSessionUtils.get_organization_id returns None for these requests,
get_subscription(None) finds no row, and verify_subscription reads that as
"nothing to enforce" and admits every one of them.
SubmitView therefore calls the cloud plugin's gate straight after key
validation, where the org IS known. The policy is NOT restated: the gate
calls the same SubscriptionHelper.get_subscription + verify_subscription the
middleware calls, so the two paths cannot drift and the 402 bodies are
identical. Absent service_class (a cloud build predating the gate) degrades
to previous behaviour; OSS-only deployments 501 at the capability probe
before reaching it.
The identifier is the subtle part and has its own test:
`agent_kv_key.organization.organization_id` -- the org SLUG that the
CharField Subscription.organization_id is keyed on, and what the deployment
URL supplies as org_name -- NOT `agent_kv_key.organization_id`, the
Organization FK primary key. The pk matches no row, which the policy reads as
"no subscription" and admits, leaving a gate that never fires while looking
correctly wired. Both that and "denial is returned verbatim" were verified to
fail when the respective mistake is reintroduced.
Scope, recorded in design §6.6: this follows what Unstract actually does
rather than that section's original fail-closed spend reserve. There is no
budget reservation (the platform has none anywhere), and an org with no
subscription row is admitted -- matching the middleware exactly, and
differing from validate_etl_run which denies. Both are properties of the
shared policy; tightening them should tighten every caller at once.
Also corrects the submit example's "status": the view returns
job.status.lower(), so it is "dispatched", not "DISPATCHED".
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
…,options}]
BREAKING: the flat submit shape is removed with no alias.
The v1 contract put everything at top level. That was fine with one extractor,
but `qa`, `challenge`, `extraction_mode`, `structured_output`, `calculations`,
`document_class` and `key_notes` are KV-extractor knobs -- meaningless to any
other extractor -- and there was no way to express "run these extractors over
this document". One result blob and one usage_summary could not say what each
extractor produced or cost. The architecture review (UN-4044) called this a
freeze risk: cheap now, a v2 endpoint once anyone codes against it.
Submit now carries `extractors: [{name, keys, options}]`; only fields describing
the REQUEST stay top level. page_start/page_end are deliberately job-level: the
range drives the shared OCR pass and the §6.1 page cap, so it cannot differ
between extractors reading the same document. Status and result are keyed by
extractor, and usage gains `by_extractor` alongside the authoritative `total`.
`status` stays top level -- it is the job's, and a job is not done until every
extractor is.
Scope is the FORMAT only. `kv` is the sole accepted name and a second entry is
a 400: accepting two and running one would bill a caller for a job that ignored
half the request. Unknown options are a 400 too rather than DRF's default
silent drop, so a knob aimed at the wrong extractor fails loudly.
No cloud change is required. OSS unpacks the entry into the frozen
executor_params contract (`schema` + `options`) with exactly the option names
the engine already receives, and the result/status namespacing is applied at
the API edge -- read-time for results, so already-stored blobs are untouched.
Deliberate divergence from APS v2, decided by Arun and recorded in spec §7.0:
that spec composes IN THE SCHEMA (one mandatory PRIMARY plus SPECIALIZED agents
substituted into a placeholder in one shared project-level schema); this API
treats extractors as peers, each with its own schema and result, because it has
no project and every call is standalone. APS v2 is NOT being changed to match.
The cost -- two composition models, and a caller wanting a table folded into an
invoice record must stitch it -- is stated in §7.0 rather than discovered later.
No alias for the old shape: the API has no consumers (no PR merged, no
customers), so a dual path would cost two validation routes, precedence rules
and tests for all of it, to protect nobody.
Tests: 144 pass. Serializer, view, status and result suites converted, plus new
coverage for unknown extractor, multiple extractors, non-list/empty
`extractors`, unknown options, and that the old flat form is rejected. The e2e
lane is converted (submit building and result unwrapping centralised in its
conftest) but NOT executed -- it needs the Docker stack, which is torn down.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
…ge moves Every existing submit-view test patches SubmitSerializer and feeds the view a `validated_data` built by this file's own helper, while the serializer suite validates payloads and never runs a view. Both suites pass even when the two disagree about the shape that crosses between them -- and that shape is exactly what the extractor-scoped format (§7.0) changed. 144 green tests said nothing about it. Adds two integration tests that build a real multipart upload and run it through the real serializer AND the real view, asserting what arrives at dispatch_job (the frozen OSS<->cloud contract on the far side): the extractor's own schema, its options with defaults filled in, and the job-level page range folded into options where the engine reads it from. The second asserts the hard switch end-to-end -- a caller still sending the pre-§7.0 flat shape gets a 400 from the real stack, not a job that quietly ran with defaults. Verified to have teeth: the first fails if the view reads the wrong key off the entry, and again if it forgets to fold the page range into options; the unknown-option test fails if the serializer goes back to DRF's silent drop. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Both were invisible to 146 passing tests and made the branch non-functional. 1. The subscription gate dereferenced `agent_kv_key.organization`, which is NULLABLE on the model (DefaultOrganizationMixin fills it from UserContext at save time). Every submit against a key without one raised AttributeError and returned 500. The code this replaced used `agent_kv_key.organization_id` and never dereferenced, so the crash path was new. Refused once at the auth boundary instead: every downstream use of a key is org-scoped -- the gate, the concurrency limiter, the storage prefix, every job lookup -- so a key with no organization is not usable and should not reach a view at all. Adds tests/_factories.py::kv_key() because the guard invalidates the bare `AgentKVKey(name=...)` fixture 32 tests used: that key is one the auth layer would reject, so anything asserted past it was unreachable in production anyway. 2. `dispatch.py` called `get_executor_dispatcher(celery_app=...)`, but UN-4046 removed that parameter when the routing dispatcher's Celery branch went with the pg_queue_enabled flag -- its docstring notes that keeping an ignored parameter is what previously "left a `headers=` argument at three call sites and broke every extraction". Every submit raised TypeError inside dispatch_job and failed the job. The architecture review checked PgExecutionDispatcher's signature and cleared it; this is a different function, and the mismatch is invisible to a mocked dispatcher. The guard for it patches the real factory with autospec=True, so our call site is asserted against the live signature rather than a permissive Mock -- verified to fail when the kwarg is put back. 147 tests pass. Found by building the stack and driving POST /agent-kv/ with curl; neither is reachable through the suite as it was written. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
A job ran to success in the executor and then sat in RUNNING forever. The executor enqueues its terminal continuation onto `agent_kv_callback` -- that callback is what persists the result, deletes the staged input, releases the concurrency slot and fires the webhook -- and nothing drained the queue, so the row sat in pg_queue_message with no error at the producer. Same shape as the executor-queue finding, one queue over: the Celery ide_callback worker has carried both queues since the callbacks landed, and the cloud chart wires its PG twin correctly (workerPgIdeCallback: "ide_callback,agent_kv_callback"), but the OSS dev compose and run-worker.sh's pg-ide-callback role listed only `ide_callback`. So only the OSS stacks were affected -- which is precisely why no chart test caught it. Extends the wiring guard to the callback queue at both sites. Found by running a real job end to end; the job's own success is what makes this invisible otherwise, since everything up to finalize behaves normally. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Self-review of the wire-format change: the serializer's accepted-name tuple and the constant the status/result documents key their payloads by were two separate literals spelling "kv". Changing one would let the serializer accept a name the responses file under something else -- the two would disagree with nothing failing. Derived from the single constant instead. 148 tests pass. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
Self-review of the e2e conversion (the one part of this change that cannot be executed here without the stack): with `keys=None` the helper omits `extractors` entirely, so any per-extractor options the caller passed had nowhere to go and were silently discarded. No call site does this today, but such a test would look like it exercised a knob it never sent -- and would still pass, since these submits are expected to fail on the missing field. Asserts instead, naming the offending fields and both ways out. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
An automated review of the unpushed range found five real issues, all in the
parts the unit suites structurally cannot see -- the e2e lane and the response
contract.
Contract:
- The COMPLETED result payload had lost its top-level `success`/`status` (they
moved inside `extractors.kv`), while FAILED and CANCELLED kept them. A client
could no longer branch on one key: on success `body["success"]` was a
KeyError, on failure it was False. Restored at top level on every terminal
payload -- they describe the JOB, and each extractor keeps its own `success`
inside its own block. This also hit the sync-`timeout` submit response, where
the caller has no status document to disambiguate with.
Validation:
- `ExtractorSerializer` rejected unknown keys inside `options` but not on the
entry itself, so `"option"` (or `"Options"`) for `"options"` was dropped whole
by DRF: `options` defaulted to {}, and the job ran with qa=True/challenge=True
-- roughly double the LLM spend asked for, returned as a 202 with no sign
anything was ignored. Now rejected, like the options block one level down.
- The per-extractor schema cap measured `json.dumps(spec)` with the default
ensure_ascii=True, counting a CJK character as a 6-byte \uXXXX escape rather
than its 3-byte UTF-8 form. A non-ASCII schema could clear the outer
`extractors` raw-byte cap and then be rejected by the inner one as "too
large". Both caps now measure the same bytes.
e2e lane (converted but not executable here, so these were live-only failures):
- Two sites read `status_doc["stages"]`, which moved under `extractors.kv`;
added a `kv_stages()` helper beside `kv_result_body()`.
- The re-read stability assertion compared the full wrapped payload against the
unwrapped kv body -- introduced by this branch, which changed the first side
and not the second. Now compares like for like.
- The sync-wait test read the engine's `record` off the wrapped payload.
(Its `success` check, and the calculations test's, work again via the
top-level restoration above.)
Docs:
- `result_payload`'s docstring still claimed COMPLETED returns the stored engine
result "unchanged", which the wrapping directly contradicts -- in the module
that declares itself the source of truth for the result shape.
- Spec §7.3 and API reference §5 updated for the top-level fields.
150 tests pass.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
The wire format already carried per-extractor entries; this makes the name mean something. Adds the table options/keys serializers, an EXTRACTOR_ROUTES table dispatch reads instead of hardcoding one executor/operation pair, and an `extractor` column so status and result key by what actually ran -- they were keyed by the hardcoded `kv` name, which would have filed a table job's output under the wrong extractor. Stage names become per-extractor too: StageReportView persists any name the executor sends, but the status document filtered through the KV-only list, so table stages would have been recorded and then silently dropped. Also updates test_dispatch.py's five dispatch_job callers to pass extractor=V1_EXTRACTOR_NAME (no default was added to the parameter, by design), and fixes two pre-existing tests that assumed `table` was an unknown extractor / that used a kv keys fixture whose own schema failure was masking what the test intended to check. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
1. The drift guard for the three per-extractor lookup tables (EXTRACTOR_ROUTES, STAGE_NAMES_BY_EXTRACTOR, _OPTIONS_SERIALIZERS) only asserted two of the three, while the comment claimed all three were kept in step. Forgetting an options serializer for a future extractor would have stayed green until a real submit hit an uncaught KeyError instead of a validation error. Added the missing assertion and renamed the test to describe what it now checks. 2. Added a submit-to-dispatch join test for the table path (test_real_serializer_through_real_view_reaches_dispatch_intact_for_table), mirroring the existing kv-path version. This guards the exact contract that already broke once in this project: a consumer reading `target_table` off the top level while the producer nested it under `schema` -- internally consistent on both sides, so unit tests on either side alone would not have caught it. Verified the new test is load-bearing by temporarily breaking the join and confirming it fails. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
The table queue now carries paid API work as well as IDE prompts, so it joins the consumer-wiring guard -- an unwired queue accepts work and drains nothing, with no error at the producer. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
…I doc
The doc is this API's contract, and it still described a kv-only
service: a caller reading it concluded the table extractor did not
exist. §7.1's name row said "Currently `kv` only; anything else => 400",
the §4 stage list omitted `table_extraction`, there was no `table`
options table beside the `kv` one, and nothing documented the table
result shape.
Worst of it: nothing documented that `enable_header_mapping` changes
`output.tables` from a list of rows to `{"header_mapping": ...,
"rows": [...]}`. That is a breaking shape change between two settings of
one option, and it was recorded only in a serializer docstring, where no
API consumer would ever see it.
Also corrects §7 (Cancel), which claimed a running job "continues to
completion in the worker ... No mid-run abort". That was already stale
for `kv` -- its node listener raises `JobCancelled` and stops the run --
and is now stale for `table` too, which gained the same cooperative stop
in the cloud repo. Both extractors are described as what they are:
stopped cooperatively at their engine's progress points, with work
already done still billed.
`_status_document` also stops subscripting `STAGE_NAMES_BY_EXTRACTOR`.
No row can hold an out-of-dict value today, but a retired extractor name
with surviving rows would 500 every `GET status` for those jobs while
`GET result` kept working, since the result payload keys by
`job.extractor` without consulting that table. It degrades to an empty
stage array and a warning instead.
agent_kv 164 -> 165.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TMcuErLyV8mGTdcawbYuBM
|
|
for more information, see https://pre-commit.ci
|
Four findings, all verified against the code before changing anything. 1. Periodic tasks were unregistered (P1). pre-commit.ci's auto-fix in a79e9d6 deleted BOTH side-effect imports from scheduler/tasks.py -- including the pre-existing dashboard-metrics one, live since UN-3796. Reproduced locally: ruff 0.3.4's isort merges the two `from scheduler import ...` statements into one parenthesised statement, which relocates each `# noqa: F401` onto a MEMBER line where it no longer suppresses the statement-level F401; pycln (`all = true`) then removes both. Re-adding them with `# noqa` would be deleted again on the next run, so they are now genuinely referenced via `_PERIODIC_TASK_MODULES`. Verified by re-running both hooks over the fixed file. The test written to catch exactly this (test_scheduler_registers_the_metrics_proxies_under_their_wire_names) imported `scheduler.dashboard_metrics_tasks` directly, which registers the tasks by itself -- so it passed regardless of what scheduler/tasks.py contained, and it passed straight through this regression. Replaced with one that goes through `scheduler.tasks`, what a booting worker loads, and asserts the Celery registry. Mutation-checked against a79e9d6's exact diff: [] vs the 5 expected names. 2. DELETE leaked a concurrency slot (P1). JobStatusView.delete terminalized an in-flight job without releasing its slot, where JobCancelView does -- and the sweep's phase-1 only targets PENDING, never CANCELLED, so nothing reclaimed it until the 6h TTL. Enough deletes exhaust the org's allowance and start rejecting new submits. Released on the `won` branch only, so a lost terminal-state race does not hand a slot back twice. 3. A failed file delete lost its retry handle (P1). delete_job_files logs and continues, but both callers blanked the refs unconditionally -- and TTL cleanup selects candidates by `input_ref > "" OR result_ref > ""`, so a blanked row drops out of the candidate set permanently and the object is orphaned in the bucket. delete_job_files now returns the ref fields confirmed to point at nothing (already-blank and FileNotFoundError count as clear; any other error does not), and both callers blank only those. run_ttl_cleanup also reports `retained`, so a backlog that will never drain is visible instead of silently counting as cleaned. Note test (7) in test_sweeps.py already asserted this invariant in its comment -- "a delete failure must not blank a ref pointing at a file that's still there" -- while only pinning call ORDER, which is irrelevant when the update blanks both refs regardless. Now actually covered, at both the storage layer and both call sites. 4. Webhook egress moved onto the shared guard (unstract.core.network.ssrf), the single place that decides whether a tenant-supplied URL may be dialled, retiring this sink's own `_host_is_public` and the hardcoded CGNAT range SonarCloud flagged. The review's specific claim does not hold and is not what this fixes: on the 3.12 runtime these workers use, `ipaddress` reports is_reserved for 64:ff9b::/96, so the local guard already refused the NAT64 loopback example (checked on 3.12.9). What it got WRONG was the other direction -- it refused 64:ff9b::5db8:d822 too, which embeds a public IPv4 and is a legitimate destination, because is_reserved cannot tell the two apart. The shared guard re-checks the embedded IPv4, and brings rules this sink never had (*.localhost without a resolver, credentials-in-URL, the urlparse/urllib3 host disagreement). The old refusal was also incidental to a version-dependent flag that nothing recorded a dependency on. Not fixed here, deliberately: the compose sandbox's lack of an egress boundary (P2). docker-compose `command:` is CMD, appended to the image ENTRYPOINT, so that service is correct as written; the containment control for the sandbox is the k8s NetworkPolicy, tightened in the cloud PR. The real issue underneath it is that the sandbox holds credentials the generated code can read -- the scrubbed subprocess env does not prevent that, since the child shares the parent's UID and /proc/<ppid>/environ is readable. Needs a least-privilege role, not a compose change. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
`workers/shared/enums/tests/` is local-only development scaffolding and does not belong in this PR. It was swept in by a `git add -A workers/shared` while staging the webhook-guard changes in eba836f — note the giveaway duplicate `tests/tests/` nesting, which is not a layout anything would use deliberately. Untracked only; the files stay on disk. Nothing in the PR referenced them, so no behaviour changes. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Greptile re-review on the previous commit, and a fair hit: retaining the ref of a file whose delete failed is what makes a retry possible, but under the plain oldest-first ordering those same rows refill the 500-row batch on every tick. A persistent object-store fault on 500 expired jobs stalled cleanup outright, and every later expired job kept its files past TTL. I had spotted this while writing that fix and settled for documenting it in a comment plus a `retained` counter. That was the wrong call, and inconsistently so: on the cloud PR I argued in the same review round that a documented blast radius is not a control. Same standard applies here. Candidates are now ordered (cleanup_failed_at NULLS FIRST, expires_at): a job that has never failed is always processed ahead of one that has, so failures cannot block fresh work, and they are still retried once the backlog clears. The stamp is refreshed on every failed attempt, so among failures the oldest failure goes first and they rotate rather than one row absorbing every retry. Cleared on recovery so the column means "currently failing", not "failed once". The stamp also has to be written when NOTHING was deleted — skipping that write would leave the column NULL, and NULL sorts first, so the row would hold the head of every batch forever. Covered by its own test. New field + composite index in 0003. A plain AddIndex is correct here, not the CONCURRENT builds used elsewhere in this codebase: agent_kv_job is created by this app's own 0001_initial and has never been deployed, so the table is empty when this runs. `makemigrations --check` reports no drift. The migration itself is exercised by CI's migrate step, not locally — there is no Postgres on this machine right now, and these tests mock the ORM. Ordering and both stamping paths mutation-checked against the previous behaviour: 3 tests fail without them. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
|
| candidates = ( | ||
| AgentKVJob.objects.filter(expires_at__lt=timezone.now()) | ||
| .filter(Q(input_ref__gt="") | Q(result_ref__gt="")) | ||
| .order_by(F("cleanup_failed_at").asc(nulls_first=True), "expires_at")[ |
There was a problem hiding this comment.
If at least 500 never-attempted jobs expire between cleanup runs, they fill every batch because rows with no failure timestamp sort first. Jobs whose file deletion failed are then never retried, so their retained files can remain in storage indefinitely. Reserve some cleanup capacity for retries even when new expirations keep arriving.
Prompt To Fix With AI
This is a comment left during a code review.
Path: backend/agent_kv/maintenance.py
Line: 167
Comment:
**Failed deletions can starve**
If at least 500 never-attempted jobs expire between cleanup runs, they fill every batch because rows with no failure timestamp sort first. Jobs whose file deletion failed are then never retried, so their retained files can remain in storage indefinitely. Reserve some cleanup capacity for retries even when new expirations keep arriving.
---
For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.| # TTL cleanup's candidate ordering is (cleanup_failed_at NULLS | ||
| # FIRST, expires_at); without this the sort is a filesort over | ||
| # every expired row each tick. | ||
| models.Index(fields=["cleanup_failed_at", "expires_at"]), |
There was a problem hiding this comment.
Index cannot serve cleanup order
The new index uses PostgreSQL’s default ascending null order, but cleanup requests cleanup_failed_at ASC NULLS FIRST. PostgreSQL therefore cannot use this index to return rows in cleanup order and must sort the matching expired jobs before applying the 500-row limit. That adds growing work to each cleanup run when there is a large backlog.
Prompt To Fix With AI
This is a comment left during a code review.
Path: backend/agent_kv/models.py
Line: 101
Comment:
**Index cannot serve cleanup order**
The new index uses PostgreSQL’s default ascending null order, but cleanup requests `cleanup_failed_at ASC NULLS FIRST`. PostgreSQL therefore cannot use this index to return rows in cleanup order and must sort the matching expired jobs before applying the 500-row limit. That adds growing work to each cleanup run when there is a large backlog.
---
For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.Note: If this suggestion doesn't match your team's coding style, reply to this and let me know. I'll remember it for next time!
Unstract test resultsPer-group results
Critical paths
|
vishnuszipstack
left a comment
There was a problem hiding this comment.
Automated multi-agent review — findings
This review was produced by a multi-agent code review of #2309 and its cloud counterpart Zipstack/unstract-cloud#1816. Each finding below was re-verified against the code (several by execution). Grouped as must-fix before merge, should-fix, and suggestions. Tenant isolation, mark_terminal CAS, the subscription gate, the TTL-starvation fix (cb6f7f2), and the no-local-execution-fallback property were all checked and confirmed correct — not repeated here.
Merge order: this PR must merge before the cloud half; the cloud plugin imports unstract-agent-kv-schema from this branch and its CI is red until this lands.
Inline comments follow.
| # row this UPDATE doesn't match is left exactly as the winning writer | ||
| # left it. `modified_at` is stamped automatically by | ||
| # BaseModelQuerySet.update() (utils/models/base_model.py). | ||
| AgentKVJob.objects.filter(id=job.id, status=JobStatus.PENDING).update( |
There was a problem hiding this comment.
[Must-fix] Stranded-job race no sweep can recover. This guarded UPDATE persists dispatched_at only for a still-PENDING row, but StageReportView (internal_views.py:112) promotes PENDING→RUNNING on the executor's first stage report. If that report wins the race, this UPDATE matches 0 rows and dispatched_at stays NULL. Neither sweep phase can then see the row: phase 1 requires status=PENDING; phase 2 requires dispatched_at__lt=cutoff, and SQL NULL < x is never true. The job reports running forever and GET result 409s for the life of the row. Fix: stamp dispatched_at for any non-terminal row, or add a dispatched_at__isnull=True arm to sweep phase 2.
| _SLOT_TTL_SECONDS, | ||
| ) | ||
| return bool(int(acquired)) | ||
| except Exception: |
There was a problem hiding this comment.
[Must-fix] Both limiters fail OPEN. check_and_acquire here and check_key_rate (:79) both return True on any Redis error. A Sentinel failover or pool exhaustion removes the concurrency ceiling and the per-key rate ceiling simultaneously, and the API keeps returning 202s for billable LLM work — with only a logger.warning per request. The Lua script itself is correct; it's the catch that's the defect. Fail closed (429), or gate the choice behind a setting so the waiver is deliberate and observable.
| return value.strip().casefold() in {e.casefold() for e in spec.enum_values} | ||
| if fmt == "regex": | ||
| try: | ||
| return re.fullmatch(spec.regex_pattern, value.strip()) is not None |
There was a problem hiding this comment.
[Must-fix] ReDoS via author-supplied regex. compile_schema caps the pattern length (200) but never compiles it, and this runs it with no time budget. Measured on this code: format: "regex:^(a+)+$" (7 chars) takes ~7s on 28 as and ~4x per added char — a 40-char value runs for hours. validate_format runs per key per value in the QA pass, with no timeout/cancellation, and AGENT_KV_CONCURRENT_LIMIT=5 lets one org pin five shared worker slots. Fix: re.compile() at submit (raise SchemaError on re.error) and run user patterns under a hard wall-clock budget (or re2). The length cap is not a mitigation — ReDoS needs ~10 chars.
| # above -- so a cancelled job's input intentionally | ||
| # rides the normal TTL sweep instead of being deleted | ||
| # here. | ||
| delete_input(job) |
There was a problem hiding this comment.
[Must-fix] Failed input-delete still blanks the ref. delete_input(job) swallows every exception (storage.py), and this update(input_ref="") runs unconditionally. TTL cleanup selects only rows with a non-blank ref, so a transient object-store error here orphans the customer's uploaded document permanently. This is exactly the bug the PR fixed for delete_job_files + its two callers (Greptile #3) — this third site was missed. Make delete_input return a confirmed-cleared signal and blank only on success; log with exc_info=True.
| sandbox-escape primitives, and the bare name itself can be aliased — | ||
| ``x = sys`` — to reach any of that through a different name, so any Load | ||
| reference to ``sys`` other than the immediate ``sys.argv`` attribute access is | ||
| rejected, not just a denylist of attributes) and ``open()`` on the runner's |
There was a problem hiding this comment.
[Must-fix: correct the stated security posture] The AST gate is fully bypassable. Verified by execution under the production python -I -S -E: statistics.random._os.system('id') ran id; collections._sys is the sys module; collections._sys._getframe(0).f_builtins['__imp'+'ort__']('os').system(...) reaches arbitrary import+exec — all GATE-PASS. Allowlisted modules re-export sys/os as non-dunder attributes the gate never inspects, and the dunder-subscript check only catches constant strings (so '__imp'+'ort__' slips through). This does not by itself break the model (layers 2–5 still contain it), but the module/PR text claiming the gate 'closes aliasing / from-sys-import / dunder subscripts' overstates it to ~zero marginal value. Represent the gate as best-effort and lean the design explicitly on layers 2–5.
| text=True, | ||
| preexec_fn=_limits(cpu_seconds, memory_mb, max_pids, max_output_bytes), | ||
| ) | ||
| stdout, stderr = proc.communicate(timeout=timeout) |
There was a problem hiding this comment.
[Should-fix] Unbounded capture can OOM the worker. communicate() buffers the full child stdout+stderr in the parent's memory; the [:_CAPTURE_LIMIT] truncation happens after. RLIMIT_FSIZE doesn't apply to pipes and RLIMIT_AS bounds only the child, so while True: print('A'*100000) (no imports, passes the gate trivially) streams GB into the parent and can OOM the Celery worker before the timeout fires — a gate-agnostic DoS. Cap captured bytes during the read, not after. Also: both broad except blocks here (:164, :229) log nothing, so an OSError/fork-EAGAIN is reported to the engine as a codegen failure.
| f"max is {settings.AGENT_KV_MAX_PAGES} (§6.1)" | ||
| } | ||
| ) | ||
| elif ext in IMAGE_LIKE: |
There was a problem hiding this comment.
[Must-fix — cross-repo contract] Image uploads are accepted here but extract nothing in the engine. ALLOWED_EXTENSIONS accepts .png/.jpg/.jpeg/.tiff and this sets pages_total=1, so the job dispatches normally. But the cloud engine's _build_agent_graph (agentic_kv kv_extractor.py) only treats .pdf/.xlsx/.xls as a document, and ImageLoader.load_pages has no call site anywhere — so an image returns success:true, every key not-found, 1 page billed. Either drop the image extensions here, or wire images through ImageLoader.load_pages on the engine side. The two repos currently disagree about what file types the product supports, and it fails open (billed empty result).
| tier: unit | ||
| workdir: workers | ||
| paths: [tests, shared/tests] | ||
| paths: [tests, shared/tests, sandbox/tests, shared/infrastructure/config/tests] |
There was a problem hiding this comment.
[Must-fix] Three new test dirs are in no CI group. This line adds sandbox/tests and shared/infrastructure/config/tests to unit-workers but omits ide_callback/tests (401-line callback suite). Separately, no group covers unstract/agent-kv-schema/tests or unstract/filesystem/tests. So test_compile.py — the regression net for the schema ReDoS/depth findings — and the callback suite run nowhere in CI, and the PR's '165 passed' is a local figure, not CI-verified. Add unit-agent-kv-schema + unit-filesystem groups and append ide_callback/tests here.
| The sandbox receives only `{record JSON + generated code}` — no document, no schema, no | ||
| tenant identifiers; output is validated JSONL, row- and byte-capped. | ||
|
|
||
| **Accepted v1 residual (R12):** generated code legitimately needs the `open()` builtin |
There was a problem hiding this comment.
[Must-fix before sign-off] The documented risk acceptance is factually wrong. This R12 text describes the residual as 'a single read of a known path … non-sensitive … not an exfiltration path.' Verified false by execution: generated code open()s /proc/<ppid>/environ (child runs as the same UID as the worker) and recovers the worker's full environment incl. DB_PASSWORD; with the sandbox pod mounting the database secret group and the NetworkPolicy allowing egress to db-proxy:5432 (cloud half), that credential is directly usable. So the residual is credential disclosure + DB reach beyond the job, not a benign file read. This is the text an admin signs off against — correcting it is a prerequisite independent of when the hardening (numeric runAsUser / scoped DB role / hidepid) lands. Same paragraph is duplicated in the sandbox design spec.



What
The Agent-KV extraction engine's OSS half: the
agent_kvDjango app (public API, auth, dispatch, storage, rate limiting, job lifecycle), theunstract-agent-kv-schemaworkspace package, the hardened code sandbox worker, and the compose/test wiring for both.Built by @Arun; merged up to current
mainand raised as part of taking the work forward. Pairs with the cloud half — see Related.Why
Agent-KV is a multi-stage key/value extraction engine. Agentic Prompt Studio v2 already declares it as the Multi-Stage Extractor —
MULTI_STAGE → executor="agentic_kv", operation="kv_extract"is onmaintoday withavailable=False. Merging this is what lets that flag flip; without it, M2 would have to rebuild the same pipeline as a fourth copy of an engine that already exists.Reference: Agent-KV — Architecture Review (UN-4044).
How
Size: 116 files, ~16.1k lines, including the
agent_kvapp, its 16 test modules, and the schema-compiler package.Security posture — generated Python runs behind five independent layers, and the AST gate is deliberately not the boundary (layers 2–5 are). With no transport configured the engine errors rather than executing locally; there is no fallback path:
from sys import, dunder subscriptspython -I -S -E,start_new_session, five rlimits, CPU backstop above the wall clockRuntimeDefaultseccomp, no service-account token[]; egress only broker, db-proxy, DNSENCRYPTION_KEY; no tenant id reaches the podThe accepted v1 residual is narrow: generated code can
open()a world-readable file in its own pod and return it to the caller who submitted the job.pathlib,os,socketandurllibare outside the allowlist, so there is no path to a secret, a cross-tenant read, or exfiltration. Layers 4–5 are what make that true and must not be relaxed — weakening either turns a local file read into an exfiltration path.Review findings — state at this head
The architecture review (6 Sep) raised four. Re-verified against this branch, post-merge:
DISPATCHEDsilentlycelery_executor_agentic_kvis in the executor role's queue list;tests/test_queue_consumer_wiring.py(13 tests) guards itmaindeletedworker-sandboxis defined by this branchpg-sandboxonWORKER_PG_QUEUE_CONSUMER_QUEUE: sandbox_codegen, with fleet-guard wiringFinding 3 is knowingly out of scope here
The review recommends moving the KV knobs under
extractors: [{name, options}], keying the result document by extractor, and promotingusage_summaryto per-extractor, with today's flat fields kept as a deprecated alias. Anextractorsfield exists (execution_serializers.py:209) butqa/challenge/extraction_mode/calculationsare still top-level.Shipping first is deliberate: nothing consumes this API yet, APS v2 keeps Multi-Stage dark behind
available=False, and a wire-format redesign tangled into a 16k-line PR is harder to review than either change alone. The review's own step 4 is "fix findings 1–2 during the rebase; open both PRs" — 1, 2 and 4 are done.Note the review's ordering constraint ("must precede M0/M0b merging") is already moot: M0/M0b merged on 30 Sep, and did so having already adopted the review's §5.3 recommendation — the registry points at
agentic_kv, not at a rebuilt pipeline.Can this PR break any existing features. If yes, please list possible items. If no, please explain why. (PS: Admins do not merge the PR without this section filled)
Low. The work is additive — a new Django app, a new workspace package, a new worker and its compose/test entries. No existing endpoint changes shape and no existing model is altered.
Two things a reviewer should check rather than take on trust:
test_queue_consumer_wiring.pyis the guard.Database Migrations
agent_kvapp migrations (new tables only; no alterations to existing tables).Env Config
AGENT_KV_STORAGE_DIR_PREFIX(defaults tounstract/agent_kv)AGENT_KV_CALCULATIONS_ENABLED— defaults false on absence, so OSS/on-prem installs stay flag-offNotes on Testing
cd backend && pytest agent_kv/tests --no-migrations— 165 passedcd workers && pytest tests/test_queue_consumer_wiring.py— 13 passedtests/compose/docker-compose.test.yamlRelated
UN-4044-agent-kv-cloud-executor(raised alongside this)Checklist
I have read and understood the Contribution Guidelines.
🤖 Generated with Claude Code