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 pub fn subscribe_channels(
31 selector: ChannelSelector,
32 ) -> cynic::StreamingOperation<SubscribeChannels, ChannelsVariables> {
33 SubscribeChannels::build(ChannelsVariables::from(selector))
34 }
35
36 pub fn subscribe_accounts(
38 selector: AccountSelector,
39 ) -> cynic::StreamingOperation<SubscribeAccounts, AccountVariables> {
40 SubscribeAccounts::build(AccountVariables::from(selector))
41 }
42
43 pub fn subscribe_graph() -> cynic::StreamingOperation<SubscribeGraph, ()> {
45 SubscribeGraph::build(())
46 }
47
48 pub fn subscribe_ticket_params() -> cynic::StreamingOperation<SubscribeTicketParams, ()> {
50 SubscribeTicketParams::build(())
51 }
52
53 pub fn subscribe_health() -> cynic::StreamingOperation<SubscribeHealth, ()> {
55 SubscribeHealth::build(())
56 }
57
58 pub fn subscribe_safe_deployments() -> cynic::StreamingOperation<SubscribeSafeDeployment, ()> {
60 SubscribeSafeDeployment::build(())
61 }
62
63 pub fn subscribe_services(
65 selector: ServiceSelector,
66 ) -> cynic::StreamingOperation<SubscribeServices, ServiceVariables> {
67 SubscribeServices::build(ServiceVariables::from(selector))
68 }
69
70 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 pub fn subscribe_service_registry_config() -> cynic::StreamingOperation<SubscribeServiceRegistryConfig, ()> {
79 SubscribeServiceRegistryConfig::build(())
80 }
81
82 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 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 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 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}