CHK-02: Машина состояний задания
Путь задания
Каждое задание в ChokaQ проходит детерминированный путь через набор состояний. Понимание этого пути - ключ к пониманию всей системы.
Состояния
| Status | Value | Location | Meaning |
|---|---|---|---|
| Pending | 0 | JobsHot | Ожидает в очереди и готово к выборке |
| Fetched | 1 | JobsHot | Загружено в буфер памяти воркера и ждет свободный execution slot |
| Processing | 2 | JobsHot | Активно выполняется, heartbeat обновляется |
| Succeeded | 3 | JobsArchive | Успешно завершено и перемещено в archive |
| Failed | 4 | JobsDLQ | Завершилось постоянной ошибкой и перемещено в DLQ для разбора |
| Cancelled | 5 | JobsDLQ | Отменено администратором и сохранено в DLQ/history workflow |
| Zombie | 6 | JobsDLQ | Воркер упал, heartbeat истек |
📝 Примечание
В режиме SQL Server статусы 3-6 являются виртуальными. У задания фактически нет колонки Status со значением 3: строка просто перемещается в таблицу JobsArchive. Эти enum-значения нужны для SignalR-уведомлений и In-Memory storage mode.
Полный lifecycle

Фаза 1: Enqueue
Когда вызывается IChokaQQueue.EnqueueAsync():
await _queue.EnqueueAsync(
new SendEmailJob("user@example.com", "Welcome!"),
priority: 20, // Higher = processed first
delay: TimeSpan.FromMinutes(5), // Optional delayed execution
idempotencyKey: "email-user-123" // Optional duplicate prevention
);Что происходит внутри:
- Генерируется уникальный ID в формате
ULID, удобном для сортировки по времени. - Job DTO сериализуется в JSON.
- Idempotency check: если ключ уже существует, возвращается существующий job ID, вставка пропускается.
- В
JobsHotвставляется строка соStatus = 0 (Pending). - Если задан delay, заполняется
ScheduledAtUtc. - Записывается метрика
chokaq.jobs.enqueued. - Через SignalR отправляется уведомление.
Idempotency Guard
-- Check BEFORE insert to avoid unique constraint violations
SELECT [Id] FROM [chokaq].[JobsHot]
WHERE [IdempotencyKey] = @Key;
-- If found → return existing ID, done
-- If not → proceed with INSERTПлюс уникальный filtered index как страховка:
CREATE UNIQUE NONCLUSTERED INDEX [IX_JobsHot_Idempotency]
ON [chokaq].[JobsHot] ([IdempotencyKey])
WHERE [IdempotencyKey] IS NOT NULLГраницы ответственности:
- Встроенная enqueue-idempotency работает только в
JobsHot. - Она предотвращает дублирование активной работы с одним business key.
- Когда исходное задание уходит в Archive или DLQ, тот же ключ можно поставить в очередь снова как новую логическую попытку.
- Долгоживущая память о том, что работа уже была завершена, относится к completion marker из optional idempotency middleware, а не к индексу Hot-очереди. Этот marker не replay'ит результат handler'а и не создает exactly-once side effects.
Фаза 2: Fetch (Pending -> Fetched)
Fetcher loop в SqlJobWorker выполняется с polling-интервалом:
// Simplified from SqlJobWorker.cs
while (!ct.IsCancellationRequested)
{
var batch = await _storage.FetchNextBatchAsync(
workerId: _workerId,
batchSize: _options.BatchSize,
allowedQueues: activeQueues
);
foreach (var job in batch)
{
// Push into bounded Channel<T> (Prefetch Buffer)
await _channel.Writer.WriteAsync(job, ct);
}
await Task.Delay(_options.PollingInterval, ct);
}Fetch query атомарно:
- Выбирает top N pending-заданий с учетом priority, schedule и bulkhead-лимитов.
- Устанавливает
Status = 1 (Fetched)и назначаетWorkerId. - Возвращает заблокированные строки.
Фаза 3: Process (Fetched -> Processing)
Consumer loop читает задания из Prefetch Buffer:
// For each job from the channel:
await _semaphore.WaitAsync(ct); // DynamicConcurrencyLimiter — concurrency control
try
{
await _processor.ProcessAsync(job, ct);
}
finally
{
_semaphore.Release();
}Внутри ProcessAsync():
- Проверка Circuit Breaker -> разрешен ли запуск этого job type.
- Mark as Processing ->
Status = 2, установкаHeartbeatUtcиStartedAtUtc. - Запуск heartbeat task -> параллельно обновляет
HeartbeatUtcкаждые N секунд. - Dispatch в handler -> вызов через скомпилированный delegate на Expression Tree.
- Middleware pipeline -> оборачивает handler onion-model стеком middleware.
- Result handling -> путь Success, Fatal или Transient.
Heartbeat Mechanism
Пока задание выполняется, параллельная task поддерживает его живым:
// Simplified heartbeat logic
_ = Task.Run(async () =>
{
while (!jobCt.IsCancellationRequested)
{
await Task.Delay(
RandomBetween(
options.Execution.HeartbeatIntervalMin,
options.Execution.HeartbeatIntervalMax),
jobCt);
await _storage.KeepAliveAsync(jobId, jobCt);
}
});Это не дает ZombieRescueService заархивировать долгую, но здоровую работу.
Фаза 4a: Success (Processing -> Archive)
var durationMs = stopwatch.Elapsed.TotalMilliseconds;
await _storage.ArchiveSucceededAsync(job.Id, durationMs);
_breaker.ReportSuccess(job.Type);
_metrics.RecordSuccess(job.Queue, job.Type, durationMs);Задание атомарно перемещается в JobsArchive, вместе с длительностью выполнения.
Фаза 4b: Transient Failure (Processing -> Pending, retry)
if (!IsFatalException(ex) && job.AttemptCount < options.Retry.MaxAttempts)
{
var backoffMs = CalculateBackoffMs(job.AttemptCount + 1);
var nextRun = DateTime.UtcNow.AddMilliseconds(backoffMs);
await _storage.RescheduleForRetryAsync(
job.Id, nextRun, job.AttemptCount + 1, ex.Message);
}Задание остается в JobsHot, но:
Statusвозвращается в0 (Pending);AttemptCountувеличивается;ScheduledAtUtcустанавливается на будущее время с учетом backoff;- задание не будет выбрано до наступления
ScheduledAtUtc.
Фаза 4c: Fatal Failure / Exhaustion (Processing -> DLQ)
await _storage.ArchiveFailedAsync(job.Id, ex.ToString());
_breaker.ReportFailure(job.Type);
_metrics.RecordFailure(job.Queue, job.Type, ex.GetType().Name);Задание атомарно перемещается в JobsDLQ, а полные сведения об exception сохраняются в ErrorDetails.
Фаза 5: Resurrection (DLQ -> Hot)
Категории ошибок, которые приводят в DLQ, описаны в Failure Taxonomy. Поведение shutdown и cancellation вокруг активной работы описано в Graceful Shutdown и Job Context And Cancellation.
Через The Deck dashboard или программно:
await _storage.ResurrectAsync(
jobId: "job-123",
updates: new JobDataUpdateDto
{
Payload = fixedJson, // Optional: fix the broken payload
Tags = "retried,manual", // Optional: add audit tags
Priority = 30 // Optional: boost priority
},
resurrectedBy: "admin@company.com"
);Задание возвращается в JobsHot с:
Status = 0 (Pending);AttemptCount = 0, то есть fresh start;- измененным payload/tags/priority, если они были переданы.
Дальше: как Bulkhead Isolation не дает тяжелым заданиям вытеснить легкие.
