Skip to content

Commit 67cfab2

Browse files
committed
[manager] move EventReport cleanup into background GC
1 parent b9d7dd3 commit 67cfab2

47 files changed

Lines changed: 4681 additions & 171 deletions

Some content is hidden

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

docs/README.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
- [高可用与选主机制](design/ha_leader_elector.md) - HA 架构、LeaderElector 状态机、CoordinationBackend、Leader 发现
99
- [CacheReclaimer 异步删除与过度逐出优化](design/cache_reclaimer_async_delete.md) - 异步删除生命周期、in-flight credit、反压与无进展退避
1010
- [后台扫描 GC](design/cache_garbage_collector.md) - 基于 authoritative cursor 的后台全量巡检;V1 清理长期 orphan WRITING 和普通 SERVING storage-missing,并提供无副作用读取、精确值条件 CAS 与 HA 生命周期
11+
- [EventReport 主动回收纳入后台扫描 GC](design/event_report_background_gc.md) - 以可合并 cleanup intent 取代 per-event 全 Instance 扫描,统一 Snapshot generation 与 down host metadata 回收
1112

1213
### 开发文档
1314
- [开发指南](develop/README.md) - 开发者入门指南和开发环境配置

docs/configuration.md

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -134,6 +134,9 @@ kvcm.cache_gc.scan_batch_size=256
134134
kvcm.cache_gc.orphan_writing_grace_period_ms=86400000
135135
# GC 在途删除请求硬上限;默认 2。一个慢请求只占一个槽位,全部槽位占满后暂停扫描。
136136
kvcm.cache_gc.max_inflight_delete_requests=2
137+
# 迁移期开关:将 EventReport Snapshot/HOST_DOWN 回收路由到 GC cleanup intent lane;
138+
# 默认关闭,开启时要求 kvcm.cache_gc.enabled=true,配置变更需重启生效。
139+
kvcm.cache_gc.event_report_cleanup_enabled=false
137140
138141
# 可选值有dummy,local,logging,kmonitor;若不配置或配置为空,默认使用local
139142
kvcm.metrics.reporter_type=local

