Metadata-Version: 2.4
Name: ofplang-run
Version: 0.6.1
Summary: Runner for Object-flow Programming Language workflows.
Author-email: Kazunari Kaizu <kwaizu@gmail.com>
License-Expression: MIT
Project-URL: Homepage, https://github.com/ofplang/run
Project-URL: Repository, https://github.com/ofplang/run
Keywords: ofplang,dataflow,workflow,runner,rolling-horizon
Classifier: Development Status :: 3 - Alpha
Classifier: Intended Audience :: Developers
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Topic :: Scientific/Engineering
Classifier: Typing :: Typed
Requires-Python: >=3.10
Description-Content-Type: text/markdown
License-File: LICENSE
Requires-Dist: PyYAML>=6.0
Requires-Dist: ofplang-schedule<0.6,>=0.5
Requires-Dist: ofplang-validate<0.3,>=0.2
Provides-Extra: test
Requires-Dist: pytest>=7.0; extra == "test"
Provides-Extra: dev
Requires-Dist: pytest>=7.0; extra == "dev"
Requires-Dist: ruff>=0.16; extra == "dev"
Requires-Dist: mypy>=1.11; extra == "dev"
Requires-Dist: types-PyYAML; extra == "dev"
Dynamic: license-file

# ofplang run

[![CI](https://github.com/ofplang/run/actions/workflows/ci.yml/badge.svg)](https://github.com/ofplang/run/actions/workflows/ci.yml)
[![PyPI](https://img.shields.io/pypi/v/ofplang-run.svg)](https://pypi.org/project/ofplang-run/)

A runner for **Object-flow Programming Language v0** — a YAML-based dataflow
workflow IR with linear Object tracking. The language is defined in the
[ofplang/spec](https://github.com/ofplang/spec) repository.

The runner drives an ofplang v0 workflow to completion against an execution
backend, emitting an execution status document (spec §6/§7) as it progresses and
routing typed **view values** through the workflow. It runs on a **simulator** — a
simulated physical backend — so a full workflow can be exercised end to end
without real hardware; the same dispatch contract targets real hardware later.

> **Status:** the simulator, the runner, and a typed (dummy) value layer are
> implemented.
>
> - **Simulator** (`ofplang.run.simulator`) — a physical backend: devices, spots,
>   transporters, and timed operations advanced on a clock. It validates every
>   dispatch (an inconsistent plan is rejected), models timed up/down for a device
>   or a transporter and injected operation failure, and at completion produces each
>   operation's output
>   view values via a **device model** (with none injected, the built-in
>   `script_device_model` runs a `script` process — see below — and otherwise falls
>   back to `default_device_model`, which fills type defaults and carries Object
>   outputs through from their `objects.map`
>   inputs; a custom / real model computes them). `Simulator` is an abstract base
>   over the runner's `Backend` contract; the concrete `VirtualTimeSimulator`
>   advances time instantly (deterministic; the default) and `RealTimeSimulator`
>   paces it to a wall clock (a hardware-free stand-in for a real backend).
> - **Runner** (`ofplang.run.runner`) — two ways to drive a backend:
>   - **`replay`** runs a given execution plan (spec §6) on the backend verbatim.
>   - **`run`** is a rolling-horizon loop: it calls
>     [`ofplang.schedule`](https://github.com/ofplang/schedule) each tick,
>     dispatches the work that can start now, advances the clock, and polls —
>     replanning from the committed history as it goes. It re-routes around a
>     downed machine — a device (its process modes and, by default, its spots'
>     transports and the refills of it), a transporter (the transports it carries)
>     or a replenisher (the refills it performs) — polls at a fixed
>     interval with completion-time estimation, absorbs duration variance (an
>     operation running longer or shorter than planned), and stops the whole run if
>     any activity fails (marking the abandoned work cancelled). Machine up/down,
>     operation failure, and duration variance are scenario concerns injected from
>     Python (not CLI flags).
> - **Value layer** — the runner resolves each port's type and view schema (§7),
>   routes typed view values along the workflow's arcs (producer output → consumer
>   input, across nested composites), contract-checks them, and assembles the
>   whole-workflow outputs. A caller supplies the whole-workflow I/O as a single
>   **run boundary** (`--boundary`): one document with a per-port `{spot, view}`
>   descriptor — `spot` places a boundary Object (§6.8), `view` supplies an input
>   value — for the workflow's entry inputs and final outputs. Unsupplied entry
>   views default. A workflow-embedded static literal (`bind: {port: {value: …}}`,
>   §11) is seeded as that consumer input's value in place of a default. At run end
>   the produced output views are echoed back into a result boundary of the same
>   schema (`--boundary-out`). Non-script values are typed but still dummy — a real
>   device backend plugs into the same seam later.
> - **Device-local consumables** (`ofplang-schedule` §4.7) — a device may declare a
>   stock it holds and a process mode what it draws per run. What each stock holds
>   *at the start of the run* is a property of the run, so it goes in the run
>   boundary (`boundary.inventories.levels`), beside the input views; the level at
>   any later moment is never stated, it is replayed by the scheduler from those
>   levels and the `consumption` each completed activity carries. `--ignore-resources`
>   switches the model off (§4.7.3) for a lab that declares stocks nobody is tracking.
>   The starting levels are *not* echoed into the result boundary — that document is
>   written to be fed back, and a second run replays no history, so echoing them
>   would hand it stock the first run already spent; feed back the status instead.
> - **Replenishment** (`ofplang-schedule` §4.7.1) — where the environment says a
>   replenisher can reach a device, a stock that would run out is **topped up rather
>   than ending the run**. The scheduler places the refill; the runner dispatches it
>   like any other activity, and it holds *both* machines while it works — the device
>   being filled and the replenisher filling it — so it can never overlap the work it
>   feeds. It moves no material and reports no level: a level is derived from what the
>   run started with plus its history, never observed. A refill that fails stops the
>   run like any activity failure. `Backend` gained `dispatch_replenishment` for this,
>   which is why 0.3.0 is a breaking release for a custom backend.
> - **Python script processes** (spec §22, `python_script_processes`) — an atomic
>   Pure-Data process may carry a `script: {language: python, code: …}` section.
>   The built-in device model runs it: the input port values are bound as locals,
>   the code returns a mapping of the declared outputs, and those become the
>   operation's computed output values — the first genuine (non-dummy) computation
>   in the value layer. A script that raises, returns the wrong output names, or
>   returns a non-conformant value fails runtime verification (§22.2) and stops the
>   run gracefully like any other activity failure (the failed activity is marked
>   `failed`, its unstarted successors `cancelled`, exit 1). The script runs inline
>   in real time but advances no simulation time; its outputs appear when the
>   operation completes at its `end`, and the environment mode's `duration` is the
>   scheduler's estimate of the compute cost. Scripts run with full Python builtins
>   and no import restriction (§22.3 permits an implementation to restrict these;
>   this one does not).
> - **Contracts** (spec §9) — an atomic process may declare `contracts` with
>   `requires` (preconditions over `inputs.*.view`) and `ensures` (postconditions
>   over `inputs.*.view` and `outputs.*.view`). The runner evaluates them at
>   runtime against the actual view values: `requires` before the operation runs
>   (a violation stops it before dispatch), `ensures` after it completes. An atomic
>   `requires` that references only run/graph-phase inputs (§5.6) is knowable at run
>   start, so it is checked there as a *preflight* — before any work is dispatched —
>   rather than waiting for that (possibly late) operation. A
>   violation — or a runtime evaluation error — stops the run gracefully like any
>   activity failure (`failed` + downstream `cancelled`, exit 1). Static /
>   graph-time contract checking (a fully-constant contract that is statically
>   false, type errors, bad references) is `ofplang-validate`'s job and is not
>   repeated here. Contracts are checked on atomic processes and on the top-level
>   entry composite (`main`): the entry composite's `requires` is evaluated at run
>   start over the whole-workflow inputs (a violation stops the run before any work
>   runs) and its `ensures` at run end over the whole-workflow inputs and outputs (a
>   violation marks the completed run failed). Nested composite contracts are checked
>   too, at each composite invocation's value boundary: its `requires` once its inputs
>   are available (for an input fed by an upstream process this is mid-run, before the
>   composite's body runs) and its `ensures` once its outputs are, a violation
>   stopping the run gracefully at the composite boundary (its not-yet-run body
>   cancelled). Because non-script process outputs are still typed defaults, `ensures`
>   bites mainly on script processes and boundary-supplied inputs until a real device
>   backend computes physical outputs.
> - **Failure observability** — when a run stops (an injected activity failure, a
>   script error, or a contract violation), the reason is exposed as a structured
>   `RollingRunner.failure` (a machine-readable `kind` code, e.g. `contract_requires`
>   / `script_error`, plus a human-readable `detail`, the `subject`, and the time),
>   and the CLI prints it to stderr. The status document itself stays a valid §6
>   document (the reason is out of band). An optional `contract_observer` callback is
>   invoked for every contract check (held or violated) — a trace hook for debugging.
> - **Static view values** (spec §7.4) — a type whose view field declares a `value:`
>   fixes that field to a constant for every value of the type. The runner projects it
>   onto every view record it routes, so a Python script reading the field and a
>   contract referencing it both see the static value (not a runtime default). A value
>   that carries a conflicting one is forced to the static value.

## Install

```sh
pip install ofplang-run
```

Requires Python 3.10+. Runtime dependencies (pulled in automatically) are PyYAML,
the sibling [`ofplang-schedule`](https://pypi.org/project/ofplang-schedule/) that
`run` (rolling-horizon) calls each tick, and
[`ofplang-validate`](https://pypi.org/project/ofplang-validate/) used by the shared
front door (`ofplang.run.front_door_check`), which validates a workflow file *or* an
already-loaded workflow document, so an embedding caller holding one in memory is
checked the same way a CLI is. The runner *library* never imports validate (and the
replans never re-validate), so it stays a one-shot front door.

For development, install editable with the test extra from a clone:

```sh
pip install -e ".[test]"
```

## Command line

```sh
ofp-run run <workflow> --env <env>
    [--boundary DOC] [--boundary-out FILE] [--observation-out FILE]
    [--poll-interval D] [--margin M] [--seed N] [--no-validate] [-o OUT]
ofp-run run --jobs <run doc> --env <env> [--on-job-failure continue|stop] [...]
ofp-run replay <plan> --env <env> [-o OUT]
```

`run` drives a v0 workflow to completion by replanning as it goes: each tick it
polls the backend and, when anything the scheduler reads has changed -- an operation
finished, a machine went down, a pending activity came due -- renders the committed
history as a status, calls the scheduler in-process -- handing it the workflow, the
environment and that status as documents, so nothing goes through a temporary file --
and dispatches the newly-runnable work. A
tick that changed none of those keeps the plan it already has, so a long protocol
costs one solve per activity event rather than one per unit of its makespan; what is
observed, and so the status produced, is the same either way. `--boundary` supplies the whole-workflow I/O as one document —
a `boundary:` mapping with a `{spot, view}` descriptor per entry input / final
output port. `spot` places a boundary Object on an environment spot (spec §6.8;
Object ports only); `view` supplies an input's view value (unsupplied entry views
default). The runner projects it into the scheduler's interface (spots only, so the
scheduler stays value-independent) and the seeded input values. `--boundary-out`
writes the result boundary — the same schema with each produced output's `view`
filled in — a run-local artifact, separate from the value-free status document. On
completion each pinned Object output is checked to have reached its declared spot.
`--observation-out` streams the **observation document** (see `docs/OBSERVATION.md`):
a YAML multi-document stream recording each *completed* activity's concrete input /
output view values (a transport's moved view), appended as each activity finishes —
the value-layer companion to the status document, also run-local.
`--poll-interval` sets the fixed polling interval (default 1). `--max-ticks` is the
non-termination guard: a run that takes more than that many ticks is given up on. One tick
is one poll interval, so the guard also caps the makespan a run can reach (the default
100000 with the default interval means 100000 time units); pass `0` for no limit when a
long virtual run needs it, which also gives up the protection against a backend whose clock
does not advance. `--margin` sets the
running-task margin: on each replan a still-running activity is pinned to end at
`max(reported end, now + margin)`, so a positive margin is what keeps an
overrunning operation's successor from being planned at `now` and dispatched onto
a value that operation has not produced yet. The default is 0, which is only safe
with a backend whose operations cannot finish later than planned (the in-process
virtual-time simulator); against a wall-clock or real backend set it to at least
the poll interval, or a successor is refused with `input_not_produced` rather than
computing on a typed default. `--on-job-failure`
decides what one job's failure does to the rest of a `--jobs` run: `continue` (the
default) stops that job alone and lets the others finish — which is why they were
planned together — while `stop` stops the whole run. A stopped job's remaining work is
reported `cancelled`, and the spots its material is still sitting on are declared in
the status (`occupied`, §6.12) so the rest of the run is planned around them rather
than onto them. A single workflow is a single job, so this makes no difference to it.
`--jobs` runs
**several workflows together** in one laboratory (schedule SPEC §6.11) in place of
the single `<workflow>` argument. Its run document names each job — an `id`, the
workflow it runs, its own `boundary`, and the `release` time before which it may not
start — plus the two things that belong to the laboratory rather than to any one job:
what its stocks hold at the start of the run (`inventories`, §6.10) and which spots it
is already holding (`occupied`, §6.12). The jobs are planned *together*, so they
compete for the same machines and draw on the same stocks: a refill neither job needs
alone can appear because the pair of them does. Each job's activities carry its `id`
in the status, and the plan's roster reports the completion the scheduler promised
each one. See `examples/shared_refill.run.yaml`. `--no-validate`
skips the one-shot `ofplang-validate` front-door check of the workflow — use it
when the workflow was already validated upstream (e.g. by the `ofp` umbrella CLI);
`$import` is still resolved and the capability gate still runs, since both are
structural rather than validation. `replay` runs a plan
produced by `ofp-schedule` verbatim on the simulator (no value layer). Both write
the final execution status as YAML (`-o`, else stdout). Exit codes: `0` success,
`1` execution failed (an activity failed, or a replan is infeasible), `2`
usage/input error.

This tool is also the `run` subcommand of the umbrella `ofp` CLI
([`ofplang`](https://pypi.org/project/ofplang/)), which forwards to it in-process
with this CLI's own subcommands intact: `ofp run run …`, `ofp run replay …`, each
with the same options and exit codes as above.

The package lives under the `ofplang` PEP 420 namespace (`ofplang.run`), shared
across the organization's tools.

## Feature support

v0 defines seven optional features (spec §4.2), and a document requiring one an
implementation does not have "is valid v0 but unsupported by that implementation"
(§4.1). So `ofp-validate` accepting a workflow does not mean this runner can
execute it:

| v0 feature | `ofplang-run` |
|---|---|
| `python_script_processes` | **Supported** — the built-in device model runs the script and verifies its outputs (see above). |
| `scheduling_policies` | Ignored, as in [`ofplang-schedule`](https://github.com/ofplang/schedule), which does the planning. |
| `generic_processes` | **Not supported.** The front door's capability gate refuses it before anything runs, naming the process. |
| `node_map`, `node_fold`, `node_do_while`, `node_branch` | **Not supported.** The front door's capability gate refuses a structured node before anything runs, naming the node and the feature — a structured node reshapes dataflow (lifting an output to an `Array`, threading a value across iterations, leaving an arm unrun) in ways neither this runner nor the scheduler it plans through represents. |

## Examples

[`examples/`](examples/README.md) holds runnable scenarios — supplied inputs and
computed outputs, views routed across a composite boundary, a script process with
contracts checked at runtime, re-routing around a device that goes down, the drift
fixed-interval polling costs, two jobs run together needing a refill neither needs
alone, and one job of three failing while the others finish. Most are Python scripts rather than CLI invocations, because what they
demonstrate is injected from code: a device model, a machine fault, a polling
interval. Their output is committed under
`examples/outputs/`, so an example can be read without being run.

## Tests

```sh
pytest
```
