From d8f926cf46afd4bf61d1f92d62787c584ba4fd2f Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Thu, 1 Oct 2026 18:15:38 +0300 Subject: [PATCH] volume server: stream needle chunks without the store lock (#11488) * volume server: read GET/HEAD needles off the store lock, and only once The GET/HEAD handler read the needle synchronously on the tokio worker while holding store.read(): first a stream-info read that loaded the whole record just to parse its meta, then, for every needle that was not streamed (small, compressed, chunk manifest, image ops), a second full read. For a tiered volume each read is an S3 GET under the store lock, and a writer queued behind it parks every other store reader. The regular-volume read now runs in spawn_blocking. Under the store guard it only resolves a NeedleReadPlan (index lookup, a freshly opened .dat handle or the remote backend, offset, size); the guard is dropped before any needle data I/O. No data-file lease is held across the read either, since a writer waits for one while holding the store write lock. The index size decides the read, as in Go's readNeedle: a HEAD, a ranged read or a needle above the stream threshold reads only its header and meta tail (ReadNeedleMeta) and hands off to StreamingBody or the range path; everything else is read in full once, with its checksum verified. A compressed or manifest needle found by the meta read is then read in full once. The range-from-source read also moves to spawn_blocking. Co-Authored-By: Claude Opus 5.5 (1M context) * volume server: stream needle chunks without the store lock StreamingBody::poll_frame took store.read() and find_volume for every chunk to compare the volume's compaction revision, dup'd the source handle, and allocated a fresh chunk buffer. With -hasSlowRead=false the stream also holds a data-file read lease for its whole life, while a writer waits for that lease under store.write(): the next chunk's store.read() then waits for the writer and the writer for the stream. The per-chunk re-lookup was also wrong. The stream reads a handle opened at plan time, which pins the .dat inode the offset was resolved against; a vacuum commit renames a new file over .dat and leaves that inode untouched. The re-looked-up offset belongs to the new file but was read from the old inode, so a stream whose needle a vacuum moved ended in a checksum error. The pinned offset stays valid, so the check, and with it every store access, is dropped, along with the now unused re_lookup_needle_data_offset and the revision fields of the read plan. The source is shared as an Arc instead of dup'd per chunk, and the chunk buffer is a BytesMut that the blocking read hands back with its result, so its allocation is reclaimed once the previous frame has been written. Co-Authored-By: Claude Opus 5.5 (1M context) * volume server: stop a needle stream once its volume becomes unavailable Taking the store lock out of StreamingBody also dropped its per-chunk unavailable_error() check. With -hasSlowRead a writer can take the data-file lease between chunks, fail its fsync and its truncate, and mark the volume unavailable; the stream then kept serving the rest of the needle from its pinned handle. The volume's io_unavailable reason is now an Arc-shared leaf mutex that the read plan hands to the stream. Each chunk checks it under its data-file lease, where the writer marks it, and fails with the same "volume is unavailable: " error the old check returned. Co-Authored-By: Claude Opus 5.5 (1M context) --------- Co-authored-by: Claude Opus 5.5 (1M context) Co-authored-by: Chris Lu --- seaweed-volume/src/server/handlers.rs | 365 ++++++++++++++++++++------ seaweed-volume/src/storage/store.rs | 11 - seaweed-volume/src/storage/volume.rs | 179 +++---------- 3 files changed, 327 insertions(+), 228 deletions(-) diff --git a/seaweed-volume/src/server/handlers.rs b/seaweed-volume/src/server/handlers.rs index 35059be26..38ef73332 100644 --- a/seaweed-volume/src/server/handlers.rs +++ b/seaweed-volume/src/server/handlers.rs @@ -218,7 +218,9 @@ fn read_needle_for_get( /// A body that streams needle data from the dat file in chunks using pread, /// avoiding loading the entire payload into memory at once. struct StreamingBody { - source: crate::storage::volume::NeedleStreamSource, + /// Opened at plan time: it pins the inode the offset was resolved against, + /// so a vacuum commit mid-stream moves nothing under it. + source: Arc, data_offset: u64, data_size: u32, pos: usize, @@ -226,20 +228,15 @@ struct StreamingBody { data_file_access_control: Arc, hold_read_lock_for_stream: bool, _held_read_lease: Option, - /// Pending result from spawn_blocking, polled to completion. - pending: Option>>, + io_unavailable: Arc, + /// Pending chunk read; it hands `buf` back with the result. + pending: Option)>>, + /// Chunk buffer, reclaimed once the previous frame has been written out. + buf: bytes::BytesMut, /// For download throttling — released on drop. state: Option>, tracked_bytes: i64, - /// Server state used to re-lookup needle offset if compaction occurs during streaming. - server_state: Arc, - /// Volume ID for compaction-revision re-lookup. - volume_id: crate::storage::types::VolumeId, - /// Needle ID for compaction-revision re-lookup. needle_id: crate::storage::types::NeedleId, - /// Compaction revision at the time of the initial read; if the volume's revision - /// changes between chunks, the needle may have moved and we must re-lookup its offset. - compaction_revision: u16, /// Checksum stored in the needle tail. expected_checksum: u32, /// Running checksum over the emitted data. @@ -261,6 +258,10 @@ impl http_body::Body for StreamingBody { std::task::Poll::Pending => return std::task::Poll::Pending, std::task::Poll::Ready(result) => { self.pending = None; + let result = result.map(|(buf, read)| { + self.buf = buf; + read + }); match result { Ok(Ok(chunk)) => { let len = chunk.len(); @@ -308,48 +309,13 @@ impl http_body::Body for StreamingBody { return std::task::Poll::Ready(None); } - // Check if compaction has changed the needle's disk location (Go parity: - // readNeedleDataInto re-reads the needle offset when CompactionRevision changes). - let relookup_result = { - let store = self.server_state.store.read().unwrap(); - if let Some((_, vol)) = store.find_volume(self.volume_id) { - if let Some(e) = vol.unavailable_error() { - return std::task::Poll::Ready(Some(Err(std::io::Error::other(e)))); - } - if vol.super_block.compaction_revision != self.compaction_revision { - // Compaction occurred — re-lookup the needle's data offset - Some(vol.re_lookup_needle_data_offset(self.needle_id)) - } else { - None - } - } else { - None - } - }; - if let Some(result) = relookup_result { - match result { - Ok((new_offset, new_rev)) => { - self.data_offset = new_offset; - self.compaction_revision = new_rev; - } - Err(_) => { - return std::task::Poll::Ready(Some(Err(std::io::Error::new( - std::io::ErrorKind::NotFound, - "needle not found after compaction", - )))); - } - } - } - let chunk_len = std::cmp::min(self.chunk_size, total - self.pos); let file_offset = self.data_offset + self.pos as u64; - - let source_clone = match self.source.clone_for_read() { - Ok(source) => source, - Err(e) => return std::task::Poll::Ready(Some(Err(e))), - }; + let source = self.source.clone(); let data_file_access_control = self.data_file_access_control.clone(); let hold_read_lock_for_stream = self.hold_read_lock_for_stream; + let io_unavailable = self.io_unavailable.clone(); + let mut buf = std::mem::take(&mut self.buf); let handle = tokio::task::spawn_blocking(move || { let _lease = if hold_read_lock_for_stream { @@ -357,9 +323,15 @@ impl http_body::Body for StreamingBody { } else { Some(data_file_access_control.read_lock()) }; - let mut buf = vec![0u8; chunk_len]; - source_clone.read_exact_at(&mut buf, file_offset)?; - Ok::(bytes::Bytes::from(buf)) + // Under the lease: a writer marks the volume holding the write lock. + if let Some(e) = io_unavailable.error() { + return (buf, Err(std::io::Error::other(e))); + } + buf.resize(chunk_len, 0); + let read = source + .read_exact_at(&mut buf, file_offset) + .map(|()| buf.split().freeze()); + (buf, read) }); self.pending = Some(handle); @@ -1571,7 +1543,7 @@ async fn get_or_head_handler_inner( }; let streaming = StreamingBody { - source: info.source, + source: Arc::new(info.source), data_offset: info.data_file_offset, data_size: info.data_size, pos: 0, @@ -1583,13 +1555,12 @@ async fn get_or_head_handler_inner( }, data_file_access_control: info.data_file_access_control, hold_read_lock_for_stream: !state.has_slow_read, + io_unavailable: info.io_unavailable, pending: None, + buf: bytes::BytesMut::new(), state: tracking_state, tracked_bytes, - server_state: state.clone(), - volume_id: info.volume_id, needle_id: info.needle_id, - compaction_revision: info.compaction_revision, expected_checksum: info.checksum, crc: CRC(0), }; @@ -4671,6 +4642,37 @@ mod tests { server.abort(); } + /// A body streaming `data_size` bytes of `file` from offset 0. + fn file_streaming_body( + file: &std::fs::File, + data_size: usize, + chunk_size: usize, + expected_checksum: u32, + ) -> StreamingBody { + StreamingBody { + source: Arc::new(crate::storage::volume::NeedleStreamSource::Local( + file.try_clone().unwrap(), + )), + data_offset: 0, + data_size: data_size as u32, + pos: 0, + chunk_size, + data_file_access_control: Arc::new( + crate::storage::volume::DataFileAccessControl::default(), + ), + hold_read_lock_for_stream: true, + _held_read_lease: None, + io_unavailable: Arc::default(), + pending: None, + buf: bytes::BytesMut::new(), + state: None, + tracked_bytes: 0, + needle_id: NeedleId(7), + expected_checksum, + crc: CRC(0), + } + } + /// Repeated target response headers (e.g. Set-Cookie) must all reach the /// client, as Go's `w.Header().Add` does; only `Server` is dropped. #[tokio::test] @@ -4804,25 +4806,8 @@ mod tests { tmp.write_all(&data).unwrap(); let file = tmp.reopen().unwrap(); - let control = Arc::new(crate::storage::volume::DataFileAccessControl::default()); - let mk_body = |expected_checksum: u32, file: &std::fs::File| StreamingBody { - source: crate::storage::volume::NeedleStreamSource::Local(file.try_clone().unwrap()), - data_offset: 0, - data_size: data.len() as u32, - pos: 0, - chunk_size: 4096, - data_file_access_control: control.clone(), - hold_read_lock_for_stream: true, - _held_read_lease: None, - pending: None, - state: None, - tracked_bytes: 0, - server_state: streaming_test_state(), - volume_id: VolumeId(1), - needle_id: NeedleId(7), - compaction_revision: 0, - expected_checksum, - crc: CRC(0), + let mk_body = |expected_checksum: u32, file: &std::fs::File| { + file_streaming_body(file, data.len(), 4096, expected_checksum) }; // Healthy needle: every byte is delivered. @@ -5129,4 +5114,234 @@ mod tests { let (status, _, _) = send_read(&state, Method::GET, &path, None).await; assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR); } + + /// Each chunk reuses the buffer of the previous frame once that frame has + /// been dropped, instead of allocating a new one. + #[tokio::test] + async fn test_streaming_body_reuses_its_chunk_buffer() { + use std::io::Write; + + const CHUNK: usize = 4096; + let data: Vec = (0..3 * CHUNK + 100).map(|i| (i % 251) as u8).collect(); + let mut tmp = tempfile::NamedTempFile::new().unwrap(); + tmp.write_all(&data).unwrap(); + let file = tmp.reopen().unwrap(); + let mut body = file_streaming_body(&file, data.len(), CHUNK, CRC::new(&data).0); + + let mut got = Vec::new(); + while let Some(frame) = std::future::poll_fn(|cx| { + http_body::Body::poll_frame(std::pin::Pin::new(&mut body), cx) + }) + .await + { + let chunk = frame.unwrap().into_data().unwrap(); + got.extend_from_slice(&chunk); + assert!( + !body.buf.try_reclaim(CHUNK), + "the buffer must be the one the live frame points into" + ); + drop(chunk); + assert!( + body.buf.try_reclaim(CHUNK), + "the chunk buffer did not come back from the read" + ); + } + assert_eq!(got, data); + } + + /// Start a GET and return its body unread. + async fn open_read(state: &Arc, path: &str) -> Body { + use tower::ServiceExt; + let resp = super::super::volume_server::build_public_router(state.clone()) + .oneshot(Request::builder().uri(path).body(Body::empty()).unwrap()) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + resp.into_body() + } + + /// A needle that streams in several chunks. + fn multi_chunk_data() -> Vec { + (0..9 * 1024 * 1024 + 123) + .map(|i| (i % 251) as u8) + .collect() + } + + /// With -hasSlowRead=false a stream holds a data-file lease for its whole + /// life, and a writer waits for that lease under store.write(). Producing + /// a frame must not need the store lock, or the two wait on each other. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_streaming_body_progresses_while_a_writer_holds_the_store() { + use futures::StreamExt; + use std::sync::mpsc; + use std::time::{Duration, Instant}; + + const ID: u64 = 0x6e7a_0501; + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + assert!(!state.has_slow_read); + let data = multi_chunk_data(); + let path = put_test_needle(&state, ID, &data); + + let mut stream = open_read(&state, &path).await.into_data_stream(); + let mut got = stream.next().await.unwrap().unwrap().to_vec(); + assert!( + got.len() < data.len(), + "the needle must span several frames" + ); + + // The writer takes store.write(), then parks on the stream's lease. + let (wrote_tx, wrote_rx) = mpsc::channel(); + std::thread::spawn({ + let state = state.clone(); + move || { + put_test_needle(&state, ID + 1, b"written meanwhile"); + let _ = wrote_tx.send(()); + } + }); + let deadline = Instant::now() + Duration::from_secs(10); + while state.store.try_read().is_ok() { + assert!(Instant::now() < deadline, "the writer never took the store"); + std::thread::sleep(Duration::from_millis(1)); + } + + // Drained off the runtime: a frame stuck on the store lock would block + // its thread, which a tokio timeout cannot interrupt. + let (done_tx, done_rx) = mpsc::channel(); + let rt = tokio::runtime::Handle::current(); + std::thread::spawn(move || { + let rest = rt.block_on(async move { + let mut rest = Vec::new(); + while let Some(chunk) = stream.next().await { + rest.extend_from_slice(&chunk?); + } + Ok::<_, axum::Error>(rest) + }); + let _ = done_tx.send(rest); + }); + let rest = + tokio::task::spawn_blocking(move || done_rx.recv_timeout(Duration::from_secs(10))) + .await + .unwrap() + .expect("the stream stalled behind a writer holding the store"); + got.extend(rest.unwrap()); + assert_eq!(got, data); + + let wrote = tokio::task::spawn_blocking(move || { + wrote_rx.recv_timeout(Duration::from_secs(10)).is_ok() + }) + .await + .unwrap(); + assert!(wrote, "the writer must get the lease once the stream ends"); + } + + /// With -hasSlowRead a writer gets the data-file lease between chunks; if + /// its failed append leaves the volume unavailable, the stream must stop. + #[tokio::test] + async fn test_streaming_body_stops_once_the_volume_becomes_unavailable() { + use futures::StreamExt; + + const ID: u64 = 0x6e7a_0701; + let tmp = tempfile::TempDir::new().unwrap(); + let mut state = volume_test_state(&tmp); + Arc::get_mut(&mut state).unwrap().has_slow_read = true; + let data = multi_chunk_data(); + let path = put_test_needle(&state, ID, &data); + + let mut stream = open_read(&state, &path).await.into_data_stream(); + let first = stream.next().await.unwrap().unwrap(); + assert!( + first.len() < data.len(), + "the needle must span several frames" + ); + + { + let mut store = state.store.write().unwrap(); + let (_, vol) = store.find_volume_mut(VolumeId(1)).unwrap(); + vol.fail_next_fsync_for_test(true); + vol.fail_next_truncate_for_test(true); + let mut n = Needle { + id: NeedleId(ID + 1), + cookie: Cookie(TEST_COOKIE), + data: b"never-landed".to_vec(), + data_size: 12, + ..Needle::default() + }; + store + .write_volume_needle(VolumeId(1), &mut n, true) + .unwrap_err(); + let (_, vol) = store.find_volume_mut(VolumeId(1)).unwrap(); + vol.fail_next_fsync_for_test(false); + vol.fail_next_truncate_for_test(false); + assert!(vol.unavailable_error().is_some()); + } + + match stream.next().await { + Some(Err(err)) => assert!( + err.to_string().contains("volume is unavailable"), + "unexpected stream error: {err}" + ), + other => panic!( + "the stream kept reading an unavailable volume: {:?}", + other.map(|chunk| chunk.map(|c| c.len())) + ), + } + } + + /// A vacuum committed mid-stream must not change the bytes served: the + /// stream's handle pins the inode its offset was resolved against. + #[cfg(unix)] + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_streaming_body_survives_a_vacuum_that_moves_the_needle() { + use futures::StreamExt; + + const FILLER: u64 = 0x6e7a_0601; + const ID: u64 = 0x6e7a_0602; + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + put_test_needle(&state, FILLER, &[7u8; 4096]); + let data = multi_chunk_data(); + let path = put_test_needle(&state, ID, &data); + let mut filler = Needle { + id: NeedleId(FILLER), + cookie: Cookie(TEST_COOKIE), + ..Needle::default() + }; + state + .store + .write() + .unwrap() + .delete_volume_needle(VolumeId(1), &mut filler) + .unwrap(); + let data_offset = || { + let mut plan = state + .store + .read() + .unwrap() + .needle_read_plan(VolumeId(1), NeedleId(ID), false) + .unwrap(); + let mut n = Needle::default(); + plan.read_meta(&mut n).unwrap(); + plan.into_stream_info(&n).data_file_offset + }; + let before = data_offset(); + + let mut stream = open_read(&state, &path).await.into_data_stream(); + let mut got = stream.next().await.unwrap().unwrap().to_vec(); + assert!( + got.len() < data.len(), + "the needle must span several frames" + ); + { + let mut store = state.store.write().unwrap(); + store.compact_volume(VolumeId(1), 0, 0, |_| true).unwrap(); + store.commit_compact_volume(VolumeId(1)).unwrap(); + } + assert_ne!(data_offset(), before, "the vacuum must move the needle"); + + while let Some(chunk) = stream.next().await { + got.extend_from_slice(&chunk.expect("the stream failed across a vacuum")); + } + assert_eq!(got, data); + } } diff --git a/seaweed-volume/src/storage/store.rs b/seaweed-volume/src/storage/store.rs index ea3da0104..7039d553a 100644 --- a/seaweed-volume/src/storage/store.rs +++ b/seaweed-volume/src/storage/store.rs @@ -726,17 +726,6 @@ impl Store { vol.needle_read_plan(id, read_deleted) } - /// Re-lookup a needle's data-file offset after compaction may have moved it. - /// Returns `(new_data_file_offset, current_compaction_revision)`. - pub fn re_lookup_needle_data_offset( - &self, - vid: VolumeId, - needle_id: NeedleId, - ) -> Result<(u64, u16), VolumeError> { - let (_, vol) = self.find_volume(vid).ok_or(VolumeError::NotFound)?; - vol.re_lookup_needle_data_offset(needle_id) - } - /// Write a needle to a volume. With `fsync` the volume flushes its .dat /// before returning, so the caller can ack a durable write. pub fn write_volume_needle( diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index e07e94898..d2991a7bd 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -536,13 +536,6 @@ pub(crate) enum NeedleStreamSource { } impl NeedleStreamSource { - pub(crate) fn clone_for_read(&self) -> io::Result { - match self { - NeedleStreamSource::Local(file) => Ok(NeedleStreamSource::Local(file.try_clone()?)), - NeedleStreamSource::Remote(remote) => Ok(NeedleStreamSource::Remote(remote.clone())), - } - } - pub(crate) fn read_exact_at(&self, buf: &mut [u8], offset: u64) -> io::Result<()> { match self { NeedleStreamSource::Local(file) => read_exact_at(file, buf, offset), @@ -719,6 +712,29 @@ impl DatScanPlan { } } +/// Why a volume refuses all I/O, shared with its in-flight needle streams so +/// they see a mark made after they left the store lock. A leaf lock. +#[derive(Debug, Default)] +pub(crate) struct IoUnavailable(Mutex>); + +impl IoUnavailable { + fn set(&self, reason: String) { + *self.0.lock().unwrap() = Some(reason); + } + + fn is_set(&self) -> bool { + self.0.lock().unwrap().is_some() + } + + pub(crate) fn error(&self) -> Option { + self.0 + .lock() + .unwrap() + .as_ref() + .map(|reason| VolumeError::Unavailable(reason.clone())) + } +} + /// A needle read resolved under a store guard and run after it is released; /// the handle pins the inode, as for `DatScanPlan`. It takes no data-file /// lease: writers wait for one while holding the store write lock. @@ -727,11 +743,10 @@ pub(crate) struct NeedleReadPlan { offset: i64, size: Size, version: Version, - volume_id: VolumeId, needle_id: NeedleId, - compaction_revision: u16, data_file_access_control: Arc, io_errors: Arc, + io_unavailable: Arc, } impl NeedleReadPlan { @@ -794,9 +809,8 @@ impl NeedleReadPlan { data_file_offset, data_size: n.data_size, data_file_access_control: self.data_file_access_control, - volume_id: self.volume_id, + io_unavailable: self.io_unavailable, needle_id: self.needle_id, - compaction_revision: self.compaction_revision, checksum: n.checksum.0, } } @@ -1070,14 +1084,9 @@ pub struct NeedleStreamInfo { pub data_size: u32, /// Per-volume file access lock used to match Go's slow-read behavior. pub data_file_access_control: Arc, - /// Volume ID — used to re-lookup needle offset if compaction occurs during streaming. - pub volume_id: VolumeId, - /// Needle ID — used to re-lookup needle offset if compaction occurs during streaming. + /// Checked before each chunk: the volume can become unavailable mid-stream. + pub(crate) io_unavailable: Arc, pub needle_id: NeedleId, - /// Compaction revision at the time of the initial read. If this changes during - /// streaming, the needle's disk offset must be re-read from the needle map because - /// compaction may have moved the needle to a different location. - pub compaction_revision: u16, /// Checksum stored in the needle tail, verified once the last chunk has /// been read — before that frame is emitted. pub checksum: u32, @@ -1165,7 +1174,7 @@ pub struct Volume { /// Set when a failed recovery leaves the .dat/index pair unverified: all /// I/O is refused and a `.unavailable` marker keeps the volume quarantined /// across restarts. Mirrors Go's ioUnavailable. - io_unavailable: Option, + io_unavailable: Arc, /// Shared flag from the parent DiskLocation indicating low disk space. /// Matches Go's `v.location.isDiskSpaceLow` checked in `IsReadOnly()`. @@ -1268,7 +1277,7 @@ impl Volume { }, no_write_or_delete: false, no_write_can_delete: false, - io_unavailable: None, + io_unavailable: Arc::default(), location_disk_space_low: Arc::new(AtomicBool::new(false)), last_modified_ts_seconds: 0, last_append_at_ns: 0, @@ -1314,7 +1323,7 @@ impl Volume { super_block: SuperBlock::default(), no_write_or_delete: false, no_write_can_delete: false, - io_unavailable: None, + io_unavailable: Arc::default(), location_disk_space_low: Arc::new(AtomicBool::new(false)), last_modified_ts_seconds: 0, last_append_at_ns: 0, @@ -2210,11 +2219,10 @@ impl Volume { offset: nv.offset.to_actual_offset(), size, version: self.version(), - volume_id: self.id, needle_id: id, - compaction_revision: self.super_block.compaction_revision, data_file_access_control: self.data_file_access_control.clone(), io_errors: self.io_errors.clone(), + io_unavailable: self.io_unavailable.clone(), }) } @@ -2237,39 +2245,6 @@ impl Volume { } } - /// Re-lookup a needle's data-file offset after compaction may have moved it. - /// - /// Returns `(new_data_file_offset, current_compaction_revision)` or an error - /// if the needle is no longer present / has been deleted. - /// - /// This matches Go's `readNeedleDataInto` behaviour: when the volume's - /// `CompactionRevision` changes between streaming chunks, the needle offset - /// is re-read from the needle map because compaction may have relocated it. - pub fn re_lookup_needle_data_offset( - &self, - needle_id: NeedleId, - ) -> Result<(u64, u16), VolumeError> { - let nm = self.nm_or_not_found()?; - let nv = nm.get(needle_id)?.ok_or(VolumeError::NotFound)?; - if nv.offset.is_zero() { - return Err(VolumeError::NotFound); - } - if nv.size.is_deleted() { - return Err(VolumeError::Deleted); - } - - let offset = nv.offset.to_actual_offset(); - let version = self.version(); - - let data_file_offset = if version == VERSION_1 { - offset as u64 + NEEDLE_HEADER_SIZE as u64 - } else { - offset as u64 + NEEDLE_HEADER_SIZE as u64 + 4 // skip DataSize (4 bytes) - }; - - Ok((data_file_offset, self.super_block.compaction_revision)) - } - // ---- Write ---- /// Write a needle to the volume (synchronous path). @@ -2884,14 +2859,14 @@ impl Volume { pub fn is_read_only(&self) -> bool { self.no_write_or_delete || self.no_write_can_delete - || self.io_unavailable.is_some() + || self.io_unavailable.is_set() || self.location_disk_space_low.load(Ordering::Relaxed) } /// Mirrors Go's ReadOnlyReasons: `no_write_or_delete` already covers the /// io_unavailable quarantine. pub fn read_only_reasons(&self) -> (bool, bool, bool, bool) { - let no_write_or_delete = self.no_write_or_delete || self.io_unavailable.is_some(); + let no_write_or_delete = self.no_write_or_delete || self.io_unavailable.is_set(); let disk_space_low = self.location_disk_space_low.load(Ordering::Relaxed); ( no_write_or_delete || self.no_write_can_delete || disk_space_low, @@ -2904,9 +2879,7 @@ impl Volume { /// The reason the volume refuses all I/O, when a failed recovery left the /// .dat/index pair unverified. Mirrors Go's unavailableError. pub fn unavailable_error(&self) -> Option { - self.io_unavailable - .as_ref() - .map(|reason| VolumeError::Unavailable(reason.clone())) + self.io_unavailable.error() } /// Fail closed after a recovery could not return the volume to a verified @@ -2914,7 +2887,7 @@ impl Volume { /// reload stays unavailable until an operator verifies the volume. fn mark_io_unavailable(&mut self, reason: String) { self.no_write_or_delete = true; - self.io_unavailable = Some(reason.clone()); + self.io_unavailable.set(reason.clone()); self.mark_io_quarantined(); if let Err(e) = self.persist_unavailable(&reason) { warn!( @@ -2956,7 +2929,7 @@ impl Volume { return; }; self.no_write_or_delete = true; - self.io_unavailable = Some(reason.trim().to_string()); + self.io_unavailable.set(reason.trim().to_string()); self.mark_io_quarantined(); warn!( volume_id = self.id.0, @@ -8207,83 +8180,7 @@ mod tests { } #[test] - fn test_compaction_revision_relookup() { - // Verifies that re_lookup_needle_data_offset returns the correct data offset - // and compaction revision, and that after compaction the offset changes. - let tmp = TempDir::new().unwrap(); - let dir = tmp.path().to_str().unwrap(); - let mut v = make_test_volume(dir); - - // Write two needles - let mut n1 = Needle { - id: NeedleId(1), - cookie: Cookie(0xAABBCCDD), - data: b"first-needle-data".to_vec(), - data_size: 17, - ..Needle::default() - }; - v.write_needle(&mut n1, true, false).unwrap(); - - let mut n2 = Needle { - id: NeedleId(2), - cookie: Cookie(0x11223344), - data: b"second-needle-data".to_vec(), - data_size: 18, - ..Needle::default() - }; - v.write_needle(&mut n2, true, false).unwrap(); - - // Get initial revision and offset for needle 1 - let initial_rev = v.super_block.compaction_revision; - let (initial_offset, rev) = v.re_lookup_needle_data_offset(NeedleId(1)).unwrap(); - assert_eq!(rev, initial_rev); - assert!(initial_offset > 0, "data offset should be positive"); - - // Delete needle 2 so compaction removes it - let mut del_n2 = Needle { - id: NeedleId(2), - cookie: Cookie(0x11223344), - ..Needle::default() - }; - v.delete_needle(&mut del_n2).unwrap(); - - // Compact the volume — this increments compaction_revision and may move needles - v.compact_by_index(0, 0, |_| true).unwrap(); - v.commit_compact().unwrap(); - - // After compaction, the revision should have changed - let new_rev = v.super_block.compaction_revision; - assert_eq!( - new_rev, - initial_rev + 1, - "compaction should increment revision" - ); - - // Re-lookup needle 1 — should still be found with the new revision - let (new_offset, relookup_rev) = v.re_lookup_needle_data_offset(NeedleId(1)).unwrap(); - assert_eq!(relookup_rev, new_rev); - assert!(new_offset > 0, "data offset should still be positive"); - - // The data should still be readable correctly after compaction - let mut read_n1 = Needle { - id: NeedleId(1), - ..Needle::default() - }; - v.read_needle(&mut read_n1).unwrap(); - assert_eq!(read_n1.data, b"first-needle-data"); - - // Deleted needle should not be found - let result = v.re_lookup_needle_data_offset(NeedleId(2)); - assert!( - result.is_err(), - "deleted needle should not be found after compaction" - ); - } - - #[test] - fn test_stream_info_includes_compaction_revision() { - // Verifies that NeedleStreamInfo carries the volume's compaction revision - // so that StreamingBody can detect when compaction has occurred. + fn test_stream_info_locates_needle_data() { let tmp = TempDir::new().unwrap(); let dir = tmp.path().to_str().unwrap(); let mut v = make_test_volume(dir); @@ -8309,9 +8206,7 @@ mod tests { plan.read_meta(&mut read_n).unwrap(); let info = plan.into_stream_info(&read_n); - assert_eq!(info.volume_id, VolumeId(1)); assert_eq!(info.needle_id, NeedleId(42)); - assert_eq!(info.compaction_revision, v.super_block.compaction_revision); assert_eq!(info.data_size, data.len() as u32); assert!(info.data_file_offset > 0); }