Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
4961b05
UN-4232 [FEAT] Agent-KV API serving the table extractor only
vishnuszipstack Oct 6, 2026
d6dc9e4
UN-4232 [FIX] Lock the cross-repo stage name; correct the API docs
vishnuszipstack Oct 6, 2026
7067790
UN-4232 [MISC] Clear the six SonarCloud cognitive-complexity findings
vishnuszipstack Oct 6, 2026
2f44cc4
UN-4232 [FIX] DELETE on a lost terminal race blanked the winner's res…
vishnuszipstack Oct 6, 2026
d01cb17
UN-4232 [FIX] Four review findings: gate fail-open, unscheduled sweep…
vishnuszipstack Oct 6, 2026
66553a1
UN-4232 [FIX] 1.2 — the table API operation had no billing backstop
vishnuszipstack Oct 6, 2026
1280337
UN-4232 [FIX] 2.1, 2.2, 2.4 — double-validated identity, an unrouted …
vishnuszipstack Oct 6, 2026
4a20ef7
UN-4232 [FIX] 2.5, 2.9, 2.11, 2.12 — an over-subscribed concurrency c…
vishnuszipstack Oct 6, 2026
835a46d
UN-4232 [FIX] 2.6, 2.7, 2.8, 2.13, 2.14, 2.15 — correctness group
vishnuszipstack Oct 6, 2026
a1c6c9d
UN-4232 [FIX] Six Greptile findings on this branch's own new code
vishnuszipstack Oct 6, 2026
bd82bc8
UN-4232 [FIX] Webhook delivery results, and two comments that describ…
vishnuszipstack Oct 6, 2026
7efbc4c
UN-4232 [FIX] Address the remaining open review threads on #2317
vishnuszipstack Oct 7, 2026
fbd0ba9
Merge origin/main into feat/agent-kv-api-table
vishnuszipstack Oct 7, 2026
f163650
UN-4232 [FIX] Give agent-kv's `extractor` column choices and no default
vishnuszipstack Oct 7, 2026
e283b9a
UN-4232 [REFACTOR] Delete agent-kv's dead `compiled` serializer attri…
vishnuszipstack Oct 7, 2026
a6df38b
UN-4232 [MISC] Clear the eight SonarCloud findings on #2317
vishnuszipstack Oct 7, 2026
d068e58
UN-4232 [FEAT] The table extractor runs on the caller's adapters, not…
vishnuszipstack Oct 7, 2026
7ee61de
UN-4232 [FIX] Two P1 review findings on the adapter change — one of t…
vishnuszipstack Oct 8, 2026
82d5181
UN-4232 [FIX] Round-2 review: a billing control the gate skipped, two…
vishnuszipstack Oct 8, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Empty file added backend/agent_kv/__init__.py
Empty file.
6 changes: 6 additions & 0 deletions backend/agent_kv/apps.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
from django.apps import AppConfig


class AgentKvConfig(AppConfig):
default_auto_field = "django.db.models.BigAutoField"
name = "agent_kv"
70 changes: 70 additions & 0 deletions backend/agent_kv/constants.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
#: The `kv` extractor's name. Kept as a constant, but NOT routable on this
#: deployment -- see `EXTRACTOR_ROUTES` below, which carries `table` only.
#:
#: An earlier version of this comment said "v1 accepts exactly one extractor
#: (`kv`), so the job row does not carry which one it ran". Both halves are now
#: wrong, on line 1 of the file that defines the routing table: `kv` is the one
#: extractor this deployment does NOT accept, and `AgentKVJob.extractor` has
#: carried the name since migration 0002.
V1_EXTRACTOR_NAME = "kv"

STAGE_NAMES = [
"document_processing",
"extraction",
"qa",
"challenge",
"normalize",
"constraints",
"codegen",
"code_execution",
]
EXECUTOR_NAME = "agentic_kv"
OPERATION_KV_EXTRACT = "kv_extract"
EXECUTION_SOURCE = "agent_kv_api"

#: The table extractor, served by the cloud `agentic_table` plugin's blind-API
#: operation. Same executor as the IDE table path (and therefore the same,
#: already-wired `celery_executor_agentic_table` queue) with a second operation
#: -- a new executor name would derive a new queue needing wiring at five
#: sites, and an unwired queue accepts work and drains nothing, silently.
TABLE_EXTRACTOR_NAME = "table"

