Skip to content

Commit c2b1280

Browse files
committed
[subscriber] add kvcache event subscriber
1 parent 09d2b01 commit c2b1280

128 files changed

Lines changed: 32933 additions & 0 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

subscriber/.gitignore

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,53 @@
1+
# === Qoder / local agent state ===
2+
.qoder/
3+
4+
# === 本地环境信息 ===
5+
env.txt
6+
7+
# === Python & uv ===
8+
__pycache__/
9+
*.py[cod]
10+
*$py.class
11+
*.so
12+
*.egg-info/
13+
*.egg
14+
dist/
15+
build/
16+
.eggs/
17+
18+
# uv 虚拟环境与锁文件缓存
19+
.venv/
20+
.python-version
21+
22+
# 测试与覆盖率
23+
.pytest_cache/
24+
.coverage
25+
htmlcov/
26+
.tox/
27+
.nox/
28+
29+
# mypy / ruff / linting 缓存
30+
.mypy_cache/
31+
.ruff_cache/
32+
.dmypy.json
33+
34+
# Jupyter Notebook
35+
.ipynb_checkpoints/
36+
37+
# === PyCharm / JetBrains IDE ===
38+
.idea/
39+
*.iml
40+
*.iws
41+
*.ipr
42+
out/
43+
44+
# === OS ===
45+
.DS_Store
46+
Thumbs.db
47+
48+
# === 环境变量与密钥(切勿提交)===
49+
.env
50+
.env.*
51+
!.env.example
52+
53+
*.log

subscriber/AGENTS.md

Lines changed: 302 additions & 0 deletions
Large diffs are not rendered by default.

subscriber/README.md

