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}
188
189impl TurnPlanner {
190    /// Starts a builder.
191    #[must_use]
192    pub fn builder() -> TurnPlannerBuilder {
193        TurnPlannerBuilder::default()
194    }
195
196    pub(crate) fn assemble(
197        workflows: WorkflowReadRegistry,
198        stores: ReadOnlyStores,
199        directory: Arc<dyn CaseDirectory>,
200        shared: SharedPlanning,
201    ) -> Self {
202        Self {
203            workflows,
204            stores,
205            directory,
206            shared,
207        }
208    }
209
210    /// The seeded planner over the same definitions and parts.
211    #[must_use]
212    pub fn seeded(&self) -> SeededTurnPlanner {
213        SeededTurnPlanner {
214            definitions: self.workflows.definitions().clone(),
215            shared: self.shared.clone(),
216        }
217    }
218
219    /// Plans `input` against the stores, stopping before the first effect.
220    ///
221    /// # Errors
222    ///
223    /// Whatever the same steps of a real turn would fail with.
224    pub async fn plan(&self, input: TurnInput) -> Result<PlannedTurn, OrchestratorError> {
225        let now = self.shared.clock.now();
226        input.validate_shape_within(&self.shared.config.understanding.turn_limits)?;
227        let account = input.actor.account_id.clone();
228        self.stores
229            .conversations()
230            .load_conversation(&account, &input.conversation_id)
231            .await
232            .map_err(|error| match error {
233                StoreError::NotFound => {
234                    OrchestratorError::Unauthorized(AuthorizationError::ConversationNotAccessible {
235                        conversation_id: input.conversation_id,
236                    })
237                }
238                other => OrchestratorError::Store(other),
239            })?;
240        let open_interactions = self
241            .stores
242            .interactions()
243            .list_open_for_conversation(&account, &input.conversation_id)
244            .await
245            .map_err(OrchestratorError::Store)?;
246        let (cases, origin_case) = self.load_cases(&input, &open_interactions).await?;
247        let open_interactions = addressable_cards(&cases, open_interactions);
248        let answered = match input.interaction_response.as_ref() {
249            None => None,
250            Some(response) => {
251                let record = self
252                    .stores
253                    .interactions()
254                    .get(&account, &response.interaction_id)
255                    .await
256                    .map_err(|error| match error {
257                        StoreError::NotFound => OrchestratorError::Interaction(
258                            InteractionError::Rejected(InteractionRejection::NotFound),
259                        ),
260                        _ => OrchestratorError::Interaction(InteractionError::NotPersisted),
261                    })?;
262                admit(&input, &cases, record, now)?
263            }
264        };
265        let turns = self
266            .stores
267            .conversations()
268            .load_recent_turns(
269                &account,
270                &input.conversation_id,
271                self.shared
272                    .config
273                    .understanding
274                    .transcript_turns
275                    .unwrap_or(usize::MAX),
276            )
277            .await
278            .unwrap_or_default();
279        let (recent, previous) = recent_messages(turns, input.turn_id);
280        plan_with(
281            &self.shared,
282            self.workflows.definitions(),
283            &input,
284            PlanningContext {
285                lookup: Some(Lookup {
286                    directory: self.directory.as_ref(),
287                    workflows: &self.workflows,
288                }),
289                now,
290                cases,
291                open_interactions,
292                origin_case,
293                answered,
294                recent,
295                previous,
296            },
297        )
298        .await
299    }
300
301    async fn load_cases(
302        &self,
303        input: &TurnInput,
304        open_interactions: &[Interaction],
305    ) -> Result<(IndexMap<CaseKey, LoadedCase>, Option<CaseKey>), OrchestratorError> {
306        let account = &input.actor.account_id;
307        let mut candidates = self
308            .directory
309            .candidates(&input.actor, &input.conversation_id)
310            .await
311            .map_err(OrchestratorError::Store)?;
312        let mut origin_case = None;
313        if let Some(origin) = input.origin.as_ref()
314            && let Some(candidate) = self
315                .directory
316                .resolve_origin(&input.actor, &origin.origin_token)
317                .await
318                .map_err(OrchestratorError::Store)?
319        {
320            origin_case = Some(candidate.key.clone());
321            candidates.push(candidate);
322        }
323        let from_cards = candidates_of_open_cards(&candidates, open_interactions);
324        let mut cases = IndexMap::new();
325        for (candidate, from_card) in candidates
326            .into_iter()
327            .map(|candidate| (candidate, false))
328            .chain(from_cards.into_iter().map(|candidate| (candidate, true)))
329        {
330            if cases.contains_key(&candidate.key) {
331                continue;
332            }
333            let Ok(definition) = self
334                .workflows
335                .definitions()
336                .require(&candidate.key.workflow)
337                .cloned()
338            else {
339                continue;
340            };
341            let loader = self.workflows.require_loader(&candidate.key.workflow)?;
342            let loaded = loader
343                .load_case(account, &candidate.key.case_id)
344                .await
345                .map_err(OrchestratorError::Store)?;
346            let candidate = if from_card {
347                let Some(authorized) = admit_card_case(
348                    self.shared.observer.as_ref(),
349                    self.directory.as_ref(),
350                    &input.actor,
351                    &input.conversation_id,
352                    candidate,
353                    loaded.value.is_some(),
354                )
355                .await
356                .map_err(OrchestratorError::Store)?
357                else {
358                    continue;
359                };
360                authorized
361            } else {
362                candidate
363            };
364            let case_ref = candidate.key.clone().at(loaded.revision);
365            cases.insert(
366                candidate.key,
367                project_case(
368                    self.shared.observer.as_ref(),
369                    &definition,
370                    case_ref,
371                    candidate.label,
372                    loaded.value,
373                    DirectoryTerms {
374                        confirm_every_write: candidate.confirm_every_write,
375                        subject_only_when_named: candidate.subject_only_when_named,
376                    },
377                )?,
378            );
379        }
380        Ok((cases, origin_case))
381    }
382}
383
384impl SeededTurnPlanner {
385    /// Starts a builder.
386    #[must_use]
387    pub fn builder() -> SeededTurnPlannerBuilder {
388        SeededTurnPlannerBuilder::default()
389    }
390
391    /// Plans `input` against `cases`, stopping before the first effect.
392    ///
393    /// # Errors
394    ///
395    /// Whatever the same steps of a real turn would fail with.
396    pub async fn plan(
397        &self,
398        input: TurnInput,
399        cases: Vec<SeededCase>,
400    ) -> Result<PlannedTurn, OrchestratorError> {
401        let now = self.shared.clock.now();
402        input.validate_shape_within(&self.shared.config.understanding.turn_limits)?;
403        let mut loaded = IndexMap::new();
404        let mut open_interactions = Vec::new();
405        for seeded in cases {
406            let key = seeded.case_ref.key();
407            let definition = self.definitions.require(&key.workflow)?.clone();
408            open_interactions.extend(seeded.open_interactions);
409            let label = seeded
410                .label
411                .unwrap_or_else(|| seeded.case_ref.case_id.to_string());
412            loaded.insert(
413                key,
414                project_case(
415                    self.shared.observer.as_ref(),
416                    &definition,
417                    seeded.case_ref,
418                    label,
419                    seeded.state,
420                    DirectoryTerms {
421                        confirm_every_write: false,
422                        subject_only_when_named: seeded.subject_only_when_named,
423                    },
424                )?,
425            );
426        }
427        let open_interactions = addressable_cards(&loaded, open_interactions);
428        let answered = match input.interaction_response.as_ref() {
429            None => None,
430            Some(response) => {
431                let Some(interaction) = open_interactions
432                    .iter()
433                    .find(|open| open.id == response.interaction_id)
434                    .cloned()
435                else {
436                    return Err(OrchestratorError::Interaction(InteractionError::Rejected(
437                        InteractionRejection::NotFound,
438                    )));
439                };
440                admit(&input, &loaded, InteractionRecord::new(interaction), now)?
441            }
442        };
443        plan_with(
444            &self.shared,
445            &self.definitions,
446            &input,
447            PlanningContext {
448                lookup: None,
449                now,
450                cases: loaded,
451                open_interactions,
452                origin_case: None,
453                answered,
454                recent: Vec::new(),
455                previous: None,
456            },
457        )
458        .await
459    }
460}
461
462struct PlanningContext<'a> {
463    lookup: Option<Lookup<'a>>,
464    now: DateTime<Utc>,
465    cases: IndexMap<CaseKey, LoadedCase>,
466    open_interactions: Vec<Interaction>,
467    origin_case: Option<CaseKey>,
468    answered: Option<AcceptedInteraction>,
469    recent: Vec<RecentMessage>,
470    previous: Option<AssistantTurn>,
471}
472
473/// Where a plan looks up a record the message named and the turn did not have.
474struct Lookup<'a> {
475    directory: &'a dyn CaseDirectory,
476    workflows: &'a WorkflowReadRegistry,
477}
478
479impl Lookup<'_> {
480    /// Loads what the directory finds for each unlisted act, and returns it per act.
481    async fn find(
482        &self,
483        shared: &SharedPlanning,
484        input: &TurnInput,
485        understanding: &Understanding,
486        cases: &mut IndexMap<CaseKey, LoadedCase>,
487    ) -> Result<std::collections::BTreeMap<ActId, Vec<CaseKey>>, OrchestratorError> {
488        let text = input.text.as_deref().unwrap_or_default();
489        let limit = shared.config.interaction.max_selection_candidates;
490        let mut found = std::collections::BTreeMap::new();
491        for (act, workflow, named) in unlisted(understanding, text) {
492            let candidates = self
493                .directory
494                .find(
495                    &input.actor,
496                    &input.conversation_id,
497                    &workflow,
498                    named.as_deref(),
499                )
500                .await
501                .map_err(OrchestratorError::Store)?;
502            let mut keys = Vec::new();
503            for candidate in candidates.into_iter().take(limit) {
504                keys.push(candidate.key.clone());
505                if let Some(case) = load_listed(shared, self.workflows, input, candidate).await? {
506                    cases.insert(case.case_ref.key(), case);
507                }
508            }
509            found.insert(act, keys);
510        }
511        Ok(found)
512    }
513}
514
515/// One candidate the directory listed, loaded and projected; `None` for a workflow
516/// this runtime does not host.
517async fn load_listed(
518    shared: &SharedPlanning,
519    workflows: &WorkflowReadRegistry,
520    input: &TurnInput,
521    candidate: crate::orchestrator::CaseCandidate,
522) -> Result<Option<LoadedCase>, OrchestratorError> {
523    let Ok(definition) = workflows
524        .definitions()
525        .require(&candidate.key.workflow)
526        .cloned()
527    else {
528        return Ok(None);
529    };
530    let loaded = workflows
531        .require_loader(&candidate.key.workflow)?
532        .load_case(&input.actor.account_id, &candidate.key.case_id)
533        .await
534        .map_err(OrchestratorError::Store)?;
535    let case_ref = candidate.key.clone().at(loaded.revision);
536    project_case(
537        shared.observer.as_ref(),
538        &definition,
539        case_ref,
540        candidate.label,
541        loaded.value,
542        DirectoryTerms {
543            confirm_every_write: candidate.confirm_every_write,
544            subject_only_when_named: candidate.subject_only_when_named,
545        },
546    )
547    .map(Some)
548}
549
550async fn plan_with(
551    shared: &SharedPlanning,
552    definitions: &WorkflowDefinitions,
553    input: &TurnInput,
554    context: PlanningContext<'_>,
555) -> Result<PlannedTurn, OrchestratorError> {
556    let PlanningContext {
557        lookup,
558        now,
559        mut cases,
560        open_interactions,
561        origin_case,
562        mut answered,
563        recent,
564        previous,
565    } = context;
566    let origin = input
567        .origin
568        .as_ref()
569        .map(|origin| &origin.origin_token)
570        .zip(origin_case.as_ref());
571    let mut resolver = build_resolver(
572        &input.actor.account_id,
573        input.turn_id,
574        definitions,
575        &shared.case_ids,
576        &cases,
577        origin,
578        blocking_summary(answered.as_ref(), &open_interactions),
579    );
580    let mut operations = crate::understand::operation_catalog(definitions, &cases)?;
581    let card = open_interactions
582        .iter()
583        .find(|interaction| interaction.blocking)
584        .filter(|_| answered.is_none());
585    let card_summary = card.map(crate::interactions::summarize);
586    let text = input.text.as_deref().unwrap_or_default();
587    let effort = crate::effort::resolve(
588        &shared.config,
589        input.effort.unwrap_or(shared.config.effort.default),
590    );
591    let understood = crate::understand::run(
592        shared.understander.as_ref(),
593        &crate::understand::Sources {
594            turn: input.turn_id,
595            definitions,
596            cases: &cases,
597            resolver: &resolver,
598            text,
599            locale: &input.locale,
600            today: now.date_naive(),
601            recent: &recent,
602            previous: previous.as_ref(),
603            card: card.zip(card_summary.as_ref()),
604            typed_answers_allowed: shared.policy.allow_text_resolution_for_low_risk,
605            config: &shared.config.understanding,
606            effort: &effort,
607        },
608        &operations,
609        &NoSteps,
610    )
611    .await?;
612    let mut understanding = understood.understanding;
613    if let Some(lookup) = lookup {
614        let found = lookup
615            .find(shared, input, &understanding, &mut cases)
616            .await?;
617        if found.values().any(|keys| !keys.is_empty()) {
618            resolver = build_resolver(
619                &input.actor.account_id,
620                input.turn_id,
621                definitions,
622                &shared.case_ids,
623                &cases,
624                origin,
625                blocking_summary(answered.as_ref(), &open_interactions),
626            );
627            operations = crate::understand::operation_catalog(definitions, &cases)?;
628            aim_at_found(&mut understanding, &found, &resolver);
629        }
630    }
631    // A typed answer to the card on screen is admitted like a click, on the channel
632    // that authorizes only what needs no confirmation.
633    if answered.is_none()
634        && let (Some(card), Some(typed)) = (card, understanding.card_answer.as_ref())
635    {
636        let response = typed_response(card, &typed.option);
637        answered = admit_typed(input, &cases, card, &response, now).ok();
638    }
639    let with_card = card_act_of(understanding, answered.as_ref(), &resolver);
640    let reduction_context = reduction_context(
641        &cases,
642        &open_interactions,
643        &resolver,
644        &operations,
645        &shared.policy,
646        shared.config.understanding.plan_limits,
647        now,
648    );
649    let reducer = turn_reducer(
650        definitions,
651        &cases,
652        &resolver,
653        &shared.policy_engine,
654        &shared.config,
655        answered.as_ref(),
656        with_card.card_act,
657    );
658    let stage = crate::signals::Stage::enter();
659    let reduced = reducer.reduce_turn(input, &with_card.understanding, &reduction_context)?;
660    stage.observe(
661        shared.observer.as_ref(),
662        Signal::ReductionDuration,
663        &SignalLabels::none(),
664    );
665    Ok(PlannedTurn {
666        views: cases.into_values().map(|case| case.view).collect(),
667        understanding: (!text.trim().is_empty() || !with_card.understanding.acts.is_empty())
668            .then(|| with_card.understanding.clone()),
669        tasks: understood.tasks,
670        target_resolutions: target_resolution_records(&reduced.plan),
671        policy_decisions: reduced.plan.policy_decisions.clone(),
672        would_persist: reduced.plan.pre_execution_interactions.clone(),
673        would_execute: command_refs(&reduced.plan),
674        would_claim: would_claim(&reduced.plan),
675        reduction: reduced.plan,
676    })
677}
678
679/// Collects everything a [`TurnPlanner`] needs.
680#[derive(Default)]
681pub struct TurnPlannerBuilder {
682    workflows: Option<WorkflowReadRegistry>,
683    stores: Option<ReadOnlyStores>,
684    directory: Option<Arc<dyn CaseDirectory>>,
685    shared: SharedBuilder,
686}
687
688impl fmt::Debug for TurnPlannerBuilder {
689    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
690        f.debug_struct("TurnPlannerBuilder")
691            .field("workflows", &self.workflows.is_some())
692            .field("stores", &self.stores.is_some())
693            .finish_non_exhaustive()
694    }
695}
696
697macro_rules! shared_setters {
698    () => {
699        /// Sets what understands each turn.
700        #[must_use]
701        pub fn understander(mut self, understander: Arc<dyn TurnUnderstander>) -> Self {
702            self.shared.understander = Some(understander);
703            self
704        }
705
706        /// Sets the policy snapshot.
707        #[must_use]
708        pub fn policy(mut self, policy: PolicySnapshot) -> Self {
709            self.shared.policy = policy;
710            self
711        }
712
713        /// Sets the configuration.
714        #[must_use]
715        pub fn config(mut self, config: OrchestratorConfig) -> Self {
716            self.shared.config = config;
717            self
718        }
719
720        /// Sets the clock.
721        #[must_use]
722        pub fn clock(mut self, clock: Arc<dyn TurnClock>) -> Self {
723            self.shared.clock = clock;
724            self
725        }
726
727        /// Sets the factory of new case identifiers.
728        #[must_use]
729        pub fn case_id_factory(mut self, factory: Arc<dyn CaseIdFactory>) -> Self {
730            self.shared.case_ids = factory;
731            self
732        }
733
734        /// Sets the observer.
735        #[must_use]
736        pub fn observer(mut self, observer: Arc<dyn Observer>) -> Self {
737            self.shared.observer = observer;
738            self
739        }
740    };
741}
742
743impl TurnPlannerBuilder {
744    /// An empty builder.
745    #[must_use]
746    pub fn new() -> Self {
747        Self::default()
748    }
749
750    /// Sets the read-only registry.
751    #[must_use]
752    pub fn workflows(mut self, workflows: WorkflowReadRegistry) -> Self {
753        self.workflows = Some(workflows);
754        self
755    }
756
757    /// Sets the read-only stores.
758    #[must_use]
759    pub fn stores(mut self, stores: ReadOnlyStores) -> Self {
760        self.stores = Some(stores);
761        self
762    }
763
764    /// Sets the case directory.
765    #[must_use]
766    pub fn case_directory(mut self, directory: Arc<dyn CaseDirectory>) -> Self {
767        self.directory = Some(directory);
768        self
769    }
770
771    shared_setters!();
772
773    /// Builds the planner.
774    ///
775    /// # Errors
776    ///
777    /// [`crate::orchestrator::BuildError`] for a missing part or an invalid configuration.
778    pub fn build(self) -> Result<TurnPlanner, crate::orchestrator::BuildError> {
779        use crate::orchestrator::BuildError;
780        self.shared.config.validate()?;
781        let workflows = self
782            .workflows
783            .ok_or(BuildError::Missing { part: "workflows" })?;
784        let stores = self.stores.ok_or(BuildError::Missing { part: "stores" })?;
785        let directory = self
786            .directory
787            .ok_or(BuildError::Missing { part: "directory" })?;
788        Ok(TurnPlanner::assemble(
789            workflows,
790            stores,
791            directory,
792            self.shared.build()?,
793        ))
794    }
795}
796
797/// Collects everything a [`SeededTurnPlanner`] needs.
798#[derive(Default)]
799pub struct SeededTurnPlannerBuilder {
800    definitions: Option<WorkflowDefinitions>,
801    shared: SharedBuilder,
802}
803
804impl fmt::Debug for SeededTurnPlannerBuilder {
805    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
806        f.debug_struct("SeededTurnPlannerBuilder")
807            .field("definitions", &self.definitions.is_some())
808            .finish_non_exhaustive()
809    }
810}
811
812impl SeededTurnPlannerBuilder {
813    /// An empty builder.
814    #[must_use]
815    pub fn new() -> Self {
816        Self::default()
817    }
818
819    /// Sets the workflow definitions.
820    #[must_use]
821    pub fn definitions(mut self, definitions: WorkflowDefinitions) -> Self {
822        self.definitions = Some(definitions);
823        self
824    }
825
826    shared_setters!();
827
828    /// Builds the planner.
829    ///
830    /// # Errors
831    ///
832    /// [`crate::orchestrator::BuildError`] for a missing part or an invalid configuration.
833    pub fn build(self) -> Result<SeededTurnPlanner, crate::orchestrator::BuildError> {
834        use crate::orchestrator::BuildError;
835        self.shared.config.validate()?;
836        let definitions = self
837            .definitions
838            .ok_or(BuildError::Missing { part: "workflows" })?;
839        Ok(SeededTurnPlanner {
840            definitions,
841            shared: self.shared.build()?,
842        })
843    }
844}
845
846struct SharedBuilder {
847    understander: Option<Arc<dyn TurnUnderstander>>,
848    policy: PolicySnapshot,
849    config: OrchestratorConfig,
850    clock: Arc<dyn TurnClock>,
851    case_ids: Arc<dyn CaseIdFactory>,
852    observer: Arc<dyn Observer>,
853}
854
855impl Default for SharedBuilder {
856    fn default() -> Self {
857        Self {
858            understander: None,
859            policy: PolicySnapshot::conservative(),
860            config: OrchestratorConfig::conservative(),
861            clock: Arc::new(SystemTurnClock),
862            case_ids: Arc::new(DerivedCaseIdFactory),
863            observer: Arc::new(NoopObserver),
864        }
865    }
866}
867
868impl SharedBuilder {
869    fn build(self) -> Result<SharedPlanning, crate::orchestrator::BuildError> {
870        use crate::orchestrator::BuildError;
871        let understander = self.understander.ok_or(BuildError::Missing {
872            part: "understander",
873        })?;
874        let policy_engine = PolicyEngine::new(&self.config);
875        Ok(SharedPlanning {
876            understander,
877            policy_engine,
878            policy: self.config.policy_snapshot(self.policy),
879            config: self.config,
880            clock: self.clock,
881            case_ids: self.case_ids,
882            observer: self.observer,
883        })
884    }
885}
886
887pub(crate) fn seeded_from_parts(
888    definitions: WorkflowDefinitions,
889    shared: SharedPlanning,
890) -> SeededTurnPlanner {
891    SeededTurnPlanner {
892        definitions,
893        shared,
894    }
895}