use super::{CaptureRow, IntervalValue, RawHandle, TextHandle};
use crate::clock;
#[cfg(test)]
use crate::collection_names::open_configured;
use crate::files::presentation::{present, ViewOptions};
use crate::out::Part;
use crate::schemas::body::{capture, DEFAULT_SCOPE_ID, KIND_CAPTURE, KIND_INTENT};
use crate::storage::FactRead;
use crate::storage::{FactArchive, Storage};
use anybytes::Bytes;
use anyhow::{bail, Context, Result};
use hifitime::Epoch;
use std::path::PathBuf;
use triblespace::core::blob::encodings::succinctarchive::{
Rank9AcceleratedSuccinctArchiveBlob, SuccinctArchiveBlob,
};
use triblespace::core::collection::{CollectionSnapshotExt, CollectionStoreExt};
use triblespace::core::metadata;
use crate::storage::{AcquiringReader, FacultySnapshot};
use triblespace::core::repo::{BlobStoreGet, SnapshotSource};
use triblespace::prelude::*;
#[derive(Clone, Debug)]
pub enum Signal {
Vision {
bytes: Bytes,
mime: String,
width: u64,
height: u64,
},
Audio {
bytes: Bytes,
mime: String,
},
Touch,
}
#[derive(Clone, Debug)]
pub struct CaptureInput {
pub signal: Signal,
pub pose: String,
pub note: Option<String>,
}
impl CaptureInput {
pub fn validate(&self) -> Result<()> {
let (mime, kind) = match &self.signal {
Signal::Vision {
mime,
width,
height,
..
} => {
if *width == 0 || *height == 0 {
bail!("vision dimensions must be positive");
}
(mime, "image")
}
Signal::Audio { mime, .. } => (mime, "audio"),
Signal::Touch => return Ok(()),
};
if mime.len() > 32 || mime.contains('\0') {
bail!("Body MIME must fit 32 UTF-8 bytes and contain no NUL");
}
let parsed: mime::Mime = mime.parse().context("invalid capture MIME type")?;
if parsed.type_().as_str() != kind || parsed.essence_str().contains('*') {
bail!("{kind} capture requires a concrete {kind} MIME type");
}
Ok(())
}
}
#[derive(Clone, Debug)]
pub struct CaptureReceipt {
pub id: Id,
pub modality: String,
pub bytes: usize,
pub dimensions: Option<(u64, u64)>,
pub note: Option<String>,
}
#[derive(Clone, Debug)]
pub struct CaptureSummary {
pub id: Id,
pub tai_ns: i128,
pub modality: String,
pub note: String,
}
#[derive(Clone, Debug)]
pub struct CaptureExport {
pub id: Id,
pub bytes: Bytes,
pub mime: Option<String>,
}
impl CaptureExport {
pub fn uri(&self) -> String {
format!("body:{:x}", self.id)
}
}
#[derive(Clone, Debug)]
pub struct IntentObservation {
pub id: Id,
pub tai_ns: i128,
pub text: String,
}
#[derive(Clone, Debug)]
pub struct Body {
storage: Storage,
}
impl Body {
pub fn new(pile: PathBuf, key: Option<PathBuf>) -> Self {
Self::with_storage(Storage::new(pile, key))
}
pub fn with_storage(storage: Storage) -> Self {
Self { storage }
}
fn storage(&self) -> BodyStorage<'_> {
BodyStorage {
storage: &self.storage,
}
}
pub fn capture(&self, input: &CaptureInput) -> Result<CaptureReceipt> {
let fragment = capture_fragment(input, clock::point_now()?)?;
let id = fragment.root().expect("capture has one intrinsic root");
let (modality, bytes, dimensions) = match &input.signal {
Signal::Vision {
bytes,
width,
height,
..
} => ("vision", bytes.len(), Some((*width, *height))),
Signal::Audio { bytes, .. } => ("audio", bytes.len(), None),
Signal::Touch => ("touch", 0, None),
};
self.storage().publish(fragment)?;
Ok(CaptureReceipt {
id,
modality: modality.to_owned(),
bytes,
dimensions,
note: input.note.clone(),
})
}
pub fn set_intent(&self, text: &str) -> Result<IntentObservation> {
let created = clock::point_now()?;
let fragment = intent_fragment(text, created);
let id = fragment.root().expect("intent has one intrinsic root");
self.storage().publish(fragment)?;
Ok(IntentObservation {
id,
tai_ns: interval_key(created),
text: text.to_owned(),
})
}
pub fn intent(&self) -> Result<Option<IntentObservation>> {
self.storage().with_indexed_view(latest_intent)
}
pub fn list(&self) -> Result<Vec<CaptureSummary>> {
self.storage().with_view(|space, reader| {
let mut rows = Vec::new();
for (id, modality, created) in find!(
(c: Id, m: String, t: IntervalValue),
pattern!(space, [{ ?c @ metadata::tag: KIND_CAPTURE, capture::modality: ?m, metadata::created_at: ?t }])
) {
let note = find!(h: TextHandle, pattern!(space, [{ id @ capture::note: ?h }]))
.next().map(|handle| reader.get::<View<str>, _>(handle)
.map(|text| text.to_string())
.map_err(|error| anyhow::anyhow!("read note for capture {id:X}: {error}")))
.transpose()?.unwrap_or_default();
rows.push(CaptureSummary { id, tai_ns: interval_key(created), modality, note });
}
rows.sort_by(|a, b| (b.tai_ns, b.id).cmp(&(a.tai_ns, a.id)));
Ok(rows)
})
}
pub fn get(&self, id: &str) -> Result<CaptureExport> {
validate_selector(id)?;
self.storage().with_view(|space, reader| {
let ids = find!(c: Id, pattern!(space, [{ ?c @ metadata::tag: KIND_CAPTURE }]));
let needle = id.to_ascii_lowercase();
let mut matches = ids.filter(|candidate| format!("{candidate:x}").starts_with(&needle));
let id = matches
.next()
.with_context(|| format!("no capture matching '{id}'"))?;
if matches.next().is_some() {
bail!("ambiguous capture prefix '{needle}'");
}
let handle = find!(h: RawHandle, pattern!(space, [{ id @ capture::frame: ?h }]))
.next()
.context("capture has no frame payload (a touch capture has no file)")?;
let mime = find!(m: String, pattern!(space, [{ id @ capture::mime: ?m }])).next();
let bytes = reader
.get(handle)
.map_err(|error| anyhow::anyhow!("read frame for capture {id:X}: {error}"))?;
Ok(CaptureExport { id, bytes, mime })
})
}
pub fn view(&self, id: &str, options: &ViewOptions) -> Result<Part> {
options.validate()?;
let export = self.get(id)?;
let mime = export
.mime
.as_deref()
.context("capture has no MIME type; use get for original bytes")?;
present(export.bytes, mime, options)
}
}
pub fn validate_selector(id: &str) -> Result<()> {
if id.is_empty() || id.len() > 32 || !id.bytes().all(|b| b.is_ascii_hexdigit()) {
bail!(
"capture selector must be a nonempty hexadecimal ID or prefix, at most 32 characters"
);
}
Ok(())
}
pub fn capture_fragment(input: &CaptureInput, created: IntervalValue) -> Result<Fragment> {
input.validate()?;
super::validate_point("capture creation time", created)?;
let mut fragment = Fragment::empty();
let (modality, frame, mime, width, height) = match &input.signal {
Signal::Vision {
bytes,
mime,
width,
height,
} => (
"vision",
Some(fragment.put::<blobencodings::RawBytes, _>(bytes.clone())),
Some(mime.clone()),
Some((*width).to_inline()),
Some((*height).to_inline()),
),
Signal::Audio { bytes, mime } => (
"audio",
Some(fragment.put::<blobencodings::RawBytes, _>(bytes.clone())),
Some(mime.clone()),
None,
None,
),
Signal::Touch => ("touch", None, None, None, None),
};
let pose = fragment.put::<blobencodings::UTF8String, _>(input.pose.clone());
let note = input
.note
.as_ref()
.map(|note| fragment.put::<blobencodings::UTF8String, _>(note.clone()));
fragment += super::capture_record(&CaptureRow {
id: KIND_CAPTURE,
created_at: created,
frame,
mime,
width,
height,
modality: modality.to_owned(),
note,
pose,
});
Ok(fragment)
}
fn intent_fragment(text: &str, created: IntervalValue) -> Fragment {
let mut fragment = Fragment::empty();
let text = fragment.put::<blobencodings::UTF8String, _>(text.to_owned());
fragment += super::intent_record(&super::IntentRow {
id: KIND_INTENT,
created_at: created,
text,
});
fragment
}
pub(super) fn interval_key(interval: IntervalValue) -> i128 {
let (lower, _): (Epoch, Epoch) = interval.try_from_inline().expect("valid TAI interval");
lower.to_tai_duration().total_nanoseconds()
}
fn latest_intent<R: BlobStoreGet>(snapshot: &super::BodySnapshot<R>) -> Result<Option<IntentObservation>> {
let Some(row) = super::latest_intent(snapshot.facts(), snapshot.intent_register())? else {
return Ok(None);
};
let text: View<str> = snapshot
.store_snapshot()
.get(row.text)
.map_err(|error| anyhow::anyhow!("read latest intent {:X}: {error}", row.id))?;
Ok(Some(IntentObservation {
id: row.id,
tai_ns: interval_key(row.created_at),
text: text.to_string(),
}))
}
#[derive(Clone, Copy)]
struct BodyStorage<'a> {
storage: &'a Storage,
}
impl BodyStorage<'_> {
fn publish(&self, fragment: Fragment) -> Result<()> {
self.storage.with_store(|store, signer, runtime| {
let collection = crate::collection_names::open_configured_acquiring(
store, DEFAULT_SCOPE_ID, signer.verifying_key(), runtime,
)?;
crate::collection_names::require_command_write_admission_acquiring(
store, collection, signer, "Body", "body show", runtime,
)?;
let pile = store;
pile.commit(collection, signer, fragment)
.context("publish native Body collection fragment")?;
runtime.block_on(async {
let intents = super::intent_register_collection(pile, signer.verifying_key())?;
crate::storage::seed_attached(pile, intents, signer).await?;
crate::storage::ensure_downstream(pile, collection, signer).await?;
Ok::<_, anyhow::Error>(())
})
.context("Body facts were committed, but ensuring its derived views failed")?;
Ok(())
})
}
fn with_view<T>(&self, f: impl FnOnce(&FactArchive, &AcquiringReader<FacultySnapshot>) -> Result<T>) -> Result<T> {
self.storage.with_store(|pile, signer, runtime| {
let source = crate::collection_names::open_configured_acquiring(
pile, DEFAULT_SCOPE_ID, signer.verifying_key(), runtime,
)?;
let collection_succinct = pile.attach::<SuccinctArchiveBlob>(source, ())?;
let collection_rank9 =
pile.attach::<Rank9AcceleratedSuccinctArchiveBlob>(source, collection_succinct)?;
let store_snapshot = AcquiringReader::new(
pile.snapshot().context("freeze Body fact collection")?, runtime.clone(),
);
let facts = crate::storage::acquire_facts(&store_snapshot, collection_rank9)
.context("read maintained Body fact collection")?;
f(&facts, &store_snapshot)
})
}
fn with_indexed_view<T>(&self, f: impl FnOnce(&super::BodySnapshot<AcquiringReader<FacultySnapshot>>) -> Result<T>) -> Result<T> {
self.storage.with_store(|pile, signer, runtime| {
let source = crate::collection_names::open_configured_acquiring(
pile, DEFAULT_SCOPE_ID, signer.verifying_key(), runtime,
)?;
let succinct = pile.attach::<SuccinctArchiveBlob>(source, ())?;
let rank9 = pile.attach::<Rank9AcceleratedSuccinctArchiveBlob>(source, succinct)?;
let target = super::intent_register_collection(pile, signer.verifying_key())?;
let store_snapshot = AcquiringReader::new(pile.snapshot()?, runtime.clone());
let facts = crate::storage::acquire_facts(&store_snapshot, rank9)?;
let intents = store_snapshot.attached_acquiring(target)?
.read_acquiring::<triblespace::core::collection::lww_register::LwwIndex>()?;
if !intents.unread().is_empty() {
bail!("Body intent register has {} unread foundations", intents.unread().len());
}
let intents = intents.into_value().query()?;
f(&super::BodySnapshot { facts, store_snapshot, intents })
})
}
}
#[cfg(test)]
#[path = "tests.rs"]
mod tests;
#[cfg(test)]
mod projection_tests {
use super::*;
use crate::storage::{initialize_signer, load_signer, open_pile_strict_as};
use triblespace::core::collection::lww_register::LwwIndex;
#[test]
fn body_publication_reaches_facts_and_intent_for_a_preparing_read() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("body.pile");
let key = directory.path().join("body.key");
std::fs::File::create(&path).unwrap();
initialize_signer(&path, Some(&key)).unwrap();
let body = Body::new(path.clone(), Some(key.clone()));
let intent = body.set_intent("observe the completed write").unwrap();
let capture = body
.capture(&CaptureInput {
signal: Signal::Touch,
pose: "{}".to_owned(),
note: None,
})
.unwrap();
let intent_id = intent.id;
let capture_id = capture.id;
let signer = load_signer(&path, Some(&key)).unwrap();
let mut pile = open_pile_strict_as(&path, signer.verifying_key()).unwrap();
let source = open_configured(&mut pile, DEFAULT_SCOPE_ID, signer.verifying_key()).unwrap();
let succinct = pile.attach::<SuccinctArchiveBlob>(source, ()).unwrap();
let rank9 = pile
.attach::<Rank9AcceleratedSuccinctArchiveBlob>(source, succinct)
.unwrap();
let intents =
crate::body::intent_register_collection(&mut pile, signer.verifying_key()).unwrap();
let snapshot = pollster::block_on(async {
drop(pile.maintain_attached(succinct, &signer).await?);
drop(pile.maintain_attached(rank9, &signer).await?);
pile.maintain_attached(intents, &signer).await
})
.unwrap();
let facts = snapshot.read_facts(rank9).unwrap();
assert!(exists!(
pattern!(&facts, [{ intent_id @ metadata::tag: &KIND_INTENT }])
));
assert!(exists!(
pattern!(&facts, [{ capture_id @ metadata::tag: &KIND_CAPTURE }])
));
let winners = snapshot
.attached(intents)
.unwrap()
.read::<LwwIndex>()
.unwrap()
.into_value()
.query()
.unwrap();
assert_eq!(winners.winner(KIND_INTENT), Some(intent.id));
pile.close().unwrap();
}
}