Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ flowchart TD
metrics -. metrics_reporter .-> manager

%% 对外分支与 optimizer
cache_event_subscriber["cache_event_subscriber"] --> py_connector
py_connector["py_connector"] --> client["client SDK"]
client -.-> config
optimizer["optimizer(仿真与优化)"] -. cache_location .-> meta
Expand Down
1 change: 1 addition & 0 deletions docs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
- [基本概念](design/basic_concepts.md) - Storage、Instance Group、Instance、Block、CacheLocation 等核心概念
- [高可用与选主机制](design/ha_leader_elector.md) - HA 架构、LeaderElector 状态机、CoordinationBackend、Leader 发现
- [CacheReclaimer 异步删除与过度逐出优化](design/cache_reclaimer_async_delete.md) - 异步删除生命周期、in-flight credit、反压与无进展退避
- [Cache Event Subscriber 设计](design/cache_event_subscriber.md) - RTP/vLLM 全量与增量同步、ACK/cursor 可靠性语义

### 开发文档
- [开发指南](develop/README.md) - 开发者入门指南和开发环境配置
Expand Down
65 changes: 65 additions & 0 deletions docs/design/cache_event_subscriber.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
# Cache Event Subscriber 设计

## 目标

在不依赖 RTP-LLM #1198、tair-kvcache #241 或 #236 代码的前提下,从 public `main` 实现相同业务能力:将 RTP-LLM/vLLM 的本地 KVCache 状态可靠同步到 KVCM,并同时支持权威全量恢复和低开销增量更新。

## 数据流

```text
RTP GetCacheStatus v2 / vLLM ZMQ
→ Source.prepare(不提交 cursor)
→ 规范化 BlockRecord
→ KVCM ReportEvent
→ ACK 成功
→ Source.commit(提交 generation/cursor/state)
```

- RTP:changefeed 返回 `generation`、队列 head、明确的分页 `next_cursor`;冷启动、generation 不匹配、history gap 和周期校准返回快照。
- vLLM:通过 sequence 检测丢包并请求 replay。replay 是有限历史,只有从 sequence 0 完整重放或收到 `AllBlocksCleared` 才能建立权威状态。
- vLLM publisher 在已提交正游标后重新出现 sequence 0 时视为新 epoch,立即清空旧候选状态并上报权威快照。
- vLLM 冷启动没有任何 batch 时保持“未就绪”,既不注册节点也不发送 HEARTBEAT;Subscriber 会主动查询 replay,直到 sequence 0 或 `AllBlocksCleared` 建立权威基线,不把空闲误判成引擎故障。
- vLLM block hash 按低 64 bit 统一映射到 KVCM signed-int64 key;这与 vLLM 的 legacy int event 表示一致。按 engine hash 查询 KVCM 的调用方必须复用同一 codec。
- KVCM:`EVENT_BLOCK_SNAPSHOT` 表示该 `instance_id + host_ip_port` 的完整集合。服务端同步清理旧位置后写入新集合;空快照表示清空。同一 instance + host 的快照、增量和异步清理通过分片锁串行化;generation 变化或节点已恢复时旧清理立即中止。

## 存储身份

引擎本地 cache 使用独立的 `ST_EVENT_REPORT` 类型和 `event-report://` URI,不复用 `ST_VINEYARD`。该后端只管理节点注册、心跳、generation 和位置元数据,不提供数据读写;实际 KV 数据仍由 RTP-LLM/vLLM 进程持有。这样 KVCM 返回位置时,调用方不会误选 Vineyard connector。

`ST_EVENT_REPORT` 不计入 KVCM 管理的存储容量,也不会进入主动回收请求。其位置只能由引擎事件、权威快照、`HOST_DOWN` 或心跳超时清理,避免 KVCM 调用一个并不拥有引擎显存的 backend 执行 `Delete`。

部署时需创建一个 `event_report` storage,并将其名字放入 instance group 的 `event_reporting_storage_candidates`。现有 Vineyard 配置和位置协议不受影响。

## 一致性语义

