Files
worldmodel/plans/PRISM/07_pipeline_D_consolidation.md
T

22 KiB
Raw Blame History

Chapter 07 — 管线 D:记忆巩固 (Memory Consolidation)

本章目标:机器人充电 / 空闲时跑一次"睡眠",把 delta/pending.jsonl反复确认的变化真正写入 LTM,并淘汰过时锚点、重训 3DGS、版本化备份。


7.1 为什么需要"睡眠"

如果 ZED 看到一次"沙发挪了"就立刻改 LTM,机器人就会变得易骗

  • 客人挪一下沙发拍照?被记成永久挪动
  • ZED 单帧检测错把椅子识成桌子?长期记忆被污染
  • 镜面区噪声偶尔产生"虚影"?写进去再也清不掉

仿照人脑:白天积累短期记忆 → 睡眠时筛选 → 仅把多次确认的信号转入长期记忆

PRISM 也一样:

flowchart LR
    subgraph DAY["白天(Online"]
        D1[("delta/pending.jsonl<br/><i>各种 DeltaEvent 累积,可能很乱</i>")]
    end
    subgraph NIGHT["夜里 / 充电时(Consolidation"]
        direction LR
        N1["pending"] --> N2["筛选"] --> N3["confirmed"] --> N4["应用到 LTM"] --> N5["版本化"]
    end
    DAY -. 触发 .-> NIGHT
    style DAY fill:#fff7d6,stroke:#c97a00
    style NIGHT fill:#d8e4ff,stroke:#1565c0

7.2 触发条件

