You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Blocking and streaming flex requests (service_tier: "flex" on /v1/chat/completions, or non-background /v1/responses) are written into fusillade and polled by the handler. Nothing cancelled the row when the client dropped the connection. The daemon claimed and dispatched the abandoned request anyway, and because upstream timeouts are retriable it could re-dispatch it repeatedly. A client with its own timeout that retries therefore stacks redundant work onto an already backlogged engine.
The daemon could not have honoured a cancel even if one had been written: only batched requests have a cancellation path (batch cancelling_at → shared token). Batchless requests got a fresh token nobody fired (daemon/mod.rs, None => CancellationToken::new()).
Change
dwctl
flex_stream_response selects on the SSE channel closing (the receiver lives in the response body, dropped on disconnect) and cancels the row.
The blocking handlers (handle_flex foreground, handle_chat_completion_flex) hold an AbandonGuard across the poll. hyper drops the handler future on client disconnect, so the guard's Drop issues the cancel. It is disarmed once a terminal state has been rendered.
Both also cancel when the poll itself gives up (1h timeout / storage error), since the client has by then been answered with an error.
fusillade daemon
Per-request cancellation tokens for in-flight batchless requests, removed when the processing task exits.
The existing cancellation poll additionally asks storage which of those rows are now canceled and fires their tokens, which aborts the HTTP task through the existing complete() select. New gauge fusillade_cancellation_poll_requests_checked.
storage
Storage::cancel_batchless_request(id) -> bool: a single guarded UPDATE mirroring the persist(Canceled) transition (only completed/failed resist). Needed because the default cancel_requests goes through get_requests, which inner-joins batches and never sees batchless rows.
Storage::get_cancelled_request_ids(ids): batchless sibling of get_cancelled_batch_ids.
canceled remains the soft terminal: a completion that lands anyway supersedes it and is billed, exactly as for batches today.
Behaviour change
A foreground flex request whose caller disconnects is now cancelled rather than left to complete. Its result can no longer be fetched later by ID. Background (background: true) requests are unaffected: they return 202 immediately and are never polled by the handler.
Tests
fusillade-arsenal: get_cancelled_request_ids_reports_only_canceled_rows covers both new storage methods, including that a completed row resists the cancel.
fusillade integration: cancelling_in_flight_batchless_request_aborts_http runs the daemon against a mock upstream that never answers, cancels the row, and asserts the in-flight call is aborted, not retried, and the row stays canceled.
Aborting the daemon's HTTP task closes the loopback hop, which drops the upstream call from onwards. Whether the inference engine actually stops generating on a dropped connection is the same open question as for realtime disconnects and is outside this PR.
Blocking and streaming flex requests are written to fusillade and then
polled; nothing cancelled the row when the client went away. The daemon
claimed and dispatched the abandoned request anyway, and if the upstream
call was slow it kept retrying it, so a client that times out and retries
piles redundant work onto an already backlogged engine.
Two halves:
- dwctl cancels the row when the caller is gone. The streaming path
selects on the SSE channel closing; the blocking paths hold a drop
guard across the poll, since hyper drops the handler future on client
disconnect. Both also cancel when the poll itself gives up, because the
client has by then been answered with an error.
- The daemon honours row-level cancels on batchless requests. Batched
work is cancelled via the batch's `cancelling_at` and a shared token;
batchless work previously got a fresh token nobody could fire. In-flight
batchless requests now have per-request tokens, and the existing
cancellation poll also asks storage which of those rows are `canceled`,
aborting the upstream call through the existing `complete()` select.
Storage gains `cancel_batchless_request` (the default `cancel_requests`
rebuilds rows via `get_requests`, which inner-joins `batches` and so
never sees batchless rows) and `get_cancelled_request_ids`. `canceled`
stays the soft terminal: a completion that lands anyway supersedes it,
as before.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
The reason will be displayed to describe this comment to others. Learn more.
1 issue found across 7 files
Prompt for AI agents (unresolved issues)
Check if these issues are valid — if so, understand the root cause of each and fix them. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="fusillade-arsenal/src/postgres.rs">
<violation number="1" location="fusillade-arsenal/src/postgres.rs:5197">
P1: A retriable upstream result can overwrite this cancellation and re-enter the retry path, because batchless rows have no batch `cancelling_at` fence. Preserve a durable cancellation fence through retry decisions or suppress retries for rows canceled by `cancel_batchless_request`; otherwise a disconnected flex request can be dispatched again.</violation>
</file>
Reply with feedback, questions, or to request a fix.
The reason will be displayed to describe this comment to others. Learn more.
P1: A retriable upstream result can overwrite this cancellation and re-enter the retry path, because batchless rows have no batch cancelling_at fence. Preserve a durable cancellation fence through retry decisions or suppress retries for rows canceled by cancel_batchless_request; otherwise a disconnected flex request can be dispatched again.
Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fusillade-arsenal/src/postgres.rs, line 5197:
<comment>A retriable upstream result can overwrite this cancellation and re-enter the retry path, because batchless rows have no batch `cancelling_at` fence. Preserve a durable cancellation fence through retry decisions or suppress retries for rows canceled by `cancel_batchless_request`; otherwise a disconnected flex request can be dispatched again.</comment>
<file context>
@@ -5156,6 +5156,56 @@ impl<P: PoolProvider> Storage for PostgresRequestManager<P> {
+ canceled_at = NOW()
+ WHERE id = $1
+ AND batch_id IS NULL
+ AND state NOT IN ('completed', 'failed', 'canceled')
+ RETURNING id
+ "#,
</file context>
The reason will be displayed to describe this comment to others. Learn more.
Not reachable. A retriable failure on a cancelled row goes through reschedule_for_retry, whose UPDATE is fenced on state IN ('claimed', 'processing') AND daemon_id = $2 AND claimed_at = $5 (postgres.rs reschedule_for_retry). A row moved to canceled by cancel_batchless_request fails that fence, the daemon logs the retry as not persisted, and the row stays canceled. When retries are exhausted, persist(Failed) does supersede canceled — that is the existing soft-terminal rule for batches as well, and it ends the work rather than re-dispatching it. The integration test cancelling_in_flight_batchless_request_aborts_http asserts the mock upstream is called exactly once after the cancel.
Review follow-ups:
- The abandon guard was armed after `create_flex`, so a disconnect while
the INSERT was in flight dropped the handler with no guard, and a row
that committed anyway was dispatched. Arm it before the enqueue on the
foreground paths (blocking, streaming); background submissions never arm
it. The streaming poll task now owns the guard instead of calling the
cancel directly.
- `get_cancelled_request_ids` reads from the primary: it is the daemon's
detector for a marker the API layer has just written, and replica lag
is engine time spent on an abandoned request. PK probe over the
in-flight set, so the cost is negligible.
- Publish the cancellation-poll gauges before the empty-list early exit so
an idle daemon reports 0 rather than the last non-empty poll.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Record detached cancellation failures in background error metrics
dwctl/src/inference/store.rs:325
This failure is handled by a detached task after the request future is dropped, so the cancellation error is not reflected in the HTTP metrics and can leave abandoned work running. The repository's background-task convention is to use background_error! for swallowed off-request-path failures (see dwctl/src/metrics/errors.rs:3-9); please record this with an appropriate component/reason label so failed cancellations are observable and alertable.
A failed cancel of an abandoned flex request runs off the request path
(the caller is gone; from the drop guard it is a detached task), so a
plain warn! was invisible to metrics. Route it through background_error!
under a new `flex_cancel` component so it shows up in
dwctl_background_errors_total. Also reword the not-cancelled debug log:
with the guard armed before the enqueue, Ok(false) now also covers a row
that never committed, not only one that is already terminal.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Re Copilot's "previously missed" note on cancel_abandoned_request (store.rs, not attached to a thread): done in 15c2458. The Err arm now goes through background_error! with a new flex_cancel component and reason cancel_abandoned, severity Error, so failed cancels land in dwctl_background_errors_total.
`background: true` returns 202 and is collected later by id, so the
caller going away must not cancel it. Only the foreground flex paths arm
the disconnect cancel; this end-to-end test submits a background flex
Responses request, drops the client response, and asserts the row is not
`canceled`.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
The reason will be displayed to describe this comment to others. Learn more.
1 issue found across 1 file (changes from recent commits).
Prompt for AI agents (unresolved issues)
Check if these issues are valid — if so, understand the root cause of each and fix them. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="dwctl/src/test/responses.rs">
<violation number="1" location="dwctl/src/test/responses.rs:1405">
P2: `drop(response)` runs after the complete 202 exchange, and the immediate mock lets the daemon finish before the detached cancellation can be observed. Hold the upstream request in flight and assert the row remains non-canceled before allowing completion; otherwise a background-cancellation regression can pass this test.</violation>
</file>
Tip: Review your code locally with the cubic CLI to iterate faster.
The reason will be displayed to describe this comment to others. Learn more.
P2: drop(response) runs after the complete 202 exchange, and the immediate mock lets the daemon finish before the detached cancellation can be observed. Hold the upstream request in flight and assert the row remains non-canceled before allowing completion; otherwise a background-cancellation regression can pass this test.
Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At dwctl/src/test/responses.rs, line 1405:
<comment>`drop(response)` runs after the complete 202 exchange, and the immediate mock lets the daemon finish before the detached cancellation can be observed. Hold the upstream request in flight and assert the row remains non-canceled before allowing completion; otherwise a background-cancellation regression can pass this test.</comment>
<file context>
@@ -1373,3 +1373,48 @@ async fn test_realtime_zdr_suppresses_stored_bodies(pool: PgPool) {
+ .expect("202 body carries a resp_<uuid> id");
+
+ // The client has its id and moves on: the HTTP exchange is over.
+ drop(response);
+ tokio::time::sleep(std::time::Duration::from_millis(500)).await;
+
</file context>
The reason will be displayed to describe this comment to others. Learn more.
The premise doesn't hold here: the test harness sets batch_daemon.enabled = DaemonEnabled::Never, so no daemon ever claims the row and the mock is never reached. The row can only leave pending if something cancels it. Tightened in d90fc22 to assert_eq!(state, "pending") so a spurious cancel fails the test outright, and added a comment saying why.
The test harness never runs the daemon, so the row can only leave
`pending` if something cancels it. Asserting equality makes a spurious
cancel fail the test outright rather than depending on timing.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This branch has not been deployed
No deployments
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Problem
Blocking and streaming flex requests (
service_tier: "flex"on/v1/chat/completions, or non-background/v1/responses) are written into fusillade and polled by the handler. Nothing cancelled the row when the client dropped the connection. The daemon claimed and dispatched the abandoned request anyway, and because upstream timeouts are retriable it could re-dispatch it repeatedly. A client with its own timeout that retries therefore stacks redundant work onto an already backlogged engine.The daemon could not have honoured a cancel even if one had been written: only batched requests have a cancellation path (batch
cancelling_at→ shared token). Batchless requests got a fresh token nobody fired (daemon/mod.rs,None => CancellationToken::new()).Change
dwctl
flex_stream_responseselects on the SSE channel closing (the receiver lives in the response body, dropped on disconnect) and cancels the row.handle_flexforeground,handle_chat_completion_flex) hold anAbandonGuardacross the poll. hyper drops the handler future on client disconnect, so the guard'sDropissues the cancel. It is disarmed once a terminal state has been rendered.fusillade daemon
canceledand fires their tokens, which aborts the HTTP task through the existingcomplete()select. New gaugefusillade_cancellation_poll_requests_checked.storage
Storage::cancel_batchless_request(id) -> bool: a single guarded UPDATE mirroring thepersist(Canceled)transition (onlycompleted/failedresist). Needed because the defaultcancel_requestsgoes throughget_requests, which inner-joinsbatchesand never sees batchless rows.Storage::get_cancelled_request_ids(ids): batchless sibling ofget_cancelled_batch_ids.canceledremains the soft terminal: a completion that lands anyway supersedes it and is billed, exactly as for batches today.Behaviour change
A foreground flex request whose caller disconnects is now cancelled rather than left to complete. Its result can no longer be fetched later by ID. Background (
background: true) requests are unaffected: they return 202 immediately and are never polled by the handler.Tests
fusillade-arsenal:get_cancelled_request_ids_reports_only_canceled_rowscovers both new storage methods, including that a completed row resists the cancel.fusilladeintegration:cancelling_in_flight_batchless_request_aborts_httpruns the daemon against a mock upstream that never answers, cancels the row, and asserts the in-flight call is aborted, not retried, and the row stayscanceled.dwctl:dropping_the_stream_cancels_the_flex_request,abandon_guard_cancels_when_handler_is_dropped,abandon_guard_leaves_delivered_result_alone.Not covered here
Aborting the daemon's HTTP task closes the loopback hop, which drops the upstream call from onwards. Whether the inference engine actually stops generating on a dropped connection is the same open question as for realtime disconnects and is outside this PR.
🤖 Generated with Claude Code