1use std::collections::{BTreeMap, BTreeSet, VecDeque};
11use std::net::{IpAddr, SocketAddr};
12use std::sync::atomic::{AtomicBool, Ordering};
13use std::sync::{Arc, Condvar, Mutex};
14use std::time::{Duration, Instant};
15
16use serde::Serialize;
17use serde_json::Value;
18
19use super::changes::diff;
20use super::{
21 LABEL_REV, LABEL_SERVICE, LABEL_SLOT, LABEL_STACK, StackDef, Store, instance_name, new_id,
22 now_secs, validate_stack_name,
23};
24use crate::balance::Balancer;
25use crate::client::{Client, encode_query, encode_segment};
26use crate::error::{Error, Result};
27use crate::org::OrgId;
28use crate::plan::{Desired, split_addr};
29use crate::sandbox::{EnsureOptions, Sandbox};
30use crate::secrets::Secrets;
31use crate::spec::{
32 DependCondition, FailureAction, HealthProbe, PortBind, RestartCondition, RestartMode,
33 SandboxSpec, UpdateConfig, UpdateOrder,
34};
35use crate::supervise;
36
37const DRAIN: Duration = Duration::from_secs(10);
39const RESTARTS_BEFORE_REPLACE: u32 = 3;
41const HEALTH_DEBOUNCE: Duration = Duration::from_secs(20);
44
45#[derive(Debug, Clone, Serialize)]
47pub struct InstanceStatus {
48 pub name: String,
49 pub slot: u32,
50 pub rev: String,
51 pub status: String,
53 pub health: String,
56 pub ip: Option<String>,
57 pub in_rotation: bool,
59 pub restarts: u32,
60 #[serde(skip_serializing_if = "String::is_empty")]
61 pub last_probe: String,
62 pub cpu_pct: Option<f32>,
64 pub cpu_history: Vec<f32>,
65 pub mem_bytes: Option<u64>,
66 pub disk_bytes: Option<u64>,
68}
69
70#[derive(Debug, Clone, Serialize)]
72pub struct PortStatus {
73 pub listen: String,
74 pub target: u16,
75 pub backends: Vec<String>,
76 #[serde(skip_serializing_if = "Option::is_none")]
77 pub error: Option<String>,
78 pub accepted: u64,
80 pub active: usize,
81 pub rate_history: Vec<f32>,
82}
83
84#[derive(Debug, Clone, Serialize, Default)]
86pub struct ServiceStatus {
87 pub service: String,
88 pub image: String,
89 pub rev: String,
90 pub replicas: u32,
91 pub running: u32,
92 pub healthy: u32,
93 pub state: String,
95 #[serde(skip_serializing_if = "Option::is_none")]
96 pub message: Option<String>,
97 pub instances: Vec<InstanceStatus>,
98 pub ports: Vec<PortStatus>,
99 #[serde(skip_serializing_if = "Option::is_none")]
101 pub rollout: Option<RolloutStatus>,
102 pub checked_at: u64,
104 #[serde(skip_serializing_if = "Vec::is_empty")]
106 pub domains: Vec<crate::ingress::DomainStatus>,
107}
108
109pub trait Observer: Send + Sync {
112 fn rotation(&self, stack: &str, service: &str, ips: &[IpAddr]);
114 fn drain(&self, stack: &str, service: &str, ip: IpAddr, timeout: Duration);
117 fn stacks_changed(&self, defs: Vec<Arc<StackDef>>);
119 fn domains(&self, stack: &str, service: &str) -> Vec<crate::ingress::DomainStatus>;
121}
122
123#[derive(Debug, Clone, Serialize, Default)]
125pub struct RolloutStatus {
126 pub to_rev: String,
127 pub order: String,
129 pub parallelism: usize,
130 pub done: usize,
131 pub total: usize,
132 pub started_at: u64,
134 pub slots: Vec<SlotRollout>,
135}
136
137#[derive(Debug, Clone, Serialize, Default)]
139pub struct SlotRollout {
140 pub slot: u32,
141 pub old: Option<String>,
142 pub old_rev: Option<String>,
143 pub old_state: String,
145 pub new: Option<String>,
146 pub new_state: String,
148}
149
150#[derive(Debug, Clone, Serialize)]
152pub struct Event {
153 pub seq: u64,
155 pub at: u64,
157 pub level: String,
159 pub stack: String,
160 #[serde(skip_serializing_if = "String::is_empty")]
161 pub service: String,
162 #[serde(skip_serializing_if = "Option::is_none")]
163 pub instance: Option<String>,
164 pub message: String,
165 #[serde(default, skip_serializing_if = "Option::is_none")]
176 pub kind: Option<String>,
177}
178
179pub type EventSink = Arc<dyn Fn(&Event) + Send + Sync>;
182
183const EVENTS_KEPT: usize = 1000;
185
186#[derive(Debug, Clone, Default)]
188pub struct Snapshot {
189 pub host: crate::metrics::HostSample,
190 pub instances: BTreeMap<String, crate::metrics::InstanceSample>,
192 pub at: u64,
194}
195
196pub fn now_ms() -> u64 {
197 std::time::SystemTime::now()
198 .duration_since(std::time::UNIX_EPOCH)
199 .map(|d| d.as_millis() as u64)
200 .unwrap_or(0)
201}
202
203#[derive(Debug, Clone, Serialize)]
205pub struct StackStatus {
206 pub name: String,
207 pub org: String,
208 pub deployed_at: u64,
209 pub deployed_by: String,
210 pub has_previous: bool,
211 pub converged: bool,
213 pub services: Vec<ServiceStatus>,
214}
215
216#[derive(Debug, Clone, Serialize)]
218pub struct DeployChange {
219 pub service: String,
220 pub change: String,
223 pub rev: String,
224 pub replicas: u32,
225}
226
227struct Slot {
229 def: Arc<StackDef>,
230 remove: bool,
232 remove_volumes: bool,
234}
235
236struct WorkerShared {
237 slot: Mutex<Slot>,
238 wake: Condvar,
239 stop: AtomicBool,
240 kick: AtomicBool,
243}
244
245struct Inner {
246 client: Client,
247 store: Store,
248 balancer: Balancer,
249 interval: Duration,
250 workers: Mutex<BTreeMap<(String, String), Arc<WorkerShared>>>,
252 stacks: Mutex<BTreeMap<String, Arc<StackDef>>>,
254 status: Mutex<BTreeMap<(String, String), ServiceStatus>>,
255 failures: super::failure::Failures,
257 events: Mutex<(u64, VecDeque<Event>)>,
258 snapshot: Mutex<Snapshot>,
259 secrets: Arc<Secrets>,
261 edit: Mutex<()>,
264 refresh: Mutex<super::secrets::RefreshSchedule>,
266 observer: Option<Arc<dyn Observer>>,
267 event_sink: Mutex<Option<EventSink>>,
269 metrics_sink: Mutex<Option<std::sync::mpsc::SyncSender<crate::metrics_history::Sample>>>,
271}
272
273impl Inner {
274 fn emit(
275 &self,
276 level: &str,
277 stack: &str,
278 service: &str,
279 instance: Option<&str>,
280 message: String,
281 ) {
282 self.emit_kind(None, level, stack, service, instance, message);
283 }
284
285 fn emit_kind(
286 &self,
287 kind: Option<&str>,
288 level: &str,
289 stack: &str,
290 service: &str,
291 instance: Option<&str>,
292 message: String,
293 ) {
294 let mut e = self.events.lock().unwrap();
295 e.0 += 1;
296 let ev = Event {
297 seq: e.0,
298 at: now_ms(),
299 level: level.into(),
300 stack: stack.into(),
301 service: service.into(),
302 instance: instance.map(String::from),
303 message,
304 kind: kind.map(String::from),
305 };
306 if e.1.len() == EVENTS_KEPT {
307 e.1.pop_front();
308 }
309 let sink = self.event_sink.lock().unwrap().clone();
310 if let Some(s) = sink {
311 s(&ev);
312 }
313 e.1.push_back(ev);
314 }
315}
316
317const SAMPLE_EVERY: Duration = Duration::from_secs(2);
319const SECRET_TICK: Duration = Duration::from_secs(10);
321
322#[derive(Clone)]
324pub struct Controller {
325 inner: Arc<Inner>,
326}
327
328impl Controller {
329 pub fn start(
332 client: Client,
333 store: Store,
334 interval: Duration,
335 secrets: Arc<Secrets>,
336 ) -> Result<Controller> {
337 Controller::start_with(client, store, interval, secrets, None)
338 }
339
340 pub fn start_with(
343 client: Client,
344 store: Store,
345 interval: Duration,
346 secrets: Arc<Secrets>,
347 observer: Option<Arc<dyn Observer>>,
348 ) -> Result<Controller> {
349 let c = Controller {
350 inner: Arc::new(Inner {
351 client,
352 store,
353 balancer: Balancer::new(),
354 interval,
355 workers: Mutex::new(BTreeMap::new()),
356 stacks: Mutex::new(BTreeMap::new()),
357 status: Mutex::new(BTreeMap::new()),
358 failures: Default::default(),
359 events: Mutex::new((0, VecDeque::new())),
360 snapshot: Mutex::new(Snapshot::default()),
361 secrets,
362 edit: Mutex::new(()),
363 refresh: Mutex::new(Default::default()),
364 observer,
365 event_sink: Mutex::new(None),
366 metrics_sink: Mutex::new(None),
367 }),
368 };
369 let weak = Arc::downgrade(&c.inner);
372 let _ = std::thread::Builder::new()
373 .name("isb-secrets".into())
374 .spawn(move || {
375 while let Some(inner) = weak.upgrade() {
376 let tick = inner.interval.clamp(Duration::from_secs(1), SECRET_TICK);
377 Controller { inner }.check_due_secrets();
378 std::thread::sleep(tick);
379 }
380 });
381 let weak = Arc::downgrade(&c.inner);
383 let _ = std::thread::Builder::new()
384 .name("isb-metrics".into())
385 .spawn(move || {
386 let mut sampler = crate::metrics::Sampler::new();
387 while let Some(inner) = weak.upgrade() {
388 match sampler.sample(&inner.client) {
389 Ok((host, insts)) => {
390 if let Some(tx) = &*inner.metrics_sink.lock().unwrap() {
391 crate::metrics_history::offer(tx, (now_ms(), insts.clone()));
392 }
393 *inner.snapshot.lock().unwrap() = Snapshot {
394 host,
395 instances: insts
396 .into_iter()
397 .map(|i| (format!("{}/{}", i.project, i.name), i))
398 .collect(),
399 at: now_ms(),
400 };
401 }
402 Err(e) => eprintln!("isb serve: metrics: {e}"),
403 }
404 drop(inner);
405 std::thread::sleep(SAMPLE_EVERY);
406 }
407 });
408 let defs = c.inner.store.load_all()?;
409 for def in &defs {
411 let dir = crate::discovery::org_dir(&def.org);
412 let keep: Vec<(String, String)> = defs
413 .iter()
414 .filter(|d| d.org == def.org)
415 .flat_map(|d| d.file.services.keys().map(|s| (d.name.clone(), s.clone())))
416 .collect();
417 crate::discovery::prune(&dir, &keep);
418 }
419 for def in defs {
420 eprintln!("isb serve: resuming stack {}", def.name);
421 c.apply(Arc::new(def));
422 }
423 for def in c.definitions() {
425 let keys: Vec<String> = def.secrets.keys().cloned().collect();
426 c.check_bindings(&def.qualified(), &keys, false);
427 }
428 Ok(c)
429 }
430
431 pub fn secret_changed(&self, org: &OrgId, name: &str) -> Vec<String> {
435 let mut rolled = Vec::new();
436 for def in self.definitions() {
437 if def.org != *org {
438 continue;
439 }
440 let keys: Vec<String> = def
441 .secrets
442 .iter()
443 .filter(|(_, b)| b.name == name)
444 .map(|(k, _)| k.clone())
445 .collect();
446 if !keys.is_empty() && self.check_bindings(&def.qualified(), &keys, false) {
447 rolled.push(def.qualified());
448 }
449 }
450 rolled
451 }
452
453 pub fn refresh_secret(
457 &self,
458 org: &OrgId,
459 name: &str,
460 ) -> Result<crate::stack::secrets::Refreshed> {
461 let mut found = Vec::new();
462 let mut rolled = Vec::new();
463 for def in self.definitions() {
464 if def.org != *org {
465 continue;
466 }
467 let q = def.qualified();
468 let keys: Vec<String> = def
469 .secrets
470 .iter()
471 .filter(|(_, b)| b.name == name)
472 .map(|(k, _)| k.clone())
473 .collect();
474 for k in &keys {
475 let b = &def.secrets[k];
476 let v = self.inner.secrets.refresh_in(&b.driver, org, name)?;
477 found.push((b.driver.clone(), v));
478 let every = def
479 .file
480 .secrets
481 .get(k)
482 .map(crate::spec::SecretDef::refresh_interval)
483 .unwrap_or(crate::spec::DEFAULT_SECRET_REFRESH);
484 self.inner
485 .refresh
486 .lock()
487 .unwrap()
488 .reset(&q, k, every, Instant::now());
489 }
490 if !keys.is_empty() && self.check_bindings(&q, &keys, true) {
491 rolled.push(q);
492 }
493 }
494 Ok((found, rolled))
495 }
496
497 fn check_due_secrets(&self) {
499 let defs: Vec<(String, Arc<StackDef>)> = self
500 .inner
501 .stacks
502 .lock()
503 .unwrap()
504 .iter()
505 .map(|(q, d)| (q.clone(), d.clone()))
506 .collect();
507 let due = self
508 .inner
509 .refresh
510 .lock()
511 .unwrap()
512 .due(defs.iter().map(|(q, d)| (q.as_str(), &**d)), Instant::now());
513 let mut by_stack: BTreeMap<String, Vec<String>> = BTreeMap::new();
514 for (q, k) in due {
515 by_stack.entry(q).or_default().push(k);
516 }
517 for (q, keys) in by_stack {
518 self.check_bindings(&q, &keys, false);
519 }
520 }
521
522 fn check_bindings(&self, q: &str, keys: &[String], quiet: bool) -> bool {
526 let _g = self.inner.edit.lock().unwrap();
527 let Ok(cur) = self.get_def(q) else {
528 return false;
529 };
530 let mut def = (*cur).clone();
531 let mut moved = Vec::new();
532 for k in keys {
533 let Some(b) = def.secrets.get_mut(k) else {
534 continue;
535 };
536 match self.inner.secrets.version_in(&b.driver, &def.org, &b.name) {
537 Ok(v) if v != b.version => {
538 moved.push(format!("{} v{} -> v{v}", b.name, b.version));
539 b.version = v;
540 }
541 Ok(_) => {}
542 Err(e) if !quiet => self.note(
544 "warn",
545 q,
546 format!("secret {}: cannot check its version: {e}", b.name),
547 ),
548 Err(_) => {}
549 }
550 }
551 if moved.is_empty() {
552 return false;
553 }
554 if let Err(e) = self.inner.store.save(&def) {
555 self.note("error", q, format!("cannot save new secret versions: {e}"));
556 return false;
557 }
558 self.apply(Arc::new(def));
559 self.note(
560 "info",
561 q,
562 format!("new secret version ({}): rolling", moved.join(", ")),
563 );
564 true
565 }
566
567 pub fn balancer(&self) -> &Balancer {
568 &self.inner.balancer
569 }
570
571 pub fn events(&self, since: u64, limit: usize) -> (u64, Vec<Event>) {
573 let e = self.inner.events.lock().unwrap();
574 let out: Vec<Event> = e.1.iter().filter(|x| x.seq > since).cloned().collect();
575 let skip = out.len().saturating_sub(limit);
576 (e.0, out.into_iter().skip(skip).collect())
577 }
578
579 pub fn wait_events(&self, since: u64, limit: usize, timeout: Duration) -> (u64, Vec<Event>) {
581 let started = Instant::now();
582 loop {
583 let r = self.events(since, limit);
584 if !r.1.is_empty() || started.elapsed() >= timeout {
585 return r;
586 }
587 std::thread::sleep(Duration::from_millis(250));
588 }
589 }
590
591 pub fn set_event_sink(&self, sink: EventSink) {
595 let e = self.inner.events.lock().unwrap();
596 for ev in &e.1 {
597 sink(ev);
598 }
599 *self.inner.event_sink.lock().unwrap() = Some(sink);
600 }
601
602 pub fn set_metrics_sink(
604 &self,
605 tx: std::sync::mpsc::SyncSender<crate::metrics_history::Sample>,
606 ) {
607 *self.inner.metrics_sink.lock().unwrap() = Some(tx);
608 }
609
610 pub fn snapshot(&self) -> Snapshot {
612 self.inner.snapshot.lock().unwrap().clone()
613 }
614
615 pub fn event(&self, kind: &str, level: &str, stack: &str, service: &str, message: String) {
619 eprintln!("isb serve: {stack}: {message}");
620 self.inner
621 .emit_kind(Some(kind), level, stack, service, None, message);
622 }
623
624 pub fn relay(
627 &self,
628 kind: Option<&str>,
629 level: &str,
630 stack: &str,
631 service: &str,
632 instance: Option<&str>,
633 message: String,
634 ) {
635 self.inner
636 .emit_kind(kind, level, stack, service, instance, message);
637 }
638
639 pub fn note(&self, level: &str, stack: &str, message: String) {
640 eprintln!("isb serve: {stack}: {message}");
641 self.inner.emit(level, stack, "", None, message);
642 }
643
644 pub fn note_service(&self, level: &str, stack: &str, service: &str, message: String) {
647 self.inner.emit(level, stack, service, None, message);
648 }
649
650 pub fn service_event(&self, level: &str, stack: &str, service: &str, message: String) {
652 eprintln!("isb serve: {stack}/{service}: {message}");
653 self.inner.emit(level, stack, service, None, message);
654 }
655
656 fn notify_stacks(&self) {
657 if let Some(o) = &self.inner.observer {
658 o.stacks_changed(self.definitions());
659 }
660 }
661
662 pub fn plan(&self, def: &StackDef) -> Result<Vec<DeployChange>> {
664 self.validate(def)?;
665 let mut def = def.clone();
666 self.pin_images(&mut def)?;
667 let def = &def;
668 let old = self
669 .inner
670 .stacks
671 .lock()
672 .unwrap()
673 .get(&def.qualified())
674 .cloned();
675 diff(old.as_deref(), def)
676 }
677
678 pub fn client(&self) -> &Client {
679 &self.inner.client
680 }
681
682 pub fn validate(&self, def: &StackDef) -> Result<()> {
685 validate_stack_name(&def.name)?;
686 let host = crate::sandbox::host_facts(&crate::org::client(&self.inner.client, &def.org))?;
688 for (svc, spec) in &def.file.services {
689 let mut s = instance_spec(def, svc, spec, 1, "0000")?;
690 s.name = Some(instance_name(&def.name, svc, 1, "0000")?);
691 crate::plan::resolve(&s, &def.file.volumes, &host, &def.base_dir)?;
692 published(spec)?;
693 crate::ingress::domain::validate(svc, &spec.domains)?;
694 }
695 Ok(())
696 }
697
698 pub fn deploy(&self, mut def: StackDef) -> Result<Vec<DeployChange>> {
701 self.validate(&def)?;
702 self.pin_images(&mut def)?;
703 let _g = self.inner.edit.lock().unwrap();
704 let old = self
705 .inner
706 .stacks
707 .lock()
708 .unwrap()
709 .get(&def.qualified())
710 .cloned();
711 if let Some(old) = &old {
712 let mut prev = (**old).clone();
713 prev.previous = None;
714 def.previous = Some(Box::new(prev));
715 for (k, v) in &old.force {
717 def.force.entry(k.clone()).or_insert(*v);
718 }
719 }
720 let changes = diff(old.as_deref(), &def)?;
721 self.inner.store.save(&def)?;
722 self.apply(Arc::new(def));
723 Ok(changes)
724 }
725
726 fn apply(&self, def: Arc<StackDef>) {
729 let name = def.qualified();
730 self.inner
731 .stacks
732 .lock()
733 .unwrap()
734 .insert(name.clone(), def.clone());
735 let mut workers = self.inner.workers.lock().unwrap();
736 for svc in def.file.services.keys() {
737 let key = (name.clone(), svc.clone());
738 match workers.get(&key) {
739 Some(w) => {
740 let mut slot = w.slot.lock().unwrap();
741 slot.def = def.clone();
742 slot.remove = false;
743 if let Some(st) = self.inner.status.lock().unwrap().get_mut(&key) {
747 if st.state == "failing" {
748 st.state = "updating".into();
749 st.message = None;
750 }
751 }
752 w.wake.notify_all();
753 }
754 None => {
755 let shared = Arc::new(WorkerShared {
756 slot: Mutex::new(Slot {
757 def: def.clone(),
758 remove: false,
759 remove_volumes: false,
760 }),
761 wake: Condvar::new(),
762 stop: AtomicBool::new(false),
763 kick: AtomicBool::new(false),
764 });
765 workers.insert(key, shared.clone());
766 spawn_worker(self.inner.clone(), &def, svc.clone(), shared);
767 }
768 }
769 }
770 for ((stack, svc), w) in workers.iter() {
771 if *stack == name && !def.file.services.contains_key(svc) {
772 let mut slot = w.slot.lock().unwrap();
773 slot.remove = true;
774 w.wake.notify_all();
775 }
776 }
777 drop(workers);
778 self.notify_stacks();
779 }
780
781 pub fn org_limits_changed(&self, org: &OrgId) -> usize {
785 let stacks = self.inner.stacks.lock().unwrap();
786 let workers = self.inner.workers.lock().unwrap();
787 let status = self.inner.status.lock().unwrap();
788 let mut n = 0;
789 for (key, w) in workers.iter() {
790 let ours = stacks.get(&key.0).is_some_and(|d| d.org == *org);
791 let limited = status.get(key).is_some_and(|s| {
792 s.state == "failing" && s.message.as_deref().is_some_and(limit_error)
793 });
794 if ours && limited {
795 w.kick.store(true, Ordering::SeqCst);
796 let _slot = w.slot.lock().unwrap();
799 w.wake.notify_all();
800 n += 1;
801 }
802 }
803 n
804 }
805
806 pub fn remove(&self, name: &str, volumes: bool, timeout: Duration) -> Result<()> {
810 {
811 let _g = self.inner.edit.lock().unwrap();
812 let Some(def) = self.inner.stacks.lock().unwrap().remove(name) else {
813 return Err(Error::NotFound(format!("stack {name}")));
814 };
815 self.inner.store.remove(&def.org, &def.name)?;
816 }
817 self.notify_stacks();
818 let ws: Vec<Arc<WorkerShared>> = self
819 .inner
820 .workers
821 .lock()
822 .unwrap()
823 .iter()
824 .filter(|((s, _), _)| s == name)
825 .map(|(_, w)| w.clone())
826 .collect();
827 for w in &ws {
828 let mut slot = w.slot.lock().unwrap();
829 slot.remove = true;
830 slot.remove_volumes = volumes;
831 w.wake.notify_all();
832 }
833 let started = Instant::now();
834 while started.elapsed() < timeout {
835 let left = self
836 .inner
837 .workers
838 .lock()
839 .unwrap()
840 .keys()
841 .any(|(s, _)| s == name);
842 if !left {
843 return Ok(());
844 }
845 std::thread::sleep(Duration::from_millis(200));
846 }
847 Err(Error::invalid(format!(
848 "stack {name}: still removing after {timeout:?}; it carries on in the background"
849 )))
850 }
851
852 pub fn rollback(&self, name: &str) -> Result<Vec<DeployChange>> {
855 let _g = self.inner.edit.lock().unwrap();
856 let cur = self.get_def(name)?;
857 let prev = cur
858 .previous
859 .clone()
860 .ok_or_else(|| Error::invalid(format!("stack {name} has no previous deployment")))?;
861 let mut def = *prev;
862 def.deployed_at = now_secs();
863 for b in def.secrets.values_mut() {
866 if let Ok(v) = self.inner.secrets.version_in(&b.driver, &def.org, &b.name) {
867 b.version = v;
868 }
869 }
870 let mut cur2 = (*cur).clone();
871 cur2.previous = None;
872 let changes = diff(Some(&cur), &def)?;
873 def.previous = Some(Box::new(cur2));
874 self.inner.store.save(&def)?;
875 self.apply(Arc::new(def));
876 Ok(changes)
877 }
878
879 pub fn scale(&self, name: &str, service: &str, replicas: u32) -> Result<()> {
881 let _g = self.inner.edit.lock().unwrap();
882 let cur = self.get_def(name)?;
883 let mut def = (*cur).clone();
884 let spec = def
885 .file
886 .services
887 .get_mut(service)
888 .ok_or_else(|| Error::NotFound(format!("service {service} in stack {name}")))?;
889 spec.deploy.get_or_insert_with(Default::default).replicas = Some(replicas);
890 self.inner.store.save(&def)?;
891 self.apply(Arc::new(def));
892 Ok(())
893 }
894
895 pub fn redeploy(&self, name: &str, service: &str) -> Result<()> {
899 let _g = self.inner.edit.lock().unwrap();
900 let cur = self.get_def(name)?;
901 cur.service(service)?;
902 let mut def = (*cur).clone();
903 *def.force.entry(service.to_string()).or_insert(0) += 1;
904 self.pin_images(&mut def)?;
906 self.inner.store.save(&def)?;
907 self.apply(Arc::new(def));
908 Ok(())
909 }
910
911 fn pin_images(&self, def: &mut StackDef) -> Result<()> {
915 def.images.clear();
916 for (svc, spec) in &def.file.services {
917 let Some(r) = spec.image.strip_prefix("registry:") else {
918 continue;
919 };
920 let r = crate::registry::ImageRef::parse(r)?;
921 let reg = crate::registry::Registry::shared(&self.inner.client)?;
922 let d = reg.resolve(&def.org, &r)?;
925 if r.digest.is_none() {
926 def.images.insert(svc.clone(), d);
927 }
928 }
929 Ok(())
930 }
931
932 fn get_def(&self, name: &str) -> Result<Arc<StackDef>> {
933 self.inner
934 .stacks
935 .lock()
936 .unwrap()
937 .get(name)
938 .cloned()
939 .ok_or_else(|| Error::NotFound(format!("stack {name}")))
940 }
941
942 pub fn definitions(&self) -> Vec<Arc<StackDef>> {
944 self.inner
945 .stacks
946 .lock()
947 .unwrap()
948 .values()
949 .cloned()
950 .collect()
951 }
952
953 pub fn definition(&self, name: &str) -> Result<StackDef> {
955 self.get_def(name).map(|d| (*d).clone())
956 }
957
958 pub fn list(&self) -> Vec<StackStatus> {
959 let names: Vec<String> = self.inner.stacks.lock().unwrap().keys().cloned().collect();
960 names.iter().filter_map(|n| self.status(n).ok()).collect()
961 }
962
963 pub fn status(&self, name: &str) -> Result<StackStatus> {
964 let def = self.get_def(name)?;
965 let st = self.inner.status.lock().unwrap();
966 let services: Vec<ServiceStatus> = def
967 .file
968 .services
969 .keys()
970 .map(|svc| {
971 st.get(&(name.to_string(), svc.clone()))
972 .cloned()
973 .unwrap_or_else(|| ServiceStatus {
974 service: svc.clone(),
975 state: "starting".into(),
976 ..Default::default()
977 })
978 })
979 .collect();
980 drop(st);
981 let mut services = services;
982 if let Some(o) = &self.inner.observer {
983 for s in &mut services {
984 s.domains = o.domains(name, &s.service);
985 }
986 }
987 let converged = services.iter().all(|s| s.state == "converged");
988 Ok(StackStatus {
989 name: def.name.clone(),
990 org: def.org.to_string(),
991 deployed_at: def.deployed_at,
992 deployed_by: def.deployed_by.clone(),
993 has_previous: def.previous.is_some(),
994 converged,
995 services,
996 })
997 }
998
999 pub fn logs(
1001 &self,
1002 name: &str,
1003 service: &str,
1004 slot: Option<u32>,
1005 lines: usize,
1006 ) -> Result<BTreeMap<String, String>> {
1007 let def = self.get_def(name)?;
1008 let oci = crate::plan::ImageSource::parse(&def.service(service)?.image)?.is_oci();
1009 let oc = crate::org::client(&self.inner.client, &def.org);
1010 super::failure::replica_logs(&oc, &def.name, service, oci, slot, lines)
1011 }
1012
1013 pub fn last_failure(&self, name: &str, service: &str) -> Option<super::failure::FailedAttempt> {
1017 let state = self.inner.status.lock().unwrap();
1018 let converged = state
1019 .get(&(name.to_string(), service.to_string()))
1020 .is_some_and(|s| s.state == "converged");
1021 drop(state);
1022 if converged {
1023 return None;
1024 }
1025 self.inner.failures.last(name, service)
1026 }
1027
1028 pub fn shutdown(&self) {
1031 for w in self.inner.workers.lock().unwrap().values() {
1032 w.stop.store(true, Ordering::SeqCst);
1033 w.wake.notify_all();
1034 }
1035 self.inner.balancer.clear();
1036 }
1037}
1038
1039fn limit_error(msg: &str) -> bool {
1042 msg.contains(" quota (") || msg.contains(" limit (")
1043}
1044
1045fn instance_spec(
1049 def: &StackDef,
1050 service: &str,
1051 spec: &SandboxSpec,
1052 slot: u32,
1053 rev: &str,
1054) -> Result<SandboxSpec> {
1055 let mut s = spec.clone();
1056 s.image = def.instance_image(service, &spec.image);
1057 s.restart = Some(RestartMode::Always);
1058 s.ports.retain(|p| p.bind == PortBind::Guest);
1059 s.domains.clear();
1060 if let Some(d) = &s.deploy {
1061 s.labels.extend(d.labels.clone());
1062 }
1063 s.labels.insert(LABEL_STACK.into(), def.name.clone());
1064 s.labels.insert(LABEL_SERVICE.into(), service.into());
1065 s.labels.insert(LABEL_SLOT.into(), slot.to_string());
1066 s.labels.insert(LABEL_REV.into(), rev.into());
1067 Ok(s)
1068}
1069
1070#[derive(Debug, Clone, PartialEq)]
1073struct Published {
1074 listen: SocketAddr,
1075 target: u16,
1076}
1077
1078fn published(spec: &SandboxSpec) -> Result<Vec<Published>> {
1079 let mut out = Vec::new();
1080 for p in &spec.ports {
1081 if p.bind == PortBind::Guest {
1082 continue;
1083 }
1084 let listen = crate::plan::normalize_addr(&p.listen, "127.0.0.1").map_err(Error::invalid)?;
1085 let connect =
1086 crate::plan::normalize_addr(&p.connect, "127.0.0.1").map_err(Error::invalid)?;
1087 let (lp, lh, lport) = split_addr(&listen).ok_or_else(|| {
1088 Error::invalid(format!(
1089 "port {listen}: a stack publishes single tcp ports (no ranges)"
1090 ))
1091 })?;
1092 let (_, _, cport) = split_addr(&connect).ok_or_else(|| {
1093 Error::invalid(format!(
1094 "port {connect}: a stack publishes single tcp ports (no ranges)"
1095 ))
1096 })?;
1097 if lp != "tcp" {
1098 return Err(Error::invalid(format!(
1099 "port {listen}: the stack balancer is tcp only"
1100 )));
1101 }
1102 if p.search.is_some() {
1103 return Err(Error::invalid(
1104 "a stack's published ports are fixed; port search is for isb up",
1105 ));
1106 }
1107 let host: IpAddr = lh
1108 .trim_start_matches('[')
1109 .trim_end_matches(']')
1110 .parse()
1111 .map_err(|_| {
1112 Error::invalid(format!("port {listen}: the host must be an IP address"))
1113 })?;
1114 out.push(Published {
1115 listen: SocketAddr::new(host, lport),
1116 target: cport,
1117 });
1118 }
1119 Ok(out)
1120}
1121
1122#[derive(Debug, Clone)]
1124#[doc(hidden)]
1125pub struct Inst {
1126 pub name: String,
1127 pub(super) slot: u32,
1128 pub rev: String,
1129 status: String,
1130}
1131
1132impl Inst {
1133 #[doc(hidden)]
1134 pub fn is_running(&self) -> bool {
1135 self.running()
1136 }
1137
1138 fn running(&self) -> bool {
1139 self.status.eq_ignore_ascii_case("running")
1140 }
1141}
1142
1143#[doc(hidden)]
1145pub fn list_instances(client: &Client, stack: &str, service: Option<&str>) -> Result<Vec<Inst>> {
1146 let mut filter = format!("config.user.{LABEL_STACK} eq {stack}");
1147 if let Some(s) = service {
1148 filter.push_str(&format!(" and config.user.{LABEL_SERVICE} eq {s}"));
1149 }
1150 let v = client.get(&format!(
1151 "/1.0/instances?recursion=1&filter={}",
1152 encode_query(&filter)
1153 ))?;
1154 let mut out = Vec::new();
1155 for i in v.as_array().into_iter().flatten() {
1156 let info = crate::sandbox::SandboxInfo::from_api(i);
1157 let c = &info.config;
1158 if c.get(&format!("user.{LABEL_STACK}")).map(String::as_str) != Some(stack) {
1160 continue;
1161 }
1162 let svc = c
1163 .get(&format!("user.{LABEL_SERVICE}"))
1164 .cloned()
1165 .unwrap_or_default();
1166 if service.is_some_and(|s| s != svc) {
1167 continue;
1168 }
1169 out.push(Inst {
1170 name: info.name.clone(),
1171 slot: c
1172 .get(&format!("user.{LABEL_SLOT}"))
1173 .and_then(|s| s.parse().ok())
1174 .unwrap_or(0),
1175 rev: c
1176 .get(&format!("user.{LABEL_REV}"))
1177 .cloned()
1178 .unwrap_or_default(),
1179 status: info.status.clone(),
1180 });
1181 }
1182 out.sort_by(|a, b| (a.slot, &a.name).cmp(&(b.slot, &b.name)));
1183 Ok(out)
1184}
1185
1186fn instance_state(client: &Client, name: &str) -> Result<(i64, Option<IpAddr>)> {
1189 let v = client.get(&format!("/1.0/instances/{}/state", encode_segment(name)))?;
1190 let pid = v.get("pid").and_then(Value::as_i64).unwrap_or(0);
1191 let mut v4 = None;
1192 let mut v6 = None;
1193 if let Some(nets) = v.get("network").and_then(Value::as_object) {
1194 for (ifname, n) in nets {
1195 if ifname == "lo" {
1196 continue;
1197 }
1198 for a in n
1199 .get("addresses")
1200 .and_then(Value::as_array)
1201 .into_iter()
1202 .flatten()
1203 {
1204 if a.get("scope").and_then(Value::as_str) != Some("global") {
1205 continue;
1206 }
1207 let Some(ip) = a
1208 .get("address")
1209 .and_then(Value::as_str)
1210 .and_then(|s| s.parse::<IpAddr>().ok())
1211 else {
1212 continue;
1213 };
1214 match ip {
1215 IpAddr::V4(_) if v4.is_none() => v4 = Some(ip),
1216 IpAddr::V6(_) if v6.is_none() => v6 = Some(ip),
1217 _ => {}
1218 }
1219 }
1220 }
1221 }
1222 Ok((pid, v4.or(v6)))
1223}
1224
1225#[derive(Debug, Default)]
1227struct InstRt {
1228 pid: i64,
1230 since: Option<Instant>,
1232 ip: Option<IpAddr>,
1233 failures: u32,
1234 healthy: Option<bool>,
1235 next_probe: Option<Instant>,
1236 last_probe: String,
1237 unhealthy_restarts: u32,
1239 restarts: VecDeque<Instant>,
1241 in_rotation: bool,
1242}
1243
1244struct Worker {
1246 inner: Arc<Inner>,
1247 stack: String,
1249 q: String,
1251 oclient: Client,
1253 service: String,
1254 shared: Arc<WorkerShared>,
1255 rt: BTreeMap<String, InstRt>,
1256 paused: Option<(String, String)>,
1259 create_backoff: Option<(Instant, Duration)>,
1261 routes: BTreeMap<String, Published>,
1262 route_errors: BTreeMap<String, String>,
1263 template: Option<(String, Desired)>,
1265 state: String,
1266 message: Option<String>,
1267 last_error: Option<String>,
1268 insts: Vec<Inst>,
1270 rollout: Option<RolloutStatus>,
1271 rates: BTreeMap<String, (u64, Instant, VecDeque<f32>)>,
1273 deps_met: bool,
1276 org: OrgId,
1277 dns_last: Option<Vec<IpAddr>>,
1279 dns_error: Option<String>,
1280 observed: Option<Vec<IpAddr>>,
1282 ever_healthy: bool,
1284 health_down: Option<Instant>,
1286 health_alarm: bool,
1288 seen: Option<Arc<StackDef>>,
1290}
1291
1292fn spawn_worker(inner: Arc<Inner>, def: &StackDef, service: String, shared: Arc<WorkerShared>) {
1293 let (stack, q, org) = (def.name.clone(), def.qualified(), def.org.clone());
1294 let oclient = crate::org::client(&inner.client, &def.org);
1295 let name = format!("isb-{q}-{service}");
1296 let r = std::thread::Builder::new().name(name).spawn(move || {
1297 let mut w = Worker::new(inner, stack, q, oclient, service, shared, org);
1298 w.run();
1299 });
1300 if let Err(e) = r {
1301 eprintln!("isb serve: cannot start a worker thread: {e}");
1302 }
1303}
1304
1305impl Worker {
1306 fn new(
1307 inner: Arc<Inner>,
1308 stack: String,
1309 q: String,
1310 oclient: Client,
1311 service: String,
1312 shared: Arc<WorkerShared>,
1313 org: OrgId,
1314 ) -> Worker {
1315 Worker {
1316 inner,
1317 stack,
1318 q,
1319 oclient,
1320 service,
1321 shared,
1322 rt: BTreeMap::new(),
1323 paused: None,
1324 create_backoff: None,
1325 routes: BTreeMap::new(),
1326 route_errors: BTreeMap::new(),
1327 template: None,
1328 state: "starting".into(),
1329 message: None,
1330 last_error: None,
1331 insts: Vec::new(),
1332 rollout: None,
1333 rates: BTreeMap::new(),
1334 deps_met: false,
1335 org,
1336 dns_last: None,
1337 dns_error: None,
1338 observed: None,
1339 ever_healthy: false,
1340 health_down: None,
1341 health_alarm: false,
1342 seen: None,
1343 }
1344 }
1345
1346 fn begin_pass(&mut self, def: &Arc<StackDef>) {
1351 let new_def = self.seen.as_ref().is_none_or(|d| !Arc::ptr_eq(d, def));
1352 let kicked = self.shared.kick.swap(false, Ordering::SeqCst);
1353 if new_def {
1354 self.seen = Some(def.clone());
1355 self.last_error = None;
1356 if self.state == "failing" {
1357 self.state = "updating".into();
1358 self.message = None;
1359 }
1360 }
1361 if new_def || kicked {
1362 self.create_backoff = None;
1363 }
1364 }
1365}
1366
1367impl Worker {
1368 fn log(&self, msg: &str) {
1369 self.event("info", None, msg);
1370 }
1371
1372 fn event(&self, level: &str, instance: Option<&str>, msg: &str) {
1373 eprintln!("isb serve: {}/{}: {msg}", self.q, self.service);
1374 self.inner
1375 .emit(level, &self.q, &self.service, instance, msg.to_string());
1376 }
1377
1378 fn client(&self) -> &Client {
1379 &self.oclient
1380 }
1381
1382 fn watch_health(&mut self, healthy: u32, replicas: u32) {
1386 if healthy > 0 {
1387 self.ever_healthy = true;
1388 self.health_down = None;
1389 if self.health_alarm {
1390 self.health_alarm = false;
1391 let msg = format!("{healthy} of {replicas} replicas healthy again");
1392 eprintln!("isb serve: {}/{}: {msg}", self.q, self.service);
1393 self.inner.emit_kind(
1394 Some("health.recovered"),
1395 "info",
1396 &self.q,
1397 &self.service,
1398 None,
1399 msg,
1400 );
1401 }
1402 return;
1403 }
1404 if replicas == 0 || !self.ever_healthy || self.rollout.is_some() {
1405 self.health_down = None;
1406 return;
1407 }
1408 let since = *self.health_down.get_or_insert_with(Instant::now);
1409 if !self.health_alarm && since.elapsed() >= HEALTH_DEBOUNCE {
1410 self.health_alarm = true;
1411 let msg = format!(
1412 "no healthy replica (of {replicas}) for {}s",
1413 since.elapsed().as_secs()
1414 );
1415 eprintln!("isb serve: {}/{}: {msg}", self.q, self.service);
1416 self.inner.emit_kind(
1417 Some("health.unhealthy"),
1418 "error",
1419 &self.q,
1420 &self.service,
1421 None,
1422 msg,
1423 );
1424 }
1425 }
1426
1427 fn key(&self) -> (String, String) {
1428 (self.q.clone(), self.service.clone())
1429 }
1430
1431 fn run(&mut self) {
1432 loop {
1433 if self.shared.stop.load(Ordering::SeqCst) {
1434 return;
1435 }
1436 let (def, remove, remove_volumes) = {
1437 let s = self.shared.slot.lock().unwrap();
1438 (s.def.clone(), s.remove, s.remove_volumes)
1439 };
1440 if remove {
1441 self.teardown(&def, remove_volumes);
1442 let mut ws = self.inner.workers.lock().unwrap();
1445 let slot = self.shared.slot.lock().unwrap();
1446 if slot.remove || self.shared.stop.load(Ordering::SeqCst) {
1447 ws.remove(&self.key());
1448 self.inner.status.lock().unwrap().remove(&self.key());
1449 return;
1450 }
1451 continue;
1452 }
1453 self.begin_pass(&def);
1454 if let Err(e) = self.pass(&def) {
1455 self.state = "failing".into();
1456 self.message = Some(e.to_string());
1457 if self.last_error.as_deref() != Some(&e.to_string()) {
1459 self.event("error", None, &e.to_string());
1460 self.last_error = Some(e.to_string());
1461 }
1462 self.publish_status(&def);
1463 let slot = self.shared.slot.lock().unwrap();
1464 if Arc::ptr_eq(&slot.def, &def) && !slot.remove {
1465 let _ = self.shared.wake.wait_timeout(slot, self.inner.interval);
1466 }
1467 continue;
1468 }
1469 self.last_error = None;
1470 let slot = self.shared.slot.lock().unwrap();
1471 if Arc::ptr_eq(&slot.def, &def) && !slot.remove {
1472 let _ = self.shared.wake.wait_timeout(slot, self.inner.interval);
1473 }
1474 }
1475 }
1476
1477 fn superseded(&self, def: &Arc<StackDef>) -> bool {
1479 let s = self.shared.slot.lock().unwrap();
1480 s.remove || !Arc::ptr_eq(&s.def, def) || self.shared.stop.load(Ordering::SeqCst)
1481 }
1482
1483 fn teardown(&mut self, def: &StackDef, volumes: bool) {
1485 if let Some(dir) = Some(crate::discovery::org_dir(&self.org)).filter(|d| d.is_dir()) {
1487 if let Err(e) =
1488 crate::discovery::publish(&dir, &self.org, &self.stack, &self.service, &[])
1489 {
1490 self.log(&format!("cannot remove the service name: {e}"));
1491 }
1492 }
1493 self.dns_last = Some(Vec::new());
1494 if let Some(o) = &self.inner.observer {
1495 o.rotation(&self.q, &self.service, &[]);
1496 }
1497 self.observed = Some(Vec::new());
1498 for (k, _) in std::mem::take(&mut self.routes) {
1499 self.inner.balancer.remove_route(&k);
1500 }
1501 match list_instances(self.client(), &self.stack, Some(&self.service)) {
1502 Ok(insts) => {
1503 for i in insts {
1504 self.log(&format!("removing {}", i.name));
1505 if let Err(e) = Sandbox::remove(self.client(), &i.name, true) {
1506 if !e.is_not_found() {
1507 self.log(&format!("cannot remove {}: {e}", i.name));
1508 }
1509 }
1510 }
1511 }
1512 Err(e) => self.log(&format!("cannot list instances to remove: {e}")),
1513 }
1514 if volumes {
1515 self.remove_volumes(def);
1516 }
1517 self.rt.clear();
1518 self.template = None;
1519 }
1520
1521 fn remove_volumes(&self, def: &StackDef) {
1522 let Ok(spec) = def.service(&self.service) else {
1523 return;
1524 };
1525 let Ok(host) = crate::sandbox::host_facts(self.client()) else {
1526 return;
1527 };
1528 let Ok(pool) = host.pick_pool(spec.storage.as_deref()) else {
1529 return;
1530 };
1531 for v in &spec.volumes {
1532 if v.mount_type != crate::spec::MountType::Volume {
1533 continue;
1534 }
1535 let d = def.file.volumes.get(&v.source);
1536 if v.external || d.is_some_and(|d| d.external) {
1537 continue;
1538 }
1539 let name = d
1540 .and_then(|d| d.name.clone())
1541 .unwrap_or_else(|| v.source.clone());
1542 let vpool = match v.pool.as_deref().or(d.and_then(|d| d.pool.as_deref())) {
1543 Some(p) if p != "auto" => p.to_string(),
1544 _ => pool.clone(),
1545 };
1546 match crate::volume::remove(self.client(), &vpool, &name) {
1547 Ok(()) => self.log(&format!("volume {name}: deleted")),
1548 Err(e) if e.is_not_found() => {}
1549 Err(e) => self.log(&format!("volume {name}: kept ({e})")),
1551 }
1552 }
1553 }
1554
1555 #[expect(
1557 clippy::too_many_lines,
1558 reason = "predates the lint ratchet; split it when next changed"
1559 )]
1560 fn pass(&mut self, def: &Arc<StackDef>) -> Result<()> {
1561 let spec = def.service(&self.service)?.clone();
1562 let rev = def.revision(&self.service)?;
1563 let replicas = spec.replicas();
1564 let oci = crate::plan::ImageSource::parse(&spec.image)?.is_oci();
1565 let probe = spec.health_probe().map_err(Error::invalid)?;
1566
1567 if !self.deps_met {
1568 if let Some(msg) = self.waiting_for(&spec) {
1569 self.state = "waiting".into();
1570 self.message = Some(msg);
1571 self.publish_status(def);
1572 return Ok(());
1573 }
1574 self.deps_met = true;
1575 }
1576 if self.template.as_ref().is_none_or(|(r, _)| *r != rev) {
1577 let mut s = instance_spec(def, &self.service, &spec, 1, &rev)?;
1578 s.name = Some(instance_name(&self.stack, &self.service, 1, "0000")?);
1579 let d = crate::sandbox::resolve(self.client(), &s, &def.file.volumes, &def.base_dir)?;
1580 self.template = Some((rev.clone(), d));
1581 }
1582 self.set_routes(&spec);
1583
1584 let mut insts = list_instances(self.client(), &self.stack, Some(&self.service))?;
1585 self.rt.retain(|n, _| insts.iter().any(|i| i.name == *n));
1586
1587 let extra: Vec<Inst> = insts
1589 .iter()
1590 .filter(|i| i.slot > replicas || i.slot == 0)
1591 .cloned()
1592 .collect();
1593 for i in extra.iter().rev() {
1594 self.log(&format!("scaling down: removing {}", i.name));
1595 self.retire(&i.name)?;
1596 }
1597 insts.retain(|i| i.slot >= 1 && i.slot <= replicas);
1598 self.insts = insts.clone();
1599
1600 for i in &insts {
1602 self.maintain(def, i, &spec, oci, probe.as_ref())?;
1603 }
1604 for slot in 1..=replicas {
1607 let current_ok = insts.iter().any(|i| {
1608 i.slot == slot
1609 && i.rev == rev
1610 && self.rt.get(&i.name).is_some_and(|r| r.in_rotation)
1611 });
1612 if current_ok {
1613 for i in insts.iter().filter(|i| i.slot == slot && i.rev != rev) {
1614 self.log(&format!("removing leftover {}", i.name));
1615 self.retire(&i.name)?;
1616 }
1617 }
1618 let mut current: Vec<&Inst> = insts
1619 .iter()
1620 .filter(|i| i.slot == slot && i.rev == rev)
1621 .collect();
1622 while current.len() > 1 {
1624 let i = current.pop().unwrap();
1625 self.log(&format!("removing duplicate {}", i.name));
1626 self.retire(&i.name)?;
1627 }
1628 }
1629 self.sync_routes();
1630 self.publish_status(def);
1631
1632 let mut pending: Vec<(u32, Option<String>)> = Vec::new();
1634 for slot in 1..=replicas {
1635 if insts.iter().any(|i| i.slot == slot && i.rev == rev) {
1636 continue;
1637 }
1638 let old = insts
1639 .iter()
1640 .find(|i| i.slot == slot)
1641 .map(|i| i.name.clone());
1642 pending.push((slot, old));
1643 }
1644 if pending.is_empty() {
1645 self.create_backoff = None;
1646 let all_ok = insts
1647 .iter()
1648 .all(|i| i.rev == rev && self.rt.get(&i.name).is_some_and(|r| r.in_rotation));
1649 self.state = if all_ok { "converged" } else { "failing" }.into();
1650 if all_ok {
1651 self.message = None;
1652 } else if self.message.is_none() {
1653 self.message = Some("some replicas are not healthy".into());
1654 }
1655 self.publish_status(def);
1656 return Ok(());
1657 }
1658 if self.paused.as_ref().is_some_and(|(r, _)| *r == rev) {
1659 self.state = "paused".into();
1660 self.message = self.paused.as_ref().map(|(_, m)| m.clone());
1661 self.publish_status(def);
1662 return Ok(());
1663 }
1664 if let Some((at, wait)) = self.create_backoff {
1665 if at.elapsed() < wait {
1666 return Ok(());
1667 }
1668 }
1669 self.state = "updating".into();
1670 self.message = None;
1671 self.publish_status(def);
1672
1673 let uc: UpdateConfig = spec
1674 .deploy
1675 .as_ref()
1676 .and_then(|d| d.update_config.clone())
1677 .unwrap_or_default();
1678 let parallel = match uc.parallelism.unwrap_or(1) {
1679 0 => pending.len(),
1680 n => n as usize,
1681 };
1682 let delay = uc
1683 .delay
1684 .as_deref()
1685 .map(crate::flex::parse_duration)
1686 .transpose()
1687 .map_err(Error::invalid)?
1688 .unwrap_or_default();
1689 let monitor = uc
1690 .monitor
1691 .as_deref()
1692 .map(crate::flex::parse_duration)
1693 .transpose()
1694 .map_err(Error::invalid)?
1695 .unwrap_or(Duration::from_secs(5));
1696 let order = uc.order.unwrap_or_default();
1697 let order_name = match order {
1698 UpdateOrder::StopFirst => "stop-first",
1699 UpdateOrder::StartFirst => "start-first",
1700 };
1701 self.rollout = Some(RolloutStatus {
1702 to_rev: rev.clone(),
1703 order: order_name.into(),
1704 parallelism: parallel,
1705 done: 0,
1706 total: pending.len(),
1707 started_at: now_secs(),
1708 slots: pending
1709 .iter()
1710 .map(|(slot, old)| SlotRollout {
1711 slot: *slot,
1712 old: old.clone(),
1713 old_rev: old
1714 .as_ref()
1715 .and_then(|o| insts.iter().find(|i| i.name == *o))
1716 .map(|i| i.rev.clone()),
1717 old_state: if old.is_some() { "serving" } else { "none" }.into(),
1718 new: None,
1719 new_state: "waiting".into(),
1720 })
1721 .collect(),
1722 });
1723 let rollout_started = Instant::now();
1724 self.log(&format!(
1725 "rolling out rev {rev} to {} slot(s), {order_name}",
1726 pending.len()
1727 ));
1728 self.publish_status(def);
1729 let r = self.roll(
1730 def,
1731 &spec,
1732 &rev,
1733 &pending,
1734 parallel,
1735 delay,
1736 monitor,
1737 order,
1738 oci,
1739 probe.as_ref(),
1740 &uc,
1741 );
1742 let rollout = self.rollout.take();
1743 if let (Ok(true), Some(ro)) = (&r, rollout) {
1744 self.log(&format!(
1745 "rollout of rev {rev} complete: {}/{} slot(s) in {:.0?}",
1746 ro.done,
1747 ro.total,
1748 rollout_started.elapsed()
1749 ));
1750 }
1751 self.publish_status(def);
1752 r.map(|_| ())
1753 }
1754
1755 #[expect(clippy::too_many_arguments)]
1758 #[expect(
1759 clippy::excessive_nesting,
1760 reason = "predates the lint ratchet; split it when next changed"
1761 )]
1762 fn roll(
1763 &mut self,
1764 def: &Arc<StackDef>,
1765 spec: &SandboxSpec,
1766 rev: &String,
1767 pending: &[(u32, Option<String>)],
1768 parallel: usize,
1769 delay: Duration,
1770 monitor: Duration,
1771 order: UpdateOrder,
1772 oci: bool,
1773 probe: Option<&HealthProbe>,
1774 uc: &UpdateConfig,
1775 ) -> Result<bool> {
1776 for (n, batch) in pending.chunks(parallel).enumerate() {
1777 if n > 0 && !delay.is_zero() {
1778 std::thread::sleep(delay);
1779 }
1780 if self.superseded(def) {
1781 return Ok(false);
1782 }
1783 for (slot, old) in batch {
1784 let r = self.replace(
1785 def,
1786 spec,
1787 rev,
1788 *slot,
1789 old.as_deref(),
1790 order,
1791 oci,
1792 probe,
1793 monitor,
1794 );
1795 if let Err(e) = r {
1796 let msg = format!("slot {slot}: {e}");
1797 self.event("error", None, &msg);
1798 if old.is_none() {
1799 let wait = self
1801 .create_backoff
1802 .map(|(_, w)| (w * 2).min(Duration::from_secs(300)))
1803 .unwrap_or(Duration::from_secs(10));
1804 self.create_backoff = Some((Instant::now(), wait));
1805 self.state = "failing".into();
1806 self.message = Some(format!("{msg}; retrying in {wait:?}"));
1807 self.publish_status(def);
1808 return Ok(false);
1809 }
1810 match uc.failure_action.unwrap_or_default() {
1811 FailureAction::Continue => continue,
1812 FailureAction::Pause => {
1813 self.event("warn", None, &format!("rollout of rev {rev} paused"));
1814 self.paused = Some((rev.clone(), format!("rollout paused: {msg}")));
1815 return Ok(false);
1816 }
1817 FailureAction::Rollback => {
1818 self.paused = Some((rev.clone(), format!("rolled back: {msg}")));
1819 let ctl = Controller {
1820 inner: self.inner.clone(),
1821 };
1822 self.event(
1823 "warn",
1824 None,
1825 &format!("rollout of rev {rev} failed; rolling back"),
1826 );
1827 if let Err(e) = ctl.rollback(&self.q) {
1828 self.event("error", None, &format!("rollback failed: {e}"));
1829 }
1830 return Ok(false);
1831 }
1832 }
1833 }
1834 }
1835 }
1836 Ok(true)
1837 }
1838
1839 fn slot_state(
1841 &mut self,
1842 def: &StackDef,
1843 slot: u32,
1844 old_state: Option<&str>,
1845 new: Option<&str>,
1846 new_state: Option<&str>,
1847 ) {
1848 if let Some(ro) = &mut self.rollout {
1849 if let Some(s) = ro.slots.iter_mut().find(|s| s.slot == slot) {
1850 if let Some(o) = old_state {
1851 s.old_state = o.into();
1852 }
1853 if let Some(n) = new {
1854 s.new = Some(n.into());
1855 }
1856 if let Some(n) = new_state {
1857 s.new_state = n.into();
1858 if n == "serving" {
1859 ro.done += 1;
1860 }
1861 }
1862 }
1863 }
1864 self.publish_status(def);
1865 }
1866
1867 fn waiting_for(&self, spec: &SandboxSpec) -> Option<String> {
1869 let st = self.inner.status.lock().unwrap();
1870 for (dep, d) in &spec.depends_on {
1871 let s = st.get(&(self.q.clone(), dep.clone()));
1872 let ok = match d.condition {
1873 DependCondition::ServiceStarted => s.is_some_and(|s| s.running > 0),
1874 DependCondition::ServiceHealthy => s.is_some_and(|s| s.healthy > 0),
1875 };
1876 if !ok {
1877 return Some(format!("waiting for {dep} ({:?})", d.condition));
1878 }
1879 }
1880 None
1881 }
1882
1883 fn handle(&self, name: &str) -> Sandbox {
1884 let (_, d) = self.template.as_ref().expect("resolved in pass");
1885 Sandbox::like(self.client(), name, d)
1886 }
1887
1888 #[expect(
1891 clippy::too_many_lines,
1892 reason = "predates the lint ratchet; split it when next changed"
1893 )]
1894 fn maintain(
1895 &mut self,
1896 def: &StackDef,
1897 i: &Inst,
1898 spec: &SandboxSpec,
1899 oci: bool,
1900 probe: Option<&HealthProbe>,
1901 ) -> Result<()> {
1902 let policy = spec.deploy.as_ref().and_then(|d| d.restart_policy.clone());
1903 let condition = policy
1904 .as_ref()
1905 .and_then(|p| p.condition)
1906 .unwrap_or_default();
1907 let sb = self.handle(&i.name);
1908 if !i.running() {
1909 self.set_rotation(&i.name, false);
1910 if condition == RestartCondition::None {
1911 return Ok(());
1912 }
1913 if self.restart_budget_spent(&i.name, policy.as_ref()) {
1914 self.message = Some(format!("{}: restart limit reached", i.name));
1915 return Ok(());
1916 }
1917 self.event(
1918 "warn",
1919 Some(&i.name),
1920 &format!("{} is {}; starting it", i.name, i.status),
1921 );
1922 self.count_restart(&i.name);
1923 if let Err(e) = sb.start() {
1924 self.log(&format!("cannot start {}: {e}", i.name));
1925 return Ok(());
1926 }
1927 }
1928 let (pid, ip) = match instance_state(self.client(), &i.name) {
1929 Ok(s) => s,
1930 Err(e) if e.is_not_found() => return Ok(()),
1931 Err(e) => return Err(e),
1932 };
1933 let rt = self.rt.entry(i.name.clone()).or_default();
1934 rt.ip = ip;
1935 if rt.pid != pid {
1936 rt.pid = pid;
1939 rt.since = Some(Instant::now());
1940 rt.failures = 0;
1941 rt.healthy = None;
1942 rt.next_probe = None;
1943 let r = self.setup(def, &sb, spec, oci);
1944 if let Err(e) = r {
1945 if let Some(rt) = self.rt.get_mut(&i.name) {
1947 rt.pid = 0;
1948 }
1949 self.set_rotation(&i.name, false);
1950 return Err(Error::OperationFailed {
1951 step: format!("set up {}", i.name),
1952 message: e.to_string(),
1953 });
1954 }
1955 }
1956 let alive = self.alive(&sb, spec, oci);
1957 let healthy = match probe {
1958 None => alive,
1959 Some(p) => {
1960 let rt = self.rt.get_mut(&i.name).unwrap();
1961 let since = rt.since.unwrap_or_else(Instant::now);
1962 let in_start = since.elapsed() < p.start_period;
1963 if rt.next_probe.is_none_or(|t| Instant::now() >= t) {
1964 let r = supervise::probe(&sb, p);
1965 let rt = self.rt.get_mut(&i.name).unwrap();
1966 rt.last_probe = r.output.clone();
1967 if r.ok {
1968 rt.failures = 0;
1969 rt.healthy = Some(true);
1970 rt.unhealthy_restarts = 0;
1971 } else if !in_start {
1972 rt.failures += 1;
1973 if rt.failures >= p.retries {
1974 rt.healthy = Some(false);
1975 }
1976 }
1977 let wait = if rt.healthy.is_none() {
1978 p.start_interval
1979 } else {
1980 p.interval
1981 };
1982 rt.next_probe = Some(Instant::now() + wait);
1983 }
1984 self.rt[&i.name].healthy == Some(true) && alive
1985 }
1986 };
1987 self.set_rotation(&i.name, healthy);
1988 let unhealthy = self.rt[&i.name].healthy == Some(false);
1989 if unhealthy && condition != RestartCondition::None {
1990 if self.restart_budget_spent(&i.name, policy.as_ref()) {
1991 self.message = Some(format!("{}: unhealthy, restart limit reached", i.name));
1992 return Ok(());
1993 }
1994 let n = self.rt[&i.name].unhealthy_restarts;
1995 if n >= RESTARTS_BEFORE_REPLACE {
1996 self.log(&format!(
1997 "{} stayed unhealthy through {n} restarts; replacing it",
1998 i.name
1999 ));
2000 let rt = self.rt.get_mut(&i.name).unwrap();
2001 rt.unhealthy_restarts = 0;
2002 drop(sb);
2003 self.retire(&i.name)?;
2004 return Ok(());
2005 }
2006 self.event(
2007 "warn",
2008 Some(&i.name),
2009 &format!(
2010 "{} is unhealthy ({}); restarting its app",
2011 i.name,
2012 self.rt[&i.name]
2013 .last_probe
2014 .lines()
2015 .last()
2016 .unwrap_or("probe failed")
2017 ),
2018 );
2019 self.count_restart(&i.name);
2020 let rt = self.rt.get_mut(&i.name).unwrap();
2021 rt.unhealthy_restarts += 1;
2022 rt.failures = 0;
2023 rt.healthy = None;
2024 rt.since = Some(Instant::now());
2025 supervise::restart_app(&sb, &self.service, oci)?;
2026 }
2027 Ok(())
2028 }
2029
2030 fn setup(&self, def: &StackDef, sb: &Sandbox, spec: &SandboxSpec, oci: bool) -> Result<()> {
2032 let keys = spec.secret_keys();
2034 let values = if keys.is_empty() {
2035 BTreeMap::new()
2036 } else {
2037 super::secrets::values(&self.inner.secrets, &def.org, &def.secrets, keys)?
2038 };
2039 if supervise::push_secrets(sb, spec, &values)? && oci {
2042 supervise::restart_app(sb, &self.service, oci)?;
2043 }
2044 if spec.command.is_some() && !oci {
2045 let mut s = spec.clone();
2046 s.restart = Some(RestartMode::Always);
2047 let env = supervise::secret_env(spec, &values)?;
2048 supervise::install(sb, &self.service, &s, !spec.secrets.is_empty(), &env)?;
2049 }
2050 Ok(())
2051 }
2052
2053 fn oci_secret_env(
2055 &self,
2056 def: &StackDef,
2057 spec: &SandboxSpec,
2058 ) -> Result<BTreeMap<String, String>> {
2059 if spec.env.secrets.is_empty() {
2060 return Ok(BTreeMap::new());
2061 }
2062 let values = super::secrets::values(
2063 &self.inner.secrets,
2064 &def.org,
2065 &def.secrets,
2066 spec.env.secrets.values().map(String::as_str),
2067 )?;
2068 supervise::secret_env(spec, &values)
2069 }
2070
2071 fn alive(&self, sb: &Sandbox, spec: &SandboxSpec, oci: bool) -> bool {
2074 if spec.command.is_some() && !oci {
2075 return supervise::unit_state(sb, &self.service).is_ok_and(|s| s == "active");
2076 }
2077 sb.info()
2078 .is_ok_and(|i| i.status.eq_ignore_ascii_case("running"))
2079 }
2080
2081 fn count_restart(&mut self, name: &str) {
2082 self.rt
2083 .entry(name.to_string())
2084 .or_default()
2085 .restarts
2086 .push_back(Instant::now());
2087 }
2088
2089 fn restart_budget_spent(
2090 &mut self,
2091 name: &str,
2092 policy: Option<&crate::spec::RestartPolicy>,
2093 ) -> bool {
2094 let Some(p) = policy else { return false };
2095 let Some(max) = p.max_attempts else {
2096 return false;
2097 };
2098 let window = p
2099 .window
2100 .as_deref()
2101 .and_then(|w| crate::flex::parse_duration(w).ok());
2102 let rt = self.rt.entry(name.to_string()).or_default();
2103 if let Some(w) = window {
2104 while rt.restarts.front().is_some_and(|t| t.elapsed() > w) {
2105 rt.restarts.pop_front();
2106 }
2107 }
2108 rt.restarts.len() as u32 >= max
2109 }
2110
2111 #[expect(clippy::too_many_arguments)]
2113 fn replace(
2114 &mut self,
2115 def: &StackDef,
2116 spec: &SandboxSpec,
2117 rev: &str,
2118 slot: u32,
2119 old: Option<&str>,
2120 order: UpdateOrder,
2121 oci: bool,
2122 probe: Option<&HealthProbe>,
2123 monitor: Duration,
2124 ) -> Result<()> {
2125 let secret_env = if oci {
2128 self.oci_secret_env(def, spec)?
2129 } else {
2130 BTreeMap::new()
2131 };
2132 if let (Some(o), UpdateOrder::StopFirst) = (old, order) {
2133 self.log(&format!("slot {slot}: replacing {o} (stop-first)"));
2134 self.slot_state(def, slot, Some("draining"), None, None);
2135 self.retire(o)?;
2136 self.slot_state(def, slot, Some("retired"), None, None);
2137 }
2138 let name = instance_name(&self.stack, &self.service, slot, &new_id())?;
2139 self.log(&format!("slot {slot}: creating {name} (rev {rev})"));
2140 self.slot_state(def, slot, None, Some(&name), Some("creating"));
2141 let mut s = instance_spec(def, &self.service, spec, slot, rev)?;
2142 s.name = Some(name.clone());
2143 s.env.vars.extend(secret_env);
2145 let d = crate::sandbox::resolve(self.client(), &s, &def.file.volumes, &def.base_dir)?;
2146 let stack = self.q.clone();
2147 let mut report = |m: &str| eprintln!("isb serve: {stack}: {m}");
2148 let created =
2149 crate::sandbox::ensure(self.client(), &d, EnsureOptions::default(), &mut report);
2150 let result = created.and_then(|_| {
2151 let inst = Inst {
2152 name: name.clone(),
2153 slot,
2154 rev: rev.to_string(),
2155 status: "Running".into(),
2156 };
2157 self.insts.push(inst.clone());
2158 self.slot_state(def, slot, None, None, Some("probing"));
2159 self.wait_serving(def, &inst, spec, oci, probe, monitor)
2160 });
2161 if let Err(e) = result {
2162 let (e, attempt) =
2164 super::failure::explain(self.client(), &name, &self.service, oci, e, now_ms());
2165 self.inner.failures.record(&self.q, &self.service, attempt);
2166 self.event(
2167 "error",
2168 Some(&name),
2169 &format!("{name} did not come up: {e}"),
2170 );
2171 self.slot_state(def, slot, None, None, Some("failed"));
2172 let _ = self.retire(&name);
2173 return Err(e);
2174 }
2175 if let (Some(o), UpdateOrder::StartFirst) = (old, order) {
2176 self.log(&format!("slot {slot}: {name} is serving; retiring {o}"));
2177 self.slot_state(def, slot, Some("draining"), None, None);
2178 self.retire(o)?;
2179 self.slot_state(def, slot, Some("retired"), None, None);
2180 }
2181 self.slot_state(def, slot, None, None, Some("serving"));
2182 Ok(())
2183 }
2184
2185 fn wait_serving(
2189 &mut self,
2190 def: &StackDef,
2191 i: &Inst,
2192 spec: &SandboxSpec,
2193 oci: bool,
2194 probe: Option<&HealthProbe>,
2195 monitor: Duration,
2196 ) -> Result<()> {
2197 let deadline = match probe {
2198 Some(p) => p.start_period + p.interval * p.retries + Duration::from_secs(30),
2199 None => Duration::from_secs(60),
2200 }
2201 .max(Duration::from_secs(60));
2202 let started = Instant::now();
2203 loop {
2204 self.maintain(def, i, spec, oci, probe)?;
2205 if self.rt.get(&i.name).is_some_and(|r| r.in_rotation) {
2206 break;
2207 }
2208 if self
2209 .rt
2210 .get(&i.name)
2211 .is_some_and(|r| r.healthy == Some(false))
2212 {
2213 return Err(Error::invalid(format!(
2214 "unhealthy: {}",
2215 self.rt[&i.name].last_probe
2216 )));
2217 }
2218 if started.elapsed() > deadline {
2219 let why = self
2220 .rt
2221 .get(&i.name)
2222 .map(|r| r.last_probe.clone())
2223 .filter(|s| !s.is_empty())
2224 .unwrap_or_else(|| "its app is not running".into());
2225 return Err(Error::invalid(format!(
2226 "not serving after {:?}: {why}",
2227 started.elapsed()
2228 )));
2229 }
2230 if let Some(rt) = self.rt.get_mut(&i.name) {
2232 rt.next_probe = None;
2233 }
2234 std::thread::sleep(
2235 probe
2236 .map(|p| p.start_interval)
2237 .unwrap_or(Duration::from_secs(1))
2238 .min(Duration::from_secs(2)),
2239 );
2240 }
2241 self.sync_routes();
2242 self.slot_state(def, i.slot, None, None, Some("monitoring"));
2243 let watch = Instant::now();
2244 while watch.elapsed() < monitor {
2245 std::thread::sleep(Duration::from_secs(1).min(monitor));
2246 if let Some(rt) = self.rt.get_mut(&i.name) {
2247 rt.next_probe = None;
2248 }
2249 self.maintain(def, i, spec, oci, probe)?;
2250 if !self.rt.get(&i.name).is_some_and(|r| r.in_rotation) {
2251 return Err(Error::invalid(format!(
2252 "failed within the {monitor:?} monitor period"
2253 )));
2254 }
2255 }
2256 Ok(())
2257 }
2258
2259 fn retire(&mut self, name: &str) -> Result<()> {
2262 let ip = self.rt.get(name).and_then(|r| r.ip);
2263 self.set_rotation(name, false);
2264 self.sync_routes();
2265 if let Some(ip) = ip {
2266 let started = Instant::now();
2267 if let Some(o) = &self.inner.observer {
2268 o.drain(&self.q, &self.service, ip, DRAIN);
2269 }
2270 let left = DRAIN.saturating_sub(started.elapsed());
2271 for (k, p) in &self.routes {
2272 self.inner
2273 .balancer
2274 .wait_drained(k, SocketAddr::new(ip, p.target), left);
2275 }
2276 }
2277 if let Ok(sb) = Sandbox::get(self.client(), name) {
2278 let _ = sb.stop(false, Duration::from_secs(10));
2279 }
2280 self.rt.remove(name);
2281 self.insts.retain(|i| i.name != name);
2282 match Sandbox::remove(self.client(), name, true) {
2283 Err(e) if !e.is_not_found() => Err(e),
2284 _ => Ok(()),
2285 }
2286 }
2287
2288 fn set_rotation(&mut self, name: &str, on: bool) {
2289 let rt = self.rt.entry(name.to_string()).or_default();
2290 if rt.in_rotation != on {
2291 rt.in_rotation = on;
2292 self.sync_routes();
2293 }
2294 }
2295
2296 fn set_routes(&mut self, spec: &SandboxSpec) {
2298 let want: BTreeMap<String, Published> = match published(spec) {
2299 Ok(ps) => ps
2300 .into_iter()
2301 .map(|p| (format!("{}/{}/{}", self.q, self.service, p.listen), p))
2302 .collect(),
2303 Err(e) => {
2304 self.message = Some(e.to_string());
2305 BTreeMap::new()
2306 }
2307 };
2308 let stale: Vec<String> = self
2309 .routes
2310 .keys()
2311 .filter(|k| !want.contains_key(*k))
2312 .cloned()
2313 .collect();
2314 for k in stale {
2315 self.inner.balancer.remove_route(&k);
2316 self.routes.remove(&k);
2317 self.route_errors.remove(&k);
2318 }
2319 for (k, p) in want {
2320 self.routes.entry(k).or_insert(p);
2321 }
2322 self.sync_routes();
2323 }
2324
2325 fn sync_routes(&mut self) {
2326 let mut errors = BTreeMap::new();
2327 for (k, p) in &self.routes {
2328 let backends: Vec<SocketAddr> = self
2329 .rt
2330 .values()
2331 .filter(|r| r.in_rotation)
2332 .filter_map(|r| r.ip)
2333 .map(|ip| SocketAddr::new(ip, p.target))
2334 .collect();
2335 if let Err(e) = self.inner.balancer.set_route(k, p.listen, backends) {
2336 errors.insert(k.clone(), e.to_string());
2337 }
2338 }
2339 for (k, e) in &errors {
2340 if self.route_errors.get(k) != Some(e) {
2341 self.log(&format!("cannot publish {k}: {e}"));
2342 }
2343 }
2344 self.route_errors = errors;
2345 self.sync_observer();
2346 self.sync_dns();
2347 }
2348
2349 fn sync_observer(&mut self) {
2351 let Some(o) = self.inner.observer.clone() else {
2352 return;
2353 };
2354 let mut ips: Vec<IpAddr> = self
2355 .rt
2356 .values()
2357 .filter(|r| r.in_rotation)
2358 .filter_map(|r| r.ip)
2359 .collect();
2360 ips.sort();
2361 if self.observed.as_ref() != Some(&ips) {
2362 o.rotation(&self.q, &self.service, &ips);
2363 self.observed = Some(ips);
2364 }
2365 }
2366
2367 fn sync_dns(&mut self) {
2371 let dir = crate::discovery::org_dir(&self.org);
2372 let mut ips: Vec<IpAddr> = self
2373 .rt
2374 .values()
2375 .filter(|r| r.in_rotation)
2376 .filter_map(|r| r.ip)
2377 .collect();
2378 ips.sort();
2379 let fresh = self.dns_last.is_none() && ips.is_empty();
2383 if fresh || self.dns_last.as_ref() == Some(&ips) || !dir.is_dir() {
2384 return;
2385 }
2386 match crate::discovery::publish(&dir, &self.org, &self.stack, &self.service, &ips) {
2387 Ok(()) => {
2388 self.dns_last = Some(ips);
2389 self.dns_error = None;
2390 }
2391 Err(e) => {
2392 let e = e.to_string();
2393 if self.dns_error.as_deref() != Some(&e) {
2394 self.event(
2395 "warn",
2396 None,
2397 &format!("cannot publish the service name: {e}"),
2398 );
2399 self.dns_error = Some(e);
2400 }
2401 }
2402 }
2403 }
2404
2405 #[expect(
2406 clippy::too_many_lines,
2407 reason = "predates the lint ratchet; split it when next changed"
2408 )]
2409 fn publish_status(&mut self, def: &StackDef) {
2410 self.sync_observer();
2412 self.sync_dns();
2413 let Ok(spec) = def.service(&self.service) else {
2414 return;
2415 };
2416 let rev = def.revision(&self.service).unwrap_or_default();
2417 let probe = matches!(spec.health_probe(), Ok(Some(_)));
2418 let snap = self.inner.snapshot.lock().unwrap().instances.clone();
2419 let instances: Vec<InstanceStatus> = self
2420 .insts
2421 .iter()
2422 .map(|i| {
2423 let rt = self.rt.get(&i.name);
2424 let m = snap.get(&format!("{}/{}", self.oclient.project_name(), i.name));
2425 let health = match (probe, rt.and_then(|r| r.healthy)) {
2426 (false, _) => "none",
2427 (true, Some(true)) => "healthy",
2428 (true, Some(false)) => "unhealthy",
2429 (true, None) => "starting",
2430 };
2431 InstanceStatus {
2432 name: i.name.clone(),
2433 slot: i.slot,
2434 rev: i.rev.clone(),
2435 status: m
2436 .map(|m| m.status.clone())
2437 .unwrap_or_else(|| i.status.clone()),
2438 health: health.into(),
2439 ip: rt.and_then(|r| r.ip).map(|ip| ip.to_string()),
2440 in_rotation: rt.is_some_and(|r| r.in_rotation),
2441 restarts: rt.map(|r| r.restarts.len() as u32).unwrap_or(0),
2442 last_probe: rt.map(|r| r.last_probe.clone()).unwrap_or_default(),
2443 cpu_pct: m.and_then(|m| m.cpu_pct),
2444 cpu_history: m.map(|m| m.cpu_history.clone()).unwrap_or_default(),
2445 mem_bytes: m.and_then(|m| m.mem_bytes),
2446 disk_bytes: m.and_then(|m| m.disk_bytes),
2447 }
2448 })
2449 .collect();
2450 let mut instances = instances;
2451 instances.sort_by(|a, b| (a.slot, &a.name).cmp(&(b.slot, &b.name)));
2452 let routes = self.inner.balancer.routes();
2453 let now = Instant::now();
2454 let mut ports = Vec::new();
2455 for (k, p) in &self.routes {
2456 let r = routes.iter().find(|r| r.key == *k);
2457 let accepted = r.map(|r| r.accepted).unwrap_or(0);
2458 let e = self
2459 .rates
2460 .entry(k.clone())
2461 .or_insert_with(|| (accepted, now, VecDeque::new()));
2462 let dt = now.duration_since(e.1).as_secs_f32();
2463 if dt >= 1.0 {
2466 let rate = accepted.saturating_sub(e.0) as f32 / dt;
2467 if e.2.len() == crate::metrics::HISTORY {
2468 e.2.pop_front();
2469 }
2470 e.2.push_back(rate);
2471 e.0 = accepted;
2472 e.1 = now;
2473 }
2474 ports.push(PortStatus {
2475 listen: p.listen.to_string(),
2476 target: p.target,
2477 backends: r
2478 .map(|r| r.backends.iter().map(|b| b.addr.to_string()).collect())
2479 .unwrap_or_default(),
2480 error: self.route_errors.get(k).cloned(),
2481 accepted,
2482 active: r
2483 .map(|r| {
2484 r.backends
2485 .iter()
2486 .chain(r.draining.iter())
2487 .map(|b| b.active)
2488 .sum()
2489 })
2490 .unwrap_or(0),
2491 rate_history: e.2.iter().copied().collect(),
2492 });
2493 }
2494 self.rates.retain(|k, _| self.routes.contains_key(k));
2495 let running = instances
2496 .iter()
2497 .filter(|i| i.status.eq_ignore_ascii_case("running"))
2498 .count() as u32;
2499 let healthy = instances.iter().filter(|i| i.in_rotation).count() as u32;
2500 self.watch_health(healthy, spec.replicas());
2501 let st = ServiceStatus {
2502 service: self.service.clone(),
2503 image: spec.image.clone(),
2504 rev,
2505 replicas: spec.replicas(),
2506 running,
2507 healthy,
2508 state: self.state.clone(),
2509 message: self.message.clone(),
2510 instances,
2511 ports,
2512 rollout: self.rollout.clone(),
2513 checked_at: now_secs(),
2514 domains: Vec::new(),
2515 };
2516 self.inner.status.lock().unwrap().insert(self.key(), st);
2517 }
2518}
2519
2520pub fn unsettled(st: &StackStatus) -> BTreeSet<String> {
2523 st.services
2524 .iter()
2525 .filter(|s| s.state != "converged")
2526 .map(|s| s.service.clone())
2527 .collect()
2528}
2529
2530#[cfg(test)]
2531#[path = "controller_tests.rs"]
2532mod tests;