use std::io::{self, Read, Seek, SeekFrom};
use super::{Frame, Frames, RunRecord, read_record};
use crate::execution::ToolOutcome;
use crate::region::EntryContent;
#[derive(Debug)]
pub struct Execution {
pub id: String,
pub call_id: String,
pub tool: String,
pub arguments: String,
pub stage_index: usize,
pub iteration: usize,
pub visit_id: String,
pub requested_by: String,
pub artifacts: Vec<crate::output::Artifact>,
pub dispatched_at: i64,
pub position: u64,
pub ended_at: Option<i64>,
pub result_position: Option<u64>,
pub outcome: Option<ToolOutcome>,
}
impl Execution {
pub fn unfinished(&self) -> bool {
self.ended_at.is_none()
}
}
pub fn read_archive_executions(r: &mut dyn Read) -> io::Result<Vec<Execution>> {
let (_, mut frames) = Frames::open(r)?;
let mut executions: Vec<Execution> = Vec::new();
while let Ok(Some((position, frame))) = frames.next_frame() {
let Frame::Record(record) = frame else {
continue;
};
match *record {
RunRecord::ToolBatch {
calls,
at,
stage_index,
iteration,
visit_id,
requested_by,
..
} => {
for call in calls {
let inline = call.result.is_some();
executions.push(Execution {
id: call.execution_id,
call_id: call.id,
tool: call.name,
arguments: call.arguments,
stage_index,
iteration,
visit_id: visit_id.clone(),
requested_by: requested_by.clone(),
artifacts: Vec::new(),
dispatched_at: at,
position,
ended_at: inline.then_some(at),
result_position: inline.then_some(position),
outcome: None,
});
}
}
RunRecord::ArtifactsProduced {
execution_id,
artifacts,
..
} => {
if let Some(execution) = executions
.iter_mut()
.find(|e| !e.id.is_empty() && e.id == execution_id)
{
execution.artifacts.extend(artifacts);
}
}
RunRecord::ToolCallDone {
iteration,
call_id,
execution_id,
outcome,
at,
..
} => {
if let Some(execution) =
match_completion(&mut executions, &execution_id, &call_id, iteration)
{
execution.ended_at = Some(at);
execution.result_position = Some(position);
execution.outcome = outcome;
}
}
_ => {}
}
}
Ok(executions)
}
fn match_completion<'e>(
executions: &'e mut [Execution],
execution_id: &str,
call_id: &str,
iteration: usize,
) -> Option<&'e mut Execution> {
if !execution_id.is_empty() {
return executions.iter_mut().find(|e| e.id == execution_id);
}
executions
.iter_mut()
.rev()
.find(|e| e.call_id == call_id && e.iteration == iteration && e.unfinished())
}
pub trait SeekRead: Read + Seek {}
impl<T: Read + Seek> SeekRead for T {}
pub fn read_result_at(
r: &mut dyn SeekRead,
position: u64,
call_id: &str,
) -> io::Result<Option<EntryContent>> {
r.seek(SeekFrom::Start(position))?;
let Some(record) = read_record(r)? else {
return Ok(None);
};
Ok(result_of(record, call_id))
}
fn result_of(record: RunRecord, call_id: &str) -> Option<EntryContent> {
match record {
RunRecord::ToolCallDone { call_id: done, .. } if done != call_id => None,
RunRecord::ToolCallDone { result, .. } => Some(result),
RunRecord::ToolBatch { calls, .. } => calls
.into_iter()
.find(|c| c.id == call_id)
.and_then(|c| c.result),
_ => None,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::run_archive::{
RUN_ARCHIVE_MAGIC, RUN_ARCHIVE_VERSION, RunRecord, ToolCallRecord, write_archive_start,
write_record,
};
fn call(id: &str, execution_id: &str, result: Option<&str>) -> ToolCallRecord {
ToolCallRecord {
execution_id: execution_id.to_string(),
id: id.to_string(),
name: "shell".to_string(),
arguments: r#"{"command":"ls"}"#.to_string(),
result: result.map(Into::into),
thought_signature: None,
}
}
fn archive(records: Vec<RunRecord>) -> Vec<u8> {
let mut bytes = Vec::new();
write_archive_start(&mut bytes, RUN_ARCHIVE_VERSION).expect("a Vec takes the preamble");
for record in &records {
write_record(&mut bytes, record).expect("a Vec takes a record");
}
bytes
}
#[test]
fn a_batch_and_its_completions_read_as_executions() {
let bytes = archive(vec![
RunRecord::ToolBatch {
calls: vec![call("c1", "x1", None), call("c2", "x2", None)],
at: 10,
stage_index: 2,
iteration: 5,
visit_id: String::new(),
requested_by: String::new(),
response: "doing two things".to_string(),
},
RunRecord::ToolCallDone {
iteration: 5,
call_id: "c2".to_string(),
execution_id: "x2".to_string(),
result: "second".into(),
outcome: None,
at: 12,
},
RunRecord::ToolCallDone {
iteration: 5,
call_id: "c1".to_string(),
execution_id: "x1".to_string(),
result: "first".into(),
outcome: Some(ToolOutcome::Succeeded),
at: 14,
},
]);
let executions =
read_archive_executions(&mut bytes.as_slice()).expect("the archive reads back");
assert_eq!(executions.len(), 2);
assert_eq!(executions[0].id, "x1");
assert_eq!(executions[0].call_id, "c1");
assert_eq!(executions[0].tool, "shell");
assert_eq!(executions[0].arguments, r#"{"command":"ls"}"#);
assert_eq!(executions[0].stage_index, 2);
assert_eq!(executions[0].iteration, 5);
assert_eq!(executions[0].dispatched_at, 10);
assert_eq!(executions[0].ended_at, Some(14));
assert_eq!(executions[0].outcome, Some(ToolOutcome::Succeeded));
assert_eq!(executions[1].id, "x2");
assert_eq!(executions[1].ended_at, Some(12));
assert_eq!(executions[0].position, executions[1].position);
assert!(!executions[0].unfinished());
}
#[test]
fn a_call_with_no_completion_stays_unfinished() {
let bytes = archive(vec![RunRecord::ToolBatch {
calls: vec![call("c1", "x1", None)],
at: 10,
stage_index: 0,
iteration: 1,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
}]);
let executions = read_archive_executions(&mut bytes.as_slice()).expect("it reads");
assert_eq!(executions.len(), 1);
assert!(executions[0].unfinished());
assert_eq!(executions[0].ended_at, None);
assert_eq!(executions[0].result_position, None);
assert_eq!(executions[0].outcome, None);
}
#[test]
fn an_inline_result_ends_with_the_batch_that_carried_it() {
let bytes = archive(vec![RunRecord::ToolBatch {
calls: vec![call("c1", "x1", Some("refused"))],
at: 10,
stage_index: 0,
iteration: 1,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
}]);
let executions = read_archive_executions(&mut bytes.as_slice()).expect("it reads");
assert_eq!(executions[0].ended_at, Some(10));
assert_eq!(
executions[0].result_position,
Some(executions[0].position),
"its result is in the record that dispatched it"
);
let mut cursor = std::io::Cursor::new(bytes);
let result = read_result_at(
&mut cursor,
executions[0].result_position.expect("a position"),
"c1",
)
.expect("the record reads");
assert_eq!(result.as_deref(), Some("refused"));
}
#[test]
fn a_result_is_read_from_the_position_the_execution_names() {
let bytes = archive(vec![
RunRecord::ToolBatch {
calls: vec![call("c1", "x1", None)],
at: 10,
stage_index: 0,
iteration: 1,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
},
RunRecord::ToolCallDone {
iteration: 1,
call_id: "c1".to_string(),
execution_id: "x1".to_string(),
result: "the file body".into(),
outcome: Some(ToolOutcome::Succeeded),
at: 11,
},
]);
let executions = read_archive_executions(&mut bytes.as_slice()).expect("it reads");
let position = executions[0].result_position.expect("it ended");
let mut cursor = std::io::Cursor::new(bytes);
assert_eq!(
read_result_at(&mut cursor, position, "c1")
.expect("the record reads")
.as_deref(),
Some("the file body")
);
assert_eq!(
read_result_at(&mut cursor, position, "c9").expect("the record reads"),
None
);
}
#[test]
fn two_attempts_at_one_call_id_stay_apart() {
let bytes = archive(vec![
RunRecord::ToolBatch {
calls: vec![call("c1", "x1", None)],
at: 10,
stage_index: 0,
iteration: 1,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
},
RunRecord::ToolCallDone {
iteration: 1,
call_id: "c1".to_string(),
execution_id: "x1".to_string(),
result: "timed out".into(),
outcome: Some(ToolOutcome::Failed),
at: 11,
},
RunRecord::ToolBatch {
calls: vec![call("c1", "x2", None)],
at: 12,
stage_index: 0,
iteration: 2,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
},
RunRecord::ToolCallDone {
iteration: 2,
call_id: "c1".to_string(),
execution_id: "x2".to_string(),
result: "done".into(),
outcome: Some(ToolOutcome::Succeeded),
at: 13,
},
]);
let executions = read_archive_executions(&mut bytes.as_slice()).expect("it reads");
assert_eq!(executions.len(), 2);
assert_eq!(executions[0].outcome, Some(ToolOutcome::Failed));
assert_eq!(executions[1].outcome, Some(ToolOutcome::Succeeded));
assert_ne!(executions[0].position, executions[1].position);
}
#[test]
fn a_journal_with_no_execution_ids_still_pairs_up() {
let bytes = archive(vec![
RunRecord::ToolBatch {
calls: vec![call("c1", "", None)],
at: 10,
stage_index: 0,
iteration: 1,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
},
RunRecord::ToolBatch {
calls: vec![call("c1", "", None)],
at: 12,
stage_index: 0,
iteration: 2,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
},
RunRecord::ToolCallDone {
iteration: 2,
call_id: "c1".to_string(),
execution_id: String::new(),
result: "second".into(),
outcome: None,
at: 13,
},
]);
let executions = read_archive_executions(&mut bytes.as_slice()).expect("it reads");
assert_eq!(executions.len(), 2);
assert!(
executions[0].unfinished(),
"the first iteration's call never got an ending"
);
assert_eq!(executions[1].ended_at, Some(13));
}
#[test]
fn a_completion_with_no_dispatch_is_dropped() {
let bytes = archive(vec![RunRecord::ToolCallDone {
iteration: 1,
call_id: "ghost".to_string(),
execution_id: "x9".to_string(),
result: "from nowhere".into(),
outcome: None,
at: 11,
}]);
let executions = read_archive_executions(&mut bytes.as_slice()).expect("it reads");
assert!(executions.is_empty(), "{executions:?}");
}
fn artifact(name: &str) -> crate::output::Artifact {
crate::output::Artifact {
name: name.to_string(),
path: format!("out/{name}"),
mime_type: crate::mime::MimeType::parse("text/markdown").expect("a type"),
size: 12,
sha256: "beef".to_string(),
}
}
#[test]
fn the_files_an_execution_produced_land_on_it() {
let bytes = archive(vec![
RunRecord::ToolBatch {
calls: vec![
call("c1", "x1", Some("recorded")),
call("c2", "x2", Some("ok")),
],
at: 10,
stage_index: 0,
iteration: 1,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
},
RunRecord::ArtifactsProduced {
execution_id: "x1".to_string(),
artifacts: vec![artifact("report")],
at: 11,
},
RunRecord::ArtifactsProduced {
execution_id: "x1".to_string(),
artifacts: vec![artifact("chart")],
at: 12,
},
RunRecord::ArtifactsProduced {
execution_id: "x-nothing-dispatched".to_string(),
artifacts: vec![artifact("orphan")],
at: 13,
},
]);
let executions = read_archive_executions(&mut bytes.as_slice()).expect("it reads");
let names: Vec<&str> = executions[0]
.artifacts
.iter()
.map(|a| a.name.as_str())
.collect();
assert_eq!(
names,
vec!["report", "chart"],
"both, in the order recorded"
);
assert!(
executions[1].artifacts.is_empty(),
"the other call produced nothing"
);
}
#[test]
fn an_unidentified_execution_takes_no_files() {
let bytes = archive(vec![
RunRecord::ToolBatch {
calls: vec![call("c1", "", Some("recorded"))],
at: 10,
stage_index: 0,
iteration: 1,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
},
RunRecord::ArtifactsProduced {
execution_id: String::new(),
artifacts: vec![artifact("report")],
at: 11,
},
]);
let executions = read_archive_executions(&mut bytes.as_slice()).expect("it reads");
assert!(executions[0].artifacts.is_empty());
}
#[test]
fn something_that_is_not_an_archive_is_refused() {
let mut bytes = b"not an archive at all".as_slice();
assert!(read_archive_executions(&mut bytes).is_err());
}
#[test]
fn a_torn_tail_keeps_the_executions_before_it() {
let mut bytes = archive(vec![RunRecord::ToolBatch {
calls: vec![call("c1", "x1", None)],
at: 10,
stage_index: 0,
iteration: 1,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
}]);
bytes.extend_from_slice(&[0, 0, 0]);
let executions = read_archive_executions(&mut bytes.as_slice()).expect("it reads");
assert_eq!(executions.len(), 1);
}
#[test]
fn an_unreadable_frame_does_not_shift_later_positions() {
let mut bytes = archive(vec![RunRecord::ToolBatch {
calls: vec![call("c1", "x1", None)],
at: 10,
stage_index: 0,
iteration: 1,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
}]);
let payload = br#"{"FromALaterBuild":{"whatever":1}}"#;
bytes.extend_from_slice(&(payload.len() as u64).to_be_bytes());
bytes.extend_from_slice(payload);
let tail = RunRecord::ToolCallDone {
iteration: 1,
call_id: "c1".to_string(),
execution_id: "x1".to_string(),
result: "after the unknown".into(),
outcome: None,
at: 15,
};
write_record(&mut bytes, &tail).expect("a Vec takes a record");
let executions = read_archive_executions(&mut bytes.as_slice()).expect("it reads");
let position = executions[0]
.result_position
.expect("the completion past the unknown frame still paired");
let mut cursor = std::io::Cursor::new(bytes);
assert_eq!(
read_result_at(&mut cursor, position, "c1")
.expect("the record reads")
.as_deref(),
Some("after the unknown")
);
}
#[test]
fn records_about_anything_else_are_passed_over() {
let bytes = archive(vec![
RunRecord::StatusChanged {
status: crate::run_meta::RunStatus::Running,
at: 1,
},
RunRecord::ToolBatch {
calls: vec![call("c1", "x1", None)],
at: 10,
stage_index: 0,
iteration: 1,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
},
RunRecord::StatusChanged {
status: crate::run_meta::RunStatus::Complete,
at: 20,
},
]);
let executions = read_archive_executions(&mut bytes.as_slice()).expect("it reads");
assert_eq!(executions.len(), 1, "one dispatch, two status changes");
assert_eq!(executions[0].call_id, "c1");
}
struct Broken {
seek_fails: bool,
}
impl std::io::Read for Broken {
fn read(&mut self, _: &mut [u8]) -> io::Result<usize> {
Err(io::Error::other("the file went away mid-read"))
}
}
impl Seek for Broken {
fn seek(&mut self, _: SeekFrom) -> io::Result<u64> {
match self.seek_fails {
true => Err(io::Error::other("the file went away mid-seek")),
false => Ok(0),
}
}
}
#[test]
fn a_result_read_that_breaks_reports_the_failure() {
let failed_seek = read_result_at(&mut Broken { seek_fails: true }, 10, "c1");
assert!(failed_seek.is_err(), "a failed seek is not an empty answer");
let failed_read = read_result_at(&mut Broken { seek_fails: false }, 10, "c1");
assert!(failed_read.is_err(), "a failed read is not an empty answer");
}
#[test]
fn a_position_on_another_kind_of_record_holds_no_result() {
let bytes = archive(vec![RunRecord::StatusChanged {
status: crate::run_meta::RunStatus::Complete,
at: 20,
}]);
let mut cursor = std::io::Cursor::new(bytes);
let position = RUN_ARCHIVE_MAGIC.len() as u64 + 2;
assert_eq!(
read_result_at(&mut cursor, position, "c1").expect("the record reads"),
None
);
}
#[test]
fn a_position_past_the_end_answers_nothing() {
let bytes = archive(Vec::new());
let mut cursor = std::io::Cursor::new(bytes);
assert_eq!(
read_result_at(&mut cursor, 4096, "c1").expect("a clean end of file"),
None
);
}
}