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