Skip to content

Inject the stage runner into package objects #5

Description

@SimonHeybrock

Summary

Package objects that drive stages, such as the SansReduction of scipp/sciline#245 (design doc section 6.1), could take the way a stage call runs as an injected dependency. The package author writes the driver once: which stages, how runs and banks are looped, which contributions are kept. essapps injects a runner that decides where each stage call executes: in the calling process, held in a session, or as a task on a cluster node.

This issue records the idea and its open questions. Nothing is decided and nothing is built.

Why

The essapps binding (PipelineAdapter(aggregations=...)) uses only the contribute half of a package's sciline.Aggregation. It builds its own final stage, handles the tuning inputs, and keeps its own incremental cache. The package object of sciline section 6.1 solves the same problems differently: it keeps contributions per file and drops runs no longer listed, while the adapter starts over on any change that does not extend the previous list. The same problem is solved twice, and a package object as esssans would ship it cannot be wrapped through the current seam.

The design principle is that the essapps interface follows what sciline supports (Aggregation, Stage, drivers that combine stages). Injecting the runner into the package object applies that principle from the other side: essapps uses the package's driver instead of being one.

Sketch

# sciline: the elementary block, as an interface
class StageRunner(Protocol):
    def stage(self, pipeline, inputs, outputs) -> StageLike: ...

class StageLike(Protocol):
    keys: ...                                       # what the stage reads, for identity
    def submit(self, inputs) -> Future[dict]: ...   # compute() is submit().result()

# esssans: the driver, written once against the interface
class SansReduction:
    def __init__(self, pipeline, runner: StageRunner = LocalRunner()): ...
    # set_runs and compute as in section 6.1, with every stage made by runner.stage(...)

# essapps: other runners
SessionRunner(cache)        # holds stages in memory, for tuning and incremental sums
ClusterRunner(backend)      # each submit() becomes a task on a node, outputs to the data store

What it would give

  • One driver. Runs times banks (section 6.3), per-member contributions, and handling a changed run list live in the package. The essapps adapter becomes "build the package object with our runner, then call it".
  • Splits come from the driver. The stage boundaries the driver already has are where work can move to another node. The author does not write a separate pair of specs such as CONTRIBUTE and COMBINE (docs/developer/aggregation.md, "Reducing the runs of a sum on separate nodes").
  • Incremental series by cache hits. If a remote runner stores each stage result under a content key, a series under a rule re-runs the whole request on each arrival, and only the new run's contribution computes. The key would be built from the stage (its keys and the parameter values they hold), the hashes of its inputs, and the code version. For sums, this removes the need for rules that fire on completed records. Two-phase fan-out (a first call finds the keys, then one request per key) still needs such rules.
  • Records stay "what". The top-level request is still spec plus parameter values. Stage calls underneath are cache entries, not records. Their keys come from Stage.keys at run time, computed by code, so the spec declares nothing about the graph and no reads table is needed.
  • Tuning is a session runner that holds stages. The driver still needs extra inputs on its final stage (the finalize_inputs= item in docs/developer/open-issues.md), so that a tuned parameter does not force a new package object. This is needed with or without the runner.

Costs and open questions

  • Futures in the interface. The driver of sciline section 6.3 is a sequential loop. To run contributions in parallel, submit returns a future and the driver pushes results as they complete. This makes sciline's driver code less plain.
  • A coordinator process per request. The package's driver runs somewhere and waits for remote calls. In batch and automatic reduction, this is a small process per request. The framework then runs package code as orchestration, not only as leaf work.
  • Serialization at each remote boundary. Contributions need a storable form. The explicit spec split has the same cost, but there it is visible in the spec.
  • Is a cache entry good enough for publication? Publication refuses results served by a held stage. A durable, content-keyed cache that includes the code version may be acceptable. This needs a decision.
  • Driver determinism. Reuse by content key requires deterministic driver code. This becomes a requirement on package authors. AiiDA's workchains and Dask's delayed have the same shape. docs/developer/prior-art/aiida.md records the provenance split between workchain and calculation nodes there.
  • Relation to current decisions. If adopted, the author-written CONTRIBUTE/COMBINE split becomes an optimisation the runner provides, and rules over completed records serve only two-phase fan-out.

Possible spike

  1. In sciline, a StageRunner protocol, with Aggregation and the validation script's SansReduction taking runner=. The default runner is local, so no existing behaviour changes.
  2. In essapps, a SessionRunner and a process-pool runner with a content-keyed store, driving the NORMALIZE example and a series under a rule.
  3. Check that the k-th arrival of a series computes one contribution, that results equal the one-request list form, and how much driver code changes.

🤖 Generated with Claude Code

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    questionFurther information is requested

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions