use std::collections::{BTreeMap, HashMap, HashSet};
use polyc_eventlog_model::Event;
use polyc_projection::trace_vocabulary::MAX_STEPS_PER_POSITION;
use polyc_proto::kinds;
use polyc_proto::proto::polychrome::events::v1::{TurnDispatchedEvent, TurnFailedEvent};
use polyc_proto::proto::polychrome::harness::v1::TurnFailureKind;
use super::{
ApprovalDecision, ApprovalPhase, AudienceView, BoundaryMarker, ConversationTraceError,
ConversationTraceFacts, FailureKind, HandoffPhase, Id, MarkerKind, MessageRole,
PaymentDirection, QuestionDecision, QuestionPhase, RoleVerdict, SignerRole, SubagentPhase,
Text, ToolDocument, ToolStatus, TraceBoundsError, TraceFailureFact, TraceSignerFact,
TraceStepFact, TraceStepPayload, TraceTrust, TraceTurnFact, TraceWarningFact, TurnStatus,
WarningCode, bounded_hash, bounded_id_list, bounded_warning,
};
pub fn fold_conversation_trace(
events: &[(u64, Event)],
excised: &BTreeMap<u64, u64>,
trust: &TraceTrust,
) -> Result<ConversationTraceFacts, ConversationTraceError> {
refuse_repeated_positions(events)?;
let mut scan = Scan::default();
for (position, event) in events {
scan.fold_event(*position, event, excised, trust)?;
}
scan.callers = crate::caller_by_turn_last_wins(crate::attribution_events_with_positions(
events.iter().map(|(position, event)| (*position, event)),
crate::AttributionScope::CallerOnly,
))
.into_iter()
.filter_map(|(turn, persona)| {
Id::new("persona_id", &persona)
.ok()
.map(|persona| (turn.to_string(), persona))
})
.collect();
scan.resolve_payment_turns();
scan.resolve_tool_calls();
let turns = scan.build_turns();
let failures = rank_failures(&turns, &scan);
let mut steps = scan.steps;
steps.sort_by_key(|step| (step.position, step.ordinal));
Ok(ConversationTraceFacts {
signers: scan.signers,
steps,
turns,
failures,
warnings: scan.warnings,
})
}
fn is_turn_bearing_kind(base: &str) -> bool {
matches!(
base,
kinds::TURN_START
| kinds::TURN_DISPATCHED
| kinds::TURN_COMPLETE
| kinds::TURN_FAILED
| kinds::MODEL_CALL
| kinds::USAGE
| kinds::USER_MSG
| kinds::OUTPUT_MSG
| kinds::APPROVAL_REQUEST
| kinds::APPROVAL_RESPONSE
| kinds::APPROVAL_DEFERRED
| kinds::QUESTION_REQUEST
| kinds::QUESTION_RESPONSE
| kinds::SUBAGENT_SPAWN
| kinds::SUBAGENT_RESULT
| kinds::SUBAGENT_MODEL_CALL
)
}
fn is_payment_ledger_kind(base: &str) -> bool {
matches!(
base,
kinds::OUTBOUND_PAYMENT_ATTEMPT | kinds::OUTBOUND_PAYMENT_RECEIPT | kinds::PAYMENT_RECEIPT
)
}
fn refuse_repeated_positions(events: &[(u64, Event)]) -> Result<(), ConversationTraceError> {
let mut seen = HashSet::with_capacity(events.len());
for (position, _) in events {
if !seen.insert(*position) {
return Err(ConversationTraceError::RepeatedPosition {
position: *position,
});
}
}
Ok(())
}
#[derive(Default, Clone)]
#[allow(
clippy::struct_excessive_bools,
reason = "each boolean is one marker the writer either wrote or did not"
)]
struct TurnMarkers {
first_position: u64,
has_start: bool,
has_dispatched: bool,
has_complete: bool,
has_failed: bool,
has_ambiguous: bool,
failed_kind: Id,
failed_message: Text,
failed_position: Option<u64>,
model_provider: Id,
model_name: Id,
system_config_ref: Id,
temperature_micros: Option<u64>,
top_p_micros: Option<u64>,
max_tokens: Option<u64>,
reasoning_level: Id,
dispatch_clock_ms: Option<u64>,
usage_input: Option<u64>,
usage_output: Option<u64>,
evidence: Vec<u64>,
}
struct PendingCall {
position: u64,
turn_id: Id,
tool_call_id: Id,
}
#[derive(Default)]
struct Scan {
steps: Vec<TraceStepFact>,
warnings: Vec<TraceWarningFact>,
turn_order: Vec<String>,
markers: HashMap<String, TurnMarkers>,
ordinals: HashMap<u64, u64>,
calls: Vec<PendingCall>,
results: HashSet<(String, String)>,
result_sites: Vec<(u64, Id, Id)>,
open_approvals: HashMap<(String, String), Id>,
denied_calls: HashSet<(String, String)>,
open_questions: HashMap<(String, String, u32), Id>,
subagent_failures: Vec<(Id, Id)>,
signers: Vec<TraceSignerFact>,
callers: HashMap<String, Id>,
wallet_link_calls: Vec<(Id, Id)>,
}
impl Scan {
fn record_signer(&mut self, step: usize, role: SignerRole, key: &[u8]) {
let Ok(signer_key) = <[u8; 32]>::try_from(key) else {
return;
};
let Some(step) = self.steps.get(step) else {
return;
};
self.signers.push(TraceSignerFact {
position: step.position,
ordinal: step.ordinal,
step_kind: step.payload.step_kind(),
role,
signer_key,
});
}
fn push(
&mut self,
position: u64,
turn_id: Id,
kind_base: Id,
trust: Id,
payload: TraceStepPayload,
) -> Result<usize, ConversationTraceError> {
let ordinal = self.ordinals.entry(position).or_default();
if usize::try_from(*ordinal).unwrap_or(usize::MAX) >= MAX_STEPS_PER_POSITION {
return Err(ConversationTraceError::OrdinalOverflow { position });
}
let assigned = *ordinal;
*ordinal += 1;
self.steps.push(TraceStepFact {
position,
ordinal: assigned,
turn_id,
kind_base,
trust,
payload,
});
Ok(self.steps.len() - 1)
}
fn warn(&mut self, position: u64, code: WarningCode, message: &str) {
self.warnings.push(TraceWarningFact {
position: Some(position),
code,
message: bounded_warning(message),
});
}
fn touch_turn(&mut self, turn: &str, position: u64) {
let entry = self.markers.entry(turn.to_owned()).or_insert_with(|| {
self.turn_order.push(turn.to_owned());
TurnMarkers {
first_position: position,
..TurnMarkers::default()
}
});
entry.evidence.push(position);
}
}
struct Ctx<'a> {
position: u64,
base: &'a str,
kind_base: Id,
trust: Id,
turn_id: Id,
}
const fn at(position: u64) -> impl Fn(TraceBoundsError) -> ConversationTraceError {
move |source| ConversationTraceError::Bounds { position, source }
}
impl Scan {
#[allow(
clippy::too_many_lines,
reason = "one arm per journal kind; splitting the dispatch hides which kinds are covered"
)]
fn fold_event(
&mut self,
position: u64,
event: &Event,
excised: &BTreeMap<u64, u64>,
trust: &TraceTrust,
) -> Result<(), ConversationTraceError> {
let (base, turn) = kinds::parse(&event.kind);
let kind_base = Id::new("kind_base", base).map_err(at(position))?;
let step_trust = Id::new("trust", event.trust.as_str()).map_err(at(position))?;
let suffix = turn.map(|turn| turn.to_string());
let resolved = if is_payment_ledger_kind(base) {
None
} else if is_turn_bearing_kind(base) {
suffix
} else {
suffix.filter(|turn| self.markers.contains_key(turn))
};
let turn_id = Id::optional("turn_id", resolved.as_deref()).map_err(at(position))?;
if let Some(turn) = &resolved {
self.touch_turn(turn, position);
}
let ctx = Ctx {
position,
base,
kind_base,
trust: step_trust,
turn_id,
};
match base {
kinds::TURN_START
| kinds::TURN_COMPLETE
| kinds::TURN_FAILED
| kinds::TURN_DISPATCHED
| kinds::TURN_AMBIGUOUS
| kinds::STEP_COMMIT
| kinds::TURN_TEXT_WITHHELD
| kinds::INGRESS_DIRECTIVE => {
self.fold_boundary(&ctx, event)?;
}
kinds::USER_MSG | kinds::OUTPUT_MSG => {
self.fold_message(&ctx, event)?;
}
kinds::MODEL_CALL => self.fold_model_call(&ctx, event),
kinds::USAGE => self.fold_usage(&ctx, event),
kinds::APPROVAL_REQUEST | kinds::APPROVAL_RESPONSE | kinds::APPROVAL_DEFERRED => {
self.fold_approval(&ctx, event, trust)?;
}
kinds::QUESTION_REQUEST | kinds::QUESTION_RESPONSE => {
self.fold_question(&ctx, event, trust)?;
}
kinds::SUBAGENT_SPAWN | kinds::SUBAGENT_RESULT | kinds::SUBAGENT_MODEL_CALL => {
self.fold_subagent(&ctx, event, trust)?;
}
kinds::HANDOFF | kinds::HANDOFF_DENIED => {
self.fold_handoff(&ctx, event, trust)?;
}
kinds::SUMMARY => {
self.fold_summary(&ctx, event)?;
}
kinds::SUMMARY_GATE_REJECTED | kinds::SUMMARY_GATE_ADMITTED => {
self.fold_summary_gate(&ctx, event)?;
}
kinds::OBSERVED_IDENTITY => {
self.fold_identity(&ctx, event)?;
}
kinds::GRANT_REPLAY => {
self.fold_grant_replay(&ctx, event, trust)?;
}
kinds::OUTBOUND_PAYMENT_ATTEMPT => {
self.fold_payment_attempt(&ctx, event)?;
}
kinds::PAYMENT_RECEIPT | kinds::OUTBOUND_PAYMENT_RECEIPT => {
self.fold_payment_receipt(&ctx, event, trust)?;
}
kinds::PAYMENT_REFUSAL => {
self.fold_payment_refusal(&ctx, event, trust)?;
}
kinds::WALLET_LINK_LIFECYCLE => {
self.fold_wallet_link(&ctx, event, trust)?;
}
crate::EXCISED_KIND => {
self.push(
position,
ctx.turn_id.clone(),
ctx.kind_base.clone(),
ctx.trust.clone(),
TraceStepPayload::Excised {
marker_position: excised.get(&position).copied(),
},
)?;
}
polyc_eventlog_model::MMR_SIGNED_ROOT_KIND | kinds::COMPACTION_CHECKPOINT => {
}
_ => self.fold_marker_or_unknown(&ctx, event)?,
}
Ok(())
}
#[allow(
clippy::too_many_lines,
reason = "one arm per boundary marker; splitting hides which markers are covered"
)]
fn fold_boundary(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
base,
kind_base,
trust: step_trust,
turn_id,
} = ctx;
let position = *position;
let base: &str = base;
let marker = match base {
kinds::TURN_START => BoundaryMarker::Start,
kinds::TURN_DISPATCHED => BoundaryMarker::Dispatched,
kinds::TURN_COMPLETE => BoundaryMarker::Complete,
kinds::TURN_FAILED => BoundaryMarker::Failed,
kinds::TURN_AMBIGUOUS => BoundaryMarker::Ambiguous,
kinds::STEP_COMMIT => BoundaryMarker::StepCommit,
kinds::TURN_TEXT_WITHHELD => BoundaryMarker::TextWithheld,
_ => BoundaryMarker::IngressDirective,
};
let mut audience = None;
let mut source_turn_id = None;
let mut failed_kind = None;
let mut failed_message = None;
if base == kinds::TURN_DISPATCHED {
match polyc_proto::events_decode::try_decode_event_payload::<TurnDispatchedEvent>(
&event.payload,
) {
Ok(decoded) => {
audience = Some(AudienceView {
recorded: Id::new(
"audience",
polyc_proto::audience_display::recorded_visibility_label(
decoded.visibility,
),
)
.map_err(at(position))?,
source: Id::new(
"audience_source",
polyc_proto::audience_display::recorded_source_label(
decoded.visibility_source,
),
)
.map_err(at(position))?,
edge_asserted: Id::new(
"edge_asserted",
polyc_proto::audience_display::asserted_visibility_label(
decoded.edge_asserted_visibility,
),
)
.map_err(at(position))?,
});
if !decoded.source_turn_id.is_empty() {
source_turn_id = Some(
Id::new("source_turn_id", &decoded.source_turn_id)
.map_err(at(position))?,
);
}
}
Err(error) => self.warn(
position,
WarningCode::DecodeFailed,
&format!("turn_dispatched did not decode: {error}"),
),
}
}
if base == kinds::TURN_FAILED {
match polyc_proto::events_decode::try_decode_event_payload::<TurnFailedEvent>(
&event.payload,
) {
Ok(decoded) => {
let label = failure_label(decoded.kind);
failed_kind = Some(Id::new("failed_kind", label).map_err(at(position))?);
failed_message = Some(Text::new(&decoded.message));
}
Err(error) => self.warn(
position,
WarningCode::DecodeFailed,
&format!("turn_failed did not decode: {error}"),
),
}
}
if let Some(entry) = self.markers.get_mut(turn_id.as_str()) {
match marker {
BoundaryMarker::Start => entry.has_start = true,
BoundaryMarker::Dispatched => entry.has_dispatched = true,
BoundaryMarker::Complete => entry.has_complete = true,
BoundaryMarker::Ambiguous => entry.has_ambiguous = true,
BoundaryMarker::Failed => {
entry.has_failed = true;
entry.failed_position = Some(position);
if let Some(kind) = failed_kind.clone() {
entry.failed_kind = kind;
}
if let Some(message) = failed_message.clone() {
entry.failed_message = message;
}
}
BoundaryMarker::StepCommit
| BoundaryMarker::TextWithheld
| BoundaryMarker::IngressDirective => {}
}
}
self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::Boundary {
marker,
audience,
source_turn_id,
failed_kind,
failed_message,
},
)?;
Ok(())
}
}
fn failure_label(kind: buffa::EnumValue<TurnFailureKind>) -> &'static str {
match kind.as_known() {
Some(TurnFailureKind::RateLimit) => "rate_limit",
Some(TurnFailureKind::Timeout) => "timeout",
Some(TurnFailureKind::Unavailable) => "unavailable",
Some(TurnFailureKind::Auth) => "auth",
Some(TurnFailureKind::BadRequest) => "bad_request",
Some(TurnFailureKind::Unspecified | TurnFailureKind::Other) | None => "other",
}
}
impl Scan {
#[allow(
clippy::too_many_lines,
reason = "the shell and its one content block are one record, read together"
)]
fn fold_message(&mut self, ctx: &Ctx<'_>, event: &Event) -> Result<(), ConversationTraceError> {
let Ctx {
position,
base,
kind_base,
trust: step_trust,
turn_id,
} = ctx;
let position = *position;
let base: &str = base;
let Some(message) = polyc_proto::events_decode::decode_event_payload::<
polyc_proto::proto::polychrome::agent::v1::Message,
>(&event.payload) else {
self.warn(
position,
WarningCode::DecodeFailed,
&format!("{base} did not decode"),
);
return Ok(());
};
let internal_only = message.internal_only;
let turn_ref = (!turn_id.is_empty()).then(|| turn_id.as_str());
let folded =
crate::fold_message_content(&message, position, turn_ref, event.trust.as_str());
let warnings = folded.warnings;
let role = if base == kinds::USER_MSG {
MessageRole::Input
} else {
MessageRole::Output
};
let mut text = Text::default();
let mut tool_call_ids = Vec::new();
let mut call_step = None;
let mut result_step = None;
match folded.content {
crate::MessageContent::Text(fact) => text = Text::new(&fact.text),
crate::MessageContent::ToolCall(call) => {
let tool_call_id =
Id::new("tool_call_id", &call.tool_call_id).map_err(at(position))?;
tool_call_ids.push(tool_call_id.clone());
call_step = Some((
tool_call_id,
Id::new("name", &call.name).map_err(at(position))?,
document_of(&call.arguments),
));
}
crate::MessageContent::ToolResult(result) => {
let tool_call_id =
Id::new("tool_call_id", &result.tool_call_id).map_err(at(position))?;
self.results
.insert((turn_id.as_str().to_owned(), result.tool_call_id.clone()));
self.result_sites
.push((position, turn_id.clone(), tool_call_id.clone()));
result_step = Some((
tool_call_id,
Id::new("name", &result.name).map_err(at(position))?,
document_of(&result.result),
result.first_party,
needs_wallet_link(&result.name, &result.result),
));
}
crate::MessageContent::None => {}
}
self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::Message {
role,
text,
tool_call_ids,
internal_only,
},
)?;
if let Some((tool_call_id, name, arguments)) = call_step {
self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::ToolCall {
tool_call_id: tool_call_id.clone(),
name,
arguments,
status: ToolStatus::Unknown,
origin_turn_id: turn_id.clone(),
result_turn_id: None,
internal_only,
},
)?;
self.calls.push(PendingCall {
position,
turn_id: turn_id.clone(),
tool_call_id,
});
}
if let Some((tool_call_id, name, result, first_party, wallet)) = result_step {
if wallet {
self.wallet_link_calls
.push((turn_id.clone(), tool_call_id.clone()));
}
self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::ToolResult {
tool_call_id,
name,
result,
first_party,
needs_wallet_link: wallet,
internal_only,
},
)?;
}
for warning in warnings {
self.warn(position, WarningCode::DecodeFailed, &warning);
}
Ok(())
}
fn fold_model_call(&mut self, ctx: &Ctx<'_>, event: &Event) {
let Ctx {
position, turn_id, ..
} = ctx;
let position = *position;
match crate::fold_model_call_event(&event.payload) {
Ok(fact) => {
if let Some(entry) = self.markers.get_mut(turn_id.as_str()) {
entry.model_provider = Id::new("provider", &fact.provider).unwrap_or_default();
entry.model_name = Id::new("model", &fact.model).unwrap_or_default();
entry.system_config_ref =
Id::new("system_config_ref", &fact.system_config_ref).unwrap_or_default();
entry.temperature_micros = fact.decode_params.temperature.map(to_micros);
entry.top_p_micros = fact.decode_params.top_p.map(to_micros);
entry.max_tokens = fact.decode_params.max_tokens.map(u64::from);
entry.reasoning_level = fact
.decode_params
.reasoning_level
.as_deref()
.map(|level| Id::new("reasoning_level", level).unwrap_or_default())
.unwrap_or_default();
entry.dispatch_clock_ms = Some(fact.captured_clock_unix_ms);
}
}
Err(error) => self.warn(
position,
WarningCode::DecodeFailed,
&format!("model_call did not decode: {error}"),
),
}
}
fn fold_usage(&mut self, ctx: &Ctx<'_>, event: &Event) {
let Ctx {
position, turn_id, ..
} = ctx;
let position = *position;
match crate::fold_usage_event(&event.payload) {
Ok(fact) => {
if let Some(entry) = self.markers.get_mut(turn_id.as_str()) {
entry.usage_input = Some(fact.input_tokens);
entry.usage_output = Some(fact.output_tokens);
}
}
Err(error) => self.warn(
position,
WarningCode::DecodeFailed,
&format!("usage did not decode: {error}"),
),
}
}
}
fn document_of(value: &serde_json::Value) -> ToolDocument {
ToolDocument::new(&serde_json::to_string(value).unwrap_or_else(|_| "{}".to_owned()))
}
fn needs_wallet_link(tool_name: &str, result: &serde_json::Value) -> bool {
if result.get("needs_wallet_link").is_some() {
return true;
}
tool_name == "paid_fetch"
&& result
.get("error")
.and_then(serde_json::Value::as_str)
.is_some_and(|error| {
error.contains("no spending wallet") || error.contains("Link a wallet")
})
}
fn to_micros(value: f64) -> u64 {
const CEILING: f64 = 18_446_744_073_709_551_616.0;
let scaled = value * 1_000_000.0;
if scaled.is_nan() || scaled <= 0.0 {
return 0;
}
if scaled >= CEILING {
return u64::MAX;
}
#[allow(
clippy::cast_possible_truncation,
clippy::cast_sign_loss,
reason = "guarded above: positive, non-NaN, and below the u64 ceiling"
)]
{
scaled as u64
}
}
impl Scan {
fn fold_approval(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
trust: &TraceTrust,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
base,
kind_base,
trust: step_trust,
turn_id,
} = ctx;
let position = *position;
let base: &str = base;
if base == kinds::APPROVAL_DEFERRED {
let request_id =
polyc_crypto::approval::decode_request_id(&event.payload).unwrap_or_default();
let request_id = Id::new("request_id", &request_id).map_err(at(position))?;
self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::Approval {
phase: ApprovalPhase::Deferred,
request_id,
tool_name: None,
args: None,
decision: ApprovalDecision::Deferred,
reason: None,
signature: None,
},
)?;
return Ok(());
}
let Some(fact) = crate::fold_approval_event(event, &trust.approval) else {
self.warn(
position,
WarningCode::DecodeFailed,
&format!("{base} did not decode"),
);
return Ok(());
};
let mut signer_key = None;
let payload = match fact {
crate::ApprovalFact::Request(request) => {
let request_id =
Id::new("request_id", &request.request_id).map_err(at(position))?;
self.open_approvals.insert(
(turn_id.as_str().to_owned(), request.request_id.clone()),
turn_id.clone(),
);
TraceStepPayload::Approval {
phase: ApprovalPhase::Request,
request_id,
tool_name: Some(
Id::new("tool_name", &request.tool_name).map_err(at(position))?,
),
args: Some(ToolDocument::new(&request.args_json)),
decision: ApprovalDecision::Pending,
reason: (!request.reason.is_empty()).then(|| Text::new(&request.reason)),
signature: None,
}
}
crate::ApprovalFact::Response(response) => {
if response.signature_status == crate::ApprovalSignatureStatus::Verified {
signer_key.clone_from(&response.signer_public_key);
}
let request_id =
Id::new("request_id", &response.request_id).map_err(at(position))?;
self.open_approvals
.remove(&(turn_id.as_str().to_owned(), response.request_id.clone()));
if !response.approved {
self.denied_calls
.insert((turn_id.as_str().to_owned(), response.request_id.clone()));
}
TraceStepPayload::Approval {
phase: ApprovalPhase::Response,
request_id,
tool_name: None,
args: None,
decision: if response.approved {
ApprovalDecision::Approved
} else {
ApprovalDecision::Denied
},
reason: response.response_reason.as_deref().map(Text::new),
signature: Some(approval_verdict(response.signature_status)),
}
}
};
let step = self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
payload,
)?;
if let Some(key) = signer_key {
self.record_signer(step, SignerRole::Approval, &key);
}
Ok(())
}
#[allow(
clippy::too_many_lines,
reason = "one block per question phase; the two read different shapes"
)]
fn fold_question(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
trust: &TraceTrust,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
base,
kind_base,
trust: step_trust,
turn_id,
} = ctx;
let position = *position;
let base: &str = base;
let Ok(value) = serde_json::from_slice::<serde_json::Value>(&event.payload) else {
self.warn(
position,
WarningCode::DecodeFailed,
&format!("{base} did not decode"),
);
return Ok(());
};
let bounded_index = |value: u64| u32::try_from(value).unwrap_or(u32::MAX);
let mut signer_key = None;
let payload = if base == kinds::QUESTION_REQUEST {
let call_id = Id::new(
"call_id",
value
.get("call_id")
.and_then(serde_json::Value::as_str)
.unwrap_or_default(),
)
.map_err(at(position))?;
let index = bounded_index(
value
.get("index")
.and_then(serde_json::Value::as_u64)
.unwrap_or_default(),
);
self.open_questions.insert(
(
turn_id.as_str().to_owned(),
call_id.as_str().to_owned(),
index,
),
turn_id.clone(),
);
TraceStepPayload::Question {
phase: QuestionPhase::Request,
call_id,
index,
header: value
.get("header")
.and_then(serde_json::Value::as_str)
.map(Text::new),
question: value
.get("question")
.and_then(serde_json::Value::as_str)
.map(Text::new),
args: Some(ToolDocument::new(&value.to_string())),
decision: QuestionDecision::Pending,
selected_index: None,
selected_label: None,
answered_by: None,
verdict: None,
}
} else {
let (verdict, verified) = question_verdict(&event.payload, &trust.question);
let Some(answer) = verified else {
self.warn(
position,
WarningCode::SignatureUnverified,
"question_response did not verify, so it names no question",
);
return Ok(());
};
signer_key = Some(answer.signer_public_key.clone());
let call_id = Id::new("call_id", &answer.call_id).map_err(at(position))?;
let index = bounded_index(u64::from(answer.index));
self.open_questions.remove(&(
turn_id.as_str().to_owned(),
call_id.as_str().to_owned(),
index,
));
TraceStepPayload::Question {
phase: QuestionPhase::Response,
call_id,
index,
header: None,
question: None,
args: None,
decision: question_decision(&answer.state, position)?,
selected_index: answer.selected_index,
selected_label: (!answer.selected_label.is_empty())
.then(|| Text::new(&answer.selected_label)),
answered_by: (!answer.answered_by.is_empty())
.then(|| Id::new("answered_by", &answer.answered_by))
.transpose()
.map_err(at(position))?,
verdict: Some(verdict),
}
};
let step = self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
payload,
)?;
if let Some(key) = signer_key {
self.record_signer(step, SignerRole::Question, &key);
}
Ok(())
}
}
const fn approval_verdict(status: crate::ApprovalSignatureStatus) -> RoleVerdict {
match status {
crate::ApprovalSignatureStatus::Verified => RoleVerdict::Verified,
crate::ApprovalSignatureStatus::Invalid => RoleVerdict::Invalid,
crate::ApprovalSignatureStatus::LegacyUnverifiable => RoleVerdict::LegacyUnverifiable,
}
}
fn question_verdict(
payload: &[u8],
question_trust: &[Vec<u8>],
) -> (
RoleVerdict,
Option<polyc_crypto::question::VerifiedQuestionAnswer>,
) {
let Some(verified) = polyc_crypto::question::verify_signed_answer(payload) else {
return (RoleVerdict::Invalid, None);
};
let verdict = if question_trust
.iter()
.any(|key| key.as_slice() == verified.signer_public_key)
{
RoleVerdict::Verified
} else {
RoleVerdict::Untrusted
};
(verdict, Some(verified))
}
fn question_decision(
label: &str,
position: u64,
) -> Result<QuestionDecision, ConversationTraceError> {
Ok(match label {
"pending" => QuestionDecision::Pending,
"answered" => QuestionDecision::Answered,
"declined" => QuestionDecision::Declined,
"auto_resolved" => QuestionDecision::AutoResolved,
other => QuestionDecision::Other(Id::new("decision", other).map_err(at(position))?),
})
}
impl Scan {
fn fold_subagent(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
trust: &TraceTrust,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
base,
kind_base,
trust: step_trust,
turn_id,
} = ctx;
let position = *position;
let base: &str = base;
let subagent_trust = &trust.subagent;
let payload = if base == kinds::SUBAGENT_SPAWN {
let Some(fact) = crate::fold_subagent_spawn(&event.payload, subagent_trust) else {
self.warn(
position,
WarningCode::DecodeFailed,
"subagent_spawn did not decode",
);
return Ok(());
};
TraceStepPayload::Subagent {
phase: SubagentPhase::Spawn,
sub_agent_id: Id::new("sub_agent_id", &fact.sub_agent_id).map_err(at(position))?,
target_agent_id: Id::new("target_agent_id", &fact.target_agent_id)
.map_err(at(position))?,
task: Some(Text::new(&fact.task)),
provider: Some(Id::new("provider", &fact.resolved_provider).map_err(at(position))?),
model: Some(Id::new("model", &fact.resolved_model).map_err(at(position))?),
succeeded: None,
error: None,
input_tokens: None,
output_tokens: None,
first_party: None,
verdict: Some(role_verdict(fact.signature_status)),
}
} else if base == kinds::SUBAGENT_RESULT {
let Some(fact) = crate::fold_subagent_result(&event.payload, subagent_trust) else {
self.warn(
position,
WarningCode::DecodeFailed,
"subagent_result did not decode",
);
return Ok(());
};
let sub_agent_id = Id::new("sub_agent_id", &fact.sub_agent_id).map_err(at(position))?;
if !fact.succeeded {
self.subagent_failures
.push((turn_id.clone(), sub_agent_id.clone()));
}
TraceStepPayload::Subagent {
phase: SubagentPhase::Result,
sub_agent_id,
target_agent_id: Id::new("target_agent_id", &fact.target_agent_id)
.map_err(at(position))?,
task: None,
provider: None,
model: None,
succeeded: Some(fact.succeeded),
error: (!fact.error.is_empty()).then(|| Text::new(&fact.error)),
input_tokens: Some(fact.input_tokens),
output_tokens: Some(fact.output_tokens),
first_party: Some(fact.first_party),
verdict: Some(role_verdict(fact.signature_status)),
}
} else {
let Some(call) = crate::fold_subagent_model_call(&event.payload) else {
self.warn(
position,
WarningCode::DecodeFailed,
"subagent_model_call did not decode",
);
return Ok(());
};
TraceStepPayload::Subagent {
phase: SubagentPhase::ModelCall,
sub_agent_id: Id::new("sub_agent_id", &call.sub_agent_id).map_err(at(position))?,
target_agent_id: Id::default(),
task: None,
provider: Some(Id::new("provider", &call.provider).map_err(at(position))?),
model: Some(Id::new("model", &call.model).map_err(at(position))?),
succeeded: None,
error: None,
input_tokens: None,
output_tokens: None,
first_party: None,
verdict: None,
}
};
self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
payload,
)?;
Ok(())
}
fn fold_handoff(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
trust: &TraceTrust,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
base,
kind_base,
trust: step_trust,
turn_id,
} = ctx;
let position = *position;
let base: &str = base;
let handoff_trust = &trust.handoff;
let Some(fact) = crate::fold_handoff_event(base, &event.payload, handoff_trust) else {
self.warn(
position,
WarningCode::DecodeFailed,
&format!("{base} did not decode"),
);
return Ok(());
};
let payload = match fact {
crate::HandoffFact::Handoff(spawn) => TraceStepPayload::Handoff {
phase: HandoffPhase::Handoff,
child_conversation_id: Some(
Id::new("child_conversation_id", &spawn.child_conversation_id)
.map_err(at(position))?,
),
child_agent_id: Id::new("child_agent_id", &spawn.child_agent_id)
.map_err(at(position))?,
carried_count: spawn.carried_count,
reason: Text::new(&spawn.reason),
parent_agent_id: None,
denial_reason: None,
allowed: Vec::new(),
verdict: role_verdict(spawn.signature_status),
},
crate::HandoffFact::Denied(denied) => {
let (allowed, clipped) =
bounded_id_list("allowed", denied.allowed.clone()).map_err(at(position))?;
if clipped {
self.warn(
position,
WarningCode::ListClipped,
"a handoff denial listed more allowed agents than the bound admits",
);
}
TraceStepPayload::Handoff {
phase: HandoffPhase::Denied,
child_conversation_id: None,
child_agent_id: Id::new("child_agent_id", &denied.child_agent_id)
.map_err(at(position))?,
carried_count: 0,
reason: Text::new(&denied.reason),
parent_agent_id: Some(
Id::new("parent_agent_id", &denied.parent_agent_id)
.map_err(at(position))?,
),
denial_reason: Some(Text::new(&denied.denial_reason)),
allowed,
verdict: role_verdict(denied.signature_status),
}
}
};
self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
payload,
)?;
Ok(())
}
fn fold_summary(&mut self, ctx: &Ctx<'_>, event: &Event) -> Result<(), ConversationTraceError> {
let Ctx {
position,
kind_base,
trust: step_trust,
turn_id,
..
} = ctx;
let position = *position;
match crate::fold_summary_event(&event.payload) {
Ok(fact) => {
self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::Summary {
text: Text::new(&fact.text),
covers_through_position: fact.covers_through_position,
},
)?;
}
Err(error) => self.warn(
position,
WarningCode::DecodeFailed,
&format!("summary did not decode: {error}"),
),
}
Ok(())
}
fn fold_summary_gate(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
base,
kind_base,
trust: step_trust,
turn_id,
} = ctx;
let position = *position;
let base: &str = base;
let rejected = base == kinds::SUMMARY_GATE_REJECTED;
let decoded = if rejected {
crate::fold_summary_gate_rejected_event(&event.payload)
} else {
crate::fold_summary_gate_admitted_event(&event.payload)
};
match decoded {
Ok(fact) => {
let (dropped_identifiers, clipped) =
bounded_id_list("dropped_identifiers", fact.dropped_identifiers)
.map_err(at(position))?;
if clipped {
self.warn(
position,
WarningCode::ListClipped,
"a summary gate named more identifiers than the bound admits",
);
}
self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::SummaryGate {
rejected,
dropped_identifiers,
count: u64::from(fact.count),
},
)?;
}
Err(error) => self.warn(
position,
WarningCode::DecodeFailed,
&format!("{base} did not decode: {error}"),
),
}
Ok(())
}
fn fold_identity(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
kind_base,
trust: step_trust,
turn_id,
..
} = ctx;
let position = *position;
let Some(decoded) = polyc_proto::events_decode::decode_event_payload::<
polyc_proto::proto::polychrome::events::v1::ObservedIdentityEvent,
>(&event.payload) else {
self.warn(
position,
WarningCode::DecodeFailed,
"observed_identity did not decode",
);
return Ok(());
};
let identity = decoded.identity.into_option().unwrap_or_default();
self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::Identity {
provider: Id::new("provider", &identity.provider).map_err(at(position))?,
scope: Id::new("scope", &identity.scope).map_err(at(position))?,
external_id: Id::new("external_id", &identity.external_id).map_err(at(position))?,
display_name: Text::new(&identity.display_name),
persona_id: Id::new("persona_id", &decoded.persona_id).map_err(at(position))?,
},
)?;
Ok(())
}
fn fold_grant_replay(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
trust: &TraceTrust,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
kind_base,
trust: step_trust,
turn_id,
..
} = ctx;
let position = *position;
let Some(fact) = crate::fold_grant_replay_event(event, &trust.approval) else {
self.warn(
position,
WarningCode::DecodeFailed,
"grant_replay did not decode",
);
return Ok(());
};
let (covered_capabilities, clipped) =
bounded_id_list("covered_capabilities", fact.covered_capabilities)
.map_err(at(position))?;
if clipped {
self.warn(
position,
WarningCode::ListClipped,
"a grant replay covered more capabilities than the bound admits",
);
}
let verdict = match fact.signature_status {
crate::GrantReplaySignatureStatus::Verified => RoleVerdict::Verified,
crate::GrantReplaySignatureStatus::Invalid => RoleVerdict::Invalid,
};
let signer_key = (fact.signature_status == crate::GrantReplaySignatureStatus::Verified)
.then(|| fact.signer_public_key.clone())
.flatten();
let step = self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::GrantReplay {
tool: Id::new("tool", &fact.tool).map_err(at(position))?,
grant_ref: Id::new("grant_ref", &fact.grant_ref).map_err(at(position))?,
covered_capabilities,
coverage_hash: bounded_hash("coverage_hash", &fact.coverage_hash)
.map_err(at(position))?,
turn_id: (!turn_id.is_empty()).then(|| turn_id.clone()),
verdict,
},
)?;
if let Some(key) = signer_key {
self.record_signer(step, SignerRole::GrantReplay, &key);
}
Ok(())
}
}
const fn role_verdict(status: polyc_crypto::signing_role::SignatureVerdict) -> RoleVerdict {
match status {
polyc_crypto::signing_role::SignatureVerdict::Verified => RoleVerdict::Verified,
polyc_crypto::signing_role::SignatureVerdict::Invalid => RoleVerdict::Invalid,
polyc_crypto::signing_role::SignatureVerdict::Untrusted => RoleVerdict::Untrusted,
}
}
impl Scan {
fn fold_marker_or_unknown(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
base,
kind_base,
trust: step_trust,
turn_id,
} = ctx;
let position = *position;
let base: &str = base;
let Some(marker) = marker_kind(base) else {
self.warn(
position,
WarningCode::UnknownKind,
&format!("`{base}` is outside the closed registry"),
);
self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::Unknown,
)?;
return Ok(());
};
let mut on = None;
let mut escalated_at_ms = None;
if base == kinds::INCOGNITO_SET {
match crate::fold_incognito_set(&event.payload) {
Some(decoded) => on = Some(decoded.on),
None => self.warn(
position,
WarningCode::DecodeFailed,
"incognito_set did not decode",
),
}
}
if base == kinds::CONVERSATION_MULTIPARTY_FLOOR {
match polyc_proto::events_decode::decode_event_payload::<
polyc_proto::proto::polychrome::events::v1::ConversationMultipartyFloorEvent,
>(&event.payload)
{
Some(decoded) => escalated_at_ms = Some(decoded.escalated_at_ms),
None => self.warn(
position,
WarningCode::DecodeFailed,
"conversation_multiparty_floor did not decode",
),
}
}
self.push(
position,
turn_id.clone(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::Marker {
marker,
on,
escalated_at_ms,
},
)?;
Ok(())
}
fn resolve_tool_calls(&mut self) {
let open_approvals: HashSet<(String, String)> =
self.open_approvals.keys().cloned().collect();
let open_question_calls: HashSet<(String, String)> = self
.open_questions
.keys()
.map(|(turn, call, _)| (turn.clone(), call.clone()))
.collect();
let mut settled: Vec<(u64, u64, ToolStatus, Option<Id>)> = Vec::new();
for step in &self.steps {
let TraceStepPayload::ToolCall {
tool_call_id,
origin_turn_id,
..
} = &step.payload
else {
continue;
};
let turn = origin_turn_id.as_str().to_owned();
let call = tool_call_id.as_str().to_owned();
let result = self
.result_sites
.iter()
.filter(|(position, _, id)| {
id.as_str() == call.as_str() && *position >= step.position
})
.min_by_key(|(position, _, _)| *position);
let status = if self.denied_calls.contains(&(turn.clone(), call.clone())) {
ToolStatus::Denied
} else if result.is_some() {
ToolStatus::Completed
} else if open_approvals.contains(&(turn.clone(), call.clone())) {
ToolStatus::WaitingApproval
} else if open_question_calls.contains(&(turn.clone(), call.clone())) {
ToolStatus::WaitingQuestion
} else {
ToolStatus::MissingResult
};
settled.push((
step.position,
step.ordinal,
status,
result.map(|(_, turn, _)| turn.clone()),
));
}
let mut settled = settled.into_iter();
for step in &mut self.steps {
if !matches!(step.payload, TraceStepPayload::ToolCall { .. }) {
continue;
}
let Some((position, ordinal, resolved, result_turn)) = settled.next() else {
break;
};
debug_assert_eq!(
(position, ordinal),
(step.position, step.ordinal),
"the two walks must visit the same tool-call steps in the same order"
);
if let TraceStepPayload::ToolCall {
status,
result_turn_id,
origin_turn_id,
..
} = &mut step.payload
{
*status = resolved;
*result_turn_id =
result_turn.filter(|turn| turn.as_str() != origin_turn_id.as_str());
}
}
}
fn resolve_payment_turns(&mut self) {
let mut unresolved = Vec::new();
for step in &mut self.steps {
if !step.turn_id.is_empty() {
continue;
}
let call = match &step.payload {
TraceStepPayload::PaymentAttempt { tool_call_id, .. }
| TraceStepPayload::PaymentReceipt { tool_call_id, .. }
| TraceStepPayload::PaymentRefusal { tool_call_id, .. } => tool_call_id.clone(),
_ => continue,
};
let owner = self
.calls
.iter()
.filter(|candidate| {
candidate.tool_call_id.as_str() == call.as_str()
&& candidate.position <= step.position
&& !candidate.turn_id.is_empty()
})
.max_by_key(|candidate| candidate.position);
match owner {
Some(call) => step.turn_id = call.turn_id.clone(),
None => unresolved.push(step.position),
}
}
for position in unresolved {
self.warn(
position,
WarningCode::UnresolvedPayment,
"a payment record names a tool call this conversation did not make",
);
}
}
fn build_turns(&self) -> Vec<TraceTurnFact> {
let open_approval_turns: HashSet<&str> = self
.open_approvals
.values()
.map(super::Id::as_str)
.collect();
let open_question_turns: HashSet<&str> = self
.open_questions
.values()
.map(super::Id::as_str)
.collect();
self.turn_order
.iter()
.filter_map(|turn| {
let markers = self.markers.get(turn)?;
let turn_id = Id::new("turn_id", turn).ok()?;
let has_open_approval = open_approval_turns.contains(turn.as_str());
let has_open_question = open_question_turns.contains(turn.as_str());
let status = derive_turn_status(markers, has_open_approval, has_open_question);
let mut evidence_positions = markers.evidence.clone();
evidence_positions.sort_unstable();
evidence_positions.dedup();
Some(TraceTurnFact {
turn_id,
first_position: markers.first_position,
status,
has_start: markers.has_start,
has_dispatched: markers.has_dispatched,
has_complete: markers.has_complete,
has_failed: markers.has_failed,
has_ambiguous: markers.has_ambiguous,
failed_kind: markers.failed_kind.clone(),
failed_message: markers.failed_message.clone(),
failed_position: markers.failed_position,
model_provider: markers.model_provider.clone(),
model_name: markers.model_name.clone(),
system_config_ref: markers.system_config_ref.clone(),
temperature_micros: markers.temperature_micros,
top_p_micros: markers.top_p_micros,
max_tokens: markers.max_tokens,
reasoning_level: markers.reasoning_level.clone(),
dispatch_clock_ms: markers.dispatch_clock_ms,
usage_input: markers.usage_input,
usage_output: markers.usage_output,
caller_persona_id: self.callers.get(turn).cloned().unwrap_or_default(),
plannable: markers.has_dispatched,
evidence_positions,
})
})
.collect()
}
}
fn marker_kind(base: &str) -> Option<MarkerKind> {
Some(match base {
kinds::INCOGNITO_SET => MarkerKind::IncognitoSet,
kinds::CONVERSATION_MULTIPARTY_FLOOR => MarkerKind::MultipartyFloor,
kinds::ADMISSION_REFUSED => MarkerKind::AdmissionRefused,
kinds::TOOL_INPUT_REWRITE => MarkerKind::ToolInputRewrite,
kinds::TOOL_CONTEXT_INJECTION => MarkerKind::ToolContextInjection,
kinds::TOOL_RESULT_REDACTION => MarkerKind::ToolResultRedaction,
kinds::CALLER => MarkerKind::Caller,
kinds::PARTICIPANT => MarkerKind::Participant,
kinds::ENROLLMENT => MarkerKind::Enrollment,
kinds::ROUTINE_GRANT => MarkerKind::RoutineGrant,
kinds::GRANT_REVOCATION => MarkerKind::GrantRevocation,
kinds::GRANT_SUSPENSION => MarkerKind::GrantSuspension,
kinds::NUDGE_SENT => MarkerKind::NudgeSent,
kinds::TAINT_EXCISION => MarkerKind::TaintExcision,
kinds::CONVERSATION_NAMESPACE => MarkerKind::ConversationNamespace,
kinds::NAMESPACE_REFUSED => MarkerKind::NamespaceRefused,
_ => return None,
})
}
const fn derive_turn_status(
markers: &TurnMarkers,
has_open_approval: bool,
has_open_question: bool,
) -> TurnStatus {
if markers.has_failed {
return TurnStatus::Failed;
}
if markers.has_ambiguous && !markers.has_complete {
return TurnStatus::Quarantined;
}
if markers.has_dispatched && !markers.has_complete {
return TurnStatus::OrphanedDispatch;
}
if has_open_approval {
return TurnStatus::WaitingApproval;
}
if has_open_question {
return TurnStatus::WaitingQuestion;
}
if markers.has_complete {
return TurnStatus::Complete;
}
TurnStatus::Unknown
}
fn rank_failures(turns: &[TraceTurnFact], scan: &Scan) -> Vec<TraceFailureFact> {
let mut failures = Vec::new();
for turn in turns {
rank_turn_failures(turn, scan, &mut failures);
}
failures
}
const fn status_failure(turn: &TraceTurnFact) -> Option<(FailureKind, &'static str)> {
match turn.status {
TurnStatus::Failed => Some((FailureKind::TurnFailed, "")),
TurnStatus::Quarantined => Some((
FailureKind::Quarantined,
"the turn recorded an ambiguous outcome and never completed",
)),
TurnStatus::OrphanedDispatch => Some((
FailureKind::OrphanedDispatch,
"the turn dispatched and never completed",
)),
TurnStatus::WaitingApproval
| TurnStatus::WaitingQuestion
| TurnStatus::Complete
| TurnStatus::Unknown => None,
}
}
fn ids_for_turn<'a>(
keys: impl Iterator<Item = (&'a str, &'a str)>,
turn: &Id,
field: &'static str,
) -> Vec<Id> {
let mut ids: Vec<Id> = keys
.filter(|(owner, _)| *owner == turn.as_str())
.map(|(_, key)| Id::new(field, key).unwrap_or_default())
.collect();
ids.sort_by(|left, right| left.as_str().cmp(right.as_str()));
ids.dedup_by(|left, right| left.as_str() == right.as_str());
ids
}
#[allow(
clippy::too_many_lines,
reason = "one block per failure signal; splitting hides which signals a turn can carry"
)]
fn rank_turn_failures(turn: &TraceTurnFact, scan: &Scan, failures: &mut Vec<TraceFailureFact>) {
let mut rank = 0;
let mut push = |kind: FailureKind, message: &str, evidence: Vec<u64>, ids: Vec<Id>| {
rank += 1;
failures.push(TraceFailureFact {
turn_id: turn.turn_id.clone(),
rank,
kind,
message: Text::new(message),
evidence_positions: evidence,
related_ids: ids,
});
};
if let Some((kind, message)) = status_failure(turn) {
let (message, evidence, ids) = if kind == FailureKind::TurnFailed {
(
turn.failed_message.as_str(),
turn.failed_position.map(|p| vec![p]).unwrap_or_default(),
vec![turn.failed_kind.clone()],
)
} else {
(message, Vec::new(), Vec::new())
};
push(kind, message, evidence, ids);
}
let open_approvals = ids_for_turn(
scan.open_approvals
.keys()
.map(|(owner, key)| (owner.as_str(), key.as_str())),
&turn.turn_id,
"request_id",
);
if !open_approvals.is_empty() {
push(
FailureKind::WaitingApproval,
"an approval is still open",
Vec::new(),
open_approvals,
);
}
let open_questions = ids_for_turn(
scan.open_questions
.keys()
.map(|(owner, key, _)| (owner.as_str(), key.as_str())),
&turn.turn_id,
"call_id",
);
if !open_questions.is_empty() {
push(
FailureKind::WaitingQuestion,
"a question is still open",
Vec::new(),
open_questions,
);
}
let calls_in_turn = |wanted: ToolStatus| -> Vec<(&Id, u64)> {
scan.steps
.iter()
.filter_map(|step| match &step.payload {
TraceStepPayload::ToolCall {
tool_call_id,
status,
origin_turn_id,
..
} if *status == wanted && origin_turn_id.as_str() == turn.turn_id.as_str() => {
Some((tool_call_id, step.position))
}
_ => None,
})
.collect()
};
let mut denied: Vec<Id> = calls_in_turn(ToolStatus::Denied)
.into_iter()
.map(|(id, _)| id.clone())
.collect();
let unattended: Vec<Id> = scan
.denied_calls
.iter()
.filter(|(owner, request)| {
owner.as_str() == turn.turn_id.as_str()
&& !denied.iter().any(|id| id.as_str() == request.as_str())
})
.filter_map(|(_, request)| Id::new("request_id", request).ok())
.collect();
denied.extend(unattended);
denied.sort_by(|left, right| left.as_str().cmp(right.as_str()));
if !denied.is_empty() {
push(
FailureKind::DeniedApproval,
"an approval denied a tool call",
Vec::new(),
denied,
);
}
let failed_subagents: Vec<Id> = scan
.subagent_failures
.iter()
.filter(|(owner, _)| owner.as_str() == turn.turn_id.as_str())
.map(|(_, id)| id.clone())
.collect();
if !failed_subagents.is_empty() {
push(
FailureKind::SubagentFailed,
"a subagent reported failure",
Vec::new(),
failed_subagents,
);
}
let missing = calls_in_turn(ToolStatus::MissingResult);
if !missing.is_empty() {
push(
FailureKind::MissingResult,
"a tool call has no recorded result",
missing.iter().map(|(_, position)| *position).collect(),
missing.iter().map(|(id, _)| (*id).clone()).collect(),
);
}
let wallet: Vec<Id> = scan
.wallet_link_calls
.iter()
.filter(|(owner, _)| owner.as_str() == turn.turn_id.as_str())
.map(|(_, call)| call.clone())
.collect();
if !wallet.is_empty() {
push(
FailureKind::WalletLinkNeeded,
"a tool result asked for a wallet link",
Vec::new(),
wallet,
);
}
let settled: HashSet<(&str, &str)> = scan
.steps
.iter()
.filter_map(|step| match &step.payload {
TraceStepPayload::PaymentReceipt { tool_call_id, .. } => {
Some((step.turn_id.as_str(), tool_call_id.as_str()))
}
_ => None,
})
.collect();
let unsettled: Vec<Id> = scan
.steps
.iter()
.filter_map(|step| match &step.payload {
TraceStepPayload::PaymentAttempt { tool_call_id, .. }
if step.turn_id.as_str() == turn.turn_id.as_str()
&& !settled.contains(&(turn.turn_id.as_str(), tool_call_id.as_str())) =>
{
Some(tool_call_id.clone())
}
_ => None,
})
.collect();
if !unsettled.is_empty() {
push(
FailureKind::PaymentAttemptUnsettled,
"a payment attempt has no matching receipt",
Vec::new(),
unsettled,
);
}
}
impl Scan {
fn fold_payment_attempt(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
kind_base,
trust: step_trust,
..
} = ctx;
let position = *position;
let Ok(value) = serde_json::from_slice::<serde_json::Value>(&event.payload) else {
self.warn(
position,
WarningCode::DecodeFailed,
"outbound_payment_attempt did not decode",
);
return Ok(());
};
let field = |name: &str| {
value
.get(name)
.and_then(serde_json::Value::as_str)
.unwrap_or_default()
.to_owned()
};
self.push(
position,
Id::default(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::PaymentAttempt {
approval_pos: field("approval_pos").parse().ok(),
tool_call_id: Id::new("tool_call_id", &field("tool_call_id"))
.map_err(at(position))?,
approved_args_hash: bounded_hash(
"approved_args_hash",
&field("approved_args_hash"),
)
.map_err(at(position))?,
url: Text::new(&field("url")),
currency: Id::new("currency", &field("currency")).map_err(at(position))?,
challenge_amount: Text::new(&field("challenge_amount")),
challenge_id: Id::new("challenge_id", &field("challenge_id"))
.map_err(at(position))?,
},
)?;
Ok(())
}
fn fold_payment_receipt(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
trust: &TraceTrust,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
base,
kind_base,
trust: step_trust,
..
} = ctx;
let position = *position;
let Some(receipt) =
crate::verified_receipts(std::slice::from_ref(event), base, &trust.approval).next()
else {
self.warn(
position,
WarningCode::SignatureUnverified,
&format!("{base} did not verify against the approval trust set"),
);
return Ok(());
};
let tool_call_id = Id::new("tool_call_id", &receipt.tool_call_id).map_err(at(position))?;
let direction = if *base == kinds::PAYMENT_RECEIPT {
PaymentDirection::Inbound
} else {
PaymentDirection::Outbound
};
let step = self.push(
position,
Id::default(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::PaymentReceipt {
direction,
reference: Id::new("reference", &receipt.reference).map_err(at(position))?,
amount: Text::new(&receipt.amount),
currency: Id::new("currency", &receipt.currency).map_err(at(position))?,
recipient: Text::new(&receipt.recipient),
method: Id::new("method", &receipt.method).map_err(at(position))?,
tool_call_id,
approval_pos: receipt.approval_pos.parse().ok(),
approved_args_hash: bounded_hash("approved_args_hash", &receipt.approved_args_hash)
.map_err(at(position))?,
subject: Id::new("subject", &receipt.subject).map_err(at(position))?,
kind: Id::new("kind", &receipt.kind).map_err(at(position))?,
payer_kind: Id::new("payer_kind", &receipt.payer_kind).map_err(at(position))?,
paying_account: Id::new("paying_account", &receipt.paying_account)
.map_err(at(position))?,
},
)?;
self.record_signer(step, SignerRole::Payment, &receipt.signer_public_key);
Ok(())
}
fn fold_payment_refusal(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
trust: &TraceTrust,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
kind_base,
trust: step_trust,
..
} = ctx;
let position = *position;
let Some(refusal) =
crate::verified_refusals(std::slice::from_ref(event), &trust.approval).next()
else {
self.warn(
position,
WarningCode::SignatureUnverified,
"payment_refusal did not verify against the approval trust set",
);
return Ok(());
};
let step = self.push(
position,
Id::default(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::PaymentRefusal {
reason: Id::new("reason", &refusal.reason).map_err(at(position))?,
reason_detail: Text::new(&refusal.reason_detail),
merchant_host: Text::new(&refusal.merchant_host),
requested_base_units: refusal.requested_base_units.parse().unwrap_or_default(),
permitted_base_units: refusal.permitted_base_units.parse().unwrap_or_default(),
tool_call_id: Id::new("tool_call_id", &refusal.tool_call_id)
.map_err(at(position))?,
subject: Id::new("subject", &refusal.subject).map_err(at(position))?,
timestamp: refusal.timestamp.parse().unwrap_or_default(),
},
)?;
self.record_signer(step, SignerRole::Payment, &refusal.signer_public_key);
Ok(())
}
fn fold_wallet_link(
&mut self,
ctx: &Ctx<'_>,
event: &Event,
trust: &TraceTrust,
) -> Result<(), ConversationTraceError> {
let Ctx {
position,
kind_base,
trust: step_trust,
..
} = ctx;
let position = *position;
let Some(link) = crate::verified_wallet_link_lifecycle_events(
std::slice::from_ref(event),
&trust.approval,
)
.next() else {
self.warn(
position,
WarningCode::SignatureUnverified,
"wallet_link_lifecycle did not verify against the approval trust set",
);
return Ok(());
};
let (recipients, clipped) = bounded_id_list(
"recipients",
link.recipients
.split(',')
.filter(|value| !value.is_empty())
.map(str::to_owned),
)
.map_err(at(position))?;
if clipped {
self.warn(
position,
WarningCode::ListClipped,
"a wallet link named more recipients than the bound admits",
);
}
let step = self.push(
position,
Id::default(),
kind_base.clone(),
step_trust.clone(),
TraceStepPayload::WalletLink {
transition: Id::new("transition", &link.transition).map_err(at(position))?,
subject: Id::new("subject", &link.subject).map_err(at(position))?,
wallet_address: Id::new("wallet_address", &link.wallet_address)
.map_err(at(position))?,
currency: Id::new("currency", &link.currency).map_err(at(position))?,
chain_id: link.chain_id.parse().unwrap_or_default(),
limit_base_units: link.limit_base_units.parse().unwrap_or_default(),
limit_human: Text::new(&link.limit_human),
period_secs: link.period_secs.parse().unwrap_or_default(),
expiry_unix: link.expiry_unix.parse().unwrap_or_default(),
recipients,
conversation_id: Id::new("conversation_id", &link.conversation_id)
.map_err(at(position))?,
timestamp: link.timestamp.parse().unwrap_or_default(),
},
)?;
self.record_signer(step, SignerRole::WalletLink, &link.signer_public_key);
Ok(())
}
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs, clippy::unwrap_used)]
use super::*;
use polyc_crypto::signing_role::{HandoffRole, RoleSigner, RoleTrustSet, SubagentRole};
use uuid::Uuid;
fn turn() -> Uuid {
Uuid::from_u128(0x7ACE)
}
fn tagged(base: &str) -> String {
kinds::tagged(base, &turn())
}
fn trust() -> TraceTrust {
let handoff = RoleSigner::<HandoffRole>::from_seed(1);
let subagent = RoleSigner::<SubagentRole>::from_seed(2);
TraceTrust {
approval: vec![vec![7u8; 32]],
handoff: RoleTrustSet::from_public_keys(vec![handoff.public_key_bytes()]).unwrap(),
subagent: RoleTrustSet::from_public_keys(vec![subagent.public_key_bytes()]).unwrap(),
question: vec![vec![9u8; 32]],
}
}
fn event(kind: &str) -> Event {
Event::new(kind.to_owned(), Vec::new())
}
const PARTITION: &str = "conv-trace-1";
fn excision_marker(positions: &[u64]) -> Event {
use polyc_crypto::approval::{
ApprovalSigner, EXCISION_SCOPE_SOURCE_ONLY, excision_payload,
};
let (payload, _, _) = excision_payload(
"trace-1",
EXCISION_SCOPE_SOURCE_ONLY,
positions,
"persona-1",
"test excision",
&ApprovalSigner::from_seed(11),
);
Event::new(kinds::TAINT_EXCISION.to_owned(), payload)
}
fn fold(events: Vec<(u64, Event)>) -> ConversationTraceFacts {
fold_conversation_trace(&events, &BTreeMap::new(), &trust()).expect("the fold")
}
#[test]
fn an_uncommitted_turn_still_folds() {
let facts = fold(vec![
(1, event(&tagged(kinds::TURN_START))),
(2, event(&tagged(kinds::TURN_DISPATCHED))),
]);
assert_eq!(facts.turns.len(), 1, "the turn is present without a commit");
assert_eq!(facts.turns[0].status, TurnStatus::OrphanedDispatch);
assert_eq!(facts.steps.len(), 2);
}
#[test]
fn an_ambiguous_turn_is_quarantined_rather_than_orphaned() {
let facts = fold(vec![
(1, event(&tagged(kinds::TURN_START))),
(2, event(&tagged(kinds::TURN_DISPATCHED))),
(3, event(&tagged(kinds::TURN_AMBIGUOUS))),
]);
assert_eq!(facts.turns[0].status, TurnStatus::Quarantined);
assert!(facts.turns[0].has_ambiguous);
assert_eq!(
facts.failures[0].kind,
FailureKind::Quarantined,
"and it ranks as its own failure signal"
);
}
#[test]
fn an_ambiguous_turn_that_completes_is_complete() {
let facts = fold(vec![
(1, event(&tagged(kinds::TURN_START))),
(2, event(&tagged(kinds::TURN_AMBIGUOUS))),
(3, event(&tagged(kinds::TURN_COMPLETE))),
]);
assert_eq!(facts.turns[0].status, TurnStatus::Complete);
}
#[test]
fn the_six_previously_untyped_kinds_are_typed_boundaries() {
let facts = fold(vec![
(1, event(&tagged(kinds::TURN_TEXT_WITHHELD))),
(2, event(&tagged(kinds::STEP_COMMIT))),
(3, event(&tagged(kinds::TURN_AMBIGUOUS))),
(4, event(&tagged(kinds::INGRESS_DIRECTIVE))),
(5, event(&tagged(kinds::CONVERSATION_NAMESPACE))),
(6, event(&tagged(kinds::NAMESPACE_REFUSED))),
]);
let markers: Vec<BoundaryMarker> = facts
.steps
.iter()
.filter_map(|step| match &step.payload {
TraceStepPayload::Boundary { marker, .. } => Some(*marker),
_ => None,
})
.collect();
assert_eq!(
markers,
vec![
BoundaryMarker::TextWithheld,
BoundaryMarker::StepCommit,
BoundaryMarker::Ambiguous,
BoundaryMarker::IngressDirective,
],
"the four boundary-shaped kinds are boundaries"
);
let marker_kinds: Vec<MarkerKind> = facts
.steps
.iter()
.filter_map(|step| match &step.payload {
TraceStepPayload::Marker { marker, .. } => Some(*marker),
_ => None,
})
.collect();
assert_eq!(
marker_kinds,
vec![
MarkerKind::ConversationNamespace,
MarkerKind::NamespaceRefused
]
);
assert!(
facts.warnings.is_empty(),
"none of the six is an unknown kind any more: {:?}",
facts.warnings
);
}
#[test]
fn an_unknown_kind_is_a_spine_row_with_no_payload() {
let mut unknown = event("a_kind_from_a_later_schema");
unknown.payload = b"{\"secret\":\"do not copy me\"}".to_vec();
let facts = fold(vec![(1, unknown)]);
assert_eq!(facts.steps.len(), 1);
assert_eq!(facts.steps[0].payload, TraceStepPayload::Unknown);
assert!(facts.steps[0].payload.is_spine_only());
assert_eq!(facts.warnings.len(), 1);
assert_eq!(facts.warnings[0].code, WarningCode::UnknownKind);
assert!(facts.turns.is_empty(), "an unknown kind never mints a turn");
}
#[test]
fn an_excised_position_is_a_tombstone_with_no_body() {
let secret = "the model was told to exfiltrate a key";
let mut events = vec![
(1, event(&tagged(kinds::TURN_START))),
(4, {
let mut poisoned = event(&tagged(kinds::OUTPUT_MSG));
poisoned.payload = secret.as_bytes().to_vec();
poisoned
}),
(9, excision_marker(&[4])),
];
let prepared = crate::prepare_conversation_core(&mut events, PARTITION);
assert_eq!(
prepared.excised.get(&4).copied(),
Some(9),
"the prepare step attributes the strip to the marker at 9"
);
let facts = fold_conversation_trace(&events, &prepared.excised, &trust()).unwrap();
let tombstone = facts
.steps
.iter()
.find(|step| matches!(step.payload, TraceStepPayload::Excised { .. }))
.expect("the stripped position is recorded");
assert_eq!(tombstone.position, 4);
assert_eq!(
tombstone.payload,
TraceStepPayload::Excised {
marker_position: Some(9)
}
);
assert!(tombstone.payload.is_spine_only());
assert_eq!(
tombstone.turn_id.as_str(),
turn().to_string(),
"the stand-in keeps the turn the excised content belonged to"
);
let rendered = format!("{facts:?}");
assert!(
!rendered.contains(secret),
"the excised body must not survive into the facts"
);
}
#[test]
fn an_unattributed_stand_in_is_still_a_tombstone() {
let mut events = vec![
(1, event(&tagged(kinds::TURN_START))),
(4, event(&tagged(kinds::OUTPUT_MSG))),
];
crate::strip_excised(&mut events, &[4].into_iter().collect());
let facts = fold_conversation_trace(&events, &BTreeMap::new(), &trust()).unwrap();
let tombstone = facts
.steps
.iter()
.find(|step| matches!(step.payload, TraceStepPayload::Excised { .. }))
.expect("the stripped position is still recorded");
assert_eq!(tombstone.position, 4);
assert_eq!(
tombstone.payload,
TraceStepPayload::Excised {
marker_position: None
}
);
}
#[test]
fn a_kind_the_fold_models_never_reaches_the_opaque_path() {
let metadata_only = [kinds::MODEL_CALL, kinds::USAGE];
let typed = [
kinds::TURN_START,
kinds::TURN_COMPLETE,
kinds::TURN_FAILED,
kinds::TURN_DISPATCHED,
kinds::TURN_AMBIGUOUS,
kinds::STEP_COMMIT,
kinds::TURN_TEXT_WITHHELD,
kinds::INGRESS_DIRECTIVE,
kinds::USER_MSG,
kinds::OUTPUT_MSG,
kinds::APPROVAL_REQUEST,
kinds::APPROVAL_RESPONSE,
kinds::APPROVAL_DEFERRED,
kinds::QUESTION_REQUEST,
kinds::QUESTION_RESPONSE,
kinds::SUBAGENT_SPAWN,
kinds::SUBAGENT_RESULT,
kinds::SUBAGENT_MODEL_CALL,
kinds::HANDOFF,
kinds::HANDOFF_DENIED,
kinds::SUMMARY,
kinds::SUMMARY_GATE_REJECTED,
kinds::SUMMARY_GATE_ADMITTED,
kinds::OBSERVED_IDENTITY,
kinds::GRANT_REPLAY,
kinds::OUTBOUND_PAYMENT_ATTEMPT,
kinds::PAYMENT_RECEIPT,
kinds::OUTBOUND_PAYMENT_RECEIPT,
kinds::PAYMENT_REFUSAL,
kinds::WALLET_LINK_LIFECYCLE,
];
for (index, base) in typed.iter().enumerate() {
let position = index as u64;
let facts = fold(vec![(position, event(&kinds::tagged(base, &turn())))]);
for step in &facts.steps {
assert!(
!matches!(step.payload, TraceStepPayload::Unknown),
"`{base}` reached the unknown path"
);
assert!(
!matches!(step.payload, TraceStepPayload::Marker { .. }),
"`{base}` reached the marker bucket"
);
}
assert!(
!facts.steps.is_empty()
|| facts.warnings.iter().any(|warning| matches!(
warning.code,
WarningCode::DecodeFailed | WarningCode::SignatureUnverified
)),
"`{base}` produced neither a step nor a warning, so this \
case proved nothing"
);
}
for base in metadata_only {
let facts = fold(vec![(0, event(&kinds::tagged(base, &turn())))]);
assert!(
facts.steps.is_empty(),
"`{base}` fills the turn's columns and mints no step"
);
let mut undecodable = event(&kinds::tagged(base, &turn()));
undecodable.payload = vec![0xFF; 8];
let facts = fold(vec![(0, undecodable)]);
assert!(
facts.steps.iter().all(|step| !matches!(
step.payload,
TraceStepPayload::Unknown | TraceStepPayload::Marker { .. }
)),
"`{base}` reached the opaque path"
);
assert!(
facts
.warnings
.iter()
.any(|warning| warning.code == WarningCode::DecodeFailed),
"`{base}` swallowed a payload it could not decode"
);
}
}
#[test]
fn each_bound_acts_at_the_value_one_past_it() {
use polyc_projection::trace_vocabulary::{
MAX_ID_BYTES, MAX_LIST_ITEM_BYTES, MAX_LIST_ITEMS, MAX_TEXT_BYTES,
};
assert!(
Id::new("field", &"a".repeat(MAX_ID_BYTES)).is_ok(),
"the bound admits its own limit"
);
assert!(
matches!(
Id::new("field", &"a".repeat(MAX_ID_BYTES + 1)),
Err(TraceBoundsError::Identifier { field: "field" })
),
"one byte past the identifier bound is refused, not truncated"
);
let at_limit = Text::new(&"a".repeat(MAX_TEXT_BYTES));
assert!(!at_limit.truncated() && at_limit.as_str().len() == MAX_TEXT_BYTES);
let one_over = Text::new(&"a".repeat(MAX_TEXT_BYTES + 1));
assert!(
one_over.truncated() && one_over.as_str().len() == MAX_TEXT_BYTES,
"one byte past the text bound truncates AND records that it did"
);
let (kept, clipped) =
bounded_id_list("list", (0..MAX_LIST_ITEMS).map(|item| item.to_string()))
.expect("a list at its own limit is admitted");
assert!(!clipped && kept.len() == MAX_LIST_ITEMS);
let (kept, clipped) =
bounded_id_list("list", (0..=MAX_LIST_ITEMS).map(|item| item.to_string()))
.expect("one item past the count bound clips rather than refusing");
assert!(
clipped && kept.len() == MAX_LIST_ITEMS,
"and the clip is reported, so a short list is not read as a whole one"
);
assert!(
bounded_id_list("list", ["a".repeat(MAX_LIST_ITEM_BYTES)]).is_ok(),
"an item at its own limit is admitted"
);
assert!(
matches!(
bounded_id_list("list", ["a".repeat(MAX_LIST_ITEM_BYTES + 1)]),
Err(TraceBoundsError::ListItem { field: "list" })
),
"one byte past the list-item bound refuses; a half id names the \
wrong thing"
);
let mut values: Vec<String> = (0..MAX_LIST_ITEMS).map(|item| item.to_string()).collect();
values.push("a".repeat(MAX_LIST_ITEM_BYTES + 1));
assert!(
bounded_id_list("list", values).is_ok(),
"an oversized item beyond the count bound is dropped, not read"
);
}
#[test]
fn one_prefix_always_folds_to_the_same_facts() {
let turns: Vec<uuid::Uuid> = (0..12).map(uuid::Uuid::from_u128).collect();
let mut events = Vec::new();
let mut position = 0u64;
for id in &turns {
for base in [
kinds::TURN_START,
kinds::TURN_DISPATCHED,
kinds::USER_MSG,
kinds::OUTPUT_MSG,
kinds::APPROVAL_REQUEST,
kinds::TURN_COMPLETE,
] {
events.push((position, event(&kinds::tagged(base, id))));
position += 1;
}
}
let first = fold(events.clone());
for _ in 0..8 {
assert_eq!(
fold(events.clone()),
first,
"the same prefix folded to different facts"
);
}
assert!(
!first.steps.is_empty() && first.turns.len() == turns.len(),
"precondition: the fixture actually produced facts"
);
}
#[test]
fn every_step_is_ordered_and_names_a_position_the_prefix_holds() {
let id = turn();
let events: Vec<(u64, Event)> = [
kinds::TURN_START,
kinds::TURN_DISPATCHED,
kinds::USER_MSG,
kinds::OUTPUT_MSG,
kinds::APPROVAL_REQUEST,
kinds::APPROVAL_RESPONSE,
kinds::SUMMARY,
kinds::OBSERVED_IDENTITY,
kinds::TAINT_EXCISION,
kinds::TURN_COMPLETE,
]
.into_iter()
.enumerate()
.map(|(index, base)| (index as u64 * 7 + 3, event(&kinds::tagged(base, &id))))
.collect();
let held: std::collections::BTreeSet<u64> =
events.iter().map(|(position, _)| *position).collect();
let facts = fold(events);
assert!(!facts.steps.is_empty(), "precondition: steps were produced");
let mut previous = (0u64, 0u64);
for step in &facts.steps {
assert!(
held.contains(&step.position),
"step at position {} is not in the prefix",
step.position
);
let current = (step.position, step.ordinal);
assert!(
current >= previous,
"steps are out of order at {current:?} after {previous:?}"
);
previous = current;
}
for warning in &facts.warnings {
if let Some(position) = warning.position {
assert!(
held.contains(&position),
"a warning names position {position}, which the prefix \
does not hold"
);
}
}
}
fn tool_call_event(id: &str) -> Event {
use buffa::Message as _;
use polyc_proto::proto::polychrome::agent::v1::{
Content, Message, ToolCallContent, content,
};
let message = Message {
role: "model".to_owned(),
content: buffa::MessageField::some(Content {
r#type: Some(content::Type::ToolCall(Box::new(ToolCallContent {
id: id.to_owned(),
..Default::default()
}))),
..Default::default()
}),
..Default::default()
};
let mut event = event(&kinds::tagged(kinds::OUTPUT_MSG, &turn()));
event.payload = message.encode_to_vec();
event
}
fn tool_result_event(id: &str, turn_uuid: &uuid::Uuid) -> Event {
use buffa::Message as _;
use polyc_proto::proto::polychrome::agent::v1::{
Content, Message, ToolResultContent, content,
};
let message = Message {
role: "user".to_owned(),
content: buffa::MessageField::some(Content {
r#type: Some(content::Type::ToolResult(Box::new(ToolResultContent {
call_id: id.to_owned(),
..Default::default()
}))),
..Default::default()
}),
..Default::default()
};
let mut event = event(&kinds::tagged(kinds::OUTPUT_MSG, turn_uuid));
event.payload = message.encode_to_vec();
event
}
fn tool_status_at(facts: &ConversationTraceFacts, position: u64) -> Option<ToolStatus> {
facts.steps.iter().find_map(|step| match &step.payload {
TraceStepPayload::ToolCall { status, .. } if step.position == position => Some(*status),
_ => None,
})
}
#[test]
fn a_tool_call_carries_the_outcome_the_prefix_settles() {
let facts = fold(vec![
(1, event(&tagged(kinds::TURN_START))),
(2, tool_call_event("call-1")),
(3, tool_result_event("call-1", &turn())),
(4, event(&tagged(kinds::TURN_COMPLETE))),
]);
assert_eq!(tool_status_at(&facts, 2), Some(ToolStatus::Completed));
let facts = fold(vec![
(1, event(&tagged(kinds::TURN_START))),
(2, tool_call_event("call-1")),
(3, event(&tagged(kinds::TURN_COMPLETE))),
]);
assert_eq!(tool_status_at(&facts, 2), Some(ToolStatus::MissingResult));
let facts = fold(vec![
(1, event(&tagged(kinds::TURN_START))),
(2, tool_call_event("call-1")),
(3, {
let mut question = event(&tagged(kinds::QUESTION_REQUEST));
question.payload = br#"{"call_id":"call-1","index":0}"#.to_vec();
question
}),
]);
assert_eq!(
tool_status_at(&facts, 2),
Some(ToolStatus::WaitingQuestion),
"an open question on this call outranks a missing result"
);
let facts = fold(vec![
(1, event(&tagged(kinds::TURN_START))),
(2, tool_call_event("call-1")),
(3, {
let mut question = event(&tagged(kinds::QUESTION_REQUEST));
question.payload = br#"{"call_id":"another-call","index":0}"#.to_vec();
question
}),
]);
assert_eq!(
tool_status_at(&facts, 2),
Some(ToolStatus::MissingResult),
"an unrelated open question must not mask a lost result"
);
assert_ne!(tool_status_at(&facts, 2), Some(ToolStatus::Unknown));
}
#[test]
fn a_reused_call_id_in_a_later_turn_does_not_settle_the_earlier_one() {
let first = uuid::Uuid::from_u128(0x7ACE_0001);
let second = uuid::Uuid::from_u128(0x7ACE_0002);
let approval = |turn: &uuid::Uuid, position: u64| {
let mut request = event(&kinds::tagged(kinds::APPROVAL_REQUEST, turn));
request.payload =
br#"{"request_id":"call-1","tool_name":"read_file","args_json":"{}"}"#.to_vec();
(position, request)
};
let facts = fold(vec![
(1, event(&kinds::tagged(kinds::TURN_START, &first))),
approval(&first, 2),
(3, event(&kinds::tagged(kinds::TURN_START, &second))),
approval(&second, 4),
]);
let waiting: Vec<&str> = facts
.turns
.iter()
.filter(|turn| turn.status == TurnStatus::WaitingApproval)
.map(|turn| turn.turn_id.as_str())
.collect();
assert_eq!(
waiting.len(),
2,
"both turns opened an approval and neither closed one: {:?}",
facts
.turns
.iter()
.map(|turn| (turn.turn_id.as_str(), turn.status))
.collect::<Vec<_>>()
);
}
#[test]
fn a_synthetic_kind_suffix_never_mints_a_turn() {
let real = turn();
let payment_marker = uuid::Uuid::from_u128(0x9A1D_0001);
let stray = uuid::Uuid::from_u128(0x9A1D_0002);
let facts = fold(vec![
(1, event(&kinds::tagged(kinds::TURN_START, &real))),
(2, event(&kinds::tagged(kinds::TURN_COMPLETE, &real))),
(
3,
event(&kinds::tagged(
kinds::OUTBOUND_PAYMENT_RECEIPT,
&payment_marker,
)),
),
(
4,
event(&kinds::tagged("a_kind_from_a_later_schema", &stray)),
),
]);
let minted: Vec<&str> = facts.turns.iter().map(|t| t.turn_id.as_str()).collect();
assert_eq!(
minted,
vec![real.to_string()],
"only the turn-bearing kinds mint a turn"
);
let real_turn = &facts.turns[0];
assert!(
!real_turn.evidence_positions.contains(&3)
&& !real_turn.evidence_positions.contains(&4),
"an unrelated suffix must not become evidence for the real turn: {:?}",
real_turn.evidence_positions
);
}
#[test]
fn an_unmodelled_kind_attaches_to_a_turn_already_minted() {
let real = turn();
let facts = fold(vec![
(1, event(&kinds::tagged(kinds::TURN_START, &real))),
(
2,
event(&kinds::tagged("a_kind_from_a_later_schema", &real)),
),
(3, event(&kinds::tagged(kinds::TURN_COMPLETE, &real))),
]);
assert_eq!(facts.turns.len(), 1);
assert!(
facts.turns[0].evidence_positions.contains(&2),
"the unmodelled position is evidence for the turn it names"
);
let at_two = facts
.steps
.iter()
.find(|step| step.position == 2)
.expect("the position still produces a row");
assert_eq!(at_two.turn_id.as_str(), real.to_string());
}
#[test]
fn a_completed_call_never_also_reports_a_missing_result() {
let first = uuid::Uuid::from_u128(0x7ACE_0011);
let second = uuid::Uuid::from_u128(0x7ACE_0012);
let mut call = tool_call_event("call-1");
call.kind = kinds::tagged(kinds::OUTPUT_MSG, &first);
let facts = fold(vec![
(1, event(&kinds::tagged(kinds::TURN_START, &first))),
(2, call),
(3, event(&kinds::tagged(kinds::TURN_START, &second))),
(4, tool_result_event("call-1", &second)),
]);
let status = facts.steps.iter().find_map(|step| match &step.payload {
TraceStepPayload::ToolCall { status, .. } => Some(*status),
_ => None,
});
assert_eq!(status, Some(ToolStatus::Completed));
assert!(
!facts
.failures
.iter()
.any(|failure| failure.kind == FailureKind::MissingResult),
"a completed call must not also carry a missing-result failure: {:?}",
facts.failures
);
}
#[test]
fn a_signed_answer_closes_the_question_it_names() {
let signer = polyc_crypto::signing_role::ApprovalSigner::from_seed(31);
let mut trust = trust();
trust.question = vec![signer.public_key_bytes()];
let mut request = event(&tagged(kinds::QUESTION_REQUEST));
request.payload = br#"{"call_id":"ask-1","index":0,"header":"Target"}"#.to_vec();
let (payload, _, _) = polyc_crypto::question::answer_payload(
&turn().to_string(),
"ask-1",
0,
"{}",
"declined",
None,
"",
"persona-1",
"trace-1",
"nonce-1",
&signer,
);
let mut response = event(&tagged(kinds::QUESTION_RESPONSE));
response.payload = payload;
let facts = fold_conversation_trace(
&[
(1, event(&tagged(kinds::TURN_START))),
(2, request),
(3, response),
(4, event(&tagged(kinds::TURN_COMPLETE))),
],
&BTreeMap::new(),
&trust,
)
.expect("the prefix folds");
let answered = facts
.steps
.iter()
.find_map(|step| match &step.payload {
TraceStepPayload::Question {
phase: QuestionPhase::Response,
call_id,
index,
decision,
verdict,
..
} => Some((call_id.clone(), *index, decision.clone(), *verdict)),
_ => None,
})
.expect("the response folds to a question step");
assert_eq!(answered.0.as_str(), "ask-1", "the answer names its call");
assert_eq!(answered.1, 0);
assert_eq!(
answered.2,
QuestionDecision::Declined,
"the recorded state is read, not defaulted to `answered`"
);
assert_eq!(answered.3, Some(RoleVerdict::Verified));
let turn_fact = facts.turns.first().expect("one turn");
assert_ne!(
turn_fact.status,
TurnStatus::WaitingQuestion,
"an answered question must not leave its turn waiting"
);
assert!(
!facts
.failures
.iter()
.any(|failure| failure.kind == FailureKind::WaitingQuestion),
"and must not leave a waiting_question failure behind"
);
assert_eq!(
facts
.signers
.iter()
.map(|signer| signer.role)
.collect::<Vec<_>>(),
vec![SignerRole::Question]
);
}
#[test]
fn an_unverified_record_contributes_no_signer() {
let mut forged = event(&tagged(kinds::APPROVAL_RESPONSE));
forged.payload =
br#"{"request_id":"call-1","approved":true,"signed_by":"aa","signature_hex":"00"}"#
.to_vec();
let facts = fold(vec![(1, event(&tagged(kinds::TURN_START))), (2, forged)]);
assert!(
facts.signers.is_empty(),
"a signature this build did not accept names nobody: {:?}",
facts.signers
);
}
#[test]
fn the_question_role_tells_a_forgery_from_an_unknown_signer() {
let signer = polyc_crypto::signing_role::ApprovalSigner::from_seed(21);
let (payload, _, _) = polyc_crypto::question::answer_payload(
&turn().to_string(),
"call-1",
0,
"{}",
"answered",
Some(0),
"yes",
"persona-1",
"trace-1",
"nonce-1",
&signer,
);
assert_eq!(
question_verdict(&payload, &[signer.public_key_bytes()]).0,
RoleVerdict::Verified,
"a good signature from a trusted signer"
);
assert_eq!(
question_verdict(&payload, &[vec![9u8; 32]]).0,
RoleVerdict::Untrusted,
"a good signature from a signer this deployment does not know"
);
assert_eq!(
question_verdict(&payload, &[]).0,
RoleVerdict::Untrusted,
"an empty trust set knows no signer; it does not make a good \
signature bad"
);
let mut decoded: serde_json::Value = serde_json::from_slice(&payload).unwrap();
decoded["selected_label"] = serde_json::json!("no");
let tampered = decoded.to_string().into_bytes();
assert_eq!(
question_verdict(&tampered, &[signer.public_key_bytes()]).0,
RoleVerdict::Invalid,
"the answer was changed after signing"
);
assert_eq!(
question_verdict(b"not json at all", &[signer.public_key_bytes()]).0,
RoleVerdict::Invalid,
"a payload that does not decode is not a signature this build trusts"
);
assert_eq!(
question_verdict(&payload, &[vec![9u8; 32]])
.1
.map(|answer| answer.signer_public_key),
Some(signer.public_key_bytes()),
"an untrusted answer still names its signer"
);
assert!(
question_verdict(&tampered, &[signer.public_key_bytes()])
.1
.is_none(),
"a signature that does not verify names nobody"
);
}
#[test]
fn no_verdict_mapping_collapses_two_statuses_into_one() {
use polyc_crypto::signing_role::SignatureVerdict;
let role = [
role_verdict(SignatureVerdict::Verified),
role_verdict(SignatureVerdict::Invalid),
role_verdict(SignatureVerdict::Untrusted),
];
assert_eq!(
role,
[
RoleVerdict::Verified,
RoleVerdict::Invalid,
RoleVerdict::Untrusted
],
"the handoff and subagent roles carry three distinct verdicts"
);
let approval = [
approval_verdict(crate::ApprovalSignatureStatus::Verified),
approval_verdict(crate::ApprovalSignatureStatus::Invalid),
approval_verdict(crate::ApprovalSignatureStatus::LegacyUnverifiable),
];
assert_eq!(
approval,
[
RoleVerdict::Verified,
RoleVerdict::Invalid,
RoleVerdict::LegacyUnverifiable
],
"an approval written before signing is not the same fact as a \
forged one"
);
assert!(
!role.contains(&RoleVerdict::LegacyUnverifiable),
"`LegacyUnverifiable` is approval-only; a handoff or subagent \
record has no pre-binding era to explain"
);
assert!(
!approval.contains(&RoleVerdict::Untrusted),
"the approval status enum has no unknown-signer state, so the \
mapping must not invent one"
);
}
#[test]
fn journal_plumbing_is_skipped_without_a_warning() {
let facts = fold(vec![
(1, event(&tagged(kinds::TURN_START))),
(2, event(polyc_eventlog_model::MMR_SIGNED_ROOT_KIND)),
(3, event(kinds::COMPACTION_CHECKPOINT)),
(4, event(&tagged(kinds::TURN_COMPLETE))),
]);
assert!(
facts
.steps
.iter()
.all(|step| step.position != 2 && step.position != 3),
"neither plumbing kind mints a step"
);
assert!(
facts.warnings.is_empty(),
"and neither mints a warning: {:?}",
facts.warnings
);
assert_eq!(facts.turns.len(), 1, "the turn around them is unaffected");
}
#[test]
fn a_stripped_position_is_not_an_other_marker() {
let mut events = vec![
(1, event(&tagged(kinds::TURN_START))),
(4, event(&tagged(kinds::OUTPUT_MSG))),
(9, excision_marker(&[4])),
];
let prepared = crate::prepare_conversation_core(&mut events, PARTITION);
let facts = fold_conversation_trace(&events, &prepared.excised, &trust()).unwrap();
let at_four: Vec<&str> = facts
.steps
.iter()
.filter(|step| step.position == 4)
.map(|step| step.payload.step_kind())
.collect();
assert_eq!(at_four, vec!["excised"], "one row, and it is the tombstone");
assert!(
facts.warnings.is_empty(),
"a stripped position is expected, not an unknown kind"
);
}
#[test]
fn a_repeated_source_position_is_refused() {
let error = fold_conversation_trace(
&[
(2, event(&tagged(kinds::TURN_START))),
(2, event(&tagged(kinds::TURN_COMPLETE))),
],
&BTreeMap::new(),
&trust(),
)
.unwrap_err();
assert_eq!(
error,
ConversationTraceError::RepeatedPosition { position: 2 }
);
}
#[test]
fn steps_are_ordered_by_position_then_ordinal() {
let facts = fold(vec![
(5, event(&tagged(kinds::TURN_COMPLETE))),
(1, event(&tagged(kinds::TURN_START))),
(3, event(&tagged(kinds::TURN_DISPATCHED))),
]);
let order: Vec<(u64, u64)> = facts
.steps
.iter()
.map(|step| (step.position, step.ordinal))
.collect();
let mut sorted = order.clone();
sorted.sort_unstable();
assert_eq!(order, sorted);
assert_eq!(order.first(), Some(&(1, 0)));
}
}