触发 频率 模式
充电桩 + 静止 > 10 min 通常每天 1 次 完整巩固
手动命令 prism consolidate 按需 完整巩固
pending.jsonl > 5000 条 紧急 轻量巩固(只筛选不重训)
LTM 版本年龄 > 30 天 月级 深度巩固(含重算 anchors + CLIP

7.3 巩固总流程

⚠️ v1.5 升级:本节描述的"筛选 → 仲裁 → 应用 → 重训"线性流程及其 内部 deduplicate / merge 决策器(Step 13)的训练范式已被 § 7.13.1 增强为 self-augmentation 监督学习;本节流程本身仍可作为 baseline 使用, 决策器若由简单规则(PROMOTION_RULES + arbitrate())实现则完全不受影响。 仅当决策器升级为可训练策略网络(v1.5 → v2.0 路线图)时,§ 7.13.1 的 self-aug 范式才生效。

flowchart TB
    P[("pending.jsonl")]
    S1["<b>Step 1: 事件分类与筛选</b><br/>observation_count > N<br/>时间跨度 > T<br/>多 keyframe 多视角"]
    S2["<b>Step 2: 冲突仲裁</b><br/>iPhone 标 vs ZED 改"]
    S3["<b>Step 3: 应用到 LTM</b><br/>add / update / remove node"]
    S4["<b>Step 4: 锚点 / 索引 / CLIP</b><br/>重新计算受影响项"]
    S5["<b>Step 5: 增量 3DGS 重训</b><br/>只重训受影响房间"]
    S6["<b>Step 6: 版本化 + 健康检查</b><br/>snapshots/, validate()"]
    P --> S1 --> S2 --> S3 --> S4 --> S5 --> S6
    style P fill:#f5e1ff,stroke:#7b1fa2
    style S6 fill:#d4f0d4,stroke:#2e7d32

7.4 Step 1:事件筛选

# consolidation/filter.py
from datetime import timedelta

# 不同事件类型有不同"晋升门槛"
PROMOTION_RULES = {
    "object_moved":    dict(min_obs=5,  min_span_s=300,  min_views=2),
    "object_removed":  dict(min_obs=10, min_span_s=600,  min_views=3),
    "object_added":    dict(min_obs=8,  min_span_s=300,  min_views=2),
    "geometry_changed":dict(min_obs=15, min_span_s=900,  min_views=4),
}

def classify(ev: DeltaEvent, mem: SpatialMemory) -> str:
    rule = PROMOTION_RULES[ev.event_type]
    span = ev.last_observed - ev.first_observed
    unique_views = count_unique_viewpoints(ev.evidence, mem)

    # 若是 iPhone 的高 confidence 节点,门槛 1.5x
    target = mem.nodes.get(ev.target_uid) if ev.target_uid else None
    if target and target.source == "iphone" and target.confidence > 0.85:
        rule = {k: v*1.5 for k, v in rule.items()}

    if ev.observation_count >= rule["min_obs"] and \
       span >= rule["min_span_s"] and \
       unique_views >= rule["min_views"]:
        return "confirm"
    if span > 7*86400 and ev.observation_count < rule["min_obs"]//2:
        return "reject"           # 7 天还没攒够观测 → 拒绝
    return "keep"                  # 继续等待

关键直觉

  • min_obs:要"看过很多次"
  • min_span_s:必须跨越足够长时间(防止瞬时假象)
  • min_views:必须来自多个视点(防止单点死磕)

7.5 Step 2:冲突仲裁

不同 delta 之间,或 delta 与 LTM 之间可能冲突:

# consolidation/arbiter.py
def arbitrate(events: List[DeltaEvent], mem: SpatialMemory) -> List[DeltaEvent]:
    """对同一 target 的冲突事件做仲裁"""
    by_target = group_by_target(events)
    final = []
    for uid, evs in by_target.items():
        if len(evs) == 1:
            final.append(evs[0]); continue
        # 多个事件涉及同一节点
        sorted_evs = sorted(evs, key=lambda e: e.observation_count, reverse=True)
        primary = sorted_evs[0]
        # "搬走"+"挪到新位置" → 合并为一个 moved 事件
        if {e.event_type for e in evs} == {"object_removed","object_added"}:
            primary = merge_remove_add(evs)
        final.append(primary)
    return final


def merge_remove_add(evs):
    rem = next(e for e in evs if e.event_type == "object_removed")
    add = next(e for e in evs if e.event_type == "object_added")
    return DeltaEvent(
        event_id=uuid(),
        event_type="object_moved",
        target_uid=rem.target_uid,
        new_pose=add.new_pose,
        new_bbox=add.new_bbox,
        evidence=rem.evidence + add.evidence,
        observation_count=rem.observation_count + add.observation_count,
        first_observed=min(rem.first_observed, add.first_observed),
        last_observed=max(rem.last_observed, add.last_observed))

7.6 Step 3:应用到 LTM

# consolidation/apply.py
import shutil
from copy import deepcopy

def apply_events(events: List[DeltaEvent], mem: SpatialMemory) -> SpatialMemory:
    """返回应用后的新 SpatialMemory;不修改原对象"""
    new = deepcopy(mem)
    for ev in events:
        if ev.event_type == "object_moved":
            n = new.nodes[ev.target_uid]
            n.pose = ev.new_pose
            n.bbox_3d = ev.new_bbox
            n.last_seen = ev.last_observed
            n.confidence = min(1.0, n.confidence * 0.95)  # 轻微降低(被动过)
            # 重新计算 parent_room
            n.parent_room = find_room_for_point(n.pose.position[:2], new)
        elif ev.event_type == "object_removed":
            uid = ev.target_uid
            n = new.nodes[uid]
            # 不立刻硬删除,标记 deprecated 一段时间
            n.attributes["state"] = "removed"
            n.confidence *= 0.3
            # 切断关系
            new.edges = [e for e in new.edges
                         if e.src_uid != uid and e.dst_uid != uid]
        elif ev.event_type == "object_added":
            node = SpatialNode(**ev.payload)
            node.source = Source.FUSED   # 经过 consolidation
            node.confidence = 0.7        # 提升
            new.add_node(node)
            new.add_edge(SpatialEdge(node.parent_room, node.uid,
                                     "contains", source=Source.FUSED))
        elif ev.event_type == "geometry_changed":
            # 在 L2 TSDF / OctoMap 上打孔重建,不改场景图
            update_dense_layer_at_bbox(new, ev.new_bbox)
    return new

→ 注意 deepcopy + 整体替换 保证原子性:要么全部成功,要么回滚。


7.7 Step 4:锚点 / 索引 / CLIP 更新

# consolidation/refresh.py
def refresh_anchors_and_index(new_mem: SpatialMemory):
    # 1) 锚点重算
    new_mem.anchors = build_anchors(new_mem)        # 同 Chapter 04

    # 2) 受影响节点的 CLIP 重算
    for ev in events_just_applied:
        if ev.event_type in ("object_moved","object_added"):
            uid = ev.target_uid or ev.payload["uid"]
            node = new_mem.nodes[uid]
            # 渲染或从近期 keyframe 裁剪
            crops = collect_recent_crops(uid)
            node.clip_embedding = mean_clip_embedding(crops)

    # 3) Faiss 索引重建(L4 全量)
    build_faiss_index(new_mem)

    # 4) L3 房间 CLIP 受影响时也重算
    affected_rooms = {ev.target_uid_room for ev in events_just_applied}
    for r_uid in affected_rooms:
        new_mem.nodes[r_uid].clip_embedding = recompute_room_clip(r_uid)

7.8 Step 5:增量 3DGS / TSDF 重训

完全重训 3DGS 太贵(一房间 15 k iter × 5 min),所以只重训受影响房间 + 复用未变区域

def incremental_retrain_3dgs(new_mem: SpatialMemory, events):
    affected_rooms = set(find_room_for_point(ev.new_pose.position[:2], new_mem)
                         for ev in events if ev.new_pose)
    for r_uid in affected_rooms:
        # 收集该房间最近 24h 的关键帧
        kfs = list_keyframes_in_room(r_uid, since=time.time()-86400)
        # 从老 ckpt warmstart
        cmd = f"""ns-train splatfacto-bigtraining \
                  --pipeline.model.warmstart-ckpt outputs/{r_uid}/last.ckpt \
                  --max-num-iterations 3000 \
                  --data {kfs_dir}"""
        subprocess.run(cmd, shell=True)
        # 替换 LTM 中的 .ply
        shutil.move(f"outputs/{r_uid}/final.ply",
                    f"robot_memory/ltm/dense/3dgs/{r_uid}.ply")

def incremental_update_tsdf(new_mem, events):
    """TSDF 更便宜,直接在受影响 bbox 内重新融合最近关键帧"""
    for ev in events:
        bbox = ev.new_bbox if ev.new_bbox is not None else \
               new_mem.nodes[ev.target_uid].bbox_3d
        # 把 LTM TSDF 在该 bbox 内"清零"
        new_mem.dense.tsdf.reset_in_bbox(bbox)
        # 用最近 keyframes 重融合
        for kf in keyframes_observing_bbox(bbox, since_h=24):
            new_mem.dense.tsdf.integrate(kf)

7.9 Step 6:版本化 + 健康检查 + 切换

巩固结果不直接覆盖在线版本,先写到一个 staging,校验通过才切换:

# consolidation/commit.py
def commit(new_mem, old_dir, new_dir):
    # 1) 写 staging
    save(new_mem, f"{new_dir}/spatial_memory.json")
    copy_dense_dir(f"{old_dir}/dense", f"{new_dir}/dense")
    apply_dense_updates(new_dir, new_mem)

    # 2) 校验
    errs = validate(new_mem)
    sanity = run_sanity_relocalize(new_dir)  # 用 10 张历史关键帧重定位
    if errs or sanity.success_rate < 0.9:
        log.error(f"consolidation rejected: {errs}, success={sanity}")
        return False

    # 3) 原子切换
    timestamp = datetime.now().strftime("%Y%m%d_%H%M")
    snapshot = f"robot_memory/snapshots/{timestamp}.tar.zst"
    archive(old_dir, snapshot)
    os.rename(old_dir, f"{old_dir}.prev")
    os.rename(new_dir, old_dir)
    return True

7.10 锚点淘汰策略

def prune_stale_anchors(mem, max_age_days=60, min_recent_validations=2):
    keep = []
    for anc in mem.anchors:
        node = mem.nodes[anc.anchor_uid]
        age_days = (time.time() - anc.last_validated) / 86400
        # 仍是不可移动 + 最近 60 天被巩固确认过 → 保留
        if not anc.is_mobile and \
           node.attributes.get("state","") != "removed" and \
           age_days < max_age_days:
            keep.append(anc)
        elif node.observation_count >= min_recent_validations:
            keep.append(anc)
    mem.anchors = keep

→ 每次巩固跑一次。


7.11 完整 CLIprism consolidate

# tools/prism_consolidate.py
import click

@click.command()
@click.option("--mode", default="full",
              type=click.Choice(["full","light","deep"]))
@click.option("--ltm",  default="robot_memory/ltm")
@click.option("--delta", default="robot_memory/delta/pending.jsonl")
def main(mode, ltm, delta):
    mem = load(f"{ltm}/spatial_memory.json")
    raw_events = load_pending(delta)
    print(f"[consolidate-{mode}] start, ltm_nodes={len(mem.nodes)}, "
          f"pending={len(raw_events)}")

    # Step 1: filter
    decisions = [(ev, classify(ev, mem)) for ev in raw_events]
    confirmed = [ev for ev, d in decisions if d == "confirm"]
    rejected  = [ev for ev, d in decisions if d == "reject"]
    kept      = [ev for ev, d in decisions if d == "keep"]

    # Step 2: arbitrate
    final = arbitrate(confirmed, mem)

    # Step 3-4: apply + refresh
    staging = f"{ltm}.staging"
    new_mem = apply_events(final, mem)
    refresh_anchors_and_index(new_mem)

    # Step 5: densemode=full/deep 时才跑)
    if mode in ("full","deep"):
        incremental_update_tsdf(new_mem, final)
    if mode == "deep":
        incremental_retrain_3dgs(new_mem, final)

    # Step 6: commit
    ok = commit(new_mem, ltm, staging)
    if not ok:
        print("[consolidate] rolled back; pending.jsonl unchanged")
        return

    # 重写 pending:把 confirmed 移走,rejected 归档,kept 留下
    write_jsonl(f"robot_memory/delta/confirmed.jsonl", final, append=True)
    write_jsonl(f"robot_memory/delta/rejected.jsonl", rejected, append=True)
    write_jsonl(delta, kept)   # 覆盖
    print(f"[consolidate-{mode}] applied={len(final)} "
          f"rejected={len(rejected)} kept={len(kept)}")

运行:

# 充电时自动触发
systemd-timer --on=20:00 --weekly --command "prism consolidate --mode full"
# 月度深度巩固
systemd-timer --on=monthly --command "prism consolidate --mode deep"

7.12 安全机制

风险 缓解
巩固后机器人 reloc 失败 Step 6 的 sanity check 不通过 → 自动回滚
巩固过程中断电 staging 写完之前不替换 ltm/;半成品 staging 启动时清理
误删 iPhone 高 conf 节点 应用 remove 前再次检查 source==iphone && conf>0.85,若是则需 ≥ 10 倍证据
数据竞争(在线感知与巩固同时写) 巩固开始前发 set_readonly(true),在线感知此期间只允许写 delta/,不能切换 LTM 句柄

7.13 巩固后的指标记录

def log_metrics(events, mem_before, mem_after):
    with open("robot_memory/logs/consolidation.log", "a") as f:
        f.write(json.dumps({
            "ts": datetime.now().isoformat(),
            "applied": len(events),
            "nodes_before": len(mem_before.nodes),
            "nodes_after":  len(mem_after.nodes),
            "anchors_before": len(mem_before.anchors),
            "anchors_after":  len(mem_after.anchors),
            "ltm_version_new": mem_after.schema_version,
        }) + "\n")

便于后续画"机器人记忆演化"曲线。


7.13.1 v1.5 升级:Self-Augmentation 巩固训练

背景:train-test discrepancy

§ 7.47.7 的 dedup / merge / arbitrate 决策当前由简单规则 PROMOTION_RULES 阈值 + arbitrate() 启发式)做。v1.5 → v2.0 路线 图中,我们计划把这些规则升级为可训练的策略网络 policy_net: SceneRepr → {confirm, reject, merge, ...})。一旦走到 策略网络,就立刻撞上 Lyra 2.0 § 3.3 描述的 train-test discrepancy

  • 训练时:策略网络看到的是"干净的当前 L3 snapshot"——离线流水线 bundle adjustment 完毕、节点位姿无累计漂移、CLIP 嵌入由完整观测重算。
  • 推理时:策略网络面对的是"有累积误差的 snapshot"——白天 Pipeline C 在线写入、L2 已漂、CLIP 由匆忙的 ZED crop 算出、还混着第 50 次巩固 之后才逐渐显形的系统性偏差。

