From 103596f2b421d3501143f03e77fe86ddab76ea3b Mon Sep 17 00:00:00 2001 From: redshift <213178690+1ftredsh@users.noreply.github.com> Date: Wed, 30 Sep 2026 20:01:43 +0800 Subject: [PATCH] fix: bound billed request lifetimes and recover abandoned reservations --- .../a73d19b6c204_reservation_deadlines.py | 42 ++++++ repro/IMPLEMENTATION.md | 46 ++++++ repro/dummy_upstream.py | 45 ++++++ repro/probe.py | 60 ++++++++ repro/results-final.txt | 19 +++ repro/results.txt | 19 +++ repro/router-final.log | 108 ++++++++++++++ repro/router-first.log | 115 +++++++++++++++ routstr/auth.py | 35 ++++- routstr/core/db.py | 12 +- routstr/core/lifecycle.py | 135 ++++++++++++++++++ routstr/core/main.py | 5 + routstr/core/settings.py | 10 ++ routstr/upstream/stream_ownership.py | 7 +- tests/unit/test_request_lifecycle.py | 57 ++++++++ tests/unit/test_stale_reservations.py | 59 ++++++++ 16 files changed, 770 insertions(+), 4 deletions(-) create mode 100644 migrations/versions/a73d19b6c204_reservation_deadlines.py create mode 100644 repro/IMPLEMENTATION.md create mode 100644 repro/dummy_upstream.py create mode 100644 repro/probe.py create mode 100644 repro/results-final.txt create mode 100644 repro/results.txt create mode 100644 repro/router-final.log create mode 100644 repro/router-first.log create mode 100644 routstr/core/lifecycle.py create mode 100644 tests/unit/test_request_lifecycle.py diff --git a/migrations/versions/a73d19b6c204_reservation_deadlines.py b/migrations/versions/a73d19b6c204_reservation_deadlines.py new file mode 100644 index 00000000..74702d61 --- /dev/null +++ b/migrations/versions/a73d19b6c204_reservation_deadlines.py @@ -0,0 +1,42 @@ +"""Immutable reservation start and absolute recovery deadline. + +Revision ID: a73d19b6c204 +Revises: e4c7a1b9d520 +""" + +import time + +import sqlalchemy as sa +from alembic import op + +revision = "a73d19b6c204" +down_revision = "e4c7a1b9d520" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.add_column( + "reservation_releases", sa.Column("started_at", sa.Integer(), nullable=True) + ) + op.add_column( + "reservation_releases", sa.Column("expires_at", sa.Integer(), nullable=True) + ) + op.create_index( + "ix_reservation_releases_expires_at", "reservation_releases", ["expires_at"] + ) + # Original ages are unknowable for renewed legacy rows. Give them a finite + # migration grace period; deploy only after draining old workers. + op.execute( + sa.text( + "UPDATE reservation_releases SET expires_at = :expiry WHERE status = 'active'" + ).bindparams(expiry=int(time.time()) + 1830) + ) + + +def downgrade() -> None: + op.drop_index( + "ix_reservation_releases_expires_at", table_name="reservation_releases" + ) + op.drop_column("reservation_releases", "expires_at") + op.drop_column("reservation_releases", "started_at") diff --git a/repro/IMPLEMENTATION.md b/repro/IMPLEMENTATION.md new file mode 100644 index 00000000..1c231f1f --- /dev/null +++ b/repro/IMPLEMENTATION.md @@ -0,0 +1,46 @@ +# Reservation lifecycle implementation and validation + +Branch: fix/reservation-lifecycle. Baseline: 96c8e2f7. + +## Implemented + +- Outermost pure-ASGI lifecycle supervision with one coordinated receive consumer, explicit disconnect monitoring, cancellation, and exact reservation fallback cleanup. +- Finite overall request lifetime (MAX_REQUEST_LIFETIME_SECONDS, default 1800), downstream send timeout (DOWNSTREAM_SEND_TIMEOUT_SECONDS, default 60), and cleanup timeout (REQUEST_CLEANUP_TIMEOUT_SECONDS, default 30). +- Lifecycle identity shared through context across middleware tasks; reservation replacements are registered for exact cleanup. +- Heartbeats stop on lifecycle termination or local maximum age. +- Persistent stream finalization has a finite cleanup budget. +- Durable immutable started_at and expires_at columns; expiry covers remaining request lifetime plus settlement grace, including provider fallback without restarting the original deadline. +- Renewal and charge claims refuse expired reservations. Sweeping can release absolute-expired reservations even when their renewable timestamp is fresh. +- Migration grants legacy active rows 1830 seconds of grace; original ages are not fabricated. Drain old workers before deployment. + +## Verification + +Run from worktree with PYTHONPATH=$PWD because the shared root virtual environment's editable install points at the original checkout: + +PYTHONPATH=$PWD ../../.venv/bin/pytest tests/unit/test_request_lifecycle.py tests/unit/test_stale_reservations.py tests/unit/test_streaming_billing_finalization.py tests/integration/test_negative_available_balance_repro.py -q + +64 tests passed. Ruff checks passed on changed files. Full-project mypy was attempted but did not finish within the tool timeout; no successful typecheck is claimed. + +Final built image: localhost/routstr-reserved-repro:fix, ef81426ad79e3d14ec462a39ab1f7481fd0cb410a9cb93cb42de41e8b3523869. + +Container tests used real TCP, full middleware stack, frozen image dependencies, isolated SQLite and synthetic balances. Read timeout 3s, lifetime 15s, delivery timeout 2s, cleanup timeout 3s, stale timeout 6s. + +Reused the main probe on ports 18100/18101. Results in results-final.txt and router-final.log: + +- Finite and silent streams settled. +- Header wait released its reservation. +- Disconnected endless stream no longer retained its reservation. +- Non-reading flood client hit bounded delivery/cleanup. +- Connected keepalive-only stream terminated at maximum lifetime. +- After the background-sweep interval and all client closures: every key reserved_balance=0, no active durable reservations. Explicit database assertions passed. +- Router shut down within the 10-second grace without SIGKILL. Dummy upstream still required SIGKILL: its fixture deliberately sleeps/open-streams and is not patched router code. + +Actual mint payout was not tested. Protocol errors on already-started streams when deadlines interrupt them are expected; an HTTP status cannot be replaced after headers are sent. + +## Financial policy / limitations + +The lifecycle first lets existing finalization run within a bounded budget. If still active, fallback releases only that reservation; late charge is fenced by terminal state. This can forgo charging observed output on failed settlement. It prioritizes freeing customer funds over leaving them locked; review this policy before deployment. Upstream compute may continue remotely even after local connection closure. + +This implementation does not complete every proposed hardening idea: provider cancellation APIs, full observability, per-record unexpected DB-failure isolation, legacy NULL aggregate background reconciliation, multi-worker/alternate-route network matrix and DB-outage injection remain follow-up work. No dependency upgrade was needed for the tested cases because explicit disconnect supervision avoids relying solely on send errors. + +All reproduction containers are stopped. Original node data/configuration is untouched. Source changes are uncommitted in the worktree for review. diff --git a/repro/dummy_upstream.py b/repro/dummy_upstream.py new file mode 100644 index 00000000..1be75ecc --- /dev/null +++ b/repro/dummy_upstream.py @@ -0,0 +1,45 @@ +"""Loopback-only streaming fixture; no router monkeypatches.""" +import asyncio +import json +import time +from fastapi import FastAPI, Request +from fastapi.responses import StreamingResponse + +app = FastAPI() +events = [] + +@app.get('/events') +async def history(): + return events + +@app.get('/v1/models') +async def models(): + return {'object': 'list', 'data': [{'id': 'gpt-4o-mini', 'object': 'model', 'created': 1, 'owned_by': 'repro'}]} + +@app.post('/v1/chat/completions') +async def completions(request: Request): + body = await request.json() + mode = body.get('messages', [{}])[0].get('content', 'finite') + events.append({'event': 'start', 'mode': mode, 'time': time.time()}) + if mode.startswith('header'): + await asyncio.sleep(3600) + async def stream(): + count = 0 + try: + while True: + if mode.startswith('keepalive'): + yield ': ping\n\n' + else: + chunk = {'id': 'repro', 'object': 'chat.completion.chunk', 'created': int(time.time()), 'model': 'gpt-4o-mini', 'choices': [{'index': 0, 'delta': {'content': 'x' * (65536 if mode.startswith('flood') else 1)}, 'finish_reason': None}]} + yield 'data: ' + json.dumps(chunk) + '\n\n' + count += 1 + if mode == 'finite' and count >= 3: + yield 'data: ' + json.dumps({'id': 'repro', 'object': 'chat.completion.chunk', 'model': 'gpt-4o-mini', 'choices': [], 'usage': {'prompt_tokens': 1, 'completion_tokens': count, 'total_tokens': count + 1}}) + '\n\n' + yield 'data: [DONE]\n\n' + return + await asyncio.sleep(3600 if mode.startswith('silent') else (0.001 if mode.startswith('flood') else 0.5)) + finally: + event = {'event': 'close', 'mode': mode, 'chunks': count, 'time': time.time()} + events.append(event) + print(json.dumps(event), flush=True) + return StreamingResponse(stream(), media_type='text/event-stream') diff --git a/repro/probe.py b/repro/probe.py new file mode 100644 index 00000000..93d71a75 --- /dev/null +++ b/repro/probe.py @@ -0,0 +1,60 @@ +import asyncio +import json +import socket +import subprocess +import time +import httpx + +BASE='http://127.0.0.1:18100' + +def snapshot(): + code="import sqlite3,json,time; c=sqlite3.connect('/tmp/reserved-fix.db'); c.row_factory=sqlite3.Row; print(json.dumps({'time':time.time(),'keys':[dict(r) for r in c.execute(\"select hashed_key,balance,reserved_balance,reserved_at from api_keys where hashed_key like 'main-%'\")],'rows':[dict(r) for r in c.execute(\"select * from reservation_releases where key_hash like 'main-%'\")]}))" + return json.loads(subprocess.check_output(['podman','exec','reserved-router-fix','/.venv/bin/python','-c',code],text=True)) + +async def consume(mode): + try: + async with httpx.AsyncClient(timeout=None) as c: + async with c.stream('POST',BASE+'/v1/chat/completions',headers={'Authorization':'Bearer sk-main-'+mode},json={'model':'gpt-4o-mini','messages':[{'role':'user','content':mode}],'stream':True,'max_tokens':10}) as r: + print('STREAM',mode,r.status_code,flush=True) + async for _ in r.aiter_bytes(): pass + print('ENDED',mode,flush=True) + except asyncio.CancelledError: + print('CLIENT_DISCONNECTED',mode,flush=True) + raise + except Exception as e: + print('CLIENT_ERROR',mode,type(e).__name__,str(e),flush=True) + +async def report(label): + print(label,json.dumps(snapshot()),flush=True) + async with httpx.AsyncClient(timeout=5) as c: + for mode in ['silent-disconnect','endless-disconnect','keepalive','flood','header']: + # Only attempt payout while reserved: avoid requiring a real mint. + if next(k for k in snapshot()['keys'] if k['hashed_key']=='main-'+mode)['reserved_balance']: + r=await c.post(BASE+'/v1/wallet/refund',headers={'Authorization':'Bearer sk-main-'+mode}) + print('REFUND',mode,r.status_code,r.text,flush=True) + print('UPSTREAM_EVENTS',json.dumps((await c.get('http://127.0.0.1:18101/events')).json()),flush=True) + +async def main(): + modes=['finite','silent','silent-disconnect','endless-disconnect','keepalive','header'] + tasks={m:asyncio.create_task(consume(m)) for m in modes} + # Real client with a small receive buffer, never draining the HTTP response. + sock=socket.socket(); sock.setsockopt(socket.SOL_SOCKET,socket.SO_RCVBUF,1024); sock.connect(('127.0.0.1',18100)) + body=json.dumps({'model':'gpt-4o-mini','messages':[{'role':'user','content':'flood'}],'stream':True,'max_tokens':10}).encode() + sock.sendall(b'POST /v1/chat/completions HTTP/1.1\r\nHost: localhost\r\nAuthorization: Bearer sk-main-flood\r\nContent-Type: application/json\r\nContent-Length: '+str(len(body)).encode()+b'\r\n\r\n'+body) + await asyncio.sleep(1) + for m in ['silent-disconnect','endless-disconnect']: + tasks[m].cancel() + await asyncio.gather(tasks['silent-disconnect'],tasks['endless-disconnect'],return_exceptions=True) + await asyncio.sleep(9) + await report('AT_10_SECONDS') + await asyncio.sleep(60) + await report('AFTER_SWEEP') + sock.close() + tasks['keepalive'].cancel() + await asyncio.gather(tasks['keepalive'],return_exceptions=True) + await asyncio.sleep(8) + await report('AFTER_ALL_CLIENTS_CLOSED') + for task in tasks.values(): task.cancel() + await asyncio.gather(*tasks.values(),return_exceptions=True) + +asyncio.run(main()) diff --git a/repro/results-final.txt b/repro/results-final.txt new file mode 100644 index 00000000..6545bf6d --- /dev/null +++ b/repro/results-final.txt @@ -0,0 +1,19 @@ +STREAM endless-disconnect 200 +STREAM keepalive 200 +STREAM finite 200 +STREAM silent 200 +STREAM silent-disconnect 200 +CLIENT_DISCONNECTED silent-disconnect +CLIENT_DISCONNECTED endless-disconnect +ENDED finite +ENDED silent +STREAM header 424 +ENDED header +AT_10_SECONDS {"time": 1790767971.1792026, "keys": [{"hashed_key": "main-finite", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-silent", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-silent-disconnect", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-endless-disconnect", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-keepalive", "balance": 1000000000, "reserved_balance": 12, "reserved_at": 1790767960}, {"hashed_key": "main-flood", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-header", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}], "rows": [{"id": "9b5b9063f6f045a290611e4144c897de", "key_hash": "main-flood", "billing_key_hash": "main-flood", "reserved_msats": 12, "status": "released", "created_at": 1790767962, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "efc8563301ba458890ee3ab2335e4e49", "key_hash": "main-endless-disconnect", "billing_key_hash": "main-endless-disconnect", "reserved_msats": 13, "status": "released", "created_at": 1790767960, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "cbed89b7ec444fea9789cf53b3e0f476", "key_hash": "main-keepalive", "billing_key_hash": "main-keepalive", "reserved_msats": 12, "status": "active", "created_at": 1790767971, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "1c56b25265df4743b4cff70dc57544c6", "key_hash": "main-finite", "billing_key_hash": "main-finite", "reserved_msats": 12, "status": "charged", "created_at": 1790767960, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "cc7f113eef004b7ba27bc761c5d9b9a1", "key_hash": "main-silent", "billing_key_hash": "main-silent", "reserved_msats": 12, "status": "charged", "created_at": 1790767963, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "6a1b5506b8ab4747884475075ebc38da", "key_hash": "main-silent-disconnect", "billing_key_hash": "main-silent-disconnect", "reserved_msats": 13, "status": "released", "created_at": 1790767961, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "1dfedc1c2cee407fb2cd1a28b1253c6d", "key_hash": "main-header", "billing_key_hash": "main-header", "reserved_msats": 12, "status": "released", "created_at": 1790767963, "started_at": 1790767960, "expires_at": 1790767978}]} +REFUND keepalive 400 {"detail":"Cannot refund key. There are ongoing requests for this api key.","request_id":"45f783cc-4c0b-4222-be13-8e726fc7cebc"} +UPSTREAM_EVENTS [{"event": "start", "mode": "flood", "time": 1790767717.2871263}, {"event": "start", "mode": "silent-disconnect", "time": 1790767717.2960703}, {"event": "start", "mode": "keepalive", "time": 1790767717.3026786}, {"event": "start", "mode": "header", "time": 1790767717.3104746}, {"event": "start", "mode": "endless-disconnect", "time": 1790767717.318341}, {"event": "start", "mode": "silent", "time": 1790767717.3653235}, {"event": "start", "mode": "finite", "time": 1790767717.3924189}, {"event": "close", "mode": "finite", "chunks": 3, "time": 1790767718.397855}, {"event": "start", "mode": "flood", "time": 1790767960.9594278}, {"event": "start", "mode": "endless-disconnect", "time": 1790767961.008998}, {"event": "start", "mode": "keepalive", "time": 1790767961.036595}, {"event": "start", "mode": "finite", "time": 1790767961.060875}, {"event": "start", "mode": "silent", "time": 1790767961.0869172}, {"event": "start", "mode": "silent-disconnect", "time": 1790767961.1580715}, {"event": "start", "mode": "header", "time": 1790767961.1944675}, {"event": "close", "mode": "finite", "chunks": 3, "time": 1790767962.064207}] +CLIENT_ERROR keepalive RemoteProtocolError peer closed connection without sending complete message body (incomplete chunked read) +AFTER_SWEEP {"time": 1790768033.9280283, "keys": [{"hashed_key": "main-finite", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-silent", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-silent-disconnect", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-endless-disconnect", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-keepalive", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-flood", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-header", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}], "rows": [{"id": "9b5b9063f6f045a290611e4144c897de", "key_hash": "main-flood", "billing_key_hash": "main-flood", "reserved_msats": 12, "status": "released", "created_at": 1790767962, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "efc8563301ba458890ee3ab2335e4e49", "key_hash": "main-endless-disconnect", "billing_key_hash": "main-endless-disconnect", "reserved_msats": 13, "status": "released", "created_at": 1790767960, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "cbed89b7ec444fea9789cf53b3e0f476", "key_hash": "main-keepalive", "billing_key_hash": "main-keepalive", "reserved_msats": 12, "status": "released", "created_at": 1790767975, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "1c56b25265df4743b4cff70dc57544c6", "key_hash": "main-finite", "billing_key_hash": "main-finite", "reserved_msats": 12, "status": "charged", "created_at": 1790767960, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "cc7f113eef004b7ba27bc761c5d9b9a1", "key_hash": "main-silent", "billing_key_hash": "main-silent", "reserved_msats": 12, "status": "charged", "created_at": 1790767963, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "6a1b5506b8ab4747884475075ebc38da", "key_hash": "main-silent-disconnect", "billing_key_hash": "main-silent-disconnect", "reserved_msats": 13, "status": "released", "created_at": 1790767961, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "1dfedc1c2cee407fb2cd1a28b1253c6d", "key_hash": "main-header", "billing_key_hash": "main-header", "reserved_msats": 12, "status": "released", "created_at": 1790767963, "started_at": 1790767960, "expires_at": 1790767978}]} +UPSTREAM_EVENTS [{"event": "start", "mode": "flood", "time": 1790767717.2871263}, {"event": "start", "mode": "silent-disconnect", "time": 1790767717.2960703}, {"event": "start", "mode": "keepalive", "time": 1790767717.3026786}, {"event": "start", "mode": "header", "time": 1790767717.3104746}, {"event": "start", "mode": "endless-disconnect", "time": 1790767717.318341}, {"event": "start", "mode": "silent", "time": 1790767717.3653235}, {"event": "start", "mode": "finite", "time": 1790767717.3924189}, {"event": "close", "mode": "finite", "chunks": 3, "time": 1790767718.397855}, {"event": "start", "mode": "flood", "time": 1790767960.9594278}, {"event": "start", "mode": "endless-disconnect", "time": 1790767961.008998}, {"event": "start", "mode": "keepalive", "time": 1790767961.036595}, {"event": "start", "mode": "finite", "time": 1790767961.060875}, {"event": "start", "mode": "silent", "time": 1790767961.0869172}, {"event": "start", "mode": "silent-disconnect", "time": 1790767961.1580715}, {"event": "start", "mode": "header", "time": 1790767961.1944675}, {"event": "close", "mode": "finite", "chunks": 3, "time": 1790767962.064207}] +AFTER_ALL_CLIENTS_CLOSED {"time": 1790768044.521246, "keys": [{"hashed_key": "main-finite", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-silent", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-silent-disconnect", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-endless-disconnect", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-keepalive", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-flood", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-header", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}], "rows": [{"id": "9b5b9063f6f045a290611e4144c897de", "key_hash": "main-flood", "billing_key_hash": "main-flood", "reserved_msats": 12, "status": "released", "created_at": 1790767962, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "efc8563301ba458890ee3ab2335e4e49", "key_hash": "main-endless-disconnect", "billing_key_hash": "main-endless-disconnect", "reserved_msats": 13, "status": "released", "created_at": 1790767960, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "cbed89b7ec444fea9789cf53b3e0f476", "key_hash": "main-keepalive", "billing_key_hash": "main-keepalive", "reserved_msats": 12, "status": "released", "created_at": 1790767975, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "1c56b25265df4743b4cff70dc57544c6", "key_hash": "main-finite", "billing_key_hash": "main-finite", "reserved_msats": 12, "status": "charged", "created_at": 1790767960, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "cc7f113eef004b7ba27bc761c5d9b9a1", "key_hash": "main-silent", "billing_key_hash": "main-silent", "reserved_msats": 12, "status": "charged", "created_at": 1790767963, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "6a1b5506b8ab4747884475075ebc38da", "key_hash": "main-silent-disconnect", "billing_key_hash": "main-silent-disconnect", "reserved_msats": 13, "status": "released", "created_at": 1790767961, "started_at": 1790767960, "expires_at": 1790767978}, {"id": "1dfedc1c2cee407fb2cd1a28b1253c6d", "key_hash": "main-header", "billing_key_hash": "main-header", "reserved_msats": 12, "status": "released", "created_at": 1790767963, "started_at": 1790767960, "expires_at": 1790767978}]} +UPSTREAM_EVENTS [{"event": "start", "mode": "flood", "time": 1790767717.2871263}, {"event": "start", "mode": "silent-disconnect", "time": 1790767717.2960703}, {"event": "start", "mode": "keepalive", "time": 1790767717.3026786}, {"event": "start", "mode": "header", "time": 1790767717.3104746}, {"event": "start", "mode": "endless-disconnect", "time": 1790767717.318341}, {"event": "start", "mode": "silent", "time": 1790767717.3653235}, {"event": "start", "mode": "finite", "time": 1790767717.3924189}, {"event": "close", "mode": "finite", "chunks": 3, "time": 1790767718.397855}, {"event": "start", "mode": "flood", "time": 1790767960.9594278}, {"event": "start", "mode": "endless-disconnect", "time": 1790767961.008998}, {"event": "start", "mode": "keepalive", "time": 1790767961.036595}, {"event": "start", "mode": "finite", "time": 1790767961.060875}, {"event": "start", "mode": "silent", "time": 1790767961.0869172}, {"event": "start", "mode": "silent-disconnect", "time": 1790767961.1580715}, {"event": "start", "mode": "header", "time": 1790767961.1944675}, {"event": "close", "mode": "finite", "chunks": 3, "time": 1790767962.064207}] diff --git a/repro/results.txt b/repro/results.txt new file mode 100644 index 00000000..01417f00 --- /dev/null +++ b/repro/results.txt @@ -0,0 +1,19 @@ +STREAM silent-disconnect 200 +STREAM keepalive 200 +STREAM endless-disconnect 200 +STREAM silent 200 +STREAM finite 200 +CLIENT_DISCONNECTED silent-disconnect +CLIENT_DISCONNECTED endless-disconnect +ENDED finite +STREAM header 424 +ENDED header +ENDED silent +AT_10_SECONDS {"time": 1790767727.5333533, "keys": [{"hashed_key": "main-finite", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-silent", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-silent-disconnect", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-endless-disconnect", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-keepalive", "balance": 1000000000, "reserved_balance": 12, "reserved_at": 1790767717}, {"hashed_key": "main-flood", "balance": 999893370, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-header", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}], "rows": [{"id": "79c7feb2274140748da2a97180f56d2c", "key_hash": "main-flood", "billing_key_hash": "main-flood", "reserved_msats": 12, "status": "charged", "created_at": 1790767719, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "ecf20c01870b4ce49fc81bd300ee35df", "key_hash": "main-silent-disconnect", "billing_key_hash": "main-silent-disconnect", "reserved_msats": 13, "status": "charged", "created_at": 1790767717, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "ec6f854d6a8647ac8b3bba50752bc647", "key_hash": "main-keepalive", "billing_key_hash": "main-keepalive", "reserved_msats": 12, "status": "active", "created_at": 1790767725, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "6da67d4c5cbd4ab19cc30f2f0fac6aaa", "key_hash": "main-header", "billing_key_hash": "main-header", "reserved_msats": 12, "status": "released", "created_at": 1790767719, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "f4a35c5e89c343bfbe21415dce28a4d5", "key_hash": "main-endless-disconnect", "billing_key_hash": "main-endless-disconnect", "reserved_msats": 13, "status": "released", "created_at": 1790767717, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "50a37cffba7745fa84d03b4070c86066", "key_hash": "main-silent", "billing_key_hash": "main-silent", "reserved_msats": 12, "status": "charged", "created_at": 1790767719, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "63ca23fa71024428baaf8baeb7ee3ede", "key_hash": "main-finite", "billing_key_hash": "main-finite", "reserved_msats": 12, "status": "charged", "created_at": 1790767717, "started_at": 1790767717, "expires_at": 1790767735}]} +REFUND keepalive 400 {"detail":"Cannot refund key. There are ongoing requests for this api key.","request_id":"f0bcc404-dffe-4860-8091-308f721ba053"} +UPSTREAM_EVENTS [{"event": "start", "mode": "flood", "time": 1790767717.2871263}, {"event": "start", "mode": "silent-disconnect", "time": 1790767717.2960703}, {"event": "start", "mode": "keepalive", "time": 1790767717.3026786}, {"event": "start", "mode": "header", "time": 1790767717.3104746}, {"event": "start", "mode": "endless-disconnect", "time": 1790767717.318341}, {"event": "start", "mode": "silent", "time": 1790767717.3653235}, {"event": "start", "mode": "finite", "time": 1790767717.3924189}, {"event": "close", "mode": "finite", "chunks": 3, "time": 1790767718.397855}] +CLIENT_ERROR keepalive RemoteProtocolError peer closed connection without sending complete message body (incomplete chunked read) +AFTER_SWEEP {"time": 1790767790.5300848, "keys": [{"hashed_key": "main-finite", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-silent", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-silent-disconnect", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-endless-disconnect", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-keepalive", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-flood", "balance": 999893370, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-header", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}], "rows": [{"id": "79c7feb2274140748da2a97180f56d2c", "key_hash": "main-flood", "billing_key_hash": "main-flood", "reserved_msats": 12, "status": "charged", "created_at": 1790767719, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "ecf20c01870b4ce49fc81bd300ee35df", "key_hash": "main-silent-disconnect", "billing_key_hash": "main-silent-disconnect", "reserved_msats": 13, "status": "charged", "created_at": 1790767717, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "ec6f854d6a8647ac8b3bba50752bc647", "key_hash": "main-keepalive", "billing_key_hash": "main-keepalive", "reserved_msats": 12, "status": "released", "created_at": 1790767731, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "6da67d4c5cbd4ab19cc30f2f0fac6aaa", "key_hash": "main-header", "billing_key_hash": "main-header", "reserved_msats": 12, "status": "released", "created_at": 1790767719, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "f4a35c5e89c343bfbe21415dce28a4d5", "key_hash": "main-endless-disconnect", "billing_key_hash": "main-endless-disconnect", "reserved_msats": 13, "status": "released", "created_at": 1790767717, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "50a37cffba7745fa84d03b4070c86066", "key_hash": "main-silent", "billing_key_hash": "main-silent", "reserved_msats": 12, "status": "charged", "created_at": 1790767719, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "63ca23fa71024428baaf8baeb7ee3ede", "key_hash": "main-finite", "billing_key_hash": "main-finite", "reserved_msats": 12, "status": "charged", "created_at": 1790767717, "started_at": 1790767717, "expires_at": 1790767735}]} +UPSTREAM_EVENTS [{"event": "start", "mode": "flood", "time": 1790767717.2871263}, {"event": "start", "mode": "silent-disconnect", "time": 1790767717.2960703}, {"event": "start", "mode": "keepalive", "time": 1790767717.3026786}, {"event": "start", "mode": "header", "time": 1790767717.3104746}, {"event": "start", "mode": "endless-disconnect", "time": 1790767717.318341}, {"event": "start", "mode": "silent", "time": 1790767717.3653235}, {"event": "start", "mode": "finite", "time": 1790767717.3924189}, {"event": "close", "mode": "finite", "chunks": 3, "time": 1790767718.397855}] +AFTER_ALL_CLIENTS_CLOSED {"time": 1790767800.920076, "keys": [{"hashed_key": "main-finite", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-silent", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-silent-disconnect", "balance": 999999997, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-endless-disconnect", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-keepalive", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-flood", "balance": 999893370, "reserved_balance": 0, "reserved_at": null}, {"hashed_key": "main-header", "balance": 1000000000, "reserved_balance": 0, "reserved_at": null}], "rows": [{"id": "79c7feb2274140748da2a97180f56d2c", "key_hash": "main-flood", "billing_key_hash": "main-flood", "reserved_msats": 12, "status": "charged", "created_at": 1790767719, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "ecf20c01870b4ce49fc81bd300ee35df", "key_hash": "main-silent-disconnect", "billing_key_hash": "main-silent-disconnect", "reserved_msats": 13, "status": "charged", "created_at": 1790767717, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "ec6f854d6a8647ac8b3bba50752bc647", "key_hash": "main-keepalive", "billing_key_hash": "main-keepalive", "reserved_msats": 12, "status": "released", "created_at": 1790767731, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "6da67d4c5cbd4ab19cc30f2f0fac6aaa", "key_hash": "main-header", "billing_key_hash": "main-header", "reserved_msats": 12, "status": "released", "created_at": 1790767719, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "f4a35c5e89c343bfbe21415dce28a4d5", "key_hash": "main-endless-disconnect", "billing_key_hash": "main-endless-disconnect", "reserved_msats": 13, "status": "released", "created_at": 1790767717, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "50a37cffba7745fa84d03b4070c86066", "key_hash": "main-silent", "billing_key_hash": "main-silent", "reserved_msats": 12, "status": "charged", "created_at": 1790767719, "started_at": 1790767717, "expires_at": 1790767735}, {"id": "63ca23fa71024428baaf8baeb7ee3ede", "key_hash": "main-finite", "billing_key_hash": "main-finite", "reserved_msats": 12, "status": "charged", "created_at": 1790767717, "started_at": 1790767717, "expires_at": 1790767735}]} +UPSTREAM_EVENTS [{"event": "start", "mode": "flood", "time": 1790767717.2871263}, {"event": "start", "mode": "silent-disconnect", "time": 1790767717.2960703}, {"event": "start", "mode": "keepalive", "time": 1790767717.3026786}, {"event": "start", "mode": "header", "time": 1790767717.3104746}, {"event": "start", "mode": "endless-disconnect", "time": 1790767717.318341}, {"event": "start", "mode": "silent", "time": 1790767717.3653235}, {"event": "start", "mode": "finite", "time": 1790767717.3924189}, {"event": "close", "mode": "finite", "chunks": 3, "time": 1790767718.397855}] diff --git a/repro/router-final.log b/repro/router-final.log new file mode 100644 index 00000000..dee5fe67 --- /dev/null +++ b/repro/router-final.log @@ -0,0 +1,108 @@ +/.venv/lib/python3.14/site-packages/anyio/from_thread.py:119: SyntaxWarning: 'return' in a 'finally' block + return result +2026-09-30 11:32:06 WARNING routstr.core.main UI dist directory not found at /app/ui_out; serving API only. Run `make ui-build` to build the static UI served from here, or `make ui-dev` for the Next.js dev server with hot reload on :3000 (it targets this backend on :8000). +2026-09-30 11:32:06 INFO uvicorn.error Started server process [1] +2026-09-30 11:32:06 INFO uvicorn.error Waiting for application startup. +2026-09-30 11:32:06 INFO routstr.core.main Application startup initiated +2026-09-30 11:32:10 INFO routstr.core.db Database migrations completed successfully +2026-09-30 11:32:10 INFO routstr.core.db Reset reserved balances on startup +2026-09-30 11:32:11 INFO routstr.upstream.helpers Seeding custom provider +2026-09-30 11:32:11 INFO routstr.upstream.helpers Seeded 1 upstream providers from settings +2026-09-30 11:32:12 INFO routstr.proxy Initialized 1 upstream providers +2026-09-30 11:32:12 INFO routstr.nostr.listing Nostr private key not configured (NSEC); waiting for one to be set before announcing this provider +2026-09-30 11:32:12 INFO routstr.nostr.analytics Usage analytics sharing task started +2026-09-30 11:32:12 INFO routstr.nostr.analytics NSEC is not configured; skipping analytics sharing to Nostr +2026-09-30 11:32:12 INFO routstr.auth Dead-key pruning disabled (interval <= 0) +2026-09-30 11:32:12 INFO uvicorn.error Application startup complete. +2026-09-30 11:32:12 INFO uvicorn.error Uvicorn running on http://127.0.0.1:18100 (Press CTRL+C to quit) +2026-09-30 11:32:40 INFO routstr.auth Existing sk- API key found +2026-09-30 11:32:40 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:32:40 INFO routstr.auth Processing payment for request +2026-09-30 11:32:40 INFO routstr.auth Existing sk- API key found +2026-09-30 11:32:40 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:32:40 INFO routstr.auth Processing payment for request +2026-09-30 11:32:40 INFO routstr.auth Existing sk- API key found +2026-09-30 11:32:40 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:32:40 INFO routstr.auth Processing payment for request +2026-09-30 11:32:40 INFO routstr.auth Existing sk- API key found +2026-09-30 11:32:40 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:32:40 INFO routstr.auth Processing payment for request +2026-09-30 11:32:40 INFO routstr.auth Existing sk- API key found +2026-09-30 11:32:40 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:32:40 INFO routstr.auth Processing payment for request +2026-09-30 11:32:40 INFO routstr.auth Existing sk- API key found +2026-09-30 11:32:40 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:32:40 INFO routstr.auth Processing payment for request +2026-09-30 11:32:40 INFO routstr.auth Existing sk- API key found +2026-09-30 11:32:40 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:32:40 INFO routstr.auth Processing payment for request +2026-09-30 11:32:40 INFO routstr.auth Payment processed successfully +2026-09-30 11:32:40 INFO routstr.payments RESERVE +2026-09-30 11:32:40 INFO routstr.auth Payment processed successfully +2026-09-30 11:32:40 INFO routstr.payments RESERVE +2026-09-30 11:32:41 INFO routstr.auth Payment processed successfully +2026-09-30 11:32:41 INFO routstr.payments RESERVE +2026-09-30 11:32:41 INFO routstr.auth Payment processed successfully +2026-09-30 11:32:41 INFO routstr.payments RESERVE +2026-09-30 11:32:41 INFO routstr.auth Payment processed successfully +2026-09-30 11:32:41 INFO routstr.payments RESERVE +2026-09-30 11:32:41 INFO routstr.auth Payment processed successfully +2026-09-30 11:32:41 INFO routstr.payments RESERVE +2026-09-30 11:32:41 INFO routstr.auth Payment processed successfully +2026-09-30 11:32:41 INFO routstr.payments RESERVE +2026-09-30 11:32:41 INFO routstr.payment.cost_calculation Applied model-specific pricing +2026-09-30 11:32:41 INFO routstr.payment.cost_calculation Calculated token-based cost +2026-09-30 11:32:41 INFO routstr.payment.cost_calculation Applied model-specific pricing +2026-09-30 11:32:41 INFO routstr.payment.cost_calculation Calculated token-based cost +2026-09-30 11:32:41 INFO routstr.auth Payment settlement finished +2026-09-30 11:32:41 INFO routstr.auth Payment settlement finished +2026-09-30 11:32:42 INFO routstr.payment.cost_calculation Applied model-specific pricing +2026-09-30 11:32:42 INFO routstr.payment.cost_calculation Calculated token-based cost +2026-09-30 11:32:42 INFO routstr.auth Calculated token-based cost +2026-09-30 11:32:42 INFO routstr.auth Refunding excess payment +2026-09-30 11:32:42 INFO routstr.auth Refund processed successfully +2026-09-30 11:32:42 INFO routstr.payments FINALIZE +2026-09-30 11:32:42 INFO routstr.auth Payment settlement finished +2026-09-30 11:32:42 INFO routstr.upstream.auto_topup Auto top-up worker started +2026-09-30 11:32:43 INFO routstr.payment.cost_calculation Applied model-specific pricing +2026-09-30 11:32:43 INFO routstr.payment.cost_calculation Calculated token-based cost +2026-09-30 11:32:43 ERROR routstr.core.exceptions Unhandled exception +asyncio.exceptions.CancelledError + +The above exception was the direct cause of the following exception: + +TimeoutError +2026-09-30 11:32:43 ERROR uvicorn.error Exception in ASGI application +asyncio.exceptions.CancelledError + +The above exception was the direct cause of the following exception: + +TimeoutError +2026-09-30 11:32:43 INFO routstr.auth Payment settlement finished +2026-09-30 11:32:44 WARNING routstr.upstream.base Streaming interrupted; finalizing before closing upstream +2026-09-30 11:32:44 INFO routstr.payment.cost_calculation Applied model-specific pricing +2026-09-30 11:32:44 INFO routstr.payment.cost_calculation Calculated token-based cost +2026-09-30 11:32:44 INFO routstr.auth Calculated token-based cost +2026-09-30 11:32:44 INFO routstr.auth Refunding excess payment +2026-09-30 11:32:44 INFO routstr.auth Refund processed successfully +2026-09-30 11:32:44 INFO routstr.payments FINALIZE +2026-09-30 11:32:44 INFO routstr.auth Payment settlement finished +2026-09-30 11:32:44 ERROR routstr.core.exceptions Unhandled exception +httpcore.ReadTimeout + +The above exception was the direct cause of the following exception: + +httpx.ReadTimeout +2026-09-30 11:32:44 ERROR uvicorn.error Exception in ASGI application +httpcore.ReadTimeout + +The above exception was the direct cause of the following exception: + +httpx.ReadTimeout +2026-09-30 11:32:44 ERROR routstr.upstream.base HTTP request error to upstream +2026-09-30 11:32:44 WARNING routstr.proxy Upstream base failed for model=gpt-4o-mini: Upstream service request timed out +2026-09-30 11:32:52 INFO routstr.core.exceptions HTTP 400 on /v1/wallet/refund: Cannot refund key. There are ongoing requests for this api key. +2026-09-30 11:32:55 INFO routstr.payment.cost_calculation Applied model-specific pricing +2026-09-30 11:32:55 INFO routstr.payment.cost_calculation Calculated token-based cost +2026-09-30 11:32:55 ERROR uvicorn.error ASGI callable returned without completing response. +2026-09-30 11:32:55 INFO routstr.auth Payment settlement finished diff --git a/repro/router-first.log b/repro/router-first.log new file mode 100644 index 00000000..b8979fc0 --- /dev/null +++ b/repro/router-first.log @@ -0,0 +1,115 @@ +/.venv/lib/python3.14/site-packages/anyio/from_thread.py:119: SyntaxWarning: 'return' in a 'finally' block + return result +2026-09-30 11:28:16 WARNING routstr.core.main UI dist directory not found at /app/ui_out; serving API only. Run `make ui-build` to build the static UI served from here, or `make ui-dev` for the Next.js dev server with hot reload on :3000 (it targets this backend on :8000). +2026-09-30 11:28:16 INFO uvicorn.error Started server process [1] +2026-09-30 11:28:16 INFO uvicorn.error Waiting for application startup. +2026-09-30 11:28:16 INFO routstr.core.main Application startup initiated +2026-09-30 11:28:20 INFO routstr.core.db Database migrations completed successfully +2026-09-30 11:28:21 INFO routstr.core.db Reset reserved balances on startup +2026-09-30 11:28:21 INFO routstr.upstream.helpers Seeding custom provider +2026-09-30 11:28:21 INFO routstr.upstream.helpers Seeded 1 upstream providers from settings +2026-09-30 11:28:22 INFO routstr.proxy Initialized 1 upstream providers +2026-09-30 11:28:22 INFO routstr.nostr.listing Nostr private key not configured (NSEC); waiting for one to be set before announcing this provider +2026-09-30 11:28:22 INFO routstr.nostr.analytics Usage analytics sharing task started +2026-09-30 11:28:22 INFO routstr.nostr.analytics NSEC is not configured; skipping analytics sharing to Nostr +2026-09-30 11:28:22 INFO routstr.auth Dead-key pruning disabled (interval <= 0) +2026-09-30 11:28:22 INFO uvicorn.error Application startup complete. +2026-09-30 11:28:22 INFO uvicorn.error Uvicorn running on http://127.0.0.1:18100 (Press CTRL+C to quit) +2026-09-30 11:28:37 INFO routstr.auth Existing sk- API key found +2026-09-30 11:28:37 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:28:37 INFO routstr.auth Processing payment for request +2026-09-30 11:28:37 INFO routstr.auth Existing sk- API key found +2026-09-30 11:28:37 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:28:37 INFO routstr.auth Processing payment for request +2026-09-30 11:28:37 INFO routstr.auth Existing sk- API key found +2026-09-30 11:28:37 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:28:37 INFO routstr.auth Processing payment for request +2026-09-30 11:28:37 INFO routstr.auth Existing sk- API key found +2026-09-30 11:28:37 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:28:37 INFO routstr.auth Processing payment for request +2026-09-30 11:28:37 INFO routstr.auth Payment processed successfully +2026-09-30 11:28:37 INFO routstr.payments RESERVE +2026-09-30 11:28:37 INFO routstr.auth Existing sk- API key found +2026-09-30 11:28:37 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:28:37 INFO routstr.auth Processing payment for request +2026-09-30 11:28:37 INFO routstr.auth Existing sk- API key found +2026-09-30 11:28:37 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:28:37 INFO routstr.auth Processing payment for request +2026-09-30 11:28:37 INFO routstr.auth Existing sk- API key found +2026-09-30 11:28:37 INFO routstr.proxy Bearer token validated successfully +2026-09-30 11:28:37 INFO routstr.auth Processing payment for request +2026-09-30 11:28:37 INFO routstr.auth Payment processed successfully +2026-09-30 11:28:37 INFO routstr.payments RESERVE +2026-09-30 11:28:37 INFO routstr.auth Payment processed successfully +2026-09-30 11:28:37 INFO routstr.payments RESERVE +2026-09-30 11:28:37 INFO routstr.auth Payment processed successfully +2026-09-30 11:28:37 INFO routstr.payments RESERVE +2026-09-30 11:28:37 INFO routstr.auth Payment processed successfully +2026-09-30 11:28:37 INFO routstr.payments RESERVE +2026-09-30 11:28:37 INFO routstr.auth Payment processed successfully +2026-09-30 11:28:37 INFO routstr.payments RESERVE +2026-09-30 11:28:37 INFO routstr.auth Payment processed successfully +2026-09-30 11:28:37 INFO routstr.payments RESERVE +2026-09-30 11:28:38 INFO routstr.payment.cost_calculation Applied model-specific pricing +2026-09-30 11:28:38 INFO routstr.payment.cost_calculation Calculated token-based cost +2026-09-30 11:28:38 INFO routstr.payment.cost_calculation Applied model-specific pricing +2026-09-30 11:28:38 INFO routstr.payment.cost_calculation Calculated token-based cost +2026-09-30 11:28:38 INFO routstr.auth Payment settlement finished +2026-09-30 11:28:38 INFO routstr.auth Calculated token-based cost +2026-09-30 11:28:38 INFO routstr.auth Refunding excess payment +2026-09-30 11:28:38 INFO routstr.auth Refund processed successfully +2026-09-30 11:28:38 INFO routstr.payments FINALIZE +2026-09-30 11:28:38 INFO routstr.auth Payment settlement finished +2026-09-30 11:28:38 INFO routstr.payment.cost_calculation Applied model-specific pricing +2026-09-30 11:28:38 INFO routstr.payment.cost_calculation Calculated token-based cost +2026-09-30 11:28:38 INFO routstr.auth Calculated token-based cost +2026-09-30 11:28:38 INFO routstr.auth Refunding excess payment +2026-09-30 11:28:38 INFO routstr.auth Refund processed successfully +2026-09-30 11:28:38 INFO routstr.payments FINALIZE +2026-09-30 11:28:38 INFO routstr.auth Payment settlement finished +2026-09-30 11:28:39 INFO routstr.payment.cost_calculation Applied model-specific pricing +2026-09-30 11:28:39 INFO routstr.payment.cost_calculation Calculated token-based cost +2026-09-30 11:28:39 INFO routstr.auth Calculated token-based cost +2026-09-30 11:28:39 INFO routstr.auth Finalized payment with additional charge +2026-09-30 11:28:39 INFO routstr.payments FINALIZE +2026-09-30 11:28:39 INFO routstr.auth Payment settlement finished +2026-09-30 11:28:39 ERROR routstr.core.exceptions Unhandled exception +asyncio.exceptions.CancelledError + +The above exception was the direct cause of the following exception: + +TimeoutError +2026-09-30 11:28:39 ERROR uvicorn.error Exception in ASGI application +asyncio.exceptions.CancelledError + +The above exception was the direct cause of the following exception: + +TimeoutError +2026-09-30 11:28:40 ERROR routstr.upstream.base HTTP request error to upstream +2026-09-30 11:28:40 WARNING routstr.proxy Upstream base failed for model=gpt-4o-mini: Upstream service request timed out +2026-09-30 11:28:40 WARNING routstr.upstream.base Streaming interrupted; finalizing before closing upstream +2026-09-30 11:28:40 INFO routstr.payment.cost_calculation Applied model-specific pricing +2026-09-30 11:28:40 INFO routstr.payment.cost_calculation Calculated token-based cost +2026-09-30 11:28:40 INFO routstr.auth Calculated token-based cost +2026-09-30 11:28:40 INFO routstr.auth Refunding excess payment +2026-09-30 11:28:40 INFO routstr.auth Refund processed successfully +2026-09-30 11:28:40 INFO routstr.payments FINALIZE +2026-09-30 11:28:40 INFO routstr.auth Payment settlement finished +2026-09-30 11:28:40 ERROR routstr.core.exceptions Unhandled exception +httpcore.ReadTimeout + +The above exception was the direct cause of the following exception: + +httpx.ReadTimeout +2026-09-30 11:28:40 ERROR uvicorn.error Exception in ASGI application +httpcore.ReadTimeout + +The above exception was the direct cause of the following exception: + +httpx.ReadTimeout +2026-09-30 11:28:49 INFO routstr.core.exceptions HTTP 400 on /v1/wallet/refund: Cannot refund key. There are ongoing requests for this api key. +2026-09-30 11:28:52 INFO routstr.payment.cost_calculation Applied model-specific pricing +2026-09-30 11:28:52 INFO routstr.payment.cost_calculation Calculated token-based cost +2026-09-30 11:28:52 ERROR uvicorn.error ASGI callable returned without completing response. +2026-09-30 11:28:52 INFO routstr.auth Payment settlement finished +2026-09-30 11:28:52 INFO routstr.upstream.auto_topup Auto top-up worker started diff --git a/routstr/auth.py b/routstr/auth.py index 9bfc7b3d..b7c8e16a 100644 --- a/routstr/auth.py +++ b/routstr/auth.py @@ -637,6 +637,14 @@ async def pay_for_request( ) # Charge the base cost for the request atomically to avoid race conditions + from .core.lifecycle import request_lifetime + + lifetime = request_lifetime.get() + remaining_lifetime = ( + max(0, lifetime.deadline - asyncio.get_running_loop().time()) + if lifetime is not None + else settings.max_request_lifetime_seconds + ) reserved_at_now = int(time.time()) stmt = ( update(ApiKey) @@ -686,6 +694,9 @@ async def pay_for_request( billing_key_hash=reservation.billing_key_hash, reserved_msats=reservation.reserved_msats, status="active", + started_at=reserved_at_now, + expires_at=reserved_at_now + + math.ceil(remaining_lifetime + settings.request_cleanup_timeout_seconds), ) ) # Publish the identity before commit. If the commit succeeds but its @@ -726,6 +737,11 @@ async def pay_for_request( # The reservation is durable; keep its lease fresh for the whole request # lifetime (upstream header waits, non-streaming and streaming alike). + from .core.lifecycle import request_lifetime + + lifetime = request_lifetime.get() + if lifetime is not None: + lifetime.reservations.append(reservation) _start_reservation_heartbeat(reservation) try: @@ -875,6 +891,10 @@ async def renew_reservation( update(ReservationRelease) .where(col(ReservationRelease.id) == snapshot.release_id) .where(col(ReservationRelease.status) == "active") + .where( + (col(ReservationRelease.expires_at).is_(None)) + | (col(ReservationRelease.expires_at) > int(time.time())) + ) .values(created_at=int(time.time())) ) await session.commit() @@ -902,12 +922,21 @@ def _start_reservation_heartbeat(snapshot: ReservationSnapshot) -> None: """ interval = max(1, settings.stale_reservation_timeout_seconds // 3) owner = asyncio.current_task() + from .core.lifecycle import request_lifetime + + lifetime = request_lifetime.get() + deadline = asyncio.get_running_loop().time() + settings.max_request_lifetime_seconds async def beat() -> None: try: while True: await asyncio.sleep(interval) - if owner is None or owner.done(): + if ( + owner is None + or owner.done() + or (lifetime is not None and lifetime.stopped) + or asyncio.get_running_loop().time() >= deadline + ): # Request control is gone; let the lease expire so the # sweeper can release the reservation if no terminal # transition ever ran. @@ -1091,6 +1120,10 @@ async def _claim_reservation_for_charge( update(ReservationRelease) .where(col(ReservationRelease.id) == snapshot.release_id) .where(col(ReservationRelease.status) == "active") + .where( + col(ReservationRelease.expires_at).is_(None) + | (col(ReservationRelease.expires_at) > int(time.time())) + ) .where(col(ReservationRelease.key_hash) == snapshot.key_hash) .where(col(ReservationRelease.billing_key_hash) == snapshot.billing_key_hash) .where(col(ReservationRelease.reserved_msats) == snapshot.reserved_msats) diff --git a/routstr/core/db.py b/routstr/core/db.py index 500420d5..84307a7e 100644 --- a/routstr/core/db.py +++ b/routstr/core/db.py @@ -174,7 +174,10 @@ async def _transition_stale_reservation( update(ReservationRelease) .where(col(ReservationRelease.id) == reservation_id) .where(col(ReservationRelease.status) == "active") - .where(col(ReservationRelease.created_at) < cutoff) + .where( + (col(ReservationRelease.created_at) < cutoff) + | (col(ReservationRelease.expires_at) <= int(time.time())) + ) .values(status="released") ) return bool(transition.rowcount == 1) @@ -221,7 +224,10 @@ async def release_stale_reservations( query = ( select(ReservationRelease) .where(col(ReservationRelease.status) == "active") - .where(col(ReservationRelease.created_at) < cutoff) + .where( + (col(ReservationRelease.created_at) < cutoff) + | (col(ReservationRelease.expires_at) <= int(time.time())) + ) ) if key_hash is not None: query = query.where( @@ -784,6 +790,8 @@ class ReservationRelease(SQLModel, table=True): # type: ignore key_hash: str = Field(index=True) billing_key_hash: str = Field(index=True) reserved_msats: int + started_at: int | None = Field(default=None) + expires_at: int | None = Field(default=None, index=True) status: str = Field(default="active") created_at: int = Field(default_factory=lambda: int(time.time())) diff --git a/routstr/core/lifecycle.py b/routstr/core/lifecycle.py new file mode 100644 index 00000000..897df5b7 --- /dev/null +++ b/routstr/core/lifecycle.py @@ -0,0 +1,135 @@ +"""Supervise the real downstream connection, outside HTTP middleware wrappers.""" + +from __future__ import annotations + +import asyncio +from contextvars import ContextVar +from dataclasses import dataclass, field +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from ..auth import ReservationSnapshot + +from starlette.types import ASGIApp, Message, Receive, Scope, Send + +from . import get_logger +from .settings import settings + +logger = get_logger(__name__) + + +@dataclass +class RequestLifetime: + deadline: float = 0 + stopped: bool = False + reservations: list[ReservationSnapshot] = field(default_factory=list) + + +request_lifetime: ContextVar[RequestLifetime | None] = ContextVar( + "request_lifetime", default=None +) + + +class RequestLifecycleMiddleware: + def __init__(self, app: ASGIApp) -> None: + self.app = app + + async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None: + if scope["type"] != "http": + await self.app(scope, receive, send) + return + lifetime = RequestLifetime( + deadline=asyncio.get_running_loop().time() + + settings.max_request_lifetime_seconds + ) + token = request_lifetime.set(lifetime) + disconnected = asyncio.Event() + # One receive consumer. Backpressure uploads until consumed; after the + # final body message, continue listening independently of the app. + messages: asyncio.Queue[Message] = asyncio.Queue(maxsize=1) + response_started = False + + async def pump() -> None: + while True: + message = await receive() + if message["type"] == "http.disconnect": + disconnected.set() + return + await messages.put(message) + + async def downstream_receive() -> Message: + if disconnected.is_set(): + return {"type": "http.disconnect"} + get = asyncio.create_task(messages.get()) + gone = asyncio.create_task(disconnected.wait()) + try: + await asyncio.wait((get, gone), return_when=asyncio.FIRST_COMPLETED) + if disconnected.is_set(): + return {"type": "http.disconnect"} + return get.result() + finally: + for task in (get, gone): + task.cancel() + await asyncio.gather(get, gone, return_exceptions=True) + + async def downstream_send(message: Message) -> None: + nonlocal response_started + if disconnected.is_set() or lifetime.stopped: + raise OSError("Downstream request terminated") + async with asyncio.timeout(settings.downstream_send_timeout_seconds): + await send(message) + if message["type"] == "http.response.start": + response_started = True + + receiver = asyncio.create_task(pump()) + work = asyncio.create_task(self.app(scope, downstream_receive, downstream_send)) + gone = asyncio.create_task(disconnected.wait()) + try: + done, _ = await asyncio.wait( + (work, gone), + timeout=settings.max_request_lifetime_seconds, + return_when=asyncio.FIRST_COMPLETED, + ) + if work in done: + await work + elif not disconnected.is_set() and not response_started: + await downstream_send( + {"type": "http.response.start", "status": 504, "headers": []} + ) + await downstream_send( + {"type": "http.response.body", "body": b"Request deadline exceeded"} + ) + finally: + lifetime.stopped = True + for task in (receiver, gone, work): + task.cancel() + # Cancellation/close is bounded: an uncooperative finalizer must not + # hold ownership or renewal indefinitely. + done, pending = await asyncio.wait( + (receiver, gone, work), timeout=settings.request_cleanup_timeout_seconds + ) + for task in done: + if not task.cancelled(): + task.exception() + for task in pending: + task.cancel() + task.add_done_callback( + lambda t: t.exception() if not t.cancelled() else None + ) + try: + async with asyncio.timeout(settings.request_cleanup_timeout_seconds): + from ..auth import _stop_reservation_heartbeat, release_reservation + from .db import create_session + + for snapshot in lifetime.reservations: + await _stop_reservation_heartbeat(snapshot.release_id) + async with create_session() as session: + await release_reservation( + snapshot, session, snapshot.reserved_msats + ) + except Exception: + logger.exception( + "Request cleanup failed; durable expiry will recover reservations" + ) + finally: + request_lifetime.reset(token) diff --git a/routstr/core/main.py b/routstr/core/main.py index 584fa736..46aad1fc 100644 --- a/routstr/core/main.py +++ b/routstr/core/main.py @@ -45,6 +45,7 @@ from .exceptions import ( http_exception_handler, validation_exception_handler, ) +from .lifecycle import RequestLifecycleMiddleware from .logging import get_logger, setup_logging from .middleware import LoggingMiddleware from .not_found import _NOT_FOUND_HTML, not_found_catch_all # noqa: F401 @@ -315,6 +316,10 @@ app.add_middleware( # Add logging middleware app.add_middleware(LoggingMiddleware) +# Outermost: observe the actual downstream connection, not middleware streams. + +app.add_middleware(RequestLifecycleMiddleware) + # Add exception handlers app.add_exception_handler(HTTPException, http_exception_handler) # type: ignore app.add_exception_handler(RequestValidationError, validation_exception_handler) diff --git a/routstr/core/settings.py b/routstr/core/settings.py index 5dd2d57d..d9324f13 100644 --- a/routstr/core/settings.py +++ b/routstr/core/settings.py @@ -117,6 +117,16 @@ class Settings(BaseSettings): default=604_800, env="DEAD_KEY_MIN_AGE_SECONDS" ) + max_request_lifetime_seconds: float = Field( + default=1800, gt=0, env="MAX_REQUEST_LIFETIME_SECONDS" + ) + downstream_send_timeout_seconds: float = Field( + default=60, gt=0, env="DOWNSTREAM_SEND_TIMEOUT_SECONDS" + ) + request_cleanup_timeout_seconds: float = Field( + default=30, gt=0, env="REQUEST_CLEANUP_TIMEOUT_SECONDS" + ) + # Network cors_origins: list[str] = Field(default_factory=lambda: ["*"], env="CORS_ORIGINS") # Comma-separated METHOD:path pairs adding to the proxy's canonical diff --git a/routstr/upstream/stream_ownership.py b/routstr/upstream/stream_ownership.py index e0e3e51d..ea06e400 100644 --- a/routstr/upstream/stream_ownership.py +++ b/routstr/upstream/stream_ownership.py @@ -10,6 +10,7 @@ from fastapi.responses import StreamingResponse from starlette.types import Receive, Scope, Send from ..core import get_logger +from ..core.settings import settings logger = get_logger(__name__) @@ -74,10 +75,14 @@ class PersistentStreamFinalizer: self._task: asyncio.Future[None] | None = None self._lock = asyncio.Lock() + async def _bounded_finalize(self) -> None: + async with asyncio.timeout(settings.request_cleanup_timeout_seconds): + await self._finalize() + async def run(self) -> None: async with self._lock: if self._task is None: - self._task = asyncio.ensure_future(self._finalize()) + self._task = asyncio.ensure_future(self._bounded_finalize()) task = self._task await asyncio.shield(task) diff --git a/tests/unit/test_request_lifecycle.py b/tests/unit/test_request_lifecycle.py new file mode 100644 index 00000000..03cf9897 --- /dev/null +++ b/tests/unit/test_request_lifecycle.py @@ -0,0 +1,57 @@ +import asyncio +from unittest.mock import patch + +import pytest + +from routstr.core.lifecycle import RequestLifecycleMiddleware +from routstr.core.settings import settings + + +@pytest.mark.asyncio +@pytest.mark.parametrize("reason", ["disconnect", "deadline", "send"]) +async def test_lifecycle_stops_live_work(reason): + closed = asyncio.Event() + receive_queue = asyncio.Queue() + await receive_queue.put({"type": "http.request", "body": b"", "more_body": False}) + sent = [] + + async def app(scope, receive, send): + try: + assert (await receive())["type"] == "http.request" + await send({"type": "http.response.start", "status": 200, "headers": []}) + while True: + await send( + {"type": "http.response.body", "body": b"x", "more_body": True} + ) + await asyncio.sleep(0.01) + finally: + closed.set() + + async def send(message): + sent.append(message) + if reason == "send" and message["type"] == "http.response.body": + await asyncio.sleep(100) + + async def disconnect(): + await asyncio.sleep(0.02) + await receive_queue.put({"type": "http.disconnect"}) + + task = asyncio.create_task(disconnect()) if reason == "disconnect" else None + with ( + patch.object(settings, "max_request_lifetime_seconds", 0.08), + patch.object(settings, "downstream_send_timeout_seconds", 0.03), + patch.object(settings, "request_cleanup_timeout_seconds", 0.1), + ): + try: + await asyncio.wait_for( + RequestLifecycleMiddleware(app)( + {"type": "http"}, receive_queue.get, send + ), + 1, + ) + except TimeoutError: + assert reason == "send" + if task: + await task + assert closed.is_set() + assert sent diff --git a/tests/unit/test_stale_reservations.py b/tests/unit/test_stale_reservations.py index 558fdda5..6cd59682 100644 --- a/tests/unit/test_stale_reservations.py +++ b/tests/unit/test_stale_reservations.py @@ -428,3 +428,62 @@ async def test_proxy_reverts_reservation_on_client_disconnect() -> None: await proxy_module.proxy(request, "v1/chat/completions") revert_mock.assert_awaited_once_with(key, session, 1000, reservation_snapshot) + + +@pytest.mark.asyncio +async def test_absolute_expiry_releases_fresh_lease(session: AsyncSession) -> None: + now = int(time.time()) + key = ApiKey( + hashed_key="expired-deadline", + balance=5000, + reserved_balance=1000, + reserved_at=now, + ) + session.add(key) + session.add( + ReservationRelease( + id="expired", + key_hash=key.hashed_key, + billing_key_hash=key.hashed_key, + reserved_msats=1000, + created_at=now, + started_at=now - 100, + expires_at=now - 1, + ) + ) + await session.commit() + assert await release_stale_reservations(session, 300) == 1 + await session.refresh(key) + assert key.reserved_balance == 0 + assert key.balance == 5000 + + +@pytest.mark.asyncio +async def test_expired_reservation_cannot_renew_or_claim_charge( + session: AsyncSession, +) -> None: + from routstr.auth import ( + ReservationSnapshot, + _claim_reservation_for_charge, + renew_reservation, + ) + + snapshot = ReservationSnapshot( + release_id="fenced", + key_hash="fenced-key", + billing_key_hash="fenced-key", + reserved_msats=1000, + ) + session.add(ApiKey(hashed_key="fenced-key", balance=5000, reserved_balance=1000)) + session.add( + ReservationRelease( + id="fenced", + key_hash="fenced-key", + billing_key_hash="fenced-key", + reserved_msats=1000, + expires_at=int(time.time()) - 1, + ) + ) + await session.commit() + assert not await renew_reservation(snapshot, session) + assert not await _claim_reservation_for_charge(snapshot, session)