Skip to main content

turnframe_runtime/reduce/
mod.rs

1//! The whole-turn reducer (spec §13).
2//!
3//! [`DefaultTurnReducer`] is the deterministic decision engine: every understood act
4//! comes out with an explicit result (I11), executable commands are grouped into
5//! batches, and the plan hashes to a value a replay reproduces (I20). It does no I/O,
6//! reads no clock and consults no randomness.
7//!
8//! | Step | What it does |
9//! | --- | --- |
10//! | limits | an understanding over the deployment's limits is refused whole |
11//! | understanding | an act that asks for a value or is held keeps that result |
12//! | catalog | an unoffered operation, or arguments its schema refuses, refuses the act |
13//! | prerequisites | an act on a record an earlier act opens waits for that act |
14//! | targets | one target per act, through [`TargetResolver`] |
15//! | compilation, validation, policy | per command, the domain first |
16//! | batching | by case and atomicity scope; mutations on one case commit together |
17//! | answers | one [`AnswerTask`] per question, an unsafe basis overridden |
18//! | self-check | [`ReductionPlan::validate`] runs before the plan is returned |
19//!
20//! Build one reducer per turn: the states it compiles against are the turn's own.
21
22mod 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/// The deterministic whole-turn reducer of §13. Build one per turn.
71#[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    /// Builds a reducer for one turn: `resolver` carries the tokens and the card on
85    /// screen, `policy` the mode and the card copy, `config` the execution and
86    /// narration sections.
87    #[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    /// Adds the loaded state of one case. A case with no entry compiles against `None`,
107    /// which is what a case that does not exist yet looks like.
108    #[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    /// Adds several loaded states.
115    #[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    /// Declares the origin a card answer minted, for the one act the card itself put
125    /// into the plan. No other act receives it (I12).
126    #[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    /// Replaces the notices the reducer writes itself. The shipped copy is English.
133    #[must_use]
134    pub fn with_copy(mut self, copy: NoticeCopy) -> Self {
135        self.copy = copy;
136        self
137    }
138
139    /// The resolver this reducer was built with.
140    #[must_use]
141    pub const fn resolver(&self) -> &TargetResolver {
142        &self.resolver
143    }
144}
145
146/// A reduction plus the envelopes its confirmation cards will authorize (§15.3): a
147/// click resumes exactly the commands that were reviewed.
148#[derive(Debug, Clone, PartialEq, Eq)]
149#[non_exhaustive]
150pub struct ReducedTurn {
151    /// The reduction plan.
152    pub plan: ReductionPlan,
153    /// Batches that only execute once a confirmation card is answered.
154    pub pending: Vec<CommandBatch<serde_json::Value>>,
155}
156
157impl ReducedTurn {
158    /// Every envelope waiting behind a confirmation, in batch order.
159    #[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    /// Reduces one turn, keeping the envelopes its confirmation cards will authorize.
170    ///
171    /// # Errors
172    ///
173    /// Whatever [`TurnReducer::reduce`] returns.
174    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
198/// Everything one reduction accumulates.
199struct 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    /// Cases this turn opens, in act order, so a second opening is asked about the first.
217    minted: Vec<CaseRef>,
218    /// The case each opening act opens, for the acts that apply to it.
219    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    /// The user's words a range covers, as the turn carries them.
259    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}