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 std::sync::Arc;
36
37use crate::CloseReason;
38use crate::messages::{ControllerMessage, ServiceMessage};
39
40/// Opaque send-failure marker for [`ServiceTransport`].
41///
42/// Layer 1 deliberately erases transport-specific failure details. Callers
43/// should treat any `TransportError` as a disconnect/reconnect signal.
44#[derive(Debug, thiserror::Error)]
45#[error("transport send failed")]
46pub struct TransportError;
47
48/// Policy the event loop should follow when the transport receive stream ends.
49///
50/// Each transport implementation returns the appropriate policy via
51/// [`ServiceTransport::close_policy`].
52#[derive(Debug, Clone, PartialEq, Eq)]
53pub enum TransportClosePolicy {
54    /// Reconnect to the controller, optionally including the cached close reason
55    /// from the last close frame.
56    Reconnect { reason: Option<CloseReason> },
57    /// Shut down the event loop entirely.
58    Shutdown,
59}
60
61/// Transport-agnostic interface for sending/receiving wire messages.
62///
63/// Method names are prefixed with `transport_` to avoid collision with
64/// existing `send()`/`recv()` methods on concrete types.
65#[async_trait::async_trait]
66pub trait ServiceTransport: Send {
67    /// Send a message reliably. Returns error on transport failure.
68    async fn transport_send(&mut self, msg: ServiceMessage) -> Result<(), TransportError>;
69
70    /// Send a message on a best-effort basis. Drops silently on failure.
71    async fn transport_send_best_effort(&mut self, msg: ServiceMessage);
72
73    /// Send a message, auto-paginating large report payloads.
74    ///
75    /// For WebSocket transports this splits large payloads into pages.
76    /// For in-process transports this delegates to [`transport_send`](Self::transport_send)
77    /// since channel-based transports have no frame size limits.
78    async fn transport_send_auto_paginate(
79        &mut self,
80        msg: ServiceMessage,
81    ) -> Result<(), TransportError>;
82
83    /// Receive the next controller message.
84    ///
85    /// # Contract
86    ///
87    /// `None` is the terminal condition for this abstraction:
88    ///
89    /// - clean close
90    /// - peer-closed stream
91    /// - receive-side transport failure
92    ///
93    /// Layer 1 callers must treat all three the same and stop reading.
94    ///
95    /// In production, standalone services bypass this lossy method and call
96    /// `ControllerConnection::recv()` directly for richer error handling.
97    /// The `None` contract here is exercised by embedded channel transports.
98    async fn transport_recv(&mut self) -> Option<ControllerMessage>;
99
100    /// Policy the event loop should follow when `transport_recv` returns `None`.
101    fn close_policy(&self) -> TransportClosePolicy {
102        TransportClosePolicy::Reconnect { reason: None }
103    }
104
105    /// Whether this transport is currently yielded to an external counterpart.
106    fn is_yielded(&self) -> bool {
107        false
108    }
109
110    /// Optional notifier fired when the yield state changes.
111    ///
112    /// Returns `Some` only for transports that support embedded yield signalling
113    /// (i.e. `EmbeddedTransport`). `ControllerConnection` uses the default `None`.
114    ///
115    /// Called once before the `run_embedded_service` loop to obtain a stable
116    /// notifier handle. The `Arc<Notify>` is valid for the transport lifetime.
117    fn yield_change_notifier(&self) -> Option<Arc<tokio::sync::Notify>> {
118        None
119    }
120}
121
122#[cfg(test)]
123mod tests {
124    use super::{ServiceTransport, TransportClosePolicy, TransportError};
125    use crate::CloseReason;
126    use crate::messages::{ControllerMessage, ServiceMessage};
127
128    struct DefaultTransport;
129
130    #[async_trait::async_trait]
131    impl ServiceTransport for DefaultTransport {
132        async fn transport_send(&mut self, _msg: ServiceMessage) -> Result<(), TransportError> {
133            Ok(())
134        }
135
136        async fn transport_send_best_effort(&mut self, _msg: ServiceMessage) {}
137
138        async fn transport_send_auto_paginate(
139            &mut self,
140            _msg: ServiceMessage,
141        ) -> Result<(), TransportError> {
142            Ok(())
143        }
144
145        async fn transport_recv(&mut self) -> Option<ControllerMessage> {
146            None
147        }
148    }
149
150    #[test]
151    fn default_close_policy_is_reconnect_with_no_reason() {
152        let transport = DefaultTransport;
153        assert_eq!(
154            transport.close_policy(),
155            TransportClosePolicy::Reconnect { reason: None }
156        );
157    }
158
159    #[test]
160    fn default_is_yielded_is_false() {
161        let transport = DefaultTransport;
162        assert!(!transport.is_yielded());
163    }
164
165    #[test]
166    fn close_policy_reconnect_none_eq_reconnect_none() {
167        let a = TransportClosePolicy::Reconnect { reason: None };
168        let b = TransportClosePolicy::Reconnect { reason: None };
169        assert_eq!(a, b);
170    }
171
172    #[test]
173    fn close_policy_shutdown_eq_shutdown() {
174        assert_eq!(
175            TransportClosePolicy::Shutdown,
176            TransportClosePolicy::Shutdown
177        );
178    }
179
180    #[test]
181    fn close_policy_reconnect_ne_shutdown() {
182        assert_ne!(
183            TransportClosePolicy::Reconnect { reason: None },
184            TransportClosePolicy::Shutdown
185        );
186    }
187
188    #[test]
189    fn close_policy_clone_preserves_equality() {
190        let original = TransportClosePolicy::Reconnect {
191            reason: Some(CloseReason::CertificateRotated),
192        };
193        let cloned = original.clone();
194        assert_eq!(original, cloned);
195    }
196
197    #[test]
198    fn close_policy_debug_contains_variant_name() {
199        let policy = TransportClosePolicy::Shutdown;
200        let debug = format!("{policy:?}");
201        assert!(debug.contains("Shutdown"), "Debug output was: {debug}");
202    }
203
204    #[test]
205    fn default_yield_change_notifier_returns_none() {
206        let t = DefaultTransport;
207        assert!(t.yield_change_notifier().is_none());
208    }
209}