Pull requests / #436
#436 serve: a non-streaming request stops when its client disconnects (#430, part of #431)
closed · @homeofe · 0 Kommentare · Auf GitHub
BenchmarksSetup & installServer & APIAMD / HIPNVIDIA / CUDAModels & quantsWindows
Beschreibung
Fixes #430. Part of #431: a stream that hangs up during a long prompt read is now found one prompt chunk sooner. See *What remains for #431*.
## Problem
A non-streaming answer writes nothing until it's finished. `openai_collect` / `anthropic_collect` read the whole generator, and `cancel` is only set in the streaming paths' `except OSError`, when a write fails. So a client that gave up (an agent's request timeout, a closed tab, a killed script) went unnoticed: the engine read the whole prompt and generated up to `max_tokens`, while the next request waited in the FIFO behind it. The same applies to `/v1/messages`, not only `/v1/chat/completions`.
## Change (`serve/server.py`)
- **`Handler._client_gone()`** runs `select([conn], [], [], 0)` and then `conn.recv(1, socket.MSG_PEEK)`.
- Only `b''` (EOF) or a `ConnectionError` (reset) means the client is gone.
- Bytes that are there, such as a pipelined next request, are left where they are and mean "alive". MSG_PEEK never consumes them, and the `select` makes sure the peek never blocks.
- Any other `OSError`/`ValueError` (e.g. an fd `select` can't take) means "alive": a later write still finds out, as before.
- **`Handler._watched(items, cancel)`** wraps the request's chunks/events. It checks the socket on every heartbeat (`None`: one per prompt chunk, or every 10 s of quiet) and at most once a second between tokens. When the client is gone:
- It sets `cancel` and raises the new `ClientGone(ConnectionAbortedError)`.
- Its `finally` closes the inner generator. That's the chain the streaming path already uses: `Service.run` records finish `disconnect` and closes the run, and `StrataEngine.generate` sends `STOP` and drains to `DONE`.
- `ClientGone` is an `OSError`, so the existing disconnect handling in the stream loops and in `do_POST` takes it unchanged. It's a `ConnectionError`, so `Server.handle_error` prints no traceback for it.
- **Where it's applied:** both APIs, streaming and not: `chunks = self._watched(self._capture(chunks, "openai"), cancel)`, and the same for the Anthropic `events`.
There are no engine changes. The engine already stops at the next prompt chunk once it gets `STOP`: `generate.cpp` reads stdin on its own thread, and `prefill.cpp:1049` checks `should_stop` before every chunk.
## Tests (`serve/test_server.py`)
`BeatingEngine(MockEngine)` simulates a long prompt read: `beats` heartbeats `beat_s` apart, then the scripted tokens. It records heartbeats, tokens and when each request was closed. The new `ClientDisconnect` class runs a real HTTP server with raw-socket clients, and runs every case on both `/v1/chat/completions` and `/v1/messages`.
| test | before | after |
|---|---|---|
| non-streaming, the client hangs up while tokens are generated (2,000 tokens at 20 ms, client closes at 0.5 s): the engine stops within 5 s with < 200 tokens, the next request is answered within 5 s, history says `disconnect` | **FAIL**: "the engine still runs for a client that is gone" (Anthropic: **ERROR**, the abandoned request still holds the FIFO) | pass (about 50 tokens) |
| non-streaming, the client hangs up while the prompt is read (40 heartbeats at 0.3 s, closes at 0.45 s): 0 tokens, ≤ 3 heartbeats | **FAIL** / **ERROR** as above | pass (2 heartbeats) |
| streaming, the client hangs up during the prompt read (#431): stopped after the first heartbeat | **FAIL**: `(2, 0) != (1, 0)`, two heartbeats | pass |
| a client that waits gets its whole answer (slow 300-token non-streaming answer with several checks, with and without bytes of a next request sent on the same socket mid-answer): HTTP 200, 300 tokens, finish `length` | pass | pass |
```
.venv/bin/python -B -m unittest serve.test_server # 89 tests OK (85 + 4 new)
.venv/bin/python -B -m unittest serve.test_lifecycle serve.test_mcp serve.test_structured serve.test_monitor \
serve.test_detok serve.test_winjob # OK
python3 -m unittest tools.test_setup_{amd,choices,draft_vocab,lowram,pins,prompts,risk,rope,unsloth} # OK
```
`ClientDisconnect` passed in 5 consecutive runs with no tracebacks in the server output (`(disconnect, cancel=True)`).
## Portability
- **Windows:** `select.select` works on sockets there, and `socket.MSG_PEEK` exists. A graceful close reads as `b''`. An abortive close raises WSAECONNRESET / WSAECONNABORTED, which Python maps to `ConnectionResetError` / `ConnectionAbortedError`; both are `ConnectionError`, so both mean "gone".
- **No socket mode changes:** no `fcntl`, `poll` or `MSG_DONTWAIT`, and no non-blocking or timeout changes on the socket.
- **One thread per socket:** the check runs on the handler's own thread, so no other thread touches the socket.
## Known limitation
A client that *half-closes*, sending its request and then `shutdown(SHUT_WR)` while still waiting for the answer, reads as EOF and would be cancelled. cpp-httplib's `is_socket_alive`, which llama.cpp's server uses, behaves the same way. curl, browsers, `requests`, `httpx`, the OpenAI/Anthropic SDKs and Node's fetch (undici) don't do this.
## What remains for #431
- **Why it was late:** the stream path found a hang-up only when a write failed, and on TCP the first write to a closed socket still succeeds. The peek now finds it at the first heartbeat after the disconnect.
- **What still waits:** detection still waits for the next PP line or the 10 s heartbeat, because the HTTP thread blocks in `lines.get(timeout=10)` until the engine prints something. The worst case is two chunk boundaries after the disconnect, where it was three.
- **Unexplained:** the reported 22 s (nearly the whole read) is more than this one-chunk delay explains, so something else may be involved on the HIP build. An `STRATA_TRACE=1` log from the reporter would show when `STOP` arrives.
- **A full fix** would watch the socket independently of engine output, for example a per-request watcher feeding a 1 s cancel poll in `generate()`. That's a larger change, left out here.
### Test machine
Ubuntu 24.04.4 (kernel 6.8), Python 3.12.3 (the install's venv), AMD Ryzen Threadripper 3960X (24C/48T), 128 GB DDR4, RTX 2080 Ti 11 GB. The tests use the mock engine (no GPU) and ran beside a live Strata server on the same PC, which was not touched.
### On the real engine
I applied the change to this PC's live install (engine 0.1.33, Qwen3.8-Flash-Next IQ3_S, 256K context, CUDA, RTX 2080 Ti) and ran a non-streaming request with `max_tokens: 2000`. The client gave up after 4 s (`curl -m 4`, exit 28). Then I sent a short request.
```
[strata] done: 143 tokens in 5 s (35.1 tok/s) (disconnect, cancel=True), expert cache 44.7% hit
[strata] done: 2 tokens in 0 s (20.9 tok/s) (stop, cancel=False), expert cache 4.4% hit
```
The engine stopped about 1 s after the hang-up. The next request was answered 5.3 s after the first one started, where before it would have waited for all 2,000 tokens (about 55 s at 35 tok/s).
**Not run:** cancelling a non-streaming request during a long prompt read on the real engine. The mock test covers it, and the engine-side stop at a chunk boundary is existing behaviour.
Mehr auf der Site
Links zu Install, Modellen, Releases.