pointlock_store/projection/
overview.rs1use std::collections::BTreeMap;
9
10use pointlock_ir::{AlignmentClass, PathFrame, RunLogPayload, StepState, render_run_path};
11use schemars::JsonSchema;
12use serde::{Deserialize, Serialize};
13use serde_json::Value;
14
15use super::ProjectionVersion;
16use crate::error::StoreError;
17use crate::store::Store;
18
19#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
21#[serde(rename_all = "camelCase", deny_unknown_fields)]
22pub struct AlignmentSummary {
23 pub reusable: u32,
25 pub judge_dirty: u32,
27 pub effect_dirty: u32,
29 pub new: u32,
31 pub orphaned: u32,
33}
34
35#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
37#[serde(rename_all = "camelCase", deny_unknown_fields)]
38pub struct StepStateSummary {
39 pub state: StepState,
41 #[serde(skip_serializing_if = "Option::is_none")]
43 pub verdict_status: Option<String>,
44 #[serde(skip_serializing_if = "Option::is_none")]
46 pub degraded: Option<bool>,
47 #[serde(skip_serializing_if = "Option::is_none")]
54 pub act_chain_marks: Option<Vec<ActChainMark>>,
55}
56
57#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
63#[serde(rename_all = "camelCase", deny_unknown_fields)]
64pub struct ActChainMark {
65 pub chain_index: u32,
67 pub mark: String,
69 #[serde(skip_serializing_if = "Option::is_none")]
71 pub execution_mode: Option<String>,
72 #[serde(skip_serializing_if = "Option::is_none")]
74 pub fallback_reason: Option<String>,
75}
76
77#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
79#[serde(rename_all = "camelCase", deny_unknown_fields)]
80pub struct RunOverview {
81 pub projection_version: ProjectionVersion,
83 pub run_id: String,
85 pub flow_id: String,
87 pub ir_hash: String,
89 pub lockfile_digest: String,
92 pub device_id: String,
94 pub session_lineage: Vec<String>,
96 pub status: String,
98 pub revision: u64,
101 pub created_at_ms: u64,
103 #[serde(skip_serializing_if = "Option::is_none")]
105 pub finished_at_ms: Option<u64>,
106 #[serde(skip_serializing_if = "Option::is_none")]
109 pub flow_verdict_status: Option<String>,
110 #[serde(skip_serializing_if = "Option::is_none")]
112 pub flow_verdict_degraded: Option<bool>,
113 pub supervise_policy: Option<String>,
116 #[serde(skip_serializing_if = "Option::is_none")]
118 pub alignment: Option<AlignmentSummary>,
119 pub awaiting_human: bool,
121 #[serde(skip_serializing_if = "Option::is_none")]
127 pub last_suspension_provider_state_summary: Option<pointlock_ir::ProviderStateSummary>,
128 pub steps: BTreeMap<String, StepStateSummary>,
132}
133
134fn wire<T: Serialize>(value: &T) -> String {
136 serde_json::to_value(value)
137 .ok()
138 .and_then(|v| v.as_str().map(str::to_owned))
139 .unwrap_or_default()
140}
141
142fn instance_path(path: &[PathFrame]) -> Vec<PathFrame> {
144 path.iter()
145 .filter(|frame| {
146 !matches!(
147 frame,
148 PathFrame::Attempt { .. } | PathFrame::Phase { .. } | PathFrame::Assertion { .. }
149 )
150 })
151 .cloned()
152 .collect()
153}
154
155pub fn run_overview(store: &Store, run_id: &str) -> Result<RunOverview, StoreError> {
157 let meta = store.run_meta(run_id)?;
158 let status = store.run_status(run_id)?;
159 let events = store.events(run_id)?;
160 let revision = events.last().map(|event| event.seq).unwrap_or(0);
161
162 let mut steps: BTreeMap<String, StepStateSummary> = BTreeMap::new();
163 let mut finished_at_ms = None;
164 let mut flow_verdict: Option<(String, bool)> = None;
165 let mut supervise_policy: Option<String> = None;
166 let mut alignment = None;
167 let mut awaiting: Option<(String, pointlock_ir::RunPath)> = None;
168 let mut last_suspension_summary: Option<pointlock_ir::ProviderStateSummary> = None;
169 let mut session_lineage = meta.binding.session_lineage.clone();
170 let mut chain_marks: BTreeMap<String, Vec<ActChainMark>> = BTreeMap::new();
171 let mut intent_index: BTreeMap<String, (String, u32)> = BTreeMap::new();
172 let mut boundary_pending: std::collections::BTreeSet<String> =
173 std::collections::BTreeSet::new();
174
175 for event in &events {
176 let key = || render_run_path(&instance_path(&event.run_path));
177 match &event.payload {
178 RunLogPayload::RunStarted {
179 supervise_policy: policy,
180 ..
181 } => {
182 supervise_policy = policy.as_ref().map(wire);
183 }
184 RunLogPayload::RunResumed {
185 alignment_report,
186 supervise_policy: policy,
187 event_cursor,
188 } => {
189 if let Some(cursor) = event_cursor
193 && session_lineage.last() != Some(&cursor.session_id)
194 {
195 session_lineage.push(cursor.session_id.clone());
196 }
197 supervise_policy = policy.as_ref().map(wire);
200 last_suspension_summary = None;
201 let count = |class: AlignmentClass| {
202 alignment_report
203 .entries
204 .iter()
205 .filter(|entry| entry.class == class)
206 .count() as u32
207 };
208 alignment = Some(AlignmentSummary {
209 reusable: count(AlignmentClass::Reusable),
210 judge_dirty: count(AlignmentClass::JudgeDirty),
211 effect_dirty: count(AlignmentClass::EffectDirty),
212 new: count(AlignmentClass::New),
213 orphaned: count(AlignmentClass::Orphaned),
214 });
215 }
216 RunLogPayload::StepEntered { .. } => {
217 chain_marks.remove(&key());
219 steps.insert(
220 key(),
221 StepStateSummary {
222 state: StepState::Ready,
223 verdict_status: None,
224 degraded: None,
225 act_chain_marks: None,
226 },
227 );
228 }
229 RunLogPayload::ActionIntent {
230 call_id,
231 chain_index: Some(index),
232 ..
233 } => {
234 let step_key = key();
235 if boundary_pending.remove(&step_key) {
236 chain_marks.remove(&step_key);
238 } else if let Some(marks) = chain_marks.get_mut(&step_key) {
239 let max_marked = marks.iter().map(|mark| mark.chain_index).max();
240 if max_marked.is_some_and(|max| *index < max) {
241 marks.clear();
245 } else {
246 marks.retain(|mark| mark.chain_index != *index);
249 }
250 }
251 intent_index.insert(call_id.clone(), (step_key, *index));
252 }
253 RunLogPayload::ActionSettled { call_id, outcome } => {
254 if let Some((step_key, index)) = intent_index.remove(call_id) {
255 let (mark, execution_mode, fallback_reason) = match outcome {
256 pointlock_ir::ActionOutcome::Succeeded { result } => {
257 let (mode, reason) = match &result.execution {
258 Some(pointlock_ir::ActionExecution::NativeSemantic { .. }) => {
259 (Some("nativeSemantic".to_owned()), None)
260 }
261 Some(pointlock_ir::ActionExecution::WebSemantic { .. }) => {
262 (Some("webSemantic".to_owned()), None)
263 }
264 Some(pointlock_ir::ActionExecution::CoordinateFallback {
265 fallback_reason,
266 ..
267 }) => (
268 Some("coordinateFallback".to_owned()),
269 Some(wire(fallback_reason)),
270 ),
271 None => (None, None),
272 };
273 ("succeeded", mode, reason)
274 }
275 _ => ("crossed", None, None),
276 };
277 chain_marks.entry(step_key).or_default().push(ActChainMark {
278 chain_index: index,
279 mark: mark.to_owned(),
280 execution_mode,
281 fallback_reason,
282 });
283 }
284 }
285 RunLogPayload::HandlerTriggered { hook, .. } => {
286 if matches!(
295 hook,
296 pointlock_ir::HandlerHook::OnFail
297 | pointlock_ir::HandlerHook::OnError
298 | pointlock_ir::HandlerHook::OnUnknown
299 ) {
300 boundary_pending.insert(key());
301 }
302 }
303 RunLogPayload::StepExited { state, .. } => {
304 if let Some(cell) = steps.get_mut(&key()) {
305 cell.state = *state;
306 }
307 if awaiting.as_ref().is_some_and(|(_, path)| {
311 crate::fold::exit_settles_pending(&event.run_path, path)
312 }) {
313 awaiting = None;
314 }
315 }
316 RunLogPayload::VerdictRecorded { verdict, .. } => {
317 if let Some(cell) = steps.get_mut(&key()) {
318 cell.verdict_status = Some(wire(&verdict.status));
319 cell.degraded = Some(verdict.degraded);
320 }
321 }
322 RunLogPayload::HumanRequested { request_id, .. } => {
323 awaiting = Some((request_id.clone(), event.run_path.clone()));
324 if let Some(cell) = steps.get_mut(&key()) {
325 cell.state = StepState::AwaitingHuman;
326 }
327 }
328 RunLogPayload::HumanResponded {
329 request_id,
330 purpose,
331 response,
332 ..
333 } => {
334 let non_final = *purpose == pointlock_ir::HumanPurpose::Supervision
335 && response.get("decision").and_then(Value::as_str) == Some("suspend");
336 if !non_final && awaiting.as_ref().is_some_and(|(id, _)| id == request_id) {
337 awaiting = None;
338 }
339 }
340 RunLogPayload::RunSuspended {
341 provider_state_summary,
342 ..
343 } => {
344 last_suspension_summary = provider_state_summary.clone();
345 }
346 RunLogPayload::RunFinished {
347 verdict,
348 remote_archival_error: _,
349 } => {
350 finished_at_ms = Some(event.at_ms);
351 flow_verdict = verdict
352 .as_ref()
353 .map(|verdict| (wire(&verdict.status), verdict.degraded));
354 }
355 _ => {}
356 }
357 }
358
359 for (step_key, marks) in chain_marks {
361 if let Some(cell) = steps.get_mut(&step_key) {
362 cell.act_chain_marks = Some(marks);
363 }
364 }
365
366 if let Some((_, view)) = store.materialized_checkpoint(run_id)? {
368 let key = render_run_path(&instance_path(&view.frontier.run_path));
369 if let Some(cell) = steps.get_mut(&key) {
370 cell.state = view.frontier.state;
371 }
372 }
373
374 let (flow_verdict_status, flow_verdict_degraded) = match flow_verdict {
375 Some((status, degraded)) => (Some(status), Some(degraded)),
376 None => (None, None),
377 };
378
379 Ok(RunOverview {
380 projection_version: ProjectionVersion,
381 run_id: run_id.to_owned(),
382 flow_id: meta.flow_id.to_string(),
383 ir_hash: meta.ir_hash.to_string(),
384 lockfile_digest: meta.lockfile_digest.to_string(),
385 device_id: meta.binding.device_id.clone(),
386 session_lineage,
387 status: status.as_str().to_owned(),
388 revision,
389 created_at_ms: meta.created_at_ms,
390 finished_at_ms,
391 flow_verdict_status,
392 flow_verdict_degraded,
393 supervise_policy,
394 alignment,
395 awaiting_human: awaiting.is_some(),
396 last_suspension_provider_state_summary: match status {
397 crate::RunStatus::Suspended | crate::RunStatus::AwaitingHuman => {
398 last_suspension_summary
399 }
400 _ => None,
401 },
402 steps,
403 })
404}