1. STORED 只在引擎真实 cache index 写入成功后产生,REMOVED 只在真实删除后产生。
2. KVCM 只接受已注册且可用节点的缓存变更;节点失活后必须重新注册再重放未确认更新。
3. Subscriber 同一时刻最多保留一个未确认更新;KVCM 不可用时不继续拉取。
4. 增量 ADD/DELETE 是幂等的;多批请求中途失败时,可从旧 cursor 重放。
5. 全量快照是恢复与校准能力,不是稳态传输方式;稳态仅发送变化块。
6. 所有索引都按 `instance_id` 隔离;不同 DP endpoint 必须全部拉取成功后才提交聚合状态。
7. KVCM 的 key 域只有 64 bit;vLLM 的完整 digest 会被截断,部署方需接受相应碰撞模型,并以 `instance_id` 做租户/模型隔离。
8. Source 短暂失败期间暂停 HEARTBEAT;连续失败达到阈值后先上报 `HOST_DOWN` 再退出。RTP Launcher 在可选模式下限频重启 Subscriber,在 required 模式下由 `ProcessManager` 传播失败。
9. 对无法从 ZMQ 空闲与断连中区分存活状态的引擎,可配置 HTTP health URL;连续探测失败会停止 HEARTBEAT 并触发 Subscriber 下线。
10. 冷启动先 `prepare` 权威全量,再执行“节点注册 → 快照 ACK → cursor commit → 启动 HEARTBEAT”;首快照失败时立即 `HOST_DOWN` 取消该次注册,并以同一个未提交快照重新注册重试。

小快照通过单个 `EVENT_BLOCK_SNAPSHOT` 替换;超过单请求上限时采用“空快照清理 + 分批 ADD”。后者在传输期间允许短暂的部分可见,但任一批失败都不会提交 source cursor,重试会再次从清理开始,最终收敛且不会产生超大 HTTP 请求。

RTP 配置必须提供与 `dp_size` 等量且互不重复的 endpoint。vLLM 当前明确限制为单 DP;这与其每个 DP rank 使用独立 ZMQ 端口的发布模型一致,禁止在没有多 publisher 聚合器时静默只消费 rank 0。

## 与原 PR 的关键区别

- 基线直接来自两个仓库 public `main`,没有 cherry-pick、commit 依赖或隐藏的 #236 前置条件。
- 在 main 已有 `EventReportingBackend` 抽象上增加最小 `ST_EVENT_REPORT` 后端,不引入 #236 的大范围 storage 重构。
- 全量和增量是一个协议的两种模式,而不是两套互相独立的上报链路。
- RTP 分页区分 queue head 与 next cursor,避免分页时跳过事件。
- cursor 只在 KVCM ACK 后提交;Manager 故障形成反压,而不是继续堆积无界事件。
- vLLM 历史不足时显式失败,不把残缺 replay 当作完整缓存状态。

## 验收

- 冷启动空/非空快照、增量新增/删除、分页、generation 变化、history gap 均可收敛。
- KVCM 请求失败后 cursor 不推进,恢复后重试不丢事件。
- 空权威快照能清除该 host 的旧位置;其他 host 和其他 instance 不受影响。
- RTP 多 endpoint 任一失败时不提交任何 endpoint 状态。
- 主动回收只选择 KVCM 管理的存储位置,不选择 `ST_EVENT_REPORT` 位置。
8 changes: 8 additions & 0 deletions docs/design/module_architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ KVCache Manager 采用中心化部署,负责 KVCache 的全局元数据管理
|---|---|---|
| **client** | `client/` | C++/Python 客户端 SDK,是推理引擎与 KVCM 之间的桥梁。对外提供 `ManagerClient`/`RTPLLMClient` 门面,内部由两条链路组成(见下):**元数据面** `MetaClient`(经 gRPC 桩 `internal/stub` 调用 KVCM 服务)与**数据面** `TransferClient`(经 `internal/sdk` 在推理引擎显存/内存与存储后端之间搬运 KVCache 数据)。面向外部,不被服务端核心调用。 |
| **py_connector** | `py_connector/` | 推理框架集成(Python)。将 client 接入 vLLM/SGLang/TRT-LLM,含 CUDA kernel 辅助,负责在引擎的推理流程中按正确顺序调用元数据面与数据面接口。此外自带一个纯 Python 的 HTTP 元数据面客户端 `KvCacheManagerClient`(`common/manager_client.py`),可通过统一服务发现 URL 获取 Manager 入口,并作为 C++ `MetaClient` 之外的另一条元数据面通路。位于 Python 侧栈顶。 |
| **cache_event_subscriber** | `cache_event_subscriber/` | 独立的缓存元数据同步进程。消费 RTP-LLM 的全量/增量 gRPC changefeed 或 vLLM 的 ZMQ 事件流,使用 prepare→KVCM ACK→commit 两阶段推进 cursor,并通过 `py_connector/common/KvCacheManagerClient` 调用 `RegisterInstance`/`ReportEvent`。它不参与推理热路径,位于 Python 侧栈顶。 |

