Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions src/free_claude_code/api/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@
get_request_id,
)
from .request_lifetime import ClientRequestLifetimeMiddleware
from .request_outcomes import RequestOutcomeMiddleware
from .routes import router
from .validation_log import summarize_request_validation_body

Expand All @@ -47,6 +48,7 @@ def create_app(services: ApiServices) -> FastAPI:
app.state.services = services
app.add_middleware(AdminNoStoreMiddleware)
app.add_middleware(ClientRequestLifetimeMiddleware)
app.add_middleware(RequestOutcomeMiddleware)
app.add_middleware(RequestCorrelationMiddleware)

app.include_router(admin_router)
Expand Down
5 changes: 5 additions & 0 deletions src/free_claude_code/api/handlers/messages.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
from free_claude_code.core.diagnostics import safe_exception_message
from free_claude_code.core.failures import ExecutionFailure, find_execution_failure
from free_claude_code.core.reasoning import ReasoningControl, ReasoningPolicy
from free_claude_code.core.request_outcomes import record_request_route
from free_claude_code.core.trace import trace_event

from .classifier_response import classifier_response
Expand Down Expand Up @@ -103,6 +104,10 @@ async def create(
require_non_empty_messages(request_data.messages)
routed = self._model_router.resolve_messages_request(request_data)
routed = self._apply_message_routing_policies(routed)
record_request_route(
routed.resolved.primary.provider_id,
routed.resolved.primary.provider_model,
)
tool_body = self._web_tools.try_stream_messages(
routed, request_id=request_id
)
Expand Down
5 changes: 5 additions & 0 deletions src/free_claude_code/api/handlers/responses.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
openai_error_type_for_failure,
openai_failure_payload,
)
from free_claude_code.core.request_outcomes import record_request_route


