fix(simple-engine): close delegated streams - #687
Conversation
janhilgard
left a comment
There was a problem hiding this comment.
Approving. I reproduced #682 on main, confirmed the fix, and checked it against #679 since both rewrite the same streaming routes.
Reproduction on b998776, matching the issue exactly:
backend_closed : False
active_requests : 2
num_running : 1
At your head:
backend_closed : True
active_requests : 0
num_running : 0
aclosing needs 3.10 and requires-python is >=3.10, so the floor is right — worth saying out loud because it is the kind of stdlib addition that passes locally and breaks the oldest matrix entry.
The ownership split reads correctly to me: stream_generate closes the tracked stream it created, _track_request_stream closes the source it consumes, and each nested delegation closes its own. Double-closing is harmless since aclose() is idempotent, so the overlap between the re-entrant branch and the _stream_generate_text fallback is not a problem.
Cross-check against #679
These two conflict, but trivially — one import line, aclosing versus ThreadPoolExecutor, independent additions on both sides. Keeping both resolves it. I merged them locally and ran everything:
tests/test_simple_engine_thread_pinning.py: 9 passed- full suite merged: 2306 passed (the four failures I see are reproduced on pristine
main— three Python 3.14, one missingffmpeg)
More interesting than the absence of breakage: the two fixes compose into something neither achieves alone, and I measured it rather than assuming. #679 moves the backend generator's close() onto the pinned generation thread, because MLX buffers can only be evaluated where their stream was built. But that finally only runs when the generator actually closes — which, before this PR, it did not on a client disconnect. It ran at garbage-collection time instead, on whatever thread the collector happened to be on, which is precisely the thread-ownership violation #679 exists to prevent.
With both applied, a mid-stream aclose() gives:
event loop = 8314968256
generation worker = 6119976960
generation ran on = {6119976960}
close ran on = {6119976960}
active_requests = 0, num_running = 0
warnings = none
So cleanup lands deterministically on the owning thread at disconnect. No GeneratorExit trouble from awaiting inside the finally either, which was my main worry going in.
One gap
The specprefill delegation at _stream_generate_impl is wrapped, but none of the three new tests exercise it — they cover direct generation, the nested chat fallback and the MLLM text route. That path is the one I would most want covered, because it is the only delegation with its own cancellation machinery (_SpecPrefillCancelled, abort_event, on_cancel through _run_blocking_serialized), so a GeneratorExit arriving there interacts with a second mechanism rather than just unwinding. The two pre-existing specprefill cancellation tests are also the ones that already behave differently on 3.14, so that corner is not the one to leave to inference.
Not a blocker — the wrapping is the same shape as the three that are tested, and leaving it unwrapped would be worse.
Production note
Both our services run SimpleEngine routes behind opencode, which drops streams on user cancel constantly, so this is the everyday path rather than an edge case. We had been reading the generator didn't stop after athrow() lines as harmless noise; on the evidence above they were marking real leaked backend iterators and stale _active_requests entries, which also means /v1/status was over-reporting in-flight requests. Good catch.
dcffb7e to
7b14e13
Compare
|
Rebased the single stream-closure commit onto current |
…waybarrios#687 aclosing kept Rebase onto upstream `22efb47` (5 commits past `4b654c0`). Six conflict stops, all in three files. This commit carries the non-mechanical parts. - Upstream waybarrios#677 adds `not _thinking_disabled(...)` to the chat-completions streaming gate — exactly the skip patch #27 exists to prevent (gemma-4 emits `<|channel>thought` with thinking off; skipping the parser leaks raw markers into content). Clause dropped in both places; upstream's `or output_finished` and `_extract_streaming_reasoning_delta` helper kept, since waybarrios#677 now withholds a trailing partial tag and flushes it via `finalize_stream()`. - The Responses-path marker latch is rewired rather than re-applied: upstream's `use_reasoning` is computed once per output as "parser and thinking on", which would silently defeat a latch that can engage mid-stream. It is now derived from the latch and re-armed when the latch trips. - Upstream waybarrios#687 (`aclosing` on delegated streams) is KEPT as a wanted delta — it is not superseded by our patches, and it closes the inner generator when a client disconnects. Fork's `strip_effort_fallback` (waybarrios#76) re-applied inside it. - `test_public_stream_chat_close_cleans_nested_generate` given a fork-aware fixture instead of a skip: it stubs `engine._model.stream_generate`, but the fork's text route calls `mlx_lm.stream_generate` directly to pass a `prompt_cache`, so the real entry point ran and died on a MagicMock prompt. With both stubbed it passes and verifies fork abort cleanup. - Removed three upstream imports confirmed dead in the fork's `simple.py`: `OrderedDict` (upstream's own system-KV LRU; we delegate to SystemKVManager) and `ThreadPoolExecutor` (upstream's generation-executor machinery). Suite: 3315 passed / 31 skipped / 30 deselected; ruff clean. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01QFxkYqZwQgVp5QDNfpU6Hp
…waybarrios#687 aclosing kept Rebase onto upstream `22efb47` (5 commits past `4b654c0`). Six conflict stops, all in three files. This commit carries the non-mechanical parts. - Upstream waybarrios#677 adds `not _thinking_disabled(...)` to the chat-completions streaming gate — exactly the skip patch #27 exists to prevent (gemma-4 emits `<|channel>thought` with thinking off; skipping the parser leaks raw markers into content). Clause dropped in both places; upstream's `or output_finished` and `_extract_streaming_reasoning_delta` helper kept, since waybarrios#677 now withholds a trailing partial tag and flushes it via `finalize_stream()`. - The Responses-path marker latch is rewired rather than re-applied: upstream's `use_reasoning` is computed once per output as "parser and thinking on", which would silently defeat a latch that can engage mid-stream. It is now derived from the latch and re-armed when the latch trips. - Upstream waybarrios#687 (`aclosing` on delegated streams) is KEPT as a wanted delta — it is not superseded by our patches, and it closes the inner generator when a client disconnects. Fork's `strip_effort_fallback` (waybarrios#76) re-applied inside it. - `test_public_stream_chat_close_cleans_nested_generate` given a fork-aware fixture instead of a skip: it stubs `engine._model.stream_generate`, but the fork's text route calls `mlx_lm.stream_generate` directly to pass a `prompt_cache`, so the real entry point ran and died on a MagicMock prompt. With both stubbed it passes and verifies fork abort cleanup. - Removed three upstream imports confirmed dead in the fork's `simple.py`: `OrderedDict` (upstream's own system-KV LRU; we delegate to SystemKVManager) and `ThreadPoolExecutor` (upstream's generation-executor machinery). Suite: 3315 passed / 31 skipped / 30 deselected; ruff clean. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01QFxkYqZwQgVp5QDNfpU6Hp
…waybarrios#687 aclosing kept Rebase onto upstream `22efb47` (5 commits past `4b654c0`). Six conflict stops, all in three files. This commit carries the non-mechanical parts. - Upstream waybarrios#677 adds `not _thinking_disabled(...)` to the chat-completions streaming gate — exactly the skip patch #27 exists to prevent (gemma-4 emits `<|channel>thought` with thinking off; skipping the parser leaks raw markers into content). Clause dropped in both places; upstream's `or output_finished` and `_extract_streaming_reasoning_delta` helper kept, since waybarrios#677 now withholds a trailing partial tag and flushes it via `finalize_stream()`. - The Responses-path marker latch is rewired rather than re-applied: upstream's `use_reasoning` is computed once per output as "parser and thinking on", which would silently defeat a latch that can engage mid-stream. It is now derived from the latch and re-armed when the latch trips. - Upstream waybarrios#687 (`aclosing` on delegated streams) is KEPT as a wanted delta — it is not superseded by our patches, and it closes the inner generator when a client disconnects. Fork's `strip_effort_fallback` (waybarrios#76) re-applied inside it. - `test_public_stream_chat_close_cleans_nested_generate` given a fork-aware fixture instead of a skip: it stubs `engine._model.stream_generate`, but the fork's text route calls `mlx_lm.stream_generate` directly to pass a `prompt_cache`, so the real entry point ran and died on a MagicMock prompt. With both stubbed it passes and verifies fork abort cleanup. - Removed three upstream imports confirmed dead in the fork's `simple.py`: `OrderedDict` (upstream's own system-KV LRU; we delegate to SystemKVManager) and `ThreadPoolExecutor` (upstream's generation-executor machinery). Suite: 3315 passed / 31 skipped / 30 deselected; ruff clean. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01QFxkYqZwQgVp5QDNfpU6Hp
Closing a public SimpleEngine stream did not close the async generators it delegates to. Backend iterators and request state could remain live until finalization, which also produced
generator didn't stop after athrow()errors.This uses
aclosingat each public and nested streaming delegation boundary, including the MLLM text route and SpecPrefill. Normal exhaustion is unchanged.Tests cover direct generation, nested chat fallback, and the MLLM text route. They verify backend cleanup, request counters, generation-lock release, and tracker context reset.
Verification: 2,300 passed, 24 skipped.
Refs #682