mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-05 22:12:04 +02:00
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) <noreply@anthropic.com> * 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) <noreply@anthropic.com> * 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: <reason>" error the old check returned. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com>
This commit is contained in:
3 files changed
+327
-228
No files matched your search
@@ -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<crate::storage::volume::NeedleStreamSource>,
|
||||
data_offset: u64,
|
||||
data_size: u32,
|
||||
pos: usize,
|
||||
@@ -226,20 +228,15 @@ struct StreamingBody {
|
||||
data_file_access_control: Arc<crate::storage::volume::DataFileAccessControl>,
|
||||
hold_read_lock_for_stream: bool,
|
||||
_held_read_lease: Option<crate::storage::volume::DataFileReadLease>,
|
||||
/// Pending result from spawn_blocking, polled to completion.
|
||||
pending: Option<tokio::task::JoinHandle<Result<bytes::Bytes, std::io::Error>>>,
|
||||
io_unavailable: Arc<crate::storage::volume::IoUnavailable>,
|
||||
/// Pending chunk read; it hands `buf` back with the result.
|
||||
pending: Option<tokio::task::JoinHandle<(bytes::BytesMut, std::io::Result<bytes::Bytes>)>>,
|
||||
/// Chunk buffer, reclaimed once the previous frame has been written out.
|
||||
buf: bytes::BytesMut,
|
||||
/// For download throttling — released on drop.
|
||||
state: Option<Arc<VolumeServerState>>,
|
||||
tracked_bytes: i64,
|
||||
/// Server state used to re-lookup needle offset if compaction occurs during streaming.
|
||||
server_state: Arc<VolumeServerState>,
|
||||
/// 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, std::io::Error>(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<u8> = (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<VolumeServerState>, 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<u8> {
|
||||
(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);
|
||||
}
|
||||
}
|
||||
@@ -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(
|
||||
|
||||
@@ -536,13 +536,6 @@ pub(crate) enum NeedleStreamSource {
|
||||
}
|
||||
|
||||
impl NeedleStreamSource {
|
||||
pub(crate) fn clone_for_read(&self) -> io::Result<Self> {
|
||||
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<Option<String>>);
|
||||
|
||||
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<VolumeError> {
|
||||
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<DataFileAccessControl>,
|
||||
io_errors: Arc<IoErrorTracker>,
|
||||
io_unavailable: Arc<IoUnavailable>,
|
||||
}
|
||||
|
||||
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<DataFileAccessControl>,
|
||||
/// 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<IoUnavailable>,
|
||||
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<String>,
|
||||
io_unavailable: Arc<IoUnavailable>,
|
||||
|
||||
/// 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<VolumeError> {
|
||||
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);
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user