Skip to main content

futu_opend/startup/
phase4.rs

1//! v1.4.110 Layer 3 A: startup Phase 4 — surface server (WS/REST/gRPC/Telnet)
2//! spawn / card_num expand / SIGHUP unified handler / TCP fail-closed gate /
3//! 主 `tokio::select!` 循环 + 清理. 抽自原 `mod.rs::run_daemon` 533..1075 行段.
4//!
5//! Phase 4 取走 `Phase1Out` 所有权 (保证 `_audit_guard` 直到 fn 返回才 drop),
6//! 同时取 `bridge: Arc<GatewayBridge>` 与 `Phase3Out` 启动各 surface server.
7
8use anyhow::Result;
9use std::sync::Arc;
10
11use futu_gateway_core::bridge::GatewayBridge;
12use futu_server::listener_status::ListenerBindEvent;
13use futu_server::ws_listener::{WsServer, WsServerDeps};
14
15use crate::config::RuntimeConfig;
16use crate::startup::phase1::Phase1Out;
17use crate::startup::phase2::{AuthPlan, execute_auth_plan};
18use crate::startup::phase3::Phase3Out;
19
20mod listener_readiness;
21mod network_exposure;
22mod shutdown;
23mod summary;
24#[cfg(test)]
25mod tests;
26
27use listener_readiness::{
28    ActivatedListenerReadiness, EnabledListenerSurfaces, ListenerReadiness, SurfaceTasks,
29    await_listeners_opened, wrap_grpc_server_result,
30};
31#[cfg(test)]
32use network_exposure::legacy_network_exposure_warnings;
33use network_exposure::warn_legacy_network_exposure;
34use shutdown::{
35    BackgroundJoinHandle, SHUTDOWN_GRACE_PERIOD, SurfaceJoinHandle, await_background_shutdown,
36    await_surface_shutdown, log_surface_task_result, request_surface_shutdown,
37    surface_task_result_or_pending,
38};
39#[cfg(test)]
40use summary::http_host_for_bind;
41use summary::{RestTransportKind, log_startup_summary, merge_push_health_snapshots_for_rest};
42
43include!("phase4/preauth.rs");
44
45pub(super) async fn run_phase4(
46    config: &RuntimeConfig,
47    phase1: Phase1Out,
48    bridge: Arc<GatewayBridge>,
49    mut auth_plan: Option<AuthPlan>,
50    phase3: Phase3Out,
51    shutdown_tx: tokio::sync::watch::Sender<bool>,
52    mut shutdown_rx: tokio::sync::watch::Receiver<bool>,
53) -> Result<()> {
54    let Phase1Out {
55        _audit_guard,
56        shared_counters,
57        listen_addr,
58        rest_keys_file,
59        rest_key_store,
60        rest_transport,
61        ws_keys_file,
62        ws_key_store,
63        grpc_keys_file,
64        grpc_key_store,
65        allow_tcp_unauthenticated,
66    } = phase1;
67    let rest_transport_kind = match config.rest_port {
68        None => RestTransportKind::Disabled,
69        Some(_) if rest_transport.is_tls() => RestTransportKind::Tls,
70        Some(_) => RestTransportKind::Plaintext,
71    };
72    let Phase3Out {
73        server,
74        server_config,
75        ws_broadcaster,
76        grpc_broadcaster,
77    } = phase3;
78    let push_sinks: Vec<std::sync::Arc<dyn futu_server::push::ExternalPushSink>> = vec![
79        std::sync::Arc::clone(&ws_broadcaster) as _,
80        std::sync::Arc::clone(&grpc_broadcaster) as _,
81    ];
82    let prepared_push_dispatcher = bridge.create_push_dispatcher(&server, push_sinks);
83    let (listener_events_tx, mut listener_events_rx) =
84        tokio::sync::mpsc::unbounded_channel::<ListenerBindEvent>();
85
86    // v1.4.103 (B10): 把 3 个 server 的 key_store Option 提升到外层作用域,
87    // 让后面的 expand_allowed_card_nums spawn task 能跨 server 块捕获.
88    // 各 server 块仍各自加载 (允许独立 keys.json), 这里只共享 Arc clone.
89    let mut ws_key_store_holder: Option<std::sync::Arc<futu_auth::KeyStore>> = None;
90    let mut rest_key_store_holder: Option<std::sync::Arc<futu_auth::KeyStore>> = None;
91    let mut grpc_key_store_holder: Option<std::sync::Arc<futu_auth::KeyStore>> = None;
92    let any_keys_configured =
93        rest_keys_file.is_some() || grpc_keys_file.is_some() || ws_keys_file.is_some();
94    warn_legacy_network_exposure(config, any_keys_configured);
95
96    // 7. 启动 WebSocket 服务(可选)
97    let mut ws_handle = if let Some(ws_port) = config.websocket_port {
98        let ws_addr = format!("{}:{}", config.ip, ws_port);
99        // v1.0:WS 握手鉴权 —— 复用 REST 的 key store 设计(`--ws-keys-file` 独立
100        // 指定,不指定时 legacy 放行)
101        // v1.4.102 BUG-007 fix (P1, leaf v1.4.100 报告): keys-file load 失败
102        // 必须 fail-closed (abort daemon), 不再 silent fallback to legacy mode.
103        // 历史: 用户传 `--ws-keys-file` 表示明确意图启用 auth, 文件 typo / parse
104        // 错时 daemon 只 log error 然后继续无 auth → 用户以为有门锁, 实际门锁
105        // 没装上 (security misconfig 比纯 legacy 更危险).
106        let ws_key_store = ws_key_store.clone();
107        let ws_counters = std::sync::Arc::clone(&shared_counters);
108        // v1.4.103 (B10): clone Arc 给外层 holder 持有, 同时让 ws_server 仍接管原 Arc.
109        ws_key_store_holder = ws_key_store.as_ref().map(std::sync::Arc::clone);
110        let ws_server = WsServer::with_auth(
111            ws_addr.clone(),
112            server_config.clone(),
113            WsServerDeps::new(
114                std::sync::Arc::clone(server.connections()),
115                std::sync::Arc::clone(server.router()),
116                Some(bridge.subscription_runtime().manager()),
117            ),
118            ws_key_store,
119            Some(ws_counters),
120        )
121        .with_server_time_store(bridge.server_clock().anchor_store());
122        tracing::info!(addr = %ws_addr, "starting WebSocket server");
123        let ws_shutdown_rx = shutdown_rx.clone();
124        let ws_listener_events = listener_events_tx.clone();
125        Some(tokio::spawn(async move {
126            ws_server
127                .run_until_shutdown_with_listener_events(ws_shutdown_rx, Some(ws_listener_events))
128                .await
129        }))
130    } else {
131        None
132    };
133
134    // 8. 启动 REST API 服务(可选,含 WebSocket 推送)
135    let mut rest_handle = if let Some(rest_port) = config.rest_port {
136        let rest_addr = format!("{}:{}", config.ip, rest_port);
137        let router = std::sync::Arc::clone(server.router());
138        let broadcaster = std::sync::Arc::clone(&ws_broadcaster);
139        // v1.4.102 BUG-007 fix: 同 WS 路径 (fail-closed when keys-file 显式指定).
140        let rest_key_store = rest_key_store
141            .clone()
142            .unwrap_or_else(|| std::sync::Arc::new(futu_auth::KeyStore::empty()));
143        let rest_transport_name = match rest_transport_kind {
144            RestTransportKind::Tls => "https",
145            RestTransportKind::Plaintext => "http",
146            RestTransportKind::Disabled => unreachable!("REST port is configured"),
147        };
148        tracing::info!(
149            addr = %rest_addr,
150            transport = rest_transport_name,
151            "starting REST API server (WebSocket: /ws)"
152        );
153
154        // v1.4.103 (B10): clone Arc 给外层 holder 持有 (跨 server 块共享给
155        // expand_allowed_card_nums spawn task). 仅在 keys 实际配置时填充
156        // (empty store 不需要 expand).
157        if rest_key_store.is_configured() {
158            rest_key_store_holder = Some(std::sync::Arc::clone(&rest_key_store));
159        }
160
161        // v1.4.103 codex F3.1 (P1) round 3: REST 单独的 SIGHUP reload listener
162        // 已**移除** — 防与 card_num expand 间的 race (多 listener 并发顺序乱
163        // 可能让 reload 覆盖 expand 写入的 sentinel, 受限 key 在 reload 窗口
164        // 内 silent unrestricted). 现在 reload + expand 由文件末尾的 unified
165        // SIGHUP handler 顺序操作 (调 card_num_reload_and_expand_fn(true)).
166
167        let rest_counters = std::sync::Arc::clone(&shared_counters);
168        // v1.4.32+ admin snapshot provider:closure 捕获 `Arc<GatewayBridge>`
169        // 每次被 admin_status handler 调用时返回实时 StatusSnapshot JSON。
170        // bridge 已经在上面 Arc 化,这里只做一次 clone。
171        let bridge_for_status = std::sync::Arc::clone(&bridge);
172        let admin_status_provider: futu_rest::adapter::AdminStatusProvider =
173            std::sync::Arc::new(move || {
174                serde_json::to_value(bridge_for_status.snapshot_status())
175                    .unwrap_or_else(|_| serde_json::json!({"error": "snapshot serialize failed"}))
176            });
177        let rest_shutdown_tx = shutdown_tx.clone();
178        let admin_shutdown_handler: futu_rest::adapter::AdminShutdownHandler =
179            std::sync::Arc::new(move || {
180                rest_shutdown_tx
181                    .send(true)
182                    .map_err(|e| format!("shutdown receiver dropped: {e}"))
183            });
184        // v1.4.32+ admin reload handler:closure 调 Bridge::reload() 清 cipher cache
185        // v1.4.34: reload 升级为 async(内部刷 credentials 走网络 I/O),
186        // handler 返 Future,axum admin_reload async handler await 之
187        // v1.4.106 codex 0554 F3 [P2]: reload 拆两阶段后变 sync (sync clear +
188        // tokio::spawn refresh). closure 仍返 Future 类型不变 (向后兼容
189        // AdminReloadHandler API), 但内部不 await — sync 阶段已完成 + spawn
190        // 已派发, ReloadReport 立即可用.
191        let bridge_for_reload = std::sync::Arc::clone(&bridge);
192        let admin_reload_handler: futu_rest::adapter::AdminReloadHandler =
193            std::sync::Arc::new(move || {
194                let bridge = std::sync::Arc::clone(&bridge_for_reload);
195                Box::pin(async move {
196                    serde_json::to_value(bridge.reload())
197                        .unwrap_or_else(|_| serde_json::json!({"error": "reload serialize failed"}))
198                })
199            });
200        // v1.4.83 §9 Phase 2 F5: push health snapshot provider —— closure
201        // 捕获 bridge.push_runtime().push_health Arc, 每次
202        // `/api/push-subscriber-info` 调用时返当前真实 snapshot.
203        // v1.4.91 P1-D wiring: closure 同时捕获 bridge.push_runtime().qot_login_health,
204        // 在返 push_health 同时多带一个 qot_login_health 字段供 ops 看
205        // qot_logined self-heal counter (修 P1-D non-deterministic gap).
206        let bridge_for_push_health = std::sync::Arc::clone(&bridge);
207        let push_health_snapshot_provider: futu_rest::adapter::PushHealthSnapshotProvider =
208            std::sync::Arc::new(move || {
209                let push = serde_json::to_value(
210                    bridge_for_push_health
211                        .push_runtime()
212                        .push_health()
213                        .snapshot(),
214                )
215                .unwrap_or_else(
216                    |_| serde_json::json!({"error": "push_health snapshot serialize failed"}),
217                );
218                let qot_login = serde_json::to_value(
219                    bridge_for_push_health
220                        .push_runtime()
221                        .qot_login_health()
222                        .snapshot(),
223                )
224                .unwrap_or_else(
225                    |_| serde_json::json!({"error": "qot_login_health snapshot serialize failed"}),
226                );
227                let backend_connected =
228                    bridge_for_push_health.broker_runtime().platform_connected();
229                let stock_list_status = bridge_for_push_health
230                    .caches()
231                    .static_cache
232                    .stock_list_sync_status();
233                let stock_list_static_query_ready = matches!(
234                    futu_domain_static_data::stock_list_static_query_gate_action_for_status(
235                        &stock_list_status.domain_facts(),
236                    ),
237                    futu_domain_static_data::StockListStaticQueryGateAction::Allow
238                );
239                merge_push_health_snapshots_for_rest(
240                    push,
241                    qot_login,
242                    backend_connected,
243                    stock_list_static_query_ready,
244                )
245            });
246        // v1.4.105 D12 (Phase 2): 注入 card_num resolver 让 REST trade
247        // handler (place_order / modify_order / cancel_all_order) 能解析
248        // user 传 `card_num` 字段 → acc_id (覆盖 c2s.header.acc_id).
249        // closure 捕获 bridge.caches().trd_cache, 调
250        // `find_acc_ids_by_card_num(input) -> Vec<u64>`.
251        let bridge_for_card_num = std::sync::Arc::clone(&bridge);
252        let card_num_resolver: futu_rest::adapter::CardNumResolver =
253            std::sync::Arc::new(move |cn: &str| {
254                bridge_for_card_num
255                    .caches()
256                    .trd_cache
257                    .find_acc_ids_by_card_num(cn)
258            });
259        let rest_shutdown_rx = shutdown_rx.clone();
260        let rest_listener_events = listener_events_tx.clone();
261        Some(tokio::spawn(async move {
262            futu_rest::server::start_with_auth_full_admin_until_shutdown_with_transport_and_listener_events(
263                &rest_addr,
264                router,
265                broadcaster,
266                rest_key_store,
267                rest_counters,
268                futu_rest::server::RestAdminHooks {
269                    admin_status_provider: Some(admin_status_provider),
270                    admin_shutdown_handler: Some(admin_shutdown_handler),
271                    admin_reload_handler: Some(admin_reload_handler),
272                    push_health_snapshot_provider: Some(push_health_snapshot_provider),
273                    card_num_resolver: Some(card_num_resolver),
274                },
275                rest_transport,
276                rest_shutdown_rx,
277                Some(rest_listener_events),
278            )
279            .await
280            .map_err(anyhow::Error::from)
281        }))
282    } else {
283        None
284    };
285
286    // 9. 启动 gRPC 服务(可选,含流式推送)
287    let mut grpc_handle = if let Some(grpc_port) = config.grpc_port {
288        let grpc_addr = format!("{}:{}", config.ip, grpc_port);
289        let router = std::sync::Arc::clone(server.router());
290        let broadcaster = std::sync::Arc::clone(&grpc_broadcaster);
291        // v1.4.102 BUG-007 fix: 同 WS / REST 路径 (fail-closed when keys-file 显式指定).
292        let grpc_key_store = grpc_key_store
293            .clone()
294            .unwrap_or_else(|| std::sync::Arc::new(futu_auth::KeyStore::empty()));
295        tracing::info!(addr = %grpc_addr, "starting gRPC server (SubscribePush: streaming)");
296
297        // v1.4.103 (B10): clone Arc 给外层 holder 持有.
298        if grpc_key_store.is_configured() {
299            grpc_key_store_holder = Some(std::sync::Arc::clone(&grpc_key_store));
300        }
301
302        // v1.4.103 codex F3.1 (P1) round 3: gRPC 单独的 SIGHUP reload listener
303        // 已**移除** — 同 REST 块, 由文件末尾的 unified SIGHUP handler 顺序
304        // reload + expand 防 race.
305
306        let grpc_counters = std::sync::Arc::clone(&shared_counters);
307        let grpc_shutdown_rx = shutdown_rx.clone();
308        let grpc_listener_events = listener_events_tx.clone();
309        Some(tokio::spawn(async move {
310            let result = futu_grpc::server::start_with_auth_until_shutdown_with_listener_events(
311                &grpc_addr,
312                router,
313                broadcaster,
314                grpc_key_store,
315                grpc_counters,
316                grpc_shutdown_rx,
317                Some(grpc_listener_events),
318            )
319            .await;
320            wrap_grpc_server_result(result)
321        }))
322    } else {
323        None
324    };
325
326    // 10. 启动 Telnet 管理服务(可选)
327    let telnet_addr = telnet_bind_addr(&config.telnet_ip, config.telnet_port);
328    let mut telnet_handle = if let Some(telnet_addr) = telnet_addr {
329        // v1.4.97 P1-D-F: relogin callback closes over bridge to clear
330        // login_cache; next 30s P1-D health tick triggers relogin.
331        // Aligns with C++ GTWCmd_ReLogin (FTGateway/FTGTW_Define_Key.h:5).
332        let bridge_for_relogin = std::sync::Arc::clone(&bridge);
333        let relogin_fn: futu_server::telnet::ReloginFn = std::sync::Arc::new(move || {
334            tracing::warn!(
335                "v1.4.97 P1-D-F: telnet relogin clearing login_cache; \
336                     next P1-D tick will trigger central auth-session refresh"
337            );
338            bridge_for_relogin.caches().login_cache.clear();
339        });
340        let telnet_server = futu_server::telnet::TelnetServer::new(
341            telnet_addr.clone(),
342            std::sync::Arc::clone(server.connections()),
343            Some(bridge.subscription_runtime().manager()),
344            Some(std::sync::Arc::clone(server.metrics())),
345            shutdown_tx.clone(),
346        )
347        .with_relogin_fn(relogin_fn);
348        tracing::info!(addr = %telnet_addr, "starting Telnet server");
349        let telnet_shutdown_rx = shutdown_rx.clone();
350        let telnet_listener_events = listener_events_tx.clone();
351        Some(tokio::spawn(async move {
352            telnet_server
353                .run_until_shutdown_with_listener_events(
354                    telnet_shutdown_rx,
355                    Some(telnet_listener_events),
356                )
357                .await
358        }))
359    } else {
360        None
361    };
362
363    // v1.4.103 (B10) + codex F2 (P1): 立即跑首次 expand + background retry
364    // + SIGHUP reload 钩子 (fail-closed sentinel 即时生效, 无 startup window).
365    //
366    // 行为:
367    // 1. **立即跑首次 expand** — 即便 cache 空, fail-closed sentinel (codex F1)
368    //    也写进 allowed_acc_ids, 短路 startup window 的 silent unrestricted.
369    // 2. **background retry**: 每 10s 检查 cache, 加载后 re-expand 真 acc_id
370    //    覆盖 sentinel. 60s 上限.
371    // 3. **SIGHUP reload 钩子**: 重载 keys.json 后再 expand 一次 (codex F2 reload window).
372    // v1.4.103 codex F3.1 (P1) round 3: 合并 reload + expand 为单一 ordered op,
373    // 防 race. 当 SIGHUP 触发时, 此 fn 同步顺序: (1) reload 所有 store (从
374    // keys.json 读 raw allowed_card_nums + ArcSwap.store) (2) expand (resolve
375    // card_num → acc_ids + ArcSwap.store with sentinel/resolved). 因为是单线
376    // 调用, 不存在多 SIGHUP listener 间的乱序 race.
377    //
378    // `do_reload` 控制是否先 reload (启动时第一次 / 60s retry 不需 reload,
379    // 直接 expand; SIGHUP 路径需要先 reload).
380    let card_num_reload_and_expand_fn: std::sync::Arc<dyn Fn(bool) + Send + Sync> = {
381        let bridge_for_expand = std::sync::Arc::clone(&bridge);
382        let ws_ks = ws_key_store_holder.clone();
383        let rest_ks = rest_key_store_holder.clone();
384        let grpc_ks = grpc_key_store_holder.clone();
385        std::sync::Arc::new(move |do_reload: bool| {
386            // (1) Reload phase — 仅 SIGHUP 路径调
387            if do_reload {
388                for (ks_name, ks_opt) in [("ws", &ws_ks), ("rest", &rest_ks), ("grpc", &grpc_ks)] {
389                    let Some(ks) = ks_opt.as_ref() else { continue };
390                    match ks.reload() {
391                        Ok(()) => tracing::warn!(
392                            ks = ks_name,
393                            keys_loaded = ks.len(),
394                            "v1.4.103 F3.1: keys reloaded on SIGHUP (before card_num expand)"
395                        ),
396                        Err(e) => tracing::error!(
397                            ks = ks_name,
398                            error = %e,
399                            "v1.4.103 F3.1: keys reload failed (skipping expand for this store)"
400                        ),
401                    }
402                }
403            }
404            // (2) Expand phase
405            let trd_cache = std::sync::Arc::clone(&bridge_for_expand.caches().trd_cache);
406            let resolver = {
407                let cache_clone = std::sync::Arc::clone(&trd_cache);
408                move |cn: &str| cache_clone.find_acc_ids_by_card_num(cn)
409            };
410            for (ks_name, ks_opt) in [("ws", &ws_ks), ("rest", &rest_ks), ("grpc", &grpc_ks)] {
411                let Some(ks) = ks_opt.as_ref() else { continue };
412                let (resolved, unresolved, ambiguous) = ks.expand_allowed_card_nums(
413                    &resolver,
414                    |key_id, cn| {
415                        tracing::warn!(
416                            key_id = %key_id,
417                            card_num = %cn,
418                            "v1.4.103 B10/F1 fail-closed: card_num not found in trd_cache; \
419                             writing sentinel acc_id=0 to enforce restrictive denylist \
420                             (limits.contains check 永远 false → reject 真账户)"
421                        );
422                    },
423                    |key_id, cn, candidates| {
424                        tracing::warn!(
425                            key_id = %key_id,
426                            card_num = %cn,
427                            candidates = ?candidates,
428                            "v1.4.103 B10/F1 fail-closed: ambiguous card_num suffix \
429                             matched multiple accounts (skipped, write 完整 16 位 / specific 4 位)"
430                        );
431                    },
432                );
433                tracing::info!(
434                    ks = ks_name,
435                    resolved,
436                    unresolved,
437                    ambiguous,
438                    "v1.4.103 B10: expanded allowed_card_nums into allowed_acc_ids"
439                );
440            }
441        })
442    };
443    // 老 alias (不 reload, 仅 expand) 用于启动 + retry 路径
444    let card_num_expand_fn: std::sync::Arc<dyn Fn() + Send + Sync> = {
445        let inner = std::sync::Arc::clone(&card_num_reload_and_expand_fn);
446        std::sync::Arc::new(move || (inner)(false))
447    };
448
449    // codex F2 (P1) 立即跑首次: fail-closed sentinel 即时生效, 不留 startup window
450    (card_num_expand_fn)();
451
452    // background retry loop: cache 加载后 re-expand 覆盖 sentinel
453    let card_num_retry_handle: BackgroundJoinHandle = {
454        let card_num_expand_fn_loop = std::sync::Arc::clone(&card_num_expand_fn);
455        let bridge_for_check = std::sync::Arc::clone(&bridge);
456        let mut card_num_retry_shutdown_rx = shutdown_rx.clone();
457        tokio::spawn(async move {
458            let trd_cache = std::sync::Arc::clone(&bridge_for_check.caches().trd_cache);
459            let mut attempts = 0u32;
460            let max_attempts = 6u32; // 6 × 10s = 60s
461            loop {
462                tokio::select! {
463                    changed = card_num_retry_shutdown_rx.changed() => {
464                        if changed.is_err() || *card_num_retry_shutdown_rx.borrow() {
465                            tracing::debug!(
466                                "v1.4.111: card_num retry loop received shutdown signal"
467                            );
468                            return;
469                        }
470                    }
471                    _ = tokio::time::sleep(std::time::Duration::from_secs(10)) => {}
472                }
473                attempts += 1;
474                let accounts = trd_cache.get_accounts();
475                if accounts.is_empty() {
476                    if attempts >= max_attempts {
477                        tracing::warn!(
478                            "v1.4.103 B10: trd_cache 仍空 (after {max_attempts} × 10s); \
479                             受限 key 仍走 fail-closed sentinel reject 直到 SIGHUP / cache 加载."
480                        );
481                        return;
482                    }
483                    continue;
484                }
485                (card_num_expand_fn_loop)();
486                return;
487            }
488        })
489    };
490
491    // codex F2 (P1) + F3.1 (P1) round 3: 单一 SIGHUP 钩子 — reload + expand
492    // 顺序操作避免 race. 之前每 server (REST/gRPC) 各自 SIGHUP 监听 reload,
493    // 加上 card_num expand 单独监听 — 多 SIGHUP listener 并发执行无序, 可能
494    // expand 先跑 (写 sentinel) 然后 reload 后跑 (overwrite sentinel 用 raw
495    // allowed_card_nums) → 受限 key 在 reload window 内 silent unrestricted.
496    //
497    // 现在: **唯一 SIGHUP listener**, 顺序 reload → expand. 删除 REST/gRPC 各
498    // 自的 SIGHUP reload listener (上面 server 块内已注释).
499    #[cfg(unix)]
500    let sighup_handle: Option<BackgroundJoinHandle> = {
501        let unified_sighup_fn = std::sync::Arc::clone(&card_num_reload_and_expand_fn);
502        let mut sighup_shutdown_rx = shutdown_rx.clone();
503        Some(tokio::spawn(async move {
504            use tokio::signal::unix::{SignalKind, signal};
505            let mut sig = match signal(SignalKind::hangup()) {
506                Ok(s) => s,
507                Err(e) => {
508                    tracing::error!(error = %e, "SIGHUP install failed (unified reload+expand)");
509                    return;
510                }
511            };
512            tracing::info!(
513                "v1.4.103 F3.1: unified SIGHUP handler installed (reload all keys + expand card_num)"
514            );
515            loop {
516                tokio::select! {
517                    signal = sig.recv() => {
518                        if signal.is_none() {
519                            return;
520                        }
521                        tracing::info!(
522                            "v1.4.103 F3.1: SIGHUP received — running reload_all_stores + \
523                             expand_allowed_card_nums (single ordered op, no race)"
524                        );
525                        (unified_sighup_fn)(true); // do_reload=true
526                    }
527                    changed = sighup_shutdown_rx.changed() => {
528                        if changed.is_err() || *sighup_shutdown_rx.borrow() {
529                            tracing::debug!(
530                                "v1.4.111: unified SIGHUP handler received shutdown signal"
531                            );
532                            return;
533                        }
534                    }
535                }
536            }
537        }))
538    };
539    #[cfg(not(unix))]
540    let sighup_handle: Option<BackgroundJoinHandle> = None;
541
542    // 11. v1.4.104 external reviewer S-001 (P0): native TCP keystore guard.
543    //
544    // 配任一 keys file (rest / grpc / ws) → 用户**意图**启用 scope mode.
545    // 但 TCP FTAPI 协议 InitConnect 没有 Bearer 字段, 无法做 caller-specific
546    // scope check (S-001). 默认 fail-closed: 不启 TCP listener, 用户需要
547    // 通过 REST/gRPC/WS 与 daemon 交互 (这些 surface 都已加 pipeline auth).
548    //
549    // 显式 `--allow-tcp-unauthenticated` opt-in → 启 TCP, 但 daemon 启动
550    // loud warn 该端口完全无 auth, agent skill 等 local process 可任意调用.
551    // 用之前 captured 的 local 变量, args 已被 merge_config 消费
552    let tcp_disabled = !super::phase1::ftapi_listener_enabled(config);
553
554    if tcp_disabled {
555        tracing::warn!(
556            listen_addr = %listen_addr,
557            "v1.4.104 external report S-001 (P0) fix: TCP listener (port {}) NOT started — \
558             keys file configured but --allow-tcp-unauthenticated not set. \
559             native TCP FTAPI protocol has no Bearer field, cannot enforce \
560             caller-specific scope check; defaulting to fail-closed (skip TCP). \
561             Use REST/gRPC/WS endpoints for authenticated access. \
562             To restore TCP (legacy Python SDK clients) add --allow-tcp-unauthenticated, \
563             but be aware that port {} will accept ANY local connection without \
564             scope check (跨账户 leak risk).",
565            config.port,
566            config.port,
567        );
568        eprintln!(
569            "⚠️  TCP listener (port {}) DISABLED (v1.4.104 external report S-001 fix): \
570             keys file configured + no --allow-tcp-unauthenticated. \
571             Pass --allow-tcp-unauthenticated to restore (with security warning).",
572            config.port,
573        );
574    } else if any_keys_configured && allow_tcp_unauthenticated {
575        tracing::warn!(
576            listen_addr = %listen_addr,
577            "⚠️  v1.4.104: TCP listener running WITHOUT scope check despite keys configured \
578             (--allow-tcp-unauthenticated set). Port {} accepts ANY local connection — \
579             跨账户 leak risk. Use REST/gRPC/WS for authenticated clients; reserve \
580             TCP only for legacy Python SDK / C++ OpenD where Bearer not feasible.",
581            config.port,
582        );
583        eprintln!(
584            "⚠️  TCP port {} ACCEPTS UNAUTHENTICATED connections (--allow-tcp-unauthenticated). \
585             受限 keys 不在该 surface 强制. 推荐改用 REST/gRPC/WS.",
586            config.port,
587        );
588    }
589
590    let tcp_shutdown_rx = shutdown_rx.clone();
591    let tcp_handle: Option<SurfaceJoinHandle> = if tcp_disabled {
592        None
593    } else {
594        let tcp_listener_events = listener_events_tx.clone();
595        Some(tokio::spawn(async move {
596            server
597                .run_until_shutdown_with_listener_events(tcp_shutdown_rx, Some(tcp_listener_events))
598                .await
599        }))
600    };
601
602    let enabled_listeners = EnabledListenerSurfaces {
603        ftapi: !tcp_disabled,
604        websocket: config.websocket_port.is_some(),
605        rest: config.rest_port.is_some(),
606        grpc: config.grpc_port.is_some(),
607        telnet: config.telnet_port.is_some(),
608    };
609    if let Some(plan) = auth_plan.as_mut() {
610        let loopback = bind_includes_loopback(&config.ip);
611        let available = external_verification_owner_available(VerificationOwnerFacts {
612            ftapi_loopback_legacy: enabled_listeners.ftapi && loopback,
613            rest_loopback_legacy: enabled_listeners.rest
614                && loopback
615                && rest_key_store
616                    .as_ref()
617                    .is_none_or(|store| !store.is_configured()),
618            ws_loopback_legacy: enabled_listeners.websocket
619                && loopback
620                && ws_key_store
621                    .as_ref()
622                    .is_none_or(|store| !store.is_configured()),
623            grpc_loopback_legacy: enabled_listeners.grpc
624                && loopback
625                && grpc_key_store
626                    .as_ref()
627                    .is_none_or(|store| !store.is_configured()),
628            rest_auth_setup: enabled_listeners.rest
629                && rest_key_store
630                    .as_ref()
631                    .is_some_and(|store| store.has_usable_scope(futu_auth::Scope::AuthSetup)),
632            ws_auth_setup: enabled_listeners.websocket
633                && ws_key_store
634                    .as_ref()
635                    .is_some_and(|store| store.has_usable_scope(futu_auth::Scope::AuthSetup)),
636            grpc_auth_setup: enabled_listeners.grpc
637                && grpc_key_store
638                    .as_ref()
639                    .is_some_and(|store| store.has_usable_scope(futu_auth::Scope::AuthSetup)),
640        });
641        plan.set_external_verification_available(available);
642        tracing::info!(
643            external_verification_owner_available = available,
644            "computed pre-login verification ownership from bound listeners"
645        );
646    }
647    let mut surface_tasks = SurfaceTasks {
648        ftapi: tcp_handle,
649        websocket: ws_handle.take(),
650        rest: rest_handle.take(),
651        grpc: grpc_handle.take(),
652        telnet: telnet_handle.take(),
653    };
654    drop(listener_events_tx);
655    let listener_readiness = await_listeners_opened(
656        &mut listener_events_rx,
657        enabled_listeners,
658        &mut surface_tasks,
659        &mut shutdown_rx,
660    )
661    .await;
662
663    // 11. 等待任一 surface 退出、Ctrl+C 或 telnet shutdown。Bind 聚合失败
664    // 也流入同一个 cleanup 尾部,避免已成功打开的 sibling surface 被 detach。
665    let phase4_result: anyhow::Result<()> = match listener_readiness {
666        Err(error) => Err(error),
667        Ok(readiness) => {
668            // Bind success is not enough for an external Verification owner:
669            // the accept/serve loops must be active before authentication can
670            // block awaiting a pre-login 1006 request. Every surface keeps its
671            // shared readiness gate, so activation exposes only the explicit
672            // pre-login allowlist until the identity transition reaches Ready.
673            match activate_restricted_listeners_before_auth(
674                readiness,
675                enabled_listeners,
676                &mut shutdown_rx,
677            )
678            .await
679            {
680                Err(error) => Err(error),
681                Ok(listener_activation) => {
682                    let auth_error = if let Some(auth_plan) = auth_plan {
683                        match execute_auth_plan_until_shutdown(
684                            Arc::clone(&bridge),
685                            auth_plan,
686                            &mut shutdown_rx,
687                            &mut surface_tasks,
688                        )
689                        .await
690                        {
691                            Ok(push_rx) => {
692                                bridge.start_prepared_push_dispatcher(
693                                    prepared_push_dispatcher,
694                                    push_rx,
695                                );
696                                tracing::info!(
697                                    "authentication promoted atomically; push dispatcher started"
698                                );
699                                None
700                            }
701                            Err(error) => {
702                                let startup = bridge.startup_readiness().snapshot();
703                                if startup.state
704                                    == futu_server::identity::StartupState::Authenticating
705                                {
706                                    let _ = bridge.startup_readiness().transition(
707                                        startup.generation,
708                                        futu_server::identity::StartupEvent::AuthFailed,
709                                    );
710                                }
711                                crate::hints::print_auth_error_hint(&error, config);
712                                Some(startup_auth_failure_after_restricted_bind(error))
713                            }
714                        }
715                    } else {
716                        None
717                    };
718                    if let Some(error) = auth_error {
719                        Err(error)
720                    } else {
721                        match listener_activation
722                            .activate_ready_only_and_confirm(&mut shutdown_rx, |marker| *marker)
723                            .await
724                        {
725                            Err(error) => Err(error),
726                            Ok(listeners_marker) => {
727                                tracing::info!(marker = %listeners_marker, "authentication complete on active listeners");
728                                log_startup_summary(
729                                    config,
730                                    &listen_addr,
731                                    &bridge,
732                                    rest_transport_kind,
733                                );
734                                tracing::info!(
735                                    "gateway ready, accepting connections on {listen_addr}"
736                                );
737                                tracing::info!("press Ctrl+C to exit");
738                                tokio::select! {
739                                result = surface_task_result_or_pending(&mut surface_tasks.ftapi) => {
740                                    surface_tasks.ftapi = None;
741                                    log_surface_task_result("API server", result)
742                                }
743                                result = surface_task_result_or_pending(&mut surface_tasks.websocket) => {
744                                    surface_tasks.websocket = None;
745                                    log_surface_task_result("WebSocket", result)
746                                }
747                                result = surface_task_result_or_pending(&mut surface_tasks.rest) => {
748                                    surface_tasks.rest = None;
749                                    log_surface_task_result("REST API", result)
750                                }
751                                result = surface_task_result_or_pending(&mut surface_tasks.grpc) => {
752                                    surface_tasks.grpc = None;
753                                    log_surface_task_result("gRPC", result)
754                                }
755                                result = surface_task_result_or_pending(&mut surface_tasks.telnet) => {
756                                    surface_tasks.telnet = None;
757                                    log_surface_task_result("Telnet", result)
758                                }
759                                _ = tokio::signal::ctrl_c() => {
760                                    tracing::info!("received Ctrl+C, shutting down gracefully...");
761                                    Ok(())
762                                }
763                                _ = wait_for_sigterm() => {
764                                    tracing::info!("received SIGTERM, shutting down gracefully...");
765                                    Ok(())
766                                }
767                                _ = async {
768                                    while shutdown_rx.changed().await.is_ok() {
769                                        if *shutdown_rx.borrow() {
770                                            break;
771                                        }
772                                    }
773                                } => {
774                                    tracing::info!("shutdown requested via telnet");
775                                    Ok(())
776                                }
777                                        }
778                            }
779                        }
780                    }
781                }
782            }
783        }
784    };
785
786    // Publish the gateway-init failure/suspension outcome before stopping the
787    // surface loops. Deferred InitConnect completions observe this generation
788    // transition and enqueue their failure response while each connection's
789    // bounded send queue is still owned by the active listener.
790    transition_startup_to_shutdown(bridge.startup_readiness());
791    if tokio::time::timeout(
792        SHUTDOWN_GRACE_PERIOD,
793        bridge
794            .startup_readiness()
795            .await_pending_init_connects_drained(),
796    )
797    .await
798    .is_err()
799    {
800        tracing::warn!(
801            pending_init_connects = bridge.startup_readiness().pending_init_connect_count(),
802            "timed out flushing deferred InitConnect completions before surface shutdown"
803        );
804    }
805
806    // 清理后台任务
807    request_surface_shutdown(&shutdown_tx, "phase4 exit");
808    if !bridge
809        .qot_runtime()
810        .indicator_calc_runtime()
811        .shutdown_and_wait(SHUTDOWN_GRACE_PERIOD)
812        .await
813    {
814        tracing::warn!("timed out cancelling indicator calculations during shutdown");
815    }
816    await_background_shutdown(
817        "card_num retry",
818        card_num_retry_handle,
819        SHUTDOWN_GRACE_PERIOD,
820    )
821    .await;
822    if let Some(handle) = sighup_handle {
823        await_background_shutdown("SIGHUP reload", handle, SHUTDOWN_GRACE_PERIOD).await;
824    }
825    if let Some(handle) = surface_tasks.ftapi {
826        await_surface_shutdown("API server", handle, SHUTDOWN_GRACE_PERIOD).await;
827    }
828    if let Some(handle) = surface_tasks.websocket {
829        await_surface_shutdown("WebSocket", handle, SHUTDOWN_GRACE_PERIOD).await;
830    }
831    if let Some(handle) = surface_tasks.rest {
832        await_surface_shutdown("REST API", handle, SHUTDOWN_GRACE_PERIOD).await;
833    }
834    if let Some(handle) = surface_tasks.grpc {
835        await_surface_shutdown("gRPC", handle, SHUTDOWN_GRACE_PERIOD).await;
836    }
837    if let Some(handle) = surface_tasks.telnet {
838        await_surface_shutdown("Telnet", handle, SHUTDOWN_GRACE_PERIOD).await;
839    }
840
841    phase4_result?;
842    tracing::info!("gateway stopped");
843    Ok(())
844}