3232import json
3333import os
3434import shlex
35+ import subprocess
3536import sys
37+ import time
3638from pathlib import Path
39+ from collections .abc import Mapping
3740from dataclasses import dataclass
3841from typing import IO , Any , Callable , Iterable
3942
4447_DEFAULT_REAL_CODEX = "/usr/local/libexec/codex-real"
4548_DEFAULT_STATE_PATH = "/tmp/studio-codex-shim-state.json"
4649
50+ # The migration CLI settles an attempt with its own deterministic contract and opens
51+ # another attempt when a finding blocks that contract, which is why a migration turn
52+ # could claim success and still be followed by a second one. The shim runs the same
53+ # contract inside the exec it already owns and hands the blocking findings back to the
54+ # same thread, so the repair lands in the turn that made the claim.
55+ _CONTRACT_OUTPUT_ENV = "AGENTKIT_MIGRATE_OUTPUT_DIR"
56+ _CONTRACT_ASSET_ENV = "AGENTKIT_MIGRATE_ASSET_DIR"
57+ _CONTRACT_SKILL_ENV = "AGENTKIT_MIGRATE_SKILL_PATH"
58+ _CONTRACT_SCRIPT = "scripts/validate_runtime.sh"
59+ _CONTRACT_SKILL_DIR = "source-to-veadk"
60+ _DEFAULT_SKILL_PATH = "/home/gem/.codex/skills"
61+ _CONTRACT_ROW_NAME = "确定性校验"
62+ _CONTRACT_JUDGED_MARKER = "Validation finished:"
63+ _MAX_CONTRACT_REPAIRS = 2
64+ _CONTRACT_TIMEOUT_SECONDS = 300.0
65+ _CONTRACT_OUTPUT_CHARS = 4_000
66+
4767# Codex' own labels, kept identical to the ones the app-server driver writes so a
4868# migration turn reads like an intelligent-build turn.
4969_COMMAND_NAME = "运行命令"
@@ -176,6 +196,163 @@ def sandbox_policy(invocation: Invocation) -> dict[str, object]:
176196 return {"type" : "readOnly" }
177197
178198
199+ @dataclass (frozen = True )
200+ class ContractVerdict :
201+ """One run of the CLI's deterministic contract, as the shim read it."""
202+
203+ passed : bool
204+ judged : bool
205+ command : str
206+ output : str
207+ exit_code : int
208+ duration_ms : int
209+ findings : tuple [str , ...] = ()
210+
211+
212+ def contract_target (environ : Mapping [str , str ]) -> tuple [str , str ] | None :
213+ """The CLI's validation script and the output directory it judges, if any.
214+
215+ The CLI exports both to the Codex process whose prompt tells the model to run
216+ that script, so the shim reads the same two variables instead of guessing at the
217+ layout. A CLI that exports neither (or a script that is not there) keeps its own
218+ attempt loop: the shim only ever adds a check it can actually run.
219+ """
220+ output = str (environ .get (_CONTRACT_OUTPUT_ENV ) or "" ).strip ()
221+ asset = str (environ .get (_CONTRACT_ASSET_ENV ) or "" ).strip ()
222+ if not asset :
223+ skill = (
224+ str (environ .get (_CONTRACT_SKILL_ENV ) or "" ).strip () or _DEFAULT_SKILL_PATH
225+ )
226+ asset = os .path .join (skill , _CONTRACT_SKILL_DIR )
227+ if not output or not asset or not os .path .isdir (output ):
228+ return None
229+ script = os .path .join (asset , _CONTRACT_SCRIPT )
230+ if not os .path .isfile (script ):
231+ return None
232+ return script , output
233+
234+
235+ def blocking_findings (output_dir : str , * , limit : int = 6 ) -> tuple [str , ...]:
236+ """The fatal and repairable findings one validation run left behind."""
237+ try :
238+ with open (
239+ os .path .join (output_dir , "validation_findings.json" ), encoding = "utf-8"
240+ ) as handle :
241+ value = json .load (handle )
242+ except (OSError , ValueError ):
243+ return ()
244+ if not isinstance (value , dict ):
245+ return ()
246+ lines : list [str ] = []
247+ for severity in ("fatal" , "repairable" ):
248+ entries = value .get (severity )
249+ if not isinstance (entries , list ):
250+ continue
251+ for entry in entries :
252+ if not isinstance (entry , dict ):
253+ continue
254+ name = str (entry .get ("name" ) or "" ).strip ()
255+ detail = str (entry .get ("detail" ) or "" ).strip ()
256+ if name or detail :
257+ lines .append (f"- { name } [{ severity } ] { detail } " .strip ())
258+ if len (lines ) >= limit :
259+ return tuple (lines )
260+ return tuple (lines )
261+
262+
263+ def _tail (text : str , limit : int ) -> str :
264+ """The end of a command's output, which is where a validator puts its verdict."""
265+ text = text .strip ()
266+ return text if len (text ) <= limit else text [- limit :]
267+
268+
269+ def run_contract (script : str , output_dir : str ) -> ContractVerdict :
270+ """Run the CLI's deterministic contract the way the migration prompt does.
271+
272+ A verdict counts only when the validator reported one: a script that could not
273+ run at all must not cost the CLI a turn, so ``judged`` gates the repair loop.
274+ """
275+ command = 'bash "$AGENTKIT_MIGRATE_ASSET_DIR/scripts/validate_runtime.sh"'
276+ started = time .monotonic ()
277+ try :
278+ completed = subprocess .run (
279+ ["bash" , script ],
280+ cwd = output_dir ,
281+ capture_output = True ,
282+ text = True ,
283+ timeout = _CONTRACT_TIMEOUT_SECONDS ,
284+ check = False ,
285+ )
286+ except (OSError , subprocess .SubprocessError ) as error :
287+ return ContractVerdict (
288+ passed = False ,
289+ judged = False ,
290+ command = command ,
291+ output = str (error ),
292+ exit_code = - 1 ,
293+ duration_ms = int ((time .monotonic () - started ) * 1000 ),
294+ )
295+ text = "\n " .join (part for part in (completed .stdout , completed .stderr ) if part )
296+ judged = _CONTRACT_JUDGED_MARKER in (completed .stdout or "" )
297+ return ContractVerdict (
298+ passed = judged and completed .returncode == 0 ,
299+ judged = judged ,
300+ command = command ,
301+ output = _tail (text , _CONTRACT_OUTPUT_CHARS ),
302+ exit_code = completed .returncode ,
303+ duration_ms = int ((time .monotonic () - started ) * 1000 ),
304+ findings = blocking_findings (output_dir ),
305+ )
306+
307+
308+ def contract_row (verdict : ContractVerdict , * , index : int ) -> dict [str , object ]:
309+ """One validation run as the command row the migration page already draws."""
310+ return {
311+ "type" : "item.completed" ,
312+ "item" : {
313+ "id" : f"studio-contract-{ index } " ,
314+ "type" : "command_execution" ,
315+ "name" : _CONTRACT_ROW_NAME ,
316+ "command" : verdict .command ,
317+ "aggregated_output" : verdict .output ,
318+ "exit_code" : verdict .exit_code ,
319+ "duration_ms" : verdict .duration_ms ,
320+ "status" : "completed" if verdict .passed else "failed" ,
321+ },
322+ }
323+
324+
325+ def contract_feedback (verdict : ContractVerdict ) -> str :
326+ """The repair instructions handed back into the same turn."""
327+ findings = "\n " .join (verdict .findings ) or "- 见校验输出。"
328+ return "\n " .join (
329+ [
330+ "# 确定性校验未通过:在本回合内修复" ,
331+ "" ,
332+ "CLI 的迁移契约刚刚在这个输出目录上失败,所以这次迁移还不能结束。" ,
333+ "不要重开迁移,也不要改写 CLI 初始化生成的 `.agentkit/agentkit.yaml`:" ,
334+ "它的 sha256 就是 `migration_metadata.json` 里记录的配置基线," ,
335+ "应用名由已确认的迁移设置决定,不是本回合可以更改的内容。" ,
336+ "在当前输出目录里修掉下面的阻断项,然后重跑校验,直到它通过。" ,
337+ "" ,
338+ "## 阻断项" ,
339+ findings ,
340+ "" ,
341+ "## 校验输出(末尾)" ,
342+ "```" ,
343+ verdict .output or "(校验脚本没有输出)" ,
344+ "```" ,
345+ "" ,
346+ "## 完成条件" ,
347+ '- 重跑 `bash "$AGENTKIT_MIGRATE_ASSET_DIR/scripts/validate_runtime.sh"`,' ,
348+ " 退出码为 0,且 `validation_findings.json` 的 `fatal`、`repairable` 都为空" ,
349+ " (`degraded` 可以保留,但要在报告里如实说明)。" ,
350+ "- 没有通过校验之前,不要输出迁移完成的结论。" ,
351+ "" ,
352+ ]
353+ )
354+
355+
179356def read_output_schema (path : str ) -> object :
180357 """The JSON Schema the CLI pinned the model's final message to, if readable."""
181358 if not path :
@@ -403,6 +580,32 @@ def add_usage(left: dict[str, int], right: dict[str, int]) -> dict[str, int]:
403580 return total
404581
405582
583+ def _second (value : object ) -> float | None :
584+ """An app-server timestamp in seconds, or ``None`` when it is not one."""
585+ if isinstance (value , bool ) or not isinstance (value , (int , float )):
586+ return None
587+ number = float (value )
588+ return number if number >= 0 else None
589+
590+
591+ def span_turns (turns : list [dict [str , object ]], turn : dict [str , object ]) -> None :
592+ """Report a whole exec as one turn when the contract kept it open.
593+
594+ A repaired contract makes one `codex exec` carry several app-server turns, but
595+ the migration page draws one turn per exec, so its elapsed time has to cover the
596+ model work of every sub-turn and the validation between them.
597+ """
598+ if len (turns ) < 2 :
599+ return
600+ started = _second (turns [0 ].get ("startedAt" ))
601+ completed = _second (turns [- 1 ].get ("completedAt" ))
602+ if started is None or completed is None or completed < started :
603+ return
604+ turn ["startedAt" ] = turns [0 ]["startedAt" ]
605+ turn ["completedAt" ] = turns [- 1 ]["completedAt" ]
606+ turn ["durationMs" ] = int (round ((completed - started ) * 1000 ))
607+
608+
406609def turn_line (
407610 turn : dict [str , object ],
408611 * ,
@@ -474,6 +677,7 @@ def __init__(self, socket: Any, emit: Callable[[dict[str, object]], None]) -> No
474677 self .failure = ""
475678 self ._completed = asyncio .Event ()
476679 self ._turn : dict [str , object ] = {}
680+ self ._turns : list [dict [str , object ]] = []
477681 self ._usage : dict [str , int ] = {}
478682 self ._thread_total : dict [str , int ] | None = None
479683
@@ -499,6 +703,10 @@ async def start(self) -> asyncio.Task[None]:
499703 """Start reading this connection: every request needs the reader running."""
500704 return asyncio .create_task (self ._read ())
501705
706+ def begin_turn (self ) -> None :
707+ """Arm the connection for another turn on the thread it already holds."""
708+ self ._completed = asyncio .Event ()
709+
502710 async def run_turn (
503711 self ,
504712 * ,
@@ -618,6 +826,7 @@ def _notification(self, method: str, params: object) -> None:
618826 turn = payload .get ("turn" )
619827 if isinstance (turn , dict ):
620828 self ._turn = turn
829+ self ._turns .append (dict (turn ))
621830 raw_status = turn .get ("status" )
622831 if isinstance (raw_status , dict ):
623832 raw_status = raw_status .get ("type" )
@@ -668,6 +877,7 @@ def summary_line(self, model: str) -> dict[str, object]:
668877 reported = turn .get ("model" )
669878 if not model and isinstance (reported , str ):
670879 model = reported
880+ span_turns (self ._turns , turn )
671881 return turn_line (turn , usage = self ._usage , model = model or self .model )
672882
673883
@@ -787,6 +997,7 @@ async def _drive_turn(
787997 sandbox = sandbox_policy (invocation ),
788998 output_schema = read_output_schema (invocation .output_schema_path ),
789999 )
1000+ await settle_contract (turn , invocation , emit = emit )
7901001 emit (turn .summary_line (invocation .model ))
7911002 if invocation .last_message_path and turn .final_text :
7921003 try :
@@ -799,6 +1010,46 @@ async def _drive_turn(
7991010 return 0 if turn .status in {"" , "completed" } else 1
8001011
8011012
1013+ async def settle_contract (
1014+ turn : _AppServerTurn ,
1015+ invocation : Invocation ,
1016+ * ,
1017+ emit : Callable [[dict [str , object ]], None ],
1018+ ) -> None :
1019+ """Hold one exec inside a single turn until the CLI's own contract passes.
1020+
1021+ The migration CLI validates the output after Codex exits and opens a new attempt
1022+ when a finding blocks the delivery, so a turn could claim the migration was done
1023+ and still be followed by another one. The contract is the deterministic script
1024+ the migration prompt already tells the model to run, so the shim runs that same
1025+ script inside the exec and hands the blocking findings back into the same thread:
1026+ the repair lands in the turn that made the claim, and the CLI's attempt loop
1027+ stays a backstop.
1028+
1029+ Every bail-out is deliberate. A contract the shim cannot run or cannot read a
1030+ verdict from, a sub-turn that failed, and the repair budget all end the loop
1031+ without touching the turn, because a turn must never be held open by the shim.
1032+ """
1033+ target = contract_target (os .environ )
1034+ if target is None or turn .status not in {"" , "completed" }:
1035+ return
1036+ script , output_dir = target
1037+ for repair in range (_MAX_CONTRACT_REPAIRS + 1 ):
1038+ verdict = await asyncio .to_thread (run_contract , script , output_dir )
1039+ emit (contract_row (verdict , index = repair + 1 ))
1040+ if verdict .passed or not verdict .judged or repair == _MAX_CONTRACT_REPAIRS :
1041+ return
1042+ turn .begin_turn ()
1043+ await turn .run_turn (
1044+ thread_id = turn .turn_id ,
1045+ prompt = contract_feedback (verdict ),
1046+ model = invocation .model ,
1047+ sandbox = sandbox_policy (invocation ),
1048+ )
1049+ if turn .status not in {"" , "completed" }:
1050+ return
1051+
1052+
8021053def read_state (path : str ) -> dict [str , object ]:
8031054 try :
8041055 with open (path , encoding = "utf-8" ) as handle :
0 commit comments