use std::collections::HashSet;
use buffa::Message as BuffaMessage;
use polyc_eventlog_model::Event;
use polyc_proto::proto::polychrome::{
agent::v1::Message as WireMessage,
events::v1::{SummaryEvent, TurnDispatchedEvent, TurnFailedEvent},
harness::v1::TurnFailureKind,
};
use polyc_proto::{events_decode::try_decode_event_payload, kinds};
use crate::{
MessageContent, PreparedSource, committed_turn_ids, fold_message_content,
fold_model_call_event, fold_usage_event, prepare_conversation_core,
};
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ConversationExecutionFacts {
pub usage: Vec<UsageHistoryFact>,
pub model_calls: Vec<ModelCallHistoryFact>,
pub tool_calls: Vec<ToolHistoryFact>,
pub failures: Vec<TurnFailureHistoryFact>,
pub summaries: Vec<SummaryHistoryFact>,
pub turn_dispatches: Vec<TurnDispatchHistoryFact>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UsageHistoryFact {
pub position: u64,
pub turn_id: String,
pub input_tokens: u64,
pub output_tokens: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ModelCallHistoryFact {
pub position: u64,
pub turn_id: String,
pub provider: String,
pub model: String,
pub captured_clock_unix_ms: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ToolHistoryKind {
Call,
Result,
}
impl ToolHistoryKind {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Call => "call",
Self::Result => "result",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ToolHistoryFact {
pub position: u64,
pub turn_id: String,
pub tool_call_id: String,
pub kind: ToolHistoryKind,
pub name: String,
pub arguments: Option<String>,
pub result: Option<String>,
pub first_party: Option<bool>,
pub internal_only: bool,
pub trust: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TurnFailureHistoryFact {
pub position: u64,
pub turn_id: String,
pub failure_kind: String,
pub message: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TurnDispatchHistoryFact {
pub position: u64,
pub turn_id: String,
pub occurrence: String,
pub visibility: &'static str,
pub visibility_source: &'static str,
pub source_turn_id: String,
pub edge_asserted_visibility: &'static str,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SummaryHistoryFact {
pub position: u64,
pub summary_id: String,
pub text: String,
pub covers_through_position: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum ConversationExecutionError {
#[error("position {position} appears more than once in one source prefix")]
RepeatedPosition {
position: u64,
},
#[error("the committed {kind} fact at position {position} is malformed")]
Malformed {
kind: &'static str,
position: u64,
},
#[error("the committed tool fact at position {position} is malformed")]
ToolJson {
position: u64,
},
}
#[allow(
clippy::too_many_lines,
reason = "one match arm per kind this family folds; turn_dispatch (POLY-361) pushed it over \
the line budget, not added logic — splitting it would only move the same list"
)]
pub fn fold_conversation_execution(
events: &[(u64, Event)],
) -> Result<ConversationExecutionFacts, ConversationExecutionError> {
refuse_repeated_positions(events)?;
let committed = committed_turn_ids(events.iter().map(|(_, event)| event));
let mut facts = ConversationExecutionFacts::default();
for (position, event) in events {
let (base, turn) = kinds::parse(&event.kind);
if base == kinds::SUMMARY {
let summary = decode::<SummaryEvent>(kinds::SUMMARY, *position, &event.payload)?;
facts.summaries.push(SummaryHistoryFact {
position: *position,
summary_id: turn.map_or_else(String::new, |id| id.to_string()),
text: summary.text,
covers_through_position: summary.covers_through_position,
});
continue;
}
let Some(turn) = committed_turn(turn, &committed) else {
continue;
};
match base {
kinds::USAGE => {
let usage = fold_usage_event(&event.payload)
.map_err(|_| malformed(kinds::USAGE, *position))?;
facts.usage.push(UsageHistoryFact {
position: *position,
turn_id: turn,
input_tokens: usage.input_tokens,
output_tokens: usage.output_tokens,
});
}
kinds::MODEL_CALL => {
let model = fold_model_call_event(&event.payload)
.map_err(|_| malformed(kinds::MODEL_CALL, *position))?;
facts.model_calls.push(ModelCallHistoryFact {
position: *position,
turn_id: turn,
provider: model.provider,
model: model.model,
captured_clock_unix_ms: model.captured_clock_unix_ms,
});
}
kinds::TURN_FAILED => {
let failure =
decode::<TurnFailedEvent>(kinds::TURN_FAILED, *position, &event.payload)?;
facts.failures.push(TurnFailureHistoryFact {
position: *position,
turn_id: turn,
failure_kind: failure_kind(failure.kind).to_owned(),
message: failure.message,
});
}
kinds::TURN_DISPATCHED => {
let dispatch = decode::<TurnDispatchedEvent>(
kinds::TURN_DISPATCHED,
*position,
&event.payload,
)?;
facts.turn_dispatches.push(TurnDispatchHistoryFact {
position: *position,
turn_id: turn,
occurrence: dispatch.occurrence,
visibility: polyc_proto::audience_display::recorded_visibility_label(
dispatch.visibility,
),
visibility_source: polyc_proto::audience_display::recorded_source_label(
dispatch.visibility_source,
),
source_turn_id: dispatch.source_turn_id,
edge_asserted_visibility:
polyc_proto::audience_display::asserted_visibility_label(
dispatch.edge_asserted_visibility,
),
});
}
kinds::USER_MSG | kinds::OUTPUT_MSG => {
let message_kind = if base == kinds::USER_MSG {
kinds::USER_MSG
} else {
kinds::OUTPUT_MSG
};
let message = decode::<WireMessage>(message_kind, *position, &event.payload)?;
let internal_only = message.internal_only;
let folded =
fold_message_content(&message, *position, Some(&turn), event.trust.as_str());
match folded.content {
MessageContent::ToolCall(call) => facts.tool_calls.push(ToolHistoryFact {
position: call.position,
turn_id: turn,
tool_call_id: call.tool_call_id,
kind: ToolHistoryKind::Call,
name: call.name,
arguments: Some(json(&call.arguments, *position)?),
result: None,
first_party: None,
internal_only,
trust: call.trust,
}),
MessageContent::ToolResult(result) => {
facts.tool_calls.push(ToolHistoryFact {
position: result.position,
turn_id: turn,
tool_call_id: result.tool_call_id,
kind: ToolHistoryKind::Result,
name: result.name,
arguments: None,
result: Some(json(&result.result, *position)?),
first_party: Some(result.first_party),
internal_only,
trust: result.trust,
});
}
MessageContent::Text(_) | MessageContent::None => {}
}
}
_ => {}
}
}
Ok(facts)
}
pub fn prepare_and_fold_conversation_execution(
events: &mut [(u64, Event)],
partition: &str,
) -> Result<(PreparedSource, ConversationExecutionFacts), ConversationExecutionError> {
let prepared = prepare_conversation_core(events, partition);
let facts = fold_conversation_execution(events)?;
Ok((prepared, facts))
}
fn committed_turn(turn: Option<uuid::Uuid>, committed: &HashSet<uuid::Uuid>) -> Option<String> {
turn.filter(|turn| committed.contains(turn))
.map(|turn| turn.to_string())
}
fn refuse_repeated_positions(events: &[(u64, Event)]) -> Result<(), ConversationExecutionError> {
let mut positions = HashSet::with_capacity(events.len());
for (position, _) in events {
if !positions.insert(*position) {
return Err(ConversationExecutionError::RepeatedPosition {
position: *position,
});
}
}
Ok(())
}
fn decode<T: BuffaMessage + Default>(
kind: &'static str,
position: u64,
payload: &[u8],
) -> Result<T, ConversationExecutionError> {
try_decode_event_payload(payload).map_err(|_| malformed(kind, position))
}
const fn malformed(kind: &'static str, position: u64) -> ConversationExecutionError {
ConversationExecutionError::Malformed { kind, position }
}
fn json(value: &serde_json::Value, position: u64) -> Result<String, ConversationExecutionError> {
serde_json::to_string(value).map_err(|_| ConversationExecutionError::ToolJson { position })
}
fn failure_kind(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::Other) => "other",
Some(TurnFailureKind::Unspecified) | None => "unspecified",
}
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs)]
use polyc_proto::proto::polychrome::{
agent::v1::{
Content, FunctionCallContent, FunctionResultContent, Message, TextContent,
ToolCallContent, ToolResultContent, content, function_result_content,
tool_call_content, tool_result_content,
},
events::v1::{
ModelCallEvent, RecordedTurnVisibility, RecordedVisibilitySource, SummaryEvent,
TurnDispatchedEvent, TurnFailedEvent, UsageEvent,
},
};
use serde_json::json;
use super::*;
fn turn() -> uuid::Uuid {
uuid::Uuid::parse_str("01950000-0000-7000-8000-00000000e201").unwrap()
}
fn tagged(base: &str) -> String {
kinds::tagged(base, &turn())
}
fn tool_call(internal_only: bool) -> Message {
let arguments =
serde_json::from_value::<buffa_types::google::protobuf::Struct>(json!({"q": "needle"}))
.map(buffa::MessageField::some)
.unwrap_or_default();
Message {
role: "model".to_owned(),
content: buffa::MessageField::some(Content {
r#type: Some(content::Type::ToolCall(Box::new(ToolCallContent {
id: "call-1".to_owned(),
r#type: Some(tool_call_content::Type::FunctionCall(Box::new(
FunctionCallContent {
name: "search".to_owned(),
arguments,
..Default::default()
},
))),
..Default::default()
}))),
..Default::default()
}),
internal_only,
..Default::default()
}
}
fn tool_result(internal_only: bool) -> Message {
let response =
serde_json::from_value::<buffa_types::google::protobuf::Struct>(json!({"hits": 1}))
.ok()
.map(|value| function_result_content::Result::Response(Box::new(value)));
Message {
role: "tool".to_owned(),
content: buffa::MessageField::some(Content {
r#type: Some(content::Type::ToolResult(Box::new(ToolResultContent {
call_id: "call-1".to_owned(),
first_party: true,
r#type: Some(tool_result_content::Type::FunctionResult(Box::new(
FunctionResultContent {
name: "search".to_owned(),
result: response,
..Default::default()
},
))),
..Default::default()
}))),
..Default::default()
}),
internal_only,
..Default::default()
}
}
#[test]
fn an_uncommitted_fact_never_reaches_the_projection() {
let usage = UsageEvent {
input_tokens: 7,
output_tokens: 3,
..Default::default()
};
let events = vec![
(1, Event::new(tagged(kinds::TURN_START), Vec::new())),
(2, Event::new(tagged(kinds::USAGE), usage.encode_to_vec())),
];
assert!(
fold_conversation_execution(&events)
.unwrap()
.usage
.is_empty()
);
}
#[test]
fn committed_usage_keeps_its_authoritative_position() {
let usage = UsageEvent {
input_tokens: 7,
output_tokens: 3,
..Default::default()
};
let events = vec![
(11, Event::new(tagged(kinds::TURN_START), Vec::new())),
(12, Event::new(tagged(kinds::USAGE), usage.encode_to_vec())),
(13, Event::new(tagged(kinds::TURN_COMPLETE), Vec::new())),
];
let facts = fold_conversation_execution(&events).unwrap();
assert_eq!(facts.usage.len(), 1);
assert_eq!(facts.usage[0].position, 12);
assert_eq!(facts.usage[0].turn_id, turn().to_string());
assert_eq!(
(facts.usage[0].input_tokens, facts.usage[0].output_tokens),
(7, 3)
);
}
#[test]
fn a_malformed_relevant_fact_refuses_the_whole_fold() {
let events = vec![
(1, Event::new(tagged(kinds::TURN_START), Vec::new())),
(2, Event::new(tagged(kinds::USAGE), vec![0xff])),
(3, Event::new(tagged(kinds::TURN_COMPLETE), Vec::new())),
];
assert_eq!(
fold_conversation_execution(&events),
Err(ConversationExecutionError::Malformed {
kind: kinds::USAGE,
position: 2,
})
);
}
#[test]
fn a_repeated_source_position_refuses_the_whole_fold() {
let events = vec![
(1, Event::new(tagged(kinds::TURN_START), Vec::new())),
(1, Event::new(tagged(kinds::TURN_COMPLETE), Vec::new())),
];
assert_eq!(
fold_conversation_execution(&events),
Err(ConversationExecutionError::RepeatedPosition { position: 1 })
);
}
#[test]
fn text_is_not_reclassified_as_a_tool_fact() {
let message = Message {
content: buffa::MessageField::some(Content {
r#type: Some(content::Type::Text(Box::new(TextContent {
text: "plain".to_owned(),
..Default::default()
}))),
..Default::default()
}),
..Default::default()
};
let events = vec![
(1, Event::new(tagged(kinds::TURN_START), Vec::new())),
(
2,
Event::new(tagged(kinds::OUTPUT_MSG), message.encode_to_vec()),
),
(3, Event::new(tagged(kinds::TURN_COMPLETE), Vec::new())),
];
assert!(
fold_conversation_execution(&events)
.unwrap()
.tool_calls
.is_empty()
);
}
#[test]
fn one_committed_turn_preserves_execution_facts_and_tool_classification() {
let model = ModelCallEvent {
provider: "provider-a".to_owned(),
model: "model-a".to_owned(),
captured_clock_unix_ms: 1_700_000_000_123,
..Default::default()
};
let failure = TurnFailedEvent {
kind: TurnFailureKind::Timeout.into(),
message: "deadline".to_owned(),
..Default::default()
};
let events = vec![
(1, Event::new(tagged(kinds::TURN_START), Vec::new())),
(
2,
Event::new(tagged(kinds::MODEL_CALL), model.encode_to_vec()),
),
(
3,
Event::trusted(tagged(kinds::OUTPUT_MSG), tool_call(false).encode_to_vec()),
),
(
4,
Event::new(tagged(kinds::USER_MSG), tool_result(true).encode_to_vec()),
),
(
5,
Event::new(tagged(kinds::TURN_FAILED), failure.encode_to_vec()),
),
(6, Event::new(tagged(kinds::TURN_COMPLETE), Vec::new())),
];
let facts = fold_conversation_execution(&events).unwrap();
assert_eq!(facts.model_calls.len(), 1);
assert_eq!(facts.model_calls[0].position, 2);
assert_eq!(facts.model_calls[0].provider, "provider-a");
assert_eq!(
facts.model_calls[0].captured_clock_unix_ms,
1_700_000_000_123
);
assert_eq!(facts.tool_calls.len(), 2);
assert_eq!(facts.tool_calls[0].kind, ToolHistoryKind::Call);
assert_eq!(
facts.tool_calls[0].arguments.as_deref(),
Some(r#"{"q":"needle"}"#)
);
assert!(!facts.tool_calls[0].internal_only);
assert_eq!(
facts.tool_calls[0].trust,
polyc_eventlog_model::TrustTag::TrustedUser.as_str()
);
assert_eq!(facts.tool_calls[1].kind, ToolHistoryKind::Result);
assert_eq!(
facts.tool_calls[1].result.as_deref(),
Some(r#"{"hits":1.0}"#)
);
assert_eq!(facts.tool_calls[1].first_party, Some(true));
assert!(facts.tool_calls[1].internal_only);
assert_eq!(facts.failures.len(), 1);
assert_eq!(facts.failures[0].failure_kind, "timeout");
assert_eq!(facts.failures[0].message, "deadline");
}
#[test]
fn a_summary_is_independently_committed_fleet_history() {
let summary = SummaryEvent {
text: "older history".to_owned(),
covers_through_position: 41,
..Default::default()
};
let summary_id = uuid::Uuid::parse_str("01950000-0000-7000-8000-00000000e202").unwrap();
let facts = fold_conversation_execution(&[(
42,
Event::new(
kinds::tagged(kinds::SUMMARY, &summary_id),
summary.encode_to_vec(),
),
)])
.unwrap();
assert_eq!(facts.summaries.len(), 1);
assert_eq!(facts.summaries[0].position, 42);
assert_eq!(facts.summaries[0].summary_id, summary_id.to_string());
assert_eq!(facts.summaries[0].covers_through_position, 41);
}
#[test]
fn a_committed_dispatch_keeps_its_occurrence_and_audience_columns() {
let dispatch = TurnDispatchedEvent {
occurrence: "daily-standup-28461600".to_owned(),
visibility: buffa::EnumValue::Known(
RecordedTurnVisibility::RECORDED_TURN_VISIBILITY_DIRECT,
),
visibility_source: buffa::EnumValue::Known(
RecordedVisibilitySource::RECORDED_VISIBILITY_SOURCE_ASSERTED,
),
..Default::default()
};
let events = vec![
(1, Event::new(tagged(kinds::TURN_START), Vec::new())),
(
2,
Event::new(tagged(kinds::TURN_DISPATCHED), dispatch.encode_to_vec()),
),
(3, Event::new(tagged(kinds::TURN_COMPLETE), Vec::new())),
];
let facts = fold_conversation_execution(&events).unwrap();
assert_eq!(facts.turn_dispatches.len(), 1);
let row = &facts.turn_dispatches[0];
assert_eq!(row.position, 2);
assert_eq!(row.turn_id, turn().to_string());
assert_eq!(row.occurrence, "daily-standup-28461600");
assert_eq!(row.visibility, "direct");
assert_eq!(row.visibility_source, "asserted");
assert_eq!(row.source_turn_id, "");
}
#[test]
fn an_uncommitted_dispatch_never_reaches_the_projection() {
let events = vec![
(1, Event::new(tagged(kinds::TURN_START), Vec::new())),
(
2,
Event::new(
tagged(kinds::TURN_DISPATCHED),
TurnDispatchedEvent::default().encode_to_vec(),
),
),
];
assert!(
fold_conversation_execution(&events)
.unwrap()
.turn_dispatches
.is_empty()
);
}
#[test]
fn a_malformed_dispatch_refuses_the_whole_fold() {
let events = vec![
(1, Event::new(tagged(kinds::TURN_START), Vec::new())),
(2, Event::new(tagged(kinds::TURN_DISPATCHED), vec![0xff])),
(3, Event::new(tagged(kinds::TURN_COMPLETE), Vec::new())),
];
assert_eq!(
fold_conversation_execution(&events),
Err(ConversationExecutionError::Malformed {
kind: kinds::TURN_DISPATCHED,
position: 2,
})
);
}
}