UN-4078 [MISC] Remove the Celery execution transport from the workers, backend and SDK - #2284
Conversation
… its tests First slice of the Celery transport removal. `ExecutionDispatcher` published `execute_extraction` to RabbitMQ and blocked on an `AsyncResult`; nothing has selected it since UN-4046 made `get_executor_dispatcher()` return the PG request-reply dispatcher unconditionally. It had zero non-test importers. Deleted: `unstract/sdk1/.../execution/dispatcher.py` (316 lines), its export from the package `__init__`, and the `TestExecutionDispatcher` suite. The 11 workers test files that imported it are retargeted rather than dropped — they encode live product behaviour (which queue each operation lands on, what the payload carries), and that behaviour survives on `PgExecutionDispatcher`: - Queue naming moves to `unstract.workflow_execution.executor_rpc.QUEUE_PREFIX`, the surviving definition. The `celery_executor_*` names are NOT dead Celery surface — `worker-pg-executor` subscribes to exactly those strings — so the assertions pin the literal wire names. - New `tests/executor_dispatch_fakes.py` holds the shared fake transport, an eager variant that runs the task in-process (replacing the `send_task` monkey-patching four round-trip tests did), and a correctly-shaped callback signature builder. - The header-forwarding suite is replaced, not ported: PG dispatch takes no `headers=` by design, so the new tests assert org routing rides the payload AND that passing `headers=` still raises — the exact mistake that broke three call sites during UN-4046. - The raw-`send_task` canary is kept and its message updated; it matters more now, since a raw send_task publishes to a broker nothing drains. Two docstrings that claimed executor calls "go through Celery" were corrected — they became false with this change, not merely dated. Verification: workers 1517 passed / 1 skipped. sdk1 unchanged at 20 pre-existing failures (`test_llm_compat`, `test_litellm_cohere_timeout`), reproduced on a stashed tree to confirm they predate this work. Pre-commit clean on all 17 files. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Removes the Celery fan-out substrate and the per-execution transport field. Nothing has selected either since UN-4046 made `select_backend()` return PG unconditionally, but both were still compiled, imported and shipped — and importable Celery surface has already cost two incidents (UN-3779's hardcoded dispatcher, and the `headers=` breakage at three call sites during UN-4046). Deleted outright: - `queue_backend/redis_barrier.py` (761 lines) — the DECR-counter substrate - `CeleryChordBarrier` + `BarrierBackend` + `get_barrier()` + the `WORKER_BARRIER_BACKEND` env selector - `WorkflowTransport`, `DEFAULT_WORKFLOW_TRANSPORT`, `normalize_transport`, `is_pg_transport` (88 lines of `unstract/core/data_models.py`) - `barrier_pg_decr_and_check` — a thin @worker_task wrapper reachable only as a Celery `.link`; the decrement runs in-body via `_barrier_pg_decrement` Collapsed: the 10 `is_pg_transport` branch sites, the `transport` parameter threaded through the general/api/scheduler orchestrators, the `transport` field on `WorkflowContextData` and `FileProcessingContext`, the `is_pg` flag on the file-batch path, and the `transport` key the backend wrote into the dispatch payload. `_barrier_for_transport` becomes `PgBarrier()`. Two guards become unconditional rather than PG-gated: the terminal-execution skip and the duplicate-destination-write check. Both were only ever no-ops on Celery. This also starts retiring a documented live hazard. `DEFAULT_WORKFLOW_TRANSPORT` was `"celery"`, so a payload that lost the field during a rolling deploy selected a Celery fan-out with no consumers. Worth correcting the in-code note while deleting it: it claimed `CeleryChordBarrier` was selected and no PG rows were written, so the reaper could not see the strand. Neither held — every PG worker sets `WORKER_BARRIER_BACKEND=pg`, so `PgBarrier` was selected, and its `_reset_barrier` UPSERT is unconditional, so the barrier row was written and the reaper did see it. Real behaviour was a ~2.5h delayed ERROR. The reads are removed here, but the writes are NOT: a pre-UN-4078 worker still defaults an absent field to Celery, so during a rolling deploy it would publish to a RabbitMQ with no consumers. The next commit keeps writing the field for one release; the hazard is gone only once that shim is removed. `EXECUTION_EXCLUDED_PARAMS` deliberately keeps its `"transport"` entry: a rolling deploy can still have an un-upgraded producer emitting the field, and without the exclusion it reaches the legacy `execute_workflow` signature and raises. Tests: four whole suites deleted with their substrate (redis barrier, barrier backend selection, barrier differential, workflow-context transport). The rest are retargeted rather than dropped — the chord canary now asserts *zero* chord call sites (with a non-vacuity lock, since a zero-expectation assertion is exactly what a path bug passes silently), and the Celery-gated guard tests are replaced by the unconditional behaviour they now have. Verification: workers 1438 passed / 1 skipped / 0 failed. Pre-commit clean. The backend suite cannot start in this environment (`AppRegistryNotReady`, reproduced on a stashed tree — pre-existing), so the four backend edits are CI-verified only. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
…rolling-deploy shim The previous commit removed every read of the `transport` payload field and every write. Removing the writes in the same release is unsafe: a pre-UN-4078 worker still treats an absent field as "celery" (DEFAULT_WORKFLOW_TRANSPORT) and, during a Kubernetes RollingUpdate, can receive a payload from an upgraded producer. It then publishes to RabbitMQ, which has no Celery consumers, so: - general workflows are marked ERROR by the reaper after ~2.5h, even with every file processed; - API deployments take no orchestration claim, so nothing sweeps them and they stay EXECUTING. The window is real: old file-processing pods keep draining claimed batches for their full grace period (up to 9120s). The four writes are restored as `transport: "pg_queue"`, all via one constant, LEGACY_TRANSPORT_KEY / LEGACY_TRANSPORT_VALUE in unstract.core.data_models, so the follow-up removal is a single grep: - backend WorkflowHelper orchestrator dispatch payload - backend create_workflow_execution response - workers scheduler async_execute_bin kwargs - workers PgBarrier callback descriptor (CallbackDescriptor gains `transport: NotRequired[str]`) Nothing on the new side reads the field. Tests pin each write so the shim cannot be dropped early; the scheduler cases also assert a stale "celery" in the backend response is never forwarded. Remove in the release after UN-4078: the writes, the constants, the `transport` entry in EXECUTION_EXCLUDED_PARAMS, and test_legacy_transport_shim.py. Verification: workers 1478 passed / 1 skipped / 0 failed; backend dispatch tests 5 passed (with backend/sample.env loaded); pre-commit clean. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
|
muhammad-ali-e
left a comment
There was a problem hiding this comment.
Standardized review — REQUEST CHANGES
Mode: INITIAL · Head: 789b36b1 · run under unstract plugin v0.18.1 (16-lens rubric), executed via the PR Review Toolkit specialist agents.
Summary — Critical: 0 · High: 2 · Medium: 9 · Low: 3 · Lenses run: 16/16
The core change holds up. The deletion sweep leaves no live importer of any removed symbol (checked across backend/, workers/, unstract/, prompt-service/, platform-service/, x2text-service/, tools/, runner/); every collapsed branch kept the PG side, not the Celery one; all four shim writes survive serialization to the wire; and both unconditional-ised guards shed only an early return. The two High findings are about the safety net around the shim and a docstring that would mislead a maintainer into deleting an alarm — not about the deletion itself.
CI note: the three red checks (e2e, test (integration), report) all fail on pull access denied for minio/minio, a Docker Hub pull failure on the runner. Not caused by this PR. test (unit) passes.
Lens checklist (16/16)
| # | Lens | Result |
|---|---|---|
| 1 | Spec & intent | See findings #3, #4, #5 |
| 2 | Architectural fit | See findings #7, #8, #10 |
| 3 | Correctness & edge cases | See findings #3, #5 |
| 4 | Security | Clean — no auth, tenant-scoping, secrets or input-validation surface in the diff |
| 5 | Data integrity & migrations | Clean — backend/pg_queue/models.py changes are docstring-only, so no migration is owed and none is added |
| 6 | Concurrency | Clean — the single-UPDATE … RETURNING decrement, own-transaction requirement and DELETE … RETURNING abort are unchanged; collapsing the dispatch branch removed only the .link arm |
| 7 | API & contract compatibility | See findings #1, #7 — both skew directions otherwise sound; all four shim writes traced to the wire, and a stray transport kwarg is inert on new workers |
| 8 | Reliability & resilience | See finding #5 — otherwise both collapsed error handlers strictly improve (the api-deployment Celery arm had an unwrapped status update that could mask the original exception) |
| 9 | Performance & cost | Clean — no new queries, loops or hot-path I/O; PgBarrier() per call replaces a module singleton but holds no state |
| 10 | Observability | See finding #3 |
| 11 | Operational safety | See findings #1, #9 — worker-pg-reaper is a real, unprofiled compose service, so the reaper's new hard-dependency status is genuinely deployed |
| 12 | LLM/agent | Clean — the diff touches the LLM path only in docstrings and retargeted dispatch tests; no prompt template, model config, tool authorization or fallback logic changed |
| 13 | Testing | See findings #1, #10, #12 |
| 14 | Dependencies & build | Clean — zero changes to any pyproject.toml, lockfile or Dockerfile; sdk1 declares no celery dependency and retains no celery imports |
| 15 | Code quality | See findings #5, #8 |
| 16 | Doc & comment accuracy | See findings #2, #6, #8, #9, #11, #13, #14 |
Lenses 4, 5, 6, 9, 11, 12 and 14 were assessed by the orchestrator directly (no specialist agent leads on them).
Unanchored findings
These four touch files the diff does not modify, so they have no line to attach to.
[Medium] [Lens 1, 16] — QueueBackend.CELERY and its live send_task publish survive, while the docstring says they go with this change
workers/queue_backend/routing.py:44; workers/queue_backend/dispatch.py:9-11, :101-111, :36
dispatch.py:9-11 states "reachable only by an explicit backend=QueueBackend.CELERY override. Nothing passes one; the branch goes with the rest of the Celery transport." This PR is that removal, and the branch stayed: dispatch() still falls through to current_app.send_task(...), and from celery import current_app is still imported. So the PR's stated rationale — "while this code exists, the Celery library and broker settings are still imported", blocking UN-4077 — is not met at this site. A caller passing the member strands work with no pg_queue_message row, so neither the reaper nor the undispatched sweep can see it; scheduler/tasks.py:191 now passes backend=QueueBackend.PG explicitly, making a copy-paste with the wrong member plausible. Fix: drop the member, collapse dispatch() to _enqueue_pg(...), delete the celery import — or say in the description that it is deferred to UN-4077 and correct the docstring. Confidence: High on the facts; normalized down from one agent's High because no live caller exists today.
[Medium] [Lens 11, 16] — workers/sample.env ships a live WORKER_BARRIER_BACKEND=chord for a knob this PR made dead
workers/sample.env:98-117 (value at :107), :119-131; docker/sample.env:121-122; docker/docker-compose.yaml:210,246,280,321,360,413,492,544,585
This PR deletes the only reader (get_barrier()/BarrierBackend), yet the sample config an operator copies still sets it, uncommented, to the value that previously selected the deleted Celery barrier, and documents a redis option and "Default chord keeps existing behaviour exactly." The previous code raised at worker startup on an unrecognised value; now any value is silently ignored — the exact silent-misconfiguration posture the deleted factory existed to prevent. Separately WORKER_BARRIER_KEY_TTL_SECONDS (:119-131) is still read (barrier.py:64) but is documented entirely in terms of Redis DEL and a Celery .link that no longer exist, and states an expiry consequence contradicting pg_barrier.py:54-57. Fix: delete the dead blocks in this PR — pure deletion, no code risk. Confidence: High. (Four of five agents converged; the PR body acknowledges and defers this.)
[Low] [Lens 13] — Dead mock_chord fixture patches a symbol this PR deleted
workers/tests/test_barrier.py:51-62 does patch("queue_backend.barrier.chord"). barrier.py retains no chord symbol after this PR (only docstring text at :14, :165). No test requests the fixture, so CI stays green — it fails with AttributeError the moment anyone uses it. Also test_barrier.py:146 keeps a section banner for the suite that was deleted beneath it, and test_chord_callback_boundary.py:797 still asserts the sentinel {"r": "celery"}. Fix: delete the fixture; retitle the banner; rename the sentinel. Confidence: High.
[Low] [Lens 16] — docs/local-dev-setup-executor-migration.md documents the deleted ExecutionDispatcher
docs/local-dev-setup-executor-migration.md:31-48, :55 shows a "Post-Migration" architecture built on ExecutionDispatcher.dispatch() [Celery], AsyncResult.get() and a Celery result backend. This PR deletes that class and drops it from the SDK's __all__. A developer following the doc searches for a symbol that is gone and provisions a result backend the flow no longer uses. Fix: update to PgExecutionDispatcher over the PG queue, matching the already-corrected workers/file_processing/structure_tool_task.py:1-16. Confidence: High.
Open questions
- Is
process_file_batch_django_compatstill reachable from any producer outside this repo? If not it should go with the rest of the Celery transport; if it is, finding #3 is a live duplicate-billing path. - Is leaving
QueueBackend.CELERYdeliberate deferral to UN-4077? If so thedispatch.py:9-11docstring should say that rather than claim the branch already went. - Is
integration-workersintended to stayoptional: true? That is what makes thePgBarriershim assertion non-gating.
Assumptions
- The Celery worker fleet is genuinely at zero, so messages published to RabbitMQ are stranded rather than mis-executed. If any Celery consumer still runs, findings #4 and #5 escalate.
unstract-sdk1is not published to PyPI from this repo (no publish workflow under.github/), so removingExecutionDispatcherfrom its__all__is not an external breaking change; all in-repo consumers are path-based editable installs.
…ranches, stale comments Review findings on #2284, highest first. #1 (High) — half the rolling-deploy shim had no gating test that runs in the default lane. The create_workflow_execution response write had none at all, and the PgBarrier descriptor was pinned only by a test behind the barrier_db fixture, which skips wherever Postgres is absent. Descriptor construction moves to build_callback_descriptor(), pinned by a DB-free suite (workers/tests/test_legacy_transport_shim.py); the response write is pinned by a new backend case. All four writes now fail a required build if dropped early. #2 (High) — _process_file_batch_core's Args entries still described barrier_context=None as the supported Celery chord path while the body treats it as a malformed payload. Both entries rewritten; the log line now carries execution/workflow ids and the batch size. #3 — process_file_batch_django_compat delegates with two arguments, so after this PR every call landed in that malformed-payload branch: a false ERROR, and a full batch run with no claim (no pg_batch_dedup marker, so a redelivery re-runs it — LLM spend twice and a duplicate destination write with use_file_history=False). Its only producer was the Celery chord this PR removes, so it now refuses, the same fail-fast choice step_execution makes. Its MRQ helpers are left in place per the no-removal rule. #5 — PG_TRANSPORT_CALLBACK_KWARG became unconditional but its consumers still branched on it, leaving an else branch that swallows a failed finalization and strands the execution. The marker is now popped for wire compatibility only; the duplicate guard and the re-raise are unconditional. is_pg on _update_execution_status_unified becomes raise_on_failure (default True); the one caller inside an except block opts out so a raise cannot mask the original error. #6 — the removal checklist now lists all five artefact groups, and EXECUTION_EXCLUDED_PARAMS names the key through LEGACY_TRANSPORT_KEY instead of a hardcoded literal grep would miss. #7 — CallbackDescriptor.transport is Literal["pg_queue"] and required, not NotRequired[str]: it is mandatory this release, "celery" must never be written, and the follow-up removal becomes a type error at the write site. #8 — the Barrier protocol's "two call sites still program against it" was untrue; both construct PgBarrier concretely. Corrected, with what actually keeps it. #10 — executor_dispatch_fakes reimplemented the queue-naming rule and accepted any enqueue signature. queue_for now delegates to a new production helper queue_for_executor() (also used by the three dispatch sites), and the fake mirrors the QueueTransport protocol's keyword-only key set. #11, #14 — stale "PG only" / "No-op on Celery" comments at the sites this PR made unconditional, a citation of the deleted Celery max_retries=0 wrapper, and a cross-reference to a chart file that lives in the cloud repo. Five tests that asserted the removed Celery branches are restated against the collapsed behaviour rather than deleted. Verification: workers 1481 passed / 1 skipped; backend shim, dispatch-orchestrator and async-wait tests pass; pre-commit clean. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
|
Unstract test resultsPer-group results
Critical paths
|
Syncs with main, which has since removed the Celery execution transport from the workers, backend and SDK (UN-4078 / #2284). No conflicts: that change and this one touch different halves — it removed a transport, this one changes how the Redis CLIENT is built. Re-checked rather than assumed, because "merged cleanly" says nothing about whether the change still makes sense: * Socket.IO still rides kombu over Redis (backend/utils/log_events.py, workers/log_consumer/tasks.py) — so the rediss:// manager URL is still needed and still on the live log-streaming path. * The settings TLS block and the sidecar/tool env forwarding both survived intact. * Tests: core 160, runner 10, green. * Live re-probe against the local TLS Redis: managed-like URL mode and plain local Redis both connect.



What
ExecutionDispatcher(no non-test importers since UN-4046).redis_barrier.py,CeleryChordBarrier,BarrierBackend/get_barrier()/WORKER_BARRIER_BACKEND,WorkflowTransport,DEFAULT_WORKFLOW_TRANSPORT,normalize_transport,is_pg_transport, andbarrier_pg_decr_and_check.is_pg_transportbranch sites and thetransport/is_pgparameters threaded through the orchestrators,WorkflowContextDataandFileProcessingContext.transport: "pg_queue"in four places for one release, as a rolling-deploy shim (commit 3).Why
headers=breakage during UN-4046.How
Three commits, meant to be reviewed separately:
transport.LEGACY_TRANSPORT_KEY/LEGACY_TRANSPORT_VALUEinunstract.core.data_models:workflow_helper.py)create_workflow_executionresponse (internal_api_views.py)async_execute_binkwargs (workers/scheduler/tasks.py)PgBarriercallback descriptor (CallbackDescriptor.transport: NotRequired[str])Why the shim: a pre-UN-4078 worker treats a missing
transportas"celery". During a KubernetesRollingUpdate, old and new PG pods run together, and old file-processing pods keep draining claimed batches for up to 9120s. An old pod receiving a new producer's payload would publish to RabbitMQ, which has no consumers:Nothing on the new side reads the field.
Can this PR break any existing features? If yes, please list possible items. If no, please explain why.
No, for the following reasons:
select_backend()andget_executor_dispatcher()return PG unconditionally, and every PG worker already ranPgBarrier.transportkwarg lands in**kwargsand is ignored.EXECUTION_EXCLUDED_PARAMSkeeps its"transport"entry so it never reachesexecute_workflow."pg_queue".barrier_pg_decr_and_checkremoval is safe. It was only attached as a Celery.linkand never sent through PGdispatch(), so nopg_queue_messagerow can reference it._barrier_context, so both guards were already active.queue_backend.worker_task,pg_queue.executor_rpcandPgQueueClient, which all still exist.On-prem: docker-compose restarts all services together, so versions never mix. Upgrades from ≤ v0.179.0 must already go through v0.180.0, as documented in the release notes.
Relevant Docs
LEGACY_TRANSPORT_KEYcomment inunstract/core/src/unstract/core/data_models.py.Related Issues or PRs
LEGACY_TRANSPORT_*writes, the constants, the"transport"entry inEXECUTION_EXCLUDED_PARAMS, andtest_legacy_transport_shim.pyDependencies Versions / Env Variables
WORKER_BARRIER_BACKENDandWORKER_BARRIER_REDIS_*are now ignored. They're still present inworkers/sample.env,docker/sample.envanddocker-compose.yaml, which is harmless; cleanup can follow.Notes on Testing
tests/test_execution.py: 60 passed.test_legacy_transport_shim.py,test_execute_workflow_async_wait.py,test_dispatch_orchestrator.py): 5 passed, run withbackend/sample.envloaded andDJANGO_SETTINGS_MODULE=backend.settings.test."celery"in the backend response never being forwarded;PgBarriercallback descriptor;internal_api_views.pywrite has no unit test, because that view needs a DB.Screenshots
N/A
Checklist
I have read and understood the Contribution Guidelines.
🤖 Generated with Claude Code