mod artifact;
mod attachment;
pub mod config;
mod dynamic;
mod logging;
mod parameter;
mod run;
mod sealed;
#[cfg(test)]
mod tests;
pub use attachment::{detect_file_media_type, AttachmentTable, DEFAULT_FILE_MEDIA_TYPE};
pub use dynamic::{ExperimentDyn, RunDyn, SealedRunDyn, SolveDyn};
pub use logging::AttachmentLogger;
pub use parameter::{ParameterValue, RunParameterCell};
pub use run::{FailedSolveRecord, FinishedSolveRecord};
pub use sealed::{SealedRun, Solve};
use crate::artifact::local_registry::{LocalRegistry, StoredDescriptor, TempLocalRegistry};
use crate::artifact::{media_types, ImageRef, LocalArtifact};
use anyhow::{ensure, Context, Result};
use oci_spec::image::Descriptor;
use parameter::ParameterSet;
use rmpv::Value as MessagePackValue;
use std::sync::{Mutex, MutexGuard};
use std::{collections::BTreeMap, io::Cursor};
const EXPERIMENT_STATUS_FINISHED: &str = "finished";
const EXPERIMENT_STATUS_DRAFT: &str = "draft";
const EXPERIMENT_STATUS_FAILED: &str = "failed";
const EXPERIMENT_STATUS_INTERRUPTED: &str = "interrupted";
const RUN_PARAMETERS_MEDIA_TYPE: &str = "application/org.ommx.v1.experiment.run-parameters+json";
const EXPERIMENT_CONFIG_MEDIA_TYPE: &str = "application/org.ommx.v1.experiment.config+json";
const RUN_STATUS_FINISHED: &str = "finished";
const RUN_STATUS_FAILED: &str = "failed";
const RUN_STATUS_INTERRUPTED: &str = "interrupted";
const SOLVE_STATUS_FINISHED: &str = "finished";
const SOLVE_STATUS_FAILED: &str = "failed";
const SOLVE_STATUS_INTERRUPTED: &str = "interrupted";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ExperimentStatus {
Finished,
Draft,
Failed,
Interrupted,
}
impl ExperimentStatus {
pub fn as_str(&self) -> &'static str {
match self {
Self::Finished => EXPERIMENT_STATUS_FINISHED,
Self::Draft => EXPERIMENT_STATUS_DRAFT,
Self::Failed => EXPERIMENT_STATUS_FAILED,
Self::Interrupted => EXPERIMENT_STATUS_INTERRUPTED,
}
}
fn from_config(status: &str) -> Result<Self> {
match status {
EXPERIMENT_STATUS_FINISHED => Ok(Self::Finished),
EXPERIMENT_STATUS_DRAFT => Ok(Self::Draft),
EXPERIMENT_STATUS_FAILED => Ok(Self::Failed),
EXPERIMENT_STATUS_INTERRUPTED => Ok(Self::Interrupted),
_ => {
crate::bail!(
"Experiment status is {status}, expected {EXPERIMENT_STATUS_FINISHED}, \
{EXPERIMENT_STATUS_DRAFT}, {EXPERIMENT_STATUS_FAILED}, or \
{EXPERIMENT_STATUS_INTERRUPTED}"
)
}
}
}
}
impl std::fmt::Display for ExperimentStatus {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RunStatus {
Finished,
Failed,
Interrupted,
}
impl RunStatus {
pub fn as_str(&self) -> &'static str {
match self {
Self::Finished => RUN_STATUS_FINISHED,
Self::Failed => RUN_STATUS_FAILED,
Self::Interrupted => RUN_STATUS_INTERRUPTED,
}
}
fn from_config(status: &str) -> Result<Self> {
match status {
RUN_STATUS_FINISHED => Ok(Self::Finished),
RUN_STATUS_FAILED => Ok(Self::Failed),
RUN_STATUS_INTERRUPTED => Ok(Self::Interrupted),
_ => {
crate::bail!(
"Run status is {status}, expected {RUN_STATUS_FINISHED}, \
{RUN_STATUS_FAILED}, or {RUN_STATUS_INTERRUPTED}"
)
}
}
}
}
impl std::fmt::Display for RunStatus {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SolveStatus {
Finished,
Failed,
Interrupted,
}
impl SolveStatus {
pub fn as_str(&self) -> &'static str {
match self {
Self::Finished => SOLVE_STATUS_FINISHED,
Self::Failed => SOLVE_STATUS_FAILED,
Self::Interrupted => SOLVE_STATUS_INTERRUPTED,
}
}
fn from_config(status: &str) -> Result<Self> {
match status {
SOLVE_STATUS_FINISHED => Ok(Self::Finished),
SOLVE_STATUS_FAILED => Ok(Self::Failed),
SOLVE_STATUS_INTERRUPTED => Ok(Self::Interrupted),
_ => {
crate::bail!(
"Solve status is {status}, expected {SOLVE_STATUS_FINISHED}, \
{SOLVE_STATUS_FAILED}, or {SOLVE_STATUS_INTERRUPTED}"
)
}
}
}
}
impl std::fmt::Display for SolveStatus {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug)]
pub struct Experiment<'reg> {
registry: &'reg LocalRegistry,
state: Mutex<UnsealedExperimentState<'reg>>,
}
#[derive(Debug, Clone)]
pub struct SealedExperiment<'reg> {
status: ExperimentStatus,
artifact: LocalArtifact<'reg>,
attachments: AttachmentTable<StoredDescriptor<'reg>>,
runs: BTreeMap<u64, sealed::SealedRun<'reg>>,
run_parameters: parameter::RunParameterTable,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Trace {
bytes: Vec<u8>,
}
impl Trace {
pub fn from_bytes(bytes: impl Into<Vec<u8>>) -> Self {
Self {
bytes: bytes.into(),
}
}
pub fn as_bytes(&self) -> &[u8] {
&self.bytes
}
pub fn into_bytes(self) -> Vec<u8> {
self.bytes
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Name {
Named(ImageRef),
Anonymous,
}
impl Name {
fn resolve(self, registry: &LocalRegistry) -> Result<ImageRef> {
match self {
Self::Named(image_name) => Ok(image_name),
Self::Anonymous => registry.synthesize_anonymous_experiment_image_name(),
}
}
}
impl From<ImageRef> for Name {
fn from(image_name: ImageRef) -> Self {
Self::Named(image_name)
}
}
#[derive(Debug)]
pub struct Run<'exp, 'reg> {
experiment: &'exp Experiment<'reg>,
run_id: u64,
attachments: AttachmentTable<StoredDescriptor<'reg>>,
trace: Option<StoredDescriptor<'reg>>,
solves: Vec<SolveEntry<'reg>>,
next_solve_id: u64,
parameters: ParameterSet,
}
#[derive(Debug)]
struct RunEntry<'reg> {
run_id: u64,
status: RunStatus,
attachments: AttachmentTable<StoredDescriptor<'reg>>,
trace: Option<StoredDescriptor<'reg>>,
solves: Vec<SolveEntry<'reg>>,
parameters: ParameterSet,
}
#[derive(Debug, Clone)]
struct SolveEntry<'reg> {
solve_id: u64,
status: SolveStatus,
input: StoredDescriptor<'reg>,
output: Option<StoredDescriptor<'reg>>,
adapter: String,
adapter_options: String,
diagnostics: Option<StoredDescriptor<'reg>>,
}
#[derive(Debug, Clone)]
pub struct SolveDiagnosticPayload {
value: MessagePackValue,
}
impl SolveDiagnosticPayload {
pub fn new(bytes: Vec<u8>) -> Result<Self> {
let mut cursor = Cursor::new(&bytes);
let value = rmpv::decode::read_value(&mut cursor)
.context("Solve diagnostic payload must be valid MessagePack")?;
ensure!(
cursor.position() == bytes.len() as u64,
"Solve diagnostic payload must contain exactly one MessagePack value",
);
Self::from_value(value)
}
pub fn from_value(value: MessagePackValue) -> Result<Self> {
ensure!(
matches!(value, MessagePackValue::Array(_)),
"Solve diagnostic payload must decode to a MessagePack array",
);
Ok(Self { value })
}
pub fn value(&self) -> &MessagePackValue {
&self.value
}
pub(crate) fn to_msgpack_bytes(&self) -> Result<Vec<u8>> {
let mut bytes = Vec::new();
rmpv::encode::write_value(&mut bytes, &self.value)
.context("Failed to encode Solve diagnostic payload as MessagePack")?;
Ok(bytes)
}
}
fn read_solve_diagnostic_payload(
solve_id: u64,
descriptor: &StoredDescriptor<'_>,
) -> Result<(Vec<u8>, SolveDiagnosticPayload)> {
descriptor.ensure_media_type(&media_types::diagnostic_msgpack())?;
let bytes = descriptor.registry().get_blob(descriptor)?;
let payload = SolveDiagnosticPayload::new(bytes.clone())
.with_context(|| format!("Invalid Solve {solve_id} diagnostic payload"))?;
Ok((bytes, payload))
}
#[derive(Debug)]
struct UnsealedExperimentState<'reg> {
image_name: ImageRef,
subject: Option<oci_spec::image::Descriptor>,
attachments: AttachmentTable<StoredDescriptor<'reg>>,
runs: BTreeMap<u64, RunEntry<'reg>>,
next_run_id: u64,
}
impl Experiment<'static> {
pub fn new(name: impl Into<Name>) -> Result<Self> {
let registry = LocalRegistry::shared_default()?;
Self::with_registry(registry, name)
}
}
impl<'reg> Experiment<'reg> {
pub fn with_temp_local_registry<T>(
name: impl Into<Name>,
f: impl FnOnce(Experiment<'_>) -> anyhow::Result<T>,
) -> Result<T> {
let temp = TempLocalRegistry::new()?;
let experiment = Experiment::with_registry(temp.registry(), name)?;
f(experiment)
}
pub fn with_registry(registry: &'reg LocalRegistry, name: impl Into<Name>) -> Result<Self> {
let image_name = name.into().resolve(registry)?;
Ok(Experiment {
registry,
state: Mutex::new(UnsealedExperimentState {
image_name,
subject: None,
attachments: AttachmentTable::new(),
runs: BTreeMap::new(),
next_run_id: 0,
}),
})
}
pub fn image_name(&self) -> ImageRef {
self.lock_state().image_name.clone()
}
pub fn run(&self) -> Result<Run<'_, 'reg>> {
let mut state = self.lock_state();
let run_id = allocate_next_run_id(&mut state.next_run_id)?;
Ok(Run {
experiment: self,
run_id,
attachments: AttachmentTable::new(),
trace: None,
solves: Vec::new(),
next_solve_id: 0,
parameters: ParameterSet::new(),
})
}
fn push_closed_run(&self, run: RunEntry<'reg>) -> Result<()> {
let mut state = self.lock_state();
if state.runs.contains_key(&run.run_id) {
crate::bail!("Run {} has already been registered", run.run_id);
}
state.runs.insert(run.run_id, run);
if let Err(error) = state.autosave_checkpoint(self.registry) {
tracing::warn!(
error = %error,
"Failed to publish Experiment autosave checkpoint after Run close"
);
}
Ok(())
}
fn lock_state(&self) -> MutexGuard<'_, UnsealedExperimentState<'reg>> {
match self.state.lock() {
Ok(state) => state,
Err(poisoned) => {
tracing::warn!("Experiment state mutex was poisoned; continuing with inner state");
poisoned.into_inner()
}
}
}
pub fn commit(self) -> Result<SealedExperiment<'reg>> {
let state = match self.state.into_inner() {
Ok(state) => state,
Err(poisoned) => {
tracing::warn!("Experiment state mutex was poisoned; committing inner state");
poisoned.into_inner()
}
};
let artifact = state.commit(self.registry)?;
SealedExperiment::from_artifact(artifact)
}
}
impl<'reg> logging::AttachmentLoggerStorage for &Experiment<'reg> {
type Descriptor = StoredDescriptor<'reg>;
fn with_local_registry<R>(&self, f: impl FnOnce(&LocalRegistry) -> Result<R>) -> Result<R> {
f(self.registry)
}
fn with_attachment_table<R>(
&mut self,
f: impl FnOnce(&mut AttachmentTable<Self::Descriptor>) -> Result<R>,
) -> Result<R> {
let mut state = self.lock_state();
f(&mut state.attachments)
}
fn descriptor_for_attachment_table(&self, descriptor: Descriptor) -> Result<Self::Descriptor> {
self.registry.stored_descriptor(descriptor)
}
}
impl<'reg> SealedExperiment<'reg> {
pub fn artifact(&self) -> LocalArtifact<'reg> {
self.artifact.clone()
}
pub fn into_artifact(self) -> LocalArtifact<'reg> {
self.artifact
}
pub fn fork(&self, name: impl Into<Name>) -> Result<Experiment<'reg>> {
let registry = self.artifact.registry();
let image_name = name.into().resolve(registry)?;
let subject = Some(self.artifact.stored_manifest_descriptor()?.into());
let mut runs = BTreeMap::new();
let mut parameters_by_run = self.run_parameters.parameter_sets()?;
for run in self.runs.values() {
let parameters = parameters_by_run
.remove(&run.run_id())
.unwrap_or_else(ParameterSet::new);
let solves = run
.solves()
.iter()
.map(|solve| SolveEntry {
solve_id: solve.solve_id(),
status: solve.status().clone(),
input: solve.input_descriptor().clone(),
output: solve.output_descriptor().cloned(),
adapter: solve.adapter().to_string(),
adapter_options: solve.adapter_options().to_string(),
diagnostics: solve.diagnostic_descriptor().cloned(),
})
.collect();
runs.insert(
run.run_id(),
RunEntry {
run_id: run.run_id(),
status: run.status().clone(),
attachments: run.attachment_table().clone(),
trace: run.trace_descriptor().cloned(),
solves,
parameters,
},
);
}
Ok(Experiment {
registry,
state: Mutex::new(UnsealedExperimentState {
image_name,
subject,
attachments: self.attachments.clone(),
next_run_id: next_run_id(runs.keys().copied())?,
runs,
}),
})
}
}
fn next_run_id(run_ids: impl Iterator<Item = u64>) -> Result<u64> {
match run_ids.max() {
Some(max) => max
.checked_add(1)
.ok_or_else(|| anyhow::anyhow!("Run ID space is exhausted")),
None => Ok(0),
}
}
fn allocate_next_run_id(next_run_id: &mut u64) -> Result<u64> {
let run_id = *next_run_id;
*next_run_id = next_run_id
.checked_add(1)
.ok_or_else(|| anyhow::anyhow!("Run ID space is exhausted"))?;
Ok(run_id)
}