arq-worker executes batches, and campaign-orchestrator decides when the next batch should run and when a campaign is finished. Neither serves HTTP. Both build from the same image as the API.
This page is about the two processes — what they run, how they coordinate, and how to scale them. Campaign semantics — states, retries, scheduling windows, the circuit breaker — are in Campaigns.
Two containers, one image
api, arq-worker, and campaign-orchestrator are all build: apps/api/Dockerfile. They mount the same code and differ only in command:
Their environment blocks differ too. The orchestrator gets the smallest set — FerretDB connection,
SECRET_KEY, REDIS_URL, and CAMPAIGN_BATCH_SIZE — because it never touches MinIO or places a call itself. The worker additionally gets the MinIO variables, INTERNAL_API_KEY, PROVIDER_AUTH_ENCRYPTION_KEY, and VOICE_SERVER_BASE_URL, because it dispatches real calls.
Redis is the only thing between them.
Neither worker calls the API over HTTP. Both import the same service layer and talk to FerretDB and Redis directly.
The ARQ worker
app/tasks/arq.py defines WorkerSettings, with these verified values:
max_jobs = 10 is the concurrency ceiling for one worker process: at most ten jobs run at once. It is not the campaign batch size — that is CAMPAIGN_BATCH_SIZE, default 10, from app/config.py.
The two registered functions live in app/tasks/campaign_tasks.py, and their names are pinned in the FunctionNames enum in app/tasks/function_names.py so producer and consumer cannot drift.
sync_campaign_source
Loads the campaign, picks a sync service by source_type (default csv), and pulls the rows in.
- Zero rows synced → the campaign goes straight to
completed, withsource_sync_status: "completed". - Rows synced → state becomes
runningand async_completedevent is published, which is what wakes the orchestrator. - Any exception → state
failed,source_sync_status: "failed",source_sync_errorset, a campaign log line appended, and the exception re-raised so ARQ records the failure.
process_campaign_batch
Calls campaign_call_dispatcher.process_batch(campaign_id, batch_size) with batch_size defaulting to settings.CAMPAIGN_BATCH_SIZE, then publishes the outcome.
The pool-exhausted path is a deliberate soft retry:
MAX_PHONE_POOL_ATTEMPTS = 3, counted in the campaign document under phone_number_pool_exhausted_attempts. Publishing batch_completed with zero processed rows makes the orchestrator schedule another batch, giving numbers time to free up. A batch that processes anything resets the counter.
The campaign orchestrator
CampaignOrchestrator is a long-lived asyncio process. run() starts two tasks and gathers them:
_listen_for_events()subscribes toCAMPAIGN_EVENTS_CHANNEL, defined as"campaign_events"inapp/constants/campaign.py._monitor_completion()sweeps running campaigns on a timer.
__init__:
Event handling
_handle_event() dispatches on the parsed event type from campaign_event_protocol.py:
_schedule_next_batch() re-reads the campaign, checks the state is running or syncing, checks the schedule window with _is_within_schedule(), checks the circuit breaker — pausing the campaign and publishing circuit_breaker_tripped if it is open — confirms there is pending or processing work, and only then enqueues process_campaign_batch.
Retries
_handle_retry_event() reads retry_config from the campaign document and honours enabled, retry_on_busy, retry_on_no_answer, retry_on_voicemail, max_retries, and retry_delay_seconds. Defaults live in DEFAULT_CAMPAIGN_RETRY_CONFIG in app/constants/campaign.py: two retries, 120 seconds apart, retrying busy and no-answer but not voicemail. A run past max_retries increments failed_rows instead. Retry runs get a derived source_uuid of {original}_retry_{n} and a parent_queued_run_id.
Completion detection
Every 60 seconds_check_stale_campaigns() walks campaigns in state running. For each one, it reschedules if a tracked batch has been in progress longer than 300 seconds, schedules a batch if there is pending work and no batch running, and otherwise tests _should_mark_complete(): no batch in progress, no pending or processing runs, and no activity for completion_timeout. When there is no in-memory activity timestamp it falls back to last_activity_at, last_batch_scheduled_at, or started_at on the document. Completion re-checks pending work once more before writing state: "completed" and publishing campaign_completed.
Redis keys
Both processes share one Redis instance, addressed byREDIS_URL.
Slot and rate-limit keys are manipulated by Lua scripts so acquisition is atomic across processes. See Call concurrency and rate limiting.
Scaling
The ARQ worker scales horizontally. ARQ hands each queued job to exactly one worker, andprocess_campaign_batch acquires concurrency slots through the atomic Lua scripts in rate_limiter.py. Run as many replicas as your provider concurrency allows; max_jobs = 10 multiplies per replica.
Failure and restart behaviour
Both containers carryrestart: unless-stopped, and both gate on ferretdb starting and redis passing its healthcheck. The worker additionally waits for api to start.
The orchestrator handles
SIGTERM and SIGINT, sets _running to false, cancels its task, unsubscribes, and closes the Redis client — so docker compose stop is a clean shutdown, not a kill.
The 60-second sweep is the only recovery path for events missed during an orchestrator restart. A campaign can therefore stall for up to a minute after the container comes back. If it stalls for longer, check the logs below before assuming the campaign is broken — see Campaign troubleshooting.
Logs to watch
A worker with no
Processing batch lines while a campaign sits in running usually means the orchestrator is not enqueueing — check that its container is up and subscribed before looking at the worker.
Related
- Campaigns — states, retries, scheduling, and the circuit breaker
- API (apps/api) — the same package, running as HTTP
- Call concurrency and rate limiting
- Running a campaign · Campaign troubleshooting