Просмотр исходного кода

重构(追踪): 独立认知记录存储职责

从 FileSystemTraceStore 中抽出 Cognition 读写实现,保留原有 TraceStore 方法、文件布局和导出路径,并新增认知存储的兼容与持久化测试。
SamLee 20 часов назад
Родитель
Сommit
8786491407

+ 6 - 0
agent/agent/trace/__init__.py

@@ -15,6 +15,7 @@ from .attachments import AttachmentRef
 from .protocols import TraceAttachmentStore, TraceStore
 from .protocols import TraceAttachmentStore, TraceStore
 from .store import FileSystemTraceStore
 from .store import FileSystemTraceStore
 from .trace_id import generate_trace_id, generate_sub_trace_id, parse_parent_trace_id
 from .trace_id import generate_trace_id, generate_sub_trace_id, parse_parent_trace_id
+from .tree_dump import StepTreeDumper, dump_json, dump_markdown, dump_tree
 
 
 __all__ = [
 __all__ = [
     # Models
     # Models
@@ -34,4 +35,9 @@ __all__ = [
     "generate_trace_id",
     "generate_trace_id",
     "generate_sub_trace_id",
     "generate_sub_trace_id",
     "parse_parent_trace_id",
     "parse_parent_trace_id",
+    # Debug helpers (kept for the documented compatibility import)
+    "StepTreeDumper",
+    "dump_tree",
+    "dump_json",
+    "dump_markdown",
 ]
 ]

+ 255 - 0
agent/agent/trace/_cognition_store.py

@@ -0,0 +1,255 @@
+"""Cognition-log persistence behind the FileSystemTraceStore facade."""
+
+from __future__ import annotations
+
+import json
+from datetime import datetime
+from pathlib import Path
+from typing import Any, Dict, List
+
+
+class TraceCognitionStore:
+    """Own cognition-log behavior without changing the public TraceStore API."""
+
+    def __init__(self, owner: Any) -> None:
+        self._owner = owner
+
+    def _get_cognition_log_file(self, trace_id: str) -> Path:
+        """获取 cognition_log.json 文件路径"""
+        return self._owner._get_trace_dir(trace_id) / "cognition_log.json"
+
+    def _get_knowledge_log_file(self, trace_id: str) -> Path:
+        """兼容旧接口:优先使用 cognition_log,回退到 knowledge_log"""
+        cognition_file = self._owner._get_cognition_log_file(trace_id)
+        if cognition_file.exists():
+            return cognition_file
+        legacy_file = self._owner._get_trace_dir(trace_id) / "knowledge_log.json"
+        if legacy_file.exists():
+            return legacy_file
+        return cognition_file  # 新建时用 cognition_log
+
+    async def get_cognition_log(self, trace_id: str) -> Dict[str, Any]:
+        """读取认知日志"""
+        log_file = self._owner._get_cognition_log_file(trace_id)
+        if log_file.exists():
+            return json.loads(log_file.read_text(encoding="utf-8"))
+        # 兼容旧格式:如果只有 knowledge_log.json,读取并转换
+        legacy_file = self._owner._get_trace_dir(trace_id) / "knowledge_log.json"
+        if legacy_file.exists():
+            return json.loads(legacy_file.read_text(encoding="utf-8"))
+        return {"trace_id": trace_id, "events": []}
+
+    async def get_knowledge_log(self, trace_id: str) -> Dict[str, Any]:
+        """兼容旧接口"""
+        log = await self._owner.get_cognition_log(trace_id)
+        # 旧格式用 entries,新格式用 events
+        if "entries" not in log and "events" in log:
+            log["entries"] = log["events"]
+        return log
+
+    async def append_cognition_event(
+        self,
+        trace_id: str,
+        event: Dict[str, Any],
+    ) -> None:
+        """追加认知事件到 cognition_log.json。
+
+        所有事件共有字段:
+            type: str         事件类型(见下表)
+            timestamp: str    ISO 格式时间戳(框架自动写入)
+
+        已定义的事件类型及典型字段:
+
+            type="query" — 知识注入查询(goal focus 时触发)
+                sequence, goal_id, query, response, source_ids, sources
+
+            type="evaluation" — 知识评估(Goal 完成/压缩前/任务结束触发)
+                knowledge_id, eval_result{relevance, utility, notes}, trigger_event
+
+            type="extraction_pending" — 反思侧分支暂存的待审核提取(Phase 1.2+)
+                extraction_id, sequence, goal_id, branch_id, payload
+                (payload 字段与 knowledge_save 参数一一对应)
+
+            type="extraction_reviewed" — 人工审核决策(CLI / HTTP API 写入)
+                extraction_id, decision("approve"/"edit"/"discard"), edited_payload?
+
+            type="extraction_committed" — 已上传到 KnowHub
+                extraction_id, knowledge_id
+
+            type="reflection" — Dream 的 per-trace 反思摘要(Phase 2.4 / 3.1)
+                sequence_range: [start, end]    本次反思覆盖的消息区间
+                summary: str                    LLM 生成的反思摘要
+                consumed_at: 可选, ISO 时间戳   当跨 trace 整合已消化此反思时写入
+
+        其他字段可按需附加,不做强校验(演进友好)。
+        """
+        log = await self._owner.get_cognition_log(trace_id)
+        if "events" not in log:
+            log["events"] = log.pop("entries", [])
+        event["timestamp"] = datetime.now().isoformat()
+        log["events"].append(event)
+        log_file = self._owner._get_cognition_log_file(trace_id)
+        log_file.write_text(
+            json.dumps(log, indent=2, ensure_ascii=False), encoding="utf-8"
+        )
+
+    async def append_knowledge_entry(
+        self,
+        trace_id: str,
+        knowledge_id: str,
+        goal_id: str,
+        injected_at_sequence: int,
+        task: str,
+        content: str,
+    ) -> None:
+        """兼容旧接口:追加知识注入记录(转换为 query 事件)"""
+        await self._owner.append_cognition_event(
+            trace_id=trace_id,
+            event={
+                "type": "query",
+                "sequence": injected_at_sequence,
+                "goal_id": goal_id,
+                "query": task,
+                "response": "",
+                "source_ids": [knowledge_id],
+                "sources": [
+                    {"id": knowledge_id, "task": task, "content": content[:500]}
+                ],
+            },
+        )
+
+    async def update_knowledge_evaluation(
+        self,
+        trace_id: str,
+        knowledge_id: str,
+        eval_result: Dict[str, Any],
+        trigger_event: str,
+    ) -> None:
+        """更新知识评估结果(兼容旧格式 + 新 cognition_log 格式)
+
+        旧格式:更新 entries[] 中匹配 knowledge_id 的条目的 eval_result
+        新格式:追加 evaluation 事件到 events[]
+        """
+        log = await self._owner.get_cognition_log(trace_id)
+        events = log.get("events", log.get("entries", []))
+
+        # 旧格式兼容:直接更新 entries 中的 eval_result 字段
+        if "entries" in log:
+            matching = [
+                (i, e)
+                for i, e in enumerate(log["entries"])
+                if e.get("knowledge_id") == knowledge_id
+                and e.get("eval_result") is None
+            ]
+            if matching:
+                matching.sort(
+                    key=lambda x: x[1].get("injected_at_sequence", 0), reverse=True
+                )
+                _, entry = matching[0]
+                entry["eval_result"] = eval_result
+                entry["evaluated_at"] = datetime.now().isoformat()
+                entry["evaluated_at_trigger"] = trigger_event
+                log_file = self._owner._get_knowledge_log_file(trace_id)
+                log_file.write_text(
+                    json.dumps(log, indent=2, ensure_ascii=False), encoding="utf-8"
+                )
+                return
+
+        # 新格式:追加 evaluation 事件
+        # 找到包含该 knowledge_id 的最近 query 事件
+        query_events = [
+            e
+            for e in events
+            if e.get("type") == "query" and knowledge_id in e.get("source_ids", [])
+        ]
+        query_sequence = query_events[-1]["sequence"] if query_events else None
+
+        await self._owner.append_cognition_event(
+            trace_id=trace_id,
+            event={
+                "type": "evaluation",
+                "sequence": max((e.get("sequence", 0) for e in events), default=0) + 1,
+                "query_sequence": query_sequence,
+                "trigger": trigger_event,
+                "assessments": [
+                    {
+                        "source_id": knowledge_id,
+                        "status": eval_result.get("eval_status", ""),
+                        "reason": eval_result.get("reason", ""),
+                    }
+                ],
+            },
+        )
+
+    async def get_pending_knowledge_entries(
+        self, trace_id: str
+    ) -> List[Dict[str, Any]]:
+        """获取所有待评估的知识条目(兼容旧格式 + 新格式)"""
+        log = await self._owner.get_cognition_log(trace_id)
+
+        # 旧格式
+        if "entries" in log:
+            return [e for e in log["entries"] if e.get("eval_result") is None]
+
+        # 新格式:找没有对应 evaluation 事件的 query 事件
+        events = log.get("events", [])
+        query_events = [e for e in events if e.get("type") == "query"]
+        eval_events = [e for e in events if e.get("type") == "evaluation"]
+
+        # 已评估的 query sequences
+        evaluated_sequences = {e.get("query_sequence") for e in eval_events}
+
+        pending = []
+        for qe in query_events:
+            if qe.get("sequence") not in evaluated_sequences:
+                # 转为旧格式兼容(runner 中的评估逻辑期望此格式)
+                for source in qe.get("sources", []):
+                    pending.append(
+                        {
+                            "knowledge_id": source.get("id", ""),
+                            "goal_id": qe.get("goal_id", ""),
+                            "injected_at_sequence": qe.get("sequence", 0),
+                            "task": source.get("task", ""),
+                            "content": source.get("content", ""),
+                            "query_sequence": qe.get("sequence"),
+                        }
+                    )
+        return pending
+
+    async def update_user_feedback(
+        self, trace_id: str, knowledge_id: str, user_feedback: Dict[str, Any]
+    ) -> None:
+        """记录用户对知识的反馈(confirm/override)"""
+        log = await self._owner.get_cognition_log(trace_id)
+
+        # 旧格式
+        if "entries" in log:
+            matching = [
+                (i, e)
+                for i, e in enumerate(log["entries"])
+                if e.get("knowledge_id") == knowledge_id
+            ]
+            if matching:
+                matching.sort(
+                    key=lambda x: x[1].get("injected_at_sequence", 0), reverse=True
+                )
+                _, entry = matching[0]
+                entry["user_feedback"] = user_feedback
+            log_file = self._owner._get_knowledge_log_file(trace_id)
+            log_file.write_text(
+                json.dumps(log, indent=2, ensure_ascii=False), encoding="utf-8"
+            )
+            return
+
+        # 新格式:追加 user_feedback 事件(或直接记录在 evaluation 上)
+        await self._owner.append_cognition_event(
+            trace_id=trace_id,
+            event={
+                "type": "user_feedback",
+                "knowledge_id": knowledge_id,
+                "feedback": user_feedback,
+            },
+        )
+
+
+__all__ = ["TraceCognitionStore"]

+ 57 - 211
agent/agent/trace/store.py

@@ -32,6 +32,7 @@ from typing import Dict, List, Optional, Any
 from .attachments import AttachmentRef
 from .attachments import AttachmentRef
 from .models import Trace, Message
 from .models import Trace, Message
 from .goal_models import GoalTree, Goal, GoalStats
 from .goal_models import GoalTree, Goal, GoalStats
+from ._cognition_store import TraceCognitionStore
 
 
 logger = logging.getLogger(__name__)
 logger = logging.getLogger(__name__)
 
 
@@ -42,6 +43,7 @@ class FileSystemTraceStore:
     def __init__(self, base_path: str = ".trace"):
     def __init__(self, base_path: str = ".trace"):
         self.base_path = Path(base_path)
         self.base_path = Path(base_path)
         self.base_path.mkdir(exist_ok=True)
         self.base_path.mkdir(exist_ok=True)
+        self._cognition = TraceCognitionStore(self)
 
 
     def _get_trace_dir(self, trace_id: str) -> Path:
     def _get_trace_dir(self, trace_id: str) -> Path:
         """获取 trace 目录"""
         """获取 trace 目录"""
@@ -299,7 +301,8 @@ class FileSystemTraceStore:
                     data["completed_at"] = datetime.fromisoformat(data["completed_at"])
                     data["completed_at"] = datetime.fromisoformat(data["completed_at"])
 
 
                 traces.append(Trace.from_dict(data))
                 traces.append(Trace.from_dict(data))
-            except Exception:
+            except Exception as exc:
+                logger.warning("Skipping unreadable trace metadata %s: %s", meta_file, exc)
                 continue
                 continue
 
 
         # 排序(最新的在前)
         # 排序(最新的在前)
@@ -318,7 +321,8 @@ class FileSystemTraceStore:
         try:
         try:
             data = json.loads(goal_file.read_text(encoding="utf-8"))
             data = json.loads(goal_file.read_text(encoding="utf-8"))
             return GoalTree.from_dict(data)
             return GoalTree.from_dict(data)
-        except Exception:
+        except Exception as exc:
+            logger.warning("Failed to load GoalTree %s: %s", goal_file, exc)
             return None
             return None
 
 
     async def update_goal_tree(self, trace_id: str, tree: GoalTree) -> None:
     async def update_goal_tree(self, trace_id: str, tree: GoalTree) -> None:
@@ -341,19 +345,22 @@ class FileSystemTraceStore:
         event_data = {"goal": goal.to_dict(), "parent_id": goal.parent_id}
         event_data = {"goal": goal.to_dict(), "parent_id": goal.parent_id}
         await self.append_event(trace_id, "goal_added", event_data)
         await self.append_event(trace_id, "goal_added", event_data)
 
 
-        # 打印详细的 goal 信息
         desc_preview = (
         desc_preview = (
             goal.description[:80] + "..."
             goal.description[:80] + "..."
             if len(goal.description) > 80
             if len(goal.description) > 80
             else goal.description
             else goal.description
         )
         )
-        print(f"[Goal Added] ID={goal.id}, Parent={goal.parent_id or 'root'}")
-        print(f"  📝 {desc_preview}")
+        logger.info(
+            "Goal added: id=%s parent=%s description=%s",
+            goal.id,
+            goal.parent_id or "root",
+            desc_preview,
+        )
         if goal.reason:
         if goal.reason:
             reason_preview = (
             reason_preview = (
                 goal.reason[:60] + "..." if len(goal.reason) > 60 else goal.reason
                 goal.reason[:60] + "..." if len(goal.reason) > 60 else goal.reason
             )
             )
-            print(f"  💡 {reason_preview}")
+            logger.debug("Goal reason: %s", reason_preview)
 
 
     async def update_goal(
     async def update_goal(
         self,
         self,
@@ -397,8 +404,11 @@ class FileSystemTraceStore:
             "goal_updated",
             "goal_updated",
             {"goal_id": goal_id, "updates": updates, "affected_goals": affected_goals},
             {"goal_id": goal_id, "updates": updates, "affected_goals": affected_goals},
         )
         )
-        print(
-            f"[DEBUG] Pushed goal_updated event: goal_id={goal_id}, updates={updates}, affected={len(affected_goals)}"
+        logger.debug(
+            "Goal update event appended: goal_id=%s updates=%s affected=%d",
+            goal_id,
+            updates,
+            len(affected_goals),
         )
         )
 
 
         # Goal 完成时触发知识评估
         # Goal 完成时触发知识评估
@@ -589,8 +599,6 @@ class FileSystemTraceStore:
             goal.self_stats.total_tokens += message.tokens
             goal.self_stats.total_tokens += message.tokens
         if message.cost:
         if message.cost:
             goal.self_stats.total_cost += message.cost
             goal.self_stats.total_cost += message.cost
-        # TODO: 更新 preview(工具调用摘要)
-
         # 更新自身 cumulative_stats
         # 更新自身 cumulative_stats
         goal.cumulative_stats.message_count += 1
         goal.cumulative_stats.message_count += 1
         if message.tokens:
         if message.tokens:
@@ -670,8 +678,8 @@ class FileSystemTraceStore:
                 try:
                 try:
                     data = json.loads(message_file.read_text(encoding="utf-8"))
                     data = json.loads(message_file.read_text(encoding="utf-8"))
                     return Message.from_dict(data)
                     return Message.from_dict(data)
-                except Exception:
-                    pass
+                except Exception as exc:
+                    logger.warning("Skipping unreadable message %s: %s", message_file, exc)
 
 
         return None
         return None
 
 
@@ -691,7 +699,8 @@ class FileSystemTraceStore:
                 data = json.loads(message_file.read_text(encoding="utf-8"))
                 data = json.loads(message_file.read_text(encoding="utf-8"))
                 msg = Message.from_dict(data)
                 msg = Message.from_dict(data)
                 messages.append(msg)
                 messages.append(msg)
-            except Exception:
+            except Exception as exc:
+                logger.warning("Skipping unreadable message %s: %s", message_file, exc)
                 continue
                 continue
 
 
         # 按 sequence 排序
         # 按 sequence 排序
@@ -899,7 +908,12 @@ class FileSystemTraceStore:
                     event = json.loads(line.strip())
                     event = json.loads(line.strip())
                     if event.get("event_id", 0) > since_event_id:
                     if event.get("event_id", 0) > since_event_id:
                         events.append(event)
                         events.append(event)
-                except Exception:
+                except Exception as exc:
+                    logger.warning(
+                        "Skipping unreadable event in %s: %s",
+                        events_file,
+                        exc,
+                    )
                     continue
                     continue
 
 
         return events
         return events
@@ -934,86 +948,26 @@ class FileSystemTraceStore:
 
 
         return event_id
         return event_id
 
 
-    # ===== Cognition Log 管理 =====
+    # ===== Cognition Log compatibility facade =====
 
 
     def _get_cognition_log_file(self, trace_id: str) -> Path:
     def _get_cognition_log_file(self, trace_id: str) -> Path:
-        """获取 cognition_log.json 文件路径"""
-        return self._get_trace_dir(trace_id) / "cognition_log.json"
+        return self._cognition._get_cognition_log_file(trace_id)
 
 
     def _get_knowledge_log_file(self, trace_id: str) -> Path:
     def _get_knowledge_log_file(self, trace_id: str) -> Path:
-        """兼容旧接口:优先使用 cognition_log,回退到 knowledge_log"""
-        cognition_file = self._get_cognition_log_file(trace_id)
-        if cognition_file.exists():
-            return cognition_file
-        legacy_file = self._get_trace_dir(trace_id) / "knowledge_log.json"
-        if legacy_file.exists():
-            return legacy_file
-        return cognition_file  # 新建时用 cognition_log
+        return self._cognition._get_knowledge_log_file(trace_id)
 
 
     async def get_cognition_log(self, trace_id: str) -> Dict[str, Any]:
     async def get_cognition_log(self, trace_id: str) -> Dict[str, Any]:
-        """读取认知日志"""
-        log_file = self._get_cognition_log_file(trace_id)
-        if log_file.exists():
-            return json.loads(log_file.read_text(encoding="utf-8"))
-        # 兼容旧格式:如果只有 knowledge_log.json,读取并转换
-        legacy_file = self._get_trace_dir(trace_id) / "knowledge_log.json"
-        if legacy_file.exists():
-            return json.loads(legacy_file.read_text(encoding="utf-8"))
-        return {"trace_id": trace_id, "events": []}
+        return await self._cognition.get_cognition_log(trace_id)
 
 
     async def get_knowledge_log(self, trace_id: str) -> Dict[str, Any]:
     async def get_knowledge_log(self, trace_id: str) -> Dict[str, Any]:
-        """兼容旧接口"""
-        log = await self.get_cognition_log(trace_id)
-        # 旧格式用 entries,新格式用 events
-        if "entries" not in log and "events" in log:
-            log["entries"] = log["events"]
-        return log
+        return await self._cognition.get_knowledge_log(trace_id)
 
 
     async def append_cognition_event(
     async def append_cognition_event(
         self,
         self,
         trace_id: str,
         trace_id: str,
         event: Dict[str, Any],
         event: Dict[str, Any],
     ) -> None:
     ) -> None:
-        """追加认知事件到 cognition_log.json。
-
-        所有事件共有字段:
-            type: str         事件类型(见下表)
-            timestamp: str    ISO 格式时间戳(框架自动写入)
-
-        已定义的事件类型及典型字段:
-
-            type="query" — 知识注入查询(goal focus 时触发)
-                sequence, goal_id, query, response, source_ids, sources
-
-            type="evaluation" — 知识评估(Goal 完成/压缩前/任务结束触发)
-                knowledge_id, eval_result{relevance, utility, notes}, trigger_event
-
-            type="extraction_pending" — 反思侧分支暂存的待审核提取(Phase 1.2+)
-                extraction_id, sequence, goal_id, branch_id, payload
-                (payload 字段与 knowledge_save 参数一一对应)
-
-            type="extraction_reviewed" — 人工审核决策(CLI / HTTP API 写入)
-                extraction_id, decision("approve"/"edit"/"discard"), edited_payload?
-
-            type="extraction_committed" — 已上传到 KnowHub
-                extraction_id, knowledge_id
-
-            type="reflection" — Dream 的 per-trace 反思摘要(Phase 2.4 / 3.1)
-                sequence_range: [start, end]    本次反思覆盖的消息区间
-                summary: str                    LLM 生成的反思摘要
-                consumed_at: 可选, ISO 时间戳   当跨 trace 整合已消化此反思时写入
-
-        其他字段可按需附加,不做强校验(演进友好)。
-        """
-        log = await self.get_cognition_log(trace_id)
-        if "events" not in log:
-            log["events"] = log.pop("entries", [])
-        event["timestamp"] = datetime.now().isoformat()
-        log["events"].append(event)
-        log_file = self._get_cognition_log_file(trace_id)
-        log_file.write_text(
-            json.dumps(log, indent=2, ensure_ascii=False), encoding="utf-8"
-        )
+        await self._cognition.append_cognition_event(trace_id, event)
 
 
     async def append_knowledge_entry(
     async def append_knowledge_entry(
         self,
         self,
@@ -1024,20 +978,13 @@ class FileSystemTraceStore:
         task: str,
         task: str,
         content: str,
         content: str,
     ) -> None:
     ) -> None:
-        """兼容旧接口:追加知识注入记录(转换为 query 事件)"""
-        await self.append_cognition_event(
-            trace_id=trace_id,
-            event={
-                "type": "query",
-                "sequence": injected_at_sequence,
-                "goal_id": goal_id,
-                "query": task,
-                "response": "",
-                "source_ids": [knowledge_id],
-                "sources": [
-                    {"id": knowledge_id, "task": task, "content": content[:500]}
-                ],
-            },
+        await self._cognition.append_knowledge_entry(
+            trace_id,
+            knowledge_id,
+            goal_id,
+            injected_at_sequence,
+            task,
+            content,
         )
         )
 
 
     async def update_knowledge_evaluation(
     async def update_knowledge_evaluation(
@@ -1047,130 +994,29 @@ class FileSystemTraceStore:
         eval_result: Dict[str, Any],
         eval_result: Dict[str, Any],
         trigger_event: str,
         trigger_event: str,
     ) -> None:
     ) -> None:
-        """更新知识评估结果(兼容旧格式 + 新 cognition_log 格式)
-
-        旧格式:更新 entries[] 中匹配 knowledge_id 的条目的 eval_result
-        新格式:追加 evaluation 事件到 events[]
-        """
-        log = await self.get_cognition_log(trace_id)
-        events = log.get("events", log.get("entries", []))
-
-        # 旧格式兼容:直接更新 entries 中的 eval_result 字段
-        if "entries" in log:
-            matching = [
-                (i, e)
-                for i, e in enumerate(log["entries"])
-                if e.get("knowledge_id") == knowledge_id
-                and e.get("eval_result") is None
-            ]
-            if matching:
-                matching.sort(
-                    key=lambda x: x[1].get("injected_at_sequence", 0), reverse=True
-                )
-                _, entry = matching[0]
-                entry["eval_result"] = eval_result
-                entry["evaluated_at"] = datetime.now().isoformat()
-                entry["evaluated_at_trigger"] = trigger_event
-                log_file = self._get_knowledge_log_file(trace_id)
-                log_file.write_text(
-                    json.dumps(log, indent=2, ensure_ascii=False), encoding="utf-8"
-                )
-                return
-
-        # 新格式:追加 evaluation 事件
-        # 找到包含该 knowledge_id 的最近 query 事件
-        query_events = [
-            e
-            for e in events
-            if e.get("type") == "query" and knowledge_id in e.get("source_ids", [])
-        ]
-        query_sequence = query_events[-1]["sequence"] if query_events else None
-
-        await self.append_cognition_event(
-            trace_id=trace_id,
-            event={
-                "type": "evaluation",
-                "sequence": max((e.get("sequence", 0) for e in events), default=0) + 1,
-                "query_sequence": query_sequence,
-                "trigger": trigger_event,
-                "assessments": [
-                    {
-                        "source_id": knowledge_id,
-                        "status": eval_result.get("eval_status", ""),
-                        "reason": eval_result.get("reason", ""),
-                    }
-                ],
-            },
+        await self._cognition.update_knowledge_evaluation(
+            trace_id,
+            knowledge_id,
+            eval_result,
+            trigger_event,
         )
         )
 
 
     async def get_pending_knowledge_entries(
     async def get_pending_knowledge_entries(
-        self, trace_id: str
+        self,
+        trace_id: str,
     ) -> List[Dict[str, Any]]:
     ) -> List[Dict[str, Any]]:
-        """获取所有待评估的知识条目(兼容旧格式 + 新格式)"""
-        log = await self.get_cognition_log(trace_id)
-
-        # 旧格式
-        if "entries" in log:
-            return [e for e in log["entries"] if e.get("eval_result") is None]
-
-        # 新格式:找没有对应 evaluation 事件的 query 事件
-        events = log.get("events", [])
-        query_events = [e for e in events if e.get("type") == "query"]
-        eval_events = [e for e in events if e.get("type") == "evaluation"]
-
-        # 已评估的 query sequences
-        evaluated_sequences = {e.get("query_sequence") for e in eval_events}
-
-        pending = []
-        for qe in query_events:
-            if qe.get("sequence") not in evaluated_sequences:
-                # 转为旧格式兼容(runner 中的评估逻辑期望此格式)
-                for source in qe.get("sources", []):
-                    pending.append(
-                        {
-                            "knowledge_id": source.get("id", ""),
-                            "goal_id": qe.get("goal_id", ""),
-                            "injected_at_sequence": qe.get("sequence", 0),
-                            "task": source.get("task", ""),
-                            "content": source.get("content", ""),
-                            "query_sequence": qe.get("sequence"),
-                        }
-                    )
-        return pending
+        return await self._cognition.get_pending_knowledge_entries(trace_id)
 
 
     async def update_user_feedback(
     async def update_user_feedback(
-        self, trace_id: str, knowledge_id: str, user_feedback: Dict[str, Any]
+        self,
+        trace_id: str,
+        knowledge_id: str,
+        user_feedback: Dict[str, Any],
     ) -> None:
     ) -> None:
-        """记录用户对知识的反馈(confirm/override)"""
-        log = await self.get_cognition_log(trace_id)
-
-        # 旧格式
-        if "entries" in log:
-            matching = [
-                (i, e)
-                for i, e in enumerate(log["entries"])
-                if e.get("knowledge_id") == knowledge_id
-            ]
-            if matching:
-                matching.sort(
-                    key=lambda x: x[1].get("injected_at_sequence", 0), reverse=True
-                )
-                _, entry = matching[0]
-                entry["user_feedback"] = user_feedback
-            log_file = self._get_knowledge_log_file(trace_id)
-            log_file.write_text(
-                json.dumps(log, indent=2, ensure_ascii=False), encoding="utf-8"
-            )
-            return
-
-        # 新格式:追加 user_feedback 事件(或直接记录在 evaluation 上)
-        await self.append_cognition_event(
-            trace_id=trace_id,
-            event={
-                "type": "user_feedback",
-                "knowledge_id": knowledge_id,
-                "feedback": user_feedback,
-            },
+        await self._cognition.update_user_feedback(
+            trace_id,
+            knowledge_id,
+            user_feedback,
         )
         )
 
 
 
 

+ 77 - 0
agent/tests/test_trace_cognition_store.py

@@ -0,0 +1,77 @@
+from __future__ import annotations
+
+import json
+
+import pytest
+
+from agent.trace.models import Trace
+from agent.trace.store import FileSystemTraceStore
+
+
+@pytest.mark.asyncio
+async def test_cognition_facade_preserves_new_event_flow(tmp_path) -> None:
+    store = FileSystemTraceStore(str(tmp_path))
+    await store.create_trace(Trace(trace_id="trace-new", mode="agent"))
+
+    await store.append_cognition_event(
+        "trace-new",
+        {
+            "type": "query",
+            "sequence": 4,
+            "goal_id": "goal-1",
+            "source_ids": ["knowledge-1"],
+            "sources": [
+                {
+                    "id": "knowledge-1",
+                    "task": "task",
+                    "content": "content",
+                }
+            ],
+        },
+    )
+
+    pending = await store.get_pending_knowledge_entries("trace-new")
+    assert [item["knowledge_id"] for item in pending] == ["knowledge-1"]
+
+    await store.update_knowledge_evaluation(
+        "trace-new",
+        "knowledge-1",
+        {"eval_status": "helpful", "reason": "used"},
+        "goal_complete",
+    )
+
+    assert await store.get_pending_knowledge_entries("trace-new") == []
+    compatibility_log = await store.get_knowledge_log("trace-new")
+    assert compatibility_log["entries"] is compatibility_log["events"]
+
+
+@pytest.mark.asyncio
+async def test_cognition_facade_preserves_legacy_log_update(tmp_path) -> None:
+    store = FileSystemTraceStore(str(tmp_path))
+    await store.create_trace(Trace(trace_id="trace-legacy", mode="agent"))
+    legacy_file = tmp_path / "trace-legacy" / "knowledge_log.json"
+    legacy_file.write_text(
+        json.dumps(
+            {
+                "trace_id": "trace-legacy",
+                "entries": [
+                    {
+                        "knowledge_id": "knowledge-legacy",
+                        "injected_at_sequence": 2,
+                        "eval_result": None,
+                    }
+                ],
+            }
+        ),
+        encoding="utf-8",
+    )
+
+    await store.update_knowledge_evaluation(
+        "trace-legacy",
+        "knowledge-legacy",
+        {"eval_status": "unused", "reason": "not needed"},
+        "goal_complete",
+    )
+
+    persisted = json.loads(legacy_file.read_text(encoding="utf-8"))
+    assert persisted["entries"][0]["eval_result"]["eval_status"] == "unused"