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#[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 #[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 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 #[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 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 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 #[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 #[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 #[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 #[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#[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 #[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 #[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 #[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 #[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 #[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 pub async fn start(&self) -> Option<Vec<EventParseResult>> {
383 self.start_internal(None).await
384 }
385
386 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 let url = url_with_start_from(&self.events_url, Some(0));
414
415 #[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 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 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
660pub struct EventHandlerFn(Arc<Mutex<dyn Fn(EventParseResult) + Send + Sync>>);
662
663impl EventHandlerFn {
664 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 pub fn call(&self, event_result: EventParseResult) {
686 let func = self.0.lock().unwrap();
687 (*func)(event_result); }
689}
690
691impl fmt::Debug for EventHandlerFn {
692 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
694 write!(f, "EventHandlerFn")
695 }
696}
697
698impl Clone for EventHandlerFn {
699 fn clone(&self) -> Self {
701 EventHandlerFn(self.0.clone())
702 }
703}
704
705impl Default for EventHandlerFn {
706 fn default() -> Self {
708 EventHandlerFn(Arc::new(Mutex::new(|_event_result| {})))
709 }
710}
711
712#[cfg(not(all(feature = "js", target_arch = "wasm32")))]
714#[derive(Debug, Clone, Default)]
716pub struct Subscription {
717 pub target_hash: String,
719 pub event_handler_fn: EventHandlerFn,
721}
722
723#[cfg(all(feature = "js", target_arch = "wasm32"))]
724#[derive(Debug, Clone, Default)]
726#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
727pub struct Subscription {
728 #[cfg_attr(feature = "js", wasm_bindgen(js_name = "targetHash"))]
730 pub target_hash: String,
731 #[cfg_attr(feature = "js", wasm_bindgen(js_name = "eventHandlerFn"))]
733 pub event_handler_fn: js_sys::Function,
734}
735
736impl Subscription {
737 #[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 #[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#[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#[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 }
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#[derive(Debug, Deserialize, Clone, Default, Serialize)]
798#[cfg_attr(feature = "js", wasm_bindgen(getter_with_clone))]
799pub struct ExecutionResult {
800 #[serde(rename = "Version2")]
802 #[cfg_attr(feature = "js", wasm_bindgen(js_name = "Success"))]
803 pub success: Option<Version2>,
804 #[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#[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 pub execution_result: ExecutionResult,
901 pub messages: Vec<Messages>,
902}
903
904#[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#[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#[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 let (_, events_url, _, _, _) = get_network_constants();
951 let timeout_duration = 5000;
952
953 let watcher = Watcher::new(events_url.clone(), Some(timeout_duration));
955
956 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 let (_, events_url, _, _, _) = get_network_constants();
970
971 let watcher = Watcher::new(events_url.clone(), None);
973
974 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 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 let result = watcher.process_events(DEPLOY_MOCK, target_deploy_hash);
996
997 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 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 let result = watcher.process_events(TRANSACTION_MOCK, target_transaction_hash);
1021
1022 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 let (_, events_url, _, _, _) = get_network_constants();
1050 let watcher = Watcher::new(events_url, Some(1));
1051
1052 let result = watcher.start().await;
1054
1055 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 let (_, events_url, _, _, _) = get_network_constants();
1067 let watcher = Watcher::new(events_url, None);
1068 assert!(*watcher.active.lock().unwrap());
1069
1070 watcher.stop();
1072
1073 assert!(!(*watcher.active.lock().unwrap()));
1075 }
1076
1077 #[test]
1078 fn test_subscribe() {
1079 let (_, events_url, _, _, _) = get_network_constants();
1081 let mut watcher = Watcher::new(events_url, None);
1082 let transaction_hash = "8c6823d9480eee9fe0cfb5ed1fbf77f928cc6af21121298c05b4e3d87a328271";
1083
1084 let subscription =
1086 Subscription::new(transaction_hash.to_string(), EventHandlerFn::default());
1087
1088 let result = watcher.subscribe(vec![subscription]);
1090
1091 assert!(result.is_ok());
1093
1094 let duplicate_subscription =
1096 Subscription::new(transaction_hash.to_string(), EventHandlerFn::default());
1097 let result_duplicate = watcher.subscribe(vec![duplicate_subscription]);
1098
1099 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 let (_, events_url, _, _, _) = get_network_constants();
1111 let mut watcher = Watcher::new(events_url, None);
1112 let transaction_hash = "8c6823d9480eee9fe0cfb5ed1fbf77f928cc6af21121298c05b4e3d87a328271";
1113
1114 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!(watcher
1124 .subscriptions
1125 .iter()
1126 .any(|s| s.target_hash == transaction_hash_to_subscribe));
1127
1128 watcher.unsubscribe(transaction_hash_to_subscribe.clone());
1130
1131 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 let sdk = SDK::new(None, None, None);
1143 let (_, events_url, _, _, _) = get_network_constants();
1144 let timeout_duration = 5000;
1145
1146 let watcher = sdk.watch_deploy(&events_url, Some(timeout_duration));
1148
1149 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 let sdk = SDK::new(None, None, None);
1163 let (_, events_url, _, _, _) = get_network_constants();
1164 let timeout_duration = 5000;
1165
1166 let watcher = sdk.watch_transaction(&events_url, Some(timeout_duration));
1168
1169 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 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 let result = sdk
1190 .wait_deploy(&events_url, deploy_hash, timeout_duration)
1191 .await;
1192
1193 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 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 let result = sdk
1209 .wait_transaction(&events_url, transaction_hash, timeout_duration)
1210 .await;
1211
1212 assert!(result.is_ok());
1214 let event_parse_result = result.unwrap();
1215 assert_wait_or_connect_err(&event_parse_result.err);
1216 }
1217}