这种分布偏移导致策略网络在第 1–10 次巩固时表现良好、到第 50 次开始 误删高 confidence 节点或把两个真实独立的家具误合为一个。Lyra 2.0 § 3.3(b) 给出的解法是 self-augmentation:训练时主动把模型自己 之前的不完美输出当作输入,但监督信号仍用干净 ground-truth——模型 反复见到"自己会犯的错",因此学会自我纠错。完整原则陈述与 p_{\text{aug}} 取值讨论见 18_lyra_inspirations.md § 18.4

算法:consolidate_with_self_aug()

# consolidation/self_aug.py  (v1.5 新增,配合 policy_net 升级使用)
import numpy as np

def consolidate_with_self_aug(
    l3_clean:    SceneRepr,            # 当前离线巩固出的"干净"L3 快照
    history_l3s: List[SceneRepr],      # 过去 N 次巩固留下的不完美快照
    policy_net,                        # 待训练的 dedup/merge 决策器
    p_aug:       float = 0.7,
    t_max:       float = 0.5,
) -> torch.Tensor:
    """
    单步训练:以 p_aug 概率把"历史不完美 L3"当输入,
    监督信号始终是从 l3_clean 推得的 dedup/merge 目标。
    """
    # ---- 1. 决定本步是否使用 self-aug ----
    if np.random.rand() < p_aug and len(history_l3s) > 0:
        # 从历史中采样一份"曾经的不完美 snapshot"做基底
        t = np.random.uniform(0.0, t_max)   # 噪声强度,见下 t 的物理意义
        base = sample_history(history_l3s)
        corrupt_input = inject_noise_from_history(
            l3_clean, base, t=t)            # 见下"噪声注入策略"
    else:
        # 1 - p_aug = 0.3 概率用真实的"干净"输入做兜底
        corrupt_input = l3_clean

    # ---- 2. 前向 + 统一向"干净目标"对齐 ----
    pred    = policy_net(corrupt_input)                       # logits over actions
    targets = dedup_targets_from(l3_clean)                    # 监督一律取 clean
    loss    = cross_entropy(pred, targets)
    return loss   # 调用方负责 backward + step;典型 1000 iter / 场景(见参数表)

