Skip to content

Cleanup/fix reply mechanism to request from frontend #445

Description

@SimonHeybrock

There is still some legacy WorkflowStatus code floating about, as "replies" to WorkflowConfig messages (which trigger workflow starts). This are currently not handled by the frontend any more (which relies on periodic JobStatus since #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_counts
  • cave_monitor_1:sliding_window_counts

Since the device_name:signal_name can 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, clear
  • message_id: A unique identifier for the message
  • device: The device to be configured, for example endcap_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

Key Type Description
sliding_window int The sliding window size in milliseconds
update_every int The time in milliseconds to update the devices plots
time_of_arrival_bins int The number of bins for the time of arrival histogram

Example 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 message
  • device: The device that the command was sent to, for example endcap_backwards_detector:accumulated_counts
  • response: The response to the command, ACK, ERR

Optional key, specific to the ERR response:

  • message: A string with the error message

Heartbeat message (x5f2)

Schema: https://github.com/ess-dmsc/streaming-data-types/blob/master/schemas/x5f2_status.fbs

Most parameters are self-explanatory, the status_json should contain a JSON string with the status of the device.

The status_json should contain a status key with the status code and a message key with a message.

We follow the NICOS status constants:

OK = 200
WARN = 210
BUSY = 220
NOTREACHED = 230
DISABLED = 235
ERROR = 240
UNKNOWN = 999

The service_id should be the same as the device in the command JSON so we can match the status to the device.

Example code:

msg = serialise_x5f2(
    software_name = "beamlime",
    software_version = "1.0.1",
    service_id = "endcap_backwards_detector:accumulated_counts",
    host_name = "beamlimeserver.esss.dk",
    process_id = 8529,
    update_interval = 2000,
    status_json =  "{'status': 200, 'message': 'Idle'}",
)

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 signal and the axes list linked to the other Variable axes values.

The signal Variable data should be a 1D array or a 2D array.

The label in signal Variable should specify the plot type, for example hist-1d, hist-2d, scatter-2d.

All the Variable that contain the axes values should be 1D arrays.

source should specify the name of the plot that will be selected within NICOS

For example

  • A 1D histogram would have a signal Variable of length N and an bin edge axis Variable of length N+1.
  • a 2D histogram would have a signal Variable of shape (N, M) and two bin edge axes Variable of length N+1 and M+1.
  • (TBD WIP) A 2D scatter plot would have a signal Variable of length N and two axis Variable of length N. (TBD WIP)

Kafka topics

For each [INSTRUMENT] we have the following topics:

  • name: [INSTRUMENT]_beamlime_commands
    • description: Commands from NICOS to be executed by the BeamLime server
  • name: [INSTRUMENT]_beamlime_heartbeat
    • description: Heartbeat messages from the BeamLime server to NICOS
  • name: [INSTRUMENT]_beamlime_responses
    • description: Responses to the commands sent from NICOS to the BeamLime server
  • name: [INSTRUMENT]_beamlime_data
    • description: Data messages from the BeamLime server to be displayed in NICOS

Activity

  1. 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
  2. SimonHeybrock commented on Oct 2, 2025

    @SimonHeybrock
    MemberAuthor

    Partial plan forming in my head:

    WorkflowConfig contains 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 (based JobSchedule, which we can include in the config part of the message)?

    "Device" naming (e.g., for a clear command) 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?

  3. added 2 commits that reference this issue on Nov 19, 2025
  4. SimonHeybrock commented on Jan 5, 2026

    @SimonHeybrock
    MemberAuthor

    Design Notes: Splitting WorkflowConfig into Config + Start

    During work on the command acknowledgement feature, we analyzed the relationship between message_id and job_number in WorkflowConfig. This comment documents the findings for future reference when implementing the config/start split.

    Current State

    WorkflowConfig currently 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_id Command/response ACK correlation Transient (until ACK received) Frontend
    job_number Job 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_id serves double duty:

    1. ACK correlation (standard use)
    2. 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

    1. Independent jobs with shared config: Frontend and NICOS can both send WorkflowStart referencing the same config_ref=abc but with different job_numbers

    2. Cleaner semantics: Configuration is stable, jobs are instances of that config

    3. Aligns with NICOS interface: Matches the config / start command 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_jobs pattern in JobOrchestrator already provides the "staging area" concept

    Code comments have been added to WorkflowConfig and JobOrchestrator referencing this design.

  5. SimonHeybrock commented on Feb 26, 2026

    @SimonHeybrock
    MemberAuthor

    Config/Start Split Design

    Context

    WorkflowConfig currently 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.UUID type identifies a configuration, analogous to how
    JobNumber identifies a job. This replaces the earlier idea of reusing
    WorkflowConfig.message_id as a "config handle".

    Reasons for a dedicated type:

    • Lifetime mismatch: message_id is 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. With ConfigNumber, 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 JobNumber pattern in the codebase.

    WorkflowConfig is a single atomic message

    A WorkflowConfig message contains the full config set — all per-source configs in one
    message, identified by a single ConfigNumber. 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 its jobs dict.

    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 start command currently has no config reference — it implicitly uses the
    latest config for a device. We will extend the spec to include ConfigNumber. As a
    fallback, omitting config_number on 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 same JobCommand model,
    aligning with the NICOS spec where start, stop, config, clear are variants of one
    command structure.

    The key enabler is targeting by JobNumber instead of JobId. Currently, stop and
    reset send N individual commands (one per JobId, i.e., per source). But since all jobs
    in a set share the same JobNumber, a single command targeting the JobNumber is
    sufficient. This makes start structurally 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 start carries 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_id field on JobCommand is unused in practice — every call site in the
    codebase targets by job_id or broadcasts to all jobs. No frontend or test code ever
    sets workflow_id on a JobCommand.

    More importantly, workflow_id targeting is the wrong granularity for a multi-client world.
    If NICOS sends stop(workflow_id=W), it would kill the dashboard's jobs for that workflow
    too. JobNumber is the correct scope: each client creates its own JobNumbers, so
    stop(job_number=J) affects only that client's jobs.

    The backend dispatch simplifies from three branches to two:

    • job_number specified → 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 (new ConfigNumber). 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 its ConfigNumber.

    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 C1
    

    Per-worker ACK per command, with per-source error detail

    Each command (config, start, stop, reset) produces one CommandAcknowledgement per 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 (ERR with 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. ACK if
      all succeed, ERR with source_errors if 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. ACK if all jobs affected by the JobNumber were
      successfully stopped/reset, ERR with 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 current PendingCommandTracker exists precisely because N
      commands produce N ACKs that must be accumulated — tracking expected_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 from PendingCommandTracker'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 (JobStatus with state=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. CommandAcknowledgement gains worker_id and sources_handled fields (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_id and sources_handled are None, 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 in source_errors is 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 ConfigProcessor deduplicates 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. ConfigProcessor will 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 their ConfigNumbers. 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?

  6. added
    area:wireKafka topics, message schemas, command/reply contracts, incl. toward NICOS/ECDC
    on Aug 11, 2026
  7. SimonHeybrock commented on Aug 31, 2026

    @SimonHeybrock
    MemberAuthor

    The original ask here is done, and the design that grew on top of it was overtaken by
    ADR 0006.

    Done. The legacy WorkflowStatus reply type is gone from the tree.
    CommandAcknowledgement implements the command-acknowledgement shape this issue linked,
    and the dashboard consumes it through PendingCommandTracker. ADR 0006 went further
    than the plan in the first comment: JobCommand is the command NICOS produces, and
    WorkflowId accepts its instrument/name/version string 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 send WorkflowStart referencing the same
    config_ref but with different job_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 is reset keyed by WorkflowId, 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_id from JobCommand" 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, and stop_workflow's
    orphan sweep leans on it.

    Smaller drift: the "ConfigProcessor replacement" open question names a class that no
    longer exists (CommandDispatcher replaced it), and auto-start "based on JobSchedule"
    rests on a field no producer in src ever 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.

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

Metadata

Metadata

Assignees

Labels

area:wireKafka topics, message schemas, command/reply contracts, incl. toward NICOS/ECDC

Type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions