pub enum ControlRequest {
Show 38 variants
Build {
session_id: SessionId,
dataflow: Descriptor,
git_sources: BTreeMap<NodeId, GitSource>,
prev_git_sources: BTreeMap<NodeId, GitSource>,
local_working_dir: Option<PathBuf>,
uv: bool,
},
WaitForBuild {
build_id: BuildId,
},
Start {
build_id: Option<BuildId>,
session_id: SessionId,
dataflow: Descriptor,
name: Option<String>,
local_working_dir: Option<PathBuf>,
uv: bool,
write_events_to: Option<PathBuf>,
},
WaitForSpawn {
dataflow_id: Uuid,
},
Reload {
dataflow_id: Uuid,
node_id: NodeId,
operator_id: Option<OperatorId>,
},
Check {
dataflow_uuid: Uuid,
},
Stop {
dataflow_uuid: Uuid,
grace_duration: Option<Duration>,
force: bool,
},
StopByName {
name: String,
grace_duration: Option<Duration>,
force: bool,
},
Restart {
dataflow_uuid: Uuid,
grace_duration: Option<Duration>,
force: bool,
},
RestartByName {
name: String,
grace_duration: Option<Duration>,
force: bool,
},
Logs {
uuid: Option<Uuid>,
name: Option<String>,
node: String,
tail: Option<usize>,
},
Destroy,
List,
Clean,
Info {
dataflow_uuid: Uuid,
},
DaemonConnected,
ConnectedMachines,
LogSubscribe {
dataflow_id: Uuid,
level: LevelFilter,
},
BuildLogSubscribe {
build_id: BuildId,
level: LevelFilter,
},
CliAndDefaultDaemonOnSameMachine,
GetNodeInfo,
TopicSubscribe {
dataflow_id: Uuid,
topics: Vec<(NodeId, DataId)>,
protocol_version: Option<u16>,
},
TopicUnsubscribe {
subscription_id: Uuid,
},
GetTraces,
GetTraceSpans {
trace_id: String,
},
RestartNode {
dataflow_id: Uuid,
node_id: NodeId,
grace_duration: Option<Duration>,
},
StopNode {
dataflow_id: Uuid,
node_id: NodeId,
grace_duration: Option<Duration>,
},
TopicPublish {
dataflow_id: Uuid,
node_id: NodeId,
output_id: DataId,
data_json: String,
},
GetParams {
dataflow_id: Uuid,
node_id: NodeId,
},
GetParam {
dataflow_id: Uuid,
node_id: NodeId,
key: String,
},
SetParam {
dataflow_id: Uuid,
node_id: NodeId,
key: String,
value: Value,
},
DeleteParam {
dataflow_id: Uuid,
node_id: NodeId,
key: String,
},
AddNode {
dataflow_id: Uuid,
node: Node,
},
RemoveNode {
dataflow_id: Uuid,
node_id: NodeId,
grace_duration: Option<Duration>,
},
ReplaceNode {
dataflow_id: Uuid,
node: Node,
grace_duration: Option<Duration>,
},
AddMapping {
dataflow_id: Uuid,
source_node: NodeId,
source_output: DataId,
target_node: NodeId,
target_input: DataId,
},
RemoveMapping {
dataflow_id: Uuid,
source_node: NodeId,
source_output: DataId,
target_node: NodeId,
target_input: DataId,
},
Hello {
dora_version: Version,
},
}Variants§
Build
Fields
dataflow: DescriptorWaitForBuild
Start
Fields
dataflow: DescriptorWaitForSpawn
Reload
Check
Stop
StopByName
Restart
RestartByName
Logs
Destroy
List
Clean
Remove fully-completed dataflows from the coordinator’s state.
A dataflow is considered fully completed when no daemon is still
running it (i.e. it’s no longer in running_dataflows). Multi-daemon
dataflows that are still finishing — where some daemons have reported
results but others haven’t — are intentionally skipped so their final
status is computed correctly when the last daemon completes.
Candidates are enumerated from BOTH the in-memory
dataflow_results map AND the persisted store
(Succeeded / Failed records). This lets a restarted
coordinator still reap historical rows that exist only on
disk — the recovery loop intentionally does not reload them
into memory, so without the persisted-store pass they would
otherwise sit in redb forever and never become reachable for
dora clean. If the persisted-store enumeration itself
errors, the entire request fails with
ControlRequestReply::Error and no in-memory state is
mutated — degrading silently to in-memory-only would let the
CLI claim “nothing to clean” while historical rows are still
sitting on disk.
For each cleaned dataflow the coordinator removes its persisted
record first and only then mutates in-memory state, so the reply
reflects what was actually persisted. The response is
ControlRequestReply::CleanResult carrying two separate
lists: cleaned for dataflows whose redb row (and every
dora param row owned by it — the persisted-store delete
cascades) is gone, and failed for dataflows whose
persisted-store delete errored. In-memory entries for failed
candidates are preserved so a later dora clean can retry;
they show up in failed, not cleaned, so the CLI can tell
“nothing eligible” apart from “all candidates failed to
clean”. Logs and archived descriptors for successfully cleaned
dataflows are no longer available afterward. Cached build
results (finished_builds) are intentionally not touched —
clearing them would break concurrent dora build calls with
“unknown build id” errors.
Info
DaemonConnected
ConnectedMachines
LogSubscribe
BuildLogSubscribe
CliAndDefaultDaemonOnSameMachine
GetNodeInfo
TopicSubscribe
Fields
protocol_version: Option<u16>Binary-frame encoding the client speaks, see
TOPIC_DATA_PROTOCOL_VERSION.
None means the client predates the handshake and therefore speaks
the bincode encoding; the coordinator rejects it rather than let it
misparse postcard frames.
TopicUnsubscribe
GetTraces
GetTraceSpans
RestartNode
Restart a specific node without stopping the entire dataflow.
StopNode
Stop a specific node without stopping the entire dataflow.
TopicPublish
Publish a message to a topic (for debugging/testing).
The coordinator serializes the JSON data into Arrow format and publishes it to Zenoh on the appropriate topic key.
Fields
GetParams
List runtime parameters for a node.
GetParam
Get a single runtime parameter value.
SetParam
Set a runtime parameter on a node.
DeleteParam
Delete a runtime parameter from a node.
AddNode
Add a node to a running dataflow.
RemoveNode
Remove a node from a running dataflow.
ReplaceNode
Atomically replace a running node with a new definition under the same id (dora-rs/dora#2927). The replacement must keep the node’s edges (same input mappings, outputs covering every mapped output); a spawn failure leaves the current incarnation running.
AddMapping
Add a mapping (connection) between two nodes in a running dataflow.
Fields
RemoveMapping
Remove a mapping (connection) between two nodes in a running dataflow.
Fields
Hello
Protocol version handshake. Sent by the CLI as its first request after connecting so the coordinator can reject version-mismatched clients before they exchange incompatible messages (dora-rs/adora#151).
The coordinator replies with either
ControlRequestReply::HelloOk carrying its own crate version,
or ControlRequestReply::Error with a human-readable mismatch
message.
Implementations§
Trait Implementations§
Source§impl Clone for ControlRequest
impl Clone for ControlRequest
Source§fn clone(&self) -> ControlRequest
fn clone(&self) -> ControlRequest
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read more