1use std::cell::{Cell, RefCell};
7use std::rc::Rc;
8use std::time::Duration;
9
10use crate::async_driver::{DriverCancel, LocalAsyncDriver};
11use crate::ctx::{Ctx, DeferredCtx, DepTerminal};
12use crate::environment::{
13 HttpRequest, HttpResponse, LocalWebSocketSession, ProcessCommand, ProcessResult,
14 WebSocketDriverEvent, WebSocketEvent, WebSocketRequest, WebSocketSend, WebSocketSendResult,
15};
16use crate::graph::{Graph, GraphNodeOpts};
17use crate::node::Node;
18use crate::protocol::{GraphError, Message};
19use crate::resilience::RetryPolicy;
20
21type CancelSlot = Rc<RefCell<Option<DriverCancel>>>;
22type CancelSlots = Rc<RefCell<Vec<CancelSlot>>>;
23type OutboundSend<T, R> =
24 Rc<dyn Fn(T, Box<dyn FnOnce(Result<R, GraphError>)>) -> Option<DriverCancel>>;
25
26#[derive(Debug, Clone, PartialEq, Eq)]
27pub enum OutboundEvent<T, R> {
29 Attempt {
31 value: T,
33 attempt: u32,
35 },
36 Retry {
38 value: T,
40 attempt: u32,
42 delay_ms: u64,
44 error: String,
46 },
47 Sent {
49 value: T,
51 attempt: u32,
53 result: R,
55 },
56 Exhausted {
58 value: T,
60 attempt: u32,
62 error: String,
64 },
65 UpstreamComplete,
67 UpstreamError {
69 error: String,
71 },
72}
73
74#[derive(Debug, Clone, PartialEq, Eq)]
75pub enum OutboundState {
77 Idle,
79 Running,
81 Waiting,
83 Succeeded,
85 Exhausted,
87 Failed,
89 Completed,
91}
92
93#[derive(Debug, Clone, PartialEq, Eq)]
94pub struct OutboundStatus {
96 pub state: OutboundState,
98 pub in_flight: u32,
100 pub attempt: u32,
102 pub sent: u64,
104 pub failed: u64,
106 pub last_delay_ms: Option<u64>,
108}
109
110impl Default for OutboundStatus {
111 fn default() -> Self {
112 Self {
113 state: OutboundState::Idle,
114 in_flight: 0,
115 attempt: 0,
116 sent: 0,
117 failed: 0,
118 last_delay_ms: None,
119 }
120 }
121}
122
123pub struct OutboundBundle<T: 'static, R: 'static> {
125 pub events: Node<OutboundEvent<T, R>>,
127 pub status: Node<OutboundStatus>,
129 pub attempts: Node<u32>,
131 pub errors: Node<String>,
133}
134
135#[derive(Clone, Default)]
136pub struct OutboundAdapterOptions {
138 pub name: Option<String>,
140 pub retry: RetryPolicy,
142}
143
144#[derive(Debug, Clone, PartialEq, Eq)]
145pub enum WebSocketSessionCommand {
147 Start,
149 Send(WebSocketSend),
151 Close {
153 code: Option<u16>,
155 reason: Option<String>,
157 },
158}
159
160#[derive(Debug, Clone, PartialEq, Eq)]
161pub enum WebSocketSessionInbound {
163 Text(String),
165 Binary(Vec<u8>),
167}
168
169#[derive(Debug, Clone, PartialEq, Eq)]
170pub enum WebSocketSessionLifecycle {
172 Starting {
174 attempt: u32,
176 max_attempts: u32,
178 },
179 Open {
181 attempt: u32,
183 },
184 Sent {
186 message: WebSocketSend,
188 },
189 Closing {
191 code: Option<u16>,
193 reason: Option<String>,
195 },
196 Closed {
198 code: Option<u16>,
200 reason: Option<String>,
202 },
203 Retrying {
205 attempt: u32,
207 next_attempt: u32,
209 delay_ms: u64,
211 error: String,
213 },
214 Exhausted {
216 attempt: u32,
218 error: String,
220 },
221}
222
223#[derive(Debug, Clone, PartialEq, Eq)]
224pub enum WebSocketSessionStateKind {
226 Idle,
228 Connecting,
230 Open,
232 Closing,
234 Closed,
236 Waiting,
238 Exhausted,
240 Errored,
242}
243
244#[derive(Debug, Clone, PartialEq, Eq)]
245pub struct WebSocketSessionStatus {
247 pub state: WebSocketSessionStateKind,
249 pub attempt: u32,
251 pub max_attempts: u32,
253 pub sent: u64,
255 pub received: u64,
257 pub errors: u64,
259 pub last_delay_ms: Option<u64>,
261}
262
263#[derive(Debug, Clone, PartialEq, Eq)]
264pub enum WebSocketSessionOutbound {
266 Queued {
268 seq: u64,
270 message: WebSocketSend,
272 },
273 Sending {
275 seq: u64,
277 message: WebSocketSend,
279 },
280 Sent {
282 seq: u64,
284 message: WebSocketSend,
286 },
287 Rejected {
289 seq: u64,
291 message: WebSocketSend,
293 error: String,
295 },
296 Canceled {
298 seq: u64,
300 message: WebSocketSend,
302 reason: String,
304 },
305}
306
307pub struct WebSocketSessionBundle {
309 pub command: Node<WebSocketSessionCommand>,
311 pub inbound: Node<WebSocketSessionInbound>,
313 pub lifecycle: Node<WebSocketSessionLifecycle>,
315 pub outbound: Node<WebSocketSessionOutbound>,
317 pub status: Node<WebSocketSessionStatus>,
319 pub errors: Node<String>,
321 pub attempts: Node<u32>,
323}
324
325impl WebSocketSessionBundle {
326 pub fn start(&self) {
328 self.command.set(WebSocketSessionCommand::Start);
329 }
330
331 pub fn send(&self, message: WebSocketSend) {
333 self.command.set(WebSocketSessionCommand::Send(message));
334 }
335
336 pub fn send_text(&self, text: impl Into<String>) {
338 self.send(WebSocketSend::text(text));
339 }
340
341 pub fn send_binary(&self, bytes: impl Into<Vec<u8>>) {
343 self.send(WebSocketSend::binary(bytes));
344 }
345
346 pub fn close(&self, code: Option<u16>, reason: Option<String>) {
348 self.command
349 .set(WebSocketSessionCommand::Close { code, reason });
350 }
351}
352
353#[derive(Debug, Clone, PartialEq, Eq, Default)]
354pub enum WebSocketSessionSendPolicy {
356 #[default]
357 Reject,
359 Buffer {
361 max_pending: usize,
363 },
364}
365
366#[derive(Clone, Default)]
367pub struct WebSocketSessionOptions {
369 pub name: Option<String>,
371 pub retry: RetryPolicy,
373 pub send_policy: WebSocketSessionSendPolicy,
375}
376
377#[derive(Debug, Clone, PartialEq, Eq)]
378enum WebSocketSessionEvent {
379 Attempt {
380 attempt: u32,
381 max_attempts: u32,
382 },
383 Open {
384 attempt: u32,
385 },
386 Message(WebSocketSessionInbound),
387 Sent {
388 message: WebSocketSend,
389 },
390 Closing {
391 code: Option<u16>,
392 reason: Option<String>,
393 },
394 Closed {
395 code: Option<u16>,
396 reason: Option<String>,
397 },
398 Retry {
399 attempt: u32,
400 next_attempt: u32,
401 delay_ms: u64,
402 error: String,
403 },
404 Outbound(WebSocketSessionOutbound),
405 Error {
406 attempt: Option<u32>,
407 error: String,
408 },
409 Exhausted {
410 attempt: u32,
411 error: String,
412 },
413}
414
415#[derive(Clone)]
416struct WebSocketSessionStateCell(Rc<RefCell<WebSocketSessionState>>);
417
418struct WebSocketSessionState {
419 node_active: bool,
420 connected: bool,
421 waiting_retry: bool,
422 current_attempt: u32,
423 current_session: Option<Rc<dyn LocalWebSocketSession>>,
424 next_send_id: u64,
425 next_outbound_seq: u64,
426 pending_outbound: Vec<PendingWebSocketOutbound>,
427 live_outbound: Vec<LiveWebSocketOutbound>,
428 send_cancels: Vec<WebSocketSendCancel>,
429 retry_cancels: Vec<CancelSlot>,
430}
431
432#[derive(Clone)]
433struct WebSocketSendCancel {
434 id: u64,
435 slot: CancelSlot,
436}
437
438#[derive(Debug, Clone, PartialEq, Eq)]
439struct PendingWebSocketOutbound {
440 seq: u64,
441 message: WebSocketSend,
442}
443
444#[derive(Debug, Clone, PartialEq, Eq)]
445struct LiveWebSocketOutbound {
446 id: u64,
447 item: PendingWebSocketOutbound,
448}
449
450pub fn websocket_session(graph: &Graph, request: WebSocketRequest) -> WebSocketSessionBundle {
452 websocket_session_with_options(graph, request, WebSocketSessionOptions::default())
453}
454
455pub fn websocket_session_with_options(
457 graph: &Graph,
458 request: WebSocketRequest,
459 opts: WebSocketSessionOptions,
460) -> WebSocketSessionBundle {
461 assert!(
462 !request.url.is_empty(),
463 "websocket_session: url must be non-empty"
464 );
465 let name = opts
466 .name
467 .clone()
468 .unwrap_or_else(|| "websocketSession".to_owned());
469 validate_websocket_session_send_policy(&opts.send_policy);
470 let command = graph.state_empty_opts::<WebSocketSessionCommand>(GraphNodeOpts::named(format!(
471 "{name}/command"
472 )));
473 let events = websocket_session_events(
474 graph,
475 &command,
476 request,
477 name.clone(),
478 opts.retry.clone(),
479 opts.send_policy.clone(),
480 );
481 websocket_session_nodes(graph, command, events, name, opts.retry)
482}
483
484fn validate_websocket_session_send_policy(policy: &WebSocketSessionSendPolicy) {
485 if let WebSocketSessionSendPolicy::Buffer { max_pending } = policy {
486 assert!(
487 *max_pending > 0,
488 "websocket_session: buffer send_policy requires finite max_pending >= 1"
489 );
490 }
491}
492
493fn websocket_session_events(
494 graph: &Graph,
495 command: &Node<WebSocketSessionCommand>,
496 request: WebSocketRequest,
497 name: String,
498 policy: RetryPolicy,
499 send_policy: WebSocketSessionSendPolicy,
500) -> Node<WebSocketSessionEvent> {
501 let events_name = name.clone();
502 graph.node_opts::<WebSocketSessionEvent, _>(
503 vec![command.erased()],
504 move |ctx| {
505 let state = init_websocket_session_state(ctx);
506 let out = Rc::new(ctx.defer());
507 let driver = ctx.environment().websocket_driver();
508 let async_driver = ctx.local_async_driver();
509 for command in ctx.batch::<WebSocketSessionCommand>(0) {
510 match command.as_ref() {
511 WebSocketSessionCommand::Start => {
512 start_websocket_session_attempt(WebSocketSessionStart {
513 state: state.clone(),
514 out: out.clone(),
515 driver: driver.clone(),
516 async_driver: async_driver.clone(),
517 request: request.clone(),
518 name: name.clone(),
519 policy: policy.clone(),
520 attempt: 1,
521 });
522 }
523 WebSocketSessionCommand::Send(message) => {
524 send_websocket_session_message(
525 &state,
526 &out,
527 &name,
528 message.clone(),
529 &send_policy,
530 );
531 }
532 WebSocketSessionCommand::Close { code, reason } => {
533 cancel_websocket_retry_timers(&state);
534 out.emit(WebSocketSessionEvent::Closing {
535 code: *code,
536 reason: reason.clone(),
537 });
538 let cleanup = close_websocket_session_connection(&state, false);
539 emit_outbound_canceled(
540 &out,
541 cleanup.canceled,
542 format!("{name}: session closed"),
543 );
544 let pending = drain_pending_websocket_outbound(&state);
545 emit_outbound_canceled(&out, pending, format!("{name}: session closed"));
546 if let Some(session) = cleanup.session {
547 session.close(*code, reason.clone());
548 }
549 out.emit(WebSocketSessionEvent::Closed {
550 code: *code,
551 reason: reason.clone(),
552 });
553 }
554 }
555 }
556 },
557 GraphNodeOpts::named(format!("{events_name}/events")),
558 )
559}
560
561struct WebSocketSessionStart {
562 state: WebSocketSessionStateCell,
563 out: Rc<DeferredCtx>,
564 driver: Option<Rc<dyn crate::environment::LocalWebSocketDriver>>,
565 async_driver: Option<Rc<dyn LocalAsyncDriver>>,
566 request: WebSocketRequest,
567 name: String,
568 policy: RetryPolicy,
569 attempt: u32,
570}
571
572fn start_websocket_session_attempt(args: WebSocketSessionStart) {
573 {
574 let mut state = args.state.0.borrow_mut();
575 if !state.node_active
576 || state.connected
577 || state.current_attempt != 0
578 || state.waiting_retry
579 {
580 return;
581 }
582 state.current_attempt = args.attempt;
583 }
584 args.out.emit(WebSocketSessionEvent::Attempt {
585 attempt: args.attempt,
586 max_attempts: args.policy.max_attempts,
587 });
588 let Some(driver) = args.driver.clone() else {
589 exhaust_websocket_session(
590 &args.state,
591 &args.out,
592 args.attempt,
593 format!("{}: missing websocket driver", args.name),
594 );
595 return;
596 };
597 let state_for_callback = args.state.clone();
598 let out_for_callback = args.out.clone();
599 let name_for_callback = args.name.clone();
600 let policy_for_callback = args.policy.clone();
601 let request_for_callback = args.request.clone();
602 let driver_for_callback = args.driver.clone();
603 let async_for_callback = args.async_driver.clone();
604 let session = driver.connect_session(
605 args.request.clone(),
606 Rc::new(move |event| {
607 if !is_current_websocket_attempt(&state_for_callback, args.attempt) {
608 return;
609 }
610 match event {
611 WebSocketDriverEvent::Event(WebSocketEvent::Open) => {
612 state_for_callback.0.borrow_mut().connected = true;
613 out_for_callback.emit(WebSocketSessionEvent::Open {
614 attempt: args.attempt,
615 });
616 flush_pending_websocket_outbound(
617 &state_for_callback,
618 &out_for_callback,
619 &name_for_callback,
620 );
621 }
622 WebSocketDriverEvent::Event(WebSocketEvent::Text(text)) => {
623 out_for_callback.emit(WebSocketSessionEvent::Message(
624 WebSocketSessionInbound::Text(text),
625 ));
626 }
627 WebSocketDriverEvent::Event(WebSocketEvent::Binary(bytes)) => {
628 out_for_callback.emit(WebSocketSessionEvent::Message(
629 WebSocketSessionInbound::Binary(bytes),
630 ));
631 }
632 WebSocketDriverEvent::Event(WebSocketEvent::Close { code, reason }) => {
633 let normal_close = code == Some(1000);
634 let cleanup = close_websocket_session_connection(&state_for_callback, false);
635 emit_outbound_canceled(
636 &out_for_callback,
637 cleanup.canceled,
638 format!("{name_for_callback}: session closed"),
639 );
640 out_for_callback.emit(WebSocketSessionEvent::Closed { code, reason });
641 if normal_close {
642 let pending = drain_pending_websocket_outbound(&state_for_callback);
643 emit_outbound_canceled(
644 &out_for_callback,
645 pending,
646 format!("{name_for_callback}: session closed"),
647 );
648 }
649 if !normal_close {
650 retry_or_exhaust_websocket_session(WebSocketSessionRetry {
651 state: state_for_callback.clone(),
652 out: out_for_callback.clone(),
653 driver: driver_for_callback.clone(),
654 async_driver: async_for_callback.clone(),
655 request: request_for_callback.clone(),
656 name: name_for_callback.clone(),
657 policy: policy_for_callback.clone(),
658 attempt: args.attempt,
659 error: format!("{}: websocket closed", name_for_callback),
660 });
661 }
662 }
663 WebSocketDriverEvent::Error(error) => {
664 retry_or_exhaust_websocket_session(WebSocketSessionRetry {
665 state: state_for_callback.clone(),
666 out: out_for_callback.clone(),
667 driver: driver_for_callback.clone(),
668 async_driver: async_for_callback.clone(),
669 request: request_for_callback.clone(),
670 name: name_for_callback.clone(),
671 policy: policy_for_callback.clone(),
672 attempt: args.attempt,
673 error: error.to_string(),
674 });
675 }
676 WebSocketDriverEvent::Complete => {
677 retry_or_exhaust_websocket_session(WebSocketSessionRetry {
678 state: state_for_callback.clone(),
679 out: out_for_callback.clone(),
680 driver: driver_for_callback.clone(),
681 async_driver: async_for_callback.clone(),
682 request: request_for_callback.clone(),
683 name: name_for_callback.clone(),
684 policy: policy_for_callback.clone(),
685 attempt: args.attempt,
686 error: format!("{}: connection completed", name_for_callback),
687 });
688 }
689 }
690 }),
691 );
692 match session {
693 Some(session) => {
694 if is_current_websocket_attempt(&args.state, args.attempt) {
695 args.state.0.borrow_mut().current_session = Some(session);
696 } else {
697 session.cancel();
698 }
699 }
700 None => {
701 exhaust_websocket_session(
702 &args.state,
703 &args.out,
704 args.attempt,
705 format!("{}: missing WebSocket session capability", args.name),
706 );
707 }
708 }
709}
710
711struct WebSocketSessionRetry {
712 state: WebSocketSessionStateCell,
713 out: Rc<DeferredCtx>,
714 driver: Option<Rc<dyn crate::environment::LocalWebSocketDriver>>,
715 async_driver: Option<Rc<dyn LocalAsyncDriver>>,
716 request: WebSocketRequest,
717 name: String,
718 policy: RetryPolicy,
719 attempt: u32,
720 error: String,
721}
722
723fn retry_or_exhaust_websocket_session(args: WebSocketSessionRetry) {
724 let cleanup = close_websocket_session_connection(&args.state, true);
725 emit_outbound_canceled(
726 &args.out,
727 cleanup.canceled,
728 format!("{}: connection cleanup", args.name),
729 );
730 if !args.policy.should_retry(args.attempt) {
731 let pending = drain_pending_websocket_outbound(&args.state);
732 emit_outbound_rejected(&args.out, pending, args.error.clone());
733 exhaust_websocket_session(&args.state, &args.out, args.attempt, args.error);
734 return;
735 }
736 let next_attempt = args.attempt.saturating_add(1);
737 let delay_ms = args.policy.next_delay_ms(next_attempt).unwrap_or_default();
738 args.out.emit(WebSocketSessionEvent::Retry {
739 attempt: args.attempt,
740 next_attempt,
741 delay_ms,
742 error: args.error,
743 });
744 if delay_ms == 0 {
745 start_websocket_session_attempt(WebSocketSessionStart {
746 state: args.state,
747 out: args.out,
748 driver: args.driver,
749 async_driver: args.async_driver,
750 request: args.request,
751 name: args.name,
752 policy: args.policy,
753 attempt: next_attempt,
754 });
755 return;
756 }
757 let Some(async_driver) = args.async_driver.clone() else {
758 exhaust_websocket_session(
759 &args.state,
760 &args.out,
761 args.attempt,
762 format!("{}: missing async driver for delayed reconnect", args.name),
763 );
764 return;
765 };
766 let slot: CancelSlot = Rc::new(RefCell::new(None));
767 let wake_slot = slot.clone();
768 let state = args.state.clone();
769 let state_for_slot = args.state.clone();
770 args.state.0.borrow_mut().waiting_retry = true;
771 let cancel = async_driver.sleep(
772 Duration::from_millis(delay_ms),
773 Box::new(move || {
774 let _ = wake_slot.borrow_mut().take();
775 state.0.borrow_mut().waiting_retry = false;
776 start_websocket_session_attempt(WebSocketSessionStart {
777 state,
778 out: args.out,
779 driver: args.driver,
780 async_driver: args.async_driver,
781 request: args.request,
782 name: args.name,
783 policy: args.policy,
784 attempt: next_attempt,
785 });
786 }),
787 );
788 *slot.borrow_mut() = Some(cancel);
789 state_for_slot.0.borrow_mut().retry_cancels.push(slot);
790}
791
792fn send_websocket_session_message(
793 state: &WebSocketSessionStateCell,
794 out: &Rc<DeferredCtx>,
795 name: &str,
796 message: WebSocketSend,
797 send_policy: &WebSocketSessionSendPolicy,
798) {
799 let item = next_websocket_outbound_item(state, message);
800 if !state.0.borrow().connected {
801 match send_policy {
802 WebSocketSessionSendPolicy::Reject => {
803 let error = format!("{name}: session is not open");
804 out.emit(WebSocketSessionEvent::Outbound(
805 WebSocketSessionOutbound::Rejected {
806 seq: item.seq,
807 message: item.message,
808 error: error.clone(),
809 },
810 ));
811 out.emit(WebSocketSessionEvent::Error {
812 attempt: (state.0.borrow().current_attempt > 0)
813 .then_some(state.0.borrow().current_attempt),
814 error,
815 });
816 }
817 WebSocketSessionSendPolicy::Buffer { max_pending } => {
818 if state.0.borrow().pending_outbound.len() >= *max_pending {
819 let error = format!("{name}: outbound buffer full");
820 out.emit(WebSocketSessionEvent::Outbound(
821 WebSocketSessionOutbound::Rejected {
822 seq: item.seq,
823 message: item.message,
824 error: error.clone(),
825 },
826 ));
827 out.emit(WebSocketSessionEvent::Error {
828 attempt: (state.0.borrow().current_attempt > 0)
829 .then_some(state.0.borrow().current_attempt),
830 error,
831 });
832 return;
833 }
834 out.emit(WebSocketSessionEvent::Outbound(
835 WebSocketSessionOutbound::Queued {
836 seq: item.seq,
837 message: item.message.clone(),
838 },
839 ));
840 state.0.borrow_mut().pending_outbound.push(item);
841 }
842 }
843 return;
844 }
845 send_websocket_session_item(state, out, name, item);
846}
847
848fn send_websocket_session_item(
849 state: &WebSocketSessionStateCell,
850 out: &Rc<DeferredCtx>,
851 name: &str,
852 item: PendingWebSocketOutbound,
853) {
854 let session = {
855 let state_ref = state.0.borrow();
856 if !state_ref.connected {
857 out.emit(WebSocketSessionEvent::Outbound(
858 WebSocketSessionOutbound::Rejected {
859 seq: item.seq,
860 message: item.message,
861 error: format!("{name}: session is not open"),
862 },
863 ));
864 out.emit(WebSocketSessionEvent::Error {
865 attempt: (state_ref.current_attempt > 0).then_some(state_ref.current_attempt),
866 error: format!("{name}: session is not open"),
867 });
868 return;
869 }
870 state_ref.current_session.clone()
871 };
872 let Some(session) = session else {
873 out.emit(WebSocketSessionEvent::Outbound(
874 WebSocketSessionOutbound::Rejected {
875 seq: item.seq,
876 message: item.message,
877 error: format!("{name}: missing WebSocket session capability"),
878 },
879 ));
880 out.emit(WebSocketSessionEvent::Error {
881 attempt: None,
882 error: format!("{name}: missing WebSocket session capability"),
883 });
884 return;
885 };
886 let state_for_callback = state.clone();
887 let out_for_callback = out.clone();
888 let send_id = {
889 let mut state_ref = state.0.borrow_mut();
890 let send_id = state_ref.next_send_id;
891 state_ref.next_send_id = state_ref.next_send_id.saturating_add(1);
892 state_ref.live_outbound.push(LiveWebSocketOutbound {
893 id: send_id,
894 item: item.clone(),
895 });
896 send_id
897 };
898 out.emit(WebSocketSessionEvent::Outbound(
899 WebSocketSessionOutbound::Sending {
900 seq: item.seq,
901 message: item.message.clone(),
902 },
903 ));
904 let send_active = Rc::new(Cell::new(true));
905 let send_active_for_callback = send_active.clone();
906 let cancel = session.send(
907 item.message,
908 Box::new(move |result| {
909 if !send_active_for_callback.get() || !state_for_callback.0.borrow().node_active {
910 return;
911 }
912 send_active_for_callback.set(false);
913 remove_websocket_send_cancel(&state_for_callback, send_id);
914 let Some(live_item) = remove_websocket_live_outbound(&state_for_callback, send_id)
915 else {
916 return;
917 };
918 match result {
919 Ok(_) => {
920 out_for_callback.emit(WebSocketSessionEvent::Outbound(
921 WebSocketSessionOutbound::Sent {
922 seq: live_item.seq,
923 message: live_item.message.clone(),
924 },
925 ));
926 out_for_callback.emit(WebSocketSessionEvent::Sent {
927 message: live_item.message,
928 });
929 }
930 Err(error) => {
931 let error = error.to_string();
932 out_for_callback.emit(WebSocketSessionEvent::Outbound(
933 WebSocketSessionOutbound::Rejected {
934 seq: live_item.seq,
935 message: live_item.message,
936 error: error.clone(),
937 },
938 ));
939 out_for_callback.emit(WebSocketSessionEvent::Error {
940 attempt: (state_for_callback.0.borrow().current_attempt > 0)
941 .then_some(state_for_callback.0.borrow().current_attempt),
942 error,
943 });
944 }
945 }
946 }),
947 );
948 let send_active_for_cancel = send_active.clone();
949 let slot: CancelSlot = Rc::new(RefCell::new(Some(Box::new(move || {
950 send_active_for_cancel.set(false);
951 cancel();
952 }))));
953 if send_active.get() && slot.borrow().is_some() {
954 state
955 .0
956 .borrow_mut()
957 .send_cancels
958 .push(WebSocketSendCancel { id: send_id, slot });
959 } else {
960 let _ = remove_websocket_live_outbound(state, send_id);
961 }
962}
963
964fn websocket_session_nodes(
965 graph: &Graph,
966 command: Node<WebSocketSessionCommand>,
967 events: Node<WebSocketSessionEvent>,
968 name: String,
969 policy: RetryPolicy,
970) -> WebSocketSessionBundle {
971 let inbound = graph.node_opts::<WebSocketSessionInbound, _>(
972 vec![events.erased()],
973 move |ctx| {
974 for event in ctx.batch::<WebSocketSessionEvent>(0) {
975 if let WebSocketSessionEvent::Message(message) = event.as_ref() {
976 ctx.emit(message.clone());
977 }
978 }
979 },
980 GraphNodeOpts::named(format!("{name}/inbound")),
981 );
982 let lifecycle = graph.node_opts::<WebSocketSessionLifecycle, _>(
983 vec![events.erased()],
984 move |ctx| {
985 for event in ctx.batch::<WebSocketSessionEvent>(0) {
986 match event.as_ref() {
987 WebSocketSessionEvent::Attempt {
988 attempt,
989 max_attempts,
990 } => ctx.emit(WebSocketSessionLifecycle::Starting {
991 attempt: *attempt,
992 max_attempts: *max_attempts,
993 }),
994 WebSocketSessionEvent::Open { attempt } => {
995 ctx.emit(WebSocketSessionLifecycle::Open { attempt: *attempt });
996 }
997 WebSocketSessionEvent::Sent { message } => {
998 ctx.emit(WebSocketSessionLifecycle::Sent {
999 message: message.clone(),
1000 });
1001 }
1002 WebSocketSessionEvent::Closing { code, reason } => {
1003 ctx.emit(WebSocketSessionLifecycle::Closing {
1004 code: *code,
1005 reason: reason.clone(),
1006 });
1007 }
1008 WebSocketSessionEvent::Closed { code, reason } => {
1009 ctx.emit(WebSocketSessionLifecycle::Closed {
1010 code: *code,
1011 reason: reason.clone(),
1012 });
1013 }
1014 WebSocketSessionEvent::Retry {
1015 attempt,
1016 next_attempt,
1017 delay_ms,
1018 error,
1019 } => ctx.emit(WebSocketSessionLifecycle::Retrying {
1020 attempt: *attempt,
1021 next_attempt: *next_attempt,
1022 delay_ms: *delay_ms,
1023 error: error.clone(),
1024 }),
1025 WebSocketSessionEvent::Exhausted { attempt, error } => {
1026 ctx.emit(WebSocketSessionLifecycle::Exhausted {
1027 attempt: *attempt,
1028 error: error.clone(),
1029 });
1030 }
1031 WebSocketSessionEvent::Message(_) | WebSocketSessionEvent::Error { .. } => {}
1032 WebSocketSessionEvent::Outbound(_) => {}
1033 }
1034 }
1035 },
1036 GraphNodeOpts::named(format!("{name}/lifecycle")),
1037 );
1038 let outbound = graph.node_opts::<WebSocketSessionOutbound, _>(
1039 vec![events.erased()],
1040 move |ctx| {
1041 for event in ctx.batch::<WebSocketSessionEvent>(0) {
1042 if let WebSocketSessionEvent::Outbound(fact) = event.as_ref() {
1043 ctx.emit(fact.clone());
1044 }
1045 }
1046 },
1047 GraphNodeOpts::named(format!("{name}/outbound")),
1048 );
1049 let status = graph.node_opts::<WebSocketSessionStatus, _>(
1050 vec![events.erased()],
1051 move |ctx| {
1052 let mut next = ctx.state_get::<WebSocketSessionStatus>().map_or_else(
1053 || WebSocketSessionStatus {
1054 state: WebSocketSessionStateKind::Idle,
1055 attempt: 0,
1056 max_attempts: policy.max_attempts,
1057 sent: 0,
1058 received: 0,
1059 errors: 0,
1060 last_delay_ms: None,
1061 },
1062 |status| (*status).clone(),
1063 );
1064 for event in ctx.batch::<WebSocketSessionEvent>(0) {
1065 next = reduce_websocket_session_status(&next, event.as_ref());
1066 }
1067 ctx.state_set(next.clone());
1068 ctx.emit(next);
1069 },
1070 GraphNodeOpts::named(format!("{name}/status")),
1071 );
1072 let errors = graph.node_opts::<String, _>(
1073 vec![events.erased()],
1074 move |ctx| {
1075 for event in ctx.batch::<WebSocketSessionEvent>(0) {
1076 match event.as_ref() {
1077 WebSocketSessionEvent::Retry { error, .. }
1078 | WebSocketSessionEvent::Error { error, .. }
1079 | WebSocketSessionEvent::Exhausted { error, .. } => ctx.emit(error.clone()),
1080 WebSocketSessionEvent::Attempt { .. }
1081 | WebSocketSessionEvent::Open { .. }
1082 | WebSocketSessionEvent::Message(_)
1083 | WebSocketSessionEvent::Sent { .. }
1084 | WebSocketSessionEvent::Closing { .. }
1085 | WebSocketSessionEvent::Closed { .. }
1086 | WebSocketSessionEvent::Outbound(_) => {}
1087 }
1088 }
1089 },
1090 GraphNodeOpts::named(format!("{name}/errors")),
1091 );
1092 let attempts = graph.node_opts::<u32, _>(
1093 vec![events.erased()],
1094 move |ctx| {
1095 for event in ctx.batch::<WebSocketSessionEvent>(0) {
1096 if let WebSocketSessionEvent::Attempt { attempt, .. } = event.as_ref() {
1097 ctx.emit(*attempt);
1098 }
1099 }
1100 },
1101 GraphNodeOpts::named(format!("{name}/attempts")),
1102 );
1103 WebSocketSessionBundle {
1104 command,
1105 inbound,
1106 lifecycle,
1107 outbound,
1108 status,
1109 errors,
1110 attempts,
1111 }
1112}
1113
1114fn reduce_websocket_session_status(
1115 current: &WebSocketSessionStatus,
1116 event: &WebSocketSessionEvent,
1117) -> WebSocketSessionStatus {
1118 match event {
1119 WebSocketSessionEvent::Attempt {
1120 attempt,
1121 max_attempts,
1122 } => WebSocketSessionStatus {
1123 state: WebSocketSessionStateKind::Connecting,
1124 attempt: *attempt,
1125 max_attempts: *max_attempts,
1126 last_delay_ms: None,
1127 ..current.clone()
1128 },
1129 WebSocketSessionEvent::Open { attempt } => WebSocketSessionStatus {
1130 state: WebSocketSessionStateKind::Open,
1131 attempt: *attempt,
1132 last_delay_ms: None,
1133 ..current.clone()
1134 },
1135 WebSocketSessionEvent::Message(_) => WebSocketSessionStatus {
1136 received: current.received.saturating_add(1),
1137 ..current.clone()
1138 },
1139 WebSocketSessionEvent::Sent { .. } => WebSocketSessionStatus {
1140 sent: current.sent.saturating_add(1),
1141 ..current.clone()
1142 },
1143 WebSocketSessionEvent::Closing { .. } => WebSocketSessionStatus {
1144 state: WebSocketSessionStateKind::Closing,
1145 ..current.clone()
1146 },
1147 WebSocketSessionEvent::Closed { .. } => WebSocketSessionStatus {
1148 state: WebSocketSessionStateKind::Closed,
1149 ..current.clone()
1150 },
1151 WebSocketSessionEvent::Retry {
1152 attempt, delay_ms, ..
1153 } => WebSocketSessionStatus {
1154 state: WebSocketSessionStateKind::Waiting,
1155 attempt: *attempt,
1156 errors: current.errors.saturating_add(1),
1157 last_delay_ms: Some(*delay_ms),
1158 ..current.clone()
1159 },
1160 WebSocketSessionEvent::Error { attempt, .. } => WebSocketSessionStatus {
1161 state: WebSocketSessionStateKind::Errored,
1162 attempt: attempt.unwrap_or(current.attempt),
1163 errors: current.errors.saturating_add(1),
1164 ..current.clone()
1165 },
1166 WebSocketSessionEvent::Exhausted { attempt, .. } => WebSocketSessionStatus {
1167 state: WebSocketSessionStateKind::Exhausted,
1168 attempt: *attempt,
1169 errors: current.errors.saturating_add(1),
1170 ..current.clone()
1171 },
1172 WebSocketSessionEvent::Outbound(_) => current.clone(),
1173 }
1174}
1175
1176fn init_websocket_session_state(ctx: &Ctx) -> WebSocketSessionStateCell {
1177 if let Some(state) = ctx.state_get::<WebSocketSessionStateCell>() {
1178 let state = (*state).clone();
1179 state.0.borrow_mut().node_active = true;
1180 return state;
1181 }
1182 let state = WebSocketSessionStateCell(Rc::new(RefCell::new(WebSocketSessionState {
1183 node_active: true,
1184 connected: false,
1185 waiting_retry: false,
1186 current_attempt: 0,
1187 current_session: None,
1188 next_send_id: 0,
1189 next_outbound_seq: 0,
1190 pending_outbound: Vec::new(),
1191 live_outbound: Vec::new(),
1192 send_cancels: Vec::new(),
1193 retry_cancels: Vec::new(),
1194 })));
1195 let cleanup_state = state.clone();
1196 ctx.on_deactivation(move || {
1197 deactivate_websocket_session_state(&cleanup_state);
1198 });
1199 ctx.state_set(state.clone());
1200 state
1201}
1202
1203fn is_current_websocket_attempt(state: &WebSocketSessionStateCell, attempt: u32) -> bool {
1204 let state = state.0.borrow();
1205 state.node_active && state.current_attempt == attempt
1206}
1207
1208fn next_websocket_outbound_item(
1209 state: &WebSocketSessionStateCell,
1210 message: WebSocketSend,
1211) -> PendingWebSocketOutbound {
1212 let mut state_ref = state.0.borrow_mut();
1213 let seq = state_ref.next_outbound_seq;
1214 state_ref.next_outbound_seq = state_ref.next_outbound_seq.saturating_add(1);
1215 PendingWebSocketOutbound { seq, message }
1216}
1217
1218fn flush_pending_websocket_outbound(
1219 state: &WebSocketSessionStateCell,
1220 out: &Rc<DeferredCtx>,
1221 name: &str,
1222) {
1223 let mut pending = drain_pending_websocket_outbound(state);
1224 while !pending.is_empty() {
1225 if !state.0.borrow().connected || state.0.borrow().current_session.is_none() {
1226 emit_outbound_canceled(out, pending, format!("{name}: session closed"));
1227 return;
1228 }
1229 let item = pending.remove(0);
1230 send_websocket_session_item(state, out, name, item);
1231 }
1232}
1233
1234fn drain_pending_websocket_outbound(
1235 state: &WebSocketSessionStateCell,
1236) -> Vec<PendingWebSocketOutbound> {
1237 state.0.borrow_mut().pending_outbound.drain(..).collect()
1238}
1239
1240fn remove_websocket_live_outbound(
1241 state: &WebSocketSessionStateCell,
1242 send_id: u64,
1243) -> Option<PendingWebSocketOutbound> {
1244 let mut state_ref = state.0.borrow_mut();
1245 let index = state_ref
1246 .live_outbound
1247 .iter()
1248 .position(|send| send.id == send_id)?;
1249 Some(state_ref.live_outbound.remove(index).item)
1250}
1251
1252fn emit_outbound_canceled(out: &DeferredCtx, items: Vec<PendingWebSocketOutbound>, reason: String) {
1253 for item in items {
1254 out.emit(WebSocketSessionEvent::Outbound(
1255 WebSocketSessionOutbound::Canceled {
1256 seq: item.seq,
1257 message: item.message,
1258 reason: reason.clone(),
1259 },
1260 ));
1261 }
1262}
1263
1264fn emit_outbound_rejected(out: &DeferredCtx, items: Vec<PendingWebSocketOutbound>, error: String) {
1265 for item in items {
1266 out.emit(WebSocketSessionEvent::Outbound(
1267 WebSocketSessionOutbound::Rejected {
1268 seq: item.seq,
1269 message: item.message,
1270 error: error.clone(),
1271 },
1272 ));
1273 }
1274}
1275
1276fn exhaust_websocket_session(
1277 state: &WebSocketSessionStateCell,
1278 out: &DeferredCtx,
1279 attempt: u32,
1280 error: String,
1281) {
1282 let cleanup = close_websocket_session_connection(state, true);
1283 emit_outbound_canceled(
1284 out,
1285 cleanup.canceled,
1286 "websocketSession: connection cleanup".to_owned(),
1287 );
1288 let pending = drain_pending_websocket_outbound(state);
1289 emit_outbound_rejected(out, pending, error.clone());
1290 out.emit(WebSocketSessionEvent::Exhausted { attempt, error });
1291}
1292
1293struct WebSocketConnectionCleanup {
1294 canceled: Vec<PendingWebSocketOutbound>,
1295 session: Option<Rc<dyn LocalWebSocketSession>>,
1296}
1297
1298fn close_websocket_session_connection(
1299 state: &WebSocketSessionStateCell,
1300 cancel_session: bool,
1301) -> WebSocketConnectionCleanup {
1302 let (session, send_cancels, canceled) = {
1303 let mut state = state.0.borrow_mut();
1304 state.connected = false;
1305 state.current_attempt = 0;
1306 let session = state.current_session.take();
1307 let send_cancels = state.send_cancels.drain(..).collect::<Vec<_>>();
1308 let canceled = state
1309 .live_outbound
1310 .drain(..)
1311 .map(|send| send.item)
1312 .collect();
1313 (session, send_cancels, canceled)
1314 };
1315 for send in send_cancels {
1316 if let Some(cancel) = send.slot.borrow_mut().take() {
1317 cancel();
1318 }
1319 }
1320 if cancel_session {
1321 if let Some(session) = &session {
1322 session.cancel();
1323 }
1324 }
1325 WebSocketConnectionCleanup { canceled, session }
1326}
1327
1328fn remove_websocket_send_cancel(state: &WebSocketSessionStateCell, send_id: u64) {
1329 let mut state = state.0.borrow_mut();
1330 if let Some(index) = state
1331 .send_cancels
1332 .iter()
1333 .position(|send| send.id == send_id)
1334 {
1335 let send = state.send_cancels.remove(index);
1336 let _ = send.slot.borrow_mut().take();
1337 }
1338}
1339
1340fn cancel_websocket_retry_timers(state: &WebSocketSessionStateCell) {
1341 let mut state = state.0.borrow_mut();
1342 state.waiting_retry = false;
1343 for slot in state.retry_cancels.drain(..) {
1344 if let Some(cancel) = slot.borrow_mut().take() {
1345 cancel();
1346 }
1347 }
1348}
1349
1350fn deactivate_websocket_session_state(state: &WebSocketSessionStateCell) {
1351 {
1352 let mut state_ref = state.0.borrow_mut();
1353 state_ref.node_active = false;
1354 }
1355 cancel_websocket_retry_timers(state);
1356 let _cleanup = close_websocket_session_connection(state, true);
1357}
1358
1359pub fn to_http<T, F>(
1361 graph: &Graph,
1362 source: &Node<T>,
1363 request_of: F,
1364) -> OutboundBundle<T, HttpResponse>
1365where
1366 T: Clone + 'static,
1367 F: Fn(&T) -> HttpRequest + 'static,
1368{
1369 to_http_with_options(graph, source, request_of, OutboundAdapterOptions::default())
1370}
1371
1372pub fn to_http_with_options<T, F>(
1374 graph: &Graph,
1375 source: &Node<T>,
1376 request_of: F,
1377 opts: OutboundAdapterOptions,
1378) -> OutboundBundle<T, HttpResponse>
1379where
1380 T: Clone + 'static,
1381 F: Fn(&T) -> HttpRequest + 'static,
1382{
1383 let name = opts.name.clone().unwrap_or_else(|| "toHttp".to_owned());
1384 let request_of = Rc::new(request_of);
1385 let events = outbound_node(graph, source, name.clone(), opts.retry, move |ctx| {
1386 let driver = ctx.environment().http_driver()?;
1387 let request_of = request_of.clone();
1388 Some(Rc::new(move |value: T, callback| {
1389 let request = request_of(&value);
1390 Some(driver.request(request, callback))
1391 }) as OutboundSend<T, HttpResponse>)
1392 });
1393 outbound_bundle(graph, events, name)
1394}
1395
1396pub fn to_process<T, F>(
1398 graph: &Graph,
1399 source: &Node<T>,
1400 command_of: F,
1401) -> OutboundBundle<T, ProcessResult>
1402where
1403 T: Clone + 'static,
1404 F: Fn(&T) -> ProcessCommand + 'static,
1405{
1406 to_process_with_options(graph, source, command_of, OutboundAdapterOptions::default())
1407}
1408
1409pub fn to_process_with_options<T, F>(
1411 graph: &Graph,
1412 source: &Node<T>,
1413 command_of: F,
1414 opts: OutboundAdapterOptions,
1415) -> OutboundBundle<T, ProcessResult>
1416where
1417 T: Clone + 'static,
1418 F: Fn(&T) -> ProcessCommand + 'static,
1419{
1420 let name = opts.name.clone().unwrap_or_else(|| "toProcess".to_owned());
1421 let command_of = Rc::new(command_of);
1422 let events = outbound_node(graph, source, name.clone(), opts.retry, move |ctx| {
1423 let driver = ctx.environment().process_driver()?;
1424 let command_of = command_of.clone();
1425 Some(Rc::new(move |value: T, callback| {
1426 let command = command_of(&value);
1427 Some(driver.run(command, callback))
1428 }) as OutboundSend<T, ProcessResult>)
1429 });
1430 outbound_bundle(graph, events, name)
1431}
1432
1433pub fn to_websocket<T, F>(
1435 graph: &Graph,
1436 source: &Node<T>,
1437 request: WebSocketRequest,
1438 send_of: F,
1439) -> OutboundBundle<T, WebSocketSendResult>
1440where
1441 T: Clone + 'static,
1442 F: Fn(&T) -> WebSocketSend + 'static,
1443{
1444 to_websocket_with_options(
1445 graph,
1446 source,
1447 request,
1448 send_of,
1449 OutboundAdapterOptions::default(),
1450 )
1451}
1452
1453pub fn to_websocket_with_options<T, F>(
1455 graph: &Graph,
1456 source: &Node<T>,
1457 request: WebSocketRequest,
1458 send_of: F,
1459 opts: OutboundAdapterOptions,
1460) -> OutboundBundle<T, WebSocketSendResult>
1461where
1462 T: Clone + 'static,
1463 F: Fn(&T) -> WebSocketSend + 'static,
1464{
1465 let name = opts
1466 .name
1467 .clone()
1468 .unwrap_or_else(|| "toWebSocket".to_owned());
1469 let request = Rc::new(request);
1470 let send_of = Rc::new(send_of);
1471 let events = outbound_node(graph, source, name.clone(), opts.retry, move |ctx| {
1472 let driver = ctx.environment().websocket_driver()?;
1473 let request = request.clone();
1474 let send_of = send_of.clone();
1475 Some(Rc::new(move |value: T, callback| {
1476 let message = send_of(&value);
1477 driver.send((*request).clone(), message, callback)
1478 }) as OutboundSend<T, WebSocketSendResult>)
1479 });
1480 outbound_bundle(graph, events, name)
1481}
1482
1483fn outbound_node<T, R, S>(
1484 graph: &Graph,
1485 source: &Node<T>,
1486 name: String,
1487 policy: RetryPolicy,
1488 send: S,
1489) -> Node<OutboundEvent<T, R>>
1490where
1491 T: Clone + 'static,
1492 R: Clone + 'static,
1493 S: Fn(&Ctx) -> Option<OutboundSend<T, R>> + 'static,
1494{
1495 let missing_driver_name = name.clone();
1496 graph.node_opts::<OutboundEvent<T, R>, _>(
1497 vec![source.erased()],
1498 move |ctx| {
1499 let Some(send) = send(ctx) else {
1500 ctx.down(vec![Message::Error(
1501 format!("{missing_driver_name}: missing environment driver").into(),
1502 )]);
1503 return;
1504 };
1505 let active = Rc::new(Cell::new(true));
1506 let cancels = Rc::new(RefCell::new(Vec::<CancelSlot>::new()));
1507 let cleanup_active = active.clone();
1508 let cleanup_cancels = cancels.clone();
1509 ctx.on_deactivation(move || {
1510 cleanup_active.set(false);
1511 for slot in cleanup_cancels.borrow_mut().drain(..) {
1512 if let Some(cancel) = slot.borrow_mut().take() {
1513 cancel();
1514 }
1515 }
1516 });
1517 let out = Rc::new(ctx.defer());
1518 let async_driver = ctx.local_async_driver();
1519 for value in ctx.batch::<T>(0) {
1520 start_attempt(StartArgs {
1521 value: (*value).clone(),
1522 attempt: 1,
1523 active: active.clone(),
1524 out: out.clone(),
1525 cancels: cancels.clone(),
1526 policy: policy.clone(),
1527 async_driver: async_driver.clone(),
1528 send: send.clone(),
1529 });
1530 }
1531 match ctx.terminal(0) {
1532 Some(DepTerminal::Complete) => ctx.emit(OutboundEvent::<T, R>::UpstreamComplete),
1533 Some(DepTerminal::Error(error)) => ctx.emit(OutboundEvent::<T, R>::UpstreamError {
1534 error: error.to_string(),
1535 }),
1536 None => {}
1537 }
1538 },
1539 GraphNodeOpts::named(name),
1540 )
1541}
1542
1543struct StartArgs<T: Clone + 'static, R: Clone + 'static> {
1544 value: T,
1545 attempt: u32,
1546 active: Rc<Cell<bool>>,
1547 out: Rc<DeferredCtx>,
1548 cancels: CancelSlots,
1549 policy: RetryPolicy,
1550 async_driver: Option<Rc<dyn LocalAsyncDriver>>,
1551 send: OutboundSend<T, R>,
1552}
1553
1554fn start_attempt<T, R>(args: StartArgs<T, R>)
1555where
1556 T: Clone + 'static,
1557 R: Clone + 'static,
1558{
1559 if !args.active.get() {
1560 return;
1561 }
1562 args.out.emit(OutboundEvent::<T, R>::Attempt {
1563 value: args.value.clone(),
1564 attempt: args.attempt,
1565 });
1566 let send_value = args.value.clone();
1567 let cancel_slot: CancelSlot = Rc::new(RefCell::new(None));
1568 let done = Rc::new(Cell::new(false));
1569 let callback_args = StartArgs {
1570 value: args.value.clone(),
1571 attempt: args.attempt,
1572 active: args.active.clone(),
1573 out: args.out.clone(),
1574 cancels: args.cancels.clone(),
1575 policy: args.policy.clone(),
1576 async_driver: args.async_driver.clone(),
1577 send: args.send.clone(),
1578 };
1579 let out = args.out.clone();
1580 let callback_slot = cancel_slot.clone();
1581 let callback_done = done.clone();
1582 let cancel = (args.send)(
1583 send_value,
1584 Box::new(move |result| {
1585 callback_done.set(true);
1586 let _ = callback_slot.borrow_mut().take();
1587 if !callback_args.active.get() {
1588 return;
1589 }
1590 match result {
1591 Ok(result) => callback_args.out.emit(OutboundEvent::<T, R>::Sent {
1592 value: callback_args.value,
1593 attempt: callback_args.attempt,
1594 result,
1595 }),
1596 Err(error) => {
1597 let error = error.to_string();
1598 if callback_args.policy.should_retry(callback_args.attempt) {
1599 let next_attempt = callback_args.attempt.saturating_add(1);
1600 let delay_ms = callback_args
1601 .policy
1602 .next_delay_ms(next_attempt)
1603 .unwrap_or_default();
1604 if delay_ms > 0 && callback_args.async_driver.is_none() {
1605 let scheduler_error =
1606 "outbound adapter: missing async driver for delayed retry"
1607 .to_owned();
1608 callback_args.out.emit(OutboundEvent::<T, R>::Exhausted {
1609 value: callback_args.value,
1610 attempt: callback_args.attempt,
1611 error: scheduler_error.clone(),
1612 });
1613 callback_args
1614 .out
1615 .down(vec![Message::Error(scheduler_error.into())]);
1616 return;
1617 }
1618 let active = callback_args.active.clone();
1619 let out = callback_args.out.clone();
1620 let cancels = callback_args.cancels.clone();
1621 let policy = callback_args.policy.clone();
1622 let async_driver = callback_args.async_driver.clone();
1623 let send = callback_args.send.clone();
1624 callback_args.out.emit(OutboundEvent::<T, R>::Retry {
1625 value: callback_args.value.clone(),
1626 attempt: callback_args.attempt,
1627 delay_ms,
1628 error: error.clone(),
1629 });
1630 let next_args = StartArgs {
1631 value: callback_args.value,
1632 attempt: next_attempt,
1633 active,
1634 out,
1635 cancels: cancels.clone(),
1636 policy,
1637 async_driver: async_driver.clone(),
1638 send,
1639 };
1640 if delay_ms == 0 {
1641 start_attempt(next_args);
1642 } else if let Some(driver) = async_driver {
1643 let sleep_slot: CancelSlot = Rc::new(RefCell::new(None));
1644 let wake_slot = sleep_slot.clone();
1645 let cancel = driver.sleep(
1646 Duration::from_millis(delay_ms),
1647 Box::new(move || {
1648 let _ = wake_slot.borrow_mut().take();
1649 start_attempt(next_args);
1650 }),
1651 );
1652 *sleep_slot.borrow_mut() = Some(cancel);
1653 cancels.borrow_mut().push(sleep_slot);
1654 }
1655 } else {
1656 callback_args.out.emit(OutboundEvent::<T, R>::Exhausted {
1657 value: callback_args.value,
1658 attempt: callback_args.attempt,
1659 error,
1660 });
1661 }
1662 }
1663 }
1664 }),
1665 );
1666 match cancel {
1667 Some(cancel) => {
1668 if !done.get() {
1669 *cancel_slot.borrow_mut() = Some(cancel);
1670 args.cancels.borrow_mut().push(cancel_slot);
1671 }
1672 }
1673 None => {
1674 if !done.get() {
1675 let error = "outbound adapter: missing driver send capability".to_owned();
1676 out.emit(OutboundEvent::<T, R>::Exhausted {
1677 value: args.value,
1678 attempt: args.attempt,
1679 error: error.clone(),
1680 });
1681 out.down(vec![Message::Error(error.into())]);
1682 }
1683 }
1684 }
1685}
1686
1687fn outbound_bundle<T, R>(
1688 graph: &Graph,
1689 events: Node<OutboundEvent<T, R>>,
1690 name: String,
1691) -> OutboundBundle<T, R>
1692where
1693 T: Clone + 'static,
1694 R: Clone + 'static,
1695{
1696 let status = graph.node_opts::<OutboundStatus, _>(
1697 vec![events.erased()],
1698 move |ctx| {
1699 let mut next = ctx
1700 .state_get::<OutboundStatus>()
1701 .map_or_else(OutboundStatus::default, |value| (*value).clone());
1702 for event in ctx.batch::<OutboundEvent<T, R>>(0) {
1703 next = reduce_status(&next, event.as_ref());
1704 }
1705 ctx.state_set(next.clone());
1706 ctx.emit(next);
1707 },
1708 GraphNodeOpts::named(format!("{name}/status")),
1709 );
1710 let attempts = graph.node_opts::<u32, _>(
1711 vec![events.erased()],
1712 move |ctx| {
1713 for event in ctx.batch::<OutboundEvent<T, R>>(0) {
1714 match event.as_ref() {
1715 OutboundEvent::Attempt { attempt, .. }
1716 | OutboundEvent::Retry { attempt, .. }
1717 | OutboundEvent::Sent { attempt, .. }
1718 | OutboundEvent::Exhausted { attempt, .. } => ctx.emit(*attempt),
1719 OutboundEvent::UpstreamComplete | OutboundEvent::UpstreamError { .. } => {}
1720 }
1721 }
1722 },
1723 GraphNodeOpts::named(format!("{name}/attempts")),
1724 );
1725 let errors = graph.node_opts::<String, _>(
1726 vec![events.erased()],
1727 move |ctx| {
1728 for event in ctx.batch::<OutboundEvent<T, R>>(0) {
1729 match event.as_ref() {
1730 OutboundEvent::Retry { error, .. }
1731 | OutboundEvent::Exhausted { error, .. }
1732 | OutboundEvent::UpstreamError { error } => ctx.emit(error.clone()),
1733 OutboundEvent::Attempt { .. }
1734 | OutboundEvent::Sent { .. }
1735 | OutboundEvent::UpstreamComplete => {}
1736 }
1737 }
1738 },
1739 GraphNodeOpts::named(format!("{name}/errors")),
1740 );
1741 OutboundBundle {
1742 events,
1743 status,
1744 attempts,
1745 errors,
1746 }
1747}
1748
1749fn reduce_status<T, R>(current: &OutboundStatus, event: &OutboundEvent<T, R>) -> OutboundStatus {
1750 match event {
1751 OutboundEvent::Attempt { attempt, .. } => OutboundStatus {
1752 state: OutboundState::Running,
1753 in_flight: current.in_flight.saturating_add(1),
1754 attempt: *attempt,
1755 sent: current.sent,
1756 failed: current.failed,
1757 last_delay_ms: current.last_delay_ms,
1758 },
1759 OutboundEvent::Retry {
1760 attempt, delay_ms, ..
1761 } => OutboundStatus {
1762 state: OutboundState::Waiting,
1763 in_flight: current.in_flight.saturating_sub(1),
1764 attempt: *attempt,
1765 sent: current.sent,
1766 failed: current.failed,
1767 last_delay_ms: Some(*delay_ms),
1768 },
1769 OutboundEvent::Sent { attempt, .. } => OutboundStatus {
1770 state: OutboundState::Succeeded,
1771 in_flight: current.in_flight.saturating_sub(1),
1772 attempt: *attempt,
1773 sent: current.sent.saturating_add(1),
1774 failed: current.failed,
1775 last_delay_ms: None,
1776 },
1777 OutboundEvent::Exhausted { attempt, .. } => OutboundStatus {
1778 state: OutboundState::Exhausted,
1779 in_flight: current.in_flight.saturating_sub(1),
1780 attempt: *attempt,
1781 sent: current.sent,
1782 failed: current.failed.saturating_add(1),
1783 last_delay_ms: current.last_delay_ms,
1784 },
1785 OutboundEvent::UpstreamComplete => OutboundStatus {
1786 state: OutboundState::Completed,
1787 ..current.clone()
1788 },
1789 OutboundEvent::UpstreamError { .. } => OutboundStatus {
1790 state: OutboundState::Failed,
1791 failed: current.failed.saturating_add(1),
1792 ..current.clone()
1793 },
1794 }
1795}
1796
1797#[cfg(test)]
1798mod tests {
1799 use super::*;
1800 use crate::environment::{
1801 EnvironmentDrivers, LocalHttpDriver, LocalWebSocketDriver, LocalWebSocketSession,
1802 WebSocketDriverEvent, WebSocketEvent,
1803 };
1804 use crate::graph::{graph_opts, GraphNodeOpts, GraphOptions};
1805
1806 struct ManualHttpDriver {
1807 requests: Rc<RefCell<Vec<HttpRequest>>>,
1808 }
1809
1810 impl LocalHttpDriver for ManualHttpDriver {
1811 fn request(
1812 &self,
1813 request: HttpRequest,
1814 callback: Box<dyn FnOnce(Result<HttpResponse, GraphError>)>,
1815 ) -> DriverCancel {
1816 self.requests.borrow_mut().push(request);
1817 callback(Ok(HttpResponse {
1818 status: 204,
1819 headers: Vec::new(),
1820 body: Vec::new(),
1821 }));
1822 Box::new(|| {})
1823 }
1824 }
1825
1826 struct ConnectOnlyWebSocketDriver;
1827
1828 impl LocalWebSocketDriver for ConnectOnlyWebSocketDriver {
1829 fn connect(
1830 &self,
1831 _request: WebSocketRequest,
1832 _callback: Rc<dyn Fn(WebSocketDriverEvent)>,
1833 ) -> DriverCancel {
1834 Box::new(|| {})
1835 }
1836 }
1837
1838 #[derive(Default)]
1839 struct ManualSessionWebSocketDriver {
1840 sessions: RefCell<Vec<Rc<ManualWebSocketSession>>>,
1841 }
1842
1843 impl ManualSessionWebSocketDriver {
1844 fn session(&self, index: usize) -> Rc<ManualWebSocketSession> {
1845 self.sessions.borrow()[index].clone()
1846 }
1847 }
1848
1849 impl LocalWebSocketDriver for ManualSessionWebSocketDriver {
1850 fn connect(
1851 &self,
1852 _request: WebSocketRequest,
1853 _callback: Rc<dyn Fn(WebSocketDriverEvent)>,
1854 ) -> DriverCancel {
1855 Box::new(|| {})
1856 }
1857
1858 fn connect_session(
1859 &self,
1860 request: WebSocketRequest,
1861 callback: Rc<dyn Fn(WebSocketDriverEvent)>,
1862 ) -> Option<Rc<dyn LocalWebSocketSession>> {
1863 let session = Rc::new(ManualWebSocketSession {
1864 url: request.url,
1865 active: Cell::new(true),
1866 callback,
1867 sent: RefCell::new(Vec::new()),
1868 closes: RefCell::new(Vec::new()),
1869 });
1870 self.sessions.borrow_mut().push(session.clone());
1871 Some(session)
1872 }
1873 }
1874
1875 struct ManualWebSocketSession {
1876 url: String,
1877 active: Cell<bool>,
1878 callback: Rc<dyn Fn(WebSocketDriverEvent)>,
1879 sent: RefCell<Vec<WebSocketSend>>,
1880 closes: RefCell<Vec<(Option<u16>, Option<String>)>>,
1881 }
1882
1883 impl ManualWebSocketSession {
1884 fn emit(&self, event: WebSocketDriverEvent) {
1885 if self.active.get() {
1886 (self.callback)(event);
1887 }
1888 }
1889
1890 fn emit_ignoring_cancel(&self, event: WebSocketDriverEvent) {
1891 (self.callback)(event);
1892 }
1893 }
1894
1895 #[derive(Default)]
1896 struct PendingSendWebSocketDriver {
1897 sessions: RefCell<Vec<Rc<PendingSendWebSocketSession>>>,
1898 }
1899
1900 impl PendingSendWebSocketDriver {
1901 fn session(&self, index: usize) -> Rc<PendingSendWebSocketSession> {
1902 self.sessions.borrow()[index].clone()
1903 }
1904 }
1905
1906 impl LocalWebSocketDriver for PendingSendWebSocketDriver {
1907 fn connect(
1908 &self,
1909 _request: WebSocketRequest,
1910 _callback: Rc<dyn Fn(WebSocketDriverEvent)>,
1911 ) -> DriverCancel {
1912 Box::new(|| {})
1913 }
1914
1915 fn connect_session(
1916 &self,
1917 request: WebSocketRequest,
1918 callback: Rc<dyn Fn(WebSocketDriverEvent)>,
1919 ) -> Option<Rc<dyn LocalWebSocketSession>> {
1920 let session = Rc::new(PendingSendWebSocketSession {
1921 url: request.url,
1922 callback,
1923 sent: RefCell::new(Vec::new()),
1924 send_callbacks: RefCell::new(Vec::new()),
1925 closes: RefCell::new(Vec::new()),
1926 canceled_sends: Rc::new(Cell::new(0)),
1927 });
1928 self.sessions.borrow_mut().push(session.clone());
1929 Some(session)
1930 }
1931 }
1932
1933 type PendingSendCallback = Box<dyn FnOnce(Result<WebSocketSendResult, GraphError>)>;
1934
1935 struct PendingSendWebSocketSession {
1936 url: String,
1937 callback: Rc<dyn Fn(WebSocketDriverEvent)>,
1938 sent: RefCell<Vec<WebSocketSend>>,
1939 send_callbacks: RefCell<Vec<PendingSendCallback>>,
1940 closes: RefCell<Vec<(Option<u16>, Option<String>)>>,
1941 canceled_sends: Rc<Cell<u32>>,
1942 }
1943
1944 impl PendingSendWebSocketSession {
1945 fn emit(&self, event: WebSocketDriverEvent) {
1946 (self.callback)(event);
1947 }
1948
1949 fn complete_next_send(&self, result: Result<WebSocketSendResult, GraphError>) {
1950 let callback = self.send_callbacks.borrow_mut().remove(0);
1951 callback(result);
1952 }
1953 }
1954
1955 impl LocalWebSocketSession for PendingSendWebSocketSession {
1956 fn send(
1957 &self,
1958 message: WebSocketSend,
1959 callback: Box<dyn FnOnce(Result<WebSocketSendResult, GraphError>)>,
1960 ) -> DriverCancel {
1961 self.sent.borrow_mut().push(message);
1962 self.send_callbacks.borrow_mut().push(callback);
1963 let canceled = self.canceled_sends.clone();
1964 Box::new(move || {
1965 canceled.set(canceled.get().saturating_add(1));
1966 })
1967 }
1968
1969 fn close(&self, code: Option<u16>, reason: Option<String>) {
1970 self.closes.borrow_mut().push((code, reason));
1971 }
1972
1973 fn cancel(&self) {}
1974 }
1975
1976 impl LocalWebSocketSession for ManualWebSocketSession {
1977 fn send(
1978 &self,
1979 message: WebSocketSend,
1980 callback: Box<dyn FnOnce(Result<WebSocketSendResult, GraphError>)>,
1981 ) -> DriverCancel {
1982 if self.active.get() {
1983 self.sent.borrow_mut().push(message);
1984 callback(Ok(WebSocketSendResult { sent: true }));
1985 }
1986 Box::new(|| {})
1987 }
1988
1989 fn close(&self, code: Option<u16>, reason: Option<String>) {
1990 self.closes.borrow_mut().push((code, reason));
1991 self.active.set(false);
1992 }
1993
1994 fn cancel(&self) {
1995 self.active.set(false);
1996 }
1997 }
1998
1999 struct FailOnceWebSocketDriver {
2000 attempts: Cell<u32>,
2001 }
2002
2003 impl LocalWebSocketDriver for FailOnceWebSocketDriver {
2004 fn connect(
2005 &self,
2006 _request: WebSocketRequest,
2007 _callback: Rc<dyn Fn(WebSocketDriverEvent)>,
2008 ) -> DriverCancel {
2009 Box::new(|| {})
2010 }
2011
2012 fn connect_session(
2013 &self,
2014 _request: WebSocketRequest,
2015 callback: Rc<dyn Fn(WebSocketDriverEvent)>,
2016 ) -> Option<Rc<dyn LocalWebSocketSession>> {
2017 let attempt = self.attempts.get().saturating_add(1);
2018 self.attempts.set(attempt);
2019 let session = Rc::new(ManualWebSocketSession {
2020 url: format!("attempt-{attempt}"),
2021 active: Cell::new(true),
2022 callback: callback.clone(),
2023 sent: RefCell::new(Vec::new()),
2024 closes: RefCell::new(Vec::new()),
2025 });
2026 if attempt == 1 {
2027 callback(WebSocketDriverEvent::Error("boom".into()));
2028 } else {
2029 callback(WebSocketDriverEvent::Event(WebSocketEvent::Open));
2030 }
2031 Some(session)
2032 }
2033 }
2034
2035 type PendingSleep = (Rc<Cell<bool>>, Box<dyn FnOnce()>);
2036
2037 #[derive(Default)]
2038 struct ManualAsyncDriver {
2039 sleeps: RefCell<Vec<PendingSleep>>,
2040 }
2041
2042 impl ManualAsyncDriver {
2043 fn fire_next(&self) {
2044 let Some((active, callback)) = self.sleeps.borrow_mut().pop() else {
2045 panic!("expected a pending sleep");
2046 };
2047 if active.get() {
2048 callback();
2049 }
2050 }
2051 }
2052
2053 impl LocalAsyncDriver for ManualAsyncDriver {
2054 fn sleep(&self, _duration: Duration, callback: Box<dyn FnOnce()>) -> DriverCancel {
2055 let active = Rc::new(Cell::new(true));
2056 self.sleeps.borrow_mut().push((active.clone(), callback));
2057 Box::new(move || active.set(false))
2058 }
2059
2060 fn interval(&self, _period: Duration, _callback: Rc<dyn Fn()>) -> DriverCancel {
2061 Box::new(|| {})
2062 }
2063
2064 fn spawn_local(
2065 &self,
2066 _fut: std::pin::Pin<Box<dyn std::future::Future<Output = ()> + 'static>>,
2067 ) -> DriverCancel {
2068 Box::new(|| {})
2069 }
2070 }
2071
2072 fn collect_node_data<T: Clone + 'static>(node: &Node<T>) -> Rc<RefCell<Vec<T>>> {
2073 let seen = Rc::new(RefCell::new(Vec::new()));
2074 let seen_sink = seen.clone();
2075 let _keep = node.subscribe(move |msg| {
2076 if let Message::Data(value) = msg {
2077 if let Some(value) = value.as_ref().downcast_ref::<T>() {
2078 seen_sink.borrow_mut().push(value.clone());
2079 }
2080 }
2081 });
2082 seen
2083 }
2084
2085 #[test]
2086 fn to_http_emits_graph_visible_attempt_sent_and_status() {
2087 let requests = Rc::new(RefCell::new(Vec::new()));
2088 let environment = EnvironmentDrivers::new().with_http(Rc::new(ManualHttpDriver {
2089 requests: requests.clone(),
2090 }));
2091 let g = graph_opts(GraphOptions {
2092 environment,
2093 ..GraphOptions::default()
2094 });
2095 let source = g.state_empty_opts::<String>(GraphNodeOpts::named("source"));
2096 let bundle = to_http(&g, &source, |value| {
2097 HttpRequest::new("POST", format!("https://example.test/{value}"))
2098 });
2099 let _events = bundle.events.subscribe(|_| {});
2100 let _status = bundle.status.subscribe(|_| {});
2101
2102 source.set("order".to_owned());
2103
2104 assert_eq!(requests.borrow()[0].url, "https://example.test/order");
2105 assert_eq!(
2106 bundle.events.cache(),
2107 Some(OutboundEvent::Sent {
2108 value: "order".to_owned(),
2109 attempt: 1,
2110 result: HttpResponse {
2111 status: 204,
2112 headers: Vec::new(),
2113 body: Vec::new(),
2114 },
2115 })
2116 );
2117 assert_eq!(
2118 bundle.status.cache(),
2119 Some(OutboundStatus {
2120 state: OutboundState::Succeeded,
2121 in_flight: 0,
2122 attempt: 1,
2123 sent: 1,
2124 failed: 0,
2125 last_delay_ms: None,
2126 })
2127 );
2128 let snap = g.describe();
2129 assert!(snap
2130 .edges
2131 .iter()
2132 .any(|edge| edge.from == "source" && edge.to == "toHttp"));
2133 }
2134
2135 #[test]
2136 fn missing_send_capability_closes_status_ledger() {
2137 let environment =
2138 EnvironmentDrivers::new().with_websocket(Rc::new(ConnectOnlyWebSocketDriver));
2139 let g = graph_opts(GraphOptions {
2140 environment,
2141 ..GraphOptions::default()
2142 });
2143 let source = g.state_empty_opts::<String>(GraphNodeOpts::named("source"));
2144 let bundle = to_websocket(
2145 &g,
2146 &source,
2147 WebSocketRequest::new("wss://example.test"),
2148 |value| WebSocketSend::text(value.clone()),
2149 );
2150 let _events = bundle.events.subscribe(|_| {});
2151 let _status = bundle.status.subscribe(|_| {});
2152
2153 source.set("order".to_owned());
2154
2155 assert_eq!(
2156 bundle.events.cache(),
2157 Some(OutboundEvent::Exhausted {
2158 value: "order".to_owned(),
2159 attempt: 1,
2160 error: "outbound adapter: missing driver send capability".to_owned(),
2161 })
2162 );
2163 assert_eq!(
2164 bundle.status.cache(),
2165 Some(OutboundStatus {
2166 state: OutboundState::Exhausted,
2167 in_flight: 0,
2168 attempt: 1,
2169 sent: 0,
2170 failed: 1,
2171 last_delay_ms: None,
2172 })
2173 );
2174 }
2175
2176 #[test]
2177 fn websocket_session_missing_session_capability_is_graph_visible() {
2178 let environment =
2179 EnvironmentDrivers::new().with_websocket(Rc::new(ConnectOnlyWebSocketDriver));
2180 let g = graph_opts(GraphOptions {
2181 environment,
2182 ..GraphOptions::default()
2183 });
2184 let bundle = websocket_session(&g, WebSocketRequest::new("wss://example.test/session"));
2185 let errors = collect_node_data(&bundle.errors);
2186 let attempts = collect_node_data(&bundle.attempts);
2187 let lifecycle = collect_node_data(&bundle.lifecycle);
2188 let _status = bundle.status.subscribe(|_| {});
2189
2190 bundle.start();
2191
2192 assert_eq!(*attempts.borrow(), vec![1]);
2193 assert_eq!(
2194 *errors.borrow(),
2195 vec!["websocketSession: missing WebSocket session capability".to_owned()]
2196 );
2197 assert_eq!(
2198 bundle.status.cache(),
2199 Some(WebSocketSessionStatus {
2200 state: WebSocketSessionStateKind::Exhausted,
2201 attempt: 1,
2202 max_attempts: 1,
2203 sent: 0,
2204 received: 0,
2205 errors: 1,
2206 last_delay_ms: None,
2207 })
2208 );
2209 assert!(matches!(
2210 lifecycle.borrow().last(),
2211 Some(WebSocketSessionLifecycle::Exhausted { attempt: 1, .. })
2212 ));
2213 }
2214
2215 #[test]
2216 fn websocket_session_helpers_publish_command_facts_only() {
2217 let g = graph_opts(GraphOptions::default());
2218 let bundle = websocket_session(&g, WebSocketRequest::new("wss://example.test/session"));
2219 let commands = collect_node_data(&bundle.command);
2220
2221 bundle.start();
2222 bundle.send_text("hello");
2223 bundle.close(Some(1000), Some("done".to_owned()));
2224
2225 assert_eq!(
2226 *commands.borrow(),
2227 vec![
2228 WebSocketSessionCommand::Start,
2229 WebSocketSessionCommand::Send(WebSocketSend::text("hello")),
2230 WebSocketSessionCommand::Close {
2231 code: Some(1000),
2232 reason: Some("done".to_owned()),
2233 },
2234 ]
2235 );
2236 }
2237
2238 #[test]
2239 fn websocket_session_default_send_before_open_rejects_outbound() {
2240 let driver = Rc::new(ManualSessionWebSocketDriver::default());
2241 let g = graph_opts(GraphOptions {
2242 environment: EnvironmentDrivers::new().with_websocket(driver),
2243 ..GraphOptions::default()
2244 });
2245 let bundle = websocket_session_with_options(
2246 &g,
2247 WebSocketRequest::new("wss://example.test/session"),
2248 WebSocketSessionOptions {
2249 name: Some("ws".to_owned()),
2250 ..WebSocketSessionOptions::default()
2251 },
2252 );
2253 let outbound = collect_node_data(&bundle.outbound);
2254 let errors = collect_node_data(&bundle.errors);
2255 let _status = bundle.status.subscribe(|_| {});
2256
2257 bundle.send_text("early");
2258
2259 assert_eq!(
2260 *outbound.borrow(),
2261 vec![WebSocketSessionOutbound::Rejected {
2262 seq: 0,
2263 message: WebSocketSend::text("early"),
2264 error: "ws: session is not open".to_owned(),
2265 }]
2266 );
2267 assert_eq!(*errors.borrow(), vec!["ws: session is not open".to_owned()]);
2268 assert_eq!(
2269 bundle.status.cache(),
2270 Some(WebSocketSessionStatus {
2271 state: WebSocketSessionStateKind::Errored,
2272 attempt: 0,
2273 max_attempts: 1,
2274 sent: 0,
2275 received: 0,
2276 errors: 1,
2277 last_delay_ms: None,
2278 })
2279 );
2280 }
2281
2282 #[test]
2283 fn websocket_session_buffer_policy_is_bounded_and_flushes_fifo_on_open() {
2284 let driver = Rc::new(ManualSessionWebSocketDriver::default());
2285 let g = graph_opts(GraphOptions {
2286 environment: EnvironmentDrivers::new().with_websocket(driver.clone()),
2287 ..GraphOptions::default()
2288 });
2289 let bundle = websocket_session_with_options(
2290 &g,
2291 WebSocketRequest::new("wss://example.test/session"),
2292 WebSocketSessionOptions {
2293 name: Some("ws".to_owned()),
2294 send_policy: WebSocketSessionSendPolicy::Buffer { max_pending: 2 },
2295 ..WebSocketSessionOptions::default()
2296 },
2297 );
2298 let outbound = collect_node_data(&bundle.outbound);
2299
2300 bundle.send_text("a");
2301 bundle.send_text("b");
2302 bundle.send_text("c");
2303 bundle.start();
2304 let session = driver.session(0);
2305
2306 assert_eq!(
2307 *outbound.borrow(),
2308 vec![
2309 WebSocketSessionOutbound::Queued {
2310 seq: 0,
2311 message: WebSocketSend::text("a"),
2312 },
2313 WebSocketSessionOutbound::Queued {
2314 seq: 1,
2315 message: WebSocketSend::text("b"),
2316 },
2317 WebSocketSessionOutbound::Rejected {
2318 seq: 2,
2319 message: WebSocketSend::text("c"),
2320 error: "ws: outbound buffer full".to_owned(),
2321 },
2322 ]
2323 );
2324 assert!(session.sent.borrow().is_empty());
2325
2326 session.emit(WebSocketDriverEvent::Event(WebSocketEvent::Open));
2327
2328 assert_eq!(
2329 session.sent.borrow().as_slice(),
2330 &[WebSocketSend::text("a"), WebSocketSend::text("b")]
2331 );
2332 assert_eq!(
2333 *outbound.borrow(),
2334 vec![
2335 WebSocketSessionOutbound::Queued {
2336 seq: 0,
2337 message: WebSocketSend::text("a"),
2338 },
2339 WebSocketSessionOutbound::Queued {
2340 seq: 1,
2341 message: WebSocketSend::text("b"),
2342 },
2343 WebSocketSessionOutbound::Rejected {
2344 seq: 2,
2345 message: WebSocketSend::text("c"),
2346 error: "ws: outbound buffer full".to_owned(),
2347 },
2348 WebSocketSessionOutbound::Sending {
2349 seq: 0,
2350 message: WebSocketSend::text("a"),
2351 },
2352 WebSocketSessionOutbound::Sent {
2353 seq: 0,
2354 message: WebSocketSend::text("a"),
2355 },
2356 WebSocketSessionOutbound::Sending {
2357 seq: 1,
2358 message: WebSocketSend::text("b"),
2359 },
2360 WebSocketSessionOutbound::Sent {
2361 seq: 1,
2362 message: WebSocketSend::text("b"),
2363 },
2364 ]
2365 );
2366 }
2367
2368 #[test]
2369 fn websocket_session_close_cancels_pending_and_fences_late_send_callback() {
2370 let driver = Rc::new(PendingSendWebSocketDriver::default());
2371 let g = graph_opts(GraphOptions {
2372 environment: EnvironmentDrivers::new().with_websocket(driver.clone()),
2373 ..GraphOptions::default()
2374 });
2375 let bundle = websocket_session_with_options(
2376 &g,
2377 WebSocketRequest::new("wss://example.test/session"),
2378 WebSocketSessionOptions {
2379 name: Some("ws".to_owned()),
2380 send_policy: WebSocketSessionSendPolicy::Buffer { max_pending: 1 },
2381 ..WebSocketSessionOptions::default()
2382 },
2383 );
2384 let outbound = collect_node_data(&bundle.outbound);
2385 let lifecycle = collect_node_data(&bundle.lifecycle);
2386
2387 bundle.send_text("queued");
2388 bundle.start();
2389 let session = driver.session(0);
2390 assert_eq!(session.url, "wss://example.test/session");
2391 session.emit(WebSocketDriverEvent::Event(WebSocketEvent::Open));
2392 bundle.send_text("live");
2393 bundle.close(Some(1000), Some("done".to_owned()));
2394 session.complete_next_send(Ok(WebSocketSendResult { sent: true }));
2395
2396 assert_eq!(session.canceled_sends.get(), 2);
2397 assert!(outbound
2398 .borrow()
2399 .contains(&WebSocketSessionOutbound::Canceled {
2400 seq: 1,
2401 message: WebSocketSend::text("live"),
2402 reason: "ws: session closed".to_owned(),
2403 }));
2404 assert!(
2405 !outbound.borrow().contains(&WebSocketSessionOutbound::Sent {
2406 seq: 1,
2407 message: WebSocketSend::text("live"),
2408 })
2409 );
2410 assert!(!lifecycle
2411 .borrow()
2412 .contains(&WebSocketSessionLifecycle::Sent {
2413 message: WebSocketSend::text("live"),
2414 }));
2415 }
2416
2417 #[test]
2418 fn websocket_session_uses_same_connection_send_close_and_fences_late_callbacks() {
2419 let driver = Rc::new(ManualSessionWebSocketDriver::default());
2420 let g = graph_opts(GraphOptions {
2421 environment: EnvironmentDrivers::new().with_websocket(driver.clone()),
2422 ..GraphOptions::default()
2423 });
2424 let bundle = websocket_session_with_options(
2425 &g,
2426 WebSocketRequest::new("wss://example.test/session"),
2427 WebSocketSessionOptions {
2428 name: Some("ws".to_owned()),
2429 ..WebSocketSessionOptions::default()
2430 },
2431 );
2432 let inbound = collect_node_data(&bundle.inbound);
2433 let lifecycle = collect_node_data(&bundle.lifecycle);
2434 let _status = bundle.status.subscribe(|_| {});
2435
2436 bundle.start();
2437 let session = driver.session(0);
2438 assert_eq!(session.url, "wss://example.test/session");
2439 session.emit(WebSocketDriverEvent::Event(WebSocketEvent::Open));
2440 bundle.send_text("hello");
2441 session.emit(WebSocketDriverEvent::Event(WebSocketEvent::Text(
2442 "server".to_owned(),
2443 )));
2444 bundle.close(Some(1000), Some("done".to_owned()));
2445 session.emit_ignoring_cancel(WebSocketDriverEvent::Event(WebSocketEvent::Text(
2446 "late".to_owned(),
2447 )));
2448
2449 assert_eq!(
2450 session.sent.borrow().as_slice(),
2451 &[WebSocketSend::text("hello")]
2452 );
2453 assert_eq!(
2454 session.closes.borrow().as_slice(),
2455 &[(Some(1000), Some("done".to_owned()))]
2456 );
2457 assert_eq!(
2458 *inbound.borrow(),
2459 vec![WebSocketSessionInbound::Text("server".to_owned())]
2460 );
2461 assert!(lifecycle
2462 .borrow()
2463 .contains(&WebSocketSessionLifecycle::Sent {
2464 message: WebSocketSend::text("hello"),
2465 }));
2466 assert_eq!(
2467 bundle.status.cache(),
2468 Some(WebSocketSessionStatus {
2469 state: WebSocketSessionStateKind::Closed,
2470 attempt: 1,
2471 max_attempts: 1,
2472 sent: 1,
2473 received: 1,
2474 errors: 0,
2475 last_delay_ms: None,
2476 })
2477 );
2478 }
2479
2480 #[test]
2481 fn one_driver_serves_multiple_independent_websocket_sessions() {
2482 let driver = Rc::new(ManualSessionWebSocketDriver::default());
2483 let g = graph_opts(GraphOptions {
2484 environment: EnvironmentDrivers::new().with_websocket(driver.clone()),
2485 ..GraphOptions::default()
2486 });
2487 let first = websocket_session_with_options(
2488 &g,
2489 WebSocketRequest::new("wss://example.test/a"),
2490 WebSocketSessionOptions {
2491 name: Some("first".to_owned()),
2492 ..WebSocketSessionOptions::default()
2493 },
2494 );
2495 let second = websocket_session_with_options(
2496 &g,
2497 WebSocketRequest::new("wss://example.test/b"),
2498 WebSocketSessionOptions {
2499 name: Some("second".to_owned()),
2500 ..WebSocketSessionOptions::default()
2501 },
2502 );
2503 let first_inbound = collect_node_data(&first.inbound);
2504 let second_inbound = collect_node_data(&second.inbound);
2505 let _first_status = first.status.subscribe(|_| {});
2506 let _second_status = second.status.subscribe(|_| {});
2507
2508 first.start();
2509 second.start();
2510 let first_session = driver.session(0);
2511 let second_session = driver.session(1);
2512 first_session.emit(WebSocketDriverEvent::Event(WebSocketEvent::Open));
2513 second_session.emit(WebSocketDriverEvent::Event(WebSocketEvent::Open));
2514 first_session.emit(WebSocketDriverEvent::Event(WebSocketEvent::Text(
2515 "a".to_owned(),
2516 )));
2517 second_session.emit(WebSocketDriverEvent::Event(WebSocketEvent::Text(
2518 "b".to_owned(),
2519 )));
2520 first.close(Some(1000), None);
2521 second_session.emit(WebSocketDriverEvent::Event(WebSocketEvent::Text(
2522 "b2".to_owned(),
2523 )));
2524
2525 assert_eq!(
2526 *first_inbound.borrow(),
2527 vec![WebSocketSessionInbound::Text("a".to_owned())]
2528 );
2529 assert_eq!(
2530 *second_inbound.borrow(),
2531 vec![
2532 WebSocketSessionInbound::Text("b".to_owned()),
2533 WebSocketSessionInbound::Text("b2".to_owned()),
2534 ]
2535 );
2536 assert_eq!(driver.sessions.borrow().len(), 2);
2537 assert!(!first_session.active.get());
2538 assert!(second_session.active.get());
2539 }
2540
2541 #[test]
2542 fn websocket_session_retry_is_visible_and_bounded() {
2543 let driver = Rc::new(FailOnceWebSocketDriver {
2544 attempts: Cell::new(0),
2545 });
2546 let g = graph_opts(GraphOptions {
2547 environment: EnvironmentDrivers::new().with_websocket(driver.clone()),
2548 ..GraphOptions::default()
2549 });
2550 let bundle = websocket_session_with_options(
2551 &g,
2552 WebSocketRequest::new("wss://example.test/retry"),
2553 WebSocketSessionOptions {
2554 retry: RetryPolicy::new(2, crate::resilience::BackoffPolicy::None),
2555 ..WebSocketSessionOptions::default()
2556 },
2557 );
2558 let attempts = collect_node_data(&bundle.attempts);
2559 let errors = collect_node_data(&bundle.errors);
2560 let lifecycle = collect_node_data(&bundle.lifecycle);
2561 let _status = bundle.status.subscribe(|_| {});
2562
2563 bundle.start();
2564
2565 assert_eq!(*attempts.borrow(), vec![1, 2]);
2566 assert_eq!(*errors.borrow(), vec!["boom".to_owned()]);
2567 assert!(lifecycle.borrow().iter().any(|event| matches!(
2568 event,
2569 WebSocketSessionLifecycle::Retrying {
2570 attempt: 1,
2571 next_attempt: 2,
2572 ..
2573 }
2574 )));
2575 assert_eq!(
2576 bundle.status.cache(),
2577 Some(WebSocketSessionStatus {
2578 state: WebSocketSessionStateKind::Open,
2579 attempt: 2,
2580 max_attempts: 2,
2581 sent: 0,
2582 received: 0,
2583 errors: 1,
2584 last_delay_ms: None,
2585 })
2586 );
2587 assert_eq!(driver.attempts.get(), 2);
2588 }
2589
2590 #[test]
2591 fn websocket_session_delayed_retry_reconnects_after_timer() {
2592 let driver = Rc::new(FailOnceWebSocketDriver {
2593 attempts: Cell::new(0),
2594 });
2595 let async_driver = Rc::new(ManualAsyncDriver::default());
2596 let g = graph_opts(GraphOptions {
2597 environment: EnvironmentDrivers::new()
2598 .with_websocket(driver.clone())
2599 .with_local_async(async_driver.clone()),
2600 ..GraphOptions::default()
2601 });
2602 let bundle = websocket_session_with_options(
2603 &g,
2604 WebSocketRequest::new("wss://example.test/retry"),
2605 WebSocketSessionOptions {
2606 retry: RetryPolicy::new(
2607 2,
2608 crate::resilience::BackoffPolicy::Constant { delay_ms: 10 },
2609 ),
2610 ..WebSocketSessionOptions::default()
2611 },
2612 );
2613 let attempts = collect_node_data(&bundle.attempts);
2614 let _status = bundle.status.subscribe(|_| {});
2615
2616 bundle.start();
2617
2618 assert_eq!(*attempts.borrow(), vec![1]);
2619 assert_eq!(
2620 bundle.status.cache().map(|status| status.state),
2621 Some(WebSocketSessionStateKind::Waiting)
2622 );
2623
2624 async_driver.fire_next();
2625
2626 assert_eq!(*attempts.borrow(), vec![1, 2]);
2627 assert_eq!(
2628 bundle.status.cache().map(|status| status.state),
2629 Some(WebSocketSessionStateKind::Open)
2630 );
2631 assert_eq!(driver.attempts.get(), 2);
2632 }
2633}