噪声注入策略:把历史 L3 当 "corruption source"

inject_noise_from_history(l3_clean, base, t) 不是凭空加 Gaussian,而是 l3_clean 退化成"看起来像 base 那个时点的中间状态"。具体扰动 按 t ∈ [0, 0.5] 线性加权应用以下三类:

# 扰动名 实现 模拟的真实失败模式
1 节点位置抖动 对每个节点 pose.positionN(0, σ_xyz)σ_xyz = t · 0.10 m L2 长走廊漂移导致 L3 节点中心偏移 ±10 cm
2 CLIP 嵌入扰动 clip_embeddingN(0, σ_clip) 后重新 L2 归一化,σ_clip = t · 0.05 ZED 暗光 / 模糊 crop 导致嵌入轻微偏离
3 节点随机丢弃 以概率 p_drop = t · 0.10(即 t=0.5 时丢 5%)随机移除非锚点节点 在线管线漏检小物品 / 被遮挡未上报

实现时三种扰动独立采样、叠加施加,并对应 history_l3s 中真实出现过 的偏差量级做了线性归一(最大扰动幅度对齐"第 50 次巩固"时的统计观测)。

参数表

参数 默认值 物理意义 / 来源
p_aug 0.7 Lyra 2.0 § 3.3(b) 报告的最佳点;> 0.5 让模型主要见自己的错,< 1.0 保留 30% 真实兜底,避免沉迷自生失败模式
t_max 0.5 t \in [0, 0.5] 表示"最多让历史看起来像巩固到一半时的中间状态";t=1.0 会让样本退化到几乎全噪声,监督信号失效
σ_xyz t · 0.10 m 与典型 L2 长走廊漂移上限对齐
σ_clip t · 0.05 与 ZED 暗光 crop 实测 CLIP 偏移分位数对齐
p_drop t · 0.10 与 Pipeline C 漏检率(v1.4 实测 ~5%)对齐
收敛迭代数 ~1000 iter / 场景 Lyra 用 7000 iter(视频扩散数据规模大);PRISM 单场景规模约小 5–10×,1000 iter 经验上够
训练批大小 1 场景 / step 场景图本身就是一个 graph batch,无需 mini-batch

