OpenViking Usage Reporter / Sink 设计:从 Session Commit 解耦记忆使用事件上报的完整技术方案
【免费下载链接】OpenVikingSelf-evolving Context Database for AI Agents. Unify Agent Memory, Knowledge RAG and Skills.项目地址: https://gitcode.com/GitHub_Trending/op/OpenViking
本文基于 OpenViking 设计文档 openviking-usage-reporter-sink-design.md 展开,系统讲解 OpenViking 如何从 session commit 中解析出"记忆资源被检索/注入"这类使用行为,并把它转成结构化UsageEvent投递到可扩展的下游 Sink(自定义类或内置文件日志)。读完本文,你可以掌握:UsageEvent 协议与event_id幂等设计、MemoryUsageExtractor的识别规则、UsageSink动态加载与 5 秒超时隔离机制、内置文件日志 Sink 的计量协议,以及server.usage_reporter配置段的完整写法,并能在源码层面理解这套机制对 session commit 主链路的侵入边界。
1. 背景:Session Commit 中的使用行为线索
OpenViking 在 session commit 时会收到完整的 session messages。Agent runtime 调用 tool 后,会在 session messages 中留下 tool parts。部分 tool parts 可以表达"某个记忆文件被检索、被读取、被注入"等行为,这些正是记忆资源真实被消费的证据。
OpenViking 需要从 session 中解析出这些使用行为,并将结构化事件交给可扩展的下游。内核不绑定具体消息队列、Webhook、日志系统或数据库,而是提供通用的 Usage Reporter / Sink 扩展机制。
2. 设计选型:向 Observability 基础设施借鉴
"核心定义事件 + 插件式 Sink 输出"是开源基础设施里常见的设计模式:
- OpenTelemetry Collector把数据链路拆成 receivers、processors、exporters、pipelines,其中 exporters 专门负责把数据发送到不同 backend 或 destination(OTLP、Kafka、Prometheus、file 等),且配置 exporter 本身不代表启用,需要在 pipeline 中声明;
- Vector使用 sinks 概念,把 observability data 投递到不同目的地,例如 File、HTTP、Kafka、S3、ClickHouse、Prometheus remote write 等;
- Fluent Bit使用 Outputs 概念,官方定义里 Outputs 用来定义数据目的地,常见目的地包括远端服务、本地文件系统或标准接口,并且 Outputs 以 plugin 形式实现。
OpenViking 采用同一思路,核心价值是把"事件产生"和"事件去哪"解耦:内核只负责定义与抽取事件,投递目标完全由部署方通过 Sink 扩展决定。
3. 设计目标
- OpenViking 内核只定义
UsageEvent标准结构; - OpenViking 内核负责从 session commit 中解析
UsageEvent; - OpenViking 内核不绑定具体下游及其依赖;
- 部署方可以通过自定义 Sink 接入目标系统;
- Sink 失败默认不影响 session commit;
- 默认不上报完整 session,只上报结构化事件,并保留定位原始 session ToolPart 所需的证据信息。
4. 总体架构
数据流如下:
Agent runtime 调用 tool -> session message 留下 tool part -> client 上传 session 并 commit -> OpenViking archive session -> UsageExtractor 从 session messages 解析 UsageEvent -> UsageReporter 分发 UsageEvent -> UsageSink 写入目标系统模块拆分为四层:
UsageExtractor:负责从 session 里解析事件 UsageEvent:标准结构化事件 UsageReporter:负责分发事件 UsageSink:负责写入不同下游四个模块在仓库中的对应实现:
UsageEvent/UsageContext数据模型:models.py;UsageExtractor协议与MemoryUsageExtractor:extractors.py;UsageSink协议:sinks.py;UsageReporter分发器:reporter.py;- 配置构建与动态加载:config.py。
5. UsageEvent:内核与外部 Sink 之间的稳定协议
UsageEvent是 OpenViking 内核和外部 Sink 之间的稳定协议。示例:
{ "schema_version": "v1", "event_id": "ue_<sha256>", "event_type": "memory.injected", "resource_uri": "viking://user/test/memories/experiences/xxx.md", "resource_type": "experience", "account_id": "new", "user_id": "test", "session_id": "510bb5f9-4671-498e-adf4-27bb1b3691fe", "task_id": "b174eb56-e7d4-4fee-98a6-c53c0ddf62ed", "occurred_at": "2026-07-09T12:00:00Z", "evidence": { "archive_uri": "viking://user/test/sessions/510bb5f9/history/archive_001", "message_id": "msg_xxx", "tool_call_id": "call_xxx", "tool_name": "mcp__openviking__read" }, "attributes": {} }默认只上报结构化事件,不上报完整 session 内容。
字段语义上有两点值得注意:
resource_uri和resource_type描述被使用的资源,不限定为记忆文件;事件类型特有的数据写入attributes。UsageEvent是可独立传输和消费的完整事件。UsageContext(定义于 models.py 的account_id / user_id / session_id / archive_uri / task_id)只用于 Extractor 构造事件,不再重复传给 Sink。
当前MemoryUsageExtractor只接受属于UsageContext.user_id的规范 experience URI,其他用户 URI 不生成 UsageEvent。ToolPart 必须包含非空tool_id,无法稳定标识具体调用的 ToolPart 不进入统计。
5.1 event_id 的幂等性设计
event_id是以下字段规范序列化后的 SHA-256:
schema_version + event_type + account_id + user_id + session_id + evidence.message_id + evidence.tool_call_id + resource_uri在 models.py 中,UsageEvent.__post_init__将上述字段以json.dumps(..., separators=(",", ":"))紧凑序列化后做 SHA-256,并加ue_前缀。occurred_at、task_id、archive_uri和attributes不参与计算,避免重放时间差、任务恢复或附加属性变化破坏幂等性。同一 session message 中同一 tool call 对同一资源产生的事件,在 phase2 重放后仍得到相同event_id。
5.2 Experience 使用事件识别
插件使用 OpenViking 原生通用工具消费 Experience,不额外注册 Experience 专用工具:
- 成功的
find、search、list调用结果中出现 Experience URI,产生memory.recalled; - 成功的
read、multi_read调用实际读取 Experience URI,产生memory.injected; - 支持 OpenCode 的
openviking_*、OpenClaw 的ov_*和mcp__openviking__*命名空间形式;仅按受支持的工具名精确识别; - 通用工具返回其他记忆类型时,只保留当前用户
memories/experiences/目录下的规范文件 URI; - 只识别上述正式通用工具,不为未发布的专用工具名提供解析或配置兼容。
从源码看,这些规则在 extractors.py 中落实为两张精确映射表和一个 MCP 命名空间正则:
_RECALL_TOOL_OPERATIONS = { "find": "find", "search": "search", "list": "list", "openviking_find": "find", "openviking_search": "search", "openviking_list": "list", "ov_search": "search", "ov_list": "list", } _INJECTION_TOOL_OPERATIONS = { "read": "read", "multi_read": "multi_read", "openviking_read": "read", "openviking_multi_read": "multi_read", "ov_read": "read", "ov_multi_read": "multi_read", } _MCP_OPENVIKING_TOOL_RE = re.compile( r"^mcp__(?:openviking|plugin_.+_openviking)__(find|search|list|read|multi_read)$", re.IGNORECASE, )(见 extractors.py)
抽取过程中的过滤逻辑(MemoryUsageExtractor.extract,见 extractors.py):
- 遍历所有 message 的 parts,只处理
ToolPart; - 跳过
tool_id为空的 part,跳过tool_status != "completed"的 part(只统计成功调用); - 用
_tool_usage()精确匹配工具名,未命中直接跳过; - recall 类事件从
find/search/list的结构化输出或文本行中解析候选 URI,注入类事件从read/multi_read的tool_input/ 输出中解析 URI,并检查multi_read的逐 URI 成功标志与"读取失败"文本标记,失败的 URI 不产生memory.injected; - 所有候选 URI 经
_unique_uris()去重并归一化(_canonicalize_usage_uri会把历史的viking://user/memories/...简写与viking://~/memories/...主目录别名统一为调用者用户根下的规范 URI),再用_is_experience_uri()(extractors.py)做最终校验:必须位于当前用户memories/experiences/目录下、路径合法、且不是.abstract.md/.overview.md/.relations.json等 sidecar 文件。
6. UsageSink 机制
OpenViking 开源包定义统一的 Sink 抽象,实现见 sinks.py:
class UsageSink(Protocol): async def write(self, *, events: list[UsageEvent]) -> None: ...具体 Sink 可以作为外部扩展通过class_path动态加载,加载逻辑在 config.py 中:
import importlib def _load_class(class_path: str) -> type: module_name, class_name = class_path.rsplit(".", 1) module = importlib.import_module(module_name) cls = getattr(module, class_name) if not isinstance(cls, type): raise TypeError(f"{class_path} does not resolve to a class") return cls只有配置了该 Sink 时才 import 对应模块,OpenViking 内核不 import 或安装具体下游依赖。_build_sink会按type分派:file_log构造内置的FileLogUsageSink,custom则要求必填class_path并以config段作为构造参数传入。
6.1 超时、隔离与生命周期
在 reporter.py 中,UsageReporter的行为细节与设计目标一一对应:
- 每个 Sink 的
write()调用最多等待 5 秒(sink_timeout_seconds: float = 5.0),由asyncio.wait_for强制超时,见 reporter.py; - 超时或异常只记录日志(
logger.warning/logger.exception),不影响其他 Sink;多个 Sink 通过asyncio.gather并发分发、相互隔离; - Extractor 失败同样被捕获并记录日志,不会中断其余 Extractor(reporter.py);
- Reporter 在应用生命周期内只创建一次,应用退出时调用 Sink 可选的
close()方法,见 reporter.py。同步和异步close()均受同一 5 秒超时限制;同步 close hook 在独立 daemon 线程中执行,超时后不会阻塞事件循环、后续 Sink 清理或进程退出。
服务层的接线在 app.py:usage_reporter惰性构建(build_usage_reporter(config.usage_reporter)),通过set_usage_reporter注入 session 服务;session_service.py 将其转发给 Session,并在提交时透传(session_service.py)。应用退出时在 app.py 调用await usage_reporter.close()。
6.2 内置文件日志 Sink
开源包同时提供不依赖第三方消息队列 SDK 的内置文件日志 Sink(FileLogUsageSink,实现见 file_log_sink.py),供日志采集系统读取。它将事件立即追加到专用日志文件,不写默认 stdout,也不发起 HTTP 请求。
关键实现特性:
- UTC 小时滚动:使用
TimedRotatingFileHandler(when="H"、utc=True、delay=True),rotation_interval_hours默认 1,backup_count默认 168,即默认保留 168 个小时文件(file_log_sink.py); - 多 worker 安全滚动:多个 server worker 写入同一路径时,通过进程间文件锁(POSIX 下
fcntl.flock,Windows 下msvcrt.locking)串行化写入和滚动,锁文件为.<文件名>.lock(file_log_sink.py)。自定义 handler_ProcessSafeTimedRotatingFileHandler在每次 emit 前检测基文件是否被其他进程滚动重命名,并据此重新打开文件、同步 rollover 截止时间(file_log_sink.py); - Windows 兼容:Windows worker 写入后主动关闭文件句柄,避免其他进程滚动重命名失败(file_log_sink.py);
- 环境校验:构造时要求
resource_id_env指定的环境变量非空,rotation_interval_hours必须为正、backup_count不能为负,否则直接抛ValueError。
部署侧应将专用目录挂载到日志采集系统可见的宿主机路径。
6.3 文件日志计量协议
UsageEvent继续作为 OpenViking 内部抽取结果和自定义 Sink 的稳定协议。内置文件日志 Sink 将UsageEvent转换为计量接收端使用的扁平 JSON,每个事件写成一行:
{"event_time":"2026-08-05 11:30:00","tenant_id":"resource_id:ov-xxx;account_id:new;user_id:test;resource_uri:viking://user/test/memories/experiences/example.md","event_name":"experience.recall.count","object_id":"ue_<sha256>","count":1,"tags":{"resource_type":"experience"}}字段映射规则由_to_log_record(file_log_sink.py)实现:
| 字段 | 映射规则 |
|---|---|
event_time | 由occurred_at转换为 UTCYYYY-MM-DD HH:MM:SS |
tenant_id | 固定为resource_id:<resource_id>;account_id:<account_id>;user_id:<user_id>;resource_uri:<resource_uri>,其中resource_id取自环境变量 |
event_name | memory.recalled→experience.recall.count;memory.injected→experience.inject.count |
object_id | 使用稳定的event_id。下游以(tenant_id, object_id)作为复合去重键,不跨 tenant 单独按object_id去重 |
count | 固定为1 |
tags.resource_type | 固定记录被使用资源的类型;本期只产生experience使用事件 |
接收端按tenant_id、event_name和event_time范围过滤,并通过sum(count)聚合使用次数。
无法识别的event_type不生成含义不明确的计量记录,转换时抛出错误(ValueError: unsupported usage event type)并由 Reporter 的 best-effort 隔离机制处理。
7. 配置设计
配置模型定义在 server/config.py:UsageReporterConfig包含enabled(默认False)、extractors(默认["memory_usage"])、sinks列表,且extra: "forbid",即不支持未声明字段;UsageReporterSinkConfig包含type(custom或file_log)、可选class_path与config字典。
默认关闭:
server: usage_reporter: enabled: false自定义 Sink:
server: usage_reporter: enabled: true sinks: - type: custom class_path: example_usage.custom_sink.CustomUsageSink config: endpoint: https://usage.example.com/events内置文件日志 Sink:
server: usage_reporter: enabled: true sinks: - type: file_log config: path: /var/log/openviking_usage/usage.log resource_id_env: OV_RESOURCE_ID rotation_interval_hours: 1 backup_count: 168部署时必须设置resource_id_env指定的环境变量。Sink 使用其值构造tenant_id,保证不同 OpenViking resource 的数据相互隔离。
从 config.py 的构建逻辑看:enabled: false时build_usage_reporter直接返回None,服务不会创建任何 Extractor/Sink 实例;extractors目前只接受memory_usage,出现未知提取器名会直接抛ValueError。
8. 对 OpenViking 主链路的侵入点
侵入点控制在 4 个地方:
新增 config:增加
usage_reporter配置段(server/config.py);新增数据模型:增加
UsageEvent、UsageContext(models.py);session commit 增加 hook:在 session archive 成功后触发:
archive session success -> usage extractor -> usage reporter实现见 session.py 的
_run_usage_reporting:构造UsageContext(account、user、session、archive_uri、task_id)后调用reporter.extract_and_report。测试 test_session_usage_reporter.py 验证了端到端行为:在 session 中加入一次对 experience URI 的成功readToolPart,commit 完成后 task 结果为completed,usage_events_extracted == 1,Sink 收到一个memory.injected事件,且evidence.archive_uri、task_id与 commit 结果一致;新增 reporter/sink 模块:新增通用扩展点和 custom sink 动态加载能力。
不侵入的地方:
- 不改
find/search语义; - 不改
read语义; - 不把具体下游写死进 session commit;
- 不默认上传完整 session;
- 不强制写 MEMORY_FIELDS;
- 不强制写 search_tags;
- 不影响 snapshot;
- Sink 调用具有 5 秒超时边界,失败不会中断 phase2。
整体侵入属于低到中等,核心主链路只增加一个旁路 hook。
9. 可靠性策略:best-effort 投递语义
Usage Reporter 采用 best-effort 投递语义:
- Sink 成功:正常返回;
- Sink 失败或超时:记录日志,不影响 session commit。自定义 Sink 是否重试由其实现决定;
- 多个 Sink 相互隔离,某个 Sink 失败不影响其他 Sink;
- Sink 失败时事件可能丢失,因此本机制不保证 at-least-once;
- 如果 Sink 已写入成功,但进程在 phase2 写入完成标记前退出,phase2 恢复执行时可能重复发送同一事件;
- 每个事件包含稳定的
event_id。Sink 可将其作为 Kafka message key;消费端按(tenant_id, object_id)复合键去重; event_id只用于识别重复事件,不代表事件一定成功送达。
内置文件日志 Sink 在write()返回前完成本地追加,但不负责 TLS 采集、Kafka 投递或下游确认。文件写入、TLS 采集和下游消费任一阶段都可能在故障时产生丢失或重复,消费端需按(tenant_id, object_id)复合键去重,整体保持 best-effort 语义。
10. 小结
OpenViking 的 Usage Reporter / Sink 机制用"事件协议 + 插件式输出"回答了记忆资源使用度量的接入问题:内核只在 session commit 的旁路 hook 上做一次结构化抽取(本期仅 Experience 的 recall/injection 两类事件),通过 SHA-256 稳定的event_id保证重放幂等,再用 5 秒超时 + 异常隔离把 Sink 故障挡在 session commit 主链路之外。对于部署方,最小可用路径是启用file_logSink 并把日志目录交给既有采集系统;对于有自研计量平台的场景,实现async write(events=...)并通过class_path挂入自定义 Sink 即可,无需触碰 OpenViking 内核依赖。相关实现可继续深入 openviking/usage_reporter/ 目录、配置校验测试 test_usage_reporter_config.py 以及端到端测试 test_session_usage_reporter.py 查看。
【免费下载链接】OpenVikingSelf-evolving Context Database for AI Agents. Unify Agent Memory, Knowledge RAG and Skills.项目地址: https://gitcode.com/GitHub_Trending/op/OpenViking
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考