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    /// Whether a knowledge source is configured to answer questions about the domain.
334    pub(crate) const fn has_knowledge(&self) -> bool {
335        self.knowledge.is_some()
336    }
337
338    /// Attaches a knowledge provider (spec §19.2).
339    #[must_use]
340    pub fn with_knowledge(mut self, knowledge: Arc<dyn KnowledgeProvider>) -> Self {
341        self.knowledge = Some(knowledge);
342        self
343    }
344
345    /// Replaces the copy composition writes itself. English and Italian by default; a
346    /// [`LocalizedText`] adds languages with `with`, and the copy's `translated` a whole
347    /// language at once.
348    #[must_use]
349    pub fn with_copy(mut self, copy: CompositionCopy) -> Self {
350        self.copy = copy;
351        self
352    }
353
354    /// Replaces the questions code writes when no model does.
355    #[must_use]
356    pub fn with_ask_copy(mut self, copy: AskCopy) -> Self {
357        self.ask_copy = copy;
358        self
359    }
360
361    /// The server's own sentences this composer writes, for the languages check.
362    pub(crate) fn server_copy(&self) -> [&dyn crate::copy::ServerCopy; 2] {
363        [&self.copy, &self.ask_copy]
364    }
365
366    /// Sets how many knowledge chunks one answer may retrieve.
367    #[must_use]
368    pub const fn with_max_chunks(mut self, max_chunks: Option<usize>) -> Self {
369        self.max_chunks = max_chunks;
370        self
371    }
372
373    /// Renders receipts from a ledger read, carrying every erasure through.
374    ///
375    /// # Errors
376    ///
377    /// [`OrchestratorError::Erasure`] when a workflow cannot read back its own events.
378    pub fn ledger_receipts(
379        &self,
380        groups: &[LedgerReceiptGroup],
381        locale: &Locale,
382    ) -> Result<Vec<OperationalReceipt>, OrchestratorError> {
383        let mut receipts = Vec::new();
384        let mut seen = BTreeSet::new();
385        for group in groups {
386            let registered = self.workflows.require(&group.case_key.workflow)?;
387            for receipt in registered.definition.receipts(&group.events, locale)? {
388                if seen.insert(receipt.receipt_id) {
389                    receipts.push(receipt);
390                }
391            }
392        }
393        Ok(receipts)
394    }
395
396    /// Renders the receipts of a turn from the events that just committed (§17.3,
397    /// I16). Deterministic, so a regenerated response is the one that was lost.
398    ///
399    /// # Errors
400    ///
401    /// [`OrchestratorError::Erasure`] when a workflow cannot read back its own events.
402    pub fn receipts(
403        &self,
404        events: &[EventBatch],
405        locale: &Locale,
406    ) -> Result<Vec<OperationalReceipt>, OrchestratorError> {
407        let mut receipts = Vec::new();
408        let mut seen = BTreeSet::new();
409        for batch in events {
410            let registered = self.workflows.require(&batch.case_key.workflow)?;
411            let committed: Vec<ReceiptEvent<serde_json::Value>> = batch
412                .events
413                .iter()
414                .cloned()
415                .map(ReceiptEvent::Committed)
416                .collect();
417            for receipt in registered.definition.receipts(&committed, locale)? {
418                if seen.insert(receipt.receipt_id) {
419                    receipts.push(receipt);
420                }
421            }
422        }
423        Ok(receipts)
424    }
425
426    /// Says each step `steps` delivers as it arrives, concurrently, and returns the
427    /// calls made. Ends when the sender is dropped.
428    pub(crate) async fn say_steps(
429        &self,
430        steps: futures::channel::mpsc::UnboundedReceiver<turnframe_understand::Step>,
431        locale: &turnframe_core::locale::Locale,
432        turn: turnframe_core::ids::TurnId,
433        publisher: &crate::stream::TurnPublisher,
434    ) -> Vec<turnframe_core::replay::TaskRecord> {
435        use futures::StreamExt as _;
436        let scope =
437            TaskScope::new(self.narration.budget, locale.clone()).for_turn(turn.to_string());
438        let narrator = Narrator {
439            engine: &self.engine,
440            scope: &scope,
441            max_chars: None,
442        };
443        steps
444            .for_each_concurrent(None, |step| {
445                let narrator = &narrator;
446                async move {
447                    if let Some(text) = narrator.step(locale.as_str(), &step.describe()).await {
448                        publisher.step_said(step, text);
449                    }
450                }
451            })
452            .await;
453        scope.records()
454    }
455
456    /// Whether a model writes this turn's prose: not past the budget, and not when the
457    /// state a failure left could not be read back, since the prose would guess it.
458    fn narrates(&self, input: &CompositionInput<'_>) -> bool {
459        self.narration.enabled
460            && !input.outcomes.budget_exhausted
461            && !input.outcomes.case_refresh_unavailable
462    }
463
464    /// The workflow's guidance for `stage` on the views of `cases`.
465    fn guidance(
466        &self,
467        input: &CompositionInput<'_>,
468        cases: &[CaseKey],
469        stage: WritingStage,
470    ) -> Vec<String> {
471        input
472            .views
473            .iter()
474            .filter(|view| cases.contains(&view.case_ref.key()))
475            .filter_map(|view| {
476                let workflow = self.workflows.require(&view.case_ref.workflow).ok()?;
477                workflow
478                    .definition
479                    .narration_briefing(stage, view)
480                    .ok()
481                    .flatten()
482            })
483            .collect()
484    }
485
486    /// Composes the whole answer (spec §23 steps Q to T).
487    ///
488    /// # Errors
489    ///
490    /// [`OrchestratorError::Erasure`] when receipts cannot be rendered, and
491    /// [`OrchestratorError::Internal`] when the assembled turn fails
492    /// [`claim_guard::verify`], a defect in this module.
493    pub async fn compose(
494        &self,
495        input: CompositionInput<'_>,
496    ) -> Result<Composition, OrchestratorError> {
497        let locale = &input.turn.locale;
498        let receipts = if input.ledger.is_empty() {
499            self.receipts(input.events, locale)?
500        } else {
501            self.ledger_receipts(input.ledger, locale)?
502        };
503        let mut scope = TaskScope::new(
504            input
505                .effort
506                .map_or(self.narration.budget, |effort| effort.reply_budget),
507            locale.clone(),
508        )
509        .for_turn(input.turn.turn_id.to_string());
510        if let Some(effort) = input.effort {
511            scope = scope
512                .with_profiles(effort.tasks.clone())
513                .with_effort(effort.effort);
514        }
515        let narrator = Narrator {
516            engine: &self.engine,
517            scope: &scope,
518            max_chars: self.narration.max_answer_chars,
519        };
520        // A workflow the turn named and cannot start says why, in its own words.
521        let mut facts = input.refusals.to_vec();
522        facts.extend(
523            input
524                .unavailable
525                .iter()
526                .filter(|blocked| input.named_workflows.contains(&blocked.workflow))
527                .map(|blocked| NarratableFact::WorkflowUnavailable {
528                    workflow: blocked.workflow.clone(),
529                    reason: blocked.reason.clone(),
530                }),
531        );
532        let outcome = Material {
533            receipts: &receipts,
534            facts: &facts,
535            interactions: input.interactions,
536            views: input.views,
537            touched: input.touched,
538            beside: input.beside,
539            labels: input.case_labels,
540            disputes: input.disputes,
541            started: input.started,
542            contested: input.contested,
543            next_steps: input.next_steps,
544            locale,
545            copy: &self.ask_copy,
546        }
547        .outcome();
548
549        let answers = answers::Answers {
550            composer: self,
551            narrator: &narrator,
552            input: &input,
553        };
554        let answered = answers.all().await;
555        let notices = self.notices(&input);
556        let mut blocks: Vec<ResponseBlock> = Vec::new();
557        let mut asked = false;
558        if self.narrates(&input) {
559            // What the facts answered is given; a question they do not answer is named, so
560            // the reply says so in its own words, or owns a mistake the user points at.
561            let answer_texts: Vec<String> = answered
562                .iter()
563                .filter(|answer| answer.status == AnswerStatus::Answered)
564                .map(|answer| answer.text.clone())
565                .collect();
566            let unanswered: Vec<String> = answered
567                .iter()
568                .filter(|answer| answer.status != AnswerStatus::Answered)
569                .map(|answer| {
570                    let question = input
571                        .answer_tasks
572                        .iter()
573                        .find(|task| Some(&task.question_id) == answer.question_id.as_ref())
574                        .map_or("", |task| task.question.as_str());
575                    format!("«{question}»")
576                })
577                .collect();
578            let notice_texts: Vec<String> = notices
579                .iter()
580                .map(|notice| notice.text.resolve(locale).to_owned())
581                .collect();
582            let only_answers = !crate::narrate::speaks(&outcome)
583                && notice_texts.is_empty()
584                && unanswered.is_empty();
585            // The turn's one reply: the answers alone as written, else a reply that gives
586            // what was done, the answers, the notices and the ask.
587            let reply = if only_answers {
588                // Answers alone, or nothing: the question to go on closes the reply.
589                let parts: Vec<String> = answer_texts
590                    .iter()
591                    .cloned()
592                    .chain(outcome.closing.clone())
593                    .collect();
594                (!parts.is_empty()).then(|| parts.join("\n\n"))
595            } else if let Some(text) = self
596                .acknowledge(
597                    &input,
598                    &outcome,
599                    Carried {
600                        answers: &answer_texts,
601                        unanswered: &unanswered,
602                        notices: &notice_texts,
603                    },
604                    &narrator,
605                )
606                .await
607            {
608                asked = outcome.ask.is_some();
609                Some(text)
610            } else {
611                // Code's own words stand in for a reply that failed, in the reply's order,
612                // and end on the way forward, as every reply does.
613                let done: Vec<&str> = receipts
614                    .iter()
615                    .map(|receipt| receipt.body.resolve(locale))
616                    .collect();
617                let mut parts: Vec<String> = Vec::new();
618                if !done.is_empty() {
619                    parts.push(done.join(" "));
620                }
621                if !outcome.not_done.is_empty() {
622                    parts.push(outcome.not_done.join(" "));
623                }
624                parts.extend(answered.iter().map(|answer| answer.text.clone()));
625                parts.extend(notice_texts);
626                if let Some(ask) = &outcome.ask {
627                    asked = true;
628                    parts.push(ask.question.clone());
629                } else if !outcome.next.is_empty() {
630                    let go_on = self.ask_copy.go_on.resolve(locale);
631                    parts.push(format!("{go_on} {}", outcome.next.join(" ")));
632                } else if let Some(closing) = &outcome.closing {
633                    parts.push(closing.clone());
634                }
635                (!parts.is_empty()).then(|| parts.join("\n\n"))
636            };
637            if let Some(text) = reply {
638                blocks.push(self.transition(text, &receipts, &answered, &input));
639            }
640        }
641        blocks.extend(answered.into_iter().map(ResponseBlock::Answer));
642
643        for receipt in &receipts {
644            blocks.push(ResponseBlock::Receipt(ReceiptBlock {
645                block_id: BlockId::from(format!("receipt:{}", receipt.receipt_id)),
646                receipt: receipt.clone(),
647            }));
648        }
649        self.artifacts(&input, &mut blocks);
650        blocks.extend(notices.into_iter().map(ResponseBlock::Notice));
651        // The cards last: what happened is read before what is asked.
652        for interaction in input.interactions {
653            blocks.push(ResponseBlock::Interaction(InteractionBlock {
654                block_id: BlockId::from(format!("interaction:{}", interaction.id)),
655                view: interaction.view(),
656            }));
657        }
658
659        let turn = AssistantTurn {
660            turn_id: input.turn.turn_id,
661            conversation_id: input.turn.conversation_id,
662            blocks,
663            subjects: input
664                .views
665                .iter()
666                .filter(|view| input.subjects.contains(&view.case_ref.key()))
667                .map(|view| view.case_ref.clone())
668                .collect(),
669            replay_token: derive_replay_token(&input.turn.turn_id),
670            expectations: Vec::new(),
671            done: Vec::new(),
672        };
673        claim_guard::verify(&turn).map_err(|violation| {
674            tracing::error!(
675                target: "turnframe.compose",
676                violation = %violation,
677                "composed turn failed the claim guard"
678            );
679            OrchestratorError::Internal {
680                code: "claim_guard".to_owned(),
681            }
682        })?;
683        Ok(Composition {
684            turn,
685            tasks: scope.records(),
686            budget: scope.budget_report(),
687            expectation: outcome
688                .ask
689                .filter(|_| asked)
690                .and_then(|ask| ask.expectation),
691        })
692    }
693
694    async fn acknowledge(
695        &self,
696        input: &CompositionInput<'_>,
697        outcome: &TurnOutcome,
698        carried: Carried<'_>,
699        narrator: &Narrator<'_>,
700    ) -> Option<String> {
701        let locale = &input.turn.locale;
702        let window = input.recent.len().saturating_sub(TRANSCRIPT_WINDOW);
703        let guidance = self.guidance(input, input.touched, WritingStage::Transition);
704        let on_screen: Vec<String> = outcome.card.clone().into_iter().collect();
705        // Beside something left undone, the words that asked for it read as done; and
706        // a question's words, shown, get answered twice.
707        let message = input
708            .turn
709            .text
710            .as_deref()
711            .filter(|_| outcome.not_done.is_empty())
712            .and_then(|text| without_questions(text, input.answer_tasks));
713        let acknowledge = AcknowledgeInput {
714            outcome,
715            locale: locale.as_str(),
716            tone: self.narration.tone,
717            message: message.as_deref(),
718            on_screen: &on_screen,
719            answers: carried.answers,
720            unanswered: carried.unanswered,
721            notices: carried.notices,
722            transcript: &input.recent[window..],
723            guidance: &guidance,
724        };
725        narrator.acknowledge(&acknowledge).await
726    }
727
728    /// The acknowledgement block, citing the receipts and cards it stands beside.
729    fn transition(
730        &self,
731        text: String,
732        receipts: &[OperationalReceipt],
733        answered: &[turnframe_core::response::GeneratedAnswer],
734        input: &CompositionInput<'_>,
735    ) -> ResponseBlock {
736        let mut facts_used: Vec<NarratableFact> = receipts
737            .iter()
738            .map(|receipt| NarratableFact::OperationalOutcome {
739                receipt_id: receipt.receipt_id,
740                event_ids: receipt.event_ids.clone(),
741                status_code: receipt.status_code.clone(),
742            })
743            .collect();
744        facts_used.extend(input.interactions.iter().map(|interaction| {
745            NarratableFact::InteractionAvailable {
746                interaction_id: interaction.id,
747                interaction_kind: interaction.kind,
748            }
749        }));
750        facts_used.extend(input.refusals.iter().cloned());
751        // The answers it gives rest on the facts they rest on.
752        for answer in answered {
753            for fact in &answer.facts_used {
754                if !facts_used.contains(fact) {
755                    facts_used.push(fact.clone());
756                }
757            }
758        }
759        ResponseBlock::Transition(GeneratedTransition {
760            block_id: BlockId::from("transition:0"),
761            text,
762            facts_used,
763        })
764    }
765
766    /// The documents of the turn's subjects, each once per revision.
767    fn artifacts(&self, input: &CompositionInput<'_>, blocks: &mut Vec<ResponseBlock>) {
768        for view in input.views {
769            if !input.subjects.contains(&view.case_ref.key()) {
770                continue;
771            }
772            let Ok(workflow) = self.workflows.require(&view.case_ref.workflow) else {
773                continue;
774            };
775            match workflow.definition.artifacts(view) {
776                Ok(artifacts) => {
777                    for artifact in artifacts {
778                        let block_id = BlockId::from(format!(
779                            "artifact:{}:{}:{}:{}",
780                            view.case_ref.workflow,
781                            view.case_ref.case_id,
782                            view.case_ref.expected_revision,
783                            artifact.artifact_id
784                        ));
785                        if !input.artifacts_shown.contains(&block_id) {
786                            blocks
787                                .push(ResponseBlock::Artifact(ArtifactView { block_id, artifact }));
788                        }
789                    }
790                }
791                Err(error) => tracing::warn!(
792                    target: "turnframe.compose",
793                    workflow = view.case_ref.workflow.as_str(),
794                    error = %error,
795                    "a workflow's artifacts could not be read from its own view"
796                ),
797            }
798        }
799    }
800
801    /// The turn's notices, then composition's own, each code once.
802    fn notices(&self, input: &CompositionInput<'_>) -> Vec<ServerNotice> {
803        let flags = input.outcomes;
804        let copy = &self.copy;
805        let own = [
806            (
807                flags.revision_conflict,
808                notice::REVISION_CONFLICT,
809                NoticeSeverity::Warning,
810                &copy.revision_conflict,
811            ),
812            (
813                flags.had_failure,
814                notice::COMMAND_FAILED,
815                NoticeSeverity::Error,
816                &copy.command_failed,
817            ),
818            (
819                flags.outcome_unknown,
820                notice::VERIFICATION_IN_PROGRESS,
821                NoticeSeverity::Warning,
822                &copy.verification_in_progress,
823            ),
824            (
825                flags.interaction_unavailable,
826                notice::INTERACTION_UNAVAILABLE,
827                NoticeSeverity::Error,
828                &copy.interaction_unavailable,
829            ),
830            (
831                flags.case_refresh_unavailable,
832                notice::CASE_REFRESH_UNAVAILABLE,
833                NoticeSeverity::Error,
834                &copy.case_refresh_unavailable,
835            ),
836            (
837                flags.instruction_declined,
838                notice::INSTRUCTION_DECLINED,
839                NoticeSeverity::Info,
840                &copy.instruction_declined,
841            ),
842            (
843                flags.budget_exhausted,
844                notice::BUDGET_EXHAUSTED,
845                NoticeSeverity::Warning,
846                &copy.budget_exhausted,
847            ),
848        ];
849        let mut notices: Vec<ServerNotice> = Vec::new();
850        let raised = own
851            .into_iter()
852            .filter(|(raised, ..)| *raised)
853            .map(|(_, code, severity, text)| notice_of(code, severity, text.clone()));
854        for notice in input.notices.iter().cloned().chain(raised) {
855            if !notices.iter().any(|existing| existing.code == notice.code) {
856                notices.push(notice);
857            }
858        }
859        notices
860    }
861}
862
863fn notice_of(code: &str, severity: NoticeSeverity, text: LocalizedText) -> ServerNotice {
864    ServerNotice {
865        block_id: BlockId::from(format!("notice:{code}")),
866        code: code.to_owned(),
867        severity,
868        text,
869    }
870}
871
872/// `text` without the words its questions were asked in, and the separators they
873/// leave; `None` when nothing else is left.
874fn without_questions(text: &str, tasks: &[AnswerTask]) -> Option<String> {
875    let mut spans: Vec<(usize, usize)> = tasks
876        .iter()
877        .filter_map(|task| task.asked_at.map(|span| (span.start_byte, span.end_byte)))
878        .collect();
879    spans.sort_unstable();
880    let mut kept = String::new();
881    let mut at = 0;
882    for (start, end) in spans {
883        let Some(before) = text.get(at..start) else {
884            continue;
885        };
886        kept.push_str(before);
887        kept.push(' ');
888        at = end;
889    }
890    kept.push_str(text.get(at..).unwrap_or_default());
891    let kept = kept.split_whitespace().collect::<Vec<_>>().join(" ");
892    let kept = kept.trim_matches(|c: char| c.is_whitespace() || matches!(c, ',' | ';' | ':'));
893    (!kept.is_empty()).then(|| kept.to_owned())
894}
895
896#[cfg(test)]
897mod tests {
898    use super::*;
899
900    #[test]
901    fn replay_tokens_are_derived_and_stable() {
902        let turn = TurnId::nil();
903        assert_eq!(derive_replay_token(&turn), derive_replay_token(&turn));
904    }
905}