diff --git a/Cargo.toml b/Cargo.toml index 4ade9061..0dd68992 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -92,6 +92,10 @@ harness = false name = "oracle_batch" harness = false +[[bench]] +name = "xact_spill" +harness = false + [dev-dependencies] tempfile = "3" # Enable paused time for batch deadline tests diff --git a/bench/README.md b/bench/README.md index 567f82e3..d0e940e3 100644 --- a/bench/README.md +++ b/bench/README.md @@ -1,5 +1,27 @@ # Replication benchmarks +## Transaction spill + +Run buffer admission, spill encoding, commit drain, and cleanup without servers: + +```sh +cargo bench --bench xact_spill -- --buffer-bytes 1048576 --rows 7821 +cargo bench --bench xact_spill -- --buffer-bytes 1048576 --rows 781 --transactions 80 +cargo bench --bench xact_spill -- --rows 111108 --transactions 2 +``` + +Use `--dir` to select filesystem; default `target`. Avoid tmpfs for disk benchmarks. +Defaults use 1 KiB text payloads and three repeats. Verify every row's payload, +xid, LSN order, descriptor, and file cleanup. Write timing includes allocation +and admission; read timing includes reader setup, verified drain, and unlink. +Report spill bytes alongside throughput: final memory tails can avoid disk. +Compare equal spill bytes when isolating I/O improvements. + +Files follow disposable-spill semantics: no fsync, immediate replay, then unlink. +Results include page cache and codec costs; they do not measure sustained device +bandwidth or PostgreSQL-to-ClickHouse lag. Buffer peak reports post-admission +payload estimates; use process RSS to measure allocator and I/O buffer overhead. + ## Local runs Start PostgreSQL, ClickHouse, and walshadow with SQL runtime config installed diff --git a/benches/xact_spill.rs b/benches/xact_spill.rs new file mode 100644 index 00000000..34c41add --- /dev/null +++ b/benches/xact_spill.rs @@ -0,0 +1,143 @@ +//! Transaction buffering and verified commit drain, without PostgreSQL or ClickHouse. +//! `cargo bench --bench xact_spill -- --buffer-bytes 1048576` + +use std::path::PathBuf; +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use clap::Parser; +use walrus::pg::walparser::RelFileNode; +use walshadow::heap_decoder::{ColumnValue, DecodedHeap, DecodedTuple, DescribedHeap, HeapOp}; +use walshadow::schema::{RelDescriptor, RelName, ReplIdent}; +use walshadow::xact_buffer::{XactBuffer, XactBufferConfig}; + +#[global_allocator] +static GLOBAL: mimalloc::MiMalloc = mimalloc::MiMalloc; + +#[derive(Debug, Parser)] +struct Args { + #[arg(long, default_value_t = 7812)] + rows: usize, + #[arg(long, default_value_t = 8)] + transactions: u32, + #[arg(long, default_value_t = 1024)] + payload_bytes: usize, + #[arg(long, default_value_t = 64 << 20)] + buffer_bytes: usize, + #[arg(long, default_value_t = 3)] + repeats: usize, + #[arg(long, default_value = "target")] + dir: PathBuf, +} + +fn main() -> anyhow::Result<()> { + if std::env::args().any(|arg| arg == "--list") { + return Ok(()); + } + let args = Args::parse_from(std::env::args().filter(|arg| arg != "--bench")); + anyhow::ensure!(args.rows > 0 && args.transactions > 0 && args.repeats > 0); + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build()?; + runtime.block_on(run(args)) +} + +async fn run(args: Args) -> anyhow::Result<()> { + let descriptor = Arc::new(RelDescriptor { + rfn: RelFileNode { + spc_node: 1663, + db_node: 5, + rel_node: 16385, + }, + oid: 16385, + toast_oid: 0, + namespace_oid: 2200, + rel_name: RelName::new("public", "spill_bench"), + kind: 'r', + persistence: 'p', + replident: ReplIdent::Default { pk_attnums: None }, + attributes: Vec::new(), + }); + let payload = "x".repeat(args.payload_bytes); + println!("{args:?}"); + for repeat in 0..args.repeats { + let dir = tempfile::tempdir_in(&args.dir)?; + let mut config = XactBufferConfig::new(dir.path().to_path_buf()); + config.xact_buffer_max = args.buffer_bytes; + let mut buffer = XactBuffer::new(config)?; + let mut write_time = Duration::ZERO; + let mut read_time = Duration::ZERO; + let mut spill_bytes = 0; + let mut peak_memory = 0; + for xid in 1..=args.transactions { + let start = Instant::now(); + for row in 0..args.rows { + buffer + .on_heap(DescribedHeap { + decoded: DecodedHeap { + rfn: descriptor.rfn, + xid, + source_lsn: u64::from(xid) * (args.rows as u64 + 1) + row as u64, + op: HeapOp::Insert, + new: Some(DecodedTuple { + columns: vec![Some(ColumnValue::Text(payload.clone()))], + partial: false, + }), + old: None, + }, + descriptor: descriptor.clone(), + descriptor_valid_from: 1, + }) + .await?; + peak_memory = peak_memory.max(buffer.stats().bytes_in_memory); + } + write_time += start.elapsed(); + spill_bytes += buffer.stats().spill_bytes_active; + let start = Instant::now(); + let commit_lsn = u64::from(xid + 1) * (args.rows as u64 + 1); + let mut drain = buffer + .drain_committed(xid, 0, commit_lsn, &[], false) + .await?; + let mut rows = 0; + while let Some(batch) = drain.next_batch(1024, 1 << 20, None).await? { + for heap in batch.heaps { + assert_eq!(heap.decoded.xid, xid); + assert_eq!( + heap.decoded.source_lsn, + u64::from(xid) * (args.rows as u64 + 1) + rows + ); + assert_eq!(heap.descriptor.as_ref(), descriptor.as_ref()); + assert_eq!(heap.descriptor_valid_from, 1); + let tuple = heap.decoded.new.unwrap(); + assert_eq!(tuple.columns.len(), 1); + assert!( + matches!(&tuple.columns[0], Some(ColumnValue::Text(s)) if s == &payload) + ); + rows += 1; + } + } + assert_eq!(rows, args.rows as u64); + drain.finish().await?; + buffer.resume_safe_lsn(walshadow::pos::Pos::new(commit_lsn)); + read_time += start.elapsed(); + } + let rows = args.rows as f64 * f64::from(args.transactions); + println!( + "repeat={} rows={} spill_bytes={} evictions={} peak_buffer_bytes={} write_s={:.6} read_s={:.6} rows_per_s={:.0} spill_write_mib_s={:.2} spill_read_mib_s={:.2}", + repeat + 1, + rows as u64, + spill_bytes, + buffer.stats().spill_evictions_total, + peak_memory, + write_time.as_secs_f64(), + read_time.as_secs_f64(), + rows / (write_time + read_time).as_secs_f64(), + spill_bytes as f64 / (1 << 20) as f64 / write_time.as_secs_f64(), + spill_bytes as f64 / (1 << 20) as f64 / read_time.as_secs_f64(), + ); + assert_eq!(buffer.stats().bytes_in_memory, 0); + assert_eq!(buffer.stats().spill_bytes_active, 0); + assert_eq!(std::fs::read_dir(buffer.scratch_dir())?.count(), 0); + } + Ok(()) +} diff --git a/src/xact/spill.rs b/src/xact/spill.rs index 5c639811..318959ee 100644 --- a/src/xact/spill.rs +++ b/src/xact/spill.rs @@ -54,7 +54,8 @@ use std::sync::Arc; use thiserror::Error; use tokio::fs::{File, OpenOptions}; -use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::io::AsyncWriteExt; +use tokio::task::JoinHandle; use walrus::pg::walparser::RelFileNode; use crate::catalog::desc_log::{decode_descriptor_bytes, encode_descriptor_bytes}; @@ -367,8 +368,9 @@ impl SpillStore { .open(&path) .await?; file.write_all(&header).await?; + file.flush().await?; Ok(SpillWriter { - file, + file: Arc::new(file.into_std().await), path, byte_count: header.len() as u64, dict: HashMap::new(), @@ -536,13 +538,15 @@ impl BodySpoolWriter { } pub struct SpillWriter { - file: File, + file: Arc, path: PathBuf, byte_count: u64, /// Descriptor dictionary ids by log identity, assigned in first-use order dict: HashMap<(RelFileNode, u64), u32>, } +const SPILL_IO_BUFFER: usize = 256 << 10; + impl SpillWriter { pub fn path(&self) -> &Path { &self.path @@ -554,47 +558,73 @@ impl SpillWriter { pub async fn write(&mut self, entry: &SpillEntry) -> Result<()> { let mut body = Vec::with_capacity(128); + self.encode(entry, &mut body); + wait_write(self.spawn_write(body)).await?; + Ok(()) + } + + /// Encode next buffer while previous one writes on blocking pool + pub async fn write_batch(&mut self, entries: Vec) -> Result<()> { + let mut body = Vec::with_capacity(SPILL_IO_BUFFER); + let mut pending = None; + for entry in entries { + self.encode(&entry, &mut body); + if body.len() >= SPILL_IO_BUFFER { + let spare = match pending.take() { + Some(write) => wait_write(write).await?, + None => Vec::with_capacity(SPILL_IO_BUFFER), + }; + pending = Some(self.spawn_write(std::mem::replace(&mut body, spare))); + } + } + if let Some(write) = pending { + wait_write(write).await?; + } + wait_write(self.spawn_write(body)).await?; + Ok(()) + } + + fn spawn_write(&mut self, body: Vec) -> JoinHandle>> { + self.byte_count += body.len() as u64; + let file = self.file.clone(); + tokio::task::spawn_blocking(move || { + use std::io::Write; + (&*file).write_all(&body)?; + Ok(body) + }) + } + + fn encode(&mut self, entry: &SpillEntry, body: &mut Vec) { match entry { SpillEntry::Heap(h) => { let key = (h.descriptor.rfn, h.descriptor_valid_from); let next_id = self.dict.len() as u32; let id = *self.dict.entry(key).or_insert(next_id); if id == next_id { - frame_into(&mut body, TAG_DESCRIPTOR, |out| { + frame_into(body, TAG_DESCRIPTOR, |out| { push_u64(out, h.descriptor_valid_from); encode_descriptor_bytes(out, &h.descriptor); }); } - frame_into(&mut body, TAG_HEAP, |out| { + frame_into(body, TAG_HEAP, |out| { push_u32(out, id); encode_heap_into(out, &h.decoded); }); } - SpillEntry::Chunk(c) => { - frame_into(&mut body, TAG_CHUNK, |out| encode_chunk_into(out, c)) - } - SpillEntry::ToastDelete(d) => frame_into(&mut body, TAG_TOAST_DELETE, |out| { + SpillEntry::Chunk(c) => frame_into(body, TAG_CHUNK, |out| encode_chunk_into(out, c)), + SpillEntry::ToastDelete(d) => frame_into(body, TAG_TOAST_DELETE, |out| { encode_toast_delete_into(out, d) }), - SpillEntry::Raw(r) => frame_into(&mut body, TAG_RAW, |out| encode_raw_into(out, r)), + SpillEntry::Raw(r) => frame_into(body, TAG_RAW, |out| encode_raw_into(out, r)), } - self.file.write_all(&body).await?; - self.byte_count += body.len() as u64; - Ok(()) } - /// Flush + close, return a reader at the start. Caller drives `next()` to + /// Close, return a reader at the start. Caller drives `next()` to /// `Ok(None)` then `unlink()`. No fsync: restart wipes spill unread - pub async fn finish(mut self) -> Result { - self.file.flush().await?; + pub async fn finish(self) -> Result { drop(self.file); - let file = OpenOptions::new().read(true).open(&self.path).await?; - Ok(SpillReader { - file, - path: self.path, - header_checked: false, - dict: Vec::new(), - }) + let chunk_size = self.byte_count.min(SPILL_IO_BUFFER as u64) as usize; + SpillReader::open(self.path, chunk_size).await } /// Abort path: drop the file unread @@ -603,6 +633,12 @@ impl SpillWriter { } } +async fn wait_write(write: JoinHandle>>) -> Result> { + let mut body = write.await.map_err(io::Error::other)??; + body.clear(); + Ok(body) +} + /// `[tag][u32 len LE][body]`, length back-patched after `f` appends in place fn frame_into(out: &mut Vec, tag: u8, f: impl FnOnce(&mut Vec)) { out.push(tag); @@ -616,7 +652,7 @@ fn frame_into(out: &mut Vec, tag: u8, f: impl FnOnce(&mut Vec)) { /// Drop `file`, remove `path`, tolerating already-gone (abort races, /// crash-cleanup re-runs) -async fn unlink_file(file: File, path: &Path) -> Result<()> { +async fn unlink_file(file: T, path: &Path) -> Result<()> { drop(file); match tokio::fs::remove_file(path).await { Ok(()) => Ok(()), @@ -626,8 +662,12 @@ async fn unlink_file(file: File, path: &Path) -> Result<()> { } pub struct SpillReader { - file: File, path: PathBuf, + buf: Vec, + pos: usize, + chunk_size: usize, + /// Next chunk reads on blocking pool while current chunk decodes + next_chunk: Option, /// Lazy header check: first `next()` verifies [`SPILL_MAGIC`] + version, /// so a stale on-disk spill fails cleanly with [`SpillError::Format`] header_checked: bool, @@ -637,21 +677,57 @@ pub struct SpillReader { } impl SpillReader { + async fn open(path: PathBuf, chunk_size: usize) -> Result { + let file = File::open(&path).await?.into_std().await; + Ok(Self { + path, + buf: Vec::new(), + pos: 0, + chunk_size, + next_chunk: Some(read_chunk(file, Vec::new(), chunk_size)), + header_checked: false, + dict: Vec::new(), + }) + } + + pub(crate) fn buffered_bytes(&self) -> usize { + 2 * self.chunk_size + } + + /// Buffer at least `n` unread bytes, false on EOF first + async fn fill(&mut self, n: usize) -> Result { + while self.buf.len() - self.pos < n { + let Some(pending) = self.next_chunk.take() else { + return Ok(false); + }; + let (file, mut chunk) = pending.await.map_err(io::Error::other)??; + if chunk.is_empty() { + return Ok(false); + } + if self.pos == self.buf.len() { + std::mem::swap(&mut self.buf, &mut chunk); + } else { + self.buf.drain(..self.pos); + self.buf.extend_from_slice(&chunk); + } + self.pos = 0; + self.next_chunk = Some(read_chunk(file, chunk, self.chunk_size)); + } + Ok(true) + } + pub fn path(&self) -> &Path { &self.path } async fn check_header(&mut self) -> Result<()> { - let mut buf = [0u8; 4]; - if let Err(e) = self.file.read_exact(&mut buf).await { - return Err(match e.kind() { - io::ErrorKind::UnexpectedEof => SpillError::Format { - offset: 0, - detail: "spill file shorter than 4-byte header".into(), - }, - _ => e.into(), + if !self.fill(4).await? { + return Err(SpillError::Format { + offset: 0, + detail: "spill file shorter than 4-byte header".into(), }); } + let buf = &self.buf[..4]; if buf[..2] != SPILL_MAGIC { return Err(SpillError::Format { offset: 0, @@ -665,6 +741,7 @@ impl SpillReader { detail: format!("unsupported spill version {version}, expected {SPILL_VERSION}"), }); } + self.pos = 4; self.header_checked = true; Ok(()) } @@ -676,19 +753,23 @@ impl SpillReader { self.check_header().await?; } loop { - let mut tag_buf = [0u8; 1]; - match self.file.read_exact(&mut tag_buf).await { - Ok(_) => {} - Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => return Ok(None), - Err(e) => return Err(e.into()), + if !self.fill(5).await? { + if self.pos == self.buf.len() { + return Ok(None); + } + return Err(io::Error::from(io::ErrorKind::UnexpectedEof).into()); + } + let tag = self.buf[self.pos]; + let len_bytes = self.buf[self.pos + 1..self.pos + 5].try_into().unwrap(); + let len = u32::from_le_bytes(len_bytes) as usize; + if !self.fill(5 + len).await? { + return Err(io::Error::from(io::ErrorKind::UnexpectedEof).into()); } - let mut len_buf = [0u8; 4]; - self.file.read_exact(&mut len_buf).await?; - let len = u32::from_le_bytes(len_buf) as usize; - let mut body = vec![0u8; len]; - self.file.read_exact(&mut body).await?; - let mut cur = Cursor::new(&body); - let entry = match tag_buf[0] { + let start = self.pos + 5; + self.pos = start + len; + let body = &self.buf[start..self.pos]; + let mut cur = Cursor::new(body); + let entry = match tag { TAG_HEAP => { let dict_id = cur.u32()? as usize; let h = decode_heap(&mut cur)?; @@ -740,10 +821,22 @@ impl SpillReader { } pub async fn unlink(self) -> Result<()> { - unlink_file(self.file, &self.path).await + unlink_file(self.next_chunk, &self.path).await } } +type ChunkRead = JoinHandle)>>; + +fn read_chunk(mut file: std::fs::File, mut chunk: Vec, size: usize) -> ChunkRead { + tokio::task::spawn_blocking(move || { + use std::io::Read; + chunk.resize(size, 0); + let n = file.read(&mut chunk)?; + chunk.truncate(n); + Ok((file, chunk)) + }) +} + // ── encoding ──────────────────────────────────────────────────────── pub(crate) fn push_u8(out: &mut Vec, v: u8) { @@ -1348,6 +1441,70 @@ mod tests { } } + #[tokio::test(flavor = "current_thread")] + async fn batch_round_trip_across_buffers_and_descriptor_reuse() { + let tmp = tempdir().unwrap(); + let store = SpillStore::new(tmp.path().to_path_buf()).unwrap(); + let mut writer = store.writer(42, 0x1000).await.unwrap(); + let mut expected = Vec::new(); + for i in 0..600 { + let mut heap = sample_heap(42, 0x2000 + i); + heap.decoded.new.as_mut().unwrap().columns = vec![Some(ColumnValue::Bytea(vec![ + i as u8; + if i == 300 { + SPILL_IO_BUFFER + 17 + } else { + 1021 + } + ]))]; + heap.descriptor_valid_from += i % 2; + expected.push(SpillEntry::Heap(Box::new(heap))); + } + writer.write_batch(Vec::new()).await.unwrap(); + for batch in expected.chunks(137) { + writer.write_batch(batch.to_vec()).await.unwrap(); + } + let bytes = writer.byte_count(); + assert_eq!(writer.dict.len(), 2); + let mut reader = writer.finish().await.unwrap(); + assert_eq!(std::fs::metadata(reader.path()).unwrap().len(), bytes); + for entry in expected { + assert_eq!(reader.next().await.unwrap(), Some(entry)); + } + assert!(reader.next().await.unwrap().is_none()); + reader.unlink().await.unwrap(); + } + + #[tokio::test(flavor = "current_thread")] + async fn buffered_reader_rejects_truncated_frame() { + let tmp = tempdir().unwrap(); + let store = SpillStore::new(tmp.path().to_path_buf()).unwrap(); + for remaining in [1, 3, 5, 31] { + let mut writer = store.writer(42, remaining).await.unwrap(); + writer + .write_batch(vec![SpillEntry::Chunk(sample_chunk( + 99, + 0, + 0x2000, + &vec![7; SPILL_IO_BUFFER + 1], + ))]) + .await + .unwrap(); + let mut reader = writer.finish().await.unwrap(); + let file = OpenOptions::new() + .write(true) + .open(reader.path()) + .await + .unwrap(); + file.set_len(4 + remaining).await.unwrap(); + assert!(matches!( + reader.next().await, + Err(SpillError::Io(e)) if e.kind() == io::ErrorKind::UnexpectedEof + )); + reader.unlink().await.unwrap(); + } + } + #[tokio::test(flavor = "current_thread")] async fn round_trip_heap_and_chunk() { let tmp = tempdir().unwrap(); @@ -1623,12 +1780,7 @@ mod tests { bytes.extend_from_slice(&SPILL_VERSION.to_le_bytes()); bytes.extend_from_slice(&[255u8, 0u8, 0u8, 0u8, 0u8]); tokio::fs::write(&path, &bytes).await.unwrap(); - let mut r = SpillReader { - file: OpenOptions::new().read(true).open(&path).await.unwrap(), - path, - header_checked: false, - dict: Vec::new(), - }; + let mut r = SpillReader::open(path, SPILL_IO_BUFFER).await.unwrap(); let err = r.next().await.expect_err("must error on bad tag"); match err { SpillError::Format { detail, .. } => { @@ -1648,12 +1800,7 @@ mod tests { bytes.extend_from_slice(&SPILL_MAGIC); bytes.extend_from_slice(&5u16.to_le_bytes()); tokio::fs::write(&path, &bytes).await.unwrap(); - let mut r = SpillReader { - file: OpenOptions::new().read(true).open(&path).await.unwrap(), - path, - header_checked: false, - dict: Vec::new(), - }; + let mut r = SpillReader::open(path, SPILL_IO_BUFFER).await.unwrap(); let err = r.next().await.expect_err("must error on old version"); match err { SpillError::Format { detail, .. } => { @@ -1670,12 +1817,7 @@ mod tests { tokio::fs::write(&path, &[0u8, 0u8, 0u8, 0u8]) .await .unwrap(); - let mut r = SpillReader { - file: OpenOptions::new().read(true).open(&path).await.unwrap(), - path, - header_checked: false, - dict: Vec::new(), - }; + let mut r = SpillReader::open(path, SPILL_IO_BUFFER).await.unwrap(); let err = r.next().await.expect_err("must error on missing magic"); match err { SpillError::Format { detail, .. } => assert!(detail.contains("bad magic"), "{detail}"), diff --git a/src/xact/xact_buffer.rs b/src/xact/xact_buffer.rs index acdfeefe..a1f120cc 100644 --- a/src/xact/xact_buffer.rs +++ b/src/xact/xact_buffer.rs @@ -27,7 +27,8 @@ //! //! Once `memory_used > config.xact_buffer_max`, flush the largest //! in-memory xact to a [`SpillWriter`]; the xact stays open and later -//! records append to the file. Mirrors PG `ReorderBufferLargestTXN` +//! records accumulate another memory-budgeted tail before eviction. +//! Mirrors PG `ReorderBufferLargestTXN` //! (`src/backend/replication/logical/reorderbuffer.c`). //! //! Drain: spilled entries first (older), then in-mem. Eviction always @@ -1129,19 +1130,10 @@ impl XactBuffer { let sz = approximate_size(&entry); let raw_sz = matches!(&entry, SpillEntry::Raw(_)).then_some(sz as u64); let st = self.state_for(xid, first_lsn); - if let Some(spill) = st.spill.as_mut() { - // Already spilling: append straight to disk - spill.write(&entry).await?; - let bc = spill.byte_count(); - let prev = std::mem::replace(&mut st.spill_bytes, bc); - self.stats.spill_bytes_active += bc - prev; - self.stats.raw_stash_bytes_spill += raw_sz.unwrap_or(0); - } else { - st.in_mem.push(entry); - st.in_mem_bytes += sz; - self.bytes_in_memory += sz; - self.stats.raw_stash_bytes_mem += raw_sz.unwrap_or(0); - } + st.in_mem.push(entry); + st.in_mem_bytes += sz; + self.bytes_in_memory += sz; + self.stats.raw_stash_bytes_mem += raw_sz.unwrap_or(0); self.stats.bytes_in_memory = self.bytes_in_memory as u64; self.maybe_evict().await?; Ok(()) @@ -1156,8 +1148,6 @@ impl XactBuffer { .max_by_key(|(_, s)| s.in_mem_bytes) .map(|(xid, _)| *xid); let Some(xid) = largest else { - // All active xacts already on disk; caller pushing into - // spilled xacts faster than budget allows break; }; self.evict_xact(xid).await?; @@ -1174,9 +1164,7 @@ impl XactBuffer { let writer = st.spill.as_mut().unwrap(); let drained: Vec = std::mem::take(&mut st.in_mem); let freed = std::mem::take(&mut st.in_mem_bytes); - for entry in drained { - writer.write(&entry).await?; - } + writer.write_batch(drained).await?; let bc = writer.byte_count(); let new_spill_bytes = bc - st.spill_bytes; st.spill_bytes = bc; @@ -1541,6 +1529,9 @@ impl MergeSource { for e in &in_mem { gauge.add(approximate_size(e)); } + if let Some(r) = &reader { + gauge.add(r.buffered_bytes()); + } let mut src = Self { head: None, reader, @@ -3660,6 +3651,44 @@ mod tests { } } + #[tokio::test(flavor = "current_thread")] + async fn repeated_eviction_preserves_spilled_prefix_and_memory_tail() { + let tmp = tempdir().unwrap(); + let mut config = cfg(tmp.path().to_path_buf()); + let row_bytes = approximate_size(&SpillEntry::Heap(Box::new(heap_with_value(7, 1, 128)))); + config.xact_buffer_max = 4 * row_bytes; + let mut buffer = XactBuffer::new(config).unwrap(); + for lsn in 1..=12 { + for xid in [7, 8] { + buffer + .on_heap(heap_with_value(xid, 2 * lsn + u64::from(xid == 8), 128)) + .await + .unwrap(); + assert!(buffer.bytes_in_memory <= 4 * row_bytes); + } + } + assert!(buffer.stats().spill_evictions_total > 2); + assert!( + buffer + .inflight + .values() + .any(|s| s.spill.is_some() && !s.in_mem.is_empty()) + ); + let mut drain = buffer + .drain_committed(7, 0, 100, &[8], false) + .await + .unwrap(); + assert_eq!(buffer.stats().bytes_in_memory, 0); + let mut lsns = Vec::new(); + while let Some(batch) = drain.next_batch(3, usize::MAX, None).await.unwrap() { + lsns.extend(batch.heaps.into_iter().map(|heap| heap.decoded.source_lsn)); + } + assert_eq!(lsns, (2..=25).collect::>()); + drain.finish().await.unwrap(); + assert_eq!(buffer.drain_resident_bytes(), 0); + assert!(spill_files(tmp.path()).is_empty()); + } + #[tokio::test(flavor = "current_thread")] async fn abort_drops_xact_and_unlinks_spill() { let tmp = tempdir().unwrap(); @@ -5512,11 +5541,20 @@ mod tests { } let total = n * 512; let mut drain = b.drain_committed(4, 0, 0x9000, &[], false).await.unwrap(); + let read_buffers = drain + .merged + .as_ref() + .unwrap() + .sources + .iter() + .filter_map(|s| s.reader.as_ref()) + .map(|r| r.buffered_bytes() as u64) + .sum::(); let mut rows = 0usize; while let Some(batch) = drain.next_batch(4, usize::MAX, None).await.unwrap() { rows += batch.heaps.len(); assert!( - b.drain_resident_bytes() < total / 4, + b.drain_resident_bytes() < read_buffers + total / 4, "resident {} vs xact {total}", b.drain_resident_bytes(), ); @@ -5527,7 +5565,7 @@ mod tests { assert_eq!(rows as u64, n); assert!(b.drain_resident_peak() > 0, "gauge saw the merge heads"); assert!( - b.drain_resident_peak() < total / 4, + b.drain_resident_peak() < read_buffers + total / 4, "peak {} vs xact {total}", b.drain_resident_peak(), );