From 22a4758ba8b8b58b9d792899c9e96644c104ea3f Mon Sep 17 00:00:00 2001 From: Jaeyeon Lee Date: Tue, 29 Sep 2026 15:20:08 +0200 Subject: [PATCH 1/4] 198: Track client-closed requests as 499 in usage logging --- aqueduct/gateway/tests/test_stream_status.py | 105 ++++++++++++++++++ aqueduct/gateway/views/decorators.py | 4 +- aqueduct/gateway/views/utils.py | 37 +++--- .../management/tests/test_usage_dashboard.py | 61 ++++++++++ 4 files changed, 189 insertions(+), 18 deletions(-) create mode 100644 aqueduct/gateway/tests/test_stream_status.py create mode 100644 aqueduct/management/tests/test_usage_dashboard.py diff --git a/aqueduct/gateway/tests/test_stream_status.py b/aqueduct/gateway/tests/test_stream_status.py new file mode 100644 index 00000000..43c29020 --- /dev/null +++ b/aqueduct/gateway/tests/test_stream_status.py @@ -0,0 +1,105 @@ +"""Tests for the streaming status-code recording in the request log. + +Covers the mapping the stream generator writes to ``Request.status_code``: + +- Clean completion -> 200 +- Client closes -> 499 (client-closed-request) +- Upstream failure -> 500 +""" + +import json +import unittest + +from asgiref.sync import async_to_sync +from django.contrib.auth import get_user_model +from django.contrib.auth.models import Group +from django.test import TestCase + +from gateway.views.utils import _openai_stream +from management.models import Org, Request, Token, UserGroup, UserProfile + +User = get_user_model() + + +class _Chunk: + """Minimal stand-in for a streamed chunk that only needs model_dump_json.""" + + def model_dump_json(self, *args, **kwargs) -> str: + return json.dumps({"choices": [{"delta": {"content": "hi"}}]}) + + +class OpenAIStreamStatusTests(TestCase): + def setUp(self): + self.org = Org.objects.create(name="stream-org") + self.user = User.objects.create_user(username="streamuser", email="stream@example.com") + UserProfile.objects.create(user=self.user, org=self.org) + Group.objects.get_or_create(name=UserGroup.USER.value) + self.token = Token(name="stream-token", user=self.user) + self.token._set_new_key() + self.token.save() + + def _request_log(self) -> Request: + request_log = Request(token=self.token, model="gpt-4.1-nano") + request_log.save() + return request_log + + def test_clean_completion_records_200(self): + async def fake_stream(): + yield _Chunk() + yield _Chunk() + + request_log = self._request_log() + stream = _openai_stream(fake_stream(), request_log) + + async def consume_all(): + async for _ in stream: + pass + + async_to_sync(consume_all)() + request_log.refresh_from_db() + self.assertEqual(request_log.status_code, 200) + + def test_client_close_records_499(self): + """When the client closes the connection mid-stream we record 499.""" + + async def fake_stream(): + yield _Chunk() + yield _Chunk() + yield _Chunk() + + request_log = self._request_log() + stream = _openai_stream(fake_stream(), request_log) + + async def consume_then_disconnect(): + it = stream.__aiter__() + await it.__anext__() + await it.__anext__() + await stream.aclose() + + async_to_sync(consume_then_disconnect)() + request_log.refresh_from_db() + self.assertEqual(request_log.status_code, 499) + + def test_upstream_failure_records_500(self): + """An upstream error part-way through the stream is recorded as 500.""" + + async def failing_stream(): + yield _Chunk() + raise RuntimeError("upstream exploded") + + request_log = self._request_log() + stream = _openai_stream(failing_stream(), request_log) + + async def consume_until_failure(): + async for _ in stream: + pass + + with self.assertRaises(RuntimeError): + async_to_sync(consume_until_failure)() + + request_log.refresh_from_db() + self.assertEqual(request_log.status_code, 500) + + +if __name__ == "__main__": + unittest.main() diff --git a/aqueduct/gateway/views/decorators.py b/aqueduct/gateway/views/decorators.py index c2d3473e..15f5ac0a 100644 --- a/aqueduct/gateway/views/decorators.py +++ b/aqueduct/gateway/views/decorators.py @@ -427,7 +427,9 @@ async def wrapper(request: ASGIRequest, *args: Any, **kwargs: Any) -> ViewResult ) request_log.processing_time_ms = int((response_start_time - kwargs["request_start"]) * 1000) request_log.response_time_ms = int((end_time - response_start_time) * 1000) - request_log.status_code = result.status_code + + if not isinstance(result, StreamingHttpResponse): + request_log.status_code = result.status_code await request_log.asave() return result diff --git a/aqueduct/gateway/views/utils.py b/aqueduct/gateway/views/utils.py index cda617d9..fdb45c0f 100644 --- a/aqueduct/gateway/views/utils.py +++ b/aqueduct/gateway/views/utils.py @@ -1,3 +1,4 @@ +import asyncio import json import logging import time @@ -61,27 +62,29 @@ def _openai_stream( async def _stream() -> AsyncGenerator[str, None]: token_usage = Usage(0, 0) - async for chunk in stream: - chunk_str = chunk.model_dump_json(exclude_none=True, exclude_unset=True) + try: + async for chunk in stream: + chunk_str = chunk.model_dump_json(exclude_none=True, exclude_unset=True) - # Extract token usage from this chunk - chunk_usage = _get_token_usage(chunk_str.encode("utf-8")) + chunk_usage = _get_token_usage(chunk_str.encode("utf-8")) - # Only update if we got actual usage data (non-zero tokens) - if chunk_usage.input_tokens > 0 or chunk_usage.output_tokens > 0: - token_usage = chunk_usage + if chunk_usage.input_tokens > 0 or chunk_usage.output_tokens > 0: + token_usage = chunk_usage - try: yield f"data: {chunk_str}\n\n" - except Exception as e: - yield f"data: {e!s}\n\n" - - end_time = time.monotonic() - request_log.token_usage = token_usage - request_log.response_time_ms = int((end_time - start_time) * 1000) - await request_log.asave() - # Streaming is done, yield the [DONE] chunk - yield "data: [DONE]\n\n" + + request_log.status_code = 200 + yield "data: [DONE]\n\n" + except (asyncio.CancelledError, GeneratorExit): + request_log.status_code = 499 + raise + except Exception: + request_log.status_code = 500 + raise + finally: + request_log.token_usage = token_usage + request_log.response_time_ms = int((time.monotonic() - start_time) * 1000) + await request_log.asave() return _stream() diff --git a/aqueduct/management/tests/test_usage_dashboard.py b/aqueduct/management/tests/test_usage_dashboard.py new file mode 100644 index 00000000..4c0df823 --- /dev/null +++ b/aqueduct/management/tests/test_usage_dashboard.py @@ -0,0 +1,61 @@ +"""Tests the usage dashboard counts 499 (client-closed) as a failed request. + +A 499 status code is >= 400, so the dashboard's ``failed_requests`` stat +(``status_code__gte=400``) must include it — confirming client-closed requests +show up in the dashboard's failed-requests count. +""" + +from pathlib import Path + +from django.contrib.auth import get_user_model +from django.contrib.auth.models import Group +from django.test import TestCase, override_settings +from django.urls import reverse + +from management.models import Org, Request, Token, UserGroup, UserProfile + +User = get_user_model() + +ROOT = Path(__file__).resolve().parents[3] + +ALLOWED_MODEL = "gpt-4.1-nano" + + +@override_settings(LITELLM_ROUTER_CONFIG_FILE_PATH=str(ROOT / "example_router_config.yaml")) +class UsageDashboardFailedRequestTests(TestCase): + def setUp(self): + self.org = Org.objects.create(name="usage-org") + self.user = User.objects.create_user(username="usageuser", email="usage@example.com") + UserProfile.objects.create(user=self.user, org=self.org) + Group.objects.get_or_create(name=UserGroup.USER.value) + self.token = Token(name="usage-token", user=self.user) + self.token._set_new_key() + self.token.save() + self.client.force_login(self.user) + + def _add_request(self, status_code: int) -> Request: + return Request.objects.create( + token=self.token, + model=ALLOWED_MODEL, + status_code=status_code, + user_id=self.user.email, + path="/chat/completions", + ) + + def test_499_counts_as_failed_request(self): + self._add_request(status_code=499) + + resp = self.client.get(reverse("usage")) + self.assertEqual(resp.status_code, 200) + self.assertEqual(resp.context["failed_requests"], 1) + self.assertEqual(resp.context["total_requests"], 1) + + def test_499_and_500_both_count_as_failed(self): + self._add_request(status_code=200) + self._add_request(status_code=499) + self._add_request(status_code=500) + + resp = self.client.get(reverse("usage")) + self.assertEqual(resp.status_code, 200) + self.assertEqual(resp.context["total_requests"], 3) + self.assertEqual(resp.context["failed_requests"], 2) From 9a55f6b37e8e023165abd7ec6ea99d82aa3c31c9 Mon Sep 17 00:00:00 2001 From: Jaeyeon Lee Date: Mon, 5 Oct 2026 15:14:46 +0200 Subject: [PATCH 2/4] Add non-streaming 499 error --- aqueduct/gateway/views/decorators.py | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/aqueduct/gateway/views/decorators.py b/aqueduct/gateway/views/decorators.py index 15f5ac0a..94a25058 100644 --- a/aqueduct/gateway/views/decorators.py +++ b/aqueduct/gateway/views/decorators.py @@ -1,3 +1,4 @@ +import asyncio import base64 import io import json @@ -419,7 +420,15 @@ async def wrapper(request: ASGIRequest, *args: Any, **kwargs: Any) -> ViewResult log.debug("Initial request log object created.") response_start_time = time.monotonic() - result: HttpResponse | StreamingHttpResponse = await view_func(request, *args, **kwargs) + + try: + result: HttpResponse | StreamingHttpResponse = await view_func(request, *args, **kwargs) + except asyncio.CancelledError: + request_log.status_code = 499 + request_log.response_time_ms = int((time.monotonic() - response_start_time) * 1000) + await request_log.asave() + raise + end_time = time.monotonic() assert "request_start" in kwargs, ( From 235c4d6a57313e9d45edd9d6c4f28026064ca681 Mon Sep 17 00:00:00 2001 From: Jaeyeon Lee Date: Mon, 5 Oct 2026 18:39:57 +0200 Subject: [PATCH 3/4] Apply review --- .../gateway/tests/test_log_request_cancel.py | 70 +++++++++++++++++++ aqueduct/gateway/tests/test_stream_status.py | 5 -- 2 files changed, 70 insertions(+), 5 deletions(-) create mode 100644 aqueduct/gateway/tests/test_log_request_cancel.py diff --git a/aqueduct/gateway/tests/test_log_request_cancel.py b/aqueduct/gateway/tests/test_log_request_cancel.py new file mode 100644 index 00000000..b631e44f --- /dev/null +++ b/aqueduct/gateway/tests/test_log_request_cancel.py @@ -0,0 +1,70 @@ +"""Tests for non-streaming client-disconnect handling in ``log_request``. + +When a client disconnects while the view is still awaiting the upstream LLM +(before any response has been returned), Django's ASGI handler cancels the +request task, which raises ``asyncio.CancelledError`` inside the view. +``log_request`` catches that and records status 499 on the request log. + +This asserts the non-streaming behaviour, complementing the streaming tests +in ``test_stream_status.py`` (which exercise ``_openai_stream`` directly). +""" + +import asyncio +import time +from types import SimpleNamespace + +from asgiref.sync import async_to_sync +from django.contrib.auth import get_user_model +from django.contrib.auth.models import Group +from django.test import TestCase + +from gateway.views.decorators import log_request +from management.models import Org, Request, Token, UserGroup, UserProfile + +User = get_user_model() + + +def _make_request(path: str = "/v1/chat/completions") -> SimpleNamespace: + """Minimal stand-in for an ASGI request, exposing only what log_request reads.""" + return SimpleNamespace( + path=path, + method="POST", + headers=SimpleNamespace(get=lambda key, default="": default), + META=SimpleNamespace(get=lambda key, default="": default), + ) + + +class LogRequestCancelTests(TestCase): + def setUp(self): + self.org = Org.objects.create(name="cancel-org") + self.user = User.objects.create_user(username="canceluser", email="cancel@example.com") + UserProfile.objects.create(user=self.user, org=self.org) + Group.objects.get_or_create(name=UserGroup.USER.value) + self.token = Token(name="cancel-token", user=self.user) + self.token._set_new_key() + self.token.save() + + def test_non_streaming_client_disconnect_records_499(self): + """A disconnect while awaiting the upstream is recorded as 499.""" + started = asyncio.Event() + + async def hanging_view(request, *args, **kwargs): + started.set() + await asyncio.sleep(3600) # simulate awaiting the upstream LLM + + wrapped = log_request(hanging_view) + + async def run(): + task = asyncio.create_task( + wrapped(_make_request(), token=self.token, request_start=time.monotonic()) + ) + await started.wait() + task.cancel() + with self.assertRaises(asyncio.CancelledError): + await task + + async_to_sync(run)() + + request_log = Request.objects.get(token=self.token) + self.assertEqual(request_log.status_code, 499) + self.assertIsNotNone(request_log.response_time_ms) diff --git a/aqueduct/gateway/tests/test_stream_status.py b/aqueduct/gateway/tests/test_stream_status.py index 43c29020..1b795aae 100644 --- a/aqueduct/gateway/tests/test_stream_status.py +++ b/aqueduct/gateway/tests/test_stream_status.py @@ -8,7 +8,6 @@ """ import json -import unittest from asgiref.sync import async_to_sync from django.contrib.auth import get_user_model @@ -99,7 +98,3 @@ async def consume_until_failure(): request_log.refresh_from_db() self.assertEqual(request_log.status_code, 500) - - -if __name__ == "__main__": - unittest.main() From bd0e7a99c35dfa13696085f90c1eb30063599a25 Mon Sep 17 00:00:00 2001 From: Jaeyeon Lee Date: Mon, 5 Oct 2026 19:17:25 +0200 Subject: [PATCH 4/4] Apply review --- .../gateway/tests/test_log_request_cancel.py | 17 +++++------------ aqueduct/gateway/tests/test_stream_status.py | 17 +++++------------ 2 files changed, 10 insertions(+), 24 deletions(-) diff --git a/aqueduct/gateway/tests/test_log_request_cancel.py b/aqueduct/gateway/tests/test_log_request_cancel.py index b631e44f..856c1ccb 100644 --- a/aqueduct/gateway/tests/test_log_request_cancel.py +++ b/aqueduct/gateway/tests/test_log_request_cancel.py @@ -12,16 +12,13 @@ import asyncio import time from types import SimpleNamespace +from typing import ClassVar from asgiref.sync import async_to_sync -from django.contrib.auth import get_user_model -from django.contrib.auth.models import Group from django.test import TestCase from gateway.views.decorators import log_request -from management.models import Org, Request, Token, UserGroup, UserProfile - -User = get_user_model() +from management.models import Request, Token def _make_request(path: str = "/v1/chat/completions") -> SimpleNamespace: @@ -35,14 +32,10 @@ def _make_request(path: str = "/v1/chat/completions") -> SimpleNamespace: class LogRequestCancelTests(TestCase): + fixtures: ClassVar[list[str]] = ["gateway_data.json"] + def setUp(self): - self.org = Org.objects.create(name="cancel-org") - self.user = User.objects.create_user(username="canceluser", email="cancel@example.com") - UserProfile.objects.create(user=self.user, org=self.org) - Group.objects.get_or_create(name=UserGroup.USER.value) - self.token = Token(name="cancel-token", user=self.user) - self.token._set_new_key() - self.token.save() + self.token = Token.objects.get(name="My Token") def test_non_streaming_client_disconnect_records_499(self): """A disconnect while awaiting the upstream is recorded as 499.""" diff --git a/aqueduct/gateway/tests/test_stream_status.py b/aqueduct/gateway/tests/test_stream_status.py index 1b795aae..0804d421 100644 --- a/aqueduct/gateway/tests/test_stream_status.py +++ b/aqueduct/gateway/tests/test_stream_status.py @@ -8,16 +8,13 @@ """ import json +from typing import ClassVar from asgiref.sync import async_to_sync -from django.contrib.auth import get_user_model -from django.contrib.auth.models import Group from django.test import TestCase from gateway.views.utils import _openai_stream -from management.models import Org, Request, Token, UserGroup, UserProfile - -User = get_user_model() +from management.models import Request, Token class _Chunk: @@ -28,14 +25,10 @@ def model_dump_json(self, *args, **kwargs) -> str: class OpenAIStreamStatusTests(TestCase): + fixtures: ClassVar[list[str]] = ["gateway_data.json"] + def setUp(self): - self.org = Org.objects.create(name="stream-org") - self.user = User.objects.create_user(username="streamuser", email="stream@example.com") - UserProfile.objects.create(user=self.user, org=self.org) - Group.objects.get_or_create(name=UserGroup.USER.value) - self.token = Token(name="stream-token", user=self.user) - self.token._set_new_key() - self.token.save() + self.token = Token.objects.get(name="My Token") def _request_log(self) -> Request: request_log = Request(token=self.token, model="gpt-4.1-nano")