use assay_runner_core::RunSpec;
use assay_runner_schema::SDK_EVENT_SCHEMA;
use clap::{Args, Subcommand};
use std::fs::File;
use std::path::PathBuf;
#[derive(Debug, Clone, Args)]
pub struct RunnerSpikeArgs {
#[command(subcommand)]
pub cmd: RunnerSpikeCommand,
}
#[derive(Debug, Clone, Subcommand)]
pub enum RunnerSpikeCommand {
Run(RunnerSpikeRunArgs),
}
#[derive(Debug, Clone, Args)]
pub struct RunnerSpikeRunArgs {
#[arg(long, default_value = "none")]
pub agent_shim: String,
#[arg(long)]
pub run_id: Option<String>,
#[arg(long, short = 'o')]
pub output: Option<PathBuf>,
#[arg(long, hide = true)]
pub kernel_capture: bool,
#[arg(long, hide = true)]
pub ebpf: Option<PathBuf>,
#[arg(long, hide = true, default_value_t = 100)]
pub kernel_drain_ms: u64,
#[arg(long, hide = true)]
pub policy_decision_log: Option<PathBuf>,
#[arg(long, hide = true)]
pub sdk_event_log: Option<PathBuf>,
#[arg(long, hide = true)]
pub phase_timing_log: Option<PathBuf>,
#[arg(allow_hyphen_values = true, required = true, trailing_var_arg = true)]
pub command: Vec<String>,
}
pub async fn run(args: RunnerSpikeArgs) -> anyhow::Result<i32> {
match args.cmd {
RunnerSpikeCommand::Run(args) => cmd_run(args).await,
}
}
async fn cmd_run(args: RunnerSpikeRunArgs) -> anyhow::Result<i32> {
validate_runner_spike_args(&args)?;
if args.kernel_capture {
return cmd_run_with_kernel_capture(args).await;
}
cmd_run_contract_only(args)
}
fn validate_runner_spike_args(args: &RunnerSpikeRunArgs) -> anyhow::Result<()> {
if args.sdk_event_log.is_some() && args.agent_shim == "none" {
anyhow::bail!("runner-spike --sdk-event-log requires an SDK agent shim");
}
Ok(())
}
fn build_spec(args: &RunnerSpikeRunArgs) -> RunSpec {
let mut spec = RunSpec::new(args.command.clone()).with_agent_shim(args.agent_shim.clone());
if let Some(run_id) = &args.run_id {
spec = spec.with_run_id(run_id.clone());
}
if let Some(path) = &args.sdk_event_log {
let run_id = spec.run_id.clone();
spec = spec
.with_env("ASSAY_RUNNER_SDK_EVENT_LOG", path.display().to_string())
.with_env("ASSAY_RUNNER_RUN_ID", run_id)
.with_env("ASSAY_RUNNER_SDK_EVENT_SCHEMA", SDK_EVENT_SCHEMA);
}
spec
}
fn bundle_output_path(args: &RunnerSpikeRunArgs, run_id: &str) -> PathBuf {
args.output
.clone()
.unwrap_or_else(|| PathBuf::from(format!("assay-runner-spike-{run_id}.tar.gz")))
}
fn cmd_run_contract_only(args: RunnerSpikeRunArgs) -> anyhow::Result<i32> {
let spec = build_spec(&args);
let output = bundle_output_path(&args, &spec.run_id);
let mut outcome = spec.run_contract_only()?;
apply_policy_then_sdk_logs_if_requested(&spec, &args, &mut outcome.archive)?;
let mut file = File::create(&output)?;
outcome.archive.write(&mut file)?;
let exit_status = exit_status_label(outcome.exit_code, outcome.signal);
println!(
"wrote runner-spike bundle: {} (run_id={}, status={})",
output.display(),
spec.run_id,
exit_status
);
Ok(exit_status_code(outcome.exit_code, outcome.signal))
}
#[cfg(not(target_os = "linux"))]
async fn cmd_run_with_kernel_capture(_args: RunnerSpikeRunArgs) -> anyhow::Result<i32> {
eprintln!("Error: runner-spike --kernel-capture is only supported on Linux.");
Ok(40)
}
#[cfg(target_os = "linux")]
async fn cmd_run_with_kernel_capture(args: RunnerSpikeRunArgs) -> anyhow::Result<i32> {
use assay_monitor::Monitor;
use assay_runner_core::KernelLayerBuilder;
use assay_runner_linux::CgroupManager;
use assay_runner_schema::CgroupCorrelationStatus;
use std::collections::BTreeMap;
use std::time::{Duration, Instant};
use tokio_stream::StreamExt;
let total_start = Instant::now();
let mut phases = BTreeMap::new();
let phase_log = args.phase_timing_log.clone();
let spec = build_spec(&args);
spec.validate()?;
let output = bundle_output_path(&args, &spec.run_id);
let ebpf_path = args
.ebpf
.clone()
.unwrap_or_else(|| PathBuf::from("target/assay-ebpf.o"));
if !ebpf_path.exists() {
eprintln!(
"Error: eBPF object not found at {}. Build it with 'cargo xtask build-ebpf' or provide --ebpf <path>.",
ebpf_path.display()
);
record_phase(&mut phases, "preflight_ms", total_start);
write_phase_timing_log(
phase_log.as_ref(),
&spec,
&phases,
Some(40),
None,
Some("ebpf_object_missing"),
)?;
return Ok(40);
}
record_phase(&mut phases, "preflight_ms", total_start);
let monitor_start = Instant::now();
let mut monitor = match Monitor::load_file(&ebpf_path) {
Ok(monitor) => monitor,
Err(error) => {
eprintln!("Failed to load eBPF: {error}");
record_phase(&mut phases, "monitor_attach_ms", monitor_start);
write_phase_timing_log(
phase_log.as_ref(),
&spec,
&phases,
Some(40),
None,
Some("ebpf_load_failed"),
)?;
return Ok(40);
}
};
if let Err(error) = monitor.configure_defaults() {
eprintln!("Failed to configure eBPF defaults: {error}");
record_phase(&mut phases, "monitor_attach_ms", monitor_start);
write_phase_timing_log(
phase_log.as_ref(),
&spec,
&phases,
Some(40),
None,
Some("ebpf_configure_failed"),
)?;
return Ok(40);
}
if let Err(error) = monitor.set_emit_inode_resolved(false) {
eprintln!("Failed to disable runner-spike inode telemetry: {error}");
record_phase(&mut phases, "monitor_attach_ms", monitor_start);
write_phase_timing_log(
phase_log.as_ref(),
&spec,
&phases,
Some(40),
None,
Some("ebpf_inode_telemetry_config_failed"),
)?;
return Ok(40);
}
if let Err(error) = monitor.set_dedup_open_paths(true) {
eprintln!("Failed to enable runner-spike open path dedupe: {error}");
record_phase(&mut phases, "monitor_attach_ms", monitor_start);
write_phase_timing_log(
phase_log.as_ref(),
&spec,
&phases,
Some(40),
None,
Some("ebpf_open_path_dedupe_config_failed"),
)?;
return Ok(40);
}
if let Err(error) = monitor.attach() {
eprintln!("Failed to attach eBPF probes: {error}");
record_phase(&mut phases, "monitor_attach_ms", monitor_start);
write_phase_timing_log(
phase_log.as_ref(),
&spec,
&phases,
Some(40),
None,
Some("ebpf_attach_failed"),
)?;
return Ok(40);
}
record_phase(&mut phases, "monitor_attach_ms", monitor_start);
let cgroup_start = Instant::now();
let cgroup_manager = match CgroupManager::new() {
Ok(manager) => manager,
Err(error) => {
eprintln!("Failed to initialize runner cgroup manager: {error}");
record_phase(&mut phases, "cgroup_prepare_ms", cgroup_start);
write_phase_timing_log(
phase_log.as_ref(),
&spec,
&phases,
Some(40),
None,
Some("cgroup_manager_init_failed"),
)?;
return Ok(40);
}
};
let session_cgroup = match cgroup_manager.create_session() {
Ok(cgroup) => cgroup,
Err(error) => {
eprintln!("Failed to create runner cgroup session: {error}");
record_phase(&mut phases, "cgroup_prepare_ms", cgroup_start);
write_phase_timing_log(
phase_log.as_ref(),
&spec,
&phases,
Some(40),
None,
Some("cgroup_session_create_failed"),
)?;
return Ok(40);
}
};
if let Err(error) = monitor.set_monitored_cgroups(&[session_cgroup.id()]) {
eprintln!("Failed to populate runner cgroup map: {error}");
record_phase(&mut phases, "cgroup_prepare_ms", cgroup_start);
write_phase_timing_log(
phase_log.as_ref(),
&spec,
&phases,
Some(40),
None,
Some("cgroup_monitor_map_failed"),
)?;
return Ok(40);
}
record_phase(&mut phases, "cgroup_prepare_ms", cgroup_start);
let before_stats = monitor.snapshot_stats()?;
let mut stream = monitor.listen()?;
let mut builder = KernelLayerBuilder::new(&spec.run_id)?;
let mut archive = spec.skeleton_archive()?;
let clock = Instant::now();
spec.append_run_started(&mut archive, 0, Duration::ZERO)?;
let child_spawn_start = Instant::now();
let mut child = match spawn_child_in_cgroup(&spec, &session_cgroup) {
Ok(child) => {
record_phase(&mut phases, "child_spawn_ms", child_spawn_start);
child
}
Err(error) => {
record_phase(&mut phases, "child_spawn_ms", child_spawn_start);
write_phase_timing_log(
phase_log.as_ref(),
&spec,
&phases,
None,
None,
Some("child_spawn_failed"),
)?;
return Err(error);
}
};
let mut cgroup_correlation = CgroupCorrelationStatus::Clean;
let child_runtime_start = Instant::now();
let status = loop {
tokio::select! {
status = child.wait() => break status?,
event = stream.next() => {
match event {
Some(Ok(event)) => builder.push_monitor_event(&event)?,
Some(Err(error)) => {
eprintln!("Warning: failed to parse kernel event: {error}");
cgroup_correlation = CgroupCorrelationStatus::Partial;
}
None => {
eprintln!("Warning: kernel event stream closed before child exit.");
cgroup_correlation = CgroupCorrelationStatus::Partial;
break child.wait().await?;
}
}
}
}
};
record_phase(&mut phases, "child_runtime_ms", child_runtime_start);
let event_flush_start = Instant::now();
let drain_complete = drain_kernel_events(
&mut stream,
&mut builder,
Duration::from_millis(args.kernel_drain_ms),
)
.await?;
if !drain_complete {
cgroup_correlation = CgroupCorrelationStatus::Partial;
}
drop(stream);
let after_stats = monitor.snapshot_stats()?;
let capture = builder.finish(&before_stats, &after_stats);
capture.apply_to_archive(&mut archive, cgroup_correlation)?;
apply_policy_then_sdk_logs_if_requested(&spec, &args, &mut archive)?;
spec.append_run_finished(&mut archive, 1, &status, clock.elapsed())?;
record_phase(&mut phases, "event_flush_ms", event_flush_start);
let archive_write_start = Instant::now();
let mut file = File::create(&output)?;
archive.write(&mut file)?;
record_phase(&mut phases, "archive_write_ms", archive_write_start);
let exit_code = status.code();
let signal = exit_signal(&status);
let exit_status = exit_status_label(exit_code, signal);
write_phase_timing_log(phase_log.as_ref(), &spec, &phases, exit_code, signal, None)?;
println!(
"wrote runner-spike bundle: {} (run_id={}, status={}, kernel_capture={})",
output.display(),
spec.run_id,
exit_status,
cgroup_correlation_label(cgroup_correlation)
);
Ok(exit_status_code(exit_code, signal))
}
#[cfg(target_os = "linux")]
fn record_phase(
phases: &mut std::collections::BTreeMap<&'static str, f64>,
name: &'static str,
start: std::time::Instant,
) {
phases.insert(name, start.elapsed().as_secs_f64() * 1000.0);
}
#[cfg(target_os = "linux")]
fn write_phase_timing_log(
path: Option<&PathBuf>,
spec: &RunSpec,
phases: &std::collections::BTreeMap<&'static str, f64>,
exit_code: Option<i32>,
signal: Option<i32>,
error: Option<&str>,
) -> anyhow::Result<()> {
let Some(path) = path else {
return Ok(());
};
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let payload = serde_json::json!({
"schema": "assay.experiment.runner_phase_timing.v0",
"run_id": &spec.run_id,
"agent_shim": &spec.agent_shim,
"phases_ms": phases,
"exit_code": exit_code,
"signal": signal,
"error": error,
});
std::fs::write(path, serde_json::to_vec_pretty(&payload)?)?;
Ok(())
}
#[cfg(target_os = "linux")]
fn spawn_child_in_cgroup(
spec: &RunSpec,
cgroup: &assay_runner_linux::SessionCgroup,
) -> anyhow::Result<tokio::process::Child> {
use std::ffi::CString;
use std::os::unix::ffi::OsStrExt;
let procs_path = CString::new(cgroup.procs_path().as_os_str().as_bytes())?;
let mut command = tokio::process::Command::new(&spec.command[0]);
command.args(&spec.command[1..]);
apply_kernel_capture_child_env(&mut command, spec);
unsafe {
command.pre_exec(move || write_self_to_cgroup(&procs_path));
}
command
.spawn()
.map_err(|error| anyhow::anyhow!("failed to spawn child in runner cgroup: {error}"))
}
#[cfg(target_os = "linux")]
fn apply_kernel_capture_child_env(command: &mut tokio::process::Command, spec: &RunSpec) {
for key in [
"LD_AUDIT",
"LD_LIBRARY_PATH",
"LD_PRELOAD",
"LOCPATH",
"GCONV_PATH",
] {
command.env_remove(key);
}
command.env("LC_ALL", "C");
command.env("LANG", "C");
command.envs(&spec.env);
}
#[cfg(target_os = "linux")]
async fn drain_kernel_events(
stream: &mut assay_monitor::EventStream,
builder: &mut assay_runner_core::KernelLayerBuilder,
duration: std::time::Duration,
) -> anyhow::Result<bool> {
use tokio_stream::StreamExt;
let mut complete = true;
let deadline = tokio::time::sleep(duration);
tokio::pin!(deadline);
loop {
tokio::select! {
_ = &mut deadline => break,
event = stream.next() => {
match event {
Some(Ok(event)) => builder.push_monitor_event(&event)?,
Some(Err(error)) => {
eprintln!("Warning: failed to parse kernel event while draining: {error}");
complete = false;
}
None => {
complete = false;
break;
}
}
}
}
}
Ok(complete)
}
#[cfg(target_os = "linux")]
fn write_self_to_cgroup(procs_path: &std::ffi::CStr) -> std::io::Result<()> {
let fd = retry_open_write_only(procs_path)?;
let pid = unsafe { libc::getpid() } as u32;
let mut buf = [0_u8; 32];
let len = write_u32_decimal(pid, &mut buf);
let write_result = retry_write_all(fd, &buf[..len]);
let close_result = unsafe { libc::close(fd) };
match (write_result.err(), close_result) {
(Some(error), _) => Err(error),
(None, -1) => Err(std::io::Error::last_os_error()),
(None, _) => Ok(()),
}
}
#[cfg(target_os = "linux")]
fn retry_open_write_only(path: &std::ffi::CStr) -> std::io::Result<i32> {
loop {
let fd = unsafe { libc::open(path.as_ptr(), libc::O_WRONLY | libc::O_CLOEXEC) };
if fd >= 0 {
return Ok(fd);
}
let error = std::io::Error::last_os_error();
if error.raw_os_error() != Some(libc::EINTR) {
return Err(error);
}
}
}
#[cfg(target_os = "linux")]
fn retry_write_all(fd: i32, mut bytes: &[u8]) -> std::io::Result<()> {
while !bytes.is_empty() {
let written = unsafe { libc::write(fd, bytes.as_ptr().cast(), bytes.len()) };
if written < 0 {
let error = std::io::Error::last_os_error();
if error.raw_os_error() == Some(libc::EINTR) {
continue;
}
return Err(error);
}
if written == 0 {
return Err(std::io::Error::from_raw_os_error(libc::EIO));
}
bytes = &bytes[written as usize..];
}
Ok(())
}
#[cfg(target_os = "linux")]
fn write_u32_decimal(value: u32, buf: &mut [u8; 32]) -> usize {
let mut n = value;
if n == 0 {
buf[0] = b'0';
return 1;
}
let mut scratch = [0_u8; 10];
let mut len = 0;
while n > 0 {
scratch[len] = b'0' + (n % 10) as u8;
n /= 10;
len += 1;
}
for idx in 0..len {
buf[idx] = scratch[len - idx - 1];
}
len
}
fn apply_policy_then_sdk_logs_if_requested(
spec: &RunSpec,
args: &RunnerSpikeRunArgs,
archive: &mut assay_runner_core::RunnerSpikeArchive,
) -> anyhow::Result<()> {
apply_policy_decision_log_if_requested(spec, args, archive)?;
apply_sdk_event_log_if_requested(spec, args, archive)?;
Ok(())
}
fn apply_policy_decision_log_if_requested(
spec: &RunSpec,
args: &RunnerSpikeRunArgs,
archive: &mut assay_runner_core::RunnerSpikeArchive,
) -> anyhow::Result<()> {
let Some(path) = args.policy_decision_log.as_ref() else {
return Ok(());
};
let bytes = std::fs::read(path).map_err(|error| {
anyhow::anyhow!(
"failed to read runner-spike policy decision log {}: {error}",
path.display()
)
})?;
let capture =
assay_runner_core::PolicyLayerCapture::from_decision_ndjson(spec.run_id.clone(), &bytes)?;
capture.apply_to_archive(archive)?;
Ok(())
}
fn apply_sdk_event_log_if_requested(
spec: &RunSpec,
args: &RunnerSpikeRunArgs,
archive: &mut assay_runner_core::RunnerSpikeArchive,
) -> anyhow::Result<()> {
let Some(path) = args.sdk_event_log.as_ref() else {
return Ok(());
};
let bytes = std::fs::read(path).map_err(|error| {
anyhow::anyhow!(
"failed to read runner-spike SDK event log {}: {error}",
path.display()
)
})?;
let capture = assay_runner_core::SdkLayerCapture::from_sdk_ndjson(spec.run_id.clone(), &bytes)?;
capture.apply_to_archive(archive)?;
Ok(())
}
#[cfg(target_os = "linux")]
fn cgroup_correlation_label(status: assay_runner_schema::CgroupCorrelationStatus) -> &'static str {
use assay_runner_schema::CgroupCorrelationStatus;
match status {
CgroupCorrelationStatus::Clean => "clean",
CgroupCorrelationStatus::Partial => "partial",
CgroupCorrelationStatus::Failed => "failed",
}
}
#[cfg(target_os = "linux")]
fn exit_signal(status: &std::process::ExitStatus) -> Option<i32> {
use std::os::unix::process::ExitStatusExt;
status.signal()
}
fn exit_status_label(exit_code: Option<i32>, signal: Option<i32>) -> String {
match (exit_code, signal) {
(Some(code), _) => format!("exit_code:{code}"),
(None, Some(signal)) => format!("signal:{signal}"),
(None, None) => "unknown".to_string(),
}
}
fn exit_status_code(exit_code: Option<i32>, signal: Option<i32>) -> i32 {
match (exit_code, signal) {
(Some(code), _) => code,
(None, Some(signal)) => 128 + signal,
(None, None) => 1,
}
}
#[cfg(all(test, target_os = "linux"))]
mod tests {
use super::*;
#[test]
fn write_u32_decimal_writes_pid_bytes_without_allocation() {
let mut buf = [0_u8; 32];
let len = write_u32_decimal(12345, &mut buf);
assert_eq!(&buf[..len], b"12345");
}
#[test]
fn phase_timing_log_is_experiment_scoped_json() {
let tmp = tempfile::tempdir().unwrap();
let path = tmp.path().join("phase-timing.json");
let spec = RunSpec::new(vec!["true".to_string()])
.with_run_id("run_001")
.with_agent_shim("openai-agents");
let mut phases = std::collections::BTreeMap::new();
phases.insert("child_runtime_ms", 12.5);
write_phase_timing_log(Some(&path), &spec, &phases, Some(0), None, None).unwrap();
let payload: serde_json::Value =
serde_json::from_slice(&std::fs::read(path).unwrap()).unwrap();
assert_eq!(payload["schema"], "assay.experiment.runner_phase_timing.v0");
assert_eq!(payload["run_id"], "run_001");
assert_eq!(payload["phases_ms"]["child_runtime_ms"], 12.5);
assert_eq!(payload["exit_code"], 0);
assert!(payload["error"].is_null());
}
}