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}