use std::collections::BTreeMap;
use pointlock_ir::{AlignmentClass, PathFrame, RunLogPayload, StepState, render_run_path};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use super::ProjectionVersion;
use crate::error::StoreError;
use crate::store::Store;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct AlignmentSummary {
pub reusable: u32,
pub judge_dirty: u32,
pub effect_dirty: u32,
pub new: u32,
pub orphaned: u32,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct StepStateSummary {
pub state: StepState,
#[serde(skip_serializing_if = "Option::is_none")]
pub verdict_status: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub degraded: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub act_chain_marks: Option<Vec<ActChainMark>>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct ActChainMark {
pub chain_index: u32,
pub mark: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub execution_mode: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub fallback_reason: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct RunOverview {
pub projection_version: ProjectionVersion,
pub run_id: String,
pub flow_id: String,
pub ir_hash: String,
pub lockfile_digest: String,
pub device_id: String,
pub session_lineage: Vec<String>,
pub status: String,
pub revision: u64,
pub created_at_ms: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub finished_at_ms: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub flow_verdict_status: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub flow_verdict_degraded: Option<bool>,
pub supervise_policy: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub alignment: Option<AlignmentSummary>,
pub awaiting_human: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub last_suspension_provider_state_summary: Option<pointlock_ir::ProviderStateSummary>,
pub steps: BTreeMap<String, StepStateSummary>,
}
fn wire<T: Serialize>(value: &T) -> String {
serde_json::to_value(value)
.ok()
.and_then(|v| v.as_str().map(str::to_owned))
.unwrap_or_default()
}
fn instance_path(path: &[PathFrame]) -> Vec<PathFrame> {
path.iter()
.filter(|frame| {
!matches!(
frame,
PathFrame::Attempt { .. } | PathFrame::Phase { .. } | PathFrame::Assertion { .. }
)
})
.cloned()
.collect()
}
pub fn run_overview(store: &Store, run_id: &str) -> Result<RunOverview, StoreError> {
let meta = store.run_meta(run_id)?;
let status = store.run_status(run_id)?;
let events = store.events(run_id)?;
let revision = events.last().map(|event| event.seq).unwrap_or(0);
let mut steps: BTreeMap<String, StepStateSummary> = BTreeMap::new();
let mut finished_at_ms = None;
let mut flow_verdict: Option<(String, bool)> = None;
let mut supervise_policy: Option<String> = None;
let mut alignment = None;
let mut awaiting: Option<(String, pointlock_ir::RunPath)> = None;
let mut last_suspension_summary: Option<pointlock_ir::ProviderStateSummary> = None;
let mut session_lineage = meta.binding.session_lineage.clone();
let mut chain_marks: BTreeMap<String, Vec<ActChainMark>> = BTreeMap::new();
let mut intent_index: BTreeMap<String, (String, u32)> = BTreeMap::new();
let mut boundary_pending: std::collections::BTreeSet<String> =
std::collections::BTreeSet::new();
for event in &events {
let key = || render_run_path(&instance_path(&event.run_path));
match &event.payload {
RunLogPayload::RunStarted {
supervise_policy: policy,
..
} => {
supervise_policy = policy.as_ref().map(wire);
}
RunLogPayload::RunResumed {
alignment_report,
supervise_policy: policy,
event_cursor,
} => {
if let Some(cursor) = event_cursor
&& session_lineage.last() != Some(&cursor.session_id)
{
session_lineage.push(cursor.session_id.clone());
}
supervise_policy = policy.as_ref().map(wire);
last_suspension_summary = None;
let count = |class: AlignmentClass| {
alignment_report
.entries
.iter()
.filter(|entry| entry.class == class)
.count() as u32
};
alignment = Some(AlignmentSummary {
reusable: count(AlignmentClass::Reusable),
judge_dirty: count(AlignmentClass::JudgeDirty),
effect_dirty: count(AlignmentClass::EffectDirty),
new: count(AlignmentClass::New),
orphaned: count(AlignmentClass::Orphaned),
});
}
RunLogPayload::StepEntered { .. } => {
chain_marks.remove(&key());
steps.insert(
key(),
StepStateSummary {
state: StepState::Ready,
verdict_status: None,
degraded: None,
act_chain_marks: None,
},
);
}
RunLogPayload::ActionIntent {
call_id,
chain_index: Some(index),
..
} => {
let step_key = key();
if boundary_pending.remove(&step_key) {
chain_marks.remove(&step_key);
} else if let Some(marks) = chain_marks.get_mut(&step_key) {
let max_marked = marks.iter().map(|mark| mark.chain_index).max();
if max_marked.is_some_and(|max| *index < max) {
marks.clear();
} else {
marks.retain(|mark| mark.chain_index != *index);
}
}
intent_index.insert(call_id.clone(), (step_key, *index));
}
RunLogPayload::ActionSettled { call_id, outcome } => {
if let Some((step_key, index)) = intent_index.remove(call_id) {
let (mark, execution_mode, fallback_reason) = match outcome {
pointlock_ir::ActionOutcome::Succeeded { result } => {
let (mode, reason) = match &result.execution {
Some(pointlock_ir::ActionExecution::NativeSemantic { .. }) => {
(Some("nativeSemantic".to_owned()), None)
}
Some(pointlock_ir::ActionExecution::WebSemantic { .. }) => {
(Some("webSemantic".to_owned()), None)
}
Some(pointlock_ir::ActionExecution::CoordinateFallback {
fallback_reason,
..
}) => (
Some("coordinateFallback".to_owned()),
Some(wire(fallback_reason)),
),
None => (None, None),
};
("succeeded", mode, reason)
}
_ => ("crossed", None, None),
};
chain_marks.entry(step_key).or_default().push(ActChainMark {
chain_index: index,
mark: mark.to_owned(),
execution_mode,
fallback_reason,
});
}
}
RunLogPayload::HandlerTriggered { hook, .. } => {
if matches!(
hook,
pointlock_ir::HandlerHook::OnFail
| pointlock_ir::HandlerHook::OnError
| pointlock_ir::HandlerHook::OnUnknown
) {
boundary_pending.insert(key());
}
}
RunLogPayload::StepExited { state, .. } => {
if let Some(cell) = steps.get_mut(&key()) {
cell.state = *state;
}
if awaiting
.as_ref()
.is_some_and(|(_, path)| *path == event.run_path)
{
awaiting = None;
}
}
RunLogPayload::VerdictRecorded { verdict, .. } => {
if let Some(cell) = steps.get_mut(&key()) {
cell.verdict_status = Some(wire(&verdict.status));
cell.degraded = Some(verdict.degraded);
}
}
RunLogPayload::HumanRequested { request_id, .. } => {
awaiting = Some((request_id.clone(), event.run_path.clone()));
if let Some(cell) = steps.get_mut(&key()) {
cell.state = StepState::AwaitingHuman;
}
}
RunLogPayload::HumanResponded {
request_id,
purpose,
response,
..
} => {
let non_final = *purpose == pointlock_ir::HumanPurpose::Supervision
&& response.get("decision").and_then(Value::as_str) == Some("suspend");
if !non_final && awaiting.as_ref().is_some_and(|(id, _)| id == request_id) {
awaiting = None;
}
}
RunLogPayload::RunSuspended {
provider_state_summary,
..
} => {
last_suspension_summary = provider_state_summary.clone();
}
RunLogPayload::RunFinished {
verdict,
remote_archival_error: _,
} => {
finished_at_ms = Some(event.at_ms);
flow_verdict = verdict
.as_ref()
.map(|verdict| (wire(&verdict.status), verdict.degraded));
}
_ => {}
}
}
for (step_key, marks) in chain_marks {
if let Some(cell) = steps.get_mut(&step_key) {
cell.act_chain_marks = Some(marks);
}
}
if let Some((_, view)) = store.materialized_checkpoint(run_id)? {
let key = render_run_path(&instance_path(&view.frontier.run_path));
if let Some(cell) = steps.get_mut(&key) {
cell.state = view.frontier.state;
}
}
let (flow_verdict_status, flow_verdict_degraded) = match flow_verdict {
Some((status, degraded)) => (Some(status), Some(degraded)),
None => (None, None),
};
Ok(RunOverview {
projection_version: ProjectionVersion,
run_id: run_id.to_owned(),
flow_id: meta.flow_id.to_string(),
ir_hash: meta.ir_hash.to_string(),
lockfile_digest: meta.lockfile_digest.to_string(),
device_id: meta.binding.device_id.clone(),
session_lineage,
status: status.as_str().to_owned(),
revision,
created_at_ms: meta.created_at_ms,
finished_at_ms,
flow_verdict_status,
flow_verdict_degraded,
supervise_policy,
alignment,
awaiting_human: awaiting.is_some(),
last_suspension_provider_state_summary: match status {
crate::RunStatus::Suspended | crate::RunStatus::AwaitingHuman => {
last_suspension_summary
}
_ => None,
},
steps,
})
}