Files
deepseek-harness/python/sdk/tests/test_client.py
T
Tianyi Cui d5e894edf4 fix(sdk): harden runtime lifecycle and JSON-RPC
Keep DeepSeekHarness.run() reusable, but make ownership of its lazy
runtime process explicit. Document the context-manager/close contract and
update every construction example to use a context manager so repeated runs
remain valid without encouraging leaked subprocesses.

Contain notification predicate failures at the subscription boundary. Remove
only the subscriber whose callback raised, deliver that exception through its
queue, and continue dispatching to healthy subscribers so arbitrary callback
code cannot terminate the shared reader thread or strand later requests.

Enforce one in-flight prompt per server session with an atomic activePrompt
guard. Route overlap through the existing -32603 handler-error response and
clear the guard in finally, preserving parallel prompts across sessions and
sequential reuse without changing JSON-RPC request or notification shapes.

Use StringDecoder for line framing so a UTF-8 code point split across Buffer
chunks is not corrupted. Add a queued-write flush barrier, and make memoized
shutdown await it before disposal and exit while retaining exactly-once
cleanup when shutdown calls race or flushing fails.

Cover callback isolation, same-session exclusion, cross-session concurrency,
split multibyte input, delayed writes, racing shutdown, and flush failure with
deterministic tests.
2026-07-13 21:40:12 +08:00

768 lines
29 KiB
Python

