use std::collections::BTreeMap;
use std::path::Path;
use pointlock_ir::{
AlignmentClass, AlignmentReport, PathFrame, RunLogPayload, RunPath, StepRecord, render_run_path,
};
use pointlock_store::Store;
use serde::Serialize;
use crate::commands::{store_failure, usage_failure, wire_str};
use crate::{Failure, OutputFormat, exit};
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct RunReport {
pointlock_report: u32,
run_id: String,
flow_id: String,
ir_hash: String,
lockfile_digest: String,
device_id: String,
session_lineage: Vec<String>,
status: String,
created_at_ms: u64,
#[serde(skip_serializing_if = "Option::is_none")]
finished_at_ms: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
flow_verdict: Option<VerdictLine>,
counts: ReportCounts,
steps: Vec<StepLine>,
humans: Vec<HumanLine>,
handlers: Vec<HandlerLine>,
segments: Vec<SegmentLine>,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct VerdictLine {
status: String,
degraded: bool,
summary: String,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct ReportCounts {
steps: u32,
pass: u32,
fail: u32,
unknown: u32,
unverified: u32,
degraded: u32,
remote_archival_failed: u32,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct StepLine {
run_path: String,
step_id: String,
#[serde(skip_serializing_if = "Option::is_none")]
verdict_status: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
degraded: Option<bool>,
unverified: bool,
#[serde(skip_serializing_if = "Option::is_none")]
state: Option<String>,
attempts: u32,
#[serde(skip_serializing_if = "Option::is_none")]
error_class: Option<String>,
#[serde(skip_serializing_if = "is_zero")]
evidence_gaps: u32,
superseded: u32,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct HumanLine {
request_id: String,
purpose: String,
#[serde(skip_serializing_if = "Option::is_none")]
mode: Option<String>,
prompt: String,
requested_at_ms: u64,
#[serde(skip_serializing_if = "Option::is_none")]
response: Option<serde_json::Value>,
#[serde(skip_serializing_if = "Option::is_none")]
actor: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
responded_at_ms: Option<u64>,
#[serde(skip_serializing_if = "std::ops::Not::not")]
closed_unanswered: bool,
#[serde(skip)]
raw_path: String,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct HandlerLine {
run_path: String,
hook: String,
triggers: u64,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct SegmentLine {
kind: String,
at_ms: u64,
supervise_policy: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
alignment: Option<AlignmentCounts>,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct AlignmentCounts {
reusable: u32,
judge_dirty: u32,
effect_dirty: u32,
new: u32,
orphaned: u32,
}
pub fn report(store_dir: &Path, run_id: &str, format: OutputFormat) -> Result<i32, Failure> {
let report_failure = |err: pointlock_store::StoreError| match err {
pointlock_store::StoreError::UnknownRun(_) => usage_failure(err.to_string()),
other => store_failure(other),
};
let store = Store::open(store_dir).map_err(store_failure)?;
let overview =
pointlock_store::projection::run_overview(&store, run_id).map_err(report_failure)?;
let view = store.rebuild_checkpoint(run_id).map_err(report_failure)?;
let events = store.events(run_id).map_err(report_failure)?;
let mut lineage: BTreeMap<String, u32> = BTreeMap::new();
let mut gap_tally: BTreeMap<String, u32> = BTreeMap::new();
let mut exit_states: BTreeMap<String, String> = BTreeMap::new();
let mut humans: Vec<HumanLine> = Vec::new();
let mut handlers: BTreeMap<(String, String), (String, u64)> = BTreeMap::new();
let mut segments: Vec<SegmentLine> = Vec::new();
for event in &events {
match &event.payload {
RunLogPayload::VerdictRecorded {
localization_gaps, ..
} => {
let site = site_key(&instance_path(&event.run_path));
*gap_tally
.entry(render_run_path(&instance_path(&event.run_path)))
.or_insert(0) += localization_gaps.len() as u32;
*lineage.entry(site).or_insert(0) += 1;
}
RunLogPayload::StepExited { state, .. } => {
exit_states.insert(site_key(&instance_path(&event.run_path)), wire_str(state));
let exited = render_run_path(&event.run_path);
for line in humans
.iter_mut()
.filter(|line| line.raw_path == exited && line.response.is_none())
{
line.closed_unanswered = true;
}
}
RunLogPayload::HumanRequested {
request_id,
purpose,
mode,
prompt,
..
} => humans.push(HumanLine {
request_id: request_id.clone(),
purpose: wire_str(purpose),
mode: mode.as_ref().map(wire_str),
prompt: prompt.clone(),
requested_at_ms: event.at_ms,
response: None,
actor: None,
responded_at_ms: None,
closed_unanswered: false,
raw_path: render_run_path(&event.run_path),
}),
RunLogPayload::HumanResponded {
request_id,
response,
actor,
..
} => {
if let Some(line) = humans.iter_mut().find(|line| {
line.request_id == *request_id
&& (line.response.is_none() || nonfinal_suspend(line))
}) {
line.response = Some(response.clone());
line.actor = Some(actor.clone());
line.responded_at_ms = Some(event.at_ms);
}
}
RunLogPayload::HandlerTriggered { hook, trigger, .. } => {
let key = (site_key(&event.run_path), wire_str(hook));
let entry = handlers
.entry(key)
.or_insert_with(|| (render_run_path(&event.run_path), 0));
entry.1 = entry.1.max(*trigger);
}
RunLogPayload::RunStarted {
supervise_policy, ..
} => segments.push(SegmentLine {
kind: "started".to_owned(),
at_ms: event.at_ms,
supervise_policy: supervise_policy.as_ref().map(wire_str),
alignment: None,
}),
RunLogPayload::RunResumed {
alignment_report,
supervise_policy,
..
} => segments.push(SegmentLine {
kind: "resumed".to_owned(),
at_ms: event.at_ms,
supervise_policy: supervise_policy.as_ref().map(wire_str),
alignment: Some(alignment_counts(alignment_report)),
}),
_ => {}
}
}
let steps: Vec<StepLine> = view
.completed
.iter()
.map(|record| step_line(record, &lineage, &exit_states, &gap_tally))
.collect();
let mut counts = tally(&steps);
counts.remote_archival_failed = count_remote_archival_failures(&events);
let run_report = RunReport {
pointlock_report: 1,
run_id: overview.run_id,
flow_id: overview.flow_id,
ir_hash: overview.ir_hash,
lockfile_digest: overview.lockfile_digest,
device_id: overview.device_id,
session_lineage: overview.session_lineage,
status: overview.status,
created_at_ms: overview.created_at_ms,
finished_at_ms: overview.finished_at_ms,
flow_verdict: events
.iter()
.rev()
.find_map(|event| match &event.payload {
RunLogPayload::RunFinished { verdict, .. } => {
Some(verdict.as_ref().map(|verdict| VerdictLine {
status: wire_str(&verdict.status),
degraded: verdict.degraded,
summary: verdict.summary.clone(),
}))
}
_ => None,
})
.flatten(),
counts,
steps,
humans,
handlers: handlers
.into_iter()
.map(|((_site, hook), (run_path, triggers))| HandlerLine {
run_path,
hook,
triggers,
})
.collect(),
segments,
};
match format {
OutputFormat::Json => {
let body = serde_json::to_string_pretty(&run_report)
.map_err(|err| Failure::new(exit::INTERNAL, format!("serialize report: {err}")))?;
println!("{body}");
}
OutputFormat::Text => print_text(&run_report),
}
Ok(exit::PASS)
}
fn site_key(path: &RunPath) -> String {
path.iter()
.map(|frame| match frame {
PathFrame::Flow { flow_id, .. } => format!("f:{flow_id}"),
PathFrame::Step { step_id } => format!("s:{step_id}"),
PathFrame::Call {
step_id,
callee_flow_id,
..
} => format!(
"c:{}→{callee_flow_id}",
step_id
.as_ref()
.map(|step_id| step_id.as_str())
.unwrap_or("")
),
other => pointlock_ir::render_run_path(std::slice::from_ref(other)),
})
.collect::<Vec<_>>()
.join("/")
}
fn nonfinal_suspend(line: &HumanLine) -> bool {
line.purpose == "supervision"
&& line
.response
.as_ref()
.and_then(|response| response.get("decision"))
.and_then(serde_json::Value::as_str)
== Some("suspend")
}
fn is_zero(count: &u32) -> bool {
*count == 0
}
fn instance_path(path: &RunPath) -> RunPath {
path.iter()
.filter(|frame| {
!matches!(
frame,
PathFrame::Attempt { .. } | PathFrame::Phase { .. } | PathFrame::Assertion { .. }
)
})
.cloned()
.collect()
}
fn step_line(
record: &StepRecord,
lineage: &BTreeMap<String, u32>,
exit_states: &BTreeMap<String, String>,
gap_tally: &BTreeMap<String, u32>,
) -> StepLine {
let site = site_key(&instance_path(&record.run_path));
let error_class = record
.attempts
.iter()
.rev()
.find_map(|attempt| attempt.error_class.as_ref().map(wire_str));
let state = exit_states.get(&site).cloned();
StepLine {
step_id: record.step_id.as_str().to_owned(),
verdict_status: record
.verdict
.as_ref()
.map(|verdict| wire_str(&verdict.status)),
degraded: record.verdict.as_ref().map(|verdict| verdict.degraded),
unverified: record.verdict.is_none()
&& !record.attempts.is_empty()
&& state.as_deref() == Some("judged"),
state,
attempts: record.attempts.len() as u32,
error_class,
superseded: lineage
.get(&site)
.map(|count| count.saturating_sub(1))
.unwrap_or(0),
evidence_gaps: gap_tally
.get(&render_run_path(&instance_path(&record.run_path)))
.copied()
.unwrap_or(0),
run_path: render_run_path(&record.run_path),
}
}
fn tally(steps: &[StepLine]) -> ReportCounts {
let mut counts = ReportCounts {
steps: steps.len() as u32,
pass: 0,
fail: 0,
unknown: 0,
unverified: 0,
degraded: 0,
remote_archival_failed: 0,
};
for step in steps {
match step.verdict_status.as_deref() {
Some("pass") => counts.pass += 1,
Some("fail") => counts.fail += 1,
Some("unknown") => counts.unknown += 1,
_ => {}
}
if step.unverified {
counts.unverified += 1;
}
if step.degraded == Some(true) {
counts.degraded += 1;
}
}
counts
}
fn count_remote_archival_failures(events: &[pointlock_ir::RunLogEvent]) -> u32 {
events
.iter()
.filter(|event| {
matches!(
&event.payload,
RunLogPayload::VerdictRecorded {
remote_archival_error: Some(_),
..
} | RunLogPayload::RunFinished {
remote_archival_error: Some(_),
..
}
)
})
.count() as u32
}
fn alignment_counts(report: &AlignmentReport) -> AlignmentCounts {
let mut counts = AlignmentCounts {
reusable: 0,
judge_dirty: 0,
effect_dirty: 0,
new: 0,
orphaned: 0,
};
for entry in &report.entries {
match entry.class {
AlignmentClass::Reusable => counts.reusable += 1,
AlignmentClass::JudgeDirty => counts.judge_dirty += 1,
AlignmentClass::EffectDirty => counts.effect_dirty += 1,
AlignmentClass::New => counts.new += 1,
AlignmentClass::Orphaned => counts.orphaned += 1,
}
}
counts
}
fn print_text(report: &RunReport) {
println!("run: {}", report.run_id);
println!("flow: {} ({})", report.flow_id, report.ir_hash);
println!("lockfile: {}", report.lockfile_digest);
println!(
"device: {} (sessions: {})",
report.device_id,
report.session_lineage.join(" → ")
);
print!("status: {}", report.status);
match report.finished_at_ms {
Some(finished) => println!(
" (created {} ms, finished {finished} ms)",
report.created_at_ms
),
None => println!(" (created {} ms)", report.created_at_ms),
}
match &report.flow_verdict {
Some(verdict) => println!(
"flow verdict: {}{} — {}",
verdict.status,
if verdict.degraded { " [degraded]" } else { "" },
verdict.summary
),
None => println!("flow verdict: none (unfinished, finished unverified, or aborted)"),
}
let counts = &report.counts;
println!(
"steps: {} — {} pass, {} fail, {} unknown; {} unverified; {} degraded",
counts.steps, counts.pass, counts.fail, counts.unknown, counts.unverified, counts.degraded
);
if counts.remote_archival_failed > 0 {
println!(
"remote archival failed: {} verdict write-back(s) — local verdicts \
unaffected (04 §5)",
counts.remote_archival_failed
);
}
println!();
println!("steps:");
for step in &report.steps {
let mut line = format!(" {} ", step.run_path);
match &step.verdict_status {
Some(status) => {
line.push_str(status);
if step.degraded == Some(true) {
line.push_str(" [degraded]");
}
}
None if step.unverified => line.push_str("unverified (executed, no assertions)"),
None => match &step.state {
Some(state) if state == "skipped" => line.push_str("skipped (branch not taken)"),
Some(state) => line.push_str(&format!("no verdict (state {state})")),
None => line.push_str("no verdict"),
},
}
if step.superseded > 0 {
line.push_str(&format!(" (supersedes {} prior)", step.superseded));
}
if let Some(error_class) = &step.error_class {
line.push_str(&format!(" [last error: {error_class}]"));
}
if step.evidence_gaps > 0 {
line.push_str(&format!(" [evidence gaps: {}]", step.evidence_gaps));
}
println!("{line}");
}
if !report.humans.is_empty() {
println!();
println!("humans:");
for human in &report.humans {
let mode = human.mode.as_deref().unwrap_or("-");
let mut line = format!(
" {} {} ({mode}) \"{}\" requested at {} ms",
human.request_id, human.purpose, human.prompt, human.requested_at_ms
);
match (&human.response, &human.actor, human.responded_at_ms) {
(Some(response), Some(actor), Some(at_ms)) => {
line.push_str(&format!(" → {response} by {actor} at {at_ms} ms"));
}
_ if human.closed_unanswered => {
line.push_str(" → closed unanswered (settled by step exit, 06 §5.3)");
}
_ => line.push_str(" → pending"),
}
println!("{line}");
}
}
if !report.handlers.is_empty() {
println!();
println!("handlers:");
for handler in &report.handlers {
println!(
" {} {} ×{}",
handler.run_path, handler.hook, handler.triggers
);
}
}
println!();
println!("segments:");
for (index, segment) in report.segments.iter().enumerate() {
let supervise = segment.supervise_policy.as_deref().unwrap_or("null");
let mut line = format!(
" {}. {} at {} ms, supervise: {supervise}",
index + 1,
segment.kind,
segment.at_ms
);
if let Some(alignment) = &segment.alignment {
line.push_str(&format!(
", alignment: {} reusable / {} judgeDirty / {} effectDirty / {} new / {} orphaned",
alignment.reusable,
alignment.judge_dirty,
alignment.effect_dirty,
alignment.new,
alignment.orphaned
));
}
println!("{line}");
}
}
#[cfg(test)]
mod tests {
use super::*;
use pointlock_ir::{RunLogEvent, Verdict, VerdictStatus};
fn event(seq: u64, payload: RunLogPayload) -> RunLogEvent {
RunLogEvent {
run_id: "run-1".to_owned(),
seq,
at_ms: 1_000 + seq,
run_path: Vec::new(),
payload,
}
}
fn verdict() -> Verdict {
Verdict {
status: VerdictStatus::Pass,
degraded: false,
summary: "pass".to_owned(),
evidence: Vec::new(),
supersedes: None,
}
}
#[test]
fn remote_archival_failures_are_tallied_across_both_carriers() {
let annotated = vec![
event(
1,
RunLogPayload::VerdictRecorded {
verdict: verdict(),
localized: Vec::new(),
localization_gaps: Vec::new(),
remote_archival_error: Some("remote archival failed: gone".to_owned()),
},
),
event(
2,
RunLogPayload::VerdictRecorded {
verdict: verdict(),
localized: Vec::new(),
localization_gaps: Vec::new(),
remote_archival_error: None,
},
),
event(
3,
RunLogPayload::RunFinished {
verdict: Some(verdict()),
remote_archival_error: Some("remote archival failed: gone".to_owned()),
},
),
];
assert_eq!(count_remote_archival_failures(&annotated), 2);
let clean = vec![
event(
1,
RunLogPayload::VerdictRecorded {
verdict: verdict(),
localized: Vec::new(),
localization_gaps: Vec::new(),
remote_archival_error: None,
},
),
event(
2,
RunLogPayload::RunFinished {
verdict: Some(verdict()),
remote_archival_error: None,
},
),
];
assert_eq!(count_remote_archival_failures(&clean), 0);
}
}