use crate::bundle::types::{
as_artifact_ref, DefinitionSnapshot, Manifest, RunState, SessionBinding, SessionCapture,
SessionEntryRecord, SessionEventRecord, TraceEvent, RUN_BUNDLE_SCHEMA,
};
use anyhow::{Context, Result};
use serde_json::Value;
use std::path::{Path, PathBuf};
#[derive(Debug, Clone)]
pub struct BundlePaths {
pub dir: PathBuf,
pub workflow: PathBuf,
pub state: PathBuf,
pub trace: PathBuf,
pub session: Option<PathBuf>,
pub artifacts: Option<PathBuf>,
}
impl BundlePaths {
pub fn from_manifest(dir: &Path, manifest: &Manifest) -> Self {
Self {
dir: dir.to_path_buf(),
workflow: dir.join(&manifest.paths.workflow),
state: dir.join(&manifest.paths.state),
trace: dir.join(&manifest.paths.trace),
session: Some(match manifest.paths.session.as_ref() {
Some(p) => dir.join(p),
None => dir.join("session"),
}),
artifacts: manifest.paths.artifacts.as_ref().map(|p| dir.join(p)),
}
}
pub fn session_binding(&self) -> Option<PathBuf> {
self.session.as_ref().map(|dir| dir.join("binding.json"))
}
pub fn session_entries(&self) -> Option<PathBuf> {
self.session.as_ref().map(|dir| dir.join("entries.ndjson"))
}
pub fn session_events(&self) -> Option<PathBuf> {
self.session.as_ref().map(|dir| dir.join("events.ndjson"))
}
pub fn session_capture(&self) -> Option<PathBuf> {
self.session.as_ref().map(|dir| dir.join("capture.json"))
}
}
#[derive(Debug, Clone)]
pub struct LoadedBundle {
pub manifest: Manifest,
pub paths: BundlePaths,
pub state: RunState,
pub snapshot: Option<DefinitionSnapshot>,
pub trace: Vec<TraceEvent>,
pub session_binding: Option<SessionBinding>,
pub session_entries: Vec<SessionEntryRecord>,
pub session_events: Vec<SessionEventRecord>,
pub session_capture: Option<SessionCapture>,
}
fn is_contained(relative: &str) -> bool {
let path = Path::new(relative);
!relative.is_empty()
&& path.is_relative()
&& path
.components()
.all(|component| matches!(component, std::path::Component::Normal(_)))
}
pub fn contained_path(bundle_dir: &Path, path: &Path) -> Option<PathBuf> {
let canonical = path.canonicalize().ok()?;
let base = bundle_dir.canonicalize().ok()?;
canonical.starts_with(&base).then_some(canonical)
}
pub fn read_contained(bundle_dir: &Path, path: &Path) -> Option<String> {
std::fs::read_to_string(contained_path(bundle_dir, path)?).ok()
}
pub fn read_manifest(dir: &Path) -> Result<Manifest> {
read_manifest_value(dir).map(|(_, manifest)| manifest)
}
pub fn read_manifest_value(dir: &Path) -> Result<(Value, Manifest)> {
let path = dir.join("manifest.json");
let raw = read_contained(dir, &path)
.with_context(|| format!("reading {} inside the bundle", path.display()))?;
let raw: Value =
serde_json::from_str(&raw).with_context(|| format!("parsing {}", path.display()))?;
let manifest: Manifest = serde_json::from_value(raw.clone())
.with_context(|| format!("parsing {}", path.display()))?;
anyhow::ensure!(
manifest.schema == RUN_BUNDLE_SCHEMA,
"unsupported bundle schema {:?} in {}",
manifest.schema,
path.display()
);
let entries = [
Some(&manifest.paths.workflow),
Some(&manifest.paths.state),
Some(&manifest.paths.trace),
manifest.paths.session.as_ref(),
manifest.paths.artifacts.as_ref(),
];
for entry in entries.into_iter().flatten() {
anyhow::ensure!(
is_contained(entry),
"manifest path {entry:?} escapes the bundle in {}",
path.display()
);
}
Ok((raw, manifest))
}
fn read_json<T: serde::de::DeserializeOwned>(bundle_dir: &Path, path: &Path) -> Result<T> {
let raw = read_contained(bundle_dir, path)
.with_context(|| format!("reading {} inside the bundle", path.display()))?;
serde_json::from_str(&raw).with_context(|| format!("parsing {}", path.display()))
}
pub fn parse_ndjson<T: serde::de::DeserializeOwned>(raw: &str) -> Vec<T> {
raw.lines()
.filter(|line| !line.trim().is_empty())
.filter_map(|line| serde_json::from_str(line).ok())
.collect()
}
pub fn read_bundle(dir: &Path) -> Result<LoadedBundle> {
let manifest = read_manifest(dir)?;
let paths = BundlePaths::from_manifest(dir, &manifest);
let state: RunState = read_json(dir, &paths.state)?;
let snapshot: Option<DefinitionSnapshot> = read_json(dir, &paths.workflow).ok();
let trace: Vec<TraceEvent> = read_contained(dir, &paths.trace)
.map(|raw| parse_ndjson(&raw))
.unwrap_or_default();
let session_binding = paths
.session_binding()
.and_then(|path| read_json(dir, &path).ok());
let session_entries = paths
.session_entries()
.and_then(|path| read_contained(dir, &path))
.map(|raw| parse_ndjson(&raw))
.unwrap_or_default();
let session_events = paths
.session_events()
.and_then(|path| read_contained(dir, &path))
.map(|raw| parse_ndjson(&raw))
.unwrap_or_default();
let session_capture = paths
.session_capture()
.and_then(|path| read_json(dir, &path).ok());
Ok(LoadedBundle {
manifest,
paths,
state,
snapshot,
trace,
session_binding,
session_entries,
session_events,
session_capture,
})
}
pub fn list_bundles(runs_dir: &Path) -> Vec<(PathBuf, Manifest)> {
let Ok(entries) = std::fs::read_dir(runs_dir) else {
return Vec::new();
};
let mut bundles: Vec<(PathBuf, Manifest)> = entries
.filter_map(|entry| entry.ok())
.filter(|entry| entry.file_type().map(|t| t.is_dir()).unwrap_or(false))
.filter_map(|entry| {
let dir = entry.path();
read_manifest(&dir).ok().map(|manifest| (dir, manifest))
})
.collect();
bundles.sort_by(|a, b| {
b.1.started_at
.cmp(&a.1.started_at)
.then_with(|| b.1.run_id.cmp(&a.1.run_id))
});
bundles
}
pub fn artifact_placeholder(path: &str, bytes: u64) -> String {
let size = if bytes < 1024 {
format!("{bytes}B")
} else {
format!("{:.1}KB", bytes as f64 / 1024.0)
};
format!("«artifact {size} {path}»")
}
pub fn resolve_artifacts(value: &Value, bundle_dir: &Path, max_bytes: u64) -> Value {
match value {
Value::Object(_) => {
if let Some(reference) = as_artifact_ref(value) {
let text = if reference.bytes <= max_bytes {
read_artifact_checked(bundle_dir, &reference.path, max_bytes)
} else {
None
};
return match text {
Some(text) => Value::String(text),
None => Value::String(artifact_placeholder(&reference.path, reference.bytes)),
};
}
if let Some(inner) = crate::bundle::types::as_escaped(value) {
return match inner.as_object() {
Some(object) => Value::Object(
object
.iter()
.map(|(key, value)| {
(key.clone(), resolve_artifacts(value, bundle_dir, max_bytes))
})
.collect(),
),
None => inner.clone(),
};
}
let object = value.as_object().unwrap();
Value::Object(
object
.iter()
.map(|(key, value)| {
(key.clone(), resolve_artifacts(value, bundle_dir, max_bytes))
})
.collect(),
)
}
Value::Array(items) => Value::Array(
items
.iter()
.map(|item| resolve_artifacts(item, bundle_dir, max_bytes))
.collect(),
),
other => other.clone(),
}
}
pub fn read_artifact_checked(bundle_dir: &Path, relative: &str, max_bytes: u64) -> Option<String> {
let canonical = contained_path(bundle_dir, &bundle_dir.join(relative))?;
if std::fs::metadata(&canonical).ok()?.len() > max_bytes {
return None;
}
std::fs::read_to_string(canonical).ok()
}
pub fn read_declared_artifact_checked(
bundle_dir: &Path,
artifact_dir: &str,
relative: &str,
max_bytes: u64,
) -> Option<String> {
let artifact_dir = Path::new(artifact_dir);
let relative = Path::new(relative);
if artifact_dir.as_os_str().is_empty() || !relative.starts_with(artifact_dir) {
return None;
}
let canonical_root = contained_path(bundle_dir, &bundle_dir.join(artifact_dir))?;
let canonical_file = contained_path(bundle_dir, &bundle_dir.join(relative))?;
if !canonical_file.starts_with(&canonical_root)
|| std::fs::metadata(&canonical_file).ok()?.len() > max_bytes
{
return None;
}
std::fs::read_to_string(canonical_file).ok()
}
pub fn with_artifact_placeholders(value: &Value) -> Value {
match value {
Value::Object(_) => {
if let Some(reference) = as_artifact_ref(value) {
return Value::String(artifact_placeholder(&reference.path, reference.bytes));
}
if let Some(inner) = crate::bundle::types::as_escaped(value) {
return match inner.as_object() {
Some(object) => Value::Object(
object
.iter()
.map(|(key, value)| (key.clone(), with_artifact_placeholders(value)))
.collect(),
),
None => inner.clone(),
};
}
let object = value.as_object().unwrap();
Value::Object(
object
.iter()
.map(|(key, value)| (key.clone(), with_artifact_placeholders(value)))
.collect(),
)
}
Value::Array(items) => Value::Array(items.iter().map(with_artifact_placeholders).collect()),
other => other.clone(),
}
}