Skip to main content

casper_rust_wasm_sdk/sdk/sse/watcher/
mod.rs

1#[cfg(test)]
2pub(crate) mod deploy_mock;
3#[cfg(test)]
4pub(crate) mod transaction_mock;
5
6use crate::sdk::sse::framing::{extract_frames, url_with_start_from};
7use crate::SDK;
8use chrono::{Duration, Utc};
9use futures_util::StreamExt;
10#[cfg(all(feature = "js", target_arch = "wasm32"))]
11use gloo_utils::format::JsValueSerdeExt;
12#[cfg(all(feature = "js", target_arch = "wasm32"))]
13use js_sys::Promise;
14use serde::{Deserialize, Deserializer, Serialize};
15use serde_json::Value;
16use std::{
17    fmt,
18    sync::{Arc, Mutex},
19};
20#[cfg(feature = "js")]
21use wasm_bindgen::prelude::*;
22#[cfg(all(feature = "js", target_arch = "wasm32"))]
23use wasm_bindgen_futures::future_to_promise;
24
25const DEFAULT_TIMEOUT_MS: u64 = 60000;
26
27/// HTTP client for native SSE. Disables idle keep-alive so a dropped tokio
28/// runtime cannot leave pooled connections (hyperium/hyper#2136).
29#[cfg(not(target_arch = "wasm32"))]
30fn sse_http_client() -> Result<reqwest::Client, String> {
31    reqwest::Client::builder()
32        .pool_max_idle_per_host(0)
33        .build()
34        .map_err(|e| format!("SSE HTTP client build failed: {e}"))
35}
36
37#[cfg(target_arch = "wasm32")]
38fn sse_http_client() -> Result<reqwest::Client, String> {
39    reqwest::Client::builder()
40        .build()
41        .map_err(|e| format!("SSE HTTP client build failed: {e}"))
42}
43
44impl SDK {
45    /// Creates a new Watcher instance to watch deploys.
46    /// Legacy alias
47    ///
48    /// # Arguments
49    ///
50    /// * `events_url` - The URL to monitor for transaction events.
51    /// * `timeout_duration` - An optional timeout duration in seconds.
52    ///
53    /// # Returns
54    ///
55    /// A `Watcher` instance.
56    #[deprecated(note = "prefer 'watch_transaction'")]
57    pub fn watch_deploy(&self, events_url: &str, timeout_duration: Option<u64>) -> Watcher {
58        Watcher::new(events_url.to_string(), timeout_duration)
59    }
60
61    /// Creates a new Watcher instance to watch deploys.
62    ///
63    /// # Arguments
64    ///
65    /// * `events_url` - The URL to monitor for transaction events.
66    /// * `timeout_duration` - An optional timeout duration in seconds.
67    ///
68    /// # Returns
69    ///
70    /// A `Watcher` instance.
71    pub fn watch_transaction(&self, events_url: &str, timeout_duration: Option<u64>) -> Watcher {
72        Watcher::new(events_url.to_string(), timeout_duration)
73    }
74
75    /// Waits for a deploy event to be processed asynchronously.
76    /// Legacy alias
77    ///
78    /// # Arguments
79    ///
80    /// * `events_url` - The URL to monitor for transaction events.
81    /// * `deploy_hash` - The deploy hash to wait for.
82    /// * `timeout_duration` - An optional timeout duration in milliseconds.
83    ///
84    /// # Returns
85    ///
86    /// A `Result` containing either the processed `EventParseResult` or an error message.
87    #[deprecated(note = "prefer 'wait_transaction' with transaction")]
88    pub async fn wait_deploy(
89        &self,
90        events_url: &str,
91        deploy_hash: &str,
92        timeout_duration: Option<u64>,
93    ) -> Result<EventParseResult, String> {
94        Self::wait_transaction_internal(
95            events_url.to_string(),
96            deploy_hash.to_string(),
97            timeout_duration,
98        )
99        .await
100    }
101
102    /// Alias for wait_deploy Waits for a deploy event to be processed asynchronously.
103    ///
104    /// # Arguments
105    ///
106    /// * `events_url` - The URL to monitor for transaction events.
107    /// * `target_hash` - The transaction hash to wait for.
108    /// * `timeout_duration` - An optional timeout duration in milliseconds.
109    ///
110    /// # Returns
111    ///
112    /// A `Result` containing either the processed `EventParseResult` or an error message
113    pub async fn wait_transaction(
114        &self,
115        events_url: &str,
116        target_hash: &str,
117        timeout_duration: Option<u64>,
118    ) -> Result<EventParseResult, String> {
119        Self::wait_transaction_internal(
120            events_url.to_string(),
121            target_hash.to_string(),
122            timeout_duration,
123        )
124        .await
125    }
126
127    /// Internal function to wait for a deploy event.
128    ///
129    /// # Arguments
130    ///
131    /// * `events_url` - The URL to monitor for transaction events.
132    /// * `target_hash` - The transaction hash to wait for.
133    /// * `timeout_duration` - An optional timeout duration in milliseconds.
134    ///
135    /// # Returns
136    ///
137    /// A `Result` containing either the processed `EventParseResult` or an error message.
138    async fn wait_transaction_internal(
139        events_url: String,
140        target_hash: String,
141        timeout_duration: Option<u64>,
142    ) -> Result<EventParseResult, String> {
143        let watcher = Watcher::new(events_url, timeout_duration);
144        let result = watcher.start_internal(Some(target_hash)).await;
145        match result {
146            Some(event_parse_results) => {
147                if let Some(event_parse_result) = event_parse_results.first() {
148                    return Ok(event_parse_result.clone());
149                }
150                Err("No first event result".to_string())
151            }
152            None => Err("No event result found".to_string()),
153        }
154    }
155}
156
157#[cfg_attr(feature = "js", wasm_bindgen)]
158impl SDK {
159    /// Creates a new Watcher instance to watch deploys (JavaScript-friendly).
160    /// Legacy alias
161    ///
162    /// # Arguments
163    ///
164    /// * `events_url` - The URL to monitor for transaction events.
165    /// * `timeout_duration` - An optional timeout duration in seconds.
166    ///
167    /// # Returns
168    ///
169    /// A `Watcher` instance.
170    #[cfg(all(feature = "js", target_arch = "wasm32"))]
171    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "watchDeploy"))]
172    #[deprecated(note = "prefer 'watchTransaction'")]
173    #[allow(deprecated)]
174    pub fn watch_deploy_js_alias(
175        &self,
176        events_url: &str,
177        timeout_duration: Option<u32>,
178    ) -> Watcher {
179        self.watch_deploy(events_url, timeout_duration.map(Into::into))
180    }
181
182    /// Creates a new Watcher instance to watch deploys (JavaScript-friendly).
183    ///
184    /// # Arguments
185    ///
186    /// * `events_url` - The URL to monitor for transaction events.
187    /// * `timeout_duration` - An optional timeout duration in seconds.
188    ///
189    /// # Returns
190    ///
191    /// A `Watcher` instance.
192    #[cfg(all(feature = "js", target_arch = "wasm32"))]
193    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "watchTransaction"))]
194    pub fn watch_transaction_js_alias(
195        &self,
196        events_url: &str,
197        timeout_duration: Option<u32>,
198    ) -> Watcher {
199        self.watch_transaction(events_url, timeout_duration.map(Into::into))
200    }
201
202    /// Waits for a deploy event to be processed asynchronously (JavaScript-friendly).
203    /// Legacy alias
204    ///
205    /// # Arguments
206    ///
207    /// * `events_url` - The URL to monitor for transaction events.
208    /// * `deploy_hash` - The deploy hash to wait for.
209    /// * `timeout_duration` - An optional timeout duration in seconds.
210    ///
211    /// # Returns
212    ///
213    /// A JavaScript `Promise` resolving to either the processed `EventParseResult` or an error message.
214    #[cfg(all(feature = "js", target_arch = "wasm32"))]
215    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "waitDeploy"))]
216    #[deprecated(note = "prefer 'waitTransaction' with transaction")]
217    #[allow(deprecated)]
218    pub async fn wait_deploy_js_alias(
219        &self,
220        events_url: &str,
221        deploy_hash: &str,
222        timeout_duration: Option<u32>,
223    ) -> Promise {
224        self.wait_transaction_js_alias(events_url, deploy_hash, timeout_duration)
225            .await
226    }
227
228    /// Waits for a deploy event to be processed asynchronously (JavaScript-friendly).
229    ///
230    /// # Arguments
231    ///
232    /// * `events_url` - The URL to monitor for transaction events.
233    /// * `target_hash` - The transaction hash to wait for.
234    /// * `timeout_duration` - An optional timeout duration in seconds.
235    ///
236    /// # Returns
237    ///
238    /// A JavaScript `Promise` resolving to either the processed `EventParseResult` or an error message.
239    #[cfg(all(feature = "js", target_arch = "wasm32"))]
240    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "waitTransaction"))]
241    pub async fn wait_transaction_js_alias(
242        &self,
243        events_url: &str,
244        target_hash: &str,
245        timeout_duration: Option<u32>,
246    ) -> Promise {
247        let events_url = events_url.to_string();
248        let target_hash = target_hash.to_string();
249        let future = async move {
250            let result = Self::wait_transaction_internal(
251                events_url,
252                target_hash,
253                timeout_duration.map(Into::into),
254            )
255            .await;
256            match result {
257                Ok(event_parse_result) => JsValue::from_serde(&event_parse_result)
258                    .map_err(|err| JsValue::from_str(&format!("{err}"))),
259                Err(err) => Err(JsValue::from_str(&err)),
260            }
261        };
262
263        future_to_promise(future)
264    }
265}
266
267/// Represents a deploy watcher responsible for monitoring transaction events.
268///
269/// This struct allows clients to subscribe to transaction events, start watching for events,
270/// or wait for an event and handle the received deploy event data.
271///
272/// # Fields
273///
274/// * `events_url` - The URL for transaction events.
275/// * `subscriptions` - Vector containing deploy subscriptions.
276/// * `active` - Reference-counted cell indicating whether the deploy watcher is active.
277/// * `timeout_duration` - Duration representing the optional timeout for watching events.
278#[derive(Clone)]
279#[cfg_attr(feature = "js", wasm_bindgen)]
280pub struct Watcher {
281    events_url: String,
282    subscriptions: Vec<Subscription>,
283    active: Arc<Mutex<bool>>,
284    timeout_duration: Duration,
285}
286
287#[cfg_attr(feature = "js", wasm_bindgen)]
288impl Watcher {
289    /// Creates a new `Watcher` instance.
290    ///
291    /// # Arguments
292    ///
293    /// * `events_url` - The URL for transaction events.
294    /// * `timeout_duration` - Optional duration in milliseconds for watching events. If not provided,
295    ///   a default timeout of 60,000 milliseconds (1 minute) is used.
296    ///
297    /// # Returns
298    ///
299    /// A new `Watcher` instance.
300    #[cfg_attr(feature = "js", wasm_bindgen(constructor))]
301    pub fn new(events_url: String, timeout_duration: Option<u64>) -> Self {
302        let timeout_duration = Duration::try_milliseconds(
303            timeout_duration
304                .unwrap_or(DEFAULT_TIMEOUT_MS)
305                .try_into()
306                .unwrap(),
307        )
308        .unwrap_or_default();
309
310        Watcher {
311            events_url,
312            subscriptions: Vec::new(),
313            active: Arc::new(Mutex::new(true)),
314            timeout_duration,
315        }
316    }
317
318    /// Subscribes to transaction events.
319    ///
320    /// # Arguments
321    ///
322    /// * `subscriptions` - Vector of deploy subscriptions to be added.
323    ///
324    /// # Returns
325    ///
326    /// Result indicating success or an error message.
327    #[cfg(all(feature = "js", target_arch = "wasm32"))]
328    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "subscribe"))]
329    pub fn subscribe_js_alias(&mut self, subscriptions: Vec<Subscription>) -> Result<(), String> {
330        self.subscribe(subscriptions)
331    }
332
333    /// Unsubscribes from transaction events based on the provided transaction hash.
334    ///
335    /// # Arguments
336    ///
337    /// * `transaction_hash` - The transaction hash to unsubscribe.
338    ///
339    /// This method removes the deploy subscription associated with the provided transaction hash.
340    #[cfg_attr(feature = "js", wasm_bindgen)]
341    pub fn unsubscribe(&mut self, target_hash: String) {
342        self.subscriptions.retain(|s| s.target_hash != target_hash);
343    }
344
345    /// Starts watching for transaction events (JavaScript-friendly).
346    ///
347    /// # Returns
348    ///
349    /// Result containing the serialized transaction events data or an error message.
350    #[cfg(all(feature = "js", target_arch = "wasm32"))]
351    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "start"))]
352    pub async fn start_js_alias(&self) -> Result<JsValue, JsError> {
353        let result = match self.start_internal(None).await {
354            Some(res) => res,
355            None => return Ok(JsValue::NULL),
356        };
357
358        let serialized = JsValue::from_serde(&result)
359            .map_err(|err| JsError::new(&format!("Error serializing events: {err:?}")))?;
360
361        Ok(serialized)
362    }
363
364    /// Stops watching for transaction events.
365    ///
366    /// This method sets the deploy watcher as inactive and stops the event listener if it exists.
367    #[cfg_attr(feature = "js", wasm_bindgen)]
368    pub fn stop(&self) {
369        {
370            let mut active = self.active.lock().unwrap();
371            *active = false;
372        }
373    }
374}
375
376impl Watcher {
377    /// Asynchronously starts watching for transaction events and execute callback handler functions from deploy subscriptions
378    ///
379    /// # Returns
380    ///
381    /// An `Option` containing the serialized deploy event data or `None` if no events are received.
382    pub async fn start(&self) -> Option<Vec<EventParseResult>> {
383        self.start_internal(None).await
384    }
385
386    /// Asynchronously starts watching for transaction events
387    ///
388    /// # Arguments
389    ///
390    /// * `transaction_hash` - Optional transaction hash to directly return processed event. If provided, it directly returns matched events without executing callback handler functions from deploy subscriptions. If `None`, it executes callback handler functions from deploy subscriptions.
391    ///
392    /// # Returns
393    ///
394    /// An `Option` containing the serialized deploy event data or `None` if no events are received.
395    async fn start_internal(&self, target_hash: Option<String>) -> Option<Vec<EventParseResult>> {
396        {
397            let mut active = self.active.lock().unwrap();
398            *active = true;
399        }
400
401        let client = match sse_http_client() {
402            Ok(client) => client,
403            Err(err) => {
404                let event_parse_result = EventParseResult {
405                    err: Some(err),
406                    body: None,
407                };
408                return Some(vec![event_parse_result]);
409            }
410        };
411        // Replay from event id 0 so a TransactionProcessed already emitted
412        // before connect is still visible (parity with SSEClient).
413        let url = url_with_start_from(&self.events_url, Some(0));
414
415        // Clippy false positive until rust-lang/rust-clippy#11034 is fixed.
416        #[allow(clippy::arc_with_non_send_sync)]
417        let watcher = Arc::new(Mutex::new(self.clone()));
418
419        let start_time = Utc::now();
420        let timeout_duration = self.timeout_duration;
421
422        let response = match client.get(&url).send().await {
423            Ok(res) => res,
424            Err(err) => {
425                let err = err.to_string();
426                let event_parse_result = EventParseResult {
427                    err: Some(err.to_string()),
428                    body: None,
429                };
430                return Some([event_parse_result].to_vec());
431            }
432        };
433
434        if !response.status().is_success() {
435            let event_parse_result = EventParseResult {
436                err: Some("Failed to fetch stream".to_string()),
437                body: None,
438            };
439            return Some([event_parse_result].to_vec());
440        }
441
442        let buffer_size = 1;
443        let mut buffer = Vec::with_capacity(buffer_size);
444        let mut bytes_stream = response.bytes_stream();
445
446        loop {
447            let chunk = {
448                #[cfg(not(target_arch = "wasm32"))]
449                {
450                    let elapsed = Utc::now() - start_time;
451                    let remaining = timeout_duration - elapsed;
452                    if remaining <= Duration::zero() {
453                        return Some(
454                            [EventParseResult {
455                                err: Some("Timeout expired".to_string()),
456                                body: None,
457                            }]
458                            .to_vec(),
459                        );
460                    }
461                    let wait = std::time::Duration::from_millis(
462                        remaining.num_milliseconds().max(0) as u64,
463                    );
464                    match tokio::time::timeout(wait, bytes_stream.next()).await {
465                        Ok(chunk) => chunk,
466                        Err(_) => {
467                            return Some(
468                                [EventParseResult {
469                                    err: Some("Timeout expired".to_string()),
470                                    body: None,
471                                }]
472                                .to_vec(),
473                            );
474                        }
475                    }
476                }
477                #[cfg(target_arch = "wasm32")]
478                {
479                    bytes_stream.next().await
480                }
481            };
482
483            let Some(chunk) = chunk else {
484                break;
485            };
486
487            match chunk {
488                Ok(bytes) => {
489                    let this_clone = Arc::clone(&watcher);
490                    if !*this_clone.lock().unwrap().active.lock().unwrap() {
491                        return None;
492                    }
493
494                    if Utc::now() - start_time >= timeout_duration {
495                        let event_parse_result = EventParseResult {
496                            err: Some("Timeout expired".to_string()),
497                            body: None,
498                        };
499                        return Some([event_parse_result].to_vec());
500                    }
501
502                    buffer.extend_from_slice(&bytes);
503
504                    while let Some(index) = buffer.iter().position(|&b| b == b'\n') {
505                        let message = buffer.drain(..=index).collect::<Vec<_>>();
506
507                        if let Ok(message) = std::str::from_utf8(&message) {
508                            let watcher_guard = this_clone.lock().unwrap();
509                            let result = watcher_guard
510                                .clone()
511                                .process_events(message, target_hash.as_deref());
512                            match result {
513                                Some(event_parse_result) => return Some(event_parse_result),
514                                None => {
515                                    continue;
516                                }
517                            };
518                        } else {
519                            let event_parse_result = EventParseResult {
520                                err: Some("Error decoding UTF-8 data".to_string()),
521                                body: None,
522                            };
523                            return Some([event_parse_result].to_vec());
524                        }
525                    }
526                }
527                Err(err) => {
528                    let event_parse_result = EventParseResult {
529                        err: Some(format!("Error reading chunk: {err}")),
530                        body: None,
531                    };
532                    return Some([event_parse_result].to_vec());
533                }
534            }
535        }
536        None
537    }
538
539    /// Subscribes to transaction events.
540    ///
541    /// # Arguments
542    ///
543    /// * `subscriptions` - Vector of subscriptions to be added.
544    ///
545    /// # Returns
546    ///
547    /// Result indicating success or an error message.
548    pub fn subscribe(&mut self, subscriptions: Vec<Subscription>) -> Result<(), String> {
549        for new_subscription in &subscriptions {
550            if self
551                .subscriptions
552                .iter()
553                .any(|s| s.target_hash == new_subscription.target_hash)
554            {
555                return Err(String::from("Already subscribed to this event"));
556            }
557        }
558        self.subscriptions.extend(subscriptions);
559        Ok(())
560    }
561
562    /// Processes events received from the stream and notifies subscribers.
563    ///
564    /// # Arguments
565    ///
566    /// * `message` - The raw message received from the event stream.
567    /// * `target_transaction_hash` - Optional transaction hash to directly return. If provided, it directly returns matched events without executing callback handler functions from subscriptions. If `None`, it executes callback handler functions from subscriptions.
568    ///
569    /// # Returns
570    ///
571    /// An `Option` containing the serialized transaction/deploy event data or `None` if an error occurs.
572    fn process_events(
573        mut self,
574        message: &str,
575        target_hash: Option<&str>,
576    ) -> Option<Vec<EventParseResult>> {
577        for frame in extract_frames(message) {
578            let trimmed_item = frame.data.trim();
579            let transaction_processed_str = "TransactionProcessed";
580
581            if !trimmed_item.contains(transaction_processed_str) {
582                continue;
583            }
584
585            if let Ok(parsed_json) = serde_json::from_str::<Value>(trimmed_item) {
586                let transaction = parsed_json.get(transaction_processed_str);
587                if let Some(transaction_processed) =
588                    transaction.and_then(|transaction| transaction.as_object())
589                {
590                    if let Some(transaction_hash_processed) = transaction_processed
591                        .get("transaction_hash")
592                        .and_then(|transaction_hash| {
593                            transaction_hash
594                                .get("Version1")
595                                .or_else(|| transaction_hash.get("Deploy"))
596                                .and_then(|transaction_hash| transaction_hash.as_str())
597                        })
598                    {
599                        let mut transaction_hash_found =
600                            target_hash == Some(transaction_hash_processed);
601
602                        let transaction_processed: Option<TransactionProcessed> =
603                            serde_json::from_value(transaction.unwrap().clone()).ok();
604
605                        let body = Some(Body {
606                            transaction_processed,
607                        });
608
609                        let event_parse_result = EventParseResult { err: None, body };
610
611                        if transaction_hash_found {
612                            self.unsubscribe(target_hash.unwrap().to_string());
613                            self.stop();
614                            return Some([event_parse_result].to_vec());
615                        }
616
617                        let mut results: Vec<EventParseResult> = [].to_vec();
618                        for subscription in self.subscriptions.clone().iter() {
619                            if transaction_hash_processed == subscription.target_hash {
620                                let event_handler = &subscription.event_handler_fn;
621
622                                #[cfg(not(all(feature = "js", target_arch = "wasm32")))]
623                                {
624                                    event_handler.call(event_parse_result.clone());
625                                }
626                                #[cfg(all(feature = "js", target_arch = "wasm32"))]
627                                {
628                                    let this = JsValue::null();
629                                    let args = js_sys::Array::new();
630                                    args.push(
631                                        &JsValue::from_serde(&event_parse_result.clone()).unwrap(),
632                                    );
633                                    event_handler.apply(&this, &args).unwrap();
634                                }
635
636                                self.unsubscribe(transaction_hash_processed.to_string());
637                                transaction_hash_found = true;
638                                results.push(event_parse_result.clone())
639                            }
640                        }
641
642                        if transaction_hash_found && self.subscriptions.is_empty() {
643                            self.stop();
644                            return Some(results);
645                        }
646                    }
647                }
648            } else {
649                let event_parse_result = EventParseResult {
650                    err: Some("Failed to parse JSON data.".to_string()),
651                    body: None,
652                };
653                return Some([event_parse_result].to_vec());
654            }
655        }
656        None
657    }
658}
659
660/// A wrapper for an event handler function, providing synchronization and cloning capabilities.
661pub struct EventHandlerFn(Arc<Mutex<dyn Fn(EventParseResult) + Send + Sync>>);
662
663impl EventHandlerFn {
664    /// Creates a new `EventHandlerFn` with the specified event handling function.
665    ///
666    /// # Arguments
667    ///
668    /// * `func` - A function that takes an `EventParseResult` as an argument.
669    ///
670    /// # Returns
671    ///
672    /// A new `EventHandlerFn` instance.
673    pub fn new<F>(func: F) -> Self
674    where
675        F: Fn(EventParseResult) + Send + Sync + 'static,
676    {
677        EventHandlerFn(Arc::new(Mutex::new(func)))
678    }
679
680    /// Calls the stored event handling function with the provided `EventParseResult`.
681    ///
682    /// # Arguments
683    ///
684    /// * `event_result` - The result of an event to be passed to the stored event handling function.
685    pub fn call(&self, event_result: EventParseResult) {
686        let func = self.0.lock().unwrap();
687        (*func)(event_result); // Call the stored function with arguments
688    }
689}
690
691impl fmt::Debug for EventHandlerFn {
692    /// Implements the `Debug` trait for better debugging support.
693    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
694        write!(f, "EventHandlerFn")
695    }
696}
697
698impl Clone for EventHandlerFn {
699    /// Implements the `Clone` trait for creating a cloned instance with shared underlying data.
700    fn clone(&self) -> Self {
701        EventHandlerFn(self.0.clone())
702    }
703}
704
705impl Default for EventHandlerFn {
706    /// Implements the `Default` trait, creating a default instance with a no-op event handling function.
707    fn default() -> Self {
708        EventHandlerFn(Arc::new(Mutex::new(|_event_result| {})))
709    }
710}
711
712// Define Subscription struct with different configurations based on JS/wasm surface.
713#[cfg(not(all(feature = "js", target_arch = "wasm32")))]
714/// Represents a subscription to transaction events (native / Rust-only wasm).
715#[derive(Debug, Clone, Default)]
716pub struct Subscription {
717    /// Transaction target hash to identify the subscription.
718    pub target_hash: String,
719    /// Handler function for transaction events.
720    pub event_handler_fn: EventHandlerFn,
721}
722
723#[cfg(all(feature = "js", target_arch = "wasm32"))]
724/// Represents a subscription to transaction events for wasm32 target architecture.
725#[derive(Debug, Clone, Default)]
726#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
727pub struct Subscription {
728    /// Transaction target hash to identify the subscription.
729    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "targetHash"))]
730    pub target_hash: String,
731    /// Handler function for transaction events.
732    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "eventHandlerFn"))]
733    pub event_handler_fn: js_sys::Function,
734}
735
736impl Subscription {
737    /// Constructor for Subscription for non-wasm32 target architecture.
738    ///
739    /// # Arguments
740    ///
741    /// * `target_hash` - Transaction target hash to identify the subscription.
742    /// * `event_handler_fn` - Handler function for transaction events.
743    #[cfg(not(all(feature = "js", target_arch = "wasm32")))]
744    pub fn new(target_hash: String, event_handler_fn: EventHandlerFn) -> Self {
745        Self {
746            target_hash,
747            event_handler_fn,
748        }
749    }
750}
751
752#[cfg_attr(feature = "js", wasm_bindgen)]
753impl Subscription {
754    /// Constructor for Subscription for wasm32 target architecture.
755    ///
756    /// # Arguments
757    ///
758    /// * `transaction_hash` - Transaction hash to identify the subscription.
759    /// * `event_handler_fn` - Handler function for transaction events.
760    #[cfg(all(feature = "js", target_arch = "wasm32"))]
761    #[cfg_attr(feature = "js", wasm_bindgen(constructor))]
762    pub fn new(target_hash: String, event_handler_fn: js_sys::Function) -> Self {
763        Self {
764            target_hash,
765            event_handler_fn,
766        }
767    }
768}
769
770/// Represents a failure response containing an error message.
771#[derive(Debug, Deserialize, Clone, Default, Serialize)]
772#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
773pub struct Failure {
774    pub cost: String,
775    pub error_message: String,
776}
777
778/// Represents a success response containing a cost value.
779#[derive(Debug, Deserialize, Clone, Default, Serialize)]
780#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
781pub struct Version2 {
782    pub initiator: PublicKeyString,
783    pub error_message: Option<String>,
784    pub limit: String,
785    pub consumed: String,
786    pub cost: String,
787    // pub payment: Vec<Payment>,
788}
789
790#[derive(Debug, Deserialize, Clone, Default, Serialize)]
791#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
792pub struct Payment {
793    pub source: String,
794}
795
796/// Represents the result of an execution, either Success or Failure.
797#[derive(Debug, Deserialize, Clone, Default, Serialize)]
798#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
799pub struct ExecutionResult {
800    /// Optional Success information.
801    #[serde(rename = "Version2")]
802    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "Success"))]
803    pub success: Option<Version2>,
804    /// Optional Failure information.
805    #[serde(rename = "Failure")]
806    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "Failure"))]
807    pub failure: Option<Failure>,
808}
809
810#[derive(Debug, Clone, Serialize, Default)]
811#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
812pub struct HashString {
813    pub hash: String,
814}
815
816impl HashString {
817    fn from_hash(hash: String) -> Self {
818        HashString { hash }
819    }
820}
821
822impl<'de> Deserialize<'de> for HashString {
823    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
824    where
825        D: Deserializer<'de>,
826    {
827        let map: std::collections::HashMap<String, String> =
828            Deserialize::deserialize(deserializer)?;
829
830        if let Some(hash) = map.get("Version1").or_else(|| map.get("Deploy")) {
831            Ok(HashString::from_hash(hash.clone()))
832        } else {
833            Err(serde::de::Error::missing_field("Deploy or Version1"))
834        }
835    }
836}
837
838#[cfg_attr(feature = "js", wasm_bindgen)]
839impl HashString {
840    #[cfg_attr(feature = "js", wasm_bindgen(getter, js_name = "Deploy"))]
841    pub fn deploy(&self) -> String {
842        self.hash.clone()
843    }
844
845    #[cfg_attr(feature = "js", wasm_bindgen(getter, js_name = "Version1"))]
846    pub fn version1(&self) -> String {
847        self.hash.clone()
848    }
849
850    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "toString"))]
851    pub fn to_string_js(&self) -> String {
852        self.to_string()
853    }
854}
855
856impl fmt::Display for HashString {
857    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
858        write!(f, "{}", self.hash)
859    }
860}
861
862#[derive(Debug, Deserialize, Clone, Serialize, Default)]
863#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
864pub struct PublicKeyString {
865    #[serde(rename = "PublicKey")]
866    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "PublicKey"))]
867    pub public_key: String,
868}
869
870#[derive(Debug, Deserialize, Clone, Serialize, Default)]
871#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
872pub struct Message {
873    #[serde(rename = "String")]
874    #[cfg_attr(feature = "js", wasm_bindgen(js_name = "String"))]
875    pub string: String,
876}
877
878#[derive(Debug, Deserialize, Clone, Serialize, Default)]
879#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
880pub struct Messages {
881    pub entity_hash: String,
882    pub message: Message,
883    pub topic_name: String,
884    pub topic_name_hash: String,
885    pub topic_index: u32,
886    pub block_index: u64,
887}
888
889/// Represents processed deploy information.
890#[derive(Debug, Deserialize, Clone, Default, Serialize)]
891#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
892pub struct TransactionProcessed {
893    #[serde(alias = "transaction_hash")]
894    pub hash: HashString,
895    pub initiator_addr: PublicKeyString,
896    pub timestamp: String,
897    pub ttl: String,
898    pub block_hash: String,
899    /// Result of the execution, either Success or Failure.
900    pub execution_result: ExecutionResult,
901    pub messages: Vec<Messages>,
902}
903
904/// Represents the body of an event, containing processed deploy information.
905#[derive(Debug, Deserialize, Clone, Default, Serialize)]
906#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
907pub struct Body {
908    #[serde(rename = "TransactionProcessed")]
909    pub transaction_processed: Option<TransactionProcessed>,
910}
911
912// Implementing methods to get the field using different aliases
913#[cfg_attr(feature = "js", wasm_bindgen)]
914impl Body {
915    #[cfg_attr(feature = "js", wasm_bindgen(getter, js_name = "get_deploy_processed"))]
916    #[deprecated(note = "prefer 'get_transaction_processed'")]
917    #[allow(deprecated)]
918    pub fn get_deploy_processed(&self) -> Option<TransactionProcessed> {
919        self.transaction_processed.clone()
920    }
921
922    #[cfg_attr(
923        feature = "js",
924        wasm_bindgen(getter, js_name = "get_transaction_processed")
925    )]
926    pub fn get_transaction_processed(&self) -> Option<TransactionProcessed> {
927        self.transaction_processed.clone()
928    }
929}
930
931/// Represents the result of parsing an event, containing error information and the event body.
932#[derive(Debug, Deserialize, Clone, Default, Serialize)]
933#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
934pub struct EventParseResult {
935    pub err: Option<String>,
936    pub body: Option<Body>,
937}
938
939#[cfg(test)]
940mod tests {
941    use super::*;
942    use crate::sdk::sse::watcher::deploy_mock::DEPLOY_MOCK;
943    use crate::sdk::sse::watcher::transaction_mock::TRANSACTION_MOCK;
944    use sdk_tests::tests::helpers::get_network_constants;
945    use tokio;
946
947    #[test]
948    fn test_new() {
949        // Arrange
950        let (_, events_url, _, _, _) = get_network_constants();
951        let timeout_duration = 5000;
952
953        // Act
954        let watcher = Watcher::new(events_url.clone(), Some(timeout_duration));
955
956        // Assert
957        assert_eq!(watcher.events_url, events_url);
958        assert_eq!(watcher.subscriptions.len(), 0);
959        assert!(*watcher.active.lock().unwrap());
960        assert_eq!(
961            watcher.timeout_duration,
962            Duration::try_milliseconds(timeout_duration.try_into().unwrap()).unwrap()
963        );
964    }
965
966    #[test]
967    fn test_new_default_timeout() {
968        // Arrange
969        let (_, events_url, _, _, _) = get_network_constants();
970
971        // Act
972        let watcher = Watcher::new(events_url.clone(), None);
973
974        // Assert
975        assert_eq!(watcher.events_url, events_url);
976        assert_eq!(watcher.subscriptions.len(), 0);
977        assert!(*watcher.active.lock().unwrap());
978        assert_eq!(
979            watcher.timeout_duration,
980            Duration::try_milliseconds(DEFAULT_TIMEOUT_MS.try_into().unwrap()).unwrap()
981        );
982    }
983
984    #[tokio::test]
985    #[allow(deprecated)]
986    async fn test_process_events_legacy() {
987        // Arrange
988        let (_, events_url, _, _, _) = get_network_constants();
989        let watcher = Watcher::new(events_url, None);
990        let deploy_hash = "19dbf9bdcd821e55392393c74c86deede02d9434d62d0bc72ab381ce7ea1c4f2";
991
992        let target_deploy_hash = Some(deploy_hash);
993
994        // Act
995        let result = watcher.process_events(DEPLOY_MOCK, target_deploy_hash);
996
997        // Assert
998        assert!(result.is_some());
999        let results = result.unwrap();
1000        assert_eq!(results.len(), 1);
1001
1002        let event_parse_result = &results[0];
1003        assert!(event_parse_result.err.is_none());
1004
1005        let body = event_parse_result.body.as_ref().unwrap();
1006        let get_deploy_processed = body.get_deploy_processed().unwrap();
1007        assert_eq!(get_deploy_processed.hash.to_string(), deploy_hash);
1008    }
1009
1010    #[tokio::test]
1011    async fn test_process_events() {
1012        // Arrange
1013        let (_, events_url, _, _, _) = get_network_constants();
1014        let watcher = Watcher::new(events_url, None);
1015        let transaction_hash = "8c6823d9480eee9fe0cfb5ed1fbf77f928cc6af21121298c05b4e3d87a328271";
1016
1017        let target_transaction_hash = Some(transaction_hash);
1018
1019        // Act
1020        let result = watcher.process_events(TRANSACTION_MOCK, target_transaction_hash);
1021
1022        // Assert
1023        assert!(result.is_some());
1024        let results = result.unwrap();
1025        assert_eq!(results.len(), 1);
1026
1027        let event_parse_result = &results[0];
1028        assert!(event_parse_result.err.is_none());
1029
1030        let body = event_parse_result.body.as_ref().unwrap();
1031        let transaction_processed = body.get_transaction_processed().unwrap();
1032        assert_eq!(transaction_processed.hash.to_string(), transaction_hash);
1033    }
1034
1035    fn assert_wait_or_connect_err(err: &Option<String>) {
1036        let msg = err.as_ref().expect("expected error");
1037        assert!(
1038            msg == "Timeout expired"
1039                || msg.contains("error sending request")
1040                || msg.contains("Failed to fetch stream")
1041                || msg.contains("error"),
1042            "unexpected err: {msg}"
1043        );
1044    }
1045
1046    #[tokio::test]
1047    async fn test_start_timeout() {
1048        // Arrange
1049        let (_, events_url, _, _, _) = get_network_constants();
1050        let watcher = Watcher::new(events_url, Some(1));
1051
1052        // Act
1053        let result = watcher.start().await;
1054
1055        // Assert
1056        assert!(result.is_some());
1057        let results = result.unwrap();
1058        assert_eq!(results.len(), 1);
1059        assert_wait_or_connect_err(&results[0].err);
1060        assert!(results[0].body.is_none());
1061    }
1062
1063    #[test]
1064    fn test_stop() {
1065        // Arrange
1066        let (_, events_url, _, _, _) = get_network_constants();
1067        let watcher = Watcher::new(events_url, None);
1068        assert!(*watcher.active.lock().unwrap());
1069
1070        // Act
1071        watcher.stop();
1072
1073        // Assert
1074        assert!(!(*watcher.active.lock().unwrap()));
1075    }
1076
1077    #[test]
1078    fn test_subscribe() {
1079        // Arrange
1080        let (_, events_url, _, _, _) = get_network_constants();
1081        let mut watcher = Watcher::new(events_url, None);
1082        let transaction_hash = "8c6823d9480eee9fe0cfb5ed1fbf77f928cc6af21121298c05b4e3d87a328271";
1083
1084        // Create a subscription
1085        let subscription =
1086            Subscription::new(transaction_hash.to_string(), EventHandlerFn::default());
1087
1088        // Act
1089        let result = watcher.subscribe(vec![subscription]);
1090
1091        // Assert
1092        assert!(result.is_ok());
1093
1094        // Try subscribing to the same deploy hash again
1095        let duplicate_subscription =
1096            Subscription::new(transaction_hash.to_string(), EventHandlerFn::default());
1097        let result_duplicate = watcher.subscribe(vec![duplicate_subscription]);
1098
1099        // Assert
1100        assert!(result_duplicate.is_err());
1101        assert_eq!(
1102            result_duplicate.err().unwrap(),
1103            "Already subscribed to this event"
1104        );
1105    }
1106
1107    #[test]
1108    fn test_unsubscribe() {
1109        // Arrange
1110        let (_, events_url, _, _, _) = get_network_constants();
1111        let mut watcher = Watcher::new(events_url, None);
1112        let transaction_hash = "8c6823d9480eee9fe0cfb5ed1fbf77f928cc6af21121298c05b4e3d87a328271";
1113
1114        // Subscribe to a transaction hash
1115        let transaction_hash_to_subscribe = transaction_hash.to_string();
1116        let subscription = Subscription::new(
1117            transaction_hash_to_subscribe.clone(),
1118            EventHandlerFn::default(),
1119        );
1120        let _ = watcher.subscribe(vec![subscription]);
1121
1122        // Assert that the deploy hash is initially subscribed
1123        assert!(watcher
1124            .subscriptions
1125            .iter()
1126            .any(|s| s.target_hash == transaction_hash_to_subscribe));
1127
1128        // Act
1129        watcher.unsubscribe(transaction_hash_to_subscribe.clone());
1130
1131        // Assert that the deploy hash is unsubscribed after calling unsubscribe
1132        assert!(!watcher
1133            .subscriptions
1134            .iter()
1135            .any(|s| s.target_hash == transaction_hash_to_subscribe));
1136    }
1137
1138    #[test]
1139    #[allow(deprecated)]
1140    fn test_sdk_watch_deploy_retunrs_instance() {
1141        // Arrange
1142        let sdk = SDK::new(None, None, None);
1143        let (_, events_url, _, _, _) = get_network_constants();
1144        let timeout_duration = 5000;
1145
1146        // Act
1147        let watcher = sdk.watch_deploy(&events_url, Some(timeout_duration));
1148
1149        // Assert
1150        assert_eq!(watcher.events_url, events_url);
1151        assert_eq!(watcher.subscriptions.len(), 0);
1152        assert!(*watcher.active.lock().unwrap());
1153        assert_eq!(
1154            watcher.timeout_duration,
1155            Duration::try_milliseconds(timeout_duration.try_into().unwrap()).unwrap()
1156        );
1157    }
1158
1159    #[test]
1160    fn test_sdk_watch_transaction_retunrs_instance() {
1161        // Arrange
1162        let sdk = SDK::new(None, None, None);
1163        let (_, events_url, _, _, _) = get_network_constants();
1164        let timeout_duration = 5000;
1165
1166        // Act
1167        let watcher = sdk.watch_transaction(&events_url, Some(timeout_duration));
1168
1169        // Assert
1170        assert_eq!(watcher.events_url, events_url);
1171        assert_eq!(watcher.subscriptions.len(), 0);
1172        assert!(*watcher.active.lock().unwrap());
1173        assert_eq!(
1174            watcher.timeout_duration,
1175            Duration::try_milliseconds(timeout_duration.try_into().unwrap()).unwrap()
1176        );
1177    }
1178
1179    #[tokio::test]
1180    #[allow(deprecated)]
1181    async fn test_wait_deploy_timeout() {
1182        // Arrange
1183        let sdk = SDK::new(None, None, None);
1184        let (_, events_url, _, _, _) = get_network_constants();
1185        let deploy_hash = "19dbf9bdcd821e55392393c74c86deede02d9434d62d0bc72ab381ce7ea1c4f2";
1186        let timeout_duration = Some(5000);
1187
1188        // Act
1189        let result = sdk
1190            .wait_deploy(&events_url, deploy_hash, timeout_duration)
1191            .await;
1192
1193        // Assert
1194        assert!(result.is_ok());
1195        let event_parse_result = result.unwrap();
1196        assert_wait_or_connect_err(&event_parse_result.err);
1197    }
1198
1199    #[tokio::test]
1200    async fn test_wait_transaction_timeout() {
1201        // Arrange
1202        let sdk = SDK::new(None, None, None);
1203        let (_, events_url, _, _, _) = get_network_constants();
1204        let transaction_hash = "8c6823d9480eee9fe0cfb5ed1fbf77f928cc6af21121298c05b4e3d87a328271";
1205        let timeout_duration = Some(5000);
1206
1207        // Act
1208        let result = sdk
1209            .wait_transaction(&events_url, transaction_hash, timeout_duration)
1210            .await;
1211
1212        // Assert
1213        assert!(result.is_ok());
1214        let event_parse_result = result.unwrap();
1215        assert_wait_or_connect_err(&event_parse_result.err);
1216    }
1217}