everscale_network/subscriber/
mod.rs1use std::borrow::Cow;
2use std::sync::Arc;
3
4use anyhow::Result;
5use tl_proto::TlRead;
6
7use crate::adnl;
8
9#[async_trait::async_trait]
11pub trait MessageSubscriber: Send + Sync {
12 async fn try_consume_custom<'a>(
13 &self,
14 ctx: SubscriberContext<'a>,
15 constructor: u32,
16 data: &'a [u8],
17 ) -> Result<bool>;
18}
19
20#[async_trait::async_trait]
22pub trait QuerySubscriber: Send + Sync {
23 async fn try_consume_query<'a>(
24 &self,
25 ctx: SubscriberContext<'a>,
26 constructor: u32,
27 query: Cow<'a, [u8]>,
28 ) -> Result<QueryConsumingResult<'a>>;
29}
30
31#[derive(Copy, Clone)]
35pub struct SubscriberContext<'a> {
36 pub adnl: &'a Arc<adnl::Node>,
37 pub local_id: &'a adnl::NodeIdShort,
38 pub peer_id: &'a adnl::NodeIdShort,
39}
40
41pub enum QueryConsumingResult<'a> {
43 Consumed(Option<Vec<u8>>),
45 Rejected(Cow<'a, [u8]>),
47}
48
49impl QueryConsumingResult<'_> {
50 pub fn consume<T>(answer: T) -> Result<Self>
51 where
52 T: tl_proto::TlWrite<Repr = tl_proto::Boxed>,
53 {
54 Ok(Self::Consumed(Some(tl_proto::serialize(answer))))
55 }
56}
57
58pub(crate) async fn process_query<'a>(
59 ctx: SubscriberContext<'a>,
60 subscribers: &[Arc<dyn QuerySubscriber>],
61 mut query: Cow<'_, [u8]>,
62) -> Result<QueryProcessingResult<Vec<u8>>> {
63 let constructor = u32::read_from(&query, &mut 0)?;
64
65 for subscriber in subscribers {
66 query = match subscriber
67 .try_consume_query(ctx, constructor, query)
68 .await?
69 {
70 QueryConsumingResult::Consumed(answer) => {
71 return Ok(QueryProcessingResult::Processed(answer))
72 }
73 QueryConsumingResult::Rejected(query) => query,
74 };
75 }
76
77 Ok(QueryProcessingResult::Rejected)
78}
79
80pub(crate) enum QueryProcessingResult<T> {
81 Processed(Option<T>),
82 Rejected,
83}