← 代码课堂

第 14 课:流式会话进阶 —— 锁、快照回滚、打断

难度:★★★(需要第 08、09 课基础)
教材:worry_debate_game/worry_debate_game/app.py(294 行)+ nodes.py(355 行)
预计时间:讲解 70 分钟 + 练习 50 分钟
学完你能:把"玩具级流式接口"升级成"能多轮、防并发、失败可重试、还能演剧情"的会话式接口


课前须知:从"一次请求"到"一局游戏"

第 09 课的流式是"发一次请求,收一串事件"。
但真实的一局辩论要持续好几轮,于是冒出一堆新问题:

这一课逐个解决。


第一步:体验这四种情况

cd worry_debate_game
start.bat
  1. 多轮:玩到第二轮,在输入框写一句话再继续——凡人插话会被带进剧情
  2. 打断:多玩几局,注意天使有时说到一半被恶魔"停停停!"打断
  3. 并发:点"下一轮"之后立刻再点一次,观察提示
  4. 失败重试:在设置里把 API Key 改错,触发一次失败,

观察错误提示后还能不能重试同一轮


第二步:模块地图

app.py
├── Session 类                第 61-67 行  一局游戏的状态 + 锁
├── sessions / _sessions_lock  第 70-71 行  所有会话 + 保护它的锁
├── _purge_sessions()          第 74-82 行  过期清理 + 超量淘汰
├── _play_round()              第 123-161 行 一轮演出(天使→打断→恶魔→裁判→结束)
├── /api/start-stream          第 184-243 行 开局(含快照与回滚)
└── /api/next-stream           第 246-294 行 下一轮(含并发锁与回滚)

nodes.py
├── resolve_api_key()          第 32-36 行   Key 的优先级
├── DemoLLM                    第 94-120 行  无 Key 也能离线演示
├── _transcript()              第 176-189 行 把历史拼成上下文
├── _demon_messages()          第 207-225 行 含"打断后第一句要嚣张"的指令
├── _stream_text()             第 240-258 行 流式 + 截断(cut_at)
└── evaluate_round()           第 312-349 行 裁判打分 + 容错

第三步:逐块精讲

块 1:Session 对象 —— 一局游戏装进一个类(第 61-67 行)

class Session:
    def __init__(self, state: GameState, allow_interrupt: bool, judge: bool):
        self.state = state              # 游戏数据(烦恼、历史台词、倾向值…)
        self.allow_interrupt = allow_interrupt   # 这一局允不允许打断
        self.judge = judge              # 要不要裁判打分
        self.touched = time.time()      # 最后一次活动时间(用来清理)
        self.lock = threading.Lock()    # 这一局的"独占锁"

对比第 08 课:那时会话就是裸的 dict;现在包成了一个类,因为要挂配置(打断/裁判)
和运行期资源(锁、活动时间)。状态多了,就值得用对象封装。

块 2:会话生命周期 —— 清理与淘汰(第 57-58、74-82 行)

SESSION_TTL = 2 * 60 * 60  # 2 小时无活动即清理
MAX_SESSIONS = 500

def _purge_sessions():
    now = time.time()
    with _sessions_lock:
        expired = [sid for sid, s in sessions.items() if now - s.touched > SESSION_TTL]
        for sid in expired:
            sessions.pop(sid, None)
        if len(sessions) > MAX_SESSIONS:
            for sid, _ in sorted(sessions.items(), key=lambda kv: kv[1].touched)[: len(sessions) - MAX_SESSIONS]:
                sessions.pop(sid, None)

内存字典当存储的老问题:只进不出会爆内存。这里用两级策略:

  1. TTL(生存时间):2 小时没动过 → 过期删除(now - touched > TTL)
  2. LRU 淘汰:数量超过 500 → 按 touched 从旧到新删到 500

(sorted(...)[:超出数量] 就是把最旧的挑出来)

每次开局前调用一次(第 187 行 _purge_sessions()),
不用定时器也能保持内存不涨——这种"顺手清理"是轻量服务常用的小技巧。

真·生产的做法是把会话放 Redis(自动过期、多进程共享)。
但能说清"内存 + TTL + LRU 的适用边界",本身就是工程判断力。

块 3:并发锁 —— 一轮只许跑一个(第 246-257、288-292 行)

