Skip to main content

mongreldb_core/
spill.rs

1//! Spill manager for query execution (spec section 10.5, S1E-004).
2//!
3//! Implemented in the Stage 1E wave: when the node
4//! [`MemoryGovernor`](crate::memory::MemoryGovernor) enters escalation step 3,
5//! spill-eligible operators move working memory to disk through this manager.
6//! Per the spec, spill files are:
7//!
8//! - **Query-ID namespaced** — every [`QueryId`] spills under its own
9//!   subdirectory `temp/spill/q-<hex>/` of the database root, so one query's
10//!   files never interleave with another's and cancellation removes exactly
11//!   one directory.
12//! - **Checksummed** — every frame carries a CRC32C (the same Castagnoli CRC
13//!   as the WAL) over its kind, sequence, and stored payload; the sealing
14//!   trailer carries a SHA-256 over all plaintext data payloads plus the
15//!   frame count, verified on read.
16//! - **Bounded** — a per-query cap (fed from
17//!   [`ResourceGroup::temporary_disk_bytes`]) and a node-global cap, enforced
18//!   with the memory governor's add-then-validate-rollback protocol; overflow
19//!   is the typed [`SpillError::BudgetExceeded`].
20//! - **Deleted on success/error/cancel/startup cleanup** — the [`SpillHandle`]
21//!   RAII guard deletes its file on drop, an unfinished [`SpillWriter`]
22//!   deletes its partial file on drop (the error/cancel path), a dropped
23//!   [`SpillSession`] removes the whole per-query directory, and
24//!   [`SpillManager::open`] sweeps every stale entry left by a prior process
25//!   run (spill files never outlive the process that created them).
26//! - **Encrypted when database encryption is enabled** — with a meta DEK
27//!   present ([`crate::encryption::meta_dek_for`]) every frame payload is
28//!   sealed through the page-cipher stack (`encrypt_blob`: AES-256-GCM with a
29//!   fresh random nonce per frame), mirroring the catalog/PITR seal idiom;
30//!   otherwise frames are stored plaintext.
31//!
32//! Layout of one spill file:
33//!
34//! ```text
35//! [magic: 8B "MDBSPILL"][version: u16][enc flag: u8][reserved: u8]
36//! frame*: [payload len: u32][crc32c: u32][kind: u8][seq: u64][payload]
37//! trailer frame (kind = 1): [data frames: u64][data bytes: u64][sha256: 32B]
38//! ```
39//!
40//! The on-disk tree lives under a descriptor-pinned
41//! [`DurableRoot`](crate::durable_file::DurableRoot); every file operation is
42//! descriptor-relative and `temp/spill` itself is created lazily on first
43//! use. Open one `SpillManager` per database. After a process restart, opening
44//! the manager sweeps files whose prior-process handles were closed by exit.
45
46use std::fmt;
47use std::io::{self, Read, Write};
48use std::path::Path;
49use std::sync::atomic::{AtomicU64, Ordering};
50use std::sync::Arc;
51
52use crc::{Crc, CRC_32_ISCSI};
53use mongreldb_types::ids::QueryId;
54use sha2::{Digest, Sha256};
55
56use crate::durable_file::DurableRoot;
57use crate::encryption::DEK_LEN;
58use crate::resource::ResourceGroup;
59use crate::MongrelError;
60
61/// Spill tree location relative to the database root.
62const SPILL_DIR_REL: &str = "temp/spill";
63/// File magic (8 bytes, the [`MongrelError::MagicMismatch`] idiom).
64const SPILL_MAGIC: &[u8; 8] = b"MDBSPILL";
65/// On-disk format version this build reads and writes.
66const FORMAT_VERSION: u16 = 1;
67/// Encryption flag stored in the header: frames are plaintext.
68const ENC_PLAINTEXT: u8 = 0;
69/// Encryption flag stored in the header: frames are AES-256-GCM sealed.
70const ENC_AES_GCM: u8 = 1;
71/// Header length: magic + version + enc flag + reserved.
72const HEADER_LEN: usize = 8 + 2 + 1 + 1;
73/// Frame head length: payload len + crc + kind + seq.
74const FRAME_HEAD_LEN: usize = 4 + 4 + 1 + 8;
75/// Kind byte of a data frame.
76const FRAME_DATA: u8 = 0;
77/// Kind byte of the sealing trailer frame.
78const FRAME_TRAILER: u8 = 1;
79/// Trailer payload length: frame count + byte count + SHA-256.
80const TRAILER_LEN: usize = 8 + 8 + 32;
81/// Maximum bytes of one frame's plaintext payload (bounds reader allocation
82/// on corrupt input, mirroring the PITR chunk cap).
83const MAX_FRAME_PAYLOAD: u64 = 64 * 1024 * 1024;
84
85const CRC32C: Crc<u32> = Crc::<u32>::new(&CRC_32_ISCSI);
86
87/// Errors of spill configuration, I/O, budget enforcement, and verification.
88#[derive(Debug, thiserror::Error)]
89pub enum SpillError {
90    /// Filesystem failure on the spill tree.
91    #[error("io error: {0}")]
92    Io(#[from] std::io::Error),
93    /// The [`SpillConfig`] failed validation.
94    #[error("invalid spill configuration: {0}")]
95    InvalidConfig(&'static str),
96    /// A frame append would exceed the per-query or the node-global spill
97    /// budget (S1E-004 "bounded").
98    #[error(
99        "spill budget exceeded for query {query_id}: requested {requested} bytes \
100         ({query_remaining} per-query, {global_remaining} global remaining)"
101    )]
102    BudgetExceeded {
103        /// The query that hit its bound.
104        query_id: QueryId,
105        /// Bytes the frame needed.
106        requested: u64,
107        /// Per-query budget remaining before the attempt.
108        query_remaining: u64,
109        /// Node-global budget remaining before the attempt.
110        global_remaining: u64,
111    },
112    /// One frame's plaintext payload exceeded [`MAX_FRAME_PAYLOAD`].
113    #[error("spill frame of {bytes} bytes exceeds the {limit}-byte frame limit")]
114    FrameTooLarge {
115        /// Attempted payload size.
116        bytes: u64,
117        /// Configured bound.
118        limit: u64,
119    },
120    /// A frame's stored CRC32C did not match its bytes.
121    #[error("checksum mismatch for {context}: expected {expected}, got {actual}")]
122    ChecksumMismatch {
123        /// Which frame/file failed verification.
124        context: String,
125        /// Stored checksum.
126        expected: u32,
127        /// Computed checksum.
128        actual: u32,
129    },
130    /// Structural corruption: bad magic, version, framing, sequence gap, or
131    /// a trailer that does not match the streamed frames. Always fail-closed.
132    #[error("corrupt spill file: {0}")]
133    Corrupt(String),
134    /// An encrypted spill file was opened by a manager without a meta DEK.
135    #[error("encrypted spill file requires the database encryption key")]
136    EncryptionRequired,
137    /// A meta DEK was configured (encrypted path).
138    #[error("spill encryption requires the `encryption` feature")]
139    EncryptionDisabled,
140    /// Sealing a frame failed.
141    #[error("spill encryption error: {0}")]
142    Encryption(String),
143    /// Opening a sealed frame failed (wrong key or tampering).
144    #[error("spill decryption error: {0}")]
145    Decryption(String),
146}
147
148impl From<SpillError> for MongrelError {
149    fn from(error: SpillError) -> Self {
150        match error {
151            SpillError::Io(error) => MongrelError::Io(error),
152            SpillError::InvalidConfig(message) => MongrelError::InvalidArgument(message.into()),
153            SpillError::BudgetExceeded {
154                requested,
155                query_remaining,
156                global_remaining,
157                ..
158            } => MongrelError::ResourceLimitExceeded {
159                resource: "spill temporary disk",
160                requested: usize::try_from(requested).unwrap_or(usize::MAX),
161                limit: usize::try_from(
162                    requested.saturating_add(query_remaining.min(global_remaining)),
163                )
164                .unwrap_or(usize::MAX),
165            },
166            SpillError::FrameTooLarge { bytes, limit } => MongrelError::ResourceLimitExceeded {
167                resource: "spill frame",
168                requested: usize::try_from(bytes).unwrap_or(usize::MAX),
169                limit: usize::try_from(limit).unwrap_or(usize::MAX),
170            },
171            SpillError::ChecksumMismatch {
172                context,
173                expected,
174                actual,
175            } => MongrelError::ChecksumMismatch {
176                expected: u64::from(expected),
177                actual: u64::from(actual),
178                context,
179            },
180            SpillError::Corrupt(message) => {
181                MongrelError::Other(format!("corrupt spill file: {message}"))
182            }
183            SpillError::EncryptionRequired | SpillError::EncryptionDisabled => {
184                MongrelError::EncryptionDisabled
185            }
186            SpillError::Encryption(message) => MongrelError::Encryption(message),
187            SpillError::Decryption(message) => MongrelError::Decryption(message),
188        }
189    }
190}
191
192/// Node-global spill configuration (S1E-004). The per-query bound is supplied
193/// per session from [`ResourceGroup::temporary_disk_bytes`].
194#[derive(Debug, Clone, Copy, PartialEq, Eq)]
195pub struct SpillConfig {
196    /// Total bytes of live spill files across every query on this node.
197    pub global_bytes: u64,
198}
199
200impl SpillConfig {
201    /// A config with the given node-global cap.
202    pub fn new(global_bytes: u64) -> Self {
203        Self { global_bytes }
204    }
205
206    fn validate(&self) -> Result<(), SpillError> {
207        if self.global_bytes == 0 {
208            return Err(SpillError::InvalidConfig(
209                "global spill budget must be nonzero",
210            ));
211        }
212        Ok(())
213    }
214}
215
216/// A point-in-time snapshot of spill-manager state (telemetry and tests).
217#[derive(Debug, Clone, PartialEq, Eq)]
218pub struct SpillStats {
219    /// Cumulative stored bytes written (headers, frames, trailers).
220    pub bytes_written: u64,
221    /// Cumulative stored bytes read back.
222    pub bytes_read: u64,
223    /// Finished spill files currently live (held by a [`SpillHandle`]).
224    pub files_live: u64,
225    /// Live bytes currently charged against the global budget.
226    pub global_used: u64,
227    /// Configured node-global budget.
228    pub global_budget_bytes: u64,
229    /// Global budget remaining (`global_budget_bytes - global_used`).
230    pub budget_remaining: u64,
231}
232
233struct ManagerInner {
234    db_root: DurableRoot,
235    /// Lazily created pinned `temp/spill` root (`None` until the first spill
236    /// file is created; the startup sweep never creates it).
237    spill_root: parking_lot::Mutex<Option<DurableRoot>>,
238    config: SpillConfig,
239    meta_dek: Option<[u8; DEK_LEN]>,
240    global_used: AtomicU64,
241    bytes_written: AtomicU64,
242    bytes_read: AtomicU64,
243    files_live: AtomicU64,
244}
245
246impl ManagerInner {
247    /// The pinned `temp/spill` root, creating it (and `temp/`) on first use.
248    fn spill_root(&self) -> Result<DurableRoot, SpillError> {
249        let mut guard = self.spill_root.lock();
250        if let Some(root) = guard.as_ref() {
251            return Ok(root.try_clone()?);
252        }
253        let root = self.db_root.create_directory_all_pinned(SPILL_DIR_REL)?;
254        *guard = Some(root.try_clone()?);
255        Ok(root)
256    }
257
258    /// The pinned `temp/spill` root only if it already exists (never
259    /// creates). Used by cleanup paths, which must not resurrect the tree.
260    fn spill_root_if_created(&self) -> Option<DurableRoot> {
261        self.spill_root
262            .lock()
263            .as_ref()
264            .and_then(|root| root.try_clone().ok())
265    }
266
267    /// Removes every stale entry a prior process run left in `temp/spill`.
268    /// Spill files are process-local by construction (live files are held by
269    /// handles), so everything present at open is garbage.
270    fn sweep_stale(&self) -> Result<(), SpillError> {
271        match self.db_root.entry_exists(SPILL_DIR_REL) {
272            Ok(true) => {}
273            // No `temp/spill` (or no `temp` yet): nothing to sweep. The tree
274            // is created lazily on the first spill.
275            Ok(false) => return Ok(()),
276            Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()),
277            Err(error) => return Err(SpillError::Io(error)),
278        }
279        let root = self.db_root.open_directory(SPILL_DIR_REL)?;
280        // `io_path` is the descriptor-pinned operational path of `root`, so
281        // the enumeration itself cannot wander off the pinned directory.
282        for entry in std::fs::read_dir(root.io_path()?)? {
283            let entry = entry?;
284            let name = entry.file_name();
285            if entry.file_type()?.is_dir() {
286                root.remove_directory_all(Path::new(&name))?;
287            } else {
288                root.remove_file(Path::new(&name))?;
289            }
290        }
291        Ok(())
292    }
293
294    fn header(&self) -> [u8; HEADER_LEN] {
295        let mut header = [0u8; HEADER_LEN];
296        header[..8].copy_from_slice(SPILL_MAGIC);
297        header[8..10].copy_from_slice(&FORMAT_VERSION.to_le_bytes());
298        header[10] = if self.meta_dek.is_some() {
299            ENC_AES_GCM
300        } else {
301            ENC_PLAINTEXT
302        };
303        header
304    }
305
306    /// Charges `bytes` against the per-query and global budgets with the
307    /// governor's add-then-validate-rollback protocol, so a granted set never
308    /// exceeds either bound.
309    fn try_charge(&self, session: &SessionInner, bytes: u64) -> Result<(), SpillError> {
310        let new_query = session.used.fetch_add(bytes, Ordering::Relaxed) + bytes;
311        let new_global = self.global_used.fetch_add(bytes, Ordering::Relaxed) + bytes;
312        if new_query <= session.cap && new_global <= self.config.global_bytes {
313            Ok(())
314        } else {
315            session.used.fetch_sub(bytes, Ordering::Relaxed);
316            self.global_used.fetch_sub(bytes, Ordering::Relaxed);
317            Err(SpillError::BudgetExceeded {
318                query_id: session.query_id,
319                requested: bytes,
320                query_remaining: session.cap.saturating_sub(new_query - bytes),
321                global_remaining: self.config.global_bytes.saturating_sub(new_global - bytes),
322            })
323        }
324    }
325
326    /// Releases a previous charge (exact inverse of [`try_charge`](Self::try_charge)).
327    fn release(&self, session: &SessionInner, bytes: u64) {
328        session.used.fetch_sub(bytes, Ordering::Relaxed);
329        self.global_used.fetch_sub(bytes, Ordering::Relaxed);
330    }
331}
332
333/// The node-level spill manager (S1E-004). Cheap to clone (one `Arc`);
334/// thread-safe; budget accounting is lock-free.
335pub struct SpillManager {
336    inner: Arc<ManagerInner>,
337}
338
339impl Clone for SpillManager {
340    fn clone(&self) -> Self {
341        Self {
342            inner: Arc::clone(&self.inner),
343        }
344    }
345}
346
347impl fmt::Debug for SpillManager {
348    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
349        f.debug_struct("SpillManager")
350            .field("global_budget_bytes", &self.inner.config.global_bytes)
351            .field(
352                "global_used",
353                &self.inner.global_used.load(Ordering::Relaxed),
354            )
355            .field("files_live", &self.inner.files_live.load(Ordering::Relaxed))
356            .field("encrypted", &self.inner.meta_dek.is_some())
357            .finish()
358    }
359}
360
361impl SpillManager {
362    /// Opens the spill manager on a database root, sweeping every stale entry
363    /// a prior process run left in `temp/spill` (S1E-004 startup cleanup).
364    /// The `temp/spill` tree itself is created lazily on the first spill.
365    ///
366    /// `meta_dek` is the database meta DEK
367    /// ([`crate::encryption::meta_dek_for`]): `Some` seals every frame with
368    /// AES-256-GCM, `None` stores frames plaintext. Without the `encryption`
369    /// feature a `Some` DEK is rejected (fail closed).
370    pub fn open(
371        db_root: &DurableRoot,
372        config: SpillConfig,
373        meta_dek: Option<[u8; DEK_LEN]>,
374    ) -> Result<Self, SpillError> {
375        config.validate()?;
376        let manager = Self {
377            inner: Arc::new(ManagerInner {
378                db_root: db_root.try_clone()?,
379                spill_root: parking_lot::Mutex::new(None),
380                config,
381                meta_dek,
382                global_used: AtomicU64::new(0),
383                bytes_written: AtomicU64::new(0),
384                bytes_read: AtomicU64::new(0),
385                files_live: AtomicU64::new(0),
386            }),
387        };
388        manager.inner.sweep_stale()?;
389        Ok(manager)
390    }
391
392    /// The manager's configuration.
393    pub fn config(&self) -> &SpillConfig {
394        &self.inner.config
395    }
396
397    /// Starts a spill session for one query with an explicit per-query cap.
398    /// One session per `query_id` at a time: sessions share the per-query
399    /// directory name, and chunk creation is exclusive.
400    pub fn begin_query(
401        &self,
402        query_id: QueryId,
403        per_query_bytes: u64,
404    ) -> Result<SpillSession, SpillError> {
405        Ok(SpillSession {
406            inner: Arc::new(SessionInner {
407                manager: self.clone(),
408                query_id,
409                dir_name: format!("q-{}", query_id.to_hex()),
410                cap: per_query_bytes,
411                used: AtomicU64::new(0),
412                next_chunk: AtomicU64::new(0),
413            }),
414        })
415    }
416
417    /// Starts a spill session for one query admitted into `group`, with the
418    /// per-query cap fed from [`ResourceGroup::temporary_disk_bytes`]
419    /// (S1E-002/S1E-004).
420    pub fn begin_query_in_group(
421        &self,
422        query_id: QueryId,
423        group: &ResourceGroup,
424    ) -> Result<SpillSession, SpillError> {
425        self.begin_query(query_id, group.temporary_disk_bytes)
426    }
427
428    /// A point-in-time snapshot of spill state.
429    pub fn stats(&self) -> SpillStats {
430        let global_used = self.inner.global_used.load(Ordering::Relaxed);
431        SpillStats {
432            bytes_written: self.inner.bytes_written.load(Ordering::Relaxed),
433            bytes_read: self.inner.bytes_read.load(Ordering::Relaxed),
434            files_live: self.inner.files_live.load(Ordering::Relaxed),
435            global_used,
436            global_budget_bytes: self.inner.config.global_bytes,
437            budget_remaining: self.inner.config.global_bytes.saturating_sub(global_used),
438        }
439    }
440}
441
442struct SessionInner {
443    manager: SpillManager,
444    query_id: QueryId,
445    /// Per-query directory name under `temp/spill` (`q-<hex>`).
446    dir_name: String,
447    cap: u64,
448    used: AtomicU64,
449    next_chunk: AtomicU64,
450}
451
452impl SessionInner {
453    /// The pinned per-query directory, creating it on first use.
454    fn query_dir(&self) -> Result<DurableRoot, SpillError> {
455        Ok(self
456            .manager
457            .inner
458            .spill_root()?
459            .create_directory_all_pinned(&self.dir_name)?)
460    }
461}
462
463impl Drop for SessionInner {
464    /// Removes the whole per-query directory (S1E-004 cleanup on
465    /// success/error/cancel). Never creates the tree to delete it.
466    fn drop(&mut self) {
467        if let Some(root) = self.manager.inner.spill_root_if_created() {
468            let _ = root.remove_directory_all(Path::new(&self.dir_name));
469        }
470    }
471}
472
473/// A query's spill session: namespaced, budgeted factory of spill files.
474/// Dropping the session removes the whole per-query directory.
475pub struct SpillSession {
476    inner: Arc<SessionInner>,
477}
478
479impl fmt::Debug for SpillSession {
480    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
481        f.debug_struct("SpillSession")
482            .field("query_id", &self.inner.query_id)
483            .field("cap", &self.inner.cap)
484            .field("used", &self.inner.used.load(Ordering::Relaxed))
485            .finish()
486    }
487}
488
489impl SpillSession {
490    /// The query this session spills for.
491    pub fn query_id(&self) -> QueryId {
492        self.inner.query_id
493    }
494
495    /// The per-query spill cap in bytes.
496    pub fn cap(&self) -> u64 {
497        self.inner.cap
498    }
499
500    /// Live bytes currently charged to this query.
501    pub fn used(&self) -> u64 {
502        self.inner.used.load(Ordering::Relaxed)
503    }
504
505    /// Per-query budget remaining.
506    pub fn budget_remaining(&self) -> u64 {
507        self.inner.cap.saturating_sub(self.used())
508    }
509
510    /// Creates the next spill file of this query and returns its writer.
511    /// Dropping the writer before [`SpillWriter::finish`] deletes the partial
512    /// file (the error/cancel path).
513    pub fn new_writer(&self) -> Result<SpillWriter, SpillError> {
514        let seq = self.inner.next_chunk.fetch_add(1, Ordering::Relaxed);
515        let name = format!("chunk-{seq:06}.spill");
516        let dir = self.inner.query_dir()?;
517        let manager = &self.inner.manager;
518        manager.inner.try_charge(&self.inner, HEADER_LEN as u64)?;
519        let result = (|| {
520            let mut file = dir.create_regular_new(&name)?;
521            file.write_all(&manager.inner.header())?;
522            io::Result::Ok(file)
523        })();
524        match result {
525            Ok(file) => {
526                manager
527                    .inner
528                    .bytes_written
529                    .fetch_add(HEADER_LEN as u64, Ordering::Relaxed);
530                Ok(SpillWriter {
531                    inner: Some(WriterInner {
532                        session: Arc::clone(&self.inner),
533                        dir,
534                        name,
535                        file,
536                        bytes_on_disk: HEADER_LEN as u64,
537                        data_frames: 0,
538                        data_bytes: 0,
539                        digest: Sha256::new(),
540                        next_seq: 0,
541                    }),
542                })
543            }
544            Err(error) => {
545                manager.inner.release(&self.inner, HEADER_LEN as u64);
546                Err(SpillError::Io(error))
547            }
548        }
549    }
550}
551
552struct WriterInner {
553    session: Arc<SessionInner>,
554    /// Pinned per-query directory (descriptor-relative delete on abort).
555    dir: DurableRoot,
556    name: String,
557    file: std::fs::File,
558    bytes_on_disk: u64,
559    data_frames: u64,
560    data_bytes: u64,
561    digest: Sha256,
562    next_seq: u64,
563}
564
565impl WriterInner {
566    /// Seals (when encrypted), frames, charges, and appends one frame.
567    fn write_frame(&mut self, kind: u8, payload: &[u8]) -> Result<(), SpillError> {
568        if payload.len() as u64 > MAX_FRAME_PAYLOAD {
569            return Err(SpillError::FrameTooLarge {
570                bytes: payload.len() as u64,
571                limit: MAX_FRAME_PAYLOAD,
572            });
573        }
574        let stored = seal_payload(self.session.manager.inner.meta_dek.as_ref(), payload)?;
575        let seq = self.next_seq;
576        let mut digest = CRC32C.digest();
577        digest.update(&[kind]);
578        digest.update(&seq.to_le_bytes());
579        digest.update(&stored);
580        let crc = digest.finalize();
581        let frame_bytes = FRAME_HEAD_LEN as u64 + stored.len() as u64;
582        let manager = &self.session.manager;
583        manager.inner.try_charge(&self.session, frame_bytes)?;
584        let result = (|| {
585            self.file.write_all(&(stored.len() as u32).to_le_bytes())?;
586            self.file.write_all(&crc.to_le_bytes())?;
587            self.file.write_all(&[kind])?;
588            self.file.write_all(&seq.to_le_bytes())?;
589            self.file.write_all(&stored)
590        })();
591        match result {
592            Ok(()) => {
593                self.next_seq += 1;
594                self.bytes_on_disk += frame_bytes;
595                manager
596                    .inner
597                    .bytes_written
598                    .fetch_add(frame_bytes, Ordering::Relaxed);
599                Ok(())
600            }
601            Err(error) => {
602                manager.inner.release(&self.session, frame_bytes);
603                Err(SpillError::Io(error))
604            }
605        }
606    }
607
608    /// Writes the sealing trailer (frame count, byte count, SHA-256 over all
609    /// plaintext data payloads), then fsyncs file and directory.
610    fn write_trailer_and_sync(&mut self) -> Result<(), SpillError> {
611        let hash: [u8; 32] = self.digest.clone().finalize().into();
612        let mut trailer = Vec::with_capacity(TRAILER_LEN);
613        trailer.extend_from_slice(&self.data_frames.to_le_bytes());
614        trailer.extend_from_slice(&self.data_bytes.to_le_bytes());
615        trailer.extend_from_slice(&hash);
616        self.write_frame(FRAME_TRAILER, &trailer)?;
617        self.file.sync_all()?;
618        self.dir.sync_entry_parent(Path::new(&self.name))?;
619        Ok(())
620    }
621}
622
623/// Deletes the partial file of an unfinished writer and releases its budget
624/// (the error/cancel path; S1E-004 "deleted on error/cancel").
625fn abort_writer(inner: WriterInner) {
626    let WriterInner {
627        session,
628        dir,
629        name,
630        file,
631        bytes_on_disk,
632        ..
633    } = inner;
634    drop(file);
635    let _ = dir.remove_file(Path::new(&name));
636    session.manager.inner.release(&session, bytes_on_disk);
637}
638
639/// Streaming writer of one spill file. Append data frames with
640/// [`append`](Self::append); seal and fsync with [`finish`](Self::finish).
641/// Dropping an unfinished writer deletes the partial file.
642pub struct SpillWriter {
643    inner: Option<WriterInner>,
644}
645
646impl fmt::Debug for SpillWriter {
647    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
648        f.debug_struct("SpillWriter")
649            .field("name", &self.inner.as_ref().map(|inner| &inner.name))
650            .field(
651                "bytes_on_disk",
652                &self.inner.as_ref().map(|inner| inner.bytes_on_disk),
653            )
654            .finish()
655    }
656}
657
658impl SpillWriter {
659    /// Appends one data frame. `payload` is opaque to the manager (an
660    /// operator's serialized run); it is checksummed and, when the database
661    /// is encrypted, sealed before it touches disk.
662    pub fn append(&mut self, payload: &[u8]) -> Result<(), SpillError> {
663        let inner = self
664            .inner
665            .as_mut()
666            .expect("spill writer is live until finish");
667        inner.write_frame(FRAME_DATA, payload)?;
668        inner.digest.update(payload);
669        inner.data_frames += 1;
670        inner.data_bytes += payload.len() as u64;
671        Ok(())
672    }
673
674    /// Stored bytes written so far (header and frames, post-seal).
675    pub fn bytes_on_disk(&self) -> u64 {
676        self.inner.as_ref().map_or(0, |inner| inner.bytes_on_disk)
677    }
678
679    /// Seals the file with its checksum trailer, fsyncs it, and returns the
680    /// RAII handle. On error the partial file is deleted, exactly as if the
681    /// writer had been dropped.
682    pub fn finish(mut self) -> Result<SpillHandle, SpillError> {
683        let mut inner = self
684            .inner
685            .take()
686            .expect("spill writer is live until finish");
687        if let Err(error) = inner.write_trailer_and_sync() {
688            abort_writer(inner);
689            return Err(error);
690        }
691        inner
692            .session
693            .manager
694            .inner
695            .files_live
696            .fetch_add(1, Ordering::Relaxed);
697        let WriterInner {
698            session,
699            dir,
700            name,
701            bytes_on_disk,
702            data_frames,
703            ..
704        } = inner;
705        Ok(SpillHandle {
706            inner: Some(HandleInner {
707                session,
708                dir,
709                name,
710                bytes_on_disk,
711                data_frames,
712            }),
713        })
714    }
715
716    /// Abandons the file: deletes the partial write and releases its budget.
717    /// Equivalent to dropping the writer, but reports I/O failures.
718    pub fn abort(mut self) -> Result<(), SpillError> {
719        let inner = self
720            .inner
721            .take()
722            .expect("spill writer is live until finish");
723        let WriterInner {
724            session,
725            dir,
726            name,
727            file,
728            bytes_on_disk,
729            ..
730        } = inner;
731        drop(file);
732        dir.remove_file(Path::new(&name))?;
733        session.manager.inner.release(&session, bytes_on_disk);
734        Ok(())
735    }
736}
737
738impl Drop for SpillWriter {
739    fn drop(&mut self) {
740        if let Some(inner) = self.inner.take() {
741            abort_writer(inner);
742        }
743    }
744}
745
746struct HandleInner {
747    session: Arc<SessionInner>,
748    dir: DurableRoot,
749    name: String,
750    bytes_on_disk: u64,
751    data_frames: u64,
752}
753
754/// Deletes the finished file and releases its budget and live-file count.
755fn delete_file(inner: HandleInner) -> Result<(), SpillError> {
756    inner.dir.remove_file(Path::new(&inner.name))?;
757    inner
758        .session
759        .manager
760        .inner
761        .release(&inner.session, inner.bytes_on_disk);
762    inner
763        .session
764        .manager
765        .inner
766        .files_live
767        .fetch_sub(1, Ordering::Relaxed);
768    Ok(())
769}
770
771/// RAII guard for one sealed spill file: dropping it deletes the file and
772/// releases its budget (S1E-004 "deleted on success"). Open a verify-on-read
773/// [`SpillReader`] with [`reader`](Self::reader).
774#[must_use = "a spill handle deletes its file on drop"]
775pub struct SpillHandle {
776    inner: Option<HandleInner>,
777}
778
779impl fmt::Debug for SpillHandle {
780    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
781        f.debug_struct("SpillHandle")
782            .field("name", &self.inner.as_ref().map(|inner| &inner.name))
783            .field(
784                "bytes_on_disk",
785                &self.inner.as_ref().map(|inner| inner.bytes_on_disk),
786            )
787            .finish()
788    }
789}
790
791impl SpillHandle {
792    /// The query whose session created this file.
793    pub fn query_id(&self) -> QueryId {
794        self.inner
795            .as_ref()
796            .expect("spill handle is live")
797            .session
798            .query_id
799    }
800
801    /// Stored file size in bytes (header, frames, trailer).
802    pub fn bytes_on_disk(&self) -> u64 {
803        self.inner.as_ref().map_or(0, |inner| inner.bytes_on_disk)
804    }
805
806    /// Number of data frames sealed into the file.
807    pub fn frames(&self) -> u64 {
808        self.inner.as_ref().map_or(0, |inner| inner.data_frames)
809    }
810
811    /// Opens a streaming, verify-on-read reader over the file.
812    pub fn reader(&self) -> Result<SpillReader, SpillError> {
813        let inner = self.inner.as_ref().expect("spill handle is live");
814        let file = inner.dir.open_regular(Path::new(&inner.name))?;
815        reader_from(file, &inner.session.manager)
816    }
817
818    /// Deletes the file now, releasing its budget. Equivalent to dropping the
819    /// handle, but reports I/O failures.
820    pub fn delete(mut self) -> Result<(), SpillError> {
821        let inner = self.inner.take().expect("spill handle is live");
822        delete_file(inner)
823    }
824}
825
826impl Drop for SpillHandle {
827    fn drop(&mut self) {
828        if let Some(inner) = self.inner.take() {
829            let _ = delete_file(inner);
830        }
831    }
832}
833
834/// Parses and validates a spill file header and builds the reader (shared by
835/// [`SpillHandle::reader`] and tests). Fails closed on bad magic, an
836/// unsupported version, an unknown encryption flag, or a DEK mismatch with
837/// the manager (mirroring the PITR chunk envelope rules).
838fn reader_from(mut file: std::fs::File, manager: &SpillManager) -> Result<SpillReader, SpillError> {
839    let mut header = [0u8; HEADER_LEN];
840    file.read_exact(&mut header).map_err(|error| {
841        if error.kind() == io::ErrorKind::UnexpectedEof {
842            SpillError::Corrupt("spill file is shorter than its header".into())
843        } else {
844            SpillError::Io(error)
845        }
846    })?;
847    manager
848        .inner
849        .bytes_read
850        .fetch_add(HEADER_LEN as u64, Ordering::Relaxed);
851    if header[..8] != SPILL_MAGIC[..] {
852        return Err(SpillError::Corrupt("bad spill magic".into()));
853    }
854    let version = u16::from_le_bytes(header[8..10].try_into().expect("slice length"));
855    if version != FORMAT_VERSION {
856        return Err(SpillError::Corrupt(format!(
857            "unsupported spill format version {version}"
858        )));
859    }
860    let dek = match (header[10], manager.inner.meta_dek.as_ref()) {
861        (ENC_PLAINTEXT, None) => None,
862        (ENC_AES_GCM, Some(dek)) => Some(*dek),
863        (ENC_AES_GCM, None) => return Err(SpillError::EncryptionRequired),
864        (ENC_PLAINTEXT, Some(_)) => {
865            return Err(SpillError::Corrupt(
866                "plaintext spill file opened with an encryption key".into(),
867            ))
868        }
869        (other, _) => {
870            return Err(SpillError::Corrupt(format!(
871                "unknown spill encryption flag {other}"
872            )))
873        }
874    };
875    Ok(SpillReader {
876        file,
877        manager: manager.clone(),
878        dek,
879        next_seq: 0,
880        data_frames: 0,
881        data_bytes: 0,
882        digest: Sha256::new(),
883        done: false,
884    })
885}
886
887/// Streaming, verify-on-read reader of one sealed spill file. Every frame's
888/// CRC32C is verified before its payload is returned (and decrypted, when
889/// the file is sealed); the closing trailer re-checks the frame count, byte
890/// count, and the SHA-256 over all data payloads. Any failure is terminal.
891pub struct SpillReader {
892    file: std::fs::File,
893    manager: SpillManager,
894    dek: Option<[u8; DEK_LEN]>,
895    next_seq: u64,
896    data_frames: u64,
897    data_bytes: u64,
898    digest: Sha256,
899    done: bool,
900}
901
902impl fmt::Debug for SpillReader {
903    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
904        f.debug_struct("SpillReader")
905            .field("next_seq", &self.next_seq)
906            .field("data_frames", &self.data_frames)
907            .field("done", &self.done)
908            .finish()
909    }
910}
911
912impl SpillReader {
913    /// Reads the next data frame, or `None` once the sealing trailer has been
914    /// verified. Errors are terminal: after one, the reader yields no more
915    /// frames.
916    pub fn next_frame(&mut self) -> Result<Option<Vec<u8>>, SpillError> {
917        if self.done {
918            return Ok(None);
919        }
920        match self.next_frame_inner() {
921            Ok(frame) => Ok(frame),
922            Err(error) => {
923                self.done = true;
924                Err(error)
925            }
926        }
927    }
928
929    fn next_frame_inner(&mut self) -> Result<Option<Vec<u8>>, SpillError> {
930        let mut head = [0u8; FRAME_HEAD_LEN];
931        self.file.read_exact(&mut head).map_err(|error| {
932            if error.kind() == io::ErrorKind::UnexpectedEof {
933                SpillError::Corrupt("truncated spill file: missing trailer".into())
934            } else {
935                SpillError::Io(error)
936            }
937        })?;
938        let len = u64::from(u32::from_le_bytes(
939            head[0..4].try_into().expect("slice length"),
940        ));
941        let expected_crc = u32::from_le_bytes(head[4..8].try_into().expect("slice length"));
942        let kind = head[8];
943        let seq = u64::from_le_bytes(head[9..17].try_into().expect("slice length"));
944        if len > MAX_FRAME_PAYLOAD {
945            return Err(SpillError::Corrupt(format!(
946                "spill frame of {len} bytes exceeds the {MAX_FRAME_PAYLOAD}-byte limit"
947            )));
948        }
949        if seq != self.next_seq {
950            return Err(SpillError::Corrupt(format!(
951                "spill frame sequence gap: expected {}, found {seq}",
952                self.next_seq
953            )));
954        }
955        let mut stored = vec![0u8; len as usize];
956        self.file.read_exact(&mut stored).map_err(|error| {
957            if error.kind() == io::ErrorKind::UnexpectedEof {
958                SpillError::Corrupt("truncated spill frame payload".into())
959            } else {
960                SpillError::Io(error)
961            }
962        })?;
963        let mut digest = CRC32C.digest();
964        digest.update(&[kind]);
965        digest.update(&seq.to_le_bytes());
966        digest.update(&stored);
967        let actual_crc = digest.finalize();
968        if actual_crc != expected_crc {
969            return Err(SpillError::ChecksumMismatch {
970                context: format!("spill frame {seq}"),
971                expected: expected_crc,
972                actual: actual_crc,
973            });
974        }
975        self.manager
976            .inner
977            .bytes_read
978            .fetch_add(FRAME_HEAD_LEN as u64 + len, Ordering::Relaxed);
979        let plaintext = open_payload(self.dek.as_ref(), &stored)?;
980        self.next_seq += 1;
981        match kind {
982            FRAME_DATA => {
983                self.digest.update(&plaintext);
984                self.data_frames += 1;
985                self.data_bytes += plaintext.len() as u64;
986                Ok(Some(plaintext))
987            }
988            FRAME_TRAILER => {
989                if plaintext.len() != TRAILER_LEN {
990                    return Err(SpillError::Corrupt(
991                        "spill trailer has the wrong length".into(),
992                    ));
993                }
994                let frames = u64::from_le_bytes(plaintext[0..8].try_into().expect("slice length"));
995                let bytes = u64::from_le_bytes(plaintext[8..16].try_into().expect("slice length"));
996                let hash: [u8; 32] = plaintext[16..48].try_into().expect("slice length");
997                let actual_hash: [u8; 32] = self.digest.clone().finalize().into();
998                if frames != self.data_frames || bytes != self.data_bytes || hash != actual_hash {
999                    return Err(SpillError::Corrupt(
1000                        "spill trailer does not match the streamed frames".into(),
1001                    ));
1002                }
1003                self.done = true;
1004                Ok(None)
1005            }
1006            other => Err(SpillError::Corrupt(format!(
1007                "unknown spill frame kind {other}"
1008            ))),
1009        }
1010    }
1011}
1012
1013impl Iterator for SpillReader {
1014    type Item = Result<Vec<u8>, SpillError>;
1015
1016    fn next(&mut self) -> Option<Self::Item> {
1017        match self.next_frame() {
1018            Ok(Some(frame)) => Some(Ok(frame)),
1019            Ok(None) => None,
1020            Err(error) => Some(Err(error)),
1021        }
1022    }
1023}
1024
1025/// Maps the crypto stack's `MongrelError` onto the typed spill errors.
1026fn map_crypto(error: MongrelError) -> SpillError {
1027    match error {
1028        MongrelError::Encryption(message) => SpillError::Encryption(message),
1029        MongrelError::Decryption(message) => SpillError::Decryption(message),
1030        other => SpillError::Corrupt(other.to_string()),
1031    }
1032}
1033
1034/// Seals one frame's plaintext payload with the page-cipher stack when a meta
1035/// DEK is present (fresh random nonce per frame, the `encrypt_blob` idiom);
1036/// passes plaintext through otherwise.
1037fn seal_payload(dek: Option<&[u8; DEK_LEN]>, plaintext: &[u8]) -> Result<Vec<u8>, SpillError> {
1038    match dek {
1039        Some(dek) => crate::encryption::encrypt_blob(dek, plaintext).map_err(map_crypto),
1040        None => Ok(plaintext.to_vec()),
1041    }
1042}
1043
1044/// Inverse of [`seal_payload`]: authenticates and opens a sealed frame.
1045fn open_payload(dek: Option<&[u8; DEK_LEN]>, stored: &[u8]) -> Result<Vec<u8>, SpillError> {
1046    match dek {
1047        Some(dek) => crate::encryption::decrypt_blob(dek, stored).map_err(map_crypto),
1048        None => Ok(stored.to_vec()),
1049    }
1050}
1051
1052#[cfg(test)]
1053mod tests {
1054    use super::*;
1055    use std::io::{Seek, SeekFrom};
1056    use std::path::PathBuf;
1057
1058    fn manager(dir: &tempfile::TempDir, global_bytes: u64) -> SpillManager {
1059        let root = DurableRoot::open(dir.path()).unwrap();
1060        SpillManager::open(&root, SpillConfig::new(global_bytes), None).unwrap()
1061    }
1062
1063    fn query_dir(dir: &tempfile::TempDir, query_id: QueryId) -> PathBuf {
1064        dir.path()
1065            .join("temp")
1066            .join("spill")
1067            .join(format!("q-{}", query_id.to_hex()))
1068    }
1069
1070    fn only_file(dir: &tempfile::TempDir, query_id: QueryId) -> PathBuf {
1071        let entries: Vec<_> = std::fs::read_dir(query_dir(dir, query_id))
1072            .unwrap()
1073            .map(|entry| entry.unwrap().path())
1074            .collect();
1075        assert_eq!(entries.len(), 1, "expected exactly one spill file");
1076        entries[0].clone()
1077    }
1078
1079    /// Serializes one frame exactly as `WriterInner::write_frame` does
1080    /// (plaintext mode), for crafting corrupt files.
1081    fn crafted_frame(kind: u8, seq: u64, payload: &[u8]) -> Vec<u8> {
1082        let mut digest = CRC32C.digest();
1083        digest.update(&[kind]);
1084        digest.update(&seq.to_le_bytes());
1085        digest.update(payload);
1086        let crc = digest.finalize();
1087        let mut out = Vec::with_capacity(FRAME_HEAD_LEN + payload.len());
1088        out.extend_from_slice(&(payload.len() as u32).to_le_bytes());
1089        out.extend_from_slice(&crc.to_le_bytes());
1090        out.push(kind);
1091        out.extend_from_slice(&seq.to_le_bytes());
1092        out.extend_from_slice(payload);
1093        out
1094    }
1095
1096    fn crafted_header(enc: u8) -> Vec<u8> {
1097        let mut header = vec![0u8; HEADER_LEN];
1098        header[..8].copy_from_slice(SPILL_MAGIC);
1099        header[8..10].copy_from_slice(&FORMAT_VERSION.to_le_bytes());
1100        header[10] = enc;
1101        header
1102    }
1103
1104    #[test]
1105    fn frame_round_trip_streams_in_order() {
1106        let dir = tempfile::tempdir().unwrap();
1107        let manager = manager(&dir, 1 << 20);
1108        let session = manager.begin_query(QueryId::new_random(), 1 << 20).unwrap();
1109        let payloads: Vec<Vec<u8>> = vec![
1110            b"first".to_vec(),
1111            Vec::new(), // empty frames are legal
1112            vec![0xAB; 1000],
1113            b"last".to_vec(),
1114        ];
1115        let mut writer = session.new_writer().unwrap();
1116        for payload in &payloads {
1117            writer.append(payload).unwrap();
1118        }
1119        let handle = writer.finish().unwrap();
1120        assert_eq!(handle.frames(), 4);
1121        assert_eq!(handle.query_id(), session.query_id());
1122        assert_eq!(handle.bytes_on_disk(), session.used());
1123
1124        let frames: Vec<Vec<u8>> = handle.reader().unwrap().collect::<Result<_, _>>().unwrap();
1125        assert_eq!(frames, payloads);
1126
1127        let stats = manager.stats();
1128        assert_eq!(stats.files_live, 1);
1129        assert_eq!(stats.bytes_written, handle.bytes_on_disk());
1130        assert_eq!(stats.bytes_read, handle.bytes_on_disk());
1131        assert_eq!(stats.global_used, handle.bytes_on_disk());
1132        assert_eq!(stats.budget_remaining, (1 << 20) - handle.bytes_on_disk());
1133    }
1134
1135    #[test]
1136    fn large_multi_chunk_round_trip() {
1137        let dir = tempfile::tempdir().unwrap();
1138        let manager = manager(&dir, 1 << 24);
1139        let session = manager.begin_query(QueryId::new_random(), 1 << 24).unwrap();
1140        // Two chunk files, three frames of mixed large sizes each.
1141        let mut handles = Vec::new();
1142        let mut expected = Vec::new();
1143        for chunk in 0..2u8 {
1144            let mut writer = session.new_writer().unwrap();
1145            let mut payloads = Vec::new();
1146            for (index, len) in [(1usize << 20), (1 << 20) + 7, 333_333].iter().enumerate() {
1147                let payload = vec![chunk * 16 + index as u8; *len];
1148                writer.append(&payload).unwrap();
1149                payloads.push(payload);
1150            }
1151            expected.push(payloads);
1152            handles.push(writer.finish().unwrap());
1153        }
1154        assert_eq!(manager.stats().files_live, 2);
1155        for (handle, payloads) in handles.iter().zip(expected.iter()) {
1156            let frames: Vec<Vec<u8>> = handle.reader().unwrap().collect::<Result<_, _>>().unwrap();
1157            assert_eq!(&frames, payloads);
1158        }
1159    }
1160
1161    #[test]
1162    fn checksum_mismatch_is_detected_on_read() {
1163        let dir = tempfile::tempdir().unwrap();
1164        let manager = manager(&dir, 1 << 20);
1165        let query_id = QueryId::new_random();
1166        let session = manager.begin_query(query_id, 1 << 20).unwrap();
1167        let mut writer = session.new_writer().unwrap();
1168        writer.append(&vec![0x11; 256]).unwrap();
1169        writer.append(b"second").unwrap();
1170        let handle = writer.finish().unwrap();
1171
1172        // Flip one byte in the first frame's payload (header 12 + head 17).
1173        let path = only_file(&dir, query_id);
1174        let mut file = std::fs::OpenOptions::new()
1175            .read(true)
1176            .write(true)
1177            .open(&path)
1178            .unwrap();
1179        file.seek(SeekFrom::Start((HEADER_LEN + FRAME_HEAD_LEN) as u64))
1180            .unwrap();
1181        file.write_all(&[0x99]).unwrap();
1182        drop(file);
1183
1184        let mut reader = handle.reader().unwrap();
1185        let error = reader.next_frame().unwrap_err();
1186        assert!(
1187            matches!(error, SpillError::ChecksumMismatch { .. }),
1188            "expected ChecksumMismatch, got {error:?}"
1189        );
1190        // The failure is terminal: no further frames are yielded.
1191        assert!(reader.next_frame().unwrap().is_none());
1192    }
1193
1194    #[test]
1195    fn trailer_mismatch_is_detected() {
1196        let dir = tempfile::tempdir().unwrap();
1197        let manager = manager(&dir, 1 << 20);
1198        // A crafted file: one valid data frame, then a trailer claiming two
1199        // frames — all CRCs valid, only the trailer semantics are wrong.
1200        let mut bytes = crafted_header(ENC_PLAINTEXT);
1201        bytes.extend_from_slice(&crafted_frame(FRAME_DATA, 0, b"abc"));
1202        let mut trailer = Vec::with_capacity(TRAILER_LEN);
1203        trailer.extend_from_slice(&2u64.to_le_bytes()); // wrong frame count
1204        trailer.extend_from_slice(&3u64.to_le_bytes());
1205        trailer.extend_from_slice(&[0u8; 32]);
1206        bytes.extend_from_slice(&crafted_frame(FRAME_TRAILER, 1, &trailer));
1207        let path = dir.path().join("crafted.spill");
1208        std::fs::write(&path, &bytes).unwrap();
1209
1210        let file = std::fs::File::open(&path).unwrap();
1211        let mut reader = reader_from(file, &manager).unwrap();
1212        assert_eq!(reader.next_frame().unwrap(), Some(b"abc".to_vec()));
1213        let error = reader.next_frame().unwrap_err();
1214        assert!(
1215            matches!(error, SpillError::Corrupt(_)),
1216            "expected Corrupt, got {error:?}"
1217        );
1218    }
1219
1220    #[test]
1221    fn bad_magic_and_version_are_rejected() {
1222        let dir = tempfile::tempdir().unwrap();
1223        let manager = manager(&dir, 1 << 20);
1224        for (label, mut header) in [
1225            ("magic", crafted_header(ENC_PLAINTEXT)),
1226            ("version", crafted_header(ENC_PLAINTEXT)),
1227            ("enc flag", crafted_header(ENC_PLAINTEXT)),
1228        ] {
1229            match label {
1230                "magic" => header[0] ^= 0xFF,
1231                "version" => header[8..10].copy_from_slice(&99u16.to_le_bytes()),
1232                _ => header[10] = 77,
1233            }
1234            let path = dir.path().join(format!("{label}.spill"));
1235            std::fs::write(&path, &header).unwrap();
1236            let file = std::fs::File::open(&path).unwrap();
1237            let error = reader_from(file, &manager).unwrap_err();
1238            assert!(
1239                matches!(error, SpillError::Corrupt(_)),
1240                "{label}: expected Corrupt, got {error:?}"
1241            );
1242        }
1243    }
1244
1245    #[test]
1246    fn per_query_budget_is_enforced() {
1247        let dir = tempfile::tempdir().unwrap();
1248        let manager = manager(&dir, 1 << 20);
1249        let query_id = QueryId::new_random();
1250        let session = manager.begin_query(query_id, 100).unwrap();
1251        let mut writer = session.new_writer().unwrap(); // header: 12 bytes
1252        writer.append(&[0u8; 50]).unwrap(); // frame: 17 + 50 = 67 → 79 used
1253        assert_eq!(session.used(), 79);
1254        assert_eq!(session.budget_remaining(), 21);
1255        let error = writer.append(&[0u8; 50]).unwrap_err();
1256        assert!(
1257            matches!(
1258                error,
1259                SpillError::BudgetExceeded {
1260                    query_id: id,
1261                    requested,
1262                    query_remaining: 21,
1263                    ..
1264                } if id == query_id && requested == 67
1265            ),
1266            "expected BudgetExceeded, got {error:?}"
1267        );
1268        // The failed frame charged nothing; a fitting frame still lands.
1269        assert_eq!(session.used(), 79);
1270        writer.append(&[0u8; 4]).unwrap(); // 17 + 4 = 21 → exactly the cap
1271        assert_eq!(session.used(), 100);
1272        assert_eq!(session.budget_remaining(), 0);
1273        let handle = writer.finish().unwrap_err();
1274        // Even the trailer no longer fits.
1275        assert!(matches!(handle, SpillError::BudgetExceeded { .. }));
1276    }
1277
1278    #[test]
1279    fn global_budget_is_enforced_across_queries() {
1280        let dir = tempfile::tempdir().unwrap();
1281        let manager = manager(&dir, 200);
1282        let a = manager.begin_query(QueryId::new_random(), 1 << 20).unwrap();
1283        let b = manager.begin_query(QueryId::new_random(), 1 << 20).unwrap();
1284        let mut writer_a = a.new_writer().unwrap();
1285        let mut writer_b = b.new_writer().unwrap();
1286        // Headers: 2 × 12 = 24 bytes of the global budget.
1287        writer_a.append(&[0u8; 100]).unwrap(); // +117 → 141 global
1288        let error = writer_b.append(&[0u8; 100]).unwrap_err(); // +117 > 200
1289        assert!(
1290            matches!(
1291                error,
1292                SpillError::BudgetExceeded {
1293                    global_remaining: 59,
1294                    ..
1295                }
1296            ),
1297            "expected global BudgetExceeded, got {error:?}"
1298        );
1299        assert_eq!(manager.stats().global_used, 141);
1300        // Releasing the first query's file re-opens the global budget.
1301        drop(writer_a);
1302        assert_eq!(manager.stats().global_used, 12);
1303        writer_b.append(&[0u8; 100]).unwrap();
1304    }
1305
1306    #[test]
1307    fn unfinished_writer_drop_deletes_file_and_releases_budget() {
1308        let dir = tempfile::tempdir().unwrap();
1309        let manager = manager(&dir, 1 << 20);
1310        let query_id = QueryId::new_random();
1311        let session = manager.begin_query(query_id, 1 << 20).unwrap();
1312        let mut writer = session.new_writer().unwrap();
1313        writer.append(&[0u8; 100]).unwrap();
1314        let path = only_file(&dir, query_id);
1315        assert!(path.exists());
1316        let used = session.used();
1317        assert!(used > 0);
1318        drop(writer); // the cancel path
1319        assert!(!path.exists(), "partial spill file must be deleted");
1320        assert_eq!(session.used(), 0);
1321        assert_eq!(manager.stats().global_used, 0);
1322        assert_eq!(manager.stats().files_live, 0);
1323    }
1324
1325    #[test]
1326    fn explicit_abort_deletes_and_reports() {
1327        let dir = tempfile::tempdir().unwrap();
1328        let manager = manager(&dir, 1 << 20);
1329        let query_id = QueryId::new_random();
1330        let session = manager.begin_query(query_id, 1 << 20).unwrap();
1331        let mut writer = session.new_writer().unwrap();
1332        writer.append(&[0u8; 64]).unwrap();
1333        let path = only_file(&dir, query_id);
1334        writer.abort().unwrap();
1335        assert!(!path.exists());
1336        assert_eq!(session.used(), 0);
1337        assert_eq!(manager.stats().global_used, 0);
1338    }
1339
1340    #[test]
1341    fn finished_handle_drop_and_delete_release_everything() {
1342        let dir = tempfile::tempdir().unwrap();
1343        let manager = manager(&dir, 1 << 20);
1344        let query_id = QueryId::new_random();
1345        let session = manager.begin_query(query_id, 1 << 20).unwrap();
1346
1347        let mut writer = session.new_writer().unwrap();
1348        writer.append(b"payload").unwrap();
1349        let handle = writer.finish().unwrap();
1350        let path = only_file(&dir, query_id);
1351        assert_eq!(manager.stats().files_live, 1);
1352        let used = session.used();
1353        drop(handle);
1354        assert!(!path.exists(), "sealed spill file must be deleted on drop");
1355        assert_eq!(session.used(), 0);
1356        assert_eq!(manager.stats().files_live, 0);
1357        assert_eq!(manager.stats().global_used, 0);
1358
1359        // Explicit delete reports and releases the same way.
1360        let mut writer = session.new_writer().unwrap();
1361        writer.append(b"payload").unwrap();
1362        let handle = writer.finish().unwrap();
1363        assert_eq!(session.used(), used);
1364        handle.delete().unwrap();
1365        assert_eq!(session.used(), 0);
1366        assert_eq!(manager.stats().files_live, 0);
1367    }
1368
1369    #[test]
1370    fn session_drop_removes_the_query_directory() {
1371        let dir = tempfile::tempdir().unwrap();
1372        let manager = manager(&dir, 1 << 20);
1373        let query_id = QueryId::new_random();
1374        let session = manager.begin_query(query_id, 1 << 20).unwrap();
1375        let mut writer = session.new_writer().unwrap();
1376        writer.append(b"x").unwrap();
1377        drop(writer);
1378        assert!(query_dir(&dir, query_id).exists());
1379        drop(session);
1380        assert!(
1381            !query_dir(&dir, query_id).exists(),
1382            "session drop must remove the per-query directory"
1383        );
1384        // A session that never spilled leaves nothing behind either.
1385        let quiet = manager.begin_query(QueryId::new_random(), 1 << 20).unwrap();
1386        drop(quiet);
1387        assert!(
1388            !dir.path().join("temp").join("spill").exists() || {
1389                std::fs::read_dir(dir.path().join("temp").join("spill"))
1390                    .unwrap()
1391                    .next()
1392                    .is_none()
1393            }
1394        );
1395    }
1396
1397    #[test]
1398    fn open_sweeps_stale_entries_from_prior_runs() {
1399        let dir = tempfile::tempdir().unwrap();
1400        let query_id = QueryId::new_random();
1401        let stale_dir = query_dir(&dir, query_id);
1402        let stale_path = stale_dir.join("chunk-000000.spill");
1403        std::fs::create_dir_all(&stale_dir).unwrap();
1404        std::fs::write(&stale_path, b"from a previous process").unwrap();
1405        // A stray file directly in the spill root must be swept too.
1406        std::fs::write(dir.path().join("temp").join("spill").join("stray"), b"x").unwrap();
1407        assert!(stale_path.exists());
1408
1409        // No prior-process handles remain: process exit closed them.
1410        let second = manager(&dir, 1 << 20);
1411        assert!(
1412            !stale_path.exists(),
1413            "startup sweep must remove stale files"
1414        );
1415        assert!(!dir.path().join("temp").join("spill").join("stray").exists());
1416        assert!(!query_dir(&dir, query_id).exists());
1417        assert_eq!(second.stats().global_used, 0);
1418        assert_eq!(second.stats().files_live, 0);
1419    }
1420
1421    #[test]
1422    fn frame_larger_than_the_limit_is_rejected() {
1423        let dir = tempfile::tempdir().unwrap();
1424        let manager = manager(&dir, u64::MAX);
1425        let session = manager
1426            .begin_query(QueryId::new_random(), u64::MAX)
1427            .unwrap();
1428        let mut writer = session.new_writer().unwrap();
1429        let huge = vec![0u8; MAX_FRAME_PAYLOAD as usize + 1];
1430        let error = writer.append(&huge).unwrap_err();
1431        assert!(
1432            matches!(error, SpillError::FrameTooLarge { .. }),
1433            "expected FrameTooLarge, got {error:?}"
1434        );
1435    }
1436
1437    #[test]
1438    fn invalid_config_is_rejected() {
1439        let dir = tempfile::tempdir().unwrap();
1440        let root = DurableRoot::open(dir.path()).unwrap();
1441        let error = SpillManager::open(&root, SpillConfig::new(0), None).unwrap_err();
1442        assert!(matches!(error, SpillError::InvalidConfig(_)));
1443    }
1444
1445    #[test]
1446    fn spill_error_maps_to_mongrel_error() {
1447        let io: MongrelError = SpillError::Io(io::Error::other("x")).into();
1448        assert!(matches!(io, MongrelError::Io(_)));
1449        let budget: MongrelError = SpillError::BudgetExceeded {
1450            query_id: QueryId::new_random(),
1451            requested: 10,
1452            query_remaining: 5,
1453            global_remaining: 90,
1454        }
1455        .into();
1456        assert!(matches!(budget, MongrelError::ResourceLimitExceeded { .. }));
1457        let checksum: MongrelError = SpillError::ChecksumMismatch {
1458            context: "frame".into(),
1459            expected: 1,
1460            actual: 2,
1461        }
1462        .into();
1463        assert!(matches!(checksum, MongrelError::ChecksumMismatch { .. }));
1464        let corrupt: MongrelError = SpillError::Corrupt("bad".into()).into();
1465        assert!(matches!(corrupt, MongrelError::Other(_)));
1466    }
1467
1468    mod encrypted {
1469        use super::*;
1470        use crate::encryption::{meta_dek_for, Kek, SALT_LEN};
1471
1472        fn encrypted_manager(dir: &tempfile::TempDir, dek: [u8; DEK_LEN]) -> SpillManager {
1473            let root = DurableRoot::open(dir.path()).unwrap();
1474            SpillManager::open(&root, SpillConfig::new(1 << 20), Some(dek)).unwrap()
1475        }
1476
1477        fn test_dek(passphrase: &str) -> [u8; DEK_LEN] {
1478            let salt = [7u8; SALT_LEN];
1479            let kek = Kek::derive(passphrase, &salt).unwrap();
1480            meta_dek_for(Some(&kek)).unwrap()
1481        }
1482
1483        #[test]
1484        fn encrypted_round_trip_seals_every_frame_on_disk() {
1485            let dir = tempfile::tempdir().unwrap();
1486            let manager = encrypted_manager(&dir, test_dek("pw"));
1487            let query_id = QueryId::new_random();
1488            let session = manager.begin_query(query_id, 1 << 20).unwrap();
1489            let marker = b"highly-recognizable-plaintext-marker";
1490            let mut writer = session.new_writer().unwrap();
1491            writer.append(marker).unwrap();
1492            writer.append(&vec![0x5A; 4096]).unwrap();
1493            let handle = writer.finish().unwrap();
1494
1495            // Nothing on disk is the plaintext: marker and payload absent.
1496            let raw = std::fs::read(only_file(&dir, query_id)).unwrap();
1497            assert_eq!(raw[10], ENC_AES_GCM);
1498            assert!(!raw
1499                .windows(marker.len())
1500                .any(|window| window == marker.as_slice()));
1501            assert!(!raw.windows(64).any(|window| window == [0x5A; 64]));
1502
1503            let frames: Vec<Vec<u8>> = handle.reader().unwrap().collect::<Result<_, _>>().unwrap();
1504            assert_eq!(frames, vec![marker.to_vec(), vec![0x5A; 4096]]);
1505        }
1506
1507        #[test]
1508        fn encrypted_file_requires_the_key_and_detects_tampering() {
1509            let dir = tempfile::tempdir().unwrap();
1510            // Every manager is opened before any spill file exists: opening a
1511            // manager sweeps stale entries, so opening one after the file
1512            // lands would delete it (that behavior has its own test).
1513            let plaintext_manager = manager(&dir, 1 << 20);
1514            let wrong = encrypted_manager(&dir, test_dek("wrong"));
1515            let manager = encrypted_manager(&dir, test_dek("pw"));
1516            let query_id = QueryId::new_random();
1517            let session = manager.begin_query(query_id, 1 << 20).unwrap();
1518            let mut writer = session.new_writer().unwrap();
1519            writer.append(b"secret").unwrap();
1520            let handle = writer.finish().unwrap();
1521            let path = only_file(&dir, query_id);
1522            // A manager without the DEK refuses the encrypted file.
1523            let error =
1524                reader_from(std::fs::File::open(&path).unwrap(), &plaintext_manager).unwrap_err();
1525            assert!(
1526                matches!(error, SpillError::EncryptionRequired),
1527                "expected EncryptionRequired, got {error:?}"
1528            );
1529
1530            // A wrong DEK fails at the GCM tag (decryption), not at parse.
1531            let mut reader = reader_from(std::fs::File::open(&path).unwrap(), &wrong).unwrap();
1532            let error = reader.next_frame().unwrap_err();
1533            assert!(
1534                matches!(error, SpillError::Decryption(_)),
1535                "expected Decryption, got {error:?}"
1536            );
1537
1538            // The right key still reads.
1539            let frames: Vec<Vec<u8>> = handle.reader().unwrap().collect::<Result<_, _>>().unwrap();
1540            assert_eq!(frames, vec![b"secret".to_vec()]);
1541        }
1542    }
1543}