diff --git a/apps/ingest/src/main.rs b/apps/ingest/src/main.rs index ec7c6ba9b..b5423ad1d 100644 --- a/apps/ingest/src/main.rs +++ b/apps/ingest/src/main.rs @@ -2490,7 +2490,7 @@ async fn main() { "Shutdown drain deadline hit with WAL backlog remaining" ); } - // Even a clean drain runs this: it retires the owner heartbeat, so a + // Even a clean drain runs this: it marks the owner retired, so a // successor claims anything left instead of waiting out the staleness // window. With a backlog it also seals and ships the tail, which is the // difference between "replays if this task's storage survives" (it does diff --git a/apps/ingest/src/telemetry.rs b/apps/ingest/src/telemetry.rs index 2ca2762ff..308af5383 100644 --- a/apps/ingest/src/telemetry.rs +++ b/apps/ingest/src/telemetry.rs @@ -876,9 +876,9 @@ impl TelemetryPipeline { /// this task as alive. /// /// Shutdown calls this once the drain has done what it can: whatever did not - /// export is now in the object store under an owner whose heartbeat is gone, - /// so the next task claims it immediately instead of waiting out the - /// staleness window. Returns the bytes shipped. + /// export is now in the object store under a retired owner, which the next + /// task to boot claims without waiting out the staleness window. A clean + /// drain ships nothing. Returns the bytes shipped. pub async fn flush_wal_to_object_store(&self) -> u64 { let shipped = self.inner.wal.flush_to_object_store().await; if let Some(store) = self.inner.wal.store.get() { @@ -1563,20 +1563,32 @@ impl WalLane { /// Whether the exporter has already moved past this segment. fn is_exported(&self, seq: u64) -> bool { - self.export.lock().is_ok_and(|state| seq < state.cursor.seq) + self.export + .lock() + .is_ok_and(|state| self.is_behind(state.cursor, seq)) + } + + /// Whether `cursor` has consumed all of sealed segment `seq`. A cursor + /// parked at the end of a segment that sealed after its last export (the + /// shutdown seal, typically) counts: it holds nothing left to replay. + fn is_behind(&self, cursor: SegmentCursor, seq: u64) -> bool { + seq < cursor.seq + || (seq == cursor.seq + && seq < self.active_seq.load(Ordering::Acquire) + && cursor.offset >= file_len(&segment_path(&self.dir, seq))) } /// Sealed segments this lane still owes, oldest first. fn unexported_segments(&self) -> Vec { - let cursor_seq = match self.export.lock() { - Ok(state) => state.cursor.seq, + let cursor = match self.export.lock() { + Ok(state) => state.cursor, Err(_) => return Vec::new(), }; let active = self.active_seq.load(Ordering::Acquire); list_segments(&self.dir) .unwrap_or_default() .into_iter() - .filter(|seq| *seq >= cursor_seq && *seq < active) + .filter(|seq| *seq < active && !self.is_behind(cursor, *seq)) .collect() } @@ -1922,6 +1934,9 @@ impl ShardedWal { } }; let mut recovered_frames = 0usize; + // Segments still in the bucket after this pass. Any at all keeps the + // owner discoverable so a later boot retries them. + let mut left_behind = 0usize; let mut keys = Vec::with_capacity(segments.len()); for segment in segments { let Some(lane) = Self::lane_for_key(&segment.lane_key, cfg) else { @@ -1930,6 +1945,7 @@ impl ShardedWal { lane_key = segment.lane_key, "Skipping an orphaned WAL segment for an unknown lane" ); + left_behind += 1; continue; }; let frames = decode_segment_frames(&segment.bytes); @@ -1956,10 +1972,14 @@ impl ShardedWal { } if committed { keys.push(segment.key); + } else { + left_behind += 1; } } if recovered_frames == 0 && keys.is_empty() { - drop(store.release_owner(&owner).await); + if left_behind == 0 { + drop(store.release_owner(&owner).await); + } continue; } // Sealed and re-shipped under our own owner id *before* the source @@ -1969,9 +1989,16 @@ impl ShardedWal { for key in keys { if let Err(error) = store.release_key(&key).await { warn!(owner, key, error = %error, "Failed to delete a recovered WAL segment"); + left_behind += 1; } } - if let Err(error) = store.release_owner(&owner).await { + if left_behind > 0 { + warn!( + owner, + left_behind, + "Keeping a WAL owner so a later boot retries its remaining segments" + ); + } else if let Err(error) = store.release_owner(&owner).await { warn!(owner, error = %error, "Failed to retire a recovered WAL owner"); } metrics::wal_frames_recovered(recovered_frames as u64); @@ -7331,6 +7358,134 @@ mod tests { ); } + #[tokio::test] + async fn a_clean_drain_ships_nothing_but_an_unexported_tail_still_ships() { + // The shutdown seal leaves the cursor parked at the end of a segment it + // fully exported; shipping that one only buys duplicates on recovery. + let queue_dir = unique_test_dir("wal-clean-drain"); + std::fs::create_dir_all(&queue_dir).unwrap(); + let cfg = segmented_cfg(queue_dir.clone(), 64 * 1024, WAL_SEGMENT_MAX_BYTES); + let (endpoint, bucket) = fake_s3::spawn("maple-wal-test").await; + let wal = ShardedWal::open(&cfg).expect("open WAL"); + wal.attach_object_store(&test_wal_store(&endpoint)); + let frame = wal_test_frame(200); + + let (seq, _, end) = wal.append(0, &frame).await.expect("append"); + wal.mark_exported(0, SegmentCursor { seq, offset: end }, end) + .await + .expect("mark_exported"); + assert_eq!( + wal.flush_to_object_store().await, + 0, + "a drained lane owes nothing" + ); + + let (_, _, tail) = wal.append(0, &frame).await.expect("append tail"); + assert_eq!( + wal.flush_to_object_store().await, + tail, + "the unexported tail ships" + ); + sleep(Duration::from_millis(100)).await; + let segments: Vec<_> = bucket + .keys() + .into_iter() + .filter(|key| key.contains("/segments/")) + .collect(); + assert_eq!( + segments.len(), + 1, + "only the tail is in the bucket: {segments:?}" + ); + + drop(std::fs::remove_dir_all(queue_dir)); + } + + #[tokio::test] + async fn segments_shipped_at_shutdown_are_claimed_by_the_next_boot() { + // Retiring used to only delete the heartbeat, which left the owner + // invisible to recovery and its shutdown segments stranded in the bucket. + let (endpoint, bucket) = fake_s3::spawn("maple-wal-test").await; + let frame = wal_test_frame(64); + + let old_dir = unique_test_dir("wal-retired-owner"); + std::fs::create_dir_all(&old_dir).unwrap(); + let old_cfg = segmented_cfg(old_dir.clone(), 64 * 1024, WAL_SEGMENT_MAX_BYTES); + let old_wal = ShardedWal::open(&old_cfg).expect("open WAL"); + // No shipper attached, so its async upload cannot race the cleanup. + let old_store = test_wal_store(&endpoint); + old_store.heartbeat().await.expect("heartbeat"); + old_wal.append(0, &frame).await.expect("append"); + assert!( + old_wal.flush_segments_to(&old_store).await > 0, + "the undrained tail ships" + ); + old_store.retire().await.expect("retire"); + let retired = old_store.owner().to_owned(); + + // Booting right away: no staleness window to wait out. + let new_dir = unique_test_dir("wal-retired-successor"); + std::fs::create_dir_all(&new_dir).unwrap(); + let new_cfg = segmented_cfg(new_dir.clone(), 64 * 1024, WAL_SEGMENT_MAX_BYTES); + let new_wal = ShardedWal::open(&new_cfg).expect("open WAL"); + new_wal + .recover_orphans(&new_cfg, &test_wal_store(&endpoint)) + .await; + + let replayed = new_wal.replay(0).await.expect("replay"); + assert_eq!(replayed.len(), 1, "the shipped frame is recovered"); + let keys = bucket.keys(); + assert!( + keys.iter().all(|key| !key.contains(&retired)), + "the retired owner's segments and markers are cleaned up: {keys:?}" + ); + + drop(std::fs::remove_dir_all(old_dir)); + drop(std::fs::remove_dir_all(new_dir)); + } + + #[tokio::test] + async fn an_owner_with_unrecovered_segments_stays_claimable() { + // Releasing the owner after a partial recovery dropped its last discovery + // path, so the segments it still held were never retried. + let (endpoint, bucket) = fake_s3::spawn("maple-wal-test").await; + let frame = wal_test_frame(64); + + let old_dir = unique_test_dir("wal-partial-owner"); + std::fs::create_dir_all(&old_dir).unwrap(); + let old_cfg = segmented_cfg(old_dir.clone(), 64 * 1024, WAL_SEGMENT_MAX_BYTES); + let old_wal = ShardedWal::open(&old_cfg).expect("open WAL"); + let old_store = test_wal_store(&endpoint); + old_store.heartbeat().await.expect("heartbeat"); + old_wal.append(0, &frame).await.expect("append"); + assert!(old_wal.flush_segments_to(&old_store).await > 0); + // A lane this binary cannot place, so recovery has to leave it behind. + old_store + .put_segment("shard-000-unknown", 0, b"unplaceable".to_vec()) + .await + .expect("put unknown-lane segment"); + old_store.retire().await.expect("retire"); + let retired = old_store.owner().to_owned(); + + let new_dir = unique_test_dir("wal-partial-successor"); + std::fs::create_dir_all(&new_dir).unwrap(); + let new_cfg = segmented_cfg(new_dir.clone(), 64 * 1024, WAL_SEGMENT_MAX_BYTES); + let new_wal = ShardedWal::open(&new_cfg).expect("open WAL"); + new_wal + .recover_orphans(&new_cfg, &test_wal_store(&endpoint)) + .await; + + assert_eq!(new_wal.replay(0).await.expect("replay").len(), 1); + let keys = bucket.keys(); + assert!( + keys.iter().any(|key| key.ends_with(&format!("/retired/{retired}"))), + "the retired marker survives while a segment is left: {keys:?}" + ); + + drop(std::fs::remove_dir_all(old_dir)); + drop(std::fs::remove_dir_all(new_dir)); + } + /// Cross-language contract with the Prometheus scraper (apps/scraper). /// /// The scraper converts scraped exposition text into OTLP/JSON diff --git a/apps/ingest/src/wal_store.rs b/apps/ingest/src/wal_store.rs index 175a66f50..4f0b6f683 100644 --- a/apps/ingest/src/wal_store.rs +++ b/apps/ingest/src/wal_store.rs @@ -130,6 +130,14 @@ impl WalSegmentStore { format!("{}/{KEY_SCHEME}/claims/{owner}", self.prefix) } + fn retired_key(&self, owner: &str) -> String { + format!("{}/{KEY_SCHEME}/retired/{owner}", self.prefix) + } + + fn listing(&self, kind: &str) -> String { + format!("{}/{KEY_SCHEME}/{kind}/", self.prefix) + } + /// Ship one sealed segment. Overwrites unconditionally: the only writer for /// this key is this process, and a retry of a partial PUT must win. pub async fn put_segment( @@ -157,29 +165,34 @@ impl WalSegmentStore { .await } - /// Remove this owner's heartbeat, so a successor does not have to wait out - /// `orphan_after` before claiming whatever we failed to drain. + /// Mark this owner as gone for good, so a successor claims whatever we + /// shipped at shutdown right away. The marker lands before the heartbeat is + /// dropped, so the owner is never invisible to `stale_owners`. pub async fn retire(&self) -> Result<(), S3Error> { + self.s3 + .put(&self.retired_key(&self.owner), Vec::new(), false) + .await?; self.s3.delete(&self.owner_key(&self.owner)).await } - /// Owners whose heartbeat has gone stale, newest-stale first. + /// Owners whose segments are claimable now: retired ones at any age, and + /// ones whose heartbeat went stale. pub async fn stale_owners(&self, now: DateTime) -> Result, S3Error> { - let prefix = format!("{}/{KEY_SCHEME}/owners/", self.prefix); - let objects = self.s3.list(&prefix, MAX_LIST_PAGES).await?; - Ok(objects - .into_iter() - .filter(|object| !object.key.ends_with(&self.owner)) - .filter(|object| is_older_than(object, now, self.orphan_after)) - .filter_map(|object| { - object - .key - .rsplit('/') - .next() - .filter(|owner| !owner.is_empty()) - .map(str::to_owned) - }) - .collect()) + let heartbeats = self + .s3 + .list(&self.listing("owners"), MAX_LIST_PAGES) + .await?; + let retired = self + .s3 + .list(&self.listing("retired"), MAX_LIST_PAGES) + .await?; + Ok(claimable_owners( + &self.owner, + &heartbeats, + &retired, + now, + self.orphan_after, + )) } /// Take ownership of a dead task's segments, or report that someone else @@ -238,13 +251,15 @@ impl WalSegmentStore { } /// Finish a claim: drop the recovered object and, once every segment is - /// gone, the owner's heartbeat and claim marker. + /// gone, the owner's heartbeat, retired marker and claim marker. pub async fn release_key(&self, key: &str) -> Result<(), S3Error> { self.s3.delete(key).await } pub async fn release_owner(&self, owner: &str) -> Result<(), S3Error> { self.s3.delete(&self.owner_key(owner)).await?; + self.s3.delete(&self.retired_key(owner)).await?; + // The claim goes last: until it does, no other task re-claims the owner. self.s3.delete(&self.claim_key(owner)).await } } @@ -260,6 +275,32 @@ fn is_older_than(object: &S3Object, now: DateTime, age: Duration) -> bool { .is_ok_and(|elapsed| elapsed >= age) } +/// Which owners other than `own` may be claimed, in key order. +fn claimable_owners( + own: &str, + heartbeats: &[S3Object], + retired: &[S3Object], + now: DateTime, + orphan_after: Duration, +) -> Vec { + let stale = heartbeats + .iter() + .filter(|object| is_older_than(object, now, orphan_after)) + .filter_map(|object| last_component(&object.key)); + // A retired owner shut down on purpose and will never write again, so its + // age does not matter. + let retired = retired + .iter() + .filter_map(|object| last_component(&object.key)); + let claimable: std::collections::BTreeSet<&str> = + stale.chain(retired).filter(|owner| *owner != own).collect(); + claimable.into_iter().map(str::to_owned).collect() +} + +fn last_component(key: &str) -> Option<&str> { + key.rsplit('/').next().filter(|part| !part.is_empty()) +} + /// `…/segments///.seg` → ``. fn lane_key_of(key: &str) -> Option { let mut parts = key.rsplit('/'); @@ -303,6 +344,27 @@ mod tests { )); } + #[test] + fn retired_owners_are_claimable_at_any_age() { + let heartbeats = [ + object("wal/v1/owners/live", 1), + object("wal/v1/owners/dead", 30), + object("wal/v1/owners/me", 30), + ]; + // Retired a moment ago: no staleness window to wait out. + let retired = [object("wal/v1/retired/just-retired", 0)]; + assert_eq!( + claimable_owners( + "me", + &heartbeats, + &retired, + Utc::now(), + DEFAULT_ORPHAN_AFTER + ), + vec!["dead", "just-retired"] + ); + } + #[test] fn a_listing_without_a_timestamp_reads_as_stale() { // Otherwise an unparsed LastModified would strand the segments under it