use std::collections::BTreeSet;
use std::fmt::Write as _;
use turnframe_core::case::CaseKey;
use turnframe_core::flow::WorkflowRegistry;
use turnframe_core::ids::{AccountId, CaseId, InteractionId, TurnId};
use turnframe_core::interaction::{InteractionKind, InteractionStatus};
use turnframe_core::replay::{DiscardedAnswer, ProviderAttemptOutcome, ReplayRecord, TurnPhase};
use turnframe_core::response::{AssistantTurn, ResponseBlock};
use turnframe_core::target::TargetResolution;
use turnframe_core::understanding::UnderstoodAct;
use turnframe_runtime::resume::CARD_UNIT;
use turnframe_store::conversation::ConversationReader;
use turnframe_store::events::{EventCursor, EventJournalReader};
use turnframe_store::interaction::InteractionReader;
use turnframe_store::journal::CommandJournalReader;
use turnframe_store::replay::ReplayReader;
use turnframe_store::stores::Stores;
use crate::corpus::{BlockKind, CaseSeed};
const EVENT_PAGE_SIZE: usize = 512;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ObservedState {
pub before: serde_json::Value,
pub after: Option<serde_json::Value>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ObservedAct {
pub kind: String,
pub operation: Option<String>,
pub outcome: Option<String>,
}
impl ObservedAct {
fn of(act: &UnderstoodAct, outcome: Option<&String>) -> Self {
Self {
kind: act.kind_name().to_owned(),
operation: act.operation().map(|key| key.as_str().to_owned()),
outcome: outcome.cloned(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ObservedResolution {
pub act_index: usize,
pub resolution: String,
pub case_id: Option<CaseId>,
}
impl ObservedResolution {
fn of(act_index: usize, resolution: &TargetResolution) -> Self {
let (name, case_id) = match resolution {
TargetResolution::Exact { case_ref } => ("exact", Some(case_ref.case_id.clone())),
TargetResolution::Ambiguous { .. } => ("ambiguous", None),
TargetResolution::Missing => ("missing", None),
TargetResolution::Unauthorized => ("unauthorized", None),
TargetResolution::Stale { case_ref, .. } => ("stale", Some(case_ref.case_id.clone())),
_ => ("other", None),
};
Self {
act_index,
resolution: name.to_owned(),
case_id,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ObservedInteraction {
pub id: InteractionId,
pub case: CaseKey,
pub kind: InteractionKind,
pub status: InteractionStatus,
pub blocking: bool,
pub created_by_turn: TurnId,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Observation {
pub turn_id: TurnId,
pub error_code: Option<String>,
pub acts: Vec<ObservedAct>,
pub target_resolutions: Vec<ObservedResolution>,
pub understanding: Option<turnframe_core::understanding::Understanding>,
pub commands: Vec<String>,
pub events: Vec<String>,
pub events_truncated: bool,
pub revisions: Vec<(CaseKey, u64)>,
pub states: Vec<(CaseKey, ObservedState)>,
pub interactions: Vec<ObservedInteraction>,
pub blocks: Vec<BlockKind>,
pub phase: Option<TurnPhase>,
pub provider_failures: usize,
pub discarded_answers: Vec<DiscardedAnswer>,
pub cards_created: usize,
pub answer: String,
}
impl Observation {
pub async fn with_conversation_cases(
mut self,
stores: &Stores,
workflows: &WorkflowRegistry,
account: &AccountId,
earlier: &[TurnId],
) -> Self {
let mut keys: Vec<CaseKey> = Vec::new();
for turn in earlier.iter().copied().chain(std::iter::once(self.turn_id)) {
let entries = CommandJournalReader::for_turn(stores.journal().as_ref(), account, &turn)
.await
.unwrap_or_default();
for entry in entries {
let key = entry.case_ref.key();
let seen =
keys.contains(&key) || self.states.iter().any(|(known, _)| *known == key);
if !seen {
keys.push(key);
}
}
}
for key in keys {
let Some(registered) = workflows.get(&key.workflow) else {
continue;
};
if let Ok(loaded) = registered.executor.load(account, &key.case_id).await {
self.states.push((
key,
ObservedState {
before: serde_json::Value::Null,
after: loaded.value,
},
));
}
}
self
}
pub async fn collect(
stores: &Stores,
workflows: &WorkflowRegistry,
account: &AccountId,
turn_id: TurnId,
cases: &[CaseSeed],
outcome: Result<&AssistantTurn, String>,
) -> Self {
Self::collect_bounded(stores, workflows, account, turn_id, cases, outcome, None).await
}
#[allow(clippy::too_many_arguments)]
pub async fn collect_bounded(
stores: &Stores,
workflows: &WorkflowRegistry,
account: &AccountId,
turn_id: TurnId,
cases: &[CaseSeed],
outcome: Result<&AssistantTurn, String>,
max_events: Option<usize>,
) -> Self {
let (turn, error_code) = match outcome {
Ok(turn) => (Some(turn), None),
Err(code) => (None, Some(code)),
};
let record = ReplayReader::get(stores.replay().as_ref(), account, &turn_id)
.await
.ok();
let entries = CommandJournalReader::for_turn(stores.journal().as_ref(), account, &turn_id)
.await
.unwrap_or_default();
let mut case_keys: Vec<CaseKey> = cases
.iter()
.map(|seed| CaseKey::new(seed.workflow.clone(), seed.case_id.clone()))
.collect();
for entry in &entries {
let key = entry.case_ref.key();
if !case_keys.contains(&key) {
case_keys.push(key);
}
}
let command_ids: BTreeSet<_> = entries.iter().map(|entry| entry.command_id).collect();
let ledger = events_of(stores, account, &case_keys, &command_ids, max_events).await;
Self {
turn_id,
error_code,
acts: acts_of(record.as_ref()),
target_resolutions: resolutions_of(record.as_ref()),
understanding: record
.as_ref()
.and_then(|record| record.understanding.clone()),
commands: entries
.iter()
.map(|entry| entry.command_type.clone())
.collect(),
events: ledger.events,
events_truncated: ledger.truncated,
revisions: revisions_of(workflows, account, &case_keys).await,
states: states_of(workflows, account, cases).await,
interactions: interactions_of(stores, account, turn_id, &case_keys, record.as_ref())
.await,
blocks: turn
.map(|turn| turn.blocks.iter().map(BlockKind::of).collect())
.unwrap_or_default(),
phase: ConversationReader::turn_phase(
stores.conversations().as_ref(),
account,
&turn_id,
)
.await
.ok()
.map(|marker| marker.phase),
provider_failures: provider_failures_of(record.as_ref()),
discarded_answers: record
.as_ref()
.map(ReplayRecord::discarded_answers)
.unwrap_or_default(),
cards_created: record
.as_ref()
.map_or(0, |record| record.interactions_created.len()),
answer: turn.map(narration).unwrap_or_default(),
}
}
#[must_use]
pub fn revision_of(&self, case: &CaseKey) -> Option<u64> {
self.revisions
.iter()
.find(|(key, _)| key == case)
.map(|(_, revision)| *revision)
}
#[must_use]
pub fn interactions_of(&self, case: &CaseKey) -> Vec<&ObservedInteraction> {
self.interactions
.iter()
.filter(|card| &card.case == case)
.collect()
}
#[must_use]
pub fn signature(&self) -> String {
let mut out = String::new();
if let Some(code) = &self.error_code {
let _ = writeln!(out, "error={code}");
}
for act in &self.acts {
let _ = writeln!(
out,
"act={} op={}",
act.kind,
act.operation.as_deref().unwrap_or("-")
);
}
for resolution in &self.target_resolutions {
let _ = writeln!(
out,
"target[{}]={} case={}",
resolution.act_index,
resolution.resolution,
resolution.case_id.as_ref().map_or("-", CaseId::as_str)
);
}
let _ = writeln!(out, "commands={}", self.commands.join(","));
let _ = writeln!(out, "events={}", self.events.join(","));
if self.events_truncated {
let _ = writeln!(out, "events_truncated=true");
}
for (case, revision) in &self.revisions {
let _ = writeln!(out, "rev {}/{}={revision}", case.workflow, case.case_id);
}
for card in &self.interactions {
let _ = writeln!(
out,
"card {}/{} {:?} {:?}",
card.case.workflow, card.case.case_id, card.kind, card.status
);
}
let blocks: Vec<&str> = self.blocks.iter().map(|kind| kind.as_str()).collect();
let _ = writeln!(out, "blocks={}", blocks.join(","));
let _ = writeln!(out, "phase={:?}", self.phase);
out
}
#[must_use]
pub fn discard_codes(&self) -> Vec<&str> {
self.discarded_answers
.iter()
.map(|discarded| discarded.code.as_str())
.collect()
}
}
fn acts_of(record: Option<&ReplayRecord>) -> Vec<ObservedAct> {
let Some(record) = record else {
return Vec::new();
};
message_acts(record)
.map(|(at, act)| ObservedAct::of(act, record.act_outcomes.get(at)))
.collect()
}
fn resolutions_of(record: Option<&ReplayRecord>) -> Vec<ObservedResolution> {
let Some(record) = record else {
return Vec::new();
};
let acts: Vec<&UnderstoodAct> = message_acts(record).map(|(_, act)| act).collect();
record
.target_resolutions
.iter()
.filter_map(|entry| {
let at = acts.iter().position(|act| act.id == entry.act)?;
Some(ObservedResolution::of(at, &entry.resolution))
})
.collect()
}
fn message_acts(record: &ReplayRecord) -> impl Iterator<Item = (usize, &UnderstoodAct)> {
record
.understanding
.iter()
.flat_map(|understanding| understanding.acts.iter().enumerate())
.filter(|(_, act)| act.id.unit != CARD_UNIT)
}
fn provider_failures_of(record: Option<&ReplayRecord>) -> usize {
record.map_or(0, |record| {
record
.provider_attempts
.iter()
.filter(|attempt| {
matches!(
attempt.outcome,
ProviderAttemptOutcome::Failed { .. } | ProviderAttemptOutcome::FellBack { .. }
)
})
.count()
})
}
struct Ledger {
events: Vec<String>,
truncated: bool,
}
async fn events_of(
stores: &Stores,
account: &AccountId,
cases: &[CaseKey],
command_ids: &BTreeSet<turnframe_core::ids::CommandId>,
max_events: Option<usize>,
) -> Ledger {
let mut events = Vec::new();
let mut truncated = false;
let mut cursor = EventCursor::START;
'paging: loop {
let Ok(page) = EventJournalReader::read_from(
stores.events().as_ref(),
account,
cursor,
EVENT_PAGE_SIZE,
)
.await
else {
break;
};
if page.is_empty() {
break;
}
cursor = page.next_cursor;
for event in page.events {
if !command_ids.contains(&event.command_id) || !cases.contains(&event.case_key) {
continue;
}
if max_events.is_some_and(|limit| events.len() >= limit) {
truncated = true;
break 'paging;
}
events.push(event.event_type);
}
}
Ledger { events, truncated }
}
async fn states_of(
workflows: &WorkflowRegistry,
account: &AccountId,
cases: &[CaseSeed],
) -> Vec<(CaseKey, ObservedState)> {
let mut found = Vec::new();
for seed in cases {
let key = CaseKey::new(seed.workflow.clone(), seed.case_id.clone());
let Some(registered) = workflows.get(&seed.workflow) else {
continue;
};
let after = match registered.executor.load(account, &seed.case_id).await {
Ok(loaded) => loaded.value,
Err(_) => continue,
};
found.push((
key,
ObservedState {
before: seed.state.clone(),
after,
},
));
}
found
}
async fn revisions_of(
workflows: &WorkflowRegistry,
account: &AccountId,
cases: &[CaseKey],
) -> Vec<(CaseKey, u64)> {
let mut found = Vec::new();
for case in cases {
let Some(registered) = workflows.get(&case.workflow) else {
continue;
};
if let Ok(loaded) = registered.executor.load(account, &case.case_id).await {
found.push((case.clone(), loaded.revision.0));
}
}
found
}
async fn interactions_of(
stores: &Stores,
account: &AccountId,
turn_id: TurnId,
cases: &[CaseKey],
record: Option<&ReplayRecord>,
) -> Vec<ObservedInteraction> {
let mut ids: Vec<InteractionId> = Vec::new();
for id in record
.iter()
.flat_map(|record| &record.interactions_created)
{
if !ids.contains(id) {
ids.push(*id);
}
}
if let Ok(stored) =
ConversationReader::load_turn(stores.conversations().as_ref(), account, &turn_id).await
&& let Some(answered) = stored.user.input.interaction_response.as_ref()
&& !ids.contains(&answered.interaction_id)
{
ids.push(answered.interaction_id);
}
for case in cases {
let open =
InteractionReader::list_open_for_case(stores.interactions().as_ref(), account, case)
.await
.unwrap_or_default();
for card in open {
if !ids.contains(&card.id) {
ids.push(card.id);
}
}
}
let mut found = Vec::new();
for id in ids {
let Ok(record) = InteractionReader::get(stores.interactions().as_ref(), account, &id).await
else {
continue;
};
let card = record.interaction;
found.push((
card.created_at,
ObservedInteraction {
id: card.id,
case: card.case_ref.key(),
kind: card.kind,
status: card.status,
blocking: card.blocking,
created_by_turn: card.created_by_turn,
},
));
}
found.sort_by(|left, right| left.0.cmp(&right.0).then(left.1.id.cmp(&right.1.id)));
found.into_iter().map(|(_, card)| card).collect()
}
fn narration(turn: &AssistantTurn) -> String {
turn.blocks
.iter()
.filter_map(|block| match block {
ResponseBlock::Answer(answer) => Some(answer.text.as_str()),
ResponseBlock::Transition(transition) => Some(transition.text.as_str()),
_ => None,
})
.collect::<Vec<_>>()
.join(" ")
}