评测建议(回链 13_evaluation.md

13_evaluation.md 的"长时一致性"指标族 (典型项:第 N 次巩固后的节点误删率、误合率、ghost-node 残留率)下, 工程预估:使用 self-aug 训练的策略网络在第 50 次巩固之后,

  • 误删高 confidence 节点率:下降 3050%baseline 无 self-aug 训练时通常 ~8% → 预期 45.5%)
  • 误合相邻独立家具率:下降 3040%
  • 第 1–10 次巩固指标基本持平(self-aug 主要解决长尾分布偏移,不解决初期能力问题)

ablation 计划:固定其他条件,对 p_aug ∈ {0.3, 0.5, 0.7, 0.9}t_max ∈ {0.25, 0.5, 0.75} 跑 2D 扫描,验证 (0.7, 0.5) 是否对 PRISM 真实数据也最优——如果最优点漂到 (0.5, 0.5),说明 Pipeline C 已足够稳定, self-aug 信号可以降权。

与原巩固管线的关系

本小节不替换 § 7.4–7.7 的任何步骤,只升级其中 dedup/merge 决策器的 训练范式——决策器接口、apply_events() 的下游调用、commit / 回滚 机制全部不变。旧版规则式决策器(PROMOTION_RULES + arbitrate())仍 可作为 baseline 与 self-aug 训练出的策略网络做 A/B。回链 18_lyra_inspirations.md § 18.4 "原则三:巩固期 Self-Augmentation"。


