use std::path::PathBuf;
use std::sync::Arc;
use std::time::{Instant, SystemTime};
pub type Slot = u16;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TaskOutcome {
Completed,
Failed(String),
Cancelled,
TimedOut,
}
#[derive(Debug, Clone, Default)]
pub struct RunSummary {
pub agents_spawned: u32,
pub programs_spawned: u32,
pub terminal_tasks: usize,
pub total_tasks: usize,
pub accounting: Option<AccountingRunSummary>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MessageLevel {
Info,
Warn,
Error,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AgentStream {
Stdout,
Stderr,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum DimensionStatus {
Measured,
Partial,
Unsupported,
Omitted,
Unknown,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct DimensionSummary {
pub value: Option<u64>,
pub status: DimensionStatus,
pub missing_count: u64,
pub measured_count: u64,
}
impl Default for DimensionSummary {
fn default() -> Self {
Self { value: None, status: DimensionStatus::Unknown, missing_count: 0, measured_count: 0 }
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum UsageCoverage {
Complete,
Partial,
Unpriced,
None,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum UsageStatus {
Measured,
UnsupportedAgent,
ExtractorUnavailable,
ExtractorFailed,
NoUsageEmitted,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum PricingStatus {
Priced,
PartialPrice,
Unpriced,
NotApplicable,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct UsageSummary {
pub invocation_id: String,
pub state: String,
pub agent: String,
pub provider: Option<String>,
pub model: Option<String>,
pub total: DimensionSummary,
pub input_total: DimensionSummary,
pub input_cached_read: DimensionSummary,
pub input_cache_write: DimensionSummary,
pub output_total: DimensionSummary,
pub output_cached_read: DimensionSummary,
pub output_cache_write: DimensionSummary,
pub cost_micro: Option<u64>,
pub priced_cost_micro: Option<u64>,
pub currency: Option<String>,
pub coverage: UsageCoverage,
pub status: UsageStatus,
pub pricing_status: PricingStatus,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct AccountingRunSummary {
pub total: DimensionSummary,
pub input_total: DimensionSummary,
pub input_cached_read: DimensionSummary,
pub input_cache_write: DimensionSummary,
pub output_total: DimensionSummary,
pub output_cached_read: DimensionSummary,
pub output_cache_write: DimensionSummary,
pub cost_micro: Option<u64>,
pub priced_cost_micro: Option<u64>,
pub currency: Option<String>,
pub coverage: UsageCoverage,
pub pricing_status: PricingStatus,
pub invocation_count: u64,
pub measured_invocation_count: u64,
pub missing_invocation_count: u64,
}
pub fn summarize_usage_summaries<'a>(
usages: impl IntoIterator<Item = &'a UsageSummary>,
) -> Option<AccountingRunSummary> {
let usages: Vec<&UsageSummary> = usages.into_iter().collect();
if usages.is_empty() {
return None;
}
let measured_invocation_count =
usages.iter().filter(|usage| usage.status == UsageStatus::Measured).count() as u64;
let missing_invocation_count = usages.len() as u64 - measured_invocation_count;
let priced_cost_micro = sum_options(usages.iter().map(|usage| usage.priced_cost_micro));
let cost_micro = if usages.iter().all(|usage| usage.cost_micro.is_some()) {
Some(usages.iter().filter_map(|usage| usage.cost_micro).sum())
} else {
None
};
let currency = usages.iter().find_map(|usage| usage.currency.clone());
let pricing_status = summarize_pricing_status(&usages);
let coverage = summarize_coverage(&usages, cost_micro, priced_cost_micro);
Some(AccountingRunSummary {
total: summarize_dimension(usages.iter().map(|usage| &usage.total)),
input_total: summarize_dimension(usages.iter().map(|usage| &usage.input_total)),
input_cached_read: summarize_dimension(usages.iter().map(|usage| &usage.input_cached_read)),
input_cache_write: summarize_dimension(usages.iter().map(|usage| &usage.input_cache_write)),
output_total: summarize_dimension(usages.iter().map(|usage| &usage.output_total)),
output_cached_read: summarize_dimension(
usages.iter().map(|usage| &usage.output_cached_read),
),
output_cache_write: summarize_dimension(
usages.iter().map(|usage| &usage.output_cache_write),
),
cost_micro,
priced_cost_micro,
currency,
coverage,
pricing_status,
invocation_count: usages.len() as u64,
measured_invocation_count,
missing_invocation_count,
})
}
fn summarize_dimension<'a>(
dimensions: impl IntoIterator<Item = &'a DimensionSummary>,
) -> DimensionSummary {
let mut value = 0u64;
let mut saw_value = false;
let mut missing_count = 0u64;
let mut measured_count = 0u64;
let mut unavailable_status = None;
for dimension in dimensions {
if let Some(v) = dimension.value {
value = value.saturating_add(v);
saw_value = true;
}
measured_count = measured_count.saturating_add(dimension.measured_count);
missing_count = missing_count.saturating_add(dimension.missing_count);
if dimension.status != DimensionStatus::Measured {
unavailable_status = Some(dimension.status);
}
}
let status = if saw_value && missing_count == 0 {
DimensionStatus::Measured
} else if saw_value {
DimensionStatus::Partial
} else {
unavailable_status.unwrap_or(DimensionStatus::Unknown)
};
DimensionSummary { value: saw_value.then_some(value), status, missing_count, measured_count }
}
fn sum_options(values: impl IntoIterator<Item = Option<u64>>) -> Option<u64> {
let mut total = 0u64;
let mut saw = false;
for value in values.into_iter().flatten() {
total = total.saturating_add(value);
saw = true;
}
saw.then_some(total)
}
fn summarize_pricing_status(usages: &[&UsageSummary]) -> PricingStatus {
let mut saw_priced = false;
let mut saw_partial = false;
let mut saw_unpriced = false;
let mut saw_applicable = false;
for usage in usages {
match usage.pricing_status {
PricingStatus::Priced => {
saw_priced = true;
saw_applicable = true;
}
PricingStatus::PartialPrice => {
saw_partial = true;
saw_applicable = true;
}
PricingStatus::Unpriced => {
saw_unpriced = true;
saw_applicable = true;
}
PricingStatus::NotApplicable => {}
}
}
if !saw_applicable {
PricingStatus::NotApplicable
} else if saw_partial || (saw_priced && saw_unpriced) {
PricingStatus::PartialPrice
} else if saw_priced {
PricingStatus::Priced
} else {
PricingStatus::Unpriced
}
}
fn summarize_coverage(
usages: &[&UsageSummary],
cost_micro: Option<u64>,
priced_cost_micro: Option<u64>,
) -> UsageCoverage {
if usages.iter().all(|usage| usage.coverage == UsageCoverage::None) {
return UsageCoverage::None;
}
if usages.iter().any(|usage| {
matches!(usage.coverage, UsageCoverage::Partial) || usage.status != UsageStatus::Measured
}) {
return UsageCoverage::Partial;
}
if cost_micro.is_some() {
UsageCoverage::Complete
} else if priced_cost_micro.is_some() {
UsageCoverage::Partial
} else if usages.iter().any(|usage| usage.coverage == UsageCoverage::Unpriced) {
UsageCoverage::Unpriced
} else {
UsageCoverage::None
}
}
#[derive(Debug, Clone)]
pub enum RunEvent {
RunStarted {
workspace: PathBuf,
parallel: u16,
total_tasks: usize,
},
PassStarted {
pass: u32,
ready: Vec<String>,
},
SlotAssigned {
slot: Slot,
task: String,
from: String,
to: String,
agent: Option<String>,
template_context: Option<crate::rhei_viz_model::TemplateContext>,
log_path: PathBuf,
started_at: Instant,
wall_clock: SystemTime,
},
SlotReleased {
slot: Slot,
task: String,
from: String,
to: String,
log_path: PathBuf,
outcome: TaskOutcome,
finished_at: Instant,
wall_clock: SystemTime,
exit_code: Option<i32>,
duration_ms: u64,
},
PassEnded {
pass: u32,
progressed: bool,
},
TasksDeferred {
pass: u32,
tasks: Vec<String>,
},
RunFinished {
summary: RunSummary,
},
Message {
level: MessageLevel,
text: String,
},
RunLink {
label: String,
url: String,
},
AgentOutput {
slot: Slot,
task: String,
stream: AgentStream,
line: String,
wall_clock: SystemTime,
},
UsageReported {
slot: Option<Slot>,
task: String,
invocation_id: String,
usage: UsageSummary,
},
TaskOutputsMissing {
task: String,
state: String,
entries: Vec<String>,
},
}
pub trait EventSink: Send + Sync {
fn emit(&self, event: RunEvent);
}
#[derive(Clone)]
pub struct Tee {
inners: Arc<Vec<Arc<dyn EventSink>>>,
}
impl Tee {
pub fn new(sinks: Vec<Arc<dyn EventSink>>) -> Self {
Self { inners: Arc::new(sinks) }
}
}
impl EventSink for Tee {
fn emit(&self, event: RunEvent) {
for sink in self.inners.iter() {
sink.emit(event.clone());
}
}
}
pub struct NullSink;
impl EventSink for NullSink {
fn emit(&self, _event: RunEvent) {}
}