Skip to main content

glass/browser/session/
workflow.rs

1//! Versioned, declarative workflow definitions.
2//!
3//! This module contains the validated workflow contract, execution evidence,
4//! bounded checkpointing, and resume reconciliation.
5
6use super::types::{BatchMode, BatchStep, BrowserResult, VerificationPredicate};
7use super::{
8    INTENT_RESOLUTION_SCHEMA_VERSION, IntentConfidence, IntentConstraints, IntentPolicyDecision,
9    IntentScope, SemanticIntentAction, SemanticIntentExecutionRequest, SemanticIntentRequest,
10    SemanticIntentResult, SemanticResolution, SemanticResolutionPolicy, SemanticRouteIdentity,
11    target_fingerprint_digest,
12};
13use serde::de::Error as DeError;
14use serde::{Deserialize, Deserializer, Serialize, Serializer};
15use serde_json::Value;
16use sha2::{Digest, Sha256};
17use std::collections::{BTreeMap, BTreeSet};
18use std::fmt;
19use std::time::{Duration, Instant};
20use url::Url;
21
22mod checkpoint;
23pub use checkpoint::*;
24mod recorder;
25pub use recorder::*;
26
27/// The workflow definition schema understood by this crate.
28pub const WORKFLOW_SCHEMA_VERSION: u32 = 1;
29const MAX_NAME_BYTES: usize = 128;
30const MAX_DESCRIPTION_BYTES: usize = 4 * 1024;
31const MAX_INPUTS: usize = 64;
32const MAX_STEPS: usize = 64;
33const MAX_DURATION_MS: u64 = 15 * 60 * 1_000;
34const MAX_RETRIES: u32 = 8;
35const MAX_EXTRACTED_BYTES: usize = 4 * 1024 * 1024;
36const MAX_TEXT_BYTES: usize = 64 * 1024;
37const MAX_TARGET_BYTES: usize = 1_024;
38const MAX_WAIT_CONDITION_BYTES: usize = 4 * 1024;
39const MAX_STEP_REPETITIONS: u32 = 8;
40const MAX_WORKFLOW_TRACE_EVENTS: usize = 2_048;
41const WORKFLOW_TRACE_SCHEMA_VERSION: u8 = 1;
42const WORKFLOW_CHECKPOINT_SCHEMA_VERSION: u8 = 1;
43const MAX_WORKFLOW_CHECKPOINT_BYTES: usize = 8 * 1024;
44const MAX_CHECKPOINT_HISTORY_STATES: usize = 64;
45const MAX_CHECKPOINT_EXECUTION_IDS: usize = 8;
46const MAX_INTENT_PURPOSE_BYTES: usize = 256;
47
48/// A complete declarative workflow.
49#[derive(Debug, Clone, Serialize, Deserialize)]
50#[serde(rename_all = "camelCase")]
51pub struct WorkflowDefinition {
52    /// Schema version, independent of the workflow's business version.
53    pub schema_version: u32,
54    /// Stable human-readable workflow name.
55    pub name: String,
56    /// Caller-owned version of this workflow definition.
57    pub workflow_version: String,
58    #[serde(default, skip_serializing_if = "Option::is_none")]
59    pub description: Option<String>,
60    pub inputs: BTreeMap<String, WorkflowInput>,
61    pub budgets: WorkflowBudgets,
62    #[serde(default)]
63    pub preconditions: Vec<VerificationPredicate>,
64    pub steps: Vec<WorkflowStep>,
65    pub terminal_condition: VerificationPredicate,
66    pub outputs: BTreeMap<String, WorkflowOutputDeclaration>,
67}
68
69impl WorkflowDefinition {
70    /// Parse and validate a JSON workflow definition.
71    pub fn from_json(input: &str) -> Result<Self, WorkflowValidationError> {
72        let value: Value = serde_json::from_str(input)
73            .map_err(|error| WorkflowValidationError::new("$", format!("invalid JSON: {error}")))?;
74        Self::from_value(value)
75    }
76
77    /// Deserialize and validate a workflow definition from JSON data.
78    pub fn from_value(value: Value) -> Result<Self, WorkflowValidationError> {
79        require_object_fields(
80            &value,
81            &[
82                "schemaVersion",
83                "name",
84                "workflowVersion",
85                "inputs",
86                "budgets",
87                "steps",
88                "terminalCondition",
89                "outputs",
90            ],
91        )?;
92        reject_unknown_fields(
93            &value,
94            "$",
95            &[
96                "schemaVersion",
97                "name",
98                "workflowVersion",
99                "description",
100                "inputs",
101                "budgets",
102                "preconditions",
103                "steps",
104                "terminalCondition",
105                "outputs",
106            ],
107        )?;
108        if let Some(inputs) = value.get("inputs").and_then(Value::as_object) {
109            for (name, input) in inputs {
110                reject_unknown_fields(
111                    input,
112                    &format!("inputs.{name}"),
113                    &["valueType", "type", "required", "maxLength", "sensitive"],
114                )?;
115            }
116        }
117        if let Some(budgets) = value.get("budgets") {
118            reject_unknown_fields(
119                budgets,
120                "budgets",
121                &[
122                    "maxSteps",
123                    "maxDurationMs",
124                    "maxRetries",
125                    "maxExtractedBytes",
126                ],
127            )?;
128        }
129        if let Some(outputs) = value.get("outputs").and_then(Value::as_object) {
130            for (name, output) in outputs {
131                reject_unknown_fields(
132                    output,
133                    &format!("outputs.{name}"),
134                    &["valueType", "type", "source", "required", "sensitive"],
135                )?;
136            }
137        }
138        if let Some(preconditions) = value.get("preconditions").and_then(Value::as_array) {
139            for (index, predicate) in preconditions.iter().enumerate() {
140                reject_predicate_fields(predicate, &format!("preconditions[{index}]"))?;
141            }
142        }
143        if let Some(predicate) = value.get("terminalCondition") {
144            reject_predicate_fields(predicate, "terminalCondition")?;
145        }
146        if let Some(steps) = value.get("steps").and_then(Value::as_array) {
147            for (index, step) in steps.iter().enumerate() {
148                reject_unknown_fields(
149                    step,
150                    &format!("steps[{index}]"),
151                    &[
152                        "id",
153                        "action",
154                        "intent",
155                        "when",
156                        "expect",
157                        "beforeRetry",
158                        "transaction",
159                        "idempotencyKey",
160                        "maxRetries",
161                        "repeat",
162                        "url",
163                        "timeoutMs",
164                        "target",
165                        "text",
166                        "value",
167                        "condition",
168                        "dx",
169                        "dy",
170                        "includeDom",
171                        "includeScreenshot",
172                        "includeFormValues",
173                    ],
174                )?;
175                if let Some(intent) = step.get("intent") {
176                    if step.get("action").is_some() {
177                        return Err(WorkflowValidationError::new(
178                            format!("steps[{index}].action"),
179                            "semantic intent steps cannot also declare a batch action",
180                        ));
181                    }
182                    reject_unknown_fields(
183                        intent,
184                        &format!("steps[{index}].intent"),
185                        &[
186                            "action",
187                            "purpose",
188                            "intent",
189                            "scope",
190                            "constraints",
191                            "resolutionPolicy",
192                            "value",
193                        ],
194                    )?;
195                }
196                for field in ["when", "expect", "beforeRetry"] {
197                    if let Some(predicate) = step.get(field) {
198                        reject_predicate_fields(predicate, &format!("steps[{index}].{field}"))?;
199                    }
200                }
201            }
202        }
203        let definition: Self = serde_json::from_value(value).map_err(|error| {
204            WorkflowValidationError::new("$", format!("invalid workflow shape: {error}"))
205        })?;
206        definition.validate()?;
207        Ok(definition)
208    }
209
210    /// Validate all structural and resource-boundary constraints.
211    pub fn validate(&self) -> Result<(), WorkflowValidationError> {
212        if self.schema_version != WORKFLOW_SCHEMA_VERSION {
213            return Err(WorkflowValidationError::new(
214                "schemaVersion",
215                format!(
216                    "unsupported schema version {}; expected {}",
217                    self.schema_version, WORKFLOW_SCHEMA_VERSION
218                ),
219            ));
220        }
221        validate_name("name", &self.name)?;
222        validate_name("workflowVersion", &self.workflow_version)?;
223        if let Some(description) = &self.description {
224            validate_bytes("description", description, 1, MAX_DESCRIPTION_BYTES)?;
225        }
226        if self.inputs.len() > MAX_INPUTS {
227            return Err(WorkflowValidationError::new(
228                "inputs",
229                format!("must contain at most {MAX_INPUTS} entries"),
230            ));
231        }
232        for (name, input) in &self.inputs {
233            validate_map_key("inputs", name)?;
234            input.validate(&format!("inputs.{name}"))?;
235        }
236        self.budgets.validate("budgets")?;
237        if self.steps.is_empty() {
238            return Err(WorkflowValidationError::new(
239                "steps",
240                "must contain at least one step",
241            ));
242        }
243        if self.steps.len() > self.budgets.max_steps as usize {
244            return Err(WorkflowValidationError::new(
245                "steps",
246                "step count exceeds budgets.maxSteps",
247            ));
248        }
249        let expanded_steps: usize = self.steps.iter().map(|step| step.repeat as usize).sum();
250        if expanded_steps > self.budgets.max_steps as usize {
251            return Err(WorkflowValidationError::new(
252                "steps",
253                "expanded repetition count exceeds budgets.maxSteps",
254            ));
255        }
256
257        let mut ids = BTreeSet::new();
258        let mut idempotency_keys = BTreeSet::new();
259        for (index, step) in self.steps.iter().enumerate() {
260            let path = format!("steps[{index}]");
261            step.validate(&path, self.budgets.max_retries)?;
262            if !ids.insert(step.id.as_str()) {
263                return Err(WorkflowValidationError::new(
264                    format!("{path}.id"),
265                    format!("duplicate step ID {:?}", step.id),
266                ));
267            }
268            if let Some(key) = &step.idempotency_key
269                && !idempotency_keys.insert(key.as_str())
270            {
271                return Err(WorkflowValidationError::new(
272                    format!("{path}.idempotencyKey"),
273                    format!("duplicate idempotency key {:?}", key),
274                ));
275            }
276        }
277        for (index, predicate) in self.preconditions.iter().enumerate() {
278            validate_predicate(predicate, &format!("preconditions[{index}]"))?;
279        }
280        validate_predicate(&self.terminal_condition, "terminalCondition")?;
281        for (name, output) in &self.outputs {
282            validate_map_key("outputs", name)?;
283            output.validate(&format!("outputs.{name}"))?;
284        }
285        Ok(())
286    }
287
288    /// Return stable JSON suitable for hashing, caching, or audit records.
289    pub fn to_canonical_json(&self) -> Result<String, WorkflowValidationError> {
290        self.validate()?;
291        serde_json::to_string(self).map_err(|error| {
292            WorkflowValidationError::new("$", format!("cannot serialize workflow: {error}"))
293        })
294    }
295
296    /// Validate caller-provided input values before execution starts.
297    pub fn validate_inputs(
298        &self,
299        values: &BTreeMap<String, Value>,
300    ) -> Result<(), WorkflowValidationError> {
301        for name in values.keys() {
302            if !self.inputs.contains_key(name) {
303                return Err(WorkflowValidationError::new(
304                    format!("inputs.{name}"),
305                    "value has no declared input",
306                ));
307            }
308        }
309        for (name, declaration) in &self.inputs {
310            match values.get(name) {
311                Some(value) => declaration.validate_value(&format!("inputs.{name}"), value)?,
312                None if declaration.required => {
313                    return Err(WorkflowValidationError::new(
314                        format!("inputs.{name}"),
315                        "required input is missing",
316                    ));
317                }
318                None => {}
319            }
320        }
321        Ok(())
322    }
323
324    /// Resolve bounded `${inputs.name}` placeholders in declared actions.
325    /// Resolution happens before browser startup or dispatch and never
326    /// evaluates arbitrary expressions.
327    pub fn resolve_actions(
328        &self,
329        values: &BTreeMap<String, Value>,
330    ) -> Result<Vec<WorkflowStep>, WorkflowValidationError> {
331        self.validate_inputs(values)?;
332        self.steps
333            .iter()
334            .enumerate()
335            .map(|(index, step)| {
336                let mut resolved = step.clone();
337                if let Some(intent) = &step.intent {
338                    resolved.intent = Some(resolve_workflow_intent(
339                        intent,
340                        values,
341                        &format!("steps[{index}].intent"),
342                    )?);
343                } else {
344                    resolved.action = resolve_batch_step(
345                        &step.action,
346                        values,
347                        &format!("steps[{index}].action"),
348                    )?;
349                }
350                resolved.validate(&format!("steps[{index}]"), self.budgets.max_retries)?;
351                Ok(resolved)
352            })
353            .collect()
354    }
355}
356
357/// A declared workflow input and its accepted JSON type.
358#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
359#[serde(rename_all = "camelCase")]
360pub struct WorkflowInput {
361    #[serde(alias = "type")]
362    pub value_type: WorkflowValueType,
363    #[serde(default = "default_true")]
364    pub required: bool,
365    #[serde(default, skip_serializing_if = "Option::is_none")]
366    pub max_length: Option<usize>,
367    #[serde(default, skip_serializing_if = "Option::is_none")]
368    pub sensitive: Option<bool>,
369}
370
371impl WorkflowInput {
372    fn validate(&self, path: &str) -> Result<(), WorkflowValidationError> {
373        if self.max_length == Some(0) || self.max_length.is_some_and(|value| value > MAX_TEXT_BYTES)
374        {
375            return Err(WorkflowValidationError::new(
376                format!("{path}.maxLength"),
377                format!("must be 1..={MAX_TEXT_BYTES}"),
378            ));
379        }
380        Ok(())
381    }
382
383    fn validate_value(&self, path: &str, value: &Value) -> Result<(), WorkflowValidationError> {
384        let valid = match self.value_type {
385            WorkflowValueType::String => value.is_string(),
386            WorkflowValueType::Integer => value.as_i64().is_some() || value.as_u64().is_some(),
387            WorkflowValueType::Number => value.is_number(),
388            WorkflowValueType::Boolean => value.is_boolean(),
389            WorkflowValueType::Url => value
390                .as_str()
391                .is_some_and(|value| Url::parse(value).is_ok()),
392        };
393        if !valid {
394            return Err(WorkflowValidationError::new(
395                path,
396                format!("expected {}", self.value_type),
397            ));
398        }
399        if let Some(max_length) = self.max_length {
400            let length = value
401                .as_str()
402                .map_or_else(|| value.to_string().len(), str::len);
403            if length > max_length {
404                return Err(WorkflowValidationError::new(
405                    path,
406                    format!("value exceeds maxLength {max_length}"),
407                ));
408            }
409        }
410        Ok(())
411    }
412}
413
414fn default_true() -> bool {
415    true
416}
417
418/// Runtime resource limits for a workflow.
419#[derive(Debug, Clone, Serialize, Deserialize)]
420#[serde(rename_all = "camelCase")]
421pub struct WorkflowBudgets {
422    pub max_steps: u32,
423    pub max_duration_ms: u64,
424    #[serde(default)]
425    pub max_retries: u32,
426    pub max_extracted_bytes: usize,
427}
428
429impl WorkflowBudgets {
430    fn validate(&self, path: &str) -> Result<(), WorkflowValidationError> {
431        if self.max_steps == 0 || self.max_steps as usize > MAX_STEPS {
432            return Err(WorkflowValidationError::new(
433                format!("{path}.maxSteps"),
434                format!("must be 1..={MAX_STEPS}"),
435            ));
436        }
437        if self.max_duration_ms == 0 || self.max_duration_ms > MAX_DURATION_MS {
438            return Err(WorkflowValidationError::new(
439                format!("{path}.maxDurationMs"),
440                format!("must be 1..={MAX_DURATION_MS}"),
441            ));
442        }
443        if self.max_retries > MAX_RETRIES {
444            return Err(WorkflowValidationError::new(
445                format!("{path}.maxRetries"),
446                format!("must be <= {MAX_RETRIES}"),
447            ));
448        }
449        if self.max_extracted_bytes == 0 || self.max_extracted_bytes > MAX_EXTRACTED_BYTES {
450            return Err(WorkflowValidationError::new(
451                format!("{path}.maxExtractedBytes"),
452                format!("must be 1..={MAX_EXTRACTED_BYTES}"),
453            ));
454        }
455        Ok(())
456    }
457}
458
459/// JSON value types supported by workflow inputs and outputs.
460#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
461#[serde(rename_all = "snake_case")]
462pub enum WorkflowValueType {
463    String,
464    Integer,
465    Number,
466    Boolean,
467    Url,
468}
469
470impl fmt::Display for WorkflowValueType {
471    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
472        formatter.write_str(match self {
473            Self::String => "string",
474            Self::Integer => "integer",
475            Self::Number => "number",
476            Self::Boolean => "boolean",
477            Self::Url => "url",
478        })
479    }
480}
481
482/// A named action with an optional postcondition.
483#[derive(Debug, Clone)]
484pub struct WorkflowStep {
485    pub id: String,
486    pub action: BatchStep,
487    pub intent: Option<WorkflowIntentStep>,
488    pub when: Option<VerificationPredicate>,
489    pub expect: Option<VerificationPredicate>,
490    pub before_retry: Option<VerificationPredicate>,
491    pub transaction: WorkflowTransactionClass,
492    pub idempotency_key: Option<String>,
493    pub max_retries: u32,
494    pub repeat: u32,
495}
496
497/// A semantic workflow action resolved against fresh page evidence at runtime.
498/// Existing locator-based workflow actions remain unchanged.
499#[derive(Debug, Clone, Serialize, Deserialize)]
500#[serde(rename_all = "camelCase", deny_unknown_fields)]
501pub struct WorkflowIntentStep {
502    pub action: SemanticIntentAction,
503    #[serde(default, skip_serializing_if = "Option::is_none")]
504    pub purpose: Option<String>,
505    #[serde(default, skip_serializing_if = "Option::is_none")]
506    pub intent: Option<String>,
507    #[serde(default)]
508    pub scope: IntentScope,
509    #[serde(default)]
510    pub constraints: IntentConstraints,
511    #[serde(default = "default_workflow_resolution_policy")]
512    pub resolution_policy: SemanticResolutionPolicy,
513    #[serde(default, skip_serializing_if = "Option::is_none")]
514    pub value: Option<String>,
515}
516
517impl WorkflowIntentStep {
518    fn validate(&self, path: &str) -> Result<(), WorkflowValidationError> {
519        let supplied = self.purpose.is_some() as u8 + self.intent.is_some() as u8;
520        if supplied != 1 {
521            return Err(WorkflowValidationError::new(
522                format!("{path}.purpose"),
523                "provide exactly one of purpose or intent",
524            ));
525        }
526        if let Some(purpose) = &self.purpose {
527            validate_bytes(
528                &format!("{path}.purpose"),
529                purpose,
530                1,
531                MAX_INTENT_PURPOSE_BYTES,
532            )?;
533        }
534        if let Some(intent) = &self.intent {
535            validate_bytes(
536                &format!("{path}.intent"),
537                intent,
538                1,
539                MAX_INTENT_PURPOSE_BYTES * 2,
540            )?;
541        }
542        if let Some(value) = &self.value {
543            validate_bytes(&format!("{path}.value"), value, 0, 4_096)?;
544        }
545        self.execution_request(path).map(|_| ())
546    }
547
548    fn execution_request(
549        &self,
550        path: &str,
551    ) -> Result<SemanticIntentExecutionRequest, WorkflowValidationError> {
552        let intent = self
553            .intent
554            .clone()
555            .or_else(|| self.purpose.as_deref().map(purpose_to_intent))
556            .ok_or_else(|| {
557                WorkflowValidationError::new(format!("{path}.purpose"), "intent phrase is required")
558            })?;
559        let request = SemanticIntentRequest {
560            schema_version: INTENT_RESOLUTION_SCHEMA_VERSION,
561            intent,
562            action: self.action,
563            scope: self.scope.clone(),
564            constraints: self.constraints.clone(),
565            resolution_policy: self.resolution_policy,
566            expected_revision: None,
567        };
568        let execution = SemanticIntentExecutionRequest {
569            request,
570            candidate_id: "workflow-selected-candidate".into(),
571            value: self.value.clone(),
572        };
573        execution.validate().map_err(|error| {
574            WorkflowValidationError::new(format!("{path}.{}", error.path), error.reason)
575        })?;
576        Ok(execution)
577    }
578}
579
580fn default_workflow_resolution_policy() -> SemanticResolutionPolicy {
581    SemanticResolutionPolicy::RequireUniqueHighConfidence
582}
583
584fn purpose_to_intent(purpose: &str) -> String {
585    let mut result = String::with_capacity(purpose.len() + 8);
586    for (index, character) in purpose.chars().enumerate() {
587        if character.is_ascii_uppercase() && index > 0 {
588            result.push(' ');
589        }
590        result.push(character.to_ascii_lowercase());
591    }
592    result
593}
594
595/// Effect classification used to decide whether a failed attempt may be
596/// replayed before dispatch.
597#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
598#[serde(rename_all = "snake_case")]
599pub enum WorkflowTransactionClass {
600    ReadOnly,
601    Idempotent,
602    ConditionallyIdempotent,
603    NonIdempotent,
604    #[default]
605    Unknown,
606}
607
608impl WorkflowTransactionClass {
609    /// Whether this class permits a retry known to have happened before
610    /// dispatch.
611    pub fn permits_pre_dispatch_retry(self) -> bool {
612        matches!(
613            self,
614            Self::ReadOnly | Self::Idempotent | Self::ConditionallyIdempotent
615        )
616    }
617}
618
619impl Serialize for WorkflowStep {
620    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
621    where
622        S: Serializer,
623    {
624        let mut action = if let Some(intent) = &self.intent {
625            serde_json::json!({ "intent": intent })
626        } else {
627            serde_json::to_value(&self.action).map_err(serde::ser::Error::custom)?
628        };
629        if self.intent.is_none() {
630            let object = action
631                .as_object_mut()
632                .ok_or_else(|| serde::ser::Error::custom("workflow action must be an object"))?;
633            for (internal, public) in [
634                ("timeout_ms", "timeoutMs"),
635                ("include_dom", "includeDom"),
636                ("include_screenshot", "includeScreenshot"),
637                ("include_form_values", "includeFormValues"),
638            ] {
639                if let Some(value) = object.remove(internal) {
640                    object.insert(public.to_string(), value);
641                }
642            }
643        }
644
645        let mut workflow = serde_json::Map::new();
646        workflow.insert("id".into(), Value::String(self.id.clone()));
647        if let Value::Object(action) = action {
648            workflow.extend(action);
649        }
650        if let Some(when) = &self.when {
651            workflow.insert(
652                "when".into(),
653                serde_json::to_value(when).map_err(serde::ser::Error::custom)?,
654            );
655        }
656        if let Some(expect) = &self.expect {
657            workflow.insert(
658                "expect".into(),
659                serde_json::to_value(expect).map_err(serde::ser::Error::custom)?,
660            );
661        }
662        if let Some(before_retry) = &self.before_retry {
663            workflow.insert(
664                "beforeRetry".into(),
665                serde_json::to_value(before_retry).map_err(serde::ser::Error::custom)?,
666            );
667        }
668        workflow.insert(
669            "transaction".into(),
670            serde_json::to_value(self.transaction).map_err(serde::ser::Error::custom)?,
671        );
672        if let Some(key) = &self.idempotency_key {
673            workflow.insert("idempotencyKey".into(), Value::String(key.clone()));
674        }
675        if self.max_retries > 0 {
676            workflow.insert("maxRetries".into(), Value::from(self.max_retries));
677        }
678        if self.repeat > 1 {
679            workflow.insert("repeat".into(), Value::from(self.repeat));
680        }
681        Value::Object(workflow).serialize(serializer)
682    }
683}
684
685impl<'de> Deserialize<'de> for WorkflowStep {
686    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
687    where
688        D: Deserializer<'de>,
689    {
690        let mut workflow = serde_json::Map::<String, Value>::deserialize(deserializer)?;
691        let id = workflow
692            .remove("id")
693            .ok_or_else(|| D::Error::custom("workflow step is missing id"))?;
694        let id = serde_json::from_value(id).map_err(D::Error::custom)?;
695        let when = workflow
696            .remove("when")
697            .map(serde_json::from_value)
698            .transpose()
699            .map_err(D::Error::custom)?;
700        let expect = workflow
701            .remove("expect")
702            .map(serde_json::from_value)
703            .transpose()
704            .map_err(D::Error::custom)?;
705        let before_retry = workflow
706            .remove("beforeRetry")
707            .map(serde_json::from_value)
708            .transpose()
709            .map_err(D::Error::custom)?;
710        let transaction = workflow
711            .remove("transaction")
712            .map(serde_json::from_value)
713            .transpose()
714            .map_err(D::Error::custom)?
715            .unwrap_or_default();
716        let idempotency_key = workflow
717            .remove("idempotencyKey")
718            .map(serde_json::from_value)
719            .transpose()
720            .map_err(D::Error::custom)?;
721        let max_retries = workflow
722            .remove("maxRetries")
723            .map(serde_json::from_value)
724            .transpose()
725            .map_err(D::Error::custom)?
726            .unwrap_or(0);
727        let repeat = workflow
728            .remove("repeat")
729            .map(serde_json::from_value)
730            .transpose()
731            .map_err(D::Error::custom)?
732            .unwrap_or(1);
733        let intent = workflow
734            .remove("intent")
735            .map(serde_json::from_value)
736            .transpose()
737            .map_err(D::Error::custom)?;
738        for (public, internal) in [
739            ("timeoutMs", "timeout_ms"),
740            ("includeDom", "include_dom"),
741            ("includeScreenshot", "include_screenshot"),
742            ("includeFormValues", "include_form_values"),
743        ] {
744            if let Some(value) = workflow.remove(public) {
745                workflow.insert(internal.into(), value);
746            }
747        }
748        let action = if intent.is_some() {
749            BatchStep::Observe {
750                include_dom: false,
751                include_screenshot: false,
752                include_form_values: false,
753            }
754        } else {
755            serde_json::from_value(Value::Object(workflow)).map_err(D::Error::custom)?
756        };
757        Ok(Self {
758            id,
759            action,
760            intent,
761            when,
762            expect,
763            before_retry,
764            transaction,
765            idempotency_key,
766            max_retries,
767            repeat,
768        })
769    }
770}
771
772impl WorkflowStep {
773    fn validate(
774        &self,
775        path: &str,
776        workflow_max_retries: u32,
777    ) -> Result<(), WorkflowValidationError> {
778        validate_name(&format!("{path}.id"), &self.id)?;
779        if let Some(intent) = &self.intent {
780            intent.validate(&format!("{path}.intent"))?;
781        } else {
782            validate_batch_step(&self.action, &format!("{path}.action"))?;
783        }
784        if let Some(predicate) = &self.when {
785            validate_predicate(predicate, &format!("{path}.when"))?;
786        }
787        if let Some(predicate) = &self.expect {
788            validate_predicate(predicate, &format!("{path}.expect"))?;
789        }
790        if let Some(predicate) = &self.before_retry {
791            validate_predicate(predicate, &format!("{path}.beforeRetry"))?;
792        }
793        if let Some(key) = &self.idempotency_key {
794            validate_bytes(&format!("{path}.idempotencyKey"), key, 1, 256)?;
795        }
796        if self.repeat == 0 || self.repeat > MAX_STEP_REPETITIONS {
797            return Err(WorkflowValidationError::new(
798                format!("{path}.repeat"),
799                format!("must be 1..={MAX_STEP_REPETITIONS}"),
800            ));
801        }
802        if self.repeat > 1 && self.when.is_some() {
803            return Err(WorkflowValidationError::new(
804                format!("{path}.repeat"),
805                "conditional steps cannot repeat automatically",
806            ));
807        }
808        if self.repeat > 1 && !self.transaction.permits_pre_dispatch_retry() {
809            return Err(WorkflowValidationError::new(
810                format!("{path}.repeat"),
811                "unknown or non-idempotent steps cannot repeat automatically",
812            ));
813        }
814        if self.max_retries > workflow_max_retries {
815            return Err(WorkflowValidationError::new(
816                format!("{path}.maxRetries"),
817                "step retry count exceeds budgets.maxRetries",
818            ));
819        }
820        match self.transaction {
821            WorkflowTransactionClass::ConditionallyIdempotent if self.idempotency_key.is_none() => {
822                return Err(WorkflowValidationError::new(
823                    format!("{path}.idempotencyKey"),
824                    "conditionally idempotent steps require an idempotency key",
825                ));
826            }
827            WorkflowTransactionClass::NonIdempotent | WorkflowTransactionClass::Unknown
828                if self.max_retries > 0 =>
829            {
830                return Err(WorkflowValidationError::new(
831                    format!("{path}.maxRetries"),
832                    "non-idempotent or unknown steps cannot be retried automatically",
833                ));
834            }
835            _ => {}
836        }
837        Ok(())
838    }
839}
840
841/// Durable state of one workflow step.
842#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
843#[serde(rename_all = "snake_case")]
844pub enum WorkflowStepState {
845    Pending,
846    Ready,
847    Preflight,
848    Resolving,
849    NotDispatched,
850    Dispatched,
851    EffectObserved,
852    Verified,
853    OutputsExtracted,
854    Committed,
855    FailedBeforeDispatch,
856    FailedAfterDispatch,
857    Indeterminate,
858    Skipped,
859}
860
861impl WorkflowStepState {
862    /// Return whether a state transition is valid for the linear runner.
863    pub fn can_transition_to(self, next: Self) -> bool {
864        matches!(
865            (self, next),
866            (Self::Pending, Self::Ready | Self::Skipped)
867                | (Self::Ready, Self::Preflight | Self::Skipped)
868                | (
869                    Self::Preflight,
870                    Self::Resolving | Self::EffectObserved | Self::FailedBeforeDispatch
871                )
872                | (Self::Resolving, Self::NotDispatched | Self::Dispatched)
873                | (Self::NotDispatched, Self::FailedBeforeDispatch)
874                | (
875                    Self::Dispatched,
876                    Self::EffectObserved | Self::FailedAfterDispatch | Self::Indeterminate
877                )
878                | (
879                    Self::EffectObserved,
880                    Self::Verified | Self::FailedAfterDispatch | Self::Indeterminate
881                )
882                | (Self::Verified, Self::OutputsExtracted)
883                | (Self::OutputsExtracted, Self::Committed)
884                | (Self::FailedBeforeDispatch, Self::Ready)
885                | (Self::Committed, Self::Ready)
886        )
887    }
888}
889
890/// The retained state and transition history for one workflow step.
891#[derive(Debug, Clone, Serialize, Deserialize)]
892#[serde(rename_all = "camelCase")]
893pub struct WorkflowStepRecord {
894    pub id: String,
895    pub state: WorkflowStepState,
896    pub history: Vec<WorkflowStepState>,
897    pub attempts: u32,
898    #[serde(default, skip_serializing_if = "Vec::is_empty")]
899    pub execution_ids: Vec<String>,
900    #[serde(default)]
901    pub dispatch_acknowledged: bool,
902    #[serde(default)]
903    pub effect_observed: bool,
904    #[serde(default)]
905    pub postcondition_verified: bool,
906    #[serde(default)]
907    pub retry_safe: bool,
908    #[serde(default, skip_serializing_if = "Option::is_none")]
909    pub previous_revision: Option<u64>,
910    #[serde(default, skip_serializing_if = "Option::is_none")]
911    pub current_revision: Option<u64>,
912    #[serde(skip_serializing_if = "Option::is_none")]
913    pub branch_decision: Option<WorkflowBranchDecision>,
914    #[serde(skip_serializing_if = "Option::is_none")]
915    pub intent_evidence: Option<WorkflowIntentEvidence>,
916    #[serde(skip_serializing_if = "Option::is_none")]
917    pub error: Option<String>,
918}
919
920/// Bounded evidence linking a semantic workflow step to its accepted target.
921#[derive(Debug, Clone, Serialize, Deserialize)]
922#[serde(rename_all = "camelCase")]
923pub struct WorkflowIntentEvidence {
924    pub resolution_id: String,
925    pub candidate_id: String,
926    pub revision: u64,
927    pub resolution: super::SemanticResolution,
928    pub policy_decision: super::IntentPolicyDecision,
929    pub confidence: super::IntentConfidence,
930    pub fingerprint: Option<super::SemanticTargetFingerprint>,
931}
932
933impl WorkflowStepRecord {
934    fn new(id: &str) -> Self {
935        Self {
936            id: id.to_string(),
937            state: WorkflowStepState::Pending,
938            history: vec![WorkflowStepState::Pending],
939            attempts: 0,
940            execution_ids: Vec::new(),
941            dispatch_acknowledged: false,
942            effect_observed: false,
943            postcondition_verified: false,
944            retry_safe: false,
945            previous_revision: None,
946            current_revision: None,
947            branch_decision: None,
948            intent_evidence: None,
949            error: None,
950        }
951    }
952
953    fn transition(&mut self, next: WorkflowStepState) -> Result<(), String> {
954        if !self.state.can_transition_to(next) {
955            return Err(format!(
956                "invalid workflow step transition {} -> {}",
957                state_name(self.state),
958                state_name(next)
959            ));
960        }
961        self.state = next;
962        self.history.push(next);
963        Ok(())
964    }
965
966    fn fail(&mut self, state: WorkflowStepState, error: &str) {
967        let _ = self.transition(state);
968        self.error = Some(bound_workflow_text(error, 512));
969    }
970}
971
972/// Overall outcome of a workflow run.
973#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
974#[serde(rename_all = "snake_case")]
975pub enum WorkflowRunStatus {
976    Completed,
977    Failed,
978    BudgetExhausted,
979    ResumeRequired,
980}
981
982/// Evidence that the workflow's terminal condition was satisfied.
983#[derive(Debug, Clone, Serialize, Deserialize)]
984#[serde(rename_all = "camelCase")]
985pub struct WorkflowTerminalProof {
986    pub predicate: VerificationPredicate,
987    pub revision: u64,
988    pub state: String,
989}
990
991/// Result of a linear workflow run.
992#[derive(Debug, Clone, Serialize, Deserialize)]
993#[serde(rename_all = "camelCase")]
994pub struct WorkflowRunResult {
995    pub run_id: String,
996    pub name: String,
997    pub workflow_version: String,
998    pub status: WorkflowRunStatus,
999    pub steps: Vec<WorkflowStepRecord>,
1000    pub trace: WorkflowTrace,
1001    pub outputs: BTreeMap<String, WorkflowOutput>,
1002    #[serde(skip_serializing_if = "Option::is_none")]
1003    pub terminal_proof: Option<WorkflowTerminalProof>,
1004    #[serde(skip_serializing_if = "Option::is_none")]
1005    pub failed_step: Option<String>,
1006    #[serde(skip_serializing_if = "Option::is_none")]
1007    pub failure: Option<String>,
1008    pub initial_revision: u64,
1009    pub final_revision: u64,
1010}
1011
1012/// Evidence for a declarative step condition evaluated before dispatch.
1013#[derive(Debug, Clone, Serialize, Deserialize)]
1014#[serde(rename_all = "camelCase")]
1015pub struct WorkflowBranchDecision {
1016    pub step_id: String,
1017    pub predicate: VerificationPredicate,
1018    pub matched: bool,
1019}
1020
1021/// Deterministic state-transition trace for replay and debugging.
1022#[derive(Debug, Clone, Serialize, Deserialize)]
1023#[serde(rename_all = "camelCase")]
1024pub struct WorkflowTrace {
1025    #[serde(default = "default_workflow_trace_schema_version")]
1026    pub schema_version: u8,
1027    #[serde(default, skip_serializing_if = "Option::is_none")]
1028    pub run_id: Option<String>,
1029    pub events: Vec<WorkflowTraceEvent>,
1030    #[serde(default, skip_serializing_if = "Vec::is_empty")]
1031    pub branch_decisions: Vec<WorkflowBranchDecision>,
1032    #[serde(default, skip_serializing_if = "Vec::is_empty")]
1033    pub intent_resolutions: Vec<WorkflowIntentEvidence>,
1034}
1035
1036/// One ordered state transition in a workflow trace.
1037#[derive(Debug, Clone, Serialize, Deserialize)]
1038#[serde(rename_all = "camelCase")]
1039pub struct WorkflowTraceEvent {
1040    pub sequence: u64,
1041    pub step_id: String,
1042    pub state: WorkflowStepState,
1043    pub attempt: u32,
1044}
1045
1046impl WorkflowTrace {
1047    /// Build a stable trace from retained step histories.
1048    pub fn from_steps(steps: &[WorkflowStepRecord]) -> Self {
1049        let mut events = Vec::new();
1050        for step in steps {
1051            let mut attempt = 0u32;
1052            for state in &step.history {
1053                if events.len() == MAX_WORKFLOW_TRACE_EVENTS {
1054                    break;
1055                }
1056                if *state == WorkflowStepState::Preflight {
1057                    attempt = attempt.saturating_add(1);
1058                }
1059                events.push(WorkflowTraceEvent {
1060                    sequence: events.len() as u64,
1061                    step_id: step.id.clone(),
1062                    state: *state,
1063                    attempt,
1064                });
1065            }
1066        }
1067        let branch_decisions = steps
1068            .iter()
1069            .filter_map(|step| step.branch_decision.clone())
1070            .collect();
1071        let intent_resolutions = steps
1072            .iter()
1073            .filter_map(|step| step.intent_evidence.clone())
1074            .collect();
1075        Self {
1076            schema_version: WORKFLOW_TRACE_SCHEMA_VERSION,
1077            run_id: None,
1078            events,
1079            branch_decisions,
1080            intent_resolutions,
1081        }
1082    }
1083
1084    /// Validate sequence ordering and the trace event budget.
1085    pub fn validate(&self) -> Result<(), WorkflowValidationError> {
1086        if self.schema_version != WORKFLOW_TRACE_SCHEMA_VERSION {
1087            return Err(WorkflowValidationError::new(
1088                "trace.schemaVersion",
1089                format!(
1090                    "unsupported trace schema version {}; expected {}",
1091                    self.schema_version, WORKFLOW_TRACE_SCHEMA_VERSION
1092                ),
1093            ));
1094        }
1095        if self
1096            .run_id
1097            .as_ref()
1098            .is_some_and(|run_id| run_id.is_empty() || run_id.len() > 128)
1099        {
1100            return Err(WorkflowValidationError::new(
1101                "trace.runId",
1102                "run ID must contain 1 to 128 bytes",
1103            ));
1104        }
1105        if self.events.len() > MAX_WORKFLOW_TRACE_EVENTS {
1106            return Err(WorkflowValidationError::new(
1107                "trace.events",
1108                format!("must contain at most {MAX_WORKFLOW_TRACE_EVENTS} events"),
1109            ));
1110        }
1111        for (index, event) in self.events.iter().enumerate() {
1112            if event.sequence != index as u64 {
1113                return Err(WorkflowValidationError::new(
1114                    format!("trace.events[{index}].sequence"),
1115                    "sequence must be contiguous from zero",
1116                ));
1117            }
1118        }
1119        for (index, decision) in self.branch_decisions.iter().enumerate() {
1120            if decision.step_id.is_empty() {
1121                return Err(WorkflowValidationError::new(
1122                    format!("trace.branchDecisions[{index}].stepId"),
1123                    "step ID must not be empty",
1124                ));
1125            }
1126            validate_predicate(
1127                &decision.predicate,
1128                &format!("trace.branchDecisions[{index}].predicate"),
1129            )?;
1130        }
1131        Ok(())
1132    }
1133
1134    /// Replay the trace into step records without dispatching browser work.
1135    ///
1136    /// Replay is intentionally an inspection operation: it checks that the
1137    /// trace references the declared steps in order, follows the workflow
1138    /// state machine, and carries monotonic attempt numbers. A truncated
1139    /// trace is accepted as a valid prefix when it ends at the event budget.
1140    pub fn replay(
1141        &self,
1142        workflow: &WorkflowDefinition,
1143    ) -> Result<Vec<WorkflowStepRecord>, WorkflowValidationError> {
1144        workflow.validate()?;
1145        self.validate()?;
1146        let mut records: Vec<_> = workflow
1147            .steps
1148            .iter()
1149            .map(|step| WorkflowStepRecord::new(&step.id))
1150            .collect();
1151        let mut seen = vec![false; records.len()];
1152        let mut highest_step = 0usize;
1153        let mut started = false;
1154
1155        for (event_index, event) in self.events.iter().enumerate() {
1156            let step_index = workflow
1157                .steps
1158                .iter()
1159                .position(|step| step.id == event.step_id)
1160                .ok_or_else(|| {
1161                    WorkflowValidationError::new(
1162                        format!("trace.events[{event_index}].stepId"),
1163                        "step ID is not declared by the workflow",
1164                    )
1165                })?;
1166            if !started {
1167                if step_index != 0 {
1168                    return Err(WorkflowValidationError::new(
1169                        format!("trace.events[{event_index}].stepId"),
1170                        "trace must begin with the first declared step",
1171                    ));
1172                }
1173                started = true;
1174            }
1175            if step_index > highest_step.saturating_add(1) {
1176                return Err(WorkflowValidationError::new(
1177                    format!("trace.events[{event_index}].stepId"),
1178                    "trace skips a declared step",
1179                ));
1180            }
1181            if step_index < highest_step {
1182                return Err(WorkflowValidationError::new(
1183                    format!("trace.events[{event_index}].stepId"),
1184                    "trace returns to an earlier step",
1185                ));
1186            }
1187            highest_step = highest_step.max(step_index);
1188
1189            let record = &mut records[step_index];
1190            if !seen[step_index] {
1191                if event.state != WorkflowStepState::Pending || event.attempt != 0 {
1192                    return Err(WorkflowValidationError::new(
1193                        format!("trace.events[{event_index}]"),
1194                        "each step must begin with pending at attempt zero",
1195                    ));
1196                }
1197                seen[step_index] = true;
1198                continue;
1199            }
1200            if event.state == WorkflowStepState::Preflight {
1201                if event.attempt != record.attempts.saturating_add(1) {
1202                    return Err(WorkflowValidationError::new(
1203                        format!("trace.events[{event_index}].attempt"),
1204                        "preflight attempt must increment by one",
1205                    ));
1206                }
1207                record.attempts = event.attempt;
1208            } else if event.attempt != record.attempts {
1209                return Err(WorkflowValidationError::new(
1210                    format!("trace.events[{event_index}].attempt"),
1211                    "event attempt does not match the current step attempt",
1212                ));
1213            }
1214            record.transition(event.state).map_err(|reason| {
1215                WorkflowValidationError::new(format!("trace.events[{event_index}].state"), reason)
1216            })?;
1217        }
1218        for decision in &self.branch_decisions {
1219            let index = workflow
1220                .steps
1221                .iter()
1222                .position(|step| step.id == decision.step_id)
1223                .ok_or_else(|| {
1224                    WorkflowValidationError::new(
1225                        "trace.branchDecisions",
1226                        "branch decision references an undeclared step",
1227                    )
1228                })?;
1229            if workflow.steps[index].when.as_ref() != Some(&decision.predicate) {
1230                return Err(WorkflowValidationError::new(
1231                    "trace.branchDecisions",
1232                    "branch decision predicate does not match the workflow",
1233                ));
1234            }
1235            records[index].branch_decision = Some(decision.clone());
1236        }
1237        Ok(records)
1238    }
1239}
1240
1241fn default_workflow_trace_schema_version() -> u8 {
1242    WORKFLOW_TRACE_SCHEMA_VERSION
1243}
1244
1245impl WorkflowRunResult {
1246    fn failed(
1247        workflow: &WorkflowDefinition,
1248        run_id: String,
1249        steps: Vec<WorkflowStepRecord>,
1250        failed_step: Option<String>,
1251        failure: impl Into<String>,
1252        initial_revision: u64,
1253        final_revision: u64,
1254    ) -> Self {
1255        let mut trace = WorkflowTrace::from_steps(&steps);
1256        trace.run_id = Some(run_id.clone());
1257        Self {
1258            run_id,
1259            name: workflow.name.clone(),
1260            workflow_version: workflow.workflow_version.clone(),
1261            status: WorkflowRunStatus::Failed,
1262            steps,
1263            trace,
1264            outputs: BTreeMap::new(),
1265            terminal_proof: None,
1266            failed_step,
1267            failure: Some(bound_workflow_text(&failure.into(), 512)),
1268            initial_revision,
1269            final_revision,
1270        }
1271    }
1272
1273    fn budget_exhausted(
1274        workflow: &WorkflowDefinition,
1275        run_id: String,
1276        steps: Vec<WorkflowStepRecord>,
1277        failed_step: Option<String>,
1278        reason: impl Into<String>,
1279        initial_revision: u64,
1280        final_revision: u64,
1281    ) -> Self {
1282        let mut trace = WorkflowTrace::from_steps(&steps);
1283        trace.run_id = Some(run_id.clone());
1284        Self {
1285            run_id,
1286            name: workflow.name.clone(),
1287            workflow_version: workflow.workflow_version.clone(),
1288            status: WorkflowRunStatus::BudgetExhausted,
1289            steps,
1290            trace,
1291            outputs: BTreeMap::new(),
1292            terminal_proof: None,
1293            failed_step,
1294            failure: Some(bound_workflow_text(&reason.into(), 512)),
1295            initial_revision,
1296            final_revision,
1297        }
1298    }
1299
1300    fn resume_required(
1301        workflow: &WorkflowDefinition,
1302        run_id: String,
1303        steps: Vec<WorkflowStepRecord>,
1304        failed_step: Option<String>,
1305        reason: impl Into<String>,
1306        initial_revision: u64,
1307        final_revision: u64,
1308    ) -> Self {
1309        let mut trace = WorkflowTrace::from_steps(&steps);
1310        trace.run_id = Some(run_id.clone());
1311        Self {
1312            run_id,
1313            name: workflow.name.clone(),
1314            workflow_version: workflow.workflow_version.clone(),
1315            status: WorkflowRunStatus::ResumeRequired,
1316            steps,
1317            trace,
1318            outputs: BTreeMap::new(),
1319            terminal_proof: None,
1320            failed_step,
1321            failure: Some(bound_workflow_text(&reason.into(), 512)),
1322            initial_revision,
1323            final_revision,
1324        }
1325    }
1326}
1327
1328impl super::BrowserSession {
1329    async fn execute_workflow_intent_step(
1330        &self,
1331        intent: &WorkflowIntentStep,
1332        expected_revision: u64,
1333    ) -> BrowserResult<(super::types::BatchOutcome, WorkflowIntentEvidence)> {
1334        let mut execution = intent.execution_request("workflow.intent")?;
1335        execution.request.expected_revision = Some(expected_revision);
1336        let resolution = self.resolve_intent(&execution.request).await?;
1337        let candidate_id = resolution.selected_candidate.clone().ok_or_else(|| {
1338            format!(
1339                "semantic workflow intent was not uniquely authorized: resolution={}, policy={}",
1340                serde_json::to_string(&resolution.resolution).unwrap_or_else(|_| "unknown".into()),
1341                serde_json::to_string(&resolution.policy_decision)
1342                    .unwrap_or_else(|_| "unknown".into())
1343            )
1344        })?;
1345        execution.candidate_id = candidate_id;
1346        let result = self.execute_intent(&execution).await?;
1347        let accepted_candidate = result
1348            .resolution
1349            .candidates
1350            .iter()
1351            .find(|candidate| candidate.id == result.candidate_id)
1352            .ok_or("executed intent result omitted its accepted candidate")?;
1353        let evidence = WorkflowIntentEvidence {
1354            resolution_id: result.resolution_id.clone(),
1355            candidate_id: result.candidate_id.clone(),
1356            revision: result
1357                .resolution
1358                .revision
1359                .ok_or("executed intent result omitted revision")?,
1360            resolution: result.resolution.resolution,
1361            policy_decision: result.resolution.policy_decision,
1362            confidence: accepted_candidate.confidence,
1363            fingerprint: accepted_candidate.fingerprint.clone(),
1364        };
1365        let Some(action) = result.action else {
1366            return Err(result
1367                .reason
1368                .unwrap_or_else(|| "semantic workflow intent was not executed".into())
1369                .into());
1370        };
1371        let execution_id = action.execution_id.clone();
1372        Ok((
1373            super::types::BatchOutcome {
1374                mode: BatchMode::Fixed,
1375                initial_revision: expected_revision,
1376                final_revision: current_revision(self),
1377                steps: vec![super::types::BatchStepOutcome::Success {
1378                    index: 0,
1379                    action: "intent".into(),
1380                    response_bytes: serde_json::to_string(&action).ok().map(|value| value.len()),
1381                    execution_id: Some(execution_id),
1382                }],
1383                completed: 1,
1384                failed: 0,
1385                total: 1,
1386                success: true,
1387            },
1388            evidence,
1389        ))
1390    }
1391
1392    /// Execute a validated workflow linearly through the existing batch
1393    /// policy and action runtime, retaining bounded per-step evidence.
1394    pub async fn run_workflow(
1395        &self,
1396        workflow: &WorkflowDefinition,
1397        inputs: &BTreeMap<String, Value>,
1398    ) -> BrowserResult<WorkflowRunResult> {
1399        workflow.validate()?;
1400        workflow.validate_inputs(inputs)?;
1401        let resolved_steps = workflow.resolve_actions(inputs)?;
1402        let initial_revision = self
1403            .page_revision
1404            .load(std::sync::atomic::Ordering::Relaxed);
1405        let run_id = format!("run_{}", self.next_execution_id());
1406        let started = Instant::now();
1407        let duration_budget = Duration::from_millis(workflow.budgets.max_duration_ms);
1408        let mut executed_steps = 0u32;
1409        let mut records: Vec<_> = resolved_steps
1410            .iter()
1411            .map(|step| WorkflowStepRecord::new(&step.id))
1412            .collect();
1413
1414        for predicate in &workflow.preconditions {
1415            if workflow_budget_expired(started, duration_budget) {
1416                skip_remaining(&mut records, 0);
1417                return Ok(WorkflowRunResult::budget_exhausted(
1418                    workflow,
1419                    run_id.clone(),
1420                    records,
1421                    None,
1422                    "workflow maxDurationMs budget exhausted before a precondition",
1423                    initial_revision,
1424                    current_revision(self),
1425                ));
1426            }
1427            if let Err(error) = self
1428                .verify(
1429                    predicate.clone(),
1430                    workflow_budget_remaining(started, duration_budget),
1431                )
1432                .await
1433            {
1434                skip_remaining(&mut records, 0);
1435                if workflow_budget_expired(started, duration_budget) {
1436                    return Ok(WorkflowRunResult::budget_exhausted(
1437                        workflow,
1438                        run_id.clone(),
1439                        records,
1440                        None,
1441                        "workflow maxDurationMs budget exhausted while checking a precondition",
1442                        initial_revision,
1443                        current_revision(self),
1444                    ));
1445                }
1446                return Ok(WorkflowRunResult::failed(
1447                    workflow,
1448                    run_id.clone(),
1449                    records,
1450                    None,
1451                    format!("workflow precondition failed: {error}"),
1452                    initial_revision,
1453                    current_revision(self),
1454                ));
1455            }
1456        }
1457
1458        for (index, step) in resolved_steps.iter().enumerate() {
1459            for repetition in 0..step.repeat {
1460                if executed_steps >= workflow.budgets.max_steps
1461                    || workflow_budget_expired(started, duration_budget)
1462                {
1463                    skip_remaining(&mut records, index);
1464                    return Ok(WorkflowRunResult::budget_exhausted(
1465                        workflow,
1466                        run_id.clone(),
1467                        records,
1468                        Some(step.id.clone()),
1469                        if executed_steps >= workflow.budgets.max_steps {
1470                            "workflow maxSteps budget exhausted before dispatch"
1471                        } else {
1472                            "workflow maxDurationMs budget exhausted before dispatch"
1473                        },
1474                        initial_revision,
1475                        current_revision(self),
1476                    ));
1477                }
1478                executed_steps = executed_steps.saturating_add(1);
1479                if repetition > 0 {
1480                    let record = &mut records[index];
1481                    let _ = record.transition(WorkflowStepState::Ready);
1482                }
1483                if let Some(predicate) = &step.when {
1484                    {
1485                        let record = &mut records[index];
1486                        if record.state == WorkflowStepState::Pending {
1487                            let _ = record.transition(WorkflowStepState::Ready);
1488                        }
1489                    }
1490                    match self.evaluate_predicate_once(predicate).await {
1491                        Ok((matched, _state)) => {
1492                            let record = &mut records[index];
1493                            record.branch_decision = Some(WorkflowBranchDecision {
1494                                step_id: step.id.clone(),
1495                                predicate: predicate.clone(),
1496                                matched,
1497                            });
1498                            if !matched {
1499                                let _ = record.transition(WorkflowStepState::Skipped);
1500                                break;
1501                            }
1502                        }
1503                        Err(error) => {
1504                            let message = bound_workflow_text(&error.to_string(), 512);
1505                            let record = &mut records[index];
1506                            let _ = record.transition(WorkflowStepState::Preflight);
1507                            record.fail(WorkflowStepState::FailedBeforeDispatch, &message);
1508                            skip_remaining(&mut records, index + 1);
1509                            if workflow_budget_expired(started, duration_budget) {
1510                                return Ok(WorkflowRunResult::budget_exhausted(
1511                                    workflow,
1512                                    run_id.clone(),
1513                                    records,
1514                                    Some(step.id.clone()),
1515                                    "workflow maxDurationMs budget exhausted while evaluating a branch",
1516                                    initial_revision,
1517                                    current_revision(self),
1518                                ));
1519                            }
1520                            return Ok(WorkflowRunResult::failed(
1521                                workflow,
1522                                run_id.clone(),
1523                                records,
1524                                Some(step.id.clone()),
1525                                message,
1526                                initial_revision,
1527                                current_revision(self),
1528                            ));
1529                        }
1530                    }
1531                }
1532                let mut attempt_number = 0;
1533                let mut effect_marker_completed = false;
1534                let mut intent_evidence = None;
1535                let outcome = loop {
1536                    let attempt_revision = current_revision(self);
1537                    {
1538                        let record = &mut records[index];
1539                        if record.previous_revision.is_none() {
1540                            record.previous_revision = Some(attempt_revision);
1541                        }
1542                        record.retry_safe = step.transaction.permits_pre_dispatch_retry();
1543                        if record.state == WorkflowStepState::Pending {
1544                            let _ = record.transition(WorkflowStepState::Ready);
1545                        }
1546                        let _ = record.transition(WorkflowStepState::Preflight);
1547                        let _ = record.transition(WorkflowStepState::Resolving);
1548                        attempt_number += 1;
1549                        record.attempts = record.attempts.saturating_add(1);
1550                    }
1551
1552                    let outcome = if let Some(intent) = &step.intent {
1553                        match self
1554                            .execute_workflow_intent_step(intent, attempt_revision)
1555                            .await
1556                        {
1557                            Ok((outcome, evidence)) => {
1558                                intent_evidence = Some(evidence);
1559                                Ok(outcome)
1560                            }
1561                            Err(error) => Err(error),
1562                        }
1563                    } else {
1564                        match workflow_target(&step.action) {
1565                            Some(target) => match self.resolve_element(target).await {
1566                                Ok(_) => {
1567                                    self.run_batch_with_mode(
1568                                        std::slice::from_ref(&step.action),
1569                                        false,
1570                                        BatchMode::Unguarded,
1571                                        None,
1572                                    )
1573                                    .await
1574                                }
1575                                Err(error) => Err(error),
1576                            },
1577                            None => {
1578                                self.run_batch_with_mode(
1579                                    std::slice::from_ref(&step.action),
1580                                    false,
1581                                    BatchMode::Unguarded,
1582                                    None,
1583                                )
1584                                .await
1585                            }
1586                        }
1587                    };
1588                    match outcome {
1589                        Ok(outcome) => break outcome,
1590                        Err(error) => {
1591                            let message = bound_workflow_text(&error.to_string(), 512);
1592                            let retry = can_retry_before_dispatch(
1593                                step.transaction,
1594                                false,
1595                                attempt_number,
1596                                step.max_retries,
1597                            );
1598                            {
1599                                let record = &mut records[index];
1600                                record.current_revision = Some(current_revision(self));
1601                                let _ = record.transition(WorkflowStepState::NotDispatched);
1602                                record.fail(WorkflowStepState::FailedBeforeDispatch, &message);
1603                                if retry {
1604                                    let _ = record.transition(WorkflowStepState::Ready);
1605                                }
1606                            }
1607                            if retry {
1608                                let marker_matches = match &step.before_retry {
1609                                    Some(predicate) => {
1610                                        match self.evaluate_predicate_once(predicate).await {
1611                                            Ok((matched, _)) => Some(matched),
1612                                            Err(error) => {
1613                                                let message = bound_workflow_text(
1614                                                    &format!(
1615                                                        "effect marker could not be evaluated: {error}"
1616                                                    ),
1617                                                    512,
1618                                                );
1619                                                let record = &mut records[index];
1620                                                let _ =
1621                                                    record.transition(WorkflowStepState::Preflight);
1622                                                record.fail(
1623                                                    WorkflowStepState::FailedBeforeDispatch,
1624                                                    &message,
1625                                                );
1626                                                skip_remaining(&mut records, index + 1);
1627                                                return Ok(WorkflowRunResult::failed(
1628                                                    workflow,
1629                                                    run_id.clone(),
1630                                                    records,
1631                                                    Some(step.id.clone()),
1632                                                    message,
1633                                                    initial_revision,
1634                                                    current_revision(self),
1635                                                ));
1636                                            }
1637                                        }
1638                                    }
1639                                    None => None,
1640                                };
1641                                if marker_matches == Some(true) {
1642                                    effect_marker_completed = true;
1643                                    break super::types::BatchOutcome {
1644                                        mode: BatchMode::Unguarded,
1645                                        initial_revision: current_revision(self),
1646                                        final_revision: current_revision(self),
1647                                        steps: Vec::new(),
1648                                        completed: 0,
1649                                        failed: 0,
1650                                        total: 0,
1651                                        success: true,
1652                                    };
1653                                }
1654                                continue;
1655                            }
1656                            let message = records[index]
1657                                .error
1658                                .clone()
1659                                .unwrap_or_else(|| "workflow step failed".into());
1660                            skip_remaining(&mut records, index + 1);
1661                            return Ok(WorkflowRunResult::failed(
1662                                workflow,
1663                                run_id.clone(),
1664                                records,
1665                                Some(step.id.clone()),
1666                                message,
1667                                initial_revision,
1668                                current_revision(self),
1669                            ));
1670                        }
1671                    }
1672                };
1673
1674                if effect_marker_completed {
1675                    commit_workflow_effect_marker(&mut records[index], current_revision(self));
1676                    continue;
1677                }
1678
1679                if !outcome.success {
1680                    let message = outcome
1681                        .steps
1682                        .last()
1683                        .and_then(|step| match step {
1684                            super::types::BatchStepOutcome::Error { message, .. } => {
1685                                Some(message.as_str())
1686                            }
1687                            super::types::BatchStepOutcome::Success { .. } => None,
1688                        })
1689                        .unwrap_or("workflow action failed after dispatch");
1690                    let message = {
1691                        let record = &mut records[index];
1692                        record.dispatch_acknowledged = true;
1693                        record.current_revision = Some(current_revision(self));
1694                        let _ = record.transition(WorkflowStepState::Dispatched);
1695                        record.fail(WorkflowStepState::Indeterminate, message);
1696                        record
1697                            .error
1698                            .clone()
1699                            .unwrap_or_else(|| "workflow step failed".into())
1700                    };
1701                    skip_remaining(&mut records, index + 1);
1702                    return Ok(WorkflowRunResult::resume_required(
1703                        workflow,
1704                        run_id.clone(),
1705                        records,
1706                        Some(step.id.clone()),
1707                        message,
1708                        initial_revision,
1709                        current_revision(self),
1710                    ));
1711                }
1712
1713                {
1714                    records[index].intent_evidence = intent_evidence;
1715                    for execution_id in outcome.steps.iter().filter_map(|step| match step {
1716                        super::types::BatchStepOutcome::Success {
1717                            execution_id: Some(execution_id),
1718                            ..
1719                        } => Some(execution_id.clone()),
1720                        _ => None,
1721                    }) {
1722                        records[index].execution_ids.push(execution_id);
1723                    }
1724                    let record = &mut records[index];
1725                    record.dispatch_acknowledged = true;
1726                    record.effect_observed = true;
1727                    record.current_revision = Some(current_revision(self));
1728                    let _ = record.transition(WorkflowStepState::Dispatched);
1729                    let _ = record.transition(WorkflowStepState::EffectObserved);
1730                }
1731                if let Some(predicate) = &step.expect
1732                    && let Err(error) = self
1733                        .verify(
1734                            predicate.clone(),
1735                            workflow_budget_remaining(started, duration_budget),
1736                        )
1737                        .await
1738                {
1739                    let message = {
1740                        let record = &mut records[index];
1741                        record.current_revision = Some(current_revision(self));
1742                        record.fail(WorkflowStepState::FailedAfterDispatch, &error.to_string());
1743                        record
1744                            .error
1745                            .clone()
1746                            .unwrap_or_else(|| "workflow verification failed".into())
1747                    };
1748                    skip_remaining(&mut records, index + 1);
1749                    if workflow_budget_expired(started, duration_budget) {
1750                        return Ok(WorkflowRunResult::budget_exhausted(
1751                            workflow,
1752                            run_id.clone(),
1753                            records,
1754                            Some(step.id.clone()),
1755                            "workflow maxDurationMs budget exhausted while verifying a step",
1756                            initial_revision,
1757                            current_revision(self),
1758                        ));
1759                    }
1760                    return Ok(WorkflowRunResult::resume_required(
1761                        workflow,
1762                        run_id.clone(),
1763                        records,
1764                        Some(step.id.clone()),
1765                        message,
1766                        initial_revision,
1767                        current_revision(self),
1768                    ));
1769                }
1770                let record = &mut records[index];
1771                record.postcondition_verified = true;
1772                let _ = record.transition(WorkflowStepState::Verified);
1773                let _ = record.transition(WorkflowStepState::OutputsExtracted);
1774                let _ = record.transition(WorkflowStepState::Committed);
1775            }
1776        }
1777
1778        if workflow_budget_expired(started, duration_budget) {
1779            return Ok(WorkflowRunResult::budget_exhausted(
1780                workflow,
1781                run_id.clone(),
1782                records,
1783                None,
1784                "workflow maxDurationMs budget exhausted before terminal verification",
1785                initial_revision,
1786                current_revision(self),
1787            ));
1788        }
1789        let terminal_proof = match self
1790            .verify(
1791                workflow.terminal_condition.clone(),
1792                workflow_budget_remaining(started, duration_budget),
1793            )
1794            .await
1795        {
1796            Ok(outcome) => WorkflowTerminalProof {
1797                predicate: outcome.predicate,
1798                revision: current_revision(self),
1799                state: outcome.state,
1800            },
1801            Err(error) => {
1802                if workflow_budget_expired(started, duration_budget) {
1803                    return Ok(WorkflowRunResult::budget_exhausted(
1804                        workflow,
1805                        run_id.clone(),
1806                        records,
1807                        None,
1808                        "workflow maxDurationMs budget exhausted while verifying the terminal condition",
1809                        initial_revision,
1810                        current_revision(self),
1811                    ));
1812                }
1813                return Ok(WorkflowRunResult::failed(
1814                    workflow,
1815                    run_id.clone(),
1816                    records,
1817                    None,
1818                    format!("workflow terminal condition was not proven: {error}"),
1819                    initial_revision,
1820                    current_revision(self),
1821                ));
1822            }
1823        };
1824        if workflow_budget_expired(started, duration_budget) {
1825            return Ok(WorkflowRunResult::budget_exhausted(
1826                workflow,
1827                run_id.clone(),
1828                records,
1829                None,
1830                "workflow maxDurationMs budget exhausted before output extraction",
1831                initial_revision,
1832                current_revision(self),
1833            ));
1834        }
1835        let outputs = match extract_workflow_outputs(self, workflow).await {
1836            Ok(outputs) => outputs,
1837            Err(error) => {
1838                if workflow_budget_expired(started, duration_budget) {
1839                    return Ok(WorkflowRunResult::budget_exhausted(
1840                        workflow,
1841                        run_id.clone(),
1842                        records,
1843                        None,
1844                        "workflow maxDurationMs budget exhausted while extracting outputs",
1845                        initial_revision,
1846                        current_revision(self),
1847                    ));
1848                }
1849                return Ok(WorkflowRunResult::failed(
1850                    workflow,
1851                    run_id.clone(),
1852                    records,
1853                    None,
1854                    format!("workflow output extraction failed: {error}"),
1855                    initial_revision,
1856                    current_revision(self),
1857                ));
1858            }
1859        };
1860        if workflow_budget_expired(started, duration_budget) {
1861            return Ok(WorkflowRunResult::budget_exhausted(
1862                workflow,
1863                run_id.clone(),
1864                records,
1865                None,
1866                "workflow maxDurationMs budget exhausted after output extraction",
1867                initial_revision,
1868                current_revision(self),
1869            ));
1870        }
1871
1872        let mut trace = WorkflowTrace::from_steps(&records);
1873        trace.run_id = Some(run_id.clone());
1874        Ok(WorkflowRunResult {
1875            run_id,
1876            name: workflow.name.clone(),
1877            workflow_version: workflow.workflow_version.clone(),
1878            status: WorkflowRunStatus::Completed,
1879            steps: records,
1880            trace,
1881            outputs,
1882            terminal_proof: Some(terminal_proof),
1883            failed_step: None,
1884            failure: None,
1885            initial_revision,
1886            final_revision: current_revision(self),
1887        })
1888    }
1889}
1890
1891impl super::BrowserSession {
1892    /// Export a deterministic, redacted workflow checkpoint bounded to 8 KiB.
1893    pub async fn export_workflow_checkpoint(
1894        &self,
1895        workflow: &WorkflowDefinition,
1896        result: &WorkflowRunResult,
1897    ) -> BrowserResult<WorkflowCheckpoint> {
1898        workflow.validate()?;
1899        if result.name != workflow.name
1900            || result.workflow_version != workflow.workflow_version
1901            || result.steps.len() != workflow.steps.len()
1902        {
1903            return Err(WorkflowResumeError::CheckpointShape(
1904                "run result does not belong to workflow definition".into(),
1905            )
1906            .into());
1907        }
1908        let page = self.page_info().await?;
1909        let next_step_index = result
1910            .steps
1911            .iter()
1912            .position(|step| step.state != WorkflowStepState::Committed)
1913            .unwrap_or(result.steps.len());
1914        let checkpoint = WorkflowCheckpoint {
1915            schema_version: WORKFLOW_CHECKPOINT_SCHEMA_VERSION,
1916            run_id: result.run_id.clone(),
1917            workflow_name: workflow.name.clone(),
1918            workflow_version: workflow.workflow_version.clone(),
1919            definition_hash: workflow_definition_hash(workflow)?,
1920            status: result.status,
1921            next_step_index,
1922            steps: result
1923                .steps
1924                .iter()
1925                .map(|step| WorkflowCheckpointStep {
1926                    id: step.id.clone(),
1927                    state: step.state,
1928                    attempts: step.attempts,
1929                    history: step.history.clone(),
1930                    execution_ids: step.execution_ids.clone(),
1931                    dispatch_acknowledged: step.dispatch_acknowledged,
1932                    effect_observed: step.effect_observed,
1933                    postcondition_verified: step.postcondition_verified,
1934                    retry_safe: step.retry_safe,
1935                    previous_revision: step.previous_revision,
1936                    current_revision: step.current_revision,
1937                    branch_decision: step.branch_decision.clone(),
1938                    intent_evidence: step.intent_evidence.clone(),
1939                })
1940                .collect(),
1941            page: WorkflowCheckpointPage {
1942                target_id: bound_workflow_text(&page.target_id, 256),
1943                frame_id: bound_workflow_text(&page.frame_id, 256),
1944                url: bound_workflow_text(&page.url, 1_024),
1945                title: bound_workflow_text(&page.title, 1_024),
1946                revision: current_revision(self),
1947            },
1948        };
1949        checkpoint.validate_size()?;
1950        Ok(checkpoint)
1951    }
1952
1953    /// Parse and validate a workflow checkpoint without contacting Chrome.
1954    pub fn parse_workflow_checkpoint(
1955        input: &str,
1956    ) -> Result<WorkflowCheckpoint, WorkflowResumeError> {
1957        let checkpoint: WorkflowCheckpoint = serde_json::from_str(input)
1958            .map_err(|error| WorkflowResumeError::CheckpointShape(error.to_string()))?;
1959        checkpoint.validate_size()?;
1960        Ok(checkpoint)
1961    }
1962
1963    /// Reconcile a checkpoint with the current route and return the only safe
1964    /// next step. This method never dispatches a browser action.
1965    pub async fn reconcile_workflow_checkpoint(
1966        &self,
1967        workflow: &WorkflowDefinition,
1968        checkpoint: &WorkflowCheckpoint,
1969    ) -> BrowserResult<WorkflowResumePlan> {
1970        workflow.validate()?;
1971        checkpoint.validate_size()?;
1972        if checkpoint.schema_version != WORKFLOW_CHECKPOINT_SCHEMA_VERSION {
1973            return Err(WorkflowResumeError::SchemaVersionMismatch {
1974                expected: WORKFLOW_CHECKPOINT_SCHEMA_VERSION,
1975                found: checkpoint.schema_version,
1976            }
1977            .into());
1978        }
1979        if checkpoint.workflow_name != workflow.name
1980            || checkpoint.workflow_version != workflow.workflow_version
1981            || checkpoint.definition_hash != workflow_definition_hash(workflow)?
1982        {
1983            return Err(WorkflowResumeError::DefinitionMismatch.into());
1984        }
1985        if checkpoint.steps.len() != workflow.steps.len()
1986            || checkpoint
1987                .steps
1988                .iter()
1989                .zip(&workflow.steps)
1990                .any(|(checkpoint_step, step)| checkpoint_step.id != step.id)
1991        {
1992            return Err(WorkflowResumeError::CheckpointShape(
1993                "checkpoint steps do not match workflow steps".into(),
1994            )
1995            .into());
1996        }
1997
1998        let page = self.page_info().await?;
1999        if checkpoint.page.target_id != page.target_id
2000            || checkpoint.page.frame_id != page.frame_id
2001            || checkpoint.page.url != bound_workflow_text(&page.url, 1_024)
2002            || checkpoint.page.title != bound_workflow_text(&page.title, 1_024)
2003        {
2004            return Err(WorkflowResumeError::RouteChanged.into());
2005        }
2006
2007        let mut next_step_index = checkpoint.steps.len();
2008        for (index, checkpoint_step) in checkpoint.steps.iter().enumerate() {
2009            if checkpoint_step.state != WorkflowStepState::Committed {
2010                next_step_index = index;
2011                break;
2012            }
2013        }
2014        if checkpoint.next_step_index != next_step_index {
2015            return Err(WorkflowResumeError::CheckpointShape(
2016                "nextStepIndex does not match step states".into(),
2017            )
2018            .into());
2019        }
2020        if checkpoint.status == WorkflowRunStatus::Completed
2021            && next_step_index != workflow.steps.len()
2022        {
2023            return Err(WorkflowResumeError::CheckpointShape(
2024                "completed checkpoint has uncommitted steps".into(),
2025            )
2026            .into());
2027        }
2028        if next_step_index < workflow.steps.len() {
2029            let step = &workflow.steps[next_step_index];
2030            let state = checkpoint.steps[next_step_index].state;
2031            if state == WorkflowStepState::FailedBeforeDispatch
2032                && !step.transaction.permits_pre_dispatch_retry()
2033            {
2034                return Err(WorkflowResumeError::InvalidState {
2035                    step_id: step.id.clone(),
2036                    state,
2037                }
2038                .into());
2039            }
2040            if !matches!(
2041                state,
2042                WorkflowStepState::Pending | WorkflowStepState::FailedBeforeDispatch
2043            ) {
2044                return Err(WorkflowResumeError::InvalidState {
2045                    step_id: step.id.clone(),
2046                    state,
2047                }
2048                .into());
2049            }
2050            for checkpoint_step in checkpoint.steps.iter().skip(next_step_index + 1) {
2051                if checkpoint_step.state != WorkflowStepState::Skipped {
2052                    return Err(WorkflowResumeError::InvalidState {
2053                        step_id: checkpoint_step.id.clone(),
2054                        state: checkpoint_step.state,
2055                    }
2056                    .into());
2057                }
2058            }
2059        }
2060
2061        Ok(WorkflowResumePlan {
2062            workflow_name: workflow.name.clone(),
2063            workflow_version: workflow.workflow_version.clone(),
2064            next_step_index,
2065            current_revision: current_revision(self),
2066            reconciled: true,
2067        })
2068    }
2069
2070    /// Reconcile a checkpoint and execute only its safe pending suffix.
2071    /// Previously committed steps are never re-dispatched by this method.
2072    pub async fn resume_workflow(
2073        &self,
2074        workflow: &WorkflowDefinition,
2075        inputs: &BTreeMap<String, Value>,
2076        checkpoint: &WorkflowCheckpoint,
2077    ) -> BrowserResult<WorkflowRunResult> {
2078        let plan = self
2079            .reconcile_workflow_checkpoint(workflow, checkpoint)
2080            .await?;
2081        if plan.next_step_index >= workflow.steps.len() {
2082            return Err(WorkflowResumeError::CheckpointShape(
2083                "workflow checkpoint is already complete".into(),
2084            )
2085            .into());
2086        }
2087        let mut suffix = workflow.clone();
2088        suffix.steps = workflow.steps[plan.next_step_index..].to_vec();
2089        let mut result = self.run_workflow(&suffix, inputs).await?;
2090        let mut prefix = checkpoint.steps[..plan.next_step_index]
2091            .iter()
2092            .map(checkpoint_step_to_record)
2093            .collect::<Vec<_>>();
2094        prefix.append(&mut result.steps);
2095        result.steps = prefix;
2096        result.trace = WorkflowTrace::from_steps(&result.steps);
2097        result.trace.run_id = Some(result.run_id.clone());
2098        Ok(result)
2099    }
2100}
2101
2102fn checkpoint_step_to_record(step: &WorkflowCheckpointStep) -> WorkflowStepRecord {
2103    let history = if step.history.is_empty() {
2104        vec![
2105            WorkflowStepState::Pending,
2106            WorkflowStepState::Ready,
2107            WorkflowStepState::Preflight,
2108            WorkflowStepState::Resolving,
2109            WorkflowStepState::Dispatched,
2110            WorkflowStepState::EffectObserved,
2111            WorkflowStepState::Verified,
2112            WorkflowStepState::OutputsExtracted,
2113            WorkflowStepState::Committed,
2114        ]
2115    } else {
2116        step.history.clone()
2117    };
2118    WorkflowStepRecord {
2119        id: step.id.clone(),
2120        state: step.state,
2121        history,
2122        attempts: step.attempts,
2123        execution_ids: step.execution_ids.clone(),
2124        dispatch_acknowledged: step.dispatch_acknowledged,
2125        effect_observed: step.effect_observed,
2126        postcondition_verified: step.postcondition_verified,
2127        retry_safe: step.retry_safe,
2128        previous_revision: step.previous_revision,
2129        current_revision: step.current_revision,
2130        branch_decision: step.branch_decision.clone(),
2131        intent_evidence: step.intent_evidence.clone(),
2132        error: None,
2133    }
2134}
2135
2136impl WorkflowCheckpoint {
2137    /// Serialize a checkpoint after enforcing its size and schema bounds.
2138    pub fn to_canonical_json(&self) -> Result<String, WorkflowResumeError> {
2139        self.validate_size()?;
2140        serde_json::to_string(self)
2141            .map_err(|error| WorkflowResumeError::CheckpointShape(error.to_string()))
2142    }
2143
2144    fn validate_size(&self) -> Result<(), WorkflowResumeError> {
2145        if self.schema_version != WORKFLOW_CHECKPOINT_SCHEMA_VERSION {
2146            return Err(WorkflowResumeError::SchemaVersionMismatch {
2147                expected: WORKFLOW_CHECKPOINT_SCHEMA_VERSION,
2148                found: self.schema_version,
2149            });
2150        }
2151        if self.steps.len() > MAX_STEPS {
2152            return Err(WorkflowResumeError::CheckpointShape(
2153                "checkpoint contains too many steps".into(),
2154            ));
2155        }
2156        if self.run_id.len() > 128 {
2157            return Err(WorkflowResumeError::CheckpointShape(
2158                "checkpoint runId is too long".into(),
2159            ));
2160        }
2161        for (index, step) in self.steps.iter().enumerate() {
2162            if step.history.len() > MAX_CHECKPOINT_HISTORY_STATES {
2163                return Err(WorkflowResumeError::CheckpointShape(format!(
2164                    "checkpoint step {index} contains too much state history"
2165                )));
2166            }
2167            if !step.history.is_empty() && step.history.last() != Some(&step.state) {
2168                return Err(WorkflowResumeError::CheckpointShape(format!(
2169                    "checkpoint step {index} history does not end at its state"
2170                )));
2171            }
2172            if step.execution_ids.len() > MAX_CHECKPOINT_EXECUTION_IDS
2173                || step.execution_ids.iter().any(|id| id.len() > 128)
2174            {
2175                return Err(WorkflowResumeError::CheckpointShape(format!(
2176                    "checkpoint step {index} contains too many execution IDs"
2177                )));
2178            }
2179        }
2180        let bytes = serde_json::to_vec(self)
2181            .map_err(|error| WorkflowResumeError::CheckpointShape(error.to_string()))?;
2182        if bytes.len() > MAX_WORKFLOW_CHECKPOINT_BYTES {
2183            return Err(WorkflowResumeError::CheckpointTooLarge);
2184        }
2185        Ok(())
2186    }
2187}
2188
2189fn workflow_definition_hash(
2190    workflow: &WorkflowDefinition,
2191) -> Result<String, WorkflowValidationError> {
2192    let canonical = workflow.to_canonical_json()?;
2193    let digest = Sha256::digest(canonical.as_bytes());
2194    Ok(digest.iter().map(|byte| format!("{byte:02x}")).collect())
2195}
2196
2197fn workflow_target(action: &BatchStep) -> Option<&str> {
2198    match action {
2199        BatchStep::Click { target }
2200        | BatchStep::Check { target }
2201        | BatchStep::Uncheck { target }
2202        | BatchStep::Select { target, .. }
2203        | BatchStep::Clear { target } => Some(target),
2204        BatchStep::Type { target, .. } => target.as_deref(),
2205        _ => None,
2206    }
2207}
2208
2209fn resolve_workflow_intent(
2210    intent: &WorkflowIntentStep,
2211    inputs: &BTreeMap<String, Value>,
2212    path: &str,
2213) -> Result<WorkflowIntentStep, WorkflowValidationError> {
2214    let resolve = |field: &str, value: &str, maximum: usize| {
2215        resolve_input_template(value, inputs, &format!("{path}.{field}"), maximum)
2216    };
2217    let resolve_optional = |field: &str, value: &Option<String>, maximum: usize| {
2218        value
2219            .as_deref()
2220            .map(|value| resolve(field, value, maximum))
2221            .transpose()
2222    };
2223    let mut resolved = intent.clone();
2224    resolved.purpose = resolve_optional("purpose", &intent.purpose, MAX_INTENT_PURPOSE_BYTES)?;
2225    resolved.intent = resolve_optional("intent", &intent.intent, MAX_INTENT_PURPOSE_BYTES * 2)?;
2226    resolved.value = resolve_optional("value", &intent.value, 4_096)?;
2227    resolved.scope.region_id = resolve_optional("scope.regionId", &intent.scope.region_id, 128)?;
2228    resolved.scope.form_label = resolve_optional("scope.formLabel", &intent.scope.form_label, 256)?;
2229    resolved.constraints.role = resolve_optional("constraints.role", &intent.constraints.role, 64)?;
2230    resolved.constraints.name = resolve_optional(
2231        "constraints.name",
2232        &intent.constraints.name,
2233        MAX_TARGET_BYTES,
2234    )?;
2235    resolved.constraints.name_contains = resolve_optional(
2236        "constraints.nameContains",
2237        &intent.constraints.name_contains,
2238        MAX_TARGET_BYTES,
2239    )?;
2240    resolved.constraints.exclude_text = intent
2241        .constraints
2242        .exclude_text
2243        .iter()
2244        .enumerate()
2245        .map(|(index, value)| {
2246            resolve_input_template(
2247                value,
2248                inputs,
2249                &format!("{path}.constraints.excludeText[{index}]"),
2250                MAX_TARGET_BYTES,
2251            )
2252        })
2253        .collect::<Result<_, _>>()?;
2254    resolved.validate(path)?;
2255    Ok(resolved)
2256}
2257
2258fn resolve_batch_step(
2259    action: &BatchStep,
2260    inputs: &BTreeMap<String, Value>,
2261    path: &str,
2262) -> Result<BatchStep, WorkflowValidationError> {
2263    let resolve = |field: &str, value: &str, maximum: usize| {
2264        resolve_input_template(value, inputs, &format!("{path}.{field}"), maximum)
2265    };
2266    Ok(match action {
2267        BatchStep::Navigate { url, timeout_ms } => BatchStep::Navigate {
2268            url: resolve("url", url, MAX_TARGET_BYTES)?,
2269            timeout_ms: *timeout_ms,
2270        },
2271        BatchStep::Click { target } => BatchStep::Click {
2272            target: resolve("target", target, MAX_TARGET_BYTES)?,
2273        },
2274        BatchStep::Type { text, target } => BatchStep::Type {
2275            text: resolve("text", text, MAX_TEXT_BYTES)?,
2276            target: target
2277                .as_deref()
2278                .map(|value| resolve("target", value, MAX_TARGET_BYTES))
2279                .transpose()?,
2280        },
2281        BatchStep::Check { target } => BatchStep::Check {
2282            target: resolve("target", target, MAX_TARGET_BYTES)?,
2283        },
2284        BatchStep::Uncheck { target } => BatchStep::Uncheck {
2285            target: resolve("target", target, MAX_TARGET_BYTES)?,
2286        },
2287        BatchStep::Select { target, value } => BatchStep::Select {
2288            target: resolve("target", target, MAX_TARGET_BYTES)?,
2289            value: resolve("value", value, MAX_TEXT_BYTES)?,
2290        },
2291        BatchStep::Clear { target } => BatchStep::Clear {
2292            target: resolve("target", target, MAX_TARGET_BYTES)?,
2293        },
2294        BatchStep::Wait {
2295            condition,
2296            timeout_ms,
2297        } => BatchStep::Wait {
2298            condition: resolve("condition", condition, MAX_WAIT_CONDITION_BYTES)?,
2299            timeout_ms: *timeout_ms,
2300        },
2301        BatchStep::Scroll { dx, dy } => BatchStep::Scroll { dx: *dx, dy: *dy },
2302        BatchStep::Observe {
2303            include_dom,
2304            include_screenshot,
2305            include_form_values,
2306        } => BatchStep::Observe {
2307            include_dom: *include_dom,
2308            include_screenshot: *include_screenshot,
2309            include_form_values: *include_form_values,
2310        },
2311        BatchStep::Screenshot => BatchStep::Screenshot,
2312        BatchStep::Evaluate { expression } => BatchStep::Evaluate {
2313            expression: resolve("expression", expression, MAX_TEXT_BYTES)?,
2314        },
2315        BatchStep::AcceptDialog => BatchStep::AcceptDialog,
2316        BatchStep::DismissDialog => BatchStep::DismissDialog,
2317    })
2318}
2319
2320fn resolve_input_template(
2321    value: &str,
2322    inputs: &BTreeMap<String, Value>,
2323    path: &str,
2324    maximum: usize,
2325) -> Result<String, WorkflowValidationError> {
2326    let marker = "${inputs.";
2327    let mut resolved = String::with_capacity(value.len());
2328    let mut remainder = value;
2329    while let Some(start) = remainder.find(marker) {
2330        resolved.push_str(&remainder[..start]);
2331        let placeholder = &remainder[start..];
2332        let end = placeholder.find('}').ok_or_else(|| {
2333            WorkflowValidationError::new(path, "input placeholder is missing a closing brace")
2334        })?;
2335        let name = &placeholder[marker.len()..end];
2336        if name.is_empty() || name.contains(['{', '}', '$']) {
2337            return Err(WorkflowValidationError::new(
2338                path,
2339                "input placeholder name is invalid",
2340            ));
2341        }
2342        let input = inputs.get(name).ok_or_else(|| {
2343            WorkflowValidationError::new(path, format!("input {name:?} is missing or not declared"))
2344        })?;
2345        let text = match input {
2346            Value::String(text) => text.clone(),
2347            Value::Bool(value) => value.to_string(),
2348            Value::Number(value) => value.to_string(),
2349            _ => {
2350                return Err(WorkflowValidationError::new(
2351                    path,
2352                    format!("input {name:?} cannot be inserted into an action string"),
2353                ));
2354            }
2355        };
2356        resolved.push_str(&text);
2357        remainder = &placeholder[end + 1..];
2358    }
2359    resolved.push_str(remainder);
2360    validate_bytes(path, &resolved, 0, maximum)?;
2361    Ok(resolved)
2362}
2363
2364async fn extract_workflow_outputs(
2365    session: &super::BrowserSession,
2366    workflow: &WorkflowDefinition,
2367) -> BrowserResult<BTreeMap<String, WorkflowOutput>> {
2368    let page = if workflow
2369        .outputs
2370        .values()
2371        .any(|output| !matches!(output.source, WorkflowOutputSource::VisibleText))
2372    {
2373        Some(session.page_info().await?)
2374    } else {
2375        None
2376    };
2377    let visible_text = if workflow
2378        .outputs
2379        .values()
2380        .any(|output| matches!(output.source, WorkflowOutputSource::VisibleText))
2381    {
2382        Some(session.text().await?)
2383    } else {
2384        None
2385    };
2386    let mut outputs = BTreeMap::new();
2387    let mut extracted_bytes = 0usize;
2388    let revision = current_revision(session);
2389    for (name, declaration) in &workflow.outputs {
2390        let text = match declaration.source {
2391            WorkflowOutputSource::PageUrl => page.as_ref().map(|page| page.url.as_str()),
2392            WorkflowOutputSource::PageTitle => page.as_ref().map(|page| page.title.as_str()),
2393            WorkflowOutputSource::VisibleText => visible_text.as_deref(),
2394        }
2395        .ok_or_else(|| format!("output {name:?} has no extraction source"))?;
2396        extracted_bytes = extracted_bytes.saturating_add(text.len());
2397        if extracted_bytes > workflow.budgets.max_extracted_bytes {
2398            return Err(format!(
2399                "outputs exceed maxExtractedBytes {}",
2400                workflow.budgets.max_extracted_bytes
2401            )
2402            .into());
2403        }
2404        let value = typed_output_value(name, declaration.value_type, text)?;
2405        let redacted = declaration.sensitive;
2406        outputs.insert(
2407            name.clone(),
2408            WorkflowOutput {
2409                value_type: declaration.value_type,
2410                value: if redacted { Value::Null } else { value },
2411                redacted,
2412                evidence: WorkflowOutputEvidence {
2413                    source: declaration.source,
2414                    revision,
2415                },
2416            },
2417        );
2418    }
2419    Ok(outputs)
2420}
2421
2422fn typed_output_value(
2423    name: &str,
2424    value_type: WorkflowValueType,
2425    text: &str,
2426) -> BrowserResult<Value> {
2427    let trimmed = text.trim();
2428    match value_type {
2429        WorkflowValueType::String => Ok(Value::String(text.to_string())),
2430        WorkflowValueType::Url => {
2431            if Url::parse(trimmed).is_err() {
2432                return Err(format!("output {name:?} is not a valid URL").into());
2433            }
2434            Ok(Value::String(text.to_string()))
2435        }
2436        WorkflowValueType::Integer => trimmed
2437            .parse::<i64>()
2438            .map(Value::from)
2439            .map_err(|_| format!("output {name:?} cannot be parsed as an integer").into()),
2440        WorkflowValueType::Number => {
2441            let number = trimmed
2442                .parse::<f64>()
2443                .map_err(|_| format!("output {name:?} cannot be parsed as a number"))?;
2444            serde_json::Number::from_f64(number)
2445                .map(Value::Number)
2446                .ok_or_else(|| format!("output {name:?} is not a finite number").into())
2447        }
2448        WorkflowValueType::Boolean => match trimmed {
2449            "true" => Ok(Value::Bool(true)),
2450            "false" => Ok(Value::Bool(false)),
2451            _ => Err(format!("output {name:?} cannot be parsed as a boolean").into()),
2452        },
2453    }
2454}
2455
2456fn current_revision(session: &super::BrowserSession) -> u64 {
2457    session
2458        .page_revision
2459        .load(std::sync::atomic::Ordering::Relaxed)
2460}
2461
2462fn workflow_budget_expired(started: Instant, budget: Duration) -> bool {
2463    started.elapsed() >= budget
2464}
2465
2466fn workflow_budget_remaining(started: Instant, budget: Duration) -> Duration {
2467    budget.saturating_sub(started.elapsed())
2468}
2469
2470fn can_retry_before_dispatch(
2471    transaction: WorkflowTransactionClass,
2472    dispatch_observed: bool,
2473    attempt_number: u32,
2474    max_retries: u32,
2475) -> bool {
2476    !dispatch_observed && attempt_number <= max_retries && transaction.permits_pre_dispatch_retry()
2477}
2478
2479fn skip_remaining(records: &mut [WorkflowStepRecord], start: usize) {
2480    for record in records.iter_mut().skip(start) {
2481        if record.state == WorkflowStepState::Pending {
2482            let _ = record.transition(WorkflowStepState::Skipped);
2483        }
2484    }
2485}
2486
2487fn commit_workflow_effect_marker(record: &mut WorkflowStepRecord, revision: u64) {
2488    if record.state == WorkflowStepState::Ready {
2489        let _ = record.transition(WorkflowStepState::Preflight);
2490    }
2491    if record.state == WorkflowStepState::Preflight {
2492        let _ = record.transition(WorkflowStepState::EffectObserved);
2493    }
2494    record.effect_observed = true;
2495    record.postcondition_verified = true;
2496    record.current_revision = Some(revision);
2497    let _ = record.transition(WorkflowStepState::Verified);
2498    let _ = record.transition(WorkflowStepState::OutputsExtracted);
2499    let _ = record.transition(WorkflowStepState::Committed);
2500}
2501
2502fn state_name(state: WorkflowStepState) -> &'static str {
2503    match state {
2504        WorkflowStepState::Pending => "pending",
2505        WorkflowStepState::Ready => "ready",
2506        WorkflowStepState::Preflight => "preflight",
2507        WorkflowStepState::Resolving => "resolving",
2508        WorkflowStepState::NotDispatched => "not_dispatched",
2509        WorkflowStepState::Dispatched => "dispatched",
2510        WorkflowStepState::EffectObserved => "effect_observed",
2511        WorkflowStepState::Verified => "verified",
2512        WorkflowStepState::OutputsExtracted => "outputs_extracted",
2513        WorkflowStepState::Committed => "committed",
2514        WorkflowStepState::FailedBeforeDispatch => "failed_before_dispatch",
2515        WorkflowStepState::FailedAfterDispatch => "failed_after_dispatch",
2516        WorkflowStepState::Indeterminate => "indeterminate",
2517        WorkflowStepState::Skipped => "skipped",
2518    }
2519}
2520
2521fn bound_workflow_text(value: &str, max_bytes: usize) -> String {
2522    if value.len() <= max_bytes {
2523        return value.to_string();
2524    }
2525    let mut end = max_bytes.saturating_sub(13);
2526    while end > 0 && !value.is_char_boundary(end) {
2527        end -= 1;
2528    }
2529    format!("{}\n[truncated]", &value[..end])
2530}
2531
2532/// A declared output captured after terminal verification.
2533#[derive(Debug, Clone, Serialize, Deserialize)]
2534#[serde(rename_all = "camelCase")]
2535pub struct WorkflowOutputDeclaration {
2536    #[serde(alias = "type")]
2537    pub value_type: WorkflowValueType,
2538    pub source: WorkflowOutputSource,
2539    #[serde(default)]
2540    pub required: bool,
2541    #[serde(default)]
2542    pub sensitive: bool,
2543}
2544
2545impl WorkflowOutputDeclaration {
2546    fn validate(&self, path: &str) -> Result<(), WorkflowValidationError> {
2547        let compatible = match self.source {
2548            WorkflowOutputSource::PageUrl => {
2549                matches!(
2550                    self.value_type,
2551                    WorkflowValueType::String | WorkflowValueType::Url
2552                )
2553            }
2554            WorkflowOutputSource::PageTitle | WorkflowOutputSource::VisibleText => true,
2555        };
2556        if !compatible {
2557            return Err(WorkflowValidationError::new(
2558                format!("{path}.valueType"),
2559                format!(
2560                    "{} output source requires a compatible value type",
2561                    self.source
2562                ),
2563            ));
2564        }
2565        Ok(())
2566    }
2567}
2568
2569/// Bounded, non-JavaScript sources for workflow outputs.
2570#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
2571#[serde(rename_all = "snake_case")]
2572pub enum WorkflowOutputSource {
2573    PageUrl,
2574    PageTitle,
2575    VisibleText,
2576}
2577
2578impl fmt::Display for WorkflowOutputSource {
2579    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
2580        formatter.write_str(match self {
2581            Self::PageUrl => "page_url",
2582            Self::PageTitle => "page_title",
2583            Self::VisibleText => "visible_text",
2584        })
2585    }
2586}
2587
2588/// A typed output value with its declared extraction source.
2589#[derive(Debug, Clone, Serialize, Deserialize)]
2590#[serde(rename_all = "camelCase")]
2591pub struct WorkflowOutput {
2592    pub value_type: WorkflowValueType,
2593    pub value: Value,
2594    #[serde(default, skip_serializing_if = "is_false")]
2595    pub redacted: bool,
2596    pub evidence: WorkflowOutputEvidence,
2597}
2598
2599fn is_false(value: &bool) -> bool {
2600    !*value
2601}
2602
2603/// Bounded provenance for a typed workflow output.
2604#[derive(Debug, Clone, Serialize, Deserialize)]
2605#[serde(rename_all = "camelCase")]
2606pub struct WorkflowOutputEvidence {
2607    pub source: WorkflowOutputSource,
2608    pub revision: u64,
2609}
2610
2611/// A path-aware validation failure.
2612#[derive(Debug, Clone, Serialize)]
2613pub struct WorkflowValidationError {
2614    pub path: String,
2615    pub reason: String,
2616}
2617
2618impl WorkflowValidationError {
2619    pub(crate) fn new(path: impl Into<String>, reason: impl Into<String>) -> Self {
2620        Self {
2621            path: path.into(),
2622            reason: reason.into(),
2623        }
2624    }
2625}
2626
2627impl fmt::Display for WorkflowValidationError {
2628    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
2629        write!(formatter, "{}: {}", self.path, self.reason)
2630    }
2631}
2632
2633impl std::error::Error for WorkflowValidationError {}
2634
2635fn require_object_fields(value: &Value, fields: &[&str]) -> Result<(), WorkflowValidationError> {
2636    let object = value.as_object().ok_or_else(|| {
2637        WorkflowValidationError::new("$", "workflow definition must be a JSON object")
2638    })?;
2639    for field in fields {
2640        if !object.contains_key(*field) {
2641            return Err(WorkflowValidationError::new(
2642                format!("$.{field}"),
2643                "required field is missing",
2644            ));
2645        }
2646    }
2647    Ok(())
2648}
2649
2650fn reject_unknown_fields(
2651    value: &Value,
2652    path: &str,
2653    allowed: &[&str],
2654) -> Result<(), WorkflowValidationError> {
2655    let object = value
2656        .as_object()
2657        .ok_or_else(|| WorkflowValidationError::new(path, "expected a JSON object"))?;
2658    if let Some(field) = object
2659        .keys()
2660        .find(|field| !allowed.contains(&field.as_str()))
2661    {
2662        return Err(WorkflowValidationError::new(
2663            format!("{path}.{field}"),
2664            "unknown workflow field",
2665        ));
2666    }
2667    Ok(())
2668}
2669
2670fn reject_predicate_fields(value: &Value, path: &str) -> Result<(), WorkflowValidationError> {
2671    let object = value
2672        .as_object()
2673        .ok_or_else(|| WorkflowValidationError::new(path, "expected a predicate object"))?;
2674    let allowed = [
2675        "urlEquals",
2676        "titleContains",
2677        "visible",
2678        "textContains",
2679        "popupOpened",
2680        "dialogOpen",
2681        "downloadStarted",
2682        "revisionEquals",
2683        "all",
2684        "any",
2685        "not",
2686    ];
2687    if object.len() != 1 {
2688        return Err(WorkflowValidationError::new(
2689            path,
2690            "predicate must contain exactly one recognized field",
2691        ));
2692    }
2693    let Some((field, nested)) = object.iter().next() else {
2694        return Err(WorkflowValidationError::new(
2695            path,
2696            "predicate must not be empty",
2697        ));
2698    };
2699    if !allowed.contains(&field.as_str()) {
2700        return Err(WorkflowValidationError::new(
2701            format!("{path}.{field}"),
2702            "unknown predicate field",
2703        ));
2704    }
2705    match field.as_str() {
2706        "all" | "any" => {
2707            let predicates = nested.as_array().ok_or_else(|| {
2708                WorkflowValidationError::new(format!("{path}.{field}"), "must be an array")
2709            })?;
2710            for (index, predicate) in predicates.iter().enumerate() {
2711                reject_predicate_fields(predicate, &format!("{path}.{field}[{index}]"))?;
2712            }
2713        }
2714        "not" => reject_predicate_fields(nested, &format!("{path}.not"))?,
2715        _ => {}
2716    }
2717    Ok(())
2718}
2719
2720fn validate_name(path: &str, value: &str) -> Result<(), WorkflowValidationError> {
2721    validate_bytes(path, value, 1, MAX_NAME_BYTES)?;
2722    if value.trim() != value {
2723        return Err(WorkflowValidationError::new(
2724            path,
2725            "must not have leading or trailing whitespace",
2726        ));
2727    }
2728    Ok(())
2729}
2730
2731fn validate_map_key(scope: &str, key: &str) -> Result<(), WorkflowValidationError> {
2732    validate_name(&format!("{scope}.{key}"), key)
2733}
2734
2735fn validate_bytes(
2736    path: &str,
2737    value: &str,
2738    minimum: usize,
2739    maximum: usize,
2740) -> Result<(), WorkflowValidationError> {
2741    if value.len() < minimum || value.len() > maximum {
2742        return Err(WorkflowValidationError::new(
2743            path,
2744            format!("must be {minimum}..={maximum} UTF-8 bytes"),
2745        ));
2746    }
2747    Ok(())
2748}
2749
2750fn validate_target(path: &str, target: &str) -> Result<(), WorkflowValidationError> {
2751    validate_bytes(path, target, 1, MAX_TARGET_BYTES)
2752}
2753
2754fn validate_predicate(
2755    predicate: &VerificationPredicate,
2756    path: &str,
2757) -> Result<(), WorkflowValidationError> {
2758    predicate
2759        .validate(0)
2760        .map_err(|error| WorkflowValidationError::new(path, error.to_string()))
2761}
2762
2763fn validate_batch_step(action: &BatchStep, path: &str) -> Result<(), WorkflowValidationError> {
2764    match action {
2765        BatchStep::Navigate { url, timeout_ms } => {
2766            if Url::parse(url).is_err() {
2767                return Err(WorkflowValidationError::new(
2768                    format!("{path}.url"),
2769                    "must be an absolute URL",
2770                ));
2771            }
2772            if *timeout_ms == 0 || *timeout_ms > MAX_DURATION_MS {
2773                return Err(WorkflowValidationError::new(
2774                    format!("{path}.timeoutMs"),
2775                    format!("must be 1..={MAX_DURATION_MS}"),
2776                ));
2777            }
2778        }
2779        BatchStep::Click { target }
2780        | BatchStep::Check { target }
2781        | BatchStep::Uncheck { target }
2782        | BatchStep::Clear { target } => validate_target(&format!("{path}.target"), target)?,
2783        BatchStep::Type { text, target } => {
2784            validate_bytes(&format!("{path}.text"), text, 0, MAX_TEXT_BYTES)?;
2785            if let Some(target) = target {
2786                validate_target(&format!("{path}.target"), target)?;
2787            }
2788        }
2789        BatchStep::Select { target, value } => {
2790            validate_target(&format!("{path}.target"), target)?;
2791            validate_bytes(&format!("{path}.value"), value, 1, MAX_TEXT_BYTES)?;
2792        }
2793        BatchStep::Scroll { dx, dy } => {
2794            if !dx.is_finite()
2795                || !dy.is_finite()
2796                || dx.abs() > 1_000_000.0
2797                || dy.abs() > 1_000_000.0
2798            {
2799                return Err(WorkflowValidationError::new(
2800                    path,
2801                    "scroll deltas must be finite and bounded",
2802                ));
2803            }
2804        }
2805        BatchStep::Wait {
2806            condition,
2807            timeout_ms,
2808        } => {
2809            validate_bytes(
2810                &format!("{path}.condition"),
2811                condition,
2812                1,
2813                MAX_WAIT_CONDITION_BYTES,
2814            )?;
2815            if *timeout_ms == 0 || *timeout_ms > MAX_DURATION_MS {
2816                return Err(WorkflowValidationError::new(
2817                    format!("{path}.timeoutMs"),
2818                    format!("must be 1..={MAX_DURATION_MS}"),
2819                ));
2820            }
2821        }
2822        BatchStep::Observe { .. }
2823        | BatchStep::Screenshot
2824        | BatchStep::AcceptDialog
2825        | BatchStep::DismissDialog => {}
2826        BatchStep::Evaluate { .. } => {
2827            return Err(WorkflowValidationError::new(
2828                path,
2829                "evaluate is not permitted in a declarative workflow",
2830            ));
2831        }
2832    }
2833    Ok(())
2834}
2835
2836#[cfg(test)]
2837mod tests {
2838    use super::*;
2839    use crate::browser::session::{
2840        FingerprintInvalidation, SemanticIntentCandidate, SemanticIntentPurpose,
2841        SemanticRegionKind, SemanticTargetFingerprint,
2842    };
2843    use serde_json::json;
2844
2845    fn definition() -> WorkflowDefinition {
2846        WorkflowDefinition {
2847            schema_version: WORKFLOW_SCHEMA_VERSION,
2848            name: "example".into(),
2849            workflow_version: "1.0.0".into(),
2850            description: None,
2851            inputs: BTreeMap::from([(
2852                "url".into(),
2853                WorkflowInput {
2854                    value_type: WorkflowValueType::Url,
2855                    required: true,
2856                    max_length: Some(2_048),
2857                    sensitive: None,
2858                },
2859            )]),
2860            budgets: WorkflowBudgets {
2861                max_steps: 2,
2862                max_duration_ms: 30_000,
2863                max_retries: 1,
2864                max_extracted_bytes: 4_096,
2865            },
2866            preconditions: vec![],
2867            steps: vec![WorkflowStep {
2868                id: "open".into(),
2869                action: BatchStep::Navigate {
2870                    url: "https://example.com".into(),
2871                    timeout_ms: 20_000,
2872                },
2873                intent: None,
2874                when: None,
2875                expect: Some(VerificationPredicate::TitleContains {
2876                    value: "Example".into(),
2877                }),
2878                before_retry: None,
2879                transaction: WorkflowTransactionClass::ReadOnly,
2880                idempotency_key: None,
2881                max_retries: 0,
2882                repeat: 1,
2883            }],
2884            terminal_condition: VerificationPredicate::UrlEquals {
2885                value: "https://example.com/".into(),
2886            },
2887            outputs: BTreeMap::new(),
2888        }
2889    }
2890
2891    #[test]
2892    fn canonical_json_is_stable_and_uses_camel_case() {
2893        let workflow = definition();
2894        let first = workflow.to_canonical_json().unwrap();
2895        let second = workflow.to_canonical_json().unwrap();
2896        assert_eq!(first, second);
2897        assert!(first.contains("\"schemaVersion\":1"));
2898        assert!(first.contains("\"timeoutMs\":20000"));
2899        assert!(first.contains("\"workflowVersion\":\"1.0.0\""));
2900    }
2901
2902    #[test]
2903    fn before_retry_marker_is_canonical_and_validated() {
2904        let mut workflow = definition();
2905        workflow.steps[0].before_retry = Some(VerificationPredicate::TitleContains {
2906            value: "Already saved".into(),
2907        });
2908        let json = workflow.to_canonical_json().unwrap();
2909        assert!(json.contains("\"beforeRetry\":{"));
2910        let parsed = WorkflowDefinition::from_json(&json).unwrap();
2911        assert!(parsed.steps[0].before_retry.is_some());
2912    }
2913
2914    #[test]
2915    fn type_alias_is_accepted_but_canonical_json_uses_value_type() {
2916        let mut value = serde_json::to_value(definition()).unwrap();
2917        let input = value["inputs"]["url"].as_object_mut().unwrap();
2918        let value_type = input.remove("valueType").unwrap();
2919        input.insert("type".into(), value_type);
2920
2921        let parsed = WorkflowDefinition::from_value(value).unwrap();
2922        let canonical = parsed.to_canonical_json().unwrap();
2923        assert!(canonical.contains("\"valueType\":\"url\""));
2924        assert!(!canonical.contains("\"type\":\"url\""));
2925    }
2926
2927    #[test]
2928    fn effect_marker_commit_follows_evidence_states_without_dispatch() {
2929        let mut record = WorkflowStepRecord::new("save");
2930        record.transition(WorkflowStepState::Ready).unwrap();
2931        commit_workflow_effect_marker(&mut record, 9);
2932        assert_eq!(record.state, WorkflowStepState::Committed);
2933        assert!(!record.dispatch_acknowledged);
2934        assert!(record.effect_observed);
2935        assert!(record.postcondition_verified);
2936        assert_eq!(record.current_revision, Some(9));
2937    }
2938
2939    #[test]
2940    fn from_json_reports_missing_top_level_path() {
2941        let mut value = serde_json::to_value(definition()).unwrap();
2942        value.as_object_mut().unwrap().remove("steps");
2943        let error = WorkflowDefinition::from_value(value).unwrap_err();
2944        assert_eq!(error.path, "$.steps");
2945    }
2946
2947    #[test]
2948    fn unknown_definition_fields_fail_before_deserialization() {
2949        let mut value = serde_json::to_value(definition()).unwrap();
2950        value
2951            .as_object_mut()
2952            .unwrap()
2953            .insert("futureField".into(), json!(true));
2954        let error = WorkflowDefinition::from_value(value).unwrap_err();
2955        assert_eq!(error.path, "$.futureField");
2956    }
2957
2958    #[test]
2959    fn semantic_intent_steps_round_trip_with_bounded_defaults() {
2960        let mut value = serde_json::to_value(definition()).unwrap();
2961        value["steps"][0] = json!({
2962            "id": "continue",
2963            "intent": {
2964                "action": "click",
2965                "purpose": "continueCheckout",
2966                "scope": {"regionKind": "checkoutSummary"},
2967                "resolutionPolicy": "requireUniqueHighConfidence"
2968            },
2969            "transaction": "idempotent"
2970        });
2971        let workflow = WorkflowDefinition::from_value(value).unwrap();
2972        let intent = workflow.steps[0].intent.as_ref().unwrap();
2973        assert_eq!(intent.action, SemanticIntentAction::Click);
2974        assert_eq!(intent.purpose.as_deref(), Some("continueCheckout"));
2975        assert_eq!(
2976            intent
2977                .execution_request("steps[0].intent")
2978                .unwrap()
2979                .request
2980                .intent,
2981            "continue checkout"
2982        );
2983        let canonical = workflow.to_canonical_json().unwrap();
2984        assert!(canonical.contains("\"resolutionPolicy\":\"requireUniqueHighConfidence\""));
2985        assert!(!canonical.contains("\"action\":\"navigate\""));
2986    }
2987
2988    #[test]
2989    fn semantic_intent_steps_reject_ambiguous_shape_and_missing_values() {
2990        let mut value = serde_json::to_value(definition()).unwrap();
2991        value["steps"][0] = json!({
2992            "id": "bad",
2993            "intent": {
2994                "action": "type",
2995                "purpose": "enterSearch",
2996                "intent": "enter search",
2997                "resolutionPolicy": "requireUniqueHighConfidence"
2998            },
2999            "transaction": "idempotent"
3000        });
3001        let error = WorkflowDefinition::from_value(value).unwrap_err();
3002        assert_eq!(error.path, "steps[0].intent.purpose");
3003
3004        let mut value = serde_json::to_value(definition()).unwrap();
3005        value["steps"][0] = json!({
3006            "id": "bad",
3007            "intent": {
3008                "action": "type",
3009                "purpose": "enterSearch",
3010                "resolutionPolicy": "requireUniqueHighConfidence"
3011            },
3012            "transaction": "idempotent"
3013        });
3014        let error = WorkflowDefinition::from_value(value).unwrap_err();
3015        assert_eq!(error.path, "steps[0].intent.value");
3016    }
3017
3018    #[test]
3019    fn unknown_predicate_fields_fail_with_their_json_path() {
3020        let mut value = serde_json::to_value(definition()).unwrap();
3021        value["terminalCondition"] = json!({"titleContains": "Example", "future": true});
3022        let error = WorkflowDefinition::from_value(value).unwrap_err();
3023        assert_eq!(error.path, "terminalCondition");
3024    }
3025
3026    #[test]
3027    fn duplicate_step_ids_are_rejected() {
3028        let mut workflow = definition();
3029        workflow.steps.push(WorkflowStep {
3030            id: "open".into(),
3031            action: BatchStep::Screenshot,
3032            intent: None,
3033            when: None,
3034            expect: None,
3035            before_retry: None,
3036            transaction: WorkflowTransactionClass::ReadOnly,
3037            idempotency_key: None,
3038            max_retries: 0,
3039            repeat: 1,
3040        });
3041        let error = workflow.validate().unwrap_err();
3042        assert_eq!(error.path, "steps[1].id");
3043    }
3044
3045    #[test]
3046    fn invalid_budget_reports_exact_path() {
3047        let mut workflow = definition();
3048        workflow.budgets.max_steps = 0;
3049        let error = workflow.validate().unwrap_err();
3050        assert_eq!(error.path, "budgets.maxSteps");
3051    }
3052
3053    #[test]
3054    fn conditional_idempotency_requires_a_key() {
3055        let mut workflow = definition();
3056        workflow.steps[0].transaction = WorkflowTransactionClass::ConditionallyIdempotent;
3057        let error = workflow.validate().unwrap_err();
3058        assert_eq!(error.path, "steps[0].idempotencyKey");
3059    }
3060
3061    #[test]
3062    fn duplicate_idempotency_keys_are_rejected() {
3063        let mut workflow = definition();
3064        workflow.steps[0].transaction = WorkflowTransactionClass::ConditionallyIdempotent;
3065        workflow.steps[0].idempotency_key = Some("save-once".into());
3066        workflow.steps.push(WorkflowStep {
3067            id: "second".into(),
3068            action: BatchStep::Screenshot,
3069            intent: None,
3070            when: None,
3071            expect: None,
3072            before_retry: None,
3073            transaction: WorkflowTransactionClass::ConditionallyIdempotent,
3074            idempotency_key: Some("save-once".into()),
3075            max_retries: 0,
3076            repeat: 1,
3077        });
3078        workflow.budgets.max_steps = 2;
3079        let error = workflow.validate().unwrap_err();
3080        assert_eq!(error.path, "steps[1].idempotencyKey");
3081    }
3082
3083    #[test]
3084    fn unknown_steps_cannot_request_automatic_retries() {
3085        let mut workflow = definition();
3086        workflow.steps[0].transaction = WorkflowTransactionClass::Unknown;
3087        workflow.budgets.max_retries = 1;
3088        workflow.steps[0].max_retries = 1;
3089        let error = workflow.validate().unwrap_err();
3090        assert_eq!(error.path, "steps[0].maxRetries");
3091    }
3092
3093    #[test]
3094    fn bounded_repetition_requires_retry_safe_class() {
3095        let mut workflow = definition();
3096        workflow.budgets.max_steps = 2;
3097        workflow.steps[0].repeat = 2;
3098        workflow.steps[0].transaction = WorkflowTransactionClass::Unknown;
3099        let error = workflow.validate().unwrap_err();
3100        assert_eq!(error.path, "steps[0].repeat");
3101    }
3102
3103    #[test]
3104    fn retry_policy_never_replays_after_dispatch() {
3105        assert!(can_retry_before_dispatch(
3106            WorkflowTransactionClass::Idempotent,
3107            false,
3108            1,
3109            1
3110        ));
3111        assert!(!can_retry_before_dispatch(
3112            WorkflowTransactionClass::NonIdempotent,
3113            false,
3114            1,
3115            1
3116        ));
3117        assert!(!can_retry_before_dispatch(
3118            WorkflowTransactionClass::Idempotent,
3119            true,
3120            1,
3121            1
3122        ));
3123    }
3124
3125    #[test]
3126    fn workflow_checkpoint_is_deterministic_and_redacted() {
3127        let checkpoint = WorkflowCheckpoint {
3128            schema_version: WORKFLOW_CHECKPOINT_SCHEMA_VERSION,
3129            run_id: "run_test".into(),
3130            workflow_name: "example".into(),
3131            workflow_version: "1.0.0".into(),
3132            definition_hash: "a".repeat(64),
3133            status: WorkflowRunStatus::Failed,
3134            next_step_index: 1,
3135            steps: vec![
3136                WorkflowCheckpointStep {
3137                    id: "open".into(),
3138                    state: WorkflowStepState::Committed,
3139                    attempts: 1,
3140                    history: Vec::new(),
3141                    execution_ids: Vec::new(),
3142                    dispatch_acknowledged: false,
3143                    effect_observed: false,
3144                    postcondition_verified: false,
3145                    retry_safe: false,
3146                    previous_revision: None,
3147                    current_revision: None,
3148                    branch_decision: None,
3149                    intent_evidence: None,
3150                },
3151                WorkflowCheckpointStep {
3152                    id: "save".into(),
3153                    state: WorkflowStepState::FailedBeforeDispatch,
3154                    attempts: 1,
3155                    history: Vec::new(),
3156                    execution_ids: Vec::new(),
3157                    dispatch_acknowledged: false,
3158                    effect_observed: false,
3159                    postcondition_verified: false,
3160                    retry_safe: false,
3161                    previous_revision: None,
3162                    current_revision: None,
3163                    branch_decision: None,
3164                    intent_evidence: None,
3165                },
3166            ],
3167            page: WorkflowCheckpointPage {
3168                target_id: "target".into(),
3169                frame_id: "frame".into(),
3170                url: "https://example.com".into(),
3171                title: "Example".into(),
3172                revision: 3,
3173            },
3174        };
3175        let first = checkpoint.to_canonical_json().unwrap();
3176        let second = checkpoint.to_canonical_json().unwrap();
3177        assert_eq!(first, second);
3178        assert!(!first.contains("password"));
3179        assert_eq!(
3180            crate::BrowserSession::parse_workflow_checkpoint(&first)
3181                .unwrap()
3182                .next_step_index,
3183            1
3184        );
3185    }
3186
3187    #[test]
3188    fn minimal_workflow_fixture_round_trips() {
3189        let workflow = WorkflowDefinition::from_json(include_str!(
3190            "../../../tests/fixtures/workflow-minimal.json"
3191        ))
3192        .unwrap();
3193        assert_eq!(workflow.name, "open-example");
3194        assert!(workflow.preconditions.is_empty());
3195        assert!(workflow.to_canonical_json().is_ok());
3196    }
3197
3198    #[test]
3199    fn evaluate_is_rejected_before_execution() {
3200        let mut workflow = definition();
3201        workflow.steps[0].action = BatchStep::Evaluate {
3202            expression: "document.title".into(),
3203        };
3204        let error = workflow.validate().unwrap_err();
3205        assert!(error.reason.contains("not permitted"));
3206    }
3207
3208    #[test]
3209    fn inputs_are_type_checked_and_bounded() {
3210        let workflow = definition();
3211        let values = BTreeMap::from([("url".into(), json!("https://example.com"))]);
3212        workflow.validate_inputs(&values).unwrap();
3213
3214        let bad_values = BTreeMap::from([("url".into(), json!("not a url"))]);
3215        let error = workflow.validate_inputs(&bad_values).unwrap_err();
3216        assert_eq!(error.path, "inputs.url");
3217    }
3218
3219    #[test]
3220    fn resolve_actions_substitutes_bounded_inputs() {
3221        let mut workflow = definition();
3222        workflow.steps[0].action = BatchStep::Navigate {
3223            url: "${inputs.url}".into(),
3224            timeout_ms: 20_000,
3225        };
3226        let values = BTreeMap::from([("url".into(), json!("https://docs.example.com"))]);
3227        let steps = workflow.resolve_actions(&values).unwrap();
3228        let BatchStep::Navigate { url, .. } = &steps[0].action else {
3229            panic!("expected navigate action");
3230        };
3231        assert_eq!(url, "https://docs.example.com");
3232    }
3233
3234    #[test]
3235    fn resolve_actions_rejects_unknown_placeholders_before_dispatch() {
3236        let mut workflow = definition();
3237        workflow.steps[0].action = BatchStep::Click {
3238            target: "name=${inputs.missing}".into(),
3239        };
3240        let values = BTreeMap::from([("url".into(), json!("https://example.com"))]);
3241        let error = workflow.resolve_actions(&values).unwrap_err();
3242        assert_eq!(error.path, "steps[0].action.target");
3243    }
3244
3245    #[test]
3246    fn typed_output_values_require_strict_scalar_conversions() {
3247        assert_eq!(
3248            typed_output_value("count", WorkflowValueType::Integer, " 42 ").unwrap(),
3249            json!(42)
3250        );
3251        assert_eq!(
3252            typed_output_value("ready", WorkflowValueType::Boolean, "false").unwrap(),
3253            json!(false)
3254        );
3255        assert!(typed_output_value("count", WorkflowValueType::Integer, "4.2").is_err());
3256        assert!(typed_output_value("ready", WorkflowValueType::Boolean, "yes").is_err());
3257    }
3258
3259    #[test]
3260    fn sensitive_output_serialization_contains_no_literal_value() {
3261        let output = WorkflowOutput {
3262            value_type: WorkflowValueType::String,
3263            value: Value::Null,
3264            redacted: true,
3265            evidence: WorkflowOutputEvidence {
3266                source: WorkflowOutputSource::VisibleText,
3267                revision: 4,
3268            },
3269        };
3270        let serialized = serde_json::to_string(&output).unwrap();
3271        assert!(serialized.contains("\"redacted\":true"));
3272        assert!(!serialized.contains("secret"));
3273    }
3274
3275    #[test]
3276    fn recorder_keeps_semantic_targets_and_redacts_typed_values() {
3277        let mut recorder = WorkflowRecorder::new("checkout", "1.0.0");
3278        recorder
3279            .record_click("continue", "button", "Continue", None)
3280            .unwrap();
3281        recorder
3282            .record_type_input("email", "textbox", "Email", "email")
3283            .unwrap();
3284        let draft = recorder.draft();
3285        let BatchStep::Click { target } = &draft.steps[0].action else {
3286            panic!("expected semantic click draft");
3287        };
3288        assert_eq!(target, "role=button;name=Continue");
3289        let serialized = serde_json::to_string(draft).unwrap();
3290        assert!(!serialized.contains("password-value"));
3291        assert!(serialized.contains("${inputs.email}"));
3292        assert!(draft.steps[1].review_required);
3293        let inferred = recorder.inferred_inputs();
3294        assert_eq!(inferred["email"].value_type, WorkflowValueType::String);
3295        assert_eq!(inferred["email"].sensitive, None);
3296    }
3297
3298    #[test]
3299    fn recorder_captures_semantic_evidence_without_replay_handles_or_query_values() {
3300        let request = SemanticIntentRequest {
3301            schema_version: INTENT_RESOLUTION_SCHEMA_VERSION,
3302            intent: "submit order".into(),
3303            action: SemanticIntentAction::Click,
3304            scope: IntentScope::default(),
3305            constraints: IntentConstraints::default(),
3306            resolution_policy: SemanticResolutionPolicy::RequireUniqueHighConfidence,
3307            expected_revision: None,
3308        };
3309        let route = SemanticRouteIdentity {
3310            target_id: "target-7".into(),
3311            frame_id: "frame-2".into(),
3312            url: "https://shop.example/orders?token=secret#confirmation".into(),
3313        };
3314        let result = SemanticIntentResult {
3315            schema_version: INTENT_RESOLUTION_SCHEMA_VERSION,
3316            intent: request.intent.clone(),
3317            action: request.action,
3318            normalized_intent: "submit order".into(),
3319            resolution: SemanticResolution::UniqueHighConfidence,
3320            policy_decision: IntentPolicyDecision::Allowed,
3321            route: Some(route.clone()),
3322            revision: Some(7),
3323            candidates: vec![SemanticIntentCandidate {
3324                id: "candidate-1".into(),
3325                reference: "r7:backend-secret".into(),
3326                role: "button".into(),
3327                name: "Submit order".into(),
3328                input_type: None,
3329                region_id: Some("checkout-region".into()),
3330                region_kind: Some(SemanticRegionKind::CheckoutSummary),
3331                confidence: IntentConfidence::High,
3332                evidence: Vec::new(),
3333                fingerprint: Some(SemanticTargetFingerprint {
3334                    revision: 7,
3335                    route: route.clone(),
3336                    role: "button".into(),
3337                    name: "Submit order".into(),
3338                    input_type: None,
3339                    region_id: Some("checkout-region".into()),
3340                    region_kind: Some(SemanticRegionKind::CheckoutSummary),
3341                    purpose: SemanticIntentPurpose::Submit,
3342                    invalidated_by: vec![FingerprintInvalidation::Revision],
3343                }),
3344            }],
3345            excluded_candidates: Vec::new(),
3346            excluded_count: 0,
3347            selected_candidate: Some("candidate-1".into()),
3348            suggested_constraints: Vec::new(),
3349            reason: None,
3350        };
3351        let mut recorder = WorkflowRecorder::new("checkout", "1.0.0");
3352        recorder
3353            .record_semantic_intent(
3354                "submit",
3355                &request,
3356                &result,
3357                None::<String>,
3358                WorkflowTransactionClass::NonIdempotent,
3359                Some(VerificationPredicate::TextContains {
3360                    value: "Order submitted".into(),
3361                }),
3362            )
3363            .unwrap();
3364
3365        let draft = recorder.draft();
3366        let step = &draft.steps[0];
3367        assert_eq!(
3368            step.target.as_ref().unwrap().region_kind,
3369            Some(SemanticRegionKind::CheckoutSummary)
3370        );
3371        assert_eq!(step.semantic.as_ref().unwrap().revision, Some(7));
3372        assert!(step.semantic.as_ref().unwrap().target_fingerprint.is_some());
3373        let serialized = serde_json::to_string(draft).unwrap();
3374        assert!(!serialized.contains("backend-secret"));
3375        assert!(!serialized.contains("token=secret"));
3376        assert!(!serialized.contains("target-7"));
3377        assert!(!serialized.contains("frame-2"));
3378
3379        let definition = recorder
3380            .into_definition(
3381                BTreeMap::new(),
3382                WorkflowBudgets {
3383                    max_steps: 1,
3384                    max_duration_ms: 30_000,
3385                    max_retries: 0,
3386                    max_extracted_bytes: 4_096,
3387                },
3388                VerificationPredicate::TextContains {
3389                    value: "Order submitted".into(),
3390                },
3391                BTreeMap::new(),
3392            )
3393            .unwrap();
3394        assert!(definition.steps[0].intent.is_some());
3395        assert!(definition.steps[0].expect.is_some());
3396    }
3397
3398    #[test]
3399    fn recorder_keeps_ambiguous_semantic_results_unselected() {
3400        let request = SemanticIntentRequest {
3401            schema_version: INTENT_RESOLUTION_SCHEMA_VERSION,
3402            intent: "open settings".into(),
3403            action: SemanticIntentAction::Click,
3404            scope: IntentScope::default(),
3405            constraints: IntentConstraints::default(),
3406            resolution_policy: SemanticResolutionPolicy::RequireUniqueHighConfidence,
3407            expected_revision: None,
3408        };
3409        let candidate = |id: &str| SemanticIntentCandidate {
3410            id: id.into(),
3411            reference: format!("{id}:ref"),
3412            role: "button".into(),
3413            name: "Settings".into(),
3414            input_type: None,
3415            region_id: None,
3416            region_kind: None,
3417            confidence: IntentConfidence::High,
3418            evidence: Vec::new(),
3419            fingerprint: None,
3420        };
3421        let result = SemanticIntentResult {
3422            schema_version: INTENT_RESOLUTION_SCHEMA_VERSION,
3423            intent: request.intent.clone(),
3424            action: request.action,
3425            normalized_intent: request.intent.clone(),
3426            resolution: SemanticResolution::Ambiguous,
3427            policy_decision: IntentPolicyDecision::Rejected,
3428            route: None,
3429            revision: Some(3),
3430            candidates: vec![candidate("one"), candidate("two")],
3431            excluded_candidates: Vec::new(),
3432            excluded_count: 0,
3433            selected_candidate: None,
3434            suggested_constraints: Vec::new(),
3435            reason: Some("two candidates remain".into()),
3436        };
3437        let mut recorder = WorkflowRecorder::new("settings", "1.0.0");
3438        recorder
3439            .record_semantic_intent(
3440                "open-settings",
3441                &request,
3442                &result,
3443                None::<String>,
3444                WorkflowTransactionClass::Unknown,
3445                None,
3446            )
3447            .unwrap();
3448        let step = &recorder.draft().steps[0];
3449        assert!(step.target.is_none());
3450        assert_eq!(step.confidence, WorkflowRecordingConfidence::Low);
3451        assert!(step.semantic.as_ref().unwrap().ambiguous);
3452        assert!(step.review_required);
3453    }
3454
3455    #[test]
3456    fn workflow_step_flattens_batch_action() {
3457        let value = serde_json::to_value(&definition().steps[0]).unwrap();
3458        assert_eq!(value["id"], "open");
3459        assert_eq!(value["action"], "navigate");
3460        assert_eq!(value["timeoutMs"], 20_000);
3461    }
3462
3463    #[test]
3464    fn state_machine_accepts_linear_commit_and_rejects_invalid_jump() {
3465        assert!(WorkflowStepState::Pending.can_transition_to(WorkflowStepState::Ready));
3466        assert!(WorkflowStepState::Resolving.can_transition_to(WorkflowStepState::Dispatched));
3467        assert!(!WorkflowStepState::Pending.can_transition_to(WorkflowStepState::Committed));
3468
3469        let mut record = WorkflowStepRecord::new("open");
3470        for state in [
3471            WorkflowStepState::Ready,
3472            WorkflowStepState::Preflight,
3473            WorkflowStepState::Resolving,
3474            WorkflowStepState::Dispatched,
3475            WorkflowStepState::EffectObserved,
3476            WorkflowStepState::Verified,
3477            WorkflowStepState::OutputsExtracted,
3478            WorkflowStepState::Committed,
3479        ] {
3480            record.transition(state).unwrap();
3481        }
3482        assert_eq!(record.state, WorkflowStepState::Committed);
3483        assert_eq!(record.history.len(), 9);
3484        assert!(record.transition(WorkflowStepState::Preflight).is_err());
3485        let trace = WorkflowTrace::from_steps(&[record]);
3486        trace.validate().unwrap();
3487        assert_eq!(trace.events[0].sequence, 0);
3488        assert_eq!(
3489            trace.events.last().unwrap().state,
3490            WorkflowStepState::Committed
3491        );
3492    }
3493
3494    #[test]
3495    fn step_record_serializes_bounded_execution_evidence() {
3496        let mut record = WorkflowStepRecord::new("open");
3497        record.execution_ids = vec!["act_7".into(), "act_8".into()];
3498        record.dispatch_acknowledged = true;
3499        record.effect_observed = true;
3500        record.postcondition_verified = true;
3501        record.retry_safe = false;
3502        record.previous_revision = Some(7);
3503        record.current_revision = Some(8);
3504        record.intent_evidence = Some(WorkflowIntentEvidence {
3505            resolution_id: "res_1".into(),
3506            candidate_id: "candidate_1".into(),
3507            revision: 8,
3508            resolution: crate::browser::session::SemanticResolution::Exact,
3509            policy_decision: crate::browser::session::IntentPolicyDecision::Allowed,
3510            confidence: crate::browser::session::IntentConfidence::Exact,
3511            fingerprint: None,
3512        });
3513
3514        let trace = WorkflowTrace::from_steps(std::slice::from_ref(&record));
3515        assert_eq!(trace.intent_resolutions.len(), 1);
3516        let value = serde_json::to_value(record).unwrap();
3517        assert_eq!(value["executionIds"], serde_json::json!(["act_7", "act_8"]));
3518        assert_eq!(value["dispatchAcknowledged"], true);
3519        assert_eq!(value["effectObserved"], true);
3520        assert_eq!(value["postconditionVerified"], true);
3521        assert_eq!(value["retrySafe"], false);
3522        assert_eq!(value["previousRevision"], 7);
3523        assert_eq!(value["currentRevision"], 8);
3524        assert_eq!(value["intentEvidence"]["candidateId"], "candidate_1");
3525    }
3526
3527    #[test]
3528    fn budget_exhaustion_is_a_typed_run_status() {
3529        let result = WorkflowRunResult::budget_exhausted(
3530            &definition(),
3531            "run_test".into(),
3532            vec![WorkflowStepRecord::new("open")],
3533            Some("open".into()),
3534            "maxSteps exhausted",
3535            3,
3536            3,
3537        );
3538        assert_eq!(result.status, WorkflowRunStatus::BudgetExhausted);
3539        assert_eq!(result.trace.run_id.as_deref(), Some("run_test"));
3540        let value = serde_json::to_value(result).unwrap();
3541        assert_eq!(value["status"], "budget_exhausted");
3542    }
3543
3544    #[test]
3545    fn trace_schema_version_is_independent_and_validated() {
3546        let trace = WorkflowTrace::from_steps(&[]);
3547        assert_eq!(trace.schema_version, WORKFLOW_TRACE_SCHEMA_VERSION);
3548        trace.validate().unwrap();
3549
3550        let mut unsupported = trace;
3551        unsupported.schema_version = WORKFLOW_TRACE_SCHEMA_VERSION + 1;
3552        let error = unsupported.validate().unwrap_err();
3553        assert_eq!(error.path, "trace.schemaVersion");
3554    }
3555
3556    #[test]
3557    fn post_dispatch_failures_require_resume_reconciliation() {
3558        let result = WorkflowRunResult::resume_required(
3559            &definition(),
3560            "run_test".into(),
3561            vec![WorkflowStepRecord::new("open")],
3562            Some("open".into()),
3563            "postcondition was not proven",
3564            3,
3565            4,
3566        );
3567        assert_eq!(result.status, WorkflowRunStatus::ResumeRequired);
3568        assert_eq!(
3569            serde_json::to_value(result).unwrap()["status"],
3570            "resume_required"
3571        );
3572    }
3573
3574    #[test]
3575    fn trace_replay_preserves_attempt_boundaries() {
3576        let mut record = WorkflowStepRecord::new("open");
3577        record.transition(WorkflowStepState::Ready).unwrap();
3578        record.transition(WorkflowStepState::Preflight).unwrap();
3579        record.attempts = 1;
3580        record.transition(WorkflowStepState::Resolving).unwrap();
3581        record.transition(WorkflowStepState::NotDispatched).unwrap();
3582        record.fail(WorkflowStepState::FailedBeforeDispatch, "before dispatch");
3583        record.transition(WorkflowStepState::Ready).unwrap();
3584        record.transition(WorkflowStepState::Preflight).unwrap();
3585        record.attempts = 2;
3586        record.transition(WorkflowStepState::Resolving).unwrap();
3587        record.transition(WorkflowStepState::Dispatched).unwrap();
3588        record
3589            .transition(WorkflowStepState::EffectObserved)
3590            .unwrap();
3591        record.transition(WorkflowStepState::Verified).unwrap();
3592        record
3593            .transition(WorkflowStepState::OutputsExtracted)
3594            .unwrap();
3595        record.transition(WorkflowStepState::Committed).unwrap();
3596
3597        let workflow = definition();
3598        let trace = WorkflowTrace::from_steps(&[record]);
3599        let replayed = trace.replay(&workflow).unwrap();
3600        assert_eq!(trace.events[2].attempt, 1);
3601        assert_eq!(trace.events[7].attempt, 2);
3602        assert_eq!(replayed[0].state, WorkflowStepState::Committed);
3603        assert_eq!(replayed[0].attempts, 2);
3604    }
3605
3606    #[test]
3607    fn trace_replay_rejects_a_prefix_that_skips_the_first_step() {
3608        let mut workflow = definition();
3609        workflow.steps.push(WorkflowStep {
3610            id: "save".into(),
3611            action: BatchStep::Screenshot,
3612            intent: None,
3613            when: None,
3614            expect: None,
3615            before_retry: None,
3616            transaction: WorkflowTransactionClass::ReadOnly,
3617            idempotency_key: None,
3618            max_retries: 0,
3619            repeat: 1,
3620        });
3621        let trace = WorkflowTrace::from_steps(&[WorkflowStepRecord::new("save")]);
3622        let error = trace.replay(&workflow).unwrap_err();
3623        assert_eq!(error.path, "trace.events[0].stepId");
3624    }
3625
3626    #[test]
3627    fn conditional_steps_record_and_replay_branch_decisions() {
3628        let mut workflow = definition();
3629        let predicate = VerificationPredicate::TitleContains {
3630            value: "Example".into(),
3631        };
3632        workflow.steps[0].when = Some(predicate.clone());
3633        let mut record = WorkflowStepRecord::new("open");
3634        record.transition(WorkflowStepState::Ready).unwrap();
3635        record.branch_decision = Some(WorkflowBranchDecision {
3636            step_id: "open".into(),
3637            predicate: predicate.clone(),
3638            matched: false,
3639        });
3640        record.transition(WorkflowStepState::Skipped).unwrap();
3641        let trace = WorkflowTrace::from_steps(&[record]);
3642        assert!(!trace.branch_decisions[0].matched);
3643        let replayed = trace.replay(&workflow).unwrap();
3644        assert_eq!(replayed[0].state, WorkflowStepState::Skipped);
3645        assert_eq!(
3646            replayed[0].branch_decision.as_ref().unwrap().predicate,
3647            predicate
3648        );
3649    }
3650}