1mod 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
54const REPLAY_TOKEN_DOMAIN: &str = "turnframe.replay_token.v1";
56
57const TRANSCRIPT_WINDOW: usize = 4;
59
60struct Carried<'a> {
62 answers: &'a [String],
63 unanswered: &'a [String],
64 notices: &'a [String],
65}
66
67#[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#[derive(Debug, Clone)]
80#[non_exhaustive]
81pub struct CompositionInput<'a> {
82 pub turn: &'a TurnInput,
84 pub attachments: Vec<ContentPart>,
86 pub answer_tasks: &'a [AnswerTask],
88 pub events: &'a [EventBatch],
90 pub ledger: &'a [LedgerReceiptGroup],
93 pub interactions: &'a [Interaction],
95 pub views: &'a [ErasedWorkflowView],
97 pub subjects: &'a [CaseKey],
99 pub touched: &'a [CaseKey],
101 pub beside: &'a [CaseKey],
103 pub reachable_only: &'a [CaseKey],
105 pub recent: &'a [RecentMessage],
107 pub preceding_reply: Option<&'a str>,
109 pub artifacts_shown: &'a [BlockId],
111 pub case_labels: &'a [CaseLabel],
113 pub next_steps: &'a [(CaseKey, Vec<LocalizedText>)],
115 pub named_workflows: &'a [WorkflowKey],
117 pub unavailable: &'a [UnavailableWorkflow],
119 pub disputes: &'a [String],
121 pub started: &'a [WorkflowKey],
123 pub contested: &'a [String],
125 pub notices: &'a [ServerNotice],
127 pub refusals: &'a [NarratableFact],
129 pub outcomes: OutcomeFlags,
131 pub effort: Option<&'a crate::effort::EffortProfile>,
133}
134
135use turnframe_core::ids::WorkflowKey;
136
137#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
139#[non_exhaustive]
140pub struct OutcomeFlags {
141 pub had_failure: bool,
143 pub outcome_unknown: bool,
145 pub revision_conflict: bool,
147 pub interaction_unavailable: bool,
149 pub case_refresh_unavailable: bool,
151 pub budget_exhausted: bool,
153 pub instruction_declined: bool,
155}
156
157impl<'a> CompositionInput<'a> {
158 #[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 with_attachments: attachments: Vec<ContentPart>;
209 with_answer_tasks: answer_tasks: &'a [AnswerTask];
211 with_events: events: &'a [EventBatch];
213 with_ledger: ledger: &'a [LedgerReceiptGroup];
215 with_interactions: interactions: &'a [Interaction];
217 with_views: views: &'a [ErasedWorkflowView];
219 with_subjects: subjects: &'a [CaseKey];
221 with_touched: touched: &'a [CaseKey];
223 with_beside: beside: &'a [CaseKey];
225 with_reachable_only: reachable_only: &'a [CaseKey];
227 with_recent: recent: &'a [RecentMessage];
229 with_preceding_reply: preceding_reply: Option<&'a str>;
231 with_artifacts_shown: artifacts_shown: &'a [BlockId];
233 with_case_labels: case_labels: &'a [CaseLabel];
235 with_next_steps: next_steps: &'a [(CaseKey, Vec<LocalizedText>)];
237 with_named_workflows: named_workflows: &'a [WorkflowKey];
239 with_unavailable: unavailable: &'a [UnavailableWorkflow];
241 with_disputes: disputes: &'a [String];
243 with_started: started: &'a [WorkflowKey];
245 with_contested: contested: &'a [String];
247 with_notices: notices: &'a [ServerNotice];
249 with_refusals: refusals: &'a [NarratableFact];
251 with_outcomes: outcomes: OutcomeFlags;
253 with_effort: effort: Option<&'a crate::effort::EffortProfile>;
255}
256
257#[derive(Debug, Clone)]
259#[non_exhaustive]
260pub struct Composition {
261 pub turn: AssistantTurn,
263 pub tasks: Vec<TaskRecord>,
265 pub budget: BudgetReport,
267 pub expectation: Option<Expectation>,
269}
270
271#[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 #[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 #[must_use]
315 pub fn with_tasks(mut self, engine: TaskEngine) -> Self {
316 self.engine = engine;
317 self
318 }
319
320 #[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 #[must_use]
335 pub fn with_knowledge(mut self, knowledge: Arc<dyn KnowledgeProvider>) -> Self {
336 self.knowledge = Some(knowledge);
337 self
338 }
339
340 #[must_use]
344 pub fn with_copy(mut self, copy: CompositionCopy) -> Self {
345 self.copy = copy;
346 self
347 }
348
349 #[must_use]
351 pub fn with_ask_copy(mut self, copy: AskCopy) -> Self {
352 self.ask_copy = copy;
353 self
354 }
355
356 pub(crate) fn server_copy(&self) -> [&dyn crate::copy::ServerCopy; 2] {
358 [&self.copy, &self.ask_copy]
359 }
360
361 #[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 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 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 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 fn narrates(&self, input: &CompositionInput<'_>) -> bool {
454 self.narration.enabled
455 && !input.outcomes.budget_exhausted
456 && !input.outcomes.case_refresh_unavailable
457 }
458
459 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 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 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 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 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: ¬ice_texts,
592 },
593 &narrator,
594 )
595 .await
596 {
597 asked = outcome.ask.is_some();
598 Some(text)
599 } else {
600 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 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 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 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 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 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 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 ©.revision_conflict,
791 ),
792 (
793 flags.had_failure,
794 notice::COMMAND_FAILED,
795 NoticeSeverity::Error,
796 ©.command_failed,
797 ),
798 (
799 flags.outcome_unknown,
800 notice::VERIFICATION_IN_PROGRESS,
801 NoticeSeverity::Warning,
802 ©.verification_in_progress,
803 ),
804 (
805 flags.interaction_unavailable,
806 notice::INTERACTION_UNAVAILABLE,
807 NoticeSeverity::Error,
808 ©.interaction_unavailable,
809 ),
810 (
811 flags.case_refresh_unavailable,
812 notice::CASE_REFRESH_UNAVAILABLE,
813 NoticeSeverity::Error,
814 ©.case_refresh_unavailable,
815 ),
816 (
817 flags.instruction_declined,
818 notice::INSTRUCTION_DECLINED,
819 NoticeSeverity::Info,
820 ©.instruction_declined,
821 ),
822 (
823 flags.budget_exhausted,
824 notice::BUDGET_EXHAUSTED,
825 NoticeSeverity::Warning,
826 ©.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
852fn 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}