Skip to main content

stasis/domain/runtime/
cluster_node.rs

1use chrono::{DateTime, Duration, Utc};
2
3#[derive(Clone, Copy, Debug, Eq, PartialEq)]
4pub enum ClusterNodeRole {
5    Coordinator,
6    Scheduler,
7    Worker,
8}
9
10#[derive(Clone, Copy, Debug, Eq, PartialEq)]
11pub enum QueueOwnershipMode {
12    MultiOwner,
13    SingleOwner,
14}
15
16#[derive(Clone, Copy, Debug, Eq, PartialEq)]
17pub enum ClusterNodeHealth {
18    Healthy,
19    Degraded,
20    Offline,
21}
22
23#[derive(Clone, Debug, Eq, PartialEq)]
24pub struct ClusterNode {
25    pub node_id: String,
26    pub role: ClusterNodeRole,
27    pub region: String,
28    pub queue_ownership: Vec<String>,
29    pub capability_tags: Vec<String>,
30    pub heartbeat_at: DateTime<Utc>,
31    pub lease_expires_at: DateTime<Utc>,
32    pub metadata: Option<String>,
33    pub created_at: DateTime<Utc>,
34    pub updated_at: DateTime<Utc>,
35}
36
37#[derive(Clone, Debug, Eq, PartialEq)]
38pub struct NewClusterNode {
39    pub node_id: String,
40    pub role: ClusterNodeRole,
41    pub region: String,
42    pub queue_ownership: Vec<String>,
43    pub capability_tags: Vec<String>,
44    pub heartbeat_at: DateTime<Utc>,
45    pub lease_ttl_seconds: i64,
46    pub metadata: Option<String>,
47}
48
49impl NewClusterNode {
50    pub fn into_record(self) -> ClusterNode {
51        let now = self.heartbeat_at;
52        let lease_ttl_seconds = self.lease_ttl_seconds.max(1);
53        ClusterNode {
54            node_id: self.node_id,
55            role: self.role,
56            region: self.region,
57            queue_ownership: self.queue_ownership,
58            capability_tags: self.capability_tags,
59            heartbeat_at: now,
60            lease_expires_at: now + Duration::seconds(lease_ttl_seconds),
61            metadata: self.metadata,
62            created_at: now,
63            updated_at: now,
64        }
65    }
66}
67
68#[derive(Clone, Debug, Eq, PartialEq)]
69pub struct ClusterNodeHeartbeat {
70    pub node_id: String,
71    pub heartbeat_at: DateTime<Utc>,
72    pub lease_ttl_seconds: i64,
73    pub queue_ownership: Option<Vec<String>>,
74    pub capability_tags: Option<Vec<String>>,
75    pub metadata: Option<String>,
76}
77
78#[derive(Clone, Debug, Eq, PartialEq)]
79pub struct ClusterNodeHealthSnapshot {
80    pub node: ClusterNode,
81    pub health: ClusterNodeHealth,
82}
83
84#[derive(Clone, Debug, Eq, PartialEq)]
85pub struct ClusterForwardCommand {
86    pub target_region: String,
87    pub command_name: String,
88    pub payload: String,
89    pub correlation_id: Option<String>,
90    pub issued_at: DateTime<Utc>,
91}
92
93#[derive(Clone, Debug, Eq, PartialEq)]
94pub struct ClusterForwardOutcome {
95    pub target_region: String,
96    pub command_name: String,
97    pub correlation_id: Option<String>,
98    pub accepted: bool,
99    pub attempts: u32,
100    pub error: Option<String>,
101    pub completed_at: DateTime<Utc>,
102}
103
104#[derive(Clone, Debug, Eq, PartialEq)]
105pub enum ClusterControlEvent {
106    ExpiredNodesPruned {
107        pruned_count: u64,
108        occurred_at: DateTime<Utc>,
109    },
110}