Lines changed: 194 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,194 @@
1+
# KVCacheEventSubscriber
2+
3+
A process co-located with an inference engine that subscribes to the engine's KV cache
4+
events, forwards them to the kvcm service, and ties its health to the co-located
5+
DashServing instance.
6+
7+
## Modules
8+
9+
| 模块 | 说明 |
10+
|---|---|
11+
| `subscriber/cli.py``main.py` | 进程入口与生命周期:启动 7 步、serving 监督、关停 6 步 |
12+
| `subscriber/config.py` | 配置定义、CLI 注册、校验、派生属性 |
13+
| `subscriber/forwarding.py``pipeline/` | 增量 / 快照双 pipeline 的转发链路 |
14+
| `subscriber/engine/` | 引擎适配层(vLLM / SGLang 实现 + 共用 ZMQ source + gRPC 控制面) |
15+
| `subscriber/health/` | 引擎存活协调(epoch 门控)+ DashServing 状态上报 |
16+
| `subscriber/kvcm/` | KVCM 客户端栈(注册、心跳、事件上报、gRPC/HTTP 传输、服务发现) |
17+
| `subscriber/metrics/` | 指标统一出口(主路径 telemetry + lifecycle 点事件) |
18+
| `subscriber/proto/` | engine / KVCM pb 绑定(手工维护) |
19+
20+
## Tech Stack
21+
22+
| 类别 | 技术 |
23+
|---|---|
24+
| 语言 | Python ≥ 3.11(asyncio 全异步) |
25+
| 引擎事件通道 | ZeroMQ(`pyzmq`:SUB 实时流 + DEALER replay) |
26+
| 引擎控制面 | gRPC(`grpcio` 1.75.1 + `protobuf` 3.20.3) |
27+
| 序列化 | `msgspec`(msgpack) |
28+
| KVCM | 默认 gRPC/protobuf(`grpcio`),HTTP/JSON(`httpx`)兼容回退 |
29+
| DashServing | HTTP JSON(`httpx`|
30+
| 观测 | dashlog(可选依赖,缺失时降级) |
31+
| 工具链 | `uv``ruff``mypy --strict``pytest` |
32+
33+
## 设计文档
34+
35+
稳态架构文档在 `docs/architecture/`
36+
37+
| 文档 | 内容 |
38+
|---|---|
39+
| [00-overview](docs/architecture/00-overview.md) | 模块划分、启动 / 退出时序、数据流总览、核心不变量 |
40+
| [01-forwarding-pipeline](docs/architecture/01-forwarding-pipeline.md) | 双 pipeline 数据流、凑批、门控、丢弃语义 |
41+
| [02-engine-adapter](docs/architecture/02-engine-adapter.md) | adapter 契约、代际语义、vLLM 实现、replay、快照信号 |
42+
| [03-health-and-liveness](docs/architecture/03-health-and-liveness.md) | epoch、判死、HostDown、同生同死探针 |
43+
| [04-kvcm-client](docs/architecture/04-kvcm-client.md) | 传输、服务发现、两步注册、心跳、location spec |
44+
| [05-observability](docs/architecture/05-observability.md) | span 阶段、指标出口、日志、trace id |
45+
| [06-configuration](docs/architecture/06-configuration.md) | 配置优先级、字段分组、环境变量、校验 |
46+
47+
变更设计文档在 `docs/specs/`,实施计划在 `docs/plans/`,评审与测试报告在 `docs/reviews/`
48+
指标目录在 `docs/metrics.json`。开发约束与流程见 [AGENTS.md](AGENTS.md)
49+
50+
## Installation
51+
52+
```bash
53+
uv sync --dev
54+
uv run pre-commit install # clone 后必须执行一次
55+
```
56+
57+
## Running
58+
59+
Required at startup:
60+
61+
- `SPECTRUM_DEPLOYMENT_NAME` — unique deployment identity used to build the
62+
KVCM `instance_id`; startup fails if missing or blank.
63+
- KVCM base URL — via `--kvcm-base-url` or the config file; a blank value is
64+
rejected by config validation. Default KVCM protocol is gRPC, so the port
65+
must be KVCM `meta_rpc_port`; use `--kvcm-protocol http` with
66+
`meta_http_port` for rollback. For gRPC, bare `host:port`, `grpc://`, and
67+
`http(s)://` are direct channel targets; `static://` and `spectrum://` are
68+
resolved through service discovery.
69+
- `--host-port` — worker identity port advertised to KVCM as `host_ip_port`;
70+
no default value, must match the engine endpoint port that FlexLB discovers
71+
via Spectrum (e.g. 8080 on PAI-EAS). Config validation rejects a missing or
72+
out-of-range value.
73+
74+
```bash
75+
# With CLI args
76+
SPECTRUM_DEPLOYMENT_NAME=my-deployment uv run python -m subscriber \
77+
--kvcm-base-url spectrum://vs-example:6381 \
78+
--host-port 8080 \
79+
--engine-type vllm
80+
81+
# With config file (kvcm_base_url and host_port can be set in the yaml)
82+
SPECTRUM_DEPLOYMENT_NAME=my-deployment uv run python -m subscriber --config config.yaml
83+
```
84+
85+
## Development
86+
87+
```bash
88+
# Run tests
89+
uv run pytest
90+
uv run pytest tests/engine/vllm/test_incremental.py # 单文件
91+
92+
# Lint / format
93+
uv run ruff check subscriber/ tests/ harness/
94+
uv run ruff format subscriber/ tests/ harness/
95+
96+
# Type check
97+
uv run mypy subscriber/
98+
99+
# Manifest-driven cross-repository gates and evolution review
100+
harness/run_local_checks.sh baseline all
101+
harness/run_local_checks.sh protocol all
102+
harness/run_local_checks.sh quality all
103+
uv run python harness/loop.py review
104+
```
105+
106+
### Quality Gates
107+
108+
- **pre-commit(本地强制)**:commit 时按暂存文件路由执行——Python/harness 变更触发
109+
ruff check / ruff format --check,subscriber/config/toolchain 变更触发 mypy 与全量
110+
`pytest --cov`(覆盖率门禁 `fail_under = 90`);`subscriber/proto/` 文件额外触发
111+
全部 pb import、authoritative proto/runtime/pyi parity 与 focused proto/client 测试;
112+
metrics catalog 与 staged whitespace/conflict marker 也有独立检查。不得用 `--no-verify`
113+
跳过。
114+
- **pytest(开发中随手跑)**:commit 时 pre-commit 会强制全量测试,
115+
开发过程中仍应随手 `uv run pytest` 快速反馈。测试未通过的修改视为未完成。
116+
- **AoneCI(MR 触发)**:MR opened 时触发「Python单元测试」(`pip install '.[dev]'`
117+
后跑 pytest)与「代码质量扫描(多语言)」两条流水线;dev 依赖必须写在
118+
`[project.optional-dependencies].dev`,CI 状态用 `a1 ci run get` 查询。
119+
- **e2e(显式 opt-in)**`tests/e2e/` 需要设置 `DSV_BINARY` 指向已构建的
120+
`dashservingd` 才会运行,默认全部跳过,保证各机器上默认测试行为一致。
121+
- **harness(跨仓复用)**`harness/manifest.yaml` 声明 repository/profile/gate,
122+
`harness/loop.py` 不经 shell 执行检查并把原子 JSON 证据写入已忽略的
123+
`harness/records/runs/`;反馈采用 append-only ledger,`review` 只生成演进建议,
124+
不会自动修改 gate 或源码。
125+
126+
## 基于 Agent 的开发流程
127+
128+
以下是本仓库的标准开发流程。**Agent 每完成一步必须停下来等确认,不要自动进入下一步。**
129+
130+
### 1. 明确需求与范围
131+
132+
**Developer:** 描述需求与边界。
133+
134+
**Agent:** 复述需求、列出不做什么、把不清楚的点问出来(不要猜)。
135+
136+
### 2. 分析代码
137+
138+
**Agent:** 先读 `docs/architecture/` 对应文档,再读源码,输出数据流、关键调用路径、受影响的
139+
不变量;发现文档与代码不一致先报告。
140+
141+
**Developer:** 补充背景与注意事项。
142+
143+
### 3. 确定方案
144+
145+
**Agent:** 输出方案并写入 `docs/specs/YYYY-MM-DD_slug.md`(方案、受影响的不变量、测试策略、
146+
风险)。
147+
148+
**Developer:** 确认方案,或指出问题让 Agent 修改。
149+
150+
### 4. 实施修改
151+
152+
**Agent:** 从最新 `master` 建分支(`feat/<description>` / `fix/<description>`),测试先行,
153+
小步提交。
154+
155+
**Developer:** 检查改动是否合理。
156+
157+
### 5. 质量门禁
158+
159+
**Agent:**`uv run pytest``uv run ruff check subscriber/ tests/ harness/`
160+
`uv run mypy subscriber/`,贴出实际结论。
161+
162+
**Developer:** 确认全绿。
163+
164+
### 6. 提交与评审
165+
166+
**Agent:** commit → push 分支 → 创建 MR(标题与 commit message 一致,描述含背景、改动点、
167+
验证方式、影响面与回滚方式;跨仓库改动列出对应 MR 链接)。
168+
169+
**Developer:** 完成代码评审。
170+
171+
### 7. 回顾与更新
172+
173+
**Agent:** 检查并更新 `docs/architecture/``AGENTS.md``README.md``docs/metrics.json`
174+
把落地方案与设计文档的差异回写到 spec。
175+
176+
## Protobuf Compatibility Maintenance
177+
178+
Production runs `protobuf==3.20.3`, and code emitted by `grpc_tools.protoc` imports
179+
`google.protobuf.runtime_version` / `_builder`, which cannot load on that runtime.
180+
The Python pb files in `subscriber/proto/` (`_pb2.py`, `_pb2_grpc.py`, `_pb2.pyi`) are
181+
therefore **hand-maintained** — do **not** run `grpc_tools.protoc` and commit its output.
182+
183+
`subscriber/proto/engine_service_rpc.proto` remains the authoritative engine wire
184+
schema. KVCM meta-service bindings are a hand-maintained subset of KVCM's
185+
authoritative `meta_service.proto`. After changing either schema, manually sync the
186+
pb files following the procedure in [AGENTS.md](AGENTS.md) (section "Protobuf
187+
兼容代码维护"): keep the `descriptor_pb2` / `message_factory` implementation
188+
compatible with protobuf 3.20.3, never introduce `runtime_version` or `_builder`,
189+
and only touch `_pb2_grpc.py` when the RPC service/method/request/response types
190+
change.
191+
192+
The pb files are committed to the repository and are never generated at build time.
193+
After any change, run the proto/client tests and `uv run mypy subscriber/` to confirm
194+
they import cleanly under protobuf 3.20.3.

0 commit comments

Comments
 (0)