diff options
| author | Somhairle H. Marisol <[email protected]> | 2026-09-17 14:32:37 +0800 |
|---|---|---|
| committer | Somhairle H. Marisol <[email protected]> | 2026-09-17 14:32:37 +0800 |
| commit | 5c0ba37eda80d39e6ceca59bb1d5f4942f858995 (patch) | |
| tree | 948723f9cedf7ccb0707fa6ee516bd30fe20fd10 /server/src/store.rs | |
| download | strategy-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.rs | 206 |
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))) + } +} |
