Skip to content

Commit 47af2fc

Browse files
authored
[event] harden optimizer event stream lifecycle
1 parent 764ee05 commit 47af2fc

14 files changed

Lines changed: 245 additions & 32 deletions

AGENTS.md

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,11 +19,12 @@ flowchart TD
1919
data_storage --> common["common"]
2020
2121
manager --> event["event"]
22+
event --> protocol["protocol(proto/grpc)"]
2223
manager --> metrics["metrics"]
2324
service --> metrics
2425
service --> config
2526
service --> data_storage
26-
manager --> protocol["protocol(proto/grpc)"]
27+
manager --> protocol
2728
2829
%% 有意的反向边(近似环,改动需谨慎)
2930
common -. request_context .-> metrics

docs/design/module_architecture.md

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -43,8 +43,8 @@ KVCache Manager 采用中心化部署,负责 KVCache 的全局元数据管理
4343
|---|---|---|
4444
| **common** | `common/` | 基础设施层:日志、JSON、错误码、Redis 客户端、`RequestContext`(逐请求追踪上下文)、`concurrent_hash_map``lru_cache``loop_thread`、服务发现、崩溃处理等。几乎所有 C++ 模块都依赖它。 |
4545
| **metrics** | `metrics/` | 可观测性。`MetricsRegistry`/`MetricsCollector` 收集指标,多种 reporter(kmonitor/local/logging/dummy)上报,`PrometheusExporter` 通过 HTTP 暴露。 |
46-
| **event** | `event/` | 轻量事件总线。`EventManager` 将领域事件(如 cache 回收事件)分发给注册的 `EventPublisher`(默认 `LogEventPublisher`|
47-
| **protocol** | `protocol/protobuf/` | gRPC/proto 契约。定义 meta/admin/debug/kv_meta 服务,生成 C++ 与 Python 桩。 |
46+
| **event** | `event/` | 轻量事件总线。`EventManager` 将领域事件分发给注册的 `EventPublisher`;除默认日志发布器外,`OptimizerEventPublisher` 会将缓存读取事件转换为 protocol 中的 `TraceQueryRequest`,再经 `SubscriptionEventSink` 交给 service 层的 gRPC 流|
47+
| **protocol** | `protocol/protobuf/` | gRPC/proto 契约。定义 meta/admin/debug/kv_meta 以及 optimizer 事件流服务,生成 C++ 与 Python 桩。 |
4848

4949
### 客户端与连接器
5050

@@ -83,7 +83,7 @@ client 通过 `InitParams.role_type` 区分角色:**SCHEDULER**(调度节点
8383
service → manager → meta → config → data_storage → common
8484
```
8585

86-
- `common``protocol` 是最底层的通用模块,被各层广泛依赖。
86+
- `common``protocol` 是最底层的通用模块,被各层广泛依赖`event` 的 optimizer 发布链路直接使用 `protocol` 定义的 `TraceQueryRequest`
8787
- `metrics``event` 是通用支撑模块,被 `manager``service` 复用。
8888
- `service` 在启动时实例化 `CacheManager`(注入 `MetricsRegistry` + `RegistryManager`),并通过 `config``LeaderElector` 门控 recover/cleanup。
8989

@@ -138,6 +138,7 @@ flowchart TD
138138
139139
%% 通用支撑依赖
140140
manager --> event
141+
event --> protocol
141142
manager --> metrics
142143
service --> metrics
143144
service --> event

kv_cache_manager/event/optimizer_event_publisher.cc

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -67,9 +67,8 @@ bool OptimizerEventPublisher::Stop() {
6767
}
6868

6969
void OptimizerEventPublisher::WorkerThread() {
70-
// Single threaded on purpose: the replay requires per-instance events to
71-
// arrive in non-decreasing timestamp order, and one worker draining one
72-
// queue preserves that ordering.
70+
// A single worker keeps conversion and delivery off serving threads
71+
// without adding synchronization between multiple consumers of the queue.
7372
while (running_) {
7473
BasicWait();
7574

kv_cache_manager/event/optimizer_event_publisher.h

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -28,9 +28,8 @@ class TraceQueryRequest;
2828
// the event is dropped rather than the caller blocked: these are analysis
2929
// samples, and a slow or absent consumer must never slow down serving.
3030
//
31-
// Dropping biases the replayed hit rate downwards (an unseen request never
32-
// touches the simulated LRU). That is accepted, but it means the drop count
33-
// has to be readable before anyone trusts a capacity curve.
31+
// Dropping samples may slightly affect replay accuracy, which is accepted for
32+
// this best-effort analysis path.
3433
class OptimizerEventPublisher : public EventPublisher {
3534
public:
3635
OptimizerEventPublisher(std::shared_ptr<EventSink> sink, const OptimizerEventPublisherConfig &config);

kv_cache_manager/event/optimizer_stream/subscription_event_sink.cc

Lines changed: 24 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,7 @@ SubscriptionEventSink::~SubscriptionEventSink() { Stop(); }
5858

5959
std::shared_ptr<SubscriptionEventSink::Subscription> SubscriptionEventSink::Subscribe(const std::string &consumer_id) {
6060
std::lock_guard<std::mutex> lock(subscriptions_mutex_);
61-
if (stopped_ || subscriptions_.size() >= config_.max_subscribers()) {
61+
if (stopped_ || !accepting_subscriptions_ || subscriptions_.size() >= config_.max_subscribers()) {
6262
return nullptr;
6363
}
6464
auto subscription = std::shared_ptr<Subscription>(
@@ -87,15 +87,33 @@ void SubscriptionEventSink::Unsubscribe(const std::shared_ptr<Subscription> &sub
8787
subscribers);
8888
}
8989

90-
bool SubscriptionEventSink::Send(const proto::optimizer::TraceQueryRequest &event) {
91-
if (stopped_) {
92-
dropped_.fetch_add(1);
93-
return false;
90+
void SubscriptionEventSink::EnableSubscriptions() {
91+
std::lock_guard<std::mutex> lock(subscriptions_mutex_);
92+
if (!stopped_) {
93+
accepting_subscriptions_ = true;
9494
}
95+
}
9596

97+
void SubscriptionEventSink::DisableSubscriptions() {
9698
std::vector<std::shared_ptr<Subscription>> subscriptions;
9799
{
98100
std::lock_guard<std::mutex> lock(subscriptions_mutex_);
101+
accepting_subscriptions_ = false;
102+
subscriptions.swap(subscriptions_);
103+
}
104+
for (const auto &subscription : subscriptions) {
105+
subscription->Close();
106+
}
107+
}
108+
109+
bool SubscriptionEventSink::Send(const proto::optimizer::TraceQueryRequest &event) {
110+
std::vector<std::shared_ptr<Subscription>> subscriptions;
111+
{
112+
std::lock_guard<std::mutex> lock(subscriptions_mutex_);
113+
if (stopped_ || !accepting_subscriptions_) {
114+
dropped_.fetch_add(1);
115+
return false;
116+
}
99117
subscriptions = subscriptions_;
100118
}
101119
if (subscriptions.empty()) {
@@ -118,6 +136,7 @@ void SubscriptionEventSink::Stop() {
118136
if (stopped_.exchange(true)) {
119137
return;
120138
}
139+
accepting_subscriptions_ = false;
121140
std::vector<std::shared_ptr<Subscription>> subscriptions;
122141
{
123142
std::lock_guard<std::mutex> lock(subscriptions_mutex_);

kv_cache_manager/event/optimizer_stream/subscription_event_sink.h

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,10 @@ class SubscriptionEventSink final : public EventSink {
5656
std::shared_ptr<Subscription> Subscribe(const std::string &consumer_id);
5757
void Unsubscribe(const std::shared_ptr<Subscription> &subscription);
5858

59+
void EnableSubscriptions();
60+
void DisableSubscriptions();
61+
bool accepting_subscriptions() const { return accepting_subscriptions_.load(); }
62+
5963
bool Send(const proto::optimizer::TraceQueryRequest &event) override;
6064
void Stop() override;
6165
std::size_t DroppedCount() const;
@@ -67,6 +71,7 @@ class SubscriptionEventSink final : public EventSink {
6771
OptimizerEventPublisherConfig config_;
6872
mutable std::mutex subscriptions_mutex_;
6973
std::vector<std::shared_ptr<Subscription>> subscriptions_;
74+
std::atomic<bool> accepting_subscriptions_{true};
7075
std::atomic<bool> stopped_{false};
7176
std::atomic<std::size_t> dropped_{0};
7277
};

kv_cache_manager/event/test/subscription_event_sink_test.cc

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,27 @@ TEST_F(SubscriptionEventSinkTest, TestRejectsExcessSubscribers) {
7777
EXPECT_EQ(nullptr, sink.Subscribe("second"));
7878
}
7979

80+
TEST_F(SubscriptionEventSinkTest, TestDisableClosesSubscribersAndEnableAcceptsNewOnes) {
81+
SubscriptionEventSink sink(MakeSinkConfig(1, 1));
82+
auto subscription = sink.Subscribe("optimizer");
83+
ASSERT_NE(nullptr, subscription);
84+
85+
sink.DisableSubscriptions();
86+
proto::optimizer::TraceQueryRequest event;
87+
EXPECT_EQ(SubscriptionEventSink::Subscription::WaitResult::kClosed,
88+
subscription->WaitNext(&event, std::chrono::milliseconds(10)));
89+
EXPECT_EQ(nullptr, sink.Subscribe("disabled"));
90+
EXPECT_FALSE(sink.Send(MakeRequest("disabled")));
91+
92+
sink.EnableSubscriptions();
93+
auto resumed = sink.Subscribe("resumed");
94+
ASSERT_NE(nullptr, resumed);
95+
EXPECT_TRUE(sink.Send(MakeRequest("resumed")));
96+
EXPECT_EQ(SubscriptionEventSink::Subscription::WaitResult::kEvent,
97+
resumed->WaitNext(&event, std::chrono::milliseconds(10)));
98+
EXPECT_EQ("resumed", event.trace_id());
99+
}
100+
80101
TEST_F(SubscriptionEventSinkTest, TestStopWakesSubscriberAndIsIdempotent) {
81102
SubscriptionEventSink sink(MakeSinkConfig(1, 1));
82103
auto subscription = sink.Subscribe("optimizer");

kv_cache_manager/service/grpc_service/BUILD

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@ cc_library(
3636
"//kv_cache_manager/common",
3737
"//kv_cache_manager/common:logger",
3838
"//kv_cache_manager/config",
39+
"//kv_cache_manager/config:leader_elector_base",
3940
"//kv_cache_manager/event/optimizer_stream:subscription_event_sink",
4041
"//kv_cache_manager/protocol/protobuf:service_cc_grpc",
4142
"//kv_cache_manager/service:grpc_wrapper",

kv_cache_manager/service/grpc_service/optimizer_event_service_grpc.cc

Lines changed: 36 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -7,22 +7,46 @@
77
#include "kv_cache_manager/common/request_context.h"
88
#include "kv_cache_manager/config/instance_group.h"
99
#include "kv_cache_manager/config/instance_info.h"
10+
#include "kv_cache_manager/config/leader_elector.h"
1011
#include "kv_cache_manager/config/registry_manager.h"
1112
#include "kv_cache_manager/event/optimizer_stream/subscription_event_sink.h"
1213

1314
namespace kv_cache_manager {
1415

1516
OptimizerEventServiceGRpc::OptimizerEventServiceGRpc(std::shared_ptr<SubscriptionEventSink> sink,
16-
std::shared_ptr<RegistryManager> registry_manager)
17-
: sink_(std::move(sink)), registry_manager_(std::move(registry_manager)) {}
17+
std::shared_ptr<RegistryManager> registry_manager,
18+
std::shared_ptr<LeaderElector> leader_elector)
19+
: sink_(std::move(sink))
20+
, registry_manager_(std::move(registry_manager))
21+
, leader_elector_(std::move(leader_elector)) {
22+
DisableSubscriptions();
23+
}
24+
25+
void OptimizerEventServiceGRpc::EnableSubscriptions() {
26+
if (sink_) {
27+
sink_->EnableSubscriptions();
28+
}
29+
}
30+
31+
void OptimizerEventServiceGRpc::DisableSubscriptions() {
32+
if (sink_) {
33+
sink_->DisableSubscriptions();
34+
}
35+
}
36+
37+
bool OptimizerEventServiceGRpc::IsAvailable() const {
38+
return leader_elector_ && leader_elector_->GetRoleState() == RoleState::LEADER &&
39+
leader_elector_->IsStableState() && registry_manager_ && registry_manager_->IsRecoverComplete() && sink_ &&
40+
sink_->accepting_subscriptions() && !sink_->stopped();
41+
}
1842

1943
grpc::Status OptimizerEventServiceGRpc::GetConfiguration(grpc::ServerContext *,
2044
const proto::optimizer::KvcmConfigurationRequest *request,
2145
proto::optimizer::KvcmConfigurationResponse *response) {
2246
auto *status = response->mutable_header()->mutable_status();
23-
if (!registry_manager_) {
47+
if (!IsAvailable()) {
2448
status->set_code(proto::optimizer::SERVICE_NOT_READY);
25-
status->set_message("KVCM registry manager is unavailable");
49+
status->set_message("KVCM is unavailable");
2650
return grpc::Status::OK;
2751
}
2852

@@ -76,24 +100,26 @@ grpc::Status
76100
OptimizerEventServiceGRpc::SubscribeEvents(grpc::ServerContext *context,
77101
const proto::optimizer::OptimizerEventSubscriptionRequest *request,
78102
grpc::ServerWriter<proto::optimizer::TraceQueryRequest> *writer) {
79-
if (!sink_ || sink_->stopped()) {
80-
return grpc::Status(grpc::StatusCode::UNAVAILABLE, "optimizer event publisher is unavailable");
103+
if (!IsAvailable()) {
104+
return grpc::Status(grpc::StatusCode::UNAVAILABLE, "KVCM is unavailable");
81105
}
82106
auto subscription = sink_->Subscribe(request->consumer_id());
83107
if (!subscription) {
84-
return grpc::Status(grpc::StatusCode::RESOURCE_EXHAUSTED, "optimizer subscriber limit reached");
108+
return grpc::Status(grpc::StatusCode::UNAVAILABLE, "KVCM is unavailable");
85109
}
86110

87111
KVCM_LOG_INFO("OptimizerEventServiceGRpc: stream opened, consumer_id=%s, peer=%s",
88112
subscription->consumer_id().c_str(),
89113
context->peer().c_str());
114+
bool subscription_closed = false;
90115
while (!context->IsCancelled()) {
91116
proto::optimizer::TraceQueryRequest event;
92117
const auto result = subscription->WaitNext(&event, std::chrono::milliseconds(100));
93118
if (result == SubscriptionEventSink::Subscription::WaitResult::kTimeout) {
94119
continue;
95120
}
96121
if (result == SubscriptionEventSink::Subscription::WaitResult::kClosed) {
122+
subscription_closed = true;
97123
break;
98124
}
99125
if (!writer->Write(event)) {
@@ -102,6 +128,9 @@ OptimizerEventServiceGRpc::SubscribeEvents(grpc::ServerContext *context,
102128
}
103129
sink_->Unsubscribe(subscription);
104130
KVCM_LOG_INFO("OptimizerEventServiceGRpc: stream closed, consumer_id=%s", subscription->consumer_id().c_str());
131+
if (subscription_closed) {
132+
return grpc::Status(grpc::StatusCode::UNAVAILABLE, "KVCM is unavailable");
133+
}
105134
return grpc::Status::OK;
106135
}
107136

kv_cache_manager/service/grpc_service/optimizer_event_service_grpc.h

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,12 +9,17 @@ namespace kv_cache_manager {
99

1010
class SubscriptionEventSink;
1111
class RegistryManager;
12+
class LeaderElector;
1213

1314
// gRPC service for optimizer event subscriptions.
1415
class OptimizerEventServiceGRpc final : public proto::optimizer::OptimizerEventStreamService::Service {
1516
public:
1617
OptimizerEventServiceGRpc(std::shared_ptr<SubscriptionEventSink> sink,
17-
std::shared_ptr<RegistryManager> registry_manager);
18+
std::shared_ptr<RegistryManager> registry_manager,
19+
std::shared_ptr<LeaderElector> leader_elector);
20+
21+
void EnableSubscriptions();
22+
void DisableSubscriptions();
1823

1924
grpc::Status GetConfiguration(grpc::ServerContext *context,
2025
const proto::optimizer::KvcmConfigurationRequest *request,
@@ -25,8 +30,11 @@ class OptimizerEventServiceGRpc final : public proto::optimizer::OptimizerEventS
2530
grpc::ServerWriter<proto::optimizer::TraceQueryRequest> *writer) override;
2631

2732
private:
33+
bool IsAvailable() const;
34+
2835
std::shared_ptr<SubscriptionEventSink> sink_;
2936
std::shared_ptr<RegistryManager> registry_manager_;
37+
std::shared_ptr<LeaderElector> leader_elector_;
3038
};
3139

3240
} // namespace kv_cache_manager

0 commit comments

Comments
 (0)