futu_grpc/server/
startup.rs1use std::sync::Arc;
2
3use tokio::sync::watch;
4
5use futu_auth::{KeyStore, RuntimeCounters};
6use futu_server::listener_status::{
7 ListenerBindEventSender, ListenerSurface, notify_listener_failed, notify_listener_opened,
8};
9use futu_server::router::RequestRouter;
10use tonic::transport::server::TcpIncoming;
11
12use crate::proto::futu_open_d_server::FutuOpenDServer;
13
14use super::{FutuGrpcService, GRPC_MAX_MESSAGE_SIZE_BYTES, GrpcPushBroadcaster};
15
16#[cfg(test)]
18#[allow(dead_code)]
19pub(super) fn build_service(
20 router: Arc<RequestRouter>,
21 push_broadcaster: Arc<GrpcPushBroadcaster>,
22) -> FutuOpenDServer<FutuGrpcService> {
23 apply_grpc_message_limits(FutuOpenDServer::new(FutuGrpcService::new(
24 router,
25 push_broadcaster,
26 )))
27}
28
29pub fn build_service_with_auth(
31 router: Arc<RequestRouter>,
32 push_broadcaster: Arc<GrpcPushBroadcaster>,
33 key_store: Arc<KeyStore>,
34 counters: Arc<RuntimeCounters>,
35) -> FutuOpenDServer<FutuGrpcService> {
36 apply_grpc_message_limits(FutuOpenDServer::new(FutuGrpcService::with_auth(
37 router,
38 push_broadcaster,
39 key_store,
40 counters,
41 )))
42}
43
44fn apply_grpc_message_limits(
45 service: FutuOpenDServer<FutuGrpcService>,
46) -> FutuOpenDServer<FutuGrpcService> {
47 service
48 .max_decoding_message_size(GRPC_MAX_MESSAGE_SIZE_BYTES)
49 .max_encoding_message_size(GRPC_MAX_MESSAGE_SIZE_BYTES)
50}
51
52pub async fn start(
54 listen_addr: &str,
55 router: Arc<RequestRouter>,
56 push_broadcaster: Arc<GrpcPushBroadcaster>,
57) -> Result<(), Box<dyn std::error::Error>> {
58 start_with_auth(
59 listen_addr,
60 router,
61 push_broadcaster,
62 Arc::new(KeyStore::empty()),
63 Arc::new(RuntimeCounters::new()),
64 )
65 .await
66}
67
68pub async fn start_with_auth(
70 listen_addr: &str,
71 router: Arc<RequestRouter>,
72 push_broadcaster: Arc<GrpcPushBroadcaster>,
73 key_store: Arc<KeyStore>,
74 counters: Arc<RuntimeCounters>,
75) -> Result<(), Box<dyn std::error::Error>> {
76 let addr: std::net::SocketAddr = listen_addr.parse().map_err(|error| {
77 futu_server::bind_hint::invalid_addr_error("gRPC", "--grpc-port", listen_addr, error)
78 })?;
79 if !key_store.is_configured() {
80 tracing::warn!(
81 "gRPC server running WITHOUT API key auth (legacy mode); \
82 all RPCs are open. Pass --grpc-keys-file to enable scope-based auth."
83 );
84 }
85 let service = build_service_with_auth(router, push_broadcaster, key_store, counters);
86 tracing::info!(addr = %listen_addr, "gRPC 服务已启动");
87 tonic::transport::Server::builder()
88 .add_service(service)
89 .serve(addr)
90 .await
91 .map_err(|error| {
92 futu_server::bind_hint::listener_runtime_error("gRPC", listen_addr, error)
93 })?;
94 Ok(())
95}
96
97pub async fn start_with_auth_until_shutdown(
100 listen_addr: &str,
101 router: Arc<RequestRouter>,
102 push_broadcaster: Arc<GrpcPushBroadcaster>,
103 key_store: Arc<KeyStore>,
104 counters: Arc<RuntimeCounters>,
105 shutdown_rx: watch::Receiver<bool>,
106) -> Result<(), Box<dyn std::error::Error>> {
107 start_with_auth_until_shutdown_with_listener_events(
108 listen_addr,
109 router,
110 push_broadcaster,
111 key_store,
112 counters,
113 shutdown_rx,
114 None,
115 )
116 .await?;
117 Ok(())
118}
119
120pub async fn start_with_auth_until_shutdown_with_listener_events(
122 listen_addr: &str,
123 router: Arc<RequestRouter>,
124 push_broadcaster: Arc<GrpcPushBroadcaster>,
125 key_store: Arc<KeyStore>,
126 counters: Arc<RuntimeCounters>,
127 shutdown_rx: watch::Receiver<bool>,
128 listener_events: Option<ListenerBindEventSender>,
129) -> std::io::Result<()> {
130 let addr: std::net::SocketAddr = listen_addr.parse().map_err(|error| {
131 notify_listener_failed(&listener_events, ListenerSurface::Grpc);
132 futu_server::bind_hint::invalid_addr_error("gRPC", "--grpc-port", listen_addr, error)
133 })?;
134 if !key_store.is_configured() {
135 tracing::warn!(
136 "gRPC server running WITHOUT API key auth (legacy mode); \
137 all RPCs are open. Pass --grpc-keys-file to enable scope-based auth."
138 );
139 }
140 let service = build_service_with_auth(router, push_broadcaster, key_store, counters);
141 let listener = tokio::net::TcpListener::bind(addr).await.map_err(|error| {
142 notify_listener_failed(&listener_events, ListenerSurface::Grpc);
143 futu_server::bind_hint::io_bind_error("gRPC", "--grpc-port", listen_addr, error)
144 })?;
145 let _serving =
146 notify_listener_opened(&listener_events, ListenerSurface::Grpc, &shutdown_rx).await?;
147 drop(listener_events);
148 tracing::info!(addr = %listen_addr, "gRPC 服务已启动");
149 tonic::transport::Server::builder()
150 .add_service(service)
151 .serve_with_incoming_shutdown(
152 grpc_tcp_incoming(listener),
153 grpc_shutdown_requested(shutdown_rx),
154 )
155 .await
156 .map_err(|error| {
157 futu_server::bind_hint::listener_runtime_error("gRPC", listen_addr, error)
158 })?;
159 Ok(())
160}
161
162fn grpc_tcp_incoming(listener: tokio::net::TcpListener) -> TcpIncoming {
163 TcpIncoming::from(listener).with_nodelay(Some(true))
164}
165
166async fn grpc_shutdown_requested(mut shutdown_rx: watch::Receiver<bool>) {
167 loop {
168 if *shutdown_rx.borrow() {
169 tracing::info!("gRPC server stopped by shutdown signal");
170 return;
171 }
172 if shutdown_rx.changed().await.is_err() {
173 tracing::info!("gRPC server stopped after shutdown sender dropped");
174 return;
175 }
176 }
177}
178
179#[cfg(test)]
180mod tests;