Skip to content

Commit f03d725

Browse files
committed
fix(studio): report each migration turn's own cost
The migration page showed what Codex did but never what it cost: no turn duration, no tool time, no tokens. The intelligent build reports all three from the app-server's turn lifecycle, and the migration feed renders the same summary component, so the reader now settles one turn per source. A tool-driven migration turn ends the moment its result lands, which is why the app-server never saw turn_completed and the usage stayed invisible. Turns now settle from the app-server's own turn record read back before the session closes, and the in-Sandbox `codex exec --json` turn — which writes a bare terminal line with snake_case usage and no timestamps — settles from that line, with the total computed the way the app-server reports it. The sandbox exec stream carries no timing at all, so those rows say the duration was not reported rather than inventing one.
1 parent efbe842 commit f03d725

12 files changed

Lines changed: 1056 additions & 8 deletions

‎frontend/server/migration/activity.py‎

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,17 @@
4444
_COMPLETED_STATUSES = {"completed", "done"}
4545
_FAILED_STATUSES = {"failed", "error", "declined"}
4646

47+
# The turn's own verdict, which the reader renders as one summary row the way the
48+
# intelligent build renders its ``turn-summary`` block.
49+
_TERMINAL_TURN_STATUSES = {"completed", "failed", "cancelled", "interrupted"}
50+
_TURN_EVENT_TYPES = {
51+
"completed": "turn.completed",
52+
"failed": "turn.failed",
53+
"cancelled": "turn.interrupted",
54+
"interrupted": "turn.interrupted",
55+
}
56+
_TURN_FIELDS = ("startedAt", "completedAt", "durationMs", "model")
57+
4758
_TOOL_ITEM_TYPES = {
4859
"commandExecution": "command_execution",
4960
"fileChange": "file_change",
@@ -62,6 +73,13 @@
6273
}
6374

6475

76+
def _turn_status(value: object) -> str:
77+
"""Read the turn status the app-server reports as a string or a tagged object."""
78+
if isinstance(value, dict):
79+
value = value.get("type")
80+
return str(value or "").strip().lower()
81+
82+
6583
def _text(value: object, limit: int = _MAX_TEXT_CHARS) -> str:
6684
if not isinstance(value, str):
6785
return ""
@@ -107,13 +125,18 @@ def __init__(
107125
self._outputs: dict[str, str] = {}
108126
self._names: dict[str, str] = {}
109127
self._commands: dict[str, str] = {}
128+
self._turn: dict[str, object] = {}
129+
self._usage: dict[str, object] = {}
110130

111131
@property
112132
def lines(self) -> list[str]:
113133
return list(self._lines)
114134

115135
def line(self, event: CodexAppServerEvent) -> dict[str, object] | None:
116136
"""One ``codex exec --json`` line for ``event``; ``None`` when it has no item."""
137+
turn = self._turn_line(event)
138+
if turn is not None:
139+
return turn
117140
item = self._item(event)
118141
if item is None:
119142
return None
@@ -196,6 +219,49 @@ def _content(self) -> bytes:
196219
size -= len(lines.pop(0)) + 1
197220
return ("\n".join(lines) + "\n").encode("utf-8") if lines else b""
198221

222+
def _turn_line(self, event: CodexAppServerEvent) -> dict[str, object] | None:
223+
"""One line for the turn's own cost, written when the turn settles.
224+
225+
The page reports the same turn metrics the intelligent build does — how long
226+
the turn took, how many tools it ran, what it cost in tokens — and the
227+
intelligent build reads them off the app-server's turn lifecycle and usage
228+
events rather than off any item. Items never carry them, so they are
229+
accumulated here and written as one ``turn.*`` line that the reader turns into
230+
a summary of everything logged before it.
231+
"""
232+
kind = str(event.kind or "")
233+
if kind == "usage":
234+
if event.usage is not None:
235+
self._usage["usage"] = event.usage.public_dict()
236+
if event.thread_total is not None:
237+
self._usage["thread_total"] = event.thread_total.public_dict()
238+
window = event.model_context_window
239+
if isinstance(window, int) and not isinstance(window, bool):
240+
self._usage["model_context_window"] = window
241+
return None
242+
if kind not in {"turn_started", "turn_completed"}:
243+
return None
244+
response = event.response if isinstance(event.response, dict) else {}
245+
turn = {**self._turn, **response}
246+
if event.turn_id:
247+
turn["id"] = event.turn_id
248+
status = _turn_status(event.status or turn.get("status"))
249+
if kind == "turn_started" or status not in _TERMINAL_TURN_STATUSES:
250+
self._turn = turn
251+
return None
252+
summary: dict[str, object] = {
253+
"type": _TURN_EVENT_TYPES.get(status, "turn.completed"),
254+
"turn": {
255+
"id": str(turn.get("id") or ""),
256+
"status": status,
257+
**{key: turn[key] for key in _TURN_FIELDS if key in turn},
258+
},
259+
}
260+
summary.update(self._usage)
261+
self._turn = {}
262+
self._usage = {}
263+
return summary
264+
199265
def _event_type(self, event: CodexAppServerEvent) -> str:
200266
status = str(event.status or "").lower()
201267
if status in _FAILED_STATUSES:

‎frontend/server/migration/codex_tool_turn.py‎

Lines changed: 93 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
from __future__ import annotations
2626

2727
import asyncio
28+
import logging
2829
from collections.abc import Awaitable, Callable, Sequence
2930
from dataclasses import dataclass
3031

@@ -39,6 +40,15 @@
3940
"CodexDynamicToolResult | Awaitable[CodexDynamicToolResult]",
4041
]
4142

