Phase 1: The Pulse

[!NOTE]

文档定位:本文档是 000-roadmap.md Phase 1 的详细工程实施方案,用于指导「The Pulse (脉搏引擎)」的完整落地验证工作。涵盖技术调研、架构设计、代码实现、测试验证等全流程。


#1. 执行摘要

#1.1 定位与目标 (Phase 1)

Phase 1: Foundation & The Pulse 是整个验证计划的基石阶段,核心目标是:

  1. 构建统一存储基座:部署 PostgreSQL 16+ 生态,建立 Unified Schema
  2. 验证 Session Engine:实现对标 Google ADK SessionService 的会话管理能力
  3. 验证核心机制:原子状态流转、乐观并发控制 (OCC)、实时事件流

#1.2 核心设计:ADK Session 机制复刻

基于 Google ADK 官方文档[1]的深度分析,The Pulse 确立了以 Session 为核心容器,StateEvent 为双轮驱动的架构模式。

#1.2.1 核心概念映射

我们采用 PostgreSQL 全栈生态来承载 ADK 的抽象模型,实现像素级对标:

ADK 核心概念定义PostgreSQL 落地策略
Session单次用户-Agent 交互的容器,包含 eventsstatethreads 表 (主容器)
State会话内的 Key-Value 数据,支持分层作用域JSONB + 前缀解析 (Scoped)
Event交互中的原子操作记录events 表 (Append-only)
SessionServiceSession 生命周期管理接口OpenSessionService 类实现

#1.2.2 状态作用域与生命周期 (State Scopes)

针对不同维度的状态管理需求,我们实现了 ADK 定义的分层作用域机制:

前缀作用域生命周期存储策略
(Default)Session Scope随会话存续threads.state (JSONB)
user:User Scope跨会话持久化user_states
app:App ScopeGlobal 持久化app_states
temp:Invocation Scope仅当前思维链路有效内存缓存 (Volatile)

#1.2.3 状态颗粒度 (State Granularity)

[!IMPORTANT] 对标 Roadmap Pillar I:状态颗粒度设计决定了系统的记忆密度与回溯能力。

层次表名核心职责生命周期架构价值
Threadthreads交互历史的主容器 (Human-Agent Interaction)长期持久化长期记忆的输入源
Runruns单次推理过程的思维链 (Thinking Steps / Tool Calls)执行期间存活推理过程的可观测性
Eventevents不可变的原子事件流 (Message, ToolCall, StateUpdate)Append-only确定的系统状态回溯
Messagemessages语义负载容器持久化向量检索的核心语料
Snapshotsnapshots状态检查点策略性清理快速灾难恢复 (Fast Recovery)

#1.3 执行导图 (Execution Map)

为确保 Phase 1 的精准落地,我们将实施任务与技术文档进行了二维映射,并制定了基于 SOP 的工期计划。

#1.3.1 任务-文档锚定

[!NOTE] 关联文档:001-task-checklist.md

任务模块任务 ID 范围核心章节索引
FoundationP1-1-1 ~ P1-1-94.1 Step 1: 环境部署
Schema DesignP1-2-1 ~ P1-2-143. 架构设计 / 4.2 Schema 部署
Pulse EngineP1-3-1 ~ P1-3-174.3 核心实现
Event BridgeP1-5-1 ~ P1-5-54.4 AG-UI 事件桥接
VerificationP1-4-1 ~ P1-4-44.5 测试 / 5. Phase 1 验证 SOP

#1.3.2 工期规划 (3 Days)

阶段任务模块任务 ID预估工期关键交付物 (Deliverables)
1.1环境部署P1-1-1 ~ P1-1-90.5 DayPostgreSQL 16+ (pgvector/pg_cron)
1.2Schema 设计P1-2-1 ~ P1-2-140.5 Dayagent_schema.sql (Unified Model)
1.3Pulse Engine 实现P1-3-1 ~ P1-3-171.0 DayStateManager / PgNotifyListener
1.4AG-UI 事件桥接P1-5-1 ~ P1-5-50.5 DayEventBridge / StateDebugService
1.5全链路验收P1-4-1 ~ P1-4-40.5 Day自动化测试报告 / 技术白皮书

