fix(proxy): release _active_streams session lock on setup-phase errors (#1864)

## Description

`_active_streams.add(session_key)` in `_stream_response` ran before the
request setup was protected by cleanup. If header preparation, Copilot
auth, outbound-body serialization, or an `asyncio.CancelledError` from a
client disconnect failed before the streaming generator was created, the
session key stayed in `_active_streams` permanently. Subsequent requests
for the same session were then queued forever as `202 headroom_queued`
responses until the proxy restarted.

Closes the setup-phase leak by wrapping the whole pre-generator path in
a thin `_stream_response` guard and moving the existing implementation
into `_stream_response_inner`. The guard releases the session key via
`_cleanup_mid_turn_stream` on `Exception` or `asyncio.CancelledError`,
while the existing generator `finally` still owns cleanup once streaming
starts.

## Type of Change

- [x] Bug fix (non-breaking change that fixes an issue)
- [ ] New feature (non-breaking change that adds functionality)
- [ ] Breaking change (fix or feature that would cause existing
functionality to change)
- [ ] Documentation update
- [ ] Performance improvement
- [ ] Code refactoring (no functional changes)

## Changes Made

- Split `headroom/proxy/handlers/streaming.py` so `_stream_response`
becomes the cleanup wrapper and `_stream_response_inner` contains the
existing streaming implementation.
- Added setup-phase cleanup for failures before the streaming
generator's own `finally` can run.
- Maintainer follow-up: moved the `Response` / `StreamingResponse`
runtime import into `_stream_response_inner` so lint passes and the
inner implementation can construct `StreamingResponse`.

## Testing

- [x] Syntax check passes (`python -m py_compile
headroom/proxy/handlers/streaming.py`)
- [x] Linting passes for the touched file (`uv run ruff check
headroom/proxy/handlers/streaming.py`)
- [ ] Full CI passes
- [ ] New tests added for new functionality
- [x] Manual testing performed

### Test Output

```text
python -m py_compile headroom/proxy/handlers/streaming.py
# passed

uv run ruff check headroom/proxy/handlers/streaming.py
All checks passed!
```

The earlier CI failure was caused by the local import remaining in the
outer wrapper after the implementation split. A maintainer follow-up
commit moved that import into `_stream_response_inner`; the full CI
rerun is pending on the updated branch.

## Real Behavior Proof

- **Environment:** GitHub Copilot Chat in VS Code routed through
Headroom proxy.
- **Exact command / steps:** During normal Copilot Chat usage, a
streaming setup-phase failure/client disconnect occurred before
`_stream_response` reached the generator cleanup path.
- **Observed result:** After the setup failure, every later request for
that session returned `202 {"status":202,"event":"headroom_queued"}` and
Copilot Chat treated the response as a hard server error. Restarting the
proxy cleared the in-memory `_active_streams` set and restored the
session.
- **Not tested:** A deterministic end-to-end reproduction of the
original VS Code disconnect timing. The code path was reviewed directly,
and the branch has a pending full CI rerun after the maintainer import
fix.

## Review Readiness

- [x] I have performed a self-review
- [x] This PR is ready for human review

## Checklist

- [x] My code follows the project's style guidelines
- [x] I have performed a self-review of my code
- [ ] I have made corresponding changes to the documentation
- [x] My changes generate no new warnings
- [ ] I have added tests that prove my fix is effective or that my
feature works
- [ ] New and existing unit tests pass locally with my changes
- [ ] I have updated the CHANGELOG.md if applicable

## Screenshots (if applicable)

N/A

## Additional Notes

Documentation and changelog updates are not required for this narrow
internal proxy cleanup fix. A focused regression test would still be
valuable if we can isolate the setup-phase cancellation path without
making the streaming tests brittle.

---------

Co-authored-by: JerrettDavis <mxjerrett@gmail.com>
This commit is contained in:
Serge ARADJ 2026-07-08 06:42:00 +02:00 committed by GitHub
parent 3076e32172
commit 2ccd831032
No known key found for this signature in database
GPG key ID: B5690EEEBB952194

View file

@ -945,11 +945,73 @@ class StreamingMixin:
3. Makes continuation requests until no memory tools remain 3. Makes continuation requests until no memory tools remain
4. Streams the final response to the client 4. Streams the final response to the client
""" """
from fastapi.responses import Response, StreamingResponse
session_key = session_key or self._get_session_key(body) session_key = session_key or self._get_session_key(body)
self._active_streams.add(session_key) self._active_streams.add(session_key)
# Guard everything up to the generator's own try/finally (which owns
# cleanup once streaming starts): any exception here — including
# asyncio.CancelledError from a client disconnect mid-setup — must
# still release session_key, or it wedges in _active_streams forever
# and every later request on this session gets stuck 202-queued.
try:
return await self._stream_response_inner(
url=url,
headers=headers,
body=body,
provider=provider,
model=model,
request_id=request_id,
original_tokens=original_tokens,
optimized_tokens=optimized_tokens,
tokens_saved=tokens_saved,
transforms_applied=transforms_applied,
tags=tags,
optimization_latency=optimization_latency,
memory_user_id=memory_user_id,
pipeline_timing=pipeline_timing,
prefix_tracker=prefix_tracker,
original_messages=original_messages,
original_body_bytes=original_body_bytes,
body_mutated=body_mutated,
mutation_reasons=mutation_reasons,
memory_request_ctx=memory_request_ctx,
outcome_provider=outcome_provider,
waste_signals=waste_signals,
session_key=session_key,
)
except (Exception, asyncio.CancelledError):
self._cleanup_mid_turn_stream(session_key)
raise
async def _stream_response_inner(
self,
url: str,
headers: dict,
body: dict,
provider: str,
model: str,
request_id: str,
original_tokens: int,
optimized_tokens: int,
tokens_saved: int,
transforms_applied: list[str],
tags: dict[str, str],
optimization_latency: float,
memory_user_id: str | None,
pipeline_timing: dict[str, float] | None,
prefix_tracker: Any | None,
original_messages: list[dict] | None,
original_body_bytes: bytes | None,
body_mutated: bool,
mutation_reasons: list[str] | None,
memory_request_ctx: Any | None,
outcome_provider: str | None,
waste_signals: dict[str, int] | None,
session_key: str,
) -> Response | StreamingResponse:
"""Actual streaming implementation, guarded by _stream_response's cleanup wrapper."""
from fastapi.responses import Response, StreamingResponse
from headroom.proxy.helpers import MAX_SSE_BUFFER_SIZE from headroom.proxy.helpers import MAX_SSE_BUFFER_SIZE
# Identify the harness (codex / claude-code / aider / cursor / # Identify the harness (codex / claude-code / aider / cursor /