deepstrike_sdk/runtime/
payload_store.rs1use std::fs;
2use std::io::Write;
3use std::path::{Path, PathBuf};
4use std::sync::atomic::{AtomicU64, Ordering};
5
6use deepstrike_core::runtime::kernel::wire::canonical_digest;
7
8static NEXT_TEMP_ID: AtomicU64 = AtomicU64::new(1);
9
10pub trait PayloadStore: Send + Sync {
11 fn persist(&self, session_id: &str, payload_ref: &str, content: &str) -> std::io::Result<()>;
12 fn load(&self, session_id: &str, payload_ref: &str) -> std::io::Result<Option<String>>;
13}
14
15pub struct FilePayloadStore {
16 root: PathBuf,
17}
18
19impl FilePayloadStore {
20 pub fn new(root: impl Into<PathBuf>) -> Self {
21 Self { root: root.into() }
22 }
23
24 fn path(&self, session_id: &str, payload_ref: &str) -> PathBuf {
25 let identity = format!("{session_id}\0{payload_ref}");
26 let digest = canonical_digest(identity.as_bytes());
27 self.root.join(format!(
28 "{}.payload",
29 digest.as_str().trim_start_matches("sha256:")
30 ))
31 }
32}
33
34impl PayloadStore for FilePayloadStore {
35 fn persist(&self, session_id: &str, payload_ref: &str, content: &str) -> std::io::Result<()> {
36 fs::create_dir_all(&self.root)?;
37 let target = self.path(session_id, payload_ref);
38 let temporary = temporary_path(&target);
39 let result = (|| {
40 let mut file = fs::OpenOptions::new()
41 .create_new(true)
42 .write(true)
43 .open(&temporary)?;
44 file.write_all(content.as_bytes())?;
45 file.sync_all()?;
46 drop(file);
47 fs::rename(&temporary, target)
48 })();
49 if result.is_err() {
50 let _ = fs::remove_file(temporary);
51 }
52 result
53 }
54
55 fn load(&self, session_id: &str, payload_ref: &str) -> std::io::Result<Option<String>> {
56 match fs::read_to_string(self.path(session_id, payload_ref)) {
57 Ok(content) => Ok(Some(content)),
58 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
59 Err(error) => Err(error),
60 }
61 }
62}
63
64fn temporary_path(target: &Path) -> PathBuf {
65 target.with_extension(format!(
66 "tmp-{}-{}",
67 std::process::id(),
68 NEXT_TEMP_ID.fetch_add(1, Ordering::Relaxed)
69 ))
70}
71
72#[cfg(test)]
73mod tests {
74 use super::{FilePayloadStore, PayloadStore};
75 use std::fs;
76
77 #[test]
78 fn opaque_locators_are_hashed_and_session_scoped() {
79 let root = std::env::temp_dir().join(format!(
80 "deepstrike-payload-test-{}-{}",
81 std::process::id(),
82 super::NEXT_TEMP_ID.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
83 ));
84 let store = FilePayloadStore::new(&root);
85 store
86 .persist("session-a", "../../payload", "alpha")
87 .expect("persist a");
88 store
89 .persist("session-b", "../../payload", "beta")
90 .expect("persist b");
91
92 assert_eq!(
93 store.load("session-a", "../../payload").expect("load a"),
94 Some("alpha".to_string())
95 );
96 assert_eq!(
97 store.load("session-b", "../../payload").expect("load b"),
98 Some("beta".to_string())
99 );
100 assert_eq!(
101 store.load("session-c", "../../payload").expect("load c"),
102 None
103 );
104 assert_eq!(fs::read_dir(&root).expect("list").count(), 2);
105 fs::remove_dir_all(root).expect("cleanup");
106 }
107}