Skip to main content

mkit_core/
protocol.rs

1//! Cross-transport types: error taxonomy, the [`Transport`] trait, the
2//! [`PackKey`] digest wrapper, and the retry/backoff helpers used by
3//! every transport implementation (memory, file, HTTP, S3, SSH, enc).
4//! [`retrying`] is the single shared driver of the SPEC-TRANSPORT §7
5//! ladder — every transport's `retrying`/`with_retry` method is a thin
6//! wrapper around it, so the backoff/classification policy lives in
7//! exactly one place instead of being reimplemented per crate.
8//!
9//! The SSH wire format is defined in `mkit-rpc`'s `ssh.proto` and
10//! lives in `mkit_rpc::mkit::rpc::v1::ssh`; transport-ssh consumes
11//! the schema directly. The hand-rolled `OP_HELLO` byte format that
12//! used to live in this module has been retired.
13
14// SPEC-TRANSPORT §7 calls out the exponential ladder in seconds
15// (1, 2, 4, …, 300). Expressing those values with `Duration::from_secs`
16// is deliberate — switching to `from_mins` loses the one-to-one match
17// with the spec text.
18#![allow(clippy::duration_suboptimal_units)]
19
20use core::fmt;
21use core::time::Duration;
22
23use crate::hash::{FromHexError, Hash, to_hex};
24use crate::refs::Ref;
25pub use crate::refs::RefWriteCondition;
26
27// ---------------------------------------------------------------------------
28// Error taxonomy
29// ---------------------------------------------------------------------------
30
31/// Errors that any transport may surface across the [`Transport`]
32/// boundary. Implementations MAY wrap transport-specific errors
33/// internally but MUST map them to one of these variants before
34/// returning.
35#[derive(Debug, thiserror::Error)]
36#[non_exhaustive]
37pub enum TransportError {
38    /// `download_pack` called on a digest the remote does not hold.
39    #[error("pack not found on remote")]
40    PackNotFound,
41    /// Authentication or ACL failure (HTTP 401/403, SSH auth refusal,
42    /// S3 `SignatureDoesNotMatch`, …).
43    #[error("access denied by remote")]
44    AccessDenied,
45    /// A remote requires admission before this write may proceed.
46    #[error("{0}")]
47    AdmissionRequired(Box<AdmissionRequired>),
48    /// User admission configuration or helper headers were rejected.
49    #[error("admission configuration: {0}")]
50    AdmissionConfiguration(String),
51    /// Local HTTPS trust configuration could not be loaded or validated.
52    #[error("HTTPS trust configuration: {0}")]
53    TlsConfiguration(String),
54    /// The configured admission helper failed.
55    #[error("admission helper failed: {0}")]
56    AdmissionHelperFailed(String),
57    /// Catch-all remote-side failure carrying an advisory message. The
58    /// message is for operators; programs MUST NOT pattern-match on its
59    /// contents.
60    #[error("remote error: {0}")]
61    RemoteError(String),
62    /// `update_ref` CAS precondition was not satisfied. Per
63    /// SPEC-TRANSPORT §7, callers MUST treat this as
64    /// "possibly-success on retry" for `.missing` / `.match` and
65    /// confirm with `read_ref`.
66    #[error("ref CAS precondition failed")]
67    RefConflict,
68    /// Caller passed a ref name failing SPEC-REFS §3.
69    #[error("invalid ref name: {0}")]
70    InvalidRef(String),
71    /// Network-level failure: DNS, TCP connect, TLS handshake, SSH
72    /// subprocess spawn. Retryable (see [`is_retryable`]).
73    #[error("connection to remote failed")]
74    ConnectionFailed,
75    /// Unexpected HTTP status or transport-protocol error. 5xx and 429
76    /// are retryable; 4xx (except 401/403/404/409/412) is not.
77    #[error("server error (status {status})")]
78    ServerError {
79        /// Numeric status code. HTTP uses its native codes; transports
80        /// without a status integer use `0`.
81        status: u16,
82    },
83    /// Server response did not match the wire contract (truncated
84    /// frame, unknown opcode, bad JSON, …).
85    #[error("invalid response from remote")]
86    InvalidResponse,
87    /// Generic protocol-level failure — malformed frame, unexpected
88    /// opcode order, or failed handshake.
89    #[error("protocol error")]
90    ProtocolError,
91    /// Payload exceeded a transport-specific cap.
92    #[error("payload too large: {0} bytes")]
93    PayloadTooLarge(usize),
94    /// An insecure URL scheme (plain `http://`) was supplied for a
95    /// non-loopback host. Plain HTTP is restricted to loopback addresses
96    /// (`127.0.0.1`, `::1`, `localhost`) so production traffic is never
97    /// transported in the clear.
98    #[error("insecure scheme: plain http:// is allowed only for loopback hosts")]
99    InsecureScheme,
100}
101
102/// One opaque challenge advertised by the remote.
103#[derive(Clone, PartialEq, Eq)]
104pub struct AdmissionChallengeEntry {
105    /// Lowercase scheme identifier.
106    pub scheme: String,
107    /// Opaque challenge value. Do not display or log it.
108    pub value: String,
109}
110
111impl fmt::Debug for AdmissionChallengeEntry {
112    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
113        f.debug_struct("AdmissionChallengeEntry")
114            .field("scheme", &self.scheme)
115            .field("value_len", &self.value.len())
116            .finish()
117    }
118}
119
120/// A remote admission challenge and the bounded response headers needed by a helper.
121#[derive(Clone, PartialEq, Eq)]
122#[non_exhaustive]
123pub struct AdmissionRequired {
124    /// Challenges in server preference order.
125    pub challenges: Vec<AdmissionChallengeEntry>,
126    /// Untrusted server description, retained for the helper.
127    pub description: String,
128    /// Bounded `WWW-Authenticate` header values.
129    pub www_authenticate: Vec<String>,
130    /// Bounded `PAYMENT-REQUIRED` header values.
131    pub payment_required: Vec<String>,
132    /// Terminal admission state, if a helper already ran or hit its cap.
133    pub reason: Option<&'static str>,
134}
135
136impl AdmissionRequired {
137    /// Construct a remote challenge from already-bounded fields. The Connect
138    /// client validates the bounds before calling this.
139    #[must_use]
140    pub fn new(
141        challenges: Vec<AdmissionChallengeEntry>,
142        description: String,
143        www_authenticate: Vec<String>,
144        payment_required: Vec<String>,
145    ) -> Self {
146        Self {
147            challenges,
148            description,
149            www_authenticate,
150            payment_required,
151            reason: None,
152        }
153    }
154
155    /// Add a safe terminal explanation without exposing challenge values.
156    #[must_use]
157    pub fn with_reason(mut self, reason: &'static str) -> Self {
158        self.reason = Some(reason);
159        self
160    }
161}
162
163impl fmt::Debug for AdmissionRequired {
164    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
165        f.debug_struct("AdmissionRequired")
166            .field("challenges", &self.challenges)
167            .field("description_len", &self.description.len())
168            .field("www_authenticate_count", &self.www_authenticate.len())
169            .field("payment_required_count", &self.payment_required.len())
170            .field("reason", &self.reason)
171            .finish_non_exhaustive()
172    }
173}
174
175impl fmt::Display for AdmissionRequired {
176    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
177        f.write_str("admission required by remote: ")?;
178        write_safe_remote_text(f, &self.description)?;
179        if self.challenges.is_empty() {
180            f.write_str(" (no challenges)")?;
181        } else {
182            f.write_str(" (schemes: ")?;
183            for (i, challenge) in self.challenges.iter().enumerate() {
184                if i != 0 {
185                    f.write_str(", ")?;
186                }
187                write_safe_remote_text(f, &challenge.scheme)?;
188            }
189            f.write_str(")")?;
190        }
191        if let Some(reason) = self.reason {
192            write!(f, ": {reason}")?;
193        }
194        Ok(())
195    }
196}
197
198fn write_safe_remote_text(f: &mut fmt::Formatter<'_>, value: &str) -> fmt::Result {
199    for c in value.chars() {
200        if c.is_control() || matches!(c, '\u{202a}'..='\u{202e}' | '\u{2066}'..='\u{2069}') {
201            write!(f, "\\u{{{:x}}}", c as u32)?;
202        } else {
203            write!(f, "{c}")?;
204        }
205    }
206    Ok(())
207}
208
209/// Result alias used throughout this module.
210pub type TransportResult<T> = Result<T, TransportError>;
211
212// ---------------------------------------------------------------------------
213// PackKey — 32-byte digest wrapper
214// ---------------------------------------------------------------------------
215
216/// A 32-byte pack digest used as the content-address for an uploaded
217/// pack. This is the same 32 bytes as [`Hash`](tyalias@Hash) but wrapped so pack
218/// digests and object hashes do not silently cross purposes at API
219/// boundaries.
220#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
221pub struct PackKey(pub [u8; 32]);
222
223impl PackKey {
224    /// Build a [`PackKey`] from a raw 32-byte digest.
225    #[must_use]
226    pub const fn new(bytes: [u8; 32]) -> Self {
227        Self(bytes)
228    }
229
230    /// Borrow the underlying 32 bytes.
231    #[must_use]
232    pub const fn as_bytes(&self) -> &[u8; 32] {
233        &self.0
234    }
235
236    /// Lowercase 64-char hex.
237    #[must_use]
238    pub fn to_hex(&self) -> String {
239        to_hex(&self.0)
240    }
241
242    /// Build a [`PackKey`] from a [`Hash`](tyalias@Hash) (alias for [`From`]).
243    #[must_use]
244    pub const fn from_hash(h: Hash) -> Self {
245        Self(h)
246    }
247
248    /// Convert back to a plain [`Hash`](tyalias@Hash).
249    #[must_use]
250    pub const fn into_hash(self) -> Hash {
251        self.0
252    }
253
254    /// Verify downloaded bytes against this independently requested address.
255    /// Covers the entire byte sequence, including a pack's internal trailer.
256    /// Also applies to auxiliary blobs: a lookup URL or a self-consistent
257    /// manifest cannot authenticate the returned bytes on its own.
258    ///
259    /// # Errors
260    /// Returns [`TransportError::InvalidResponse`] on a digest mismatch.
261    pub fn verify_bytes(self, bytes: &[u8]) -> TransportResult<()> {
262        if crate::hash::hash(bytes) != self.0 {
263            return Err(TransportError::InvalidResponse);
264        }
265        Ok(())
266    }
267}
268
269impl From<Hash> for PackKey {
270    fn from(h: Hash) -> Self {
271        Self(h)
272    }
273}
274
275impl From<PackKey> for Hash {
276    fn from(k: PackKey) -> Hash {
277        k.0
278    }
279}
280
281impl fmt::Display for PackKey {
282    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
283        f.write_str(&self.to_hex())
284    }
285}
286
287/// Parse a [`PackKey`] from a 64-char lowercase hex string.
288///
289/// Accepts uppercase too (matches the permissive [`crate::hash::from_hex`]
290/// semantics); callers that require lowercase MUST validate the input
291/// independently.
292pub fn pack_key_from_hex(s: &str) -> Result<PackKey, FromHexError> {
293    let h = crate::hash::from_hex(s)?;
294    Ok(PackKey(h))
295}
296
297// ---------------------------------------------------------------------------
298// Retry / backoff
299// ---------------------------------------------------------------------------
300
301/// Return `true` if a transport should retry after seeing `err`.
302///
303/// Retryable per SPEC-TRANSPORT §7:
304/// - [`TransportError::ConnectionFailed`]
305/// - [`TransportError::ServerError`] with a 5xx status OR HTTP 429.
306///
307/// Explicitly non-retryable:
308/// - [`TransportError::PackNotFound`]
309/// - [`TransportError::AccessDenied`]
310/// - [`TransportError::AdmissionRequired`] — requires an explicit helper response
311/// - [`TransportError::RefConflict`] (CAS retry is a caller-level policy)
312/// - [`TransportError::InvalidRef`]
313/// - [`TransportError::InvalidResponse`] / [`TransportError::ProtocolError`]
314/// - [`TransportError::PayloadTooLarge`]
315/// - [`TransportError::TlsConfiguration`] — local trust setup failed.
316/// - [`TransportError::RemoteError`] — the remote chose not to be specific;
317///   we do not guess.
318/// - [`TransportError::ServerError`] with any 4xx status.
319#[must_use]
320pub fn is_retryable(err: &TransportError) -> bool {
321    match err {
322        TransportError::ConnectionFailed => true,
323        TransportError::ServerError { status } => *status >= 500 || *status == 429,
324        _ => false,
325    }
326}
327
328/// Max attempts for the default backoff ladder.
329///
330/// SPEC-TRANSPORT §7: `attempt = 1; while attempt ≤ 5`.
331pub const BACKOFF_MAX_ATTEMPTS: u32 = 5;
332
333/// Initial sleep between attempts.
334pub const BACKOFF_INITIAL: Duration = Duration::from_secs(1);
335
336/// Upper bound on any individual sleep.
337pub const BACKOFF_CAP: Duration = Duration::from_secs(300);
338
339/// Per-pack body size ceiling enforced by every transport that ingests
340/// pack bytes (HTTP `Content-Length`, S3 `GetObject`, SSH
341/// `DownloadPackHeader.total_bytes`). On 64-bit targets, 4 GiB matches
342/// the pack-format addressable range; pointer-width-limited targets cap
343/// at their maximum addressable buffer size instead of failing to compile.
344#[cfg(target_pointer_width = "64")]
345pub const PACK_BODY_LIMIT: u64 = 4 * 1024 * 1024 * 1024;
346#[cfg(not(target_pointer_width = "64"))]
347pub const PACK_BODY_LIMIT: u64 = usize::MAX as u64;
348
349/// `usize`-typed mirror of [`PACK_BODY_LIMIT`] for `Vec`-shaped buffer
350/// caps. The assertion below prevents silent truncation on any target.
351#[allow(clippy::cast_possible_truncation)]
352pub const PACK_BODY_LIMIT_USIZE: usize = PACK_BODY_LIMIT as usize;
353const _: () = assert!(
354    (PACK_BODY_LIMIT_USIZE as u64) == PACK_BODY_LIMIT,
355    "PACK_BODY_LIMIT does not fit in usize on this target",
356);
357
358/// Exponential-backoff iterator used by all transports.
359///
360/// Yields `[1s, 2s, 4s, 8s, 16s]` (5 attempts) for the default ladder,
361/// doubling each step and capping at 300s. This is the ladder mandated
362/// by SPEC-TRANSPORT §7 for `ConnectionFailed`, 5xx, and HTTP 429.
363///
364/// The iterator is self-contained — it holds no reference to a clock,
365/// so it can be constructed in tests and exhaustively enumerated.
366#[derive(Debug, Clone)]
367pub struct BackoffIterator {
368    next_delay: Duration,
369    attempts_remaining: u32,
370    cap: Duration,
371}
372
373impl BackoffIterator {
374    /// Default ladder: 5 attempts, starting at 1s, doubling, capped at 300s.
375    #[must_use]
376    pub const fn new() -> Self {
377        Self {
378            next_delay: BACKOFF_INITIAL,
379            attempts_remaining: BACKOFF_MAX_ATTEMPTS,
380            cap: BACKOFF_CAP,
381        }
382    }
383
384    /// Custom ladder for tests.
385    #[must_use]
386    pub const fn with(initial: Duration, cap: Duration, attempts: u32) -> Self {
387        Self {
388            next_delay: initial,
389            attempts_remaining: attempts,
390            cap,
391        }
392    }
393}
394
395impl Default for BackoffIterator {
396    fn default() -> Self {
397        Self::new()
398    }
399}
400
401impl Iterator for BackoffIterator {
402    type Item = Duration;
403
404    fn next(&mut self) -> Option<Self::Item> {
405        if self.attempts_remaining == 0 {
406            return None;
407        }
408        self.attempts_remaining -= 1;
409        let current = self.next_delay;
410        let doubled = current.saturating_mul(2);
411        self.next_delay = if doubled > self.cap {
412            self.cap
413        } else {
414            doubled
415        };
416        Some(current)
417    }
418}
419
420/// Transport-agnostic retry driver shared by every [`Transport`]
421/// implementation, so the SPEC-TRANSPORT §7 ladder lives in exactly one
422/// place instead of being reimplemented per crate. Extracted from what
423/// was `HttpTransport::retrying` in `mkit-transport-http`.
424///
425/// `op` is re-invoked from scratch on every attempt — it MUST perform a
426/// fresh, self-contained unit of work each call (a new HTTP request on
427/// the existing connection pool, a freshly-reconnected SSH child, a
428/// redialed encrypted session, …) rather than assuming any state left
429/// over from a failed prior attempt is still valid. This matters most
430/// for connection-oriented transports: a frame-level failure can leave
431/// a stream mid-message-desynced, so `op` reconnecting before retrying
432/// (rather than resuming on the same broken handle) is what makes the
433/// retry safe, not just present.
434///
435/// `backoff` is a ladder *factory* (not a live iterator) so a fresh
436/// ladder starts on every call to `retrying` — production uses
437/// [`BackoffIterator::new`], tests inject a short/deterministic ladder.
438/// `sleep` is the delay hook between attempts; production sleeps for
439/// the full duration, tests typically inject a no-op or recorder.
440///
441/// Retries only the classes [`is_retryable`] accepts
442/// (`ConnectionFailed`, `ServerError{5xx}`, `ServerError{429}`); every
443/// other error returns immediately on the first attempt.
444pub fn retrying<T>(
445    mut op: impl FnMut() -> TransportResult<T>,
446    backoff: fn() -> BackoffIterator,
447    sleep: fn(Duration),
448) -> TransportResult<T> {
449    let mut ladder = backoff();
450    loop {
451        match op() {
452            Ok(v) => return Ok(v),
453            Err(err) => {
454                if is_retryable(&err)
455                    && let Some(delay) = ladder.next()
456                {
457                    sleep(delay);
458                    continue;
459                }
460                return Err(err);
461            }
462        }
463    }
464}
465
466// ---------------------------------------------------------------------------
467// PackChunk — transport-agnostic streaming segment
468// ---------------------------------------------------------------------------
469
470/// One bounded-size segment of a streamed pack transfer.
471///
472/// Mirrors the wire-level `PackChunk` shape shared by the SSH and enc
473/// transports (`offset`, `data`, `last` — see
474/// `mkit-rpc/proto/mkit/rpc/v1/ssh/ssh.proto`) and the `mkit.transport.v1`
475/// Connect proto being designed for the HTTP
476/// reference worker, without this crate depending on any
477/// protobuf-generated type: `mkit-core` is the dependency root that
478/// `mkit-rpc` builds on, not the reverse (see this module's header
479/// comment), so the canonical protobuf `PackChunk` cannot be named here.
480/// Transports convert 1:1 between this type and their own wire
481/// representation.
482#[derive(Debug, Clone, PartialEq, Eq)]
483pub struct PackChunk {
484    /// Byte offset of `data` within the pack. Consecutive chunks in one
485    /// transfer MUST have ascending, contiguous offsets starting at 0 —
486    /// i.e. chunk *n*'s `offset` equals the sum of every prior chunk's
487    /// `data.len()`.
488    pub offset: u64,
489    /// Chunk payload. Transports typically bound this to a fixed
490    /// per-frame maximum (e.g. `mkit_rpc::CHUNK_DATA_MAX`, 800 KiB for
491    /// SSH/enc) so no single chunk forces a large allocation.
492    pub data: Vec<u8>,
493    /// `true` on the final chunk of the stream. An empty pack is still
494    /// represented as exactly one chunk with `last = true` and empty
495    /// `data` — a stream MUST NOT end silently without a `last = true`
496    /// chunk.
497    pub last: bool,
498}
499
500// ---------------------------------------------------------------------------
501// Transport trait
502// ---------------------------------------------------------------------------
503
504/// The mkit transport vtable.
505///
506/// Every transport (memory, file, HTTP, S3, SSH) implements this trait.
507/// Methods are synchronous and take `&self`; transports that need
508/// interior mutability (e.g. connection pools) MUST use a `Mutex` /
509/// `RwLock` internally. This keeps the trait object-safe.
510///
511/// All implementations MUST honour the retry policy in
512/// SPEC-TRANSPORT §7 internally OR document that the caller is
513/// responsible — the abstract trait takes no position. The
514/// [`is_retryable`] and [`BackoffIterator`] helpers are provided for
515/// implementations that embed the policy.
516/// The repository a transport addresses, for remote error context.
517#[derive(Debug, Clone, Copy, PartialEq, Eq)]
518#[non_exhaustive]
519pub struct RepositoryAddress<'a> {
520    /// Repository identity as sent on the wire (STC §7.4).
521    pub repository: &'a str,
522    /// Origin (scheme and authority) of the remote.
523    pub origin: &'a str,
524}
525
526impl<'a> RepositoryAddress<'a> {
527    /// Build an address from a repository identity and a remote origin.
528    #[must_use]
529    pub const fn new(repository: &'a str, origin: &'a str) -> Self {
530        Self { repository, origin }
531    }
532}
533
534pub trait Transport: Send + Sync {
535    /// Repository identity and origin used for remote error context, if available.
536    /// Existing transports carry no addressing context by default.
537    fn repository_address(&self) -> Option<RepositoryAddress<'_>> {
538        None
539    }
540
541    /// Upload a pack. The digest is computed by the caller (BLAKE3 of
542    /// the full pack bytes) and used as the object key — servers MAY
543    /// dedupe on this key.
544    fn upload_pack(&self, bytes: &[u8], key: &PackKey) -> TransportResult<()>;
545
546    /// Upload a pack for an intended head ref. Transports without upload
547    /// tickets ignore the hint and retain their existing upload behavior.
548    fn upload_pack_via_ref(
549        &self,
550        bytes: &[u8],
551        key: &PackKey,
552        _head_ref: &str,
553    ) -> TransportResult<()> {
554        self.upload_pack(bytes, key)
555    }
556
557    /// Download a pack by its digest.
558    ///
559    /// Returns [`TransportError::PackNotFound`] if the remote does not
560    /// hold this digest.
561    fn download_pack(&self, key: &PackKey) -> TransportResult<Vec<u8>>;
562
563    /// Download a pack with a read-your-writes ref hint. Defaults to ignoring it.
564    fn download_pack_via_ref(&self, key: &PackKey, _ref_name: &str) -> TransportResult<Vec<u8>> {
565        self.download_pack(key)
566    }
567
568    /// Upload a pack by streaming bounded-size [`PackChunk`]s instead of
569    /// requiring the whole pack materialized as one `&[u8]` up front.
570    ///
571    /// `total_bytes` is the caller-declared pack length — the wire
572    /// header most streaming transports send before the first chunk.
573    /// `chunks` MUST yield its segments in ascending contiguous `offset`
574    /// order and end with exactly one item whose `last` field is `true`
575    /// (an empty pack still yields one `last = true` chunk with empty
576    /// `data`); the accumulated `data` length across every yielded chunk
577    /// MUST equal `total_bytes`. The digest (`key`) is still computed by
578    /// the caller up front, exactly as for [`Self::upload_pack`] — this
579    /// method does not hash the stream itself.
580    ///
581    /// This is an additive, opt-in entry point: no existing transport is
582    /// forced to implement real streaming. The default impl buffers
583    /// `chunks` into one `Vec` (bounded by [`PACK_BODY_LIMIT_USIZE`]) and
584    /// delegates to [`Self::upload_pack`], so every transport gets a
585    /// working implementation with zero code — callers may always use
586    /// this entry point, even against a transport that has not opted
587    /// into streaming. Transports that can forward chunks directly to
588    /// their own wire (SSH, enc — see `mkit-transport-ssh`'s existing
589    /// `PackChunk` frame loop) SHOULD override this to avoid the buffer
590    /// and stay in bounded memory regardless of pack size.
591    ///
592    /// # Errors
593    ///
594    /// Returns [`TransportError::ProtocolError`] if `chunks` never
595    /// yields a `last = true` item, or if the accumulated byte count
596    /// does not equal `total_bytes`. Returns
597    /// [`TransportError::PayloadTooLarge`] if `total_bytes` (or the
598    /// accumulated count) would exceed [`PACK_BODY_LIMIT`]. Propagates
599    /// any error yielded by `chunks` itself (e.g. the caller's own I/O
600    /// error while reading a pack off disk).
601    fn upload_pack_streaming(
602        &self,
603        key: &PackKey,
604        total_bytes: u64,
605        chunks: &mut dyn Iterator<Item = TransportResult<PackChunk>>,
606    ) -> TransportResult<()> {
607        if total_bytes > PACK_BODY_LIMIT {
608            return Err(TransportError::PayloadTooLarge(PACK_BODY_LIMIT_USIZE));
609        }
610        // `total_bytes <= PACK_BODY_LIMIT` was just checked, and
611        // `PACK_BODY_LIMIT_USIZE as u64 == PACK_BODY_LIMIT` is asserted
612        // at the constant's definition, so this conversion never
613        // truncates — `try_from` (rather than `as`) makes that provable
614        // to clippy instead of asserted in a comment.
615        let initial = usize::try_from(total_bytes).unwrap_or(PACK_BODY_LIMIT_USIZE);
616        let mut buf = Vec::with_capacity(initial);
617        let mut saw_last = false;
618        for chunk in chunks {
619            let c = chunk?;
620            if c.data.len() > PACK_BODY_LIMIT_USIZE.saturating_sub(buf.len()) {
621                return Err(TransportError::PayloadTooLarge(PACK_BODY_LIMIT_USIZE));
622            }
623            buf.extend_from_slice(&c.data);
624            if c.last {
625                saw_last = true;
626                break;
627            }
628        }
629        if !saw_last || buf.len() as u64 != total_bytes {
630            return Err(TransportError::ProtocolError);
631        }
632        self.upload_pack(&buf, key)
633    }
634
635    /// Download a pack as a lazy stream of bounded-size [`PackChunk`]s
636    /// instead of one big `Vec<u8>`.
637    ///
638    /// This is an additive, opt-in entry point mirroring
639    /// [`Self::upload_pack_streaming`]. The default impl calls
640    /// [`Self::download_pack`] eagerly (so it does not save memory by
641    /// itself) and wraps the whole result as a single `last = true`
642    /// chunk — every transport gets a working implementation with zero
643    /// code. Transports that can read their own wire incrementally (SSH,
644    /// enc) SHOULD override this to yield each wire chunk as it arrives,
645    /// keeping memory bounded to roughly one chunk at a time regardless
646    /// of total pack size.
647    ///
648    /// # Errors
649    ///
650    /// Returns [`TransportError::PackNotFound`] immediately if the
651    /// remote does not hold `key` — a conforming implementation never
652    /// returns a stream that then fails its first item with
653    /// `PackNotFound`. Errors surfacing mid-stream (a malformed frame, a
654    /// connection drop) are yielded as `Err` items from the returned
655    /// iterator rather than failing this call itself, since an
656    /// overridden implementation may not know the transfer will fail
657    /// until partway through.
658    fn download_pack_streaming(
659        &self,
660        key: &PackKey,
661    ) -> TransportResult<Box<dyn Iterator<Item = TransportResult<PackChunk>> + '_>> {
662        let bytes = self.download_pack(key)?;
663        Ok(Box::new(core::iter::once(Ok(PackChunk {
664            offset: 0,
665            data: bytes,
666            last: true,
667        }))))
668    }
669
670    /// HEAD-check a pack. Cheaper than [`Self::download_pack`] on
671    /// network transports.
672    fn pack_exists(&self, key: &PackKey) -> TransportResult<bool>;
673
674    /// Check a pack with a read-your-writes ref hint. Defaults to ignoring it.
675    fn pack_exists_via_ref(&self, key: &PackKey, _ref_name: &str) -> TransportResult<bool> {
676        self.pack_exists(key)
677    }
678
679    /// Upload a content-addressed **auxiliary blob** — transfer metadata
680    /// that is NOT a packfile (e.g. a packlist chain node, SPEC-PACKFILE is
681    /// silent on these). The key is BLAKE3 of `bytes`, exactly like a pack.
682    ///
683    /// Auxiliary blobs share the digest-keyed content-addressed store with
684    /// packs (the store is a general blob store; "pack" is just the primary
685    /// content kind), so the default impl delegates to [`Self::upload_pack`].
686    /// The distinct verb keeps the *kind* explicit at the call site so a
687    /// caller never has to infer "is this blob a packfile or metadata?".
688    fn upload_blob(&self, bytes: &[u8], key: &PackKey) -> TransportResult<()> {
689        self.upload_pack(bytes, key)
690    }
691
692    /// Upload auxiliary content for an intended head ref. The default keeps
693    /// the transport's existing blob behavior.
694    fn upload_blob_via_ref(
695        &self,
696        bytes: &[u8],
697        key: &PackKey,
698        _head_ref: &str,
699    ) -> TransportResult<()> {
700        self.upload_blob(bytes, key)
701    }
702
703    /// Download an auxiliary blob by digest. Counterpart to
704    /// [`Self::upload_blob`]; default impl delegates to
705    /// [`Self::download_pack`]. Returns [`TransportError::PackNotFound`] if
706    /// the remote does not hold this digest.
707    fn download_blob(&self, key: &PackKey) -> TransportResult<Vec<u8>> {
708        self.download_pack(key)
709    }
710
711    /// Download auxiliary metadata with a ref hint. Defaults to ignoring it.
712    fn download_blob_via_ref(&self, key: &PackKey, _ref_name: &str) -> TransportResult<Vec<u8>> {
713        self.download_blob(key)
714    }
715
716    /// Unconditional ref write — equivalent to
717    /// `update_ref(name, RefWriteCondition::Any, hash)`.
718    ///
719    /// Default impl delegates to [`Self::update_ref`] so transports only
720    /// implement one entry point.
721    fn write_ref(&self, name: &str, hash: &Hash) -> TransportResult<()> {
722        self.update_ref(name, RefWriteCondition::Any, hash)
723    }
724
725    /// CAS ref write. See [`RefWriteCondition`].
726    ///
727    /// On `.missing` / `.match` CAS failure, returns
728    /// [`TransportError::RefConflict`]. Callers retrying after a
729    /// timeout MUST follow up with [`Self::read_ref`] to confirm
730    /// whether the first attempt actually landed (SPEC-TRANSPORT §7).
731    fn update_ref(
732        &self,
733        name: &str,
734        condition: RefWriteCondition,
735        hash: &Hash,
736    ) -> TransportResult<()>;
737
738    /// Read the current value of a ref, or `None` if it does not exist.
739    fn read_ref(&self, name: &str) -> TransportResult<Option<Hash>>;
740
741    /// List refs whose full name starts with `prefix`. Returned names
742    /// have `prefix` stripped per SPEC-REFS §4. An empty prefix lists
743    /// every ref.
744    fn list_refs(&self, prefix: &str) -> TransportResult<Vec<Ref>>;
745
746    /// Advance a branch by updating its **head ref and its packmap ref
747    /// together**, each under its own CAS precondition.
748    ///
749    /// This exists so the delta-transfer invariant — "if `head_ref` resolves
750    /// to T, the packmap reconstructs `closure(T)`" — can be upheld without a
751    /// window where the two refs disagree. A transport backed by a
752    /// transactional ref store SHOULD override this to apply both writes in
753    /// ONE transaction: then a failed advance changes nothing, and the head
754    /// is never observed past a packmap that can't yet reconstruct it.
755    ///
756    /// The default impl is the safe non-transactional approximation used by
757    /// stores without multi-ref transactions: it writes the **packmap first**
758    /// (durable before the head moves) then the head. A crash in between
759    /// leaves the head at its prior value and the packmap a superset — still
760    /// consistent for fetch. The [`AdvanceOutcome`] distinguishes a packmap
761    /// precondition failure (caller re-reads and retries the chain) from a
762    /// head precondition failure (caller treats it as non-fast-forward),
763    /// which a single `RefConflict` could not.
764    fn advance_refs(
765        &self,
766        head_ref: &str,
767        head_condition: RefWriteCondition,
768        head_value: &Hash,
769        packmap_ref: &str,
770        packmap_condition: RefWriteCondition,
771        packmap_value: &Hash,
772    ) -> TransportResult<AdvanceOutcome> {
773        match self.update_ref(packmap_ref, packmap_condition, packmap_value) {
774            Ok(()) => {}
775            Err(TransportError::RefConflict) => return Ok(AdvanceOutcome::PackmapConflict),
776            Err(e) => return Err(e),
777        }
778        match self.update_ref(head_ref, head_condition, head_value) {
779            Ok(()) => Ok(AdvanceOutcome::Committed),
780            Err(TransportError::RefConflict) => Ok(AdvanceOutcome::HeadConflict),
781            Err(e) => Err(e),
782        }
783    }
784
785    /// Advance both refs while committing any upload tickets for `commit`.
786    /// Other transports ignore the pack keys and use their existing advance.
787    #[allow(clippy::too_many_arguments)]
788    fn advance_refs_committing(
789        &self,
790        head_ref: &str,
791        head_condition: RefWriteCondition,
792        head_value: &Hash,
793        packmap_ref: &str,
794        packmap_condition: RefWriteCondition,
795        packmap_value: &Hash,
796        _commit: &[PackKey],
797    ) -> TransportResult<CommitOutcome> {
798        self.advance_refs(
799            head_ref,
800            head_condition,
801            head_value,
802            packmap_ref,
803            packmap_condition,
804            packmap_value,
805        )
806        .map(CommitOutcome::Advanced)
807    }
808
809    /// Limits the remote advertises to the push planner. `None` means the
810    /// transport does not impose an additional limit.
811    fn upload_limits(&self) -> UploadLimits {
812        UploadLimits::default()
813    }
814
815    /// Whether [`Self::advance_refs`] commits the head + packmap advance as
816    /// one indivisible transaction, rather than the default's ordered
817    /// packmap-then-head writes.
818    ///
819    /// The default (non-transactional) `advance_refs` is safe for an
820    /// **appending** packmap write: per its doc comment, a crash or lost
821    /// head-CAS race between the two writes leaves the packmap a strict
822    /// superset of what the (unmoved) head needs — still reconstructable.
823    /// That safety argument does NOT extend to a packmap **reset** (a fresh
824    /// node with `prev = None`, produced by the pack-chain re-baseline,
825    /// mkit #406): a reset is not a superset of the prior chain, so a
826    /// packmap write that commits while the paired head write loses its CAS
827    /// would strand the (still-unmoved) head pointing at a commit whose
828    /// closure the reset packmap can no longer reconstruct (mkit #521).
829    ///
830    /// Callers MUST treat `false` (the default) as "never request a
831    /// packmap reset against this transport" — see
832    /// `remote_dispatch::push_branch`'s re-baseline gate. Override to
833    /// `true` ONLY when [`Self::advance_refs`] is overridden with a
834    /// genuinely transactional implementation (e.g. the HTTP transport's
835    /// single-request `/refs/advance` endpoint, mkit #408).
836    fn supports_atomic_advance(&self) -> bool {
837        false
838    }
839}
840
841/// Result of [`Transport::advance_refs`] — a two-ref branch advance.
842#[derive(Debug, Clone, Copy, PartialEq, Eq)]
843pub enum AdvanceOutcome {
844    /// Both refs were updated.
845    Committed,
846    /// The head precondition did not hold. An atomic transport leaves
847    /// nothing changed. Callers re-read the head (SPEC-TRANSPORT §7): if it
848    /// already holds their target, a retried write landed and the advance
849    /// succeeded; otherwise the branch moved under them (non-fast-forward).
850    HeadConflict,
851    /// The packmap precondition did not hold (a concurrent pusher advanced
852    /// the chain). Callers re-read the packmap and retry.
853    PackmapConflict,
854}
855
856/// Result of an advance that can consume upload tickets.
857#[derive(Debug, Clone, Copy, PartialEq, Eq)]
858#[non_exhaustive]
859pub enum CommitOutcome {
860    /// The advance reached the ordinary two-ref outcome.
861    Advanced(AdvanceOutcome),
862    /// A ticket was invalid, expired, incomplete, or was already consumed.
863    TicketRejected,
864    /// A delta base is not yet available to this repository.
865    DeltaBaseUnavailable,
866    /// The packlist names content not in this repository.
867    PacklistNotInRepository,
868}
869
870/// Server limits relevant to splitting a push into packs and advances.
871#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
872pub struct UploadLimits {
873    /// Maximum accepted serialized pack bytes, if advertised.
874    pub max_pack_bytes: Option<u64>,
875    /// Maximum number of upload tickets one advance may consume.
876    pub tickets_per_advance: Option<usize>,
877    /// First serialized pack size that requires a ticket, when known.
878    /// A transport may stop requiring tickets after discovery fallback.
879    pub ticket_threshold_bytes: Option<u64>,
880}
881
882// ---------------------------------------------------------------------------
883// async_shim — sync/async bridge for transports that wrap an async cipher
884// ---------------------------------------------------------------------------
885
886/// Sync-over-async shim for transports whose underlying cipher / I/O is
887/// async (e.g. `commonware-stream::encrypted`) but whose
888/// [`Transport`] trait surface is intentionally sync.
889///
890/// Lives in `mkit-core` (the trait crate) because it is generic
891/// infrastructure — multiple transports and sparse-checkout's transport
892/// layer will reuse the same plug-in point. It does **not** depend on
893/// `tokio`, `commonware-runtime`, or any concrete executor; callers
894/// pick the runner.
895///
896/// # Why a trait
897///
898/// `mkit-transport-enc` and (once its transport layer lands) `mkit-core::sparse` need to
899/// drive `async fn` bodies from a sync method. Hard-coding
900/// `tokio::runtime::Handle::block_on` would bleed tokio across the
901/// workspace; hard-coding `commonware_runtime::deterministic` would
902/// mean production = tests. A pluggable `Executor` keeps the
903/// runtime-choice at the consumer crate.
904pub mod async_shim {
905    /// Drives an async future to completion synchronously. Pluggable so
906    /// callers can choose between `tokio`, `commonware-runtime`'s
907    /// deterministic runner (tests), or the planned production tokio
908    /// runner without `mkit-core` having to compile-time depend on a
909    /// specific runtime crate.
910    ///
911    /// Implementations MUST be re-entrancy-safe in the sense expected
912    /// by the chosen runtime — calling `block_on` from inside an
913    /// already-running task on the same runtime will typically panic
914    /// or deadlock. The shim's contract is "synchronous external API
915    /// wraps async internals", not "arbitrary async-from-sync
916    /// recursion".
917    pub trait Executor: Send + Sync {
918        /// Block the current thread until `fut` resolves.
919        fn block_on<F, T>(&self, fut: F) -> T
920        where
921            F: core::future::Future<Output = T> + Send,
922            T: Send;
923    }
924}
925
926// ---------------------------------------------------------------------------
927// Tests
928// ---------------------------------------------------------------------------
929
930#[cfg(test)]
931mod tests {
932    use super::*;
933
934    #[test]
935    fn pack_key_hex_roundtrip() {
936        let bytes = [0x42u8; 32];
937        let pk = PackKey::new(bytes);
938        let hex = pk.to_hex();
939        assert_eq!(hex.len(), 64);
940        let pk2 = pack_key_from_hex(&hex).unwrap();
941        assert_eq!(pk, pk2);
942    }
943
944    #[test]
945    fn is_retryable_matches_spec() {
946        assert!(is_retryable(&TransportError::ConnectionFailed));
947        assert!(is_retryable(&TransportError::ServerError { status: 500 }));
948        assert!(is_retryable(&TransportError::ServerError { status: 503 }));
949        assert!(is_retryable(&TransportError::ServerError { status: 429 }));
950        assert!(!is_retryable(&TransportError::ServerError { status: 404 }));
951        assert!(!is_retryable(&TransportError::ServerError { status: 401 }));
952        assert!(!is_retryable(&TransportError::PackNotFound));
953        assert!(!is_retryable(&TransportError::AccessDenied));
954        assert!(!is_retryable(&TransportError::TlsConfiguration(
955            "bad CA file".into()
956        )));
957        assert!(!is_retryable(&TransportError::RefConflict));
958    }
959
960    #[test]
961    fn backoff_default_ladder_is_1_2_4_8_16() {
962        let delays: Vec<Duration> = BackoffIterator::new().collect();
963        assert_eq!(
964            delays,
965            vec![
966                Duration::from_secs(1),
967                Duration::from_secs(2),
968                Duration::from_secs(4),
969                Duration::from_secs(8),
970                Duration::from_secs(16),
971            ]
972        );
973    }
974
975    #[test]
976    fn backoff_caps_at_max() {
977        let cap = Duration::from_secs(10);
978        let delays: Vec<Duration> = BackoffIterator::with(Duration::from_secs(8), cap, 5).collect();
979        // 8s, then cap (16s would exceed 10s cap; clamped to 10s)
980        assert_eq!(delays[0], Duration::from_secs(8));
981        for d in &delays[1..] {
982            assert!(*d <= cap);
983        }
984    }
985
986    // -----------------------------------------------------------------
987    // upload_pack_streaming / download_pack_streaming default impls
988    // -----------------------------------------------------------------
989
990    /// Minimal in-memory [`Transport`] that only implements the
991    /// required whole-buffer methods, so its `*_streaming` behavior is
992    /// entirely the trait's default impl under test.
993    #[derive(Default)]
994    struct RecordingTransport {
995        uploaded: std::sync::Mutex<Option<(Vec<u8>, PackKey)>>,
996        stored: std::sync::Mutex<std::collections::HashMap<[u8; 32], Vec<u8>>>,
997    }
998
999    impl Transport for RecordingTransport {
1000        fn upload_pack(&self, bytes: &[u8], key: &PackKey) -> TransportResult<()> {
1001            *self.uploaded.lock().unwrap() = Some((bytes.to_vec(), *key));
1002            self.stored
1003                .lock()
1004                .unwrap()
1005                .insert(*key.as_bytes(), bytes.to_vec());
1006            Ok(())
1007        }
1008
1009        fn download_pack(&self, key: &PackKey) -> TransportResult<Vec<u8>> {
1010            self.stored
1011                .lock()
1012                .unwrap()
1013                .get(key.as_bytes())
1014                .cloned()
1015                .ok_or(TransportError::PackNotFound)
1016        }
1017
1018        fn pack_exists(&self, _key: &PackKey) -> TransportResult<bool> {
1019            unimplemented!("not exercised by these tests")
1020        }
1021
1022        fn update_ref(
1023            &self,
1024            _name: &str,
1025            _condition: RefWriteCondition,
1026            _hash: &Hash,
1027        ) -> TransportResult<()> {
1028            unimplemented!("not exercised by these tests")
1029        }
1030
1031        fn read_ref(&self, _name: &str) -> TransportResult<Option<Hash>> {
1032            unimplemented!("not exercised by these tests")
1033        }
1034
1035        fn list_refs(&self, _prefix: &str) -> TransportResult<Vec<Ref>> {
1036            unimplemented!("not exercised by these tests")
1037        }
1038    }
1039
1040    #[test]
1041    fn ticket_hooks_default_to_existing_transport_behavior() {
1042        let transport = RecordingTransport::default();
1043        let bytes = b"test pack";
1044        let key = PackKey::new(crate::hash::hash(bytes));
1045        transport
1046            .upload_pack_via_ref(bytes, &key, "refs/heads/main")
1047            .unwrap();
1048        assert_eq!(transport.download_pack(&key).unwrap(), bytes);
1049        transport
1050            .upload_blob_via_ref(bytes, &key, "refs/heads/main")
1051            .unwrap();
1052        assert_eq!(transport.upload_limits(), UploadLimits::default());
1053    }
1054
1055    fn chunks_of(data: &[u8], chunk_len: usize) -> Vec<PackChunk> {
1056        if data.is_empty() {
1057            return vec![PackChunk {
1058                offset: 0,
1059                data: Vec::new(),
1060                last: true,
1061            }];
1062        }
1063        let mut out = Vec::new();
1064        let mut offset = 0usize;
1065        while offset < data.len() {
1066            let end = core::cmp::min(offset + chunk_len, data.len());
1067            out.push(PackChunk {
1068                offset: offset as u64,
1069                data: data[offset..end].to_vec(),
1070                last: end == data.len(),
1071            });
1072            offset = end;
1073        }
1074        out
1075    }
1076
1077    #[test]
1078    fn upload_pack_streaming_default_delegates_to_upload_pack() {
1079        let t = RecordingTransport::default();
1080        let payload = b"hello mkit pack bytes".repeat(100);
1081        let key = PackKey::new([0x11; 32]);
1082        let mut it = chunks_of(&payload, 7).into_iter().map(Ok);
1083
1084        t.upload_pack_streaming(&key, payload.len() as u64, &mut it)
1085            .expect("streaming upload via default impl");
1086
1087        let (got_bytes, got_key) = t.uploaded.lock().unwrap().clone().expect("upload recorded");
1088        assert_eq!(got_bytes, payload);
1089        assert_eq!(got_key, key);
1090    }
1091
1092    #[test]
1093    fn upload_pack_streaming_default_rejects_missing_last_chunk() {
1094        let t = RecordingTransport::default();
1095        let key = PackKey::new([0x22; 32]);
1096        // No chunk at all — total_bytes = 0 still requires one `last =
1097        // true` chunk per the trait contract.
1098        let mut it = core::iter::empty();
1099
1100        let err = t
1101            .upload_pack_streaming(&key, 0, &mut it)
1102            .expect_err("must reject a stream with no last=true chunk");
1103        assert!(matches!(err, TransportError::ProtocolError));
1104    }
1105
1106    #[test]
1107    fn upload_pack_streaming_default_rejects_total_bytes_mismatch() {
1108        let t = RecordingTransport::default();
1109        let key = PackKey::new([0x33; 32]);
1110        let mut it = core::iter::once(Ok(PackChunk {
1111            offset: 0,
1112            data: vec![1, 2, 3],
1113            last: true,
1114        }));
1115
1116        // Declared total (10) does not match the 3 bytes actually
1117        // streamed.
1118        let err = t
1119            .upload_pack_streaming(&key, 10, &mut it)
1120            .expect_err("must reject a total_bytes/accumulated-length mismatch");
1121        assert!(matches!(err, TransportError::ProtocolError));
1122    }
1123
1124    #[test]
1125    fn upload_pack_streaming_default_propagates_chunk_error() {
1126        let t = RecordingTransport::default();
1127        let key = PackKey::new([0x44; 32]);
1128        let mut it = core::iter::once(Err(TransportError::ConnectionFailed));
1129
1130        let err = t
1131            .upload_pack_streaming(&key, 0, &mut it)
1132            .expect_err("must propagate an error yielded mid-stream");
1133        assert!(matches!(err, TransportError::ConnectionFailed));
1134    }
1135
1136    #[test]
1137    fn upload_pack_streaming_default_rejects_oversize_total() {
1138        let t = RecordingTransport::default();
1139        let key = PackKey::new([0x55; 32]);
1140        let mut it = core::iter::empty();
1141
1142        let err = t
1143            .upload_pack_streaming(&key, PACK_BODY_LIMIT + 1, &mut it)
1144            .expect_err("must reject total_bytes above PACK_BODY_LIMIT");
1145        assert!(matches!(err, TransportError::PayloadTooLarge(_)));
1146    }
1147
1148    #[test]
1149    fn download_pack_streaming_default_wraps_whole_pack() {
1150        let t = RecordingTransport::default();
1151        let key = PackKey::new([0x66; 32]);
1152        let payload = vec![9u8; 4096];
1153        t.upload_pack(&payload, &key).unwrap();
1154
1155        let mut stream = t.download_pack_streaming(&key).expect("stream opens");
1156        let first = stream.next().expect("one chunk").expect("no error");
1157        assert_eq!(first.data, payload);
1158        assert!(first.last);
1159        assert!(
1160            stream.next().is_none(),
1161            "default impl yields exactly one chunk"
1162        );
1163    }
1164
1165    #[test]
1166    fn download_pack_streaming_default_propagates_not_found() {
1167        let t = RecordingTransport::default();
1168        let key = PackKey::new([0x77; 32]);
1169        // `Box<dyn Iterator<..>>`'s `Ok` type isn't `Debug`, so match
1170        // instead of `expect_err`.
1171        match t.download_pack_streaming(&key) {
1172            Err(TransportError::PackNotFound) => {}
1173            Err(other) => panic!("expected PackNotFound, got {other:?}"),
1174            Ok(_) => panic!("missing pack must fail before any chunk is produced"),
1175        }
1176    }
1177
1178    // -----------------------------------------------------------------
1179    // retrying() driver
1180    // -----------------------------------------------------------------
1181
1182    fn test_backoff() -> BackoffIterator {
1183        BackoffIterator::with(Duration::from_millis(1), Duration::from_millis(1), 5)
1184    }
1185
1186    fn no_sleep(_delay: Duration) {}
1187
1188    #[test]
1189    fn retrying_succeeds_on_first_try_without_sleeping() {
1190        use core::sync::atomic::{AtomicUsize, Ordering};
1191
1192        fn record_sleep(_delay: Duration) {
1193            SLEEPS.fetch_add(1, Ordering::SeqCst);
1194        }
1195
1196        static CALLS: AtomicUsize = AtomicUsize::new(0);
1197        static SLEEPS: AtomicUsize = AtomicUsize::new(0);
1198        CALLS.store(0, Ordering::SeqCst);
1199        SLEEPS.store(0, Ordering::SeqCst);
1200
1201        let result = retrying::<u32>(
1202            || {
1203                CALLS.fetch_add(1, Ordering::SeqCst);
1204                Ok(7)
1205            },
1206            test_backoff,
1207            record_sleep,
1208        );
1209
1210        assert_eq!(result.unwrap(), 7);
1211        assert_eq!(CALLS.load(Ordering::SeqCst), 1);
1212        assert_eq!(SLEEPS.load(Ordering::SeqCst), 0);
1213    }
1214
1215    #[test]
1216    fn retrying_recovers_after_transient_connection_failures() {
1217        use core::sync::atomic::{AtomicUsize, Ordering};
1218        static CALLS: AtomicUsize = AtomicUsize::new(0);
1219        CALLS.store(0, Ordering::SeqCst);
1220
1221        let result = retrying::<u32>(
1222            || {
1223                let n = CALLS.fetch_add(1, Ordering::SeqCst);
1224                if n < 3 {
1225                    Err(TransportError::ConnectionFailed)
1226                } else {
1227                    Ok(42)
1228                }
1229            },
1230            test_backoff,
1231            no_sleep,
1232        );
1233
1234        assert_eq!(result.unwrap(), 42);
1235        assert_eq!(CALLS.load(Ordering::SeqCst), 4);
1236    }
1237
1238    #[test]
1239    fn retrying_gives_up_after_ladder_exhausts() {
1240        use core::sync::atomic::{AtomicUsize, Ordering};
1241        static CALLS: AtomicUsize = AtomicUsize::new(0);
1242        CALLS.store(0, Ordering::SeqCst);
1243
1244        let result = retrying::<u32>(
1245            || {
1246                CALLS.fetch_add(1, Ordering::SeqCst);
1247                Err(TransportError::ConnectionFailed)
1248            },
1249            test_backoff,
1250            no_sleep,
1251        );
1252
1253        assert!(matches!(result, Err(TransportError::ConnectionFailed)));
1254        // 5-attempt ladder => 1 initial + 5 retries = 6 total calls.
1255        assert_eq!(CALLS.load(Ordering::SeqCst), 6);
1256    }
1257
1258    #[test]
1259    fn retrying_does_not_retry_non_retryable_errors() {
1260        use core::sync::atomic::{AtomicUsize, Ordering};
1261        static CALLS: AtomicUsize = AtomicUsize::new(0);
1262        CALLS.store(0, Ordering::SeqCst);
1263
1264        let result = retrying::<u32>(
1265            || {
1266                CALLS.fetch_add(1, Ordering::SeqCst);
1267                Err(TransportError::PackNotFound)
1268            },
1269            test_backoff,
1270            no_sleep,
1271        );
1272
1273        assert!(matches!(result, Err(TransportError::PackNotFound)));
1274        assert_eq!(CALLS.load(Ordering::SeqCst), 1);
1275    }
1276}