From 6b0a51fbe88c12323c1dbd2539d04c9f3bfb68bd Mon Sep 17 00:00:00 2001 From: DeepCode Date: Fri, 31 Jul 2026 13:55:41 +0800 Subject: [PATCH] feat: Skills ecosystem - Hooks/Knowledge/Streaming (clean version) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 重建 PR #142:仅包含 3 个有实际代码的 skill,排除运行时数据文件: - deepcode-hooks: 通用 Pre/Post 钩子框架 (hooks.py) - deepcode-knowledge: Obsidian 风格笔记 + KV 记忆系统 (knowledge_server.py + memory_manager.py + rebuild_vault_index.py) - deepcode-streaming: SSE 流式 API 客户端 (streaming_api.py) 不含: .db SQLite 数据文件 / .jsonl 日志 / 空文件 / 嵌套重复目录 --- .deepcode/skills/deepcode-hooks/SKILL.md | 220 ++++ .deepcode/skills/deepcode-hooks/hooks.py | 968 ++++++++++++++++++ .deepcode/skills/deepcode-hooks/plugin.json | 11 + .deepcode/skills/deepcode-knowledge/SKILL.md | 89 ++ .../deepcode-knowledge/knowledge_server.py | 539 ++++++++++ .../deepcode-knowledge/memory_manager.py | 643 ++++++++++++ .../skills/deepcode-knowledge/plugin.json | 13 + .../deepcode-knowledge/rebuild_vault_index.py | 207 ++++ .deepcode/skills/deepcode-streaming/SKILL.md | 91 ++ .../skills/deepcode-streaming/plugin.json | 11 + .../deepcode-streaming/streaming_api.py | 514 ++++++++++ 11 files changed, 3306 insertions(+) create mode 100644 .deepcode/skills/deepcode-hooks/SKILL.md create mode 100644 .deepcode/skills/deepcode-hooks/hooks.py create mode 100644 .deepcode/skills/deepcode-hooks/plugin.json create mode 100644 .deepcode/skills/deepcode-knowledge/SKILL.md create mode 100644 .deepcode/skills/deepcode-knowledge/knowledge_server.py create mode 100644 .deepcode/skills/deepcode-knowledge/memory_manager.py create mode 100644 .deepcode/skills/deepcode-knowledge/plugin.json create mode 100644 .deepcode/skills/deepcode-knowledge/rebuild_vault_index.py create mode 100644 .deepcode/skills/deepcode-streaming/SKILL.md create mode 100644 .deepcode/skills/deepcode-streaming/plugin.json create mode 100644 .deepcode/skills/deepcode-streaming/streaming_api.py diff --git a/.deepcode/skills/deepcode-hooks/SKILL.md b/.deepcode/skills/deepcode-hooks/SKILL.md new file mode 100644 index 00000000..9b7691f4 --- /dev/null +++ b/.deepcode/skills/deepcode-hooks/SKILL.md @@ -0,0 +1,220 @@ +--- +name: deepcode-hooks +description: > + DeepCode Hook System — 移植自 Claude Code 的 Hook 架构. + 支持 Pre/Post 工具钩子、会话钩子、错误钩子、路由钩子。 + 对标 Claude Code: hook-handler.cjs, pre-bash, post-edit, + session-restore, session-end, pre-task, post-task, route. +version: 1.0.0 +author: DeepCode + Ghidra RE (Claude Code v2.1.216) +date: 2026-07-26 +tags: [hooks, pipeline, automation, security] +--- + +# DeepCode Hook System + +移植自 **Claude Code v2.1.216** 的完整 Hook 架构 — 通过 Ghidra + v8asm 从 V8 字节码中逆向提取。 + +## 逆向成果 + +通过 `v8asm` 反汇编 + `bun` 段字符串提取,还原了 Claude Code 完整的 Hook 系统设计: + +| 来源 | 工具 | 发现 | +|:----|:----|:----| +| `.bun` V8 字节码 | `v8asm disassembler` | 13 种 Hook 事件 | +| V8 heap snapshot | `extract_bun_hooks.py` | 4 种 Hook 类型 (command/prompt/agent/mcp_tool) | +| PE 字符串提取 | Ghidra + 二进制搜索 | Hook 输入/输出 JSON Schema | + +## Claude Code Hook 事件 (完整列表) + +| # | 事件 | 时机 | 用途 | +|:-|:----|:----|:-----| +| 1 | **PreToolUse** | 工具执行前 | 可阻塞工具调用—对标 `beforeCommand`/`beforeWrite` | +| 2 | **PostToolUse** | 工具执行成功 | 记录结果—对标 `afterCommand`/`afterWrite` | +| 3 | **PostToolUseFailure** | 工具执行失败 | 错误处理—对标 `onError` | +| 4 | **PermissionRequest** | 权限检查 | 权限决策注入 | +| 5 | **Notification** | 通知事件 | 通知类型过滤 | +| 6 | **Stop** | 中止请求 | 安全停止—对标 `onCancel` | +| 7 | **UserPromptSubmit** | 用户提交 | 提示注入检测 | +| 8 | **SessionStart** | 会话开始 | 状态恢复—对标 `sessionStart` | +| 9 | **Setup** | 插件安装 | 插件初始化 | +| 10 | **UserPromptExpansion** | 提示展开 | 用户输入预处理 | +| 11 | **SubagentStop** | 子 Agent 停止 | Subagent 生命周期 | +| 12 | **WorktreeCreate** | Git Worktree 创建 | Git 隔离工作区 | +| 13 | **WorktreeRemove** | Git Worktree 删除 | 清理工作区 | +| 14 | **PreCompact** | 压缩/紧凑前 | 状态持久化保护 | + +## Hook 类型 + +| 类型 | 说明 | JSON type 值 | +|:----|:----|:------------| +| **Command Hook** | 执行 Shell 命令: `{ "type": "command", "command": "prettier --write $FILE" }` | `command` | +| **Prompt Hook** | LLM 评估条件: `{ "type": "prompt", "prompt": "Is this safe?" }` | `prompt` | +| **Agent Hook** | 运行 Agent 工具链: `{ "type": "agent", "tools": ["read", "write"] }` | `agent` | +| **MCP Tool Hook** | 调用 MCP 工具: `{ "type": "mcp_tool", "mcp_tool": "tool_name" }` | `mcp_tool` | + +## Hook 协议 (stdin/stdout JSON) + +### Hook 输入 (stdin JSON) + +```json +{ + "session_id": "abc123", + "tool_name": "Bash", + "tool_input": { "command": "ls -la" }, + "tool_response": { "exit_code": 0 } // PostToolUse 才有 +} +``` + +### Hook 输出 (stdout JSON) + +```json +{ + "systemMessage": "显示给用户的消息", + "continue": true, + "decision": "block", // PreToolUse: 阻塞工具 + "reason": "安全策略阻止", + "hookSpecificOutput": { + "additionalContext": "额外上下文", + "permissionDecision": "allowed" + } +} +``` + +## 对标项 + +| Claude Code | DeepCode Hooks | +|:-----------|:---------------| +| `PreToolUse` | `beforeCommand` / `beforeWrite` / `beforeEdit` | +| `PostToolUse` | `afterCommand` / `afterWrite` / `afterEdit` | +| `PostToolUseFailure` | `onError` | +| `PermissionRequest` | (权限门控集成) | +| `SessionStart` | `sessionStart` | +| `Setup` | `startup` | +| `Stop` | (停止事件) | +| `UserPromptSubmit` | (提示注入) | +| `pre-bash` / `post-edit` / `session-restore` | 通用事件系统覆盖 | +| `hook-handler.cjs` | `hooks.py` | +| hook 异常检测 | `exit_code` + `error_message` | + +## 支持的事件 + +| 事件 | 对应 Claude Code | 触发时机 | +|:----|:----------------|:--------| +| **beforeCommand** | PreToolUse (Bash) | 命令执行前 | +| **afterCommand** | PostToolUse (Bash) | 命令执行后 | +| **beforeWrite** | PreToolUse (Write) | 文件写入前 | +| **afterWrite** | PostToolUse (Write) | 文件写入后 | +| **beforeEdit** | PreToolUse (Edit) | 文件编辑前 | +| **afterEdit** | PostToolUse (Edit) | 文件编辑后 | +| **beforeRead** | PreToolUse (Read) | 文件读取前 | +| **onError** | PostToolUseFailure | 发生错误 | +| **preTask** | (任务级) | 任务开始 | +| **postTask** | (任务级) | 任务结束 | +| **sessionStart** | SessionStart | 会话开始 | +| **sessionEnd** | SessionEnd | 会话结束 | +| **route** | (路由决策) | 路由决策 | +| **startup** | Setup | 系统启动 | +| **shutdown** | — | 系统关闭 | + +## 用法 + +### CLI + +```bash +# 注册 — 文件写入前自动 lint +python hooks.py register \ + --name pre-lint \ + --event beforeWrite \ + --handler "npx eslint {filePath}" \ + --type shell \ + --priority 10 + +# 注册 — 命令执行前安全校验 +python hooks.py register \ + --name safety-check \ + --event beforeCommand \ + --handler "python safety_check.py" \ + --type python \ + --priority 100 + +# 列出 +python hooks.py list + +# 触发 +python hooks.py trigger --event beforeWrite \ + --ctx '{"filePath":"/path/to/file.py","toolName":"Write"}' +``` + +### settings.json 集成 + +兼容现有配置格式,新增 Hook 配置节: + +```json +{ + "hooks": { + "beforeWrite": "echo [hook] 即将修改: {filePath}", + "afterEdit": "git add {filePath}", + "beforeCommand": "python safety_check.py", + "onError": "python notify_error.py", + "sessionStart": "python restore_state.py", + "sessionEnd": "python persist_state.py" + } +} +``` + +模板变量: `{filePath}`, `{command}`, `{exitCode}`, `{error}`, `{timestamp}` + +### MCP Server + +```json +"deepcode-hooks": { + "command": "python", + "args": [ + "F:/DEEPCODE/.deepcode/skills/deepcode-hooks/hooks.py", + "--mcp" + ] +} +``` + +### Python 嵌入 + +```python +from hooks import HookManager, Hook, HookContext, HookEvent + +mgr = HookManager() + +# 注册 +mgr.register(Hook( + name="auto-commit", + event=HookEvent.AFTER_EDIT, + handler="git add {filePath} && git commit -m 'auto: {filePath}'", + type="shell", + priority=50, +)) + +# 触发 +ctx = HookContext( + event=HookEvent.AFTER_EDIT, + tool_name="Edit", + file_path="/path/to/file.py", + exit_code=0, +) +results = await mgr.trigger(HookEvent.AFTER_EDIT, ctx) +for r in results: + print(f"{r['hook']}: {r['status']} ({r['duration']}s)") +``` + +## Hook 类型 + +| 类型 | 执行方式 | 适用场景 | +|:----|:--------|:--------| +| `shell` | `subprocess` | 简单命令、git 操作 | +| `python` | `subprocess` | 复杂逻辑校验 | +| `node` | `subprocess` | npm 生态工具 | + +## 安全 + +- 每个 Hook 有独立超时保护(默认 30s) +- 异常不会级联 — 一个 Hook 失败不影响其他 Hook +- Hook 通过环境变量接收上下文(`DEEPCODE_*` 系列变量) diff --git a/.deepcode/skills/deepcode-hooks/hooks.py b/.deepcode/skills/deepcode-hooks/hooks.py new file mode 100644 index 00000000..c3bb2b40 --- /dev/null +++ b/.deepcode/skills/deepcode-hooks/hooks.py @@ -0,0 +1,968 @@ +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- +""" +DeepCode Hook System — 移植自 Claude Code 的 Hook 架构 +══════════════════════════════════════════════════════════ +通用 Hook 管理器,支持 Pre/Post 工具钩子、会话钩子、错误钩子。 + +对标 Claude Code hook-handler.cjs: + - pre-bash / post-edit / pre-task / post-task + - session-restore / session-end + - route / stats + - toolFailed 检测 / intelligence 反馈 + +用法: + # 注册 Hook + python hooks.py register --event beforeWrite --cmd "python validate.py {filePath}" + + # 触发 Hook + python hooks.py trigger --event beforeWrite --ctx '{"filePath":"/path/to/file"}' + + # 列出 Hook + python hooks.py list + + # MCP Server + python hooks.py --mcp +""" + +import asyncio +import json +import os +import platform +import subprocess +import sys +import tempfile +import time +import uuid +from dataclasses import dataclass, field +from datetime import datetime +from pathlib import Path +from typing import Dict, List, Optional, Any, Callable +from enum import Enum + +SYSTEM = platform.system() + +# ── Codex 决策模型 ────────────────────────────────────── + +class HookDecision(str, Enum): + """ + Hook 决策类型 — 移植自 codex.exe Hook 系统 + + 参考 codex.exe 的 DecisionWire 类型: + - PreToolUseDecisionWire: approve / block / allow / deny + - PermissionRequestBehaviorWire: behavior / updatedInput / updatedPermissions / interrupt + - PreToolUsePermissionDecisionWire: additionalContext / permissionDecision + """ + APPROVE = "approve" # 允许继续 (PreToolUse) + BLOCK = "block" # 阻止执行 (PreToolUse) + ALLOW = "allow" # 允许 (PermissionRequest) + DENY = "deny" # 拒绝 (PermissionRequest) + INTERRUPT = "interrupt" # 中断 (PermissionRequest) + CONTINUE = "continue" # 继续 (PostToolUse/通用) + STOP = "stop" # 停止 (Stop hook) + SUPPRESS = "suppress" # 抑制输出 + +class HookBehavior(str, Enum): + """Hook 行为 — 移植自 codex.exe PermissionRequestBehaviorWire""" + DEFAULT = "default" + UPDATED_INPUT = "updatedInput" + UPDATED_PERMISSIONS = "updatedPermissions" + INTERRUPT = "interrupt" + ASK = "ask" + +class HookWire: + """ + Hook Wire 类型 — 移植自 codex.exe 的 XXXHookSpecificOutputWire 类型 + + codex.exe 中每个事件有专属的 Wire 类型,例如: + - PreToolUseHookSpecificOutputWire + - PostToolUseHookSpecificOutputWire + - PermissionRequestHookSpecificOutputWire + - SessionStartHookSpecificOutputWire + """ + + @staticmethod + def pre_tool_use(decision: HookDecision = HookDecision.APPROVE, + reason: str = "", + additional_context: str = "", + permission_decision: str = "") -> Dict: + """PreToolUseDecisionWire — 工具执行前决策""" + return { + "decision": decision.value, + "reason": reason, + "hookSpecificOutput": { + "additionalContext": additional_context, + "permissionDecision": permission_decision, + }, + } + + @staticmethod + def permission_request(behavior: HookBehavior = HookBehavior.DEFAULT, + decision: HookDecision = HookDecision.ALLOW, + reason: str = "") -> Dict: + """PermissionRequestBehaviorWire — 权限请求""" + return { + "behavior": behavior.value, + "decision": decision.value, + "reason": reason, + } + + @staticmethod + def post_tool_use(continue_flag: bool = True, + system_message: str = "", + additional_context: str = "") -> Dict: + """PostToolUseHookSpecificOutputWire — 工具执行后""" + return { + "continue": continue_flag, + "systemMessage": system_message, + "hookSpecificOutput": { + "additionalContext": additional_context, + }, + } + + @staticmethod + def session_start(system_message: str = "", + additional_context: str = "") -> Dict: + """SessionStartHookSpecificOutputWire — 会话启动""" + return { + "systemMessage": system_message, + "hookSpecificOutput": { + "additionalContext": additional_context, + }, + } + + @staticmethod + def stop(continue_flag: bool = False, + reason: str = "", + suppress_output: bool = False) -> Dict: + """StopCommandOutputWire — 停止""" + return { + "continue": continue_flag, + "reason": reason, + "suppressOutput": suppress_output, + } + + +# ── 事件定义 (扩展) ────────────────────────────────────── + +class HookEvent: + """Hook 事件类型 — 对标 Claude Code 全部 hook 点""" + + # 工具执行前 (Pre-Tool) + BEFORE_WRITE = "beforeWrite" # 文件写入前 + BEFORE_EDIT = "beforeEdit" # 文件编辑前 + BEFORE_COMMAND = "beforeCommand" # 命令执行前 + BEFORE_READ = "beforeRead" # 文件读取前 + BEFORE_SEARCH = "beforeSearch" # 搜索前 + + # 工具执行后 (Post-Tool) + AFTER_WRITE = "afterWrite" # 文件写入后 + AFTER_EDIT = "afterEdit" # 文件编辑后 — 对标 post-edit + AFTER_COMMAND = "afterCommand" # 命令执行后 + AFTER_READ = "afterRead" # 文件读取后 + AFTER_SEARCH = "afterSearch" # 搜索后 + + # 任务 (Task) + PRE_TASK = "preTask" # 任务开始 — 对标 pre-task + POST_TASK = "postTask" # 任务结束 — 对标 post-task + + # 会话 (Session) + SESSION_START = "sessionStart" # 会话开始 — 对标 session-restore + SESSION_END = "sessionEnd" # 会话结束 — 对标 session-end + + # 路由 (Route) + ROUTE = "route" # 路由决策前 — 对标 route + + # 错误 + ON_ERROR = "onError" # 发生错误 + + # 系统 + STARTUP = "startup" # DeepCode 启动 + SHUTDOWN = "shutdown" # DeepCode 关闭 + + # 高级别抽象事件(settings.json 兼容)— 运行时由 matcher 分发到具体工具 + PRE_TOOL_USE = "PreToolUse" # 任意工具执行前 + POST_TOOL_USE = "PostToolUse" # 任意工具执行后 + POST_TOOL_USE_FAILURE = "PostToolUseFailure" # 工具执行失败 + PERMISSION_REQUEST = "PermissionRequest" # 权限请求 + NOTIFICATION = "Notification" # 通知 + STOP_EVENT = "Stop" # 中止请求 + USER_PROMPT_SUBMIT = "UserPromptSubmit" # 用户提交提示 + USER_PROMPT_EXPANSION = "UserPromptExpansion" # 提示展开 + SUBAGENT_STOP = "SubagentStop" # 子Agent停止 + WORKTREE_CREATE = "WorktreeCreate" # Git Worktree 创建 + WORKTREE_REMOVE = "WorktreeRemove" # Git Worktree 删除 + PRE_COMPACT = "PreCompact" # 压缩前 + + @classmethod + def all(cls) -> List[str]: + return [v for k, v in vars(cls).items() + if not k.startswith("_") and isinstance(v, str)] + + @classmethod + def pre_events(cls) -> List[str]: + """所有 Pre 事件""" + return [e for e in cls.all() if e.startswith("before")] + + @classmethod + def post_events(cls) -> List[str]: + """所有 Post 事件""" + return [e for e in cls.all() if e.startswith("after")] + + +# ── Hook 上下文 ─────────────────────────────────────────── + +@dataclass +class HookContext: + """Hook 上下文 — 移植自 codex.exe Hook 系统 + + codex.exe Wire 类型字段: + - decision: approve/block/allow/deny (PreToolUseDecisionWire) + - permission_decision: 权限决策 (PreToolUsePermissionDecisionWire) + - behavior: default/updatedInput/updatedPermissions/interrupt + - additional_context: 额外上下文 (XXXHookSpecificOutputWire) + - suppress_output: 抑制输出 (StopCommandOutputWire) + """ + event: str + tool_name: str = "" + file_path: str = "" + file_content: str = "" + old_string: str = "" + new_string: str = "" + command: str = "" + command_stdout: str = "" + command_stderr: str = "" + exit_code: int = 0 + error_message: str = "" + + # Codex 决策模型字段 + decision: str = "" # approve/block/allow/deny + reason: str = "" # 决策理由 + additional_context: str = "" # 额外上下文 (hookSpecificOutput) + permission_decision: str = "" # 权限决策 + behavior: str = "" # 行为模式 + suppress_output: bool = False # 抑制输出 + + session_id: str = "" + task_id: str = "" + project_root: str = "" + workspace: str = "" + timestamp: str = "" + agent_id: str = "" + tool_use_id: str = "" + hook_event_name: str = "" + model: str = "" + permission_mode: str = "" + transcript_path: str = "" + extra: Dict[str, Any] = field(default_factory=dict) + + def to_wire(self) -> Dict: + """转换为 Wire 格式 — 移植自 codex.exe HookSpecificOutputWire""" + wire = { + "session_id": self.session_id, + "turn_id": self.task_id, + "agent_id": self.agent_id, + "tool_name": self.tool_name, + "tool_input": {}, + "tool_use_id": self.tool_use_id, + "hook_event_name": self.hook_event_name or self.event, + "model": self.model, + "permission_mode": self.permission_mode, + "trigger": "hook", + "transcript_path": self.transcript_path, + } + if self.command: + wire["tool_input"]["command"] = self.command + if self.file_path: + wire["tool_input"]["path"] = self.file_path + return wire + + def to_env(self) -> Dict[str, str]: + """转换为环境变量 (子进程可用)""" + env = { + "DEEPCODE_HOOK_EVENT": self.event, + "DEEPCODE_TOOL_NAME": self.tool_name, + "DEEPCODE_FILE_PATH": self.file_path, + "DEEPCODE_COMMAND": self.command, + "DEEPCODE_EXIT_CODE": str(self.exit_code), + "DEEPCODE_ERROR": self.error_message, + "DEEPCODE_SESSION_ID": self.session_id, + "DEEPCODE_TASK_ID": self.task_id, + "DEEPCODE_PROJECT_ROOT": self.project_root or os.getcwd(), + "DEEPCODE_WORKSPACE": self.workspace or os.getcwd(), + "DEEPCODE_TIMESTAMP": self.timestamp or datetime.now().isoformat(), + } + # 非空字段才设置 (避免污染) + return {k: v for k, v in env.items() if v} + + def to_dict(self) -> Dict[str, Any]: + return { + "event": self.event, + "tool_name": self.tool_name, + "file_path": self.file_path, + "command": self.command, + "exit_code": self.exit_code, + "error": self.error_message, + "session_id": self.session_id, + "task_id": self.task_id, + "project_root": self.project_root, + "timestamp": self.timestamp or datetime.now().isoformat(), + **self.extra, + } + + +# ── Hook 定义 ───────────────────────────────────────────── + +@dataclass +class Hook: + """单个 Hook 定义""" + name: str + event: str + handler: str # 脚本路径或内联命令 + type: str = "shell" # shell / python / node + priority: int = 0 # 越大越先执行 + timeout: int = 30 # 超时秒数 + enabled: bool = True + description: str = "" + created_at: str = field(default_factory=lambda: datetime.now().isoformat()) + + def expand_vars(self, text: str, ctx: HookContext) -> str: + """展开模板变量""" + import string + try: + return string.Formatter().vformat(text, (), { + "event": ctx.event, + "toolName": ctx.tool_name, + "filePath": ctx.file_path, + "command": ctx.command, + "exitCode": str(ctx.exit_code), + "error": ctx.error_message, + "sessionId": ctx.session_id, + "taskId": ctx.task_id, + "projectRoot": ctx.project_root or os.getcwd(), + "workspace": ctx.workspace or os.getcwd(), + "timestamp": ctx.timestamp or datetime.now().isoformat(), + }) + except Exception: + return text + + +# ── Hook 管理器 ─────────────────────────────────────────── + +class HookManager: + """ + Hook 管理器 — 对标 Claude Code hook-handler.cjs + + - 支持 Pre/Post/Error/Session 事件 + - 多脚本链式执行 + - 超时 + 异常保护 + - 变量模板展开 + - 配置持久化 + """ + + def __init__(self, config_path: Optional[str] = None): + self.config_path = config_path or os.path.expanduser( + "~/.deepcode/hooks_config.json" + ) + self._hooks: List[Hook] = [] + self._load() + + def _load(self): + """从配置加载 hooks""" + p = Path(self.config_path) + if p.exists(): + try: + data = json.loads(p.read_text(encoding="utf-8")) + self._hooks = [Hook(**h) for h in data.get("hooks", [])] + except Exception: + self._hooks = [] + + def _save(self): + """持久化 hooks""" + Path(self.config_path).parent.mkdir(parents=True, exist_ok=True) + Path(self.config_path).write_text( + json.dumps({ + "hooks": [h.__dict__ for h in self._hooks], + "updated_at": datetime.now().isoformat(), + }, indent=2, ensure_ascii=False), + encoding="utf-8", + ) + + def register(self, hook: Hook): + """注册一个 Hook""" + # 同名替换 + self._hooks = [h for h in self._hooks + if not (h.name == hook.name and h.event == hook.event)] + self._hooks.append(hook) + self._hooks.sort(key=lambda h: (-h.priority, h.name)) + self._save() + + def unregister(self, name: str, event: str = None): + """注销一个 Hook""" + if event: + self._hooks = [h for h in self._hooks + if not (h.name == name and h.event == event)] + else: + self._hooks = [h for h in self._hooks if h.name != name] + self._save() + + def get_hooks(self, event: str) -> List[Hook]: + """获取某事件的所有 Hook (按优先级排序)""" + return [h for h in self._hooks + if h.event == event and h.enabled] + + def list_hooks(self, event: str = None) -> List[Dict]: + """列出 Hook""" + hooks = self._hooks + if event: + hooks = [h for h in hooks if h.event == event] + return [ + { + "name": h.name, + "event": h.event, + "type": h.type, + "handler": h.handler, + "priority": h.priority, + "enabled": h.enabled, + "timeout": h.timeout, + "description": h.description, + } + for h in hooks + ] + + def get_hooks_by_group(self) -> Dict[str, List[Dict]]: + """按事件分组列出""" + groups = {} + for h in self._hooks: + groups.setdefault(h.event, []).append({ + "name": h.name, + "type": h.type, + "handler": h.handler, + "priority": h.priority, + "enabled": h.enabled, + }) + return groups + + def enable(self, name: str, event: str = None): + """启用 Hook""" + for h in self._hooks: + if h.name == name and (not event or h.event == event): + h.enabled = True + self._save() + + def disable(self, name: str, event: str = None): + """禁用 Hook""" + for h in self._hooks: + if h.name == name and (not event or h.event == event): + h.enabled = False + self._save() + + # ── Hook 执行 ────────────────────────────────────── + + async def trigger( + self, + event: str, + ctx: Optional[HookContext] = None, + timeout: int = 30, + ) -> List[Dict]: + """ + 触发事件的所有 Hook — 对标 Claude Code hook dispatch + + Args: + event: 事件名 + ctx: Hook 上下文 + timeout: 每个 hook 的超时 + + Returns: + [{"hook": name, "status": "ok"|"error"|"timeout", "output": str, ...}] + """ + hooks = self.get_hooks(event) + if not hooks: + return [] + + ctx = ctx or HookContext(event=event) + ctx.timestamp = datetime.now().isoformat() + results = [] + + for hook in hooks: + result = await self._run_hook(hook, ctx, timeout) + results.append(result) + + return results + + async def _run_hook( + self, hook: Hook, ctx: HookContext, default_timeout: int + ) -> Dict: + """执行单个 Hook""" + start = time.time() + try: + # 展开模板变量 + handler = hook.expand_vars(hook.handler, ctx) + actual_timeout = hook.timeout or default_timeout + + if hook.type == "shell": + output = await self._run_shell(handler, ctx, actual_timeout) + elif hook.type == "python": + output = await self._run_python(handler, ctx, actual_timeout) + elif hook.type == "node": + output = await self._run_node(handler, ctx, actual_timeout) + else: + output = f"Unknown hook type: {hook.type}" + + return { + "hook": hook.name, + "event": hook.event, + "status": "ok", + "output": output[:2000] if output else "", + "duration": round(time.time() - start, 3), + } + except asyncio.TimeoutError: + return { + "hook": hook.name, + "event": hook.event, + "status": "timeout", + "output": f"Timed out after {hook.timeout or default_timeout}s", + "duration": hook.timeout or default_timeout, + } + except Exception as e: + return { + "hook": hook.name, + "event": hook.event, + "status": "error", + "output": str(e)[:500], + "duration": round(time.time() - start, 3), + } + + async def _run_shell( + self, cmd: str, ctx: HookContext, timeout: int + ) -> str: + """执行 Shell 命令 hook""" + env = {**os.environ, **ctx.to_env()} + proc = await asyncio.wait_for( + asyncio.create_subprocess_shell( + cmd, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + env=env, + cwd=ctx.project_root or os.getcwd(), + ), + timeout=timeout, + ) + stdout, stderr = await asyncio.wait_for( + proc.communicate(), timeout=timeout + ) + output = (stdout or b"").decode("utf-8", errors="replace") + if stderr: + output += "\n[stderr] " + stderr.decode("utf-8", errors="replace") + return output.strip() + + async def _run_python( + self, script: str, ctx: HookContext, timeout: int + ) -> str: + """执行 Python 脚本 hook""" + if os.path.isfile(script): + cmd = [sys.executable, script] + else: + # 内联 Python 代码 + tmp = tempfile.NamedTemporaryFile( + mode="w", suffix=".py", delete=False, encoding="utf-8" + ) + tmp.write( + "import os, json\n" + f"# Hook: {ctx.event}\n" + f"print('executing hook...')\n" + ) + if "print" not in script: + tmp.write(f"print({repr(script)})\n") + else: + tmp.write(script + "\n") + tmp.close() + cmd = [sys.executable, tmp.name] + # 清理 + try: + os.unlink(tmp.name) + except Exception: + pass + + env = {**os.environ, **ctx.to_env()} + proc = await asyncio.wait_for( + asyncio.create_subprocess_exec( + *cmd, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + env=env, + cwd=ctx.project_root or os.getcwd(), + ), + timeout=timeout, + ) + stdout, stderr = await asyncio.wait_for( + proc.communicate(), timeout=timeout + ) + output = (stdout or b"").decode("utf-8", errors="replace") + if stderr: + output += "\n[stderr] " + stderr.decode("utf-8", errors="replace") + return output.strip() + + async def _run_node( + self, script: str, ctx: HookContext, timeout: int + ) -> str: + """执行 Node.js 脚本 hook""" + cmd = ["node", script] if os.path.isfile(script) else ["node", "-e", script] + env = {**os.environ, **ctx.to_env()} + proc = await asyncio.wait_for( + asyncio.create_subprocess_exec( + *cmd, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + env=env, + ), + timeout=timeout, + ) + stdout, stderr = await asyncio.wait_for( + proc.communicate(), timeout=timeout + ) + output = (stdout or b"").decode("utf-8", errors="replace") + if stderr: + output += "\n[stderr] " + stderr.decode("utf-8", errors="replace") + return output.strip() + + # ── 与 settings.json hooks 兼容 ───────────────────── + + @classmethod + def from_settings_dict(cls, hooks_dict: Dict[str, str]) -> "HookManager": + """从 settings.json 的 hooks 配置创建""" + mgr = cls() + event_map = { + "beforeWrite": HookEvent.BEFORE_WRITE, + "afterWrite": HookEvent.AFTER_WRITE, + "beforeEdit": HookEvent.BEFORE_EDIT, + "afterEdit": HookEvent.AFTER_EDIT, + "beforeCommand": HookEvent.BEFORE_COMMAND, + "afterCommand": HookEvent.AFTER_COMMAND, + "onError": HookEvent.ON_ERROR, + } + for key, handler in hooks_dict.items(): + event = event_map.get(key, key) + mgr.register(Hook( + name=f"settings_{key}", + event=event, + handler=handler, + type="shell", + priority=100, # settings hooks 优先 + )) + return mgr + + +# ── MCP Server 模式 ────────────────────────────────────── + +def _mcp_send(obj): + """JSON-RPC 单行输出(MCP 标准 stdio 传输协议要求单行 JSON)""" + print(json.dumps(obj, ensure_ascii=False), flush=True) + + +def _load_hooks_from_settings(mgr: HookManager) -> int: + """从 settings.json 自动加载 hooks 配置。返回加载数量。""" + candidates = [ + os.path.join(os.getcwd(), ".deepcode", "settings.json"), + os.path.join(os.getcwd(), ".deepcode", "settings.local.json"), + os.path.expanduser("~/.deepcode/settings.json"), + ] + loaded = 0 + for cfg_path in candidates: + if not os.path.isfile(cfg_path): + continue + try: + with open(cfg_path, "r", encoding="utf-8") as fh: + cfg = json.load(fh) + hooks_cfg = cfg.get("hooks", {}) + if not isinstance(hooks_cfg, dict): + continue + for event_name, entries in hooks_cfg.items(): + if not isinstance(entries, list): + continue + # 验证事件名合法 + if event_name not in HookEvent.all(): + continue + for i, entry in enumerate(entries): + if not isinstance(entry, dict): + continue + cmd = entry.get("command", entry.get("handler", "")) + if not cmd: + continue + hook_name = entry.get("name") or f"auto:{event_name}:{i}" + kind = entry.get("type", "shell") + if kind == "command": + kind = "shell" + mgr.register(Hook( + name=hook_name, + event=event_name, + handler=cmd, + type=kind, + priority=entry.get("priority", 0), + timeout=entry.get("timeout", 30), + description=entry.get("description", entry.get("matcher", "")), + )) + loaded += 1 + except Exception: + pass + return loaded + + +async def run_mcp(): + """作为 MCP Server 运行(标准 MCP 协议)""" + mgr = HookManager() + n_loaded = _load_hooks_from_settings(mgr) + if n_loaded: + # 写 stderr 避免污染 MCP stdio + sys.stderr.write(f"[deepcode-hooks] auto-loaded {n_loaded} hooks from settings\n") + sys.stderr.flush() + + for line in sys.stdin: + line = line.strip() + if not line: + continue + try: + req = json.loads(line) + method = req.get("method", "") + params = req.get("params", {}) + rid = req.get("id", "") + + # ── 标准 MCP 握手:initialize ────────────────── + if method == "initialize": + _mcp_send({ + "jsonrpc": "2.0", "id": rid, + "result": { + "protocolVersion": "2024-11-05", + "capabilities": {"tools": {"listChanged": True}}, + "serverInfo": { + "name": "deepcode-hooks", + "version": "1.0.0", + }, + }, + }) + + # ── notifications/initialized:静默忽略 ────── + elif method == "notifications/initialized": + pass + + # ── tools/list ────────────────────────────── + elif method == "tools/list": + _mcp_send({ + "jsonrpc": "2.0", "id": rid, + "result": { + "tools": [ + { + "name": "hook_list", + "description": "列出所有 Hook", + "inputSchema": { + "type": "object", + "properties": { + "event": {"type": "string"}, + }, + }, + }, + { + "name": "hook_register", + "description": "注册新 Hook", + "inputSchema": { + "type": "object", + "properties": { + "name": {"type": "string"}, + "event": {"type": "string", "enum": HookEvent.all()}, + "handler": {"type": "string"}, + "type": {"type": "string", "enum": ["shell", "python", "node"]}, + "priority": {"type": "integer"}, + "timeout": {"type": "integer"}, + "description": {"type": "string"}, + }, + "required": ["name", "event", "handler"], + }, + }, + { + "name": "hook_trigger", + "description": "触发事件 Hook", + "inputSchema": { + "type": "object", + "properties": { + "event": {"type": "string"}, + "ctx": {"type": "object"}, + }, + "required": ["event"], + }, + }, + { + "name": "hook_unregister", + "description": "注销 Hook", + "inputSchema": { + "type": "object", + "properties": { + "name": {"type": "string"}, + "event": {"type": "string"}, + }, + "required": ["name"], + }, + }, + { + "name": "hook_enable", + "description": "启用 Hook", + "inputSchema": { + "type": "object", + "properties": { + "name": {"type": "string"}, + "event": {"type": "string"}, + }, + "required": ["name"], + }, + }, + { + "name": "hook_disable", + "description": "禁用 Hook", + "inputSchema": { + "type": "object", + "properties": { + "name": {"type": "string"}, + "event": {"type": "string"}, + }, + "required": ["name"], + }, + }, + ], + }, + }) + + # ── tools/call ───────────────────────────── + elif method == "tools/call": + name = params.get("name", "") + args = params.get("arguments", {}) + result = {} + + if name == "hook_list": + result["hooks"] = mgr.list_hooks(args.get("event")) + elif name == "hook_register": + mgr.register(Hook( + name=args["name"], + event=args["event"], + handler=args["handler"], + type=args.get("type", "shell"), + priority=args.get("priority", 0), + timeout=args.get("timeout", 30), + description=args.get("description", ""), + )) + result["status"] = "registered" + elif name == "hook_trigger": + ctx = HookContext( + event=args["event"], + **(args.get("ctx", {})) + ) + results = await mgr.trigger(args["event"], ctx) + result["results"] = results + elif name == "hook_unregister": + mgr.unregister(args["name"], args.get("event")) + result["status"] = "unregistered" + elif name == "hook_enable": + mgr.enable(args["name"], args.get("event")) + result["status"] = "enabled" + elif name == "hook_disable": + mgr.disable(args["name"], args.get("event")) + result["status"] = "disabled" + + _mcp_send({ + "jsonrpc": "2.0", "id": rid, + "result": { + "content": [ + {"type": "text", "text": json.dumps(result, ensure_ascii=False)} + ], + }, + }) + + except json.JSONDecodeError: + pass + + +# ── CLI 入口 ────────────────────────────────────────────── + +def main(): + import argparse + parser = argparse.ArgumentParser(description="DeepCode Hook System") + parser.add_argument("--mcp", action="store_true", help="MCP Server 模式") + sub = parser.add_subparsers(dest="mode") + + # list + list_p = sub.add_parser("list", help="列出 Hook") + list_p.add_argument("--event", help="事件类型") + + # register + reg_p = sub.add_parser("register", help="注册 Hook") + reg_p.add_argument("--name", required=True) + reg_p.add_argument("--event", required=True, choices=HookEvent.all()) + reg_p.add_argument("--handler", required=True) + reg_p.add_argument("--type", default="shell", choices=["shell", "python", "node"]) + reg_p.add_argument("--priority", type=int, default=0) + reg_p.add_argument("--timeout", type=int, default=30) + reg_p.add_argument("--desc", default="") + + # unregister + unreg_p = sub.add_parser("unregister", help="注销 Hook") + unreg_p.add_argument("--name", required=True) + unreg_p.add_argument("--event") + + # trigger + trig_p = sub.add_parser("trigger", help="触发 Hook") + trig_p.add_argument("--event", required=True) + trig_p.add_argument("--ctx", default="{}", + help="JSON 上下文") + + # enable/disable + en_p = sub.add_parser("enable", help="启用 Hook") + en_p.add_argument("--name", required=True) + en_p.add_argument("--event") + + dis_p = sub.add_parser("disable", help="禁用 Hook") + dis_p.add_argument("--name", required=True) + dis_p.add_argument("--event") + + args = parser.parse_args() + + if args.mcp: + asyncio.run(run_mcp()) + return + + mgr = HookManager() + + if args.mode == "list": + hooks = mgr.list_hooks(args.event) + print(json.dumps(hooks, indent=2, ensure_ascii=False)) + print(f"\nTotal: {len(hooks)} hooks") + + elif args.mode == "register": + mgr.register(Hook( + name=args.name, event=args.event, + handler=args.handler, type=args.type, + priority=args.priority, timeout=args.timeout, + description=args.desc, + )) + print(f"[OK] Hook '{args.name}' registered on event '{args.event}'") + + elif args.mode == "unregister": + mgr.unregister(args.name, args.event) + print(f"[OK] Hook '{args.name}' unregistered") + + elif args.mode == "trigger": + ctx = HookContext(event=args.event, **json.loads(args.ctx)) + results = asyncio.run(mgr.trigger(args.event, ctx)) + print(json.dumps(results, indent=2, ensure_ascii=False)) + + elif args.mode == "enable": + mgr.enable(args.name, args.event) + print(f"[OK] Hook '{args.name}' enabled") + + elif args.mode == "disable": + mgr.disable(args.name, args.event) + print(f"[OK] Hook '{args.name}' disabled") + + else: + parser.print_help() + + +if __name__ == "__main__": + main() diff --git a/.deepcode/skills/deepcode-hooks/plugin.json b/.deepcode/skills/deepcode-hooks/plugin.json new file mode 100644 index 00000000..7b890339 --- /dev/null +++ b/.deepcode/skills/deepcode-hooks/plugin.json @@ -0,0 +1,11 @@ +{ + "name": "deepcode-hooks", + "version": "1.0.0", + "description": "DeepCode Hook System — 通用 Pre/Post 钩子框架 (对标 Claude Code hook-handler.cjs)", + "author": "DeepCode + Ghidra RE (Claude Code v2.1.216)", + "entry": "SKILL.md", + "dependencies": [], + "modules": [ + "hooks.py" + ] +} diff --git a/.deepcode/skills/deepcode-knowledge/SKILL.md b/.deepcode/skills/deepcode-knowledge/SKILL.md new file mode 100644 index 00000000..6332478f --- /dev/null +++ b/.deepcode/skills/deepcode-knowledge/SKILL.md @@ -0,0 +1,89 @@ +--- +name: deepcode-knowledge +description: > + DeepCode 统一知识引擎 — Obsidian式知识库 + KV记忆系统二合一。 + 知识库: 结构化Markdown笔记、模板(stock/daily/strategy/note)、[[wikilink]]双向链接、 + 知识图谱、每日复盘。记忆: KV持久化、多后端自动路由、全文搜索。 + Use when saving analysis results, searching knowledge base, managing memories, + generating daily reviews, or building personal knowledge systems. +version: 1.0.0 +author: DeepCode +date: 2026-07-29 +tags: [knowledge, vault, memory, obsidian, notes, templates, graph] +--- + +# DeepCode Knowledge Engine + +合并 `deepcode-vault` (知识库) + `deepcode-memory` (记忆) 的统一知识引擎。 + +## 架构 + +``` +AI 调用 + │ + ▼ +knowledge_server.py (MCP, 14 tools) + ├── Vault 引擎 (8 tools) + │ ├── Markdown 笔记 CRUD + │ ├── [[wikilink]] 双向链接 + │ ├── 知识图谱 + │ ├── 模板渲染 (stock/daily/strategy/note) + │ └── Obsidian 双向同步 + └── Memory 引擎 (6 tools) + ├── KV 存储 (save/load/forget) + ├── 全文搜索 + └── 5 后端自动路由 (ruflo/claude/flow/tokensave/plan) +``` + +## 14 工具一览 + +### Vault (知识库) +| 工具 | 说明 | +|:------|:-----| +| `knowledge__vault_status` | 知识库统计 — 笔记数、类型分布、磁盘占用 | +| `knowledge__save_analysis` | 保存分析结果 — 自动模板 + 标签 + [[wikilink]] | +| `knowledge__search_notes` | 全文搜索笔记 | +| `knowledge__read_note` | 读取笔记完整内容 (frontmatter + body) | +| `knowledge__show_graph` | 知识图谱 — 笔记关联网络 | +| `knowledge__daily_log` | 每日复盘 — generate/show/list | +| `knowledge__list_templates` | 列出 4 种分析模板 | +| `knowledge__sync_to_obsidian` | 同步到 Obsidian 仓库 | + +### Memory (记忆) +| 工具 | 说明 | +|:------|:-----| +| `knowledge__memory_save` | 保存 KV 记忆 | +| `knowledge__memory_load` | 加载记忆值 | +| `knowledge__memory_search` | 全文搜索记忆 | +| `knowledge__memory_list` | 列出所有键名 | +| `knowledge__memory_forget` | 删除记忆 | +| `knowledge__memory_stats` | 记忆系统统计 | + +## MCP 注册 + +```json +"deepcode-knowledge": { + "command": "python", + "args": ["F:/DEEPCODE/.deepcode/skills/deepcode-knowledge/knowledge_server.py"] +} +``` + +## 使用示例 + +``` +# 保存茅台分析 +mcp__deepcode-knowledge__knowledge__save_analysis + type=stock, symbol=600519, title="茅台技术面分析" + +# 搜索历史分析 +mcp__deepcode-knowledge__knowledge__search_notes query=茅台 + +# 查看知识图谱 +mcp__deepcode-knowledge__knowledge__show_graph symbol=600519 + +# 保存会话记忆 +mcp__deepcode-knowledge__knowledge__memory_save key=last_task value="分析茅台" + +# 加载记忆 +mcp__deepcode-knowledge__knowledge__memory_load key=last_task +``` diff --git a/.deepcode/skills/deepcode-knowledge/knowledge_server.py b/.deepcode/skills/deepcode-knowledge/knowledge_server.py new file mode 100644 index 00000000..ad9429ca --- /dev/null +++ b/.deepcode/skills/deepcode-knowledge/knowledge_server.py @@ -0,0 +1,539 @@ +#!/usr/bin/env python3 +""" +DeepCode Knowledge Engine — 统一知识引擎 MCP Server +═══════════════════════════════════════════════════ +合并 deepcode-vault (Obsidian式知识库) + deepcode-memory (KV记忆系统) + +Vault 工具 (8): + vault_status, save_analysis, search_notes, read_note, show_graph, + daily_log, list_templates, sync_to_obsidian + +Memory 工具 (6): + memory_save, memory_load, memory_search, memory_list, memory_forget, memory_stats + +MCP 注册: + "deepcode-knowledge": { + "command": "python", + "args": ["F:/DEEPCODE/.deepcode/skills/deepcode-knowledge/knowledge_server.py"] + } +""" +import json, sys, os, re, time, uuid, shutil, sqlite3 +from datetime import datetime, date +from pathlib import Path +from collections import defaultdict + +# ── Windows GBK 编码兼容 ── +if sys.platform == "win32" and sys.stdout.encoding and sys.stdout.encoding.lower() in ("gbk", "gb2312", "cp936"): + sys.stdout.reconfigure(encoding="utf-8") + sys.stderr.reconfigure(encoding="utf-8") + +# ── 路径配置 ── +SKILL_DIR = os.path.dirname(os.path.abspath(__file__)) +DEFAULT_VAULT = os.path.join(SKILL_DIR, "data", "vault") +NOTES_DIR = os.path.join(DEFAULT_VAULT, "notes") +DAILY_DIR = os.path.join(DEFAULT_VAULT, "daily") +TEMPLATES_DIR = os.path.join(SKILL_DIR, "data", "templates") +DB_PATH = os.path.join(SKILL_DIR, "data", "knowledge.db") +os.makedirs(NOTES_DIR, exist_ok=True) +os.makedirs(DAILY_DIR, exist_ok=True) +os.makedirs(TEMPLATES_DIR, exist_ok=True) + +# ── Memory Manager 导入 ── +_scripts_dir = os.path.join(os.path.dirname(SKILL_DIR), "..", "..", "scripts") +if _scripts_dir not in sys.path: + sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) +from memory_manager import MemoryManager +mm = MemoryManager() + + +# ══════════════════════════════════════════════ +# 数据库层 — SQLite 索引 +# ══════════════════════════════════════════════ + +def get_db(): + conn = sqlite3.connect(DB_PATH) + conn.row_factory = sqlite3.Row + conn.executescript(""" + CREATE TABLE IF NOT EXISTS notes ( + id TEXT PRIMARY KEY, path TEXT UNIQUE NOT NULL, + title TEXT NOT NULL, type TEXT DEFAULT 'note', + tags TEXT DEFAULT '', symbol TEXT DEFAULT '', + created_at REAL NOT NULL, updated_at REAL NOT NULL, + content_preview TEXT DEFAULT '' + ); + CREATE TABLE IF NOT EXISTS links ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + source_id TEXT NOT NULL, target_id TEXT NOT NULL, + link_type TEXT DEFAULT 'related', + UNIQUE(source_id, target_id) + ); + CREATE TABLE IF NOT EXISTS daily_logs ( + id TEXT PRIMARY KEY, date TEXT UNIQUE NOT NULL, + summary TEXT DEFAULT '', mood TEXT DEFAULT '', + created_at REAL NOT NULL + ); + CREATE INDEX IF NOT EXISTS idx_notes_type ON notes(type); + CREATE INDEX IF NOT EXISTS idx_notes_symbol ON notes(symbol); + CREATE INDEX IF NOT EXISTS idx_notes_tags ON notes(tags); + CREATE INDEX IF NOT EXISTS idx_links_source ON links(source_id); + CREATE INDEX IF NOT EXISTS idx_links_target ON links(target_id); + """) + conn.commit() + return conn + + +# ══════════════════════════════════════════════ +# 工具函数 +# ══════════════════════════════════════════════ + +def now_ts(): return time.time() +def date_str(): return date.today().isoformat() +def slugify(text): + s = re.sub(r'[^\w\s-]', '', text).strip().lower() + return re.sub(r'[-\s]+', '-', s)[:60] +def extract_tags(text): + return re.findall(r'#([\w\u4e00-\u9fff\-.]+)', text) +def extract_wikilinks(text): + return re.findall(r'\[\[([^\]]+)\]\]', text) +def extract_frontmatter(text): + m = re.match(r'^---\s*\n(.*?)\n---\s*\n', text, re.DOTALL) + if not m: return {}, text + fm = {} + for line in m.group(1).strip().split('\n'): + if ':' in line: + k, v = line.split(':', 1) + fm[k.strip()] = v.strip().strip('"').strip("'") + return fm, text[m.end():] + + +# ══════════════════════════════════════════════ +# 模板引擎 +# ══════════════════════════════════════════════ + +TEMPLATES = { + "stock": """--- +title: "{title}" +type: stock +symbol: {symbol} +tags: {tags} +created: {date} +--- + +# {title} + +## 基本面 +{fundamental} + +## 技术面 +{technical} + +## 资金面 +{fund_flow} + +## 综合判断 +{verdict} + +## 关联 +{links} +""", + "daily": """--- +title: "交易复盘 - {date}" +type: daily +date: {date} +tags: 复盘 +--- + +# 交易复盘 - {date} + +## 大盘概况 +{market_overview} + +## 今日操作 +{operations} + +## 持仓分析 +{holdings} + +## 明日计划 +{plan} +""", + "strategy": """--- +title: "{title}" +type: strategy +tags: {tags} +created: {date} +--- + +# {title} + +## 策略逻辑 +{logic} + +## 参数配置 +{params} + +## 回测结果 +{backtest} + +## 适用场景 +{scenarios} +""", + "note": """--- +title: "{title}" +type: note +tags: {tags} +created: {date} +--- + +# {title} + +{content} + +## 相关笔记 +{links} +""" +} + +def render_template(template_name, **kwargs): + tpl = TEMPLATES.get(template_name, TEMPLATES["note"]) + kwargs.setdefault("date", date_str()) + kwargs.setdefault("tags", "") + kwargs.setdefault("links", "") + for k, v in kwargs.items(): + if v is None: kwargs[k] = "" + return tpl.format(**kwargs) + + +# ══════════════════════════════════════════════ +# VAULT 工具 (8) +# ══════════════════════════════════════════════ + +def tool_vault_status(args): + conn = get_db() + try: + total = conn.execute("SELECT COUNT(*) FROM notes").fetchone()[0] + by_type = conn.execute("SELECT type, COUNT(*) as cnt FROM notes GROUP BY type").fetchall() + daily_count = conn.execute("SELECT COUNT(*) FROM daily_logs").fetchone()[0] + link_count = conn.execute("SELECT COUNT(*) FROM links").fetchone()[0] + recent = conn.execute("SELECT title, type, symbol, updated_at FROM notes ORDER BY updated_at DESC LIMIT 5").fetchall() + return { + "vault_path": DEFAULT_VAULT, + "total_notes": total, "daily_logs": daily_count, "total_links": link_count, + "by_type": {r["type"]: r["cnt"] for r in by_type}, + "recent": [{"title": r["title"], "type": r["type"], "symbol": r["symbol"], + "updated": datetime.fromtimestamp(r["updated_at"]).strftime("%m-%d %H:%M")} for r in recent], + "disk_usage_kb": sum(os.path.getsize(os.path.join(dp, f)) for dp, _, fs in os.walk(DEFAULT_VAULT) for f in fs if f.endswith(".md")) // 1024 if os.path.exists(DEFAULT_VAULT) else 0 + } + finally: + conn.close() + + +def tool_save_analysis(args): + note_type = args.get("type", "note") + title = args.get("title", "未命名笔记") + symbol = args.get("symbol", "") + tags = args.get("tags", "") + content = args.get("content", {}) + if isinstance(content, str): + try: content = json.loads(content) + except: content = {"content": content} + + if note_type == "stock": + md = render_template("stock", title=title, symbol=symbol, + tags=tags or f"股票,{symbol}", + fundamental=content.get("fundamental", ""), + technical=content.get("technical", ""), + fund_flow=content.get("fund_flow", ""), + verdict=content.get("verdict", ""), + links=content.get("links", "")) + elif note_type == "strategy": + md = render_template("strategy", title=title, tags=tags, + logic=content.get("logic", ""), params=content.get("params", ""), + backtest=content.get("backtest", ""), scenarios=content.get("scenarios", "")) + else: + md = render_template("note", title=title, tags=tags, + content=content.get("content", str(content))) + + filename = f"{slugify(title)}.md" + note_path = os.path.join(NOTES_DIR, filename) + existing_fm = {} + if os.path.exists(note_path): + with open(note_path, "r", encoding="utf-8") as f: + existing_fm, _ = extract_frontmatter(f.read()) + with open(note_path, "w", encoding="utf-8") as f: + f.write(md) + + note_id = slugify(title) + "_" + str(int(now_ts()))[-6:] + extracted_tags = extract_tags(md) + all_tags = ",".join(set(extracted_tags + ([t.strip() for t in tags.split(",") if t.strip()] if tags else []))) + links = extract_wikilinks(md) + + conn = get_db() + try: + conn.execute("""INSERT OR REPLACE INTO notes (id, path, title, type, tags, symbol, created_at, updated_at, content_preview) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)""", + (note_id, note_path, title, note_type, all_tags, symbol, + existing_fm.get("created", now_ts()) if existing_fm else now_ts(), now_ts(), md[:200])) + for link_title in links: + target = conn.execute("SELECT id FROM notes WHERE title LIKE ? LIMIT 1", (f"%{link_title}%",)).fetchone() + if target: + try: conn.execute("INSERT OR IGNORE INTO links (source_id, target_id, link_type) VALUES (?, ?, ?)", (note_id, target["id"], "related")) + except: pass + conn.commit() + finally: + conn.close() + return {"status": "saved", "note_id": note_id, "path": note_path, "title": title, "type": note_type, "tags": all_tags, "links_found": len(links)} + + +def tool_search_notes(args): + query = args.get("query", ""); note_type = args.get("type", ""); tag = args.get("tag", ""); limit = min(args.get("limit", 20), 100) + conn = get_db() + try: + sql, params = "SELECT * FROM notes WHERE 1=1", [] + if query: sql += " AND (title LIKE ? OR content_preview LIKE ? OR tags LIKE ?)"; q = f"%{query}%"; params.extend([q, q, q]) + if note_type: sql += " AND type = ?"; params.append(note_type) + if tag: sql += " AND tags LIKE ?"; params.append(f"%{tag}%") + sql += " ORDER BY updated_at DESC LIMIT ?"; params.append(limit) + rows = conn.execute(sql, params).fetchall() + return {"query": query, "count": len(rows), "results": [ + {"id": r["id"], "title": r["title"], "type": r["type"], "tags": r["tags"], + "symbol": r["symbol"], "path": r["path"], + "updated": datetime.fromtimestamp(r["updated_at"]).strftime("%Y-%m-%d %H:%M"), + "preview": r["content_preview"][:150]} for r in rows]} + finally: + conn.close() + + +def tool_read_note(args): + note_id = args.get("id", ""); filepath = args.get("path", "") + if filepath and os.path.exists(filepath): path = filepath + elif note_id: + conn = get_db() + try: + row = conn.execute("SELECT path FROM notes WHERE id = ?", (note_id,)).fetchone() + path = row["path"] if row else None + finally: conn.close() + if not path or not os.path.exists(path): return {"error": "笔记不存在", "id": note_id} + else: return {"error": "需要提供 id 或 path"} + with open(path, "r", encoding="utf-8") as f: content = f.read() + fm, body = extract_frontmatter(content) + return {"frontmatter": fm, "body": body.strip(), "path": path, "length": len(content)} + + +def tool_show_graph(args): + symbol = args.get("symbol", ""); limit = min(args.get("limit", 50), 200) + conn = get_db() + try: + if symbol: + nodes = conn.execute("""SELECT DISTINCT n.* FROM notes n WHERE n.symbol = ? OR n.id IN ( + SELECT l.target_id FROM links l JOIN notes s ON s.id = l.source_id WHERE s.symbol = ?) LIMIT ?""", + (symbol, symbol, limit)).fetchall() + edges = conn.execute("""SELECT DISTINCT l.* FROM links l WHERE l.source_id IN (SELECT id FROM notes WHERE symbol = ?) + OR l.target_id IN (SELECT id FROM notes WHERE symbol = ?) LIMIT ?""", (symbol, symbol, limit*2)).fetchall() + else: + nodes = conn.execute("SELECT * FROM notes ORDER BY updated_at DESC LIMIT ?", (limit,)).fetchall() + edges = conn.execute("SELECT * FROM links LIMIT ?", (limit*2,)).fetchall() + return {"nodes": [{"id": n["id"], "title": n["title"], "type": n["type"], "symbol": n["symbol"], "tags": n["tags"]} for n in nodes], + "edges": [{"source": e["source_id"], "target": e["target_id"], "type": e["link_type"]} for e in edges], + "stats": {"node_count": len(nodes), "edge_count": len(edges)}} + finally: conn.close() + + +def tool_daily_log(args): + action = args.get("action", "show"); log_date = args.get("date", date_str()) + filepath = os.path.join(DAILY_DIR, f"{log_date}.md") + if action == "generate": + md = render_template("daily", date=log_date, + market_overview=args.get("market_overview", "待补充"), + operations=args.get("operations", "今日无操作"), + holdings=args.get("holdings", "待补充"), + plan=args.get("plan", "待补充")) + with open(filepath, "w", encoding="utf-8") as f: f.write(md) + conn = get_db() + try: + conn.execute("INSERT OR REPLACE INTO daily_logs (id, date, summary, created_at) VALUES (?, ?, ?, ?)", + (log_date, log_date, args.get("summary", "")[:500], now_ts())); conn.commit() + finally: conn.close() + return {"status": "generated", "date": log_date, "path": filepath} + elif action == "show": + if os.path.exists(filepath): + with open(filepath, "r", encoding="utf-8") as f: content = f.read() + fm, body = extract_frontmatter(content) + return {"date": log_date, "frontmatter": fm, "body": body.strip(), "exists": True} + return {"date": log_date, "exists": False, "message": f"{log_date} 还没有复盘日志,使用 action=generate 创建"} + elif action == "list": + conn = get_db() + try: + rows = conn.execute("SELECT * FROM daily_logs ORDER BY date DESC LIMIT 30").fetchall() + return {"count": len(rows), "logs": [{"date": r["date"], "summary": r["summary"][:100], + "created": datetime.fromtimestamp(r["created_at"]).strftime("%Y-%m-%d %H:%M")} for r in rows]} + finally: conn.close() + return {"error": f"未知 action: {action}"} + + +def tool_list_templates(args): + return {"templates": [ + {"name": "stock", "description": "个股分析报告", "fields": ["symbol", "title", "fundamental", "technical", "fund_flow", "verdict"]}, + {"name": "daily", "description": "每日复盘日志", "fields": ["date", "market_overview", "operations", "holdings", "plan"]}, + {"name": "strategy", "description": "交易策略文档", "fields": ["title", "logic", "params", "backtest", "scenarios"]}, + {"name": "note", "description": "通用笔记", "fields": ["title", "content", "tags"]} + ], "total": 4} + + +def tool_sync_to_obsidian(args): + obsidian_vault = args.get("vault", "") + if not obsidian_vault or not os.path.isdir(obsidian_vault): + return {"error": f"Obsidian 仓库路径无效: {obsidian_vault}"} + for subdir in ["stocks", "daily", "strategies", "notes"]: + os.makedirs(os.path.join(obsidian_vault, subdir), exist_ok=True) + stats = {"exported": 0, "skipped": 0, "errors": 0} + conn = get_db() + try: + for note in conn.execute("SELECT * FROM notes ORDER BY updated_at DESC").fetchall(): + target_dir = os.path.join(obsidian_vault, {"stock":"stocks","strategy":"strategies"}.get(note["type"],"notes")) + if os.path.exists(note["path"]): + try: shutil.copy2(note["path"], os.path.join(target_dir, os.path.basename(note["path"]))); stats["exported"] += 1 + except: stats["errors"] += 1 + finally: conn.close() + for f in os.listdir(DAILY_DIR): + if f.endswith(".md"): + try: shutil.copy2(os.path.join(DAILY_DIR, f), os.path.join(obsidian_vault, "daily", f)); stats["exported"] += 1 + except: pass + return {"status": "synced", "obsidian_vault": obsidian_vault, "stats": stats, + "note": "在 Obsidian 中可打开该仓库查看,支持 [[wikilink]] 跳转和 Graph View"} + + +# ══════════════════════════════════════════════ +# MEMORY 工具 (6) — 委托给 MemoryManager +# ══════════════════════════════════════════════ + +def tool_memory_save(args): + key = args.get("key", ""); value = args.get("value", "") + if not key: return {"ok": False, "error": "Missing 'key'"} + mm.save(key, value, backend=args.get("backend", "auto"), tags=args.get("tags", [])) + return {"ok": True, "result": f"saved: {key}"} + +def tool_memory_load(args): + key = args.get("key", "") + if not key: return {"ok": False, "error": "Missing 'key'"} + return {"ok": True, "result": mm.load(key)} + +def tool_memory_search(args): + query = args.get("query", "") + limit = min(args.get("limit", 10), 100) + + # 1) 搜索 Memory 后端 (KV) + memory_results = mm.search(query, limit=limit) + + # 2) 搜索 Vault 笔记 (Markdown) + vault_results = [] + if query: + conn = get_db() + try: + sql = "SELECT * FROM notes WHERE title LIKE ? OR content_preview LIKE ? OR tags LIKE ?" + q = f"%{query}%" + rows = conn.execute(sql + " ORDER BY updated_at DESC LIMIT ?", (q, q, q, limit)).fetchall() + for r in rows: + vault_results.append({ + "key": f"vault:{r['id']}", + "value": r["title"], + "backend": "vault", + "preview": r["content_preview"][:200], + "type": r["type"], + "tags": r["tags"], + "symbol": r["symbol"], + "path": r["path"], + "updated": datetime.fromtimestamp(r["updated_at"]).strftime("%Y-%m-%d %H:%M"), + }) + finally: + conn.close() + + # 3) 合并结果:先 Memory 后 Vault + merged = memory_results + vault_results + return { + "ok": True, + "result": merged[:limit], + "count": {"memory": len(memory_results), "vault": len(vault_results), "total": len(merged)}, + } + +def tool_memory_list(args): + return {"ok": True, "result": mm.list()} + +def tool_memory_forget(args): + key = args.get("key", "") + if not key: return {"ok": False, "error": "Missing 'key'"} + return {"ok": True, "result": mm.forget(key)} + +def tool_memory_stats(args): + return {"ok": True, "result": mm.stats()} + + +# ══════════════════════════════════════════════ +# 工具注册表 +# ══════════════════════════════════════════════ + +TOOLS = { + "knowledge__vault_status": {"handler": tool_vault_status, "desc": "查看知识库统计 — 笔记数、类型分布、最近更新"}, + "knowledge__save_analysis": {"handler": tool_save_analysis, "desc": "保存分析结果到知识库 — 自动套用模板、提取标签、建立 [[wikilink]] 链接"}, + "knowledge__search_notes": {"handler": tool_search_notes, "desc": "全文搜索知识库笔记 — 支持标题/正文/标签"}, + "knowledge__read_note": {"handler": tool_read_note, "desc": "读取笔记完整内容 — 解析 YAML frontmatter 和 Markdown 正文"}, + "knowledge__show_graph": {"handler": tool_show_graph, "desc": "查看知识图谱 — 笔记之间的关联关系网络"}, + "knowledge__daily_log": {"handler": tool_daily_log, "desc": "每日复盘日志 — 生成/查看/列出"}, + "knowledge__list_templates": {"handler": tool_list_templates, "desc": "列出所有可用分析模板 (stock/daily/strategy/note)"}, + "knowledge__sync_to_obsidian": {"handler": tool_sync_to_obsidian, "desc": "同步知识库到 Obsidian 仓库"}, + "knowledge__memory_save": {"handler": tool_memory_save, "desc": "保存一条记忆 (KV) — 自动选择最优后端"}, + "knowledge__memory_load": {"handler": tool_memory_load, "desc": "加载指定键的记忆值"}, + "knowledge__memory_search": {"handler": tool_memory_search, "desc": "全文搜索记忆内容"}, + "knowledge__memory_list": {"handler": tool_memory_list, "desc": "列出所有记忆键名"}, + "knowledge__memory_forget": {"handler": tool_memory_forget, "desc": "删除一条记忆"}, + "knowledge__memory_stats": {"handler": tool_memory_stats, "desc": "记忆系统统计 — 各后端状态"}, +} + + +# ══════════════════════════════════════════════ +# MCP JSON-RPC 2.0 stdio +# ══════════════════════════════════════════════ + +def make_result(request, result): + return {"jsonrpc": "2.0", "id": request.get("id"), "result": result} + +def make_error(request, code, message): + return {"jsonrpc": "2.0", "id": request.get("id"), "error": {"code": code, "message": message}} + +def handle_request(request): + method = request.get("method", "") + if method == "initialize": + return make_result(request, {"protocolVersion": "2025-03-26", "capabilities": {"tools": {"listChanged": True}}, + "serverInfo": {"name": "deepcode-knowledge", "version": "1.0.0"}}) + elif method == "tools/list": + return make_result(request, {"tools": [ + {"name": name, "description": info["desc"], "inputSchema": {"type": "object", "properties": {}}} + for name, info in TOOLS.items()]}) + elif method == "tools/call": + tool_name = request.get("params", {}).get("name", "") + args = request.get("params", {}).get("arguments", {}) + if tool_name in TOOLS: + try: + result = TOOLS[tool_name]["handler"](args) + return make_result(request, {"content": [{"type": "text", "text": json.dumps(result, ensure_ascii=False, indent=2)}]}) + except Exception as e: + return make_result(request, {"content": [{"type": "text", "text": json.dumps({"error": str(e)}, ensure_ascii=False)}], "isError": True}) + return make_error(request, -32601, f"Unknown tool: {tool_name}") + elif method == "notifications/initialized": + return None + elif method == "ping": + return make_result(request, {}) + return make_error(request, -32601, f"Unknown method: {method}") + + +if __name__ == "__main__": + for line in sys.stdin: + line = line.strip() + if not line: continue + try: + request = json.loads(line) + response = handle_request(request) + if response is not None: + print(json.dumps(response, ensure_ascii=False), flush=True) + except json.JSONDecodeError: + continue diff --git a/.deepcode/skills/deepcode-knowledge/memory_manager.py b/.deepcode/skills/deepcode-knowledge/memory_manager.py new file mode 100644 index 00000000..7e82bbef --- /dev/null +++ b/.deepcode/skills/deepcode-knowledge/memory_manager.py @@ -0,0 +1,643 @@ +#!/usr/bin/env python3 +""" +DeepCode Unified Memory Manager +═══════════════════════════════════ +整合 5 套记忆系统为统一 API。 + +后 端 位 置 格 式 状 态 +────────────────────────────────────────────────────────────── +ruflo data/memory/memory.db SQLite 37表 ⚠️ 空 +claude .claude/memory.db SQLite 11表 ⚠️ 空 +flow .claude-flow/data/*.json JSON文件 ⚠️ 空 +tokensave .deepcode/skills/token-saver/data/ SQLite ✅ 有数据 +plan task_plan.md + progress.md Markdown ✅ 有数据 + +API: + save(key, value, backend="auto", tags=None) → str id + load(key, backend="auto") → value or None + search(query, backend="auto", limit=10) → list[dict] + list(backend="auto") → list[str] + forget(key, backend="auto") → bool + stats() → dict + +CLI: + python scripts/memory_manager.py save [--backend] + python scripts/memory_manager.py load + python scripts/memory_manager.py search + python scripts/memory_manager.py list [--backend] + python scripts/memory_manager.py forget + python scripts/memory_manager.py stats + python scripts/memory_manager.py seed # 写入示例数据启动记忆 +""" + +import json +import os +import sqlite3 +import time +from datetime import datetime +from pathlib import Path +from typing import Optional + +PROJECT_ROOT = Path(__file__).parent.parent + +# ── 后端路径 ── +BACKEND_PATHS = { + "ruflo": PROJECT_ROOT / "data" / "memory" / "memory.db", + "claude": PROJECT_ROOT / ".claude" / "memory.db", + "flow": PROJECT_ROOT / ".claude-flow" / "data", + "tokensave": PROJECT_ROOT / ".deepcode" / "skills" / "token-saver" / "data" / "token_saver_memory.db", + "plan_dir": PROJECT_ROOT, # task_plan.md / progress.md / findings.md +} + + +# ═══════════════════════════════════════════════════ +# 核心 API +# ═══════════════════════════════════════════════════ + +class MemoryManager: + """统一记忆管理器""" + + def __init__(self): + self._backends = { + "ruflo": RuFloBackend(), + "claude": ClaudeNativeBackend(), + "flow": ClaudeFlowBackend(), + "tokensave": TokenSaverBackend(), + "plan": PlanFilesBackend(), + } + + def _resolve_backend(self, backend: str): + """解析 backend 参数为具体后端实例列表""" + if backend == "auto": + return list(self._backends.values()) + if backend in self._backends: + return [self._backends[backend]] + raise ValueError(f"Unknown backend: {backend}. Options: auto, {', '.join(self._backends.keys())}") + + def save(self, key: str, value: str, backend: str = "auto", + tags: Optional[list] = None) -> str: + """保存记忆到最优后端""" + record = { + "key": key, + "value": value, + "tags": tags or [], + "timestamp": datetime.now().isoformat(), + } + + # auto: 根据内容类型选后端 + if backend == "auto": + if tags and any(t in ["task", "progress", "plan"] for t in tags): + target = "plan" + elif tags and any(t in ["compress", "token", "saving"] for t in tags): + target = "tokensave" + else: + target = "claude" # 默认存 claude(最快) + else: + target = backend + + backends = self._resolve_backend(target) + for b in backends: + b.save(record) + return key + + def load(self, key: str, backend: str = "auto"): + """加载记忆""" + for b in self._resolve_backend(backend): + result = b.load(key) + if result is not None: + return result + return None + + def search(self, query: str, backend: str = "auto", limit: int = 10) -> list: + """搜索记忆""" + results = [] + for b in self._resolve_backend(backend): + results.extend(b.search(query, limit)) + if len(results) >= limit: + break + return results[:limit] + + def list(self, backend: str = "auto") -> list: + """列出所有记忆键名""" + all_keys = [] + for b in self._resolve_backend(backend): + all_keys.extend(b.list_keys()) + return sorted(set(all_keys)) + + def forget(self, key: str, backend: str = "auto") -> bool: + """删除记忆""" + found = False + for b in self._resolve_backend(backend): + if b.forget(key): + found = True + return found + + def stats(self) -> dict: + """所有后端统计""" + return {name: bk.stats() for name, bk in self._backends.items()} + + def seed(self): + """写入示例记忆,启动记忆系统""" + seeds = [ + ("user_preference", json.dumps({ + "model": "deepseek-v4-pro", + "reasoning_effort": "max", + "thinking": True, + }), ["preference", "user"], "claude"), + ("workflow_pattern", "股票分析流程: scan -> analyze -> chanlun -> report", + ["pattern", "workflow"], "claude"), + ("last_task", "整合记忆系统", ["task", "progress"], "plan"), + ("mcp_server_list", ",".join([ + "tushareMcp", "playwright", "filesystem", "sqlite", + "duckdb", "github", "ghidra-mcp", "router-mcp", + ]), ["config", "mcp"], "claude"), + ] + count = 0 + for key, value, tags, backend in seeds: + try: + self.save(key, value, backend=backend, tags=tags) + count += 1 + except Exception as e: + print(f" [seed] Failed: {key}: {e}") + print(f"[memory_manager] Seeded {count} memories") + return count + + +# ═══════════════════════════════════════════════════ +# 后端适配器 +# ═══════════════════════════════════════════════════ + +class BaseBackend: + """后端基类""" + name = "base" + + def save(self, record: dict): + raise NotImplementedError + + def load(self, key: str): + raise NotImplementedError + + def search(self, query: str, limit: int): + raise NotImplementedError + + def list_keys(self) -> list: + raise NotImplementedError + + def forget(self, key: str) -> bool: + raise NotImplementedError + + def stats(self) -> dict: + return {"status": "unknown"} + + +class RuFloBackend(BaseBackend): + """RuFlo V3 记忆引擎 — data/memory/memory.db""" + name = "ruflo" + + def __init__(self): + self.db_path = BACKEND_PATHS["ruflo"] + + def _connect(self): + return sqlite3.connect(str(self.db_path)) + + def save(self, record: dict): + try: + conn = self._connect() + conn.execute( + "INSERT OR REPLACE INTO memory_entries (id, key, content, tags, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?)", + (record["key"], record["key"], record["value"], + json.dumps(record.get("tags", [])), + int(time.time()), int(time.time())), + ) + conn.commit() + conn.close() + except sqlite3.OperationalError as e: + pass + + def load(self, key: str): + try: + conn = self._connect() + cur = conn.execute("SELECT content FROM memory_entries WHERE key=?", (key,)) + row = cur.fetchone() + conn.close() + return row[0] if row else None + except sqlite3.OperationalError: + pass + return None + + def search(self, query: str, limit: int = 10): + try: + conn = self._connect() + cur = conn.execute( + "SELECT key, content FROM memory_entries WHERE key LIKE ? OR content LIKE ? LIMIT ?", + (f"%{query}%", f"%{query}%", limit), + ) + rows = cur.fetchall() + conn.close() + return [{"key": r[0], "value": r[1], "backend": "ruflo"} for r in rows] + except sqlite3.OperationalError: + return [] + + def list_keys(self): + try: + conn = self._connect() + cur = conn.execute("SELECT key FROM memory_entries") + keys = [r[0] for r in cur.fetchall()] + conn.close() + return keys + except sqlite3.OperationalError: + return [] + + def forget(self, key: str) -> bool: + try: + conn = self._connect() + conn.execute("DELETE FROM memory_entries WHERE key=?", (key,)) + conn.commit() + affected = conn.total_changes + conn.close() + return affected > 0 + except sqlite3.OperationalError: + return False + + def stats(self): + try: + conn = self._connect() + cur = conn.execute("SELECT COUNT(*) FROM memory_entries") + count = cur.fetchone()[0] + conn.close() + return {"status": "active", "entries": count} + except sqlite3.OperationalError: + return {"status": "empty (table not ready)", "entries": 0} + + +class ClaudeNativeBackend(BaseBackend): + """Claude Code 原生记忆 — .claude/memory.db""" + name = "claude" + + def __init__(self): + self.db_path = BACKEND_PATHS["claude"] + + def _connect(self): + self.db_path.parent.mkdir(parents=True, exist_ok=True) + conn = sqlite3.connect(str(self.db_path)) + conn.execute(""" + CREATE TABLE IF NOT EXISTS unified_memory ( + key TEXT PRIMARY KEY, + value TEXT, + tags TEXT, + created_at TEXT, + updated_at TEXT + ) + """) + conn.commit() + return conn + + def save(self, record: dict): + conn = self._connect() + conn.execute( + "INSERT OR REPLACE INTO unified_memory (key, value, tags, created_at, updated_at) VALUES (?, ?, ?, ?, ?)", + (record["key"], record["value"], json.dumps(record.get("tags", [])), + record["timestamp"], record["timestamp"]), + ) + conn.commit() + conn.close() + + def load(self, key: str): + conn = self._connect() + cur = conn.execute("SELECT value FROM unified_memory WHERE key=?", (key,)) + row = cur.fetchone() + conn.close() + return row[0] if row else None + + def search(self, query: str, limit: int = 10): + conn = self._connect() + cur = conn.execute( + "SELECT key, value FROM unified_memory WHERE key LIKE ? OR value LIKE ? LIMIT ?", + (f"%{query}%", f"%{query}%", limit), + ) + rows = cur.fetchall() + conn.close() + return [{"key": r[0], "value": r[1], "backend": "claude"} for r in rows] + + def list_keys(self): + conn = self._connect() + cur = conn.execute("SELECT key FROM unified_memory ORDER BY updated_at DESC") + keys = [r[0] for r in cur.fetchall()] + conn.close() + return keys + + def forget(self, key: str) -> bool: + conn = self._connect() + conn.execute("DELETE FROM unified_memory WHERE key=?", (key,)) + conn.commit() + affected = conn.total_changes + conn.close() + return affected > 0 + + def stats(self): + conn = self._connect() + cur = conn.execute("SELECT COUNT(*) FROM unified_memory") + count = cur.fetchone()[0] + conn.close() + return {"status": "active", "entries": count} + + +class ClaudeFlowBackend(BaseBackend): + """Claude Flow 记忆 — .claude-flow/data/*.json""" + name = "flow" + + def __init__(self): + self.data_dir = BACKEND_PATHS["flow"] + self.data_dir.mkdir(parents=True, exist_ok=True) + self._file = self.data_dir / "memory.json" + + def _load_all(self) -> dict: + if self._file.exists(): + return json.loads(self._file.read_text(encoding="utf-8")) + return {} + + def _save_all(self, data: dict): + self._file.write_text(json.dumps(data, indent=2, ensure_ascii=False), encoding="utf-8") + + def save(self, record: dict): + data = self._load_all() + data[record["key"]] = { + "value": record["value"], + "tags": record.get("tags", []), + "timestamp": record["timestamp"], + } + self._save_all(data) + + def load(self, key: str): + data = self._load_all() + entry = data.get(key) + return entry["value"] if entry else None + + def search(self, query: str, limit: int = 10): + data = self._load_all() + results = [] + for key, entry in data.items(): + if query.lower() in key.lower() or query.lower() in str(entry).lower(): + results.append({"key": key, "value": entry["value"], "backend": "flow"}) + if len(results) >= limit: + break + return results + + def list_keys(self): + return list(self._load_all().keys()) + + def forget(self, key: str) -> bool: + data = self._load_all() + if key in data: + del data[key] + self._save_all(data) + return True + return False + + def stats(self): + data = self._load_all() + return {"status": "active", "entries": len(data), "file": str(self._file)} + + +class TokenSaverBackend(BaseBackend): + """Token Saver 记忆 — 压缩历史""" + name = "tokensave" + + def __init__(self): + self.db_path = BACKEND_PATHS["tokensave"] + + def _connect(self): + try: + return sqlite3.connect(str(self.db_path)) + except Exception: + return None + + def save(self, record: dict): + conn = self._connect() + if not conn: + return + try: + conn.execute( + "INSERT OR REPLACE INTO memory (key, value) VALUES (?, ?)", + (record["key"], record["value"]), + ) + conn.commit() + except sqlite3.OperationalError: + pass + conn.close() + + def load(self, key: str): + conn = self._connect() + if not conn: + return None + try: + cur = conn.execute("SELECT value FROM memory WHERE key=?", (key,)) + row = cur.fetchone() + conn.close() + return row[0] if row else None + except sqlite3.OperationalError: + conn.close() + return None + + def search(self, query: str, limit: int = 10): + return [] # token_saver DB 结构不确定,跳过搜索 + + def list_keys(self): + conn = self._connect() + if not conn: + return [] + try: + cur = conn.execute("SELECT key FROM memory") + keys = [r[0] for r in cur.fetchall()] + conn.close() + return keys + except sqlite3.OperationalError: + conn.close() + return [] + + def forget(self, key: str) -> bool: + conn = self._connect() + if not conn: + return False + try: + conn.execute("DELETE FROM memory WHERE key=?", (key,)) + conn.commit() + affected = conn.total_changes + conn.close() + return affected > 0 + except sqlite3.OperationalError: + conn.close() + return False + + def stats(self): + return {"status": "readonly (token saver internal)", "note": "use for compression history only"} + + +class PlanFilesBackend(BaseBackend): + """Planning-with-files — task_plan.md / progress.md / findings.md""" + name = "plan" + + def __init__(self): + self.root = BACKEND_PATHS["plan_dir"] + + def save(self, record: dict): + key = record["key"] + value = record["value"] + tags = record.get("tags", []) + + if "progress" in tags or "task" in tags: + # 追加到 progress.md + path = self.root / "progress.md" + timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S") + entry = f"\n### {key} — {timestamp}\n{value}\n" + path.write_text(entry, encoding="utf-8") if not path.exists() else open(path, "a", encoding="utf-8").write(entry) + + elif "plan" in tags: + # 更新 task_plan.md + path = self.root / "task_plan.md" + if not path.exists(): + path.write_text(f"# {key}\n\n{value}\n", encoding="utf-8") + + else: + # 存 findings.md + path = self.root / "findings.md" + timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S") + entry = f"\n## {key} ({timestamp})\n{value}\n" + path.write_text(entry, encoding="utf-8") if not path.exists() else open(path, "a", encoding="utf-8").write(entry) + + def load(self, key: str): + for fname in ["task_plan.md", "progress.md", "findings.md"]: + path = self.root / fname + if path.exists(): + content = path.read_text(encoding="utf-8") + if key in content: + return content + return None + + def search(self, query: str, limit: int = 10): + results = [] + for fname in ["task_plan.md", "progress.md", "findings.md"]: + path = self.root / fname + if not path.exists(): + continue + content = path.read_text(encoding="utf-8") + if query.lower() in content.lower(): + # 提取匹配段落 + for line in content.split("\n"): + if query.lower() in line.lower(): + results.append({ + "key": fname, + "value": line.strip()[:200], + "backend": "plan", + }) + if len(results) >= limit: + break + return results + + def list_keys(self): + keys = [] + for fname in ["task_plan.md", "progress.md", "findings.md"]: + path = self.root / fname + if path.exists(): + keys.append(fname) + return keys + + def forget(self, key: str) -> bool: + return False # 不能直接删除 plan 文件中的内容 + + def stats(self): + files = {} + for fname in ["task_plan.md", "progress.md", "findings.md"]: + path = self.root / fname + if path.exists(): + files[fname] = f"{path.stat().st_size / 1024:.1f}KB" + return {"status": "active", "files": files} + + +# ═══════════════════════════════════════════════════ +# CLI 入口 +# ═══════════════════════════════════════════════════ + +def main(): + import argparse + parser = argparse.ArgumentParser(description="DeepCode Unified Memory Manager") + parser.add_argument("command", choices=["save", "load", "search", "list", "forget", "stats", "seed"], + help="操作") + parser.add_argument("key", nargs="?", help="记忆键名") + parser.add_argument("value", nargs="?", help="记忆值 (仅 save)") + parser.add_argument("--backend", default="auto", help="后端 (auto/ruflo/claude/flow/tokensave/plan)") + parser.add_argument("--tags", nargs="*", default=[], help="标签 (仅 save)") + parser.add_argument("--limit", type=int, default=10, help="搜索限制 (仅 search)") + parser.add_argument("--json", action="store_true", help="JSON 输出") + args = parser.parse_args() + + mm = MemoryManager() + + if args.command == "save": + if not args.key or args.value is None: + print("Usage: memory_manager.py save [--tags ...] [--backend ...]") + sys.exit(1) + result = mm.save(args.key, args.value, backend=args.backend, tags=args.tags) + print(f"[memory_manager] Saved: {result}" if not args.json else json.dumps({"saved": result})) + + elif args.command == "load": + if not args.key: + print("Usage: memory_manager.py load ") + sys.exit(1) + result = mm.load(args.key, backend=args.backend) + if args.json: + print(json.dumps({"key": args.key, "value": result})) + else: + print(f"Value: {result}" if result else f"Not found: {args.key}") + + elif args.command == "search": + if not args.key: + print("Usage: memory_manager.py search ") + sys.exit(1) + results = mm.search(args.key, backend=args.backend, limit=args.limit) + if args.json: + print(json.dumps(results, indent=2)) + else: + print(f"Found {len(results)} results:") + for r in results: + print(f" [{r['backend']}] {r['key']}: {str(r['value'])[:80]}") + + elif args.command == "list": + keys = mm.list(backend=args.backend) + if args.json: + print(json.dumps(keys)) + else: + print(f"Memory keys ({len(keys)}):") + for k in keys: + print(f" - {k}") + + elif args.command == "forget": + if not args.key: + print("Usage: memory_manager.py forget ") + sys.exit(1) + result = mm.forget(args.key, backend=args.backend) + print(f"[memory_manager] Deleted: {result}") + + elif args.command == "stats": + stats = mm.stats() + if args.json: + print(json.dumps(stats, indent=2)) + else: + print("Memory Backend Stats:") + print(f"{'Backend':<12} {'Status':<20} {'Entries':<10}") + print("-" * 45) + for name, s in stats.items(): + entries = s.get("entries", s.get("files", "N/A")) + status = s.get("status", "?") + print(f"{name:<12} {str(status):<20} {str(entries):<10}") + + elif args.command == "seed": + count = mm.seed() + print(f"[memory_manager] Seeded {count} memories") + if args.json: + print(json.dumps({"seeded": count})) + + +if __name__ == "__main__": + import sys + main() diff --git a/.deepcode/skills/deepcode-knowledge/plugin.json b/.deepcode/skills/deepcode-knowledge/plugin.json new file mode 100644 index 00000000..db5305a1 --- /dev/null +++ b/.deepcode/skills/deepcode-knowledge/plugin.json @@ -0,0 +1,13 @@ +{ + "name": "deepcode-knowledge", + "version": "1.0.0", + "description": "DeepCode Knowledge Base — Obsidian 风格笔记 + KV 记忆系统", + "author": "DeepCode + Ghidra RE (Claude Code v2.1.216)", + "entry": "SKILL.md", + "dependencies": [], + "modules": [ + "knowledge_server.py", + "memory_manager.py", + "rebuild_vault_index.py" + ] +} diff --git a/.deepcode/skills/deepcode-knowledge/rebuild_vault_index.py b/.deepcode/skills/deepcode-knowledge/rebuild_vault_index.py new file mode 100644 index 00000000..a54cde84 --- /dev/null +++ b/.deepcode/skills/deepcode-knowledge/rebuild_vault_index.py @@ -0,0 +1,207 @@ +#!/usr/bin/env python3 +"""Vault 索引重建 — 扫描磁盘 .md 文件批量导入 SQLite""" +import json, sys, os, re, time, sqlite3 +from datetime import datetime, date +from pathlib import Path +from collections import defaultdict + +# ── 路径 ── +SKILL_DIR = Path(__file__).parent +NOTES_DIR = SKILL_DIR / "data" / "vault" / "notes" +DAILY_DIR = SKILL_DIR / "data" / "vault" / "daily" +DB_PATH = SKILL_DIR / "data" / "knowledge.db" + +# ── 工具函数 (复用 knowledge_server 的逻辑) ── +def now_ts(): return time.time() + +def extract_frontmatter(text): + m = re.match(r'^---\s*\n(.*?)\n---\s*\n', text, re.DOTALL) + if not m: return {}, text + fm = {} + for line in m.group(1).strip().split('\n'): + if ':' in line: + k, v = line.split(':', 1) + fm[k.strip()] = v.strip().strip('"').strip("'") + return fm, text[m.end():] + +def extract_tags(text): + return re.findall(r'#([\w\u4e00-\u9fff\-.]+)', text) + +def extract_wikilinks(text): + return re.findall(r'\[\[([^\]]+)\]\]', text) + +def slugify(text): + s = re.sub(r'[^\w\s-]', '', text).strip().lower() + return re.sub(r'[-\s]+', '-', s)[:60] + +def get_db(): + conn = sqlite3.connect(str(DB_PATH)) + conn.row_factory = sqlite3.Row + conn.executescript(""" + CREATE TABLE IF NOT EXISTS notes ( + id TEXT PRIMARY KEY, path TEXT UNIQUE NOT NULL, + title TEXT NOT NULL, type TEXT DEFAULT 'note', + tags TEXT DEFAULT '', symbol TEXT DEFAULT '', + created_at REAL NOT NULL, updated_at REAL NOT NULL, + content_preview TEXT DEFAULT '' + ); + CREATE TABLE IF NOT EXISTS links ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + source_id TEXT NOT NULL, target_id TEXT NOT NULL, + link_type TEXT DEFAULT 'related', + UNIQUE(source_id, target_id) + ); + CREATE TABLE IF NOT EXISTS daily_logs ( + id TEXT PRIMARY KEY, date TEXT UNIQUE NOT NULL, + summary TEXT DEFAULT '', mood TEXT DEFAULT '', + created_at REAL NOT NULL + ); + CREATE INDEX IF NOT EXISTS idx_notes_type ON notes(type); + CREATE INDEX IF NOT EXISTS idx_notes_symbol ON notes(symbol); + CREATE INDEX IF NOT EXISTS idx_notes_tags ON notes(tags); + CREATE INDEX IF NOT EXISTS idx_links_source ON links(source_id); + CREATE INDEX IF NOT EXISTS idx_links_target ON links(target_id); + """) + conn.commit() + return conn + +def parse_tags(raw_tags): + """处理 tags 字段,支持 [tag1, tag2] 列表格式和逗号分隔""" + if not raw_tags: + return "" + raw = raw_tags.strip() + # 处理 [tag1, tag2, ...] 格式 + if raw.startswith('[') and raw.endswith(']'): + inner = raw[1:-1] + tags = [t.strip().strip('"').strip("'") for t in inner.split(',') if t.strip()] + return ','.join(tags) + return raw + +# ══════════════════════════════════════════ +# 主流程 +# ══════════════════════════════════════════ + +def rebuild(): + conn = get_db() + ts_now = now_ts() + + # 清空旧索引 + conn.execute("DELETE FROM links") + conn.execute("DELETE FROM notes") + conn.execute("DELETE FROM daily_logs") + conn.commit() + + stats = {"notes": 0, "dailies": 0, "links": 0, "errors": 0} + + # ── 扫描 notes ── + if NOTES_DIR.exists(): + for fpath in sorted(NOTES_DIR.glob("*.md")): + try: + raw = fpath.read_text(encoding="utf-8") + fm, body = extract_frontmatter(raw) + + title = fm.get("title", fpath.stem) + note_type = fm.get("type", "note") + symbol = fm.get("symbol", "") + tags = parse_tags(fm.get("tags", "")) + created = fm.get("created", str(date.today())) + updated = fm.get("updated", created) + + # 转换日期到时间戳 + try: + created_ts = datetime.strptime(created, "%Y-%m-%d").timestamp() + except: + created_ts = ts_now + try: + updated_ts = datetime.strptime(updated, "%Y-%m-%d").timestamp() + except: + updated_ts = ts_now + + note_id = slugify(title) + "_" + str(int(ts_now))[-6:] + preview = body[:200].strip() + + # 额外从 body 提取标签 + body_tags = extract_tags(body) + all_tags = set(tags.split(",") if tags else []) + all_tags.update(body_tags) + all_tags.discard("") + tags_str = ",".join(sorted(all_tags)) + + conn.execute( + """INSERT OR REPLACE INTO notes + (id, path, title, type, tags, symbol, created_at, updated_at, content_preview) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)""", + (note_id, str(fpath), title, note_type, tags_str, symbol, + created_ts, updated_ts, preview) + ) + + # 提取 wikilinks + links = extract_wikilinks(body) + for link_title in links: + target = conn.execute( + "SELECT id FROM notes WHERE title LIKE ? LIMIT 1", + (f"%{link_title}%",) + ).fetchone() + if target: + try: + conn.execute( + "INSERT OR IGNORE INTO links (source_id, target_id, link_type) VALUES (?, ?, ?)", + (note_id, target["id"], "related") + ) + stats["links"] += 1 + except: + pass + + stats["notes"] += 1 + print(f" ✓ [{note_type:6s}] {title}") + + except Exception as e: + stats["errors"] += 1 + print(f" ✗ ERROR: {fpath.name} → {e}") + + # ── 扫描 daily ── + if DAILY_DIR.exists(): + for fpath in sorted(DAILY_DIR.glob("*.md")): + try: + raw = fpath.read_text(encoding="utf-8") + fm, body = extract_frontmatter(raw) + + d = fm.get("date", fpath.stem.replace("daily_", "")) + summary = fm.get("title", "") + mood = fm.get("mood", "") + + conn.execute( + """INSERT OR REPLACE INTO daily_logs (id, date, summary, mood, created_at) + VALUES (?, ?, ?, ?, ?)""", + (f"daily_{d}", d, summary, mood, ts_now) + ) + stats["dailies"] += 1 + print(f" ✓ [daily ] {d} — {summary}") + + except Exception as e: + stats["errors"] += 1 + print(f" ✗ ERROR: {fpath.name} → {e}") + + conn.commit() + conn.close() + + return stats + +if __name__ == "__main__": + print("=" * 50) + print(" Vault 索引重建") + print(f" 笔记目录: {NOTES_DIR}") + print(f" 日志目录: {DAILY_DIR}") + print(f" 数据库: {DB_PATH}") + print("=" * 50) + + stats = rebuild() + + print() + print("=" * 50) + print(f" 重建完成!") + print(f" ✅ 笔记: {stats['notes']} 篇") + print(f" ✅ 日志: {stats['dailies']} 篇") + print(f" 🔗 链接: {stats['links']} 条") + print(f" ❌ 错误: {stats['errors']}") + print("=" * 50) diff --git a/.deepcode/skills/deepcode-streaming/SKILL.md b/.deepcode/skills/deepcode-streaming/SKILL.md new file mode 100644 index 00000000..d2ea9ba4 --- /dev/null +++ b/.deepcode/skills/deepcode-streaming/SKILL.md @@ -0,0 +1,91 @@ +--- +name: deepcode-streaming +description: > + DeepCode Streaming API — 移植自 Claude Code v2.1.216 的流式 API 支持。 + SSE (Server-Sent Events) / Chunked / WebSocket 流式传输。 + 对标 Claude Code: calculateNonstreamingTimeout, buildRequest, + buildHeaders, buildBody, shouldRetry 等。 +version: 1.0.0 +author: DeepCode + Ghidra RE (Claude Code v2.1.216) +date: 2026-07-26 +tags: [streaming, sse, api, llm] +--- + +# DeepCode Streaming API + +移植自 **Claude Code v2.1.216** 的流式 API 能力。 + +## 对标项 + +| Claude Code | DeepCode Streaming | +|:-----------|:-------------------| +| `_calculateNonstreamingTimeout` | `StreamingClient._calculate_timeout()` | +| `buildRequest` | `_build_headers()` + `_build_body()` | +| `buildHeaders` | `_build_headers()` | +| `buildBody` | `_build_body()` | +| `shouldRetry` | 指数退避重试逻辑 | +| `retryRequest` | while 循环 + backoff | +| SSE 解析 | `_parse_sse_line()` + `_extract_delta()` | + +## 用法 + +### CLI + +```bash +# 流式输出 +python streaming_api.py chat --prompt "用Python写个快排" --stream + +# 非流式 +python streaming_api.py chat --prompt "1+1=?" + +# 指定模型 +python streaming_api.py chat --prompt "分析" --model deepseek-reasoner --stream +``` + +### MCP Server + +```json +"deepcode-streaming": { + "command": "python", + "args": [ + "F:/DEEPCODE/.deepcode/skills/deepcode-streaming/streaming_api.py", + "--mcp" + ], + "env": { + "DEEPSEEK_API_KEY": "${DEEPSEEK_API_KEY}" + } +} +``` + +### Python + +```python +from streaming_api import StreamingClient + +client = StreamingClient(api_key="sk-xxx") + +# 流式 +async for event in client.chat_stream( + messages=[{"role":"user","content":"你好"}], + model="deepseek-chat", +): + if event.type == StreamEventType.CHUNK: + print(event.content, end="", flush=True) + elif event.type == StreamEventType.DONE: + print(f"\n用量: {event.usage}") + +# 非流式 +result = await client.chat(messages) +print(result["choices"][0]["message"]["content"]) + +# 用量统计 +print(client.get_usage_stats()) +``` + +## 特性 + +- **双模超时**: 流式 (300s) vs 非流式 (60s) 自动切换 +- **自动重试**: 指数退避 (2s/4s/8s...),最多 3 次 +- **背压通知**: 重试时发送 THROTTLE 事件 +- **Token 统计**: 自动累计 prompt/completion token +- **Delta 增量**: 纯增量输出,适合实时展示 diff --git a/.deepcode/skills/deepcode-streaming/plugin.json b/.deepcode/skills/deepcode-streaming/plugin.json new file mode 100644 index 00000000..61786964 --- /dev/null +++ b/.deepcode/skills/deepcode-streaming/plugin.json @@ -0,0 +1,11 @@ +{ + "name": "deepcode-streaming", + "version": "1.0.0", + "description": "DeepCode Streaming API — SSE 流式 API 客户端,支持流式 Chat 补全", + "author": "DeepCode + Ghidra RE (Claude Code v2.1.216)", + "entry": "SKILL.md", + "dependencies": [], + "modules": [ + "streaming_api.py" + ] +} diff --git a/.deepcode/skills/deepcode-streaming/streaming_api.py b/.deepcode/skills/deepcode-streaming/streaming_api.py new file mode 100644 index 00000000..71c23e8a --- /dev/null +++ b/.deepcode/skills/deepcode-streaming/streaming_api.py @@ -0,0 +1,514 @@ +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- +""" +DeepCode Streaming API — Claude Code v2.1.216 Streaming API 移植 +══════════════════════════════════════════════════════════════════ +Streaming (流式) API 模块,支持 SSE / Chunked / WebSocket 流。 + +移植自 Claude Code streaming 支持: + - _calculateNonstreamingTimeout / calculateNonstreamingTimeout + - buildRequest / buildHeaders / buildBody (流式版本) + - shouldRetry / retryRequest (流式重试) + - SSE 事件解析 + Chunked 传输 + +用法: + # CLI 流式调用 DeepSeek + python streaming_api.py chat --prompt "你好" --stream + + # MCP Server 模式 + python streaming_api.py --mcp + + # Python 嵌入 + from streaming_api import StreamingClient, StreamEvent + client = StreamingClient(api_key="sk-xxx") + async for chunk in client.chat_stream([{"role":"user","content":"hello"}]): + print(chunk.content, end="") +""" + +import asyncio +import json +import os +import sys +import time +import uuid +from dataclasses import dataclass, field +from datetime import datetime +from enum import Enum +from typing import AsyncIterator, Dict, List, Optional, Callable, Any, Union + +# ── 流事件类型 ──────────────────────────────────────────── + +class StreamEventType(Enum): + """流事件类型""" + CHUNK = "chunk" # 普通内容块 + DONE = "done" # 完成 + ERROR = "error" # 错误 + THROTTLE = "throttle" # 限流 (背压) + METADATA = "metadata" # 元数据 (token用量等) + + +@dataclass +class StreamEvent: + """流事件 — 对标 Claude Code stream chunk""" + type: StreamEventType + content: Optional[str] = None + delta: Optional[str] = None # 增量内容 (SSE) + finish_reason: Optional[str] = None + usage: Optional[Dict] = None # token 用量 + error: Optional[str] = None + metadata: Optional[Dict] = None + timestamp: float = field(default_factory=time.time) + + @classmethod + def chunk(cls, content: str, delta: str = None): + return cls(type=StreamEventType.CHUNK, content=content, delta=delta) + + @classmethod + def done(cls, reason: str = "stop", usage: Dict = None): + return cls(type=StreamEventType.DONE, finish_reason=reason, usage=usage) + + @classmethod + def error(cls, msg: str): + return cls(type=StreamEventType.ERROR, error=msg) + + +# ── 流式客户端 —──────────────────────────────────────────── + +class StreamingClient: + """ + 流式 API 客户端 — 对标 Claude Code streaming request 栈 + + 特性: + - SSE (Server-Sent Events) 解析 + - Chunked transfer 支持 + - 超时管理 (streaming vs non-streaming) + - 自动重试 (shouldRetry) + - 背压控制 (throttle) + - Token 用量跟踪 + """ + + def __init__( + self, + api_key: str = None, + base_url: str = "https://api.deepseek.com", + timeout: int = 60, + stream_timeout: int = 300, + max_retries: int = 3, + ): + self.api_key = api_key or os.environ.get("DEEPSEEK_API_KEY", "") + self.base_url = base_url.rstrip("/") + # 超时: 对标 Claude Code calculateNonstreamingTimeout / _calculateNonstreamingTimeout + self.timeout = timeout # 非流式超时 (short) + self.stream_timeout = stream_timeout # 流式超时 (long) + self.max_retries = max_retries + self._total_prompt_tokens = 0 + self._total_completion_tokens = 0 + + def _build_headers(self, stream: bool = True) -> Dict[str, str]: + """构建请求头 — 对标 Claude Code buildHeaders""" + return { + "Authorization": f"Bearer {self.api_key}", + "Content-Type": "application/json", + "Accept": "text/event-stream" if stream else "application/json", + "User-Agent": "DeepCode-Streaming/1.0", + } + + def _build_body( + self, + messages: List[Dict], + model: str = "deepseek-v4-flash", + stream: bool = True, + temperature: float = 0.6, + max_tokens: int = 4096, + **kwargs, + ) -> bytes: + """构建请求体 — 对标 Claude Code buildBody""" + body = { + "model": model, + "messages": messages, + "stream": stream, + "temperature": temperature, + "max_tokens": max_tokens, + } + body.update(kwargs) + return json.dumps(body).encode("utf-8") + + @staticmethod + def _calculate_timeout(streaming: bool, non_streaming_val: int, streaming_val: int) -> int: + """计算超时 — 对标 Claude Code _calculateNonstreamingTimeout""" + return streaming_val if streaming else non_streaming_val + + # ── SSE 解析 ─────────────────────────────────────── + + @staticmethod + def _parse_sse_line(line: str) -> Optional[Dict]: + """解析 SSE 行 — 对标 Claude Code SSE parser""" + line = line.strip() + if not line or line.startswith(":"): + return None # 注释 / 空行 + if line.startswith("data: "): + data = line[6:] + if data == "[DONE]": + return {"type": "done"} + try: + return json.loads(data) + except json.JSONDecodeError: + return {"type": "data", "raw": data} + return None + + @staticmethod + def _extract_delta(sse_data: Dict) -> Optional[str]: + """从 SSE 数据中提取增量内容 — 对标 Claude Code delta extraction""" + choices = sse_data.get("choices", []) + if not choices: + return None + delta = choices[0].get("delta", {}) + return delta.get("content", "") + + @staticmethod + def _extract_finish_reason(sse_data: Dict) -> Optional[str]: + choices = sse_data.get("choices", []) + if choices: + return choices[0].get("finish_reason") + return None + + @staticmethod + def _extract_usage(sse_data: Dict) -> Optional[Dict]: + return sse_data.get("usage") + + # ── 核心流式请求 ──────────────────────────────────── + + async def chat_stream( + self, + messages: List[Dict], + model: str = "deepseek-v4-flash", + temperature: float = 0.6, + max_tokens: int = 4096, + **kwargs, + ) -> AsyncIterator[StreamEvent]: + """ + 流式 Chat 补全 — 对标 Claude Code streaming makeRequest + + Args: + messages: [{"role":"user","content":"..."}] + model: 模型名 + temperature: 温度 + max_tokens: 最大输出 token + + Yields: + StreamEvent: 流事件 (chunk / done / error) + """ + url = f"{self.base_url}/v1/chat/completions" + body = self._build_body(messages, model, True, temperature, max_tokens, **kwargs) + headers = self._build_headers(stream=True) + actual_timeout = self._calculate_timeout( + True, self.timeout, self.stream_timeout + ) + + retries = 0 + last_error = None + + while retries <= self.max_retries: + try: + reader, writer = await asyncio.wait_for( + self._connect(url, headers, body), + timeout=self.timeout, + ) + buffer = "" + async for line in self._read_lines(reader, writer, actual_timeout): + buffer += line + if buffer.endswith("\n\n"): + for sse_line in buffer.strip().split("\n"): + parsed = self._parse_sse_line(sse_line) + if parsed is None: + continue + if parsed.get("type") == "done": + yield StreamEvent.done() + writer.close() + return + delta = self._extract_delta(parsed) + finish = self._extract_finish_reason(parsed) + usage = self._extract_usage(parsed) + if delta: + yield StreamEvent.chunk(content=delta, delta=delta) + if finish: + yield StreamEvent.done(reason=finish, usage=usage) + writer.close() + return + if usage: + self._total_prompt_tokens += usage.get("prompt_tokens", 0) + self._total_completion_tokens += usage.get("completion_tokens", 0) + buffer = "" + + # 读完缓冲区 + if buffer.strip(): + for sse_line in buffer.strip().split("\n"): + parsed = self._parse_sse_line(sse_line) + if parsed and parsed.get("type") == "done": + yield StreamEvent.done() + writer.close() + return + + except asyncio.TimeoutError: + last_error = "timeout" + retries += 1 + if retries > self.max_retries: + yield StreamEvent.error(f"Stream timed out after {actual_timeout}s") + return + # 指数退避 — 对标 Claude Code shouldRetry + wait = min(2 ** retries, 30) + yield StreamEvent( + type=StreamEventType.THROTTLE, + content=f"retry in {wait}s (attempt {retries}/{self.max_retries})", + ) + await asyncio.sleep(wait) + + except Exception as e: + last_error = str(e) + retries += 1 + if retries > self.max_retries: + yield StreamEvent.error(f"Stream error after {retries} retries: {e}") + return + wait = min(2 ** retries, 30) + await asyncio.sleep(wait) + + async def _connect(self, url: str, headers: Dict, body: bytes): + """建立 HTTP 连接""" + import http.client + parsed = url.replace("https://", "").replace("http://", "") + use_ssl = url.startswith("https") + host = parsed.split("/")[0] + path = "/" + "/".join(parsed.split("/")[1:]) + + conn = http.client.HTTPSConnection(host, timeout=self.timeout) if use_ssl \ + else http.client.HTTPConnection(host, timeout=self.timeout) + conn.request("POST", path, body=body, headers=headers) + resp = conn.getresponse() + return resp, conn + + async def _read_lines(self, reader, writer, timeout): + """逐行读取流式响应""" + import select + import socket + + async def read_one_line(): + loop = asyncio.get_event_loop() + return await asyncio.wait_for( + loop.run_in_executor(None, reader.readline), + timeout=timeout, + ) + + while True: + try: + line = await read_one_line() + if not line: + break + decoded = line.decode("utf-8", errors="replace") + yield decoded + except asyncio.TimeoutError: + raise + except Exception: + break + + # ── 非流式请求 ────────────────────────────────────── + + async def chat( + self, + messages: List[Dict], + model: str = "deepseek-v4-flash", + temperature: float = 0.6, + max_tokens: int = 4096, + **kwargs, + ) -> Dict: + """ + 非流式 Chat 补全 — 对标 Claude Code non-streaming makeRequest + + 使用较短的超时 (calculateNonstreamingTimeout). + """ + url = f"{self.base_url}/v1/chat/completions" + body = self._build_body(messages, model, False, temperature, max_tokens, **kwargs) + headers = self._build_headers(stream=False) + actual_timeout = self._calculate_timeout(False, self.timeout, self.stream_timeout) + + import urllib.request + req = urllib.request.Request(url, data=body, headers=headers, method="POST") + try: + with urllib.request.urlopen(req, timeout=actual_timeout) as resp: + data = json.loads(resp.read().decode()) + self._total_prompt_tokens += data.get("usage", {}).get("prompt_tokens", 0) + self._total_completion_tokens += data.get("usage", {}).get("completion_tokens", 0) + return data + except Exception as e: + return {"error": str(e)} + + # ── 工具函数 ──────────────────────────────────────── + + def get_usage_stats(self) -> Dict: + """获取累计 token 用量统计""" + return { + "total_prompt_tokens": self._total_prompt_tokens, + "total_completion_tokens": self._total_completion_tokens, + "total_tokens": self._total_prompt_tokens + self._total_completion_tokens, + } + + def reset_usage(self): + self._total_prompt_tokens = 0 + self._total_completion_tokens = 0 + + +# ── 辅助: 组装流为完整响应 ────────────────────────────── + +async def stream_to_completion(stream: AsyncIterator[StreamEvent]) -> str: + """将流式事件组装为完整字符串""" + result = [] + async for event in stream: + if event.type == StreamEventType.CHUNK and event.content: + result.append(event.content) + elif event.type == StreamEventType.ERROR: + return f"Error: {event.error}" + return "".join(result) + + +# ── CLI 入口 ────────────────────────────────────────────── + +async def main_cli(): + import argparse + parser = argparse.ArgumentParser(description="DeepCode Streaming API") + parser.add_argument("--mcp", action="store_true", help="MCP Server 模式") + sub = parser.add_subparsers(dest="mode") + + chat_parser = sub.add_parser("chat", help="Chat 补全") + chat_parser.add_argument("--prompt", required=True, help="用户输入") + chat_parser.add_argument("--model", default="deepseek-v4-flash") + chat_parser.add_argument("--stream", action="store_true", help="启用流式输出") + chat_parser.add_argument("--temperature", type=float, default=0.6) + chat_parser.add_argument("--max-tokens", type=int, default=4096) + + args = parser.parse_args() + + client = StreamingClient() + + if args.mode == "chat": + messages = [{"role": "user", "content": args.prompt}] + if args.stream: + print(f"[streaming] 模型: {args.model}", file=sys.stderr) + async for event in client.chat_stream( + messages, model=args.model, + temperature=args.temperature, + max_tokens=args.max_tokens, + ): + if event.type == StreamEventType.CHUNK: + print(event.content or "", end="", flush=True) + elif event.type == StreamEventType.DONE: + print() + if event.usage: + print(f"\n[usage] {event.usage}", file=sys.stderr) + elif event.type == StreamEventType.ERROR: + print(f"\n[error] {event.error}", file=sys.stderr) + elif event.type == StreamEventType.THROTTLE: + print(f"\n[throttle] {event.content}", file=sys.stderr) + print(f"\n[stats] {client.get_usage_stats()}", file=sys.stderr) + else: + result = await client.chat(messages, model=args.model) + print(json.dumps(result, ensure_ascii=False, indent=2)) + elif args.mcp: + await run_mcp() + else: + parser.print_help() + + +def _mcp_send(obj): + print(json.dumps(obj, ensure_ascii=False), flush=True) + + +async def run_mcp(): + """MCP Server 模式(标准 MCP 协议)""" + client = StreamingClient() + + for line in sys.stdin: + line = line.strip() + if not line: + continue + try: + req = json.loads(line) + method = req.get("method", "") + params = req.get("params", {}) + req_id = req.get("id", "") + + # 标准 MCP 握手 + if method == "initialize": + _mcp_send({ + "jsonrpc": "2.0", "id": req_id, + "result": { + "protocolVersion": "2024-11-05", + "capabilities": {"tools": {"listChanged": True}}, + "serverInfo": {"name": "deepcode-streaming", "version": "1.0.0"}, + }, + }) + elif method == "notifications/initialized": + pass + elif method == "tools/list": + _mcp_send({ + "jsonrpc": "2.0", "id": req_id, + "result": { + "tools": [ + { + "name": "streaming_chat", + "description": "流式 Chat 补全", + "inputSchema": { + "type": "object", + "properties": { + "prompt": {"type": "string"}, + "system": {"type": "string"}, + "model": {"type": "string"}, + "temperature": {"type": "number"}, + "max_tokens": {"type": "integer"}, + "stream": {"type": "boolean"}, + }, + "required": ["prompt"], + }, + }, + ], + }, + }) + elif method == "tools/call": + name = params.get("name", "") + args_dict = params.get("arguments", {}) + if name == "streaming_chat": + messages = [] + if args_dict.get("system"): + messages.append({"role": "system", "content": args_dict["system"]}) + messages.append({"role": "user", "content": args_dict["prompt"]}) + stream = args_dict.get("stream", True) + if stream: + full_text = "" + async for event in client.chat_stream( + messages, + model=args_dict.get("model", "deepseek-v4-flash"), + temperature=args_dict.get("temperature", 0.6), + max_tokens=args_dict.get("max_tokens", 4096), + ): + if event.type == StreamEventType.CHUNK: + full_text += event.content or "" + elif event.type == StreamEventType.ERROR: + full_text += f"\n[error] {event.error}" + result = {"response": full_text} + else: + chat_kw = dict(args_dict) + chat_kw.pop('stream', None) + chat_kw.pop('prompt', None) + chat_kw.pop('system', None) + resp = await client.chat(messages, **chat_kw) + result = {"response": resp.get("choices", [{}])[0].get("message", {}).get("content", "")} + else: + result = {"error": f"Unknown: {name}"} + _mcp_send({ + "jsonrpc": "2.0", "id": req_id, + "result": {"content": [{"type": "text", "text": json.dumps(result, ensure_ascii=False)}]}, + }) + except json.JSONDecodeError: + pass + + +if __name__ == "__main__": + asyncio.run(main_cli())