diff --git a/.gitleaks.toml b/.gitleaks.toml index 9b0ee39f2..baef030cb 100644 --- a/.gitleaks.toml +++ b/.gitleaks.toml @@ -50,3 +50,11 @@ regexes = [ paths = [ '''(^|/)contracts/fixtures/random-stream-vectors/.*''', ] + +[[allowlists]] +description = "Issue #1350: exact synthetic DBOS workflow ID in the source-hashed experiment; not a credential. Other values, paths and scanner rules remain checked." +targetRules = ["generic-api-key"] +condition = "AND" +regexTarget = "secret" +regexes = ['''^dbos-r2-cancel$'''] +paths = ['''(^|/)docs/research/execution-architecture/experiments/probe_dbos\.py$'''] diff --git a/docs/decisions/adrs/README.md b/docs/decisions/adrs/README.md index 4cb2bbfca..e8123fb5b 100644 --- a/docs/decisions/adrs/README.md +++ b/docs/decisions/adrs/README.md @@ -145,10 +145,18 @@ adr-098-portable-artifact-requirement-satisfaction adr-099-participant-relative-predicate-opacity adr-100-participant-crossing-bisimulation adr-101-adversarial-participant-flow-control +adr-102-mixed-cross-backend-participant-control +adr-103-branch-aware-python-coverage-policy +adr-104-runtime-control-plane-architecture adr-105-recursive-partial-description-semantics adr-106-developer-package-and-artifact-management adr-107-artifact-promotion-and-release-admission adr-108-modular-participant-control-and-governed-effects +adr-109-participant-identity-and-objective-assignment +adr-110-reusable-mixed-control-policies-and-occurrences +adr-111-control-applicability-and-effect-decisions +adr-112-external-inject-triggering-and-execution +adr-113-reusable-execution-machinery ``` | ADR | Title | Status | Date | @@ -265,3 +273,4 @@ adr-108-modular-participant-control-and-governed-effects | [110](adr-110-reusable-mixed-control-policies-and-occurrences.md) | Reusable Mixed-Control Policies and Occurrences | accepted | 2026-09-22 | | [111](adr-111-control-applicability-and-effect-decisions.md) | Control Applicability and Effect Decisions | accepted | 2026-09-22 | | [112](adr-112-external-inject-triggering-and-execution.md) | External Inject Triggering and Execution | accepted | 2026-09-22 | +| [113](adr-113-reusable-execution-machinery.md) | Reusable Execution Machinery Under RAE Authority | accepted | 2026-09-23 | diff --git a/docs/decisions/adrs/adr-113-reusable-execution-machinery.md b/docs/decisions/adrs/adr-113-reusable-execution-machinery.md new file mode 100644 index 000000000..b8aaf819c --- /dev/null +++ b/docs/decisions/adrs/adr-113-reusable-execution-machinery.md @@ -0,0 +1,325 @@ +# ADR-113: Reusable Execution Machinery Under RAE Authority + +## Status + +accepted + +## Date + +2026-09-23 + +## Classification + +Classification: FM3 + +Required artifacts: the architecture selection for issue #1350, under +[ADR-104](adr-104-runtime-control-plane-architecture.md) and the +[#1348 supervision semantics](../../../specs/formal/runtime-control-plane/supervision.md). +The existing abstract supervision model supplies the semantic invariants; +the [bounded experiments](../../research/execution-architecture/experiment-report.md) +supply mechanism observations and counterexamples. + +Waivers: neither evidence set establishes a distributed refinement proof or +production conformance. This decision selects +components and integration rules; it adds no runtime dependency, published +carrier, selectable profile or executable guarantee. P3 remains unavailable. + +## Context + +RAE must drive ordinary CTFs, OT digital twins, AI security research and +sandboxes, product tests, and mixed IT/OT disaster recovery without giving each +backend a separate scenario interpreter. The local quickstart is one deployment, +not the product's ceiling. The evaluation's illustrative 15-tenant, +hundreds-of-nodes workload is a design case, not a measured capacity requirement. + +The incumbent logical mutation permit spans external calls and cannot provide +the supervision defined by #1348. A new task loop alone would leave RAE building +distributed dispatch, timers, recovery and operational tooling. Conversely, +adopting a framework's success, cancellation or retry semantics as RAE's own +would change authored meaning. The [comparison](../../research/execution-architecture/candidate-comparison.md) +evaluates these trade-offs rather than selecting by framework category. + +## Decision + +### 1. Reuse execution machinery without transferring responsibility + +Select **self-hosted Temporal Server and its Python SDK, backed by PostgreSQL** +for the distributed composition. Use Temporal's durable dispatch, task queues, +workflow history, timers, worker management and replay mechanisms. RAE-owned +code makes the semantic decisions through these mechanisms. Do not build a +replacement durable scheduler, broker, actor platform or workflow language. + +Preserve **P0's in-process library composition and P1/P2's existing local +transactional store and ownership boundaries**. Reuse Python/AnyIO concurrency +and SQLite there; do not require a Temporal service for an ordinary local run. +The same RAE semantic functions and contracts serve both compositions, through +small execution adapters. This is not permission to implement two sets of +workflow, retry, time, participant or outcome rules. Selecting remote workers +or PostgreSQL must never silently upgrade a P0/P1/P2 profile. + +Temporal Server/SDK and AnyIO use MIT licenses; PostgreSQL uses the PostgreSQL +License; SQLite is public domain. No paid software, managed Temporal service, +enterprise feature or cloud account is required by this architecture. Operating +machines, storage and support still costs resources. Pin supported versions, +transitive licenses, platform compatibility and artifacts through the existing +package/tooling policy when implementation adds dependencies. The experiment +versions are an evidence record, not a production support matrix. + +Sources: [Temporal license](https://github.com/temporalio/temporal/blob/main/LICENSE), +[SDK license](https://github.com/temporalio/sdk-python/blob/main/LICENSE), +[AnyIO license](https://github.com/agronholm/anyio/blob/master/LICENSE), +[PostgreSQL license](https://www.postgresql.org/about/licence/), +[SQLite copyright](https://sqlite.org/copyright.html), +[Temporal persistence](https://docs.temporal.io/temporal-service/persistence). + +### 2. Keep the ecosystem and runtime boundaries + +[Hub #3](https://github.com/OpenRAE/hub/issues/3) and #1348 govern the following +allocation; component reuse cannot alter it. + +| Owner | Retained responsibility | +| --- | --- | +| RAES | Semantics, portable contracts and conformance. | +| Shared RAE runtime | Admission, authored execution order and time interpretation, retry/continuation permission, supervised dispatch, conflicting-effect reservations, validated outcomes and requirement satisfaction. | +| LilRAE | Complete personal/local backend: realization, readiness, evidence, rollback and teardown. Its small default pack does not cap capability. | +| BigRAE | Complete organizational backend: organizational control plane, tenancy, authentication, policy, secrets, resource scheduling, audit and operations. These duties do not replace RAE's scenario semantics or operation audit. | +| Hub, Catalog, env-packs | Journey/sequencing/release gate; discoverable reusable assets; bounded environment pack, respectively. | +| Adapters and embedders | Translation only where required; deployment, process lifecycle and selected composition. | + +Backends execute concrete effects, establish observations and cessation, and +may legitimately refuse. RAE checks required capability and contextual +willingness, and decides whether the observed result satisfies the admitted +contract. A backend cannot weaken that contract; RAE cannot infer unobserved +reality. Resource availability never supplies permission for another attempt. + +### 3. Separate state authority, deterministic orchestration and effect workers + +The diagram shows the selected distributed composition, not today's P2 topology. +RAE components may be co-deployed, but external work must not occupy the state +authority's short mutation section or its reserved control capacity. + +```mermaid +flowchart TD + A["Authored plans and authorized supervisor"] --> G["RAE admission and control API"] + subgraph RAE["Shared RAE runtime authority"] + G --> O["One state writer per admitted target/run scope"] + W["RAE semantic workflow code"] -->|"bounded state requests"| O + O <-->|"atomic state, invocation claims and audit"| L["RAE operation ledger"] + E["Effect activity workers"] -->|"current invocation authorization"| O + C["Control and observation workers"] -->|"correlated evidence"| O + end + O -->|"stable start and control identities"| T["Temporal service and durable queues"] + W -->|"permitted commands and timers"| T + T -->|"workflow tasks"| W + T -->|"effect queue"| E + T -->|"separate control queue"| C + T --> H["Temporal PostgreSQL history and visibility"] + L --> P["Separate PostgreSQL ledger database and role"] + E <-->|"validated backend calls and results"| B["LilRAE, BigRAE or other backend"] + C <-->|"stop, refusal, observation, cessation"| B +``` + +The ledger is authoritative for portable operations, snapshots, idempotency, +invocation reservations and operational audit. Temporal history is authoritative +for engine progress. These are different facts, with no dual-write success +assumption. Only the scoped RAE state writer commits semantic state, using the +existing atomic-store contract. Effect workers have no direct ledger mutation +credentials. Temporal workflow workers call RAE semantic logic; nondeterministic +store, clock, secret and backend interactions use activities and the admitted +state boundary, never workflow replay code. + +Existing workflow and participant schedulers decide *which* work is admissible; +Temporal supplies delivery and execution capacity. BigRAE supplies organizational +resource policy. Do not copy those schedulers into a new queue service. + +| Interface | Required information and authority | +| --- | --- | +| Admission to ledger | Immutable actor, target/run, operation kind, validated request commitment, resolved author policy with provenance, contract version and required guarantees. | +| Engine command | Opaque operation/invocation references, expected generation and carrier version; no credentials, live resolver, planner object or pickle. | +| Worker to state authority | Authenticated service identity scoped to tenant, target/run, operation and invocation; current authorization and willingness before effect admission. An engine task token is not RAE authorization. | +| Backend execution/control | Scoped invocation/control identity, admitted constraints, capability/willingness disposition, budgets and correlated evidence. Control has its own authorized actor. | +| Result to state authority | Existing bounded decoding and native validation against a trusted predecessor; effect knowledge, cessation/residual evidence and final release gates where required. | + +These are interface obligations, not new wire shapes. Extend the owning +contracts through governed schema publication; current recovery observation +carriers cannot express every cessation, partial-effect or continuation claim. + +### 4. Contextual author policy governs failure, retries and fresh trials + +Adopt the [retry and scope design](../../research/execution-architecture/authored-retry-policy.md). +Authors can set scenario defaults, nested lexical scope defaults and explicit +local overrides, using the same specificity principle as open/closed scopes. +RAE resolves and validates one effective policy before admission and records its +origin. Capability, idempotency and persistence do not imply permission to retry. +The policy also determines failure response and validity consequences under the +authored context: invalidating an experiment, holding/reconciling a twin's time +and state, invoking an admitted OT safe-state procedure, or restoring an IT +resource are different decisions with different evidence. Neither an IT/OT +label nor a common exception type chooses that response automatically. Mixed +scenarios can select different policies by scope; admission must check their +shared-resource, state and time dependencies. Intentional in-world faults must +not be erased by generic infrastructure recovery. + +The choices include never repeating an effect, bounded authorized repetition +under an admitted safety condition, and terminating this trial before separately +admitting a fresh trial. A technically idempotent effect may still be forbidden +to repeat for experimental validity. Reset and compensation require their own +verified obligations. Engine redelivery, effect invocation, authored workflow +attempt and experimental trial are distinct identities and counters. + +Engine retries may implement a policy only when each dispatch passes the +current RAE gate and preserves its limits and identity. Otherwise use a +single-attempt effect activity and let RAE schedule a separately authorized +attempt through Temporal. Pure/idempotent bookkeeping can use bounded engine +retries independently. No automatic retry policy applies universally to +in-world effects. In the absence of an authored permission, do not authorize a +second effect; this fallback is not a ban on explicit scoped retry policies. + +### 5. Use a conservative ledger/engine protocol + +The [execution protocol](../../research/execution-architecture/execution-protocol.md) +specifies admission, start acknowledgement, one-use invocation claims, result +settlement, ownership loss and restore handling. The essential rule is that +replaying a command or recovering an engine never recreates a consumed effect +permission. Resolve ambiguous state commits by authoritative readback; resolve +ambiguous external effects by admitted observation and cessation evidence. + +No distributed transaction spans Temporal and the RAE ledger. Reuse PostgreSQL +transactions, uniqueness and revision checks for ledger changes, and Temporal's +stable workflow identity and duplicate-start controls for engine submission. +Do not implement another durable transport queue. A retained admitted command +can be re-submitted for delivery; its effect gate still decides whether an +invocation is permitted. If safe recovery cannot be established, preserve +quarantine and report indeterminacy rather than repeat an uncertain effect. + +Distributed owner transfer is conservative: a replacement may observe and +reconcile, but it cannot release the old owner's reservations merely because +its lease expired. PostgreSQL CAS rejects stale state writes; it does not fence +external work. Automatic active/active scope ownership and disconnected +multi-writer operation are excluded from this selection. + +### 6. Reserve supervision capacity and retain time authority + +Use separate effect and control worker pools and queues, with separately bounded +HTTP/IPC admission, connection pools, audit capacity and process resources. +Bound callbacks, observation, state commits and drain independently; propagate +remaining budgets through nested calls. A shared saturated database or failed +Temporal service remains a common failure boundary. Return bounded unavailable +or indeterminate results there, not successful cancellation. Backend-local +containment and independent emergency controls must be admitted when an authored +interruption guarantee requires them. They do not transfer scenario authority. + +Scenario time belongs to the existing clock/domain/segment model. Temporal +timers measure apparatus scheduling, not simulated physics or a paused scenario +clock. RAE maps semantic deadlines explicitly and rechecks the clock segment +and remaining budget on wake/recovery. Persist restart-interpretable time data, +not raw monotonic timestamps. If continuity cannot be established, refuse the +continuation requiring it. Tight control loops and co-simulation stepping belong +in suitable existing backend machinery under the admitted time contract; +Temporal is not a hard-real-time or safety controller. + +### 7. Bound federation, isolation, deployment and upgrades + +The initial distributed topology has connected trusted workers, one home +authority for each target/run and tenant-scoped services/credentials. Federation +means authorized cooperation between these homes, not shared mutable ownership. +Prefer separate tenant ledger databases/roles, namespaces and worker fleets; +use separate Temporal clusters for mutually distrustful administrative domains. +Namespaces and queue labels route work; authentication and authorization enforce +access. Configure Temporal's supported mTLS/claim-mapping/authorization surfaces; +its default permissive authorizer is unsuitable for untrusted access. +BigRAE owns that deployment policy; RAE still authorizes every operation and +supervisory request. [Temporal security](https://docs.temporal.io/self-hosted-guide/security) + +On a partition, a worker that cannot obtain current admission starts no new +effect. An already admitted effect may continue: quarantine the conflicting +scope until observation establishes its outcome. Do not promote another home +or silently fall back to a local profile. Air-gapped deployments can run the +selected services within the enclave or select a supported local composition; +offline multi-writer reconciliation is excluded. + +Workers run trusted integration code. Hostile CTF guests, AI-generated code, +models and product-under-test workloads execute behind backend-owned VM, +container or other admitted containment. A Temporal workflow sandbox is a +determinism mechanism, not that security boundary. Keep control/store credentials +outside those workloads. Validate sizes and types before SDK decoding reaches +domain consumers; reuse JSON ingress, closed carriers, environment bindings, +participant information-flow/final-sink checks and value-free error envelopes. +Use configured endpoints only. Secret references are resolved at authorized +execution; raw secrets stay out of history, search attributes, heartbeat data, +retry records, dashboards and audit. Stored outputs also obey disclosure policy. + +Deploy the service, its supported PostgreSQL persistence and visibility stores, +RAE's separate ledger, trusted workers and existing host supervision. A separate +search cluster or Kubernetes installation is not required by this decision. +Apply existing artifact locks, least-privilege roles, backups and restore drills. +Keep compatible workflow workers available for old histories; use supported +versioning and Continue-As-New at validated boundaries, carrying operation +identity, policy, remaining budgets, reservations and pending controls. History +rollover is not a new RAE attempt or trial. Refuse unsupported versions rather +than replay changed semantics. See [versioning](https://docs.temporal.io/develop/python/versioning) +and [Continue-As-New](https://docs.temporal.io/develop/python/continue-as-new). + +### 8. Keep the use cases broad + +| Use case | Composition and semantic requirement | +| --- | --- | +| Ordinary CTF | Small local profile or shared distributed deployment as needed; bounded setup/reset/teardown and optional retry defaults. No compulsory robotics, experiment wrapper or physical-OT safety stack. | +| OT digital twin | Existing simulation/plant-model backend, explicit time/fidelity/coupling contract and evidence. SimPy can schedule discrete events; it is not a physics solver or a universal twin. ROS actions are useful only where that backend already uses them. | +| AI security and sandboxes | Isolate hostile workloads from trusted workers, protect participant disclosure, preserve model/data/tool versions and trial identity. Ray may serve a backend's distributed compute need without owning RAE outcomes. | +| Product testing | Reuse the product's test/provisioning machinery at backend seams; preserve scenario assertions, reset evidence and exact policy. A flaky-test retry must remain visible; idempotency cannot turn a failed trial into an unreported retry. | +| Mixed IT/OT disaster recovery | Execute authored restoration dependencies and verify recovery through backend observations. Restoring orchestration records does not restore the world. Physical-state reconciliation, backend refusal and separate safety systems remain explicit; no universal rollback or safety certification is claimed. | + +### 9. Acceptance and implementation boundary + +The architecture selection is complete; its executable realization follows +governed contracts and tests. [#1360](https://github.com/OpenRAE/rae/issues/1360) +and [#1362](https://github.com/OpenRAE/rae/issues/1362) must consume the interface, +retry-policy and recovery obligations. Public authoring of inherited retry +defaults, their compilation and the fresh-trial link require explicit delivery; +they are not features of today's cleanup-policy carrier. A coordinated profile +additionally needs an API-404-C4 ADR/formal-model extension and tenant/owner +conformance before P3 can become selectable. This ADR does not make that change. + +The [canonical disposition](../../research/execution-architecture/decision-discussion.md) +maps API-404 C1-C4 and the retained workflow, time, trial, result-validation and +observability owners. API-404 remains ACTIVE for its existing guarantees. Design +and probe evidence is recorded as DOCUMENTS, without new fulfillment claims. + +Before deployment claims, verify the actual RAE/backend composition under +duplicate delivery, stale workers, lost commit acknowledgements, service/store +loss, saturated control paths, unauthorized messages, secret-bearing history, +code upgrades and independent backup restores. Test policy inheritance and +fresh-trial validity, not just Temporal status. Measure load using event rates, +payloads, durations and control latency; tenant/node counts do not size a system. + +## Alternatives Considered + +- **DBOS with PostgreSQL:** a credible free durable-workflow alternative with + distributed queues. Its probes do not show it incapable. Temporal is selected + for its explicit service/worker separation and integrated execution history, + routing and operational model across connected worker fleets. DBOS's embedded + model reduces service footprint; changing this selection needs new evidence + that the operational trade-off matters, not another general survey. +- **Celery/task queues or a custom AnyIO platform:** mature task delivery is + useful, but building durable scenario orchestration, recovery coordination and + history around a task queue leaves more generic machinery for RAE to maintain. + Keep AnyIO for the local/worker seam. Do not create a second distributed engine. +- **ROS actions, Ray actors, SimPy, BehaviorTree.CPP:** retain as complementary + backend mechanisms where required. None supplies the complete portable + operation contract, and no evidence justifies making all users install them. +- **One mandatory distributed service:** violates the local composition promise. + **Engine history as the operation ledger:** conflates execution and semantic + outcomes and loses the existing atomic snapshot/operation/audit boundary. + +## Consequences + +RAE reuses established free execution and storage components while retaining +its duties. Integration work remains: semantic adapters, governed invocation +carriers, PostgreSQL store/provider, authenticated workers and policy compilation. +These are domain boundaries, not a new general-purpose execution platform. + +The distributed deployment is heavier than local SQLite and requires database, +service, worker-version and backup operations. Conservative reconciliation can +leave scopes unavailable when effects cannot be observed. That loss of liveness +is explicit; framework restart cannot manufacture permission or cessation. diff --git a/docs/decisions/adrs/adr-index.yaml b/docs/decisions/adrs/adr-index.yaml index 4bff78905..5bf040cf3 100644 --- a/docs/decisions/adrs/adr-index.yaml +++ b/docs/decisions/adrs/adr-index.yaml @@ -643,3 +643,6 @@ adrs: - id: ADR-112 path: docs/decisions/adrs/adr-112-external-inject-triggering-and-execution.md pin: 65e54addfb3a2dba74941f3c8030dec8c7a39661cc5aebb456e59c03af120ea3 + - id: ADR-113 + path: docs/decisions/adrs/adr-113-reusable-execution-machinery.md + pin: 6d904befc0d37d44b07ed083d297c7a0e3ba0d520d79f9ddce1999477c78e91d diff --git a/docs/decisions/issue-1350-execution-architecture-preflight.md b/docs/decisions/issue-1350-execution-architecture-preflight.md new file mode 100644 index 000000000..dbc05d300 --- /dev/null +++ b/docs/decisions/issue-1350-execution-architecture-preflight.md @@ -0,0 +1,296 @@ +# Issue 1350 — Execution Architecture Preflight + +Date: 2026-09-23. Status: preflight guidance; no executor, dependency, or +architecture has been accepted here. Issue #1350 is the contract for the later +decision. Its prerequisite is the [#1348 decision](issue-1348-operation-lifecycle.md), +[ADR-104](adrs/adr-104-runtime-control-plane-architecture.md), and the +[supervision semantics](../../specs/formal/runtime-control-plane/supervision.md). +The issue's historical [research](https://github.com/OpenRAE/rae/blob/077f7d04/docs/research/runtime-refactor/research.md) +and [diagnosis](https://github.com/OpenRAE/rae/blob/077f7d04/docs/research/runtime-refactor/diagnosis.md) +are available at the pinned Git revision, not in this branch's working tree. +Their candidate list is exploratory, not an adoption decision. + +Subsequent selection: [ADR-113](adrs/adr-113-reusable-execution-machinery.md) +records the accepted decision. This preflight retains the constraints used to +evaluate it; statements about pending selection below describe preflight state. + +The existing [candidate evaluation](../research/execution-architecture/README.md) +and [decision discussion](../research/execution-architecture/decision-discussion.md) +are also exploratory; their Temporal/PostgreSQL recommendation is not accepted. +This preflight neither repeats those experiments nor selects their proposal. +The supplied #1350 issue assigns [API-404](../requirements/API-404/requirement.md). +Its C1–C4 clauses and runtime profile catalog constrain the selection; this note +is design guidance, not implementation evidence or a fulfillment claim. Related +requirements retain the dispositions recorded by #1348. + +## Decision boundary + +The selected machinery must implement the existing operation contract, not +become its semantic owner. RAE owns admission, one state writer, supervised +dispatch, conflicting-effect reservations, result validation, atomic terminal +publication, and honest recovery classification. Authored workflow, time, +attempt, compensation, cleanup, and trial rules decide what work is permitted. +Backends own concrete effects, contextual willingness, interruption and effect +observation. Embedders own process deployment. A control receipt, durable +checkpoint, process death, timeout, or store lease is not evidence that remote +effects stopped. P0 must remain an in-process library composition; P1/P2 keep +their declared persistence and single-worker boundaries; P3 is unavailable. + +[Hub #3](https://github.com/OpenRAE/hub/issues/3) assigns complete local and +organizational backend responsibilities to LilRAE and BigRAE; the #1348 decision +establishes their shared RAE runtime. RAE retains authority over what must happen, +when and why. Backend resource scheduling, realization, observations and +contextual limits remain backend authority. Cooperation, including legitimate +refusal, is expected. These responsibilities cannot be traded for easier upkeep. + +Reusable machinery may implement substantial runtime logic while RAE retains +these duties. Evaluate complete compositions, with explicit handling of their +inevitable limits. The scope of reuse need not be restricted in advance to a +worker adapter. These are candidate claims to test, not accepted choices: + +The existing decision discussion assumes 15 federated tenants, hundreds of +nodes and tens of agents per scenario; these counts are not capacity evidence. +Preserve local profile compatibility without treating it as a platform ceiling. +If distributed execution is selected, evaluate mature execution components +before constructing equivalent infrastructure from concurrency primitives. +API-404's missing coordinated/tenant guarantees require explicit design and +contract extensions; neither that deployment scope nor its fulfillment follows +from the issue's requirement assignment or a library choice. + +| Pattern and concrete seam | Useful guarantee from the candidate | Limitation against #1348; deployment/failure boundary | +| --- | --- | --- | +| Incumbent Python/AnyIO worker and bounded HTTP offload, with a short state-owner section | Already composes the core, HTTP adapter, and in-memory/SQLite stores without another service. | Current logical mutation permit spans backend work; `_ControlPlaneCallExecutor` and runtime close drain indefinitely. Read, audit and mutation paths all use `run_in_threadpool`; separate methods do not reserve control capacity. Thread cancellation cannot stop a blocked callback. A process crash loses P0 work; P1/P2 need startup reconciliation. This is a comparison baseline, not proof that a new event loop is needed. | +| Durable execution: [Temporal](https://docs.temporal.io/activity-execution) or [DBOS](https://docs.dbos.dev/architecture), implementing runtime orchestration and backend tasks | Reuse durable execution records, scheduling, queues and worker recovery. Runtime-owned code specifies the semantic decisions using engine APIs. | Recovery must obey authored retry/observation rules; cancellation still requires backend cooperation. Temporal adds a service and persistence; DBOS Python supports SQLite or Postgres. Distinguish engine progress from portable operation state and specify their reconciliation. Evaluate complete compositions, not only per-operation opt-in. | +| Robotics/action protocol, e.g. [ROS 2 actions](https://design.ros2.org/articles/actions.html), at a backend adapter | Goal acceptance, feedback, result and a distinct cancel request are useful for a backend that already exposes them. | Action `CANCELING`/`CANCELED` are not RAE operation states or independent effect/cessation proof. Middleware, executor and backend deployment are extra dependencies; no ROS stack is required of non-robotic backends. Mapping must validate correlated backend evidence, not copy ROS states into portable DTOs. | +| Simulation/event scheduler, e.g. [SimPy events](https://simpy.readthedocs.io/en/4.1.1/api_reference/simpy.events.html), at semantic-time scheduling | Deterministic event ordering and cooperative process interrupts can serve a simulation backend. | A SimPy interrupt is delivered to its process, not to a blocked native/remote call. Its clock cannot replace the authored time-domain authority or apparatus monotonic deadline; no durable external-effect recovery follows from the event queue. | +| Actor mailbox/supervision, e.g. [Akka typed supervision](https://doc.akka.io/libraries/akka-core/current/typed/fault-tolerance.html), or a small in-process serialized owner | Mailbox serialization and explicit stop/restart policies can isolate local failures. | Mailbox saturation can strand urgent supervision; restarting an actor can abandon an active external effect. A separate framework/runtime adds ownership and deployment complexity. Reuse only if it preserves the existing store writer, reserved control capacity and quarantine; actor restart is not operation retry. | + +Behavior trees, statecharts, game loops and reconciliation controllers may help +particular authored or backend-local control, but none by itself supplies the +operation store, effect fencing, or authorization boundary. A desired-state +reconciler that re-applies effects is especially unsafe after an indeterminate +result. Mapping authored workflow states into engine steps must preserve their +meaning and permissions, including after recovery. + +The following is the required logical boundary from #1348, not a selected +framework, new public interface, or claim about today's executable supervision: + +```mermaid +flowchart LR + A[Authorized caller and admitted authored plan] --> G[Existing admission and policy gates] + G --> O[RAE state authority: claim, reserve, settle] + O <-->|atomic state, operation, audit and CAS| S[Admitted control-plane store] + O -->|scoped invocation; release state permit| E[Execution machinery and effect workers] + G -->|bounded supervisory admission| O + O -->|reserved control capacity| C[Control and observation workers] + E --> B[Registered backend: effects and willingness] + C --> B + B -->|correlated result and cessation evidence| V[Existing result and disclosure validation] + V -->|generation and revision checked settlement| O + E -. candidate-specific progress and reconciliation .-> H[Optional engine history] + H -. no independent operation-state authority .-> O +``` + +RAE retains the effect reservation while the state permit is free. Workers have +no direct store-write authority. Engine history, if present, is an additional +failure and disclosure boundary; the diagram asserts no atomic transaction +between it and the RAE store. + +## Existing interfaces and whole-repo gates + +Python package paths below are relative to `implementations/python/packages/`; +unqualified module names refer to `raes_runtime`. + +| Layer | Canonical incumbent and required fit | +| --- | --- | +| Contracts and semantic validation | `raes_contracts/operation_lifecycle.py`, `raes_contracts/contracts/operation_carriers.py`, closed `ContractModel` carriers, the FM3 model, authored `raes_contracts/workflow/`, `TimeCoordinator`, and SCE-007 trial/attempt/cleanup models own states, identities and permission. Reuse the SDL parser/compiler, plan admission and `backend_input_contracts.py` before dispatch, and native workflow/result/time validators on return. An executor gets an admitted command, not an alternate schema or validation policy. | +| Component admission | `ControlPlaneOptions`/`ControlPlaneConfiguration`, `control_plane_profiles.py`, `RuntimeTarget`, registry shape checks, `registry_target_validation.py`, and backend manifests distinguish profile, installed component, operation-kind support and contextual willingness. If a candidate needs an optional worker or backend control method, declare and shape-check it here; recheck willingness after queueing. No implicit profile upgrade or fallback. | +| State and persistence | `RuntimeMutationAuthority`, `RuntimeLifecycleMixin`, `RuntimeDurabilityMixin`, `ControlPlaneStore`/`AtomicControlPlaneStore`, `ControlPlaneStoreCommitAdapter`, in-memory and SQLite stores, revision CAS, strict codecs/migrations, `RuntimeOwnerLease` and path hardening own claim, readback, quarantine, audit and publication. A scheduler must neither hold the state permit across a callback nor release the effect reservation on local timeout. Engine history governs execution progress; its reconciliation with RAE operation records must preserve one semantic authority. Distributed storage and owner fencing need explicit contract extensions. | +| Backend and failure isolation | `backend_calls.py`, `_validated_backend_result()`, `backend_result_diagnostics.py`, `control_plane_recovery.py` and recovery observer own isolated inputs, trusted predecessor and recovery classification. Extend late-result admission through those boundaries; workers need scoped invocation/generation identities. Failure or malformed return must preserve portable state without claiming absent effects. Reserve target/run-wide conflicting effects unless admitted independence proves narrower isolation. | +| Security and HTTP validation | `ControlPlaneSecurityConfig.strict_defaults()`, `_ControlPlaneApiAuth`, `ControlPlaneIdentity`, `control_plane_plan_authorization.py`, `operation_admission_context()`, `RequestSizeLimitMiddleware`, API Pydantic models and bounded `_ControlPlaneCallExecutor` are the gates. Any new control route reauthorizes target/run/subject and operation scope, records its own supervisor actor, and retains bounded admission even when effect workers are saturated. IDs, receipts and engine task tokens do not authorize a RAE operation; protect engine tokens according to their native privileges too. | +| Secrets and host exposure | `SecretReferenceId`, runtime fact binding/dispatch, `raes/runtime_environment.py`, stateful resource projections and planner admission own closed env shapes, freshness, sensitivity and scoped resolution. Keep resolved credentials ephemeral, out of checkpoints, logs, audit, HTTP errors and worker payloads unless an existing validated carrier permits them. No token in process argv, inherited broad environment, request-supplied command or new shell boundary. TLS, proxy header stripping, filesystem permissions and process supervision stay with the embedder. A local process kill does not fence a remote job. | +| Observability and error envelopes | Existing `Diagnostic`/`DiagnosticModel`, `portable_diagnostic_payload()`, `AuditEvent`, rejection audit, module loggers, `control_plane_health.py`, API `_conflict_detail()` and `_operation_routes.py`'s global 422/500 handlers own value-free diagnostics and redacted responses. Keep operational readiness, participant observations, workflow history and experiment evidence distinct. Preserve denial status and headers if audit fails. Do not surface native exception text, request input, paths, credentials, traceback or engine payloads. | +| Participant scheduling and disclosure | `participant_scheduler_concurrency.py`, `participant_scheduler_concurrent_dispatch.py`, `participant_resource_budgets.py` and `participant_resource_reservation.py` already admit bounded batches, semantic independence and shared resource use. Preserve their accounting, ordering and indeterminate settlement instead of adding a competing scheduler. `participant_crossing_*`, `participant_control_*`, `participant_opacity_enforcement.py` and `participant_flow_sink.py` retain selected actor/subject, state-cut and final-sink checks before effects or disclosure. An engine's serialization/history is an additional disclosure boundary, not a reason to bypass these checks. | +| Package and deployment boundaries | ADR-036 and `tools/policy/adr_policy.yaml` own imports: backend protocols depend on contracts; runtime must not import concrete backend, CLI or MCP packages; MCP is an authoring surface. Keep engine-specific SDK objects out of portable contracts and lower layers. Runtime dependencies belong in `implementations/python/pyproject.toml` and its `uv.lock`; experiment tooling is distinct from shipped dependencies. Development executables/images use ADR-106, `implementations/tooling/artifacts.lock.json` and `tools/check_tooling_artifact_policy.py`. Record license, supported Python/OS versions, service/storage versions, migration, backup/restore and worker drain obligations before adoption. | +| Publication and verification | `.ground-control.yaml`, `.gc/plan-rules.md`, `noxfile.py`, repo policy, ADR-059 pins, requirement governance, and ADR-009/061 schema manifest and `schema_bundle()` are the gates. A portable carrier change needs its published schema, ledger, compatibility checks and strict migration together. Local verification stays targeted; CI owns full suites. | + +At a new serialized worker boundary, reuse `raes_contracts/json_ingress.py` +for bounded, duplicate-rejecting, finite JSON parsing before closed +`ContractModel` validation; preserve `runtime_value_limits.py` aggregate bounds +and domain validators. SDK deserialization or a matching Python type alone is +insufficient. Revalidation across a trust boundary is necessary; duplicating +the schema or its validation rules in an engine DTO is not. Translate SDK +failures at the adapter into the existing diagnostic/conflict vocabulary without +losing the distinction between denied admission, uncertain effects and uncertain +store acknowledgement. SDK failures must not leak through logs or error causes. + +Environment binding must pass **both** `raes/runtime_environment.py` and +`raes_processor/planner/stateful_admission.py`: literal `value` and `value_from` +are exclusive; generated output cannot claim `operator_secret`; delivery mode, +exact projection keys, node/output identity, sensitivity and consumer match must +agree. Reuse `raes_processor/planner/prepared_node_projection.py` and +`raes_runtime/backend_account_credentials.py` on the result path. A generic +engine environment dictionary cannot replace these shapes. Reuse +`runtime_fact_binding_policy.py` and `RuntimeFactDispatchCommand` for scoped, +fresh, protected-sink resolution; their in-process one-shot object is not a +durable deduplication record. Service connection credentials are embedder +configuration, not authored scenario environment. +Keep durable tasks reference-based and resolve credentials at the authorized +execution boundary; include engine history, search attributes, heartbeat data, +retry/dead-letter records and dashboards in the disclosure review. Redacting the +HTTP response alone does not protect those sinks. + +Connection endpoints must come from trusted composition, not an authored plan +or worker message. Reuse `raes_contracts/uri_safety.py` for portable URI fields; +its rejection of embedded credentials is not network authorization, TLS +verification or an SSRF defense. Selected clients still need authenticated +endpoints and explicit deployment trust. Child workers must not inherit owner +lease/store handles, broad credentials or live resolver objects. Process +isolation needs bounded IPC, explicit resource limits and owned shutdown; it +does not make hostile backend code safe or contain remote effects. + +The dependency seam must be parameterized by **operation kind and admitted +guarantees**, with a typed worker/backend capability and per-invocation control +identity. This lets a future backend offer cooperative stop, observation or +continuation without forcing every P0 backend to install an execution service. +Operational queue, callback, observation, commit and drain budgets belong in +typed runtime configuration; authored semantic time remains with its existing +clock/domain/segment model. Version any durable command or evidence carrier +through the contract publication rules, and define upgrade behavior before +using framework checkpoints for recovery. + +The current `RecoveryObservationResult` carries absent/applied/indeterminate +classification; it does not establish general partial-effect, cessation or +continuation evidence. Extend its owning contracts through governed publication +where required; do not hide stronger claims in diagnostics or engine status. + +The exact credential-sensitive retry proof in `control_plane_operation_context.py` +and `control_plane_admission.py` is deliberately ephemeral. A public request +commitment after restart cannot reauthorize credential-bearing input or supply +replay material. `RuntimeManager` is a separate direct-execution facade; it +cannot dispatch into a target/run quarantined by the control plane. + +## Open boundaries for a distributed composition + +These obligations apply to any candidate, including the existing proposal; they +are decision constraints, not a delivery sequence or a new protocol definition. + +- **Admission across a worker boundary.** The current trusted-embedder identity, + process-local planner digest registration and injected resolver/service objects + are not portable authorization. Define how a worker establishes the admitted + operation, actor, target/run, generation and current authority before an effect. + Do not pickle a live control plane, resolver, callback or mutable snapshot to + preserve process-local trust. Reuse bounded contract decoding and native result + validators on return. Unsupported carrier/worker versions fail closed. +- **Ledger and execution history.** Account for loss between claim, enqueue, + invocation, effect, result receipt, terminal commit and engine acknowledgement. + Preserve one semantic writer and atomic snapshot/operation/audit publication; + no shared database name implies a transaction across two systems. Duplicated + or out-of-order delivery must resolve through the admitted identity and + authoritative readback without another invocation or rewritten terminal parent. + Retention, restore and worker upgrades must preserve deduplication and the + evidence needed to reconcile still-active effects; history expiry is not + permission to repeat them. +- **Replay and identity.** Distinguish an engine workflow/run/task/attempt from + a RAE operation, execution generation, authored workflow attempt and trial + `run_id`. Engine reset, retry, continuation, signal redelivery or manual + redrive cannot allocate a fresh effect permission. Deterministic orchestration + must not directly call mutable stores, clocks, secret resolvers or backends; + map those interactions through the selected engine's supported boundary and + existing RAE admission. A replayed permission or willingness result is + historical evidence, not current dispatch authority. Code upgrades and + history rollover must retain operation identity, remaining budgets, + reservations and pending supervision without re-executing effects. +- **Ownership and federation.** `RuntimeOwnerLease`, its process/fork checks, + `require_single_worker_configuration()` and local private store paths remain + the P1/P2 boundary. A remote store or worker fleet needs explicit scope binding, + ownership transfer, stale-worker handling and partition behavior before it can + claim coordination. Tenant/namespace/queue names are routing labels, not access + control. Define authenticated service/worker identities, least-privilege store + and secret access, trust and credential revocation, and authorization of + callbacks/observations. Neither owner CAS nor a network partition proves that + the previous worker or backend stopped; retain exclusion when that is unknown. +- **Bounded operation.** State how effect, control, audit and store capacity stay + available under saturation and service failure. Separate queues alone do not + isolate shared thread pools, database connections, locks, CPU or memory. Reuse + resource accounting, require finite stage budgets, and preserve remaining + budgets across nested calls; scenario pause cannot pause supervision. Scope + endpoint, queue and worker settings to the selected composition through typed + configuration, with no implicit P3 upgrade or hard-coded tenant count. OS + process termination and remote backend containment require separate evidence. + +## Evidence required before selecting a component + +The retained [experiment report](../research/execution-architecture/experiment-report.md) +reports worker loss against an in-memory Temporal development server and +synthetic effect witnesses. That is not evidence for PostgreSQL persistence, +service loss, federation or end-to-end RAE supervision. Its invocation marker +is experiment scaffolding, not a reusable production fence. Review the retained +evidence for each named uncertainty before commissioning another experiment; +do not infer guarantees from the proposed stack diagram. + +Use bounded, instrumented experiments at each candidate's proposed seam, +including the incumbent. Inject a blocking callback, a cooperative and a +refusing cancel implementation, a crash after an external effect but before +acknowledgement, and a duplicate delivery. Observe backend effects independently +of futures, engine status and RAE store records. Bound each experiment and +record whether control remains reachable under worker/queue saturation, who +still owns the effect reservation, what persists across restart, and whether a +second effect is possible. Include cancellation-versus-completion and unknown +commit-ack races. A cancelled future, heartbeat timeout or durable retry alone +does not pass the interruption or duplicate-effect gate. If a candidate cannot +expose cessation/effect evidence, the honest result is retained quarantine or +`INDETERMINATE`, subject to admitted requirements. + +Build on the existing acceptance oracles, rather than replacing them with engine +status assertions: `test_issue_1348_operation_supervision.py` (abstract semantics), +`test_issue_1181_unified_control_plane_mutations.py` (one writer), +`test_issue_1187_control_plane_process_loss.py` (crash/no-replay), +`test_issue_1092_control_plane_crash_consistency.py` (atomic cuts/readback), +`test_issue_1187_control_plane_security_conformance.py` (identity and redaction), +and `test_participant_concurrent_batch_reservations.py` (reservation settlement), +all under `implementations/python/tests/`. Env/result regressions already have +`test_issue_1074_generated_artifact_env_consumers.py`, +`test_issue_1204_prepared_credentials.py` and +`test_issue_1003_final_sink_flow_enforcement.py`. Existing tests establish their +stated local contracts; they do not certify a distributed design. Record exact +candidate versions/configurations, independent effect observations and untested +failure windows. Also cover stopped-owner replacement while its backend still +acts, service/store unavailability, unauthorized completion/control messages, +credential disclosure through engine history, and incompatible worker upgrades +when claiming those boundaries. No fresh experiment or runtime test is required +for this documentation-only preflight. + +The final #1350 decision should publish **one** accepted architecture decision +and a component/interface diagram identifying state writer, worker, backend, +store and supervisor control flow; record measured guarantees, limitations, +dependency/deployment cost and alternatives. Selection is not ready until the +chosen composition resolves worker authorization, ledger/engine reconciliation +and stale-owner/partition behavior, or explicitly excludes that deployment. +These are decision gates, not an implementation sequence. Disposition the +canonical owners +API-404 (C1–C4), API-402/403, RUN-300/304/316, DSL-113, SEM-203/204/227–229, +SCE-006/007, EXP-706/712, SCE-002 and ASR-532 by reference to their existing +`docs/requirements//requirement.md` files. Do not copy their clauses into +a second catalog. An ADR-104 amendment or new ADR must obey ADR-059's amendment +and pin rules. Acceptance is a design decision, not a claim that executable +supervision, recovery or backend containment has shipped. + +Use `RAES_REQUIREMENT_UID=API-404` and the repository-backed requirement +governance configured in `tools/policy/requirement_order.yaml`. The eventual +decision needs canonical `DOCUMENTS` traceability; add `IMPLEMENTS`/`TESTS` +only for their actual code, spec or evidence claims. Do not change requirement +status or claim a new profile from this preflight or candidate experiments. + +## Non-goals and anti-patterns + +No runtime code, new schema, dependency, deployment service or accepted ADR is +created by this preflight. The later implementation must not add competing +operation authorities, unauthorized worker writes, duplicate exception hierarchy, +uncontrolled replay/compensation, free-form policy metadata, mandatory ROS, +or P3 coordination by implication. P0 stays usable without a durable service; +a distributed composition may require one explicitly. Do not conflate operation +with workflow/attempt/trial, scheduling with semantic time, CAS with external +fencing, audit with experiment evidence, or a backend cancel acknowledgement +with effect cessation. Physical-OT protections remain selectable future +backend capabilities, not universal requirements or certification claims. diff --git a/docs/requirements/API-404/requirement.md b/docs/requirements/API-404/requirement.md index aafbe21bc..c398db5f4 100644 --- a/docs/requirements/API-404/requirement.md +++ b/docs/requirements/API-404/requirement.md @@ -6,7 +6,7 @@ type: FUNCTIONAL priority: MUST wave: 1 created_at: 2026-04-03T05:55:58.825305Z -updated_at: 2026-09-22T00:00:00.000000Z +updated_at: 2026-09-23T00:00:00.000000Z --- # API-404 — Secure, Durable, And Idempotent Control-Plane Semantics @@ -104,6 +104,13 @@ identifies those implementation gaps and the retained canonical requirements. ## Traceability +- DOCUMENTS → GITHUB_ISSUE `1350` (Execution architecture selection; no new executable profile claim) +- DOCUMENTS → ADR `docs/decisions/adrs/adr-113-reusable-execution-machinery.md` (Reusable machinery, retained RAE authority and deployment boundaries) +- DOCUMENTS → DOCUMENTATION `docs/research/execution-architecture/execution-protocol.md` (Ledger/engine reconciliation, scoped worker authorization and conservative ownership recovery design) +- DOCUMENTS → DOCUMENTATION `docs/research/execution-architecture/authored-retry-policy.md` (Contextual author failure/retry policy, scoped defaults and fresh-trial distinctions; public support requires implementation) +- DOCUMENTS → DOCUMENTATION `docs/research/execution-architecture/experiment-report.md` (Bounded mechanism observations, not distributed conformance evidence) +- TESTS → TEST `implementations/python/tests/test_issue_1350_fixture_secret_scan.py` (Design-evidence publication checks: retained experiment source hashes match their record, and the exact synthetic fixture exception preserves detection of other values, paths and rules through the real-scanner integration lane; no runtime conformance claim) + - DOCUMENTS → GITHUB_ISSUE `1348` (Operation supervision decision; no new executable profile claim) - DOCUMENTS → DOCUMENTATION `docs/decisions/issue-1348-operation-lifecycle-preflight.md` (Supervision architecture guardrails) - DOCUMENTS → DOCUMENTATION `docs/decisions/issue-1348-operation-lifecycle.md` (Decision, requirement dispositions and implementation boundaries) diff --git a/docs/research/execution-architecture/README.md b/docs/research/execution-architecture/README.md new file mode 100644 index 000000000..8208ca51f --- /dev/null +++ b/docs/research/execution-architecture/README.md @@ -0,0 +1,53 @@ +# Execution architecture — issue #1350 + +The accepted decision is [ADR-113](../../decisions/adrs/adr-113-reusable-execution-machinery.md): +self-hosted Temporal and PostgreSQL for distributed execution, existing library +and SQLite compositions for local profiles, with RAE retaining semantic authority. +No paid software is required. This delivery selects and documents the design; +it does not deploy a service, add runtime dependencies or enable P3. + +The scope includes CTFs, OT twins, AI security research/sandboxes, product testing +and mixed IT/OT disaster recovery. Their failure responses and validity rules +can differ, including within one scenario. Authors select contextual policies, +scenario defaults, nested defaults and local overrides; engine retries cannot +supply those decisions. + +## Reading order + +1. [ADR-113](../../decisions/adrs/adr-113-reusable-execution-machinery.md): selection, + component/interface diagram, responsibilities, deployment and consequences. +2. [Authored failure/retry policy](authored-retry-policy.md): contextual responses, + scoped defaults, attempts, fresh trials and the implementation boundary. +3. [Execution protocol](execution-protocol.md): authorization, ledger/engine + reconciliation, stale ownership, partitions and restore. +4. [Candidate comparison](candidate-comparison.md): mechanisms, free-software + observations, retained duties and selection rationale. +5. [Experiment report](experiment-report.md), [sources](experiments/README.md), + [original evidence](evidence.json), and [local recheck](local-recheck.json). +6. [Requirement disposition](decision-discussion.md) and + [preflight](../../decisions/issue-1350-execution-architecture-preflight.md). + +## Evidence provenance + +The workspace contained the seven-candidate experiment bundle when this delivery +began. Its thirteen source hashes match the retained files. The original report +and [cleanup record](cloud-cleanup.json) describe that earlier synthetic cloud +run; this delivery did not provision cloud resources. A new local run rechecked +AnyIO, SimPy, DBOS, Temporal and SQLite publication with the same source bytes, +and ran the three witness tests. A further Temporal run checked terminal status +and result/error types through a verifier with four negative-control tests, +while keeping the original thirteen source files unchanged. +ROS/Ray/BehaviorTree.CPP results remain from the +original retained run. The two evidence records must not be conflated. + +The [issue](https://github.com/OpenRAE/rae/issues/1350), +[#1348 decision](../../decisions/issue-1348-operation-lifecycle.md) and +[Hub #3](https://github.com/OpenRAE/hub/issues/3) establish scope and responsibility. +RAE determines admitted execution and validates outcomes; backends realize +resources, report effects and enforce contextual limits. A cancellation receipt, +engine success, checkpoint or process death cannot establish physical cessation. + +Production persistence, tenant isolation, distributed ownership, co-simulation +fidelity and backend containment are not demonstrated by these bounded probes. +The design states how to handle their limits and what executable delivery must +verify. It does not use the local quickstart to limit either backend product. diff --git a/docs/research/execution-architecture/authored-retry-policy.md b/docs/research/execution-architecture/authored-retry-policy.md new file mode 100644 index 000000000..3da8741b7 --- /dev/null +++ b/docs/research/execution-architecture/authored-retry-policy.md @@ -0,0 +1,169 @@ +# Author-controlled retries and scoped defaults + +Design adopted by [ADR-113](../../decisions/adrs/adr-113-reusable-execution-machinery.md) +for #1350. This is a semantic integration design, not published SDL syntax or a +new executable retry contract. Extend the existing owning contracts before +claiming support. No engine-specific retry fields enter the authored language. +Retry is one part of a contextual failure-response policy, not a universal +failure handler shared indiscriminately by experiments, twins, OT and IT. + +## Context determines the permitted response + +The common operation lifecycle records facts; it does not impose the same +recovery action or validity meaning on every world. Resolve policy against the +authored purpose, effect contract, operating mode, time/fidelity requirements, +and actual backend guarantees. The following are examples of selectable +policies, not automatic rules inferred from an `IT` or `OT` label: + +| Context | Example response and required evidence | +| --- | --- | +| Controlled experiment | Mark this trial failed/invalid under the experiment contract, retain every attempt and observation, then request a separately admitted fresh trial if allowed. An idempotent retry may still bias results and is not automatically permitted. | +| Discrete or continuous digital twin | Hold advancement at a supported boundary, invalidate or qualify the affected state/time interval, and reconcile model/plant state before continuation. Re-reading a sensor, rewinding time or restoring a model checkpoint may change the experiment or break coupling; none is an ordinary transport retry. | +| Live OT or hardware-in-the-loop | Request the backend's admitted safe-hold/controlled-stop procedure, or a specified continuing mode where stopping is unsafe. Require process-state, interlock and operator authorization evidence where the authored contract demands it. Reissuing an actuator command, rebooting a controller or restoring an IT snapshot is not a generic recovery strategy. | +| Disposable IT test/CTF | An author may allow rebuilding a VM, restoring a known snapshot, then retrying after verified clean-state admission. That policy is inappropriate for resources the scenario is required to preserve. | +| Stateful IT/product or disaster recovery | Reconcile transactions, external services, data consistency and recovery objectives before selective retry/restore. IT is not inherently reversible or idempotent; externally visible effects may be as non-repeatable as OT actions. | +| Mixed IT/OT scenario | Apply each scope's effect and failure policy while preserving cross-scope dependencies, shared resources and time coupling. Restarting an IT gateway must not silently authorize another plant action or invalidate a twin's admitted state. | + +Failure handling must distinguish infrastructure/delivery failure, backend +refusal, operation failure, deliberate in-world faults, loss of fidelity/time +continuity, and invalid experimental evidence. Do not let an infrastructure +reconciler repair an intentionally injected outage, or let a cleanup failure +disappear behind a primary success. RAE selects only the permitted response +using the canonical workflow/time/trial owners; the backend establishes whether +the concrete response is possible and what it actually did. + +Convenience defaults can select reusable context policies and narrower scopes +can choose different ones. At admission, validate their interactions: one +scope's reset can affect another scope's equipment, data or clock. An +incompatible composition is refused, not resolved by whichever scope retries +first. No backend or engine supplies an unrecorded domain policy. + +## Choices and identities + +Idempotency describes an effect's repeat behavior; permission describes whether +the author wants repetition. Both must hold where required. An idempotent call +can still invalidate an experiment through repeated observations, elapsed time, +cost or participant exposure. RAE must support these distinct authored choices: + +| Author choice | Meaning | +| --- | --- | +| Never repeat | At most one effect invocation for this admitted occurrence. A failed or uncertain invocation cannot cause another effect merely through retry/recovery. | +| Bounded retry in this trial | Retry only for selected failure/effect classes, within attempt and time budgets, after the declared safety and cleanup conditions hold. This can require verified absence, scoped backend idempotency, reset or compensation. | +| End this trial; request a fresh trial | Preserve the failed/invalid trial and its evidence. Admit a new trial through the experiment authority, with a new run identity, declared initialization/variation and verified isolation/cleanup. Do not increment a workflow retry counter and relabel the old trial as successful. | + +Fresh-trial policy applies only where a trial context exists and its owning +experiment plan permits allocation; reject it on an ordinary operation lacking +that authority. Authors can allow bounded retries followed by a fresh trial on +exhaustion, or choose a fresh trial immediately. The trial allocation budget is +independent of within-trial attempts; neither can be unbounded by omission. + +Keep four identities separate: engine delivery/task attempt; one permitted +backend invocation; authored workflow/execution attempt; experimental trial. +Redelivering a reference to the same invocation never allocates another one. +An explicitly admitted retry gets its own invocation/attempt identity and +provenance; a fresh trial gets its own run identity. Engine reset, restart and +Continue-As-New are not any of those author decisions. + +## Scope resolution + +Use the existing [lexical scope principle](../../../specs/sdl/recursive-realization-constraints.md#2-records-scopes-and-closure): +scenario default, enclosing scope defaults, then the most-specific explicit +operation/occurrence policy. Reuse canonical semantic addresses and bounded +reference resolution. This reuses scope mechanics, not the open/closed value +vocabulary: closure does not imply a retry policy. + +1. A missing local policy inherits. An explicit never-repeat policy is a value, + not omission; a locally stricter or more permissive default affects only that + subtree. A concrete child choice can override a convenience default, subject + to binding requirements and admission. Siblings keep their inherited policy. +2. Resolve a **complete policy** at the nearest defining scope. Do not combine + unrelated fields from different levels into an accidental policy, such as a + child's reset mode with a parent's idempotency assumptions. Reusable named + policies provide concise authoring; any later partial-overlay syntax must + normalize to the same complete, validated value with field provenance. +3. Duplicate definitions at the same semantic scope, ambiguous targets, cycles, + unsupported versions and conflicting binding constraints are admission + errors. Import order, list order and engine defaults never break a tie. + Reusable definitions retain lexical binding; execution at a call site cannot + silently change them. An explicit application-site policy is validated as a + new local choice before the compiled plan is admitted. +4. Defaults are convenience, not a way to weaken a required guarantee, bypass + backend policy or grant authority. Resolve the author's desired policy, then + validate the entire composition. Refuse incompatibility instead of silently + clamping to fewer retries or substituting a new trial. +5. If no scope supplies permission to repeat, permit no second effect. This is + a specified conservative fallback, not an implementation library's default. + The author remains free to select a different scenario-wide default. +6. Materialize the resolved policy and defining scope/reference/version into the + admitted plan and operation commitment. Record changes as new admissions; + editing a scenario default cannot retroactively alter an active invocation, + a restored operation or historical evidence. + +Conceptual examples, **not YAML syntax**: + +| Scenario default | Nearer scope or local choice | Effective result | +| --- | --- | --- | +| Never repeat | Setup scope: at most three total attempts, admitted idempotency required | Setup may retry only under that condition; other operations remain one-shot. | +| Setup policy above | One actuator operation: never repeat | That operation remains one-shot, including after worker loss. | +| Bounded idempotent retry | Measurement scope: end this trial and request a fresh trial | Measurements never repeat within the same trial, even if technically idempotent. | +| Never repeat | No local choice | Inherit never repeat; no annotations are needed on every operation. | + +## Effective policy and runtime admission + +The normalized contract must identify its context/profile and revision, +failure classification, validity consequences, permitted response (for example +abort, hold, reconcile, continue at a verified boundary, retry or fresh trial), +and required operator/backend evidence. It must also identify the retry unit, +permitted effect-knowledge classes, total-attempt limit including the first attempt, +delay/backoff and clock basis, total budget, required safety/cleanup evidence, +and exhaustion disposition. Fresh-trial selection additionally binds trial +allocation authority and limits. Reuse `ExecutionRetryPolicyModel`'s existing +`max_attempts` and `after_effect_policy` meanings (`disallow`, `idempotent`, +`reset`, `compensate`) where applicable. Do not infer never-repeat solely from +`after_effect_policy=disallow`: that field addresses after-effect repetition, +while verified absence can have different admitted semantics. + +Current SCE-007 cleanup policies do not provide general lexical inheritance, +failure-class/backoff selection or fresh-trial policy. These need governed +extensions under DSL-113/SEM-203 and SCE-002/006/007/EXP-706, not a free-form +metadata dictionary or a second workflow interpreter. Existing explicit workflow +retry bounds remain binding; nested policies must account for every actual +effect and must not multiply hidden retries at transport, activity and workflow +layers. Exhausted budgets do not reset on process restart or scope inheritance. + +Before each permitted invocation, RAE rechecks current scope authorization, +policy identity, remaining budgets, backend capability/willingness, reservation +conflicts and required evidence. Backend idempotency needs a defined key, scope, +retention interval and behavior under duplicate/concurrent requests; an author +label is insufficient. Known absence is not cessation if an old caller can +still act later. Concurrent repetition requires explicit proven-safe semantics; +sequential attempts remain the default. Unknown effect state never becomes +safe merely because a worker disappeared. + +Reset, compensation and fresh-trial initialization are effects with their own +admission, failure, observation and cleanup obligations. A requested reset does +not establish clean state. An unreconciled old trial cannot share resources with +a new trial unless admitted isolation actually separates their effects. Preserve +the old failure, cleanup result, lineage and any experiment validity decision. + +## Engine mapping and verification obligations + +Use Temporal's bounded scheduling and retry facilities when they can enforce +the effective policy through current RAE admission on every dispatch. Otherwise +schedule separately admitted single-attempt activities. Never turn engine retry +numbers into the author-visible attempt history. Bookkeeping retry and workflow +replay must not invoke an effect again. Observation can itself affect a system, +so classify it by its admitted effect contract rather than assuming all reads +are harmless. + +The public-contract and implementation deliveries must test scenario defaults, +nested overrides, sibling isolation, explicit never-repeat, duplicate scopes, +definition/call-site binding, bounded normalization, incompatible guarantees, +scope reordering, missing policy, policy edits during recovery, budget exhaustion, +and no retry multiplication. Exercise different failure responses for the same +technical exception under experiment, twin, OT and IT policies, plus conflicting +responses across mixed scopes and preservation of intentional faults. +Exercise the same idempotent backend under both +retry-allowed and retry-forbidden policies, plus fresh-trial allocation and +cleanup failure. Those are behavioral gates for the later executable surface; +this design does not claim they pass today. diff --git a/docs/research/execution-architecture/candidate-comparison.md b/docs/research/execution-architecture/candidate-comparison.md new file mode 100644 index 000000000..cc79a72dd --- /dev/null +++ b/docs/research/execution-architecture/candidate-comparison.md @@ -0,0 +1,134 @@ +# Candidate comparison + +Status: supporting comparison for the selection in +[ADR-113](../../decisions/adrs/adr-113-reusable-execution-machinery.md). +Sources checked 2026-09-23; exact tested +versions are recorded in [evidence](evidence.json). Integration assessments below +are engineering inferences from those sources, the probes, and RAE's existing +contracts. They are not throughput measurements or a complete package audit. + +## Criteria + +RAE must discharge its duties: authored retry, continuation, trial, time, refusal +and cleanup rules; one operation-state authority; bounded supervision; and +validated external-effect evidence. Backend cooperation is expected. Backends +retain realization, observations and contextual limits under the +[responsibility boundary](../../decisions/adrs/adr-113-reusable-execution-machinery.md#2-keep-the-ecosystem-and-runtime-boundaries). + +Every candidate has limits. Evaluate the complete composition: useful mechanism, +remaining runtime duty, backend obligation, required handling and residual limit. +A limitation alone is not a rejection; an unfulfilled required duty is. Keep +untested handling distinct from demonstrated handling. Maintenance and deployment +work must be recorded, but cannot override the separation of duties. The bounded +probes do not justify a numerical ranking or prove any composition complete. + +P0 remains an in-process library profile. An optional external worker does not +by itself introduce multiple state owners, but a mandatory remote orchestration +service would change that deployment contract. P1/P2 preserve their existing +declared authority boundaries. Neither distributed workers nor framework HA +establish RAE's unavailable P3 profile. + +The earlier evaluation uses a design case of 15 federated tenants, hundreds of +nodes and tens of agents per scenario, without measured capacity evidence. +Local profiles are compatibility obligations, not the +platform's ceiling. The distributed design must extend the currently missing +coordination and tenant guarantees explicitly. Prefer established components +for distributed execution; custom code should implement RAE-specific semantics +and integration, not replace a mature queue, scheduler or recovery engine. + +## Mechanisms and responsibilities + +| Candidate | State and worker ownership; supervision/scheduling | Persistence and failure isolation | Reuse and retained RAE duties | +| --- | --- | --- | --- | +| Python asyncio/AnyIO | RAE retains its owner; task groups, capacity limiters and worker APIs manage execution. Control capacity must be separated from occupied effect capacity. | No new durable history. Existing RAE stores retain operation authority. Threads share the host process; process workers isolate local execution but not already dispatched remote effects. | Reuse concurrency and process primitives. RAE still implements dispatch/reservation lifecycle, scoped cancellation, bounded observation, truthful outcome mapping and startup reconciliation. Small dependency surface does not make that work trivial. | +| DBOS | An embedded engine owns workflow/step bookkeeping and queues. RAE would bind those records to its admitted operation rather than equate their states. | SQLite or Postgres system DB; interrupted steps can run again on recovery. Application process loss was tested, not machine power loss or a distributed deployment. | Reuse checkpoints, queueing and workflow management. RAE must guard each effect against unauthorized replay, classify uncertainty and reconcile engine history with its operation store. Reusing SQLite does not make their transactions automatically atomic together. | +| Temporal | Service history and task queues coordinate application workers. Workflow and activity state remain engine execution state, distinct from RAE's operation outcome. | Service persistence is separate from worker memory. Worker loss was tested with a live, in-memory development server; service/storage loss and production HA were not tested. | Reuse durable dispatch, workflow history, retry controls and heartbeat cancellation. RAE still owns effect admission, outcome evidence, quarantine and cross-store reconciliation. Disable or guard retries where repetition is unauthorized. | +| Celery/task queue (source comparison only) | Task routing and workers provide mature dispatch. Acknowledgement, retry and revocation settings require deliberate configuration. | A broker/result store adds operations; task status is not a durable RAE semantic history or effect witness. No Celery experiment was run. | Reuse at an existing backend task boundary; building the complete orchestration/recovery layer around it leaves more general machinery for RAE to implement than the selected composition. | +| ROS 2 actions | Backend action server owns goal handling; executor/callback configuration governs responsiveness. Client receives goal/cancel acknowledgements, progress and results. | The tested action protocol is not RAE's durable operation store. Node/middleware loss needs separately defined recovery and effect observation. | Reuse an action protocol and clients/servers at suitable backend boundaries. RAE must correlate goal/control identities, validate backend claims and map results. An action status cannot replace domain evidence. | +| Ray actors | Ray schedules actor processes; actor methods can serialize work. Same-actor control can queue behind a blocked method. Alternate concurrency/control arrangements need separate tests. | Configurable restart and task retries; application state recovery is the application's responsibility. Actor death is distinct from cessation of downstream effects. | Reuse placement, process actors and worker recovery. RAE retains durable authority, retry permission and reservations. Value increases if distributed computation is a real workload requirement, which these probes did not establish. | +| SimPy | Cooperative, deterministic event processing can provide a semantic-time scheduler under RAE's declared time authority. | In-memory event queue is not durable operation history. A blocking callback blocks event stepping. | Reuse event ordering and process interrupts. RAE retains multi-domain time interpretation, external execution, independent apparatus supervision and any restart representation. | +| BehaviorTree.CPP | Tick-driven reactive composition with lifecycle/halting hooks. Action implementations supply nonblocking execution and actual interruption behavior. | Tree state and an action's worker/effects are distinct. Persistence and RAE settlement are not supplied by the tested halt operation. | Reuse behavior composition where semantics match. Translating SDL workflows would need semantic validation, not just an API adapter. Backend-local use has a narrower integration boundary. | + +## Deployment obligations + +| Candidate | Added deployment responsibility if selected | Initial license observation | +| --- | --- | --- | +| AnyIO | Python library; local child processes where selected; packaging/IPC and host supervision. Existing HTTP offload already uses AnyIO indirectly. | MIT | +| DBOS | Python library plus system database. SQLite is supported; vendor recommends Postgres for production. Distributed recovery adds coordination choices. | MIT | +| Temporal | Python workers plus the selected self-hosted service and its persistence. Version compatibility and workflow upgrades become operational concerns. | Python SDK and server: MIT | +| ROS 2 | ROS distribution, client libraries, generated action types and selected middleware; discovery/executor configuration and upgrades. | rclpy: Apache-2.0; not an audit of the ROS distribution | +| Ray | Ray runtime processes, worker resource configuration and potentially a cluster; application checkpoint/upgrade policy. | Apache-2.0 | +| SimPy | Python library; RAE integration with semantic time and execution boundaries. | MIT | +| BehaviorTree.CPP | C++ library/build toolchain or backend process boundary; behavior-node integration. | MIT | + +License identifiers were checked against upstream repository metadata, not +assumed from the category. Transitive licenses, vulnerability posture, supported +platforms across RAE's Python range, production operations and total ownership +cost remain adoption checks. No candidate's documentation benchmark is treated +as a RAE performance measurement. No production dependency was added. + +## Selection rationale and alternatives + +Select Temporal Server/Python SDK/PostgreSQL for the connected distributed +composition, with the authority and failure rules in ADR-113. Its service/worker +separation, durable execution history, queues and operational tooling match that +composition directly. Preserve existing AnyIO/SQLite local compositions; do not +build a distributed scheduler from those primitives. Local library suitability +and distributed execution suitability are different selection questions. + +DBOS is a credible alternative, including distributed queues; neither the probes +nor this assessment establish a throughput ranking or show it incapable. Its +embedded orchestration model offers a smaller service footprint. Temporal's +separate execution service and worker lifecycle are the selected operational +trade-off. Celery is established task machinery, but would leave more durable +orchestration/history/recovery integration to RAE. Reopening the decision needs +a concrete unsupported requirement or measured operational problem. + +ROS actions, Ray actors, SimPy and BehaviorTree.CPP address complementary needs. +Adopt them at a backend seam only when its effect, compute, time or control +contract needs them; their presence cannot create a second scenario interpreter. +No package is selected because all digital twins are assumed to be robots, all +AI research requires Ray, or all failures should trigger workflow recovery. + +[Contextual policy](authored-retry-policy.md) is common admission machinery with +different author-selected responses: experiment invalidation/fresh trials, +twin time/state reconciliation, OT safe-process response, and IT restoration or +retry where appropriate. Defaults resolve by scenario and lexical scope; +idempotency does not imply author permission. The candidate's own retry defaults +must implement that resolved policy or be disabled for the effect concerned. + +Deployment cost is a trade-off, not grounds to move RAE responsibilities into a +backend. The selected components require no paid software or hosted service. +Transitive dependencies and platform support remain package-admission checks; +no claim that every ecosystem package or optional commercial feature is free +follows from the directly checked licenses. + +## Primary sources + +- AnyIO: [threads](https://anyio.readthedocs.io/en/stable/threads.html), + [processes](https://anyio.readthedocs.io/en/stable/subprocesses.html), + [repository](https://github.com/agronholm/anyio). +- DBOS: [architecture/recovery](https://docs.dbos.dev/architecture), + [Python databases](https://docs.dbos.dev/python/tutorials/database-connection), + [cancellation](https://docs.dbos.dev/python/tutorials/workflow-management), + [repository](https://github.com/dbos-inc/dbos-transact-py). +- Temporal: [activities](https://docs.temporal.io/activities), + [service](https://docs.temporal.io/temporal-service), + [Python cancellation](https://docs.temporal.io/develop/python/workflows/cancellation), + [SDK](https://github.com/temporalio/sdk-python), + [server](https://github.com/temporalio/temporal). +- ROS: [action protocol](https://design.ros2.org/articles/actions.html), + [callback groups](https://docs.ros.org/en/ros2_documentation/kilted/How-To-Guides/Using-callback-groups.html), + [rclpy](https://github.com/ros2/rclpy). The action design document is historical; + executable observations use the recorded Jazzy packages, not an assumed latest distribution. +- Ray: [actors](https://docs.ray.io/en/latest/ray-core/actors.html), + [fault tolerance](https://docs.ray.io/en/latest/ray-core/fault_tolerance/actors.html), + [cancellation](https://docs.ray.io/en/latest/ray-core/api/doc/ray.cancel.html), + [repository](https://github.com/ray-project/ray). +- SimPy: [scheduling](https://simpy.readthedocs.io/en/stable/topical_guides/time_and_scheduling.html), + [repository](https://github.com/simpx/simpy). +- BehaviorTree.CPP: [asynchronous actions](https://www.behaviortree.dev/docs/tutorial-basics/tutorial_04_sequence/), + [repository](https://github.com/BehaviorTree/BehaviorTree.CPP). +- Celery (source-only): [task execution and acknowledgements](https://docs.celeryq.dev/en/stable/userguide/tasks.html), + [worker revocation](https://docs.celeryq.dev/en/stable/userguide/workers.html). +- PostgreSQL: [license](https://www.postgresql.org/about/licence/). diff --git a/docs/research/execution-architecture/cloud-cleanup.json b/docs/research/execution-architecture/cloud-cleanup.json new file mode 100644 index 000000000..77a97719d --- /dev/null +++ b/docs/research/execution-architecture/cloud-cleanup.json @@ -0,0 +1,56 @@ +{ + "schema": "rae.execution-experiment-cleanup/v1", + "profile": "catalyst-dev", + "region": "us-east-2", + "stack_name": "rae-1350-experiments-20260923", + "stack_status": "DELETE_COMPLETE", + "instance_launch_utc": "2026-09-23T02:21:38Z", + "stack_deletion_complete_utc": "2026-09-23T02:35:03Z", + "resources": { + "instance": "i-077968adda5c8517e", + "volume": "vol-08169e8000d3189a1", + "network_interface": "eni-09ff348f1eb4762f3", + "subnet": "subnet-0adc6d31269ad61d1", + "security_group": "sg-00d6cb6b32a857e89", + "role": "rae-1350-experiments-20260923-Role-jooMOFJuHj2R", + "instance_profile": "rae-1350-experiments-20260923-Profile-aDzafOad5qOF" + }, + "verification": { + "instance": [ + "terminated" + ], + "volume": [], + "interface": [], + "subnet": [], + "security_group": [], + "role": { + "result": "NoSuchEntity", + "exit_code": 254 + }, + "profile": { + "result": "NoSuchEntity", + "exit_code": 254 + }, + "default_subnets": [ + { + "Id": "subnet-08827c5dbc4520d4c", + "Cidr": "172.31.0.0/20" + }, + { + "Id": "subnet-0ebdef2813812d8df", + "Cidr": "172.31.16.0/20" + }, + { + "Id": "subnet-00684b34f43fec00c", + "Cidr": "172.31.32.0/20" + } + ] + }, + "notes": [ + "Verified with direct EC2 and IAM reads after CloudFormation DELETE_COMPLETE.", + "Three original default VPC subnets remain unchanged.", + "Root disk and ENI deleted; public address was auto-assigned, no EIP allocation.", + "CloudFormation and SSM retain service history records; no live experiment resources remain.", + "Local experiment sources, evidence, and private lifecycle records retained." + ] +} diff --git a/docs/research/execution-architecture/decision-discussion.md b/docs/research/execution-architecture/decision-discussion.md new file mode 100644 index 000000000..3044cb712 --- /dev/null +++ b/docs/research/execution-architecture/decision-discussion.md @@ -0,0 +1,48 @@ +# Decision disposition and canonical traceability + +[ADR-113](../../decisions/adrs/adr-113-reusable-execution-machinery.md) is the +single accepted execution-architecture decision for #1350. It replaces the +provisional Temporal/PostgreSQL discussion previously kept here. Acceptance +selects an architecture; it does not claim a runtime implementation or P3. + +The design resolves the earlier worker-authorization, ledger/history and +stale-owner questions in the [execution protocol](execution-protocol.md). +[Contextual failure and retry policy](authored-retry-policy.md) records author +control, scenario/scoped defaults, and distinct experiment, twin, IT and OT +responses. The [comparison](candidate-comparison.md) and +[experiments](experiment-report.md) retain alternatives and evidence limits. + +## Canonical requirement mapping + +API-404 is the issue's assigned requirement and remains ACTIVE for its existing +P0–P2 clauses. Add DOCUMENTS traceability for this design and mechanism evidence; +retain existing IMPLEMENTS/TESTS assertions at their stated scope. Other owners +below retain their statements and statuses, including all DRAFT time owners. +This disposition does not claim their missing implementation has shipped. + +| Canonical owner | Evaluation relationship and remaining boundary | +| --- | --- | +| [API-404](../../requirements/API-404/requirement.md), C1–C4 | Single state authority, atomic settlement, startup without implicit replay, bounded served supervision and profile nonclaims. Toy publication and candidate probes are not production conformance evidence. | +| [API-402](../../requirements/API-402/requirement.md), [API-403](../../requirements/API-403/requirement.md), [RUN-304](../../requirements/RUN-304/requirement.md) | Submission/status/history remain portable; any new invocation/control/evidence carriers need governed publication. | +| [RUN-300](../../requirements/RUN-300/requirement.md), [ASR-532](../../requirements/ASR-532/requirement.md) | Shared RAE runtime drives backends and validates results; worker/library adoption does not transfer semantic authority. | +| [DSL-113](../../requirements/DSL-113/requirement.md), [SEM-203](../../requirements/SEM-203/requirement.md), [SEM-204](../../requirements/SEM-204/requirement.md) | Authored workflow/attempt/retry and compensation rules govern effects, rather than engine retry defaults. | +| [SEM-227](../../requirements/SEM-227/requirement.md), [SEM-228](../../requirements/SEM-228/requirement.md), [SEM-229](../../requirements/SEM-229/requirement.md), [RUN-317](../../requirements/RUN-317/requirement.md), [RUN-318](../../requirements/RUN-318/requirement.md), [API-421](../../requirements/API-421/requirement.md) | Scenario-time pause and apparatus time are distinct. SimPy demonstrates a mechanism only; clock mappings, restart continuity and multiple time domains remain untested. DRAFT owners remain DRAFT. | +| [SCE-006](../../requirements/SCE-006/requirement.md), [SCE-007](../../requirements/SCE-007/requirement.md), [EXP-706](../../requirements/EXP-706/requirement.md), [EXP-712](../../requirements/EXP-712/requirement.md), [SCE-002](../../requirements/SCE-002/requirement.md) | No framework restart may silently allocate another attempt/trial or establish clean state. The probes do not implement trial admission or cleanup verification. | +| [RUN-310](../../requirements/RUN-310/requirement.md), [RUN-311](../../requirements/RUN-311/requirement.md), [SEM-222](../../requirements/SEM-222/requirement.md) | Participant episode/termination authority is separate from an engine task or backend goal lifecycle. No participant semantics changed. | +| [RUN-316](../../requirements/RUN-316/requirement.md) | Operational status is distinct from experiment evidence. The independent witness deliberately measures effects separately from framework status. | + +## Acceptance mapping + +| Issue criterion | Delivered artifact and verification | +| --- | --- | +| Compare applicable patterns and reusable seams | Candidate comparison and ADR-113 alternatives; primary sources, explicit deployment/license observations and seven retained candidate probes. | +| State/worker ownership, supervision, scheduling, persistence, backend/deployment/failure boundaries | ADR-113 component/interface diagram and execution protocol; single semantic authority and explicit loss-window table. | +| Bounded evidence for blocked calls, cancellation, recovery and duplicate effects | Original source-hashed evidence/report plus independent local recheck; witness assertions distinguish effects from framework status. | +| One accepted decision with canonical traceability | ADR-113, ADR index/pin and API-404 DOCUMENTS entries; targeted governance and pin checks. | +| Author retry/default/context clarifications | Authored retry-policy design, scenario/nested/local resolution, context-specific failure response and mixed-scope compatibility. Future schema/compiler/runtime support is explicitly unclaimed. | + +API-404-C1 is preserved by one state writer and atomic publication; C2 by +retained claims and observation before a newly authorized effect; C3 by scoped +supervisory authorization and bounded control paths; C4 by keeping P3 unavailable +and excluding implicit tenant/coordination claims. Existing supervision and +profile tests check their current boundaries, not distributed conformance. diff --git a/docs/research/execution-architecture/evidence.json b/docs/research/execution-architecture/evidence.json new file mode 100644 index 000000000..86265a525 --- /dev/null +++ b/docs/research/execution-architecture/evidence.json @@ -0,0 +1,658 @@ +{ + "schema": "rae.execution-experiment-evidence/v1", + "date": "2026-09-23", + "python": "3.12.3", + "platform": "Linux-7.0.0-1012-aws-x86_64-with-glibc2.39", + "python_packages": [ + "anyio==4.15.1", + "attrs==26.1.0", + "certifi==2026.7.22", + "charset-normalizer==3.5.1", + "click==8.5.0", + "dbos==3.0.0", + "filelock==4.0.1", + "greenlet==3.5.6", + "idna==3.20", + "iniconfig==2.3.0", + "jsonschema==4.26.0", + "jsonschema-specifications==2025.9.1", + "msgpack==1.2.2", + "nexus-rpc==1.4.0", + "packaging==26.3", + "pluggy==1.6.0", + "protobuf==7.36.2", + "psycopg==3.3.6", + "psycopg-binary==3.3.6", + "Pygments==2.21.0", + "pytest==9.1.1", + "python-dateutil==2.9.0.post0", + "PyYAML==6.0.3", + "ray==2.58.0", + "referencing==0.37.0", + "requests==2.34.2", + "rpds-py==2026.6.3", + "simpy==4.1.2", + "six==1.17.0", + "SQLAlchemy==2.0.54", + "temporalio==1.33.0", + "types-protobuf==7.35.1.20260906", + "typing_extensions==4.16.0", + "urllib3==2.8.0", + "websockets==17.1" + ], + "ros_packages": [ + "ros-jazzy-action-tutorials-interfaces\t0.33.11-1noble.20260902.024042", + "ros-jazzy-behaviortree-cpp\t4.10.0-1noble.20260902.082337", + "ros-jazzy-rclpy\t7.1.12-1noble.20260902.053513" + ], + "source_sha256": [ + "113dbc62df2be37924b124ee2c2abdaea7e0ea35a15d8beb5c928d5e49aaca2c dbos_worker.py", + "0533b361058deb1813d8e364b04eef27ab35aa692b0474eae9e98cf814adbc68 probe_dbos.py", + "d874b4c6ba98e7999ee1cbcfc3379917ac70a3cae9f527f873ef16737acc6b61 probe_publication.py", + "982ed3fe873176f5a4ba3b0c30d8d9c823fdd593af62e6e643f17fac7d1dee9f probe_ray.py", + "46f1409a772f007a4dcf788c77821ed870cc80f74e5ee5d75f7c6bee555e2c1b probe_ros.py", + "9335ab9b9129ddb6cea683a0e94cc35c3ca7d3a868216fe80c9fac617ba30861 probe_simpy.py", + "7b12c63f91a0a138553acd7c2f5bad4b0016804581e68eba19caf0604a746517 probe_temporal.py", + "85e03d1d1f43e560e67dc5eba57fb80a016391ebdd4dc83cc9354d78988dc788 probe_workers.py", + "56529d06d2a23fc4c6d4d20968744e6805a71b45c46b61009277609886f44592 temporal_worker.py", + "d6613b18a645861fd6b7322e1650a62b237b6be5e64d9a93f00aae02abd7ac69 test_witness.py", + "26c63959b52ea4c6f4e77902a9a18acf0ba97304dd91f7a9ff6feac4571ce51f witness.py", + "5ddee1930148fdf0ac2301cf153070965571efe32df3cbf559448783d67a7074 probe_bt.cpp", + "2076c13bcaf289d344665d7c24c8ecef4fc43d09121a710f1d7300fb24f158b8 CMakeLists.txt" + ], + "witness_tests": "...\n----------------------------------------------------------------------\nRan 3 tests in 0.578s\n\nOK\n", + "temporal_environment": { + "cli": "1.9.1", + "server": "1.32.0", + "persistence": "in-memory", + "loss_boundary": "worker process, not service" + }, + "results": { + "anyio": { + "thread": { + "at_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 525.66545347 + ] + ], + "caller_pids": [ + 4688 + ] + }, + "local_cancel_return_seconds": 0.0016588409999940268, + "after_local_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 525.66545347 + ] + ], + "caller_pids": [ + 4688 + ] + }, + "caller_process_alive_after_cancel": true, + "later": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 525.66545347 + ], + [ + "effect", + 526.365592469 + ] + ], + "caller_pids": [ + 4688 + ] + } + }, + "process": { + "at_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 526.718613595 + ] + ], + "caller_pids": [ + 4697 + ] + }, + "local_cancel_return_seconds": 0.0031917039999598273, + "after_local_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 526.718613595 + ] + ], + "caller_pids": [ + 4697 + ] + }, + "caller_process_alive_after_cancel": false, + "later": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 526.718613595 + ], + [ + "effect", + 527.418764818 + ] + ], + "caller_pids": [ + 4697 + ] + } + }, + "saturation": { + "shared_capacity_control_expired": true, + "separate_capacity_status": "status", + "witness_at_separate_status": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 527.634117148 + ] + ], + "caller_pids": [ + 4688 + ] + } + } + }, + "dbos": { + "unguarded": { + "worker_crash_exit": 73, + "after_crash": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 529.4805973 + ], + [ + "effect", + 529.48068209 + ] + ], + "caller_pids": [ + 4726 + ] + }, + "after_recovery": { + "requests": 2, + "effects": 2, + "events": [ + [ + "request", + 529.4805973 + ], + [ + "effect", + 529.48068209 + ], + [ + "request", + 531.262332609 + ], + [ + "effect", + 531.262407148 + ] + ], + "caller_pids": [ + 4726, + 4736 + ] + }, + "worker": { + "result": "effect-returned", + "status": "SUCCESS" + } + }, + "guarded": { + "worker_crash_exit": 73, + "after_crash": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 534.336986398 + ], + [ + "effect", + 534.337016095 + ] + ], + "caller_pids": [ + 4748 + ] + }, + "after_recovery": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 534.336986398 + ], + [ + "effect", + 534.337016095 + ] + ], + "caller_pids": [ + 4748 + ] + }, + "worker": { + "result": "indeterminate-no-second-invocation", + "status": "SUCCESS" + } + }, + "cancellation": { + "at_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 538.116096225 + ] + ], + "caller_pids": [ + 4767 + ] + }, + "status_after_cancel": "CANCELLED", + "later": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 538.116096225 + ], + [ + "effect", + 538.916246631 + ] + ], + "caller_pids": [ + 4767 + ] + }, + "following_step_executed": false + } + }, + "temporal": { + "temporal-crash-2": { + "worker_crash_exit": 73, + "after_crash": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 548.418652553 + ], + [ + "effect", + 548.418724706 + ] + ], + "caller_pids": [ + 5034 + ] + }, + "after_recovery": { + "requests": 2, + "effects": 2, + "events": [ + [ + "request", + 548.418652553 + ], + [ + "effect", + 548.418724706 + ], + [ + "request", + 551.419929791 + ], + [ + "effect", + 551.420034051 + ] + ], + "caller_pids": [ + 5034, + 5096 + ] + }, + "result": "effect-returned" + }, + "temporal-crash-1": { + "worker_crash_exit": 73, + "after_crash": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 562.00034272 + ], + [ + "effect", + 562.000412499 + ] + ], + "caller_pids": [ + 5111 + ] + }, + "after_recovery": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 562.00034272 + ], + [ + "effect", + 562.000412499 + ] + ], + "caller_pids": [ + 5111 + ] + }, + "error_type": "WorkflowFailureError" + }, + "temporal-cancel-uncooperative": { + "at_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 574.580207594 + ] + ], + "caller_pids": [ + 5184 + ] + }, + "later": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 574.580207594 + ], + [ + "effect", + 575.380334751 + ] + ], + "caller_pids": [ + 5184 + ] + }, + "settlement_seconds": 0.8114476210000703, + "workflow_status": "COMPLETED", + "result": "effect-returned" + }, + "temporal-cancel-cooperative": { + "at_cancel": { + "requests": 0, + "effects": 0, + "events": [], + "caller_pids": [] + }, + "later": { + "requests": 0, + "effects": 0, + "events": [], + "caller_pids": [] + }, + "settlement_seconds": 0.043978476000006594, + "workflow_status": "CANCELED", + "error_type": "WorkflowFailureError" + } + }, + "ray": { + "configured_retry": { + "result": "completed", + "witness": { + "requests": 2, + "effects": 2, + "events": [ + [ + "request", + 543.343976226 + ], + [ + "effect", + 543.344088061 + ], + [ + "request", + 544.3138101 + ], + [ + "effect", + 544.313912033 + ] + ], + "caller_pids": [ + 4940, + 4939 + ] + } + }, + "control_on_same_actor_blocked": true, + "at_kill": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 544.349559141 + ] + ], + "caller_pids": [ + 4939 + ] + }, + "after_kill": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 544.349559141 + ], + [ + "effect", + 545.349701126 + ] + ], + "caller_pids": [ + 4939 + ] + } + }, + "ros": { + "ros-accept": { + "at_request": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 579.460144349 + ] + ], + "caller_pids": [ + 23 + ] + }, + "at_cancel_reply": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 579.460144349 + ] + ], + "caller_pids": [ + 23 + ] + }, + "cancel_reply_seconds": 0.005308084000034796, + "cancel_return_code": 0, + "goals_canceling": 1, + "terminal_status": 5, + "later": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 579.460144349 + ], + [ + "effect", + 580.260273152 + ] + ], + "caller_pids": [ + 23 + ] + } + }, + "ros-refuse": { + "at_request": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 580.270511766 + ] + ], + "caller_pids": [ + 23 + ] + }, + "at_cancel_reply": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 580.270511766 + ] + ], + "caller_pids": [ + 23 + ] + }, + "cancel_reply_seconds": 0.005622269000014057, + "cancel_return_code": 0, + "goals_canceling": 0, + "terminal_status": 4, + "later": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 580.270511766 + ], + [ + "effect", + 581.070629401 + ] + ], + "caller_pids": [ + 23 + ] + } + } + }, + "simpy": { + "paused_virtual_time": 0, + "apparatus_elapsed_seconds": 0.10013592300003893, + "interrupt_trace": [ + [ + "interrupted", + 0 + ] + ], + "blocking_callback_holds_step": true + }, + "publication": { + "race": [ + { + "requested": "SUCCEEDED", + "won": true + }, + { + "requested": "CANCELLED", + "won": false + } + ], + "terminal": [ + "SUCCEEDED", + 1 + ], + "audit": [ + [ + "SUCCEEDED" + ] + ], + "lost_ack_readback": [ + "SUCCEEDED", + 1 + ], + "scope": "SQLite toy transaction boundary; no external fencing claim" + }, + "behaviortree": { + "probe": "behaviortree", + "initial_running": true, + "halt_callback_count": 1, + "effects_at_halt": 0, + "effects_after_halt": 1, + "halt_return_seconds": 0.000005533 + } + } +} diff --git a/docs/research/execution-architecture/execution-protocol.md b/docs/research/execution-architecture/execution-protocol.md new file mode 100644 index 000000000..0b0c01978 --- /dev/null +++ b/docs/research/execution-architecture/execution-protocol.md @@ -0,0 +1,136 @@ +# Ledger, execution and recovery protocol + +Integration rules adopted by [ADR-113](../../decisions/adrs/adr-113-reusable-execution-machinery.md). +These describe the selected design and implementation obligations, not a shipped +PostgreSQL provider, wire protocol or P3 conformance result. The existing +[supervision model](../../../specs/formal/runtime-control-plane/supervision.md) +continues to own operation outcomes. Internal claim phases below are not new +portable operation states. + +## Admission and dispatch + +1. The authenticated RAE admission boundary resolves the authored policy and + verifies the selected composition, plan, actor, target/run, capability and + willingness. The scoped state writer atomically records the immutable + operation/request commitment and admitted command reference. This is ledger + authority, not a second scheduler or a replacement broker. +2. Start Temporal using a stable identity bound to that admitted operation and + execution generation. Reject duplicate/reused identities; an existing + execution must match the ledger binding. A timeout starting the workflow + triggers readback or resubmission of that same identity, not a new operation. + Start reconciliation uses the ledger's retained nonterminal commands and + Temporal's duplicate-start controls, not a custom durable delivery queue. +3. The trusted effect worker decodes a bounded versioned command and authenticates + to RAE's state authority. Fetch and verify the admitted commitment, artifact + digest and selected capability there. An opaque engine ID, replayed activity + output, process-local planner registration or possession of a secret + reference cannot establish worker authority. Resolve credentials only after + authorization through the backend's scoped secret mechanism. +4. Immediately before external execution, ask the current state writer to + consume a **one-use invocation claim**, tied to operation, authored attempt, + worker identity, expected owner generation and conflicting-effect scope. + The transaction rechecks policy/budget and admission state, records the + invocation and retains its effect reservation. Cancellation committed before + this cut prevents dispatch. The first confirmed consumer alone may invoke; + duplicated delivery can observe the existing claim, never receive another + invocation permission for it. No database transaction remains open across + the backend call. +5. A lost claim acknowledgement is not a reusable grant. Fail closed and + reconcile unless the protocol proves that this exact live consumer received + and has not used a valid grant. In the initial conservative design, ambiguous + grant delivery admits no invocation. A confirmed grant can still outlive a + caller or owner: retain its reservation until cessation/effect evidence + accounts for it. This deliberately trades liveness for honest effect scope. +6. Workers return correlated backend results/evidence to the state authority; + they do not write snapshots. Validate against the trusted predecessor, + admitted policy and native result/release validators. Atomically commit the + permitted snapshot change, terminal operation and matching audit, using + revision comparison. Engine acknowledgement happens after that cut. Repeated + terminal-commit requests resolve through the same immutable identity/readback. + +Temporal replay may repeat bookkeeping activities and may redeliver an effect +task, but it cannot consume the same invocation claim twice. An explicitly +permitted retry needs a new admitted attempt/invocation and the +[resolved author policy](authored-retry-policy.md). It is not forbidden merely +because it is a retry. Conversely, backend idempotency alone cannot authorize it. + +## Failure windows + +| Lost boundary | Recovery action; effect authority | +| --- | --- | +| Before admission commit | Read by scoped request identity to resolve commit uncertainty; do not enqueue an unconfirmed operation. | +| Admission committed, start absent/ack lost | Reconcile/start the same engine identity from the ledger. No effect exists solely from this command record. | +| Start recorded, invocation not consumed | Redelivery still needs the current claim gate. A cancellation or owner change can deny it. | +| Claim consumed, worker dies before or after backend call | Treat the effect as unknown until validated observation plus cessation establishes otherwise. Do not infer absence from missing engine result. | +| Effect applied, worker/result acknowledgement lost | Observe using the same invocation identity; never replay it to discover the answer. A different authorized retry follows the policy and retained reservation rules. | +| Result received, ledger commit uncertain | Close mutation readiness, read the ledger's atomic cut, and reconcile. No receipt says success before a confirmed commit. | +| Ledger terminal committed, engine acknowledgement lost | Return the committed outcome to a duplicate bookkeeping request; never invoke the backend to regenerate it. | +| Cancel and completion race | Serialize their admissible state cuts. A control receipt records a request; only evidence permits cancellation settlement. A validated completion can win. | +| Late result after terminal or owner replacement | Validate scope/generation and attach through the existing linked-resolution authority where admissible. Do not rewrite the immutable terminal parent or allow a stale writer to commit. | +| Engine history expired/reset/manually redriven | Consult retained ledger identity/claims first. History absence is no new permission. Missing or inconsistent binding closes admission for that scope. | + +The probe evidence covers representative duplicate effects, cancellation and +publication cuts, not this entire end-to-end protocol. The claim mechanism +requires a transactional provider and tests against the actual engine/backend +composition before it can be advertised. + +## Supervision and bounded ownership + +An independently authorized supervisor binds its own actor and unique request +identity to the original operation scope. Persist the request at the same state +authority, then dispatch it through the separate control workers. Backend +refusal is a disposition; timeout is uncertainty; neither is cessation. Stop +admission immediately where permitted while separately accounting for active +effects. An observation/control callback must not acquire the blocked effect +worker's capacity or long-lived permit. + +Provision independent bounded admission, worker, store-connection and rejection- +audit capacity for control. Give every callback/commit/drain a finite stage +budget and retain unresolved reservations after budget expiry. Service/store +unavailability can prevent semantic mutation; it must still yield a bounded +unavailable response. No service queue promises physical interruption during its +own outage. Use backend containment evidence when a stronger admitted guarantee +requires it, or refuse that operation before execution. + +Exactly one current state writer owns each admitted target/run scope. For the +future remote provider, owner admission/renewal uses PostgreSQL transactions and +a monotonically changing generation; every mutation compares it. A partitioned +writer that cannot validate current ownership cannot commit or issue a new +grant. This is a store fence only. A worker holding an earlier confirmed grant +may still act, so ownership replacement preserves all consumed claims and +effect reservations. A new owner may reconcile/control those invocations, but +cannot admit conflicting effects until the required cessation or isolation is +established. No lease timeout alone releases them. + +Federated homes cannot both own the same effect scope. Migration requires an +explicit admission-stop and ownership handoff, retained state/claims, scoped +credential revocation, and reconciliation of outstanding effects. If a backend +cannot establish a required fence/cessation, keep the scope quarantined. Parallel +execution uses the existing admitted independence/resource-reservation rules; +it cannot be enabled simply by adding workers. + +## Restore, retention and upgrades + +Temporal and the ledger have independent persistence transactions and backups, +even when hosted on one PostgreSQL installation. Restore into **closed effect +admission** with fresh deployment authority. Establish backend inventory and +effect scope, restore compatible command/history bindings, and reconcile each +potentially active invocation before reopening. A ledger rolled back before a +consumed claim is especially dangerous: its apparently absent record is not +absence of an external effect. Reconcile the whole affected deployment scope, +including resources no longer represented in the backup. If that inventory or +containment cannot be established, do not reopen that scope. + +Do not allow restored engine tasks to dispatch against an empty/new ledger. +Mismatch of deployment generation, operation commitment or version fails closed. +Backend state backups, trial evidence and semantic-clock continuity have their +own restore obligations; an orchestration restore cannot certify them. Fresh +trials require their own isolation/clean-state admission after recovery. + +Retain deduplication/claim and resolution records at least as long as any +possible effect, replay or admitted recovery requires them. Engine history +retention alone cannot govern their deletion. Unsupported recovery windows are +refused, not bridged by blind execution. Versioned workflow workers and ledger +carrier migrations must preserve these bindings and remaining budgets. Do not +roll over a history until pending controls and claimed effects have an explicit +handoff; retain old workers for outstanding compatible histories. diff --git a/docs/research/execution-architecture/experiment-report.md b/docs/research/execution-architecture/experiment-report.md new file mode 100644 index 000000000..65f3f0e64 --- /dev/null +++ b/docs/research/execution-architecture/experiment-report.md @@ -0,0 +1,162 @@ +# Experiment report — 2026-09-23 + +Status: bounded mechanism evidence supporting +[ADR-113](../../decisions/adrs/adr-113-reusable-execution-machinery.md). The raw +[evidence](evidence.json) contains counters, event ordering, observed PIDs, +framework results, versions, and SHA-256 digests of all thirteen executed source +files. Those digests were checked against the retained source after retrieval. +This is one final complete run, preceded by setup/debug runs; it is not a +statistical performance study. + +## Interpretation for the design + +Backend cooperation is a prerequisite, including truthful refusal when a +request exceeds capability or contextual limits. Deliberately non-interrupting +fixtures test the meaning of cancellation signals; they do not establish a +disadvantage of needing backend participation. Runtime responsibility for +requirements, supervision and settlement remains unchanged. + +The observations identify limits to manage in a composition. For example, DBOS +successfully returning an indeterminate domain result can be correct: RAE must +retain the uncertainty and decide the permitted next action. Likewise, duplicate +effects under configured recovery identify where repeat permission must be +enforced. These results do not by themselves reject either engine. + +The missing evidence is a complete cooperating runtime/backend interaction. +The [selected design](../../decisions/adrs/adr-113-reusable-execution-machinery.md) +names the next vertical slice and +its checks. Further investigation should resolve a specific design uncertainty; +no candidate is expected to remove every limitation. + +The declared federated workload changes the architectural recommendation, not +these measurements. No probe tested 15 tenants or the stated scenario scale. +Temporal's selected use rests on the fit of its distributed execution mechanisms; +capacity, isolation and the RAE integration still need their own evidence. + +## Environment and measurement + +AWS profile `catalyst-dev`, region `us-east-2`, one `t3.xlarge` with standard CPU +credits, Ubuntu 24.04 image `ami-00adec9774170bad2`, Python 3.12.3, and one encrypted +40-GiB gp3 root disk. A dedicated `172.31.240.0/24` subnet was created in the +default VPC. Systems Manager provided access; the security group had no inbound +rules and only outbound TCP 80/443. Existing default routes were reused without +modification. No NAT gateway, EIP, load balancer, bucket or production backend +was created. The host had IMDSv2 required and a temporary SSM-only instance role. + +Python versions: AnyIO 4.15.1, DBOS 3.0.0, Temporal SDK 1.33.0, Ray 2.58.0, +SimPy 4.1.2. Temporal's SDK started CLI 1.9.1 / server 1.32.0 with **in-memory +persistence**. ROS used Jazzy rclpy 7.1.12 and action tutorial interfaces 0.33.11; +BehaviorTree.CPP was 4.10.0, built with GNU C++ 13.3.0. The ROS base image was +`ros@sha256:c3706ef0a0aa45413c07803cf433602f543b22e45b4855f6fca955c2d8ecc4e8`. +Full Python/ROS package versions are retained in the evidence. Container-local +image builds added the documented packages; the base digest alone does not pin +every apt dependency. A complete supply-chain lock was not produced. + +Most probes send synthetic HTTP requests to a loopback service in a separate +process. That witness counts receipt separately from effect application and +continues after caller loss. It is outside the candidate worker but on the same +host: this tests a process/acknowledgement boundary, not network partitions or +independent machine failure. Its state is not crash-durable. The BehaviorTree +probe instead uses an independently running thread and atomic counter; it does +not claim the stronger separate-process witness boundary. + +Each candidate command had a host-side timeout, plus an SSM execution timeout. +Cloud lifecycle ownership had a two-hour deadline and teardown in `finally`; +the host also had a 150-minute shutdown/termination timer. Deadline expiry is a +cost/cleanup backstop, not a result or proof of backend cessation. The cleanup +record is retained separately in [cloud-cleanup.json](cloud-cleanup.json). + +## Observations and their limits + +| Probe | Observed result | What follows, and what does not | +| --- | --- | --- | +| AnyIO thread abandonment | Cancellation returned with zero effects observed; later the witness counted one effect. Caller process remained alive. | Waiting can be abandoned; the thread/remote work must not be classified as stopped from that return. | +| AnyIO cancellable process | Cancellation returned; the recorded caller PID no longer existed; the witness later counted one effect. | Local process termination was established in this run and still did not fence the dispatched external request. | +| AnyIO occupied capacity | A status callback sharing the sole occupied limiter expired at the chosen 100-ms bound. An independent limiter returned status before the effect. | Separate capacity is a useful mechanism. This does not test RAE HTTP admission, store locks, audit queues, CPU starvation or end-to-end bounded supervision. | +| DBOS unguarded recovery | Worker exited with code 73 after one effect but before step checkpoint. Restart recovered the workflow and produced a second effect; engine status was `SUCCESS`. | Recovery works; a non-idempotent unfinished step can repeat. It does not authorize that repetition under RAE rules. | +| DBOS experimental guard | Same crash boundary, but a retained invocation marker prevented the second external call. Returned value was `indeterminate-no-second-invocation`; DBOS status was still `SUCCESS`. | An application guard can prevent replay in this scenario. Engine success and RAE operation success cannot be equated. The marker is test scaffolding, not a complete admission/reconciliation protocol, durable power-loss proof, or multi-owner fence. | +| DBOS cancellation | Status became `CANCELLED` before the in-flight synchronous step's effect. The effect subsequently occurred; the following step did not execute. | Cancellation prevented further steps, not the already running effect. Preemptible async steps were not tested. | +| Temporal worker crash, two attempts allowed | First worker died after one effect. Replacement worker processed the retried activity; total effects became two and the workflow completed. | Durable dispatch survived worker loss with the service alive. This is configured retry, not an unavoidable behavior and not service crash-recovery evidence. | +| Temporal worker crash, one attempt allowed | Effect count remained one; workflow retrieval raised `WorkflowFailureError`. | Retry configuration prevented repetition. Failure status alone still did not describe the externally applied effect. | +| Temporal uncooperative cancellation | With `WAIT_CANCELLATION_COMPLETED`, no heartbeats, and a pending effect, the request did not stop work. The effect occurred and the workflow ended `COMPLETED`. | Cancellation request and final completion are distinct. This configuration is not evidence about all SDK activity types or cancellation policies. | +| Temporal cooperative cancellation | The driver waited for an actual heartbeat before requesting cancellation. Workflow ended `CANCELED`, with no synthetic effect invoked. | Heartbeat/cooperative interruption is useful when the activity participates. This zero-effect activity does not prove interruption after a remote effect was dispatched. | +| Ray configured actor restart/retry | `max_restarts=1`, `max_task_retries=1`: crash after effect led to actor replacement and two effects. | Reuse of actor recovery is real, but RAE must control repeat invocation. Ray actor defaults were not being tested. | +| Ray blocked regular actor | A ping on the same actor did not complete within the selected 100-ms observation. Killing that actor did not prevent the witness's later effect. | Same-mailbox control can queue behind work; actor death is not a downstream fence. Threaded/async actors and concurrency groups were not tested. | +| ROS accepted cancel | The action client observed one cancelling goal while effects were still zero. The deliberately non-interrupting server later applied the effect and reported terminal status 5 (`CANCELED`). | Protocol acknowledgement and server-declared terminal status are not independent effect evidence. Known residual effects are possible. This is not a defect claim against ROS. | +| ROS refused cancel | The callback rejected cancellation, the cancelling-goal list was empty, and the operation later reported status 4 (`SUCCEEDED`) with one effect. | Contextual refusal can be carried. In this run the response's numeric return code was zero in both cases: inspect the complete response, not that field alone. | +| SimPy time/interrupt | Virtual time remained zero while approximately 100 ms of apparatus time passed. An interrupt was processed at virtual time zero. A blocking callback held `step()` until explicitly released. | Cooperative virtual-time control can be reused; apparatus supervision needs independent progress. No multi-clock integration or restart model was tested. | +| BehaviorTree.CPP halt | Tree initially returned `RUNNING`; `haltTree()` invoked the action's halt hook once. The deliberately non-stopping hook returned before a background effect incremented the counter. | The hook supplies lifecycle control, not automatic interruption of arbitrary work. An application-specific stopping hook remains possible. | +| SQLite publication race | Competing success/cancel updates used one conditional transition and atomic audit transaction. One won; one terminal state and its matching audit persisted. A simulated lost acknowledgement after a separate commit was resolved by reopening and reading the committed row. | Demonstrates a small storage mechanism. It is not RAE's production store, an unbounded race proof, an injected disk failure or an external-effect fence. | + +The two race contenders were both admissible in the storage fixture. It does +not implement the preceding domain question of whether cancellation has enough +evidence to be admissible. Raw monotonic timestamps are host-relative event +coordinates, not portable restart timestamps. Reported durations are observations +on a burstable host, never latency guarantees or comparative benchmarks. + +## Verification and setup corrections + +Three witness unit tests establish distinct request/effect counts, independently +counted duplicate invocations, and an effect occurring after caller timeout. The +tests first failed because the witness implementation did not exist, then passed +locally and on the experiment host. Candidate scripts assert their specific +measurement expectations and fail if those observations differ. The final run +completed all seven candidates plus the publication probe, including a real +BehaviorTree.CPP compile and actual ROS action client/server exchanges. + +The initial DBOS probe incorrectly passed a `timeout` keyword to a polling +handle; its public call was corrected and the outer process watchdog retained. +The ROS-packaged BehaviorTree.CPP exports `behaviortree_cpp::behaviortree_cpp`, +not the initially assumed CMake target. That fixture was corrected. Neither +setup failure was treated as an architectural finding. The final run used a +fresh directory and fresh synthetic state. Worker code, witness instrumentation, +and final source hashes correspond to the retained evidence. + +## Outstanding evidence before any stronger claim + +No production RAE adapter was implemented. No real backend was interrupted. +Multi-host partitions, machine/store loss, credential handling through framework +checkpoints, tenant isolation, code/schema upgrade recovery, production service +operations, broad workload coverage and full HTTP/store/audit saturation remain +untested. No deployment-cost or throughput ranking follows from these probes. +They are sufficient to expose the selected semantic distinctions and inform the +architecture selection, not to certify its production realization. + +## Independent local recheck during delivery + +On 2026-09-23 the same source bytes were executed locally with Python 3.14.4, +AnyIO 4.15.1, SimPy 4.1.2, DBOS 3.0.0 and Temporal SDK 1.33.0. A fresh disposable +directory and virtual environment were used; each command had a 160-second +outer timeout. All three witness tests and the AnyIO, SimPy, DBOS, Temporal and +SQLite publication assertions passed. [local-recheck.json](local-recheck.json) +records commands, versions, source hashes and actual observations independently +of the earlier cloud run. Temporal again used CLI 1.9.1/server 1.32.0 with +in-memory persistence. Its worker fixture uses UnsandboxedWorkflowRunner for +this synthetic probe, not an isolation demonstration. No paid service or cloud +resource was used for the recheck. + +ROS, Ray and BehaviorTree.CPP were not rerun locally; their observations are the +original retained evidence. Neither run measured PostgreSQL crash persistence, +distributed ownership, tenant isolation or a real OT process. The probes expose +mechanism limits under named configurations; they do not choose a failure +response for an experiment, twin, IT or OT deployment. That response belongs to +the admitted contextual author policy, independently of technical retry safety. + +The delivery review identified an assertion gap in the Temporal collector: +`outcome()` records exceptions rather than failing, and zero effects alone can +also result from the cooperative fixture finishing normally. The added +`verify_temporal.py` entry point preserves the original source and validates +terminal status, result/error type and effect count on fresh output. Negative +controls first demonstrated that effect-only checks accepted four injected +timeouts and a cooperative completion without cancellation; all four regression +tests pass with the stronger verifier. A fresh live Temporal run also passed. +The `verified_temporal_followup` record in `local-recheck.json` binds that run to +its source hashes separately from both earlier evidence records. + +The publish hook classified the synthetic `dbos-r2-cancel` workflow identifier +as a generic API key. The repository scanner exception matches only that exact +value in `experiments/probe_dbos.py` under that one rule, preserving the original +source hash. A regression invokes the real scanner and verifies that a different +credential-shaped value and a different rule still trigger in the same file, +and that the fixture identifier still triggers outside that file. The test +failed before the exception and passes after it; no actual credential is used. diff --git a/docs/research/execution-architecture/experiments/CMakeLists.txt b/docs/research/execution-architecture/experiments/CMakeLists.txt new file mode 100644 index 000000000..56bc7d51c --- /dev/null +++ b/docs/research/execution-architecture/experiments/CMakeLists.txt @@ -0,0 +1,6 @@ +cmake_minimum_required(VERSION 3.16) +project(rae1350_bt_probe LANGUAGES CXX) +find_package(behaviortree_cpp REQUIRED) +add_executable(probe_bt probe_bt.cpp) +target_compile_features(probe_bt PRIVATE cxx_std_17) +target_link_libraries(probe_bt PRIVATE behaviortree_cpp::behaviortree_cpp) diff --git a/docs/research/execution-architecture/experiments/README.md b/docs/research/execution-architecture/experiments/README.md new file mode 100644 index 000000000..3ff8968f2 --- /dev/null +++ b/docs/research/execution-architecture/experiments/README.md @@ -0,0 +1,76 @@ +# Running the bounded probes + +These files are non-production experiment fixtures. They use only synthetic +effects and intentionally crash their own worker processes. Use a fresh, +disposable working directory: `.probe-state/` and `results/` are generated there. +The seven candidates are not production dependencies of RAE. + +## Python probes + +Use Python 3.12 with the versions in [evidence.json](../evidence.json). The tested +direct dependencies are AnyIO 4.15.1, DBOS 3.0.0, Temporal SDK 1.33.0, Ray 2.58.0 +and SimPy 4.1.2. Run each probe in a fresh copied source directory or preserve +its initial empty `.probe-state/`; reusing crash markers invalidates the intended +first-attempt boundary. All programs should be supervised by an outer timeout. + +```sh +python -m unittest test_witness +timeout 150 python probe_workers.py +timeout 150 python probe_simpy.py +timeout 150 python probe_dbos.py +timeout 150 python probe_ray.py +timeout 150 python verify_temporal.py +timeout 15 python probe_publication.py +``` + +The Temporal probe starts a local development service through its SDK. Initial +startup downloads its development binary; record the selected server version. +This run used CLI 1.9.1/server 1.32.0 with in-memory persistence. Only worker +process loss is injected. The SQLite probe is a separate toy transaction model, +not the production RAE store. DBOS uses its own SQLite files; its guard is a +synthetic retained claim marker owned by the experiment. + +`verify_temporal.py` runs the original Temporal probe, then checks its expected +terminal statuses and result/error types as well as effect counts. Use this +entry point for new runs: the original collector alone can absorb a timeout or +record zero effects without establishing cooperative cancellation. The original +thirteen sources remain unchanged for historical hash reproducibility. +From the repository root, its negative-control regressions run with: + +```sh +python3 -m unittest discover -s docs/research/execution-architecture/experiments -p test_temporal_verification.py -v +``` + +## ROS and BehaviorTree.CPP + +Use the ROS Jazzy environment and package versions in the report. Install the +action tutorial interfaces and BehaviorTree.CPP development package along with +CMake and a C++17 compiler. Source the distribution's setup, restrict discovery +to localhost, then run: + +```sh +timeout 45 python3 probe_ros.py +cmake -S . -B build +cmake --build build -j2 +timeout 10 build/probe_bt +``` + +The tested ROS CMake export is `behaviortree_cpp::behaviortree_cpp`. The behavior +tree action intentionally leaves its synthetic background thread running when +halted; this isolates the contract of the halt hook from application-provided +stopping logic. The ROS action server likewise deliberately does not interrupt +the already sent synthetic request. Neither is a production backend example. + +## Evidence handling + +Python probes write JSON records under `results/`; the C++ probe emits JSON on +stdout. Preserve exact source hashes, package versions and the external witness +observations before terminating the disposable host. Record failed setup attempts +separately from candidate behavior. Never treat a command's successful exit alone +as proof of cancellation, crash durability or exactly-once external effects. + +Cloud execution used a dedicated, tagged stack with no inbound network rules, +temporary SSM management identity, encrypted delete-on-termination storage, and +bounded teardown. The [report](../experiment-report.md) and cleanup evidence +describe the executed environment; running these files does not provision AWS +resources automatically. diff --git a/docs/research/execution-architecture/experiments/dbos_worker.py b/docs/research/execution-architecture/experiments/dbos_worker.py new file mode 100644 index 000000000..e24af0320 --- /dev/null +++ b/docs/research/execution-architecture/experiments/dbos_worker.py @@ -0,0 +1,100 @@ +"""Disposable DBOS worker; deliberately crashes after its synthetic effect.""" + +import json +import os +from pathlib import Path +import sys +import time + +from dbos import DBOS, SetWorkflowID + +from witness import effect, observe, wait_for + + +def durable_marker(path): + with path.open("x") as stream: + stream.write("claimed\n") + stream.flush() + os.fsync(stream.fileno()) + + +@DBOS.step() +def invoke(key, url, guarded, crash, delay): + root = Path(".probe-state") + claim = root / (key + "-claim") + if guarded: + try: + durable_marker(claim) + except FileExistsError: + return "indeterminate-no-second-invocation" + effect(url, key, delay=delay) + crash_marker = root / (key + "-crashed") + if crash and not crash_marker.exists(): + durable_marker(crash_marker) + os._exit(73) + return "effect-returned" + + +@DBOS.step() +def following_step(key): + Path(".probe-state", key + "-following").touch() + + +@DBOS.workflow() +def operation(key, url, guarded, crash, delay): + result = invoke(key, url, guarded, crash, delay) + following_step(key) + return result + + +def main(): + phase, key, url, guard = sys.argv[1:] + Path(".probe-state").mkdir(exist_ok=True) + DBOS( + config={ + "name": "rae1350", + "system_database_url": f"sqlite:///.probe-state/{key}.sqlite", + "run_admin_server": False, + "log_level": "ERROR", + } + ) + DBOS.launch() + try: + if phase == "recover": + handle = DBOS.retrieve_workflow(key) + else: + with SetWorkflowID(key): + handle = DBOS.start_workflow( + operation, + key, + url, + guard == "guarded", + phase != "cancel", + 0.8 if phase == "cancel" else 0, + ) + if phase == "cancel": + wait_for(url, key) + before = observe(url, key) + DBOS.cancel_workflow(key) + status = DBOS.get_workflow_status(key).status + time.sleep(1.1) + result = { + "at_cancel": before, + "status_after_cancel": status, + "later": observe(url, key), + "following_step_executed": Path( + ".probe-state", key + "-following" + ).exists(), + } + else: + result = { + "result": handle.get_result(), + "status": DBOS.get_workflow_status(key).status, + } + Path("results", key + "-worker.json").write_text(json.dumps(result, indent=2)) + finally: + DBOS.destroy() + + +if __name__ == "__main__": + main() diff --git a/docs/research/execution-architecture/experiments/probe_bt.cpp b/docs/research/execution-architecture/experiments/probe_bt.cpp new file mode 100644 index 000000000..0735c928a --- /dev/null +++ b/docs/research/execution-architecture/experiments/probe_bt.cpp @@ -0,0 +1,48 @@ +// Test fixture: halting the tree does not stop an independently running effect. +#include +#include +#include +#include +#include + +std::atomic effects{0}; +std::atomic halts{0}; +std::thread worker; + +class Action : public BT::StatefulActionNode { + public: + Action(const std::string& name, const BT::NodeConfig& config) + : BT::StatefulActionNode(name, config) {} + static BT::PortsList providedPorts() { return {}; } + BT::NodeStatus onStart() override { + worker = std::thread([] { + std::this_thread::sleep_for(std::chrono::milliseconds(500)); + effects++; + }); + return BT::NodeStatus::RUNNING; + } + BT::NodeStatus onRunning() override { + return effects.load() ? BT::NodeStatus::SUCCESS : BT::NodeStatus::RUNNING; + } + void onHalted() override { halts++; } +}; + +int main() { + BT::BehaviorTreeFactory factory; + factory.registerNodeType("Effect"); + auto tree = factory.createTreeFromText( + R"()"); + auto status = tree.tickOnce(); + auto start = std::chrono::steady_clock::now(); + tree.haltTree(); + auto elapsed = std::chrono::duration(std::chrono::steady_clock::now() - start).count(); + const int at_halt = effects.load(); + worker.join(); + std::cout << "{\"probe\":\"behaviortree\",\"initial_running\":" + << (status == BT::NodeStatus::RUNNING ? "true" : "false") + << ",\"halt_callback_count\":" << halts.load() + << ",\"effects_at_halt\":" << at_halt + << ",\"effects_after_halt\":" << effects.load() + << ",\"halt_return_seconds\":" << elapsed << "}\n"; + return (halts == 1 && at_halt == 0 && effects == 1) ? 0 : 1; +} diff --git a/docs/research/execution-architecture/experiments/probe_dbos.py b/docs/research/execution-architecture/experiments/probe_dbos.py new file mode 100644 index 000000000..be8d25d71 --- /dev/null +++ b/docs/research/execution-architecture/experiments/probe_dbos.py @@ -0,0 +1,50 @@ +import json +from pathlib import Path +import subprocess +import sys + +from witness import Witness, observe, record + + +def worker(phase, key, url, guard): + return subprocess.run( + [sys.executable, "dbos_worker.py", phase, key, url, guard], + capture_output=True, + text=True, + timeout=45, + ) + + +def main(): + Path("results").mkdir(exist_ok=True) + results = {} + with Witness() as witness: + for guard in ("unguarded", "guarded"): + key = "dbos-r2-" + guard + first = worker("start", key, witness.url, guard) + assert first.returncode == 73, first.stderr + after_crash = observe(witness.url, key) + second = worker("recover", key, witness.url, guard) + assert second.returncode == 0, second.stderr + results[guard] = { + "worker_crash_exit": first.returncode, + "after_crash": after_crash, + "after_recovery": observe(witness.url, key), + "worker": json.loads(Path("results", key + "-worker.json").read_text()), + } + assert results[guard]["after_recovery"]["effects"] == ( + 2 if guard == "unguarded" else 1 + ) + key = "dbos-r2-cancel" + cancellation = worker("cancel", key, witness.url, "unguarded") + assert cancellation.returncode == 0, cancellation.stderr + results["cancellation"] = json.loads( + Path("results", key + "-worker.json").read_text() + ) + assert results["cancellation"]["later"]["effects"] == 1 + assert not results["cancellation"]["following_step_executed"] + record("dbos", results) + + +if __name__ == "__main__": + main() diff --git a/docs/research/execution-architecture/experiments/probe_publication.py b/docs/research/execution-architecture/experiments/probe_publication.py new file mode 100644 index 000000000..fb6296e5f --- /dev/null +++ b/docs/research/execution-architecture/experiments/probe_publication.py @@ -0,0 +1,81 @@ +"""Small SQLite boundary experiment, NOT a production RAE refinement proof.""" + +from pathlib import Path +import sqlite3 +import threading + +from witness import record + + +def main(): + Path(".probe-state").mkdir(exist_ok=True) + path = ".probe-state/publication.sqlite" + connection = sqlite3.connect(path) + connection.executescript(""" + PRAGMA journal_mode=WAL; + CREATE TABLE outcome (id INTEGER PRIMARY KEY, state TEXT, revision INTEGER); + CREATE TABLE audit (state TEXT); + INSERT INTO outcome VALUES (1, 'RUNNING', 0); + """) + connection.close() + barrier = threading.Barrier(3) + cuts = [] + lock = threading.Lock() + + def settle(state): + db = sqlite3.connect(path, timeout=3) + barrier.wait(timeout=3) + with db: + cursor = db.execute( + "UPDATE outcome SET state=?,revision=revision+1 WHERE id=1 AND state='RUNNING' AND revision=0", + (state,), + ) + if cursor.rowcount: + db.execute("INSERT INTO audit VALUES (?)", (state,)) + with lock: + cuts.append({"requested": state, "won": cursor.rowcount == 1}) + db.close() + + threads = [ + threading.Thread(target=settle, args=(state,)) + for state in ("SUCCEEDED", "CANCELLED") + ] + for thread in threads: + thread.start() + barrier.wait(timeout=3) + for thread in threads: + thread.join(5) + assert not thread.is_alive() + db = sqlite3.connect(path) + state = db.execute("SELECT state,revision FROM outcome").fetchone() + audit = db.execute("SELECT state FROM audit").fetchall() + assert sum(cut["won"] for cut in cuts) == 1 + assert audit == [(state[0],)] and state[1] == 1 + # A lost acknowledgement after commit must be resolved from persisted state. + try: + with db: + db.execute("INSERT INTO outcome VALUES (2,'SUCCEEDED',1)") + db.execute("INSERT INTO audit VALUES ('ACK-LOST-SUCCESS')") + raise TimeoutError("synthetic acknowledgement loss AFTER durable commit") + except TimeoutError: + db.close() + reopened = sqlite3.connect(path) + readback = reopened.execute( + "SELECT state,revision FROM outcome WHERE id=2" + ).fetchone() + reopened.close() + assert readback == ("SUCCEEDED", 1) + record( + "publication", + { + "race": cuts, + "terminal": state, + "audit": audit, + "lost_ack_readback": readback, + "scope": "SQLite toy transaction boundary; no external fencing claim", + }, + ) + + +if __name__ == "__main__": + main() diff --git a/docs/research/execution-architecture/experiments/probe_ray.py b/docs/research/execution-architecture/experiments/probe_ray.py new file mode 100644 index 000000000..31623ef07 --- /dev/null +++ b/docs/research/execution-architecture/experiments/probe_ray.py @@ -0,0 +1,61 @@ +import os +from pathlib import Path +import time + +import ray + +from witness import Witness, effect, observe, record, wait_for + + +@ray.remote(max_restarts=1, max_task_retries=1, num_cpus=1) +class Worker: + def __init__(self, marker_dir): + self.marker_dir = marker_dir + + def execute(self, url, key, crash=False, delay=0): + effect(url, key, delay=delay) + marker = Path(self.marker_dir, key) + if crash and not marker.exists(): + marker.touch() + os._exit(73) + return "completed" + + def ping(self): + return "responsive" + + +def main(): + Path(".probe-state").mkdir(exist_ok=True) + with Witness() as witness: + ray.init(num_cpus=2, include_dashboard=False, logging_level="ERROR") + try: + worker = Worker.remote(str(Path(".probe-state").resolve())) + result = ray.get( + worker.execute.remote(witness.url, "ray-retry", True), timeout=45 + ) + retry = observe(witness.url, "ray-retry") + assert retry["effects"] == 2 + worker.execute.remote(witness.url, "ray-kill", False, 1) + wait_for(witness.url, "ray-kill") + ping = worker.ping.remote() + ready, _ = ray.wait([ping], timeout=0.1) + at_kill = observe(witness.url, "ray-kill") + ray.kill(worker, no_restart=True) + time.sleep(1.2) + after = observe(witness.url, "ray-kill") + assert after["effects"] == 1 + record( + "ray", + { + "configured_retry": {"result": result, "witness": retry}, + "control_on_same_actor_blocked": not ready, + "at_kill": at_kill, + "after_kill": after, + }, + ) + finally: + ray.shutdown() + + +if __name__ == "__main__": + main() diff --git a/docs/research/execution-architecture/experiments/probe_ros.py b/docs/research/execution-architecture/experiments/probe_ros.py new file mode 100644 index 000000000..7fd5a3a76 --- /dev/null +++ b/docs/research/execution-architecture/experiments/probe_ros.py @@ -0,0 +1,97 @@ +"""Actual ROS action messages against a deliberately non-interrupting server.""" + +import threading +import time + +from action_tutorials_interfaces.action import Fibonacci +import rclpy +from rclpy.action import ActionClient, ActionServer, CancelResponse, GoalResponse +from rclpy.callback_groups import ReentrantCallbackGroup +from rclpy.executors import MultiThreadedExecutor +from rclpy.node import Node + +from witness import Witness, effect, observe, record, wait_for + + +def finish(future, timeout=10): + end = time.monotonic() + timeout + while not future.done(): + if time.monotonic() >= end: + raise TimeoutError("ROS future") + time.sleep(0.005) + return future.result() + + +def main(): + with Witness() as witness: + rclpy.init() + node = Node("rae1350_probe") + group = ReentrantCallbackGroup() + results = {} + + def execute(goal): + key = "ros-accept" if goal.request.order == 1 else "ros-refuse" + effect(witness.url, key, delay=0.8) + result = Fibonacci.Result() + result.sequence = [1] + if goal.is_cancel_requested: + goal.canceled() + else: + goal.succeed() + return result + + server = ActionServer( + node, + Fibonacci, + "rae1350", + execute, + goal_callback=lambda _goal: GoalResponse.ACCEPT, + cancel_callback=lambda goal: ( + CancelResponse.ACCEPT + if goal.request.order == 1 + else CancelResponse.REJECT + ), + callback_group=group, + ) + client = ActionClient(node, Fibonacci, "rae1350", callback_group=group) + executor = MultiThreadedExecutor(num_threads=4) + executor.add_node(node) + thread = threading.Thread(target=executor.spin) + thread.start() + try: + assert client.wait_for_server(timeout_sec=10) + for order, key in ((1, "ros-accept"), (2, "ros-refuse")): + goal = Fibonacci.Goal() + goal.order = order + handle = finish(client.send_goal_async(goal)) + assert handle.accepted + wait_for(witness.url, key) + before = observe(witness.url, key) + start = time.monotonic() + cancel = finish(handle.cancel_goal_async()) + elapsed = time.monotonic() - start + at_ack = observe(witness.url, key) + result = finish(handle.get_result_async()) + results[key] = { + "at_request": before, + "at_cancel_reply": at_ack, + "cancel_reply_seconds": elapsed, + "cancel_return_code": cancel.return_code, + "goals_canceling": len(cancel.goals_canceling), + "terminal_status": result.status, + "later": observe(witness.url, key), + } + assert results[key]["later"]["effects"] == 1 + assert len(cancel.goals_canceling) == (1 if order == 1 else 0) + finally: + executor.shutdown(timeout_sec=5) + thread.join(5) + client.destroy() + server.destroy() + node.destroy_node() + rclpy.shutdown() + record("ros", results) + + +if __name__ == "__main__": + main() diff --git a/docs/research/execution-architecture/experiments/probe_simpy.py b/docs/research/execution-architecture/experiments/probe_simpy.py new file mode 100644 index 000000000..6262e33cc --- /dev/null +++ b/docs/research/execution-architecture/experiments/probe_simpy.py @@ -0,0 +1,59 @@ +"""Virtual-time pause, cooperative interrupt, and a blocking event callback.""" + +import threading +import time + +import simpy + +from witness import record + + +def main(): + env = simpy.Environment() + log = [] + + def simulation(): + try: + yield env.timeout(10) + log.append(["finished", env.now]) + except simpy.Interrupt: + log.append(["interrupted", env.now]) + + task = env.process(simulation()) + env.step() + start = time.monotonic() + time.sleep(0.1) # Simulation deliberately not advanced. + elapsed = time.monotonic() - start + paused_time = env.now + task.interrupt("experiment") + env.step() + blocking = simpy.Environment() + callback_started = threading.Event() + callback_released = threading.Event() + + def blocked(): + callback_started.set() + callback_released.wait(2) + yield blocking.timeout(1) + + blocking.process(blocked()) + thread = threading.Thread(target=blocking.step) + thread.start() + assert callback_started.wait(1) + callback_blocks_step = thread.is_alive() + callback_released.set() + thread.join(3) + assert paused_time == 0 and log == [["interrupted", 0]] + record( + "simpy", + { + "paused_virtual_time": paused_time, + "apparatus_elapsed_seconds": elapsed, + "interrupt_trace": log, + "blocking_callback_holds_step": callback_blocks_step, + }, + ) + + +if __name__ == "__main__": + main() diff --git a/docs/research/execution-architecture/experiments/probe_temporal.py b/docs/research/execution-architecture/experiments/probe_temporal.py new file mode 100644 index 000000000..8e027b8e7 --- /dev/null +++ b/docs/research/execution-architecture/experiments/probe_temporal.py @@ -0,0 +1,141 @@ +import asyncio +from datetime import timedelta +from pathlib import Path +import subprocess +import sys +import time + +from temporalio.testing import WorkflowEnvironment + +from temporal_worker import Operation +from witness import Witness, observe, record, wait_for + + +def start_worker(address, key): + log = Path("results", key + ".log").open("w") + process = subprocess.Popen( + [sys.executable, "temporal_worker.py", address], stdout=log, stderr=log + ) + log.close() + return process + + +def stop_worker(process): + if process.poll() is None: + process.terminate() + try: + process.wait(timeout=5) + except subprocess.TimeoutExpired: + process.kill() + process.wait(timeout=5) + + +async def outcome(handle): + try: + return {"result": await asyncio.wait_for(handle.result(), timeout=30)} + except Exception as exc: + return {"error_type": type(exc).__name__} + + +async def main(): + Path(".probe-state").mkdir(exist_ok=True) + Path("results").mkdir(exist_ok=True) + results = {} + with Witness() as witness: + async with await WorkflowEnvironment.start_local() as env: + address = env.client.service_client.config.target_host + for attempts in (2, 1): + key = f"temporal-crash-{attempts}" + spec = { + "key": key, + "url": witness.url, + "mode": "crash", + "attempts": attempts, + } + process = start_worker(address, key + "-first") + try: + handle = await env.client.start_workflow( + Operation.run, + spec, + id=key, + task_queue="rae1350", + execution_timeout=timedelta(seconds=35), + ) + await asyncio.to_thread( + wait_for, witness.url, key, "effects", 1, 30 + ) + exit_code = await asyncio.to_thread(process.wait, 10) + assert exit_code == 73 + after_crash = observe(witness.url, key) + finally: + stop_worker(process) + process = start_worker(address, key + "-second") + try: + result = await outcome(handle) + results[key] = { + "worker_crash_exit": exit_code, + "after_crash": after_crash, + "after_recovery": observe(witness.url, key), + **result, + } + assert results[key]["after_recovery"]["effects"] == attempts + finally: + stop_worker(process) + for mode in ("uncooperative", "cooperative"): + key = "temporal-cancel-" + mode + process = start_worker(address, key) + try: + spec = { + "key": key, + "url": witness.url, + "mode": mode, + "attempts": 1, + "delay": 0.8, + } + handle = await env.client.start_workflow( + Operation.run, + spec, + id=key, + task_queue="rae1350", + execution_timeout=timedelta(seconds=15), + ) + if mode == "uncooperative": + await asyncio.to_thread( + wait_for, witness.url, key, "requests", 1, 30 + ) + else: + # Wait for an actual heartbeat, rather than timing worker startup. + end = time.monotonic() + 20 + while time.monotonic() < end: + desc = await handle.describe() + if any( + a.HasField("last_heartbeat_time") + for a in desc.raw_description.pending_activities + ): + break + await asyncio.sleep(0.1) + else: + raise TimeoutError("cooperative activity heartbeat") + before = observe(witness.url, key) + start = time.monotonic() + await handle.cancel() + disposition = await outcome(handle) + elapsed = time.monotonic() - start + await asyncio.sleep(1) + results[key] = { + "at_cancel": before, + "later": observe(witness.url, key), + "settlement_seconds": elapsed, + "workflow_status": (await handle.describe()).status.name, + **disposition, + } + assert results[key]["later"]["effects"] == ( + 1 if mode == "uncooperative" else 0 + ) + finally: + stop_worker(process) + record("temporal", results) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/docs/research/execution-architecture/experiments/probe_workers.py b/docs/research/execution-architecture/experiments/probe_workers.py new file mode 100644 index 000000000..4152edf6f --- /dev/null +++ b/docs/research/execution-architecture/experiments/probe_workers.py @@ -0,0 +1,88 @@ +"""AnyIO thread/process cancellation and occupied execution capacity.""" + +import functools +import os +import time + +import anyio +from anyio import to_process + +from witness import Witness, effect, observe, record, wait_for + + +async def run_case(url, kind): + key = "anyio-" + kind + limiter = anyio.CapacityLimiter(1) + measurement = {} + scope = anyio.CancelScope() + + async def worker(): + with scope: + if kind == "thread": + await anyio.to_thread.run_sync( + functools.partial(effect, url, key, delay=0.7), + abandon_on_cancel=True, + limiter=limiter, + ) + else: + await to_process.run_sync( + effect, url, key, 0.7, cancellable=True, limiter=limiter + ) + + async with anyio.create_task_group() as group: + group.start_soon(worker) + # Witness reads use independent control capacity, not the occupied limiter. + await anyio.to_thread.run_sync(wait_for, url, key) + start = time.monotonic() + measurement["at_cancel"] = observe(url, key) + scope.cancel() + measurement["local_cancel_return_seconds"] = time.monotonic() - start + measurement["after_local_cancel"] = observe(url, key) + pid = measurement["after_local_cancel"]["caller_pids"][0] + try: + os.kill(pid, 0) + measurement["caller_process_alive_after_cancel"] = True + except ProcessLookupError: + measurement["caller_process_alive_after_cancel"] = False + if kind == "process": + assert not measurement["caller_process_alive_after_cancel"] + await anyio.sleep(0.9) + measurement["later"] = observe(url, key) + assert measurement["at_cancel"]["effects"] == 0 + assert measurement["later"]["effects"] == 1 + return measurement + + +async def main(): + with Witness() as witness: + results = { + kind: await run_case(witness.url, kind) for kind in ("thread", "process") + } + limiter = anyio.CapacityLimiter(1) + async with anyio.create_task_group() as group: + + async def occupy(): + await anyio.to_thread.run_sync( + functools.partial(effect, witness.url, "saturation", delay=0.6), + limiter=limiter, + ) + + group.start_soon(occupy) + await anyio.to_thread.run_sync(wait_for, witness.url, "saturation") + with anyio.move_on_after(0.1) as same_lane: + await anyio.to_thread.run_sync(lambda: "status", limiter=limiter) + separate_status = await anyio.to_thread.run_sync( + lambda: "status", limiter=anyio.CapacityLimiter(1) + ) + before_effect = observe(witness.url, "saturation") + results["saturation"] = { + "shared_capacity_control_expired": same_lane.cancel_called, + "separate_capacity_status": separate_status, + "witness_at_separate_status": before_effect, + } + assert same_lane.cancel_called and before_effect["effects"] == 0 + record("anyio", results) + + +if __name__ == "__main__": + anyio.run(main) diff --git a/docs/research/execution-architecture/experiments/temporal_worker.py b/docs/research/execution-architecture/experiments/temporal_worker.py new file mode 100644 index 000000000..b07d2164c --- /dev/null +++ b/docs/research/execution-architecture/experiments/temporal_worker.py @@ -0,0 +1,65 @@ +"""Worker process, deliberately separable from Temporal service and witness.""" + +import asyncio +from datetime import timedelta +import os +from pathlib import Path +import sys + +from temporalio import activity, workflow +from temporalio.client import Client +from temporalio.common import RetryPolicy +from temporalio.worker import Worker, UnsandboxedWorkflowRunner + +from witness import effect + + +@activity.defn +async def invoke(spec: dict) -> str: + if spec["mode"] == "cooperative": + for _ in range(30): + activity.heartbeat("synthetic-progress") + await asyncio.sleep(0.1) + return "finished-without-effect" + await asyncio.to_thread(effect, spec["url"], spec["key"], spec.get("delay", 0)) + marker = Path(".probe-state", spec["key"] + "-crashed") + if spec["mode"] == "crash" and not marker.exists(): + marker.touch() + os._exit(73) + return "effect-returned" + + +@workflow.defn +class Operation: + @workflow.run + async def run(self, spec: dict) -> str: + return await workflow.execute_activity( + invoke, + spec, + start_to_close_timeout=timedelta(seconds=2), + heartbeat_timeout=timedelta(seconds=0.5) + if spec["mode"] == "cooperative" + else None, + retry_policy=RetryPolicy( + maximum_attempts=spec["attempts"], + initial_interval=timedelta(seconds=0.1), + ), + cancellation_type=workflow.ActivityCancellationType.WAIT_CANCELLATION_COMPLETED, + ) + + +async def main(): + client = await Client.connect(sys.argv[1]) + async with Worker( + client, + task_queue="rae1350", + workflows=[Operation], + activities=[invoke], + workflow_runner=UnsandboxedWorkflowRunner(), + max_heartbeat_throttle_interval=timedelta(seconds=0.1), + ): + await asyncio.Event().wait() + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/docs/research/execution-architecture/experiments/test_temporal_verification.py b/docs/research/execution-architecture/experiments/test_temporal_verification.py new file mode 100644 index 000000000..521aee40b --- /dev/null +++ b/docs/research/execution-architecture/experiments/test_temporal_verification.py @@ -0,0 +1,47 @@ +"""Negative controls for the experiment's cancellation and recovery claims.""" + +import copy +import json +from pathlib import Path +import unittest + +from verify_temporal import verify + + +class TemporalVerificationTests(unittest.TestCase): + def setUp(self): + evidence = Path(__file__).resolve().parents[1] / "evidence.json" + self.results = json.loads(evidence.read_text())["results"]["temporal"] + + def test_recorded_observations_satisfy_the_claims(self): + verify(self.results) + + def test_cooperative_completion_without_cancel_is_rejected(self): + result = self.results["temporal-cancel-cooperative"] + result["workflow_status"] = "COMPLETED" + result.pop("error_type") + result["result"] = "finished-without-effect" + with self.assertRaises(AssertionError): + verify(self.results) + + def test_absorbed_timeout_is_rejected_for_every_case(self): + for key in self.results: + with self.subTest(key=key): + results = copy.deepcopy(self.results) + results[key].pop("result", None) + results[key]["error_type"] = "TimeoutError" + with self.assertRaises(AssertionError): + verify(results) + + def test_inconsistent_effect_count_is_rejected_for_every_case(self): + for key in self.results: + with self.subTest(key=key): + results = copy.deepcopy(self.results) + observation = "after_recovery" if "crash" in key else "later" + results[key][observation]["effects"] += 1 + with self.assertRaises(AssertionError): + verify(results) + + +if __name__ == "__main__": + unittest.main() diff --git a/docs/research/execution-architecture/experiments/test_witness.py b/docs/research/execution-architecture/experiments/test_witness.py new file mode 100644 index 000000000..8e8d9d495 --- /dev/null +++ b/docs/research/execution-architecture/experiments/test_witness.py @@ -0,0 +1,37 @@ +"""Narrow measurement checks; these do not certify any candidate framework.""" + +import concurrent.futures +import time +import unittest + +from witness import Witness, effect, observe + + +class WitnessTests(unittest.TestCase): + def test_effect_survives_caller_timeout(self): + with Witness() as witness: + with self.assertRaises(TimeoutError): + effect(witness.url, "lost-ack", delay=0.15, timeout=0.03) + time.sleep(0.25) + self.assertEqual(observe(witness.url, "lost-ack")["effects"], 1) + + def test_duplicate_attempts_are_independently_counted(self): + with Witness() as witness: + with concurrent.futures.ThreadPoolExecutor(2) as pool: + list(pool.map(lambda _: effect(witness.url, "repeat"), range(2))) + self.assertEqual(observe(witness.url, "repeat")["effects"], 2) + + def test_requests_and_effects_are_distinct(self): + with Witness() as witness: + with concurrent.futures.ThreadPoolExecutor(1) as pool: + task = pool.submit(effect, witness.url, "blocked", delay=0.2) + deadline = time.monotonic() + 1 + while not observe(witness.url, "blocked")["requests"]: + self.assertLess(time.monotonic(), deadline) + self.assertEqual(observe(witness.url, "blocked")["effects"], 0) + task.result(timeout=1) + self.assertEqual(observe(witness.url, "blocked")["effects"], 1) + + +if __name__ == "__main__": + unittest.main() diff --git a/docs/research/execution-architecture/experiments/verify_temporal.py b/docs/research/execution-architecture/experiments/verify_temporal.py new file mode 100644 index 000000000..552b9e746 --- /dev/null +++ b/docs/research/execution-architecture/experiments/verify_temporal.py @@ -0,0 +1,44 @@ +"""Verify engine outcomes as well as the retained probe's effect observations.""" + +import asyncio +import json +from pathlib import Path + + +def verify(results): + for attempts in (2, 1): + result = results[f"temporal-crash-{attempts}"] + assert result["after_recovery"]["effects"] == attempts + if attempts == 2: + assert result.get("result") == "effect-returned" + assert "error_type" not in result + else: + assert result.get("error_type") == "WorkflowFailureError" + assert "result" not in result + for mode, count in (("uncooperative", 1), ("cooperative", 0)): + result = results[f"temporal-cancel-{mode}"] + assert result["later"]["effects"] == count + if mode == "cooperative": + assert result["workflow_status"] == "CANCELED" + assert result.get("error_type") == "WorkflowFailureError" + assert "result" not in result + else: + assert result["workflow_status"] == "COMPLETED" + assert result.get("result") == "effect-returned" + assert "error_type" not in result + + +def main(): + # Keep the original source byte-for-byte reproducible; verify its fresh output. + from probe_temporal import main as run_probe + + output = Path("results/temporal.json") + if output.exists(): + raise RuntimeError("use a fresh experiment directory") + asyncio.run(run_probe()) + verify(json.loads(output.read_text())) + print("Temporal recovery and cancellation outcomes verified", flush=True) + + +if __name__ == "__main__": + main() diff --git a/docs/research/execution-architecture/experiments/witness.py b/docs/research/execution-architecture/experiments/witness.py new file mode 100644 index 000000000..ffa2fa058 --- /dev/null +++ b/docs/research/execution-architecture/experiments/witness.py @@ -0,0 +1,114 @@ +"""Synthetic effect service in a process independent of candidate workers.""" + +from collections import defaultdict +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +import json +import multiprocessing +import os +from pathlib import Path +import threading +import time +import urllib.parse +import urllib.request + + +def _serve(connection): + state = defaultdict( + lambda: {"requests": 0, "effects": 0, "events": [], "caller_pids": []} + ) + lock = threading.Lock() + + class Handler(BaseHTTPRequestHandler): + def log_message(self, *_args): + pass + + def reply(self, value): + data = json.dumps(value).encode() + self.send_response(200) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(data))) + self.end_headers() + try: + self.wfile.write(data) + except (BrokenPipeError, ConnectionResetError): + pass + + def do_GET(self): + key = urllib.parse.unquote(self.path[1:]) + with lock: + result = json.loads(json.dumps(state[key])) + self.reply(result) + + def do_POST(self): + data = json.loads(self.rfile.read(int(self.headers["Content-Length"]))) + key = data["key"] + with lock: + state[key]["requests"] += 1 + state[key]["caller_pids"].append(data["caller_pid"]) + state[key]["events"].append(["request", time.monotonic()]) + time.sleep(data.get("delay", 0)) + with lock: + state[key]["effects"] += 1 + state[key]["events"].append(["effect", time.monotonic()]) + time.sleep(data.get("ack_delay", 0)) + self.reply({"applied": True}) + + server = ThreadingHTTPServer(("127.0.0.1", 0), Handler) + connection.send(f"http://127.0.0.1:{server.server_port}") + connection.close() + server.serve_forever() + + +class Witness: + def __enter__(self): + parent, child = multiprocessing.Pipe() + self.process = multiprocessing.Process( + target=_serve, args=(child,), daemon=True + ) + self.process.start() + if not parent.poll(5): + self.process.terminate() + self.process.join(5) + raise TimeoutError("witness startup") + self.url = parent.recv() + parent.close() + child.close() + return self + + def __exit__(self, *_args): + self.process.terminate() + self.process.join(5) + + +def effect(url, key, delay=0, ack_delay=0, timeout=10): + data = json.dumps( + {"key": key, "delay": delay, "ack_delay": ack_delay, "caller_pid": os.getpid()} + ).encode() + req = urllib.request.Request( + url + "/effect", data=data, headers={"Content-Type": "application/json"} + ) + with urllib.request.urlopen(req, timeout=timeout) as response: + return json.load(response) + + +def observe(url, key): + with urllib.request.urlopen( + url + "/" + urllib.parse.quote(key), timeout=3 + ) as response: + return json.load(response) + + +def wait_for(url, key, field="requests", count=1, timeout=10): + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + state = observe(url, key) + if state[field] >= count: + return state + time.sleep(0.01) + raise TimeoutError(f"witness {key} {field}<{count}") + + +def record(name, value): + Path("results").mkdir(exist_ok=True) + Path("results", name + ".json").write_text(json.dumps(value, indent=2) + "\n") + print(json.dumps({"probe": name, **value}), flush=True) diff --git a/docs/research/execution-architecture/local-recheck.json b/docs/research/execution-architecture/local-recheck.json new file mode 100644 index 000000000..73688cda6 --- /dev/null +++ b/docs/research/execution-architecture/local-recheck.json @@ -0,0 +1,702 @@ +{ + "schema": "rae.execution-experiment-local-recheck/v1", + "date": "2026-09-23", + "python": "3.14.4", + "platform": "Linux-6.8.0-117-generic-x86_64-with-glibc2.39", + "packages": { + "anyio": "4.15.1", + "dbos": "3.0.0", + "temporalio": "1.33.0", + "simpy": "4.1.2" + }, + "commands": [ + { + "args": [ + "-m", + "unittest", + "test_witness", + "-v" + ], + "returncode": 0 + }, + { + "args": [ + "probe_workers.py" + ], + "returncode": 0 + }, + { + "args": [ + "probe_simpy.py" + ], + "returncode": 0 + }, + { + "args": [ + "probe_dbos.py" + ], + "returncode": 0 + }, + { + "args": [ + "probe_temporal.py" + ], + "returncode": 0 + }, + { + "args": [ + "probe_publication.py" + ], + "returncode": 0 + } + ], + "source_sha256": { + "probe_temporal.py": "7b12c63f91a0a138553acd7c2f5bad4b0016804581e68eba19caf0604a746517", + "temporal_worker.py": "56529d06d2a23fc4c6d4d20968744e6805a71b45c46b61009277609886f44592", + "test_witness.py": "d6613b18a645861fd6b7322e1650a62b237b6be5e64d9a93f00aae02abd7ac69", + "witness.py": "26c63959b52ea4c6f4e77902a9a18acf0ba97304dd91f7a9ff6feac4571ce51f", + "dbos_worker.py": "113dbc62df2be37924b124ee2c2abdaea7e0ea35a15d8beb5c928d5e49aaca2c", + "probe_dbos.py": "0533b361058deb1813d8e364b04eef27ab35aa692b0474eae9e98cf814adbc68", + "probe_ray.py": "982ed3fe873176f5a4ba3b0c30d8d9c823fdd593af62e6e643f17fac7d1dee9f", + "probe_simpy.py": "9335ab9b9129ddb6cea683a0e94cc35c3ca7d3a868216fe80c9fac617ba30861", + "probe_ros.py": "46f1409a772f007a4dcf788c77821ed870cc80f74e5ee5d75f7c6bee555e2c1b", + "probe_workers.py": "85e03d1d1f43e560e67dc5eba57fb80a016391ebdd4dc83cc9354d78988dc788", + "probe_publication.py": "d874b4c6ba98e7999ee1cbcfc3379917ac70a3cae9f527f873ef16737acc6b61" + }, + "results": { + "dbos-r2-guarded-worker": { + "result": "indeterminate-no-second-invocation", + "status": "SUCCESS" + }, + "dbos": { + "unguarded": { + "worker_crash_exit": 73, + "after_crash": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11240221.53309358 + ], + [ + "effect", + 11240221.533158053 + ] + ], + "caller_pids": [ + 3450563 + ] + }, + "after_recovery": { + "requests": 2, + "effects": 2, + "events": [ + [ + "request", + 11240221.53309358 + ], + [ + "effect", + 11240221.533158053 + ], + [ + "request", + 11240223.071707277 + ], + [ + "effect", + 11240223.071775736 + ] + ], + "caller_pids": [ + 3450563, + 3450599 + ] + }, + "worker": { + "result": "effect-returned", + "status": "SUCCESS" + } + }, + "guarded": { + "worker_crash_exit": 73, + "after_crash": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11240226.248305984 + ], + [ + "effect", + 11240226.248369824 + ] + ], + "caller_pids": [ + 3450654 + ] + }, + "after_recovery": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11240226.248305984 + ], + [ + "effect", + 11240226.248369824 + ] + ], + "caller_pids": [ + 3450654 + ] + }, + "worker": { + "result": "indeterminate-no-second-invocation", + "status": "SUCCESS" + } + }, + "cancellation": { + "at_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 11240231.077331567 + ] + ], + "caller_pids": [ + 3451565 + ] + }, + "status_after_cancel": "CANCELLED", + "later": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11240231.077331567 + ], + [ + "effect", + 11240231.87742288 + ] + ], + "caller_pids": [ + 3451565 + ] + }, + "following_step_executed": false + } + }, + "anyio": { + "thread": { + "at_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 11240216.585383965 + ] + ], + "caller_pids": [ + 3450452 + ] + }, + "local_cancel_return_seconds": 0.0007227752357721329, + "after_local_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 11240216.585383965 + ] + ], + "caller_pids": [ + 3450452 + ] + }, + "caller_process_alive_after_cancel": true, + "later": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11240216.585383965 + ], + [ + "effect", + 11240217.28547216 + ] + ], + "caller_pids": [ + 3450452 + ] + } + }, + "process": { + "at_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 11240217.595452571 + ] + ], + "caller_pids": [ + 3450470 + ] + }, + "local_cancel_return_seconds": 0.00459783710539341, + "after_local_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 11240217.595452571 + ] + ], + "caller_pids": [ + 3450470 + ] + }, + "caller_process_alive_after_cancel": false, + "later": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11240217.595452571 + ], + [ + "effect", + 11240218.29551064 + ] + ], + "caller_pids": [ + 3450470 + ] + } + }, + "saturation": { + "shared_capacity_control_expired": true, + "separate_capacity_status": "status", + "witness_at_separate_status": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 11240218.505874995 + ] + ], + "caller_pids": [ + 3450452 + ] + } + } + }, + "temporal": { + "temporal-crash-2": { + "worker_crash_exit": 73, + "after_crash": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11240235.445636837 + ], + [ + "effect", + 11240235.445696658 + ] + ], + "caller_pids": [ + 3452296 + ] + }, + "after_recovery": { + "requests": 2, + "effects": 2, + "events": [ + [ + "request", + 11240235.445636837 + ], + [ + "effect", + 11240235.445696658 + ], + [ + "request", + 11240238.451852543 + ], + [ + "effect", + 11240238.451910865 + ] + ], + "caller_pids": [ + 3452296, + 3452363 + ] + }, + "result": "effect-returned" + }, + "temporal-crash-1": { + "worker_crash_exit": 73, + "after_crash": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11240248.785291463 + ], + [ + "effect", + 11240248.785349553 + ] + ], + "caller_pids": [ + 3452980 + ] + }, + "after_recovery": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11240248.785291463 + ], + [ + "effect", + 11240248.785349553 + ] + ], + "caller_pids": [ + 3452980 + ] + }, + "error_type": "WorkflowFailureError" + }, + "temporal-cancel-uncooperative": { + "at_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 11240261.121049624 + ] + ], + "caller_pids": [ + 3453487 + ] + }, + "later": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11240261.121049624 + ], + [ + "effect", + 11240261.921132945 + ] + ], + "caller_pids": [ + 3453487 + ] + }, + "settlement_seconds": 0.8042703829705715, + "workflow_status": "COMPLETED", + "result": "effect-returned" + }, + "temporal-cancel-cooperative": { + "at_cancel": { + "requests": 0, + "effects": 0, + "events": [], + "caller_pids": [] + }, + "later": { + "requests": 0, + "effects": 0, + "events": [], + "caller_pids": [] + }, + "settlement_seconds": 0.05188502557575703, + "workflow_status": "CANCELED", + "error_type": "WorkflowFailureError" + } + }, + "simpy": { + "paused_virtual_time": 0, + "apparatus_elapsed_seconds": 0.10007679834961891, + "interrupt_trace": [ + [ + "interrupted", + 0 + ] + ], + "blocking_callback_holds_step": true + }, + "dbos-r2-cancel-worker": { + "at_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 11240231.077331567 + ] + ], + "caller_pids": [ + 3451565 + ] + }, + "status_after_cancel": "CANCELLED", + "later": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11240231.077331567 + ], + [ + "effect", + 11240231.87742288 + ] + ], + "caller_pids": [ + 3451565 + ] + }, + "following_step_executed": false + }, + "publication": { + "race": [ + { + "requested": "CANCELLED", + "won": true + }, + { + "requested": "SUCCEEDED", + "won": false + } + ], + "terminal": [ + "CANCELLED", + 1 + ], + "audit": [ + [ + "CANCELLED" + ] + ], + "lost_ack_readback": [ + "SUCCEEDED", + 1 + ], + "scope": "SQLite toy transaction boundary; no external fencing claim" + }, + "dbos-r2-unguarded-worker": { + "result": "effect-returned", + "status": "SUCCESS" + } + }, + "temporal_environment": { + "cli": "1.9.1", + "server": "1.32.0", + "persistence": "in-memory", + "loss_boundary": "worker process only", + "workflow_runner": "UnsandboxedWorkflowRunner" + }, + "verified_temporal_followup": { + "command": "timeout 150 python verify_temporal.py", + "returncode": 0, + "python": "3.14.4", + "packages": { + "anyio": "4.15.1", + "dbos": "3.0.0", + "temporalio": "1.33.0", + "simpy": "4.1.2" + }, + "temporal_environment": { + "cli": "1.9.1", + "server": "1.32.0", + "persistence": "in-memory", + "loss_boundary": "worker process only", + "workflow_runner": "UnsandboxedWorkflowRunner" + }, + "source_sha256": { + "witness.py": "26c63959b52ea4c6f4e77902a9a18acf0ba97304dd91f7a9ff6feac4571ce51f", + "temporal_worker.py": "56529d06d2a23fc4c6d4d20968744e6805a71b45c46b61009277609886f44592", + "probe_temporal.py": "7b12c63f91a0a138553acd7c2f5bad4b0016804581e68eba19caf0604a746517", + "verify_temporal.py": "628ca0bd7ca28930d46ed1475d5ac38e379e0839df2c2727c5e34bdb5938ba14", + "test_temporal_verification.py": "813620aad181c67b2491e849020ee4f9cc41804386140d29bb4a0e6dc6ef212c" + }, + "negative_control_tests": { + "test_module": "test_temporal_verification", + "passed": 4, + "before_fix": "Effect-only verification incorrectly accepted all four injected timeouts and a cooperative completion without cancellation." + }, + "results": { + "temporal-crash-2": { + "worker_crash_exit": 73, + "after_crash": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11241669.679975357 + ], + [ + "effect", + 11241669.680008808 + ] + ], + "caller_pids": [ + 3495544 + ] + }, + "after_recovery": { + "requests": 2, + "effects": 2, + "events": [ + [ + "request", + 11241669.679975357 + ], + [ + "effect", + 11241669.680008808 + ], + [ + "request", + 11241672.686307032 + ], + [ + "effect", + 11241672.686372045 + ] + ], + "caller_pids": [ + 3495544, + 3495599 + ] + }, + "result": "effect-returned" + }, + "temporal-crash-1": { + "worker_crash_exit": 73, + "after_crash": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11241682.999058144 + ], + [ + "effect", + 11241682.999115843 + ] + ], + "caller_pids": [ + 3495768 + ] + }, + "after_recovery": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11241682.999058144 + ], + [ + "effect", + 11241682.999115843 + ] + ], + "caller_pids": [ + 3495768 + ] + }, + "error_type": "WorkflowFailureError" + }, + "temporal-cancel-uncooperative": { + "at_cancel": { + "requests": 1, + "effects": 0, + "events": [ + [ + "request", + 11241695.3042283 + ] + ], + "caller_pids": [ + 3496116 + ] + }, + "later": { + "requests": 1, + "effects": 1, + "events": [ + [ + "request", + 11241695.3042283 + ], + [ + "effect", + 11241696.10431067 + ] + ], + "caller_pids": [ + 3496116 + ] + }, + "settlement_seconds": 0.8040521163493395, + "workflow_status": "COMPLETED", + "result": "effect-returned" + }, + "temporal-cancel-cooperative": { + "at_cancel": { + "requests": 0, + "effects": 0, + "events": [], + "caller_pids": [] + }, + "later": { + "requests": 0, + "effects": 0, + "events": [], + "caller_pids": [] + }, + "settlement_seconds": 0.07258031889796257, + "workflow_status": "CANCELED", + "error_type": "WorkflowFailureError" + } + } + } +} diff --git a/docs/research/runtime-control-plane/index.md b/docs/research/runtime-control-plane/index.md index 551477ac4..e6471deb9 100644 --- a/docs/research/runtime-control-plane/index.md +++ b/docs/research/runtime-control-plane/index.md @@ -30,3 +30,8 @@ distributed or highly available topology. Profile P3 (coordinated multi-process ownership) is recorded as a seam and an explicit nonclaim; no consistency, coordination, or recovery guarantee named in this set is claimable before its implementation issue lands with tests. + +The subsequent [execution architecture selection (#1350)](../execution-architecture/README.md) +reuses Temporal/PostgreSQL for distributed execution while preserving local +profiles, authored failure/retry policy and RAE's semantic responsibilities. +It does not enable P3 or certify the distributed implementation. diff --git a/implementations/python/tests/test_issue_1350_fixture_secret_scan.py b/implementations/python/tests/test_issue_1350_fixture_secret_scan.py new file mode 100644 index 000000000..e603a160e --- /dev/null +++ b/implementations/python/tests/test_issue_1350_fixture_secret_scan.py @@ -0,0 +1,83 @@ +"""The retained DBOS fixture exception must not hide other credentials.""" + +import hashlib +import json +import subprocess +from pathlib import Path + +import pytest + + +def test_retained_experiment_sources_match_recorded_hashes() -> None: + research = Path(__file__).resolve().parents[3] / "docs/research/execution-architecture" + evidence = json.loads((research / "evidence.json").read_text(encoding="utf-8")) + for record in evidence["source_sha256"]: + expected, filename = record.split(maxsplit=1) + assert hashlib.sha256((research / "experiments" / filename).read_bytes()).hexdigest() == expected, filename + + +@pytest.mark.integration +def test_dbos_fixture_exception_preserves_secret_detection(tmp_path: Path) -> None: + repo = Path(__file__).resolve().parents[3] + scan = tmp_path / "scan" + fixture = scan / "docs/research/execution-architecture/experiments/probe_dbos.py" + fixture.parent.mkdir(parents=True) + known_fixture = 'key = "' + "dbos-" + "r2-cancel" + '"\n' + # Synthetic scanner canaries, never usable credentials. Construct them so + # this regression source itself does not contain a credential-shaped literal. + generic_canary = "Ab9x" + "Q7mR" + "4vN2" + "kL8p" + "W6zT" + aws_canary = "AKIA" + "QW7R" + "TM3P" + "X5NH" + "Z2VS" + fixture.write_text( + known_fixture + f'api_key = "{generic_canary}"\naccess_id = "{aws_canary}"\n', + encoding="utf-8", + ) + elsewhere = scan / "other.py" + elsewhere.write_text(known_fixture, encoding="utf-8") + report = tmp_path / "report.json" + # The installer belongs to the frozen tooling closure, not the runtime's + # Python compatibility environment. Use the existing external-tool boundary. + installed = subprocess.run( + [ + "uv", + "run", + "--project", + str(repo / "implementations/tooling/python"), + "--frozen", + "--no-default-groups", + "python", + "-c", + "from tools.gitleaks_tool import ensure_gitleaks; print(ensure_gitleaks())", + ], + cwd=repo, + capture_output=True, + text=True, + check=True, + timeout=90, + ) + result = subprocess.run( + [ + installed.stdout.strip(), + "dir", + "--config", + str(repo / ".gitleaks.toml"), + "--no-banner", + "--redact", + "--log-level", + "error", + "--report-format", + "json", + "--report-path", + str(report), + str(scan), + ], + capture_output=True, + check=False, + timeout=30, + ) + assert result.returncode == 1 + findings = json.loads(report.read_text(encoding="utf-8")) + observed = {(Path(item["File"]).name, item["StartLine"], item["RuleID"]) for item in findings} + assert (fixture.name, 1, "generic-api-key") not in observed + assert (fixture.name, 2, "generic-api-key") in observed + assert (fixture.name, 3, "aws-access-token") in observed + assert (elsewhere.name, 1, "generic-api-key") in observed diff --git a/tools/policy/historical_identity_records.json b/tools/policy/historical_identity_records.json index cf93c235b..951c556b3 100644 --- a/tools/policy/historical_identity_records.json +++ b/tools/policy/historical_identity_records.json @@ -47,7 +47,7 @@ "record_class": "historical-index", "rationale": "Indexes immutable pre-cutover ADR titles, paths, pins, and amendment summaries without making them current identity surfaces.", "occurrences": 4, - "content_sha256": "5dd455dc3052c4223162f505f25beb53d4c0f24d704b05dfeed08662c8091add" + "content_sha256": "2968ca2de5d6ee5b1eafd3dd949ca5b89694d5686d4c1912676166b11f912a8e" }, { "path": "docs/decisions/adrs/adr-000-use-adrs.md", @@ -495,7 +495,7 @@ "record_class": "historical-index", "rationale": "Indexes immutable pre-cutover ADR titles, paths, pins, and amendment summaries without making them current identity surfaces.", "occurrences": 4, - "content_sha256": "095069130a2b2c48ce336abb4729a3357b586bce4c4abaccf74e2ddb307d9c69" + "content_sha256": "612392c0d80f2c29b3a76340f4c599f269e0194a3fba64dfbd5fde54a777a9c8" }, { "path": "docs/decisions/cage-2-replication-design.md",