Spaces:
Running
Running
Upload 9 files
Browse files- .codex/auth.json +11 -0
- .codex/config.toml +1 -0
- app.py +95 -18
.codex/auth.json
ADDED
|
@@ -0,0 +1,11 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
{
|
| 2 |
+
"auth_mode": "chatgpt",
|
| 3 |
+
"OPENAI_API_KEY": null,
|
| 4 |
+
"tokens": {
|
| 5 |
+
"id_token": "eyJhbGciOiJSUzI1NiIsImtpZCI6ImIxZGQzZjhmLTlhYWQtNDdmZS1iMGU3LWVkYjAwOTc3N2Q2YiIsInR5cCI6IkpXVCJ9.eyJhdF9oYXNoIjoiZ2VZMnVaUmc0alR6MW9RbHhUVFRtUSIsImF1ZCI6WyJhcHBfRU1vYW1FRVo3M2YwQ2tYYVhwN2hyYW5uIl0sImF1dGhfcHJvdmlkZXIiOiJnb29nbGUiLCJhdXRoX3RpbWUiOjE3ODA3MzU1NTksImVtYWlsIjoiZGV2YXJzaGlhNUBnbWFpbC5jb20iLCJlbWFpbF92ZXJpZmllZCI6dHJ1ZSwiZXhwIjoxNzgwNzM5MTYxLCJodHRwczovL2FwaS5vcGVuYWkuY29tL2F1dGgiOnsiY2hhdGdwdF9hY2NvdW50X2lkIjoiY2UwN2I3YzgtZGZlYy00M2EyLWE1YzctNDc0YjVmYjgwNTU0IiwiY2hhdGdwdF9wbGFuX3R5cGUiOiJnbyIsImNoYXRncHRfc3Vic2NyaXB0aW9uX2FjdGl2ZV9zdGFydCI6IjIwMjUtMTEtMDRUMTU6Mzg6MzUrMDA6MDAiLCJjaGF0Z3B0X3N1YnNjcmlwdGlvbl9hY3RpdmVfdW50aWwiOiIyMDI2LTA3LTA0VDE1OjM4OjM1KzAwOjAwIiwiY2hhdGdwdF9zdWJzY3JpcHRpb25fbGFzdF9jaGVja2VkIjoiMjAyNi0wNi0wNlQwODo0NTo1OS42MDAzODUrMDA6MDAiLCJjaGF0Z3B0X3VzZXJfaWQiOiJ1c2VyLWR6UFF0TmtBMlkxakFBcmlVWUhWTHR2bCIsImdyb3VwcyI6W10sImxvY2FsaG9zdCI6dHJ1ZSwib3JnYW5pemF0aW9ucyI6W3siaWQiOiJvcmctMlg0SUZveUJLTVpWMktkTDhYUW9PeXJtIiwiaXNfZGVmYXVsdCI6dHJ1ZSwicm9sZSI6Im93bmVyIiwidGl0bGUiOiJjdXJzZW9md2l0Y2hlciJ9LHsiaWQiOiJvcmctdXYyQWd3VHE2VkRjNkZaZFpiSWd6aTN6IiwiaXNfZGVmYXVsdCI6ZmFsc2UsInJvbGUiOiJvd25lciIsInRpdGxlIjoiUGVyc29uYWwifV0sInVzZXJfaWQiOiJ1c2VyLWR6UFF0TmtBMlkxakFBcmlVWUhWTHR2bCJ9LCJpYXQiOjE3ODA3MzU1NjEsImlzcyI6Imh0dHBzOi8vYXV0aC5vcGVuYWkuY29tIiwianRpIjoiYjBjMjE4MjUtMWQwZi00YzgxLThjMGMtNWQzNzhiOGRlODU1IiwibmFtZSI6IkFkaXR5YSBEZXZhcnNoaSIsInJhdCI6MTc4MDczNTU0OSwic2lkIjoiY2E0NzQzZmEtMmZlNS00OGNhLWI1NmEtMjEwY2Y0M2RhNmJmIiwic3ViIjoiZ29vZ2xlLW9hdXRoMnwxMTExNTE3MzYzMTIzOTgyNjg3OTcifQ.XCd2DYKUpnHW7DJroTNN3CR2-WzxFsJrZNOtPSZ_lJd4-HnqVN9KjzIcaIHrZKKOLEjJQraavI-FV9pM_2OqP_utp7y69gsaHiuEQFeL8tPikM388EWijFqjlqcsq5JyruLP64L3YKcl63ODr5STy9fukgaIZZJq-o5QzgJrYvhQnAlH-V_2qRbRcIFVI7RyRnsyYToT2Dbvl4eo4dZun6a7Xkh3glznYR0VYHfL2iz3hFuKxRLbtV39ce_v2Nc2JXp1VHVEW-9j6p07X_bHRAY92E41xmpS5oB5FMpKCC9L7Dn_MfMtvzWVXtIhIrTctdJvp8vV5MxBhPWCkymd0obR-FwzfKLw6oZsBHKY6Vq5xfgM7LfrJ6M6cyw3No3EbbDTiYQ5nccE7YrGZDKMplp7kYAW36Fj3FKVlGlgcLhFTxtEZRUK3uGGs4hP1p8DUCCnAkWlNLqmwUjUwrpOkM8DyUKH2RXLjCq2sV_NsqXdDrE17WySYsfaLq8aaN-4uqbJfKSSpc5XE4ZeOL4MKMpTO6UCTDYIFvP0m4rFhdg5mgptF1-lJja0WovQEPq8LXRY6m9NXd8ZKpYDz3JatgYYJJVsuhV3CaMTQ51TjQpg0_pA6XRMTmr8jqz7HQcQ-jAfIbLl48EgD-qpqjJT2nTc4qHPxX24Xa4M9pPnNpw",
|
| 6 |
+
"access_token": "eyJhbGciOiJSUzI1NiIsImtpZCI6IjE5MzQ0ZTY1LWJiYzktNDRkMS1hOWQwLWY5NTdiMDc5YmQwZSIsInR5cCI6IkpXVCJ9.eyJhdWQiOlsiaHR0cHM6Ly9hcGkub3BlbmFpLmNvbS92MSJdLCJjbGllbnRfaWQiOiJhcHBfRU1vYW1FRVo3M2YwQ2tYYVhwN2hyYW5uIiwiZXhwIjoxNzgxNTk5NTYyLCJodHRwczovL2FwaS5vcGVuYWkuY29tL2F1dGgiOnsiY2hhdGdwdF9hY2NvdW50X2lkIjoiY2UwN2I3YzgtZGZlYy00M2EyLWE1YzctNDc0YjVmYjgwNTU0IiwiY2hhdGdwdF9hY2NvdW50X3VzZXJfaWQiOiJ1c2VyLWR6UFF0TmtBMlkxakFBcmlVWUhWTHR2bF9fY2UwN2I3YzgtZGZlYy00M2EyLWE1YzctNDc0YjVmYjgwNTU0IiwiY2hhdGdwdF9jb21wdXRlX3Jlc2lkZW5jeSI6Im5vX2NvbnN0cmFpbnQiLCJjaGF0Z3B0X3BsYW5fdHlwZSI6ImdvIiwiY2hhdGdwdF91c2VyX2lkIjoidXNlci1kelBRdE5rQTJZMWpBQXJpVVlIVkx0dmwiLCJsb2NhbGhvc3QiOnRydWUsInVzZXJfaWQiOiJ1c2VyLWR6UFF0TmtBMlkxakFBcmlVWUhWTHR2bCJ9LCJodHRwczovL2FwaS5vcGVuYWkuY29tL3Byb2ZpbGUiOnsiZW1haWwiOiJkZXZhcnNoaWE1QGdtYWlsLmNvbSIsImVtYWlsX3ZlcmlmaWVkIjp0cnVlfSwiaWF0IjoxNzgwNzM1NTYxLCJpc3MiOiJodHRwczovL2F1dGgub3BlbmFpLmNvbSIsImp0aSI6ImIyNzQ5MmY1LTczODktNDRkYi1hMWI1LTU1OWRhZmFlMmFkNCIsIm5iZiI6MTc4MDczNTU2MSwicHdkX2F1dGhfdGltZSI6MTc4MDczNTU1OTYwMCwic2NwIjpbIm9wZW5pZCIsInByb2ZpbGUiLCJlbWFpbCIsIm9mZmxpbmVfYWNjZXNzIiwiYXBpLmNvbm5lY3RvcnMucmVhZCIsImFwaS5jb25uZWN0b3JzLmludm9rZSJdLCJzZXNzaW9uX2lkIjoiYXV0aHNlc3NfdVNNdGphSWJacVZqQW5xT254VzJWRUowIiwic2wiOnRydWUsInN1YiI6Imdvb2dsZS1vYXV0aDJ8MTExMTUxNzM2MzEyMzk4MjY4Nzk3In0.KvxgZQwe2_bSPwwIRsG0Xg5_GSWHdQru2xtVPXY1GOuYNaMgdbuaO6CDTtF1aEfhCXfzNpu1TTFvEMkHlsYCmZ_OK5leY3e4a01sQ7lUNdVCJ50E0QFsPxxfqPkziip7yzb76d_ufPwgkw0xUa_ierxNnc-fuaw2OgveQndp4yRT6u4WuHrmZfKVXgK9WgoZYiggmEzL-4D4SKpRv8T-IwlkT25J9hT1tb_Sdt19YLue2Xhpuaq0xtW6_G7g1gKp0CjSQGSF8-xzVfc6kJkqAnK_BQ-i3izQHFz2keFeSIRMoITA27PeDDKWhO0WnKR4IkWEEu8tOaLHO10GZMmi5P1vGIaJMbGM8USTVSGLFZtqZ1CAxvWxd7xp6QTdtXsjmwdmID4CrxvU-RgAReXqfcUp3dX5FxfFeWahwNo5d19UlbGVVz7gJ999rlSg2EQhcWy39omje8nifzEenEpnJX15tJ-AqCv6Wy2x68-STrZbmkPZAkPtFvRIRlsO1hcPnWak1wQ9sR2FhSyoCb8YWgVjghNKe6XKeTtjnJeZrg6gQ04q_OBZ62rQcqwIiyIJEg6pZGOEyHot7sL-Z1v7068MVZWUh4psa_WNb93m8nWMwfRI_7IrDhjywhHhnwB-iB6CfLJbhsWOa6wxNeAoAJCDe43pV2Ida0qU0wY0sbA",
|
| 7 |
+
"refresh_token": "rt.1.AAArhYuyl3xt4V7ysuMdJufBWZQKMmzT7yUQu3r_5GH74-rKwSqVqFqveZ4b0R0BrV6eYZ-3U2eRkXF95pqlWczEQmarc-hnZMhfYvAqJt-HCLbVW2_OOjNq00Ka20D2bHRtTJ1Mdktd59oSVbhY7ugyJJWXXnQtp5O3URFezLJqnGUFgV55VnfjIrVSdFNEAkzSyfy1-UbumyAELJcxKyRiWEUezEgURE8bVnA01dPgbh_8VeuGFt1zlZURqXCTm7ugatfhAj69qrVhA0JCp_QygWdRqLcvD_GQLnMSWPaOzsNEsmogX2MWAsYFByOdnKk",
|
| 8 |
+
"account_id": "ce07b7c8-dfec-43a2-a5c7-474b5fb80554"
|
| 9 |
+
},
|
| 10 |
+
"last_refresh": "2026-06-06T08:46:04.174823900Z"
|
| 11 |
+
}
|
.codex/config.toml
ADDED
|
@@ -0,0 +1 @@
|
|
|
|
|
|
|
| 1 |
+
personality = "pragmatic"
|
app.py
CHANGED
|
@@ -17,11 +17,13 @@ Sessions: pass `X-Session-Id` (or the OpenAI `user` field) for a persistent
|
|
| 17 |
workdir at /data/sessions/<id>/workspace + Codex thread resume. None -> ephemeral.
|
| 18 |
"""
|
| 19 |
|
|
|
|
| 20 |
import json
|
| 21 |
import os
|
| 22 |
import re
|
| 23 |
import time
|
| 24 |
import uuid
|
|
|
|
| 25 |
from pathlib import Path
|
| 26 |
from typing import Any, Optional
|
| 27 |
|
|
@@ -51,11 +53,68 @@ API_TOKEN = os.environ.get("API_TOKEN", "") # HF secret; if empty, auth is OPEN
|
|
| 51 |
DEFAULT_SANDBOX = os.environ.get("CODEX_SANDBOX", "workspace-write") # or read-only
|
| 52 |
CODEX_MODEL = os.environ.get("CODEX_MODEL", "").strip() # optional override
|
| 53 |
READ_TIMEOUT = float(os.environ.get("CODEX_TIMEOUT", "180")) # per-output-gap secs
|
|
|
|
|
|
|
|
|
|
|
|
|
| 54 |
DEFAULT_MODEL_NAME = "codex"
|
| 55 |
|
| 56 |
SESSION_ID_RE = re.compile(r"[^A-Za-z0-9_.-]")
|
| 57 |
|
| 58 |
-
app = FastAPI(title="Codex-as-API", version="2.
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 59 |
|
| 60 |
|
| 61 |
# --------------------------------------------------------------------------- #
|
|
@@ -197,6 +256,8 @@ async def health():
|
|
| 197 |
"auth_required": bool(API_TOKEN),
|
| 198 |
"sandbox": DEFAULT_SANDBOX,
|
| 199 |
"engine": "app-server",
|
|
|
|
|
|
|
| 200 |
}
|
| 201 |
|
| 202 |
|
|
@@ -233,6 +294,11 @@ async def chat_completions(
|
|
| 233 |
if not prompt:
|
| 234 |
raise HTTPException(status_code=400, detail="Empty prompt after parsing.")
|
| 235 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 236 |
turn = run_turn(
|
| 237 |
codex_bin=CODEX_BIN,
|
| 238 |
codex_home=CODEX_HOME,
|
|
@@ -246,8 +312,9 @@ async def chat_completions(
|
|
| 246 |
|
| 247 |
if req.stream:
|
| 248 |
include_usage = bool((req.stream_options or {}).get("include_usage"))
|
|
|
|
| 249 |
return StreamingResponse(
|
| 250 |
-
_sse_stream(turn, model_name, session_dir, include_usage),
|
| 251 |
media_type="text/event-stream",
|
| 252 |
)
|
| 253 |
|
|
@@ -255,22 +322,29 @@ async def chat_completions(
|
|
| 255 |
content_parts: list[str] = []
|
| 256 |
usage = {"prompt_tokens": 0, "completion_tokens": 0, "total_tokens": 0}
|
| 257 |
try:
|
| 258 |
-
async
|
| 259 |
-
|
| 260 |
-
|
| 261 |
-
|
| 262 |
-
|
| 263 |
-
|
| 264 |
-
|
| 265 |
-
|
|
|
|
| 266 |
except CodexError as e:
|
| 267 |
raise HTTPException(status_code=502, detail=f"Codex engine: {e}")
|
|
|
|
|
|
|
| 268 |
|
| 269 |
return JSONResponse(_completion_payload("".join(content_parts), model_name, usage))
|
| 270 |
|
| 271 |
|
| 272 |
-
async def _sse_stream(turn, model: str, session_dir, include_usage: bool):
|
| 273 |
-
"""OpenAI-compatible SSE: role chunk, live content deltas, finish, [DONE].
|
|
|
|
|
|
|
|
|
|
|
|
|
| 274 |
cid = f"chatcmpl-{uuid.uuid4().hex}"
|
| 275 |
created = int(time.time())
|
| 276 |
|
|
@@ -289,15 +363,18 @@ async def _sse_stream(turn, model: str, session_dir, include_usage: bool):
|
|
| 289 |
yield chunk({"role": "assistant"})
|
| 290 |
final_usage = {"prompt_tokens": 0, "completion_tokens": 0, "total_tokens": 0}
|
| 291 |
try:
|
| 292 |
-
async
|
| 293 |
-
|
| 294 |
-
|
| 295 |
-
|
| 296 |
-
|
| 297 |
-
|
|
|
|
| 298 |
except CodexError as e:
|
| 299 |
# Surface the error inside the stream, then close cleanly.
|
| 300 |
yield chunk({"content": f"\n\n[codex error: {e}]"})
|
|
|
|
|
|
|
| 301 |
|
| 302 |
yield chunk({}, finish="stop")
|
| 303 |
if include_usage:
|
|
|
|
| 17 |
workdir at /data/sessions/<id>/workspace + Codex thread resume. None -> ephemeral.
|
| 18 |
"""
|
| 19 |
|
| 20 |
+
import asyncio
|
| 21 |
import json
|
| 22 |
import os
|
| 23 |
import re
|
| 24 |
import time
|
| 25 |
import uuid
|
| 26 |
+
from contextlib import aclosing
|
| 27 |
from pathlib import Path
|
| 28 |
from typing import Any, Optional
|
| 29 |
|
|
|
|
| 53 |
DEFAULT_SANDBOX = os.environ.get("CODEX_SANDBOX", "workspace-write") # or read-only
|
| 54 |
CODEX_MODEL = os.environ.get("CODEX_MODEL", "").strip() # optional override
|
| 55 |
READ_TIMEOUT = float(os.environ.get("CODEX_TIMEOUT", "180")) # per-output-gap secs
|
| 56 |
+
# Max Codex turns running at once across all sessions (each is a heavy process).
|
| 57 |
+
MAX_CONCURRENCY = int(os.environ.get("CODEX_MAX_CONCURRENCY", "4"))
|
| 58 |
+
# How long a request may wait in the queue before we give up with 429.
|
| 59 |
+
QUEUE_TIMEOUT = float(os.environ.get("CODEX_QUEUE_TIMEOUT", "90"))
|
| 60 |
DEFAULT_MODEL_NAME = "codex"
|
| 61 |
|
| 62 |
SESSION_ID_RE = re.compile(r"[^A-Za-z0-9_.-]")
|
| 63 |
|
| 64 |
+
app = FastAPI(title="Codex-as-API", version="2.1.0")
|
| 65 |
+
|
| 66 |
+
# --------------------------------------------------------------------------- #
|
| 67 |
+
# Concurrency control
|
| 68 |
+
# - _GLOBAL_SEM caps total simultaneous Codex processes (resource guard).
|
| 69 |
+
# - _SESSION_LOCKS serialize requests that share a session id, so two calls
|
| 70 |
+
# never resume/operate on the same thread + workspace at once (corruption).
|
| 71 |
+
# Different sessions still run fully in parallel (up to the global cap).
|
| 72 |
+
# --------------------------------------------------------------------------- #
|
| 73 |
+
_GLOBAL_SEM = asyncio.Semaphore(MAX_CONCURRENCY)
|
| 74 |
+
_SESSION_LOCKS: dict[str, asyncio.Lock] = {}
|
| 75 |
+
|
| 76 |
+
|
| 77 |
+
def _session_lock(session_id: Optional[str]) -> Optional[asyncio.Lock]:
|
| 78 |
+
if not session_id:
|
| 79 |
+
return None # ephemeral request: unique workspace, no shared state
|
| 80 |
+
# setdefault is atomic between awaits in asyncio's single thread.
|
| 81 |
+
return _SESSION_LOCKS.setdefault(session_id, asyncio.Lock())
|
| 82 |
+
|
| 83 |
+
|
| 84 |
+
class _TurnGuard:
|
| 85 |
+
"""Acquire (session lock -> global slot) with a bounded wait; release both."""
|
| 86 |
+
|
| 87 |
+
def __init__(self, session_id: Optional[str]):
|
| 88 |
+
self._lock = _session_lock(session_id)
|
| 89 |
+
self._have_lock = False
|
| 90 |
+
self._have_slot = False
|
| 91 |
+
|
| 92 |
+
async def acquire(self) -> None:
|
| 93 |
+
loop = asyncio.get_event_loop()
|
| 94 |
+
deadline = loop.time() + QUEUE_TIMEOUT
|
| 95 |
+
try:
|
| 96 |
+
if self._lock is not None:
|
| 97 |
+
await asyncio.wait_for(
|
| 98 |
+
self._lock.acquire(), timeout=QUEUE_TIMEOUT
|
| 99 |
+
)
|
| 100 |
+
self._have_lock = True
|
| 101 |
+
remaining = max(0.1, deadline - loop.time())
|
| 102 |
+
await asyncio.wait_for(_GLOBAL_SEM.acquire(), timeout=remaining)
|
| 103 |
+
self._have_slot = True
|
| 104 |
+
except (asyncio.TimeoutError, TimeoutError):
|
| 105 |
+
self.release()
|
| 106 |
+
raise HTTPException(
|
| 107 |
+
status_code=429,
|
| 108 |
+
detail="Server busy (concurrency/session limit). Retry shortly.",
|
| 109 |
+
)
|
| 110 |
+
|
| 111 |
+
def release(self) -> None:
|
| 112 |
+
if self._have_slot:
|
| 113 |
+
self._have_slot = False
|
| 114 |
+
_GLOBAL_SEM.release()
|
| 115 |
+
if self._have_lock:
|
| 116 |
+
self._have_lock = False
|
| 117 |
+
self._lock.release()
|
| 118 |
|
| 119 |
|
| 120 |
# --------------------------------------------------------------------------- #
|
|
|
|
| 256 |
"auth_required": bool(API_TOKEN),
|
| 257 |
"sandbox": DEFAULT_SANDBOX,
|
| 258 |
"engine": "app-server",
|
| 259 |
+
"max_concurrency": MAX_CONCURRENCY,
|
| 260 |
+
"active_sessions": len(_SESSION_LOCKS),
|
| 261 |
}
|
| 262 |
|
| 263 |
|
|
|
|
| 294 |
if not prompt:
|
| 295 |
raise HTTPException(status_code=400, detail="Empty prompt after parsing.")
|
| 296 |
|
| 297 |
+
# Acquire concurrency guard BEFORE starting work, so we can fail fast with
|
| 298 |
+
# 429 (for streaming, headers are sent once we return — can't 429 later).
|
| 299 |
+
guard = _TurnGuard(session_id)
|
| 300 |
+
await guard.acquire()
|
| 301 |
+
|
| 302 |
turn = run_turn(
|
| 303 |
codex_bin=CODEX_BIN,
|
| 304 |
codex_home=CODEX_HOME,
|
|
|
|
| 312 |
|
| 313 |
if req.stream:
|
| 314 |
include_usage = bool((req.stream_options or {}).get("include_usage"))
|
| 315 |
+
# _sse_stream owns the guard from here and releases it when done.
|
| 316 |
return StreamingResponse(
|
| 317 |
+
_sse_stream(turn, model_name, session_dir, include_usage, guard),
|
| 318 |
media_type="text/event-stream",
|
| 319 |
)
|
| 320 |
|
|
|
|
| 322 |
content_parts: list[str] = []
|
| 323 |
usage = {"prompt_tokens": 0, "completion_tokens": 0, "total_tokens": 0}
|
| 324 |
try:
|
| 325 |
+
async with aclosing(turn) as t:
|
| 326 |
+
async for evt in t:
|
| 327 |
+
if evt["type"] == "delta":
|
| 328 |
+
content_parts.append(evt["text"])
|
| 329 |
+
elif evt["type"] == "final":
|
| 330 |
+
if evt.get("text"):
|
| 331 |
+
content_parts = [evt["text"]] # authoritative full text
|
| 332 |
+
usage = evt.get("usage", usage)
|
| 333 |
+
_persist_thread(session_dir, evt.get("thread_id"))
|
| 334 |
except CodexError as e:
|
| 335 |
raise HTTPException(status_code=502, detail=f"Codex engine: {e}")
|
| 336 |
+
finally:
|
| 337 |
+
guard.release()
|
| 338 |
|
| 339 |
return JSONResponse(_completion_payload("".join(content_parts), model_name, usage))
|
| 340 |
|
| 341 |
|
| 342 |
+
async def _sse_stream(turn, model: str, session_dir, include_usage: bool, guard):
|
| 343 |
+
"""OpenAI-compatible SSE: role chunk, live content deltas, finish, [DONE].
|
| 344 |
+
|
| 345 |
+
Owns `guard`: releases the concurrency slot/session lock when the turn ends
|
| 346 |
+
OR the client disconnects (aclosing -> run_turn finally kills the process).
|
| 347 |
+
"""
|
| 348 |
cid = f"chatcmpl-{uuid.uuid4().hex}"
|
| 349 |
created = int(time.time())
|
| 350 |
|
|
|
|
| 363 |
yield chunk({"role": "assistant"})
|
| 364 |
final_usage = {"prompt_tokens": 0, "completion_tokens": 0, "total_tokens": 0}
|
| 365 |
try:
|
| 366 |
+
async with aclosing(turn) as t:
|
| 367 |
+
async for evt in t:
|
| 368 |
+
if evt["type"] == "delta":
|
| 369 |
+
yield chunk({"content": evt["text"]})
|
| 370 |
+
elif evt["type"] == "final":
|
| 371 |
+
final_usage = evt.get("usage", final_usage)
|
| 372 |
+
_persist_thread(session_dir, evt.get("thread_id"))
|
| 373 |
except CodexError as e:
|
| 374 |
# Surface the error inside the stream, then close cleanly.
|
| 375 |
yield chunk({"content": f"\n\n[codex error: {e}]"})
|
| 376 |
+
finally:
|
| 377 |
+
guard.release()
|
| 378 |
|
| 379 |
yield chunk({}, finish="stop")
|
| 380 |
if include_usage:
|