ReMe/tests4/integration/test_auto_memory.py
jinliyl 83831ec90c
feat(core): enhance reme4 (#281)
### 1. Agent Wrapper(统一 Agent 后端抽象)
- **`base_agent_wrapper.py`**:`reply()` 返回值从 `tuple[str, Any]` 改为 `dict`(含 `session_id` / `last_message` / `result` / 可选 `structured_output`);`reply_stream()` 改为产出统一的 `StreamChunk`。废弃 `add_tools()`,改为 `add_job_tools(names: list[str])`(按名解析 BaseJob)与 `add_skills()`;新增 `_resolve_job_tools()`、`_merged_kwargs()`、`_chunk()` 辅助方法及 `project_path` / `project_skills_root` 属性。
- **`as_agent_wrapper.py`(AgentScope 后端)**:
  - 会话持久化重写:`session_path` 落地到 `<vault>/<session_dir>/agentscope/`,`_load_state` 支持 `resume` / `session_id` / `fork_session`,并做 UUID 校验(`_validate_session_id`);`_cleanup_expired_sessions` 按天数清理过期会话。
  - 新增内置工具集(`BypassAnalysisBash` + Edit/Glob/Grep/Read/Write),`BypassAnalysisBash` 绕过 AgentScope 自带 Bash 静态分析以让 permission_mode 生效;`_resolve_skills()` 把配置的 skill 暴露给后端,`_load_tool_env()` 注入项目 `.env`。
  - `_event_to_chunk()` 把 20+ 种 AgentScope 事件(Reply/Text/Thinking/Data/ToolCall/ToolResult/ModelCall/ExceedMaxIters)归一化为 `StreamChunk`。
- **`cc_agent_wrapper.py`(Claude Code SDK 后端,+551 行)**:
  - 新增 `_CcFileSessionStore`:基于 vault 的文件型会话存储,实现 append(按 uuid 去重)/ load / list / delete / list_subkeys,并对路径做 `_safe_parts` + `resolve()` 防越界校验。
  - `_build_options()`:统一构建 `ClaudeAgentOptions`,处理 skills、disallowed_tools(默认禁 `WebSearch`)、`.env` 注入、Claude Code 的 API 凭据解析(`_claude_code_api_env`,多级 base_url/api_key 回退)、`CLAUDE_CONFIG_DIR` 设置、skill 目录软链接(`_ensure_claude_skill_dir`)。
  - `_raw_event_to_chunk()` / `_message_content_to_chunks()`:把 Anthropic 流式事件(message_start/delta/stop、content_block_*)与 SDK 消息块(AssistantMessage/UserMessage/ResultMessage/RateLimitEvent)转换为统一 `StreamChunk`;跟踪 block_id/block_type/tool_call_name 做关联;处理尾部 `"success"` 误报异常的吞掉逻辑。

### 2. 统一流式协议(StreamChunk / ChunkEnum)
- **`stream_chunk.py`**:`StreamChunk` 扩展为承载 AS + CC 双后端完整信息的统一结构,新增 `session_id` / `block_id` / `tool_call_id` / `tool_call_name` / `media_type` / `input_tokens` / `output_tokens` 等字段,纯文本流仍保持轻量。
- **`chunk_enum.py`**:补全生命周期标记 `REPLY_START` / `REPLY_END`,并文档化两套后端事件 → ChunkEnum 的映射。

### 3. Index 模块重构(变化批次化 + dispatch)
- 新增 `_change_batch.py`:`coalesce_changes()` 把同路径多次事件折叠为最终状态(结合 path 存在性判定),`bucket_changes()` 按 watchfiles.Change 分桶。
- 新增 `init_changes.py`(`InitChangesStep`):一次性扫描,对比 file_store / file_catalog 已索引节点计算 added/modified/deleted,写入 `context["changes"]` 后 dispatch。
- 新增 `update_changes.py`:抽象基类 `ChangeApplyStep` 统一 added/modified/deleted 处理与错误收集;`UpdateCatalogStep`(写 file_catalog)、`UpdateIndexStep`(写 file_store,含按后缀解析 chunker)。
- **`watch_changes.py`**:改用 `dispatch_step_specs`(基类提供的 `dispatch_steps()`),每批先 `coalesce_changes` 再 dispatch;默认参数调整(debounce 5000ms / step 1000ms / poll 5000ms)并暴露常量。
- 删除旧步骤:`clear_and_scan` / `foreach_dispatch` / `scan_changes` / `update_catalog`(旧) / `update_index`(旧);`clear_store.py` 取代 clear_and_scan。

### 4. Evolve / Dream 模块(拆分为多步 pipeline)
- 删除旧的单体 `auto_dream.py` / `dream.py` / `dream.yaml`,新增 `dream/` 子包,按 5 个步骤组织:
  - **`extract.py`**:扫描当日 day-index + daily 笔记,对比 file_catalog 找出 changed/deleted,调用 LLM 全局抽取 `units`(procedure/personal/wiki 三桶)与 `topics`,路径与桶做清洗/路由。
  - **`integrate.py`**:逐个 unit 调用 LLM 写入 digest,结构化输出 `IntegrateOutcome`(CREATE/CORROBORATE/REFINE/CORRECT),失败 unit/路径收集回写。
  - **`topics.py`**:写 `daily/<date>/interests.yaml`,结合当天已有 + 近 N 天做去重(`normalize_topic`),可走 LLM 或纯规则去重两条路径。
  - **`proactive.py`**:读取当日 `interests.yaml`,作为主动推荐话题的入口。
  - **`finish.py`**:把变更路径落盘到 dream file_catalog(checkpoint),渲染最终汇总摘要。
- 新增 `schema.py`(`DreamState` 等跨步骤共享状态与结构化输出模型)与 `utils.py`(状态存取、扫描打包、YAML 读写、结构化回复解析等公共函数)。
- `evolve/__init__.py` 导出全部新 step。

### 5. auto_memory / auto_resource(适配新 Agent API)
- **`auto_memory.py`**:会话路径迁移到 `<session_dir>/dialog/<session_id>.jsonl`;改用 `job_tools`;新增 `source_conversation` frontmatter 反向链接(`_session_link`);执行后刷新 day 索引(`refresh_day_index`),并对 session_id 做合法性校验。
- **`auto_resource.py`**:资源改用「同名 daily note」方案(`_compute_note_stem` 取文件 stem);批量处理 `changes: list[dict]`(`_handle_change` 逐项处理,返回逐项结果摘要);agent 会话 id 用稳定的 `uuid5`;同样刷新 day 索引。

### 6. BaseStep 基类增强
- 新增 `dispatch_steps` / `dispatch_step_specs` 机制:`_resolve_dispatch_step()` 支持字符串或 dict 形式的 step spec,`dispatch_steps()` 复用当前 context 调用下游 step。
- 新增 `config_value()`:按 key 取 app config,缺失时回退 `ApplicationConfig` 默认值。
- 小幅清理:`language` 初始化、`copy()`、`Ref.__init__` 签名精简。

### 7. Components 改动
- **`file_store/local_file_store.py`**:持久化改用 zstd 压缩(`.jsonl.zst`,通过新 `utils/jsonl_zst.py`);upsert 时先删除旧 chunk 的 keyword 文档;embedding 复用改为 `(text, embedding)` 键控,要求文本一致才复用;新增 `_matches_search_filter()` 对 vector/keyword 搜索做 path/path_prefix/metadata 的统一后过滤。
- **`keyword_index/bm25_index.py`**:索引文件名加入组件名 + tokenizer 指纹(sha256 前 12 位),快照/恢复时校验指纹防配置漂移;空索引 dump 时删除文件,加载失败抛错而非静默。
- **`file_chunker/markdown_file_chunker.py`**:弃用 `python-frontmatter`,改用内置 YAML 解析(非法 YAML 不阻断正文索引),并修正因 frontmatter 占用行号导致的 AST 行号偏移(`line_offset`)。
- **`cron_job.py`**:大幅简化(-187 行),由原来「dispatch 外部 job/step + 多种调度模式」改为「在自身 steps 上跑 cron 表达式」;`Application` 启动顺序随之调整为 base > stream > background > cron。
- 其余小调整:service(base/http/mcp)、file_graph、file_catalog、as_llm、as_embedding、tokenizer、prompt_handler、base_component 的签名/接口微调。

### 8. Application 生命周期
- `_start()` 启动顺序明确为 components → base → stream → background → cron,启动失败会触发 `_close()` 回滚并 re-raise(不再吞异常)。
- 启动时创建 `session_dir` 目录;新增 `update_component()`(按类型/名就地更新已存在组件,不存在则报错)。

### 9. File IO / 路径安全
- **`_path.py`**:`resolve_path` 增加 vault 越界防护(`is_relative_to` 校验),禁止 `.` / `..` 路径分量,支持 `allow_empty`。
- **`read.py`**:大文件(超过 `MAX_FILE_READ_BYTES`)走按行读取 `read_file_lines_safe`,避免一次性载入内存。
- **`_file_io.py` / `_daily_index.py` / `_path.py`** 等支持函数补齐(如 `refresh_day_index`、`read_file_lines_safe`)。
- **`env_utils.py`**:新增 `parse_env_file()`,`load_env()` 返回加载到的键值、支持 `override`、对无路径调用做幂等缓存。

### 10. Config
- `ApplicationConfig` 新增 `session_dir`(默认 `reme_session`)。
- `config_parser.py`:环境变量展开后做类型转换(`_convert_value`)、dot-notation 与 key=value 参数校验更严格、配置文件路径支持相对 `_CONFIG_DIR` 查找、根非 dict 报错。
- `default.yaml`:作业编排改用 `init_changes_step` + `dispatch_steps`(index/resource/digest 三个 watch loop 与 reindex);新增 `auto_dream`(4 步)、`proactive` 作业,移除旧 `dream`;file_catalog 增配 `resource` / `digest` / `dream` 实例;LLM 默认值与 Claude Code 凭据配置调整(tool_result_limit 50000、thinking_enable=false 等)。

### 11. 其它
- 新增 `steps/common/add.py`(`AddStep` 算术 demo)、`channel/__init__.py` 与 common `__init__` 导出整理。
- 新增 4 篇文档:`docs4/auto_dream_logic_and_step_refactor.md`、`docs4/watch_loop_step_refactor_plan.md`、`docs4/todo.md`,以及 `reme_design.md` 更新。
**
2026-06-19 01:35:31 +08:00

261 lines
9.3 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Integration test for the auto_memory job (single-step).
Drives the ``auto_memory`` step against a real LLM. Two scenarios:
1. **CREATE**: calls ``auto_memory`` with a fresh ``session_id`` and
conversation messages. Expects a new note with the key facts.
2. **UPDATE**: seeds an existing daily note, calls ``auto_memory`` with
the same ``session_id`` and new conversation messages. Expects the
old facts to survive and new facts to land.
Requires LLM_API_KEY (and optionally LLM_BASE_URL / LLM_MODEL_NAME) in the
environment or a .env file at the repo root. Hits the real LLM API.
"""
import asyncio
import sys
from pathlib import Path
INTEGRATION_DIR = Path(__file__).resolve().parent
sys.path.insert(0, str(INTEGRATION_DIR))
# pylint: disable=wrong-import-position
from _vault_fixture import vault_env # noqa: E402
SEED_STEM = "auth-middleware-rewrite"
SEED_BODY = """---
name: auth-middleware-rewrite
description: JWT auth middleware rewrite driven by legal/compliance requirements around session token storage
---
# 背景
- 项目JWT auth middleware 重写,替换旧的 session middleware
- 动机legal/compliance 要求,旧的 session token 存储方式不符合新合规要求
- 决策:采用 RS256 签名,密钥放在 KMSrefresh token 写 redis 集群
- 团队Alice 主导Bob 协助
# 时间线
- 2026-05-20 立项 kickoff
- 2026-05-23 设计评审通过
# 当下状态
- 进度:实现中
- 卡点:暂无
- 下一步:完成 refresh token 写入流程
"""
def _auth_messages() -> list[dict]:
"""Messages continuing the auth middleware thread."""
return [
{
"name": "user",
"role": "user",
"content": "状态更新PR #432auth middleware rewrite今天已经合并到 dev 分支,等待 staging 验收。",
},
{
"name": "assistant",
"role": "assistant",
"content": "好的已记录。staging 验收前要先跑回归测试吗?",
},
{
"name": "user",
"role": "user",
"content": (
"对。测试时发现 refresh token TTL 设 7d 在 redis 集群挂了——"
"redis maxmemory-policy 默认 allkeys-lru会随机驱逐 token导致用户被强制登出。"
),
},
{
"name": "assistant",
"role": "assistant",
"content": "理解,要切到 volatile-ttl 才能只驱逐带 TTL 的 key对吧",
},
{
"name": "user",
"role": "user",
"content": (
"对,下一步:周五 2026-05-29 前把 redis 配置改成 volatile-ttl 并重测," "blocked 在 SRE @lihua 的排期。"
),
},
]
def _pytorch_messages() -> list[dict]:
"""Messages about a brand-new pytorch topic."""
return [
{
"name": "user",
"role": "user",
"content": ("最近在调 pytorch 分布式训练。结论DDP 启动推荐用 torchrun" "比 mp.spawn 稳很多。"),
},
{
"name": "assistant",
"role": "assistant",
"content": "是因为信号处理的原因吗?",
},
{
"name": "user",
"role": "user",
"content": (
"主要是 NCCL backend 初始化更干净。mp.spawn 在 4 卡以上偶尔会卡死握手;"
"复现版本 pytorch 2.5.1 + nccl 2.21.5。"
),
},
{
"name": "user",
"role": "user",
"content": (
"另外 batch size 用 64*world_sizeper-rank lr 用 linear scaling rule "
"lr = base_lr * world_size"
),
},
{
"name": "user",
"role": "user",
"content": "先记一下。",
},
]
def _read_text(p: Path) -> str:
return p.read_text(encoding="utf-8")
def test_auto_memory_create():
"""CREATE a new note from scratch with a fresh session_id."""
async def run():
with vault_env() as env:
app = await env.make_app()
try:
today = env.today
print("\n" + "=" * 70)
print("[setup] vault_root =", env.vault_dir)
print("[setup] today =", today)
print("=" * 70)
pytorch_session_id = "pytorch-distributed-training"
expected_stem = pytorch_session_id
with env.record_agents(prefix="agent_create") as recorder:
response = await app.run_job(
"auto_memory",
messages=_pytorch_messages(),
session_id=pytorch_session_id,
)
dumped = await recorder.dump()
for p in dumped:
print(f"[CREATE] agent memory dumped: {p}")
assert response.success is True, f"CREATE job failed: {response.answer!r}"
meta = response.metadata or {}
assert meta.get("created") is True, f"Expected created=True, got {meta!r}"
assert meta.get("path") == f"daily/{today}/{expected_stem}.md"
pytorch_path = env.vault_dir / meta["path"]
assert pytorch_path.is_file(), f"created note not found at {pytorch_path}"
pytorch_text = _read_text(pytorch_path)
print("\n" + "=" * 70)
print(f"[CREATE] {pytorch_path} ({len(pytorch_text)} bytes)")
print(f"[CREATE] body:\n{pytorch_text}")
print("=" * 70)
topic_hits = [
needle
for needle in ("torchrun", "mp.spawn", "NCCL", "2.5.1", "linear scaling", "world_size")
if needle in pytorch_text
]
print(f"[CREATE] landed topic facts: {topic_hits}")
assert (
len(topic_hits) >= 3
), f"CREATE only captured {topic_hits!r} of expected facts\n--- CREATE ---\n{pytorch_text}"
stem = pytorch_path.stem
assert (
f"name: {stem}" in pytorch_text
), f"frontmatter name does not match stem {stem!r}\n{pytorch_text[:400]}"
print("\n" + "=" * 70)
print("test_auto_memory_create passed")
print("=" * 70)
finally:
await env.close_all()
asyncio.run(run())
def test_auto_memory_update():
"""UPDATE an existing note — old facts must survive, new facts must land."""
async def run():
with vault_env() as env:
app = await env.make_app()
try:
today = env.today
expected_stem = SEED_STEM
seed_path = env.seed_daily_note(expected_stem, SEED_BODY)
seed_before = _read_text(seed_path)
assert "legal/compliance" in seed_before
print("\n" + "=" * 70)
print("[setup] vault_root =", env.vault_dir)
print("[setup] today =", today)
print("[setup] seed_path =", seed_path)
print("=" * 70)
with env.record_agents(prefix="agent_update") as recorder:
response = await app.run_job(
"auto_memory",
messages=_auth_messages(),
session_id=SEED_STEM,
)
dumped = await recorder.dump()
for p in dumped:
print(f"[UPDATE] agent memory dumped: {p}")
assert response.success is True, f"UPDATE job failed: {response.answer!r}"
meta = response.metadata or {}
assert meta.get("created") is False, f"Expected created=False, got {meta!r}"
assert meta.get("path") == f"daily/{today}/{expected_stem}.md"
seed_after = _read_text(seed_path)
print("\n" + "=" * 70)
print(f"[UPDATE] {seed_path} ({len(seed_before)} -> {len(seed_after)} bytes)")
print(f"[UPDATE] body after:\n{seed_after}")
print("=" * 70)
for old_fact in ("legal/compliance", "RS256", "Alice"):
assert (
old_fact in seed_after
), f"UPDATE dropped pre-existing fact {old_fact!r}\n--- AFTER ---\n{seed_after}"
new_hits = [
needle
for needle in ("PR #432", "432", "volatile-ttl", "maxmemory-policy", "2026-05-29")
if needle in seed_after
]
print(f"[UPDATE] preserved old facts, landed new facts: {new_hits}")
assert (
len(new_hits) >= 2
), f"UPDATE only landed {new_hits!r} of expected new facts\n--- AFTER ---\n{seed_after}"
print("\n" + "=" * 70)
print("test_auto_memory_update passed")
print("=" * 70)
finally:
await env.close_all()
asyncio.run(run())
if __name__ == "__main__":
print("=== auto_memory integration test ===")
test_auto_memory_create()
test_auto_memory_update()
print("\nAll integration tests passed!")