> **三个面的界定**:本文档区分三个面——**元数据面**指 MetaService 的接口(`GetCacheLocation`/`StartWriteCache`/`FinishWriteCache`/`GetCacheMeta`/`RemoveCache`/`RegisterInstance` 等)及 client 侧调用这些接口的逻辑,是推理引擎读写 KVCache 的热路径;**数据面**指 KVCache 数据在引擎显存/内存与存储后端之间的实际搬运(`TransferClient`,不经过 KVCM);**管控面**仅指 AdminService 的接口(Storage 增删改、Instance Group 管理、账号、配置快照、运维监控、Leader 运维等),供运维/管理工具使用,不在推理引擎的读写热路径上。

Expand Down Expand Up @@ -97,6 +98,7 @@ service → manager → meta → config → data_storage → common
### 3.3 客户端与 Optimizer

- **client** 是独立的对外分支,仅共享 `common`、`config`、`data_storage`、`protocol` 以及 `service/util:manager_message_proto_util`;**py_connector** 通过 pybind 位于 client 之上。核心服务端不依赖 client。运行时,元数据面经 gRPC(C++ `MetaClient`)或 HTTP(py_connector 的 Python `KvCacheManagerClient`)调用 KVCM 服务,数据面经 C++ `TransferClient` 直接读写存储后端——这几条链路是理解端到端流程的关键(见第 4 节)。
- **cache_event_subscriber** 只依赖 `py_connector/common` 的 Python HTTP 客户端和引擎公开的事件协议;服务端核心不反向依赖 Subscriber。KVCM ACK 前不推进引擎 cursor,Manager 不可用时暂停拉取,从而形成有界反压。
- **optimizer** 负责 KVCache 访问 trace 的仿真与优化(命中率/容量分析、逐出与容量参数调优)。目前通过 `meta:cache_location` 类型与 `event` 的 optimizer 事件与核心关联;在线化能力正在开发中,后续会与现有 optimizer 合并。

### 3.4 模块关系图
Expand Down Expand Up @@ -129,6 +131,7 @@ flowchart TD
subgraph clientside["客户端侧(对外分支)"]
client["client SDK"]
py_connector["py_connector<br/>vLLM/SGLang/TRT-LLM"]
cache_event_subscriber["cache_event_subscriber<br/>RTP/vLLM cache events"]
end

optimizer["optimizer(仿真与优化)"]
Expand All @@ -154,6 +157,7 @@ flowchart TD
metrics -. metrics_reporter .-> manager

%% 客户端分支
cache_event_subscriber --> py_connector
py_connector --> client
client --> config
client --> data_storage
Expand Down Expand Up @@ -251,6 +255,10 @@ sequenceDiagram

`LeaderElector`(config)基于 `CoordinationBackend`(memory/file/redis)的分布式锁选主。`Server` 在成为 Leader 时调用 `CacheManager::DoRecover` 恢复状态,降级时调用 `DoCleanup` 清理运行时状态(正在进行的写入按失败处理)。Python `KvCacheManagerClient` 使用服务发现 URL 时,会在每次 Leader 刷新前重新选择一个 Manager 发现端点,避免把 Leader 查询入口固定在单个节点上。

### 4.7 引擎缓存事件同步

`cache_event_subscriber` 在引擎外部维护已获 KVCM 确认的 cursor 与缓存集合,并通过专用 `ST_EVENT_REPORT` 元数据后端发布 `event-report://` 位置。冷启动、引擎 generation 变化、历史断档和周期校准使用 `EVENT_BLOCK_SNAPSHOT` 权威替换该 host 的全部上报位置;正常运行只发送 `EVENT_BLOCK_ADD`/`EVENT_BLOCK_DELETE`。只有 KVCM 完整 ACK 后才提交 source cursor;请求失败时重试同一更新,不继续消费引擎事件。vLLM replay buffer 不是快照,冷启动只有从 sequence 0 完整重放或收到 `AllBlocksCleared` 才允许建立权威基线;基线建立前不发送 HEARTBEAT。

---

## 5. 修改功能前的提示
Expand Down
13 changes: 13 additions & 0 deletions kv_cache_manager/cache_event_subscriber/BUILD
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
load("@rules_python//python:py_library.bzl", "py_library")

package(default_visibility = ["//visibility:public"])