TABLE_EXECUTOR_NAME = "agentic_table"
OPERATION_TABLE_EXTRACT_API = "table_extract_api"

#: Which executor and operation each extractor dispatches to. `dispatch_job`
#: reads this rather than hardcoding one pair, so adding an extractor is one
#: entry here plus its options serializer and stage list.
#:
#: **`kv` is deliberately absent.** This deployment ships the `agentic_table`
#: plugin and not `agentic_kv`, so nothing drains `celery_executor_agentic_kv`.
#: `SUPPORTED_EXTRACTORS` is derived from this table's keys
#: (`execution_serializers.py`), so the omission turns a `kv` submit into a 400
#: at the serializer. Restoring the entry below is all it takes to re-enable
#: the extractor once the plugin ships -- and restoring it WITHOUT the plugin
#: is the failure this guards: the submit would be accepted with a 202,
#: dispatched to a queue with no consumer, and sit in DISPATCHED forever with
#: no error at the producer and nothing in any log to find.
#:
#: `V1_EXTRACTOR_NAME`, `STAGE_NAMES` and the KV options serializer stay in the
#: tree, dormant and still under test, so re-enabling is one line rather than a
#: content merge against the branch that carries the engine.
EXTRACTOR_ROUTES = {
TABLE_EXTRACTOR_NAME: (TABLE_EXECUTOR_NAME, OPERATION_TABLE_EXTRACT_API),
}

#: The table engine reports one coarse stage: it has no node-level progress
#: hooks (the IDE path gets `stream_log` only), so inventing finer stages here
#: would describe progress the executor cannot actually report.
TABLE_STAGE_NAMES = ["table_extraction"]

#: Stage names ARE wire format -- they are returned to clients -- and they are
#: extractor-specific (`qa`/`challenge`/`codegen` mean nothing to the table
#: extractor). `_status_document` filters a job's recorded stages through the
#: list for the extractor that ran: `StageReportView` persists whatever name
#: the executor sends, so without a per-extractor list a table job's stages
#: would be stored and then silently filtered out of every status response.
STAGE_NAMES_BY_EXTRACTOR = {
V1_EXTRACTOR_NAME: STAGE_NAMES,
TABLE_EXTRACTOR_NAME: TABLE_STAGE_NAMES,
}
244 changes: 244 additions & 0 deletions backend/agent_kv/dispatch.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,244 @@
"""Executor dispatch glue (spec §5.3). One dispatch per job; UUID task_id."""

import logging
import uuid

from celery import signature
from django.conf import settings
from django.utils import timezone

from agent_kv.constants import EXECUTION_SOURCE, EXTRACTOR_ROUTES
from agent_kv.models import AgentKVJob, JobStatus
from unstract.sdk1.execution.context import ExecutionContext

logger = logging.getLogger(__name__)

CALLBACK_QUEUE = "agent_kv_callback"


class DispatchError(Exception):
"""Enqueue failed; the caller terminalizes the job (spec §5.3)."""


def dispatch_cancelled_webhook(job) -> None:
"""Queue the terminal webhook for a job the API just cancelled.

Cancellation never reaches finalize -- ``JobCancelView`` and DELETE
terminalize the row themselves -- so without this the caller who supplied
``webhook_url`` is never told, and a late executor callback cannot tell them
either (it loses the terminal guard, and the callback declines a non-fresh
finalize rather than double-notifying). Docs §8 promises delivery on
terminal states; this is the cancel half of that promise.

Call ONLY when the guarded cancel actually won. That is what makes the two
paths mutually exclusive: a cancel that lost means a finalize won and will
send, and a cancel that won means no fresh finalize can.

Best-effort by design. The job is already cancelled and the caller already
has their 200; failing their request because a notification could not be
QUEUED would be the wrong trade, so this logs and returns rather than
raising. Delivery itself is the worker's problem.
"""
if not job.webhook_url:
return
try:
from pg_queue.producer import enqueue_task

enqueue_task(
task_name="agent_kv_cancelled",
queue=CALLBACK_QUEUE,
kwargs={
"callback_kwargs": {
"job_id": str(job.id),
"webhook_url": job.webhook_url,
}
},
org_id=str(job.organization_id),
)
except Exception:
logger.exception(
"agent-kv: could not queue the cancellation webhook for job %s; "
"the job IS cancelled, only the notification was lost",
job.id,
)


