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}