Concurrency & Fan-out
Where parallelism actually lives in the pipeline, what bounds it, and which knobs are real. The short version: the normal local stack starts one Python worker, all matching workers share one Temporal (the workflow engine) task queue, activity slots are the execution-capacity boundary, and workflows fan out primarily for isolation.
Read this if you are tuning worker capacity, wondering why two stages are not running in parallel, or hunting a throughput bottleneck.
The Worker's Capacity Model
The normal local stack runs one long-lived worker process (jobctrl worker) on the jobctrl-default task queue. If more matching workers are running, Temporal may dispatch activities to any of them and the Operations read model aggregates their fresh capacity. Two numbers define each worker's capacity, both fixed at worker startup in infrastructure/temporal/worker.py:
- Activity slots — the
worker_activity_slotsSetting inconfig.json(default4): the maximum number of Temporal activities running at once, across all workflows. - Executor threads — separate worker-owned Temporal-sync and blocking-stage
ThreadPoolExecutorpools, each sizedslots + 2, so blocking stage work never spills into the process default executor or starves Temporal's own synchronous activities. Cancellation retires the affected blocking pool generation before its grace wait, so a fresh bounded generation accepts an immediate retry even when the old provider call ignores cancellation.
The worker heartbeat records the activity slots and configured executor width. GET /v1/health exposes the health boundary, while GET /v1/pipeline/operations derives the app directory from its configured database path, filters heartbeats to that resolved database/app-dir identity, selects the task queue named by the newest matching heartbeat, and aggregates fresh schema-valid rows from that queue into configured, active, and available slots. Changing worker_activity_slots in Settings writes config.json; restart the worker to apply the new capacity.
Two knobs that look like Temporal concurrency but are not:
- The Pipelines page's Internal concurrency field and Settings → General → Pipeline internal concurrency edit one shared
config.jsonvalue. Newly started manual Pipeline actions and automatic profile-update Score → Tailor → Cover batches snapshot that value. Source adapters use it for supported internal scraping parallelism, and preparation uses it for its bounded job batch; it does not create Temporal activity slots. - Resumable JobStreaming execution inside the broad-board source family is sequential by immutable query/location/board unit (a plain loop in the compatibility-named
jobspy.py). JobStreaming owns each board adapter's internal transport/pagination and cancellation-aware waits. Parallelising durable units is a filed backlog item, not current behavior.
Where Fan-out Happens (And Why)
- DiscoverWorkflow plans sources once (
plan_discovery_sources), then runs thediscovery_source_familyactivities one family at a time by default (max_parallel_families=1in Discovery Runtime's SQLite-backed settings). Planning is intentionally limited to seeding source controls and freezing the execution schedule; it does not run historical-job hygiene before producers start. That potentially expensive audit remains in terminal enrichment, so it preserves the cleanup contract without delaying time to first discovered job. This is deliberate isolation: each family gets its own activity-level timeout, heartbeat, and retry policy, and a family failure is recorded (families_failed) while the run continues to the next family. Raising the cap runs families concurrently (R9 Phase 3, below). - Score-as-you-discover streaming (R9 Phase 1, corrected 2026-08-03). The workflow starts one execution-scoped
discovery_enrichmentactivity before source crawling. It polls the durable current-execution membership whilediscovery_source_familyactivities are still committing jobs, so a broad JobStreaming family does not have to finish all of its search units before the first job can be enriched. After producers finish, the workflow cancels that live consumer and runs a terminal reconcile enrichment + fan-out, which remains authoritative for tolerated-partial-failure folding and progress finalization. The live consumer also admits current-execution robots-blocked rows and retryable failed LinkedIn rows once per activity lifetime, so re-observing a recoverable job does not leave it stranded behind the ordinary pending-only selector; non-retryable failures remain excluded. Per-family preparation fan-outs remain as best-effort backstops. The live enrichment and backstop fan-outs are progress-silent (progress_total=0) so the Runs bar stays monotonic on the family + terminal spine — see Operations & Events. Repeated fan-out is idempotent: the deterministicprep-{idempotency_key}id plusUSE_EXISTINGmeans N invocations start exactly one workflow per job. A one-time straggler sweep (include_pending_tailor=True) runs before the family loop — the only momentpending_tailorholds only pre-existing scored-but-not-tailored work and cannot race a fresh job's in-flight SCORE_JOB workflow. Every family + terminal fan-out is score-only, so a fresh job crossingpending_score->pending_tailormid-tailor is never double-fanned. - Per-job handoff (R9 Phase 2). The live enrichment activity runs with
per_job_handoff=True: as each job is individually enriched (committed topending_score), the enrichment worker starts that job'sSCORE_JOBpreparation workflow immediately, before its siblings in the same family are even scraped, tightening Time To First Score to per-job granularity. This is a side effect inside the enrichment activity (on_job_enrichedcallback threaded toenrichment/detail.py), soDiscoverWorkflow's command history is unchanged (determinism/replay safe). Starts use the same deterministicprep-{idempotency_key}id as the fan-out, so the per-job handoff and the reconciling fan-outs converge on exactly one execution per job (USE_EXISTING). Per-job starts are serialized by a lock because_run_detail_scrapermay enrich sites in parallel threads, and the handoff is best-effort (a start failure is logged and left for the fan-out backstop, never mistaken for an enrichment failure). - Parallel source families (R9 Phase 3, gated, default off). Families are processed in batches of the Discovery Runtime
max_parallel_familiesvalue (default1= sequential = today's behavior). With a value > 1, that many families' source crawls run concurrently (asyncio.gatherover the batch) while the single execution-scoped enrichment activity consumes committed jobs. The batch's score-only fan-out then runs once afterward. The cap is resolved at planning time (inplan_discovery_sources) and threaded through the plan, so the workflow stays deterministic; results are folded in submission order, and a canceled source in any batch cooperatively cancels the whole run. See the worker-capacity analysis below before raising the cap. - JobPipelineWorkflow executes the selected stages in pipeline order and delegates
discoverandapplyto child workflows, so the risky surfaces keep their own workflow identity, history, and retry policy. - ApplyWorkflow is single-flight per job: the workflow id
apply-{tenant}-{jobKey}plusUSE_EXISTINGand a one-attempt live retry policy make a duplicate submission structurally impossible rather than merely unlikely. - Concurrent workflows (a discovery run, a scoring batch, an apply) share the worker's activity slots; Temporal queues whatever exceeds them.
What Bounds Throughput
| Bound | Mechanism | Where |
|---|---|---|
| Activity slots and executor pools | Settings worker_activity_slots; Temporal-sync and blocking-stage executors each slots + 2; canceled blocking generations retire before grace | config.json, infrastructure/temporal/worker.py, infrastructure/temporal/run_in_activity.py |
| Parallel discovery families | Discovery Runtime max_parallel_families (default 1) | SQLite, infrastructure/temporal/concurrency.py, discovery/workflow.py |
| LLM spend | check_spend_budget preflight stops spendful workflows at the daily ceiling | Spend Ceiling |
| Retries | per-activity retry policies from the error taxonomy | Envelope & Activities |
| Worker readiness | worker-backed API actions return 503 until a healthy heartbeat exists | Runtime & Processes |
| Observed capacity | fresh matching heartbeat rows; exact slots include every activity even when safe detail is omitted | Operations & Events |
| Task-queue pressure | approximate workflow/activity backlog and poller observations; unavailable/unsupported are not zero | GET /v1/pipeline/operations |
| Apply single-flight | per-job workflow id + submit-intent checkpoint | Stage Walkthrough |
Worker-Capacity Analysis (Parallel Families)
Discovery Runtime max_parallel_families is off by default (1) because browser/resource contention is the first-class risk of running families concurrently: there is operational history of uncontrolled browser concurrency destroying long runs. For a chosen cap M:
- Peak activity slots. Up to
Mdiscovery_source_familyactivities plus one livediscovery_enrichmentactivity run at once, competing with fan-out + per-jobJobPreparationWorkflows from Phases 1–2 for the sameworker_activity_slotssetting (default 4). Temporal queues the excess, soMat or above the slot count starves enrichment. KeepM + 1below the activity-slot ceiling and leave headroom for preparation. - Peak concurrent browsers. The single enrichment consumer can overlap the batch's source crawls: roughly
Mbrowser-launching families plus one enrichment browser group, each with its configured internal workers. Each headless Chromium is ~300–600 MB. SizeMagainst available memory — on a typical 16 GB developer machine,M = 2–3is a safe starting point; measure before going higher. - No new failure surface. Parallel families keep each family's own activity timeout / heartbeat / retry isolation and the exact partial-failure folding; a canceled source cancels the whole run. Long-run soak on a real workload remains an owner responsibility given the historical browser-GC incident.
Data-flow context — where each stage persists and how results reach the UI — lives in Operations & Events.
What The Operations Snapshot Can Prove
Runtime activity interception counts every active activity slot. It separately keeps an oldest-first, allowlisted detail list capped at 20 and an exact allowlisted-detail total, so activeSlots may legitimately exceed the number of rendered active items. The interceptor never reads activity arguments; unsafe identifiers are replaced with local opaque hashes, while only grammar-validated workflow/run references may remain readable.
Task-queue statistics are not another capacity total. They are approximate Temporal observations for the workflow and activity queue: pollers, backlog count/age, and add/dispatch rates. The capacity response preserves unsupported, unavailable, and stale states. The ETA estimator therefore refuses to divide domain work by nominal slots when runtime telemetry is stale or shared queue contention cannot be bounded.