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}