Skip to main content

deepstrike_sdk/runtime/
payload_store.rs

1use 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}