use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use serde::Deserialize;
use sha2::{Digest, Sha256};
use crate::layer::LayerId;
use crate::manifest::ManifestEntry;
pub const PAYLOAD_DIR: &str = "payloads";
#[derive(Debug, Clone, Copy)]
pub struct Payload<'a> {
pub name: &'a str,
pub version: Option<&'a str>,
pub dispatchable: bool,
pub bytes: &'a [u8],
}
impl<'a> Payload<'a> {
pub fn tool(name: &'a str, bytes: &'a [u8]) -> Self {
Payload {
name,
version: None,
dispatchable: true,
bytes,
}
}
}
pub fn entry_is_dispatchable(entry: &ManifestEntry) -> bool {
matches!(entry.kind(), Ok(kind) if kind.is_dispatchable())
}
pub fn entry_version(entry: &ManifestEntry) -> Option<&str> {
entry
.annotations
.get("eu.pulseengine.tool.version")
.map(String::as_str)
}
fn safe_component(what: &'static str, value: &str) -> Result<(), StoreError> {
let bad = |why: &str| {
Err(StoreError::UnsafeComponent {
what,
value: value.to_string(),
why: why.to_string(),
})
};
if value.is_empty() {
return bad("empty");
}
if value == "." || value == ".." {
return bad("a relative path element");
}
if let Some(c) = value
.chars()
.find(|c| matches!(c, '/' | '\\' | '\0') || c.is_control())
{
return bad(&format!("contains {c:?}"));
}
Ok(())
}
pub fn payload_rel_path(
dispatchable: bool,
name: &str,
version: Option<&str>,
) -> Result<PathBuf, StoreError> {
safe_component("payload name", name)?;
match (dispatchable, version) {
(false, Some(version)) => {
safe_component("payload version", version)?;
Ok(PathBuf::from(PAYLOAD_DIR).join(name).join(version))
}
_ => Ok(PathBuf::from("bin").join(name)),
}
}
#[derive(Debug, Clone, Deserialize)]
struct ManifestEnvelope {
#[serde(default)]
annotations: BTreeMap<String, String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct InstalledLayer {
pub digest: String,
pub layer: LayerId,
pub channel: String,
pub root: PathBuf,
}
#[derive(Debug, Clone)]
pub struct Store {
root: PathBuf,
}
#[derive(Debug, thiserror::Error)]
pub enum StoreError {
#[error("io error at {path}")]
Io {
path: String,
#[source]
source: std::io::Error,
},
#[error("{path}: layer.json is not a valid layer manifest: {reason}")]
BadManifest { path: String, reason: String },
#[error("{what} {value:?} is not a usable path component ({why}) — refusing to lay it down")]
UnsafeComponent {
what: &'static str,
value: String,
why: String,
},
#[error(
"two payloads of this layer both claim {path} ('{first}' and '{second}') — refusing to \
lay one down over the other. A layer may hold several versions of one name, but each \
must be a distinct payload; two entries with one identity would mean the wrong bytes \
land under the right name."
)]
Collision {
path: String,
first: String,
second: String,
},
}
impl Store {
pub fn at(root: impl Into<PathBuf>) -> Self {
Store { root: root.into() }
}
pub fn root(&self) -> &Path {
&self.root
}
pub fn lay_down(
&self,
manifest_bytes: &[u8],
tools: &[(&str, &[u8])],
) -> Result<String, StoreError> {
let payloads: Vec<Payload<'_>> = tools
.iter()
.map(|(name, bytes)| Payload::tool(name, bytes))
.collect();
self.lay_down_payloads(manifest_bytes, &payloads)
}
pub fn lay_down_payloads(
&self,
manifest_bytes: &[u8],
payloads: &[Payload<'_>],
) -> Result<String, StoreError> {
let digest = manifest_digest(manifest_bytes);
let entry = self.core_dir().join(digest.replace(':', "-"));
let io = |path: &Path, source: std::io::Error| StoreError::Io {
path: path.display().to_string(),
source,
};
let mut placed: BTreeMap<PathBuf, String> = BTreeMap::new();
let mut plan: Vec<(PathBuf, &Payload<'_>)> = Vec::new();
for payload in payloads {
let rel = payload_rel_path(payload.dispatchable, payload.name, payload.version)?;
let who = match payload.version {
Some(v) => format!("{}@{v}", payload.name),
None => payload.name.to_string(),
};
if let Some(first) = placed.get(&rel) {
return Err(StoreError::Collision {
path: rel.display().to_string(),
first: first.clone(),
second: who,
});
}
placed.insert(rel.clone(), who);
plan.push((rel, payload));
}
let bin = entry.join("bin");
std::fs::create_dir_all(&bin).map_err(|e| io(&bin, e))?;
let manifest_path = entry.join("layer.json");
std::fs::write(&manifest_path, manifest_bytes).map_err(|e| io(&manifest_path, e))?;
for (rel, payload) in plan {
let path = entry.join(rel);
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).map_err(|e| io(parent, e))?;
}
std::fs::write(&path, payload.bytes).map_err(|e| io(&path, e))?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let mode = if payload.dispatchable { 0o755 } else { 0o644 };
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(mode))
.map_err(|e| io(&path, e))?;
}
}
Ok(digest)
}
pub fn list(&self) -> Result<Vec<InstalledLayer>, StoreError> {
let core = self.core_dir();
if !core.exists() {
return Ok(Vec::new());
}
let mut names: Vec<String> = std::fs::read_dir(&core)
.map_err(|e| StoreError::Io {
path: core.display().to_string(),
source: e,
})?
.filter_map(|e| e.ok())
.filter(|e| e.path().is_dir())
.map(|e| e.file_name().to_string_lossy().into_owned())
.filter(|n| n.starts_with("sha256-"))
.collect();
names.sort();
names
.into_iter()
.map(|name| self.read_entry(&name.replacen('-', ":", 1)))
.collect()
}
pub fn varve_root(&self) -> std::path::PathBuf {
let root = self.root();
if root
.parent()
.and_then(|p| p.file_name())
.is_some_and(|n| n == "realms")
{
root.parent()
.and_then(|p| p.parent())
.map(|p| p.to_path_buf())
.unwrap_or_else(|| root.to_path_buf())
} else {
root.to_path_buf()
}
}
pub fn partitions(&self) -> Vec<(Option<String>, Store)> {
let root = self.varve_root();
let mut out: Vec<(Option<String>, Store)> = vec![(None, Store::at(&root))];
if let Ok(rd) = std::fs::read_dir(root.join("realms")) {
let mut parts: Vec<std::path::PathBuf> =
rd.filter_map(|e| e.ok()).map(|e| e.path()).collect();
parts.sort();
for p in parts {
let fp = p.file_name().map(|n| n.to_string_lossy().into_owned());
out.push((fp, Store::at(p)));
}
}
out
}
pub fn find_anywhere(
&self,
digest: &str,
) -> Result<Option<(Store, InstalledLayer)>, StoreError> {
if let Some(entry) = self.get(digest)? {
return Ok(Some((self.clone(), entry)));
}
let root = self.varve_root();
let mut candidates = vec![Store::at(&root)];
if let Ok(rd) = std::fs::read_dir(root.join("realms")) {
let mut parts: Vec<std::path::PathBuf> =
rd.filter_map(|e| e.ok()).map(|e| e.path()).collect();
parts.sort();
candidates.extend(parts.into_iter().map(Store::at));
}
for candidate in candidates {
if candidate.root() == self.root() {
continue;
}
if let Some(entry) = candidate.get(digest)? {
return Ok(Some((candidate, entry)));
}
}
Ok(None)
}
pub fn get(&self, digest: &str) -> Result<Option<InstalledLayer>, StoreError> {
let entry = self.core_dir().join(digest.replace(':', "-"));
if !entry.join("layer.json").is_file() {
return Ok(None);
}
self.read_entry(digest).map(Some)
}
pub fn tool_path(&self, layer: &InstalledLayer, tool: &str) -> Option<PathBuf> {
let path = layer.root.join("bin").join(tool);
path.is_file().then_some(path)
}
pub fn entry_path(&self, layer: &InstalledLayer, entry: &ManifestEntry) -> Option<PathBuf> {
let name = entry.annotations.get("eu.pulseengine.tool")?;
let dispatchable = entry_is_dispatchable(entry);
let rel = payload_rel_path(dispatchable, name, entry_version(entry)).ok()?;
let path = layer.root.join(rel);
if path.is_file() {
return Some(path);
}
(!dispatchable)
.then(|| layer.root.join("bin").join(name))
.filter(|legacy| legacy.is_file())
}
fn core_dir(&self) -> PathBuf {
self.root.join("core")
}
fn read_entry(&self, digest: &str) -> Result<InstalledLayer, StoreError> {
let root = self.core_dir().join(digest.replace(':', "-"));
let manifest_path = root.join("layer.json");
let bad = |reason: String| StoreError::BadManifest {
path: manifest_path.display().to_string(),
reason,
};
let bytes = std::fs::read(&manifest_path).map_err(|source| StoreError::Io {
path: manifest_path.display().to_string(),
source,
})?;
let envelope: ManifestEnvelope =
serde_json::from_slice(&bytes).map_err(|e| bad(e.to_string()))?;
let layer_str = envelope
.annotations
.get("eu.pulseengine.varve.layer")
.ok_or_else(|| bad("missing eu.pulseengine.varve.layer annotation".into()))?;
let layer: LayerId = layer_str
.parse()
.map_err(|e: crate::layer::LayerIdError| bad(e.to_string()))?;
let channel = envelope
.annotations
.get("eu.pulseengine.varve.channel")
.cloned()
.unwrap_or_default();
Ok(InstalledLayer {
digest: digest.to_string(),
layer,
channel,
root,
})
}
}
impl Store {
pub fn manifest_tool_names(&self, layer: &InstalledLayer) -> Result<Vec<String>, StoreError> {
let payload = std::fs::read(layer.root.join("layer.json")).map_err(|e| StoreError::Io {
path: layer.root.join("layer.json").display().to_string(),
source: e,
})?;
let json: serde_json::Value = match serde_json::from_slice(&payload) {
Ok(j) => j,
Err(_) => return Ok(Vec::new()),
};
Ok(json["manifests"]
.as_array()
.map(|es| {
es.iter()
.filter_map(|e| e["annotations"]["eu.pulseengine.tool"].as_str())
.map(str::to_string)
.collect()
})
.unwrap_or_default())
}
}
pub fn manifest_digest(bytes: &[u8]) -> String {
format!("sha256:{}", hex::encode(Sha256::digest(bytes)))
}
#[cfg(test)]
pub(crate) mod fixtures {
pub fn manifest(layer: &str, channel: &str) -> Vec<u8> {
format!(
r#"{{
"schemaVersion": 2,
"mediaType": "application/vnd.oci.image.index.v1+json",
"artifactType": "application/vnd.pulseengine.varve.layer.v1+json",
"annotations": {{
"eu.pulseengine.varve.layer": "{layer}",
"eu.pulseengine.varve.channel": "{channel}"
}},
"manifests": []
}}"#
)
.into_bytes()
}
pub fn manifest_without_channel(layer: &str) -> Vec<u8> {
format!(
r#"{{
"schemaVersion": 2,
"mediaType": "application/vnd.oci.image.index.v1+json",
"artifactType": "application/vnd.pulseengine.varve.layer.v1+json",
"annotations": {{
"eu.pulseengine.varve.layer": "{layer}"
}},
"manifests": []
}}"#
)
.into_bytes()
}
}
#[cfg(test)]
mod partition_tests {
use super::*;
#[test]
fn a_layer_is_found_in_any_partition_under_the_same_root() {
let tmp = tempfile::tempdir().unwrap();
let root = tmp.path();
let mine = Store::at(root.join("realms").join("aaaa"));
let theirs = Store::at(root.join("realms").join("bbbb"));
let digest = theirs
.lay_down(
&fixtures::manifest("2026.08.0", "qualified"),
&[("btool", b"b")],
)
.unwrap();
assert!(mine.get(&digest).unwrap().is_none());
let (owner, entry) = mine.find_anywhere(&digest).unwrap().expect("found");
assert_eq!(entry.digest, digest);
assert_eq!(owner.root(), theirs.root(), "found in the owning partition");
assert!(owner.tool_path(&entry, "btool").is_some());
}
#[test]
fn the_varve_root_is_recovered_from_a_realm_partition() {
let tmp = tempfile::tempdir().unwrap();
let root = tmp.path();
assert_eq!(
Store::at(root.join("realms").join("ffff")).varve_root(),
root.to_path_buf()
);
assert_eq!(Store::at(root).varve_root(), root.to_path_buf());
}
}
#[cfg(test)]
mod tests {
use super::*;
fn store() -> (tempfile::TempDir, Store) {
let tmp = tempfile::tempdir().unwrap();
let store = Store::at(tmp.path().join("varve-root"));
(tmp, store)
}
#[test]
fn two_layers_coexist_and_are_independently_addressable() {
let (_tmp, store) = store();
let july = fixtures::manifest("2026.07.0", "qualified");
let august = fixtures::manifest("2026.08.0", "qualified");
let d_july = store.lay_down(&july, &[("synth", b"july-synth")]).unwrap();
let d_august = store
.lay_down(&august, &[("synth", b"august-synth")])
.unwrap();
assert_ne!(d_july, d_august);
let listed = store.list().unwrap();
assert_eq!(listed.len(), 2);
let july_entry = store.get(&d_july).unwrap().unwrap();
let august_entry = store.get(&d_august).unwrap().unwrap();
assert_eq!(july_entry.layer.to_string(), "2026.07.0");
assert_eq!(august_entry.layer.to_string(), "2026.08.0");
let july_synth = store.tool_path(&july_entry, "synth").unwrap();
let august_synth = store.tool_path(&august_entry, "synth").unwrap();
assert_eq!(std::fs::read(july_synth).unwrap(), b"july-synth");
assert_eq!(std::fs::read(august_synth).unwrap(), b"august-synth");
}
#[test]
fn store_key_is_the_manifest_digest() {
let (_tmp, store) = store();
let bytes = fixtures::manifest("2026.07.0", "qualified");
let digest = store.lay_down(&bytes, &[]).unwrap();
assert_eq!(digest, manifest_digest(&bytes));
let entry = store.get(&digest).unwrap().unwrap();
assert!(
entry
.root
.ends_with(format!("core/{}", digest.replace(':', "-"))),
"entry rooted at digest-keyed dir, got {}",
entry.root.display()
);
assert_eq!(std::fs::read(entry.root.join("layer.json")).unwrap(), bytes);
}
#[test]
fn missing_layer_is_none_not_an_invention() {
let (_tmp, store) = store();
assert_eq!(
store
.get("sha256:0000000000000000000000000000000000000000000000000000000000000000")
.unwrap(),
None
);
assert_eq!(store.list().unwrap(), vec![]);
}
fn entry(name: &str, version: Option<&str>, kind: Option<&str>) -> ManifestEntry {
let mut annotations = BTreeMap::new();
annotations.insert("eu.pulseengine.tool".to_string(), name.to_string());
if let Some(v) = version {
annotations.insert("eu.pulseengine.tool.version".to_string(), v.to_string());
}
if let Some(k) = kind {
annotations.insert(crate::kind::ANN_KIND.to_string(), k.to_string());
}
ManifestEntry {
digest: manifest_digest(name.as_bytes()),
annotations,
}
}
fn held<'a>(name: &'a str, version: &'a str, bytes: &'a [u8]) -> Payload<'a> {
Payload {
name,
version: Some(version),
dispatchable: false,
bytes,
}
}
#[test]
fn two_versions_of_one_name_are_two_files_neither_overwriting_the_other() {
let (_tmp, store) = store();
let manifest = fixtures::manifest("2026.08.0", "qualified");
let digest = store
.lay_down_payloads(
&manifest,
&[
held("serde", "1.0.200", b"serde-200-bytes"),
held("serde", "1.0.210", b"serde-210-bytes"),
],
)
.unwrap();
let layer = store.get(&digest).unwrap().unwrap();
let two_hundred = layer.root.join("payloads/serde/1.0.200");
let two_ten = layer.root.join("payloads/serde/1.0.210");
assert_eq!(std::fs::read(&two_hundred).unwrap(), b"serde-200-bytes");
assert_eq!(std::fs::read(&two_ten).unwrap(), b"serde-210-bytes");
assert!(!layer.root.join("bin/serde").exists());
assert_eq!(
store.entry_path(&layer, &entry("serde", Some("1.0.200"), Some("crate"))),
Some(two_hundred)
);
assert_eq!(
store.entry_path(&layer, &entry("serde", Some("1.0.210"), Some("crate"))),
Some(two_ten)
);
}
#[test]
fn two_payloads_claiming_one_path_are_refused_before_anything_is_written() {
let (_tmp, store) = store();
let manifest = fixtures::manifest("2026.08.0", "qualified");
let err = store
.lay_down_payloads(
&manifest,
&[
Payload::tool("synth", b"first-bytes"),
Payload::tool("synth", b"second-bytes"),
],
)
.unwrap_err();
assert!(
matches!(&err, StoreError::Collision { path, .. } if path.contains("synth")),
"got: {err}"
);
assert!(store.list().unwrap().is_empty(), "nothing may be laid down");
let err = store
.lay_down_payloads(
&manifest,
&[
Payload {
name: "wit-pkg",
version: None,
dispatchable: false,
bytes: b"a",
},
Payload {
name: "wit-pkg",
version: None,
dispatchable: false,
bytes: b"b",
},
],
)
.unwrap_err();
assert!(matches!(err, StoreError::Collision { .. }), "got: {err}");
}
#[test]
fn a_name_or_version_that_escapes_the_layer_is_refused() {
let (_tmp, store) = store();
let manifest = fixtures::manifest("2026.08.0", "qualified");
for (name, version) in [
("../../escape", Some("1.0.0")),
("serde", Some("../../escape")),
("a/b", Some("1.0.0")),
("..", Some("1.0.0")),
("", Some("1.0.0")),
("serde", Some("")),
] {
let err = store
.lay_down_payloads(&manifest, &[held(name, version.unwrap(), b"x")])
.unwrap_err();
assert!(
matches!(err, StoreError::UnsafeComponent { .. }),
"{name:?}@{version:?} must be refused, got: {err}"
);
}
assert!(store.list().unwrap().is_empty());
assert!(!store.root().join("escape").exists());
}
#[test]
fn a_layer_installed_before_this_change_still_resolves_its_crate() {
let (_tmp, store) = store();
let manifest = fixtures::manifest("2026.08.0", "qualified");
let digest = store
.lay_down(&manifest, &[("legacy-crate", b"old-layout-bytes")])
.unwrap();
let layer = store.get(&digest).unwrap().unwrap();
let path = store
.entry_path(&layer, &entry("legacy-crate", Some("0.1.0"), Some("crate")))
.expect("a pre-REQ-STORE-002 layer must keep resolving");
assert_eq!(std::fs::read(path).unwrap(), b"old-layout-bytes");
}
#[test]
fn a_tool_keeps_bin_and_a_held_payload_never_borrows_it() {
let (_tmp, store) = store();
let manifest = fixtures::manifest("2026.08.0", "qualified");
let digest = store
.lay_down_payloads(
&manifest,
&[
Payload::tool("synth", b"synth-binary"),
held("serde", "1.0.200", b"serde-crate"),
],
)
.unwrap();
let layer = store.get(&digest).unwrap().unwrap();
assert_eq!(
std::fs::read(store.tool_path(&layer, "synth").unwrap()).unwrap(),
b"synth-binary"
);
assert_eq!(store.tool_path(&layer, "serde"), None, "not dispatchable");
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let mode =
|p: std::path::PathBuf| std::fs::metadata(p).unwrap().permissions().mode() & 0o777;
assert_eq!(mode(layer.root.join("bin/synth")), 0o755);
assert_eq!(
mode(layer.root.join("payloads/serde/1.0.200")),
0o644,
"a .crate tarball is data, not something to execute"
);
}
}
#[test]
fn an_unrecognised_kind_is_held_by_version_never_dispatched_by_name() {
let e = entry("future", Some("2.0.0"), Some("quantum-blob"));
assert!(!entry_is_dispatchable(&e));
assert_eq!(
payload_rel_path(false, "future", Some("2.0.0")).unwrap(),
PathBuf::from("payloads/future/2.0.0")
);
assert!(entry_is_dispatchable(&entry("synth", Some("0.45.0"), None)));
assert_eq!(
payload_rel_path(true, "synth", Some("0.45.0")).unwrap(),
PathBuf::from("bin/synth")
);
}
#[test]
fn a_tree_shaped_payload_is_held_by_name_and_version_like_any_other() {
let (_tmp, store) = store();
let manifest = fixtures::manifest("2026.08.0", "qualified");
let digest = store
.lay_down_payloads(
&manifest,
&[
held("poky-cortexa53", "4.0.15", b"sdk-tarball-4.0.15"),
held("poky-cortexa53", "5.0.2", b"sdk-tarball-5.0.2"),
held("zephyr-hal", "3.6.0", b"zephyr-module-tarball"),
],
)
.unwrap();
let layer = store.get(&digest).unwrap().unwrap();
for (version, want) in [
("4.0.15", &b"sdk-tarball-4.0.15"[..]),
("5.0.2", &b"sdk-tarball-5.0.2"[..]),
] {
let path = store
.entry_path(&layer, &entry("poky-cortexa53", Some(version), Some("sdk")))
.expect("an sdk resolves from its manifest entry like any held payload");
assert_eq!(
path,
layer
.root
.join(format!("payloads/poky-cortexa53/{version}"))
);
assert_eq!(std::fs::read(&path).unwrap(), want);
}
assert!(!crate::kind::PayloadKind::Sdk.is_dispatchable());
assert_eq!(store.tool_path(&layer, "poky-cortexa53"), None);
assert!(!layer.root.join("bin/poky-cortexa53").exists());
assert!(
store
.entry_path(
&layer,
&entry("zephyr-hal", Some("3.6.0"), Some("zephyr-module"))
)
.is_some()
);
}
#[test]
fn missing_tool_in_an_installed_layer_is_none() {
let (_tmp, store) = store();
let bytes = fixtures::manifest("2026.07.0", "qualified");
let digest = store.lay_down(&bytes, &[("rivet", b"r")]).unwrap();
let entry = store.get(&digest).unwrap().unwrap();
assert!(store.tool_path(&entry, "rivet").is_some());
assert_eq!(store.tool_path(&entry, "synth"), None);
}
}