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}