Skip to main content

futu_opend/startup/
phase3.rs

1//! v1.4.110 Layer 3 A: startup Phase 3 — ApiServer 构造 / handler 注册 /
2//! 推送广播器 / push dispatcher 启动. 抽自原 `mod.rs::run_daemon` 477..531 行段.
3//!
4//! Phase 3 主要副作用 (按顺序):
5//! 1. 从 `bridge.caches().login_cache` 取 user_id 构造 `ServerConfig`
6//! 2. `ApiServer::new` + `set_metrics` + `set_subscriptions`
7//! 3. `install_prometheus_extension` (server.metrics() 暴露到 `/metrics`)
8//! 4. 调 3 个域 register fn (qot / trd / sys) 注册业务 handler
9//! 5. 创建 `WsBroadcaster` + `GrpcPushBroadcaster` (各容量 1024)
10//! 6. 如 push_receiver Some → `bridge.start_push_dispatcher` (ws + grpc sinks)
11
12use std::sync::Arc;
13
14use futu_gateway_core::bridge::GatewayBridge;
15use futu_server::listener::{ApiServer, ServerConfig};
16
17use crate::config::RuntimeConfig;
18
19/// Phase 3 output — Phase 4 spawn 各 surface server 时需要这些资源.
20pub(super) struct Phase3Out {
21    pub(super) server: ApiServer,
22    pub(super) server_config: ServerConfig,
23    pub(super) ws_broadcaster: Arc<futu_rest::ws::WsBroadcaster>,
24    pub(super) grpc_broadcaster: Arc<futu_grpc::server::GrpcPushBroadcaster>,
25}
26
27pub(super) fn run_phase3(
28    config: &RuntimeConfig,
29    bridge: &Arc<GatewayBridge>,
30    listen_addr: &str,
31    shutdown_tx: tokio::sync::watch::Sender<bool>,
32) -> Phase3Out {
33    // 3. 创建 API 服务端
34    let server_config = ServerConfig {
35        listen_addr: listen_addr.to_string(),
36        server_ver: 1000,
37        // Runtime InitConnect reads the generation-bound IdentitySnapshot from
38        // RequestRouter; production no longer freezes startup uid here.
39        login_user_id: 0,
40        keepalive_interval: 10,
41        rsa_private_key: config.rsa_private_key.clone(),
42    };
43    if server_config.rsa_private_key.is_some() {
44        tracing::info!("RSA encryption enabled for InitConnect");
45    }
46    let mut server = ApiServer::new(server_config.clone());
47    server
48        .router()
49        .set_startup_readiness(bridge.startup_readiness().clone());
50    server.set_server_time_store(bridge.server_clock().anchor_store());
51    server.set_metrics(std::sync::Arc::clone(bridge.push_runtime().metrics()));
52    server.set_subscriptions(bridge.subscription_runtime().manager());
53
54    // v1.4.90 P1-B: 把 GatewayMetrics 注册为 futu_auth Registry 的 extension
55    // renderer, 让 `/metrics` HTTP 端点暴露 per-cmd / per-hour push counter.
56    // 之前 v1.4.83/84 只在 telnet `show_metrics` 暴露, /metrics 仅 3 个 auth
57    // counter (cmd_6212_quote / cmd_14716_trade_new 等被双 tester 报告 missing).
58    futu_server::metrics::install_prometheus_extension(std::sync::Arc::clone(server.metrics()));
59    // v1.4.113 optimization #4: PushHealth used to be visible only through
60    // `/api/push-subscriber-info`; register it as a Prometheus extension too
61    // so F2/F3/F4/F6 health signals can alert without polling JSON.
62    let bridge_for_push_health_metrics = std::sync::Arc::clone(bridge);
63    futu_auth::metrics::register_global_renderer(move || {
64        bridge_for_push_health_metrics
65            .push_runtime()
66            .push_health()
67            .render_prometheus()
68    });
69    // v1.4.113 optimization #3: expose per-broker TCP reconnect health on
70    // `/metrics`, using the same BrokerRuntime snapshot as `/api/admin/status`.
71    let bridge_for_broker_tcp_metrics = std::sync::Arc::clone(bridge);
72    futu_auth::metrics::register_global_renderer(move || {
73        bridge_for_broker_tcp_metrics
74            .broker_runtime()
75            .render_prometheus()
76    });
77
78    // 4. 注册业务处理器
79    //
80    // v1.4.110 P0-5 T15 P4 wire: 显式调 3 个域 register fn (替原
81    // `bridge.register_handlers(&server)` 反向调用). 3 crate split 后各域
82    // register fn 在各自 crate:
83    {
84        let router = server.router();
85        futu_gateway_qot::register_handlers(router, bridge);
86        futu_gateway_trd::register_handlers(router, bridge);
87        futu_gateway_core::handlers_sys::register_handlers_with_shutdown_and_user_info(
88            router,
89            bridge,
90            shutdown_tx,
91            futu_gateway_core::handlers_sys::UserInfoRuntimeConfig {
92                update_check_url: Some(config.update_check_url.clone()),
93                is_nn: config.platform == crate::cli::Platform::Futunn,
94            },
95        );
96        tracing::info!("all business handlers registered");
97    }
98
99    // 5. 创建推送广播器 (REST WebSocket + gRPC)
100    let ws_broadcaster = std::sync::Arc::new(futu_rest::ws::WsBroadcaster::new(1024));
101    let grpc_broadcaster = std::sync::Arc::new(futu_grpc::server::GrpcPushBroadcaster::new(1024));
102
103    Phase3Out {
104        server,
105        server_config,
106        ws_broadcaster,
107        grpc_broadcaster,
108    }
109}
110
111#[cfg(test)]
112mod tests;