1
from __future__ import annotations3
import uuid4
from dataclasses import dataclass, field5
from pathlib import Path6
from typing import Callable8
from .client import HarnessClient, HarnessConfig9
from .errors import SdkProtocolError10
from .models import JsonObject, Notification13
@dataclass(slots=True)14
class DeepSeekHarnessConfig:15
"""Configuration for launching the local DeepSeek Harness SDK runtime.17
The runtime inherits the caller's environment by default, so existing18
DEEPSEEK_API_KEY and DEEPSEEK_BASE_URL settings keep working. Use ``env`` to19
intentionally override or inject variables for a subprocess.20
"""22
provider: str = "deepseek-official"23
model: str = "deepseek-v4-flash"24
reasoning_effort: str | None = None25
max_tokens: int | None = None26
cwd: str | None = None27
runtime_cwd: str | None = None28
dsh_bin: str | None = None29
profile: str = "sdk"30
patches: tuple[str, ...] = ()31
dsh_home: str | None = None32
env: dict[str, str] = field(default_factory=dict)33
initialize_timeout_seconds: float = 30.034
request_timeout_seconds: float | None = None35
shutdown_timeout_seconds: float | None = 1.036
base_url: str | None = None37
api_key: str | None = None40
@dataclass(slots=True)41
class RunResult:42
session_id: str43
final_response: str44
finish_reason: str | None45
events: list[JsonObject]46
notifications: list[Notification]49
class DeepSeekHarness:50
"""Reusable synchronous SDK for running DeepSeek Harness agent turns.52
The runtime subprocess starts lazily and remains owned by this instance53
across calls to :meth:`run`. Use the instance as a context manager, or call54
:meth:`close` explicitly when finished, so the subprocess is always reaped.55
"""57
def __init__(58
self,59
config: DeepSeekHarnessConfig | None = None,60
*,61
_launch_args: tuple[str, ...] | None = None,62
**kwargs: object,63
) -> None:64
if config is not None and kwargs:65
raise TypeError("pass either DeepSeekHarnessConfig or keyword options, not both")66
self.config = config or DeepSeekHarnessConfig(**kwargs)67
cwd = str(Path(self.config.cwd or Path.cwd()).resolve())68
runtime_cwd = str(Path(self.config.runtime_cwd).resolve()) if self.config.runtime_cwd is not None else cwd69
self._cwd = cwd70
env = dict(self.config.env)71
if self.config.base_url is not None:72
env["DEEPSEEK_BASE_URL"] = self.config.base_url73
if self.config.api_key is not None:74
env["DEEPSEEK_API_KEY"] = self.config.api_key76
self._client = HarnessClient(77
HarnessConfig(78
dsh_bin=self.config.dsh_bin,79
profile=self.config.profile,80
patches=self.config.patches,81
dsh_home=self.config.dsh_home,82
cwd=runtime_cwd,83
env=env,84
initialize_timeout_seconds=self.config.initialize_timeout_seconds,85
request_timeout_seconds=self.config.request_timeout_seconds,86
shutdown_timeout_seconds=self.config.shutdown_timeout_seconds,87
),88
_launch_args=_launch_args,89
)90
self._initialized = False92
def __enter__(self) -> "DeepSeekHarness":93
self.start()94
return self96
def __exit__(self, _exc_type, _exc, _tb) -> None:97
self.close()99
@property100
def client(self) -> HarnessClient:101
return self._client103
def start(self) -> None:104
if self._initialized:105
return106
self._client.start()107
self._client.initialize(108
cwd=self._cwd,109
provider=self.config.provider,110
model=self.config.model,111
reasoning_effort=self.config.reasoning_effort,112
max_tokens=self.config.max_tokens,113
)114
self._initialized = True116
def close(self) -> None:117
self._client.close()118
self._initialized = False120
def start_session(self, session_id: str | None = None) -> "Session":121
self.start()122
return Session(self, session_id or f"session-{uuid.uuid4().hex}")124
def run(125
self,126
input: str | list[JsonObject],127
*,128
session_id: str | None = None,129
on_notification: Callable[[Notification], None] | None = None,130
) -> RunResult:131
return self.start_session(session_id).run(input, on_notification=on_notification)134
class Session:135
def __init__(self, harness: DeepSeekHarness, session_id: str) -> None:136
self.harness = harness137
self.id = session_id139
def run(140
self,141
input: str | list[JsonObject],142
*,143
on_notification: Callable[[Notification], None] | None = None,144
) -> RunResult:145
content_blocks = normalize_input(input)146
notifications: list[Notification] = []147
events: list[JsonObject] = []149
def collect(notification: Notification) -> None:150
notifications.append(notification)151
if on_notification is not None:152
on_notification(notification)153
if (154
notification.method == "session.event"155
and notification.payload.get("sessionId") == self.id156
):157
event = notification.payload.get("event")158
if isinstance(event, dict):159
events.append(event)161
with self.harness.client.subscribe_session_notifications(self.id) as subscription:162
message_id = self.harness.client.session_prompt(163
self.id,164
content_blocks,165
notification_subscription=subscription,166
)168
received = False169
while True:170
notification = subscription.next()171
if not received:172
if not _is_inbox_receipt(notification, self.id, message_id):173
continue174
received = True175
collect(notification)176
if (177
notification.method == "session.status"178
and notification.payload.get("sessionId") == self.id179
and notification.payload.get("status") == "idle"180
):181
break183
return RunResult(184
session_id=self.id,185
final_response=final_response(events),186
finish_reason=finish_reason(events),187
events=events,188
notifications=notifications,189
)192
def _is_inbox_receipt(notification: Notification, session_id: str, message_id: str) -> bool:193
if notification.method != "session.event" or notification.payload.get("sessionId") != session_id:194
return False195
event = notification.payload.get("event")196
if not isinstance(event, dict) or event.get("type") != "agent/inbox/spliced":197
return False198
data = event.get("data")199
inserted = data.get("inserted") if isinstance(data, dict) else None200
return isinstance(inserted, list) and any(201
isinstance(message, dict) and message.get("id") == message_id for message in inserted202
)205
def normalize_input(input: str | list[JsonObject]) -> list[JsonObject]:206
if isinstance(input, str):207
return [{"type": "text", "text": input}]208
return input211
def final_response(events: list[JsonObject]) -> str:212
for event in reversed(events):213
if event.get("type") != "assistant/message":214
continue215
data = event.get("data")216
if not isinstance(data, dict):217
continue218
message = data.get("message")219
content_owner = message if isinstance(message, dict) else data220
content = content_owner.get("content")221
if not isinstance(content, list):222
continue223
parts: list[str] = []224
for block in content:225
if isinstance(block, dict) and block.get("type") == "text":226
parts.append(str(block.get("text") or ""))227
return "".join(parts)228
return ""231
def finish_reason(events: list[JsonObject]) -> str | None:232
"""Return the last turn-ending kind.234
The input must contain root-session events from one owned run interval.236
Raises:237
SdkProtocolError: The last ``turn/end`` has no string reason kind.238
"""239
for event in reversed(events):240
if event.get("type") != "turn/end":241
continue242
data = event.get("data")243
reason = data.get("reason") if isinstance(data, dict) else None244
kind = reason.get("kind") if isinstance(reason, dict) else None245
if not isinstance(kind, str):246
raise SdkProtocolError("turn/end event requires a string data.reason.kind")247
return kind248
return None