use std::{
collections::{BTreeMap, BTreeSet},
path::PathBuf,
time::Duration,
};
use crate::{
BuildId, DataflowId, SessionId,
common::{DaemonId, GitSource},
descriptor::{Descriptor, ResolvedNode},
id::{DataId, NodeId, OperatorId},
};
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct StateCatchUpEntry {
pub sequence: u64,
pub operation: StateCatchUpOperation,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum StateCatchUpOperation {
SetParam {
node_id: NodeId,
key: String,
value: serde_json::Value,
},
DeleteParam {
node_id: NodeId,
key: String,
},
}
pub use crate::common::Timestamped;
#[derive(Debug, serde::Serialize, serde::Deserialize)]
pub enum RegisterResult {
Ok {
daemon_id: DaemonId,
},
Err(String),
}
impl RegisterResult {
pub fn to_result(self) -> eyre::Result<DaemonId> {
match self {
RegisterResult::Ok { daemon_id } => Ok(daemon_id),
RegisterResult::Err(err) => Err(eyre::eyre!(err)),
}
}
}
#[allow(clippy::large_enum_variant)]
#[derive(Debug, serde::Deserialize, serde::Serialize)]
pub enum DaemonCoordinatorEvent {
Build(BuildDataflowNodes),
Spawn(SpawnDataflowNodes),
AllNodesReady {
dataflow_id: DataflowId,
exited_before_subscribe: Vec<NodeId>,
},
StopDataflow {
dataflow_id: DataflowId,
grace_duration: Option<Duration>,
#[serde(default)]
force: bool,
},
ReloadDataflow {
dataflow_id: DataflowId,
node_id: NodeId,
operator_id: Option<OperatorId>,
},
Logs {
dataflow_id: DataflowId,
node_id: NodeId,
tail: Option<usize>,
},
RestartNode {
dataflow_id: DataflowId,
node_id: NodeId,
grace_duration: Option<Duration>,
},
StopNode {
dataflow_id: DataflowId,
node_id: NodeId,
grace_duration: Option<Duration>,
},
SetParam {
dataflow_id: DataflowId,
node_id: NodeId,
key: String,
value: serde_json::Value,
},
DeleteParam {
dataflow_id: DataflowId,
node_id: NodeId,
key: String,
},
Destroy,
Heartbeat,
PeerDaemonDisconnected {
daemon_id: DaemonId,
},
AddNode {
dataflow_id: DataflowId,
node: crate::descriptor::ResolvedNode,
uv: bool,
},
RemoveNode {
dataflow_id: DataflowId,
node_id: NodeId,
grace_duration: Option<Duration>,
},
AddMapping {
dataflow_id: DataflowId,
source_node: NodeId,
source_output: DataId,
target_node: NodeId,
target_input: DataId,
},
RemoveMapping {
dataflow_id: DataflowId,
source_node: NodeId,
source_output: DataId,
target_node: NodeId,
target_input: DataId,
},
StartTopicDebugStream {
dataflow_id: DataflowId,
outputs: Vec<(NodeId, DataId)>,
subscription_id: uuid::Uuid,
},
StopTopicDebugStream {
dataflow_id: DataflowId,
subscription_id: uuid::Uuid,
},
StateCatchUp {
dataflow_id: DataflowId,
entries: Vec<StateCatchUpEntry>,
},
}
#[derive(Debug, serde::Deserialize, serde::Serialize)]
pub struct BuildDataflowNodes {
pub build_id: BuildId,
pub session_id: SessionId,
pub local_working_dir: Option<PathBuf>,
pub git_sources: BTreeMap<NodeId, GitSource>,
pub prev_git_sources: BTreeMap<NodeId, GitSource>,
pub dataflow_descriptor: Descriptor,
pub nodes_on_machine: BTreeSet<NodeId>,
pub uv: bool,
}
#[derive(Debug, serde::Deserialize, serde::Serialize)]
pub struct SpawnDataflowNodes {
pub build_id: Option<BuildId>,
pub session_id: SessionId,
pub dataflow_id: DataflowId,
pub local_working_dir: Option<PathBuf>,
pub nodes: BTreeMap<NodeId, ResolvedNode>,
pub dataflow_descriptor: Descriptor,
pub spawn_nodes: BTreeSet<NodeId>,
pub uv: bool,
pub write_events_to: Option<PathBuf>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub artifact_base_url: Option<String>,
}