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