Skip to main content

ControlRequest

Enum ControlRequest 

Source
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

§session_id: SessionId
§dataflow: Descriptor
§git_sources: BTreeMap<NodeId, GitSource>
§prev_git_sources: BTreeMap<NodeId, GitSource>
§local_working_dir: Option<PathBuf>

Allows overwriting the base working dir when CLI and daemon are running on the same machine.

Must not be used for multi-machine dataflows.

Note that nodes with git sources still use a subdirectory of the base working dir.

§uv: bool
§

WaitForBuild

Fields

§build_id: BuildId
§

Start

Fields

§build_id: Option<BuildId>
§session_id: SessionId
§dataflow: Descriptor
§local_working_dir: Option<PathBuf>

Allows overwriting the base working dir when CLI and daemon are running on the same machine.

Must not be used for multi-machine dataflows.

Note that nodes with git sources still use a subdirectory of the base working dir.

§uv: bool
§write_events_to: Option<PathBuf>
§

WaitForSpawn

Fields

§dataflow_id: Uuid
§

Reload

Fields

§dataflow_id: Uuid
§node_id: NodeId
§operator_id: Option<OperatorId>
§

Check

Fields

§dataflow_uuid: Uuid
§

Stop

Fields

§dataflow_uuid: Uuid
§grace_duration: Option<Duration>
§force: bool
§

StopByName

Fields

§name: String
§grace_duration: Option<Duration>
§force: bool
§

Restart

Fields

§dataflow_uuid: Uuid
§grace_duration: Option<Duration>
§force: bool
§

RestartByName

Fields

§name: String
§grace_duration: Option<Duration>
§force: bool
§

Logs

Fields

§uuid: Option<Uuid>
§node: String
§

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

Fields

§dataflow_uuid: Uuid
§

DaemonConnected

§

ConnectedMachines

§

LogSubscribe

Fields

§dataflow_id: Uuid
§

BuildLogSubscribe

Fields

§build_id: BuildId
§

CliAndDefaultDaemonOnSameMachine

§

GetNodeInfo

§

TopicSubscribe

Fields

§dataflow_id: Uuid
§topics: Vec<(NodeId, DataId)>
§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

Fields

§subscription_id: Uuid
§

GetTraces

§

GetTraceSpans

Fields

§trace_id: String
§

RestartNode

Restart a specific node without stopping the entire dataflow.

Fields

§dataflow_id: Uuid
§node_id: NodeId
§grace_duration: Option<Duration>
§

StopNode

Stop a specific node without stopping the entire dataflow.

Fields

§dataflow_id: Uuid
§node_id: NodeId
§grace_duration: Option<Duration>
§

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

§dataflow_id: Uuid
§node_id: NodeId
§output_id: DataId
§data_json: String

JSON data to publish (will be converted to Arrow UInt8 array)

§

GetParams

List runtime parameters for a node.

Fields

§dataflow_id: Uuid
§node_id: NodeId
§

GetParam

Get a single runtime parameter value.

Fields

§dataflow_id: Uuid
§node_id: NodeId
§

SetParam

Set a runtime parameter on a node.

Fields

§dataflow_id: Uuid
§node_id: NodeId
§value: Value
§

DeleteParam

Delete a runtime parameter from a node.

Fields

§dataflow_id: Uuid
§node_id: NodeId
§

AddNode

Add a node to a running dataflow.

Fields

§dataflow_id: Uuid
§node: Node
§

RemoveNode

Remove a node from a running dataflow.

Fields

§dataflow_id: Uuid
§node_id: NodeId
§grace_duration: Option<Duration>
§

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.

Fields

§dataflow_id: Uuid
§node: Node
§grace_duration: Option<Duration>
§

AddMapping

Add a mapping (connection) between two nodes in a running dataflow.

Fields

§dataflow_id: Uuid
§source_node: NodeId
§source_output: DataId
§target_node: NodeId
§target_input: DataId
§

RemoveMapping

Remove a mapping (connection) between two nodes in a running dataflow.

Fields

§dataflow_id: Uuid
§source_node: NodeId
§source_output: DataId
§target_node: NodeId
§target_input: DataId
§

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.

Fields

§dora_version: Version

Implementations§

Source§

impl ControlRequest

Source

pub fn hello() -> Self

Build a Hello request stamped with the current crate version of dora-message (the wire-protocol version).

Trait Implementations§

Source§

impl Clone for ControlRequest

Source§

fn clone(&self) -> ControlRequest

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for ControlRequest

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl<'de> Deserialize<'de> for ControlRequest

Source§

fn deserialize<__D>(__deserializer: __D) -> Result<Self, __D::Error>
where __D: Deserializer<'de>,

Deserialize this value from the given Serde deserializer. Read more
Source§

impl Serialize for ControlRequest

Source§

fn serialize<__S>(&self, __serializer: __S) -> Result<__S::Ok, __S::Error>
where __S: Serializer,

Serialize this value into the given Serde serializer. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Allocation for T
where T: RefUnwindSafe + Send + Sync,

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> DeserializeOwned for T
where T: for<'de> Deserialize<'de>,

Source§

impl<T> DynClone for T
where T: Clone,

Source§

fn __clone_box(&self, _: Private) -> *mut ()

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V