Repository navigation
New async & distributed ML backend (aka "PSv2") #515
Description
Activity
Controller can be tested here:
https://preview.ami.ecoscience.dk/swagger-ui/index.html#/pipeline-controller/createRequestExample request:
{ "projectId": "ecos", "jobId": "ea12ac70-288c-11ef-9ca5-00155d926c42", "sourceImages": [ { "id": "NScxODE3NzEyMwo=", "url": "https://anon.erda.au.dk/share_redirect/DSWDMAO70L/ias/denmark/DK1/2023_07_05/20230705000135-00-07.jpg", "eventId": "1234" } ], "pipelineConfig": { "stages": [ { "stage": "OBJECT_DETECTION", "stageImplementation": "flatbug" }, { "stage": "CLASSIFICATION", "stageImplementation": "mcc24" } ] }, "callback": { "callbackUrl": "http://127.0.0.1:8080/example/callback", "callbackToken": "1234" } }The callback can be inspected using Ngrok locally:
- run
ngrok http 2222 - update the
callbackUrlin the sample request to the generated ngrok URL - open web interface http://localhost:4040
- trigger the request to the ML API controller
Image that causes a truncated request
https://static.dev.insectai.org/ami-trapdata/Panama/E43B615A/20231113004009-snapshot.jpg- run
There are a few Refactorings I'm considering in the PipelineController.
- Version on Detections in the callback is not set. The reason is that it's not obvious which value to set it to. The bounding box is done by flatbug, each classification is done by another classifier and cnn-features and crops might be from other stages.
Stagein the PipelineConfig is currently unused by thecontroller. I originally included it because I thought different stages might be handled differently but I think the current solution is cleaner - every stage use the same interface. One problem with handling different stages differently is that one stage implementation may perform multiple operations on the image.ImageCrop- Originally I thought the pipeline asObjectDetection->ImageCrop->Classification. But I realized that we don't need to crop in order to do classification. Actually it performs a lot better if the entire sourceimage is used with bounding boxes. So maybe cropping is post-pipeline step? It can also be added as a stage and set thecropUrlon the detection.
@kaviecos Have you considered adding a status endpoint on the controller? It would be nice if a callback is missed, or if the job is taking a long time, for the client to be able to request the status of a request. PROCESSING, FAILED, WAITING, 2/6 complete, etc.
Also will you document the types of failures and what those will look like in a callback?
I just added these as subtasks as well.
NOTES from call with @mihow and @kaviecos 2024-08-20
- CNN features are stored to a generic "URI" field (could reference a remote storage, or s3:// or vector database or file)
- Add tests for dummy detector that returns many detections, and no detections - in the controller repo but also a good idea in the AMI Platform to add tests with mock responses.
- Consider seeing how many 7mb requests the controller can take and send back (to callback) (load test)
- Michael consider attempting a dev setup in one compose file for local dev & testing of the AMI Platform (single stage everything & multi stage setups)
- If fetching the same image multiple times is a problem, then our current recommendation is to put a proxy cache in front of the local network where the stages are running.
- Consider a cancellation method
- Change single source image to array of images, even though we are using one
- Add "stageParams": {} to the controller request, which are then passed via URL to the stage implementation requests.
- Update the registry to support a "*" projectID that any project can use
- Think about an endpoint for checking what stages are registered. Use the projectID to determine if it's a private or shared stage implementation.
- Think about a way to see if a stage implementation is online or when it was last seen - Add some sort of basic healthcheck endpoint (see kubernetes /livez & /readyz pattern)? or can we check when a message was last taken from the queue by a stage?
docs for OpenAI
https://platform.openai.com/docs/guides/batch/getting-started@kaviecos Let's also keep in mind that scenario where we have stage implementations running behind a firewall and cannot have a publicly accessible endpoint. Our major providers (Compute Canada and the UK's JASMIN) both have way more compute available to us, but not on persistent VMs like we are using now. For a big job like Biodiversa+'s data, we can request a bunch of GPUs and provide the docker container to run via a SLURM job scheduler. Each job can access the internet, so can pull from the queue and send back results. But they won't have a publicly accessible endpoint (unless we can use a tunnel?)
@kaviecos Let's also keep in mind that scenario where we have stage implementations running behind a firewall and cannot have a publicly accessible endpoint. Our major providers (Compute Canada and the UK's JASMIN) both have way more compute available to us, but not on persistent VMs like we are using now. For a big job like Biodiversa+'s data, we can request a bunch of GPUs and provide the docker container to run via a SLURM job scheduler. Each job can access the internet, so can pull from the queue and send back results. But they won't have a publicly accessible endpoint (unless we can use a tunnel?)
@mihow The current implementation actually allows for consumers hosted on other servers. The requirements are that they can make an outbound connection to the RabbitMQ server (port 5672). Then all communication with the controller can go through RabbitMQ. Of course they also need to be able to access the source images.
When it comes to monitoring this also means that we need to monitor both the consumers and the stage implementations. And if we cannot make inbound http requests then we need to think that into the monitoring solution. Maybe a push-based solution (last_seen as you mentioned). It would be nice to know more about the restrictions - like, is it even possible to connect to RabbitMQ?
- modified the milestones: ML Pipeline Ready for Production, ML Pipeline Enhancements, More Models & Backends
on Dec 19, 2024 2 remaining items
Goal: Have a draft of the a working celery worker by the end of the month (July)
Here are some notes from a conversation in Slack
The current processing method was a placeholder method that is quite simple. It sends a single request per image batch and waits for the response. There is no parent/coordinator process that keeps track of a job and knows which images remain (and nothing to know if the job should be auto-restarted, as you suggest). In the current method, multiple processing services can't work together to finish the job in parallel.
The new implementation will be an asynchronous publisher/subscriber model. Which will work like so:
A new job is created, and a parent job will keep track of all the images that need to be processed with the configuration you selected.- Each image is added to the processing queue
- Processing services look at the queue for images to be processed and save the raw results are added to the ingestion queue for Antenna (this I believe solves the biggest thorn in the current system, the raw results in the synchronous return requests)
Antenna processes results in the results queue (converts raw results to Occurrences, Taxa, etc.) - The parent job monitors the queue until the job is complete. Handles failures in the way that is configured.
This may sound complex, but I will use the existing framework that we use for other background tasks called Celery. I setup the AMI Data Companion like this originally, and it works efficiently on a queue of local files. But the implementation for the web that you run with ami api was meant to be a temporary solution to connect Antenna to the ADC.
We can use our existing PipelineRequest & PipelineResponse objects, but use them to pass serialized data to & from the celery workers. Something like:
def process_images(source_images=[]): prediction_request = PipelineRequest( source_images=[...], pipeline_config={...}, ) # Submit task to celery queue as an argument task_id = submit_prediction_task(prediction_request) logger.info(f"Prediction task submitted: {task_id}") # Instead of sending the request directly as we do now # requests.post(endpoint, data=prediction_request)def submit_prediction_task(prediction_request: PipelineRequest) -> str: """Submit prediction task to appropriate celery queue.""" queue_name = get_pipeline_queue_name(prediction_request.pipeline_config['slug']) task_result = predict_image.apply_async( args=[prediction_request.dict()], queue=queue_name, routing_key=queue_name, # research routing_key vs. queue ) return task_result.id- Add rabbitmq to Antenna local, use as CELERY_BROKER_BACKEND
- Add new celery task called
process_pipeline_request(PipelineRequest) -> PipelineRequestin Antenna - Add new celery task with same name & signature to processing_service
- Change
process_images()in antenna to submit PipelineRequest to queue. Callprocess_pipeline_request.apply_async(args=[pipeline_request], queue=queue_name, routing_key=routing_key)(queue name or routing key should be specific to each pipeline) - Can queue names be dynamic while app is still running? When we register new pipelines in Antenna. Perhaps the queue_name needs to be static ("pipeline_requests") but the
routing_keycan be dynamic? (use the pipeline slug) - Can a worker publish a new task? save_results
- Make sure antenna DOES NOT try to consume
process_pipeline_requestswith it's dummy task signature - Design & implement a parent task to watch the entire job that monitors all image processing & saving tasks and updates the Job status.
- Could consider using callbacks from the processing service to Antenna, when processing is complete.
Current implementation
# Current "parent job" loop for image_batch in source_images_batches pipeline_request = PipelineRequest(source_images=image_batch) result = requests.post(endpoint, data=pipeline_request) save_results_task = save_results.apply_async(result) for task in save_results_task: check status # Antenna worker def save_results(data: PipelineRequestResults): creates detections, occurrences, etc from json, save to Django databaseNew implementation
# Antenna publisher (main Django app) for image_batch in source_images_batches pipeline_request = PipelineRequest(source_images=image_batch) task = process_pipeline_request.apply_async(args[pipeline_request]) tasks_to_watch.append(task) for task in tasks_to_watch: check status, update job status in UI until done # Processing service worker def process_pipeline_request(data: PipelineRequest): process images save_results.apply_async() def save_results(): pass # Antenna worker def save_results(data: PipelineRequestResults): creates detections, occurrences, etc from json, save to Django database def process_pipeline_request(): passv1
keep API as-is
add celery to processing service
Antenna callsprocess_pipeline_request.apply_async()
processing service callssave_results.apply_async()
Antenna saves the results as normalv2
Processing service can operate without public API
rename default queue to antenna queue
save results can go to the antenna queue
is a separate queue required to ensure save results is high priority?Next steps:
- Confirm through more testing that tasks are queued and processed properly (i.e. routing key, flower, etc)
- In the processing service, refactor to include
processing.pywhich contains the process logic; don't make http request, just call a function - Main goal: we don't want any long running tasks!
MLJob.run()shouldn't endlessly poll for status; we should be able to callcheck_ml_job_status()anytime to quickly check status oftask_ids--> See comment# TLDR; MLJob.run() just submits processing tasks + schedules celery beat; celery beat responsible for everything after
a. Add a new field for storing all task ids (process_pipeline_requestandsave_results) to the MLJob model. An array field of strings calledtask_ids
b. Create a separate function (check_ml_job_status) which celery beat runs periodically
c.check_ml_job_statusshould still watch the status ofprocess_pipeline_requesttasks; but only use the status to update the UI; shouldn't do anything with the results.
d. The processing service should produce thesave_resultstasks, queue it into the antenna queue (maybe change this to ensure higher priority?), and then the celerybeat watches for these tasks. Celerybeat will watch for any tasks with a parent id equal totasks_ids(celery might already attach parent ids).
GOAL for next meeting: start a job, process images, and see the job status change to success (only after all sub-tasks
process_pipeline_requestandsave_resultsall complete)Next Tasks (08/04)
- Add a celery worker to the main antenna stack
- test with larger batches
- Write some tests to consistently check my work
- make it easier to see the job progress
- minimize the logs (especially important for the jobs with many images)
- log/publish the job's summary stats, or give more details (i.e. task was picked up by this worker / status of each task / time per task)
Edge cases to consider:
- Cancelling the job; need to remove the subtasks too? and update the job progress to the current state
- If the celery workers / processing services crash, we should immediately delete the periodic task and clean up the subtasks?
- need to make the timeouts adjustable based on expected job time and processing service (or dynamically modify based on demand?)
- changed the title
[-]New async & distributed ML backend [/-][+]New async & distributed ML backend (aka "PSv2")[/+]on Mar 2, 2026 - addedPSv2Async & distributed ML backend (PSv2): job state, NATS dispatch, result handling. Umbrella #515.Async & distributed ML backend (PSv2): job state, NATS dispatch, result handling. Umbrella #515.
on Jun 16, 2026 Claude says: (posted on Michael's behalf — consolidated status map of the PSv2 effort, assembled from the full
PSv2-labeled issue/PR set.)PSv2 — implementation history and current status
This is a consolidated map of the async + distributed ML backend ("PSv2") work: how it was built, what is fixed, what was attempted and abandoned, what is still open, and the root issues that look most critical before async processing can become the default path. It was assembled by reading every issue and PR carrying the new
PSv2label (~120 items) and following the cross-references between them. Where a claim is an interpretation rather than a settled fact, it is flagged as such.All PSv2 work is now tagged with the
PSv2label for easier tracking.How we got here (the arc)
Phase 0 — controller/push era (spec). PSv2 began as an external pipeline controller designed at Aarhus (this issue, plus #505, #516, #517, #538, #539, #540). It was push-based and lived outside Antenna. Cancellation (#539) and stage monitoring (#538) were specified here and are still open — they resurface years later.
Phase 1 — in-Antenna push baseline. #632 landed V1: register a processing service, pick the lowest-latency online one, push work to it. #1060 raised the RabbitMQ ack timeout so long sync jobs stopped dying at 30 minutes. This is the synchronous baseline PSv2 replaces.
Phase 2 — the pull pivot (design). A deliberate move from push to pull: #968 (pull API spec), #969 (reference impl), #970 (implement), #971 (scheduler). Rationale, recorded in #987: per-task tracking for resilience, horizontal scale, and — critically — no public endpoint required, so workers can sit behind HPC/university firewalls. #1046 landed the API scaffold with the
tasks()/result()endpoints stubbed.Phase 3 — core implementation. #987 is the keystone: NATS JetStream chosen (after rejecting RabbitMQ/Beanstalkd for lacking the disconnected pull/ack semantics we need), task state + queue managers, real queue/pull/ack/result. Shipped disabled by default behind the
async_pipeline_workersflag for extended testing. #1109 added the NATS host alias for deployment.Phase 4 — registration, state decoupling, cleanup. Pipeline registration for pull mode (#1086 → #1076,
endpoint_urlnullable). Decoupling the Celery task state from the job state — the recurring hard problem — seeded by #1062, framed in #1084, solved in #1114 (a job is SUCCESS only when all stages complete, not when the queueing task returns), hardened in #1125. Routing via the newJob.dispatch_modeenum (#1111 → #1118). Resource cleanup on completion (#1083 → #1113).Phase 5 — progress, hardening, cancellation. Progress parity with the sync path (#1121). Integration-test hardening after real worker flooding (#1135, #1142, with test infra #1129/#1143). Service "live" status groundwork (#1122 → #1146). Cancellation was attempted in #1144 but closed unmerged and reworked later.
What is fixed and merged
- Architecture and contract: NATS pull model (Distributed Worker Architecture for ML Processing (Processing Service V2) pt. 1 #987), pipeline registration (PSV2: endpoint to register pipelines #1076),
dispatch_moderouting (fix: PSv2: Workers should not try to fetch tasks from v1 jobs #1118), resource cleanup on completion (PSv2: Implement queue clean-up upon job completion #1113), Celery-vs-job state decoupling (fix: Properly handle async job state with celery tasks #1114/PSv2 cleanup: use is_complete() and dispatch_mode in job progress handler #1125). - Reliability fixes that landed: finalize a job when zero images need processing (fix(jobs): finalize async job when 0 images need processing #1191); ack-after-save ordering to stop stranded STARTED jobs plus a counter-inflation guard (fix(jobs): prevent jobs from hanging in STARTED state with no progress #1234); stop progress being coerced to 100% and tripping premature FAILURE, and retune the reaper from 3 days to 10 minutes (fix(jobs): don't fail jobs prematurely, but fail stalled jobs sooner (3 days -> 10 min) #1235); reconcile NATS-lost images so a job lands SUCCESS/FAILURE instead of a whole-job REVOKE (fix(jobs): mark failed or lost images as failed so job can be marked complete #1244); revoke stale jobs to REVOKED rather than PENDING and extract the reaper (fix: revoke stale jobs to "revoked" status instead of "pending" #1169); actually schedule that reaper, which had existed but never ran (feat(jobs): schedule periodic stale-job check #1227); drain zombie NATS streams (feat(jobs): drain zombie NATS streams in periodic health check #1239).
- Transient-Redis resilience: retry on Redis blips instead of failing the job (fix: retry on Redis blips instead of failing the job #1231, closing
async_apijobs killed by transient Redis errors duringupdate_state—RedisErrorand "state actually missing" are conflated into a single fatal path #1219). - Observability: NATS lifecycle events in job logs (feat(jobs): make NATS queue activity visible in async ML job logs #1222, closing Surface NATS consumer/queue lifecycle events in the job logger so users can see what's happening from the UI #1220), who-fetched-which-task and throughput/ETA (feat(jobs): show who requests which tasks in ML processing job logs #1238), worker availability with a zero-workers warning (feat(jobs): show worker availability in async job logs #1240), DLQ visibility (Support for NATS dead letter queue #1175).
- Throughput / contention relief: move job logs to an append-only table (Refactor job logging to use separate table #1259, the permanent fix for feat(job): refactor job logging so it isn't a bottleneck #1256), throttle and defer the heartbeat off the request path (fix(jobs): throttle and defer pipeline heartbeat update #1258/fix(jobs): harden throttled pipeline heartbeat updates #1260), scale ML workers (fix(celery): update worker concurrency defaults #1228/Scale celeryworkers to keep up with async ML processing #1295), and speed up large-capture-set collection from ~19 min to ~42 s on 100k images (Improve the filter stage during job preparation, fix for large capture sets #1322).
What was attempted and then abandoned, reverted, or superseded
- Support cancelling ML async jobs #1144 (cancellation) — closed unmerged. First real async cancellation; abandoned (its cancel test was a stub) and reworked later. The concept traces to the still-open Pipeline requests should have a cancellation method #539.
- Update status of disconnected and stale jobs #981 (periodic status check) — closed unmerged. First attempt at "always-correct status" (Fixes Ensure that a processing job's status is always updated correctly #721). Its reaper logic was re-extracted into fix: revoke stale jobs to "revoked" status instead of "pending" #1169/feat(jobs): schedule periodic stale-job check #1227; its concurrent-write-safety half was pulled into the still-open draft [Draft] Don't overwrite logs & status in concurrent background tasks #1026.
- PSv2: Improve task fetching & web worker concurrency configuration #1142 partially superseded by fix: PSv2 follow-up fixes from integration tests #1135 during rebase (the batch-reservation API), leaving only the TTR/ack-pending/worker-count parts as net-new.
- PSv2: Use connection pooling and retries for NATS #1130 (NATS connection pooling) — still draft; fix: PSv2 follow-up fixes from integration tests #1135 took a fail-fast (
allow_reconnect=False) route instead. - RabbitMQ tuning (Increase the length a job can run in RabbitMQ #1060) is a dead branch after the NATS decision in Distributed Worker Architecture for ML Processing (Processing Service V2) pt. 1 #987.
- PSv2: Preparation for showing status of async processing services #1146/PSv2: Service "live" status for async processing services #1087 — landed but incomplete: service-status groundwork merged, but per-service online/offline identity was deferred and PSv2: Service "live" status for async processing services #1087 stays open.
A pattern worth naming: several fixes introduced the next bug. #1162 (cancellation + fail-on-missing-Redis) created the "transient Redis error kills the job" path that #1219/#1231 then had to repair, and whose residue is still open (#1232, #1241). #1234 dropped
max_deliverfrom 5 to 2, which sharpened the lost-image symptom that #1244 then had to reconcile. The slice has repeatedly fixed a symptom and surfaced the next one.What is still open
The never-closed roots (each has had multiple PRs chip at sub-cases without closing the general guarantee):
- Ensure that a processing job's status is always updated correctly #721 — "job status is not always correct" (reopened). The umbrella symptom: jobs that are done or errored still read STARTED, overall and per-stage. ~8 PRs fixed individual sub-cases; the general guarantee is still not closed.
- PSv2: Async jobs hang forever when NATS tasks exhaust max_deliver without posting results #1168 — async jobs hang forever when NATS exhausts
max_deliverwithout results. Five PRs mitigate it (Support for NATS dead letter queue #1175/feat(jobs): schedule periodic stale-job check #1227/feat(jobs): drain zombie NATS streams in periodic health check #1239/fix(jobs): mark failed or lost images as failed so job can be marked complete #1244/feat(jobs): show worker availability in async job logs #1240); the detection mechanism as specified and the root cause (Decouple ML result ingestion from the HTTP POST: large bodies, no durable handoff, no retry-record path #1223, the large-result-POST → 503 → redelivery cascade) are unbuilt. - Add COMPLETED job state for jobs that finish with partial errors #1134 — no COMPLETED / COMPLETED_WITH_ERRORS state. The structural simplification that would dissolve much of the SUCCESS/FAILURE/threshold churn the status path keeps hitting. Repeatedly worked around, never adopted.
The newly-named concurrency root:
- Saving job progress concurrently is the root of multiple issues related to incorrect job statuses #1337 — concurrent unlocked read-modify-write of
progress/status. fix(jobs): fixes for concurrent ML processing jobs #1261 droppedselect_for_updateto kill the row-lock contention that was serializing every result handler (feat(job): refactor job logging so it isn't a bottleneck #1256) and causing redelivery storms. That fixed contention but reintroduced a classic lost-update race: two handlers read the same snapshot, both write the whole JSONB blob back, one update is lost. This is the hinge — we traded a contention problem for a correctness problem — and it is the common cause behind the stuck-STARTED display bugs, the false REVOKEs of completed jobs (fix(jobs): fix dangling jobs from going to revoked #1276 mitigates), the SUCCESS→STARTED resurrection, and the scheduling crossover below.
The crossover that makes "just a display bug" unsafe: status is also a scheduling input, not only a UI output. A job wrongly left non-terminal stays claimable, so the worker
/nextendpoint keeps handing it out and starves newer jobs. #1282 (CANCELLED jobs leaking through/next) is exactly this with a different trigger. The/nextreshape in #1265 must agree with #1337 on which states are claimable, or the crossover reopens.Open PRs gating the default flip: #1276 (reaper reads Redis instead of clobbered progress — mergeable), #1279 (propagate pipeline config to pull-mode tasks), #1312 (captures marked processed with zero detections — currently conflicting), #1324 (dispatch/cancel hardening — see caution below), plus #1287 (savepoint around the log insert) and the not-yet-PR items #1265 and #1288.
Other open blockers and debt: #1263 (non-superuser worker tokens 403 on
/tasks//result, currently masked by making the worker a superuser), #1302 (Postgres connection-slot exhaustion under load — infra), #1283 (NATS stream/consumer not cleaned on REVOKE/restart → duplicate tasks), #1247 (one worker batch stranded after long jobs), the/tasksendpoint cluster (#1141/#1152/#1153/#1182), and the residual Redis-state handling (#1232/#1241).A caution on #1324: it is reasonable in intent but stacks two risky behaviors. Its
acks_late+reject_on_worker_lostredelivery-safety relies on readingjob.statusin an early guard — the very field #1337 says is unreliable under the clobber race. AndCELERY_WORKER_POOL_OPTIMIZATION = "fair"is not read by Celery from settings (fair scheduling is the-O fairCLI flag), so as written that half is a no-op while the cross-container head-of-line lever the originating incident hit is left commented out. Worth gating behind the #1337 fix (at least its conditional-status-write piece) rather than merging as standalone hardening. Also note its RabbitMQacks_lateinteraction withconsumer_timeoutdiffers by environment.Most critical remaining root issues, ranked
- Saving job progress concurrently is the root of multiple issues related to incorrect job statuses #1337 — atomic progress/status writes. The lost-update root behind Ensure that a processing job's status is always updated correctly #721/fix(jobs): fix dangling jobs from going to revoked #1276/CANCELLED jobs leak through /next filter, starve newer async_api jobs #1282/A few images are regularly stranded after long-running jobs #1247/PSv2: Async jobs hang forever when NATS tasks exhaust max_deliver without posting results #1168/Add GET /jobs/next action; deprecate list(ids_only=1) claim-semantics #1265. Until it lands, making async the default ships a known data-correctness race that gets worse with load, not better. Everything from fix(jobs): fix dangling jobs from going to revoked #1276 onward is mitigation around it.
- CANCELLED jobs leak through /next filter, starve newer async_api jobs #1282 — CANCELLED leaks through
/next, starves new jobs. A concrete, deterministic production-impacting instance of Saving job progress concurrently is the root of multiple issues related to incorrect job statuses #1337's crossover; small fix, but must align with Add GET /jobs/next action; deprecate list(ids_only=1) claim-semantics #1265. - Clean up task queue when job is revoked or re-started, fix duplicate tasks #1283 — no NATS cleanup on REVOKE/restart → duplicate tasks + counter resets. Independent lifecycle root cause; produces duplicate Detection/Classification rows.
- A few images are regularly stranded after long-running jobs #1247 — exactly-one-batch stranding after long jobs. Failure rate scales with job size; invisible at small scale, real at production scale. Partly lives in the worker.
- Fix captures are marked as processed with zero detections when they shouldn't be #1310 / Fix null detections in exports & API. Don't mark images as processed too soon #1312 — captures marked processed with zero detections. Real data-corruption root cause; fix exists but needs a rebase.
- Add COMPLETED job state for jobs that finish with partial errors #1134 — adopt a COMPLETED state. Would retire a whole class of SUCCESS/FAILURE/threshold edge cases the status path keeps relitigating.
- Add required permissions to processing service users for processing jobs #1263 — register worker permission codenames. Blocks the non-superuser / federated worker story (feat: API key auth and identity for processing services #1194/feat(ui): API key management and client info for processing services #1201).
A separate root-cause issue (#1337) now tracks item 1 in detail with proposed fix directions; #1285 tracks the diagnostic triangulation across the DB/Redis/NATS sources of truth that would make this whole class easier to debug.
- Architecture and contract: NATS pull model (Distributed Worker Architecture for ML Processing (Processing Service V2) pt. 1 #987), pipeline registration (PSV2: endpoint to register pipelines #1076),
Claude says: A proposed order for the remaining PSv2 root work, before async becomes the default path. The principle behind the ordering: fix the things that corrupt or lose state before the things that re-architect or optimize, and land each root fix on prod + demo + staging before relying on it. Sequence, not a hard dependency graph — several items can run in parallel where noted. Framed as a proposal to discuss, not a settled plan.
Stage 0 — stop the bleeding (correctness, small, parallelizable)
These are deterministic, low-risk, and each removes a way the system currently loses or corrupts state. They gate the flip.
- Saving job progress concurrently is the root of multiple issues related to incorrect job statuses #1337 Layer 1 — conditional terminal status transition. Stops status resurrection and the
/nextstarvation crossover. Small PR, ready. Everything else that readsjob.status(including Improve celery task dispatch and cancellation to prevent stuck jobs #1324's redelivery guard) is safer once this is in. - CANCELLED jobs leak through /next filter, starve newer async_api jobs #1282 — CANCELLED leaking through
/next. The concrete production instance of the crossover. Tiny fix (status__in+ add CANCELLED/terminal to the claim filter). Must be designed together with Add GET /jobs/next action; deprecate list(ids_only=1) claim-semantics #1265 so the claimable-state set is defined once. - Fix null detections in exports & API. Don't mark images as processed too soon #1312 — captures marked processed with zero detections. Real data corruption; the fix exists but needs a rebase and a human sign-off on the
valid()semantic change, plus the cleanup-command predicate patch.
Stage 1 — define the claim contract (the endpoint the workers pull from)
- Add GET /jobs/next action; deprecate list(ids_only=1) claim-semantics #1265 — dedicated
GET /jobs/next/, removing the three temporary hacks fix(jobs): fixes for concurrent ML processing jobs #1261 left inviews.py. This is where the "which states are claimable" decision lives; Saving job progress concurrently is the root of multiple issues related to incorrect job statuses #1337 Layer 1 and CANCELLED jobs leak through /next filter, starve newer async_api jobs #1282 must agree with it. Doing this after Stage 0 means the claim filter is written once, against the corrected state machine, rather than patched twice. - Add required permissions to processing service users for processing jobs #1263 — register worker permission codenames so non-superuser/federated workers can pull/post. Independent of the state work; can run in parallel with Stage 0/1. Blocks the API-key/identity work (feat: API key auth and identity for processing services #1194/feat(ui): API key management and client info for processing services #1201).
Stage 2 — the structural state fix (needs design first)
- Saving job progress concurrently is the root of multiple issues related to incorrect job statuses #1337 Layer 2 — per-task progress/history store. The counter-accumulation race. This needs a design pass first (see the plan comment on Saving job progress concurrently is the root of multiple issues related to incorrect job statuses #1337): there are other consumers of per-image job history (provenance PSv2: add Job relationship to a Classification and Detection objects #1156, debugging, re-processing), and a naive per-task table explodes on large capture sets. Decide the shape (counter columns + idempotency guard vs. a deliberately-sized per-task table vs. Redis-as-live + durable snapshot vs. event log) with those consumers in the room. Fold in idempotency on redelivery (fix(jobs): prevent jobs from hanging in STARTED state with no progress #1234) and the Redis-fragility class (Follow-ups from PR #1231 review: remaining Redis/retry gaps in async job state #1232/Make missing-Redis-state job failures loud and self-diagnosing #1241).
- Add COMPLETED job state for jobs that finish with partial errors #1134 — adopt a COMPLETED / COMPLETED_WITH_ERRORS state. Pairs naturally with the Layer 2 work, because it removes the SUCCESS/FAILURE-threshold logic that the progress path keeps relitigating. Worth deciding alongside the state-store redesign rather than bolting on after.
Stage 3 — durability and the hang/stranding class
- Decouple ML result ingestion from the HTTP POST: large bodies, no durable handoff, no retry-record path #1223 — decouple result ingestion from the HTTP POST. The large-result-body → 503 → redelivery-cascade is the root that generates the
max_deliverexhaustion. This is the architectural fix under PSv2: Async jobs hang forever when NATS tasks exhaust max_deliver without posting results #1168; it should come after the state store is settled because a durable result-handoff path interacts with how per-task state is recorded. - PSv2: Async jobs hang forever when NATS tasks exhaust max_deliver without posting results #1168 / A few images are regularly stranded after long-running jobs #1247 — async hang and last-batch stranding. Once Decouple ML result ingestion from the HTTP POST: large bodies, no durable handoff, no retry-record path #1223 and Layer 2 are in, revisit whether the remaining mitigations (feat(jobs): schedule periodic stale-job check #1227/feat(jobs): drain zombie NATS streams in periodic health check #1239/fix(jobs): mark failed or lost images as failed so job can be marked complete #1244) are still needed or can be simplified; A few images are regularly stranded after long-running jobs #1247 also has an ami-data-companion-side component.
- Clean up task queue when job is revoked or re-started, fix duplicate tasks #1283 — NATS stream/consumer cleanup on REVOKE/restart. Independent lifecycle root cause (duplicate tasks, counter resets); can move earlier if it bites, but it's lower-frequency than the state-corruption items.
Cross-cutting (run alongside, not blocking)
- Review how the state and progress of jobs are tracked #1285 — DB/Redis/NATS triangulation diagnostics. Build incrementally; it makes every stage above easier to verify and de-risks the flip itself.
- perf(jobs): defer aggregate count refresh in process_nats_pipeline_result (5-10x speedup expected) #1288 — defer the aggregate-count cascade. Latency, but it amplifies the contention that feeds Saving job progress concurrently is the root of multiple issues related to incorrect job statuses #1337, so worth doing around Stage 2.
- Behavioural gates before the flip: concurrent-load async e2e green on demo and staging (
test_ml_job_e2e+ chaos), zombie-stream sweep clean under that load, and the broker/Redis config confirmed per environment. Idle-green is not load-green.
A note on #1324
Reasonable in intent but should be gated behind #1337 Layer 1 (its
acks_lateredelivery guard readsjob.status), and its fair-scheduler setting is a no-op as written (fair scheduling is the-O fairCLI flag, not a settings value). TheJob.cancel()rewrite in it is sound and could be split out and merged early. Itsacks_lateinteraction with RabbitMQconsumer_timeoutalso differs by environment and needs the timeout raised where it currently uses the 30-minute default.- Saving job progress concurrently is the root of multiple issues related to incorrect job statuses #1337 Layer 1 — conditional terminal status transition. Stops status resurrection and the
- addedPipeline APIUpdates to the requests & responses to/from processing service workers for ML pipelinesUpdates to the requests & responses to/from processing service workers for ML pipelines
on Jun 24, 2026


@kaviecos and @mihow have designed & written the specifications for a new ML backend that orchestrates multiple types of models by different research teams, across multiple stages of processing, and is horizontally scalable. This expands on the current ML backend API defined here https://ml.dev.insectai.org/ by adding asynchronous processing, a controller & queue system, auth and many other production features.
The initial spec and notes are here, but are being re-written in the Aarhus GitLab wiki as the backend is developed.
https://docs.google.com/document/d/1caKxxfZhWhRi9Jfv9fy5fVeoM9bvhYPJ/
Docs in progress:
https://gitlab.au.dk/ecos/hoyelab/ami/ami-pipeline-controller/-/wikis/Getting-Started
https://gitlab.au.dk/ecos/hoyelab/ami/ami-pipeline-controller/-/wikis/Pipeline-Stages
https://gitlab.au.dk/ecos/hoyelab/ami/ami-pipeline-controller/-/wikis/Architecture-Overview
Known remaining tasks: