Skip to content
Open
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
40 changes: 40 additions & 0 deletions app/db/alembic/versions/20260728_000000_add_file_account_pins.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
"""add durable file account pins

Revision ID: 20260728_000000_add_file_account_pins
Revises: 20260725_000000_add_http_bridge_pending_tool_calls
Create Date: 2026-07-28
"""

from __future__ import annotations

import sqlalchemy as sa
from alembic import op

revision = "20260728_000000_add_file_account_pins"
down_revision = "20260725_000000_add_http_bridge_pending_tool_calls"
branch_labels = None
depends_on = None

_TABLE = "file_account_pins"


def upgrade() -> None:
bind = op.get_bind()
if sa.inspect(bind).has_table(_TABLE):
return
op.create_table(
_TABLE,
sa.Column("file_id", sa.String(), nullable=False),
sa.Column("account_id", sa.String(), nullable=False),
sa.Column("expires_at", sa.DateTime(timezone=True), nullable=False),
sa.PrimaryKeyConstraint("file_id"),
)
op.create_index("ix_file_account_pins_expires_at", _TABLE, ["expires_at"], unique=False)


def downgrade() -> None:
bind = op.get_bind()
if not sa.inspect(bind).has_table(_TABLE):
return
op.drop_index("ix_file_account_pins_expires_at", table_name=_TABLE)
op.drop_table(_TABLE)
10 changes: 10 additions & 0 deletions app/db/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,16 @@ class StickySessionKind(str, Enum):
PROMPT_CACHE = "prompt_cache"


class FileAccountPin(Base):
__tablename__ = "file_account_pins"

file_id: Mapped[str] = mapped_column(String, primary_key=True)
account_id: Mapped[str] = mapped_column(String, nullable=False)
expires_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)

__table_args__ = (Index("ix_file_account_pins_expires_at", "expires_at"),)


class RequestKind(str, Enum):
NORMAL = "normal"
WARMUP = "warmup"
Expand Down
54 changes: 31 additions & 23 deletions app/modules/proxy/_service/api_key_usage.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
from app.modules.proxy._service.support import (
_ApiKeyReservationTouchState,
_consume_api_key_reservation_heartbeat_result,
_signal_propagated_responses_service_cleanup_ready,
_StreamSettlement,
_WebSocketRequestState,
)
Expand Down Expand Up @@ -274,28 +275,34 @@ async def _settle_compact_api_key_usage(
)

proxy = cast(_ApiKeyUsageServiceProtocol, self)
with anyio.CancelScope(shield=True):
try:
async with proxy._repo_factory() as repos:
api_keys_service = _service_api_keys_service()(repos.api_keys)
if response is not None and input_tokens is not None and output_tokens is not None:
await api_keys_service.finalize_usage_reservation(
reservation_id,
model=model_name,
input_tokens=input_tokens,
output_tokens=output_tokens,
cached_input_tokens=cached_input_tokens or 0,
service_tier=service_tier,
)
else:
await api_keys_service.release_usage_reservation(reservation_id)
except Exception:
logger.warning(
"Failed to settle compact API key reservation key_id=%s request_id=%s",
api_key.id,
get_request_id(),
exc_info=True,
)
try:
with anyio.CancelScope(shield=True):
try:
async with proxy._repo_factory() as repos:
api_keys_service = _service_api_keys_service()(repos.api_keys)
if response is not None and input_tokens is not None and output_tokens is not None:
await api_keys_service.finalize_usage_reservation(
reservation_id,
model=model_name,
input_tokens=input_tokens,
output_tokens=output_tokens,
cached_input_tokens=cached_input_tokens or 0,
service_tier=service_tier,
)
else:
await api_keys_service.release_usage_reservation(reservation_id)
except Exception:
logger.warning(
"Failed to settle compact API key reservation key_id=%s request_id=%s",
api_key.id,
get_request_id(),
exc_info=True,
)
finally:
# The compact service has made its one cancellation-safe settlement
# attempt. A caller that created the reservation must not issue a
# second release as a fallback after this boundary.
_signal_propagated_responses_service_cleanup_ready()

async def _settle_stream_api_key_usage(
self,
Expand Down Expand Up @@ -429,7 +436,7 @@ def _schedule_cancel_safe_cleanup(
*,
action: str,
request_id: str,
) -> None:
) -> asyncio.Task[None]:
task = asyncio.create_task(coro, name=f"proxy-{action}-{request_id}")
proxy = cast(_ApiKeyUsageServiceProtocol, self)
proxy._background_cleanup_tasks.add(task)
Expand All @@ -449,6 +456,7 @@ def _cleanup_done(done_task: asyncio.Task[None]) -> None:
)

task.add_done_callback(_cleanup_done)
return task

async def _release_unsettled_stream_api_key_usage(
self,
Expand Down
Loading
Loading