Conversation
GET /api/v1/tasks?workflow_id=<id> materialised every TaskInfo in the runtime and filtered in Python on the event loop, pinning a CPU and freezing the server under a large task count. Apply the workflow_id filter inside TaskRuntime.list_tasks before building TaskInfo objects, and run the remaining store work off the loop via asyncio.to_thread. Co-Authored-By: Claude Code <noreply@anthropic.com> Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com> Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
kaiitunnz
left a comment
There was a problem hiding this comment.
While reviewing this I traced where the listing time actually goes. The workflow_id filter-before-build is what fixes the measured case. Several other paths share the same root cause, so folding them in here seems worthwhile.
Where the time goes
Local microbenchmark: in-memory TaskRuntime (no Redis), 2000 echo tasks, no query params.
| Tasks per workflow | list_tasks build |
filter_models_by_queries |
Response serialization |
|---|---|---|---|
| 5 | 0.024 s | 3.8 s | 0.04 s |
| 50 | 0.029 s | 29.6 s | 0.09 s |
| 200 | 0.053 s | 117.8 s | 0.24 s |
filter_models_by_queries calls model_dump() on every model (src/server/utils/misc.py:103-104). That runs TaskRecord's redacting serializer, which calls redact_raw_yaml (pure-Python yaml.safe_load + safe_dump) on source. source holds the full workflow YAML, copied onto every task (src/server/task/runtime.py:158). The result is cached on the instance (_redacted_source), but list_tasks builds fresh TaskInfos on every call. A counter confirmed exactly one redaction per task per listing, so the cost grows with tasks × workflow size. Response serialization stays cheap only because it reuses that cache.
So asyncio.to_thread covers about 1% of the cost. The filter and serialization still run on the event loop, and an unfiltered GET /api/v1/tasks (or ?status=...) still freezes the server.
Same root cause in submission
register calls save_task_states_async (runtime.py:231), which runs PersistedTask.model_dump_json() and hits the same redaction once per task. Persisting one 200-task workflow took 11.5 s on the loop locally. The cache isn't persisted, so after rehydrate each record redacts again on its first commit.
Suggested changes
- Redact once per workflow, at registration. Nothing in
src/readsTaskRecord.sourceexcept redaction, and the persisted copy is already redacted. Storing the redacted text on every record (one string shared across the workflow) therefore looks safe, and the serializer's per-task source redaction can go. This covers listing, submission and restart. - Query filter that reads attributes from a declared allowlist. It would read attributes instead of calling
model_dump(), skip entirely when there are no query keys, and accept only a declared set of filterable fields per endpoint. The allowlist matters: matching raw attributes instead of the redacted dump would let?task.spec.<…>.api_key=<guess>serve as an oracle for credential values. Indexed fields (workflow_id,status) can be pushed down intoTaskRuntime.list_tasks, the way this PR does forworkflow_id. - Repeated
workflow_id.request.query_params.get("workflow_id")returns only the last value (QueryParams('workflow_id=A&workflow_id=B').get('workflow_id') == 'B'). So?workflow_id=A&workflow_id=Bnow returns only B, wheremainreturned both underfilter_models_by_queries's OR semantics. Usinggetlistplus anincheck in the runtime keeps the old behavior. - Pagination.
docs/API.md:114documentslimit/before/afterfor/api/v1/workflowsand/api/v1/tasks, but neither handler reads them, and the filter ignores unknown keys, solimitis silently dropped. Implementing the documented cursor contract would bound response size. GET /api/v1/workflowsN+1. It makes 5 sequential Redis calls per workflow (workflows.py:529-534→WorkflowRegistry.get_workflow_async), across every workflow inWORKFLOWS_SET_KEY. One pipeline for the batch would make this a single round trip.- Blocking I/O on the loop.
to_threadis effective here, since file and network I/O release the GIL:download_result_bundlebuilds the tar.gz synchronously (_create_result_bundle_archive), and its size grows with checkpoints and artifacts.upload_result_fileandupload_task_traceread the whole upload withawait file.read(), then write synchronously. A chunked copy in a thread would also bound memory.analyze_workflow_tracereads all the JSONL and runsanalyzeon the loop (traces.py:71-75).WorkflowRegistry.commit_transitionuses the sync Redis pipeline and is called from the event monitor on the loop, so every task transition is a blocking round trip.
Tests
test_event_loop_stays_responsive_during_listswapslist_tasksfor a stub that sleeps, so it would still pass with the filter cost on the loop. A version whose slow step is the filtering and serialization would cover the real failure.monkeypatch.setattrwould also remove the# type: ignore[method-assign].- A test that redaction runs once per workflow rather than once per task would pin fix 1.
The A/B only covered ?workflow_id=. Rerunning it with an unfiltered GET /api/v1/tasks and a large-workflow submission would confirm the end-to-end effect.
Refs mlsys-io#183 Signed-off-by: kaiitunnz <79885172+kaiitunnz@users.noreply.github.com>
Filter matching scanned a key's accepted values per item, and the task listing matches under the runtime lock, so 3000 repeated task_id values held the lock for seconds at 20000 tasks. Each key's values are now normalized once per query, so a match is a set lookup. Each list route declares the filters its own response model serves: identity, ownership, state, small descriptors and id lists, never a credential-bearing, private, free-form or continuous field. /stack/workers gets its own set, since its model has no node_id. A declared path through an unset value reads as null, so a worker with no hardware report matches hardware.gpu.cuda_version=null rather than every version. Only a route's own paging parameters are skipped, so /workers?limit=5 is a 400, and every route parses its filter before doing any work. The task page serializes straight to bytes. Refs mlsys-io#183 Signed-off-by: kaiitunnz <79885172+kaiitunnz@users.noreply.github.com>
A workflow page read every workflow's submission time to order them, so its cost grew with every workflow ever submitted: about 0.5 s per page and a 210 ms loop stall at 50000 workflows, and an O(N^2) SDK walk. Workflows now also join a sorted set keyed by submission time and id, written in the registration transaction, so an unfiltered page reads its range in O(log N + limit). The index is derived: the root indexes any registered workflow it lacks before it serves, so workflows from before the index list too. A caller restricted to some workflows, or a workflow_id filter, orders just those candidates, and a rejecting filter scans in chunks that double up to 1000. The sync get_workflow, which nothing called, is gone. Refs mlsys-io#183 Signed-off-by: kaiitunnz <79885172+kaiitunnz@users.noreply.github.com>
Purpose
GET /api/v1/tasks?workflow_id=<id>froze the whole server. On a node server holding 13025 tasks, each call pinned one CPU for about 14 minutes. During that time/healthzand every other request timed out (the listen backlog grew to 618), the scheduler dispatched nothing, and Lumilake jobs stalled. A client that timed out and retried queued another full scan, so the freeze compounded.Changes
src/server/task/runtime.py).TaskRuntime.list_taskstakes an optionalworkflow_idand builds aTaskInfoonly for the tasks of that workflow, so building, sanitising, filtering and serialising scale with the matching tasks, not with every task the server holds.src/server/routers/v1/tasks.py). The router runslist_tasksthroughasyncio.to_thread, so a large listing no longer blocks other requests. The permission filter (resolve_accessible_ids), the other query filters and the response shape are unchanged.Test Plan
Pre-commit (isort, black, ruff, codespell, mypy) on the changed files, then
tests/server. The new tests cover:TaskInfofor other workflows' tasks;Test Result
tests/server: 1004 passed atffd55ef(containsmain1b0a320).GET /api/v1/tasks?workflow_id=for a dedicated 2-task workflow, with/healthzprobed every second:main/healthzduring the call