use std::collections::BTreeSet;
use std::sync::Arc;
use async_trait::async_trait;
use tracing::warn;
use super::nar_refs::NarRefIndex;
use super::nar_stream::{self, NarSource, NarStream};
use super::{NarResidency, StorageBackend};
use crate::StoreError;
struct TierNarSource {
tier: Arc<dyn StorageBackend>,
path: String,
}
impl TierNarSource {
fn new(tier: &Arc<dyn StorageBackend>, path: &str) -> Self {
Self { tier: Arc::clone(tier), path: path.to_string() }
}
}
#[async_trait]
impl NarSource for TierNarSource {
async fn open(&self) -> Result<NarStream, StoreError> {
self.tier.get_nar_stream(&self.path).await?.ok_or_else(|| {
StoreError::PathNotFound(format!(
"{}: vanished from the source tier mid-promotion",
self.path
))
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum WritePolicy {
#[default]
WriteThrough,
WriteBack,
WriteAround,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TieredTier {
MockParityProven,
LiveClusterProven,
}
pub const TIERED_BACKEND_TIER: TieredTier = TieredTier::MockParityProven;
pub struct TieredBackend {
l1: Arc<dyn StorageBackend>,
l2: Arc<dyn StorageBackend>,
l3: Arc<dyn StorageBackend>,
write_policy: WritePolicy,
}
impl TieredBackend {
#[must_use]
pub fn new(
l1: Arc<dyn StorageBackend>,
l2: Arc<dyn StorageBackend>,
l3: Arc<dyn StorageBackend>,
) -> Self {
Self::with_write_policy(l1, l2, l3, WritePolicy::default())
}
#[must_use]
pub fn with_write_policy(
l1: Arc<dyn StorageBackend>,
l2: Arc<dyn StorageBackend>,
l3: Arc<dyn StorageBackend>,
write_policy: WritePolicy,
) -> Self {
Self { l1, l2, l3, write_policy }
}
#[must_use]
pub fn write_policy(&self) -> WritePolicy {
self.write_policy
}
fn tiers(&self) -> [(&'static str, &Arc<dyn StorageBackend>); 3] {
[("l1", &self.l1), ("l2", &self.l2), ("l3", &self.l3)]
}
async fn warm_narinfo(tier: &Arc<dyn StorageBackend>, hash: &str, content: &str) {
if let Err(e) = tier.put_narinfo(hash, content).await {
warn!(hash = %hash, error = %e, "tiered: best-effort narinfo warm failed");
}
}
async fn warm_nar_from(tier: &Arc<dyn StorageBackend>, path: &str, src: &dyn NarSource) {
if let Err(e) = tier.put_nar_stream(path, src).await {
warn!(path = %path, error = %e, "tiered: best-effort NAR warm failed");
}
}
async fn warm_nar_from_tier(
tier: &Arc<dyn StorageBackend>,
from: &Arc<dyn StorageBackend>,
path: &str,
) {
Self::warm_nar_from(tier, path, &TierNarSource::new(from, path)).await;
}
fn note_tier_read_failure(tier: &'static str, key: &str, e: StoreError) -> StoreError {
tracing::error!(
tier = tier,
key = %key,
error = %e,
"tiered: READ FAILED on a tier — falling through to the next tier; \
this tier is degraded and needs attention",
);
e
}
fn durable_write_outcome(
kind: &'static str,
key: &str,
l2: Result<(), StoreError>,
l3: Result<(), StoreError>,
) -> Result<(), StoreError> {
match (l2, l3) {
(Ok(()), Ok(())) => Ok(()),
(Err(e), Ok(())) => {
warn!(
kind = kind, key = %key, tier = "l2", error = %e,
"tiered: durable write failed on ONE tier; the other durable tier \
accepted it, so the content is still serveable — redundancy lost",
);
Ok(())
}
(Ok(()), Err(e)) => {
warn!(
kind = kind, key = %key, tier = "l3", error = %e,
"tiered: durable write failed on ONE tier; the other durable tier \
accepted it, so the content is still serveable — redundancy lost",
);
Ok(())
}
(Err(e2), Err(e3)) => {
tracing::error!(
kind = kind, key = %key, l2_error = %e2, l3_error = %e3,
"tiered: durable write failed on EVERY durable tier — nothing was stored",
);
Err(e2)
}
}
}
}
#[async_trait]
impl StorageBackend for TieredBackend {
async fn get_narinfo(&self, hash: &str) -> Result<Option<String>, StoreError> {
let mut broken: Option<StoreError> = None;
match self.l1.get_narinfo(hash).await {
Ok(Some(v)) => return Ok(Some(v)),
Ok(None) => {}
Err(e) => broken = Some(Self::note_tier_read_failure("l1", hash, e)),
}
match self.l2.get_narinfo(hash).await {
Ok(Some(v)) => {
Self::warm_narinfo(&self.l1, hash, &v).await;
return Ok(Some(v));
}
Ok(None) => {}
Err(e) => broken = Some(Self::note_tier_read_failure("l2", hash, e)),
}
match self.l3.get_narinfo(hash).await {
Ok(Some(v)) => {
Self::warm_narinfo(&self.l2, hash, &v).await;
Self::warm_narinfo(&self.l1, hash, &v).await;
return Ok(Some(v));
}
Ok(None) => {}
Err(e) => broken = Some(Self::note_tier_read_failure("l3", hash, e)),
}
match broken {
Some(e) => Err(e),
None => Ok(None),
}
}
async fn get_nar(&self, path: &str) -> Result<Option<Vec<u8>>, StoreError> {
match self.get_nar_stream(path).await? {
Some(s) => Ok(Some(nar_stream::collect_nar(s, None).await?)),
None => Ok(None),
}
}
fn nar_residency(&self) -> NarResidency {
self.l1
.nar_residency()
.weaker(self.l2.nar_residency())
.weaker(self.l3.nar_residency())
}
async fn get_nar_stream(&self, path: &str) -> Result<Option<NarStream>, StoreError> {
let mut broken: Option<StoreError> = None;
match self.l1.get_nar_stream(path).await {
Ok(Some(s)) => return Ok(Some(s)),
Ok(None) => {}
Err(e) => broken = Some(Self::note_tier_read_failure("l1", path, e)),
}
match self.l2.get_nar_stream(path).await {
Ok(Some(s)) => {
Self::warm_nar_from_tier(&self.l1, &self.l2, path).await;
return Ok(Some(s));
}
Ok(None) => {}
Err(e) => broken = Some(Self::note_tier_read_failure("l2", path, e)),
}
match self.l3.get_nar_stream(path).await {
Ok(Some(s)) => {
Self::warm_nar_from_tier(&self.l2, &self.l3, path).await;
Self::warm_nar_from_tier(&self.l1, &self.l3, path).await;
return Ok(Some(s));
}
Ok(None) => {}
Err(e) => broken = Some(Self::note_tier_read_failure("l3", path, e)),
}
match broken {
Some(e) => Err(e),
None => Ok(None),
}
}
async fn put_narinfo_record(&self, hash: &str, content: &str) -> Result<(), StoreError> {
if self.write_policy == WritePolicy::WriteBack {
Self::warm_narinfo(&self.l1, hash, content).await;
}
let l2 = self.l2.put_narinfo_record(hash, content).await;
let l3 = self.l3.put_narinfo_record(hash, content).await;
Self::durable_write_outcome("narinfo", hash, l2, l3)?;
if self.write_policy == WritePolicy::WriteThrough {
Self::warm_narinfo(&self.l1, hash, content).await;
}
Ok(())
}
async fn delete_narinfo_record(&self, hash: &str) -> Result<(), StoreError> {
let mut failed: Option<StoreError> = None;
for (name, tier) in self.tiers() {
if let Err(e) = tier.delete_narinfo_record(hash).await {
tracing::error!(
hash = %hash, tier = name, error = %e,
"tiered: narinfo delete FAILED on a tier — reads fall through, so this \
narinfo is still servable; refusing to report the delete as complete",
);
failed = Some(e);
}
}
failed.map_or(Ok(()), Err)
}
async fn delete_nar_record(&self, nar_path: &str) -> Result<(), StoreError> {
let mut failed: Option<StoreError> = None;
for (name, tier) in self.tiers() {
if let Err(e) = tier.delete_nar_record(nar_path).await {
warn!(
path = %nar_path, tier = name, error = %e,
"tiered: NAR delete failed on a tier — a copy survives there",
);
failed = Some(e);
}
}
failed.map_or(Ok(()), Err)
}
fn nar_ref_index(&self) -> &dyn NarRefIndex {
self
}
async fn put_nar(&self, path: &str, data: &[u8]) -> Result<(), StoreError> {
self.put_nar_stream(path, &nar_stream::BytesNarSource::from(data)).await
}
async fn put_nar_stream(&self, path: &str, src: &dyn NarSource) -> Result<(), StoreError> {
if self.write_policy == WritePolicy::WriteBack {
Self::warm_nar_from(&self.l1, path, src).await;
}
let l2 = self.l2.put_nar_stream(path, src).await;
let l3 = self.l3.put_nar_stream(path, src).await;
Self::durable_write_outcome("nar", path, l2, l3)?;
if self.write_policy == WritePolicy::WriteThrough {
Self::warm_nar_from(&self.l1, path, src).await;
}
Ok(())
}
async fn list_narinfos(&self) -> Result<Vec<String>, StoreError> {
let mut set = BTreeSet::new();
set.extend(self.l2.list_narinfos().await?);
set.extend(self.l3.list_narinfos().await?);
Ok(set.into_iter().collect())
}
async fn wipe_all(&self) -> Result<usize, StoreError> {
let mut cleared = 0usize;
for (name, tier) in self.tiers() {
match tier.wipe_all().await {
Ok(n) => cleared = cleared.max(n),
Err(e) => warn!(tier = name, error = %e, "tiered: best-effort wipe failed"),
}
}
Ok(cleared)
}
}
#[async_trait]
impl NarRefIndex for TieredBackend {
async fn record(&self, nar_path: &str, hash: &str) -> Result<(), StoreError> {
let l2 = self.l2.nar_ref_index().record(nar_path, hash).await;
let l3 = self.l3.nar_ref_index().record(nar_path, hash).await;
Self::durable_write_outcome("nar-ref", nar_path, l2, l3)?;
if let Err(e) = self.l1.nar_ref_index().record(nar_path, hash).await {
warn!(path = %nar_path, error = %e, "tiered: best-effort nar-ref warm failed");
}
Ok(())
}
async fn forget(&self, nar_path: &str, hash: &str) -> Result<(), StoreError> {
for (name, tier) in self.tiers() {
if let Err(e) = tier.nar_ref_index().forget(nar_path, hash).await {
warn!(
path = %nar_path, tier = name, error = %e,
"tiered: best-effort nar-ref forget failed — the edge survives, so the \
NAR is retained rather than stranded",
);
}
}
Ok(())
}
async fn referrers(&self, nar_path: &str) -> Result<Vec<String>, StoreError> {
let mut set = BTreeSet::new();
let mut broken: Option<StoreError> = None;
for (name, tier) in self.tiers() {
match tier.nar_ref_index().referrers(nar_path).await {
Ok(hashes) => set.extend(hashes),
Err(e) => {
broken = Some(Self::note_tier_read_failure(name, nar_path, e));
}
}
}
match broken {
Some(e) => Err(e),
None => Ok(set.into_iter().collect()),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::storage::nar_refs::MemNarRefIndex;
use crate::storage::LocalStorage;
use std::collections::HashMap;
use std::sync::Mutex;
#[derive(Default)]
struct MemBackend {
narinfo: Mutex<HashMap<String, String>>,
nar: Mutex<HashMap<String, Vec<u8>>>,
refs: MemNarRefIndex,
writes_fail: Mutex<bool>,
reads_fail: Mutex<bool>,
deletes_fail: Mutex<bool>,
}
impl MemBackend {
fn has_narinfo(&self, hash: &str) -> bool {
self.narinfo.lock().unwrap().contains_key(hash)
}
fn has_nar(&self, path: &str) -> bool {
self.nar.lock().unwrap().contains_key(path)
}
fn clear(&self) {
self.narinfo.lock().unwrap().clear();
self.nar.lock().unwrap().clear();
}
fn set_writes_fail(&self, v: bool) {
*self.writes_fail.lock().unwrap() = v;
}
fn set_reads_fail(&self, v: bool) {
*self.reads_fail.lock().unwrap() = v;
}
fn fail_if_configured(&self) -> Result<(), StoreError> {
if *self.writes_fail.lock().unwrap() {
Err(StoreError::NotImplemented("mock writes disabled"))
} else {
Ok(())
}
}
fn fail_reads_if_configured(&self) -> Result<(), StoreError> {
if *self.reads_fail.lock().unwrap() {
Err(StoreError::SchemaMissing("mock: relation does not exist".to_string()))
} else {
Ok(())
}
}
fn set_deletes_fail(&self, v: bool) {
*self.deletes_fail.lock().unwrap() = v;
}
fn fail_deletes_if_configured(&self) -> Result<(), StoreError> {
if *self.deletes_fail.lock().unwrap() {
Err(StoreError::NotImplemented("mock deletes disabled"))
} else {
Ok(())
}
}
}
#[async_trait]
impl StorageBackend for MemBackend {
async fn get_narinfo(&self, hash: &str) -> Result<Option<String>, StoreError> {
self.fail_reads_if_configured()?;
Ok(self.narinfo.lock().unwrap().get(hash).cloned())
}
async fn put_narinfo_record(&self, hash: &str, content: &str) -> Result<(), StoreError> {
self.fail_if_configured()?;
self.narinfo.lock().unwrap().insert(hash.to_string(), content.to_string());
Ok(())
}
async fn delete_narinfo_record(&self, hash: &str) -> Result<(), StoreError> {
self.fail_deletes_if_configured()?;
self.narinfo.lock().unwrap().remove(hash);
Ok(())
}
async fn delete_nar_record(&self, nar_path: &str) -> Result<(), StoreError> {
self.fail_deletes_if_configured()?;
self.nar.lock().unwrap().remove(nar_path);
Ok(())
}
fn nar_ref_index(&self) -> &dyn NarRefIndex {
&self.refs
}
async fn get_nar(&self, path: &str) -> Result<Option<Vec<u8>>, StoreError> {
self.fail_reads_if_configured()?;
Ok(self.nar.lock().unwrap().get(path).cloned())
}
async fn put_nar(&self, path: &str, data: &[u8]) -> Result<(), StoreError> {
self.fail_if_configured()?;
self.nar.lock().unwrap().insert(path.to_string(), data.to_vec());
Ok(())
}
fn nar_residency(&self) -> NarResidency {
NarResidency::WholeValue
}
async fn list_narinfos(&self) -> Result<Vec<String>, StoreError> {
Ok(self.narinfo.lock().unwrap().keys().cloned().collect())
}
}
const NARINFO: &str = "StorePath: /nix/store/abc-hello\nURL: nar/abc.nar.xz\nCompression: xz\nNarHash: sha256:bbb\nNarSize: 200\nReferences: \n";
const ADVERTISED_NAR: &str = "nar/abc.nar.xz";
fn mocks() -> (Arc<MemBackend>, Arc<MemBackend>, Arc<MemBackend>, TieredBackend) {
let l1 = Arc::new(MemBackend::default());
let l2 = Arc::new(MemBackend::default());
let l3 = Arc::new(MemBackend::default());
let tiered = TieredBackend::new(l1.clone(), l2.clone(), l3.clone());
(l1, l2, l3, tiered)
}
#[tokio::test]
async fn l1_hit_returns_without_touching_lower_tiers() {
let (l1, l2, l3, tiered) = mocks();
l1.put_narinfo("h", "hot").await.unwrap();
assert_eq!(tiered.get_narinfo("h").await.unwrap().unwrap(), "hot");
assert!(!l2.has_narinfo("h"));
assert!(!l3.has_narinfo("h"));
}
#[tokio::test]
async fn l2_hit_promotes_into_l1() {
let (l1, l2, _l3, tiered) = mocks();
l2.put_narinfo("h", NARINFO).await.unwrap();
assert!(!l1.has_narinfo("h"));
let got = tiered.get_narinfo("h").await.unwrap().unwrap();
assert_eq!(got, NARINFO);
assert!(l1.has_narinfo("h"), "L2 hit must promote into L1");
}
#[tokio::test]
async fn l3_hit_promotes_into_l2_and_l1() {
let (l1, l2, l3, tiered) = mocks();
l3.put_narinfo("h", NARINFO).await.unwrap();
let got = tiered.get_narinfo("h").await.unwrap().unwrap();
assert_eq!(got, NARINFO);
assert!(l2.has_narinfo("h"), "L3 hit must promote into L2");
assert!(l1.has_narinfo("h"), "L3 hit must promote into L1");
}
#[tokio::test]
async fn nar_l3_hit_promotes_into_l2_and_l1() {
let (l1, l2, l3, tiered) = mocks();
l3.put_nar("nar/x.nar.xz", b"blob").await.unwrap();
let got = tiered.get_nar("nar/x.nar.xz").await.unwrap().unwrap();
assert_eq!(got, b"blob");
assert!(l2.has_nar("nar/x.nar.xz"));
assert!(l1.has_nar("nar/x.nar.xz"));
}
#[tokio::test]
async fn miss_at_all_tiers_is_none() {
let (_l1, _l2, _l3, tiered) = mocks();
assert!(tiered.get_narinfo("ghost").await.unwrap().is_none());
assert!(tiered.get_nar("nar/ghost.nar.xz").await.unwrap().is_none());
}
#[tokio::test]
async fn l2_read_failure_falls_through_to_l3() {
let (l1, l2, l3, tiered) = mocks();
l3.put_narinfo("h", NARINFO).await.unwrap();
l3.put_nar("nar/h.nar.xz", b"blob").await.unwrap();
l2.set_reads_fail(true);
assert_eq!(
tiered.get_narinfo("h").await.unwrap().unwrap(),
NARINFO,
"a broken L2 must not hide a healthy L3",
);
assert_eq!(tiered.get_nar("nar/h.nar.xz").await.unwrap().unwrap(), b"blob");
assert!(l1.has_narinfo("h"), "the L3 hit still warms the working hot tier");
}
#[tokio::test]
async fn broken_l1_and_l2_still_serve_from_l3() {
let (l1, l2, l3, tiered) = mocks();
l3.put_narinfo("h", NARINFO).await.unwrap();
l1.set_reads_fail(true);
l2.set_reads_fail(true);
assert_eq!(tiered.get_narinfo("h").await.unwrap().unwrap(), NARINFO);
}
#[tokio::test]
async fn every_tier_broken_and_no_hit_surfaces_an_error_not_a_false_absence() {
let (l1, l2, l3, tiered) = mocks();
for t in [&l1, &l2, &l3] {
t.set_reads_fail(true);
}
assert!(matches!(
tiered.get_narinfo("h").await.unwrap_err(),
StoreError::SchemaMissing(_),
));
assert!(matches!(
tiered.get_nar("nar/h.nar.xz").await.unwrap_err(),
StoreError::SchemaMissing(_),
));
}
#[tokio::test]
async fn all_tiers_healthy_and_empty_is_a_clean_miss_not_an_error() {
let (_l1, _l2, _l3, tiered) = mocks();
assert!(tiered.get_narinfo("ghost").await.unwrap().is_none());
}
#[tokio::test]
async fn promotion_failure_does_not_break_a_read() {
let (l1, l2, _l3, tiered) = mocks();
l2.put_narinfo("h", NARINFO).await.unwrap();
l1.set_writes_fail(true); let got = tiered.get_narinfo("h").await.unwrap();
assert_eq!(got.unwrap(), NARINFO);
assert!(!l1.has_narinfo("h"), "warm failed, so L1 stays empty — but the read still succeeded");
}
#[tokio::test]
async fn write_through_populates_all_tiers() {
let (l1, l2, l3, tiered) = mocks();
tiered.put_narinfo("h", NARINFO).await.unwrap();
assert!(l1.has_narinfo("h"), "write-through warms L1");
assert!(l2.has_narinfo("h"), "write-through persists L2");
assert!(l3.has_narinfo("h"), "write-through persists L3");
}
#[tokio::test]
async fn write_around_skips_l1_but_persists_durable() {
let l1 = Arc::new(MemBackend::default());
let l2 = Arc::new(MemBackend::default());
let l3 = Arc::new(MemBackend::default());
let tiered = TieredBackend::with_write_policy(
l1.clone(), l2.clone(), l3.clone(), WritePolicy::WriteAround,
);
tiered.put_narinfo("h", NARINFO).await.unwrap();
assert!(!l1.has_narinfo("h"), "write-around must NOT touch L1");
assert!(l2.has_narinfo("h"));
assert!(l3.has_narinfo("h"));
let _ = tiered.get_narinfo("h").await.unwrap();
assert!(l1.has_narinfo("h"), "read-through fills L1 after a write-around");
}
#[tokio::test]
async fn write_back_populates_all_tiers_and_is_durable() {
let l1 = Arc::new(MemBackend::default());
let l2 = Arc::new(MemBackend::default());
let l3 = Arc::new(MemBackend::default());
let tiered = TieredBackend::with_write_policy(
l1.clone(), l2.clone(), l3.clone(), WritePolicy::WriteBack,
);
assert_eq!(tiered.write_policy(), WritePolicy::WriteBack);
tiered.put_nar("nar/x.nar.xz", b"blob").await.unwrap();
assert!(l1.has_nar("nar/x.nar.xz"));
assert!(l2.has_nar("nar/x.nar.xz"));
assert!(l3.has_nar("nar/x.nar.xz"));
}
#[tokio::test]
async fn one_broken_durable_tier_still_lands_the_write_on_the_other() {
let l1 = Arc::new(MemBackend::default());
let l2 = Arc::new(MemBackend::default());
let l3 = Arc::new(MemBackend::default());
l2.set_writes_fail(true);
let tiered = TieredBackend::new(l1.clone(), l2.clone(), l3.clone());
tiered.put_narinfo("h", NARINFO).await.expect("one healthy durable tier must accept");
tiered.put_nar("nar/h.nar.xz", b"blob").await.expect("one healthy durable tier must accept");
assert!(!l2.has_narinfo("h"), "the broken tier holds nothing");
assert!(l3.has_narinfo("h"), "the healthy durable tier MUST have taken the write");
assert!(l3.has_nar("nar/h.nar.xz"));
assert_eq!(tiered.get_narinfo("h").await.unwrap().unwrap(), NARINFO);
}
#[tokio::test]
async fn write_fails_only_when_every_durable_tier_rejects() {
let l1 = Arc::new(MemBackend::default());
let l2 = Arc::new(MemBackend::default());
let l3 = Arc::new(MemBackend::default());
l2.set_writes_fail(true);
l3.set_writes_fail(true);
let tiered = TieredBackend::new(l1, l2, l3);
let err = tiered.put_narinfo("h", NARINFO).await.unwrap_err();
assert!(matches!(err, StoreError::NotImplemented(_)));
}
#[tokio::test]
async fn pod_roll_losing_l1_loses_nothing() {
let (l1, _l2, _l3, tiered) = mocks();
tiered.put_narinfo("h", NARINFO).await.unwrap();
tiered.put_nar("nar/h.nar.xz", b"blob").await.unwrap();
l1.clear();
assert_eq!(tiered.get_narinfo("h").await.unwrap().unwrap(), NARINFO);
assert_eq!(tiered.get_nar("nar/h.nar.xz").await.unwrap().unwrap(), b"blob");
}
#[tokio::test]
async fn delete_fans_out_to_all_tiers() {
let (l1, l2, l3, tiered) = mocks();
tiered.put_narinfo("h", NARINFO).await.unwrap();
tiered.put_nar(ADVERTISED_NAR, b"blob").await.unwrap();
tiered.delete("h").await.unwrap();
for t in [&l1, &l2, &l3] {
assert!(!t.has_narinfo("h"));
assert!(!t.has_nar(ADVERTISED_NAR));
}
}
#[tokio::test]
async fn delete_resolves_across_tiers_instead_of_guessing() {
let (l1, l2, l3, tiered) = mocks();
tiered.put_narinfo("h", NARINFO).await.unwrap();
tiered.put_nar(ADVERTISED_NAR, b"blob").await.unwrap();
tiered.put_nar("nar/h.nar.zst", b"someone else's nar").await.unwrap();
tiered.delete("h").await.unwrap();
for t in [&l1, &l2, &l3] {
assert!(!t.has_nar(ADVERTISED_NAR), "the advertised NAR must go");
}
assert_eq!(
tiered.get_nar("nar/h.nar.zst").await.unwrap().unwrap(),
b"someone else's nar",
);
}
#[tokio::test]
async fn a_co_referenced_nar_survives_the_first_delete_on_every_tier() {
let (l1, l2, l3, tiered) = mocks();
tiered.put_narinfo("pathA", NARINFO).await.unwrap();
tiered.put_narinfo("pathB", NARINFO).await.unwrap();
tiered.put_nar(ADVERTISED_NAR, b"shared").await.unwrap();
tiered.delete("pathA").await.unwrap();
for t in [&l1, &l2, &l3] {
assert!(t.has_nar(ADVERTISED_NAR), "pathB still advertises it");
}
tiered.delete("pathB").await.unwrap();
for t in [&l1, &l2, &l3] {
assert!(!t.has_nar(ADVERTISED_NAR), "the last referrer is gone");
}
}
#[tokio::test]
async fn a_failed_narinfo_delete_must_not_take_the_nar_with_it() {
let (l1, l2, l3, tiered) = mocks();
tiered.put_narinfo("h", NARINFO).await.unwrap();
tiered.put_nar(ADVERTISED_NAR, b"blob").await.unwrap();
l2.set_deletes_fail(true);
let err = tiered.delete("h").await.expect_err("a partial delete must surface");
assert!(
matches!(err, StoreError::NotImplemented(_)),
"expected the tier's own error, got {err:?}",
);
assert!(l2.has_narinfo("h"), "L2 kept the narinfo — that is the premise");
assert_eq!(
tiered.get_narinfo("h").await.unwrap().unwrap(),
NARINFO,
"and a read still serves it, because reads fall through",
);
for t in [&l1, &l2, &l3] {
assert!(
t.has_nar(ADVERTISED_NAR),
"the NAR must be untouched: its narinfo is still servable",
);
}
}
#[tokio::test]
async fn wipe_all_clears_every_tier() {
let (l1, l2, l3, tiered) = mocks();
tiered.put_narinfo("h", NARINFO).await.unwrap();
tiered.put_nar(ADVERTISED_NAR, b"blob").await.unwrap();
l2.put_narinfo("only2", "x").await.unwrap();
let removed = tiered.wipe_all().await.unwrap();
assert!(removed >= 1, "wipe reported nothing cleared");
for t in [&l1, &l2, &l3] {
assert!(t.list_narinfos().await.unwrap().is_empty(), "a tier survived the wipe");
assert!(!t.has_narinfo("h"));
assert!(!t.has_nar(ADVERTISED_NAR));
}
assert!(tiered.list_narinfos().await.unwrap().is_empty(), "cache not cold after wipe");
assert!(tiered.get_narinfo("h").await.unwrap().is_none());
}
#[tokio::test]
async fn list_narinfos_unions_durable_tiers_deduped() {
let (l1, l2, l3, tiered) = mocks();
l2.put_narinfo("shared", "x").await.unwrap();
l3.put_narinfo("shared", "x").await.unwrap();
l2.put_narinfo("only2", "y").await.unwrap();
l3.put_narinfo("only3", "z").await.unwrap();
l1.put_narinfo("hot-only", "w").await.unwrap();
let listed = tiered.list_narinfos().await.unwrap();
assert_eq!(listed, vec!["only2".to_string(), "only3".to_string(), "shared".to_string()]);
}
#[tokio::test]
async fn read_through_from_a_real_local_storage_l3() {
let dir = tempfile::tempdir().unwrap();
let l1 = Arc::new(MemBackend::default());
let l2 = Arc::new(MemBackend::default());
let l3_disk = Arc::new(LocalStorage::new(dir.path()));
l3_disk.put_narinfo("h", NARINFO).await.unwrap();
l3_disk.put_nar("nar/h.nar.xz", b"disk-blob").await.unwrap();
let tiered = TieredBackend::new(l1.clone(), l2.clone(), l3_disk);
assert_eq!(tiered.get_narinfo("h").await.unwrap().unwrap(), NARINFO);
assert_eq!(tiered.get_nar("nar/h.nar.xz").await.unwrap().unwrap(), b"disk-blob");
assert!(l1.has_narinfo("h"));
assert!(l2.has_narinfo("h"));
}
struct RecordingTier {
name: &'static str,
log: Arc<Mutex<Vec<&'static str>>>,
refuse: bool,
refs: MemNarRefIndex,
}
#[async_trait]
impl StorageBackend for RecordingTier {
async fn get_narinfo(&self, _h: &str) -> Result<Option<String>, StoreError> {
Ok(None)
}
async fn put_narinfo_record(&self, _h: &str, _c: &str) -> Result<(), StoreError> {
Ok(())
}
async fn delete_narinfo_record(&self, _h: &str) -> Result<(), StoreError> {
Ok(())
}
async fn delete_nar_record(&self, _p: &str) -> Result<(), StoreError> {
Ok(())
}
fn nar_ref_index(&self) -> &dyn NarRefIndex {
&self.refs
}
async fn get_nar(&self, _p: &str) -> Result<Option<Vec<u8>>, StoreError> {
Ok(None)
}
async fn put_nar(&self, _p: &str, _d: &[u8]) -> Result<(), StoreError> {
self.log.lock().unwrap().push(self.name);
if self.refuse {
return Err(StoreError::TooLarge { limit: 1, at_least: 2 });
}
Ok(())
}
fn nar_residency(&self) -> NarResidency {
NarResidency::WholeValue
}
async fn list_narinfos(&self) -> Result<Vec<String>, StoreError> {
Ok(vec![])
}
}
fn recording_tiers(
refuse_l1: bool,
) -> (Arc<Mutex<Vec<&'static str>>>, TieredBackend) {
let log = Arc::new(Mutex::new(Vec::new()));
let mk = |name, refuse| {
Arc::new(RecordingTier {
name,
log: Arc::clone(&log),
refuse,
refs: MemNarRefIndex::new(),
}) as Arc<dyn StorageBackend>
};
let tiered = TieredBackend::new(mk("l1", refuse_l1), mk("l2", false), mk("l3", false));
(log, tiered)
}
#[tokio::test]
async fn streamed_put_writes_l2_then_l3_then_warms_l1() {
let (log, tiered) = recording_tiers(false);
tiered.put_nar("nar/x.nar.xz", b"blob").await.unwrap();
assert_eq!(*log.lock().unwrap(), vec!["l2", "l3", "l1"]);
}
#[tokio::test]
async fn a_refused_l1_warm_never_fails_the_write() {
let (log, tiered) = recording_tiers(true);
tiered
.put_nar("nar/x.nar.xz", b"blob")
.await
.expect("a refused L1 warm must not fail the write");
assert_eq!(*log.lock().unwrap(), vec!["l2", "l3", "l1"], "L1 is still attempted, last");
}
#[tokio::test]
async fn write_back_warms_l1_before_the_durable_gate() {
let log = Arc::new(Mutex::new(Vec::new()));
let mk = |name| {
Arc::new(RecordingTier {
name,
log: Arc::clone(&log),
refuse: false,
refs: MemNarRefIndex::new(),
}) as Arc<dyn StorageBackend>
};
let tiered = TieredBackend::with_write_policy(
mk("l1"), mk("l2"), mk("l3"), WritePolicy::WriteBack,
);
tiered.put_nar("nar/x.nar.xz", b"blob").await.unwrap();
assert_eq!(*log.lock().unwrap(), vec!["l1", "l2", "l3"]);
}
#[tokio::test]
async fn a_multi_chunk_nar_round_trips_through_a_real_disk_tier() {
let dir = tempfile::tempdir().unwrap();
let l1 = Arc::new(MemBackend::default());
let l2: Arc<dyn StorageBackend> = Arc::new(LocalStorage::new(dir.path().join("l2")));
let l3: Arc<dyn StorageBackend> = Arc::new(LocalStorage::new(dir.path().join("l3")));
let tiered = TieredBackend::new(l1.clone(), l2, l3);
let nar: Vec<u8> = (0..crate::storage::NAR_CHUNK_BYTES + 777).map(|i| (i % 251) as u8).collect();
tiered.put_nar("nar/big.nar.xz", &nar).await.unwrap();
assert_eq!(tiered.get_nar("nar/big.nar.xz").await.unwrap().unwrap(), nar);
}
#[tokio::test]
async fn residency_reports_the_weakest_tier_not_the_resolver() {
let dir = tempfile::tempdir().unwrap();
let disk1: Arc<dyn StorageBackend> = Arc::new(LocalStorage::new(dir.path().join("a")));
let disk2: Arc<dyn StorageBackend> = Arc::new(LocalStorage::new(dir.path().join("b")));
let disk3: Arc<dyn StorageBackend> = Arc::new(LocalStorage::new(dir.path().join("c")));
let all_disk = TieredBackend::new(disk1.clone(), disk2.clone(), disk3);
assert_eq!(all_disk.nar_residency(), NarResidency::Streaming);
let with_double =
TieredBackend::new(Arc::new(MemBackend::default()), disk1, disk2);
assert_eq!(with_double.nar_residency(), NarResidency::WholeValue);
}
#[test]
fn honest_gate_tier_is_mock_parity_proven_not_live_cluster() {
assert_eq!(TIERED_BACKEND_TIER, TieredTier::MockParityProven);
}
}