Some checks failed
ci/hello Job 'hello' completed
ci/smoke Pipeline 'smoke' success
ci/check Job 'check' completed
ci/clippy Job 'clippy' completed
ci/ci Pipeline 'ci' failure
Content addressable datastore. Allows you to configure a node to store and stream large blobs of data, and retrieve them from any swactor-connected node. Co-authored-by: Zachery Aaron Shores-Chmielewski <zacheryasc@gmail.com> Co-committed-by: Zachery Aaron Shores-Chmielewski <zacheryasc@gmail.com>
175 lines
5.7 KiB
Rust
175 lines
5.7 KiB
Rust
//! StorageBackend trait and implementations.
|
|
//!
|
|
//! Abstracts chunk and manifest I/O so backends can be swapped
|
|
//! (filesystem for MVP, IndexedDB for browser, in-memory for tests/WASM).
|
|
|
|
pub mod in_memory;
|
|
|
|
use std::collections::HashSet;
|
|
use std::fs;
|
|
use std::io::Write;
|
|
use std::path::PathBuf;
|
|
|
|
use crate::types::{ContentHash, ObjectManifest};
|
|
|
|
pub use in_memory::InMemoryBackend;
|
|
|
|
/// Pluggable storage backend for chunks and manifests.
|
|
pub trait StorageBackend: Send {
|
|
fn write_chunk(&mut self, hash: &ContentHash, data: &[u8]) -> Result<(), std::io::Error>;
|
|
fn read_chunk(&self, hash: &ContentHash) -> Result<Option<Vec<u8>>, std::io::Error>;
|
|
fn delete_chunk(&mut self, hash: &ContentHash) -> Result<(), std::io::Error>;
|
|
fn has_chunk(&self, hash: &ContentHash) -> bool;
|
|
fn list_chunks(&self) -> Vec<ContentHash>;
|
|
fn write_manifest(&mut self, manifest: &ObjectManifest) -> Result<(), std::io::Error>;
|
|
fn read_manifest(&self, content_hash: &ContentHash) -> Result<Option<ObjectManifest>, std::io::Error>;
|
|
fn delete_manifest(&mut self, content_hash: &ContentHash) -> Result<(), std::io::Error>;
|
|
}
|
|
|
|
/// Filesystem-backed storage with 2-level directory sharding.
|
|
///
|
|
/// Layout:
|
|
/// ```text
|
|
/// {root}/
|
|
/// ├── chunks/{hex[0..2]}/{hex[2..4]}/{full_hex_hash}
|
|
/// └── manifests/{hex[0..2]}/{hex[2..4]}/{full_hex_hash}
|
|
/// ```
|
|
pub struct FilesystemBackend {
|
|
root: PathBuf,
|
|
chunk_index: HashSet<ContentHash>,
|
|
}
|
|
|
|
impl FilesystemBackend {
|
|
pub fn new(root: PathBuf) -> Self {
|
|
let mut backend = Self {
|
|
root,
|
|
chunk_index: HashSet::new(),
|
|
};
|
|
backend.scan_chunks();
|
|
backend
|
|
}
|
|
|
|
fn chunk_path(&self, hash: &ContentHash) -> PathBuf {
|
|
let hex = hash.to_hex();
|
|
self.root
|
|
.join("chunks")
|
|
.join(&hex[..2])
|
|
.join(&hex[2..4])
|
|
.join(&hex)
|
|
}
|
|
|
|
fn manifest_path(&self, hash: &ContentHash) -> PathBuf {
|
|
let hex = hash.to_hex();
|
|
self.root
|
|
.join("manifests")
|
|
.join(&hex[..2])
|
|
.join(&hex[2..4])
|
|
.join(&hex)
|
|
}
|
|
|
|
fn scan_chunks(&mut self) {
|
|
let chunks_dir = self.root.join("chunks");
|
|
if !chunks_dir.exists() {
|
|
return;
|
|
}
|
|
let Ok(level1) = fs::read_dir(&chunks_dir) else {
|
|
return;
|
|
};
|
|
for d1 in level1.flatten() {
|
|
let Ok(level2) = fs::read_dir(d1.path()) else {
|
|
continue;
|
|
};
|
|
for d2 in level2.flatten() {
|
|
let Ok(files) = fs::read_dir(d2.path()) else {
|
|
continue;
|
|
};
|
|
for file in files.flatten() {
|
|
if let Some(name) = file.file_name().to_str() {
|
|
if let Some(hash) = ContentHash::from_hex(name) {
|
|
self.chunk_index.insert(hash);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
fn write_and_sync(path: &PathBuf, data: &[u8]) -> std::io::Result<()> {
|
|
if let Some(parent) = path.parent() {
|
|
fs::create_dir_all(parent)?;
|
|
}
|
|
let mut file = fs::File::create(path)?;
|
|
file.write_all(data)?;
|
|
file.sync_all()?;
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
impl StorageBackend for FilesystemBackend {
|
|
fn write_chunk(&mut self, hash: &ContentHash, data: &[u8]) -> Result<(), std::io::Error> {
|
|
let path = self.chunk_path(hash);
|
|
Self::write_and_sync(&path, data)?;
|
|
self.chunk_index.insert(*hash);
|
|
Ok(())
|
|
}
|
|
|
|
fn read_chunk(&self, hash: &ContentHash) -> Result<Option<Vec<u8>>, std::io::Error> {
|
|
if !self.chunk_index.contains(hash) {
|
|
return Ok(None);
|
|
}
|
|
let path = self.chunk_path(hash);
|
|
match fs::read(&path) {
|
|
Ok(data) => Ok(Some(data)),
|
|
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
|
|
Err(e) => Err(e),
|
|
}
|
|
}
|
|
|
|
fn delete_chunk(&mut self, hash: &ContentHash) -> Result<(), std::io::Error> {
|
|
self.chunk_index.remove(hash);
|
|
let path = self.chunk_path(hash);
|
|
match fs::remove_file(&path) {
|
|
Ok(()) => Ok(()),
|
|
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
|
|
Err(e) => Err(e),
|
|
}
|
|
}
|
|
|
|
fn has_chunk(&self, hash: &ContentHash) -> bool {
|
|
self.chunk_index.contains(hash)
|
|
}
|
|
|
|
fn list_chunks(&self) -> Vec<ContentHash> {
|
|
self.chunk_index.iter().copied().collect()
|
|
}
|
|
|
|
fn write_manifest(&mut self, manifest: &ObjectManifest) -> Result<(), std::io::Error> {
|
|
let path = self.manifest_path(&manifest.content_hash);
|
|
let data = serde_json::to_vec(manifest)
|
|
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
|
|
Self::write_and_sync(&path, &data)
|
|
}
|
|
|
|
fn read_manifest(&self, content_hash: &ContentHash) -> Result<Option<ObjectManifest>, std::io::Error> {
|
|
let path = self.manifest_path(content_hash);
|
|
match fs::read(&path) {
|
|
Ok(data) => {
|
|
let manifest: ObjectManifest = serde_json::from_slice(&data)
|
|
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
|
|
Ok(Some(manifest))
|
|
}
|
|
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
|
|
Err(e) => Err(e),
|
|
}
|
|
}
|
|
|
|
fn delete_manifest(&mut self, content_hash: &ContentHash) -> Result<(), std::io::Error> {
|
|
let path = self.manifest_path(content_hash);
|
|
match fs::remove_file(&path) {
|
|
Ok(()) => Ok(()),
|
|
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
|
|
Err(e) => Err(e),
|
|
}
|
|
}
|
|
}
|
|
|