use std::sync::Arc;
use crate::core::{ObservedStep, RunId, RuntimeError};
use crate::journal::{Append, JournalStore, RecordKind};
#[cfg(feature = "acp")]
pub mod acp;
const OBSERVED_EPOCH: crate::core::Epoch = crate::runtime::LEASE_FREE_EPOCH;
#[derive(Debug)]
pub struct Session {
store: Arc<dyn JournalStore>,
id: String,
run: RunId,
case: Option<crate::core::CaseId>,
steps: usize,
}
impl Session {
#[must_use]
pub fn new(store: Arc<dyn JournalStore>, session: impl Into<String>) -> Self {
Self {
store,
id: session.into(),
run: RunId::generate(),
case: None,
steps: 0,
}
}
#[must_use]
pub const fn in_case(mut self, case: crate::core::CaseId) -> Self {
self.case = Some(case);
self
}
#[must_use]
pub fn source(&self) -> crate::core::SourceId {
crate::core::SourceId::new(format!("observed:{}", self.id))
}
#[must_use]
pub const fn run(&self) -> RunId {
self.run
}
pub async fn step(
&mut self,
step: ObservedStep,
detail: Option<String>,
) -> Result<(), RuntimeError> {
let mut entry = Append::new(
self.run,
RecordKind::Observed {
session: self.id.clone(),
reported: step,
detail,
},
);
if let Some(case) = self.case {
entry = entry.case(case);
}
self.store
.append(OBSERVED_EPOCH, vec![entry])
.await
.map_err(RuntimeError::from_store)?;
self.steps += 1;
Ok(())
}
pub async fn seal(self) -> Result<Option<RunId>, RuntimeError> {
if self.steps == 0 {
return Ok(None);
}
let head = self
.store
.head(self.run)
.await
.map_err(RuntimeError::from_store)?;
let mut sealed = Append::new(
self.run,
RecordKind::RunConcluded {
outcome: crate::runtime::OBSERVED_OUTCOME.to_owned(),
reason: None,
exhaustion: None,
live_spend: crate::core::Spend::default(),
chain_head: head.hash,
},
);
if let Some(case) = self.case {
sealed = sealed.case(case);
}
self.store
.append(OBSERVED_EPOCH, vec![sealed])
.await
.map_err(RuntimeError::from_store)?;
self.store
.seal(self.run, OBSERVED_EPOCH, crate::runtime::OBSERVED_OUTCOME)
.await
.map_err(RuntimeError::from_store)?;
Ok(Some(self.run))
}
}