did-btcr2-js

ADR 040: Multi-Cohort Aggregation Service Runner

Status: Accepted

Date: 2026-06-22

Branch / PR: feat/aggregation-multi-cohort

References: ADR 008, ADR 020, ADR 027, ADR 038, ADR 039

Context

The goal is one Aggregation Service advertises many cohorts simultaneously, and a participant subscribes across several of them with per-cohort opt-in. Today the high-level Runner facade drives exactly one cohort to completion and exits. This ADR records how multi-cohort is delivered and, just as importantly, where it is not needed.

The state machines and transport are already N-cohort.

Every single-cohort assumption lives in AggregationServiceRunner. The facade (ADR 020 layer-3 runner) hard-wires one cohort through:

The participant runner is already multi-subscription. AggregationParticipantRunner holds no per-cohort fields, delegates to its session, and demuxes by msg.body.cohortId in every handler. Only the static convenience joinFirst() (participant-runner.ts:180-200) collapses it to one cohort via a single once('cohort-complete'). The structural capability is present; only a multi-join entry point is missing.

Where this subsystem is headed (the lens for the choices below):

The destination is a long-lived, multi-cohort service that gets extracted. In each decision below that pulls toward the more general and more consistent option, and the cost is modest because the state machines already support it: the runner facade is only being made to express generality that already exists underneath.

Decision

Refactor AggregationServiceRunner into a long-lived multiplexer keyed by cohortId. The multi-cohort orchestration leaves the state machines and transport/* structurally unchanged; the only state-machine edit is one additive read accessor (point 6) that fixes a pre-existing participant-completion gap. Add a participant-side multi-join convenience so multi-cohort works end to end, and make every event carry a top-level cohortId.

  1. Per-cohort RunContext, stored in #contexts: Map<string, RunContext>. Each advertised cohort owns its state: cohortId, its CohortConfig, the deferred resolve / reject and the completion promise handed back to the caller, a per-cohort finalizing guard, its own ttlTimer / phaseTimer / lastObservedPhase, its stopAdvertRepeat handle, and a settled flag so a late timer or trailing message cannot double-settle it. Every singular #-field enumerated in Context moves into this struct.

  2. advertiseCohort(config): { cohortId; completion } is the additive multi-cohort entry point. It calls session.createCohort(config), builds the RunContext, starts that context’s timers, sends the advert (and its republish loop), and returns the cohort id plus a per-cohort completion promise. It is callable many times on one runner. Transport handler registration is a one-time, idempotent setup independent of any cohort (the handlers are already cohort-agnostic).

  3. Per-cohort completion is the load-bearing primitive; the runner is a long-lived multiplexer. Each RunContext.completion resolves with that cohort’s AggregationResult (which already carries cohortId) when its signing authorization lands, and rejects via a per-cohort failure path. A caller awaits one cohort via the returned completion. A thin runAll(): Promise<AggregationResult[]> convenience drains the currently outstanding contexts; its semantics are explicitly dynamic - new cohorts may be advertised between calls, and runAll() settles when the live set empties, not against a frozen snapshot. The runner is a persistent service, not a one-shot batch.

  4. Failure and teardown are per-cohort, never global-by-accident. A TTL or phase-stall expiry calls #failCohort(cohortId, err), which clears only that context’s timers and advert loop, removeCohorts its state, rejects its completion, and emits cohort-failed with the cohortId. It does not touch sibling contexts and does not unregister the shared transport handlers. stopCohort(cohortId) is the deliberate single-cohort teardown; stop() becomes stop-all (iterate contexts, then unregister the shared handlers once). A separate runner-fatal path handles transport-level failures by failing all contexts. The runner reclaims its own per-cohort bookkeeping (timers, advert loop, RunContext) on every settle, but removes the cohort from the state machine (session.removeCohort) on failure and stop only: a successfully completed cohort is left in session so callers can still read its beaconAddress / cohort via session.getCohort(result.cohortId) and reclaim it explicitly when done.

  5. run() and solo() are preserved as thin wrappers. run(): Promise<AggregationResult> becomes advertiseCohort(constructorConfig).completion - byte-for-byte behavior for the single-cohort case - so AggregationRunner.solo() and every existing await runner.run() call site keep working unchanged. The constructor’s config becomes optional (required only for the run() convenience path; omit it when driving via advertiseCohort).

  6. Add a participant multi-join convenience now, and fix the completion sidecar. AggregationParticipantRunner gains joinMatching(options, count) - the bounded N-cohort generalization of joinFirst - to deliver participant multi-subscription (one participant subscribing across several cohorts) and to write this change’s end-to-end test (one service running two cohorts end to end, a participant joining both), which joinFirst() cannot express. Open-ended subscription is already the constructor + shouldJoin + start() path, so no never-resolving joinAll() is added; joinFirst() stays as the single-cohort convenience. Folded in: the participant’s cohort-complete could not surface its sidecar (CAS Announcement map / SMT inclusion proof) because it read the phase-filtered pendingValidations (which lists only the AwaitingValidation phase) and so returned nothing once the cohort reached Complete. A small additive AggregationParticipant.getValidation(cohortId) accessor returns the retained validation regardless of phase, so the participant receives its sidecar at completion. This is the only state-machine touch and it is purely additive.

  7. Every event carries a top-level cohortId. Add cohortId to the service events that structurally lack it (participant-accepted, update-received, validation-received, signing-started, nonce-received) and surface a top-level cohortId uniformly on the participant events whose id is currently only nested. This establishes one invariant - every aggregation event payload identifies its cohort - which a service-monitoring dashboard and the extracted package both rely on. The additions are receive-only fields (a listener reads a subset), so they are source-compatible for consumers.

  8. The multi-cohort orchestration requires no state-machine or transport changes. Because AggregationService / AggregationParticipant already demux by cohortId and the transport already multiplexes per actor, the concurrency work is confined to the runner facade plus events.ts. The sole state-machine edit is the additive getValidation accessor in point 6 - a pre-existing participant-completion fix, not multi-cohort plumbing; cohort.ts, phases.ts, signing-session.ts, messages/*, and transport/* are untouched. Any structural change forced into the demux logic would signal a missed coupling and is out of scope.

Lifecycle model considered

Model Long-term fit Cost / risk
Multiplexer - per-cohort completion + dynamic runAll() (chosen) Matches the long-lived service daemon and cohorts persisting across signing rounds directly; the extracted API is the general one Runner is stateful: must reclaim each RunContext on settle and must not unregister shared transport handlers until stop(); runAll() needs an explicit drain contract
One-shot batch - runAll() over a frozen set, no advertising after Simplest for a single script Dead end for the daemon: cannot advertise after runAll(), forcing a second refactor once the service deployment lands - the churn is paid twice
Per-cohort only - no aggregate method Cleanest primitive, maximally composable Every caller (CLI, demo, tests) re-hand-rolls “wait for this batch,” repeatedly

“Batch vs long-lived” and “ship runAll() or not” are nearly orthogonal: even the per-cohort-only option is long-lived. The load-bearing primitive in all three is the per-cohort completion promise; runAll() is only a convenience on top of it. The chosen model is the multiplexer built on that primitive.

Rejected alternatives

Consequences

Positive

Negative

Accepted

References