From 44684823e0fb8a6e839950450deec79d07b3d135 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 24 Sep 2026 12:37:14 +0000 Subject: [PATCH] =?UTF-8?q?=E6=B5=81=E5=BC=8F=E8=BE=93=E5=87=BA(TALOS=5FST?= =?UTF-8?q?REAM=3D1,=E9=BB=98=E8=AE=A4=E5=85=B3):=E8=BD=AC=E5=9C=88?= =?UTF-8?q?=E9=82=A3=E4=B8=80=E8=A1=8C=E6=8A=A5=E6=A8=A1=E5=9E=8B=E5=9C=A8?= =?UTF-8?q?=E5=86=99=E4=BB=80=E4=B9=88=E3=80=81=E5=86=99=E4=BA=86=E5=A4=9A?= =?UTF-8?q?=E5=B0=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 原来一次调用要么整份回来、要么什么都没有。单次 16s 到 235s 都有,屏幕上只有转圈和 每 30 秒一句「已等 Ns」,分不出慢和卡死,活着的调用被 Ctrl-C 掉过好几次。 开了之后,主循环那一次调用(子 agent 也走这里)边生成边收,转圈上显示 「模型在写 write_file(src\app.py) · 12,000 字」「在思考 · N 字」「在写回复 · N 字」; 老 Windows 控制台不开转圈,同一句进心跳那一行。最终回答照旧收完再出 Markdown 面板, 不逐字打 —— 边打边渲染要反复重绘,仓库记过那在老控制台上不可靠。 · _collect 把流拼回跟整份返回同一个形状(content / tool_calls / reasoning_content / usage),下游一行没改。判据直接比下游结果:同一轮流式、整份各跑一遍, 返回值、历史、token 账逐项相同 · 工具调用按各家的拆法拼:有 index 按 index;Gemini 3 不给 index,按 id 分, thought_signature 跟着 extra_content 留下;每片重复同一个 id 的还是一个调用; 没给 id 的补一个(None 发回去下一轮 400) · 用量靠 include_usage 要:OpenAI / DeepSeek 在最后一块上,Kimi 在 choices[0] 里, 是 SDK 不认识的原样 dict · 拼在 _chat 的重试里面:断流跟掉线同一套认法,整次重来,半截不带过来; Ctrl-C 打断也关连接 · 进度里的 path 从 JSON 参数里抠,认转义(Windows 路径)、转义后才进 rich · 判断器、压缩摘要、复盘不流 —— 看不见的地方流了也没用 · conftest 把 STREAM 钉成关:开着它的终端里跑 pytest 不会无关地红一片 用真 openai SDK + httpx.MockTransport 喂各家形状的 SSE 核对过:OpenAI 拆片、 Gemini 无 index 带签名、Kimi 思考 + choices 里的用量、流中途的 error 事件、断连。 手上没有 key,没真调用过 —— 所以默认关。 343 判据:341 过、5 skip、1 xfail(原有);TALOS_STREAM=1 下全套同样绿。 selfcheck 绿。 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_014gdkHx6rLSMVsUmQiErVSn --- DEVELOPMENT.md | 4 +- README.md | 12 +-- agent.py | 105 +++++++++++++++++++++++- console_ui.py | 23 +++++- tests/conftest.py | 3 + tests/test_loop.py | 195 +++++++++++++++++++++++++++++++++++++++++++- tests/test_tools.py | 21 +++++ 7 files changed, 353 insertions(+), 10 deletions(-) diff --git a/DEVELOPMENT.md b/DEVELOPMENT.md index 9934d6e..14363d1 100644 --- a/DEVELOPMENT.md +++ b/DEVELOPMENT.md @@ -150,7 +150,7 @@ def check_permission(state, cls, name, args): # 决策 + ```bash python agent.py --selfcheck # 零依赖、零网络、零 key。改完先跑这个 -python -m pytest tests/ -q # 337 条。CI 跑 ubuntu/windows × 3.10/3.13 +python -m pytest tests/ -q # 343 条。CI 跑 ubuntu/windows × 3.10/3.13 ``` `--selfcheck` 覆盖 4 个工具的读/写/改、`edit_file` 的"找不到 / 不唯一"两种报错、 @@ -190,7 +190,7 @@ set TALOS_MAX_STEPS=12 | **行为类** | **模型读到守卫那条消息之后干什么** | ❌ 必须用你真在用的那个 | 「拒绝之后会不会改用 `python -c` 绕路」是那个具体模型的性格,换了模型不算数。 -而**行为类正是 live 测试唯一还值钱的部分** —— 管道通不通,337 条判据已经免费覆盖了。 +而**行为类正是 live 测试唯一还值钱的部分** —— 管道通不通,343 条判据已经免费覆盖了。 按这条界,**目标闸的 live 判据是行为类的**:离线判据已经证明判断器拿得到 `read_file`、 越界会被驳回、判不出来时不假装成功;它们证明不了的是**这个具体模型拿到只读工具之后 diff --git a/README.md b/README.md index 8274bb1..5682255 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ [![tests](https://github.com/Jerry-TZ/Talos/actions/workflows/test.yml/badge.svg)](https://github.com/Jerry-TZ/Talos/actions/workflows/test.yml) -**一个你能完整读完的编程 agent。** 5311 行 Python,337 条离线判据,一份不糊弄人的安全说明。 +**一个你能完整读完的编程 agent。** 5435 行 Python,343 条离线判据,一份不糊弄人的安全说明。 Talos 终端界面:动盘之前先弹确认框 @@ -19,8 +19,8 @@ | 文件 | 行数 | 职责 | |---|---|---| -| `agent.py` | 4209 | 循环 + 工具 + 权限门 + 自学习 | -| `console_ui.py` | 214 | 终端界面(可整体替换) | +| `agent.py` | 4312 | 循环 + 工具 + 权限门 + 自学习 | +| `console_ui.py` | 235 | 终端界面(可整体替换) | | `recall.py` | 582 | 联想记忆:扩散激活检索 | | `session.py` | 306 | 会话持久化(想换 SQLite 只改这个) | @@ -33,7 +33,7 @@ 极简 agent 赛道很挤,有人用 Zig 做到 678KB 二进制。**Talos 不比谁小,它比谁都好读。** - **能读完** — 四个文件,注释解释的是*为什么*,不是*是什么*。几乎每条防御旁边都写着它挡的那次真实翻车。 -- **能验证** — 337 个测试,**离线、免 API key、几秒跑完**,CI 在 Linux/Windows × Python 3.10/3.13 上都跑。clone 下来立刻知道它没坏。 +- **能验证** — 343 个测试,**离线、免 API key、几秒跑完**,CI 在 Linux/Windows × Python 3.10/3.13 上都跑。clone 下来立刻知道它没坏。 - **不吹牛** — [`SECURITY.md`](SECURITY.md) 明写 `create_tool` 就是进程内 RCE、正则黑名单只是减速带。**没有沙箱就是没有沙箱** —— 真要隔离,[三条现成方案](SECURITY.md#真要隔离怎么办)按代价从低到高列在那儿。 - **有考卷** — [`EXAM.md`](EXAM.md) / [`EXAM2.md`](EXAM2.md) 是两份可复现的能力测试,带标准答案和作弊检测(比如逐个核验 arXiv ID 真伪,防止编造引用)。记录的是"我怎么验证它真的有用",不是功能列表。 - **有实测** — [`FINDINGS.md`](FINDINGS.md) 记了二十二个真实任务量出来的东西:哪条提示词生效、哪条从头到尾没生效、六次翻车、检索改动的前后数字,以及**两个被数据否掉的自己的方案**。样本小,局限写在最前面。 @@ -159,6 +159,8 @@ tail -3 .talos/recall_trace.jsonl **一次性模式** — `agent.py -p "任务"` 跑完即退,方便脚本和计时。 +**流式(试验,默认关)** — `TALOS_STREAM=1` 让主循环边生成边收,转圈那一行实时显示「模型在写 write_file(src/app.py) · 12,000 字」,分得清慢和卡死。拼回来的回复跟整份返回的逐项相同(判据钉着),断流按原来的重试规矩整次重来。几家的流式形状只照文档和造出来的数据核对过,**还没真调用过** —— 跑出问题关掉,就回到原来那条路。 + **可选自动化**(默认关闭,在 `talos.bat` 里开): ```bat @@ -200,7 +202,7 @@ Ctrl+C 停下当前这轮 —— 做过的都留着,可以直接说 ```bash .venv\Scripts\python.exe -m pip install -r requirements-dev.txt -.venv\Scripts\python.exe -m pytest tests/ -q # 337 条,约 8 秒,不联网、不需要 key +.venv\Scripts\python.exe -m pytest tests/ -q # 343 条,约 8 秒,不联网、不需要 key python agent.py --selfcheck # 免依赖的冒烟检查 ``` diff --git a/agent.py b/agent.py index 28aff48..c905893 100644 --- a/agent.py +++ b/agent.py @@ -32,6 +32,7 @@ import re import sys import time +import types # .env is loaded from the launch directory, which in a coding agent is very often somebody # else's repository. Config is fine to pick up there; a command to execute is not — that turns @@ -260,6 +261,13 @@ def _env_block() -> str: # 后面也没法用。宁可这条稍紧,反正连不上会当场说清楚原因,不像超时那样一声不吭。 CONNECT_TIMEOUT = float(os.environ.get("TALOS_CONNECT_TIMEOUT", "5")) SLOW_CALL = float(os.environ.get("TALOS_SLOW_CALL", "15")) # 超过这么久的调用才报耗时 +# 流式:主循环那一次调用边生成边收,转圈那一行实时报「在写 write_file(x.py) · 12,000 字」。 +# **默认关**:六家的流式细节各不一样(Gemini 不给 index、Kimi 把用量塞进 choices[0]), +# 写它的时候手上没有 key,只照文档和造出来的块核对过。开着跑稳了再改默认。 +# 只影响看得见的那一处 —— 判断器、压缩摘要、复盘你看不到,流了也没用。 +STREAM = os.environ.get("TALOS_STREAM", "").strip() in ("1", "true", "yes", "on") +# 不带 include_usage,流式默认**不回用量**:会话预算提醒和缓存命中率就静悄悄地全是 0 +_STREAM_KW = {"stream": True, "stream_options": {"include_usage": True}} ui = None # 界面 handle, set by repl(); kept out of module scope so --selfcheck is dep-free _RUNTIME = {} # live client/model/state (+ subagent depth), set in agent_turn so tools like @@ -2436,6 +2444,97 @@ def _retry_after(e): except (TypeError, ValueError): return None +def _ns(v): + """dict → 能按属性读的对象,递归。Kimi 把流式的用量放在 `choices[0].usage`,SDK 不认识 + 这个位置,给的是原样的 dict —— 而 `_usage` 按属性读,读 dict 只会静悄悄地拿到 0。""" + if isinstance(v, dict): + return types.SimpleNamespace(**{k: _ns(x) for k, x in v.items()}) + return v + +# 参数是 JSON:Windows 路径在里面是 `src\\app.py`,抠的时候得认转义,抠完再解一遍 +_PATH_ARG = re.compile(r'"path"\s*:\s*"((?:[^"\\]|\\.)*)"') + +def _collect(stream): + """把一次流式调用拼回**跟整份返回一模一样**的形状:`resp.choices[0].message` 带 + content / tool_calls / reasoning_content,`resp.usage` 带用量。下游那几千行一行不用改 —— + 流式换的只是「怎么拿到这份回复」,不是「回复是什么」。 + + 工具调用是拼的重点,几家拆法不一样: + - OpenAI / DeepSeek:每片带 `index`,第一片给 id 和名字,后面只有一段段参数 + - Gemini 3:**不给 `index`**,每个调用一整块,`thought_signature` 挂在这一块的 + `extra_content` 上(丢了下一轮当场 400,见 `_tool_call_entry`) + - 有的每一片都重复同一个 id —— 那还是同一个调用 + 所以有 index 按 index 归位;没有就看 id:变了是新调用,没变或没给是接着上一个。 + 名字是**赋值**不是拼接(没有哪家把名字拆开发,而重复发名字的有)。 + + 一边拼一边报进度(`ui.progress`),最该报的是写大文件:参数就是整份文件,要写几分钟。 + 流收完、断掉、被 Ctrl-C 打断,连接都关。""" + text, think, calls, at, usage = [], [], [], {}, None + n_text = n_think = 0 + try: + for chunk in stream: + usage = getattr(chunk, "usage", None) or usage + for ch in (getattr(chunk, "choices", None) or ())[:1]: + usage = getattr(ch, "usage", None) or usage # Kimi 放在这儿 + d = getattr(ch, "delta", None) + if d is None: + continue + said = None + r = _reasoning(d) + if r: + think.append(r) + n_think += len(r) + said = f"在思考 · {n_think:,} 字" + if getattr(d, "content", None): + text.append(d.content) + n_text += len(d.content) + said = f"在写回复 · {n_text:,} 字" + for tc in getattr(d, "tool_calls", None) or (): + i, cid = getattr(tc, "index", None), getattr(tc, "id", None) + f = getattr(tc, "function", None) + if i is not None and i in at: + c = calls[at[i]] + elif i is None and calls and not (cid and cid != calls[-1]["id"]): + c = calls[-1] + else: + c = {"id": cid, "name": "", "args": [], "n": 0, "path": None, "extra": {}} + calls.append(c) + if i is not None: + at[i] = len(calls) - 1 + c["id"] = c["id"] or cid + if getattr(f, "name", None): + c["name"] = f.name + a = getattr(f, "arguments", None) + if a: + c["args"].append(a) + c["n"] += len(a) + # path 只在参数开头找:在后面的话要一遍遍 join 整份文件去搜,不值 + if c["path"] is None and c["n"] - len(a) < 300: + m = _PATH_ARG.search("".join(c["args"])[:300]) + try: + c["path"] = json.loads(f'"{m.group(1)}"') if m else None + except ValueError: + c["path"] = m.group(1) + c["extra"].update({k: v for k, v in + (getattr(tc, "model_extra", None) or {}).items() + if v is not None}) + where = f"({c['path']})" if c["path"] else "" + said = f"在写 {c['name']}{where} · {c['n']:,} 字" + if said and ui is not None: + ui.progress(said) + finally: + close = getattr(stream, "close", None) + if callable(close): + close() + # 没给 id 的补一个:工具结果靠 `tool_call_id` 对回调用,None 发回去是 null,下一轮 400 + tool_calls = [types.SimpleNamespace( + id=c["id"] or f"call_{k}", type="function", model_extra=c["extra"], + function=types.SimpleNamespace(name=c["name"], arguments="".join(c["args"]))) + for k, c in enumerate(calls)] or None + msg = types.SimpleNamespace(content="".join(text) or None, tool_calls=tool_calls, + reasoning_content="".join(think) or None) + return types.SimpleNamespace(choices=[types.SimpleNamespace(message=msg)], usage=_ns(usage)) + def _chat(client, **kwargs): """Call the model, retrying briefly on rate-limit / 'busy' / transient errors. Free tiers (esp. glm-4.7-flash) get congested — '当前模型用户多' is just a busy signal.""" @@ -2444,6 +2543,10 @@ def _chat(client, **kwargs): t0 = time.time() try: resp = client.chat.completions.create(**kwargs) + # 流式的错**在读到一半时**才出来,所以拼这一步必须在 try 里面:断流跟连不上 + # 走同一套认法(掉线是 "connection",重试),重来那一趟从零拼,半截不带过来。 + if kwargs.get("stream"): + resp = _collect(resp) # 慢调用才报。快的不值一行,而推理模型上"这一次比平常久得多"正是你想知道的事。 took = time.time() - t0 if ui is not None and took >= SLOW_CALL: @@ -3014,7 +3117,7 @@ def _seal_trace(steps: int, capped: bool) -> None: resp = _chat(client, model=model, messages=([{"role": "system", "content": system}] + _with_recall(messages, recalled, slot)), - tools=tool_specs()) + tools=tool_specs(), **(_STREAM_KW if STREAM else {})) state["tok"]["steps"] += 1 _in, _out, _cached = _usage(resp) for _k, _v in zip(("in", "out", "cached"), (_in, _out, _cached)): diff --git a/console_ui.py b/console_ui.py index 808b66a..217c128 100644 --- a/console_ui.py +++ b/console_ui.py @@ -48,6 +48,7 @@ def read_task(mode: str) -> str: return console.input(f"[bold {C_YOU}]你[/] [dim]({mode})[/] › ").strip() HEARTBEAT = 30 # 每这么多秒往下追加一行"还活着" +_active = None # 正在转的那个圈 —— `progress()` 往它上面写;同一时刻只有一个(子 agent 在圈外跑) class _Thinking: """转圈 + 每 HEARTBEAT 秒**追加**一行「已等 Ns」。 @@ -68,8 +69,11 @@ def __init__(self, label: str): self._stop = threading.Event() self._status = None if console.legacy_windows else console.status(label, spinner="dots") self._label = label + self._said = "" # 流式时模型写到哪了,`progress()` 填 def __enter__(self): + global _active + _active = self if self._status is not None: self._status.__enter__() else: @@ -89,9 +93,20 @@ def _beat(self) -> None: waited += HEARTBEAT # 用 :.0f —— waited 是累加出来的,HEARTBEAT 非整数时会攒出 # 0.15000000000000002 这种。生产里 HEARTBEAT=30 看不出来,测试里一眼就露。 - console.print(f"[dim] … 已等 {waited:.0f}s,还在等模型回话(Ctrl-C 停这一轮)[/]") + # 流式时换成模型写到哪了 —— legacy_windows 不开转圈,这一行是那儿唯一的进度。 + doing = f"模型{escape(self._said)}" if self._said else "还在等模型回话" + console.print(f"[dim] … 已等 {waited:.0f}s,{doing}(Ctrl-C 停这一轮)[/]") + + def progress(self, said: str) -> None: + self._said = said + if self._status is not None: + # 转圈的**文字**原地换,不追加 —— 秒数那条说过原地重绘会刷屏,但那是 + # legacy_windows 上;这里 `_status` 存在就说明不是那种控制台。 + self._status.update(f"[{C_MODEL}]模型{escape(said)}[/]") def __exit__(self, *exc): + global _active + _active = None self._stop.set() return self._status.__exit__(*exc) if self._status is not None else False @@ -100,6 +115,12 @@ def thinking(): """上下文管理器:模型思考时转个圈,并定期报"还活着"。""" return _Thinking(f"[{C_MODEL}]模型思考中…[/]") +def progress(said: str) -> None: + """流式时模型写到哪了(「在写 write_file(x.py) · 12,000 字」),显示在转圈上和心跳里。 + 转圈之外调它什么也不做。""" + if _active is not None: + _active.progress(said) + def took(seconds: float) -> None: """调用回来之后报一次耗时 —— 跟上面的心跳配套:心跳说"还活着",这条说"花了多久"。""" console.print(f"[dim] ⏱ {seconds:.0f}s[/]") diff --git a/tests/conftest.py b/tests/conftest.py index fd138f6..febdf1b 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -85,6 +85,9 @@ def _no_test_ever_writes_the_real_talos(tmp_path, monkeypatch): # 就会往真实的 `.talos/sessions/` 里塞垃圾,而那批文件正是往事检索和燃尽表的语料。 # 上面那段说「枚举永远落后一步」,这就是下一步:**加功能等于给这片面新开一条路。** monkeypatch.setattr(session, "SESS_DIR", os.path.join(d, "sessions")) + # 开着 `TALOS_STREAM=1` 的终端里跑 pytest,几百条用例的假客户端会收到 stream=True + # 而照旧回整份 —— 全套红一片,红的原因跟被测的东西毫无关系。要测流式的自己打开。 + monkeypatch.setattr(agent, "STREAM", False) @pytest.fixture(autouse=True) def _keys_stay_put(): diff --git a/tests/test_loop.py b/tests/test_loop.py index 753b3f2..a537731 100644 --- a/tests/test_loop.py +++ b/tests/test_loop.py @@ -6,7 +6,7 @@ def _ui(): n = lambda *a, **k: None - return types.SimpleNamespace(thinking=lambda: contextlib.nullcontext(), + return types.SimpleNamespace(thinking=lambda: contextlib.nullcontext(), progress=n, show_tool=n, denied=n, think=n, assistant_text=n, note=n, error=n) def _msg(content=None, tool_calls=None, usage=None): @@ -2687,3 +2687,196 @@ def turn(client, model, messages, state, top=False): assert [a["content"] for a in asks] == ["新的活"], \ f"删掉的会话里的原话还挂在 asks 上:{[a['content'] for a in asks]}" assert b_human not in asked, "已删掉的会话的原话还在参与「点名」判断" + + +# ── 流式(TALOS_STREAM=1):拼回来的必须跟整份返回的是同一个东西 ───────────────── +# 形状照 openai 2.x 的 `ChatCompletionChunk` 造(真对象在本机核对过),不 import openai。 +def _d(content=None, reasoning=None, tool_calls=None): + """一块增量。`reasoning_content` 是 DeepSeek / Kimi 的字段,SDK 不认识,挂在 extra 上。""" + delta = types.SimpleNamespace(content=content, tool_calls=tool_calls, + reasoning_content=reasoning) + return types.SimpleNamespace(choices=[types.SimpleNamespace(index=0, delta=delta, + finish_reason=None)], usage=None) + +def _frag(index, cid=None, name=None, args=None, extra=None): + """工具调用的一个碎片。OpenAI 的拆法:第一片给 index/id/name,后面只有 index + 一段参数。""" + return types.SimpleNamespace(index=index, id=cid, type="function" if cid else None, + function=types.SimpleNamespace(name=name, arguments=args), + model_extra=extra or {}) + +def _split(s, n): + return [s[i:i + n] for i in range(0, len(s), n)] + +class _Stream: + """流式返回:能迭代、有 close()。给了 `die` 就在吐完之后抛它 —— 流到一半断掉。""" + def __init__(self, chunks, die=None): + self.chunks, self.die, self.closed = chunks, die, False + def __iter__(self): + yield from self.chunks + if self.die is not None: + raise self.die + def close(self): + self.closed = True + +class _KwClient(_Client): + """跟 `_Client` 一样,外加记下每次调用带了哪些参数。""" + def __init__(self, script): + super().__init__(script) + self.kw = [] + def _c(self, **k): + self.kw.append(k) + return super()._c(**k) + + +def test_a_streamed_turn_leaves_exactly_what_a_whole_reply_would(ws, monkeypatch): + """同一轮,流式跑一遍、整份跑一遍:返回值、历史、token 账**逐项相同**。 + + 流式只该换「怎么拿到这份回复」,不该换「回复是什么」。拼回来的对象后面还有四千行 + 代码要读(工具分发、落盘、压缩、回放),哪一处的字段对不上,都是在真跑到那条路径时才炸。 + 所以判据不去逐个字段比,直接比**下游看到的全部结果**。 + + 参数故意切成 5 个字符一片、两个调用交错不了(OpenAI 的拆法是一个调用拆完再下一个); + 用量在最后一块,那一块 `choices` 是空的 —— 这是 `include_usage` 的真实形状。""" + import agent as A + monkeypatch.setattr(A, "ui", _ui()) + usage = types.SimpleNamespace(prompt_tokens=100, completion_tokens=20, + prompt_tokens_details=types.SimpleNamespace(cached_tokens=60)) + a1, a2 = '{"command": "echo alpha"}', '{"command": "echo beta"}' + whole = [_msg(content="先看看", usage=usage, + tool_calls=[_tc("run_bash", a1, "c1"), _tc("run_bash", a2, "c2")]), + _msg(content="两个都跑完了。", usage=usage)] + last = types.SimpleNamespace(choices=[], usage=usage) + streamed = [ + _Stream([_d(reasoning="要跑"), _d(reasoning="两条"), _d(content="先"), _d(content="看看"), + _d(tool_calls=[_frag(0, "c1", "run_bash", "")]), + *[_d(tool_calls=[_frag(0, args=p)]) for p in _split(a1, 5)], + _d(tool_calls=[_frag(1, "c2", "run_bash", "")]), + *[_d(tool_calls=[_frag(1, args=p)]) for p in _split(a2, 5)], + last]), + _Stream([_d(content="两个都"), _d(content="跑完了。"), last])] + runs = {} + for how, script, on in (("整份", whole, False), ("流式", streamed, True)): + monkeypatch.setattr(A, "STREAM", on) + client = _KwClient(script) + msgs = [{"role": "user", "content": "跑两条"}] + state = {"mode": "bypass", "allow": set()} + out = A.agent_turn(client, "m", msgs, state) + runs[how] = (out, msgs, state["tok"], client.kw) + assert runs["流式"][0] == runs["整份"][0] == "两个都跑完了。" + assert runs["流式"][1] == runs["整份"][1], "流式拼回来的历史跟整份返回的不一样" + assert runs["流式"][2] == runs["整份"][2], "流式的 token 账跟整份的不一样(用量那一块丢了?)" + assert "alpha" in str(runs["流式"][1]) and "beta" in str(runs["流式"][1]), "两个调用没都跑" + # 开着的时候要用量;关着的时候一个流式参数都不带 —— 默认关,就得真的是原来那条路 + assert all(k.get("stream") is True and k.get("stream_options") == {"include_usage": True} + for k in runs["流式"][3]), runs["流式"][3] + assert not any("stream" in k or "stream_options" in k for k in runs["整份"][3]) + + +def test_streamed_tool_calls_are_told_apart_the_way_each_provider_sends_them(monkeypatch): + """几家拆工具调用的方式不一样,拼法得都认: + + - OpenAI / DeepSeek:每片带 `index`,按它归位(上一条判据) + - Gemini 3:**每个调用一整块、不给 `index`**,`thought_signature` 挂在这一块的 + `extra_content` 上 —— 签名丢了,下一轮回放当场 400(`_tool_call_entry` 那条记着) + - 有的每一片都重复同一个 `id`:那还是**同一个**调用,不能拆成两个 + + 不给 index 时靠 id 分:id 变了是新调用,没变(或没给)是接着上一个。 + 这条判据还顺带钉了「不需要界面也能拼」:`ui` 是 None。""" + import agent as A + monkeypatch.setattr(A, "ui", None) + sig = {"google": {"thought_signature": "sig-1"}} + gem = _Stream([_d(tool_calls=[_frag(None, "g1", "read_file", '{"path": "a.py"}', + {"extra_content": sig})]), + _d(tool_calls=[_frag(None, "g2", "read_file", '{"path": "b.py"}')])]) + calls = A._collect(gem).choices[0].message.tool_calls + assert [(c.id, c.function.name, c.function.arguments) for c in calls] == [ + ("g1", "read_file", '{"path": "a.py"}'), ("g2", "read_file", '{"path": "b.py"}')], ( + "不给 index 的两个调用被粘成了一个(或者拆错了)") + assert A._tool_call_entry(calls[0])["extra_content"] == sig, "Gemini 的签名在流式里被抄丢了" + assert gem.closed, "流收完没关" + + rep = _Stream([_d(tool_calls=[_frag(None, "d1", "write_file", '{"path": "x.py", ')]), + _d(tool_calls=[_frag(None, "d1", None, '"content": "hi"}')])]) + calls = A._collect(rep).choices[0].message.tool_calls + assert len(calls) == 1 and calls[0].function.arguments == '{"path": "x.py", "content": "hi"}' + + # 没给 id:工具结果要拿 `tool_call_id` 对回这个调用,None 发出去就是 null,下一轮 400 + anon = A._collect(_Stream([_d(tool_calls=[_frag(0, None, "run_bash", '{"command": "ls"}')])])) + assert anon.choices[0].message.tool_calls[0].id, "没给 id 的调用拼出来 id 是空的" + + plain = A._collect(_Stream([_d(content="就这样")])).choices[0].message + assert plain.content == "就这样" and plain.tool_calls is None, ( + "没有工具调用时 tool_calls 得是 None,跟整份返回一样 —— 循环拿它判断这一轮是不是说完了") + + +def test_stream_usage_is_read_from_wherever_the_provider_puts_it(monkeypatch): + """用量不回来,会话预算提醒和缓存命中率就静悄悄地全是 0 —— 不报错,只是数不对。 + + OpenAI / DeepSeek 把它放在最后一块上(那一块 `choices` 是空的);Kimi 放在最后一块的 + `choices[0].usage` 里 —— SDK 不认识这个位置,给的是**原样的 dict**,而 `_usage` + 按属性读。两种都得认,dict 里嵌的 dict 也得能按属性读。""" + import agent as A + monkeypatch.setattr(A, "ui", None) + obj = types.SimpleNamespace(prompt_tokens=100, completion_tokens=20, + prompt_tokens_details=types.SimpleNamespace(cached_tokens=60)) + openai_style = _Stream([_d(content="好"), types.SimpleNamespace(choices=[], usage=obj)]) + assert A._usage(A._collect(openai_style)) == (100, 20, 60) + + fin = types.SimpleNamespace(index=0, delta=types.SimpleNamespace(content=None, tool_calls=None), + finish_reason="stop", + usage={"prompt_tokens": 100, "completion_tokens": 20, + "prompt_tokens_details": {"cached_tokens": 60}}) + kimi_style = _Stream([_d(content="好"), types.SimpleNamespace(choices=[fin], usage=None)]) + assert A._usage(A._collect(kimi_style)) == (100, 20, 60), "Kimi 放在 choices[0] 里的用量没认出来" + + +def test_a_stream_that_breaks_midway_starts_over_instead_of_splicing(monkeypatch): + """流到一半断了:整次重来,不接着拼。 + + 原来 `_chat` 的重试只包住 `create()`,而流式的错误**在读到一半时**才出来 —— + 拼在 `_chat` 外面的话,断线就直接抛出去了,重试那一整套(状态码、退避、余额)全不管用。 + 所以拼这一步放在重试**里面**:断了就按原来的规矩认(掉线是 "connection",会重试), + 重来的那一趟从零开始,前一趟收到的半截一个字都不带过来。 + + Ctrl-C 也在读到一半时来。它不是 Exception,不重试,但**连接得关**。""" + import agent as A + monkeypatch.setattr(A.time, "sleep", lambda _s: None) + monkeypatch.setattr(A, "ui", _ui()) + # httpx 断流时的原话 + broken = _Stream([_d(content="写到一半")], die=Exception( + "peer closed connection without sending complete message body (incomplete chunked read)")) + whole = _Stream([_d(content="完整的回答")]) + resp = A._chat(_Client([broken, whole]), stream=True) + assert resp.choices[0].message.content == "完整的回答", "断掉那一趟的半截被拼进来了" + assert broken.closed and whole.closed + + stopped = _Stream([_d(content="写到一半")], die=KeyboardInterrupt()) + with pytest.raises(KeyboardInterrupt): + A._chat(_Client([stopped]), stream=True) + assert stopped.closed, "Ctrl-C 之后连接没关" + + +def test_while_streaming_the_spinner_says_what_is_being_written(monkeypatch): + """流式要换来的就是这一行:等待的时候知道模型在干什么、干了多少。 + + 原来只有转圈和每 30 秒一句「已等 Ns」—— 单次调用 16s 到 235s 都有,分不出慢和卡死, + 活着的调用被 Ctrl-C 掉过好几次。最该看见进度的是写大文件:参数就是整份文件,几分钟, + 所以工具名和 path 都要报(`path` 在参数开头才报得出来,在后面就只报工具名)。 + + path 用 Windows 写法:参数是 JSON,`src\\app.py` 在里面是转义过的两个反斜杠, + 按「引号之间不含反斜杠」去抠的话只抠得出 `src`。""" + import json + import agent as A + said = [] + ui = _ui() + ui.progress = said.append + monkeypatch.setattr(A, "ui", ui) + args = json.dumps({"path": "src\\app.py", "content": "x" * 12000}) + A._collect(_Stream([_d(reasoning="先想想"), _d(content="好,"), + _d(tool_calls=[_frag(0, "c1", "write_file", "")]), + *[_d(tool_calls=[_frag(0, args=p)]) for p in _split(args, 500)]])) + assert said, "流式期间一次进度都没报" + assert "思考" in said[0] and "3" in said[0], said[:2] + assert "回复" in said[1] and "2" in said[1], said[:2] + assert "write_file" in said[-1] and "(src\\app.py)" in said[-1] and f"{len(args):,}" in said[-1], ( + f"写大文件时没说在写哪个文件、写了多少:{said[-1]!r}") diff --git a/tests/test_tools.py b/tests/test_tools.py index 1c2c6b5..ad1a097 100644 --- a/tests/test_tools.py +++ b/tests/test_tools.py @@ -950,6 +950,27 @@ def test_a_long_call_says_it_is_still_alive(monkeypatch): assert beats, f"长调用期间一行心跳都没有:{lines!r}" assert "0.1s" not in " ".join(beats), f"秒数没取整:{beats!r}" +def test_what_the_stream_has_written_so_far_shows_up_in_the_heartbeat(monkeypatch): + """流式时 `ui.progress` 报的那句,心跳行里也要有 —— 老 Windows 控制台上不开转圈, + 心跳是那儿**唯一**看得见的东西。 + + 那句话里有模型写的 path,而 path 是模型给的:`[red]` 这种东西不转义就被 rich 当成 + 样式吃掉(轻则字没了,重则 MarkupError 把这一轮打断)。 + 转圈之外调 `progress` 什么都不该发生 —— 流式只在转圈里跑,但别让它因为这个炸。""" + import time + ui = pytest.importorskip("console_ui", reason="需要 rich(界面层的可选依赖)") + lines = [] + monkeypatch.setattr(ui.console, "print", lambda *a, **k: lines.append(str(a[0]) if a else "")) + monkeypatch.setattr(ui, "HEARTBEAT", 0.05) + with ui.thinking(): + ui.progress("在写 write_file([red]x.py) · 12,000 字") + time.sleep(0.25) + ui.progress("转圈已经停了") + beats = [ln for ln in lines if "已等" in ln] + assert any("12,000 字" in b for b in beats), f"心跳里没有流式进度:{beats!r}" + assert any("\\[red]x.py" in b for b in beats), f"模型给的 path 没转义就进了 rich:{beats!r}" + assert not any("转圈已经停了" in ln for ln in lines) + def test_no_provider_key_survives_startup(tmp_path): """启动那一刻,六个 provider 的 key 一个都不该留在 `os.environ` 里。