PolyGenius
Contents

ExecutionScheduler

Execution Scheduler

Internal controller for one reactive execution. This class turns resource specs into a dependency-indexed task graph:

  • discovery registers dependencies once and parks tasks with unresolved inputs;
  • dependency completion wakes only direct waiters through waiting.by.dep;
  • scheduling consumes only ready task ids and never scans all waiting tasks just because resources changed.

The class owns the live SchedulerState record. The engine reads scheduler state only through narrow pull accessors (is.done(), running.entries(), status.counts(), waiting.tasks(), timing.totals(), summary.snapshot(), resolved.spec()) rather than reaching into execution$state, so there is no shared-mutable-state sync step.

Details

Execution scheduler

Requested outputs are displayed as waiting until discovery proves they are directly runnable. This keeps lazy discovery from showing thousands of undiscovered final outputs as if they were the active runnable queue.

Processes the entire discovery queue in one pass. The bulk SQL prefetch (resource.rows.bulk + relations.bulk) is done upfront so the per-task loop avoids individual round-trips for cache lookups. All depends_on relations collected during the loop are written in a single batched transaction at the end, replacing the previous one-write-per-dep pattern. max.ready can throttle discovery when only a bounded ready queue is desired; default Inf discovers everything.

Scheduling never scans waiting tasks. It consumes only the ready queue, fails tasks that can never fit the total budget, and submits the highest-priority task that fits current free resources.

Collection is the event that releases resources and wakes dependency waiters. Successful results are persisted before dependents are rediscovered so downstream cache lookup sees the completed input.

Materialized from the id-map in insertion order. The list is bounded by the core budget and shares the underlying future objects by reference -- it is never a deep copy, so per-tick reads stay cheap.

Materialized as a plain list. Safe on the zero-failed deadlock path (no vapply over the task environment).

The maps are materialized via as.list() (insertion-ordered value copies), so mutating the result cannot reach live scheduler state.

Used by the status renderer, which needs per-group elapsed seconds but not the other terminal maps. Copying only task.timings avoids the full summary.snapshot() (which also materializes failed/blocked/ resolved/dynamic.outputs) once per group per render.

await, log.state, and the timing metrics run every loop tick. This returns the live running entries (shared by reference) plus the scalar core/wall totals and the done flag, without materializing the large id-keyed maps (tasks/resolved/task.timings) the way summary.snapshot() does.

Methods

Public methods

  • ExecutionScheduler$new()
  • ExecutionScheduler$add.outputs()
  • ExecutionScheduler$process.discovery()
  • ExecutionScheduler$schedule.ready()
  • ExecutionScheduler$collect.completed()
  • ExecutionScheduler$has.work()
  • ExecutionScheduler$is.deadlocked()
  • ExecutionScheduler$state()
  • ExecutionScheduler$summary()
  • ExecutionScheduler$is.done()
  • ExecutionScheduler$mark.done()
  • ExecutionScheduler$running.entries()
  • ExecutionScheduler$status.counts()
  • ExecutionScheduler$waiting.tasks()
  • ExecutionScheduler$summary.snapshot()
  • ExecutionScheduler$task.timings.snapshot()
  • ExecutionScheduler$timing.totals()
  • ExecutionScheduler$resolved.spec()
  • ExecutionScheduler$clone()

Method new()

Create a scheduler for one execution.

Usage

ExecutionScheduler$new(execution, store, registry, state, logger, callbacks)

Arguments

execution — Mutable execution environment owned by ExecutionEngine.

store — ResourceStore used for cache lookup and persistence.

registry — RuleRegistry used to match output specs to rules.

state — SchedulerState instance for this execution.

logger — Execution logger.

callbacks — Named callback list supplied by ExecutionEngine.

Method add.outputs()

Register requested root outputs.

Usage

ExecutionScheduler$add.outputs(outputs)

Arguments

outputs — List of ResourceSpec objects.

Returns

Named summary list.

Method process.discovery()

Discover rules and dependencies for pending tasks.

Usage

ExecutionScheduler$process.discovery(max.ready = Inf)

Arguments

max.ready — Stop discovery once the ready queue reaches this size.

Returns

Named summary list with discovery counts.

Method schedule.ready()

Submit ready tasks while resources are available.

Usage

ExecutionScheduler$schedule.ready(submit.fn)

Arguments

submit.fn — Function accepting (task, req.cores, req.mem) and returning a running backend entry.

Returns

Named summary list with scheduling counts.

Method collect.completed()

Collect finished backend tasks and release dependents.

Usage

ExecutionScheduler$collect.completed(
  ready.fn,
  collect.fn,
  input.log.paths.fn,
  cleanup.run.log.fn,
  dispatch.fn = NULL
)

Arguments

ready.fn — Function returning completed running-entry indices.

collect.fn — Function collecting one running-entry result.

input.log.paths.fn — Function returning upstream log paths for a task input list.

cleanup.run.log.fn — Function removing temporary run logs.

dispatch.fn — Optional zero-argument function called once phase 1 has returned the collected tasks' cores to the budget, so freed slots refill before the write transaction instead of after discovery.

Returns

Named summary list with completed/failed counts.

Method has.work()

Test whether scheduler work remains.

Usage

ExecutionScheduler$has.work()

Returns

Logical scalar.

Method is.deadlocked()

Test for unresolved waiting work with no way to progress.

Usage

ExecutionScheduler$is.deadlocked()

Returns

Logical scalar.

Method state()

Return the mutable scheduler state.

Usage

ExecutionScheduler$state()

Returns

SchedulerState.

Method summary()

Return a compact scheduler summary.

Usage

ExecutionScheduler$summary()

Returns

Named list of status and queue counts.

Method is.done()

Whether the execution has been marked done.

Usage

ExecutionScheduler$is.done()

Returns

Logical scalar.

Method mark.done()

Mark the execution done. The scheduler owns the terminal flag; the engine calls this instead of writing state directly.

Usage

ExecutionScheduler$mark.done()

Returns

Invisibly returns NULL.

Method running.entries()

Live running backend entries (<= core budget; hold futures).

Usage

ExecutionScheduler$running.entries()

Returns

Named list of running entries.

Method status.counts()

Task status counts across the graph.

Usage

ExecutionScheduler$status.counts()

Returns

Named integer vector.

Method waiting.tasks()

Tasks currently parked in the waiting state.

Usage

ExecutionScheduler$waiting.tasks()

Returns

Named list of SchedulerTask records.

Method summary.snapshot()

Frozen, value-copied snapshot of the terminal/result maps and timing counters used by engine summaries and finalization.

Usage

ExecutionScheduler$summary.snapshot()

Returns

Named list (mutating it cannot reach live scheduler state).

Method task.timings.snapshot()

Value copy of just the per-task timing map.

Usage

ExecutionScheduler$task.timings.snapshot()

Returns

Named list of per-task timing records.

Method timing.totals()

Lean timing/running snapshot for the hot per-tick path.

Usage

ExecutionScheduler$timing.totals()

Returns

Named list with running, core.seconds.completed, wall.seconds.completed, and done.

Method resolved.spec()

Resolved canonical spec for one requested output id.

Usage

ExecutionScheduler$resolved.spec(id)

Arguments

id — Character requested-output id.

Returns

ResourceSpec or NULL when the id was not resolved.

Method clone()

The objects of this class are cloneable with this method.

Usage

ExecutionScheduler$clone(deep = FALSE)

Arguments

deep — Whether to make a deep clone.

See Also