The native adapter's route was named deepseek, colliding with pi-ai's catalog provider of the same name, so the two DeepSeek paths could never be mounted side by side. The web settings page needs both configurable at once. Compositions, fixtures, goldens, scaffolding defaults, and docs all move together (pre-release, no shim); TUI/session-query-spill/ missing-credential goldens re-recorded through their keyless refresh modes because provider-name length shifts box padding and spill truncation points.
923 lines
36 KiB
Python
923 lines
36 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, Notification
|
|
|
|
|
|
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"
|
|
init_dump = tmp_path / "init.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":
|
|
json.dump(msg.get("params"), open(os.environ["INIT_DUMP"], "w"))
|
|
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",
|
|
max_tokens=4096,
|
|
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),
|
|
"INIT_DUMP": str(init_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")
|
|
assert json.loads(init_dump.read_text()) == {
|
|
"cwd": str(tmp_path),
|
|
"provider": "deepseek-official",
|
|
"model": "deepseek-v4-flash",
|
|
"maxTokens": 4096,
|
|
}
|
|
|
|
|
|
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_collects_nested_subagent_tree_without_polluting_root_events(
|
|
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":
|
|
root = (msg.get("params") or {})["sessionId"]
|
|
print(json.dumps({"jsonrpc": "2.0", "method": "subagent.started", "params": {"parentSessionId": root, "childSessionId": "child"}}), flush=True)
|
|
print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": "child", "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "child response"}]}}}}), flush=True)
|
|
print(json.dumps({"jsonrpc": "2.0", "method": "subagent.started", "params": {"parentSessionId": "child", "childSessionId": "grandchild"}}), flush=True)
|
|
print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": "grandchild", "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "grandchild response"}]}}}}), flush=True)
|
|
print(json.dumps({"jsonrpc": "2.0", "method": "subagent.finished", "params": {"parentSessionId": "child", "childSessionId": "grandchild", "status": "ok"}}), flush=True)
|
|
print(json.dumps({"jsonrpc": "2.0", "method": "subagent.finished", "params": {"parentSessionId": root, "childSessionId": "child", "status": "ok"}}), flush=True)
|
|
print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": root, "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "root response"}]}}}}), flush=True)
|
|
print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": root, "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:
|
|
result = harness.run(
|
|
"delegate recursively",
|
|
session_id="main",
|
|
on_notification=lambda notification: seen.append(notification.method),
|
|
)
|
|
assert harness.client._notifications.qsize() == 0
|
|
|
|
assert result.status == "ok"
|
|
assert result.final_response == "root response"
|
|
assert [event["data"]["content"][0]["text"] for event in result.events] == ["root response"]
|
|
assert [notification.method for notification in result.notifications] == [
|
|
"subagent.started",
|
|
"session.event",
|
|
"subagent.started",
|
|
"session.event",
|
|
"subagent.finished",
|
|
"subagent.finished",
|
|
"session.event",
|
|
"session.finished",
|
|
]
|
|
assert seen == [notification.method for notification in result.notifications]
|
|
|
|
|
|
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(provider="deepseek-official", 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_session_subscription_keeps_descendant_relationships_across_subscriptions() -> None:
|
|
client = HarnessClient()
|
|
with client.subscribe_session_notifications("main") as first:
|
|
client._handle_message({
|
|
"jsonrpc": "2.0",
|
|
"method": "subagent.started",
|
|
"params": {"parentSessionId": "main", "childSessionId": "child"},
|
|
})
|
|
assert first.next().payload["childSessionId"] == "child"
|
|
|
|
with client.subscribe_session_notifications("main") as second:
|
|
client._handle_message({
|
|
"jsonrpc": "2.0",
|
|
"method": "subagent.started",
|
|
"params": {"parentSessionId": "child", "childSessionId": "grandchild"},
|
|
})
|
|
client._handle_message({
|
|
"jsonrpc": "2.0",
|
|
"method": "session.event",
|
|
"params": {"sessionId": "grandchild", "event": {"type": "assistant/message"}},
|
|
})
|
|
assert second.next().payload["childSessionId"] == "grandchild"
|
|
assert second.next().payload["sessionId"] == "grandchild"
|
|
|
|
assert client._notifications.qsize() == 0
|
|
|
|
|
|
def test_session_subscription_preserves_reused_child_ancestry_after_late_finish() -> None:
|
|
client = HarnessClient()
|
|
old_seen: list[Notification] = []
|
|
new_seen: list[Notification] = []
|
|
with (
|
|
client.subscribe_session_notifications("old-parent") as old_subscription,
|
|
client.subscribe_session_notifications("new-parent") as new_subscription,
|
|
):
|
|
client._handle_message({
|
|
"jsonrpc": "2.0",
|
|
"method": "subagent.started",
|
|
"params": {"parentSessionId": "old-parent", "childSessionId": "reused-child"},
|
|
})
|
|
old_subscription.drain(old_seen.append)
|
|
new_subscription.drain(new_seen.append)
|
|
assert [notification.method for notification in old_seen] == ["subagent.started"]
|
|
assert new_seen == []
|
|
|
|
client._handle_message({
|
|
"jsonrpc": "2.0",
|
|
"method": "subagent.started",
|
|
"params": {"parentSessionId": "new-parent", "childSessionId": "reused-child"},
|
|
})
|
|
old_subscription.drain(old_seen.append)
|
|
new_subscription.drain(new_seen.append)
|
|
assert [notification.method for notification in new_seen] == ["subagent.started"]
|
|
|
|
client._handle_message({
|
|
"jsonrpc": "2.0",
|
|
"method": "subagent.finished",
|
|
"params": {"parentSessionId": "old-parent", "childSessionId": "reused-child"},
|
|
})
|
|
old_subscription.drain(old_seen.append)
|
|
new_subscription.drain(new_seen.append)
|
|
assert [notification.method for notification in old_seen] == [
|
|
"subagent.started",
|
|
"subagent.finished",
|
|
]
|
|
assert [notification.method for notification in new_seen] == ["subagent.started"]
|
|
|
|
client._handle_message({
|
|
"jsonrpc": "2.0",
|
|
"method": "session.event",
|
|
"params": {"sessionId": "reused-child", "event": {"type": "assistant/message"}},
|
|
})
|
|
old_subscription.drain(old_seen.append)
|
|
new_subscription.drain(new_seen.append)
|
|
|
|
assert [notification.method for notification in old_seen] == [
|
|
"subagent.started",
|
|
"subagent.finished",
|
|
]
|
|
assert [notification.method for notification in new_seen] == [
|
|
"subagent.started",
|
|
"session.event",
|
|
]
|
|
assert client._notifications.qsize() == 0
|
|
|
|
|
|
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(provider="deepseek-official", 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(provider="deepseek-official", 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(provider="deepseek-official", 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(provider="deepseek-official", 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(provider="deepseek-official", 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(provider="deepseek-official", 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(provider="deepseek-official", 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 "max_tokens" in DeepSeekHarnessConfig.__dataclass_fields__
|
|
assert "max_tokens" in inspect.signature(HarnessClient.initialize).parameters
|
|
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(provider="deepseek-official", 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(provider="deepseek-official", 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(provider="deepseek-official", 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:
|
|
"""Install a fake runtime package that records config and serves lifecycle calls.
|
|
|
|
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(provider="deepseek-official", 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(provider="deepseek-official", 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()
|