use std::{collections::BTreeMap, path::Path, str::FromStr};
use serde::Deserialize;
use super::{plan_relative, read_json, run_relative, WorkflowStore};
use crate::workflow::{
check_schema_version, NodeId, NodeState, PlanId, RunId, RunLifecycle, ScopeInstanceId,
WorkflowOutcome, WorkflowResult, PLAN_MANIFEST_VERSION, RUN_MANIFEST_VERSION,
};
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum RecordAccess {
Executable,
ReadOnly,
}
impl RecordAccess {
pub(super) fn from_version(
kind: &'static str,
version: u32,
current: u32,
) -> WorkflowResult<Self> {
if version == super::legacy::LEGACY_MANIFEST_VERSION {
Ok(Self::ReadOnly)
} else {
check_schema_version(kind, version, current)?;
Ok(Self::Executable)
}
}
}
#[derive(Deserialize)]
struct InventoryManifest {
schema_version: u32,
#[serde(default)]
created_at_unix_nanos: u64,
workspace_identity: String,
#[serde(default)]
name: String,
#[serde(default)]
step_count: usize,
}
#[derive(Deserialize)]
struct PlanSummary {
plan_id: PlanId,
#[serde(flatten)]
metadata: InventoryManifest,
}
#[derive(Deserialize)]
struct RunSummary {
run_id: RunId,
#[serde(flatten)]
metadata: InventoryManifest,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct PlanInventoryItem {
pub(crate) plan_id: PlanId,
pub(crate) access: RecordAccess,
pub(crate) created_at_unix_nanos: u64,
pub(crate) workspace_identity: String,
pub(crate) name: String,
pub(crate) step_count: usize,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct RunInventoryItem {
pub(crate) run_id: RunId,
pub(crate) access: RecordAccess,
pub(crate) created_at_unix_nanos: u64,
pub(crate) workspace_identity: String,
pub(crate) name: String,
pub(crate) lifecycle: RunLifecycle,
pub(crate) outcome: Option<WorkflowOutcome>,
pub(crate) done_steps: usize,
pub(crate) total_steps: usize,
}
#[derive(Debug, Deserialize)]
struct InventoryStateFile {
state: InventoryWorkflowState,
}
#[derive(Debug, Deserialize)]
struct InventoryWorkflowState {
lifecycle: RunLifecycle,
scopes: BTreeMap<ScopeInstanceId, InventoryScopeState>,
}
#[derive(Debug, Deserialize)]
struct InventoryScopeState {
nodes: BTreeMap<NodeId, NodeState>,
result: Option<InventoryScopeResult>,
}
#[derive(Debug, Deserialize)]
struct InventoryScopeResult {
outcome: WorkflowOutcome,
}
#[derive(Debug, Deserialize)]
struct RevisionStateFile {
state: RevisionOnly,
}
#[derive(Debug, Deserialize)]
struct RevisionOnly {
revision: u64,
}
#[derive(Debug, Deserialize)]
struct LifecycleStateFile {
state: LifecycleOnly,
}
#[derive(Debug, Deserialize)]
struct LifecycleOnly {
lifecycle: RunLifecycle,
}
impl WorkflowStore {
pub(crate) fn list_plan_inventory(&self) -> WorkflowResult<Vec<PlanInventoryItem>> {
let mut plans = Vec::new();
for name in self.root.directory_names(Path::new("plans"))? {
let Ok(name) = name.into_string() else {
continue;
};
let Ok(id) = PlanId::from_str(&name) else {
continue;
};
plans.push(self.read_plan_inventory(id)?);
}
plans.sort_by_key(|plan| std::cmp::Reverse((plan.created_at_unix_nanos, plan.plan_id)));
Ok(plans)
}
pub(crate) fn list_run_inventory(&self) -> WorkflowResult<Vec<RunInventoryItem>> {
let mut runs = Vec::new();
for name in self.root.directory_names(Path::new("runs"))? {
let Ok(name) = name.into_string() else {
continue;
};
let Ok(id) = RunId::from_str(&name) else {
continue;
};
runs.push(self.read_run_inventory(id)?);
}
runs.sort_by_key(|run| std::cmp::Reverse((run.created_at_unix_nanos, run.run_id)));
Ok(runs)
}
pub(crate) fn read_run_revision(&self, id: RunId) -> WorkflowResult<u64> {
let state: RevisionStateFile =
read_json(&self.root, &run_relative(id, Path::new("state.json")))?;
Ok(state.state.revision)
}
pub(crate) fn read_run_lifecycle(&self, id: RunId) -> WorkflowResult<RunLifecycle> {
let state: LifecycleStateFile =
read_json(&self.root, &run_relative(id, Path::new("state.json")))?;
Ok(state.state.lifecycle)
}
pub(crate) fn read_run_inventory(&self, id: RunId) -> WorkflowResult<RunInventoryItem> {
let summary: RunSummary =
read_json(&self.root, &run_relative(id, Path::new("manifest.json")))?;
if summary.run_id != id {
return super::corrupt(
&self.layout.run_manifest(id),
"run manifest ID differs from its directory ID",
);
}
let manifest = summary.metadata;
let access = RecordAccess::from_version(
"run manifest",
manifest.schema_version,
RUN_MANIFEST_VERSION,
)?;
if access == RecordAccess::ReadOnly {
return self.read_legacy_run_inventory(id, manifest, access);
}
let state: InventoryStateFile =
read_json(&self.root, &run_relative(id, Path::new("state.json")))?;
let scope = state
.state
.scopes
.get(&ScopeInstanceId::ROOT)
.ok_or_else(|| crate::workflow::WorkflowError::Corrupt {
path: self.layout.run_state(id),
reason: "run has no root scope instance".to_owned(),
})?;
Ok(RunInventoryItem {
run_id: id,
access,
created_at_unix_nanos: manifest.created_at_unix_nanos,
workspace_identity: manifest.workspace_identity,
name: manifest.name,
lifecycle: state.state.lifecycle,
outcome: if state.state.lifecycle == RunLifecycle::Completed {
scope.result.as_ref().map(|result| result.outcome)
} else {
None
},
done_steps: terminal_count(scope.nodes.values()),
total_steps: manifest.step_count,
})
}
fn read_legacy_run_inventory(
&self,
id: RunId,
manifest: InventoryManifest,
access: RecordAccess,
) -> WorkflowResult<RunInventoryItem> {
let state = self.read_legacy_run_state(id)?.state;
Ok(RunInventoryItem {
run_id: id,
access,
created_at_unix_nanos: manifest.created_at_unix_nanos,
workspace_identity: manifest.workspace_identity,
name: manifest.name,
lifecycle: state.lifecycle,
outcome: state.outcome,
done_steps: state
.nodes
.values()
.filter(|node| {
node.get("state").and_then(serde_json::Value::as_str) == Some("terminal")
})
.count(),
total_steps: manifest.step_count,
})
}
pub(crate) fn read_plan_inventory(&self, id: PlanId) -> WorkflowResult<PlanInventoryItem> {
let summary: PlanSummary =
read_json(&self.root, &plan_relative(id, Path::new("manifest.json")))?;
if summary.plan_id != id {
return super::corrupt(
&self.layout.plan_manifest(id),
"plan manifest ID differs from its directory ID",
);
}
let manifest = summary.metadata;
let access = RecordAccess::from_version(
"plan manifest",
manifest.schema_version,
PLAN_MANIFEST_VERSION,
)?;
Ok(PlanInventoryItem {
plan_id: id,
access,
created_at_unix_nanos: manifest.created_at_unix_nanos,
workspace_identity: manifest.workspace_identity,
name: manifest.name,
step_count: manifest.step_count,
})
}
pub(super) fn require_executable_run(&self, id: RunId) -> WorkflowResult<()> {
let summary: RunSummary =
read_json(&self.root, &run_relative(id, Path::new("manifest.json")))?;
if summary.run_id != id {
return super::corrupt(
&self.layout.run_manifest(id),
"run manifest ID differs from its directory ID",
);
}
match RecordAccess::from_version(
"run manifest",
summary.metadata.schema_version,
RUN_MANIFEST_VERSION,
)? {
RecordAccess::Executable => Ok(()),
RecordAccess::ReadOnly => Err(super::legacy::legacy_record("run", id)),
}
}
}
fn terminal_count<'a>(nodes: impl Iterator<Item = &'a NodeState>) -> usize {
nodes.filter(|node| node.terminal().is_some()).count()
}