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: opening a second manager on the
44//! same root sweeps the first manager's files, exactly like a restart.
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/// Fail-closed stub: [`SpillManager::open`] rejects a DEK without the
1045/// `encryption` feature, so this is unreachable in practice.
1046
1047/// Inverse of [`seal_payload`]: authenticates and opens a sealed frame.
1048fn open_payload(dek: Option<&[u8; DEK_LEN]>, stored: &[u8]) -> Result<Vec<u8>, SpillError> {
1049    match dek {
1050        Some(dek) => crate::encryption::decrypt_blob(dek, stored).map_err(map_crypto),
1051        None => Ok(stored.to_vec()),
1052    }
1053}
1054
1055/// Fail-closed stub — see [`seal_payload`].
1056
1057#[cfg(test)]
1058mod tests {
1059    use super::*;
1060    use std::io::{Seek, SeekFrom};
1061    use std::path::PathBuf;
1062
1063    fn manager(dir: &tempfile::TempDir, global_bytes: u64) -> SpillManager {
1064        let root = DurableRoot::open(dir.path()).unwrap();
1065        SpillManager::open(&root, SpillConfig::new(global_bytes), None).unwrap()
1066    }
1067
1068    fn query_dir(dir: &tempfile::TempDir, query_id: QueryId) -> PathBuf {
1069        dir.path()
1070            .join("temp")
1071            .join("spill")
1072            .join(format!("q-{}", query_id.to_hex()))
1073    }
1074
1075    fn only_file(dir: &tempfile::TempDir, query_id: QueryId) -> PathBuf {
1076        let entries: Vec<_> = std::fs::read_dir(query_dir(dir, query_id))
1077            .unwrap()
1078            .map(|entry| entry.unwrap().path())
1079            .collect();
1080        assert_eq!(entries.len(), 1, "expected exactly one spill file");
1081        entries[0].clone()
1082    }
1083
1084    /// Serializes one frame exactly as `WriterInner::write_frame` does
1085    /// (plaintext mode), for crafting corrupt files.
1086    fn crafted_frame(kind: u8, seq: u64, payload: &[u8]) -> Vec<u8> {
1087        let mut digest = CRC32C.digest();
1088        digest.update(&[kind]);
1089        digest.update(&seq.to_le_bytes());
1090        digest.update(payload);
1091        let crc = digest.finalize();
1092        let mut out = Vec::with_capacity(FRAME_HEAD_LEN + payload.len());
1093        out.extend_from_slice(&(payload.len() as u32).to_le_bytes());
1094        out.extend_from_slice(&crc.to_le_bytes());
1095        out.push(kind);
1096        out.extend_from_slice(&seq.to_le_bytes());
1097        out.extend_from_slice(payload);
1098        out
1099    }
1100
1101    fn crafted_header(enc: u8) -> Vec<u8> {
1102        let mut header = vec![0u8; HEADER_LEN];
1103        header[..8].copy_from_slice(SPILL_MAGIC);
1104        header[8..10].copy_from_slice(&FORMAT_VERSION.to_le_bytes());
1105        header[10] = enc;
1106        header
1107    }
1108
1109    #[test]
1110    fn frame_round_trip_streams_in_order() {
1111        let dir = tempfile::tempdir().unwrap();
1112        let manager = manager(&dir, 1 << 20);
1113        let session = manager.begin_query(QueryId::new_random(), 1 << 20).unwrap();
1114        let payloads: Vec<Vec<u8>> = vec![
1115            b"first".to_vec(),
1116            Vec::new(), // empty frames are legal
1117            vec![0xAB; 1000],
1118            b"last".to_vec(),
1119        ];
1120        let mut writer = session.new_writer().unwrap();
1121        for payload in &payloads {
1122            writer.append(payload).unwrap();
1123        }
1124        let handle = writer.finish().unwrap();
1125        assert_eq!(handle.frames(), 4);
1126        assert_eq!(handle.query_id(), session.query_id());
1127        assert_eq!(handle.bytes_on_disk(), session.used());
1128
1129        let frames: Vec<Vec<u8>> = handle.reader().unwrap().collect::<Result<_, _>>().unwrap();
1130        assert_eq!(frames, payloads);
1131
1132        let stats = manager.stats();
1133        assert_eq!(stats.files_live, 1);
1134        assert_eq!(stats.bytes_written, handle.bytes_on_disk());
1135        assert_eq!(stats.bytes_read, handle.bytes_on_disk());
1136        assert_eq!(stats.global_used, handle.bytes_on_disk());
1137        assert_eq!(stats.budget_remaining, (1 << 20) - handle.bytes_on_disk());
1138    }
1139
1140    #[test]
1141    fn large_multi_chunk_round_trip() {
1142        let dir = tempfile::tempdir().unwrap();
1143        let manager = manager(&dir, 1 << 24);
1144        let session = manager.begin_query(QueryId::new_random(), 1 << 24).unwrap();
1145        // Two chunk files, three frames of mixed large sizes each.
1146        let mut handles = Vec::new();
1147        let mut expected = Vec::new();
1148        for chunk in 0..2u8 {
1149            let mut writer = session.new_writer().unwrap();
1150            let mut payloads = Vec::new();
1151            for (index, len) in [(1usize << 20), (1 << 20) + 7, 333_333].iter().enumerate() {
1152                let payload = vec![chunk * 16 + index as u8; *len];
1153                writer.append(&payload).unwrap();
1154                payloads.push(payload);
1155            }
1156            expected.push(payloads);
1157            handles.push(writer.finish().unwrap());
1158        }
1159        assert_eq!(manager.stats().files_live, 2);
1160        for (handle, payloads) in handles.iter().zip(expected.iter()) {
1161            let frames: Vec<Vec<u8>> = handle.reader().unwrap().collect::<Result<_, _>>().unwrap();
1162            assert_eq!(&frames, payloads);
1163        }
1164    }
1165
1166    #[test]
1167    fn checksum_mismatch_is_detected_on_read() {
1168        let dir = tempfile::tempdir().unwrap();
1169        let manager = manager(&dir, 1 << 20);
1170        let query_id = QueryId::new_random();
1171        let session = manager.begin_query(query_id, 1 << 20).unwrap();
1172        let mut writer = session.new_writer().unwrap();
1173        writer.append(&vec![0x11; 256]).unwrap();
1174        writer.append(b"second").unwrap();
1175        let handle = writer.finish().unwrap();
1176
1177        // Flip one byte in the first frame's payload (header 12 + head 17).
1178        let path = only_file(&dir, query_id);
1179        let mut file = std::fs::OpenOptions::new()
1180            .read(true)
1181            .write(true)
1182            .open(&path)
1183            .unwrap();
1184        file.seek(SeekFrom::Start((HEADER_LEN + FRAME_HEAD_LEN) as u64))
1185            .unwrap();
1186        file.write_all(&[0x99]).unwrap();
1187        drop(file);
1188
1189        let mut reader = handle.reader().unwrap();
1190        let error = reader.next_frame().unwrap_err();
1191        assert!(
1192            matches!(error, SpillError::ChecksumMismatch { .. }),
1193            "expected ChecksumMismatch, got {error:?}"
1194        );
1195        // The failure is terminal: no further frames are yielded.
1196        assert!(reader.next_frame().unwrap().is_none());
1197    }
1198
1199    #[test]
1200    fn trailer_mismatch_is_detected() {
1201        let dir = tempfile::tempdir().unwrap();
1202        let manager = manager(&dir, 1 << 20);
1203        // A crafted file: one valid data frame, then a trailer claiming two
1204        // frames — all CRCs valid, only the trailer semantics are wrong.
1205        let mut bytes = crafted_header(ENC_PLAINTEXT);
1206        bytes.extend_from_slice(&crafted_frame(FRAME_DATA, 0, b"abc"));
1207        let mut trailer = Vec::with_capacity(TRAILER_LEN);
1208        trailer.extend_from_slice(&2u64.to_le_bytes()); // wrong frame count
1209        trailer.extend_from_slice(&3u64.to_le_bytes());
1210        trailer.extend_from_slice(&[0u8; 32]);
1211        bytes.extend_from_slice(&crafted_frame(FRAME_TRAILER, 1, &trailer));
1212        let path = dir.path().join("crafted.spill");
1213        std::fs::write(&path, &bytes).unwrap();
1214
1215        let file = std::fs::File::open(&path).unwrap();
1216        let mut reader = reader_from(file, &manager).unwrap();
1217        assert_eq!(reader.next_frame().unwrap(), Some(b"abc".to_vec()));
1218        let error = reader.next_frame().unwrap_err();
1219        assert!(
1220            matches!(error, SpillError::Corrupt(_)),
1221            "expected Corrupt, got {error:?}"
1222        );
1223    }
1224
1225    #[test]
1226    fn bad_magic_and_version_are_rejected() {
1227        let dir = tempfile::tempdir().unwrap();
1228        let manager = manager(&dir, 1 << 20);
1229        for (label, mut header) in [
1230            ("magic", crafted_header(ENC_PLAINTEXT)),
1231            ("version", crafted_header(ENC_PLAINTEXT)),
1232            ("enc flag", crafted_header(ENC_PLAINTEXT)),
1233        ] {
1234            match label {
1235                "magic" => header[0] ^= 0xFF,
1236                "version" => header[8..10].copy_from_slice(&99u16.to_le_bytes()),
1237                _ => header[10] = 77,
1238            }
1239            let path = dir.path().join(format!("{label}.spill"));
1240            std::fs::write(&path, &header).unwrap();
1241            let file = std::fs::File::open(&path).unwrap();
1242            let error = reader_from(file, &manager).unwrap_err();
1243            assert!(
1244                matches!(error, SpillError::Corrupt(_)),
1245                "{label}: expected Corrupt, got {error:?}"
1246            );
1247        }
1248    }
1249
1250    #[test]
1251    fn per_query_budget_is_enforced() {
1252        let dir = tempfile::tempdir().unwrap();
1253        let manager = manager(&dir, 1 << 20);
1254        let query_id = QueryId::new_random();
1255        let session = manager.begin_query(query_id, 100).unwrap();
1256        let mut writer = session.new_writer().unwrap(); // header: 12 bytes
1257        writer.append(&[0u8; 50]).unwrap(); // frame: 17 + 50 = 67 → 79 used
1258        assert_eq!(session.used(), 79);
1259        assert_eq!(session.budget_remaining(), 21);
1260        let error = writer.append(&[0u8; 50]).unwrap_err();
1261        assert!(
1262            matches!(
1263                error,
1264                SpillError::BudgetExceeded {
1265                    query_id: id,
1266                    requested,
1267                    query_remaining: 21,
1268                    ..
1269                } if id == query_id && requested == 67
1270            ),
1271            "expected BudgetExceeded, got {error:?}"
1272        );
1273        // The failed frame charged nothing; a fitting frame still lands.
1274        assert_eq!(session.used(), 79);
1275        writer.append(&[0u8; 4]).unwrap(); // 17 + 4 = 21 → exactly the cap
1276        assert_eq!(session.used(), 100);
1277        assert_eq!(session.budget_remaining(), 0);
1278        let handle = writer.finish().unwrap_err();
1279        // Even the trailer no longer fits.
1280        assert!(matches!(handle, SpillError::BudgetExceeded { .. }));
1281    }
1282
1283    #[test]
1284    fn global_budget_is_enforced_across_queries() {
1285        let dir = tempfile::tempdir().unwrap();
1286        let manager = manager(&dir, 200);
1287        let a = manager.begin_query(QueryId::new_random(), 1 << 20).unwrap();
1288        let b = manager.begin_query(QueryId::new_random(), 1 << 20).unwrap();
1289        let mut writer_a = a.new_writer().unwrap();
1290        let mut writer_b = b.new_writer().unwrap();
1291        // Headers: 2 × 12 = 24 bytes of the global budget.
1292        writer_a.append(&[0u8; 100]).unwrap(); // +117 → 141 global
1293        let error = writer_b.append(&[0u8; 100]).unwrap_err(); // +117 > 200
1294        assert!(
1295            matches!(
1296                error,
1297                SpillError::BudgetExceeded {
1298                    global_remaining: 59,
1299                    ..
1300                }
1301            ),
1302            "expected global BudgetExceeded, got {error:?}"
1303        );
1304        assert_eq!(manager.stats().global_used, 141);
1305        // Releasing the first query's file re-opens the global budget.
1306        drop(writer_a);
1307        assert_eq!(manager.stats().global_used, 12);
1308        writer_b.append(&[0u8; 100]).unwrap();
1309    }
1310
1311    #[test]
1312    fn unfinished_writer_drop_deletes_file_and_releases_budget() {
1313        let dir = tempfile::tempdir().unwrap();
1314        let manager = manager(&dir, 1 << 20);
1315        let query_id = QueryId::new_random();
1316        let session = manager.begin_query(query_id, 1 << 20).unwrap();
1317        let mut writer = session.new_writer().unwrap();
1318        writer.append(&[0u8; 100]).unwrap();
1319        let path = only_file(&dir, query_id);
1320        assert!(path.exists());
1321        let used = session.used();
1322        assert!(used > 0);
1323        drop(writer); // the cancel path
1324        assert!(!path.exists(), "partial spill file must be deleted");
1325        assert_eq!(session.used(), 0);
1326        assert_eq!(manager.stats().global_used, 0);
1327        assert_eq!(manager.stats().files_live, 0);
1328    }
1329
1330    #[test]
1331    fn explicit_abort_deletes_and_reports() {
1332        let dir = tempfile::tempdir().unwrap();
1333        let manager = manager(&dir, 1 << 20);
1334        let query_id = QueryId::new_random();
1335        let session = manager.begin_query(query_id, 1 << 20).unwrap();
1336        let mut writer = session.new_writer().unwrap();
1337        writer.append(&[0u8; 64]).unwrap();
1338        let path = only_file(&dir, query_id);
1339        writer.abort().unwrap();
1340        assert!(!path.exists());
1341        assert_eq!(session.used(), 0);
1342        assert_eq!(manager.stats().global_used, 0);
1343    }
1344
1345    #[test]
1346    fn finished_handle_drop_and_delete_release_everything() {
1347        let dir = tempfile::tempdir().unwrap();
1348        let manager = manager(&dir, 1 << 20);
1349        let query_id = QueryId::new_random();
1350        let session = manager.begin_query(query_id, 1 << 20).unwrap();
1351
1352        let mut writer = session.new_writer().unwrap();
1353        writer.append(b"payload").unwrap();
1354        let handle = writer.finish().unwrap();
1355        let path = only_file(&dir, query_id);
1356        assert_eq!(manager.stats().files_live, 1);
1357        let used = session.used();
1358        drop(handle);
1359        assert!(!path.exists(), "sealed spill file must be deleted on drop");
1360        assert_eq!(session.used(), 0);
1361        assert_eq!(manager.stats().files_live, 0);
1362        assert_eq!(manager.stats().global_used, 0);
1363
1364        // Explicit delete reports and releases the same way.
1365        let mut writer = session.new_writer().unwrap();
1366        writer.append(b"payload").unwrap();
1367        let handle = writer.finish().unwrap();
1368        assert_eq!(session.used(), used);
1369        handle.delete().unwrap();
1370        assert_eq!(session.used(), 0);
1371        assert_eq!(manager.stats().files_live, 0);
1372    }
1373
1374    #[test]
1375    fn session_drop_removes_the_query_directory() {
1376        let dir = tempfile::tempdir().unwrap();
1377        let manager = manager(&dir, 1 << 20);
1378        let query_id = QueryId::new_random();
1379        let session = manager.begin_query(query_id, 1 << 20).unwrap();
1380        let mut writer = session.new_writer().unwrap();
1381        writer.append(b"x").unwrap();
1382        drop(writer);
1383        assert!(query_dir(&dir, query_id).exists());
1384        drop(session);
1385        assert!(
1386            !query_dir(&dir, query_id).exists(),
1387            "session drop must remove the per-query directory"
1388        );
1389        // A session that never spilled leaves nothing behind either.
1390        let quiet = manager.begin_query(QueryId::new_random(), 1 << 20).unwrap();
1391        drop(quiet);
1392        assert!(
1393            !dir.path().join("temp").join("spill").exists() || {
1394                std::fs::read_dir(dir.path().join("temp").join("spill"))
1395                    .unwrap()
1396                    .next()
1397                    .is_none()
1398            }
1399        );
1400    }
1401
1402    #[test]
1403    fn open_sweeps_stale_entries_from_prior_runs() {
1404        let dir = tempfile::tempdir().unwrap();
1405        let query_id = QueryId::new_random();
1406        let stale_path;
1407        let first = manager(&dir, 1 << 20);
1408        let handle;
1409        {
1410            let session = first.begin_query(query_id, 1 << 20).unwrap();
1411            let mut writer = session.new_writer().unwrap();
1412            writer.append(b"from a previous process").unwrap();
1413            handle = writer.finish().unwrap();
1414            stale_path = only_file(&dir, query_id);
1415            // A stray file directly in the spill root must be swept too.
1416            std::fs::write(dir.path().join("temp").join("spill").join("stray"), b"x").unwrap();
1417            // Simulate a crash: the session and handle leak away.
1418            std::mem::forget(session);
1419        }
1420        assert!(stale_path.exists());
1421        assert_eq!(first.stats().global_used, handle.bytes_on_disk());
1422
1423        // The next process opens its manager: everything stale is swept.
1424        let second = manager(&dir, 1 << 20);
1425        assert!(
1426            !stale_path.exists(),
1427            "startup sweep must remove stale files"
1428        );
1429        assert!(!dir.path().join("temp").join("spill").join("stray").exists());
1430        assert!(!query_dir(&dir, query_id).exists());
1431        assert_eq!(second.stats().global_used, 0);
1432        assert_eq!(second.stats().files_live, 0);
1433
1434        // The leaked handle's drop is still safe (no double delete) and its
1435        // accounting unwinds against the first manager.
1436        drop(handle);
1437        assert_eq!(first.stats().global_used, 0);
1438    }
1439
1440    #[test]
1441    fn frame_larger_than_the_limit_is_rejected() {
1442        let dir = tempfile::tempdir().unwrap();
1443        let manager = manager(&dir, u64::MAX);
1444        let session = manager
1445            .begin_query(QueryId::new_random(), u64::MAX)
1446            .unwrap();
1447        let mut writer = session.new_writer().unwrap();
1448        let huge = vec![0u8; MAX_FRAME_PAYLOAD as usize + 1];
1449        let error = writer.append(&huge).unwrap_err();
1450        assert!(
1451            matches!(error, SpillError::FrameTooLarge { .. }),
1452            "expected FrameTooLarge, got {error:?}"
1453        );
1454    }
1455
1456    #[test]
1457    fn invalid_config_is_rejected() {
1458        let dir = tempfile::tempdir().unwrap();
1459        let root = DurableRoot::open(dir.path()).unwrap();
1460        let error = SpillManager::open(&root, SpillConfig::new(0), None).unwrap_err();
1461        assert!(matches!(error, SpillError::InvalidConfig(_)));
1462    }
1463
1464    #[test]
1465    fn spill_error_maps_to_mongrel_error() {
1466        let io: MongrelError = SpillError::Io(io::Error::other("x")).into();
1467        assert!(matches!(io, MongrelError::Io(_)));
1468        let budget: MongrelError = SpillError::BudgetExceeded {
1469            query_id: QueryId::new_random(),
1470            requested: 10,
1471            query_remaining: 5,
1472            global_remaining: 90,
1473        }
1474        .into();
1475        assert!(matches!(budget, MongrelError::ResourceLimitExceeded { .. }));
1476        let checksum: MongrelError = SpillError::ChecksumMismatch {
1477            context: "frame".into(),
1478            expected: 1,
1479            actual: 2,
1480        }
1481        .into();
1482        assert!(matches!(checksum, MongrelError::ChecksumMismatch { .. }));
1483        let corrupt: MongrelError = SpillError::Corrupt("bad".into()).into();
1484        assert!(matches!(corrupt, MongrelError::Other(_)));
1485    }
1486
1487    mod encrypted {
1488        use super::*;
1489        use crate::encryption::{meta_dek_for, Kek, SALT_LEN};
1490
1491        fn encrypted_manager(dir: &tempfile::TempDir, dek: [u8; DEK_LEN]) -> SpillManager {
1492            let root = DurableRoot::open(dir.path()).unwrap();
1493            SpillManager::open(&root, SpillConfig::new(1 << 20), Some(dek)).unwrap()
1494        }
1495
1496        fn test_dek(passphrase: &str) -> [u8; DEK_LEN] {
1497            let salt = [7u8; SALT_LEN];
1498            let kek = Kek::derive(passphrase, &salt).unwrap();
1499            meta_dek_for(Some(&kek)).unwrap()
1500        }
1501
1502        #[test]
1503        fn encrypted_round_trip_seals_every_frame_on_disk() {
1504            let dir = tempfile::tempdir().unwrap();
1505            let manager = encrypted_manager(&dir, test_dek("pw"));
1506            let query_id = QueryId::new_random();
1507            let session = manager.begin_query(query_id, 1 << 20).unwrap();
1508            let marker = b"highly-recognizable-plaintext-marker";
1509            let mut writer = session.new_writer().unwrap();
1510            writer.append(marker).unwrap();
1511            writer.append(&vec![0x5A; 4096]).unwrap();
1512            let handle = writer.finish().unwrap();
1513
1514            // Nothing on disk is the plaintext: marker and payload absent.
1515            let raw = std::fs::read(only_file(&dir, query_id)).unwrap();
1516            assert_eq!(raw[10], ENC_AES_GCM);
1517            assert!(!raw
1518                .windows(marker.len())
1519                .any(|window| window == marker.as_slice()));
1520            assert!(!raw.windows(64).any(|window| window == [0x5A; 64]));
1521
1522            let frames: Vec<Vec<u8>> = handle.reader().unwrap().collect::<Result<_, _>>().unwrap();
1523            assert_eq!(frames, vec![marker.to_vec(), vec![0x5A; 4096]]);
1524        }
1525
1526        #[test]
1527        fn encrypted_file_requires_the_key_and_detects_tampering() {
1528            let dir = tempfile::tempdir().unwrap();
1529            // Every manager is opened before any spill file exists: opening a
1530            // manager sweeps stale entries, so opening one after the file
1531            // lands would delete it (that behavior has its own test).
1532            let plaintext_manager = manager(&dir, 1 << 20);
1533            let wrong = encrypted_manager(&dir, test_dek("wrong"));
1534            let manager = encrypted_manager(&dir, test_dek("pw"));
1535            let query_id = QueryId::new_random();
1536            let session = manager.begin_query(query_id, 1 << 20).unwrap();
1537            let mut writer = session.new_writer().unwrap();
1538            writer.append(b"secret").unwrap();
1539            let handle = writer.finish().unwrap();
1540            let path = only_file(&dir, query_id);
1541            // A manager without the DEK refuses the encrypted file.
1542            let error =
1543                reader_from(std::fs::File::open(&path).unwrap(), &plaintext_manager).unwrap_err();
1544            assert!(
1545                matches!(error, SpillError::EncryptionRequired),
1546                "expected EncryptionRequired, got {error:?}"
1547            );
1548
1549            // A wrong DEK fails at the GCM tag (decryption), not at parse.
1550            let mut reader = reader_from(std::fs::File::open(&path).unwrap(), &wrong).unwrap();
1551            let error = reader.next_frame().unwrap_err();
1552            assert!(
1553                matches!(error, SpillError::Decryption(_)),
1554                "expected Decryption, got {error:?}"
1555            );
1556
1557            // The right key still reads.
1558            let frames: Vec<Vec<u8>> = handle.reader().unwrap().collect::<Result<_, _>>().unwrap();
1559            assert_eq!(frames, vec![b"secret".to_vec()]);
1560        }
1561    }
1562}