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
20pub(crate) const SUPPORTED_REPLICA_COUNTS: [usize; 3] = [1, 3, 5];
23
24pub(crate) const DEFAULT_LANE_STALL_TIMEOUT: Duration = Duration::from_secs(5);
29
30pub(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)]
49pub struct ClientConfig {
54 pub max_retries: usize,
56 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
77struct 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
179struct 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 !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
223pub(crate) struct RecoveredTail {
228 canonical: CanonicalPrefix,
229 had_discarded_suffix: bool,
230}
231
232pub(crate) enum RecoveryCandidate {
234 Absent,
236 Empty {
241 reusable_writer: Option<Box<Writer>>,
242 },
243 NonEmpty(RecoveredTail),
245}
246
247impl RecoveredTail {
248 pub fn len(&self) -> usize {
249 self.canonical.len()
250 }
251
252 pub fn digest(&self) -> String {
255 digest_bytes(&self.canonical.bytes)
256 }
257
258 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 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
327struct 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 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 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 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 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 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 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 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 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 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 } 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 pub async fn seal_digest(&mut self) -> String {
1404 self.admitted.digest().await
1405 }
1406
1407 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 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 self.admitted.queue_digest(chunks);
1510 Ok(self
1511 .commits
1512 .admit_window(&batch.boundaries, &represented_zones))
1513 }
1514
1515 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 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 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 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 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 #[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
1833struct 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
1864struct 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
1936const LANE_DRAIN_GRACE: std::time::Duration = std::time::Duration::from_millis(100);
1941
1942struct 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#[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 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 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 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 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
2224async 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
2251fn 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
2276enum 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
2343fn 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 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 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 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
2564pub(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
2585fn 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 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}