首页
/ Claude-Mem Server Beta 独立 BullMQ Observation 运行时:架构决策、Postgres 存储与持久化 Outbox 全解析

Claude-Mem Server Beta 独立 BullMQ Observation 运行时:架构决策、Postgres 存储与持久化 Outbox 全解析

2026-09-06 15:58:36作者:蔡怀权

本文基于 Claude-Mem 仓库中的实现计划文档 plans/2026-05-07-server-beta-independent-bullmq-observation-runtime.md,完整解析 claude-mem 13 Server(beta) 的核心架构决策:如何让 Server 侧摆脱对 legacy worker 的依赖,以 BullMQ/Valkey 作为规范队列、Postgres 作为规范 observation 存储,构建一条从 REST/MCP/hooks 事件写入、事务化 outbox、BullMQ 执行传输,到 provider 生成与检索的端到端独立运行时。读完后,你将掌握这套"Postgres outbox 为权威、BullMQ 仅做执行传输"的设计模式、完整的建表契约与幂等规则,以及当前仓库中已落地源码的对应位置。

1. 背景与执行级决策(Executive Decision)

Claude-Mem 的定位是"为每个 Agent 提供跨会话的持久化上下文":捕获会话中 Agent 的一切行为,经 AI 压缩后,把相关上下文注入未来会话。在计划文档撰写时,系统的 AI observation 生成能力完全由本地 worker 持有:/v1 服务端路由只写不生成,observation 生成链路(SessionManager → provider → ResponseProcessor)全部在 worker 进程内。

该计划(状态:implementation plan,发布日期 2026-05-07,发布目标:claude-mem 13 Server beta)给出的执行级决策是:

Server beta must own its runtime end to end:

REST/MCP/hooks -> Server beta HTTP/API layer -> BullMQ observation jobs -> provider generation -> server storage/search

核心约束有四条:

  1. worker 保持为稳定的 legacy 运行时原样保留,但 Server beta 不得依赖 WorkerService、worker HTTP 路由、worker 队列消费者或 worker 进程生命周期来生成 observations;
  2. BullMQ/Valkey 是 Server beta 的规范队列,Postgres 是规范 observation 存储;SQLite 仅作为 legacy worker/本地兼容存储保留;
  3. Redis/Valkey 只是运行时基础设施(承载 job、重试、并发与可观测性),不是 observations 的 source of truth;
  4. 计划关系上,它扩展了更早的 server 计划(Apache + BullMQ + 团队认证),并取代了此前允许 Server beta 以"包装/拷贝 WorkerService"方式实现 worker-parity 的旧计划中与之冲突的部分。

术语决策:一切围绕 "observation" 命名

Claude-Mem 的领域对象是 observation。Server beta 必须在用户可见 API、文档、job、存储命名、测试、日志与实现计划中统一保留这一措辞;"memory" 仅用于 worker 时代已存在且无法干净改名的遗留兼容名,或外部库/API 概念。新概念的命名映射如下:

遗留命名 Server beta 规范命名
memory_items observations
memory_sources observation_sources
MemoryItemsRepository ObservationRepository
泛化的 memory generation GenerateObservationsForEventJob
/v1/memories(如保留) 仅作为 observation 之上的兼容别名,规范面是 /v1/observations 与 observation 导向的 MCP 工具

2. Phase 0:文档勘察的发现与边界

计划的 Phase 0 记录了撰写时的现状勘察,这些"发现"是整个架构设计的动因,值得逐条保留:

