490 lines
15 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""DeepSeek Harness headless 适配器。"""
from __future__ import annotations
import json
import os
import re
import signal
import subprocess
import time
from dataclasses import dataclass
from pathlib import Path
from typing import Any, Callable, Mapping, Protocol, Sequence
from framework.primitives.execution import FrameworkExecutionRequest, FrameworkExecutionResult
from .normalization import (
DshSessionOutcome,
normalize_dsh_session,
)
from .pi_ai import PiAiRouteConfig
DEFAULT_DSH_BIN = "dsh"
DEFAULT_DSH_PROFILE = "headless"
_MAX_STDERR_BYTES = 64 * 1024
_SECRET_PATTERN = re.compile(r"sk-[A-Za-z0-9_-]{12,}")
# 这些行是 DSH 当前 headless bundle 的模型可见工具或上下文自动注入入口。
# 适配器现在只承诺无工具任务;未知/新增能力不应默默进入 Muse 角色。
_DISABLED_HEADLESS_ROWS = (
"session-title-llm",
"agent-instructions",
"skill",
"skill-filesystem",
"tool-skill",
"tool-bash",
"tool-pwsh",
"tool-jobs",
"tool-fs",
"tool-fs-search",
"web",
"web-search-deepseek",
"tool-web",
"goal",
"goal-round-driver",
"command-goal",
"subagent",
"subagent-spawn-in-process",
"subagent-fork-in-process",
"tool-subagent-control",
"tool-subagent-list-agents",
"tool-subagent",
"tool-subagent-fork",
"tool-subagent-report",
"workflow-worker-thread",
"tool-workflow",
"tool-todo",
"tool-goal",
"tool-ralph",
"tool-str-replace-editor",
)
class TraceSink(Protocol):
"""Muse 注入的旁路事件接收端;框架不拥有持久化实现。"""
def emit(self, event_type: str, **kwargs: Any) -> int: ...
class DshError(RuntimeError):
"""DSH 启动、会话回放或执行失败。"""
def __init__(
self,
error_code: str,
message: str,
*,
outcome: "DshRunOutcome | None" = None,
) -> None:
super().__init__(message)
self.error_code = error_code
self.outcome = outcome
@dataclass(frozen=True)
class DshExecutionPolicy:
"""DSH 运行策略;模型路由由 Muse 派发方显式解析后传入。"""
provider: str | None = None
model: str | None = None
dsh_bin: str = DEFAULT_DSH_BIN
profile: str = DEFAULT_DSH_PROFILE
cwd: str | None = None
dsh_home: str | None = None
permission_mode: str = "read-only"
pi_ai_route: PiAiRouteConfig | None = None
def __post_init__(self) -> None:
if not isinstance(self.provider, str) or not self.provider.strip():
raise ValueError("provider 必须显式传入")
if not isinstance(self.model, str) or not self.model.strip():
raise ValueError("model 必须显式传入")
if self.profile != DEFAULT_DSH_PROFILE:
raise ValueError("当前 DSH 适配器只支持 headless profile")
if self.permission_mode not in {"read-only", "workspace-write", "danger-full-access"}:
raise ValueError("permission_mode 不受支持")
if self.pi_ai_route is not None:
if self.pi_ai_route.provider != self.provider:
raise ValueError("pi_ai_route.provider 必须与 policy.provider 一致")
if self.pi_ai_route.model != self.model:
raise ValueError("pi_ai_route.model 必须与 policy.model 一致")
@property
def framework(self) -> str:
return "dsh"
@property
def requested_model_id(self) -> str:
return f"{self.provider}/{self.model}"
@dataclass(frozen=True)
class DshProcessResult:
"""一次 dsh 子进程的无正文进程结果;正文从会话日志读取。"""
returncode: int
stdout: bytes = b""
stderr: bytes = b""
timed_out: bool = False
@dataclass(frozen=True)
class DshRunOutcome:
"""DSH 执行结果与已发布通用事件工件。"""
session: DshSessionOutcome
requested_model: str
exit_code: int
duration_ms: int
stderr_summary: str = ""
@property
def session_id(self) -> str:
return self.session.session_id
@property
def final_text(self) -> str | None:
return self.session.final_text
@property
def model_calls(self):
return self.session.model_calls
@property
def tool_calls(self):
return self.session.tool_calls
@property
def turns(self) -> int:
return self.session.turns
@property
def unknown_event_types(self):
return self.session.unknown_event_types
def as_framework_result(self) -> FrameworkExecutionResult:
"""转换为框架通用结果;不代表 Muse 候选已通过任何业务门。"""
return FrameworkExecutionResult(
status="completed",
final_text=self.final_text,
requested_model=self.requested_model,
actual_models=tuple(call.actual_model_id for call in self.session.model_calls),
session_id=self.session_id,
artifact_locator=self.session.artifact_path,
trace_digest=self.session.artifact_sha256,
)
def build_dsh_patch(
request: FrameworkExecutionRequest,
policy: DshExecutionPolicy,
*,
session_root: str | Path,
) -> list[dict[str, Any]]:
"""生成一次性 JSON/YAML patch;JSON 是 YAML 的安全子集,避免标量注入。"""
if request.tool_allowlist:
raise DshError(
"DSH_TOOL_POLICY_UNSUPPORTED",
"当前 DSH headless 适配器尚未接入 Muse 只读工具插件,非空工具白名单拒绝启动",
)
cwd = Path(policy.cwd or os.getcwd()).resolve()
patch = [
{
"id": "agent-default-model",
"config": {"provider": policy.provider, "model": policy.model},
},
]
if policy.pi_ai_route is not None:
patch.append(
{
"id": "llm-pi-ai",
"config": {
"providers": {
policy.pi_ai_route.provider: policy.pi_ai_route.as_profile(),
}
},
}
)
patch.extend([
{
"id": "system-prompt",
"config": {"persona": request.system_prompt},
},
{
"id": "session-persistence-jsonl",
"config": {
"root": str(Path(session_root).resolve()),
"compression": "none",
"packChunks": False,
},
},
{
"id": "sandbox-policy",
"config": {
"mode": policy.permission_mode,
"workspaceRoot": str(cwd),
},
},
{
"id": "tools",
"config": {"mode": "native"},
},
*({"id": row_id, "disabled": True} for row_id in _DISABLED_HEADLESS_ROWS),
])
return patch
def build_dsh_argv(
request: FrameworkExecutionRequest,
policy: DshExecutionPolicy,
patch_path: str | Path,
) -> list[str]:
"""构造 headless argv;任务只作为一个 app-owned positional 参数传入。"""
if request.session_mode != "fresh":
raise DshError(
"DSH_SESSION_CONTINUE_UNSUPPORTED",
"当前 DSH headless 适配器不支持 continue,会话续接需接入 DSH resume 入口后再开放",
)
if request.tool_allowlist:
raise DshError(
"DSH_TOOL_POLICY_UNSUPPORTED",
"当前 DSH headless 适配器尚未接入 Muse 只读工具插件,非空工具白名单拒绝启动",
)
return [
policy.dsh_bin,
"--profile",
policy.profile,
"--patch",
str(Path(patch_path).resolve()),
"--",
request.user_content,
]
def _redact_stderr(raw: bytes) -> str:
text = raw[:_MAX_STDERR_BYTES].decode("utf-8", errors="replace")
text = _SECRET_PATTERN.sub("<redacted>", text)
return text.strip()
def _run_subprocess(
argv: Sequence[str],
*,
timeout_seconds: float,
cwd: str | None,
env: Mapping[str, str],
) -> DshProcessResult:
try:
process = subprocess.Popen(
list(argv),
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
cwd=cwd,
env=dict(env),
start_new_session=(os.name != "nt"),
)
except OSError as exc:
raise DshError("DSH_START_FAILED", f"DSH 进程启动失败: {type(exc).__name__}") from exc
try:
stdout, stderr = process.communicate(timeout=timeout_seconds)
except subprocess.TimeoutExpired:
if os.name != "nt":
try:
os.killpg(process.pid, signal.SIGKILL)
except ProcessLookupError:
pass
else:
process.kill()
stdout, stderr = process.communicate()
return DshProcessResult(
returncode=process.returncode if process.returncode is not None else -9,
stdout=stdout or b"",
stderr=stderr or b"",
timed_out=True,
)
return DshProcessResult(
returncode=process.returncode,
stdout=stdout or b"",
stderr=stderr or b"",
)
class DshHeadlessRunner:
"""启动 DSH headless,并以 flush 后的会话日志作为唯一轨迹输入。"""
def __init__(
self,
launcher: Callable[..., DshProcessResult] | None = None,
*,
framework_version: str = "unknown",
) -> None:
self._launcher = launcher
self._framework_version = framework_version
def _launch(
self,
argv: Sequence[str],
*,
timeout_seconds: float,
cwd: str | None,
env: Mapping[str, str],
) -> DshProcessResult:
if self._launcher is not None:
return self._launcher(argv, timeout_seconds, cwd, env)
return _run_subprocess(
argv,
timeout_seconds=timeout_seconds,
cwd=cwd,
env=env,
)
@staticmethod
def _find_session(session_root: Path) -> Path:
candidates = sorted(session_root.rglob("session.jsonl"))
if len(candidates) != 1:
raise DshError(
"DSH_SESSION_MISSING" if not candidates else "DSH_SESSION_AMBIGUOUS",
"DSH 未产生唯一的明文 session.jsonl 工件",
)
return candidates[0]
def run(
self,
request: FrameworkExecutionRequest,
policy: DshExecutionPolicy,
sink: TraceSink,
*,
artifact_dir: str | Path,
timeout_seconds: float | None = None,
run_id: str | None = None,
) -> DshRunOutcome:
"""执行一次 DSH 任务;模型、工具和业务 Schema 由上层合同负责。"""
if request.session_mode != "fresh":
raise DshError(
"DSH_SESSION_CONTINUE_UNSUPPORTED",
"当前 DSH headless 适配器不支持 continue",
)
if request.tool_allowlist:
raise DshError(
"DSH_TOOL_POLICY_UNSUPPORTED",
"当前 DSH headless 适配器尚未接入 Muse 只读工具插件,非空工具白名单拒绝启动",
)
root = Path(artifact_dir).resolve()
root.mkdir(parents=True, mode=0o700, exist_ok=True)
root.chmod(0o700)
session_root = root / "dsh-sessions"
session_root.mkdir(parents=True, mode=0o700, exist_ok=True)
session_root.chmod(0o700)
patch_path = root / "dsh.patch.json"
patch = build_dsh_patch(request, policy, session_root=session_root)
patch_path.write_text(json.dumps(patch, ensure_ascii=False, indent=2) + "\n", encoding="utf-8")
patch_path.chmod(0o600)
argv = build_dsh_argv(request, policy, patch_path)
env = os.environ.copy()
env["DSH_PERMISSION_MODE"] = policy.permission_mode
env["DSH_TELEMETRY_DISABLED"] = "1"
if policy.dsh_home:
env["DSH_HOME"] = str(Path(policy.dsh_home).resolve())
cwd = str(Path(policy.cwd).resolve()) if policy.cwd else None
timeout = float(timeout_seconds if timeout_seconds is not None else request.timeout_seconds)
started = time.monotonic()
sink.emit(
"agent.started",
status="ok",
requested_model_id=policy.requested_model_id,
details={"framework": "dsh", "profile": policy.profile},
)
process_result = self._launch(
argv,
timeout_seconds=timeout,
cwd=cwd,
env=env,
)
duration_ms = int((time.monotonic() - started) * 1000)
stderr_summary = _redact_stderr(process_result.stderr)
if process_result.timed_out:
raise DshError(
"DSH_TIMEOUT",
f"DSH 执行超时(>{timeout}s)",
)
if process_result.returncode != 0:
suffix = f": {stderr_summary[:256]}" if stderr_summary else ""
raise DshError(
"DSH_EXIT_NONZERO",
f"DSH 进程退出码 {process_result.returncode}{suffix}",
)
session_path = self._find_session(session_root)
try:
session = normalize_dsh_session(
session_path,
root / "framework-events.jsonl",
run_id=run_id,
framework_version=self._framework_version,
)
except ValueError as exc:
raise DshError("DSH_SESSION_INVALID", str(exc)) from exc
if not session.model_calls:
raise DshError("DSH_NO_MODEL_RESPONSE", "DSH 会话没有 assistant/message 模型回合")
if session.final_text is None or not session.final_text.strip():
raise DshError("DSH_EMPTY_FINAL_MESSAGE", "DSH 会话没有最终文本")
for call in session.model_calls:
sink.emit(
"model.completed",
status="error" if call.stop_reason in {"error", "interrupted"} else "ok",
requested_model_id=policy.requested_model_id,
actual_model_id=call.actual_model_id,
usage=dict(call.usage),
cost_usd=call.cost_usd,
details={"stopReason": call.stop_reason, "provider": call.provider},
)
for call in session.tool_calls:
sink.emit(
"tool.completed",
status="error" if call.is_error else "ok",
tool_name=call.name,
details={"toolCallId": call.tool_call_id},
)
outcome = DshRunOutcome(
session=session,
requested_model=policy.requested_model_id,
exit_code=process_result.returncode,
duration_ms=duration_ms,
stderr_summary=stderr_summary,
)
sink.emit(
"agent.completed",
status="ok",
requested_model_id=policy.requested_model_id,
actual_model_id=session.model_calls[-1].actual_model_id,
details={
"sessionId": session.session_id,
"turns": session.turns,
"modelCalls": len(session.model_calls),
"toolCalls": len(session.tool_calls),
"unknownFrameworkEvents": list(session.unknown_event_types),
"durationMs": duration_ms,
},
)
return outcome
__all__ = [
"DEFAULT_DSH_BIN",
"DEFAULT_DSH_PROFILE",
"DshError",
"DshExecutionPolicy",
"DshHeadlessRunner",
"DshProcessResult",
"DshRunOutcome",
"TraceSink",
"build_dsh_argv",
"build_dsh_patch",
]