| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177 |
- #!/usr/bin/env python3
- """Keyless full-turn and snapshot smoke for the Python SDK runtime."""
- from __future__ import annotations
- import argparse
- import difflib
- import importlib
- import importlib.metadata
- import json
- import os
- import queue
- import re
- import shutil
- import subprocess
- import sys
- import sysconfig
- import tempfile
- import threading
- import time
- from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
- from pathlib import Path
- from typing import TYPE_CHECKING, Callable
- if TYPE_CHECKING:
- from deepseek_harness import RunResult
- EXPECTED_TEXT = "runtime smoke ok"
- LIVE_API_SENTINEL = "PYTHON_SDK_LIVE_OK"
- CODE_PROMPT = "Use run_code to compute the packaged worker smoke value."
- CODE_WORKER_TEXT = "code worker smoke ok"
- WORKFLOW_PROMPT = "Use workflow to compute the packaged worker smoke value without agents."
- WORKFLOW_WORKER_TEXT = "workflow worker smoke ok"
- MINIMAL_PROMPT = "Exercise the packaged minimal agent's persistent shell and string-replacement editor."
- MINIMAL_TEXT = "minimal agent smoke ok"
- MINIMAL_EDITOR_PATH_PREFIX = "Editor path: "
- FS_SEARCH_PROMPT = "Exercise the packaged filesystem search tools."
- FS_SEARCH_TEXT = "filesystem search smoke ok"
- FS_SEARCH_MARKER = "PACKAGED_FS_SEARCH_OK"
- MCP_PROMPT = "Exercise the packaged MCP client with one external stdio server."
- MCP_TEXT = "MCP client smoke ok"
- PROFILE_PLUGIN_PROMPT = "Verify the Python-installed dsh profile plugin."
- PROFILE_PLUGIN_TEXT = "profile plugin smoke ok"
- PROFILE_PLUGIN_MARKER = "PYTHON_INSTALLED_DSH_PROFILE_PLUGIN"
- IS_WINDOWS = sys.platform == "win32"
- MINIMAL_SHELL_TOOL = "pwsh" if IS_WINDOWS else "bash"
- MINIMAL_SHELL_COMMAND = (
- "$global:dshSdkCounter = [int]$global:dshSdkCounter + 1; "
- 'Write-Output "COUNT=$global:dshSdkCounter CWD=$((Get-Location).Path)"; '
- "if ($global:dshSdkCounter -eq 1) { Set-Location $env:TEMP }"
- if IS_WINDOWS
- else (
- "counter=$(( ${counter:-0} + 1 )); export counter; "
- "printf 'COUNT=%s CWD=%s\\n' \"$counter\" \"$PWD\"; "
- "if [ \"$counter\" -eq 1 ]; then cd /tmp; fi"
- )
- )
- MINIMAL_SHELL_SECOND_CWD = str(Path(tempfile.gettempdir()).resolve()) if IS_WINDOWS else "/tmp"
- SPAWN_NODE_PROMPT = "Run node --version through the packaged shell tool."
- SPAWN_NODE_TEXT = "spawn node smoke ok"
- SPAWN_NODE_CALL_ID = "spawn-node-shell"
- # The POSIX command string starts with `node ` inside the shell tool's `bash -c`
- # argv, the exact form @yao-pkg/pkg's unpatched SEA bootstrap rewrites to the
- # executable itself while stamping PKG_EXECPATH into the child environment.
- SPAWN_NODE_COMMAND = (
- 'node --version; if ($env:PKG_EXECPATH) { "PKG_EXECPATH=$env:PKG_EXECPATH" } else { "PKG_EXECPATH=ABSENT" }'
- if IS_WINDOWS
- else 'node --version; echo "PKG_EXECPATH=${PKG_EXECPATH:-ABSENT}"'
- )
- LEGACY_CUSTOM_DISABLED_ROWS = (
- "agent-instructions",
- "goal",
- "goal-round-driver",
- "command-goal",
- "plan-mode",
- "skill",
- "skill-filesystem",
- "tool-fs",
- "tool-fs-search",
- "tool-goal",
- "tool-ralph",
- "tool-skill",
- "tool-str-replace-editor",
- "tool-subagent-control",
- "tool-subagent-list-agents",
- "tool-subagent-fork",
- "tool-todo",
- "tool-web",
- )
- SNAPSHOT_PROMPT = "Run the advanced packaged-runtime snapshot scenario."
- SNAPSHOT_SESSION_ID = "advanced-executable"
- SNAPSHOT_DIRECT_CHILD_PROMPT = "Reply with exactly DIRECT_CHILD_OK and nothing else."
- SNAPSHOT_WORKFLOW_CHILD_PROMPT = "Reply with exactly WORKFLOW_CHILD_OK and nothing else."
- SNAPSHOT_FINAL_TEXT = "ADVANCED_EXECUTABLE_OK"
- RESTART_FIRST_PROMPT = "Complete the first isolated Python SDK process turn."
- RESTART_FIRST_TEXT = "PROCESS_ONE_OK"
- RESTART_SECOND_PROMPT = "Complete the second isolated Python SDK process turn."
- RESTART_SECOND_TEXT = "PROCESS_TWO_OK"
- RESTART_FIRST_SESSION_ID = "process-one"
- RESTART_SECOND_SESSION_ID = "process-two"
- SNAPSHOT_PLUGIN_CODE = """\
- return (ctx) => {
- harness.registerTool(ctx, harness.defineTool({
- name: 'snapshot_double',
- description: 'Double a number for executable snapshot verification.',
- parameters: { value: { type: 'number', required: true } },
- output: {
- schema: { type: 'number' },
- render(_args, value) {
- return [{ type: 'text', text: String(value) }]
- }
- },
- async execute(args) {
- return args.value * 2
- }
- }))
- }
- """
- SNAPSHOT_WORKFLOW_SCRIPT = (
- "phase('Delegate')\n"
- f"const reply = await agent('{SNAPSHOT_WORKFLOW_CHILD_PROMPT}', {{ label: 'workflow-child' }})\n"
- "return { reply }"
- )
- ADVANCED_SNAPSHOT_DIRECTORY = (
- Path(__file__).resolve().parent / "snapshots" / "python-sdk-single-exe" / "advanced"
- )
- ADVANCED_SNAPSHOT_FILENAMES = (
- "result.json", "session.v2.jsonl", "session.1.v2.jsonl", "session.2.v2.jsonl",
- )
- MINIMAL_SNAPSHOT_DIRECTORY = (
- Path(__file__).resolve().parent / "snapshots" / "python-sdk-single-exe" / "minimal"
- )
- if IS_WINDOWS:
- MINIMAL_SNAPSHOT_DIRECTORY /= "win-x64"
- MINIMAL_SNAPSHOT_FILENAMES = ("model-visible.json",)
- RESTART_SNAPSHOT_DIRECTORY = (
- Path(__file__).resolve().parent / "snapshots" / "python-sdk-single-exe" / "restart"
- )
- RESTART_SNAPSHOT_FILENAMES = (
- "result.json", "requests.json", "session.1.v2.jsonl", "session.2.v2.jsonl",
- )
- MCP_SERVER_SCRIPT = """\
- import json
- import os
- import sys
- import time
- log_path = os.environ.get("MCP_SMOKE_LOG")
- def send(message):
- sys.stdout.write(json.dumps(message, separators=(",", ":")) + "\\n")
- sys.stdout.flush()
- for line in sys.stdin:
- request = json.loads(line)
- if log_path is not None:
- with open(log_path, "a", encoding="utf-8") as log:
- log.write(str(request.get("method")) + "\\n")
- request_id = request.get("id")
- if request_id is None:
- continue
- method = request.get("method")
- if method == "initialize":
- send({
- "jsonrpc": "2.0",
- "id": request_id,
- "result": {
- "protocolVersion": request["params"]["protocolVersion"],
- "capabilities": {"tools": {"listChanged": False}},
- "serverInfo": {"name": "python-wheel-fixture", "version": "1.0.0"},
- },
- })
- elif method == "tools/list":
- # Keep discovery pending long enough that an SDK runtime answering
- # initialize before discovery completes makes its first model request
- # without this tool and fails deterministically.
- time.sleep(0.25)
- send({
- "jsonrpc": "2.0",
- "id": request_id,
- "result": {
- "tools": [{
- "name": "add",
- "description": "Add two numbers.",
- "inputSchema": {
- "type": "object",
- "properties": {"a": {"type": "number"}, "b": {"type": "number"}},
- "required": ["a", "b"],
- "additionalProperties": False,
- },
- }],
- },
- })
- elif method == "tools/call":
- params = request["params"]
- if params.get("name") != "add" or params.get("arguments") != {"a": 19, "b": 23}:
- send({
- "jsonrpc": "2.0",
- "id": request_id,
- "error": {"code": -32602, "message": "unexpected tool call"},
- })
- continue
- send({
- "jsonrpc": "2.0",
- "id": request_id,
- "result": {"content": [{"type": "text", "text": "42"}]},
- })
- else:
- send({
- "jsonrpc": "2.0",
- "id": request_id,
- "error": {"code": -32601, "message": f"unsupported method: {method}"},
- })
- """
- def write_profile_patch(
- root: Path,
- name: str,
- sessions: Path,
- patches: list[dict[str, object]],
- ) -> Path:
- """Write one JSON-form dsh profile patch with deterministic persistence."""
- path = root / name
- path.write_text(json.dumps([
- {
- "id": "session-persistence-jsonl",
- "config": {"root": str(sessions), "compression": "none"},
- },
- {"id": "session-telemetry-otel", "disabled": True},
- *patches,
- ], indent=2))
- return path
- def write_advanced_profile_patch(root: Path, name: str, sessions: Path) -> Path:
- """Write the shared custom, snapshot, and restart profile patch."""
- return write_profile_patch(root, name, sessions, [
- {"id": "tools", "config": {"mode": "both"}},
- {
- "id": "system-prompt",
- "config": {
- "persona": "You are a coding agent powered by the {{model}} model. Your working directory is {{cwd}}.",
- },
- },
- {"id": "session-log-deepseek", "config": {"enabled": True}},
- *({"id": row_id, "disabled": True} for row_id in LEGACY_CUSTOM_DISABLED_ROWS),
- {"id": "tool-bash", "disabled": True},
- {"id": "tool-pwsh", "disabled": True},
- {
- "id": "tool-subagent",
- "config": {
- "provider": "spawn",
- "toolName": "subagent",
- "backgroundMode": "one-shot",
- },
- },
- {"insert": [
- {"id": "code-runtime", "name": "@deepseek-ai/dsh-code-runtime-worker-thread"},
- {"id": "cordis-host-runner", "name": "@deepseek-ai/dsh-cordis-host-runner"},
- {"id": "cordis-tool", "name": "@deepseek-ai/dsh-tool-cordis"},
- ]},
- ])
- def write_mcp_patch(root: Path, sessions: Path, server_script: Path) -> Path:
- """Write a profile patch that mounts the packaged MCP client."""
- return write_profile_patch(root, "mcp.patch.yml", sessions, [{
- "insert": [{
- "id": "mcp-fixture",
- "name": "@deepseek-ai/dsh-mcp-client",
- "config": {
- "serverName": "fixture",
- "transport": "stdio",
- "command": sys.executable,
- "args": [str(server_script)],
- "env": {"MCP_SMOKE_LOG": str(server_script.with_suffix(".log"))},
- "failOnStartupError": True,
- "reconnect": {"enabled": False},
- },
- }],
- }])
- class MockModelHandler(BaseHTTPRequestHandler):
- """Return deterministic text, worker, and orchestration completions."""
- requests: list[dict[str, object]] = []
- def do_POST(self) -> None:
- content_length = int(self.headers.get("content-length", "0"))
- body = json.loads(self.rfile.read(content_length))
- self.requests.append(body)
- self.send_response(200)
- self.send_header("content-type", "text/event-stream")
- self.end_headers()
- chunks = completion_chunks(body)
- for chunk in chunks:
- self.wfile.write(f"data: {json.dumps(chunk)}\n\n".encode())
- self.wfile.write(b"data: [DONE]\n\n")
- self.wfile.flush()
- def log_message(self, _format: str, *_args: object) -> None:
- return
- def completion_chunks(body: dict[str, object]) -> list[dict[str, object]]:
- """Choose the next deterministic model response from request history."""
- messages = body.get("messages")
- if not isinstance(messages, list) or not messages:
- raise AssertionError(f"model request has no messages: {body}")
- latest = messages[-1]
- if not isinstance(latest, dict):
- raise AssertionError(f"model request has an invalid latest message: {body}")
- if latest.get("role") == "tool":
- call_id, tool_name = latest_tool_call(messages)
- tool_text = message_text(latest.get("content"))
- mcp = mcp_tool_followup(call_id, tool_name, tool_text)
- if mcp is not None:
- return mcp
- fs_search = fs_search_tool_followup(call_id, tool_name, tool_text)
- if fs_search is not None:
- return fs_search
- spawn_node = spawn_node_tool_followup(call_id, tool_name, tool_text)
- if spawn_node is not None:
- return spawn_node
- minimal = minimal_tool_followup(body, call_id, tool_name, tool_text)
- if minimal is not None:
- return minimal
- advanced = advanced_tool_followup(body, call_id, tool_name, tool_text)
- if advanced is not None:
- return advanced
- if "42" not in tool_text:
- raise AssertionError(f"{tool_name} worker returned no expected value: {latest}")
- if tool_name == "run_code":
- return text_chunks(CODE_WORKER_TEXT)
- if tool_name == "workflow":
- return text_chunks(WORKFLOW_WORKER_TEXT)
- raise AssertionError(f"unexpected tool follow-up: {tool_name}")
- user_prompts = [
- message_text(message.get("content"))
- for message in reversed(messages)
- if isinstance(message, dict) and message.get("role") == "user"
- ]
- minimal_prompt = next(
- (
- prompt
- for prompt in user_prompts
- if prompt.startswith(f"{MINIMAL_PROMPT}\n{MINIMAL_EDITOR_PATH_PREFIX}")
- ),
- None,
- )
- # The minimal composition's assembled system prompt, advertised tool schemas, and
- # model-visible messages are pinned by its snapshot, not asserted here.
- if minimal_prompt is not None:
- return tool_call_chunks(
- "minimal-bash-1",
- MINIMAL_SHELL_TOOL,
- {"command": MINIMAL_SHELL_COMMAND},
- )
- scenario_prompts = {
- SNAPSHOT_DIRECT_CHILD_PROMPT,
- SNAPSHOT_WORKFLOW_CHILD_PROMPT,
- SNAPSHOT_PROMPT,
- CODE_PROMPT,
- WORKFLOW_PROMPT,
- FS_SEARCH_PROMPT,
- SPAWN_NODE_PROMPT,
- MCP_PROMPT,
- RESTART_FIRST_PROMPT,
- RESTART_SECOND_PROMPT,
- PROFILE_PLUGIN_PROMPT,
- }
- prompt = next(
- (candidate for candidate in user_prompts if candidate in scenario_prompts),
- message_text(latest.get("content")),
- )
- if prompt == SNAPSHOT_DIRECT_CHILD_PROMPT:
- return text_chunks("DIRECT_CHILD_OK")
- if prompt == SNAPSHOT_WORKFLOW_CHILD_PROMPT:
- return text_chunks("WORKFLOW_CHILD_OK")
- if prompt == SNAPSHOT_PROMPT:
- assert_advertised_tool(body, "cordis_define")
- return tool_call_chunks(
- "advanced-define",
- "cordis_define",
- {
- "plugin": {"kind": "new", "idPrefix": "snap"},
- "name": "Snapshot Double",
- "purpose": "Expose a deterministic doubling tool for executable snapshot verification.",
- "code": {"host": SNAPSHOT_PLUGIN_CODE},
- },
- )
- if prompt == RESTART_FIRST_PROMPT:
- return text_chunks(RESTART_FIRST_TEXT)
- if prompt == RESTART_SECOND_PROMPT:
- if any(
- isinstance(message, dict)
- and RESTART_FIRST_TEXT in message_text(message.get("content"))
- for message in messages
- ):
- raise AssertionError("second isolated process inherited the first process history")
- return text_chunks(RESTART_SECOND_TEXT)
- if prompt == CODE_PROMPT:
- assert_advertised_tool(body, "run_code")
- return tool_call_chunks(
- "call-code-worker",
- "run_code",
- {"code": "return 6 * 7", "description": "Compute the smoke value"},
- )
- if prompt == WORKFLOW_PROMPT:
- assert_advertised_tool(body, "workflow")
- return tool_call_chunks(
- "call-workflow-worker",
- "workflow",
- {
- "script": "return 6 * 7",
- "meta": {
- "name": "pkg-worker-smoke",
- "description": "exercise the packaged workflow worker",
- },
- },
- )
- if prompt == FS_SEARCH_PROMPT:
- assert_advertised_tool(body, "grep")
- assert_advertised_tool(body, "glob")
- return tool_call_chunks(
- "fs-search-grep",
- "grep",
- {"pattern": FS_SEARCH_MARKER, "path": "."},
- )
- if prompt == SPAWN_NODE_PROMPT:
- assert_advertised_tool(body, MINIMAL_SHELL_TOOL)
- return tool_call_chunks(
- SPAWN_NODE_CALL_ID,
- MINIMAL_SHELL_TOOL,
- {"command": SPAWN_NODE_COMMAND, "description": "Report the reachable Node version"},
- )
- if prompt == MCP_PROMPT:
- assert_advertised_tool(body, "mcp__fixture__add")
- return tool_call_chunks(
- "mcp-add",
- "mcp__fixture__add",
- {"a": 19, "b": 23},
- )
- if prompt == PROFILE_PLUGIN_PROMPT:
- system_text = "\n".join(
- message_text(message.get("content"))
- for message in messages
- if isinstance(message, dict) and message.get("role") == "system"
- )
- if PROFILE_PLUGIN_MARKER not in system_text:
- raise AssertionError("external profile plugin contributed no model-visible marker")
- return text_chunks(PROFILE_PLUGIN_TEXT)
- return text_chunks(EXPECTED_TEXT)
- def mcp_tool_followup(
- call_id: str,
- tool_name: str,
- tool_text: str,
- ) -> list[dict[str, object]] | None:
- """Verify one tool call through the packaged MCP client."""
- if call_id != "mcp-add":
- return None
- if tool_name != "mcp__fixture__add" or "42" not in tool_text:
- raise AssertionError(f"packaged MCP call returned an unexpected result: {tool_name}: {tool_text}")
- return text_chunks(MCP_TEXT)
- def fs_search_tool_followup(
- call_id: str,
- tool_name: str,
- tool_text: str,
- ) -> list[dict[str, object]] | None:
- """Exercise both ripgrep-backed tools through the packaged executable."""
- if not call_id.startswith("fs-search-"):
- return None
- if call_id == "fs-search-grep" and tool_name == "grep":
- if "needle.txt" not in tool_text or FS_SEARCH_MARKER not in tool_text:
- raise AssertionError(f"packaged grep returned no marker: {tool_text}")
- return tool_call_chunks(
- "fs-search-glob",
- "glob",
- {"pattern": "**/*.txt"},
- )
- if call_id == "fs-search-glob" and tool_name == "glob":
- if "needle.txt" not in tool_text:
- raise AssertionError(f"packaged glob returned no fixture path: {tool_text}")
- return text_chunks(FS_SEARCH_TEXT)
- raise AssertionError(f"unexpected filesystem-search follow-up: {call_id} {tool_name}: {tool_text}")
- def host_node_version() -> str:
- """The machine's own `node --version` line, the required shell resolution target."""
- node = shutil.which("node")
- if node is None:
- raise AssertionError("the spawn-node scenario requires Node on PATH for comparison")
- return subprocess.run(
- [node, "--version"], capture_output=True, text=True, check=True,
- ).stdout.strip()
- def spawn_node_tool_followup(
- call_id: str,
- tool_name: str,
- tool_text: str,
- ) -> list[dict[str, object]] | None:
- """Verify the packaged shell reached the machine's Node with a clean environment."""
- if call_id != SPAWN_NODE_CALL_ID:
- return None
- if tool_name != MINIMAL_SHELL_TOOL:
- raise AssertionError(f"spawn-node follow-up used an unexpected tool: {tool_name}")
- expected = host_node_version()
- if expected not in tool_text:
- raise AssertionError(
- f"packaged shell did not reach the machine's node {expected}: {tool_text}"
- )
- if "PKG_EXECPATH=ABSENT" not in tool_text:
- raise AssertionError(f"PKG_EXECPATH reached the shell child environment: {tool_text}")
- return text_chunks(SPAWN_NODE_TEXT)
- def minimal_tool_followup(
- body: dict[str, object],
- call_id: str,
- tool_name: str,
- tool_text: str,
- ) -> list[dict[str, object]] | None:
- """Verify the checked-in minimal composition's PTY and editor."""
- if not call_id.startswith("minimal-"):
- return None
- if call_id == "minimal-bash-1" and tool_name == MINIMAL_SHELL_TOOL:
- if "COUNT=1" not in tool_text:
- raise AssertionError(f"first persistent shell call lost its output: {tool_text}")
- return tool_call_chunks(
- "minimal-bash-2",
- MINIMAL_SHELL_TOOL,
- {"command": MINIMAL_SHELL_COMMAND},
- )
- if call_id == "minimal-bash-2" and tool_name == MINIMAL_SHELL_TOOL:
- expected = f"COUNT=2 CWD={MINIMAL_SHELL_SECOND_CWD}"
- if expected.lower() not in tool_text.lower():
- raise AssertionError(f"persistent shell did not retain state: {tool_text}")
- messages = body.get("messages")
- if not isinstance(messages, list):
- raise AssertionError("persistent editor smoke request has no messages")
- editor_path = next(
- (
- text.split(MINIMAL_EDITOR_PATH_PREFIX, 1)[1].strip()
- for message in messages
- if isinstance(message, dict) and message.get("role") == "user"
- for text in [message_text(message.get("content"))]
- if MINIMAL_EDITOR_PATH_PREFIX in text
- ),
- None,
- )
- if editor_path is None:
- raise AssertionError("persistent editor smoke prompt has no editor path")
- return tool_call_chunks(
- "minimal-editor",
- "str_replace_editor",
- {
- "command": "create",
- "path": editor_path,
- "file_text": "created by packaged editor\n",
- },
- )
- if call_id == "minimal-editor" and tool_name == "str_replace_editor":
- if "New file created successfully" not in tool_text:
- raise AssertionError(f"packaged editor did not create its file: {tool_text}")
- return text_chunks(MINIMAL_TEXT)
- raise AssertionError(f"unexpected minimal-agent follow-up: {call_id} {tool_name}: {tool_text}")
- def advanced_tool_followup(
- body: dict[str, object],
- call_id: str,
- tool_name: str,
- tool_text: str,
- ) -> list[dict[str, object]] | None:
- """Advance the executable snapshot's deterministic parent tool chain."""
- if not call_id.startswith("advanced-"):
- return None
- if call_id == "advanced-define" and tool_name == "cordis_define":
- if "Defined snap-1/pkg-1 (Snapshot Double)" not in tool_text:
- raise AssertionError(f"cordis_define returned no dynamic Package ids: {tool_text}")
- if "snapshot_double" in advertised_tool_names(body):
- raise AssertionError("snapshot_double was advertised before cordis_run")
- assert_advertised_tool(body, "cordis_run")
- return tool_call_chunks(
- "advanced-run",
- "cordis_run",
- {"pluginId": "snap-1", "packageId": "pkg-1", "mode": "run"},
- )
- if call_id == "advanced-run" and tool_name == "cordis_run":
- if "snap-1/pkg-1 is running (run-1)" not in tool_text:
- raise AssertionError(f"cordis_run returned no running Package ids: {tool_text}")
- assert_advertised_tool(body, "run_code")
- assert_advertised_tool(body, "snapshot_double")
- return tool_call_chunks(
- "advanced-code",
- "run_code",
- {
- "code": "return await tools.snapshot_double({ value: 21 })",
- "description": "Run the temporary Plugin tool",
- },
- )
- if call_id == "advanced-code" and tool_name == "run_code":
- if "42" not in tool_text:
- raise AssertionError(f"run_code returned no dynamic-tool value: {tool_text}")
- assert_advertised_tool(body, "subagent")
- return tool_call_chunks(
- "advanced-direct-child",
- "subagent",
- {
- "description": "Check direct child",
- "prompt": SNAPSHOT_DIRECT_CHILD_PROMPT,
- },
- )
- if call_id == "advanced-direct-child" and tool_name == "subagent":
- if "DIRECT_CHILD_OK" not in tool_text:
- raise AssertionError(f"subagent returned no expected child value: {tool_text}")
- assert_advertised_tool(body, "workflow")
- return tool_call_chunks(
- "advanced-workflow",
- "workflow",
- {
- "script": SNAPSHOT_WORKFLOW_SCRIPT,
- "meta": {
- "name": "advanced-exe-snapshot",
- "description": "exercise one packaged workflow child",
- },
- },
- )
- if call_id == "advanced-workflow" and tool_name == "workflow":
- if "WORKFLOW_CHILD_OK" not in tool_text:
- raise AssertionError(f"workflow returned no expected child value: {tool_text}")
- assert_advertised_tool(body, "cordis_undefine")
- return tool_call_chunks(
- "advanced-undefine",
- "cordis_undefine",
- {"pluginId": "snap-1"},
- )
- if call_id == "advanced-undefine" and tool_name == "cordis_undefine":
- if "Removed dynamic Plugin snap-1 and all of its Packages." not in tool_text:
- raise AssertionError(f"cordis_undefine returned no removal result: {tool_text}")
- if "snapshot_double" in advertised_tool_names(body):
- raise AssertionError("snapshot_double remained advertised after cordis_undefine")
- return text_chunks(SNAPSHOT_FINAL_TEXT)
- raise AssertionError(f"unexpected advanced tool follow-up: {call_id} {tool_name}: {tool_text}")
- def text_chunks(text: str) -> list[dict[str, object]]:
- """Build a complete streaming text response."""
- return [
- {"choices": [{"delta": {"role": "assistant", "content": None, "reasoning_content": ""}}]},
- {"choices": [{"delta": {"content": text}}]},
- {
- "choices": [{"delta": {"content": ""}, "finish_reason": "stop"}],
- "usage": {"prompt_tokens": 3, "completion_tokens": 3},
- },
- ]
- def tool_call_chunks(call_id: str, name: str, arguments: dict[str, object]) -> list[dict[str, object]]:
- """Build a complete streaming function-call response."""
- return [
- {"choices": [{"delta": {"role": "assistant", "content": None, "reasoning_content": ""}}]},
- {
- "choices": [{
- "delta": {
- "tool_calls": [{
- "index": 0,
- "id": call_id,
- "type": "function",
- "function": {"name": name, "arguments": json.dumps(arguments)},
- }],
- },
- }],
- },
- {
- "choices": [{"delta": {"content": ""}, "finish_reason": "tool_calls"}],
- "usage": {"prompt_tokens": 3, "completion_tokens": 3},
- },
- ]
- def latest_tool_call(messages: list[object]) -> tuple[str, str]:
- """Find the assistant call id and name paired with the latest tool result."""
- for message in reversed(messages[:-1]):
- if not isinstance(message, dict):
- continue
- calls = message.get("tool_calls")
- if not isinstance(calls, list):
- continue
- for call in reversed(calls):
- if not isinstance(call, dict):
- continue
- function = call.get("function")
- call_id = call.get("id")
- if (
- isinstance(call_id, str)
- and isinstance(function, dict)
- and isinstance(function.get("name"), str)
- ):
- return call_id, function["name"]
- raise AssertionError(f"tool result has no preceding assistant tool call: {messages}")
- def message_text(content: object) -> str:
- """Read OpenAI text content in either string or block-list form."""
- if isinstance(content, str):
- return content
- if isinstance(content, list):
- return "".join(
- block.get("text", "")
- for block in content
- if isinstance(block, dict) and isinstance(block.get("text"), str)
- )
- return ""
- def advertised_tool_names(body: dict[str, object]) -> set[str]:
- """Return the model-facing tool names advertised on one request."""
- tools = body.get("tools")
- if not isinstance(tools, list):
- raise AssertionError(f"model request advertised no tools: {body}")
- names: set[str] = set()
- for tool in tools:
- if not isinstance(tool, dict):
- continue
- function = tool.get("function")
- if isinstance(function, dict) and isinstance(function.get("name"), str):
- names.add(function["name"])
- return names
- def assert_advertised_tool(body: dict[str, object], expected: str) -> None:
- """Require the packaged deployment to expose the requested tool."""
- names = advertised_tool_names(body)
- if expected not in names:
- raise AssertionError(f"model request did not advertise {expected}: {names}")
- class MockModel:
- def __enter__(self) -> "MockModel":
- MockModelHandler.requests.clear()
- self.server = ThreadingHTTPServer(("127.0.0.1", 0), MockModelHandler)
- self.thread = threading.Thread(target=self.server.serve_forever, daemon=True)
- self.thread.start()
- host, port = self.server.server_address
- self.url = f"http://{host}:{port}"
- return self
- def __exit__(self, _exc_type: object, _exc: object, _tb: object) -> None:
- self.server.shutdown()
- self.server.server_close()
- self.thread.join(timeout=5)
- def main() -> None:
- parser = argparse.ArgumentParser(description=__doc__)
- parser.add_argument(
- "--scenario",
- choices=("all", "sdk-default", "sdk-custom", "sdk-minimal", "sdk-fs-search", "sdk-spawn-node", "sdk-mcp", "sdk-snapshot", "sdk-restart", "sdk-profile-plugin", "sdk-live", "direct"),
- default="all",
- )
- parser.add_argument("--exe", type=Path)
- parser.add_argument(
- "--installed-wheel",
- action="store_true",
- help="require a clean virtual environment containing matching installed SDK and runtime wheels",
- )
- parser.add_argument("--update-snapshots", action="store_true")
- args = parser.parse_args()
- if args.installed_wheel and args.exe is not None:
- parser.error("--installed-wheel resolves the wheel's own runtime and cannot be combined with --exe")
- if args.scenario == "sdk-live" and not args.installed_wheel:
- parser.error("--scenario sdk-live requires --installed-wheel")
- if args.scenario == "sdk-profile-plugin" and not args.installed_wheel:
- parser.error("--scenario sdk-profile-plugin requires --installed-wheel")
- if args.installed_wheel:
- args.exe = assert_installed_wheel_environment()
- if args.scenario in {"all", "sdk-custom", "sdk-minimal", "sdk-fs-search", "sdk-spawn-node", "sdk-snapshot", "sdk-restart", "direct"} and args.exe is None:
- parser.error("--exe is required for custom, minimal, snapshot, and direct scenarios")
- if args.update_snapshots and args.scenario not in {"all", "sdk-minimal", "sdk-snapshot", "sdk-restart"}:
- parser.error("--update-snapshots requires --scenario sdk-minimal, sdk-snapshot, sdk-restart, or all")
- if args.exe is not None and not args.exe.is_file():
- parser.error(f"runtime executable does not exist: {args.exe}")
- if args.scenario == "sdk-live":
- smoke_sdk_live()
- print("smoke-python-runtime: sdk-live passed")
- return
- with MockModel() as model:
- if args.scenario in {"all", "sdk-default"}:
- smoke_sdk_default(model.url)
- if args.scenario in {"all", "sdk-custom"}:
- assert args.exe is not None
- smoke_sdk_custom(model.url, args.exe.resolve())
- if args.scenario in {"all", "sdk-minimal"}:
- assert args.exe is not None
- smoke_sdk_minimal(model.url, args.exe.resolve(), args.update_snapshots)
- if args.scenario in {"all", "sdk-fs-search"}:
- assert args.exe is not None
- smoke_sdk_fs_search(model.url, args.exe.resolve())
- if args.scenario in {"all", "sdk-spawn-node"}:
- assert args.exe is not None
- smoke_sdk_spawn_node(model.url, args.exe.resolve())
- if args.scenario in {"all", "sdk-mcp"}:
- smoke_sdk_mcp(model.url, None if args.exe is None else args.exe.resolve())
- if args.scenario in {"all", "sdk-snapshot"}:
- assert args.exe is not None
- smoke_sdk_snapshot(model.url, args.exe.resolve(), args.update_snapshots)
- if args.scenario in {"all", "sdk-restart"}:
- assert args.exe is not None
- smoke_sdk_restart_snapshot(model.url, args.exe.resolve(), args.update_snapshots)
- if args.installed_wheel and args.scenario in {"all", "sdk-profile-plugin"}:
- smoke_sdk_profile_plugin(model.url)
- if args.scenario in {"all", "direct"}:
- assert args.exe is not None
- smoke_direct(model.url, args.exe.resolve())
- if not MockModelHandler.requests:
- raise AssertionError("mock model endpoint received no requests")
- print(f"smoke-python-runtime: {args.scenario} passed")
- def assert_installed_wheel_environment() -> Path:
- """Prove that this process imports matching non-editable wheel installations."""
- if sys.prefix == sys.base_prefix:
- raise AssertionError("installed-wheel smoke must run inside a virtual environment")
- if os.environ.get("PYTHONPATH"):
- raise AssertionError("installed-wheel smoke requires PYTHONPATH to be unset")
- if os.environ.get("DSH_RUNTIME_MODE"):
- raise AssertionError("installed-wheel smoke requires DSH_RUNTIME_MODE to be unset")
- repo_root = Path(__file__).resolve().parent.parent
- cwd = Path.cwd().resolve()
- if cwd.is_relative_to(repo_root):
- raise AssertionError(f"installed-wheel smoke must run outside the repository, got {cwd}")
- sdk_version = importlib.metadata.version("deepseek-harness-sdk")
- runtime_version = importlib.metadata.version("deepseek-harness-runtime-bin")
- if sdk_version != runtime_version:
- raise AssertionError(
- f"installed SDK/runtime versions differ: {sdk_version} != {runtime_version}"
- )
- expected_runtime_requirement = f"deepseek-harness-runtime-bin=={sdk_version}"
- requirements = importlib.metadata.requires("deepseek-harness-sdk") or []
- if expected_runtime_requirement not in requirements:
- raise AssertionError(
- f"installed SDK does not require {expected_runtime_requirement}: {requirements}"
- )
- prefix = Path(sys.prefix).resolve()
- imported: dict[str, Path] = {}
- for name in ("deepseek_harness", "deepseek_harness_runtime"):
- module = importlib.import_module(name)
- module_file = getattr(module, "__file__", None)
- if not isinstance(module_file, str):
- raise AssertionError(f"installed module {name} has no filesystem location")
- path = Path(module_file).resolve()
- if not path.is_relative_to(prefix):
- raise AssertionError(f"installed module {name} came from outside the virtual environment: {path}")
- if path.is_relative_to(repo_root):
- raise AssertionError(f"installed module {name} came from the repository checkout: {path}")
- imported[name] = path
- runtime_module = sys.modules["deepseek_harness_runtime"]
- executable = runtime_module.bundled_runtime_path().resolve()
- runtime_package = imported["deepseek_harness_runtime"].parent
- if not executable.is_relative_to(runtime_package):
- raise AssertionError(f"bundled runtime came from outside the installed runtime wheel: {executable}")
- runtime_files = importlib.metadata.files("deepseek-harness-runtime-bin") or []
- if not any(Path(file).name == executable.name for file in runtime_files):
- raise AssertionError(f"runtime executable is absent from installed distribution records: {executable}")
- return executable
- def smoke_sdk_live() -> None:
- """Run a real-model, tool-using two-turn task through installed wheels."""
- from deepseek_harness import DeepSeekHarness
- api_key = os.environ.get("DEEPSEEK_API_KEY")
- base_url = os.environ.get("DEEPSEEK_BASE_URL")
- if not api_key:
- raise AssertionError("sdk-live requires DEEPSEEK_API_KEY")
- if not base_url:
- raise AssertionError("sdk-live requires an explicit DEEPSEEK_BASE_URL")
- with tempfile.TemporaryDirectory(prefix="dsh-sdk-live-") as temporary:
- root = Path(temporary).resolve()
- dsh_home = root / "home"
- sessions = dsh_home / "sessions"
- marker = root / "live-api-marker.txt"
- session_id = "installed-wheel-live-api"
- shell_tool = "pwsh" if IS_WINDOWS else "bash"
- create_prompt = (
- f"Use the {shell_tool} tool to create the file at the absolute path below with exactly one line "
- f"containing {LIVE_API_SENTINEL}. Then reply with exactly {LIVE_API_SENTINEL}.\n{marker}"
- )
- verify_prompt = (
- "Use a tool to read the file created in the previous turn. "
- f"If its only line is {LIVE_API_SENTINEL}, reply with exactly {LIVE_API_SENTINEL}."
- )
- with DeepSeekHarness(
- provider="deepseek-official",
- model="deepseek-v4-flash",
- cwd=str(root),
- dsh_home=str(dsh_home),
- env={
- "DSH_PERMISSION_MODE": "danger-full-access",
- "DSH_TELEMETRY_DISABLED": "1",
- },
- api_key=api_key,
- base_url=base_url,
- request_timeout_seconds=180,
- ) as harness:
- created = harness.run(create_prompt, session_id=session_id)
- verified = harness.run(verify_prompt, session_id=session_id)
- for label, result in (("create", created), ("verify", verified)):
- if result.finish_reason != "completed":
- event_types = [event.get("type") for event in result.events]
- turn_end_data = next(
- (event.get("data") for event in reversed(result.events) if event.get("type") == "turn/end"),
- None,
- )
- turn_end = safe_turn_end(turn_end_data)
- raise AssertionError(
- f"{label} turn ended with {result.finish_reason!r}; "
- f"final={result.final_response!r}; turn_end={turn_end!r}; events={event_types}"
- )
- if not any(event.get("type") == "tool/call" for event in result.events):
- raise AssertionError(
- f"{label} turn made no model-requested tool call; "
- f"final={result.final_response!r}"
- )
- if result.final_response.strip() != LIVE_API_SENTINEL:
- raise AssertionError(f"{label} turn returned {result.final_response!r}")
- if not marker.is_file():
- raise AssertionError(f"real-model tool turn did not create {marker}")
- if marker.read_text(encoding="utf-8").splitlines() != [LIVE_API_SENTINEL]:
- raise AssertionError(f"real-model tool turn wrote unexpected text to {marker}")
- assert_zstd_session_log(sessions)
- def safe_turn_end(value: object) -> object:
- """Project a live-provider failure without retaining credential-bearing text."""
- if not isinstance(value, dict):
- return value
- reason = value.get("reason")
- if not isinstance(reason, dict):
- return {"turn": value.get("turn"), "reason": reason}
- error = reason.get("error")
- safe_error = None
- if isinstance(error, dict):
- safe_error = {
- key: error.get(key)
- for key in ("code", "status")
- if error.get(key) is not None
- }
- return {
- "turn": value.get("turn"),
- "reason": {
- "kind": reason.get("kind"),
- **({"error": safe_error} if safe_error is not None else {}),
- },
- }
- def smoke_sdk_default(base_url: str) -> None:
- from deepseek_harness import DeepSeekHarness
- with tempfile.TemporaryDirectory(prefix="dsh-sdk-default-") as temporary:
- root = Path(temporary).resolve()
- dsh_home = root / "home"
- sessions = dsh_home / "sessions"
- with DeepSeekHarness(
- provider="deepseek-official",
- model="smoke-model",
- cwd=str(root),
- dsh_home=str(dsh_home),
- env={
- "DSH_PERMISSION_MODE": "danger-full-access",
- "DSH_TELEMETRY_DISABLED": "1",
- },
- api_key="sk-keyless-smoke",
- base_url=base_url,
- request_timeout_seconds=60,
- ) as harness:
- result = harness.run("reply with the smoke text", session_id="default-smoke")
- assert result.final_response == EXPECTED_TEXT, (
- f"final={result.final_response!r} finish={result.finish_reason!r} "
- f"events={[event.get('type') for event in result.events]!r} "
- f"turn_end={safe_turn_end(next((event.get('data', event) for event in reversed(result.events) if event.get('type') == 'turn/end'), {}))!r}"
- )
- assert_zstd_session_log(sessions)
- def smoke_sdk_custom(base_url: str, executable: Path) -> None:
- from deepseek_harness import DeepSeekHarness
- with tempfile.TemporaryDirectory(prefix="dsh-sdk-custom-") as temporary:
- root = Path(temporary).resolve()
- dsh_home = root / "home"
- sessions = dsh_home / "sessions"
- patch = write_advanced_profile_patch(root, "custom.patch.yml", sessions)
- with DeepSeekHarness(
- provider="deepseek-official",
- model="smoke-model",
- cwd=str(root),
- dsh_bin=str(executable),
- dsh_home=str(dsh_home),
- patches=(str(patch),),
- env={
- "DSH_PERMISSION_MODE": "danger-full-access",
- "DSH_TELEMETRY_DISABLED": "1",
- },
- api_key="sk-keyless-smoke",
- base_url=base_url,
- request_timeout_seconds=60,
- ) as harness:
- text_result = harness.run("reply with the smoke text", session_id="custom-smoke")
- code_result = harness.run(CODE_PROMPT, session_id="custom-smoke")
- workflow_result = harness.run(WORKFLOW_PROMPT, session_id="custom-smoke")
- assert text_result.final_response == EXPECTED_TEXT, text_result.final_response
- assert code_result.final_response == CODE_WORKER_TEXT, code_result.final_response
- assert workflow_result.final_response == WORKFLOW_WORKER_TEXT, workflow_result.final_response
- assert_session_log(sessions, root, EXPECTED_TEXT, CODE_WORKER_TEXT, WORKFLOW_WORKER_TEXT)
- def smoke_sdk_minimal(base_url: str, executable: Path, update_snapshots: bool) -> None:
- """Exercise the shipped standalone minimal profile through the packaged executable."""
- from deepseek_harness import DeepSeekHarness
- # One mock model serves every scenario of a run, so the snapshot takes this turn's slice.
- first_request = len(MockModelHandler.requests)
- with tempfile.TemporaryDirectory(prefix="dsh-sdk-minimal-") as temporary:
- root = Path(temporary).resolve()
- editor_path = root / "created.txt"
- prompt = f"{MINIMAL_PROMPT}\n{MINIMAL_EDITOR_PATH_PREFIX}{editor_path}"
- dsh_home = root / "home"
- sessions = dsh_home / "sessions"
- with DeepSeekHarness(
- provider="deepseek-official",
- model="smoke-model",
- cwd=str(root),
- dsh_bin=str(executable),
- dsh_home=str(dsh_home),
- profile="sdk-minimal",
- api_key="sk-keyless-smoke",
- base_url=base_url,
- request_timeout_seconds=60,
- ) as harness:
- result = harness.run(prompt, session_id="minimal-agent-smoke")
- event_text = json.dumps(result.events)
- if MINIMAL_TEXT not in event_text:
- raise AssertionError(f"minimal agent run emitted no final response: {result.events}")
- if editor_path.read_text() != "created by packaged editor\n":
- raise AssertionError(f"packaged editor wrote unexpected content: {editor_path.read_text()!r}")
- assert_session_log(sessions, root, MINIMAL_TEXT, "COUNT=1", "COUNT=2")
- files = build_minimal_snapshot_files(MockModelHandler.requests[first_request:], root)
- compare_snapshot_files(
- files, update_snapshots, MINIMAL_SNAPSHOT_DIRECTORY, MINIMAL_SNAPSHOT_FILENAMES,
- )
- def smoke_sdk_fs_search(base_url: str, executable: Path) -> None:
- """Exercise real grep and glob spawns through the packaged executable."""
- from deepseek_harness import DeepSeekHarness
- with tempfile.TemporaryDirectory(prefix="dsh-sdk-fs-search-") as temporary:
- root = Path(temporary).resolve()
- (root / "needle.txt").write_text(f"{FS_SEARCH_MARKER}\n")
- dsh_home = root / "home"
- sessions = dsh_home / "sessions"
- patch = write_profile_patch(root, "fs-search.patch.yml", sessions, [
- {"id": "skill-filesystem", "disabled": True},
- {"id": "tool-fs-search", "config": {"sampleOverCapGlobResults": False}},
- ])
- with DeepSeekHarness(
- provider="deepseek-official",
- model="smoke-model",
- cwd=str(root),
- dsh_bin=str(executable),
- dsh_home=str(dsh_home),
- patches=(str(patch),),
- env={
- "DSH_PERMISSION_MODE": "danger-full-access",
- "DSH_TELEMETRY_DISABLED": "1",
- },
- api_key="sk-keyless-smoke",
- base_url=base_url,
- request_timeout_seconds=60,
- ) as harness:
- result = harness.run(FS_SEARCH_PROMPT, session_id="fs-search-smoke")
- assert result.final_response == FS_SEARCH_TEXT, result.final_response
- assert_session_log(sessions, root, FS_SEARCH_TEXT, FS_SEARCH_MARKER, "needle.txt")
- def smoke_sdk_spawn_node(base_url: str, executable: Path) -> None:
- """A shell command starting with `node` must reach the machine's Node, not the executable."""
- from deepseek_harness import DeepSeekHarness
- with tempfile.TemporaryDirectory(prefix="dsh-sdk-spawn-node-") as temporary:
- root = Path(temporary).resolve()
- dsh_home = root / "home"
- sessions = dsh_home / "sessions"
- patch = write_profile_patch(root, "spawn-node.patch.yml", sessions, [])
- with DeepSeekHarness(
- provider="deepseek-official",
- model="smoke-model",
- cwd=str(root),
- dsh_bin=str(executable),
- dsh_home=str(dsh_home),
- patches=(str(patch),),
- env={
- "DSH_PERMISSION_MODE": "danger-full-access",
- "DSH_TELEMETRY_DISABLED": "1",
- },
- api_key="sk-keyless-smoke",
- base_url=base_url,
- request_timeout_seconds=60,
- ) as harness:
- result = harness.run(SPAWN_NODE_PROMPT, session_id="spawn-node-smoke")
- assert result.final_response == SPAWN_NODE_TEXT, result.final_response
- assert_session_log(sessions, root, SPAWN_NODE_TEXT, "PKG_EXECPATH=ABSENT")
- def smoke_sdk_mcp(base_url: str, executable: Path | None) -> None:
- """Discover and call an external stdio MCP tool through the packaged client."""
- from deepseek_harness import DeepSeekHarness
- with tempfile.TemporaryDirectory(prefix="dsh-sdk-mcp-") as temporary:
- root = Path(temporary).resolve()
- dsh_home = root / "home"
- sessions = dsh_home / "sessions"
- server_script = root / "mcp_server.py"
- server_script.write_text(MCP_SERVER_SCRIPT)
- patch = write_mcp_patch(root, sessions, server_script)
- discovery_log = server_script.with_suffix(".log")
- with DeepSeekHarness(
- provider="deepseek-official",
- model="smoke-model",
- cwd=str(root),
- dsh_bin=None if executable is None else str(executable),
- dsh_home=str(dsh_home),
- patches=(str(patch),),
- env={
- "DSH_PERMISSION_MODE": "danger-full-access",
- "DSH_TELEMETRY_DISABLED": "1",
- },
- api_key="sk-keyless-smoke",
- base_url=base_url,
- request_timeout_seconds=60,
- ) as harness:
- result = harness.run(MCP_PROMPT, session_id="mcp-smoke")
- assert result.final_response == MCP_TEXT, result.final_response
- assert discovery_log.read_text().splitlines() == [
- "initialize",
- "notifications/initialized",
- "tools/list",
- "tools/call",
- ]
- assert_session_log(sessions, root, MCP_TEXT, "mcp__fixture__add", "42")
- def smoke_sdk_profile_plugin(base_url: str) -> None:
- """Install an external bundle through Python's dsh command and load it in the SDK."""
- from deepseek_harness import DeepSeekHarness
- with tempfile.TemporaryDirectory(prefix="dsh-sdk-profile-plugin-") as temporary:
- root = Path(temporary).resolve()
- dsh_home = root / "home"
- plugin = root / "plugin"
- plugin.mkdir()
- (plugin / "package.json").write_text(json.dumps({
- "name": "dsh-python-blackbox-plugin",
- "version": "1.0.0",
- "private": True,
- "type": "module",
- "exports": "./index.js",
- "peerDependencies": {"@deepseek-ai/cordis": "*"},
- "dsh": {"bundle": {"patch": "./cordis.patch.yml"}},
- }, indent=2))
- (plugin / "index.js").write_text(
- "import { Context } from '@deepseek-ai/cordis'\n"
- "export const name = 'python-sdk-blackbox-plugin'\n"
- "export const inject = ['systemPrompt']\n"
- "export function apply(ctx) {\n"
- " if (!(ctx instanceof Context)) throw new Error('external plugin loaded a second Cordis instance')\n"
- " ctx.effect(() => ctx.systemPrompt.section({\n"
- " name: 'python-sdk:blackbox-plugin',\n"
- " order: 10,\n"
- f" text: '{PROFILE_PLUGIN_MARKER}',\n"
- " }))\n"
- "}\n"
- )
- (plugin / "cordis.patch.yml").write_text(json.dumps([{
- "insert": [{"id": "python-sdk-blackbox-plugin", "name": "dsh-python-blackbox-plugin"}],
- }], indent=2))
- dsh = Path(sysconfig.get_path("scripts")) / ("dsh.exe" if IS_WINDOWS else "dsh")
- environment = {**os.environ, "DSH_HOME": str(dsh_home)}
- installed = subprocess.run(
- [str(dsh), "plugin", "--profile", "sdk", "add", f"file:{plugin}"],
- cwd=root,
- env=environment,
- text=True,
- capture_output=True,
- check=False,
- )
- if installed.returncode != 0:
- raise AssertionError(
- f"Python-installed dsh could not add the external profile plugin: "
- f"stdout={installed.stdout!r} stderr={installed.stderr!r}"
- )
- manifest = json.loads((dsh_home / "profiles" / "sdk" / "package.json").read_text())
- if "dsh-python-blackbox-plugin" not in manifest.get("dependencies", {}):
- raise AssertionError(f"dsh plugin did not record the external dependency: {manifest}")
- if "dsh-python-blackbox-plugin" not in manifest["dsh"]["profile"]["bundles"]:
- raise AssertionError(f"dsh plugin did not activate the external bundle: {manifest}")
- harness = DeepSeekHarness(
- provider="deepseek-official",
- model="smoke-model",
- cwd=str(root),
- dsh_home=str(dsh_home),
- env={
- "DSH_PERMISSION_MODE": "danger-full-access",
- "DSH_TELEMETRY_DISABLED": "1",
- },
- api_key="sk-keyless-smoke",
- base_url=base_url,
- request_timeout_seconds=60,
- )
- try:
- with harness:
- result = harness.run(PROFILE_PLUGIN_PROMPT, session_id="profile-plugin-smoke")
- except Exception as error:
- raise AssertionError(
- f"external profile plugin runtime failed: {harness.client._runtime_diagnostics()}"
- ) from error
- assert result.final_response == PROFILE_PLUGIN_TEXT, result.final_response
- assert_zstd_session_log(dsh_home / "sessions")
- def smoke_sdk_snapshot(base_url: str, executable: Path, update_snapshots: bool) -> None:
- """Drive and compare the advanced SDK/executable behavioral snapshot."""
- from deepseek_harness import DeepSeekHarness
- with tempfile.TemporaryDirectory(prefix="dsh-sdk-snapshot-") as temporary:
- root = Path(temporary).resolve()
- dsh_home = root / "home"
- sessions = dsh_home / "sessions"
- patch = write_advanced_profile_patch(root, "snapshot.patch.yml", sessions)
- with DeepSeekHarness(
- provider="deepseek-official",
- model="smoke-model",
- cwd=str(root),
- dsh_bin=str(executable),
- dsh_home=str(dsh_home),
- patches=(str(patch),),
- env={
- "DSH_PERMISSION_MODE": "danger-full-access",
- "DSH_TELEMETRY_DISABLED": "1",
- },
- api_key="sk-keyless-smoke",
- base_url=base_url,
- request_timeout_seconds=60,
- ) as harness:
- result = harness.run(SNAPSHOT_PROMPT, session_id=SNAPSHOT_SESSION_ID)
- assert result.final_response == SNAPSHOT_FINAL_TEXT, result.final_response
- methods = [notification.method for notification in result.notifications]
- if methods.count("subagent.started") != 2 or methods.count("subagent.finished") != 2:
- raise AssertionError(f"advanced snapshot emitted unexpected subagent lifecycle: {methods}")
- if not any(event.get("type") == "tool/code-dispatch" for event in result.events):
- raise AssertionError("advanced snapshot emitted no tool/code-dispatch event")
- logs = read_session_logs(sessions)
- child_ids = snapshot_child_ids(result)
- expected_ids = {SNAPSHOT_SESSION_ID, *child_ids}
- if set(logs) != expected_ids:
- raise AssertionError(f"advanced snapshot expected parent plus two child logs: {sorted(logs)}")
- if "DIRECT_CHILD_OK" not in render_jsonl(logs[child_ids[0]]):
- raise AssertionError("first advanced child log has no direct-subagent result")
- if "WORKFLOW_CHILD_OK" not in render_jsonl(logs[child_ids[1]]):
- raise AssertionError("second advanced child log has no workflow-subagent result")
- files = build_snapshot_files(result, logs, child_ids, root)
- compare_snapshot_files(
- files, update_snapshots, ADVANCED_SNAPSHOT_DIRECTORY, ADVANCED_SNAPSHOT_FILENAMES,
- )
- def smoke_sdk_restart_snapshot(base_url: str, executable: Path, update_snapshots: bool) -> None:
- """Snapshot two isolated sessions across complete SDK runtime restarts."""
- from deepseek_harness import DeepSeekHarness
- with tempfile.TemporaryDirectory(prefix="dsh-sdk-restart-") as temporary:
- root = Path(temporary).resolve()
- dsh_home = root / "home"
- sessions = dsh_home / "sessions"
- patch = write_advanced_profile_patch(root, "restart.patch.yml", sessions)
- first_request = len(MockModelHandler.requests)
- def run(prompt: str, session_id: str) -> "RunResult":
- with DeepSeekHarness(
- provider="deepseek-official",
- model="smoke-model",
- cwd=str(root),
- dsh_bin=str(executable),
- dsh_home=str(dsh_home),
- patches=(str(patch),),
- env={
- "DSH_PERMISSION_MODE": "danger-full-access",
- "DSH_TELEMETRY_DISABLED": "1",
- },
- api_key="sk-keyless-smoke",
- base_url=base_url,
- request_timeout_seconds=60,
- ) as harness:
- return harness.run(prompt, session_id=session_id)
- first = run(RESTART_FIRST_PROMPT, RESTART_FIRST_SESSION_ID)
- second = run(RESTART_SECOND_PROMPT, RESTART_SECOND_SESSION_ID)
- requests = MockModelHandler.requests[first_request:]
- if len(requests) != 2:
- raise AssertionError(f"restart snapshot expected two model requests: {requests}")
- if first.final_response != RESTART_FIRST_TEXT or second.final_response != RESTART_SECOND_TEXT:
- raise AssertionError(
- f"restart snapshot responses differ: {first.final_response!r}, {second.final_response!r}"
- )
- logs = read_session_logs(sessions)
- expected_ids = {RESTART_FIRST_SESSION_ID, RESTART_SECOND_SESSION_ID}
- if set(logs) != expected_ids:
- raise AssertionError(f"restart snapshot expected two durable sessions: {sorted(logs)}")
- for session_id, expected in (
- (RESTART_FIRST_SESSION_ID, RESTART_FIRST_TEXT),
- (RESTART_SECOND_SESSION_ID, RESTART_SECOND_TEXT),
- ):
- records = logs[session_id]
- if sum(record.get("type") == "turn/end" for record in records) != 1:
- raise AssertionError(f"restart snapshot {session_id} has an unexpected turn count")
- if expected not in render_jsonl(records):
- raise AssertionError(f"restart snapshot durable log has no {expected}")
- files = build_restart_snapshot_files(
- first,
- second,
- requests,
- logs,
- root,
- sessions,
- )
- compare_snapshot_files(
- files, update_snapshots, RESTART_SNAPSHOT_DIRECTORY, RESTART_SNAPSHOT_FILENAMES,
- )
- def smoke_direct(base_url: str, executable: Path) -> None:
- with tempfile.TemporaryDirectory(prefix="dsh-direct-") as temporary:
- root = Path(temporary).resolve()
- dsh_home = root / "home"
- sessions = dsh_home / "sessions"
- patch = write_profile_patch(root, "direct.patch.yml", sessions, [])
- environment = {
- **os.environ,
- "DSH_HOME": str(dsh_home),
- "DSH_PERMISSION_MODE": "danger-full-access",
- "DSH_TELEMETRY_DISABLED": "1",
- "DEEPSEEK_API_KEY": "sk-keyless-smoke",
- "DEEPSEEK_BASE_URL": base_url,
- }
- peer = RuntimePeer(
- [str(executable), "--profile", "sdk", "--patch", str(patch)],
- root,
- environment,
- )
- try:
- peer.send({"jsonrpc": "2.0", "id": "initialize", "method": "initialize", "params": {"cwd": str(root), "provider": "deepseek-official", "model": "smoke-model"}})
- peer.read_until(lambda message: message.get("id") == "initialize")
- peer.send({
- "jsonrpc": "2.0",
- "id": "prompt",
- "method": "session/prompt",
- "params": {"sessionId": "direct-smoke", "contentBlocks": [{"type": "text", "text": "reply with the smoke text"}]},
- })
- messages = peer.read_until(lambda message: message.get("id") == "prompt")
- if not any(is_idle_notification(message) for message in messages):
- messages.extend(peer.read_until(is_idle_notification))
- event_text = json.dumps(messages)
- if EXPECTED_TEXT not in event_text:
- raise AssertionError(f"direct runtime emitted no final response: {messages}")
- peer.send({"jsonrpc": "2.0", "id": "shutdown", "method": "shutdown"})
- peer.read_until(lambda message: message.get("id") == "shutdown")
- finally:
- peer.close()
- assert_session_log(sessions, root, EXPECTED_TEXT)
- def is_idle_notification(message: dict[str, object]) -> bool:
- """Return whether a JSON-RPC notification marks a session idle."""
- params = message.get("params")
- return (
- message.get("method") == "session.status"
- and isinstance(params, dict)
- and params.get("status") == "idle"
- )
- class RuntimePeer:
- def __init__(self, argv: list[str], cwd: Path, environment: dict[str, str]) -> None:
- self.process = subprocess.Popen(
- argv,
- cwd=cwd,
- env=environment,
- stdin=subprocess.PIPE,
- stdout=subprocess.PIPE,
- stderr=subprocess.PIPE,
- text=True,
- encoding="utf-8",
- bufsize=1,
- )
- self.stdout: queue.Queue[str | None] = queue.Queue()
- self.stderr: list[str] = []
- threading.Thread(target=self._read_stdout, daemon=True).start()
- threading.Thread(target=self._read_stderr, daemon=True).start()
- def send(self, message: dict[str, object]) -> None:
- if self.process.stdin is None:
- raise RuntimeError("runtime stdin is unavailable")
- self.process.stdin.write(json.dumps(message) + "\n")
- self.process.stdin.flush()
- def read_until(self, predicate: Callable[[dict[str, object]], bool]) -> list[dict[str, object]]:
- deadline = time.monotonic() + 60
- messages: list[dict[str, object]] = []
- while time.monotonic() < deadline:
- try:
- line = self.stdout.get(timeout=min(0.25, deadline - time.monotonic()))
- except queue.Empty:
- continue
- if line is None:
- raise RuntimeError(f"runtime exited before expected message; stderr: {''.join(self.stderr)}")
- try:
- message = json.loads(line)
- except json.JSONDecodeError:
- continue
- messages.append(message)
- if predicate(message):
- return messages
- raise TimeoutError(f"runtime timed out; messages={messages}; stderr={''.join(self.stderr)}")
- def close(self) -> None:
- if self.process.stdin is not None and not self.process.stdin.closed:
- self.process.stdin.close()
- try:
- self.process.wait(timeout=10)
- except subprocess.TimeoutExpired:
- self.process.kill()
- self.process.wait()
- if self.process.returncode not in {0, -15}:
- raise RuntimeError(f"runtime exited {self.process.returncode}; stderr: {''.join(self.stderr)}")
- def _read_stdout(self) -> None:
- assert self.process.stdout is not None
- for line in self.process.stdout:
- self.stdout.put(line)
- self.stdout.put(None)
- def _read_stderr(self) -> None:
- assert self.process.stderr is not None
- self.stderr.extend(self.process.stderr)
- PERSISTED_SESSION_FILENAME = re.compile(r"^session(?:\.v([1-9]\d*))?\.jsonl(\.zstd)?$")
- SNAPSHOT_SESSION_FILENAME = re.compile(
- r"^session(?:\.([1-9]\d*))?(?:\.v([1-9]\d*))?\.jsonl$",
- )
- def persisted_session_filename_version(path: Path, compressed: bool = False) -> int | None:
- """Return one canonical persistence basename's generation for the selected encoding."""
- match = PERSISTED_SESSION_FILENAME.fullmatch(path.name)
- if match is None or (match.group(2) is not None) != compressed:
- return None
- return int(match.group(1) or 0)
- def latest_persisted_session_paths(sessions: Path, compressed: bool = False) -> list[Path]:
- """Select the numeric-highest immutable generation in each physical Session directory."""
- pattern = "*.jsonl.zstd" if compressed else "*.jsonl"
- selected: dict[Path, tuple[int, Path]] = {}
- for path in sessions.rglob(pattern):
- version = persisted_session_filename_version(path, compressed)
- if version is None:
- continue
- previous = selected.get(path.parent)
- if previous is None or version > previous[0]:
- selected[path.parent] = (version, path)
- return sorted((entry[1] for entry in selected.values()), key=lambda path: str(path))
- def session_header_version(content: str, label: str) -> int:
- """Read a non-negative physical Session generation from the first JSONL record."""
- first = next((line for line in content.splitlines() if line), None)
- if first is None:
- raise AssertionError(f"{label}: Session log is empty")
- header = json.loads(first)
- version = header.get("version") if isinstance(header, dict) and header.get("type") == "session" else None
- if not isinstance(version, int) or isinstance(version, bool) or version < 0:
- raise AssertionError(f"{label}: Session header has no non-negative integer version")
- return version
- def assert_persisted_session_version(path: Path, content: str) -> int:
- """Require a raw persistence basename and header to name the same generation."""
- filename_version = persisted_session_filename_version(path)
- if filename_version is None:
- raise AssertionError(f"non-canonical Session persistence filename: {path.name}")
- header_version = session_header_version(content, path.name)
- if filename_version != header_version:
- raise AssertionError(
- f"{path.name}: filename declares Session format v{filename_version}, "
- f"header declares v{header_version}",
- )
- return header_version
- def snapshot_session_filename(index: int, version: int) -> str:
- """Render parent/ordinal snapshot role plus an omitted-v0 generation."""
- if index < 0 or version < 0:
- raise ValueError("snapshot Session index and version must be non-negative")
- ordinal = "" if index == 0 else f".{index}"
- generation = "" if version == 0 else f".v{version}"
- return f"session{ordinal}{generation}.jsonl"
- def parse_snapshot_session_filename(name: str) -> tuple[int, int] | None:
- """Parse one canonical parent/ordinal snapshot filename."""
- match = SNAPSHOT_SESSION_FILENAME.fullmatch(name)
- if match is None:
- if name.startswith("session") and name.endswith(".jsonl"):
- raise AssertionError(f"invalid snapshot Session filename: {name}")
- return None
- return int(match.group(1) or 0), int(match.group(2) or 0)
- def selected_snapshot_session_files(directory: Path) -> dict[int, Path]:
- """Select one highest-generation expected file per parent/ordinal role."""
- selected: dict[int, tuple[int, Path]] = {}
- for path in directory.iterdir():
- if not path.is_file():
- continue
- parsed = parse_snapshot_session_filename(path.name)
- if parsed is None:
- continue
- index, version = parsed
- content = path.read_text(encoding="utf-8")
- header_version = session_header_version(content, path.name)
- if header_version != version:
- raise AssertionError(
- f"{path.name}: filename declares Session format v{version}, header declares v{header_version}",
- )
- previous = selected.get(index)
- if previous is None or version > previous[0]:
- selected[index] = (version, path)
- return {index: value[1] for index, value in selected.items()}
- def assert_session_log(sessions: Path, cwd: Path, *expected_texts: str) -> None:
- logs = latest_persisted_session_paths(sessions)
- if len(logs) != 1:
- raise AssertionError(f"expected one JSONL session log under {sessions}, found {logs}")
- content = logs[0].read_text()
- assert_persisted_session_version(logs[0], content)
- lines = content.splitlines()
- header = json.loads(lines[0])
- if header.get("cwd") != str(cwd):
- raise AssertionError(f"session header cwd is not absolute/canonical: {header}")
- rendered = "\n".join(lines)
- for expected in expected_texts:
- if expected not in rendered:
- raise AssertionError(f"session log has no {expected!r} response: {logs[0]}")
- def assert_zstd_session_log(sessions: Path) -> None:
- logs = latest_persisted_session_paths(sessions, compressed=True)
- if len(logs) != 1:
- raise AssertionError(f"expected one Zstandard JSONL session log under {sessions}, found {logs}")
- if not logs[0].read_bytes().startswith(bytes.fromhex("28b52ffd")):
- raise AssertionError(f"session log has no Zstandard magic: {logs[0]}")
- def read_session_logs(sessions: Path) -> dict[str, list[dict[str, object]]]:
- """Parse every persisted JSONL session into a map keyed by header id."""
- logs: dict[str, list[dict[str, object]]] = {}
- for path in latest_persisted_session_paths(sessions):
- content = path.read_text(encoding="utf-8")
- assert_persisted_session_version(path, content)
- records = [
- json.loads(line)
- for line in content.splitlines()
- if line
- ]
- if not records or records[0].get("type") != "session":
- raise AssertionError(f"session log has no header: {path}")
- session_id = records[0].get("id")
- if not isinstance(session_id, str):
- raise AssertionError(f"session log header has no string id: {path}")
- if session_id in logs:
- raise AssertionError(f"duplicate persisted session id: {session_id}")
- logs[session_id] = records
- return logs
- def snapshot_child_ids(result: "RunResult") -> list[str]:
- """Return the two child session ids in their SDK notification order."""
- child_ids: list[str] = []
- for notification in result.notifications:
- if notification.method != "subagent.started":
- continue
- payload = notification.payload
- if payload.get("parentSessionId") != SNAPSHOT_SESSION_ID:
- continue
- child_id = payload.get("childSessionId")
- if isinstance(child_id, str) and child_id not in child_ids:
- child_ids.append(child_id)
- if len(child_ids) != 2:
- raise AssertionError(f"advanced snapshot expected two child session ids: {child_ids}")
- return child_ids
- def build_minimal_snapshot_files(
- requests: list[dict[str, object]],
- cwd: Path,
- ) -> dict[str, str]:
- """Render the minimal composition's model-visible surface as expected output.
- Every assembled system prompt, advertised tool schema, and system or user message is
- kept verbatim: they carry what the deployment actually shows the model, so a plugin
- that contributes an unintended system section or user message cannot pass unnoticed.
- Assistant and tool payloads keep only their call identity because their text differs
- across the platforms this expected output must replay on. The shipped profile omits
- dynamic runtime context, so every message it emits is compared.
- """
- snapshot = []
- for body in requests:
- messages = body.get("messages")
- if not isinstance(messages, list):
- raise AssertionError(f"minimal model request has no messages: {body}")
- snapshot.append({
- "tools": minimal_snapshot_text(body.get("tools"), cwd),
- "messages": [
- minimal_snapshot_message(message, cwd)
- for message in messages
- ],
- })
- return {"model-visible.json": json.dumps(snapshot, indent=2, ensure_ascii=False) + "\n"}
- def minimal_snapshot_message(message: object, cwd: Path) -> dict[str, object]:
- """Reduce one model-visible message to its stable, behavior-carrying parts."""
- if not isinstance(message, dict):
- raise AssertionError(f"minimal model request has an invalid message: {message}")
- role = message.get("role")
- if role in ("system", "user"):
- return {"role": role, "text": minimal_snapshot_text(message_text(message.get("content")), cwd)}
- if role == "assistant":
- calls = message.get("tool_calls")
- if not isinstance(calls, list):
- raise AssertionError(f"minimal assistant message has no tool calls: {message}")
- return {
- "role": role,
- "toolCalls": [
- {"id": call.get("id"), "name": (call.get("function") or {}).get("name")}
- for call in calls
- if isinstance(call, dict)
- ],
- }
- if role == "tool":
- return {"role": role, "toolCallId": message.get("tool_call_id"), "text": "{{tool-result}}"}
- raise AssertionError(f"minimal model request has an unexpected message role: {message}")
- def minimal_snapshot_text(value: object, cwd: Path) -> object:
- """Replace the scenario's temporary working directory everywhere it appears."""
- if isinstance(value, str):
- return value.replace(str(cwd), "{{cwd}}")
- if isinstance(value, list):
- return [minimal_snapshot_text(item, cwd) for item in value]
- if isinstance(value, dict):
- return {key: minimal_snapshot_text(item, cwd) for key, item in value.items()}
- return value
- def build_snapshot_files(
- result: "RunResult",
- logs: dict[str, list[dict[str, object]]],
- child_ids: list[str],
- cwd: Path,
- ) -> dict[str, str]:
- """Render the SDK result and three persisted logs into stable expected outputs."""
- replacements = [(str(cwd), "{{cwd}}"), (SNAPSHOT_SESSION_ID, "{{parent}}")]
- replacements.append((snapshot_workflow_run_id(result), "{{workflow-run}}"))
- for index, child_id in enumerate(child_ids, start=1):
- replacements.append((child_id, f"{{{{child-{index}}}}}"))
- agent_id = snapshot_agent_id(result, child_id)
- replacements.append((agent_id, f"{{{{agent-{index}}}}}"))
- replacements.sort(key=lambda pair: len(pair[0]), reverse=True)
- result_value = {
- "session_id": result.session_id,
- "final_response": result.final_response,
- "events": result.events,
- "notifications": [
- {"method": notification.method, "payload": notification.payload}
- for notification in result.notifications
- ],
- }
- normalized_result = normalize_snapshot_value(result_value, replacements)
- parent_records = project_session_snapshot([
- normalize_snapshot_value(record, replacements) for record in logs[SNAPSHOT_SESSION_ID]
- ])
- files = {
- "result.json": json.dumps(normalized_result, indent=2, ensure_ascii=False) + "\n",
- snapshot_session_filename(
- 0, session_header_version(render_jsonl(parent_records), "advanced parent"),
- ): render_jsonl(parent_records),
- }
- for index, child_id in enumerate(child_ids, start=1):
- child_records = project_session_snapshot([
- normalize_snapshot_value(record, replacements) for record in logs[child_id]
- ])
- child_content = render_jsonl(child_records)
- files[snapshot_session_filename(
- index, session_header_version(child_content, f"advanced child {index}"),
- )] = child_content
- return files
- def build_restart_snapshot_files(
- first: "RunResult",
- second: "RunResult",
- requests: list[dict[str, object]],
- logs: dict[str, list[dict[str, object]]],
- cwd: Path,
- sessions: Path,
- ) -> dict[str, str]:
- """Render two SDK processes, isolated model histories, and durable logs."""
- replacements = [
- (str(sessions), "{{sessions}}"),
- (str(cwd), "{{cwd}}"),
- (RESTART_FIRST_SESSION_ID, "{{session-1}}"),
- (RESTART_SECOND_SESSION_ID, "{{session-2}}"),
- ]
- result_value = [
- {
- "session_id": result.session_id,
- "final_response": result.final_response,
- "finish_reason": result.finish_reason,
- "eventTypes": [
- event.get("type")
- for event in result.events
- ],
- "notificationMethods": [
- notification.method
- for notification in result.notifications
- ],
- }
- for result in (first, second)
- ]
- request_value = [
- {
- "model": request.get("model"),
- "messages": restart_request_messages(request),
- "toolNames": sorted(advertised_tool_names(request)),
- }
- for request in requests
- ]
- first_records = project_session_snapshot([
- normalize_snapshot_value(record, replacements) for record in logs[RESTART_FIRST_SESSION_ID]
- ])
- second_records = project_session_snapshot([
- normalize_snapshot_value(record, replacements) for record in logs[RESTART_SECOND_SESSION_ID]
- ])
- first_content = render_jsonl(first_records)
- second_content = render_jsonl(second_records)
- return {
- "result.json": json.dumps(
- normalize_snapshot_value(result_value, replacements), indent=2, ensure_ascii=False,
- ) + "\n",
- "requests.json": json.dumps(
- normalize_snapshot_value(request_value, replacements), indent=2, ensure_ascii=False,
- ) + "\n",
- snapshot_session_filename(
- 1, session_header_version(first_content, "restart Session 1"),
- ): first_content,
- snapshot_session_filename(
- 2, session_header_version(second_content, "restart Session 2"),
- ): second_content,
- }
- def restart_request_messages(request: dict[str, object]) -> list[object]:
- """Project model history while tokenizing composition-owned system prose."""
- messages = request.get("messages")
- if not isinstance(messages, list):
- raise AssertionError(f"restart snapshot request has no messages: {request}")
- return [
- {"role": "system", "content": "{{system}}"}
- if isinstance(message, dict) and message.get("role") == "system"
- else message
- for message in messages
- ]
- def snapshot_workflow_run_id(result: "RunResult") -> str:
- """Return the one workflow run id emitted by the advanced scenario."""
- run_ids: set[str] = set()
- for event in result.events:
- event_type = event.get("type")
- data = event.get("data")
- if not isinstance(event_type, str) or not event_type.startswith("tool-workflow/"):
- continue
- if isinstance(data, dict) and isinstance(data.get("runId"), str):
- run_ids.add(data["runId"])
- if len(run_ids) != 1:
- raise AssertionError(f"advanced snapshot expected one workflow run id: {sorted(run_ids)}")
- return next(iter(run_ids))
- def snapshot_agent_id(result: "RunResult", child_id: str) -> str:
- """Find the successful subagent id paired with one child session."""
- for notification in result.notifications:
- if notification.method != "subagent.finished":
- continue
- payload = notification.payload
- if payload.get("childSessionId") != child_id:
- continue
- if payload.get("provider") != "spawn" or payload.get("status") != "ok":
- raise AssertionError(f"advanced child did not finish successfully: {payload}")
- agent_id = payload.get("agentId")
- if isinstance(agent_id, str):
- return agent_id
- raise AssertionError(f"advanced snapshot has no finished agent for child {child_id}")
- def normalize_snapshot_value(
- value: object,
- replacements: list[tuple[str, str]],
- ) -> object:
- """Scrub volatile values and bulky request headers without losing behavior."""
- if isinstance(value, str):
- normalized = value
- for actual, token in replacements:
- normalized = normalized.replace(actual, token)
- return normalized
- if isinstance(value, list):
- return [normalize_snapshot_value(item, replacements) for item in value]
- if not isinstance(value, dict):
- return value
- normalized = {
- key: normalize_snapshot_value(item, replacements)
- for key, item in value.items()
- }
- if normalized.get("type") == "session" and "createdAt" in normalized:
- normalized["createdAt"] = 0
- if "seq" in normalized and "time" in normalized:
- normalized["time"] = 0
- if normalized.get("type") in ("assistant/message", "assistant/attempt"):
- data = normalized.get("data")
- stream = data.get("stream") if isinstance(data, dict) else None
- if isinstance(stream, list):
- for member in stream:
- if not isinstance(member, dict):
- continue
- if isinstance(member.get("time"), (int, float)):
- member["time"] = 0
- if isinstance(member.get("time0"), (int, float)):
- member["time0"] = 0
- dt = member.get("dt")
- if isinstance(dt, list):
- member["dt"] = [0] * len(dt)
- if isinstance(normalized.get("id"), str) and normalized.get("role") in ("assistant", "user"):
- normalized["id"] = "{{messageId}}"
- scrub_snapshot_header(normalized)
- return normalized
- def scrub_snapshot_header(value: dict[object, object]) -> None:
- """Tokenize full request-header bulk while retaining tool names."""
- data = value.get("data")
- if not isinstance(data, dict):
- return
- if value.get("type") == "request/header":
- header = data.get("header")
- if not isinstance(header, dict):
- return
- if "system" in header:
- header["system"] = "{{system}}"
- tools = header.get("tools")
- if isinstance(tools, list):
- header["tools"] = [
- tool.get("name") if isinstance(tool, dict) else "{{tools}}"
- for tool in tools
- ]
- def render_jsonl(records: list[object]) -> str:
- """Render parsed JSON values as compact, newline-terminated JSONL."""
- return "".join(
- json.dumps(record, ensure_ascii=False, separators=(",", ":")) + "\n"
- for record in records
- )
- def project_session_snapshot(records: list[dict[str, object]]) -> list[dict[str, object]]:
- """Omit storage sequence/time envelopes from snapshot body records."""
- projected = [dict(record) for record in records]
- for record in projected[1:]:
- for key in ("seq", "time", "seq0", "time0"):
- record.pop(key, None)
- return projected
- SESSION_FORMAT_PROVENANCE = "{{sessionFormatVersion}}"
- def expand_snapshot_stream_member(member: object) -> list[dict[str, object]]:
- """Expand one compact Assistant stream member into logical provider chunks."""
- if not isinstance(member, dict):
- raise AssertionError(f"snapshot Assistant stream member is not an object: {member!r}")
- member_type = member.get("type")
- if member_type == "chunk":
- chunk = member.get("chunk")
- if not isinstance(chunk, dict):
- raise AssertionError(f"snapshot Assistant chunk member has no chunk: {member!r}")
- return [chunk]
- packed_kinds = {
- "text-chunks": ("texts", "text-delta", "text"),
- "reasoning-chunks": ("texts", "reasoning-delta", "text"),
- "tool-call-chunks": ("args", "tool-call-delta", "argumentsDelta"),
- }
- packed = packed_kinds.get(member_type)
- if packed is None:
- raise AssertionError(f"snapshot Assistant stream has unknown member type: {member_type!r}")
- values_key, chunk_type, value_key = packed
- values = member.get(values_key)
- if not isinstance(values, list):
- raise AssertionError(f"snapshot Assistant stream member has no {values_key}: {member!r}")
- shared = {
- key: member[key]
- for key in ("index", "id", "name")
- if key in member
- }
- return [
- {"type": chunk_type, **shared, value_key: value}
- for value in values
- ]
- def expand_snapshot_assistant_event(value: object) -> list[object]:
- """Expand one direct or SDK-wrapped v2 settlement for generation-neutral comparison."""
- if not isinstance(value, dict):
- return [value]
- event = value
- wrapper_key: str | None = None
- wrapper: dict[str, object] | None = None
- if value.get("method") == "session.event":
- for candidate in ("payload", "params"):
- container = value.get(candidate)
- nested = container.get("event") if isinstance(container, dict) else None
- if isinstance(nested, dict):
- event = nested
- wrapper_key = candidate
- wrapper = container
- break
- if event.get("type") not in ("assistant/message", "assistant/attempt"):
- return [value]
- data = event.get("data")
- stream = data.get("stream") if isinstance(data, dict) else None
- if not isinstance(stream, list):
- return [value]
- def wrap(expanded: dict[str, object]) -> object:
- if wrapper_key is None or wrapper is None:
- return expanded
- return {**value, wrapper_key: {**wrapper, "event": expanded}}
- common = {
- key: data[key]
- for key in ("turn", "step")
- if key in data
- }
- expanded = [
- wrap({
- "type": "assistant/chunk",
- "data": {**common, "chunk": chunk},
- })
- for member in stream
- for chunk in expand_snapshot_stream_member(member)
- ]
- if event.get("type") == "assistant/message":
- expanded.append(wrap({
- **event,
- "data": {key: item for key, item in data.items() if key != "stream"},
- }))
- return expanded
- def normalize_session_format_comparison(
- value: object,
- source_session_version: int | None = None,
- ) -> object:
- """Canonicalize only generation provenance that differs across immutable Session files."""
- if isinstance(value, list):
- return [
- normalize_session_format_comparison(expanded, source_session_version)
- for item in value
- for expanded in expand_snapshot_assistant_event(item)
- ]
- if not isinstance(value, dict):
- return value
- normalized = {
- key: normalize_session_format_comparison(item, source_session_version)
- for key, item in value.items()
- }
- if normalized.get("type") == "session" and "version" in normalized:
- normalized["version"] = SESSION_FORMAT_PROVENANCE
- normalized.setdefault("isSeeded", False)
- ordered_header = {
- key: normalized[key]
- for key in ("type", "version", "id", "createdAt", "cwd", "isSeeded", "delegationDepth")
- if key in normalized
- }
- normalized = {
- **ordered_header,
- **{key: item for key, item in normalized.items() if key not in ordered_header},
- }
- if isinstance(normalized.get("type"), str) and "data" in normalized:
- normalized.pop("seq", None)
- normalized.pop("time", None)
- if source_session_version == 1 and normalized.get("type") == "assistant/message":
- normalized.pop("sourceEventSeqs", None)
- if normalized.get("type") == "session-log-deepseek/delivery-accepted":
- data = normalized.get("data")
- if isinstance(data, dict):
- data.pop("throughSeq", None)
- data.pop("sessionFormatVersion", None)
- data["sessionFormatVersion"] = SESSION_FORMAT_PROVENANCE
- if normalized.get("kind") == "session-reference":
- references = normalized.get("references")
- if isinstance(references, list):
- for reference in references:
- if isinstance(reference, dict):
- reference.pop("capturedFormatVersion", None)
- reference["capturedFormatVersion"] = SESSION_FORMAT_PROVENANCE
- return normalized
- def normalize_snapshot_comparison_text(name: str, content: str) -> str:
- """Normalize Session generation provenance only while comparing committed expected outputs."""
- if name.startswith("session") and name.endswith(".jsonl"):
- parsed = [json.loads(line) for line in content.splitlines() if line]
- header = parsed[0] if parsed else None
- source_version = header.get("version") if isinstance(header, dict) else None
- if not isinstance(source_version, int):
- raise AssertionError(f"{name}: snapshot Session header has no integer format version")
- records = [
- normalize_session_format_comparison(expanded, source_version)
- for record in parsed
- for expanded in expand_snapshot_assistant_event(record)
- ]
- return render_jsonl(records)
- if name.endswith(".json"):
- return json.dumps(
- normalize_session_format_comparison(json.loads(content)),
- indent=2,
- ensure_ascii=False,
- ) + "\n"
- return content
- def compare_snapshot_files(
- files: dict[str, str],
- update: bool,
- directory: Path,
- filenames: tuple[str, ...],
- ) -> None:
- """Write or exactly compare one scenario's expected snapshot files."""
- scenario = directory.name
- if tuple(files) != filenames:
- raise AssertionError(f"{scenario} snapshot builder produced {tuple(files)}, expected {filenames}")
- if update:
- directory.mkdir(parents=True, exist_ok=True)
- for name, content in files.items():
- (directory / name).write_text(content, encoding="utf-8", newline="\n")
- print(f"smoke-python-runtime: updated snapshots in {directory}")
- existing = [path for path in directory.iterdir() if path.is_file()] if directory.is_dir() else []
- expected_non_session = {
- name for name in filenames if parse_snapshot_session_filename(name) is None
- }
- existing_non_session = {
- path.name for path in existing if parse_snapshot_session_filename(path.name) is None
- }
- if existing_non_session != expected_non_session:
- raise AssertionError(
- f"{scenario} snapshot files differ: "
- f"missing={sorted(expected_non_session - existing_non_session)}, "
- f"unexpected={sorted(existing_non_session - expected_non_session)}"
- )
- selected_expected = selected_snapshot_session_files(directory)
- actual_sessions: dict[int, tuple[str, str]] = {}
- for name, content in files.items():
- parsed = parse_snapshot_session_filename(name)
- if parsed is None:
- continue
- index, filename_version = parsed
- header_version = session_header_version(content, name)
- if filename_version != header_version:
- raise AssertionError(
- f"{name}: filename declares Session format v{filename_version}, "
- f"header declares v{header_version}",
- )
- if index in actual_sessions:
- raise AssertionError(f"{scenario} snapshot builder produced duplicate Session role {index}")
- actual_sessions[index] = (name, content)
- if set(selected_expected) != set(actual_sessions):
- raise AssertionError(
- f"{scenario} snapshot Session roles differ: "
- f"expected={sorted(selected_expected)}, actual={sorted(actual_sessions)}",
- )
- for name, actual in files.items():
- parsed = parse_snapshot_session_filename(name)
- expected_path = directory / name if parsed is None else selected_expected[parsed[0]]
- expected_text = expected_path.read_text(encoding="utf-8")
- compared_actual = normalize_snapshot_comparison_text(name, actual)
- compared_expected = normalize_snapshot_comparison_text(expected_path.name, expected_text)
- if compared_actual == compared_expected:
- continue
- diff = "".join(difflib.unified_diff(
- compared_expected.splitlines(keepends=True),
- compared_actual.splitlines(keepends=True),
- fromfile=f"expected/{expected_path.name}",
- tofile=f"actual/{name}",
- ))
- raise AssertionError(
- f"{scenario} executable snapshot mismatch in {name}; "
- "rerun with --update-snapshots after reviewing the behavior\n"
- f"{diff}"
- )
- if __name__ == "__main__":
- main()
|