Queue — durable или operational backlog единиц работы, ожидающих выполнения. Worker — исполнитель, который забирает допустимую работу, выполняет её в ограниченном контексте и фиксирует результат.
№42 Events & Triggers определяет смысл события и реакцию; №57 хранит/исполняет resulting work unit. №58 Scheduler создаёт сигнал/работу по времени; queue исполняет её, когда есть capacity. №59 Broker / Message Bus транспортирует сообщения/pub-sub между producers/consumers; №57 фокусируется на work backlog, claim/lease/ack. №63 Retry & Circuit Breakers позже владеет общей resilience policy; №57 реализует базовые retry/redrive semantics конкретной работы. №68 Durable Workflow хранит историю и состояние многошагового процесса; queue лишь доставляет отдельные activity/work units. №69 Distributed Reliability синтезирует delivery/idempotency/locks/DLQ на системном уровне.
Prerequisites: №10 State Management, №42 Events, №46 Observability, №50 Contracts, №55 Ingestion & Sync. Forward references: №58 Scheduler, №59 Broker, №60 Artifact Store, №63 Retry/Circuit Breakers, №64 Rate Limits/Budgets, №68 Durable Workflow, №69 Distributed Reliability.
REQUEST-TIME: producer может enqueue и сразу вернуть task/job reference. CONTROL PLANE: queue definitions, worker pools, concurrency, priorities, retry/dead-letter rules. DATA PLANE: job envelopes, leases, results/status. OFFLINE: re-drive, maintenance, capacity analysis, load/failure tests.
Success: one logical work item eventually reaches terminal success or explicit terminal failure. Retryable: transient worker/tool/provider failures. Permanent: invalid contract, permission/policy deny, deterministic unsupported input. Delivery: assume at-least-once unless proven otherwise. Idempotency: handler/effect key required for repeat delivery. Persist: job_id, payload ref/version, attempts, lease owner/expiry, status, result/error refs, timestamps. Trace: enqueue→claim→attempt→result/retry/dead.
№57 не владеет business event semantics, calendar/time scheduling, pub-sub topology, workflow history/replay, global retry/circuit policy или artifact/blob storage. Она владеет THE LIFECYCLE OF EXECUTABLE WORK UNITS AND THE WORKERS THAT CLAIM THEM.
HTTP/user flow ждёт OCR, ingestion, embeddings, export или long tool call.
Crash = работа потеряна, spike = request storm.
Producer durable записывает work item и получает job_id.
Backlog становится наблюдаемым и управляемым.
Workers забирают столько работы, сколько система способна безопасно выполнить.
Fetch, parse, chunk, embed, reindex — естественные background work units.
Exports, conversions, browser jobs, large research stages, sandbox tasks.
1000 webhook events приходят за минуту; workers обрабатывают с controlled concurrency.
Provider temporarily unavailable; job waits/retries independently of user connection.
Page parsing, embedding chunks, evaluation cases, fan-out jobs.
Analytics, audit enrichment, thumbnails, secondary indexes where consistency allows.
10–100 ms deterministic transform внутри одного процесса не выигрывает от queue hop.
Если operation короткая и result нужен для продолжения request, synchronous call проще.
Пока workload мал, обычная function + persisted state может быть лучше дополнительной инфраструктуры.
{
"contract": "job.v1",
"job_id": "JOB-...",
"job_type": "document.parse",
"tenant_id": "tenant_A",
"payload": {
"source_ref": "doc://...",
"parser_profile": "documents.v2"
},
"priority": 50,
"idempotency_key": "parse:doc123:v42:p2",
"attempt": 0,
"max_attempts": 4,
"timeout_s": 120,
"not_before": null,
"trace_id": "TRACE-...",
"created_at": "...",
"contract_version": "1.0"
}Минимально полезны:
Worker не должен угадывать retry semantics из текста payload.
200 MB PDF, full trace или giant model context копируется в message/job row и на каждый retry.
Queue хранит artifact_ref, source_ref, state_ref + version/hash. №60 Artifact Store хранит bytes.
Можно claim, если not_before прошёл и queue/policy допускает.
Worker получил временное право выполнить job.
Attempt выполняется; heartbeat может продлевать lease.
Terminal success; result persisted.
Retryable failure; next_attempt_at установлен.
Attempts exhausted / poison / manual intervention.
Работа больше не должна исполняться.
READY
↓ claim atomically
LEASED / RUNNING
├─ success ───────────────→ SUCCEEDED
├─ retryable failure ─────→ RETRY_WAIT → READY
├─ permanent failure ─────→ DEAD
├─ attempts exhausted ────→ DEAD
├─ cancel requested ──────→ CANCELLED / cooperative stop
└─ worker disappears
↓ lease expires
READY again
Worker atomically marks job leased by worker_id until lease_until.
Для длинного job worker может продлевать lease, но не бесконечно без upper bounds.
Если heartbeat исчез, job снова становится claimable.
| Action | Meaning | When |
|---|---|---|
| ACK / SUCCESS | Required effects and result state are durable. | Job becomes SUCCEEDED. |
| NACK / RETRY | Attempt failed transiently. | Increment attempt, calculate next_attempt_at. |
| FAIL PERMANENT | Same input will not succeed unchanged. | DEAD / terminal error; no blind retries. |
| ABANDON | Worker lost/crashed without final transition. | Lease expiry makes job visible again. |
handle(job):
validate_contract(job)
prior = lookup_effect(job.idempotency_key)
if prior.completed:
return prior.result
input = load_refs(job.payload)
result = execute_deterministically(input)
persist_result(
idempotency_key=job.idempotency_key,
result=result
)
return result«Мы никогда не доставим job дважды» — слабый invariant.
| Error | Retry? | Typical behavior |
|---|---|---|
| Network timeout / transient 5xx | YES | Exponential backoff + jitter, bounded attempts. |
| Rate limited | YES, later | Honor Retry-After / reduce concurrency. |
| Invalid input schema | NO unchanged | Permanent fail / producer repair. |
| Permission / policy denied | NO unchanged | Do not retry to bypass security. |
| Version conflict | CONDITIONAL | Reload current state and re-evaluate. |
| Model/provider quality failure | CONDITIONAL | Only if repair/escalation policy expects improvement. |
Corrupt file, unsupported schema, deterministic bug, missing permission.
После max attempts/permanent error job переводится в DEAD с structured error.
После исправления code/config/input можно создать new attempt/redrive с audit trail.
Queue grows temporarily; workers catch up without overloading downstream.
Admission control, producer rate reduction, lower priority shedding, quotas.
Add workers only if provider/DB/CPU quotas can absorb more concurrency.
Защищает CPU/RAM/database/shared infrastructure.
OCR может иметь concurrency 2, embeddings 16, emails 4.
Один крупный customer не занимает весь worker pool.
Connector/provider concurrency respects rate and quota limits.
| Mechanism | Use | Risk |
|---|---|---|
| Simple FIFO | Homogeneous jobs. | Large slow jobs block urgent work. |
| Priority value | Interactive/high-risk recovery ahead of batch. | Low-priority starvation. |
| Separate queues/pools | Different resource classes / SLAs. | More operational complexity. |
| Tenant fair-share | Multi-tenant workloads. | Needs explicit fairness policy. |
| Aging | Old low-priority jobs gradually rise. | Implementation complexity but prevents starvation. |
Parsing, compression, deterministic processing.
Network-heavy API calls with controlled concurrency.
Embeddings/local inference jobs with scarce accelerator capacity.
Jobs requiring special containment/runtime profiles.
Worker сам запрашивает next eligible job. Естественно для DB queues и многих task queues; concurrency локально controllable.
Infrastructure вызывает consumer/доставляет message. Требует обработать duplicate delivery, timeout/ack semantics и consumer availability.
Не claim-ить job; transition to CANCELLED.
Worker periodically checks cancellation or receives signal and safely stops at cancellation points.
Cancellation may require compensating action, not mere process kill.
Producer authorized for job_type/tenant/resource; queue payload validated.
Worker pool получает только capabilities, нужные его job types.
Для sensitive side effects permission/policy may need re-evaluation at execution time, not only enqueue time.
Ready/retry/running/dead counts by job type/tenant.
Лучший сигнал sustained backlog и SLA risk.
enqueue_at → start_at p50/p95.
start_at → terminal result.
arrival rate vs completion rate.
retry distribution and retry causes.
expired leases / lost workers / heartbeat delays.
Dead jobs count, age, owner, redrive status.
| Signal | Useful? | Caveat |
|---|---|---|
| Queue depth | YES | Needs job cost context; 100 tiny != 100 huge. |
| Oldest age | VERY | Strong SLA/backlog signal. |
| Arrival/completion ratio | VERY | Shows whether backlog will grow. |
| Worker CPU/RAM | YES | Only if worker resource is bottleneck. |
| Provider 429 / DB saturation | CRITICAL | May mean scale DOWN concurrency, not up. |
Persistent executable jobs, claim/lease/ack, concurrency, retry/dead state.
At 09:00, every hour, after delay, at deadline — creates signal/job at time boundary.
Topics/pub-sub/routing/consumer groups/transport semantics between producers and consumers.
«Parse document 42», «embed chunk batch 7», «send approved message». Worker owns one activity attempt.
«Ingest source → wait → fan-out parsing → aggregate → approval → publish» с history/replay/timers/checkpoints.
Same job delivered twice → one logical effect.
Worker dies; lease expires; job recovered.
Retry does not duplicate external effect.
Timeout leads to clear failure/retry state, no orphan process.
Ends DEAD after policy-defined attempts.
Queue absorbs burst without uncontrolled downstream overload.
Graceful shutdown doesn't lose acknowledged work.
Wrong tenant/priority/resource cannot be injected through job payload.
Главный backlog/SLA signal by queue/job type.
enqueue → worker start p50/p95.
Completed jobs/sec or min, by type/pool.
Created jobs rate vs completion rate.
Attempts/job and retry causes.
Dead count, oldest age, unresolved ownership.
Expired leases / worker loss frequency.
Logical side effects duplicated after retry. Target: zero.
jobs(
job_id uuid primary key,
job_type text,
tenant_id text,
payload_json jsonb,
status text,
priority int,
idempotency_key text,
attempt int,
max_attempts int,
not_before timestamptz,
lease_owner text,
lease_until timestamptz,
created_at timestamptz,
started_at timestamptz,
finished_at timestamptz,
result_ref text,
error_json jsonb
)
claim:
SELECT ...
WHERE status='READY'
AND not_before <= now()
ORDER BY priority DESC, created_at
FOR UPDATE SKIP LOCKED
LIMIT 1
then:
set LEASED + lease_owner + lease_untilПереходить на dedicated queue/broker, когда throughput, latency, routing, isolation или operational requirements реально выходят за возможности DB-based approach.
| Signal | Potential next step |
|---|---|
| Very high job throughput / DB contention | Dedicated task queue/broker or partitioned queue architecture. |
| Complex routing / many independent consumers | Message broker / topics — see №59. |
| Strict delayed/timed delivery at scale | Dedicated scheduler/timer infrastructure — see №58. |
| Complex multi-day multi-step workflows | Durable workflow engine — see №68. |
| GPU/heterogeneous pools | Queue routing per capability/resource class. |
| Multi-region failover / consistency requirements | Distributed reliability design — see №69. |
| Вопрос | Ответ |
|---|---|
| Стоит ли реализовывать? | Да, когда появляются durable/background/long-running jobs. Для простого synchronous prototype — необязательно. |
| Separate Component? | YES. Queue/worker execution — отдельная production responsibility, хотя MVP может быть таблицей PostgreSQL и worker process. |
| Минимум 80% ценности? | Job contract, durable READY state, atomic claim/lease, idempotent handler, bounded retries, dead state, concurrency caps, metrics. |
| Когда overkill? | Kafka/RabbitMQ/Kubernetes worker fleet для пары фоновых jobs в час, которые надёжно выполняются через Postgres. |
| Trigger? | Work must survive request/process restart, absorb bursts, run later, retry independently or use separate resource pool. |
| Как измерить uplift? | Request latency reduction, completion reliability, oldest age/SLA, throughput, retry/dead rate, duplicate-effect rate, operational incidents. |
| Можно ли rule/tool/code вместо LLM-agent? | Полностью. Queue and worker lifecycle are deterministic infrastructure. LLM may run inside a worker job, but does not implement the queue. |
Передавать small job envelope + references.
At-least-once → handlers/effects idempotent.
Claim имеет expiry, чтобы crash не блокировал job навсегда.
Success фиксируется последним, после required effects.
Permanent/security failures не должны бесконечно redrive.
Protect downstream providers, DB, GPU and tenants.
Oldest work and flow rates expose real backlog pressure.
Scheduler №58 and Broker №59 are different responsibilities.
PostgreSQL queue is valid architecture until requirements prove otherwise.
USER / EVENT / SYNC / WORKFLOW
↓
CREATE WORK UNIT
↓
VALIDATE JOB CONTRACT
↓
ENQUEUE
job_id
type
tenant
priority
idempotency_key
payload_refs
↓
READY BACKLOG
↓
WORKER POOL
capability
resource class
concurrency cap
↓
ATOMIC CLAIM / LEASE
↓
EXECUTE ATTEMPT
↓
├─ SUCCESS
│ ↓
│ persist result/effect
│ ↓
│ ACK → SUCCEEDED
│
├─ RETRYABLE
│ ↓
│ attempt++
│ backoff / not_before
│ ↓
│ READY
│
├─ PERMANENT
│ ↓
│ DEAD
│
└─ WORKER CRASH
↓
lease expires
↓
READY AGAIN
CONTROL:
concurrency
priorities
tenant fairness
resource limits
retry ceilings
graceful drain
OBSERVE:
depth
oldest age
queue wait
execution time
arrival/completion rate
retries
lease expiry
dead jobs
CORE PRINCIPLE:
THE QUEUE DOES NOT
"RUN THE AI".
IT HOLDS DURABLE,
EXPLICIT UNITS OF WORK.
THE WORKER DOES NOT
"OWN THE PROCESS".
IT CLAIMS ONE UNIT,
EXECUTES IT UNDER A CONTRACT,
AND RETURNS A DURABLE RESULT.
ASSUME DUPLICATES.
BOUND CONCURRENCY.
ACK LAST.
START SIMPLE.
A. Related ECC ideas. Context-as-cache, scoped memory, lifecycle hooks, selective capabilities, feature flags, deterministic enforcement, provider-neutral adapters and eval-gated learning are applied only where relevant to №57 Queues & Workers.
B–E. Existing boundary and placement. The existing conceptual boundary, class PRODUCTION, default CONDITIONAL and owner Production Fabric remain authoritative. Runtime/control/data/offline placement is unchanged; durable state stays outside model context.
F–H. Hooks and contracts. Use bounded PRE_MODEL/POST_MODEL, PRE_TOOL/POST_TOOL, CHECKPOINT and TASK_COMPLETED events as applicable. Illustrative fields and canonical contracts are defined in NEW_CONTRACTS_SPEC.md; no universal schema is implied.
I–J. Security and evaluation. Host-side schema, permission, secret, budget, idempotency and audit checks take precedence over LLM output. Optional mechanisms require a feature flag and WITH/WITHOUT ablation; measure quality, acceptance, correction, latency, cost, escalations and severe errors.
K–L. Task profiles and cross-references. A TaskProfile selects the relevant skill, tool/context slice, memory scope and enforcement profile independently from FAST/STANDARD/DEEP. See cross-reference map, hook spec and ablation plan. Provider adapters remain outside the core.