#2. 核心参考模型:Google ADK 契约与规范

#2.1 模型定位

本节定义了 Pulse Engine 必须遵循的 Normative Reference Model (规范性参考模型)。我们的设计并非凭空创造,而是通过严格复刻 Google GenAI ADK 的 SessionService 契约,确保系统具备行业标准的可扩展性与互操作性。

#2.2 ADK 核心对象建模

基于 ADK 源码[2],我们建立了如下对象关系模型,直接指导后续 Schema 设计:

#2.3 核心数据结构契约

#2.3.1 Session (会话容器)

Session 是状态管理的主体容器,对应数据库中的 threads 表:

hljs python
@dataclass
class Session:
    """
    Session Scope: 长期记忆容器
    Mapped to: table `threads`
    """
    id: str                    # Primary Key
    app_name: str              # Partition Key (Tenant)
    user_id: str               # Partition Key (User)

    # State Container (JSONB)
    # 关键:通过 version 字段实现 OCC (Optimistic Concurrency Control)
    state: dict[str, Any]

    events: list[Event]        # Event Sourcing History

#2.3.2 Event (原子事件)

Event 是不可变的交互记录,对应数据库中的 events 表:

hljs python
@dataclass
class Event:
    """
    Append-Only Ledger: 交互历史账本
    Mapped to: table `events`
    """
    id: str
    invocation_id: str         # Trace ID for Observability
    author: str                # 'user' | 'model' | 'tool'

    content: Content           # Payload (Text/Image/...)
    actions: EventActions      # Side Effects

#2.4 服务接口契约 (Interface Contract)

OpenSessionService 必须完整实现以下抽象基类定义的操作原语:

hljs python
class BaseSessionService(ABC):
    """
    Core Abstraction: 状态管理服务标准接口
    """

    @abstractmethod
    async def create_session(
        self,
        app_name: str,
        user_id: str,
        state: dict | None = None
    ) -> Session:
        """初始化会话上下文"""
        ...

    @abstractmethod
    async def get_session(
        self,
        app_name: str,
        user_id: str,
        session_id: str
    ) -> Session | None:
        """获取强一致性会话快照"""
        ...

    @abstractmethod
    async def list_sessions(
        self,
        app_name: str,
        user_id: str
    ) -> list[Session]:
        """列出用户所有会话"""
        ...

    @abstractmethod
    async def delete_session(
        self,
        app_name: str,
        user_id: str,
        session_id: str
    ) -> None:
        """删除会话"""
        ...

    @abstractmethod
    async def append_event(
        self,
        session: Session,
        event: Event
    ) -> Event:
        """
        核心原子操作:
        1. 持久化 Event
        2. 应用 State Delta
        3. 验证 OCC Version
        """
        ...

#2.5 前端集成规范:AG-UI 事件桥接

[!NOTE]

Protocol Alignment (协议对齐):本节定义 The Pulse 与 AG-UI 可视化层之间的 Event Bridge Protocol (事件桥接协议),确保状态变更与交互事件能够以毫秒级延迟实时投影到前端。

参考资源

#2.5.1 事件流架构概览

[!TIP]

Metaphor: The Pulse of the System (系统脉搏)

可以将整个架构想象成一个生命体监测系统:

  • 心脏 (Heart): Pulse Engine (PostgreSQL),每一次数据变更 (INSERT/UPDATE) 就像一次心脏跳动。
  • 脉搏波 (Pulse Wave): NOTIFY 机制,将心脏的跳动信号实时传导出去。
  • 监护仪 (Monitor): EventBridge,捕捉微弱的脉搏信号,将其转化为可视化的波形数据 (AG-UI Events)。
  • 屏幕 (Display): AG-UI 前端,实时显示生命体征,让用户看见系统的"存活"状态。

#2.5.2 事件映射契约 (Event Mapping Contract)

Pulse 产生的内部事件必须通过 EventBridge 转换为标准的 AG-UI 协议格式:

