Files
kefu/wechat_rpa/dsh_agent.py
2026-08-18 17:25:22 +08:00

438 lines
15 KiB
Python
Raw Permalink 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.
# -*- coding: utf-8 -*-
"""Operator-facing DeepSeek Harness desk.
The public programming model follows the Python SDK guide:
result = harness.run(prompt, session_id=...)
On Linux/macOS, a bundled DeepSeek Harness runtime is started and reused.
Official runtime wheels are not published for Windows, so this module falls
back to the same DeepSeek-compatible chat API already configured in
``ai_config``, still keyed by ``session_id``.
"""
from __future__ import annotations
import json
import threading
import time
import uuid
from dataclasses import dataclass
from pathlib import Path
from typing import Any
from runtime_paths import application_data_dir, resource_path
BACKEND_SDK = "sdk"
BACKEND_COMPAT = "compat"
INDEX_VERSION = 1
@dataclass(slots=True)
class RunOutcome:
session_id: str
text: str
backend: str
finish_reason: str | None = None
@dataclass(slots=True)
class DeskSession:
id: str
title: str
preview: str
created_at: float
updated_at: float
messages: list[dict[str, Any]]
def desk_root(base: Path | None = None) -> Path:
return Path(base) if base is not None else application_data_dir() / "dsh_desk"
def desk_system_prompt() -> str:
try:
import ai_config
name = str(getattr(ai_config, "AI_AGENT_NAME", "") or "客服").strip() or "客服"
hospital = str(getattr(ai_config, "AI_HOSPITAL_NAME", "") or "").strip()
except Exception:
name = "客服"
hospital = ""
hospital_line = f"对外医院名称是「{hospital}」。" if hospital else ""
return (
"你是甄养堂客服工作台里的 DeepSeek Agent,协助坐席处理客户咨询、"
"起草回复和整理沟通思路。\n"
f"坐席当前对外身份是「{name}」。{hospital_line}\n"
"- 直接回答坐席的问题;需要给客户看的草稿,用「"
f"{name}」的真人口吻写,不要出现 AI、机器人、模型等字样。\n"
"- 不要编造检查结果、单号、物流、挂号或尚未发生的事实。\n"
"- 缺信息时明确说还缺什么,而不是猜。"
)
def _now() -> float:
return time.time()
def _title_from_prompt(prompt: str) -> str:
text = " ".join(str(prompt or "").split())
if not text:
return "新对话"
return text if len(text) <= 18 else text[:17] + "…"
class DeepSeekDesk:
"""Reusable desk: one harness (or compatible client) across many sessions."""
def __init__(self, root: Path | None = None) -> None:
self.root = desk_root(root)
self.workspace = self.root / "workspace"
self.session_root = self.root / "sessions"
self.transcript_dir = self.root / "transcripts"
self.index_path = self.root / "index.json"
self._lock = threading.RLock()
self._backend: str | None = None
self._backend_error = ""
self._harness: Any = None
def close(self) -> None:
with self._lock:
harness = self._harness
self._harness = None
if harness is not None:
try:
harness.close()
except Exception:
pass
def backend_name(self) -> str:
if self._backend == BACKEND_SDK:
return BACKEND_SDK
if self._backend == BACKEND_COMPAT:
return BACKEND_COMPAT
return BACKEND_SDK if _sdk_module_available() else BACKEND_COMPAT
def backend_label(self) -> str:
if self._backend == BACKEND_SDK:
return "DeepSeek Harness Python SDK"
if self._backend == BACKEND_COMPAT:
if self._backend_error:
return f"DeepSeek 兼容接口(SDK 运行时不可用:{self._backend_error}"
return "DeepSeek 兼容接口"
if _sdk_module_available():
return "DeepSeek Harness Python SDK(首次对话时启动运行时)"
return "DeepSeek 兼容接口(Windows 无官方 SDK 运行时,会话 API 与指南一致)"
def list_sessions(self) -> list[DeskSession]:
with self._lock:
index = self._load_index()
sessions = [self._read_session(item["id"]) for item in index.get("sessions") or []]
return [item for item in sessions if item is not None]
def get_session(self, session_id: str) -> DeskSession | None:
with self._lock:
return self._read_session(session_id)
def create_session(self) -> DeskSession:
with self._lock:
session = DeskSession(
id=f"desk-{uuid.uuid4().hex}",
title="新对话",
preview="",
created_at=_now(),
updated_at=_now(),
messages=[],
)
self._write_session(session)
self._upsert_index(session, current=True)
return session
def run(self, prompt: str, *, session_id: str | None = None) -> RunOutcome:
text = str(prompt or "").strip()
if not text:
raise ValueError("请先输入要交给 Agent 的任务")
with self._lock:
session = self._read_session(session_id) if session_id else None
if session is None:
if session_id:
session = DeskSession(
id=session_id,
title="新对话",
preview="",
created_at=_now(),
updated_at=_now(),
messages=[],
)
self._write_session(session)
self._upsert_index(session, current=True)
else:
session = self.create_session()
session.messages.append(
{"role": "user", "content": text, "ts": _now()}
)
if session.title == "新对话":
session.title = _title_from_prompt(text)
session.preview = text
session.updated_at = _now()
self._write_session(session)
self._upsert_index(session, current=True)
active_id = session.id
outcome = self._run_model(text, active_id)
reply = str(outcome.text or "").strip()
with self._lock:
session = self._read_session(active_id) or session
session.messages.append(
{
"role": "assistant",
"content": reply or "Agent 没有返回文本)",
"ts": _now(),
"backend": outcome.backend,
}
)
session.preview = reply or session.preview
session.updated_at = _now()
self._write_session(session)
self._upsert_index(session, current=True)
return RunOutcome(
session_id=active_id,
text=reply,
backend=outcome.backend,
finish_reason=outcome.finish_reason,
)
def current_session_id(self) -> str:
with self._lock:
index = self._load_index()
current = str(index.get("current_id") or "")
if current and self._read_session(current) is not None:
return current
sessions = index.get("sessions") or []
if sessions:
return str(sessions[0]["id"])
return self.create_session().id
def set_current_session(self, session_id: str) -> DeskSession | None:
with self._lock:
session = self._read_session(session_id)
if session is None:
return None
self._upsert_index(session, current=True)
return session
def _ensure_backend(self) -> None:
with self._lock:
if self._backend is not None:
return
try:
self._harness = self._start_harness()
self._backend = BACKEND_SDK
self._backend_error = ""
except Exception as exc:
self._harness = None
self._backend = BACKEND_COMPAT
self._backend_error = _short_error(exc)
def _start_harness(self) -> Any:
from deepseek_harness import DeepSeekHarness
self.workspace.mkdir(parents=True, exist_ok=True)
self.session_root.mkdir(parents=True, exist_ok=True)
model, api_key, base_url, timeout = _model_settings()
if not api_key:
raise RuntimeError("未配置 DeepSeek API Key")
cordis = resource_path("dsh_desk.cordis.yml")
if not cordis.is_file():
cordis = Path(__file__).resolve().parent / "dsh_desk.cordis.yml"
env = {
"DSH_SYSTEM_PROMPT": desk_system_prompt(),
"DSH_MODEL": model,
}
harness = DeepSeekHarness(
provider="deepseek-official",
model=model,
cwd=str(self.workspace.resolve()),
session_root=str(self.session_root.resolve()),
cordis=str(cordis.resolve()) if cordis.is_file() else None,
api_key=api_key,
base_url=base_url or None,
env=env,
request_timeout_seconds=timeout,
)
harness.start()
return harness
def _run_model(self, prompt: str, session_id: str) -> RunOutcome:
self._ensure_backend()
if self._backend == BACKEND_SDK and self._harness is not None:
result = self._harness.run(prompt, session_id=session_id)
return RunOutcome(
session_id=session_id,
text=str(getattr(result, "final_response", "") or ""),
backend=BACKEND_SDK,
finish_reason=getattr(result, "finish_reason", None),
)
return RunOutcome(
session_id=session_id,
text=self._run_compat(prompt, session_id),
backend=BACKEND_COMPAT,
finish_reason="completed",
)
def _run_compat(self, prompt: str, session_id: str) -> str:
import ai_chat
session = self._read_session(session_id)
history = []
if session is not None:
for item in session.messages[:-1]:
role = str(item.get("role") or "")
content = str(item.get("content") or "").strip()
if role in {"user", "assistant"} and content:
history.append({"role": role, "content": content})
try:
import ai_config
max_rounds = int(getattr(ai_config, "AI_CONTEXT_MAX_ROUNDS", 8) or 8)
except Exception:
max_rounds = 8
if max_rounds > 0:
history = history[-(max_rounds * 2) :]
messages = [{"role": "system", "content": desk_system_prompt()}]
messages.extend(history)
messages.append({"role": "user", "content": prompt})
msg = ai_chat._chat_completion(messages)
return str(msg.get("content") or "").strip()
def _load_index(self) -> dict[str, Any]:
self.root.mkdir(parents=True, exist_ok=True)
try:
data = json.loads(self.index_path.read_text(encoding="utf-8"))
except (OSError, ValueError, TypeError):
data = {}
if not isinstance(data, dict):
data = {}
data.setdefault("version", INDEX_VERSION)
data.setdefault("sessions", [])
data.setdefault("current_id", "")
return data
def _upsert_index(self, session: DeskSession, *, current: bool) -> None:
index = self._load_index()
items = [
item
for item in index.get("sessions") or []
if isinstance(item, dict) and item.get("id") != session.id
]
items.insert(
0,
{
"id": session.id,
"title": session.title,
"preview": session.preview,
"updated_at": session.updated_at,
},
)
items.sort(key=lambda item: float(item.get("updated_at") or 0), reverse=True)
index["sessions"] = items[:80]
if current:
index["current_id"] = session.id
self._atomic_write(self.index_path, index)
def _transcript_path(self, session_id: str) -> Path:
self.transcript_dir.mkdir(parents=True, exist_ok=True)
return self.transcript_dir / f"{session_id}.json"
def _read_session(self, session_id: str | None) -> DeskSession | None:
if not session_id:
return None
path = self._transcript_path(session_id)
try:
data = json.loads(path.read_text(encoding="utf-8"))
except (OSError, ValueError, TypeError):
return None
if not isinstance(data, dict):
return None
messages = data.get("messages") if isinstance(data.get("messages"), list) else []
return DeskSession(
id=str(data.get("id") or session_id),
title=str(data.get("title") or "新对话"),
preview=str(data.get("preview") or ""),
created_at=float(data.get("created_at") or _now()),
updated_at=float(data.get("updated_at") or _now()),
messages=[item for item in messages if isinstance(item, dict)],
)
def _write_session(self, session: DeskSession) -> None:
payload = {
"id": session.id,
"title": session.title,
"preview": session.preview,
"created_at": session.created_at,
"updated_at": session.updated_at,
"messages": session.messages,
}
self._atomic_write(self._transcript_path(session.id), payload)
@staticmethod
def _atomic_write(path: Path, payload: dict[str, Any]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
temporary = path.with_suffix(path.suffix + ".tmp")
temporary.write_text(
json.dumps(payload, ensure_ascii=False, indent=2),
encoding="utf-8",
)
temporary.replace(path)
_MODEL_DEFAULT = "deepseek-chat"
_desk: DeepSeekDesk | None = None
_desk_lock = threading.Lock()
def _model_settings() -> tuple[str, str, str, float]:
try:
import ai_config
model = str(getattr(ai_config, "AI_MODEL", "") or "").strip() or _MODEL_DEFAULT
api_key = str(getattr(ai_config, "AI_API_KEY", "") or "").strip()
base_url = str(getattr(ai_config, "AI_API_BASE", "") or "").strip()
timeout = float(getattr(ai_config, "AI_TIMEOUT", 120) or 120)
except Exception:
model, api_key, base_url, timeout = _MODEL_DEFAULT, "", "", 120.0
return model, api_key, base_url, max(30.0, timeout)
def _sdk_module_available() -> bool:
try:
import deepseek_harness # noqa: F401
except Exception:
return False
return True
def _short_error(exc: BaseException) -> str:
text = str(exc or exc.__class__.__name__).strip().replace("\n", " ")
if "No matching distribution" in text or "deepseek-harness-runtime" in text:
return "当前系统没有官方 Windows 运行时"
if "FileNotFoundError" in type(exc).__name__ or "Unable to locate" in text:
return "未找到 SDK 运行时"
return text[:80] or exc.__class__.__name__
def get_desk() -> DeepSeekDesk:
global _desk
with _desk_lock:
if _desk is None:
_desk = DeepSeekDesk()
return _desk
def close_desk() -> None:
global _desk
with _desk_lock:
desk = _desk
_desk = None
if desk is not None:
desk.close()