futu_opend/startup/phase4/
surfaces.rs1use std::sync::Arc;
9
10use futu_gateway_core::bridge::GatewayBridge;
11use futu_server::listener::{ApiServer, ServerConfig};
12use futu_server::listener_status::ListenerBindEvent;
13use futu_server::ws_listener::{WsServer, WsServerDeps};
14
15use crate::config::RuntimeConfig;
16
17use super::listener_readiness::wrap_grpc_server_result;
18use super::shutdown::SurfaceJoinHandle;
19use super::summary::{RestTransportKind, merge_push_health_snapshots_for_rest};
20
21#[allow(clippy::too_many_arguments)]
23pub(super) fn spawn_ws_server(
24 config: &RuntimeConfig,
25 server: &ApiServer,
26 server_config: &ServerConfig,
27 bridge: &Arc<GatewayBridge>,
28 ws_key_store: &Option<Arc<futu_auth::KeyStore>>,
29 shared_counters: &Arc<futu_auth::RuntimeCounters>,
30 shutdown_rx: &tokio::sync::watch::Receiver<bool>,
31 listener_events_tx: &tokio::sync::mpsc::UnboundedSender<ListenerBindEvent>,
32) -> (Option<SurfaceJoinHandle>, Option<Arc<futu_auth::KeyStore>>) {
33 let mut ws_key_store_holder: Option<std::sync::Arc<futu_auth::KeyStore>> = None;
36 let ws_handle = if let Some(ws_port) = config.websocket_port {
37 let ws_addr = format!("{}:{}", config.ip, ws_port);
38 let ws_key_store = ws_key_store.clone();
46 let ws_counters = std::sync::Arc::clone(shared_counters);
47 ws_key_store_holder = ws_key_store.as_ref().map(std::sync::Arc::clone);
49 let ws_server = WsServer::with_auth(
50 ws_addr.clone(),
51 server_config.clone(),
52 WsServerDeps::new(
53 std::sync::Arc::clone(server.connections()),
54 std::sync::Arc::clone(server.router()),
55 Some(bridge.subscription_runtime().manager()),
56 ),
57 ws_key_store,
58 Some(ws_counters),
59 )
60 .with_server_time_store(bridge.server_clock().anchor_store());
61 tracing::info!(addr = %ws_addr, "starting WebSocket server");
62 let ws_shutdown_rx = shutdown_rx.clone();
63 let ws_listener_events = listener_events_tx.clone();
64 Some(tokio::spawn(async move {
65 ws_server
66 .run_until_shutdown_with_listener_events(ws_shutdown_rx, Some(ws_listener_events))
67 .await
68 }))
69 } else {
70 None
71 };
72 (ws_handle, ws_key_store_holder)
73}
74
75#[allow(clippy::too_many_arguments)]
77pub(super) fn spawn_rest_server(
78 config: &RuntimeConfig,
79 server: &ApiServer,
80 bridge: &Arc<GatewayBridge>,
81 ws_broadcaster: &Arc<futu_rest::ws::WsBroadcaster>,
82 rest_key_store: &Option<Arc<futu_auth::KeyStore>>,
83 rest_transport: futu_rest::server::RestTransport,
84 rest_transport_kind: RestTransportKind,
85 shared_counters: &Arc<futu_auth::RuntimeCounters>,
86 shutdown_tx: &tokio::sync::watch::Sender<bool>,
87 shutdown_rx: &tokio::sync::watch::Receiver<bool>,
88 listener_events_tx: &tokio::sync::mpsc::UnboundedSender<ListenerBindEvent>,
89) -> (Option<SurfaceJoinHandle>, Option<Arc<futu_auth::KeyStore>>) {
90 let mut rest_key_store_holder: Option<std::sync::Arc<futu_auth::KeyStore>> = None;
91 let rest_handle = if let Some(rest_port) = config.rest_port {
92 let rest_addr = format!("{}:{}", config.ip, rest_port);
93 let router = std::sync::Arc::clone(server.router());
94 let broadcaster = std::sync::Arc::clone(ws_broadcaster);
95 let rest_key_store = rest_key_store
97 .clone()
98 .unwrap_or_else(|| std::sync::Arc::new(futu_auth::KeyStore::empty()));
99 let rest_transport_name = match rest_transport_kind {
100 RestTransportKind::Tls => "https",
101 RestTransportKind::Plaintext => "http",
102 RestTransportKind::Disabled => unreachable!("REST port is configured"),
103 };
104 tracing::info!(
105 addr = %rest_addr,
106 transport = rest_transport_name,
107 "starting REST API server (WebSocket: /ws)"
108 );
109
110 if rest_key_store.is_configured() {
114 rest_key_store_holder = Some(std::sync::Arc::clone(&rest_key_store));
115 }
116
117 let rest_counters = std::sync::Arc::clone(shared_counters);
124 let bridge_for_status = std::sync::Arc::clone(bridge);
128 let admin_status_provider: futu_rest::adapter::AdminStatusProvider =
129 std::sync::Arc::new(move || {
130 serde_json::to_value(bridge_for_status.snapshot_status())
131 .unwrap_or_else(|_| serde_json::json!({"error": "snapshot serialize failed"}))
132 });
133 let rest_shutdown_tx = shutdown_tx.clone();
134 let admin_shutdown_handler: futu_rest::adapter::AdminShutdownHandler =
135 std::sync::Arc::new(move || {
136 rest_shutdown_tx
137 .send(true)
138 .map_err(|e| format!("shutdown receiver dropped: {e}"))
139 });
140 let bridge_for_reload = std::sync::Arc::clone(bridge);
148 let admin_reload_handler: futu_rest::adapter::AdminReloadHandler =
149 std::sync::Arc::new(move || {
150 let bridge = std::sync::Arc::clone(&bridge_for_reload);
151 Box::pin(async move {
152 serde_json::to_value(bridge.reload())
153 .unwrap_or_else(|_| serde_json::json!({"error": "reload serialize failed"}))
154 })
155 });
156 let bridge_for_push_health = std::sync::Arc::clone(bridge);
163 let push_health_snapshot_provider: futu_rest::adapter::PushHealthSnapshotProvider =
164 std::sync::Arc::new(move || {
165 let push = serde_json::to_value(
166 bridge_for_push_health
167 .push_runtime()
168 .push_health()
169 .snapshot(),
170 )
171 .unwrap_or_else(
172 |_| serde_json::json!({"error": "push_health snapshot serialize failed"}),
173 );
174 let qot_login = serde_json::to_value(
175 bridge_for_push_health
176 .push_runtime()
177 .qot_login_health()
178 .snapshot(),
179 )
180 .unwrap_or_else(
181 |_| serde_json::json!({"error": "qot_login_health snapshot serialize failed"}),
182 );
183 let backend_connected =
184 bridge_for_push_health.broker_runtime().platform_connected();
185 let stock_list_status = bridge_for_push_health
186 .caches()
187 .static_cache
188 .stock_list_sync_status();
189 let stock_list_static_query_ready = matches!(
190 futu_domain_static_data::stock_list_static_query_gate_action_for_status(
191 &stock_list_status.domain_facts(),
192 ),
193 futu_domain_static_data::StockListStaticQueryGateAction::Allow
194 );
195 merge_push_health_snapshots_for_rest(
196 push,
197 qot_login,
198 backend_connected,
199 stock_list_static_query_ready,
200 )
201 });
202 let bridge_for_card_num = std::sync::Arc::clone(bridge);
208 let card_num_resolver: futu_rest::adapter::CardNumResolver =
209 std::sync::Arc::new(move |cn: &str| {
210 bridge_for_card_num
211 .caches()
212 .trd_cache
213 .find_acc_ids_by_card_num(cn)
214 });
215 let rest_shutdown_rx = shutdown_rx.clone();
216 let rest_listener_events = listener_events_tx.clone();
217 Some(tokio::spawn(async move {
218 futu_rest::server::start_with_auth_full_admin_until_shutdown_with_transport_and_listener_events(
219 &rest_addr,
220 router,
221 broadcaster,
222 rest_key_store,
223 rest_counters,
224 futu_rest::server::RestAdminHooks {
225 admin_status_provider: Some(admin_status_provider),
226 admin_shutdown_handler: Some(admin_shutdown_handler),
227 admin_reload_handler: Some(admin_reload_handler),
228 push_health_snapshot_provider: Some(push_health_snapshot_provider),
229 card_num_resolver: Some(card_num_resolver),
230 },
231 rest_transport,
232 rest_shutdown_rx,
233 Some(rest_listener_events),
234 )
235 .await
236 .map_err(anyhow::Error::from)
237 }))
238 } else {
239 None
240 };
241 (rest_handle, rest_key_store_holder)
242}
243
244pub(super) fn spawn_grpc_server(
246 config: &RuntimeConfig,
247 server: &ApiServer,
248 grpc_broadcaster: &Arc<futu_grpc::server::GrpcPushBroadcaster>,
249 grpc_key_store: &Option<Arc<futu_auth::KeyStore>>,
250 shared_counters: &Arc<futu_auth::RuntimeCounters>,
251 shutdown_rx: &tokio::sync::watch::Receiver<bool>,
252 listener_events_tx: &tokio::sync::mpsc::UnboundedSender<ListenerBindEvent>,
253) -> (Option<SurfaceJoinHandle>, Option<Arc<futu_auth::KeyStore>>) {
254 let mut grpc_key_store_holder: Option<std::sync::Arc<futu_auth::KeyStore>> = None;
255 let grpc_handle = if let Some(grpc_port) = config.grpc_port {
256 let grpc_addr = format!("{}:{}", config.ip, grpc_port);
257 let router = std::sync::Arc::clone(server.router());
258 let broadcaster = std::sync::Arc::clone(grpc_broadcaster);
259 let grpc_key_store = grpc_key_store
261 .clone()
262 .unwrap_or_else(|| std::sync::Arc::new(futu_auth::KeyStore::empty()));
263 tracing::info!(addr = %grpc_addr, "starting gRPC server (SubscribePush: streaming)");
264
265 if grpc_key_store.is_configured() {
267 grpc_key_store_holder = Some(std::sync::Arc::clone(&grpc_key_store));
268 }
269
270 let grpc_counters = std::sync::Arc::clone(shared_counters);
275 let grpc_shutdown_rx = shutdown_rx.clone();
276 let grpc_listener_events = listener_events_tx.clone();
277 Some(tokio::spawn(async move {
278 let result = futu_grpc::server::start_with_auth_until_shutdown_with_listener_events(
279 &grpc_addr,
280 router,
281 broadcaster,
282 grpc_key_store,
283 grpc_counters,
284 grpc_shutdown_rx,
285 Some(grpc_listener_events),
286 )
287 .await;
288 wrap_grpc_server_result(result)
289 }))
290 } else {
291 None
292 };
293 (grpc_handle, grpc_key_store_holder)
294}
295
296pub(super) fn spawn_telnet_server(
298 config: &RuntimeConfig,
299 server: &ApiServer,
300 bridge: &Arc<GatewayBridge>,
301 shutdown_tx: &tokio::sync::watch::Sender<bool>,
302 shutdown_rx: &tokio::sync::watch::Receiver<bool>,
303 listener_events_tx: &tokio::sync::mpsc::UnboundedSender<ListenerBindEvent>,
304) -> Option<SurfaceJoinHandle> {
305 let telnet_addr = super::telnet_bind_addr(&config.telnet_ip, config.telnet_port);
306 if let Some(telnet_addr) = telnet_addr {
307 let bridge_for_relogin = std::sync::Arc::clone(bridge);
311 let relogin_fn: futu_server::telnet::ReloginFn = std::sync::Arc::new(move || {
312 tracing::warn!(
313 "v1.4.97 P1-D-F: telnet relogin clearing login_cache; \
314 next P1-D tick will trigger central auth-session refresh"
315 );
316 bridge_for_relogin.caches().login_cache.clear();
317 });
318 let telnet_server = futu_server::telnet::TelnetServer::new(
319 telnet_addr.clone(),
320 std::sync::Arc::clone(server.connections()),
321 Some(bridge.subscription_runtime().manager()),
322 Some(std::sync::Arc::clone(server.metrics())),
323 shutdown_tx.clone(),
324 )
325 .with_relogin_fn(relogin_fn);
326 tracing::info!(addr = %telnet_addr, "starting Telnet server");
327 let telnet_shutdown_rx = shutdown_rx.clone();
328 let telnet_listener_events = listener_events_tx.clone();
329 Some(tokio::spawn(async move {
330 telnet_server
331 .run_until_shutdown_with_listener_events(
332 telnet_shutdown_rx,
333 Some(telnet_listener_events),
334 )
335 .await
336 }))
337 } else {
338 None
339 }
340}
341
342pub(super) fn warn_tcp_keystore_gate(
353 config: &RuntimeConfig,
354 listen_addr: &str,
355 tcp_disabled: bool,
356 any_keys_configured: bool,
357 allow_tcp_unauthenticated: bool,
358) {
359 if tcp_disabled {
360 tracing::warn!(
361 listen_addr = %listen_addr,
362 "v1.4.104 external report S-001 (P0) fix: TCP listener (port {}) NOT started — \
363 keys file configured but --allow-tcp-unauthenticated not set. \
364 native TCP FTAPI protocol has no Bearer field, cannot enforce \
365 caller-specific scope check; defaulting to fail-closed (skip TCP). \
366 Use REST/gRPC/WS endpoints for authenticated access. \
367 To restore TCP (legacy Python SDK clients) add --allow-tcp-unauthenticated, \
368 but be aware that port {} will accept ANY local connection without \
369 scope check (跨账户 leak risk).",
370 config.port,
371 config.port,
372 );
373 eprintln!(
374 "⚠️ TCP listener (port {}) DISABLED (v1.4.104 external report S-001 fix): \
375 keys file configured + no --allow-tcp-unauthenticated. \
376 Pass --allow-tcp-unauthenticated to restore (with security warning).",
377 config.port,
378 );
379 } else if any_keys_configured && allow_tcp_unauthenticated {
380 tracing::warn!(
381 listen_addr = %listen_addr,
382 "⚠️ v1.4.104: TCP listener running WITHOUT scope check despite keys configured \
383 (--allow-tcp-unauthenticated set). Port {} accepts ANY local connection — \
384 跨账户 leak risk. Use REST/gRPC/WS for authenticated clients; reserve \
385 TCP only for legacy Python SDK / C++ OpenD where Bearer not feasible.",
386 config.port,
387 );
388 eprintln!(
389 "⚠️ TCP port {} ACCEPTS UNAUTHENTICATED connections (--allow-tcp-unauthenticated). \
390 受限 keys 不在该 surface 强制. 推荐改用 REST/gRPC/WS.",
391 config.port,
392 );
393 }
394}