Skip to content

Review how the state and progress of jobs are tracked #1285

Description

@mihow

Summary

The platform tracks an async_api job's lifecycle in three independent
places — Job model in Postgres, AsyncJobStateManager in Redis, and
JetStream consumer state in NATS. Each has different failure modes, and
none alone catches every observed bug class. We should expose a unified
diagnostic view that triangulates all three.

This is a low-priority follow-up to #1276, which patched the immediate
"reaper REVOKES on clobbered Job.progress" loop by reading Redis
directly. Triangulation is the next layer up — forensic visibility and
catching adjacent bug classes that Redis alone can't see.

Why three sources of truth lie differently

Bug class DB Redis NATS Diagnostic signal
Clobbered progress (root cause of #1276) progress.is_complete() False drained ack_floor == delivered last-writer race in _update_job_progress
Lost-image stranding (#1247) stuck STARTED non-empty num_pending > 0, redelivery exhausted NATS gave up, Redis still tracks
Worker crash loop stuck STARTED non-empty num_redelivered climbing infrastructure issue
Result-handler swallow SUCCESS failed set non-empty drained error path bug
Hung in-flight worker progress 95% non-empty num_ack_pending > 0 for minutes worker stalled mid-batch

DB alone catches none reliably. Redis catches most. Redis+NATS catches all.

Proposal

Three pieces, mostly independent:

  1. Snapshot NATS state on terminal transition. On SUCCESS / FAILURE /
    REVOKED, persist consumer counters into Job.progress.diagnostics (or
    a new Job.diagnostics JSONB). Cheap. NATS consumer state is deleted
    post-completion — without a snapshot the forensic trail is gone.

  2. GET /api/v2/jobs/<id>/diagnostics endpoint returning structured
    triangulation:

    • db — Job.status, Celery state, progress.is_complete(), per-stage progress
    • redis — per-stage pending/total, failed set, all_tasks_processed, TTL
    • nats — stream/consumer, delivered, ack_floor, num_pending, num_ack_pending, num_redelivered
    • divergence — server-side computed flags (DB-vs-Redis, DB-vs-NATS, Redis-vs-NATS mismatches with severity)

    Cache 5–10s. NATS info is an async socket round-trip; cheap but not free.

  3. JobReconciler.diagnose() abstraction. Promote the tri-state
    AsyncJobStateManager.all_tasks_processed() added in fix(jobs): fix dangling jobs from going to revoked #1276 to a
    richer diagnose() returning the full triangulation. Both the
    diagnostics endpoint and the reaper guard call the same function.
    Single source of triangulation logic.

  4. UI diagnostic panel (admin / power-user). Keep default progress
    bars unchanged. Add expandable section:

    process     dispatched  acked   pending  failed  redelivered
                450/450     445     5        0       12
                ── NATS ──  ─NATS─  ─Redis─  Redis   NATS
    

    Plus header strip with consumer name, Redis TTL, divergence flags,
    and a "Force reconcile" admin button that re-runs the reaper guard
    on demand.

Tradeoffs

  • NATS info call needs JetStream credentials reachable from Django (already
    have these via the natsconn integration).
  • More numbers = more confusion for non-admins — gate the diag panel by
    permission.
  • NATS consumer state goes stale fast post-completion (consumers deleted).
    Snapshot-on-terminal mitigates this.
  • Adds a per-job NATS round-trip on endpoint hit. Cache + admin-only
    surface keeps cost bounded.

Recommended order

  1. Snapshot-on-terminal first — cheap, no UI, big forensic value.
  2. Diagnostics endpoint + JSON shape. Lets reaper + admin tooling consume
    it; no UI needed yet.
  3. Promote reaper to consume JobReconciler.diagnose() — replaces inline
    tri-state branch from fix(jobs): fix dangling jobs from going to revoked #1276.
  4. UI panel last, once endpoint shape is stable.

Related issues

Reference

Full breakdown captured in docs/claude/planning/job-state-triangulation.md
in the same PR series.

Activity

  1. changed the title [-]Job state triangulation: unify DB + Redis + NATS into a diagnostic view[/-] [+]Review how the state and progress of jobs are tracked[/+] on Apr 30, 2026
  2. added
    PSv2Async & distributed ML backend (PSv2): job state, NATS dispatch, result handling. Umbrella #515.
    on Jun 16, 2026
  3. mihow commented on Jul 29, 2026

    @mihow
    CollaboratorAuthor

    Claude says: A production sighting from 2026-07-29 relevant to this issue's bug-class table, plus a distinct failure mode the table does not yet cover.

    (Correcting my earlier version of this comment, which said "Redis was drained" and matched this to row 1. Both were wrong — the two halves of that sentence contradicted each other, and the correct row is a different one.)

    Where it sits in the table. The signature was: the database said the job was incomplete, Redis still held exactly one outstanding entry, and the NATS consumer was fully drained and acknowledged with ack_floor == delivered and num_pending = 0. That is not row 1 — row 1 requires Redis drained, and it was not; the single remaining entry is exactly why progress stopped ticking.

    It is closest to the lost-image stranding row, but differs in one field: num_pending was zero rather than greater than zero. So the message bus had genuinely finished and the only thing outstanding was a leaked tracker entry. That looks like a new variant worth adding to the table, since "Redis non-empty but the bus is fully drained" is diagnosable and the existing rows do not name it.

    The concrete outcome: one large job was revoked at 99.997 percent, a single image short, because that one leaked entry stopped the progress tick and the fixed no-progress threshold fired on an effectively finished job.

    A bug class the table does not cover. The same job was revoked twice earlier at zero percent while the log asserted that no workers had been seen for its pipeline in the past hour. That was false — a worker had logged 585 batches in that window.

    _log_worker_availability (ami/jobs/tasks.py:165-200, warning emitted at :200, called from run_job at :160) derives its claim from ProcessingService.last_seen, which is a heartbeat, not evidence of work. During a broker outage the heartbeat write path is itself a casualty, so the log reports a second symptom of the outage as though it were a fact about workers — and it points triage toward the workers, which were healthy and busy the whole time.

    That is distinct from the existing issues in this area, which all concern how service status is displayed. None of them says this log line can be actively wrong while looking authoritative. Worth a row of its own.

    Related: #1276 and #1241 are reopened, though with narrower claims than I first made for them — see the corrections on each. The infrastructure cause underneath the incident is in #1382.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    PSv2Async & distributed ML backend (PSv2): job state, NATS dispatch, result handling. Umbrella #515.enhancementNew feature or request

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions