1
0

smoke-python-runtime.py 107 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509
  1. #!/usr/bin/env python3
  2. """Keyless full-turn and snapshot smoke for the Python SDK runtime."""
  3. from __future__ import annotations
  4. import argparse
  5. import difflib
  6. import importlib
  7. import importlib.metadata
  8. import json
  9. import os
  10. import queue
  11. import re
  12. import secrets
  13. import shutil
  14. import subprocess
  15. import sys
  16. import sysconfig
  17. import tempfile
  18. import threading
  19. import time
  20. from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
  21. from pathlib import Path
  22. from typing import TYPE_CHECKING, Callable
  23. if TYPE_CHECKING:
  24. from deepseek_harness import RunResult
  25. EXPECTED_TEXT = "runtime smoke ok"
  26. LIVE_API_SENTINEL = "PYTHON_SDK_LIVE_OK"
  27. CODE_PROMPT = "Use run_code to compute the packaged worker smoke value."
  28. CODE_WORKER_TEXT = "code worker smoke ok"
  29. WORKFLOW_PROMPT = "Use workflow to compute the packaged worker smoke value without agents."
  30. WORKFLOW_WORKER_TEXT = "workflow worker smoke ok"
  31. MINIMAL_PROMPT = "Exercise the packaged minimal agent's persistent shell."
  32. MINIMAL_TEXT = "minimal agent smoke ok"
  33. FS_SEARCH_PROMPT = "Exercise the packaged filesystem search tools."
  34. FS_SEARCH_TEXT = "filesystem search smoke ok"
  35. FS_SEARCH_MARKER = "PACKAGED_FS_SEARCH_OK"
  36. MCP_PROMPT = "Exercise the packaged MCP client with one external stdio server."
  37. MCP_TEXT = "MCP client smoke ok"
  38. PROFILE_PLUGIN_PROMPT = "Verify the Python-installed dsh profile plugin."
  39. PROFILE_PLUGIN_TEXT = "profile plugin smoke ok"
  40. PROFILE_PLUGIN_MARKER = "PYTHON_INSTALLED_DSH_PROFILE_PLUGIN"
  41. IS_WINDOWS = sys.platform == "win32"
  42. MINIMAL_SHELL_TOOL = "pwsh" if IS_WINDOWS else "bash"
  43. MINIMAL_SHELL_COMMAND = (
  44. "$global:dshSdkCounter = [int]$global:dshSdkCounter + 1; "
  45. 'Write-Output "COUNT=$global:dshSdkCounter CWD=$((Get-Location).Path)"; '
  46. "if ($global:dshSdkCounter -eq 1) { Set-Location $env:TEMP }"
  47. if IS_WINDOWS
  48. else (
  49. "counter=$(( ${counter:-0} + 1 )); export counter; "
  50. "printf 'COUNT=%s CWD=%s\\n' \"$counter\" \"$PWD\"; "
  51. "if [ \"$counter\" -eq 1 ]; then cd /tmp; fi"
  52. )
  53. )
  54. MINIMAL_SHELL_SECOND_CWD = str(Path(tempfile.gettempdir()).resolve()) if IS_WINDOWS else "/tmp"
  55. SPAWN_NODE_PROMPT = "Run node --version through the packaged shell tool."
  56. SPAWN_NODE_TEXT = "spawn node smoke ok"
  57. SPAWN_NODE_CALL_ID = "spawn-node-shell"
  58. # The POSIX command string starts with `node ` inside the shell tool's `bash -c`
  59. # argv, the exact form @yao-pkg/pkg's unpatched SEA bootstrap rewrites to the
  60. # executable itself while stamping PKG_EXECPATH into the child environment.
  61. SPAWN_NODE_COMMAND = (
  62. 'node --version; if ($env:PKG_EXECPATH) { "PKG_EXECPATH=$env:PKG_EXECPATH" } else { "PKG_EXECPATH=ABSENT" }'
  63. if IS_WINDOWS
  64. else 'node --version; echo "PKG_EXECPATH=${PKG_EXECPATH:-ABSENT}"'
  65. )
  66. LEGACY_CUSTOM_DISABLED_ROWS = (
  67. "agent-instructions",
  68. "goal",
  69. "goal-round-driver",
  70. "command-goal",
  71. "plan-mode",
  72. "skill",
  73. "skill-filesystem",
  74. "tool-fs",
  75. "tool-fs-search",
  76. "tool-goal",
  77. "tool-ralph",
  78. "tool-skill",
  79. "tool-str-replace-editor",
  80. "tool-subagent-control",
  81. "tool-subagent-list-agents",
  82. "tool-subagent-fork",
  83. "tool-todo",
  84. "tool-web",
  85. )
  86. SNAPSHOT_PROMPT = "Run the advanced packaged-runtime snapshot scenario."
  87. SNAPSHOT_SESSION_ID = "advanced-executable"
  88. SNAPSHOT_DIRECT_CHILD_PROMPT = "Reply with exactly DIRECT_CHILD_OK and nothing else."
  89. SNAPSHOT_WORKFLOW_CHILD_PROMPT = "Reply with exactly WORKFLOW_CHILD_OK and nothing else."
  90. SNAPSHOT_FINAL_TEXT = "ADVANCED_EXECUTABLE_OK"
  91. RESTART_FIRST_PROMPT = "Complete the first isolated Python SDK process turn."
  92. RESTART_FIRST_TEXT = "PROCESS_ONE_OK"
  93. RESTART_SECOND_PROMPT = "Complete the second isolated Python SDK process turn."
  94. RESTART_SECOND_TEXT = "PROCESS_TWO_OK"
  95. RESTART_FIRST_SESSION_ID = "process-one"
  96. RESTART_SECOND_SESSION_ID = "process-two"
  97. SNAPSHOT_PLUGIN_CODE = """\
  98. return (ctx) => {
  99. ctx.on('tools/pre-execute', (exec, next) => {
  100. if (exec.name !== 'snapshot_double' || exec.arguments.value !== -1) return next()
  101. return {
  102. kind: 'deny',
  103. reason: 'Auto review rejected tool "snapshot_double"; its body was not executed',
  104. info: {
  105. name: 'AutoReviewDeniedError', code: 'AUTO_REVIEW_DENIED',
  106. reason: ' transport raw\\r\\nreason ',
  107. },
  108. }
  109. })
  110. harness.registerTool(ctx, harness.defineTool({
  111. name: 'snapshot_double',
  112. description: 'Double a number for executable snapshot verification.',
  113. parameters: { value: { type: 'number', required: true } },
  114. output: {
  115. schema: { type: 'number' },
  116. render(_args, value) {
  117. return [{ type: 'text', text: String(value) }]
  118. }
  119. },
  120. async execute(args) {
  121. return args.value * 2
  122. }
  123. }))
  124. }
  125. """
  126. SNAPSHOT_WORKFLOW_SCRIPT = (
  127. "phase('Delegate')\n"
  128. f"const reply = await agent('{SNAPSHOT_WORKFLOW_CHILD_PROMPT}', {{ label: 'workflow-child' }})\n"
  129. "return { reply }"
  130. )
  131. ADVANCED_SNAPSHOT_DIRECTORY = (
  132. Path(__file__).resolve().parent / "snapshots" / "python-sdk-single-exe" / "advanced"
  133. )
  134. ADVANCED_SNAPSHOT_FILENAMES = (
  135. "result.json", "session.v3.jsonl", "session.1.v3.jsonl", "session.2.v3.jsonl",
  136. )
  137. MINIMAL_SNAPSHOT_DIRECTORY = (
  138. Path(__file__).resolve().parent / "snapshots" / "python-sdk-single-exe" / "minimal"
  139. )
  140. if IS_WINDOWS:
  141. MINIMAL_SNAPSHOT_DIRECTORY /= "win-x64"
  142. MINIMAL_SNAPSHOT_FILENAMES = ("model-visible.json",)
  143. IN_HISTORY_SNAPSHOT_DIRECTORY = (
  144. Path(__file__).resolve().parent / "snapshots" / "python-sdk-single-exe" / "minimal-in-history"
  145. )
  146. IN_HISTORY_SNAPSHOT_FILENAMES = ("prompt-history.json",)
  147. RESTART_SNAPSHOT_DIRECTORY = (
  148. Path(__file__).resolve().parent / "snapshots" / "python-sdk-single-exe" / "restart"
  149. )
  150. RESTART_SNAPSHOT_FILENAMES = (
  151. "result.json", "requests.json", "session.1.v3.jsonl", "session.2.v3.jsonl",
  152. )
  153. MCP_SERVER_SCRIPT = """\
  154. import json
  155. import os
  156. import sys
  157. import time
  158. log_path = os.environ.get("MCP_SMOKE_LOG")
  159. def send(message):
  160. sys.stdout.write(json.dumps(message, separators=(",", ":")) + "\\n")
  161. sys.stdout.flush()
  162. for line in sys.stdin:
  163. request = json.loads(line)
  164. if log_path is not None:
  165. with open(log_path, "a", encoding="utf-8") as log:
  166. log.write(str(request.get("method")) + "\\n")
  167. request_id = request.get("id")
  168. if request_id is None:
  169. continue
  170. method = request.get("method")
  171. if method == "initialize":
  172. send({
  173. "jsonrpc": "2.0",
  174. "id": request_id,
  175. "result": {
  176. "protocolVersion": request["params"]["protocolVersion"],
  177. "capabilities": {"tools": {"listChanged": False}},
  178. "serverInfo": {"name": "python-wheel-fixture", "version": "1.0.0"},
  179. },
  180. })
  181. elif method == "tools/list":
  182. # Keep discovery pending long enough that an SDK runtime answering
  183. # initialize before discovery completes makes its first model request
  184. # without this tool and fails deterministically.
  185. time.sleep(0.25)
  186. send({
  187. "jsonrpc": "2.0",
  188. "id": request_id,
  189. "result": {
  190. "tools": [{
  191. "name": "add",
  192. "description": "Add two numbers.",
  193. "inputSchema": {
  194. "type": "object",
  195. "properties": {"a": {"type": "number"}, "b": {"type": "number"}},
  196. "required": ["a", "b"],
  197. "additionalProperties": False,
  198. },
  199. }],
  200. },
  201. })
  202. elif method == "tools/call":
  203. params = request["params"]
  204. if params.get("name") != "add" or params.get("arguments") != {"a": 19, "b": 23}:
  205. send({
  206. "jsonrpc": "2.0",
  207. "id": request_id,
  208. "error": {"code": -32602, "message": "unexpected tool call"},
  209. })
  210. continue
  211. send({
  212. "jsonrpc": "2.0",
  213. "id": request_id,
  214. "result": {"content": [{"type": "text", "text": "42"}]},
  215. })
  216. else:
  217. send({
  218. "jsonrpc": "2.0",
  219. "id": request_id,
  220. "error": {"code": -32601, "message": f"unsupported method: {method}"},
  221. })
  222. """
  223. def write_profile_patch(
  224. root: Path,
  225. name: str,
  226. sessions: Path,
  227. patches: list[dict[str, object]],
  228. ) -> Path:
  229. """Write one JSON-form dsh profile patch with deterministic persistence."""
  230. path = root / name
  231. path.write_text(json.dumps([
  232. {
  233. "id": "session-persistence-jsonl",
  234. "config": {"root": str(sessions), "compression": "none"},
  235. },
  236. {"id": "session-telemetry-otel", "disabled": True},
  237. *patches,
  238. ], indent=2))
  239. return path
  240. def write_advanced_profile_patch(root: Path, name: str, sessions: Path) -> Path:
  241. """Write the shared custom, snapshot, and restart profile patch."""
  242. return write_profile_patch(root, name, sessions, [
  243. {"id": "tools", "config": {"mode": "both"}},
  244. {
  245. "id": "system-prompt",
  246. "config": {
  247. "persona": "You are a coding agent powered by the {{model}} model. Your working directory is {{cwd}}.",
  248. },
  249. },
  250. {"id": "session-log-deepseek", "config": {"enabled": True}},
  251. *({"id": row_id, "disabled": True} for row_id in LEGACY_CUSTOM_DISABLED_ROWS),
  252. {"id": "tool-bash", "disabled": True},
  253. {"id": "tool-pwsh", "disabled": True},
  254. {
  255. "id": "tool-subagent",
  256. "config": {
  257. "provider": "spawn",
  258. "toolName": "subagent",
  259. "backgroundMode": "one-shot",
  260. },
  261. },
  262. {"insert": [
  263. {"id": "ptc-runtime", "name": "@deepseek-ai/dsh-ptc-runtime-node"},
  264. {"id": "cordis-host-runner", "name": "@deepseek-ai/dsh-cordis-host-runner"},
  265. {"id": "cordis-tool", "name": "@deepseek-ai/dsh-tool-cordis"},
  266. ]},
  267. ])
  268. def write_mcp_patch(root: Path, sessions: Path, server_script: Path) -> Path:
  269. """Write a profile patch that mounts the packaged MCP client."""
  270. return write_profile_patch(root, "mcp.patch.yml", sessions, [{
  271. "insert": [{
  272. "id": "mcp-fixture",
  273. "name": "@deepseek-ai/dsh-mcp-client",
  274. "config": {
  275. "serverName": "fixture",
  276. "transport": "stdio",
  277. "command": sys.executable,
  278. "args": [str(server_script)],
  279. "env": {"MCP_SMOKE_LOG": str(server_script.with_suffix(".log"))},
  280. "failOnStartupError": True,
  281. "reconnect": {"enabled": False},
  282. },
  283. }],
  284. }])
  285. class MockModelHandler(BaseHTTPRequestHandler):
  286. """Return deterministic text, worker, and orchestration completions."""
  287. requests: list[dict[str, object]] = []
  288. def do_POST(self) -> None:
  289. if self.path != "/v1/messages":
  290. self.send_error(404)
  291. return
  292. content_length = int(self.headers.get("content-length", "0"))
  293. body = json.loads(self.rfile.read(content_length))
  294. self.requests.append(body)
  295. self.send_response(200)
  296. self.send_header("content-type", "text/event-stream")
  297. self.end_headers()
  298. chunks = completion_chunks(body)
  299. for chunk in chunks:
  300. self.wfile.write(f"event: {chunk['type']}\ndata: {json.dumps(chunk)}\n\n".encode())
  301. self.wfile.flush()
  302. def log_message(self, _format: str, *_args: object) -> None:
  303. return
  304. def completion_chunks(body: dict[str, object]) -> list[dict[str, object]]:
  305. """Choose the next deterministic model response from request history."""
  306. messages = body.get("messages")
  307. if not isinstance(messages, list) or not messages:
  308. raise AssertionError(f"model request has no messages: {body}")
  309. # A system prompt update may follow the tool result without replacing it.
  310. latest = next(message for message in reversed(messages) if message.get("role") != "system")
  311. if not isinstance(latest, dict):
  312. raise AssertionError(f"model request has an invalid latest message: {body}")
  313. tool_results = [
  314. block for block in latest.get("content", [])
  315. if isinstance(block, dict) and block.get("type") == "tool_result"
  316. ]
  317. if tool_results:
  318. result = tool_results[-1]
  319. call_id, tool_name = latest_tool_call(messages, result.get("tool_use_id"))
  320. tool_text = message_text(result.get("content"))
  321. mcp = mcp_tool_followup(call_id, tool_name, tool_text)
  322. if mcp is not None:
  323. return mcp
  324. fs_search = fs_search_tool_followup(call_id, tool_name, tool_text)
  325. if fs_search is not None:
  326. return fs_search
  327. spawn_node = spawn_node_tool_followup(call_id, tool_name, tool_text)
  328. if spawn_node is not None:
  329. return spawn_node
  330. minimal = minimal_tool_followup(call_id, tool_name, tool_text)
  331. if minimal is not None:
  332. return minimal
  333. advanced = advanced_tool_followup(body, call_id, tool_name, tool_text)
  334. if advanced is not None:
  335. return advanced
  336. if "42" not in tool_text:
  337. raise AssertionError(f"{tool_name} worker returned no expected value: {latest}")
  338. if tool_name == "run_code":
  339. return text_chunks(CODE_WORKER_TEXT)
  340. if tool_name == "workflow":
  341. return text_chunks(WORKFLOW_WORKER_TEXT)
  342. raise AssertionError(f"unexpected tool follow-up: {tool_name}")
  343. user_prompts = [
  344. block["text"]
  345. for message in reversed(messages)
  346. if isinstance(message, dict) and message.get("role") == "user"
  347. for block in message.get("content", [])
  348. if isinstance(block, dict) and block.get("type") == "text"
  349. ]
  350. minimal_prompt = next((prompt for prompt in user_prompts if prompt == MINIMAL_PROMPT), None)
  351. # The minimal composition's assembled system prompt, advertised tool schemas, and
  352. # model-visible messages are pinned by its snapshot, not asserted here.
  353. if minimal_prompt is not None:
  354. return tool_call_chunks(
  355. "minimal-bash-1",
  356. MINIMAL_SHELL_TOOL,
  357. {"command": MINIMAL_SHELL_COMMAND},
  358. )
  359. scenario_prompts = {
  360. SNAPSHOT_DIRECT_CHILD_PROMPT,
  361. SNAPSHOT_WORKFLOW_CHILD_PROMPT,
  362. SNAPSHOT_PROMPT,
  363. CODE_PROMPT,
  364. WORKFLOW_PROMPT,
  365. FS_SEARCH_PROMPT,
  366. SPAWN_NODE_PROMPT,
  367. MCP_PROMPT,
  368. RESTART_FIRST_PROMPT,
  369. RESTART_SECOND_PROMPT,
  370. PROFILE_PLUGIN_PROMPT,
  371. }
  372. prompt = next(
  373. (candidate for candidate in user_prompts if candidate in scenario_prompts),
  374. message_text(latest.get("content")),
  375. )
  376. if prompt == SNAPSHOT_DIRECT_CHILD_PROMPT:
  377. return text_chunks("DIRECT_CHILD_OK")
  378. if prompt == SNAPSHOT_WORKFLOW_CHILD_PROMPT:
  379. return text_chunks("WORKFLOW_CHILD_OK")
  380. if prompt == SNAPSHOT_PROMPT:
  381. assert_advertised_tool(body, "cordis_define")
  382. return tool_call_chunks(
  383. "advanced-define",
  384. "cordis_define",
  385. {
  386. "plugin": {"kind": "new", "idPrefix": "snap"},
  387. "name": "Snapshot Double",
  388. "purpose": "Expose a deterministic doubling tool for executable snapshot verification.",
  389. "code": {"host": SNAPSHOT_PLUGIN_CODE},
  390. },
  391. )
  392. if prompt == RESTART_FIRST_PROMPT:
  393. return text_chunks(RESTART_FIRST_TEXT)
  394. if prompt == RESTART_SECOND_PROMPT:
  395. if any(
  396. isinstance(message, dict)
  397. and RESTART_FIRST_TEXT in message_text(message.get("content"))
  398. for message in messages
  399. ):
  400. raise AssertionError("second isolated process inherited the first process history")
  401. return text_chunks(RESTART_SECOND_TEXT)
  402. if prompt == CODE_PROMPT:
  403. assert_advertised_tool(body, "run_code")
  404. return tool_call_chunks(
  405. "call-code-worker",
  406. "run_code",
  407. {"code": "return 6 * 7", "description": "Compute the smoke value"},
  408. )
  409. if prompt == WORKFLOW_PROMPT:
  410. assert_advertised_tool(body, "workflow")
  411. return tool_call_chunks(
  412. "call-workflow-worker",
  413. "workflow",
  414. {
  415. "script": "return 6 * 7",
  416. "meta": {
  417. "name": "pkg-worker-smoke",
  418. "description": "exercise the packaged workflow worker",
  419. },
  420. },
  421. )
  422. if prompt == FS_SEARCH_PROMPT:
  423. assert_advertised_tool(body, "grep")
  424. assert_advertised_tool(body, "glob")
  425. return tool_call_chunks(
  426. "fs-search-grep",
  427. "grep",
  428. {"pattern": FS_SEARCH_MARKER, "path": "."},
  429. )
  430. if prompt == SPAWN_NODE_PROMPT:
  431. assert_advertised_tool(body, MINIMAL_SHELL_TOOL)
  432. return tool_call_chunks(
  433. SPAWN_NODE_CALL_ID,
  434. MINIMAL_SHELL_TOOL,
  435. {"command": SPAWN_NODE_COMMAND, "description": "Report the reachable Node version"},
  436. )
  437. if prompt == MCP_PROMPT:
  438. assert_advertised_tool(body, "mcp__fixture__add")
  439. return tool_call_chunks(
  440. "mcp-add",
  441. "mcp__fixture__add",
  442. {"a": 19, "b": 23},
  443. )
  444. if prompt == PROFILE_PLUGIN_PROMPT:
  445. system_text = message_text(body.get("system")) + "\n" + "\n".join(
  446. message_text(message.get("content"))
  447. for message in messages
  448. if isinstance(message, dict) and message.get("role") == "system"
  449. )
  450. if PROFILE_PLUGIN_MARKER not in system_text:
  451. raise AssertionError("external profile plugin contributed no model-visible marker")
  452. return text_chunks(PROFILE_PLUGIN_TEXT)
  453. return text_chunks(EXPECTED_TEXT)
  454. def mcp_tool_followup(
  455. call_id: str,
  456. tool_name: str,
  457. tool_text: str,
  458. ) -> list[dict[str, object]] | None:
  459. """Verify one tool call through the packaged MCP client."""
  460. if call_id != "mcp-add":
  461. return None
  462. if tool_name != "mcp__fixture__add" or "42" not in tool_text:
  463. raise AssertionError(f"packaged MCP call returned an unexpected result: {tool_name}: {tool_text}")
  464. return text_chunks(MCP_TEXT)
  465. def fs_search_tool_followup(
  466. call_id: str,
  467. tool_name: str,
  468. tool_text: str,
  469. ) -> list[dict[str, object]] | None:
  470. """Exercise both ripgrep-backed tools through the packaged executable."""
  471. if not call_id.startswith("fs-search-"):
  472. return None
  473. if call_id == "fs-search-grep" and tool_name == "grep":
  474. if "needle.txt" not in tool_text or FS_SEARCH_MARKER not in tool_text:
  475. raise AssertionError(f"packaged grep returned no marker: {tool_text}")
  476. return tool_call_chunks(
  477. "fs-search-glob",
  478. "glob",
  479. {"pattern": "**/*.txt"},
  480. )
  481. if call_id == "fs-search-glob" and tool_name == "glob":
  482. if "needle.txt" not in tool_text:
  483. raise AssertionError(f"packaged glob returned no fixture path: {tool_text}")
  484. return text_chunks(FS_SEARCH_TEXT)
  485. raise AssertionError(f"unexpected filesystem-search follow-up: {call_id} {tool_name}: {tool_text}")
  486. def host_node_version() -> str:
  487. """The machine's own `node --version` line, the required shell resolution target."""
  488. node = shutil.which("node")
  489. if node is None:
  490. raise AssertionError("the spawn-node scenario requires Node on PATH for comparison")
  491. return subprocess.run(
  492. [node, "--version"], capture_output=True, text=True, check=True,
  493. ).stdout.strip()
  494. def spawn_node_tool_followup(
  495. call_id: str,
  496. tool_name: str,
  497. tool_text: str,
  498. ) -> list[dict[str, object]] | None:
  499. """Verify the packaged shell reached the machine's Node with a clean environment."""
  500. if call_id != SPAWN_NODE_CALL_ID:
  501. return None
  502. if tool_name != MINIMAL_SHELL_TOOL:
  503. raise AssertionError(f"spawn-node follow-up used an unexpected tool: {tool_name}")
  504. expected = host_node_version()
  505. if expected not in tool_text:
  506. raise AssertionError(
  507. f"packaged shell did not reach the machine's node {expected}: {tool_text}"
  508. )
  509. if "PKG_EXECPATH=ABSENT" not in tool_text:
  510. raise AssertionError(f"PKG_EXECPATH reached the shell child environment: {tool_text}")
  511. return text_chunks(SPAWN_NODE_TEXT)
  512. def minimal_tool_followup(
  513. call_id: str,
  514. tool_name: str,
  515. tool_text: str,
  516. ) -> list[dict[str, object]] | None:
  517. """Verify the checked-in minimal composition's persistent PTY."""
  518. if not call_id.startswith("minimal-"):
  519. return None
  520. if call_id == "minimal-bash-1" and tool_name == MINIMAL_SHELL_TOOL:
  521. if "COUNT=1" not in tool_text:
  522. raise AssertionError(f"first persistent shell call lost its output: {tool_text}")
  523. return tool_call_chunks(
  524. "minimal-bash-2",
  525. MINIMAL_SHELL_TOOL,
  526. {"command": MINIMAL_SHELL_COMMAND},
  527. )
  528. if call_id == "minimal-bash-2" and tool_name == MINIMAL_SHELL_TOOL:
  529. expected = f"COUNT=2 CWD={MINIMAL_SHELL_SECOND_CWD}"
  530. if expected.lower() not in tool_text.lower():
  531. raise AssertionError(f"persistent shell did not retain state: {tool_text}")
  532. return text_chunks(MINIMAL_TEXT)
  533. raise AssertionError(f"unexpected minimal-agent follow-up: {call_id} {tool_name}: {tool_text}")
  534. def advanced_tool_followup(
  535. body: dict[str, object],
  536. call_id: str,
  537. tool_name: str,
  538. tool_text: str,
  539. ) -> list[dict[str, object]] | None:
  540. """Advance the executable snapshot's deterministic parent tool chain."""
  541. if not call_id.startswith("advanced-"):
  542. return None
  543. if call_id == "advanced-define" and tool_name == "cordis_define":
  544. if "Defined snap-1/pkg-1 (Snapshot Double)" not in tool_text:
  545. raise AssertionError(f"cordis_define returned no dynamic Package ids: {tool_text}")
  546. if "snapshot_double" in advertised_tool_names(body):
  547. raise AssertionError("snapshot_double was advertised before cordis_run")
  548. assert_advertised_tool(body, "cordis_run")
  549. return tool_call_chunks(
  550. "advanced-run",
  551. "cordis_run",
  552. {"pluginId": "snap-1", "packageId": "pkg-1", "mode": "run"},
  553. )
  554. if call_id == "advanced-run" and tool_name == "cordis_run":
  555. if "snap-1/pkg-1 is running (run-1)" not in tool_text:
  556. raise AssertionError(f"cordis_run returned no running Package ids: {tool_text}")
  557. assert_advertised_tool(body, "run_code")
  558. assert_advertised_tool(body, "snapshot_double")
  559. return tool_call_chunks(
  560. "advanced-code",
  561. "run_code",
  562. {
  563. "code": "return await tools.snapshot_double({ value: 21 })",
  564. "description": "Run the temporary Plugin tool",
  565. },
  566. )
  567. if call_id == "advanced-code" and tool_name == "run_code":
  568. if "42" not in tool_text:
  569. raise AssertionError(f"run_code returned no dynamic-tool value: {tool_text}")
  570. return tool_call_chunks("advanced-denied-native", "snapshot_double", {"value": -1})
  571. if call_id == "advanced-denied-native" and tool_name == "snapshot_double":
  572. if 'Auto review rejected tool "snapshot_double"; its body was not executed' not in tool_text:
  573. raise AssertionError(f"native denial did not preserve the model result: {tool_text}")
  574. if "transport raw" in tool_text:
  575. raise AssertionError("native denial leaked its user-facing reason to the model")
  576. return tool_call_chunks("advanced-denied-ptc", "run_code", {
  577. "code": "try { await tools.snapshot_double({ value: -1 }) } catch (error) { return error.message }",
  578. "description": "Catch a structured inner tool denial",
  579. })
  580. if call_id == "advanced-denied-ptc" and tool_name == "run_code":
  581. if 'Auto review rejected tool "snapshot_double"; its body was not executed' not in tool_text:
  582. raise AssertionError(f"PTC denial did not preserve the model result: {tool_text}")
  583. if "transport raw" in tool_text:
  584. raise AssertionError("PTC denial leaked its user-facing reason to the model")
  585. assert_advertised_tool(body, "subagent")
  586. return tool_call_chunks(
  587. "advanced-direct-child",
  588. "subagent",
  589. {
  590. "description": "Check direct child",
  591. "prompt": SNAPSHOT_DIRECT_CHILD_PROMPT,
  592. },
  593. )
  594. if call_id == "advanced-direct-child" and tool_name == "subagent":
  595. if "DIRECT_CHILD_OK" not in tool_text:
  596. raise AssertionError(f"subagent returned no expected child value: {tool_text}")
  597. assert_advertised_tool(body, "workflow")
  598. return tool_call_chunks(
  599. "advanced-workflow",
  600. "workflow",
  601. {
  602. "script": SNAPSHOT_WORKFLOW_SCRIPT,
  603. "meta": {
  604. "name": "advanced-exe-snapshot",
  605. "description": "exercise one packaged workflow child",
  606. },
  607. },
  608. )
  609. if call_id == "advanced-workflow" and tool_name == "workflow":
  610. if "WORKFLOW_CHILD_OK" not in tool_text:
  611. raise AssertionError(f"workflow returned no expected child value: {tool_text}")
  612. assert_advertised_tool(body, "cordis_undefine")
  613. return tool_call_chunks(
  614. "advanced-undefine",
  615. "cordis_undefine",
  616. {"pluginId": "snap-1"},
  617. )
  618. if call_id == "advanced-undefine" and tool_name == "cordis_undefine":
  619. if "Removed dynamic Plugin snap-1 and all of its Packages." not in tool_text:
  620. raise AssertionError(f"cordis_undefine returned no removal result: {tool_text}")
  621. if "snapshot_double" in advertised_tool_names(body):
  622. raise AssertionError("snapshot_double remained advertised after cordis_undefine")
  623. return text_chunks(SNAPSHOT_FINAL_TEXT)
  624. raise AssertionError(f"unexpected advanced tool follow-up: {call_id} {tool_name}: {tool_text}")
  625. def text_chunks(text: str) -> list[dict[str, object]]:
  626. """Build a complete Messages text response."""
  627. return [
  628. message_start(),
  629. {"type": "content_block_start", "index": 0, "content_block": {"type": "text", "text": ""}},
  630. {"type": "content_block_delta", "index": 0, "delta": {"type": "text_delta", "text": text}},
  631. {"type": "content_block_stop", "index": 0},
  632. {"type": "message_delta", "delta": {"stop_reason": "end_turn"}, "usage": {"output_tokens": 3}},
  633. {"type": "message_stop"},
  634. ]
  635. def tool_call_chunks(call_id: str, name: str, arguments: dict[str, object]) -> list[dict[str, object]]:
  636. """Build a complete Messages tool-use response."""
  637. return [
  638. message_start(),
  639. {
  640. "type": "content_block_start", "index": 0,
  641. "content_block": {"type": "tool_use", "id": call_id, "name": name, "input": {}},
  642. },
  643. {
  644. "type": "content_block_delta", "index": 0,
  645. "delta": {"type": "input_json_delta", "partial_json": json.dumps(arguments)},
  646. },
  647. {"type": "content_block_stop", "index": 0},
  648. {"type": "message_delta", "delta": {"stop_reason": "tool_use"}, "usage": {"output_tokens": 3}},
  649. {"type": "message_stop"},
  650. ]
  651. def message_start() -> dict[str, object]:
  652. """Start a Messages response with deterministic token usage."""
  653. return {
  654. "type": "message_start",
  655. "message": {"id": "msg_smoke", "model": "smoke-model", "usage": {"input_tokens": 3, "output_tokens": 0}},
  656. }
  657. def latest_tool_call(messages: list[object], result_id: object) -> tuple[str, str]:
  658. """Find the assistant call id and name paired with the latest tool result."""
  659. for message in reversed(messages[:-1]):
  660. if not isinstance(message, dict):
  661. continue
  662. calls = message.get("content")
  663. if not isinstance(calls, list):
  664. continue
  665. for call in reversed(calls):
  666. if not isinstance(call, dict):
  667. continue
  668. call_id = call.get("id")
  669. if (
  670. isinstance(call_id, str)
  671. and call_id == result_id
  672. and call.get("type") == "tool_use"
  673. and isinstance(call.get("name"), str)
  674. ):
  675. return call_id, call["name"]
  676. raise AssertionError(f"tool result has no preceding assistant tool call: {messages}")
  677. def message_text(content: object) -> str:
  678. """Read Messages text content in either string or block-list form."""
  679. if isinstance(content, str):
  680. return content
  681. if isinstance(content, list):
  682. return "".join(
  683. block.get("text", "")
  684. for block in content
  685. if isinstance(block, dict) and isinstance(block.get("text"), str)
  686. )
  687. return ""
  688. def advertised_tool_names(body: dict[str, object]) -> set[str]:
  689. """Return the model-facing tool names advertised on one request."""
  690. tools = body.get("tools")
  691. if not isinstance(tools, list):
  692. raise AssertionError(f"model request advertised no tools: {body}")
  693. names: set[str] = set()
  694. for tool in tools:
  695. if not isinstance(tool, dict):
  696. continue
  697. if isinstance(tool.get("name"), str):
  698. names.add(tool["name"])
  699. return names
  700. def assert_advertised_tool(body: dict[str, object], expected: str) -> None:
  701. """Require the packaged deployment to expose the requested tool."""
  702. names = advertised_tool_names(body)
  703. if expected not in names:
  704. raise AssertionError(f"model request did not advertise {expected}: {names}")
  705. class MockModel:
  706. def __enter__(self) -> "MockModel":
  707. MockModelHandler.requests.clear()
  708. self.server = ThreadingHTTPServer(("127.0.0.1", 0), MockModelHandler)
  709. self.thread = threading.Thread(target=self.server.serve_forever, daemon=True)
  710. self.thread.start()
  711. host, port = self.server.server_address
  712. self.url = f"http://{host}:{port}"
  713. return self
  714. def __exit__(self, _exc_type: object, _exc: object, _tb: object) -> None:
  715. self.server.shutdown()
  716. self.server.server_close()
  717. self.thread.join(timeout=5)
  718. def main() -> None:
  719. parser = argparse.ArgumentParser(description=__doc__)
  720. parser.add_argument(
  721. "--scenario",
  722. choices=("all", "sdk-default", "sdk-custom", "sdk-minimal", "sdk-minimal-in-history", "sdk-fs-search", "sdk-spawn-node", "sdk-mcp", "sdk-snapshot", "sdk-restart", "sdk-profile-plugin", "sdk-live", "runner", "direct"),
  723. default="all",
  724. )
  725. parser.add_argument("--exe", type=Path)
  726. parser.add_argument(
  727. "--installed-wheel",
  728. action="store_true",
  729. help="require a clean virtual environment containing matching installed SDK and runtime wheels",
  730. )
  731. parser.add_argument("--update-snapshots", action="store_true")
  732. args = parser.parse_args()
  733. if args.installed_wheel and args.exe is not None:
  734. parser.error("--installed-wheel resolves the wheel's own runtime and cannot be combined with --exe")
  735. if args.scenario == "sdk-live" and not args.installed_wheel:
  736. parser.error("--scenario sdk-live requires --installed-wheel")
  737. if args.scenario == "sdk-profile-plugin" and not args.installed_wheel:
  738. parser.error("--scenario sdk-profile-plugin requires --installed-wheel")
  739. if args.installed_wheel:
  740. args.exe = assert_installed_wheel_environment()
  741. if args.scenario in {"all", "sdk-custom", "sdk-minimal", "sdk-minimal-in-history", "sdk-fs-search", "sdk-spawn-node", "sdk-snapshot", "sdk-restart", "runner", "direct"} and args.exe is None:
  742. parser.error("--exe is required for custom, minimal, fs-search, spawn-node, snapshot, restart, runner, and direct scenarios")
  743. if args.update_snapshots and args.scenario not in {"all", "sdk-minimal", "sdk-minimal-in-history", "sdk-snapshot", "sdk-restart"}:
  744. parser.error("--update-snapshots requires --scenario sdk-minimal, sdk-minimal-in-history, sdk-snapshot, sdk-restart, or all")
  745. if args.exe is not None and not args.exe.is_file():
  746. parser.error(f"runtime executable does not exist: {args.exe}")
  747. if args.scenario in {"all", "runner"}:
  748. assert args.exe is not None
  749. smoke_packaged_runner(args.exe.resolve())
  750. if args.scenario == "runner":
  751. print("smoke-python-runtime: runner passed")
  752. return
  753. if args.scenario == "sdk-live":
  754. smoke_sdk_live()
  755. print("smoke-python-runtime: sdk-live passed")
  756. return
  757. with MockModel() as model:
  758. if args.scenario in {"all", "sdk-default"}:
  759. smoke_sdk_default(model.url)
  760. if args.scenario in {"all", "sdk-custom"}:
  761. assert args.exe is not None
  762. smoke_sdk_custom(model.url, args.exe.resolve())
  763. if args.scenario in {"all", "sdk-minimal"}:
  764. assert args.exe is not None
  765. smoke_sdk_minimal(model.url, args.exe.resolve(), args.update_snapshots)
  766. if args.scenario in {"all", "sdk-minimal-in-history"}:
  767. assert args.exe is not None
  768. smoke_sdk_minimal(model.url, args.exe.resolve(), args.update_snapshots, in_history=True)
  769. if args.scenario in {"all", "sdk-fs-search"}:
  770. assert args.exe is not None
  771. smoke_sdk_fs_search(model.url, args.exe.resolve())
  772. if args.scenario in {"all", "sdk-spawn-node"}:
  773. assert args.exe is not None
  774. smoke_sdk_spawn_node(model.url, args.exe.resolve())
  775. if args.scenario in {"all", "sdk-mcp"}:
  776. smoke_sdk_mcp(model.url, None if args.exe is None else args.exe.resolve())
  777. if args.scenario in {"all", "sdk-snapshot"}:
  778. assert args.exe is not None
  779. smoke_sdk_snapshot(model.url, args.exe.resolve(), args.update_snapshots)
  780. if args.scenario in {"all", "sdk-restart"}:
  781. assert args.exe is not None
  782. smoke_sdk_restart_snapshot(model.url, args.exe.resolve(), args.update_snapshots)
  783. if args.installed_wheel and args.scenario in {"all", "sdk-profile-plugin"}:
  784. smoke_sdk_profile_plugin(model.url)
  785. if args.scenario in {"all", "direct"}:
  786. assert args.exe is not None
  787. smoke_direct(model.url, args.exe.resolve())
  788. if not MockModelHandler.requests:
  789. raise AssertionError("mock model endpoint received no requests")
  790. print(f"smoke-python-runtime: {args.scenario} passed")
  791. def assert_installed_wheel_environment() -> Path:
  792. """Prove that this process imports matching non-editable wheel installations."""
  793. if sys.prefix == sys.base_prefix:
  794. raise AssertionError("installed-wheel smoke must run inside a virtual environment")
  795. if os.environ.get("PYTHONPATH"):
  796. raise AssertionError("installed-wheel smoke requires PYTHONPATH to be unset")
  797. if os.environ.get("DSH_RUNTIME_MODE"):
  798. raise AssertionError("installed-wheel smoke requires DSH_RUNTIME_MODE to be unset")
  799. repo_root = Path(__file__).resolve().parent.parent
  800. cwd = Path.cwd().resolve()
  801. if cwd.is_relative_to(repo_root):
  802. raise AssertionError(f"installed-wheel smoke must run outside the repository, got {cwd}")
  803. sdk_version = importlib.metadata.version("deepseek-harness-sdk")
  804. runtime_version = importlib.metadata.version("deepseek-harness-runtime-bin")
  805. if sdk_version != runtime_version:
  806. raise AssertionError(
  807. f"installed SDK/runtime versions differ: {sdk_version} != {runtime_version}"
  808. )
  809. expected_runtime_requirement = f"deepseek-harness-runtime-bin=={sdk_version}"
  810. requirements = importlib.metadata.requires("deepseek-harness-sdk") or []
  811. if expected_runtime_requirement not in requirements:
  812. raise AssertionError(
  813. f"installed SDK does not require {expected_runtime_requirement}: {requirements}"
  814. )
  815. prefix = Path(sys.prefix).resolve()
  816. imported: dict[str, Path] = {}
  817. for name in ("deepseek_harness", "deepseek_harness_runtime"):
  818. module = importlib.import_module(name)
  819. module_file = getattr(module, "__file__", None)
  820. if not isinstance(module_file, str):
  821. raise AssertionError(f"installed module {name} has no filesystem location")
  822. path = Path(module_file).resolve()
  823. if not path.is_relative_to(prefix):
  824. raise AssertionError(f"installed module {name} came from outside the virtual environment: {path}")
  825. if path.is_relative_to(repo_root):
  826. raise AssertionError(f"installed module {name} came from the repository checkout: {path}")
  827. imported[name] = path
  828. runtime_module = sys.modules["deepseek_harness_runtime"]
  829. executable = runtime_module.bundled_runtime_path().resolve()
  830. runtime_package = imported["deepseek_harness_runtime"].parent
  831. if not executable.is_relative_to(runtime_package):
  832. raise AssertionError(f"bundled runtime came from outside the installed runtime wheel: {executable}")
  833. runtime_files = importlib.metadata.files("deepseek-harness-runtime-bin") or []
  834. if not any(Path(file).name == executable.name for file in runtime_files):
  835. raise AssertionError(f"runtime executable is absent from installed distribution records: {executable}")
  836. return executable
  837. def smoke_sdk_live() -> None:
  838. """Run a real-model, tool-using two-turn task through installed wheels."""
  839. from deepseek_harness import DeepSeekHarness
  840. api_key = os.environ.get("DEEPSEEK_API_KEY")
  841. base_url = os.environ.get("DEEPSEEK_BASE_URL")
  842. if not api_key:
  843. raise AssertionError("sdk-live requires DEEPSEEK_API_KEY")
  844. if not base_url:
  845. raise AssertionError("sdk-live requires an explicit DEEPSEEK_BASE_URL")
  846. with tempfile.TemporaryDirectory(prefix="dsh-sdk-live-") as temporary:
  847. root = Path(temporary).resolve()
  848. dsh_home = root / "home"
  849. sessions = dsh_home / "sessions"
  850. marker = root / "live-api-marker.txt"
  851. session_id = "installed-wheel-live-api"
  852. shell_tool = "pwsh" if IS_WINDOWS else "bash"
  853. create_prompt = (
  854. f"Use the {shell_tool} tool to create the file at the absolute path below with exact UTF-8 "
  855. f"content {LIVE_API_SENTINEL}, with no newline or byte-order mark. "
  856. f"Then reply with exactly {LIVE_API_SENTINEL}.\n{marker}"
  857. )
  858. with DeepSeekHarness(
  859. provider="deepseek-official",
  860. model="deepseek-v4-flash",
  861. cwd=str(root),
  862. dsh_home=str(dsh_home),
  863. env={
  864. "DSH_PERMISSION_MODE": "danger-full-access",
  865. "DSH_TELEMETRY_DISABLED": "1",
  866. },
  867. api_key=api_key,
  868. base_url=base_url,
  869. request_timeout_seconds=180,
  870. ) as harness:
  871. created = harness.run(create_prompt, session_id=session_id)
  872. assert_live_turn("create", created)
  873. if not marker.is_file():
  874. raise AssertionError(f"create turn did not create {marker}")
  875. if marker.read_bytes() != LIVE_API_SENTINEL.encode("utf-8"):
  876. raise AssertionError(f"create turn wrote unexpected bytes to {marker}")
  877. # The challenge is absent from the prior turn and the verification prompt.
  878. challenge = secrets.token_hex(32).encode("ascii")
  879. marker.write_bytes(challenge)
  880. with tempfile.TemporaryDirectory(prefix="receipt-", dir=root) as receipt_directory:
  881. receipt = Path(receipt_directory) / "receipt.txt"
  882. verify_prompt = (
  883. "The file created in the previous turn has changed externally. "
  884. "Use a tool to read that same file and copy its exact current content to the "
  885. "new receipt path below, without changing the source file. "
  886. "Preserve every byte; do not add a newline or byte-order mark. "
  887. f"{LIVE_API_SENTINEL} is only the completion acknowledgement, "
  888. "not a claim about the source or receipt contents. "
  889. f"Then reply with exactly {LIVE_API_SENTINEL}.\n{receipt}"
  890. )
  891. verified = harness.run(verify_prompt, session_id=session_id)
  892. assert_live_turn("verify", verified)
  893. if not receipt.is_file():
  894. raise AssertionError(f"verify turn did not create receipt {receipt}")
  895. if receipt.read_bytes() != challenge:
  896. raise AssertionError(f"verify turn wrote unexpected bytes to receipt {receipt}")
  897. if not marker.is_file() or marker.read_bytes() != challenge:
  898. raise AssertionError(f"verify turn changed source file {marker}")
  899. assert_zstd_session_log(sessions)
  900. def assert_live_turn(label: str, result: RunResult) -> None:
  901. """Require completed model tool use and the exact smoke answer for each live turn."""
  902. if result.finish_reason != "completed":
  903. event_types = [event.get("type") for event in result.events]
  904. turn_end_data = next(
  905. (event.get("data") for event in reversed(result.events) if event.get("type") == "turn/end"),
  906. None,
  907. )
  908. turn_end = safe_turn_end(turn_end_data)
  909. raise AssertionError(
  910. f"{label} turn ended with {result.finish_reason!r}; "
  911. f"final={result.final_response!r}; turn_end={turn_end!r}; events={event_types}"
  912. )
  913. if not any(event.get("type") == "tool/call" for event in result.events):
  914. raise AssertionError(
  915. f"{label} turn made no model-requested tool call; "
  916. f"final={result.final_response!r}"
  917. )
  918. if result.final_response.strip() != LIVE_API_SENTINEL:
  919. raise AssertionError(f"{label} turn returned {result.final_response!r}")
  920. def safe_turn_end(value: object) -> object:
  921. """Project a live-provider failure without retaining credential-bearing text."""
  922. if not isinstance(value, dict):
  923. return value
  924. reason = value.get("reason")
  925. if not isinstance(reason, dict):
  926. return {"turn": value.get("turn"), "reason": reason}
  927. error = reason.get("error")
  928. safe_error = None
  929. if isinstance(error, dict):
  930. safe_error = {
  931. key: error.get(key)
  932. for key in ("code", "status")
  933. if error.get(key) is not None
  934. }
  935. return {
  936. "turn": value.get("turn"),
  937. "reason": {
  938. "kind": reason.get("kind"),
  939. **({"error": safe_error} if safe_error is not None else {}),
  940. },
  941. }
  942. def smoke_sdk_default(base_url: str) -> None:
  943. from deepseek_harness import DeepSeekHarness
  944. with tempfile.TemporaryDirectory(prefix="dsh-sdk-default-") as temporary:
  945. root = Path(temporary).resolve()
  946. dsh_home = root / "home"
  947. sessions = dsh_home / "sessions"
  948. with DeepSeekHarness(
  949. provider="deepseek-official",
  950. model="smoke-model",
  951. cwd=str(root),
  952. dsh_home=str(dsh_home),
  953. env={
  954. "DSH_PERMISSION_MODE": "danger-full-access",
  955. "DSH_TELEMETRY_DISABLED": "1",
  956. },
  957. api_key="sk-keyless-smoke",
  958. base_url=base_url,
  959. request_timeout_seconds=60,
  960. ) as harness:
  961. result = harness.run("reply with the smoke text", session_id="default-smoke")
  962. assert result.final_response == EXPECTED_TEXT, (
  963. f"final={result.final_response!r} finish={result.finish_reason!r} "
  964. f"events={[event.get('type') for event in result.events]!r} "
  965. f"turn_end={safe_turn_end(next((event.get('data', event) for event in reversed(result.events) if event.get('type') == 'turn/end'), {}))!r}"
  966. )
  967. assert_zstd_session_log(sessions)
  968. def smoke_sdk_custom(base_url: str, executable: Path) -> None:
  969. from deepseek_harness import DeepSeekHarness
  970. with tempfile.TemporaryDirectory(prefix="dsh-sdk-custom-") as temporary:
  971. root = Path(temporary).resolve()
  972. dsh_home = root / "home"
  973. sessions = dsh_home / "sessions"
  974. patch = write_advanced_profile_patch(root, "custom.patch.yml", sessions)
  975. with DeepSeekHarness(
  976. provider="deepseek-official",
  977. model="smoke-model",
  978. cwd=str(root),
  979. dsh_bin=str(executable),
  980. dsh_home=str(dsh_home),
  981. patches=(str(patch),),
  982. env={
  983. "DSH_PERMISSION_MODE": "danger-full-access",
  984. "DSH_TELEMETRY_DISABLED": "1",
  985. },
  986. api_key="sk-keyless-smoke",
  987. base_url=base_url,
  988. request_timeout_seconds=60,
  989. ) as harness:
  990. text_result = harness.run("reply with the smoke text", session_id="custom-smoke")
  991. code_result = harness.run(CODE_PROMPT, session_id="custom-smoke")
  992. workflow_result = harness.run(WORKFLOW_PROMPT, session_id="custom-smoke")
  993. assert text_result.final_response == EXPECTED_TEXT, text_result.final_response
  994. assert code_result.final_response == CODE_WORKER_TEXT, code_result.final_response
  995. assert workflow_result.final_response == WORKFLOW_WORKER_TEXT, workflow_result.final_response
  996. assert_session_log(sessions, root, EXPECTED_TEXT, CODE_WORKER_TEXT, WORKFLOW_WORKER_TEXT)
  997. def smoke_sdk_minimal(
  998. base_url: str, executable: Path, update_snapshots: bool, *, in_history: bool = False,
  999. ) -> None:
  1000. """Exercise the shipped standalone minimal profile through the packaged executable."""
  1001. from deepseek_harness import DeepSeekHarness
  1002. # One mock model serves every scenario of a run, so the snapshot takes this turn's slice.
  1003. first_request = len(MockModelHandler.requests)
  1004. with tempfile.TemporaryDirectory(prefix="dsh-sdk-minimal-") as temporary:
  1005. root = Path(temporary).resolve()
  1006. dsh_home = root / "home"
  1007. sessions = dsh_home / "sessions"
  1008. patches = ()
  1009. if in_history:
  1010. patch = root / "in-history.patch.yml"
  1011. patch.write_text(json.dumps([
  1012. {"id": "llm-deepseek", "config": {"models": [
  1013. {"id": "smoke-model", "systemPromptUpdate": "in-history"},
  1014. ]}},
  1015. {"insert": [{
  1016. "id": "in-history-prompt",
  1017. "name": (Path(__file__).resolve().parent / "fixtures/python-sdk-in-history-prompt.mjs").as_uri(),
  1018. }]},
  1019. ]))
  1020. patches = (str(patch),)
  1021. with DeepSeekHarness(
  1022. provider="deepseek-official",
  1023. model="smoke-model",
  1024. cwd=str(root),
  1025. dsh_bin=str(executable),
  1026. dsh_home=str(dsh_home),
  1027. profile="sdk-minimal",
  1028. patches=patches,
  1029. api_key="sk-keyless-smoke",
  1030. base_url=base_url,
  1031. request_timeout_seconds=60,
  1032. ) as harness:
  1033. result = harness.run(MINIMAL_PROMPT, session_id="minimal-agent-smoke")
  1034. event_text = json.dumps(result.events)
  1035. if MINIMAL_TEXT not in event_text:
  1036. raise AssertionError(f"minimal agent run emitted no final response: {result.events}")
  1037. assert_session_log(sessions, root, MINIMAL_TEXT, "COUNT=1", "COUNT=2")
  1038. requests = MockModelHandler.requests[first_request:]
  1039. if in_history:
  1040. logs = read_session_logs(sessions)
  1041. files = build_in_history_snapshot_files(result, requests, logs[result.session_id])
  1042. compare_snapshot_files(
  1043. files, update_snapshots, IN_HISTORY_SNAPSHOT_DIRECTORY, IN_HISTORY_SNAPSHOT_FILENAMES,
  1044. )
  1045. else:
  1046. files = build_minimal_snapshot_files(requests, root)
  1047. compare_snapshot_files(
  1048. files, update_snapshots, MINIMAL_SNAPSHOT_DIRECTORY, MINIMAL_SNAPSHOT_FILENAMES,
  1049. )
  1050. def smoke_sdk_fs_search(base_url: str, executable: Path) -> None:
  1051. """Exercise real grep and glob spawns through the packaged executable."""
  1052. from deepseek_harness import DeepSeekHarness
  1053. with tempfile.TemporaryDirectory(prefix="dsh-sdk-fs-search-") as temporary:
  1054. root = Path(temporary).resolve()
  1055. (root / "needle.txt").write_text(f"{FS_SEARCH_MARKER}\n")
  1056. dsh_home = root / "home"
  1057. sessions = dsh_home / "sessions"
  1058. patch = write_profile_patch(root, "fs-search.patch.yml", sessions, [
  1059. {"id": "skill-filesystem", "disabled": True},
  1060. {"id": "tool-fs-search", "config": {"sampleOverCapGlobResults": False}},
  1061. ])
  1062. with DeepSeekHarness(
  1063. provider="deepseek-official",
  1064. model="smoke-model",
  1065. cwd=str(root),
  1066. dsh_bin=str(executable),
  1067. dsh_home=str(dsh_home),
  1068. patches=(str(patch),),
  1069. env={
  1070. "DSH_PERMISSION_MODE": "danger-full-access",
  1071. "DSH_TELEMETRY_DISABLED": "1",
  1072. },
  1073. api_key="sk-keyless-smoke",
  1074. base_url=base_url,
  1075. request_timeout_seconds=60,
  1076. ) as harness:
  1077. result = harness.run(FS_SEARCH_PROMPT, session_id="fs-search-smoke")
  1078. assert result.final_response == FS_SEARCH_TEXT, result.final_response
  1079. assert_session_log(sessions, root, FS_SEARCH_TEXT, FS_SEARCH_MARKER, "needle.txt")
  1080. def smoke_sdk_spawn_node(base_url: str, executable: Path) -> None:
  1081. """A shell command starting with `node` must reach the machine's Node, not the executable."""
  1082. from deepseek_harness import DeepSeekHarness
  1083. with tempfile.TemporaryDirectory(prefix="dsh-sdk-spawn-node-") as temporary:
  1084. root = Path(temporary).resolve()
  1085. dsh_home = root / "home"
  1086. sessions = dsh_home / "sessions"
  1087. patch = write_profile_patch(root, "spawn-node.patch.yml", sessions, [])
  1088. with DeepSeekHarness(
  1089. provider="deepseek-official",
  1090. model="smoke-model",
  1091. cwd=str(root),
  1092. dsh_bin=str(executable),
  1093. dsh_home=str(dsh_home),
  1094. patches=(str(patch),),
  1095. env={
  1096. "DSH_PERMISSION_MODE": "danger-full-access",
  1097. "DSH_TELEMETRY_DISABLED": "1",
  1098. },
  1099. api_key="sk-keyless-smoke",
  1100. base_url=base_url,
  1101. request_timeout_seconds=60,
  1102. ) as harness:
  1103. result = harness.run(SPAWN_NODE_PROMPT, session_id="spawn-node-smoke")
  1104. assert result.final_response == SPAWN_NODE_TEXT, result.final_response
  1105. assert_session_log(sessions, root, SPAWN_NODE_TEXT, "PKG_EXECPATH=ABSENT")
  1106. def smoke_sdk_mcp(base_url: str, executable: Path | None) -> None:
  1107. """Discover and call an external stdio MCP tool through the packaged client."""
  1108. from deepseek_harness import DeepSeekHarness
  1109. with tempfile.TemporaryDirectory(prefix="dsh-sdk-mcp-") as temporary:
  1110. root = Path(temporary).resolve()
  1111. dsh_home = root / "home"
  1112. sessions = dsh_home / "sessions"
  1113. server_script = root / "mcp_server.py"
  1114. server_script.write_text(MCP_SERVER_SCRIPT)
  1115. patch = write_mcp_patch(root, sessions, server_script)
  1116. discovery_log = server_script.with_suffix(".log")
  1117. with DeepSeekHarness(
  1118. provider="deepseek-official",
  1119. model="smoke-model",
  1120. cwd=str(root),
  1121. dsh_bin=None if executable is None else str(executable),
  1122. dsh_home=str(dsh_home),
  1123. patches=(str(patch),),
  1124. env={
  1125. "DSH_PERMISSION_MODE": "danger-full-access",
  1126. "DSH_TELEMETRY_DISABLED": "1",
  1127. },
  1128. api_key="sk-keyless-smoke",
  1129. base_url=base_url,
  1130. request_timeout_seconds=60,
  1131. ) as harness:
  1132. result = harness.run(MCP_PROMPT, session_id="mcp-smoke")
  1133. assert result.final_response == MCP_TEXT, result.final_response
  1134. assert discovery_log.read_text().splitlines() == [
  1135. "server/discover",
  1136. "initialize",
  1137. "notifications/initialized",
  1138. "tools/list",
  1139. "tools/call",
  1140. ]
  1141. assert_session_log(sessions, root, MCP_TEXT, "mcp__fixture__add", "42")
  1142. def smoke_sdk_profile_plugin(base_url: str) -> None:
  1143. """Install an external bundle through Python's dsh command and load it in the SDK."""
  1144. from deepseek_harness import DeepSeekHarness
  1145. with tempfile.TemporaryDirectory(prefix="dsh-sdk-profile-plugin-") as temporary:
  1146. root = Path(temporary).resolve()
  1147. dsh_home = root / "home"
  1148. plugin = root / "plugin"
  1149. plugin.mkdir()
  1150. (plugin / "package.json").write_text(json.dumps({
  1151. "name": "dsh-python-blackbox-plugin",
  1152. "version": "1.0.0",
  1153. "private": True,
  1154. "type": "module",
  1155. "exports": "./index.js",
  1156. "peerDependencies": {"@deepseek-ai/cordis": "*"},
  1157. "dsh": {"bundle": {"patch": "./cordis.patch.yml"}},
  1158. }, indent=2))
  1159. (plugin / "index.js").write_text(
  1160. "import { Context } from '@deepseek-ai/cordis'\n"
  1161. "export const name = 'python-sdk-blackbox-plugin'\n"
  1162. "export const inject = ['systemPrompt']\n"
  1163. "export function apply(ctx) {\n"
  1164. " if (!(ctx instanceof Context)) throw new Error('external plugin loaded a second Cordis instance')\n"
  1165. " ctx.effect(() => ctx.systemPrompt.section({\n"
  1166. " name: 'python-sdk:blackbox-plugin',\n"
  1167. " order: 10,\n"
  1168. f" text: '{PROFILE_PLUGIN_MARKER}',\n"
  1169. " }))\n"
  1170. "}\n"
  1171. )
  1172. (plugin / "cordis.patch.yml").write_text(json.dumps([{
  1173. "insert": [{"id": "python-sdk-blackbox-plugin", "name": "dsh-python-blackbox-plugin"}],
  1174. }], indent=2))
  1175. dsh = Path(sysconfig.get_path("scripts")) / ("dsh.exe" if IS_WINDOWS else "dsh")
  1176. environment = {**os.environ, "DSH_HOME": str(dsh_home)}
  1177. installed = subprocess.run(
  1178. [str(dsh), "plugin", "--profile", "sdk", "add", f"file:{plugin}"],
  1179. cwd=root,
  1180. env=environment,
  1181. text=True,
  1182. capture_output=True,
  1183. check=False,
  1184. )
  1185. if installed.returncode != 0:
  1186. raise AssertionError(
  1187. f"Python-installed dsh could not add the external profile plugin: "
  1188. f"returncode={installed.returncode} (0x{installed.returncode & 0xffffffff:08x}) "
  1189. f"stdout={installed.stdout!r} stderr={installed.stderr!r}"
  1190. )
  1191. manifest = json.loads((dsh_home / "profiles" / "sdk" / "package.json").read_text())
  1192. if "dsh-python-blackbox-plugin" not in manifest.get("dependencies", {}):
  1193. raise AssertionError(f"dsh plugin did not record the external dependency: {manifest}")
  1194. if "dsh-python-blackbox-plugin" not in manifest["dsh"]["profile"]["bundles"]:
  1195. raise AssertionError(f"dsh plugin did not activate the external bundle: {manifest}")
  1196. harness = DeepSeekHarness(
  1197. provider="deepseek-official",
  1198. model="smoke-model",
  1199. cwd=str(root),
  1200. dsh_home=str(dsh_home),
  1201. env={
  1202. "DSH_PERMISSION_MODE": "danger-full-access",
  1203. "DSH_TELEMETRY_DISABLED": "1",
  1204. },
  1205. api_key="sk-keyless-smoke",
  1206. base_url=base_url,
  1207. request_timeout_seconds=60,
  1208. )
  1209. try:
  1210. with harness:
  1211. result = harness.run(PROFILE_PLUGIN_PROMPT, session_id="profile-plugin-smoke")
  1212. except Exception as error:
  1213. raise AssertionError(
  1214. f"external profile plugin runtime failed: {harness.client._runtime_diagnostics()}"
  1215. ) from error
  1216. assert result.final_response == PROFILE_PLUGIN_TEXT, result.final_response
  1217. assert_zstd_session_log(dsh_home / "sessions")
  1218. def smoke_sdk_snapshot(base_url: str, executable: Path, update_snapshots: bool) -> None:
  1219. """Drive and compare the advanced SDK/executable behavioral snapshot."""
  1220. from deepseek_harness import DeepSeekHarness
  1221. with tempfile.TemporaryDirectory(prefix="dsh-sdk-snapshot-") as temporary:
  1222. root = Path(temporary).resolve()
  1223. dsh_home = root / "home"
  1224. sessions = dsh_home / "sessions"
  1225. patch = write_advanced_profile_patch(root, "snapshot.patch.yml", sessions)
  1226. feedback_patch = write_profile_patch(root, "feedback.patch.yml", sessions, [{"insert": [
  1227. {"id": "snapshot-image-offload", "name": (
  1228. Path(__file__).resolve().parent / "fixtures/python-snapshot-image-offload.mjs"
  1229. ).as_uri(), "config": {"parentSessionId": SNAPSHOT_SESSION_ID}},
  1230. {"id": "snapshot-workflow-order", "name": (
  1231. Path(__file__).resolve().parent / "fixtures/python-snapshot-workflow-order.mjs"
  1232. ).as_uri(), "config": {
  1233. "parentSessionId": SNAPSHOT_SESSION_ID, "prompt": SNAPSHOT_WORKFLOW_CHILD_PROMPT,
  1234. }},
  1235. {"id": "snapshot-message-feedback", "name": "@deepseek-ai/dsh-message-feedback",
  1236. "config": {"maxNoteBytes": 1024}},
  1237. {"id": "snapshot-feedback-producer", "name": (
  1238. Path(__file__).resolve().parent.parent / "snapshots/sdk/text-turn/feedback-producer.mjs"
  1239. ).as_uri()},
  1240. ]}])
  1241. creation_patch = write_profile_patch(root, "creation.patch.yml", sessions, [
  1242. {"insert": [{
  1243. "id": "serial-created-fixture",
  1244. "name": str(Path(__file__).resolve().parents[1] / "packages/core/agent-loop/tests/fixtures/serial-created.mjs"),
  1245. }]},
  1246. ])
  1247. with DeepSeekHarness(
  1248. provider="deepseek-official",
  1249. model="smoke-model",
  1250. cwd=str(root),
  1251. dsh_bin=str(executable),
  1252. dsh_home=str(dsh_home),
  1253. patches=(str(patch), str(feedback_patch), str(creation_patch)),
  1254. env={
  1255. "DSH_PERMISSION_MODE": "danger-full-access",
  1256. "DSH_TELEMETRY_DISABLED": "1",
  1257. },
  1258. api_key="sk-keyless-smoke",
  1259. base_url=base_url,
  1260. request_timeout_seconds=60,
  1261. ) as harness:
  1262. result = harness.run(SNAPSHOT_PROMPT, session_id=SNAPSHOT_SESSION_ID)
  1263. assert result.final_response == SNAPSHOT_FINAL_TEXT, result.final_response
  1264. offloads = [event for event in result.events if event.get("type") == "image/offload"]
  1265. if len(offloads) != 1 or "surfaceOp" in offloads[0]:
  1266. raise AssertionError(f"advanced snapshot expected one standalone image offload: {offloads}")
  1267. targets = offloads[0]["data"]["targets"]
  1268. if len(targets) != 1 or targets[0]["imageIndexes"] != [0]:
  1269. raise AssertionError(f"advanced snapshot selected unexpected image occurrences: {targets}")
  1270. feedback_types = [event.get("type") for event in result.events
  1271. if str(event.get("type")).startswith("feedback/")]
  1272. if feedback_types != ["feedback/record", "feedback/record", "feedback/message-put", "feedback/message-put", "feedback/message-delete"]:
  1273. raise AssertionError(f"advanced snapshot did not exercise all feedback mutations: {feedback_types}")
  1274. methods = [notification.method for notification in result.notifications]
  1275. if methods.count("subagent.started") != 2 or methods.count("subagent.finished") != 2:
  1276. raise AssertionError(f"advanced snapshot emitted unexpected subagent lifecycle: {methods}")
  1277. ptc_events = [event for event in result.events
  1278. if event.get("type") in ("tool/ptc-dispatch-start", "tool/ptc-dispatch")]
  1279. if [event["type"] for event in ptc_events] != ["tool/ptc-dispatch-start", "tool/ptc-dispatch"] * 2:
  1280. raise AssertionError(f"advanced snapshot emitted unexpected PTC dispatch events: {ptc_events}")
  1281. for index, event in enumerate(ptc_events):
  1282. data = event["data"]
  1283. identity = (data.get("rootCallId"), data.get("parentCallId"), data.get("subCallId"))
  1284. root_call = "advanced-code" if index < 2 else "advanced-denied-ptc"
  1285. if identity != (root_call, root_call, root_call + ":ptc:1"):
  1286. raise AssertionError(f"advanced snapshot emitted unexpected PTC dispatch identity: {identity}")
  1287. if {"description", "parameters", "schema"}.intersection(data):
  1288. raise AssertionError("PTC binding schema entered the packaged SDK wire")
  1289. errors = [event["data"]["error"] for event in result.events
  1290. if event.get("type") in ("tool/result", "tool/ptc-dispatch") and "error" in event["data"]]
  1291. assert errors == [{
  1292. "name": "AutoReviewDeniedError", "code": "AUTO_REVIEW_DENIED", "reason": " transport raw\r\nreason ",
  1293. }] * 2, errors
  1294. logs = read_session_logs(sessions)
  1295. child_ids = snapshot_child_ids(result)
  1296. expected_ids = {SNAPSHOT_SESSION_ID, *child_ids}
  1297. if set(logs) != expected_ids:
  1298. raise AssertionError(f"advanced snapshot expected parent plus two child logs: {sorted(logs)}")
  1299. if "DIRECT_CHILD_OK" not in render_jsonl(logs[child_ids[0]]):
  1300. raise AssertionError("first advanced child log has no direct-subagent result")
  1301. if "WORKFLOW_CHILD_OK" not in render_jsonl(logs[child_ids[1]]):
  1302. raise AssertionError("second advanced child log has no workflow-subagent result")
  1303. files = build_snapshot_files(result, logs, child_ids, root)
  1304. compare_snapshot_files(
  1305. files, update_snapshots, ADVANCED_SNAPSHOT_DIRECTORY, ADVANCED_SNAPSHOT_FILENAMES,
  1306. )
  1307. def smoke_sdk_restart_snapshot(base_url: str, executable: Path, update_snapshots: bool) -> None:
  1308. """Snapshot two isolated sessions across complete SDK runtime restarts."""
  1309. from deepseek_harness import DeepSeekHarness
  1310. with tempfile.TemporaryDirectory(prefix="dsh-sdk-restart-") as temporary:
  1311. root = Path(temporary).resolve()
  1312. dsh_home = root / "home"
  1313. sessions = dsh_home / "sessions"
  1314. patch = write_advanced_profile_patch(root, "restart.patch.yml", sessions)
  1315. first_request = len(MockModelHandler.requests)
  1316. def run(prompt: str, session_id: str) -> "RunResult":
  1317. with DeepSeekHarness(
  1318. provider="deepseek-official",
  1319. model="smoke-model",
  1320. cwd=str(root),
  1321. dsh_bin=str(executable),
  1322. dsh_home=str(dsh_home),
  1323. patches=(str(patch),),
  1324. env={
  1325. "DSH_PERMISSION_MODE": "danger-full-access",
  1326. "DSH_TELEMETRY_DISABLED": "1",
  1327. },
  1328. api_key="sk-keyless-smoke",
  1329. base_url=base_url,
  1330. request_timeout_seconds=60,
  1331. ) as harness:
  1332. return harness.run(prompt, session_id=session_id)
  1333. first = run(RESTART_FIRST_PROMPT, RESTART_FIRST_SESSION_ID)
  1334. second = run(RESTART_SECOND_PROMPT, RESTART_SECOND_SESSION_ID)
  1335. requests = MockModelHandler.requests[first_request:]
  1336. if len(requests) != 2:
  1337. raise AssertionError(f"restart snapshot expected two model requests: {requests}")
  1338. if first.final_response != RESTART_FIRST_TEXT or second.final_response != RESTART_SECOND_TEXT:
  1339. raise AssertionError(
  1340. f"restart snapshot responses differ: {first.final_response!r}, {second.final_response!r}"
  1341. )
  1342. logs = read_session_logs(sessions)
  1343. expected_ids = {RESTART_FIRST_SESSION_ID, RESTART_SECOND_SESSION_ID}
  1344. if set(logs) != expected_ids:
  1345. raise AssertionError(f"restart snapshot expected two durable sessions: {sorted(logs)}")
  1346. for session_id, expected in (
  1347. (RESTART_FIRST_SESSION_ID, RESTART_FIRST_TEXT),
  1348. (RESTART_SECOND_SESSION_ID, RESTART_SECOND_TEXT),
  1349. ):
  1350. records = logs[session_id]
  1351. if sum(record.get("type") == "turn/end" for record in records) != 1:
  1352. raise AssertionError(f"restart snapshot {session_id} has an unexpected turn count")
  1353. if expected not in render_jsonl(records):
  1354. raise AssertionError(f"restart snapshot durable log has no {expected}")
  1355. files = build_restart_snapshot_files(
  1356. first,
  1357. second,
  1358. requests,
  1359. logs,
  1360. root,
  1361. sessions,
  1362. )
  1363. compare_snapshot_files(
  1364. files, update_snapshots, RESTART_SNAPSHOT_DIRECTORY, RESTART_SNAPSHOT_FILENAMES,
  1365. )
  1366. def smoke_direct(base_url: str, executable: Path) -> None:
  1367. with tempfile.TemporaryDirectory(prefix="dsh-direct-") as temporary:
  1368. root = Path(temporary).resolve()
  1369. dsh_home = root / "home"
  1370. sessions = dsh_home / "sessions"
  1371. patch = write_profile_patch(root, "direct.patch.yml", sessions, [])
  1372. environment = {
  1373. **os.environ,
  1374. "DSH_HOME": str(dsh_home),
  1375. "DSH_PERMISSION_MODE": "danger-full-access",
  1376. "DSH_TELEMETRY_DISABLED": "1",
  1377. "DEEPSEEK_API_KEY": "sk-keyless-smoke",
  1378. "DEEPSEEK_BASE_URL": base_url,
  1379. }
  1380. peer = RuntimePeer(
  1381. [str(executable), "--profile", "sdk", "--patch", str(patch)],
  1382. root,
  1383. environment,
  1384. )
  1385. try:
  1386. peer.send({"jsonrpc": "2.0", "id": "initialize", "method": "initialize", "params": {"cwd": str(root), "provider": "deepseek-official", "model": "smoke-model"}})
  1387. peer.read_until(lambda message: message.get("id") == "initialize")
  1388. peer.send({
  1389. "jsonrpc": "2.0",
  1390. "id": "prompt",
  1391. "method": "session/prompt",
  1392. "params": {"sessionId": "direct-smoke", "contentBlocks": [{"type": "text", "text": "reply with the smoke text"}]},
  1393. })
  1394. messages = peer.read_until(lambda message: message.get("id") == "prompt")
  1395. if not any(is_idle_notification(message) for message in messages):
  1396. messages.extend(peer.read_until(is_idle_notification))
  1397. event_text = json.dumps(messages)
  1398. if EXPECTED_TEXT not in event_text:
  1399. raise AssertionError(f"direct runtime emitted no final response: {messages}")
  1400. peer.send({"jsonrpc": "2.0", "id": "shutdown", "method": "shutdown"})
  1401. peer.read_until(lambda message: message.get("id") == "shutdown")
  1402. finally:
  1403. peer.close()
  1404. assert_session_log(sessions, root, EXPECTED_TEXT)
  1405. def smoke_packaged_runner(executable: Path) -> None:
  1406. """Exercise the private subprocess runner through the single-file entry."""
  1407. with tempfile.TemporaryDirectory(prefix="dsh-packaged-runner-") as temporary:
  1408. root = Path(temporary).resolve()
  1409. target_script = (
  1410. "import os,sys; "
  1411. "ok = (os.getcwd() == os.environ['PACKAGED_RUNNER_EXPECTED_CWD'] "
  1412. "and os.environ.get('DSH_SUBPROCESS_RUNNER') == 'target-collision-restored'); "
  1413. "sys.exit(7 if ok else 9)"
  1414. )
  1415. if not IS_WINDOWS:
  1416. request_path = root / "launch-request.json"
  1417. target_env = dict(os.environ)
  1418. target_env["DSH_SUBPROCESS_RUNNER"] = "target-collision-restored"
  1419. target_env["PACKAGED_RUNNER_EXPECTED_CWD"] = str(root)
  1420. request_path.write_text(
  1421. json.dumps({"cwd": str(root), "env": target_env}),
  1422. encoding="utf-8",
  1423. )
  1424. request_path.chmod(0o600)
  1425. environment = dict(os.environ)
  1426. environment["DSH_SUBPROCESS_RUNNER"] = str(request_path)
  1427. result = subprocess.run(
  1428. [str(executable), "--", sys.executable, "-c", target_script],
  1429. cwd=root,
  1430. env=environment,
  1431. capture_output=True,
  1432. text=True,
  1433. timeout=30,
  1434. check=False,
  1435. )
  1436. if result.returncode != 7 or request_path.exists() or (root / "startup-error.json").exists():
  1437. raise AssertionError(
  1438. "packaged POSIX runner failed: "
  1439. f"exit={result.returncode}; stdout={result.stdout!r}; stderr={result.stderr!r}"
  1440. )
  1441. return
  1442. node = shutil.which("node")
  1443. if node is None:
  1444. raise AssertionError("packaged Windows runner smoke requires node on PATH")
  1445. helper = root / "windows-runner-smoke.mjs"
  1446. helper.write_text(
  1447. """import { spawn } from 'node:child_process'
  1448. const [runtime, target, cwd, targetScript] = process.argv.slice(2)
  1449. const child = spawn(runtime, ['--', target, '-c', targetScript], {
  1450. cwd,
  1451. env: { ...process.env, DSH_SUBPROCESS_RUNNER: 'windows' },
  1452. stdio: ['ignore', 'ignore', 'ignore', 'ipc', 'pipe', 'pipe', 'pipe'],
  1453. })
  1454. const messages = []
  1455. let stdout = ''
  1456. let stderr = ''
  1457. child.stdio[4].destroy()
  1458. child.stdio[5].on('data', chunk => { stdout += chunk.toString() })
  1459. child.stdio[6].on('data', chunk => { stderr += chunk.toString() })
  1460. child.on('message', message => { messages.push(message) })
  1461. const result = await new Promise((resolve, reject) => {
  1462. child.once('error', reject)
  1463. child.once('spawn', () => {
  1464. child.send({
  1465. type: 'start',
  1466. cwd,
  1467. env: {
  1468. ...process.env,
  1469. DSH_SUBPROCESS_RUNNER: 'target-collision-restored',
  1470. PACKAGED_RUNNER_EXPECTED_CWD: cwd,
  1471. },
  1472. }, error => { if (error) reject(error) })
  1473. })
  1474. child.once('close', (exitCode, signal) => { resolve({ exitCode, signal }) })
  1475. })
  1476. process.stdout.write(JSON.stringify({ ...result, messages, stdout, stderr }))
  1477. """,
  1478. encoding="utf-8",
  1479. )
  1480. helper_result = subprocess.run(
  1481. [node, str(helper), str(executable), sys.executable, str(root), target_script],
  1482. cwd=root,
  1483. capture_output=True,
  1484. text=True,
  1485. timeout=30,
  1486. check=False,
  1487. )
  1488. if helper_result.returncode != 0:
  1489. raise AssertionError(f"packaged Windows runner helper failed: {helper_result.stderr}")
  1490. observed = json.loads(helper_result.stdout)
  1491. expected = {
  1492. "exitCode": 0,
  1493. "signal": None,
  1494. "messages": [{"type": "target-exit", "exitCode": 7}],
  1495. "stdout": "",
  1496. "stderr": "",
  1497. }
  1498. if observed != expected:
  1499. raise AssertionError(f"packaged Windows runner returned unexpected facts: {observed}")
  1500. def is_idle_notification(message: dict[str, object]) -> bool:
  1501. """Return whether a JSON-RPC notification marks a session idle."""
  1502. params = message.get("params")
  1503. return (
  1504. message.get("method") == "session.status"
  1505. and isinstance(params, dict)
  1506. and params.get("status") == "idle"
  1507. )
  1508. class RuntimePeer:
  1509. def __init__(self, argv: list[str], cwd: Path, environment: dict[str, str]) -> None:
  1510. self.process = subprocess.Popen(
  1511. argv,
  1512. cwd=cwd,
  1513. env=environment,
  1514. stdin=subprocess.PIPE,
  1515. stdout=subprocess.PIPE,
  1516. stderr=subprocess.PIPE,
  1517. text=True,
  1518. encoding="utf-8",
  1519. bufsize=1,
  1520. )
  1521. self.stdout: queue.Queue[str | None] = queue.Queue()
  1522. self.stderr: list[str] = []
  1523. threading.Thread(target=self._read_stdout, daemon=True).start()
  1524. threading.Thread(target=self._read_stderr, daemon=True).start()
  1525. def send(self, message: dict[str, object]) -> None:
  1526. if self.process.stdin is None:
  1527. raise RuntimeError("runtime stdin is unavailable")
  1528. self.process.stdin.write(json.dumps(message) + "\n")
  1529. self.process.stdin.flush()
  1530. def read_until(self, predicate: Callable[[dict[str, object]], bool]) -> list[dict[str, object]]:
  1531. deadline = time.monotonic() + 60
  1532. messages: list[dict[str, object]] = []
  1533. while time.monotonic() < deadline:
  1534. try:
  1535. line = self.stdout.get(timeout=min(0.25, deadline - time.monotonic()))
  1536. except queue.Empty:
  1537. continue
  1538. if line is None:
  1539. raise RuntimeError(f"runtime exited before expected message; stderr: {''.join(self.stderr)}")
  1540. try:
  1541. message = json.loads(line)
  1542. except json.JSONDecodeError:
  1543. continue
  1544. messages.append(message)
  1545. if predicate(message):
  1546. return messages
  1547. raise TimeoutError(f"runtime timed out; messages={messages}; stderr={''.join(self.stderr)}")
  1548. def close(self) -> None:
  1549. if self.process.stdin is not None and not self.process.stdin.closed:
  1550. self.process.stdin.close()
  1551. try:
  1552. self.process.wait(timeout=10)
  1553. except subprocess.TimeoutExpired:
  1554. self.process.kill()
  1555. self.process.wait()
  1556. if self.process.returncode not in {0, -15}:
  1557. raise RuntimeError(f"runtime exited {self.process.returncode}; stderr: {''.join(self.stderr)}")
  1558. def _read_stdout(self) -> None:
  1559. assert self.process.stdout is not None
  1560. for line in self.process.stdout:
  1561. self.stdout.put(line)
  1562. self.stdout.put(None)
  1563. def _read_stderr(self) -> None:
  1564. assert self.process.stderr is not None
  1565. self.stderr.extend(self.process.stderr)
  1566. PERSISTED_SESSION_FILENAME = re.compile(r"^session(?:\.v([1-9]\d*))?\.jsonl(\.zstd)?$")
  1567. SNAPSHOT_SESSION_FILENAME = re.compile(
  1568. r"^session(?:\.([1-9]\d*))?(?:\.v([1-9]\d*))?\.jsonl$",
  1569. )
  1570. def persisted_session_filename_version(path: Path, compressed: bool = False) -> int | None:
  1571. """Return one canonical persistence basename's generation for the selected encoding."""
  1572. match = PERSISTED_SESSION_FILENAME.fullmatch(path.name)
  1573. if match is None or (match.group(2) is not None) != compressed:
  1574. return None
  1575. return int(match.group(1) or 0)
  1576. def latest_persisted_session_paths(sessions: Path, compressed: bool = False) -> list[Path]:
  1577. """Select the numeric-highest immutable generation in each physical Session directory."""
  1578. pattern = "*.jsonl.zstd" if compressed else "*.jsonl"
  1579. selected: dict[Path, tuple[int, Path]] = {}
  1580. for path in sessions.rglob(pattern):
  1581. version = persisted_session_filename_version(path, compressed)
  1582. if version is None:
  1583. continue
  1584. previous = selected.get(path.parent)
  1585. if previous is None or version > previous[0]:
  1586. selected[path.parent] = (version, path)
  1587. return sorted((entry[1] for entry in selected.values()), key=lambda path: str(path))
  1588. def session_header_version(content: str, label: str) -> int:
  1589. """Read a non-negative physical Session generation from the first JSONL record."""
  1590. first = next((line for line in content.splitlines() if line), None)
  1591. if first is None:
  1592. raise AssertionError(f"{label}: Session log is empty")
  1593. header = json.loads(first)
  1594. version = header.get("version") if isinstance(header, dict) and header.get("type") == "session" else None
  1595. if not isinstance(version, int) or isinstance(version, bool) or version < 0:
  1596. raise AssertionError(f"{label}: Session header has no non-negative integer version")
  1597. return version
  1598. def assert_current_session_version(version: int, label: str) -> None:
  1599. """Require generated logs to use the source writer generation, independent of goldens."""
  1600. source = Path(__file__).resolve().parents[1] / "packages/core/session/src/types.ts"
  1601. declarations = re.findall(
  1602. r"^export const SESSION_FORMAT_VERSION = ([0-9]+)$",
  1603. source.read_text(encoding="utf-8"),
  1604. re.MULTILINE,
  1605. )
  1606. if len(declarations) != 1:
  1607. raise AssertionError(f"{source}: expected one literal SESSION_FORMAT_VERSION declaration")
  1608. current_version = int(declarations[0])
  1609. if version != current_version:
  1610. raise AssertionError(
  1611. f"{label}: expected current Session format v{current_version}, got v{version}",
  1612. )
  1613. def assert_persisted_session_version(path: Path, content: str) -> int:
  1614. """Require generated persistence filenames and headers to use the current generation."""
  1615. filename_version = persisted_session_filename_version(path)
  1616. if filename_version is None:
  1617. raise AssertionError(f"non-canonical Session persistence filename: {path.name}")
  1618. header_version = session_header_version(content, path.name)
  1619. if filename_version != header_version:
  1620. raise AssertionError(
  1621. f"{path.name}: filename declares Session format v{filename_version}, "
  1622. f"header declares v{header_version}",
  1623. )
  1624. assert_current_session_version(header_version, path.name)
  1625. return header_version
  1626. def snapshot_session_filename(index: int, version: int) -> str:
  1627. """Render parent/ordinal snapshot role plus an omitted-v0 generation."""
  1628. if index < 0 or version < 0:
  1629. raise ValueError("snapshot Session index and version must be non-negative")
  1630. ordinal = "" if index == 0 else f".{index}"
  1631. generation = "" if version == 0 else f".v{version}"
  1632. return f"session{ordinal}{generation}.jsonl"
  1633. def parse_snapshot_session_filename(name: str) -> tuple[int, int] | None:
  1634. """Parse one canonical parent/ordinal snapshot filename."""
  1635. match = SNAPSHOT_SESSION_FILENAME.fullmatch(name)
  1636. if match is None:
  1637. if name.startswith("session") and name.endswith(".jsonl"):
  1638. raise AssertionError(f"invalid snapshot Session filename: {name}")
  1639. return None
  1640. return int(match.group(1) or 0), int(match.group(2) or 0)
  1641. def selected_snapshot_session_files(directory: Path) -> dict[int, Path]:
  1642. """Select one highest-generation expected file per parent/ordinal role."""
  1643. selected: dict[int, tuple[int, Path]] = {}
  1644. for path in directory.iterdir():
  1645. if not path.is_file():
  1646. continue
  1647. parsed = parse_snapshot_session_filename(path.name)
  1648. if parsed is None:
  1649. continue
  1650. index, version = parsed
  1651. content = path.read_text(encoding="utf-8")
  1652. header_version = session_header_version(content, path.name)
  1653. if header_version != version:
  1654. raise AssertionError(
  1655. f"{path.name}: filename declares Session format v{version}, header declares v{header_version}",
  1656. )
  1657. previous = selected.get(index)
  1658. if previous is None or version > previous[0]:
  1659. selected[index] = (version, path)
  1660. return {index: value[1] for index, value in selected.items()}
  1661. def assert_session_log(sessions: Path, cwd: Path, *expected_texts: str) -> None:
  1662. logs = latest_persisted_session_paths(sessions)
  1663. if len(logs) != 1:
  1664. raise AssertionError(f"expected one JSONL session log under {sessions}, found {logs}")
  1665. content = logs[0].read_text()
  1666. assert_persisted_session_version(logs[0], content)
  1667. lines = content.splitlines()
  1668. header = json.loads(lines[0])
  1669. if header.get("cwd") != str(cwd):
  1670. raise AssertionError(f"session header cwd is not absolute/canonical: {header}")
  1671. rendered = "\n".join(lines)
  1672. for expected in expected_texts:
  1673. if expected not in rendered:
  1674. raise AssertionError(f"session log has no {expected!r} response: {logs[0]}")
  1675. def assert_zstd_session_log(sessions: Path) -> None:
  1676. logs = latest_persisted_session_paths(sessions, compressed=True)
  1677. if len(logs) != 1:
  1678. raise AssertionError(f"expected one Zstandard JSONL session log under {sessions}, found {logs}")
  1679. if not logs[0].read_bytes().startswith(bytes.fromhex("28b52ffd")):
  1680. raise AssertionError(f"session log has no Zstandard magic: {logs[0]}")
  1681. def read_session_logs(sessions: Path) -> dict[str, list[dict[str, object]]]:
  1682. """Parse every persisted JSONL session into a map keyed by header id."""
  1683. logs: dict[str, list[dict[str, object]]] = {}
  1684. for path in latest_persisted_session_paths(sessions):
  1685. content = path.read_text(encoding="utf-8")
  1686. assert_persisted_session_version(path, content)
  1687. records = [
  1688. json.loads(line)
  1689. for line in content.splitlines()
  1690. if line
  1691. ]
  1692. if not records or records[0].get("type") != "session":
  1693. raise AssertionError(f"session log has no header: {path}")
  1694. session_id = records[0].get("id")
  1695. if not isinstance(session_id, str):
  1696. raise AssertionError(f"session log header has no string id: {path}")
  1697. if session_id in logs:
  1698. raise AssertionError(f"duplicate persisted session id: {session_id}")
  1699. logs[session_id] = records
  1700. return logs
  1701. def snapshot_child_ids(result: "RunResult") -> list[str]:
  1702. """Return the two child session ids in their SDK notification order."""
  1703. child_ids: list[str] = []
  1704. for notification in result.notifications:
  1705. if notification.method != "subagent.started":
  1706. continue
  1707. payload = notification.payload
  1708. if payload.get("parentSessionId") != SNAPSHOT_SESSION_ID:
  1709. continue
  1710. child_id = payload.get("childSessionId")
  1711. if isinstance(child_id, str) and child_id not in child_ids:
  1712. child_ids.append(child_id)
  1713. if len(child_ids) != 2:
  1714. raise AssertionError(f"advanced snapshot expected two child session ids: {child_ids}")
  1715. return child_ids
  1716. def build_in_history_snapshot_files(
  1717. result: "RunResult",
  1718. requests: list[dict[str, object]],
  1719. log: list[dict[str, object]],
  1720. ) -> dict[str, str]:
  1721. """Assert live requests, SDK subscriptions, and persistence retain both prompts."""
  1722. systems = [event for event in result.events if event.get("type") == "system/message"]
  1723. assert len(systems) == 2, systems
  1724. assert [event.get("surfaceOp") for event in systems] == ["append", "append"], systems
  1725. prompts = [message_text(event["data"]["message"]["content"]) for event in systems]
  1726. assert "Python SDK prompt version 1." in prompts[0], prompts
  1727. assert "Python SDK prompt version 2." not in prompts[0], prompts
  1728. assert "Python SDK prompt version 2." in prompts[1], prompts
  1729. assert "Python SDK prompt version 1." not in prompts[1], prompts
  1730. assert prompts[0] != prompts[1], prompts
  1731. assert [event for event in log if event.get("type") == "system/message"] == systems
  1732. subscribed = [
  1733. notification.payload["event"]
  1734. for notification in result.notifications
  1735. if notification.method == "session.event"
  1736. and notification.payload.get("event", {}).get("type") == "system/message"
  1737. ]
  1738. assert subscribed == systems, subscribed
  1739. contexts = [event["data"] for event in result.events if event.get("type") == "request/context"]
  1740. assert contexts and all(context.get("systemPromptUpdate") == "in-history" for context in contexts), contexts
  1741. assert len([event for event in result.events if event.get("type") == "request/header"]) == 1
  1742. assert all(event.get("surfaceOp") in (None, "append") for event in result.events)
  1743. first_tool = next(index for index, event in enumerate(result.events) if event.get("type") == "tool/result")
  1744. assert result.events.index(systems[1]) > first_tool
  1745. assert len(requests) == 3, requests
  1746. request_prompts = []
  1747. for index, request in enumerate(requests):
  1748. messages = request["messages"]
  1749. assert message_text(request.get("system")) == prompts[0]
  1750. assert request["tools"] == requests[0]["tools"], "prompt update changed tool schemas"
  1751. positions = [position for position, message in enumerate(messages) if message["role"] == "system"]
  1752. texts = [message_text(request["system"]), *[
  1753. message_text(messages[position]["content"]) for position in positions
  1754. ]]
  1755. assert texts == (prompts[:1] if index == 0 else prompts), texts
  1756. if index > 0:
  1757. previous = messages[positions[0] - 1]
  1758. assert previous["role"] == "user" and any(
  1759. block.get("type") == "tool_result" for block in previous["content"]
  1760. ), messages
  1761. request_prompts.append(texts)
  1762. evidence = {
  1763. "requestSystemPrompts": request_prompts,
  1764. "systemMessageOperations": [event["surfaceOp"] for event in systems],
  1765. "subscribedSystemPrompts": prompts,
  1766. "requestContexts": contexts,
  1767. }
  1768. return {"prompt-history.json": json.dumps(evidence, indent=2, ensure_ascii=False) + "\n"}
  1769. def build_minimal_snapshot_files(
  1770. requests: list[dict[str, object]],
  1771. cwd: Path,
  1772. ) -> dict[str, str]:
  1773. """Render the minimal composition's model-visible surface as expected output.
  1774. Every assembled system prompt, advertised tool schema, and system or user message is
  1775. kept verbatim: they carry what the deployment actually shows the model, so a plugin
  1776. that contributes an unintended system section or user message cannot pass unnoticed.
  1777. Assistant and tool payloads keep only their call identity because their text differs
  1778. across the platforms this expected output must replay on. The shipped profile omits
  1779. dynamic runtime context, so every message it emits is compared.
  1780. """
  1781. snapshot = []
  1782. for body in requests:
  1783. messages = body.get("messages")
  1784. if not isinstance(messages, list):
  1785. raise AssertionError(f"minimal model request has no messages: {body}")
  1786. snapshot.append({
  1787. "system": minimal_snapshot_text(body.get("system"), cwd),
  1788. "tools": minimal_snapshot_text(body.get("tools"), cwd),
  1789. "messages": [
  1790. minimal_snapshot_message(message, cwd)
  1791. for message in messages
  1792. ],
  1793. })
  1794. return {"model-visible.json": json.dumps(snapshot, indent=2, ensure_ascii=False) + "\n"}
  1795. def minimal_snapshot_message(message: object, cwd: Path) -> dict[str, object]:
  1796. """Reduce one model-visible message to its stable, behavior-carrying parts."""
  1797. if not isinstance(message, dict):
  1798. raise AssertionError(f"minimal model request has an invalid message: {message}")
  1799. role = message.get("role")
  1800. if role == "system":
  1801. return {"role": role, "text": minimal_snapshot_text(message_text(message.get("content")), cwd)}
  1802. if role == "user":
  1803. content = []
  1804. for block in message.get("content", []):
  1805. if block.get("type") == "tool_result":
  1806. content.append({"type": "tool_result", "tool_use_id": block.get("tool_use_id"), "content": "{{tool-result}}"})
  1807. elif block.get("type") == "text":
  1808. content.append(minimal_snapshot_text(block, cwd))
  1809. else:
  1810. raise AssertionError(f"minimal user message has unexpected content: {block}")
  1811. return {"role": role, "content": content}
  1812. if role == "assistant":
  1813. calls = message.get("content")
  1814. if not isinstance(calls, list):
  1815. raise AssertionError(f"minimal assistant message has no tool calls: {message}")
  1816. return {
  1817. "role": role,
  1818. "toolCalls": [
  1819. {"id": call.get("id"), "name": call.get("name")}
  1820. for call in calls
  1821. if isinstance(call, dict) and call.get("type") == "tool_use"
  1822. ],
  1823. }
  1824. raise AssertionError(f"minimal model request has an unexpected message role: {message}")
  1825. def minimal_snapshot_text(value: object, cwd: Path) -> object:
  1826. """Replace the scenario's temporary working directory everywhere it appears."""
  1827. if isinstance(value, str):
  1828. return value.replace(str(cwd), "{{cwd}}")
  1829. if isinstance(value, list):
  1830. return [minimal_snapshot_text(item, cwd) for item in value]
  1831. if isinstance(value, dict):
  1832. return {key: minimal_snapshot_text(item, cwd) for key, item in value.items()}
  1833. return value
  1834. def build_snapshot_files(
  1835. result: "RunResult",
  1836. logs: dict[str, list[dict[str, object]]],
  1837. child_ids: list[str],
  1838. cwd: Path,
  1839. ) -> dict[str, str]:
  1840. """Render the SDK result and three persisted logs into stable expected outputs."""
  1841. replacements = [(str(cwd), "{{cwd}}"), (SNAPSHOT_SESSION_ID, "{{parent}}")]
  1842. replacements.append((snapshot_workflow_run_id(result), "{{workflow-run}}"))
  1843. for index, child_id in enumerate(child_ids, start=1):
  1844. replacements.append((child_id, f"{{{{child-{index}}}}}"))
  1845. agent_id = snapshot_agent_id(result, child_id)
  1846. replacements.append((agent_id, f"{{{{agent-{index}}}}}"))
  1847. command_index = 0
  1848. for record in logs[SNAPSHOT_SESSION_ID]:
  1849. data = record.get("data")
  1850. if record.get("type") == "command/run" and isinstance(data, dict):
  1851. command_index += 1
  1852. replacements.append((data["commandId"], f"{{{{command:{command_index}}}}}"))
  1853. if record.get("type") == "command/done" and isinstance(data, dict):
  1854. anonymous = re.search(r"Anonymous user: ([0-9a-f-]{36})", str(data.get("text")))
  1855. if anonymous is not None:
  1856. replacements.append((anonymous.group(1), "{{anonymous-user}}"))
  1857. feedback_targets = dict.fromkeys(
  1858. record["data"]["item"]["messageId"]
  1859. for record in logs[SNAPSHOT_SESSION_ID]
  1860. if record.get("type") == "feedback/message-put"
  1861. )
  1862. for index, message_id in enumerate(feedback_targets, start=1):
  1863. replacements.append((message_id, f"{{{{message:{index}}}}}"))
  1864. feedback_versions = dict.fromkeys(
  1865. record["data"]["item"]["version"]
  1866. for record in logs[SNAPSHOT_SESSION_ID]
  1867. if record.get("type") == "feedback/message-put"
  1868. )
  1869. for index, version in enumerate(feedback_versions, start=1):
  1870. replacements.append((version, f"{{{{feedback-version:{index}}}}}"))
  1871. replacements.sort(key=lambda pair: len(pair[0]), reverse=True)
  1872. result_value = {
  1873. "session_id": result.session_id,
  1874. "final_response": result.final_response,
  1875. "events": result.events,
  1876. "notifications": [
  1877. {"method": notification.method, "payload": notification.payload}
  1878. for notification in result.notifications
  1879. ],
  1880. }
  1881. normalized_result = normalize_snapshot_value(result_value, replacements)
  1882. parent_records = project_session_snapshot([
  1883. normalize_snapshot_value(record, replacements) for record in logs[SNAPSHOT_SESSION_ID]
  1884. ])
  1885. files = {
  1886. "result.json": json.dumps(normalized_result, indent=2, ensure_ascii=False) + "\n",
  1887. snapshot_session_filename(
  1888. 0, session_header_version(render_jsonl(parent_records), "advanced parent"),
  1889. ): render_jsonl(parent_records),
  1890. }
  1891. for index, child_id in enumerate(child_ids, start=1):
  1892. child_records = project_session_snapshot([
  1893. normalize_snapshot_value(record, replacements) for record in logs[child_id]
  1894. ])
  1895. child_content = render_jsonl(child_records)
  1896. files[snapshot_session_filename(
  1897. index, session_header_version(child_content, f"advanced child {index}"),
  1898. )] = child_content
  1899. return files
  1900. def build_restart_snapshot_files(
  1901. first: "RunResult",
  1902. second: "RunResult",
  1903. requests: list[dict[str, object]],
  1904. logs: dict[str, list[dict[str, object]]],
  1905. cwd: Path,
  1906. sessions: Path,
  1907. ) -> dict[str, str]:
  1908. """Render two SDK processes, isolated model histories, and durable logs."""
  1909. replacements = [
  1910. (str(sessions), "{{sessions}}"),
  1911. (str(cwd), "{{cwd}}"),
  1912. (RESTART_FIRST_SESSION_ID, "{{session-1}}"),
  1913. (RESTART_SECOND_SESSION_ID, "{{session-2}}"),
  1914. ]
  1915. result_value = [
  1916. {
  1917. "session_id": result.session_id,
  1918. "final_response": result.final_response,
  1919. "finish_reason": result.finish_reason,
  1920. "eventTypes": [
  1921. event.get("type")
  1922. for event in result.events
  1923. ],
  1924. "notificationMethods": [
  1925. notification.method
  1926. for notification in result.notifications
  1927. ],
  1928. }
  1929. for result in (first, second)
  1930. ]
  1931. request_value = [
  1932. {
  1933. "model": request.get("model"),
  1934. "system": "{{system}}" if request.get("system") else None,
  1935. "messages": restart_request_messages(request),
  1936. "toolNames": sorted(advertised_tool_names(request)),
  1937. }
  1938. for request in requests
  1939. ]
  1940. first_records = project_session_snapshot([
  1941. normalize_snapshot_value(record, replacements) for record in logs[RESTART_FIRST_SESSION_ID]
  1942. ])
  1943. second_records = project_session_snapshot([
  1944. normalize_snapshot_value(record, replacements) for record in logs[RESTART_SECOND_SESSION_ID]
  1945. ])
  1946. first_content = render_jsonl(first_records)
  1947. second_content = render_jsonl(second_records)
  1948. return {
  1949. "result.json": json.dumps(
  1950. normalize_snapshot_value(result_value, replacements), indent=2, ensure_ascii=False,
  1951. ) + "\n",
  1952. "requests.json": json.dumps(
  1953. normalize_snapshot_value(request_value, replacements), indent=2, ensure_ascii=False,
  1954. ) + "\n",
  1955. snapshot_session_filename(
  1956. 1, session_header_version(first_content, "restart Session 1"),
  1957. ): first_content,
  1958. snapshot_session_filename(
  1959. 2, session_header_version(second_content, "restart Session 2"),
  1960. ): second_content,
  1961. }
  1962. def restart_request_messages(request: dict[str, object]) -> list[object]:
  1963. """Project model history while tokenizing composition-owned system prose."""
  1964. messages = request.get("messages")
  1965. if not isinstance(messages, list):
  1966. raise AssertionError(f"restart snapshot request has no messages: {request}")
  1967. return [
  1968. {"role": "system", "content": "{{system}}"}
  1969. if isinstance(message, dict) and message.get("role") == "system"
  1970. else message
  1971. for message in messages
  1972. ]
  1973. def snapshot_workflow_run_id(result: "RunResult") -> str:
  1974. """Return the one workflow run id emitted by the advanced scenario."""
  1975. run_ids: set[str] = set()
  1976. for event in result.events:
  1977. event_type = event.get("type")
  1978. data = event.get("data")
  1979. if not isinstance(event_type, str) or not event_type.startswith("tool-workflow/"):
  1980. continue
  1981. if isinstance(data, dict) and isinstance(data.get("runId"), str):
  1982. run_ids.add(data["runId"])
  1983. if len(run_ids) != 1:
  1984. raise AssertionError(f"advanced snapshot expected one workflow run id: {sorted(run_ids)}")
  1985. return next(iter(run_ids))
  1986. def snapshot_agent_id(result: "RunResult", child_id: str) -> str:
  1987. """Find the successful subagent id paired with one child session."""
  1988. for notification in result.notifications:
  1989. if notification.method != "subagent.finished":
  1990. continue
  1991. payload = notification.payload
  1992. if payload.get("childSessionId") != child_id:
  1993. continue
  1994. if payload.get("provider") != "spawn" or payload.get("status") != "ok":
  1995. raise AssertionError(f"advanced child did not finish successfully: {payload}")
  1996. agent_id = payload.get("agentId")
  1997. if isinstance(agent_id, str):
  1998. return agent_id
  1999. raise AssertionError(f"advanced snapshot has no finished agent for child {child_id}")
  2000. def normalize_snapshot_value(
  2001. value: object,
  2002. replacements: list[tuple[str, str]],
  2003. ) -> object:
  2004. """Scrub volatile values and bulky request headers without losing behavior."""
  2005. if isinstance(value, str):
  2006. normalized = value
  2007. for actual, token in replacements:
  2008. normalized = normalized.replace(actual, token)
  2009. return normalized
  2010. if isinstance(value, list):
  2011. return [normalize_snapshot_value(item, replacements) for item in value]
  2012. if not isinstance(value, dict):
  2013. return value
  2014. normalized = {
  2015. key: normalize_snapshot_value(item, replacements)
  2016. for key, item in value.items()
  2017. }
  2018. if normalized.get("type") == "session" and "createdAt" in normalized:
  2019. normalized["createdAt"] = 0
  2020. if normalized.get("type") == "subagent/catalog":
  2021. data = normalized.get("data")
  2022. if isinstance(data, dict) and "childCreatedAt" in data:
  2023. data["childCreatedAt"] = 0
  2024. if "seq" in normalized and "time" in normalized:
  2025. normalized["time"] = 0
  2026. if normalized.get("type") in ("assistant/message", "assistant/attempt"):
  2027. data = normalized.get("data")
  2028. stream = data.get("stream") if isinstance(data, dict) else None
  2029. if isinstance(stream, list):
  2030. for member in stream:
  2031. if not isinstance(member, dict):
  2032. continue
  2033. if isinstance(member.get("time"), (int, float)):
  2034. member["time"] = 0
  2035. if isinstance(member.get("time0"), (int, float)):
  2036. member["time0"] = 0
  2037. dt = member.get("dt")
  2038. if isinstance(dt, list):
  2039. member["dt"] = [0] * len(dt)
  2040. if isinstance(normalized.get("id"), str) and normalized.get("role") in ("assistant", "system", "user"):
  2041. if not normalized["id"].startswith("{{message:"):
  2042. normalized["id"] = "{{messageId}}"
  2043. if normalized.get("type") in ("feedback/message-put", "feedback/message-delete"):
  2044. data = normalized.get("data")
  2045. if isinstance(data, dict):
  2046. item = data.get("item") if normalized["type"] == "feedback/message-put" else data
  2047. if isinstance(item, dict):
  2048. if normalized["type"] == "feedback/message-put":
  2049. item["createdAt"] = 0
  2050. item["updatedAt"] = 0
  2051. scrub_snapshot_header(normalized)
  2052. scrub_snapshot_system_message(normalized)
  2053. return normalized
  2054. def scrub_snapshot_header(value: dict[object, object]) -> None:
  2055. """Tokenize full request-header tool schemas while retaining tool names."""
  2056. data = value.get("data")
  2057. if not isinstance(data, dict):
  2058. return
  2059. if value.get("type") == "request/header":
  2060. header = data.get("header")
  2061. if not isinstance(header, dict):
  2062. return
  2063. tools = header.get("tools")
  2064. if isinstance(tools, list):
  2065. header["tools"] = [
  2066. tool.get("name") if isinstance(tool, dict) else "{{tools}}"
  2067. for tool in tools
  2068. ]
  2069. def scrub_snapshot_system_message(value: dict[object, object]) -> None:
  2070. """Tokenize the rendered prompt text of a `system/message` surface node."""
  2071. if value.get("type") != "system/message":
  2072. return
  2073. data = value.get("data")
  2074. message = data.get("message") if isinstance(data, dict) else None
  2075. content = message.get("content") if isinstance(message, dict) else None
  2076. if not isinstance(content, list):
  2077. return
  2078. for block in content:
  2079. if isinstance(block, dict) and block.get("type") == "text":
  2080. block["text"] = "{{system}}"
  2081. def render_jsonl(records: list[object]) -> str:
  2082. """Render parsed JSON values as compact, newline-terminated JSONL."""
  2083. return "".join(
  2084. json.dumps(record, ensure_ascii=False, separators=(",", ":")) + "\n"
  2085. for record in records
  2086. )
  2087. def project_session_snapshot(records: list[dict[str, object]]) -> list[dict[str, object]]:
  2088. """Omit storage sequence/time envelopes from snapshot body records."""
  2089. projected = [dict(record) for record in records]
  2090. for record in projected[1:]:
  2091. for key in ("seq", "time", "seq0", "time0"):
  2092. record.pop(key, None)
  2093. return projected
  2094. SESSION_FORMAT_TOKEN = "{{sessionFormatVersion}}"
  2095. def expand_snapshot_stream_member(member: object) -> list[dict[str, object]]:
  2096. """Expand one compact Assistant stream member into logical provider chunks."""
  2097. if not isinstance(member, dict):
  2098. raise AssertionError(f"snapshot Assistant stream member is not an object: {member!r}")
  2099. member_type = member.get("type")
  2100. if member_type == "chunk":
  2101. chunk = member.get("chunk")
  2102. if not isinstance(chunk, dict):
  2103. raise AssertionError(f"snapshot Assistant chunk member has no chunk: {member!r}")
  2104. return [chunk]
  2105. packed_kinds = {
  2106. "text-chunks": ("texts", "text-delta", "text"),
  2107. "reasoning-chunks": ("texts", "reasoning-delta", "text"),
  2108. "tool-call-chunks": ("args", "tool-call-delta", "argumentsDelta"),
  2109. }
  2110. packed = packed_kinds.get(member_type)
  2111. if packed is None:
  2112. raise AssertionError(f"snapshot Assistant stream has unknown member type: {member_type!r}")
  2113. values_key, chunk_type, value_key = packed
  2114. values = member.get(values_key)
  2115. if not isinstance(values, list):
  2116. raise AssertionError(f"snapshot Assistant stream member has no {values_key}: {member!r}")
  2117. shared = {
  2118. key: member[key]
  2119. for key in ("index", "id", "name")
  2120. if key in member
  2121. }
  2122. return [
  2123. {"type": chunk_type, **shared, value_key: value}
  2124. for value in values
  2125. ]
  2126. def expand_snapshot_assistant_event(value: object) -> list[object]:
  2127. """Expand one direct or SDK-wrapped v2 settlement for generation-neutral comparison."""
  2128. if not isinstance(value, dict):
  2129. return [value]
  2130. event = value
  2131. wrapper_key: str | None = None
  2132. wrapper: dict[str, object] | None = None
  2133. if value.get("method") == "session.event":
  2134. for candidate in ("payload", "params"):
  2135. container = value.get(candidate)
  2136. nested = container.get("event") if isinstance(container, dict) else None
  2137. if isinstance(nested, dict):
  2138. event = nested
  2139. wrapper_key = candidate
  2140. wrapper = container
  2141. break
  2142. if event.get("type") not in ("assistant/message", "assistant/attempt"):
  2143. return [value]
  2144. data = event.get("data")
  2145. stream = data.get("stream") if isinstance(data, dict) else None
  2146. if not isinstance(stream, list):
  2147. return [value]
  2148. def wrap(expanded: dict[str, object]) -> object:
  2149. if wrapper_key is None or wrapper is None:
  2150. return expanded
  2151. return {**value, wrapper_key: {**wrapper, "event": expanded}}
  2152. common = {
  2153. key: data[key]
  2154. for key in ("turn", "step")
  2155. if key in data
  2156. }
  2157. expanded = [
  2158. wrap({
  2159. "type": "assistant/chunk",
  2160. "data": {**common, "chunk": chunk},
  2161. })
  2162. for member in stream
  2163. for chunk in expand_snapshot_stream_member(member)
  2164. ]
  2165. if event.get("type") == "assistant/message":
  2166. expanded.append(wrap({
  2167. **event,
  2168. "data": {key: item for key, item in data.items() if key != "stream"},
  2169. }))
  2170. return expanded
  2171. def normalize_session_format_comparison(
  2172. value: object,
  2173. source_session_version: int | None = None,
  2174. ) -> object:
  2175. """Canonicalize only generation metadata that differs across immutable Session files."""
  2176. if isinstance(value, list):
  2177. return [
  2178. normalize_session_format_comparison(expanded, source_session_version)
  2179. for item in value
  2180. for expanded in expand_snapshot_assistant_event(item)
  2181. ]
  2182. if not isinstance(value, dict):
  2183. return value
  2184. normalized = {
  2185. key: normalize_session_format_comparison(item, source_session_version)
  2186. for key, item in value.items()
  2187. }
  2188. if normalized.get("type") == "session" and "version" in normalized:
  2189. normalized["version"] = SESSION_FORMAT_TOKEN
  2190. normalized.setdefault("isSeeded", False)
  2191. ordered_header = {
  2192. key: normalized[key]
  2193. for key in ("type", "version", "id", "createdAt", "cwd", "isSeeded", "delegationDepth")
  2194. if key in normalized
  2195. }
  2196. normalized = {
  2197. **ordered_header,
  2198. **{key: item for key, item in normalized.items() if key not in ordered_header},
  2199. }
  2200. if isinstance(normalized.get("type"), str) and "data" in normalized:
  2201. normalized.pop("seq", None)
  2202. normalized.pop("time", None)
  2203. if source_session_version == 1 and normalized.get("type") == "assistant/message":
  2204. normalized.pop("sourceEventSeqs", None)
  2205. return normalized
  2206. def normalize_snapshot_comparison_text(name: str, content: str) -> str:
  2207. """Normalize Session generation metadata only while comparing committed expected outputs."""
  2208. if name.startswith("session") and name.endswith(".jsonl"):
  2209. parsed = [json.loads(line) for line in content.splitlines() if line]
  2210. header = parsed[0] if parsed else None
  2211. source_version = header.get("version") if isinstance(header, dict) else None
  2212. if not isinstance(source_version, int):
  2213. raise AssertionError(f"{name}: snapshot Session header has no integer format version")
  2214. records = [
  2215. normalize_session_format_comparison(expanded, source_version)
  2216. for record in parsed
  2217. for expanded in expand_snapshot_assistant_event(record)
  2218. ]
  2219. return render_jsonl(records)
  2220. if name.endswith(".json"):
  2221. return json.dumps(
  2222. normalize_session_format_comparison(json.loads(content)),
  2223. indent=2,
  2224. ensure_ascii=False,
  2225. ) + "\n"
  2226. return content
  2227. def compare_snapshot_files(
  2228. files: dict[str, str],
  2229. update: bool,
  2230. directory: Path,
  2231. filenames: tuple[str, ...],
  2232. ) -> None:
  2233. """Compare ordered artifact roles and Session content across generations, or write generated filenames."""
  2234. scenario = directory.name
  2235. def role_name(name: str) -> str:
  2236. parsed = parse_snapshot_session_filename(name)
  2237. return name if parsed is None else snapshot_session_filename(parsed[0], 0)
  2238. if tuple(map(role_name, files)) != tuple(map(role_name, filenames)):
  2239. raise AssertionError(f"{scenario} snapshot builder produced {tuple(files)}, expected {filenames}")
  2240. for name, content in files.items():
  2241. if parse_snapshot_session_filename(name) is not None:
  2242. assert_current_session_version(session_header_version(content, name), name)
  2243. if update:
  2244. directory.mkdir(parents=True, exist_ok=True)
  2245. for name, content in files.items():
  2246. (directory / name).write_text(content, encoding="utf-8", newline="\n")
  2247. print(f"smoke-python-runtime: updated snapshots in {directory}")
  2248. existing = [path for path in directory.iterdir() if path.is_file()] if directory.is_dir() else []
  2249. expected_non_session = {
  2250. name for name in filenames if parse_snapshot_session_filename(name) is None
  2251. }
  2252. existing_non_session = {
  2253. path.name for path in existing if parse_snapshot_session_filename(path.name) is None
  2254. }
  2255. if existing_non_session != expected_non_session:
  2256. raise AssertionError(
  2257. f"{scenario} snapshot files differ: "
  2258. f"missing={sorted(expected_non_session - existing_non_session)}, "
  2259. f"unexpected={sorted(existing_non_session - expected_non_session)}"
  2260. )
  2261. selected_expected = selected_snapshot_session_files(directory)
  2262. actual_sessions: dict[int, tuple[str, str]] = {}
  2263. for name, content in files.items():
  2264. parsed = parse_snapshot_session_filename(name)
  2265. if parsed is None:
  2266. continue
  2267. index, filename_version = parsed
  2268. header_version = session_header_version(content, name)
  2269. if filename_version != header_version:
  2270. raise AssertionError(
  2271. f"{name}: filename declares Session format v{filename_version}, "
  2272. f"header declares v{header_version}",
  2273. )
  2274. if index in actual_sessions:
  2275. raise AssertionError(f"{scenario} snapshot builder produced duplicate Session role {index}")
  2276. actual_sessions[index] = (name, content)
  2277. if set(selected_expected) != set(actual_sessions):
  2278. raise AssertionError(
  2279. f"{scenario} snapshot Session roles differ: "
  2280. f"expected={sorted(selected_expected)}, actual={sorted(actual_sessions)}",
  2281. )
  2282. for name, actual in files.items():
  2283. parsed = parse_snapshot_session_filename(name)
  2284. expected_path = directory / name if parsed is None else selected_expected[parsed[0]]
  2285. expected_text = expected_path.read_text(encoding="utf-8")
  2286. compared_actual = normalize_snapshot_comparison_text(name, actual)
  2287. compared_expected = normalize_snapshot_comparison_text(expected_path.name, expected_text)
  2288. if compared_actual == compared_expected:
  2289. continue
  2290. diff = "".join(difflib.unified_diff(
  2291. compared_expected.splitlines(keepends=True),
  2292. compared_actual.splitlines(keepends=True),
  2293. fromfile=f"expected/{expected_path.name}",
  2294. tofile=f"actual/{name}",
  2295. ))
  2296. raise AssertionError(
  2297. f"{scenario} executable snapshot mismatch in {name}; "
  2298. "rerun with --update-snapshots after reviewing the behavior\n"
  2299. f"{diff}"
  2300. )
  2301. if __name__ == "__main__":
  2302. main()