1use 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
61const SPILL_DIR_REL: &str = "temp/spill";
63const SPILL_MAGIC: &[u8; 8] = b"MDBSPILL";
65const FORMAT_VERSION: u16 = 1;
67const ENC_PLAINTEXT: u8 = 0;
69const ENC_AES_GCM: u8 = 1;
71const HEADER_LEN: usize = 8 + 2 + 1 + 1;
73const FRAME_HEAD_LEN: usize = 4 + 4 + 1 + 8;
75const FRAME_DATA: u8 = 0;
77const FRAME_TRAILER: u8 = 1;
79const TRAILER_LEN: usize = 8 + 8 + 32;
81const MAX_FRAME_PAYLOAD: u64 = 64 * 1024 * 1024;
84
85const CRC32C: Crc<u32> = Crc::<u32>::new(&CRC_32_ISCSI);
86
87#[derive(Debug, thiserror::Error)]
89pub enum SpillError {
90 #[error("io error: {0}")]
92 Io(#[from] std::io::Error),
93 #[error("invalid spill configuration: {0}")]
95 InvalidConfig(&'static str),
96 #[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 query_id: QueryId,
105 requested: u64,
107 query_remaining: u64,
109 global_remaining: u64,
111 },
112 #[error("spill frame of {bytes} bytes exceeds the {limit}-byte frame limit")]
114 FrameTooLarge {
115 bytes: u64,
117 limit: u64,
119 },
120 #[error("checksum mismatch for {context}: expected {expected}, got {actual}")]
122 ChecksumMismatch {
123 context: String,
125 expected: u32,
127 actual: u32,
129 },
130 #[error("corrupt spill file: {0}")]
133 Corrupt(String),
134 #[error("encrypted spill file requires the database encryption key")]
136 EncryptionRequired,
137 #[error("spill encryption requires the `encryption` feature")]
139 EncryptionDisabled,
140 #[error("spill encryption error: {0}")]
142 Encryption(String),
143 #[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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
195pub struct SpillConfig {
196 pub global_bytes: u64,
198}
199
200impl SpillConfig {
201 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#[derive(Debug, Clone, PartialEq, Eq)]
218pub struct SpillStats {
219 pub bytes_written: u64,
221 pub bytes_read: u64,
223 pub files_live: u64,
225 pub global_used: u64,
227 pub global_budget_bytes: u64,
229 pub budget_remaining: u64,
231}
232
233struct ManagerInner {
234 db_root: DurableRoot,
235 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 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 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 fn sweep_stale(&self) -> Result<(), SpillError> {
271 match self.db_root.entry_exists(SPILL_DIR_REL) {
272 Ok(true) => {}
273 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 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 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 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
333pub 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 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 pub fn config(&self) -> &SpillConfig {
394 &self.inner.config
395 }
396
397 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 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 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 dir_name: String,
447 cap: u64,
448 used: AtomicU64,
449 next_chunk: AtomicU64,
450}
451
452impl SessionInner {
453 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 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
473pub 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 pub fn query_id(&self) -> QueryId {
492 self.inner.query_id
493 }
494
495 pub fn cap(&self) -> u64 {
497 self.inner.cap
498 }
499
500 pub fn used(&self) -> u64 {
502 self.inner.used.load(Ordering::Relaxed)
503 }
504
505 pub fn budget_remaining(&self) -> u64 {
507 self.inner.cap.saturating_sub(self.used())
508 }
509
510 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 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 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 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
623fn 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
639pub 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 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 pub fn bytes_on_disk(&self) -> u64 {
676 self.inner.as_ref().map_or(0, |inner| inner.bytes_on_disk)
677 }
678
679 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 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
754fn 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#[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 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 pub fn bytes_on_disk(&self) -> u64 {
803 self.inner.as_ref().map_or(0, |inner| inner.bytes_on_disk)
804 }
805
806 pub fn frames(&self) -> u64 {
808 self.inner.as_ref().map_or(0, |inner| inner.data_frames)
809 }
810
811 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 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
834fn 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
887pub 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 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
1025fn 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
1034fn 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
1044fn 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#[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 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(), 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 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 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 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 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()); 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(); writer.append(&[0u8; 50]).unwrap(); 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 assert_eq!(session.used(), 79);
1275 writer.append(&[0u8; 4]).unwrap(); assert_eq!(session.used(), 100);
1277 assert_eq!(session.budget_remaining(), 0);
1278 let handle = writer.finish().unwrap_err();
1279 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 writer_a.append(&[0u8; 100]).unwrap(); let error = writer_b.append(&[0u8; 100]).unwrap_err(); 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 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); 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 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 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 std::fs::write(dir.path().join("temp").join("spill").join("stray"), b"x").unwrap();
1417 std::mem::forget(session);
1419 }
1420 assert!(stale_path.exists());
1421 assert_eq!(first.stats().global_used, handle.bytes_on_disk());
1422
1423 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 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 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 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 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 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 let frames: Vec<Vec<u8>> = handle.reader().unwrap().collect::<Result<_, _>>().unwrap();
1559 assert_eq!(frames, vec![b"secret".to_vec()]);
1560 }
1561 }
1562}