class ResponsesHandler:
Expand Down Expand Up @@ -77,6 +78,10 @@ async def create(

try:
routed = self._model_router.resolve_responses_request(request_data)
record_request_route(
routed.resolved.primary.provider_id,
routed.resolved.primary.provider_model,
)
streamed = self._provider_executor.stream_responses(
routed,
raw_log_payload=request_payload,
Expand Down
6 changes: 6 additions & 0 deletions src/free_claude_code/api/request_errors.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,10 @@
openai_error_payload,
openai_error_type_for_failure,
)
from free_claude_code.core.request_outcomes import (
record_request_exception,
record_request_failure,
)

WireApi = Literal["messages", "responses"]

Expand All @@ -37,6 +41,7 @@ def ordinary_application_error_response(
request_id: str,
) -> JSONResponse:
"""Serialize a deterministic application error without terminal headers."""
record_request_failure(error.kind.value)
if wire_api == "responses":
return JSONResponse(
status_code=error.status_code,
Expand Down Expand Up @@ -67,6 +72,7 @@ def log_unexpected_api_exception(
request_id: str | None = None,
) -> None:
"""Log API failures without echoing exception text unless opted in."""
record_request_exception(exc)
if settings.log_api_error_tracebacks:
if request_id is not None:
logger.error(
Expand Down
126 changes: 126 additions & 0 deletions src/free_claude_code/api/request_outcomes.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,126 @@
"""Log one final inference outcome after HTTP delivery and owned cleanup."""

import asyncio
import codecs
from time import monotonic

from loguru import logger
from starlette.datastructures import Headers
from starlette.types import ASGIApp, Message, Receive, Scope, Send

from free_claude_code.core.anthropic.stream_contracts import SSEEvent
from free_claude_code.core.anthropic.streaming.decoder import AnthropicSSEDecoder
from free_claude_code.core.request_outcomes import (
RequestOutcome,
current_request_outcome,
record_request_exception,
record_request_failure,
)

_FAILURE_EVENTS = frozenset({"error", "response.error", "response.failed"})


def _observe_event(event: SSEEvent) -> None:
kind = event.event or event.data.get("type")
if not isinstance(kind, str) or kind not in _FAILURE_EVENTS:
return
response = event.data.get("response")
payload = response if isinstance(response, dict) else event.data
error = payload.get("error")
error = error if isinstance(error, dict) else payload
reason = error.get("code") or error.get("type")
record_request_failure(reason if isinstance(reason, str) else str(kind))


class RequestOutcomeMiddleware:
"""Observe inference delivery without changing its body or control flow."""

def __init__(self, app: ASGIApp) -> None:
self._app = app

async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None:
if (
scope["type"] != "http"
or scope.get("method") != "POST"
or scope.get("path") not in {"/v1/messages", "/v1/responses"}
):
await self._app(scope, receive, send)
return

started = monotonic()
outcome = RequestOutcome()
status_code: int | None = None
completed = disconnected = cancelled = False
decoder: AnthropicSSEDecoder | None = None
text_decoder = codecs.getincrementaldecoder("utf-8")(errors="replace")

async def receive_observed() -> Message:
nonlocal disconnected
message = await receive()
if message["type"] == "http.disconnect":
disconnected = True
return message

async def send_observed(message: Message) -> None:
nonlocal status_code, completed, decoder, disconnected
try:
await send(message)
except OSError:
disconnected = True
raise
if message["type"] == "http.response.start":
status_code = message["status"]
if (
Headers(raw=message.get("headers", []))
.get("content-type", "")
.startswith("text/event-stream")
):
decoder = AnthropicSSEDecoder(event_names=_FAILURE_EVENTS)
elif message["type"] == "http.response.body":
completed = not message.get("more_body", False)
if decoder is not None:
text = text_decoder.decode(
message.get("body", b""), final=completed
)
for event in decoder.feed(text):
_observe_event(event)
if completed:
for event in decoder.finish():
_observe_event(event)

token = current_request_outcome.set(outcome)
try:
await self._app(scope, receive_observed, send_observed)
except asyncio.CancelledError:
cancelled = True
raise
except BaseException as error:
if not disconnected:
record_request_exception(error)
if status_code is None:
status_code = 500
raise
finally:
current_request_outcome.reset(token)
reason = outcome.failure_reason
if reason is not None or (status_code is not None and status_code >= 400):
result = "failure"
reason = reason or f"http_{status_code}"
elif (cancelled or disconnected) and not completed:
result = "cancelled"
elif not completed:
result, reason = "failure", "incomplete_response"
else:
result = "success"
logger.bind(
event="request.completed",
wire_api="responses"
if scope["path"] == "/v1/responses"
else "messages",
provider_id=outcome.provider_id,
model=outcome.model,
status_code=status_code,
outcome=result,
duration_ms=round((monotonic() - started) * 1000, 2),
failure_reason=reason,
).info("Inference request {}", result)
8 changes: 8 additions & 0 deletions src/free_claude_code/api/response_streams.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,10 @@
committed_response_failure_frame,
openai_error_type_for_failure,
)
from free_claude_code.core.request_outcomes import (
record_request_exception,
record_request_failure,
)
from free_claude_code.core.trace import close_stream_input, trace_event

TERMINAL_EXECUTION_ERROR_HEADERS = {"x-should-retry": "false"}
Expand Down Expand Up @@ -198,6 +202,10 @@ def trace_terminal_execution_error(
error: BaseException | None = None,
) -> None:
"""Record one correlated terminal-execution decision at the HTTP boundary."""
if error is not None:
record_request_exception(error)
else:
record_request_failure(error_type)
fields: dict[str, object] = {
"stage": "egress",
"event": "free_claude_code.api.response.terminal_execution_error",
Expand Down
2 changes: 2 additions & 0 deletions src/free_claude_code/application/execution.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
estimate_responses_input_tokens,
)
from free_claude_code.core.reasoning import ReasoningPolicy
from free_claude_code.core.request_outcomes import record_request_route
from free_claude_code.core.trace import (
close_stream_input,
trace_event,
Expand Down Expand Up @@ -341,6 +342,7 @@ async def provider_body() -> AsyncIterator[str]:
loop = asyncio.get_running_loop()
progress_deadline = loop.time() + self._progress_timeout_seconds
for index, target in enumerate(candidates):
record_request_route(target.provider_id, target.provider_model)
provider_stream: AsyncIterator[str] | None = None
candidate_committed = False
candidate_failure: ExecutionFailure | None = None
Expand Down
18 changes: 13 additions & 5 deletions src/free_claude_code/core/anthropic/stream_contracts.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,10 @@ class SSEEvent:
raw: str


def parse_sse_lines(lines: Iterable[str]) -> list[SSEEvent]:
def parse_sse_lines(
lines: Iterable[str], *, event_names: frozenset[str] | None = None
) -> list[SSEEvent]:
"""Decode selected named events and all unnamed events when a filter is supplied."""
events: list[SSEEvent] = []
current_event = ""
data_parts: list[str] = []
Expand All @@ -56,7 +59,7 @@ def parse_sse_lines(lines: Iterable[str]) -> list[SSEEvent]:
for line in lines:
stripped = line.rstrip("\r\n")
if stripped == "":
_append_event(events, current_event, data_parts, raw_parts)
_append_event(events, current_event, data_parts, raw_parts, event_names)
current_event = ""
data_parts = []
raw_parts = []
Expand All @@ -67,23 +70,28 @@ def parse_sse_lines(lines: Iterable[str]) -> list[SSEEvent]:
elif stripped.startswith("data:"):
data_parts.append(stripped.split(":", 1)[1].strip())

_append_event(events, current_event, data_parts, raw_parts)
_append_event(events, current_event, data_parts, raw_parts, event_names)
return events


def parse_sse_text(text: str) -> list[SSEEvent]:
def parse_sse_text(
text: str, *, event_names: frozenset[str] | None = None
) -> list[SSEEvent]:
# SSE uses CR/LF framing; Unicode line separators can occur inside JSON text.
return parse_sse_lines(re.split(r"\r\n|\r|\n", text))
return parse_sse_lines(re.split(r"\r\n|\r|\n", text), event_names=event_names)


def _append_event(
events: list[SSEEvent],
current_event: str,
data_parts: list[str],
raw_parts: list[str],
event_names: frozenset[str] | None,
) -> None:
if not current_event and not data_parts:
return
if event_names is not None and current_event and current_event not in event_names:
return
data_text = "\n".join(data_parts)
data: dict[str, Any]
try:
Expand Down
11 changes: 6 additions & 5 deletions src/free_claude_code/core/anthropic/streaming/decoder.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,14 +8,15 @@


class AnthropicSSEDecoder:
"""Decode arbitrarily split SSE text without losing frame order."""
"""Decode split SSE text, optionally filtering named events before JSON parsing."""

def __init__(self) -> None:
def __init__(self, *, event_names: frozenset[str] | None = None) -> None:
self._event_names = event_names
self._parts: list[str] = []
self._boundary_tail = ""

def feed(self, chunk: str) -> tuple[SSEEvent, ...]:
"""Consume one wire chunk and return every complete event."""
"""Consume one wire chunk and return its selected complete events."""

events: list[SSEEvent] = []
probe = self._boundary_tail + chunk
Expand All @@ -26,7 +27,7 @@ def feed(self, chunk: str) -> tuple[SSEEvent, ...]:
self._parts.append(chunk[chunk_start:chunk_end])
raw = "".join(self._parts)
self._parts.clear()
events.extend(parse_sse_text(raw))
events.extend(parse_sse_text(raw, event_names=self._event_names))
chunk_start = chunk_end

remainder = chunk[chunk_start:]
Expand All @@ -46,4 +47,4 @@ def finish(self) -> tuple[SSEEvent, ...]:
self._boundary_tail = ""
if not remainder.strip():
return ()
return tuple(parse_sse_text(remainder))
return tuple(parse_sse_text(remainder, event_names=self._event_names))
38 changes: 38 additions & 0 deletions src/free_claude_code/core/request_outcomes.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
"""Request-scoped diagnostic fields shared by routing and HTTP delivery."""

from contextvars import ContextVar
from dataclasses import dataclass

from .failures import find_execution_failure


@dataclass
class RequestOutcome:
"""Mutable fields shared with the request's response and cleanup tasks."""

provider_id: str | None = None
model: str | None = None
failure_reason: str | None = None


current_request_outcome: ContextVar[RequestOutcome | None] = ContextVar(
"request_outcome", default=None
)


def record_request_route(provider_id: str, model: str) -> None:
outcome = current_request_outcome.get()
if outcome is not None:
outcome.provider_id = provider_id
outcome.model = model


def record_request_failure(reason: str) -> None:
outcome = current_request_outcome.get()
if outcome is not None and outcome.failure_reason is None:
outcome.failure_reason = reason


def record_request_exception(error: BaseException) -> None:
failure = find_execution_failure(error)
record_request_failure(failure.kind.value if failure else type(error).__name__)
Loading
Loading