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: 1 addition & 1 deletion src/free_claude_code/api/handlers/messages.py
Original file line number Diff line number Diff line change
Expand Up @@ -121,7 +121,7 @@ async def create(
result = _MessagesStreamResult(
self._provider_executor.stream_messages(
routed,
raw_log_payload=routed.request.model_dump(),
raw_log_payload=routed.request.model_dump,
request_id=request_id,
)
)
Expand Down
5 changes: 3 additions & 2 deletions src/free_claude_code/api/handlers/responses.py
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,6 @@ async def create(
) -> object:
"""Create a streaming OpenAI Responses-compatible response."""
request_id = request_id or new_request_id()
request_payload = request_data.model_dump(mode="json", exclude_none=True)
if request_data.stream is False:
raise InvalidRequestError(
"FCC /v1/responses supports streaming only; omit stream or set stream=true."
Expand All @@ -84,7 +83,9 @@ async def create(
)
streamed = self._provider_executor.stream_responses(
routed,
raw_log_payload=request_payload,
raw_log_payload=lambda: request_data.model_dump(
mode="json", exclude_none=True
),
request_id=request_id,
)
return await openai_responses_sse_streaming_response(
Expand Down
9 changes: 6 additions & 3 deletions src/free_claude_code/api/handlers/token_count.py
Original file line number Diff line number Diff line change
Expand Up @@ -60,16 +60,19 @@ def count(
provider_model_ref=routed.resolved.primary.provider_model_ref,
gateway_model=routed.resolved.original_model,
)
request_snapshot = anthropic_request_snapshot(routed.request)
request_snapshot["model"] = routed.resolved.original_model
trace_event(
lambda: {
"snapshot": {
**anthropic_request_snapshot(routed.request),
"model": routed.resolved.original_model,
}
},
stage="ingress",
event="free_claude_code.api.count_tokens.completed",
source="api",
request_id=request_id,
message_count=len(routed.request.messages),
input_tokens=tokens,
snapshot=request_snapshot,
)
return TokenCountResponse(input_tokens=tokens)
except ApplicationError:
Expand Down
15 changes: 8 additions & 7 deletions src/free_claude_code/api/response_streams.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
from typing import Literal

from fastapi.responses import JSONResponse, Response, StreamingResponse
from loguru import logger
from starlette.background import BackgroundTask
from starlette.responses import ContentStream
from starlette.types import Receive, Scope, Send
Expand Down Expand Up @@ -114,15 +115,15 @@ async def _cleanup(self, *, preserved_error: BaseException | None) -> None:
preserved_error=preserved_error,
)
except Exception as exc:
_trace_response_cleanup_failure("close_body", exc)
_log_response_cleanup_failure("close_body", exc)

release = self._release
if release is None:
return
try:
await release()
except Exception as exc:
_trace_response_cleanup_failure("release_resource", exc)
_log_response_cleanup_failure("release_resource", exc)


async def _wait_for_cleanup(task: asyncio.Task[None]) -> None:
Expand All @@ -134,28 +135,28 @@ async def _wait_for_cleanup(task: asyncio.Task[None]) -> None:
except asyncio.CancelledError as exc:
cancellation = exc

# Ordinary defensive failures are trace-only; cancellation remains control flow.
# Ordinary cleanup failures preserve the response; cancellation remains control flow.
try:
task.result()
except asyncio.CancelledError:
if cancellation is not None:
raise cancellation from None
raise
except Exception as exc:
_trace_response_cleanup_failure("cleanup_task", exc)
_log_response_cleanup_failure("cleanup_task", exc)

if cancellation is not None:
raise cancellation


def _trace_response_cleanup_failure(operation: str, exc: BaseException) -> None:
trace_event(
def _log_response_cleanup_failure(operation: str, exc: BaseException) -> None:
logger.bind(
stage="egress",
event="free_claude_code.api.response.cleanup_failed",
source="api",
operation=operation,
exc_type=type(exc).__name__,
)
).opt(exception=exc).warning("Response cleanup failed")


async def bind_response_lifetime(
Expand Down
47 changes: 33 additions & 14 deletions src/free_claude_code/application/code_sessions/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -183,9 +183,10 @@ async def _close(self) -> None:
await self._stop(
owner.state.session.id, owner.state.run.id, shutting_down=True
)
except Exception:
except Exception as exc:
self._mark_storage_failed(owner, "stop_run", exc)
await self._close_connection(owner)
self._storage_failure(owner)
self._notify_storage_failure(owner)
jobs = [owner.job for owner in owners if owner.job is not None]
if jobs:
await asyncio.wait(jobs, timeout=5)
Expand All @@ -199,8 +200,9 @@ async def _close(self) -> None:
"interrupted",
"FCC stopped before this turn finished.",
)
except Exception:
self._storage_failure(owner)
except Exception as exc:
self._mark_storage_failed(owner, "finish_run", exc)
self._notify_storage_failure(owner)
if self._jobs:
await asyncio.gather(*tuple(self._jobs), return_exceptions=True)
self._events.close()
Expand Down Expand Up @@ -978,16 +980,16 @@ async def _commit_progress(
)
except CodeConflictError, CodeNotFoundError:
raise
except Exception:
owner.storage_failed = True
except Exception as exc:
self._mark_storage_failed(owner, "save_progress", exc)
if owner.failure_task is None:
owner.failure_task = self._job(self._halt_storage(owner))
raise
owner.state.apply_progress(progress)

async def _halt_storage(self, owner: _SessionRuntime) -> None:
await self._close_connection(owner)
self._storage_failure(owner)
self._notify_storage_failure(owner)

def _detach_connection_locked(
self, owner: _SessionRuntime
Expand Down Expand Up @@ -1017,8 +1019,9 @@ async def _close_connection(self, owner: _SessionRuntime) -> None:
owner, owner.state.attach_thread(connection.thread_id)
)
await self._expire_prompts_locked(owner)
except Exception:
self._storage_failure(owner)
except Exception as exc:
self._mark_storage_failed(owner, "close_connection", exc)
self._notify_storage_failure(owner)

async def _fail(
self,
Expand Down Expand Up @@ -1049,11 +1052,24 @@ async def _fail(
else "failed"
)
await self._finish_locked(owner, status, message)
except Exception:
self._storage_failure(owner)
except Exception as exc:
self._mark_storage_failed(owner, "fail_run", exc)
self._notify_storage_failure(owner)

def _storage_failure(self, owner: _SessionRuntime) -> None:
def _mark_storage_failed(
self, owner: _SessionRuntime, operation: str, error: Exception
) -> None:
if owner.storage_failed:
return
owner.storage_failed = True
logger.bind(
event="code.storage_failed",
operation=operation,
session_id=owner.state.session.id,
run_id=owner.state.run.id if owner.state.run is not None else None,
).opt(exception=error).error("Code session stopped after a storage failure")

def _notify_storage_failure(self, owner: _SessionRuntime) -> None:
owner.finished.set()
self._publish(
owner,
Expand Down Expand Up @@ -1104,8 +1120,11 @@ async def _delete_native(
progress = owner.state.deletion_failed(status, _error_message(exc))
try:
await self._commit_progress(owner, progress)
except Exception:
self._storage_failure(owner)
except Exception as save_error:
self._mark_storage_failed(
owner, "save_deletion_failure", save_error
)
self._notify_storage_failure(owner)
return
self._publish(owner, "session.updated")
if status == "delete_uncertain" and not reconcile and not native_complete:
Expand Down
19 changes: 10 additions & 9 deletions src/free_claude_code/application/execution.py
Original file line number Diff line number Diff line change
Expand Up @@ -167,7 +167,7 @@ def stream_messages(
self,
routed: RoutedMessagesRequest,
*,
raw_log_payload: object,
raw_log_payload: Callable[[], object],
request_id: str,
) -> AsyncIterator[str]:
"""Execute one Anthropic Messages request."""
Expand Down Expand Up @@ -210,7 +210,7 @@ async def open_candidate(
wire_api="messages",
raw_log_label="FULL_PAYLOAD",
raw_log_payload=raw_log_payload,
request_snapshot=anthropic_request_snapshot(routed.request),
request_snapshot=lambda: anthropic_request_snapshot(routed.request),
ingress_count_name="message_count",
ingress_count=len(routed.request.messages),
request_id=request_id,
Expand All @@ -221,7 +221,7 @@ def stream_responses(
self,
routed: RoutedResponsesRequest,
*,
raw_log_payload: object,
raw_log_payload: Callable[[], object],
request_id: str,
) -> AsyncIterator[str]:
"""Execute one native OpenAI Responses request."""
Expand Down Expand Up @@ -266,7 +266,7 @@ async def open_candidate(
wire_api="responses",
raw_log_label="FULL_RESPONSES_PAYLOAD",
raw_log_payload=raw_log_payload,
request_snapshot={
request_snapshot=lambda: {
"model": routed.request.model,
"input_item_count": input_item_count,
"tool_count": len(routed.request.tools or ()),
Expand All @@ -284,8 +284,8 @@ def _stream_candidates(
reasoning: ReasoningPolicy,
wire_api: WireApi,
raw_log_label: str,
raw_log_payload: object,
request_snapshot: dict[str, object],
raw_log_payload: Callable[[], object],
request_snapshot: Callable[[], dict[str, object]],
ingress_count_name: str,
ingress_count: int,
request_id: str,
Expand Down Expand Up @@ -318,7 +318,6 @@ def _stream_candidates(
route_trace["generation_id"] = self._generation_id
trace_event(**route_trace)

request_snapshot["model"] = gateway_model
ingress_trace: dict[str, object] = {
"stage": "ingress",
"event": (
Expand All @@ -327,16 +326,18 @@ def _stream_candidates(
else "free_claude_code.api.request.received"
),
"source": "api",
"snapshot": request_snapshot,
"request_id": request_id,
ingress_count_name: ingress_count,
}
trace_event(
lambda: {"snapshot": {**request_snapshot(), "model": gateway_model}},
**ingress_trace,
)

if self._log_raw_payloads:
logger.debug(f"{raw_log_label} [{{}}]: {{}}", request_id, raw_log_payload)
logger.opt(lazy=True).debug(
f"{raw_log_label} [{{}}]: {{}}", lambda: request_id, raw_log_payload
)

async def provider_body() -> AsyncIterator[str]:
loop = asyncio.get_running_loop()
Expand Down
2 changes: 1 addition & 1 deletion src/free_claude_code/application/web_tools/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -138,7 +138,7 @@ async def _stream_automatic_search(
)
provider_stream = self._executor.stream_messages(
translated,
raw_log_payload=plan.request.model_dump(),
raw_log_payload=plan.request.model_dump,
request_id=request_id,
)
chunks: list[str] = []
Expand Down
Loading
Loading