Repository navigation
Cleanup/fix reply mechanism to request from frontend #445
Description
Activity
- changed the title
[-]Split job status reporting (periodic) from replies to request from frontend.[/-][+]Cleanup/fix reply mechanism to request from frontend[/+]on Sep 10, 2025 - added a commit that references this issue
on Oct 1, 2025 Partial plan forming in my head:
- Consider
JobCommandandWorkflowConfig: These should map to Nicos compatible commands. WorkflowStatusshould be changed to a Nicos compatible command ack.
WorkflowConfigcontains both "config" and an implicit "start" command. Should we require sending two messages here? We could first send the config, then the "start". Or configuring a job that does not exist could auto-start it (basedJobSchedule, which we can include in theconfigpart of the message)?"Device" naming (e.g., for a
clearcommand) is a bit unusual, since it needs to refer to a Job ID, but I think this is equivalent to what we have for the status heartbeat?- Consider
- added 2 commits that reference this issue
on Nov 19, 2025 Design Notes: Splitting WorkflowConfig into Config + Start
During work on the command acknowledgement feature, we analyzed the relationship between
message_idandjob_numberinWorkflowConfig. This comment documents the findings for future reference when implementing the config/start split.Current State
WorkflowConfigcurrently conflates "configure" and "start" into a single command, containing both:message_id: for command/response correlation (ACK pattern)job_number: for job identity
Why Both Are Necessary (Not Redundant)
Identifier Purpose Lifetime Generated by message_idCommand/response ACK correlation Transient (until ACK received) Frontend job_numberJob identity for results and control Persistent (entire job lifecycle) Frontend They serve fundamentally different purposes and cannot be consolidated.
Proposed Split Design
Per the NICOS-BeamLime interface spec, commands should be separate. Here is a proposed design:
WorkflowConfig (config-only message):
class WorkflowConfig(BaseModel): identifier: WorkflowId message_id: str # For ACK, AND serves as "config handle" params: dict[str, Any] aux_source_names: dict[str, str] schedule: JobSchedule # Or move to start # NO job_number - not starting a job yet
WorkflowStart (new message type):
class WorkflowStart(BaseModel): message_id: str # For THIS command's ACK config_ref: str # References config's message_id job_number: JobNumber # Identity of the new job source_names: list[str] # Optional: which sources to start
Key Insight: message_id as Config Handle
The config's
message_idserves double duty:- ACK correlation (standard use)
- Config handle for subsequent start commands
Flow Example
Frontend Backend | | |-- WorkflowConfig ----------------->| | (message_id=abc, params, aux) | | | stores config[abc] = {...} |<-- ACK (message_id=abc) ----------| | | |-- WorkflowStart ------------------>| | (message_id=xyz, | | config_ref=abc, | creates job using config[abc] | job_number=job1) | |<-- ACK (message_id=xyz) ----------|Benefits of Split
-
Independent jobs with shared config: Frontend and NICOS can both send
WorkflowStartreferencing the sameconfig_ref=abcbut with differentjob_numbers -
Cleaner semantics: Configuration is stable, jobs are instances of that config
-
Aligns with NICOS interface: Matches the
config/startcommand pattern in the spec
Implementation Notes
- Backend needs to store configs indexed by
message_id(the config handle) - Consider TTL/cleanup for unused configs
JobOrchestrator.commit_workflow()would become two steps: send config, then send start- The current
staged_jobspattern inJobOrchestratoralready provides the "staging area" concept
Code comments have been added to
WorkflowConfigandJobOrchestratorreferencing this design.- added 2 commits that reference this issue
on Jan 5, 2026 Config/Start Split Design
Context
WorkflowConfigcurrently conflates "configure" and "start" into a single command.
This document captures design decisions for splitting it into separate config and start
messages, aligning with the NICOS-BeamLime interface spec.Agreed Design Decisions
ConfigNumber (analogous to JobNumber)
A dedicated
ConfigNumber = uuid.UUIDtype identifies a configuration, analogous to how
JobNumberidentifies a job. This replaces the earlier idea of reusing
WorkflowConfig.message_idas a "config handle".Reasons for a dedicated type:
- Lifetime mismatch:
message_idis transient (discarded after ACK). Config identity
is persistent — it must outlive the ACK handshake so that a start command can reference
it later. - Retryability: Retrying a failed config command generates a new
message_id(new ACK
cycle) but represents the same config. WithConfigNumber, the retry is naturally
idempotent. - Separation of concerns: Each identifier does one thing:
message_id— ACK correlation (transient)ConfigNumber— config identity (persistent)JobNumber— job identity (persistent)
- Consistency: Follows the established
JobNumberpattern in the codebase.
WorkflowConfig is a single atomic message
A
WorkflowConfigmessage contains the full config set — all per-source configs in one
message, identified by a singleConfigNumber. This mirrors the frontend's existing
staged_jobs: dict[SourceName, JobConfig]structure, which is essentially what gets sent:class WorkflowConfig(BaseModel): message_id: str # For ACK identifier: WorkflowId config_number: ConfigNumber jobs: dict[SourceName, JobConfig] # The full config set
One message, one
ConfigNumber, one ACK. No accumulation logic needed on the backend —
the config set is atomic.This is a direct mapping from the frontend's
WorkflowState.staged_jobs. The staging
action that currently only updates local state will now also send this to Kafka.Config registry on the backend
The backend stores configs by
ConfigNumber:dict[ConfigNumber, WorkflowConfig]When a start command arrives referencing a
ConfigNumber, the backend looks up the config
and starts jobs for each source in itsjobsdict.TTL/cleanup will be needed for configs that are never started, but is not critical
initially given the low volume (config changes come from modal dialogs, not continuous
interaction).NICOS compatibility
The NICOS spec's
startcommand currently has no config reference — it implicitly uses the
latest config for a device. We will extend the spec to includeConfigNumber. As a
fallback, omittingconfig_numberon a start command could fall back to "latest config for
this device", preserving compatibility with the original NICOS model.Unified command model: start as a JobCommand action
All commands —
start,stop,reset,remove, etc. — use the sameJobCommandmodel,
aligning with the NICOS spec wherestart,stop,config,clearare variants of one
command structure.The key enabler is targeting by
JobNumberinstead ofJobId. Currently,stopand
resetsend N individual commands (one perJobId, i.e., per source). But since all jobs
in a set share the sameJobNumber, a single command targeting theJobNumberis
sufficient. This makesstartstructurally consistent with other actions:- Start:
job_number=J, config_number=C— "start jobs for config C with number J" - Stop:
job_number=J— "stop all jobs with number J" - Reset:
job_number=J— "reset all jobs with number J"
The only structural difference is that
startcarries one extra field (config_number).
This is a single optional field, not a proliferation of mutually exclusive optionals.class JobCommand(BaseModel): message_id: str | None action: JobAction # start, stop, reset, remove, ... job_number: JobNumber | None config_number: ConfigNumber | None # only for start
Per-source targeting (
JobId) is removed. If fine-grained per-source control is needed
in the future, it should be added uniformly across all actions (including start), not just
for stop/reset.Remove workflow_id from JobCommand
The
workflow_idfield onJobCommandis unused in practice — every call site in the
codebase targets byjob_idor broadcasts to all jobs. No frontend or test code ever
setsworkflow_idon aJobCommand.More importantly,
workflow_idtargeting is the wrong granularity for a multi-client world.
If NICOS sendsstop(workflow_id=W), it would kill the dashboard's jobs for that workflow
too.JobNumberis the correct scope: each client creates its ownJobNumbers, so
stop(job_number=J)affects only that client's jobs.The backend dispatch simplifies from three branches to two:
job_numberspecified → operate on all jobs with that number- Neither specified → operate on all jobs (emergency kill switch)
Config sets are immutable, full-replacement only
A config set (all configs sharing a
ConfigNumber) is immutable once sent. There is no
mechanism to add, remove, or modify individual sources within an existing config set.To change anything — parameters, aux sources, or which sources are included — the frontend
creates a new config set (newConfigNumber). This keeps the protocol simple and avoids
incremental-mutation complexity (remove-source commands, partial updates, ordering hazards).Eager config sending (config = staging)
Config messages are sent to Kafka eagerly — when the user confirms a config dialog, not
when they click "Start". The act of staging a config is sending it to the backend.Benefits:
- Early validation: The backend ACKs (or ERRs) the config before the user attempts to
start. Invalid parameters surface immediately. - Decoupled config and start lifecycles: The config exists on the backend independently
of any job. Multiple clients can reference it at different times. - NICOS reuse: NICOS can reference a config that the dashboard user sent earlier, even
if the dashboard has since moved on to different parameters. Both configs coexist in the
registry, each identified by itsConfigNumber.
The overhead is negligible: config changes come from a modal dialog (not continuous
sliders), producing at most a handful of small JSON messages per user action.Backend config registry as shared state
The backend config registry (
dict[ConfigNumber, WorkflowConfig]) serves as shared state
between clients. Configs are pinned as long as they are referenced
by a running job. Unreferenced configs can expire after a TTL, but cleanup is not critical
initially given the low volume.Example multi-client scenario:
Dashboard Backend NICOS | | | |-- config C1 ------------->| store C1 | |<-- ACK -------------------| | |-- start(C1, J1) --------->| run J1 with C1 | | | | | (user tweaks params) | | |-- config C2 ------------->| store C2 | |<-- ACK -------------------| | |-- stop(J1) -------------->| stop J1 | |-- start(C2, J2) --------->| run J2 with C2 | | | | | | (NICOS wants C1 back) | | |<-- start(C1, J3) ----------| | | run J3 with C1 | | | | | J2 runs with C2 | J2 and J3 coexist | J3 runs with C1Per-worker ACK per command, with per-source error detail
Each command (config, start, stop, reset) produces one
CommandAcknowledgementper worker
that handles sources affected by the command. In the single-worker case (the common
deployment today), this means exactly one ACK per command.class CommandAcknowledgement(BaseModel): message_id: str response: AcknowledgementResponse # ACK if all ok, ERR if any failed worker_id: str | None = None # Identifies the responding worker (multi-worker) sources_handled: frozenset[SourceName] | None = None # Sources this worker processed (multi-worker) message: str | None = None # Summary error message source_errors: dict[SourceName, str] | None = None # Per-source detail on partial failure
ACK scope by command type:
- Config: Transactional validation. Either the entire config set is valid (
ACK) or it
isn't (ERRwith validation errors). No partial acceptance — like a form submission that
returns all field errors at once. - Start: The backend iterates over sources in the config set and creates jobs.
ACKif
all succeed,ERRwithsource_errorsif any fail. Partial failure is rare in practice:
if the config was already validated and ACK'd, startup failures are unlikely for only some
sources. - Stop/Reset: Same pattern.
ACKif all jobs affected by theJobNumberwere
successfully stopped/reset,ERRwith detail if any failed.
Why per-worker ACK, not per-source ACK:
- 1 command → 1 response per worker is a cleaner protocol contract than per-source.
Each worker processes its sources in a synchronous loop — the per-worker ACK is its
natural atomic unit. The currentPendingCommandTrackerexists precisely because N
commands produce N ACKs that must be accumulated — trackingexpected_count,
is_complete, and timeouts. With per-worker ACK, accumulation is reduced to tracking
source coverage (see multi-worker extension below), and in the single-worker case it
is eliminated entirely. - Within a worker, accumulation does not exist. The worker loops over its sources
synchronously, collects errors into a dict, and returns one response. This is a bounded,
synchronous loop — qualitatively different fromPendingCommandTracker's asynchronous
message accumulation across Kafka round-trips. - Performance: One Kafka round-trip per worker for serialization/deserialization instead
of one per source. Small N in both cases, but worth not paying unnecessarily.
Error reporting boundaries:
Phase Error type Mechanism Command processing Synchronous failures (config not found, validation error, workflow not found, init crash) CommandAcknowledgement(ACK/ERR)Job runtime Asynchronous failures (data processing errors, workflow exceptions) Job heartbeats ( JobStatuswithstate=error)Startup errors fall in the first category — they happen during synchronous command
processing, before a job exists to emit heartbeats. The ACK is the correct and only
mechanism for reporting them.Multi-worker extension: per-worker ACK with source coverage
The single-ACK model assumes one backend instance handles all sources for an instrument.
For scalability, we may split a service (e.g.,detector_data) across multiple workers,
each handling a distinct subset of source names. In that case, no single worker can produce
a complete ACK for a command that spans multiple workers.The extension is straightforward: each worker ACKs independently for the sources it
handles.CommandAcknowledgementgainsworker_idandsources_handledfields (see
model above). The client tracks source coverage rather than counting responses:expected_sources = set(config.jobs.keys()) acked_sources: set[SourceName] = set() for ack in receive_acks(message_id, timeout=...): acked_sources |= ack.sources_handled if ack.response == ERR: collect_errors(ack) if acked_sources >= expected_sources: break # complete
Why per-worker, not per-source:
A worker processes its sources in a synchronous loop — it always knows the outcome for all
of them at once. There is no scenario where a worker can ACK for source A but hasn't
finished evaluating source B yet. Per-source ACKs would fragment a single synchronous
result into N messages for no reason. Per-worker is the natural atomic unit, and keeps the
number of ACK messages proportional to workers (a small, bounded number) rather than
sources (which could grow significantly for instruments like DREAM).Why not a coordinator:
A coordinator service that aggregates per-worker ACKs into a single client-facing ACK
would contain the same accumulation logic the client needs anyway — it just relocates it
to a new service. Since we expect only a handful of workers per service type, the overhead
of a dedicated coordinator is not justified.Compatibility with the single-worker case:
When
worker_idandsources_handledareNone, the ACK covers all sources implicitly —
identical to the base design. The multi-worker fields are additive. Clients that don't
need multi-worker support can ignore them.Timeout and partial failure:
If a worker is down, its sources will never be ACK'd. The client detects this via timeout:
the expected sources are known (the client sent them), and coverage that remains incomplete
after a deadline indicates a missing worker. This is surfaced as an error for the
uncovered sources.Why atomic config messages hold up under scrutiny
The single-message-per-config-set design was questioned in light of the per-source ACK
detail: if the response needs per-source granularity, should the request be per-source too?No — the granularity of the request and the response don't have to match. A batch operation
naturally produces per-item results. The per-source detail insource_errorsis error
reporting, not per-source protocol.Arguments for keeping atomic config messages:
- The dashboard's natural unit is the set. Users configure all sources in a dialog,
stage everything, start everything. Per-source messages would split a single user action
into N messages only to reassemble them on the backend. - Per-source configs reintroduce accumulation. The backend would need to know when it has
"all" per-source configs, or the start command would need to reference a list of
ConfigNumbers (or fall back to "latest config per source", which is implicit and
fragile). - A config set with one source IS a per-source config. NICOS can send a config set
containing a single source — same protocol, same ACK shape, no impedance mismatch. - Transactional validation is a feature. Rejecting the whole set on any validation error
lets the user see all problems at once and fix them before resubmitting.
Message flow (single client, single worker)
Frontend Backend | | | (user confirms config dialog) | | | |-- WorkflowConfig ------------------->| | (message_id=m1, config_number=C, | | identifier=W, | | jobs={src_a: ..., src_b: ...}) | store config[C] | | |<-- ACK (message_id=m1) --------------| (config validated and stored) | | | (user clicks "Start") | | | |-- JobCommand (start) --------------->| | (message_id=m2, | | action=start, | | config_number=C, | create jobs from config[C] | job_number=J) | |<-- ACK (message_id=m2) --------------| | | | ... time passes ... | | | |-- JobCommand (stop) ---------------->| | (message_id=m3, | | action=stop, | stop all jobs with number J | job_number=J) | |<-- ACK (message_id=m3) --------------|Open Questions
ConfigProcessor replacement
The current
ConfigProcessordeduplicates by taking only the latest config per key per
source name. With the new design (atomic config messages,JobCommand-based start), the
dispatch logic changes significantly.ConfigProcessorwill likely need a rework or
replacement, though the single-message config model simplifies things compared to
multi-message accumulation.Config discoverability for NICOS
NICOS is subscribed to the shared commands Kafka topic, so it can observe configs sent by
the dashboard (or any other client) and reference theirConfigNumbers. No special
discovery mechanism is needed beyond what the protocol already provides.Config loss if backend down
If frontend sends (stages) config while backend is down, starting a job later would reference an unknown config ID => frontend needs to be able to resend config?
- Lifetime mismatch:
- addedarea:wireKafka topics, message schemas, command/reply contracts, incl. toward NICOS/ECDCKafka topics, message schemas, command/reply contracts, incl. toward NICOS/ECDC
on Aug 11, 2026 The original ask here is done, and the design that grew on top of it was overtaken by
ADR 0006.Done. The legacy
WorkflowStatusreply type is gone from the tree.
CommandAcknowledgementimplements the command-acknowledgement shape this issue linked,
and the dashboard consumes it throughPendingCommandTracker. ADR 0006 went further
than the plan in the first comment:JobCommandis the command NICOS produces, and
WorkflowIdaccepts itsinstrument/name/versionstring form on the wire so NICOS
sends the contract string verbatim.Overtaken. The config/start split rested on independent jobs sharing a
configuration — "Frontend and NICOS can both sendWorkflowStartreferencing the same
config_refbut with differentjob_numbers", and the multi-client flow drawn from it.
ADR 0006 settled the opposite: the dashboard configures the workflow backing the
devices, NICOS's whole command surface isresetkeyed byWorkflowId, and giving
NICOS its own jobs was considered and rejected ("needs a parallel orchestrator, a
curated fixed-id yaml, a dedicated stream kind"). With one client there is no second
reader for a config that outlives its job, and since NICOS sends neither config nor
start, there is no spec surface left to align to either.Wrong as written now. "Remove
workflow_idfromJobCommand" argues the field is
unused and is the wrong granularity for a multi-client world. ADR 0006 then made it
NICOS's addressing mechanism, deliberately whole-workflow; acting on that section would
break NICOS reset. Relatedly, "per-source targeting (JobId) is removed" needs
re-deriving rather than adopting: stop is still per-JobId, andstop_workflow's
orphan sweep leans on it.Smaller drift: the "ConfigProcessor replacement" open question names a class that no
longer exists (CommandDispatcherreplaced it), and auto-start "based onJobSchedule"
rests on a field no producer insrcever sets.Carried forward as #1268: one command per worker instead of N per source. That
part stands on its own — it removes the ack accumulation without needing the
config/start split. Early validation of params at staging is the other idea worth
keeping; it does not need a message split either, so it can be filed if it turns out to
matter.- added 2 commits that reference this issue
on Sep 1, 2026
There is still some legacy
WorkflowStatuscode floating about, as "replies" toWorkflowConfigmessages (which trigger workflow starts). This are currently not handled by the frontend any more (which relies on periodicJobStatussince #440).https://gitlab.esss.lu.se/ecdc/interface-specs/-/blob/477da77510561f1a6afee874f7c8bfe907183723/nicos-beamlime.md#command-acknowledgement describes a mechanism that would be understood by NICOS. We should consider implementing it this way, then start using it on the frontend. Here is the content of that file (note BeamLime references there is ESSlivedata, i.e., this project) - also note this is not the latest document version, but some things are presented better here:
BeamLime <-> NICOS Interface
Introduction
BeamLime will be performing all the data reduction and NICOS will be responsible for displaying the data
and controlling the BeamLime services. All communication between BeamLime and NICOS will be done through Kafka.
BeamLime
BeamLime consists of multiple services each providing a different plot of reduced data.
The services are responsible for collecting data from the detectors and monitors, reducing it and sending it to NICOS
or other clients for display.
Each type of plot, even if coming from the same device, could be a different service.
NICOS
NICOS is a control system software that provides the user interface for the instruments at the ESS.
NICOS can display plots from different services and devices. In this case, NICOS will display the plots from BeamLime.
There is a central device on the NICOS server that will receive the data from BeamLime and distribute it to follower devices.
Configuration commands will be sent from the follower devices through the central device to the BeamLime server via Kafka.
Source Names
The source names should be unique per instrument and follow the format
device_name:signal_name.For example:
endcap_backwards_detector:accumulated_countscave_monitor_1:sliding_window_countsSince the
device_name:signal_namecan be parsed by NICOS, we can use it to send commands to the device.Command JSON
The command JSON is a dictionary with the following mandatory keys:
command: The command to be executed,config,start,stop,clearmessage_id: A unique identifier for the messagedevice: The device to be configured, for exampleendcap_backwards_detector:accumulated_counts,cave_monitor_1:sliding_window_counts.This should match the data source name.
Table of valid config keys
TODO: Add the correct keys that are expected by BeamLime. Below
sliding_windowupdate_everytime_of_arrival_binsExample command JSON
Start command:
{ "command": "start", "message_id": "1234", "device": "endcap_backwards_detector:accumulated_counts" }Stop command:
{ "command": "stop", "message_id": "1234", "device": "endcap_backwards_detector:accumulated_counts" }Clear command:
{ "command": "clear", "message_id": "1234", "device": "endcap_backwards_detector:accumulated_counts" }{ "command": "config", "message_id": "1234", "device": "endcap_backwards_detector:accumulated_counts", "config": { "update_every": 2000, "time_of_arrival_bins": 100 } }Command acknowledgement
The command acknowledgement JSON is a dictionary with the following mandatory keys:
message_id: The unique identifier of the messagedevice: The device that the command was sent to, for exampleendcap_backwards_detector:accumulated_countsresponse: The response to the command,ACK,ERROptional key, specific to the
ERRresponse:message: A string with the error messageHeartbeat message (x5f2)
Schema: https://github.com/ess-dmsc/streaming-data-types/blob/master/schemas/x5f2_status.fbs
Most parameters are self-explanatory, the
status_jsonshould contain a JSON string with the status of the device.The
status_jsonshould contain astatuskey with the status code and amessagekey with a message.We follow the NICOS status constants:
The
service_idshould be the same as thedevicein the command JSON so we can match the status to the device.Example code:
Data message (da00)
Schema: https://github.com/ess-dmsc/streaming-data-types/blob/master/schemas/da00_dataarray.fbs
Main plot data should have the name
signaland the axes list linked to the otherVariableaxes values.The
signalVariabledatashould be a 1D array or a 2D array.The
labelinsignalVariableshould specify the plot type, for examplehist-1d,hist-2d,scatter-2d.All the
Variablethat contain the axes values should be 1D arrays.sourceshould specify the name of the plot that will be selected within NICOSFor example
signalVariableof length N and an bin edge axisVariableof length N+1.signalVariableof shape (N, M) and two bin edge axesVariableof length N+1 and M+1.signalVariableof length N and twoaxisVariableof length N. (TBD WIP)Kafka topics
For each
[INSTRUMENT]we have the following topics:[INSTRUMENT]_beamlime_commands[INSTRUMENT]_beamlime_heartbeat[INSTRUMENT]_beamlime_responses[INSTRUMENT]_beamlime_data