docs/design/cache_garbage_collector.md

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
| 项目 | 内容 |
44
|---|---|
55
| 状态 | V1 已实现并完成相关单测、E2E 与性能冒烟;基线已包含异步删除(#234)、分层存储迁移(#209)和精确值条件删除(#233|
6-
| 更新时间 | 2026-07-30 |
6+
| 更新时间 | 2026-08-14 |
77
| 涉及模块 | `manager``meta``data_storage``config``metrics``service` |
88
| 历史参考 | [PR #184](https://github.com/alibaba/tair-kvcache/pull/184) |
99

@@ -101,7 +101,7 @@ V1 实现以下六项能力:
101101
4. 动态或大规模并发窗口、多维 bytes/Group 配额、Future deadline 和持久化任务。
102102
5. `WriteLocationManager` 三元组索引、settling guard,以及 Reclaimer/Finish 竞态修复。
103103
6. `CLS_DELETING` 自动恢复和 storage version/epoch。
104-
7. EventReport reconciliation、Block TTL 和闲置 Instance 清理
104+
7. EventReport reconciliation、Block TTL 和闲置 Instance 清理不属于本文定义的 GC V1;其中 EventReport 已作为独立扩展接入同一后台 GC,详见 [EventReport 主动回收纳入后台扫描 GC](event_report_background_gc.md)
105105
8. 查询链路 fast-submit 重构。
106106
9. GC lease/generation、跨进程 cursor 或 pending 恢复。
107107
10. data storage I/O timeout/cancel 和 CAS 后补偿状态机。
@@ -605,11 +605,11 @@ Block TTL 需要先定义 TTL 写入、刷新、authoritative `expire_at`、并
605605
4. GC task 持久化、lease/generation 和跨进程恢复;
606606
5. data storage I/O timeout/cancel;
607607
6. authoritative time source;
608-
7. EventReport snapshot/event 的后台 reconciliation;其 metadata-only、版本和所有权语义不能复用普通 storage 的物理删除规则。
608+
7. EventReport snapshot/event 的后台 reconciliation 已作为独立扩展设计和交付;其 metadata-only、版本和所有权语义不复用普通 storage 的物理删除规则,详见 [EventReport 主动回收纳入后台扫描 GC](event_report_background_gc.md)
609609

610610
### 11.6 建议交付顺序
611611

612612
1. **GC V1 PR**:本文定义的 full scan、固定 WRITING 与普通 SERVING storage-missing 判定、精确值条件删除、有界 Future/pending、leader 生命周期、指标和测试。
613613
2. **公共基础设施跟进(按需)**:若现网验证确认 Redis 命令缺少有限返回上界,再独立补齐 connect/auth/command/reconnect timeout 和坏连接淘汰。
614614
3. **独立正确性 PR**:WriteLocationManager Instance 隔离和 Reclaimer/Finish settling race。
615-
4. **后续 PR**:根据性能和业务优先级分别接入 pacing/index、显式三态 probe、TTL、EventReport 和 DELETING reconciliation。
615+
4. **后续 PR**:根据性能和业务优先级分别接入 pacing/index、显式三态 probe、TTL 和 DELETING reconciliation;EventReport reconciliation 按独立扩展文档交付

docs/design/event_report_background_gc.md

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

docs/design/module_architecture.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@
44

55
> **维护提示**:当模块的职责、依赖方向或调用关系发生变化,或新增/删除模块时,请同步更新本文档与文末的 Mermaid 图,并同步更新 [AGENTS.md](../../AGENTS.md) 中的缩略图。
66
7-
相关文档:[基本概念](basic_concepts.md)[ReportEvent Snapshot URI 版本方案](report_event_snapshot_uri_version.md)[高可用与选主机制](ha_leader_elector.md)[CacheReclaimer 异步删除设计](cache_reclaimer_async_delete.md)[后台扫描 GC 设计](cache_garbage_collector.md)[配置指南](../configuration.md)[优化器文档](../optimizer.md)
7+
相关文档:[基本概念](basic_concepts.md)[ReportEvent Snapshot URI 版本方案](report_event_snapshot_uri_version.md)[高可用与选主机制](ha_leader_elector.md)[CacheReclaimer 异步删除设计](cache_reclaimer_async_delete.md)[后台扫描 GC 设计](cache_garbage_collector.md)[EventReport 主动回收纳入后台 GC](event_report_background_gc.md)[配置指南](../configuration.md)[优化器文档](../optimizer.md)
88

99
---
1010

docs/design/report_event_snapshot_uri_version.md

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -395,6 +395,8 @@ snapshot replace + commit/abort
395395

396396
## 11. 限流与清理
397397

398+
本节保留 legacy fallback 的行为说明;启用 `kvcm.cache_gc.event_report_cleanup_enabled` 后,系统按 [EventReport 主动回收纳入后台扫描 GC](event_report_background_gc.md) 把 per-event 全 Instance 扫描替换为可合并 cleanup intent。两条路径的 Snapshot、generation、lifecycle fencing 和 metadata-only 契约保持一致。
399+
398400
`EventReportStorageSpec.snapshot_min_interval_ms` 提供 per-reporter 最小 snapshot 间隔,默认
399401
30 秒。完整维度是:
400402

integration_test/cache_garbage_collector/cache_garbage_collector_test.py

Lines changed: 113 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -214,8 +214,117 @@ def test_missing_serving_data_is_collected_without_query(self):
214214
"missing SERVING metadata should be removed and become writable",
215215
)
216216

217+
def test_event_report_host_down_is_collected_by_intent_lane(self):
218+
"""HOST_DOWN is reconciled by the metadata-only intent lane."""
219+
self.worker_manager.stop_worker(0)
220+
self.assertTrue(
221+
self.worker_manager.start_worker(
222+
0,
223+
**{
224+
"kvcm.cache_gc.enabled": "true",
225+
"kvcm.cache_gc.event_report_cleanup_enabled": "true",
226+
"kvcm.cache_gc.scan_interval_ms": 25,
227+
"kvcm.cache_gc.round_pause_ms": 100,
228+
"kvcm.cache_gc.scan_batch_size": 8,
229+
"kvcm.cache_gc.orphan_writing_grace_period_ms": self._GRACE_MS,
230+
"kvcm.cache_gc.max_inflight_delete_requests": 2,
231+
},
232+
)
233+
)
234+
self._admin_client.close()
235+
self._client.close()
236+
self._admin_client, self._client = self._get_manager_clients()
237+
238+
event_storage_name = "cache_gc_event_report"
239+
self._admin_client.add_storage(
240+
{
241+
"trace_id": f"{self._trace_id}_event_storage",
242+
"storage": {
243+
"global_unique_name": event_storage_name,
244+
"storage_type": "ST_EVENT_REPORT_L2",
245+
"event_report": {
246+
"heartbeat_timeout_ms": 60_000,
247+
"cleanup_grace_ms": 60_000,
248+
"liveness_check_interval_ms": 1_000,
249+
},
250+
"check_storage_available_when_open": False,
251+
},
252+
}
253+
)
254+
self._create_topology(
255+
event_report_storage_candidates=[event_storage_name]
256+
)
257+
instance_id = "cache_gc_event_report_instance"
258+
self._register_instance(instance_id)
259+
host = "192.168.80.1:8080"
260+
location_uri = f"vineyard://{host}/mem?source=cache_gc_e2e"
261+
self._client.report_event(
262+
{
263+
"trace_id": f"{self._trace_id}_event_add",
264+
"instance_id": instance_id,
265+
"host_ip_port": host,
266+
"storage_type": "ST_EVENT_REPORT_L2",
267+
"events": [
268+
{
269+
"event_type": "EVENT_NODE_REGISTER",
270+
"node_register": {"mediums": ["mem"]},
271+
},
272+
{
273+
"event_type": "EVENT_BLOCK_ADD",
274+
"block_add": {
275+
"block_key": str(self._block_key),
276+
"medium": "mem",
277+
"specs": [{"name": "tp0", "uri": location_uri}],
278+
},
279+
},
280+
],
281+
}
282+
)
283+
visible = self._client.get_cache_location(
284+
{
285+
"trace_id": f"{self._trace_id}_event_visible",
286+
"instance_id": instance_id,
287+
"query_type": "QT_BATCH_GET",
288+
"block_keys": [self._block_key],
289+
"block_mask": {"offset": 0},
290+
}
291+
)
292+
self.assertTrue(visible.get("locations"))
293+
294+
self._client.report_event(
295+
{
296+
"trace_id": f"{self._trace_id}_event_down",
297+
"instance_id": instance_id,
298+
"host_ip_port": host,
299+
"storage_type": "ST_EVENT_REPORT_L2",
300+
"events": [
301+
{"event_type": "EVENT_HOST_DOWN", "host_down": {}}
302+
],
303+
}
304+
)
305+
self._wait_metric_at_least(
306+
"cache_gc.event_report_delete_location_count",
307+
1,
308+
tags={"reason": "down_host", "status": "deleted"},
309+
timeout_s=10,
310+
)
311+
self.assertGreaterEqual(
312+
self._metric_value("cache_gc.event_report_scan_batch_count"), 1
313+
)
314+
self.assertEqual(
315+
0,
316+
self._metric_value(
317+
"cache_gc.event_report_intent_count",
318+
tags={"reason": "down_host"},
319+
),
320+
)
321+
self.assertEqual(0, self._metric_value("cache_gc.delete_target_count"))
322+
217323
def _create_topology(
218-
self, meta_storage_type="dummy", data_storage_type="nfs"
324+
self,
325+
meta_storage_type="dummy",
326+
data_storage_type="nfs",
327+
event_report_storage_candidates=None,
219328
):
220329
metadata_uri = (
221330
f"file://{self.get_workdir()}/cache_gc_metadata"
@@ -247,6 +356,9 @@ def _create_topology(
247356
"instance_group": {
248357
"name": self._group_name,
249358
"storage_candidates": [self._storage_name],
359+
"event_report_storage_candidates": (
360+
event_report_storage_candidates or []
361+
),
250362
"global_quota_group_name": "cache_gc_quota",
251363
"max_instance_count": 8,
252364
"quota": {

integration_test/meta_service/http_interface_test.py

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,10 @@ def finish_write_cache(self, data, check_response=True):
6161
"""Finish writing cache data"""
6262
return self._make_api_request('/api/finishWriteCache', data, check_response)
6363

64+
def report_event(self, data, check_response=True):
65+
"""Report external cache lifecycle events"""
66+
return self._make_api_request('/api/reportEvent', data, check_response)
67+
6468
def remove_cache(self, data, check_response=True):
6569
"""Remove cache data for specified block keys"""
6670
return self._make_api_request('/api/removeCache', data, check_response)

kv_cache_manager/common/cache/advanced_cache.h

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -500,6 +500,17 @@ class Cache {
500500
const std::string_view &key, ObjectPtr obj, size_t charge, const CacheItemHelper *helper)> & /*callback*/) {
501501
}
502502

503+
// Applies a callback to one existing entry without marking it as a hit or
504+
// changing its LRU position. The callback runs while the cache shard is
505+
// locked and returns the entry charge delta caused by an in-place update.
506+
// Callers must keep the callback bounded and must not call back into this
507+
// Cache. Returns false when the key is absent or unsupported.
508+
virtual bool ApplyToEntryNoTouch(
509+
const std::string_view & /*key*/,
510+
const std::function<ssize_t(ObjectPtr obj, size_t charge, const CacheItemHelper *helper)> & /*callback*/) {
511+
return false;
512+
}
513+
503514
// Insert a mapping from key->object only if the key does not already exist.
504515
// Returns EC_OK on successful insertion, EC_EXIST if the key is already
505516
// present, or other error codes on failure (e.g. EC_NOSPC).
@@ -777,6 +788,12 @@ class CacheWrapper : public Cache {
777788
target_->ApplyToSingleShard(shard_id, callback);
778789
}
779790

791+
bool ApplyToEntryNoTouch(
792+
const std::string_view &key,
793+
const std::function<ssize_t(ObjectPtr obj, size_t charge, const CacheItemHelper *helper)> &callback) override {
794+
return target_->ApplyToEntryNoTouch(key, callback);
795+
}
796+
780797
void StartAsyncLookup(AsyncLookupHandle &async_handle) override { target_->StartAsyncLookup(async_handle); }
781798

782799
void WaitAll(AsyncLookupHandle *async_handles, size_t count) override { target_->WaitAll(async_handles, count); }

kv_cache_manager/common/cache/lru_cache.cc

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -497,6 +497,52 @@ void LRUCacheShard::LookupBatch(const std::string_view *keys,
497497
}
498498
}
499499

500+
bool LRUCacheShard::ApplyToEntryNoTouch(
501+
const std::string_view &key,
502+
uint32_t hash,
503+
const std::function<ssize_t(Cache::ObjectPtr obj, size_t charge, const Cache::CacheItemHelper *helper)> &callback) {
504+
std::lock_guard<std::mutex> l(mutex_);
505+
LRUHandle *e = table_.Lookup(key, hash);
506+
if (e == nullptr) {
507+
return false;
508+
}
509+
assert(e->InCache());
510+
const bool in_lru = !e->HasRefs();
511+
const ssize_t delta = callback(e->value, e->total_charge, e->helper);
512+
if (delta > 0) {
513+
const size_t increase = static_cast<size_t>(delta);
514+
e->total_charge += increase;
515+
usage_ += increase;
516+
if (in_lru) {
517+
lru_usage_ += increase;
518+
if (e->InHighPriPool()) {
519+
high_pri_pool_usage_ += increase;
520+
} else if (e->InLowPriPool()) {
521+
low_pri_pool_usage_ += increase;
522+
}
523+
MaintainPoolSize();
524+
}
525+
} else if (delta < 0) {
526+
size_t decrease = static_cast<size_t>(-delta);
527+
decrease = std::min(decrease, e->total_charge);
528+
decrease = std::min(decrease, usage_);
529+
e->total_charge -= decrease;
530+
usage_ -= decrease;
531+
if (in_lru) {
532+
assert(lru_usage_ >= decrease);
533+
lru_usage_ -= decrease;
534+
if (e->InHighPriPool()) {
535+
assert(high_pri_pool_usage_ >= decrease);
536+
high_pri_pool_usage_ -= decrease;
537+
} else if (e->InLowPriPool()) {
538+
assert(low_pri_pool_usage_ >= decrease);
539+
low_pri_pool_usage_ -= decrease;
540+
}
541+
}
542+
}
543+
return true;
544+
}
545+
500546
bool LRUCacheShard::Ref(LRUHandle *e) {
501547
std::lock_guard<std::mutex> l(mutex_);
502548
// To create another reference - entry must be already externally referenced.

0 commit comments

Comments
 (0)