Skip to main content

graphrefly/adapters/
environment.rs

1//! Graph-visible outbound environment adapters (D132).
2//!
3//! Transport work stays behind graph-local EnvironmentDrivers. These helpers
4//! expose attempt/status/error facts through declared graph deps.
5
6use 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)]
27/// `OutboundEvent` variants.
28pub enum OutboundEvent<T, R> {
29    /// `Attempt` variant.
30    Attempt {
31        /// `value` field for value.
32        value: T,
33        /// `attempt` field for attempt.
34        attempt: u32,
35    },
36    /// `Retry` variant.
37    Retry {
38        /// `value` field for value.
39        value: T,
40        /// `attempt` field for attempt.
41        attempt: u32,
42        /// `delay_ms` field for delay ms.
43        delay_ms: u64,
44        /// `error` field for error.
45        error: String,
46    },
47    /// `Sent` variant.
48    Sent {
49        /// `value` field for value.
50        value: T,
51        /// `attempt` field for attempt.
52        attempt: u32,
53        /// `result` field for result.
54        result: R,
55    },
56    /// `Exhausted` variant.
57    Exhausted {
58        /// `value` field for value.
59        value: T,
60        /// `attempt` field for attempt.
61        attempt: u32,
62        /// `error` field for error.
63        error: String,
64    },
65    /// `UpstreamComplete` variant.
66    UpstreamComplete,
67    /// `UpstreamError` variant.
68    UpstreamError {
69        /// `error` field for error.
70        error: String,
71    },
72}
73
74#[derive(Debug, Clone, PartialEq, Eq)]
75/// `OutboundState` variants.
76pub enum OutboundState {
77    /// `Idle` variant.
78    Idle,
79    /// `Running` variant.
80    Running,
81    /// `Waiting` variant.
82    Waiting,
83    /// `Succeeded` variant.
84    Succeeded,
85    /// `Exhausted` variant.
86    Exhausted,
87    /// `Failed` variant.
88    Failed,
89    /// `Completed` variant.
90    Completed,
91}
92
93#[derive(Debug, Clone, PartialEq, Eq)]
94/// `OutboundStatus` data container.
95pub struct OutboundStatus {
96    /// `state` field for state.
97    pub state: OutboundState,
98    /// `in_flight` field for in flight.
99    pub in_flight: u32,
100    /// `attempt` field for attempt.
101    pub attempt: u32,
102    /// `sent` field for sent.
103    pub sent: u64,
104    /// `failed` field for failed.
105    pub failed: u64,
106    /// `last_delay_ms` field for last delay ms.
107    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
123/// `OutboundBundle` data container.
124pub struct OutboundBundle<T: 'static, R: 'static> {
125    /// `events` field for events.
126    pub events: Node<OutboundEvent<T, R>>,
127    /// `status` field for status.
128    pub status: Node<OutboundStatus>,
129    /// `attempts` field for attempts.
130    pub attempts: Node<u32>,
131    /// `errors` field for errors.
132    pub errors: Node<String>,
133}
134
135#[derive(Clone, Default)]
136/// `OutboundAdapterOptions` data container.
137pub struct OutboundAdapterOptions {
138    /// `name` field for name.
139    pub name: Option<String>,
140    /// `retry` field for retry.
141    pub retry: RetryPolicy,
142}
143
144#[derive(Debug, Clone, PartialEq, Eq)]
145/// `WebSocketSessionCommand` variants.
146pub enum WebSocketSessionCommand {
147    /// `Start` variant.
148    Start,
149    /// `Send` variant.
150    Send(WebSocketSend),
151    /// `Close` variant.
152    Close {
153        /// `code` field for code.
154        code: Option<u16>,
155        /// `reason` field for reason.
156        reason: Option<String>,
157    },
158}
159
160#[derive(Debug, Clone, PartialEq, Eq)]
161/// `WebSocketSessionInbound` variants.
162pub enum WebSocketSessionInbound {
163    /// `Text` variant.
164    Text(String),
165    /// `Binary` variant.
166    Binary(Vec<u8>),
167}
168
169#[derive(Debug, Clone, PartialEq, Eq)]
170/// `WebSocketSessionLifecycle` variants.
171pub enum WebSocketSessionLifecycle {
172    /// `Starting` variant.
173    Starting {
174        /// `attempt` field for attempt.
175        attempt: u32,
176        /// `max_attempts` field for max attempts.
177        max_attempts: u32,
178    },
179    /// `Open` variant.
180    Open {
181        /// `attempt` field for attempt.
182        attempt: u32,
183    },
184    /// `Sent` variant.
185    Sent {
186        /// `message` field for message.
187        message: WebSocketSend,
188    },
189    /// `Closing` variant.
190    Closing {
191        /// `code` field for code.
192        code: Option<u16>,
193        /// `reason` field for reason.
194        reason: Option<String>,
195    },
196    /// `Closed` variant.
197    Closed {
198        /// `code` field for code.
199        code: Option<u16>,
200        /// `reason` field for reason.
201        reason: Option<String>,
202    },
203    /// `Retrying` variant.
204    Retrying {
205        /// `attempt` field for attempt.
206        attempt: u32,
207        /// `next_attempt` field for next attempt.
208        next_attempt: u32,
209        /// `delay_ms` field for delay ms.
210        delay_ms: u64,
211        /// `error` field for error.
212        error: String,
213    },
214    /// `Exhausted` variant.
215    Exhausted {
216        /// `attempt` field for attempt.
217        attempt: u32,
218        /// `error` field for error.
219        error: String,
220    },
221}
222
223#[derive(Debug, Clone, PartialEq, Eq)]
224/// `WebSocketSessionStateKind` variants.
225pub enum WebSocketSessionStateKind {
226    /// `Idle` variant.
227    Idle,
228    /// `Connecting` variant.
229    Connecting,
230    /// `Open` variant.
231    Open,
232    /// `Closing` variant.
233    Closing,
234    /// `Closed` variant.
235    Closed,
236    /// `Waiting` variant.
237    Waiting,
238    /// `Exhausted` variant.
239    Exhausted,
240    /// `Errored` variant.
241    Errored,
242}
243
244#[derive(Debug, Clone, PartialEq, Eq)]
245/// `WebSocketSessionStatus` data container.
246pub struct WebSocketSessionStatus {
247    /// `state` field for state.
248    pub state: WebSocketSessionStateKind,
249    /// `attempt` field for attempt.
250    pub attempt: u32,
251    /// `max_attempts` field for max attempts.
252    pub max_attempts: u32,
253    /// `sent` field for sent.
254    pub sent: u64,
255    /// `received` field for received.
256    pub received: u64,
257    /// `errors` field for errors.
258    pub errors: u64,
259    /// `last_delay_ms` field for last delay ms.
260    pub last_delay_ms: Option<u64>,
261}
262
263#[derive(Debug, Clone, PartialEq, Eq)]
264/// `WebSocketSessionOutbound` variants.
265pub enum WebSocketSessionOutbound {
266    /// `Queued` variant.
267    Queued {
268        /// `seq` field for seq.
269        seq: u64,
270        /// `message` field for message.
271        message: WebSocketSend,
272    },
273    /// `Sending` variant.
274    Sending {
275        /// `seq` field for seq.
276        seq: u64,
277        /// `message` field for message.
278        message: WebSocketSend,
279    },
280    /// `Sent` variant.
281    Sent {
282        /// `seq` field for seq.
283        seq: u64,
284        /// `message` field for message.
285        message: WebSocketSend,
286    },
287    /// `Rejected` variant.
288    Rejected {
289        /// `seq` field for seq.
290        seq: u64,
291        /// `message` field for message.
292        message: WebSocketSend,
293        /// `error` field for error.
294        error: String,
295    },
296    /// `Canceled` variant.
297    Canceled {
298        /// `seq` field for seq.
299        seq: u64,
300        /// `message` field for message.
301        message: WebSocketSend,
302        /// `reason` field for reason.
303        reason: String,
304    },
305}
306
307/// `WebSocketSessionBundle` data container.
308pub struct WebSocketSessionBundle {
309    /// `command` field for command.
310    pub command: Node<WebSocketSessionCommand>,
311    /// `inbound` field for inbound.
312    pub inbound: Node<WebSocketSessionInbound>,
313    /// `lifecycle` field for lifecycle.
314    pub lifecycle: Node<WebSocketSessionLifecycle>,
315    /// `outbound` field for outbound.
316    pub outbound: Node<WebSocketSessionOutbound>,
317    /// `status` field for status.
318    pub status: Node<WebSocketSessionStatus>,
319    /// `errors` field for errors.
320    pub errors: Node<String>,
321    /// `attempts` field for attempts.
322    pub attempts: Node<u32>,
323}
324
325impl WebSocketSessionBundle {
326    /// Updates or reads `start`.
327    pub fn start(&self) {
328        self.command.set(WebSocketSessionCommand::Start);
329    }
330
331    /// Updates or reads `send`.
332    pub fn send(&self, message: WebSocketSend) {
333        self.command.set(WebSocketSessionCommand::Send(message));
334    }
335
336    /// Updates or reads `send_text`.
337    pub fn send_text(&self, text: impl Into<String>) {
338        self.send(WebSocketSend::text(text));
339    }
340
341    /// Updates or reads `send_binary`.
342    pub fn send_binary(&self, bytes: impl Into<Vec<u8>>) {
343        self.send(WebSocketSend::binary(bytes));
344    }
345
346    /// Updates or reads `close`.
347    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)]
354/// `WebSocketSessionSendPolicy` variants.
355pub enum WebSocketSessionSendPolicy {
356    #[default]
357    /// `Reject` variant.
358    Reject,
359    /// `Buffer` variant.
360    Buffer {
361        /// `max_pending` field for max pending.
362        max_pending: usize,
363    },
364}
365
366#[derive(Clone, Default)]
367/// `WebSocketSessionOptions` data container.
368pub struct WebSocketSessionOptions {
369    /// `name` field for name.
370    pub name: Option<String>,
371    /// `retry` field for retry.
372    pub retry: RetryPolicy,
373    /// `send_policy` field for send policy.
374    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
450/// Creates or computes `websocket_session`.
451pub fn websocket_session(graph: &Graph, request: WebSocketRequest) -> WebSocketSessionBundle {
452    websocket_session_with_options(graph, request, WebSocketSessionOptions::default())
453}
454
455/// Creates or computes `websocket_session_with_options`.
456pub 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
1359/// Creates or computes `to_http`.
1360pub 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
1372/// Creates or computes `to_http_with_options`.
1373pub 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
1396/// Creates or computes `to_process`.
1397pub 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
1409/// Creates or computes `to_process_with_options`.
1410pub 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
1433/// Creates or computes `to_websocket`.
1434pub 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
1453/// Creates or computes `to_websocket_with_options`.
1454pub 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}