@app.post("/api/next-stream")
def next_round_stream(req: NextRequest, x_api_key: Optional[str] = Header(None)):
    session = sessions.get(req.session_id)

    def generate():
        if session is None:
            yield sse("error", {"message": "会话不存在或已过期,请重新开始", "retryable": False})
            return
        if not session.lock.acquire(blocking=False):     # ← 抢锁,抢不到立刻返回
            yield sse("error", {"message": "上一轮还在进行中", "retryable": True})
            return
        ...
        finally:
            session.lock.release()                       # ← 无论如何都要还锁

为什么不用 blocking=True 排队?因为流式生成可能几十秒,
让第二个请求干等没意义,不如立刻反馈。前端也更好做提示。

块 4:快照回滚 —— 失败也能重试同一轮(第 224-243、258-292 行)

先看开局(第 225-240 行):

snapshot = copy.deepcopy(state)     # 开跑前拍一张"快照"
snapshot["round"] = 0
completed = False
with session.lock:
    try:
        yield from _play_round(session, session_id)
        completed = True
    except Exception as e:
        state.clear()
        state.update(snapshot)      # 回滚:把状态恢复到快照
        completed = True
        yield sse("error", {...})
    finally:
        if not completed:           # 客户端中途断开(生成器被关闭)
            state.clear()
            state.update(snapshot)  # 同样回滚
        session.touched = time.time()

三个关键点:

  1. copy.deepcopy:深拷贝——快照和原状态完全独立,

改原状态不会影响快照(浅拷贝会让嵌套的列表/字典还是同一份)

  1. 回滚用 state.clear() + state.update(snapshot):

注意不能用 state = snapshot(那只是换了变量指向,外面的会话对象还指着旧字典)

  1. finally 里的 completed 判断:处理"用户关页面导致生成器被提前关闭"的情况——

这时 except 不会执行,但状态已经被改了一半,所以也要回滚

为什么流式接口必须回滚?
HTTP 流一旦开始发送(200 + 数据已出门),就没法再改状态码。
失败只能以 error 事件的形式出现在流里。这时如果状态被改了半截,
用户重试就会基于错误数据——回滚让"重试"变得安全。

块 5:打断 —— 让 AI 现场演剧情(nodes.py 第 240-258 行)

def _stream_text(llm, messages, cut_at=None):
    stream = llm.stream(messages)
    total = 0
    try:
        for chunk in stream:
            content = chunk.content
            if not content:
                continue
            if cut_at is not None and total + len(content) >= cut_at:
                head = content[: max(0, cut_at - total)]   # 截到指定长度
                yield head, True                            # 标记"被截断了"
                return
            total += len(content)
            yield content, False
    finally:
        close = getattr(stream, "close", None)
        if close:
            close()          # 提前关掉模型连接,不再继续生成(省钱!)

打断的实现思路很巧:不需要"另一个进程去砍断",
只要在流式循环里累计长度,到点就停止读取并关闭流。

再看随机触发(app.py 第 128-138 行):

cut_at = None
if session.allow_interrupt:
    chance = 0.3 if round_num == 1 else 0.5
    if random.random() < chance:
        cut_at = random.randint(38, 70)

yield sse("angel_start", {"round": round_num})
for chunk in stream_angel(state, cut_at=cut_at):
    yield sse("angel", {"text": chunk})
if state.get("interrupted") and state["interrupted"][-1]:
    yield sse("interrupt", {"by": "demon"})
yield sse("angel_end", {})

而恶魔的 prompt 会配合剧情(_demon_messages 第 220-221 行):

if interrupted:
    user += "\n注意:天使的话还没说完就被你粗暴地打断了。你的第一句必须是一句短促、嚣张的打断(例如"停停停!""够了!"),然后再展开攻击。\n"

这是本课最值得抄走的设计:状态里记下"发生了打断",
再把这件事翻译成 prompt 指令,让下一次生成延续剧情——
技术事件变成了叙事元素。

块 6:裁判打分 —— 老朋友的容错版(nodes.py 第 302-349 行)

_JSON_RE = re.compile(r"\{.*\}", re.S)

def _clamp_int(value, lo, hi, default):
    try:
        return max(lo, min(hi, int(round(float(value)))))
    except (TypeError, ValueError):
        return default

