use std::sync::Arc;
use triblespace::core::collection::AttachedSnapshot;
use anyhow::{anyhow, Context, Result};
use ed25519_dalek::SigningKey;
use triblespace::core::blob::encodings::succinctarchive::{
Rank9AcceleratedSuccinctArchiveBlob, SuccinctArchiveBlob,
};
use triblespace::core::blob::encodings::{simplearchive::SimpleArchive, UnknownBlob};
use triblespace::core::blob::Blob;
use triblespace::core::collection::{
Collection, CollectionCommit, CollectionSnapshotExt, CollectionStoreExt,
};
use triblespace::core::inline::encodings::UnknownInline;
use triblespace::core::query::TriblePattern;
use triblespace::core::repo::pile::{Pile, PileSnapshot};
use triblespace::core::repo::{BlobStorePut, SnapshotSource, StorageClose, Store};
use triblespace::core::repo::async_store::{AsyncBlobStoreAcquire, AsyncBlobStoreGet};
use triblespace::prelude::*;
use crate::schemas::code::DEFAULT_SCOPE_ID;
use crate::storage::{FactArchive, FacultyStore};
pub struct CodeImportWriter<P = FacultyStore> {
pile: P,
collection: Collection<SimpleArchive>,
signer: SigningKey,
current: FactArchive,
delta: Fragment,
runtime: Arc<tokio::runtime::Runtime>,
}
impl CodeImportWriter {
pub fn open(
pile_path: &std::path::Path,
key_path: Option<&std::path::Path>,
) -> Result<Self> {
let signer = crate::storage::load_signer(pile_path, key_path)?;
let runtime = Arc::new(crate::storage::runtime()?);
let mut pile = crate::storage::open_store_as(pile_path, signer.verifying_key())?;
let result = Self::prepare(&mut pile, &signer, &runtime);
match result {
Ok((collection, current)) => {
let mut writer = Self {
pile,
collection,
signer,
current,
delta: Fragment::empty(),
runtime,
};
if let Err(error) = writer.stage_fragment(crate::code::law_fragment()) {
return close_pile(
writer.pile,
Err(error),
"closing Code pile after extraction-law staging failed",
);
}
Ok(writer)
}
Err(error) => close_pile(
pile,
Err(error),
"closing Code pile after failed open also failed",
),
}
}
}
impl<P> CodeImportWriter<P>
where
P: Store + AsyncBlobStoreAcquire + Send,
P::Snapshot: AsyncBlobStoreGet,
{
pub fn from_store(
mut pile: P,
signer: &SigningKey,
runtime: Arc<tokio::runtime::Runtime>,
) -> Result<Self> {
let (source, current) = Self::prepare(&mut pile, signer, &runtime)?;
let mut writer = Self {
pile,
collection: source,
signer: signer.clone(),
current,
delta: Fragment::empty(),
runtime,
};
writer.stage_fragment(crate::code::law_fragment())?;
Ok(writer)
}
fn prepare(
pile: &mut P,
signer: &SigningKey,
runtime: &Arc<tokio::runtime::Runtime>,
) -> Result<(Collection<SimpleArchive>, FactArchive)> {
let source = crate::collection_names::open_configured_acquiring(
pile, DEFAULT_SCOPE_ID, signer.verifying_key(), runtime,
)?;
let (succinct, rank9) = crate::storage::fact_pair(pile, source)?;
runtime.block_on(async {
crate::storage::tolerate_own_lag(pile.maintain_attached(succinct, signer).await)?;
crate::storage::tolerate_own_lag(pile.maintain_attached(rank9, signer).await)
}).context("maintain Code import facts")?;
let reader = crate::storage::AcquiringReader::new(pile.snapshot()?, runtime.clone());
let current = crate::storage::acquire_facts(&reader, rank9)
.context("read Code import facts")?;
Ok((source, current))
}
pub fn holds_entity(&self, entity: Id) -> bool {
use triblespace::core::metadata;
exists!((tag: Id), pattern!(&self.current, [{ entity @ metadata::tag: ?tag }]))
|| exists!((tag: Id), pattern!(self.delta.facts(), [{ entity @ metadata::tag: ?tag }]))
}
pub fn stage_fragment(&mut self, fragment: Fragment) -> Result<()> {
let (_, facts, metafacts, blobs) = fragment.into_parts();
if facts.iter().all(|fact| {
self.delta.facts().contains(fact) || fact_archive_contains(&self.current, fact)
}) {
return Ok(());
}
let embedded = embedded_blobs(blobs);
stage_embedded_blobs(&mut self.pile, embedded)?;
self.delta += Fragment::from_parts(facts, metafacts, Default::default());
Ok(())
}
pub fn delta_len(&self) -> usize {
self.delta.facts().len()
}
pub fn commit_unit(&mut self) -> Result<Option<CollectionCommit>> {
if self.delta.facts().is_empty() {
return Ok(None);
}
let fragment = std::mem::replace(&mut self.delta, Fragment::empty());
let published = fragment.facts().clone();
crate::collection_names::require_command_write_admission_acquiring(
&mut self.pile,
self.collection,
&self.signer,
"Code",
"code find",
&self.runtime,
)?;
let commit = self
.pile
.commit(self.collection, &self.signer, fragment)
.context("commit authored Code projection unit")?;
self.current = extend_archive(&self.current, &published);
Ok(Some(commit))
}
}
impl CodeImportWriter {
pub fn close<T>(mut self, surrounding: Result<T>) -> Result<T> {
let result = surrounding.and_then(|value| {
self.commit_unit()?;
Ok(value)
});
close_pile(
self.pile,
result,
"closing Code pile after failure also failed",
)
}
}
fn stage_embedded_blobs<S>(store: &mut S, embedded: Vec<Blob<UnknownBlob>>) -> Result<()>
where
S: BlobStorePut,
{
for blob in embedded {
store
.put::<UnknownBlob, _>(blob)
.context("stage Code embedded blob")?;
}
Ok(())
}
fn embedded_blobs(mut blobs: triblespace::core::blob::MemoryBlobStore) -> Vec<Blob<UnknownBlob>> {
let reader = blobs
.snapshot()
.expect("MemoryBlobStore reader creation is infallible");
let mut embedded: Vec<_> = reader.iter().collect();
embedded.sort_unstable_by_key(|(store_key, _)| store_key.raw);
embedded.into_iter().map(|(_, blob)| blob).collect()
}
fn fact_archive_contains(facts: &FactArchive, fact: &Trible) -> bool {
exists!(facts.pattern(
inlineencodings::GenId::inline_from(*fact.e()),
inlineencodings::GenId::inline_from(*fact.a()),
*fact.v::<UnknownInline>(),
))
}
fn extend_archive(current: &FactArchive, additions: &TribleSet) -> FactArchive {
if additions.is_empty() {
return current.clone();
}
current.with_segments([
triblespace::core::blob::encodings::succinctarchive::SuccinctArchive::from(additions),
])
}
fn close_pile<T>(pile: impl StorageClose, result: Result<T>, failure_context: &str) -> Result<T> {
match (result, pile.close()) {
(Ok(value), Ok(())) => Ok(value),
(Err(error), Ok(())) => Err(error),
(Ok(_), Err(close_error)) => Err(anyhow!("close Code pile: {close_error}")),
(Err(error), Err(close_error)) => {
Err(error.context(format!("{failure_context} also failed: {close_error}")))
}
}
}
pub async fn ensure_facts(
pile: &mut Pile,
source: Collection<SimpleArchive>,
signer: &SigningKey,
) -> Result<AttachedSnapshot<PileSnapshot, Rank9AcceleratedSuccinctArchiveBlob>> {
let succinct = pile
.attach::<SuccinctArchiveBlob>(source, ())
.context("register Succinct Code fact collection")?;
let rank9 = pile
.attach::<Rank9AcceleratedSuccinctArchiveBlob>(source, succinct)
.context("register Rank9 Code fact collection")?;
crate::storage::tolerate_own_lag(pile.maintain_attached(succinct, signer).await)
.context("maintain Succinct Code fact collection")?;
crate::storage::tolerate_own_lag(pile.maintain_attached(rank9, signer).await)
.context("maintain Rank9 Code fact collection")?;
pile.snapshot()
.context("freeze maintained Code facts")?
.attached(rank9)
.context("attach Code fact collection")
}
pub fn ensure_local_with_storage(
storage: &crate::storage::Storage,
) -> Result<AttachedSnapshot<PileSnapshot, Rank9AcceleratedSuccinctArchiveBlob>> {
storage.with_store(|store, signer, runtime| {
let source = crate::collection_names::open_configured_acquiring(
store, DEFAULT_SCOPE_ID, signer.verifying_key(), runtime,
)?;
let mut local = store.store();
let pile = &mut *local;
pollster::block_on(ensure_facts(pile, source, signer))
})
}