ADR 0003: Stages outside the graph - #245
Draft
SimonHeybrock wants to merge 55 commits into
Draft
SimonHeybrock wants to merge 55 commits into
SimonHeybrock wants to merge 55 commits into
Conversation
A design proposal with a prototype. Sciline's map/reduce put a loop over members and a combine step inside the graph; that is the source of every breaking change in the PEP 695 work and cannot express what ess.reduce.streaming.StreamProcessor and the essapps architecture need. The proposal makes the partition of a flat pipeline at named keys the primitive (Stage: static part held, inputs per call) and builds the replacement for map/reduce on it (Fold: member table, cut keys with combine functions, contribute/combine/finalize), with StreamProcessor, the essapps warm workflow, and split workflows as the same objects. The prototype runs against main and reproduces the LoKI multi-run reduction identically against the existing with_sample_runs. Related: #235, #241, #233.
…embers The first draft kept as_pipeline() so that esssans's with_* helpers could keep returning a pipeline. That bridge was the only way two folds could compose, fixed the member table at construction, and hid a loop inside a provider. It is dropped; a drop-in pipeline is a non-goal. Stage is unchanged in kind: it passes through an output that is an input, exposes the keys it reads, and gets warm(), which computes the static parts of several stages of one pipeline in one run so that shared static work, such as a mask file every group reads, is done once. Between stages sit connectors with push, value, and clear: reducers where members are combined, a Forwarder where a context is held. Drivers are loops over stages and connectors: Fold for the table-fold shape, with groups of member keys, settable members and parameters, contributions held per member so that adding a run costs one contribution, and the three entry points contribute, combine, finalize; StreamProcessor for streams. A generic network object and a boundary builder were considered and deferred, with the reasons recorded. The LoKI validation uses one Fold with two groups and the pixel masks as a list parameter, since their cut sits inside the per-run work. Results are identical to the reference and provider call counts are equal. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Fold had grown three jobs: orchestration, a cache-invalidation engine (a __setitem__ deciding per stage what to drop, set_members keeping held values by a row-equality heuristic, a held member frontier), and a member table with groups. The last two reproduced the map/reduced pipeline of esssans, and they were the part that was hard to reason about. Simon asked for parameters on the pipeline and Fold as stage orchestration only. Fold is now an immutable snapshot: contribute stage, finalize stage (optional), combine functions, the three entry points, and compute over a table that holds nothing. compute_members is a function. Groups are gone; several folds sharing a finalize are composed by the caller, and the "set runs, set parameters, compute" experience lives in a package object. The LoKI validation shows that object, SansReduction, in about thirty lines, with results identical to the reference and equal call counts. The doc records why the earlier Fold was wrong and where a network object would come from if the package objects repeat. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Condenses the design document and its three prototype passes into a proposed decision: Stage and warm in sciline, connectors and a stateless Fold in ess.reduce, package objects owning topology and state, no drop-in replacement for the map/reduced pipeline. Lists the alternatives tried and dropped along the way. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Where the branch, prototype, environment, and companion notes are; the design in ten lines; how three passes led here and why the first two were rejected; open questions with recommendations; a staging to evaluate that lands Stage and warm additively before removing map/reduce; per-package migration sites; things to verify in review. Working document, to be dropped before the branch merges. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Fold and compute_members move from ess.reduce to sciline: they need nothing beyond Stage, hold no policy, and are the documented replacement for map(...).reduce(...) that users outside ESS need. Connectors stay in ess.reduce as a push/value/clear convention without a base class, since nothing in sciline consumes one. ADR corrections from the review against the design doc, prototype, and source: the PEP 695 PRs are closed prototypes, not landed work; groupby never existed; constraints= is removed by the generics, not by this change; Fold holds its stages' frontiers; only esssans, essreflectometry and bifrost fold runs in library code; BatchProcessor lives in essreflectometry; the LoKI timing carries its dask caveat. Terms (frontier, cut, member, contribution, driver, package object) are defined at first use. Added the coexistence alternative and the missing negative consequences: silently static cut keys, held frontiers and contributions, graph walks per parameter change, lost mapped-graph visualization, progress reporting, networkx as a direct dependency. Open questions are listed under Decision. Design doc: test count is 238, not 335; two prototype tests fail on the PEP 695 branch, both inside map; DREAM's mask fold is static, not nested; esssans folds runs in four map/reduce pairs; bifrost has three folds; migration recommends shipping the additive part in a minor release so packages migrate one at a time.
…tion Decisions from the review discussion, applied to the ADR, the design doc, the handoff note, and the prototype: - The combine per accumulation key is an accumulator, not an n-ary function. Sciline defines a structural Accumulator protocol (push, value) and Buffered(func), a factory wrapping an n-ary function. Aggregation takes a factory per key, makes fresh accumulators per combine, and compute(table) pushes each contribution as it is made, so peak memory is the accumulator's choice. The ess.reduce accumulators satisfy the protocol once maybe_hist leaves their base class; clear stays an ess.reduce convention. - Parallelism over members is the caller's job; contribute is a plain call to map over rows. The guide will show an example. - Names: Fold becomes Aggregation, cut key becomes accumulation key, at= becomes accumulators=, fold.cut becomes accumulation_keys. fold clashes with scipp's reshaping fold; the contribute/combine/finalize shape is what Spark, Flink, Beam, and pandas call aggregation, and there the stateful object is the accumulator. - No experimental label; the additive minor release with esssans migrated is the trial period. - Member tables are Mapping[label, Mapping[Key, value]] only. Prototype: drops the redundant single-sink trick in _compute (both schedulers keep requested keys), rejects an empty member table, and adds tests for Buffered, custom accumulator classes, streaming push order, and snapshot semantics under grafting. LoKI validation is identical to the map/reduce reference with equal call counts.
The prototype from docs/developer/architecture-and-design/stage-prototype becomes src/sciline/stage.py (Stage, warm) and src/sciline/aggregation.py (Accumulator, Buffered, Aggregation, compute_members), exported from the package and listed in the API reference. This is the additive part of ADR 0003; map/reduce are untouched. Changes from the prototype: the concrete graph comes from to_task_graph instead of TaskGraph._graph; the scheduler default is shared with TaskGraph through scheduler_or_default, so stages use dask when it is installed; outputs not in the pipeline are rejected at build time; NumPy-style docstrings; mypy strict clean. Tests move to tests/stage_test.py and tests/aggregation_test.py, with additions for scheduler handling, snapshot semantics, and warm skipping warm stages. The LoKI validation script moves next to the design doc, imports from sciline, and runs both sides on the naive scheduler so the timing compares call structure only; results remain identical. The design doc gains a rollout plan: PR-level order across sciline, ess.reduce, the reduction packages, esslivedata, and essapps, with the file-level sites from the survey and the scheduling dependencies.
keys is what a caller tests to decide whether a parameter change touches a stage. With an intermediate result as input, the input's ancestors are cut off and a change to them does not reach the stage, so they must not be in keys. keys is now the union of the held part and the per-call part.
A user-guide page next to the parameter-tables page, using one example that grows through the sections: a stage and what it holds, warming several stages, an aggregation with Buffered, a running-sum accumulator, the three steps called separately and a member dropped by recombining held contributions, members contributed in parallel from a thread pool, compute_members, an accumulation key that turns out static, and a table with two member keys. The page does not claim to replace parameter tables; ADR 0003 is not accepted yet.
The ADR, the design document, and its rollout plan carry everything the note held.
This was referenced Sep 11, 2026
A type: ignore that only a newer mypy needs, and an Any escaping from the validation script where ess.sans is not installed.
Clearing inside the aggregation would either forbid combining in batches over several calls or silently add to stale state, and holding instances would make combine non-reentrant. With factories, whoever asks for accumulators owns their lifetime.
The dask.distributed workers cannot import the test module, so keys and providers defined there fail to unpickle and the workers die; on CI the scheduler kept retrying and the dask job never finished. Builtin keys and local providers pickle by value, as in the existing distributed test.
Buffered holds every pushed value, which is the wrong default for sums of large arrays. Running-total accumulators were hand-written in the tests and the user guide; Reduced(func) replaces them. It never updates in place, so it cannot mutate values owned by the caller, and requires an associative function so that combining in groups gives the same result. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
…ge.static Aggregation.combine takes combined values as well as contributions, and the chained test relies on it, but neither the Accumulator protocol nor Buffered said what that requires: a pushable value and a result that does not depend on grouping. The docstrings and the design document now say so. Stage computed its held part on first use without a guard, while the design document recommends mapping contribute over rows with threads; two cold calls would compute the frontier twice. A lock makes it once. Design document: the essapps paragraph described the chained series as held accumulators fed per arrival, which is the fold; the chained series is combine([previous, new]) per throwaway process. The D13 declaration of finalize parameters stays in the essapps spec, checked against the graph, since the backend validates without importing workflow code. The evidence paragraph still described a prototype reading TaskGraph._graph with a hardcoded scheduler. Rollout item F updated to match stages.md. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
SimonHeybrock
added a commit
to scipp/essapps
that referenced
this pull request
Sep 14, 2026
The sciline proposal settled its open choices: Fold is Aggregation, the combine protocol is an accumulator factory per accumulation key with Buffered for n-ary functions, Aggregation lives in sciline, the member table is a plain mapping, and parallelism over members is the caller's. The note follows those, revises one judgment (keep the D13 declaration of finalize parameters and check it against the graph, for the reason D8 keeps cheap parameters declared), adds the closure condition that chaining puts on accumulators, and lists what should go back to the sciline proposal. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
SimonHeybrock
added a commit
to scipp/essapps
that referenced
this pull request
Sep 14, 2026
Sciline's proposal (scipp/sciline#245) calls the key at which per-member values are combined an accumulation key; the sketch used accumulation point for the same thing. One word across sciline, ess.reduce, and this framework. The stages note now records "fold" as kept for the process shape, and the review log gets the eleventh pass. Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
A stage holds only the values at its frontier, not the whole static part. warm uses the scheduler of the first stage that is not yet warm. Aggregation also raises if an accumulation key or an output is unknown. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Context now shows what map/reduce does and states the reasons for removal without the mechanics of the failed prototypes. Decision introduces Stage and Aggregation with examples and a diagram, defines terms before use, and separates what sciline provides from what callers own. Semantic detail is left to the docstrings and the design document. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
The document had grown through several rounds of design and prototyping. It is now organised as survey, overview with a glossary, the three building blocks, composition patterns for each current use, design choices, validation, migration, and open questions. Motivation is left to ADR 0003. Details of the generics prototype and superseded naming history are dropped; the rollout plan is kept with its file references. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Phase 3 in essapps is interactive applications, of which splitting a workflow is one of three models; say so, and explain spec, record, and binding for readers outside essapps. Replace the per-item rollout plan and its file and line references with what changes per project and the order of releases. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
SimonHeybrock
force-pushed
the
map-reduce-outside-the-graph
branch
from
September 14, 2026 14:47
a1e56ed to
26d2f33
Compare
jl-wynen
reviewed
Sep 14, 2026
| `Reduced(func)` is a factory for an accumulator that holds only a running result of an associative binary function. | ||
| The accumulators of `ess.reduce.streaming` satisfy the protocol. | ||
| `Forwarder`, which holds the value of the latest push, is the connector for a context between stages and stays in ess.reduce. | ||
| - `Aggregation`, in sciline: the table-fold shape as two stages of one pipeline, contribute and finalize, with one accumulator per accumulation key between them and with the three entry points exposed. |
Member
There was a problem hiding this comment.
What is a 'table-fold shape'? And which three entry points?
The ADR again says that parameter values are held by reference, that the StreamProcessor rewrite is not done, and why pinning is offered instead of a separate namespace. The design document again states the obligation on histogramming accumulators, the requirements Stage places on the new Pipeline, the driver's ownership of order dependence, why Accumulator and Forwarder live where they do, and what StreamProcessor.visualize needs. The survey findings that matter for the migration and the dependencies in the release order are back in short form. The Stage docstring states the by-reference behaviour too. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Stage is mostly introspection (inputs, outputs, keys, frontier, dynamic), so a call on the object hid the one action among its members, and call sites such as `agg.contribute_stage(row)` read like a getter. `compute` matches `Aggregation.compute` and `Pipeline.compute`. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Reading the static property computed the held part on first access, which runs providers and may read files; a method makes that visible. It delegates to warm, which now takes each stage's lock, so the held part is computed once also when warm and compute run in different threads. Locks are taken in a fixed order to avoid deadlocks between overlapping warm calls. Buffered's docstring now says that the function is applied on every read of value. Also wrap test lines that exceeded the line length after the rename to compute. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
ESSlivedata selects its scheduler by replacing sciline.task_graph.DaskScheduler and reads that name at import. Moving the default into sciline.scheduler removed the name, so ESSlivedata would fail to start. scheduler_or_default now lives in task_graph and looks the name up on each call, so the replacement applies to pipelines and stages alike. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Prototypes of LoKI banks and Bifrost triplets over several runs, and of the StreamProcessor shape, showed that no real driver used Aggregation's finalize stage or compute(table), and that compute(table) invited the flat runs-times-banks table, which counts a run-level key once per bank without an error. split(pipeline, *parts) takes one Part per loop level, each with its inputs, the outputs the driver pushes, and its parent, and returns one Stage per part. A value a part reads that does not depend on its inputs is computed by the deepest ancestor whose inputs it depends on, once per iteration of that loop. An output that does not vary in its part, and a read of a value that depends on a non-ancestor's inputs, raise. This covers runs times banks, Bifrost's per-run step after combining triplets, and StreamProcessor's context and finalize boundaries, with results equal to the references and to the real StreamProcessor. Accumulator, Buffered, and Reduced stay, in sciline.accumulators. The user guide becomes "Stages"; the design document, ADR 0003, and the LoKI validation are updated, and record that esslivedata resets nothing on context updates, so one context stage costs compute only. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
split raised networkx errors for keys not in the pipeline, repeated an output of a part that a descendant also reads, repeated the inputs of a part given twice, and claimed "no part" for an output whose owner is not an ancestor. Outputs are now checked against the pipeline as Stage does, outputs are listed once, a part given twice is rejected, and the probe stage uses the given scheduler. New tests cover these cases and check that per-iteration work runs once per iteration of its loop; the scheduler tests use a recording scheduler instead of reading private attributes. The user guide example of a part after combining now reads a value from the file part, as the text says. The design document no longer quotes prototype numbers that this repository cannot reproduce, and states which drivers fail loudly when a key is pushed or passed wrongly. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
split now rejects a part that lists an input of one of its ancestors, since its stage would recompute per iteration what the ancestor holds. Error messages name parts by their position in the arguments, because parts that differ only in their parent print alike, and give advice only when an ancestor owns the key. The scheduler test runs every stage with the default scheduler replaced by one that refuses, so a fallback to the default fails. Second review of the documents: resolved contradictions about the StreamProcessor validation, the removal release, and the default scheduler; "contribution" is used only for the values of one member at all accumulation keys; "driver" and the held part are defined and used throughout; the Stage.keys claim excludes inputs; section 6.7's snippet no longer reuses the names of parts; prototype results are marked as such in the ADR. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
visualize_stages now has a default styling: what each stage computes per call is filled with a color of that stage, labeled by the stage's inputs, with the held part, inputs, and outputs marked as by Stage.visualize. This draws the boundaries split derives, such as a value computed by the file stage and passed to the bank stage, and generalizes what Aggregation.visualize drew for two stages. The argument for custom styling is renamed from parts to groups, since parts are now what split takes. The split docstring gets an example driver over runs and banks, and the user guide draws both nested-loop examples. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The nested-loops section introduced Part's parent without saying why split needs it, and the example with work after combining the banks had no driver. The section now writes each computation as plain loops first, so that a part reads as one loop and parent as the enclosing loop, and runs the second example with a driver. The error examples move to the end of the section, and "driver" is defined where the first one appears. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A driver of nested loops had to select, for each stage, the values it takes from what the stages of the enclosing loops returned. compute now ignores values for keys outside Stage.keys, so the driver passes them all on. A value for a key the stage uses but does not take as input is still rejected: the stage holds or computes that key itself, and ignoring the value would hide that, for example, a parameter passed to compute does not override the one the stage was built with. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The nested-loops guide and the Part docstring defined parent first and then said that split does not infer it, which read as a contradiction. They now start from the reason: the pipeline does not say which loop encloses which. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… out split took a Part per loop, with parent naming the enclosing loop, and derived which values each loop forwards to the loops inside it. parent was an indirect way of declaring these connections and was hard to follow, both in the user guide and in review. enclose starts from the stages of the innermost loop and puts them inside a loop over further inputs. The values those stages hold that depend on the new inputs are computed by the returned outer stage, and the rebuilt inner stages take them as inputs. Applied from the inside out, this gives the stages that split gave, using only the frontier that Stage already exposes. Nesting order is the order of construction, so nothing is declared. What split checked when called, a part reading from a loop that does not enclose it, now fails in the driver: a stage left out of a loop still holds the value and rejects it when the driver passes it on. enclose raises if the pipeline computes a key differently than when the stages were built, since it rebuilds them from the pipeline. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Both documents described split with Part and parent. They now describe enclose, built from the inside out, with split recorded as a dropped alternative. They also record the cost: mistakes that split rejected at construction surface in the driver or in warm, and nothing detects the esssans background stage holding the masks of the sample run set on the pipeline. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
enclose sees the stages of one loop, so it cannot tell that a stage was left out of a loop, or that a stage outside a loop reads a value that depends on the inputs of that loop, such as the esssans background stage holding the masks of the sample run set on the pipeline. Such a stage holds the value for the one parameter value set on the pipeline, while the driver varies it. warm sees all stages of a driver, so it now rejects a stage whose held part depends on a parameter that another stage takes as input. Only parameters count: a stage may take as input a value that another stage holds, such as a calibration. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A review of the guide, done without reading the library code, found that the rules of enclose were never stated in one place, that "held" named several different things, and that the driver of nested loops was explained after its first use. The guide now lists the rules of enclose, uses "held" only for the frontier, marks which stage computes which line of each plain-loop example, calls warm in every nested driver, shows the fix that the enclose error names, and states the trade-off of putting the calibration loop outermost. NumValues counts only the values that calibrate keeps, so the mean per file is correct. The legend of Stage.visualize names the part computed once but not kept "Computed once, not kept", and the enclose error names the outer stage as the guide does. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Since warm rejects a stage holding a value that depends on a parameter another stage takes as input, the validation script failed: the background stage held DetectorMasks built from the detector IDs of the sample run set on the pipeline, while the sample stage varies that run. enclose is not the fix, since the background loop is not inside the sample loop. The script now builds the masks from the detector IDs of the empty-beam run, one of the options listed as an open question. The results stay identical to the reference; the masks are built once instead of three times. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Detector banks are a neutron-scattering concept that readers outside ESS do not know. The enclose docstring and the stages user guide now loop over the channels of a file, as in a multi-channel recording. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The nested drivers had been validated only with split, on fake workflows. The new script runs the esssans LoKI workflow over two sample runs and the nine banks of the LoKI test file, with a bank stage enclosed in a loop over runs, and compares with with_banks over with_sample_runs. Results are identical for every bank, and run-level and bank-level providers run as often as in the reference. The second run is a modified copy of the test file, so that a run-level value wrongly held by the bank stage would change the result. The wavelength binning differs from the esssans guide, whose binning makes every value NaN on this 5-pulse file. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The error suggested only enclosing the stage in the loop. For a stage that does not run inside that loop, such as the esssans background stage holding masks built from the sample run, the fix is to change the pipeline so that the held value does not depend on the varied input. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The loop over banks reused the name of a stage result for a dict keyed by run, which mypy rejects as an incompatible assignment. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
enclose went through three designs (Aggregation, split, enclose) and is the least settled part of the proposal. It builds plain stages from the frontier, so nothing in Stage, warm, or the accumulators depends on it. No current use of map/reduce needs it: runs times banks, the one nested shape the ESS packages test, works with one stage per loop and the forwarded values derived from the frontier. ADR 0003 and the design document now describe that pattern and its pitfalls, and leave the helper to a follow-up proposal. Stage.compute again requires exactly the inputs. Ignoring values for unused keys only served drivers of enclosed stages. Relaxing the check later breaks no caller; tightening it would. warm keeps rejecting a stage that holds a value another stage varies. Without enclose it catches a stage of a nested loop that was not rebuilt. The per-stage colors of visualize_stages are now tested with stages built directly. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The ADR and design document described the nested-loops helper as a follow-up proposal with a set of features. Since that design is not settled, they now only name enclose as a likely later addition. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
With several stages inside a loop, each takes only the forwarded values at its own frontier, and since Stage.compute accepts exactly the inputs, the driver selects them per stage. The nested-loops section now says so. The validation section implied that the nested prototypes ran the recipe of section 5. They were built with an earlier helper that builds the same stages; it now says that. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Stage.compute wrapped every held value in a parameter provider on each call, although the held values do not change once computed. warm now builds the per-call graph with the held values once, and compute only adds the inputs to a copy of it. With 20 held values and the naive scheduler, a call takes about 107 us instead of 150 us. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Proposes replacing
Pipeline.map/reducewith stages built from an ordinary flat pipeline, outside the graph. AStageis the part of a pipeline from named input keys to named output keys, with the part that does not depend on the inputs computed once and held. Drivers loop over stages and combine per-member values in accumulators. This PR is the additive half: the ADR (status proposed), the design document with a rollout plan, the code, its tests, and a user-guide page.map,reduce, and cyclebane are untouched; nothing breaks.Why
Map/reduce inside the graph has been implemented twice and is the source of every breaking change and open question in the PEP 695 generics prototypes (#236, #237). Three consumers need an operation it cannot express: run one part of a graph repeatedly with some values supplied per call, combine the results outside, run the rest once.
ess.reduce.streaming.StreamProcessorbuilds that partition by hand (#241 records the requirements), the essapps reduction-service architecture needs it in three places, and every ESS reduction package folds runs at one or more keys and grafts the result back onto the pipeline.The ADR states the decision, the alternatives tried and rejected during prototyping, and the consequences. The design document has the semantics, the survey of every map/reduce site in the ESS packages and esslivedata, the evidence, and a PR-level rollout plan across sciline, ess.reduce, the reduction packages, esslivedata, and essapps.
Nested loops
Loops in loops, such as detector banks within runs or chunks within a stream context, use one stage per loop. The outer stage computes the values that the inner stage would hold and that depend on the outer inputs, and the inner stage takes them as inputs. Both are derived from
Stage.frontier, see the "Nested loops" sections of the ADR and design document. A helper that builds these stages for several levels and rejects a run-level key pushed once per bank is left to the stacked follow-up #248; no current use of map/reduce needs it.An earlier revision of this PR had
Aggregation: a contribute stage, accumulators, and a finalize stage, withcompute(table). No real driver used its finalize stage orcompute, andcompute(table)invited a flat runs-times-banks table, which counts a run-level key once per bank without an error. The ADR records this under alternatives.What is in the PR
docs/developer/adr/0003-replace-map-reduce-with-stages-outside-the-graph.md: the decision. Read this first.docs/developer/architecture-and-design/map-reduce-outside-the-graph.md: design, survey, evidence, migration, rollout plan, settled and open questions.sciline.Stage,sciline.warm: the graph partition.warm, given all stages of a driver, rejects a stage that holds a value depending on a parameter that another stage takes as input, such as a stage of a nested loop that was not rebuilt.sciline.visualize_stages: draws several stages together, by default with a color per stage for what it computes per call.sciline.Accumulator(protocol),sciline.Buffered,sciline.Reduced: accumulator factories for an n-ary function applied to all pushed values, and for a running result of an associative binary function.Reducedis the memory-cheap choice for sums of large arrays; it never updates in place, so it cannot corrupt values the driver still holds.docs/user-guide/stages.ipynb: a guide next to the parameter-tables page. It does not claim to replace anything while the ADR is proposed, but reviewers may find it the quickest way to see the semantics.docs/developer/architecture-and-design/loki_validation.py: the esssans multi-run reduction over stages against the map/reduce reference. Results identical; provider call counts equal, except that the masks are built once instead of once per run.scheduler_or_defaultinsciline.task_graph, shared byTaskGraphandStage. It looks upsciline.task_graph.DaskScheduleron each call, because esslivedata replaces that name to select its scheduler; a test holds sciline to this.Tracking issues: #246 for the sciline side, scipp/ess#746 for ess.reduce, the reduction packages, and esslivedata.
Not in the PR
The helper for nested loops (#248).
provide(key, callable)and areporterargument onStageare listed as a follow-up in the rollout plan. The breaking release that removesmap/reducewaits until the ESS packages and esslivedata have migrated; the plan gives the order. esssans must decide, when it migrates, which sample run's detector IDs mask the background runs;warmrejects the current structure for that reason (design doc section 6.1).Test plan
tests/stage_test.pyandtests/accumulators_test.py; the full suite passes, also with the dask environment.loki_validation.pyrun locally against the reference: identical results.🤖 Generated with Claude Code