def _dispatcher():
# No `celery_app`: UN-4046 removed that parameter when the routing
# dispatcher's Celery branch went with the pg_queue_enabled flag. Passing it
# raises TypeError, so every submit failed to dispatch -- and because it
# fails at the call rather than at import, nothing catches it until a real
# request is made.
from pg_queue.executor_rpc import get_executor_dispatcher

return get_executor_dispatcher()


def _platform_api_key(job) -> str:
# Lazy import: avoids Django app registry init order (mirrors
# PromptStudioHelper._get_platform_api_key).
from platform_settings_v2.platform_auth_service import (
PlatformAuthenticationService,
)

# ``get_active_platform_key`` takes the org's public *slug*
# (``Organization.organization_id``, e.g. ``org_abc123``) and resolves it
# via ``get_organization_by_org_id`` -- NOT the row's UUID primary key that
# ``job.organization_id`` holds. Passing the PK here silently resolves to
# no organization and every dispatch fails with ``ActiveKeyNotFound``
# (caught live in the Task 13b integration run).
org_slug = job.organization.organization_id
platform_key = PlatformAuthenticationService.get_active_platform_key(org_slug)
if not platform_key:
raise DispatchError(f"No active platform key for org {org_slug}")
return str(platform_key.key)


def dispatch_job(
job, *, extractor: str, schema: dict, options: dict, adapters: dict | None = None
) -> None:
executor_name, operation = EXTRACTOR_ROUTES[extractor]
org_id = str(job.organization_id)
# Everything that can fail — platform-key lookup, context construction,
# and the enqueue call itself — lives inside this try so no internal
# failure (e.g. a transient DB error resolving the platform key) can
# escape as a raw, uncaught exception. Only the post-success bookkeeping
# below runs outside it.
try:
job.task_id = uuid.uuid4()
context = ExecutionContext(
executor_name=executor_name,
operation=operation,
run_id=str(job.id),
execution_source=EXECUTION_SOURCE,
organization_id=org_id,
executor_params={
"job_id": str(job.id),
"input_ref": job.input_ref,
"schema": schema,
"options": options,
# Platform adapter instance ids, by role, already validated at
# submit against THIS job's organization and against the
# expected `AdapterTypes` (see
# `execution_serializers._validated_adapter_shape` for the
# shape half and `execution_views._resolved_adapters` for the
# tenancy half).
#
# Defence in depth, NOT the only check: the platform service
# re-scopes every lookup as
# `WHERE id=%s and organization_id=%s`
# (platform-service/.../helper/adapter_instance.py:28-30),
# where the org comes from the bearer platform key that
# `_platform_api_key(job)` below mints from THIS job's org. So
# org A cannot spend org B's credential even with the submit
# gate removed. What the submit gate buys is a clean 400 naming
# the role instead of a mid-run `SdkError`, plus the two checks
# the platform service does NOT make: `is_usable` (exhausted
# trial) and `is_available` (deprecated).
#
# Empty for an env-configured extractor (`kv`), which is why
# this is a dict rather than three params: the two credential
# models coexist, one per extractor.
"adapters": adapters or {},
"platform_api_key": _platform_api_key(job),
# The CAP the engine must enforce (spec §6.1/§6.6), not the
# measured count -- job.pages_total is None for Excel (no
# pre-OCR page concept), which would otherwise leave the
# engine with nothing to check the post-OCR virtual-page cap
# against. The measured count still rides along separately.
"max_pages": settings.AGENT_KV_MAX_PAGES,
"pages_total": job.pages_total,
},
)
# Last check before spending money. A cancel can land between the
# submit's `job.save()` and this enqueue: the cancel sees a PENDING,
# never-dispatched row, so it terminalizes it AND releases its
# concurrency slot -- correctly, because nothing had been dispatched
# yet. Enqueueing anyway would then run paid work for a job the caller
# already cancelled, with its slot already handed to someone else.
#
# Re-read rather than trusting the in-memory row, which predates the
# cancel by construction.
if AgentKVJob.objects.filter(
id=job.id, status__in=list(AgentKVJob.TERMINAL)
).exists():
logger.info(
"agent-kv: job %s was terminalized before dispatch; not enqueueing",
job.id,
)
return

