Skip to main content

futu_mcp/state/
push_subscribers.rs

1use std::collections::{HashSet, VecDeque};
2use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
3use std::sync::{Arc, Mutex as StdMutex, Weak};
4use std::time::Instant;
5
6use rmcp::{
7    RoleServer,
8    model::Resource,
9    service::{Peer, SubscriptionSink},
10};
11use tokio::sync::{OwnedSemaphorePermit, TryAcquireError};
12
13use super::ServerState;
14use super::push_filter::subscriber_visible_to_caller;
15
16mod modern_delivery;
17mod visibility;
18use visibility::current_scope_allows_subscriber;
19pub(crate) use visibility::{parse_push_resource_uri, push_resource_uri};
20
21pub(crate) const PUSH_RESOURCE_PREFIX: &str = "futu://push/";
22pub(crate) const MODERN_PUSH_QUEUE_CAPACITY: usize = 64;
23pub(crate) const MAX_MODERN_PUSH_HANDLES: usize = 128;
24pub(crate) const PUSH_SUBSCRIBER_MAX_AGE: std::time::Duration =
25    std::time::Duration::from_secs(4 * 3600);
26
27/// One rmcp service instance owns all legacy push rows created by that logical
28/// transport session. Rows keep only a weak opaque identity, so TTL and
29/// explicit unsubscribe do not leave a second, unbounded handle index behind.
30pub(crate) struct LegacyPushServiceLease {
31    state: ServerState,
32    owner: Arc<()>,
33    registration_attempted: AtomicBool,
34}
35
36impl LegacyPushServiceLease {
37    pub(crate) fn new(state: ServerState) -> Self {
38        Self {
39            state,
40            owner: Arc::new(()),
41            registration_attempted: AtomicBool::new(false),
42        }
43    }
44
45    pub(crate) async fn register_push_subscriber(
46        &self,
47        delivery: PushDeliveryTarget,
48        acc_ids: HashSet<u64>,
49        owner_key_id: Option<String>,
50        allowed_acc_ids_snapshot: Option<HashSet<u64>>,
51        allowed_markets_snapshot: Option<HashSet<String>>,
52    ) -> Result<String, String> {
53        if !matches!(delivery, PushDeliveryTarget::LegacyPeer(_, _)) {
54            return Err("legacy push service lease cannot own a modern resource".to_string());
55        }
56        // Store before the only await. A cancelled/failed registration may
57        // cause one harmless empty scan, but can never leave an inserted row
58        // without arming final service cleanup.
59        self.registration_attempted.store(true, Ordering::Release);
60        self.state
61            .register_push_subscriber_with_owner(
62                delivery,
63                acc_ids,
64                owner_key_id,
65                allowed_acc_ids_snapshot,
66                allowed_markets_snapshot,
67                Some(Arc::downgrade(&self.owner)),
68            )
69            .await
70    }
71
72    #[cfg(test)]
73    pub(crate) fn mark_registration_attempted_for_test(&self) {
74        self.registration_attempted.store(true, Ordering::Release);
75    }
76}
77
78impl Drop for LegacyPushServiceLease {
79    fn drop(&mut self) {
80        if !self.registration_attempted.load(Ordering::Acquire) {
81            return;
82        }
83        // Pointer identity remains comparable after the last strong owner is
84        // dropped. Do not use `upgrade()` in the cleanup path.
85        let owner = Arc::downgrade(&self.owner);
86        let state = self.state.clone();
87        match tokio::runtime::Handle::try_current() {
88            Ok(runtime) => {
89                runtime.spawn(async move {
90                    let removed = state
91                        .remove_legacy_push_subscribers_for_service(&owner)
92                        .await;
93                    tracing::debug!(
94                        removed,
95                        "removed legacy MCP push subscribers after service close"
96                    );
97                });
98            }
99            Err(error) => {
100                // Runtime shutdown is not evidence that async cleanup ran; the
101                // existing subscriber TTL remains the fallback. Never panic in Drop.
102                tracing::debug!(
103                    %error,
104                    "legacy MCP service dropped without a live Tokio runtime; deferred cleanup not scheduled"
105                );
106            }
107        }
108    }
109}
110
111#[derive(Clone)]
112pub(crate) enum PushDeliveryTarget {
113    LegacyPeer(Peer<RoleServer>, Arc<AtomicU8>),
114    ModernResource(Arc<OwnedSemaphorePermit>),
115}
116
117#[derive(Clone)]
118pub(super) enum SubscriberDelivery {
119    LegacyPeer(Peer<RoleServer>, Arc<AtomicU8>),
120    ModernResource(Arc<ModernResourceDelivery>),
121}
122
123struct ActiveModernListener {
124    token: u64,
125    sink: SubscriptionSink,
126    cancel: tokio_util::sync::CancellationToken,
127}
128
129struct ModernResourceInner {
130    queue: VecDeque<serde_json::Value>,
131    dropped_events: u64,
132    enqueued_generation: u64,
133    next_listener_token: u64,
134    active_listener: Option<ActiveModernListener>,
135    pending_notification: Option<u64>,
136    #[cfg(test)]
137    notification_send_pause: Option<(Arc<tokio::sync::Barrier>, Arc<tokio::sync::Barrier>)>,
138    #[cfg(test)]
139    notification_send_attempts: u64,
140    #[cfg(test)]
141    notification_send_failures: u64,
142    closed: bool,
143}
144
145pub(super) struct ModernResourceDelivery {
146    handle: String,
147    inner: StdMutex<ModernResourceInner>,
148    _permit: Arc<OwnedSemaphorePermit>,
149}
150
151pub(super) struct PendingModernNotification {
152    listener_token: u64,
153    sink: SubscriptionSink,
154    cancel: tokio_util::sync::CancellationToken,
155    uri: String,
156    generation: u64,
157}
158
159/// v1.4.38 Phase 5: 订阅了 push 通知的 MCP 客户端 session 记录。
160///
161/// 每个 session 用 `futu_sub_acc_push` 注册时登记一条。daemon push 事件来到
162/// MCP server 时,按 `acc_ids` 过滤后用 `peer.notify_logging_message()` 推回。
163///
164/// v1.4.38 100%:`acc_ids` 过滤已生效(state.rs drain loop 实装)。
165/// caller key ownership / scope 快照在注册时解析并保存,后续 push 分发不再
166/// 重新读取 bearer 明文。
167#[derive(Clone)]
168pub(super) struct PushSubscriber {
169    pub delivery: SubscriberDelivery,
170    /// 该 session 关心的 acc_id 列表。空集合表示"不过滤"(接收所有 acc 的 push)。
171    pub acc_ids: HashSet<u64>,
172    /// v1.4.39 per-key acc_id 白名单**注册时快照**(非 live-reload)。
173    ///
174    /// Some(set) + non-empty → push 的 acc_id 必须在 set 里才推。
175    /// Some(empty) / None → 不做 key 级过滤(兼容无 allowed_acc_ids 约束的 key
176    /// 或 stdio / legacy 模式)。
177    ///
178    /// 快照语义:注册后 SIGHUP 重载 keys.json 修改 allowed_acc_ids 不会立即
179    /// 反映到已注册订阅者。用户需重新 `futu_sub_acc_push` 才应用新 scope。
180    /// 这是 defense-in-depth 层(主 auth 在 tool 调用时 guard.rs),可接受。
181    pub allowed_acc_ids_snapshot: Option<HashSet<u64>>,
182    /// v1.4.105 T-C2: per-key `allowed_markets` 注册时快照, Layer 3 trade push gate.
183    ///
184    /// `Some(set)` + non-empty → push event 的 `trd_market` 必须 ∈ set 才推 (按
185    /// `Trd_Common.TrdMarket` int → 字符串映射, 与 `keys.json::allowed_markets`
186    /// 配置字符串一致, e.g. "HK"/"US"/"FUTURES").
187    /// `None` / `Some(empty)` → 不做 market gate (兼容 stdio / legacy / 未配
188    /// allowed_markets 的 key).
189    ///
190    /// 与 `allowed_acc_ids_snapshot` 同样**注册时快照**, SIGHUP 重载不影响
191    /// 已注册订阅者. 用户需重新 `futu_sub_acc_push` 才应用新 scope.
192    /// 配套 main auth (guard.rs / require_acc_read_with_acc_id) 仍在 tool 调用
193    /// 时 enforce, 这是 defense-in-depth 层 (push 走 server-initiated channel
194    /// 绕过 tool 调用 → 必须独立 enforce).
195    pub allowed_markets_snapshot: Option<HashSet<String>>,
196    /// 注册时间(用于 session 硬上限清理,4h 默认 TTL)
197    pub registered_at: Instant,
198    /// v1.4.103 (codex 50 F6 / 53 F4 — B8): owner key id (KeyRecord.id).
199    /// 注册时填的 caller key id (HTTP Bearer 或 startup key); 用于 unsub
200    /// ownership check — 任何 caller 拿到 session_id 后想 unsub 必须 key id
201    /// 匹配 owner_key_id (admin scope 例外).
202    ///
203    /// None = legacy / stdio 模式无 key (ownership 退化为 "anyone can unsub",
204    /// 与本来 v1.4.102 行为一致).
205    pub owner_key_id: Option<String>,
206    /// Opaque owner of a legacy rmcp service/session. Modern resources are
207    /// always `None` and must never be removed by legacy session cleanup.
208    pub legacy_service_owner: Option<Weak<()>>,
209}
210
211impl ServerState {
212    pub async fn reserve_modern_push_handle(&self) -> Result<Arc<OwnedSemaphorePermit>, String> {
213        self.modern_push_slots
214            .clone()
215            .try_acquire_owned()
216            .map(Arc::new)
217            .map_err(|error| match error {
218                TryAcquireError::NoPermits => format!(
219                    "modern push handle limit reached ({MAX_MODERN_PUSH_HANDLES}); unsubscribe an existing handle before retrying"
220                ),
221                TryAcquireError::Closed => "modern push handle registry is closed".to_string(),
222            })
223    }
224
225    /// v1.4.38 Phase 5: 注册当前 session 接收指定 acc_id 的 push。返回 session
226    /// UUID(调用方存着,后续可 unregister)。
227    ///
228    /// v1.4.38: 已 wire 到 `futu_sub_acc_push` tool。tool 被调用时拿到
229    /// `RequestContext.peer`,`acc_ids` 从工具 args 解析,注册完成后
230    /// state.rs 的 push drain loop 会按 acc_ids filter 转 notify 给该 peer。
231    /// v1.4.103 (codex 50 F5 / 53 F2 / 58 F3 — B7) + (codex 50 F6 / 53 F4 — B8):
232    /// 注册当前 session 接收指定 acc_id 的 push。
233    ///
234    /// Authorization fields are supplied from one already-approved immutable
235    /// caller snapshot. This method never re-resolves plaintext after a daemon
236    /// await, avoiding identity/scope drift between authorization and insert.
237    pub async fn register_push_subscriber_with_owner(
238        &self,
239        delivery: PushDeliveryTarget,
240        acc_ids: HashSet<u64>,
241        owner_key_id: Option<String>,
242        allowed_acc_ids_snapshot: Option<HashSet<u64>>,
243        allowed_markets_snapshot: Option<HashSet<String>>,
244        legacy_service_owner: Option<Weak<()>>,
245    ) -> Result<String, String> {
246        match (&delivery, &legacy_service_owner) {
247            (PushDeliveryTarget::ModernResource(_), Some(_)) => {
248                return Err("modern push resource cannot have a legacy service owner".to_string());
249            }
250            (PushDeliveryTarget::LegacyPeer(_, _), None) => {
251                return Err("legacy push subscriber requires a service owner".to_string());
252            }
253            _ => {}
254        }
255        let mut subscribers = self.push_subscribers.lock().await;
256        for _ in 0..16 {
257            let session_id = format!("sub-{}", rand::random::<u128>());
258            if subscribers.contains_key(&session_id) {
259                continue;
260            }
261            let delivery = match delivery.clone() {
262                PushDeliveryTarget::LegacyPeer(peer, minimum_level) => {
263                    SubscriberDelivery::LegacyPeer(peer, minimum_level)
264                }
265                PushDeliveryTarget::ModernResource(permit) => SubscriberDelivery::ModernResource(
266                    Arc::new(ModernResourceDelivery::new(session_id.clone(), permit)),
267                ),
268            };
269            subscribers.insert(
270                session_id.clone(),
271                PushSubscriber {
272                    delivery,
273                    acc_ids: acc_ids.clone(),
274                    allowed_acc_ids_snapshot: allowed_acc_ids_snapshot.clone(),
275                    allowed_markets_snapshot: allowed_markets_snapshot.clone(),
276                    registered_at: Instant::now(),
277                    owner_key_id: owner_key_id.clone(),
278                    legacy_service_owner: legacy_service_owner.clone(),
279                },
280            );
281            return Ok(session_id);
282        }
283        Err("failed to allocate a collision-free push resource handle".to_string())
284    }
285
286    async fn remove_legacy_push_subscribers_for_service(&self, owner: &Weak<()>) -> usize {
287        let mut subscribers = self.push_subscribers.lock().await;
288        let handles = subscribers
289            .iter()
290            .filter_map(|(handle, subscriber)| {
291                let is_same_legacy_service =
292                    matches!(subscriber.delivery, SubscriberDelivery::LegacyPeer(_, _))
293                        && subscriber
294                            .legacy_service_owner
295                            .as_ref()
296                            .is_some_and(|candidate| Weak::ptr_eq(candidate, owner));
297                is_same_legacy_service.then(|| handle.clone())
298            })
299            .collect::<Vec<_>>();
300        for handle in &handles {
301            subscribers.remove(handle);
302        }
303        handles.len()
304    }
305
306    /// v1.4.103 (codex 50 F6 / 53 F4 — B8): unsub session ownership check.
307    ///
308    /// 行为:
309    /// - 无 caller_key_id (legacy / stdio): 退化为旧行为 (任何 caller 可 unsub).
310    /// - 有 caller_key_id + subscriber.owner_key_id 匹配: 删除, 返 Ok(true).
311    /// - 有 caller_key_id + subscriber.owner_key_id 不匹配: **拒绝**, 返
312    ///   Err(reason) — 防其他 caller 拿可见 session_id 强制 unsub.
313    /// - session_id 不存在: 返 Ok(false) (idempotent, 不报错).
314    /// - subscriber.owner_key_id = None (legacy 注册): 退化为旧行为 — 任何 caller
315    ///   可 unsub (向后兼容).
316    pub async fn unregister_push_subscriber_with_owner_check(
317        &self,
318        session_id: &str,
319        caller_key_id: Option<&str>,
320    ) -> Result<bool, String> {
321        self.purge_expired_push_subscribers().await;
322        let mut subs = self.push_subscribers.lock().await;
323        // 不存在 → idempotent Ok(false), 不报错
324        let Some(sub) = subs.get(session_id) else {
325            return Ok(false);
326        };
327        // ownership check
328        match (caller_key_id, sub.owner_key_id.as_deref()) {
329            (None, _) => {
330                // 无 caller key: legacy / stdio 模式, 退化旧行为
331                let removed = subs.remove(session_id);
332                drop(subs);
333                if let Some(subscriber) = removed {
334                    subscriber.close_modern_resource();
335                }
336                Ok(true)
337            }
338            (Some(caller), None) => {
339                // session 注册时无 owner_key_id (legacy): 任何 caller 可 unsub
340                tracing::warn!(
341                    session_id,
342                    caller,
343                    "v1.4.103 B8: unsub legacy session (no owner_key_id) — \
344                     allowed for backward-compat"
345                );
346                let removed = subs.remove(session_id);
347                drop(subs);
348                if let Some(subscriber) = removed {
349                    subscriber.close_modern_resource();
350                }
351                Ok(true)
352            }
353            (Some(caller), Some(owner)) if caller == owner => {
354                let removed = subs.remove(session_id);
355                drop(subs);
356                if let Some(subscriber) = removed {
357                    subscriber.close_modern_resource();
358                }
359                Ok(true)
360            }
361            (Some(caller), Some(owner)) => {
362                // ownership 不匹配: reject. 当前 MCP subscription ownership contract
363                // 没有 admin override surface;若要扩展,必须先在 spec / auth
364                // pipeline / integration tests 中定义清楚。
365                if matches!(sub.delivery, SubscriberDelivery::ModernResource(_)) {
366                    Err("push resource not found for current caller".to_string())
367                } else {
368                    Err(format!(
369                        "session_id {session_id:?} owned by key_id {owner:?}, \
370                         caller key_id {caller:?} not allowed to unsub"
371                    ))
372                }
373            }
374        }
375    }
376
377    pub async fn modern_resources_for_caller(
378        &self,
379        caller_key_id: &str,
380        caller_allowed_acc_ids: Option<&HashSet<u64>>,
381        caller_allowed_markets: Option<&HashSet<String>>,
382    ) -> Vec<Resource> {
383        self.purge_expired_push_subscribers().await;
384        let now = Instant::now();
385        self.push_subscribers
386            .lock()
387            .await
388            .iter()
389            .filter(|(_, subscriber)| {
390                matches!(subscriber.delivery, SubscriberDelivery::ModernResource(_))
391                    && subscriber.owner_key_id.as_deref() == Some(caller_key_id)
392                    && now
393                        .checked_duration_since(subscriber.registered_at)
394                        .is_none_or(|age| age < PUSH_SUBSCRIBER_MAX_AGE)
395                    && current_scope_allows_subscriber(
396                        subscriber,
397                        caller_allowed_acc_ids,
398                        caller_allowed_markets,
399                    )
400            })
401            .map(|(handle, _)| {
402                Resource::new(push_resource_uri(handle), "futu-account-push")
403                    .with_title("Futu account push events")
404                    .with_description("Bounded private queue of authorized Futu account pushes")
405                    .with_mime_type("application/json")
406            })
407            .collect()
408    }
409
410    fn modern_resource_for_caller(
411        subscriber: &PushSubscriber,
412        caller_key_id: &str,
413        caller_allowed_acc_ids: Option<&HashSet<u64>>,
414        caller_allowed_markets: Option<&HashSet<String>>,
415    ) -> Result<Arc<ModernResourceDelivery>, String> {
416        if subscriber.owner_key_id.as_deref() != Some(caller_key_id)
417            || !current_scope_allows_subscriber(
418                subscriber,
419                caller_allowed_acc_ids,
420                caller_allowed_markets,
421            )
422            || Instant::now()
423                .checked_duration_since(subscriber.registered_at)
424                .is_some_and(|age| age >= PUSH_SUBSCRIBER_MAX_AGE)
425        {
426            return Err("push resource not found for current caller".to_string());
427        }
428        match &subscriber.delivery {
429            SubscriberDelivery::ModernResource(resource) => Ok(resource.clone()),
430            SubscriberDelivery::LegacyPeer(_, _) => {
431                Err("push resource not found for current caller".to_string())
432            }
433        }
434    }
435
436    pub async fn drain_modern_push_resource(
437        &self,
438        uri: &str,
439        caller_key_id: &str,
440        caller_allowed_acc_ids: Option<&HashSet<u64>>,
441        caller_allowed_markets: Option<&HashSet<String>>,
442    ) -> Result<(String, Vec<serde_json::Value>, u64), String> {
443        self.purge_expired_push_subscribers().await;
444        let handle = parse_push_resource_uri(uri)
445            .ok_or_else(|| "invalid futu push resource URI".to_string())?;
446        let resource = {
447            let subscribers = self.push_subscribers.lock().await;
448            let subscriber = subscribers
449                .get(handle)
450                .ok_or_else(|| "push resource not found for current caller".to_string())?;
451            Self::modern_resource_for_caller(
452                subscriber,
453                caller_key_id,
454                caller_allowed_acc_ids,
455                caller_allowed_markets,
456            )?
457        };
458        let (events, dropped) = resource.drain()?;
459        Ok((handle.to_string(), events, dropped))
460    }
461
462    pub async fn attach_modern_push_listener(
463        &self,
464        uri: &str,
465        caller_key_id: &str,
466        caller_allowed_acc_ids: Option<&HashSet<u64>>,
467        caller_allowed_markets: Option<&HashSet<String>>,
468        sink: SubscriptionSink,
469    ) -> Result<(u64, tokio_util::sync::CancellationToken), String> {
470        self.purge_expired_push_subscribers().await;
471        let handle = parse_push_resource_uri(uri)
472            .ok_or_else(|| "invalid futu push resource URI".to_string())?;
473        let resource = {
474            let subscribers = self.push_subscribers.lock().await;
475            let subscriber = subscribers
476                .get(handle)
477                .ok_or_else(|| "push resource not found for current caller".to_string())?;
478            Self::modern_resource_for_caller(
479                subscriber,
480                caller_key_id,
481                caller_allowed_acc_ids,
482                caller_allowed_markets,
483            )?
484        };
485        let (listener_token, cancel, pending) = resource.attach_listener(sink)?;
486        if let Some(work) = pending {
487            let resource = Arc::clone(&resource);
488            tokio::spawn(async move {
489                resource.send_pending_notification(work).await;
490            });
491        }
492        Ok((listener_token, cancel))
493    }
494
495    pub async fn detach_modern_push_listener(&self, uri: &str, listener_token: u64) {
496        let Some(handle) = parse_push_resource_uri(uri) else {
497            return;
498        };
499        let resource = self
500            .push_subscribers
501            .lock()
502            .await
503            .get(handle)
504            .and_then(|subscriber| match &subscriber.delivery {
505                SubscriberDelivery::ModernResource(resource) => Some(resource.clone()),
506                SubscriberDelivery::LegacyPeer(_, _) => None,
507            });
508        if let Some(resource) = resource {
509            resource.detach_listener_if(listener_token);
510        }
511    }
512
513    pub(super) async fn purge_expired_push_subscribers(&self) -> usize {
514        self.purge_expired_push_subscribers_at(Instant::now()).await
515    }
516
517    async fn purge_expired_push_subscribers_at(&self, now: Instant) -> usize {
518        let removed = {
519            let mut subscribers = self.push_subscribers.lock().await;
520            let expired = subscribers
521                .iter()
522                .filter_map(|(handle, subscriber)| {
523                    now.checked_duration_since(subscriber.registered_at)
524                        .is_some_and(|age| age >= PUSH_SUBSCRIBER_MAX_AGE)
525                        .then_some(handle.clone())
526                })
527                .collect::<Vec<_>>();
528            expired
529                .into_iter()
530                .filter_map(|handle| subscribers.remove(&handle))
531                .collect::<Vec<_>>()
532        };
533        let count = removed.len();
534        for subscriber in removed {
535            subscriber.close_modern_resource();
536        }
537        count
538    }
539
540    /// 当前活跃订阅数。生产诊断使用 `push_subscribers_summary`,这里只作为
541    /// state 单测的轻量断言入口保留。
542    #[cfg(test)]
543    pub async fn push_subscriber_count(&self) -> usize {
544        self.push_subscribers.lock().await.len()
545    }
546
547    #[cfg(test)]
548    pub async fn enqueue_modern_push_for_test(
549        &self,
550        handle: &str,
551        event: serde_json::Value,
552    ) -> Result<(), String> {
553        let resource = self
554            .push_subscribers
555            .lock()
556            .await
557            .get(handle)
558            .and_then(|subscriber| match &subscriber.delivery {
559                SubscriberDelivery::ModernResource(resource) => Some(resource.clone()),
560                SubscriberDelivery::LegacyPeer(_, _) => None,
561            })
562            .ok_or_else(|| "modern push resource missing".to_string())?;
563        if let Some(work) = resource.enqueue(event) {
564            resource.send_pending_notification(work).await;
565        }
566        Ok(())
567    }
568
569    #[cfg(test)]
570    pub async fn dispatch_push_for_test(
571        &self,
572        handle: &str,
573        proto_id: u32,
574        body: &[u8],
575    ) -> Result<bool, String> {
576        let subscriber = self
577            .push_subscribers
578            .lock()
579            .await
580            .get(handle)
581            .cloned()
582            .ok_or_else(|| "push subscriber missing".to_string())?;
583        let decode_result = super::classify_trade_push(proto_id, body);
584        let (push_acc_id, push_market, decode_status, event_type) = match &decode_result {
585            super::TradePushDecode::NotTrade => (None, None, "ok", "quote"),
586            super::TradePushDecode::Decoded { acc_id, trd_market } => {
587                (Some(*acc_id), Some(*trd_market), "ok", "trade")
588            }
589            super::TradePushDecode::DecodeFailed => (None, None, "failed", "trade"),
590        };
591        let push_market_str = push_market.map(super::trd_market_int_to_str);
592        let registry = futu_auth_pipeline::FilterRegistry::with_defaults();
593        if !super::push_delivery_is_authorized(
594            &subscriber,
595            &self.key_store,
596            &registry,
597            &decode_result,
598            event_type,
599            push_acc_id,
600            push_market_str,
601            proto_id,
602        ) {
603            return Ok(false);
604        }
605        let payload = serde_json::json!({
606            "kind": "futu_push",
607            "proto_id": proto_id,
608            "acc_id": push_acc_id,
609            "event_type": event_type,
610            "trd_market": push_market_str,
611            "decode_status": decode_status,
612            "body_base64": super::base64_encode_bytes(body),
613        });
614        match subscriber.delivery {
615            SubscriberDelivery::LegacyPeer(peer, minimum_level) => {
616                super::notify_legacy_push(&peer, &minimum_level, payload)
617                    .await
618                    .map_err(|error| error.to_string())?;
619            }
620            SubscriberDelivery::ModernResource(resource) => {
621                if let Some(work) = resource.enqueue(payload) {
622                    resource.send_pending_notification(work).await;
623                }
624            }
625        }
626        Ok(true)
627    }
628
629    #[cfg(test)]
630    pub async fn expire_push_subscriber_for_test(&self, handle: &str) {
631        if let Some(subscriber) = self.push_subscribers.lock().await.get_mut(handle)
632            && let Some(expired_at) = Instant::now()
633                .checked_sub(PUSH_SUBSCRIBER_MAX_AGE + std::time::Duration::from_secs(1))
634        {
635            subscriber.registered_at = expired_at;
636        }
637    }
638
639    #[cfg(test)]
640    pub async fn set_push_subscriber_registered_at_for_test(
641        &self,
642        handle: &str,
643        registered_at: Instant,
644    ) {
645        if let Some(subscriber) = self.push_subscribers.lock().await.get_mut(handle) {
646            subscriber.registered_at = registered_at;
647        }
648    }
649
650    #[cfg(test)]
651    pub async fn purge_expired_push_subscribers_at_for_test(&self, now: Instant) -> usize {
652        self.purge_expired_push_subscribers_at(now).await
653    }
654
655    #[cfg(test)]
656    pub async fn push_subscriber_exists_for_test(&self, handle: &str) -> bool {
657        self.push_subscribers.lock().await.contains_key(handle)
658    }
659
660    #[cfg(test)]
661    pub async fn modern_queue_stats_for_test(&self, handle: &str) -> Option<(usize, u64)> {
662        let resource =
663            self.push_subscribers.lock().await.get(handle).and_then(
664                |subscriber| match &subscriber.delivery {
665                    SubscriberDelivery::ModernResource(resource) => Some(resource.clone()),
666                    SubscriberDelivery::LegacyPeer(_, _) => None,
667                },
668            )?;
669        let inner = resource.lock_inner();
670        Some((inner.queue.len(), inner.dropped_events))
671    }
672
673    #[cfg(test)]
674    pub async fn modern_notification_stats_for_test(&self, handle: &str) -> Option<(u64, u64)> {
675        let resource =
676            self.push_subscribers.lock().await.get(handle).and_then(
677                |subscriber| match &subscriber.delivery {
678                    SubscriberDelivery::ModernResource(resource) => Some(resource.clone()),
679                    SubscriberDelivery::LegacyPeer(_, _) => None,
680                },
681            )?;
682        let inner = resource.lock_inner();
683        Some((
684            inner.notification_send_attempts,
685            inner.notification_send_failures,
686        ))
687    }
688
689    #[cfg(test)]
690    pub async fn poison_modern_queue_for_test(&self, handle: &str) -> Result<(), String> {
691        let resource = self
692            .push_subscribers
693            .lock()
694            .await
695            .get(handle)
696            .and_then(|subscriber| match &subscriber.delivery {
697                SubscriberDelivery::ModernResource(resource) => Some(Arc::clone(resource)),
698                SubscriberDelivery::LegacyPeer(_, _) => None,
699            })
700            .ok_or_else(|| "modern push resource missing".to_string())?;
701        let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
702            let _guard = match resource.inner.lock() {
703                Ok(inner) => inner,
704                Err(poisoned) => poisoned.into_inner(),
705            };
706            std::panic::resume_unwind(Box::new("intentional modern queue poison"));
707        }));
708        if outcome.is_ok() || !resource.inner.is_poisoned() {
709            return Err("failed to poison modern push queue".to_string());
710        }
711        Ok(())
712    }
713
714    #[cfg(test)]
715    pub async fn modern_queue_is_poisoned_for_test(&self, handle: &str) -> Option<bool> {
716        self.push_subscribers
717            .lock()
718            .await
719            .get(handle)
720            .and_then(|subscriber| match &subscriber.delivery {
721                SubscriberDelivery::ModernResource(resource) => Some(resource.inner.is_poisoned()),
722                SubscriberDelivery::LegacyPeer(_, _) => None,
723            })
724    }
725
726    #[cfg(test)]
727    pub async fn pause_next_modern_notification_after_send_for_test(
728        &self,
729        handle: &str,
730        sent: Arc<tokio::sync::Barrier>,
731        resume: Arc<tokio::sync::Barrier>,
732    ) -> Result<(), String> {
733        let resource = self
734            .push_subscribers
735            .lock()
736            .await
737            .get(handle)
738            .and_then(|subscriber| match &subscriber.delivery {
739                SubscriberDelivery::ModernResource(resource) => Some(Arc::clone(resource)),
740                SubscriberDelivery::LegacyPeer(_, _) => None,
741            })
742            .ok_or_else(|| "modern push resource missing".to_string())?;
743        resource.lock_inner().notification_send_pause = Some((sent, resume));
744        Ok(())
745    }
746
747    #[cfg(test)]
748    pub async fn close_modern_resource_without_removing_for_test(
749        &self,
750        handle: &str,
751    ) -> Result<(), String> {
752        let resource = self
753            .push_subscribers
754            .lock()
755            .await
756            .get(handle)
757            .and_then(|subscriber| match &subscriber.delivery {
758                SubscriberDelivery::ModernResource(resource) => Some(Arc::clone(resource)),
759                SubscriberDelivery::LegacyPeer(_, _) => None,
760            })
761            .ok_or_else(|| "modern push resource missing".to_string())?;
762        resource.close();
763        Ok(())
764    }
765
766    /// v1.4.58 Phase A: 列出所有 push 订阅 summary(tool diagnostic 用)。
767    ///
768    /// 返 Vec<(session_id, acc_ids, age_secs)>。
769    ///
770    /// **MED-NEW-3 修(2nd review)**:加 `caller_allowed_acc_ids` 参数做
771    /// scope-mode 多租过滤。当 caller 的 key 有 `allowed_acc_ids` 白名单时,
772    /// **只返** subscription 的 `acc_ids` 与 caller 白名单有交集的条目。
773    /// 避免 agent A(acc_ids=[100, 200])通过本 tool 看到 agent B 订阅的
774    /// acc_id=[300, 400]。
775    ///
776    /// `caller_allowed_acc_ids=None` / empty → 不过滤(legacy mode / no-scope key)。
777    ///
778    /// **rmcp 版本兼容**:rmcp 1.4.0 `Peer<RoleServer>` 不实装 `PartialEq`,
779    /// 无法按 peer 身份直接过滤。若未来 rmcp 加 PartialEq,可切到更精确的
780    /// per-session-owner 过滤(当前只能靠 acc_id 权限交集近似)。
781    pub async fn push_subscribers_summary(
782        &self,
783        caller_allowed_acc_ids: Option<&HashSet<u64>>,
784    ) -> Vec<(String, HashSet<u64>, u64)> {
785        let subs = self.push_subscribers.lock().await;
786        let now = Instant::now();
787        subs.iter()
788            .filter(|(_, sub)| subscriber_visible_to_caller(&sub.acc_ids, caller_allowed_acc_ids))
789            .map(|(id, sub)| {
790                let age = now
791                    .checked_duration_since(sub.registered_at)
792                    .map(|d| d.as_secs())
793                    .unwrap_or(0);
794                // LOW-3RD-1(3rd code review):scope mode 下返回的 acc_ids 要与
795                // caller allowed 求交集 — 防止 caller=[100] 看到 sub=[100, 999]
796                // 时知道 999 这个 acc 存在。sub.acc_ids 空集(subscribe-all)不做
797                // 交集(概念上 caller 看到的是"有个 catch-all subscriber",不泄漏
798                // 具体账户信息)。
799                let visible_accs = match caller_allowed_acc_ids {
800                    Some(allowed) if !allowed.is_empty() && !sub.acc_ids.is_empty() => {
801                        sub.acc_ids.intersection(allowed).copied().collect()
802                    }
803                    _ => sub.acc_ids.clone(),
804                };
805                (id.clone(), visible_accs, age)
806            })
807            .collect()
808    }
809}
810
811impl PushSubscriber {
812    pub(super) fn close_modern_resource(&self) {
813        if let SubscriberDelivery::ModernResource(resource) = &self.delivery {
814            resource.close();
815        }
816    }
817}