Skip to main content

blokli_client/client/
subscriptions.rs

1use cynic::SubscriptionBuilder;
2use futures::{Stream, TryStreamExt};
3
4use super::{BlokliClient, GraphQlQueries};
5use crate::api::{
6    AccountSelector, BlokliSubscriptionClient, ChannelSelector, Result, ServiceSelector, ServiceTypeId, TicketSelector,
7    TxId,
8    internal::{
9        AccountVariables, ChannelsVariables, ServiceTypeVariables, ServiceVariables, SubscribeAccounts,
10        SubscribeChannels, SubscribeGraph, SubscribeHealth, SubscribeSafeDeployment, SubscribeServiceRegistryConfig,
11        SubscribeServiceTypes, SubscribeServices, SubscribeTicketParams, SubscribeTicketRedeemed,
12        TicketRedeemedVariables,
13    },
14    types::{
15        Account, Channel, OpenedChannelsGraphEntry, ReadinessState, RedeemTicketDetails, Safe, ServiceRegistryConfig,
16        ServiceTypeUpdate, ServiceUpdate, TicketParameters, Transaction,
17    },
18};
19#[cfg(feature = "curvy")]
20use crate::api::{
21    internal::{
22        CurvyEventSubscriptionVariables, SubscribeCurvyCommittedNote, SubscribeCurvyCommittedNullifier,
23        SubscribeCurvyPendingNote,
24    },
25    types::{CurvyCommittedNote, CurvyCommittedNullifier, CurvyPendingNote, Uint64},
26};
27
28impl GraphQlQueries {
29    /// `SubscribeChannels` subscription GraphQL query.
30    pub fn subscribe_channels(
31        selector: ChannelSelector,
32    ) -> cynic::StreamingOperation<SubscribeChannels, ChannelsVariables> {
33        SubscribeChannels::build(ChannelsVariables::from(selector))
34    }
35
36    /// `SubscribeAccounts` subscription GraphQL query.
37    pub fn subscribe_accounts(
38        selector: AccountSelector,
39    ) -> cynic::StreamingOperation<SubscribeAccounts, AccountVariables> {
40        SubscribeAccounts::build(AccountVariables::from(selector))
41    }
42
43    /// `SubscribeGraph` subscription GraphQL query.
44    pub fn subscribe_graph() -> cynic::StreamingOperation<SubscribeGraph, ()> {
45        SubscribeGraph::build(())
46    }
47
48    /// `SubscribeTicketParams` subscription GraphQL query.
49    pub fn subscribe_ticket_params() -> cynic::StreamingOperation<SubscribeTicketParams, ()> {
50        SubscribeTicketParams::build(())
51    }
52
53    /// `SubscribeHealth` subscription GraphQL query.
54    pub fn subscribe_health() -> cynic::StreamingOperation<SubscribeHealth, ()> {
55        SubscribeHealth::build(())
56    }
57
58    /// `SubscribeSafeDeployment` subscription GraphQL query.
59    pub fn subscribe_safe_deployments() -> cynic::StreamingOperation<SubscribeSafeDeployment, ()> {
60        SubscribeSafeDeployment::build(())
61    }
62
63    /// `SubscribeServices` subscription GraphQL query.
64    pub fn subscribe_services(
65        selector: ServiceSelector,
66    ) -> cynic::StreamingOperation<SubscribeServices, ServiceVariables> {
67        SubscribeServices::build(ServiceVariables::from(selector))
68    }
69
70    /// `SubscribeServiceTypes` subscription GraphQL query.
71    pub fn subscribe_service_types(
72        service_type: Option<ServiceTypeId>,
73    ) -> cynic::StreamingOperation<SubscribeServiceTypes, ServiceTypeVariables> {
74        SubscribeServiceTypes::build(ServiceTypeVariables::from(service_type))
75    }
76
77    /// `SubscribeServiceRegistryConfig` subscription GraphQL query.
78    pub fn subscribe_service_registry_config() -> cynic::StreamingOperation<SubscribeServiceRegistryConfig, ()> {
79        SubscribeServiceRegistryConfig::build(())
80    }
81
82    /// `SubscribeTicketRedeemed` subscription GraphQL query.
83    pub fn subscribe_ticket_redeemed(
84        selector: TicketSelector,
85    ) -> cynic::StreamingOperation<SubscribeTicketRedeemed, TicketRedeemedVariables> {
86        SubscribeTicketRedeemed::build(TicketRedeemedVariables::from(selector))
87    }
88
89    #[cfg(feature = "curvy")]
90    /// Pending Curvy note subscription used for local ownership detection.
91    pub fn subscribe_curvy_pending_notes(
92        from_block: Option<u64>,
93    ) -> cynic::StreamingOperation<SubscribeCurvyPendingNote, CurvyEventSubscriptionVariables> {
94        SubscribeCurvyPendingNote::build(CurvyEventSubscriptionVariables {
95            from_block: from_block.map(|block| Uint64(block.to_string())),
96        })
97    }
98
99    #[cfg(feature = "curvy")]
100    /// Committed Curvy note subscription used for owned-note correlation.
101    pub fn subscribe_curvy_committed_notes(
102        from_block: Option<u64>,
103    ) -> cynic::StreamingOperation<SubscribeCurvyCommittedNote, CurvyEventSubscriptionVariables> {
104        SubscribeCurvyCommittedNote::build(CurvyEventSubscriptionVariables {
105            from_block: from_block.map(|block| Uint64(block.to_string())),
106        })
107    }
108
109    #[cfg(feature = "curvy")]
110    /// Committed Curvy nullifier subscription.
111    pub fn subscribe_curvy_committed_nullifiers(
112        from_block: Option<u64>,
113    ) -> cynic::StreamingOperation<SubscribeCurvyCommittedNullifier, CurvyEventSubscriptionVariables> {
114        SubscribeCurvyCommittedNullifier::build(CurvyEventSubscriptionVariables {
115            from_block: from_block.map(|block| Uint64(block.to_string())),
116        })
117    }
118}
119
120impl BlokliSubscriptionClient for BlokliClient {
121    #[tracing::instrument(level = "debug", skip(self), fields(?selector))]
122    fn subscribe_channels(&self, selector: ChannelSelector) -> Result<impl Stream<Item = Result<Channel>> + Send> {
123        Ok(self
124            .build_subscription_stream(GraphQlQueries::subscribe_channels(selector))?
125            .try_filter_map(|item| futures::future::ok(Some(item.channel_updated))))
126    }
127
128    #[tracing::instrument(level = "debug", skip(self), fields(?selector))]
129    fn subscribe_accounts(&self, selector: AccountSelector) -> Result<impl Stream<Item = Result<Account>> + Send> {
130        Ok(self
131            .build_subscription_stream(GraphQlQueries::subscribe_accounts(selector))?
132            .try_filter_map(|item| futures::future::ok(Some(item.account_updated))))
133    }
134
135    #[tracing::instrument(level = "debug", skip(self))]
136    fn subscribe_graph(&self) -> Result<impl Stream<Item = Result<OpenedChannelsGraphEntry>> + Send> {
137        Ok(self
138            .build_subscription_stream(GraphQlQueries::subscribe_graph())?
139            .try_filter_map(|item| futures::future::ok(Some(item.opened_channel_graph_updated))))
140    }
141
142    #[tracing::instrument(level = "debug", skip(self))]
143    fn subscribe_ticket_params(&self) -> Result<impl Stream<Item = Result<TicketParameters>> + Send> {
144        Ok(self
145            .build_subscription_stream(GraphQlQueries::subscribe_ticket_params())?
146            .try_filter_map(|item| futures::future::ok(Some(item.ticket_parameters_updated))))
147    }
148
149    #[tracing::instrument(level = "debug", skip(self))]
150    fn subscribe_health(&self) -> Result<impl Stream<Item = Result<ReadinessState>> + Send> {
151        Ok(self
152            .build_subscription_stream(GraphQlQueries::subscribe_health())?
153            .try_filter_map(|item| futures::future::ok(Some(item.health))))
154    }
155
156    #[tracing::instrument(level = "debug", skip(self))]
157    fn subscribe_safe_deployments(&self) -> Result<impl Stream<Item = Result<Safe>> + Send> {
158        Ok(self
159            .build_subscription_stream(GraphQlQueries::subscribe_safe_deployments())?
160            .try_filter_map(|item| futures::future::ok(Some(item.safe_deployed))))
161    }
162
163    #[tracing::instrument(level = "debug", skip(self), fields(?selector))]
164    fn subscribe_services(
165        &self,
166        selector: ServiceSelector,
167    ) -> Result<impl Stream<Item = Result<ServiceUpdate>> + Send> {
168        Ok(self
169            .build_subscription_stream(GraphQlQueries::subscribe_services(selector))?
170            .try_filter_map(|item| futures::future::ok(Some(item.service_updated))))
171    }
172
173    #[tracing::instrument(level = "debug", skip(self))]
174    fn subscribe_service_types(
175        &self,
176        service_type: Option<ServiceTypeId>,
177    ) -> Result<impl Stream<Item = Result<ServiceTypeUpdate>> + Send> {
178        Ok(self
179            .build_subscription_stream(GraphQlQueries::subscribe_service_types(service_type))?
180            .try_filter_map(|item| futures::future::ok(Some(item.service_type_updated))))
181    }
182
183    #[tracing::instrument(level = "debug", skip(self))]
184    fn subscribe_service_registry_config(
185        &self,
186    ) -> Result<impl Stream<Item = Result<ServiceRegistryConfig>> + Send + 'static> {
187        Ok(self
188            .build_subscription_stream(GraphQlQueries::subscribe_service_registry_config())?
189            .try_filter_map(|item| futures::future::ok(Some(item.service_registry_config_updated))))
190    }
191
192    #[tracing::instrument(level = "debug", skip(self))]
193    fn subscribe_track_transaction(&self, tx_id: TxId) -> Result<impl Stream<Item = Result<Transaction>> + Send> {
194        Ok(self
195            .build_subscription_stream(GraphQlQueries::subscribe_track_transaction(tx_id))?
196            .try_filter_map(|item| futures::future::ok(Some(item.transaction_updated))))
197    }
198
199    #[tracing::instrument(level = "debug", skip(self), fields(?selector))]
200    fn subscribe_ticket_redeemed(
201        &self,
202        selector: TicketSelector,
203    ) -> Result<impl futures::Stream<Item = Result<RedeemTicketDetails>> + Send> {
204        Ok(self
205            .build_subscription_stream(GraphQlQueries::subscribe_ticket_redeemed(selector))?
206            .try_filter_map(|item| futures::future::ok(Some(item.ticket_redeemed))))
207    }
208
209    #[cfg(feature = "curvy")]
210    #[tracing::instrument(level = "debug", skip(self))]
211    fn subscribe_curvy_pending_notes(
212        &self,
213        from_block: Option<u64>,
214    ) -> Result<impl futures::Stream<Item = Result<CurvyPendingNote>> + Send> {
215        Ok(self
216            .build_subscription_stream(GraphQlQueries::subscribe_curvy_pending_notes(from_block))?
217            .map_ok(|item| item.curvy_pending_note))
218    }
219
220    #[cfg(feature = "curvy")]
221    #[tracing::instrument(level = "debug", skip(self))]
222    fn subscribe_curvy_committed_notes(
223        &self,
224        from_block: Option<u64>,
225    ) -> Result<impl futures::Stream<Item = Result<CurvyCommittedNote>> + Send> {
226        Ok(self
227            .build_subscription_stream(GraphQlQueries::subscribe_curvy_committed_notes(from_block))?
228            .map_ok(|item| item.curvy_committed_note))
229    }
230
231    #[cfg(feature = "curvy")]
232    #[tracing::instrument(level = "debug", skip(self))]
233    fn subscribe_curvy_committed_nullifiers(
234        &self,
235        from_block: Option<u64>,
236    ) -> Result<impl futures::Stream<Item = Result<CurvyCommittedNullifier>> + Send> {
237        Ok(self
238            .build_subscription_stream(GraphQlQueries::subscribe_curvy_committed_nullifiers(from_block))?
239            .map_ok(|item| item.curvy_committed_nullifier))
240    }
241}
242
243#[cfg(all(test, feature = "curvy"))]
244mod tests {
245    use serde_json::json;
246
247    use super::GraphQlQueries;
248    use crate::api::types::CurvyEventPosition;
249
250    #[test]
251    fn curvy_pending_note_subscription_serializes_from_block() {
252        let operation = GraphQlQueries::subscribe_curvy_pending_notes(Some(9));
253
254        let serialized = serde_json::to_value(operation).expect("subscription operation should serialize");
255
256        assert_eq!(
257            serialized["variables"],
258            json!({
259                "fromBlock": "9",
260            })
261        );
262        assert!(
263            serialized["query"]
264                .as_str()
265                .is_some_and(|query| query.contains("curvyPendingNote"))
266        );
267    }
268
269    #[test]
270    fn event_position_type_is_publicly_constructible() {
271        let _position: Option<CurvyEventPosition> = None;
272    }
273}