Skip to main content

futu_grpc/
server.rs

1//! gRPC 服务实现
2//!
3//! FutuOpenD 服务通过通用的 proto_id + body 方式,
4//! 将所有请求转发到现有的 RequestRouter。
5//! 支持流式推送:行情、交易、广播事件通过 SubscribePush 接口推送给客户端。
6//!
7//! ## v1.4.106 codex 0517 ζ25-redo F2: stateful QOT stable identity
8//!
9//! gRPC `request()` 与 `subscribe_push()` 共享同一身份派生函数
10//! [`auth::derive_grpc_conn_id`]:从 `(Bearer token, optional grpc-session-id
11//! metadata)` 派生 deterministic stable conn_id(在
12//! [`auth::GRPC_STABLE_CONN_NAMESPACE`] 即 bit 62 namespace 内)。
13//!
14//! 这让同一 caller 的连续 RPC 命中同一 SubscriptionManager / cache 状态:
15//! - subscribe → unsubscribe / get_sub_info / query_subscription 全对齐
16//! - quote cache 命中(per-conn cache 不再每次 RPC miss)
17//! - QOT push fanout per-conn filter 能找到 caller
18//!
19//! 不同 Bearer / 不同 session_id → 自然隔离(caller 之间互不影响)。
20//! Legacy mode(KeyStore 未配置 / Bearer 缺失)→ 共享一个固定 conn_id(
21//! 全 legacy caller 共享同 sub state,对齐"无鉴权配置 = 单租户"语义)。
22//!
23//! 历史:v1.4.105 之前用自增 `conn_id_counter`,每次 RPC 拿新 ID,
24//! 等价 REST v1.4.90 P0-B 之前的 quota 永久泄漏 bug。F2 修复对齐
25//! `REST_SHARED_CONN` 设计哲学,不同点是 gRPC 按 caller 隔离(REST 共享)。
26
27use std::sync::Arc;
28use std::sync::atomic::{AtomicU32, Ordering};
29
30use bytes::Bytes;
31use prost::Message as _;
32use tokio::sync::{broadcast, mpsc};
33use tokio_stream::wrappers::ReceiverStream;
34use tonic::{Request, Response, Status};
35
36use futu_auth::{KeyStore, RuntimeCounters, Scope};
37use futu_codec::header::ProtoFmtType;
38use futu_server::conn::IncomingRequest;
39use futu_server::router::RequestRouter;
40
41use crate::auth::{
42    derive_grpc_conn_id, extract_grpc_idempotency_key, extract_grpc_session_id, extract_grpc_token,
43    grpc_audit_context, grpc_status_for,
44};
45use crate::proto::futu_open_d_server::FutuOpenD;
46use crate::proto::{FutuRequest, FutuResponse, PushEvent, SubscribePushRequest};
47use futu_auth_pipeline::{
48    AuthDecision, AuthEnvelope, Credential, Endpoint, FilterRegistry, PushEventCtx, SurfaceId,
49    authenticate_request,
50};
51
52mod generated_exposure {
53    include!(concat!(env!("OUT_DIR"), "/generated_grpc_exposure.rs"));
54}
55
56mod push;
57mod startup;
58pub use push::GrpcPushBroadcaster;
59pub use startup::{
60    build_service_with_auth, start, start_with_auth, start_with_auth_until_shutdown,
61    start_with_auth_until_shutdown_with_listener_events,
62};
63
64/// Explicit gRPC message boundary.
65///
66/// Native FTAPI frames are capped at 12 MiB; keeping gRPC at the same ceiling
67/// prevents relying on tonic defaults while preserving protocol-sized payloads.
68pub const GRPC_MAX_MESSAGE_SIZE_BYTES: usize = 12 * 1024 * 1024;
69
70/// gRPC 服务实现
71pub struct FutuGrpcService {
72    router: Arc<RequestRouter>,
73    push_broadcaster: Arc<GrpcPushBroadcaster>,
74    key_store: Arc<KeyStore>,
75    counters: Arc<RuntimeCounters>,
76    /// v1.4.104: response filter 注册中心 (proto 2001 = AccListFilter).
77    /// 加新 filter 在 `FilterRegistry::install_defaults` 注册一次, 4 surface
78    /// 自动生效.
79    filter_registry: Arc<FilterRegistry>,
80    /// v1.4.106 codex 0517 ζ25-redo F2: gRPC `conn_id_counter` 已删除 —
81    /// 自增 conn_id 让同一 caller 连续 RPC 拿不到同 sub state, 等价 REST
82    /// v1.4.90 P0-B 之前的 quota 永久泄漏 bug. 现在改用
83    /// [`auth::derive_grpc_conn_id`] 从 `(bearer, session_id)` 派生
84    /// deterministic stable conn_id, 同 caller 连续 RPC 命中同一 sub state.
85    serial_counter: AtomicU32,
86}
87
88impl FutuGrpcService {
89    pub fn new(router: Arc<RequestRouter>, push_broadcaster: Arc<GrpcPushBroadcaster>) -> Self {
90        Self::with_auth(
91            router,
92            push_broadcaster,
93            Arc::new(KeyStore::empty()),
94            Arc::new(RuntimeCounters::new()),
95        )
96    }
97
98    /// 完整构造:同时接 key_store + counters(v1.0 推荐入口)
99    ///
100    /// `counters` 应由 main 全进程共享:REST / gRPC / MCP 共用一个实例才能保证
101    /// rate limit / 日累计跨接口一致
102    pub fn with_auth(
103        router: Arc<RequestRouter>,
104        push_broadcaster: Arc<GrpcPushBroadcaster>,
105        key_store: Arc<KeyStore>,
106        counters: Arc<RuntimeCounters>,
107    ) -> Self {
108        Self {
109            router,
110            push_broadcaster,
111            key_store,
112            counters,
113            filter_registry: Arc::new(FilterRegistry::with_defaults()),
114            serial_counter: AtomicU32::new(1),
115        }
116    }
117
118    fn next_serial(&self) -> u32 {
119        self.serial_counter.fetch_add(1, Ordering::Relaxed)
120    }
121}
122
123#[tonic::async_trait]
124impl FutuOpenD for FutuGrpcService {
125    /// 通用请求-响应
126    async fn request(
127        &self,
128        request: Request<FutuRequest>,
129    ) -> Result<Response<FutuResponse>, Status> {
130        // v1.4.104: 走 futu_auth_pipeline::authenticate_request 单一函数,
131        // 不再 inline authenticate / check_scope / rate gate / body-aware /
132        // audit. surface adapter 极薄 — 仅 transport extract + RejectKind →
133        // Status 翻译.
134
135        let proto_id = request.get_ref().proto_id;
136        let caller_is_loopback = request
137            .remote_addr()
138            .is_some_and(|addr| addr.ip().is_loopback());
139        if proto_id == 0 {
140            return Err(Status::invalid_argument("proto_id is required"));
141        }
142        // v1.4.106 codex 0532 F3 (P2): daemon-internal proto_id (高位
143        // 0x8000_0000 bit) 绝不应从 gRPC 公开 surface 进入 — 仅 REST handler
144        // 内部合成给 router. 显式 reject + audit, 防探测 daemon 内部 routing.
145        if futu_auth::is_internal_proto_id(proto_id) {
146            tracing::warn!(
147                proto_id,
148                "rejecting daemon-internal proto_id at gRPC public surface (audit 0532 F3)"
149            );
150            return Err(Status::permission_denied(
151                "daemon-internal proto_id not allowed on public surface",
152            ));
153        }
154
155        // 1. extract token + session (owned String 避免 borrow vs move 冲突).
156        // v1.4.106 codex 0517 ζ25-redo F2: 同时拿 grpc-session-id metadata,
157        // 派生 stateful QOT stable conn_id (derive_grpc_conn_id).
158        let token = extract_grpc_token(&request);
159        let session_id = extract_grpc_session_id(&request);
160        let idempotency_key = extract_grpc_idempotency_key(&request);
161        let stable_conn_id = derive_grpc_conn_id(token.as_deref(), session_id.as_deref());
162        let audit_ctx = grpc_audit_context(&request, stable_conn_id);
163        // The manifest-generated match is the only public gRPC exposure registry.
164        // Unknown proto IDs cannot bypass EndpointSpec into auth/router.
165        let spec = match generated_exposure::generated_grpc_exposure(proto_id) {
166            Ok(spec) => spec,
167            Err(generated_exposure::GeneratedGrpcReject::NotExposed { endpoint, reason }) => {
168                tracing::warn!(
169                    proto_id,
170                    endpoint,
171                    reason,
172                    "rejecting gRPC request for endpoint not exposed to gRPC"
173                );
174                return Err(Status::permission_denied(reason));
175            }
176            Err(generated_exposure::GeneratedGrpcReject::Unknown) => {
177                tracing::warn!(
178                    proto_id,
179                    "rejecting gRPC request for proto_id not declared in the surface manifest"
180                );
181                return Err(Status::invalid_argument(
182                    "proto_id is not declared in the surface manifest",
183                ));
184            }
185        };
186        let endpoint_name = spec.canonical_name;
187        tracing::debug!(
188            proto_id,
189            endpoint = endpoint_name,
190            "gRPC request dispatch (Layer 2 spec lookup)"
191        );
192        let req_inner = request.into_inner();
193        let request_body = normalize_grpc_history_query_body(proto_id, req_inner.body)?;
194        let needed_scope = if spec.runtime.side_effects == futu_surface_spec::SideEffectKind::Write
195            && matches!(spec.runtime.scope, Scope::TradeReal | Scope::TradeSimulate)
196        {
197            futu_auth_pipeline::body_aware::trade_write_scope_from_body(proto_id, &request_body)
198                .or(Some(spec.runtime.scope))
199        } else {
200            Some(spec.runtime.scope)
201        };
202
203        // 2. 构 credential + envelope, 调 pipeline
204        let credential = match token.as_deref() {
205            Some(t) => Credential::Bearer(t),
206            None => Credential::None,
207        };
208        let env = AuthEnvelope {
209            surface: SurfaceId::Grpc,
210            endpoint: Endpoint::Proto(proto_id),
211            needed_scope,
212            credential,
213            proto_id: Some(proto_id),
214            body: &request_body,
215            explicit_acc_id: None,
216            explicit_ctx: None,
217            commit_rate: true, // gRPC middleware 层 commit rate (trade:real)
218            audit_emit: true,
219        };
220        let (allowed_for_filter, caller_rec) =
221            match futu_auth::audit::with_context(audit_ctx.clone(), || {
222                authenticate_request(&self.key_store, &self.counters, env)
223            }) {
224                AuthDecision::Reject { kind, reason, .. } => {
225                    return Err(grpc_status_for(kind, reason));
226                }
227                AuthDecision::Allow {
228                    allowed_acc_ids,
229                    rec,
230                    ..
231                } => (allowed_acc_ids, rec),
232            };
233
234        // 3. dispatch + response filter
235        // v1.4.105 D2 T-A1 fix: caller_allowed_acc_ids 从 pipeline allow decision
236        // 真填进 IncomingRequest, 让 dispatch handler (e.g. SubAccPushHandler)
237        // 端 enforce per-acc whitelist defense-in-depth.
238        // codex 0522 F1 v1.4.106: 同步填 caller_key_id 让 cross-surface handler
239        // 能识别 gRPC caller.
240        let caller_allowed = allowed_for_filter
241            .as_ref()
242            .map(|s| std::sync::Arc::new(s.clone()));
243        let caller_key_id = caller_rec.as_ref().map(|r| r.id.clone());
244        let caller_has_auth_setup_scope = caller_rec
245            .as_ref()
246            .is_some_and(|record| record.scopes.contains(&Scope::AuthSetup));
247        let serial_no = self.next_serial();
248        let body = futu_server::trade_packet_id::fill_omitted_trade_packet_id_bytes(
249            proto_id,
250            request_body,
251            stable_conn_id,
252            serial_no,
253        )
254        .map_err(|e| Status::internal(format!("failed to prepare trade PacketID: {e}")))?;
255        let incoming = IncomingRequest::builder(
256            stable_conn_id,
257            proto_id,
258            serial_no,
259            ProtoFmtType::Protobuf,
260            Bytes::from(body),
261        )
262        .with_idempotency_key(idempotency_key)
263        .with_caller_scope(caller_allowed, caller_key_id)
264        .with_auth_setup_admission(
265            caller_has_auth_setup_scope,
266            caller_is_loopback,
267            !self.key_store.is_configured(),
268        )
269        .build();
270
271        match self.router.dispatch(incoming.conn_id, &incoming).await {
272            Some(resp_bytes) => {
273                // v1.4.104: 用 FilterRegistry (单一注册表) 替代 inline
274                // filter_acc_list_response. 加新 filter (e.g. cash-flow) 只
275                // 在 registry 注册一次, 4 surface 自动生效.
276                let filtered_body =
277                    self.filter_registry
278                        .apply(proto_id, resp_bytes, allowed_for_filter.as_ref());
279                Ok(Response::new(FutuResponse {
280                    ret_type: 0,
281                    ret_msg: String::new(),
282                    proto_id,
283                    body: filtered_body,
284                }))
285            }
286            None => Ok(Response::new(FutuResponse {
287                ret_type: -1,
288                ret_msg: "handler returned no response".to_string(),
289                proto_id,
290                body: Vec::new(),
291            })),
292        }
293    }
294
295    type SubscribePushStream = ReceiverStream<Result<PushEvent, Status>>;
296
297    /// 流式推送:客户端建立连接后持续接收行情、交易、广播推送
298    ///
299    /// v1.1:按订阅 key 的 scope 过滤推送 —— `qot:read`-only 的 key 不会收到
300    /// `trade` 类(账户交易回报),对齐 REST `/ws` v0.9.0 加的 push filter。
301    ///
302    /// ## v1.4.104 阶段 7-2: pipeline 委托
303    ///
304    /// **handshake**: pipeline 调一次 with `Endpoint::Event("subscribe_push")` +
305    /// `needed_scope=Some(qot:read)`. Allow → 拿 rec snapshot. Reject (Bearer
306    /// invalid / expired / no qot:read) → translate to gRPC Status.
307    ///
308    /// **per-event filter**: stream 内每 event 调一次 pipeline with
309    /// `Credential::PreVerified(rec)` + `Endpoint::Event(event_type)` +
310    /// `needed_scope=Some(scope_for_event(...))` + `explicit_acc_id` (trade event
311    /// 给 event.acc_id, 其他不传) + `audit_emit=false` (避免每 event 一条 audit
312    /// 把日志冲爆) + `commit_rate=false` (push 不计 rate). Reject → silent drop +
313    /// `metrics::bump_ws_filtered`. Allow → forward.
314    ///
315    /// `SubscribePush` 是 quote-first mixed push stream: handshake 必须有
316    /// qot:read;caller 若还持有 acc:read, stream 内 trade event 再由 per-event
317    /// scope match + trade acc_id whitelist 放行。acc:read-only key 不能打开
318    /// mixed stream。
319    async fn subscribe_push(
320        &self,
321        request: Request<SubscribePushRequest>,
322    ) -> Result<Response<Self::SubscribePushStream>, Status> {
323        if !grpc_push_ready(&self.router.startup_readiness()) {
324            return Err(Status::unavailable(
325                "gateway authentication is not ready for push subscriptions",
326            ));
327        }
328        // v1.4.106 codex 1125 F6 [P2]: 显式 notify subscription opt-in.
329        // 对齐 C++ raw TCP `IsConnSubRecvNotify` 语义 (broadcast push 必须显式
330        // sub 才下发).
331        let notify_subscribe = request.get_ref().notify_subscribe;
332
333        // ── handshake: pipeline 调一次拿 caller-key + audit 一次 ──────────────────
334        //
335        // Bearer 提取 / verify / expiry / qot:read scope 都走 pipeline
336        // (单一 source)。trade push 不是单独的 acc-only stream;caller 需要
337        // qot:read 打开 stream, 再靠 per-event acc:read 过滤接收 trade event.
338        //
339        // v1.4.106 codex 0517 ζ25-redo F2: 与 `request()` 共享 stable conn_id.
340        // 这让 ops/log 能把 push stream 和同 caller 的 subscribe request 串起来:
341        // grep `conn_id=<stable>` 看到 subscribe RPC + push event filter 同 id.
342        let token = extract_grpc_token(&request);
343        let session_id = extract_grpc_session_id(&request);
344        let stable_conn_id = derive_grpc_conn_id(token.as_deref(), session_id.as_deref());
345        let audit_ctx = grpc_audit_context(&request, stable_conn_id);
346        let credential = match token.as_deref() {
347            Some(t) => Credential::Bearer(t),
348            None => Credential::None,
349        };
350        let env = AuthEnvelope {
351            surface: SurfaceId::Grpc,
352            endpoint: Endpoint::Event("subscribe_push"),
353            needed_scope: Some(Scope::QotRead),
354            credential,
355            proto_id: None,
356            body: &[],
357            explicit_acc_id: None,
358            explicit_ctx: None,
359            commit_rate: false,
360            audit_emit: true, // handshake 一次 audit
361        };
362        let rec_snapshot = match futu_auth::audit::with_context(audit_ctx.clone(), || {
363            authenticate_request(&self.key_store, &self.counters, env)
364        }) {
365            AuthDecision::Reject { kind, reason, .. } => {
366                return Err(grpc_status_for(kind, reason));
367            }
368            AuthDecision::Allow { rec, .. } => rec,
369        };
370
371        let (scopes, key_id, allowed_acc_ids) = match rec_snapshot.as_ref() {
372            Some(rec) => (
373                rec.scopes.clone(),
374                rec.id.clone(),
375                rec.allowed_acc_ids.clone(),
376            ),
377            None => (
378                // legacy: 全 scope + 无 acc 限制
379                [
380                    Scope::QotRead,
381                    Scope::AccRead,
382                    Scope::TradeSimulate,
383                    Scope::TradeReal,
384                ]
385                .into_iter()
386                .collect::<std::collections::HashSet<Scope>>(),
387                "<none>".to_string(),
388                None,
389            ),
390        };
391
392        let (tx, rx) = mpsc::channel(256);
393        let mut push_rx = self.push_broadcaster.subscribe();
394
395        tracing::info!(
396            key_id = %key_id,
397            // v1.4.106 codex 0517 ζ25-redo F2: stable conn_id 让 subscribe RPC
398            // 和 push stream 在 log 里可关联.
399            conn_id = stable_conn_id,
400            scopes = ?scopes,
401            allowed_acc_ids = ?allowed_acc_ids.as_ref().map(|s| s.len()),
402            "gRPC client subscribed to push events",
403        );
404
405        // ── per-event filter: pipeline 二次调 (audit_emit=false 防日志冲爆) ───────
406        //
407        // 每 event 调一次 pipeline:
408        //   - Credential::PreVerified(rec) 复用 handshake 已 verify 的 rec
409        //   - Endpoint::Event(event_type) 给 audit label (虽然 audit_emit=false)
410        //   - needed_scope = scope_for_event(event_type)
411        //   - explicit_acc_id = trade event 给 event.acc_id, 其他 None
412        //
413        // legacy mode (rec_snapshot=None) 不调 pipeline (没必要, 全放行).
414        let key_store_arc = self.key_store.clone();
415        let counters_arc = self.counters.clone();
416        // v1.4.105 D3 (Phase 4): clone filter_registry for per-event Layer 3
417        // (allowed_markets) check + Layer 1/2 reuse.
418        let filter_registry_arc = self.filter_registry.clone();
419        let push_readiness = self.router.startup_readiness();
420        tokio::spawn(async move {
421            loop {
422                match push_rx.recv().await {
423                    Ok(event) => {
424                        if !grpc_push_ready(&push_readiness) {
425                            futu_auth::metrics::bump_ws_filtered("startup_not_ready", &key_id);
426                            continue;
427                        }
428                        // legacy fast-path: rec=None → 全放行, scope check / acc_id
429                        // whitelist 都不需要走 pipeline.
430                        let allow_event = if let Some(rec) = rec_snapshot.as_ref() {
431                            let needed = scope_for_event(&event.event_type);
432                            let explicit_acc_id = if event.event_type == "trade" {
433                                Some(event.acc_id)
434                            } else {
435                                None
436                            };
437                            let env = AuthEnvelope {
438                                surface: SurfaceId::Grpc,
439                                endpoint: Endpoint::Event(&event.event_type),
440                                needed_scope: Some(needed),
441                                credential: Credential::PreVerified(rec.clone()),
442                                proto_id: None,
443                                body: &[],
444                                explicit_acc_id,
445                                explicit_ctx: None,
446                                commit_rate: false, // push 不计 rate
447                                audit_emit: false,  // per-event 不 audit, 避免冲爆
448                            };
449                            matches!(
450                                authenticate_request(&key_store_arc, &counters_arc, env),
451                                AuthDecision::Allow { .. }
452                            )
453                        } else {
454                            true
455                        };
456
457                        if !allow_event {
458                            // metric label 区分 trade acc_id reject vs scope reject
459                            // (近似旧行为: trade 类 reject 标 "trade_acc_id" 仅当
460                            // acc_id whitelist 非空; 否则按 event_type)
461                            let label = if event.event_type == "trade"
462                                && rec_snapshot
463                                    .as_ref()
464                                    .and_then(|r| r.allowed_acc_ids.as_ref())
465                                    .is_some_and(|s| !s.is_empty())
466                            {
467                                "trade_acc_id"
468                            } else {
469                                event.event_type.as_str()
470                            };
471                            futu_auth::metrics::bump_ws_filtered(label, &key_id);
472                            continue;
473                        }
474
475                        // v1.4.105 D3 (Phase 4): FilterRegistry::should_drop_event
476                        // Layer 3 (allowed_markets). Layer 1 (allowed_acc_ids) 已在
477                        // 上面 pipeline body-aware 跑过, 此处串行 Layer 3 双重
478                        // 防御 + 未来 Layer 扩展自动覆盖.
479                        // ⚠️ UNVERIFIED — pending real-machine verify (HK+US 双
480                        // 账户跨 market 推送). 这里是 backend-semantic 风险,
481                        // 只能靠真机跨 market trade push 验证后再升 confidence.
482                        //
483                        // v1.4.105 T-B3: trd_market 改读 PushEvent.trd_market 字段
484                        // (PushDispatcher 端一次 decode), 不再各 surface 独立
485                        // decode body. 空字符串 = 非 trade event / decode 失败 /
486                        // market 未知 → 转 None (Layer 3 不 trigger drop, 向后
487                        // 兼容).
488                        let event_trd_market =
489                            if event.event_type == "trade" && !event.trd_market.is_empty() {
490                                Some(event.trd_market.as_str())
491                            } else {
492                                None
493                            };
494                        let allowed_markets_for_filter = rec_snapshot
495                            .as_ref()
496                            .and_then(|r| r.allowed_markets.as_ref());
497                        let push_ctx = PushEventCtx {
498                            event_type: &event.event_type,
499                            event_acc: if event.event_type == "trade" {
500                                Some(event.acc_id)
501                            } else {
502                                None
503                            },
504                            // Layer 1 已 pipeline 跑, 此处传 None 防双重 reject (LE 之间 OR 短路)
505                            allowed_acc_ids: None,
506                            // gRPC 没 sub_state (REST 特有)
507                            sub_state: None,
508                            // Layer 3 — 新加
509                            event_trd_market,
510                            allowed_markets: allowed_markets_for_filter,
511                        };
512                        if filter_registry_arc.should_drop_event(&push_ctx) {
513                            futu_auth::metrics::bump_ws_filtered("trade_market", &key_id);
514                            continue;
515                        }
516
517                        // v1.4.106 codex 1125 F6 [P2]: notify_subscribe gate.
518                        // 对齐 C++ `IsConnSubRecvNotify` (raw TCP). broadcast notify
519                        // 类 push (e.g. price reminder) 必须 client 显式 sub 才下发.
520                        if event.event_type == "notify" && !notify_subscribe {
521                            futu_auth::metrics::bump_ws_filtered("notify_unsub", &key_id);
522                            continue;
523                        }
524
525                        if tx.send(Ok(event)).await.is_err() {
526                            break; // 客户端断开
527                        }
528                    }
529                    Err(broadcast::error::RecvError::Lagged(n)) => {
530                        tracing::warn!(skipped = n, "gRPC push client lagged, skipped events");
531                        // 继续接收,不断开
532                    }
533                    Err(broadcast::error::RecvError::Closed) => {
534                        break; // 广播器关闭
535                    }
536                }
537            }
538            tracing::info!("gRPC push stream ended");
539        });
540
541        Ok(Response::new(ReceiverStream::new(rx)))
542    }
543}
544
545fn normalize_grpc_history_query_body(proto_id: u32, body: Vec<u8>) -> Result<Vec<u8>, Status> {
546    match proto_id {
547        futu_core::proto_id::TRD_GET_HISTORY_ORDER_LIST => {
548            normalize_history_order_list_filter_market(body)
549        }
550        futu_core::proto_id::TRD_GET_HISTORY_ORDER_FILL_LIST => {
551            normalize_history_order_fill_list_filter_market(body)
552        }
553        _ => Ok(body),
554    }
555}
556
557fn normalize_history_order_list_filter_market(body: Vec<u8>) -> Result<Vec<u8>, Status> {
558    let Ok(mut req) = futu_proto::trd_get_history_order_list::Request::decode(body.as_slice())
559    else {
560        return Ok(body);
561    };
562    if default_history_filter_market(
563        req.c2s.header.trd_market,
564        &mut req.c2s.filter_conditions.filter_market,
565    ) {
566        return Ok(req.encode_to_vec());
567    }
568    Ok(body)
569}
570
571fn normalize_history_order_fill_list_filter_market(body: Vec<u8>) -> Result<Vec<u8>, Status> {
572    let Ok(mut req) = futu_proto::trd_get_history_order_fill_list::Request::decode(body.as_slice())
573    else {
574        return Ok(body);
575    };
576    if default_history_filter_market(
577        req.c2s.header.trd_market,
578        &mut req.c2s.filter_conditions.filter_market,
579    ) {
580        return Ok(req.encode_to_vec());
581    }
582    Ok(body)
583}
584
585fn default_history_filter_market(header_market: i32, filter_market: &mut Option<i32>) -> bool {
586    // Ref: FutuOpenD/Src/APIServer/Business/Trade/APIServer_Trd_GetHistoryOrderList.cpp:48-55
587    // and _APIServer_Trd_Comm.cpp:3233-3247. C++ uses filterMarket both for
588    // time parsing and response market filtering; missing filterMarket means
589    // Unknown and passes every market. CLI/MCP/REST high-level paths already
590    // default it from header.trd_market, so generic gRPC must do the same.
591    if filter_market.is_none() && header_market != 0 {
592        *filter_market = Some(header_market);
593        return true;
594    }
595    false
596}
597
598// v1.4.104: `grpc_handler_full_check` 已删除 — 功能搬到
599// `futu_auth_pipeline::pipeline::authenticate_request` 内的 body-aware step.
600// gRPC `request()` 调一次 pipeline 即覆盖所有 (caller-key + scope + body-aware
601// + audit + rate). 4 surface 共用同一函数, 不再有 hand-copy diverge.
602
603/// gRPC PushEvent 的 event_type → 客户端必须持有的 scope (与 REST `/ws` 对齐).
604/// 用于 SubscribePush stream 内 per-event filter.
605fn scope_for_event(event_type: &str) -> Scope {
606    match event_type {
607        "trade" => Scope::AccRead, // 账户交易回报
608        _ => Scope::QotRead,       // quote / notify / 未知都按行情门槛
609    }
610}
611
612fn grpc_push_ready(readiness: &futu_server::identity::StartupReadiness) -> bool {
613    readiness.snapshot().state == futu_server::identity::StartupState::Ready
614}
615
616// v1.4.105 D3 (Phase 4) T-B3: 旧 `extract_trd_market_from_trade_event` 已搬到
617// `futu_server::push::extract_trd_market_from_trade_body` + PushDispatcher 端
618// 一次 decode 后透传给 PushEvent.trd_market 字段, 4 surface (gRPC / REST WS /
619// MCP) 共用. gRPC subscribe_push 现读 event.trd_market 字段不再 decode body.
620//
621// 这里若想要静态 import 提示, 可保留:
622//
623//     use futu_server::push::extract_trd_market_from_trade_body;
624//
625// 但目前 server.rs 不再调用, 全凭 PushEvent 字段透传, 故不 import.
626
627#[cfg(test)]
628mod tests;