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