Skip to main content

isb_core/ingress/
mod.rs

1//! Ingress: public hostnames for stack services (docs/guides/domains.md).
2//!
3//! A service's `domains:` ([`crate::spec::DomainSpec`]) become routes on an
4//! HTTP(S) edge, Caddy ([`caddy`]), run by `isb serve`. The controller tells
5//! the [`Manager`] which replicas are in rotation (it is the controller's
6//! [`Observer`]), and the manager regenerates Caddy's whole config and loads
7//! it whenever the routes or their replicas change. Certificates come from
8//! ACME through Caddy.
9//!
10//! Per org, the domains go out through the server's public listeners
11//! (`caddy`, the default) or through the org's own Cloudflare Tunnel
12//! ([`cloudflare`]). An org may only serve names in its allowlist, and a
13//! name served by one org is refused to every other ([`domain::resolve`]).
14
15pub mod caddy;
16pub mod cloudflare;
17pub mod domain;
18mod extra;
19mod status;
20mod tunnels;
21pub use extra::{ExtraRoute, port_service, workspace_route, workspace_stack};
22
23use std::collections::{BTreeMap, BTreeSet};
24use std::net::{IpAddr, SocketAddr};
25use std::path::PathBuf;
26use std::sync::{Arc, Condvar, Mutex, OnceLock};
27use std::time::{Duration, Instant};
28
29use serde::Serialize;
30use serde_json::{Value, json};
31
32use crate::client::Client;
33use crate::error::{Error, Result};
34use crate::org::OrgId;
35use crate::secrets::Secrets;
36use crate::stack::controller::Observer;
37use crate::stack::{Controller, StackDef};
38use caddy::{Ca, Served, TunnelListener, Via};
39use domain::{Claim, Conflict, Route};
40
41/// The port each tunnel org's listener uses on its bridge address.
42pub const DEFAULT_TUNNEL_PORT: u16 = 8480;
43
44/// `isb serve`'s ingress settings.
45#[derive(Debug, Clone)]
46pub struct IngressConfig {
47    /// Public listeners. With neither, only tunnel orgs are served.
48    pub http: Option<SocketAddr>,
49    pub https: Option<SocketAddr>,
50    pub ca: Ca,
51    pub email: Option<String>,
52    /// For `host: auto` names.
53    pub public_ip: Option<IpAddr>,
54    pub tunnel_port: u16,
55    /// A Caddy binary to run instead of the pinned download.
56    pub caddy_bin: Option<PathBuf>,
57    /// The Cloudflare API (a fake one in tests).
58    pub cloudflare_api: String,
59}
60
61impl Default for IngressConfig {
62    fn default() -> Self {
63        IngressConfig {
64            http: None,
65            https: None,
66            ca: Ca::Acme(caddy::LETSENCRYPT.into()),
67            email: None,
68            public_ip: None,
69            tunnel_port: DEFAULT_TUNNEL_PORT,
70            caddy_bin: None,
71            cloudflare_api: cloudflare::API_BASE.into(),
72        }
73    }
74}
75
76/// The host's public IPv4 address, if its default route leaves from one:
77/// the source address the kernel picks toward the internet. No packet is
78/// sent. Behind NAT this is a private address and `None` comes back.
79pub fn detect_public_ip() -> Option<IpAddr> {
80    let s = std::net::UdpSocket::bind("0.0.0.0:0").ok()?;
81    s.connect("1.1.1.1:53").ok()?;
82    match s.local_addr().ok()?.ip() {
83        IpAddr::V4(v4)
84            if !(v4.is_private()
85                || v4.is_loopback()
86                || v4.is_link_local()
87                || v4.is_unspecified()
88                || v4.octets()[0] == 100 && (64..128).contains(&v4.octets()[1])) =>
89        {
90            Some(IpAddr::V4(v4))
91        }
92        _ => None,
93    }
94}
95
96/// One domain of a service, as `stack_status` reports it.
97#[derive(Debug, Clone, Serialize, Default, PartialEq)]
98pub struct DomainStatus {
99    pub host: String,
100    pub path: String,
101    #[serde(skip_serializing_if = "Option::is_none")]
102    pub url: Option<String>,
103    pub https: bool,
104    /// `caddy` or `cloudflare-tunnel`.
105    pub provider: String,
106    /// `serving`, `redirect`, `no-replicas`, `conflict`, `refused` (allowlist,
107    /// settings) or `off` (no listener for it).
108    pub state: String,
109    /// `issued`, `pending`, `failed`, `unsupported` (a wildcard needs a DNS
110    /// challenge), `cloudflare` (TLS ends at Cloudflare), or `none` (plain HTTP).
111    pub cert: String,
112    #[serde(skip_serializing_if = "Option::is_none")]
113    pub message: Option<String>,
114    #[serde(skip_serializing_if = "Vec::is_empty")]
115    pub upstreams: Vec<String>,
116}
117
118/// A certificate's state, from Caddy's log and storage.
119#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
120pub struct CertState {
121    /// `issued`, `failed`, `pending`.
122    pub state: String,
123    #[serde(skip_serializing_if = "Option::is_none")]
124    pub error: Option<String>,
125    pub issuer: String,
126    /// Unix seconds.
127    pub at: u64,
128}
129
130/// An org's ingress settings, from its project.
131#[derive(Debug, Clone, Default, PartialEq)]
132struct OrgIngress {
133    domains: Vec<String>,
134    tunnel: bool,
135    account: Option<String>,
136    zone: Option<String>,
137    /// The bridge address (`10.64.3.1/24`), for the tunnel listener.
138    subnet: Option<String>,
139}
140
141/// Per tunnel org: what isb last did for it.
142#[derive(Debug, Clone, Default, Serialize)]
143pub struct TunnelStatus {
144    pub org: String,
145    /// What cloudflared should send the org's hostnames to.
146    pub origin: Option<String>,
147    /// Whether the cloudflared stack is deployed.
148    pub stack: bool,
149    /// Whether isb manages the tunnel's ingress rules and DNS (an API token
150    /// is in the org's secrets).
151    pub api_managed: bool,
152    #[serde(skip_serializing_if = "Option::is_none")]
153    pub last_sync: Option<cloudflare::SyncReport>,
154    pub last_sync_at: u64,
155    #[serde(skip_serializing_if = "Option::is_none")]
156    pub error: Option<String>,
157    #[serde(skip)]
158    synced_hosts: Vec<String>,
159}
160
161/// Refused routes by (qualified stack, service): (host, path, reason).
162type Refused = BTreeMap<(String, String), Vec<(String, String, String)>>;
163/// In-rotation replica addresses by (qualified stack, service).
164type Rotation = BTreeMap<(String, String), Vec<IpAddr>>;
165
166#[derive(Default)]
167struct State {
168    defs: Vec<Arc<StackDef>>,
169    rotation: Rotation,
170    /// Routes outside any stack: workspace ports ([`extra`]).
171    extras: Vec<ExtraRoute>,
172    /// Bumped by every change; the applier works up to it.
173    generation: u64,
174    claims: Vec<Claim>,
175    certs: BTreeMap<String, CertState>,
176    // What the last computation found.
177    served: Vec<Served>,
178    conflicts: Vec<Conflict>,
179    /// Routes that cannot be served, by (qualified stack, service): the
180    /// reason (allowlist, a bad `auto`, a tunnel org without settings).
181    refused: Refused,
182    reported: BTreeSet<String>,
183    orgs: BTreeMap<OrgId, (Instant, OrgIngress)>,
184    tunnels: BTreeMap<OrgId, TunnelStatus>,
185    last_error: Option<String>,
186}
187
188/// The ingress: keeps Caddy's config in line with the stacks.
189pub struct Manager {
190    cfg: IngressConfig,
191    client: Client,
192    secrets: Arc<Secrets>,
193    dir: PathBuf,
194    state: Mutex<State>,
195    wake: Condvar,
196    applied: Mutex<u64>,
197    applied_cv: Condvar,
198    edge: OnceLock<Arc<caddy::Edge>>,
199    ctl: OnceLock<Controller>,
200}
201
202/// How long org settings are trusted before they are read again.
203const ORG_TTL: Duration = Duration::from_secs(15);
204/// How often the applier wakes without a change (org settings, tunnel sync).
205const TICK: Duration = Duration::from_secs(15);
206/// How often an API-managed tunnel is re-synced when nothing changed.
207const RESYNC: u64 = 600;
208
209impl Manager {
210    /// Create the manager. Caddy starts with [`Manager::start`], once the
211    /// controller exists.
212    pub fn new(
213        cfg: IngressConfig,
214        client: Client,
215        secrets: Arc<Secrets>,
216        state_dir: &std::path::Path,
217    ) -> Result<Arc<Manager>> {
218        let dir = state_dir.join("ingress");
219        std::fs::create_dir_all(&dir)?;
220        let claims = load_claims(&dir.join("claims.json"));
221        Ok(Arc::new(Manager {
222            cfg,
223            client,
224            secrets,
225            dir,
226            state: Mutex::new(State {
227                claims,
228                ..Default::default()
229            }),
230            wake: Condvar::new(),
231            applied: Mutex::new(0),
232            applied_cv: Condvar::new(),
233            edge: OnceLock::new(),
234            ctl: OnceLock::new(),
235        }))
236    }
237
238    /// Caddy's admin socket, in a directory only we can open. A unix socket
239    /// path holds at most 107 bytes: a long state directory moves it under
240    /// the runtime directory, named by a hash of the state directory.
241    fn admin_socket(&self) -> PathBuf {
242        let p = self.dir.join("run").join("admin.sock");
243        if p.as_os_str().len() <= 100 {
244            return p;
245        }
246        let mut h: u32 = 0x811c9dc5;
247        for b in self.dir.as_os_str().as_encoded_bytes() {
248            h ^= *b as u32;
249            h = h.wrapping_mul(0x01000193);
250        }
251        let base = std::env::var_os("XDG_RUNTIME_DIR")
252            .map(PathBuf::from)
253            .unwrap_or_else(|| {
254                PathBuf::from(format!("/tmp/isb-{}", rustix::process::getuid().as_raw()))
255            });
256        base.join("isb")
257            .join(format!("caddy-{h:08x}"))
258            .join("admin.sock")
259    }
260
261    fn storage(&self) -> PathBuf {
262        self.dir.join("caddy")
263    }
264
265    fn params(&self) -> caddy::Params {
266        caddy::Params {
267            admin_socket: self.admin_socket(),
268            storage: self.storage(),
269            http: self.cfg.http,
270            https: self.cfg.https,
271            ca: self.cfg.ca.clone(),
272            email: self.cfg.email.clone(),
273        }
274    }
275
276    /// Get Caddy (downloading the pinned release if needed), start it, and
277    /// start the applier.
278    pub fn start(self: &Arc<Self>, ctl: Controller) -> Result<()> {
279        let _ = self.ctl.set(ctl);
280        let bin = caddy::ensure_binary(&self.dir.join("bin"), self.cfg.caddy_bin.as_deref())?;
281        let initial = caddy::render(&self.params(), &[], &[]);
282        let weak = Arc::downgrade(self);
283        let on_log: caddy::OnLog = Arc::new(move |ev| {
284            if let Some(m) = weak.upgrade() {
285                m.cert_event(ev);
286            }
287        });
288        let edge = caddy::Edge::start(bin, &self.dir, self.admin_socket(), initial, on_log)?;
289        let _ = self.edge.set(edge);
290        let m = self.clone();
291        std::thread::Builder::new()
292            .name("isb-ingress".into())
293            .spawn(move || m.applier())?;
294        self.bump();
295        Ok(())
296    }
297
298    /// The address generated (`host: auto`) names point at.
299    pub fn public_ip(&self) -> Option<IpAddr> {
300        self.cfg.public_ip
301    }
302
303    pub fn shutdown(&self) {
304        if let Some(e) = self.edge.get() {
305            e.shutdown();
306        }
307    }
308
309    fn bump(&self) {
310        let mut st = self.state.lock().unwrap();
311        st.generation += 1;
312        self.wake.notify_all();
313    }
314
315    fn cert_event(&self, ev: caddy::LogEvent) {
316        let issuer_key = self.cfg.ca.issuer_key();
317        let (host, cs, kind, level, msg) = match ev {
318            caddy::LogEvent::Obtained { host, issuer } => {
319                if !issuer.is_empty() && issuer != issuer_key {
320                    return;
321                }
322                let msg = format!("certificate issued for {host} ({issuer})");
323                (
324                    host,
325                    CertState {
326                        state: "issued".into(),
327                        error: None,
328                        issuer,
329                        at: crate::stack::now_secs(),
330                    },
331                    "cert.issued",
332                    "info",
333                    msg,
334                )
335            }
336            caddy::LogEvent::Failed { host, error } => {
337                let msg = format!("certificate for {host} failed: {error}");
338                (
339                    host,
340                    CertState {
341                        state: "failed".into(),
342                        error: Some(error),
343                        issuer: issuer_key,
344                        at: crate::stack::now_secs(),
345                    },
346                    "cert.failed",
347                    "warn",
348                    msg,
349                )
350            }
351        };
352        let owners: Vec<(String, String)> = {
353            let mut st = self.state.lock().unwrap();
354            let same = st
355                .certs
356                .get(&host)
357                .is_some_and(|c| c.state == cs.state && c.error == cs.error);
358            st.certs.insert(host.clone(), cs);
359            if same {
360                return;
361            }
362            st.served
363                .iter()
364                .filter(|s| s.route.host == host)
365                .map(|s| (s.route.qualified_stack(), s.route.service.clone()))
366                .collect::<BTreeSet<_>>()
367                .into_iter()
368                .collect()
369        };
370        if let Some(ctl) = self.ctl.get() {
371            for (q, svc) in owners {
372                ctl.event(kind, level, &q, &svc, msg.clone());
373            }
374        }
375    }
376
377    /// Check a stack about to be deployed: its domains must be allowed for
378    /// its org and free of other orgs' (and other services') claims.
379    pub fn check(&self, def: &StackDef) -> Result<()> {
380        let mut routes = Vec::new();
381        for (svc, spec) in &def.file.services {
382            if spec.domains.is_empty() {
383                continue;
384            }
385            routes.extend(domain::routes_for(
386                &def.org,
387                &def.name,
388                svc,
389                &spec.domains,
390                self.cfg.public_ip,
391            )?);
392        }
393        if routes.is_empty() {
394            return Ok(());
395        }
396        let oi = self.org_settings(&def.org, true)?;
397        for r in &routes {
398            if !r.auto && !domain::allowed(&r.host, &oi.domains) {
399                return Err(Error::invalid(not_allowed(&r.host, &def.org, &oi.domains)));
400            }
401        }
402        let (defs, claims) = {
403            let st = self.state.lock().unwrap();
404            (st.defs.clone(), st.claims.clone())
405        };
406        let q = def.qualified();
407        let mut all: Vec<Route> = routes.clone();
408        for d in defs.iter().filter(|d| d.qualified() != q) {
409            for (svc, spec) in &d.file.services {
410                if let Ok(r) =
411                    domain::routes_for(&d.org, &d.name, svc, &spec.domains, self.cfg.public_ip)
412                {
413                    all.extend(r);
414                }
415            }
416        }
417        let res = domain::resolve(&all, &claims, crate::stack::now_secs());
418        let mine: Vec<String> = res
419            .conflicts
420            .iter()
421            .filter(|c| c.route.org == def.org && c.route.stack == def.name)
422            .map(|c| format!("service {}: {}", c.route.service, c.reason))
423            .collect();
424        if !mine.is_empty() {
425            return Err(Error::invalid(format!(
426                "domain conflict: {}",
427                mine.join("; ")
428            )));
429        }
430        Ok(())
431    }
432
433    /// An org's settings, cached for [`ORG_TTL`] unless `fresh`.
434    fn org_settings(&self, org: &OrgId, fresh: bool) -> Result<OrgIngress> {
435        if !fresh {
436            if let Some((at, oi)) = self.state.lock().unwrap().orgs.get(org) {
437                if at.elapsed() < ORG_TTL {
438                    return Ok(oi.clone());
439                }
440            }
441        }
442        let info = crate::org::get(&self.client, org)?;
443        let oi = OrgIngress {
444            domains: info.domains,
445            tunnel: info.ingress == crate::org::INGRESS_CLOUDFLARE_TUNNEL,
446            account: info.cloudflare_account,
447            zone: info.cloudflare_zone,
448            subnet: info.subnet,
449        };
450        self.state
451            .lock()
452            .unwrap()
453            .orgs
454            .insert(org.clone(), (Instant::now(), oi.clone()));
455        Ok(oi)
456    }
457
458    fn applier(self: Arc<Self>) {
459        let mut done = 0u64;
460        let mut last_cfg: Option<Value> = None;
461        loop {
462            let target_gen = {
463                let st = self.state.lock().unwrap();
464                let (st, _) = self
465                    .wake
466                    .wait_timeout_while(st, TICK, |s| s.generation == done)
467                    .unwrap();
468                st.generation
469            };
470            if let Err(e) = self.apply_once(&mut last_cfg) {
471                let msg = e.to_string();
472                let mut st = self.state.lock().unwrap();
473                if st.last_error.as_deref() != Some(&msg) {
474                    eprintln!("isb serve: ingress: {msg}");
475                }
476                st.last_error = Some(msg);
477            } else {
478                self.state.lock().unwrap().last_error = None;
479            }
480            done = target_gen;
481            *self.applied.lock().unwrap() = target_gen;
482            self.applied_cv.notify_all();
483        }
484    }
485
486    /// Compute the routes and load them into Caddy when they changed.
487    #[expect(
488        clippy::too_many_lines,
489        clippy::excessive_nesting,
490        reason = "predates the lint ratchet; split it when next changed"
491    )]
492    fn apply_once(&self, last_cfg: &mut Option<Value>) -> Result<()> {
493        let (defs, mut rotation, claims, extras) = {
494            let st = self.state.lock().unwrap();
495            let s = &st;
496            (
497                s.defs.clone(),
498                s.rotation.clone(),
499                s.claims.clone(),
500                s.extras.clone(),
501            )
502        };
503        let mut orgs: BTreeMap<OrgId, OrgIngress> = BTreeMap::new();
504        let mut org_errors: BTreeMap<OrgId, String> = BTreeMap::new();
505        let mut want: BTreeSet<OrgId> = defs
506            .iter()
507            .filter(|d| d.file.services.values().any(|s| !s.domains.is_empty()))
508            .map(|d| d.org.clone())
509            .collect();
510        // Tunnel orgs we serve now, whose settings may have changed.
511        want.extend(self.state.lock().unwrap().tunnels.keys().cloned());
512        want.extend(extras.iter().map(|e| e.route.org.clone()));
513        for o in want {
514            match self.org_settings(&o, false) {
515                Ok(oi) => {
516                    orgs.insert(o, oi);
517                }
518                Err(e) => {
519                    org_errors.insert(o, e.to_string());
520                }
521            }
522        }
523
524        let mut routes = Vec::new();
525        let mut refused: Refused = BTreeMap::new();
526        for d in &defs {
527            for (svc, spec) in &d.file.services {
528                if spec.domains.is_empty() {
529                    continue;
530                }
531                let key = (d.qualified(), svc.clone());
532                let refuse =
533                    |refused: &mut BTreeMap<_, Vec<_>>, host: &str, path: &str, why: String| {
534                        refused.entry(key.clone()).or_insert_with(Vec::new).push((
535                            host.to_string(),
536                            path.to_string(),
537                            why,
538                        ));
539                    };
540                if let Some(e) = org_errors.get(&d.org) {
541                    refuse(&mut refused, "", "", format!("org settings: {e}"));
542                    continue;
543                }
544                let oi = orgs.get(&d.org).cloned().unwrap_or_default();
545                match domain::routes_for(&d.org, &d.name, svc, &spec.domains, self.cfg.public_ip) {
546                    Ok(rs) => {
547                        for r in rs {
548                            if !r.auto && !domain::allowed(&r.host, &oi.domains) {
549                                let why = not_allowed(&r.host, &d.org, &oi.domains);
550                                refuse(&mut refused, &r.host, &r.path, why);
551                            } else {
552                                routes.push(r);
553                            }
554                        }
555                    }
556                    Err(e) => refuse(&mut refused, "", "", e.to_string()),
557                }
558            }
559        }
560        let out = (&mut routes, &mut rotation, &mut refused);
561        Self::merge_extras(&extras, &orgs, &org_errors, out);
562        let res = domain::resolve(&routes, &claims, crate::stack::now_secs());
563
564        // Tunnel listeners, one per tunnel org with an address.
565        let mut tunnels = Vec::new();
566        for (o, oi) in &orgs {
567            if !oi.tunnel {
568                continue;
569            }
570            if let Some((gw, subnet)) = oi.subnet.as_deref().and_then(gateway) {
571                tunnels.push(TunnelListener {
572                    org: o.clone(),
573                    listen: SocketAddr::new(gw, self.cfg.tunnel_port),
574                    subnet,
575                });
576            }
577        }
578        let served: Vec<Served> = res
579            .accepted
580            .iter()
581            .map(|r| {
582                let ips = rotation
583                    .get(&(r.qualified_stack(), r.service.clone()))
584                    .cloned()
585                    .unwrap_or_default();
586                let upstreams = match r.port {
587                    Some(p) => ips.iter().map(|ip| SocketAddr::new(*ip, p)).collect(),
588                    None => Vec::new(),
589                };
590                let via = if orgs.get(&r.org).is_some_and(|o| o.tunnel) {
591                    Via::Tunnel(r.org.clone())
592                } else {
593                    Via::Public
594                };
595                Served {
596                    route: r.clone(),
597                    upstreams,
598                    via,
599                }
600            })
601            .collect();
602
603        // New conflicts become events, once each.
604        let mut new_events = Vec::new();
605        {
606            let mut st = self.state.lock().unwrap();
607            let mut now_reported = BTreeSet::new();
608            for c in &res.conflicts {
609                let k = format!(
610                    "{}|{}|{}|{}",
611                    c.route.qualified_stack(),
612                    c.route.service,
613                    c.route.host,
614                    c.route.path
615                );
616                if !st.reported.contains(&k) {
617                    new_events.push((
618                        c.route.qualified_stack(),
619                        c.route.service.clone(),
620                        format!("domain conflict: {}", c.reason),
621                    ));
622                }
623                now_reported.insert(k);
624            }
625            for (k, list) in &refused {
626                for (h, p, why) in list {
627                    let key = format!("{}|{}|{h}|{p}|refused", k.0, k.1);
628                    if !st.reported.contains(&key) {
629                        new_events.push((
630                            k.0.clone(),
631                            k.1.clone(),
632                            format!("domain refused: {why}"),
633                        ));
634                    }
635                    now_reported.insert(key);
636                }
637            }
638            st.reported = now_reported;
639            if st.claims != res.claims {
640                st.claims = res.claims.clone();
641                save_claims(&self.dir.join("claims.json"), &st.claims);
642            }
643            st.served = served.clone();
644            st.conflicts = res.conflicts.clone();
645            st.refused = refused;
646        }
647        if let Some(ctl) = self.ctl.get() {
648            for (q, svc, msg) in new_events {
649                ctl.service_event("warn", &q, &svc, msg);
650            }
651        }
652
653        let cfg = caddy::render(&self.params(), &served, &tunnels);
654        let mut result = Ok(());
655        if last_cfg.as_ref() != Some(&cfg) {
656            if let Some(edge) = self.edge.get() {
657                match edge.load(cfg.clone()) {
658                    Ok(()) => *last_cfg = Some(cfg),
659                    Err(e) => result = Err(e),
660                }
661            }
662        }
663        self.tunnels(&orgs, &served, &tunnels);
664        result
665    }
666}
667
668fn not_allowed(host: &str, org: &OrgId, list: &[String]) -> String {
669    if list.is_empty() {
670        format!(
671            "{host}: org {org} may not serve wildcard hosts (no wildcard in its --allow-domain list)"
672        )
673    } else {
674        format!(
675            "{host} is outside org {org}'s domains ({})",
676            list.join(", ")
677        )
678    }
679}
680
681/// The bridge's own address and its subnet: `10.64.3.1/24` ->
682/// (10.64.3.1, 10.64.3.0/24).
683fn gateway(cidr: &str) -> Option<(IpAddr, String)> {
684    let (ip, len) = cidr.split_once('/')?;
685    let ip: std::net::Ipv4Addr = ip.parse().ok()?;
686    let len: u32 = len.parse().ok()?;
687    if len > 32 {
688        return None;
689    }
690    let mask = if len == 0 { 0 } else { u32::MAX << (32 - len) };
691    let net = std::net::Ipv4Addr::from(u32::from(ip) & mask);
692    Some((IpAddr::V4(ip), format!("{net}/{len}")))
693}
694
695fn load_claims(p: &std::path::Path) -> Vec<Claim> {
696    std::fs::read(p)
697        .ok()
698        .and_then(|b| serde_json::from_slice(&b).ok())
699        .unwrap_or_default()
700}
701
702fn save_claims(p: &std::path::Path, claims: &[Claim]) {
703    let tmp = p.with_extension("json.tmp");
704    let ok = serde_json::to_vec_pretty(claims)
705        .ok()
706        .and_then(|b| std::fs::write(&tmp, b).ok())
707        .and_then(|_| std::fs::rename(&tmp, p).ok());
708    if ok.is_none() {
709        eprintln!("isb serve: ingress: cannot save {}", p.display());
710    }
711}
712
713impl Observer for Manager {
714    fn rotation(&self, stack: &str, service: &str, ips: &[IpAddr]) {
715        let mut st = self.state.lock().unwrap();
716        let k = (stack.to_string(), service.to_string());
717        if st.rotation.get(&k).map(|v| v.as_slice()) == Some(ips) {
718            return;
719        }
720        if ips.is_empty() {
721            st.rotation.remove(&k);
722        } else {
723            st.rotation.insert(k, ips.to_vec());
724        }
725        // Only services with domains change the config.
726        let routed = st
727            .served
728            .iter()
729            .any(|s| s.route.service == service && s.route.qualified_stack() == stack)
730            || st.defs.iter().any(|d| {
731                d.qualified() == stack
732                    && d.file
733                        .services
734                        .get(service)
735                        .is_some_and(|s| !s.domains.is_empty())
736            });
737        if routed {
738            st.generation += 1;
739            self.wake.notify_all();
740        }
741    }
742
743    fn drain(&self, stack: &str, service: &str, ip: IpAddr, timeout: Duration) {
744        let started = Instant::now();
745        let target_gen = {
746            let st = self.state.lock().unwrap();
747            let routed = st
748                .served
749                .iter()
750                .any(|s| s.route.service == service && s.route.qualified_stack() == stack);
751            if !routed {
752                return;
753            }
754            st.generation
755        };
756        // The config without the replica is loaded (or Caddy is down).
757        {
758            let g = self.applied.lock().unwrap();
759            let _ = self
760                .applied_cv
761                .wait_timeout_while(g, timeout, |a| *a < target_gen)
762                .unwrap();
763        }
764        let Some(edge) = self.edge.get() else { return };
765        let prefix = format!("{ip}:");
766        while started.elapsed() < timeout {
767            match edge.upstreams() {
768                Ok(ups) => {
769                    if !ups.iter().any(|(a, n)| a.starts_with(&prefix) && *n > 0) {
770                        return;
771                    }
772                }
773                Err(_) => return,
774            }
775            std::thread::sleep(Duration::from_millis(100));
776        }
777    }
778
779    fn stacks_changed(&self, defs: Vec<Arc<StackDef>>) {
780        let mut st = self.state.lock().unwrap();
781        st.defs = defs;
782        st.generation += 1;
783        self.wake.notify_all();
784    }
785
786    fn domains(&self, stack: &str, service: &str) -> Vec<DomainStatus> {
787        let st = self.state.lock().unwrap();
788        let mut out: Vec<DomainStatus> = st
789            .served
790            .iter()
791            .filter(|s| s.route.service == service && s.route.qualified_stack() == stack)
792            .map(|s| self.domain_status_of(&st, s))
793            .collect();
794        for c in st
795            .conflicts
796            .iter()
797            .filter(|c| c.route.service == service && c.route.qualified_stack() == stack)
798        {
799            out.push(DomainStatus {
800                host: c.route.host.clone(),
801                path: c.route.path.clone(),
802                https: c.route.https,
803                state: "conflict".into(),
804                cert: "none".into(),
805                message: Some(c.reason.clone()),
806                ..Default::default()
807            });
808        }
809        if let Some(list) = st.refused.get(&(stack.to_string(), service.to_string())) {
810            for (h, p, why) in list {
811                out.push(DomainStatus {
812                    host: h.clone(),
813                    path: p.clone(),
814                    state: "refused".into(),
815                    cert: "none".into(),
816                    message: Some(why.clone()),
817                    ..Default::default()
818                });
819            }
820        }
821        out
822    }
823}
824
825#[cfg(test)]
826mod tests {
827    use super::*;
828
829    #[test]
830    fn gateways() {
831        assert_eq!(
832            gateway("10.64.3.1/24"),
833            Some(("10.64.3.1".parse().unwrap(), "10.64.3.0/24".into()))
834        );
835        assert_eq!(gateway("nope"), None);
836    }
837
838    #[test]
839    fn claims_round_trip() {
840        let d = tempfile::tempdir().unwrap();
841        let p = d.path().join("claims.json");
842        assert!(load_claims(&p).is_empty());
843        let c = vec![Claim {
844            host: "a.example.com".into(),
845            path: "/".into(),
846            org: "acme".into(),
847            stack: "s".into(),
848            service: "w".into(),
849            since: 5,
850        }];
851        save_claims(&p, &c);
852        assert_eq!(load_claims(&p), c);
853    }
854}