Skip to main content

chorus_client/
protocol.rs

1use std::collections::{BTreeMap, HashMap, VecDeque};
2use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
3use std::sync::{Arc, Mutex, OnceLock};
4use std::time::Duration;
5
6use bytes::Bytes;
7use futures::future::join_all;
8use futures::stream::{FuturesUnordered, StreamExt};
9use sha2::{Digest, Sha256};
10use tokio::sync::watch;
11
12use crate::grpc::pack_append;
13use crate::metrics::Metrics;
14use crate::record::{RecordError, RecordFrame};
15use crate::transport::{
16    AppendToken, LaneDurableChange, PackedAppend, Replica, ReplicaSnapshot, TransportCode,
17    TransportError, FORMAT_VERSION, META_FORMAT,
18};
19
20/// Replication widths the protocol supports. Each uses a strict-majority
21/// quorum: 1-of-1, 2-of-3, 3-of-5.
22pub(crate) const SUPPORTED_REPLICA_COUNTS: [usize; 3] = [1, 3, 5];
23
24/// A lane must make some durable-tail progress within this interval whenever
25/// it retains unacknowledged writes. Five seconds is deliberately above the
26/// default exponential-backoff budget while bounding a genuinely stuck lane
27/// well below the retained-byte limit at normal WAL throughputs.
28pub(crate) const DEFAULT_LANE_STALL_TIMEOUT: Duration = Duration::from_secs(5);
29
30/// Strict-majority quorum size for a replica set of `replica_count` zones.
31pub(crate) fn majority(replica_count: usize) -> usize {
32    replica_count / 2 + 1
33}
34
35fn select_recovery_size(sizes: &mut [i64], replica_count: usize) -> Option<i64> {
36    let quorum = majority(replica_count);
37    if sizes.len() < quorum || sizes.len() > replica_count {
38        return None;
39    }
40    sizes.sort_unstable();
41    let unavailable = replica_count - sizes.len();
42    let required_available_support = quorum.checked_sub(unavailable)?;
43    (required_available_support > 0).then(|| sizes[sizes.len() - required_available_support])
44}
45
46pub(crate) type AttemptedBytes = Arc<dyn Fn(u64) + Send + Sync>;
47
48#[derive(Clone, Debug)]
49/// Retry policy for transport operations used by recovery and writes.
50///
51/// Only transient transport codes are retried. `FAILED_PRECONDITION` and
52/// append-open `ABORTED` are terminal fencing signals and stop the writer.
53pub struct ClientConfig {
54    /// Number of retries after the initial attempt.
55    pub max_retries: usize,
56    /// Base exponential-backoff delay. DST uses zero; production should retain
57    /// a nonzero value to avoid synchronized retry pressure.
58    pub retry_base: Duration,
59}
60
61impl Default for ClientConfig {
62    fn default() -> Self {
63        Self {
64            max_retries: 5,
65            retry_base: Duration::from_millis(20),
66        }
67    }
68}
69
70pub(crate) struct QuorumVolume {
71    replicas: Vec<Arc<dyn Replica>>,
72    config: ClientConfig,
73    metadata: HashMap<String, String>,
74    metrics: Arc<Metrics>,
75}
76
77/// A live lane: the ordered work channel into its writer task plus the task
78/// handle that yields the final token at shutdown.
79struct LaneHandle {
80    work: tokio::sync::mpsc::UnboundedSender<LaneBatch>,
81    done: tokio::task::JoinHandle<Option<AppendToken>>,
82    budget: Arc<LaneBudget>,
83    stall_timeout: Arc<LaneStallTimeout>,
84}
85
86#[derive(Debug)]
87struct LaneBudget {
88    outstanding: AtomicUsize,
89    limit: AtomicUsize,
90}
91
92impl LaneBudget {
93    fn new() -> Arc<Self> {
94        Arc::new(Self {
95            outstanding: AtomicUsize::new(0),
96            limit: AtomicUsize::new(usize::MAX),
97        })
98    }
99
100    fn set_limit(&self, limit: usize) {
101        self.limit.store(limit, Ordering::Relaxed);
102    }
103
104    fn try_reserve(self: &Arc<Self>, bytes: usize) -> Option<Arc<LaneReservation>> {
105        let mut current = self.outstanding.load(Ordering::Relaxed);
106        loop {
107            let next = current.checked_add(bytes)?;
108            if next > self.limit.load(Ordering::Relaxed) {
109                return None;
110            }
111            match self.outstanding.compare_exchange_weak(
112                current,
113                next,
114                Ordering::Relaxed,
115                Ordering::Relaxed,
116            ) {
117                Ok(_) => {
118                    return Some(Arc::new(LaneReservation {
119                        budget: Arc::clone(self),
120                        bytes,
121                    }));
122                }
123                Err(observed) => current = observed,
124            }
125        }
126    }
127}
128
129#[derive(Debug)]
130struct LaneStallTimeout {
131    nanos: AtomicU64,
132}
133
134impl LaneStallTimeout {
135    fn new() -> Arc<Self> {
136        Arc::new(Self {
137            nanos: AtomicU64::new(Self::encode(DEFAULT_LANE_STALL_TIMEOUT)),
138        })
139    }
140
141    fn set(&self, timeout: Duration) {
142        self.nanos.store(Self::encode(timeout), Ordering::Relaxed);
143    }
144
145    fn get(&self) -> Duration {
146        Duration::from_nanos(self.nanos.load(Ordering::Relaxed))
147    }
148
149    fn encode(timeout: Duration) -> u64 {
150        u64::try_from(timeout.as_nanos()).unwrap_or(u64::MAX).max(1)
151    }
152}
153
154#[derive(Debug)]
155struct LaneReservation {
156    budget: Arc<LaneBudget>,
157    bytes: usize,
158}
159
160impl Drop for LaneReservation {
161    fn drop(&mut self) {
162        self.budget
163            .outstanding
164            .fetch_sub(self.bytes, Ordering::Relaxed);
165    }
166}
167
168pub(crate) struct Writer {
169    replicas: Vec<Arc<dyn Replica>>,
170    config: ClientConfig,
171    lanes: Vec<Option<LaneHandle>>,
172    admitted: AdmittedPrefix,
173    commits: Arc<CommitTracker>,
174    sealed: bool,
175    metadata: HashMap<String, String>,
176    metrics: Arc<Metrics>,
177}
178
179/// Lightweight description of a live segment's admitted prefix.
180///
181/// Replica lanes retain unacknowledged encoded chunks and their boundaries for
182/// retry. The writer retains only admitted byte/record counts plus the ordered
183/// seal digest and CRC32C; the commit tracker separately holds unresolved
184/// boundaries. Canonical bytes are reconstructed from storage only if degraded
185/// finalization requires enforcement. SHA-256 updates run after lane dispatch
186/// and are awaited only when rotation freezes the prefix.
187struct AdmittedPrefix {
188    records: usize,
189    bytes: usize,
190    digest: DigestState,
191    crc32c: u32,
192}
193
194enum DigestState {
195    Ready(Sha256),
196    Pending {
197        sender: tokio::sync::mpsc::UnboundedSender<Arc<[Bytes]>>,
198        task: tokio::task::JoinHandle<Sha256>,
199    },
200}
201
202#[derive(Clone, Debug, Default, Eq, PartialEq)]
203pub(crate) struct SealReport {
204    finalized: Vec<Option<ReplicaSnapshot>>,
205}
206
207impl SealReport {
208    pub fn all_replicas_finalized(&self) -> bool {
209        // the default report (slow-path enforcement) is empty and must keep
210        // requesting targeted repair
211        !self.finalized.is_empty()
212            && self.finalized.iter().flatten().count() == self.finalized.len()
213    }
214}
215
216#[derive(Clone, Debug, Default)]
217pub(crate) struct CanonicalPrefix {
218    bytes: Vec<u8>,
219    records: Vec<RecordFrame>,
220    record_ends: Vec<usize>,
221}
222
223/// The canonical content of a fenced tail segment, computed by
224/// [`QuorumVolume::recover_for_seal`]. The seal *decision* (advancing the
225/// manifest `tail_base` with the canonical digest) belongs to the caller;
226/// [`QuorumVolume::enforce_seal`] then installs and finalizes the bytes.
227pub(crate) struct RecoveredTail {
228    canonical: CanonicalPrefix,
229    had_discarded_suffix: bool,
230}
231
232/// Result of fencing one manifest candidate in recovery order.
233pub(crate) enum RecoveryCandidate {
234    /// A quorum confirmed the object name is absent.
235    Absent,
236    /// The committed complete-record prefix is empty. A writer is reusable only
237    /// when every replica supplied a live zero-length session; finalized-empty
238    /// or partial-record witnesses prove the empty frontier but require the
239    /// caller to retire this object name.
240    Empty {
241        reusable_writer: Option<Box<Writer>>,
242    },
243    /// A non-empty committed prefix, frozen for seal enforcement.
244    NonEmpty(RecoveredTail),
245}
246
247impl RecoveredTail {
248    pub fn len(&self) -> usize {
249        self.canonical.len()
250    }
251
252    /// SHA-256 hex digest of the exact canonical bytes (the manifest
253    /// `chorus.seal_digest` value).
254    pub fn digest(&self) -> String {
255        digest_bytes(&self.canonical.bytes)
256    }
257
258    /// Full-object CRC32C of the exact canonical bytes.
259    pub fn crc32c(&self) -> u32 {
260        crc32c::crc32c(&self.canonical.bytes)
261    }
262
263    pub fn canonical(&self) -> &CanonicalPrefix {
264        &self.canonical
265    }
266
267    /// Whether any fenced witness extended past the recovered complete-record
268    /// prefix. A pending successor above such a tail gap is speculative: the
269    /// engine could not have acknowledged it before the missing tail record.
270    pub fn had_discarded_suffix(&self) -> bool {
271        self.had_discarded_suffix
272    }
273}
274
275#[derive(Clone, Debug)]
276enum CommitFailure {
277    Poisoned,
278    Fenced(String),
279    Transport(TransportError),
280}
281
282impl CommitFailure {
283    fn protocol_error(&self) -> ProtocolError {
284        match self {
285            Self::Poisoned => ProtocolError::Poisoned,
286            Self::Fenced(error) => ProtocolError::Fenced(error.clone()),
287            Self::Transport(error) => ProtocolError::Transport(error.clone()),
288        }
289    }
290}
291
292#[derive(Clone, Debug, Default)]
293struct CommitSnapshot {
294    committed: usize,
295    failure: Option<CommitFailure>,
296}
297
298#[derive(Debug)]
299struct LaneCommitState {
300    durable: i64,
301    represented_end: i64,
302    finished: bool,
303    error: Option<TransportError>,
304}
305
306impl Default for LaneCommitState {
307    fn default() -> Self {
308        Self {
309            durable: 0,
310            represented_end: 0,
311            finished: true,
312            error: None,
313        }
314    }
315}
316
317#[derive(Debug)]
318struct CommitState {
319    boundaries: VecDeque<i64>,
320    admitted: usize,
321    committed: usize,
322    committed_bytes: usize,
323    failure: Option<CommitFailure>,
324    lanes: Vec<LaneCommitState>,
325}
326
327/// Segment-scoped aggregation of lane durability. Lanes publish monotonic byte
328/// offsets; this tracker resolves the quorum offset against only the unresolved
329/// record boundaries and broadcasts one contiguous record watermark.
330struct CommitTracker {
331    quorum: usize,
332    state: Mutex<CommitState>,
333    updates: watch::Sender<CommitSnapshot>,
334    metrics: Arc<Metrics>,
335}
336
337pub(crate) struct CommitRange {
338    first_offset: usize,
339    end_offset: usize,
340    updates: watch::Receiver<CommitSnapshot>,
341}
342
343#[cfg(test)]
344pub(crate) struct PendingCommit {
345    pub logical_offset: u64,
346    updates: watch::Receiver<CommitSnapshot>,
347}
348
349#[cfg(test)]
350impl PendingCommit {
351    pub async fn wait(mut self) -> Result<u64, ProtocolError> {
352        loop {
353            let snapshot = self.updates.borrow_and_update().clone();
354            if snapshot.committed > self.logical_offset as usize {
355                return Ok(self.logical_offset);
356            }
357            if let Some(failure) = snapshot.failure {
358                return Err(failure.protocol_error());
359            }
360            self.updates
361                .changed()
362                .await
363                .map_err(|_| ProtocolError::PipelineClosed)?;
364        }
365    }
366}
367
368#[derive(Debug, thiserror::Error)]
369pub(crate) enum ProtocolError {
370    #[error("the replica count must be 1, 3, or 5")]
371    ReplicaCount,
372    #[error("operation did not reach a replica quorum")]
373    NoQuorum,
374    #[error("writer is poisoned by an indeterminate record and must be recovered")]
375    Poisoned,
376    #[error("writer was fenced: {0}")]
377    Fenced(String),
378    #[error("recovery witnesses contain different bytes at record {record_index}")]
379    ConflictingPrefix { record_index: usize },
380    #[error("recovery prefix has {actual} records, expected at least {expected}")]
381    RecoveryPrefixTooShort { expected: usize, actual: usize },
382    #[error("recovered seal digest {actual} does not match committed digest {expected}")]
383    SealDigestMismatch { expected: String, actual: String },
384    #[error("recovered seal CRC32C {actual:08x} does not match committed CRC32C {expected:08x}")]
385    SealCrc32cMismatch { expected: u32, actual: u32 },
386    #[error("manifest register is invalid: {0}")]
387    InvalidManifest(String),
388    #[error(
389        "the manifest segment directory is full: truncate the WAL to free \
390         retained sealed segments before sealing again"
391    )]
392    SegmentDirectoryFull,
393    #[error(transparent)]
394    ManifestStore(#[from] crate::manifest_store::ManifestStoreError),
395    #[error("manifest register is unavailable")]
396    ManifestUnavailable,
397    #[error("commit pipeline closed before reporting a result")]
398    PipelineClosed,
399    #[error("segment writer is sealed")]
400    Finalized,
401    #[error(transparent)]
402    Record(#[from] RecordError),
403    #[error("transport error: {0}")]
404    Transport(#[from] TransportError),
405}
406
407impl CanonicalPrefix {
408    pub fn len(&self) -> usize {
409        self.records.len()
410    }
411
412    pub fn into_records(self) -> Vec<RecordFrame> {
413        self.records
414    }
415
416    fn truncate(&mut self, records: usize) {
417        self.records.truncate(records);
418        self.record_ends.truncate(records);
419        self.bytes.truncate(self.committed_bytes_len(records));
420    }
421
422    fn committed_bytes_len(&self, records: usize) -> usize {
423        records
424            .checked_sub(1)
425            .and_then(|index| self.record_ends.get(index).copied())
426            .unwrap_or(0)
427    }
428
429    fn record_bytes(&self, index: usize) -> &[u8] {
430        let start = index
431            .checked_sub(1)
432            .and_then(|previous| self.record_ends.get(previous).copied())
433            .unwrap_or(0);
434        &self.bytes[start..self.record_ends[index]]
435    }
436
437    fn from_snapshot(snapshot: &ReplicaSnapshot) -> Self {
438        let (records, consumed) = RecordFrame::decode_complete_prefix(&snapshot.bytes);
439        let mut record_ends = Vec::with_capacity(records.len());
440        let mut end = 0usize;
441        for record in &records {
442            end += record.encode().expect("a decoded record must encode").len();
443            record_ends.push(end);
444        }
445        Self {
446            bytes: snapshot.bytes[..consumed].to_vec(),
447            records,
448            record_ends,
449        }
450    }
451}
452
453impl Default for AdmittedPrefix {
454    fn default() -> Self {
455        Self {
456            records: 0,
457            bytes: 0,
458            digest: DigestState::Ready(Sha256::new()),
459            crc32c: 0,
460        }
461    }
462}
463
464impl AdmittedPrefix {
465    fn len(&self) -> usize {
466        self.records
467    }
468
469    fn is_empty(&self) -> bool {
470        self.records == 0
471    }
472
473    fn bytes_len(&self) -> usize {
474        self.bytes
475    }
476
477    fn extend_metadata(&mut self, chunks: &[Bytes]) {
478        for chunk in chunks {
479            self.crc32c = crc32c::crc32c_append(self.crc32c, chunk);
480            self.bytes += chunk.len();
481            self.records += 1;
482        }
483    }
484
485    fn queue_digest(&mut self, chunks: Arc<[Bytes]>) {
486        let previous = std::mem::replace(&mut self.digest, DigestState::Ready(Sha256::new()));
487        let pending = match previous {
488            DigestState::Ready(mut hasher) => {
489                let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel::<Arc<[Bytes]>>();
490                let task = tokio::spawn(async move {
491                    while let Some(chunks) = receiver.recv().await {
492                        for chunk in chunks.iter() {
493                            hasher.update(chunk);
494                        }
495                        tokio::task::yield_now().await;
496                    }
497                    hasher
498                });
499                sender.send(chunks).expect("digest worker failed");
500                DigestState::Pending { sender, task }
501            }
502            DigestState::Pending { sender, task } => {
503                sender.send(chunks).expect("digest worker failed");
504                DigestState::Pending { sender, task }
505            }
506        };
507        self.digest = pending;
508    }
509
510    async fn digest(&mut self) -> String {
511        let pending = std::mem::replace(&mut self.digest, DigestState::Ready(Sha256::new()));
512        let hasher = match pending {
513            DigestState::Ready(hasher) => hasher,
514            DigestState::Pending { sender, task } => {
515                drop(sender);
516                task.await.expect("ordered digest task failed")
517            }
518        };
519        let digest = digest_hex(hasher.clone().finalize());
520        self.digest = DigestState::Ready(hasher);
521        digest
522    }
523
524    async fn shutdown_digest(&mut self) {
525        let digest = std::mem::replace(&mut self.digest, DigestState::Ready(Sha256::new()));
526        if let DigestState::Pending { sender, task } = digest {
527            drop(sender);
528            let _ = task.await;
529        }
530    }
531
532    fn crc32c(&self) -> u32 {
533        self.crc32c
534    }
535}
536
537impl Drop for AdmittedPrefix {
538    fn drop(&mut self) {
539        let digest = std::mem::replace(&mut self.digest, DigestState::Ready(Sha256::new()));
540        if let DigestState::Pending { sender, task } = digest {
541            drop(sender);
542            task.abort();
543        }
544    }
545}
546
547impl CommitTracker {
548    fn new(lanes: usize, quorum: usize, metrics: Arc<Metrics>) -> Arc<Self> {
549        let snapshot = CommitSnapshot::default();
550        let (updates, _) = watch::channel(snapshot);
551        Arc::new(Self {
552            quorum,
553            state: Mutex::new(CommitState {
554                boundaries: VecDeque::new(),
555                admitted: 0,
556                committed: 0,
557                committed_bytes: 0,
558                failure: None,
559                lanes: (0..lanes).map(|_| LaneCommitState::default()).collect(),
560            }),
561            updates,
562            metrics,
563        })
564    }
565
566    fn activate_lane(&self, zone: usize, durable: i64) {
567        let mut state = self
568            .state
569            .lock()
570            .unwrap_or_else(|poisoned| poisoned.into_inner());
571        let lane = &mut state.lanes[zone];
572        let previous_lag = lane_durable_lag(lane);
573        lane.durable = durable.max(0);
574        lane.finished = false;
575        self.metrics
576            .adjust_zone_durable_lag(zone, lane_durable_lag(lane) - previous_lag);
577    }
578
579    fn admitted_len(&self) -> usize {
580        self.state
581            .lock()
582            .unwrap_or_else(|poisoned| poisoned.into_inner())
583            .admitted
584    }
585
586    fn committed_len(&self) -> usize {
587        self.state
588            .lock()
589            .unwrap_or_else(|poisoned| poisoned.into_inner())
590            .committed
591    }
592
593    fn committed_bytes(&self) -> usize {
594        self.state
595            .lock()
596            .unwrap_or_else(|poisoned| poisoned.into_inner())
597            .committed_bytes
598    }
599
600    fn is_poisoned(&self) -> bool {
601        self.state
602            .lock()
603            .unwrap_or_else(|poisoned| poisoned.into_inner())
604            .failure
605            .is_some()
606    }
607
608    fn subscribe(&self) -> watch::Receiver<CommitSnapshot> {
609        self.updates.subscribe()
610    }
611
612    fn admit_window(&self, boundaries: &[i64], represented_zones: &[usize]) -> CommitRange {
613        let mut state = self
614            .state
615            .lock()
616            .unwrap_or_else(|poisoned| poisoned.into_inner());
617        let first_offset = state.admitted;
618        let end = boundaries
619            .last()
620            .copied()
621            .expect("an admitted window is non-empty");
622        state.admitted += boundaries.len();
623        state.boundaries.extend(boundaries.iter().copied());
624        for &zone in represented_zones {
625            let lane = &mut state.lanes[zone];
626            let previous_lag = lane_durable_lag(lane);
627            lane.represented_end = lane.represented_end.max(end);
628            self.metrics
629                .adjust_zone_durable_lag(zone, lane_durable_lag(lane) - previous_lag);
630        }
631        self.recompute_locked(&mut state);
632        CommitRange {
633            first_offset,
634            end_offset: state.admitted,
635            updates: self.subscribe(),
636        }
637    }
638
639    fn publish_durable(&self, zone: usize, durable: i64) {
640        let mut state = self
641            .state
642            .lock()
643            .unwrap_or_else(|poisoned| poisoned.into_inner());
644        let lane = &mut state.lanes[zone];
645        if durable <= lane.durable {
646            return;
647        }
648        let previous_lag = lane_durable_lag(lane);
649        lane.durable = durable;
650        self.metrics
651            .adjust_zone_durable_lag(zone, lane_durable_lag(lane) - previous_lag);
652        self.recompute_locked(&mut state);
653    }
654
655    fn finish_lane(&self, zone: usize, error: Option<TransportError>) {
656        let mut state = self
657            .state
658            .lock()
659            .unwrap_or_else(|poisoned| poisoned.into_inner());
660        let fence = error
661            .as_ref()
662            .filter(|error| error.code.fences_writer())
663            .map(ToString::to_string);
664        let lane = &mut state.lanes[zone];
665        let previous_lag = lane_durable_lag(lane);
666        lane.finished = true;
667        if error.is_some() {
668            lane.error = error;
669        }
670        self.metrics
671            .adjust_zone_durable_lag(zone, lane_durable_lag(lane) - previous_lag);
672        if let Some(fence) = fence {
673            // A takeover fence is writer-wide, not a removable lane failure.
674            // Publish it even when the just-confirmed boundary drained the queue.
675            if state.failure.is_none() {
676                state.failure = Some(CommitFailure::Fenced(fence));
677                self.publish_locked(&state);
678            }
679            return;
680        }
681        self.recompute_locked(&mut state);
682    }
683
684    fn poison(&self) {
685        let mut state = self
686            .state
687            .lock()
688            .unwrap_or_else(|poisoned| poisoned.into_inner());
689        if state.failure.is_none() {
690            state.failure = Some(CommitFailure::Poisoned);
691            self.publish_locked(&state);
692        }
693    }
694
695    fn recompute_locked(&self, state: &mut CommitState) {
696        if state.failure.is_some() {
697            return;
698        }
699        let quorum_watermark = quorum_durable_watermark(&state.lanes, self.quorum);
700        let mut changed = false;
701        while state
702            .boundaries
703            .front()
704            .is_some_and(|boundary| *boundary <= quorum_watermark)
705        {
706            let boundary = state
707                .boundaries
708                .pop_front()
709                .expect("front boundary was present");
710            state.committed += 1;
711            state.committed_bytes =
712                usize::try_from(boundary).expect("record boundaries are nonnegative");
713            changed = true;
714        }
715
716        if let Some(&oldest) = state.boundaries.front() {
717            let possible = state
718                .lanes
719                .iter()
720                .filter(|lane| {
721                    lane.durable >= oldest || (!lane.finished && lane.represented_end >= oldest)
722                })
723                .count();
724            if possible < self.quorum {
725                state.failure = Some(select_commit_failure(&state.lanes, oldest));
726                changed = true;
727            }
728        }
729        if changed {
730            self.publish_locked(state);
731        }
732    }
733
734    fn publish_locked(&self, state: &CommitState) {
735        self.updates.send_replace(CommitSnapshot {
736            committed: state.committed,
737            failure: state.failure.clone(),
738        });
739    }
740}
741
742impl Drop for CommitTracker {
743    fn drop(&mut self) {
744        let state = self
745            .state
746            .get_mut()
747            .unwrap_or_else(|poisoned| poisoned.into_inner());
748        for (zone, lane) in state.lanes.iter_mut().enumerate() {
749            let lag = lane_durable_lag(lane);
750            if lag != 0 {
751                self.metrics.adjust_zone_durable_lag(zone, -lag);
752                lane.finished = true;
753            }
754        }
755    }
756}
757
758fn lane_durable_lag(lane: &LaneCommitState) -> i64 {
759    if lane.finished {
760        0
761    } else {
762        lane.represented_end.saturating_sub(lane.durable).max(0)
763    }
764}
765
766fn quorum_durable_watermark(lanes: &[LaneCommitState], quorum: usize) -> i64 {
767    let mut durables: Vec<_> = lanes.iter().map(|lane| lane.durable).collect();
768    durables.sort_unstable();
769    durables[durables.len() - quorum]
770}
771
772fn select_commit_failure(lanes: &[LaneCommitState], boundary: i64) -> CommitFailure {
773    let represented = lanes.iter().filter(|lane| lane.represented_end >= boundary);
774    let mut fenced = None;
775    let mut terminal = None;
776    for error in represented.filter_map(|lane| lane.error.clone()) {
777        if error.code.fences_writer() {
778            prefer_lower_zone(&mut fenced, error);
779        } else if !error.code.transient() {
780            prefer_lower_zone(&mut terminal, error);
781        }
782    }
783    match (fenced, terminal) {
784        (Some(error), _) => CommitFailure::Fenced(error.to_string()),
785        (None, Some(error)) => CommitFailure::Transport(error),
786        (None, None) => CommitFailure::Poisoned,
787    }
788}
789
790impl CommitRange {
791    pub(crate) fn first_offset(&self) -> usize {
792        self.first_offset
793    }
794
795    pub(crate) fn end_offset(&self) -> usize {
796        self.end_offset
797    }
798
799    #[cfg(test)]
800    pub(crate) fn into_pending(self) -> Vec<PendingCommit> {
801        (self.first_offset..self.end_offset)
802            .map(|logical_offset| PendingCommit {
803                logical_offset: logical_offset as u64,
804                updates: self.updates.clone(),
805            })
806            .collect()
807    }
808
809    pub(crate) fn progress(&mut self) -> (usize, Option<ProtocolError>) {
810        let snapshot = self.updates.borrow_and_update().clone();
811        (
812            snapshot.committed,
813            snapshot.failure.map(|failure| failure.protocol_error()),
814        )
815    }
816
817    pub(crate) async fn changed(&mut self) -> Result<(), ProtocolError> {
818        self.updates
819            .changed()
820            .await
821            .map_err(|_| ProtocolError::PipelineClosed)
822    }
823}
824
825pub(crate) fn digest_bytes(bytes: &[u8]) -> String {
826    digest_hex(Sha256::digest(bytes))
827}
828
829fn digest_hex(digest: impl AsRef<[u8]>) -> String {
830    use std::fmt::Write;
831
832    let digest = digest.as_ref();
833    let mut encoded = String::with_capacity(digest.len() * 2);
834    for byte in digest {
835        write!(&mut encoded, "{byte:02x}").expect("writing to a String cannot fail");
836    }
837    encoded
838}
839
840impl QuorumVolume {
841    /// A volume whose every created or replaced object carries `metadata` —
842    /// the constant format marker. Rewrites (canonical write-back, repair)
843    /// replace custom metadata wholesale with the same constant; chain
844    /// position lives in the manifest's segment directory, never on the
845    /// object.
846    pub fn with_metadata(
847        replicas: Vec<Arc<dyn Replica>>,
848        config: ClientConfig,
849        metadata: HashMap<String, String>,
850        metrics: Arc<Metrics>,
851    ) -> Result<Self, ProtocolError> {
852        if !SUPPORTED_REPLICA_COUNTS.contains(&replicas.len()) {
853            return Err(ProtocolError::ReplicaCount);
854        }
855        Ok(Self {
856            replicas,
857            config,
858            metadata,
859            metrics,
860        })
861    }
862
863    fn quorum(&self) -> usize {
864        majority(self.replicas.len())
865    }
866
867    /// Conditionally create one new segment generation on a quorum.
868    ///
869    /// Only objects created by this call count toward writer eligibility. An
870    /// `AlreadyExists` object never lets a racing process adopt that segment.
871    pub async fn create_writer(&self) -> Result<Writer, ProtocolError> {
872        let metadata = self.metadata.clone();
873        let creates = join_all(
874            self.replicas
875                .iter()
876                .map(|replica| create_session_with_retry(replica, metadata.clone(), &self.config)),
877        )
878        .await;
879        let tokens: Vec<_> = creates.into_iter().flatten().collect();
880        if tokens.len() < self.quorum() {
881            return Err(ProtocolError::NoQuorum);
882        }
883        Ok(Writer::new(
884            self.replicas.clone(),
885            self.config.clone(),
886            self.metadata.clone(),
887            tokens,
888            Arc::clone(&self.metrics),
889        ))
890    }
891
892    /// Fence an existing segment, select a compatible recovery quorum, and
893    /// rewrite its canonical prefix to the same witnesses. The returned tail
894    /// contains the canonical bytes and the taken-over replica sessions needed
895    /// for seal enforcement; application appends start in a new segment.
896    pub async fn recover_for_seal(
897        &self,
898        expected_records: Option<usize>,
899    ) -> Result<RecoveredTail, ProtocolError> {
900        match self.recover_candidate(expected_records).await? {
901            RecoveryCandidate::NonEmpty(tail) => Ok(tail),
902            RecoveryCandidate::Empty { .. } | RecoveryCandidate::Absent => {
903                Err(ProtocolError::RecoveryPrefixTooShort {
904                    expected: expected_records.unwrap_or(1),
905                    actual: 0,
906                })
907            }
908        }
909    }
910
911    /// Fence before observing size, then recover the committed record prefix.
912    ///
913    /// `takeover_current` resolves the latest object generation explicitly,
914    /// then opens that exact generation to revoke any stale stream and obtain
915    /// authoritative `persisted_size`. The identity stat is tail-blind and
916    /// never participates in prefix selection. Object bytes are read only after
917    /// non-empty witnesses are frozen, because they are not readable while
918    /// appendable.
919    pub(crate) async fn recover_candidate(
920        &self,
921        expected_records: Option<usize>,
922    ) -> Result<RecoveryCandidate, ProtocolError> {
923        enum Observation {
924            Live(AppendToken),
925            Finalized(ReplicaSnapshot),
926            Missing(usize),
927        }
928
929        let attempts = join_all(self.replicas.iter().map(|replica| {
930            let replica = Arc::clone(replica);
931            let config = self.config.clone();
932            async move {
933                match takeover_current_with_retry(&replica, &config).await {
934                    Ok(token) => Ok::<_, ProtocolError>(Observation::Live(token)),
935                    Err(error) if error.code == TransportCode::FailedPrecondition => {
936                        let snapshot = snapshot_with_retry(&replica, &config).await?;
937                        if valid_format(&snapshot.metadata) {
938                            Ok(Observation::Finalized(snapshot))
939                        } else {
940                            Err(ProtocolError::NoQuorum)
941                        }
942                    }
943                    Err(error) if error.code == TransportCode::NotFound => {
944                        Ok(Observation::Missing(error.zone))
945                    }
946                    Err(error)
947                        if error.code.transient()
948                            || matches!(
949                                error.code,
950                                TransportCode::DataLoss | TransportCode::Ambiguous
951                            ) =>
952                    {
953                        Err(ProtocolError::NoQuorum)
954                    }
955                    Err(error) => Err(error.into()),
956                }
957            }
958        }))
959        .await;
960
961        let mut observations = Vec::new();
962        for attempt in attempts {
963            match attempt {
964                Ok(observation) => observations.push(observation),
965                Err(ProtocolError::NoQuorum) => {}
966                Err(error) => return Err(error),
967            }
968        }
969        if observations.len() < self.quorum() {
970            return Err(ProtocolError::NoQuorum);
971        }
972        if observations
973            .iter()
974            .all(|observation| matches!(observation, Observation::Missing(_)))
975        {
976            return Ok(RecoveryCandidate::Absent);
977        }
978
979        let missing_zones: Vec<_> = observations
980            .iter()
981            .filter_map(|observation| match observation {
982                Observation::Missing(zone) => Some(*zone),
983                _ => None,
984            })
985            .collect();
986
987        let mut sizes = observations
988            .iter()
989            .map(|observation| match observation {
990                Observation::Live(token) => token.persisted_size,
991                Observation::Finalized(snapshot) => snapshot.persisted_size,
992                Observation::Missing(_) => 0,
993            })
994            .collect::<Vec<_>>();
995        if sizes.len() < self.quorum() {
996            return Err(ProtocolError::NoQuorum);
997        }
998        let max_observed_size = *sizes
999            .iter()
1000            .max()
1001            .expect("a quorum supplied at least one size");
1002        // An acknowledged write quorum can hide at most N-k of its witnesses
1003        // among the unavailable replicas. Therefore a potentially acknowledged
1004        // offset must appear on at least Q-(N-k) of the k available sizes. This
1005        // is the Qth-largest size when all replicas answer, the maximum when
1006        // only a read quorum answers, and the corresponding intermediate order
1007        // statistic for larger replica sets. Failed reads can remove evidence,
1008        // but can never lower the selected offset below a quorum that may have
1009        // acknowledged it.
1010        let committed_size = select_recovery_size(&mut sizes, self.replicas.len())
1011            .expect("a reachable majority intersects every write quorum");
1012
1013        if committed_size == 0 {
1014            let mut empty_witnesses = observations
1015                .iter()
1016                .filter(|observation| {
1017                    matches!(
1018                        observation,
1019                        Observation::Live(token) if token.persisted_size == 0
1020                    ) || matches!(
1021                        observation,
1022                        Observation::Finalized(snapshot) if snapshot.persisted_size == 0
1023                    )
1024                })
1025                .count();
1026            let mut tokens = observations
1027                .into_iter()
1028                .filter_map(|observation| match observation {
1029                    Observation::Live(token) if token.persisted_size == 0 => Some(token),
1030                    _ => None,
1031                })
1032                .collect::<Vec<_>>();
1033            if tokens.len() < self.quorum() {
1034                let creates = join_all(missing_zones.into_iter().map(|zone| {
1035                    create_session_with_retry(
1036                        &self.replicas[zone],
1037                        self.metadata.clone(),
1038                        &self.config,
1039                    )
1040                }))
1041                .await;
1042                for token in creates.into_iter().flatten() {
1043                    empty_witnesses += 1;
1044                    tokens.push(token);
1045                }
1046            }
1047            if empty_witnesses < self.quorum() {
1048                return Err(ProtocolError::NoQuorum);
1049            }
1050            let reusable_writer = (tokens.len() == self.replicas.len()).then(|| {
1051                Box::new(Writer::new(
1052                    self.replicas.clone(),
1053                    self.config.clone(),
1054                    self.metadata.clone(),
1055                    tokens,
1056                    Arc::clone(&self.metrics),
1057                ))
1058            });
1059            return Ok(RecoveryCandidate::Empty { reusable_writer });
1060        }
1061
1062        let mut live_tokens = Vec::new();
1063        let mut recovery_snapshots = Vec::new();
1064        for observation in observations {
1065            match observation {
1066                Observation::Live(token) => live_tokens.push(token),
1067                Observation::Finalized(snapshot) => recovery_snapshots.push(snapshot),
1068                Observation::Missing(_) => {}
1069            }
1070        }
1071        let frozen = join_all(live_tokens.into_iter().map(|mut token| {
1072            let replica = Arc::clone(&self.replicas[token.zone]);
1073            let config = self.config.clone();
1074            async move {
1075                let persisted_size = token.persisted_size;
1076                finalize_with_retry(&replica, &mut token, persisted_size, &config).await
1077            }
1078        }))
1079        .await;
1080        let frozen_zones = frozen
1081            .into_iter()
1082            .flatten()
1083            .map(|snapshot| snapshot.zone)
1084            .collect::<Vec<_>>();
1085        if frozen_zones.len() + recovery_snapshots.len() + missing_zones.len() < self.quorum() {
1086            return Err(ProtocolError::NoQuorum);
1087        }
1088
1089        let committed_size = usize::try_from(committed_size).map_err(|_| {
1090            ProtocolError::InvalidManifest("persisted segment size does not fit in usize".into())
1091        })?;
1092        let max_observed_size = usize::try_from(max_observed_size).map_err(|_| {
1093            ProtocolError::InvalidManifest("persisted segment size does not fit in usize".into())
1094        })?;
1095        let reads = join_all(
1096            frozen_zones
1097                .into_iter()
1098                .map(|zone| snapshot_with_retry(&self.replicas[zone], &self.config)),
1099        )
1100        .await;
1101        for snapshot in reads.into_iter().flatten() {
1102            if valid_format(&snapshot.metadata) {
1103                recovery_snapshots.push(snapshot);
1104            }
1105        }
1106        for snapshot in &mut recovery_snapshots {
1107            snapshot
1108                .bytes
1109                .truncate(snapshot.bytes.len().min(committed_size));
1110        }
1111        recovery_snapshots.extend(missing_zones.into_iter().map(|zone| ReplicaSnapshot {
1112            zone,
1113            generation: 0,
1114            metageneration: 0,
1115            persisted_size: 0,
1116            finalized: true,
1117            crc32c: Some(crc32c::crc32c(&[])),
1118            metadata: self.metadata.clone(),
1119            bytes: Vec::new(),
1120        }));
1121        let (mut canonical, _) = select_canonical_quorum(&recovery_snapshots, self.quorum())?;
1122        let had_discarded_suffix = max_observed_size > canonical.bytes.len();
1123        if let Some(expected) = expected_records {
1124            if canonical.len() < expected {
1125                return Err(ProtocolError::RecoveryPrefixTooShort {
1126                    expected,
1127                    actual: canonical.len(),
1128                });
1129            }
1130            canonical.truncate(expected);
1131        }
1132        if canonical.len() == 0 {
1133            return Ok(RecoveryCandidate::Empty {
1134                reusable_writer: None,
1135            });
1136        }
1137        Ok(RecoveryCandidate::NonEmpty(RecoveredTail {
1138            canonical,
1139            had_discarded_suffix,
1140        }))
1141    }
1142
1143    /// Install and finalize the canonical bytes on at least a quorum of
1144    /// replicas.
1145    ///
1146    /// Enforcement only: the seal decision must already be committed (through
1147    /// the manifest for the active segment, or implied by the following
1148    /// segment's base for chain repair). Replicas are processed with the
1149    /// shortest prefix first so the canonical bytes always remain durable in
1150    /// at least one object throughout the rewrite; the operation is
1151    /// idempotent and convergent under crashes and races because the bytes
1152    /// are fixed by the committed decision.
1153    pub async fn enforce_seal(&self, canonical: &CanonicalPrefix) -> Result<(), ProtocolError> {
1154        let data = Bytes::from(canonical.bytes.clone());
1155        let reads = join_all(
1156            self.replicas
1157                .iter()
1158                .map(|replica| snapshot_with_retry(replica, &self.config)),
1159        )
1160        .await;
1161        let mut witnesses: Vec<ReplicaSnapshot> = Vec::new();
1162        let mut missing: Vec<usize> = Vec::new();
1163        for (zone, read) in reads.into_iter().enumerate() {
1164            match read {
1165                Ok(snapshot) if valid_format(&snapshot.metadata) => witnesses.push(snapshot),
1166                Ok(_) => {}
1167                Err(error) if error.code == TransportCode::NotFound => missing.push(zone),
1168                Err(error) if error.code == TransportCode::DataLoss => {
1169                    match stat_with_retry(&self.replicas[zone], &self.config).await {
1170                        Ok(snapshot) if valid_format(&snapshot.metadata) => {
1171                            // `stat` is deliberately content-blind. Empty bytes
1172                            // force `enforce_witness` to guarded-replace the
1173                            // rotted generation from the canonical prefix.
1174                            witnesses.push(snapshot);
1175                        }
1176                        Ok(_) => {}
1177                        Err(stat_error) if stat_error.code == TransportCode::NotFound => {
1178                            missing.push(zone);
1179                        }
1180                        Err(stat_error) if stat_error.code.transient() => {}
1181                        Err(stat_error) => return Err(stat_error.into()),
1182                    }
1183                }
1184                Err(error) if error.code.transient() => {}
1185                Err(error) => return Err(error.into()),
1186            }
1187        }
1188        // shortest matching prefix first: the best copy is rewritten last
1189        witnesses.sort_by_key(|snapshot| {
1190            let shared = snapshot
1191                .bytes
1192                .iter()
1193                .zip(&canonical.bytes)
1194                .take_while(|(a, b)| a == b)
1195                .count();
1196            (shared == canonical.bytes.len() && snapshot.bytes.len() == canonical.bytes.len())
1197                as usize
1198                * canonical.bytes.len()
1199                + shared
1200        });
1201        let mut finalized = 0usize;
1202        for snapshot in witnesses {
1203            match self.enforce_witness(snapshot, &data).await {
1204                Ok(true) => finalized += 1,
1205                Ok(false) => {}
1206                Err(error) => return Err(error),
1207            }
1208        }
1209        for zone in missing {
1210            let replica = Arc::clone(&self.replicas[zone]);
1211            let Ok(created) =
1212                create_with_retry(&replica, self.metadata.clone(), &self.config).await
1213            else {
1214                continue;
1215            };
1216            match self.enforce_witness(created, &data).await {
1217                Ok(true) => finalized += 1,
1218                Ok(false) => {}
1219                Err(error) => return Err(error),
1220            }
1221        }
1222        if finalized >= self.quorum() {
1223            Ok(())
1224        } else {
1225            Err(ProtocolError::NoQuorum)
1226        }
1227    }
1228
1229    async fn enforce_witness(
1230        &self,
1231        mut snapshot: ReplicaSnapshot,
1232        data: &Bytes,
1233    ) -> Result<bool, ProtocolError> {
1234        let zone = snapshot.zone;
1235        let replica = Arc::clone(&self.replicas[zone]);
1236        for _ in 0..=self.config.max_retries {
1237            if snapshot.finalized {
1238                if snapshot.bytes == data[..] {
1239                    return Ok(true);
1240                }
1241                // wrong finalized content: replace with a fresh generation
1242            } else if snapshot.bytes == data[..] {
1243                let Ok(mut token) = takeover_with_retry(&replica, &snapshot, &self.config).await
1244                else {
1245                    snapshot = snapshot_with_retry(&replica, &self.config).await?;
1246                    continue;
1247                };
1248                match finalize_with_retry(&replica, &mut token, data.len() as i64, &self.config)
1249                    .await
1250                {
1251                    Ok(_) => return Ok(true),
1252                    Err(error) if error.code == TransportCode::FailedPrecondition => {
1253                        snapshot = snapshot_with_retry(&replica, &self.config).await?;
1254                        continue;
1255                    }
1256                    Err(error) if error.code.transient() => return Ok(false),
1257                    Err(error) => return Err(error.into()),
1258                }
1259            }
1260            let mut token = match replace_with_retry(
1261                &replica,
1262                snapshot.clone(),
1263                data.clone(),
1264                self.metadata.clone(),
1265                &self.config,
1266            )
1267            .await
1268            {
1269                Ok(token) => token,
1270                Err(error) if error.code == TransportCode::FailedPrecondition => {
1271                    snapshot = snapshot_with_retry(&replica, &self.config).await?;
1272                    continue;
1273                }
1274                Err(error) if error.code.transient() => return Ok(false),
1275                Err(error) => return Err(error.into()),
1276            };
1277            match finalize_with_retry(&replica, &mut token, data.len() as i64, &self.config).await {
1278                Ok(_) => return Ok(true),
1279                Err(error) if error.code == TransportCode::FailedPrecondition => {
1280                    snapshot = snapshot_with_retry(&replica, &self.config).await?;
1281                }
1282                Err(error) if error.code.transient() => return Ok(false),
1283                Err(error) => return Err(error.into()),
1284            }
1285        }
1286        Ok(false)
1287    }
1288}
1289
1290impl Drop for Writer {
1291    fn drop(&mut self) {
1292        for lane in self.lanes.iter_mut().filter_map(Option::take) {
1293            lane.done.abort();
1294        }
1295    }
1296}
1297
1298impl Writer {
1299    fn new(
1300        replicas: Vec<Arc<dyn Replica>>,
1301        config: ClientConfig,
1302        metadata: HashMap<String, String>,
1303        tokens: Vec<AppendToken>,
1304        metrics: Arc<Metrics>,
1305    ) -> Self {
1306        let mut by_zone: HashMap<usize, AppendToken> = tokens
1307            .into_iter()
1308            .map(|token| (token.zone, token))
1309            .collect();
1310        let commits = CommitTracker::new(
1311            replicas.len(),
1312            majority(replicas.len()),
1313            Arc::clone(&metrics),
1314        );
1315        let lanes = (0..replicas.len())
1316            .map(|zone| {
1317                by_zone.remove(&zone).map(|token| {
1318                    commits.activate_lane(zone, token.persisted_size);
1319                    let (work, rx) = tokio::sync::mpsc::unbounded_channel();
1320                    let replica = Arc::clone(&replicas[zone]);
1321                    let config = config.clone();
1322                    let metrics = Arc::clone(&metrics);
1323                    let budget = LaneBudget::new();
1324                    let stall_timeout = LaneStallTimeout::new();
1325                    let lane_commits = Arc::clone(&commits);
1326                    LaneHandle {
1327                        work,
1328                        done: tokio::spawn(run_lane(
1329                            replica,
1330                            token,
1331                            config,
1332                            metrics,
1333                            rx,
1334                            lane_commits,
1335                            Arc::clone(&stall_timeout),
1336                        )),
1337                        budget,
1338                        stall_timeout,
1339                    }
1340                })
1341            })
1342            .collect();
1343        Self {
1344            replicas,
1345            config,
1346            metadata,
1347            lanes,
1348            admitted: AdmittedPrefix::default(),
1349            commits,
1350            sealed: false,
1351            metrics,
1352        }
1353    }
1354
1355    fn quorum(&self) -> usize {
1356        majority(self.replicas.len())
1357    }
1358
1359    pub fn admitted_len(&self) -> usize {
1360        self.commits.admitted_len()
1361    }
1362
1363    pub(crate) fn set_max_replica_lag_bytes(&mut self, limit: usize) {
1364        for lane in self.lanes.iter().flatten() {
1365            lane.budget.set_limit(limit);
1366        }
1367    }
1368
1369    pub(crate) fn set_lane_stall_timeout(&mut self, timeout: Duration) {
1370        for lane in self.lanes.iter().flatten() {
1371            lane.stall_timeout.set(timeout);
1372        }
1373    }
1374
1375    pub(crate) async fn shutdown_background_tasks(&mut self) {
1376        let lanes = std::mem::take(&mut self.lanes);
1377        let mut tasks = Vec::new();
1378        for lane in lanes.into_iter().flatten() {
1379            lane.done.abort();
1380            tasks.push(lane.done);
1381        }
1382        let _ = join_all(tasks).await;
1383        self.admitted.shutdown_digest().await;
1384        join_all(self.replicas.iter().map(|replica| replica.shutdown())).await;
1385    }
1386
1387    pub fn committed_len(&self) -> usize {
1388        self.commits.committed_len()
1389    }
1390
1391    pub fn is_poisoned(&self) -> bool {
1392        self.commits.is_poisoned()
1393    }
1394
1395    pub fn physical_size(&self) -> usize {
1396        self.commits.committed_bytes()
1397    }
1398
1399    /// SHA-256 of every admitted encoded record. Digest work is queued only
1400    /// after lane dispatch; rotation waits for the ordered worker after
1401    /// admission freezes, and the background fold waits for the commit
1402    /// watermark to reach the admitted record count.
1403    pub async fn seal_digest(&mut self) -> String {
1404        self.admitted.digest().await
1405    }
1406
1407    /// Full-object CRC32C of every admitted encoded record. Admission freezes
1408    /// before the fold CAS, so this names the same exact byte range as
1409    /// [`Self::seal_digest`].
1410    pub fn seal_crc32c(&self) -> u32 {
1411        self.admitted.crc32c()
1412    }
1413
1414    pub async fn enqueue_data_window(
1415        &mut self,
1416        records: Vec<RecordFrame>,
1417        on_attempted: AttemptedBytes,
1418    ) -> Result<CommitRange, ProtocolError> {
1419        if records.is_empty() {
1420            return Ok(CommitRange {
1421                first_offset: self.admitted_len(),
1422                end_offset: self.admitted_len(),
1423                updates: self.commits.subscribe(),
1424            });
1425        }
1426        if self.sealed {
1427            return Err(ProtocolError::Finalized);
1428        }
1429        if self.is_poisoned() {
1430            return Err(ProtocolError::Poisoned);
1431        }
1432        let chunks: Result<Vec<_>, _> = records.iter().map(RecordFrame::encode).collect();
1433        let chunks: Arc<[Bytes]> = chunks?.into();
1434        drop(records);
1435        let quorum = self.quorum();
1436        let active_lanes = self.lanes.iter().filter(|lane| lane.is_some()).count();
1437        if active_lanes < quorum {
1438            tracing::warn!(active_lanes, required_lanes = quorum, "WAL writer poisoned");
1439            self.commits.poison();
1440            return Err(ProtocolError::NoQuorum);
1441        }
1442        self.metrics.batches_sent.increment();
1443        let batch_bytes = chunks.iter().map(Bytes::len).sum::<usize>();
1444        let mut reservations = Vec::with_capacity(self.lanes.len());
1445        let mut retired_lanes = Vec::new();
1446        for zone in 0..self.lanes.len() {
1447            let reservation = self.lanes[zone]
1448                .as_ref()
1449                .and_then(|lane| lane.budget.try_reserve(batch_bytes));
1450            if reservation.is_none() && self.lanes[zone].is_some() {
1451                let lane = self.lanes[zone].take().expect("lane checked present");
1452                lane.done.abort();
1453                retired_lanes.push(lane.done);
1454                self.commits.finish_lane(zone, None);
1455                self.metrics.lane_capacity_drops.increment();
1456                tracing::warn!(
1457                    zone,
1458                    batch_bytes,
1459                    "replica lane exceeded its retained-byte budget"
1460                );
1461            }
1462            reservations.push(reservation);
1463        }
1464        let _ = join_all(retired_lanes).await;
1465        if reservations.iter().flatten().count() < quorum {
1466            self.commits.poison();
1467            return Err(ProtocolError::NoQuorum);
1468        }
1469        let start = self.admitted.bytes_len() as i64;
1470        let mut boundaries = Vec::with_capacity(chunks.len());
1471        let mut acc = start;
1472        for chunk in chunks.iter() {
1473            acc += chunk.len() as i64;
1474            boundaries.push(acc);
1475        }
1476        let boundaries: Arc<[i64]> = boundaries.into();
1477        self.admitted.extend_metadata(&chunks);
1478        let batch = Arc::new(BatchDescriptor {
1479            start,
1480            chunks: Arc::clone(&chunks),
1481            boundaries,
1482            on_attempted: Arc::clone(&on_attempted),
1483            bytes: batch_bytes,
1484            pending_lanes: AtomicUsize::new(reservations.iter().flatten().count()),
1485            packed_groups: std::sync::Mutex::new(BTreeMap::new()),
1486        });
1487        let mut represented_zones = Vec::with_capacity(self.lanes.len());
1488        for (zone, reservation) in reservations.iter_mut().enumerate() {
1489            if let (Some(lane), Some(reservation)) = (&self.lanes[zone], reservation.take()) {
1490                let lane_batch = LaneBatch::new(Arc::clone(&batch), reservation);
1491                if let Err(error) = lane.work.send(lane_batch) {
1492                    drop(error.0);
1493                    // the lane task already terminated; its death notice
1494                    // reached earlier batches, and this batch simply never
1495                    // gains this zone's support
1496                    let lane = self.lanes[zone]
1497                        .take()
1498                        .expect("failed lane send still owns its task handle");
1499                    let _ = lane.done.await;
1500                    self.commits.finish_lane(zone, None);
1501                } else {
1502                    represented_zones.push(zone);
1503                }
1504            }
1505        }
1506        // SHA-256 is ordered behind every preceding batch but starts only
1507        // after this batch is visible to all live lanes. Rotation is the sole
1508        // consumer that waits for the digest pipeline.
1509        self.admitted.queue_digest(chunks);
1510        Ok(self
1511            .commits
1512            .admit_window(&batch.boundaries, &represented_zones))
1513    }
1514
1515    /// Finalize the rotated segment at exactly the committed bytes.
1516    ///
1517    /// Enforcement only: the caller has already committed the rotation view
1518    /// (advanced `tail_base` with the canonical digest) through the manifest
1519    /// quorum. The fast path finalizes this writer's own lanes; if fewer than
1520    /// a quorum survive, the fallback reconstructs the decided prefix from
1521    /// storage, verifies its digest, and installs it through
1522    /// [`QuorumVolume::enforce_seal`].
1523    ///
1524    /// This method is one-shot even though the fallback operation is
1525    /// idempotent: sealing takes and closes the live lane set before it can
1526    /// fail. A caller that needs retries must reconstruct from storage through
1527    /// the committed range and digest rather than calling `seal` again.
1528    pub async fn seal(&mut self) -> Result<SealReport, ProtocolError> {
1529        self.sealed = true;
1530        let total = self.admitted.len();
1531        if self.admitted.is_empty() || self.is_poisoned() {
1532            self.shutdown_background_tasks().await;
1533            return Err(ProtocolError::Poisoned);
1534        }
1535        let expected_digest = self.seal_digest().await;
1536        let expected_crc32c = self.seal_crc32c();
1537        let quorum = self.quorum();
1538        // Close every lane's work channel: each task flushes its in-flight
1539        // acknowledgments and yields the token at its durable tail. The abort
1540        // guard stops stragglers on cancellation or early return; the normal
1541        // exit below drains every lane join handle before returning.
1542        let lanes = std::mem::take(&mut self.lanes);
1543        let aborts = AbortLanesOnDrop(
1544            lanes
1545                .iter()
1546                .flatten()
1547                .map(|lane| lane.done.abort_handle())
1548                .collect(),
1549        );
1550        let mut draining: FuturesUnordered<_> = lanes
1551            .into_iter()
1552            .enumerate()
1553            .filter_map(|(zone, lane)| {
1554                lane.map(|lane| {
1555                    drop(lane.work);
1556                    let done = lane.done;
1557                    async move { done.await.ok().flatten().map(|token| (zone, token)) }
1558                })
1559            })
1560            .collect();
1561        let result = async {
1562            // Settle against the same segment-wide watermark used by normal
1563            // completions. Poison wins even if the durable watermark later moves:
1564            // sealing may not acknowledge across an indeterminate gap.
1565            let mut settling = self.commits.subscribe();
1566            loop {
1567                let snapshot = settling.borrow_and_update().clone();
1568                if snapshot.failure.is_some() {
1569                    return Err(ProtocolError::Poisoned);
1570                }
1571                if snapshot.committed >= total {
1572                    break;
1573                }
1574                settling
1575                    .changed()
1576                    .await
1577                    .map_err(|_| ProtocolError::PipelineClosed)?;
1578            }
1579            // The watermark above proves quorum durability for every admitted
1580            // record, but a per-record quorum may be stitched from different
1581            // zone pairs: finalization needs lanes whose own copy holds every
1582            // byte. Wait for a quorum of full drains (or for every lane to
1583            // settle), give stragglers one bounded grace, then abandon them to
1584            // targeted repair — a lagging lane may legally retain tens of MiB
1585            // of unacknowledged backlog, and the seal must not wait out a
1586            // drain the committed boundary does not need. An abandoned copy is
1587            // a prefix of the canonical bytes, which recovery and repair
1588            // already tolerate.
1589            let mut drained = Vec::new();
1590            while drained.len() < quorum {
1591                match draining.next().await {
1592                    Some(Some((zone, token))) => drained.push((zone, token)),
1593                    Some(None) => {}
1594                    None => break,
1595                }
1596            }
1597            if !draining.is_empty() {
1598                let grace = tokio::time::sleep(LANE_DRAIN_GRACE);
1599                tokio::pin!(grace);
1600                loop {
1601                    tokio::select! {
1602                        biased;
1603                        next = draining.next() => match next {
1604                            Some(Some((zone, token))) => drained.push((zone, token)),
1605                            Some(None) => {}
1606                            None => break,
1607                        },
1608                        () = &mut grace => break,
1609                    }
1610                }
1611            }
1612            let write_offset = self.physical_size() as i64;
1613            let finalizations = join_all(drained.into_iter().map(|(zone, mut token)| {
1614                let replica = Arc::clone(&self.replicas[zone]);
1615                let config = self.config.clone();
1616                async move {
1617                    finalize_with_retry(&replica, &mut token, write_offset, &config)
1618                        .await
1619                        .map(|snapshot| (zone, snapshot))
1620                }
1621            }))
1622            .await;
1623            let mut finalized = vec![None; self.replicas.len()];
1624            for result in finalizations {
1625                match result {
1626                    Ok((zone, snapshot)) => {
1627                        finalized[zone] = Some(snapshot);
1628                    }
1629                    Err(error) => {
1630                        tracing::warn!(
1631                            zone = error.zone,
1632                            code = ?error.code,
1633                            %error,
1634                            write_offset,
1635                            "live segment finalization failed"
1636                        );
1637                    }
1638                }
1639            }
1640            if finalized.iter().flatten().count() >= self.quorum() {
1641                return Ok(SealReport { finalized });
1642            }
1643            let volume = QuorumVolume::with_metadata(
1644                self.replicas.clone(),
1645                self.config.clone(),
1646                self.metadata.clone(),
1647                Arc::clone(&self.metrics),
1648            )?;
1649            let recovered = volume.recover_for_seal(Some(total)).await?;
1650            let actual_digest = recovered.digest();
1651            if actual_digest != expected_digest {
1652                return Err(ProtocolError::SealDigestMismatch {
1653                    expected: expected_digest,
1654                    actual: actual_digest,
1655                });
1656            }
1657            let actual_crc32c = recovered.crc32c();
1658            if actual_crc32c != expected_crc32c {
1659                return Err(ProtocolError::SealCrc32cMismatch {
1660                    expected: expected_crc32c,
1661                    actual: actual_crc32c,
1662                });
1663            }
1664            volume.enforce_seal(recovered.canonical()).await?;
1665            // Enforcement guarantees a sealed quorum, but does not promise that
1666            // every replica was reachable. Conservatively request targeted repair.
1667            Ok(SealReport::default())
1668        }
1669        .await;
1670        aborts.abort();
1671        while draining.next().await.is_some() {}
1672        self.shutdown_background_tasks().await;
1673        result
1674    }
1675
1676    /// Poison this writer so no later completion can be acknowledged.
1677    #[cfg(test)]
1678    pub fn poison(&self) {
1679        self.commits.poison();
1680    }
1681}
1682
1683fn prefer_lower_zone(current: &mut Option<TransportError>, candidate: TransportError) {
1684    if current
1685        .as_ref()
1686        .is_none_or(|existing| candidate.zone < existing.zone)
1687    {
1688        *current = Some(candidate);
1689    }
1690}
1691
1692async fn snapshot_with_retry(
1693    replica: &Arc<dyn Replica>,
1694    config: &ClientConfig,
1695) -> Result<ReplicaSnapshot, TransportError> {
1696    let mut attempt = 0usize;
1697    loop {
1698        match replica.snapshot().await {
1699            Ok(snapshot) => return Ok(snapshot),
1700            Err(error) if error.code.transient() && attempt < config.max_retries => {
1701                retry_sleep(config, attempt).await;
1702                attempt += 1;
1703            }
1704            Err(error) => return Err(error),
1705        }
1706    }
1707}
1708
1709async fn stat_with_retry(
1710    replica: &Arc<dyn Replica>,
1711    config: &ClientConfig,
1712) -> Result<ReplicaSnapshot, TransportError> {
1713    let mut attempt = 0usize;
1714    loop {
1715        match replica.stat().await {
1716            Ok(snapshot) => return Ok(snapshot),
1717            Err(error) if error.code.transient() && attempt < config.max_retries => {
1718                retry_sleep(config, attempt).await;
1719                attempt += 1;
1720            }
1721            Err(error) => return Err(error),
1722        }
1723    }
1724}
1725
1726async fn create_with_retry(
1727    replica: &Arc<dyn Replica>,
1728    metadata: HashMap<String, String>,
1729    config: &ClientConfig,
1730) -> Result<ReplicaSnapshot, TransportError> {
1731    let mut attempt = 0usize;
1732    loop {
1733        match replica.create_appendable(metadata.clone()).await {
1734            Ok(snapshot) => return Ok(snapshot),
1735            Err(error) if error.code.transient() && attempt < config.max_retries => {
1736                retry_sleep(config, attempt).await;
1737                attempt += 1;
1738            }
1739            Err(error) => return Err(error),
1740        }
1741    }
1742}
1743
1744async fn create_session_with_retry(
1745    replica: &Arc<dyn Replica>,
1746    metadata: HashMap<String, String>,
1747    config: &ClientConfig,
1748) -> Result<AppendToken, TransportError> {
1749    let mut attempt = 0usize;
1750    loop {
1751        match replica.create_append_session(metadata.clone()).await {
1752            Ok(token) => return Ok(token),
1753            Err(error) if error.code.transient() && attempt < config.max_retries => {
1754                retry_sleep(config, attempt).await;
1755                attempt += 1;
1756            }
1757            Err(error) => return Err(error),
1758        }
1759    }
1760}
1761
1762async fn takeover_with_retry(
1763    replica: &Arc<dyn Replica>,
1764    observed: &ReplicaSnapshot,
1765    config: &ClientConfig,
1766) -> Result<AppendToken, TransportError> {
1767    let mut attempt = 0usize;
1768    loop {
1769        match replica.takeover(observed).await {
1770            Ok(token) => return Ok(token),
1771            Err(error) if error.code.transient() && attempt < config.max_retries => {
1772                retry_sleep(config, attempt).await;
1773                attempt += 1;
1774            }
1775            Err(error) => return Err(error),
1776        }
1777    }
1778}
1779
1780async fn takeover_current_with_retry(
1781    replica: &Arc<dyn Replica>,
1782    config: &ClientConfig,
1783) -> Result<AppendToken, TransportError> {
1784    let mut attempt = 0usize;
1785    loop {
1786        match replica.takeover_current().await {
1787            Ok(token) => return Ok(token),
1788            Err(error) if error.code.transient() && attempt < config.max_retries => {
1789                retry_sleep(config, attempt).await;
1790                attempt += 1;
1791            }
1792            Err(error) => return Err(error),
1793        }
1794    }
1795}
1796
1797async fn replace_with_retry(
1798    replica: &Arc<dyn Replica>,
1799    mut observed: ReplicaSnapshot,
1800    data: Bytes,
1801    metadata: HashMap<String, String>,
1802    config: &ClientConfig,
1803) -> Result<AppendToken, TransportError> {
1804    let mut attempt = 0usize;
1805    loop {
1806        match replica
1807            .replace_appendable(&observed, data.clone(), metadata.clone())
1808            .await
1809        {
1810            Ok(token) => return Ok(token),
1811            Err(error) if error.code.transient() && attempt < config.max_retries => {
1812                retry_sleep(config, attempt).await;
1813                attempt += 1;
1814                observed = snapshot_with_retry(replica, config).await?;
1815                if observed.bytes == data[..]
1816                    && observed.metadata == metadata
1817                    && !observed.finalized
1818                {
1819                    return Ok(AppendToken {
1820                        zone: observed.zone,
1821                        generation: Some(observed.generation),
1822                        metageneration: Some(observed.metageneration),
1823                        persisted_size: data.len() as i64,
1824                        write_handle: None,
1825                    });
1826                }
1827            }
1828            Err(error) => return Err(error),
1829        }
1830    }
1831}
1832
1833/// Immutable data shared by every lane representing one admitted window.
1834struct BatchDescriptor {
1835    start: i64,
1836    chunks: Arc<[Bytes]>,
1837    boundaries: Arc<[i64]>,
1838    on_attempted: AttemptedBytes,
1839    bytes: usize,
1840    pending_lanes: AtomicUsize,
1841    packed_groups: Mutex<BTreeMap<i64, Arc<OnceLock<Arc<PackedAppend>>>>>,
1842}
1843
1844impl BatchDescriptor {
1845    fn end(&self) -> i64 {
1846        self.boundaries
1847            .last()
1848            .copied()
1849            .expect("an admitted batch is non-empty")
1850    }
1851
1852    fn lane_staged(&self) {
1853        let previous = self.pending_lanes.fetch_sub(1, Ordering::AcqRel);
1854        debug_assert!(previous > 0, "batch lane count underflow");
1855        if previous == 1 {
1856            self.packed_groups
1857                .lock()
1858                .unwrap_or_else(|poisoned| poisoned.into_inner())
1859                .clear();
1860        }
1861    }
1862}
1863
1864/// Lane-local ownership for one shared batch descriptor.
1865struct LaneBatch {
1866    batch: Arc<BatchDescriptor>,
1867    reservation: Option<Arc<LaneReservation>>,
1868    staged: bool,
1869}
1870
1871impl LaneBatch {
1872    fn new(batch: Arc<BatchDescriptor>, reservation: Arc<LaneReservation>) -> Self {
1873        Self {
1874            batch,
1875            reservation: Some(reservation),
1876            staged: false,
1877        }
1878    }
1879
1880    fn into_retained(mut self) -> RetainedBatch {
1881        let retained = RetainedBatch {
1882            batch: Arc::clone(&self.batch),
1883            next_chunk: 0,
1884            _reservation: self
1885                .reservation
1886                .take()
1887                .expect("an unstaged lane batch owns its reservation"),
1888        };
1889        self.staged = true;
1890        self.batch.lane_staged();
1891        retained
1892    }
1893}
1894
1895impl Drop for LaneBatch {
1896    fn drop(&mut self) {
1897        if !self.staged {
1898            self.batch.lane_staged();
1899        }
1900    }
1901}
1902
1903struct RetainedBatch {
1904    batch: Arc<BatchDescriptor>,
1905    next_chunk: usize,
1906    _reservation: Arc<LaneReservation>,
1907}
1908
1909fn packed_group(batches: &[LaneBatch]) -> Arc<PackedAppend> {
1910    let first = &batches.first().expect("coalesced group is non-empty").batch;
1911    let end = batches
1912        .last()
1913        .expect("coalesced group is non-empty")
1914        .batch
1915        .end();
1916    let cell = {
1917        let mut groups = first
1918            .packed_groups
1919            .lock()
1920            .unwrap_or_else(|poisoned| poisoned.into_inner());
1921        Arc::clone(
1922            groups
1923                .entry(end)
1924                .or_insert_with(|| Arc::new(OnceLock::new())),
1925        )
1926    };
1927    Arc::clone(cell.get_or_init(|| {
1928        let chunks = batches
1929            .iter()
1930            .flat_map(|batch| batch.batch.chunks.iter().cloned())
1931            .collect();
1932        Arc::new(pack_append(chunks))
1933    }))
1934}
1935
1936/// How long a seal waits for straggler lane drains once a quorum of lanes
1937/// has fully drained and the commit watermark has settled. In the healthy
1938/// case the last lane finishes within this window and gets finalized; a
1939/// deeply backlogged laggard is cut and left to targeted repair.
1940const LANE_DRAIN_GRACE: std::time::Duration = std::time::Duration::from_millis(100);
1941
1942/// Aborts the lane tasks a seal took ownership of when the seal exits by
1943/// any path. A finished task ignores the abort; an abandoned straggler must
1944/// stop appending so the next generation owns the object alone.
1945struct AbortLanesOnDrop(Vec<tokio::task::AbortHandle>);
1946
1947impl AbortLanesOnDrop {
1948    fn abort(&self) {
1949        for lane in &self.0 {
1950            lane.abort();
1951        }
1952    }
1953}
1954
1955impl Drop for AbortLanesOnDrop {
1956    fn drop(&mut self) {
1957        self.abort();
1958    }
1959}
1960
1961/// The per-lane writer receives batches in offset order, snapshots and drains
1962/// everything currently queued, and sends the resulting group with a flush on
1963/// its final data message without waiting for earlier flush acknowledgments.
1964/// Durable-tail movement publishes one monotonic byte watermark. On a session
1965/// disturbance it resumes via the session handle and resends the retained
1966/// unacknowledged suffix; a fence or exhausted retries publishes one terminal
1967/// lane outcome. A lane that makes no durable progress for its configured
1968/// timeout is retired without recovery: changing sessions cannot prove that a
1969/// completely stationary durable tail is healthy, and the commit tracker must
1970/// promptly decide whether the remaining lanes still form a true quorum.
1971#[derive(Debug)]
1972enum LaneDeath {
1973    Stalled,
1974    Transport(TransportError),
1975}
1976
1977impl LaneDeath {
1978    fn stalled(metrics: &Metrics) -> Self {
1979        metrics.lane_timeouts.increment();
1980        Self::Stalled
1981    }
1982}
1983
1984struct LaneRuntime {
1985    replica: Arc<dyn Replica>,
1986    token: AppendToken,
1987    config: ClientConfig,
1988    metrics: Arc<Metrics>,
1989    commits: Arc<CommitTracker>,
1990    stall_timeout: Arc<LaneStallTimeout>,
1991    durable: i64,
1992    retained: VecDeque<RetainedBatch>,
1993    attempted: Option<AttemptedBytes>,
1994    monitor_session: bool,
1995    last_progress: tokio::time::Instant,
1996}
1997
1998impl LaneRuntime {
1999    fn new(
2000        replica: Arc<dyn Replica>,
2001        token: AppendToken,
2002        config: ClientConfig,
2003        metrics: Arc<Metrics>,
2004        commits: Arc<CommitTracker>,
2005        stall_timeout: Arc<LaneStallTimeout>,
2006    ) -> Self {
2007        let durable = token.persisted_size;
2008        Self {
2009            replica,
2010            token,
2011            config,
2012            metrics,
2013            commits,
2014            stall_timeout,
2015            durable,
2016            retained: VecDeque::new(),
2017            attempted: None,
2018            // Keep observing an idle live stream so a trailing fence cannot
2019            // disappear merely because its persisted-size response drained
2020            // the retained suffix.
2021            monitor_session: true,
2022            last_progress: tokio::time::Instant::now(),
2023        }
2024    }
2025
2026    fn zone(&self) -> usize {
2027        self.token.zone
2028    }
2029
2030    fn stall_deadline(&self) -> tokio::time::Instant {
2031        self.last_progress + self.stall_timeout.get()
2032    }
2033
2034    fn publish_advance(&mut self, change: LaneDurableChange) -> Result<bool, LaneDeath> {
2035        publish_lane_advance(
2036            change,
2037            self.zone(),
2038            &mut self.durable,
2039            &mut self.last_progress,
2040            &mut self.retained,
2041            &self.commits,
2042        )
2043    }
2044
2045    /// Resolve an elapsed lane deadline against the durable-tail stream.
2046    ///
2047    /// A timeout is not itself proof of failure: the final observation may
2048    /// publish progress or a stream error. Return whether the caller still
2049    /// needs session recovery after applying that observation.
2050    async fn confirm_timeout(&mut self) -> Result<bool, LaneDeath> {
2051        let change = confirm_lane_stall(
2052            &self.replica,
2053            self.durable,
2054            self.stall_deadline(),
2055            &self.metrics,
2056        )
2057        .await?;
2058        let stream_failed = self.publish_advance(change)?;
2059        Ok(stream_failed || !self.retained.is_empty())
2060    }
2061
2062    async fn stage(&mut self, batches: Vec<LaneBatch>) -> Result<bool, LaneDeath> {
2063        match tokio::time::timeout_at(
2064            self.stall_deadline(),
2065            stage_group(
2066                &self.replica,
2067                batches,
2068                &mut self.attempted,
2069                &mut self.retained,
2070            ),
2071        )
2072        .await
2073        {
2074            Ok(failed) => Ok(failed),
2075            Err(_) => self.confirm_timeout().await,
2076        }
2077    }
2078}
2079
2080async fn run_lane(
2081    replica: Arc<dyn Replica>,
2082    token: AppendToken,
2083    config: ClientConfig,
2084    metrics: Arc<Metrics>,
2085    mut work: tokio::sync::mpsc::UnboundedReceiver<LaneBatch>,
2086    commits: Arc<CommitTracker>,
2087    stall_timeout: Arc<LaneStallTimeout>,
2088) -> Option<AppendToken> {
2089    let mut lane = LaneRuntime::new(replica, token, config, metrics, commits, stall_timeout);
2090    let mut closed = false;
2091    let death: Option<LaneDeath> = loop {
2092        if closed && lane.retained.is_empty() {
2093            break None;
2094        }
2095        tokio::select! {
2096            biased;
2097            // Progress must win ties with queued work. The writer reserves a
2098            // lane's retained bytes before dispatch; if a ready durable-tail
2099            // update sits behind an always-ready work queue, a healthy lane
2100            // looks stalled and is falsely dropped at its byte budget. `biased`
2101            // keeps DST deterministic, while the work arm still snapshots and
2102            // drains the ready queue so sustained producers cannot postpone a
2103            // flush indefinitely.
2104            changed = lane_progress(
2105                &lane.replica,
2106                lane.durable,
2107                !lane.retained.is_empty(),
2108                lane.stall_deadline(),
2109                &lane.metrics,
2110            ), if lane.monitor_session || !lane.retained.is_empty() => {
2111                match changed {
2112                    LaneProgress::Advanced(change) => {
2113                        match lane.publish_advance(change) {
2114                            Ok(false) => lane.monitor_session = true,
2115                            Ok(true) => {
2116                                if let Err(error) = lane.recover().await {
2117                                    break Some(error);
2118                                }
2119                                lane.monitor_session = true;
2120                            }
2121                            Err(error) => break Some(error),
2122                        }
2123                    }
2124                    LaneProgress::Stalled => {
2125                        break Some(LaneDeath::Stalled);
2126                    }
2127                    LaneProgress::Failed(error) if !error.code.transient() => {
2128                        break Some(LaneDeath::Transport(error));
2129                    }
2130                    LaneProgress::Failed(_) => {
2131                        lane.monitor_session = false;
2132                        if lane.retained.is_empty() {
2133                            continue;
2134                        }
2135                        if let Err(error) = lane.recover().await {
2136                            break Some(error);
2137                        }
2138                        lane.monitor_session = true;
2139                    }
2140                }
2141            }
2142            batch = work.recv(), if !closed => match batch {
2143                Some(batch) => {
2144                    // Snapshot and drain everything already queued, then flush
2145                    // on the final data message. Work arriving after the
2146                    // snapshot forms the next group, so sustained producers
2147                    // cannot postpone this flush indefinitely.
2148                    let mut batches = vec![batch];
2149                    let queued = work.len();
2150                    for _ in 0..queued {
2151                        match work.try_recv() {
2152                            Ok(batch) => batches.push(batch),
2153                            Err(tokio::sync::mpsc::error::TryRecvError::Empty) => break,
2154                            Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => {
2155                                closed = true;
2156                                break;
2157                            }
2158                        }
2159                    }
2160                    if lane.retained.is_empty() {
2161                        lane.last_progress = tokio::time::Instant::now();
2162                    }
2163                    let failed = match lane.stage(batches).await {
2164                        Ok(failed) => failed,
2165                        Err(error) => break Some(error),
2166                    };
2167                    if failed {
2168                        if let Err(error) = lane.recover().await {
2169                            break Some(error);
2170                        }
2171                    }
2172                    lane.monitor_session = true;
2173                }
2174                None => closed = true,
2175            },
2176        }
2177    };
2178    match death {
2179        None => {
2180            lane.commits.finish_lane(lane.zone(), None);
2181            lane.token.persisted_size = lane.durable;
2182            Some(lane.token)
2183        }
2184        Some(LaneDeath::Stalled) => {
2185            tracing::warn!(
2186                zone = lane.zone(),
2187                durable_offset = lane.durable,
2188                retained_chunks = retained_chunk_count(&lane.retained),
2189                stall_timeout_ms = lane.stall_timeout.get().as_millis(),
2190                "append lane shed after making no durable progress"
2191            );
2192            lane.commits.finish_lane(lane.zone(), None);
2193            lane.retained.clear();
2194            work.close();
2195            while work.try_recv().is_ok() {}
2196            None
2197        }
2198        Some(LaneDeath::Transport(error)) => {
2199            tracing::warn!(
2200                zone = error.zone,
2201                code = ?error.code,
2202                error = %error,
2203                durable_offset = lane.durable,
2204                retained_chunks = retained_chunk_count(&lane.retained),
2205                "append lane died"
2206            );
2207            tracing::debug!(
2208                zone = error.zone,
2209                code = ?error.code,
2210                message = error.message.as_str(),
2211                durable_offset = lane.durable,
2212                retained_chunks = retained_chunk_count(&lane.retained),
2213                "append lane dropped"
2214            );
2215            lane.commits.finish_lane(lane.zone(), Some(error));
2216            lane.retained.clear();
2217            work.close();
2218            while work.try_recv().is_ok() {}
2219            None
2220        }
2221    }
2222}
2223
2224/// Stage one coalesced group onto the session: record every batch's chunks for
2225/// acknowledgment matching and resend, then send all group chunks together so
2226/// the transport flushes on the final data message. Returns whether the send
2227/// failed so the caller can run lane recovery, which also flushes its resend.
2228async fn stage_group(
2229    replica: &Arc<dyn Replica>,
2230    batches: Vec<LaneBatch>,
2231    attempted: &mut Option<AttemptedBytes>,
2232    retained: &mut VecDeque<RetainedBatch>,
2233) -> bool {
2234    let group_start = batches
2235        .first()
2236        .expect("coalesced group is non-empty")
2237        .batch
2238        .start;
2239    let packed = packed_group(&batches);
2240    for batch in batches {
2241        (batch.batch.on_attempted)(batch.batch.bytes as u64);
2242        *attempted = Some(Arc::clone(&batch.batch.on_attempted));
2243        retained.push_back(batch.into_retained());
2244    }
2245    replica
2246        .lane_send_packed(group_start, &packed)
2247        .await
2248        .is_err()
2249}
2250
2251/// Advance each retained batch's byte cursor and release fully durable batches.
2252/// `partition_point` avoids per-record notification work when one persisted
2253/// offset covers a large coalesced group.
2254fn ack_through(durable: i64, retained: &mut VecDeque<RetainedBatch>) {
2255    while let Some(batch) = retained.front_mut() {
2256        batch.next_chunk = batch
2257            .batch
2258            .boundaries
2259            .partition_point(|boundary| *boundary <= durable)
2260            .max(batch.next_chunk);
2261        if batch.next_chunk == batch.batch.chunks.len() {
2262            retained.pop_front();
2263        } else {
2264            break;
2265        }
2266    }
2267}
2268
2269fn retained_chunk_count(retained: &VecDeque<RetainedBatch>) -> usize {
2270    retained
2271        .iter()
2272        .map(|batch| batch.batch.chunks.len() - batch.next_chunk)
2273        .sum()
2274}
2275
2276/// Wait for the durable tail to move past `seen` or for the live session to
2277/// fail. Idle sessions have no stall deadline, while retained writes must make
2278/// durable progress before their configured deadline.
2279enum LaneProgress {
2280    Advanced(LaneDurableChange),
2281    Stalled,
2282    Failed(TransportError),
2283}
2284
2285async fn lane_progress(
2286    replica: &Arc<dyn Replica>,
2287    seen: i64,
2288    active: bool,
2289    stall_deadline: tokio::time::Instant,
2290    metrics: &Metrics,
2291) -> LaneProgress {
2292    loop {
2293        if !active {
2294            match replica.lane_durable_change(seen).await {
2295                Ok(change) if change.persisted_size > seen => {
2296                    return LaneProgress::Advanced(change);
2297                }
2298                Ok(_) => tokio::task::yield_now().await,
2299                Err(error) => return LaneProgress::Failed(error),
2300            }
2301            continue;
2302        }
2303        tokio::select! {
2304            biased;
2305            result = replica.lane_durable_change(seen) => match result {
2306                Ok(change) if change.persisted_size > seen => {
2307                    return LaneProgress::Advanced(change);
2308                }
2309                Ok(_) if tokio::time::Instant::now() < stall_deadline => {
2310                    tokio::task::yield_now().await;
2311                }
2312                Ok(_) => {
2313                    let _ = LaneDeath::stalled(metrics);
2314                    return LaneProgress::Stalled;
2315                }
2316                Err(error) => return LaneProgress::Failed(error),
2317            },
2318            _ = tokio::time::sleep_until(stall_deadline) => {
2319                let _ = LaneDeath::stalled(metrics);
2320                return LaneProgress::Stalled;
2321            }
2322        }
2323    }
2324}
2325
2326fn publish_lane_progress(
2327    tail: i64,
2328    zone: usize,
2329    durable: &mut i64,
2330    last_progress: &mut tokio::time::Instant,
2331    retained: &mut VecDeque<RetainedBatch>,
2332    commits: &CommitTracker,
2333) {
2334    let previous = *durable;
2335    *durable = (*durable).max(tail);
2336    if *durable > previous {
2337        *last_progress = tokio::time::Instant::now();
2338    }
2339    ack_through(*durable, retained);
2340    commits.publish_durable(zone, *durable);
2341}
2342
2343/// Publish physical progress first, then report whether a transient stream
2344/// failure needs recovery or a terminal error must stop the lane.
2345fn publish_lane_advance(
2346    change: LaneDurableChange,
2347    zone: usize,
2348    durable: &mut i64,
2349    last_progress: &mut tokio::time::Instant,
2350    retained: &mut VecDeque<RetainedBatch>,
2351    commits: &CommitTracker,
2352) -> Result<bool, LaneDeath> {
2353    publish_lane_progress(
2354        change.persisted_size,
2355        zone,
2356        durable,
2357        last_progress,
2358        retained,
2359        commits,
2360    );
2361    match change.error {
2362        None => Ok(false),
2363        Some(error) if error.code.transient() => Ok(true),
2364        Some(error) => Err(LaneDeath::Transport(error)),
2365    }
2366}
2367
2368async fn confirm_lane_stall(
2369    replica: &Arc<dyn Replica>,
2370    seen: i64,
2371    stall_deadline: tokio::time::Instant,
2372    metrics: &Metrics,
2373) -> Result<LaneDurableChange, LaneDeath> {
2374    match lane_progress(replica, seen, true, stall_deadline, metrics).await {
2375        LaneProgress::Advanced(change) => Ok(change),
2376        LaneProgress::Stalled => Err(LaneDeath::Stalled),
2377        LaneProgress::Failed(error) if error.code.transient() => Err(LaneDeath::stalled(metrics)),
2378        LaneProgress::Failed(error) => Err(LaneDeath::Transport(error)),
2379    }
2380}
2381
2382impl LaneRuntime {
2383    /// Re-learn the durable tail through a session resume and resend the
2384    /// retained unacknowledged suffix, slicing a partially durable chunk at
2385    /// the durable boundary. This never reads object bytes: one
2386    /// generation-guarded stream wrote every byte at these offsets, so the
2387    /// offsets identify our data.
2388    async fn recover(&mut self) -> Result<(), LaneDeath> {
2389        let mut attempt = 0usize;
2390        loop {
2391            self.metrics.lane_retries.increment();
2392            let deadline = self.stall_deadline();
2393            let resumed =
2394                tokio::time::timeout_at(deadline, self.replica.resume_tail(&mut self.token)).await;
2395            let resumed = match resumed {
2396                Ok(resumed) => resumed,
2397                Err(_) => {
2398                    if self.confirm_timeout().await? {
2399                        continue;
2400                    }
2401                    return Ok(());
2402                }
2403            };
2404            match resumed {
2405                Ok(tail) => {
2406                    publish_lane_progress(
2407                        tail,
2408                        self.zone(),
2409                        &mut self.durable,
2410                        &mut self.last_progress,
2411                        &mut self.retained,
2412                        &self.commits,
2413                    );
2414                    let mut suffix = Vec::with_capacity(retained_chunk_count(&self.retained));
2415                    for retained_batch in &self.retained {
2416                        for index in retained_batch.next_chunk..retained_batch.batch.chunks.len() {
2417                            let chunk = &retained_batch.batch.chunks[index];
2418                            if suffix.is_empty() {
2419                                let offset = index
2420                                    .checked_sub(1)
2421                                    .and_then(|previous| {
2422                                        retained_batch.batch.boundaries.get(previous).copied()
2423                                    })
2424                                    .unwrap_or(retained_batch.batch.start);
2425                                let skip = usize::try_from((self.durable - offset).max(0))
2426                                    .unwrap_or(chunk.len());
2427                                if skip >= chunk.len() {
2428                                    continue;
2429                                }
2430                                suffix.push(chunk.slice(skip..));
2431                            } else {
2432                                suffix.push(chunk.clone());
2433                            }
2434                        }
2435                    }
2436                    if suffix.is_empty() {
2437                        return Ok(());
2438                    }
2439                    let packed = pack_append(suffix);
2440                    let resend_bytes = packed.len();
2441                    tracing::debug!(
2442                        zone = self.zone(),
2443                        durable_offset = self.durable,
2444                        chunks = packed.chunks().len(),
2445                        bytes = resend_bytes,
2446                        attempt,
2447                        "resending append lane batch after recovery"
2448                    );
2449                    if let Some(attempted) = &self.attempted {
2450                        attempted(resend_bytes as u64);
2451                    }
2452                    let sent = tokio::time::timeout_at(
2453                        self.stall_deadline(),
2454                        self.replica.lane_send_packed(self.durable, &packed),
2455                    )
2456                    .await;
2457                    let sent = match sent {
2458                        Ok(sent) => sent,
2459                        Err(_) => {
2460                            if self.confirm_timeout().await? {
2461                                continue;
2462                            }
2463                            return Ok(());
2464                        }
2465                    };
2466                    match sent {
2467                        Ok(()) => return Ok(()),
2468                        Err(error)
2469                            if error.code.transient() && attempt < self.config.max_retries =>
2470                        {
2471                            if tokio::time::timeout_at(
2472                                self.stall_deadline(),
2473                                retry_sleep(&self.config, attempt),
2474                            )
2475                            .await
2476                            .is_err()
2477                            {
2478                                if self.confirm_timeout().await? {
2479                                    continue;
2480                                }
2481                                return Ok(());
2482                            }
2483                            attempt += 1;
2484                        }
2485                        Err(error) => return Err(LaneDeath::Transport(error)),
2486                    }
2487                }
2488                Err(error) if error.code.transient() && attempt < self.config.max_retries => {
2489                    if tokio::time::timeout_at(
2490                        self.stall_deadline(),
2491                        retry_sleep(&self.config, attempt),
2492                    )
2493                    .await
2494                    .is_err()
2495                    {
2496                        if self.confirm_timeout().await? {
2497                            continue;
2498                        }
2499                        return Ok(());
2500                    }
2501                    attempt += 1;
2502                }
2503                Err(error) => return Err(LaneDeath::Transport(error)),
2504            }
2505        }
2506    }
2507}
2508
2509async fn finalize_with_retry(
2510    replica: &Arc<dyn Replica>,
2511    token: &mut AppendToken,
2512    write_offset: i64,
2513    config: &ClientConfig,
2514) -> Result<ReplicaSnapshot, TransportError> {
2515    let mut attempt = 0usize;
2516    loop {
2517        match replica.finalize(token, write_offset).await {
2518            Ok(snapshot) => return Ok(snapshot),
2519            Err(error)
2520                if error.code == TransportCode::FailedPrecondition || error.code.transient() =>
2521            {
2522                // A retained create stream has no generation identity until a
2523                // successful handle resume proves it. If its finish response
2524                // is ambiguous, canonical seal enforcement must resolve it.
2525                let Some(generation) = token.generation else {
2526                    return Err(error);
2527                };
2528                if error.code.transient() {
2529                    if attempt >= config.max_retries {
2530                        return Err(error);
2531                    }
2532                    retry_sleep(config, attempt).await;
2533                    attempt += 1;
2534                }
2535                // A finalized object's size is authoritative, so an
2536                // ambiguous finish response needs metadata only. Open-object
2537                // metrics remain tail-blind and therefore cannot falsely prove
2538                // that finalization landed.
2539                let snapshot = stat_with_retry(replica, config).await?;
2540                if snapshot.finalized
2541                    && snapshot.generation == generation
2542                    && snapshot.persisted_size == write_offset
2543                {
2544                    return Ok(snapshot);
2545                }
2546                if error.code == TransportCode::FailedPrecondition {
2547                    return Err(error);
2548                }
2549            }
2550            Err(error) => return Err(error),
2551        }
2552    }
2553}
2554
2555pub(crate) fn retry_delay(config: &ClientConfig, attempt: usize) -> Duration {
2556    let multiplier = 1u32.checked_shl(attempt.min(16) as u32).unwrap_or(u32::MAX);
2557    config.retry_base.saturating_mul(multiplier)
2558}
2559
2560pub(crate) async fn retry_sleep(config: &ClientConfig, attempt: usize) {
2561    tokio::time::sleep(retry_delay(config, attempt)).await;
2562}
2563
2564/// Select the longest exact, mutually consistent well-formed prefix visible
2565/// on a quorum of recovery witnesses.
2566///
2567/// A candidate is any quorum-sized subset of the readable snapshots whose
2568/// well-formed prefixes agree pairwise on their overlap; its prefix is the
2569/// longest member's. The witness set must be quorum-sized because a committed
2570/// record is durable on a write quorum, and only a quorum-sized read subset
2571/// is guaranteed to intersect it — a smaller consistent set could miss
2572/// committed records entirely. A record visible on only one selected witness
2573/// is retained and then written back to every witness before sealing; this
2574/// promotion is required when only a quorum of zones is readable, since an
2575/// unavailable zone may hold another copy of a committed record. A divergent
2576/// minority lane is excluded by quorum agreement; different bytes among
2577/// equally long maximal candidates remain ambiguous and fail recovery.
2578pub(crate) fn canonical_prefix(
2579    snapshots: &[ReplicaSnapshot],
2580    quorum: usize,
2581) -> Result<Vec<RecordFrame>, ProtocolError> {
2582    select_canonical_quorum(snapshots, quorum).map(|(prefix, _)| prefix.into_records())
2583}
2584
2585/// Ascending index combinations of size `quorum` out of `count` snapshots,
2586/// in lexicographic order (`count` is at most 5).
2587fn quorum_subsets(count: usize, quorum: usize) -> Vec<Vec<usize>> {
2588    fn extend(
2589        start: usize,
2590        count: usize,
2591        quorum: usize,
2592        current: &mut Vec<usize>,
2593        subsets: &mut Vec<Vec<usize>>,
2594    ) {
2595        if current.len() == quorum {
2596            subsets.push(current.clone());
2597            return;
2598        }
2599        for index in start..count {
2600            current.push(index);
2601            extend(index + 1, count, quorum, current, subsets);
2602            current.pop();
2603        }
2604    }
2605    let mut subsets = Vec::new();
2606    extend(
2607        0,
2608        count,
2609        quorum,
2610        &mut Vec::with_capacity(quorum),
2611        &mut subsets,
2612    );
2613    subsets
2614}
2615
2616fn select_canonical_quorum(
2617    snapshots: &[ReplicaSnapshot],
2618    quorum: usize,
2619) -> Result<(CanonicalPrefix, Vec<ReplicaSnapshot>), ProtocolError> {
2620    if snapshots.len() < quorum {
2621        return Err(ProtocolError::NoQuorum);
2622    }
2623    let decoded: Vec<_> = snapshots
2624        .iter()
2625        .map(CanonicalPrefix::from_snapshot)
2626        .collect();
2627    let mut conflicts = vec![false; decoded.len() * decoded.len()];
2628    let mut first_conflict = None;
2629    for left in 0..decoded.len() {
2630        for right in left + 1..decoded.len() {
2631            let overlap = decoded[left].len().min(decoded[right].len());
2632            let conflict = (0..overlap).find(|index| {
2633                decoded[left].record_bytes(*index) != decoded[right].record_bytes(*index)
2634            });
2635            if let Some(index) = conflict {
2636                first_conflict.get_or_insert(index);
2637                conflicts[left * decoded.len() + right] = true;
2638            }
2639        }
2640    }
2641    let mut candidates = Vec::new();
2642    for subset in quorum_subsets(decoded.len(), quorum) {
2643        let consistent = subset.iter().enumerate().all(|(position, &left)| {
2644            subset[position + 1..]
2645                .iter()
2646                .all(|&right| !conflicts[left * decoded.len() + right])
2647        });
2648        if !consistent {
2649            continue;
2650        }
2651        let mut longest_member = subset[0];
2652        for &member in &subset[1..] {
2653            if decoded[member].len() > decoded[longest_member].len() {
2654                longest_member = member;
2655            }
2656        }
2657        candidates.push((
2658            decoded[longest_member].clone(),
2659            subset
2660                .iter()
2661                .map(|&member| snapshots[member].clone())
2662                .collect::<Vec<_>>(),
2663        ));
2664    }
2665    let longest = candidates
2666        .iter()
2667        .map(|(candidate, _)| candidate.len())
2668        .max()
2669        .ok_or(ProtocolError::ConflictingPrefix {
2670            record_index: first_conflict.unwrap_or(0),
2671        })?;
2672    let mut longest_candidates = candidates
2673        .into_iter()
2674        .filter(|(candidate, _)| candidate.len() == longest);
2675    let first = longest_candidates
2676        .next()
2677        .expect("longest length came from a candidate");
2678    let mut equivalent = vec![first];
2679    for candidate in longest_candidates {
2680        if candidate.0.bytes != equivalent[0].0.bytes {
2681            let record_index = (0..candidate.0.len())
2682                .find(|index| {
2683                    candidate.0.record_bytes(*index) != equivalent[0].0.record_bytes(*index)
2684                })
2685                .unwrap_or(0);
2686            return Err(ProtocolError::ConflictingPrefix { record_index });
2687        }
2688        equivalent.push(candidate);
2689    }
2690    Ok(equivalent
2691        .into_iter()
2692        .max_by(|(_, left_witnesses), (_, right_witnesses)| {
2693            let left_zones: Vec<_> = left_witnesses.iter().map(|copy| copy.zone).collect();
2694            let right_zones: Vec<_> = right_witnesses.iter().map(|copy| copy.zone).collect();
2695            right_zones.cmp(&left_zones)
2696        })
2697        .expect("at least one equivalent candidate remains"))
2698}
2699
2700pub(crate) fn protocol_metadata() -> HashMap<String, String> {
2701    HashMap::from([(META_FORMAT.to_string(), FORMAT_VERSION.to_string())])
2702}
2703
2704pub(crate) fn valid_format(metadata: &HashMap<String, String>) -> bool {
2705    metadata.get(META_FORMAT).map(String::as_str) == Some(FORMAT_VERSION)
2706}
2707
2708#[cfg(test)]
2709fn encode_records(records: &[RecordFrame]) -> Result<Vec<u8>, RecordError> {
2710    let encoded: Result<Vec<_>, _> = records.iter().map(RecordFrame::encode).collect();
2711    Ok(encoded?.into_iter().flatten().collect())
2712}
2713
2714#[cfg(test)]
2715mod tests {
2716    use super::*;
2717    use std::collections::{HashMap, VecDeque};
2718    use std::sync::Arc;
2719
2720    use async_trait::async_trait;
2721    use tokio::sync::{mpsc, oneshot, watch, Mutex};
2722
2723    use crate::metrics::{test_support::TestMetricsRecorder, Metrics};
2724
2725    fn snapshot(zone: usize, records: &[RecordFrame]) -> ReplicaSnapshot {
2726        let bytes = encode_records(records).unwrap();
2727        ReplicaSnapshot {
2728            zone,
2729            generation: 1,
2730            metageneration: 1,
2731            persisted_size: 0,
2732            finalized: false,
2733            crc32c: Some(crc32c::crc32c(&bytes)),
2734            metadata: protocol_metadata(),
2735            bytes,
2736        }
2737    }
2738
2739    fn record(value: &[u8]) -> RecordFrame {
2740        RecordFrame {
2741            payload: Bytes::copy_from_slice(value),
2742        }
2743    }
2744
2745    fn batch_descriptor(
2746        start: i64,
2747        chunks: Vec<Bytes>,
2748        pending_lanes: usize,
2749    ) -> Arc<BatchDescriptor> {
2750        let mut end = start;
2751        let mut boundaries = Vec::with_capacity(chunks.len());
2752        for chunk in &chunks {
2753            end += chunk.len() as i64;
2754            boundaries.push(end);
2755        }
2756        let bytes = chunks.iter().map(Bytes::len).sum();
2757        Arc::new(BatchDescriptor {
2758            start,
2759            chunks: chunks.into(),
2760            boundaries: boundaries.into(),
2761            on_attempted: Arc::new(|_| {}),
2762            bytes,
2763            pending_lanes: AtomicUsize::new(pending_lanes),
2764            packed_groups: std::sync::Mutex::new(BTreeMap::new()),
2765        })
2766    }
2767
2768    struct ScriptedLaneReplica {
2769        zone: usize,
2770        durable: watch::Sender<i64>,
2771        send_releases: Mutex<VecDeque<oneshot::Receiver<()>>>,
2772        sends: mpsc::UnboundedSender<i64>,
2773    }
2774
2775    impl ScriptedLaneReplica {
2776        fn new(
2777            zone: usize,
2778            send_releases: VecDeque<oneshot::Receiver<()>>,
2779            sends: mpsc::UnboundedSender<i64>,
2780        ) -> Self {
2781            let (durable, _) = watch::channel(0);
2782            Self {
2783                zone,
2784                durable,
2785                send_releases: Mutex::new(send_releases),
2786                sends,
2787            }
2788        }
2789
2790        fn error(&self, code: TransportCode, message: &str) -> TransportError {
2791            TransportError {
2792                zone: self.zone,
2793                code,
2794                message: message.into(),
2795            }
2796        }
2797    }
2798
2799    struct ReaderTerminalReplica {
2800        zone: usize,
2801        reader_failed: watch::Sender<bool>,
2802        resume_calls: AtomicUsize,
2803    }
2804
2805    struct StalledReplica {
2806        zone: usize,
2807    }
2808
2809    impl ReaderTerminalReplica {
2810        fn new(zone: usize) -> Self {
2811            let (reader_failed, _) = watch::channel(false);
2812            Self {
2813                zone,
2814                reader_failed,
2815                resume_calls: AtomicUsize::new(0),
2816            }
2817        }
2818
2819        fn error(&self) -> TransportError {
2820            TransportError {
2821                zone: self.zone,
2822                code: TransportCode::PermissionDenied,
2823                message: "async reader rejected append".into(),
2824            }
2825        }
2826
2827        fn resume_calls(&self) -> usize {
2828            self.resume_calls.load(Ordering::SeqCst)
2829        }
2830    }
2831
2832    struct SingleReplicaFactory {
2833        replica: Arc<dyn Replica>,
2834    }
2835
2836    impl SingleReplicaFactory {
2837        fn new(replica: Arc<dyn Replica>) -> Self {
2838            Self { replica }
2839        }
2840    }
2841
2842    #[async_trait]
2843    impl Replica for ScriptedLaneReplica {
2844        async fn snapshot(&self) -> Result<ReplicaSnapshot, TransportError> {
2845            panic!("snapshot is not used in this test")
2846        }
2847
2848        async fn stat(&self) -> Result<ReplicaSnapshot, TransportError> {
2849            panic!("stat is not used in this test")
2850        }
2851
2852        async fn create_appendable(
2853            &self,
2854            _metadata: HashMap<String, String>,
2855        ) -> Result<ReplicaSnapshot, TransportError> {
2856            panic!("create_appendable is not used in this test")
2857        }
2858
2859        async fn create_append_session(
2860            &self,
2861            _metadata: HashMap<String, String>,
2862        ) -> Result<AppendToken, TransportError> {
2863            panic!("create_append_session is not used in this test")
2864        }
2865
2866        async fn create_register(
2867            &self,
2868            _metadata: HashMap<String, String>,
2869        ) -> Result<ReplicaSnapshot, TransportError> {
2870            panic!("create_register is not used in this test")
2871        }
2872
2873        async fn update_register(
2874            &self,
2875            _metageneration: i64,
2876            _metadata: HashMap<String, String>,
2877        ) -> Result<ReplicaSnapshot, TransportError> {
2878            panic!("update_register is not used in this test")
2879        }
2880
2881        async fn resume_tail(&self, _token: &mut AppendToken) -> Result<i64, TransportError> {
2882            Err(self.error(
2883                TransportCode::Internal,
2884                "resume_tail should not run in this test",
2885            ))
2886        }
2887
2888        async fn takeover(
2889            &self,
2890            _observed: &ReplicaSnapshot,
2891        ) -> Result<AppendToken, TransportError> {
2892            panic!("takeover is not used in this test")
2893        }
2894
2895        async fn replace_appendable(
2896            &self,
2897            _observed: &ReplicaSnapshot,
2898            _data: Bytes,
2899            _metadata: HashMap<String, String>,
2900        ) -> Result<AppendToken, TransportError> {
2901            panic!("replace_appendable is not used in this test")
2902        }
2903
2904        async fn append(
2905            &self,
2906            _token: &AppendToken,
2907            _write_offset: i64,
2908            _data: Vec<u8>,
2909        ) -> Result<i64, TransportError> {
2910            panic!("append is not used in this test")
2911        }
2912
2913        async fn lane_send(
2914            &self,
2915            write_offset: i64,
2916            chunks: &[Bytes],
2917        ) -> Result<(), TransportError> {
2918            let end = write_offset + chunks.iter().map(|chunk| chunk.len() as i64).sum::<i64>();
2919            self.durable.send_replace(end);
2920            let _ = self.sends.send(end);
2921            let release = self.send_releases.lock().await.pop_front();
2922            if let Some(release) = release {
2923                let _ = release.await;
2924            }
2925            Ok(())
2926        }
2927
2928        async fn lane_durable_change(
2929            &self,
2930            seen: i64,
2931        ) -> Result<LaneDurableChange, TransportError> {
2932            let mut durable = self.durable.subscribe();
2933            loop {
2934                let current = *durable.borrow_and_update();
2935                if current > seen {
2936                    return Ok(LaneDurableChange {
2937                        persisted_size: current,
2938                        error: None,
2939                    });
2940                }
2941                durable
2942                    .changed()
2943                    .await
2944                    .map_err(|_| self.error(TransportCode::Unavailable, "durable watch closed"))?;
2945            }
2946        }
2947
2948        async fn delete(&self, _generation: i64) -> Result<(), TransportError> {
2949            panic!("delete is not used in this test")
2950        }
2951
2952        async fn finalize(
2953            &self,
2954            _token: &mut AppendToken,
2955            _write_offset: i64,
2956        ) -> Result<ReplicaSnapshot, TransportError> {
2957            panic!("finalize is not used in this test")
2958        }
2959    }
2960
2961    #[async_trait]
2962    impl Replica for ReaderTerminalReplica {
2963        async fn snapshot(&self) -> Result<ReplicaSnapshot, TransportError> {
2964            panic!("snapshot is not used in this test")
2965        }
2966
2967        async fn stat(&self) -> Result<ReplicaSnapshot, TransportError> {
2968            panic!("stat is not used in this test")
2969        }
2970
2971        async fn create_appendable(
2972            &self,
2973            _metadata: HashMap<String, String>,
2974        ) -> Result<ReplicaSnapshot, TransportError> {
2975            panic!("create_appendable is not used in this test")
2976        }
2977
2978        async fn create_append_session(
2979            &self,
2980            _metadata: HashMap<String, String>,
2981        ) -> Result<AppendToken, TransportError> {
2982            Ok(AppendToken {
2983                zone: self.zone,
2984                generation: Some(1),
2985                metageneration: Some(1),
2986                persisted_size: 0,
2987                write_handle: None,
2988            })
2989        }
2990
2991        async fn create_register(
2992            &self,
2993            _metadata: HashMap<String, String>,
2994        ) -> Result<ReplicaSnapshot, TransportError> {
2995            panic!("create_register is not used in this test")
2996        }
2997
2998        async fn update_register(
2999            &self,
3000            _metageneration: i64,
3001            _metadata: HashMap<String, String>,
3002        ) -> Result<ReplicaSnapshot, TransportError> {
3003            panic!("update_register is not used in this test")
3004        }
3005
3006        async fn resume_tail(&self, _token: &mut AppendToken) -> Result<i64, TransportError> {
3007            self.resume_calls.fetch_add(1, Ordering::SeqCst);
3008            Err(TransportError {
3009                zone: self.zone,
3010                code: TransportCode::Unavailable,
3011                message: "recovery must not run after a terminal reader error".into(),
3012            })
3013        }
3014
3015        async fn takeover(
3016            &self,
3017            _observed: &ReplicaSnapshot,
3018        ) -> Result<AppendToken, TransportError> {
3019            panic!("takeover is not used in this test")
3020        }
3021
3022        async fn replace_appendable(
3023            &self,
3024            _observed: &ReplicaSnapshot,
3025            _data: Bytes,
3026            _metadata: HashMap<String, String>,
3027        ) -> Result<AppendToken, TransportError> {
3028            panic!("replace_appendable is not used in this test")
3029        }
3030
3031        async fn append(
3032            &self,
3033            _token: &AppendToken,
3034            _write_offset: i64,
3035            _data: Vec<u8>,
3036        ) -> Result<i64, TransportError> {
3037            panic!("append is not used in this test")
3038        }
3039
3040        async fn lane_send(
3041            &self,
3042            _write_offset: i64,
3043            _chunks: &[Bytes],
3044        ) -> Result<(), TransportError> {
3045            self.reader_failed.send_replace(true);
3046            Ok(())
3047        }
3048
3049        async fn lane_durable_change(
3050            &self,
3051            _seen: i64,
3052        ) -> Result<LaneDurableChange, TransportError> {
3053            let mut reader_failed = self.reader_failed.subscribe();
3054            loop {
3055                if *reader_failed.borrow_and_update() {
3056                    return Err(self.error());
3057                }
3058                reader_failed.changed().await.map_err(|_| TransportError {
3059                    zone: self.zone,
3060                    code: TransportCode::Unavailable,
3061                    message: "reader failure watch closed".into(),
3062                })?;
3063            }
3064        }
3065
3066        async fn delete(&self, _generation: i64) -> Result<(), TransportError> {
3067            panic!("delete is not used in this test")
3068        }
3069
3070        async fn finalize(
3071            &self,
3072            _token: &mut AppendToken,
3073            _write_offset: i64,
3074        ) -> Result<ReplicaSnapshot, TransportError> {
3075            panic!("finalize is not used in this test")
3076        }
3077    }
3078
3079    #[async_trait]
3080    impl Replica for StalledReplica {
3081        async fn snapshot(&self) -> Result<ReplicaSnapshot, TransportError> {
3082            panic!("snapshot is not used in this test")
3083        }
3084
3085        async fn stat(&self) -> Result<ReplicaSnapshot, TransportError> {
3086            panic!("stat is not used in this test")
3087        }
3088
3089        async fn create_appendable(
3090            &self,
3091            _metadata: HashMap<String, String>,
3092        ) -> Result<ReplicaSnapshot, TransportError> {
3093            panic!("create_appendable is not used in this test")
3094        }
3095
3096        async fn create_append_session(
3097            &self,
3098            _metadata: HashMap<String, String>,
3099        ) -> Result<AppendToken, TransportError> {
3100            Ok(AppendToken {
3101                zone: self.zone,
3102                generation: Some(1),
3103                metageneration: Some(1),
3104                persisted_size: 0,
3105                write_handle: None,
3106            })
3107        }
3108
3109        async fn create_register(
3110            &self,
3111            _metadata: HashMap<String, String>,
3112        ) -> Result<ReplicaSnapshot, TransportError> {
3113            panic!("create_register is not used in this test")
3114        }
3115
3116        async fn update_register(
3117            &self,
3118            _metageneration: i64,
3119            _metadata: HashMap<String, String>,
3120        ) -> Result<ReplicaSnapshot, TransportError> {
3121            panic!("update_register is not used in this test")
3122        }
3123
3124        async fn resume_tail(&self, _token: &mut AppendToken) -> Result<i64, TransportError> {
3125            panic!("a no-progress timeout must shed instead of recovering the lane")
3126        }
3127
3128        async fn takeover(
3129            &self,
3130            _observed: &ReplicaSnapshot,
3131        ) -> Result<AppendToken, TransportError> {
3132            panic!("takeover is not used in this test")
3133        }
3134
3135        async fn replace_appendable(
3136            &self,
3137            _observed: &ReplicaSnapshot,
3138            _data: Bytes,
3139            _metadata: HashMap<String, String>,
3140        ) -> Result<AppendToken, TransportError> {
3141            panic!("replace_appendable is not used in this test")
3142        }
3143
3144        async fn append(
3145            &self,
3146            _token: &AppendToken,
3147            _write_offset: i64,
3148            _data: Vec<u8>,
3149        ) -> Result<i64, TransportError> {
3150            panic!("append is not used in this test")
3151        }
3152
3153        async fn lane_send(
3154            &self,
3155            _write_offset: i64,
3156            _chunks: &[Bytes],
3157        ) -> Result<(), TransportError> {
3158            Ok(())
3159        }
3160
3161        async fn lane_durable_change(
3162            &self,
3163            _seen: i64,
3164        ) -> Result<LaneDurableChange, TransportError> {
3165            std::future::pending().await
3166        }
3167
3168        async fn delete(&self, _generation: i64) -> Result<(), TransportError> {
3169            panic!("delete is not used in this test")
3170        }
3171
3172        async fn finalize(
3173            &self,
3174            _token: &mut AppendToken,
3175            _write_offset: i64,
3176        ) -> Result<ReplicaSnapshot, TransportError> {
3177            panic!("finalize is not used in this test")
3178        }
3179    }
3180
3181    #[async_trait]
3182    impl crate::transport::ReplicaFactory for SingleReplicaFactory {
3183        fn bucket_name(&self) -> &str {
3184            "single-replica"
3185        }
3186
3187        fn replica(&self, _object: &str) -> Arc<dyn Replica> {
3188            self.replica.clone()
3189        }
3190
3191        async fn list(
3192            &self,
3193            _prefix: &str,
3194        ) -> Result<Vec<crate::transport::ListedObject>, TransportError> {
3195            Ok(Vec::new())
3196        }
3197    }
3198
3199    #[test]
3200    fn majority_matches_supported_widths() {
3201        assert_eq!(majority(1), 1);
3202        assert_eq!(majority(3), 2);
3203        assert_eq!(majority(5), 3);
3204        let lanes = |durables: &[i64]| {
3205            durables
3206                .iter()
3207                .map(|durable| LaneCommitState {
3208                    durable: *durable,
3209                    ..LaneCommitState::default()
3210                })
3211                .collect::<Vec<_>>()
3212        };
3213        assert_eq!(quorum_durable_watermark(&lanes(&[7]), majority(1)), 7);
3214        assert_eq!(quorum_durable_watermark(&lanes(&[1, 9, 5]), majority(3)), 5);
3215        assert_eq!(
3216            quorum_durable_watermark(&lanes(&[1, 9, 5, 7, 3]), majority(5)),
3217            5
3218        );
3219    }
3220
3221    #[tokio::test]
3222    async fn admitted_prefix_hashes_framed_bytes_in_order_off_path() {
3223        let first = record(b"first").encode().unwrap();
3224        let second = record(b"second").encode().unwrap();
3225        let third = record(b"third").encode().unwrap();
3226        let expected = [first.as_ref(), second.as_ref(), third.as_ref()].concat();
3227
3228        let mut admitted = AdmittedPrefix::default();
3229        let first_batch: Arc<[Bytes]> = vec![first].into();
3230        admitted.extend_metadata(&first_batch);
3231        admitted.queue_digest(first_batch);
3232        let second_batch: Arc<[Bytes]> = vec![second, third].into();
3233        admitted.extend_metadata(&second_batch);
3234        admitted.queue_digest(second_batch);
3235
3236        assert_eq!(admitted.len(), 3);
3237        assert_eq!(admitted.bytes_len(), expected.len());
3238        assert_eq!(admitted.digest().await, digest_bytes(&expected));
3239        assert_eq!(admitted.crc32c(), crc32c::crc32c(&expected));
3240    }
3241
3242    #[test]
3243    fn matching_lane_groups_share_one_packed_wire_payload() {
3244        let first = batch_descriptor(0, vec![Bytes::from_static(b"first")], 2);
3245        let second = batch_descriptor(first.end(), vec![Bytes::from_static(b"second")], 2);
3246        let first_budget = LaneBudget::new();
3247        let second_budget = LaneBudget::new();
3248        let group_one = vec![
3249            LaneBatch::new(
3250                Arc::clone(&first),
3251                first_budget.try_reserve(first.bytes).unwrap(),
3252            ),
3253            LaneBatch::new(
3254                Arc::clone(&second),
3255                first_budget.try_reserve(second.bytes).unwrap(),
3256            ),
3257        ];
3258        let group_two = vec![
3259            LaneBatch::new(
3260                Arc::clone(&first),
3261                second_budget.try_reserve(first.bytes).unwrap(),
3262            ),
3263            LaneBatch::new(
3264                Arc::clone(&second),
3265                second_budget.try_reserve(second.bytes).unwrap(),
3266            ),
3267        ];
3268
3269        let packed_one = packed_group(&group_one);
3270        let packed_two = packed_group(&group_two);
3271
3272        assert!(Arc::ptr_eq(&packed_one, &packed_two));
3273        assert_eq!(
3274            packed_one.chunks(),
3275            &[Bytes::from_static(b"first"), Bytes::from_static(b"second")]
3276        );
3277        drop(group_one);
3278        drop(group_two);
3279        assert!(first
3280            .packed_groups
3281            .lock()
3282            .unwrap_or_else(|poisoned| poisoned.into_inner())
3283            .is_empty());
3284    }
3285
3286    #[test]
3287    fn packed_wire_payload_preserves_offsets_bytes_and_checksums() {
3288        let first = Bytes::from(vec![7; 262_143]);
3289        let second = Bytes::from_static(b"xy");
3290        let packed = pack_append(vec![first.clone(), second.clone()]);
3291        let messages = packed.messages();
3292
3293        assert_eq!(messages.len(), 2);
3294        assert_eq!(messages[0].relative_offset, 0);
3295        assert_eq!(messages[0].content, first);
3296        assert_eq!(messages[0].crc32c, crc32c::crc32c(&messages[0].content));
3297        assert_eq!(messages[1].relative_offset, 262_143);
3298        assert_eq!(messages[1].content, second);
3299        assert_eq!(messages[1].crc32c, crc32c::crc32c(&messages[1].content));
3300    }
3301
3302    #[test]
3303    fn quorum_subsets_enumerate_lexicographically() {
3304        assert_eq!(quorum_subsets(1, 1), vec![vec![0]]);
3305        assert_eq!(
3306            quorum_subsets(3, 2),
3307            vec![vec![0, 1], vec![0, 2], vec![1, 2]]
3308        );
3309        assert_eq!(quorum_subsets(5, 3).len(), 10);
3310        assert_eq!(quorum_subsets(5, 3)[0], vec![0, 1, 2]);
3311    }
3312
3313    #[test]
3314    fn canonical_promotes_a_tail_visible_on_one_recovery_witness() {
3315        let first = record(b"first");
3316        let second = record(b"second");
3317        let snapshots = vec![
3318            snapshot(0, &[first.clone(), second.clone()]),
3319            snapshot(1, std::slice::from_ref(&first)),
3320        ];
3321        assert_eq!(
3322            canonical_prefix(&snapshots, 2).unwrap(),
3323            vec![first, second]
3324        );
3325    }
3326
3327    #[test]
3328    fn canonical_rejects_conflicting_bytes_at_the_same_record() {
3329        let snapshots = vec![snapshot(0, &[record(b"a")]), snapshot(1, &[record(b"b")])];
3330        assert!(matches!(
3331            canonical_prefix(&snapshots, 2),
3332            Err(ProtocolError::ConflictingPrefix { record_index: 0 })
3333        ));
3334    }
3335
3336    #[test]
3337    fn canonical_ignores_one_conflicting_lane_when_two_exact_copies_agree() {
3338        let good = record(b"good");
3339        let snapshots = vec![
3340            snapshot(0, &[record(b"bad")]),
3341            snapshot(1, std::slice::from_ref(&good)),
3342            snapshot(2, std::slice::from_ref(&good)),
3343        ];
3344        assert_eq!(canonical_prefix(&snapshots, 2).unwrap(), vec![good]);
3345    }
3346
3347    #[test]
3348    fn canonical_rejects_equal_length_candidates_without_a_quorum_choice() {
3349        let snapshots = vec![
3350            snapshot(0, &[]),
3351            snapshot(1, &[record(b"left")]),
3352            snapshot(2, &[record(b"right")]),
3353        ];
3354        assert!(matches!(
3355            canonical_prefix(&snapshots, 2),
3356            Err(ProtocolError::ConflictingPrefix { record_index: 0 })
3357        ));
3358    }
3359
3360    #[test]
3361    fn canonical_stops_at_a_partial_tail() {
3362        let first = record(b"first");
3363        let second = record(b"second");
3364        let mut damaged = encode_records(&[first.clone(), second]).unwrap();
3365        damaged.truncate(damaged.len() - 2);
3366        let snapshots = vec![
3367            ReplicaSnapshot {
3368                bytes: damaged,
3369                ..snapshot(0, &[])
3370            },
3371            snapshot(1, std::slice::from_ref(&first)),
3372        ];
3373        assert_eq!(canonical_prefix(&snapshots, 2).unwrap(), vec![first]);
3374    }
3375
3376    #[test]
3377    fn canonical_accepts_a_single_replica_witness() {
3378        let first = record(b"first");
3379        let snapshots = vec![snapshot(0, std::slice::from_ref(&first))];
3380        assert_eq!(canonical_prefix(&snapshots, 1).unwrap(), vec![first]);
3381    }
3382
3383    #[test]
3384    fn canonical_five_zone_quorum_requires_three_consistent_witnesses() {
3385        let good = record(b"good");
3386        let consistent = vec![
3387            snapshot(0, std::slice::from_ref(&good)),
3388            snapshot(1, std::slice::from_ref(&good)),
3389            snapshot(2, &[record(b"divergent")]),
3390            snapshot(3, std::slice::from_ref(&good)),
3391        ];
3392        assert_eq!(
3393            canonical_prefix(&consistent, 3).unwrap(),
3394            vec![good.clone()]
3395        );
3396
3397        // two consistent witnesses out of five are not a read quorum: a
3398        // committed record could live only on the two unreachable zones
3399        // plus the divergent lane's pre-divergence prefix
3400        let insufficient = vec![
3401            snapshot(0, std::slice::from_ref(&good)),
3402            snapshot(1, std::slice::from_ref(&good)),
3403            snapshot(2, &[record(b"divergent")]),
3404        ];
3405        assert!(matches!(
3406            canonical_prefix(&insufficient, 3),
3407            Err(ProtocolError::ConflictingPrefix { .. })
3408        ));
3409    }
3410
3411    #[test]
3412    fn canonical_five_zone_promotes_the_longest_member_of_the_quorum() {
3413        let first = record(b"first");
3414        let second = record(b"second");
3415        let snapshots = vec![
3416            snapshot(0, std::slice::from_ref(&first)),
3417            snapshot(1, &[first.clone(), second.clone()]),
3418            snapshot(2, &[]),
3419        ];
3420        assert_eq!(
3421            canonical_prefix(&snapshots, 3).unwrap(),
3422            vec![first, second]
3423        );
3424    }
3425
3426    #[test]
3427    fn recovery_size_uses_the_available_quorum_intersection_rank() {
3428        let mut all_three = [10, 20, 20];
3429        assert_eq!(select_recovery_size(&mut all_three, 3), Some(20));
3430
3431        let mut two_of_three = [10, 20];
3432        assert_eq!(select_recovery_size(&mut two_of_three, 3), Some(20));
3433
3434        let mut four_of_five = [10, 20, 30, 40];
3435        assert_eq!(select_recovery_size(&mut four_of_five, 5), Some(30));
3436
3437        let mut all_five = [10, 20, 30, 40, 50];
3438        assert_eq!(select_recovery_size(&mut all_five, 5), Some(30));
3439    }
3440
3441    fn active_commit_tracker(lanes: usize) -> Arc<CommitTracker> {
3442        let recorder = TestMetricsRecorder::default();
3443        let metrics = Arc::new(Metrics::new(&recorder, lanes));
3444        let tracker = CommitTracker::new(lanes, majority(lanes), metrics);
3445        for zone in 0..lanes {
3446            tracker.activate_lane(zone, 0);
3447        }
3448        tracker
3449    }
3450
3451    #[test]
3452    fn quorum_byte_watermark_evicts_resolved_record_boundaries() {
3453        let tracker = active_commit_tracker(3);
3454        let mut range = tracker.admit_window(&[10, 20, 30], &[0, 1, 2]);
3455
3456        tracker.publish_durable(0, 30);
3457        assert_eq!(range.progress().0, 0);
3458        tracker.publish_durable(1, 20);
3459        assert_eq!(range.progress().0, 2);
3460        {
3461            let state = tracker
3462                .state
3463                .lock()
3464                .unwrap_or_else(|poisoned| poisoned.into_inner());
3465            assert_eq!(state.committed_bytes, 20);
3466            assert_eq!(state.boundaries, VecDeque::from([30]));
3467        }
3468
3469        tracker.publish_durable(2, 30);
3470        assert_eq!(range.progress().0, 3);
3471        assert!(tracker
3472            .state
3473            .lock()
3474            .unwrap_or_else(|poisoned| poisoned.into_inner())
3475            .boundaries
3476            .is_empty());
3477    }
3478
3479    #[test]
3480    fn per_zone_durable_lag_tracks_admitted_and_persisted_bytes() {
3481        let recorder = Arc::new(TestMetricsRecorder::default());
3482        let metrics = Arc::new(Metrics::new(recorder.as_ref(), 1));
3483        let tracker = CommitTracker::new(1, 1, metrics);
3484        tracker.activate_lane(0, 0);
3485        let _range = tracker.admit_window(&[10, 20], &[0]);
3486        assert_eq!(
3487            recorder.labeled_gauge("chorus.wal.replica.durable_lag_bytes", &[("zone", "0")]),
3488            20
3489        );
3490
3491        tracker.publish_durable(0, 10);
3492        assert_eq!(
3493            recorder.labeled_gauge("chorus.wal.replica.durable_lag_bytes", &[("zone", "0")]),
3494            10
3495        );
3496        tracker.finish_lane(0, None);
3497        assert_eq!(
3498            recorder.labeled_gauge("chorus.wal.replica.durable_lag_bytes", &[("zone", "0")]),
3499            0
3500        );
3501    }
3502
3503    #[test]
3504    fn retired_lanes_cannot_support_future_admissions() {
3505        let tracker = active_commit_tracker(3);
3506        tracker.finish_lane(0, None);
3507        let mut range = tracker.admit_window(&[10], &[1, 2]);
3508
3509        tracker.finish_lane(1, None);
3510
3511        let (committed, failure) = range.progress();
3512        assert_eq!(committed, 0);
3513        assert!(matches!(failure, Some(ProtocolError::Poisoned)));
3514    }
3515
3516    #[tokio::test]
3517    async fn durable_progress_resets_the_lane_stall_timeout() {
3518        let recorder = Arc::new(TestMetricsRecorder::default());
3519        let metrics = Arc::new(Metrics::new(recorder.as_ref(), 1));
3520        let (sends, _) = mpsc::unbounded_channel();
3521        let replica = Arc::new(ScriptedLaneReplica::new(0, VecDeque::new(), sends));
3522        let replica_for_progress: Arc<dyn Replica> = replica.clone();
3523        let updater = replica.clone();
3524        let updates = tokio::spawn(async move {
3525            for durable in [10, 20, 30] {
3526                tokio::time::sleep(Duration::from_millis(5)).await;
3527                updater.durable.send_replace(durable);
3528            }
3529        });
3530        let timeout = Duration::from_millis(100);
3531        let mut seen = 0;
3532
3533        for expected in [10, 20, 30] {
3534            match lane_progress(
3535                &replica_for_progress,
3536                seen,
3537                true,
3538                tokio::time::Instant::now() + timeout,
3539                &metrics,
3540            )
3541            .await
3542            {
3543                LaneProgress::Advanced(change) => {
3544                    assert_eq!(change.persisted_size, expected);
3545                    assert!(change.error.is_none());
3546                    seen = change.persisted_size;
3547                }
3548                LaneProgress::Stalled => panic!("advancing lane was falsely shed"),
3549                LaneProgress::Failed(error) => panic!("advancing lane failed: {error}"),
3550            }
3551        }
3552        updates.await.unwrap();
3553        assert_eq!(recorder.counter("chorus.wal.lane.timeouts"), 0);
3554    }
3555
3556    #[tokio::test]
3557    async fn ready_progress_wins_an_expired_operation_deadline() {
3558        let recorder = Arc::new(TestMetricsRecorder::default());
3559        let metrics = Arc::new(Metrics::new(recorder.as_ref(), 1));
3560        let tracker = CommitTracker::new(1, 1, Arc::clone(&metrics));
3561        tracker.activate_lane(0, 0);
3562        let (sends, _) = mpsc::unbounded_channel();
3563        let replica = Arc::new(ScriptedLaneReplica::new(0, VecDeque::new(), sends));
3564        replica.durable.send_replace(10);
3565        let replica: Arc<dyn Replica> = replica;
3566        let mut durable = 0;
3567        let mut last_progress = tokio::time::Instant::now() - Duration::from_secs(1);
3568        let mut retained = VecDeque::new();
3569
3570        let change = confirm_lane_stall(&replica, durable, tokio::time::Instant::now(), &metrics)
3571            .await
3572            .expect("ready durable progress must win the deadline tie");
3573        publish_lane_advance(
3574            change,
3575            0,
3576            &mut durable,
3577            &mut last_progress,
3578            &mut retained,
3579            &tracker,
3580        )
3581        .expect("ready progress has no stream failure");
3582
3583        assert_eq!(durable, 10);
3584        assert_eq!(recorder.counter("chorus.wal.lane.timeouts"), 0);
3585    }
3586
3587    #[tokio::test]
3588    async fn fencing_lane_failure_stops_the_writer_after_publishing_progress() {
3589        let tracker = active_commit_tracker(1);
3590        let range = tracker.admit_window(&[10], &[0]);
3591
3592        tracker.publish_durable(0, 10);
3593        tracker.finish_lane(
3594            0,
3595            Some(TransportError {
3596                zone: 0,
3597                code: TransportCode::FailedPrecondition,
3598                message: "newer writer took over".into(),
3599            }),
3600        );
3601
3602        assert_eq!(range.into_pending().remove(0).wait().await.unwrap(), 0);
3603        assert!(tracker.is_poisoned());
3604        let mut updates = tracker.subscribe();
3605        let snapshot = updates.borrow_and_update().clone();
3606        assert!(matches!(snapshot.failure, Some(CommitFailure::Fenced(_))));
3607    }
3608
3609    #[tokio::test]
3610    async fn pending_commits_follow_the_prefix_watermark_and_gap_poison() {
3611        let tracker = active_commit_tracker(3);
3612        let range = tracker.admit_window(&[10, 20], &[0, 1, 2]);
3613        let mut pending = range.into_pending();
3614        let second = pending.pop().expect("second pending commit");
3615        let first = pending.pop().expect("first pending commit");
3616
3617        tracker.publish_durable(0, 20);
3618        tracker.publish_durable(1, 10);
3619
3620        tracker.finish_lane(
3621            1,
3622            Some(TransportError {
3623                zone: 1,
3624                code: TransportCode::PermissionDenied,
3625                message: "terminal".into(),
3626            }),
3627        );
3628        tracker.finish_lane(
3629            2,
3630            Some(TransportError {
3631                zone: 2,
3632                code: TransportCode::Unavailable,
3633                message: "transient".into(),
3634            }),
3635        );
3636        tracker.publish_durable(2, 20);
3637
3638        assert_eq!(first.wait().await.unwrap(), 0);
3639        assert!(matches!(
3640            second.wait().await,
3641            Err(ProtocolError::Transport(TransportError {
3642                zone: 1,
3643                code: TransportCode::PermissionDenied,
3644                ..
3645            }))
3646        ));
3647        assert_eq!(tracker.committed_len(), 1);
3648    }
3649
3650    #[tokio::test]
3651    async fn ready_progress_releases_lane_budget_before_more_work_is_staged() {
3652        let recorder = Arc::new(TestMetricsRecorder::default());
3653        let metrics = Arc::new(Metrics::new(recorder.as_ref(), 1));
3654        let (first_send_release_tx, first_send_release_rx) = oneshot::channel();
3655        let (second_send_release_tx, second_send_release_rx) = oneshot::channel();
3656        let (sends_tx, mut sends_rx) = mpsc::unbounded_channel();
3657        let replica: Arc<dyn Replica> = Arc::new(ScriptedLaneReplica::new(
3658            0,
3659            VecDeque::from([first_send_release_rx, second_send_release_rx]),
3660            sends_tx,
3661        ));
3662        let mut writer = Writer::new(
3663            vec![replica],
3664            ClientConfig::default(),
3665            protocol_metadata(),
3666            vec![AppendToken {
3667                zone: 0,
3668                generation: Some(1),
3669                metageneration: Some(1),
3670                persisted_size: 0,
3671                write_handle: None,
3672            }],
3673            metrics,
3674        );
3675        let encoded = record(b"first").encode().unwrap().len();
3676        writer.set_max_replica_lag_bytes(encoded * 2);
3677        let attempted: AttemptedBytes = Arc::new(|_| {});
3678
3679        let first = writer
3680            .enqueue_data_window(vec![record(b"first")], attempted.clone())
3681            .await
3682            .unwrap()
3683            .into_pending()
3684            .remove(0);
3685        assert_eq!(sends_rx.recv().await, Some(encoded as i64));
3686
3687        let second = writer
3688            .enqueue_data_window(vec![record(b"other")], attempted.clone())
3689            .await
3690            .unwrap()
3691            .into_pending()
3692            .remove(0);
3693        first_send_release_tx
3694            .send(())
3695            .expect("the first staged group should still be blocked");
3696        assert_eq!(sends_rx.recv().await, Some((encoded * 2) as i64));
3697
3698        let third = writer
3699            .enqueue_data_window(vec![record(b"third")], attempted)
3700            .await
3701            .unwrap()
3702            .into_pending()
3703            .remove(0);
3704        second_send_release_tx
3705            .send(())
3706            .expect("the second staged group should still be blocked");
3707
3708        assert_eq!(first.wait().await.unwrap(), 0);
3709        assert_eq!(second.wait().await.unwrap(), 1);
3710        assert_eq!(third.wait().await.unwrap(), 2);
3711        assert_eq!(writer.committed_len(), 3);
3712        assert_eq!(recorder.counter("chorus.wal.lane.capacity_drops"), 0);
3713        assert!(!writer.is_poisoned());
3714    }
3715
3716    #[tokio::test]
3717    async fn terminal_lane_failures_preserve_transport_errors_in_completions() {
3718        let replica = Arc::new(ReaderTerminalReplica::new(0));
3719        let factory: Arc<dyn crate::transport::ReplicaFactory> =
3720            Arc::new(SingleReplicaFactory::new(replica.clone()));
3721        let manifest_store =
3722            Arc::new(crate::manifest_store::test_support::InMemoryManifestStore::default());
3723        let volume = crate::segment::SegmentedVolume::new_with_factories_and_manifest_store(
3724            vec![factory],
3725            manifest_store,
3726            "terminal-completion",
3727            ClientConfig {
3728                max_retries: 0,
3729                retry_base: Duration::ZERO,
3730            },
3731        )
3732        .unwrap();
3733        let writer = volume.recover_writer().await.unwrap();
3734        let mut handle = crate::engine::WalEngine::start(
3735            writer,
3736            crate::WalEngineConfig {
3737                repair_interval: None,
3738                ..Default::default()
3739            },
3740        )
3741        .unwrap();
3742        let completion = handle
3743            .enqueue_append(
3744                crate::segment::WalSeqNo::ZERO,
3745                Bytes::from_static(b"terminal"),
3746            )
3747            .await
3748            .unwrap();
3749        let error = completion.await.unwrap_err();
3750        assert!(matches!(
3751            error,
3752            crate::Error::Transport {
3753                code: TransportCode::PermissionDenied,
3754                ..
3755            }
3756        ));
3757        assert_eq!(replica.resume_calls(), 0);
3758        let _ = tokio::time::timeout(Duration::from_secs(1), handle.shutdown())
3759            .await
3760            .expect("engine shutdown timed out");
3761    }
3762
3763    #[tokio::test]
3764    async fn stalled_lane_poison_releases_blocked_admission() {
3765        let replica: Arc<dyn Replica> = Arc::new(StalledReplica { zone: 0 });
3766        let factory: Arc<dyn crate::transport::ReplicaFactory> =
3767            Arc::new(SingleReplicaFactory::new(replica));
3768        let manifest_store =
3769            Arc::new(crate::manifest_store::test_support::InMemoryManifestStore::default());
3770        let volume = crate::segment::SegmentedVolume::new_with_factories_and_manifest_store(
3771            vec![factory],
3772            manifest_store,
3773            "stalled-admission",
3774            ClientConfig {
3775                max_retries: 0,
3776                retry_base: Duration::ZERO,
3777            },
3778        )
3779        .unwrap();
3780        let writer = volume.recover_writer().await.unwrap();
3781        let payload = Bytes::from_static(b"stalled");
3782        let encoded_bytes = payload.len() + 4;
3783        let stall_timeout = Duration::from_millis(20);
3784        let mut handle = crate::engine::WalEngine::start(
3785            writer,
3786            crate::WalEngineConfig {
3787                queue_capacity: 2,
3788                max_record_bytes: payload.len(),
3789                pipeline_window_records: 1,
3790                max_inflight_bytes: encoded_bytes,
3791                max_replica_lag_bytes: encoded_bytes,
3792                lane_stall_timeout: stall_timeout,
3793                repair_interval: None,
3794                ..Default::default()
3795            },
3796        )
3797        .unwrap();
3798        let first = handle
3799            .enqueue_append(crate::segment::WalSeqNo::ZERO, payload.clone())
3800            .await
3801            .unwrap();
3802
3803        let second = tokio::time::timeout(
3804            stall_timeout.saturating_mul(10),
3805            handle.enqueue_append(crate::segment::WalSeqNo::record(1), payload),
3806        )
3807        .await
3808        .expect("blocked admission did not wake after the writer poisoned");
3809        assert!(matches!(second, Err(crate::Error::Closed)));
3810
3811        let first = tokio::time::timeout(stall_timeout.saturating_mul(10), first)
3812            .await
3813            .expect("admitted append did not receive terminal poison");
3814        assert!(matches!(first, Err(crate::Error::Poisoned)));
3815        let _ = handle.shutdown().await;
3816    }
3817}