py_library(
name = "cache_event_subscriber",
srcs = glob(["*.py"]),
deps = [
"//kv_cache_manager/py_connector/common:common",
"@pip_cpu//grpcio",
"@pip_cpu//protobuf",
],
)
67 changes: 67 additions & 0 deletions kv_cache_manager/cache_event_subscriber/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
# KVCM Cache Event Subscriber

独立进程消费 RTP-LLM/vLLM 的 KVCache 事件,并通过 KVCM HTTP `ReportEvent` 同步元数据。设计与一致性语义见 [设计文档](../../docs/design/cache_event_subscriber.md)。

## 运行依赖

- 公共依赖:Python 3.10+、仓库内 `KvCacheManagerClient`、`requests`、`protobuf`、`grpcio`。
- RTP 路径无需额外包。
- vLLM 路径额外安装 `requirements-vllm.txt` 中的 `msgspec` 与 `pyzmq`。

Subscriber 需从 tair-kvcache 源码树运行,或将 `kv_cache_manager` Python 包安装/挂载到 RTP 容器;RTP Launcher 只管理命令生命周期,不复制另一仓库的运行时。

KVCM 侧必须先配置专用元数据后端,并绑定到 instance group:

```json
{
"global_unique_name": "engine_cache_events",
"event_report": {
"heartbeat_timeout_ms": 30000,
"cleanup_grace_ms": 300000,
"liveness_check_interval_ms": 5000
}
}
```

对应 instance group 的 `event_reporting_storage_candidates` 需包含 `engine_cache_events`。Subscriber 固定以 `ST_EVENT_REPORT` 上报,并生成 `event-report://<host>/<medium>` 位置;不要配置为 Vineyard storage。

该 storage 只保存引擎缓存的可发现元数据,不属于 KVCM 容量池,也不能放入普通 `storage_candidates`;缓存淘汰由引擎负责,KVCM 仅依据事件、快照和节点失活清理位置。

RTP 示例:

```bash
python -m kv_cache_manager.cache_event_subscriber \
--engine rtp \
--manager-uri http://manager:8080 \
--instance-group default \
--instance-id model-a \
--host-ip-port worker-a:9000 \
--rtp-endpoints worker-a:9001 \
--block-size 64 \
--model-name model-a \
--cache-group-count 1
```

vLLM 运行时还需安装 `requirements-vllm.txt`,并配置 publisher 的 PUB 与 replay endpoint。冷启动 replay 已裁剪时 Subscriber 会拒绝上报,必须与引擎共同重启或等到明确的 `AllBlocksCleared`。

RTP 要求 `--rtp-endpoints` 恰好提供 `--dp-size` 个互不重复的 endpoint,只有全部拉取成功才提交聚合状态。当前 vLLM 路径显式限制 `--dp-size=1`,避免多 DP 时静默只消费 rank 0;多 publisher 聚合需后续作为独立能力扩展。

vLLM 示例:

```bash
python -m kv_cache_manager.cache_event_subscriber \
--engine vllm \
--manager-uri http://manager:8080 \
--instance-group default \
--instance-id model-a \
--host-ip-port worker-a:9000 \
--vllm-pub-endpoint tcp://worker-a:5557 \
--vllm-replay-endpoint tcp://worker-a:5558 \
--engine-health-url http://worker-a:8000/health \
--block-size 16 \
--model-name model-a
```

由 RTP-LLM 托管时,将完整 argv 放入 `KVCM_CACHE_EVENT_SUBSCRIBER_COMMAND`。命令中的 `{rtp_endpoint}` 会在 backend 健康后替换为本机 gRPC endpoint;`KVCM_CACHE_EVENT_SUBSCRIBER_OWNER_RANK` 默认是 `0`,`KVCM_CACHE_EVENT_SUBSCRIBER_REQUIRED=false` 时异常退出会限频重启,设为 `true` 时失败由 RTP `ProcessManager` 向整组进程传播。

Subscriber 只有在 Source 已准备好首个权威快照后才注册节点,并在快照得到 KVCM ACK 后启动 HEARTBEAT。首快照失败会先上报 `HOST_DOWN` 取消该次注册,再用未提交的同一快照重新注册重试。
6 changes: 6 additions & 0 deletions kv_cache_manager/cache_event_subscriber/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
"""Reliable cache-event ingestion for RTP-LLM and vLLM."""

from .key_codec import to_signed_i64
from .models import BlockRecord, EngineUpdate, LocationSpec

__all__ = ["BlockRecord", "EngineUpdate", "LocationSpec", "to_signed_i64"]
Loading