summaryrefslogtreecommitdiff
path: root/server/src/store.rs
diff options
context:
space:
mode:
authorSomhairle H. Marisol <[email protected]>2026-09-17 14:32:37 +0800
committerSomhairle H. Marisol <[email protected]>2026-09-17 14:32:37 +0800
commit5c0ba37eda80d39e6ceca59bb1d5f4942f858995 (patch)
tree948723f9cedf7ccb0707fa6ee516bd30fe20fd10 /server/src/store.rs
downloadstrategy-lab-5c0ba37eda80d39e6ceca59bb1d5f4942f858995.tar.gz
chore: establish Strategy Lab source baseline (development, not release)
Diffstat (limited to 'server/src/store.rs')
-rw-r--r--server/src/store.rs206
1 files changed, 206 insertions, 0 deletions
diff --git a/server/src/store.rs b/server/src/store.rs
new file mode 100644
index 0000000..f371d10
--- /dev/null
+++ b/server/src/store.rs
@@ -0,0 +1,206 @@
+use std::path::{Path, PathBuf};
+
+use sha2::{Digest, Sha256};
+
+/// Immutable content-addressed object storage under {data_dir}/objects/{hash[:2]}/{hash}.{ext}.
+/// Content is deduplicated: storing identical bytes twice keeps the first object.
+pub struct ObjectStore {
+ pub root: PathBuf,
+}
+
+/// One ingested artifact (file content hashed and copied into the object store).
+#[derive(Debug, Clone)]
+pub struct StoredObject {
+ /// sha256 of content
+ pub hash: String,
+ /// path relative to the object root (immutable stored path)
+ pub stored_path: String,
+ /// file name as produced by the worker, safe for /data mounts
+ pub mount_name: String,
+ pub size: u64,
+}
+
+impl ObjectStore {
+ pub fn new(data_dir: &str) -> Self {
+ ObjectStore { root: Path::new(data_dir).join("objects") }
+ }
+
+ /// Store raw bytes by content hash. Returns (hash, stored relative path).
+ pub fn store(&self, bytes: &[u8], filename: &str) -> std::io::Result<(String, String)> {
+ let hash = format!("{:x}", Sha256::digest(bytes));
+ let dir = self.root.join(&hash[..2]);
+ std::fs::create_dir_all(&dir)?;
+ let ext = safe_ext(filename);
+ let path = dir.join(format!("{hash}.{ext}"));
+ if !path.is_file() {
+ // Atomic write in the final directory; suffix append (not with_extension)
+ // so different source extensions cannot collide on one tmp name.
+ let tmp = dir.join(format!("{hash}.{ext}.tmp"));
+ std::fs::write(&tmp, bytes)?;
+ std::fs::rename(&tmp, &path)?;
+ }
+ let rel = path
+ .strip_prefix(&self.root)
+ .map(|p| p.to_string_lossy().into_owned())
+ .unwrap_or_else(|_| path.to_string_lossy().into_owned());
+ Ok((hash, rel))
+ }
+
+ /// Map a stored relative path (server controlled) back to an absolute path.
+ pub fn absolute(&self, rel: &str) -> PathBuf {
+ self.root.join(rel)
+ }
+
+ /// Hash and ingest every regular file under an output directory (worker
+ /// artifacts). Returns one entry per file, sorted for determinism.
+ /// Rejects symlinked entries.
+ pub fn ingest_directory(&self, dir: &Path) -> std::io::Result<Vec<StoredObject>> {
+ let mut files: Vec<PathBuf> = Vec::new();
+ collect_files(dir, dir, &mut files)?;
+ files.sort();
+ let mut out = Vec::with_capacity(files.len());
+ for f in files {
+ let is_symlink = f.symlink_metadata()?.file_type().is_symlink()
+ || std::fs::symlink_metadata(&f)?.file_type().is_symlink();
+ if is_symlink {
+ return Err(std::io::Error::new(
+ std::io::ErrorKind::InvalidInput,
+ "symlinked artifact rejected",
+ ));
+ }
+ let bytes = std::fs::read(&f)?;
+ let mount_name = f
+ .file_name()
+ .and_then(|n| n.to_str())
+ .unwrap_or("artifact.bin")
+ .to_string();
+ let rel_path = f.strip_prefix(dir).expect("strip_prefix");
+ let (hash, stored_path) =
+ self.store(&bytes, &mount_name)?;
+ out.push(StoredObject {
+ hash,
+ stored_path,
+ mount_name: rel_path.display().to_string(),
+ size: bytes.len() as u64,
+ });
+ }
+ Ok(out)
+ }
+}
+
+fn collect_files(_root: &Path, dir: &Path, out: &mut Vec<PathBuf>) -> std::io::Result<()> {
+ for entry in std::fs::read_dir(dir)? {
+ let p = entry?.path();
+ let ty = p.symlink_metadata()?.file_type();
+ if ty.is_symlink() {
+ // path traversal defense: no symlinked artifacts, ever
+ return Err(std::io::Error::new(
+ std::io::ErrorKind::InvalidInput,
+ "symlink in artifact tree rejected",
+ ));
+ }
+ if ty.is_dir() {
+ collect_files(_root, &p, out)?;
+ } else {
+ out.push(p);
+ }
+ }
+ Ok(())
+}
+
+/// Sanitize a filename for use as an object extension: only the final
+/// extension survives, restricted to alphanumeric chars.
+fn safe_ext(filename: &str) -> String {
+ let base = filename.rsplit('/').next().unwrap_or("data");
+ let e = base.rsplit('.').next().unwrap_or("bin").to_string();
+ let v: String = e.chars().filter(|c| c.is_ascii_alphanumeric()).collect();
+ if v.is_empty() || v.parse::<usize>().is_ok() {
+ "bin".into()
+ } else {
+ v.to_lowercase()
+ }
+}
+
+/// Public wrapper used when a filename has no usable extension.
+pub fn filename_or_bin(name: &str, fallback: &str) -> String {
+ let base = name.rsplit('/').next().unwrap_or(fallback);
+ let v: String = base
+ .chars()
+ .filter(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '-' | '_'))
+ .collect();
+ if v.is_empty() || v == "." {
+ fallback.into()
+ } else {
+ v
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ fn tmp_store(tag: &str) -> (tempfile::TempDir, ObjectStore) {
+ let t = tempfile::tempdir_in("/tmp/opencode").unwrap();
+ let store = ObjectStore::new(t.path().join(tag).to_str().unwrap());
+ (t, store)
+ }
+
+ #[test]
+ fn store_is_content_addressed_and_deduplicated() {
+ let (_t, s) = tmp_store("objs1");
+ let (h1, p1) = s.store(b"hello world", "a.csv").unwrap();
+ let (h2, p2) = s.store(b"hello world", "b.csv").unwrap();
+ assert_eq!(h1, h2);
+ assert_eq!(p1, p2);
+ assert!(p1.starts_with(&h1[..2]), "{p1}");
+ let abs = s.absolute(&p1);
+ assert_eq!(std::fs::read(&abs).unwrap(), b"hello world".to_vec());
+ assert!(s.root.join(&p1) == abs, "stored path resolves under root");
+ }
+
+ #[test]
+ fn same_stem_different_extension_no_collision() {
+ let (_t, s) = tmp_store("objs2");
+ let (_, a) = s.store(b"csv-bytes", "obj.csv").unwrap();
+ let (_, b) = s.store(b"json-bytes", "obj.json").unwrap();
+ assert_ne!(a, b);
+ assert!(!a.ends_with(".tmp"));
+ assert!(std::fs::read(s.absolute(&a)).unwrap().starts_with(b"csv"));
+ }
+
+ #[test]
+ fn ingest_directory_walks_and_rejects_symlinks() {
+ let td = tempfile::tempdir_in("/tmp/opencode").unwrap();
+ let out = td.path().join("out");
+ let objd = out.join("objects");
+ std::fs::create_dir_all(&objd).unwrap();
+ std::fs::write(objd.join("data.csv"), b"date,close\n2024-01-02,10\n").unwrap();
+ std::fs::write(out.join("result.json"), b"{\"status\":\"ready\"}").unwrap();
+ let (_t, s) = tmp_store("objs3");
+ let stored = s.ingest_directory(&out).unwrap();
+ assert_eq!(stored.len(), 2);
+ let names: Vec<&str> = stored.iter().map(|o| o.mount_name.as_str()).collect();
+ assert!(names.contains_all(&["objects/data.csv", "result.json"]), "{names:?}");
+ // same content re-ingested maps to the same stored object
+ let again = s.ingest_directory(&out).unwrap();
+ for o in &stored {
+ assert!(again.iter().any(|n| n.hash == o.hash));
+ }
+ // symlink rejection
+ std::os::unix::fs::symlink(
+ objd.join("data.csv"),
+ objd.join("data_link.csv"),
+ )
+ .unwrap();
+ assert!(s.ingest_directory(&out).is_err());
+ }
+}
+
+trait ContainsAll {
+ fn contains_all(&self, needles: &[&str]) -> bool;
+}
+impl ContainsAll for Vec<&str> {
+ fn contains_all(&self, needles: &[&str]) -> bool {
+ needles.iter().all(|n| self.iter().any(|m| m.contains(n)))
+ }
+}