use pointlock_ir::{
ActionOutcome, HandlerHook, RunLogEvent, 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;
pub const TIMELINE_MAX_PAGE_SIZE: u32 = 50;
pub const TIMELINE_TEXT_MAX_BYTES: usize = 4 * 1024;
pub const TIMELINE_JSON_MAX_BYTES: usize = 16 * 1024;
pub const TIMELINE_JSON_MAX_DEPTH: usize = 12;
pub const TIMELINE_EVIDENCE_MAX: usize = 32;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase")]
pub enum RunTimelineFilter {
All,
Observations,
Actions,
Errors,
Verdicts,
}
impl RunTimelineFilter {
pub fn admits(self, payload: &RunLogPayload) -> bool {
match self {
RunTimelineFilter::All => true,
RunTimelineFilter::Observations => {
matches!(payload, RunLogPayload::ObservationRecorded { .. })
}
RunTimelineFilter::Actions => matches!(
payload,
RunLogPayload::ActionIntent { .. } | RunLogPayload::ActionSettled { .. }
),
RunTimelineFilter::Errors => match payload {
RunLogPayload::ActionSettled { outcome, .. } => {
!matches!(outcome, ActionOutcome::Succeeded { .. })
}
RunLogPayload::HandlerTriggered { hook, .. } => *hook == HandlerHook::OnError,
_ => false,
},
RunTimelineFilter::Verdicts => matches!(
payload,
RunLogPayload::PreflightProbed { .. }
| RunLogPayload::AssertionEvaluated { .. }
| RunLogPayload::VerdictRecorded { .. }
),
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct TimelineEvidenceRef {
pub id: String,
pub media_type: String,
pub sha256: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct TimelineErrorView {
#[serde(skip_serializing_if = "Option::is_none")]
pub error_class: Option<String>,
pub code: String,
pub message: String,
pub retryable: bool,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(tag = "type", rename_all = "camelCase")]
pub enum TimelineDetail {
#[serde(rename_all = "camelCase")]
RunStarted {
supervise_policy: Option<String>,
},
#[serde(rename_all = "camelCase")]
StepEntered {
step_id: String,
},
#[serde(rename_all = "camelCase")]
PreflightProbed {
pass: u32,
fail: u32,
unknown: u32,
#[serde(default, skip_serializing_if = "core::ops::Not::not")]
unprobed: bool,
},
#[serde(rename_all = "camelCase")]
ActionIntent {
call_id: String,
args: BoundedValue,
},
#[serde(rename_all = "camelCase")]
ActionSettled {
call_id: String,
outcome: String,
#[serde(skip_serializing_if = "Option::is_none")]
execution_mode: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
fallback_reason: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
error: Option<TimelineErrorView>,
},
#[serde(rename_all = "camelCase")]
ObservationRecorded {
observation_id: String,
captured_at_ms: u64,
#[serde(skip_serializing_if = "Option::is_none")]
screenshot_omission: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
ui_snapshot_omission: Option<String>,
},
#[serde(rename_all = "camelCase")]
AssertionEvaluated {
assert_id: String,
result: String,
#[serde(skip_serializing_if = "Option::is_none")]
channel: Option<String>,
reason: String,
},
#[serde(rename_all = "camelCase")]
VerdictRecorded {
status: String,
degraded: bool,
#[serde(skip_serializing_if = "Option::is_none")]
supersedes: Option<String>,
summary: String,
#[serde(skip_serializing_if = "Option::is_none")]
remote_archival_error: Option<String>,
},
#[serde(rename_all = "camelCase")]
StepExited {
state: StepState,
},
#[serde(rename_all = "camelCase")]
CallFramePushed {
callee: String,
#[serde(default, skip_serializing_if = "core::ops::Not::not")]
rebase: bool,
},
#[serde(rename_all = "camelCase")]
CallFramePopped {
has_outputs: bool,
},
#[serde(rename_all = "camelCase")]
HandlerTriggered {
hook: String,
trigger: u64,
#[serde(skip_serializing_if = "Option::is_none")]
disposition: Option<String>,
},
#[serde(rename_all = "camelCase")]
HumanRequested {
request_id: String,
purpose: String,
#[serde(skip_serializing_if = "Option::is_none")]
mode: Option<String>,
prompt: String,
},
#[serde(rename_all = "camelCase")]
HumanResponded {
request_id: String,
purpose: String,
actor: String,
response: BoundedValue,
},
#[serde(rename_all = "camelCase")]
RunSuspended {
reason: Option<String>,
},
#[serde(rename_all = "camelCase")]
RunResumed {
alignment: BoundedValue,
supervise_policy: Option<String>,
},
#[serde(rename_all = "camelCase")]
RunFinished {
#[serde(skip_serializing_if = "Option::is_none")]
status: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
degraded: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
remote_archival_error: Option<String>,
},
#[serde(rename_all = "camelCase")]
OverLimit {
event_type: String,
},
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct BoundedValue {
pub value: Value,
pub truncated: bool,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct RunTimelineEntry {
pub seq: u64,
pub at_ms: u64,
pub run_path: String,
pub detail: TimelineDetail,
pub evidence: Vec<TimelineEvidenceRef>,
pub evidence_omitted: u32,
pub truncated: bool,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct TimelinePage {
pub projection_version: ProjectionVersion,
pub run_id: String,
pub filter: RunTimelineFilter,
pub page: u32,
pub page_size: u32,
pub revision: u64,
pub total: u64,
pub entries: Vec<RunTimelineEntry>,
}
pub fn timeline_page(
store: &Store,
run_id: &str,
filter: RunTimelineFilter,
page: u32,
page_size: u32,
) -> Result<TimelinePage, StoreError> {
let events = store.events(run_id)?;
let revision = events.last().map(|event| event.seq).unwrap_or(0);
let mut error_classes: std::collections::BTreeMap<String, String> =
std::collections::BTreeMap::new();
if let Some((_, view)) = store.materialized_checkpoint(run_id)? {
for record in &view.completed {
for attempt in &record.attempts {
if let Some(class) = &attempt.error_class {
error_classes.insert(attempt.call_id.clone(), wire(class));
}
}
}
}
let admitted: Vec<&RunLogEvent> = events
.iter()
.filter(|event| filter.admits(&event.payload))
.collect();
let page = page.max(1);
let page_size = page_size.clamp(1, TIMELINE_MAX_PAGE_SIZE);
let start = (page as usize - 1).saturating_mul(page_size as usize);
let entries = admitted
.iter()
.skip(start)
.take(page_size as usize)
.map(|event| entry_of(event, &error_classes))
.collect();
Ok(TimelinePage {
projection_version: ProjectionVersion,
run_id: run_id.to_owned(),
filter,
page,
page_size,
revision,
total: admitted.len() as u64,
entries,
})
}
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 bound_text(text: &str) -> (String, bool) {
if text.len() <= TIMELINE_TEXT_MAX_BYTES {
return (text.to_owned(), false);
}
let mut cut = TIMELINE_TEXT_MAX_BYTES;
while !text.is_char_boundary(cut) {
cut -= 1;
}
(text[..cut].to_owned(), true)
}
fn bound_depth(value: &Value, depth: usize, truncated: &mut bool) -> Value {
if depth == 0 {
*truncated = true;
return Value::String("…depth truncated".to_owned());
}
match value {
Value::Object(map) => Value::Object(
map.iter()
.map(|(k, v)| (k.clone(), bound_depth(v, depth - 1, truncated)))
.collect(),
),
Value::Array(items) => Value::Array(
items
.iter()
.map(|v| bound_depth(v, depth - 1, truncated))
.collect(),
),
Value::String(text) => {
let (bounded, cut) = bound_text(text);
if cut {
*truncated = true;
}
Value::String(bounded)
}
other => other.clone(),
}
}
fn bound_value(value: &Value) -> BoundedValue {
let mut truncated = false;
let bounded = bound_depth(value, TIMELINE_JSON_MAX_DEPTH, &mut truncated);
let size = serde_json::to_string(&bounded)
.map(|s| s.len())
.unwrap_or(0);
if size > TIMELINE_JSON_MAX_BYTES {
return BoundedValue {
value: Value::String("…over 16 KiB, see the step dossier".to_owned()),
truncated: true,
};
}
BoundedValue {
value: bounded,
truncated,
}
}
fn entry_of(
event: &RunLogEvent,
error_classes: &std::collections::BTreeMap<String, String>,
) -> RunTimelineEntry {
let mut truncated = false;
let mut evidence: Vec<TimelineEvidenceRef> = Vec::new();
let mut evidence_omitted = 0u32;
let detail = match &event.payload {
RunLogPayload::RunStarted {
supervise_policy, ..
} => TimelineDetail::RunStarted {
supervise_policy: supervise_policy.as_ref().map(wire),
},
RunLogPayload::StepEntered { step_id, .. } => TimelineDetail::StepEntered {
step_id: step_id.to_string(),
},
RunLogPayload::PreflightProbed { outcomes } => {
let count = |status: pointlock_ir::VerdictStatus| {
outcomes.iter().filter(|o| o.result == status).count() as u32
};
TimelineDetail::PreflightProbed {
pass: count(pointlock_ir::VerdictStatus::Pass),
fail: count(pointlock_ir::VerdictStatus::Fail),
unknown: count(pointlock_ir::VerdictStatus::Unknown),
unprobed: outcomes.is_empty(),
}
}
RunLogPayload::ActionIntent {
call_id,
args_snapshot,
..
} => {
let args = bound_value(args_snapshot);
truncated |= args.truncated;
TimelineDetail::ActionIntent {
call_id: call_id.clone(),
args,
}
}
RunLogPayload::ActionSettled { call_id, outcome } => {
let (execution_mode, fallback_reason) = match outcome {
ActionOutcome::Succeeded { result } => {
for asset in &result.evidence {
push_evidence(asset, &mut evidence, &mut evidence_omitted);
}
match result.execution.as_ref() {
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),
}
}
_ => (None, None),
};
let error = match outcome {
ActionOutcome::Succeeded { .. } => None,
ActionOutcome::Failed { error }
| ActionOutcome::Cancelled { error }
| ActionOutcome::TimedOut { error } => {
let (message, cut) = bound_text(&error.message);
truncated |= cut;
Some(TimelineErrorView {
error_class: error_classes.get(call_id).cloned(),
code: error.code.clone(),
message,
retryable: error.retryable,
})
}
};
TimelineDetail::ActionSettled {
call_id: call_id.clone(),
outcome: outcome.kind().to_owned(),
execution_mode,
fallback_reason,
error,
}
}
RunLogPayload::ObservationRecorded { observation } => {
if let Some(evidence_ref) = &observation.screenshot {
push_evidence(&evidence_ref.asset, &mut evidence, &mut evidence_omitted);
}
if let Some(evidence_ref) = &observation.ui_snapshot {
push_evidence(&evidence_ref.asset, &mut evidence, &mut evidence_omitted);
}
TimelineDetail::ObservationRecorded {
observation_id: observation.observation_id.clone(),
captured_at_ms: observation.captured_at_ms,
screenshot_omission: observation.screenshot_omission.as_ref().map(wire),
ui_snapshot_omission: observation.ui_snapshot_omission.as_ref().map(wire),
}
}
RunLogPayload::AssertionEvaluated { outcome } => {
let (reason, cut) = bound_text(&outcome.reason);
truncated |= cut;
TimelineDetail::AssertionEvaluated {
assert_id: outcome.assert_id.to_string(),
result: wire(&outcome.result),
channel: outcome.channel.as_ref().map(wire),
reason,
}
}
RunLogPayload::VerdictRecorded {
verdict,
remote_archival_error,
..
} => {
let (summary, cut) = bound_text(&verdict.summary);
truncated |= cut;
for asset in &verdict.evidence {
push_evidence(asset, &mut evidence, &mut evidence_omitted);
}
let remote_archival_error = remote_archival_error.as_deref().map(|error| {
let (bounded, cut) = bound_text(error);
truncated |= cut;
bounded
});
TimelineDetail::VerdictRecorded {
status: wire(&verdict.status),
degraded: verdict.degraded,
supersedes: verdict.supersedes.clone(),
summary,
remote_archival_error,
}
}
RunLogPayload::StepExited { state, .. } => TimelineDetail::StepExited { state: *state },
RunLogPayload::CallFramePushed { frame, rebase } => TimelineDetail::CallFramePushed {
callee: format!("{}@{}", frame.flow_id, frame.ir_hash),
rebase: *rebase,
},
RunLogPayload::CallFramePopped { outputs } => TimelineDetail::CallFramePopped {
has_outputs: outputs.is_some(),
},
RunLogPayload::HandlerTriggered {
hook,
trigger,
disposition,
} => TimelineDetail::HandlerTriggered {
hook: wire(hook),
trigger: *trigger,
disposition: disposition.clone(),
},
RunLogPayload::HumanRequested {
request_id,
purpose,
mode,
prompt,
..
} => {
let (prompt, cut) = bound_text(prompt);
truncated |= cut;
TimelineDetail::HumanRequested {
request_id: request_id.clone(),
purpose: wire(purpose),
mode: mode.as_ref().map(wire),
prompt,
}
}
RunLogPayload::HumanResponded {
request_id,
purpose,
response,
actor,
} => {
let response = bound_value(response);
truncated |= response.truncated;
TimelineDetail::HumanResponded {
request_id: request_id.clone(),
purpose: wire(purpose),
actor: actor.clone(),
response,
}
}
RunLogPayload::RunSuspended { reason, .. } => TimelineDetail::RunSuspended {
reason: reason.clone(),
},
RunLogPayload::RunResumed {
alignment_report,
supervise_policy,
..
} => {
let counts = serde_json::json!({
"reusable": count_class(alignment_report, pointlock_ir::AlignmentClass::Reusable),
"judgeDirty": count_class(alignment_report, pointlock_ir::AlignmentClass::JudgeDirty),
"effectDirty": count_class(alignment_report, pointlock_ir::AlignmentClass::EffectDirty),
"new": count_class(alignment_report, pointlock_ir::AlignmentClass::New),
"orphaned": count_class(alignment_report, pointlock_ir::AlignmentClass::Orphaned),
});
TimelineDetail::RunResumed {
alignment: bound_value(&counts),
supervise_policy: supervise_policy.as_ref().map(wire),
}
}
RunLogPayload::RunFinished {
verdict,
remote_archival_error,
} => {
let remote_archival_error = remote_archival_error.as_deref().map(|error| {
let (bounded, cut) = bound_text(error);
truncated |= cut;
bounded
});
TimelineDetail::RunFinished {
status: verdict.as_ref().map(|v| wire(&v.status)),
degraded: verdict.as_ref().map(|v| v.degraded),
remote_archival_error,
}
}
};
let serialized = serde_json::to_string(&detail).map(|s| s.len()).unwrap_or(0);
let detail = if serialized > TIMELINE_JSON_MAX_BYTES {
truncated = true;
TimelineDetail::OverLimit {
event_type: event.payload.event_type().to_owned(),
}
} else {
detail
};
RunTimelineEntry {
seq: event.seq,
at_ms: event.at_ms,
run_path: render_run_path(&event.run_path),
detail,
evidence,
evidence_omitted,
truncated,
}
}
fn count_class(
report: &pointlock_ir::AlignmentReport,
class: pointlock_ir::AlignmentClass,
) -> usize {
report
.entries
.iter()
.filter(|entry| entry.class == class)
.count()
}
fn push_evidence(
asset: &pointlock_ir::AssetRef,
evidence: &mut Vec<TimelineEvidenceRef>,
omitted: &mut u32,
) {
if evidence.len() >= TIMELINE_EVIDENCE_MAX {
*omitted += 1;
return;
}
evidence.push(TimelineEvidenceRef {
id: asset.id.clone(),
media_type: asset.media_type.clone(),
sha256: asset.sha256.clone().unwrap_or_default(),
});
}