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