Pulse SourceTrigger ConditionAG-UI Event TypePayload Schema (Lite)
runsINSERT (Link Start)RUN_STARTED{ run_id, thread_id }
runsUPDATE (Finalized)RUN_FINISHED{ run_id, status, error? }
eventsINSERT (Role=user/agent)TEXT_MESSAGE_START{ message_id, role }
eventsINSERT (Chunk Delta)TEXT_MESSAGE_CONTENT{ delta_content }
threadsUPDATE (State Change)STATE_DELTA{ json_patch_diff }
eventsINSERT (Tool Call)TOOL_CALL_START{ tool_name, args_json }

#2.6 状态一致性模型 (Consistency Model)

#2.6.1 事务边界与可见性 (Transaction & Visibility)

[!IMPORTANT]

Read-Your-Writes Constraint (写后读约束)

根据 ADK 规范[3],状态变更 (state_delta) 仅在 Event 持久化事务提交后才对全局可见。这引入了Visibility Latency (可见性延迟)

  • Rule 1 (Persist-then-Visible): 任何 Agent 逻辑产生的状态变更,必须在 yield Event 被 Runner 捕获并 commit 到 DB 后,才能被新的 Session get() 操作读取。
  • Engineering Pitfall: "Airborne State" —— 开发者常错误地认为 yield UpdateState(...) 后,内存中的 state 对象会立即更新。实际上,在事务落地前,该指令处于“飞行中”状态,本地读取仍只能获取旧值。

⚠️ 常见代码误区 (The "Airborne" Trap)

hljs python
# ❌ 错误的直觉:认为 yield 后状态立刻改变
def my_agent_logic():
    # 1. 发出指令:更新计数
    yield UpdateState(key="count", value=100)

    # 2. 立刻读取
    # 此时指令还在“空中飞” (Airborne),Runner 尚未落地执行
    # 这里的 state.count 仍然是旧值(例如 0)
    if state.count == 100:
       logger.info("Success") # 永远不会执行!

#2.6.2 易失性状态与叠加视图 (Overlay View)

为解决上述延迟问题,并在长链路调用(Invocation)中支持连续的状态依赖,Pulse Engine 必须在内存中维护一个叠加视图。

[!TIP]

Analogy: The Scratchpad (草稿纸机制)

  • Scenario: 考试(Invocation)过程中,你在草稿纸(Memory Overlay)上写下中间步骤。
  • Requirement: 下一题计算必须能直接引用草稿纸上的结果,而不需要等待考试结束(Commit)后再去查阅试卷。
  • Risk: 如果考试中途被终止(Crash),草稿纸内容丢弃,不污染正式试卷(Database)。

核心实现要求StateManager 必须实现 Overlay Read 机制:

Stateeffective=Statepersistent+Delta_pending State*{effective} = State*{persistent} + \sum Delta\_{pending}

#3. 架构设计:Unified Schema

#3.1 ER 图设计

[!NOTE]

Design Principles: Protocol-First, Unified Storage, Event-Driven. 采用 "7 Tables + 2 Triggers" 的架构,实现 ADK Session 协议的完整持久化。

Schema 规格说明 (Schema Specification)

TableResponsibilitiesCore Spec & FeaturesKey Constraints / Indexes
threadsSession ContainerOCC Check: version field
Data: state (JSONB)
Unique: (app, user, id)
Idx: (app, user)
eventsImmutable StreamTrigger: notify_event_insert
Type: Append-only Log
Idx: (thread_id, sequence_num)
FK: ON DELETE CASCADE
runsExecution LoopObservability: thinking_steps
Status: Async State Machine
Idx: (thread_id), (status)
messagesSemantic ContentAI Ready: vector(1536) field
Role: user / assistant
Idx: (thread_id), (role)
(Phase 2: HNSW Index)
snapshotsState CheckpointsRecovery: Fast-forward restore
Freq: Per N events / Runs
Unique: (thread_id, version)
user_statesUser PersistenceScope: Cross-session memory
Query: GIN Indexing
PK: (user_id, app_name)
Idx: GIN(state)
app_statesGlobal ConfigScope: App-level config
Query: GIN Indexing
PK: (app_name)
Idx: GIN(state)

#3.2 系统交互设计 (System Interaction)

[!TIP]