43+
logger = logging.getLogger(__name__)
44+
45+
# A turn that is interrupted the moment its result lands needs its own settlement, and
46+
# the app-server needs a moment to record the interruption before the turn reads back
47+
# as terminal.
48+
_TURN_SETTLE_ATTEMPTS = 6
49+
_TURN_SETTLE_SECONDS = 0.5
50+
_TERMINAL_TURN_STATUSES = {"completed", "failed", "interrupted", "cancelled"}
51+
4252
__all__ = [
4353
"DynamicTool",
4454
"ToolHandler",
@@ -126,6 +136,7 @@ async def run_tool_turn(
126136
tool.handler,
127137
)
128138
used_thread = thread_id
139+
turn_id = ""
129140
loop = asyncio.get_running_loop()
130141
try:
131142
try:
@@ -144,7 +155,11 @@ async def run_tool_turn(
144155
async for event in session.stream_turn(
145156
prompt,
146157
timeout_seconds=idle_timeout,
158+
# 本轮耗时/模型要跟智能构建一样报给页面,所以即使这不是 Studio
159+
# 任务回合也要收生命周期事件。
160+
emit_turn_lifecycle=True,
147161
):
162+
turn_id = str(getattr(event, "turn_id", "") or "") or turn_id
148163
if event_sink is not None:
149164
event_sink(event)
150165
if has_result():
@@ -170,5 +185,83 @@ async def run_tool_turn(
170185
raise ToolTurnUnavailable(str(error)) from error
171186
used_thread = session.thread_id or thread_id
172187
finally:
188+
# 结果一到手就打断的回合(以及超时收尾的回合)都走不到 app-server 的
189+
# turn_completed,而本轮耗时/模型只挂在那条事件上:会话还在的时候回读这一轮,
190+
# 替它补一条结算,页面才能像智能构建那样报出本轮的成本。
191+
await _settle_turn(
192+
session,
193+
event_sink,
194+
turn_id=turn_id or str(getattr(session, "active_turn_id", "") or ""),
195+
accepted=has_result(),
196+
)
173197
await session.close()
174198
return used_thread
199+
200+
201+
def _turn_status(turn: dict[str, object]) -> str:
202+
status = turn.get("status")
203+
if isinstance(status, dict):
204+
status = status.get("type")
205+
return str(status or "").strip().lower()
206+
207+
208+
async def _settle_turn(
209+
session: CodexAppServerSession,
210+
event_sink: Callable[[object], None] | None,
211+
*,
212+
turn_id: str,
213+
accepted: bool,
214+
) -> None:
215+
"""Report a turn's own timing when the caller stopped it before the app-server did.
216+
217+
Breaking out of the stream once the result arrives (or at the caller's deadline)
218+
leaves the turn without a ``turn_completed`` event, so codex' native timing
219+
(``startedAt`` / ``completedAt`` / ``durationMs`` / ``model``) never reaches the
220+
page. Reading the turn back keeps those numbers, and ``accepted`` says the caller
221+
took the result: the turn delivered what it was asked for, whatever codex calls the
222+
interruption the caller requested.
223+
224+
Settlement is decoration on top of the turn's real outcome, so it never raises.
225+
"""
226+
if event_sink is None or not turn_id:
227+
return
228+
try:
229+
read_turn = getattr(session, "read_turn", None)
230+
lifecycle = getattr(session, "turn_lifecycle_event", None)
231+
if not callable(read_turn) or not callable(lifecycle):
232+
return
233+
turn: dict[str, object] | None = None
234+
for attempt in range(_TURN_SETTLE_ATTEMPTS):
235+
candidate = await read_turn(turn_id)
236+
if not isinstance(candidate, dict):
237+
return
238+
turn = candidate
239+
if _turn_status(turn) in _TERMINAL_TURN_STATUSES:
240+
break
241+
if attempt + 1 < _TURN_SETTLE_ATTEMPTS:
242+
await asyncio.sleep(_TURN_SETTLE_SECONDS)
243+
if turn is None:
244+
return
245+
status = _turn_status(turn)
246+
if accepted:
247+
turn = {**turn, "status": "completed"}
248+
elif status not in _TERMINAL_TURN_STATUSES:
249+
# 既没拿到结果、这一轮又还在跑:没有可以报的终态,不编一个。
250+
return
251+
if "durationMs" not in turn:
252+
# 回合自己的时间戳就是权威值,缺 durationMs 时由它俩相减得出。
253+
started, completed = turn.get("startedAt"), turn.get("completedAt")
254+
if (
255+
isinstance(started, (int, float))
256+
and not isinstance(started, bool)
257+
and isinstance(completed, (int, float))
258+
and not isinstance(completed, bool)
259+
and completed >= started
260+
):
261+
turn = {**turn, "durationMs": completed - started}
262+
event_sink(lifecycle("turn_completed", turn))
263+
except Exception as error: # noqa: BLE001 - 读数是装饰,不能改变回合结果
264+
logger.warning(
265+
"Studio migration turn settlement failed error_type=%s",
266+
type(error).__name__,
267+
)

0 commit comments

Comments
 (0)