现状发现:

  • 当时的 /v1 服务端路由把事件与直传 observation 记录存放在遗留 "memory" 命名的路由/仓储下:POST /v1/eventsPOST /v1/events/batchPOST /v1/memories 分别调用 AgentEventsRepository.create(...)MemoryItemsRepository.create(...),且当前不触发任何 provider 生成 job。当前仓库的服务端 v1 路由位于 src/server/routes/v1/,可对照查看。
  • 当时的 AI observation 生成路径完全由 worker 持有:SessionManager 通过 getMessageIterator(...) 消费排队消息;worker-service.ts 通过 startSessionProcessor(...) 启动 provider 会话;ResponseProcessorparseAgentXml(...) 解析 provider XML,并经 sessionStore.storeObservations(...) 落盘(对应源码位于 src/services/worker-service.ts 及 worker 子目录)。
  • BullMQ 官方文档确立了 Server beta 应直接使用的原语:Worker 将成功 job 移入 completed、失败 job 移入 failed;worker 必须挂 error 监听;支持 autorun: false;支持通过 options 设置并发;多 worker 是提升可用性的推荐方式;active job 在 worker 停止续约锁后可能 stalled 并被重试。
  • Better Auth 的 Express 集成要求 auth handler 必须挂载在 express.json() 之前,Express 5 下使用 /api/auth/*splat

允许复用的 API 与模式(计划明确列出):

  • src/services/server/Server.ts 与 Better Auth 文档借鉴 pre-body 路由挂载;
  • src/server/middleware/auth.tssrc/server/auth/api-key-service.ts 借鉴 API-key 认证(当前仓库 src/server/auth/src/server/middleware/ 已落地对应实现);
  • 仓储行为可借鉴,但 Server beta 仓储必须面向 Postgres 实现,不得复用 worker 的 legacy SessionStore;
  • 从 worker 各 provider(ClaudeProvider/GeminiProvider/OpenRouterProvider)复制请求构造,再把共享逻辑移入 src/server/generation(现已存在 src/server/generation/);
  • src/sdk/parser.ts 复制 XML 解析,从 ResponseProcessor 复制后处理规则;
  • 直接、原生地使用 BullMQ 的 QueueWorkerQueueEvents;
  • 保留 src/server/queue/redis-config.ts 中的 Valkey/Redis 健康检查与既有 Docker E2E 设施(当前仓库 E2E 脚本为 docker/e2e/server-e2e.mjsscripts/e2e-server-docker.sh)。

**反模式护栏(Anti-Pattern Guards)**是贯穿全部 Phase 的硬约束,Phase 0 的清单包括:

  • 不允许 Server beta 调用 new WorkerService();
  • 不允许 Server beta 的生成依赖 worker HTTP 路由类;
  • 不允许 /v1 成为"只写事件归档"却宣称 Server beta 能生成 observation;
  • 不允许 Server beta 生成使用 legacy SQLite pending-message 队列;
  • 不允许把规范 observation 记录存进 Redis;
  • 不允许移除或 destabilize 现有 worker;
  • 不允许在显式 BullMQ 模式下静默回退 SQLite;
  • 不允许把 Better Auth 挂载到 express.json() 之后。

3. 目标架构:运行时分离与生成主链路

计划定义的运行时分离边界是:

src/services/worker-service.ts
  Legacy worker runtime. 稳定的兼容路径,日后可引入共享 core 模块。

src/server/runtime/ServerBetaService.ts
  独立 server 运行时。持有 HTTP server、BullMQ 队列、provider 生成 worker、
  server 存储仓储、认证、健康检查与 Docker 部署。

Server beta 的完整生成主链路:

POST /v1/events
POST /v1/events/batch
Claude Code hook routed to Server beta
MCP observation_record_* tool
        |
        v
AgentEventsRepository transaction
        |
        v
ObservationGenerationJobRepository outbox row
        |
        v
BullMQ Queue.add(...)
        |
        v
BullMQ Worker processor
        |
        v
ProviderObservationGenerator
        |
        v
parseAgentXml / structured parser
        |
        v
ObservationRepository.create(...) + ObservationSourcesRepository.addSource(...)
        |
        v
QueueEvents/SSE/audit/search index update

这条链路的关键语义是:事件写入与 outbox job 行在同一事务中提交,提交后才 enqueue BullMQ job;Postgres outbox 记录"应该生成什么",BullMQ 只负责"执行"

4. Phase 1:Postgres Observation 存储地基

4.1 存储代码组织

Phase 1 要求把 Postgres 存储集中在(当前仓库均已落地于 src/storage/postgres/):

  • config.ts — 环境变量解析、连接池规模、超时与 SSL 设置(环境入口变量 CLAUDE_MEM_SERVER_DATABASE_URL);
  • pool.ts — 共享 pg.Pool 工厂、健康检查、事务与优雅关停(关停时 drain 并关闭连接池);
  • schema.ts — 迁移/引导 SQL 与 schema 版本常量;
  • index.ts — Server beta 运行时装配使用的导出。

另有两个硬性要求:启动校验在"需要 Postgres 但不可用"时直接让 Server beta 启动失败(不得静默回退 SQLite);Node/Bun 包清单中加入 pg@types/pg 依赖。迁移/引导 helper 要求创建全部 schema/表/索引、记录已应用版本,且可重复执行(启动时与测试中均安全)。

4.2 规范表清单

规范表共 10 张:teamsprojectsteam_membersapi_keysaudit_logserver_sessionsagent_eventsobservationsobservation_sourcesobservation_generation_jobsobservation_generation_job_events

4.3 初始 Schema 契约(完整 SQL)

计划明确要求该契约在 Phase 1 迁移中显式实现,列名只有在同步更新全部仓储契约与测试时才允许调整。以下为计划文档中的完整建表 DDL,原样继承:

CREATE TABLE teams (
  id TEXT PRIMARY KEY,
  name TEXT NOT NULL,
  metadata JSONB NOT NULL DEFAULT '{}'::jsonb,
  created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);

CREATE TABLE projects (
  id TEXT PRIMARY KEY,
  team_id TEXT NOT NULL REFERENCES teams(id) ON DELETE CASCADE,
  name TEXT NOT NULL,
  metadata JSONB NOT NULL DEFAULT '{}'::jsonb,
  created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  UNIQUE (id, team_id)
);

CREATE TABLE team_members (
  team_id TEXT NOT NULL REFERENCES teams(id) ON DELETE CASCADE,
  user_id TEXT NOT NULL,
  role TEXT NOT NULL,
  metadata JSONB NOT NULL DEFAULT '{}'::jsonb,
  created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  PRIMARY KEY (team_id, user_id)
);

CREATE TABLE api_keys (
  id TEXT PRIMARY KEY,
  key_hash TEXT NOT NULL UNIQUE,
  team_id TEXT REFERENCES teams(id) ON DELETE CASCADE,
  project_id TEXT REFERENCES projects(id) ON DELETE CASCADE,
  actor_id TEXT NOT NULL,
  scopes JSONB NOT NULL DEFAULT '[]'::jsonb,
  revoked_at TIMESTAMPTZ,
  expires_at TIMESTAMPTZ,
  created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  CHECK (project_id IS NULL OR team_id IS NOT NULL),
  FOREIGN KEY (project_id, team_id) REFERENCES projects(id, team_id) ON DELETE CASCADE
);

CREATE TABLE audit_log (
  id TEXT PRIMARY KEY,
  team_id TEXT REFERENCES teams(id) ON DELETE SET NULL,
  project_id TEXT REFERENCES projects(id) ON DELETE SET NULL,
  actor_id TEXT,
  api_key_id TEXT REFERENCES api_keys(id) ON DELETE SET NULL,
  action TEXT NOT NULL,
  resource_type TEXT NOT NULL,
  resource_id TEXT,
  details JSONB NOT NULL DEFAULT '{}'::jsonb,
  created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  CHECK (project_id IS NULL OR team_id IS NOT NULL),
  FOREIGN KEY (project_id, team_id) REFERENCES projects(id, team_id) ON DELETE SET NULL
);

CREATE TABLE server_sessions (
  id TEXT PRIMARY KEY,
  project_id TEXT NOT NULL REFERENCES projects(id) ON DELETE CASCADE,
  team_id TEXT NOT NULL REFERENCES teams(id) ON DELETE CASCADE,
  external_session_id TEXT,
  content_session_id TEXT,
  agent_id TEXT,
  agent_type TEXT,
  platform_source TEXT,
  generation_status TEXT NOT NULL DEFAULT 'idle',
  metadata JSONB NOT NULL DEFAULT '{}'::jsonb,
  started_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  ended_at TIMESTAMPTZ,
  last_generated_at TIMESTAMPTZ,
  created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  UNIQUE (project_id, external_session_id),
  FOREIGN KEY (project_id, team_id) REFERENCES projects(id, team_id) ON DELETE CASCADE
);

CREATE TABLE agent_events (
  id TEXT PRIMARY KEY,
  project_id TEXT NOT NULL REFERENCES projects(id) ON DELETE CASCADE,
  team_id TEXT NOT NULL REFERENCES teams(id) ON DELETE CASCADE,
  server_session_id TEXT REFERENCES server_sessions(id) ON DELETE SET NULL,
  source_adapter TEXT NOT NULL,
  source_event_id TEXT,
  idempotency_key TEXT NOT NULL,
  event_type TEXT NOT NULL,
  payload JSONB NOT NULL,
  metadata JSONB NOT NULL DEFAULT '{}'::jsonb,
  occurred_at TIMESTAMPTZ NOT NULL,
  received_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  UNIQUE (idempotency_key),
  UNIQUE (id, project_id, team_id),
  FOREIGN KEY (project_id, team_id) REFERENCES projects(id, team_id) ON DELETE CASCADE
);

CREATE TABLE observation_generation_jobs (
  id TEXT PRIMARY KEY,
  project_id TEXT NOT NULL REFERENCES projects(id) ON DELETE CASCADE,
  team_id TEXT NOT NULL REFERENCES teams(id) ON DELETE CASCADE,
  agent_event_id TEXT REFERENCES agent_events(id) ON DELETE CASCADE,
  source_type TEXT NOT NULL CHECK (source_type IN ('agent_event', 'session_summary', 'observation_reindex')),
  source_id TEXT NOT NULL,
  server_session_id TEXT REFERENCES server_sessions(id) ON DELETE SET NULL,
  job_type TEXT NOT NULL,
  status TEXT NOT NULL CHECK (status IN ('queued', 'processing', 'completed', 'failed', 'cancelled')),
  idempotency_key TEXT NOT NULL UNIQUE,
  bullmq_job_id TEXT UNIQUE,
  attempts INTEGER NOT NULL DEFAULT 0,
  max_attempts INTEGER NOT NULL DEFAULT 3,
  next_attempt_at TIMESTAMPTZ,
  locked_at TIMESTAMPTZ,
  locked_by TEXT,
  completed_at TIMESTAMPTZ,
  failed_at TIMESTAMPTZ,
  cancelled_at TIMESTAMPTZ,
  last_error JSONB,
  payload JSONB NOT NULL DEFAULT '{}'::jsonb,
  created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  UNIQUE (team_id, project_id, source_type, source_id, job_type),
  CHECK (
    (source_type = 'agent_event' AND agent_event_id IS NOT NULL AND source_id = agent_event_id)
    OR
    (source_type = 'session_summary' AND agent_event_id IS NULL AND server_session_id IS NOT NULL AND source_id = server_session_id)
    OR
    (source_type = 'observation_reindex' AND agent_event_id IS NULL)
  ),
  FOREIGN KEY (agent_event_id, project_id, team_id) REFERENCES agent_events(id, project_id, team_id) ON DELETE CASCADE,
  FOREIGN KEY (project_id, team_id) REFERENCES projects(id, team_id) ON DELETE CASCADE
);

CREATE TABLE observations (
  id TEXT PRIMARY KEY,
  project_id TEXT NOT NULL REFERENCES projects(id) ON DELETE CASCADE,
  team_id TEXT NOT NULL REFERENCES teams(id) ON DELETE CASCADE,
  server_session_id TEXT REFERENCES server_sessions(id) ON DELETE SET NULL,
  kind TEXT NOT NULL DEFAULT 'observation',
  content TEXT NOT NULL,
  content_search TSVECTOR GENERATED ALWAYS AS (to_tsvector('english', content)) STORED,
  generation_key TEXT,
  metadata JSONB NOT NULL DEFAULT '{}'::jsonb,
  embedding JSONB,
  created_by_job_id TEXT REFERENCES observation_generation_jobs(id) ON DELETE SET NULL,
  created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  UNIQUE (team_id, project_id, generation_key),
  FOREIGN KEY (project_id, team_id) REFERENCES projects(id, team_id) ON DELETE CASCADE
);

CREATE TABLE observation_sources (
  id TEXT PRIMARY KEY,
  observation_id TEXT NOT NULL REFERENCES observations(id) ON DELETE CASCADE,
  agent_event_id TEXT REFERENCES agent_events(id) ON DELETE CASCADE,
  generation_job_id TEXT REFERENCES observation_generation_jobs(id) ON DELETE SET NULL,
  source_type TEXT NOT NULL CHECK (source_type IN ('agent_event', 'session_summary', 'observation_reindex', 'manual')),
  source_id TEXT NOT NULL,
  metadata JSONB NOT NULL DEFAULT '{}'::jsonb,
  created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
  UNIQUE (observation_id, source_type, source_id),
  UNIQUE (source_type, source_id, generation_job_id, observation_id),
  CHECK (
    (source_type = 'agent_event' AND agent_event_id IS NOT NULL AND source_id = agent_event_id)
    OR
    (source_type <> 'agent_event' AND agent_event_id IS NULL)
  )
);

CREATE TABLE observation_generation_job_events (
  id TEXT PRIMARY KEY,
  generation_job_id TEXT NOT NULL REFERENCES observation_generation_jobs(id) ON DELETE CASCADE,
  event_type TEXT NOT NULL CHECK (event_type IN ('queued', 'enqueued', 'processing', 'retry_scheduled', 'completed', 'failed', 'cancelled')),
  status_after TEXT NOT NULL CHECK (status_after IN ('queued', 'processing', 'completed', 'failed', 'cancelled')),
  attempt INTEGER NOT NULL DEFAULT 0,
  details JSONB NOT NULL DEFAULT '{}'::jsonb,
  created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);

CREATE INDEX idx_agent_events_project_session ON agent_events(project_id, server_session_id, occurred_at);
CREATE INDEX idx_projects_team ON projects(team_id, id);
CREATE INDEX idx_agent_events_team_project ON agent_events(team_id, project_id, occurred_at);
CREATE INDEX idx_observations_project_session ON observations(project_id, server_session_id, created_at);
CREATE INDEX idx_observations_team_project ON observations(team_id, project_id, created_at);
CREATE INDEX idx_observations_content_search ON observations USING GIN (content_search);
CREATE INDEX idx_observation_sources_event ON observation_sources(agent_event_id);
CREATE INDEX idx_observation_sources_source ON observation_sources(source_type, source_id);
CREATE INDEX idx_observation_jobs_status_next_attempt ON observation_generation_jobs(status, next_attempt_at, created_at);
CREATE INDEX idx_observation_jobs_team_project ON observation_generation_jobs(team_id, project_id, status, created_at);
CREATE INDEX idx_observation_jobs_event ON observation_generation_jobs(agent_event_id);
CREATE INDEX idx_observation_jobs_source ON observation_generation_jobs(source_type, source_id);
CREATE INDEX idx_observation_job_events_job_created ON observation_generation_job_events(generation_job_id, created_at);
CREATE INDEX idx_audit_log_scope_created ON audit_log(project_id, team_id, created_at);

4.4 事件/Outbox 关系与所有权规则

Schema 之外,计划定义了一组语义规则,它们决定了整套系统的正确性:

  • agent_events 是原始摄入事件的规范表;每个 project 恰好归属一个 team(经 projects.team_id),Postgres 规范存储中没有无主/默认 project 模式;
  • 仓储与路由必须从 projects.team_id 解析所有权,要求调用方 team/API-key scope 匹配,拒绝任何 team_id 与 project 属主不符的写入;
  • 同时携带 project_idteam_id 的行必须使用 FOREIGN KEY (project_id, team_id) REFERENCES projects(id, team_id) 做 FK 级所有权校验——DDL 中大量出现的复合外键正是为此;
  • 三类生成源由 source_type/source_id 表达,不挤占 event-only 列:
    • 事件生成 job:source_type = 'agent_event'source_id = agent_event_id,且 agent_event_id 外键非空;
    • 会话摘要 job:source_type = 'session_summary'source_id = server_session_idagent_event_id = NULL;
    • 重建索引 job:source_type = 'observation_reindex',source_id 为目标 observation ID 或确定性 reindex scope ID,agent_event_id = NULL;
  • 非事件类 job 的 source_id 在插入前必须经仓储校验所有权(summary job 要在同一 project/team 下加载 server_sessions 行;reindex job 要加载目标 observation 或文档化 scope);
  • observation_generation_job_events 记录每个生成 job 的持久化生命周期/outbox 事件(enqueue、processing、retry、完成、失败等状态迁移),它不是 agent_events 的替代品,不得把原始事件 payload 作为规范事件记录存放。

4.5 Outbox 状态机与幂等规则

  • observation_generation_jobs.status 限定为 queuedprocessingcompletedfailedcancelled;
  • 合法生命周期:queued -> processing -> completedqueued -> processing -> failedqueued -> cancelled,以及"当 attempts < max_attempts 时从 stale/可重试 failed 工作回到 queued"的重试迁移;
  • attempts 仅在 worker 把 job 迁入 processing 时自增;next_attempt_at 是重试/对账的门控;locked_at/locked_by 在 worker 持有处理期间设置,在终态或 stale-lock 恢复时清除;
  • completed_at/failed_at/cancelled_at 为终态时间戳,终态 job 恰好一个非空;
  • 事件幂等:agent_events.source_event_id 仅是可选的适配器元数据,不得单独作为幂等权威;agent_events.idempotency_key 必填且确定性——有 source_event_id 时由 team_idproject_idsource_adaptersource_event_id 派生;缺省时由 team_idproject_idsource_adapterserver_session_idevent_typeoccurred_atpayload 的规范 JSON 哈希派生;UNIQUE (idempotency_key) 抑制重复摄入;
  • Job 幂等:job 的 idempotency_keyteam_idproject_idsource_typesource_idjob_type 确定性派生,UNIQUE (idempotency_key) 抑制重复 outbox 行;UNIQUE (team_id, project_id, source_type, source_id, job_type) 保证同一归属范围内每类生成源唯一;
  • bullmq_job_id 存在时必须确定性且唯一,使对账可以安全地重新添加或替换终态 BullMQ job;
  • 生成结果幂等:observations.generation_key 对直传/manual observation 可为空,对 provider 生成结果必填,格式为 generation:v1:{generation_job_id}:{parsed_observation_index}:{canonical_content_fingerprint}(fingerprint 在 parser 归一化之后、持久化之前计算);UNIQUE (team_id, project_id, generation_key) 是主要的重试幂等护栏——重试同一 job 与同一条解析出的 observation 必须 upsert/重载已有行,而非新建;
  • observations.created_by_job_idobservation_sources.generation_job_id 是到 observation_generation_jobs(id) 的可空外键,provider 生成的 observation 必须回填创建它的持久化 job;
  • UNIQUE (observation_id, source_type, source_id) 保证同一 source 不会重复链接到同一 observation;
  • 触碰 observation sources、job 状态或生命周期事件的变更类 API,必须要求 project_idteam_id 并写入变更 SQL 谓词;
  • ObservationRepository.search(...) 必须使用生成的 observations.content_search tsvector、其 GIN 索引与 websearch_to_tsquery('english', query) 实现范围化全文检索(DDL 中的 content_search TSVECTOR GENERATED ALWAYS AS (to_tsvector('english', content)) STOREDidx_observations_content_search 即为此设计);
  • provider 重试前必须重载 Postgres job 行与权威 source 行(事件 job 重载 agent_events,摘要 job 重载 server_sessions,reindex job 重载目标 observation/scope);BullMQ payload 只是"咨询性"执行数据,不是权威。

4.6 仓储接口与验证要求

计划的仓储接口面为:ProjectRepositoryTeamRepositoryObservationRepositoryObservationSourcesRepositoryObservationGenerationJobRepositoryObservationGenerationJobEventsRepositoryAgentEventsRepository(均基于 Server beta 的 Postgres 连接)。遗留 memory_items 数据只能作为迁移/视图存在;现有 MemoryItemsRepository 仅作兼容参考,不是 Server beta 的仓储契约。测试 helper 要求在未配置测试 Postgres URL 时干净跳过 Postgres 集成测试。

Phase 1 验证清单(保留原计划要点):仓储接口的单元测试(fake 适配器);Postgres 集成测试覆盖——迁移/引导幂等性;ProjectRepository.create(...) 要求合法 team_id 且 project 作用域写入拒绝错配的 team_id/project_id 组合;ObservationRepository.create(...) 与按 project/session/team 的查询;ObservationRepository.search(...) 走 GIN + websearch_to_tsquery 路径且只返回请求作用域的行;addSource(...) 幂等性与错作用域拒绝;AgentEventsRepository 的创建、批量插入/重载、两种幂等 key 派生路径,以及"省略 source event ID 时同一事件摄入两次不得产生重复 agent_events 行或重复生成 job";job 仓储的创建/状态迁移/重载与三类 job 的重复抑制;transitionStatus(...) 在条件更新与 fallback 重载中均要求 project/team 作用域,错作用域不得改行;generation_key 重试幂等;job 生命周期事件 append/list 及作用域校验。最后还有一条源码级检查:rg -n "MemoryItemsRepository" src/server 在新实现中不得命中(兼容适配器除外)。

Phase 1 反模式:不得让 SQLite 成为规范存储;不得新增名为 memory_items 的表或 MemoryItemsRepository 的仓储;不得让 BullMQ/Valkey 成为 observation 或 outbox 历史的 source of truth;不得以"静默回退 worker SQLite"掩盖 Postgres 缺失。

5. Phase 2:定义 Server 运行时边界

Phase 2 只创建生命周期/运行时边界,不实现 BullMQ 处理、provider 生成、生成 worker 或 SSE 广播(队列实现在 Phase 3 启动,生成在后续生成阶段,事件广播在其所属阶段接入):

  • 新增 src/server/runtime/ 下的运行时服务、工厂与服务图类型文件(当前仓库中即 src/server/runtime/ServerService.tssrc/server/runtime/create-server-service.tssrc/server/runtime/types.ts;计划原文命名为 ServerBetaService.ts/create-server-beta-service.ts,从源码结构看,实现落地时命名收敛掉了 "Beta" 后缀,仓库中另保存有后续更名计划 plans/2026-05-25-cmem-sdk-and-server-rename.md);
  • 服务图包含:Postgres 连接池;Phase 1 存储引导/迁移状态;认证模式;以及四个 inert(no-op)边界接口——队列管理器、生成 worker 管理器、provider 注册表、SSE/事件广播器,均以禁用/no-op 适配器占位,后续 Phase 再替换为真实实现;
  • claude-mem server start|stop|restart|status 路由到 Server 运行时服务,worker 命令继续路由到 WorkerService;
  • 独立运行时状态文件:.server-beta.pid.server-beta.port.server-beta.runtime.json(不复用 worker 的 PID/port 文件);
  • Server beta 的 /v1/info/api/health 需带 runtime = "server-beta" 标记。

实现参考与验证:从 Server.ts 复制路由处理器组合风格;只从 worker-service.ts 复制生命周期原语(不拷贝整个 worker 类);从 ProcessManager 复制 PID 文件安全模式。验证清单:rg -n "WorkerService|services/worker-service|worker/http" 在 server 运行时源码中不得命中;server status 独立于 worker 报告 server 状态;worker 命令仍可用;worker 停止时 Server beta 可启动,停止 Server beta 不触碰 worker。反模式:不得"后台拉起 worker"来实现 Server beta,不得以 worker 健康作为 server 健康源。

6. Phase 3:BullMQ 优先的 Server 队列

6.1 Job 类型与队列封装

计划要求在 src/server/jobs/ 定义 job 类型(当前仓库已落地 src/server/jobs/types.tssrc/server/jobs/ServerJobQueue.tssrc/server/jobs/job-id.ts)。计划列出的 job 族为:ServerGenerationJob(基类)、GenerateObservationsForEventJobGenerateObservationsForEventBatchJobGenerateSessionSummaryJobReindexObservationJob;每个 job 类型必须携带 team_idproject_idsource_typesource_idgeneration_job_id,事件 job 额外携带 agent_event_id,摘要 job 携带 server_session_id,reindex job 携带目标 observation ID 或确定性 scope ID。

当前仓库的实现把该契约以 Zod 判别联合固化在队列边界:基础字段 schema 强制 team_id/project_id/source_type(枚举 agent_event/session_summary/observation_reindex)/source_id/generation_job_id/source_adapter,并允许 api_key_id/actor_id 为空但字段必须存在(保证审计记录形状稳定);request_id 为 Phase 12 引入的请求关联 ID。源码注释把 Phase 11 的安全语义写得非常直白:job 中的 team_id/project_idadvisory 的,worker 必须先重载 Postgres 中的规范 outbox 行再比对,"把它当认证权威就是绕过——比对是篡改检测器,不是认证闸门"。两条队列名为 server_beta_generate_eventserver_beta_generate_summary,当前实现聚焦 event/summary 两类 job(reindex 保留为 source_type 枚举值)。

6.2 确定性 Job ID

job-id.ts 要求生成"确定性、无冒号"的 job ID。当前实现用 SHA-256 对 {kind, team_id, project_id, source_type, source_id} 的规范 JSON 求摘要,格式为 ${kindPrefix}_${sha256hex}(前缀 evt/sum)。源码注释解释了原因:SHA-256 派生 ID 避免跨租户 Redis key 碰撞,保持 BullMQ jobId 去重在进程重启间有效;而 BullMQ 内部用 : 作为 key 分隔符,jobId 中嵌入冒号会造成 scan/状态混乱。

6.3 Outbox、启动对账与健康

  • src/server/jobs/outbox 语义:durable 行在 observation_generation_jobs,source 身份在 source_type/source_id,生命周期事件在 observation_generation_job_events;outbox 是"应生成什么"的持久化权威,BullMQ 只是执行传输;
  • 启动对账:把 queued 或 stale processing 的 outbox 行重新入队;已完成 job 的行不再入队;在确定性 job ID 复用前,先移除或替换终态 BullMQ job;
  • 队列健康并入 /v1/info/api/healthclaude-mem server status

当前仓库的 ServerJobQueue 封装了 BullMQ Queue + Worker + QueueEvents,并把 BullMQ 文档要求落实为硬约束:autorun: false(显式 start())、默认并发 1(按计划"每 provider/session lane 保守值 1,后续可调")、每个 Worker 挂 error 监听(避免 job 抛错引发未处理崩溃)、默认 job 选项为 attempts: 3 + 指数退避(基线 5s)、锁时长默认 5 分钟,completed 记录保留 7 天/1000 条、failed 保留 30 天/1000 条;此外实现了 per-process 的 stalled/errored 计数器(BullMQ 的 getJobCounts 不暴露 stalled 数),并对 worker 与 QueueEvents 双路触发的 stalled 信号做 30 秒窗口去重,防止重复回调;还提供 queueFactory/workerFactory 测试缝隙,允许注入 fake 而不触碰 Redis。

Phase 3 验证清单:单元测试覆盖 job ID 稳定性、重复入队抑制、终态 job 替换、outbox 重启对账、失败 job 在 Postgres 与 BullMQ 中均保留、"选择 BullMQ 时 Redis 不可用必须让 Server beta 启动失败";集成测试用 fake processor 验证"重启后对账恢复 job 且 outbox 恰好标记一次"。反模式:不得把 BullMQ 的 completed/failed 状态当规范历史;本阶段通过不依赖事件路由与 provider 生成;重试不得产生重复 processor 副作用(后续 observation 写入必须按确定性生成键幂等);不得使用 BullMQ Pro-only 的 groups;不得把 pending 工作只留在 Redis。

7. Phase 4:Server 自有的"事件到生成 Job"管道

POST /v1/eventsPOST /v1/events/batch 的语义从"写入归档"升级为四步:

  1. 校验认证与 project/team 作用域;
  2. 事务性插入事件;
  3. 在同一事务中创建 server outbox 生成 job;
  4. 提交后 enqueue 对应 BullMQ job。

请求侧控制为渐进式 opt-in:

  • 默认:异步 enqueue 生成;
  • ?generate=false:仅存储事件,不生成;
  • ?wait=true(若本阶段实现):只等待有界的队列接收或 job 状态,返回 queued/accepted/job 状态,不得宣称 observation 已生成

配套接口:GET /v1/jobs/:id 查询生成状态;POST /v1/memories 保留为 manual/direct observation 插入的兼容别名,且不得调用生成器

验证清单(原计划要点):事件请求返回 eventgenerationJob;?generate=false 不返回生成 job;事件插入与 outbox job 创建事务一致——"没有无 outbox 行的事件,也没有无事件链接的 outbox 行";成功请求在提交后入队 BullMQ job;混合 project 的批量请求在任何事件/outbox/BullMQ 副作用发生前整单预校验拒绝;?wait=true 若实现,只返回 queued/accepted/job 状态,不返回生成的 observation ID;project 作用域的 API key 不得为其他 project 入队生成。反模式:不得调用 worker 的 /api/sessions/observations;不得让 /v1/events 依赖 Claude Code 特有的 hook payload 形状;不得在 HTTP 请求内(绕过队列)生成 observation;本阶段验证不要求 provider 生成、生成 ID 或重复检查。

8. Phase 5:剥离 provider 生成,解除 worker 耦合

  • 新增 src/server/generation/ProviderObservationGenerator.ts(当前仓库已存在该文件)与 provider 适配器目录 src/server/generation/providers/,规划三个适配器:ClaudeObservationProviderGeminiObservationProviderOpenRouterObservationProvider;
  • 从 worker provider 抽取公共 prompt 构造与 provider 调用代码为可复用模块,worker provider 退化为可调共享适配器的兼容包装;
  • 新增 src/server/generation/processGeneratedResponse.ts(当前仓库已存在该文件),职责:用 parseAgentXml(...) 解析响应 → 映射到 server observation 创建 schema/仓储输入 → 经 ObservationRepository 落库 → 把 source 链接到事件/job ID → 更新 outbox 状态 → 审计 observation 生成;
  • 新增 GET /v1/events/:id/observations 查看某事件的生成结果;
  • 稳定 server 生成 prompt:输入为 AgentEvent 记录列表 + project/session 元数据,输出为现有 parser 可接受的 XML 或结构化 JSON,并保留 <private> 跳过行为。

实现参考:解析/存储行为来自 worker 的 ResponseProcessor;各 provider 的认证与请求构造分别来自 ClaudeProvider/GeminiProvider/OpenRouterProvider;兼容字段约束取自 legacy schema(src/core/schemas/ 下),但 Server beta 的创建契约以 observation schema 暴露;保留 provider-errors.ts 的错误分类语义。

验证清单:fake provider 单元测试(合法 XML 产生 observation;skip/private 响应使 job 无 observation 而完成;畸形响应按策略失败或可重试;生成 observation 保留 project/session/source 元数据);?wait=true 只有在 Phase 5 的 provider 生成与持久化接通后、且 job 在超时内完成时才返回生成的 observation ID;同一事件/job 在重启后重放不得产生重复 observation;provider 分类测试与 worker 响应处理器测试保持通过;源码检查 rg -n "services/worker/(ClaudeProvider|GeminiProvider|OpenRouterProvider|agents/ResponseProcessor)" src/server 必须无命中。反模式:server 生成代码不得导入 WorkerRef/ActiveSession 等 legacy worker 会话类型;不得改写 legacy SessionStore 表;不得假设 Claude Code 转写。

9. Phase 6–9:会话语义、Hook 路由、MCP 与兼容层

Phase 6:独立于 worker 会话的 Server 会话语义。 server_sessions 是规范会话模型;为生成需要补充 contentSessionId(或通用外部会话 ID)、agentIdagentTypeplatformSourcegenerationStatuslastGeneratedAtEpoch(对应 DDL 中 server_sessions 列);ServerSessionRuntimeRepository 提供"取活跃会话、列未处理事件、标记生成开始/完成/失败"三个 helper;会话级生成策略支持:逐事件生成、短去抖窗口内合并小批量事件突发、/v1/sessions/:id/end 时生成 summary,且策略可由 server 配置调整。摘要以 kind/type"summary" 的 observation 记录存储。验证:server 会话的起止除显式迁移/导入代码外不触碰 legacy worker 会话行;结束会话入队 summary job;重复结束会话幂等;session 作用域的 API key 仍是 project 作用域。反模式:生成不得要求 legacy worker 会话 ID,不得把 worker ActiveSession 当作 server 运行时状态对象。

Phase 7:Hook 路由直达 Server beta。 安装器选择 Server beta 时,hook 直接调用 Server 端点:SessionStart → /v1/sessions/start(或兼容端点);PostToolUse → /v1/events;Stop/Summarize → /v1/sessions/:id/end。worker 仅保留为 fallback:Server beta 被选中但不健康时,hook 可回退 worker 并记录可观测的警告;现有 hook JSON 输出保持不变。同时要求本地 hook 的 API-key 引导:安装时创建 scope 到本地 project/user 的 hook key,以正确文件权限存于本地配置,并提供 key 轮换命令。hook 命令与输出契约参考 plugin/hooks/hooks.json,安装器 prompt/设置模式参考 src/npx-cli/commands/install.ts。验证:worker 模式与 server-beta 模式的生命周期 hook 测试均通过;server 宕机时回退 worker 且只记一条警告;server 健康时不启动 worker;仅靠 Server beta 即可在 PostToolUse hook 后看到生成的 observation。反模式:不得把 Server beta hook 经 worker 路由转发;server 健康时不得静默拉起 worker;hook key 不得写进生成的 bundle。

Phase 8:MCP 直接使用 Server 运行时。 新增由 Server beta API/core 逻辑支撑的 MCP 工具:observation_addobservation_record_eventobservation_searchobservation_contextobservation_generation_status;既有 memory_* 名称只能作为 observation 工具之上的兼容别名;Server beta 模式不得要求 worker 运行;MCP 写工具必须走与 REST 相同的服务方法。工具 schema 风格参考 src/servers/mcp-server.ts,REST schema 参考 src/core/schemas/。验证:MCP 客户端在 worker 未运行时即可记录事件并取回生成的上下文,可搜索生成的 observation;既有 MCP 搜索测试保持通过。反模式:不得在 MCP 工具中复制生成逻辑,不得向 MCP server 模式导入 WorkerService

Phase 9:无耦合的兼容性层。 兼容路由仅作适配器:/api/sessions/observations → 把 legacy payload 转换为 AgentEvent 再入队 Server 生成 job;/api/sessions/summarize → 转换为会话结束/summary job;legacy data/search 路由 → 读 Server 仓储或显式迁移视图。适配器位于 src/server/compat/(当前仓库 src/server/compat/),必须调用 Server 服务而非 worker 路由类,且必须维护一份 parity map,逐条记录每个 legacy 路由是"原生实现 / 适配器实现 / 有意不支持"。payload 归一化参考 worker http 的 shared 模块,映射风格参考 Claude Code adapter。验证:rg -n "services/worker/http/routes|WorkerService" src/server/compat src/server/runtime 无导入命中;Server beta 上的 legacy PostToolUse 路由创建事件与生成 job;viewer 兼容路由不要求 worker。反模式:不得整段拷贝 worker 路由类,兼容适配器不得反客为主成为规范 Server API。

10. Phase 10:Docker 与可部署运行时

  • Docker 镜像只启动 Server beta:无 worker 进程、无 worker PID、无 worker 健康依赖(镜像布局参考 docker/claude-mem/Dockerfile,E2E 风格参考 scripts/e2e-server-docker.sh);
  • Compose 栈(docker-compose.yml)包含三个容器:Server beta、承载规范 observation/job/session 存储的 Postgres、承载 BullMQ 的 Valkey;
  • 环境变量校验:CLAUDE_MEM_RUNTIME=server-betaCLAUDE_MEM_QUEUE_ENGINE=bullmq、Postgres URL 必填、Redis/Valkey URL 必填、API-key 认证默认强制;
  • 可选的独立生成 worker 进程模式:claude-mem server worker start——同一代码库、独立进程、共享 BullMQ 队列。

验证:E2E 不启动任何 worker;docker compose ps 显示 server + Postgres + Valkey 三个容器;/v1/events?wait=true 能产出 generated observations;job 执行中途重启 server 验证重试/幂等;吊销 API key 后验证写入与搜索被拒绝。反模式:容器内不得安装或拉起 worker;Docker 中不得使用 local-dev 认证;不得使用进程内队列。

11. Phase 11–12:团队感知生成与可观测性

Phase 11(团队感知)。 每个生成 job 必须携带 team_idproject_id、actor/API-key ID、source adapter;作用域校验发生在事件插入前与 job 执行前;生成的 observation 带 team/project 元数据;审计覆盖五个节点——事件接收、job 入队、provider 生成开始、observation 生成、observation 对外提供;新增团队级队列状态端点 /v1/teams/:id/jobs/v1/projects/:id/jobs。API-key/team 存储模式参考 src/storage/sqlite/ 下的 legacy 实现,project scope 护栏参考 v1 路由。验证:team 作用域的 key 不得越界读写/生成;project 作用域的 key 不得为其他 project 入队;生成 observation 带正确 team/project ID;审计记录含生成 job ID。反模式:BullMQ job data 不得成为认证绕过;job payload 中的 project/team ID 必须通过重载 Postgres outbox 行来核实(与 src/server/jobs/types.ts 中 "advisory + tampering detector" 注释一致)。

Phase 12(可观测性与运维)。 CLI 新增 claude-mem server jobs statusjobs retry <id>jobs cancel <id>jobs failed;队列指标覆盖 waiting/active/completed/failed/delayed/stalled 事件计数(对应 ServerJobQueue 暴露的 ServerJobCounts 与 stalled/errored 计数器);日志带 request ID/job ID 关联(即 job payload 中的 request_id);新增 /v1/jobs 列表端点。验证:失败的 provider 响应出现在 server jobs failed;retry 使 job 回到 queued 且 observation 恰好生成一次;cancel 阻止后续生成;stalled 事件带 job ID 记录在日志。反模式:队列状态默认不得暴露完整敏感事件 payload;不得非幂等重试。

12. Phase 13:最终验证门与退出标准

Phase 13 不是实现阶段,而是证明独立 Server beta 运行时"完整、持久、且与 legacy worker 兼容"的最终发布门。

要求的自动化测试。 单元:provider 生成 parser、事件到 job 的事务、job ID/幂等、生成上的 team/project 认证、兼容路由适配器。集成:Server beta 无 worker 启动;/v1/events 生成 observations;PostToolUse hook 经 Server beta 生成;MCP 事件写入经 Server beta 生成;生成进行中重启可安全重试。Docker:Server beta + Postgres + Valkey;API-key 认证;事件生成;重启持久化;被吊销 key 的拒绝;无 worker 进程。

要求的源码检查(原计划给出的 grep 集合):

rg -n "new WorkerService|services/worker-service|services/worker/http/routes" src/server
rg -n "PendingMessageStore|SessionQueueProcessor" src/server
rg -n "CLAUDE_MEM_AUTH_MODE=local-dev|ALLOW_LOCAL_DEV_BYPASS" docker docs/server.md
rg -n "POST /v1/events|generationJob|wait=true" docs README.md

期望:前两条 grep 在 Server beta 运行时中零命中;Docker 文档不推荐 local-dev 认证;文档提及事件生成语义。

手动验证九步: 启动 worker 确认既有流程正常 → 停 worker → 带 Valkey 启动 Server beta → 提交通用 REST 事件 → 确认无 worker 时 observation 出现 → 经 Server beta hook 路由提交 Claude Code PostToolUse payload → 再次确认无 worker 时 observation 出现 → 在 provider 调用期间重启 Server beta → 确认 job 重试且只生成一次。

退出标准: 当且仅当以下全部成立,Server beta 才算"独立":worker 停止时 Server beta 能生成 observations;Docker 镜像不拉起 worker;/v1/events 能入队并生成 observation;hook 路由在 server 健康时生成 observation;BullMQ 队列状态跨重启存活且重试安全;Postgres 是 observation 与生成 job 历史的 source of truth;worker 仍作为独立稳定运行时可用。

13. 当前仓库中的落地对照

从源码结构看,该计划的大部分内容已在当前仓库中实现,且目录布局与计划高度一致,可作为继续深挖的入口:

  • Postgres 存储层:src/storage/postgres/ 下已有 config.tspool.tsschema.tsindex.ts 四件套,以及 observations.tsagent-events.tsserver-sessions.tsgeneration-jobs.tsprojects.tsteams.ts 等仓储实现;其中 generation-jobs.tsPostgresObservationGenerationJobRepository 中,status 五值、source_type 三值、生命周期事件七值的类型定义与计划 DDL 中的 CHECK 约束一一对应;
  • 队列与 Job:src/server/jobs/types.ts(Zod 判别联合 + advisory 注释)、job-id.ts(SHA-256 无冒号 ID)、ServerJobQueue.ts(BullMQ 封装、stalled 去重、测试注入缝隙)即 Phase 3 的落地;
  • 运行时服务:src/server/runtime/ServerService.tscreate-server-service.tstypes.ts,以及 ActiveServerQueueManager.tsActiveServerGenerationWorkerManager.ts 两个把 Phase 2 中"inert 边界"替换为真实实现的适配器,还有 SessionGenerationPolicy.ts 对应 Phase 6 的会话级生成策略;
  • 生成层:src/server/generation/ProviderObservationGenerator.tsprocessGeneratedResponse.tsproviders/ 目录对应 Phase 5;
  • 测试面:tests/server/ 下有 auth-api-key.test.tsserver-boot.test.tsserver-runtime-guard.test.tsserver-runtime-smoke.test.tsdata-deletion.test.tsv1-routes.test.tsjobs/runtime/generation/ 子目录测试,tests/storage/postgres/ 覆盖 Postgres 仓储,与 Phase 13 的验证门要求相呼应。

值得注意的演化点:计划文档中的 ServerBetaService.tsserver-beta-e2e.mjs/e2e-server-beta-docker.sh 等名称,在当前仓库分别收敛为 ServerService.tsdocker/e2e/server-e2e.mjsscripts/e2e-server-docker.sh;命名差异源于后续的产品更名(见 plans/2026-05-25-cmem-sdk-and-server-rename.md),但架构语义——运行时边界、outbox 权威性、确定性 ID、作用域护栏——与计划保持一致。

14. 小结:这套架构可复用的设计要点

这份计划对任何"把 AI 后处理从单体 worker 拆成独立服务"的场景都有参考价值,核心要点可归纳为:

  1. 权威与传输分离:Postgres outbox 回答"应该生成什么",BullMQ 只回答"现在执行谁";BullMQ 的 completed/failed 状态不是规范历史;
  2. 幂等贯穿三层:事件层(agent_events.idempotency_key)、job 层(确定性 idempotency_key + UNIQUE(team_id, project_id, source_type, source_id, job_type) + 确定性 bullmq_job_id)、结果层(observations.generation_keygeneration:v1:... 指纹唯一约束),使重启、stalled 重放与手动 retry 都安全;
  3. Job payload 是 advisory 的:worker 执行前必须重载 Postgres 权威行并比对 team/project 字段,把 job 数据当作篡改检测对象而非信任来源,堵死"构造恶意 job payload 绕过认证"的路径;
  4. 运行时独立性以可验证的负断言定义:rg 检查 WorkerService 导入、独立 PID 文件、runtime = "server-beta" 健康标记、"worker 停止仍能生成"的集成测试——"独立"不是口号,而是一组可自动化的检查;
  5. 兼容只走适配器:legacy 路由/payload 在 src/server/compat/ 内转换成规范领域对象,配合 parity map 明确"原生 / 适配 / 不支持"三态,防止兼容层侵蚀规范 API。

配合仓库中的 CLAUDE.mdREADME.md 了解项目全貌后,沿上文各 Phase 对应的源码路径深入阅读,即可完整复现 Server beta 独立 BullMQ observation 运行时从计划到实现的每一条设计决策。

登录后查看全文
热门项目推荐
相关项目推荐