1use 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#[derive(Debug, Clone, PartialEq, Eq)]
53#[non_exhaustive]
54pub struct SeededCase {
55 pub case_ref: CaseRef,
57 pub state: Option<serde_json::Value>,
59 pub label: Option<String>,
61 pub open_interactions: Vec<Interaction>,
63 pub subject_only_when_named: bool,
65}
66
67impl SeededCase {
68 #[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 #[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 #[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 #[must_use]
101 pub const fn reachable_only(mut self) -> Self {
102 self.subject_only_when_named = true;
103 self
104 }
105
106 #[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#[derive(Debug, Clone)]
116#[non_exhaustive]
117pub struct PlannedTurn {
118 pub views: Vec<ErasedWorkflowView>,
120 pub understanding: Option<Understanding>,
122 pub tasks: Vec<TaskRecord>,
124 pub target_resolutions: Vec<TargetResolutionRecord>,
126 pub reduction: ReductionPlan,
128 pub policy_decisions: Vec<PolicyDecision>,
130 pub would_persist: Vec<InteractionSpec>,
132 pub would_execute: Vec<CommandRef>,
134 pub would_claim: Vec<ClaimClass>,
136}
137
138impl PlannedTurn {
139 #[must_use]
141 pub fn summary(&self) -> TurnSummary {
142 TurnSummary::from_planned(self)
143 }
144}
145
146pub 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
163pub 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#[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 pub(crate) knowledge: bool,
189}
190
191impl TurnPlanner {
192 #[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 #[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 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 #[must_use]
389 pub fn builder() -> SeededTurnPlannerBuilder {
390 SeededTurnPlannerBuilder::default()
391 }
392
393 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
475struct Lookup<'a> {
477 directory: &'a dyn CaseDirectory,
478 workflows: &'a WorkflowReadRegistry,
479}
480
481impl Lookup<'_> {
482 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
517async 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 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#[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 #[must_use]
704 pub fn understander(mut self, understander: Arc<dyn TurnUnderstander>) -> Self {
705 self.shared.understander = Some(understander);
706 self
707 }
708
709 #[must_use]
711 pub fn policy(mut self, policy: PolicySnapshot) -> Self {
712 self.shared.policy = policy;
713 self
714 }
715
716 #[must_use]
718 pub fn config(mut self, config: OrchestratorConfig) -> Self {
719 self.shared.config = config;
720 self
721 }
722
723 #[must_use]
725 pub fn clock(mut self, clock: Arc<dyn TurnClock>) -> Self {
726 self.shared.clock = clock;
727 self
728 }
729
730 #[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 #[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 #[must_use]
749 pub fn new() -> Self {
750 Self::default()
751 }
752
753 #[must_use]
755 pub fn workflows(mut self, workflows: WorkflowReadRegistry) -> Self {
756 self.workflows = Some(workflows);
757 self
758 }
759
760 #[must_use]
762 pub fn stores(mut self, stores: ReadOnlyStores) -> Self {
763 self.stores = Some(stores);
764 self
765 }
766
767 #[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 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#[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 #[must_use]
818 pub fn new() -> Self {
819 Self::default()
820 }
821
822 #[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 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 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}