Skip to main content

mesh_client/client/
control_plane.rs

1use crate::client::builder::MeshClient;
2use crate::crypto::OwnerKeypair;
3use crate::proto::node::{
4    NodeConfigSnapshot, OwnerControlApplyConfigRequest, OwnerControlApplyConfigResponse,
5    OwnerControlConfigSnapshot, OwnerControlConfigUpdate, OwnerControlDrainModelRequest,
6    OwnerControlDrainModelResponse, OwnerControlEnsureModelRequest,
7    OwnerControlEnsureModelResponse, OwnerControlEnvelope, OwnerControlError,
8    OwnerControlErrorCode, OwnerControlGetConfigRequest, OwnerControlHandshake,
9    OwnerControlLoadModelRequest, OwnerControlLoadModelResponse, OwnerControlModelRef,
10    OwnerControlRefreshInventory, OwnerControlRefreshInventoryRequest, OwnerControlRequest,
11    OwnerControlResponse, OwnerControlUnloadModelRequest, OwnerControlUnloadModelResponse,
12    OwnerControlWatchAccepted, OwnerControlWatchConfigRequest, OwnerControlWatchConfigResponse,
13    SignedNodeOwnership,
14};
15use crate::protocol::{
16    ALPN_CONTROL_V1, ALPN_V1, NODE_PROTOCOL_GENERATION, decode_owner_control_envelope,
17    write_len_prefixed,
18};
19use anyhow::Context;
20use base64::Engine;
21use iroh::{Endpoint, EndpointAddr};
22use prost::Message;
23use std::fmt;
24use std::sync::atomic::{AtomicU64, Ordering};
25use thiserror::Error;
26
27const DEFAULT_NODE_CERT_LIFETIME_SECS: u64 = 7 * 24 * 60 * 60;
28const NODE_OWNERSHIP_VERSION: u32 = 1;
29const SIGNING_DOMAIN_TAG: &[u8] = b"mesh-llm-node-ownership-v1:";
30const OWNER_CONTROL_CONNECT_TIMEOUT_SECS: u64 = 8;
31const OWNER_CONTROL_OPEN_TIMEOUT_SECS: u64 = 2;
32const OWNER_CONTROL_HANDSHAKE_TIMEOUT_SECS: u64 = 2;
33const OWNER_CONTROL_REQUEST_WRITE_TIMEOUT_SECS: u64 = 2;
34const OWNER_CONTROL_SERVER_UNARY_DEADLINE_SECS_FOR_CLIENT_MARGIN: u64 = 5;
35const OWNER_CONTROL_UNARY_RESPONSE_TIMEOUT_SECS: u64 =
36    OWNER_CONTROL_SERVER_UNARY_DEADLINE_SECS_FOR_CLIENT_MARGIN + 5;
37const OWNER_CONTROL_SERVER_SCAN_DEADLINE_SECS_FOR_CLIENT_MARGIN: u64 = 30;
38const OWNER_CONTROL_INVENTORY_RESPONSE_TIMEOUT_SECS: u64 =
39    OWNER_CONTROL_SERVER_SCAN_DEADLINE_SECS_FOR_CLIENT_MARGIN + 5;
40const OWNER_CONTROL_WATCH_ACCEPT_TIMEOUT_SECS: u64 = 5;
41const FAILED_BOOTSTRAP_CLOSE_TIMEOUT_MILLIS: u64 = 250;
42
43fn owner_control_client_bind_addr() -> std::net::SocketAddr {
44    std::net::SocketAddr::from(([0, 0, 0, 0], 0))
45}
46
47/// Explicit owner-control bootstrap policy for new config clients.
48///
49/// Negotiation matrix:
50/// - new client + explicit control endpoint -> use `mesh-llm-control/1`; configured control
51///   failures stay on the control lane and return structured errors.
52/// - new client + no control endpoint -> return `ControlEndpointRequired`.
53///
54/// Config and inventory mutation is intentionally exclusive to `mesh-llm-control/1`.
55/// The legacy mesh-plane config stream IDs remain reserved, but no client bootstrap path
56/// falls back to them.
57#[derive(Clone, Debug, PartialEq, Eq)]
58pub struct ControlPlaneBootstrapOptions {
59    control_endpoint: Option<String>,
60    connect_timeout: std::time::Duration,
61}
62
63impl Default for ControlPlaneBootstrapOptions {
64    fn default() -> Self {
65        Self {
66            control_endpoint: None,
67            connect_timeout: std::time::Duration::from_secs(OWNER_CONTROL_CONNECT_TIMEOUT_SECS),
68        }
69    }
70}
71
72impl ControlPlaneBootstrapOptions {
73    pub fn new() -> Self {
74        Self::default()
75    }
76
77    pub fn with_control_endpoint(mut self, control_endpoint: impl Into<String>) -> Self {
78        self.control_endpoint = Some(control_endpoint.into());
79        self
80    }
81
82    pub fn control_endpoint(&self) -> Option<&str> {
83        self.control_endpoint.as_deref()
84    }
85
86    /// Bound owner-control endpoint bootstrap without changing request deadlines.
87    pub fn with_connect_timeout(mut self, timeout: std::time::Duration) -> Self {
88        self.connect_timeout = timeout;
89        self
90    }
91
92    pub fn select_transport(
93        &self,
94    ) -> Result<ConfigTransportSelection, ControlPlaneNegotiationError> {
95        match self.control_endpoint() {
96            Some(endpoint) => Ok(ConfigTransportSelection::OwnerControl {
97                endpoint: endpoint.to_string(),
98                retry_policy: ControlPlaneRetryPolicy::NoSilentLegacyDowngrade,
99            }),
100            None => Err(ControlPlaneNegotiationError::endpoint_required()),
101        }
102    }
103
104    pub fn configured_endpoint_failure(
105        &self,
106        code: OwnerControlErrorCode,
107        message: impl Into<String>,
108    ) -> ControlPlaneNegotiationError {
109        debug_assert!(self.control_endpoint.is_some());
110        ControlPlaneNegotiationError::structured(code, message, false)
111    }
112}
113
114#[derive(Clone, Debug, PartialEq, Eq)]
115pub enum ConfigTransportSelection {
116    OwnerControl {
117        endpoint: String,
118        retry_policy: ControlPlaneRetryPolicy,
119    },
120}
121
122#[derive(Clone, Copy, Debug, PartialEq, Eq)]
123pub enum ControlPlaneRetryPolicy {
124    NoSilentLegacyDowngrade,
125}
126
127#[derive(Clone, Debug, PartialEq, Eq)]
128pub struct ControlPlaneNegotiationError {
129    pub code: OwnerControlErrorCode,
130    pub message: String,
131    pub legacy_retry_allowed: bool,
132}
133
134impl fmt::Display for ControlPlaneNegotiationError {
135    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
136        write!(f, "{:?}: {}", self.code, self.message)
137    }
138}
139
140impl std::error::Error for ControlPlaneNegotiationError {}
141
142impl ControlPlaneNegotiationError {
143    pub fn endpoint_required() -> Self {
144        Self {
145            code: OwnerControlErrorCode::ControlEndpointRequired,
146            message: "owner-control endpoint must be provided explicitly".to_string(),
147            legacy_retry_allowed: false,
148        }
149    }
150
151    pub fn structured(
152        code: OwnerControlErrorCode,
153        message: impl Into<String>,
154        legacy_retry_allowed: bool,
155    ) -> Self {
156        Self {
157            code,
158            message: message.into(),
159            legacy_retry_allowed,
160        }
161    }
162}
163
164#[derive(Debug, Error)]
165pub enum ControlPlaneClientError {
166    #[error(transparent)]
167    Negotiation(#[from] ControlPlaneNegotiationError),
168    #[error(transparent)]
169    Remote(#[from] OwnerControlRemoteError),
170    #[error("control transport error: {0}")]
171    Transport(String),
172    #[error("control protocol error: {0}")]
173    Protocol(String),
174}
175
176#[derive(Clone, Debug, PartialEq, Eq)]
177pub struct OwnerControlRemoteError {
178    pub code: OwnerControlErrorCode,
179    pub message: String,
180    pub request_id: Option<u64>,
181    pub current_revision: Option<u64>,
182}
183
184impl fmt::Display for OwnerControlRemoteError {
185    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
186        write!(f, "{:?}: {}", self.code, self.message)
187    }
188}
189
190impl std::error::Error for OwnerControlRemoteError {}
191
192impl From<OwnerControlError> for OwnerControlRemoteError {
193    fn from(error: OwnerControlError) -> Self {
194        Self {
195            code: OwnerControlErrorCode::try_from(error.code)
196                .unwrap_or(OwnerControlErrorCode::BadRequest),
197            message: error.message,
198            request_id: error.request_id,
199            current_revision: error.current_revision,
200        }
201    }
202}
203
204/// Control-plane bootstrap is explicit and out-of-band.
205///
206/// Callers either receive an owner-control session bound to a configured endpoint,
207/// or a structured error. The client never performs a silent downgrade.
208pub enum ControlPlaneConnection {
209    OwnerControl(Box<OwnerControlClient>),
210}
211
212pub struct OwnerControlClient {
213    endpoint_token: String,
214    endpoint: Endpoint,
215    connection: iroh::endpoint::Connection,
216    owner_keypair: OwnerKeypair,
217    next_request_id: AtomicU64,
218}
219
220pub struct OwnerControlWatchStream {
221    send: iroh::endpoint::SendStream,
222    recv: iroh::endpoint::RecvStream,
223    request_id: u64,
224    pending: Option<OwnerControlWatchEvent>,
225    closed: bool,
226}
227
228pub enum OwnerControlWatchEvent {
229    Accepted(OwnerControlWatchAccepted),
230    Snapshot(OwnerControlConfigSnapshot),
231    Update(OwnerControlConfigUpdate),
232}
233
234/// Result of a completed owner-control inventory scan.
235///
236/// Older servers return only the refreshed config snapshot. In that compatibility
237/// case, `inventory` is `None` while the command itself still succeeds.
238#[derive(Clone, Debug, PartialEq)]
239pub struct OwnerControlScanRefreshResult {
240    pub snapshot: OwnerControlConfigSnapshot,
241    pub inventory: Option<OwnerControlRefreshInventory>,
242}
243
244impl MeshClient {
245    /// Bootstrap config transport using the explicit owner-control endpoint policy.
246    ///
247    /// Owner-control endpoints are not discovered through gossip or status APIs;
248    /// callers must provide them explicitly through out-of-band bootstrap.
249    pub async fn connect_control_plane(
250        &self,
251        options: ControlPlaneBootstrapOptions,
252    ) -> Result<ControlPlaneConnection, ControlPlaneClientError> {
253        match options.select_transport()? {
254            ConfigTransportSelection::OwnerControl { endpoint, .. } => {
255                OwnerControlClient::connect(endpoint, self.config.owner_keypair.clone(), &options)
256                    .await
257                    .map(Box::new)
258                    .map(ControlPlaneConnection::OwnerControl)
259            }
260        }
261    }
262}
263
264fn validate_lifecycle_acceptance(
265    operation: &str,
266    intent_id: &str,
267    accepted_state: &str,
268    expected_state: &str,
269    target: Option<&crate::proto::node::OwnerControlModelRef>,
270    expected_model_ref: &str,
271    expected_instance_id: Option<&str>,
272) -> Result<(), ControlPlaneClientError> {
273    if intent_id.is_empty() {
274        return Err(ControlPlaneClientError::Protocol(format!(
275            "owner-control {operation} response missing intent id"
276        )));
277    }
278    if accepted_state != expected_state {
279        return Err(ControlPlaneClientError::Protocol(format!(
280            "owner-control {operation} response has invalid accepted state"
281        )));
282    }
283    let target = target.ok_or_else(|| {
284        ControlPlaneClientError::Protocol(format!(
285            "owner-control {operation} response missing target"
286        ))
287    })?;
288    if target.canonical_model_ref != expected_model_ref
289        || target.instance_id.as_deref() != expected_instance_id
290    {
291        return Err(ControlPlaneClientError::Protocol(format!(
292            "owner-control {operation} response target does not match request"
293        )));
294    }
295    Ok(())
296}
297
298fn map_legacy_lifecycle_unsupported(
299    operation: &str,
300    error: ControlPlaneClientError,
301) -> ControlPlaneClientError {
302    const LEGACY_UNKNOWN_COMMAND_MESSAGE: &str =
303        "owner control request requires exactly one command variant";
304    match error {
305        ControlPlaneClientError::Remote(mut remote)
306            if matches!(
307                remote.code,
308                OwnerControlErrorCode::BadRequest | OwnerControlErrorCode::UnknownCommand
309            ) && remote.message == LEGACY_UNKNOWN_COMMAND_MESSAGE =>
310        {
311            remote.code = OwnerControlErrorCode::ControlUnsupported;
312            remote.message = format!("remote owner-control endpoint does not support {operation}");
313            ControlPlaneClientError::Remote(remote)
314        }
315        other => other,
316    }
317}
318
319enum LifecycleCommand {
320    Load(OwnerControlLoadModelRequest),
321    Unload(OwnerControlUnloadModelRequest),
322    Ensure(OwnerControlEnsureModelRequest),
323    Drain(OwnerControlDrainModelRequest),
324}
325
326impl LifecycleCommand {
327    fn operation(&self) -> &'static str {
328        match self {
329            Self::Load(_) => "load_model",
330            Self::Unload(_) => "unload_model",
331            Self::Ensure(_) => "ensure_model",
332            Self::Drain(_) => "drain_model",
333        }
334    }
335
336    fn into_request(self, request_id: u64) -> OwnerControlRequest {
337        let mut request = OwnerControlRequest {
338            request_id,
339            ..Default::default()
340        };
341        match self {
342            Self::Load(command) => request.load_model = Some(command),
343            Self::Unload(command) => request.unload_model = Some(command),
344            Self::Ensure(command) => request.ensure_model = Some(command),
345            Self::Drain(command) => request.drain_model = Some(command),
346        }
347        request
348    }
349}
350
351impl OwnerControlClient {
352    async fn connect(
353        endpoint_token: String,
354        owner_keypair: OwnerKeypair,
355        options: &ControlPlaneBootstrapOptions,
356    ) -> Result<Self, ControlPlaneClientError> {
357        let control_addr = decode_endpoint_addr_token(&endpoint_token).map_err(|error| {
358            ControlPlaneClientError::Negotiation(options.configured_endpoint_failure(
359                OwnerControlErrorCode::ControlUnavailable,
360                format!("invalid owner-control endpoint token: {error}"),
361            ))
362        })?;
363        let mut builder = Endpoint::builder(iroh::endpoint::presets::Minimal)
364            .secret_key(iroh::SecretKey::generate())
365            .alpns(vec![ALPN_CONTROL_V1.to_vec()])
366            .bind_addr(owner_control_client_bind_addr())
367            .map_err(|error| ControlPlaneClientError::Transport(error.to_string()))?;
368        builder = builder.relay_mode(relay_mode_from_endpoint_addr(&control_addr));
369        let endpoint = builder
370            .bind()
371            .await
372            .map_err(|error| ControlPlaneClientError::Transport(error.to_string()))?;
373        if control_addr.relay_urls().next().is_some() {
374            let _ = tokio::time::timeout(options.connect_timeout, endpoint.online()).await;
375        }
376        let connection = match tokio::time::timeout(
377            options.connect_timeout,
378            endpoint.connect(control_addr.clone(), ALPN_CONTROL_V1),
379        )
380        .await
381        {
382            Ok(Ok(connection)) => connection,
383            Ok(Err(error)) => {
384                let error =
385                    configured_endpoint_connect_error(&endpoint, control_addr, options, error)
386                        .await;
387                close_failed_bootstrap_endpoint(&endpoint).await;
388                return Err(error);
389            }
390            Err(_) => {
391                close_failed_bootstrap_endpoint(&endpoint).await;
392                return Err(ControlPlaneClientError::Negotiation(options.configured_endpoint_failure(
393                    OwnerControlErrorCode::ControlUnavailable,
394                    format!(
395                        "remote owner-control endpoint is unavailable or unreachable: connect timed out after {:.3}s",
396                        options.connect_timeout.as_secs_f64()
397                    ),
398                )));
399            }
400        };
401        Ok(Self {
402            endpoint_token,
403            endpoint,
404            connection,
405            owner_keypair,
406            next_request_id: AtomicU64::new(1),
407        })
408    }
409
410    pub fn endpoint_token(&self) -> &str {
411        &self.endpoint_token
412    }
413
414    pub fn local_node_id(&self) -> [u8; 32] {
415        *self.endpoint.id().as_bytes()
416    }
417
418    pub fn target_node_id(&self) -> [u8; 32] {
419        *self.connection.remote_id().as_bytes()
420    }
421
422    pub async fn close(&self) {
423        self.connection
424            .close(0u32.into(), b"owner-control-client-close");
425        self.endpoint.close().await;
426    }
427
428    pub async fn get_config(&self) -> Result<OwnerControlConfigSnapshot, ControlPlaneClientError> {
429        let response = self
430            .send_unary_request(
431                std::time::Duration::from_secs(OWNER_CONTROL_UNARY_RESPONSE_TIMEOUT_SECS),
432                |request_id, requester_node_id, target_node_id| OwnerControlRequest {
433                    request_id,
434                    get_config: Some(OwnerControlGetConfigRequest {
435                        requester_node_id,
436                        target_node_id,
437                    }),
438                    watch_config: None,
439                    apply_config: None,
440                    refresh_inventory: None,
441                    load_model: None,
442                    unload_model: None,
443                    ensure_model: None,
444                    drain_model: None,
445                },
446            )
447            .await?;
448        response
449            .get_config
450            .and_then(|response| response.snapshot)
451            .ok_or_else(|| {
452                ControlPlaneClientError::Protocol(
453                    "owner-control get_config response missing snapshot payload".to_string(),
454                )
455            })
456    }
457
458    pub async fn apply_config(
459        &self,
460        expected_revision: u64,
461        config: NodeConfigSnapshot,
462    ) -> Result<OwnerControlApplyConfigResponse, ControlPlaneClientError> {
463        let response = self
464            .send_unary_request(
465                std::time::Duration::from_secs(OWNER_CONTROL_UNARY_RESPONSE_TIMEOUT_SECS),
466                |request_id, requester_node_id, target_node_id| OwnerControlRequest {
467                    request_id,
468                    get_config: None,
469                    watch_config: None,
470                    apply_config: Some(OwnerControlApplyConfigRequest {
471                        requester_node_id,
472                        target_node_id,
473                        expected_revision,
474                        config: Some(config),
475                    }),
476                    refresh_inventory: None,
477                    load_model: None,
478                    unload_model: None,
479                    ensure_model: None,
480                    drain_model: None,
481                },
482            )
483            .await?;
484        response.apply_config.ok_or_else(|| {
485            ControlPlaneClientError::Protocol(
486                "owner-control apply_config response missing apply payload".to_string(),
487            )
488        })
489    }
490
491    pub async fn refresh_inventory(
492        &self,
493    ) -> Result<OwnerControlConfigSnapshot, ControlPlaneClientError> {
494        self.scan_refresh().await.map(|result| result.snapshot)
495    }
496
497    pub async fn scan_refresh(
498        &self,
499    ) -> Result<OwnerControlScanRefreshResult, ControlPlaneClientError> {
500        let response = self
501            .send_unary_request(
502                std::time::Duration::from_secs(OWNER_CONTROL_INVENTORY_RESPONSE_TIMEOUT_SECS),
503                |request_id, requester_node_id, target_node_id| OwnerControlRequest {
504                    request_id,
505                    get_config: None,
506                    watch_config: None,
507                    apply_config: None,
508                    refresh_inventory: Some(OwnerControlRefreshInventoryRequest {
509                        requester_node_id,
510                        target_node_id,
511                    }),
512                    load_model: None,
513                    unload_model: None,
514                    ensure_model: None,
515                    drain_model: None,
516                },
517            )
518            .await?;
519        let response = response.refresh_inventory.ok_or_else(|| {
520            ControlPlaneClientError::Protocol(
521                "owner-control refresh_inventory response missing refresh payload".to_string(),
522            )
523        })?;
524        let snapshot = response.snapshot.ok_or_else(|| {
525            ControlPlaneClientError::Protocol(
526                "owner-control refresh_inventory response missing snapshot payload".to_string(),
527            )
528        })?;
529        Ok(OwnerControlScanRefreshResult {
530            snapshot,
531            inventory: response.inventory,
532        })
533    }
534
535    pub async fn load_model(
536        &self,
537        model_ref: String,
538        profile: Option<String>,
539    ) -> Result<OwnerControlLoadModelResponse, ControlPlaneClientError> {
540        let expected_model_ref = model_ref.clone();
541        let response = self
542            .send_lifecycle_request(LifecycleCommand::Load(OwnerControlLoadModelRequest {
543                requester_node_id: self.endpoint.id().as_bytes().to_vec(),
544                target_node_id: self.connection.remote_id().as_bytes().to_vec(),
545                model: Some(OwnerControlModelRef {
546                    canonical_model_ref: model_ref,
547                    instance_id: None,
548                }),
549                profile,
550            }))
551            .await?;
552        let response = response.load_model.ok_or_else(|| {
553            ControlPlaneClientError::Protocol(
554                "owner-control load_model response missing payload".to_string(),
555            )
556        })?;
557        validate_lifecycle_acceptance(
558            "load_model",
559            &response.intent_id,
560            &response.accepted_state,
561            "present",
562            response.target.as_ref(),
563            &expected_model_ref,
564            None,
565        )?;
566        Ok(response)
567    }
568
569    pub async fn unload_model(
570        &self,
571        model_ref: String,
572        instance_id: Option<String>,
573    ) -> Result<OwnerControlUnloadModelResponse, ControlPlaneClientError> {
574        let (expected_model_ref, expected_instance_id) =
575            validate_absent_model_target(model_ref, instance_id)?;
576        let response = self
577            .send_lifecycle_request(LifecycleCommand::Unload(OwnerControlUnloadModelRequest {
578                requester_node_id: self.endpoint.id().as_bytes().to_vec(),
579                target_node_id: self.connection.remote_id().as_bytes().to_vec(),
580                model: Some(OwnerControlModelRef {
581                    canonical_model_ref: expected_model_ref.clone(),
582                    instance_id: expected_instance_id.clone(),
583                }),
584            }))
585            .await?;
586        let response = response.unload_model.ok_or_else(|| {
587            ControlPlaneClientError::Protocol(
588                "owner-control unload_model response missing payload".to_string(),
589            )
590        })?;
591        validate_lifecycle_acceptance(
592            "unload_model",
593            &response.intent_id,
594            &response.accepted_state,
595            "absent",
596            response.target.as_ref(),
597            &expected_model_ref,
598            expected_instance_id.as_deref(),
599        )?;
600        Ok(response)
601    }
602
603    pub async fn ensure_model(
604        &self,
605        model_ref: String,
606        profile: Option<String>,
607    ) -> Result<OwnerControlEnsureModelResponse, ControlPlaneClientError> {
608        let expected_model_ref = model_ref.clone();
609        let response = self
610            .send_lifecycle_request(LifecycleCommand::Ensure(OwnerControlEnsureModelRequest {
611                requester_node_id: self.endpoint.id().as_bytes().to_vec(),
612                target_node_id: self.connection.remote_id().as_bytes().to_vec(),
613                model: Some(OwnerControlModelRef {
614                    canonical_model_ref: model_ref,
615                    instance_id: None,
616                }),
617                profile,
618            }))
619            .await?;
620        let response = response.ensure_model.ok_or_else(|| {
621            ControlPlaneClientError::Protocol(
622                "owner-control ensure_model response missing payload".to_string(),
623            )
624        })?;
625        validate_lifecycle_acceptance(
626            "ensure_model",
627            &response.intent_id,
628            &response.accepted_state,
629            "present",
630            response.target.as_ref(),
631            &expected_model_ref,
632            None,
633        )?;
634        Ok(response)
635    }
636
637    pub async fn drain_model(
638        &self,
639        model_ref: String,
640        instance_id: Option<String>,
641    ) -> Result<OwnerControlDrainModelResponse, ControlPlaneClientError> {
642        let (expected_model_ref, expected_instance_id) =
643            validate_absent_model_target(model_ref, instance_id)?;
644        let response = self
645            .send_lifecycle_request(LifecycleCommand::Drain(OwnerControlDrainModelRequest {
646                requester_node_id: self.endpoint.id().as_bytes().to_vec(),
647                target_node_id: self.connection.remote_id().as_bytes().to_vec(),
648                model: Some(OwnerControlModelRef {
649                    canonical_model_ref: expected_model_ref.clone(),
650                    instance_id: expected_instance_id.clone(),
651                }),
652                drain_timeout_secs: None,
653            }))
654            .await?;
655        let response = response.drain_model.ok_or_else(|| {
656            ControlPlaneClientError::Protocol(
657                "owner-control drain_model response missing payload".to_string(),
658            )
659        })?;
660        validate_lifecycle_acceptance(
661            "drain_model",
662            &response.intent_id,
663            &response.accepted_state,
664            "draining",
665            response.target.as_ref(),
666            &expected_model_ref,
667            expected_instance_id.as_deref(),
668        )?;
669        Ok(response)
670    }
671
672    async fn send_lifecycle_request(
673        &self,
674        command: LifecycleCommand,
675    ) -> Result<OwnerControlResponse, ControlPlaneClientError> {
676        let operation = command.operation();
677        self.send_unary_request(
678            std::time::Duration::from_secs(OWNER_CONTROL_UNARY_RESPONSE_TIMEOUT_SECS),
679            move |request_id, _, _| command.into_request(request_id),
680        )
681        .await
682        .map_err(|error| map_legacy_lifecycle_unsupported(operation, error))
683    }
684
685    pub async fn watch_config(
686        &self,
687        include_snapshot: bool,
688    ) -> Result<OwnerControlWatchStream, ControlPlaneClientError> {
689        let request_id = self.next_request_id();
690        let (mut send, recv) = self.open_authenticated_stream().await?;
691        let envelope = OwnerControlEnvelope {
692            r#gen: NODE_PROTOCOL_GENERATION,
693            handshake: None,
694            request: Some(OwnerControlRequest {
695                request_id,
696                get_config: None,
697                watch_config: Some(OwnerControlWatchConfigRequest {
698                    requester_node_id: self.endpoint.id().as_bytes().to_vec(),
699                    target_node_id: self.connection.remote_id().as_bytes().to_vec(),
700                    include_snapshot,
701                }),
702                apply_config: None,
703                refresh_inventory: None,
704                load_model: None,
705                unload_model: None,
706                ensure_model: None,
707                drain_model: None,
708            }),
709            response: None,
710            error: None,
711        };
712        write_owner_control_request(&mut send, &envelope).await?;
713        let mut stream = OwnerControlWatchStream {
714            send,
715            recv,
716            request_id,
717            pending: None,
718            closed: false,
719        };
720        let accepted = tokio::time::timeout(
721            std::time::Duration::from_secs(OWNER_CONTROL_WATCH_ACCEPT_TIMEOUT_SECS),
722            stream.next(),
723        )
724        .await
725        .map_err(|_| {
726            ControlPlaneClientError::Transport(format!(
727                "owner-control watch accept timed out after {OWNER_CONTROL_WATCH_ACCEPT_TIMEOUT_SECS}s"
728            ))
729        })??;
730        stream.pending = Some(accepted);
731        Ok(stream)
732    }
733
734    async fn send_unary_request<F>(
735        &self,
736        response_timeout: std::time::Duration,
737        build_request: F,
738    ) -> Result<OwnerControlResponse, ControlPlaneClientError>
739    where
740        F: FnOnce(u64, Vec<u8>, Vec<u8>) -> OwnerControlRequest,
741    {
742        let request_id = self.next_request_id();
743        let (mut send, mut recv) = self.open_authenticated_stream().await?;
744        let envelope = OwnerControlEnvelope {
745            r#gen: NODE_PROTOCOL_GENERATION,
746            handshake: None,
747            request: Some(build_request(
748                request_id,
749                self.endpoint.id().as_bytes().to_vec(),
750                self.connection.remote_id().as_bytes().to_vec(),
751            )),
752            response: None,
753            error: None,
754        };
755        write_owner_control_request(&mut send, &envelope).await?;
756        let envelope =
757            tokio::time::timeout(response_timeout, read_owner_control_message(&mut recv))
758                .await
759                .map_err(|_| {
760                    ControlPlaneClientError::Transport(format!(
761                        "owner-control unary response timed out after {}s",
762                        response_timeout.as_secs()
763                    ))
764                })??;
765        let _ = send.finish();
766        decode_response_envelope(request_id, envelope)
767    }
768
769    fn next_request_id(&self) -> u64 {
770        next_nonzero_request_id(&self.next_request_id)
771    }
772
773    async fn open_authenticated_stream(
774        &self,
775    ) -> Result<(iroh::endpoint::SendStream, iroh::endpoint::RecvStream), ControlPlaneClientError>
776    {
777        let (mut send, recv) = tokio::time::timeout(
778            std::time::Duration::from_secs(OWNER_CONTROL_OPEN_TIMEOUT_SECS),
779            self.connection.open_bi(),
780        )
781        .await
782        .map_err(|_| {
783            ControlPlaneClientError::Transport(format!(
784                "owner-control stream open timed out after {OWNER_CONTROL_OPEN_TIMEOUT_SECS}s"
785            ))
786        })?
787        .map_err(|error| ControlPlaneClientError::Transport(error.to_string()))?;
788        let handshake = OwnerControlEnvelope {
789            r#gen: NODE_PROTOCOL_GENERATION,
790            handshake: Some(OwnerControlHandshake {
791                ownership: Some(sign_node_ownership_proto(
792                    &self.owner_keypair,
793                    self.endpoint.id().as_bytes(),
794                )),
795            }),
796            request: None,
797            response: None,
798            error: None,
799        };
800        tokio::time::timeout(
801            std::time::Duration::from_secs(OWNER_CONTROL_HANDSHAKE_TIMEOUT_SECS),
802            write_len_prefixed(&mut send, &handshake.encode_to_vec()),
803        )
804        .await
805        .map_err(|_| {
806            ControlPlaneClientError::Transport(format!(
807                "owner-control handshake timed out after {OWNER_CONTROL_HANDSHAKE_TIMEOUT_SECS}s"
808            ))
809        })?
810        .map_err(|error| ControlPlaneClientError::Transport(error.to_string()))?;
811        Ok((send, recv))
812    }
813}
814
815fn validate_absent_model_target(
816    model_ref: String,
817    instance_id: Option<String>,
818) -> Result<(String, Option<String>), ControlPlaneClientError> {
819    let has_model_ref = !model_ref.trim().is_empty();
820    let has_instance_id = instance_id
821        .as_deref()
822        .is_some_and(|id| !id.trim().is_empty());
823
824    match (has_model_ref, has_instance_id) {
825        (true, false) => Ok((model_ref, None)),
826        (false, true) => Ok((String::new(), instance_id)),
827        _ => Err(ControlPlaneClientError::Protocol(
828            "unload_model and drain_model require exactly one model reference or instance id"
829                .to_string(),
830        )),
831    }
832}
833
834async fn close_failed_bootstrap_endpoint(endpoint: &Endpoint) {
835    let _ = tokio::time::timeout(
836        std::time::Duration::from_millis(FAILED_BOOTSTRAP_CLOSE_TIMEOUT_MILLIS),
837        endpoint.close(),
838    )
839    .await;
840}
841
842async fn write_owner_control_request(
843    send: &mut iroh::endpoint::SendStream,
844    envelope: &OwnerControlEnvelope,
845) -> Result<(), ControlPlaneClientError> {
846    tokio::time::timeout(
847        std::time::Duration::from_secs(OWNER_CONTROL_REQUEST_WRITE_TIMEOUT_SECS),
848        write_len_prefixed(send, &envelope.encode_to_vec()),
849    )
850    .await
851    .map_err(|_| {
852        ControlPlaneClientError::Transport(format!(
853            "owner-control request write timed out after {OWNER_CONTROL_REQUEST_WRITE_TIMEOUT_SECS}s"
854        ))
855    })?
856    .map_err(|error| ControlPlaneClientError::Transport(error.to_string()))
857}
858
859impl OwnerControlWatchStream {
860    pub fn request_id(&self) -> u64 {
861        self.request_id
862    }
863
864    pub async fn next(&mut self) -> Result<OwnerControlWatchEvent, ControlPlaneClientError> {
865        if let Some(event) = self.pending.take() {
866            return Ok(event);
867        }
868        let envelope = read_owner_control_message(&mut self.recv).await?;
869        let response = decode_response_envelope(self.request_id, envelope)?;
870        let watch = response.watch_config.ok_or_else(|| {
871            ControlPlaneClientError::Protocol(
872                "owner-control watch response missing watch_config payload".to_string(),
873            )
874        })?;
875        decode_watch_event(watch)
876    }
877
878    pub async fn close(&mut self) -> Result<(), ControlPlaneClientError> {
879        if self.closed {
880            return Ok(());
881        }
882        self.send
883            .finish()
884            .map_err(|error| ControlPlaneClientError::Transport(error.to_string()))?;
885        self.closed = true;
886        Ok(())
887    }
888
889    pub async fn cancel(&mut self) -> Result<(), ControlPlaneClientError> {
890        self.close().await
891    }
892}
893
894fn next_nonzero_request_id(counter: &AtomicU64) -> u64 {
895    loop {
896        let request_id = counter.fetch_add(1, Ordering::Relaxed);
897        if request_id != 0 {
898            return request_id;
899        }
900    }
901}
902
903impl Drop for OwnerControlWatchStream {
904    fn drop(&mut self) {
905        if !self.closed {
906            let _ = self.send.finish();
907            self.closed = true;
908        }
909    }
910}
911
912fn decode_watch_event(
913    watch: OwnerControlWatchConfigResponse,
914) -> Result<OwnerControlWatchEvent, ControlPlaneClientError> {
915    if let Some(accepted) = watch.accepted {
916        return Ok(OwnerControlWatchEvent::Accepted(accepted));
917    }
918    if let Some(snapshot) = watch.snapshot {
919        return Ok(OwnerControlWatchEvent::Snapshot(snapshot));
920    }
921    if let Some(update) = watch.update {
922        return Ok(OwnerControlWatchEvent::Update(update));
923    }
924    Err(ControlPlaneClientError::Protocol(
925        "owner-control watch response missing accepted/snapshot/update payload".to_string(),
926    ))
927}
928
929fn decode_response_envelope(
930    expected_request_id: u64,
931    envelope: OwnerControlEnvelope,
932) -> Result<OwnerControlResponse, ControlPlaneClientError> {
933    if let Some(error) = envelope.error {
934        return Err(ControlPlaneClientError::Remote(error.into()));
935    }
936    let response = envelope.response.ok_or_else(|| {
937        ControlPlaneClientError::Protocol(
938            "owner-control response envelope missing response payload".to_string(),
939        )
940    })?;
941    if response.request_id != expected_request_id {
942        return Err(ControlPlaneClientError::Protocol(format!(
943            "owner-control response request_id mismatch: expected {expected_request_id}, got {}",
944            response.request_id
945        )));
946    }
947    Ok(response)
948}
949
950async fn read_owner_control_message(
951    recv: &mut iroh::endpoint::RecvStream,
952) -> Result<OwnerControlEnvelope, ControlPlaneClientError> {
953    let bytes = crate::protocol::read_len_prefixed(recv)
954        .await
955        .map_err(|error| ControlPlaneClientError::Transport(error.to_string()))?;
956    decode_owner_control_envelope(&bytes)
957        .map_err(|error| ControlPlaneClientError::Protocol(error.to_string()))
958}
959
960async fn configured_endpoint_connect_error(
961    endpoint: &Endpoint,
962    control_addr: EndpointAddr,
963    options: &ControlPlaneBootstrapOptions,
964    error: iroh::endpoint::ConnectError,
965) -> ControlPlaneClientError {
966    let message = error.to_string();
967    let disposition = connect_error_probe_disposition(&error);
968    let legacy_mesh_reachable = match disposition {
969        ConnectProbeDisposition::ProbeLegacyMesh => legacy_mesh_probe(endpoint, control_addr).await,
970        ConnectProbeDisposition::SkipUnavailable => false,
971        ConnectProbeDisposition::Unsupported => true,
972    };
973    let (code, rendered) = if legacy_mesh_reachable {
974        control_unsupported_message(&message)
975    } else {
976        control_unavailable_message(&message)
977    };
978    ControlPlaneClientError::Negotiation(options.configured_endpoint_failure(code, rendered))
979}
980
981#[derive(Clone, Copy, Debug, Eq, PartialEq)]
982enum ConnectProbeDisposition {
983    SkipUnavailable,
984    ProbeLegacyMesh,
985    Unsupported,
986}
987
988fn connect_error_probe_disposition(
989    error: &iroh::endpoint::ConnectError,
990) -> ConnectProbeDisposition {
991    match error {
992        iroh::endpoint::ConnectError::Connect { source, .. } => match source {
993            iroh::endpoint::ConnectWithOptsError::SelfConnect { .. }
994            | iroh::endpoint::ConnectWithOptsError::NoAddress { .. }
995            | iroh::endpoint::ConnectWithOptsError::Noq { .. }
996            | iroh::endpoint::ConnectWithOptsError::InternalConsistencyError { .. }
997            | iroh::endpoint::ConnectWithOptsError::LocallyRejected { .. }
998            | iroh::endpoint::ConnectWithOptsError::EndpointClosed { .. } => {
999                ConnectProbeDisposition::SkipUnavailable
1000            }
1001            _ => fallback_probe_disposition(&source.to_string()),
1002        },
1003        iroh::endpoint::ConnectError::Connecting { source, .. } => match source {
1004            iroh::endpoint::ConnectingError::ConnectionError { source, .. } => {
1005                connection_error_probe_disposition(&source.to_string())
1006            }
1007            iroh::endpoint::ConnectingError::HandshakeFailure { source, .. } => match source {
1008                iroh::endpoint::AuthenticationError::NoAlpn { .. } => {
1009                    ConnectProbeDisposition::Unsupported
1010                }
1011                iroh::endpoint::AuthenticationError::RemoteId { .. } => {
1012                    ConnectProbeDisposition::SkipUnavailable
1013                }
1014                _ => fallback_probe_disposition(&source.to_string()),
1015            },
1016            iroh::endpoint::ConnectingError::InternalConsistencyError { .. }
1017            | iroh::endpoint::ConnectingError::LocallyRejected { .. } => {
1018                ConnectProbeDisposition::SkipUnavailable
1019            }
1020            _ => fallback_probe_disposition(&source.to_string()),
1021        },
1022        iroh::endpoint::ConnectError::Connection { source, .. } => {
1023            connection_error_probe_disposition(&source.to_string())
1024        }
1025        _ => fallback_probe_disposition(&error.to_string()),
1026    }
1027}
1028
1029fn connection_error_probe_disposition(message: &str) -> ConnectProbeDisposition {
1030    if is_alpn_mismatch_message(message) {
1031        ConnectProbeDisposition::Unsupported
1032    } else {
1033        ConnectProbeDisposition::SkipUnavailable
1034    }
1035}
1036
1037fn fallback_probe_disposition(message: &str) -> ConnectProbeDisposition {
1038    if is_alpn_mismatch_message(message) {
1039        ConnectProbeDisposition::ProbeLegacyMesh
1040    } else {
1041        ConnectProbeDisposition::SkipUnavailable
1042    }
1043}
1044
1045fn control_unsupported_message(message: &str) -> (OwnerControlErrorCode, String) {
1046    (
1047        OwnerControlErrorCode::ControlUnsupported,
1048        format!("remote endpoint did not negotiate mesh-llm-control/1: {message}"),
1049    )
1050}
1051
1052fn control_unavailable_message(message: &str) -> (OwnerControlErrorCode, String) {
1053    (
1054        OwnerControlErrorCode::ControlUnavailable,
1055        format!("remote owner-control endpoint is unavailable or unreachable: {message}"),
1056    )
1057}
1058
1059async fn legacy_mesh_probe(_endpoint: &Endpoint, control_addr: EndpointAddr) -> bool {
1060    let Ok(probe_endpoint) = Endpoint::builder(iroh::endpoint::presets::Minimal)
1061        .secret_key(iroh::SecretKey::generate())
1062        .alpns(vec![ALPN_V1.to_vec()])
1063        .relay_mode(relay_mode_from_endpoint_addr(&control_addr))
1064        .bind_addr(owner_control_client_bind_addr())
1065    else {
1066        return false;
1067    };
1068    let Ok(probe_endpoint) = probe_endpoint.bind().await else {
1069        return false;
1070    };
1071    if control_addr.relay_urls().next().is_some() {
1072        let _ =
1073            tokio::time::timeout(std::time::Duration::from_secs(3), probe_endpoint.online()).await;
1074    }
1075    let reachable = match tokio::time::timeout(
1076        std::time::Duration::from_secs(3),
1077        probe_endpoint.connect(control_addr, ALPN_V1),
1078    )
1079    .await
1080    {
1081        Ok(Ok(connection)) => {
1082            connection.close(0u32.into(), b"owner-control-legacy-probe-complete");
1083            true
1084        }
1085        _ => false,
1086    };
1087    probe_endpoint.close().await;
1088    reachable
1089}
1090
1091fn relay_mode_from_endpoint_addr(addr: &EndpointAddr) -> iroh::endpoint::RelayMode {
1092    match relay_map_from_endpoint_addr(addr) {
1093        Some(relay_map) => iroh::endpoint::RelayMode::Custom(relay_map),
1094        None => iroh::endpoint::RelayMode::Disabled,
1095    }
1096}
1097
1098fn is_alpn_mismatch_message(message: &str) -> bool {
1099    let lowered = message.to_ascii_lowercase();
1100    lowered.contains("alpn mismatch")
1101        || lowered.contains("no application protocol")
1102        || lowered.contains("application protocol selected")
1103}
1104
1105fn decode_endpoint_addr_token(invite_token: &str) -> anyhow::Result<EndpointAddr> {
1106    let json = base64::engine::general_purpose::URL_SAFE_NO_PAD
1107        .decode(invite_token)
1108        .context("invalid endpoint encoding")?;
1109    serde_json::from_slice(&json).context("invalid endpoint JSON")
1110}
1111
1112fn relay_map_from_endpoint_addr(addr: &EndpointAddr) -> Option<iroh::RelayMap> {
1113    let configs: Vec<_> = addr
1114        .relay_urls()
1115        .cloned()
1116        // Preserve iroh's default QUIC Address Discovery (QAD). `new(url, None)`
1117        // disables it, preventing reflexive candidate discovery and direct-path
1118        // upgrades across NAT (see issue #1065). `RelayUrl::into()` keeps QAD on.
1119        .map(|url| -> iroh::RelayConfig { url.into() })
1120        .collect();
1121    if configs.is_empty() {
1122        None
1123    } else {
1124        Some(iroh::RelayMap::from_iter(configs))
1125    }
1126}
1127
1128fn sign_node_ownership_proto(
1129    owner: &OwnerKeypair,
1130    node_endpoint_id: &[u8; 32],
1131) -> SignedNodeOwnership {
1132    let issued_at_unix_ms = current_time_unix_ms();
1133    let expires_at_unix_ms =
1134        issued_at_unix_ms + DEFAULT_NODE_CERT_LIFETIME_SECS.saturating_mul(1000);
1135    let cert_id = uuid::Uuid::new_v4().simple().to_string();
1136    let owner_sign_public_key = owner.verifying_key().as_bytes().to_vec();
1137    let owner_id = owner.owner_id();
1138    let signature_payload = canonical_claim_bytes(CanonicalClaim {
1139        version: NODE_OWNERSHIP_VERSION,
1140        cert_id: &cert_id,
1141        owner_id: &owner_id,
1142        owner_sign_public_key: &owner_sign_public_key,
1143        node_endpoint_id,
1144        issued_at_unix_ms,
1145        expires_at_unix_ms,
1146        node_label: None,
1147        hostname_hint: None,
1148    });
1149    SignedNodeOwnership {
1150        version: NODE_OWNERSHIP_VERSION,
1151        cert_id,
1152        owner_id,
1153        owner_sign_public_key,
1154        node_endpoint_id: node_endpoint_id.to_vec(),
1155        issued_at_unix_ms,
1156        expires_at_unix_ms,
1157        node_label: None,
1158        hostname_hint: None,
1159        signature: owner.sign_bytes(&signature_payload).to_vec(),
1160    }
1161}
1162
1163struct CanonicalClaim<'a> {
1164    version: u32,
1165    cert_id: &'a str,
1166    owner_id: &'a str,
1167    owner_sign_public_key: &'a [u8],
1168    node_endpoint_id: &'a [u8; 32],
1169    issued_at_unix_ms: u64,
1170    expires_at_unix_ms: u64,
1171    node_label: Option<&'a str>,
1172    hostname_hint: Option<&'a str>,
1173}
1174
1175fn canonical_claim_bytes(claim: CanonicalClaim<'_>) -> Vec<u8> {
1176    let mut buf = Vec::with_capacity(256);
1177    buf.extend_from_slice(SIGNING_DOMAIN_TAG);
1178    buf.extend_from_slice(&claim.version.to_le_bytes());
1179    write_string(&mut buf, claim.cert_id);
1180    write_string(&mut buf, claim.owner_id);
1181    buf.extend_from_slice(claim.owner_sign_public_key);
1182    buf.extend_from_slice(claim.node_endpoint_id);
1183    buf.extend_from_slice(&claim.issued_at_unix_ms.to_le_bytes());
1184    buf.extend_from_slice(&claim.expires_at_unix_ms.to_le_bytes());
1185    write_optional_string(&mut buf, claim.node_label);
1186    write_optional_string(&mut buf, claim.hostname_hint);
1187    buf
1188}
1189
1190fn write_string(buf: &mut Vec<u8>, value: &str) {
1191    buf.extend_from_slice(&(value.len() as u64).to_le_bytes());
1192    buf.extend_from_slice(value.as_bytes());
1193}
1194
1195fn write_optional_string(buf: &mut Vec<u8>, value: Option<&str>) {
1196    match value {
1197        Some(value) => {
1198            buf.push(1);
1199            write_string(buf, value);
1200        }
1201        None => buf.push(0),
1202    }
1203}
1204
1205fn current_time_unix_ms() -> u64 {
1206    std::time::SystemTime::now()
1207        .duration_since(std::time::UNIX_EPOCH)
1208        .unwrap_or_default()
1209        .as_millis() as u64
1210}
1211
1212#[cfg(test)]
1213mod tests {
1214    use super::*;
1215    use std::str::FromStr;
1216
1217    #[test]
1218    fn owner_control_client_binds_wildcard_for_direct_remote_endpoints() {
1219        let bind_addr = owner_control_client_bind_addr();
1220
1221        assert_eq!(bind_addr.port(), 0);
1222        assert!(
1223            bind_addr.ip().is_unspecified(),
1224            "owner-control clients must not be loopback-bound when dialing explicit remote endpoints"
1225        );
1226    }
1227
1228    #[test]
1229    fn relay_mode_uses_custom_relays_from_endpoint_addr() {
1230        let addr = EndpointAddr::new(iroh::SecretKey::generate().public()).with_relay_url(
1231            iroh::RelayUrl::from_str("https://relay.example.com").expect("relay URL parses"),
1232        );
1233
1234        assert!(matches!(
1235            relay_mode_from_endpoint_addr(&addr),
1236            iroh::endpoint::RelayMode::Custom(_)
1237        ));
1238    }
1239
1240    #[test]
1241    fn endpoint_addr_relays_preserve_default_qad() {
1242        let addr = EndpointAddr::new(iroh::SecretKey::generate().public())
1243            .with_relay_url(
1244                iroh::RelayUrl::from_str("https://relay-a.example.com").expect("relay URL parses"),
1245            )
1246            .with_relay_url(
1247                iroh::RelayUrl::from_str("https://relay-b.example.com").expect("relay URL parses"),
1248            );
1249
1250        let map = relay_map_from_endpoint_addr(&addr).expect("relay map should be enabled");
1251        let configs = map.relays::<Vec<_>>();
1252
1253        assert_eq!(configs.len(), 2);
1254        assert!(
1255            configs
1256                .iter()
1257                .all(|config| { config.quic.as_ref().is_some_and(|quic| quic.port == 7842) })
1258        );
1259    }
1260
1261    #[test]
1262    fn relay_mode_is_disabled_without_endpoint_relays() {
1263        let addr = EndpointAddr::new(iroh::SecretKey::generate().public());
1264
1265        assert!(matches!(
1266            relay_mode_from_endpoint_addr(&addr),
1267            iroh::endpoint::RelayMode::Disabled
1268        ));
1269    }
1270
1271    #[test]
1272    fn request_id_generator_skips_zero_after_wraparound() {
1273        let counter = AtomicU64::new(u64::MAX);
1274
1275        assert_eq!(next_nonzero_request_id(&counter), u64::MAX);
1276        assert_eq!(next_nonzero_request_id(&counter), 1);
1277    }
1278
1279    #[test]
1280    fn inventory_response_timeout_exceeds_server_scan_deadline() {
1281        const {
1282            assert!(
1283                OWNER_CONTROL_INVENTORY_RESPONSE_TIMEOUT_SECS
1284                    > OWNER_CONTROL_SERVER_SCAN_DEADLINE_SECS_FOR_CLIENT_MARGIN
1285            );
1286        }
1287        assert_eq!(
1288            OWNER_CONTROL_INVENTORY_RESPONSE_TIMEOUT_SECS
1289                - OWNER_CONTROL_SERVER_SCAN_DEADLINE_SECS_FOR_CLIENT_MARGIN,
1290            5
1291        );
1292    }
1293
1294    #[test]
1295    fn unary_response_timeout_exceeds_server_command_deadline() {
1296        const {
1297            assert!(
1298                OWNER_CONTROL_UNARY_RESPONSE_TIMEOUT_SECS
1299                    > OWNER_CONTROL_SERVER_UNARY_DEADLINE_SECS_FOR_CLIENT_MARGIN
1300            );
1301        }
1302        assert_eq!(
1303            OWNER_CONTROL_UNARY_RESPONSE_TIMEOUT_SECS
1304                - OWNER_CONTROL_SERVER_UNARY_DEADLINE_SECS_FOR_CLIENT_MARGIN,
1305            5
1306        );
1307    }
1308
1309    #[test]
1310    fn lifecycle_acceptance_requires_intent_id_and_exact_state() {
1311        let target = crate::proto::node::OwnerControlModelRef {
1312            canonical_model_ref: "model/test".to_string(),
1313            instance_id: None,
1314        };
1315        assert!(
1316            validate_lifecycle_acceptance(
1317                "load_model",
1318                "owner-1",
1319                "present",
1320                "present",
1321                Some(&target),
1322                "model/test",
1323                None,
1324            )
1325            .is_ok()
1326        );
1327        assert!(
1328            validate_lifecycle_acceptance(
1329                "load_model",
1330                "",
1331                "present",
1332                "present",
1333                Some(&target),
1334                "model/test",
1335                None,
1336            )
1337            .is_err()
1338        );
1339        assert!(
1340            validate_lifecycle_acceptance(
1341                "load_model",
1342                "owner-1",
1343                "absent",
1344                "present",
1345                Some(&target),
1346                "model/test",
1347                None,
1348            )
1349            .is_err()
1350        );
1351    }
1352
1353    #[test]
1354    fn absent_model_target_requires_exactly_one_reference() {
1355        assert_eq!(
1356            validate_absent_model_target("model/test".to_string(), None)
1357                .expect("model-only target"),
1358            ("model/test".to_string(), None)
1359        );
1360        assert_eq!(
1361            validate_absent_model_target(String::new(), Some("runtime-2".to_string()))
1362                .expect("instance-only target"),
1363            (String::new(), Some("runtime-2".to_string()))
1364        );
1365        assert!(matches!(
1366            validate_absent_model_target("model/test".to_string(), Some("runtime-2".to_string())),
1367            Err(ControlPlaneClientError::Protocol(_))
1368        ));
1369        assert!(matches!(
1370            validate_absent_model_target(String::new(), None),
1371            Err(ControlPlaneClientError::Protocol(_))
1372        ));
1373    }
1374
1375    #[test]
1376    fn legacy_unknown_lifecycle_command_maps_to_control_unsupported() {
1377        for code in [
1378            OwnerControlErrorCode::BadRequest,
1379            OwnerControlErrorCode::UnknownCommand,
1380        ] {
1381            let legacy = ControlPlaneClientError::Remote(OwnerControlRemoteError {
1382                code,
1383                message: "owner control request requires exactly one command variant".to_string(),
1384                request_id: Some(7),
1385                current_revision: None,
1386            });
1387            let mapped = map_legacy_lifecycle_unsupported("load_model", legacy);
1388            let ControlPlaneClientError::Remote(mapped) = mapped else {
1389                panic!("legacy response should remain a structured remote error");
1390            };
1391            assert_eq!(mapped.code, OwnerControlErrorCode::ControlUnsupported);
1392            assert_eq!(
1393                mapped.message,
1394                "remote owner-control endpoint does not support load_model"
1395            );
1396            assert_eq!(mapped.request_id, Some(7));
1397        }
1398
1399        let unrelated = ControlPlaneClientError::Remote(OwnerControlRemoteError {
1400            code: OwnerControlErrorCode::BadRequest,
1401            message: "invalid model target".to_string(),
1402            request_id: Some(8),
1403            current_revision: None,
1404        });
1405        let ControlPlaneClientError::Remote(unrelated) =
1406            map_legacy_lifecycle_unsupported("load_model", unrelated)
1407        else {
1408            panic!("unrelated response should remain a structured remote error");
1409        };
1410        assert_eq!(unrelated.code, OwnerControlErrorCode::BadRequest);
1411        assert_eq!(unrelated.message, "invalid model target");
1412    }
1413
1414    #[test]
1415    fn connect_error_fallback_is_narrow_to_alpn_negotiation() {
1416        assert!(is_alpn_mismatch_message("no application protocol selected"));
1417        assert!(is_alpn_mismatch_message("ALPN mismatch"));
1418        assert!(!is_alpn_mismatch_message("connection refused"));
1419        assert!(!is_alpn_mismatch_message("endpoint is closed"));
1420        assert!(!is_alpn_mismatch_message(
1421            "ALPN configuration is unavailable locally"
1422        ));
1423        assert_eq!(
1424            fallback_probe_disposition("connection refused"),
1425            ConnectProbeDisposition::SkipUnavailable
1426        );
1427        assert_eq!(
1428            fallback_probe_disposition("no application protocol selected"),
1429            ConnectProbeDisposition::ProbeLegacyMesh
1430        );
1431    }
1432}