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;