1use std::{
2 collections::{BTreeMap, BTreeSet, HashMap},
3 fmt,
4 net::SocketAddr,
5 sync::{Arc, Mutex},
6 time::{Duration, SystemTime, UNIX_EPOCH},
7};
8
9use datum::{NotUsed, Sink, StreamCompletion, Topic, TopicOverflow, TopicTryPublishError};
10use datum_cluster::{
11 ClusterConfig, ClusterNode, ClusterState, Member, MemberEvent, MemberEventKind, MemberState,
12 Signal,
13};
14use datum_net::quic::quinn;
15use tokio::{
16 sync::{mpsc, oneshot, watch},
17 task::JoinHandle,
18 time::{Instant, sleep},
19};
20
21use crate::{
22 Agent, AgentConfig, AgentError, AgentHandle, ClusterJobMetadata, ClusterPlacementHistory,
23 DesiredJobState, JobRegistryHandle, JobState, PlacementSpec as RegistryPlacementSpec,
24 PlacementStrategy as RegistryPlacementStrategy,
25 dcp::{
26 ClientKind, ClusterEvent, ClusterJobList, ClusterJobNode, ClusterJobStart,
27 ClusterNodeError, ClusterNodeList, ClusterNodeStatus, ClusterViewProvider,
28 CompleteShardingAsk, DcpClient, DcpError, DcpJobFactories, DcpServer, DcpServerConfig,
29 DcpServerHandle, ForwardShardEnvelopes, Hello, PlacementSpec, PlacementStrategy,
30 RememberClusterAssignment, ResponseStatus, ShardAllocation, ShardAllocationTable,
31 ShardEnvelopeBatchResult, ShardPipeClient, ShardPipeFrame, SubmitClusterJob,
32 client::MetricSubscription,
33 proto::MetricSample,
34 server::{
35 cluster_metadata_from_wire, placement_spec_from_wire, wire_cluster_job_start,
36 wire_job_status,
37 },
38 },
39};
40
41pub const AGENT_ROLE: &str = "agent";
43
44pub type ClusterAgentResult<T> = Result<T, ClusterAgentError>;
46
47#[derive(Debug, thiserror::Error)]
49pub enum ClusterAgentError {
50 #[error("invalid cluster-agent config: {0}")]
51 InvalidConfig(String),
52 #[error(transparent)]
53 Agent(#[from] crate::AgentError),
54 #[error(transparent)]
55 Dcp(#[from] DcpError),
56 #[error(transparent)]
57 Cluster(#[from] datum_cluster::ClusterError),
58}
59
60#[derive(Clone)]
62pub enum NodeSessionTransport {
63 TcpLoopback,
65 QuicMtls {
67 server_name: String,
68 client_config: quinn::ClientConfig,
69 },
70}
71
72impl fmt::Debug for NodeSessionTransport {
73 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
74 match self {
75 Self::TcpLoopback => formatter.write_str("TcpLoopback"),
76 Self::QuicMtls { server_name, .. } => formatter
77 .debug_struct("QuicMtls")
78 .field("server_name", server_name)
79 .finish_non_exhaustive(),
80 }
81 }
82}
83
84#[derive(Clone, Debug)]
86pub struct NodeSessionConfig {
87 pub agent_role: String,
89 pub transport: NodeSessionTransport,
91 pub reconnect_min_backoff: Duration,
93 pub reconnect_max_backoff: Duration,
95 pub request_timeout: Duration,
97 pub command_buffer: usize,
99}
100
101impl Default for NodeSessionConfig {
102 fn default() -> Self {
103 Self {
104 agent_role: AGENT_ROLE.to_owned(),
105 transport: NodeSessionTransport::TcpLoopback,
106 reconnect_min_backoff: Duration::from_millis(50),
107 reconnect_max_backoff: Duration::from_secs(2),
108 request_timeout: Duration::from_millis(750),
109 command_buffer: 32,
110 }
111 }
112}
113
114#[derive(Clone, Default)]
116pub struct ClusterAgentConfig {
117 pub agent: AgentConfig,
119 pub cluster: ClusterConfig,
121 pub dcp: DcpServerConfig,
123 pub sessions: NodeSessionConfig,
125}
126
127impl fmt::Debug for ClusterAgentConfig {
128 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
129 formatter
130 .debug_struct("ClusterAgentConfig")
131 .field("agent", &self.agent)
132 .field("cluster", &self.cluster)
133 .field("dcp", &"<DcpServerConfig>")
134 .field("sessions", &self.sessions)
135 .finish()
136 }
137}
138
139pub struct ClusterAgent;
142
143impl ClusterAgent {
144 pub async fn start(
145 config: ClusterAgentConfig,
146 factories: DcpJobFactories,
147 ) -> ClusterAgentResult<ClusterAgentHandle> {
148 let mut config = config;
149 validate_cluster_agent_config(&config)?;
150 ensure_agent_role(&mut config.cluster, &config.sessions.agent_role);
151 config.dcp.node_id = config.cluster.node_id.clone();
152
153 let cluster_events = ClusterEventPublisher::new(config.agent.event_buffer)?;
154 let agent = Agent::start_with_config(config.agent.clone())?;
155 let server = DcpServer::from_agent(&agent, factories.clone(), config.dcp.clone());
156 let server_handle = server.start().await?;
157 let agent_addr = advertised_agent_addr(&config.sessions.transport, &server_handle)?;
158 config.cluster.agent_addr = Some(agent_addr);
159
160 let cluster = ClusterNode::start(config.cluster.clone()).await?;
161 let membership_events = cluster
162 .events()
163 .changes()
164 .run_with(Sink::foreach({
165 let cluster_events = cluster_events.clone();
166 move |event| cluster_events.publish_member(&event)
167 }))
168 .map_err(crate::AgentError::from)?;
169 let sessions = NodeSessionManagerHandle::start(
170 config.sessions.clone(),
171 cluster.node_id().to_owned(),
172 cluster.state(),
173 )?;
174 let placement = PlacementCoordinatorHandle::start(
175 cluster.node_id().to_owned(),
176 config.sessions.agent_role.clone(),
177 cluster.state(),
178 PlacementDependencies {
179 registry: agent.registry().clone(),
180 sessions: sessions.clone(),
181 factories: factories.clone(),
182 cluster_events: cluster_events.clone(),
183 },
184 config.sessions.request_timeout,
185 )?;
186 let provider = Arc::new(ClusterView {
187 registry: agent.registry().clone(),
188 state: cluster.state(),
189 sessions: sessions.clone(),
190 placement: placement.clone(),
191 self_node: cluster.node_id().to_owned(),
192 agent_role: config.sessions.agent_role.clone(),
193 cluster_events: cluster_events.clone(),
194 });
195 server.set_cluster_view(provider);
196
197 Ok(ClusterAgentHandle {
198 agent,
199 cluster,
200 server,
201 server_handle: Some(server_handle),
202 sessions,
203 placement,
204 cluster_events,
205 _membership_events: membership_events,
206 })
207 }
208}
209
210pub struct ClusterAgentHandle {
212 agent: AgentHandle,
213 cluster: ClusterNode,
214 server: DcpServer,
215 server_handle: Option<DcpServerHandle>,
216 sessions: NodeSessionManagerHandle,
217 placement: PlacementCoordinatorHandle,
218 cluster_events: ClusterEventPublisher,
219 _membership_events: StreamCompletion<NotUsed>,
220}
221
222impl ClusterAgentHandle {
223 #[must_use]
224 pub fn agent(&self) -> &AgentHandle {
225 &self.agent
226 }
227
228 #[must_use]
229 pub fn cluster(&self) -> &ClusterNode {
230 &self.cluster
231 }
232
233 #[must_use]
234 pub fn sessions(&self) -> &NodeSessionManagerHandle {
235 &self.sessions
236 }
237
238 #[must_use]
239 pub fn server(&self) -> &DcpServer {
240 &self.server
241 }
242
243 #[must_use]
244 pub fn tcp_addr(&self) -> Option<SocketAddr> {
245 self.server_handle
246 .as_ref()
247 .and_then(DcpServerHandle::tcp_addr)
248 }
249
250 #[must_use]
251 pub fn quic_addr(&self) -> Option<SocketAddr> {
252 self.server_handle
253 .as_ref()
254 .and_then(DcpServerHandle::quic_addr)
255 }
256
257 pub async fn shutdown(mut self) -> ClusterAgentResult<()> {
258 let node_id = self.cluster.node_id().to_owned();
259 self.server.clear_cluster_view();
260 self.server.clear_sharding_view();
261 self.placement.shutdown().await;
262 eprintln!(
263 "datum-agent INFO shutdown_component_stopped node_id={node_id} component=placement"
264 );
265 self.sessions.shutdown().await;
266 eprintln!(
267 "datum-agent INFO shutdown_component_stopped node_id={node_id} component=sessions"
268 );
269 let _ = self.cluster.abort().await;
270 eprintln!(
271 "datum-agent INFO shutdown_component_stopped node_id={node_id} component=cluster"
272 );
273 let _ = self.cluster_events.close();
274 if let Some(handle) = self.server_handle.take() {
275 handle.shutdown().await;
276 }
277 eprintln!(
278 "datum-agent INFO shutdown_component_stopped node_id={node_id} component=dcp_server"
279 );
280 self.agent.registry().shutdown()?;
281 eprintln!(
282 "datum-agent INFO shutdown_component_stopped node_id={node_id} component=registry"
283 );
284 eprintln!("datum-agent INFO shutdown_done node_id={node_id}");
285 Ok(())
286 }
287}
288
289impl Drop for ClusterAgentHandle {
290 fn drop(&mut self) {
291 self.server.clear_cluster_view();
292 self.server.clear_sharding_view();
293 let _ = self.cluster_events.close();
294 }
295}
296
297#[derive(Clone)]
298struct ClusterEventPublisher {
299 topic: Topic<ClusterEvent>,
300 next_sequence: Arc<Mutex<u64>>,
301 subscriber_buffer: usize,
302}
303
304impl ClusterEventPublisher {
305 fn new(buffer: usize) -> crate::AgentResult<Self> {
306 let subscriber_buffer = buffer.max(1);
307 Ok(Self {
308 topic: Topic::new(subscriber_buffer, TopicOverflow::Sliding)?,
309 next_sequence: Arc::new(Mutex::new(1)),
310 subscriber_buffer,
311 })
312 }
313
314 fn subscribe(&self) -> crate::dcp::DcpResult<mpsc::Receiver<ClusterEvent>> {
315 let (sender, receiver) = mpsc::channel(self.subscriber_buffer);
316 let completion = self
317 .topic
318 .subscribe()
319 .run_with(Sink::foreach_result(move |event| {
320 sender
321 .blocking_send(event)
322 .map_err(|_| datum::StreamError::Cancelled)
323 }))?;
324 tokio::task::spawn_blocking(move || {
325 let _ = completion.wait();
326 });
327 Ok(receiver)
328 }
329
330 fn publish_member(&self, event: &MemberEvent) {
331 self.publish(
332 system_time_ms(event.at),
333 member_event_kind(event.kind),
334 event.member.node_id.clone(),
335 format!(
336 "state={:?} address={} unreachable={} incarnation={}",
337 event.member.state,
338 event.member.address,
339 event.member.unreachable,
340 event.member.incarnation
341 ),
342 Some(event.member.incarnation),
343 );
344 }
345
346 fn publish_coordinator_changed(
347 &self,
348 previous: Option<&str>,
349 current: Option<&str>,
350 observer: &str,
351 ) {
352 self.publish(
353 system_time_ms(SystemTime::now()),
354 "CoordinatorChanged",
355 current.unwrap_or(observer).to_owned(),
356 format!(
357 "previous={} current={}",
358 previous.unwrap_or("none"),
359 current.unwrap_or("none")
360 ),
361 None,
362 );
363 }
364
365 fn publish_job_replaced(
366 &self,
367 at: SystemTime,
368 name: &str,
369 old_node: &str,
370 target: &str,
371 reason: &str,
372 generation: u64,
373 ) {
374 self.publish(
375 system_time_ms(at),
376 "JobReplaced",
377 target.to_owned(),
378 format!("job={name} from={old_node} to={target} reason={reason}"),
379 Some(generation),
380 );
381 }
382
383 fn publish(
384 &self,
385 timestamp_ms: u64,
386 kind: &str,
387 node_id: String,
388 detail: String,
389 generation: Option<u64>,
390 ) {
391 let mut next_sequence = self
392 .next_sequence
393 .lock()
394 .unwrap_or_else(|poison| poison.into_inner());
395 let event = ClusterEvent {
396 sequence: *next_sequence,
397 timestamp_ms,
398 kind: kind.to_owned(),
399 node_id,
400 detail,
401 generation,
402 };
403 *next_sequence = next_sequence.wrapping_add(1).max(1);
404 match self.topic.try_publish(event) {
405 Ok(())
406 | Err(TopicTryPublishError::Full(_))
407 | Err(TopicTryPublishError::Busy(_))
408 | Err(TopicTryPublishError::Closed(_)) => {}
409 }
410 }
411
412 fn close(&self) -> datum::StreamResult<()> {
413 if self.topic.is_closed() {
414 Ok(())
415 } else {
416 self.topic.close()
417 }
418 }
419}
420
421const fn member_event_kind(kind: MemberEventKind) -> &'static str {
422 match kind {
423 MemberEventKind::Initialized => "Initialized",
424 MemberEventKind::MemberJoining => "MemberJoining",
425 MemberEventKind::MemberUp => "MemberUp",
426 MemberEventKind::MemberUpdated => "MemberUpdated",
427 MemberEventKind::MemberUnreachable => "MemberUnreachable",
428 MemberEventKind::MemberReachable => "MemberReachable",
429 MemberEventKind::MemberLeaving => "MemberLeaving",
430 MemberEventKind::MemberExiting => "MemberExiting",
431 MemberEventKind::MemberDown => "MemberDown",
432 MemberEventKind::MemberRemoved => "MemberRemoved",
433 MemberEventKind::MemberRejoined => "MemberRejoined",
434 }
435}
436
437fn system_time_ms(time: SystemTime) -> u64 {
438 time.duration_since(UNIX_EPOCH)
439 .map(|duration| duration.as_millis().min(u128::from(u64::MAX)) as u64)
440 .unwrap_or(0)
441}
442
443#[derive(Clone)]
444pub struct NodeSessionManagerHandle {
445 inner: Arc<NodeSessionManagerInner>,
446}
447
448struct NodeSessionManagerInner {
449 config: NodeSessionConfig,
450 self_node: String,
451 state: Signal<ClusterState>,
452 sessions: tokio::sync::Mutex<BTreeMap<String, PeerSession>>,
453 tasks: Mutex<Vec<JoinHandle<()>>>,
454}
455
456struct PeerSession {
457 member: Member,
458 commands: mpsc::Sender<SessionCommand>,
459 pipe_commands: mpsc::Sender<ShardPipeCommand>,
460 stop: watch::Sender<bool>,
461 state: watch::Receiver<PeerSessionState>,
462 task: JoinHandle<()>,
463}
464
465#[derive(Debug, Clone, PartialEq, Eq)]
466enum PeerSessionState {
467 Connecting,
468 Connected(u64),
469 BackingOff(String),
470 Closed,
471}
472
473impl PeerSessionState {
474 fn as_text(&self) -> String {
475 match self {
476 Self::Connecting => "connecting".to_owned(),
477 Self::Connected(_) => "connected".to_owned(),
478 Self::BackingOff(error) => format!("backing_off:{error}"),
479 Self::Closed => "closed".to_owned(),
480 }
481 }
482}
483
484#[derive(Debug, Clone, PartialEq, Eq)]
485struct PeerSessionSnapshot {
486 member_incarnation: u64,
487 agent_addr: Option<SocketAddr>,
488 state: PeerSessionState,
489}
490
491#[derive(Debug, Clone, PartialEq, Eq)]
492struct MetricRouteTarget {
493 member_incarnation: u64,
494 agent_addr: Option<SocketAddr>,
495 connection_generation: u64,
496}
497
498impl PeerSessionSnapshot {
499 fn metric_route_target(&self) -> Option<MetricRouteTarget> {
500 match &self.state {
501 PeerSessionState::Connected(connection_generation) => Some(MetricRouteTarget {
502 member_incarnation: self.member_incarnation,
503 agent_addr: self.agent_addr,
504 connection_generation: *connection_generation,
505 }),
506 PeerSessionState::Connecting
507 | PeerSessionState::BackingOff(_)
508 | PeerSessionState::Closed => None,
509 }
510 }
511}
512
513enum SessionCommand {
514 ListJobs {
515 timeout: Duration,
516 reply: oneshot::Sender<Result<Vec<crate::dcp::proto::JobStatus>, String>>,
517 },
518 SubmitClusterJob {
519 request: SubmitClusterJob,
520 timeout: Duration,
521 reply: oneshot::Sender<Result<crate::dcp::proto::JobStatus, String>>,
522 },
523 StartClusterJob {
524 factory_name: String,
525 instance_name: String,
526 params: HashMap<String, String>,
527 assignment: ClusterJobStart,
528 timeout: Duration,
529 reply: oneshot::Sender<Result<crate::dcp::proto::JobStatus, String>>,
530 },
531 ClusterJobStatus {
532 name: String,
533 cluster: bool,
534 timeout: Duration,
535 reply: oneshot::Sender<Result<crate::dcp::proto::JobStatus, String>>,
536 },
537 DrainClusterJob {
538 name: String,
539 cluster: bool,
540 timeout: Duration,
541 reply: oneshot::Sender<Result<crate::dcp::proto::JobStatus, String>>,
542 },
543 StopClusterJob {
544 name: String,
545 cluster: bool,
546 timeout: Duration,
547 reply: oneshot::Sender<Result<crate::dcp::proto::JobStatus, String>>,
548 },
549 ListClusterJobs {
550 timeout: Duration,
551 reply: oneshot::Sender<Result<ClusterJobList, String>>,
552 },
553 SubscribeMetrics {
554 interval: Duration,
555 job_names: Vec<String>,
556 timeout: Duration,
557 reply: oneshot::Sender<Result<MetricSubscription, String>>,
558 },
559 RememberClusterAssignment {
560 instance_name: String,
561 assignment: Option<ClusterJobStart>,
562 tombstone: Option<AssignmentTombstone>,
563 timeout: Duration,
564 reply: oneshot::Sender<Result<(), String>>,
565 },
566 AllocateShard {
567 type_name: String,
568 shard_id: String,
569 timeout: Duration,
570 reply: oneshot::Sender<Result<ShardAllocation, String>>,
571 },
572 RememberShardAllocations {
573 table: ShardAllocationTable,
574 timeout: Duration,
575 reply: oneshot::Sender<Result<(), String>>,
576 },
577 GetShardAllocations {
578 type_name: String,
579 timeout: Duration,
580 reply: oneshot::Sender<Result<ShardAllocationTable, String>>,
581 },
582 ForwardShardEnvelopes {
583 batch: ForwardShardEnvelopes,
584 timeout: Duration,
585 reply: oneshot::Sender<Result<ShardEnvelopeBatchResult, String>>,
586 },
587 CompleteShardingAsk {
588 response: CompleteShardingAsk,
589 timeout: Duration,
590 reply: oneshot::Sender<Result<(), String>>,
591 },
592}
593
594enum ShardPipeCommand {
595 Forward { batch: ForwardShardEnvelopes },
596 Reply { response: CompleteShardingAsk },
597}
598
599impl NodeSessionManagerHandle {
600 fn start(
601 config: NodeSessionConfig,
602 self_node: String,
603 state: Signal<ClusterState>,
604 ) -> ClusterAgentResult<Self> {
605 let inner = Arc::new(NodeSessionManagerInner {
606 config,
607 self_node,
608 state,
609 sessions: tokio::sync::Mutex::new(BTreeMap::new()),
610 tasks: Mutex::new(Vec::new()),
611 });
612
613 let manager_task = {
614 let inner = Arc::clone(&inner);
615 tokio::spawn(async move {
616 run_node_session_manager(inner).await;
617 })
618 };
619 {
620 let mut tasks = inner.tasks.lock().expect("node-session tasks poisoned");
621 tasks.push(manager_task);
622 }
623
624 Ok(Self { inner })
625 }
626
627 pub async fn shutdown(&self) {
628 let sessions = {
629 let mut locked = self.inner.sessions.lock().await;
630 std::mem::take(&mut *locked)
631 };
632 for (_, session) in sessions {
633 let _ = session.stop.send(true);
634 session.task.abort();
635 }
636 let tasks = {
637 let mut locked = self
638 .inner
639 .tasks
640 .lock()
641 .expect("node-session tasks poisoned");
642 std::mem::take(&mut *locked)
643 };
644 for task in tasks {
645 task.abort();
646 }
647 }
648
649 async fn list_jobs(
650 &self,
651 node_id: &str,
652 timeout: Duration,
653 ) -> Result<Vec<crate::dcp::proto::JobStatus>, String> {
654 let sender = {
655 let locked = self.inner.sessions.lock().await;
656 locked.get(node_id).map(|session| session.commands.clone())
657 }
658 .ok_or_else(|| "node session unavailable".to_owned())?;
659 let (reply, receiver) = oneshot::channel();
660 sender
661 .try_send(SessionCommand::ListJobs { timeout, reply })
662 .map_err(|error| format!("node session command queue unavailable: {error}"))?;
663 tokio::time::timeout(timeout, receiver)
664 .await
665 .map_err(|_| "node session request timed out".to_owned())?
666 .map_err(|_| "node session closed".to_owned())?
667 }
668
669 async fn submit_cluster_job(
670 &self,
671 node_id: &str,
672 request: SubmitClusterJob,
673 timeout: Duration,
674 ) -> Result<crate::dcp::proto::JobStatus, String> {
675 self.request_status(node_id, timeout, |reply| SessionCommand::SubmitClusterJob {
676 request,
677 timeout,
678 reply,
679 })
680 .await
681 }
682
683 async fn start_cluster_job(
684 &self,
685 node_id: &str,
686 factory_name: String,
687 instance_name: String,
688 params: HashMap<String, String>,
689 assignment: ClusterJobStart,
690 timeout: Duration,
691 ) -> Result<crate::dcp::proto::JobStatus, String> {
692 self.request_status(node_id, timeout, |reply| SessionCommand::StartClusterJob {
693 factory_name,
694 instance_name,
695 params,
696 assignment,
697 timeout,
698 reply,
699 })
700 .await
701 }
702
703 async fn cluster_job_status(
704 &self,
705 node_id: &str,
706 name: String,
707 timeout: Duration,
708 ) -> Result<crate::dcp::proto::JobStatus, String> {
709 self.request_status(node_id, timeout, |reply| SessionCommand::ClusterJobStatus {
710 name,
711 cluster: false,
712 timeout,
713 reply,
714 })
715 .await
716 }
717
718 async fn drain_cluster_job(
719 &self,
720 node_id: &str,
721 name: String,
722 timeout: Duration,
723 ) -> Result<crate::dcp::proto::JobStatus, String> {
724 self.request_status(node_id, timeout, |reply| SessionCommand::DrainClusterJob {
725 name,
726 cluster: false,
727 timeout,
728 reply,
729 })
730 .await
731 }
732
733 async fn stop_cluster_job(
734 &self,
735 node_id: &str,
736 name: String,
737 timeout: Duration,
738 ) -> Result<crate::dcp::proto::JobStatus, String> {
739 self.request_status(node_id, timeout, |reply| SessionCommand::StopClusterJob {
740 name,
741 cluster: false,
742 timeout,
743 reply,
744 })
745 .await
746 }
747
748 async fn forward_cluster_job_status(
749 &self,
750 node_id: &str,
751 name: String,
752 timeout: Duration,
753 ) -> Result<crate::dcp::proto::JobStatus, String> {
754 self.request_status(node_id, timeout, |reply| SessionCommand::ClusterJobStatus {
755 name,
756 cluster: true,
757 timeout,
758 reply,
759 })
760 .await
761 }
762
763 async fn forward_drain_cluster_job(
764 &self,
765 node_id: &str,
766 name: String,
767 timeout: Duration,
768 ) -> Result<crate::dcp::proto::JobStatus, String> {
769 self.request_status(node_id, timeout, |reply| SessionCommand::DrainClusterJob {
770 name,
771 cluster: true,
772 timeout,
773 reply,
774 })
775 .await
776 }
777
778 async fn forward_stop_cluster_job(
779 &self,
780 node_id: &str,
781 name: String,
782 timeout: Duration,
783 ) -> Result<crate::dcp::proto::JobStatus, String> {
784 self.request_status(node_id, timeout, |reply| SessionCommand::StopClusterJob {
785 name,
786 cluster: true,
787 timeout,
788 reply,
789 })
790 .await
791 }
792
793 async fn list_cluster_jobs(
794 &self,
795 node_id: &str,
796 timeout: Duration,
797 ) -> Result<ClusterJobList, String> {
798 let sender = self.session_sender(node_id).await?;
799 let (reply, receiver) = oneshot::channel();
800 sender
801 .try_send(SessionCommand::ListClusterJobs { timeout, reply })
802 .map_err(|error| format!("node session command queue unavailable: {error}"))?;
803 wait_session_reply(timeout, receiver).await
804 }
805
806 async fn subscribe_metrics(
807 &self,
808 node_id: &str,
809 interval: Duration,
810 job_names: Vec<String>,
811 timeout: Duration,
812 ) -> Result<MetricSubscription, String> {
813 let sender = self.session_sender(node_id).await?;
814 let (reply, receiver) = oneshot::channel();
815 sender
816 .try_send(SessionCommand::SubscribeMetrics {
817 interval,
818 job_names,
819 timeout,
820 reply,
821 })
822 .map_err(|error| format!("node session command queue unavailable: {error}"))?;
823 wait_session_reply(timeout, receiver).await
824 }
825
826 async fn remember_cluster_assignment(
827 &self,
828 node_id: &str,
829 instance_name: String,
830 assignment: Option<ClusterJobStart>,
831 tombstone: Option<AssignmentTombstone>,
832 timeout: Duration,
833 ) -> Result<(), String> {
834 let sender = self.session_sender(node_id).await?;
835 let (reply, receiver) = oneshot::channel();
836 sender
837 .try_send(SessionCommand::RememberClusterAssignment {
838 instance_name,
839 assignment,
840 tombstone,
841 timeout,
842 reply,
843 })
844 .map_err(|error| format!("node session command queue unavailable: {error}"))?;
845 wait_session_reply(timeout, receiver).await
846 }
847
848 pub async fn allocate_shard(
849 &self,
850 node_id: &str,
851 type_name: String,
852 shard_id: String,
853 timeout: Duration,
854 ) -> Result<ShardAllocation, String> {
855 let sender = self.session_sender(node_id).await?;
856 let (reply, receiver) = oneshot::channel();
857 sender
858 .try_send(SessionCommand::AllocateShard {
859 type_name,
860 shard_id,
861 timeout,
862 reply,
863 })
864 .map_err(|error| format!("node session command queue unavailable: {error}"))?;
865 wait_session_reply(timeout, receiver).await
866 }
867
868 pub async fn remember_shard_allocations(
869 &self,
870 node_id: &str,
871 table: ShardAllocationTable,
872 timeout: Duration,
873 ) -> Result<(), String> {
874 let sender = self.session_sender(node_id).await?;
875 let (reply, receiver) = oneshot::channel();
876 sender
877 .try_send(SessionCommand::RememberShardAllocations {
878 table,
879 timeout,
880 reply,
881 })
882 .map_err(|error| format!("node session command queue unavailable: {error}"))?;
883 wait_session_reply(timeout, receiver).await
884 }
885
886 pub async fn get_shard_allocations(
887 &self,
888 node_id: &str,
889 type_name: String,
890 timeout: Duration,
891 ) -> Result<ShardAllocationTable, String> {
892 let sender = self.session_sender(node_id).await?;
893 let (reply, receiver) = oneshot::channel();
894 sender
895 .try_send(SessionCommand::GetShardAllocations {
896 type_name,
897 timeout,
898 reply,
899 })
900 .map_err(|error| format!("node session command queue unavailable: {error}"))?;
901 wait_session_reply(timeout, receiver).await
902 }
903
904 pub async fn forward_shard_envelopes(
905 &self,
906 node_id: &str,
907 batch: ForwardShardEnvelopes,
908 timeout: Duration,
909 ) -> Result<ShardEnvelopeBatchResult, String> {
910 let sender = self.session_sender(node_id).await?;
911 let (reply, receiver) = oneshot::channel();
912 sender
913 .try_send(SessionCommand::ForwardShardEnvelopes {
914 batch,
915 timeout,
916 reply,
917 })
918 .map_err(|error| format!("node session command queue unavailable: {error}"))?;
919 wait_session_reply(timeout, receiver).await
920 }
921
922 pub async fn complete_sharding_ask(
923 &self,
924 node_id: &str,
925 response: CompleteShardingAsk,
926 timeout: Duration,
927 ) -> Result<(), String> {
928 let sender = self.session_sender(node_id).await?;
929 let (reply, receiver) = oneshot::channel();
930 sender
931 .try_send(SessionCommand::CompleteShardingAsk {
932 response,
933 timeout,
934 reply,
935 })
936 .map_err(|error| format!("node session command queue unavailable: {error}"))?;
937 wait_session_reply(timeout, receiver).await
938 }
939
940 pub async fn forward_shard_pipe_envelopes(
941 &self,
942 node_id: &str,
943 batch: ForwardShardEnvelopes,
944 ) -> Result<(), String> {
945 let sender = self.pipe_sender(node_id).await?;
946 sender
947 .try_send(ShardPipeCommand::Forward { batch })
948 .map_err(|error| format!("node session shard pipe queue unavailable: {error}"))
949 }
950
951 pub async fn complete_sharding_ask_pipe(
952 &self,
953 node_id: &str,
954 response: CompleteShardingAsk,
955 ) -> Result<(), String> {
956 let sender = self.pipe_sender(node_id).await?;
957 sender
958 .try_send(ShardPipeCommand::Reply { response })
959 .map_err(|error| format!("node session shard pipe queue unavailable: {error}"))
960 }
961
962 pub fn try_complete_sharding_ask_pipe(
963 &self,
964 node_id: &str,
965 response: CompleteShardingAsk,
966 ) -> Result<(), String> {
967 let locked = self
968 .inner
969 .sessions
970 .try_lock()
971 .map_err(|_| "node session table busy".to_owned())?;
972 let sender = locked
973 .get(node_id)
974 .map(|session| session.pipe_commands.clone())
975 .ok_or_else(|| "node session unavailable".to_owned())?;
976 sender
977 .try_send(ShardPipeCommand::Reply { response })
978 .map_err(|error| format!("node session shard pipe queue unavailable: {error}"))
979 }
980
981 async fn request_status<F>(
982 &self,
983 node_id: &str,
984 timeout: Duration,
985 make_command: F,
986 ) -> Result<crate::dcp::proto::JobStatus, String>
987 where
988 F: FnOnce(oneshot::Sender<Result<crate::dcp::proto::JobStatus, String>>) -> SessionCommand,
989 {
990 let sender = self.session_sender(node_id).await?;
991 let (reply, receiver) = oneshot::channel();
992 sender
993 .try_send(make_command(reply))
994 .map_err(|error| format!("node session command queue unavailable: {error}"))?;
995 wait_session_reply(timeout, receiver).await
996 }
997
998 async fn session_sender(&self, node_id: &str) -> Result<mpsc::Sender<SessionCommand>, String> {
999 let locked = self.inner.sessions.lock().await;
1000 locked
1001 .get(node_id)
1002 .map(|session| session.commands.clone())
1003 .ok_or_else(|| "node session unavailable".to_owned())
1004 }
1005
1006 async fn pipe_sender(&self, node_id: &str) -> Result<mpsc::Sender<ShardPipeCommand>, String> {
1007 let locked = self.inner.sessions.lock().await;
1008 locked
1009 .get(node_id)
1010 .map(|session| session.pipe_commands.clone())
1011 .ok_or_else(|| "node session unavailable".to_owned())
1012 }
1013
1014 async fn session_snapshots(&self) -> BTreeMap<String, PeerSessionSnapshot> {
1015 let locked = self.inner.sessions.lock().await;
1016 locked
1017 .iter()
1018 .map(|(node_id, session)| {
1019 (
1020 node_id.clone(),
1021 PeerSessionSnapshot {
1022 member_incarnation: session.member.incarnation,
1023 agent_addr: session.member.agent_addr,
1024 state: session.state.borrow().clone(),
1025 },
1026 )
1027 })
1028 .collect()
1029 }
1030}
1031
1032async fn wait_session_reply<T>(
1033 timeout: Duration,
1034 receiver: oneshot::Receiver<Result<T, String>>,
1035) -> Result<T, String> {
1036 tokio::time::timeout(timeout, receiver)
1037 .await
1038 .map_err(|_| "node session request timed out".to_owned())?
1039 .map_err(|_| "node session closed".to_owned())?
1040}
1041
1042async fn run_node_session_manager(inner: Arc<NodeSessionManagerInner>) {
1043 let interval = inner
1044 .config
1045 .reconnect_min_backoff
1046 .min(Duration::from_millis(100))
1047 .max(Duration::from_millis(10));
1048 loop {
1049 reconcile_sessions(&inner).await;
1050 tokio::time::sleep(interval).await;
1051 }
1052}
1053
1054async fn reconcile_sessions(inner: &Arc<NodeSessionManagerInner>) {
1055 let state = inner.state.get();
1056 let mut wanted = BTreeSet::new();
1057 for member in state.members.values() {
1058 if eligible_member(inner, member) {
1059 wanted.insert(member.node_id.clone());
1060 ensure_session(inner, member.clone()).await;
1061 }
1062 }
1063 close_unwanted_sessions(inner, &wanted).await;
1064}
1065
1066fn eligible_member(inner: &NodeSessionManagerInner, member: &Member) -> bool {
1067 member.node_id != inner.self_node
1068 && member.state == MemberState::Up
1069 && member.has_role(&inner.config.agent_role)
1070 && member.agent_addr.is_some()
1071}
1072
1073async fn ensure_session(inner: &Arc<NodeSessionManagerInner>, member: Member) {
1074 let mut locked = inner.sessions.lock().await;
1075 if let Some(existing) = locked.get(&member.node_id)
1076 && existing.member.agent_addr == member.agent_addr
1077 && existing.member.incarnation == member.incarnation
1078 {
1079 return;
1080 }
1081
1082 if let Some(existing) = locked.remove(&member.node_id) {
1083 let _ = existing.stop.send(true);
1084 existing.task.abort();
1085 }
1086
1087 let (commands, receiver) = mpsc::channel(inner.config.command_buffer.max(1));
1088 let pipe_buffer = inner.config.command_buffer.saturating_mul(4).max(1);
1089 let (pipe_commands, pipe_receiver) = mpsc::channel(pipe_buffer);
1090 let (stop, stop_receiver) = watch::channel(false);
1091 let (state_sender, state_receiver) = watch::channel(PeerSessionState::Connecting);
1092 let config = inner.config.clone();
1093 let local_node = inner.self_node.clone();
1094 let task_member = member.clone();
1095 let task = tokio::spawn(async move {
1096 run_peer_session(
1097 config,
1098 local_node,
1099 task_member,
1100 receiver,
1101 pipe_receiver,
1102 stop_receiver,
1103 state_sender,
1104 )
1105 .await;
1106 });
1107
1108 locked.insert(
1109 member.node_id.clone(),
1110 PeerSession {
1111 member,
1112 commands,
1113 pipe_commands,
1114 stop,
1115 state: state_receiver,
1116 task,
1117 },
1118 );
1119}
1120
1121async fn close_unwanted_sessions(inner: &Arc<NodeSessionManagerInner>, wanted: &BTreeSet<String>) {
1122 let to_close = {
1123 let locked = inner.sessions.lock().await;
1124 locked
1125 .keys()
1126 .filter(|node_id| !wanted.contains(*node_id))
1127 .cloned()
1128 .collect::<Vec<_>>()
1129 };
1130 for node_id in to_close {
1131 close_session(inner, &node_id).await;
1132 }
1133}
1134
1135async fn close_session(inner: &Arc<NodeSessionManagerInner>, node_id: &str) {
1136 let removed = {
1137 let mut locked = inner.sessions.lock().await;
1138 locked.remove(node_id)
1139 };
1140 if let Some(session) = removed {
1141 let _ = session.stop.send(true);
1142 session.task.abort();
1143 }
1144}
1145
1146async fn run_peer_session(
1147 config: NodeSessionConfig,
1148 local_node: String,
1149 member: Member,
1150 mut commands: mpsc::Receiver<SessionCommand>,
1151 mut pipe_commands: mpsc::Receiver<ShardPipeCommand>,
1152 mut stop: watch::Receiver<bool>,
1153 state: watch::Sender<PeerSessionState>,
1154) {
1155 let mut backoff = config.reconnect_min_backoff;
1156 let mut connection_generation = 0_u64;
1157 loop {
1158 if *stop.borrow() {
1159 break;
1160 }
1161 let _ = state.send(PeerSessionState::Connecting);
1162 match connect_peer(&config.transport, &local_node, &member).await {
1163 Ok(client) => match connect_shard_pipe(&config.transport, &local_node, &member).await {
1164 Ok(pipe) => {
1165 connection_generation = connection_generation.saturating_add(1).max(1);
1166 let _ = state.send(PeerSessionState::Connected(connection_generation));
1167 eprintln!(
1168 "datum-agent INFO peer_session_connected node_id={local_node} peer={}",
1169 member.node_id
1170 );
1171 backoff = config.reconnect_min_backoff;
1172 if !run_connected_peer(
1173 client,
1174 pipe,
1175 &mut commands,
1176 &mut pipe_commands,
1177 &mut stop,
1178 )
1179 .await
1180 {
1181 break;
1182 }
1183 eprintln!(
1184 "datum-agent WARN peer_session_lost node_id={local_node} peer={}",
1185 member.node_id
1186 );
1187 }
1188 Err(error) => {
1189 log_peer_connect_failed(&local_node, &member.node_id, &error);
1190 let _ = state.send(PeerSessionState::BackingOff(error.clone()));
1191 if !backoff_or_stop(
1192 backoff,
1193 error,
1194 &mut commands,
1195 &mut pipe_commands,
1196 &mut stop,
1197 )
1198 .await
1199 {
1200 break;
1201 }
1202 backoff = next_backoff(backoff, config.reconnect_max_backoff);
1203 log_peer_reconnecting(&local_node, &member.node_id, backoff);
1204 }
1205 },
1206 Err(error) => {
1207 log_peer_connect_failed(&local_node, &member.node_id, &error);
1208 let _ = state.send(PeerSessionState::BackingOff(error.clone()));
1209 if !backoff_or_stop(backoff, error, &mut commands, &mut pipe_commands, &mut stop)
1210 .await
1211 {
1212 break;
1213 }
1214 backoff = next_backoff(backoff, config.reconnect_max_backoff);
1215 log_peer_reconnecting(&local_node, &member.node_id, backoff);
1216 }
1217 }
1218 }
1219 let _ = state.send(PeerSessionState::Closed);
1220 eprintln!(
1221 "datum-agent INFO peer_session_closed node_id={local_node} peer={}",
1222 member.node_id
1223 );
1224}
1225
1226fn log_peer_connect_failed(local_node: &str, peer: &str, error: &str) {
1227 eprintln!(
1228 "datum-agent WARN peer_session_connect_failed node_id={local_node} peer={peer} error={error}"
1229 );
1230}
1231
1232fn log_peer_reconnecting(local_node: &str, peer: &str, backoff: Duration) {
1233 eprintln!(
1234 "datum-agent INFO peer_session_reconnecting node_id={local_node} peer={peer} backoff_ms={}",
1235 backoff.as_millis()
1236 );
1237}
1238
1239async fn connect_peer(
1240 transport: &NodeSessionTransport,
1241 local_node: &str,
1242 member: &Member,
1243) -> Result<DcpClient, String> {
1244 let addr = member
1245 .agent_addr
1246 .ok_or_else(|| "member did not advertise a DCP agent address".to_owned())?;
1247 let mut hello = Hello::new(local_node.to_owned(), ClientKind::ClusterNode);
1248 hello.capabilities.push("cluster-node-sessions".to_owned());
1249 match transport {
1250 NodeSessionTransport::TcpLoopback => {
1251 if !addr.ip().is_loopback() {
1252 return Err(format!(
1253 "plaintext node-session TCP requires loopback, got {addr}"
1254 ));
1255 }
1256 DcpClient::connect_tcp(addr, hello)
1257 .await
1258 .map_err(|error| error.to_string())
1259 }
1260 NodeSessionTransport::QuicMtls {
1261 server_name,
1262 client_config,
1263 } => DcpClient::connect_quic(addr, server_name, client_config.clone(), hello)
1264 .await
1265 .map_err(|error| error.to_string()),
1266 }
1267}
1268
1269async fn connect_shard_pipe(
1270 transport: &NodeSessionTransport,
1271 local_node: &str,
1272 member: &Member,
1273) -> Result<ShardPipeClient, String> {
1274 let addr = member
1275 .agent_addr
1276 .ok_or_else(|| "member did not advertise a DCP agent address".to_owned())?;
1277 let mut hello = Hello::new(local_node.to_owned(), ClientKind::ClusterNode);
1278 hello.capabilities.push("cluster-node-sessions".to_owned());
1279 hello.capabilities.push("cluster-sharding-pipe".to_owned());
1280 match transport {
1281 NodeSessionTransport::TcpLoopback => {
1282 if !addr.ip().is_loopback() {
1283 return Err(format!(
1284 "plaintext node-session TCP requires loopback, got {addr}"
1285 ));
1286 }
1287 ShardPipeClient::connect_tcp(addr, hello)
1288 .await
1289 .map_err(|error| error.to_string())
1290 }
1291 NodeSessionTransport::QuicMtls {
1292 server_name,
1293 client_config,
1294 } => ShardPipeClient::connect_quic(addr, server_name, client_config.clone(), hello)
1295 .await
1296 .map_err(|error| error.to_string()),
1297 }
1298}
1299
1300async fn run_connected_peer(
1301 client: DcpClient,
1302 pipe: ShardPipeClient,
1303 commands: &mut mpsc::Receiver<SessionCommand>,
1304 pipe_commands: &mut mpsc::Receiver<ShardPipeCommand>,
1305 stop: &mut watch::Receiver<bool>,
1306) -> bool {
1307 loop {
1308 tokio::select! {
1309 changed = stop.changed() => {
1310 return changed.is_ok() && !*stop.borrow();
1311 }
1312 command = pipe_commands.recv() => {
1313 let Some(command) = command else {
1314 return false;
1315 };
1316 let frame = build_shard_pipe_frame(command, pipe_commands);
1317 if pipe.send(frame).await.is_err() {
1318 return true;
1319 }
1320 }
1321 command = commands.recv() => {
1322 let Some(command) = command else {
1323 return false;
1324 };
1325 match command {
1326 SessionCommand::ListJobs { timeout, reply } => {
1327 let (result, reconnect) = classify_peer_request(
1328 tokio::time::timeout(timeout, client.list_jobs()).await,
1329 "ListJobs",
1330 );
1331 let _ = reply.send(result);
1332 if reconnect {
1333 return true;
1334 }
1335 }
1336 SessionCommand::SubmitClusterJob {
1337 request,
1338 timeout,
1339 reply,
1340 } => {
1341 let timeout_ms = duration_millis_u64(timeout);
1342 let result = tokio::time::timeout(
1343 timeout,
1344 client.submit_cluster_job(
1345 request.factory_name,
1346 request.instance_name,
1347 request.params,
1348 request.placement.unwrap_or_else(default_wire_placement),
1349 timeout_ms,
1350 ),
1351 )
1352 .await;
1353 let (result, reconnect) =
1354 classify_peer_request(result, "SubmitClusterJob");
1355 let _ = reply.send(result);
1356 if reconnect {
1357 return true;
1358 }
1359 }
1360 SessionCommand::StartClusterJob {
1361 factory_name,
1362 instance_name,
1363 params,
1364 assignment,
1365 timeout,
1366 reply,
1367 } => {
1368 let result = tokio::time::timeout(
1369 timeout,
1370 client.start_cluster_job_on_node(
1371 factory_name,
1372 instance_name,
1373 params,
1374 assignment,
1375 ),
1376 )
1377 .await;
1378 let (result, reconnect) = classify_peer_request(result, "StartJob");
1379 let _ = reply.send(result);
1380 if reconnect {
1381 return true;
1382 }
1383 }
1384 SessionCommand::ClusterJobStatus {
1385 name,
1386 cluster,
1387 timeout,
1388 reply,
1389 } => {
1390 let request = async {
1391 if cluster {
1392 client
1393 .cluster_job_status(name, duration_millis_u64(timeout))
1394 .await
1395 } else {
1396 client.job_status(name).await
1397 }
1398 };
1399 let (result, reconnect) = classify_peer_request(
1400 tokio::time::timeout(timeout, request).await,
1401 "JobStatus",
1402 );
1403 let _ = reply.send(result);
1404 if reconnect {
1405 return true;
1406 }
1407 }
1408 SessionCommand::DrainClusterJob {
1409 name,
1410 cluster,
1411 timeout,
1412 reply,
1413 } => {
1414 let request = async {
1415 if cluster {
1416 client
1417 .drain_cluster_job(name, duration_millis_u64(timeout))
1418 .await
1419 } else {
1420 client.drain_job(name).await
1421 }
1422 };
1423 let (result, reconnect) = classify_peer_request(
1424 tokio::time::timeout(timeout, request).await,
1425 "DrainJob",
1426 );
1427 let _ = reply.send(result);
1428 if reconnect {
1429 return true;
1430 }
1431 }
1432 SessionCommand::StopClusterJob {
1433 name,
1434 cluster,
1435 timeout,
1436 reply,
1437 } => {
1438 let request = async {
1439 if cluster {
1440 client
1441 .stop_cluster_job(name, duration_millis_u64(timeout))
1442 .await
1443 } else {
1444 client.stop_job(name).await
1445 }
1446 };
1447 let (result, reconnect) = classify_peer_request(
1448 tokio::time::timeout(timeout, request).await,
1449 "StopJob",
1450 );
1451 let _ = reply.send(result);
1452 if reconnect {
1453 return true;
1454 }
1455 }
1456 SessionCommand::ListClusterJobs { timeout, reply } => {
1457 let result = tokio::time::timeout(
1458 timeout,
1459 client.list_cluster_jobs(duration_millis_u64(timeout)),
1460 )
1461 .await;
1462 let (result, reconnect) =
1463 classify_peer_request(result, "ListClusterJobs");
1464 let _ = reply.send(result);
1465 if reconnect {
1466 return true;
1467 }
1468 }
1469 SessionCommand::SubscribeMetrics {
1470 interval,
1471 job_names,
1472 timeout,
1473 reply,
1474 } => {
1475 let result = tokio::time::timeout(
1476 timeout,
1477 client.subscribe_local_metrics(
1478 duration_millis_u64(interval),
1479 job_names,
1480 ),
1481 )
1482 .await;
1483 let (result, reconnect) =
1484 classify_peer_request(result, "SubscribeMetrics");
1485 let _ = reply.send(result);
1486 if reconnect {
1487 return true;
1488 }
1489 }
1490 SessionCommand::RememberClusterAssignment {
1491 instance_name,
1492 assignment,
1493 tombstone,
1494 timeout,
1495 reply,
1496 } => {
1497 let (tombstone_generation, tombstone_coordinator) = tombstone
1498 .map(|tombstone| {
1499 (tombstone.placement_generation, tombstone.coordinator_node)
1500 })
1501 .unwrap_or_default();
1502 let result = tokio::time::timeout(
1503 timeout,
1504 client.remember_cluster_assignment(
1505 instance_name,
1506 assignment,
1507 tombstone_generation,
1508 tombstone_coordinator,
1509 ),
1510 )
1511 .await;
1512 let (result, reconnect) =
1513 classify_peer_request(result, "RememberClusterAssignment");
1514 let _ = reply.send(result);
1515 if reconnect {
1516 return true;
1517 }
1518 }
1519 SessionCommand::AllocateShard {
1520 type_name,
1521 shard_id,
1522 timeout,
1523 reply,
1524 } => {
1525 let timeout_ms = duration_millis_u64(timeout);
1526 let result = tokio::time::timeout(
1527 timeout,
1528 client.allocate_shard(type_name, shard_id, timeout_ms),
1529 )
1530 .await;
1531 let (result, reconnect) = classify_peer_request(result, "AllocateShard");
1532 let _ = reply.send(result);
1533 if reconnect {
1534 return true;
1535 }
1536 }
1537 SessionCommand::RememberShardAllocations {
1538 table,
1539 timeout,
1540 reply,
1541 } => {
1542 let result = tokio::time::timeout(
1543 timeout,
1544 client.remember_shard_allocations(table),
1545 )
1546 .await;
1547 let (result, reconnect) =
1548 classify_peer_request(result, "RememberShardAllocations");
1549 let _ = reply.send(result);
1550 if reconnect {
1551 return true;
1552 }
1553 }
1554 SessionCommand::GetShardAllocations {
1555 type_name,
1556 timeout,
1557 reply,
1558 } => {
1559 let timeout_ms = duration_millis_u64(timeout);
1560 let result = tokio::time::timeout(
1561 timeout,
1562 client.get_shard_allocations(type_name, timeout_ms),
1563 )
1564 .await;
1565 let (result, reconnect) =
1566 classify_peer_request(result, "GetShardAllocations");
1567 let _ = reply.send(result);
1568 if reconnect {
1569 return true;
1570 }
1571 }
1572 SessionCommand::ForwardShardEnvelopes {
1573 batch,
1574 timeout,
1575 reply,
1576 } => {
1577 let timeout_ms = duration_millis_u64(timeout);
1578 let result = tokio::time::timeout(
1579 timeout,
1580 client.forward_shard_envelopes(batch, timeout_ms),
1581 )
1582 .await;
1583 let (result, reconnect) =
1584 classify_peer_request(result, "ForwardShardEnvelopes");
1585 let _ = reply.send(result);
1586 if reconnect {
1587 return true;
1588 }
1589 }
1590 SessionCommand::CompleteShardingAsk {
1591 response,
1592 timeout,
1593 reply,
1594 } => {
1595 let result =
1596 tokio::time::timeout(timeout, client.complete_sharding_ask(response))
1597 .await;
1598 let (result, reconnect) =
1599 classify_peer_request(result, "CompleteShardingAsk");
1600 let _ = reply.send(result);
1601 if reconnect {
1602 return true;
1603 }
1604 }
1605 }
1606 }
1607 }
1608 }
1609}
1610
1611fn classify_peer_request<T>(
1612 result: Result<Result<T, DcpError>, tokio::time::error::Elapsed>,
1613 operation: &str,
1614) -> (Result<T, String>, bool) {
1615 match result {
1616 Ok(Ok(value)) => (Ok(value), false),
1617 Err(_) => (Err(format!("peer {operation} timed out")), false),
1618 Ok(Err(error)) => {
1619 let reconnect = dcp_error_requires_reconnect(&error);
1620 (Err(error.to_string()), reconnect)
1621 }
1622 }
1623}
1624
1625fn dcp_error_requires_reconnect(error: &DcpError) -> bool {
1626 match error {
1627 DcpError::Closed | DcpError::Protocol(_) | DcpError::Io(_) | DcpError::Decode(_) => true,
1628 DcpError::Response { .. }
1629 | DcpError::Encode(_)
1630 | DcpError::Agent(_)
1631 | DcpError::Stream(_)
1632 | DcpError::Join(_) => false,
1633 }
1634}
1635
1636const SHARD_PIPE_MAX_BATCH_COMMANDS: usize = 1024;
1637
1638fn build_shard_pipe_frame(
1639 first: ShardPipeCommand,
1640 pipe_commands: &mut mpsc::Receiver<ShardPipeCommand>,
1641) -> ShardPipeFrame {
1642 let mut frame = ShardPipeFrame {
1643 forwards: Vec::new(),
1644 replies: Vec::new(),
1645 };
1646 push_shard_pipe_command(&mut frame, first);
1647 for _ in 1..SHARD_PIPE_MAX_BATCH_COMMANDS {
1648 match pipe_commands.try_recv() {
1649 Ok(command) => push_shard_pipe_command(&mut frame, command),
1650 Err(_) => break,
1651 }
1652 }
1653 frame
1654}
1655
1656fn push_shard_pipe_command(frame: &mut ShardPipeFrame, command: ShardPipeCommand) {
1657 match command {
1658 ShardPipeCommand::Forward { batch } => {
1659 if let Some(existing) = frame
1660 .forwards
1661 .iter_mut()
1662 .find(|existing| existing.type_name == batch.type_name)
1663 {
1664 existing.envelopes.extend(batch.envelopes);
1665 } else {
1666 frame.forwards.push(batch);
1667 }
1668 }
1669 ShardPipeCommand::Reply { response } => {
1670 frame.replies.push(response);
1671 }
1672 }
1673}
1674
1675async fn backoff_or_stop(
1676 backoff: Duration,
1677 error: String,
1678 commands: &mut mpsc::Receiver<SessionCommand>,
1679 pipe_commands: &mut mpsc::Receiver<ShardPipeCommand>,
1680 stop: &mut watch::Receiver<bool>,
1681) -> bool {
1682 let sleep = sleep(backoff);
1683 tokio::pin!(sleep);
1684 loop {
1685 tokio::select! {
1686 () = &mut sleep => return true,
1687 changed = stop.changed() => {
1688 return changed.is_ok() && !*stop.borrow();
1689 }
1690 command = pipe_commands.recv() => {
1691 if command.is_none() {
1692 return false;
1693 }
1694 }
1695 command = commands.recv() => {
1696 let Some(command) = command else {
1697 return false;
1698 };
1699 match command {
1700 SessionCommand::ListJobs { reply, .. } => {
1701 let _ = reply.send(Err(format!("node session unreachable: {error}")));
1702 }
1703 SessionCommand::SubmitClusterJob { reply, .. }
1704 | SessionCommand::StartClusterJob { reply, .. }
1705 | SessionCommand::ClusterJobStatus { reply, .. }
1706 | SessionCommand::DrainClusterJob { reply, .. }
1707 | SessionCommand::StopClusterJob { reply, .. } => {
1708 let _ = reply.send(Err(format!("node session unreachable: {error}")));
1709 }
1710 SessionCommand::ListClusterJobs { reply, .. } => {
1711 let _ = reply.send(Err(format!("node session unreachable: {error}")));
1712 }
1713 SessionCommand::SubscribeMetrics { reply, .. } => {
1714 let _ = reply.send(Err(format!("node session unreachable: {error}")));
1715 }
1716 SessionCommand::RememberClusterAssignment { reply, .. } => {
1717 let _ = reply.send(Err(format!("node session unreachable: {error}")));
1718 }
1719 SessionCommand::AllocateShard { reply, .. } => {
1720 let _ = reply.send(Err(format!("node session unreachable: {error}")));
1721 }
1722 SessionCommand::RememberShardAllocations { reply, .. } => {
1723 let _ = reply.send(Err(format!("node session unreachable: {error}")));
1724 }
1725 SessionCommand::GetShardAllocations { reply, .. } => {
1726 let _ = reply.send(Err(format!("node session unreachable: {error}")));
1727 }
1728 SessionCommand::ForwardShardEnvelopes { reply, .. } => {
1729 let _ = reply.send(Err(format!("node session unreachable: {error}")));
1730 }
1731 SessionCommand::CompleteShardingAsk { reply, .. } => {
1732 let _ = reply.send(Err(format!("node session unreachable: {error}")));
1733 }
1734 }
1735 }
1736 }
1737 }
1738}
1739
1740fn default_wire_placement() -> PlacementSpec {
1741 PlacementSpec {
1742 role_constraint: String::new(),
1743 strategy: PlacementStrategy::LeastJobs as i32,
1744 pinned_node_id: String::new(),
1745 }
1746}
1747
1748fn duration_millis_u64(duration: Duration) -> u64 {
1749 duration.as_millis().min(u128::from(u64::MAX)) as u64
1750}
1751
1752fn next_backoff(current: Duration, max: Duration) -> Duration {
1753 current.saturating_mul(2).min(max)
1754}
1755
1756#[derive(Clone)]
1757struct PlacementCoordinatorHandle {
1758 commands: mpsc::Sender<PlacementCommand>,
1759 tasks: Arc<Mutex<Vec<JoinHandle<()>>>>,
1760}
1761
1762enum PlacementCommand {
1763 Submit {
1764 request: SubmitClusterJob,
1765 timeout: Duration,
1766 reply: oneshot::Sender<Result<crate::dcp::proto::JobStatus, DcpError>>,
1767 },
1768 Status {
1769 name: String,
1770 timeout: Duration,
1771 reply: oneshot::Sender<Result<crate::dcp::proto::JobStatus, DcpError>>,
1772 },
1773 Drain {
1774 name: String,
1775 timeout: Duration,
1776 reply: oneshot::Sender<Result<crate::dcp::proto::JobStatus, DcpError>>,
1777 },
1778 Stop {
1779 name: String,
1780 timeout: Duration,
1781 reply: oneshot::Sender<Result<crate::dcp::proto::JobStatus, DcpError>>,
1782 },
1783 Remember {
1784 instance_name: String,
1785 assignment: Option<ClusterJobStart>,
1786 tombstone: Option<AssignmentTombstone>,
1787 reply: oneshot::Sender<Result<(), DcpError>>,
1788 },
1789 RegisterRestarted {
1790 instance_name: String,
1791 assignment: ClusterJobStart,
1792 timeout: Duration,
1793 reply: oneshot::Sender<Result<(), DcpError>>,
1794 },
1795 AssignmentSyncCompleted {
1796 node_id: String,
1797 incarnation: u64,
1798 instance_name: String,
1799 sync_mark: AssignmentSyncMark,
1800 ok: bool,
1801 },
1802 Tick,
1803 Shutdown,
1804}
1805
1806struct PlacementActor {
1807 commands: mpsc::Sender<PlacementCommand>,
1808 self_node: String,
1809 agent_role: String,
1810 state: Signal<ClusterState>,
1811 registry: JobRegistryHandle,
1812 sessions: NodeSessionManagerHandle,
1813 factories: DcpJobFactories,
1814 request_timeout: Duration,
1815 active_coordinator: bool,
1816 last_known_coordinator: Option<String>,
1817 assignments: BTreeMap<String, ClusterJobMetadata>,
1818 assignment_tombstones: BTreeMap<String, AssignmentTombstone>,
1819 replicated_assignments: BTreeMap<String, PeerAssignmentSync>,
1820 inflight_assignment_syncs: BTreeSet<AssignmentSyncKey>,
1821 pending_replacements: BTreeMap<String, ClusterJobMetadata>,
1822 assignment_updates: watch::Sender<Arc<BTreeMap<String, String>>>,
1823 cluster_events: ClusterEventPublisher,
1824}
1825
1826#[derive(Clone, Debug, PartialEq, Eq)]
1827struct AssignmentTombstone {
1828 placement_generation: u64,
1829 coordinator_node: String,
1830}
1831
1832impl AssignmentTombstone {
1833 fn from_metadata(metadata: &ClusterJobMetadata) -> Self {
1834 Self {
1835 placement_generation: metadata.placement_generation,
1836 coordinator_node: metadata.coordinator_node.clone(),
1837 }
1838 }
1839
1840 fn from_wire(request: &RememberClusterAssignment) -> Self {
1841 Self {
1842 placement_generation: request.tombstone_placement_generation,
1843 coordinator_node: request.tombstone_coordinator_node_id.clone(),
1844 }
1845 }
1846}
1847
1848#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord)]
1849enum AssignmentSyncMark {
1850 Active {
1851 placement_generation: u64,
1852 assigned_node: String,
1853 coordinator_node: String,
1854 },
1855 Tombstone {
1856 placement_generation: u64,
1857 coordinator_node: String,
1858 },
1859}
1860
1861impl AssignmentSyncMark {
1862 fn active(metadata: &ClusterJobMetadata) -> Self {
1863 Self::Active {
1864 placement_generation: metadata.placement_generation,
1865 assigned_node: metadata.assigned_node.clone(),
1866 coordinator_node: metadata.coordinator_node.clone(),
1867 }
1868 }
1869
1870 fn tombstone(tombstone: &AssignmentTombstone) -> Self {
1871 Self::Tombstone {
1872 placement_generation: tombstone.placement_generation,
1873 coordinator_node: tombstone.coordinator_node.clone(),
1874 }
1875 }
1876}
1877
1878#[derive(Clone, Debug, PartialEq, Eq)]
1879struct PeerAssignmentSync {
1880 incarnation: u64,
1881 assignments: BTreeMap<String, AssignmentSyncMark>,
1882}
1883
1884#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord)]
1885struct AssignmentSyncKey {
1886 node_id: String,
1887 incarnation: u64,
1888 instance_name: String,
1889 sync_mark: AssignmentSyncMark,
1890}
1891
1892struct PlacementDependencies {
1893 registry: JobRegistryHandle,
1894 sessions: NodeSessionManagerHandle,
1895 factories: DcpJobFactories,
1896 cluster_events: ClusterEventPublisher,
1897}
1898
1899impl PlacementCoordinatorHandle {
1900 fn start(
1901 self_node: String,
1902 agent_role: String,
1903 state: Signal<ClusterState>,
1904 dependencies: PlacementDependencies,
1905 request_timeout: Duration,
1906 ) -> ClusterAgentResult<Self> {
1907 let (commands, receiver) = mpsc::channel(128);
1908 let (assignment_updates, _assignments) = watch::channel(Arc::new(BTreeMap::new()));
1909 let tasks = Arc::new(Mutex::new(Vec::new()));
1910 let actor = PlacementActor {
1911 commands: commands.clone(),
1912 self_node,
1913 agent_role,
1914 state,
1915 registry: dependencies.registry,
1916 sessions: dependencies.sessions,
1917 factories: dependencies.factories,
1918 request_timeout,
1919 active_coordinator: false,
1920 last_known_coordinator: None,
1921 assignments: BTreeMap::new(),
1922 assignment_tombstones: BTreeMap::new(),
1923 replicated_assignments: BTreeMap::new(),
1924 inflight_assignment_syncs: BTreeSet::new(),
1925 pending_replacements: BTreeMap::new(),
1926 assignment_updates,
1927 cluster_events: dependencies.cluster_events,
1928 };
1929 let actor_task = tokio::spawn(run_placement_actor(actor, receiver));
1930 let tick_commands = commands.clone();
1931 let tick_task = tokio::spawn(async move {
1932 let mut interval = tokio::time::interval(Duration::from_millis(50));
1933 loop {
1934 interval.tick().await;
1935 if tick_commands.send(PlacementCommand::Tick).await.is_err() {
1936 break;
1937 }
1938 }
1939 });
1940 {
1941 let mut locked = tasks.lock().expect("placement tasks poisoned");
1942 locked.push(actor_task);
1943 locked.push(tick_task);
1944 }
1945 Ok(Self { commands, tasks })
1946 }
1947
1948 async fn submit_cluster_job(
1949 &self,
1950 request: SubmitClusterJob,
1951 timeout: Duration,
1952 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
1953 self.request(timeout, |reply| PlacementCommand::Submit {
1954 request,
1955 timeout,
1956 reply,
1957 })
1958 .await
1959 }
1960
1961 async fn cluster_job_status(
1962 &self,
1963 name: String,
1964 timeout: Duration,
1965 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
1966 self.request(timeout, |reply| PlacementCommand::Status {
1967 name,
1968 timeout,
1969 reply,
1970 })
1971 .await
1972 }
1973
1974 async fn drain_cluster_job(
1975 &self,
1976 name: String,
1977 timeout: Duration,
1978 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
1979 self.request(timeout, |reply| PlacementCommand::Drain {
1980 name,
1981 timeout,
1982 reply,
1983 })
1984 .await
1985 }
1986
1987 async fn stop_cluster_job(
1988 &self,
1989 name: String,
1990 timeout: Duration,
1991 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
1992 self.request(timeout, |reply| PlacementCommand::Stop {
1993 name,
1994 timeout,
1995 reply,
1996 })
1997 .await
1998 }
1999
2000 async fn remember_cluster_assignment(
2001 &self,
2002 instance_name: String,
2003 assignment: Option<ClusterJobStart>,
2004 tombstone: Option<AssignmentTombstone>,
2005 ) -> Result<(), DcpError> {
2006 self.request(Duration::from_millis(250), |reply| {
2007 PlacementCommand::Remember {
2008 instance_name,
2009 assignment,
2010 tombstone,
2011 reply,
2012 }
2013 })
2014 .await
2015 }
2016
2017 async fn register_restarted_cluster_assignment(
2018 &self,
2019 instance_name: String,
2020 assignment: ClusterJobStart,
2021 timeout: Duration,
2022 ) -> Result<(), DcpError> {
2023 self.request(timeout, |reply| PlacementCommand::RegisterRestarted {
2024 instance_name,
2025 assignment,
2026 timeout,
2027 reply,
2028 })
2029 .await
2030 }
2031
2032 async fn shutdown(&self) {
2033 let _ = self.commands.send(PlacementCommand::Shutdown).await;
2034 let tasks = {
2035 let mut locked = self.tasks.lock().expect("placement tasks poisoned");
2036 std::mem::take(&mut *locked)
2037 };
2038 for task in tasks {
2039 task.abort();
2040 }
2041 }
2042
2043 async fn request<T, F>(&self, timeout: Duration, make: F) -> Result<T, DcpError>
2044 where
2045 T: Send + 'static,
2046 F: FnOnce(oneshot::Sender<Result<T, DcpError>>) -> PlacementCommand,
2047 {
2048 let (reply, receiver) = oneshot::channel();
2049 self.commands
2050 .send(make(reply))
2051 .await
2052 .map_err(|_| DcpError::response(ResponseStatus::Failed, "placement actor stopped"))?;
2053 tokio::time::timeout(timeout.saturating_add(Duration::from_millis(250)), receiver)
2054 .await
2055 .map_err(|_| {
2056 DcpError::response(
2057 ResponseStatus::DeadlineExceeded,
2058 "placement request timed out",
2059 )
2060 })?
2061 .map_err(|_| DcpError::response(ResponseStatus::Failed, "placement actor stopped"))?
2062 }
2063}
2064
2065async fn run_placement_actor(
2066 mut actor: PlacementActor,
2067 mut receiver: mpsc::Receiver<PlacementCommand>,
2068) {
2069 while let Some(command) = receiver.recv().await {
2070 match command {
2071 PlacementCommand::Submit {
2072 request,
2073 timeout,
2074 reply,
2075 } => {
2076 let result = actor.submit(request, timeout).await;
2077 let _ = reply.send(result);
2078 }
2079 PlacementCommand::Status {
2080 name,
2081 timeout,
2082 reply,
2083 } => {
2084 let result = actor.cluster_job_status(name, timeout).await;
2085 let _ = reply.send(result);
2086 }
2087 PlacementCommand::Drain {
2088 name,
2089 timeout,
2090 reply,
2091 } => {
2092 let result = actor.drain_cluster_job(name, timeout).await;
2093 let _ = reply.send(result);
2094 }
2095 PlacementCommand::Stop {
2096 name,
2097 timeout,
2098 reply,
2099 } => {
2100 let result = actor.stop_cluster_job(name, timeout).await;
2101 let _ = reply.send(result);
2102 }
2103 PlacementCommand::Remember {
2104 instance_name,
2105 assignment,
2106 tombstone,
2107 reply,
2108 } => {
2109 let result = actor.remember(instance_name, assignment, tombstone).await;
2110 let _ = reply.send(result);
2111 }
2112 PlacementCommand::RegisterRestarted {
2113 instance_name,
2114 assignment,
2115 timeout,
2116 reply,
2117 } => {
2118 let result = actor
2119 .register_restarted(instance_name, assignment, timeout)
2120 .await;
2121 let _ = reply.send(result);
2122 }
2123 PlacementCommand::AssignmentSyncCompleted {
2124 node_id,
2125 incarnation,
2126 instance_name,
2127 sync_mark,
2128 ok,
2129 } => {
2130 actor.assignment_sync_completed(node_id, incarnation, instance_name, sync_mark, ok);
2131 }
2132 PlacementCommand::Tick => {
2133 let _ = actor.reconcile().await;
2134 }
2135 PlacementCommand::Shutdown => break,
2136 }
2137 }
2138}
2139
2140impl PlacementActor {
2141 async fn submit(
2142 &mut self,
2143 request: SubmitClusterJob,
2144 timeout: Duration,
2145 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
2146 self.ensure_active_coordinator(timeout).await?;
2147 validate_submit(&request)?;
2148 if self.assignments.contains_key(&request.instance_name) {
2149 return Err(DcpError::response(
2150 ResponseStatus::Conflict,
2151 format!("cluster job already exists: {}", request.instance_name),
2152 ));
2153 }
2154
2155 let placement = placement_spec_from_wire(request.placement.clone())?;
2156 let target = self.choose_submit_target(&placement)?;
2157 let history = ClusterPlacementHistory {
2158 generation: 1,
2159 from_node: None,
2160 to_node: target.clone(),
2161 reason: "submitted".to_owned(),
2162 timestamp: SystemTime::now(),
2163 };
2164 let metadata = ClusterJobMetadata {
2165 factory_name: request.factory_name.clone(),
2166 params: request.params.clone().into_iter().collect(),
2167 placement,
2168 coordinator_node: self.self_node.clone(),
2169 assigned_node: target.clone(),
2170 placement_generation: 1,
2171 history: vec![history],
2172 };
2173 let status = self
2174 .start_on_node(
2175 &target,
2176 request.factory_name,
2177 request.instance_name.clone(),
2178 request.params,
2179 metadata.clone(),
2180 timeout,
2181 )
2182 .await?;
2183 self.assignments
2184 .insert(request.instance_name.clone(), metadata.clone());
2185 self.assignment_tombstones.remove(&request.instance_name);
2186 self.pending_replacements.remove(&request.instance_name);
2187 self.publish_assignments();
2188 self.publish_assignment(request.instance_name, metadata, timeout)
2189 .await?;
2190 Ok(status)
2191 }
2192
2193 async fn cluster_job_status(
2194 &mut self,
2195 name: String,
2196 timeout: Duration,
2197 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
2198 self.ensure_active_coordinator(timeout).await?;
2199 let metadata = self.assignment_or_inactive(&name, timeout).await?;
2200 let mut status = self
2201 .status_on_node(&metadata.assigned_node, name, timeout)
2202 .await?;
2203 status.coordinator_node_id.clone_from(&self.self_node);
2204 Ok(status)
2205 }
2206
2207 async fn drain_cluster_job(
2208 &mut self,
2209 name: String,
2210 timeout: Duration,
2211 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
2212 self.ensure_active_coordinator(timeout).await?;
2213 let metadata = self.assignment_or_rebuild(&name, timeout).await?;
2214 let status = self
2215 .drain_on_node(&metadata.assigned_node, name.clone(), timeout)
2216 .await?;
2217 if let Some(removed) = self.assignments.remove(&name) {
2218 self.record_assignment_tombstone(&name, &removed);
2219 }
2220 self.pending_replacements.remove(&name);
2221 self.publish_assignments();
2222 self.replicate_assignment_mutation(name, None, timeout);
2223 Ok(status)
2224 }
2225
2226 async fn stop_cluster_job(
2227 &mut self,
2228 name: String,
2229 timeout: Duration,
2230 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
2231 self.ensure_active_coordinator(timeout).await?;
2232 let metadata = self.assignment_or_rebuild(&name, timeout).await?;
2233 let status = self
2234 .stop_on_node(&metadata.assigned_node, name.clone(), timeout)
2235 .await?;
2236 if let Some(removed) = self.assignments.remove(&name) {
2237 self.record_assignment_tombstone(&name, &removed);
2238 }
2239 self.pending_replacements.remove(&name);
2240 self.publish_assignments();
2241 self.replicate_assignment_mutation(name, None, timeout);
2242 Ok(status)
2243 }
2244
2245 async fn remember(
2246 &mut self,
2247 instance_name: String,
2248 assignment: Option<ClusterJobStart>,
2249 tombstone: Option<AssignmentTombstone>,
2250 ) -> Result<(), DcpError> {
2251 let current_coordinator = self.current_placement_coordinator();
2252 let Some(assignment) = assignment else {
2253 let tombstone = tombstone.unwrap_or(AssignmentTombstone {
2254 placement_generation: 0,
2255 coordinator_node: String::new(),
2256 });
2257 ensure_mutation_from_current_coordinator(
2258 &tombstone.coordinator_node,
2259 current_coordinator.as_deref(),
2260 )?;
2261 if let Some(existing) = self.assignments.get(&instance_name)
2262 && !tombstone_dominates_assignment(
2263 &tombstone,
2264 existing,
2265 current_coordinator.as_deref(),
2266 )
2267 {
2268 return Ok(());
2269 }
2270 self.record_assignment_tombstone_value(
2271 &instance_name,
2272 tombstone.clone(),
2273 current_coordinator.as_deref(),
2274 );
2275 self.stop_local_cluster_assignment(
2276 &instance_name,
2277 &tombstone,
2278 current_coordinator.as_deref(),
2279 )
2280 .await?;
2281 let removed = self.assignments.remove(&instance_name);
2282 if removed.is_some() {
2283 self.publish_assignments();
2284 }
2285 self.pending_replacements.remove(&instance_name);
2286 return Ok(());
2287 };
2288 let metadata = cluster_metadata_from_wire(assignment)?;
2289 ensure_mutation_from_current_coordinator(
2290 &metadata.coordinator_node,
2291 current_coordinator.as_deref(),
2292 )?;
2293 if let Some(tombstone) = self.assignment_tombstones.get(&instance_name)
2294 && !active_assignment_supersedes_tombstone(&metadata, tombstone)
2295 {
2296 return Err(stale_assignment_conflict(
2297 &instance_name,
2298 metadata.placement_generation,
2299 tombstone,
2300 ));
2301 }
2302 if !remember_assignment(
2303 &mut self.assignments,
2304 instance_name.clone(),
2305 metadata.clone(),
2306 current_coordinator.as_deref(),
2307 ) {
2308 return Ok(());
2309 }
2310 self.assignment_tombstones.remove(&instance_name);
2311 self.pending_replacements.remove(&instance_name);
2312 self.publish_assignments();
2313 self.reconcile_local_assignment(&instance_name, &metadata)
2314 .await?;
2315 Ok(())
2316 }
2317
2318 async fn register_restarted(
2319 &mut self,
2320 instance_name: String,
2321 assignment: ClusterJobStart,
2322 timeout: Duration,
2323 ) -> Result<(), DcpError> {
2324 let mut metadata = cluster_metadata_from_wire(assignment)?;
2325 if metadata.assigned_node != self.self_node {
2326 return Err(DcpError::response(
2327 ResponseStatus::Conflict,
2328 "only the owning node can re-register restarted cluster intent",
2329 ));
2330 }
2331 if let Some(existing) = self.assignments.get(&instance_name).cloned()
2332 && active_assignment_blocks_local_restart(&metadata, &existing)
2333 {
2334 self.stop_stale_restarted_local_assignment(&instance_name, &existing)
2335 .await?;
2336 return Err(stale_restart_conflict(&instance_name, &metadata, &existing));
2337 }
2338 if let Some(tombstone) = self.assignment_tombstones.get(&instance_name)
2339 && metadata.placement_generation <= tombstone.placement_generation
2340 {
2341 metadata.placement_generation = tombstone.placement_generation.saturating_add(1);
2342 metadata.history.push(ClusterPlacementHistory {
2343 generation: metadata.placement_generation,
2344 from_node: Some(self.self_node.clone()),
2345 to_node: self.self_node.clone(),
2346 reason: "local_restart_after_terminal_intent".to_owned(),
2347 timestamp: SystemTime::now(),
2348 });
2349 }
2350 let snapshot = self.state.get();
2351 metadata.coordinator_node = snapshot
2352 .placement_coordinator()
2353 .map(|member| member.node_id.clone())
2354 .ok_or_else(|| {
2355 DcpError::response(
2356 ResponseStatus::Failed,
2357 "no placement coordinator is available",
2358 )
2359 })?;
2360 drop(snapshot);
2361 self.assignments
2362 .insert(instance_name.clone(), metadata.clone());
2363 self.assignment_tombstones.remove(&instance_name);
2364 self.pending_replacements.remove(&instance_name);
2365 self.publish_assignments();
2366 self.publish_assignment(instance_name, metadata, timeout)
2367 .await
2368 }
2369
2370 async fn reconcile(&mut self) -> Result<(), DcpError> {
2371 let snapshot = self.state.get();
2372 let coordinator_id = snapshot
2373 .placement_coordinator()
2374 .map(|member| member.node_id.clone());
2375 let is_coordinator = snapshot.is_placement_coordinator(&self.self_node);
2376 drop(snapshot);
2377
2378 if coordinator_id != self.last_known_coordinator {
2379 self.cluster_events.publish_coordinator_changed(
2380 self.last_known_coordinator.as_deref(),
2381 coordinator_id.as_deref(),
2382 &self.self_node,
2383 );
2384 if let Some(coordinator_id) = &coordinator_id {
2385 eprintln!(
2386 "datum-agent INFO coordinator_changed node_id={} new_coordinator={coordinator_id}",
2387 self.self_node
2388 );
2389 }
2390 self.last_known_coordinator = coordinator_id;
2391 }
2392
2393 if !is_coordinator {
2394 if self.active_coordinator {
2395 eprintln!(
2396 "datum-agent INFO placement_coordinator_inactive node_id={}",
2397 self.self_node
2398 );
2399 }
2400 self.active_coordinator = false;
2401 return Ok(());
2402 }
2403 self.ensure_active_coordinator(self.request_timeout).await?;
2404 self.replace_down_members(self.request_timeout).await?;
2405 self.sync_assignments_to_peers(self.request_timeout);
2406 Ok(())
2407 }
2408
2409 async fn ensure_active_coordinator(&mut self, timeout: Duration) -> Result<(), DcpError> {
2410 let snapshot = self.state.get();
2411 if !snapshot.is_placement_coordinator(&self.self_node) {
2412 return Err(DcpError::response(
2413 ResponseStatus::Failed,
2414 "this node is not the placement coordinator",
2415 ));
2416 }
2417 drop(snapshot);
2418 if !self.active_coordinator {
2419 self.rebuild_from_registries(timeout).await?;
2420 self.active_coordinator = true;
2421 eprintln!(
2422 "datum-agent INFO placement_coordinator_active node_id={}",
2423 self.self_node
2424 );
2425 }
2426 Ok(())
2427 }
2428
2429 async fn assignment_or_rebuild(
2430 &mut self,
2431 name: &str,
2432 timeout: Duration,
2433 ) -> Result<ClusterJobMetadata, DcpError> {
2434 if let Some(metadata) = self.assignments.get(name).cloned() {
2435 return Ok(metadata);
2436 }
2437 self.rebuild_from_registries(timeout).await?;
2438 self.assignments.get(name).cloned().ok_or_else(|| {
2439 DcpError::response(
2440 ResponseStatus::NotFound,
2441 format!("cluster job not found: {name}"),
2442 )
2443 })
2444 }
2445
2446 async fn assignment_or_inactive(
2447 &mut self,
2448 name: &str,
2449 timeout: Duration,
2450 ) -> Result<ClusterJobMetadata, DcpError> {
2451 if let Some(metadata) = self.assignments.get(name).cloned() {
2452 return Ok(metadata);
2453 }
2454 let inactive = self.rebuild_from_registries(timeout).await?;
2455 self.assignments
2456 .get(name)
2457 .or_else(|| inactive.get(name))
2458 .cloned()
2459 .ok_or_else(|| {
2460 DcpError::response(
2461 ResponseStatus::NotFound,
2462 format!("cluster job not found: {name}"),
2463 )
2464 })
2465 }
2466
2467 async fn rebuild_from_registries(
2468 &mut self,
2469 timeout: Duration,
2470 ) -> Result<BTreeMap<String, ClusterJobMetadata>, DcpError> {
2471 let statuses = self.collect_job_statuses(timeout).await?;
2477 let mut running = BTreeSet::new();
2478 let mut inactive = BTreeSet::new();
2479 let mut inactive_metadata = BTreeMap::new();
2480 for status in statuses {
2481 if !status.cluster_job {
2482 continue;
2483 }
2484 if status.desired_state == "Running"
2485 && let Some(mut metadata) = metadata_from_status(&status)?
2486 {
2487 metadata.coordinator_node.clone_from(&self.self_node);
2488 if let Some(tombstone) = self.assignment_tombstones.get(&status.name)
2489 && registry_observation_is_tombstoned(tombstone, &metadata)
2490 {
2491 let retained_assignment_is_newer =
2492 self.assignments.get(&status.name).is_some_and(|existing| {
2493 active_assignment_supersedes_tombstone(existing, tombstone)
2494 });
2495 if !retained_assignment_is_newer && !running.contains(&status.name) {
2496 inactive.insert(status.name.clone());
2497 }
2498 continue;
2499 }
2500 running.insert(status.name.clone());
2501 inactive.remove(&status.name);
2502 inactive_metadata.remove(&status.name);
2503 let replace = self
2504 .assignments
2505 .get(&status.name)
2506 .map(|existing| {
2507 metadata.placement_generation > existing.placement_generation
2508 || (metadata.placement_generation == existing.placement_generation
2509 && metadata.assigned_node < existing.assigned_node)
2510 })
2511 .unwrap_or(true);
2512 if replace {
2513 self.assignment_tombstones.remove(&status.name);
2514 self.pending_replacements.remove(&status.name);
2515 self.assignments.insert(status.name.clone(), metadata);
2516 self.publish_assignments();
2517 }
2518 } else if !running.contains(&status.name) {
2519 if let Some(mut metadata) = metadata_from_status(&status)? {
2520 metadata.coordinator_node.clone_from(&self.self_node);
2521 inactive_metadata.insert(status.name.clone(), metadata);
2522 }
2523 inactive.insert(status.name);
2524 }
2525 }
2526 for name in &inactive {
2527 let removed = self.assignments.remove(name);
2528 if let Some(metadata) = inactive_metadata.get(name) {
2529 self.record_assignment_tombstone(name, metadata);
2530 } else if let Some(removed) = removed.as_ref() {
2531 self.record_assignment_tombstone(name, removed);
2532 }
2533 if removed.is_some() {
2534 self.pending_replacements.remove(name);
2535 self.publish_assignments();
2536 }
2537 }
2538 for metadata in self.assignments.values_mut() {
2539 metadata.coordinator_node.clone_from(&self.self_node);
2540 }
2541 self.publish_assignments();
2542
2543 let assignments = self
2544 .assignments
2545 .iter()
2546 .map(|(name, metadata)| (name.clone(), metadata.clone()))
2547 .collect::<Vec<_>>();
2548 for (name, metadata) in assignments {
2549 self.publish_assignment(name, metadata, timeout).await?;
2550 }
2551 for name in inactive {
2552 if let Some(metadata) = inactive_metadata.get(&name).cloned() {
2553 self.publish_assignment(name.clone(), metadata, timeout)
2554 .await?;
2555 }
2556 self.replicate_assignment_mutation(name, None, timeout);
2557 }
2558 Ok(inactive_metadata)
2559 }
2560
2561 async fn replace_down_members(&mut self, timeout: Duration) -> Result<(), DcpError> {
2562 let snapshot = self.state.get();
2563 let down_nodes = snapshot
2564 .members
2565 .values()
2566 .filter(|member| member.state == MemberState::Down)
2567 .map(|member| member.node_id.clone())
2568 .collect::<BTreeSet<_>>();
2569 drop(snapshot);
2570 if down_nodes.is_empty() {
2571 return Ok(());
2572 }
2573
2574 let to_replace = self
2575 .assignments
2576 .iter()
2577 .filter(|(_name, metadata)| down_nodes.contains(&metadata.assigned_node))
2578 .map(|(name, metadata)| (name.clone(), metadata.clone()))
2579 .collect::<Vec<_>>();
2580
2581 for (name, current) in to_replace {
2582 let metadata = if let Some(pending) = self.pending_replacements.get(&name)
2583 && pending
2584 .history
2585 .last()
2586 .and_then(|event| event.from_node.as_ref())
2587 .is_some_and(|node| down_nodes.contains(node))
2588 && self.is_replacement_target_eligible(pending)
2589 {
2590 pending.clone()
2591 } else {
2592 let mut metadata = current.clone();
2593 let old_node = metadata.assigned_node.clone();
2594 let target = self.choose_replacement_target(&metadata.placement, &old_node)?;
2595 metadata.placement_generation = metadata.placement_generation.saturating_add(1);
2596 metadata.assigned_node = target.clone();
2597 metadata.coordinator_node = self.self_node.clone();
2598 metadata.history.push(ClusterPlacementHistory {
2599 generation: metadata.placement_generation,
2600 from_node: Some(old_node.clone()),
2601 to_node: target,
2602 reason: format!("node_down:{old_node}"),
2603 timestamp: SystemTime::now(),
2604 });
2605 self.pending_replacements
2606 .insert(name.clone(), metadata.clone());
2607 metadata
2608 };
2609
2610 let params = metadata.params.clone().into_iter().collect();
2611 let factory_name = metadata.factory_name.clone();
2612 self.ensure_started_on_node(
2613 &metadata.assigned_node,
2614 factory_name,
2615 name.clone(),
2616 params,
2617 metadata.clone(),
2618 timeout,
2619 )
2620 .await?;
2621 self.assignments.insert(name.clone(), metadata.clone());
2622 self.assignment_tombstones.remove(&name);
2623 self.pending_replacements.remove(&name);
2624 self.publish_assignments();
2625 self.publish_committed_replacement(&name, &metadata);
2626 self.publish_assignment(name, metadata, timeout).await?;
2627 }
2628 Ok(())
2629 }
2630
2631 fn is_replacement_target_eligible(&self, metadata: &ClusterJobMetadata) -> bool {
2632 let snapshot = self.state.get();
2633 snapshot
2634 .member(&metadata.assigned_node)
2635 .is_some_and(|member| member.state == MemberState::Up && !member.unreachable)
2636 }
2637
2638 fn current_placement_coordinator(&self) -> Option<String> {
2639 self.state
2640 .get()
2641 .placement_coordinator()
2642 .map(|member| member.node_id.clone())
2643 }
2644
2645 fn publish_committed_replacement(&self, name: &str, metadata: &ClusterJobMetadata) {
2646 let Some(history) = metadata.history.last() else {
2647 return;
2648 };
2649 let Some(old_node) = history.from_node.as_deref() else {
2650 return;
2651 };
2652 eprintln!(
2653 "datum-agent INFO job_replaced job={name} from_node={old_node} to_node={} reason={} generation={}",
2654 metadata.assigned_node, history.reason, metadata.placement_generation
2655 );
2656 self.cluster_events.publish_job_replaced(
2657 history.timestamp,
2658 name,
2659 old_node,
2660 &metadata.assigned_node,
2661 &history.reason,
2662 metadata.placement_generation,
2663 );
2664 }
2665
2666 fn publish_assignments(&self) {
2667 let assignments = self
2668 .assignments
2669 .iter()
2670 .map(|(name, metadata)| (name.clone(), metadata.assigned_node.clone()))
2671 .collect();
2672 self.assignment_updates.send_replace(Arc::new(assignments));
2673 }
2674
2675 fn choose_submit_target(&self, placement: &RegistryPlacementSpec) -> Result<String, DcpError> {
2676 let snapshot = self.state.get();
2677 let candidates = eligible_members(
2678 &snapshot,
2679 &self.agent_role,
2680 placement.role_constraint.as_deref(),
2681 );
2682 if candidates.is_empty() {
2683 return Err(no_eligible_node_error(placement));
2684 }
2685 match &placement.strategy {
2686 RegistryPlacementStrategy::LeastJobs => self.least_loaded(candidates),
2687 RegistryPlacementStrategy::Pinned { node_id } => candidates
2688 .into_iter()
2689 .find(|member| member.node_id == *node_id)
2690 .map(|member| member.node_id.clone())
2691 .ok_or_else(|| {
2692 DcpError::response(
2693 ResponseStatus::NotFound,
2694 format!("pinned node is not eligible for placement: {node_id}"),
2695 )
2696 }),
2697 }
2698 }
2699
2700 fn choose_replacement_target(
2701 &self,
2702 placement: &RegistryPlacementSpec,
2703 old_node: &str,
2704 ) -> Result<String, DcpError> {
2705 let snapshot = self.state.get();
2706 let candidates = eligible_members(
2707 &snapshot,
2708 &self.agent_role,
2709 placement.role_constraint.as_deref(),
2710 )
2711 .into_iter()
2712 .filter(|member| member.node_id != old_node)
2713 .collect::<Vec<_>>();
2714 if candidates.is_empty() {
2715 return Err(no_eligible_node_error(placement));
2716 }
2717 self.least_loaded(candidates)
2718 }
2719
2720 fn least_loaded(&self, candidates: Vec<&Member>) -> Result<String, DcpError> {
2721 candidates
2722 .into_iter()
2723 .min_by(|left, right| {
2724 let left_count = self.jobs_on_node(&left.node_id);
2725 let right_count = self.jobs_on_node(&right.node_id);
2726 left_count
2727 .cmp(&right_count)
2728 .then_with(|| left.node_id.cmp(&right.node_id))
2729 })
2730 .map(|member| member.node_id.clone())
2731 .ok_or_else(|| DcpError::response(ResponseStatus::NotFound, "no eligible node"))
2732 }
2733
2734 fn jobs_on_node(&self, node_id: &str) -> usize {
2735 self.assignments
2736 .values()
2737 .filter(|metadata| metadata.assigned_node == node_id)
2738 .count()
2739 }
2740
2741 async fn start_on_node(
2742 &self,
2743 node_id: &str,
2744 factory_name: String,
2745 instance_name: String,
2746 params: HashMap<String, String>,
2747 metadata: ClusterJobMetadata,
2748 timeout: Duration,
2749 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
2750 if node_id == self.self_node {
2751 let mut spec = self
2752 .factories
2753 .build(&factory_name, instance_name.clone(), params)?;
2754 spec = spec.with_cluster_metadata(metadata);
2755 let registry = self.registry.clone();
2756 let status = tokio::task::spawn_blocking(move || {
2757 registry.submit(spec)?;
2758 registry.start(instance_name)
2759 })
2760 .await??;
2761 return Ok(wire_job_status(&status));
2762 }
2763
2764 self.sessions
2765 .start_cluster_job(
2766 node_id,
2767 factory_name,
2768 instance_name,
2769 params,
2770 wire_cluster_job_start(&metadata),
2771 timeout,
2772 )
2773 .await
2774 .map_err(session_error)
2775 }
2776
2777 async fn ensure_started_on_node(
2778 &self,
2779 node_id: &str,
2780 factory_name: String,
2781 instance_name: String,
2782 params: HashMap<String, String>,
2783 metadata: ClusterJobMetadata,
2784 timeout: Duration,
2785 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
2786 match self
2787 .start_on_node(
2788 node_id,
2789 factory_name,
2790 instance_name.clone(),
2791 params,
2792 metadata.clone(),
2793 timeout,
2794 )
2795 .await
2796 {
2797 Ok(status) => Ok(status),
2798 Err(error) if placement_start_may_have_committed(&error) => {
2799 match self.status_on_node(node_id, instance_name, timeout).await {
2800 Ok(status) if status_matches_assignment(&status, &metadata) => Ok(status),
2801 _ => Err(error),
2802 }
2803 }
2804 Err(error) => Err(error),
2805 }
2806 }
2807
2808 async fn status_on_node(
2809 &self,
2810 node_id: &str,
2811 name: String,
2812 timeout: Duration,
2813 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
2814 if node_id == self.self_node {
2815 let registry = self.registry.clone();
2816 let status = tokio::task::spawn_blocking(move || registry.status(name)).await??;
2817 return Ok(wire_job_status(&status));
2818 }
2819 self.sessions
2820 .cluster_job_status(node_id, name, timeout)
2821 .await
2822 .map_err(session_error)
2823 }
2824
2825 async fn drain_on_node(
2826 &self,
2827 node_id: &str,
2828 name: String,
2829 timeout: Duration,
2830 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
2831 if node_id == self.self_node {
2832 let registry = self.registry.clone();
2833 let status = tokio::task::spawn_blocking(move || registry.drain(name)).await??;
2834 return Ok(wire_job_status(&status));
2835 }
2836 self.sessions
2837 .drain_cluster_job(node_id, name, timeout)
2838 .await
2839 .map_err(session_error)
2840 }
2841
2842 async fn stop_on_node(
2843 &self,
2844 node_id: &str,
2845 name: String,
2846 timeout: Duration,
2847 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
2848 if node_id == self.self_node {
2849 let registry = self.registry.clone();
2850 let status = tokio::task::spawn_blocking(move || registry.stop(name)).await??;
2851 return Ok(wire_job_status(&status));
2852 }
2853 self.sessions
2854 .stop_cluster_job(node_id, name, timeout)
2855 .await
2856 .map_err(session_error)
2857 }
2858
2859 async fn collect_job_statuses(
2860 &self,
2861 timeout: Duration,
2862 ) -> Result<Vec<crate::dcp::proto::JobStatus>, DcpError> {
2863 let snapshot = self.state.get();
2864 let mut statuses = Vec::new();
2865 let registry = self.registry.clone();
2866 let local_jobs = tokio::task::spawn_blocking(move || registry.list()).await??;
2867 statuses.extend(local_jobs.iter().map(wire_job_status));
2868
2869 let mut pending = Vec::new();
2870 for member in snapshot.members.values() {
2871 if member.node_id == self.self_node
2872 || !member.has_role(&self.agent_role)
2873 || member.state != MemberState::Up
2874 || member.unreachable
2875 {
2876 continue;
2877 }
2878 let node_id = member.node_id.clone();
2879 let sessions = self.sessions.clone();
2880 pending.push(tokio::spawn(async move {
2881 sessions.list_jobs(&node_id, timeout).await
2882 }));
2883 }
2884 drop(snapshot);
2885
2886 let deadline = Instant::now() + timeout;
2891 for task in pending {
2892 let remaining = deadline.saturating_duration_since(Instant::now());
2893 if remaining.is_zero() {
2894 return Err(DcpError::response(
2895 ResponseStatus::DeadlineExceeded,
2896 "registry reconciliation timed out",
2897 ));
2898 }
2899 match tokio::time::timeout(remaining, task).await {
2900 Ok(Ok(Ok(mut peer_statuses))) => statuses.append(&mut peer_statuses),
2901 Ok(Ok(Err(message))) => return Err(session_error(message)),
2902 Ok(Err(error)) => {
2903 return Err(DcpError::response(
2904 ResponseStatus::Failed,
2905 format!("registry reconciliation task failed: {error}"),
2906 ));
2907 }
2908 Err(_) => {
2909 return Err(DcpError::response(
2910 ResponseStatus::DeadlineExceeded,
2911 "registry reconciliation timed out",
2912 ));
2913 }
2914 }
2915 }
2916 Ok(statuses)
2917 }
2918
2919 async fn publish_assignment(
2920 &mut self,
2921 instance_name: String,
2922 metadata: ClusterJobMetadata,
2923 timeout: Duration,
2924 ) -> Result<(), DcpError> {
2925 self.reconcile_local_assignment(&instance_name, &metadata)
2926 .await?;
2927 self.replicate_assignment_mutation(
2928 instance_name,
2929 Some(wire_cluster_job_start(&metadata)),
2930 timeout,
2931 );
2932 Ok(())
2933 }
2934
2935 async fn reconcile_local_assignment(
2936 &self,
2937 instance_name: &str,
2938 metadata: &ClusterJobMetadata,
2939 ) -> Result<(), DcpError> {
2940 if metadata.assigned_node == self.self_node {
2941 return self.update_local_assignment(instance_name, metadata).await;
2942 }
2943 self.stop_stale_local_assignment(instance_name, metadata)
2944 .await
2945 }
2946
2947 async fn update_local_assignment(
2948 &self,
2949 instance_name: &str,
2950 metadata: &ClusterJobMetadata,
2951 ) -> Result<(), DcpError> {
2952 if metadata.assigned_node != self.self_node {
2953 return Ok(());
2954 }
2955 let registry = self.registry.clone();
2956 let instance_name = instance_name.to_owned();
2957 let metadata = metadata.clone();
2958 tokio::task::spawn_blocking(move || {
2959 registry.update_cluster_metadata(instance_name, metadata)
2960 })
2961 .await??;
2962 Ok(())
2963 }
2964
2965 async fn stop_stale_local_assignment(
2966 &self,
2967 instance_name: &str,
2968 metadata: &ClusterJobMetadata,
2969 ) -> Result<(), DcpError> {
2970 let registry = self.registry.clone();
2971 let instance_name = instance_name.to_owned();
2972 let self_node = self.self_node.clone();
2973 let placement_generation = metadata.placement_generation;
2974 tokio::task::spawn_blocking(move || -> Result<(), AgentError> {
2975 let status = match registry.status(instance_name.clone()) {
2976 Ok(status) => status,
2977 Err(AgentError::JobNotFound(_)) => return Ok(()),
2978 Err(error) => return Err(error),
2979 };
2980 let Some(local_metadata) = status.cluster.as_ref() else {
2981 return Ok(());
2982 };
2983 if local_metadata.assigned_node == self_node
2984 && local_metadata.placement_generation < placement_generation
2985 && should_stop_stale_local_job(status.state, status.desired_state)
2986 {
2987 let _ = registry.stop(instance_name)?;
2988 }
2989 Ok(())
2990 })
2991 .await??;
2992 Ok(())
2993 }
2994
2995 async fn stop_local_cluster_assignment(
2996 &self,
2997 instance_name: &str,
2998 tombstone: &AssignmentTombstone,
2999 current_coordinator: Option<&str>,
3000 ) -> Result<(), DcpError> {
3001 let registry = self.registry.clone();
3002 let instance_name = instance_name.to_owned();
3003 let tombstone = tombstone.clone();
3004 let current_coordinator = current_coordinator.map(str::to_owned);
3005 tokio::task::spawn_blocking(move || -> Result<(), AgentError> {
3006 let status = match registry.status(instance_name.clone()) {
3007 Ok(status) => status,
3008 Err(AgentError::JobNotFound(_)) => return Ok(()),
3009 Err(error) => return Err(error),
3010 };
3011 let Some(mut local_metadata) = status.cluster.clone() else {
3012 return Ok(());
3013 };
3014 if !tombstone_dominates_assignment(
3015 &tombstone,
3016 &local_metadata,
3017 current_coordinator.as_deref(),
3018 ) {
3019 return Ok(());
3020 }
3021 local_metadata
3022 .coordinator_node
3023 .clone_from(&tombstone.coordinator_node);
3024 if local_metadata.placement_generation < tombstone.placement_generation {
3025 local_metadata.placement_generation = tombstone.placement_generation;
3026 }
3027 registry.update_cluster_metadata(instance_name.clone(), local_metadata)?;
3028 if should_stop_stale_local_job(status.state, status.desired_state) {
3029 let _ = registry.stop(instance_name)?;
3030 }
3031 Ok(())
3032 })
3033 .await??;
3034 Ok(())
3035 }
3036
3037 async fn stop_stale_restarted_local_assignment(
3038 &self,
3039 instance_name: &str,
3040 known_assignment: &ClusterJobMetadata,
3041 ) -> Result<(), DcpError> {
3042 let registry = self.registry.clone();
3043 let instance_name = instance_name.to_owned();
3044 let self_node = self.self_node.clone();
3045 let known_assignment = known_assignment.clone();
3046 tokio::task::spawn_blocking(move || -> Result<(), AgentError> {
3047 let status = match registry.status(instance_name.clone()) {
3048 Ok(status) => status,
3049 Err(AgentError::JobNotFound(_)) => return Ok(()),
3050 Err(error) => return Err(error),
3051 };
3052 let Some(local_metadata) = status.cluster.as_ref() else {
3053 return Ok(());
3054 };
3055 let stale_local_owner = local_metadata.assigned_node == self_node
3056 && (local_metadata.placement_generation < known_assignment.placement_generation
3057 || (local_metadata.placement_generation
3058 == known_assignment.placement_generation
3059 && known_assignment.assigned_node != self_node));
3060 if stale_local_owner && should_stop_stale_local_job(status.state, status.desired_state)
3061 {
3062 let _ = registry.stop(instance_name)?;
3063 }
3064 Ok(())
3065 })
3066 .await??;
3067 Ok(())
3068 }
3069
3070 fn replicate_assignment_mutation(
3071 &mut self,
3072 instance_name: String,
3073 assignment: Option<ClusterJobStart>,
3074 timeout: Duration,
3075 ) {
3076 let peers = self.assignment_sync_peers();
3077 let sync_mark = self.assignment_sync_mark(&instance_name, assignment.as_ref());
3078 let timeout = timeout.min(Duration::from_millis(250));
3079 for (node_id, incarnation) in peers {
3080 self.sync_assignment_to_peer(
3081 node_id,
3082 incarnation,
3083 instance_name.clone(),
3084 assignment.clone(),
3085 sync_mark.clone(),
3086 timeout,
3087 )
3088 }
3089 }
3090
3091 fn sync_assignments_to_peers(&mut self, timeout: Duration) {
3092 let peers = self.assignment_sync_peers();
3093 if peers.is_empty() {
3094 return;
3095 }
3096 let timeout = timeout.min(Duration::from_millis(250));
3097 let assignments = self
3098 .assignments
3099 .iter()
3100 .map(|(name, metadata)| {
3101 (
3102 name.clone(),
3103 Some(wire_cluster_job_start(metadata)),
3104 AssignmentSyncMark::active(metadata),
3105 )
3106 })
3107 .collect::<Vec<_>>();
3108 let tombstones = self
3109 .assignment_tombstones
3110 .iter()
3111 .filter(|(name, _tombstone)| !self.assignments.contains_key(*name))
3112 .map(|(name, tombstone)| (name.clone(), None, AssignmentSyncMark::tombstone(tombstone)))
3113 .collect::<Vec<_>>();
3114
3115 for (node_id, incarnation) in peers {
3116 for (name, assignment, sync_mark) in assignments.iter().chain(tombstones.iter()) {
3117 self.sync_assignment_to_peer(
3118 node_id.clone(),
3119 incarnation,
3120 name.clone(),
3121 assignment.clone(),
3122 sync_mark.clone(),
3123 timeout,
3124 )
3125 }
3126 }
3127 }
3128
3129 fn assignment_sync_peers(&mut self) -> Vec<(String, u64)> {
3130 let snapshot = self.state.get();
3131 let peers = snapshot
3132 .members
3133 .values()
3134 .filter(|member| {
3135 member.node_id != self.self_node
3136 && member.has_role(&self.agent_role)
3137 && member.state == MemberState::Up
3138 && !member.unreachable
3139 })
3140 .map(|member| (member.node_id.clone(), member.incarnation))
3141 .collect::<BTreeMap<_, _>>();
3142 drop(snapshot);
3143 self.replicated_assignments.retain(|node_id, sync| {
3144 peers
3145 .get(node_id)
3146 .is_some_and(|incarnation| *incarnation == sync.incarnation)
3147 });
3148 self.inflight_assignment_syncs.retain(|key| {
3149 peers
3150 .get(&key.node_id)
3151 .is_some_and(|incarnation| *incarnation == key.incarnation)
3152 });
3153 peers.into_iter().collect()
3154 }
3155
3156 fn sync_assignment_to_peer(
3157 &mut self,
3158 node_id: String,
3159 incarnation: u64,
3160 instance_name: String,
3161 assignment: Option<ClusterJobStart>,
3162 sync_mark: AssignmentSyncMark,
3163 timeout: Duration,
3164 ) {
3165 if self
3166 .replicated_assignments
3167 .get(&node_id)
3168 .filter(|sync| sync.incarnation == incarnation)
3169 .and_then(|sync| sync.assignments.get(&instance_name))
3170 == Some(&sync_mark)
3171 {
3172 return;
3173 }
3174 let sync_key = AssignmentSyncKey {
3175 node_id: node_id.clone(),
3176 incarnation,
3177 instance_name: instance_name.clone(),
3178 sync_mark: sync_mark.clone(),
3179 };
3180 if !self.inflight_assignment_syncs.insert(sync_key) {
3181 return;
3182 }
3183
3184 let sessions = self.sessions.clone();
3185 let commands = self.commands.clone();
3186 let tombstone = self
3187 .assignment_tombstone_for_sync(&instance_name, assignment.as_ref())
3188 .cloned();
3189 tokio::spawn(async move {
3190 let ok = sessions
3191 .remember_cluster_assignment(
3192 &node_id,
3193 instance_name.clone(),
3194 assignment,
3195 tombstone,
3196 timeout,
3197 )
3198 .await
3199 .is_ok();
3200 let _ = commands
3201 .send(PlacementCommand::AssignmentSyncCompleted {
3202 node_id,
3203 incarnation,
3204 instance_name,
3205 sync_mark,
3206 ok,
3207 })
3208 .await;
3209 });
3210 }
3211
3212 fn assignment_sync_completed(
3213 &mut self,
3214 node_id: String,
3215 incarnation: u64,
3216 instance_name: String,
3217 sync_mark: AssignmentSyncMark,
3218 ok: bool,
3219 ) {
3220 self.inflight_assignment_syncs.remove(&AssignmentSyncKey {
3221 node_id: node_id.clone(),
3222 incarnation,
3223 instance_name: instance_name.clone(),
3224 sync_mark: sync_mark.clone(),
3225 });
3226 if ok {
3227 if self.current_assignment_sync_mark(&instance_name) != sync_mark {
3228 return;
3229 }
3230 self.replicated_assignments
3231 .entry(node_id)
3232 .and_modify(|sync| {
3233 if sync.incarnation != incarnation {
3234 sync.incarnation = incarnation;
3235 sync.assignments.clear();
3236 }
3237 sync.assignments
3238 .insert(instance_name.clone(), sync_mark.clone());
3239 })
3240 .or_insert_with(|| PeerAssignmentSync {
3241 incarnation,
3242 assignments: BTreeMap::from([(instance_name, sync_mark)]),
3243 });
3244 }
3245 }
3246
3247 fn current_assignment_sync_mark(&self, instance_name: &str) -> AssignmentSyncMark {
3248 self.assignments
3249 .get(instance_name)
3250 .map(AssignmentSyncMark::active)
3251 .unwrap_or_else(|| {
3252 self.assignment_tombstones
3253 .get(instance_name)
3254 .map(AssignmentSyncMark::tombstone)
3255 .unwrap_or_else(|| {
3256 AssignmentSyncMark::tombstone(&AssignmentTombstone {
3257 placement_generation: 0,
3258 coordinator_node: String::new(),
3259 })
3260 })
3261 })
3262 }
3263
3264 fn assignment_tombstone_for_sync(
3265 &self,
3266 instance_name: &str,
3267 assignment: Option<&ClusterJobStart>,
3268 ) -> Option<&AssignmentTombstone> {
3269 assignment
3270 .is_none()
3271 .then(|| self.assignment_tombstones.get(instance_name))
3272 .flatten()
3273 }
3274
3275 fn assignment_sync_mark(
3276 &self,
3277 instance_name: &str,
3278 assignment: Option<&ClusterJobStart>,
3279 ) -> AssignmentSyncMark {
3280 match assignment {
3281 Some(assignment) => AssignmentSyncMark::Active {
3282 placement_generation: assignment.placement_generation,
3283 assigned_node: assignment.assigned_node_id.clone(),
3284 coordinator_node: assignment.coordinator_node_id.clone(),
3285 },
3286 None => self
3287 .assignment_tombstones
3288 .get(instance_name)
3289 .map(AssignmentSyncMark::tombstone)
3290 .unwrap_or_else(|| {
3291 AssignmentSyncMark::tombstone(&AssignmentTombstone {
3292 placement_generation: 0,
3293 coordinator_node: String::new(),
3294 })
3295 }),
3296 }
3297 }
3298
3299 fn record_assignment_tombstone(&mut self, instance_name: &str, metadata: &ClusterJobMetadata) {
3300 let current_coordinator = self.current_placement_coordinator();
3301 self.record_assignment_tombstone_value(
3302 instance_name,
3303 AssignmentTombstone::from_metadata(metadata),
3304 current_coordinator.as_deref(),
3305 );
3306 }
3307
3308 fn record_assignment_tombstone_value(
3309 &mut self,
3310 instance_name: &str,
3311 tombstone: AssignmentTombstone,
3312 current_coordinator: Option<&str>,
3313 ) {
3314 self.assignment_tombstones
3315 .entry(instance_name.to_owned())
3316 .and_modify(|existing| {
3317 if tombstone_supersedes(&tombstone, existing, current_coordinator) {
3318 *existing = tombstone.clone();
3319 }
3320 })
3321 .or_insert(tombstone);
3322 }
3323}
3324
3325fn remember_assignment(
3326 assignments: &mut BTreeMap<String, ClusterJobMetadata>,
3327 instance_name: String,
3328 metadata: ClusterJobMetadata,
3329 current_coordinator: Option<&str>,
3330) -> bool {
3331 let accept = assignments.get(&instance_name).is_none_or(|existing| {
3332 metadata.placement_generation > existing.placement_generation
3333 || (metadata.placement_generation == existing.placement_generation
3334 && metadata.assigned_node == existing.assigned_node
3335 && (metadata.coordinator_node == existing.coordinator_node
3336 || current_coordinator
3337 .is_some_and(|coordinator| coordinator == metadata.coordinator_node)))
3338 });
3339 if accept {
3340 assignments.insert(instance_name, metadata);
3341 }
3342 accept
3343}
3344
3345fn ensure_mutation_from_current_coordinator(
3346 mutation_coordinator: &str,
3347 current_coordinator: Option<&str>,
3348) -> Result<(), DcpError> {
3349 if mutation_coordinator.is_empty() {
3350 return Err(DcpError::response(
3351 ResponseStatus::Conflict,
3352 "assignment mutation is missing coordinator fence",
3353 ));
3354 }
3355 if current_coordinator == Some(mutation_coordinator) {
3356 return Ok(());
3357 }
3358 Err(DcpError::response(
3359 ResponseStatus::Conflict,
3360 format!(
3361 "assignment mutation from non-current coordinator: mutation={mutation_coordinator} current={}",
3362 current_coordinator.unwrap_or("<none>")
3363 ),
3364 ))
3365}
3366
3367fn tombstone_dominates_assignment(
3368 tombstone: &AssignmentTombstone,
3369 metadata: &ClusterJobMetadata,
3370 current_coordinator: Option<&str>,
3371) -> bool {
3372 tombstone.placement_generation > metadata.placement_generation
3373 || (tombstone.placement_generation == metadata.placement_generation
3374 && current_coordinator == Some(tombstone.coordinator_node.as_str()))
3375}
3376
3377fn active_assignment_supersedes_tombstone(
3378 metadata: &ClusterJobMetadata,
3379 tombstone: &AssignmentTombstone,
3380) -> bool {
3381 metadata.placement_generation > tombstone.placement_generation
3382}
3383
3384fn registry_observation_is_tombstoned(
3385 tombstone: &AssignmentTombstone,
3386 metadata: &ClusterJobMetadata,
3387) -> bool {
3388 !active_assignment_supersedes_tombstone(metadata, tombstone)
3389}
3390
3391fn active_assignment_blocks_local_restart(
3392 restarted: &ClusterJobMetadata,
3393 existing: &ClusterJobMetadata,
3394) -> bool {
3395 existing.placement_generation > restarted.placement_generation
3396 || (existing.placement_generation == restarted.placement_generation
3397 && existing.assigned_node != restarted.assigned_node)
3398}
3399
3400fn stale_assignment_conflict(
3401 instance_name: &str,
3402 placement_generation: u64,
3403 tombstone: &AssignmentTombstone,
3404) -> DcpError {
3405 DcpError::response(
3406 ResponseStatus::Conflict,
3407 format!(
3408 "cluster assignment {instance_name} generation {placement_generation} does not supersede tombstone generation {} from coordinator {}",
3409 tombstone.placement_generation, tombstone.coordinator_node
3410 ),
3411 )
3412}
3413
3414fn stale_restart_conflict(
3415 instance_name: &str,
3416 restarted: &ClusterJobMetadata,
3417 existing: &ClusterJobMetadata,
3418) -> DcpError {
3419 DcpError::response(
3420 ResponseStatus::Conflict,
3421 format!(
3422 "local restart for cluster assignment {instance_name} generation {} on {} is superseded by active generation {} on {}",
3423 restarted.placement_generation,
3424 restarted.assigned_node,
3425 existing.placement_generation,
3426 existing.assigned_node
3427 ),
3428 )
3429}
3430
3431fn placement_start_may_have_committed(error: &DcpError) -> bool {
3432 match error {
3433 DcpError::Closed | DcpError::Protocol(_) | DcpError::Io(_) | DcpError::Decode(_) => true,
3434 DcpError::Response { status, message } => {
3435 matches!(
3436 status,
3437 ResponseStatus::Conflict
3438 | ResponseStatus::DeadlineExceeded
3439 | ResponseStatus::Failed
3440 ) || message.contains("peer StartJob timed out")
3441 || message.contains("DCP connection closed")
3442 }
3443 DcpError::Agent(AgentError::JobAlreadyExists(_))
3444 | DcpError::Agent(AgentError::JobAlreadyRunning(_)) => true,
3445 DcpError::Encode(_) | DcpError::Agent(_) | DcpError::Stream(_) | DcpError::Join(_) => false,
3446 }
3447}
3448
3449fn status_matches_assignment(
3450 status: &crate::dcp::proto::JobStatus,
3451 metadata: &ClusterJobMetadata,
3452) -> bool {
3453 status.cluster_job
3454 && status.desired_state == "Running"
3455 && status.placement_node_id == metadata.assigned_node
3456 && status.placement_generation == metadata.placement_generation
3457 && status.coordinator_node_id == metadata.coordinator_node
3458}
3459
3460const fn should_stop_stale_local_job(state: JobState, desired_state: DesiredJobState) -> bool {
3461 matches!(desired_state, DesiredJobState::Running)
3462 && matches!(
3463 state,
3464 JobState::Starting | JobState::Running | JobState::BackingOff
3465 )
3466}
3467
3468fn tombstone_supersedes(
3469 candidate: &AssignmentTombstone,
3470 existing: &AssignmentTombstone,
3471 current_coordinator: Option<&str>,
3472) -> bool {
3473 candidate.placement_generation > existing.placement_generation
3474 || (candidate.placement_generation == existing.placement_generation
3475 && !candidate.coordinator_node.is_empty()
3476 && (candidate.coordinator_node == existing.coordinator_node
3477 || current_coordinator == Some(candidate.coordinator_node.as_str())))
3478}
3479
3480fn validate_submit(request: &SubmitClusterJob) -> Result<(), DcpError> {
3481 if request.factory_name.trim().is_empty() || request.instance_name.trim().is_empty() {
3482 return Err(DcpError::response(
3483 ResponseStatus::BadRequest,
3484 "SubmitClusterJob requires factory_name and instance_name",
3485 ));
3486 }
3487 Ok(())
3488}
3489
3490fn eligible_members<'a>(
3491 snapshot: &'a ClusterState,
3492 agent_role: &str,
3493 role_constraint: Option<&str>,
3494) -> Vec<&'a Member> {
3495 snapshot
3496 .members
3497 .values()
3498 .filter(|member| member.state == MemberState::Up && !member.unreachable)
3499 .filter(|member| member.has_role(agent_role))
3500 .filter(|member| role_constraint.is_none_or(|role| member.has_role(role)))
3501 .collect()
3502}
3503
3504fn no_eligible_node_error(placement: &RegistryPlacementSpec) -> DcpError {
3505 let role = placement
3506 .role_constraint
3507 .as_deref()
3508 .map(|role| format!(" with role '{role}'"))
3509 .unwrap_or_default();
3510 DcpError::response(
3511 ResponseStatus::NotFound,
3512 format!("no eligible Up placement nodes{role}"),
3513 )
3514}
3515
3516fn metadata_from_status(
3517 status: &crate::dcp::proto::JobStatus,
3518) -> Result<Option<ClusterJobMetadata>, DcpError> {
3519 if !status.cluster_job {
3520 return Ok(None);
3521 }
3522 let placement = placement_spec_from_wire(status.placement.clone())?;
3523 Ok(Some(ClusterJobMetadata {
3524 factory_name: status.factory_name.clone(),
3525 params: status.params.clone().into_iter().collect(),
3526 placement,
3527 coordinator_node: status.coordinator_node_id.clone(),
3528 assigned_node: status.placement_node_id.clone(),
3529 placement_generation: status.placement_generation,
3530 history: status
3531 .placement_history
3532 .iter()
3533 .map(|history| ClusterPlacementHistory {
3534 generation: history.generation,
3535 from_node: if history.from_node_id.is_empty() {
3536 None
3537 } else {
3538 Some(history.from_node_id.clone())
3539 },
3540 to_node: history.to_node_id.clone(),
3541 reason: history.reason.clone(),
3542 timestamp: std::time::UNIX_EPOCH + Duration::from_millis(history.timestamp_ms),
3543 })
3544 .collect(),
3545 }))
3546}
3547
3548fn session_error(message: String) -> DcpError {
3549 let status = if message.contains("DCP response NotFound") {
3550 ResponseStatus::NotFound
3551 } else if message.contains("DCP response Conflict") {
3552 ResponseStatus::Conflict
3553 } else if message.contains("DCP response BadRequest") {
3554 ResponseStatus::BadRequest
3555 } else if message.contains("DCP response DeadlineExceeded") {
3556 ResponseStatus::DeadlineExceeded
3557 } else {
3558 ResponseStatus::Failed
3559 };
3560 DcpError::response(status, message)
3561}
3562
3563struct ClusterView {
3564 registry: JobRegistryHandle,
3565 state: Signal<ClusterState>,
3566 sessions: NodeSessionManagerHandle,
3567 placement: PlacementCoordinatorHandle,
3568 self_node: String,
3569 agent_role: String,
3570 cluster_events: ClusterEventPublisher,
3571}
3572
3573impl ClusterViewProvider for ClusterView {
3574 fn subscribe_cluster_events(&self) -> crate::dcp::DcpResult<mpsc::Receiver<ClusterEvent>> {
3575 self.cluster_events.subscribe()
3576 }
3577
3578 fn subscribe_cluster_metrics(
3579 &self,
3580 interval: Duration,
3581 ) -> Option<mpsc::Receiver<MetricSample>> {
3582 Some(spawn_cluster_metrics_proxy(
3583 self.self_node.clone(),
3584 self.sessions.clone(),
3585 interval,
3586 ))
3587 }
3588
3589 fn submit_cluster_job(
3590 &self,
3591 request: SubmitClusterJob,
3592 timeout: Duration,
3593 ) -> crate::dcp::server::ClusterViewFuture<'_, crate::dcp::proto::JobStatus> {
3594 Box::pin(async move { self.submit_cluster_job_inner(request, timeout).await })
3595 }
3596
3597 fn list_cluster_jobs(
3598 &self,
3599 timeout: Duration,
3600 ) -> crate::dcp::server::ClusterViewFuture<'_, ClusterJobList> {
3601 Box::pin(async move { self.list_cluster_jobs_inner(timeout).await })
3602 }
3603
3604 fn cluster_node_info(
3605 &self,
3606 timeout: Duration,
3607 ) -> crate::dcp::server::ClusterViewFuture<'_, ClusterNodeList> {
3608 Box::pin(async move { self.cluster_node_info_inner(timeout).await })
3609 }
3610
3611 fn cluster_job_status(
3612 &self,
3613 name: String,
3614 timeout: Duration,
3615 ) -> crate::dcp::server::ClusterViewFuture<'_, crate::dcp::proto::JobStatus> {
3616 Box::pin(async move { self.cluster_job_status_inner(name, timeout).await })
3617 }
3618
3619 fn drain_cluster_job(
3620 &self,
3621 name: String,
3622 timeout: Duration,
3623 ) -> crate::dcp::server::ClusterViewFuture<'_, crate::dcp::proto::JobStatus> {
3624 Box::pin(async move { self.drain_cluster_job_inner(name, timeout).await })
3625 }
3626
3627 fn stop_cluster_job(
3628 &self,
3629 name: String,
3630 timeout: Duration,
3631 ) -> crate::dcp::server::ClusterViewFuture<'_, crate::dcp::proto::JobStatus> {
3632 Box::pin(async move { self.stop_cluster_job_inner(name, timeout).await })
3633 }
3634
3635 fn remember_cluster_assignment(
3636 &self,
3637 request: RememberClusterAssignment,
3638 ) -> crate::dcp::server::ClusterViewFuture<'_, ()> {
3639 Box::pin(async move { self.remember_cluster_assignment_inner(request).await })
3640 }
3641
3642 fn register_restarted_cluster_assignment(
3643 &self,
3644 instance_name: String,
3645 assignment: ClusterJobStart,
3646 timeout: Duration,
3647 ) -> crate::dcp::server::ClusterViewFuture<'_, ()> {
3648 Box::pin(async move {
3649 self.placement
3650 .register_restarted_cluster_assignment(instance_name, assignment, timeout)
3651 .await
3652 })
3653 }
3654}
3655
3656struct RemoteMetricRoute {
3657 target: MetricRouteTarget,
3658 cancel: oneshot::Sender<()>,
3659 task: JoinHandle<()>,
3660}
3661
3662struct RemoteMetricUpdate {
3663 node_id: String,
3664 sample: MetricSample,
3665}
3666
3667struct RemoteMetricSubscriptionAttempt {
3668 node_id: String,
3669 target: MetricRouteTarget,
3670 result: Result<MetricSubscription, String>,
3671}
3672
3673#[derive(Default)]
3674struct RemoteMetricProxyState {
3675 routes: BTreeMap<String, RemoteMetricRoute>,
3676 latest: BTreeMap<String, MetricSample>,
3677 pending: BTreeMap<String, MetricRouteTarget>,
3678 wanted: BTreeMap<String, MetricRouteTarget>,
3679}
3680
3681fn spawn_cluster_metrics_proxy(
3682 self_node: String,
3683 sessions: NodeSessionManagerHandle,
3684 interval: Duration,
3685) -> mpsc::Receiver<MetricSample> {
3686 let (output, receiver) = mpsc::channel(8);
3687 tokio::spawn(run_cluster_metrics_proxy(
3688 self_node, sessions, interval, output,
3689 ));
3690 receiver
3691}
3692
3693async fn run_cluster_metrics_proxy(
3694 self_node: String,
3695 sessions: NodeSessionManagerHandle,
3696 interval: Duration,
3697 output: mpsc::Sender<MetricSample>,
3698) {
3699 let (updates, mut update_receiver) = mpsc::channel(32);
3700 let (attempts, mut attempt_receiver) = mpsc::channel(32);
3701 let mut proxy = RemoteMetricProxyState::default();
3702 let retry_interval = interval
3703 .max(Duration::from_millis(50))
3704 .min(Duration::from_secs(1));
3705 let mut retry = tokio::time::interval(retry_interval);
3706 let mut emit = tokio::time::interval(interval);
3707
3708 reconcile_metric_routes(&self_node, &sessions, interval, &attempts, &mut proxy).await;
3709
3710 loop {
3711 tokio::select! {
3712 _ = output.closed() => break,
3713 _ = retry.tick() => {
3714 reconcile_metric_routes(
3715 &self_node,
3716 &sessions,
3717 interval,
3718 &attempts,
3719 &mut proxy,
3720 ).await;
3721 }
3722 attempt = attempt_receiver.recv() => {
3723 if let Some(attempt) = attempt {
3724 install_metric_route(
3725 attempt,
3726 &mut proxy,
3727 &updates,
3728 ).await;
3729 }
3730 }
3731 update = update_receiver.recv() => {
3732 if let Some(update) = update
3733 && proxy.routes.contains_key(&update.node_id)
3734 {
3735 proxy.latest.insert(update.node_id, update.sample);
3736 }
3737 }
3738 _ = emit.tick() => {
3739 let streams = proxy.latest
3740 .values()
3741 .flat_map(|sample| sample.streams.iter().cloned())
3742 .collect();
3743 let nodes = proxy.latest
3744 .values()
3745 .flat_map(|sample| sample.nodes.iter().cloned())
3746 .collect();
3747 let timestamp_ms = proxy.latest
3748 .values()
3749 .map(|sample| sample.timestamp_ms)
3750 .max()
3751 .unwrap_or(0);
3752 match output.try_send(MetricSample {
3753 timestamp_ms,
3754 streams,
3755 nodes,
3756 }) {
3757 Ok(()) | Err(mpsc::error::TrySendError::Full(_)) => {}
3758 Err(mpsc::error::TrySendError::Closed(_)) => break,
3759 }
3760 }
3761 }
3762 }
3763
3764 drop(attempt_receiver);
3765 drop(attempts);
3766 cancel_metric_routes(proxy.routes).await;
3767}
3768
3769async fn reconcile_metric_routes(
3770 self_node: &str,
3771 sessions: &NodeSessionManagerHandle,
3772 interval: Duration,
3773 attempts: &mpsc::Sender<RemoteMetricSubscriptionAttempt>,
3774 proxy: &mut RemoteMetricProxyState,
3775) {
3776 let next_wanted = sessions
3777 .session_snapshots()
3778 .await
3779 .into_iter()
3780 .filter_map(|(node_id, snapshot)| {
3781 if node_id == self_node {
3782 return None;
3783 }
3784 snapshot
3785 .metric_route_target()
3786 .map(|target| (node_id, target))
3787 })
3788 .collect::<BTreeMap<_, _>>();
3789 proxy.wanted = next_wanted;
3790
3791 let stale = proxy
3792 .routes
3793 .iter()
3794 .filter(|(node_id, route)| {
3795 proxy.wanted.get(*node_id) != Some(&route.target) || route.task.is_finished()
3796 })
3797 .map(|(node_id, _)| node_id.clone())
3798 .collect::<Vec<_>>();
3799 for node_id in stale {
3800 proxy.latest.remove(&node_id);
3801 if let Some(route) = proxy.routes.remove(&node_id) {
3802 cancel_metric_route(route).await;
3803 }
3804 }
3805
3806 for (node_id, target) in &proxy.wanted {
3807 if proxy.routes.contains_key(node_id) || proxy.pending.get(node_id) == Some(target) {
3808 continue;
3809 }
3810 proxy.pending.insert(node_id.clone(), target.clone());
3811 let timeout = sessions.inner.config.request_timeout;
3812 let sessions = sessions.clone();
3813 let attempts = attempts.clone();
3814 let node_id = node_id.clone();
3815 let target = target.clone();
3816 let _task = tokio::spawn(async move {
3817 let result = sessions
3821 .subscribe_metrics(&node_id, interval, Vec::new(), timeout)
3822 .await;
3823 let attempt = RemoteMetricSubscriptionAttempt {
3824 node_id,
3825 target,
3826 result,
3827 };
3828 if let Err(error) = attempts.send(attempt).await
3829 && let Ok(subscription) = error.0.result
3830 {
3831 let _ = subscription.cancel().await;
3832 }
3833 });
3834 }
3835}
3836
3837async fn install_metric_route(
3838 attempt: RemoteMetricSubscriptionAttempt,
3839 proxy: &mut RemoteMetricProxyState,
3840 updates: &mpsc::Sender<RemoteMetricUpdate>,
3841) {
3842 let RemoteMetricSubscriptionAttempt {
3843 node_id,
3844 target,
3845 result,
3846 } = attempt;
3847 if proxy.pending.get(&node_id) == Some(&target) {
3848 proxy.pending.remove(&node_id);
3849 }
3850 let is_current =
3851 proxy.wanted.get(&node_id) == Some(&target) && !proxy.routes.contains_key(&node_id);
3852 let Ok(subscription) = result else {
3853 return;
3854 };
3855 if !is_current {
3856 let _ = subscription.cancel().await;
3857 return;
3858 }
3859
3860 let (cancel, cancel_receiver) = oneshot::channel();
3861 let task = tokio::spawn(pump_remote_metrics(
3862 node_id.clone(),
3863 subscription,
3864 updates.clone(),
3865 cancel_receiver,
3866 ));
3867 proxy.routes.insert(
3868 node_id,
3869 RemoteMetricRoute {
3870 target,
3871 cancel,
3872 task,
3873 },
3874 );
3875}
3876
3877async fn pump_remote_metrics(
3878 node_id: String,
3879 mut subscription: MetricSubscription,
3880 updates: mpsc::Sender<RemoteMetricUpdate>,
3881 mut cancel: oneshot::Receiver<()>,
3882) {
3883 loop {
3884 tokio::select! {
3885 _ = &mut cancel => {
3886 let _ = subscription.cancel().await;
3887 break;
3888 }
3889 sample = subscription.recv() => {
3890 let Some(sample) = sample else {
3891 break;
3892 };
3893 match updates.try_send(RemoteMetricUpdate {
3894 node_id: node_id.clone(),
3895 sample,
3896 }) {
3897 Ok(()) | Err(mpsc::error::TrySendError::Full(_)) => {}
3898 Err(mpsc::error::TrySendError::Closed(_)) => {
3899 let _ = subscription.cancel().await;
3900 break;
3901 }
3902 }
3903 }
3904 }
3905 }
3906}
3907
3908async fn cancel_metric_route(route: RemoteMetricRoute) {
3909 let _ = route.cancel.send(());
3910 let _ = route.task.await;
3911}
3912
3913async fn cancel_metric_routes(routes: BTreeMap<String, RemoteMetricRoute>) {
3914 for (_, route) in routes {
3915 cancel_metric_route(route).await;
3916 }
3917}
3918
3919impl ClusterView {
3920 async fn submit_cluster_job_inner(
3921 &self,
3922 request: SubmitClusterJob,
3923 timeout: Duration,
3924 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
3925 if self.is_local_coordinator() {
3926 return self.placement.submit_cluster_job(request, timeout).await;
3927 }
3928 let coordinator = self.coordinator_node_id()?;
3929 self.sessions
3930 .submit_cluster_job(&coordinator, request, timeout)
3931 .await
3932 .map_err(session_error)
3933 }
3934
3935 async fn list_cluster_jobs_inner(&self, timeout: Duration) -> Result<ClusterJobList, DcpError> {
3936 if !self.is_local_coordinator() {
3937 let coordinator = self.coordinator_node_id()?;
3938 return self
3939 .sessions
3940 .list_cluster_jobs(&coordinator, timeout)
3941 .await
3942 .map_err(session_error);
3943 }
3944
3945 let timeout = if timeout.is_zero() {
3946 self.sessions.inner.config.request_timeout
3947 } else {
3948 timeout
3949 };
3950 let snapshot = self.state.get();
3951 let local_member = snapshot.member(&self.self_node).cloned();
3952 let local_jobs = list_local_jobs(self.registry.clone()).await?;
3953 let mut nodes = Vec::new();
3954 nodes.push(ClusterJobNode {
3955 node_id: self.self_node.clone(),
3956 address: local_member
3957 .as_ref()
3958 .map(|member| member.address.to_string())
3959 .unwrap_or_default(),
3960 local: true,
3961 jobs: local_jobs,
3962 });
3963
3964 let mut errors = Vec::new();
3965 let mut pending = Vec::new();
3966 for member in snapshot.members.values() {
3967 if member.node_id == self.self_node || !member.has_role(&self.agent_role) {
3968 continue;
3969 }
3970 if member.state == MemberState::Removed {
3971 continue;
3972 }
3973 if member.state != MemberState::Up || member.unreachable {
3974 errors.push(ClusterNodeError {
3975 node_id: member.node_id.clone(),
3976 message: if member.unreachable {
3977 "member is unreachable".to_owned()
3978 } else {
3979 format!("member is {:?}", member.state)
3980 },
3981 });
3982 continue;
3983 }
3984 let node_id = member.node_id.clone();
3985 let address = member.address.to_string();
3986 let sessions = self.sessions.clone();
3987 pending.push(tokio::spawn(async move {
3988 let result = sessions.list_jobs(&node_id, timeout).await;
3989 (node_id, address, result)
3990 }));
3991 }
3992
3993 let deadline = Instant::now() + timeout;
3994 for task in pending {
3995 let remaining = deadline.saturating_duration_since(Instant::now());
3996 match tokio::time::timeout(remaining, task).await {
3997 Ok(Ok((node_id, address, Ok(jobs)))) => nodes.push(ClusterJobNode {
3998 node_id,
3999 address,
4000 local: false,
4001 jobs,
4002 }),
4003 Ok(Ok((node_id, _address, Err(message)))) => {
4004 errors.push(ClusterNodeError { node_id, message });
4005 }
4006 Ok(Err(error)) => errors.push(ClusterNodeError {
4007 node_id: "unknown".to_owned(),
4008 message: format!("cluster fan-out task failed: {error}"),
4009 }),
4010 Err(_) => errors.push(ClusterNodeError {
4011 node_id: "unknown".to_owned(),
4012 message: "cluster fan-out timed out".to_owned(),
4013 }),
4014 }
4015 }
4016
4017 nodes.sort_by(|left, right| left.node_id.cmp(&right.node_id));
4018 errors.sort_by(|left, right| left.node_id.cmp(&right.node_id));
4019 Ok(ClusterJobList {
4020 partial: !errors.is_empty(),
4021 nodes,
4022 errors,
4023 })
4024 }
4025
4026 async fn cluster_node_info_inner(
4027 &self,
4028 _timeout: Duration,
4029 ) -> Result<ClusterNodeList, DcpError> {
4030 let snapshot = self.state.get();
4031 let coordinator_node_id = snapshot
4032 .placement_coordinator()
4033 .map(|member| member.node_id.clone())
4034 .unwrap_or_default();
4035 let session_states = self.sessions.session_snapshots().await;
4036 let mut nodes = snapshot
4037 .members
4038 .values()
4039 .map(|member| ClusterNodeStatus {
4040 node_id: member.node_id.clone(),
4041 member_state: format!("{:?}", member.state),
4042 address: member.address.to_string(),
4043 agent_addr: member
4044 .agent_addr
4045 .map(|addr| addr.to_string())
4046 .unwrap_or_default(),
4047 roles: member.roles.clone(),
4048 unreachable: member.unreachable,
4049 local: member.node_id == self.self_node,
4050 session_state: if member.node_id == self.self_node {
4051 "local".to_owned()
4052 } else {
4053 session_states
4054 .get(&member.node_id)
4055 .map(|snapshot| snapshot.state.as_text())
4056 .unwrap_or_else(|| "not_connected".to_owned())
4057 },
4058 })
4059 .collect::<Vec<_>>();
4060 nodes.sort_by(|left, right| left.node_id.cmp(&right.node_id));
4061
4062 let mut errors = nodes
4063 .iter()
4064 .filter(|node| !node.local)
4065 .filter(|node| node.roles.iter().any(|role| role == &self.agent_role))
4066 .filter(|node| {
4067 node.member_state == "Up" && !node.unreachable && node.session_state != "connected"
4068 })
4069 .map(|node| ClusterNodeError {
4070 node_id: node.node_id.clone(),
4071 message: format!("node session {}", node.session_state),
4072 })
4073 .collect::<Vec<_>>();
4074 errors.sort_by(|left, right| left.node_id.cmp(&right.node_id));
4075
4076 Ok(ClusterNodeList {
4077 partial: !errors.is_empty(),
4078 nodes,
4079 errors,
4080 coordinator_node_id,
4081 })
4082 }
4083
4084 async fn cluster_job_status_inner(
4085 &self,
4086 name: String,
4087 timeout: Duration,
4088 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
4089 if self.is_local_coordinator() {
4090 return self.placement.cluster_job_status(name, timeout).await;
4091 }
4092 let coordinator = self.coordinator_node_id()?;
4093 self.sessions
4094 .forward_cluster_job_status(&coordinator, name, timeout)
4095 .await
4096 .map_err(session_error)
4097 }
4098
4099 async fn drain_cluster_job_inner(
4100 &self,
4101 name: String,
4102 timeout: Duration,
4103 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
4104 if self.is_local_coordinator() {
4105 return self.placement.drain_cluster_job(name, timeout).await;
4106 }
4107 let coordinator = self.coordinator_node_id()?;
4108 self.sessions
4109 .forward_drain_cluster_job(&coordinator, name, timeout)
4110 .await
4111 .map_err(session_error)
4112 }
4113
4114 async fn stop_cluster_job_inner(
4115 &self,
4116 name: String,
4117 timeout: Duration,
4118 ) -> Result<crate::dcp::proto::JobStatus, DcpError> {
4119 if self.is_local_coordinator() {
4120 return self.placement.stop_cluster_job(name, timeout).await;
4121 }
4122 let coordinator = self.coordinator_node_id()?;
4123 self.sessions
4124 .forward_stop_cluster_job(&coordinator, name, timeout)
4125 .await
4126 .map_err(session_error)
4127 }
4128
4129 async fn remember_cluster_assignment_inner(
4130 &self,
4131 request: RememberClusterAssignment,
4132 ) -> Result<(), DcpError> {
4133 let tombstone = request
4134 .assignment
4135 .is_none()
4136 .then(|| AssignmentTombstone::from_wire(&request));
4137 self.placement
4138 .remember_cluster_assignment(request.instance_name, request.assignment, tombstone)
4139 .await
4140 }
4141
4142 fn is_local_coordinator(&self) -> bool {
4143 self.state.get().is_placement_coordinator(&self.self_node)
4144 }
4145
4146 fn coordinator_node_id(&self) -> Result<String, DcpError> {
4147 self.state
4148 .get()
4149 .placement_coordinator()
4150 .map(|member| member.node_id.clone())
4151 .ok_or_else(|| {
4152 DcpError::response(
4153 ResponseStatus::Failed,
4154 "no placement coordinator is available",
4155 )
4156 })
4157 }
4158}
4159
4160async fn list_local_jobs(
4161 registry: JobRegistryHandle,
4162) -> Result<Vec<crate::dcp::proto::JobStatus>, DcpError> {
4163 let statuses = tokio::task::spawn_blocking(move || registry.list()).await??;
4164 Ok(statuses.iter().map(wire_job_status).collect())
4165}
4166
4167fn validate_cluster_agent_config(config: &ClusterAgentConfig) -> ClusterAgentResult<()> {
4168 if config.sessions.agent_role.trim().is_empty() {
4169 return Err(ClusterAgentError::InvalidConfig(
4170 "sessions.agent_role must not be empty".to_owned(),
4171 ));
4172 }
4173 if config.sessions.reconnect_min_backoff.is_zero()
4174 || config.sessions.reconnect_max_backoff < config.sessions.reconnect_min_backoff
4175 {
4176 return Err(ClusterAgentError::InvalidConfig(
4177 "node-session reconnect backoff must be non-zero and ordered".to_owned(),
4178 ));
4179 }
4180 match config.sessions.transport {
4181 NodeSessionTransport::TcpLoopback if config.dcp.tcp.is_none() => {
4182 Err(ClusterAgentError::InvalidConfig(
4183 "node-session TCP transport requires a DCP TCP listener".to_owned(),
4184 ))
4185 }
4186 NodeSessionTransport::QuicMtls { .. } if config.dcp.quic.is_none() => {
4187 Err(ClusterAgentError::InvalidConfig(
4188 "node-session QUIC transport requires a DCP QUIC listener".to_owned(),
4189 ))
4190 }
4191 _ => Ok(()),
4192 }
4193}
4194
4195fn advertised_agent_addr(
4196 transport: &NodeSessionTransport,
4197 handle: &DcpServerHandle,
4198) -> ClusterAgentResult<SocketAddr> {
4199 match transport {
4200 NodeSessionTransport::TcpLoopback => handle.tcp_addr().ok_or_else(|| {
4201 ClusterAgentError::InvalidConfig("DCP TCP listener did not start".to_owned())
4202 }),
4203 NodeSessionTransport::QuicMtls { .. } => handle.quic_addr().ok_or_else(|| {
4204 ClusterAgentError::InvalidConfig("DCP QUIC listener did not start".to_owned())
4205 }),
4206 }
4207}
4208
4209fn ensure_agent_role(config: &mut ClusterConfig, role: &str) {
4210 if !config.roles.iter().any(|candidate| candidate == role) {
4211 config.roles.push(role.to_owned());
4212 }
4213}
4214
4215#[cfg(test)]
4216mod tests {
4217 use super::*;
4218 use crate::{JobMat, JobRegistry, JobSpec};
4219 use datum::{Keep, Source};
4220 use prost::Message as _;
4221 use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _, DuplexStream};
4222
4223 #[tokio::test]
4224 async fn connected_peer_keeps_session_after_request_timeout() {
4225 let (client_io, mut peer_io) = tokio::io::duplex(4096);
4226 let (client_reader, client_writer) = tokio::io::split(client_io);
4227 let client = DcpClient::from_test_parts(client_reader, client_writer);
4228 let peer_task = tokio::spawn(async move {
4229 let first_request = read_test_request(&mut peer_io).await;
4230 tokio::time::sleep(Duration::from_millis(100)).await;
4231 write_empty_job_list(&mut peer_io, first_request).await;
4232
4233 let second_request = read_test_request(&mut peer_io).await;
4234 write_empty_job_list(&mut peer_io, second_request).await;
4235 });
4236
4237 let (pipe_io, _pipe_peer_io) = tokio::io::duplex(4096);
4238 let (pipe_reader, pipe_writer) = tokio::io::split(pipe_io);
4239 let pipe = ShardPipeClient::from_test_parts(pipe_reader, pipe_writer);
4240 let (command_sender, mut commands) = mpsc::channel(2);
4241 let (_pipe_sender, mut pipe_commands) = mpsc::channel(1);
4242 let (stop_sender, mut stop) = watch::channel(false);
4243 let session_task = tokio::spawn(async move {
4244 run_connected_peer(client, pipe, &mut commands, &mut pipe_commands, &mut stop).await
4245 });
4246
4247 let (first_reply, first_result) = oneshot::channel();
4248 command_sender
4249 .send(SessionCommand::ListJobs {
4250 timeout: Duration::from_millis(20),
4251 reply: first_reply,
4252 })
4253 .await
4254 .expect("first command accepted");
4255 assert_eq!(
4256 first_result.await.expect("first reply sent"),
4257 Err("peer ListJobs timed out".to_owned())
4258 );
4259
4260 let (second_reply, second_result) = oneshot::channel();
4261 command_sender
4262 .send(SessionCommand::ListJobs {
4263 timeout: Duration::from_secs(1),
4264 reply: second_reply,
4265 })
4266 .await
4267 .expect("same session accepts second command");
4268 assert_eq!(
4269 second_result.await.expect("second reply sent"),
4270 Ok(Vec::new())
4271 );
4272
4273 stop_sender.send(true).expect("session still running");
4274 assert!(
4275 !session_task.await.expect("session task joins"),
4276 "orderly stop must not be reported as a lost peer session"
4277 );
4278 peer_task.await.expect("test peer joins");
4279 }
4280
4281 #[tokio::test]
4282 async fn connected_peer_reconnects_after_transport_error() {
4283 let (client_io, mut peer_io) = tokio::io::duplex(4096);
4284 let (client_reader, client_writer) = tokio::io::split(client_io);
4285 let client = DcpClient::from_test_parts(client_reader, client_writer);
4286 let peer_task = tokio::spawn(async move {
4287 let _request_id = read_test_request(&mut peer_io).await;
4288 drop(peer_io);
4289 });
4290
4291 let (pipe_io, _pipe_peer_io) = tokio::io::duplex(4096);
4292 let (pipe_reader, pipe_writer) = tokio::io::split(pipe_io);
4293 let pipe = ShardPipeClient::from_test_parts(pipe_reader, pipe_writer);
4294 let (command_sender, mut commands) = mpsc::channel(1);
4295 let (_pipe_sender, mut pipe_commands) = mpsc::channel(1);
4296 let (_stop_sender, mut stop) = watch::channel(false);
4297 let session_task = tokio::spawn(async move {
4298 run_connected_peer(client, pipe, &mut commands, &mut pipe_commands, &mut stop).await
4299 });
4300
4301 let (reply, result) = oneshot::channel();
4302 command_sender
4303 .send(SessionCommand::ListJobs {
4304 timeout: Duration::from_secs(1),
4305 reply,
4306 })
4307 .await
4308 .expect("command accepted");
4309 let error = result
4310 .await
4311 .expect("transport error reply sent")
4312 .expect_err("closed transport fails request");
4313 assert!(!error.contains("timed out"), "unexpected error: {error}");
4314 assert!(
4315 tokio::time::timeout(Duration::from_secs(1), session_task)
4316 .await
4317 .expect("connected loop exits promptly")
4318 .expect("session task joins"),
4319 "transport error must ask the outer peer-session loop to reconnect"
4320 );
4321 peer_task.await.expect("test peer joins");
4322 }
4323
4324 #[test]
4325 fn peer_response_errors_do_not_request_reconnect() {
4326 for status in [
4327 ResponseStatus::Conflict,
4328 ResponseStatus::NotFound,
4329 ResponseStatus::DeadlineExceeded,
4330 ResponseStatus::BadRequest,
4331 ] {
4332 let (result, reconnect) = classify_peer_request::<()>(
4333 Ok(Err(DcpError::response(status, "application response"))),
4334 "StartJob",
4335 );
4336 assert!(result.is_err());
4337 assert!(
4338 !reconnect,
4339 "DCP response {status:?} must not tear down the peer session"
4340 );
4341 }
4342
4343 let (_result, reconnect) =
4344 classify_peer_request::<()>(Ok(Err(DcpError::Closed)), "StartJob");
4345 assert!(reconnect, "closed transports still reconnect");
4346 }
4347
4348 async fn read_test_request(stream: &mut DuplexStream) -> u64 {
4349 let frame = read_test_frame(stream).await;
4350 let request = match frame.frame {
4351 Some(crate::dcp::proto::dcp_frame::Frame::Request(request)) => request,
4352 _ => panic!("expected DCP request frame"),
4353 };
4354 assert!(matches!(
4355 request.command,
4356 Some(crate::dcp::proto::request::Command::ListJobs(_))
4357 ));
4358 request.request_id
4359 }
4360
4361 async fn read_test_frame(stream: &mut DuplexStream) -> crate::dcp::proto::DcpFrame {
4362 let mut header = [0_u8; 4];
4363 stream
4364 .read_exact(&mut header)
4365 .await
4366 .expect("DCP frame header");
4367 let mut payload = vec![0_u8; u32::from_be_bytes(header) as usize];
4368 stream
4369 .read_exact(&mut payload)
4370 .await
4371 .expect("DCP frame payload");
4372 crate::dcp::proto::DcpFrame::decode(payload.as_slice()).expect("valid DCP frame")
4373 }
4374
4375 async fn write_empty_job_list(stream: &mut DuplexStream, request_id: u64) {
4376 let payload = crate::dcp::proto::JobList { jobs: Vec::new() }.encode_to_vec();
4377 let frame = crate::dcp::proto::DcpFrame::response(crate::dcp::proto::Response::ok(
4378 request_id, payload,
4379 ));
4380 let payload = frame.encode_to_vec();
4381 let len = u32::try_from(payload.len()).expect("test DCP frame length");
4382 stream
4383 .write_all(&len.to_be_bytes())
4384 .await
4385 .expect("write DCP frame header");
4386 stream
4387 .write_all(&payload)
4388 .await
4389 .expect("write DCP frame payload");
4390 stream.flush().await.expect("flush DCP frame");
4391 }
4392
4393 #[test]
4394 fn remembered_assignment_generation_is_monotonic() {
4395 let mut assignments = BTreeMap::new();
4396 assert!(remember_assignment(
4397 &mut assignments,
4398 "job".to_owned(),
4399 assignment(2, "current-owner"),
4400 Some("coordinator"),
4401 ));
4402
4403 assert!(!remember_assignment(
4404 &mut assignments,
4405 "job".to_owned(),
4406 assignment(1, "stale-owner"),
4407 Some("coordinator"),
4408 ));
4409 assert_eq!(assignments["job"].assigned_node, "current-owner");
4410 assert_eq!(assignments["job"].placement_generation, 2);
4411
4412 assert!(!remember_assignment(
4413 &mut assignments,
4414 "job".to_owned(),
4415 assignment(2, "equal-generation-owner"),
4416 Some("coordinator"),
4417 ));
4418 assert_eq!(assignments["job"].assigned_node, "current-owner");
4419
4420 let mut stale_same_owner = assignment(2, "current-owner");
4421 stale_same_owner.coordinator_node = "old-coordinator".to_owned();
4422 assert!(!remember_assignment(
4423 &mut assignments,
4424 "job".to_owned(),
4425 stale_same_owner,
4426 Some("new-coordinator"),
4427 ));
4428 assert_eq!(assignments["job"].coordinator_node, "coordinator");
4429
4430 let mut same_owner_refresh = assignment(2, "current-owner");
4431 same_owner_refresh.coordinator_node = "new-coordinator".to_owned();
4432 assert!(remember_assignment(
4433 &mut assignments,
4434 "job".to_owned(),
4435 same_owner_refresh,
4436 Some("new-coordinator"),
4437 ));
4438 assert_eq!(assignments["job"].assigned_node, "current-owner");
4439 assert_eq!(assignments["job"].coordinator_node, "new-coordinator");
4440
4441 assert!(remember_assignment(
4442 &mut assignments,
4443 "job".to_owned(),
4444 assignment(3, "new-owner"),
4445 Some("new-coordinator"),
4446 ));
4447 assert_eq!(assignments["job"].assigned_node, "new-owner");
4448 assert_eq!(assignments["job"].placement_generation, 3);
4449 }
4450
4451 #[test]
4452 fn registry_rebuild_preserves_tombstone_against_stale_running_observations() {
4453 let tombstone = AssignmentTombstone {
4454 placement_generation: 2,
4455 coordinator_node: "new-coordinator".to_owned(),
4456 };
4457
4458 let older = assignment_with_coordinator(1, "stale-owner", "old-coordinator");
4459 assert!(registry_observation_is_tombstoned(&tombstone, &older));
4460
4461 let equal = assignment_with_coordinator(2, "stale-owner", "new-coordinator");
4462 assert!(registry_observation_is_tombstoned(&tombstone, &equal));
4463
4464 let newer = assignment_with_coordinator(3, "new-owner", "new-coordinator");
4465 assert!(!registry_observation_is_tombstoned(&tombstone, &newer));
4466 }
4467
4468 #[tokio::test]
4469 async fn stale_tombstone_does_not_stop_newer_local_owner() {
4470 let registry = JobRegistry::start(AgentConfig::default()).expect("registry starts");
4471 let sessions = test_sessions("node-a");
4472 let mut actor = test_placement_actor_with_coordinator(
4473 "node-a",
4474 "new-coordinator",
4475 registry.clone(),
4476 sessions.clone(),
4477 );
4478 let mut current = assignment(2, "node-a");
4479 current.coordinator_node = "new-coordinator".to_owned();
4480 actor
4481 .assignments
4482 .insert("local-job".to_owned(), current.clone());
4483 submit_running_cluster_job(®istry, "local-job", current.clone());
4484
4485 let stale = actor
4486 .remember(
4487 "local-job".to_owned(),
4488 None,
4489 Some(AssignmentTombstone {
4490 placement_generation: 1,
4491 coordinator_node: "old-coordinator".to_owned(),
4492 }),
4493 )
4494 .await;
4495 assert!(
4496 matches!(
4497 stale,
4498 Err(DcpError::Response {
4499 status: ResponseStatus::Conflict,
4500 ..
4501 })
4502 ),
4503 "non-current coordinator tombstone must not be falsely acknowledged: {stale:?}"
4504 );
4505 let status = registry.status("local-job").expect("local job status");
4506 assert_eq!(status.state, JobState::Running);
4507 assert!(actor.assignments.contains_key("local-job"));
4508
4509 actor
4510 .remember(
4511 "local-job".to_owned(),
4512 None,
4513 Some(AssignmentTombstone {
4514 placement_generation: 2,
4515 coordinator_node: "new-coordinator".to_owned(),
4516 }),
4517 )
4518 .await
4519 .expect("current tombstone applies");
4520 let status = registry.status("local-job").expect("local job status");
4521 assert_eq!(status.state, JobState::Stopped);
4522 assert!(!actor.assignments.contains_key("local-job"));
4523
4524 sessions.shutdown().await;
4525 registry.shutdown().expect("registry shuts down");
4526 }
4527
4528 #[tokio::test]
4529 async fn delayed_active_after_recorded_tombstone_is_rejected() {
4530 let registry = JobRegistry::start(AgentConfig::default()).expect("registry starts");
4531 let sessions = test_sessions("node-a");
4532 let mut actor = test_placement_actor_with_coordinator(
4533 "node-a",
4534 "new-coordinator",
4535 registry.clone(),
4536 sessions.clone(),
4537 );
4538
4539 actor
4540 .remember(
4541 "local-job".to_owned(),
4542 None,
4543 Some(AssignmentTombstone {
4544 placement_generation: 2,
4545 coordinator_node: "new-coordinator".to_owned(),
4546 }),
4547 )
4548 .await
4549 .expect("valid no-existing-assignment tombstone is retained");
4550 assert_eq!(
4551 actor
4552 .assignment_tombstones
4553 .get("local-job")
4554 .map(|tombstone| (
4555 tombstone.placement_generation,
4556 tombstone.coordinator_node.as_str()
4557 )),
4558 Some((2, "new-coordinator"))
4559 );
4560
4561 let mut stale = assignment(2, "node-b");
4562 stale.coordinator_node = "new-coordinator".to_owned();
4563 let result = actor
4564 .remember(
4565 "local-job".to_owned(),
4566 Some(wire_cluster_job_start(&stale)),
4567 None,
4568 )
4569 .await;
4570 assert!(matches!(
4571 result,
4572 Err(DcpError::Response {
4573 status: ResponseStatus::Conflict,
4574 ..
4575 })
4576 ));
4577 assert!(!actor.assignments.contains_key("local-job"));
4578 assert!(actor.assignment_tombstones.contains_key("local-job"));
4579
4580 let mut newer = assignment(3, "node-b");
4581 newer.coordinator_node = "new-coordinator".to_owned();
4582 actor
4583 .remember(
4584 "local-job".to_owned(),
4585 Some(wire_cluster_job_start(&newer)),
4586 None,
4587 )
4588 .await
4589 .expect("strictly newer active assignment supersedes tombstone");
4590 assert_eq!(actor.assignments["local-job"].placement_generation, 3);
4591 assert!(!actor.assignment_tombstones.contains_key("local-job"));
4592
4593 sessions.shutdown().await;
4594 registry.shutdown().expect("registry shuts down");
4595 }
4596
4597 #[tokio::test]
4598 async fn new_coordinator_tombstone_applies_before_same_generation_refresh() {
4599 let registry = JobRegistry::start(AgentConfig::default()).expect("registry starts");
4600 let sessions = test_sessions("node-a");
4601 let mut actor = test_placement_actor_with_coordinator(
4602 "node-a",
4603 "new-coordinator",
4604 registry.clone(),
4605 sessions.clone(),
4606 );
4607 let mut local = assignment(2, "node-a");
4608 local.coordinator_node = "old-coordinator".to_owned();
4609 actor
4610 .assignments
4611 .insert("local-job".to_owned(), local.clone());
4612 submit_running_cluster_job(®istry, "local-job", local);
4613
4614 actor
4615 .remember(
4616 "local-job".to_owned(),
4617 None,
4618 Some(AssignmentTombstone {
4619 placement_generation: 2,
4620 coordinator_node: "new-coordinator".to_owned(),
4621 }),
4622 )
4623 .await
4624 .expect("current coordinator tombstone applies across handoff");
4625 let status = registry.status("local-job").expect("local job status");
4626 assert_eq!(status.state, JobState::Stopped);
4627 assert!(!actor.assignments.contains_key("local-job"));
4628 assert!(actor.assignment_tombstones.contains_key("local-job"));
4629
4630 let mut same_generation_refresh = assignment(2, "node-a");
4631 same_generation_refresh.coordinator_node = "new-coordinator".to_owned();
4632 let result = actor
4633 .remember(
4634 "local-job".to_owned(),
4635 Some(wire_cluster_job_start(&same_generation_refresh)),
4636 None,
4637 )
4638 .await;
4639 assert!(matches!(
4640 result,
4641 Err(DcpError::Response {
4642 status: ResponseStatus::Conflict,
4643 ..
4644 })
4645 ));
4646 assert!(!actor.assignments.contains_key("local-job"));
4647 assert!(actor.assignment_tombstones.contains_key("local-job"));
4648
4649 sessions.shutdown().await;
4650 registry.shutdown().expect("registry shuts down");
4651 }
4652
4653 #[tokio::test]
4654 async fn current_tombstone_refreshes_terminal_local_metadata() {
4655 let registry = JobRegistry::start(AgentConfig::default()).expect("registry starts");
4656 let sessions = test_sessions("node-a");
4657 let mut actor = test_placement_actor_with_coordinator(
4658 "node-a",
4659 "new-coordinator",
4660 registry.clone(),
4661 sessions.clone(),
4662 );
4663 let mut local = assignment(2, "node-a");
4664 local.coordinator_node = "old-coordinator".to_owned();
4665 submit_running_cluster_job(®istry, "local-job", local);
4666 registry.stop("local-job").expect("local job stopped");
4667
4668 actor
4669 .remember(
4670 "local-job".to_owned(),
4671 None,
4672 Some(AssignmentTombstone {
4673 placement_generation: 2,
4674 coordinator_node: "new-coordinator".to_owned(),
4675 }),
4676 )
4677 .await
4678 .expect("current coordinator tombstone refreshes terminal metadata");
4679
4680 let status = registry.status("local-job").expect("local job status");
4681 assert_eq!(status.state, JobState::Stopped);
4682 let metadata = status.cluster.expect("cluster metadata retained");
4683 assert_eq!(metadata.placement_generation, 2);
4684 assert_eq!(metadata.coordinator_node, "new-coordinator");
4685 assert!(!actor.assignments.contains_key("local-job"));
4686 assert!(actor.assignment_tombstones.contains_key("local-job"));
4687
4688 sessions.shutdown().await;
4689 registry.shutdown().expect("registry shuts down");
4690 }
4691
4692 #[tokio::test]
4693 async fn delayed_active_replacement_stops_stale_local_owner() {
4694 let registry = JobRegistry::start(AgentConfig::default()).expect("registry starts");
4695 let sessions = test_sessions("node-a");
4696 let mut actor = test_placement_actor_with_coordinator(
4697 "node-a",
4698 "new-coordinator",
4699 registry.clone(),
4700 sessions.clone(),
4701 );
4702 let mut stale = assignment(1, "node-a");
4703 stale.coordinator_node = "old-coordinator".to_owned();
4704 actor
4705 .assignments
4706 .insert("local-job".to_owned(), stale.clone());
4707 submit_running_cluster_job(®istry, "local-job", stale);
4708
4709 let mut replacement = assignment(2, "node-b");
4710 replacement.coordinator_node = "new-coordinator".to_owned();
4711 actor
4712 .remember(
4713 "local-job".to_owned(),
4714 Some(wire_cluster_job_start(&replacement)),
4715 None,
4716 )
4717 .await
4718 .expect("delayed active replacement is accepted");
4719
4720 let status = registry.status("local-job").expect("local job status");
4721 assert_eq!(status.state, JobState::Stopped);
4722 assert_eq!(actor.assignments["local-job"].assigned_node, "node-b");
4723 assert_eq!(actor.assignments["local-job"].placement_generation, 2);
4724
4725 sessions.shutdown().await;
4726 registry.shutdown().expect("registry shuts down");
4727 }
4728
4729 #[tokio::test]
4730 async fn register_restarted_rejects_stale_local_restart_over_active_assignment() {
4731 let registry = JobRegistry::start(AgentConfig::default()).expect("registry starts");
4732 let sessions = test_sessions("node-a");
4733 let mut actor = test_placement_actor_with_coordinator(
4734 "node-a",
4735 "new-coordinator",
4736 registry.clone(),
4737 sessions.clone(),
4738 );
4739
4740 let mut stale_higher = assignment(1, "node-a");
4741 stale_higher.coordinator_node = "new-coordinator".to_owned();
4742 actor.assignments.insert(
4743 "higher-job".to_owned(),
4744 assignment_with_coordinator(2, "node-b", "new-coordinator"),
4745 );
4746 submit_running_cluster_job(®istry, "higher-job", stale_higher.clone());
4747 let result = actor
4748 .register_restarted(
4749 "higher-job".to_owned(),
4750 wire_cluster_job_start(&stale_higher),
4751 Duration::from_millis(50),
4752 )
4753 .await;
4754 assert!(matches!(
4755 result,
4756 Err(DcpError::Response {
4757 status: ResponseStatus::Conflict,
4758 ..
4759 })
4760 ));
4761 assert_eq!(actor.assignments["higher-job"].assigned_node, "node-b");
4762 assert_eq!(actor.assignments["higher-job"].placement_generation, 2);
4763 assert_eq!(
4764 registry.status("higher-job").expect("higher status").state,
4765 JobState::Stopped
4766 );
4767
4768 let mut stale_equal = assignment(2, "node-a");
4769 stale_equal.coordinator_node = "new-coordinator".to_owned();
4770 actor.assignments.insert(
4771 "equal-job".to_owned(),
4772 assignment_with_coordinator(2, "node-b", "new-coordinator"),
4773 );
4774 submit_running_cluster_job(®istry, "equal-job", stale_equal.clone());
4775 let result = actor
4776 .register_restarted(
4777 "equal-job".to_owned(),
4778 wire_cluster_job_start(&stale_equal),
4779 Duration::from_millis(50),
4780 )
4781 .await;
4782 assert!(matches!(
4783 result,
4784 Err(DcpError::Response {
4785 status: ResponseStatus::Conflict,
4786 ..
4787 })
4788 ));
4789 assert_eq!(actor.assignments["equal-job"].assigned_node, "node-b");
4790 assert_eq!(actor.assignments["equal-job"].placement_generation, 2);
4791 assert_eq!(
4792 registry.status("equal-job").expect("equal status").state,
4793 JobState::Stopped
4794 );
4795
4796 sessions.shutdown().await;
4797 registry.shutdown().expect("registry shuts down");
4798 }
4799
4800 #[tokio::test]
4801 async fn tombstone_stop_error_keeps_assignment_for_retry() {
4802 let registry = JobRegistry::start(AgentConfig::default()).expect("registry starts");
4803 let sessions = test_sessions("node-a");
4804 let mut actor = test_placement_actor_with_coordinator(
4805 "node-a",
4806 "new-coordinator",
4807 registry.clone(),
4808 sessions.clone(),
4809 );
4810 let mut current = assignment(2, "node-a");
4811 current.coordinator_node = "new-coordinator".to_owned();
4812 actor
4813 .assignments
4814 .insert("local-job".to_owned(), current.clone());
4815 registry
4816 .shutdown()
4817 .expect("registry shuts down before cleanup");
4818
4819 let result = actor
4820 .remember(
4821 "local-job".to_owned(),
4822 None,
4823 Some(AssignmentTombstone {
4824 placement_generation: 2,
4825 coordinator_node: "new-coordinator".to_owned(),
4826 }),
4827 )
4828 .await;
4829 assert!(result.is_err());
4830 assert!(
4831 actor.assignments.contains_key("local-job"),
4832 "failed cleanup must leave assignment for the next retry"
4833 );
4834
4835 sessions.shutdown().await;
4836 }
4837
4838 #[tokio::test]
4839 async fn membership_coordinator_and_placement_sources_publish_cluster_events() {
4840 let publisher = ClusterEventPublisher::new(8).expect("cluster event topic");
4841 let mut events = publisher.subscribe().expect("cluster event subscription");
4842 let at = UNIX_EPOCH + Duration::from_millis(1_720_000_000_000);
4843 publisher.publish_member(&MemberEvent {
4844 sequence: 11,
4845 kind: MemberEventKind::MemberDown,
4846 member: Member {
4847 node_id: "node-1".to_owned(),
4848 address: "127.0.0.1:2551".parse().expect("member address"),
4849 agent_addr: Some("127.0.0.1:9555".parse().expect("agent address")),
4850 roles: vec![AGENT_ROLE.to_owned()],
4851 state: MemberState::Down,
4852 unreachable: true,
4853 incarnation: 4,
4854 seen_at: at,
4855 unreachable_since: Some(at),
4856 },
4857 at,
4858 });
4859 publisher.publish_coordinator_changed(Some("node-1"), Some("node-2"), "node-2");
4860 publisher.publish_job_replaced(at, "ingest", "node-1", "node-2", "node_down:node-1", 3);
4861
4862 let mut received = Vec::new();
4863 for _ in 0..3 {
4864 received.push(
4865 tokio::time::timeout(Duration::from_secs(1), events.recv())
4866 .await
4867 .expect("cluster event arrives")
4868 .expect("cluster event subscription remains open"),
4869 );
4870 }
4871
4872 assert_eq!(
4873 received
4874 .iter()
4875 .map(|event| (event.sequence, event.kind.as_str(), event.node_id.as_str()))
4876 .collect::<Vec<_>>(),
4877 vec![
4878 (1, "MemberDown", "node-1"),
4879 (2, "CoordinatorChanged", "node-2"),
4880 (3, "JobReplaced", "node-2"),
4881 ]
4882 );
4883 assert_eq!(received[0].generation, Some(4));
4884 assert_eq!(received[2].generation, Some(3));
4885 assert!(received[2].detail.contains("job=ingest"));
4886 }
4887
4888 fn assignment(placement_generation: u64, assigned_node: &str) -> ClusterJobMetadata {
4889 ClusterJobMetadata {
4890 factory_name: "ticker".to_owned(),
4891 params: BTreeMap::new(),
4892 placement: RegistryPlacementSpec::least_jobs(None),
4893 coordinator_node: "coordinator".to_owned(),
4894 assigned_node: assigned_node.to_owned(),
4895 placement_generation,
4896 history: Vec::new(),
4897 }
4898 }
4899
4900 fn assignment_with_coordinator(
4901 placement_generation: u64,
4902 assigned_node: &str,
4903 coordinator_node: &str,
4904 ) -> ClusterJobMetadata {
4905 let mut metadata = assignment(placement_generation, assigned_node);
4906 metadata.coordinator_node = coordinator_node.to_owned();
4907 metadata
4908 }
4909
4910 fn submit_running_cluster_job(
4911 registry: &JobRegistryHandle,
4912 name: &str,
4913 metadata: ClusterJobMetadata,
4914 ) {
4915 let spec = JobSpec::new(name, |context| {
4916 let control = context.control();
4917 Ok(Source::tick(Duration::ZERO, Duration::from_secs(60), 1_u64)
4918 .via_mat(context.drain_flow(), Keep::right)
4919 .to_mat(Sink::ignore(), move |_switch, completion| {
4920 JobMat::new(completion, control.clone())
4921 }))
4922 })
4923 .with_cluster_metadata(metadata);
4924 registry.submit(spec).expect("job submitted");
4925 registry.start(name).expect("job starts");
4926 }
4927
4928 fn test_sessions(node_id: &str) -> NodeSessionManagerHandle {
4929 let state = Signal::new(test_cluster_state(node_id)).expect("cluster state");
4930 NodeSessionManagerHandle::start(NodeSessionConfig::default(), node_id.to_owned(), state)
4931 .expect("sessions start")
4932 }
4933
4934 fn test_placement_actor_with_coordinator(
4935 node_id: &str,
4936 coordinator_node: &str,
4937 registry: JobRegistryHandle,
4938 sessions: NodeSessionManagerHandle,
4939 ) -> PlacementActor {
4940 let state = Signal::new(test_cluster_state_with_coordinator(
4941 node_id,
4942 coordinator_node,
4943 ))
4944 .expect("cluster state");
4945 let (commands, _receiver) = mpsc::channel(8);
4946 let (assignment_updates, _assignments) = watch::channel(Arc::new(BTreeMap::new()));
4947 PlacementActor {
4948 commands,
4949 self_node: node_id.to_owned(),
4950 agent_role: AGENT_ROLE.to_owned(),
4951 state,
4952 registry,
4953 sessions,
4954 factories: DcpJobFactories::new(),
4955 request_timeout: Duration::from_millis(50),
4956 active_coordinator: false,
4957 last_known_coordinator: None,
4958 assignments: BTreeMap::new(),
4959 assignment_tombstones: BTreeMap::new(),
4960 replicated_assignments: BTreeMap::new(),
4961 inflight_assignment_syncs: BTreeSet::new(),
4962 pending_replacements: BTreeMap::new(),
4963 assignment_updates,
4964 cluster_events: ClusterEventPublisher::new(8).expect("cluster events"),
4965 }
4966 }
4967
4968 fn test_cluster_state(node_id: &str) -> ClusterState {
4969 test_cluster_state_with_coordinator(node_id, node_id)
4970 }
4971
4972 fn test_cluster_state_with_coordinator(node_id: &str, coordinator_node: &str) -> ClusterState {
4973 let mut members = BTreeMap::new();
4974 members.insert(node_id.to_owned(), test_member_with_incarnation(node_id, 1));
4975 members.insert(
4976 coordinator_node.to_owned(),
4977 test_member_with_incarnation(coordinator_node, 0),
4978 );
4979 ClusterState {
4980 self_node: node_id.to_owned(),
4981 members,
4982 }
4983 }
4984
4985 fn test_member_with_incarnation(node_id: &str, incarnation: u64) -> Member {
4986 Member {
4987 node_id: node_id.to_owned(),
4988 address: "127.0.0.1:2551".parse().expect("member address"),
4989 agent_addr: Some("127.0.0.1:9555".parse().expect("agent address")),
4990 roles: vec![AGENT_ROLE.to_owned()],
4991 state: MemberState::Up,
4992 unreachable: false,
4993 incarnation,
4994 seen_at: UNIX_EPOCH,
4995 unreachable_since: None,
4996 }
4997 }
4998}