Skip to main content

uptrakit_wire/
transport.rs

1//! Transport abstraction shared by service business logic.
2//!
3//! # Layered Error Contract
4//!
5//! The service runtime uses three error layers:
6//!
7//! - **Layer 1 (this module)**: [`ServiceTransport`] + [`TransportError`].
8//!   This layer is intentionally lossy: send failures collapse to an opaque
9//!   `TransportError`, and receive failures collapse to `None`.
10//! - **Layer 2 (`service-sdk`)**: `ControllerConnection::recv()`.
11//!   This layer keeps rich receive fidelity via
12//!   `Result<Option<ControllerMessage>, Report<EnrollmentError>>`.
13//! - **Layer 3 (`service-sdk`)**: event-loop/lifecycle classification.
14//!   This layer maps `EnrollmentError` to `LoopError` and decides reconnect
15//!   semantics (`Disconnected`, backoff, or fatal propagation).
16//!
17//! `agent-core` functions accept `&mut dyn ServiceTransport` so the same
18//! logic can run with either WebSocket or in-process channel transports.
19//!
20//! # `transport_recv()` Terminal Contract
21//!
22//! `transport_recv()` returns `Option<ControllerMessage>`:
23//!
24//! - `Some(message)` => keep processing.
25//! - `None` => transport is closed or broken; stop reading.
26//!
27//! Layer 1 callers must not distinguish clean close from transport error.
28//!
29//! # Unresolved Gaps
30//!
31//! `ServiceTransport` intentionally has no `transport_close()` method.
32//! Close semantics are transport-specific and managed by lifecycle owners,
33//! not by shared business logic.
34
35use crate::CloseReason;
36use crate::messages::{ControllerMessage, ServiceMessage};
37
38/// Opaque send-failure marker for [`ServiceTransport`].
39///
40/// Layer 1 deliberately erases transport-specific failure details. Callers
41/// should treat any `TransportError` as a disconnect/reconnect signal.
42#[derive(Debug, thiserror::Error)]
43#[error("transport send failed")]
44pub struct TransportError;
45
46/// Policy the event loop should follow when the transport receive stream ends.
47///
48/// Each transport implementation returns the appropriate policy via
49/// [`ServiceTransport::close_policy`].
50#[derive(Debug, Clone, PartialEq, Eq)]
51pub enum TransportClosePolicy {
52    /// Reconnect to the controller, optionally including the cached close reason
53    /// from the last close frame.
54    Reconnect { reason: Option<CloseReason> },
55    /// Shut down the event loop entirely.
56    Shutdown,
57}
58
59/// Transport-agnostic interface for sending/receiving wire messages.
60///
61/// Method names are prefixed with `transport_` to avoid collision with
62/// existing `send()`/`recv()` methods on concrete types.
63#[async_trait::async_trait]
64pub trait ServiceTransport: Send {
65    /// Send a message reliably. Returns error on transport failure.
66    async fn transport_send(&mut self, msg: ServiceMessage) -> Result<(), TransportError>;
67
68    /// Send a message on a best-effort basis. Drops silently on failure.
69    async fn transport_send_best_effort(&mut self, msg: ServiceMessage);
70
71    /// Send a message, auto-paginating large report payloads.
72    ///
73    /// For WebSocket transports this splits large payloads into pages.
74    /// For in-process transports this delegates to [`transport_send`](Self::transport_send)
75    /// since channel-based transports have no frame size limits.
76    async fn transport_send_auto_paginate(
77        &mut self,
78        msg: ServiceMessage,
79    ) -> Result<(), TransportError>;
80
81    /// Receive the next controller message.
82    ///
83    /// # Contract
84    ///
85    /// `None` is the terminal condition for this abstraction:
86    ///
87    /// - clean close
88    /// - peer-closed stream
89    /// - receive-side transport failure
90    ///
91    /// Layer 1 callers must treat all three the same and stop reading.
92    ///
93    /// In production, standalone services bypass this lossy method and call
94    /// `ControllerConnection::recv()` directly for richer error handling.
95    /// The `None` contract here is exercised by embedded channel transports.
96    async fn transport_recv(&mut self) -> Option<ControllerMessage>;
97
98    /// Policy the event loop should follow when `transport_recv` returns `None`.
99    fn close_policy(&self) -> TransportClosePolicy {
100        TransportClosePolicy::Reconnect { reason: None }
101    }
102
103    /// Whether this transport is currently yielded to an external counterpart.
104    fn is_yielded(&self) -> bool {
105        false
106    }
107}
108
109#[cfg(test)]
110mod tests {
111    use super::{ServiceTransport, TransportClosePolicy, TransportError};
112    use crate::CloseReason;
113    use crate::messages::{ControllerMessage, ServiceMessage};
114
115    struct DefaultTransport;
116
117    #[async_trait::async_trait]
118    impl ServiceTransport for DefaultTransport {
119        async fn transport_send(&mut self, _msg: ServiceMessage) -> Result<(), TransportError> {
120            Ok(())
121        }
122
123        async fn transport_send_best_effort(&mut self, _msg: ServiceMessage) {}
124
125        async fn transport_send_auto_paginate(
126            &mut self,
127            _msg: ServiceMessage,
128        ) -> Result<(), TransportError> {
129            Ok(())
130        }
131
132        async fn transport_recv(&mut self) -> Option<ControllerMessage> {
133            None
134        }
135    }
136
137    #[test]
138    fn default_close_policy_is_reconnect_with_no_reason() {
139        let transport = DefaultTransport;
140        assert_eq!(
141            transport.close_policy(),
142            TransportClosePolicy::Reconnect { reason: None }
143        );
144    }
145
146    #[test]
147    fn default_is_yielded_is_false() {
148        let transport = DefaultTransport;
149        assert!(!transport.is_yielded());
150    }
151
152    #[test]
153    fn close_policy_reconnect_none_eq_reconnect_none() {
154        let a = TransportClosePolicy::Reconnect { reason: None };
155        let b = TransportClosePolicy::Reconnect { reason: None };
156        assert_eq!(a, b);
157    }
158
159    #[test]
160    fn close_policy_shutdown_eq_shutdown() {
161        assert_eq!(
162            TransportClosePolicy::Shutdown,
163            TransportClosePolicy::Shutdown
164        );
165    }
166
167    #[test]
168    fn close_policy_reconnect_ne_shutdown() {
169        assert_ne!(
170            TransportClosePolicy::Reconnect { reason: None },
171            TransportClosePolicy::Shutdown
172        );
173    }
174
175    #[test]
176    fn close_policy_clone_preserves_equality() {
177        let original = TransportClosePolicy::Reconnect {
178            reason: Some(CloseReason::CertificateRotated),
179        };
180        let cloned = original.clone();
181        assert_eq!(original, cloned);
182    }
183
184    #[test]
185    fn close_policy_debug_contains_variant_name() {
186        let policy = TransportClosePolicy::Shutdown;
187        let debug = format!("{policy:?}");
188        assert!(debug.contains("Shutdown"), "Debug output was: {debug}");
189    }
190}