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