deerflow-code/offline-backend-20260512/backend/scripts/migrate_memory_to_v2.py
2026-09-07 18:24:55 +08:00

260 lines
9.3 KiB
Python

"""一次性迁移:把旧版 memory.json 转换为 Memory V2 的 USER.md / MEMORY.md。
用法
====
::
# dry-run,只打印将要做什么,不写盘
PYTHONPATH=. python scripts/migrate_memory_to_v2.py --dry-run
# 迁移单个用户
PYTHONPATH=. python scripts/migrate_memory_to_v2.py --user-id default
# 迁移所有用户 (.deer-flow/users/*/memory.json 及 per-agent)
PYTHONPATH=. python scripts/migrate_memory_to_v2.py --all
转换规则 (详见 docs/MEMORY_V2_DESIGN_ZH.md §13)
================================================
- ``user.{workContext,personalContext,topOfMind}.summary`` -> USER.md 条目
- ``history.{recentMonths,earlierContext,longTermBackground}.summary`` -> USER.md 条目
(设计上 history 属于 Hindsight user bank,但 P1 期无 Hindsight,先落到 USER.md 保数据)
- ``facts[]`` category=preference -> USER.md;其余 (knowledge/context/behavior/goal) -> MEMORY.md
- per-agent memory.json (``users/{uid}/agents/{agent}/memory.json``) 的 facts 落到该 agent 的 MEMORY.md
迁移完成后,原 ``memory.json`` 重命名为 ``memory.json.bak`` 供回滚。
脚本幂等:已迁移 (存在 .bak) 的文件会被跳过。
Hindsight 灌库 (P2)
===================
本脚本只迁移本地文件。把这些记忆灌进 Hindsight bank 的步骤在 P2 期实现。
"""
from __future__ import annotations
import argparse
import json
import logging
import sys
from pathlib import Path
from typing import Any
from deerflow.agents.memory.providers.builtin import ENTRY_DELIMITER
from deerflow.config.paths import Paths, get_paths
logger = logging.getLogger(__name__)
# 落到 USER.md 的 fact category;其余都落到 MEMORY.md。
_USER_FACT_CATEGORIES: frozenset[str] = frozenset({"preference", "personal_info", "identity"})
# user / history 各 summary 字段及其在 USER.md 里的前缀。
_CONTEXT_SECTIONS: dict[str, dict[str, str]] = {
"user": {
"workContext": "[工作背景]",
"personalContext": "[个人背景]",
"topOfMind": "[当前关注]",
},
"history": {
"recentMonths": "[近期历史]",
"earlierContext": "[较早历史]",
"longTermBackground": "[长期背景]",
},
}
def _load_memory_json(path: Path) -> dict[str, Any] | None:
"""读取一个旧版 memory.json;不存在或解析失败返回 ``None``。"""
if not path.exists():
return None
try:
data = json.loads(path.read_text(encoding="utf-8"))
return data if isinstance(data, dict) else None
except (OSError, json.JSONDecodeError) as exc:
logger.warning("读取 %s 失败:%s", path, exc)
return None
def _extract_entries(data: dict[str, Any]) -> tuple[list[str], list[str]]:
"""把一份旧版 memory.json 拆成 (USER.md 条目, MEMORY.md 条目)。"""
user_entries: list[str] = []
memory_entries: list[str] = []
# user.* 与 history.* 的 summary
for section_key, fields in _CONTEXT_SECTIONS.items():
section = data.get(section_key) or {}
if not isinstance(section, dict):
continue
for field, prefix in fields.items():
node = section.get(field) or {}
summary = (node.get("summary") if isinstance(node, dict) else "") or ""
summary = summary.strip()
if summary:
user_entries.append(f"{prefix} {summary}")
# facts[] 按 category 路由
for fact in data.get("facts", []) or []:
if not isinstance(fact, dict):
continue
content = (fact.get("content") or "").strip()
if not content:
continue
category = (fact.get("category") or "context").strip().lower()
if category in _USER_FACT_CATEGORIES:
user_entries.append(content)
else:
memory_entries.append(content)
# 去重,保序
return list(dict.fromkeys(user_entries)), list(dict.fromkeys(memory_entries))
def _append_md_file(path: Path, new_entries: list[str], *, dry_run: bool) -> int:
"""把条目追加进一个 V2 markdown 文件 (USER.md / MEMORY.md),返回实际新增数。
已存在的条目跳过。dry-run 时不写盘。
"""
if not new_entries:
return 0
existing: list[str] = []
if path.exists():
raw = path.read_text(encoding="utf-8")
existing = [e.strip() for e in raw.split(ENTRY_DELIMITER) if e.strip()]
merged = list(existing)
added = 0
for entry in new_entries:
if entry not in merged:
merged.append(entry)
added += 1
if added and not dry_run:
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(ENTRY_DELIMITER.join(merged), encoding="utf-8")
return added
def _migrate_one(
memory_json: Path,
user_md: Path,
memory_md: Path,
*,
dry_run: bool,
) -> dict[str, Any]:
"""迁移一份 memory.json。返回一条迁移报告。"""
report: dict[str, Any] = {"source": str(memory_json), "action": "", "user_added": 0, "memory_added": 0}
bak = memory_json.with_suffix(memory_json.suffix + ".bak")
if bak.exists():
report["action"] = "skipped (已存在 .bak,视为已迁移)"
return report
data = _load_memory_json(memory_json)
if data is None:
report["action"] = "skipped (文件不存在或无法解析)"
return report
user_entries, memory_entries = _extract_entries(data)
report["user_added"] = _append_md_file(user_md, user_entries, dry_run=dry_run)
report["memory_added"] = _append_md_file(memory_md, memory_entries, dry_run=dry_run)
if not dry_run:
memory_json.rename(bak)
report["action"] = f"migrated -> {user_md.name}/{memory_md.name};原文件改名 {bak.name}"
else:
report["action"] = f"would migrate -> {user_md.name}/{memory_md.name}"
return report
def _discover_user_ids(paths: Paths) -> list[str]:
"""枚举 .deer-flow/users/ 下的所有 user_id。"""
users_root = paths.base_dir / "users"
if not users_root.is_dir():
return []
return sorted(d.name for d in users_root.iterdir() if d.is_dir())
def migrate_user(paths: Paths, user_id: str, *, dry_run: bool) -> list[dict[str, Any]]:
"""迁移单个用户的全局 memory.json 与所有 per-agent memory.json。"""
reports: list[dict[str, Any]] = []
# 全局 memory.json -> USER.md + default agent 的 MEMORY.md
global_json = paths.user_memory_file(user_id)
reports.append(
_migrate_one(
global_json,
paths.user_profile_md_file(user_id),
paths.user_agent_memory_md_file(user_id, "default"),
dry_run=dry_run,
)
)
# per-agent memory.json -> USER.md (preference facts) + 该 agent 的 MEMORY.md
agents_root = paths.user_dir(user_id) / "agents"
if agents_root.is_dir():
for agent_dir in sorted(agents_root.iterdir()):
if not agent_dir.is_dir():
continue
agent_json = agent_dir / "memory.json"
if not agent_json.exists():
continue
reports.append(
_migrate_one(
agent_json,
paths.user_profile_md_file(user_id),
paths.user_agent_memory_md_file(user_id, agent_dir.name),
dry_run=dry_run,
)
)
return reports
def main(argv: list[str] | None = None) -> int:
"""命令行入口。"""
logging.basicConfig(level=logging.INFO, format="%(message)s")
parser = argparse.ArgumentParser(description="把旧版 memory.json 迁移为 Memory V2 的 USER.md / MEMORY.md")
parser.add_argument("--dry-run", action="store_true", help="只打印将要做什么,不写盘")
parser.add_argument("--user-id", help="只迁移指定用户")
parser.add_argument("--all", action="store_true", help="迁移所有用户")
parser.add_argument("--base-dir", help="覆盖 .deer-flow 根目录 (默认从 DEER_FLOW_HOME 解析)")
parser.add_argument(
"--skip-hindsight",
action="store_true",
help="只迁移本地文件,不灌 Hindsight (P1 期本就如此,该选项为 P2 预留)",
)
args = parser.parse_args(argv)
if not args.all and not args.user_id:
parser.error("必须指定 --all 或 --user-id")
paths = Paths(base_dir=args.base_dir) if args.base_dir else get_paths()
if args.all:
user_ids = _discover_user_ids(paths)
if not user_ids:
logger.info("未发现任何用户目录,无需迁移。")
return 0
else:
user_ids = [args.user_id]
all_reports: list[dict[str, Any]] = []
for user_id in user_ids:
logger.info("=== 用户 %s ===", user_id)
for report in migrate_user(paths, user_id, dry_run=args.dry_run):
all_reports.append(report)
logger.info(
" %s | %s | USER +%d, MEMORY +%d",
report["action"],
report["source"],
report["user_added"],
report["memory_added"],
)
total_user = sum(r["user_added"] for r in all_reports)
total_memory = sum(r["memory_added"] for r in all_reports)
mode = "DRY-RUN(未写盘)" if args.dry_run else "已完成"
logger.info("迁移%s:USER.md 共新增 %d 条,MEMORY.md 共新增 %d 条。", mode, total_user, total_memory)
return 0
if __name__ == "__main__":
sys.exit(main())