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
27pub(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 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 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 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#[derive(Clone)]
168pub(super) struct PushSubscriber {
169 pub delivery: SubscriberDelivery,
170 pub acc_ids: HashSet<u64>,
172 pub allowed_acc_ids_snapshot: Option<HashSet<u64>>,
182 pub allowed_markets_snapshot: Option<HashSet<String>>,
196 pub registered_at: Instant,
198 pub owner_key_id: Option<String>,
206 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 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 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 let Some(sub) = subs.get(session_id) else {
325 return Ok(false);
326 };
327 match (caller_key_id, sub.owner_key_id.as_deref()) {
329 (None, _) => {
330 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 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 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 #[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 ®istry,
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 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 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}