2026-07-11 14:08:24 +08:00
|
|
|
from __future__ import annotations
|
|
|
|
|
|
|
|
|
|
import json
|
|
|
|
|
import os
|
|
|
|
|
import queue
|
|
|
|
|
import subprocess
|
|
|
|
|
import threading
|
|
|
|
|
import time
|
|
|
|
|
import uuid
|
|
|
|
|
from collections import deque
|
|
|
|
|
from dataclasses import dataclass
|
2026-07-13 16:34:03 +08:00
|
|
|
from pathlib import Path
|
2026-07-30 16:48:28 +08:00
|
|
|
from typing import Callable, TypeAlias, TypeVar
|
2026-07-11 14:08:24 +08:00
|
|
|
|
|
|
|
|
from pydantic import BaseModel
|
|
|
|
|
|
|
|
|
|
from .errors import JsonRpcError, TransportClosedError
|
|
|
|
|
from .models import IncomingRequest, InitializeResponse, JsonObject, JsonValue, Notification
|
|
|
|
|
|
|
|
|
|
ModelT = TypeVar("ModelT", bound=BaseModel)
|
|
|
|
|
NotificationFilter: TypeAlias = Callable[[Notification], bool]
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@dataclass(slots=True)
|
|
|
|
|
class HarnessConfig:
|
|
|
|
|
"""Configuration for launching the local DeepSeek Harness SDK runtime."""
|
|
|
|
|
|
feat(python-sdk): launch dsh profiles from explicit homes
Replace complete-config, session_root, runtime-bin, bridge-bin, and public argv override options with dsh_bin, profile, ordered patches, and dsh_home. Resolve executable/home/patch/cwd paths before spawn, select the sdk profile by default, and fail before launch unless dsh_home or non-empty DSH_HOME is explicit; Python never inherits ~/.dsh silently.
Remove Python-owned DSH_CORDIS_CONFIG, DSH_SESSION_ROOT, and DSH_CWD injection and drop session_root from RunResult. Keep arbitrary argv only as an underscore-prefixed fake-runtime adapter, retain provider/model/token and process controls, and append subprocess stderr to initialization JSON-RPC errors so profile boot failures name their actual plugin cause. Unit and carrier tests cover both exe and Node modes.
2026-08-23 14:56:16 +08:00
|
|
|
dsh_bin: str | None = None
|
|
|
|
|
profile: str = "sdk"
|
|
|
|
|
patches: tuple[str, ...] = ()
|
|
|
|
|
dsh_home: str | None = None
|
2026-07-11 14:08:24 +08:00
|
|
|
cwd: str | None = None
|
|
|
|
|
env: dict[str, str] | None = None
|
|
|
|
|
request_timeout_seconds: float | None = None
|
|
|
|
|
shutdown_timeout_seconds: float | None = 1.0
|
feat(python-sdk): launch dsh profiles from explicit homes
Replace complete-config, session_root, runtime-bin, bridge-bin, and public argv override options with dsh_bin, profile, ordered patches, and dsh_home. Resolve executable/home/patch/cwd paths before spawn, select the sdk profile by default, and fail before launch unless dsh_home or non-empty DSH_HOME is explicit; Python never inherits ~/.dsh silently.
Remove Python-owned DSH_CORDIS_CONFIG, DSH_SESSION_ROOT, and DSH_CWD injection and drop session_root from RunResult. Keep arbitrary argv only as an underscore-prefixed fake-runtime adapter, retain provider/model/token and process controls, and append subprocess stderr to initialization JSON-RPC errors so profile boot failures name their actual plugin cause. Unit and carrier tests cover both exe and Node modes.
2026-08-23 14:56:16 +08:00
|
|
|
_launch_args: tuple[str, ...] | None = None
|
2026-07-11 14:08:24 +08:00
|
|
|
|
|
|
|
|
|
|
|
|
|
class HarnessClient:
|
|
|
|
|
"""Synchronous JSON-RPC client for the DeepSeek Harness SDK runtime over stdio."""
|
|
|
|
|
|
feat(python-sdk): launch dsh profiles from explicit homes
Replace complete-config, session_root, runtime-bin, bridge-bin, and public argv override options with dsh_bin, profile, ordered patches, and dsh_home. Resolve executable/home/patch/cwd paths before spawn, select the sdk profile by default, and fail before launch unless dsh_home or non-empty DSH_HOME is explicit; Python never inherits ~/.dsh silently.
Remove Python-owned DSH_CORDIS_CONFIG, DSH_SESSION_ROOT, and DSH_CWD injection and drop session_root from RunResult. Keep arbitrary argv only as an underscore-prefixed fake-runtime adapter, retain provider/model/token and process controls, and append subprocess stderr to initialization JSON-RPC errors so profile boot failures name their actual plugin cause. Unit and carrier tests cover both exe and Node modes.
2026-08-23 14:56:16 +08:00
|
|
|
def __init__(
|
|
|
|
|
self,
|
|
|
|
|
config: HarnessConfig | None = None,
|
|
|
|
|
*,
|
|
|
|
|
_launch_args: tuple[str, ...] | None = None,
|
|
|
|
|
) -> None:
|
2026-07-11 14:08:24 +08:00
|
|
|
self.config = config or HarnessConfig()
|
feat(python-sdk): launch dsh profiles from explicit homes
Replace complete-config, session_root, runtime-bin, bridge-bin, and public argv override options with dsh_bin, profile, ordered patches, and dsh_home. Resolve executable/home/patch/cwd paths before spawn, select the sdk profile by default, and fail before launch unless dsh_home or non-empty DSH_HOME is explicit; Python never inherits ~/.dsh silently.
Remove Python-owned DSH_CORDIS_CONFIG, DSH_SESSION_ROOT, and DSH_CWD injection and drop session_root from RunResult. Keep arbitrary argv only as an underscore-prefixed fake-runtime adapter, retain provider/model/token and process controls, and append subprocess stderr to initialization JSON-RPC errors so profile boot failures name their actual plugin cause. Unit and carrier tests cover both exe and Node modes.
2026-08-23 14:56:16 +08:00
|
|
|
self._launch_args = _launch_args or self.config._launch_args
|
2026-07-11 14:08:24 +08:00
|
|
|
self._proc: subprocess.Popen[str] | None = None
|
|
|
|
|
self._lock = threading.Lock()
|
|
|
|
|
self._write_lock = threading.Lock()
|
|
|
|
|
self._responses: dict[str, queue.Queue[JsonValue | BaseException]] = {}
|
|
|
|
|
self._notifications: queue.Queue[Notification | BaseException] = queue.Queue()
|
|
|
|
|
self._notification_subscribers: dict[
|
|
|
|
|
str, tuple[queue.Queue[Notification | BaseException], NotificationFilter | None]
|
|
|
|
|
] = {}
|
2026-07-24 11:23:16 +08:00
|
|
|
self._session_parents: dict[str, str] = {}
|
2026-07-11 14:08:24 +08:00
|
|
|
self._requests: queue.Queue[IncomingRequest | BaseException] = queue.Queue()
|
|
|
|
|
self._stderr_lines: deque[str] = deque(maxlen=400)
|
|
|
|
|
self._reader_thread: threading.Thread | None = None
|
|
|
|
|
self._stderr_thread: threading.Thread | None = None
|
|
|
|
|
|
|
|
|
|
def __enter__(self) -> "HarnessClient":
|
|
|
|
|
self.start()
|
|
|
|
|
return self
|
|
|
|
|
|
|
|
|
|
def __exit__(self, _exc_type, _exc, _tb) -> None:
|
|
|
|
|
self.close()
|
|
|
|
|
|
|
|
|
|
def start(self) -> None:
|
|
|
|
|
if self._proc is not None:
|
|
|
|
|
return
|
2026-07-24 11:23:16 +08:00
|
|
|
with self._lock:
|
|
|
|
|
self._session_parents.clear()
|
2026-07-11 14:08:24 +08:00
|
|
|
env = os.environ.copy()
|
|
|
|
|
if self.config.env:
|
|
|
|
|
env.update(self.config.env)
|
feat(python-sdk): launch dsh profiles from explicit homes
Replace complete-config, session_root, runtime-bin, bridge-bin, and public argv override options with dsh_bin, profile, ordered patches, and dsh_home. Resolve executable/home/patch/cwd paths before spawn, select the sdk profile by default, and fail before launch unless dsh_home or non-empty DSH_HOME is explicit; Python never inherits ~/.dsh silently.
Remove Python-owned DSH_CORDIS_CONFIG, DSH_SESSION_ROOT, and DSH_CWD injection and drop session_root from RunResult. Keep arbitrary argv only as an underscore-prefixed fake-runtime adapter, retain provider/model/token and process controls, and append subprocess stderr to initialization JSON-RPC errors so profile boot failures name their actual plugin cause. Unit and carrier tests cover both exe and Node modes.
2026-08-23 14:56:16 +08:00
|
|
|
args = list(self._launch_args or self._default_launch_args(env))
|
2026-07-11 14:08:24 +08:00
|
|
|
self._proc = subprocess.Popen(
|
|
|
|
|
args,
|
|
|
|
|
stdin=subprocess.PIPE,
|
|
|
|
|
stdout=subprocess.PIPE,
|
|
|
|
|
stderr=subprocess.PIPE,
|
|
|
|
|
text=True,
|
|
|
|
|
encoding="utf-8",
|
2026-07-13 16:34:03 +08:00
|
|
|
cwd=None if self.config.cwd is None else str(Path(self.config.cwd).resolve()),
|
2026-07-11 14:08:24 +08:00
|
|
|
env=env,
|
|
|
|
|
bufsize=1,
|
|
|
|
|
)
|
|
|
|
|
self._start_reader_thread()
|
|
|
|
|
self._start_stderr_thread()
|
|
|
|
|
|
|
|
|
|
def close(self) -> None:
|
|
|
|
|
proc = self._proc
|
|
|
|
|
if proc is None:
|
|
|
|
|
return
|
|
|
|
|
try:
|
|
|
|
|
self.request("shutdown", None, response_model=_ShutdownResponse, timeout_seconds=self.config.shutdown_timeout_seconds)
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
self._stderr_lines.append(f"shutdown request failed: {exc}")
|
|
|
|
|
if proc.stdin:
|
|
|
|
|
try:
|
|
|
|
|
proc.stdin.close()
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
self._stderr_lines.append(f"stdin close failed: {exc}")
|
2026-07-13 16:34:03 +08:00
|
|
|
if proc.poll() is None:
|
|
|
|
|
try:
|
2026-07-11 14:08:24 +08:00
|
|
|
proc.terminate()
|
2026-07-13 16:34:03 +08:00
|
|
|
except ProcessLookupError:
|
|
|
|
|
pass
|
|
|
|
|
try:
|
|
|
|
|
proc.wait(timeout=self.config.shutdown_timeout_seconds)
|
|
|
|
|
except subprocess.TimeoutExpired:
|
2026-07-11 14:08:24 +08:00
|
|
|
proc.kill()
|
2026-07-13 16:34:03 +08:00
|
|
|
proc.wait()
|
|
|
|
|
self._proc = None
|
2026-07-11 14:08:24 +08:00
|
|
|
self._fail_waiters(self._runtime_closed_error("DeepSeek Harness runtime closed"))
|
|
|
|
|
if self._reader_thread and self._reader_thread.is_alive():
|
|
|
|
|
self._reader_thread.join(timeout=0.5)
|
|
|
|
|
if self._stderr_thread and self._stderr_thread.is_alive():
|
|
|
|
|
self._stderr_thread.join(timeout=0.5)
|
|
|
|
|
|
|
|
|
|
def initialize(
|
|
|
|
|
self,
|
|
|
|
|
*,
|
|
|
|
|
cwd: str,
|
2026-07-14 21:57:52 +08:00
|
|
|
provider: str,
|
2026-07-11 14:08:24 +08:00
|
|
|
model: str,
|
2026-07-28 17:36:44 +08:00
|
|
|
max_tokens: int | None = None,
|
2026-07-11 14:08:24 +08:00
|
|
|
) -> InitializeResponse:
|
|
|
|
|
payload: JsonObject = {
|
2026-07-13 16:34:03 +08:00
|
|
|
"cwd": str(Path(cwd).resolve()),
|
2026-07-14 21:57:52 +08:00
|
|
|
"provider": provider,
|
2026-07-11 14:08:24 +08:00
|
|
|
"model": model,
|
|
|
|
|
}
|
2026-07-28 17:36:44 +08:00
|
|
|
if max_tokens is not None:
|
|
|
|
|
payload["maxTokens"] = max_tokens
|
2026-07-13 16:34:03 +08:00
|
|
|
try:
|
|
|
|
|
return self.request("initialize", payload, response_model=InitializeResponse)
|
feat(python-sdk): launch dsh profiles from explicit homes
Replace complete-config, session_root, runtime-bin, bridge-bin, and public argv override options with dsh_bin, profile, ordered patches, and dsh_home. Resolve executable/home/patch/cwd paths before spawn, select the sdk profile by default, and fail before launch unless dsh_home or non-empty DSH_HOME is explicit; Python never inherits ~/.dsh silently.
Remove Python-owned DSH_CORDIS_CONFIG, DSH_SESSION_ROOT, and DSH_CWD injection and drop session_root from RunResult. Keep arbitrary argv only as an underscore-prefixed fake-runtime adapter, retain provider/model/token and process controls, and append subprocess stderr to initialization JSON-RPC errors so profile boot failures name their actual plugin cause. Unit and carrier tests cover both exe and Node modes.
2026-08-23 14:56:16 +08:00
|
|
|
except BaseException as error:
|
2026-07-13 16:34:03 +08:00
|
|
|
self.close()
|
feat(python-sdk): launch dsh profiles from explicit homes
Replace complete-config, session_root, runtime-bin, bridge-bin, and public argv override options with dsh_bin, profile, ordered patches, and dsh_home. Resolve executable/home/patch/cwd paths before spawn, select the sdk profile by default, and fail before launch unless dsh_home or non-empty DSH_HOME is explicit; Python never inherits ~/.dsh silently.
Remove Python-owned DSH_CORDIS_CONFIG, DSH_SESSION_ROOT, and DSH_CWD injection and drop session_root from RunResult. Keep arbitrary argv only as an underscore-prefixed fake-runtime adapter, retain provider/model/token and process controls, and append subprocess stderr to initialization JSON-RPC errors so profile boot failures name their actual plugin cause. Unit and carrier tests cover both exe and Node modes.
2026-08-23 14:56:16 +08:00
|
|
|
diagnostics = self._runtime_diagnostics()
|
|
|
|
|
if isinstance(error, JsonRpcError) and diagnostics:
|
|
|
|
|
raise JsonRpcError(
|
|
|
|
|
error.code,
|
|
|
|
|
f"{error.message}\n{diagnostics}",
|
|
|
|
|
error.data,
|
|
|
|
|
) from error
|
2026-07-13 16:34:03 +08:00
|
|
|
raise
|
2026-07-11 14:08:24 +08:00
|
|
|
|
|
|
|
|
def session_prompt(
|
|
|
|
|
self,
|
|
|
|
|
session_id: str,
|
|
|
|
|
content_blocks: list[JsonObject],
|
|
|
|
|
*,
|
|
|
|
|
on_notification: Callable[[Notification], None] | None = None,
|
|
|
|
|
notification_subscription: "NotificationSubscription | None" = None,
|
2026-07-30 16:48:28 +08:00
|
|
|
) -> str:
|
2026-07-11 14:08:24 +08:00
|
|
|
payload: JsonObject = {"sessionId": session_id, "contentBlocks": content_blocks}
|
2026-07-30 16:48:28 +08:00
|
|
|
response = self.request(
|
2026-07-11 14:08:24 +08:00
|
|
|
"session/prompt",
|
|
|
|
|
payload,
|
|
|
|
|
response_model=_SessionPromptResponse,
|
|
|
|
|
on_notification=on_notification,
|
2026-07-24 11:23:16 +08:00
|
|
|
notification_filter=self._notification_belongs_to_session_tree(session_id),
|
2026-07-11 14:08:24 +08:00
|
|
|
notification_subscription=notification_subscription,
|
|
|
|
|
)
|
2026-07-30 16:48:28 +08:00
|
|
|
return response.messageId
|
2026-07-11 14:08:24 +08:00
|
|
|
|
|
|
|
|
def request(
|
|
|
|
|
self,
|
|
|
|
|
method: str,
|
|
|
|
|
params: JsonObject | None,
|
|
|
|
|
*,
|
|
|
|
|
response_model: type[ModelT],
|
|
|
|
|
timeout_seconds: float | None = None,
|
|
|
|
|
on_notification: Callable[[Notification], None] | None = None,
|
|
|
|
|
notification_filter: NotificationFilter | None = None,
|
|
|
|
|
notification_subscription: "NotificationSubscription | None" = None,
|
|
|
|
|
) -> ModelT:
|
|
|
|
|
result = self._request_raw(
|
|
|
|
|
method,
|
|
|
|
|
params,
|
|
|
|
|
timeout_seconds=timeout_seconds,
|
|
|
|
|
on_notification=on_notification,
|
|
|
|
|
notification_filter=notification_filter,
|
|
|
|
|
notification_subscription=notification_subscription,
|
|
|
|
|
)
|
|
|
|
|
if not isinstance(result, dict):
|
|
|
|
|
raise TypeError(f"{method} response must be a JSON object")
|
|
|
|
|
return response_model.model_validate(result)
|
|
|
|
|
|
|
|
|
|
def notify(self, method: str, params: JsonObject | None = None) -> None:
|
|
|
|
|
message: JsonObject = {"jsonrpc": "2.0", "method": method}
|
|
|
|
|
if params is not None:
|
|
|
|
|
message["params"] = params
|
|
|
|
|
self._write_message(message)
|
|
|
|
|
|
|
|
|
|
def next_notification(self) -> Notification:
|
|
|
|
|
item = self._notifications.get()
|
|
|
|
|
if isinstance(item, BaseException):
|
|
|
|
|
raise item
|
|
|
|
|
return item
|
|
|
|
|
|
|
|
|
|
def subscribe_notifications(
|
|
|
|
|
self,
|
|
|
|
|
notification_filter: NotificationFilter | None = None,
|
|
|
|
|
) -> "NotificationSubscription":
|
|
|
|
|
subscription_id = str(uuid.uuid4())
|
|
|
|
|
notifications: queue.Queue[Notification | BaseException] = queue.Queue()
|
|
|
|
|
with self._lock:
|
|
|
|
|
self._notification_subscribers[subscription_id] = (notifications, notification_filter)
|
|
|
|
|
return NotificationSubscription(self, subscription_id, notifications)
|
|
|
|
|
|
|
|
|
|
def subscribe_session_notifications(self, session_id: str) -> "NotificationSubscription":
|
2026-07-24 11:23:16 +08:00
|
|
|
"""Subscribe to a session and descendants discovered from subagent lifecycle edges."""
|
|
|
|
|
return self.subscribe_notifications(self._notification_belongs_to_session_tree(session_id))
|
2026-07-11 14:08:24 +08:00
|
|
|
|
|
|
|
|
def next_request(self) -> IncomingRequest:
|
|
|
|
|
item = self._requests.get()
|
|
|
|
|
if isinstance(item, BaseException):
|
|
|
|
|
raise item
|
|
|
|
|
return item
|
|
|
|
|
|
|
|
|
|
def respond(self, request_id: str | int, result: JsonValue) -> None:
|
|
|
|
|
self._write_message({"jsonrpc": "2.0", "id": request_id, "result": result})
|
|
|
|
|
|
|
|
|
|
def respond_error(
|
|
|
|
|
self,
|
|
|
|
|
request_id: str | int,
|
|
|
|
|
*,
|
|
|
|
|
code: int,
|
|
|
|
|
message: str,
|
|
|
|
|
data: JsonValue | None = None,
|
|
|
|
|
) -> None:
|
|
|
|
|
error: JsonObject = {"code": code, "message": message}
|
|
|
|
|
if data is not None:
|
|
|
|
|
error["data"] = data
|
|
|
|
|
self._write_message({"jsonrpc": "2.0", "id": request_id, "error": error})
|
|
|
|
|
|
|
|
|
|
def _request_raw(
|
|
|
|
|
self,
|
|
|
|
|
method: str,
|
|
|
|
|
params: JsonObject | None = None,
|
|
|
|
|
*,
|
|
|
|
|
timeout_seconds: float | None = None,
|
|
|
|
|
on_notification: Callable[[Notification], None] | None = None,
|
|
|
|
|
notification_filter: NotificationFilter | None = None,
|
|
|
|
|
notification_subscription: "NotificationSubscription | None" = None,
|
|
|
|
|
) -> JsonValue:
|
|
|
|
|
request_id = str(uuid.uuid4())
|
|
|
|
|
waiter: queue.Queue[JsonValue | BaseException] = queue.Queue(maxsize=1)
|
|
|
|
|
temp_subscription: NotificationSubscription | None = None
|
|
|
|
|
subscription = notification_subscription
|
|
|
|
|
with self._lock:
|
|
|
|
|
self._responses[request_id] = waiter
|
|
|
|
|
if on_notification is not None and subscription is None:
|
|
|
|
|
temp_subscription = self.subscribe_notifications(notification_filter)
|
|
|
|
|
subscription = temp_subscription
|
|
|
|
|
try:
|
|
|
|
|
message: JsonObject = {"jsonrpc": "2.0", "id": request_id, "method": method}
|
|
|
|
|
if params is not None:
|
|
|
|
|
message["params"] = params
|
|
|
|
|
self._write_message(message)
|
|
|
|
|
except BaseException:
|
|
|
|
|
with self._lock:
|
|
|
|
|
self._responses.pop(request_id, None)
|
|
|
|
|
if temp_subscription is not None:
|
|
|
|
|
temp_subscription.close()
|
|
|
|
|
raise
|
|
|
|
|
timeout = self.config.request_timeout_seconds if timeout_seconds is None else timeout_seconds
|
|
|
|
|
deadline = None if timeout is None else time.monotonic() + timeout
|
|
|
|
|
try:
|
|
|
|
|
while True:
|
|
|
|
|
if on_notification is not None and subscription is not None:
|
|
|
|
|
subscription.drain(on_notification)
|
|
|
|
|
wait_timeout = None
|
|
|
|
|
if on_notification is not None:
|
|
|
|
|
wait_timeout = 0.05
|
|
|
|
|
if deadline is not None:
|
|
|
|
|
remaining = deadline - time.monotonic()
|
|
|
|
|
if remaining <= 0:
|
|
|
|
|
with self._lock:
|
|
|
|
|
self._responses.pop(request_id, None)
|
2026-08-11 14:27:59 +08:00
|
|
|
diagnostics = self._runtime_diagnostics()
|
|
|
|
|
suffix = f"\n{diagnostics}" if diagnostics else ""
|
|
|
|
|
raise TimeoutError(
|
|
|
|
|
f"{method} timed out waiting for DeepSeek Harness runtime{suffix}"
|
|
|
|
|
)
|
2026-07-11 14:08:24 +08:00
|
|
|
wait_timeout = remaining if wait_timeout is None else min(wait_timeout, remaining)
|
|
|
|
|
try:
|
|
|
|
|
item = waiter.get(timeout=wait_timeout)
|
|
|
|
|
if on_notification is not None and subscription is not None:
|
|
|
|
|
subscription.drain(on_notification)
|
|
|
|
|
break
|
|
|
|
|
except queue.Empty:
|
|
|
|
|
continue
|
|
|
|
|
except BaseException:
|
|
|
|
|
with self._lock:
|
|
|
|
|
self._responses.pop(request_id, None)
|
|
|
|
|
if temp_subscription is not None:
|
|
|
|
|
temp_subscription.close()
|
|
|
|
|
raise
|
|
|
|
|
finally:
|
|
|
|
|
if temp_subscription is not None:
|
|
|
|
|
temp_subscription.close()
|
|
|
|
|
if isinstance(item, BaseException):
|
|
|
|
|
raise item
|
|
|
|
|
return item
|
|
|
|
|
|
|
|
|
|
def _write_message(self, message: JsonObject) -> None:
|
|
|
|
|
proc = self._proc
|
|
|
|
|
if proc is None or proc.stdin is None:
|
|
|
|
|
raise TransportClosedError("DeepSeek Harness runtime is not running")
|
|
|
|
|
try:
|
|
|
|
|
payload = json.dumps(message, separators=(",", ":")) + "\n"
|
|
|
|
|
with self._write_lock:
|
|
|
|
|
proc.stdin.write(payload)
|
|
|
|
|
proc.stdin.flush()
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
raise self._runtime_closed_error("Failed to write to DeepSeek Harness runtime") from exc
|
|
|
|
|
|
|
|
|
|
def _start_reader_thread(self) -> None:
|
|
|
|
|
self._reader_thread = threading.Thread(target=self._reader_loop, name="dsh-runtime-reader", daemon=True)
|
|
|
|
|
self._reader_thread.start()
|
|
|
|
|
|
|
|
|
|
def _start_stderr_thread(self) -> None:
|
|
|
|
|
self._stderr_thread = threading.Thread(target=self._stderr_loop, name="dsh-runtime-stderr", daemon=True)
|
|
|
|
|
self._stderr_thread.start()
|
|
|
|
|
|
|
|
|
|
def _reader_loop(self) -> None:
|
|
|
|
|
proc = self._proc
|
|
|
|
|
if proc is None or proc.stdout is None:
|
|
|
|
|
return
|
|
|
|
|
try:
|
|
|
|
|
for line in proc.stdout:
|
|
|
|
|
if not line.strip():
|
|
|
|
|
continue
|
|
|
|
|
try:
|
|
|
|
|
message = json.loads(line)
|
|
|
|
|
except json.JSONDecodeError:
|
|
|
|
|
continue
|
|
|
|
|
self._handle_message(message)
|
|
|
|
|
except BaseException as exc:
|
|
|
|
|
self._fail_waiters(exc)
|
|
|
|
|
finally:
|
|
|
|
|
self._fail_waiters(self._runtime_closed_error("DeepSeek Harness runtime stdout closed"))
|
|
|
|
|
|
|
|
|
|
def _stderr_loop(self) -> None:
|
|
|
|
|
proc = self._proc
|
|
|
|
|
if proc is None or proc.stderr is None:
|
|
|
|
|
return
|
|
|
|
|
for line in proc.stderr:
|
|
|
|
|
self._stderr_lines.append(line.rstrip())
|
|
|
|
|
|
|
|
|
|
def _handle_message(self, message: object) -> None:
|
|
|
|
|
if not isinstance(message, dict):
|
|
|
|
|
return
|
|
|
|
|
msg_id = message.get("id")
|
|
|
|
|
method = message.get("method")
|
|
|
|
|
if isinstance(msg_id, (str, int)) and isinstance(method, str):
|
|
|
|
|
params = message.get("params")
|
|
|
|
|
self._requests.put(IncomingRequest(id=msg_id, method=method, payload=params if isinstance(params, dict) else {}))
|
|
|
|
|
return
|
|
|
|
|
if isinstance(msg_id, (str, int)):
|
|
|
|
|
with self._lock:
|
|
|
|
|
waiter = self._responses.pop(str(msg_id), None)
|
|
|
|
|
if waiter is None:
|
|
|
|
|
return
|
|
|
|
|
if isinstance(message.get("error"), dict):
|
|
|
|
|
err = message["error"]
|
|
|
|
|
waiter.put(JsonRpcError(_int_or_none(err.get("code")), str(err.get("message", "JSON-RPC error")), err.get("data")))
|
|
|
|
|
else:
|
|
|
|
|
waiter.put(message.get("result"))
|
|
|
|
|
return
|
|
|
|
|
if isinstance(method, str):
|
|
|
|
|
params = message.get("params")
|
|
|
|
|
notification = Notification(method=method, payload=params if isinstance(params, dict) else {})
|
|
|
|
|
with self._lock:
|
2026-07-24 11:23:16 +08:00
|
|
|
self._record_session_relationship_locked(notification)
|
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 20:53:20 +08:00
|
|
|
subscribers = list(self._notification_subscribers.items())
|
2026-07-11 14:08:24 +08:00
|
|
|
delivered = False
|
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 20:53:20 +08:00
|
|
|
for subscription_id, (subscriber, predicate) in subscribers:
|
|
|
|
|
try:
|
|
|
|
|
matches = predicate is None or predicate(notification)
|
|
|
|
|
except BaseException as exc:
|
|
|
|
|
with self._lock:
|
|
|
|
|
current = self._notification_subscribers.get(subscription_id)
|
|
|
|
|
if current is not None and current[0] is subscriber:
|
|
|
|
|
self._notification_subscribers.pop(subscription_id, None)
|
|
|
|
|
subscriber.put(exc)
|
|
|
|
|
continue
|
|
|
|
|
if matches:
|
2026-07-11 14:08:24 +08:00
|
|
|
subscriber.put(notification)
|
|
|
|
|
delivered = True
|
|
|
|
|
if not delivered:
|
|
|
|
|
self._notifications.put(notification)
|
|
|
|
|
|
|
|
|
|
def _fail_waiters(self, exc: BaseException) -> None:
|
|
|
|
|
with self._lock:
|
|
|
|
|
waiters = list(self._responses.values())
|
|
|
|
|
self._responses.clear()
|
|
|
|
|
subscribers = list(self._notification_subscribers.values())
|
|
|
|
|
self._notification_subscribers.clear()
|
|
|
|
|
for waiter in waiters:
|
|
|
|
|
waiter.put(exc)
|
|
|
|
|
for subscriber, _predicate in subscribers:
|
|
|
|
|
subscriber.put(exc)
|
|
|
|
|
self._notifications.put(exc)
|
|
|
|
|
self._requests.put(exc)
|
|
|
|
|
|
|
|
|
|
def _runtime_closed_error(self, reason: str) -> TransportClosedError:
|
2026-08-11 14:27:59 +08:00
|
|
|
diagnostics = self._runtime_diagnostics()
|
|
|
|
|
return TransportClosedError(f"{reason}\n{diagnostics}" if diagnostics else reason)
|
|
|
|
|
|
|
|
|
|
def _runtime_diagnostics(self) -> str:
|
|
|
|
|
"""Return available subprocess state for transport failures and timeouts."""
|
2026-07-11 14:08:24 +08:00
|
|
|
proc = self._proc
|
|
|
|
|
if (
|
|
|
|
|
proc is not None
|
|
|
|
|
and proc.poll() is not None
|
|
|
|
|
and self._stderr_thread is not None
|
|
|
|
|
and self._stderr_thread.is_alive()
|
|
|
|
|
and threading.current_thread() is not self._stderr_thread
|
|
|
|
|
):
|
|
|
|
|
self._stderr_thread.join(timeout=0.1)
|
|
|
|
|
|
2026-08-11 14:27:59 +08:00
|
|
|
parts: list[str] = []
|
2026-07-11 14:08:24 +08:00
|
|
|
if proc is not None:
|
|
|
|
|
exit_code = proc.poll()
|
|
|
|
|
if exit_code is not None:
|
|
|
|
|
parts.append(f"exit code: {exit_code}")
|
|
|
|
|
if self._stderr_lines:
|
|
|
|
|
parts.append("stderr tail:\n" + "\n".join(self._stderr_lines))
|
2026-08-11 14:27:59 +08:00
|
|
|
return "\n".join(parts)
|
2026-07-11 14:08:24 +08:00
|
|
|
|
feat(python-sdk): launch dsh profiles from explicit homes
Replace complete-config, session_root, runtime-bin, bridge-bin, and public argv override options with dsh_bin, profile, ordered patches, and dsh_home. Resolve executable/home/patch/cwd paths before spawn, select the sdk profile by default, and fail before launch unless dsh_home or non-empty DSH_HOME is explicit; Python never inherits ~/.dsh silently.
Remove Python-owned DSH_CORDIS_CONFIG, DSH_SESSION_ROOT, and DSH_CWD injection and drop session_root from RunResult. Keep arbitrary argv only as an underscore-prefixed fake-runtime adapter, retain provider/model/token and process controls, and append subprocess stderr to initialization JSON-RPC errors so profile boot failures name their actual plugin cause. Unit and carrier tests cover both exe and Node modes.
2026-08-23 14:56:16 +08:00
|
|
|
def _default_launch_args(self, env: dict[str, str]) -> tuple[str, ...]:
|
|
|
|
|
if self.config.dsh_bin is None:
|
|
|
|
|
try:
|
|
|
|
|
from deepseek_harness_runtime import resolve_bundled_launch_args
|
|
|
|
|
except ImportError as exc:
|
|
|
|
|
raise FileNotFoundError(
|
|
|
|
|
"Unable to locate the bundled DeepSeek Harness dsh runtime. "
|
|
|
|
|
"Install deepseek-harness-runtime-bin."
|
|
|
|
|
) from exc
|
|
|
|
|
base = resolve_bundled_launch_args()
|
|
|
|
|
else:
|
|
|
|
|
base = (str(Path(self.config.dsh_bin).expanduser().resolve()),)
|
|
|
|
|
|
|
|
|
|
if self.config.dsh_home is not None:
|
|
|
|
|
if not self.config.dsh_home.strip():
|
|
|
|
|
raise ValueError("HarnessConfig requires a non-empty dsh_home")
|
|
|
|
|
env["DSH_HOME"] = str(Path(self.config.dsh_home).expanduser().resolve())
|
|
|
|
|
elif not env.get("DSH_HOME", "").strip():
|
|
|
|
|
raise ValueError(
|
|
|
|
|
"HarnessConfig requires an explicit dsh_home or non-empty DSH_HOME; "
|
|
|
|
|
"the Python SDK never uses ~/.dsh implicitly"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
patches = tuple(
|
|
|
|
|
argument
|
|
|
|
|
for patch in self.config.patches
|
|
|
|
|
for argument in ("--patch", str(Path(patch).expanduser().resolve()))
|
2026-07-11 22:54:42 +08:00
|
|
|
)
|
feat(python-sdk): launch dsh profiles from explicit homes
Replace complete-config, session_root, runtime-bin, bridge-bin, and public argv override options with dsh_bin, profile, ordered patches, and dsh_home. Resolve executable/home/patch/cwd paths before spawn, select the sdk profile by default, and fail before launch unless dsh_home or non-empty DSH_HOME is explicit; Python never inherits ~/.dsh silently.
Remove Python-owned DSH_CORDIS_CONFIG, DSH_SESSION_ROOT, and DSH_CWD injection and drop session_root from RunResult. Keep arbitrary argv only as an underscore-prefixed fake-runtime adapter, retain provider/model/token and process controls, and append subprocess stderr to initialization JSON-RPC errors so profile boot failures name their actual plugin cause. Unit and carrier tests cover both exe and Node modes.
2026-08-23 14:56:16 +08:00
|
|
|
return (*base, "--profile", self.config.profile, *patches)
|
2026-07-11 22:54:42 +08:00
|
|
|
|
2026-07-11 14:08:24 +08:00
|
|
|
def _unsubscribe_notifications(self, subscription_id: str) -> None:
|
|
|
|
|
with self._lock:
|
|
|
|
|
self._notification_subscribers.pop(subscription_id, None)
|
|
|
|
|
|
2026-07-24 11:23:16 +08:00
|
|
|
def _record_session_relationship_locked(self, notification: Notification) -> None:
|
2026-07-24 12:21:10 +08:00
|
|
|
if notification.method != "subagent.started":
|
2026-07-24 11:23:16 +08:00
|
|
|
return
|
|
|
|
|
parent_id = notification.payload.get("parentSessionId")
|
|
|
|
|
child_id = notification.payload.get("childSessionId")
|
|
|
|
|
if (
|
|
|
|
|
isinstance(parent_id, str)
|
|
|
|
|
and parent_id
|
|
|
|
|
and isinstance(child_id, str)
|
|
|
|
|
and child_id
|
|
|
|
|
and parent_id != child_id
|
|
|
|
|
):
|
|
|
|
|
self._session_parents[child_id] = parent_id
|
|
|
|
|
|
|
|
|
|
def _notification_belongs_to_session_tree(self, session_id: str) -> NotificationFilter:
|
|
|
|
|
def belongs(notification: Notification) -> bool:
|
|
|
|
|
payload = notification.payload
|
2026-07-24 12:21:10 +08:00
|
|
|
if notification.method in {"subagent.started", "subagent.finished"}:
|
|
|
|
|
parent_id = payload.get("parentSessionId")
|
|
|
|
|
if (
|
|
|
|
|
isinstance(parent_id, str)
|
|
|
|
|
and self._session_is_descendant_of(parent_id, session_id)
|
|
|
|
|
):
|
|
|
|
|
return True
|
|
|
|
|
return payload.get("childSessionId") == session_id
|
|
|
|
|
related_id = payload.get("sessionId")
|
|
|
|
|
return (
|
2026-07-24 11:23:16 +08:00
|
|
|
isinstance(related_id, str)
|
|
|
|
|
and self._session_is_descendant_of(related_id, session_id)
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
return belongs
|
|
|
|
|
|
|
|
|
|
def _session_is_descendant_of(self, session_id: str, root_session_id: str) -> bool:
|
|
|
|
|
current = session_id
|
|
|
|
|
visited: set[str] = set()
|
|
|
|
|
while current not in visited:
|
|
|
|
|
if current == root_session_id:
|
|
|
|
|
return True
|
|
|
|
|
visited.add(current)
|
|
|
|
|
parent = self._session_parents.get(current)
|
|
|
|
|
if parent is None:
|
|
|
|
|
return False
|
|
|
|
|
current = parent
|
|
|
|
|
return False
|
|
|
|
|
|
2026-07-11 14:08:24 +08:00
|
|
|
|
|
|
|
|
class NotificationSubscription:
|
|
|
|
|
def __init__(
|
|
|
|
|
self,
|
|
|
|
|
client: HarnessClient,
|
|
|
|
|
subscription_id: str,
|
|
|
|
|
notifications: queue.Queue[Notification | BaseException],
|
|
|
|
|
) -> None:
|
|
|
|
|
self._client = client
|
|
|
|
|
self._subscription_id = subscription_id
|
|
|
|
|
self._notifications = notifications
|
|
|
|
|
self._closed = False
|
|
|
|
|
|
|
|
|
|
def __enter__(self) -> "NotificationSubscription":
|
|
|
|
|
return self
|
|
|
|
|
|
|
|
|
|
def __exit__(self, _exc_type, _exc, _tb) -> None:
|
|
|
|
|
self.close()
|
|
|
|
|
|
|
|
|
|
def close(self) -> None:
|
|
|
|
|
if self._closed:
|
|
|
|
|
return
|
|
|
|
|
self._closed = True
|
|
|
|
|
self._client._unsubscribe_notifications(self._subscription_id)
|
|
|
|
|
|
|
|
|
|
def next(self) -> Notification:
|
|
|
|
|
item = self._notifications.get()
|
|
|
|
|
if isinstance(item, BaseException):
|
|
|
|
|
raise item
|
|
|
|
|
return item
|
|
|
|
|
|
|
|
|
|
def drain(self, on_notification: Callable[[Notification], None]) -> None:
|
|
|
|
|
while True:
|
|
|
|
|
try:
|
|
|
|
|
item = self._notifications.get_nowait()
|
|
|
|
|
except queue.Empty:
|
|
|
|
|
return
|
|
|
|
|
if isinstance(item, BaseException):
|
|
|
|
|
raise item
|
|
|
|
|
on_notification(item)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class _SessionPromptResponse(BaseModel):
|
2026-07-30 16:48:28 +08:00
|
|
|
messageId: str
|
2026-07-11 14:08:24 +08:00
|
|
|
|
|
|
|
|
|
|
|
|
|
class _ShutdownResponse(BaseModel):
|
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _int_or_none(value: object) -> int | None:
|
|
|
|
|
return value if isinstance(value, int) else None
|