Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions docs/api-reference/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
:recursive:

compute_mapped
enclose
get_mapped_node_names
visualize_stages
warm
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ Remove from sciline: `map`, `reduce`, `index_names`, `indices`, `get_mapped_node
Do not add `groupby`.
`constraints=` is removed in the same release, because the PEP 695 generics do not need it.

Add building blocks that work on an ordinary flat pipeline and add nothing to it: `Stage` with `warm`, and accumulators.
Add building blocks that work on an ordinary flat pipeline and add nothing to it: `Stage` with `warm`, accumulators, and `enclose` for nested loops.

### `Stage`: a part of a pipeline, from chosen inputs to chosen outputs

Expand Down Expand Up @@ -106,29 +106,34 @@ finalize.compute({DetectorData: acc.value})

Contributing, combining, and finalizing can run at different times or in different processes, since a contribution is a plain dict; the same loop without the accumulator replaces `compute_mapped`.

### Nested loops
### `enclose`: the stages of nested loops

Banks within runs, or chunks of a stream within a context, are loops in loops.
The inner level reads values that depend on the outer level only, such as a run's monitor, and those should be computed once per outer iteration.
A stage built for the inner loop holds these values at its frontier.
The driver builds one stage per loop: the outer stage computes the values at the frontier of the inner stage that depend on the outer inputs, and the rebuilt inner stage takes these *forwarded values* as inputs:
`enclose` puts the stage inside a loop over the outer inputs, and derives the boundary from the graph:

```python
bank_stage = Stage(pipeline, inputs=(Bank,), outputs=(Numerator, Denominator))
forwarded = Stage(pipeline, inputs=(Filename,), outputs=bank_stage.frontier).dynamic_outputs
run_stage = Stage(pipeline, inputs=(Filename,), outputs=forwarded)
bank_stage = Stage(pipeline, inputs=(*forwarded, Bank), outputs=(Numerator, Denominator))
run_stage, bank_stage = enclose(pipeline, [bank_stage], inputs=(Filename,))
final_stage = Stage(pipeline, inputs=(Numerator, Denominator), outputs=(IofQ,))
```

The run stage computes, once per run, the values that the bank stage held and that depend on `Filename`.
The rebuilt bank stage takes these *forwarded values* as inputs.
Nested loops are built from the inside out: enclosing the result again adds a further outer loop, so each value is computed by the deepest loop whose inputs it depends on.
`enclose` returns plain stages.
The driver writes the loops and holds the values between them:

```python
for filename in filenames:
held = run_stage.compute({Filename: filename})
for bank in banks:
out = bank_stage.compute({**held, Bank: bank})
... # push out[Numerator] and out[Denominator] into accumulators
```

A helper that builds these stages for several levels and rejects a run-level output pushed once per bank is left to a follow-up proposal.
No current use of map/reduce needs it.
`enclose` raises an error where a combined value would otherwise be silently wrong: an output of a stage that does not vary in that stage (a run-level key pushed once per bank).

### Who owns what

Expand All @@ -140,13 +145,13 @@ Everything stateful belongs to the driver: which members exist, which contributi
- Each reduction package returns its own small object in place of the map/reduced pipeline.
For esssans it holds a contribute stage per run type, a finalize stage, and the contributions by filename.
A drop-in replacement for the map/reduced pipeline is a non-goal.
- `StreamProcessor` in ess.reduce becomes a driver over a context stage, with chunk stages and a finalize stage inside its loop, with its existing accumulators and an object that holds the current context.
- `StreamProcessor` in ess.reduce becomes a driver over stages built with `enclose` (a context stage, with chunk stages and a finalize stage inside its loop), with its existing accumulators and an object that holds the current context.
The accumulators fit the `Accumulator` protocol once histogramming moves out of their base class.
- Nested structures, such as banks within runs, are loops over one stage per loop; no graph contains another graph.
- Nested structures, such as banks within runs, are loops over stages built with `enclose`; no graph contains another graph.

### Rollout

`Stage`, `warm`, and the accumulators are added in a minor release.
`Stage`, `warm`, `enclose`, and the accumulators are added in a minor release.
The ESS packages then migrate one at a time while map/reduce still exists.
The removal comes last, in a major release (recommended; a `sciline.v2` namespace is the open alternative, see below).
Users who depend on map/reduce and do not need the new generics can stay on the last release before the removal.
Expand All @@ -164,6 +169,11 @@ A `sciline.v2` namespace that keeps the old `Pipeline` would serve them equally,
`compute(table)` invited a flat runs-times-banks table, which counts run-level keys once per bank without an error, and its one guarantee is `Stage.dynamic_outputs`.
- **Choosing the per-run keys of nested loops by hand.**
Prototyped for LoKI and Bifrost: a key left out is recomputed per bank with a correct result, so the mistake goes unnoticed.
- **`split(pipeline, *parts)`, with one `Part` per loop** that names the part of its enclosing loop as `parent`.
Prototyped and dropped: `parent` declared the forwarded values only indirectly, which made the API hard to understand.
It built the same stages as `enclose`.
Since it saw all loops at once, it rejected some mistakes at construction, such as a read from a loop that does not enclose the reader.
`enclose` builds on the frontier, which users of `Stage` already know.
- **Let an aggregation turn back into a pipeline** (`as_pipeline()`), so that the `with_*` helpers could keep returning a pipeline.
Prototyped and dropped: it became the only way to combine two aggregations, fixed the members at construction, and hid a loop inside a provider.
- **Let the sciline objects hold state** (member tables, kept contributions, rules for what a parameter change invalidates).
Expand All @@ -189,18 +199,19 @@ A `sciline.v2` namespace that keeps the old `Pipeline` would serve them equally,
Accumulator, contribution, and contribute/combine follow Beam, Flink, and Spark.
- Every parameter is set on one flat pipeline, and everything held is a plain object that the driver can inspect, clear, or serialize.
- The prototype reproduces the LoKI multi-run reduction with identical results and the same or fewer provider calls, and adding a run costs only that run's contribution.
`enclose` reproduces the LoKI reduction over runs times banks with identical results and provider calls.
Prototype drivers, which are not in this repository, gave the results of plain loops for banks or triplets times runs, and a prototype `StreamProcessor` gave those of the real class.
They were built with `split` (see alternatives), which builds the same stages as `enclose` for these shapes.

### Negative

- Breaking for the `with_*` helpers in esssans, essreflectometry, and bifrost, for `ess.reduce.parameter_mappers` and the widgets built on it, for essreflectometry's `BatchProcessor`, for notebooks, and for the bifrost bank fold in esslivedata.
- The widgets need one interface across the package objects, a base class or one generic object, decided when the second package migrates.
- Parallelism over members is the driver's job; with map/reduce, dask ran members in threads for free (about 1.7 s of 7.5 s in the LoKI validation script).
- A map/reduce inside the per-member work of another one (the pixel masks in esssans) becomes a list parameter and a provider.
- A driver must push only `stage.dynamic_outputs`; other keys are counted once per member.
In nested loops this is not enough: a run-level key declared as an output of the bank stage is dynamic there and is counted once per bank, without an error.
`warm` catches a stage inside a loop that still holds a value depending on the inputs of the loop, if the driver warms all its stages together.
- With one context stage, a `StreamProcessor` context update recomputes all context-derived values (in esslivedata compute only; no accumulators reset).
- A one-level driver that builds stages by hand must push only `stage.dynamic_outputs`; other keys are counted once per member. For nested loops, `enclose` is needed: a run-level key in a stage with inputs `(Filename, Bank)` is dynamic and still counted once per bank.
- Each loop is inside at most one other: bank-only work is computed once per run and bank, and a `StreamProcessor` context update recomputes all context-derived values (in esslivedata compute only; no accumulators reset).
- `enclose` sees one loop at a time, so some mistakes are found when the stages are warmed or run, not when they are built: `warm` rejects a stage left out of an `enclose` call, and a stage enclosed twice over the same inputs fails in the driver.
- `warm` rejects the esssans background stage reading the masks of the one sample run set on the pipeline; esssans has to decide which run's detector IDs the background masks use.
- esslivedata selects its scheduler by replacing `sciline.task_graph.DaskScheduler`; sciline keeps that working for stages, until it offers a public way.
- Stages keep their held values, and package objects keep contributions (binned events for some workflows), so package objects need a way to clear them.
Expand Down
227 changes: 227 additions & 0 deletions docs/developer/architecture-and-design/loki_banks_validation.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,227 @@
"""Validate enclose on the esssans LoKI reduction over runs times detector banks.

Reference: with_banks(with_sample_runs(...)) (sciline map/reduce), which gives
IntensityQ[SampleRun] per bank, combined over the sample runs. esssans does not
combine over banks, since banks have different Q resolution.
Prototype: a stage per bank from NeXusDetectorName to the accumulation keys, enclosed
in a loop over Filename[SampleRun], and a final stage from the accumulation keys to
IntensityQ[SampleRun]. The driver keeps one set of accumulators per bank.

The only multi-bank test file is loki_coda_file. The second sample run is a copy with
the detector event IDs reversed and the monitor event time offsets scaled, so that
both the per-bank and the per-run values differ between the runs. A value that
depends on the run but is held by the bank stage would then give a wrong result.

Run from this directory with
python loki_banks_validation.py
"""

# ruff: noqa: T201, S101

from __future__ import annotations

import shutil
import tempfile
import time
from pathlib import Path
from typing import Any

import ess.loki.data # noqa: F401
import h5py
import numpy as np
import scipp as sc
from ess import loki, sans
from ess.reduce.nexus.types import NeXusData
from ess.reduce.nexus.workflow import assemble_detector_data
from ess.sans.normalization import norm_monitor_term
from ess.sans.types import (
BeamCenter,
CorrectedMonitor,
CorrectForGravity,
Denominator,
DetectorMasks,
DirectBeam,
EmptyBeamRun,
EmptyDetector,
Filename,
Incident,
IntensityQ,
LookupTableFilename,
MonitorTerm,
NeXusDetectorName,
NormalizedQ,
Numerator,
QBins,
RawDetector,
ReturnEvents,
SampleRun,
TransmissionFraction,
TransmissionRun,
UncertaintyBroadcastMode,
WavelengthBins,
)
from ess.sans.workflow import merge_contributions
from scipp.testing import assert_allclose, assert_identical
from scippnexus import NXdetector

import sciline
from sciline import Buffered, Stage, enclose, warm

TARGET = IntensityQ[SampleRun]
ACCUMULATED = (NormalizedQ[SampleRun, Numerator], NormalizedQ[SampleRun, Denominator])
BANKS = [f'loki_detector_{i}' for i in range(9)]
SCHEDULER = sciline.scheduler.NaiveScheduler()

calls: dict[str, int] = {}


def _hit(name: str) -> None:
calls[name] = calls.get(name, 0) + 1


# Counted copies of providers, with the signatures for SampleRun only, so that the
# counts are per sample run (monitor term) and per sample run and bank (detector).
def counted_norm_monitor_term(
incident_monitor: CorrectedMonitor[SampleRun, Incident],
transmission_fraction: TransmissionFraction[SampleRun],
) -> MonitorTerm[SampleRun]:
_hit('monitor_term')
return norm_monitor_term(incident_monitor, transmission_fraction)


def counted_assemble_detector_data(
detector: EmptyDetector[SampleRun], neutron_data: NeXusData[NXdetector, SampleRun]
) -> RawDetector[SampleRun]:
_hit('assemble_detector')
return assemble_detector_data(detector, neutron_data)


def make_workflow() -> sciline.Pipeline:
# As in the esssans user guide loki-reduction-ess.ipynb, except for the wavelength
# bins. The test file holds 5 pulses, and with the 200 bins of the guide some
# bins have no monitor counts. The transmission fraction is NaN there, and since
# I(Q) sums over wavelength, every value of I(Q) would be NaN.
wf: sciline.Pipeline = loki.LokiWorkflow()
wf[WavelengthBins] = sc.linspace('wavelength', 1.0, 10.0, 21, unit='angstrom')
wf[QBins] = sc.linspace('Q', start=0.01, stop=0.3, num=101, unit='1/angstrom')
wf[CorrectForGravity] = True
wf[UncertaintyBroadcastMode] = UncertaintyBroadcastMode.upper_bound
wf[ReturnEvents] = False
wf[BeamCenter] = sc.vector([0.0, 0.0, 0.0], unit='m')
wf[DirectBeam] = None
wf[DetectorMasks] = {}
wf[LookupTableFilename] = loki.data.loki_lookup_table_no_choppers()
wf[Filename[EmptyBeamRun]] = loki.data.loki_coda_file()
wf[Filename[TransmissionRun[SampleRun]]] = loki.data.loki_coda_file()
wf.insert(counted_norm_monitor_term)
wf.insert(counted_assemble_detector_data)
return wf


def make_second_run(directory: str) -> str:
"""Copy the test file with different detector and monitor events."""
path = str(Path(directory) / 'second_run.hdf')
shutil.copy(loki.data.loki_coda_file(), path)
with h5py.File(path, 'r+') as f:
instrument = f['entry/instrument']
for bank in BANKS:
ids = instrument[f'{bank}/detector_events/event_id']
ids[...] = ids[()][::-1]
for name in instrument:
if name.startswith('beam_monitor'):
offset = instrument[f'{name}/monitor_events/event_time_offset']
offset[...] = (offset[()] * 0.9).astype(offset.dtype)
return path


def compare(name: str, a: sc.DataArray, b: sc.DataArray) -> None:
# assert_identical treats NaN as equal, so check that there is something to compare.
# NaN remains in Q bins without counts.
finite = np.isfinite(a.values).mean()
assert finite > 0, name
try:
assert_identical(a, b)
print(f' {name}: identical ({finite:.0%} finite)')
except AssertionError:
assert_allclose(a, b, rtol=sc.scalar(1e-12))
print(
f' {name}: allclose (rtol 1e-12) but not identical ({finite:.0%} finite)'
)


def show(label: str, t0: float) -> dict[str, int]:
print(f' {label}: {time.perf_counter() - t0:.1f} s, calls {calls}')
return dict(calls)


def main() -> None:
wf = make_workflow()
with tempfile.TemporaryDirectory() as tmp:
runs = [str(loki.data.loki_coda_file()), make_second_run(tmp)]

# --- Reference: map/reduce ------------------------------------------------
calls.clear()
t0 = time.perf_counter()
ref = sans.with_banks(sans.with_sample_runs(wf, runs=runs), banks=BANKS)
# compute_mapped takes no scheduler, so compute the mapped nodes directly.
names = sciline.get_mapped_node_names(ref, TARGET)
computed = ref.compute(names, scheduler=SCHEDULER)
ref_results = {bank: computed[name] for bank, name in names.items()}
t_ref = time.perf_counter() - t0
ref_calls = show('reference', t0)

# --- Prototype: a bank stage enclosed in a loop over runs --------------------
calls.clear()
t0 = time.perf_counter()
bank_stage = Stage(
wf, outputs=ACCUMULATED, inputs=(NeXusDetectorName,), scheduler=SCHEDULER
)
run_stage, bank_stage = enclose(
wf, [bank_stage], inputs=(Filename[SampleRun],), scheduler=SCHEDULER
)
final = Stage(wf, outputs=(TARGET,), inputs=ACCUMULATED, scheduler=SCHEDULER)
warm(run_stage, bank_stage, final)
acc = {
b: {k: Buffered(merge_contributions)() for k in ACCUMULATED} for b in BANKS
}
held_per_run: dict[str, dict[Any, Any]] = {}
per_run: dict[str, dict[str, Any]] = {}
for run in runs:
held = run_stage.compute({Filename[SampleRun]: run})
held_per_run[run] = held
for bank in BANKS:
values = bank_stage.compute({**held, NeXusDetectorName: bank})
per_run.setdefault(bank, {})[run] = values
for key in bank_stage.dynamic_outputs:
acc[bank][key].push(values[key])
results = {
b: final.compute({k: a.value for k, a in acc[b].items()})[TARGET]
for b in BANKS
}
t_proto = time.perf_counter() - t0
proto_calls = show('prototype', t0)

print('forwarded by enclose (run_stage.outputs):')
for key in run_stage.outputs:
print(f' {key}')

print('the two runs differ in the monitor term and in every bank:')
a, b = (held[MonitorTerm[SampleRun]] for held in held_per_run.values())
assert not sc.identical(a, b)
numerator = ACCUMULATED[0]
for bank, by_run in per_run.items():
a, b = (v[numerator] for v in by_run.values())
assert not sc.identical(a, b), bank
print(' yes')

print('per-bank results, prototype vs reference:')
for bank in BANKS:
compare(bank, results[bank], ref_results[bank])

print(f'\nSUMMARY: reference {t_ref:.1f} s, prototype {t_proto:.1f} s')
print(f' reference calls: {ref_calls}')
print(f' prototype calls: {proto_calls}')


if __name__ == '__main__':
main()
Loading
Loading