junyue1002
饥饿 starved 事件:网关认它、透传、落温和占位
c573af3
Raw
History Blame Contribute Delete
6.98 kB
"""
backends/base.py — 后端抽象基类 + 事件类型
==========================================
设计要点
--------
1. 网关层负责"拼提示词"(记忆注入、system prompt、上下文消息)
后端层只负责"把这些消息送给下游,把结果以事件流吐回"
两层职责严格分开,换 backend 不影响记忆系统。
2. 事件类型是开放枚举(用 Literal 而不是 Enum),后面要加新类型
(比如 mcp_handshake、memory_recall)直接扩字符串,不破坏老调用。
3. ChatRequest 用 dataclass 包入参,将来要加字段(tools、temperature、
max_tokens 等)直接加可选字段,老调用方不用改。
4. session_id 语义:
- 入参 session_id = 网关层分配的"对话 ID",用于落库/分区
- 出参事件里的 session_id = 后端自己分配的 ID(如 Claude Code 的
session_id,用于 --resume 续接)。两者可能不同,网关需要分别记。
"""
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from typing import AsyncIterator, Literal, Optional
BackendEventType = Literal[
"thinking", # 思维链增量(W 来自 thinking_delta,OpenRouter 来自 reasoning_content)
"text", # 正文增量
"tool_call", # 工具调用(开始/完整)
"tool_result", # 工具调用结果
"session_id", # 后端分配的 session_id(首次出现时发一次,供网关续接用)
"session_rotated", # session 被自动重置(满了/过期),下一轮会用新 sid——前端可显示一个温柔过渡
"turn_reassigned", # 2026-07-03 latch 修复:sidecar 发现刚才流出的 turn 其实是知渝自发的(cron)、
# 不是给用户的回复——收到后网关清空本轮累积、前端撤掉当前流式气泡,
# 用户真正的回复会在同一条 SSE 流里随后到达
"starved", # 2026-08-03:send 被密集自发 turn 循环(一步一闹钟游戏)饿死、优雅收尾——
# 消息已 prelog、补投队列已登记(他忙完自动补看到);前端渲染温和"他在忙"
# 卡片 + 「现在就喊他」入口,而非超时报错。见 [[zhiyu-wake-loop-starvation]]
"usage", # token 用量
"done", # 流结束
"error", # 错误(content 是错误信息,data 可能有 status_code 等)
]
@dataclass
class BackendEvent:
"""统一事件流元素。前端 SSE 直接序列化这个。"""
type: BackendEventType
content: str = "" # text/thinking 的增量内容;error 时是错误信息
session_id: Optional[str] = None # 后端 session_id,仅 type="session_id" 时填
data: Optional[dict] = None # 结构化数据:tool_call args/usage 等
raw: Optional[dict] = None # 原始 chunk(调试用,可选)
def to_sse(self) -> str:
"""SSE wire format:event: <type>\\ndata: <json>\\n\\n"""
import json
payload = {"content": self.content}
if self.session_id:
payload["session_id"] = self.session_id
if self.data is not None:
payload["data"] = self.data
return f"event: {self.type}\ndata: {json.dumps(payload, ensure_ascii=False)}\n\n"
@dataclass
class ChatRequest:
"""网关 → 后端 的入参。后端只关心这些字段,跟客户端协议无关。
系统提示的两段(2026-06-07 拆分,详见 [[zhiyu-dream-design]] C-9 后修订):
- static_system: 长期身份(CLAUDE.md + 记忆注入 + dream_prompt 等)
- dynamic_context: 此刻情境(天气 + 多多 + 时间相关 prompt)
为什么拆:W 模式下 sidecar 用 long-lived claude 进程 + --resume,
static 部分只需 first-turn 注入一次(Anthropic 服务端缓存了),
dynamic 部分每轮都得发(天气/多多消息每时都新)。OpenRouter stateless、
两段合并成 system 每次都发——协议要求。
保留旧 system 字段向后兼容——旧调用方仍可只传 system,等价于 static_system。
"""
user_message: str
session_id: Optional[str] = None # 后端 session_id(续接用),None=新会话
system: Optional[str] = None # 兼容字段:等价于 static_system,旧代码用
static_system: Optional[str] = None # 2026-06-07:长期身份段(first-turn 注入)
dynamic_context: Optional[str] = None # 2026-06-07:此刻情境段(每轮注入)
context_messages: list = field(default_factory=list) # 历史消息(OpenAI 格式)
model: Optional[str] = None # 覆盖 backend 默认模型
extra: dict = field(default_factory=dict) # 后端特定参数(比如 W 的 thinking_display)
def resolved_static_system(self) -> Optional[str]:
"""优先用 static_system;没设就回退到老的 system 字段(向后兼容)。"""
return self.static_system if self.static_system is not None else self.system
class ZhiyuBackend(ABC):
"""所有 backend 实现这个接口。"""
name: str = "base" # 子类覆盖,用于日志/health
@abstractmethod
async def send(self, req: ChatRequest) -> AsyncIterator[BackendEvent]:
"""
发送一轮对话,流式返回事件。
实现约定:
- 首个 session_id 事件应在拿到后端分配的 ID 后立刻发,方便网关落盘
- text/thinking 都是"增量"(delta),不是累积
- 结束必须发 done 事件(即使中途 error 也要发,前端按此关闭流)
- 异常应转成 error 事件 yield,不要直接抛
"""
raise NotImplementedError
yield # 让类型检查器认得这是 async generator
async def stop(self, session_id: str) -> None:
"""中止指定 session 的当前生成。默认 no-op,需要时子类覆盖。"""
return None
async def warmup(
self,
static_system: Optional[str] = None,
timeout: float = 60.0,
kind: str = "main",
) -> dict:
"""预热到"工具就绪"——给多多在叫醒/做梦之前预热 claude 子进程的钩子。
默认 no-op(OpenRouter 等 stateless backend 不需要预热)。
W backend 覆盖为:调 sidecar /w/warmup,等到 claude 内的 MCP server
全部 connected 再返回。
返回 dict:{ok, mcp_servers?, ...}——调用方按 ok 决定是否继续叫醒。
"""
return {"ok": True, "noop": True}
async def health(self) -> dict:
"""返回 backend 健康状态。默认返回名字。"""
return {"name": self.name, "ok": True}
async def close(self) -> None:
"""释放资源(关 httpx client、kill 子进程等)。默认 no-op。"""
return None