Data Flow: Transaction (Write + Update) -> Notify -> Bridge -> Push 本节描述从事件产生到端到端推送的完整时序。注意:所有的状态变更通知都是由 events 表的插入触发的。

#3.3 状态管理与 OCC 机制

[!IMPORTANT]

乐观并发控制 (OCC):为了防止多 Agent 同时修改状态导致的数据覆盖,我们引入 version 字段进行 CAS (Compare-And-Swap) 控制。 关键点state 的更新必须与记录该变更的 event 在同一个事务中提交。

核心逻辑 (Atomic Transaction)

hljs sql
BEGIN;

-- 1. 尝试更新状态 (CAS)
UPDATE threads
SET
  state = state || $new_state,
  version = version + 1
WHERE
  id = $thread_id AND version = $expected_version
RETURNING version; -- 如果返回空,说明版本不匹配(冲突)

-- 2. 插入事件 record (触发 Notify)
INSERT INTO events (thread_id, event_type, content)
VALUES ($thread_id, 'state_update', $new_state);

COMMIT;

状态流转图

#3.4 Schema 部署

参见:src/cognizes/engine/schema/agent_schema.sql


#4. 实施指南

#4.1 Step 1: 环境部署与基础设施

[!TIP]

Metaphor: 夯实地基 (Building the Foundation)

在开始编码前,我们需要先整理好土壤 (DB)、搭建好脚手架 (Python Env) 并拿到钥匙 (Secrets)。

#4.1.1 The Soil: PostgreSQL 生态部署

核心目标:为 Pulse Engine 准备肥沃的土壤。

任务清单

任务 ID任务描述验收标准参考命令
P1-1-1部署 PG 16+SELECT version() 16.x+brew install postgresql@16
P1-1-2初始化数据库cognizes-engine 存在createdb cognizes-engine
P1-1-3安装扩展vector, uuid-ossp 就绪见下文安装指南
P1-1-4安装 pg_cron定时任务调度器就绪见下文安装指南

关键安装指南

hljs bash
# 1. 基础安装 (macOS)
brew install postgresql@16 pgvector
brew services start postgresql@16

# 2. 创建数据库 (The Container)
createdb cognizes-engine

# 3. 启用基础扩展 (The Nutrients)
psql -d cognizes-engine -c "CREATE EXTENSION IF NOT EXISTS \"uuid-ossp\";"
psql -d cognizes-engine -c "CREATE EXTENSION IF NOT EXISTS vector;"

pg_cron (The Clock) 深度部署指南

[!IMPORTANT]

Prerequisite: The Heartbeat Mechanism pg_cron 作为系统的后台调度器,必需在 postgresql.conf 中预加载 (shared_preload_libraries) 才能启动其后台进程。

关键路径:源码编译与配置 (macOS)

  1. 修正编译参数 (Apple Silicon Only)

    • 现象: 链接报错 Undefined symbols: _libintl_ngettext
    • 对策: 编辑 Makefile (约第 22 行),显式链接 gettext 库:
      hljs make
      # Old: SHLIB_LINK = $(libpq)
      SHLIB_LINK = $(libpq) -L/opt/homebrew/opt/gettext/lib -lintl
      
  2. 执行安装流水线

    hljs bash
    # 1. Build & Install
    git clone https://github.com/citusdata/pg_cron.git && cd pg_cron
    # (Execute Makefile fix here if on M1/M2/M3)
    export PATH="/opt/homebrew/opt/postgresql@16/bin:$PATH"
    make clean && make && make install
    
    # 2. Config (Find path: psql -c "SHOW config_file;")
    # Add to postgresql.conf:
    # shared_preload_libraries = 'pg_cron'
    # cron.database_name = 'cognizes-engine'
    
    # 3. Restart & Enable
    brew services restart postgresql@16
    psql -d postgres -c "CREATE EXTENSION IF NOT EXISTS pg_cron;"
    

#4.1.2 The Scaffold: 开发环境配置

核心目标:搭建稳固的 Python 开发脚手架。

hljs bash
# 1. 初始化项目 (The Frame)
mkdir -p src/cognizes/engine/{pulse,schema} tests/pulse
uv init --no-workspace .

