diff --git a/core/crates/solstone-core-spl/src/loopback_pipe.rs b/core/crates/solstone-core-spl/src/loopback_pipe.rs new file mode 100644 index 000000000..a73cb1bf3 --- /dev/null +++ b/core/crates/solstone-core-spl/src/loopback_pipe.rs @@ -0,0 +1,81 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +//! Full-duplex forwarding between a relay tunnel and the local SPL listener. + +use std::io; + +use tokio::io::{AsyncRead, AsyncWrite, AsyncWriteExt}; + +/// Replay `initial_prefix` to the loopback listener, then forward both streams. +/// +/// The prefix has already been consumed by the relay-side protocol classifier. +/// Replaying it before either forwarding direction starts preserves the original +/// local byte stream exactly. `copy_bidirectional` shuts down the opposite +/// writer after an EOF and returns the underlying I/O error unchanged. +pub async fn pipe_loopback( + mut tunnel: Tunnel, + mut loopback: Loopback, + initial_prefix: &[u8], +) -> io::Result<(u64, u64)> +where + Tunnel: AsyncRead + AsyncWrite + Unpin, + Loopback: AsyncRead + AsyncWrite + Unpin, +{ + loopback.write_all(initial_prefix).await?; + loopback.flush().await?; + tokio::io::copy_bidirectional(&mut tunnel, &mut loopback).await +} + +#[cfg(test)] +mod tests { + use std::io; + + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + use super::pipe_loopback; + + #[tokio::test] + async fn replays_initial_prefix_once_before_tunnel_payload() -> io::Result<()> { + let (mut tunnel_peer, tunnel) = tokio::io::duplex(64); + let (loopback, mut loopback_peer) = tokio::io::duplex(64); + + let pipe = tokio::spawn(async move { pipe_loopback(tunnel, loopback, b"TLS").await }); + + tunnel_peer.write_all(b"-payload").await?; + tunnel_peer.shutdown().await?; + + let mut received = [0_u8; 11]; + loopback_peer.read_exact(&mut received).await?; + assert_eq!(&received, b"TLS-payload"); + loopback_peer.shutdown().await?; + + let forwarded = pipe + .await + .map_err(|error| io::Error::other(error.to_string()))??; + assert_eq!(forwarded, (8, 0)); + Ok(()) + } + + #[tokio::test] + async fn forwards_loopback_responses_to_the_tunnel() -> io::Result<()> { + let (mut tunnel_peer, tunnel) = tokio::io::duplex(64); + let (loopback, mut loopback_peer) = tokio::io::duplex(64); + + let pipe = tokio::spawn(async move { pipe_loopback(tunnel, loopback, &[]).await }); + + loopback_peer.write_all(b"response").await?; + loopback_peer.shutdown().await?; + + let mut received = [0_u8; 8]; + tunnel_peer.read_exact(&mut received).await?; + assert_eq!(&received, b"response"); + tunnel_peer.shutdown().await?; + + let forwarded = pipe + .await + .map_err(|error| io::Error::other(error.to_string()))??; + assert_eq!(forwarded, (0, 8)); + Ok(()) + } +} diff --git a/core/crates/solstone-core-spl/src/reconnect_backoff.rs b/core/crates/solstone-core-spl/src/reconnect_backoff.rs new file mode 100644 index 000000000..c38220b4e --- /dev/null +++ b/core/crates/solstone-core-spl/src/reconnect_backoff.rs @@ -0,0 +1,157 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +use std::time::Duration; + +/// The initial unjittered delay after a failed relay connection loop. +pub const INITIAL_RECONNECT_BASE: Duration = Duration::from_secs(1); + +/// The largest unjittered reconnect base. +pub const MAX_RECONNECT_BASE: Duration = Duration::from_secs(60); + +/// A deterministic reconnect schedule for one failed connection loop. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ReconnectSchedule { + /// The capped base to use after this retry has been scheduled. + pub next_base: Duration, + /// The jittered, nonnegative delay to wait before the current retry. + pub delay: Duration, +} + +/// A caller supplied an invalid normalized jitter sample. +#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)] +pub enum ReconnectBackoffError { + /// The sample must be finite and in the inclusive 0.0 through 1.0 range. + #[error("normalized jitter sample must be finite and within 0.0..=1.0")] + InvalidJitterSample, +} + +/// Schedules a deterministic reconnect delay without reading a clock or RNG. +/// +/// Pass `Duration::ZERO` before the first failure; it schedules the contract's +/// one-second base and returns a two-second base for the following failure. +/// Subsequent callers pass the preceding [`ReconnectSchedule::next_base`]. +/// `normalized_jitter` maps 0.0 to -25%, 0.5 to no change, and 1.0 to +25%. +pub fn schedule_reconnect( + current_base: Duration, + normalized_jitter: f64, +) -> Result { + if !normalized_jitter.is_finite() || !(0.0..=1.0).contains(&normalized_jitter) { + return Err(ReconnectBackoffError::InvalidJitterSample); + } + + let base = current_base.clamp(INITIAL_RECONNECT_BASE, MAX_RECONNECT_BASE); + let multiplier = 0.75 + (normalized_jitter * 0.5); + let delay_nanos = ((base.as_nanos() as f64) * multiplier).round() as u64; + + Ok(ReconnectSchedule { + next_base: base.saturating_mul(2).min(MAX_RECONNECT_BASE), + delay: Duration::from_nanos(delay_nanos), + }) +} + +#[cfg(test)] +mod tests { + use std::time::Duration; + + use super::{ + INITIAL_RECONNECT_BASE, MAX_RECONNECT_BASE, ReconnectBackoffError, ReconnectSchedule, + schedule_reconnect, + }; + + #[test] + fn first_retry_uses_a_one_second_base_and_advances_to_two_seconds() { + assert_eq!( + schedule_reconnect(Duration::ZERO, 0.5), + Ok(ReconnectSchedule { + next_base: Duration::from_secs(2), + delay: INITIAL_RECONNECT_BASE, + }) + ); + } + + #[test] + fn retries_double_the_base_until_the_cap() { + let second = schedule_reconnect(Duration::from_secs(2), 0.5); + let fourth = schedule_reconnect(Duration::from_secs(4), 0.5); + + assert_eq!( + second, + Ok(ReconnectSchedule { + next_base: Duration::from_secs(4), + delay: Duration::from_secs(2), + }) + ); + assert_eq!( + fourth, + Ok(ReconnectSchedule { + next_base: Duration::from_secs(8), + delay: Duration::from_secs(4), + }) + ); + } + + #[test] + fn base_saturates_at_sixty_seconds() { + assert_eq!( + schedule_reconnect(MAX_RECONNECT_BASE, 0.5), + Ok(ReconnectSchedule { + next_base: MAX_RECONNECT_BASE, + delay: MAX_RECONNECT_BASE, + }) + ); + assert_eq!( + schedule_reconnect(Duration::from_secs(600), 0.5), + Ok(ReconnectSchedule { + next_base: MAX_RECONNECT_BASE, + delay: MAX_RECONNECT_BASE, + }) + ); + } + + #[test] + fn jitter_bounds_are_minus_and_plus_twenty_five_percent() { + assert_eq!( + schedule_reconnect(MAX_RECONNECT_BASE, 0.0), + Ok(ReconnectSchedule { + next_base: MAX_RECONNECT_BASE, + delay: Duration::from_secs(45), + }) + ); + assert_eq!( + schedule_reconnect(MAX_RECONNECT_BASE, 1.0), + Ok(ReconnectSchedule { + next_base: MAX_RECONNECT_BASE, + delay: Duration::from_secs(75), + }) + ); + } + + #[test] + fn rejects_out_of_range_and_non_finite_jitter_samples() { + assert_eq!( + schedule_reconnect(Duration::ZERO, -0.01), + Err(ReconnectBackoffError::InvalidJitterSample) + ); + assert_eq!( + schedule_reconnect(Duration::ZERO, 1.01), + Err(ReconnectBackoffError::InvalidJitterSample) + ); + assert_eq!( + schedule_reconnect(Duration::ZERO, f64::NAN), + Err(ReconnectBackoffError::InvalidJitterSample) + ); + assert_eq!( + schedule_reconnect(Duration::ZERO, f64::INFINITY), + Err(ReconnectBackoffError::InvalidJitterSample) + ); + } + + #[test] + fn the_same_inputs_always_produce_the_same_schedule() { + let first = schedule_reconnect(Duration::from_secs(16), 0.37); + let second = schedule_reconnect(Duration::from_secs(16), 0.37); + + assert_eq!(first, second); + } +} diff --git a/core/crates/solstone-core-spl/src/relay_control.rs b/core/crates/solstone-core-spl/src/relay_control.rs new file mode 100644 index 000000000..4a77992cc --- /dev/null +++ b/core/crates/solstone-core-spl/src/relay_control.rs @@ -0,0 +1,196 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +//! Pure helpers for relay WebSocket control traffic. + +use serde_json::Value; + +/// Converts a relay HTTP endpoint to its WebSocket equivalent. +pub fn websocket_endpoint(endpoint: &str) -> String { + if let Some(rest) = endpoint.strip_prefix("http://") { + return format!("ws://{rest}"); + } + if let Some(rest) = endpoint.strip_prefix("https://") { + return format!("wss://{rest}"); + } + endpoint.to_owned() +} + +/// Builds a relay tunnel or listen URL, with the required token query field. +/// +/// Call [`bearer_authorization_value`] separately when opening the WebSocket. +pub fn relay_tunnel_url(endpoint: &str, path: &str, instance_id: &str, token: &str) -> String { + let websocket_endpoint = websocket_endpoint(endpoint); + let endpoint_without_slash = websocket_endpoint.trim_end_matches('/'); + format!( + "{endpoint_without_slash}{path}?instance={}&token={}", + encode_query_component(instance_id), + encode_query_component(token), + ) +} + +/// Builds the authorization value required in addition to the token query field. +pub fn bearer_authorization_value(token: &str) -> String { + format!("Bearer {token}") +} + +/// A listen-channel control message that is either actionable or safely ignored. +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum ListenControl { + /// A validated tunnel identifier for a newly offered tunnel. + Incoming { tunnel_id: String }, + /// A well-formed control message that does not offer a tunnel. + Ignore, + /// Malformed input or an incomplete/invalid incoming-tunnel offer. This is nonfatal. + Invalid, +} + +/// Parses text or binary WebSocket input as a relay listen control message. +pub fn parse_listen_control(message: impl AsRef<[u8]>) -> ListenControl { + let text = match std::str::from_utf8(message.as_ref()) { + Ok(text) => text, + Err(_) => return ListenControl::Invalid, + }; + let parsed: Value = match serde_json::from_str(text) { + Ok(parsed) => parsed, + Err(_) => return ListenControl::Invalid, + }; + let object = match parsed.as_object() { + Some(object) => object, + None => return ListenControl::Invalid, + }; + match object.get("type").and_then(Value::as_str) { + Some("incoming") => {} + Some(_) => return ListenControl::Ignore, + None => return ListenControl::Invalid, + } + + let tunnel_id = match object.get("tunnel_id") { + Some(Value::String(value)) if !value.is_empty() => Some(value.clone()), + Some(Value::Number(value)) => integer_tunnel_id(value), + _ => None, + }; + match tunnel_id { + Some(tunnel_id) => ListenControl::Incoming { tunnel_id }, + None => ListenControl::Invalid, + } +} + +/// Converts only JSON integers whose decimal rendering exactly matches Python's `str()`. +/// +/// JSON floating-point values are deliberately rejected: their exponent and precision +/// formatting cannot be guaranteed to match Python's `str()` across the full range. +fn integer_tunnel_id(value: &serde_json::Number) -> Option { + if let Some(integer) = value.as_i64() { + return Some(integer.to_string()); + } + value.as_u64().map(|integer| integer.to_string()) +} + +fn encode_query_component(value: &str) -> String { + const HEX: &[u8; 16] = b"0123456789ABCDEF"; + + let mut encoded = String::with_capacity(value.len()); + for byte in value.bytes() { + match byte { + b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'.' | b'_' | b'~' => { + encoded.push(char::from(byte)); + } + b' ' => encoded.push('+'), + _ => { + encoded.push('%'); + encoded.push(char::from(HEX[usize::from(byte >> 4)])); + encoded.push(char::from(HEX[usize::from(byte & 0x0f)])); + } + } + } + encoded +} + +#[cfg(test)] +mod tests { + use super::{ + ListenControl, bearer_authorization_value, parse_listen_control, relay_tunnel_url, + websocket_endpoint, + }; + + #[test] + fn converts_http_endpoints_and_preserves_websocket_endpoints() { + assert_eq!(websocket_endpoint("http://relay.test"), "ws://relay.test"); + assert_eq!(websocket_endpoint("https://relay.test"), "wss://relay.test"); + assert_eq!(websocket_endpoint("ws://relay.test"), "ws://relay.test"); + assert_eq!(websocket_endpoint("wss://relay.test"), "wss://relay.test"); + } + + #[test] + fn builds_a_url_and_header_through_separate_apis() { + let url = relay_tunnel_url( + "https://relay.test/", + "/session/listen", + "home/a b", + "token +/?", + ); + let authorization = bearer_authorization_value("token +/?"); + + assert_eq!( + url, + "wss://relay.test/session/listen?instance=home%2Fa+b&token=token+%2B%2F%3F" + ); + assert_eq!(authorization, "Bearer token +/?"); + } + + #[test] + fn parses_text_and_binary_incoming_controls() { + assert_eq!( + parse_listen_control("{\"type\":\"incoming\",\"tunnel_id\":\"abc\"}"), + ListenControl::Incoming { + tunnel_id: "abc".to_owned() + } + ); + assert_eq!( + parse_listen_control(b"{\"type\":\"incoming\",\"tunnel_id\":42}"), + ListenControl::Incoming { + tunnel_id: "42".to_owned() + } + ); + } + + #[test] + fn invalidates_malformed_and_incomplete_incoming_controls() { + for message in [ + b"\xff".as_slice(), + b"not json".as_slice(), + b"[]".as_slice(), + b"{\"type\":\"incoming\"}".as_slice(), + b"{\"type\":\"incoming\",\"tunnel_id\":\"\"}".as_slice(), + b"{\"type\":\"incoming\",\"tunnel_id\":1.5}".as_slice(), + b"{\"type\":\"incoming\",\"tunnel_id\":false}".as_slice(), + ] { + assert_eq!(parse_listen_control(message), ListenControl::Invalid); + } + } + + #[test] + fn ignores_well_formed_non_incoming_controls() { + assert_eq!( + parse_listen_control("{\"type\":\"connected\"}"), + ListenControl::Ignore + ); + } + + #[test] + fn integer_tunnel_ids_use_python_decimal_rendering() { + assert_eq!( + parse_listen_control("{\"type\":\"incoming\",\"tunnel_id\":-42}"), + ListenControl::Incoming { + tunnel_id: "-42".to_owned() + } + ); + assert_eq!( + parse_listen_control("{\"type\":\"incoming\",\"tunnel_id\":18446744073709551615}"), + ListenControl::Incoming { + tunnel_id: "18446744073709551615".to_owned() + } + ); + } +} diff --git a/core/crates/solstone-core-spl/src/relay_health.rs b/core/crates/solstone-core-spl/src/relay_health.rs new file mode 100644 index 000000000..90ad9ee18 --- /dev/null +++ b/core/crates/solstone-core-spl/src/relay_health.rs @@ -0,0 +1,217 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +use serde_json::{Value, json}; + +use crate::{ + REASON_HOME_MISSING_MOBILE, REASON_LOCAL_PRIVATE_LISTENER_UNREACHABLE, + REASON_RELAY_ADMISSION_SATURATED, REASON_RELAY_TUNNEL_REJECTED, + REASON_RELAY_TUNNEL_UNREACHABLE, REASON_SERVICE_TOKEN_REJECTED, +}; + +/// The owner-visible connection state for the relay listener. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum RelayHealthState { + Connecting, + Connected, + Reconnecting, +} + +impl RelayHealthState { + /// Returns the stable owner-visible state spelling. + pub fn as_str(self) -> &'static str { + match self { + Self::Connecting => "connecting", + Self::Connected => "connected", + Self::Reconnecting => "reconnecting", + } + } +} + +/// A relay-tunnel outcome that has a stable owner-visible reason. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum RelayTunnelFailure { + HomeMissingMobile, + ServiceTokenRejected, + RelayTunnelRejected { status: u16 }, + RelayTunnelUnreachable, + LocalPrivateListenerUnreachable, + RelayAdmissionSaturated, +} + +impl RelayTunnelFailure { + /// Returns the U3 reason for this tunnel failure. + pub fn reason(self) -> &'static str { + match self { + Self::HomeMissingMobile => REASON_HOME_MISSING_MOBILE, + Self::ServiceTokenRejected => REASON_SERVICE_TOKEN_REJECTED, + Self::RelayTunnelRejected { .. } => REASON_RELAY_TUNNEL_REJECTED, + Self::RelayTunnelUnreachable => REASON_RELAY_TUNNEL_UNREACHABLE, + Self::LocalPrivateListenerUnreachable => REASON_LOCAL_PRIVATE_LISTENER_UNREACHABLE, + Self::RelayAdmissionSaturated => REASON_RELAY_ADMISSION_SATURATED, + } + } + + /// Returns the relay HTTP status when the relay rejected the tunnel. + pub fn status(self) -> Option { + match self { + Self::RelayTunnelRejected { status } => Some(status), + Self::HomeMissingMobile + | Self::ServiceTokenRejected + | Self::RelayTunnelUnreachable + | Self::LocalPrivateListenerUnreachable + | Self::RelayAdmissionSaturated => None, + } + } +} + +/// Pure in-memory state for a relay listener's owner-visible health payload. +/// +/// This type performs no I/O, runtime work, clock reads, or event emission. +/// Callers provide the observation timestamps and admission saturation count. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct RelayHealth { + state: RelayHealthState, + listen_generation: u64, + last_successful_relay_tunnel_at: Option, + last_relay_tunnel_error: Option<&'static str>, + last_relay_tunnel_error_at: Option, + relay_tunnel_error_status: Option, + relay_admission_saturated_count: u64, +} + +impl RelayHealth { + /// Creates an initial health record before any listener attempt has begun. + pub fn new() -> Self { + Self { + state: RelayHealthState::Connecting, + listen_generation: 0, + last_successful_relay_tunnel_at: None, + last_relay_tunnel_error: None, + last_relay_tunnel_error_at: None, + relay_tunnel_error_status: None, + relay_admission_saturated_count: 0, + } + } + + /// Starts a new relay listener attempt. + pub fn begin_listen_attempt(&mut self) { + self.listen_generation = self.listen_generation.saturating_add(1); + } + + /// Updates the owner-visible listener state. + pub fn set_state(&mut self, state: RelayHealthState) { + self.state = state; + } + + /// Records a successful relay tunnel and clears the prior tunnel error. + pub fn record_tunnel_success(&mut self, timestamp_ms: u64) { + self.last_successful_relay_tunnel_at = Some(timestamp_ms); + self.last_relay_tunnel_error = None; + self.last_relay_tunnel_error_at = None; + self.relay_tunnel_error_status = None; + } + + /// Records a failed relay tunnel without changing the last success. + pub fn record_tunnel_failure(&mut self, failure: RelayTunnelFailure, timestamp_ms: u64) { + self.last_relay_tunnel_error = Some(failure.reason()); + self.last_relay_tunnel_error_at = Some(timestamp_ms); + self.relay_tunnel_error_status = failure.status(); + } + + /// Replaces the cumulative relay-admission saturation count. + pub fn set_relay_admission_saturated_count(&mut self, count: u64) { + self.relay_admission_saturated_count = count; + } + + /// Returns the complete owner-visible relay health payload. + pub fn payload(&self) -> Value { + json!({ + "state": self.state.as_str(), + "listen_generation": self.listen_generation, + "last_successful_relay_tunnel_at": self.last_successful_relay_tunnel_at, + "last_relay_tunnel_error": self.last_relay_tunnel_error, + "last_relay_tunnel_error_at": self.last_relay_tunnel_error_at, + "relay_tunnel_error_status": self.relay_tunnel_error_status, + "relay_admission_saturated_count": self.relay_admission_saturated_count, + }) + } +} + +impl Default for RelayHealth { + fn default() -> Self { + Self::new() + } +} + +#[cfg(test)] +mod tests { + use serde_json::json; + + use super::{RelayHealth, RelayHealthState, RelayTunnelFailure}; + + #[test] + fn new_payload_has_exact_owner_visible_keys_and_null_error_fields() { + let health = RelayHealth::new(); + + assert_eq!( + health.payload(), + json!({ + "state": "connecting", + "listen_generation": 0, + "last_successful_relay_tunnel_at": null, + "last_relay_tunnel_error": null, + "last_relay_tunnel_error_at": null, + "relay_tunnel_error_status": null, + "relay_admission_saturated_count": 0, + }) + ); + } + + #[test] + fn rejected_tunnel_payload_includes_its_http_status() { + let mut health = RelayHealth::new(); + health.begin_listen_attempt(); + health.set_state(RelayHealthState::Reconnecting); + health.record_tunnel_failure( + RelayTunnelFailure::RelayTunnelRejected { status: 503 }, + 1_700_000_000_123, + ); + + assert_eq!( + health.payload(), + json!({ + "state": "reconnecting", + "listen_generation": 1, + "last_successful_relay_tunnel_at": null, + "last_relay_tunnel_error": "relay_tunnel_rejected", + "last_relay_tunnel_error_at": 1_700_000_000_123_u64, + "relay_tunnel_error_status": 503, + "relay_admission_saturated_count": 0, + }) + ); + } + + #[test] + fn successful_tunnel_after_failure_clears_only_error_fields() { + let mut health = RelayHealth::new(); + health.begin_listen_attempt(); + health.set_state(RelayHealthState::Connected); + health.set_relay_admission_saturated_count(7); + health.record_tunnel_failure(RelayTunnelFailure::ServiceTokenRejected, 1_000); + health.record_tunnel_success(1_001); + + assert_eq!( + health.payload(), + json!({ + "state": "connected", + "listen_generation": 1, + "last_successful_relay_tunnel_at": 1_001, + "last_relay_tunnel_error": null, + "last_relay_tunnel_error_at": null, + "relay_tunnel_error_status": null, + "relay_admission_saturated_count": 7, + }) + ); + } +} diff --git a/core/crates/solstone-core-spl/src/relay_status_failure.rs b/core/crates/solstone-core-spl/src/relay_status_failure.rs new file mode 100644 index 000000000..ac9fc0b2a --- /dev/null +++ b/core/crates/solstone-core-spl/src/relay_status_failure.rs @@ -0,0 +1,69 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +//! Pure classification of relay-tunnel connection failures. + +use crate::relay_health::RelayTunnelFailure; + +/// The externally observed category of a failed relay-tunnel connection. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum RelayTunnelFailureSignal { + /// The relay accepted the HTTP connection attempt but returned this status. + HttpStatus(u16), + /// The connection failed before the relay could return an HTTP status. + TransportFailure, +} + +/// Maps a relay connection failure signal to its owner-visible failure reason. +pub fn classify_relay_tunnel_failure(signal: RelayTunnelFailureSignal) -> RelayTunnelFailure { + match signal { + RelayTunnelFailureSignal::HttpStatus(404) => RelayTunnelFailure::HomeMissingMobile, + RelayTunnelFailureSignal::HttpStatus(401 | 403) => RelayTunnelFailure::ServiceTokenRejected, + RelayTunnelFailureSignal::HttpStatus(status) => { + RelayTunnelFailure::RelayTunnelRejected { status } + } + RelayTunnelFailureSignal::TransportFailure => RelayTunnelFailure::RelayTunnelUnreachable, + } +} + +#[cfg(test)] +mod tests { + use super::{RelayTunnelFailureSignal, classify_relay_tunnel_failure}; + use crate::relay_health::RelayTunnelFailure; + + #[test] + fn maps_not_found_to_missing_mobile_home() { + assert_eq!( + classify_relay_tunnel_failure(RelayTunnelFailureSignal::HttpStatus(404)), + RelayTunnelFailure::HomeMissingMobile, + ); + } + + #[test] + fn maps_authentication_and_authorization_rejections_to_token_rejected() { + for status in [401, 403] { + assert_eq!( + classify_relay_tunnel_failure(RelayTunnelFailureSignal::HttpStatus(status)), + RelayTunnelFailure::ServiceTokenRejected, + ); + } + } + + #[test] + fn retains_every_other_relay_rejection_status() { + for status in [400, 402, 500, 503] { + assert_eq!( + classify_relay_tunnel_failure(RelayTunnelFailureSignal::HttpStatus(status)), + RelayTunnelFailure::RelayTunnelRejected { status }, + ); + } + } + + #[test] + fn maps_connection_failures_without_http_status_to_unreachable() { + assert_eq!( + classify_relay_tunnel_failure(RelayTunnelFailureSignal::TransportFailure), + RelayTunnelFailure::RelayTunnelUnreachable, + ); + } +} diff --git a/core/crates/solstone-core-spl/src/tunnel_route.rs b/core/crates/solstone-core-spl/src/tunnel_route.rs new file mode 100644 index 000000000..7ee28ef1e --- /dev/null +++ b/core/crates/solstone-core-spl/src/tunnel_route.rs @@ -0,0 +1,73 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +//! Routing decision for the first four bytes of an inbound relay tunnel. + +/// The action selected from a bounded four-byte tunnel prefix peek. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum TunnelRoute { + /// Fewer than four bytes are available, so the caller must continue peeking. + NeedMorePrefix, + /// A TLS ClientHello prefix that belongs on the local loopback connection. + TlsLoopback, + /// An SPL blob-transfer offer beginning with `SBO1`. + BlobReceive, + /// A complete four-byte prefix that this service does not support. + Unsupported, +} + +/// Routes an inbound tunnel from its bounded four-byte prefix. +/// +/// The caller supplies the bytes already obtained from its peek operation. Only +/// the first four bytes participate in the decision, so a caller that has read +/// farther than its bounded peek still gets the same route. +pub fn route_tunnel_prefix(prefix: &[u8]) -> TunnelRoute { + if prefix.len() < 4 { + return TunnelRoute::NeedMorePrefix; + } + + if prefix.first() == Some(&0x16) { + return TunnelRoute::TlsLoopback; + } + + if prefix.starts_with(b"SBO1") { + return TunnelRoute::BlobReceive; + } + + TunnelRoute::Unsupported +} + +#[cfg(test)] +mod tests { + use super::{TunnelRoute, route_tunnel_prefix}; + + #[test] + fn routes_a_tls_client_hello_prefix_to_loopback() { + assert_eq!( + route_tunnel_prefix(&[0x16, 0x03, 0x01, 0x00]), + TunnelRoute::TlsLoopback + ); + } + + #[test] + fn routes_an_exact_sbo1_prefix_to_blob_receive() { + assert_eq!(route_tunnel_prefix(b"SBO1"), TunnelRoute::BlobReceive); + } + + #[test] + fn rejects_an_unsupported_complete_prefix() { + assert_eq!(route_tunnel_prefix(b"NOPE"), TunnelRoute::Unsupported); + } + + #[test] + fn waits_until_all_four_prefix_bytes_are_available() { + for prefix in [ + b"".as_slice(), + b"S".as_slice(), + b"SB".as_slice(), + b"SBO".as_slice(), + ] { + assert_eq!(route_tunnel_prefix(prefix), TunnelRoute::NeedMorePrefix); + } + } +}