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