1pub 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
40pub const DEFAULT_TUNNEL_PORT: u16 = 8480;
42
43#[derive(Debug, Clone)]
45pub struct IngressConfig {
46 pub http: Option<SocketAddr>,
48 pub https: Option<SocketAddr>,
49 pub ca: Ca,
50 pub email: Option<String>,
51 pub public_ip: Option<IpAddr>,
53 pub tunnel_port: u16,
54 pub caddy_bin: Option<PathBuf>,
56 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
75pub 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#[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 pub provider: String,
105 pub state: String,
108 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#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
119pub struct CertState {
120 pub state: String,
122 #[serde(skip_serializing_if = "Option::is_none")]
123 pub error: Option<String>,
124 pub issuer: String,
125 pub at: u64,
127}
128
129#[derive(Debug, Clone, Default, PartialEq)]
131struct OrgIngress {
132 domains: Vec<String>,
133 tunnel: bool,
134 account: Option<String>,
135 zone: Option<String>,
136 subnet: Option<String>,
138}
139
140#[derive(Debug, Clone, Default, Serialize)]
142pub struct TunnelStatus {
143 pub org: String,
144 pub origin: Option<String>,
146 pub stack: bool,
148 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
160type Refused = BTreeMap<(String, String), Vec<(String, String, String)>>;
162type Rotation = BTreeMap<(String, String), Vec<IpAddr>>;
164
165#[derive(Default)]
166struct State {
167 defs: Vec<Arc<StackDef>>,
168 rotation: Rotation,
169 extras: Vec<ExtraRoute>,
171 generation: u64,
173 claims: Vec<Claim>,
174 certs: BTreeMap<String, CertState>,
175 served: Vec<Served>,
177 conflicts: Vec<Conflict>,
178 refused: Refused,
181 reported: BTreeSet<String>,
182 orgs: BTreeMap<OrgId, (Instant, OrgIngress)>,
183 tunnels: BTreeMap<OrgId, TunnelStatus>,
184 last_error: Option<String>,
185}
186
187pub 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
201const ORG_TTL: Duration = Duration::from_secs(15);
203const TICK: Duration = Duration::from_secs(15);
205const RESYNC: u64 = 600;
207
208impl Manager {
209 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 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 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 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 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 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 #[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 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 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 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 #[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 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 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
843const 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
859fn 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 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 {
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}