Skip to main content

futu_opend/startup/phase4/
card_num.rs

1//! v1.4.110 Layer 3 A: startup Phase 4 — card_num expand / background retry /
2//! unified SIGHUP reload handler.
3//!
4//! 纯搬移自 `phase4.rs::run_phase4` 的 card_num 段(v1.8.0 large-file split,
5//! 零行为变化): 首次 expand 立即执行、retry loop 与 SIGHUP handler 的 spawn 顺序
6//! 与原文一致; 两个 JoinHandle 交回 `run_phase4` 在 shutdown 尾部 join.
7
8use std::sync::Arc;
9
10use futu_gateway_core::bridge::GatewayBridge;
11
12use super::shutdown::BackgroundJoinHandle;
13
14// v1.4.103 (B10) + codex F2 (P1): 立即跑首次 expand + background retry
15// + SIGHUP reload 钩子 (fail-closed sentinel 即时生效, 无 startup window).
16//
17// 行为:
18// 1. **立即跑首次 expand** — 即便 cache 空, fail-closed sentinel (codex F1)
19//    也写进 allowed_acc_ids, 短路 startup window 的 silent unrestricted.
20// 2. **background retry**: 每 10s 检查 cache, 加载后 re-expand 真 acc_id
21//    覆盖 sentinel. 60s 上限.
22// 3. **SIGHUP reload 钩子**: 重载 keys.json 后再 expand 一次 (codex F2 reload window).
23// v1.4.103 codex F3.1 (P1) round 3: 合并 reload + expand 为单一 ordered op,
24// 防 race. 当 SIGHUP 触发时, 此 fn 同步顺序: (1) reload 所有 store (从
25// keys.json 读 raw allowed_card_nums + ArcSwap.store) (2) expand (resolve
26// card_num → acc_ids + ArcSwap.store with sentinel/resolved). 因为是单线
27// 调用, 不存在多 SIGHUP listener 间的乱序 race.
28//
29// `do_reload` 控制是否先 reload (启动时第一次 / 60s retry 不需 reload,
30// 直接 expand; SIGHUP 路径需要先 reload).
31pub(super) fn spawn_card_num_expand_tasks(
32    bridge: &Arc<GatewayBridge>,
33    ws_key_store_holder: &Option<Arc<futu_auth::KeyStore>>,
34    rest_key_store_holder: &Option<Arc<futu_auth::KeyStore>>,
35    grpc_key_store_holder: &Option<Arc<futu_auth::KeyStore>>,
36    shutdown_rx: &tokio::sync::watch::Receiver<bool>,
37) -> (BackgroundJoinHandle, Option<BackgroundJoinHandle>) {
38    let card_num_reload_and_expand_fn: std::sync::Arc<dyn Fn(bool) + Send + Sync> = {
39        let bridge_for_expand = std::sync::Arc::clone(bridge);
40        let ws_ks = ws_key_store_holder.clone();
41        let rest_ks = rest_key_store_holder.clone();
42        let grpc_ks = grpc_key_store_holder.clone();
43        std::sync::Arc::new(move |do_reload: bool| {
44            // (1) Reload phase — 仅 SIGHUP 路径调
45            if do_reload {
46                for (ks_name, ks_opt) in [("ws", &ws_ks), ("rest", &rest_ks), ("grpc", &grpc_ks)] {
47                    let Some(ks) = ks_opt.as_ref() else { continue };
48                    match ks.reload() {
49                        Ok(()) => tracing::warn!(
50                            ks = ks_name,
51                            keys_loaded = ks.len(),
52                            "v1.4.103 F3.1: keys reloaded on SIGHUP (before card_num expand)"
53                        ),
54                        Err(e) => tracing::error!(
55                            ks = ks_name,
56                            error = %e,
57                            "v1.4.103 F3.1: keys reload failed (skipping expand for this store)"
58                        ),
59                    }
60                }
61            }
62            // (2) Expand phase
63            let trd_cache = std::sync::Arc::clone(&bridge_for_expand.caches().trd_cache);
64            let resolver = {
65                let cache_clone = std::sync::Arc::clone(&trd_cache);
66                move |cn: &str| cache_clone.find_acc_ids_by_card_num(cn)
67            };
68            for (ks_name, ks_opt) in [("ws", &ws_ks), ("rest", &rest_ks), ("grpc", &grpc_ks)] {
69                let Some(ks) = ks_opt.as_ref() else { continue };
70                let (resolved, unresolved, ambiguous) = ks.expand_allowed_card_nums(
71                    &resolver,
72                    |key_id, cn| {
73                        tracing::warn!(
74                            key_id = %key_id,
75                            card_num = %cn,
76                            "v1.4.103 B10/F1 fail-closed: card_num not found in trd_cache; \
77                             writing sentinel acc_id=0 to enforce restrictive denylist \
78                             (limits.contains check 永远 false → reject 真账户)"
79                        );
80                    },
81                    |key_id, cn, candidates| {
82                        tracing::warn!(
83                            key_id = %key_id,
84                            card_num = %cn,
85                            candidates = ?candidates,
86                            "v1.4.103 B10/F1 fail-closed: ambiguous card_num suffix \
87                             matched multiple accounts (skipped, write 完整 16 位 / specific 4 位)"
88                        );
89                    },
90                );
91                tracing::info!(
92                    ks = ks_name,
93                    resolved,
94                    unresolved,
95                    ambiguous,
96                    "v1.4.103 B10: expanded allowed_card_nums into allowed_acc_ids"
97                );
98            }
99        })
100    };
101    // 老 alias (不 reload, 仅 expand) 用于启动 + retry 路径
102    let card_num_expand_fn: std::sync::Arc<dyn Fn() + Send + Sync> = {
103        let inner = std::sync::Arc::clone(&card_num_reload_and_expand_fn);
104        std::sync::Arc::new(move || (inner)(false))
105    };
106
107    // codex F2 (P1) 立即跑首次: fail-closed sentinel 即时生效, 不留 startup window
108    (card_num_expand_fn)();
109
110    // background retry loop: cache 加载后 re-expand 覆盖 sentinel
111    let card_num_retry_handle: BackgroundJoinHandle = {
112        let card_num_expand_fn_loop = std::sync::Arc::clone(&card_num_expand_fn);
113        let bridge_for_check = std::sync::Arc::clone(bridge);
114        let mut card_num_retry_shutdown_rx = shutdown_rx.clone();
115        tokio::spawn(async move {
116            let trd_cache = std::sync::Arc::clone(&bridge_for_check.caches().trd_cache);
117            let mut attempts = 0u32;
118            let max_attempts = 6u32; // 6 × 10s = 60s
119            loop {
120                tokio::select! {
121                    changed = card_num_retry_shutdown_rx.changed() => {
122                        if changed.is_err() || *card_num_retry_shutdown_rx.borrow() {
123                            tracing::debug!(
124                                "v1.4.111: card_num retry loop received shutdown signal"
125                            );
126                            return;
127                        }
128                    }
129                    _ = tokio::time::sleep(std::time::Duration::from_secs(10)) => {}
130                }
131                attempts += 1;
132                let accounts = trd_cache.get_accounts();
133                if accounts.is_empty() {
134                    if attempts >= max_attempts {
135                        tracing::warn!(
136                            "v1.4.103 B10: trd_cache 仍空 (after {max_attempts} × 10s); \
137                             受限 key 仍走 fail-closed sentinel reject 直到 SIGHUP / cache 加载."
138                        );
139                        return;
140                    }
141                    continue;
142                }
143                (card_num_expand_fn_loop)();
144                return;
145            }
146        })
147    };
148
149    // codex F2 (P1) + F3.1 (P1) round 3: 单一 SIGHUP 钩子 — reload + expand
150    // 顺序操作避免 race. 之前每 server (REST/gRPC) 各自 SIGHUP 监听 reload,
151    // 加上 card_num expand 单独监听 — 多 SIGHUP listener 并发执行无序, 可能
152    // expand 先跑 (写 sentinel) 然后 reload 后跑 (overwrite sentinel 用 raw
153    // allowed_card_nums) → 受限 key 在 reload window 内 silent unrestricted.
154    //
155    // 现在: **唯一 SIGHUP listener**, 顺序 reload → expand. 删除 REST/gRPC 各
156    // 自的 SIGHUP reload listener (上面 server 块内已注释).
157    #[cfg(unix)]
158    let sighup_handle: Option<BackgroundJoinHandle> = {
159        let unified_sighup_fn = std::sync::Arc::clone(&card_num_reload_and_expand_fn);
160        let mut sighup_shutdown_rx = shutdown_rx.clone();
161        Some(tokio::spawn(async move {
162            use tokio::signal::unix::{SignalKind, signal};
163            let mut sig = match signal(SignalKind::hangup()) {
164                Ok(s) => s,
165                Err(e) => {
166                    tracing::error!(error = %e, "SIGHUP install failed (unified reload+expand)");
167                    return;
168                }
169            };
170            tracing::info!(
171                "v1.4.103 F3.1: unified SIGHUP handler installed (reload all keys + expand card_num)"
172            );
173            loop {
174                tokio::select! {
175                    signal = sig.recv() => {
176                        if signal.is_none() {
177                            return;
178                        }
179                        tracing::info!(
180                            "v1.4.103 F3.1: SIGHUP received — running reload_all_stores + \
181                             expand_allowed_card_nums (single ordered op, no race)"
182                        );
183                        (unified_sighup_fn)(true); // do_reload=true
184                    }
185                    changed = sighup_shutdown_rx.changed() => {
186                        if changed.is_err() || *sighup_shutdown_rx.borrow() {
187                            tracing::debug!(
188                                "v1.4.111: unified SIGHUP handler received shutdown signal"
189                            );
190                            return;
191                        }
192                    }
193                }
194            }
195        }))
196    };
197    #[cfg(not(unix))]
198    let sighup_handle: Option<BackgroundJoinHandle> = None;
199
200    (card_num_retry_handle, sighup_handle)
201}