diff --git a/README.md b/README.md index 9e3fb06..75f52ab 100644 --- a/README.md +++ b/README.md @@ -30,7 +30,9 @@ product orchestrator, CLI, real provider adapter, operating-system sandbox, general scenario executor, automatic recovery scheduler, or public resume CLI. Prepared process execution now has [durable run state](docs/integrations/durable-state.md) with immutable plans, atomic snapshots, verified checkpoints, status inspection, -and explicit recovery decisions. +explicit recovery decisions, and fresh dependency-graph assessment before +caller-selected downstream execution. Stale prerequisites invalidate their +descendants for reuse while unrelated valid branches retain their eligibility. ## Executable checkpoint diff --git a/ROADMAP.md b/ROADMAP.md index d29f94a..b05f07e 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -3,12 +3,12 @@ schema: aether.architecture-document/v1 id: flow-roadmap title: Flow Roadmap kind: architecture-document -version: 1.3.2 +version: 1.3.3 status: draft owners: - egohygiene created: 2026-08-13 -updated: 2026-09-25 +updated: 2026-09-26 governed_by: - architecture-roadmap depends_on: @@ -23,7 +23,7 @@ supersedes: [] # Flow Roadmap -## 2026-09-25 live suite handoff +## 2026-09-26 live suite handoff > [!IMPORTANT] > This is the current near-term execution handoff for agents. It supersedes @@ -36,10 +36,18 @@ supersedes: [] Its [acceptance matrix](docs/integrations/acceptance-scenarios.md) proves 81 scenarios twice on Rust 1.85 and stable. -The #49 review candidate adds [durable prepared execution](docs/integrations/durable-state.md): +#49 merged through [PR #63](https://github.com/egohygiene/flow/pull/63) at +`711fbe3c19c0e080ab6c67b74d7945cd29675a74`, with +[default-branch CI green](https://github.com/egohygiene/flow/actions/runs/36213631491). +It adds [durable prepared execution](docs/integrations/durable-state.md): versioned plans and state, atomic immutable snapshots, workspace locking, accepted checkpoints, fresh resume eligibility, and explicit recovery decisions. -After this candidate merges and the default-branch gate passes, #31 is next. +#31 is active through three dependency-ordered review checkpoints: +[#64](https://github.com/egohygiene/flow/issues/64) adds fresh graph assessment and +safe dependent execution; [#65](https://github.com/egohygiene/flow/issues/65) adds +the deterministic durable lifecycle corpus; [#66](https://github.com/egohygiene/flow/issues/66) +proves authority, duplicate-effect prevention, and recovery residuals. Land one +review PR and verify default-branch CI before starting the next checkpoint. FLO-Q03 remains active until real released-provider adapters satisfy its remaining exit criteria. [#62](https://github.com/egohygiene/flow/issues/62) is later documentation visualization work and does not block this sequence. @@ -53,6 +61,7 @@ later documentation visualization work and does not block this sequence. recovery state. 3. [#31](https://github.com/egohygiene/flow/issues/31) — prove interruption, retry, resume, invalidation, and authority transitions against durable state. + Ordered children: #64 → #65 → #66; the parent remains open until all pass. 4. After #49, integrate immutable provider releases as they become available: [#50](https://github.com/egohygiene/flow/issues/50) for Optiflow, [#52](https://github.com/egohygiene/flow/issues/52) for Renderflow, and diff --git a/contracts/README.md b/contracts/README.md index 8db5cf2..28bc073 100644 --- a/contracts/README.md +++ b/contracts/README.md @@ -1,7 +1,7 @@ # Flow contract set This directory contains Flow-owned suite interchange contracts. The initial -contract set is version `0.7.0`, status `provisional`, in the v1 compatibility +contract set is version `0.8.0`, status `provisional`, in the v1 compatibility family. Provisional means versioned and testable, not stable for production. | Contract | Purpose | @@ -30,6 +30,7 @@ family. Provisional means versioned and testable, not stable for production. | `flow.run-validation/v1` | accepted evidence and validator implementation identities | | `flow.run-recovery/v1` | explicit retry/abandon decision and uncertainty acknowledgement | | `flow.run-snapshot/v1` | atomic-storage envelope with state integrity digest | +| `flow.run-assessment/v1` | read-only current graph eligibility correlated to exact saved state | `contract-set.v1.json` is the machine-readable index. Schemas live in `schemas/`; deterministic examples live in `examples/`; extension compatibility diff --git a/contracts/contract-set.v1.json b/contracts/contract-set.v1.json index d865957..c4f595d 100644 --- a/contracts/contract-set.v1.json +++ b/contracts/contract-set.v1.json @@ -1,6 +1,6 @@ { "schema_version": "flow.contract-set/v1", - "contract_set_version": "0.7.0", + "contract_set_version": "0.8.0", "status": "provisional", "contracts": [ { @@ -191,6 +191,14 @@ "invalid_examples": [ "fixtures/state/run-snapshot.v2.invalid.json" ] + }, + { + "id": "flow.run-assessment/v1", + "schema": "schemas/run-assessment.v1.schema.json", + "example": "examples/run-assessment.v1.example.json", + "invalid_examples": [ + "fixtures/state/run-assessment.v2.invalid.json" + ] } ] } diff --git a/contracts/examples/run-assessment.v1.example.json b/contracts/examples/run-assessment.v1.example.json new file mode 100644 index 0000000..67e2df7 --- /dev/null +++ b/contracts/examples/run-assessment.v1.example.json @@ -0,0 +1,15 @@ +{ + "schema_version": "flow.run-assessment/v1", + "plan_digest": "0e91d3a2af624fd222d6cf2aa1b005b8d9b3ad2f3e9951cde35b2e9e99b102d1", + "state_digest": "8796604491ccdbab553b73b319608fbbc824da503645927e810b7577a40442d3", + "sequence": 0, + "steps": [ + { + "step_id": "step:inspect", + "recorded_status": "pending", + "eligibility": "ready", + "stale_boundary": null, + "blocked_by": [] + } + ] +} diff --git a/contracts/fixtures/state/run-assessment.v2.invalid.json b/contracts/fixtures/state/run-assessment.v2.invalid.json new file mode 100644 index 0000000..3eeb00b --- /dev/null +++ b/contracts/fixtures/state/run-assessment.v2.invalid.json @@ -0,0 +1,15 @@ +{ + "schema_version": "flow.run-assessment/v2", + "plan_digest": "0e91d3a2af624fd222d6cf2aa1b005b8d9b3ad2f3e9951cde35b2e9e99b102d1", + "state_digest": "8796604491ccdbab553b73b319608fbbc824da503645927e810b7577a40442d3", + "sequence": 0, + "steps": [ + { + "step_id": "step:inspect", + "recorded_status": "pending", + "eligibility": "ready", + "stale_boundary": null, + "blocked_by": [] + } + ] +} diff --git a/contracts/schemas/run-assessment.v1.schema.json b/contracts/schemas/run-assessment.v1.schema.json new file mode 100644 index 0000000..456c70d --- /dev/null +++ b/contracts/schemas/run-assessment.v1.schema.json @@ -0,0 +1,106 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "https://egohygiene.github.io/flow/contracts/run-assessment.v1.schema.json", + "title": "Flow run assessment v1", + "type": "object", + "additionalProperties": false, + "required": [ + "schema_version", + "plan_digest", + "state_digest", + "sequence", + "steps" + ], + "properties": { + "schema_version": { + "const": "flow.run-assessment/v1" + }, + "plan_digest": { + "type": "string", + "pattern": "^[a-f0-9]{64}$" + }, + "state_digest": { + "type": "string", + "pattern": "^[a-f0-9]{64}$" + }, + "sequence": { + "type": "integer", + "minimum": 0, + "maximum": 4095 + }, + "steps": { + "type": "array", + "minItems": 1, + "maxItems": 256, + "items": { + "type": "object", + "additionalProperties": false, + "required": [ + "step_id", + "recorded_status", + "eligibility", + "stale_boundary", + "blocked_by" + ], + "properties": { + "step_id": { + "type": "string", + "maxLength": 160, + "pattern": "^step:[a-z0-9][a-z0-9._-]*$" + }, + "recorded_status": { + "enum": [ + "pending", + "running", + "succeeded", + "failed", + "cancelled", + "denied", + "abandoned" + ] + }, + "eligibility": { + "enum": [ + "ready", + "reusable", + "invalidated", + "dependency-blocked", + "approval-required", + "abandoned" + ] + }, + "stale_boundary": { + "anyOf": [ + { + "enum": [ + "invocation", + "provider", + "capability", + "configuration", + "authority", + "bindings", + "inputs", + "artifacts", + "validation" + ] + }, + { + "type": "null" + } + ] + }, + "blocked_by": { + "type": "array", + "maxItems": 256, + "uniqueItems": true, + "items": { + "type": "string", + "maxLength": 160, + "pattern": "^step:[a-z0-9][a-z0-9._-]*$" + } + } + } + } + } + } +} diff --git a/docs/architecture/foundation/ARCHITECTURE.md b/docs/architecture/foundation/ARCHITECTURE.md index 931c39a..a4fe739 100644 --- a/docs/architecture/foundation/ARCHITECTURE.md +++ b/docs/architecture/foundation/ARCHITECTURE.md @@ -3,12 +3,12 @@ schema: aether.architecture-document/v1 id: flow-architecture title: Flow Architecture kind: architecture-document -version: 0.12.0 +version: 0.13.0 status: draft owners: - egohygiene created: 2026-08-13 -updated: 2026-09-25 +updated: 2026-09-26 governed_by: - architecture-architecture depends_on: @@ -121,9 +121,17 @@ State contains typed decisions, identities, relative artifact bindings, and allowlisted validation evidence. Secret values, configuration values, raw process streams, and provider-authored messages remain outside the store. The workspace is trusted local state, not a cryptographically authenticated or sandboxed store. -ADR-0011 owns the persistence, migration, and recovery rationale. Graph scheduling, -automatic retry, provider-native checkpoints, and downstream invalidation policy -remain later orchestration work. +ADR-0011 owns the persistence, migration, and recovery rationale. ADR-0012 adds +read-only graph assessment: a complete current context inventory and the exact +saved plan are required before classifying every step in dependency order. +Completed prerequisites must be freshly reusable, not merely historically +succeeded. Stale evidence invalidates that step and every descendant for reuse; +unrelated branches retain their independently assessed eligibility. Pending or +unresolved prerequisites block their descendants. These assessments never rewrite +accepted history or grant execution authority. Dependent execution recomputes +the assessment, and single-step entry points refuse dependent steps. Automatic +scheduling, retry, cross-plan migration, and provider-native checkpoints remain +later orchestration work. ### External adapters @@ -271,7 +279,9 @@ workspace, and two fresh roots must produce equal portable evidence and output bytes. This is conformance infrastructure, not a scenario-manifest executor or production graph scheduler. Issue #49 adds the versioned durable coordinator, immutable prepared plans, atomic snapshots, accepted checkpoints, fresh reuse -assessment, and explicit recovery decisions described above. Flow does not yet supply +assessment, and explicit recovery decisions described above. Issue #64 requires +fresh prerequisite evidence for graph assessment and dependent execution without +adding a scheduler or rewriting accepted history. Flow does not yet supply the public CLI, real holon adapters, signature or transparency verification, provider-native artifact validation, an atomic filesystem snapshot, an operating-system sandbox or authenticated enforcement evidence, diff --git a/docs/architecture/governance/DECISIONS.md b/docs/architecture/governance/DECISIONS.md index 2382a2d..2b4b98f 100644 --- a/docs/architecture/governance/DECISIONS.md +++ b/docs/architecture/governance/DECISIONS.md @@ -3,12 +3,12 @@ schema: aether.architecture-document/v1 id: flow-decisions title: Flow Decisions kind: architecture-document -version: 0.9.0 +version: 0.10.0 status: draft owners: - egohygiene created: 2026-08-13 -updated: 2026-09-25 +updated: 2026-09-26 governed_by: - architecture-decisions depends_on: @@ -65,10 +65,13 @@ justifies separate ADRs. | [ADR-0009](decisions/ADR-0009-process-authority-isolation.md) | Bind process authority to explicit isolation evidence | Accepted | 2026-09-21 | None | A real sandbox, authenticated host evidence, or new authority dimension changes the preflight boundary | | [ADR-0010](decisions/ADR-0010-bounded-direct-process-supervision.md) | Bound direct provider launch and supervision | Accepted | 2026-09-21 | None | Sandbox enforcement, descriptor-bound launch, process-tree containment, or durable interruption changes the runner boundary | | [ADR-0011](decisions/ADR-0011-durable-run-state.md) | Persist prepared intent and acceptance in immutable local snapshots | Proposed | Pending review | None | Migration, distributed writers, automatic recovery, authenticated state, or stronger durability changes the boundary | +| [ADR-0012](decisions/ADR-0012-fresh-graph-assessment.md) | Require fresh prerequisite evidence for durable graph execution | Proposed | Pending review | None | Scheduling, cross-plan reuse, or concurrent artifact mutation changes the assessment boundary | ## Active decisions ADR-0011 is proposed with Flow #49 and remains subject to maintainer review. +Its implementation merged in PR #63. ADR-0012 is proposed with Flow #64, the +first bounded checkpoint under #31. The indexed ADRs are authoritative. Summaries in other documents must link back to them rather than recreate rationale. diff --git a/docs/architecture/governance/decisions/ADR-0012-fresh-graph-assessment.md b/docs/architecture/governance/decisions/ADR-0012-fresh-graph-assessment.md new file mode 100644 index 0000000..accfd1d --- /dev/null +++ b/docs/architecture/governance/decisions/ADR-0012-fresh-graph-assessment.md @@ -0,0 +1,99 @@ +--- +schema: aether.architecture-decision/v1 +id: adr-0012 +title: Require fresh prerequisite evidence for durable graph execution +kind: architecture-decision +status: proposed +accepted: null +owners: + - egohygiene +scope: + - flow +governed_by: + - architecture-decisions +supersedes: [] +superseded_by: [] +related: + - flow-architecture + - flow-roadmap + - adr-0011 +--- + +# ADR-0012 — Require fresh prerequisite evidence + +## Context and authority + +Flow #49 persists successful steps, but a historical success does not establish +that its current inputs, provider, authority, validator, or outputs still match. +Checking only a dependent step can miss stale upstream evidence. The maintainer +authorized the next dependency-ready checkpoint on 2026-09-26. #64 is the first +checkpoint under #31; this decision is proposed for review with its implementation. + +## Decision + +Add a closed, versioned `flow.run-assessment/v1` report bound to the exact saved +plan digest, current state digest, and snapshot sequence. Require one explicitly +identified current context for every planned step in plan order; reject missing, +extra, duplicated, or reordered entries. Reject a different expected plan and +corrupt retained history before producing a report. + +Visit steps in topological plan order. Freshly check each eligible step through +the existing context, input, validator, and artifact gates. A local stale boundary +invalidates that step. An invalidated prerequisite invalidates every descendant; +other non-reusable prerequisites block descendants. Record direct blockers in +declared dependency order and local stale boundaries as typed data. A blocked +step's future inputs need not exist yet and are not observed. Independent branches +remain eligible according to their own evidence. + +Assessment is read-only. Historical success and accepted checkpoints remain +intact even when currently invalidated for reuse. Serialized reports are ordinary +descriptive data, never capabilities or accepted execution evidence. Dependent +execution accepts the complete expected plan and context inventory and computes +a new assessment immediately before entering the existing durable launch path. +The single-step assessment and execution APIs refuse any step with dependencies. + +## Rationale and alternatives + +Checking only recorded dependency status misses changed upstream evidence. Taking +a caller-supplied report as permission would allow stale or forged readiness. +Rewriting succeeded records to pending would erase historical acceptance and +silently permit repeat effects. The bounded alternative keeps immutable history, +fresh eligibility, and explicit execution separate. Full inventory checking also +avoids implying that omitted branches were validated. + +## Limits, recovery, and compatibility + +The expected plan must match exactly. Selective reuse within that plan does not +authorize migration to a changed plan or rerunning an invalidated success. Retry +and abandonment continue to use explicit recorded recovery decisions. There is +no scheduler, automatic retry, artifact deletion, or exactly-once external-effect +claim. #65 owns the broad lifecycle corpus and #66 owns effect/cleanup proofs. + +The store retains ADR-0011's trusted local filesystem assumptions. Observations +are sequential, not an atomic snapshot of external files. Callers must keep inputs, +artifacts, and provider packages quiescent during assessment and execution; the +workspace lock does not contain external writers. A future concurrency guarantee +requires a stronger artifact/process boundary. + +The report contract is additive; saved plan and state schemas remain unchanged. +The changed validation implementation digest intentionally refuses reuse of +checkpoints accepted by a different implementation until a future explicit +compatibility path exists. Unknown report versions/fields are refused. Reports +retain identifiers and digests but no roots, configuration values, tokens, raw +provider messages, or source bytes. They require the same privacy handling as +the workspace metadata. + +## Validation and review triggers + +Exercise real hermetic providers, fresh reopen, branching and transitive stale +propagation, unchanged independent work, blocked dependencies, complete inventory +refusals, forged/stale report non-authority, and the single-step bypass refusal. +Assert no launch events or history changes on refusal and no repeated successful +launch. Repeat deterministic report comparisons in fresh roots. Revisit for +automatic scheduling, cross-plan reuse, concurrent mutation, or report authority. + +## Related artifacts + +Flow #11, #31, and #64; `flow-architecture`; +[`durable-state.md`](../../../integrations/durable-state.md); and +`flow.run-assessment/v1`. diff --git a/docs/integrations/durable-state.md b/docs/integrations/durable-state.md index 25c76cc..e2677af 100644 --- a/docs/integrations/durable-state.md +++ b/docs/integrations/durable-state.md @@ -2,8 +2,10 @@ Flow #49 adds a library API for persistent execution of fully prepared process steps. It builds on the exact subject, authority, transcript, and artifact gates. -There is no product CLI or graph scheduler in this checkpoint. #31 owns the -broader lifecycle scenario matrix; #53 owns the supported CLI. +Flow #64 adds fresh graph eligibility and safe caller-selected dependent +execution. There is no product CLI or graph scheduler. #31 continues through +#65 (the lifecycle corpus) and #66 (authority/effect and residual-state proofs); +#53 owns the supported CLI. ## Public entry points @@ -15,19 +17,27 @@ broader lifecycle scenario matrix; #53 owns the supported CLI. Plans contain fully pinned expected inputs; unresolved future outputs require later planning support. 3. Call `RunStore::create` with a new workspace path. Its parent must exist. -4. Call `execute` for a ready step with explicit secrets, cancellation, and event - adapters. The store persists intent before launch, then an accepted checkpoint - or a bounded typed failure. It holds its exclusive workspace lock throughout. +4. Supply a `RunStepContext { step_id, context }` for every step in plan order. + Call `execute_in_plan` with the expected plan, selected step ID, that complete + inventory, and explicit secrets, cancellation, and event adapters. It freshly + assesses the graph and executes only a ready step. The store persists intent + before launch, then an accepted checkpoint or a bounded typed failure. It + holds its exclusive workspace lock throughout. 5. Drop the handle and use `RunStore::open` after restart. `state()` exposes the recorded status without writing files or launching a process. -6. Use `assess` with the complete expected plan and fresh runtime context. - `Reusable` requires matching intent, validator identity, and fresh input/output - observations. It is eligibility evidence, not a reconstructed acceptance token. +6. Use `assess_run` with the complete expected plan and context inventory. + `Reusable` requires matching intent, validator identity, fresh input/output + observations, and freshly reusable prerequisites. The returned versioned + report describes eligibility without reconstructing acceptance or authority. + +The original `assess` and `execute` entry points remain supported for root steps +only. They return `DependencyEvidenceRequired` for any dependent step, regardless +of its recorded prerequisites. Use the graph APIs for dependent execution. `RunState::new` remains an in-memory model. `Orchestrator` and `LocalProcessRunner` remain useful low-level seams for unit tests and controlled integrations. Callers -requiring restart semantics must use `RunStore::execute` rather than treating a -low-level runner result as a persisted run. +requiring restart semantics must enter the durable coordinator instead of +treating a low-level runner result as a persisted run. ## Records and identity @@ -41,6 +51,7 @@ low-level runner result as a persisted run. | `flow.run-validation/v1` | Accepted evidence digests and validation implementation identity | | `flow.run-recovery/v1` | Explicit retry/abandon decision and uncertainty acknowledgement | | `flow.run-snapshot/v1` | Storage envelope and state integrity digest | +| `flow.run-assessment/v1` | Read-only graph eligibility bound to current state and plan digests | Every checkpoint binds the complete plan digest and its step context digest. That context includes provider identity/version/integrity, capability definition, @@ -60,6 +71,48 @@ object keys as compact UTF-8 JSON; arrays retain their declared order. The Rust semantic validators additionally enforce identity and state relationships that JSON Schema alone does not express. See `contracts/fixtures/state/` for examples. +## Fresh graph assessment + +`assess_run` verifies retained history and the exact expected plan before +inspecting current evidence. The context inventory must identify every step +exactly once in plan order. Missing, extra, duplicate, or reordered entries are +refused. Contexts for blocked descendants are required, but their future input +bytes are not observed until all prerequisites are reusable. + +| Eligibility | Current meaning | +| --- | --- | +| `ready` | Pending, matching current evidence, with all prerequisites freshly reusable | +| `reusable` | Historically accepted, with matching current local and prerequisite evidence | +| `invalidated` | Local evidence is stale, or a prerequisite is invalidated | +| `dependency-blocked` | A prerequisite is pending, awaiting recovery, abandoned, or itself blocked | +| `approval-required` | A running, failed, cancelled, or denied step needs explicit recovery | +| `abandoned` | The operator abandoned this step and current evidence still matches | + +Dependency classification takes precedence over local checks. Each report entry +retains the recorded status, a typed local stale boundary when checked, and all +non-reusable direct prerequisites in declared order. An invalidated prerequisite +takes precedence over other blockers. Follow `blocked_by` entries to the root +cause. A descendant's absent local reason means its local evidence was not checked. + +For a graph with `A → B → C` and independent `D`, stale A inputs invalidate A, B, +and C for reuse while D keeps its own eligibility. A historically succeeded step +stays succeeded in the immutable history. Assessment never rewrites checkpoints, +reruns work, deletes artifacts, or claims migration into a changed plan. + +Reports include plan/state digests and snapshot sequence, with no timestamp or +absolute roots. `RunAssessment::validate_against` checks shape and correlation +only; it cannot prove external evidence is still fresh. Saving, editing, or +deserializing a report grants no authority. `execute_in_plan` accepts no report: +it recomputes eligibility and enters the ordinary durable launch path only for +the caller-selected ready step. It returns typed `Ineligible` otherwise. + +Keep provider packages, sources, and artifacts quiescent during assessment and +execution. The workspace lock protects state coordination; these sequential +observations do not lock external writers or form an atomic filesystem snapshot. +Validation implementation changes invalidate old checkpoints conservatively; +no compatibility override or cross-plan migration is shipped here. See +[ADR-0012](../architecture/governance/decisions/ADR-0012-fresh-graph-assessment.md). + ## Commit and restart behavior A workspace contains `workspace.lock`, numbered immutable `*.json` snapshots, @@ -98,8 +151,9 @@ and acknowledgement of uncertain effects. Retry records a pending step without launching it. Operators must inspect and preserve/quarantine residual outputs and resolve surviving processes before an explicit new attempt. Flow never deletes or overwrites those artifacts automatically. A succeeded step cannot be retried -through this API. Automatic reuse scheduling, targeted downstream invalidation, -provider-native checkpoint restoration, and exactly-once external effects are +through this API. Invalidation reports describe current reuse eligibility without +authorizing reruns. Automatic reuse scheduling, cross-plan reuse/migration, +provider-native checkpoint restoration, and exactly-once external effects remain outside v1. ## Privacy, budgets, and retention @@ -138,6 +192,7 @@ The rationale and review triggers live in cargo test --lib state::store::tests --locked cargo test --test durable_state --locked cargo test --test hermetic_provider_kit durable_execution --locked +cargo test --test hermetic_provider_kit graph_recovery --locked python3 tools/validate_contracts.py ``` @@ -145,5 +200,10 @@ The portable store suite covers restart, abrupt process exit, concurrent opens, schema/corruption refusal, partial writes, and immutable history. Unit fault injection covers failures after pending-file creation, after file sync, and after rename. Real provider tests exercise the durable success/failure path, interrupted -acceptance, stale evidence, explicit recovery, and privacy. macOS/Windows CI runs -the portable store suite; Linux also runs the full provider matrix. +acceptance, stale evidence, explicit recovery, and privacy. Eight graph tests +cover deterministic fresh-root/reopen reports, transitive and branch invalidation, +all local identity boundaries, complete inventories, stale-report non-authority, +the single-step bypass, blocked future inputs, and preserved history. The +acceptance driver includes these test sources in its evidence identity and runs +them with `--all-targets`; the broader lifecycle receipt catalog belongs to #65. +macOS/Windows CI runs the portable store suite; Linux runs the full provider matrix. diff --git a/src/state/assessment.rs b/src/state/assessment.rs new file mode 100644 index 0000000..baf36ac --- /dev/null +++ b/src/state/assessment.rs @@ -0,0 +1,200 @@ +use serde::{Deserialize, Serialize}; + +use crate::{AcceptedArtifactSet, CancellationSignal, EventSink, SecretResolver}; + +use super::{ + DurableExecutionError, ProcessStepContext, ResumeEligibility, RunPlan, RunState, RunStepStatus, + RunStore, StateBoundary, StateError, digest, require, schema, +}; + +pub const RUN_ASSESSMENT_V1: &str = "flow.run-assessment/v1"; + +/// One explicitly identified current context, supplied in immutable plan order. +pub struct RunStepContext<'a> { + pub step_id: &'a str, + pub context: ProcessStepContext<'a>, +} + +/// Read-only eligibility evidence. Deserializing this report grants no authority. +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(deny_unknown_fields)] +pub struct RunAssessment { + pub schema_version: String, + pub plan_digest: String, + pub state_digest: String, + pub sequence: u64, + pub steps: Vec, +} + +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(deny_unknown_fields)] +pub struct RunStepAssessment { + pub step_id: String, + pub recorded_status: RunStepStatus, + pub eligibility: ResumeEligibility, + /// A local mismatch, absent when a prerequisite prevents local assessment. + pub stale_boundary: Option, + /// Non-reusable direct prerequisites, in the plan's dependency order. + pub blocked_by: Vec, +} + +impl RunAssessment { + /// Validate report structure and correlation, never the freshness of its claims. + /// + /// # Errors + /// Refuses unsupported schemas, wrong state identities, and contradictory + /// ordering, dependency, status, or explanation fields. + pub fn validate_against(&self, state: &RunState) -> Result<(), StateError> { + schema(&self.schema_version, RUN_ASSESSMENT_V1)?; + state.validate()?; + require( + self.plan_digest == state.plan_digest + && self.state_digest == digest(state)? + && self.sequence == state.sequence, + "assessment state identity", + )?; + require( + self.steps.len() == state.steps.len(), + "assessment inventory", + )?; + for (index, step) in self.steps.iter().enumerate() { + require( + step.step_id == state.steps[index].step_id + && step.recorded_status == state.steps[index].status, + "assessment step identity", + )?; + let (blocked_by, invalidated) = blockers(&state.plan, &self.steps[..index], index); + require( + step.blocked_by == blocked_by, + "assessment dependency evidence", + )?; + let expected = if !blocked_by.is_empty() { + require(step.stale_boundary.is_none(), "blocked local evidence")?; + blocked_eligibility(invalidated) + } else if let Some(boundary) = step.stale_boundary { + require(boundary != StateBoundary::Plan, "assessment exact plan")?; + ResumeEligibility::Invalidated + } else { + match step.recorded_status { + RunStepStatus::Pending => ResumeEligibility::Ready, + RunStepStatus::Succeeded => ResumeEligibility::Reusable, + RunStepStatus::Abandoned => ResumeEligibility::Abandoned, + _ => ResumeEligibility::ApprovalRequired, + } + }; + require(step.eligibility == expected, "assessment eligibility")?; + } + Ok(()) + } +} + +impl RunStore { + /// Freshly assess all steps in the exact saved plan without changing history. + /// + /// Non-reusable prerequisites prevent observing a descendant's future inputs. + /// Stale evidence invalidates descendants; other unresolved prerequisites block + /// them. Independent branches are assessed with their own current evidence. + /// External files must remain quiescent; observations are not an atomic snapshot. + /// + /// # Errors + /// Refuses corrupt history, changed plans, incomplete/reordered context inventories, + /// and malformed contexts. Local stale evidence becomes a typed report entry. + pub fn assess_run( + &self, + expected_plan: &RunPlan, + contexts: &[RunStepContext<'_>], + ) -> Result { + self.check_expected_plan(expected_plan)?; + if contexts.len() != self.state.steps.len() + || contexts + .iter() + .zip(&self.state.steps) + .any(|(current, saved)| current.step_id != saved.step_id) + { + return Err(StateError::ContextInventory); + } + let mut steps = Vec::with_capacity(contexts.len()); + for (index, current) in contexts.iter().enumerate() { + let (blocked_by, invalidated) = blockers(&self.state.plan, &steps, index); + let (eligibility, stale_boundary) = if blocked_by.is_empty() { + match self.assess_local(index, ¤t.context) { + Ok(eligibility) => (eligibility, None), + Err(StateError::Stale { boundary }) => { + (ResumeEligibility::Invalidated, Some(boundary)) + } + Err(error) => return Err(error), + } + } else { + (blocked_eligibility(invalidated), None) + }; + steps.push(RunStepAssessment { + step_id: current.step_id.to_owned(), + recorded_status: self.state.steps[index].status, + eligibility, + stale_boundary, + blocked_by, + }); + } + Ok(RunAssessment { + schema_version: RUN_ASSESSMENT_V1.to_owned(), + plan_digest: self.state.plan_digest.clone(), + state_digest: digest(&self.state)?, + sequence: self.state.sequence, + steps, + }) + } + + /// Reassess the complete graph, then execute one caller-selected ready step. + /// A saved or caller-edited assessment is never an input to this operation. + /// + /// # Errors + /// Refuses any step that is not freshly ready, including completed work, + /// stale ancestors, missing evidence, and unresolved dependencies. Execution + /// uses the same durable intent and acceptance path as root-step `execute`. + pub fn execute_in_plan( + &mut self, + expected_plan: &RunPlan, + step_id: &str, + contexts: &[RunStepContext<'_>], + secrets: &dyn SecretResolver, + cancellation: &dyn CancellationSignal, + events: &mut dyn EventSink, + ) -> Result { + let assessment = self.assess_run(expected_plan, contexts)?; + let index = self.state.step_index(step_id)?; + let eligibility = assessment.steps[index].eligibility; + if eligibility != ResumeEligibility::Ready { + return Err(StateError::Ineligible { eligibility }.into()); + } + self.execute_ready( + step_id, + &contexts[index].context, + secrets, + cancellation, + events, + ) + } +} + +fn blockers(plan: &RunPlan, assessed: &[RunStepAssessment], index: usize) -> (Vec, bool) { + let mut blocked_by = Vec::new(); + let mut invalidated = false; + for id in &plan.steps[index].depends_on { + // Validated plans only depend on preceding steps; reports follow plan order. + if let Some(step) = assessed.iter().find(|step| &step.step_id == id) { + if step.eligibility != ResumeEligibility::Reusable { + blocked_by.push(id.clone()); + invalidated |= step.eligibility == ResumeEligibility::Invalidated; + } + } + } + (blocked_by, invalidated) +} + +const fn blocked_eligibility(invalidated: bool) -> ResumeEligibility { + if invalidated { + ResumeEligibility::Invalidated + } else { + ResumeEligibility::DependencyBlocked + } +} diff --git a/src/state/execution.rs b/src/state/execution.rs index 49f3167..d60b447 100644 --- a/src/state/execution.rs +++ b/src/state/execution.rs @@ -1,5 +1,6 @@ use std::path::Path; +use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; use thiserror::Error; @@ -29,10 +30,12 @@ pub struct ProcessStepContext<'a> { pub bindings: &'a ArtifactBindingSet, } -#[derive(Clone, Copy, Debug, Eq, PartialEq)] +#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(rename_all = "kebab-case")] pub enum ResumeEligibility { Ready, Reusable, + Invalidated, ApprovalRequired, DependencyBlocked, Abandoned, @@ -170,37 +173,45 @@ impl ProcessStepContext<'_> { } impl RunStore { - /// Assess current evidence against the caller's expected complete plan. + /// Assess a root step against the caller's expected complete plan. /// /// `Reusable` is an eligibility decision about recorded acceptance, not a /// reconstructed `AcceptedArtifactSet` or permission to launch a provider. /// /// # Errors /// Refuses changed plans, configuration, inputs, subjects, authority, validation - /// implementations, and output evidence. This method writes no state. + /// implementations, and output evidence. Dependent steps require `assess_run`. + /// This method writes no state. pub fn assess( &self, expected_plan: &RunPlan, step_id: &str, context: &ProcessStepContext<'_>, ) -> Result { - self.ensure_usable()?; + self.check_expected_plan(expected_plan)?; + let index = self.state.step_index(step_id)?; + if !self.state.plan.steps[index].depends_on.is_empty() { + return Err(StateError::DependencyEvidenceRequired); + } + self.assess_local(index, context) + } + + pub(super) fn check_expected_plan(&self, expected_plan: &RunPlan) -> Result<(), StateError> { + self.validate_current()?; if expected_plan.digest()? != self.state.plan_digest { return Err(StateError::Stale { boundary: StateBoundary::Plan, }); } - let index = self.state.step_index(step_id)?; + Ok(()) + } + + pub(super) fn assess_local( + &self, + index: usize, + context: &ProcessStepContext<'_>, + ) -> Result { self.check_context(index, context)?; - let planned = &self.state.plan.steps[index]; - if planned.depends_on.iter().any(|id| { - self.state - .steps - .iter() - .any(|s| &s.step_id == id && s.status != RunStepStatus::Succeeded) - }) { - return Ok(ResumeEligibility::DependencyBlocked); - } let step = &self.state.steps[index]; match step.status { RunStepStatus::Pending => Ok(ResumeEligibility::Ready), @@ -293,11 +304,11 @@ impl RunStore { /// The store remains locked throughout the attempt. Any failure to commit the /// terminal result is an error; the prior running record remains uncertain. /// Existing completed steps are never executed again by this method. + /// Dependent steps require `execute_in_plan` with complete current contexts. /// /// # Errors /// Returns typed coordination, cancellation, process, observation, or acceptance /// failures. Raw provider errors remain caller-local and are never serialized. - #[allow(clippy::too_many_lines)] pub fn execute( &mut self, step_id: &str, @@ -309,6 +320,18 @@ impl RunStore { if self.assess(&self.state.plan, step_id, context)? != ResumeEligibility::Ready { return Err(StateError::Transition.into()); } + self.execute_ready(step_id, context, secrets, cancellation, events) + } + + #[allow(clippy::too_many_lines)] + pub(super) fn execute_ready( + &mut self, + step_id: &str, + context: &ProcessStepContext<'_>, + secrets: &dyn SecretResolver, + cancellation: &dyn CancellationSignal, + events: &mut dyn EventSink, + ) -> Result { if cancellation.is_cancelled() { self.cancel_pending(step_id)?; return Err(DurableExecutionError::CancelledBeforeLaunch); @@ -547,6 +570,7 @@ pub fn validation_implementation_digest() -> String { include_bytes!("model.rs"), include_bytes!("store.rs"), include_bytes!("execution.rs"), + include_bytes!("assessment.rs"), include_bytes!("../../Cargo.lock"), ]; let mut hasher = Sha256::new(); diff --git a/src/state/mod.rs b/src/state/mod.rs index 5d3fc30..31df310 100644 --- a/src/state/mod.rs +++ b/src/state/mod.rs @@ -4,10 +4,12 @@ //! Reopening never launches a provider or recreates opaque trust tokens. See //! `docs/integrations/durable-state.md` for storage and recovery limits. +mod assessment; mod execution; mod model; mod store; +pub use assessment::{RUN_ASSESSMENT_V1, RunAssessment, RunStepAssessment, RunStepContext}; pub use execution::{ DurableExecutionError, ProcessStepContext, RecoveryApproval, ResumeEligibility, validation_implementation_digest, @@ -15,13 +17,14 @@ pub use execution::{ pub use model::*; pub use store::RunStore; -use serde::Serialize; +use serde::{Deserialize, Serialize}; use serde_json::Value; use sha2::{Digest, Sha256}; use thiserror::Error; /// A precise boundary whose current evidence disagrees with saved intent. -#[derive(Clone, Copy, Debug, Eq, PartialEq)] +#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(rename_all = "kebab-case")] pub enum StateBoundary { Plan, Invocation, @@ -62,6 +65,12 @@ pub enum StateError { Transition, #[error("step identifier does not exist in this plan")] UnknownStep, + #[error("dependent steps require fresh evidence for the complete plan")] + DependencyEvidenceRequired, + #[error("current contexts must identify every step exactly once in plan order")] + ContextInventory, + #[error("fresh graph assessment refused execution: {eligibility:?}")] + Ineligible { eligibility: ResumeEligibility }, #[error("the store must be reopened after a failed commit")] ReopenRequired, #[error("workspace I/O failed during {operation}")] diff --git a/src/state/store.rs b/src/state/store.rs index 7104745..fc5412a 100644 --- a/src/state/store.rs +++ b/src/state/store.rs @@ -121,6 +121,15 @@ impl RunStore { } } + pub(super) fn validate_current(&self) -> Result<(), StateError> { + self.ensure_usable()?; + let (on_disk, _) = load_history(&self.root)?; + if on_disk != self.state { + return Err(StateError::Corrupt); + } + Ok(()) + } + pub(super) fn commit(&mut self, mut next: RunState) -> Result<(), StateError> { if self.poisoned { return Err(StateError::ReopenRequired); diff --git a/tests/durable_execution/mod.rs b/tests/durable_execution/mod.rs index 60d2212..8011b0e 100644 --- a/tests/durable_execution/mod.rs +++ b/tests/durable_execution/mod.rs @@ -333,12 +333,10 @@ fn durable_dependency_status_blocks_launch_and_abandonment_survives_reopen() { plan.steps[0].depends_on = vec![predecessor.step_id.clone()]; plan.steps.insert(0, predecessor); let mut store = RunStore::create(&fixture.workspace(), plan.clone()).unwrap(); - assert_eq!( - store - .assess(&plan, "step:inspect", &fixture.context()) - .unwrap(), - ResumeEligibility::DependencyBlocked - ); + assert!(matches!( + store.assess(&plan, "step:inspect", &fixture.context()), + Err(StateError::DependencyEvidenceRequired) + )); assert!( store .execute( diff --git a/tests/graph_recovery/mod.rs b/tests/graph_recovery/mod.rs new file mode 100644 index 0000000..7145d19 --- /dev/null +++ b/tests/graph_recovery/mod.rs @@ -0,0 +1,555 @@ +use super::{ + BINDINGS_LOCATOR, CAPABILITIES, INPUT_LOCATOR, KitFixture, NoSecrets, PACKAGE_LOCATOR, + PreparedLifecycleRun, WORKSPACE_LOCATOR, provider_binary_snapshot, +}; +use flow::{ + DurableExecutionError, EventKind, ExtensionEvent, NeverCancelled, PlannedStep, + ProcessStepContext, RUN_PLAN_V1, RecoveryApproval, ResumeEligibility, RunAssessment, RunPlan, + RunRecoveryAction, RunState, RunStepContext, RunStepStatus, RunStore, StateBoundary, + StateError, +}; +use std::fs; +use std::path::PathBuf; + +const IDS: [&str; 5] = ["step:a", "step:b", "step:c", "step:d", "step:e"]; +const DEPENDENCIES: [&[&str]; 5] = [&[], &["step:a"], &["step:b"], &[], &["step:b", "step:d"]]; + +struct Node { + kit: KitFixture, + prepared: PreparedLifecycleRun, + artifacts: PathBuf, +} + +impl Node { + fn context(&self) -> ProcessStepContext<'_> { + ProcessStepContext { + execution_root: self.kit.root.path(), + artifact_root: &self.artifacts, + resolved: self.prepared.resolved(), + invocation: &self.prepared.invocation, + subjects: &self.prepared.subject_lock, + authority: &self.prepared.authority, + bindings: &self.kit.bindings, + } + } +} + +struct Graph { + nodes: Vec, + plan: RunPlan, +} + +impl Graph { + fn new() -> Self { + let (bytes, name) = provider_binary_snapshot(); + // Real processes with independent source artifacts and explicit ordering + // dependencies. Separate roots make local stale evidence distinguishable. + let nodes: Vec<_> = (0..5) + .map(|index| { + let mut kit = KitFixture::new(CAPABILITIES[0], &bytes, &name); + kit.bindings.outputs[0].artifact_id = format!("artifact:graph-{index}"); + let artifacts = kit.root.path().join(WORKSPACE_LOCATOR); + fs::write( + artifacts.join(BINDINGS_LOCATOR), + serde_json::to_vec(&kit.bindings).unwrap(), + ) + .unwrap(); + let prepared = kit.prepare_lifecycle_named( + CAPABILITIES[0], + "success", + false, + &format!("graph-{index}"), + ); + Node { + kit, + prepared, + artifacts, + } + }) + .collect(); + let plan = RunPlan { + schema_version: RUN_PLAN_V1.to_owned(), + plan_id: "plan:graph-recovery".to_owned(), + run_id: nodes[0].prepared.invocation.run_id.clone(), + steps: nodes + .iter() + .enumerate() + .map(|(index, node)| { + PlannedStep::prepare( + IDS[index].to_owned(), + DEPENDENCIES[index] + .iter() + .map(|id| (*id).to_owned()) + .collect(), + &node.context(), + ) + .unwrap() + }) + .collect(), + }; + plan.validate().unwrap(); + Self { nodes, plan } + } + + fn contexts(&self) -> Vec> { + self.nodes + .iter() + .enumerate() + .map(|(index, node)| RunStepContext { + step_id: IDS[index], + context: node.context(), + }) + .collect() + } + + fn workspace(&self) -> PathBuf { + self.nodes[0].kit.root.path().join("graph state café") + } + + fn create(&self) -> RunStore { + RunStore::create(&self.workspace(), self.plan.clone()).unwrap() + } + + fn execute(&self, store: &mut RunStore, index: usize, events: &mut Vec) { + store + .execute_in_plan( + &self.plan, + IDS[index], + &self.contexts(), + &NoSecrets, + &NeverCancelled, + events, + ) + .unwrap(); + } + + fn assess(&self, store: &RunStore) -> RunAssessment { + let report = store.assess_run(&self.plan, &self.contexts()).unwrap(); + report.validate_against(store.state()).unwrap(); + report + } + + fn complete(&self, store: &mut RunStore) { + for index in 0..5 { + self.execute(store, index, &mut Vec::new()); + } + } +} + +fn eligibility(report: &RunAssessment) -> Vec { + report.steps.iter().map(|step| step.eligibility).collect() +} + +fn history(graph: &Graph) -> Vec<(PathBuf, Vec)> { + let mut files: Vec<_> = fs::read_dir(graph.workspace()) + .unwrap() + .map(|entry| { + let path = entry.unwrap().path(); + let name = path.file_name().unwrap().to_owned(); + // Never read a byte-range locked file on Windows. + let bytes = if name == "workspace.lock" { + Vec::new() + } else { + fs::read(path).unwrap() + }; + (PathBuf::from(name), bytes) + }) + .collect(); + files.sort(); + files +} + +#[test] +fn graph_reopen_preserves_success_and_never_repeats_a_completed_launch() { + use ResumeEligibility::{DependencyBlocked, Ready, Reusable}; + let mut deterministic_reports = Vec::new(); + for _ in 0..2 { + let graph = Graph::new(); + let mut store = graph.create(); + assert_eq!( + eligibility(&graph.assess(&store)), + [ + Ready, + DependencyBlocked, + DependencyBlocked, + Ready, + DependencyBlocked + ] + ); + let mut events = Vec::new(); + graph.execute(&mut store, 0, &mut events); + graph.execute(&mut store, 3, &mut events); + drop(store); + let mut store = RunStore::open(&graph.workspace()).unwrap(); + let partial = graph.assess(&store); + assert_eq!( + eligibility(&partial), + [ + Reusable, + Ready, + DependencyBlocked, + Reusable, + DependencyBlocked + ] + ); + deterministic_reports.push(serde_json::to_vec(&partial).unwrap()); + for index in [1, 2, 4] { + graph.execute(&mut store, index, &mut events); + } + assert_eq!( + events + .iter() + .filter(|event| event.kind == EventKind::PhaseStarted) + .count(), + 5 + ); + drop(store); + let mut store = RunStore::open(&graph.workspace()).unwrap(); + assert_eq!(eligibility(&graph.assess(&store)), [Reusable; 5]); + let before = history(&graph); + let mut repeated = Vec::new(); + for id in IDS { + assert!(matches!( + store.execute_in_plan( + &graph.plan, + id, + &graph.contexts(), + &NoSecrets, + &NeverCancelled, + &mut repeated + ), + Err(DurableExecutionError::State(StateError::Ineligible { + eligibility: Reusable + })) + )); + } + assert!(repeated.is_empty()); + assert_eq!(history(&graph), before); + assert!(store.state().steps.iter().all(|step| step.attempt == 1)); + } + assert_eq!(deterministic_reports[0], deterministic_reports[1]); +} + +#[test] +fn graph_stale_evidence_invalidates_only_affected_descendants_and_preserves_history() { + use ResumeEligibility::{Invalidated, Reusable}; + for boundary in [ + StateBoundary::Inputs, + StateBoundary::Artifacts, + StateBoundary::Provider, + StateBoundary::Configuration, + StateBoundary::Authority, + StateBoundary::Bindings, + StateBoundary::Invocation, + StateBoundary::Capability, + ] { + let mut graph = Graph::new(); + let mut store = graph.create(); + graph.complete(&mut store); + drop(store); + let store = RunStore::open(&graph.workspace()).unwrap(); + let before = history(&graph); + let node = &mut graph.nodes[0]; + match boundary { + StateBoundary::Inputs => fs::write( + node.artifacts.join(INPUT_LOCATOR), + b"changed private source", + ) + .unwrap(), + StateBoundary::Artifacts => { + fs::write(node.kit.output_path(), b"damaged output").unwrap(); + } + StateBoundary::Provider => fs::write( + node.kit.root.path().join(PACKAGE_LOCATOR).join("LICENSE"), + b"changed package", + ) + .unwrap(), + StateBoundary::Configuration => { + node.prepared + .invocation + .configuration + .values + .insert("private".to_owned(), "configuration-canary".into()); + } + StateBoundary::Authority => { + node.prepared.invocation.authorization.grants_digest = "9".repeat(64); + } + StateBoundary::Bindings => { + node.kit.bindings.outputs[0].locator = "outputs/changed.json".to_owned(); + } + StateBoundary::Invocation => { + node.prepared.invocation.cancellation_id = "cancel:changed".to_owned(); + } + StateBoundary::Capability => { + node.prepared.invocation.capability_id = "flow/changed".to_owned(); + } + _ => unreachable!(), + } + let report = graph.assess(&store); + assert_eq!( + eligibility(&report), + [Invalidated, Invalidated, Invalidated, Reusable, Invalidated], + "{boundary:?}" + ); + assert_eq!(report.steps[0].stale_boundary, Some(boundary)); + assert_eq!(report.steps[1].blocked_by, ["step:a"]); + assert_eq!(report.steps[2].blocked_by, ["step:b"]); + assert_eq!(report.steps[4].blocked_by, ["step:b"]); + assert!( + report + .steps + .iter() + .all(|step| step.recorded_status == RunStepStatus::Succeeded) + ); + assert_eq!(history(&graph), before); + let serialized = serde_json::to_string(&report).unwrap(); + for private in [ + "configuration-canary", + "changed private source", + graph.workspace().to_str().unwrap(), + ] { + assert!(!serialized.contains(private)); + } + } +} + +#[test] +fn graph_execution_rechecks_ancestors_and_refuses_single_step_bypass() { + use ResumeEligibility::{Invalidated, Ready}; + let graph = Graph::new(); + let mut store = graph.create(); + graph.execute(&mut store, 0, &mut Vec::new()); + let mut stale_report = graph.assess(&store); + assert_eq!(stale_report.steps[1].eligibility, Ready); + fs::write( + graph.nodes[0].artifacts.join(INPUT_LOCATOR), + b"changed after assessment", + ) + .unwrap(); + // A structurally valid forged claim remains mere data; execution has no report argument. + stale_report.steps[0].eligibility = ResumeEligibility::Reusable; + stale_report.validate_against(store.state()).unwrap(); + let before = history(&graph); + let mut events = Vec::new(); + assert!(matches!( + store.execute_in_plan( + &graph.plan, + IDS[1], + &graph.contexts(), + &NoSecrets, + &NeverCancelled, + &mut events + ), + Err(DurableExecutionError::State(StateError::Ineligible { + eligibility: Invalidated + })) + )); + assert!(matches!( + store.assess(&graph.plan, IDS[1], &graph.nodes[1].context()), + Err(StateError::DependencyEvidenceRequired) + )); + assert!(matches!( + store.execute( + IDS[1], + &graph.nodes[1].context(), + &NoSecrets, + &NeverCancelled, + &mut events + ), + Err(DurableExecutionError::State( + StateError::DependencyEvidenceRequired + )) + )); + assert!(events.is_empty()); + assert!(!graph.nodes[1].kit.output_path().exists()); + assert_eq!(history(&graph), before); + // Unrelated ready work still runs; the stale branch does not poison the whole plan. + graph.execute(&mut store, 3, &mut events); +} + +#[test] +fn graph_inventory_and_plan_identity_are_required_before_any_launch() { + let graph = Graph::new(); + let mut store = graph.create(); + let before = history(&graph); + let mut events = Vec::new(); + for mutation in 0..4 { + let mut contexts = graph.contexts(); + match mutation { + 0 => { + contexts.pop(); + } + 1 => contexts.swap(0, 1), + 2 => contexts[1].step_id = IDS[0], + 3 => contexts.push(RunStepContext { + step_id: "step:unknown", + context: graph.nodes[0].context(), + }), + _ => unreachable!(), + } + assert!(matches!( + store.execute_in_plan( + &graph.plan, + IDS[0], + &contexts, + &NoSecrets, + &NeverCancelled, + &mut events + ), + Err(DurableExecutionError::State(StateError::ContextInventory)) + )); + } + let mut changed = graph.plan.clone(); + changed.plan_id = "plan:changed".to_owned(); + assert!(matches!( + store.execute_in_plan( + &changed, + IDS[0], + &graph.contexts(), + &NoSecrets, + &NeverCancelled, + &mut events + ), + Err(DurableExecutionError::State(StateError::Stale { + boundary: StateBoundary::Plan + })) + )); + assert!(events.is_empty()); + assert_eq!(history(&graph), before); +} + +#[test] +fn graph_unresolved_dependencies_block_without_observing_future_inputs() { + use ResumeEligibility::{Abandoned, ApprovalRequired, DependencyBlocked, Ready}; + let graph = Graph::new(); + let mut store = graph.create(); + fs::remove_file(graph.nodes[1].artifacts.join(INPUT_LOCATOR)).unwrap(); + store.deny_pending(IDS[0]).unwrap(); + assert_eq!( + eligibility(&graph.assess(&store)), + [ + ApprovalRequired, + DependencyBlocked, + DependencyBlocked, + Ready, + DependencyBlocked + ] + ); + let before = history(&graph); + let mut events = Vec::new(); + assert!(matches!( + store.execute_in_plan( + &graph.plan, + IDS[1], + &graph.contexts(), + &NoSecrets, + &NeverCancelled, + &mut events + ), + Err(DurableExecutionError::State(StateError::Ineligible { + eligibility: DependencyBlocked + })) + )); + assert!(events.is_empty()); + assert_eq!(history(&graph), before); + store + .decide_recovery( + IDS[0], + &graph.nodes[0].context(), + RecoveryApproval { + decision_id: "decision:abandon-root".to_owned(), + action: RunRecoveryAction::Abandon, + acknowledge_uncertain_effects: true, + }, + ) + .unwrap(); + drop(store); + let store = RunStore::open(&graph.workspace()).unwrap(); + assert_eq!( + eligibility(&graph.assess(&store)), + [ + Abandoned, + DependencyBlocked, + DependencyBlocked, + Ready, + DependencyBlocked + ] + ); + assert_eq!( + graph.assess(&store).steps[4].blocked_by, + ["step:b", "step:d"] + ); +} + +#[test] +fn graph_assessment_rejects_history_corruption_even_on_an_open_handle() { + let graph = Graph::new(); + let store = graph.create(); + fs::write( + graph.workspace().join("00000000000000000000.json"), + b"corrupt", + ) + .unwrap(); + assert!(matches!( + store.assess_run(&graph.plan, &graph.contexts()), + Err(StateError::Malformed) + )); + drop(store); + assert!(RunStore::open(&graph.workspace()).is_err()); +} + +#[test] +fn graph_changed_validator_blocks_descendants_despite_recorded_success() { + let graph = Graph::new(); + let mut store = graph.create(); + graph.execute(&mut store, 0, &mut Vec::new()); + drop(store); + let path = graph.workspace().join("00000000000000000002.json"); + let mut value: serde_json::Value = serde_json::from_slice(&fs::read(&path).unwrap()).unwrap(); + value["state"]["steps"][0]["checkpoint"]["validation"]["implementation_digest"] = + "0".repeat(64).into(); + value["state_digest"] = super::digest_json(&value["state"]).into(); + fs::write(path, serde_json::to_vec(&value).unwrap()).unwrap(); + let store = RunStore::open(&graph.workspace()).unwrap(); + let report = graph.assess(&store); + assert_eq!( + report.steps[0].stale_boundary, + Some(StateBoundary::Validation) + ); + assert_eq!(report.steps[0].recorded_status, RunStepStatus::Succeeded); + assert_eq!(report.steps[2].eligibility, ResumeEligibility::Invalidated); + assert_eq!(report.steps[3].eligibility, ResumeEligibility::Ready); +} + +#[test] +fn graph_report_contract_refuses_unknown_or_contradictory_evidence() { + let state: RunState = serde_json::from_str(include_str!( + "../../contracts/examples/run-state.v1.example.json" + )) + .unwrap(); + let report: RunAssessment = serde_json::from_str(include_str!( + "../../contracts/examples/run-assessment.v1.example.json" + )) + .unwrap(); + report.validate_against(&state).unwrap(); + for mutation in 0..7 { + let mut value = report.clone(); + match mutation { + 0 => value.schema_version = "flow.run-assessment/v2".to_owned(), + 1 => value.state_digest = "0".repeat(64), + 2 => value.steps.clear(), + 3 => value.steps[0].eligibility = ResumeEligibility::Reusable, + 4 => value.steps[0].blocked_by.push("step:unknown".to_owned()), + 5 => value.steps[0].stale_boundary = Some(StateBoundary::Plan), + 6 => value.sequence += 1, + _ => unreachable!(), + } + assert!( + value.validate_against(&state).is_err(), + "mutation {mutation}" + ); + } + let mut value = serde_json::to_value(report).unwrap(); + value["steps"][0]["authorization_token"] = "private".into(); + assert!(serde_json::from_value::(value).is_err()); +} diff --git a/tests/hermetic_provider_kit.rs b/tests/hermetic_provider_kit.rs index 714b3d4..bba0513 100644 --- a/tests/hermetic_provider_kit.rs +++ b/tests/hermetic_provider_kit.rs @@ -6,6 +6,9 @@ mod scenario_matrix; #[path = "durable_execution/mod.rs"] mod durable_execution; +#[path = "graph_recovery/mod.rs"] +mod graph_recovery; + use std::collections::BTreeMap; use std::fs; use std::path::{Path, PathBuf}; @@ -2039,6 +2042,21 @@ impl KitFixture { capability: CapabilitySpec, mode: &str, lifecycle_control: bool, + ) -> PreparedLifecycleRun { + self.prepare_lifecycle_named( + capability, + mode, + lifecycle_control, + &format!("hermetic-{mode}"), + ) + } + + fn prepare_lifecycle_named( + &self, + capability: CapabilitySpec, + mode: &str, + lifecycle_control: bool, + identity: &str, ) -> PreparedLifecycleRun { let resolution = self.resolve_case( capability, @@ -2059,7 +2077,7 @@ impl KitFixture { ]); let invocation = ExtensionInvocation { schema_version: flow::EXTENSION_INVOCATION_V1.to_owned(), - invocation_id: format!("invocation:hermetic-{mode}"), + invocation_id: format!("invocation:{identity}"), run_id: format!("run:hermetic-{mode}"), phase: InvocationPhase::Execute, extension: InvocationExtension { @@ -2075,7 +2093,7 @@ impl KitFixture { protocol: resolved.execution_mode().protocol.clone(), }, input_artifacts: vec![InputArtifact { - artifact_id: INPUT_ID.to_owned(), + artifact_id: self.bindings.inputs[0].artifact_id.clone(), digest: self.bindings.inputs[0].expected_digest.clone(), }], expected_output_types: vec![capability.output_media_type.to_owned()], diff --git a/tools/run_acceptance_scenarios.py b/tools/run_acceptance_scenarios.py index 744c085..712e605 100644 --- a/tools/run_acceptance_scenarios.py +++ b/tools/run_acceptance_scenarios.py @@ -28,7 +28,7 @@ def source_identity(): paths = {"Cargo.toml", "Cargo.lock", "LICENSE", "tests/hermetic_provider_kit.rs", "tools/run_acceptance_scenarios.py"} paths.add("tests/durable_state.rs") - for directory in ["src", "contracts", "tests/scenario_matrix", "tests/durable_execution", "tests/fixtures", "tests/common"]: + for directory in ["src", "contracts", "tests/scenario_matrix", "tests/durable_execution", "tests/graph_recovery", "tests/fixtures", "tests/common"]: paths.update(str(path.relative_to(ROOT)) for path in (ROOT / directory).rglob("*") if path.is_file() and "__pycache__" not in path.parts) entries = {path: digest((ROOT / path).read_bytes()) for path in sorted(paths)} diff --git a/tools/test_durable_contracts.py b/tools/test_durable_contracts.py index 48a415a..2273c9c 100644 --- a/tools/test_durable_contracts.py +++ b/tools/test_durable_contracts.py @@ -47,6 +47,23 @@ def test_external_reference_preserves_closed_nested_objects(self): value["observation"]["source_bytes"] = "private" self.assertTrue(self.validate("run-artifact", value)) + def test_assessment_closed_vocabulary_and_bounded_dependency_inventory(self): + original = load_object(CONTRACTS / "examples/run-assessment.v1.example.json") + self.assertEqual(self.validate("run-assessment", original), []) + for mutation in ["version", "status", "private", "duplicate", "boundary"]: + value = copy.deepcopy(original) + if mutation == "version": + value["schema_version"] = "flow.run-assessment/v2" + elif mutation == "status": + value["steps"][0]["eligibility"] = "authorized" + elif mutation == "private": + value["steps"][0]["provider_message"] = "private canary" + elif mutation == "duplicate": + value["steps"][0]["blocked_by"] = ["step:a", "step:a"] + else: + value["steps"][0]["stale_boundary"] = "plan" + self.assertTrue(self.validate("run-assessment", value), mutation) + if __name__ == "__main__": unittest.main() diff --git a/tools/validate_contracts.py b/tools/validate_contracts.py index 15f333c..ace609c 100644 --- a/tools/validate_contracts.py +++ b/tools/validate_contracts.py @@ -1770,6 +1770,7 @@ def main() -> int: "flow.run-recovery/v1", "flow.run-state/v1", "flow.run-snapshot/v1", + "flow.run-assessment/v1", "flow.artifact-bindings/v1", "flow.artifact-observations/v1", "flow.artifact/v1",