Skip to main content

datum_agent/
cluster.rs

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
41/// Role used by cluster-aware agents to discover peers.
42pub const AGENT_ROLE: &str = "agent";
43
44/// Result type for cluster-aware agent startup and shutdown.
45pub type ClusterAgentResult<T> = Result<T, ClusterAgentError>;
46
47/// Errors returned by the cluster-aware agent entrypoint.
48#[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/// Transport used for node-to-node DCP sessions.
61#[derive(Clone)]
62pub enum NodeSessionTransport {
63    /// Plaintext TCP is accepted only for loopback development and tests.
64    TcpLoopback,
65    /// QUIC with mTLS. The same client config is used for every peer.
66    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/// Configuration for DCP node sessions maintained by a cluster-aware agent.
85#[derive(Clone, Debug)]
86pub struct NodeSessionConfig {
87    /// Role that marks a member as running a DCP agent endpoint.
88    pub agent_role: String,
89    /// Node-to-node DCP transport.
90    pub transport: NodeSessionTransport,
91    /// Initial reconnect backoff.
92    pub reconnect_min_backoff: Duration,
93    /// Maximum reconnect backoff.
94    pub reconnect_max_backoff: Duration,
95    /// Default per-peer fan-out timeout.
96    pub request_timeout: Duration,
97    /// Bounded per-session command queue capacity.
98    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/// Combined configuration for one cluster-aware Datum agent.
115#[derive(Clone, Default)]
116pub struct ClusterAgentConfig {
117    /// Local job registry configuration.
118    pub agent: AgentConfig,
119    /// Membership configuration.
120    pub cluster: ClusterConfig,
121    /// DCP listener configuration.
122    pub dcp: DcpServerConfig,
123    /// Node-session configuration.
124    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
139/// Entrypoint for starting membership, the local agent, DCP, and node sessions
140/// together.
141pub 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
210/// Running cluster-aware agent.
211pub 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        // Live registries are authoritative for both running intent and stopped tombstones. The
2472        // new coordinator's replicated table is retained for jobs whose owning node is already
2473        // down and therefore cannot participate in this handoff scan. Peer replica maps are not
2474        // exchanged here, so a dead-owner job is recoverable only when its upsert reached the
2475        // specific node that adopts coordination.
2476        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        // Adoption is deliberately fail-closed until every member still considered Up responds.
2887        // An unresponsive former coordinator or peer delays adoption until the failure detector
2888        // marks it Down (bounded by suspect_timeout), which prevents a surviving job from being
2889        // omitted. A genuinely Up but slow peer, such as one in a GC pause, defers adoption too.
2890        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            // Do not filter by the local placement assignment cache here. A restarted observer can
3818            // rejoin with connected peer sessions before any assignment replay or later placement
3819            // mutation reaches its local placement actor, but cluster metrics must still bootstrap.
3820            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(&registry, "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(&registry, "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(&registry, "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(&registry, "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(&registry, "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(&registry, "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}