Skip to content

Commit b3f25d9

Browse files
UN-3563 [FEAT] PG Queue 9e PR 2c — live PG fan-out/barrier/callback (fire-and-forget) (#2069)
* UN-3563 [FEAT] PG Queue 9e PR 2c — live PG fan-out/barrier/callback (fire-and-forget) Wires the coupled pipeline's fan-out → barrier → callback onto the PG queue for a transport=="pg_queue" execution. Gated: resolve_transport() still returns celery (PR3 Flipt flips it), so the whole PG branch is present-but-unreachable — default path byte-identical. Orchestrator task (async_execute_bin) stays on Celery (hybrid); routing it onto PG is a 2d follow-up. - barrier.py: Barrier Protocol + CeleryChordBarrier/RedisDecrBarrier accept (and ignore) a `transport` param; CallbackDescriptor gains an optional `backend`. - orchestration_utils._barrier_for_transport: pg_queue → fresh PgBarrier() (bypasses the WORKER_BARRIER_BACKEND singleton), else the singleton. - pg_barrier.PgBarrier.enqueue(transport): pg_queue → fire-and-forget mode — _dispatch_header_pg sends each header via dispatch(backend=PG) with an injected _barrier_context {execution_id, batch_index, callback_descriptor}, no .link; descriptor marked backend=pg_queue; UPSERT block also clears pg_batch_dedup (greptile #2068 reuse-reset). _fire_barrier_callback self-chains the callback onto PG when backend==pg_queue. clear_execution_batches at finalise + abort. run_batch_with_barrier(): claim → work → in-body _barrier_pg_decrement; redelivery skips; exception → barrier_pg_abort. - file_processing.process_file_batch(_barrier_context=None): core routes None → _run_batch_stages (celery chord path), else → run_batch_with_barrier. - general/api fan-outs thread transport into create_chord_execution. Tests: +8 PgBarrier fire-and-forget + 2 orchestration routing + 2 process_file_batch routing. Each test file green alone; ruff clean. End-to-end forced-pg dev-test pending before PR. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * UN-3563 fix SonarCloud S1172: drop unused task_instance from _run_batch_stages The extracted _run_batch_stages never uses task_instance — its only purpose (deriving celery_task_id) happens in _process_file_batch_core before the call. Removed the param + updated both call sites. _process_file_batch_core keeps task_instance (it reads .request.id). Routing test mocks with *a, unaffected. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * UN-3563 address review (muhammad-ali-e, 15): strand-on-failure hardening + typing/dedup/docs/tests Decision (with reviewer): reaper-as-safety-net for the un-catchable strand windows + fix what's catchable + document + gate on PR3. Failure handling: - [#69 Critical] run_batch_with_barrier wraps BOTH work + decrement in the abort: a decrement-side failure (guard / DB / last-batch callback dispatch) tears the barrier down in-body instead of stranding to expiry. - [#79] extracted _abort_barrier_in_body — logs when the teardown itself fails (was silently suppressed under a misleading "torn down" message). - [#74/#81] documented the two un-catchable strand windows (hard-crash-during-work, post-commit callback-dispatch-fail) as a HARD reaper dependency for PR3. - [#86] finalise cleanup split into independent try/excepts with distinct logs. Typing / clarity: - [#1] BarrierContext(TypedDict) for _barrier_context (header fan-out, run_batch_with_barrier, process_file_batch). - [#3] renamed CallbackDescriptor "backend" -> "transport" (WorkflowTransport value; avoids the QueueBackend "pg" collision). - [#27] is_pg_transport() predicate in core; used in orchestration_utils + pg_barrier. - [#20] extracted _dispatch_pg() — single home for cycle-avoiding local import + backend=PG. - [#35] normalize_transport() at the general worker entry (parity w/ api/scheduler). - [#94] log when a header has no queue option. - [#9/#13] fixed born-stale comment + kwargs-not-args docstring. Tests (+#37/#41): last-batch self-chains callback to PG + cleans up barrier/dedup; decrement-failure aborts; PG-branch mid-loop dispatch-failure deletes row; header args/queue/pre-existing-kwargs preservation. 137 barrier/dedup/routing tests green; bootstrap clean under WORKER_BARRIER_BACKEND=pg. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * UN-3563 fix run_batch_with_barrier strand-window doc inconsistency (review) The second "NOT catchable" bullet conflated two different things: it described the in-body catchable abort ("the abort here removes the row") and a *software* callback-dispatch failure — but that failure is already caught + torn down by step 3's wrap (paragraph 1), so it doesn't belong under the un-catchable heading, and on the PG path _fire_barrier_callback IS the enqueue so "committed but before the enqueue" couldn't both hold. Rewrote the bullet to the genuinely un-catchable window: a hard crash BETWEEN the decrement committing (remaining→0) and the callback enqueue completing — decrement committed (redelivery blocked by the marker), process gone before the callback enqueues or any abort runs, row survives to expiry, reaper-only recovery. Explicitly notes a software dispatch failure is the catchable case. Keeps this list an accurate spec for the PR-3 reaper-recovery dependency. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * UN-3563 address greptile (#2069, 2): clear dedup on mid-loop PG failure + carry fairness on PG callback Both in the gated PG path (greptile 4/5, safe to merge). - Issue 1: PgBarrier.enqueue mid-loop dispatch-failure handler now also calls clear_execution_batches on the PG path. Earlier headers may have committed a claim_batch marker; with the barrier row deleted, their in-flight barrier_pg_abort is a no-op (already_aborted) and never reaches the clear inside it, so reclaim the markers directly here. - Issue 2: the PG callback now carries the producer's fairness. Added _fairness_from_headers() to reconstruct the FairnessKey from the stored x-fairness-key headers and pass it to _dispatch_pg, so the callback rides the same org/priority as the Celery path (was always default priority). Tests: +fairness-carried / +fairness-none-safe on _fire_barrier_callback; extended the PG mid-loop test to assert an already-claimed marker is reclaimed. 75 barrier/dedup tests green; bootstrap clean under WORKER_BARRIER_BACKEND=pg. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * UN-3563 fix SonarCloud S3776: reduce PgBarrier.enqueue cognitive complexity (17→under 15) Extracted the per-header dispatch loop into PgBarrier._dispatch_headers — the deeply-nested for→try/except→if/else→if (PG-vs-celery branch + mid-loop failure teardown + PG dedup-clear) was the complexity driver. enqueue now calls the helper; behaviour identical. radon: enqueue C(11)→B(6); ruff C901 passes. 75 barrier/dedup tests green; ruff + ruff-format clean. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * UN-3563 fix greptile #2069: mid-loop dedup-clear test passed for the wrong reason The pre-seeded claim_batch marker was wiped by enqueue's UPSERT block (the reuse-reset DELETE) before the dispatch loop, so the mid-loop clear_execution_batches deleted 0 rows — the count==0 assertion passed on the UPSERT, not the guard under test. Now the first dispatch side-effect claims the marker AFTER the UPSERT (simulating a fast PG consumer), so the mid-loop clear is what removes it. Verified: with the clear disabled the marker orphans (count=1). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 5f2d18d commit b3f25d9

11 files changed

Lines changed: 894 additions & 93 deletions

File tree

‎unstract/core/src/unstract/core/data_models.py‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -249,6 +249,17 @@ def normalize_transport(value: object, *, logger: Any = None, context: str = "")
249249
return DEFAULT_WORKFLOW_TRANSPORT
250250

251251

252+
def is_pg_transport(transport: str | None) -> bool:
253+
"""True if ``transport`` is the Postgres-queue transport.
254+
255+
Single source for "what counts as PG transport" — centralises the
256+
``== WorkflowTransport.PG_QUEUE.value`` comparison scattered across the
257+
worker fan-out / barrier code, and the seam to extend if a second
258+
PG-family transport is ever added.
259+
"""
260+
return transport == WorkflowTransport.PG_QUEUE.value
261+
262+
252263
class FileListingResult:
253264
"""Result of listing files from a source."""
254265

‎workers/api-deployment/tasks.py‎

Lines changed: 16 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,13 @@
2626
from shared.workflow.execution.tool_validation import validate_workflow_tool_instances
2727
from worker import app
2828

29-
from unstract.core.data_models import ExecutionStatus, FileHashData, WorkerFileData
29+
from unstract.core.data_models import (
30+
DEFAULT_WORKFLOW_TRANSPORT,
31+
ExecutionStatus,
32+
FileHashData,
33+
WorkerFileData,
34+
normalize_transport,
35+
)
3036
from unstract.core.worker_models import ApiDeploymentResultStatus
3137

3238
logger = WorkerLogger.get_logger(__name__)
@@ -683,6 +689,14 @@ def _run_workflow_api(
683689
org_id=str(schema_name),
684690
workload_type=WorkloadType.API,
685691
)
692+
# Transport rides in via the dispatched task's kwargs (PR 1 seam); the
693+
# fan-out honours it (PG path → fire-and-forget PgBarrier). Fail-closed
694+
# to celery on any unrecognised value.
695+
transport = normalize_transport(
696+
kwargs.get("transport", DEFAULT_WORKFLOW_TRANSPORT),
697+
logger=logger,
698+
context=f" [exec:{execution_id}]",
699+
)
686700
result = WorkflowOrchestrationUtils.create_chord_execution(
687701
batch_tasks=batch_tasks,
688702
callback_task_name="process_batch_callback_api",
@@ -694,6 +708,7 @@ def _run_workflow_api(
694708
callback_queue=file_processing_callback_queue,
695709
app_instance=app,
696710
fairness=api_fairness,
711+
transport=transport,
697712
)
698713

699714
if result is None:

‎workers/file_processing/tasks.py‎

Lines changed: 58 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,8 @@
1111
from typing import Any
1212

1313
from queue_backend import worker_task
14+
from queue_backend.barrier import BarrierContext
15+
from queue_backend.pg_barrier import run_batch_with_barrier
1416

1517
# Import shared worker infrastructure
1618
from shared.api import InternalAPIClient
@@ -222,25 +224,16 @@ def _enhance_batch_with_mrq_flags(
222224
)
223225

224226

225-
def _process_file_batch_core(
226-
task_instance, file_batch_data: dict[str, Any]
227+
def _run_batch_stages(
228+
file_batch_data: dict[str, Any], celery_task_id: str
227229
) -> dict[str, Any]:
228-
"""Core implementation of file batch processing.
229-
230-
This function contains the actual processing logic that both the new task
231-
and Django compatibility task will use.
230+
"""The actual batch work (validate → setup → pre-create → process → compile).
232231
233-
Args:
234-
task_instance: The Celery task instance (self)
235-
file_batch_data: Dictionary that will be converted to FileBatchData dataclass
236-
237-
Returns:
238-
Dictionary with successful_files and failed_files counts
232+
Transport-agnostic: identical on the Celery chord path and the PG
233+
fire-and-forget path. The task instance isn't needed here — its only use
234+
(deriving ``celery_task_id``) happens in the caller. Returns the
235+
JSON-serialisable batch result.
239236
"""
240-
celery_task_id = (
241-
task_instance.request.id if hasattr(task_instance, "request") else "unknown"
242-
)
243-
244237
# Step 1: Validate and parse input data
245238
batch_data = _validate_and_parse_batch_data(file_batch_data)
246239

@@ -260,6 +253,44 @@ def _process_file_batch_core(
260253
return _compile_batch_result(context)
261254

262255

256+
def _process_file_batch_core(
257+
task_instance,
258+
file_batch_data: dict[str, Any],
259+
barrier_context: BarrierContext | None = None,
260+
) -> dict[str, Any]:
261+
"""Core implementation of file batch processing.
262+
263+
This function contains the actual processing logic that both the new task
264+
and Django compatibility task will use.
265+
266+
Args:
267+
task_instance: The Celery task instance (self)
268+
file_batch_data: Dictionary that will be converted to FileBatchData dataclass
269+
barrier_context: Present only on the 9e PG fire-and-forget path — carries
270+
``execution_id`` / ``batch_index`` / ``callback_descriptor`` so the
271+
batch claims its slot and runs the barrier decrement in-body (a
272+
PG-consumed task fires no Celery ``.link``). ``None`` on the Celery
273+
chord path, where the chord ``.link`` drives the decrement instead.
274+
275+
Returns:
276+
Dictionary with successful_files and failed_files counts
277+
"""
278+
celery_task_id = (
279+
task_instance.request.id if hasattr(task_instance, "request") else "unknown"
280+
)
281+
282+
if barrier_context is None:
283+
# Celery chord path — the chord's .link runs the decrement after this.
284+
return _run_batch_stages(file_batch_data, celery_task_id)
285+
286+
# PG fire-and-forget path — claim the batch (idempotent on redelivery), run
287+
# the stages, then decrement the barrier in-body / self-chain the callback.
288+
return run_batch_with_barrier(
289+
barrier_context,
290+
lambda: _run_batch_stages(file_batch_data, celery_task_id),
291+
)
292+
293+
263294
@worker_task(
264295
bind=True,
265296
name=TaskName.PROCESS_FILE_BATCH,
@@ -272,18 +303,26 @@ def _process_file_batch_core(
272303
# Timeout inherited from global Celery config (FILE_PROCESSING_TASK_TIME_LIMIT env var)
273304
)
274305
@monitor_performance
275-
def process_file_batch(self, file_batch_data: dict[str, Any]) -> dict[str, Any]:
276-
"""Process a batch of files in parallel using Celery.
306+
def process_file_batch(
307+
self,
308+
file_batch_data: dict[str, Any],
309+
_barrier_context: BarrierContext | None = None,
310+
) -> dict[str, Any]:
311+
"""Process a batch of files in parallel.
277312
278313
This is the main task entry point for new workers.
279314
280315
Args:
281316
file_batch_data: Dictionary that will be converted to FileBatchData dataclass
317+
_barrier_context: Injected only when this task is dispatched onto the PG
318+
queue (9e fire-and-forget path) by ``PgBarrier`` — carries the barrier
319+
coordination context (``execution_id`` / ``batch_index`` /
320+
``callback_descriptor``). Absent on the Celery chord path.
282321
283322
Returns:
284323
Dictionary with successful_files and failed_files counts
285324
"""
286-
return _process_file_batch_core(self, file_batch_data)
325+
return _process_file_batch_core(self, file_batch_data, _barrier_context)
287326

288327

289328
def _validate_and_parse_batch_data(file_batch_data: dict[str, Any]) -> FileBatchData:

‎workers/general/tasks.py‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,7 @@
5656
FileBatchData,
5757
FileHashData,
5858
WorkerFileData,
59+
normalize_transport,
5960
)
6061

6162
# Import common workflow utilities
@@ -488,6 +489,12 @@ def _execute_general_workflow(
488489
"""
489490
start_time = time.time()
490491

492+
# Fail-closed coercion, for parity with the api/scheduler workers (which run
493+
# normalize_transport at their entry): a typo'd transport degrades to Celery
494+
# with a warning rather than silently routing onto an unknown substrate. The
495+
# coerced value feeds both WorkflowContextData and the fan-out below.
496+
transport = normalize_transport(transport, logger=logger)
497+
491498
logger.info("Executing general workflow logic for ETL/TASK workflow")
492499

493500
try:
@@ -710,6 +717,7 @@ def _execute_general_workflow(
710717
execution_mode=execution_mode,
711718
use_file_history=use_file_history,
712719
organization_id=api_client.organization_id,
720+
transport=transport,
713721
**kwargs,
714722
)
715723

@@ -764,6 +772,7 @@ def _orchestrate_file_processing_general(
764772
execution_mode: tuple | None,
765773
use_file_history: bool,
766774
organization_id: str,
775+
transport: str = DEFAULT_WORKFLOW_TRANSPORT,
767776
**kwargs: dict[str, Any],
768777
) -> dict[str, Any]:
769778
"""Orchestrate file processing for general workflows using the same pattern as API worker.
@@ -957,6 +966,7 @@ def _orchestrate_file_processing_general(
957966
org_id=organization_id,
958967
workload_type=WorkloadType.NON_API,
959968
),
969+
transport=transport,
960970
)
961971

962972
if not result:

‎workers/queue_backend/barrier.py‎

Lines changed: 34 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,10 +36,12 @@
3636

3737
import logging
3838
import os
39-
from typing import TYPE_CHECKING, Any, Protocol, TypedDict
39+
from typing import TYPE_CHECKING, Any, NotRequired, Protocol, TypedDict
4040

4141
from celery import chord
4242

43+
from unstract.core.data_models import DEFAULT_WORKFLOW_TRANSPORT
44+
4345
from .fairness import FairnessKey
4446
from .handle import BarrierHandle
4547

@@ -92,6 +94,29 @@ class CallbackDescriptor(TypedDict):
9294
kwargs: dict[str, Any]
9395
queue: str
9496
fairness_headers: dict[str, Any] | None
97+
# 9e: the WorkflowTransport the aggregating callback is fired on when the
98+
# barrier completes. Absent / ``None`` → legacy Celery dispatch
99+
# (``current_app.apply_async`` — the ``.link`` path). ``"pg_queue"`` → the
100+
# fire-and-forget PG path self-chains the callback via dispatch onto PG.
101+
# Named ``transport`` (a WorkflowTransport value, e.g. ``"pg_queue"``) — NOT
102+
# ``backend``, to avoid confusion with ``QueueBackend`` (``"pg"``).
103+
transport: NotRequired[str | None]
104+
105+
106+
class BarrierContext(TypedDict):
107+
"""Per-batch barrier coordination injected into a PG-dispatched header task.
108+
109+
The 9e fire-and-forget sibling of :class:`CallbackDescriptor`, and typed for
110+
the same reason: it crosses producer → PG-queue → consumer, so the contract
111+
is pinned here to catch a typo/rename at the type layer rather than as a
112+
remote ``KeyError`` mid-batch. Carried as the ``_barrier_context`` kwarg on
113+
``process_file_batch`` so the consumer can claim its slot and run the barrier
114+
decrement in-body (a PG-consumed task fires no Celery ``.link``).
115+
"""
116+
117+
execution_id: str
118+
batch_index: int
119+
callback_descriptor: CallbackDescriptor
95120

96121

97122
class Barrier(Protocol):
@@ -112,9 +137,15 @@ def enqueue(
112137
callback_queue: str,
113138
app_instance: Any,
114139
fairness: FairnessKey | None = None,
140+
transport: str = DEFAULT_WORKFLOW_TRANSPORT,
115141
) -> BarrierHandle | None:
116142
"""Enqueue ``header_tasks`` and a single callback to fire on completion.
117143
144+
``transport`` is the per-execution transport (9e). Only ``PgBarrier``
145+
acts on it (``pg_queue`` → fire-and-forget PG fan-out instead of Celery
146+
``.link``); the Celery/Redis substrates accept it for Protocol parity
147+
and ignore it (they are only ever reached on the ``celery`` transport).
148+
118149
``None`` is the **sole signal** that no work was enqueued — it
119150
is returned exclusively when ``header_tasks`` is empty. Any
120151
substrate-level failure (broker outage, serialisation error,
@@ -166,8 +197,10 @@ def enqueue(
166197
callback_queue: str,
167198
app_instance: Any,
168199
fairness: FairnessKey | None = None,
200+
transport: str = DEFAULT_WORKFLOW_TRANSPORT,
169201
) -> BarrierHandle | None:
170202
"""See :class:`Barrier.enqueue`."""
203+
del transport # Celery chord is only reached on the celery transport.
171204
# Empty-header guard goes FIRST so a zero-task run skips
172205
# signature construction / fairness serialisation entirely.
173206
# Any failure in those paths is now constrained to the

0 commit comments

Comments
 (0)