Skip to main content

everscale_network/subscriber/
mod.rs

1use std::borrow::Cow;
2use std::sync::Arc;
3
4use anyhow::Result;
5use tl_proto::TlRead;
6
7use crate::adnl;
8
9/// ADNL custom messages subscriber
10#[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/// ADNL, RLDP or overlay queries subscriber
21#[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/// Message or query context.
32///
33/// See [`MessageSubscriber::try_consume_custom`] and [`QuerySubscriber::try_consume_query`]
34#[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
41/// Subscriber response for consumed query
42pub enum QueryConsumingResult<'a> {
43    /// Query is accepted and processed
44    Consumed(Option<Vec<u8>>),
45    /// Query rejected and will be processed by the next subscriber
46    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}