# 2. 安装依赖 (The Tools)
# Core: asyncpg (Driver), pydantic (Validation)
# SDK: google-adk (Protocol)
uv add asyncpg 'psycopg[binary]' google-adk pydantic
uv add --dev pytest pytest-asyncio

#4.1.3 The Keys: 配置与密钥管理

核心目标:注入启动引擎所需的燃料。

配置清单

hljs bash
# .env 文件模板
# 1. Database Connection (Soil Access)
DATABASE_URL="postgresql://user:pass@localhost:5432/cognizes-engine"

# 2. Google ADK Auth (Identity)
GOOGLE_API_KEY="your-gemini-api-key"

验证命令

hljs bash
uv run python -c "
import os
print(f'✓ DB:  {os.getenv(\"DATABASE_URL\", \"Not Set\")}')
print(f'✓ Key: {os.getenv(\"GOOGLE_API_KEY\", \"Not Set\")[:5]}...')
"

#4.2 Step 2: The Blueprint - Schema 部署

[!TIP]

Metaphor: 绘制蓝图 (Drawing the Blueprint) 在土壤准备好后,我们需要在其中绘制数据的蓝图 (Schema),定义 7 张核心表与 2 个触发器。

核心目标:将 Unified Schema 物理化到数据库中。

部署流水线

hljs bash
# 1. Deploy (The Blueprint Execution)
psql -d cognizes-engine -f src/cognizes/engine/schema/agent_schema.sql

# 2. Verify Tables (7 Core Tables)
psql -d cognizes-engine -c "\dt" | grep -E "threads|events|runs|messages|snapshots|user_states|app_states"

# 3. Verify Triggers (The Pulse Mechanism)
psql -d cognizes-engine -c "\df notify_event_insert"

验收标准

  • ✓ 7 张表全部创建成功
  • notify_event_insert 触发器就绪
  • update_updated_at 触发器就绪

#4.3 Step 3: The Heart - Pulse Engine 核心实现

[!TIP]

Metaphor: 安装心脏 (Installing the Heart) Schema 是骨架,Pulse Engine 是跳动的心脏。它负责状态管理 (StateManager) 和脉搏传导 (PgNotifyListener)。

#4.3.1 The State Keeper: StateManager 实现

核心职责

  • 原子状态流转 (Atomic State Transitions)
  • 乐观并发控制 (OCC with Version Check)
  • 前缀作用域解析 (user:, app:, temp:)

实现参考src/cognizes/engine/pulse/state_manager.py

#4.3.2 The Pulse Conductor: PgNotifyListener 实现

核心职责

  • 监听 PostgreSQL NOTIFY 信号
  • 解析事件 Payload 并分发
  • 驱动 EventBridge 进行实时推送

实现参考src/cognizes/engine/pulse/pg_notify_listener.py

#4.3.3 The Gateway: WebSocket 推送接口

核心目标:为前端提供实时事件流的订阅通道。

关键组件

  • FastAPI Endpoint: src/cognizes/engine/api/main.py
  • Event Router: 基于 thread_id 的订阅路由
  • Protocol: WebSocket (实时双向) / SSE (单向流)

#4.4 Step 4: The Meridian System - AG-UI 事件桥接

[!TIP]

Metaphor: 构建经络系统 (Building the Meridian System)

如果把 Pulse Engine 比作心脏,EventBridge 则是遍布全身的经络——它把每一次心跳(DB 事件)转化为脉络中的血液流动(AG‑UI 协议),从而驱动可视化界面的即时响应。

核心目标:实现 PostgreSQL 事件流与 AG-UI 协议的无缝转换。

#4.4.1 The Translator: EventBridge 核心实现

核心职责

  • Protocol Translation: 将 PG NOTIFY Payload 转换为 AG-UI 标准事件格式
  • Event Routing: 根据 event_type 路由到不同的 UI 组件
  • Type Safety: 6 种核心事件类型的强类型定义

事件映射矩阵