from __future__ import annotations
import json
import inspect
import sys
import threading
import time
from pathlib import Path
import pytest
from deepseek_harness import DeepSeekHarness, HarnessClient, HarnessConfig
def test_high_level_sdk_runs_turn_and_collects_final_response(tmp_path: Path) -> None:
script = tmp_path / "fake_runtime.py"
env_dump = tmp_path / "env.json"
script.write_text(
"""
import json
import os
import sys
env_dump = os.environ["ENV_DUMP"]
json.dump({
"DEEPSEEK_API_KEY": os.environ.get("DEEPSEEK_API_KEY"),
"DEEPSEEK_BASE_URL": os.environ.get("DEEPSEEK_BASE_URL"),
"DSH_CWD": os.environ.get("DSH_CWD"),
"DSH_SESSION_ROOT": os.environ.get("DSH_SESSION_ROOT"),
"DSH_CORDIS_CONFIG": os.environ.get("DSH_CORDIS_CONFIG"),
}, open(env_dump, "w"))
for line in sys.stdin:
msg = json.loads(line)
method = msg.get("method")
if method == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-runtime"}}}), flush=True)
elif method == "session/prompt":
params = msg.get("params") or {}
print(json.dumps({
"jsonrpc": "2.0",
"method": "session.event",
"params": {
"sessionId": params["sessionId"],
"event": {
"type": "assistant/message",
"data": {"content": [{"type": "text", "text": "hello from runtime"}]},
},
},
}), flush=True)
print(json.dumps({
"jsonrpc": "2.0",
"method": "session.finished",
"params": {"sessionId": params["sessionId"], "status": "ok"},
}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True)
elif method == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
with DeepSeekHarness(
model="deepseek-v4-flash",
cwd=str(tmp_path),
cordis=str(tmp_path / "cordis.yml"),
session_root=str(tmp_path / "sessions"),
launch_args_override=(sys.executable, str(script)),
env={
"ENV_DUMP": str(env_dump),
"DEEPSEEK_API_KEY": "env-key",
"DEEPSEEK_BASE_URL": "http://127.0.0.1:4321",
},
) as harness:
result = harness.run("say hello", session_id="main")
assert result.status == "ok"
assert result.final_response == "hello from runtime"
assert result.events[0]["type"] == "assistant/message"
dumped_env = json.loads(env_dump.read_text())
assert dumped_env["DEEPSEEK_API_KEY"] == "env-key"
assert dumped_env["DEEPSEEK_BASE_URL"] == "http://127.0.0.1:4321"
assert dumped_env["DSH_CWD"] == str(tmp_path)
assert dumped_env["DSH_SESSION_ROOT"] == str(tmp_path / "sessions")
assert dumped_env["DSH_CORDIS_CONFIG"] == str(tmp_path / "cordis.yml")
def test_session_run_invokes_notification_callback_before_returning(tmp_path: Path) -> None:
script = tmp_path / "fake_runtime.py"
script.write_text(
"""
import json
import sys
for line in sys.stdin:
msg = json.loads(line)
method = msg.get("method")
if method == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-runtime"}}}), flush=True)
elif method == "session/prompt":
print(json.dumps({"jsonrpc": "2.0", "method": "subagent.started", "params": {"parentSessionId": "main", "childSessionId": "child"}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": "main", "status": "ok"}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True)
elif method == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
seen: list[str] = []
with DeepSeekHarness(
launch_args_override=(sys.executable, str(script)),
cwd=str(tmp_path),
) as harness:
session = harness.start_session("main")
result = session.run(
"spawn a helper",
on_notification=lambda notification: seen.append(notification.method),
)
assert result.status == "ok"
assert seen == ["subagent.started", "session.finished"]
def test_relative_cwd_is_absolute_in_process_environment_and_wire(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
script = tmp_path / "capture_cwd.py"
capture = tmp_path / "cwd.json"
script.write_text(
"""
import json
import os
import sys
for line in sys.stdin:
msg = json.loads(line)
if msg.get("method") == "initialize":
json.dump({"process": os.getcwd(), "environment": os.environ.get("DSH_CWD"), "wire": msg["params"]["cwd"]}, open(os.environ["CAPTURE"], "w"))
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-runtime"}}}), flush=True)
elif msg.get("method") == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
monkeypatch.chdir(tmp_path)
with DeepSeekHarness(
cwd=".",
runtime_cwd=".",
launch_args_override=(sys.executable, str(script)),
env={"CAPTURE": str(capture)},
):
pass
expected = str(tmp_path.resolve())
assert json.loads(capture.read_text()) == {
"process": expected,
"environment": expected,
"wire": expected,
}
def test_session_run_includes_subagent_finished_for_parent_session(tmp_path: Path) -> None:
script = tmp_path / "fake_runtime.py"
script.write_text(
"""
import json
import sys
for line in sys.stdin:
msg = json.loads(line)
method = msg.get("method")
if method == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-runtime"}}}), flush=True)
elif method == "session/prompt":
print(json.dumps({"jsonrpc": "2.0", "method": "subagent.started", "params": {"parentSessionId": "main", "childSessionId": "child"}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "method": "subagent.finished", "params": {"parentSessionId": "main", "childSessionId": "child", "status": "ok", "stopReason": "completed"}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": "main", "status": "ok"}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True)
elif method == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
with DeepSeekHarness(
launch_args_override=(sys.executable, str(script)),
cwd=str(tmp_path),
) as harness:
result = harness.run("spawn a helper", session_id="main")
assert result.status == "ok"
assert [notification.method for notification in result.notifications] == [
"subagent.started",
"subagent.finished",
"session.finished",
]
def test_session_run_ignores_notifications_for_other_sessions(tmp_path: Path) -> None:
script = tmp_path / "fake_runtime.py"
script.write_text(
"""
import json
import sys
for line in sys.stdin:
msg = json.loads(line)
method = msg.get("method")
if method == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-runtime"}}}), flush=True)
elif method == "session/prompt":
params = msg.get("params") or {}
print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": "other", "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "wrong session"}]}}}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": "other", "status": "ok"}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": params["sessionId"], "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "right session"}]}}}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": params["sessionId"], "status": "ok"}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True)
elif method == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
with DeepSeekHarness(
launch_args_override=(sys.executable, str(script)),
cwd=str(tmp_path),
) as harness:
result = harness.run("stay in your lane", session_id="main")
assert result.status == "ok"
assert result.final_response == "right session"
assert [notification.payload.get("sessionId") for notification in result.notifications] == ["main", "main"]
def test_high_level_session_run_does_not_accumulate_global_notifications(tmp_path: Path) -> None:
script = tmp_path / "fake_runtime.py"
script.write_text(
"""
import json
import sys
for line in sys.stdin:
msg = json.loads(line)
method = msg.get("method")
if method == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-runtime"}}}), flush=True)
elif method == "session/prompt":
params = msg.get("params") or {}
print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": params["sessionId"], "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "ok"}]}}}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": params["sessionId"], "status": "ok"}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True)
elif method == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
with DeepSeekHarness(launch_args_override=(sys.executable, str(script)), cwd=str(tmp_path)) as harness:
result = harness.run("one turn", session_id="main")
assert result.status == "ok"
assert harness.client._notifications.qsize() == 0
def test_session_run_waits_for_late_finished_without_replaying_stale_notifications(tmp_path: Path) -> None:
script = tmp_path / "fake_runtime.py"
script.write_text(
"""
import json
import sys
import time
turn = 0
for line in sys.stdin:
msg = json.loads(line)
method = msg.get("method")
if method == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-runtime"}}}), flush=True)
elif method == "session/prompt":
turn += 1
params = msg.get("params") or {}
session_id = params["sessionId"]
if turn == 1:
print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": session_id, "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "first"}]}}}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": session_id, "status": "ok"}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True)
else:
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True)
time.sleep(0.05)
print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": session_id, "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "second"}]}}}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": session_id, "status": "ok"}}), flush=True)
elif method == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
with DeepSeekHarness(launch_args_override=(sys.executable, str(script)), cwd=str(tmp_path)) as harness:
first = harness.run("first turn", session_id="main")
second = harness.run("second turn", session_id="main")
assert first.final_response == "first"
assert second.final_response == "second"
assert [notification.payload.get("sessionId") for notification in second.notifications] == ["main", "main"]
def test_client_starts_subprocess_sends_requests_and_routes_notifications(tmp_path: Path) -> None:
script = tmp_path / "fake_bridge.py"
script.write_text(
"""
import json
import sys
for line in sys.stdin:
msg = json.loads(line)
method = msg.get("method")
if method == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-dsh"}}}), flush=True)
elif method == "session/prompt":
params = msg.get("params") or {}
print(json.dumps({"jsonrpc": "2.0", "method": "llm/request", "params": {"requestId": "req-1", "sessionId": params["sessionId"], "model": "dsagent", "messages": []}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True)
elif method == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
with HarnessClient(
HarnessConfig(launch_args_override=(sys.executable, str(script)))
) as client:
init = client.initialize(cwd="/workspace", model="dsagent")
assert init.serverInfo.name == "fake-dsh"
client.session_prompt("main", [{"type": "text", "text": "fix it"}])
notification = client.next_notification()
assert notification.method == "llm/request"
assert notification.payload["requestId"] == "req-1"
assert notification.payload["sessionId"] == "main"
def test_client_keeps_unmatched_notifications_available_globally_while_subscribed() -> None:
client = HarnessClient()
with client.subscribe_session_notifications("main"):
client._handle_message({
"jsonrpc": "2.0",
"method": "session.event",
"params": {"sessionId": "other", "event": {"type": "assistant/message"}},
})
assert client._notifications.qsize() == 1
notification = client._notifications.get_nowait()
assert not isinstance(notification, BaseException)
assert notification.method == "session.event"
assert notification.payload["sessionId"] == "other"
def test_client_contains_notification_filter_failure_to_its_subscription(tmp_path: Path) -> None:
script = tmp_path / "fake_bridge.py"
script.write_text(
"""
import json
import sys
for line in sys.stdin:
msg = json.loads(line)
method = msg.get("method")
if method == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-dsh"}}}), flush=True)
elif method in {"emit-first", "emit-second"}:
print(json.dumps({"jsonrpc": "2.0", "method": "tick", "params": {"source": method}}), flush=True)
elif method == "session/prompt":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True)
elif method == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
def broken_filter(_notification: object) -> bool:
raise RuntimeError("bad notification filter")
with HarnessClient(HarnessConfig(launch_args_override=(sys.executable, str(script)))) as client:
client.initialize(cwd="/workspace", model="dsagent")
with (
client.subscribe_notifications(broken_filter) as broken,
client.subscribe_notifications(lambda notification: notification.method == "tick") as healthy,
):
client.notify("emit-first")
with pytest.raises(RuntimeError, match="bad notification filter"):
broken.next()
assert healthy.next().payload == {"source": "emit-first"}
assert client._notifications.qsize() == 0
client.session_prompt("main", [{"type": "text", "text": "reader still works"}])
client.notify("emit-second")
assert healthy.next().payload == {"source": "emit-second"}
def test_client_rejects_unaccepted_session_prompt_response(tmp_path: Path) -> None:
script = tmp_path / "fake_bridge.py"
script.write_text(
"""
import json
import sys
for line in sys.stdin:
msg = json.loads(line)
method = msg.get("method")
if method == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-dsh"}}}), flush=True)
elif method == "session/prompt":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": False}}), flush=True)
elif method == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
with HarnessClient(HarnessConfig(launch_args_override=(sys.executable, str(script)))) as client:
client.initialize(cwd="/workspace", model="dsagent")
with pytest.raises(ValueError):
client.session_prompt("main", [{"type": "text", "text": "fix it"}])
def test_client_routes_bridge_requests_and_sends_responses(tmp_path: Path) -> None:
script = tmp_path / "fake_bridge.py"
script.write_text(
"""
import json
import sys
for line in sys.stdin:
msg = json.loads(line)
method = msg.get("method")
if method == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-dsh"}}}), flush=True)
print(json.dumps({"jsonrpc": "2.0", "id": "bridge-req-1", "method": "llm.request", "params": {"requestId": "req-1", "sessionId": "main", "model": "dsagent", "messages": []}}), flush=True)
elif "id" in msg and "method" not in msg:
print(json.dumps({"jsonrpc": "2.0", "method": "response/seen", "params": {"result": msg.get("result")}}), flush=True)
elif method == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
with HarnessClient(
HarnessConfig(launch_args_override=(sys.executable, str(script)))
) as client:
client.initialize(cwd="/workspace", model="dsagent")
request = client.next_request()
assert request.id == "bridge-req-1"
assert request.method == "llm.request"
assert request.payload["requestId"] == "req-1"
client.respond(request.id, {"content_blocks": [{"type": "text", "text": "done"}]})
notification = client.next_notification()
assert notification.method == "response/seen"
assert notification.payload["result"]["content_blocks"][0]["text"] == "done"
def test_client_ignores_non_json_stdout_lines(tmp_path: Path) -> None:
script = tmp_path / "fake_bridge.py"
script.write_text(
"""
import json
import sys
print("node warning: experimental loader", flush=True)
for line in sys.stdin:
msg = json.loads(line)
if msg.get("method") == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-dsh"}}}), flush=True)
elif msg.get("method") == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
with HarnessClient(
HarnessConfig(launch_args_override=(sys.executable, str(script)))
) as client:
init = client.initialize(cwd="/workspace", model="dsagent")
assert init.serverInfo.name == "fake-dsh"
def test_client_request_times_out_when_bridge_does_not_respond(tmp_path: Path) -> None:
script = tmp_path / "fake_bridge.py"
script.write_text(
"""
import time
time.sleep(60)
""".strip()
)
with HarnessClient(
HarnessConfig(
launch_args_override=(sys.executable, str(script)),
request_timeout_seconds=0.1,
)
) as client:
start = time.monotonic()
try:
client.initialize(cwd="/workspace", model="dsagent")
except TimeoutError:
assert time.monotonic() - start < 2
else:
raise AssertionError("initialize should time out")
def test_client_close_times_out_when_shutdown_does_not_respond(tmp_path: Path) -> None:
script = tmp_path / "fake_bridge.py"
script.write_text(
"""
import json
import signal
import sys
import time
signal.signal(signal.SIGTERM, signal.SIG_IGN)
for line in sys.stdin:
msg = json.loads(line)
if msg.get("method") == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-dsh"}}}), flush=True)
elif msg.get("method") == "shutdown":
time.sleep(60)
""".strip()
)
client = HarnessClient(
HarnessConfig(
launch_args_override=(sys.executable, str(script)),
shutdown_timeout_seconds=0.1,
)
)
client.start()
proc = client._proc
assert proc is not None
client.initialize(cwd="/workspace", model="dsagent")
start = time.monotonic()
client.close()
assert time.monotonic() - start < 2
assert proc.poll() is not None
assert client._proc is None
def test_initialize_failure_reaps_started_runtime(tmp_path: Path) -> None:
script = tmp_path / "rejecting_runtime.py"
script.write_text(
"""
import json
import sys
for line in sys.stdin:
msg = json.loads(line)
if msg.get("method") == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "error": {"code": -32000, "message": "bad initialize"}}), flush=True)
elif msg.get("method") == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
client = HarnessClient(HarnessConfig(launch_args_override=(sys.executable, str(script))))
client.start()
proc = client._proc
assert proc is not None
with pytest.raises(Exception, match="bad initialize"):
client.initialize(cwd=".", model="dsagent")
assert proc.wait(timeout=1) is not None
assert client._proc is None
def test_public_signatures_omit_unsupported_wire_parameters() -> None:
from deepseek_harness import DeepSeekHarnessConfig, Session
assert "session_root" not in inspect.signature(HarnessClient.initialize).parameters
assert "system_prompt" not in inspect.signature(HarnessClient.initialize).parameters
assert "profile" not in inspect.signature(HarnessClient.session_prompt).parameters
assert "profile" not in inspect.signature(DeepSeekHarness.run).parameters
assert "profile" not in inspect.signature(Session.run).parameters
assert "system_prompt" not in DeepSeekHarnessConfig.__dataclass_fields__
assert "client_name" not in HarnessConfig.__dataclass_fields__
assert "client_version" not in HarnessConfig.__dataclass_fields__
def test_client_close_is_idempotent_before_and_after_start(tmp_path: Path) -> None:
HarnessClient().close()
script = tmp_path / "fake_bridge.py"
script.write_text(
"""
import json
import sys
for line in sys.stdin:
msg = json.loads(line)
if msg.get("method") == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-dsh"}}}), flush=True)
elif msg.get("method") == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
client = HarnessClient(HarnessConfig(launch_args_override=(sys.executable, str(script))))
client.start()
client.initialize(cwd="/workspace", model="dsagent")
client.close()
client.close()
def test_runtime_closed_error_includes_stderr_tail(tmp_path: Path) -> None:
script = tmp_path / "crashing_runtime.py"
script.write_text(
"""
import sys
print("fatal bridge exploded", file=sys.stderr, flush=True)
sys.exit(42)
""".strip()
)
with HarnessClient(
HarnessConfig(
launch_args_override=(sys.executable, str(script)),
request_timeout_seconds=2,
)
) as client:
with pytest.raises(Exception, match="fatal bridge exploded"):
client.initialize(cwd="/workspace", model="dsagent")
def test_client_serializes_concurrent_writes(tmp_path: Path) -> None:
script = tmp_path / "fake_bridge.py"
output = tmp_path / "seen.jsonl"
script.write_text(
"""
import json
import os
import sys
with open(os.environ["SEEN"], "w") as seen:
for line in sys.stdin:
seen.write(line)
seen.flush()
msg = json.loads(line)
if "id" in msg and msg.get("method") == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-dsh"}}}), flush=True)
elif "id" in msg and msg.get("method") == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
with HarnessClient(
HarnessConfig(
launch_args_override=(sys.executable, str(script)),
env={"SEEN": str(output)},
)
) as client:
client.initialize(cwd="/workspace", model="dsagent")
threads = [
threading.Thread(target=client.notify, args=(f"notice-{index}", {"index": index}))
for index in range(50)
]
for thread in threads:
thread.start()
for thread in threads:
thread.join()
for line in output.read_text().splitlines():
json.loads(line)
def _install_fake_bundled_runtime(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> Path:
"""Fake the deepseek-harness-runtime-bin package on sys.path.
A stub exe that dumps DSH_CORDIS_CONFIG to $ENV_DUMP before serving
initialize/shutdown, plus a module exposing the resolution surface the
client consumes. Returns the fake bundled default config path.
"""
runtime = tmp_path / "dsh-jsonrpc-agent"
runtime.write_text(
"""#!/usr/bin/env python3
import json
import os
import sys
json.dump({"DSH_CORDIS_CONFIG": os.environ.get("DSH_CORDIS_CONFIG")}, open(os.environ["ENV_DUMP"], "w"))
for line in sys.stdin:
msg = json.loads(line)
if msg.get("method") == "initialize":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "bundled-runtime"}}}), flush=True)
elif msg.get("method") == "shutdown":
print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True)
break
""".strip()
)
runtime.chmod(0o755)
default_config = tmp_path / "default-cordis.yml"
module_dir = tmp_path / "deepseek_harness_runtime"
module_dir.mkdir()
(module_dir / "__init__.py").write_text(
f"""
def resolve_bundled_launch_args(mode=None):
return ({str(runtime)!r},)
def bundled_default_config_path():
return {str(default_config)!r}
""".strip()
)
monkeypatch.syspath_prepend(str(tmp_path))
monkeypatch.delitem(sys.modules, "deepseek_harness_runtime", raising=False)
return default_config
@pytest.mark.parametrize("ambient_config", [None, ""], ids=["unset", "empty-counts-as-absent"])
def test_client_default_launch_uses_bundled_runtime_and_injects_default_config(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ambient_config: str | None
) -> None:
env_dump = tmp_path / "env.json"
default_config = _install_fake_bundled_runtime(tmp_path, monkeypatch)
if ambient_config is None:
monkeypatch.delenv("DSH_CORDIS_CONFIG", raising=False)
else:
monkeypatch.setenv("DSH_CORDIS_CONFIG", ambient_config)
with HarnessClient(HarnessConfig(env={"ENV_DUMP": str(env_dump)})) as client:
init = client.initialize(cwd="/workspace", model="deepseek-v4-pro")
assert init.serverInfo.name == "bundled-runtime"
assert json.loads(env_dump.read_text())["DSH_CORDIS_CONFIG"] == str(default_config)
def test_client_respects_explicit_config_over_bundled_default(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
env_dump = tmp_path / "env.json"
_install_fake_bundled_runtime(tmp_path, monkeypatch)
monkeypatch.delenv("DSH_CORDIS_CONFIG", raising=False)
with HarnessClient(
HarnessConfig(env={"ENV_DUMP": str(env_dump), "DSH_CORDIS_CONFIG": "./explicit.yml"})
) as client:
client.initialize(cwd="/workspace", model="deepseek-v4-pro")
assert json.loads(env_dump.read_text())["DSH_CORDIS_CONFIG"] == "./explicit.yml"
def test_client_reports_missing_bundled_runtime_dependency(monkeypatch: pytest.MonkeyPatch) -> None:
monkeypatch.delitem(sys.modules, "deepseek_harness_runtime", raising=False)
monkeypatch.setattr(sys, "path", [])
with pytest.raises(FileNotFoundError, match="Install deepseek-harness-runtime-bin"):
HarnessClient().start()