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