PG Event SourceAG-UI Event TypeTrigger Condition
events (role=user/agent)TEXT_MESSAGE_START新消息创建
events (content delta)TEXT_MESSAGE_CONTENT流式内容追加
runs (INSERT)RUN_STARTED执行链路启动
runs (UPDATE complete)RUN_FINISHED执行完成
threads.state (via event)STATE_DELTA状态变更
events (tool_call)TOOL_CALL_START工具调用

实现参考src/cognizes/engine/pulse/event_bridge.py

#4.4.2 The Monitor: 状态调试面板

核心职责

  • State Inspection: 按前缀分组展示状态 (user:, app:, session)
  • History Tracking: 状态变更历史追溯
  • Debug API: RESTful 接口供开发工具调用

实现参考src/cognizes/engine/pulse/state_debug.py

#4.4.3 The Channel: SSE 事件流端点

核心目标:为前端提供持久化的单向事件流通道。

关键特性

  • Protocol: Server-Sent Events (HTTP Long-Polling)
  • Endpoint: /api/runs/{run_id}/events
  • Latency: Target < 100ms (端到端)

实现参考src/cognizes/engine/api/main.py

验收标准

  • ✓ 6 种事件类型正确映射
  • ✓ SSE 流延迟 < 100ms
  • ✓ 状态调试 API 完整可用
  • ✓ 单元测试覆盖率 > 80%

#4.5 Step 5: 测试

#4.5.1 单元测试套件

hljs bash
uv run pytest tests/unittests/pulse/test_state_manager.py -v

#4.5.2 端到端延迟测试

hljs bash
uv run pytest tests/integration/pulse/test_notify_latency.py -v -s

#统一执行

如需一次性运行全部 Pulse 相关测试:

hljs bash
uv run pytest tests/unittests/pulse tests/integration/pulse -v

#5. 验证 SOP (Phase 1:生命体征监测)

[!IMPORTANT]

本节提供 Phase 1: The Pulse 完整验收流程,请按顺序逐步执行。

Phase 1 验证是一次完整的外科体检。我们需要依次确认环境无菌 (Environment)、器官功能正常 (Organ Function)、血液循环畅通 (Systemic Circulation) 以及心脏抗压能力 (Cardiac Stress)。

#5.1 Step 1: 环境检查 (Clinical Environment Check)

Objective: 确保手术室(运行环境)的各项指标符合生命维持标准。

检查清单:

  1. 数据库 (The Soil): 确保统一存储基座 (cognizes-engine) 就绪。
  2. 扩展 (The Nutrients): 确保 vector (记忆) 与 pg_cron (心跳) 能力加载。
  3. 连接 (The Vessel): 确保 Python 到 PostgreSQL 的 asyncpg 通路顺畅。
hljs bash
# 1. 基础环境 (PG Version > 16)
psql -d 'cognizes-engine' -c "SELECT version();"

# 2. 核心机能 (Extensions Check)
psql -d 'cognizes-engine' -c "SELECT * FROM pg_available_extensions WHERE name IN ('vector', 'pg_cron');"

# 3. 血管通路 (Connection Check)
psql -d cognizes-engine -c "\dt"

# 4. 脉搏接口 (Python Import)
uv run python -c "from cognizes.engine.pulse.state_manager import StateManager; print('✓ Import OK')"

#5.1.1 激活心脏起搏器 (Service Startup)

启动 API 服务,激活整个系统的脉搏监听器 (PgNotifyListener)。

hljs bash
# 终端 1:启动 FastAPI 服务 (Start the Heartbeat Monitor)
uv run uvicorn cognizes.engine.api.main:app --reload --host 0.0.0.0 --port 8000

预期生命体征

hljs log
INFO:     Uvicorn running on http://0.0.0.0:8000 (Press CTRL+C to quit)
INFO:     ✓ PgNotifyListener started  <-- 监听器启动,脉搏开始传导

#5.1.2 基础反射测试 (Health Check)

hljs bash
# 终端 2:测试膝跳反射 (Health Endpoint)
curl http://localhost:8000/health

预期反应

hljs json
{ "status": "ok", "listener_running": true }

#5.2 Step 2: 器官功能测试 (Organ Function Test)

Objective: 验证核心器官 (StateManager) 的收缩与舒张逻辑是否准确。

