1use 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
27pub 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#[derive(Debug, Clone, Serialize, Deserialize)]
50#[serde(rename_all = "camelCase")]
51pub struct WorkflowDefinition {
52 pub schema_version: u32,
54 pub name: String,
56 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 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 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 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 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 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 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#[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#[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#[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#[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#[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#[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 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#[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 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#[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#[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#[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#[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#[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#[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#[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#[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 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 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 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 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 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 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 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 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 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#[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#[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#[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#[derive(Debug, Clone, Serialize, Deserialize)]
2605#[serde(rename_all = "camelCase")]
2606pub struct WorkflowOutputEvidence {
2607 pub source: WorkflowOutputSource,
2608 pub revision: u64,
2609}
2610
2611#[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}