1use 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
44const 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 #[serde(default, skip_serializing_if = "Option::is_none")]
71 pub ca_cert: Option<String>,
72 #[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
90pub struct ResolvedServer {
92 pub base_url: String,
94 pub management_url: String,
96 pub token: Option<String>,
97 pub source: &'static str,
98 pub run_max_requests: Option<u64>,
99 pub ca_cert: Option<PathBuf>,
101 pub management_ca_cert: Option<PathBuf>,
103 _lease: Option<ManagedLease>,
104}
105
106impl ResolvedServer {
107 #[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 #[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 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 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
172pub 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 #[must_use]
192 pub fn id(&self) -> Option<String> {
193 token_subject(&self.token).ok()
194 }
195
196 #[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 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
232pub 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 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 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 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 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
397pub 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
408pub 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
422pub 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 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 "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
602pub 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
629pub 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#[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
703pub 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 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 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#[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;