1mod act;
23mod cards;
24mod copy;
25mod notices;
26mod police;
27
28pub use self::copy::{CONDITION_HOLDS_ANSWER, NoticeCopy, notice, rejection};
29
30use std::collections::{BTreeMap, BTreeSet};
31use std::sync::Arc;
32
33use indexmap::IndexMap;
34use turnframe_core::case::{CaseKey, CaseRef};
35use turnframe_core::command::{
36 AtomicityScope, CommandBatch, CommandEnvelope, CommandOrigin, CommandPolicy,
37 ConfirmationPolicy, IdempotencyKey, RiskClass,
38};
39use turnframe_core::error::{DomainRejection, ErasedCallError, ReductionError, RejectionCode};
40use turnframe_core::flow::{
41 DomainEnumeration, ErasedWorkflow, ErasedWorkflowView, WorkflowDefinitions,
42};
43use turnframe_core::hash::{Digest, canonical_digest};
44use turnframe_core::ids::{BatchId, BlockId, CommandId, OptionId, QuestionId, TurnId};
45use turnframe_core::interaction::{
46 InteractionKind, InteractionOption, InteractionPayload, InteractionSpec, OptionStyle,
47 ReviewDiffEntry, StoredInteractionAction, TextResolutionPolicy,
48};
49use turnframe_core::locale::LocalizedText;
50use turnframe_core::operation::OperationSpec;
51use turnframe_core::plan::{AnswerBasis, TargetPolicy};
52use turnframe_core::policy::PolicyDecision;
53use turnframe_core::reduce::{
54 AnswerTask, Capability, CommandRef, PlannedAct, PlannedActResult, ReductionContext,
55 ReductionPlan, SourcePolicy, TextSpan, TurnReducer,
56};
57use turnframe_core::response::{AnswerProgress, NarratableFact, NoticeSeverity, ServerNotice};
58use turnframe_core::target::{ResolvedAct, ResolvedActKind, TargetCandidate, TargetResolution};
59use turnframe_core::turn::TurnInput;
60use turnframe_core::understanding::{
61 ActAction, ActId, ActStatus, ArgumentValue, ConstraintKind, NotUnderstoodReason, QuestionTopic,
62 RecordValue, Understanding, UnderstoodAct, WordRange,
63};
64
65use crate::config::{ExecutionConfig, NarrationConfig, OrchestratorConfig};
66use crate::policy::{BlockReason, PolicyEngine, PolicyOutcome, PolicyRequest};
67use crate::resolve::{TargetOutcome, TargetResolver};
68use crate::resume::DeferredAct;
69
70#[derive(Debug, Clone)]
72pub struct DefaultTurnReducer {
73 workflows: WorkflowDefinitions,
74 resolver: TargetResolver,
75 policy: PolicyEngine,
76 states: IndexMap<CaseKey, serde_json::Value>,
77 execution: ExecutionConfig,
78 narration: NarrationConfig,
79 copy: NoticeCopy,
80 confirmed_origin: Option<(CommandOrigin, ActId)>,
81}
82
83impl DefaultTurnReducer {
84 #[must_use]
88 pub fn new(
89 workflows: impl Into<WorkflowDefinitions>,
90 resolver: TargetResolver,
91 policy: PolicyEngine,
92 config: &OrchestratorConfig,
93 ) -> Self {
94 Self {
95 workflows: workflows.into(),
96 resolver,
97 policy,
98 states: IndexMap::new(),
99 execution: config.execution,
100 narration: config.narration,
101 copy: NoticeCopy::standard(),
102 confirmed_origin: None,
103 }
104 }
105
106 #[must_use]
109 pub fn with_state(mut self, case: CaseKey, state: serde_json::Value) -> Self {
110 self.states.insert(case, state);
111 self
112 }
113
114 #[must_use]
116 pub fn with_states(
117 mut self,
118 states: impl IntoIterator<Item = (CaseKey, serde_json::Value)>,
119 ) -> Self {
120 self.states.extend(states);
121 self
122 }
123
124 #[must_use]
127 pub fn with_confirmed_origin(mut self, origin: CommandOrigin, act: ActId) -> Self {
128 self.confirmed_origin = Some((origin, act));
129 self
130 }
131
132 #[must_use]
134 pub fn with_copy(mut self, copy: NoticeCopy) -> Self {
135 self.copy = copy;
136 self
137 }
138
139 #[must_use]
141 pub const fn resolver(&self) -> &TargetResolver {
142 &self.resolver
143 }
144}
145
146#[derive(Debug, Clone, PartialEq, Eq)]
149#[non_exhaustive]
150pub struct ReducedTurn {
151 pub plan: ReductionPlan,
153 pub pending: Vec<CommandBatch<serde_json::Value>>,
155}
156
157impl ReducedTurn {
158 #[must_use]
160 pub fn pending_envelopes(&self) -> Vec<&CommandEnvelope<serde_json::Value>> {
161 self.pending
162 .iter()
163 .flat_map(|batch| batch.envelopes.iter())
164 .collect()
165 }
166}
167
168impl DefaultTurnReducer {
169 pub fn reduce_turn(
175 &self,
176 input: &TurnInput,
177 understanding: &Understanding,
178 context: &ReductionContext,
179 ) -> Result<ReducedTurn, ReductionError> {
180 let (plan, pending) = Session::new(self, input, context).run(understanding)?;
181 Ok(ReducedTurn { plan, pending })
182 }
183}
184
185impl TurnReducer for DefaultTurnReducer {
186 fn reduce(
187 &self,
188 input: &TurnInput,
189 understanding: &Understanding,
190 context: &ReductionContext,
191 ) -> Result<ReductionPlan, ReductionError> {
192 Session::new(self, input, context)
193 .run(understanding)
194 .map(|(plan, _)| plan)
195 }
196}
197
198struct Session<'a> {
200 reducer: &'a DefaultTurnReducer,
201 input: &'a TurnInput,
202 context: &'a ReductionContext,
203 acts: Vec<UnderstoodAct>,
204 constraints: Vec<ConstraintKind>,
205 evidence_digests: Vec<Digest>,
206 outcomes: Vec<Option<TargetOutcome>>,
207 results: Vec<Option<PlannedActResult>>,
208 batches: IndexMap<(CaseKey, String), CommandBatch<serde_json::Value>>,
209 pending: IndexMap<(CaseKey, String), CommandBatch<serde_json::Value>>,
210 decisions: Vec<PolicyDecision>,
211 specs: Vec<InteractionSpec>,
212 notices: IndexMap<&'static str, ServerNotice>,
213 refusals: Vec<NarratableFact>,
214 changed_nothing: Vec<NarratableFact>,
215 awaiting_confirmation: Vec<NarratableFact>,
216 minted: Vec<CaseRef>,
218 opened_by: BTreeMap<ActId, CaseRef>,
220 unborn: BTreeMap<CaseKey, serde_json::Value>,
222 applied: Vec<ConstraintKind>,
223 command_count: usize,
224}
225
226impl<'a> Session<'a> {
227 fn new(
228 reducer: &'a DefaultTurnReducer,
229 input: &'a TurnInput,
230 context: &'a ReductionContext,
231 ) -> Self {
232 Self {
233 reducer,
234 input,
235 context,
236 acts: Vec::new(),
237 constraints: Vec::new(),
238 evidence_digests: Vec::new(),
239 outcomes: Vec::new(),
240 results: Vec::new(),
241 batches: IndexMap::new(),
242 pending: IndexMap::new(),
243 decisions: Vec::new(),
244 specs: Vec::new(),
245 notices: IndexMap::new(),
246 refusals: Vec::new(),
247 changed_nothing: Vec::new(),
248 awaiting_confirmation: Vec::new(),
249 minted: Vec::new(),
250 opened_by: BTreeMap::new(),
251 unborn: BTreeMap::new(),
252 applied: Vec::new(),
253 command_count: 0,
254 }
255 }
256
257 fn turn_id(&self) -> TurnId {
258 self.input.turn_id
259 }
260
261 fn words(&self, range: WordRange) -> &str {
263 self.input
264 .text
265 .as_deref()
266 .and_then(|text| text.get(range.start..range.end))
267 .unwrap_or_default()
268 }
269
270 fn run(
271 mut self,
272 understanding: &Understanding,
273 ) -> Result<(ReductionPlan, Vec<CommandBatch<serde_json::Value>>), ReductionError> {
274 self.context.limits.enforce(understanding)?;
275 self.acts = understanding.acts.clone();
276 self.constraints = understanding.constraints.iter().map(|c| c.kind).collect();
277 self.results = vec![None; self.acts.len()];
278 self.outcomes = vec![None; self.acts.len()];
279 self.evidence_digests = self
280 .acts
281 .iter()
282 .map(|act| {
283 canonical_digest(&(&act.words, &act.arguments, &act.action))
284 .map_err(|_| ReductionError::Hash)
285 })
286 .collect::<Result<_, _>>()?;
287
288 for index in 0..self.acts.len() {
289 self.plan_act(index, understanding)?;
290 }
291 self.attach_dependents();
292 if self.command_count > self.reducer.execution.max_commands_per_turn {
293 return Err(ReductionError::CommandBudgetExceeded {
294 limit: self.reducer.execution.max_commands_per_turn,
295 actual: self.command_count,
296 });
297 }
298 self.add_partial_result_notice();
299 self.note_not_understood(understanding);
300
301 let answer_tasks = self.answer_tasks(understanding);
302 let ids: Vec<ActId> = self.acts.iter().map(|act| act.id).collect();
303 let acts = self.planned_acts()?;
304 let plan = ReductionPlan {
305 turn_id: self.turn_id(),
306 acts,
307 batches: self.batches.into_values().collect(),
308 policy_decisions: self.decisions,
309 answer_tasks,
310 superseded_operations: understanding
311 .superseded
312 .iter()
313 .filter_map(|superseded| match &superseded.action {
314 ActAction::Apply { operation } => Some(operation.clone()),
315 ActAction::Start { .. } => None,
316 })
317 .collect(),
318 refusals: self.refusals,
319 changed_nothing: self.changed_nothing,
320 awaiting_confirmation: self.awaiting_confirmation,
321 constraints_applied: self.applied,
322 notices: self.notices.into_values().collect(),
323 pre_execution_interactions: self.specs,
324 plan_hash: Digest::of_bytes(b""),
325 }
326 .with_hash()
327 .map_err(|_| ReductionError::Hash)?;
328 plan.validate(&ids)?;
329 Ok((plan, self.pending.into_values().collect()))
330 }
331
332 fn planned_acts(&mut self) -> Result<Vec<PlannedAct>, ReductionError> {
333 let mut planned = Vec::with_capacity(self.acts.len());
334 for (index, act) in self.acts.iter().enumerate() {
335 let result =
336 self.results[index]
337 .clone()
338 .ok_or_else(|| ReductionError::InconsistentPlan {
339 detail: format!("act {} left without a result", act.id),
340 })?;
341 planned.push(PlannedAct {
342 act: act.clone(),
343 target: self.outcomes[index]
344 .as_ref()
345 .and_then(TargetOutcome::resolution),
346 result,
347 });
348 }
349 Ok(planned)
350 }
351
352 fn spec(&self, act: &UnderstoodAct) -> Option<&'a OperationSpec> {
353 match &act.action {
354 ActAction::Apply { operation } => self.context.operations.get(operation),
355 ActAction::Start { .. } => None,
356 }
357 }
358
359 fn operation_name(act: &UnderstoodAct) -> String {
360 match &act.action {
361 ActAction::Apply { operation } => operation.as_str().to_owned(),
362 ActAction::Start { workflow } => format!("{workflow}.start"),
363 }
364 }
365}