From 5c9c424a84fdd6c2ebc427f19cd94f690b68f82f Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Sun, 27 Sep 2026 14:38:43 +0300 Subject: [PATCH] volume server: stream ReadAllNeedles without holding the store lock (#11481) * volume server: stream ReadAllNeedles without holding the store lock read_all_needles held store.read() while Volume::read_all_needles read every live needle of the volume into a Vec, and kept holding it through the whole blocking_send loop. Memory grew with the volume, and a slow client parked the scan in a send with the guard held; needle writes and the heartbeat take store.write() on a writer-preferring lock, so the node stopped serving until the client caught up. Take a DatScanPlan (fresh .dat open, end bound) under a short guard and walk it with the guard released, sending one needle at a time. Each record is checked against the live needle map under a brief read guard, as the scan reaches it, and only a live record is parsed, so a damaged stale copy does not fail the stream. Records appended while a pass ran are walked by a follow-up plan, so a needle overwritten during the scan is streamed once, as its new copy. A vacuum commit or re-create of the volume during the scan fails the stream, since the map's offsets no longer describe the pinned file; the plan carries the volume instance and compaction revision for that check. DatScanPlan::scan_records yields records unparsed; scan keeps its behaviour on top of it. Volume::read_all_needles has no caller left and is removed; its tests move to the RPC. Co-Authored-By: Claude Opus 5.5 (1M context) * volume server: check ReadAllNeedles liveness only once the send can proceed The per-record liveness check ran before blocking_send, so a scan parked on a full channel held a record it had already judged live. An overwrite landing during that park left the old copy in the stream, and the continuation over appended records then streamed the new copy as well. Reserve channel space first, then take the store read guard, check the record against the needle map and enqueue it through the permit before releasing the guard. The wait for space still happens without the lock; the record is parsed before the guard is taken, and its parse error only counts if the record turns out to be live. Co-Authored-By: Claude Opus 5.5 (1M context) --------- Co-authored-by: Claude Opus 5.5 (1M context) --- seaweed-volume/src/server/grpc_server.rs | 403 ++++++++++++++++++++--- seaweed-volume/src/storage/volume.rs | 226 ++++++------- 2 files changed, 457 insertions(+), 172 deletions(-) diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index a984d5036..8c9fde49f 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -2802,48 +2802,15 @@ impl VolumeServer for VolumeGrpcService { let (tx, rx) = tokio::sync::mpsc::channel(32); - // Stream needles lazily via a blocking task (matches Go's scanner pattern) tokio::task::spawn_blocking(move || { - let store = state.store.read().unwrap(); for &raw_vid in &req.volume_ids { - let vid = VolumeId(raw_vid); - let v = match store.find_volume(vid) { - Some((_, v)) => v, - None => { - let _ = tx.blocking_send(Err(Status::not_found(format!( - "not found volume id {}", - vid - )))); + match read_all_needles_of_volume(&state, VolumeId(raw_vid), &tx) { + Ok(()) => {} + Err(Some(status)) => { + let _ = tx.blocking_send(Err(status)); return; } - }; - - let needles = match v.read_all_needles() { - Ok(n) => n, - Err(e) => { - let _ = tx.blocking_send(Err(Status::internal(e.to_string()))); - return; - } - }; - - for n in needles { - let compressed = n.is_compressed(); - if tx - .blocking_send(Ok(volume_server_pb::ReadAllNeedlesResponse { - volume_id: raw_vid, - needle_id: n.id.into(), - cookie: n.cookie.0, - needle_blob: n.data, - needle_blob_compressed: compressed, - last_modified: n.last_modified, - crc: n.checksum.0, - name: n.name, - mime: n.mime, - })) - .is_err() - { - return; // receiver dropped - } + Err(None) => return, // receiver dropped } } }); @@ -6082,6 +6049,97 @@ fn tail_pass( } } +/// Stream the live needles of one volume in .dat order, one at a time and +/// without holding the store lock while reading or waiting for channel space. +/// Liveness is checked against the needle map per record, and the scan follows +/// appends made while it ran. `Err(None)` means the receiver hung up. +fn read_all_needles_of_volume( + state: &VolumeServerState, + vid: VolumeId, + tx: &tokio::sync::mpsc::Sender>, +) -> Result<(), Option> { + let not_found = || Some(Status::not_found(format!("not found volume id {}", vid))); + let internal = |e: crate::storage::volume::VolumeError| Some(Status::internal(e.to_string())); + let mut plan = { + let store = state.store.read().unwrap(); + let (_, v) = store.find_volume(vid).ok_or_else(not_found)?; + v.dat_scan_plan(v.super_block.block_size() as u64) + .map_err(internal)? + }; + loop { + let mut next = None; + let mut failed = None; + let mut client_gone = false; + let scanned = plan.scan_records(|record| { + if tx.is_closed() { + client_gone = true; + return Ok(ControlFlow::Break(())); + } + next = Some(record.offset + record.bytes.len() as u64); + if record.size.0 <= 0 { + return Ok(ControlFlow::Continue(())); + } + // A stale copy's parse error must not fail the stream. + let parsed = record.parse(); + // Wait for channel space before the liveness check, then enqueue + // under the guard, so no overwrite can land between the two. + let Ok(permit) = futures::executor::block_on(tx.reserve()) else { + client_gone = true; + return Ok(ControlFlow::Break(())); + }; + let store = state.store.read().unwrap(); + let live = match store.find_volume(vid) { + Some((_, v)) => v + .is_live_in_plan(&plan, record.id, record.offset) + .map_err(internal), + None => Err(not_found()), + }; + match live { + Ok(true) => {} + Ok(false) => return Ok(ControlFlow::Continue(())), + Err(status) => { + failed = status; + return Ok(ControlFlow::Break(())); + } + } + let n = parsed?; + let compressed = n.is_compressed(); + let msg = volume_server_pb::ReadAllNeedlesResponse { + volume_id: vid.0, + needle_id: n.id.into(), + cookie: n.cookie.0, + needle_blob: n.data, + needle_blob_compressed: compressed, + last_modified: n.last_modified, + crc: n.checksum.0, + name: n.name, + mime: n.mime, + }; + permit.send(Ok(msg)); + Ok(ControlFlow::Continue(())) + }); + if client_gone { + return Err(None); + } + if failed.is_some() { + return Err(failed); + } + scanned.map_err(internal)?; + + // Records appended while the pass ran lie past its end bound; an + // overwrite's new copy is among them, so continue until none are left. + let after = { + let store = state.store.read().unwrap(); + let (_, v) = store.find_volume(vid).ok_or_else(not_found)?; + v.dat_scan_plan_after(&plan, next).map_err(internal)? + }; + match after { + Some(p) => plan = p, + None => return Ok(()), + } + } +} + /// Get disk usage (total, free) in bytes for the given path. fn get_disk_usage(path: &str) -> (u64, u64) { use sysinfo::Disks; @@ -7437,6 +7495,275 @@ mod tests { ); } + fn write_volume_needle(service: &VolumeGrpcService, id: u64, data: &[u8]) { + let mut store = service.state.store.write().unwrap(); + let (_, v) = store.find_volume_mut(VolumeId(1)).unwrap(); + let mut n = Needle { + id: NeedleId(id), + cookie: Cookie(0x1234), + data: data.to_vec(), + data_size: data.len() as u32, + ..Needle::default() + }; + v.write_needle(&mut n, true, false).unwrap(); + } + + async fn read_all_needles_stream( + service: &VolumeGrpcService, + volume_ids: Vec, + ) -> BoxStream { + service + .read_all_needles(Request::new(volume_server_pb::ReadAllNeedlesRequest { + volume_ids, + })) + .await + .unwrap() + .into_inner() + } + + async fn collect_read_all( + mut stream: BoxStream, + ) -> Vec> { + tokio::time::timeout(std::time::Duration::from_secs(30), async { + let mut items = Vec::new(); + while let Some(item) = stream.next().await { + items.push(item); + } + items + }) + .await + .expect("the stream must end") + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_read_all_needles_uses_dat_order_for_live_offsets() { + let (service, _tmp) = make_local_service_with_volume("", None); + write_volume_needle(&service, 10, b"first"); + write_volume_needle(&service, 20, b"second"); + write_volume_needle(&service, 10, b"first-overwrite"); + + let items = collect_read_all(read_all_needles_stream(&service, vec![1]).await).await; + let needles: Vec<_> = items.into_iter().map(|r| r.unwrap()).collect(); + let ids: Vec = needles.iter().map(|n| n.needle_id).collect(); + let bodies: Vec<&[u8]> = needles.iter().map(|n| n.needle_blob.as_slice()).collect(); + assert_eq!(ids, vec![11, 20, 10]); + assert_eq!( + bodies, + vec![ + b"ec-generate".as_slice(), + b"second".as_slice(), + b"first-overwrite".as_slice() + ] + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn read_all_needles_fails_at_a_header_it_cannot_advance_past() { + let (service, _tmp) = make_local_service_with_volume("", None); + let dat_path = { + let store = service.state.store.read().unwrap(); + let (_, v) = store.find_volume(VolumeId(1)).unwrap(); + v.file_name(".dat") + }; + let mut corrupt = [0u8; NEEDLE_HEADER_SIZE]; + NeedleId(99).to_bytes(&mut corrupt[4..12]); + Size(-100).to_bytes(&mut corrupt[12..16]); + let mut dat = std::fs::OpenOptions::new() + .append(true) + .open(dat_path) + .unwrap(); + std::io::Write::write_all(&mut dat, &corrupt).unwrap(); + + let items = collect_read_all(read_all_needles_stream(&service, vec![1]).await).await; + assert_eq!( + items.len(), + 2, + "the needle before the corrupt header, then the error" + ); + assert_eq!(items[0].as_ref().unwrap().needle_id, 11); + let err = items[1].as_ref().unwrap_err(); + assert_eq!(err.code(), tonic::Code::Internal); + assert!(err.message().contains("corrupt needle header"), "{err:?}"); + } + + // A damaged body in a copy that was since overwritten must not fail the + // stream: only the live copy is parsed. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_read_all_needles_skips_a_corrupt_stale_copy() { + let (service, _tmp) = make_local_service_with_volume("", None); + let (stale_offset, dat_path) = { + let store = service.state.store.read().unwrap(); + let (_, v) = store.find_volume(VolumeId(1)).unwrap(); + (v.dat_file_size().unwrap(), v.file_name(".dat")) + }; + write_volume_needle(&service, 50, b"old-data"); + write_volume_needle(&service, 50, b"new-data"); + // The first body field is DataSize; u32::MAX cannot fit the record. + let mut dat = std::fs::OpenOptions::new() + .write(true) + .open(dat_path) + .unwrap(); + std::io::Seek::seek( + &mut dat, + std::io::SeekFrom::Start(stale_offset + NEEDLE_HEADER_SIZE as u64), + ) + .unwrap(); + std::io::Write::write_all(&mut dat, &u32::MAX.to_be_bytes()).unwrap(); + drop(dat); + + let items = collect_read_all(read_all_needles_stream(&service, vec![1]).await).await; + let needles: Vec<_> = items.into_iter().map(|r| r.unwrap()).collect(); + let got: Vec<(u64, &[u8])> = needles + .iter() + .map(|n| (n.needle_id, n.needle_blob.as_slice())) + .collect(); + assert_eq!( + got, + vec![ + (11, b"ec-generate".as_slice()), + (50, b"new-data".as_slice()) + ] + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_read_all_needles_reports_a_missing_volume_after_the_found_ones() { + let (service, _tmp) = make_local_service_with_volume("", None); + let items = collect_read_all(read_all_needles_stream(&service, vec![1, 77]).await).await; + assert_eq!(items.len(), 2); + assert_eq!(items[0].as_ref().unwrap().needle_id, 11); + assert_eq!(items[1].as_ref().unwrap_err().code(), tonic::Code::NotFound); + } + + // Open a ReadAllNeedles stream over more needles than its channel holds, + // read one message so the scan is known to be running (it then fills the + // channel and parks in a send), and apply `f` under store.write(). Returns + // the first message and the rest of the stream. + async fn park_read_all_then_write( + service: &VolumeGrpcService, + ids: &[u64], + f: impl FnOnce(&mut crate::storage::volume::Volume), + ) -> ( + volume_server_pb::ReadAllNeedlesResponse, + BoxStream, + ) { + for &id in ids { + write_volume_needle(service, id, format!("body-{id}").as_bytes()); + } + let mut stream = read_all_needles_stream(service, vec![1]).await; + let first = stream.next().await.unwrap().unwrap(); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + loop { + // The scan takes a read guard briefly per record, so retry. + if let Ok(mut store) = service.state.store.try_write() { + let (_, v) = store.find_volume_mut(VolumeId(1)).unwrap(); + f(v); + return (first, stream); + } + if std::time::Instant::now() > deadline { + drop(stream); + panic!("store.write() stayed blocked while the scan was parked in a send"); + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + } + + // A client that stops reading parks the scan waiting for channel space. The + // store lock must not be held there, or needle writes and the heartbeat, + // which take store.write(), stall behind it. A needle overwritten before the + // scan reaches it is streamed once, as its new copy, the way Go's scan reads + // on to EOF. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_read_all_needles_releases_the_store_lock_while_parked() { + let (service, _tmp) = make_local_service_with_volume("", None); + let ids: Vec = (100..140).collect(); + let (first, stream) = park_read_all_then_write(&service, &ids, |v| { + let mut n = Needle { + id: NeedleId(138), + cookie: Cookie(0x1234), + data: b"overwritten".to_vec(), + data_size: b"overwritten".len() as u32, + ..Needle::default() + }; + v.write_needle(&mut n, true, false).unwrap(); + }) + .await; + + let items = collect_read_all(stream).await; + let needles: Vec<_> = std::iter::once(first) + .chain(items.into_iter().map(|r| r.unwrap())) + .collect(); + let got: Vec = needles.iter().map(|n| n.needle_id).collect(); + let mut want = vec![11]; + want.extend(ids.iter().copied().filter(|&id| id != 138)); + want.push(138); + assert_eq!(got, want); + assert_eq!(needles.last().unwrap().needle_blob, b"overwritten"); + } + + // The record the scan is parked on is checked for liveness only once it has + // channel space, so overwriting it meanwhile streams it once, as its new + // copy, not the old copy followed by the new one. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_read_all_needles_skips_a_record_overwritten_while_parked() { + let (service, _tmp) = make_local_service_with_volume("", None); + for id in 100..104 { + write_volume_needle(&service, id, b"old"); + } + let (tx, mut rx) = tokio::sync::mpsc::channel(1); + let state = service.state.clone(); + let producer = tokio::task::spawn_blocking(move || { + read_all_needles_of_volume(&state, VolumeId(1), &tx).map_err(|e| e.map(|s| s.code())) + }); + // Needle 11 fills the channel and the scan parks on needle 100. + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + while rx.len() < rx.max_capacity() { + assert!( + std::time::Instant::now() < deadline, + "the scan never filled the channel" + ); + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + // Lets a scan that checks before waiting reach its wait; the fixed + // order is correct under any timing. + tokio::time::sleep(std::time::Duration::from_millis(500)).await; + write_volume_needle(&service, 100, b"new"); + + let mut got = Vec::new(); + while let Some(item) = rx.recv().await { + let n = item.unwrap(); + got.push((n.needle_id, n.needle_blob)); + } + assert_eq!(producer.await.unwrap(), Ok(())); + let want: Vec<(u64, Vec)> = vec![ + (11, b"ec-generate".to_vec()), + (101, b"old".to_vec()), + (102, b"old".to_vec()), + (103, b"old".to_vec()), + (100, b"new".to_vec()), + ]; + assert_eq!(got, want); + } + + // After a vacuum commit the needle map describes the new .dat, not the + // one the scan pinned, so the stream fails instead of guessing. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_read_all_needles_fails_when_the_volume_is_compacted_mid_scan() { + let (service, _tmp) = make_local_service_with_volume("", None); + let ids: Vec = (100..140).collect(); + let (_first, stream) = park_read_all_then_write(&service, &ids, |v| { + v.compact_by_index(0, 0, |_| true).unwrap(); + v.commit_compact().unwrap(); + }) + .await; + + let items = collect_read_all(stream).await; + let err = items.last().unwrap().as_ref().unwrap_err(); + assert_eq!(err.code(), tonic::Code::Internal); + assert!(err.message().contains("compacted or replaced"), "{err:?}"); + assert!(items[..items.len() - 1].iter().all(|r| r.is_ok())); + } + // copy_file must stop exactly at stop_offset, never streaming past it. #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn test_copy_file_respects_stop_offset() { diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index cdbb64c6b..1d02bfb57 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -476,6 +476,10 @@ pub(crate) struct DatScanPlan { version: Version, from: u64, end: u64, + /// The volume instance and revision the plan was taken from: after a + /// vacuum or a re-create, the needle map's offsets name another file. + owner: Arc, + compaction_revision: u16, } /// One record visited by `DatScanPlan::scan`: the raw on-disk bytes and the @@ -486,6 +490,24 @@ pub(crate) struct RawNeedle<'a> { pub append_at_ns: u64, } +/// One record visited by `DatScanPlan::scan_records`, not yet parsed. +pub(crate) struct DatRecord<'a> { + pub offset: u64, + pub id: NeedleId, + pub size: Size, + /// Header and body. + pub bytes: &'a [u8], + version: Version, +} + +impl DatRecord<'_> { + pub(crate) fn parse(&self) -> Result { + let mut n = Needle::default(); + n.read_bytes(self.bytes, self.offset as i64, self.size, self.version)?; + Ok(n) + } +} + impl DatScanPlan { /// Visit the records in `[from, end)` in file order, the way Go's /// `ScanVolumeFileFrom` feeds a `VolumeFileScanner`: only the record being @@ -499,6 +521,22 @@ impl DatScanPlan { pub(crate) fn scan( &self, mut visit: impl FnMut(RawNeedle<'_>) -> ControlFlow<()>, + ) -> Result<(), VolumeError> { + self.scan_records(|record| { + let append_at_ns = record.parse()?.append_at_ns; + Ok(visit(RawNeedle { + header: &record.bytes[..NEEDLE_HEADER_SIZE], + body: &record.bytes[NEEDLE_HEADER_SIZE..], + append_at_ns, + })) + }) + } + + /// `scan` without the parse: a visitor that skips a record never fails on + /// a body it did not need. + pub(crate) fn scan_records( + &self, + mut visit: impl FnMut(DatRecord<'_>) -> Result, VolumeError>, ) -> Result<(), VolumeError> { let mut offset = self.from; while offset + NEEDLE_HEADER_SIZE as u64 <= self.end { @@ -552,17 +590,14 @@ impl DatScanPlan { ) .map_err(|e| self.read_error(e, offset + NEEDLE_HEADER_SIZE as u64))?; - let append_at_ns = { - let mut n = Needle::default(); - n.read_bytes(&record, offset as i64, size, self.version)?; - n.append_at_ns + let record = DatRecord { + offset, + id, + size, + bytes: &record, + version: self.version, }; - let needle = RawNeedle { - header: &record[..NEEDLE_HEADER_SIZE], - body: &record[NEEDLE_HEADER_SIZE..], - append_at_ns, - }; - if visit(needle).is_break() { + if visit(record)?.is_break() { break; } offset += total_size; @@ -2438,71 +2473,6 @@ impl Volume { .as_secs() } - /// Read all live needles from the volume (for ReadAllNeedles streaming RPC). - pub fn read_all_needles(&self) -> Result, VolumeError> { - let _guard = self.data_file_access_control.read_lock(); - if let Some(e) = self.unavailable_error() { - return Err(e); - } - let nm = self.nm_or_not_found()?; - let version = self.version(); - let dat_size = self.current_dat_file_size()? as i64; - let mut needles = Vec::new(); - let mut offset = self.super_block.block_size() as i64; - - while offset < dat_size { - let mut header = [0u8; NEEDLE_HEADER_SIZE]; - match self.read_exact_at_backend(&mut header, offset as u64) { - Ok(()) => {} - Err(VolumeError::Io(e)) if e.kind() == io::ErrorKind::UnexpectedEof => break, - Err(e) => return Err(e), - } - - let (_cookie, key, size) = Needle::parse_header(&header); - if size.0 == 0 && key.is_empty() { - break; - } - - let total_size = get_actual_size(size, version); - // A corrupt header can make the record length zero or negative; - // the scan cannot advance past it. - if total_size <= 0 { - return Err(VolumeError::Io(io::Error::new( - io::ErrorKind::InvalidData, - format!( - "corrupt needle header at offset {offset}: size {}, record length {total_size}", - size.0 - ), - ))); - } - - if size.is_deleted() || size.0 <= 0 { - offset += total_size; - continue; - } - - let Some(nv) = nm.get(key)? else { - offset += total_size; - continue; - }; - if nv.offset.to_actual_offset() != offset { - offset += total_size; - continue; - } - - let mut n = Needle { - id: key, - ..Needle::default() - }; - let mut read_option = ReadOption::default(); - self.read_needle_data_at_unlocked(&mut n, offset, size, &mut read_option)?; - needles.push(n); - - offset += total_size; - } - Ok(needles) - } - /// Check volume data integrity by verifying the last index entries against the .dat file. /// Matches Go's CheckVolumeDataIntegrity (volume_checking.go L117-141). /// Reads the last few index entries, verifies each needle header is readable and @@ -3038,9 +3008,56 @@ impl Volume { version: self.version(), from: from_offset, end: self.current_dat_file_size()?, + owner: self.data_file_access_control.clone(), + compaction_revision: self.super_block.compaction_revision, }) } + /// Fail unless `plan` was taken from this volume at its current + /// compaction revision. + fn check_plan_source(&self, plan: &DatScanPlan) -> Result<(), VolumeError> { + if let Some(e) = self.unavailable_error() { + return Err(e); + } + if !Arc::ptr_eq(&plan.owner, &self.data_file_access_control) + || plan.compaction_revision != self.super_block.compaction_revision + { + return Err(VolumeError::Io(io::Error::other(format!( + "volume {} was compacted or replaced during the scan", + self.id + )))); + } + Ok(()) + } + + /// Whether the record of `key` at `offset` in `plan` is the one the + /// needle map points at now. + pub(crate) fn is_live_in_plan( + &self, + plan: &DatScanPlan, + key: NeedleId, + offset: u64, + ) -> Result { + self.check_plan_source(plan)?; + let nv = self.nm_or_not_found()?.get(key)?; + Ok(nv.is_some_and(|nv| nv.offset.to_actual_offset() == offset as i64)) + } + + /// A plan for the records appended since `plan` was taken, or `None` when + /// there are none. `next` is the offset after the last record the pass + /// visited; a pass that stopped short of its end bound is not continued. + pub(crate) fn dat_scan_plan_after( + &self, + plan: &DatScanPlan, + next: Option, + ) -> Result, VolumeError> { + self.check_plan_source(plan)?; + if next.unwrap_or(plan.from) != plan.end || self.current_dat_file_size()? <= plan.end { + return Ok(None); + } + self.dat_scan_plan(plan.end).map(Some) + } + /// Insert or update a needle index entry (for low-level blob writes). pub fn put_needle_index( &mut self, @@ -5477,21 +5494,6 @@ mod tests { ); } - #[test] - fn read_all_needles_fails_at_a_header_it_cannot_advance_past() { - let tmp = TempDir::new().unwrap(); - let dir = tmp.path().to_str().unwrap(); - let mut v = make_test_volume(dir); - write_test_needle(&mut v, 1, b"good"); - append_corrupt_header(&v, -100); - - let err = v.read_all_needles().unwrap_err(); - assert!( - matches!(&err, VolumeError::Io(e) if e.kind() == io::ErrorKind::InvalidData), - "expected a corrupt-data error, got {err:?}" - ); - } - #[test] fn test_volume_write_read() { let tmp = TempDir::new().unwrap(); @@ -7069,50 +7071,6 @@ mod tests { assert!(!Path::new(&idx_path).exists()); } - #[test] - fn test_read_all_needles_uses_dat_order_for_live_offsets() { - let tmp = TempDir::new().unwrap(); - let dir = tmp.path().to_str().unwrap(); - let mut v = make_test_volume(dir); - - let mut first = Needle { - id: NeedleId(10), - cookie: Cookie(0x11223344), - data: b"first".to_vec(), - data_size: 5, - ..Needle::default() - }; - v.write_needle(&mut first, true, false).unwrap(); - - let mut second = Needle { - id: NeedleId(20), - cookie: Cookie(0x55667788), - data: b"second".to_vec(), - data_size: 6, - ..Needle::default() - }; - v.write_needle(&mut second, true, false).unwrap(); - - let mut first_overwrite = Needle { - id: NeedleId(10), - cookie: Cookie(0x11223344), - data: b"first-overwrite".to_vec(), - data_size: 15, - ..Needle::default() - }; - v.write_needle(&mut first_overwrite, true, false).unwrap(); - - let needles = v.read_all_needles().unwrap(); - let ids: Vec = needles.iter().map(|n| u64::from(n.id)).collect(); - let bodies: Vec<&[u8]> = needles.iter().map(|n| n.data.as_slice()).collect(); - - assert_eq!(ids, vec![20, 10]); - assert_eq!( - bodies, - vec![b"second".as_slice(), b"first-overwrite".as_slice()] - ); - } - #[test] fn test_get_append_at_ns() { let t1 = get_append_at_ns(0);