Skip to main content

turnframe_runtime/
planning.rs

1//! Planning a turn without performing it (spec §23 steps A–K).
2//!
3//! An application migrating from another agent runs both paths on the same turns while
4//! the old one stays authoritative. [`TurnPlanner::plan`] runs steps A through K — accept,
5//! load and project, judge a card answer, issue tokens, understand, resolve, reduce and
6//! police — and stops before step L: no card, no journal entry, no event, no replay
7//! record is written. That is a property of the types: a planner holds
8//! [`ReadOnlyStores`] and loaders, never an executor.
9//!
10//! [`SeededTurnPlanner`] takes the case state as input instead of loading it, so a
11//! recorded turn replays against the state that preceded it. Traffic routing and kill
12//! switches belong to the application; [`crate::divergence`] names what two paths
13//! disagreed about.
14
15use std::fmt;
16use std::sync::Arc;
17
18use chrono::{DateTime, Utc};
19use indexmap::IndexMap;
20use turnframe_core::case::{CaseKey, CaseRef};
21use turnframe_core::error::{AuthorizationError, InteractionError, OrchestratorError, StoreError};
22use turnframe_core::flow::{ErasedWorkflowView, WorkflowDefinitions, WorkflowReadRegistry};
23use turnframe_core::interaction::{Interaction, InteractionRejection, InteractionSpec};
24use turnframe_core::observe::{NoopObserver, Observer, Signal, SignalLabels};
25use turnframe_core::policy::{PolicyDecision, PolicySnapshot};
26use turnframe_core::reduce::{CommandRef, ReductionPlan};
27use turnframe_core::replay::{TargetResolutionRecord, TaskRecord};
28use turnframe_core::response::{AssistantTurn, ClaimClass};
29use turnframe_core::turn::TurnInput;
30use turnframe_core::understanding::{ActId, Understanding};
31use turnframe_store::interaction::InteractionRecord;
32use turnframe_store::stores::ReadOnlyStores;
33use turnframe_understand::{NoSteps, TurnUnderstander};
34
35use crate::config::OrchestratorConfig;
36use crate::conversation::RecentMessage;
37use crate::divergence::TurnSummary;
38use crate::interactions::AcceptedInteraction;
39use crate::orchestrator::{CaseDirectory, SystemTurnClock, TurnClock};
40use crate::policy::PolicyEngine;
41use crate::resolve::{CaseIdFactory, DerivedCaseIdFactory};
42pub(crate) use crate::turn::{DirectoryTerms, LoadedCase};
43use crate::turn::{
44    addressable_cards, admit, admit_card_case, admit_typed, aim_at_found, blocking_summary,
45    build_resolver, candidates_of_open_cards, card_act_of, command_refs, project_case,
46    recent_messages, reduction_context, target_resolution_records, turn_reducer, typed_response,
47    unlisted, would_claim,
48};
49
50/// A case handed to [`SeededTurnPlanner::plan`] instead of being loaded: the state a
51/// recorded turn preceded, in the erased form the registry boundary uses.
52#[derive(Debug, Clone, PartialEq, Eq)]
53#[non_exhaustive]
54pub struct SeededCase {
55    /// The case and the revision the state was read at.
56    pub case_ref: CaseRef,
57    /// The state, or `None` for a case that does not exist yet.
58    pub state: Option<serde_json::Value>,
59    /// The label the model sees; the case identifier when absent.
60    pub label: Option<String>,
61    /// Cards open on the case.
62    pub open_interactions: Vec<Interaction>,
63    /// Whether the case is in view only because the actor may reach it.
64    pub subject_only_when_named: bool,
65}
66
67impl SeededCase {
68    /// A case at `case_ref` holding `state`.
69    #[must_use]
70    pub fn new(case_ref: CaseRef, state: serde_json::Value) -> Self {
71        Self {
72            case_ref,
73            state: Some(state),
74            label: None,
75            open_interactions: Vec::new(),
76            subject_only_when_named: false,
77        }
78    }
79
80    /// A case that does not exist yet.
81    #[must_use]
82    pub fn absent(case_ref: CaseRef) -> Self {
83        Self {
84            case_ref,
85            state: None,
86            label: None,
87            open_interactions: Vec::new(),
88            subject_only_when_named: false,
89        }
90    }
91
92    /// Sets the label.
93    #[must_use]
94    pub fn with_label(mut self, label: impl Into<String>) -> Self {
95        self.label = Some(label.into());
96        self
97    }
98
99    /// Declares the case in view only because the actor may reach it.
100    #[must_use]
101    pub const fn reachable_only(mut self) -> Self {
102        self.subject_only_when_named = true;
103        self
104    }
105
106    /// Adds a card open on the case.
107    #[must_use]
108    pub fn with_open_interaction(mut self, interaction: Interaction) -> Self {
109        self.open_interactions.push(interaction);
110        self
111    }
112}
113
114/// What a turn would have done, stopped before its first effect.
115#[derive(Debug, Clone)]
116#[non_exhaustive]
117pub struct PlannedTurn {
118    /// The projections of the cases the turn addressed.
119    pub views: Vec<ErasedWorkflowView>,
120    /// What the turn was understood to say; `None` for a click with no text.
121    pub understanding: Option<Understanding>,
122    /// Every model task understanding ran.
123    pub tasks: Vec<TaskRecord>,
124    /// How each act's target resolved.
125    pub target_resolutions: Vec<TargetResolutionRecord>,
126    /// The reduction.
127    pub reduction: ReductionPlan,
128    /// The policy decision for every command.
129    pub policy_decisions: Vec<PolicyDecision>,
130    /// The cards the turn would write before anything runs.
131    pub would_persist: Vec<InteractionSpec>,
132    /// The commands it would execute now.
133    pub would_execute: Vec<CommandRef>,
134    /// The claim classes its answer would be entitled to.
135    pub would_claim: Vec<ClaimClass>,
136}
137
138impl PlannedTurn {
139    /// The summary two paths are compared on.
140    #[must_use]
141    pub fn summary(&self) -> TurnSummary {
142        TurnSummary::from_planned(self)
143    }
144}
145
146/// Plans turns against live, read-only stores.
147pub struct TurnPlanner {
148    workflows: WorkflowReadRegistry,
149    stores: ReadOnlyStores,
150    directory: Arc<dyn CaseDirectory>,
151    shared: SharedPlanning,
152}
153
154impl fmt::Debug for TurnPlanner {
155    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
156        f.debug_struct("TurnPlanner")
157            .field("workflows", &self.workflows)
158            .field("stores", &self.stores)
159            .finish_non_exhaustive()
160    }
161}
162
163/// Plans turns against state it is handed.
164pub struct SeededTurnPlanner {
165    definitions: WorkflowDefinitions,
166    shared: SharedPlanning,
167}
168
169impl fmt::Debug for SeededTurnPlanner {
170    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
171        f.debug_struct("SeededTurnPlanner")
172            .field("definitions", &self.definitions)
173            .finish_non_exhaustive()
174    }
175}
176
177/// What both planners share with the orchestrator that built them.
178#[derive(Clone)]
179pub(crate) struct SharedPlanning {
180    pub(crate) understander: Arc<dyn TurnUnderstander>,
181    pub(crate) policy_engine: PolicyEngine,
182    pub(crate) policy: PolicySnapshot,
183    pub(crate) config: OrchestratorConfig,
184    pub(crate) clock: Arc<dyn TurnClock>,
185    pub(crate) case_ids: Arc<dyn CaseIdFactory>,
186    pub(crate) observer: Arc<dyn Observer>,
187    /// Whether a knowledge source can answer a question about the domain in general.
188    pub(crate) knowledge: bool,
189}
190
191impl TurnPlanner {
192    /// Starts a builder.
193    #[must_use]
194    pub fn builder() -> TurnPlannerBuilder {
195        TurnPlannerBuilder::default()
196    }
197
198    pub(crate) fn assemble(
199        workflows: WorkflowReadRegistry,
200        stores: ReadOnlyStores,
201        directory: Arc<dyn CaseDirectory>,
202        shared: SharedPlanning,
203    ) -> Self {
204        Self {
205            workflows,
206            stores,
207            directory,
208            shared,
209        }
210    }
211
212    /// The seeded planner over the same definitions and parts.
213    #[must_use]
214    pub fn seeded(&self) -> SeededTurnPlanner {
215        SeededTurnPlanner {
216            definitions: self.workflows.definitions().clone(),
217            shared: self.shared.clone(),
218        }
219    }
220
221    /// Plans `input` against the stores, stopping before the first effect.
222    ///
223    /// # Errors
224    ///
225    /// Whatever the same steps of a real turn would fail with.
226    pub async fn plan(&self, input: TurnInput) -> Result<PlannedTurn, OrchestratorError> {
227        let now = self.shared.clock.now();
228        input.validate_shape_within(&self.shared.config.understanding.turn_limits)?;
229        let account = input.actor.account_id.clone();
230        self.stores
231            .conversations()
232            .load_conversation(&account, &input.conversation_id)
233            .await
234            .map_err(|error| match error {
235                StoreError::NotFound => {
236                    OrchestratorError::Unauthorized(AuthorizationError::ConversationNotAccessible {
237                        conversation_id: input.conversation_id,
238                    })
239                }
240                other => OrchestratorError::Store(other),
241            })?;
242        let open_interactions = self
243            .stores
244            .interactions()
245            .list_open_for_conversation(&account, &input.conversation_id)
246            .await
247            .map_err(OrchestratorError::Store)?;
248        let (cases, origin_case) = self.load_cases(&input, &open_interactions).await?;
249        let open_interactions = addressable_cards(&cases, open_interactions);
250        let answered = match input.interaction_response.as_ref() {
251            None => None,
252            Some(response) => {
253                let record = self
254                    .stores
255                    .interactions()
256                    .get(&account, &response.interaction_id)
257                    .await
258                    .map_err(|error| match error {
259                        StoreError::NotFound => OrchestratorError::Interaction(
260                            InteractionError::Rejected(InteractionRejection::NotFound),
261                        ),
262                        _ => OrchestratorError::Interaction(InteractionError::NotPersisted),
263                    })?;
264                admit(&input, &cases, record, now)?
265            }
266        };
267        let turns = self
268            .stores
269            .conversations()
270            .load_recent_turns(
271                &account,
272                &input.conversation_id,
273                self.shared
274                    .config
275                    .understanding
276                    .transcript_turns
277                    .unwrap_or(usize::MAX),
278            )
279            .await
280            .unwrap_or_default();
281        let (recent, previous) = recent_messages(turns, input.turn_id);
282        plan_with(
283            &self.shared,
284            self.workflows.definitions(),
285            &input,
286            PlanningContext {
287                lookup: Some(Lookup {
288                    directory: self.directory.as_ref(),
289                    workflows: &self.workflows,
290                }),
291                now,
292                cases,
293                open_interactions,
294                origin_case,
295                answered,
296                recent,
297                previous,
298            },
299        )
300        .await
301    }
302
303    async fn load_cases(
304        &self,
305        input: &TurnInput,
306        open_interactions: &[Interaction],
307    ) -> Result<(IndexMap<CaseKey, LoadedCase>, Option<CaseKey>), OrchestratorError> {
308        let account = &input.actor.account_id;
309        let mut candidates = self
310            .directory
311            .candidates(&input.actor, &input.conversation_id)
312            .await
313            .map_err(OrchestratorError::Store)?;
314        let mut origin_case = None;
315        if let Some(origin) = input.origin.as_ref()
316            && let Some(candidate) = self
317                .directory
318                .resolve_origin(&input.actor, &origin.origin_token)
319                .await
320                .map_err(OrchestratorError::Store)?
321        {
322            origin_case = Some(candidate.key.clone());
323            candidates.push(candidate);
324        }
325        let from_cards = candidates_of_open_cards(&candidates, open_interactions);
326        let mut cases = IndexMap::new();
327        for (candidate, from_card) in candidates
328            .into_iter()
329            .map(|candidate| (candidate, false))
330            .chain(from_cards.into_iter().map(|candidate| (candidate, true)))
331        {
332            if cases.contains_key(&candidate.key) {
333                continue;
334            }
335            let Ok(definition) = self
336                .workflows
337                .definitions()
338                .require(&candidate.key.workflow)
339                .cloned()
340            else {
341                continue;
342            };
343            let loader = self.workflows.require_loader(&candidate.key.workflow)?;
344            let loaded = loader
345                .load_case(account, &candidate.key.case_id)
346                .await
347                .map_err(OrchestratorError::Store)?;
348            let candidate = if from_card {
349                let Some(authorized) = admit_card_case(
350                    self.shared.observer.as_ref(),
351                    self.directory.as_ref(),
352                    &input.actor,
353                    &input.conversation_id,
354                    candidate,
355                    loaded.value.is_some(),
356                )
357                .await
358                .map_err(OrchestratorError::Store)?
359                else {
360                    continue;
361                };
362                authorized
363            } else {
364                candidate
365            };
366            let case_ref = candidate.key.clone().at(loaded.revision);
367            cases.insert(
368                candidate.key,
369                project_case(
370                    self.shared.observer.as_ref(),
371                    &definition,
372                    case_ref,
373                    candidate.label,
374                    loaded.value,
375                    DirectoryTerms {
376                        confirm_every_write: candidate.confirm_every_write,
377                        subject_only_when_named: candidate.subject_only_when_named,
378                    },
379                )?,
380            );
381        }
382        Ok((cases, origin_case))
383    }
384}
385
386impl SeededTurnPlanner {
387    /// Starts a builder.
388    #[must_use]
389    pub fn builder() -> SeededTurnPlannerBuilder {
390        SeededTurnPlannerBuilder::default()
391    }
392
393    /// Plans `input` against `cases`, stopping before the first effect.
394    ///
395    /// # Errors
396    ///
397    /// Whatever the same steps of a real turn would fail with.
398    pub async fn plan(
399        &self,
400        input: TurnInput,
401        cases: Vec<SeededCase>,
402    ) -> Result<PlannedTurn, OrchestratorError> {
403        let now = self.shared.clock.now();
404        input.validate_shape_within(&self.shared.config.understanding.turn_limits)?;
405        let mut loaded = IndexMap::new();
406        let mut open_interactions = Vec::new();
407        for seeded in cases {
408            let key = seeded.case_ref.key();
409            let definition = self.definitions.require(&key.workflow)?.clone();
410            open_interactions.extend(seeded.open_interactions);
411            let label = seeded
412                .label
413                .unwrap_or_else(|| seeded.case_ref.case_id.to_string());
414            loaded.insert(
415                key,
416                project_case(
417                    self.shared.observer.as_ref(),
418                    &definition,
419                    seeded.case_ref,
420                    label,
421                    seeded.state,
422                    DirectoryTerms {
423                        confirm_every_write: false,
424                        subject_only_when_named: seeded.subject_only_when_named,
425                    },
426                )?,
427            );
428        }
429        let open_interactions = addressable_cards(&loaded, open_interactions);
430        let answered = match input.interaction_response.as_ref() {
431            None => None,
432            Some(response) => {
433                let Some(interaction) = open_interactions
434                    .iter()
435                    .find(|open| open.id == response.interaction_id)
436                    .cloned()
437                else {
438                    return Err(OrchestratorError::Interaction(InteractionError::Rejected(
439                        InteractionRejection::NotFound,
440                    )));
441                };
442                admit(&input, &loaded, InteractionRecord::new(interaction), now)?
443            }
444        };
445        plan_with(
446            &self.shared,
447            &self.definitions,
448            &input,
449            PlanningContext {
450                lookup: None,
451                now,
452                cases: loaded,
453                open_interactions,
454                origin_case: None,
455                answered,
456                recent: Vec::new(),
457                previous: None,
458            },
459        )
460        .await
461    }
462}
463
464struct PlanningContext<'a> {
465    lookup: Option<Lookup<'a>>,
466    now: DateTime<Utc>,
467    cases: IndexMap<CaseKey, LoadedCase>,
468    open_interactions: Vec<Interaction>,
469    origin_case: Option<CaseKey>,
470    answered: Option<AcceptedInteraction>,
471    recent: Vec<RecentMessage>,
472    previous: Option<AssistantTurn>,
473}
474
475/// Where a plan looks up a record the message named and the turn did not have.
476struct Lookup<'a> {
477    directory: &'a dyn CaseDirectory,
478    workflows: &'a WorkflowReadRegistry,
479}
480
481impl Lookup<'_> {
482    /// Loads what the directory finds for each unlisted act, and returns it per act.
483    async fn find(
484        &self,
485        shared: &SharedPlanning,
486        input: &TurnInput,
487        understanding: &Understanding,
488        cases: &mut IndexMap<CaseKey, LoadedCase>,
489    ) -> Result<std::collections::BTreeMap<ActId, Vec<CaseKey>>, OrchestratorError> {
490        let text = input.text.as_deref().unwrap_or_default();
491        let limit = shared.config.interaction.max_selection_candidates;
492        let mut found = std::collections::BTreeMap::new();
493        for (act, workflow, named) in unlisted(understanding, text) {
494            let candidates = self
495                .directory
496                .find(
497                    &input.actor,
498                    &input.conversation_id,
499                    &workflow,
500                    named.as_deref(),
501                )
502                .await
503                .map_err(OrchestratorError::Store)?;
504            let mut keys = Vec::new();
505            for candidate in candidates.into_iter().take(limit) {
506                keys.push(candidate.key.clone());
507                if let Some(case) = load_listed(shared, self.workflows, input, candidate).await? {
508                    cases.insert(case.case_ref.key(), case);
509                }
510            }
511            found.insert(act, keys);
512        }
513        Ok(found)
514    }
515}
516
517/// One candidate the directory listed, loaded and projected; `None` for a workflow
518/// this runtime does not host.
519async fn load_listed(
520    shared: &SharedPlanning,
521    workflows: &WorkflowReadRegistry,
522    input: &TurnInput,
523    candidate: crate::orchestrator::CaseCandidate,
524) -> Result<Option<LoadedCase>, OrchestratorError> {
525    let Ok(definition) = workflows
526        .definitions()
527        .require(&candidate.key.workflow)
528        .cloned()
529    else {
530        return Ok(None);
531    };
532    let loaded = workflows
533        .require_loader(&candidate.key.workflow)?
534        .load_case(&input.actor.account_id, &candidate.key.case_id)
535        .await
536        .map_err(OrchestratorError::Store)?;
537    let case_ref = candidate.key.clone().at(loaded.revision);
538    project_case(
539        shared.observer.as_ref(),
540        &definition,
541        case_ref,
542        candidate.label,
543        loaded.value,
544        DirectoryTerms {
545            confirm_every_write: candidate.confirm_every_write,
546            subject_only_when_named: candidate.subject_only_when_named,
547        },
548    )
549    .map(Some)
550}
551
552async fn plan_with(
553    shared: &SharedPlanning,
554    definitions: &WorkflowDefinitions,
555    input: &TurnInput,
556    context: PlanningContext<'_>,
557) -> Result<PlannedTurn, OrchestratorError> {
558    let PlanningContext {
559        lookup,
560        now,
561        mut cases,
562        open_interactions,
563        origin_case,
564        mut answered,
565        recent,
566        previous,
567    } = context;
568    let origin = input
569        .origin
570        .as_ref()
571        .map(|origin| &origin.origin_token)
572        .zip(origin_case.as_ref());
573    let mut resolver = build_resolver(
574        &input.actor.account_id,
575        input.turn_id,
576        definitions,
577        &shared.case_ids,
578        &cases,
579        origin,
580        blocking_summary(answered.as_ref(), &open_interactions),
581    );
582    let mut operations = crate::understand::operation_catalog(definitions, &cases)?;
583    let card = open_interactions
584        .iter()
585        .find(|interaction| interaction.blocking)
586        .filter(|_| answered.is_none());
587    let card_summary = card.map(crate::interactions::summarize);
588    let text = input.text.as_deref().unwrap_or_default();
589    let effort = crate::effort::resolve(
590        &shared.config,
591        input.effort.unwrap_or(shared.config.effort.default),
592    );
593    let understood = crate::understand::run(
594        shared.understander.as_ref(),
595        &crate::understand::Sources {
596            turn: input.turn_id,
597            definitions,
598            cases: &cases,
599            resolver: &resolver,
600            text,
601            locale: &input.locale,
602            today: now.date_naive(),
603            recent: &recent,
604            previous: previous.as_ref(),
605            card: card.zip(card_summary.as_ref()),
606            typed_answers_allowed: shared.policy.allow_text_resolution_for_low_risk,
607            config: &shared.config.understanding,
608            effort: &effort,
609            knowledge: shared.knowledge,
610        },
611        &operations,
612        &NoSteps,
613    )
614    .await?;
615    let mut understanding = understood.understanding;
616    if let Some(lookup) = lookup {
617        let found = lookup
618            .find(shared, input, &understanding, &mut cases)
619            .await?;
620        if found.values().any(|keys| !keys.is_empty()) {
621            resolver = build_resolver(
622                &input.actor.account_id,
623                input.turn_id,
624                definitions,
625                &shared.case_ids,
626                &cases,
627                origin,
628                blocking_summary(answered.as_ref(), &open_interactions),
629            );
630            operations = crate::understand::operation_catalog(definitions, &cases)?;
631            aim_at_found(&mut understanding, &found, &resolver);
632        }
633    }
634    // A typed answer to the card on screen is admitted like a click, on the channel
635    // that authorizes only what needs no confirmation.
636    if answered.is_none()
637        && let (Some(card), Some(typed)) = (card, understanding.card_answer.as_ref())
638    {
639        let response = typed_response(card, &typed.option);
640        answered = admit_typed(input, &cases, card, &response, now).ok();
641    }
642    let with_card = card_act_of(understanding, answered.as_ref(), &resolver);
643    let reduction_context = reduction_context(
644        &cases,
645        &open_interactions,
646        &resolver,
647        &operations,
648        &shared.policy,
649        shared.config.understanding.plan_limits,
650        now,
651    );
652    let reducer = turn_reducer(
653        definitions,
654        &cases,
655        &resolver,
656        &shared.policy_engine,
657        &shared.config,
658        answered.as_ref(),
659        with_card.card_act,
660    );
661    let stage = crate::signals::Stage::enter();
662    let reduced = reducer.reduce_turn(input, &with_card.understanding, &reduction_context)?;
663    stage.observe(
664        shared.observer.as_ref(),
665        Signal::ReductionDuration,
666        &SignalLabels::none(),
667    );
668    Ok(PlannedTurn {
669        views: cases.into_values().map(|case| case.view).collect(),
670        understanding: (!text.trim().is_empty() || !with_card.understanding.acts.is_empty())
671            .then(|| with_card.understanding.clone()),
672        tasks: understood.tasks,
673        target_resolutions: target_resolution_records(&reduced.plan),
674        policy_decisions: reduced.plan.policy_decisions.clone(),
675        would_persist: reduced.plan.pre_execution_interactions.clone(),
676        would_execute: command_refs(&reduced.plan),
677        would_claim: would_claim(&reduced.plan),
678        reduction: reduced.plan,
679    })
680}
681
682/// Collects everything a [`TurnPlanner`] needs.
683#[derive(Default)]
684pub struct TurnPlannerBuilder {
685    workflows: Option<WorkflowReadRegistry>,
686    stores: Option<ReadOnlyStores>,
687    directory: Option<Arc<dyn CaseDirectory>>,
688    shared: SharedBuilder,
689}
690
691impl fmt::Debug for TurnPlannerBuilder {
692    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
693        f.debug_struct("TurnPlannerBuilder")
694            .field("workflows", &self.workflows.is_some())
695            .field("stores", &self.stores.is_some())
696            .finish_non_exhaustive()
697    }
698}
699
700macro_rules! shared_setters {
701    () => {
702        /// Sets what understands each turn.
703        #[must_use]
704        pub fn understander(mut self, understander: Arc<dyn TurnUnderstander>) -> Self {
705            self.shared.understander = Some(understander);
706            self
707        }
708
709        /// Sets the policy snapshot.
710        #[must_use]
711        pub fn policy(mut self, policy: PolicySnapshot) -> Self {
712            self.shared.policy = policy;
713            self
714        }
715
716        /// Sets the configuration.
717        #[must_use]
718        pub fn config(mut self, config: OrchestratorConfig) -> Self {
719            self.shared.config = config;
720            self
721        }
722
723        /// Sets the clock.
724        #[must_use]
725        pub fn clock(mut self, clock: Arc<dyn TurnClock>) -> Self {
726            self.shared.clock = clock;
727            self
728        }
729
730        /// Sets the factory of new case identifiers.
731        #[must_use]
732        pub fn case_id_factory(mut self, factory: Arc<dyn CaseIdFactory>) -> Self {
733            self.shared.case_ids = factory;
734            self
735        }
736
737        /// Sets the observer.
738        #[must_use]
739        pub fn observer(mut self, observer: Arc<dyn Observer>) -> Self {
740            self.shared.observer = observer;
741            self
742        }
743    };
744}
745
746impl TurnPlannerBuilder {
747    /// An empty builder.
748    #[must_use]
749    pub fn new() -> Self {
750        Self::default()
751    }
752
753    /// Sets the read-only registry.
754    #[must_use]
755    pub fn workflows(mut self, workflows: WorkflowReadRegistry) -> Self {
756        self.workflows = Some(workflows);
757        self
758    }
759
760    /// Sets the read-only stores.
761    #[must_use]
762    pub fn stores(mut self, stores: ReadOnlyStores) -> Self {
763        self.stores = Some(stores);
764        self
765    }
766
767    /// Sets the case directory.
768    #[must_use]
769    pub fn case_directory(mut self, directory: Arc<dyn CaseDirectory>) -> Self {
770        self.directory = Some(directory);
771        self
772    }
773
774    shared_setters!();
775
776    /// Builds the planner.
777    ///
778    /// # Errors
779    ///
780    /// [`crate::orchestrator::BuildError`] for a missing part or an invalid configuration.
781    pub fn build(self) -> Result<TurnPlanner, crate::orchestrator::BuildError> {
782        use crate::orchestrator::BuildError;
783        self.shared.config.validate()?;
784        let workflows = self
785            .workflows
786            .ok_or(BuildError::Missing { part: "workflows" })?;
787        let stores = self.stores.ok_or(BuildError::Missing { part: "stores" })?;
788        let directory = self
789            .directory
790            .ok_or(BuildError::Missing { part: "directory" })?;
791        Ok(TurnPlanner::assemble(
792            workflows,
793            stores,
794            directory,
795            self.shared.build()?,
796        ))
797    }
798}
799
800/// Collects everything a [`SeededTurnPlanner`] needs.
801#[derive(Default)]
802pub struct SeededTurnPlannerBuilder {
803    definitions: Option<WorkflowDefinitions>,
804    shared: SharedBuilder,
805}
806
807impl fmt::Debug for SeededTurnPlannerBuilder {
808    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
809        f.debug_struct("SeededTurnPlannerBuilder")
810            .field("definitions", &self.definitions.is_some())
811            .finish_non_exhaustive()
812    }
813}
814
815impl SeededTurnPlannerBuilder {
816    /// An empty builder.
817    #[must_use]
818    pub fn new() -> Self {
819        Self::default()
820    }
821
822    /// Sets the workflow definitions.
823    #[must_use]
824    pub fn definitions(mut self, definitions: WorkflowDefinitions) -> Self {
825        self.definitions = Some(definitions);
826        self
827    }
828
829    shared_setters!();
830
831    /// Builds the planner.
832    ///
833    /// # Errors
834    ///
835    /// [`crate::orchestrator::BuildError`] for a missing part or an invalid configuration.
836    pub fn build(self) -> Result<SeededTurnPlanner, crate::orchestrator::BuildError> {
837        use crate::orchestrator::BuildError;
838        self.shared.config.validate()?;
839        let definitions = self
840            .definitions
841            .ok_or(BuildError::Missing { part: "workflows" })?;
842        Ok(SeededTurnPlanner {
843            definitions,
844            shared: self.shared.build()?,
845        })
846    }
847}
848
849struct SharedBuilder {
850    understander: Option<Arc<dyn TurnUnderstander>>,
851    policy: PolicySnapshot,
852    config: OrchestratorConfig,
853    clock: Arc<dyn TurnClock>,
854    case_ids: Arc<dyn CaseIdFactory>,
855    observer: Arc<dyn Observer>,
856}
857
858impl Default for SharedBuilder {
859    fn default() -> Self {
860        Self {
861            understander: None,
862            policy: PolicySnapshot::conservative(),
863            config: OrchestratorConfig::conservative(),
864            clock: Arc::new(SystemTurnClock),
865            case_ids: Arc::new(DerivedCaseIdFactory),
866            observer: Arc::new(NoopObserver),
867        }
868    }
869}
870
871impl SharedBuilder {
872    fn build(self) -> Result<SharedPlanning, crate::orchestrator::BuildError> {
873        use crate::orchestrator::BuildError;
874        let understander = self.understander.ok_or(BuildError::Missing {
875            part: "understander",
876        })?;
877        let policy_engine = PolicyEngine::new(&self.config);
878        Ok(SharedPlanning {
879            understander,
880            policy_engine,
881            policy: self.config.policy_snapshot(self.policy),
882            config: self.config,
883            clock: self.clock,
884            case_ids: self.case_ids,
885            observer: self.observer,
886            // A planner alone answers nothing: it frames questions as a turn with a source.
887            knowledge: true,
888        })
889    }
890}
891
892pub(crate) fn seeded_from_parts(
893    definitions: WorkflowDefinitions,
894    shared: SharedPlanning,
895) -> SeededTurnPlanner {
896    SeededTurnPlanner {
897        definitions,
898        shared,
899    }
900}