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 applied: Vec<ConstraintKind>,
221 command_count: usize,
222}
223
224impl<'a> Session<'a> {
225 fn new(
226 reducer: &'a DefaultTurnReducer,
227 input: &'a TurnInput,
228 context: &'a ReductionContext,
229 ) -> Self {
230 Self {
231 reducer,
232 input,
233 context,
234 acts: Vec::new(),
235 constraints: Vec::new(),
236 evidence_digests: Vec::new(),
237 outcomes: Vec::new(),
238 results: Vec::new(),
239 batches: IndexMap::new(),
240 pending: IndexMap::new(),
241 decisions: Vec::new(),
242 specs: Vec::new(),
243 notices: IndexMap::new(),
244 refusals: Vec::new(),
245 changed_nothing: Vec::new(),
246 awaiting_confirmation: Vec::new(),
247 minted: Vec::new(),
248 opened_by: BTreeMap::new(),
249 applied: Vec::new(),
250 command_count: 0,
251 }
252 }
253
254 fn turn_id(&self) -> TurnId {
255 self.input.turn_id
256 }
257
258 fn words(&self, range: WordRange) -> &str {
260 self.input
261 .text
262 .as_deref()
263 .and_then(|text| text.get(range.start..range.end))
264 .unwrap_or_default()
265 }
266
267 fn run(
268 mut self,
269 understanding: &Understanding,
270 ) -> Result<(ReductionPlan, Vec<CommandBatch<serde_json::Value>>), ReductionError> {
271 self.context.limits.enforce(understanding)?;
272 self.acts = understanding.acts.clone();
273 self.constraints = understanding.constraints.iter().map(|c| c.kind).collect();
274 self.results = vec![None; self.acts.len()];
275 self.outcomes = vec![None; self.acts.len()];
276 self.evidence_digests = self
277 .acts
278 .iter()
279 .map(|act| {
280 canonical_digest(&(&act.words, &act.arguments, &act.action))
281 .map_err(|_| ReductionError::Hash)
282 })
283 .collect::<Result<_, _>>()?;
284
285 for index in 0..self.acts.len() {
286 self.plan_act(index, understanding)?;
287 }
288 self.attach_dependents();
289 if self.command_count > self.reducer.execution.max_commands_per_turn {
290 return Err(ReductionError::CommandBudgetExceeded {
291 limit: self.reducer.execution.max_commands_per_turn,
292 actual: self.command_count,
293 });
294 }
295 self.add_partial_result_notice();
296 self.note_not_understood(understanding);
297
298 let answer_tasks = self.answer_tasks(understanding);
299 let ids: Vec<ActId> = self.acts.iter().map(|act| act.id).collect();
300 let acts = self.planned_acts()?;
301 let plan = ReductionPlan {
302 turn_id: self.turn_id(),
303 acts,
304 batches: self.batches.into_values().collect(),
305 policy_decisions: self.decisions,
306 answer_tasks,
307 superseded_operations: understanding
308 .superseded
309 .iter()
310 .filter_map(|superseded| match &superseded.action {
311 ActAction::Apply { operation } => Some(operation.clone()),
312 ActAction::Start { .. } => None,
313 })
314 .collect(),
315 refusals: self.refusals,
316 changed_nothing: self.changed_nothing,
317 awaiting_confirmation: self.awaiting_confirmation,
318 constraints_applied: self.applied,
319 notices: self.notices.into_values().collect(),
320 pre_execution_interactions: self.specs,
321 plan_hash: Digest::of_bytes(b""),
322 }
323 .with_hash()
324 .map_err(|_| ReductionError::Hash)?;
325 plan.validate(&ids)?;
326 Ok((plan, self.pending.into_values().collect()))
327 }
328
329 fn planned_acts(&mut self) -> Result<Vec<PlannedAct>, ReductionError> {
330 let mut planned = Vec::with_capacity(self.acts.len());
331 for (index, act) in self.acts.iter().enumerate() {
332 let result =
333 self.results[index]
334 .clone()
335 .ok_or_else(|| ReductionError::InconsistentPlan {
336 detail: format!("act {} left without a result", act.id),
337 })?;
338 planned.push(PlannedAct {
339 act: act.clone(),
340 target: self.outcomes[index]
341 .as_ref()
342 .and_then(TargetOutcome::resolution),
343 result,
344 });
345 }
346 Ok(planned)
347 }
348
349 fn spec(&self, act: &UnderstoodAct) -> Option<&'a OperationSpec> {
350 match &act.action {
351 ActAction::Apply { operation } => self.context.operations.get(operation),
352 ActAction::Start { .. } => None,
353 }
354 }
355
356 fn operation_name(act: &UnderstoodAct) -> String {
357 match &act.action {
358 ActAction::Apply { operation } => operation.as_str().to_owned(),
359 ActAction::Start { workflow } => format!("{workflow}.start"),
360 }
361 }
362}