Skip to content

Commit 764ee05

Browse files
authored
[service] expose optimizer events over gRPC streaming
1 parent 9275345 commit 764ee05

10 files changed

Lines changed: 482 additions & 10 deletions

File tree

docs/configuration.md

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -150,8 +150,11 @@ kvcm.metrics.enable_prometheus=true
150150
# Prometheus metrics名称前缀,默认kvcm
151151
kvcm.metrics.prometheus_prefix=kvcm
152152
153-
# log event publisher的初始化配置值,暂未启用
154-
kvcm.event.event_publishers_configs
153+
# Event publisher 配置。log 默认开启;optimizer.enable=true 时,会在现有
154+
# kvcm.service.rpc_port 上注册 OptimizerEventStreamService,不新增监听端口。
155+
# 每个 publisher 的 queue_size 是它自己的发布队列上限;max_subscribers 是并发订阅数上限,
156+
# subscriber_queue_size 是每个订阅者的独立缓冲上限;字段省略时使用下列默认值。
157+
kvcm.event.event_publishers_configs={"log":{"enable":true,"queue_size":10000},"optimizer":{"enable":true,"queue_size":100000,"max_subscribers":4,"subscriber_queue_size":10000}}
155158
```
156159

157160
### SchedulePlanExecutor 线程与迁移预算

kv_cache_manager/service/BUILD

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,10 @@ cc_library(
3030
"//kv_cache_manager/config:coordination_backend",
3131
"//kv_cache_manager/config:leader_elector_base",
3232
"//kv_cache_manager/event:event_publisher",
33+
"//kv_cache_manager/event:event_publishers_config",
34+
"//kv_cache_manager/event/optimizer_stream:subscription_event_sink",
3335
"//kv_cache_manager/manager:cache_garbage_collector",
36+
"//kv_cache_manager/event:optimizer_event_publisher",
3437
"//kv_cache_manager/manager:cache_location_view",
3538
"//kv_cache_manager/manager:cache_manager",
3639
"//kv_cache_manager/manager:write_location_manager",
@@ -39,6 +42,7 @@ cc_library(
3942
"//kv_cache_manager/metrics:metrics_reporter",
4043
"//kv_cache_manager/metrics:metrics_reporter_factory",
4144
"//kv_cache_manager/service/grpc_service",
45+
"//kv_cache_manager/service/grpc_service:optimizer_event_service_grpc",
4246
"//kv_cache_manager/service/http_service",
4347
"//kv_cache_manager/service/util:manager_message_proto_util",
4448
"//kv_cache_manager/service/util:service_call_guard",

kv_cache_manager/service/grpc_service/BUILD

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,3 +23,21 @@ cc_library(
2323
"//kv_cache_manager/service:meta_service_metrics_base",
2424
],
2525
)
26+
27+
cc_library(
28+
name = "optimizer_event_service_grpc",
29+
srcs = [
30+
"optimizer_event_service_grpc.cc",
31+
],
32+
hdrs = [
33+
"optimizer_event_service_grpc.h",
34+
],
35+
deps = [
36+
"//kv_cache_manager/common",
37+
"//kv_cache_manager/common:logger",
38+
"//kv_cache_manager/config",
39+
"//kv_cache_manager/event/optimizer_stream:subscription_event_sink",
40+
"//kv_cache_manager/protocol/protobuf:service_cc_grpc",
41+
"//kv_cache_manager/service:grpc_wrapper",
42+
],
43+
)
Lines changed: 108 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,108 @@
1+
#include "kv_cache_manager/service/grpc_service/optimizer_event_service_grpc.h"
2+
3+
#include <chrono>
4+
#include <utility>
5+
6+
#include "kv_cache_manager/common/logger.h"
7+
#include "kv_cache_manager/common/request_context.h"
8+
#include "kv_cache_manager/config/instance_group.h"
9+
#include "kv_cache_manager/config/instance_info.h"
10+
#include "kv_cache_manager/config/registry_manager.h"
11+
#include "kv_cache_manager/event/optimizer_stream/subscription_event_sink.h"
12+
13+
namespace kv_cache_manager {
14+
15+
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)) {}
18+
19+
grpc::Status OptimizerEventServiceGRpc::GetConfiguration(grpc::ServerContext *,
20+
const proto::optimizer::KvcmConfigurationRequest *request,
21+
proto::optimizer::KvcmConfigurationResponse *response) {
22+
auto *status = response->mutable_header()->mutable_status();
23+
if (!registry_manager_) {
24+
status->set_code(proto::optimizer::SERVICE_NOT_READY);
25+
status->set_message("KVCM registry manager is unavailable");
26+
return grpc::Status::OK;
27+
}
28+
29+
RequestContext request_context(request->trace_id());
30+
const auto [group_ec, instance_groups] = registry_manager_->ListInstanceGroup(&request_context);
31+
if (group_ec != EC_OK) {
32+
status->set_code(proto::optimizer::INTERNAL_ERROR);
33+
status->set_message("Failed to list KVCM instance groups");
34+
return grpc::Status::OK;
35+
}
36+
37+
for (const auto &instance_group : instance_groups) {
38+
auto *group = response->add_instance_groups();
39+
group->set_name(instance_group->name());
40+
group->set_capacity_bytes(instance_group->quota().capacity());
41+
42+
const auto [instance_ec, instances] =
43+
registry_manager_->ListInstanceInfo(&request_context, instance_group->name());
44+
if (instance_ec != EC_OK) {
45+
response->clear_instance_groups();
46+
response->clear_instances();
47+
status->set_code(proto::optimizer::INTERNAL_ERROR);
48+
status->set_message("Failed to list KVCM instances");
49+
return grpc::Status::OK;
50+
}
51+
for (const auto &instance_info : instances) {
52+
auto *instance = response->add_instances();
53+
instance->set_instance_group_name(instance_info->instance_group_name());
54+
instance->set_instance_id(instance_info->instance_id());
55+
instance->set_block_size(instance_info->block_size());
56+
for (const auto &spec_info : instance_info->location_spec_infos()) {
57+
auto *spec = instance->add_location_spec_infos();
58+
spec->set_name(spec_info.name());
59+
spec->set_size(spec_info.size());
60+
}
61+
for (const auto &spec_group : instance_info->location_spec_groups()) {
62+
auto *group_config = instance->add_location_spec_groups();
63+
group_config->set_name(spec_group.name());
64+
for (const auto &spec_name : spec_group.spec_names()) {
65+
group_config->add_spec_names(spec_name);
66+
}
67+
}
68+
}
69+
}
70+
71+
status->set_code(proto::optimizer::OK);
72+
return grpc::Status::OK;
73+
}
74+
75+
grpc::Status
76+
OptimizerEventServiceGRpc::SubscribeEvents(grpc::ServerContext *context,
77+
const proto::optimizer::OptimizerEventSubscriptionRequest *request,
78+
grpc::ServerWriter<proto::optimizer::TraceQueryRequest> *writer) {
79+
if (!sink_ || sink_->stopped()) {
80+
return grpc::Status(grpc::StatusCode::UNAVAILABLE, "optimizer event publisher is unavailable");
81+
}
82+
auto subscription = sink_->Subscribe(request->consumer_id());
83+
if (!subscription) {
84+
return grpc::Status(grpc::StatusCode::RESOURCE_EXHAUSTED, "optimizer subscriber limit reached");
85+
}
86+
87+
KVCM_LOG_INFO("OptimizerEventServiceGRpc: stream opened, consumer_id=%s, peer=%s",
88+
subscription->consumer_id().c_str(),
89+
context->peer().c_str());
90+
while (!context->IsCancelled()) {
91+
proto::optimizer::TraceQueryRequest event;
92+
const auto result = subscription->WaitNext(&event, std::chrono::milliseconds(100));
93+
if (result == SubscriptionEventSink::Subscription::WaitResult::kTimeout) {
94+
continue;
95+
}
96+
if (result == SubscriptionEventSink::Subscription::WaitResult::kClosed) {
97+
break;
98+
}
99+
if (!writer->Write(event)) {
100+
break;
101+
}
102+
}
103+
sink_->Unsubscribe(subscription);
104+
KVCM_LOG_INFO("OptimizerEventServiceGRpc: stream closed, consumer_id=%s", subscription->consumer_id().c_str());
105+
return grpc::Status::OK;
106+
}
107+
108+
} // namespace kv_cache_manager
Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,32 @@
1+
#pragma once
2+
3+
#include <grpcpp/grpcpp.h>
4+
#include <memory>
5+
6+
#include "kv_cache_manager/protocol/protobuf/optimizer_service.grpc.pb.h"
7+
8+
namespace kv_cache_manager {
9+
10+
class SubscriptionEventSink;
11+
class RegistryManager;
12+
13+
// gRPC service for optimizer event subscriptions.
14+
class OptimizerEventServiceGRpc final : public proto::optimizer::OptimizerEventStreamService::Service {
15+
public:
16+
OptimizerEventServiceGRpc(std::shared_ptr<SubscriptionEventSink> sink,
17+
std::shared_ptr<RegistryManager> registry_manager);
18+
19+
grpc::Status GetConfiguration(grpc::ServerContext *context,
20+
const proto::optimizer::KvcmConfigurationRequest *request,
21+
proto::optimizer::KvcmConfigurationResponse *response) override;
22+
23+
grpc::Status SubscribeEvents(grpc::ServerContext *context,
24+
const proto::optimizer::OptimizerEventSubscriptionRequest *request,
25+
grpc::ServerWriter<proto::optimizer::TraceQueryRequest> *writer) override;
26+
27+
private:
28+
std::shared_ptr<SubscriptionEventSink> sink_;
29+
std::shared_ptr<RegistryManager> registry_manager_;
30+
};
31+
32+
} // namespace kv_cache_manager

kv_cache_manager/service/server.cc

Lines changed: 43 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,10 @@
1212
#include "kv_cache_manager/config/node_endpoint_info.h"
1313
#include "kv_cache_manager/config/registry_manager.h"
1414
#include "kv_cache_manager/event/event_manager.h"
15+
#include "kv_cache_manager/event/event_publishers_config.h"
1516
#include "kv_cache_manager/event/log_event_publisher.h"
17+
#include "kv_cache_manager/event/optimizer_event_publisher.h"
18+
#include "kv_cache_manager/event/optimizer_stream/subscription_event_sink.h"
1619
#include "kv_cache_manager/manager/cache_manager.h"
1720
#include "kv_cache_manager/manager/startup_config_loader.h"
1821
#include "kv_cache_manager/metrics/metrics_lifecycle.h"
@@ -25,6 +28,7 @@
2528
#include "kv_cache_manager/service/grpc_service/admin_service_grpc.h"
2629
#include "kv_cache_manager/service/grpc_service/debug_service_grpc.h"
2730
#include "kv_cache_manager/service/grpc_service/meta_service_grpc.h"
31+
#include "kv_cache_manager/service/grpc_service/optimizer_event_service_grpc.h"
2832
#include "kv_cache_manager/service/http_service/admin_service_http.h"
2933
#include "kv_cache_manager/service/http_service/debug_service_http.h"
3034
#include "kv_cache_manager/service/http_service/meta_service_http.h"
@@ -254,6 +258,9 @@ bool Server::StartRpcServer() {
254258
grpc::ServerBuilder builder;
255259
builder.AddListeningPort(server_address, grpc::InsecureServerCredentials());
256260
builder.RegisterService(meta_service_.get());
261+
if (optimizer_event_service_) {
262+
builder.RegisterService(optimizer_event_service_.get());
263+
}
257264
if (!use_separate_admin_server) {
258265
builder.RegisterService(admin_service_.get());
259266
}
@@ -372,18 +379,46 @@ void Server::CreateAndRegisterEventPublisher() {
372379
KVCM_LOG_WARN("do not have event manager, skip create and register event publisher.");
373380
return;
374381
}
375-
auto log_publisher = std::make_shared<LogEventPublisher>();
376-
// 这里的logpublisher初始化配置需要修改
377-
auto event_publishers_configs = config_.event_publishers_configs();
378-
if (!log_publisher->Init(event_publishers_configs)) {
379-
KVCM_LOG_ERROR("init log event publisher failed");
382+
RegisterEventPublishers(event_manager);
383+
}
384+
385+
void Server::RegisterEventPublishers(const std::shared_ptr<EventManager> &event_manager) {
386+
const auto &event_publishers_configs = config_.event_publishers_configs();
387+
EventPublishersConfig publishers_config;
388+
if (!event_publishers_configs.empty() && !publishers_config.FromJsonString(event_publishers_configs)) {
389+
KVCM_LOG_ERROR("parse event publisher config failed; event publishers disabled");
390+
return;
391+
}
392+
393+
if (publishers_config.enable_log_event_publisher()) {
394+
auto log_publisher = std::make_shared<LogEventPublisher>(publishers_config.log_event_publisher_config());
395+
if (!log_publisher->Init("")) {
396+
KVCM_LOG_ERROR("init log event publisher failed");
397+
return;
398+
}
399+
if (!event_manager->RegisterPublisher("log_event_publisher", log_publisher)) {
400+
KVCM_LOG_ERROR("add log event publisher failed");
401+
return;
402+
}
403+
KVCM_LOG_INFO("create and register log event publisher OK");
404+
}
405+
406+
if (!publishers_config.enable_optimizer_event_publisher()) {
407+
return;
408+
}
409+
const auto &config = publishers_config.optimizer_event_publisher_config();
410+
auto sink = std::make_shared<SubscriptionEventSink>(config);
411+
auto optimizer_publisher = std::make_shared<OptimizerEventPublisher>(sink, config);
412+
if (!optimizer_publisher->Init("")) {
413+
KVCM_LOG_ERROR("init optimizer event publisher failed; publisher disabled");
380414
return;
381415
}
382-
if (!event_manager->RegisterPublisher("log_event_publisher", log_publisher)) {
383-
KVCM_LOG_ERROR("add log event publisher failed");
416+
if (!event_manager->RegisterPublisher("optimizer_event_publisher", optimizer_publisher)) {
417+
KVCM_LOG_ERROR("add optimizer event publisher failed; publisher disabled");
384418
return;
385419
}
386-
KVCM_LOG_INFO("create and register event publisher OK");
420+
optimizer_event_service_ = std::make_shared<OptimizerEventServiceGRpc>(sink, registry_manager_);
421+
KVCM_LOG_INFO("create and register optimizer gRPC event publisher OK, rpc_port=%d", config_.GetServiceRpcPort());
387422
}
388423
bool Server::CreateLeaderElector() {
389424
auto coordination_uri = config_.GetCoordinationUri();

kv_cache_manager/service/server.h

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,10 +19,12 @@ class CoordinationBackend;
1919
class LeaderElector;
2020
class RegistryManager;
2121
class CacheManager;
22+
class EventManager;
2223

2324
class MetaServiceGRpc;
2425
class AdminServiceGRpc;
2526
class DebugServiceGRpc;
27+
class OptimizerEventServiceGRpc;
2628
class MetaServiceHttp;
2729
class AdminServiceHttp;
2830
class DebugServiceHttp;
@@ -46,6 +48,7 @@ class Server {
4648
void CreateMetricsReporter();
4749
bool StartMetricsReportThread();
4850
void CreateAndRegisterEventPublisher();
51+
void RegisterEventPublishers(const std::shared_ptr<EventManager> &event_manager);
4952
bool CreateLeaderElector();
5053

5154
void OnBecomeLeader();
@@ -63,6 +66,7 @@ class Server {
6366
std::shared_ptr<MetaServiceGRpc> meta_service_;
6467
std::shared_ptr<AdminServiceGRpc> admin_service_;
6568
std::shared_ptr<DebugServiceGRpc> debug_service_;
69+
std::shared_ptr<OptimizerEventServiceGRpc> optimizer_event_service_;
6670
std::shared_ptr<grpc::Server> rpc_server_;
6771
std::shared_ptr<grpc::Server> admin_rpc_server_;
6872
std::shared_ptr<MetaServiceHttp> meta_http_service_;

kv_cache_manager/service/test/BUILD

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,43 @@ cc_test(
4343
],
4444
)
4545

46+
cc_test(
47+
name = "ServerEventPublisherTest",
48+
srcs = [
49+
"server_event_publisher_test.cc",
50+
],
51+
copts = ["-fno-access-control"],
52+
data = [],
53+
deps = [
54+
"//kv_cache_manager/common:unittest",
55+
"//kv_cache_manager/event:event_publisher",
56+
"//kv_cache_manager/event:optimizer_event_publisher",
57+
"//kv_cache_manager/event/optimizer_stream:subscription_event_sink",
58+
"//kv_cache_manager/service",
59+
],
60+
)
61+
62+
cc_test(
63+
name = "OptimizerEventServiceGRpcTest",
64+
srcs = [
65+
"optimizer_event_service_grpc_test.cc",
66+
],
67+
copts = ["-fno-access-control"],
68+
data = [],
69+
deps = [
70+
"//kv_cache_manager/common:unittest",
71+
"//kv_cache_manager/config",
72+
"//kv_cache_manager/event:event_publishers_config",
73+
"//kv_cache_manager/event:optimizer_event_publisher",
74+
"//kv_cache_manager/event/optimizer_stream:subscription_event_sink",
75+
"//kv_cache_manager/event/spec_events",
76+
"//kv_cache_manager/metrics:metrics_registry",
77+
"//kv_cache_manager/protocol/protobuf:service_cc_grpc",
78+
"//kv_cache_manager/service:grpc_wrapper",
79+
"//kv_cache_manager/service/grpc_service:optimizer_event_service_grpc",
80+
],
81+
)
82+
4683
cc_test(
4784
name = "AdminServiceRemoveInstanceGroupTest",
4885
srcs = [

0 commit comments

Comments
 (0)