Repository navigation
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 101 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
|
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>
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.
…d ref-blanking site Review round 3. Four findings, two of them in code I added last round. **TTL cleanup starved retries (Greptile, on my own fix).** Last round's `cleanup_failed_at NULLS FIRST` ordering fixed failures starving fresh work by inverting it: 500 new expirations per tick now fill every batch and the failures are never retried, so their files sit in storage indefinitely. One ordering cannot express "neither side starves the other", so the batch is split into two lanes -- retries get up to `_TTL_RETRY_RESERVE` (100) slots, fresh rows take the remainder. Each side capped, so each side guaranteed capacity whenever it has work. **The index could not serve that ordering (Greptile).** A btree index is NULLS LAST ascending, so `ASC NULLS FIRST` could not use `(cleanup_failed_at, expires_at)` at all and Postgres sorted every matching expired row before applying the limit -- work growing with the backlog. The split removes the NULLS FIRST entirely: each lane is one ascending column behind an IS NULL / IS NOT NULL predicate on the index's leading column, which the existing index serves directly. No new migration. **Stranded jobs no sweep could recover.** `StageReportView` promotes PENDING -> RUNNING on the executor's first stage report, which can land before dispatch's post-enqueue bookkeeping. The PENDING-guarded UPDATE then matched 0 rows and `dispatched_at` stayed NULL -- and a non-terminal row with a NULL `dispatched_at` was invisible to BOTH sweep phases: phase 1 requires `status=PENDING`, phase 2 filters `dispatched_at__lt` and SQL `NULL < x` is never true. The job reported `running` forever and `GET result` 409'd for the life of the row. Fixed at both ends: dispatch stamps the bookkeeping for any still-non-terminal row (never touching `status` -- the row was dispatched, but moving RUNNING back to DISPATCHED would discard the executor's progress), and sweep phase 2 gained an `OR dispatched_at IS NULL AND created_at < cutoff` arm as a backstop for the window that remains if a worker dies between the enqueue and the stamp. **A third ref-blanking site was missed.** `delete_input` swallowed every exception while `FinalizeView` blanked `input_ref` unconditionally, so a transient object-store error orphaned the customer's uploaded document permanently. Identical to the bug fixed for `delete_job_files` and its two callers last round -- missed because that fix followed one function's callers instead of grepping for every site that blanks a ref. `delete_input` now returns the same confirmed-clear signal and logs with `exc_info=True`. Grepped the app afterwards: no fourth site. Tests: the two-lane split is mutation-checked against the NULLS FIRST version (the flood-of-new-expiries case fails without it), and the ttl test helper was rebuilt out of small real stand-ins -- wrapping a Mock's `__getitem__` to record the slice made it call itself. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
…orrect R12
Three review findings, none of them a code bug in the engine -- two are the
product claiming something it does not do.
**Images were accepted and extracted nothing.** `ALLOWED_EXTENSIONS` took
`.png/.jpg/.jpeg/.tiff` and stamped `pages_total=1`, so an image dispatched
normally. But the cloud engine's `_build_agent_graph` treats only
`.pdf/.xlsx/.xls` as a document -- anything else takes a branch that skips
`document_processor` outright -- and `ImageLoader.load_pages`, the only thing
that would populate pages for an image, has no call site anywhere in the
plugin (verified by grep; the engine's own comment there says images are "out
of P2 scope"). So an image returned `success: true`, every key not-found, and
a page billed. It failed OPEN.
Dropped from the allowlist rather than wired, because wiring it is engine work
and this is the half that stops charging for empty results today. The two repos
now agree on the supported file types. `docs/agent-kv-api.md` updated in the
same change, including the `AGENT_KV_MAX_PAGES` row that described a page cap
"for PDFs/images".
**Three test directories were in no CI group.** `unit-workers` listed
`sandbox/tests` and `shared/infrastructure/config/tests` but omitted
`ide_callback/tests`, and nothing covered `unstract/agent-kv-schema/tests` or
`unstract/filesystem/tests`. 34 tests ran nowhere -- including
`test_compile.py`, which is the regression net for the schema depth/ReDoS caps,
so those caps were only ever asserted on a developer's machine. Any "N passed"
figure quoted for this package was local, not CI-verified. Added
`unit-agent-kv-schema` and `unit-filesystem`, appended `ide_callback/tests`,
and ran all three locally first: 15 + 18 + 1, all green.
**R12's risk acceptance was factually wrong, in both documents that state it.**
It described the residual as "a single read of a known path … non-sensitive …
not an exfiltration path", contained by a layer-3 "no secrets" invariant.
Verified false by execution: the sandboxed child runs as the same UID as the
worker, so `open('/proc/<ppid>/environ')` returns the worker's environment
including `DB_PASSWORD` -- no import needed, so the allowlist bounds nothing --
and the pod does mount the `database` group, because on the PG transport the
queue IS Postgres. The sandbox design spec's own containment list even
contradicted itself, claiming "no secrets" while its parenthetical admitted the
DB credential.
Corrected worst case in both docs: credential disclosure plus database reach
beyond the job's own input. "No secrets" is no longer cited as an invariant for
this pod, the AST gate is described as best-effort rather than containment, and
the scoped DB role is reclassified from deferred to required (UN-4218). This is
the text an admin signs off against, so correcting it is independent of when
the hardening lands.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
…the sandbox capture
Four review findings, three of them security.
**Both limiters failed OPEN.** `check_and_acquire` and `check_key_rate` each
`return True` on any Redis exception, so one Sentinel failover or pool
exhaustion removed the concurrency ceiling AND the per-key rate ceiling at the
same time while the API kept returning 202s for billable LLM work -- with a
per-request `logger.warning` as the only signal. Now fails CLOSED by default,
behind `AGENT_KV_LIMITER_FAIL_OPEN` so the waiver is a deliberate, visible
config choice rather than the implicit behaviour of an `except` block. Callers
already raise `RateLimited`, so this surfaces as a clean 429. Logged with
`logger.exception` -- a limiter being gone is an error, and the old level is
part of why this went unnoticed. `release()` stays best-effort: it runs on
terminal paths where raising would abort the caller's work, and a slot it
cannot free expires on its own TTL.
Replaced `test_redis_error_fails_open`, which pinned the defect as the
contract.
**ReDoS via author-supplied regex.** `compile_schema` capped the pattern's
LENGTH (200) but never compiled it, and `validate_format` ran it per key per
value in the QA pass with no time budget. Reproduced: `^(a+)+$` is 7 characters
and takes 1.88s against 26 `a`s, ~4x per character added -- so a 40-char value
runs for hours, and `AGENT_KV_CONCURRENT_LIMIT=5` lets one org pin five shared
worker slots from a single submit. The length cap is not a mitigation.
Patterns are now compiled at submit (`SchemaError` on `re.error`, instead of
being discovered per-value in the engine where `_check_one` swallows it and
returns True), and a nested-quantifier check refuses the shape behind the
realistic cases. Explicitly a conservative heuristic, documented as such: it
will refuse some safe-but-similar patterns, and it is not a proof. The complete
fix is a linear-time engine (RE2), which cannot go here because this package
deliberately has zero dependencies and both repos install it -- UN-4221.
Verified both directions: four catastrophic shapes refused, four real-world
patterns (`^\\d{3}-\\d{4}$` etc.) still compile.
**Unbounded capture could OOM the worker.** `communicate()` buffered the
child's full stdout+stderr in the PARENT's memory and truncated afterwards, so
`while True: print('A'*10000)` -- no imports, passes the gate trivially --
streamed gigabytes in and could OOM the worker before the timeout fired.
RLIMIT_FSIZE does not apply to pipes and RLIMIT_AS bounds only the child, so
this was gate-agnostic. Output now goes to files in the per-job tempdir, where
the EXISTING RLIMIT_FSIZE does bite, and the parent reads back at most
`_CAPTURE_LIMIT`. Both broad `except` blocks now log; they were silent, so an
OSError or fork EAGAIN reached the engine as an ordinary codegen failure.
**The gate's stated posture overstated it.** Its careful `sys`-aliasing rules
are worth materially less than they read, because allowlisted modules
re-export `sys`/`os` as non-dunder attributes the gate never inspects.
Recorded the three measured bypasses in the module header, and corrected the
Layer-3 line that claimed the pod runs "with no mounted secrets" -- it mounts
the `database` group by necessity. Same false claim as R12; both documents
already corrected in the previous commit.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
…entry-point claim Two review findings on the shared schema package. **Constraint equality was exact float equality**, so the feature's own headline example reported a false violation on a CORRECT invoice: three line items of 8230.40 sum to 24691.199999999997, which `operator.eq` says is not 24691.20 (reproduced). A consistency check that fires on correct documents is worse than no check -- it trains reviewers to ignore the output. `==`/`!=` now use `math.isclose` with both bounds set: `rel_tol=1e-9` (~15 significant figures, far tighter than any extracted value's real precision, and magnitude-independent, which an absolute epsilon cannot be across invoice totals and balance sheets alike) and `abs_tol=1e-9` so comparisons against exact zero still work, where relative tolerance alone is useless. Ordering comparisons stay exact on purpose: a tolerant `<` would make `a < b` and `a == b` both true at the boundary, and "a total that is too large" is not a rounding artefact. This fixes the COMPARISON. The value path is still float end to end -- `coerce_number` returns float and `_fmt_number` stringifies it, so `"1,234,567,890,123,456.78"` loses its cents -- which needs Decimal and is tracked as UN-4222. The tolerance stays correct after that lands. Mutation-checked: two of the five new tests fail against `operator.eq`. **compile.py claimed to be "the single entry point both the API and the cloud engine use".** It is not: the engine calls the raw `kv_schema.compile`/`compile_arrays`, so `SchemaCaps` (max_leaves=200, max_depth=6, max_columns_per_array=40) is enforced on the submit path only. Sound while the backend is the sole producer, but it makes the caps a gate rather than an invariant the engine re-checks, and the docstring hid that. Corrected, with the condition under which it would have to change stated. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
…xist The previous commits referenced UN-4221/UN-4222 in code comments before those tickets were created, and Jira allocated UN-4225 (linear-time regex engine) and UN-4226 (Decimal money path) instead. Writing a key before creating it was the mistake; repointed, and checked there are no other stale references. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Five findings on the previous three commits. All mine; all valid. **A bookkeeping failure could fail an already-queued dispatch.** My fallback UPDATE runs after the enqueue, and an exception from it propagates into `SubmitView`, which turns a DispatchError into a FAILED job -- so the executor runs, calls back, and finds a terminal row it cannot write to, while the caller is told nothing was billed for work that did run. The single pre-review UPDATE had the same exposure; adding a second widened it. Post-enqueue bookkeeping is now wrapped and logged, and losing it entirely is recoverable: sweep phase 1 reaps a still-PENDING row with no `dispatched_at`, and phase 2's new `dispatched_at IS NULL` arm covers the non-PENDING case. **The ReDoS heuristic missed the overlapping-alternation family.** `^(a|aa)+$` passes the nested-quantifier check and still backtracks catastrophically -- at each position the engine can consume one `a` or two and must try both on failure. I had named this family as uncovered on UN-4225 rather than closing it, which was the wrong call for a security gap that cheap to express. Now refused, by the prefix relation over LITERAL branches only: one branch being a prefix of another (`a` of `aa`) makes the group ambiguous. `(foo|bar)+` is untouched -- distinct first characters mean no position admits two parses -- and branches containing metacharacters are not analysed, because that needs real regex analysis, which is still what UN-4225 is for. Seven cases verified either way. **Timed-out children were not reaped.** Moving capture from pipes to files dropped the `communicate(timeout=1)` that incidentally waited after SIGKILL, so the runner read the files while the child might still be writing (truncated diagnostics) and left it unreaped -- a zombie per timed-out job in a long-lived worker. Added an explicit `proc.wait(timeout=5)`. **The tolerance hid a cent at scale.** `rel_tol=1e-9` is 0.1 at a $100,000,000 total, so it absorbed a one-cent reconciliation error -- the exact failure the check exists to catch -- and my test only covered a ~$25,000 total, so it could not see it. Recomputed against both bounds that actually constrain this: the float noise floor (summing N values of magnitude M accumulates ~N*2.2e-16*M, so ~2.2e-5 at $1e8 over 1,000 rows) and a cent. `rel_tol=1e-12` gives 1e-4 at $1e8: above the noise, two orders below a cent. `abs_tol=1e-6` keeps zero comparisons working. Known ceiling past ~$1e12 with ~10,000 rows recorded in the code -- that is a float problem, not a tolerance problem, and it is UN-4226. Three tests at $100M scale, including the noise-absorption floor. **A stale doc line.** `docs/agent-kv-api.md` still said "Images count as 1 page" one sentence after saying images are refused. My own earlier edit had also split that bullet mid-sentence; the whole passage is rewritten. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
CI regression from the fail-closed limiter change, and the tests were resting on the defect. Three tests in test_validate_view.py relied on the limiter's ambient behaviour: it failed OPEN on any Redis error, so with no Redis in the unit lane `check_key_rate` returned True and the tests passed for a reason unrelated to what they assert. Now that it fails closed, an unreachable Redis is a 429 and they became `assert 429 == 200`. They pass locally only because a Redis happens to be listening there. Mocked explicitly, matching test_submit_view.py's existing pattern. That is the right fix independent of the limiter change: these are schema-validation tests and the rate limiter should never have been in their path. Verified against the whole agent_kv suite with Redis unreachable: 189 passed. Audited the other suites that reach a rate-limited view -- test_auth.py and test_job_views.py touch SubmitView, but auth runs before the rate check so their 401 cases short-circuit, which matches CI reporting only these three. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Review raised that the tightened comparison tolerance could report a CORRECT total as a violation on a long array, citing `sum([8230.40] * 100_000) == 823039999.9985306`. That figure is not what CPython computes -- it is exactly `823040000.0`, because for identical addends the partial sums stay representable -- and the concern did not reproduce through the real code path: 60 randomised value/row combinations up to 100,000 rows at ~$100k per line produced zero false positives, while a one-cent error was still caught at every size. Fixed anyway, because "does not reproduce across 60 samples" is a weaker guarantee than "cannot happen". The aggregates now use `math.fsum`, which is exactly rounded, so accumulation error no longer grows with row count at all and the tolerance no longer has to absorb it. One function call, and it removes the variable instead of arguing about its size -- worth it on the third review pass over these same few lines. min/max/count are unchanged (no accumulation); avg switches with sum since it divides one. Four new tests: a correct total at 1k/10k/100k rows is not a violation, and a one-cent error at 100k rows still is. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
The quality gate was already passing, so none of these blocked. Clearing the ones worth clearing and recording why the rest are left. Fixed (29): - **S5778 ×15** — `pytest.raises` blocks containing more than one call that could throw, so a pass did not prove the function under test raised. Fixture construction hoisted out of every block (`SchemaCaps(...)`, `_nest(...)`, `_schema(...)`, `_job()`, `_wrapped()`/`_request(...)`), leaving only the call being tested inside. A real improvement, not a lint appeasement: nine of these are depth/leaf-cap tests where a `SchemaError` from building the fixture would have looked identical to one from the compiler. - **S9073 ×5** — composite assertions split, so a failure says which half broke. `assert bucket and rest` in particular meant two different defects (no bucket segment vs. a bucket with no path) reported identically. - **S8572 ×3** — `logger.error(..., exc_info=True)` -> `logger.exception(...)`, matching what the rest of this PR now does. - **S1192 ×1 (CRITICAL)** — `"Invalid api key"` was duplicated across three raise sites in key_validator. Named once, with a note on why all three deliberately say the same thing: a caller holding a bad key must not learn whether it was malformed, unknown, inactive or org-less, so the three must not drift apart and start leaking that distinction. - **S7632 ×1** — my own prose in scheduler/tasks.py was parsed as a malformed suppression directive: the comment explaining the pre-commit interaction quoted the literal directive token. Reworded. - **S5799, S7500, S9409, S9083** — one each: two implicitly concatenated f-strings merged (the shape a missing comma hides in), a redundant comprehension inside `b"".join`, consecutive `append()` -> `extend()`, empty parentheses on a `@pytest.fixture`. Not fixed (10), deliberately: - **S3776 ×8, cognitive complexity.** `check_code_safe` (55), `run_code` (28), `kv_schema._walk` (25), `FinalizeView.post` (23), `compile_schema` (20), `_aggregate`/`_truth` (18), `SubmitView.post` (16). All were over the threshold before this review round and all are security- or correctness-critical; splitting them is a refactor with real regression risk, which is not what belongs in a review round on a PR about to merge. Two I did make worse and will own: `FinalizeView.post` gained a branch from the `delete_input` confirmation fix, and `run_code` gained the capture closure plus the post-SIGKILL wait. Both additions were required fixes to real bugs. - **shelldre:S1192 ×2** — `'sandbox'` as a literal in `run-worker.sh` / `run-worker-docker.sh`. It is a `case` label there; a shell variable is not usable as a case pattern without making the dispatch harder to read. No SonarCloud analysis exists for the cloud repo (no check on #1816, and the project key returns nothing), so this is OSS-only. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
…still inside Correction to the previous commit, which claimed 29 of 39 Sonar issues cleared. The real figure was 26: the three test_auth.py cases survived, and Sonar's re-analysis showed them at shifted line numbers rather than gone. I had hoisted `_wrapped()` and `_request(...)` out of the `pytest.raises` block but left `mock.Mock()` inside, which is itself a call that could throw — so the rule still applied and for the same reason it always did: a pass would not prove the raise came from the view. All three now hoist the instance too. Verified structurally rather than by eye this time: an AST walk over the three edited files confirms every `pytest.raises` block contains exactly one call. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Was deferred as UN-4220; doing it here instead. Collapsing rows identical in every extracted cell is not lossless, and the docstring claiming it was has already been corrected. A document that genuinely contains two identical line items came back with one, so the row count and any `_constraints` aggregate over those rows were wrong -- and the codegen path consumes exactly these rows. `"_dedup": false` on an array node now skips it (`ArraySpec.dedup_rows`, added to RESERVED_NODE alongside `_key`). Default stays `true`: removing it inflates counts on the replicate-column OCR corpus the extractor was built against, so both behaviours are wrong for some documents and flipping the default would trade a known-wrong case for an untested one. What changes is that an author who knows their documents can now say so, which is the part that was missing. Validated as a real boolean rather than truthily: `"_dedup": "false"` is a non-empty string and would read as TRUE under a truthiness check -- silently leaving dedup on for someone who explicitly asked for it off, and silently dropping their rows. That is a schema error now, not a default. Documented in docs/agent-kv-api.md next to `_key`, including why the default is `true` and that it is not the safer value for every document. Engine side and tests in the cloud commit. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Was part of the deferred UN-4218; doing it here instead. `groupadd -r worker && useradd -r -g worker worker` takes the next free system id descending from 999. On this base (python:3.12.9-slim) that is 999 today -- verified by running the same two commands in it -- so pinning to 999 changes nothing for existing images and matches any volume ownership already assuming it. What it buys is determinism. A future base image shipping one more system user would silently shift this to 998, and any manifest or volume ownership assuming 999 would break on a rebuild with no code change. It is also what lets a pod set a numeric `runAsUser` honestly rather than guessing: the cloud sandbox worker now asserts 999, so `runAsNonRoot: true` no longer has to depend on this file keeping its `USER` line. I declined this during review saying there was "no correct number to write yet". That was right about not guessing and wrong about the conclusion -- the number was determinable and I had not determined it. Two hadolint findings on this file came with it, both pre-existing and both blocking any commit that touches it, so they are fixed here rather than worked around: - DL3066 `USER worker` -> `USER 999`. Equivalent now the uid is pinned, needs no /etc/passwd lookup, and is the id the pod's `runAsUser` must match. - SC2086 an unquoted `$plugin_dir` inside `basename` in the plugin-install loop. A plugin directory with a space would have split there. hadolint now reports only the pre-existing DL3008 (apt version pinning), which is out of scope for this change. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Eight python:S3776 (cognitive complexity) and two shelldre:S1192 (repeated
string literal). The quality gate was already passing; these are the ones
previously left open.
Measured rather than guessed: flake8-cognitive-complexity reproduces
Sonar's numbers on these files (16/23/20/18/26/28 match exactly; it reads
gate.py as 46 where Sonar says 55), so each refactor was checked against
a <=15 ceiling locally instead of waiting for a CI round trip.
S3776 -- each split along a seam the function already had, no behaviour
change intended:
SubmitView.post 16 subscription gate + synchronous-wait poll
JobFinalizeView.post 23 success / failure / input-ref-blanking arms
compile_schema 20 shape caps, per-key caps, constraints
_aggregate 18 call-parts validation, count, numeric
_truth 19 the comparison chain
_walk 26 array node and leaf node arms
run_code 28 output validation, group kill, reap
check_code_safe 46 the nine-rule if/elif chain -> _RULES tuple
check_code_safe is a safety gate and its if/elif ORDER decides which
rejection reason an input reports, so "tests still pass" was not accepted
as evidence. The old and new implementations were diffed directly over a
142-input corpus -- every string literal in test_gate.py plus handcrafted
per-rule and ordering cases -- for zero differences in the full
(ok, reason) tuple, re-confirmed after ruff format. The harness was itself
mutation-checked: dropping _rule_dunder_attribute makes `x.__class__` pass
the gate, and the diff catches it. _RULES carries a comment saying order is
behaviour, because it no longer looks load-bearing the way an elif chain
does.
S1192: both scripts already name worker types they reference repeatedly
(EXECUTOR_WORKER_TYPE, IDE_CALLBACK_WORKER_TYPE) and use them as both
associative-array keys and case labels, so SANDBOX_WORKER_TYPE follows the
existing idiom rather than inventing one. Proven inert by extracting the
`declare -rA` blocks from both revisions and comparing the resolved arrays
-- byte-identical in both files -- plus `bash -n` and an explicit check
that `"${SANDBOX_WORKER_TYPE}")` matches `sandbox` and nothing else.
Added two tests asserting every _rule_* function is registered in _RULES
and in the documented order -- a detached rule is a failure mode the elif
chain could not have. Their docstring reports a measurement rather than an
assumption: dropping any one of the eight rules fails between 3 and 10 of
the existing behavioural tests, so these guard the rules that do not exist
yet, not a gap today. (First draft claimed no other test would notice;
that was wrong and is corrected.)
Suites: schema 42, sandbox 80, backend agent_kv 189, workers 1379.
unstract/core is untouched and its local harness cannot resolve its own
package here -- covered by CI's unit-core group.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
…ext form Not an agent-kv change. This test is main's (UN-3883, #2276), untouched by this PR, and it failed this PR's integration lane: django.db.utils.DataError: invalid input syntax for type interval: "9.59465953354055e-06 days" The fixture spread 4000 rows' created_at with `(random() * 60 || ' days')::interval`. `random()` is float8 and is evaluated PER ROW, so each row formats its own value as text -- and float8's text form switches to exponential notation below 1e-4. A single row drawing random() < 1.67e-6 yields '9.59e-06 days', which interval cannot parse, and the DataError aborts the whole UPDATE and so the class's setUpTestData. Measured rather than assumed, against postgres:15: - the observed CI value reproduces the exact error string, but only when typed float8 -- written as a bare literal it is numeric, which prints in decimal and parses fine (my first reproduction attempt, which passed) - `(random() * 60 || ' days')::interval` over 4000 rows: 6 of 400 batches failed, about 1.5% of runs - `(random() * 60) * interval '1 day'` over the same 400 x 4000: 0 failures - the two forms are exactly equal wherever the old one parsed (max abs difference 0.000000 across 1000 points spanning the 60-day range) float8 * interval never round-trips through text, so the failure mode is removed rather than made rarer. Same distribution, same fractional-day precision. Carried here rather than split out because it was the one remaining red lane on this PR and it is a one-line fix with no agent-kv coupling; happy to move it to its own PR if preferred. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
|
Unstract test resultsPer-group results
Critical paths
|



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