use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use crate::{ActivityId, WorkflowId};
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Debug, PartialEq, Eq)]
pub struct ClusterEventMeta {
pub cluster_seq: u64,
pub observed_at: DateTime<Utc>,
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Copy, Debug, PartialEq, Eq)]
#[serde(tag = "transport")]
pub enum WorkerTransport {
Grpc,
Liminal,
}
impl WorkerTransport {
#[must_use]
pub const fn name(self) -> &'static str {
match self {
Self::Grpc => "grpc",
Self::Liminal => "liminal",
}
}
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Copy, Debug, PartialEq, Eq)]
#[serde(tag = "reason")]
pub enum WorkerDeathReason {
Disconnect,
Timeout,
Deregistered,
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Copy, Debug, PartialEq, Eq)]
#[serde(rename_all = "kebab-case")]
pub enum DesiredState {
Running,
Stopped,
}
impl DesiredState {
#[must_use]
pub const fn token(self) -> &'static str {
match self {
Self::Running => "running",
Self::Stopped => "stopped",
}
}
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Copy, Debug, PartialEq, Eq)]
#[serde(rename_all = "kebab-case")]
pub enum PutOutcome {
Created,
Replaced,
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Copy, Debug, PartialEq, Eq)]
#[serde(rename_all = "kebab-case")]
pub enum DeploymentAssociation {
Known,
Absent,
Unchecked,
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Debug, PartialEq, Eq)]
#[serde(tag = "type")]
pub enum ClusterEvent {
PeerAdded {
meta: ClusterEventMeta,
peer_name: String,
forward_addr: Option<String>,
},
PeerConnected {
meta: ClusterEventMeta,
peer_name: String,
forward_addr: Option<String>,
},
PeerDisconnected {
meta: ClusterEventMeta,
peer_name: String,
consecutive_down: u32,
confirmed: bool,
},
ShardAdopted {
meta: ClusterEventMeta,
shards: Vec<usize>,
from_peer: String,
adopted_by: String,
},
ShardAdoptionFailed {
meta: ClusterEventMeta,
shards: Vec<usize>,
from_peer: String,
error: String,
},
ShardAdoptionSkipped {
meta: ClusterEventMeta,
shards: Vec<usize>,
from_peer: String,
held_by: String,
},
WorkerConnected {
meta: ClusterEventMeta,
worker_id: String,
namespaces: Vec<String>,
task_queue: String,
transport: WorkerTransport,
node: Option<String>,
deployment: Option<String>,
deployment_association: Option<DeploymentAssociation>,
},
WorkerDisconnected {
meta: ClusterEventMeta,
worker_id: String,
namespaces: Vec<String>,
reason: WorkerDeathReason,
},
DispatchParked {
meta: ClusterEventMeta,
namespace: String,
task_queue: String,
activity_type: String,
node: Option<String>,
reason: String,
policy: String,
workflow_id: WorkflowId,
activity_id: ActivityId,
waited_ms: u64,
workers_in_pool: usize,
workers_serving_activity: usize,
compatible_workers: usize,
last_compatible_poller_age_ms: Option<u64>,
},
SupervisorStarted {
meta: ClusterEventMeta,
node: String,
},
SupervisorStopped {
meta: ClusterEventMeta,
node: String,
},
NamespaceCreated {
meta: ClusterEventMeta,
name: String,
created_at: DateTime<Utc>,
origin: String,
},
NamespacePlacementChanged {
meta: ClusterEventMeta,
name: String,
placement: NamespacePlacementWire,
},
NamespaceQuotaState {
meta: ClusterEventMeta,
namespace: String,
in_flight: u64,
ceiling: u32,
},
WorkerDeploymentPut {
meta: ClusterEventMeta,
name: String,
outcome: PutOutcome,
desired_state: DesiredState,
binary_version: String,
binary_content_hash: String,
},
WorkerDeploymentDesiredStateChanged {
meta: ClusterEventMeta,
name: String,
desired_state: DesiredState,
},
WorkerDeploymentDeleted {
meta: ClusterEventMeta,
name: String,
},
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Debug, PartialEq, Eq)]
pub struct NamespacePlacementWire {
pub kind: String,
pub nodes: Vec<String>,
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Debug, PartialEq, Eq)]
pub struct ClusterPeer {
pub peer_name: String,
pub forward_addr: Option<String>,
pub connected: bool,
pub consecutive_down: u32,
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Debug, PartialEq, Eq)]
pub struct ClusterShard {
pub shard: usize,
pub owner: String,
pub epoch: u64,
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Debug, PartialEq, Eq)]
pub struct ClusterWorker {
pub worker_id: String,
pub namespaces: Vec<String>,
pub task_queue: String,
pub transport: WorkerTransport,
pub node: Option<String>,
pub deployment: Option<String>,
pub deployment_association: Option<DeploymentAssociation>,
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Debug, PartialEq, Eq)]
pub struct ClusterDeployment {
pub name: String,
pub desired_state: DesiredState,
pub binary_version: String,
pub binary_content_hash: String,
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Debug, PartialEq, Eq)]
pub struct ClusterSnapshot {
pub node: String,
pub as_of_seq: u64,
pub peers: Vec<ClusterPeer>,
pub shards: Vec<ClusterShard>,
pub workers: Vec<ClusterWorker>,
pub deployments: Vec<ClusterDeployment>,
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Debug, PartialEq, Eq)]
#[serde(tag = "command")]
pub enum ClusterCommand {
RequestClusterSnapshot {},
CancelWorkflow {
namespace: String,
workflow_id: String,
},
ReopenWorkflow {
namespace: String,
workflow_id: String,
},
RedriveOutboxRow {
namespace: String,
workflow_id: String,
ordinal: u64,
},
DrainNode {
node: String,
},
PlannedHandoff {
shard: usize,
target_node: String,
},
ChaosKillNode {
node: String,
},
}
#[derive(Serialize, Deserialize, ts_rs::TS, Clone, Debug, PartialEq, Eq)]
#[serde(tag = "kind")]
pub enum ClusterStreamError {
ClusterLagged {
skipped: u64,
},
}