def evaluate_round(state):
    ...
    angel_score, demon_score, intensity, comment = 5, 5, 5, "势均力敌,胜负未分。"
    try:
        response = _llm_for(state, temperature=0.3).invoke([...])
        match = _JSON_RE.search(response.content or "")
        data = json.loads(match.group(0)) if match else {}
        angel_score = _clamp_int(data.get("angel"), 0, 10, 5)
        ...
    except Exception:
        pass          # 出错就用默认分,游戏继续
    shift = (angel_score - demon_score) * 4
    new_tendency = max(-100, min(100, state.get("tendency", 0) + shift))

和第 10 课的 evaluate_tendency 相比,这里更成熟:

块 7:Key 的三级优先(nodes.py 第 32-36 行)

def resolve_api_key(explicit: Optional[str] = None) -> str:
    key = (explicit or "").strip() or _api_key_override or os.getenv("DEEPSEEK_API_KEY")
    if not key:
        raise ValueError("请先在右上角设置中填写 DeepSeek API Key")
    return key

优先级:请求头带的 X-API-Key > 服务端 set-key 设置的 > 环境变量。
前端每次请求带自己的 Key(client.ts 第 61 行),
服务端不必保存任何人的 Key——这也是"把 Key 交给用户自己保管"的一种产品设计。


第四步:对照开源,别人怎么设计会话接口

对照第 09 课的流式协议

事件第 09 课本课新增
会话无(用内存 dict)session(返回 session_id)
边界start/end同左 + interrupt
判定round_donejudge / judge_start
结束round_done 里带 done单独的 game_over(带 winner)
错误errorerror + retryable 标记

协议是会长大的。设计时预留"边界事件"(start/end)和"可扩展字段",
后面加功能就不用改前端解析逻辑。

对照真实产品里的会话

方案特点适用
内存 dict + TTL/LRU(你的)零依赖、单进程、重启即丢单机原型
Redis过期自动、多进程共享、可持久化生产环境
数据库强持久、可审计需要留档

对照"可重试"的通用套路

这里是内存版的同一思想——出错时回到干净状态


动手练习(做完才算过关)

练习 1(热身):关掉打断做对比
把前端开局的"允许打断"关掉(或在 StartRequest 里把 allow_interrupt 默认改成 False),
跑一局,在 Network 里对比事件序列——少了哪个事件?

练习 2(必做):调整打断参数
把 app.py 第 130-132 行的概率(0.3/0.5)和截断点(38-70)改大,
让打断更频繁、更早,观察恶魔的"第一句嚣张"是不是更常见了。

练习 3(必做):亲手制造一次失败并重试
故意把 API Key 改错(设置里),开始一局 → 触发 error → 观察前端是否提示可重试;
用正确的 Key 再点"下一轮",确认这局还能继续(说明回滚生效了)。
如果没回滚,重试时会发生什么?说说看。

练习 4(必做):并发实验
连续快速点击"下一轮"两次,观察第二次收到的错误消息。
然后在 app.py 第 255 行把 blocking=False 改成 blocking=True,
再试一次,感受两种策略的区别(一个立刻报"进行中",一个默默等待)。

练习 5(思考题):
为什么回滚用 state.clear(); state.update(snapshot) 而不是 state = snapshot?

看答案

session.state 和函数里的 state 是什么关系?

练习 6(选做,进阶):
给 _play_round 加一个新事件 mood:当裁判 intensity ≥ 8 时,
额外 yield sse("mood", {"level": "high"})。前端在 ServerEvent 里加类型并处理它
(可以只是 console.log)。做完你就把"后端状态 → 前端表现"打通了一遍。


本课小结

客户端断开也要回滚(finally 里的 completed 判断)

术语表

术语人话解释
会话(session)一段有状态、跨多次请求的交互
TTL生存时间,超时自动失效
LRU最久未使用优先淘汰
快照(snapshot)某一刻状态的完整副本,用于回滚
深拷贝 / 浅拷贝副本完全独立 / 只复制了外层
回滚(rollback)出错时恢复到之前的状态
可重入同一个操作还没结束又被调用
幂等做一次和做多次结果相同
可辨识联合用某个字段区分多种类型的 TypeScript 写法

下节预告

第 15 课:结构化输出实战:AI 总结 + 信源提取 + 关联图谱。
回到「查资料」,看它的新功能:把多轮对话一键"总结成卡片"(/api/summarize)、
从回答里自动提取信源链接、以及怎么用一张邻接表高效查出卡片之间的关联。
这课讲的是"把 AI 的自由文本变成结构化数据"的完整套路。