use std::collections::BTreeSet;
use std::panic::{catch_unwind, resume_unwind, AssertUnwindSafe};
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use anyhow::{anyhow, Context, Result};
use ed25519_dalek::{SigningKey, VerifyingKey};
use triblespace::core::blob::encodings::simplearchive::SimpleArchive;
use triblespace::core::blob::encodings::succinctarchive::{
OrderedUniverse, Rank9AcceleratedSuccinctArchiveBlob, SuccinctArchiveBlob, UnionArchive,
};
use triblespace::core::blob::encodings::UnknownBlob;
use triblespace::core::collection::{
ensure_downstream as core_ensure_downstream, maintain_downstream as core_maintain_downstream,
realize_attached_as, succinctarchive_union, Attached, AttachedRead, AttachedSnapshot,
Collection, CollectionAttachment, CollectionCommit, CollectionData, CollectionDerive,
CollectionEncoding, CollectionHandle, CollectionMerge, CollectionRead,
CollectionRealizationError, CollectionRecord, CollectionRecordSelector, CollectionSnapshotExt,
CollectionStoreExt, CoreRealizer, Derived, RealizeDerived, Realized, SourceLocator, Support,
Upkeep, UpkeepReport,
};
use triblespace::core::id::Id;
use triblespace::core::inline::encodings::hash::Handle;
use triblespace::core::inline::InlineEncoding;
use triblespace::core::metadata::MetaDescribe;
use triblespace::core::repo::async_store::AsyncBlobStoreAcquire;
use triblespace::core::repo::pile::{Pile, ReadError};
use triblespace::core::repo::{
BlobStoreGet, BlobStoreList, CapabilityProofRead, MissingBlob, SnapshotSource, StorageClose,
Store, StoreRead,
};
use triblespace::core::signing_key_file;
use triblespace::core::trible::{Fragment, TribleSet};
use triblespace_search::portable_bm25::PortableBM25Blob;
pub type FactArchive = UnionArchive<OrderedUniverse>;
pub type FacultyStore = triblespace_net::peer::Leech<Pile>;
pub type FacultySnapshot = <FacultyStore as SnapshotSource>::Snapshot;
pub use triblespace::core::repo::async_store::AcquiringReader;
#[derive(Clone)]
pub struct Storage {
pile: PathBuf,
key: Option<PathBuf>,
shared: Option<Arc<Mutex<Option<Session>>>>,
}
struct Session {
store: Option<FacultyStore>,
runtime: Arc<tokio::runtime::Runtime>,
host: VerifyingKey,
}
impl Session {
fn close(mut self) -> Result<()> {
self.store
.take()
.expect("an open session owns its store")
.close()
.context("close shared faculty store")
}
}
impl Drop for Session {
fn drop(&mut self) {
if let Some(store) = self.store.take() {
if let Err(error) = store.close() {
eprintln!("closing shared faculty store failed: {error}");
}
}
}
}
impl std::fmt::Debug for Storage {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Storage")
.field("pile", &self.pile)
.field("key", &self.key)
.field("shared", &self.shared.is_some())
.finish()
}
}
impl Storage {
pub fn retained(&self) -> Self {
if self.shared.is_some() {
self.clone()
} else {
Self::shared(self.pile.clone(), self.key.clone())
}
}
pub fn new(pile: PathBuf, key: Option<PathBuf>) -> Self {
Self {
pile,
key,
shared: None,
}
}
pub fn shared(pile: PathBuf, key: Option<PathBuf>) -> Self {
Self {
pile,
key,
shared: Some(Arc::new(Mutex::new(None))),
}
}
pub fn path(&self) -> &Path {
&self.pile
}
pub fn key_path(&self) -> Option<&Path> {
self.key.as_deref()
}
pub fn scope<T>(&self, operation: impl FnOnce(&Self) -> Result<T>) -> Result<T> {
if self.shared.is_some() {
return operation(self);
}
let storage = Self::shared(self.pile.clone(), self.key.clone());
let result = operation(&storage);
storage.finish(result)
}
fn with_session<T>(
&self,
host: VerifyingKey,
operation: impl FnOnce(&mut Session) -> Result<T>,
) -> Result<T> {
let mut session = self
.shared
.as_ref()
.expect("shared storage owner")
.lock()
.map_err(|_| anyhow!("shared faculty store is poisoned"))?;
if session.is_none() {
let runtime = Arc::new(runtime()?);
let store = open_store_as(&self.pile, host)?;
*session = Some(Session {
store: Some(store),
runtime,
host,
});
}
let opened_as = session.as_ref().expect("initialized above").host;
if opened_as != host {
anyhow::bail!(
"the durable signing key changed while the shared faculty store was open \
(opened as {}, now {}); restart the process to open it as the new key",
hex::encode(opened_as.as_bytes()),
hex::encode(host.as_bytes())
);
}
let result = catch_unwind(AssertUnwindSafe(|| {
operation(session.as_mut().expect("initialized above"))
}));
drop(session);
match result {
Ok(result) => result,
Err(panic) => resume_unwind(panic),
}
}
pub fn with_store<T>(
&self,
operation: impl FnOnce(
&mut FacultyStore,
&SigningKey,
&Arc<tokio::runtime::Runtime>,
) -> Result<T>,
) -> Result<T> {
let signer = load_signer(&self.pile, self.key.as_deref())?;
if self.shared.is_some() {
return self.with_session(signer.verifying_key(), |session| {
operation(
session.store.as_mut().expect("open store"),
&signer,
&session.runtime,
)
});
}
let runtime = Arc::new(runtime()?);
let mut store = open_store_as(&self.pile, signer.verifying_key())?;
let result = operation(&mut store, &signer, &runtime);
finish_close(result, store.close().context("close faculty store"))
}
pub fn with_pile<T>(
&self,
operation: impl FnOnce(&mut Pile, &SigningKey) -> Result<T>,
) -> Result<T> {
let signer = load_signer(&self.pile, self.key.as_deref())?;
if self.shared.is_some() {
return self.with_session(signer.verifying_key(), |session| {
let mut pile = session.store.as_ref().expect("open store").store();
pile.refresh()
.map_err(|error| pile_read_error(&self.pile, error))?;
let result = catch_unwind(AssertUnwindSafe(|| operation(&mut pile, &signer)));
drop(pile);
match result {
Ok(result) => result,
Err(panic) => resume_unwind(panic),
}
});
}
let mut pile = open_pile_strict_as(&self.pile, signer.verifying_key())?;
let result = operation(&mut pile, &signer);
finish_pile(pile, result)
}
pub fn close(&self) -> Result<()> {
let Some(shared) = &self.shared else {
return Ok(());
};
let session = shared
.lock()
.unwrap_or_else(|error| error.into_inner())
.take();
match session {
Some(session) => session.close(),
None => Ok(()),
}
}
pub fn finish<T>(&self, result: Result<T>) -> Result<T> {
finish_close(result, self.close())
}
}
fn finish_close<T>(result: Result<T>, close: Result<()>) -> Result<T> {
match (result, close) {
(Ok(value), Ok(())) => Ok(value),
(Ok(_), Err(error)) | (Err(error), Ok(())) => Err(error),
(Err(error), Err(close_error)) => {
Err(error.context(format!("closing storage also failed: {close_error:#}")))
}
}
}
pub fn runtime() -> Result<tokio::runtime::Runtime> {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.context("create faculty I/O runtime")
}
pub async fn read<S, T>(
store: &mut S,
snapshot: &S::Snapshot,
mut read: impl FnMut(&S::Snapshot) -> Result<T>,
) -> Result<T>
where
S: SnapshotSource + AsyncBlobStoreAcquire,
S::Snapshot: BlobStoreList,
{
let mut reader = snapshot.clone();
loop {
let error = match read(&reader) {
Ok(value) => return Ok(value),
Err(error) => error,
};
let Some(missing) = error
.chain()
.find_map(|error| error.downcast_ref::<MissingBlob>())
else {
return Err(error);
};
if reader.contains_blob(missing.handle)? {
return Err(error);
}
if store.acquire(missing.handle).await?.is_none() {
return Err(error);
}
reader = store.snapshot()?;
}
}
pub fn open_store(path: &Path) -> Result<FacultyStore> {
lazy_store(|| open_pile_strict(path))
}
pub fn open_store_as(path: &Path, host: VerifyingKey) -> Result<FacultyStore> {
lazy_store(|| open_pile_strict_as(path, host))
}
fn lazy_store(open: impl FnOnce() -> Result<Pile>) -> Result<FacultyStore> {
use iroh_base::{EndpointAddr, EndpointId};
use iroh_tickets::endpoint::EndpointTicket;
use rand_core::RngCore;
use triblespace_net::peer::{PeerConfig, ReconcileDirection, ReconcileQos};
let routes = std::env::var("TRIBLESPACE_PEERS").or_else(|error| match error {
std::env::VarError::NotPresent => Ok(String::new()),
error => Err(error),
})?;
let peers = routes
.split(',')
.map(str::trim)
.filter(|route| !route.is_empty())
.map(|route| {
if let Ok(ticket) = route.parse::<EndpointTicket>() {
return Ok(EndpointAddr::from(ticket));
}
route
.parse::<EndpointId>()
.map(EndpointAddr::from)
.with_context(|| format!("invalid TRIBLESPACE_PEERS endpoint {route:?}"))
})
.collect::<Result<Vec<_>>>()?;
let mut secret = [0; 32];
rand_core::OsRng
.try_fill_bytes(&mut secret)
.context("generate foreground transport identity")?;
let key = SigningKey::from_bytes(&secret);
use zeroize::Zeroize;
secret.zeroize();
Ok(FacultyStore::lazy(
open()?,
key,
PeerConfig {
peers,
qos: ReconcileQos {
direction: ReconcileDirection::ReadOnly,
},
provider_publication_budget: Some(0),
bind: None,
},
))
}
pub fn open_secrets_collection<S>(
store: &mut S,
subject: VerifyingKey,
) -> Result<crate::secrets::storage::SecretsCollection>
where
S: CollectionStoreExt + SnapshotSource,
S::Snapshot: BlobStoreGet,
{
let scope = crate::secrets::DEFAULT_SCOPE_ID;
let Some(handle) = crate::collection_names::configured_handle(scope)? else {
return crate::secrets::storage::SecretsCollection::register(
store,
crate::collection_names::require_name(scope),
crate::collection_names::private_policy(subject).with_capability(
crate::secrets::key_delivery_definition(),
triblespace::core::collection::AdmissionPolicy::direct(subject),
),
)
.context("register signer-private Secrets descriptor with key-delivery policy");
};
let snapshot = store
.snapshot()
.context("freeze configured Secrets descriptor")?;
let source = crate::collection_names::open_exact_in(&snapshot, scope, handle)
.context("open configured Secrets source collection")?;
drop(snapshot);
crate::secrets::storage::SecretsCollection::from_source(store, source)
.context("register maintained Secrets collection descriptors")
}
pub fn open_secrets_collection_read<S>(
store: &mut S,
subject: VerifyingKey,
) -> Result<crate::secrets::storage::SecretsCollection>
where
S: CollectionStoreExt + SnapshotSource,
S::Snapshot: BlobStoreGet,
{
open_secrets_collection(store, subject)
}
pub fn open_secrets_collection_acquiring<S>(
store: &mut S,
subject: VerifyingKey,
runtime: &Arc<tokio::runtime::Runtime>,
) -> Result<crate::secrets::storage::SecretsCollection>
where
S: CollectionStoreExt + SnapshotSource,
S::Snapshot: BlobStoreGet + triblespace::core::repo::async_store::AsyncBlobStoreGet,
{
let scope = crate::secrets::DEFAULT_SCOPE_ID;
let Some(handle) = crate::collection_names::configured_handle(scope)? else {
return open_secrets_collection(store, subject);
};
let snapshot = store
.snapshot()
.context("freeze configured Secrets descriptor")?;
let reader = AcquiringReader::new(snapshot, Arc::clone(runtime));
let source = crate::collection_names::open_exact_in(&reader, scope, handle)
.context("open configured Secrets source collection")?;
crate::secrets::storage::SecretsCollection::from_source(store, source)
.context("register maintained Secrets collection descriptors")
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TargetDiscovery {
commits: Vec<CollectionCommit>,
merges: Vec<CollectionMerge>,
derives: Vec<CollectionDerive>,
}
impl TargetDiscovery {
pub fn commits(&self) -> &[CollectionCommit] {
&self.commits
}
pub fn merges(&self) -> &[CollectionMerge] {
&self.merges
}
pub fn derives(&self) -> &[CollectionDerive] {
&self.derives
}
}
pub fn discover_target<S>(
store: &mut S,
scope: Id,
authority: VerifyingKey,
) -> Result<TargetDiscovery>
where
S: CollectionStoreExt + SnapshotSource,
<S as SnapshotSource>::Snapshot: BlobStoreGet + CapabilityProofRead + CollectionRead,
{
let collection = crate::collection_names::open_configured(store, scope, authority)
.context("open target collection descriptor")?;
let snapshot = store
.snapshot()
.context("freeze target collection store snapshot")?;
let selectors = std::collections::BTreeSet::from([CollectionRecordSelector::Collection(
collection.handle(),
)]);
let records = snapshot
.select_records(&selectors)
.context("discover native collection records")?;
let mut commits = Vec::new();
let mut merges = Vec::new();
let mut derives = Vec::new();
for record in records {
match record {
CollectionRecord::Commit(commit) => commits.push(commit),
CollectionRecord::Merge(merge) => merges.push(merge),
CollectionRecord::Derive(derive) => derives.push(derive),
CollectionRecord::Map(_) => {}
}
}
Ok(TargetDiscovery {
commits,
merges,
derives,
})
}
pub fn read_fact_collection<S>(
collection: Collection<SimpleArchive>,
snapshot: &S,
) -> Result<(TribleSet, Support)>
where
S: StoreRead,
{
let observed = snapshot
.collection(collection)
.context("attach realized collection")?;
let support = observed
.support()
.context("resolve realized fact collection support")?
.clone();
let facts = observed
.view::<TribleSet>()
.context("read authorized collection facts")?;
Ok((facts, support))
}
pub fn underived<R, S, T>(snapshot: &R, source: Collection<S>, view: Collection<T>) -> Result<usize>
where
R: StoreRead,
S: CollectionEncoding,
T: CollectionEncoding,
{
let foundations = source
.admitted(snapshot)
.context("read the foundations a view's source stands for")?;
let coverage = snapshot
.coverage(&BTreeSet::from([view.handle()]))
.context("read a view's leaves")?;
let mut underived = 0;
for foundation in foundations.members() {
let mut usable = false;
for output in coverage.leaf_outputs(view.handle(), SourceLocator::of(foundation.raw)) {
if snapshot
.metadata(Handle::<UnknownBlob>::from_hash(output))
.context("inspect a view leaf's output")?
.is_some()
{
usable = true;
break;
}
}
if !usable {
underived += 1;
}
}
Ok(underived)
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct FactLag {
pub succinct: usize,
pub rank9: usize,
}
impl FactLag {
pub fn of<R>(
snapshot: &R,
_source: Collection<SimpleArchive>,
succinct: Collection<SuccinctArchiveBlob>,
rank9: Collection<Rank9AcceleratedSuccinctArchiveBlob>,
) -> Result<Self>
where
R: StoreRead,
{
Ok(Self {
succinct: snapshot
.attached(succinct)
.context("read the Succinct attachments")?
.residual()
.len(),
rank9: snapshot
.attached(rank9)
.context("read the Rank9 attachments")?
.residual()
.len(),
})
}
pub const fn is_current(self) -> bool {
self.succinct == 0 && self.rank9 == 0
}
}
impl std::fmt::Display for FactLag {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
formatter,
"{} source foundation(s) not yet attached to Succinct, {} not yet attached to Rank9",
self.succinct, self.rank9
)
}
}
pub trait FactRead: StoreRead + Sized {
fn read_facts(
&self,
rank9: Collection<Rank9AcceleratedSuccinctArchiveBlob>,
) -> Result<FactArchive> {
attached_facts(
&self
.attached(rank9)
.context("attach the Rank9 fact cover")?,
)
}
}
impl<R: StoreRead> FactRead for R {}
pub fn acquire_facts<R: StoreRead>(
reader: &R,
rank9: Collection<Rank9AcceleratedSuccinctArchiveBlob>,
) -> Result<FactArchive> {
acquire_attached_facts(
&reader
.attached_acquiring(rank9)
.context("select Rank9 fact cover")?,
)
}
pub fn acquire_attached_facts<R: StoreRead>(
attached: &AttachedSnapshot<R, Rank9AcceleratedSuccinctArchiveBlob>,
) -> Result<FactArchive> {
require_complete_attached_read(
succinctarchive_union::read_attached_acquiring(attached)
.context("acquire the selected Rank9 facts")?,
)
}
pub fn require_complete_attached_read<V>(read: AttachedRead<V>) -> Result<V> {
let (value, unread) = read.into_parts();
if unread.is_empty() {
Ok(value)
} else {
Err(IncompleteAttachedRead { unread }.into())
}
}
#[derive(Debug)]
pub struct IncompleteAttachedRead {
pub unread: Support<SimpleArchive>,
}
impl std::fmt::Display for IncompleteAttachedRead {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let named = self
.unread
.members()
.take(8)
.map(|handle| hex::encode(handle.raw))
.collect::<Vec<_>>()
.join(", ");
write!(
formatter,
"cannot read {} selected foundation(s) of collection {}: {}{}",
self.unread.len(),
hex::encode(self.unread.collection().handle().raw),
named,
if self.unread.len() > 8 { ", ..." } else { "" },
)
}
}
impl std::error::Error for IncompleteAttachedRead {}
pub async fn acquire_attached_facts_async<R>(
attached: &AttachedSnapshot<R, Rank9AcceleratedSuccinctArchiveBlob>,
) -> Result<FactArchive>
where
R: StoreRead + triblespace::core::repo::async_store::AsyncBlobStoreGet + 'static,
{
let attached = attached.clone();
let runtime = tokio::runtime::Handle::current();
tokio::task::spawn_blocking(move || acquire_attached_facts(&attached.acquiring(runtime)))
.await
.context("join selected fact acquisition")?
}
pub async fn acquire_attached_read_async<E, V>(
attached: &AttachedSnapshot<
impl StoreRead + triblespace::core::repo::async_store::AsyncBlobStoreGet + 'static,
E,
>,
) -> Result<AttachedRead<V>>
where
E: CollectionAttachment + Send,
V: triblespace::core::collection::TryFromCover<E> + Send + 'static,
{
let attached = attached.clone();
let runtime = tokio::runtime::Handle::current();
tokio::task::spawn_blocking(move || {
attached
.acquiring(runtime)
.read_acquiring::<V>()
.map_err(anyhow::Error::new)
})
.await
.context("join selected attached acquisition")?
}
pub trait FactView {
fn facts(&self) -> Result<FactArchive>;
}
impl<R: StoreRead> FactView for AttachedSnapshot<R, Rank9AcceleratedSuccinctArchiveBlob> {
fn facts(&self) -> Result<FactArchive> {
attached_facts(self)
}
}
pub fn attached_facts<R>(
attached: &AttachedSnapshot<R, Rank9AcceleratedSuccinctArchiveBlob>,
) -> Result<FactArchive>
where
R: StoreRead,
{
attached_facts_read(attached).map(AttachedRead::into_value)
}
pub fn attached_facts_read<R>(
attached: &AttachedSnapshot<R, Rank9AcceleratedSuccinctArchiveBlob>,
) -> Result<AttachedRead<FactArchive>>
where
R: StoreRead,
{
succinctarchive_union::read_attached(attached)
.map_err(|error| anyhow!("read the attached Rank9 facts: {error}"))
}
pub fn tolerate_own_lag<T>(
result: std::result::Result<T, CollectionRealizationError>,
) -> std::result::Result<(), CollectionRealizationError> {
match result {
Ok(_)
| Err(CollectionRealizationError::Unmappable { .. })
| Err(CollectionRealizationError::UnauthorizedProducer { .. })
| Err(CollectionRealizationError::HostMismatch { .. }) => Ok(()),
Err(error) => Err(error),
}
}
pub fn signer_path(pile: &Path, explicit: Option<&Path>) -> PathBuf {
signing_key_file::resolve_path(explicit, pile)
}
pub fn load_signer(pile: &Path, explicit: Option<&Path>) -> Result<SigningKey> {
let path = signer_path(pile, explicit);
signing_key_file::load_existing(&path)
.with_context(|| format!("load durable signing key {}", path.display()))
}
pub fn initialize_signer(pile: &Path, explicit: Option<&Path>) -> Result<SigningKey> {
let path = signer_path(pile, explicit);
signing_key_file::init(&path)
.with_context(|| format!("initialize durable signing key {}", path.display()))
}
pub fn open_pile_strict(path: &Path) -> Result<Pile> {
refreshed(
path,
Pile::open(path).with_context(|| format!("open pile {}", path.display()))?,
)
}
pub fn open_pile_strict_as(path: &Path, host: VerifyingKey) -> Result<Pile> {
refreshed(
path,
Pile::open_as(path, host).with_context(|| format!("open pile {}", path.display()))?,
)
}
pub fn open_pile_signed(pile: &Path, key: Option<&Path>) -> Result<(Pile, SigningKey)> {
let signer = load_signer(pile, key)?;
let opened = open_pile_strict_as(pile, signer.verifying_key())?;
Ok((opened, signer))
}
fn refreshed(path: &Path, mut pile: Pile) -> Result<Pile> {
if let Err(error) = pile.refresh() {
let close = pile.close();
let mut failure = pile_read_error(path, error);
if let Err(close_error) = close {
failure = failure.context(format!(
"closing pile after failed refresh also failed: {close_error}"
));
}
return Err(failure);
}
Ok(pile)
}
pub fn publish_fragment(
pile_path: &Path,
key_path: Option<&Path>,
scope: Id,
fragment: Fragment,
) -> Result<CollectionCommit> {
let mut commits = publish_fragments(pile_path, key_path, scope, [fragment])?;
Ok(commits
.pop()
.expect("one input fragment produces one collection commit"))
}
pub fn carry_scope(pile_path: &Path, key_path: Option<&Path>, scope: Id) -> Result<()> {
let (mut pile, signer) = open_pile_signed(pile_path, key_path)?;
let collection =
crate::collection_names::open_configured(&mut pile, scope, signer.verifying_key())
.context("open native collection descriptor")?;
carry_facts(&mut pile, collection, &signer);
finish_pile(pile, Ok(()))
}
pub fn publish_fragments(
pile_path: &Path,
key_path: Option<&Path>,
scope: Id,
fragments: impl IntoIterator<Item = Fragment>,
) -> Result<Vec<CollectionCommit>> {
let (mut pile, signer) = open_pile_signed(pile_path, key_path)?;
let collection =
crate::collection_names::open_configured(&mut pile, scope, signer.verifying_key())
.context("open native collection descriptor")?;
let result = (|| {
let mut commits = Vec::new();
for fragment in fragments {
commits.push(
pile.commit(collection, &signer, fragment)
.context("publish native collection fragment")?,
);
}
if !commits.is_empty() {
drop(
pollster::block_on(ensure_downstream(&mut pile, collection, &signer))
.context("fragments were committed, but ensuring their derived views failed")?,
);
}
Ok(commits)
})();
finish_pile(pile, result)
}
pub fn pile_read_error(path: &Path, error: ReadError) -> anyhow::Error {
match error {
ReadError::CorruptPile { valid_length } => anyhow!(
"pile {} has a malformed or incomplete known record at byte {valid_length}; this \
reader cannot prove that the remaining bytes are a disposable torn write. The pile \
was left unchanged. Upgrade `trible` to the matching current source cohort, then \
inspect that boundary with `trible pile diagnose record-at {} {valid_length}` before \
considering any destructive action",
path.display(),
path.display()
),
ReadError::UnsupportedRecord { .. } => anyhow!(
"pile {} contains a record format unsupported by this binary ({error}); this is \
likely version skew. Upgrade to a reader that recognizes the marker. The pile was \
left unchanged",
path.display()
),
other => anyhow!("refresh pile {}: {other}", path.display()),
}
}
fn finish_pile<T>(pile: Pile, result: Result<T>) -> Result<T> {
let close = pile.close();
match (result, close) {
(Ok(value), Ok(())) => Ok(value),
(Ok(_), Err(error)) => Err(anyhow!("close pile: {error}")),
(Err(error), Ok(())) => Err(error),
(Err(error), Err(close_error)) => {
Err(error.context(format!("closing pile also failed: {close_error}")))
}
}
}
#[cfg(test)]
mod tests {
use std::fs::{self, File};
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use anybytes::View;
use ed25519_dalek::SigningKey;
use triblespace::core::blob::encodings::simplearchive::SimpleArchive;
use triblespace::core::blob::encodings::utf8string::UTF8String;
use triblespace::core::collection::{empty_metadata_handle, CollectionRecord, CollectionStore};
use triblespace::core::inline::encodings::hash::Handle;
use triblespace::core::inline::Inline;
use triblespace::core::metadata;
use triblespace::core::repo::memoryrepo::MemoryRepo;
use triblespace::core::repo::{BlobStoreGet, SnapshotSource};
use triblespace::core::trible::TribleSet;
use triblespace::macros::entity;
use super::*;
static NEXT_TEST: AtomicU64 = AtomicU64::new(0);
struct TestFiles {
directory: PathBuf,
pile: PathBuf,
key: PathBuf,
}
impl TestFiles {
fn new() -> Self {
let serial = NEXT_TEST.fetch_add(1, Ordering::Relaxed);
let directory = std::env::temp_dir().join(format!(
"faculties-native-collection-{}-{serial}",
std::process::id()
));
let _ = fs::remove_dir_all(&directory);
fs::create_dir_all(&directory).unwrap();
let pile = directory.join("test.pile");
File::create(&pile).unwrap();
let key = directory.join("test.key");
Self {
directory,
pile,
key,
}
}
}
impl Drop for TestFiles {
fn drop(&mut self) {
let _ = fs::remove_dir_all(&self.directory);
}
}
fn id(byte: u8) -> Id {
Id::new([byte; 16]).unwrap()
}
#[test]
fn shared_configuration_is_inert_and_thread_shareable() {
fn assert_send_sync<T: Send + Sync>() {}
assert_send_sync::<Storage>();
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("not-opened.pile");
let storage = Storage::shared(path.clone(), Some(dir.path().join("missing.key")));
let clone = storage.clone();
assert_eq!(storage.path(), path);
clone.close().unwrap();
assert_eq!(dir.path().read_dir().unwrap().count(), 0);
}
#[test]
fn shared_operations_keep_the_peer_runtime_and_local_backend() {
use triblespace::core::repo::BlobStorePut;
let files = TestFiles::new();
initialize_signer(&files.pile, Some(&files.key)).unwrap();
let storage = Storage::shared(files.pile.clone(), Some(files.key.clone()));
let clone = storage.clone();
let (peer, runtime) = storage
.with_store(|store, _, runtime| Ok((store.id(), Arc::clone(runtime))))
.unwrap();
let moved = files.directory.join("still-open.pile");
fs::rename(&files.pile, &moved).unwrap();
let handle = clone
.with_pile(|pile, _| Ok(pile.put::<UTF8String, _>("shared resident bytes")?))
.unwrap();
storage
.scope(|storage| {
storage.with_store(|store, _, current_runtime| {
assert_eq!(store.id(), peer);
assert!(Arc::ptr_eq(current_runtime, &runtime));
let snapshot = store.snapshot()?;
let text: View<str> = current_runtime.block_on(snapshot.get(handle))?;
assert_eq!(&*text, "shared resident bytes");
Ok(())
})
})
.unwrap();
clone
.with_store(|store, _, _| {
assert_eq!(
store.id(),
peer,
"nested scope must not close its application owner"
);
Ok(())
})
.unwrap();
assert!(
!files.pile.exists(),
"no operation may reopen the configured pathname"
);
storage.close().unwrap();
clone.close().unwrap();
let mut reopened = open_pile_strict(&moved).unwrap();
assert!(reopened.snapshot().unwrap().contains_blob(handle).unwrap());
reopened.close().unwrap();
}
#[test]
fn retained_store_observes_external_appends_without_changing_old_snapshots() {
use triblespace::core::repo::BlobStorePut;
let files = TestFiles::new();
initialize_signer(&files.pile, Some(&files.key)).unwrap();
let storage = Storage::shared(files.pile.clone(), Some(files.key.clone()));
let before = storage
.with_store(|store, _, _| Ok(store.snapshot()?))
.unwrap();
let mut other = open_pile_strict(&files.pile).unwrap();
let handle = other
.put::<UTF8String, _>("appended by another process")
.unwrap();
other.close().unwrap();
let after = storage
.with_store(|store, _, _| Ok(store.snapshot()?))
.unwrap();
assert!(!before.contains_blob(handle).unwrap());
assert!(after.contains_blob(handle).unwrap());
assert!(
!before.contains_blob(handle).unwrap(),
"refresh must not mutate an earlier observation"
);
storage.close().unwrap();
}
#[test]
fn shared_operation_errors_do_not_discard_the_connection() {
let files = TestFiles::new();
initialize_signer(&files.pile, Some(&files.key)).unwrap();
let storage = Storage::shared(files.pile.clone(), Some(files.key.clone()));
let peer = storage.with_store(|store, _, _| Ok(store.id())).unwrap();
let error = storage
.with_store::<()>(|_, _, _| anyhow::bail!("operation failed"))
.unwrap_err();
assert_eq!(error.to_string(), "operation failed");
storage
.with_store(|store, _, _| {
assert_eq!(store.id(), peer);
Ok(())
})
.unwrap();
storage.close().unwrap();
}
#[test]
fn panicking_local_handler_does_not_poison_shared_storage() {
let files = TestFiles::new();
initialize_signer(&files.pile, Some(&files.key)).unwrap();
let storage = Storage::shared(files.pile.clone(), Some(files.key.clone()));
let peer = storage.with_store(|store, _, _| Ok(store.id())).unwrap();
let result = catch_unwind(AssertUnwindSafe(|| {
storage.with_pile::<()>(|_, _| panic!("handler panic"))
}));
assert!(result.is_err());
storage
.with_store(|store, _, _| {
assert_eq!(store.id(), peer);
Ok(())
})
.unwrap();
storage.close().unwrap();
}
#[test]
fn independent_owners_do_not_share_by_path() {
let files = TestFiles::new();
initialize_signer(&files.pile, Some(&files.key)).unwrap();
let first = Storage::shared(files.pile.clone(), Some(files.key.clone()));
let second = Storage::shared(files.pile.clone(), Some(files.key.clone()));
let first_peer = first.with_store(|store, _, _| Ok(store.id())).unwrap();
let second_peer = second.with_store(|store, _, _| Ok(store.id())).unwrap();
assert_ne!(
first_peer, second_peer,
"sharing is explicit, never an ambient path cache"
);
first.close().unwrap();
second.close().unwrap();
}
#[test]
fn compound_cli_scope_reuses_and_closes_its_own_store() {
let files = TestFiles::new();
initialize_signer(&files.pile, Some(&files.key)).unwrap();
let cli = Storage::new(files.pile.clone(), Some(files.key.clone()));
let owner = cli
.scope(|storage| {
let peer = storage.with_store(|store, _, _| Ok(store.id()))?;
storage.with_store(|store, _, _| {
assert_eq!(store.id(), peer);
Ok(())
})?;
Ok(storage.clone())
})
.unwrap();
assert!(owner.shared.as_ref().unwrap().lock().unwrap().is_none());
assert!(cli.shared.is_none());
}
#[test]
fn live_read_fetches_only_demanded_bytes_and_preserves_its_observation() {
use anybytes::Bytes;
use triblespace::core::blob::encodings::UnknownBlob;
use triblespace::core::repo::{BlobStorePut, WantRead};
struct Supply {
pile: Pile,
payload: Bytes,
requested: Vec<Inline<Handle<UnknownBlob>>>,
}
impl SnapshotSource for Supply {
type Snapshot = triblespace::core::repo::pile::PileSnapshot;
type SnapshotError = <Pile as SnapshotSource>::SnapshotError;
fn snapshot(&mut self) -> Result<Self::Snapshot, Self::SnapshotError> {
self.pile.snapshot()
}
}
impl AsyncBlobStoreAcquire for Supply {
type AcquireError = std::io::Error;
async fn acquire(
&mut self,
handle: Inline<Handle<UnknownBlob>>,
) -> Result<Option<Bytes>, Self::AcquireError> {
self.requested.push(handle);
let stored = self
.pile
.put::<UnknownBlob, _>(self.payload.clone())
.unwrap();
assert_eq!(stored, handle);
Ok(Some(self.payload.clone()))
}
}
let files = TestFiles::new();
let payload: Bytes = Vec::from("selected body").into();
let mut source = MemoryRepo::default();
let handle = source.put::<UnknownBlob, _>(payload.clone()).unwrap();
let mut store = Supply {
pile: open_pile_strict(&files.pile).unwrap(),
payload,
requested: Vec::new(),
};
let before = store.snapshot().unwrap();
let value = pollster::block_on(read(&mut store, &before, |reader| {
let bytes = reader.get::<Bytes, UnknownBlob>(handle)?;
Ok(bytes)
}))
.unwrap();
assert_eq!(&*value, b"selected body");
assert_eq!(store.requested, [handle]);
assert!(!before.contains_blob(handle).unwrap());
assert!(store.snapshot().unwrap().wants().unwrap().next().is_none());
let resident = store.snapshot().unwrap();
let result = pollster::block_on(read(&mut store, &resident, |_| {
Ok(before.get::<Bytes, UnknownBlob>(handle)?)
}));
assert!(result.is_err());
assert_eq!(store.requested, [handle]);
drop(resident);
drop(before);
store.pile.close().unwrap();
}
#[test]
fn strict_open_reports_evidence_without_prescribing_data_loss() {
let files = TestFiles::new();
fs::write(&files.pile, [0xFF; 8]).unwrap();
let before = fs::read(&files.pile).unwrap();
let error = open_pile_strict(&files.pile)
.err()
.expect("malformed pile must fail strict open");
let rendered = format!("{error:#}");
assert!(rendered.contains("malformed or incomplete known record at byte 0"));
assert!(rendered.contains("cannot prove"));
assert!(rendered.contains("matching current source cohort"));
assert!(rendered.contains("pile diagnose record-at"));
assert!(!rendered.contains("pile amputate"));
assert_eq!(fs::read(&files.pile).unwrap(), before);
let mut unsupported = [0u8; 256];
unsupported[..16].fill(0xA5);
fs::write(&files.pile, unsupported).unwrap();
let error = open_pile_strict(&files.pile)
.err()
.expect("unsupported marker must fail strict open");
let rendered = format!("{error:#}");
assert!(rendered.contains("unsupported by this binary"));
assert!(rendered.contains("likely version skew"));
assert!(!rendered.contains("pile amputate"));
assert_eq!(fs::read(&files.pile).unwrap(), unsupported);
}
#[test]
fn target_discovery_registers_descriptor_without_definition_record() {
let signer = SigningKey::from_bytes(&[7; 32]);
let team = signer.verifying_key();
let target_scope = crate::schemas::wiki::DEFAULT_SCOPE_ID;
let other_scope = crate::schemas::compass::DEFAULT_SCOPE_ID;
let mut store = MemoryRepo::default();
let target = crate::collection_names::open(&mut store, target_scope, team)
.unwrap()
.handle();
let other = crate::collection_names::open(&mut store, other_scope, team)
.unwrap()
.handle();
let target_commit = CollectionCommit::sign(
&signer,
target,
Inline::new([1; 32]),
empty_metadata_handle(),
);
let other_commit = CollectionCommit::sign(
&signer,
other,
Inline::new([2; 32]),
empty_metadata_handle(),
);
let target_merge = CollectionMerge::sign(
&signer,
target,
[target_commit.data(), Inline::new([4; 32])],
Inline::new([5; 32]),
)
.unwrap();
let other_merge = CollectionMerge::sign(
&signer,
other,
[other_commit.data(), Inline::new([7; 32])],
Inline::new([8; 32]),
)
.unwrap();
let derive_to_target = CollectionDerive::sign(
&signer,
target,
triblespace::core::collection::SourceLocator::of(other_commit.data().raw),
Inline::new([10; 32]),
);
let derive_from_target = CollectionDerive::sign(
&signer,
other,
triblespace::core::collection::SourceLocator::of(target_commit.data().raw),
Inline::new([12; 32]),
);
for record in [
CollectionRecord::Commit(target_commit),
CollectionRecord::Commit(other_commit),
CollectionRecord::Merge(target_merge),
CollectionRecord::Merge(other_merge),
CollectionRecord::Derive(derive_to_target),
CollectionRecord::Derive(derive_from_target),
] {
store.insert(record).unwrap();
}
let discovered = discover_target(&mut store, target_scope, team).unwrap();
assert_eq!(discovered.commits(), &[target_commit]);
assert_eq!(discovered.merges(), &[target_merge]);
assert_eq!(discovered.derives(), &[derive_to_target]);
assert!(
!store.blobs.is_empty(),
"registration retains the descriptor attachment closure"
);
}
#[test]
fn operational_fact_read_names_unread_selected_foundations() {
use triblespace::core::repo::BlobStorePut;
let signer = SigningKey::from_bytes(&[7; 32]);
let mut source = MemoryRepo::for_host(signer.verifying_key());
let root = crate::collection_names::open(
&mut source,
crate::schemas::wiki::DEFAULT_SCOPE_ID,
signer.verifying_key(),
)
.unwrap();
let (_, rank9) = fact_pair(&mut source, root).unwrap();
let committed = source
.commit(root, &signer, entity! { metadata::tag: &id(9) })
.unwrap();
let selected_payload = committed.data();
let full = source.snapshot().unwrap();
let mut local = MemoryRepo::for_host(signer.verifying_key());
for blob in full.blobs() {
let handle = blob.unwrap().handle;
if handle.raw != selected_payload.raw {
local
.put::<UnknownBlob, _>(full.get::<anybytes::Bytes, _>(handle).unwrap())
.unwrap();
}
}
for record in full.records().unwrap() {
local.insert(record.unwrap()).unwrap();
}
let frozen = local.snapshot().unwrap();
let selected = frozen.attached(rank9).unwrap();
let partial = attached_facts_read(&selected).unwrap();
assert_eq!(partial.unread().len(), 1);
let error = acquire_attached_facts(&selected).err().unwrap();
assert!(error.downcast_ref::<IncompleteAttachedRead>().is_some());
assert!(error
.to_string()
.contains(&hex::encode(selected_payload.raw)));
assert!(selected.support().is_empty());
assert_eq!(selected.residual().len(), 1);
}
#[test]
fn the_attached_pair_maintains_a_shard_preserving_rank9_view() {
let signer = SigningKey::from_bytes(&[7; 32]);
let mut store = MemoryRepo::for_host(signer.verifying_key());
let source = crate::collection_names::open(
&mut store,
crate::schemas::wiki::DEFAULT_SCOPE_ID,
signer.verifying_key(),
)
.unwrap();
let (_succinct, rank9) = fact_pair(&mut store, source).unwrap();
let fragment = entity! {
metadata::tag: &id(9),
metadata::name: "maintained facts",
};
let expected = fragment.facts().clone();
store.commit(source, &signer, fragment).unwrap();
let after = pollster::block_on(async {
drop(store.ensure(source, &signer).await.unwrap());
maintain_downstream(&mut store, source, &signer)
.await
.unwrap();
store.snapshot().unwrap()
});
let observed = after.attached(rank9).unwrap();
let view = observed.view::<FactArchive>().unwrap();
let actual: TribleSet = view.iter().collect();
assert_eq!(actual, expected);
assert_eq!(observed.support().len(), 1);
assert!(observed.residual().is_empty());
assert_eq!(view.segment_count(), 1);
assert_eq!(discovered_records(&after).unwrap().maps().len(), 2);
}
#[test]
fn every_commit_is_readable_through_the_pair_whoever_wrote_it() {
use triblespace::core::collection::grant_collection_write;
let owner = SigningKey::from_bytes(&[11; 32]);
let writer = SigningKey::from_bytes(&[12; 32]);
let mut store = MemoryRepo::for_host(owner.verifying_key());
let source = crate::collection_names::open(
&mut store,
crate::schemas::wiki::DEFAULT_SCOPE_ID,
owner.verifying_key(),
)
.unwrap();
let (succinct, rank9) = fact_pair(&mut store, source).unwrap();
let lag = |store: &mut MemoryRepo| {
FactLag::of(&store.snapshot().unwrap(), source, succinct, rank9).unwrap()
};
let facts = |store: &mut MemoryRepo| -> TribleSet {
store
.snapshot()
.unwrap()
.read_facts(rank9)
.unwrap()
.iter()
.collect()
};
grant_collection_write(&mut store, source.handle(), &owner, writer.verifying_key())
.unwrap();
let written = entity! { metadata::name: "another key's write" };
store.commit(source, &writer, written.clone()).unwrap();
assert_eq!(
lag(&mut store),
FactLag {
succinct: 1,
rank9: 1
}
);
assert_eq!(facts(&mut store), written.facts().clone());
let own = entity! { metadata::name: "owner" };
store.commit(source, &owner, own.clone()).unwrap();
drop(pollster::block_on(ensure_downstream(&mut store, source, &owner)).unwrap());
assert!(lag(&mut store).is_current(), "{:?}", lag(&mut store));
assert_eq!(
facts(&mut store),
(written.clone() + own.clone()).facts().clone()
);
let later = entity! { metadata::name: "later" };
store.commit(source, &writer, later.clone()).unwrap();
let before = store.snapshot().unwrap().records().unwrap().count();
drop(pollster::block_on(ensure_downstream(&mut store, source, &writer)).unwrap());
assert_eq!(store.snapshot().unwrap().records().unwrap().count(), before);
assert_eq!(
facts(&mut store),
(written.clone() + own.clone() + later).facts().clone()
);
}
#[test]
fn a_write_is_done_whatever_an_unrelated_attached_collection_does() {
use anybytes::Bytes;
use triblespace::core::blob::encodings::succinctarchive::SuccinctArchiveBlob;
use triblespace::core::blob::encodings::UnknownBlob;
use triblespace::core::blob::{Blob, IntoBlob};
use triblespace::core::collection::records::{
collection_mapping, collection_parent, collection_representation,
KIND_COLLECTION_DESCRIPTOR,
};
use triblespace::core::collection::{CollectionAttachment, CollectionCommit};
use triblespace::core::metadata::MetaDescribe;
use triblespace::core::repo::BlobStorePut;
use triblespace::core::trible::Fragment;
use triblespace_search::portable_bm25::PortableBM25Blob;
use triblespace_search::text_bm25::{Bm25Tokenizer, TextAttributeToBm25};
let owner = SigningKey::from_bytes(&[13; 32]);
let mut store = MemoryRepo::for_host(owner.verifying_key());
let source = crate::collection_names::open(
&mut store,
crate::schemas::wiki::DEFAULT_SCOPE_ID,
owner.verifying_key(),
)
.unwrap();
let (succinct, rank9) = fact_pair(&mut store, source).unwrap();
let described = |store: &mut MemoryRepo, text: Blob<UTF8String>| {
let text = store.put::<UTF8String, _>(text).unwrap();
let mut fragment = Fragment::empty();
fragment += entity! { metadata::description: text };
store.commit(source, &owner, fragment).unwrap();
};
let write = |store: &mut MemoryRepo, name: &'static str| {
store
.commit(source, &owner, entity! { metadata::name: name })
.unwrap();
pollster::block_on(ensure_downstream(store, source, &owner))
};
let second = crate::collection_names::open(
&mut store,
crate::schemas::compass::DEFAULT_SCOPE_ID,
owner.verifying_key(),
)
.unwrap();
let (facts, mut blobs) = entity! {
metadata::tag: KIND_COLLECTION_DESCRIPTOR,
collection_parent*: [source.handle(), second.handle()],
collection_representation*: <SuccinctArchiveBlob as MetaDescribe>::describe(),
collection_mapping*: <SuccinctArchiveBlob as CollectionAttachment>::fragment(&()),
}
.into_facts_and_blobs();
for (_, blob) in blobs.snapshot().unwrap() {
store.put::<UnknownBlob, _>(blob).unwrap();
}
let both = store.put::<SimpleArchive, _>(facts).unwrap();
let stray = store
.put::<SimpleArchive, _>(TribleSet::new().to_blob())
.unwrap();
store
.insert(CollectionRecord::Commit(CollectionCommit::sign(
&owner,
both,
Handle::<SimpleArchive>::to_hash(stray),
empty_metadata_handle(),
)))
.unwrap();
assert!(store
.snapshot()
.unwrap()
.collections()
.unwrap()
.contains(&both));
let report = write(&mut store, "beside a second parent").unwrap();
assert!(!report.realized.contains(&both));
assert!(report.failed_attached.is_empty());
assert!(
FactLag::of(&store.snapshot().unwrap(), source, succinct, rank9)
.unwrap()
.is_current()
);
let index = store
.attach::<PortableBM25Blob>(
source,
TextAttributeToBm25 {
attribute: metadata::description.id(),
tokenizer: Bm25Tokenizer::Word,
},
)
.unwrap();
described(&mut store, "well formed".to_blob());
pollster::block_on(seed_attached(&mut store, index, &owner)).unwrap();
described(&mut store, Blob::new(Bytes::from(vec![0xFF, 0xFE])));
let report = write(&mut store, "beside a failing index").unwrap();
assert_eq!(
report
.failed_attached
.iter()
.map(|(attached, _)| attached.handle)
.collect::<Vec<_>>(),
[index.handle()]
);
assert!(report.realized.contains(&rank9.handle()));
assert!(
FactLag::of(&store.snapshot().unwrap(), source, succinct, rank9)
.unwrap()
.is_current()
);
}
#[test]
fn fact_read_stands_for_nothing_beneath_unadmitted_commits_and_reads_the_rollup_once_they_are_admitted(
) {
use triblespace::core::blob::IntoBlob;
use triblespace::core::repo::BlobStorePut;
let owner = SigningKey::from_bytes(&[7; 32]);
let author = SigningKey::from_bytes(&[8; 32]);
let mut store = MemoryRepo::for_host(owner.verifying_key());
let collection = crate::collection_names::open(
&mut store,
crate::schemas::wiki::DEFAULT_SCOPE_ID,
owner.verifying_key(),
)
.unwrap();
let left = entity! { metadata::name: "left" };
let right = entity! { metadata::name: "right" };
let expected = (left.clone() + right.clone()).facts().clone();
let commits: Vec<_> = [left, right]
.into_iter()
.map(|fragment| {
let blob = IntoBlob::<SimpleArchive>::to_blob(fragment.facts().clone());
CollectionCommit::sign(
&author,
collection.handle(),
Handle::<SimpleArchive>::to_hash(blob.get_handle()),
empty_metadata_handle(),
)
})
.collect();
for commit in &commits {
store.insert(CollectionRecord::Commit(*commit)).unwrap();
}
let joined = store.put::<SimpleArchive, _>(expected.clone()).unwrap();
store
.insert(CollectionRecord::Merge(
CollectionMerge::sign(
&owner,
collection.handle(),
[commits[0].data(), commits[1].data()],
Handle::<SimpleArchive>::to_hash(joined),
)
.unwrap(),
))
.unwrap();
let snapshot = store.snapshot().unwrap();
assert!(collection.admitted(&snapshot).unwrap().is_empty());
for commit in &commits {
assert!(!snapshot
.contains_blob(Handle::<SimpleArchive>::from_hash(commit.data()))
.unwrap());
}
let (facts, support) = read_fact_collection(collection, &snapshot).unwrap();
assert_eq!(facts.len(), 0);
assert!(support.is_empty());
for commit in &commits {
store
.insert(CollectionRecord::Commit(CollectionCommit::sign(
&owner,
collection.handle(),
commit.data(),
empty_metadata_handle(),
)))
.unwrap();
}
let snapshot = store.snapshot().unwrap();
assert_eq!(collection.admitted(&snapshot).unwrap().len(), 2);
for commit in &commits {
assert!(!snapshot
.contains_blob(Handle::<SimpleArchive>::from_hash(commit.data()))
.unwrap());
}
let (facts, support) = read_fact_collection(collection, &snapshot).unwrap();
assert_eq!(facts, expected);
assert_eq!(support.len(), 2);
}
#[test]
fn storage_opens_as_its_signer_so_its_own_merges_are_read_and_a_keyless_open_reads_wider() {
use triblespace::core::collection::MERGE_FAN_IN;
let files = TestFiles::new();
let signer = initialize_signer(&files.pile, Some(&files.key)).unwrap();
let scope = crate::schemas::wiki::DEFAULT_SCOPE_ID;
let names = ["a", "b", "c", "d", "e", "f", "g", "h"];
assert_eq!(names.len(), MERGE_FAN_IN);
publish_fragments(
&files.pile,
Some(&files.key),
scope,
names.map(|name| entity! { _ @ metadata::name: name }),
)
.unwrap();
let mut pile = open_pile_strict_as(&files.pile, signer.verifying_key()).unwrap();
let collection =
crate::collection_names::open_configured(&mut pile, scope, signer.verifying_key())
.unwrap();
drop(pollster::block_on(pile.maintain(collection, &signer)).unwrap());
pile.close().unwrap();
let observe = |pile: &mut Pile| -> Result<(usize, TribleSet)> {
let snapshot = pile.snapshot()?;
let observed = snapshot.collection(collection)?;
Ok((observed.cover().len(), observed.view::<TribleSet>()?))
};
let (width, facts) = Storage::new(files.pile.clone(), Some(files.key.clone()))
.with_pile(|pile, _| observe(pile))
.unwrap();
assert_eq!(width, 1, "a one-shot open reads the signer's merge");
let shared = Storage::shared(files.pile.clone(), Some(files.key.clone()));
let (shared_width, shared_facts) = shared.with_pile(|pile, _| observe(pile)).unwrap();
shared.close().unwrap();
assert_eq!(
shared_width, 1,
"the shared session store is opened as the signer"
);
assert_eq!(shared_facts, facts);
let mut keyless = open_pile_strict(&files.pile).unwrap();
let (wide, wide_facts) = observe(&mut keyless).unwrap();
keyless.close().unwrap();
assert_eq!(wide, MERGE_FAN_IN);
assert_eq!(wide_facts, facts);
assert!(!facts.is_empty());
}
#[test]
fn carry_scope_opens_as_the_signer_so_the_fact_pair_attaches_its_root_merge() {
use triblespace::core::collection::{CoverageRead, MERGE_FAN_IN};
let files = TestFiles::new();
let signer = initialize_signer(&files.pile, Some(&files.key)).unwrap();
let scope = crate::schemas::wiki::DEFAULT_SCOPE_ID;
let names = ["a", "b", "c", "d", "e", "f", "g", "h"];
assert_eq!(names.len(), MERGE_FAN_IN);
publish_fragments(
&files.pile,
Some(&files.key),
scope,
names.map(|name| entity! { _ @ metadata::name: name }),
)
.unwrap();
let (mut pile, loaded) = open_pile_signed(&files.pile, Some(&files.key)).unwrap();
assert_eq!(loaded.verifying_key(), signer.verifying_key());
let host = pile
.snapshot()
.unwrap()
.index(&BTreeSet::new())
.unwrap()
.host();
assert_eq!(
host.map(|host| host.raw),
Some(signer.verifying_key().to_bytes())
);
let collection =
crate::collection_names::open_configured(&mut pile, scope, signer.verifying_key())
.unwrap();
drop(pollster::block_on(pile.maintain(collection, &signer)).unwrap());
pile.close().unwrap();
carry_scope(&files.pile, Some(&files.key), scope).unwrap();
let rank9_cover = |pile: &mut Pile| -> (usize, usize) {
let (_, rank9) = fact_pair(pile, collection).unwrap();
let snapshot = pile.snapshot().unwrap();
let attached = snapshot.attached(rank9).unwrap();
(attached.cover().len(), attached.residual().len())
};
let mut pile = open_pile_strict_as(&files.pile, signer.verifying_key()).unwrap();
assert_eq!(
rank9_cover(&mut pile),
(1, 0),
"one attachment of the merged node covers the tier"
);
pile.close().unwrap();
let mut keyless = open_pile_strict(&files.pile).unwrap();
assert_eq!(rank9_cover(&mut keyless), (0, MERGE_FAN_IN));
keyless.close().unwrap();
}
#[test]
fn a_shared_store_refuses_a_signing_key_replaced_under_it() {
let files = TestFiles::new();
initialize_signer(&files.pile, Some(&files.key)).unwrap();
let shared = Storage::shared(files.pile.clone(), Some(files.key.clone()));
shared.with_pile(|_, _| Ok(())).unwrap();
let other = files.directory.join("other.key");
initialize_signer(&files.pile, Some(&other)).unwrap();
fs::copy(&other, &files.key).unwrap();
let error = shared.with_pile(|_, _| Ok(())).unwrap_err().to_string();
assert!(
error.contains("signing key changed while the shared faculty store was open"),
"{error}"
);
let error = shared.with_store(|_, _, _| Ok(())).unwrap_err().to_string();
assert!(error.contains("restart the process"), "{error}");
shared.close().unwrap();
Storage::shared(files.pile.clone(), Some(files.key.clone()))
.with_pile(|_, _| Ok(()))
.unwrap();
}
#[test]
fn publication_conserves_both_fact_channels_and_attachments_and_replays_idempotently() {
let files = TestFiles::new();
initialize_signer(&files.pile, Some(&files.key)).unwrap();
let mut fragment = entity! { _ @ metadata::name: "content attachment" };
let content_root = fragment.root().unwrap();
let description = entity! { _ @ metadata::name: "metadata attachment" };
let metadata_root = description.root().unwrap();
fragment.describe_with(description);
let expected_facts = fragment.facts().clone();
let expected_metafacts = fragment.metafacts().clone();
assert!(!expected_facts.is_empty());
assert!(!expected_metafacts.is_empty());
let team = load_signer(&files.pile, Some(&files.key))
.unwrap()
.verifying_key();
let target_scope = crate::schemas::wiki::DEFAULT_SCOPE_ID;
let other_scope = crate::schemas::compass::DEFAULT_SCOPE_ID;
let first = publish_fragment(
&files.pile,
Some(&files.key),
target_scope,
fragment.clone(),
)
.unwrap();
let after_first = fs::metadata(&files.pile).unwrap().len();
let unrelated = entity! { _ @ metadata::tag: &id(9) };
publish_fragment(&files.pile, Some(&files.key), other_scope, unrelated).unwrap();
let before_replay = fs::metadata(&files.pile).unwrap().len();
let repeated =
publish_fragment(&files.pile, Some(&files.key), target_scope, fragment).unwrap();
let after_replay = fs::metadata(&files.pile).unwrap().len();
assert_eq!(repeated, first);
assert!(before_replay > after_first);
assert_eq!(after_replay, before_replay);
let mut pile = open_pile_strict(&files.pile).unwrap();
let target_collection =
crate::collection_names::open(&mut pile, target_scope, team).unwrap();
let target = discover_target(&mut pile, target_scope, team).unwrap();
assert_eq!(target.commits(), &[first]);
assert_eq!(target.commits()[0].collection(), target_collection.handle());
assert!(target.merges().is_empty());
assert!(target.derives().is_empty());
let unrelated_target = discover_target(&mut pile, other_scope, team).unwrap();
assert_eq!(unrelated_target.commits().len(), 1);
let reader = pile.snapshot().unwrap();
let data_handle = Handle::<SimpleArchive>::from_hash(first.data());
let actual_facts: TribleSet = reader.get(data_handle).unwrap();
let actual_metafacts: TribleSet = reader.get(first.metadata()).unwrap();
assert_eq!(actual_facts, expected_facts);
assert_eq!(actual_metafacts, expected_metafacts);
let content_handle = actual_facts
.iter()
.find(|fact| fact.e() == &content_root && fact.a() == &metadata::name.id())
.map(|fact| *fact.v::<Handle<UTF8String>>())
.expect("content attachment handle");
let content: View<str> = reader.get(content_handle).unwrap();
assert_eq!(&*content, "content attachment");
let metadata_handle = actual_metafacts
.iter()
.find(|fact| fact.e() == &metadata_root && fact.a() == &metadata::name.id())
.map(|fact| *fact.v::<Handle<UTF8String>>())
.expect("metadata attachment handle");
let metadata_text: View<str> = reader.get(metadata_handle).unwrap();
assert_eq!(&*metadata_text, "metadata attachment");
pile.close().unwrap();
}
#[test]
fn missing_signer_fails_before_the_pile_is_touched() {
let files = TestFiles::new();
let missing = files.directory.join("missing.key");
let before = fs::metadata(&files.pile).unwrap().len();
let error = publish_fragment(
&files.pile,
Some(&missing),
id(1),
entity! { _ @ metadata::tag: &id(2) },
)
.unwrap_err();
assert!(format!("{error:#}").contains("load durable signing key"));
assert!(!missing.exists());
assert_eq!(fs::metadata(&files.pile).unwrap().len(), before);
}
#[cfg(feature = "local-embed")]
#[test]
fn a_leaf_whose_output_is_not_here_leaves_its_foundation_underived() {
use triblespace::core::collection::{AdmissionPolicy, CollectionPolicy};
use triblespace::core::repo::BlobStorePut;
use triblespace::core::trible::Trible;
use triblespace_search::nvfp4::{NvFp4CosineSet, NvFp4EmbeddingAttribute};
use triblespace_search::schemas::Embedding;
let key = SigningKey::from_bytes(&[0x78; 32]);
let policy = CollectionPolicy::new(
AdmissionPolicy::direct(key.verifying_key()),
AdmissionPolicy::direct(key.verifying_key()),
);
let attribute = Id::new([0xA7; 16]).unwrap();
let mut store = MemoryRepo::default();
let source = store.collection("vectors", policy.clone()).unwrap();
let view = store
.derive::<NvFp4CosineSet<Embedding>>(
source,
NvFp4EmbeddingAttribute::new(attribute, 3).unwrap(),
policy,
)
.unwrap();
let commit = |store: &mut MemoryRepo, entity: u8, vector: Vec<f32>| {
let embedding = store.put::<Embedding, _>(vector).unwrap();
let mut facts = TribleSet::new();
facts.insert(&Trible::force(
&Id::new([entity; 16]).unwrap(),
&attribute,
&embedding,
));
store
.commit(source, &key, Fragment::from(facts))
.unwrap()
.data()
};
commit(&mut store, 1, vec![1.0, 0.0, 0.0]);
commit(&mut store, 2, vec![0.0, 1.0, 0.0]);
drop(pollster::block_on(store.ensure(view, &key)).unwrap());
let count =
|store: &mut MemoryRepo| underived(&store.snapshot().unwrap(), source, view).unwrap();
assert_eq!(count(&mut store), 0);
let waiting = commit(&mut store, 3, vec![0.0, 0.0, 1.0]);
store
.insert(CollectionRecord::Derive(CollectionDerive::sign(
&key,
view.handle(),
SourceLocator::of(waiting.raw),
Inline::new([0x77; 32]),
)))
.unwrap();
assert_eq!(count(&mut store), 1);
commit(&mut store, 4, vec![0.6, 0.8, 0.0]);
assert_eq!(count(&mut store), 2);
}
}
#[cfg(test)]
pub(crate) struct DiscoveredRecords {
commits: Vec<triblespace::core::collection::CollectionCommit>,
merges: Vec<triblespace::core::collection::CollectionMerge>,
derives: Vec<triblespace::core::collection::CollectionDerive>,
maps: Vec<triblespace::core::collection::CollectionMap>,
}
#[cfg(test)]
impl DiscoveredRecords {
pub(crate) fn maps(&self) -> &[triblespace::core::collection::CollectionMap] {
&self.maps
}
pub(crate) fn commits(&self) -> &[triblespace::core::collection::CollectionCommit] {
&self.commits
}
pub(crate) fn merges(&self) -> &[triblespace::core::collection::CollectionMerge] {
&self.merges
}
pub(crate) fn derives(&self) -> &[triblespace::core::collection::CollectionDerive] {
&self.derives
}
}
#[cfg(test)]
pub(crate) fn discovered_records<S: triblespace::core::collection::CollectionRead>(
snapshot: &S,
) -> anyhow::Result<DiscoveredRecords> {
use triblespace::core::collection::CollectionRecord;
let mut discovered = DiscoveredRecords {
commits: Vec::new(),
merges: Vec::new(),
derives: Vec::new(),
maps: Vec::new(),
};
for record in snapshot
.records()
.map_err(|error| anyhow::anyhow!("{error}"))?
{
match record.map_err(|error| anyhow::anyhow!("{error}"))? {
CollectionRecord::Commit(commit) => discovered.commits.push(commit),
CollectionRecord::Merge(merge) => discovered.merges.push(merge),
CollectionRecord::Derive(derive) => discovered.derives.push(derive),
CollectionRecord::Map(map) => discovered.maps.push(map),
}
}
Ok(discovered)
}
pub fn carry_facts<S>(pile: &mut S, source: Collection<SimpleArchive>, signer: &SigningKey)
where
S: Store + AsyncBlobStoreAcquire + Send,
{
drop(pollster::block_on(maintain_downstream(pile, source, signer)).unwrap());
}
#[derive(Clone, Debug, Default)]
pub struct FacultiesRealizer {
pub lagging: Vec<(CollectionHandle, Vec<(CollectionData, String)>)>,
}
impl<S> RealizeDerived<S> for FacultiesRealizer
where
S: Store + AsyncBlobStoreAcquire + Send,
{
async fn realize(
&mut self,
store: &mut S,
derived: &Derived,
signer: &SigningKey,
upkeep: Upkeep,
) -> Result<Realized, CollectionRealizationError> {
match CoreRealizer.realize(store, derived, signer, upkeep).await {
Err(CollectionRealizationError::Unmappable { blocked }) => {
self.lagging.push((derived.handle, blocked));
Ok(Realized::Done)
}
realized => realized,
}
}
async fn realize_attached(
&mut self,
store: &mut S,
attached: &Attached,
signer: &SigningKey,
upkeep: Upkeep,
) -> Result<Realized, CollectionRealizationError> {
if attached.representation == <PortableBM25Blob as MetaDescribe>::id() {
realize_attached_as::<S, PortableBM25Blob>(store, attached, signer, upkeep).await
} else {
CoreRealizer
.realize_attached(store, attached, signer, upkeep)
.await
}
}
}
pub fn fact_pair<S>(
pile: &mut S,
source: Collection<SimpleArchive>,
) -> Result<(
Collection<SuccinctArchiveBlob>,
Collection<Rank9AcceleratedSuccinctArchiveBlob>,
)>
where
S: Store + AsyncBlobStoreAcquire + Send,
{
let succinct = pile
.attach::<SuccinctArchiveBlob>(source, ())
.map_err(|error| anyhow!("attach the Succinct collection: {error}"))?;
let rank9 = pile
.attach::<Rank9AcceleratedSuccinctArchiveBlob>(source, succinct)
.map_err(|error| anyhow!("attach the Rank9 collection: {error}"))?;
Ok((succinct, rank9))
}
pub async fn seed_attached<S, T>(
pile: &mut S,
target: Collection<T>,
signer: &SigningKey,
) -> Result<()>
where
S: Store + AsyncBlobStoreAcquire + Send,
T: CollectionAttachment,
Handle<T>: InlineEncoding,
{
if listed(pile, target.handle())? {
return Ok(());
}
match pile.ensure_attached(target, signer).await {
Ok(_) | Err(CollectionRealizationError::HostMismatch { .. }) => Ok(()),
Err(error) => Err(anyhow!("seed a newly attached collection: {error}")),
}
}
fn listed<S>(pile: &mut S, collection: CollectionHandle) -> Result<bool>
where
S: Store,
{
Ok(pile
.snapshot()
.context("freeze the store for its collection listing")?
.collections()
.map_err(|error| anyhow!("list the store's collections: {error}"))?
.contains(&collection))
}
async fn seed_fact_pair<S>(
pile: &mut S,
succinct: Collection<SuccinctArchiveBlob>,
rank9: Collection<Rank9AcceleratedSuccinctArchiveBlob>,
signer: &SigningKey,
) -> Result<()>
where
S: Store + AsyncBlobStoreAcquire + Send,
{
seed_attached(pile, succinct, signer).await?;
seed_attached(pile, rank9, signer).await
}
pub async fn ensure_downstream<S>(
pile: &mut S,
source: Collection<SimpleArchive>,
signer: &SigningKey,
) -> Result<UpkeepReport>
where
S: Store + AsyncBlobStoreAcquire + Send,
{
let (succinct, rank9) = fact_pair(pile, source)?;
seed_fact_pair(pile, succinct, rank9, signer).await?;
let mut realizer = FacultiesRealizer::default();
match core_ensure_downstream(pile, source.handle(), signer, &mut realizer).await {
Err(CollectionRealizationError::HostMismatch { .. }) => Ok(UpkeepReport::default()),
result => result
.map_err(|error| anyhow!("ensure the collections downstream of the source: {error}")),
}
}
pub async fn maintain_downstream<S>(
pile: &mut S,
source: Collection<SimpleArchive>,
signer: &SigningKey,
) -> Result<UpkeepReport>
where
S: Store + AsyncBlobStoreAcquire + Send,
{
let (succinct, rank9) = fact_pair(pile, source)?;
seed_fact_pair(pile, succinct, rank9, signer).await?;
let mut realizer = FacultiesRealizer::default();
core_maintain_downstream(pile, source.handle(), signer, &mut realizer)
.await
.map_err(|error| anyhow!("maintain the collections downstream of the source: {error}"))
}