cb_kwargs = {"callback_kwargs": {"job_id": str(job.id), "org_id": org_id}}
_dispatcher().dispatch_with_callback(
context,
on_success=signature(
"agent_kv_complete", kwargs=cb_kwargs, queue=CALLBACK_QUEUE
),
on_error=signature("agent_kv_error", kwargs=cb_kwargs, queue=CALLBACK_QUEUE),
task_id=str(job.task_id),
)
except DispatchError:
raise
except Exception as e:
raise DispatchError(str(e)) from e
job.status = JobStatus.DISPATCHED
job.dispatched_at = timezone.now()
# Guarded queryset UPDATE, not job.save(): a plain save would blindly
# overwrite whatever status this job already raced to. Concretely: the
# executor can fail (or the job be cancelled) essentially instantly
# after enqueue, and its finalize callback can land -- marking the row
# FAILED/CANCELLED -- before this post-enqueue bookkeeping runs. An
# unconditional save() here would rewrite that terminal status back to
# DISPATCHED, un-terminalizing the job forever (nothing else ever
# revisits a DISPATCHED row). Only a still-PENDING row is advanced; a
# 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).
# Everything below is POST-ENQUEUE bookkeeping. The task is already on the
# queue, so a failure here must never be reported as a failed dispatch:
# `SubmitView` turns a DispatchError into a FAILED job, and the executor
# would then run, callback, and find a terminal row it cannot write to --
# the caller told nothing was billed for work that did run. Wrapped rather
# than left to propagate, which is what the single pre-review UPDATE did.
#
# Losing the bookkeeping entirely is recoverable: sweep phase 1 reaps a
# still-PENDING row with no `dispatched_at`, and phase 2's
# `dispatched_at IS NULL` arm covers the non-PENDING case.
try:
_record_dispatch(job)
except Exception:
logger.exception(
"agent-kv: dispatch bookkeeping failed for job %s after enqueue "
"(task is queued; the sweep will reconcile)",
job.id,
)


def _record_dispatch(job) -> None:
"""Persist task_id/status/dispatched_at against whatever the row raced to."""
advanced = AgentKVJob.objects.filter(id=job.id, status=JobStatus.PENDING).update(
task_id=job.task_id,
status=job.status,
dispatched_at=job.dispatched_at,
)
if not advanced:
# The row moved off PENDING between the enqueue above and this write.
# The benign case is a terminal status (the guard's whole purpose) --
# but there is a non-terminal one: StageReportView promotes
# PENDING -> RUNNING on the executor's FIRST stage report, which can
# easily land before this bookkeeping. The guard above then matches 0
# rows and `dispatched_at` stays NULL -- and a non-terminal row with a
# NULL `dispatched_at` is invisible to BOTH sweep phases: phase 1
# requires `status=PENDING`, phase 2 filters `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, with nothing able to
# recover it.
#
# So stamp the dispatch bookkeeping for any still-non-terminal row,
# WITHOUT touching `status`: the row genuinely was dispatched, and
# moving RUNNING back to DISPATCHED would lose the executor's own
# progress. `dispatched_at__isnull=True` keeps this idempotent and
# stops a retry overwriting the original dispatch time.
AgentKVJob.objects.filter(id=job.id, dispatched_at__isnull=True).exclude(
status__in=list(AgentKVJob.TERMINAL)
).update(task_id=job.task_id, dispatched_at=job.dispatched_at)
33 changes: 33 additions & 0 deletions backend/agent_kv/exceptions.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
from rest_framework.exceptions import APIException


class EngineUnavailable(APIException):
status_code = 501
default_detail = "agent-kv engine not available on this deployment"


class RateLimited(APIException):
status_code = 429
default_detail = "Too many requests"


class JobNotFound(APIException):
status_code = 404
default_detail = "Job not found"


class SubscriptionGateUnavailable(APIException):
"""The engine plugin is installed but exposes no subscription gate.

Deliberately NOT a 402: the org's subscription was never evaluated, so
claiming it was denied would be a lie. This is a deployment fault — a build
that can run billable work but cannot check entitlement — and it fails
CLOSED, because the alternative is admitting unmetered paid work on a route
whose URL carries no org segment for the middleware to fall back on.
"""

status_code = 503
default_detail = (
"agent-kv cannot verify subscription entitlement on this deployment; "
"the engine plugin exposes no subscription gate"
)
Loading
Loading