Skip to main content

stasis/application/use_cases/
manage_cluster_nodes.rs

1use std::collections::BTreeMap;
2
3use chrono::{DateTime, Duration, Utc};
4
5use crate::application::dto::{
6    ClusterForwardOutcomeRow, ClusterNodeHealthRow, ForwardClusterCommandRequest,
7    ForwardClusterCommandResponse, HeartbeatClusterNodeRequest, InitiateCoordinatorFailoverRequest,
8    InitiateCoordinatorFailoverResponse, InitiateCoordinatorHandoffRequest,
9    InitiateCoordinatorHandoffResponse, ListClusterForwardOutcomesRequest,
10    ListClusterNodeHealthRequest, ListQueueOwnershipHealthRequest, PruneExpiredClusterNodesRequest,
11    QueueOwnershipHealthRow, RebalanceQueueOwnershipRequest, RebalanceQueueOwnershipResponse,
12    RegisterClusterNodeRequest, RunClusterHeartbeatSweepRequest, RunClusterHeartbeatSweepResponse,
13};
14use crate::domain::errors::{Result, StasisError};
15use crate::domain::runtime::cluster_node::{
16    ClusterControlEvent, ClusterForwardCommand, ClusterNode, ClusterNodeHealth,
17    ClusterNodeHealthSnapshot, ClusterNodeHeartbeat, NewClusterNode, QueueOwnershipMode,
18};
19use crate::ports::outbound::runtime::cluster_command_forwarder::ClusterCommandForwarder;
20use crate::ports::outbound::runtime::cluster_control_event_sink::ClusterControlEventSink;
21use crate::ports::outbound::runtime::cluster_forward_outcome_store::ClusterForwardOutcomeStore;
22use crate::ports::outbound::runtime::cluster_node_store::ClusterNodeStore;
23use serde_json::json;
24
25#[derive(Clone)]
26pub struct RegisterClusterNode<S>
27where
28    S: ClusterNodeStore,
29{
30    store: S,
31}
32
33impl<S> RegisterClusterNode<S>
34where
35    S: ClusterNodeStore,
36{
37    pub fn new(store: S) -> Self {
38        Self { store }
39    }
40
41    pub async fn execute(&self, request: RegisterClusterNodeRequest) -> Result<ClusterNode> {
42        if request.node_id.trim().is_empty() {
43            return Err(StasisError::PortFailure(
44                "node_id must not be empty".to_string(),
45            ));
46        }
47        if request.region.trim().is_empty() {
48            return Err(StasisError::PortFailure(
49                "region must not be empty".to_string(),
50            ));
51        }
52
53        let mode = request
54            .queue_ownership_mode
55            .unwrap_or(QueueOwnershipMode::MultiOwner);
56        validate_queue_ownership(
57            &self.store,
58            &request.node_id,
59            &request.queue_ownership,
60            mode,
61            request.heartbeat_at,
62        )
63        .await?;
64
65        self.store
66            .register(NewClusterNode {
67                node_id: request.node_id,
68                role: request.role,
69                region: request.region,
70                queue_ownership: request.queue_ownership,
71                capability_tags: request.capability_tags,
72                heartbeat_at: request.heartbeat_at,
73                lease_ttl_seconds: request.lease_ttl_seconds,
74                metadata: request.metadata,
75            })
76            .await
77    }
78}
79
80#[derive(Clone)]
81pub struct HeartbeatClusterNode<S>
82where
83    S: ClusterNodeStore,
84{
85    store: S,
86}
87
88impl<S> HeartbeatClusterNode<S>
89where
90    S: ClusterNodeStore,
91{
92    pub fn new(store: S) -> Self {
93        Self { store }
94    }
95
96    pub async fn execute(&self, request: HeartbeatClusterNodeRequest) -> Result<ClusterNode> {
97        if let Some(queue_ownership) = &request.queue_ownership {
98            let mode = request
99                .queue_ownership_mode
100                .unwrap_or(QueueOwnershipMode::MultiOwner);
101            validate_queue_ownership(
102                &self.store,
103                &request.node_id,
104                queue_ownership,
105                mode,
106                request.heartbeat_at,
107            )
108            .await?;
109        }
110
111        let updated = self
112            .store
113            .heartbeat(ClusterNodeHeartbeat {
114                node_id: request.node_id.clone(),
115                heartbeat_at: request.heartbeat_at,
116                lease_ttl_seconds: request.lease_ttl_seconds,
117                queue_ownership: request.queue_ownership,
118                capability_tags: request.capability_tags,
119                metadata: request.metadata,
120            })
121            .await?;
122
123        updated.ok_or_else(|| {
124            StasisError::PortFailure(format!("cluster node not found: {}", request.node_id))
125        })
126    }
127}
128
129#[derive(Clone)]
130pub struct RunClusterHeartbeatSweep<S, E>
131where
132    S: ClusterNodeStore,
133    E: ClusterControlEventSink,
134{
135    store: S,
136    event_sink: E,
137}
138
139#[derive(Clone)]
140pub struct ForwardClusterControlCommand<F>
141where
142    F: ClusterCommandForwarder,
143{
144    forwarder: F,
145}
146
147#[derive(Clone)]
148pub struct InitiateCoordinatorHandoff<F>
149where
150    F: ClusterCommandForwarder,
151{
152    forwarder: F,
153}
154
155#[derive(Clone)]
156pub struct InitiateCoordinatorFailover<F>
157where
158    F: ClusterCommandForwarder,
159{
160    forwarder: F,
161}
162
163#[derive(Clone)]
164pub struct RebalanceQueueOwnership<F>
165where
166    F: ClusterCommandForwarder,
167{
168    forwarder: F,
169}
170
171impl<F> InitiateCoordinatorHandoff<F>
172where
173    F: ClusterCommandForwarder,
174{
175    pub fn new(forwarder: F) -> Self {
176        Self { forwarder }
177    }
178
179    pub async fn execute(
180        &self,
181        request: InitiateCoordinatorHandoffRequest,
182    ) -> Result<InitiateCoordinatorHandoffResponse> {
183        if request.target_region.trim().is_empty() {
184            return Err(StasisError::PortFailure(
185                "target_region must not be empty".to_string(),
186            ));
187        }
188        if request.coordinator_node_id.trim().is_empty() {
189            return Err(StasisError::PortFailure(
190                "coordinator_node_id must not be empty".to_string(),
191            ));
192        }
193
194        let payload = json!({
195            "coordinator_node_id": request.coordinator_node_id,
196            "queue_scope": request.queue_scope,
197            "reason": request.reason,
198        })
199        .to_string();
200
201        let command_name = "coordinator.handoff".to_string();
202        let accepted = self
203            .forwarder
204            .forward(ClusterForwardCommand {
205                target_region: request.target_region,
206                command_name: command_name.clone(),
207                payload,
208                correlation_id: request.correlation_id,
209                issued_at: request.issued_at,
210            })
211            .await?;
212
213        Ok(InitiateCoordinatorHandoffResponse {
214            accepted,
215            command_name,
216        })
217    }
218}
219
220impl<F> InitiateCoordinatorFailover<F>
221where
222    F: ClusterCommandForwarder,
223{
224    pub fn new(forwarder: F) -> Self {
225        Self { forwarder }
226    }
227
228    pub async fn execute(
229        &self,
230        request: InitiateCoordinatorFailoverRequest,
231    ) -> Result<InitiateCoordinatorFailoverResponse> {
232        if request.target_region.trim().is_empty() {
233            return Err(StasisError::PortFailure(
234                "target_region must not be empty".to_string(),
235            ));
236        }
237        if request.coordinator_node_id.trim().is_empty() {
238            return Err(StasisError::PortFailure(
239                "coordinator_node_id must not be empty".to_string(),
240            ));
241        }
242
243        let payload = json!({
244            "coordinator_node_id": request.coordinator_node_id,
245            "failover_to_node_id": request.failover_to_node_id,
246            "queue_scope": request.queue_scope,
247            "reason": request.reason,
248        })
249        .to_string();
250
251        let command_name = "coordinator.failover".to_string();
252        let accepted = self
253            .forwarder
254            .forward(ClusterForwardCommand {
255                target_region: request.target_region,
256                command_name: command_name.clone(),
257                payload,
258                correlation_id: request.correlation_id,
259                issued_at: request.issued_at,
260            })
261            .await?;
262
263        Ok(InitiateCoordinatorFailoverResponse {
264            accepted,
265            command_name,
266        })
267    }
268}
269
270impl<F> RebalanceQueueOwnership<F>
271where
272    F: ClusterCommandForwarder,
273{
274    pub fn new(forwarder: F) -> Self {
275        Self { forwarder }
276    }
277
278    pub async fn execute(
279        &self,
280        request: RebalanceQueueOwnershipRequest,
281    ) -> Result<RebalanceQueueOwnershipResponse> {
282        if request.target_region.trim().is_empty() {
283            return Err(StasisError::PortFailure(
284                "target_region must not be empty".to_string(),
285            ));
286        }
287        if request.queue.trim().is_empty() {
288            return Err(StasisError::PortFailure(
289                "queue must not be empty".to_string(),
290            ));
291        }
292        if request.desired_owners.is_empty() {
293            return Err(StasisError::PortFailure(
294                "desired_owners must not be empty".to_string(),
295            ));
296        }
297
298        let payload = json!({
299            "queue": request.queue,
300            "desired_owners": request.desired_owners,
301            "strategy": request.strategy,
302            "reason": request.reason,
303        })
304        .to_string();
305
306        let command_name = "queue_ownership.rebalance".to_string();
307        let accepted = self
308            .forwarder
309            .forward(ClusterForwardCommand {
310                target_region: request.target_region,
311                command_name: command_name.clone(),
312                payload,
313                correlation_id: request.correlation_id,
314                issued_at: request.issued_at,
315            })
316            .await?;
317
318        Ok(RebalanceQueueOwnershipResponse {
319            accepted,
320            command_name,
321        })
322    }
323}
324
325#[derive(Clone)]
326pub struct ListClusterForwardOutcomes<S>
327where
328    S: ClusterForwardOutcomeStore,
329{
330    store: S,
331}
332
333impl<S> ListClusterForwardOutcomes<S>
334where
335    S: ClusterForwardOutcomeStore,
336{
337    pub fn new(store: S) -> Self {
338        Self { store }
339    }
340
341    pub async fn execute(
342        &self,
343        request: &ListClusterForwardOutcomesRequest,
344    ) -> Result<Vec<ClusterForwardOutcomeRow>> {
345        let limit = request.limit.max(1);
346        let outcomes = self.store.list_recent(limit).await?;
347
348        Ok(outcomes
349            .into_iter()
350            .filter(|outcome| {
351                request
352                    .target_region
353                    .as_ref()
354                    .map(|region| &outcome.target_region == region)
355                    .unwrap_or(true)
356            })
357            .filter(|outcome| {
358                request
359                    .command_name
360                    .as_ref()
361                    .map(|command_name| &outcome.command_name == command_name)
362                    .unwrap_or(true)
363            })
364            .filter(|outcome| {
365                request
366                    .accepted
367                    .map(|accepted| outcome.accepted == accepted)
368                    .unwrap_or(true)
369            })
370            .map(|outcome| ClusterForwardOutcomeRow {
371                target_region: outcome.target_region,
372                command_name: outcome.command_name,
373                correlation_id: outcome.correlation_id,
374                accepted: outcome.accepted,
375                attempts: outcome.attempts,
376                error: outcome.error,
377                completed_at: outcome.completed_at,
378            })
379            .collect())
380    }
381}
382
383impl<F> ForwardClusterControlCommand<F>
384where
385    F: ClusterCommandForwarder,
386{
387    pub fn new(forwarder: F) -> Self {
388        Self { forwarder }
389    }
390
391    pub async fn execute(
392        &self,
393        request: ForwardClusterCommandRequest,
394    ) -> Result<ForwardClusterCommandResponse> {
395        if request.target_region.trim().is_empty() {
396            return Err(StasisError::PortFailure(
397                "target_region must not be empty".to_string(),
398            ));
399        }
400        if request.command_name.trim().is_empty() {
401            return Err(StasisError::PortFailure(
402                "command_name must not be empty".to_string(),
403            ));
404        }
405        if request.payload.trim().is_empty() {
406            return Err(StasisError::PortFailure(
407                "payload must not be empty".to_string(),
408            ));
409        }
410
411        let accepted = self
412            .forwarder
413            .forward(ClusterForwardCommand {
414                target_region: request.target_region,
415                command_name: request.command_name,
416                payload: request.payload,
417                correlation_id: request.correlation_id,
418                issued_at: request.issued_at,
419            })
420            .await?;
421
422        Ok(ForwardClusterCommandResponse { accepted })
423    }
424}
425
426impl<S, E> RunClusterHeartbeatSweep<S, E>
427where
428    S: ClusterNodeStore,
429    E: ClusterControlEventSink,
430{
431    pub fn new(store: S, event_sink: E) -> Self {
432        Self { store, event_sink }
433    }
434
435    pub async fn execute(
436        &self,
437        request: RunClusterHeartbeatSweepRequest,
438    ) -> Result<RunClusterHeartbeatSweepResponse> {
439        let pruned = self.store.prune_expired(request.now).await?;
440        if pruned == 0 {
441            return Ok(RunClusterHeartbeatSweepResponse {
442                pruned_nodes: 0,
443                emitted_events: 0,
444            });
445        }
446
447        self.event_sink
448            .emit(ClusterControlEvent::ExpiredNodesPruned {
449                pruned_count: pruned,
450                occurred_at: request.now,
451            })
452            .await?;
453
454        Ok(RunClusterHeartbeatSweepResponse {
455            pruned_nodes: pruned,
456            emitted_events: 1,
457        })
458    }
459}
460
461#[derive(Clone)]
462pub struct ListClusterNodeHealth<S>
463where
464    S: ClusterNodeStore,
465{
466    store: S,
467}
468
469impl<S> ListClusterNodeHealth<S>
470where
471    S: ClusterNodeStore,
472{
473    pub fn new(store: S) -> Self {
474        Self { store }
475    }
476
477    pub async fn execute(
478        &self,
479        request: &ListClusterNodeHealthRequest,
480    ) -> Result<Vec<ClusterNodeHealthRow>> {
481        let now = Utc::now();
482        let mut rows = self
483            .store
484            .list()
485            .await?
486            .into_iter()
487            .filter(|node| {
488                request
489                    .role
490                    .as_ref()
491                    .map(|role| &node.role == role)
492                    .unwrap_or(true)
493            })
494            .filter(|node| {
495                request
496                    .region
497                    .as_ref()
498                    .map(|region| &node.region == region)
499                    .unwrap_or(true)
500            })
501            .filter(|node| {
502                request
503                    .capability_tag
504                    .as_ref()
505                    .map(|tag| node.capability_tags.iter().any(|value| value == tag))
506                    .unwrap_or(true)
507            })
508            .filter(|node| {
509                request
510                    .queue
511                    .as_ref()
512                    .map(|queue| node.queue_ownership.iter().any(|value| value == queue))
513                    .unwrap_or(true)
514            })
515            .map(|node| {
516                let health = classify_health(&node, now);
517                ClusterNodeHealthRow {
518                    snapshot: ClusterNodeHealthSnapshot { node, health },
519                }
520            })
521            .filter(|row| {
522                request
523                    .health
524                    .as_ref()
525                    .map(|health| &row.snapshot.health == health)
526                    .unwrap_or(true)
527            })
528            .collect::<Vec<_>>();
529
530        rows.sort_by(|left, right| {
531            right
532                .snapshot
533                .node
534                .updated_at
535                .cmp(&left.snapshot.node.updated_at)
536                .then_with(|| left.snapshot.node.node_id.cmp(&right.snapshot.node.node_id))
537        });
538
539        let offset = request.offset;
540        if offset >= rows.len() {
541            return Ok(Vec::new());
542        }
543
544        let limit = request.limit.unwrap_or(rows.len().saturating_sub(offset));
545        Ok(rows.into_iter().skip(offset).take(limit).collect())
546    }
547}
548
549#[derive(Clone)]
550pub struct ListQueueOwnershipHealth<S>
551where
552    S: ClusterNodeStore,
553{
554    store: S,
555}
556
557impl<S> ListQueueOwnershipHealth<S>
558where
559    S: ClusterNodeStore,
560{
561    pub fn new(store: S) -> Self {
562        Self { store }
563    }
564
565    pub async fn execute(
566        &self,
567        request: &ListQueueOwnershipHealthRequest,
568    ) -> Result<Vec<QueueOwnershipHealthRow>> {
569        let now = Utc::now();
570        let nodes = self.store.list().await?;
571
572        let mut by_queue: BTreeMap<String, QueueOwnershipHealthRow> = BTreeMap::new();
573        for node in nodes {
574            let health = classify_health(&node, now);
575
576            for queue in &node.queue_ownership {
577                if let Some(prefix) = &request.queue_prefix
578                    && !queue.starts_with(prefix)
579                {
580                    continue;
581                }
582
583                let row =
584                    by_queue
585                        .entry(queue.clone())
586                        .or_insert_with(|| QueueOwnershipHealthRow {
587                            queue: queue.clone(),
588                            owners: Vec::new(),
589                            healthy_owners: 0,
590                            degraded_owners: 0,
591                            offline_owners: 0,
592                        });
593
594                row.owners.push(node.node_id.clone());
595                match health {
596                    ClusterNodeHealth::Healthy => row.healthy_owners += 1,
597                    ClusterNodeHealth::Degraded => row.degraded_owners += 1,
598                    ClusterNodeHealth::Offline => row.offline_owners += 1,
599                }
600            }
601        }
602
603        let mut rows = by_queue.into_values().collect::<Vec<_>>();
604        for row in &mut rows {
605            row.owners.sort();
606        }
607        rows.sort_by(|left, right| {
608            right
609                .offline_owners
610                .cmp(&left.offline_owners)
611                .then_with(|| right.degraded_owners.cmp(&left.degraded_owners))
612                .then_with(|| left.queue.cmp(&right.queue))
613        });
614        Ok(rows)
615    }
616}
617
618#[derive(Clone)]
619pub struct PruneExpiredClusterNodes<S>
620where
621    S: ClusterNodeStore,
622{
623    store: S,
624}
625
626impl<S> PruneExpiredClusterNodes<S>
627where
628    S: ClusterNodeStore,
629{
630    pub fn new(store: S) -> Self {
631        Self { store }
632    }
633
634    pub async fn execute(&self, request: PruneExpiredClusterNodesRequest) -> Result<u64> {
635        self.store.prune_expired(request.now).await
636    }
637}
638
639fn classify_health(node: &ClusterNode, now: DateTime<Utc>) -> ClusterNodeHealth {
640    if node.lease_expires_at < now {
641        return ClusterNodeHealth::Offline;
642    }
643
644    let stale_threshold = node.heartbeat_at + Duration::seconds(30);
645    if stale_threshold < now {
646        ClusterNodeHealth::Degraded
647    } else {
648        ClusterNodeHealth::Healthy
649    }
650}
651
652async fn validate_queue_ownership<S>(
653    store: &S,
654    node_id: &str,
655    requested_queues: &[String],
656    mode: QueueOwnershipMode,
657    at: DateTime<Utc>,
658) -> Result<()>
659where
660    S: ClusterNodeStore,
661{
662    if requested_queues.is_empty() || mode == QueueOwnershipMode::MultiOwner {
663        return Ok(());
664    }
665
666    let existing = store.list().await?;
667    for node in existing {
668        if node.node_id == node_id || node.lease_expires_at < at {
669            continue;
670        }
671
672        if let Some(conflict) = requested_queues
673            .iter()
674            .find(|queue| node.queue_ownership.iter().any(|owned| owned == *queue))
675        {
676            return Err(StasisError::PortFailure(format!(
677                "queue ownership conflict for queue={} with active node={}",
678                conflict, node.node_id
679            )));
680        }
681    }
682
683    Ok(())
684}
685
686#[cfg(test)]
687mod tests {
688    use chrono::{Duration, Utc};
689
690    use crate::application::dto::{
691        HeartbeatClusterNodeRequest, ListClusterNodeHealthRequest, ListQueueOwnershipHealthRequest,
692        PruneExpiredClusterNodesRequest, RegisterClusterNodeRequest,
693        RunClusterHeartbeatSweepRequest,
694    };
695    use crate::domain::runtime::cluster_node::{
696        ClusterControlEvent, ClusterNodeHealth, ClusterNodeRole, QueueOwnershipMode,
697    };
698    use crate::infrastructure::runtime::in_memory_cluster_control_event_sink::InMemoryClusterControlEventSink;
699    use crate::infrastructure::runtime::in_memory_cluster_node_store::InMemoryClusterNodeStore;
700
701    use super::{
702        HeartbeatClusterNode, ListClusterNodeHealth, ListQueueOwnershipHealth,
703        PruneExpiredClusterNodes, RegisterClusterNode, RunClusterHeartbeatSweep,
704    };
705
706    #[tokio::test]
707    async fn register_heartbeat_and_health_views_work() {
708        let store = InMemoryClusterNodeStore::default();
709        let register = RegisterClusterNode::new(store.clone());
710        let heartbeat = HeartbeatClusterNode::new(store.clone());
711        let health = ListClusterNodeHealth::new(store.clone());
712        let queue_health = ListQueueOwnershipHealth::new(store.clone());
713        let prune = PruneExpiredClusterNodes::new(store);
714
715        let now = Utc::now();
716        register
717            .execute(RegisterClusterNodeRequest {
718                node_id: "node.worker.a".to_string(),
719                role: ClusterNodeRole::Worker,
720                region: "us-east".to_string(),
721                queue_ownership: vec!["default".to_string(), "priority".to_string()],
722                capability_tags: vec!["gpu".to_string()],
723                heartbeat_at: now,
724                lease_ttl_seconds: 15,
725                queue_ownership_mode: None,
726                metadata: None,
727            })
728            .await
729            .expect("registration should succeed");
730
731        let rows = health
732            .execute(&ListClusterNodeHealthRequest {
733                role: Some(ClusterNodeRole::Worker),
734                region: Some("us-east".to_string()),
735                capability_tag: Some("gpu".to_string()),
736                queue: Some("default".to_string()),
737                health: Some(ClusterNodeHealth::Healthy),
738                offset: 0,
739                limit: Some(10),
740            })
741            .await
742            .expect("health list should succeed");
743        assert_eq!(rows.len(), 1);
744
745        heartbeat
746            .execute(HeartbeatClusterNodeRequest {
747                node_id: "node.worker.a".to_string(),
748                heartbeat_at: now + Duration::seconds(5),
749                lease_ttl_seconds: 30,
750                queue_ownership_mode: None,
751                queue_ownership: Some(vec!["default".to_string()]),
752                capability_tags: None,
753                metadata: Some("v2".to_string()),
754            })
755            .await
756            .expect("heartbeat should succeed");
757
758        let queue_rows = queue_health
759            .execute(&ListQueueOwnershipHealthRequest { queue_prefix: None })
760            .await
761            .expect("queue health should succeed");
762        assert_eq!(queue_rows.len(), 1);
763        assert_eq!(queue_rows[0].queue, "default");
764
765        let deleted = prune
766            .execute(PruneExpiredClusterNodesRequest {
767                now: now + Duration::minutes(10),
768            })
769            .await
770            .expect("prune should succeed");
771        assert_eq!(deleted, 1);
772
773        let event_sink = InMemoryClusterControlEventSink::default();
774        let sweep =
775            RunClusterHeartbeatSweep::new(InMemoryClusterNodeStore::default(), event_sink.clone());
776        let response = sweep
777            .execute(RunClusterHeartbeatSweepRequest { now: Utc::now() })
778            .await
779            .expect("sweep should succeed");
780        assert_eq!(response.pruned_nodes, 0);
781        assert_eq!(response.emitted_events, 0);
782
783        let events = event_sink.events().expect("event list should succeed");
784        assert!(events.is_empty());
785    }
786
787    #[tokio::test]
788    async fn single_owner_mode_rejects_conflicting_queue_registration() {
789        let store = InMemoryClusterNodeStore::default();
790        let register = RegisterClusterNode::new(store.clone());
791        let now = Utc::now();
792
793        register
794            .execute(RegisterClusterNodeRequest {
795                node_id: "node.worker.a".to_string(),
796                role: ClusterNodeRole::Worker,
797                region: "us-east".to_string(),
798                queue_ownership: vec!["default".to_string()],
799                capability_tags: vec![],
800                heartbeat_at: now,
801                lease_ttl_seconds: 60,
802                queue_ownership_mode: Some(QueueOwnershipMode::SingleOwner),
803                metadata: None,
804            })
805            .await
806            .expect("initial registration should succeed");
807
808        let err = register
809            .execute(RegisterClusterNodeRequest {
810                node_id: "node.worker.b".to_string(),
811                role: ClusterNodeRole::Worker,
812                region: "us-east".to_string(),
813                queue_ownership: vec!["default".to_string()],
814                capability_tags: vec![],
815                heartbeat_at: now,
816                lease_ttl_seconds: 60,
817                queue_ownership_mode: Some(QueueOwnershipMode::SingleOwner),
818                metadata: None,
819            })
820            .await
821            .expect_err("conflicting single owner queue should fail");
822
823        assert!(
824            err.to_string()
825                .contains("queue ownership conflict for queue=default")
826        );
827
828        let sweep_sink = InMemoryClusterControlEventSink::default();
829        let sweep = RunClusterHeartbeatSweep::new(store, sweep_sink.clone());
830        let response = sweep
831            .execute(RunClusterHeartbeatSweepRequest {
832                now: now + Duration::hours(2),
833            })
834            .await
835            .expect("sweep should succeed");
836        assert_eq!(response.pruned_nodes, 1);
837        assert_eq!(response.emitted_events, 1);
838
839        let events = sweep_sink.events().expect("event list should succeed");
840        assert_eq!(events.len(), 1);
841        assert!(matches!(
842            events[0],
843            ClusterControlEvent::ExpiredNodesPruned {
844                pruned_count: 1,
845                ..
846            }
847        ));
848    }
849}