Note: 虽然归类为 Unit Test,但为了确保原子性与 OCC 机制的真实可靠性,以下测试会连接真实数据库进行验证。

测试范围

  • State Transitions: 状态能否正确更新、合并。
  • Logic Validity: 前缀作用域解析 (user:, app:) 是否精准。
  • Atomic Transaction: 确保 Event 追加与 State 更新在同一事务中提交。
hljs bash
# 1. 全面器官检查 (44 个测试用例)
uv run pytest tests/unittests/pulse/ -v

# 2. 核心瓣膜专项检查 (StateManager Focus)
# 包含 Session CRUD, OCC 冲突检测, 原子性验证
uv run pytest tests/unittests/pulse/test_state_manager.py -v --tb=short

#5.3 Step 3: 全身循环测试 (Systemic Circulation Test)

Objective: 验证血液(事件)能否从心脏(DB)泵出,经由动脉(EventBridge),最终到达末梢神经(UI Client)。这对应于 Integration & E2E Testing

#5.3.1 经络通路测试 (WebSocket)

验证前端能否通过 WebSocket 建立长连接。

hljs bash
# 使用 websocat 探针 (需安装: brew install websocat)
websocat ws://localhost:8000/ws/events/test-thread

#5.3.2 脉搏传导测试 (Notify -> Client)

验证 "数据库插入 -> 触发 Notify -> 推送 WebSocket" 的完整反射弧。

hljs bash
# 终端 3: 注入刺激信号 (Trigger Test Event)
curl http://localhost:8000/api/test-notify

预期反应:WebSocket 客户端应在 < 100ms 内“抽动”一下(收到 JSON 数据)。

#5.3.3 微循环测试 (SSE Stream)

验证 SSE (Server-Sent Events) 单向高频流的稳定性,模拟 Run 执行过程中的 Token 流输出。

hljs bash
# 终端 1: 接入微循环监测仪 (Listen to SSE)
curl -N http://localhost:8000/api/runs/test-run/events

# 终端 2: 注入造影剂 (Trigger SSE Event)
curl http://localhost:8000/api/test-sse-notify/test-run

预期反应

  1. 连接时:立即收到 connected 事件。
  2. 注入后:立即流出 RAW 事件 payload。
  3. 静置后:每 30 秒收到有力的一次 heartbeat 跳动。

#5.4 Step 4: 心脏负荷测试 (Cardiac Stress Test)

Objective: 在高压从下,验证心脏能否维持 50ms 以内的泵血延迟,且不发生心室颤动(数据竞争/丢失)。这对应于 Performance Testing

关键指标 (Vital Signs Target)

  • Systolic Latency (收缩延迟): End-to-End < 50ms
  • Cardiac Output (心输出量): > 100 msg/s throughput
  • Rhythm Stability (节律稳定性): 并发写入无冲突丢失 (OCC Effective)
hljs bash
# 1. 延迟与吞吐量专项测试 (Latency Check)
# 覆盖: test_end_to_end_latency (<50ms), test_100_msg_per_second_throughput
uv run pytest tests/integration/pulse/test_notify_latency.py -v -s

# 2. 并发压力测试 (Stress Test)
# 模拟 10 个 Agent 同时并发写入同一 Session,验证 OCC 乐观锁机制 (Zero Data Loss)
# 覆盖: test_10_concurrent_writes_no_data_loss, test_100_qps_session_creation
uv run pytest tests/unittests/pulse/test_state_manager.py -v -s -k "Concurrency or Performance"

验收标准 (Discharge Criteria)

指标目标验证用例 (Verified in Tests)
Reflex< 50msintegration/test_notify_latency.py::test_end_to_end_latency
Throughput> 100 msg/sintegration/test_notify_latency.py::test_100_msg_per_second_throughput
Endurance100 QPSunittests/test_state_manager.py::test_100_qps_session_creation
Rhythm0 Lossunittests/test_state_manager.py::test_10_concurrent_writes_no_data_loss

#6. 验收基准 (出院评估, Discharge Summary)

#6.1 功能验收矩阵

对标任务: 001-task-checklist.md

