返回源码地图

python/sdk/src/deepseek_harness/api.py

main snapshot · da00f7f5358f · 正文引用章节 14 / 14;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1from __future__ import annotations
2
3import uuid
4from dataclasses import dataclass, field
5from pathlib import Path
6from typing import Callable
7
8from .client import HarnessClient, HarnessConfig
9from .errors import SdkProtocolError
10from .models import JsonObject, Notification
11
12
13@dataclass(slots=True)
14class DeepSeekHarnessConfig:
15 """Configuration for launching the local DeepSeek Harness SDK runtime.
16
17 The runtime inherits the caller's environment by default, so existing
18 DEEPSEEK_API_KEY and DEEPSEEK_BASE_URL settings keep working. Use ``env`` to
19 intentionally override or inject variables for a subprocess.
20 """
21
22 provider: str = "deepseek-official"
23 model: str = "deepseek-v4-flash"
24 reasoning_effort: str | None = None
25 max_tokens: int | None = None
26 cwd: str | None = None
27 runtime_cwd: str | None = None
28 dsh_bin: str | None = None
29 profile: str = "sdk"
30 patches: tuple[str, ...] = ()
31 dsh_home: str | None = None
32 env: dict[str, str] = field(default_factory=dict)
33 initialize_timeout_seconds: float = 30.0
34 request_timeout_seconds: float | None = None
35 shutdown_timeout_seconds: float | None = 1.0
36 base_url: str | None = None
37 api_key: str | None = None
38
39
40@dataclass(slots=True)
41class RunResult:
42 session_id: str
43 final_response: str
44 finish_reason: str | None
45 events: list[JsonObject]
46 notifications: list[Notification]
47
48
49class DeepSeekHarness:
50 """Reusable synchronous SDK for running DeepSeek Harness agent turns.
51
52 The runtime subprocess starts lazily and remains owned by this instance
53 across calls to :meth:`run`. Use the instance as a context manager, or call
54 :meth:`close` explicitly when finished, so the subprocess is always reaped.
55 """
56
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 cwd
69 self._cwd = cwd
70 env = dict(self.config.env)
71 if self.config.base_url is not None:
72 env["DEEPSEEK_BASE_URL"] = self.config.base_url
73 if self.config.api_key is not None:
74 env["DEEPSEEK_API_KEY"] = self.config.api_key
75
76 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 = False
91
92 def __enter__(self) -> "DeepSeekHarness":
93 self.start()
94 return self
95
96 def __exit__(self, _exc_type, _exc, _tb) -> None:
97 self.close()
98
99 @property
100 def client(self) -> HarnessClient:
101 return self._client
102
103 def start(self) -> None:
104 if self._initialized:
105 return
106 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 = True
115
116 def close(self) -> None:
117 self._client.close()
118 self._initialized = False
119
120 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}")
123
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)
132
133
134class Session:
135 def __init__(self, harness: DeepSeekHarness, session_id: str) -> None:
136 self.harness = harness
137 self.id = session_id
138
139 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] = []
148
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.id
156 ):
157 event = notification.payload.get("event")
158 if isinstance(event, dict):
159 events.append(event)
160
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 )
167
168 received = False
169 while True:
170 notification = subscription.next()
171 if not received:
172 if not _is_inbox_receipt(notification, self.id, message_id):
173 continue
174 received = True
175 collect(notification)
176 if (
177 notification.method == "session.status"
178 and notification.payload.get("sessionId") == self.id
179 and notification.payload.get("status") == "idle"
180 ):
181 break
182
183 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 )
190
191
192def _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 False
195 event = notification.payload.get("event")
196 if not isinstance(event, dict) or event.get("type") != "agent/inbox/spliced":
197 return False
198 data = event.get("data")
199 inserted = data.get("inserted") if isinstance(data, dict) else None
200 return isinstance(inserted, list) and any(
201 isinstance(message, dict) and message.get("id") == message_id for message in inserted
202 )
203
204
205def normalize_input(input: str | list[JsonObject]) -> list[JsonObject]:
206 if isinstance(input, str):
207 return [{"type": "text", "text": input}]
208 return input
209
210
211def final_response(events: list[JsonObject]) -> str:
212 for event in reversed(events):
213 if event.get("type") != "assistant/message":
214 continue
215 data = event.get("data")
216 if not isinstance(data, dict):
217 continue
218 message = data.get("message")
219 content_owner = message if isinstance(message, dict) else data
220 content = content_owner.get("content")
221 if not isinstance(content, list):
222 continue
223 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 ""
229
230
231def finish_reason(events: list[JsonObject]) -> str | None:
232 """Return the last turn-ending kind.
233
234 The input must contain root-session events from one owned run interval.
235
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 continue
242 data = event.get("data")
243 reason = data.get("reason") if isinstance(data, dict) else None
244 kind = reason.get("kind") if isinstance(reason, dict) else None
245 if not isinstance(kind, str):
246 raise SdkProtocolError("turn/end event requires a string data.reason.kind")
247 return kind
248 return None