ReMe/reme4/utils/common_utils.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

249 lines
8.7 KiB
Python

"""Common utilities: hashing, async stream task execution, HTTP helpers."""
import asyncio
import hashlib
import json
import socket
import subprocess
import sys
import time
from collections.abc import AsyncGenerator, Callable
from contextlib import asynccontextmanager
from typing import Any, Literal
from .logger_utils import get_logger
from ..constants import REME_DEFAULT_HOST, REME_DEFAULT_PORT
from ..enumeration import ChunkEnum
from ..schema import StreamChunk
def hash_text(text: str, encoding: str = "utf-8") -> str:
"""Return SHA-256 hex digest of text."""
return hashlib.sha256(text.encode(encoding)).hexdigest()
def _format_chunk(
chunk: StreamChunk,
output_format: Literal["str", "bytes", "chunk"],
) -> str | bytes | StreamChunk:
"""Render a StreamChunk in the requested transport format."""
if output_format == "chunk":
return chunk
data = "data:[DONE]\n\n" if chunk.done else f"data:{chunk.model_dump_json()}\n\n"
return data.encode() if output_format == "bytes" else data
async def execute_stream_task(
stream_queue: asyncio.Queue[StreamChunk],
task: asyncio.Task[Any],
task_name: str | None = None,
output_format: Literal["str", "bytes", "chunk"] = "str",
) -> AsyncGenerator[str | bytes | StreamChunk, None]:
"""Yield chunks from stream_queue while monitoring task; cancels task on exit.
output_format: "str"/"bytes" emit SSE frames, "chunk" emits raw StreamChunk.
"""
logger = get_logger()
consumer: asyncio.Task[StreamChunk] | None = None
try:
while True:
consumer = get_chunk = asyncio.create_task(stream_queue.get())
done, _pending = await asyncio.wait({get_chunk, task}, return_when=asyncio.FIRST_COMPLETED)
# Producer still running — relay the next chunk and continue.
if task not in done:
chunk = get_chunk.result()
yield _format_chunk(chunk, output_format)
if chunk.done:
return
continue
# Producer finished. Capture any pending chunk, then stop the consumer wait
# so we can inspect task state safely.
pending_chunk: StreamChunk | None = None
if get_chunk in done:
pending_chunk = get_chunk.result()
else:
get_chunk.cancel()
try:
await get_chunk
except asyncio.CancelledError:
pass
# Surface task failure first — an exception trumps trailing data.
if task.cancelled():
msg = f"Task cancelled: {task_name}" if task_name else "Task cancelled"
raise asyncio.CancelledError(msg)
exc = task.exception()
if exc is not None:
log_msg = f"Task error in {task_name}: {exc}" if task_name else f"Task error: {exc}"
logger.error(log_msg, exc_info=exc)
raise exc
# Producer ended cleanly — flush pending + drain queue so no chunk is lost,
# then emit the terminal sentinel.
if pending_chunk is not None:
yield _format_chunk(pending_chunk, output_format)
if pending_chunk.done:
return
while not stream_queue.empty():
chunk = stream_queue.get_nowait()
yield _format_chunk(chunk, output_format)
if chunk.done:
return
yield _format_chunk(StreamChunk(chunk_type=ChunkEnum.DONE, chunk="", done=True), output_format)
return
finally:
# Cancel consumer wait if still pending (e.g. on consumer aclose).
if consumer is not None and not consumer.done():
consumer.cancel()
try:
await consumer
except asyncio.CancelledError:
pass
# Cancel producer task if still running to avoid resource leaks.
if not task.done():
task.cancel()
try:
await task
except asyncio.CancelledError:
pass
def _pick_free_port(host: str = REME_DEFAULT_HOST) -> int:
"""Bind to port 0 and return the OS-assigned free port."""
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
s.bind((host, 0))
return s.getsockname()[1]
async def _wait_reme_ready(host: str, port: int, timeout: float) -> None:
"""Poll find_reme until it reports 'reme' or timeout elapses."""
from .service_utils import find_reme
deadline = time.time() + timeout
while time.time() < deadline:
status = await find_reme(host, port)
if status == "reme":
return
await asyncio.sleep(0.2)
raise TimeoutError(f"ReMe service did not become ready at {host}:{port} within {timeout}s")
@asynccontextmanager
async def mock_reme_server(
host: str = REME_DEFAULT_HOST,
port: int | None = None,
config: str | None = None,
extra_args: list[str] | None = None,
startup_timeout: float = 120.0,
shutdown_timeout: float = 10.0,
log_to_file: bool = False,
enable_logo: bool = False,
):
"""Spawn `reme start` as a subprocess and yield (host, port) once ready.
Auto-picks a free port when port is None. Subprocess is terminated on exit.
"""
logger = get_logger()
if port is None:
port = _pick_free_port(host)
cmd: list[str] = [
sys.executable,
"-m",
"reme4.reme",
"start",
f"service.host={host}",
f"service.port={port}",
f"log_to_file={'true' if log_to_file else 'false'}",
f"enable_logo={'true' if enable_logo else 'false'}",
]
if config:
cmd.append(f"config={config}")
if extra_args:
cmd.extend(extra_args)
logger.info(f"Launching mock reme server: {' '.join(cmd)}")
proc = subprocess.Popen(
cmd,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
)
try:
await _wait_reme_ready(host, port, startup_timeout)
yield host, port
except Exception:
# Capture early-exit output for diagnostics.
if proc.poll() is not None and proc.stdout is not None:
tail = proc.stdout.read()
logger.error(f"reme server exited early. output:\n{tail}")
raise
finally:
if proc.poll() is None:
proc.terminate()
try:
proc.wait(timeout=shutdown_timeout)
except subprocess.TimeoutExpired:
logger.warning("reme server did not terminate gracefully, killing")
proc.kill()
proc.wait(timeout=shutdown_timeout)
if proc.stdout is not None:
try:
proc.stdout.close()
except Exception:
pass
async def call_action(
action: str,
host: str = REME_DEFAULT_HOST,
port: int = REME_DEFAULT_PORT,
timeout: float = 30.0,
**kwargs,
) -> dict | str:
"""POST to /{action}; return parsed JSON (dict) for JSON endpoints, raw text for SSE."""
from ..components.client.http_client import HttpClient
pieces: list[str] = []
async with HttpClient(host=host, port=port, timeout=timeout) as client:
async for chunk in client.stream_chunks(action, **kwargs):
payload = chunk.chunk
pieces.append(payload if isinstance(payload, str) else json.dumps(payload, ensure_ascii=False))
raw = "".join(pieces)
try:
return json.loads(raw)
except (ValueError, json.JSONDecodeError):
return raw
async def call_and_check(
action: str,
host: str = REME_DEFAULT_HOST,
port: int = REME_DEFAULT_PORT,
validator: Callable[[Any], bool] | None = None,
expected: Any = None,
timeout: float = 30.0,
**kwargs,
) -> Any:
"""Call action and verify response. Raises AssertionError on mismatch.
- validator(result) -> bool: custom predicate.
- expected: deep-equality target (compared to result, or to result[key] when expected is dict).
"""
result = await call_action(action, host=host, port=port, timeout=timeout, **kwargs)
if validator is not None and not validator(result):
raise AssertionError(f"validator rejected response for action={action!r}: {result!r}")
if expected is not None:
if isinstance(expected, dict) and isinstance(result, dict):
for k, v in expected.items():
if result.get(k) != v:
raise AssertionError(
f"action={action!r} expected {k}={v!r}, got {result.get(k)!r} (full: {result!r})",
)
elif result != expected:
raise AssertionError(f"action={action!r} expected {expected!r}, got {result!r}")
return result