CategoryIndicator (验收项)Target (目标值)Verification Method (验证方法)
AnatomySchema Integrity7 Tables + 2 Triggers\dt (Tables) + \df (Triggers)
PhysiologyAtomic State0 Dirty Reads / Lost Updatestest_state_manager.py (Unit/DB)
OCC Resilience100% Concurrent Write Successtest_10_concurrent_writes_no_data_loss
CirculationSystolic LatencyP99 < 50mstest_end_to_end_latency
Stream Continuity100% Delivery (No Loss)test_100_msg_per_second_throughput
NeurologyEvent Integration6 Event Types Correctly Mappedtest_event_bridge.py
SSE ReflexResponse < 100mstest_event_bridge_e2e.py

#6.2 Stress Test Limits (负荷极限)

Objective: 确立系统的安全运行边界 (Operational Boundaries)。

MetricThresholdConditionCorresponding Test
Session Gen Rate> 100 QPSSingle Node, Empty Statetest_100_qps_session_creation
Event Pulse Rate> 500 QPSSingle Node, Small Payloadtest_append_event_with_state_delta (Est.)
Notification Lag< 50 ms@ 100 msg/stest_end_to_end_latency
Write Concurrency10 AgentsSame Session, No Data Losstest_multi_agent_concurrency
Stream Latency< 100 msEvent -> SSE Clienttest_sse_reflex (E2E)

#6.3. 交付物清单 (Deliverables)

CategoryFile PathDescriptionTask ID
Theorydocs/010-the-pulse.md本实施方案与 SOPP1-4-1
Blueprintsrc/cognizes/engine/schema/agent_schema.sql统一建表脚本 (7 表 + 2 触发器)P1-2-12
Organssrc/cognizes/engine/pulse/state_manager.pyStateManager (Heart)P1-4-2
src/cognizes/engine/pulse/pg_notify_listener.pyNotify Listener (Pulse Node)P1-3-14
src/cognizes/engine/pulse/event_bridge.pyEvent Bridge (Translator)P1-5-1
src/cognizes/engine/pulse/state_debug.pyDebug Service (Monitor)P1-5-4
src/cognizes/engine/api/main.pyFastAPI Entry (Gateway)P1-3-15
DiagnosticsType: Unit / Logic
tests/unittests/pulse/test_state_manager.pySession CRUD, OCC, ConcurrencyP1-4-3
tests/unittests/pulse/test_pg_notify_listener.pyListener LogicP1-3-15
tests/unittests/pulse/test_event_bridge.pyEvent Mapping & SchemaP1-5-2
Type: Integration / Systemic
tests/integration/pulse/test_state_manager_db.pyFull DB TransactionsP1-4-4
tests/integration/pulse/test_notify_latency.pyEnd-to-End Latency & ThroughputP1-3-16
tests/integration/pulse/test_event_bridge_e2e.pySSE Stream VerificationP1-5-3
tests/integration/pulse/test_state_debug_db.pyState History QueryP1-5-6

#7. 限制与未来规划

[!WARNING]

Phase 1 工程边界:以下限制是当前架构设计的已知约束,将在后续 Phase 2 中优化。

组件/领域限制描述影响评估Phase 2 优化方向
PostgreSQLNOTIFY payload 最大 8000 bytes大消息可能被截断,需走回查机制引入 Redis Pub/Sub 或 Hybrid Hybrid Queue
pg_cron调度精度最小 1 分钟无法支持秒级定时任务引入专用 Job Scheduler (如 Temporal)
StateJSONB 整体读写状态过大时(>1MB)性能下降引入 jsonb_set 局部更新或拆表存储
Throughput单节点 DB 瓶颈预估上限 ~5k TPS引入 Read Replica 或 Sharding

#8. 参考文献

1. Google. (2025). ADK Sessions Documentation. https://google.github.io/adk-docs/sessions/

2. Google. (2025). ADK Session Overview. https://google.github.io/adk-docs/sessions/session/

3. Google. (2025). ADK State Documentation. https://google.github.io/adk-docs/sessions/state/

4. Google. (2025). ADK Context Documentation. https://google.github.io/adk-docs/context/