Execution Backends#
oxo-flow computes a fully determined static plan — topological order, parallel groups, per-rule resource declarations — before any execution begins. Execution backends do not re-understand the DAG: they map that single plan onto a scheduler API. Because plan computation, invalidation, and the checkpoint format stay single-implementation, the dry-run preview is valid for every executor, and local and remote execution share the same checkpoint semantics.
.oxoflow file ──► wildcard expansion ──► static plan ──┬─► LocalExecutor (run)
(single implementation) └─► ExecutorBackend (scheduler)
The static plan#
oxo_flow_core::backend::ScheduledPlan is the executor-agnostic plan:
order— topological execution order (WorkflowDag::execution_order),groups— parallel groups (WorkflowDag::parallel_groups),rules— oneScheduledRuleper schedulable instance: the expanded rule (single source for script rendering and resource directives), the environment-wrapped shell command withcd <workdir>folded in, the workdir, instance-level dependencies, and{config.x}bindings.
ExecutorBackend trait#
#[async_trait::async_trait]
pub trait ExecutorBackend: Send + Sync {
fn name(&self) -> &'static str;
fn render_script(&self, rule: &ScheduledRule) -> Result<String>;
async fn submit(&self, script_path: &Path) -> Result<String>;
async fn poll(&self, job_ids: &[String]) -> Result<HashMap<String, BackendJobStatus>>;
async fn cancel(&self, job_id: &str) -> Result<()>;
async fn logs(&self, job_id: &str) -> Result<String>;
// Provided methods (defaults shown):
fn render_array_script(&self, rule: &Rule, cmd_dir: &str, count: usize)
-> Result<String> // default: Err "does not support job arrays"
— render one ARRAY script for `count` same-rule instances whose
per-index commands live in `cmd_dir`; backends without array
support error and the driver falls back to per-job submission
async fn terminal_status(&self, job_id: &str) -> Option<TerminalRecord>
// default None — terminal record from the scheduler's accounting
// store; None keeps the job in flight
fn polls_elements_directly(&self) -> bool
// default false — SLURM's squeue reports array elements by their
// `{base}_{index}` ids (true); PBS/SGE/LSF list only the base id
}
The first implementation, ClusterExecutor (SLURM/PBS/SGE/LSF), renders
scripts with the existing directive generator (core/cluster.rs, unchanged)
and performs real submission and tracking:
submit—sbatch --parsable <script>for SLURM (bare job id on stdout, fixing the--dependency=afterok:Submitted batch job Ntrap the old wrapper script had); PBS prints a bare id; SGE/LSF ids are parsed out of their submission sentences. All id capture shares one helper (parse_job_id), as does status-line parsing (parse_status_line), including array elements like12345_12.poll—squeue -j <ids> --noheader -o "%i|%T"and per-backend equivalents.cancel/logs—scancel/sacctand equivalents.
BackendDriver#
oxo_flow_core::backend::driver::BackendDriver executes a ScheduledPlan
through a backend. It never decides what to run — the caller computes the
will-run set from the shared invalidation predicates (the same ones the
dry-run preview uses). The driver:
- submits waves of at most
max_submittedin-flight jobs (pending + running), - polls at
poll_intervaland settles terminal jobs intoJobRecords — the same record type the local executor produces, so checkpoint recording is identical in shape, - propagates failures: dependents of a failed rule are skipped
(
blocked by failed upstream dependency), - cancels everything in flight on submit errors and poll timeouts,
- writes a greppable run directory:
events.jsonl(append-only, source of truth) andjobs/<rule>/withjob.sh,job.id, andstatus.json. Content-addressing stays inside the checkpoint — never in paths, - supports checkpoint re-entry via two hooks (
on_checkpoint,merge) so dynamically discovered instances execute in the same run (see Checkpoint re-entry below).
DriverConfig supplies the knobs, with these defaults:
| Field | Default | Meaning |
|---|---|---|
max_submitted |
50 |
In-flight jobs (pending + running) before submissions top up as slots free |
max_array_size |
1001 |
Scheduler array limit — larger groups are chunked into several arrays |
no_arrays |
false |
Set to submit every instance as its own job |
poll_interval |
5s |
Delay between scheduler polls |
poll_timeout |
None |
Overall wall-clock budget; unset means run until terminal states |
A [cluster] profile block overrides max_submitted, max_array_size, and
poll_interval (see
Cluster submission); no_arrays,
poll_timeout, and the defaults themselves are driver-internal.
Testing without a cluster#
The workspace ships a mock SLURM scheduler at
tests/fixtures/mock-scheduler/ — shell shims (sbatch, squeue, scancel,
sacct) that execute submitted scripts for real in the background and record
job state under $MOCK_SCHEDULER_DIR. Integration tests construct a
ClusterExecutor with .with_scheduler_dir(...) and
.with_env("MOCK_SCHEDULER_DIR", ...) to run whole workflows through the
driver in CI without any cluster. The shims reproduce two real-scheduler
behaviours worth knowing about: sbatch returns immediately (the job's
output pipe is fully detached) and failing jobs still record a terminal
state.
The P1 acceptance tests (tests/cluster_backend.rs) assert the core claim:
- local execution and backend execution of the same workflow produce the same checkpoint semantics (completed sets, benchmarks, input manifests) and identical outputs,
- the dry-run preview's will-run set equals the set the driver submits,
- poll timeouts cancel in-flight jobs.
run --profile <NAME> now drives the backend when the profile carries a
[cluster] block (see Cluster submission).
Its submission set is the dry-run preview's will-run set unioned with the
run's force_rules — the same bypass set the local executor receives, since
run has already applied and persisted its invalidation analysis by the time
the cluster path runs. tests/cluster_run_profile.rs pins that combination to
what the local executor actually executes from identical state.
Real-scheduler validation: SLURM is live-verified end-to-end (including
GPU). SGE (Son of Grid Engine 8.1.9) and OpenPBS 23.06 are verified on
containerized real schedulers — master + 2 exec hosts built from source —
with the chain and array matrices passing 6/6 and 20/20 per backend
(issue #356). The harness, workflows, and runners live in
scripts/real-cluster-matrix/. LSF cannot be stood up the same way (the
community edition is IBMid-gated); its directive and status semantics are
cross-checked against IBM's bsub reference and the jokergoo/bsub client.
Re-attach to in-flight jobs is a follow-up.
Accounting is read once per job as it settles, not once per poll, and feeds
both the run directory and the checkpoint benchmarks — see
what a finished job records.
Backends report only what their store proves: SLURM reads exit code, elapsed
time, peak RSS and total CPU from sacct (peak memory comes from the step
rows, which is where SLURM puts it — the allocation row leaves it blank), PBS
and SGE report the equivalent resources_used / ru_* fields, and LSF
reports state alone.
Assumptions#
- Outputs, workdir, and logs live on shared storage — every compute node sees the same files, so jobs can run anywhere and the driver's checkpoint bookkeeping stays meaningful.
- Jobs are POSIX shell scripts; the scheduler runs them with
/bin/bash.
Checkpoint re-entry#
See Workflow Format for the config
surface. A checkpoint = true rule writes a re-entry manifest at runtime;
when it completes, the engine merges the new values, re-expands the rule
templates, and executes the new instances in the same run — on the local
path (run) and on the driver path alike. Every round is still a static
plan, so previews stay deterministic and resumes reconstruct the same plan.