Skip to main content

turnframe_runtime/compose/
mod.rs

1//! Response composition: what the assistant is allowed to say (spec §10, §17.3, §18,
2//! §23 steps Q to U).
3//!
4//! By the time this runs everything that could happen has happened, and the job is to
5//! say so accurately. Receipts, notices and cards are rendered by the server from what
6//! committed; the words around them are written by the narration tasks from facts code
7//! selected, and [`claim_guard::verify`] checks the assembled
8//! turn against the record before it is returned. Every [`AnswerTask`] produces exactly
9//! one block, answered or explicitly not.
10//!
11//! | Module | What it holds |
12//! | --- | --- |
13//! | [`copy`] | the notice codes and the copy composition writes itself |
14//! | `answers` | one block per question |
15
16mod answers;
17pub mod copy;
18
19use std::collections::BTreeSet;
20use std::fmt;
21use std::sync::Arc;
22
23use turnframe_core::case::CaseKey;
24use turnframe_core::error::OrchestratorError;
25use turnframe_core::event::{OperationalReceipt, ReceiptEvent};
26use turnframe_core::flow::{ErasedWorkflowView, WorkflowRegistry, WritingStage};
27use turnframe_core::hash::derive_uuid;
28use turnframe_core::ids::{BlockId, TurnId};
29use turnframe_core::interaction::Interaction;
30use turnframe_core::knowledge::KnowledgeProvider;
31use turnframe_core::locale::{Locale, LocalizedText};
32use turnframe_core::reduce::AnswerTask;
33use turnframe_core::replay::{BudgetReport, TaskRecord};
34use turnframe_core::response::{
35    AnswerStatus, ArtifactView, AssistantTurn, CaseLabel, Expectation, GeneratedTransition,
36    InteractionBlock, NarratableFact, NoticeSeverity, ReceiptBlock, ReplayToken, ResponseBlock,
37    ServerNotice, claim_guard,
38};
39use turnframe_core::turn::TurnInput;
40use turnframe_provider::capabilities::CapabilityRequirements;
41use turnframe_provider::request::ContentPart;
42use turnframe_provider::router::{ProviderRouter, RoutingPolicy};
43use turnframe_store::events::{EventBatch, LedgerReceiptGroup};
44use turnframe_tasks::{TaskEngine, TaskKind, TaskScope};
45
46pub use self::copy::{CompositionCopy, notice};
47use crate::config::NarrationConfig;
48use crate::conversation::{RecentMessage, UnavailableWorkflow};
49use crate::narrate::Narrator;
50pub use crate::narrate::outcome::AskCopy;
51use crate::narrate::outcome::{Material, TurnOutcome};
52use crate::narrate::tasks::AcknowledgeInput;
53
54/// Domain separation of the derived replay tokens.
55const REPLAY_TOKEN_DOMAIN: &str = "turnframe.replay_token.v1";
56
57/// How many earlier messages the acknowledgement is shown.
58const TRANSCRIPT_WINDOW: usize = 4;
59
60/// What the one reply gives beside the outcome.
61struct Carried<'a> {
62    answers: &'a [String],
63    unanswered: &'a [String],
64    notices: &'a [String],
65}
66
67/// Derives the opaque token a client uses to fetch a turn's replay record; derived, so a
68/// regenerated response (spec §23.1) carries the token of the one that was lost.
69#[must_use]
70pub fn derive_replay_token(turn_id: &TurnId) -> ReplayToken {
71    ReplayToken::from(
72        derive_uuid(REPLAY_TOKEN_DOMAIN, &[&turn_id.to_string()])
73            .simple()
74            .to_string(),
75    )
76}
77
78/// Everything one composition reads.
79#[derive(Debug, Clone)]
80#[non_exhaustive]
81pub struct CompositionInput<'a> {
82    /// The turn being answered.
83    pub turn: &'a TurnInput,
84    /// The turn's files, for the stage that answers.
85    pub attachments: Vec<ContentPart>,
86    /// The questions the reducer produced, with their basis.
87    pub answer_tasks: &'a [AnswerTask],
88    /// Event batches that committed, in commit order.
89    pub events: &'a [EventBatch],
90    /// Ledger groups read back from the journal, used instead of [`Self::events`]
91    /// when a response is regenerated: a stored event may have been redacted since.
92    pub ledger: &'a [LedgerReceiptGroup],
93    /// Cards on screen after the turn.
94    pub interactions: &'a [Interaction],
95    /// The cases as they now stand.
96    pub views: &'a [ErasedWorkflowView],
97    /// The cases the turn was about.
98    pub subjects: &'a [CaseKey],
99    /// The cases an act of the turn reached, in act order: where the ask comes from.
100    pub touched: &'a [CaseKey],
101    /// Other cases in view the ask may move on to once the touched ones need nothing.
102    pub beside: &'a [CaseKey],
103    /// Cases in view only because the actor may reach them.
104    pub reachable_only: &'a [CaseKey],
105    /// The earlier messages, oldest first.
106    pub recent: &'a [RecentMessage],
107    /// The reply just before this turn.
108    pub preceding_reply: Option<&'a str>,
109    /// Artifact blocks the conversation already carries.
110    pub artifacts_shown: &'a [BlockId],
111    /// What each case is called.
112    pub case_labels: &'a [CaseLabel],
113    /// What each case lets the user do next once it owes nothing, by case.
114    pub next_steps: &'a [(CaseKey, Vec<LocalizedText>)],
115    /// Workflows the turn named.
116    pub named_workflows: &'a [WorkflowKey],
117    /// Workflows that cannot start, with their reasons.
118    pub unavailable: &'a [UnavailableWorkflow],
119    /// What the user said the assistant got wrong.
120    pub disputes: &'a [String],
121    /// Workflows the turn started without writing anything yet.
122    pub started: &'a [WorkflowKey],
123    /// Receipts of the last reply the user contested, as they were shown.
124    pub contested: &'a [String],
125    /// Notices the turn already carries.
126    pub notices: &'a [ServerNotice],
127    /// What the reduction and the execution decided, as facts.
128    pub refusals: &'a [NarratableFact],
129    /// Deterministic outcome flags.
130    pub outcomes: OutcomeFlags,
131    /// The effort the turn runs at: its reply budget and task profiles.
132    pub effort: Option<&'a crate::effort::EffortProfile>,
133}
134
135use turnframe_core::ids::WorkflowKey;
136
137/// The outcomes that raise a notice of composition's own.
138#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
139#[non_exhaustive]
140pub struct OutcomeFlags {
141    /// A command did not commit.
142    pub had_failure: bool,
143    /// An external outcome is unknown.
144    pub outcome_unknown: bool,
145    /// A command met a moved revision.
146    pub revision_conflict: bool,
147    /// A card could not be written.
148    pub interaction_unavailable: bool,
149    /// A touched case could not be read back.
150    pub case_refresh_unavailable: bool,
151    /// The turn's budget is spent: no model writes anything.
152    pub budget_exhausted: bool,
153    /// The user declined the instruction a card guarded.
154    pub instruction_declined: bool,
155}
156
157impl<'a> CompositionInput<'a> {
158    /// An input with nothing but the turn.
159    #[must_use]
160    pub fn new(turn: &'a TurnInput) -> Self {
161        Self {
162            turn,
163            attachments: Vec::new(),
164            answer_tasks: &[],
165            events: &[],
166            ledger: &[],
167            interactions: &[],
168            views: &[],
169            subjects: &[],
170            touched: &[],
171            beside: &[],
172            reachable_only: &[],
173            recent: &[],
174            preceding_reply: None,
175            artifacts_shown: &[],
176            case_labels: &[],
177            next_steps: &[],
178            named_workflows: &[],
179            unavailable: &[],
180            disputes: &[],
181            started: &[],
182            contested: &[],
183            notices: &[],
184            refusals: &[],
185            outcomes: OutcomeFlags::default(),
186            effort: None,
187        }
188    }
189}
190
191macro_rules! setters {
192    ($($(#[$doc:meta])* $name:ident: $field:ident: $ty:ty;)*) => {
193        impl<'a> CompositionInput<'a> {
194            $(
195                $(#[$doc])*
196                #[must_use]
197                pub fn $name(mut self, value: $ty) -> Self {
198                    self.$field = value;
199                    self
200                }
201            )*
202        }
203    };
204}
205
206setters! {
207    /// Sets the turn's files.
208    with_attachments: attachments: Vec<ContentPart>;
209    /// Sets the questions.
210    with_answer_tasks: answer_tasks: &'a [AnswerTask];
211    /// Sets the committed events.
212    with_events: events: &'a [EventBatch];
213    /// Sets the ledger groups of a regenerated response.
214    with_ledger: ledger: &'a [LedgerReceiptGroup];
215    /// Sets the cards on screen.
216    with_interactions: interactions: &'a [Interaction];
217    /// Sets the views.
218    with_views: views: &'a [ErasedWorkflowView];
219    /// Sets the subjects.
220    with_subjects: subjects: &'a [CaseKey];
221    /// Sets the cases an act reached.
222    with_touched: touched: &'a [CaseKey];
223    /// Sets the cases the ask may move on to.
224    with_beside: beside: &'a [CaseKey];
225    /// Sets the cases in view only because they are reachable.
226    with_reachable_only: reachable_only: &'a [CaseKey];
227    /// Sets the earlier messages.
228    with_recent: recent: &'a [RecentMessage];
229    /// Sets the reply just before this turn.
230    with_preceding_reply: preceding_reply: Option<&'a str>;
231    /// Sets the artifact blocks already shown.
232    with_artifacts_shown: artifacts_shown: &'a [BlockId];
233    /// Sets the case labels.
234    with_case_labels: case_labels: &'a [CaseLabel];
235    /// Sets what each case lets the user do next.
236    with_next_steps: next_steps: &'a [(CaseKey, Vec<LocalizedText>)];
237    /// Sets the workflows the turn named.
238    with_named_workflows: named_workflows: &'a [WorkflowKey];
239    /// Sets the workflows that cannot start.
240    with_unavailable: unavailable: &'a [UnavailableWorkflow];
241    /// Sets what the user disputed.
242    with_disputes: disputes: &'a [String];
243    /// Sets the workflows started without a write.
244    with_started: started: &'a [WorkflowKey];
245    /// Sets the receipts the user contested.
246    with_contested: contested: &'a [String];
247    /// Sets the notices.
248    with_notices: notices: &'a [ServerNotice];
249    /// Sets the facts of what was decided.
250    with_refusals: refusals: &'a [NarratableFact];
251    /// Sets the outcome flags.
252    with_outcomes: outcomes: OutcomeFlags;
253    /// Sets the effort the turn runs at.
254    with_effort: effort: Option<&'a crate::effort::EffortProfile>;
255}
256
257/// A composed turn and the model calls that wrote it.
258#[derive(Debug, Clone)]
259#[non_exhaustive]
260pub struct Composition {
261    /// The turn, as it is returned and persisted.
262    pub turn: AssistantTurn,
263    /// Every narration task call.
264    pub tasks: Vec<TaskRecord>,
265    /// What those calls spent.
266    pub budget: BudgetReport,
267    /// What the reply asked for, for the next turn to expect.
268    pub expectation: Option<Expectation>,
269}
270
271/// Builds the ordered response blocks of one turn (spec §18.1).
272#[derive(Clone)]
273pub struct Composer {
274    workflows: Arc<WorkflowRegistry>,
275    router: Arc<dyn ProviderRouter>,
276    engine: TaskEngine,
277    knowledge: Option<Arc<dyn KnowledgeProvider>>,
278    narration: NarrationConfig,
279    copy: CompositionCopy,
280    ask_copy: AskCopy,
281    max_chunks: Option<usize>,
282}
283
284impl fmt::Debug for Composer {
285    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
286        f.debug_struct("Composer")
287            .field("narration", &self.narration)
288            .field("knowledge", &self.knowledge.is_some())
289            .finish_non_exhaustive()
290    }
291}
292
293impl Composer {
294    /// A composer whose narration tasks run over `router` with the shipped profiles.
295    #[must_use]
296    pub fn new(
297        workflows: Arc<WorkflowRegistry>,
298        router: Arc<dyn ProviderRouter>,
299        narration: NarrationConfig,
300    ) -> Self {
301        Self {
302            engine: TaskEngine::builder(Arc::clone(&router)).build(),
303            workflows,
304            router,
305            knowledge: None,
306            narration,
307            copy: CompositionCopy::standard(),
308            ask_copy: AskCopy::standard(),
309            max_chunks: None,
310        }
311    }
312
313    /// Runs the narration tasks on `engine`: its profiles, prompts and records.
314    #[must_use]
315    pub fn with_tasks(mut self, engine: TaskEngine) -> Self {
316        self.engine = engine;
317        self
318    }
319
320    /// Whether a model that answers can be shown `part`.
321    #[must_use]
322    pub fn carries(&self, part: &ContentPart) -> bool {
323        let requirements = match part {
324            ContentPart::Image { .. } => CapabilityRequirements::none().with_vision(),
325            ContentPart::Document { .. } => CapabilityRequirements::none().with_documents(),
326            _ => return true,
327        };
328        self.router
329            .select(TaskKind::Answer, &requirements, &RoutingPolicy::new())
330            .is_ok()
331    }
332
333    /// Attaches a knowledge provider (spec §19.2).
334    #[must_use]
335    pub fn with_knowledge(mut self, knowledge: Arc<dyn KnowledgeProvider>) -> Self {
336        self.knowledge = Some(knowledge);
337        self
338    }
339
340    /// Replaces the copy composition writes itself. English and Italian by default; a
341    /// [`LocalizedText`] adds languages with `with`, and the copy's `translated` a whole
342    /// language at once.
343    #[must_use]
344    pub fn with_copy(mut self, copy: CompositionCopy) -> Self {
345        self.copy = copy;
346        self
347    }
348
349    /// Replaces the questions code writes when no model does.
350    #[must_use]
351    pub fn with_ask_copy(mut self, copy: AskCopy) -> Self {
352        self.ask_copy = copy;
353        self
354    }
355
356    /// The server's own sentences this composer writes, for the languages check.
357    pub(crate) fn server_copy(&self) -> [&dyn crate::copy::ServerCopy; 2] {
358        [&self.copy, &self.ask_copy]
359    }
360
361    /// Sets how many knowledge chunks one answer may retrieve.
362    #[must_use]
363    pub const fn with_max_chunks(mut self, max_chunks: Option<usize>) -> Self {
364        self.max_chunks = max_chunks;
365        self
366    }
367
368    /// Renders receipts from a ledger read, carrying every erasure through.
369    ///
370    /// # Errors
371    ///
372    /// [`OrchestratorError::Erasure`] when a workflow cannot read back its own events.
373    pub fn ledger_receipts(
374        &self,
375        groups: &[LedgerReceiptGroup],
376        locale: &Locale,
377    ) -> Result<Vec<OperationalReceipt>, OrchestratorError> {
378        let mut receipts = Vec::new();
379        let mut seen = BTreeSet::new();
380        for group in groups {
381            let registered = self.workflows.require(&group.case_key.workflow)?;
382            for receipt in registered.definition.receipts(&group.events, locale)? {
383                if seen.insert(receipt.receipt_id) {
384                    receipts.push(receipt);
385                }
386            }
387        }
388        Ok(receipts)
389    }
390
391    /// Renders the receipts of a turn from the events that just committed (§17.3,
392    /// I16). Deterministic, so a regenerated response is the one that was lost.
393    ///
394    /// # Errors
395    ///
396    /// [`OrchestratorError::Erasure`] when a workflow cannot read back its own events.
397    pub fn receipts(
398        &self,
399        events: &[EventBatch],
400        locale: &Locale,
401    ) -> Result<Vec<OperationalReceipt>, OrchestratorError> {
402        let mut receipts = Vec::new();
403        let mut seen = BTreeSet::new();
404        for batch in events {
405            let registered = self.workflows.require(&batch.case_key.workflow)?;
406            let committed: Vec<ReceiptEvent<serde_json::Value>> = batch
407                .events
408                .iter()
409                .cloned()
410                .map(ReceiptEvent::Committed)
411                .collect();
412            for receipt in registered.definition.receipts(&committed, locale)? {
413                if seen.insert(receipt.receipt_id) {
414                    receipts.push(receipt);
415                }
416            }
417        }
418        Ok(receipts)
419    }
420
421    /// Says each step `steps` delivers as it arrives, concurrently, and returns the
422    /// calls made. Ends when the sender is dropped.
423    pub(crate) async fn say_steps(
424        &self,
425        steps: futures::channel::mpsc::UnboundedReceiver<turnframe_understand::Step>,
426        locale: &turnframe_core::locale::Locale,
427        turn: turnframe_core::ids::TurnId,
428        publisher: &crate::stream::TurnPublisher,
429    ) -> Vec<turnframe_core::replay::TaskRecord> {
430        use futures::StreamExt as _;
431        let scope =
432            TaskScope::new(self.narration.budget, locale.clone()).for_turn(turn.to_string());
433        let narrator = Narrator {
434            engine: &self.engine,
435            scope: &scope,
436            max_chars: None,
437        };
438        steps
439            .for_each_concurrent(None, |step| {
440                let narrator = &narrator;
441                async move {
442                    if let Some(text) = narrator.step(locale.as_str(), &step.describe()).await {
443                        publisher.step_said(step, text);
444                    }
445                }
446            })
447            .await;
448        scope.records()
449    }
450
451    /// Whether a model writes this turn's prose: not past the budget, and not when the
452    /// state a failure left could not be read back, since the prose would guess it.
453    fn narrates(&self, input: &CompositionInput<'_>) -> bool {
454        self.narration.enabled
455            && !input.outcomes.budget_exhausted
456            && !input.outcomes.case_refresh_unavailable
457    }
458
459    /// The workflow's guidance for `stage` on the views of `cases`.
460    fn guidance(
461        &self,
462        input: &CompositionInput<'_>,
463        cases: &[CaseKey],
464        stage: WritingStage,
465    ) -> Vec<String> {
466        input
467            .views
468            .iter()
469            .filter(|view| cases.contains(&view.case_ref.key()))
470            .filter_map(|view| {
471                let workflow = self.workflows.require(&view.case_ref.workflow).ok()?;
472                workflow
473                    .definition
474                    .narration_briefing(stage, view)
475                    .ok()
476                    .flatten()
477            })
478            .collect()
479    }
480
481    /// Composes the whole answer (spec §23 steps Q to T).
482    ///
483    /// # Errors
484    ///
485    /// [`OrchestratorError::Erasure`] when receipts cannot be rendered, and
486    /// [`OrchestratorError::Internal`] when the assembled turn fails
487    /// [`claim_guard::verify`], a defect in this module.
488    pub async fn compose(
489        &self,
490        input: CompositionInput<'_>,
491    ) -> Result<Composition, OrchestratorError> {
492        let locale = &input.turn.locale;
493        let receipts = if input.ledger.is_empty() {
494            self.receipts(input.events, locale)?
495        } else {
496            self.ledger_receipts(input.ledger, locale)?
497        };
498        let mut scope = TaskScope::new(
499            input
500                .effort
501                .map_or(self.narration.budget, |effort| effort.reply_budget),
502            locale.clone(),
503        )
504        .for_turn(input.turn.turn_id.to_string());
505        if let Some(effort) = input.effort {
506            scope = scope
507                .with_profiles(effort.tasks.clone())
508                .with_effort(effort.effort);
509        }
510        let narrator = Narrator {
511            engine: &self.engine,
512            scope: &scope,
513            max_chars: self.narration.max_answer_chars,
514        };
515        // A workflow the turn named and cannot start says why, in its own words.
516        let mut facts = input.refusals.to_vec();
517        facts.extend(
518            input
519                .unavailable
520                .iter()
521                .filter(|blocked| input.named_workflows.contains(&blocked.workflow))
522                .map(|blocked| NarratableFact::WorkflowUnavailable {
523                    workflow: blocked.workflow.clone(),
524                    reason: blocked.reason.clone(),
525                }),
526        );
527        let outcome = Material {
528            receipts: &receipts,
529            facts: &facts,
530            interactions: input.interactions,
531            views: input.views,
532            touched: input.touched,
533            beside: input.beside,
534            labels: input.case_labels,
535            disputes: input.disputes,
536            started: input.started,
537            contested: input.contested,
538            next_steps: input.next_steps,
539            locale,
540            copy: &self.ask_copy,
541        }
542        .outcome();
543
544        let answers = answers::Answers {
545            composer: self,
546            narrator: &narrator,
547            input: &input,
548        };
549        let answered = answers.all().await;
550        let notices = self.notices(&input);
551        let mut blocks: Vec<ResponseBlock> = Vec::new();
552        let mut asked = false;
553        if self.narrates(&input) {
554            // What the facts answered is given; a question they do not answer is named, so
555            // the reply says so in its own words, or owns a mistake the user points at.
556            let answer_texts: Vec<String> = answered
557                .iter()
558                .filter(|answer| answer.status == AnswerStatus::Answered)
559                .map(|answer| answer.text.clone())
560                .collect();
561            let unanswered: Vec<String> = answered
562                .iter()
563                .filter(|answer| answer.status != AnswerStatus::Answered)
564                .map(|answer| {
565                    let question = input
566                        .answer_tasks
567                        .iter()
568                        .find(|task| Some(&task.question_id) == answer.question_id.as_ref())
569                        .map_or("", |task| task.question.as_str());
570                    format!("«{question}»")
571                })
572                .collect();
573            let notice_texts: Vec<String> = notices
574                .iter()
575                .map(|notice| notice.text.resolve(locale).to_owned())
576                .collect();
577            let only_answers = !crate::narrate::speaks(&outcome)
578                && notice_texts.is_empty()
579                && unanswered.is_empty();
580            // The turn's one reply: the answers alone as written, else a reply that gives
581            // what was done, the answers, the notices and the ask.
582            let reply = if only_answers {
583                (!answer_texts.is_empty()).then(|| answer_texts.join("\n\n"))
584            } else if let Some(text) = self
585                .acknowledge(
586                    &input,
587                    &outcome,
588                    Carried {
589                        answers: &answer_texts,
590                        unanswered: &unanswered,
591                        notices: &notice_texts,
592                    },
593                    &narrator,
594                )
595                .await
596            {
597                asked = outcome.ask.is_some();
598                Some(text)
599            } else {
600                // Code's own words stand in for a reply that failed, in the reply's order.
601                let done: Vec<&str> = receipts
602                    .iter()
603                    .map(|receipt| receipt.body.resolve(locale))
604                    .collect();
605                let mut parts: Vec<String> = Vec::new();
606                if !done.is_empty() {
607                    parts.push(done.join(" "));
608                }
609                parts.extend(answered.iter().map(|answer| answer.text.clone()));
610                parts.extend(notice_texts);
611                if let Some(ask) = &outcome.ask {
612                    asked = true;
613                    parts.push(ask.question.clone());
614                }
615                (!parts.is_empty()).then(|| parts.join("\n\n"))
616            };
617            if let Some(text) = reply {
618                blocks.push(self.transition(text, &receipts, &answered, &input));
619            }
620        }
621        blocks.extend(answered.into_iter().map(ResponseBlock::Answer));
622
623        for receipt in &receipts {
624            blocks.push(ResponseBlock::Receipt(ReceiptBlock {
625                block_id: BlockId::from(format!("receipt:{}", receipt.receipt_id)),
626                receipt: receipt.clone(),
627            }));
628        }
629        self.artifacts(&input, &mut blocks);
630        blocks.extend(notices.into_iter().map(ResponseBlock::Notice));
631        // The cards last: what happened is read before what is asked.
632        for interaction in input.interactions {
633            blocks.push(ResponseBlock::Interaction(InteractionBlock {
634                block_id: BlockId::from(format!("interaction:{}", interaction.id)),
635                view: interaction.view(),
636            }));
637        }
638
639        let turn = AssistantTurn {
640            turn_id: input.turn.turn_id,
641            conversation_id: input.turn.conversation_id,
642            blocks,
643            subjects: input
644                .views
645                .iter()
646                .filter(|view| input.subjects.contains(&view.case_ref.key()))
647                .map(|view| view.case_ref.clone())
648                .collect(),
649            replay_token: derive_replay_token(&input.turn.turn_id),
650            expectations: Vec::new(),
651            done: Vec::new(),
652        };
653        claim_guard::verify(&turn).map_err(|violation| {
654            tracing::error!(
655                target: "turnframe.compose",
656                violation = %violation,
657                "composed turn failed the claim guard"
658            );
659            OrchestratorError::Internal {
660                code: "claim_guard".to_owned(),
661            }
662        })?;
663        Ok(Composition {
664            turn,
665            tasks: scope.records(),
666            budget: scope.budget_report(),
667            expectation: outcome
668                .ask
669                .filter(|_| asked)
670                .and_then(|ask| ask.expectation),
671        })
672    }
673
674    async fn acknowledge(
675        &self,
676        input: &CompositionInput<'_>,
677        outcome: &TurnOutcome,
678        carried: Carried<'_>,
679        narrator: &Narrator<'_>,
680    ) -> Option<String> {
681        let locale = &input.turn.locale;
682        let window = input.recent.len().saturating_sub(TRANSCRIPT_WINDOW);
683        let guidance = self.guidance(input, input.touched, WritingStage::Transition);
684        let on_screen: Vec<String> = outcome.card.clone().into_iter().collect();
685        // Beside something left undone, the words that asked for it read as done; and
686        // a question's words, shown, get answered twice.
687        let message = input
688            .turn
689            .text
690            .as_deref()
691            .filter(|_| outcome.not_done.is_empty())
692            .and_then(|text| without_questions(text, input.answer_tasks));
693        let acknowledge = AcknowledgeInput {
694            outcome,
695            locale: locale.as_str(),
696            tone: self.narration.tone,
697            message: message.as_deref(),
698            on_screen: &on_screen,
699            answers: carried.answers,
700            unanswered: carried.unanswered,
701            notices: carried.notices,
702            transcript: &input.recent[window..],
703            guidance: &guidance,
704        };
705        narrator.acknowledge(&acknowledge).await
706    }
707
708    /// The acknowledgement block, citing the receipts and cards it stands beside.
709    fn transition(
710        &self,
711        text: String,
712        receipts: &[OperationalReceipt],
713        answered: &[turnframe_core::response::GeneratedAnswer],
714        input: &CompositionInput<'_>,
715    ) -> ResponseBlock {
716        let mut facts_used: Vec<NarratableFact> = receipts
717            .iter()
718            .map(|receipt| NarratableFact::OperationalOutcome {
719                receipt_id: receipt.receipt_id,
720                event_ids: receipt.event_ids.clone(),
721                status_code: receipt.status_code.clone(),
722            })
723            .collect();
724        facts_used.extend(input.interactions.iter().map(|interaction| {
725            NarratableFact::InteractionAvailable {
726                interaction_id: interaction.id,
727                interaction_kind: interaction.kind,
728            }
729        }));
730        facts_used.extend(input.refusals.iter().cloned());
731        // The answers it gives rest on the facts they rest on.
732        for answer in answered {
733            for fact in &answer.facts_used {
734                if !facts_used.contains(fact) {
735                    facts_used.push(fact.clone());
736                }
737            }
738        }
739        ResponseBlock::Transition(GeneratedTransition {
740            block_id: BlockId::from("transition:0"),
741            text,
742            facts_used,
743        })
744    }
745
746    /// The documents of the turn's subjects, each once per revision.
747    fn artifacts(&self, input: &CompositionInput<'_>, blocks: &mut Vec<ResponseBlock>) {
748        for view in input.views {
749            if !input.subjects.contains(&view.case_ref.key()) {
750                continue;
751            }
752            let Ok(workflow) = self.workflows.require(&view.case_ref.workflow) else {
753                continue;
754            };
755            match workflow.definition.artifacts(view) {
756                Ok(artifacts) => {
757                    for artifact in artifacts {
758                        let block_id = BlockId::from(format!(
759                            "artifact:{}:{}:{}:{}",
760                            view.case_ref.workflow,
761                            view.case_ref.case_id,
762                            view.case_ref.expected_revision,
763                            artifact.artifact_id
764                        ));
765                        if !input.artifacts_shown.contains(&block_id) {
766                            blocks
767                                .push(ResponseBlock::Artifact(ArtifactView { block_id, artifact }));
768                        }
769                    }
770                }
771                Err(error) => tracing::warn!(
772                    target: "turnframe.compose",
773                    workflow = view.case_ref.workflow.as_str(),
774                    error = %error,
775                    "a workflow's artifacts could not be read from its own view"
776                ),
777            }
778        }
779    }
780
781    /// The turn's notices, then composition's own, each code once.
782    fn notices(&self, input: &CompositionInput<'_>) -> Vec<ServerNotice> {
783        let flags = input.outcomes;
784        let copy = &self.copy;
785        let own = [
786            (
787                flags.revision_conflict,
788                notice::REVISION_CONFLICT,
789                NoticeSeverity::Warning,
790                &copy.revision_conflict,
791            ),
792            (
793                flags.had_failure,
794                notice::COMMAND_FAILED,
795                NoticeSeverity::Error,
796                &copy.command_failed,
797            ),
798            (
799                flags.outcome_unknown,
800                notice::VERIFICATION_IN_PROGRESS,
801                NoticeSeverity::Warning,
802                &copy.verification_in_progress,
803            ),
804            (
805                flags.interaction_unavailable,
806                notice::INTERACTION_UNAVAILABLE,
807                NoticeSeverity::Error,
808                &copy.interaction_unavailable,
809            ),
810            (
811                flags.case_refresh_unavailable,
812                notice::CASE_REFRESH_UNAVAILABLE,
813                NoticeSeverity::Error,
814                &copy.case_refresh_unavailable,
815            ),
816            (
817                flags.instruction_declined,
818                notice::INSTRUCTION_DECLINED,
819                NoticeSeverity::Info,
820                &copy.instruction_declined,
821            ),
822            (
823                flags.budget_exhausted,
824                notice::BUDGET_EXHAUSTED,
825                NoticeSeverity::Warning,
826                &copy.budget_exhausted,
827            ),
828        ];
829        let mut notices: Vec<ServerNotice> = Vec::new();
830        let raised = own
831            .into_iter()
832            .filter(|(raised, ..)| *raised)
833            .map(|(_, code, severity, text)| notice_of(code, severity, text.clone()));
834        for notice in input.notices.iter().cloned().chain(raised) {
835            if !notices.iter().any(|existing| existing.code == notice.code) {
836                notices.push(notice);
837            }
838        }
839        notices
840    }
841}
842
843fn notice_of(code: &str, severity: NoticeSeverity, text: LocalizedText) -> ServerNotice {
844    ServerNotice {
845        block_id: BlockId::from(format!("notice:{code}")),
846        code: code.to_owned(),
847        severity,
848        text,
849    }
850}
851
852/// `text` without the words its questions were asked in, and the separators they
853/// leave; `None` when nothing else is left.
854fn without_questions(text: &str, tasks: &[AnswerTask]) -> Option<String> {
855    let mut spans: Vec<(usize, usize)> = tasks
856        .iter()
857        .filter_map(|task| task.asked_at.map(|span| (span.start_byte, span.end_byte)))
858        .collect();
859    spans.sort_unstable();
860    let mut kept = String::new();
861    let mut at = 0;
862    for (start, end) in spans {
863        let Some(before) = text.get(at..start) else {
864            continue;
865        };
866        kept.push_str(before);
867        kept.push(' ');
868        at = end;
869    }
870    kept.push_str(text.get(at..).unwrap_or_default());
871    let kept = kept.split_whitespace().collect::<Vec<_>>().join(" ");
872    let kept = kept.trim_matches(|c: char| c.is_whitespace() || matches!(c, ',' | ';' | ':'));
873    (!kept.is_empty()).then(|| kept.to_owned())
874}
875
876#[cfg(test)]
877mod tests {
878    use super::*;
879
880    #[test]
881    fn replay_tokens_are_derived_and_stable() {
882        let turn = TurnId::nil();
883        assert_eq!(derive_replay_token(&turn), derive_replay_token(&turn));
884    }
885}