mod answers;
pub mod copy;
use std::collections::BTreeSet;
use std::fmt;
use std::sync::Arc;
use turnframe_core::case::CaseKey;
use turnframe_core::error::OrchestratorError;
use turnframe_core::event::{OperationalReceipt, ReceiptEvent};
use turnframe_core::flow::{ErasedWorkflowView, WorkflowRegistry, WritingStage};
use turnframe_core::hash::derive_uuid;
use turnframe_core::ids::{BlockId, TurnId};
use turnframe_core::interaction::Interaction;
use turnframe_core::knowledge::KnowledgeProvider;
use turnframe_core::locale::{Locale, LocalizedText};
use turnframe_core::reduce::AnswerTask;
use turnframe_core::replay::{BudgetReport, TaskRecord};
use turnframe_core::response::{
AnswerStatus, ArtifactView, AssistantTurn, CaseLabel, Expectation, GeneratedTransition,
InteractionBlock, NarratableFact, NoticeSeverity, ReceiptBlock, ReplayToken, ResponseBlock,
ServerNotice, claim_guard,
};
use turnframe_core::turn::TurnInput;
use turnframe_provider::capabilities::CapabilityRequirements;
use turnframe_provider::request::ContentPart;
use turnframe_provider::router::{ProviderRouter, RoutingPolicy};
use turnframe_store::events::{EventBatch, LedgerReceiptGroup};
use turnframe_tasks::{TaskEngine, TaskKind, TaskScope};
pub use self::copy::{CompositionCopy, notice};
use crate::config::NarrationConfig;
use crate::conversation::{RecentMessage, UnavailableWorkflow};
use crate::narrate::Narrator;
pub use crate::narrate::outcome::AskCopy;
use crate::narrate::outcome::{Material, TurnOutcome};
use crate::narrate::tasks::AcknowledgeInput;
const REPLAY_TOKEN_DOMAIN: &str = "turnframe.replay_token.v1";
const TRANSCRIPT_WINDOW: usize = 4;
struct Carried<'a> {
answers: &'a [String],
unanswered: &'a [String],
notices: &'a [String],
}
#[must_use]
pub fn derive_replay_token(turn_id: &TurnId) -> ReplayToken {
ReplayToken::from(
derive_uuid(REPLAY_TOKEN_DOMAIN, &[&turn_id.to_string()])
.simple()
.to_string(),
)
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct CompositionInput<'a> {
pub turn: &'a TurnInput,
pub attachments: Vec<ContentPart>,
pub answer_tasks: &'a [AnswerTask],
pub events: &'a [EventBatch],
pub ledger: &'a [LedgerReceiptGroup],
pub interactions: &'a [Interaction],
pub views: &'a [ErasedWorkflowView],
pub subjects: &'a [CaseKey],
pub touched: &'a [CaseKey],
pub beside: &'a [CaseKey],
pub reachable_only: &'a [CaseKey],
pub recent: &'a [RecentMessage],
pub preceding_reply: Option<&'a str>,
pub artifacts_shown: &'a [BlockId],
pub case_labels: &'a [CaseLabel],
pub next_steps: &'a [(CaseKey, Vec<LocalizedText>)],
pub named_workflows: &'a [WorkflowKey],
pub unavailable: &'a [UnavailableWorkflow],
pub disputes: &'a [String],
pub started: &'a [WorkflowKey],
pub contested: &'a [String],
pub notices: &'a [ServerNotice],
pub refusals: &'a [NarratableFact],
pub outcomes: OutcomeFlags,
pub effort: Option<&'a crate::effort::EffortProfile>,
}
use turnframe_core::ids::WorkflowKey;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
#[non_exhaustive]
pub struct OutcomeFlags {
pub had_failure: bool,
pub outcome_unknown: bool,
pub revision_conflict: bool,
pub interaction_unavailable: bool,
pub case_refresh_unavailable: bool,
pub budget_exhausted: bool,
pub instruction_declined: bool,
}
impl<'a> CompositionInput<'a> {
#[must_use]
pub fn new(turn: &'a TurnInput) -> Self {
Self {
turn,
attachments: Vec::new(),
answer_tasks: &[],
events: &[],
ledger: &[],
interactions: &[],
views: &[],
subjects: &[],
touched: &[],
beside: &[],
reachable_only: &[],
recent: &[],
preceding_reply: None,
artifacts_shown: &[],
case_labels: &[],
next_steps: &[],
named_workflows: &[],
unavailable: &[],
disputes: &[],
started: &[],
contested: &[],
notices: &[],
refusals: &[],
outcomes: OutcomeFlags::default(),
effort: None,
}
}
}
macro_rules! setters {
($($(#[$doc:meta])* $name:ident: $field:ident: $ty:ty;)*) => {
impl<'a> CompositionInput<'a> {
$(
$(#[$doc])*
#[must_use]
pub fn $name(mut self, value: $ty) -> Self {
self.$field = value;
self
}
)*
}
};
}
setters! {
with_attachments: attachments: Vec<ContentPart>;
with_answer_tasks: answer_tasks: &'a [AnswerTask];
with_events: events: &'a [EventBatch];
with_ledger: ledger: &'a [LedgerReceiptGroup];
with_interactions: interactions: &'a [Interaction];
with_views: views: &'a [ErasedWorkflowView];
with_subjects: subjects: &'a [CaseKey];
with_touched: touched: &'a [CaseKey];
with_beside: beside: &'a [CaseKey];
with_reachable_only: reachable_only: &'a [CaseKey];
with_recent: recent: &'a [RecentMessage];
with_preceding_reply: preceding_reply: Option<&'a str>;
with_artifacts_shown: artifacts_shown: &'a [BlockId];
with_case_labels: case_labels: &'a [CaseLabel];
with_next_steps: next_steps: &'a [(CaseKey, Vec<LocalizedText>)];
with_named_workflows: named_workflows: &'a [WorkflowKey];
with_unavailable: unavailable: &'a [UnavailableWorkflow];
with_disputes: disputes: &'a [String];
with_started: started: &'a [WorkflowKey];
with_contested: contested: &'a [String];
with_notices: notices: &'a [ServerNotice];
with_refusals: refusals: &'a [NarratableFact];
with_outcomes: outcomes: OutcomeFlags;
with_effort: effort: Option<&'a crate::effort::EffortProfile>;
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct Composition {
pub turn: AssistantTurn,
pub tasks: Vec<TaskRecord>,
pub budget: BudgetReport,
pub expectation: Option<Expectation>,
}
#[derive(Clone)]
pub struct Composer {
workflows: Arc<WorkflowRegistry>,
router: Arc<dyn ProviderRouter>,
engine: TaskEngine,
knowledge: Option<Arc<dyn KnowledgeProvider>>,
narration: NarrationConfig,
copy: CompositionCopy,
ask_copy: AskCopy,
max_chunks: Option<usize>,
}
impl fmt::Debug for Composer {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Composer")
.field("narration", &self.narration)
.field("knowledge", &self.knowledge.is_some())
.finish_non_exhaustive()
}
}
impl Composer {
#[must_use]
pub fn new(
workflows: Arc<WorkflowRegistry>,
router: Arc<dyn ProviderRouter>,
narration: NarrationConfig,
) -> Self {
Self {
engine: TaskEngine::builder(Arc::clone(&router)).build(),
workflows,
router,
knowledge: None,
narration,
copy: CompositionCopy::standard(),
ask_copy: AskCopy::standard(),
max_chunks: None,
}
}
#[must_use]
pub fn with_tasks(mut self, engine: TaskEngine) -> Self {
self.engine = engine;
self
}
#[must_use]
pub fn carries(&self, part: &ContentPart) -> bool {
let requirements = match part {
ContentPart::Image { .. } => CapabilityRequirements::none().with_vision(),
ContentPart::Document { .. } => CapabilityRequirements::none().with_documents(),
_ => return true,
};
self.router
.select(TaskKind::Answer, &requirements, &RoutingPolicy::new())
.is_ok()
}
pub(crate) const fn has_knowledge(&self) -> bool {
self.knowledge.is_some()
}
#[must_use]
pub fn with_knowledge(mut self, knowledge: Arc<dyn KnowledgeProvider>) -> Self {
self.knowledge = Some(knowledge);
self
}
#[must_use]
pub fn with_copy(mut self, copy: CompositionCopy) -> Self {
self.copy = copy;
self
}
#[must_use]
pub fn with_ask_copy(mut self, copy: AskCopy) -> Self {
self.ask_copy = copy;
self
}
pub(crate) fn server_copy(&self) -> [&dyn crate::copy::ServerCopy; 2] {
[&self.copy, &self.ask_copy]
}
#[must_use]
pub const fn with_max_chunks(mut self, max_chunks: Option<usize>) -> Self {
self.max_chunks = max_chunks;
self
}
pub fn ledger_receipts(
&self,
groups: &[LedgerReceiptGroup],
locale: &Locale,
) -> Result<Vec<OperationalReceipt>, OrchestratorError> {
let mut receipts = Vec::new();
let mut seen = BTreeSet::new();
for group in groups {
let registered = self.workflows.require(&group.case_key.workflow)?;
for receipt in registered.definition.receipts(&group.events, locale)? {
if seen.insert(receipt.receipt_id) {
receipts.push(receipt);
}
}
}
Ok(receipts)
}
pub fn receipts(
&self,
events: &[EventBatch],
locale: &Locale,
) -> Result<Vec<OperationalReceipt>, OrchestratorError> {
let mut receipts = Vec::new();
let mut seen = BTreeSet::new();
for batch in events {
let registered = self.workflows.require(&batch.case_key.workflow)?;
let committed: Vec<ReceiptEvent<serde_json::Value>> = batch
.events
.iter()
.cloned()
.map(ReceiptEvent::Committed)
.collect();
for receipt in registered.definition.receipts(&committed, locale)? {
if seen.insert(receipt.receipt_id) {
receipts.push(receipt);
}
}
}
Ok(receipts)
}
pub(crate) async fn say_steps(
&self,
steps: futures::channel::mpsc::UnboundedReceiver<turnframe_understand::Step>,
locale: &turnframe_core::locale::Locale,
turn: turnframe_core::ids::TurnId,
publisher: &crate::stream::TurnPublisher,
) -> Vec<turnframe_core::replay::TaskRecord> {
use futures::StreamExt as _;
let scope =
TaskScope::new(self.narration.budget, locale.clone()).for_turn(turn.to_string());
let narrator = Narrator {
engine: &self.engine,
scope: &scope,
max_chars: None,
};
steps
.for_each_concurrent(None, |step| {
let narrator = &narrator;
async move {
if let Some(text) = narrator.step(locale.as_str(), &step.describe()).await {
publisher.step_said(step, text);
}
}
})
.await;
scope.records()
}
fn narrates(&self, input: &CompositionInput<'_>) -> bool {
self.narration.enabled
&& !input.outcomes.budget_exhausted
&& !input.outcomes.case_refresh_unavailable
}
fn guidance(
&self,
input: &CompositionInput<'_>,
cases: &[CaseKey],
stage: WritingStage,
) -> Vec<String> {
input
.views
.iter()
.filter(|view| cases.contains(&view.case_ref.key()))
.filter_map(|view| {
let workflow = self.workflows.require(&view.case_ref.workflow).ok()?;
workflow
.definition
.narration_briefing(stage, view)
.ok()
.flatten()
})
.collect()
}
pub async fn compose(
&self,
input: CompositionInput<'_>,
) -> Result<Composition, OrchestratorError> {
let locale = &input.turn.locale;
let receipts = if input.ledger.is_empty() {
self.receipts(input.events, locale)?
} else {
self.ledger_receipts(input.ledger, locale)?
};
let mut scope = TaskScope::new(
input
.effort
.map_or(self.narration.budget, |effort| effort.reply_budget),
locale.clone(),
)
.for_turn(input.turn.turn_id.to_string());
if let Some(effort) = input.effort {
scope = scope
.with_profiles(effort.tasks.clone())
.with_effort(effort.effort);
}
let narrator = Narrator {
engine: &self.engine,
scope: &scope,
max_chars: self.narration.max_answer_chars,
};
let mut facts = input.refusals.to_vec();
facts.extend(
input
.unavailable
.iter()
.filter(|blocked| input.named_workflows.contains(&blocked.workflow))
.map(|blocked| NarratableFact::WorkflowUnavailable {
workflow: blocked.workflow.clone(),
reason: blocked.reason.clone(),
}),
);
let outcome = Material {
receipts: &receipts,
facts: &facts,
interactions: input.interactions,
views: input.views,
touched: input.touched,
beside: input.beside,
labels: input.case_labels,
disputes: input.disputes,
started: input.started,
contested: input.contested,
next_steps: input.next_steps,
locale,
copy: &self.ask_copy,
}
.outcome();
let answers = answers::Answers {
composer: self,
narrator: &narrator,
input: &input,
};
let answered = answers.all().await;
let notices = self.notices(&input);
let mut blocks: Vec<ResponseBlock> = Vec::new();
let mut asked = false;
if self.narrates(&input) {
let answer_texts: Vec<String> = answered
.iter()
.filter(|answer| answer.status == AnswerStatus::Answered)
.map(|answer| answer.text.clone())
.collect();
let unanswered: Vec<String> = answered
.iter()
.filter(|answer| answer.status != AnswerStatus::Answered)
.map(|answer| {
let question = input
.answer_tasks
.iter()
.find(|task| Some(&task.question_id) == answer.question_id.as_ref())
.map_or("", |task| task.question.as_str());
format!("«{question}»")
})
.collect();
let notice_texts: Vec<String> = notices
.iter()
.map(|notice| notice.text.resolve(locale).to_owned())
.collect();
let only_answers = !crate::narrate::speaks(&outcome)
&& notice_texts.is_empty()
&& unanswered.is_empty();
let reply = if only_answers {
let parts: Vec<String> = answer_texts
.iter()
.cloned()
.chain(outcome.closing.clone())
.collect();
(!parts.is_empty()).then(|| parts.join("\n\n"))
} else if let Some(text) = self
.acknowledge(
&input,
&outcome,
Carried {
answers: &answer_texts,
unanswered: &unanswered,
notices: ¬ice_texts,
},
&narrator,
)
.await
{
asked = outcome.ask.is_some();
Some(text)
} else {
let done: Vec<&str> = receipts
.iter()
.map(|receipt| receipt.body.resolve(locale))
.collect();
let mut parts: Vec<String> = Vec::new();
if !done.is_empty() {
parts.push(done.join(" "));
}
if !outcome.not_done.is_empty() {
parts.push(outcome.not_done.join(" "));
}
parts.extend(answered.iter().map(|answer| answer.text.clone()));
parts.extend(notice_texts);
if let Some(ask) = &outcome.ask {
asked = true;
parts.push(ask.question.clone());
} else if !outcome.next.is_empty() {
let go_on = self.ask_copy.go_on.resolve(locale);
parts.push(format!("{go_on} {}", outcome.next.join(" ")));
} else if let Some(closing) = &outcome.closing {
parts.push(closing.clone());
}
(!parts.is_empty()).then(|| parts.join("\n\n"))
};
if let Some(text) = reply {
blocks.push(self.transition(text, &receipts, &answered, &input));
}
}
blocks.extend(answered.into_iter().map(ResponseBlock::Answer));
for receipt in &receipts {
blocks.push(ResponseBlock::Receipt(ReceiptBlock {
block_id: BlockId::from(format!("receipt:{}", receipt.receipt_id)),
receipt: receipt.clone(),
}));
}
self.artifacts(&input, &mut blocks);
blocks.extend(notices.into_iter().map(ResponseBlock::Notice));
for interaction in input.interactions {
blocks.push(ResponseBlock::Interaction(InteractionBlock {
block_id: BlockId::from(format!("interaction:{}", interaction.id)),
view: interaction.view(),
}));
}
let turn = AssistantTurn {
turn_id: input.turn.turn_id,
conversation_id: input.turn.conversation_id,
blocks,
subjects: input
.views
.iter()
.filter(|view| input.subjects.contains(&view.case_ref.key()))
.map(|view| view.case_ref.clone())
.collect(),
replay_token: derive_replay_token(&input.turn.turn_id),
expectations: Vec::new(),
done: Vec::new(),
};
claim_guard::verify(&turn).map_err(|violation| {
tracing::error!(
target: "turnframe.compose",
violation = %violation,
"composed turn failed the claim guard"
);
OrchestratorError::Internal {
code: "claim_guard".to_owned(),
}
})?;
Ok(Composition {
turn,
tasks: scope.records(),
budget: scope.budget_report(),
expectation: outcome
.ask
.filter(|_| asked)
.and_then(|ask| ask.expectation),
})
}
async fn acknowledge(
&self,
input: &CompositionInput<'_>,
outcome: &TurnOutcome,
carried: Carried<'_>,
narrator: &Narrator<'_>,
) -> Option<String> {
let locale = &input.turn.locale;
let window = input.recent.len().saturating_sub(TRANSCRIPT_WINDOW);
let guidance = self.guidance(input, input.touched, WritingStage::Transition);
let on_screen: Vec<String> = outcome.card.clone().into_iter().collect();
let message = input
.turn
.text
.as_deref()
.filter(|_| outcome.not_done.is_empty())
.and_then(|text| without_questions(text, input.answer_tasks));
let acknowledge = AcknowledgeInput {
outcome,
locale: locale.as_str(),
tone: self.narration.tone,
message: message.as_deref(),
on_screen: &on_screen,
answers: carried.answers,
unanswered: carried.unanswered,
notices: carried.notices,
transcript: &input.recent[window..],
guidance: &guidance,
};
narrator.acknowledge(&acknowledge).await
}
fn transition(
&self,
text: String,
receipts: &[OperationalReceipt],
answered: &[turnframe_core::response::GeneratedAnswer],
input: &CompositionInput<'_>,
) -> ResponseBlock {
let mut facts_used: Vec<NarratableFact> = receipts
.iter()
.map(|receipt| NarratableFact::OperationalOutcome {
receipt_id: receipt.receipt_id,
event_ids: receipt.event_ids.clone(),
status_code: receipt.status_code.clone(),
})
.collect();
facts_used.extend(input.interactions.iter().map(|interaction| {
NarratableFact::InteractionAvailable {
interaction_id: interaction.id,
interaction_kind: interaction.kind,
}
}));
facts_used.extend(input.refusals.iter().cloned());
for answer in answered {
for fact in &answer.facts_used {
if !facts_used.contains(fact) {
facts_used.push(fact.clone());
}
}
}
ResponseBlock::Transition(GeneratedTransition {
block_id: BlockId::from("transition:0"),
text,
facts_used,
})
}
fn artifacts(&self, input: &CompositionInput<'_>, blocks: &mut Vec<ResponseBlock>) {
for view in input.views {
if !input.subjects.contains(&view.case_ref.key()) {
continue;
}
let Ok(workflow) = self.workflows.require(&view.case_ref.workflow) else {
continue;
};
match workflow.definition.artifacts(view) {
Ok(artifacts) => {
for artifact in artifacts {
let block_id = BlockId::from(format!(
"artifact:{}:{}:{}:{}",
view.case_ref.workflow,
view.case_ref.case_id,
view.case_ref.expected_revision,
artifact.artifact_id
));
if !input.artifacts_shown.contains(&block_id) {
blocks
.push(ResponseBlock::Artifact(ArtifactView { block_id, artifact }));
}
}
}
Err(error) => tracing::warn!(
target: "turnframe.compose",
workflow = view.case_ref.workflow.as_str(),
error = %error,
"a workflow's artifacts could not be read from its own view"
),
}
}
}
fn notices(&self, input: &CompositionInput<'_>) -> Vec<ServerNotice> {
let flags = input.outcomes;
let copy = &self.copy;
let own = [
(
flags.revision_conflict,
notice::REVISION_CONFLICT,
NoticeSeverity::Warning,
©.revision_conflict,
),
(
flags.had_failure,
notice::COMMAND_FAILED,
NoticeSeverity::Error,
©.command_failed,
),
(
flags.outcome_unknown,
notice::VERIFICATION_IN_PROGRESS,
NoticeSeverity::Warning,
©.verification_in_progress,
),
(
flags.interaction_unavailable,
notice::INTERACTION_UNAVAILABLE,
NoticeSeverity::Error,
©.interaction_unavailable,
),
(
flags.case_refresh_unavailable,
notice::CASE_REFRESH_UNAVAILABLE,
NoticeSeverity::Error,
©.case_refresh_unavailable,
),
(
flags.instruction_declined,
notice::INSTRUCTION_DECLINED,
NoticeSeverity::Info,
©.instruction_declined,
),
(
flags.budget_exhausted,
notice::BUDGET_EXHAUSTED,
NoticeSeverity::Warning,
©.budget_exhausted,
),
];
let mut notices: Vec<ServerNotice> = Vec::new();
let raised = own
.into_iter()
.filter(|(raised, ..)| *raised)
.map(|(_, code, severity, text)| notice_of(code, severity, text.clone()));
for notice in input.notices.iter().cloned().chain(raised) {
if !notices.iter().any(|existing| existing.code == notice.code) {
notices.push(notice);
}
}
notices
}
}
fn notice_of(code: &str, severity: NoticeSeverity, text: LocalizedText) -> ServerNotice {
ServerNotice {
block_id: BlockId::from(format!("notice:{code}")),
code: code.to_owned(),
severity,
text,
}
}
fn without_questions(text: &str, tasks: &[AnswerTask]) -> Option<String> {
let mut spans: Vec<(usize, usize)> = tasks
.iter()
.filter_map(|task| task.asked_at.map(|span| (span.start_byte, span.end_byte)))
.collect();
spans.sort_unstable();
let mut kept = String::new();
let mut at = 0;
for (start, end) in spans {
let Some(before) = text.get(at..start) else {
continue;
};
kept.push_str(before);
kept.push(' ');
at = end;
}
kept.push_str(text.get(at..).unwrap_or_default());
let kept = kept.split_whitespace().collect::<Vec<_>>().join(" ");
let kept = kept.trim_matches(|c: char| c.is_whitespace() || matches!(c, ',' | ';' | ':'));
(!kept.is_empty()).then(|| kept.to_owned())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn replay_tokens_are_derived_and_stable() {
let turn = TurnId::nil();
assert_eq!(derive_replay_token(&turn), derive_replay_token(&turn));
}
}