Skip to main content

turnframe_core/
reduce.rs

1//! The whole-turn reduction contract (spec §13).
2//!
3//! The reducer is the deterministic decision engine: it takes the turn's
4//! [`Understanding`], the projected views, the active interactions, the target map and
5//! the policy snapshot, and produces a [`ReductionPlan`] that gives **every act an
6//! explicit result** (I11) and groups executable commands into batches. It performs no
7//! I/O and has no side effects.
8//!
9//! # Same-turn precedence (normative default, spec §13.2)
10//!
11//! 1. A correction or cancellation supersedes the act it names; understanding links
12//!    them, and the reducer never guesses a correction from a repeated operation.
13//! 2. `DoNotSubmit` blocks every submission act in the turn.
14//! 3. A question never becomes an action.
15//! 4. An ambiguous target blocks only the acts that depend on it.
16//! 5. Independent questions remain answerable.
17//! 6. A click binds more strongly than a typed answer to the same card.
18//! 7. A high-risk card is never resolved from typed text.
19
20use std::collections::{BTreeMap, BTreeSet};
21use std::fmt;
22
23use chrono::{DateTime, Utc};
24use indexmap::IndexMap;
25use serde::{Deserialize, Serialize};
26
27use crate::case::{CaseKey, CaseRef};
28use crate::command::{AtomicityScope, CommandBatch, RiskClass, origin_satisfies};
29use crate::error::{DomainRejection, ReductionError};
30use crate::flow::{DomainEnumeration, ErasedWorkflowView};
31use crate::hash::{Digest, HashError, canonical_digest};
32use crate::ids::{
33    BatchId, CommandId, InteractionId, OperationKey, OptionId, QuestionId, TurnId, WorkflowKey,
34};
35use crate::interaction::{InteractionKind, InteractionSpec, TextResolutionPolicy};
36use crate::operation::OperationCatalog;
37use crate::plan::AnswerBasis;
38use crate::plan::limits::PlanLimits;
39use crate::policy::{PolicyDecision, PolicySnapshot};
40use crate::response::{NarratableFact, ServerNotice};
41use crate::target::{TargetResolution, TargetTokenMap};
42use crate::turn::TurnInput;
43use crate::understanding::{
44    ActAction, ActId, ConstraintKind, Understanding, UnderstoodAct, UnitId,
45};
46
47/// The pure whole-turn reducer (spec §13).
48pub trait TurnReducer: Send + Sync {
49    /// Reduces one understood turn into an execution plan. Must be deterministic for
50    /// the same inputs and must not perform I/O.
51    ///
52    /// # Errors
53    ///
54    /// A [`ReductionError`] for an understanding over the limits or a plan that fails
55    /// its own consistency checks.
56    fn reduce(
57        &self,
58        input: &TurnInput,
59        understanding: &Understanding,
60        context: &ReductionContext,
61    ) -> Result<ReductionPlan, ReductionError>;
62}
63
64/// What the reducer needs to know about an active interaction.
65#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
66pub struct ActiveInteractionSummary {
67    /// The interaction.
68    pub interaction_id: InteractionId,
69    /// Case and bound revision.
70    pub case_ref: CaseRef,
71    /// Shape.
72    pub kind: InteractionKind,
73    /// Whether it owns unqualified answers for the case.
74    pub blocking: bool,
75    /// Stored option ids (for validating interpreted options).
76    pub option_ids: Vec<OptionId>,
77    /// Whether typed text may resolve it.
78    pub text_resolution: TextResolutionPolicy,
79    /// Highest risk class an answer authorizes, so the reducer can apply rule 8
80    /// without loading the journaled commands. Conservative by default.
81    #[serde(default = "RiskClass::conservative")]
82    pub confirms_risk: RiskClass,
83    /// Hash of the payload the user saw.
84    pub payload_hash: Digest,
85}
86
87impl ActiveInteractionSummary {
88    /// Returns `true` when typed text may resolve this card at all: the stored
89    /// policy allows it **and** what it confirms is low risk (spec §13.2 rule
90    /// 8, §15.7).
91    #[must_use]
92    pub fn accepts_text_resolution(&self) -> bool {
93        self.text_resolution != TextResolutionPolicy::Never
94            && !self.confirms_risk.needs_trusted_origin()
95            && !self.kind.authorizes_commands()
96    }
97}
98
99/// Everything the reducer sees besides the plan and the turn.
100#[derive(Debug, Clone)]
101pub struct ReductionContext {
102    /// Projected views of every loaded case.
103    pub views: IndexMap<CaseKey, ErasedWorkflowView>,
104    /// Active interactions in the conversation.
105    pub active_interactions: Vec<ActiveInteractionSummary>,
106    /// Tokens issued for this turn.
107    pub target_map: TargetTokenMap,
108    /// Operations offered for this turn.
109    pub operations: OperationCatalog,
110    /// Policy configuration.
111    pub policy: PolicySnapshot,
112    /// Plan limits.
113    pub limits: PlanLimits,
114    /// The clock value the reducer must use (it must not read the clock).
115    pub now: DateTime<Utc>,
116    /// The cases every write on which has to pass through a click.
117    ///
118    /// Set by the application, per case, through its case directory. A command
119    /// on one of these that would otherwise have run silently gets the
120    /// confirmation raised instead, on a card that names the case — the same
121    /// thing [`ConstraintKind::AskBeforeApplying`] does when a user asks for it. Empty is the ordinary case and changes
122    /// nothing.
123    pub confirm_every_write: BTreeSet<CaseKey>,
124    /// The cases that are in view only because the actor may reach them.
125    ///
126    /// Set by the application, per case, through its case directory: a record
127    /// another conversation is filling in, one left open in a thread that no
128    /// longer exists, anything the turn has in view because it is reachable
129    /// rather than because this turn is about it.
130    ///
131    /// What it does is keep such a case from acting as a subject on its own.
132    /// It is not hidden and not unreachable — it stays in the catalog under its
133    /// own names, and the moment an act of this turn lands on it it is a
134    /// subject like any other, which is the only way a draft left in a deleted
135    /// conversation can ever be finished. What it stops is the three things a
136    /// case does merely by being in the room: holding the door against a
137    /// second case beside it, briefing the writer in the imperative, and
138    /// putting its own outstanding fields in front of a reader answering about
139    /// something else.
140    ///
141    /// Empty is the ordinary case and changes nothing.
142    pub subject_only_when_named: BTreeSet<CaseKey>,
143}
144
145impl ReductionContext {
146    /// View of a case, if loaded.
147    #[must_use]
148    pub fn view_for(&self, key: &CaseKey) -> Option<&ErasedWorkflowView> {
149        self.views.get(key)
150    }
151
152    /// Whether every write on `key` has to pass through a click.
153    #[must_use]
154    pub fn confirms_every_write(&self, key: &CaseKey) -> bool {
155        self.confirm_every_write.contains(key)
156    }
157
158    /// Whether `key` is in view only because the actor may reach it, so it is
159    /// a subject of this turn only if the turn names it.
160    ///
161    /// See [`Self::subject_only_when_named`].
162    #[must_use]
163    pub fn is_subject_only_when_named(&self, key: &CaseKey) -> bool {
164        self.subject_only_when_named.contains(key)
165    }
166
167    /// The active blocking interaction of a case, if any.
168    #[must_use]
169    pub fn blocking_interaction_for(&self, key: &CaseKey) -> Option<&ActiveInteractionSummary> {
170        self.active_interactions
171            .iter()
172            .find(|i| i.blocking && i.case_ref.key() == *key)
173    }
174}
175
176/// Reference to a command inside a reduction plan.
177#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
178pub struct CommandRef {
179    /// The batch.
180    pub batch_id: BatchId,
181    /// The command.
182    pub command_id: CommandId,
183}
184
185impl fmt::Display for CommandRef {
186    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
187        write!(f, "{}/{}", self.batch_id, self.command_id)
188    }
189}
190
191/// The explicit result of one act (spec §13.3, I11).
192///
193/// New results are expected as the reducer learns to say more, so downstream
194/// matches need a wildcard arm.
195#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
196#[serde(tag = "kind", rename_all = "snake_case")]
197#[non_exhaustive]
198pub enum PlannedActResult {
199    /// Compiled into commands that may execute this turn.
200    ReadyToExecute {
201        /// The commands.
202        command_refs: Vec<CommandRef>,
203    },
204    /// Compiled, but policy requires a confirmation first.
205    AwaitingConfirmation {
206        /// The interaction to create.
207        interaction_spec: InteractionSpec,
208    },
209    /// The target or intent is ambiguous.
210    NeedsClarification {
211        /// The interaction to create.
212        interaction_spec: InteractionSpec,
213    },
214    /// Rejected deterministically. Carries the whole [`DomainRejection`],
215    /// including the structured `details` the UI needs, instead of re-declaring
216    /// its code.
217    Rejected {
218        /// Why the domain refused.
219        rejection: DomainRejection,
220    },
221    /// A later act in the same turn cancelled or corrected it.
222    SupersededByCorrection,
223    /// Valid but changes nothing (already in the requested state).
224    NoChange,
225    /// Arguments are missing, unstated or refused; the user is asked, nothing runs.
226    NeedsValue {
227        /// The arguments to ask for.
228        arguments: Vec<String>,
229        /// The domain's explanation, when it refused a value.
230        reason: Option<String>,
231    },
232    /// Another unit aimed at the same record was not understood.
233    Held {
234        /// That unit.
235        because: UnitId,
236    },
237    /// Waits for an earlier act of the turn that is itself waiting for a click.
238    AwaitingPrerequisite {
239        /// The act it waits for.
240        act: ActId,
241    },
242}
243
244impl PlannedActResult {
245    /// Snake-case variant name, for a record or a report.
246    ///
247    /// The distinction a measurement needs is here and nowhere else: only
248    /// `rejected` is the structure refusing a reading. `no_change` is a valid
249    /// act on a record already in the requested state, and
250    /// `awaiting_confirmation` is one waiting for a person — both journal no
251    /// commands, and counting either as a refusal reports the design working
252    /// as the model failing.
253    #[must_use]
254    pub const fn name(&self) -> &'static str {
255        match self {
256            Self::ReadyToExecute { .. } => "ready_to_execute",
257            Self::AwaitingConfirmation { .. } => "awaiting_confirmation",
258            Self::NeedsClarification { .. } => "needs_clarification",
259            Self::Rejected { .. } => "rejected",
260            Self::SupersededByCorrection => "superseded_by_correction",
261            Self::NoChange => "no_change",
262            Self::NeedsValue { .. } => "needs_value",
263            Self::Held { .. } => "held",
264            Self::AwaitingPrerequisite { .. } => "awaiting_prerequisite",
265        }
266    }
267}
268
269impl From<DomainRejection> for PlannedActResult {
270    fn from(rejection: DomainRejection) -> Self {
271        Self::Rejected { rejection }
272    }
273}
274
275/// One act with its resolution and result.
276#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
277pub struct PlannedAct {
278    /// The act as understood.
279    pub act: UnderstoodAct,
280    /// Target resolution, when the act has a target.
281    pub target: Option<TargetResolution>,
282    /// The result.
283    pub result: PlannedActResult,
284}
285
286/// Which sources an answer must rest on (spec §19.1).
287#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
288#[serde(rename_all = "snake_case")]
289pub enum SourcePolicy {
290    /// Any available source.
291    AnySource,
292    /// Only authoritative sources (case state, approved knowledge).
293    AuthoritativeOnly,
294    /// Sources must be cited.
295    RequireCitations,
296    /// No retrieval; answer from state and general knowledge only.
297    NoRetrieval,
298}
299
300/// A question to answer, with its explicit state basis (spec §19.1).
301#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
302pub struct AnswerTask {
303    /// Stable id within the turn.
304    pub question_id: QuestionId,
305    /// The question.
306    pub question: String,
307    /// State basis the reducer decided (it may override the model's preference).
308    pub basis: AnswerBasis,
309    /// Cases the question is about.
310    pub case_refs: Vec<CaseRef>,
311    /// Reference to a proposed diff when `basis` is `ProposedState`.
312    #[serde(default, skip_serializing_if = "Option::is_none")]
313    pub proposed_diff_ref: Option<String>,
314    /// Source requirements.
315    pub required_sources: SourcePolicy,
316    /// Where in the normalized message the question's words are.
317    ///
318    /// The question's text is derived from that span, and carrying the span
319    /// too is what lets composition keep those words away from the stage that
320    /// must not answer them. `None` for a task rebuilt from a replay record,
321    /// where the message is not at hand.
322    #[serde(default, skip_serializing_if = "Option::is_none")]
323    pub asked_at: Option<TextSpan>,
324    /// Complete value sets a workflow declared for what this question is
325    /// about.
326    ///
327    /// Non-empty means the deterministic layer already holds the answer, so
328    /// composition settles the question from these rather than asking a model
329    /// to describe a set it would have to remember.
330    #[serde(default, skip_serializing_if = "Vec::is_empty")]
331    pub enumerations: Vec<DomainEnumeration>,
332    /// What the user can do now, for a question about that.
333    #[serde(default, skip_serializing_if = "Vec::is_empty")]
334    pub capabilities: Vec<Capability>,
335    /// Whether it follows up the assistant's last message.
336    #[serde(default)]
337    pub continues_previous: bool,
338}
339
340/// One thing the user can do now: an operation on offer, in its workflow's words.
341#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
342pub struct Capability {
343    /// The workflow.
344    pub workflow: WorkflowKey,
345    /// The operation.
346    pub operation: OperationKey,
347    /// What it does, as the workflow says it.
348    pub summary: String,
349}
350
351/// A range of the normalized user message, in bytes.
352#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
353pub struct TextSpan {
354    /// First byte of the range.
355    pub start_byte: usize,
356    /// One past its last byte.
357    pub end_byte: usize,
358}
359
360/// The reducer's output (spec §13).
361#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
362pub struct ReductionPlan {
363    /// The turn.
364    pub turn_id: TurnId,
365    /// One entry per input act, by index (I11).
366    pub acts: Vec<PlannedAct>,
367    /// Command batches ready for execution, grouped by atomicity scope.
368    /// Commands are erased to JSON; the registry converts them at the boundary.
369    pub batches: Vec<CommandBatch<serde_json::Value>>,
370    /// Policy decisions for every command.
371    pub policy_decisions: Vec<PolicyDecision>,
372    /// Questions to answer.
373    pub answer_tasks: Vec<AnswerTask>,
374    /// Constraints that were applied.
375    pub constraints_applied: Vec<ConstraintKind>,
376    /// Notices to show (e.g. "Nothing has been submitted").
377    pub notices: Vec<ServerNotice>,
378    /// Interactions to persist **before** any command executes (clarifications
379    /// and confirmations). Deduplicated by `InteractionSpec::key`.
380    pub pre_execution_interactions: Vec<InteractionSpec>,
381    /// Digest of everything above (see [`Self::compute_hash`]).
382    pub plan_hash: Digest,
383    /// Operations of acts a correction or a cancel in the same message replaced.
384    ///
385    /// Supersession is the one reduction rule that makes a turn do less than
386    /// its plan said, and until it was counted a message carrying four data
387    /// that produced one command was only discoverable by reading the
388    /// conversation that followed. One entry per dropped act, in plan order.
389    #[serde(default)]
390    pub superseded_operations: Vec<OperationKey>,
391    /// Domain refusals, as facts the narration stage may rest on.
392    ///
393    /// The notice beside them tells the user deterministically; these tell the
394    /// narrator, so its prose does not ask for something else as though the
395    /// refusal had not happened.
396    #[serde(default)]
397    pub refusals: Vec<NarratableFact>,
398    /// Acts that were accepted and changed nothing, as facts the narration
399    /// stage may rest on.
400    ///
401    /// Beside [`Self::refusals`] and not inside it, because the two are
402    /// different outcomes and one of them is counted: the runtime raises its
403    /// refusal signal once per entry there, and a no-op filed among them would
404    /// be a refusal in every dashboard that reads it. They travel together only
405    /// at the point where both become facts for the writer.
406    #[serde(default)]
407    pub changed_nothing: Vec<NarratableFact>,
408    /// Acts this turn prepared and held behind a confirmation card, as facts
409    /// the narration stage may rest on.
410    ///
411    /// Separate from [`Self::refusals`] and not folded into it, because these
412    /// are not refusals: that list is counted as
413    /// [`Signal::ActRefused`](crate::observe::Signal::ActRefused), and an act
414    /// waiting for a click is one the server intends to run. See
415    /// [`NarratableFact::ActAwaitingConfirmation`].
416    #[serde(default)]
417    pub awaiting_confirmation: Vec<NarratableFact>,
418}
419
420#[derive(Serialize)]
421struct PlanHashInput<'a> {
422    turn_id: &'a TurnId,
423    acts: &'a [PlannedAct],
424    batches: &'a [CommandBatch<serde_json::Value>],
425    policy_decisions: &'a [PolicyDecision],
426    answer_tasks: &'a [AnswerTask],
427    constraints_applied: &'a [ConstraintKind],
428    notices: &'a [ServerNotice],
429    pre_execution_interactions: &'a [InteractionSpec],
430}
431
432impl ReductionPlan {
433    /// Computes the digest of every field except `plan_hash`.
434    pub fn compute_hash(&self) -> Result<Digest, HashError> {
435        canonical_digest(&PlanHashInput {
436            turn_id: &self.turn_id,
437            acts: &self.acts,
438            batches: &self.batches,
439            policy_decisions: &self.policy_decisions,
440            answer_tasks: &self.answer_tasks,
441            constraints_applied: &self.constraints_applied,
442            notices: &self.notices,
443            pre_execution_interactions: &self.pre_execution_interactions,
444        })
445    }
446
447    /// Sets `plan_hash` from the current content.
448    pub fn with_hash(mut self) -> Result<Self, HashError> {
449        self.plan_hash = self.compute_hash()?;
450        Ok(self)
451    }
452
453    /// Returns `true` when `plan_hash` matches the content.
454    pub fn verify_hash(&self) -> Result<bool, HashError> {
455        Ok(self.compute_hash()? == self.plan_hash)
456    }
457
458    /// Structural consistency checks.
459    ///
460    /// Everything here is a property a correct reducer already has; the point
461    /// is that a plan which fails one of them must never reach execution, so a
462    /// reducer bug becomes a refused turn instead of an unauthorized command:
463    ///
464    /// * every act of the understanding appears exactly once (I11);
465    /// * every `ReadyToExecute` command reference points at a batch command;
466    /// * a `ReadyToExecute` act resolved its target exactly, unless it starts a
467    ///   workflow or picks a target (spec §12.2, I8), and its commands run on
468    ///   the case it resolved to;
469    /// * every batched command has exactly one [`PolicyDecision`], that
470    ///   decision allows it, and its origin satisfies the policy the decision
471    ///   recorded ([`origin_satisfies`], I9, I12);
472    /// * no refused decision names a command that is nevertheless batched;
473    /// * every spec embedded in an act result is listed in
474    ///   `pre_execution_interactions` (by key);
475    /// * `PerCase` batches target a single case.
476    pub fn validate(&self, understood: &[ActId]) -> Result<(), ReductionError> {
477        let mut seen = BTreeSet::new();
478        for planned in &self.acts {
479            if !understood.contains(&planned.act.id) || !seen.insert(planned.act.id) {
480                return Err(ReductionError::InconsistentPlan {
481                    detail: format!("act {} unknown or duplicated", planned.act.id),
482                });
483            }
484        }
485        if seen.len() != understood.len() {
486            return Err(ReductionError::InconsistentPlan {
487                detail: format!("{} of {} acts have a result", seen.len(), understood.len()),
488            });
489        }
490        let mut commands: BTreeMap<CommandRef, &CaseRef> = BTreeMap::new();
491        for batch in &self.batches {
492            for envelope in &batch.envelopes {
493                commands.insert(
494                    CommandRef {
495                        batch_id: batch.batch_id,
496                        command_id: envelope.command_id,
497                    },
498                    &envelope.case_ref,
499                );
500            }
501        }
502        let command_refs: BTreeSet<CommandRef> = commands.keys().copied().collect();
503        let spec_keys: BTreeSet<&str> = self
504            .pre_execution_interactions
505            .iter()
506            .map(|s| s.key.as_str())
507            .collect();
508        for planned in &self.acts {
509            match &planned.result {
510                PlannedActResult::ReadyToExecute { command_refs: refs } => {
511                    if let Some(missing) = refs.iter().find(|r| !command_refs.contains(r)) {
512                        return Err(ReductionError::InconsistentPlan {
513                            detail: format!("dangling command reference {missing}"),
514                        });
515                    }
516                    let needs_exact_target = !matches!(planned.act.action, ActAction::Start { .. });
517                    let resolved = planned.target.as_ref().and_then(TargetResolution::exact);
518                    if needs_exact_target && resolved.is_none() {
519                        return Err(ReductionError::InconsistentPlan {
520                            detail: format!(
521                                "act {} executes without an exact target",
522                                planned.act.id
523                            ),
524                        });
525                    }
526                    if let Some(resolved) = resolved
527                        && let Some(elsewhere) = refs
528                            .iter()
529                            .filter_map(|r| commands.get(r))
530                            .find(|case_ref| !case_ref.same_case(resolved))
531                    {
532                        return Err(ReductionError::InconsistentPlan {
533                            detail: format!(
534                                "act {} resolved to {}/{} but a command targets {}/{}",
535                                planned.act.id,
536                                resolved.workflow,
537                                resolved.case_id,
538                                elsewhere.workflow,
539                                elsewhere.case_id
540                            ),
541                        });
542                    }
543                }
544                PlannedActResult::AwaitingConfirmation { interaction_spec }
545                | PlannedActResult::NeedsClarification { interaction_spec } => {
546                    if !spec_keys.contains(interaction_spec.key.as_str()) {
547                        return Err(ReductionError::InconsistentPlan {
548                            detail: format!(
549                                "interaction spec {} not listed for creation",
550                                interaction_spec.key
551                            ),
552                        });
553                    }
554                }
555                PlannedActResult::Rejected { .. }
556                | PlannedActResult::SupersededByCorrection
557                | PlannedActResult::NoChange
558                | PlannedActResult::NeedsValue { .. }
559                | PlannedActResult::Held { .. }
560                | PlannedActResult::AwaitingPrerequisite { .. } => {}
561            }
562        }
563        if let Some(bad) = self
564            .batches
565            .iter()
566            .find(|b| matches!(b.scope, AtomicityScope::PerCase) && !b.is_single_case())
567        {
568            return Err(ReductionError::InconsistentPlan {
569                detail: format!("per-case batch {} spans several cases", bad.batch_id),
570            });
571        }
572        self.validate_policy_coverage(&command_refs)
573    }
574
575    /// Every batched command is policed, allowed and authorized by its own
576    /// origin. See [`Self::validate`].
577    fn validate_policy_coverage(
578        &self,
579        command_refs: &BTreeSet<CommandRef>,
580    ) -> Result<(), ReductionError> {
581        let mut decisions: BTreeMap<CommandRef, &PolicyDecision> = BTreeMap::new();
582        for decision in &self.policy_decisions {
583            if decisions.insert(decision.command_ref, decision).is_some() {
584                return Err(ReductionError::InconsistentPlan {
585                    detail: format!(
586                        "command {} has several policy decisions",
587                        decision.command_ref
588                    ),
589                });
590            }
591            if !decision.allowed && command_refs.contains(&decision.command_ref) {
592                return Err(ReductionError::InconsistentPlan {
593                    detail: format!("refused command {} is batched", decision.command_ref),
594                });
595            }
596        }
597        for batch in &self.batches {
598            for envelope in &batch.envelopes {
599                let command_ref = CommandRef {
600                    batch_id: batch.batch_id,
601                    command_id: envelope.command_id,
602                };
603                let Some(decision) = decisions.get(&command_ref) else {
604                    return Err(ReductionError::InconsistentPlan {
605                        detail: format!("command {command_ref} has no policy decision"),
606                    });
607                };
608                if !origin_satisfies(&envelope.origin, &decision.policy) {
609                    return Err(ReductionError::InconsistentPlan {
610                        detail: format!("command {command_ref} has an origin its policy refuses"),
611                    });
612                }
613            }
614        }
615        Ok(())
616    }
617
618    /// All command references in batch order.
619    #[must_use]
620    pub fn command_refs(&self) -> Vec<CommandRef> {
621        self.batches
622            .iter()
623            .flat_map(|b| {
624                b.envelopes.iter().map(move |e| CommandRef {
625                    batch_id: b.batch_id,
626                    command_id: e.command_id,
627                })
628            })
629            .collect()
630    }
631
632    /// Returns `true` when nothing will execute this turn.
633    #[must_use]
634    pub fn has_no_effects(&self) -> bool {
635        self.batches.iter().all(CommandBatch::is_empty)
636    }
637}
638
639#[cfg(test)]
640mod tests {
641    use super::*;
642    use crate::case::CaseRef;
643    use crate::command::{
644        CommandEnvelope, CommandOrigin, CommandPolicy, IdempotencyKey, ResolutionChannel,
645    };
646    use crate::ids::{AccountId, CaseRevision, InteractionId, OperationKey, WorkflowKey};
647    use crate::interaction::{
648        ActionClass, InteractionOption, InteractionPayload, StoredInteractionAction,
649    };
650    use crate::policy::reason;
651    use crate::turn::ActorContext;
652
653    fn id(index: usize) -> ActId {
654        ActId::new(UnitId(u16::try_from(index + 1).unwrap()), 1)
655    }
656
657    fn ids(count: usize) -> Vec<ActId> {
658        (0..count).map(id).collect()
659    }
660
661    fn understood(index: usize, action: ActAction) -> UnderstoodAct {
662        UnderstoodAct {
663            id: id(index),
664            action,
665            target: crate::understanding::ActTarget::Card,
666            arguments: BTreeMap::new(),
667            words: crate::understanding::WordRange {
668                first: 0,
669                last: 0,
670                start: 0,
671                end: 2,
672            },
673            depends_on: vec![],
674            status: crate::understanding::ActStatus::Ready,
675        }
676    }
677
678    fn act(index: usize, result: PlannedActResult) -> PlannedAct {
679        PlannedAct {
680            act: understood(
681                index,
682                ActAction::Start {
683                    workflow: WorkflowKey::from("w"),
684                },
685            ),
686            target: None,
687            result,
688        }
689    }
690
691    fn apply_act(
692        index: usize,
693        target: Option<TargetResolution>,
694        result: PlannedActResult,
695    ) -> PlannedAct {
696        PlannedAct {
697            act: understood(
698                index,
699                ActAction::Apply {
700                    operation: OperationKey::from("w.op"),
701                },
702            ),
703            target,
704            result,
705        }
706    }
707
708    fn case() -> CaseRef {
709        CaseRef::new("w", "c1", CaseRevision(1))
710    }
711
712    fn plan(acts: Vec<PlannedAct>) -> ReductionPlan {
713        ReductionPlan {
714            turn_id: TurnId::nil(),
715            acts,
716            batches: vec![],
717            policy_decisions: vec![],
718            answer_tasks: vec![],
719            superseded_operations: vec![],
720            refusals: vec![],
721            changed_nothing: vec![],
722            awaiting_confirmation: vec![],
723            constraints_applied: vec![],
724            notices: vec![],
725            pre_execution_interactions: vec![],
726            plan_hash: Digest::of_bytes(b""),
727        }
728    }
729
730    fn envelope(
731        command_id: CommandId,
732        case_ref: CaseRef,
733        origin: CommandOrigin,
734    ) -> CommandEnvelope<serde_json::Value> {
735        CommandEnvelope {
736            command_id,
737            turn_id: TurnId::nil(),
738            actor: ActorContext::new("acct", "u1"),
739            case_ref,
740            idempotency_key: IdempotencyKey::new("k"),
741            origin,
742            command: serde_json::json!({"do": true}),
743        }
744    }
745
746    fn confirmed_origin() -> CommandOrigin {
747        CommandOrigin::ConfirmedInteraction {
748            interaction_id: InteractionId::nil(),
749            payload_hash: Digest::of_bytes(b"p"),
750            interaction_kind: InteractionKind::ConfirmCommand,
751            action_class: ActionClass::ConfirmsCommands,
752            channel: ResolutionChannel::Click,
753        }
754    }
755
756    fn direct_origin() -> CommandOrigin {
757        CommandOrigin::DirectSafeUserAct {
758            evidence_digest: Digest::of_bytes(b"e"),
759        }
760    }
761
762    fn decision(command_ref: CommandRef, policy: CommandPolicy, allowed: bool) -> PolicyDecision {
763        PolicyDecision {
764            command_ref,
765            policy,
766            requires_interaction: None,
767            allowed,
768            reason_key: if allowed {
769                reason::ALLOWED.to_owned()
770            } else {
771                reason::CONFIRMATION_REQUIRED.to_owned()
772            },
773        }
774    }
775
776    /// A plan with one batched command, its act and its decision.
777    fn executing_plan(
778        origin: CommandOrigin,
779        policy: CommandPolicy,
780        allowed: bool,
781    ) -> ReductionPlan {
782        let batch_id = BatchId::derive(&TurnId::nil(), &case().key(), &AtomicityScope::PerCase);
783        let command_id = CommandId::derive(&TurnId::nil(), id(0), 0);
784        let command_ref = CommandRef {
785            batch_id,
786            command_id,
787        };
788        let mut p = plan(vec![apply_act(
789            0,
790            Some(TargetResolution::Exact { case_ref: case() }),
791            PlannedActResult::ReadyToExecute {
792                command_refs: vec![command_ref],
793            },
794        )]);
795        p.batches = vec![CommandBatch {
796            batch_id,
797            scope: AtomicityScope::PerCase,
798            envelopes: vec![envelope(command_id, case(), origin)],
799        }];
800        p.policy_decisions = vec![decision(command_ref, policy, allowed)];
801        p
802    }
803
804    #[test]
805    fn every_act_needs_a_result() {
806        let p = plan(vec![act(0, PlannedActResult::NoChange)]);
807        assert!(p.validate(&ids(1)).is_ok());
808        assert!(p.validate(&ids(2)).is_err());
809        let dup = plan(vec![
810            act(0, PlannedActResult::NoChange),
811            act(0, PlannedActResult::NoChange),
812        ]);
813        assert!(dup.validate(&ids(2)).is_err());
814    }
815
816    #[test]
817    fn dangling_refs_are_detected() {
818        let p = plan(vec![act(
819            0,
820            PlannedActResult::ReadyToExecute {
821                command_refs: vec![CommandRef {
822                    batch_id: BatchId::nil(),
823                    command_id: CommandId::nil(),
824                }],
825            },
826        )]);
827        assert!(matches!(
828            p.validate(&ids(1)),
829            Err(ReductionError::InconsistentPlan { .. })
830        ));
831    }
832
833    #[test]
834    fn an_executing_act_must_have_resolved_its_target_exactly() {
835        let ok = executing_plan(confirmed_origin(), CommandPolicy::conservative(), true);
836        assert_eq!(ok.validate(&ids(1)), Ok(()));
837        for target in [
838            None,
839            Some(TargetResolution::Missing),
840            Some(TargetResolution::Ambiguous { candidates: vec![] }),
841            Some(TargetResolution::Stale {
842                case_ref: case(),
843                current_revision: CaseRevision(2),
844            }),
845        ] {
846            let mut p = ok.clone();
847            p.acts[0].target = target;
848            assert!(
849                matches!(
850                    p.validate(&ids(1)),
851                    Err(ReductionError::InconsistentPlan { .. })
852                ),
853                "an ambiguous or missing target may not execute (I8)"
854            );
855        }
856        // Starting a workflow has no target to resolve.
857        let mut start = plan(vec![act(
858            0,
859            PlannedActResult::ReadyToExecute {
860                command_refs: vec![],
861            },
862        )]);
863        start.policy_decisions = vec![];
864        assert_eq!(start.validate(&ids(1)), Ok(()));
865    }
866
867    #[test]
868    fn every_batched_command_is_policed_allowed_and_authorized() {
869        let mut no_decision =
870            executing_plan(confirmed_origin(), CommandPolicy::conservative(), true);
871        no_decision.policy_decisions.clear();
872        assert!(matches!(
873            no_decision.validate(&ids(1)),
874            Err(ReductionError::InconsistentPlan { .. })
875        ));
876
877        let refused = executing_plan(confirmed_origin(), CommandPolicy::conservative(), false);
878        assert!(matches!(
879            refused.validate(&ids(1)),
880            Err(ReductionError::InconsistentPlan { .. })
881        ));
882
883        // The decision says "allowed" but the envelope carries an origin the
884        // recorded policy refuses: the plan is not trustworthy.
885        let lying = executing_plan(direct_origin(), CommandPolicy::conservative(), true);
886        assert!(matches!(
887            lying.validate(&ids(1)),
888            Err(ReductionError::InconsistentPlan { .. })
889        ));
890
891        let mut twice = executing_plan(confirmed_origin(), CommandPolicy::conservative(), true);
892        let duplicate = twice.policy_decisions[0].clone();
893        twice.policy_decisions.push(duplicate);
894        assert!(matches!(
895            twice.validate(&ids(1)),
896            Err(ReductionError::InconsistentPlan { .. })
897        ));
898
899        // A command that runs on another case than the one the act resolved to
900        // is exactly the mix-up an exact target is supposed to prevent.
901        let mut elsewhere = executing_plan(confirmed_origin(), CommandPolicy::conservative(), true);
902        elsewhere.batches[0].envelopes[0].case_ref = CaseRef::new("w", "other", CaseRevision(1));
903        assert!(matches!(
904            elsewhere.validate(&ids(1)),
905            Err(ReductionError::InconsistentPlan { .. })
906        ));
907
908        // A refused command that is *not* batched is exactly how a plan records
909        // "this needs a confirmation first".
910        let mut awaiting = executing_plan(direct_origin(), CommandPolicy::conservative(), false);
911        awaiting.batches.clear();
912        awaiting.acts[0].result = PlannedActResult::NoChange;
913        assert_eq!(awaiting.validate(&ids(1)), Ok(()));
914    }
915
916    #[test]
917    fn interaction_specs_must_be_listed_for_creation() {
918        let spec = InteractionSpec::new(
919            "confirm:acts[0]",
920            case(),
921            InteractionKind::ConfirmCommand,
922            InteractionPayload::new("Send?")
923                .with_option(InteractionOption::new(
924                    "yes",
925                    "Send",
926                    StoredInteractionAction::ConfirmCommands {
927                        command_refs: vec![],
928                    },
929                ))
930                .with_option(InteractionOption::new(
931                    "no",
932                    "Cancel",
933                    StoredInteractionAction::DeclineCommands,
934                )),
935        );
936        for result in [
937            PlannedActResult::AwaitingConfirmation {
938                interaction_spec: spec.clone(),
939            },
940            PlannedActResult::NeedsClarification {
941                interaction_spec: spec.clone(),
942            },
943        ] {
944            let orphan = plan(vec![act(0, result.clone())]);
945            assert!(
946                matches!(
947                    orphan.validate(&ids(1)),
948                    Err(ReductionError::InconsistentPlan { .. })
949                ),
950                "a card nobody creates leaves the act unanswerable"
951            );
952            let mut listed = plan(vec![act(0, result)]);
953            listed.pre_execution_interactions = vec![spec.clone()];
954            assert_eq!(listed.validate(&ids(1)), Ok(()));
955        }
956    }
957
958    #[test]
959    fn per_case_batches_may_not_span_cases() {
960        let mut p = executing_plan(confirmed_origin(), CommandPolicy::conservative(), true);
961        let other = envelope(
962            CommandId::derive(&TurnId::nil(), id(0), 1),
963            CaseRef::new("w", "c2", CaseRevision(1)),
964            confirmed_origin(),
965        );
966        let command_ref = CommandRef {
967            batch_id: p.batches[0].batch_id,
968            command_id: other.command_id,
969        };
970        p.batches[0].envelopes.push(other);
971        p.policy_decisions
972            .push(decision(command_ref, CommandPolicy::conservative(), true));
973        assert!(matches!(
974            p.validate(&ids(1)),
975            Err(ReductionError::InconsistentPlan { .. })
976        ));
977    }
978
979    #[test]
980    fn hash_tracks_content() {
981        let p = plan(vec![act(0, PlannedActResult::NoChange)])
982            .with_hash()
983            .unwrap();
984        assert!(p.verify_hash().unwrap());
985        let mut changed = p.clone();
986        changed.acts[0].result = PlannedActResult::SupersededByCorrection;
987        assert!(!changed.verify_hash().unwrap());
988    }
989
990    /// A reducer that only uses derived identifiers, as shipped reducers must.
991    struct FixedReducer;
992
993    impl TurnReducer for FixedReducer {
994        fn reduce(
995            &self,
996            input: &TurnInput,
997            understanding: &Understanding,
998            _context: &ReductionContext,
999        ) -> Result<ReductionPlan, ReductionError> {
1000            let turn_id = input.turn_id;
1001            let batch_id = BatchId::derive(&turn_id, &case().key(), &AtomicityScope::PerCase);
1002            let mut acts = Vec::new();
1003            let mut envelopes = Vec::new();
1004            let mut policy_decisions = Vec::new();
1005            for (index, act) in understanding.acts.iter().enumerate() {
1006                let command_id = CommandId::derive(&turn_id, act.id, 0);
1007                let command_ref = CommandRef {
1008                    batch_id,
1009                    command_id,
1010                };
1011                envelopes.push(envelope(command_id, case(), confirmed_origin()));
1012                policy_decisions.push(decision(command_ref, CommandPolicy::conservative(), true));
1013                acts.push(apply_act(
1014                    index,
1015                    Some(TargetResolution::Exact { case_ref: case() }),
1016                    PlannedActResult::ReadyToExecute {
1017                        command_refs: vec![command_ref],
1018                    },
1019                ));
1020            }
1021            ReductionPlan {
1022                turn_id,
1023                acts,
1024                batches: vec![CommandBatch {
1025                    batch_id,
1026                    scope: AtomicityScope::PerCase,
1027                    envelopes,
1028                }],
1029                policy_decisions,
1030                answer_tasks: vec![],
1031                superseded_operations: vec![],
1032                refusals: vec![],
1033                changed_nothing: vec![],
1034                awaiting_confirmation: vec![],
1035                constraints_applied: understanding.constraints.iter().map(|c| c.kind).collect(),
1036                notices: vec![],
1037                pre_execution_interactions: vec![],
1038                plan_hash: Digest::of_bytes(b""),
1039            }
1040            .with_hash()
1041            .map_err(|_| ReductionError::Hash)
1042        }
1043    }
1044
1045    #[test]
1046    fn two_reductions_of_the_same_inputs_agree_on_the_plan_hash() {
1047        let input = TurnInput {
1048            turn_id: TurnId::nil(),
1049            conversation_id: crate::ids::ConversationId::nil(),
1050            actor: ActorContext::new("acct", "u1"),
1051            text: Some("do it".into()),
1052            interaction_response: None,
1053            attachments: vec![],
1054            origin: None,
1055            locale: crate::locale::Locale::from("it-IT"),
1056            effort: None,
1057        };
1058        let understanding = Understanding {
1059            acts: vec![understood(
1060                0,
1061                ActAction::Start {
1062                    workflow: WorkflowKey::from("w"),
1063                },
1064            )],
1065            ..Understanding::default()
1066        };
1067        let context = ReductionContext {
1068            views: IndexMap::new(),
1069            active_interactions: vec![],
1070            target_map: TargetTokenMap::new(AccountId::from("acct"), TurnId::nil()),
1071            operations: OperationCatalog::default(),
1072            policy: PolicySnapshot::conservative(),
1073            limits: PlanLimits::conservative(),
1074            now: chrono::DateTime::from_timestamp(1_700_000_000, 0).unwrap(),
1075            confirm_every_write: BTreeSet::new(),
1076            subject_only_when_named: BTreeSet::new(),
1077        };
1078        let first = FixedReducer
1079            .reduce(&input, &understanding, &context)
1080            .unwrap();
1081        let second = FixedReducer
1082            .reduce(&input, &understanding, &context)
1083            .unwrap();
1084        assert_eq!(first.plan_hash, second.plan_hash);
1085        assert_eq!(first, second);
1086        assert_eq!(first.validate(&ids(1)), Ok(()));
1087        assert!(first.verify_hash().unwrap());
1088    }
1089
1090    #[test]
1091    fn a_high_risk_card_never_accepts_typed_text() {
1092        let summary = ActiveInteractionSummary {
1093            interaction_id: InteractionId::nil(),
1094            case_ref: case(),
1095            kind: InteractionKind::SingleSelect,
1096            blocking: true,
1097            option_ids: vec![OptionId::from("a")],
1098            text_resolution: TextResolutionPolicy::ModelInterpretedLowRisk,
1099            confirms_risk: RiskClass::ReversibleLowRisk,
1100            payload_hash: Digest::of_bytes(b"p"),
1101        };
1102        assert!(summary.accepts_text_resolution());
1103        let risky = ActiveInteractionSummary {
1104            confirms_risk: RiskClass::Irreversible,
1105            ..summary.clone()
1106        };
1107        assert!(!risky.accepts_text_resolution());
1108        let confirming = ActiveInteractionSummary {
1109            kind: InteractionKind::ConfirmCommand,
1110            ..summary.clone()
1111        };
1112        assert!(!confirming.accepts_text_resolution());
1113        let never = ActiveInteractionSummary {
1114            text_resolution: TextResolutionPolicy::Never,
1115            ..summary
1116        };
1117        assert!(!never.accepts_text_resolution());
1118    }
1119}