Skip to main content

link_assistant_router/
managed_server.rs

1//! Server selection, managed Docker lifecycle, and per-run token handling.
2
3use std::fs::{self};
4use std::io::{Read as _, Write as _};
5use std::net::{TcpListener, TcpStream};
6use std::path::PathBuf;
7use std::process::{Child, Command, ExitCode, Stdio};
8use std::thread;
9use std::time::{Duration, Instant};
10
11use base64::Engine as _;
12use serde::{Deserialize, Serialize};
13use serde_json::Value;
14
15use crate::clients::{ClientKind, RouterModel};
16
17mod bootstrap;
18mod catalog;
19mod diagnostics;
20mod discovery;
21mod docker;
22mod http;
23mod origin;
24mod process;
25mod selection;
26
27use diagnostics::compact;
28use discovery::discover_local_router;
29pub use discovery::{discovered_local_router, effective_source};
30use docker::{
31    check_docker_output, docker_checked, docker_container_state, docker_subscription_status,
32    ensure_docker,
33};
34pub use http::revoke;
35use http::{client as http_client, revoke_with_client, verify_health, verify_health_with_client};
36pub use origin::canonical_server_origin;
37use origin::{normalize_server, same_origin};
38use process::process_alive;
39pub use selection::{
40    clear_persisted, configured_source, load_persisted, save_persisted, save_persisted_with_trust,
41    selected_server,
42};
43
44/// The default port a router binds when nothing else is specified.
45const DEFAULT_LOCAL_PORT: u16 = 8080;
46
47use catalog::fetch_models;
48
49const CONFIG_DIRECTORY: &str = "link-assistant-router";
50const SERVER_CONFIG: &str = "server.json";
51const MANAGED_STATE: &str = "managed-server.json";
52const MANAGED_LOCK: &str = "managed-server.lock";
53const CONTAINER: &str = "link-assistant-router-managed";
54const VOLUME: &str = "link-assistant-router-managed-data";
55const IMAGE: &str = "ghcr.io/link-assistant/router:latest";
56const MANAGED_LABEL: &str = "com.link-assistant.router.managed=1";
57
58type AnyError = Box<dyn std::error::Error + Send + Sync>;
59
60#[derive(Clone, Debug, Default, Deserialize, Serialize)]
61pub struct PersistedServer {
62    pub server: String,
63    #[serde(default, skip_serializing_if = "Option::is_none")]
64    pub management_server: Option<String>,
65    #[serde(default, skip_serializing_if = "Option::is_none")]
66    pub token: Option<String>,
67    #[serde(default, skip_serializing_if = "Option::is_none")]
68    pub run_max_requests: Option<u64>,
69    /// Router-owned certificate bundle associated with the inference origin.
70    #[serde(default, skip_serializing_if = "Option::is_none")]
71    pub ca_cert: Option<String>,
72    /// Router-owned certificate bundle associated with the management origin.
73    #[serde(default, skip_serializing_if = "Option::is_none")]
74    pub management_ca_cert: Option<String>,
75}
76
77#[derive(Clone, Debug, Deserialize, Serialize)]
78struct ManagedState {
79    port: u16,
80    #[serde(default = "managed_secret")]
81    token_secret: String,
82    #[serde(default)]
83    references: Vec<u32>,
84    #[serde(default)]
85    keep_running: bool,
86    #[serde(default)]
87    claimed: bool,
88}
89
90/// The effective server and the optional lease that controls a managed local one.
91pub struct ResolvedServer {
92    /// Public origin used for health, service catalogs, and generated clients.
93    pub base_url: String,
94    /// Private origin used only for management API calls.
95    pub management_url: String,
96    pub token: Option<String>,
97    pub source: &'static str,
98    pub run_max_requests: Option<u64>,
99    /// Router-owned CA bundle for requests and clients using `base_url`.
100    pub ca_cert: Option<PathBuf>,
101    /// Router-owned CA bundle for requests using `management_url`.
102    pub management_ca_cert: Option<PathBuf>,
103    _lease: Option<ManagedLease>,
104}
105
106impl ResolvedServer {
107    /// A server reference for an already-known origin.
108    ///
109    /// Holds no managed-container lease, so it neither starts nor keeps one
110    /// alive — for a router that is simply already running at `base_url`.
111    #[must_use]
112    pub fn at(base_url: impl Into<String>, token: Option<String>, source: &'static str) -> Self {
113        let base_url = base_url.into();
114        Self {
115            management_url: base_url.clone(),
116            base_url,
117            token,
118            source,
119            run_max_requests: None,
120            ca_cert: None,
121            management_ca_cert: None,
122            _lease: None,
123        }
124    }
125
126    /// A server whose management and inference listeners are intentionally
127    /// disjoint.
128    #[must_use]
129    pub fn at_origins(
130        base_url: impl Into<String>,
131        management_url: impl Into<String>,
132        token: Option<String>,
133        source: &'static str,
134    ) -> Self {
135        Self {
136            base_url: base_url.into(),
137            management_url: management_url.into(),
138            token,
139            source,
140            run_max_requests: None,
141            ca_cert: None,
142            management_ca_cert: None,
143            _lease: None,
144        }
145    }
146
147    /// HTTP client carrying the trust associated with the inference origin.
148    pub fn inference_client(&self) -> Result<reqwest::Client, AnyError> {
149        self.inference_client_with_timeout(Duration::from_secs(10))
150    }
151
152    pub(crate) fn inference_client_with_timeout(
153        &self,
154        timeout: Duration,
155    ) -> Result<reqwest::Client, AnyError> {
156        http_client(self.ca_cert.as_deref(), timeout)
157    }
158
159    /// HTTP client carrying the trust associated with the management origin.
160    pub fn management_client(&self) -> Result<reqwest::Client, AnyError> {
161        self.management_client_with_timeout(Duration::from_secs(10))
162    }
163
164    pub(crate) fn management_client_with_timeout(
165        &self,
166        timeout: Duration,
167    ) -> Result<reqwest::Client, AnyError> {
168        http_client(self.management_ca_cert.as_deref(), timeout)
169    }
170}
171
172/// An ordinary token suitable for a wrapped client.
173pub struct RunCredential {
174    pub token: String,
175    available_models: Vec<RouterModel>,
176    revocation: Option<Revocation>,
177    principal_id: String,
178}
179
180impl RunCredential {
181    pub(crate) fn models(&self) -> &[RouterModel] {
182        &self.available_models
183    }
184
185    pub(crate) fn principal_id(&self) -> &str {
186        &self.principal_id
187    }
188
189    /// The record id this credential was issued under, retained so a
190    /// persistent credential remains revocable later (issue #190).
191    #[must_use]
192    pub fn id(&self) -> Option<String> {
193        token_subject(&self.token).ok()
194    }
195
196    /// Whether this command minted the credential. Supplied tokens may be
197    /// shared with other machines and must not be revoked implicitly (#296).
198    #[must_use]
199    pub const fn was_minted(&self) -> bool {
200        self.revocation.is_some()
201    }
202}
203
204struct Revocation {
205    base_url: String,
206    admin_token: String,
207    id: String,
208    client: reqwest::Client,
209}
210
211struct ManagedLease {
212    pid: u32,
213    reaper: Child,
214}
215
216impl Drop for ManagedLease {
217    fn drop(&mut self) {
218        if let Err(error) = release_reference(self.pid) {
219            eprintln!("warning: could not release managed router reference: {error}");
220        }
221        // Closing the pipe tells the crash reaper that this owner has finished
222        // normally. Wait for its idempotent cleanup so detached subprocesses
223        // (and their coverage profiles) cannot outlive the wrapper.
224        drop(self.reaper.stdin.take());
225        let status = self.reaper.wait();
226        if !status.as_ref().is_ok_and(std::process::ExitStatus::success) {
227            eprintln!("warning: managed router crash reaper failed: {status:?}");
228        }
229    }
230}
231
232/// Resolve flags, environment, persisted selection, a router already running
233/// locally, then managed local Docker.
234///
235/// `force_managed` skips the discovery step, for workflows that want a
236/// disposable instance on purpose (issue #250).
237pub async fn resolve(
238    explicit_server: Option<&str>,
239    explicit_management_server: Option<&str>,
240    explicit_token: Option<String>,
241    run_max_requests: Option<u64>,
242    force_managed: bool,
243) -> Result<ResolvedServer, AnyError> {
244    // An explicit origin remains usable even if unrelated saved state is
245    // unreadable. Without an explicit target, persisted state is the target
246    // and must fail closed instead of silently starting or selecting another
247    // deployment.
248    let persisted = if explicit_server.is_some() {
249        load_persisted().ok().flatten()
250    } else {
251        load_persisted()?
252    };
253    let environment_server = std::env::var("LINK_ASSISTANT_ROUTER_URL")
254        .or_else(|_| std::env::var("ROUTER_URL"))
255        .ok();
256    let environment_token = std::env::var("LINK_ASSISTANT_ROUTER_TOKEN")
257        .or_else(|_| std::env::var("LINK_ASSISTANT_TOKEN"))
258        .ok();
259    let environment_management_server = std::env::var("LINK_ASSISTANT_ROUTER_MANAGEMENT_URL")
260        .or_else(|_| std::env::var("ROUTER_MANAGEMENT_URL"))
261        .ok();
262    let (base_url, management_url, source) = if let Some(server) = explicit_server {
263        let base_url = normalize_server(server)?;
264        let management_url = explicit_management_server
265            .map(normalize_server)
266            .transpose()?
267            .unwrap_or_else(|| base_url.clone());
268        (base_url, management_url, "flag")
269    } else if let Some(server) = environment_server {
270        let base_url = normalize_server(&server)?;
271        let management_url = if let Some(server) = explicit_management_server {
272            normalize_server(server)?
273        } else if let Some(server) = environment_management_server.as_deref() {
274            normalize_server(server)?
275        } else {
276            base_url.clone()
277        };
278        (base_url, management_url, "environment")
279    } else if let Some(config) = persisted.as_ref() {
280        let base_url = normalize_server(&config.server)?;
281        let management_url = explicit_management_server
282            .map(normalize_server)
283            .transpose()?
284            .or_else(|| config.management_server.clone())
285            .unwrap_or_else(|| base_url.clone());
286        (base_url, management_url, "persisted configuration")
287    } else if let Some(discovered) = discover_local_router(force_managed).await {
288        // Nothing was selected explicitly and a router is already listening
289        // here, so use it rather than starting a second one. Starting one was
290        // both the expensive branch — an image pull and a container start on a
291        // command the operator expects to be instant — and the surprising one:
292        // the new container has its own credential directory and token store,
293        // so a subscription authorized through it is invisible to the instance
294        // already running (issue #250). Every explicit mechanism above this
295        // point still wins, and `--managed` forces a fresh container.
296        let management = explicit_management_server
297            .map(normalize_server)
298            .transpose()?
299            .unwrap_or_else(|| discovered.clone());
300        (discovered, management, "already-running local server")
301    } else {
302        let (state, lease) = acquire_managed()?;
303        let base_url = format!("http://127.0.0.1:{}", state.port);
304        verify_health(&base_url).await?;
305        let token = match explicit_token.or(environment_token) {
306            Some(token) => Some(token),
307            None if !state.claimed => Some(bootstrap::read_token(CONTAINER)?),
308            None => None,
309        };
310        return Ok(ResolvedServer {
311            management_url: base_url.clone(),
312            base_url,
313            token,
314            source: "managed local container",
315            run_max_requests,
316            ca_cert: None,
317            management_ca_cert: None,
318            _lease: Some(lease),
319        });
320    };
321    // Matched by *origin*, not by how the origin was supplied. The condition
322    // used to be on the source, so writing down the address of the very router
323    // that was already selected threw away the token stored for it — and the
324    // advice in the resulting error was to run the command the user had
325    // already run. A router found by discovery got no credential at all, for
326    // the same reason, though it is the same listener at the same address
327    // (issue #311). Explicit `--token` and the environment still win.
328    let token = explicit_token.or(environment_token).or_else(|| {
329        persisted
330            .as_ref()
331            .filter(|config| same_origin(&config.server, &base_url))
332            .and_then(|config| config.token.clone())
333    });
334    let budget = run_max_requests.or_else(|| {
335        persisted
336            .as_ref()
337            .and_then(|config| config.run_max_requests)
338    });
339    let ca_cert = persisted
340        .as_ref()
341        .filter(|config| same_origin(&config.server, &base_url))
342        .map(|config| selection::certificate_path(config.ca_cert.as_deref()))
343        .transpose()?
344        .flatten();
345    let management_ca_cert = persisted
346        .as_ref()
347        .filter(|config| {
348            let selected_management = config
349                .management_server
350                .as_deref()
351                .unwrap_or(&config.server);
352            same_origin(selected_management, &management_url)
353        })
354        .map(|config| {
355            let name = config.management_ca_cert.as_deref().or_else(|| {
356                same_origin(&config.server, &management_url)
357                    .then_some(config.ca_cert.as_deref())
358                    .flatten()
359            });
360            selection::certificate_path(name)
361        })
362        .transpose()?
363        .flatten();
364    // A selected server that is not answering is an error rather than a
365    // silent fallback -- using a different router than the one the operator
366    // chose is its own surprise. But the message has to say which server,
367    // and what to do about it: the report that prompted this got docker's
368    // words about an internal container it had never heard of (issue #333).
369    let inference_client = http_client(ca_cert.as_deref(), Duration::from_secs(10))?;
370    verify_health_with_client(&inference_client, &base_url)
371        .await
372        .map_err(|error| -> AnyError {
373            match source {
374            "flag" => error,
375            _ => format!(
376                "{error}\nnote: {base_url} is the router selected by {source}.\nnote: {}",
377                concat!(
378                    "pass --local to use a router on this machine, --managed to start a ",
379                    "disposable one, or run `router server use <URL>` to select a different one."
380                )
381            )
382            .into(),
383        }
384        })?;
385    Ok(ResolvedServer {
386        base_url,
387        management_url,
388        token,
389        source,
390        run_max_requests: budget,
391        ca_cert,
392        management_ca_cert,
393        _lease: None,
394    })
395}
396
397/// Validate an ordinary token or exchange an admin credential for a run token.
398pub async fn prepare_run_credential(
399    server: &ResolvedServer,
400    client_kind: ClientKind,
401    label: &str,
402    ttl_hours: i64,
403    sliding: bool,
404) -> Result<RunCredential, AnyError> {
405    prepare_credential(server, client_kind, label, ttl_hours, sliding, true, true).await
406}
407
408/// Mint the client-bound credential used by a permanent repair.
409///
410/// Repair is a trust takeover, not a one-shot launch. It must never persist a
411/// supplied ordinary token merely because the selected listener cannot mint a
412/// replacement: only a candidate minted for this exact client is eligible.
413pub async fn prepare_repair_credential(
414    server: &ResolvedServer,
415    client_kind: ClientKind,
416    label: &str,
417    ttl_hours: i64,
418) -> Result<RunCredential, AnyError> {
419    prepare_credential(server, client_kind, label, ttl_hours, false, false, false).await
420}
421
422/// Mint or reuse a credential that remains after this command exits.
423pub async fn prepare_persistent_credential(
424    server: &ResolvedServer,
425    client_kind: ClientKind,
426    label: &str,
427    ttl_hours: i64,
428) -> Result<RunCredential, AnyError> {
429    prepare_credential(server, client_kind, label, ttl_hours, false, true, false).await
430}
431
432async fn prepare_credential(
433    server: &ResolvedServer,
434    client_kind: ClientKind,
435    label: &str,
436    ttl_hours: i64,
437    sliding: bool,
438    allow_supplied: bool,
439    ephemeral: bool,
440) -> Result<RunCredential, AnyError> {
441    let token = server.token.as_deref().ok_or_else(|| {
442        if server.source == "managed local container" {
443            format!(
444                "{} is claimed and no token is available; pass --token, use --token-stdin, set LINK_ASSISTANT_ROUTER_TOKEN, or issue an ordinary token with `docker exec {CONTAINER} link-assistant-router tokens issue --ttl-hours 24 --label with-router`",
445                server.base_url
446            )
447        } else {
448            // Naming which origin has a stored token, rather than
449            // recommending the command the user already ran: that difference
450            // is the whole content of the error (issue #311).
451            let stored = load_persisted()
452                .ok()
453                .flatten()
454                .filter(|persisted| persisted.token.is_some())
455                .map(|persisted| persisted.server);
456            let held = stored.map_or_else(String::new, |origin| {
457                format!(" A token is stored for {origin}, which is a different origin.")
458            });
459            format!(
460                "{} selected from {}, but no token is available.{held} Pass --token, use --token-stdin, set LINK_ASSISTANT_ROUTER_TOKEN, or run `link-assistant-router server use {} --token-stdin`",
461                server.base_url, server.source, server.base_url
462            )
463        }
464    })?;
465    let management_client = server.management_client()?;
466    let inference_client = server.inference_client()?;
467    let list_url = crate::route_contract::management_endpoint(
468        &server.management_url,
469        crate::route_contract::RouteId::Tokens,
470    );
471    let list = management_client
472        .get(&list_url)
473        .bearer_auth(token)
474        .send()
475        .await;
476    match list {
477        Ok(response) if response.status().is_success() => {
478            let issue_url = crate::route_contract::management_endpoint(
479                &server.management_url,
480                crate::route_contract::RouteId::ClientTokens,
481            );
482            let response = management_client
483                .post(&issue_url)
484                .bearer_auth(token)
485                .json(&serde_json::json!({
486                    "ttl_hours": ttl_hours,
487                    "label": label,
488                    "client_kind": client_kind.canonical_name(),
489                    "max_requests": server.run_max_requests,
490                    // The run is revoked when the client exits, so the clock
491                    // is a backstop for a client that never got to exit --
492                    // not a limit on how long a live session may run
493                    // (issue #354).
494                    "sliding_expiry": sliding,
495                    "ephemeral": ephemeral,
496                }))
497                .send()
498                .await
499                .map_err(|error| {
500                    format!("could not mint a per-run token at {issue_url}: {error}")
501                })?;
502            let status = response.status();
503            let body = response.text().await.unwrap_or_default();
504            if !status.is_success() {
505                return Err(format!(
506                    "per-run token minting failed at {issue_url} ({status}): {}",
507                    compact(&body)
508                )
509                .into());
510            }
511            let value: Value = serde_json::from_str(&body)
512                .map_err(|error| format!("token endpoint returned invalid JSON: {error}"))?;
513            let run_token = value
514                .get("token")
515                .and_then(Value::as_str)
516                .ok_or("token endpoint response did not contain a token")?
517                .to_string();
518            let id = token_subject(&run_token)?;
519            let principal_id = exact_token_binding(&run_token, client_kind)?;
520            let mut credential = RunCredential {
521                token: run_token,
522                available_models: Vec::new(),
523                revocation: Some(Revocation {
524                    base_url: server.management_url.clone(),
525                    admin_token: token.to_string(),
526                    id,
527                    client: management_client.clone(),
528                }),
529                principal_id,
530            };
531            match fetch_models(
532                &inference_client,
533                client_kind,
534                &server.base_url,
535                &credential.token,
536            )
537            .await
538            {
539                Ok(models) => {
540                    credential.available_models = models;
541                    Ok(credential)
542                }
543                Err(error) => match cleanup_run_credential(credential).await {
544                    Ok(()) => Err(error),
545                    Err(cleanup) => Err(format!(
546                        "{error}; the unused minted credential could not be revoked: {cleanup}"
547                    )
548                    .into()),
549                },
550            }
551        }
552        Ok(response) if response.status().as_u16() == 401 || response.status().as_u16() == 403 => {
553            if !allow_supplied {
554                return Err(format!(
555                    "client repair requires an administrator credential that can mint a token bound to `{}`; the selected credential is inference-only",
556                    client_kind.canonical_name()
557                )
558                .into());
559            }
560            let principal_id = exact_token_binding(token, client_kind)?;
561            let available_models =
562                fetch_models(&inference_client, client_kind, &server.base_url, token).await?;
563            Ok(RunCredential {
564                token: token.to_string(),
565                available_models,
566                revocation: None,
567                principal_id,
568            })
569        }
570        Ok(response) if response.status().as_u16() == 404 => {
571            if !allow_supplied {
572                return Err(format!(
573                    "client repair requires the administrator listener so Router can mint a token bound to `{}`; the selected listener exposes inference only",
574                    client_kind.canonical_name()
575                )
576                .into());
577            }
578            let principal_id = exact_token_binding(token, client_kind).map_err(|_| {
579                format!(
580                    "the selected listener exposes inference only, so its supplied token must carry the exact `{}` client binding and a subscriber principal; use the matching client token or select the administrator listener",
581                    client_kind.canonical_name()
582                )
583            })?;
584            let available_models =
585                fetch_models(&inference_client, client_kind, &server.base_url, token).await?;
586            Ok(RunCredential {
587                token: token.to_string(),
588                available_models,
589                revocation: None,
590                principal_id,
591            })
592        }
593        Ok(response) => Err(format!(
594            "could not determine token scope at {list_url}: server returned {}",
595            response.status()
596        )
597        .into()),
598        Err(error) => Err(format!("could not inspect token scope at {list_url}: {error}").into()),
599    }
600}
601
602/// Refuse to launch a client with a model the selected router cannot serve.
603pub fn ensure_model_available(
604    credential: &RunCredential,
605    client: ClientKind,
606    model: &str,
607) -> Result<(), AnyError> {
608    if crate::clients::model_is_authorized(client, &credential.available_models, model) {
609        return Ok(());
610    }
611    if credential.available_models.is_empty() {
612        return Err(
613            "the router has no available models; authorize a subscription with `link-assistant-router auth <claude|codex|gemini|qwen>` on the router host and retry"
614                .into(),
615        );
616    }
617    Err(format!(
618        "model `{model}` is not available from the selected router; available models: {}",
619        credential
620            .available_models
621            .iter()
622            .map(|item| item.id.as_str())
623            .collect::<Vec<_>>()
624            .join(", ")
625    )
626    .into())
627}
628
629/// Revoke an automatically minted token. Explicit ordinary tokens are untouched.
630pub async fn cleanup_run_credential(credential: RunCredential) -> Result<(), AnyError> {
631    let Some(revocation) = credential.revocation else {
632        return Ok(());
633    };
634    revoke_with_client(
635        &revocation.client,
636        &revocation.base_url,
637        &revocation.admin_token,
638        &revocation.id,
639    )
640    .await
641}
642
643#[path = "managed_server_token.rs"]
644mod token_helpers;
645use token_helpers::exact_token_binding;
646pub(crate) use token_helpers::{token_client_binding, token_subject};
647
648pub fn managed_status() -> Result<String, AnyError> {
649    let lock = lock_state()?;
650    let mut state = load_managed()?;
651    if let Some(state) = state.as_mut() {
652        prune_references(state);
653        save_managed(state)?;
654    }
655    drop(lock);
656    let lifecycle =
657        docker_container_state().unwrap_or_else(|error| format!("unavailable ({error})"));
658    Ok(match state {
659        Some(state) => {
660            let subscriptions = if lifecycle == "running" {
661                docker_subscription_status()
662            } else {
663                "not queried while stopped".to_string()
664            };
665            format!(
666                "{lifecycle}; administrator={}; container={CONTAINER}; volume={VOLUME}; url=http://127.0.0.1:{}; users={}; subscriptions={subscriptions}",
667                if state.claimed {
668                    "claimed"
669                } else {
670                    "unclaimed"
671                },
672                state.port,
673                state.references.len()
674            )
675        }
676        None => format!("absent; container={CONTAINER}; volume={VOLUME}"),
677    })
678}
679
680/// Explain a managed-container disappearance after a client-side failure.
681#[must_use]
682pub fn managed_failure_hint() -> Option<String> {
683    match docker_container_state() {
684        Ok(state) if state != "running" => Some(format!(
685            "managed router container is {state}; run `link-assistant-router server start` and retry"
686        )),
687        Err(error) => Some(format!("managed router state is unavailable: {error}")),
688        _ => None,
689    }
690}
691
692pub fn start_managed() -> Result<String, AnyError> {
693    let lock = lock_state()?;
694    let mut state = load_or_create_managed()?;
695    prune_references(&mut state);
696    ensure_container_running(&state)?;
697    state.keep_running = true;
698    save_managed(&state)?;
699    drop(lock);
700    Ok(format!("http://127.0.0.1:{}", state.port))
701}
702
703/// Explicitly hand the managed bootstrap administrator to its owner.
704///
705/// Before this transition the wrapper may use the credential only to mint a
706/// short-lived ordinary token. Claiming is deliberately one-way: subsequent
707/// unattended runs must receive a credential from the user.
708pub fn claim_managed() -> Result<String, AnyError> {
709    let lock = lock_state()?;
710    let Some(mut state) = load_managed()? else {
711        return Err(
712            "managed router is absent; run `link-assistant-router server start` first".into(),
713        );
714    };
715    if state.claimed {
716        return Err(
717            "managed router administrator is already claimed; the bootstrap credential is not printed twice"
718                .into(),
719        );
720    }
721    let token = bootstrap::read_token(CONTAINER)?;
722    state.claimed = true;
723    save_managed(&state)?;
724    drop(lock);
725    Ok(token)
726}
727
728pub fn stop_managed() -> Result<(), AnyError> {
729    let lock = lock_state()?;
730    let Some(mut state) = load_managed()? else {
731        return Err("managed router is absent; run `link-assistant-router server start`".into());
732    };
733    ensure_docker()?;
734    match docker_container_state()?.as_str() {
735        "running" => docker_checked(["stop", CONTAINER])?,
736        "stopped" => {}
737        "absent" => {
738            return Err(
739                "managed router container is absent; run `link-assistant-router server start` to recreate it"
740                    .into(),
741            );
742        }
743        other => return Err(format!("unexpected managed container state: {other}").into()),
744    }
745    state.keep_running = false;
746    state.references.clear();
747    save_managed(&state)?;
748    drop(lock);
749    Ok(())
750}
751
752pub fn remove_managed(yes: bool) -> Result<(), AnyError> {
753    if !yes {
754        return Err(format!(
755            "refusing to remove {CONTAINER} and volume {VOLUME}: issued tokens, request logs, and any authorized Claude/ChatGPT/Gemini/Qwen subscriptions will be permanently lost; rerun with --yes"
756        )
757        .into());
758    }
759    let lock = lock_state()?;
760    if load_managed()?.is_none() {
761        return Err("managed router is absent; no owned state was removed".into());
762    }
763    ensure_docker()?;
764    match docker_container_state()?.as_str() {
765        "running" | "stopped" => docker_checked(["rm", "-f", CONTAINER])?,
766        "absent" => {}
767        other => return Err(format!("unexpected managed container state: {other}").into()),
768    }
769    let output = Command::new("docker")
770        .args(["volume", "rm", VOLUME])
771        .output()?;
772    if !output.status.success()
773        && !String::from_utf8_lossy(&output.stderr).contains("No such volume")
774    {
775        return Err(format!(
776            "could not remove managed router volume: {}",
777            compact(&String::from_utf8_lossy(&output.stderr))
778        )
779        .into());
780    }
781    let path = state_directory()?.join(MANAGED_STATE);
782    let _ = fs::remove_file(path);
783    drop(lock);
784    Ok(())
785}
786
787#[must_use]
788pub fn reap(pid: u32) -> ExitCode {
789    // The owner holds this pipe open for the lifetime of its managed lease.
790    // EOF is delivered both on orderly teardown and if the wrapper is killed,
791    // without PID polling races or waiting for a reused PID to disappear.
792    let pipe_result = std::io::copy(&mut std::io::stdin().lock(), &mut std::io::sink());
793    if let Err(error) = pipe_result {
794        eprintln!("warning: managed router crash-reaper pipe failed: {error}");
795    }
796    match release_reference(pid) {
797        Ok(()) => ExitCode::SUCCESS,
798        Err(error) => {
799            eprintln!("error: could not reap managed router reference {pid}: {error}");
800            ExitCode::from(1)
801        }
802    }
803}
804
805fn acquire_managed() -> Result<(ManagedState, ManagedLease), AnyError> {
806    let lock = lock_state()?;
807    let mut state = load_or_create_managed()?;
808    prune_references(&mut state);
809    ensure_container_running(&state)?;
810    let pid = std::process::id();
811    if !state.references.contains(&pid) {
812        state.references.push(pid);
813    }
814    save_managed(&state)?;
815    let reaper = spawn_reaper(pid)?;
816    drop(lock);
817    Ok((state, ManagedLease { pid, reaper }))
818}
819
820fn release_reference(pid: u32) -> Result<(), AnyError> {
821    let lock = lock_state()?;
822    let Some(mut state) = load_managed()? else {
823        return Ok(());
824    };
825    state.references.retain(|reference| *reference != pid);
826    prune_references(&mut state);
827    if state.references.is_empty() && !state.keep_running {
828        match docker_container_state()?.as_str() {
829            "running" => docker_checked(["stop", CONTAINER])?,
830            "stopped" | "absent" => {}
831            other => return Err(format!("unexpected managed container state: {other}").into()),
832        }
833    }
834    save_managed(&state)?;
835    drop(lock);
836    Ok(())
837}
838
839fn spawn_reaper(pid: u32) -> Result<Child, AnyError> {
840    let current = std::env::current_exe()?;
841    let executable = if current
842        .file_stem()
843        .and_then(|value| value.to_str())
844        .is_some_and(|name| name == "link-assistant-router")
845    {
846        current
847    } else {
848        current.with_file_name(format!(
849            "link-assistant-router{}",
850            std::env::consts::EXE_SUFFIX
851        ))
852    };
853    Command::new(executable)
854        .args(["server", "reap", &pid.to_string()])
855        .stdin(Stdio::piped())
856        .stdout(Stdio::null())
857        .stderr(Stdio::null())
858        .spawn()
859        .map_err(|error| format!("could not start managed-server crash reaper: {error}").into())
860}
861
862fn ensure_container_running(state: &ManagedState) -> Result<(), AnyError> {
863    ensure_docker()?;
864    match docker_container_state()?.as_str() {
865        "running" => return wait_for_health(state.port),
866        "stopped" => {
867            docker_checked(["start", CONTAINER])?;
868        }
869        "absent" => {
870            let port_mapping = format!("127.0.0.1:{}:8080", state.port);
871            let volume = format!("{VOLUME}:/data");
872            let output = Command::new("docker")
873                .env("TOKEN_SECRET", &state.token_secret)
874                .args([
875                    "run",
876                    "-d",
877                    "--name",
878                    CONTAINER,
879                    "--label",
880                    MANAGED_LABEL,
881                    "-p",
882                    &port_mapping,
883                    "-e",
884                    "TOKEN_SECRET",
885                    "-e",
886                    "DATA_DIR=/data/router",
887                    "-e",
888                    "STORAGE_POLICY=text",
889                    "-e",
890                    "CLAUDE_CODE_HOME=/data/claude",
891                    "-v",
892                    &volume,
893                    IMAGE,
894                    "serve",
895                ])
896                .output()?;
897            check_docker_output(&output)?;
898        }
899        other => {
900            return Err(format!("unexpected managed container state: {other}").into());
901        }
902    }
903    wait_for_health(state.port)
904}
905
906fn wait_for_health(port: u16) -> Result<(), AnyError> {
907    let deadline = Instant::now() + Duration::from_secs(20);
908    while Instant::now() < deadline {
909        if tcp_health(port) {
910            return Ok(());
911        }
912        thread::sleep(Duration::from_millis(100));
913    }
914    let logs = Command::new("docker")
915        .args(["logs", "--tail", "20", CONTAINER])
916        .output()
917        .map(|output| String::from_utf8_lossy(&output.stderr).into_owned())
918        .unwrap_or_default();
919    Err(format!(
920        "managed router container did not become healthy on 127.0.0.1:{port}: {}",
921        compact(&logs)
922    )
923    .into())
924}
925
926fn tcp_health(port: u16) -> bool {
927    let Ok(mut stream) = TcpStream::connect(("127.0.0.1", port)) else {
928        return false;
929    };
930    let _ = stream.set_read_timeout(Some(Duration::from_secs(1)));
931    let _ = stream.set_write_timeout(Some(Duration::from_secs(1)));
932    if stream
933        .write_all(b"GET /api/health HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n")
934        .is_err()
935    {
936        return false;
937    }
938    let mut response = String::new();
939    stream.read_to_string(&mut response).is_ok()
940        && response.starts_with("HTTP/1.1 200")
941        && response.ends_with("ok")
942}
943
944fn new_managed_state() -> ManagedState {
945    ManagedState {
946        port: choose_port(),
947        token_secret: managed_secret(),
948        references: Vec::new(),
949        keep_running: false,
950        claimed: false,
951    }
952}
953
954fn load_or_create_managed() -> Result<ManagedState, AnyError> {
955    if let Some(state) = load_managed()? {
956        return Ok(state);
957    }
958    let state = new_managed_state();
959    // Persist before Docker creation. If creation succeeds but readiness
960    // fails, the next attempt must reuse the same port and credentials rather
961    // than orphaning an unrecoverable container.
962    save_managed(&state)?;
963    Ok(state)
964}
965
966fn managed_secret() -> String {
967    format!(
968        "{}_{}",
969        uuid::Uuid::new_v4().simple(),
970        uuid::Uuid::new_v4().simple()
971    )
972}
973
974fn choose_port() -> u16 {
975    TcpListener::bind(("127.0.0.1", DEFAULT_LOCAL_PORT)).map_or_else(
976        |_| {
977            TcpListener::bind(("127.0.0.1", 0))
978                .and_then(|listener| listener.local_addr())
979                .map_or(18080, |address| address.port())
980        },
981        |_| DEFAULT_LOCAL_PORT,
982    )
983}
984
985fn prune_references(state: &mut ManagedState) {
986    state.references.retain(|pid| process_alive(*pid));
987}
988
989#[path = "managed_server_state.rs"]
990mod state;
991/// The test-only state-root claim, so any test in the crate can isolate
992/// itself rather than operating on whoever ran it (issue #343).
993#[cfg(test)]
994pub(crate) use state::claim_state_root;
995
996use state::{load_managed, lock_state, save_managed, state_directory};
997
998#[cfg(test)]
999#[path = "managed_server_tests.rs"]
1000mod tests;