7.14 与 Chapter 06 的接口

Chapter 06 Chapter 07
写 LTM 严禁 唯一允许的写者
写 delta append 读取 + 重写 confirmed/rejected/kept
读 LTM 只读 只读 + 写 staging
频率 实时 离线(充电时)

两者通过 文件系统 解耦:在线感知不必知道巩固何时跑,巩固也不必停下感知。


7.15 本章小结

关键点 一句话
何时跑 充电 / 空闲;典型每天 1 次
筛选规则 observation_count + 时间跨度 + 视点多样性
冲突仲裁 同 target 多事件合并;iPhone 高 conf 项有更高门槛
应用 deepcopy + 整体替换,保证原子性
重训 仅受影响房间 + warmstart
安全 staging + sanity check + 回滚机制
比喻 像人脑睡眠——白天乱记,夜里整理

读完本章你应能:

  • 编写一个能跑通的 prism consolidate 脚本
  • 解释为什么要"延迟修改 LTM"
  • 设计安全的回滚机制

下一章 08_runtime_timeline.md 把 Chapter 04–07 的所有模块串成一个完整的端到端时序剧本,从 T0 到 T4。


章节版本v1.0 估计阅读时间18 分钟 关键收获:把"短期 delta"变成"长期 LTM"的完整安全流程