mirror of
https://github.com/agentscope-ai/ReMe.git
synced 2026-09-09 22:31:05 +00:00
### 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` 更新。 **
230 lines
7.7 KiB
Python
230 lines
7.7 KiB
Python
"""Parser for YAML config with CLI argument overrides."""
|
|
|
|
import json
|
|
import os
|
|
import re
|
|
from pathlib import Path
|
|
from typing import Any
|
|
|
|
import yaml
|
|
|
|
# Config files are looked up relative to this module's directory
|
|
_CONFIG_DIR = Path(__file__).parent
|
|
# Extensions in priority order: yaml > yml > json when stems collide
|
|
_SUPPORTED_EXTS = (".yaml", ".yml", ".json")
|
|
_ENV_VAR_RE = re.compile(r"\$\{([A-Za-z_][A-Za-z0-9_]*)(?::-([^}]*))?}")
|
|
# Strings like "007" / "00501" must stay as strings, not be coerced to numbers
|
|
_LEADING_ZERO_RE = re.compile(r"^-?0\d")
|
|
|
|
|
|
def _repl(m: re.Match) -> str:
|
|
name: str = m.group(1)
|
|
# group(2) is None when the placeholder has no `:-default` part
|
|
default: str | None = m.group(2)
|
|
v = os.environ.get(name)
|
|
if v is None:
|
|
if default is not None:
|
|
return default
|
|
raise ValueError(f"Config references undefined env var: {name}")
|
|
return v
|
|
|
|
|
|
def _expand_env_vars(value: Any) -> Any:
|
|
"""Recursively expand `${VAR}` / `${VAR:-default}` placeholders in strings."""
|
|
if isinstance(value, str):
|
|
expanded = _ENV_VAR_RE.sub(_repl, value)
|
|
return _convert_value(expanded) if expanded != value else value
|
|
if isinstance(value, dict):
|
|
return {k: _expand_env_vars(v) for k, v in value.items()}
|
|
if isinstance(value, list):
|
|
return [_expand_env_vars(v) for v in value]
|
|
return value
|
|
|
|
|
|
def _discover_configs() -> dict[str, Path]:
|
|
"""Pre-scan config directory: maps file stem (name without ext) -> Path."""
|
|
discovered: dict[str, Path] = {}
|
|
if _CONFIG_DIR.is_dir():
|
|
# Sort by ext priority so registration order is deterministic across filesystems
|
|
files = sorted(
|
|
(p for p in _CONFIG_DIR.iterdir() if p.is_file() and p.suffix in _SUPPORTED_EXTS),
|
|
key=lambda p: (_SUPPORTED_EXTS.index(p.suffix), p.name),
|
|
)
|
|
for p in files:
|
|
discovered.setdefault(p.stem, p)
|
|
return discovered
|
|
|
|
|
|
_CONFIG_REGISTRY = _discover_configs()
|
|
|
|
|
|
def parse_dot_notation(dot_list: list[str]) -> dict:
|
|
"""Parse "key.subkey=value" strings into nested dict."""
|
|
result: dict = {}
|
|
for item in dot_list:
|
|
if "=" not in item:
|
|
raise ValueError(f"Invalid dot notation format (missing '='): {item}")
|
|
key_path, value_str = item.split("=", 1)
|
|
keys = key_path.split(".")
|
|
if not key_path or any(not key for key in keys):
|
|
raise ValueError(f"Invalid dot notation key: {key_path!r}")
|
|
current = result
|
|
for key in keys[:-1]:
|
|
if key in current and not isinstance(current[key], dict):
|
|
raise ValueError(f"Cannot set nested key '{key_path}': '{key}' is already a value")
|
|
current = current.setdefault(key, {})
|
|
# Symmetric to the prefix check above: refuse scalar-over-dict overwrite
|
|
last_key = keys[-1]
|
|
if last_key in current and isinstance(current[last_key], dict):
|
|
raise ValueError(f"Cannot overwrite nested dict at '{key_path}' with scalar value")
|
|
current[last_key] = _convert_value(value_str)
|
|
return result
|
|
|
|
|
|
def _convert_value(value_str: str) -> Any:
|
|
"""Convert string to appropriate Python type.
|
|
|
|
Only converts "true"/"false" (case-insensitive) to boolean.
|
|
Use JSON format (e.g., '"yes"', '"no"') to preserve these as strings.
|
|
Leading-zero strings (e.g., "007", "00501") are kept as strings.
|
|
"""
|
|
s = value_str.strip()
|
|
lower = s.lower()
|
|
|
|
# Handle special values (null, bool)
|
|
if lower in ("none", "null"):
|
|
return None
|
|
if lower == "true":
|
|
return True
|
|
if lower == "false":
|
|
return False
|
|
|
|
# Skip int/float for leading-zero strings to keep zip codes / ids intact
|
|
if not _LEADING_ZERO_RE.match(s):
|
|
for converter in (int, float):
|
|
try:
|
|
return converter(s)
|
|
except ValueError:
|
|
continue
|
|
|
|
# JSON handles lists, dicts, and explicitly-quoted strings
|
|
try:
|
|
return json.loads(s)
|
|
except (ValueError, json.JSONDecodeError):
|
|
pass
|
|
|
|
# Fallback to original string
|
|
return s
|
|
|
|
|
|
def _load_config(name_or_path: str, encoding: str = "utf-8") -> dict:
|
|
"""Load a YAML or JSON config file.
|
|
|
|
First check if name_or_path matches a pre-discovered config (key in _CONFIG_REGISTRY).
|
|
If not, treat as a file path and load directly.
|
|
"""
|
|
# 1. Try pre-discovered configs first
|
|
if name_or_path in _CONFIG_REGISTRY:
|
|
return _read_config_file(_CONFIG_REGISTRY[name_or_path], encoding)
|
|
|
|
# 2. Treat as file path
|
|
p = Path(name_or_path)
|
|
if p.suffix in _SUPPORTED_EXTS:
|
|
candidates = [p]
|
|
if not p.is_absolute():
|
|
candidates.append(_CONFIG_DIR / p)
|
|
for candidate in candidates:
|
|
if candidate.exists():
|
|
return _read_config_file(candidate, encoding)
|
|
raise FileNotFoundError(f"Config file not found: {p}")
|
|
|
|
known = ", ".join(sorted(_CONFIG_REGISTRY)) if _CONFIG_REGISTRY else "none"
|
|
raise FileNotFoundError(f"Config file not found: {name_or_path}. Available: {known}")
|
|
|
|
|
|
def _read_config_file(path: Path, encoding: str = "utf-8") -> dict:
|
|
"""Read YAML or JSON file based on extension. Expands ${ENV_VAR}."""
|
|
with path.open(encoding=encoding) as f:
|
|
if path.suffix == ".json":
|
|
result = json.load(f)
|
|
else:
|
|
result = yaml.safe_load(f)
|
|
if result is None:
|
|
return {}
|
|
if not isinstance(result, dict):
|
|
raise ValueError(f"Config root must be a mapping/object: {path}")
|
|
return _expand_env_vars(result)
|
|
|
|
|
|
def _deep_merge(base: dict, update: dict) -> dict:
|
|
"""Recursively merge dicts."""
|
|
result = base.copy()
|
|
for k, v in update.items():
|
|
if k in result and isinstance(result[k], dict) and isinstance(v, dict):
|
|
result[k] = _deep_merge(result[k], v)
|
|
else:
|
|
result[k] = v
|
|
return result
|
|
|
|
|
|
def _strip_arg_dashes(arg: str) -> str:
|
|
"""Strip a single leading `--` or `-` prefix (not all leading dashes)."""
|
|
if arg.startswith("--"):
|
|
return arg[2:]
|
|
if arg.startswith("-"):
|
|
return arg[1:]
|
|
return arg
|
|
|
|
|
|
def parse_args(*args) -> tuple[str, dict]:
|
|
"""Parse CLI args: first arg is action, rest are key=value pairs.
|
|
|
|
Usage: reme app config=paw.yaml service.name=test
|
|
Returns: (action, parsed_kv_dict)
|
|
"""
|
|
if not args:
|
|
raise ValueError("No arguments provided")
|
|
|
|
first = _strip_arg_dashes(args[0])
|
|
if "=" in first:
|
|
raise ValueError(f"First argument must be action, got: {args[0]}")
|
|
|
|
kvs: list[str] = []
|
|
for raw in args[1:]:
|
|
arg = _strip_arg_dashes(raw)
|
|
if "=" in arg:
|
|
kvs.append(arg)
|
|
else:
|
|
raise ValueError(f"Invalid argument format (expected key=value): {raw}")
|
|
|
|
parsed = parse_dot_notation(kvs) if kvs else {}
|
|
return first, parsed
|
|
|
|
|
|
def resolve_app_config(**kwargs) -> dict:
|
|
"""Resolve full app-start config: load `config=path` file, fall back to
|
|
`default`, then deep-merge with the remaining kwargs as overrides.
|
|
"""
|
|
from ..utils import get_logger
|
|
|
|
logger = get_logger()
|
|
configs: list[dict] = []
|
|
|
|
# `config=path` arrives as a string here; `config.foo=bar` arrives as a
|
|
# nested dict and is left in `kwargs` to be merged as a normal override.
|
|
config_value = kwargs.get("config")
|
|
if isinstance(config_value, str):
|
|
kwargs.pop("config")
|
|
logger.info(f"Loading config: {config_value}")
|
|
configs.append(_load_config(config_value))
|
|
elif "default" in _CONFIG_REGISTRY:
|
|
logger.info("No config specified, loading 'default'")
|
|
configs.append(_load_config("default"))
|
|
|
|
configs.append(kwargs)
|
|
|
|
merged: dict = {}
|
|
for cfg in configs:
|
|
merged = _deep_merge(merged, cfg)
|
|
|
|
return merged
|