use std::collections::{HashMap, HashSet};
use std::sync::{Arc, RwLock};
use std::time::Instant;
use crate::journal::{JournalError, PartitionJournal};
use polyc_eventlog::Event;
use polyc_payments::amount::SettlementDirection;
use polyc_proto::events_decode::{decode_event_payload, try_decode_event_payload};
use polyc_proto::kinds;
use polyc_proto::proto::polychrome::agent::v1::Message;
use polyc_proto::proto::polychrome::events::v1::SummaryEvent;
use polyc_state::feed::FeedRecord;
use uuid::Uuid;
use crate::feed::{PartitionChange, commit_events};
const CONVERSATION_PARTITION_PREFIX: &str = "conv-";
const PREVIEW_BYTES: usize = 120;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DashboardSettlement {
pub persona: String,
pub spend_base_units: u128,
pub charged_base_units: u128,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct DashboardRow {
pub conversation_id: String,
pub total_events: usize,
pub committed_turns: usize,
pub input_tokens: u64,
pub output_tokens: u64,
pub summary_text: Option<String>,
pub last_turn_id: Option<String>,
pub edges: Vec<String>,
pub settlements: Vec<DashboardSettlement>,
pub created_at_ms: Option<u64>,
pub last_activity_ms: Option<u64>,
pub persona_id: Option<String>,
pub first_message_preview: Option<String>,
}
impl DashboardRow {
#[must_use]
pub const fn summary_active(&self) -> bool {
self.summary_text.is_some()
}
}
#[derive(Debug, Clone, Default)]
struct TurnAccum {
started: bool,
completed: bool,
clock_ms: Option<u64>,
caller_id: Option<String>,
first_user_msg_preview: Option<String>,
}
#[derive(Debug, Clone, Default)]
struct TrackedRow {
row: DashboardRow,
applied_position: Option<u64>,
turns: HashMap<Uuid, TurnAccum>,
}
impl TrackedRow {
fn new(conversation_id: String) -> Self {
Self {
row: DashboardRow {
conversation_id,
..Default::default()
},
..Default::default()
}
}
fn apply<'a>(
&mut self,
events: impl IntoIterator<Item = (u64, &'a Event)>,
trusted_signers: &[Vec<u8>],
settlement_decimals: u32,
) {
let new_events: Vec<(u64, &Event)> = events
.into_iter()
.filter(|(position, _)| self.applied_position.is_none_or(|hw| *position > hw))
.collect();
if new_events.is_empty() {
return;
}
let batch: Vec<Event> = new_events.iter().map(|(_, e)| (*e).clone()).collect();
for fact in polyc_facts::attribution_events(&batch, polyc_facts::AttributionScope::Both) {
let provider = fact.identity.map(|i| i.provider).unwrap_or_default();
if !provider.is_empty() && !self.row.edges.contains(&provider) {
self.row.edges.push(provider);
}
}
for (turn, persona_id) in polyc_facts::caller_by_turn_last_wins(
polyc_facts::attribution_events(&batch, polyc_facts::AttributionScope::CallerOnly),
) {
self.turns.entry(turn).or_default().caller_id = Some(persona_id);
}
let receipts = polyc_facts::verified_receipts(
&batch,
kinds::OUTBOUND_PAYMENT_RECEIPT,
trusted_signers,
)
.map(|r| (SettlementDirection::Outbound, r))
.chain(
polyc_facts::verified_receipts(&batch, kinds::PAYMENT_RECEIPT, trusted_signers)
.map(|r| (SettlementDirection::Inbound, r)),
);
for (direction, receipt) in receipts {
let read = match direction {
SettlementDirection::Outbound => {
polyc_payments::amount::read_base_unit_amount(&receipt.amount)
}
SettlementDirection::Inbound => {
polyc_payments::amount::read_dollar_amount(&receipt.amount, settlement_decimals)
}
};
let amount = match read {
Ok(amount) => amount,
Err(error) => {
polyc_payments::amount::record_unreadable_amount(
&polyc_payments::amount::UnreadableAmount {
site: polyc_payments::amount::AmountReadSite::DashboardRollup,
direction,
error,
scope: &self.row.conversation_id,
reference: &receipt.reference,
tool_call_id: &receipt.tool_call_id,
approval_pos: &receipt.approval_pos,
subject: &receipt.subject,
timestamp: &receipt.timestamp,
},
);
continue;
}
};
let persona = if receipt.subject.is_empty() {
"(unattributed)".to_owned()
} else {
receipt.subject
};
let entry = if let Some(existing) = self
.row
.settlements
.iter_mut()
.find(|s| s.persona == persona)
{
existing
} else {
self.row.settlements.push(DashboardSettlement {
persona,
spend_base_units: 0,
charged_base_units: 0,
});
self.row
.settlements
.last_mut()
.expect("just-pushed settlement row")
};
match direction {
SettlementDirection::Outbound => {
entry.spend_base_units = entry.spend_base_units.saturating_add(amount);
}
SettlementDirection::Inbound => {
entry.charged_base_units = entry.charged_base_units.saturating_add(amount);
}
}
}
for (position, event) in &new_events {
self.apply_one(*position, event);
self.applied_position = Some(
self.applied_position
.map_or(*position, |hw| hw.max(*position)),
);
}
}
fn apply_one(&mut self, position: u64, event: &Event) {
let _ = position; self.row.total_events += 1;
let (base, turn_id) = kinds::parse(&event.kind);
if base == kinds::USAGE {
let usage = polyc_facts::fold_usage_event(&event.payload).unwrap_or_default();
self.row.input_tokens = self.row.input_tokens.saturating_add(usage.input_tokens);
self.row.output_tokens = self.row.output_tokens.saturating_add(usage.output_tokens);
}
if base == kinds::SUMMARY
&& let Ok(summary) = try_decode_event_payload::<SummaryEvent>(&event.payload)
&& !summary.text.is_empty()
{
self.row.summary_text = Some(summary.text);
}
if base == kinds::MODEL_CALL
&& let Some(id) = turn_id
&& let Ok(model_call) = polyc_facts::fold_model_call_event(&event.payload)
{
self.turns.entry(id).or_default().clock_ms = Some(model_call.captured_clock_unix_ms);
}
if base == kinds::USER_MSG
&& let Some(id) = turn_id
{
let accum = self.turns.entry(id).or_default();
if accum.first_user_msg_preview.is_none()
&& let Some(msg) = decode_event_payload::<Message>(&event.payload)
&& !msg.internal_only
{
accum.first_user_msg_preview = Some(truncate_preview(&message_preview_text(&msg)));
}
}
if base == kinds::TURN_START
&& let Some(id) = turn_id
{
let now_complete = {
let accum = self.turns.entry(id).or_default();
accum.started = true;
accum.completed
};
if now_complete {
self.commit_turn(id);
}
}
if base == kinds::TURN_COMPLETE
&& let Some(id) = turn_id
{
let now_ready = {
let accum = self.turns.entry(id).or_default();
accum.completed = true;
accum.started
};
if now_ready {
self.commit_turn(id);
}
}
}
fn commit_turn(&mut self, id: Uuid) {
let Some(accum) = self.turns.remove(&id) else {
return;
};
self.row.committed_turns += 1;
self.row.last_turn_id = Some(id.as_simple().to_string());
if self.row.persona_id.is_none()
&& let Some(caller_id) = &accum.caller_id
{
self.row.persona_id = Some(caller_id.clone());
}
if self.row.created_at_ms.is_none()
&& let Some(clock) = accum.clock_ms
&& clock != 0
{
self.row.created_at_ms = Some(clock);
}
if let Some(clock) = accum.clock_ms.filter(|&c| c != 0) {
self.row.last_activity_ms = Some(
self.row
.last_activity_ms
.map_or(clock, |prev| prev.max(clock)),
);
}
if self.row.first_message_preview.is_none()
&& let Some(preview) = accum.first_user_msg_preview
{
self.row.first_message_preview = Some(preview);
}
}
}
fn message_preview_text(msg: &Message) -> String {
match polyc_facts::fold_message_content(msg, 0, None, "").content {
polyc_facts::MessageContent::Text(t) => t.text,
polyc_facts::MessageContent::ToolCall(c) => {
let headline = if c.name.is_empty() {
format!("tool · {}", c.tool_call_id)
} else {
format!("tool · {} · {}", c.name, c.tool_call_id)
};
format_tool_block(&headline, &c.arguments)
}
polyc_facts::MessageContent::ToolResult(r) => {
let headline = if r.name.is_empty() {
format!("result · {}", r.tool_call_id)
} else {
format!("result · {} · {}", r.name, r.tool_call_id)
};
format_tool_block(&headline, &r.result)
}
polyc_facts::MessageContent::None => String::new(),
}
}
fn format_tool_block(headline: &str, value: &serde_json::Value) -> String {
if value.is_null()
|| value.as_object().is_some_and(serde_json::Map::is_empty)
|| value.as_array().is_some_and(Vec::is_empty)
{
return headline.to_owned();
}
let pretty = serde_json::to_string_pretty(value).unwrap_or_else(|_| value.to_string());
format!("{headline}\n{pretty}")
}
fn truncate_preview(text: &str) -> String {
let trimmed = text.trim();
if trimmed.len() <= PREVIEW_BYTES {
return trimmed.to_owned();
}
let mut end = PREVIEW_BYTES;
while !trimmed.is_char_boundary(end) && end > 0 {
end -= 1;
}
format!("{}…", &trimmed[..end])
}
fn conversation_id_from_partition(partition: &str) -> Option<String> {
partition
.strip_prefix(CONVERSATION_PARTITION_PREFIX)
.filter(|id| !id.is_empty())
.map(str::to_owned)
}
fn merge_row(rows: &mut HashMap<String, TrackedRow>, id: String, incoming: TrackedRow) {
match rows.get(&id) {
Some(live) if live.applied_position >= incoming.applied_position => {
}
_ => {
rows.insert(id, incoming);
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RebuildStatus {
Pending,
Complete {
conversation_count: usize,
duration_ms: u64,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct RebuildSummary {
pub conversation_count: usize,
pub duration_ms: u64,
}
pub struct DashboardProjection {
trusted_signers: Vec<Vec<u8>>,
settlement_decimals: u32,
journal: Arc<dyn PartitionJournal>,
rows: RwLock<HashMap<String, TrackedRow>>,
rebuild_status: RwLock<RebuildStatus>,
#[cfg(test)]
race_hook: RaceHook,
}
#[cfg(test)]
#[derive(Default)]
struct RaceHook {
armed: std::sync::atomic::AtomicBool,
reached: tokio::sync::Notify,
resume: tokio::sync::Notify,
}
#[cfg(test)]
impl RaceHook {
fn arm(&self) {
self.armed.store(true, std::sync::atomic::Ordering::SeqCst);
}
async fn wait_for_pause(&self) {
self.reached.notified().await;
}
fn resume(&self) {
self.resume.notify_one();
}
async fn pause_if_armed(&self) {
if self.armed.load(std::sync::atomic::Ordering::SeqCst) {
self.reached.notify_one();
self.resume.notified().await;
}
}
}
pub type DashboardCell = Arc<DashboardProjection>;
impl DashboardProjection {
#[must_use]
pub fn new(
trusted_signers: Vec<Vec<u8>>,
journal: Arc<dyn PartitionJournal>,
settlement_decimals: u32,
) -> DashboardCell {
Arc::new(Self {
trusted_signers,
settlement_decimals,
journal,
rows: RwLock::new(HashMap::new()),
rebuild_status: RwLock::new(RebuildStatus::Pending),
#[cfg(test)]
race_hook: RaceHook::default(),
})
}
#[must_use]
pub fn rebuild_status(&self) -> RebuildStatus {
*self.rebuild_status.read().expect("poison")
}
#[must_use]
pub fn row(&self, conversation_id: &str) -> Option<DashboardRow> {
self.rows
.read()
.expect("poison")
.get(conversation_id)
.map(|tracked| tracked.row.clone())
}
#[must_use]
pub fn rows(&self) -> Vec<DashboardRow> {
let mut rows: Vec<DashboardRow> = self
.rows
.read()
.expect("poison")
.values()
.map(|tracked| tracked.row.clone())
.collect();
rows.sort_by(|a, b| a.conversation_id.cmp(&b.conversation_id));
rows
}
#[must_use]
pub fn covers(&self, partition: &str, position: u64) -> bool {
let Some(id) = conversation_id_from_partition(partition) else {
return false;
};
self.rows
.read()
.expect("poison")
.get(&id)
.and_then(|tracked| tracked.applied_position)
.is_some_and(|applied| applied >= position)
}
pub async fn bootstrap_partition(&self, partition: &str) -> Result<(), JournalError> {
self.replay_one_partition(partition).await
}
pub async fn rebuild_from_full_fleet(&self) -> Result<RebuildSummary, JournalError> {
let start = Instant::now();
let partitions = self.journal.list_partitions().await?;
let applied_at_listing: HashMap<String, Option<u64>> = self
.rows
.read()
.expect("poison")
.iter()
.map(|(id, tracked)| (id.clone(), tracked.applied_position))
.collect();
let mut fresh: HashMap<String, TrackedRow> = HashMap::new();
let mut unread: HashSet<String> = HashSet::new();
let mut conversation_count = 0usize;
for partition in partitions {
let Some(id) = conversation_id_from_partition(&partition) else {
continue;
};
let Ok(events) = self.journal.replay_with_positions(partition).await else {
unread.insert(id);
continue;
};
let mut tracked = TrackedRow::new(id.clone());
tracked.apply(
events.iter().map(|(position, event)| (*position, event)),
&self.trusted_signers,
self.settlement_decimals,
);
fresh.insert(id, tracked);
conversation_count += 1;
}
#[cfg(test)]
self.race_hook.pause_if_armed().await;
let pruned = {
let mut rows = self.rows.write().expect("poison");
for (id, tracked) in &fresh {
merge_row(&mut rows, id.clone(), tracked.clone());
}
let before = rows.len();
rows.retain(|id, tracked| {
fresh.contains_key(id)
|| unread.contains(id)
|| applied_at_listing.get(id) != Some(&tracked.applied_position)
});
before - rows.len()
};
let duration_ms = u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX);
*self.rebuild_status.write().expect("poison") = RebuildStatus::Complete {
conversation_count,
duration_ms,
};
tracing::info!(
conversation_count,
pruned,
duration_ms,
"dashboard projection rebuilt from a full-fleet replay"
);
Ok(RebuildSummary {
conversation_count,
duration_ms,
})
}
async fn rebuild_one_partition(&self, partition: &str) {
if let Err(error) = self.replay_one_partition(partition).await {
tracing::warn!(
%error,
partition,
"dashboard projection: one-partition rebuild replay failed; the row stays stale \
until the next full-fleet rebuild"
);
}
}
async fn replay_one_partition(&self, partition: &str) -> Result<(), JournalError> {
let Some(id) = conversation_id_from_partition(partition) else {
return Ok(());
};
let events = self
.journal
.replay_with_positions(partition.to_owned())
.await?;
let mut tracked = TrackedRow::new(id.clone());
tracked.apply(
events.iter().map(|(position, event)| (*position, event)),
&self.trusted_signers,
self.settlement_decimals,
);
merge_row(&mut self.rows.write().expect("poison"), id, tracked);
Ok(())
}
fn remove_row(&self, partition: &str) {
if let Some(id) = conversation_id_from_partition(partition) {
self.rows.write().expect("poison").remove(&id);
}
}
pub fn apply_commits(&self, partition: &str, commits: &[FeedRecord]) {
let Some(id) = conversation_id_from_partition(partition) else {
return;
};
let events = commit_events(commits);
if events.is_empty() {
return;
}
self.rows
.write()
.expect("poison")
.entry(id.clone())
.or_insert_with(|| TrackedRow::new(id))
.apply(
events.iter().map(|(position, event)| (*position, event)),
&self.trusted_signers,
self.settlement_decimals,
);
}
pub async fn note_partition_change(&self, partition: &str, change: PartitionChange) {
match change {
PartitionChange::Destroyed | PartitionChange::MigratedAway => {
self.remove_row(partition);
}
PartitionChange::Rewritten => self.rebuild_one_partition(partition).await,
}
}
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs, clippy::unwrap_used)]
use std::path::PathBuf;
use std::sync::Arc;
use buffa::Message as _;
use polyc_crypto::approval::{ApprovalSigner, ReceiptPayload, receipt_payload};
use polyc_eventlog::Event;
use polyc_eventlog_host::EventLogHost;
use polyc_proto::kinds;
use polyc_proto::proto::polychrome::agent::v1::{
Content, Message as WireMessage, TextContent, content,
};
use polyc_proto::proto::polychrome::events::v1::{
AttributionEvent, ModelCallEvent, SummaryEvent, UsageEvent,
};
use polyc_proto::proto::polychrome::persona::v1::ExternalIdentity;
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
use super::*;
struct Fixture {
eventlog: Arc<EventLogHost>,
shutdown: CancellationToken,
dir: PathBuf,
signer: ApprovalSigner,
}
impl Fixture {
async fn build(test_name: &str) -> Self {
let signer = ApprovalSigner::from_seed(7);
let dir = std::env::temp_dir().join(format!(
"polyc-query-dashboard-{test_name}-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let eventlog = Arc::new(
EventLogHost::spawn(dir.clone(), shutdown.clone(), signer.relabel_for_test())
.expect("spawn eventlog host"),
);
Self {
eventlog,
shutdown,
dir,
signer,
}
}
fn dashboard(&self) -> DashboardCell {
self.dashboard_at(polyc_payments::amount::DEFAULT_DECIMALS)
}
fn dashboard_at(&self, settlement_decimals: u32) -> DashboardCell {
DashboardProjection::new(
vec![self.signer.public_key_bytes()],
self.eventlog.clone(),
settlement_decimals,
)
}
async fn commit_to(
&self,
dashboard: &DashboardCell,
partition: &str,
events: Vec<Event>,
) -> Vec<FeedRecord> {
let positions = self
.eventlog
.append_batch(partition.to_owned(), events.clone())
.await
.expect("append");
let chunk = vec![crate::feed::test_commit(partition, &events, &positions)];
dashboard.apply_commits(partition, &chunk);
chunk
}
}
impl Drop for Fixture {
fn drop(&mut self) {
self.shutdown.cancel();
let _ = std::fs::remove_dir_all(&self.dir);
}
}
fn text_message(role: &str, text: &str, internal_only: bool) -> Vec<u8> {
WireMessage {
role: role.to_owned(),
content: buffa::MessageField::some(Content {
r#type: Some(content::Type::Text(Box::new(TextContent {
text: text.to_owned(),
..Default::default()
}))),
..Default::default()
}),
internal_only,
..Default::default()
}
.encode_to_vec()
}
fn caller_payload(persona_id: &str, provider: &str) -> Vec<u8> {
AttributionEvent {
persona_id: persona_id.to_owned(),
role: "initiator".to_owned(),
identity: buffa::MessageField::some(ExternalIdentity {
provider: provider.to_owned(),
scope: "s".to_owned(),
external_id: "u1".to_owned(),
display_name: "A".to_owned(),
..Default::default()
}),
..Default::default()
}
.encode_to_vec()
}
fn usage_payload(input_tokens: u64, output_tokens: u64) -> Vec<u8> {
UsageEvent {
input_tokens,
output_tokens,
..Default::default()
}
.encode_to_vec()
}
fn model_call_payload(captured_clock_unix_ms: u64) -> Vec<u8> {
ModelCallEvent {
captured_clock_unix_ms,
..Default::default()
}
.encode_to_vec()
}
fn summary_payload(text: &str) -> Vec<u8> {
SummaryEvent {
text: text.to_owned(),
covers_through_position: 0,
..Default::default()
}
.encode_to_vec()
}
fn receipt_bytes_of_kind(
signer: &ApprovalSigner,
kind: &str,
subject: &str,
amount: &str,
) -> Vec<u8> {
let (payload, _sig, _pk) = receipt_payload(
&ReceiptPayload {
kind,
reference: "tx-1",
amount,
currency: "USDC",
recipient: "0xrecipient",
method: "tempo",
timestamp: "2026-07-20T00:00:00Z",
tool_call_id: "call-1",
approval_pos: "1",
approved_args_hash: "hash",
subject,
payer_kind: "linked_wallet",
paying_account: "0xpayer",
},
signer,
);
payload
}
fn receipt_bytes(signer: &ApprovalSigner, subject: &str, amount: &str) -> Vec<u8> {
receipt_bytes_of_kind(signer, kinds::OUTBOUND_PAYMENT_RECEIPT, subject, amount)
}
fn inbound_receipt_bytes(signer: &ApprovalSigner, subject: &str, amount: &str) -> Vec<u8> {
receipt_bytes_of_kind(signer, kinds::PAYMENT_RECEIPT, subject, amount)
}
#[tokio::test]
async fn parity_row_matches_compute_stats_semantics_for_a_seeded_partition() {
let fx = Fixture::build("parity").await;
let dashboard = fx.dashboard();
let partition = "conv-parity-1".to_owned();
let turn = Uuid::now_v7();
fx.eventlog
.append_batch(
partition.clone(),
vec![Event::new(
kinds::tagged(kinds::CALLER, &turn),
caller_payload("persona-a", "web"),
)],
)
.await
.expect("append caller");
fx.eventlog
.append_batch(
partition.clone(),
vec![Event::new(
kinds::SUMMARY,
summary_payload("condensed so far"),
)],
)
.await
.expect("append summary");
fx.eventlog
.append_batch(
partition.clone(),
vec![
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
Event::new(
kinds::tagged(kinds::USER_MSG, &turn),
text_message("user", "hello there", false),
),
Event::new(kinds::tagged(kinds::USAGE, &turn), usage_payload(10, 5)),
Event::new(
kinds::tagged(kinds::MODEL_CALL, &turn),
model_call_payload(1_700_000_000_000),
),
Event::new(
kinds::tagged(kinds::OUTPUT_MSG, &turn),
text_message("model", "hi", false),
),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
],
)
.await
.expect("append persist_turn batch");
fx.eventlog
.append_batch(
partition.clone(),
vec![
Event::new(
kinds::OUTBOUND_PAYMENT_RECEIPT,
receipt_bytes(&fx.signer, "persona-a", "1500"),
),
Event::new(
kinds::OUTBOUND_PAYMENT_RECEIPT,
receipt_bytes(&fx.signer, "persona-a", ""),
),
Event::new(
kinds::PAYMENT_RECEIPT,
inbound_receipt_bytes(&fx.signer, "", "0.01"),
),
],
)
.await
.expect("append receipts");
dashboard
.rebuild_from_full_fleet()
.await
.expect("full-fleet rebuild");
let row = dashboard.row("parity-1").expect("row exists");
assert_eq!(
row.total_events,
11 + 4,
"every appended event counts, including the host's own per-commit MMR marker \
(matches compute_stats' unfiltered events.len())"
);
assert_eq!(row.committed_turns, 1);
assert_eq!(row.input_tokens, 10);
assert_eq!(row.output_tokens, 5);
assert_eq!(row.summary_text.as_deref(), Some("condensed so far"));
assert!(row.summary_active());
assert_eq!(
row.last_turn_id.as_deref(),
Some(turn.as_simple().to_string().as_str())
);
assert_eq!(row.edges, vec!["web".to_owned()]);
assert_eq!(
row.settlements,
vec![
DashboardSettlement {
persona: "persona-a".to_owned(),
spend_base_units: 1500,
charged_base_units: 0,
},
DashboardSettlement {
persona: "(unattributed)".to_owned(),
spend_base_units: 0,
charged_base_units: 10_000,
},
]
);
assert_eq!(row.created_at_ms, Some(1_700_000_000_000));
assert_eq!(row.last_activity_ms, Some(1_700_000_000_000));
assert_eq!(row.persona_id.as_deref(), Some("persona-a"));
assert_eq!(row.first_message_preview.as_deref(), Some("hello there"));
}
#[tokio::test]
async fn persona_id_falls_back_to_a_later_committed_turns_caller() {
let fx = Fixture::build("persona-fallback").await;
let dashboard = fx.dashboard();
let partition = "conv-persona-fallback".to_owned();
let turn1 = Uuid::now_v7();
let turn2 = Uuid::now_v7();
fx.commit_to(
&dashboard,
&partition,
vec![
Event::new(kinds::tagged(kinds::TURN_START, &turn1), Vec::new()),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn1), Vec::new()),
],
)
.await;
fx.commit_to(
&dashboard,
&partition,
vec![
Event::new(
kinds::tagged(kinds::CALLER, &turn2),
caller_payload("persona-a", "web"),
),
Event::new(kinds::tagged(kinds::TURN_START, &turn2), Vec::new()),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn2), Vec::new()),
],
)
.await;
let row = dashboard
.row("persona-fallback")
.expect("row exists after two committed turns");
assert_eq!(row.committed_turns, 2);
assert_eq!(
row.persona_id.as_deref(),
Some("persona-a"),
"the first committed turn had no caller event; persona_id must fall back \
to the second committed turn's own attribution instead of freezing at None"
);
let turn3 = Uuid::now_v7();
let bare_partition = "conv-persona-fallback-bare".to_owned();
fx.commit_to(
&dashboard,
&bare_partition,
vec![
Event::new(kinds::tagged(kinds::TURN_START, &turn3), Vec::new()),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn3), Vec::new()),
],
)
.await;
let bare_row = dashboard
.row("persona-fallback-bare")
.expect("row exists after one committed turn");
assert!(bare_row.persona_id.is_none());
}
#[tokio::test]
async fn settlement_rollup_reads_each_direction_in_its_own_unit() {
let fx = Fixture::build("units").await;
let dashboard = fx.dashboard();
let partition = "conv-units-1".to_owned();
fx.eventlog
.append_batch(
partition.clone(),
vec![
Event::new(
kinds::OUTBOUND_PAYMENT_RECEIPT,
receipt_bytes(&fx.signer, "persona-a", "10000"),
),
Event::new(
kinds::PAYMENT_RECEIPT,
inbound_receipt_bytes(&fx.signer, "", "0.01"),
),
Event::new(
kinds::PAYMENT_RECEIPT,
inbound_receipt_bytes(&fx.signer, "", "0.25"),
),
],
)
.await
.expect("append receipts");
dashboard
.rebuild_from_full_fleet()
.await
.expect("full-fleet rebuild");
let row = dashboard.row("units-1").expect("row exists");
assert_eq!(
row.settlements,
vec![
DashboardSettlement {
persona: "persona-a".to_owned(),
spend_base_units: 10_000,
charged_base_units: 0,
},
DashboardSettlement {
persona: "(unattributed)".to_owned(),
spend_base_units: 0,
charged_base_units: 260_000,
},
],
"inbound decimal charges reach the rollup, normalized to base units, in the charged \
total rather than added to the outbound caller's spend"
);
}
#[tokio::test]
async fn one_bucket_holding_both_directions_keeps_them_apart() {
let fx = Fixture::build("mixed").await;
let dashboard = fx.dashboard();
let partition = "conv-mixed-1".to_owned();
fx.eventlog
.append_batch(
partition,
vec![
Event::new(
kinds::OUTBOUND_PAYMENT_RECEIPT,
receipt_bytes(&fx.signer, "", "10000"),
),
Event::new(
kinds::PAYMENT_RECEIPT,
inbound_receipt_bytes(&fx.signer, "", "0.5"),
),
],
)
.await
.expect("append receipts");
dashboard
.rebuild_from_full_fleet()
.await
.expect("full-fleet rebuild");
assert_eq!(
dashboard.row("mixed-1").expect("row exists").settlements,
vec![DashboardSettlement {
persona: "(unattributed)".to_owned(),
spend_base_units: 10_000,
charged_base_units: 500_000,
}],
"one bucket, two totals: 510_000 would be a number describing nothing"
);
}
#[tokio::test]
async fn the_inbound_rollup_scales_at_the_configured_settlement_decimals() {
let fx = Fixture::build("decimals").await;
let partition = "conv-decimals-1".to_owned();
fx.eventlog
.append_batch(
partition,
vec![Event::new(
kinds::PAYMENT_RECEIPT,
inbound_receipt_bytes(&fx.signer, "", "0.01"),
)],
)
.await
.expect("append the inbound receipt");
let at_default = fx.dashboard_at(polyc_payments::amount::DEFAULT_DECIMALS);
at_default
.rebuild_from_full_fleet()
.await
.expect("full-fleet rebuild at the default scale");
assert_eq!(
at_default
.row("decimals-1")
.expect("row exists")
.settlements,
vec![DashboardSettlement {
persona: "(unattributed)".to_owned(),
spend_base_units: 0,
charged_base_units: 10_000,
}],
"the default scale reads 0.01 as 10_000 base units"
);
let at_eight = fx.dashboard_at(8);
at_eight
.rebuild_from_full_fleet()
.await
.expect("full-fleet rebuild at eight decimals");
assert_eq!(
at_eight.row("decimals-1").expect("row exists").settlements,
vec![DashboardSettlement {
persona: "(unattributed)".to_owned(),
spend_base_units: 0,
charged_base_units: 1_000_000,
}],
"a deployment configured at eight decimals scales the SAME stored amount at eight, \
not at the default"
);
}
#[tokio::test]
#[tracing_test::traced_test]
async fn an_unreadable_receipt_amount_is_counted_not_silently_skipped() {
let fx = Fixture::build("unreadable").await;
let dashboard = fx.dashboard();
let partition = "conv-unreadable-1".to_owned();
let absent_before = polyc_payments::amount::unreadable_amount_count(
SettlementDirection::Outbound,
polyc_payments::amount::AmountReadError::Absent,
);
let malformed_before = polyc_payments::amount::unreadable_amount_count(
SettlementDirection::Inbound,
polyc_payments::amount::AmountReadError::MalformedDollars,
);
fx.eventlog
.append_batch(
partition.clone(),
vec![
Event::new(
kinds::OUTBOUND_PAYMENT_RECEIPT,
receipt_bytes(&fx.signer, "persona-a", ""),
),
Event::new(
kinds::PAYMENT_RECEIPT,
inbound_receipt_bytes(&fx.signer, "", "not-a-number"),
),
],
)
.await
.expect("append receipts");
dashboard
.rebuild_from_full_fleet()
.await
.expect("full-fleet rebuild");
let row = dashboard.row("unreadable-1").expect("row exists");
assert!(
row.settlements.is_empty(),
"no amount is invented for a receipt that recorded none"
);
assert_eq!(
polyc_payments::amount::unreadable_amount_count(
SettlementDirection::Outbound,
polyc_payments::amount::AmountReadError::Absent,
) - absent_before,
1,
"the uncaptured outbound charge is counted"
);
assert_eq!(
polyc_payments::amount::unreadable_amount_count(
SettlementDirection::Inbound,
polyc_payments::amount::AmountReadError::MalformedDollars,
) - malformed_before,
1,
"the unreadable inbound charge is counted"
);
assert!(
logs_contain("below what really settled"),
"this call site words the dashboard-rollup consequence"
);
assert!(
!logs_contain("the reseeded budget is below actual spend"),
"the committed-spend-floor consequence belongs to another call site"
);
assert!(
!logs_contain("does not appear in their own history at all"),
"the wallet-history consequence belongs to another call site"
);
}
#[tokio::test]
async fn feed_delta_applies_across_three_different_commit_paths() {
let fx = Fixture::build("three-paths").await;
let dashboard = fx.dashboard();
let partition = "conv-three-paths".to_owned();
let turn = Uuid::now_v7();
fx.commit_to(
&dashboard,
&partition,
vec![
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
Event::new(
kinds::tagged(kinds::USER_MSG, &turn),
text_message("user", "hi there", false),
),
Event::new(kinds::tagged(kinds::USAGE, &turn), usage_payload(3, 4)),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
],
)
.await;
let after_a = dashboard
.row("three-paths")
.expect("row exists after path A");
assert_eq!(after_a.committed_turns, 1);
assert_eq!(after_a.input_tokens, 3);
assert_eq!(after_a.total_events, 4);
fx.commit_to(
&dashboard,
&partition,
vec![Event::new(kinds::SUMMARY, summary_payload("first summary"))],
)
.await;
let after_b = dashboard
.row("three-paths")
.expect("row exists after path B");
assert_eq!(after_b.summary_text.as_deref(), Some("first summary"));
assert_eq!(after_b.total_events, 5, "the summary event also counts");
assert_eq!(
after_b.committed_turns, 1,
"unaffected by the summary append"
);
let approval_chunk = fx
.commit_to(
&dashboard,
&partition,
vec![Event::new(
kinds::tagged(kinds::APPROVAL_RESPONSE, &turn),
Vec::new(),
)],
)
.await;
let after_c = dashboard
.row("three-paths")
.expect("row exists after path C");
assert_eq!(after_c.total_events, 6, "the approval response also counts");
assert_eq!(
after_c.committed_turns, 1,
"unaffected by the approval response"
);
assert_eq!(
after_c.input_tokens, 3,
"unaffected by the approval response"
);
dashboard.apply_commits(&partition, &approval_chunk);
assert_eq!(
dashboard.row("three-paths").expect("row still exists"),
after_c,
"a redelivered chunk contributes nothing"
);
}
#[test]
fn reapplying_the_same_positions_does_not_double_count() {
let turn = Uuid::now_v7();
let events = [
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
Event::new(kinds::tagged(kinds::USAGE, &turn), usage_payload(7, 2)),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
];
let positioned: Vec<(u64, &Event)> = events
.iter()
.enumerate()
.map(|(i, e)| (i as u64, e))
.collect();
let mut tracked = TrackedRow::new("redelivery".to_owned());
tracked.apply(
positioned.iter().copied(),
&[],
polyc_payments::amount::DEFAULT_DECIMALS,
);
assert_eq!(tracked.row.total_events, 3);
assert_eq!(tracked.row.input_tokens, 7);
assert_eq!(tracked.row.committed_turns, 1);
tracked.apply(
positioned.iter().copied(),
&[],
polyc_payments::amount::DEFAULT_DECIMALS,
);
assert_eq!(
tracked.row.total_events, 3,
"redelivery must not double-count events"
);
assert_eq!(
tracked.row.input_tokens, 7,
"redelivery must not double-count usage"
);
assert_eq!(
tracked.row.committed_turns, 1,
"redelivery must not double-count turns"
);
}
#[test]
fn corrupt_usage_payload_counts_as_zero_tokens() {
let turn = Uuid::now_v7();
let events = [
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
Event::new(kinds::tagged(kinds::USAGE, &turn), vec![0xFF, 0xFE, 0xFD]),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
];
let positioned: Vec<(u64, &Event)> = events
.iter()
.enumerate()
.map(|(i, e)| (i as u64, e))
.collect();
let mut tracked = TrackedRow::new("corrupt-usage".to_owned());
tracked.apply(
positioned.iter().copied(),
&[],
polyc_payments::amount::DEFAULT_DECIMALS,
);
assert_eq!(
tracked.row.total_events, 3,
"a corrupt usage event still counts toward total_events"
);
assert_eq!(
tracked.row.input_tokens, 0,
"a corrupt usage payload contributes zero input tokens, not a skipped row"
);
assert_eq!(
tracked.row.output_tokens, 0,
"a corrupt usage payload contributes zero output tokens, not a skipped row"
);
}
#[test]
fn usage_accumulation_saturates_at_the_boundary_instead_of_panicking() {
let turn = Uuid::now_v7();
let events = [
Event::new(
kinds::tagged(kinds::USAGE, &turn),
usage_payload(u64::MAX, u64::MAX),
),
Event::new(kinds::tagged(kinds::USAGE, &turn), usage_payload(5, 5)),
];
let positioned: Vec<(u64, &Event)> = events
.iter()
.enumerate()
.map(|(i, e)| (i as u64, e))
.collect();
let mut tracked = TrackedRow::new("saturating".to_owned());
tracked.apply(
positioned.iter().copied(),
&[],
polyc_payments::amount::DEFAULT_DECIMALS,
);
assert_eq!(tracked.row.input_tokens, u64::MAX);
assert_eq!(tracked.row.output_tokens, u64::MAX);
}
#[test]
fn corrupt_model_call_payload_is_skipped() {
let turn = Uuid::now_v7();
let events = [
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
Event::new(
kinds::tagged(kinds::MODEL_CALL, &turn),
vec![0xFF, 0xFE, 0xFD],
),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
];
let positioned: Vec<(u64, &Event)> = events
.iter()
.enumerate()
.map(|(i, e)| (i as u64, e))
.collect();
let mut tracked = TrackedRow::new("corrupt-model-call".to_owned());
tracked.apply(
positioned.iter().copied(),
&[],
polyc_payments::amount::DEFAULT_DECIMALS,
);
assert_eq!(tracked.row.committed_turns, 1);
assert_eq!(
tracked.row.last_activity_ms, None,
"a corrupt model_call payload never resolves a dispatch clock for this turn"
);
}
#[tokio::test]
async fn a_reported_rewrite_rebuilds_the_row_with_erased_content_gone() {
let fx = Fixture::build("rewrite").await;
let dashboard = fx.dashboard();
let partition = "conv-rewrite".to_owned();
let turn = Uuid::now_v7();
fx.commit_to(
&dashboard,
&partition,
vec![
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
Event::new(
kinds::tagged(kinds::USER_MSG, &turn),
text_message("user", "secret to erase", false),
),
Event::new(kinds::tagged(kinds::USAGE, &turn), usage_payload(9, 1)),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
],
)
.await;
let before = dashboard.row("rewrite").expect("row exists before rewrite");
assert_eq!(before.total_events, 4);
assert_eq!(before.input_tokens, 9);
assert_eq!(
before.first_message_preview.as_deref(),
Some("secret to erase")
);
fx.eventlog
.rewrite_partition(
partition,
"test-dashboard-erase-usage".to_owned(),
Box::new(|event: &Event| {
let (base, _) = kinds::parse(&event.kind);
if base == kinds::USAGE {
polyc_eventlog_host::RewriteDecision::Replace(usage_payload(0, 0))
} else {
polyc_eventlog_host::RewriteDecision::Keep
}
}),
)
.await
.expect("rewrite");
dashboard
.note_partition_change("conv-rewrite", PartitionChange::Rewritten)
.await;
{
let row = dashboard.row("rewrite").expect("row still exists");
assert_eq!(row.input_tokens, 0);
assert_eq!(
row.total_events,
5,
"rewrite drops nothing here, only replaces — the partition is re-rooted rather \
than lengthened"
);
}
}
#[tokio::test]
async fn repair_partition_rebuilds_the_row_with_quarantined_content_gone() {
let dir = std::env::temp_dir().join(format!(
"polyc-query-dashboard-repair-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let signer = ApprovalSigner::from_seed(7);
let partition = "conv-repair".to_owned();
let marker = b"UNIQUE_QUARANTINE_MARKER".to_vec();
{
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(dir.clone(), shutdown, signer.relabel_for_test())
.expect("spawn");
host.append_batch(
partition.clone(),
vec![
Event::new(kinds::TURN_START, Vec::new()),
Event::new(kinds::USER_MSG, marker.clone()),
Event::new(kinds::TURN_COMPLETE, Vec::new()),
],
)
.await
.expect("append");
drop(host);
}
let data_file = dir
.join(format!("{partition}_data"))
.join("0000000000000000");
let mut bytes = std::fs::read(&data_file).expect("read section 0");
let payload_at = bytes
.windows(marker.len())
.position(|w| w == marker.as_slice())
.expect("the marker payload bytes are present on disk");
bytes[payload_at] ^= 0xFF;
std::fs::write(&data_file, &bytes).expect("write corrupted section 0");
let shutdown2 = CancellationToken::new();
let eventlog = Arc::new(
EventLogHost::spawn(dir.clone(), shutdown2.clone(), signer.relabel_for_test())
.expect("reopen"),
);
let dashboard = DashboardProjection::new(
vec![signer.public_key_bytes()],
eventlog.clone(),
polyc_payments::amount::DEFAULT_DECIMALS,
);
assert!(
dashboard.row("repair").is_none(),
"no row should exist before this projection has been told anything"
);
let quarantined = eventlog
.repair_partition(partition.clone())
.await
.expect("repair completes");
assert!(
!quarantined.is_empty(),
"the corrupted event must quarantine for this test to be meaningful"
);
let ground_truth = eventlog
.replay_with_positions(partition.clone())
.await
.expect("replay after repair");
let expected_total_events = ground_truth.len();
dashboard
.note_partition_change(&partition, PartitionChange::Rewritten)
.await;
assert_eq!(
dashboard
.row("repair")
.expect("the reported repair built the row")
.total_events,
expected_total_events,
"the row agrees with a fresh replay of the repaired partition"
);
shutdown2.cancel();
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_reported_destroy_removes_the_row() {
let fx = Fixture::build("destroy").await;
let dashboard = fx.dashboard();
let partition = "conv-destroy".to_owned();
let turn = Uuid::now_v7();
fx.commit_to(
&dashboard,
&partition,
vec![
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
],
)
.await;
assert!(dashboard.row("destroy").is_some());
fx.eventlog
.destroy_partition(partition.clone())
.await
.expect("destroy");
dashboard
.note_partition_change(&partition, PartitionChange::Destroyed)
.await;
assert!(
dashboard.row("destroy").is_none(),
"a reported destroy removes the row"
);
dashboard
.note_partition_change(&partition, PartitionChange::Destroyed)
.await;
assert!(
dashboard.row("destroy").is_none(),
"reporting the same destroy again is a no-op"
);
}
#[tokio::test]
async fn a_full_fleet_rebuild_prunes_a_row_no_notification_ever_reported() {
let fx = Fixture::build("prune").await;
let dashboard = fx.dashboard();
let turn = Uuid::now_v7();
for partition in ["conv-prune-gone", "conv-prune-kept"] {
fx.commit_to(
&dashboard,
partition,
vec![
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
],
)
.await;
}
assert!(dashboard.row("prune-gone").is_some());
fx.eventlog
.destroy_partition("conv-prune-gone".to_owned())
.await
.expect("destroy");
dashboard
.rebuild_from_full_fleet()
.await
.expect("full-fleet rebuild");
assert!(
dashboard.row("prune-gone").is_none(),
"the reconcile prunes a conversation the fleet no longer holds"
);
assert!(
dashboard.row("prune-kept").is_some(),
"a conversation the fleet still holds is never pruned"
);
}
struct FailingReplays {
inner: Arc<EventLogHost>,
unreadable: Vec<String>,
}
impl FailingReplays {
fn refuses(&self, partition: &str) -> Option<JournalError> {
self.unreadable
.iter()
.any(|p| p == partition)
.then(|| JournalError::Unreachable("the state plane restarted mid-pass".to_owned()))
}
}
#[async_trait::async_trait]
impl PartitionJournal for FailingReplays {
async fn list_partitions(&self) -> Result<Vec<String>, JournalError> {
PartitionJournal::list_partitions(self.inner.as_ref()).await
}
async fn partition_event_count(&self, partition: String) -> Result<u64, JournalError> {
if let Some(error) = self.refuses(&partition) {
return Err(error);
}
PartitionJournal::partition_event_count(self.inner.as_ref(), partition).await
}
async fn replay_with_positions(
&self,
partition: String,
) -> Result<Vec<(u64, Event)>, JournalError> {
if let Some(error) = self.refuses(&partition) {
return Err(error);
}
PartitionJournal::replay_with_positions(self.inner.as_ref(), partition).await
}
async fn replay_with_positions_bounded(
&self,
partition: String,
max_bytes: u64,
) -> Result<polyc_eventlog::BoundedReplay, JournalError> {
if let Some(error) = self.refuses(&partition) {
return Err(error);
}
PartitionJournal::replay_with_positions_bounded(
self.inner.as_ref(),
partition,
max_bytes,
)
.await
}
async fn replay_from_with_positions_bounded(
&self,
partition: String,
start: u64,
max_bytes: u64,
) -> Result<polyc_eventlog::BoundedReplay, JournalError> {
if let Some(error) = self.refuses(&partition) {
return Err(error);
}
PartitionJournal::replay_from_with_positions_bounded(
self.inner.as_ref(),
partition,
start,
max_bytes,
)
.await
}
async fn replay_range_with_positions_bounded(
&self,
partition: String,
start: u64,
end: u64,
max_bytes: u64,
) -> Result<polyc_eventlog::BoundedReplay, JournalError> {
if let Some(error) = self.refuses(&partition) {
return Err(error);
}
PartitionJournal::replay_range_with_positions_bounded(
self.inner.as_ref(),
partition,
start,
end,
max_bytes,
)
.await
}
async fn append_batch(
&self,
partition: String,
events: Vec<Event>,
) -> Result<(), JournalError> {
PartitionJournal::append_batch(self.inner.as_ref(), partition, events).await
}
fn is_stopping(&self) -> bool {
PartitionJournal::is_stopping(self.inner.as_ref())
}
}
#[tokio::test]
async fn a_full_fleet_rebuild_keeps_a_row_whose_replay_failed_this_pass() {
let fx = Fixture::build("unread").await;
let dashboard = DashboardProjection::new(
vec![fx.signer.public_key_bytes()],
Arc::new(FailingReplays {
inner: fx.eventlog.clone(),
unreadable: vec!["conv-unread-quiet".to_owned()],
}),
polyc_payments::amount::DEFAULT_DECIMALS,
);
for (partition, tokens) in [("conv-unread-quiet", 13u64), ("conv-unread-readable", 4u64)] {
let turn = Uuid::now_v7();
fx.commit_to(
&dashboard,
partition,
vec![
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
Event::new(kinds::tagged(kinds::USAGE, &turn), usage_payload(tokens, 0)),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
],
)
.await;
}
let before = dashboard.row("unread-quiet").expect("row exists");
assert_eq!(before.committed_turns, 1);
assert_eq!(before.input_tokens, 13);
let summary = dashboard
.rebuild_from_full_fleet()
.await
.expect("full-fleet rebuild");
assert_eq!(
summary.conversation_count, 1,
"only the readable conversation was folded this pass"
);
assert_eq!(
dashboard.row("unread-quiet").as_ref(),
Some(&before),
"a conversation this pass could not read keeps its row, counters intact — an \
unreadable partition is not an absent one"
);
assert_eq!(
dashboard
.row("unread-readable")
.expect("the readable row survives")
.input_tokens,
4
);
}
#[tokio::test]
async fn boot_rebuild_populates_rows_for_every_partition() {
let fx = Fixture::build("boot").await;
let dashboard = fx.dashboard();
for (id, tokens) in [("conv-boot-a", 5u64), ("conv-boot-b", 11u64)] {
let turn = Uuid::now_v7();
fx.eventlog
.append_batch(
id.to_owned(),
vec![
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
Event::new(kinds::tagged(kinds::USAGE, &turn), usage_payload(tokens, 0)),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
],
)
.await
.expect("append");
}
fx.eventlog
.append_batch(
"audit-log".to_owned(),
vec![Event::new("marker", Vec::new())],
)
.await
.expect("append audit event");
assert_eq!(dashboard.rebuild_status(), RebuildStatus::Pending);
let summary = dashboard
.rebuild_from_full_fleet()
.await
.expect("full-fleet rebuild");
println!(
"boot rebuild: {} conversation(s) in {} ms",
summary.conversation_count, summary.duration_ms
);
assert_eq!(summary.conversation_count, 2);
let a = dashboard.row("boot-a").expect("row a");
assert_eq!(a.input_tokens, 5);
assert_eq!(a.committed_turns, 1);
let b = dashboard.row("boot-b").expect("row b");
assert_eq!(b.input_tokens, 11);
assert_eq!(b.committed_turns, 1);
assert!(
dashboard.row("audit-log").is_none() && dashboard.rows().len() == 2,
"the non-conversation partition must never surface as a row"
);
assert_eq!(
dashboard.rebuild_status(),
RebuildStatus::Complete {
conversation_count: 2,
duration_ms: summary.duration_ms,
}
);
}
#[tokio::test]
async fn rebuild_never_loses_a_concurrent_delta_or_a_partition_created_after_the_snapshot() {
let fx = Fixture::build("race").await;
let dashboard = fx.dashboard();
let turn1 = Uuid::now_v7();
fx.commit_to(
&dashboard,
"conv-race-existing",
vec![
Event::new(kinds::tagged(kinds::TURN_START, &turn1), Vec::new()),
Event::new(kinds::tagged(kinds::USAGE, &turn1), usage_payload(1, 1)),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn1), Vec::new()),
],
)
.await;
dashboard.race_hook.arm();
let rebuild_task = tokio::spawn({
let dashboard = dashboard.clone();
async move { dashboard.rebuild_from_full_fleet().await }
});
dashboard.race_hook.wait_for_pause().await;
let turn2 = Uuid::now_v7();
fx.commit_to(
&dashboard,
"conv-race-existing",
vec![
Event::new(kinds::tagged(kinds::TURN_START, &turn2), Vec::new()),
Event::new(kinds::tagged(kinds::USAGE, &turn2), usage_payload(9, 9)),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn2), Vec::new()),
],
)
.await;
let turn3 = Uuid::now_v7();
fx.commit_to(
&dashboard,
"conv-race-new",
vec![
Event::new(kinds::tagged(kinds::TURN_START, &turn3), Vec::new()),
Event::new(kinds::tagged(kinds::USAGE, &turn3), usage_payload(5, 5)),
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn3), Vec::new()),
],
)
.await;
dashboard.race_hook.resume();
rebuild_task
.await
.expect("rebuild task did not panic")
.expect("rebuild_from_full_fleet succeeds");
let existing = dashboard
.row("race-existing")
.expect("the pre-existing row must survive the merge");
assert_eq!(
existing.committed_turns, 2,
"the interleaved second turn on an already-replayed conversation must not be lost"
);
assert_eq!(existing.input_tokens, 1 + 9);
let new_row = dashboard.row("race-new").expect(
"a conversation created after the partition-list snapshot must survive the merge, \
not be deleted by it",
);
assert_eq!(new_row.committed_turns, 1);
assert_eq!(new_row.input_tokens, 5);
}
}