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)]
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")]
196 pub kind: Option<String>,
197}
198
199pub type EventSink = Arc<dyn Fn(&Event) + Send + Sync>;
202
203const EVENTS_KEPT: usize = 1000;
205
206#[derive(Debug, Clone, Default)]
208pub struct Snapshot {
209 pub host: crate::metrics::HostSample,
210 pub instances: BTreeMap<String, crate::metrics::InstanceSample>,
212 pub at: u64,
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 while let Some(inner) = weak.upgrade() {
411 match sampler.sample(&inner.client) {
412 Ok((host, insts)) => {
413 if let Some(tx) = &*inner.metrics_sink.lock().unwrap() {
414 crate::metrics_history::offer(tx, (now_ms(), insts.clone()));
415 }
416 *inner.snapshot.lock().unwrap() = Snapshot {
417 host,
418 instances: insts
419 .into_iter()
420 .map(|i| (format!("{}/{}", i.project, i.name), i))
421 .collect(),
422 at: now_ms(),
423 };
424 }
425 Err(e) => eprintln!("isb serve: metrics: {e}"),
426 }
427 drop(inner);
428 std::thread::sleep(SAMPLE_EVERY);
429 }
430 });
431 let defs = c.inner.store.load_all()?;
432 dns::prune(&defs);
434 for def in defs {
435 eprintln!("isb serve: resuming stack {}", def.name);
436 c.apply(Arc::new(def));
437 }
438 let all: Vec<(String, String)> = c
440 .definitions()
441 .iter()
442 .flat_map(|def| {
443 let q = def.qualified();
444 def.secrets.keys().map(move |k| (q.clone(), k.clone()))
445 })
446 .collect();
447 c.poll(all);
448 Ok(c)
449 }
450
451 pub fn balancer(&self) -> &Balancer {
452 &self.inner.balancer
453 }
454
455 pub fn events(&self, since: u64, limit: usize) -> (u64, Vec<Event>) {
457 let e = self.inner.events.lock().unwrap();
458 let out: Vec<Event> = e.1.iter().filter(|x| x.seq > since).cloned().collect();
459 let skip = out.len().saturating_sub(limit);
460 (e.0, out.into_iter().skip(skip).collect())
461 }
462
463 pub fn resume_from(&self, since: u64) -> u64 {
468 if since > self.inner.events.lock().unwrap().0 {
469 0
470 } else {
471 since
472 }
473 }
474
475 pub fn wait_events(&self, since: u64, limit: usize, timeout: Duration) -> (u64, Vec<Event>) {
477 let started = Instant::now();
478 loop {
479 let r = self.events(since, limit);
480 if !r.1.is_empty() || started.elapsed() >= timeout {
481 return r;
482 }
483 std::thread::sleep(Duration::from_millis(250));
484 }
485 }
486
487 pub fn set_event_sink(&self, sink: EventSink) {
491 let e = self.inner.events.lock().unwrap();
492 for ev in &e.1 {
493 sink(ev);
494 }
495 *self.inner.event_sink.lock().unwrap() = Some(sink);
496 }
497
498 pub fn set_metrics_sink(
500 &self,
501 tx: std::sync::mpsc::SyncSender<crate::metrics_history::Sample>,
502 ) {
503 *self.inner.metrics_sink.lock().unwrap() = Some(tx);
504 }
505
506 pub fn snapshot(&self) -> Snapshot {
508 self.inner.snapshot.lock().unwrap().clone()
509 }
510
511 pub fn event(&self, kind: &str, level: &str, stack: &str, service: &str, message: String) {
515 eprintln!("isb serve: {stack}: {message}");
516 self.inner
517 .emit_kind(Some(kind), level, stack, service, None, message);
518 }
519
520 pub fn relay(
523 &self,
524 kind: Option<&str>,
525 level: &str,
526 stack: &str,
527 service: &str,
528 instance: Option<&str>,
529 message: String,
530 ) {
531 self.inner
532 .emit_kind(kind, level, stack, service, instance, message);
533 }
534
535 pub fn note(&self, level: &str, stack: &str, message: String) {
536 eprintln!("isb serve: {stack}: {message}");
537 self.inner.emit(level, stack, "", None, message);
538 }
539
540 pub fn note_service(&self, level: &str, stack: &str, service: &str, message: String) {
543 self.inner.emit(level, stack, service, None, message);
544 }
545
546 pub fn service_event(&self, level: &str, stack: &str, service: &str, message: String) {
548 eprintln!("isb serve: {stack}/{service}: {message}");
549 self.inner.emit(level, stack, service, None, message);
550 }
551
552 fn notify_stacks(&self) {
553 if let Some(o) = &self.inner.observer {
554 o.stacks_changed(self.definitions());
555 }
556 }
557
558 pub fn plan(&self, def: &StackDef) -> Result<Vec<DeployChange>> {
560 self.validate(def)?;
561 let mut def = def.clone();
562 self.pin_images(&mut def)?;
563 let def = &def;
564 let old = self
565 .inner
566 .stacks
567 .lock()
568 .unwrap()
569 .get(&def.qualified())
570 .cloned();
571 diff(old.as_deref(), def)
572 }
573
574 pub fn client(&self) -> &Client {
575 &self.inner.client
576 }
577
578 pub fn validate(&self, def: &StackDef) -> Result<()> {
581 validate_stack_name(&def.name)?;
582 let host = crate::sandbox::host_facts(&crate::org::client(&self.inner.client, &def.org))?;
584 for (svc, spec) in &def.file.services {
585 let mut s = instance_spec(def, svc, spec, 1, "0000")?;
586 s.name = Some(instance_name(&def.name, svc, 1, "0000")?);
587 instance_name(&def.name, svc, spec.replicas().max(1), "0000")?;
588 crate::plan::resolve(&s, &def.file.volumes, &host, &def.base_dir)?;
589 published(spec)?;
590 crate::ingress::domain::validate(svc, &spec.domains)?;
591 }
592 super::ports::validate_udp(&self.inner.client, def, &self.definitions())
593 }
594
595 pub fn deploy(&self, mut def: StackDef) -> Result<Vec<DeployChange>> {
598 self.validate(&def)?;
599 self.pin_images(&mut def)?;
600 let _g = self.inner.edit.lock().unwrap();
601 let old = self
602 .inner
603 .stacks
604 .lock()
605 .unwrap()
606 .get(&def.qualified())
607 .cloned();
608 if let Some(old) = &old {
609 let mut prev = (**old).clone();
610 prev.previous = None;
611 def.previous = Some(Box::new(prev));
612 for (k, v) in &old.force {
614 def.force.entry(k.clone()).or_insert(*v);
615 }
616 }
617 let changes = diff(old.as_deref(), &def)?;
618 self.inner.store.save(&def)?;
619 self.apply(Arc::new(def));
620 Ok(changes)
621 }
622
623 fn apply(&self, def: Arc<StackDef>) {
626 let name = def.qualified();
627 self.inner
628 .stacks
629 .lock()
630 .unwrap()
631 .insert(name.clone(), def.clone());
632 let mut workers = self.inner.workers.lock().unwrap();
633 for svc in def.file.services.keys() {
634 let key = (name.clone(), svc.clone());
635 match workers.get(&key) {
636 Some(w) => {
637 let mut slot = w.slot.lock().unwrap();
638 slot.def = def.clone();
639 slot.remove = false;
640 if let Some(st) = self.inner.status.lock().unwrap().get_mut(&key) {
644 if st.state == "failing" {
645 st.state = "updating".into();
646 st.message = None;
647 }
648 }
649 w.wake.notify_all();
650 }
651 None => {
652 let shared = Arc::new(WorkerShared {
653 slot: Mutex::new(Slot {
654 def: def.clone(),
655 remove: false,
656 remove_volumes: false,
657 }),
658 wake: Condvar::new(),
659 stop: AtomicBool::new(false),
660 kick: AtomicBool::new(false),
661 });
662 workers.insert(key, shared.clone());
663 spawn_worker(self.inner.clone(), &def, svc.clone(), shared);
664 }
665 }
666 }
667 for ((stack, svc), w) in workers.iter() {
668 if *stack == name && !def.file.services.contains_key(svc) {
669 let mut slot = w.slot.lock().unwrap();
670 slot.remove = true;
671 w.wake.notify_all();
672 }
673 }
674 drop(workers);
675 self.notify_stacks();
676 }
677
678 pub fn org_limits_changed(&self, org: &OrgId) -> usize {
682 let stacks = self.inner.stacks.lock().unwrap();
683 let workers = self.inner.workers.lock().unwrap();
684 let status = self.inner.status.lock().unwrap();
685 let mut n = 0;
686 for (key, w) in workers.iter() {
687 let ours = stacks.get(&key.0).is_some_and(|d| d.org == *org);
688 let limited = status.get(key).is_some_and(|s| {
689 s.state == "failing" && s.message.as_deref().is_some_and(limit_error)
690 });
691 if ours && limited {
692 w.kick.store(true, Ordering::SeqCst);
693 let _slot = w.slot.lock().unwrap();
696 w.wake.notify_all();
697 n += 1;
698 }
699 }
700 n
701 }
702
703 pub fn remove(&self, name: &str, volumes: bool, timeout: Duration) -> Result<()> {
707 {
708 let _g = self.inner.edit.lock().unwrap();
709 let Some(def) = self.inner.stacks.lock().unwrap().remove(name) else {
710 return Err(Error::NotFound(format!("stack {name}")));
711 };
712 self.inner.store.remove(&def.org, &def.name)?;
713 }
714 self.notify_stacks();
715 let ws: Vec<Arc<WorkerShared>> = self
716 .inner
717 .workers
718 .lock()
719 .unwrap()
720 .iter()
721 .filter(|((s, _), _)| s == name)
722 .map(|(_, w)| w.clone())
723 .collect();
724 for w in &ws {
725 let mut slot = w.slot.lock().unwrap();
726 slot.remove = true;
727 slot.remove_volumes = volumes;
728 w.wake.notify_all();
729 }
730 let started = Instant::now();
731 while started.elapsed() < timeout {
732 let left = self
733 .inner
734 .workers
735 .lock()
736 .unwrap()
737 .keys()
738 .any(|(s, _)| s == name);
739 if !left {
740 return Ok(());
741 }
742 std::thread::sleep(Duration::from_millis(200));
743 }
744 Err(Error::invalid(format!(
745 "stack {name}: still removing after {timeout:?}; it carries on in the background"
746 )))
747 }
748
749 pub fn rollback(&self, name: &str) -> Result<Vec<DeployChange>> {
752 let _g = self.inner.edit.lock().unwrap();
753 let cur = self.get_def(name)?;
754 let prev = cur
755 .previous
756 .clone()
757 .ok_or_else(|| Error::invalid(format!("stack {name} has no previous deployment")))?;
758 let mut def = *prev;
759 def.deployed_at = now_secs();
760 for b in def.secrets.values_mut() {
763 if let Ok(v) = self.inner.secrets.version_in(&b.driver, &def.org, &b.name) {
764 b.version = v;
765 }
766 }
767 let mut cur2 = (*cur).clone();
768 cur2.previous = None;
769 let changes = diff(Some(&cur), &def)?;
770 def.previous = Some(Box::new(cur2));
771 self.inner.store.save(&def)?;
772 self.apply(Arc::new(def));
773 Ok(changes)
774 }
775
776 pub fn scale(&self, name: &str, service: &str, replicas: u32) -> Result<()> {
778 let _g = self.inner.edit.lock().unwrap();
779 let cur = self.get_def(name)?;
780 let mut def = (*cur).clone();
781 let spec = def
782 .file
783 .services
784 .get_mut(service)
785 .ok_or_else(|| Error::NotFound(format!("service {service} in stack {name}")))?;
786 spec.deploy.get_or_insert_with(Default::default).replicas = Some(replicas);
787 super::ports::check_replicas(service, spec)?;
788 instance_name(&def.name, service, replicas.max(1), "0000")?;
790 self.inner.store.save(&def)?;
791 self.apply(Arc::new(def));
792 Ok(())
793 }
794
795 pub fn redeploy(&self, name: &str, service: &str) -> Result<()> {
799 let _g = self.inner.edit.lock().unwrap();
800 let cur = self.get_def(name)?;
801 cur.service(service)?;
802 let mut def = (*cur).clone();
803 *def.force.entry(service.to_string()).or_insert(0) += 1;
804 self.pin_images(&mut def)?;
806 self.inner.store.save(&def)?;
807 self.apply(Arc::new(def));
808 Ok(())
809 }
810
811 fn pin_images(&self, def: &mut StackDef) -> Result<()> {
815 def.images.clear();
816 for (svc, spec) in &def.file.services {
817 let Some(r) = spec.image.strip_prefix("registry:") else {
818 continue;
819 };
820 let r = crate::registry::ImageRef::parse(r)?;
821 let reg = crate::registry::Registry::shared(&self.inner.client)?;
822 let d = reg.resolve(&def.org, &r)?;
825 if r.digest.is_none() {
826 def.images.insert(svc.clone(), d);
827 }
828 }
829 Ok(())
830 }
831
832 fn get_def(&self, name: &str) -> Result<Arc<StackDef>> {
833 self.inner
834 .stacks
835 .lock()
836 .unwrap()
837 .get(name)
838 .cloned()
839 .ok_or_else(|| Error::NotFound(format!("stack {name}")))
840 }
841
842 pub fn definitions(&self) -> Vec<Arc<StackDef>> {
844 self.inner
845 .stacks
846 .lock()
847 .unwrap()
848 .values()
849 .cloned()
850 .collect()
851 }
852
853 pub fn definition(&self, name: &str) -> Result<StackDef> {
855 self.get_def(name).map(|d| (*d).clone())
856 }
857
858 pub fn list(&self) -> Vec<StackStatus> {
859 let names: Vec<String> = self.inner.stacks.lock().unwrap().keys().cloned().collect();
860 names.iter().filter_map(|n| self.status(n).ok()).collect()
861 }
862
863 pub fn status(&self, name: &str) -> Result<StackStatus> {
864 let def = self.get_def(name)?;
865 let st = self.inner.status.lock().unwrap();
866 let services: Vec<ServiceStatus> = def
867 .file
868 .services
869 .keys()
870 .map(|svc| {
871 st.get(&(name.to_string(), svc.clone()))
872 .cloned()
873 .unwrap_or_else(|| ServiceStatus {
874 service: svc.clone(),
875 state: "starting".into(),
876 ..Default::default()
877 })
878 })
879 .collect();
880 drop(st);
881 let mut services = services;
882 if let Some(o) = &self.inner.observer {
883 for s in &mut services {
884 s.domains = o.domains(name, &s.service);
885 }
886 }
887 for s in services.iter_mut().filter(|s| s.state != "converged") {
888 s.last_failed_attempt = self
889 .inner
890 .failures
891 .last(name, &s.service)
892 .map(|f| f.summary());
893 }
894 let converged = services.iter().all(|s| s.state == "converged");
895 Ok(StackStatus {
896 name: def.name.clone(),
897 org: def.org.to_string(),
898 deployed_at: def.deployed_at,
899 deployed_by: def.deployed_by.clone(),
900 has_previous: def.previous.is_some(),
901 converged,
902 services,
903 })
904 }
905
906 pub fn logs(
908 &self,
909 name: &str,
910 service: &str,
911 slot: Option<u32>,
912 lines: usize,
913 ) -> Result<BTreeMap<String, String>> {
914 let def = self.get_def(name)?;
915 let oci = crate::plan::ImageSource::parse(&def.service(service)?.image)?.is_oci();
916 let oc = crate::org::client(&self.inner.client, &def.org);
917 super::failure::replica_logs(&oc, &def.name, service, oci, slot, lines)
918 }
919
920 pub fn last_failure(
924 &self,
925 name: &str,
926 service: &str,
927 always: bool,
928 ) -> Option<super::failure::FailedAttempt> {
929 let state = self.inner.status.lock().unwrap();
930 let converged = state
931 .get(&(name.to_string(), service.to_string()))
932 .is_some_and(|s| s.state == "converged");
933 drop(state);
934 if converged && !always {
935 return None;
936 }
937 self.inner.failures.last(name, service)
938 }
939
940 pub fn failed_attempts(&self, name: &str) -> BTreeMap<String, super::failure::FailedAttempt> {
942 self.inner.failures.of_stack(name)
943 }
944
945 pub fn shutdown(&self) {
948 for w in self.inner.workers.lock().unwrap().values() {
949 w.stop.store(true, Ordering::SeqCst);
950 w.wake.notify_all();
951 }
952 self.inner.balancer.clear();
953 }
954}
955
956fn limit_error(msg: &str) -> bool {
959 msg.contains(" quota (") || msg.contains(" limit (")
960}
961
962fn instance_spec(
967 def: &StackDef,
968 service: &str,
969 spec: &SandboxSpec,
970 slot: u32,
971 rev: &str,
972) -> Result<SandboxSpec> {
973 let mut s = spec.clone();
974 s.image = def.instance_image(service, &spec.image);
975 s.restart = Some(RestartMode::Always);
976 (s.ports, s.stack_udp) = super::ports::instance_ports(spec)?;
977 s.domains.clear();
978 if let Some(d) = &s.deploy {
979 s.labels.extend(d.labels.clone());
980 }
981 s.labels.insert(LABEL_STACK.into(), def.name.clone());
982 s.labels.insert(LABEL_SERVICE.into(), service.into());
983 s.labels.insert(LABEL_SLOT.into(), slot.to_string());
984 s.labels.insert(LABEL_REV.into(), rev.into());
985 let live = live_versions(def, service);
986 if !live.is_empty() {
987 s.labels
988 .insert(LABEL_SECRETS.into(), super::secrets::versions_label(&live));
989 }
990 Ok(s)
991}
992
993fn live_versions(def: &StackDef, service: &str) -> BTreeMap<String, u64> {
995 def.live_secrets(service)
996 .into_iter()
997 .map(|(k, (_, v))| (k, v))
998 .collect()
999}
1000
1001#[derive(Debug, Clone)]
1003#[doc(hidden)]
1004pub struct Inst {
1005 pub name: String,
1006 pub(super) slot: u32,
1007 pub rev: String,
1008 status: String,
1009 secrets: Option<BTreeMap<String, u64>>,
1012}
1013
1014impl Inst {
1015 #[doc(hidden)]
1016 pub fn is_running(&self) -> bool {
1017 self.running()
1018 }
1019
1020 fn running(&self) -> bool {
1021 self.status.eq_ignore_ascii_case("running")
1022 }
1023}
1024
1025struct Worker {
1027 inner: Arc<Inner>,
1028 stack: String,
1030 q: String,
1032 oclient: Client,
1034 service: String,
1035 shared: Arc<WorkerShared>,
1036 rt: BTreeMap<String, InstRt>,
1037 paused: Option<(String, String)>,
1040 create_backoff: Option<(Instant, Duration)>,
1042 routes: BTreeMap<String, Published>,
1043 route_errors: BTreeMap<String, String>,
1044 template: Option<(String, Desired)>,
1046 state: String,
1047 message: Option<String>,
1048 last_error: Option<String>,
1049 insts: Vec<Inst>,
1051 rollout: Option<RolloutStatus>,
1052 rates: BTreeMap<String, (u64, Instant, VecDeque<f32>)>,
1054 deps_met: bool,
1057 org: OrgId,
1058 dns_last: Option<(Vec<IpAddr>, Option<String>)>,
1061 dns_error: Option<String>,
1062 observed: Option<Vec<IpAddr>>,
1064 ever_healthy: bool,
1066 health_down: Option<Instant>,
1068 health_alarm: bool,
1070 seen: Option<Arc<StackDef>>,
1072 restart_failed: Option<(BTreeMap<String, u64>, String)>,
1075}
1076
1077fn spawn_worker(inner: Arc<Inner>, def: &StackDef, service: String, shared: Arc<WorkerShared>) {
1078 let (stack, q, org) = (def.name.clone(), def.qualified(), def.org.clone());
1079 let oclient = crate::org::client(&inner.client, &def.org);
1080 let name = format!("isb-{q}-{service}");
1081 let r = std::thread::Builder::new().name(name).spawn(move || {
1082 let mut w = Worker::new(inner, stack, q, oclient, service, shared, org);
1083 w.run();
1084 });
1085 if let Err(e) = r {
1086 eprintln!("isb serve: cannot start a worker thread: {e}");
1087 }
1088}
1089
1090impl Worker {
1091 fn new(
1092 inner: Arc<Inner>,
1093 stack: String,
1094 q: String,
1095 oclient: Client,
1096 service: String,
1097 shared: Arc<WorkerShared>,
1098 org: OrgId,
1099 ) -> Worker {
1100 Worker {
1101 inner,
1102 stack,
1103 q,
1104 oclient,
1105 service,
1106 shared,
1107 rt: BTreeMap::new(),
1108 paused: None,
1109 create_backoff: None,
1110 routes: BTreeMap::new(),
1111 route_errors: BTreeMap::new(),
1112 template: None,
1113 state: "starting".into(),
1114 message: None,
1115 last_error: None,
1116 insts: Vec::new(),
1117 rollout: None,
1118 rates: BTreeMap::new(),
1119 deps_met: false,
1120 org,
1121 dns_last: None,
1122 dns_error: None,
1123 observed: None,
1124 ever_healthy: false,
1125 health_down: None,
1126 health_alarm: false,
1127 seen: None,
1128 restart_failed: None,
1129 }
1130 }
1131
1132 fn begin_pass(&mut self, def: &Arc<StackDef>) {
1137 let new_def = self.seen.as_ref().is_none_or(|d| !Arc::ptr_eq(d, def));
1138 let kicked = self.shared.kick.swap(false, Ordering::SeqCst);
1139 if new_def {
1140 self.seen = Some(def.clone());
1141 self.last_error = None;
1142 if self.state == "failing" {
1143 self.state = "updating".into();
1144 self.message = None;
1145 }
1146 }
1147 if new_def || kicked {
1148 self.create_backoff = None;
1149 }
1150 }
1151}
1152
1153impl Worker {
1154 fn log(&self, msg: &str) {
1155 self.event("info", None, msg);
1156 }
1157
1158 fn event(&self, level: &str, instance: Option<&str>, msg: &str) {
1159 eprintln!("isb serve: {}/{}: {msg}", self.q, self.service);
1160 self.inner
1161 .emit(level, &self.q, &self.service, instance, msg.to_string());
1162 }
1163
1164 fn client(&self) -> &Client {
1165 &self.oclient
1166 }
1167
1168 fn watch_health(&mut self, healthy: u32, replicas: u32) {
1172 if healthy > 0 {
1173 self.ever_healthy = true;
1174 self.health_down = None;
1175 if self.health_alarm {
1176 self.health_alarm = false;
1177 let msg = format!("{healthy} of {replicas} replicas healthy again");
1178 eprintln!("isb serve: {}/{}: {msg}", self.q, self.service);
1179 self.inner.emit_kind(
1180 Some("health.recovered"),
1181 "info",
1182 &self.q,
1183 &self.service,
1184 None,
1185 msg,
1186 );
1187 }
1188 return;
1189 }
1190 if replicas == 0 || !self.ever_healthy || self.rollout.is_some() {
1191 self.health_down = None;
1192 return;
1193 }
1194 let since = *self.health_down.get_or_insert_with(Instant::now);
1195 if !self.health_alarm && since.elapsed() >= HEALTH_DEBOUNCE {
1196 self.health_alarm = true;
1197 let msg = format!(
1198 "no healthy replica (of {replicas}) for {}s",
1199 since.elapsed().as_secs()
1200 );
1201 eprintln!("isb serve: {}/{}: {msg}", self.q, self.service);
1202 self.inner.emit_kind(
1203 Some("health.unhealthy"),
1204 "error",
1205 &self.q,
1206 &self.service,
1207 None,
1208 msg,
1209 );
1210 }
1211 }
1212
1213 fn key(&self) -> (String, String) {
1214 (self.q.clone(), self.service.clone())
1215 }
1216
1217 fn run(&mut self) {
1218 loop {
1219 if self.shared.stop.load(Ordering::SeqCst) {
1220 return;
1221 }
1222 let (def, remove, remove_volumes) = {
1223 let s = self.shared.slot.lock().unwrap();
1224 (s.def.clone(), s.remove, s.remove_volumes)
1225 };
1226 if remove {
1227 self.teardown(&def, remove_volumes);
1228 let mut ws = self.inner.workers.lock().unwrap();
1231 let slot = self.shared.slot.lock().unwrap();
1232 if slot.remove || self.shared.stop.load(Ordering::SeqCst) {
1233 ws.remove(&self.key());
1234 self.inner.status.lock().unwrap().remove(&self.key());
1235 return;
1236 }
1237 continue;
1238 }
1239 self.begin_pass(&def);
1240 if let Err(e) = self.pass(&def) {
1241 self.state = "failing".into();
1242 self.message = Some(e.to_string());
1243 if self.last_error.as_deref() != Some(&e.to_string()) {
1245 self.event("error", None, &e.to_string());
1246 self.last_error = Some(e.to_string());
1247 }
1248 self.publish_status(&def);
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 continue;
1254 }
1255 self.last_error = None;
1256 let slot = self.shared.slot.lock().unwrap();
1257 if Arc::ptr_eq(&slot.def, &def) && !slot.remove {
1258 let _ = self.shared.wake.wait_timeout(slot, self.inner.interval);
1259 }
1260 }
1261 }
1262
1263 fn superseded(&self, def: &Arc<StackDef>) -> bool {
1265 let s = self.shared.slot.lock().unwrap();
1266 s.remove || !Arc::ptr_eq(&s.def, def) || self.shared.stop.load(Ordering::SeqCst)
1267 }
1268
1269 fn teardown(&mut self, def: &StackDef, volumes: bool) {
1271 self.unpublish_dns();
1273 if let Some(o) = &self.inner.observer {
1274 o.rotation(&self.q, &self.service, &[]);
1275 }
1276 self.observed = Some(Vec::new());
1277 for (k, _) in std::mem::take(&mut self.routes) {
1278 self.inner.balancer.remove_route(&k);
1279 }
1280 match list_instances(self.client(), &self.stack, Some(&self.service)) {
1281 Ok(insts) => {
1282 for i in insts {
1283 self.log(&format!("removing {}", i.name));
1284 if let Err(e) = Sandbox::remove(self.client(), &i.name, true) {
1285 if !e.is_not_found() {
1286 self.log(&format!("cannot remove {}: {e}", i.name));
1287 }
1288 }
1289 }
1290 }
1291 Err(e) => self.log(&format!("cannot list instances to remove: {e}")),
1292 }
1293 if volumes {
1294 self.remove_volumes(def);
1295 }
1296 self.rt.clear();
1297 self.template = None;
1298 }
1299
1300 fn remove_volumes(&self, def: &StackDef) {
1301 let Ok(spec) = def.service(&self.service) else {
1302 return;
1303 };
1304 let Ok(host) = crate::sandbox::host_facts(self.client()) else {
1305 return;
1306 };
1307 let Ok(pool) = host.pick_pool(spec.storage.as_deref()) else {
1308 return;
1309 };
1310 for v in &spec.volumes {
1311 if v.mount_type != crate::spec::MountType::Volume {
1312 continue;
1313 }
1314 let d = def.file.volumes.get(&v.source);
1315 if v.external || d.is_some_and(|d| d.external) {
1316 continue;
1317 }
1318 let name = d
1319 .and_then(|d| d.name.clone())
1320 .unwrap_or_else(|| v.source.clone());
1321 let vpool = match v.pool.as_deref().or(d.and_then(|d| d.pool.as_deref())) {
1322 Some(p) if p != "auto" => p.to_string(),
1323 _ => pool.clone(),
1324 };
1325 match crate::volume::remove(self.client(), &vpool, &name) {
1326 Ok(()) => self.log(&format!("volume {name}: deleted")),
1327 Err(e) if e.is_not_found() => {}
1328 Err(e) => self.log(&format!("volume {name}: kept ({e})")),
1330 }
1331 }
1332 }
1333
1334 #[expect(
1336 clippy::too_many_lines,
1337 reason = "predates the lint ratchet; split it when next changed"
1338 )]
1339 fn pass(&mut self, def: &Arc<StackDef>) -> Result<()> {
1340 let spec = def.service(&self.service)?.clone();
1341 let rev = def.revision(&self.service)?;
1342 let replicas = spec.replicas();
1343 let oci = crate::plan::ImageSource::parse(&spec.image)?.is_oci();
1344 let probe = spec.health_probe().map_err(Error::invalid)?;
1345
1346 if !self.deps_met {
1347 if let Some(msg) = self.waiting_for(&spec) {
1348 self.state = "waiting".into();
1349 self.message = Some(msg);
1350 self.publish_status(def);
1351 return Ok(());
1352 }
1353 self.deps_met = true;
1354 }
1355 if self.template.as_ref().is_none_or(|(r, _)| *r != rev) {
1356 let mut s = instance_spec(def, &self.service, &spec, 1, &rev)?;
1357 s.name = Some(instance_name(&self.stack, &self.service, 1, "0000")?);
1358 let d = crate::sandbox::resolve(self.client(), &s, &def.file.volumes, &def.base_dir)?;
1359 self.template = Some((rev.clone(), d));
1360 }
1361 self.set_routes(&spec);
1362
1363 let mut insts = list_instances(self.client(), &self.stack, Some(&self.service))?;
1364 self.rt.retain(|n, _| insts.iter().any(|i| i.name == *n));
1365
1366 let extra: Vec<Inst> = insts
1368 .iter()
1369 .filter(|i| i.slot > replicas || i.slot == 0)
1370 .cloned()
1371 .collect();
1372 for i in extra.iter().rev() {
1373 self.log(&format!("scaling down: removing {}", i.name));
1374 self.retire(&i.name)?;
1375 }
1376 insts.retain(|i| i.slot >= 1 && i.slot <= replicas);
1377 self.insts = insts.clone();
1378
1379 for i in &insts {
1381 self.maintain(def, i, &spec, oci, probe.as_ref())?;
1382 }
1383 for slot in 1..=replicas {
1386 let current_ok = insts.iter().any(|i| {
1387 i.slot == slot
1388 && i.rev == rev
1389 && self.rt.get(&i.name).is_some_and(|r| r.in_rotation)
1390 });
1391 if current_ok {
1392 for i in insts.iter().filter(|i| i.slot == slot && i.rev != rev) {
1393 self.log(&format!("removing leftover {}", i.name));
1394 self.retire(&i.name)?;
1395 }
1396 }
1397 let mut current: Vec<&Inst> = insts
1398 .iter()
1399 .filter(|i| i.slot == slot && i.rev == rev)
1400 .collect();
1401 while current.len() > 1 {
1403 let i = current.pop().unwrap();
1404 self.log(&format!("removing duplicate {}", i.name));
1405 self.retire(&i.name)?;
1406 }
1407 }
1408 self.sync_routes();
1409 self.publish_status(def);
1410
1411 let mut pending: Vec<(u32, Option<String>)> = Vec::new();
1413 for slot in 1..=replicas {
1414 if insts.iter().any(|i| i.slot == slot && i.rev == rev) {
1415 continue;
1416 }
1417 let old = insts
1418 .iter()
1419 .find(|i| i.slot == slot)
1420 .map(|i| i.name.clone());
1421 pending.push((slot, old));
1422 }
1423 if pending.is_empty() {
1424 self.create_backoff = None;
1425 let uc = spec
1427 .deploy
1428 .as_ref()
1429 .and_then(|d| d.update_config.clone())
1430 .unwrap_or_default();
1431 self.cycle_in_place(def, &spec, &uc)?;
1432 let insts = self.insts.clone();
1433 let all_ok = insts
1434 .iter()
1435 .all(|i| i.rev == rev && self.rt.get(&i.name).is_some_and(|r| r.in_rotation));
1436 self.state = if all_ok { "converged" } else { "failing" }.into();
1437 if all_ok {
1438 self.message = self.restart_failed.as_ref().map(|(_, m)| m.clone());
1439 } else if self.message.is_none() {
1440 self.message = Some("some replicas are not healthy".into());
1441 }
1442 self.publish_status(def);
1443 return Ok(());
1444 }
1445 if self.paused.as_ref().is_some_and(|(r, _)| *r == rev) {
1446 self.state = "paused".into();
1447 self.message = self.paused.as_ref().map(|(_, m)| m.clone());
1448 self.publish_status(def);
1449 return Ok(());
1450 }
1451 if let Some((at, wait)) = self.create_backoff {
1452 if at.elapsed() < wait {
1453 return Ok(());
1454 }
1455 }
1456 self.state = "updating".into();
1457 self.message = None;
1458 self.publish_status(def);
1459
1460 let uc: UpdateConfig = spec
1461 .deploy
1462 .as_ref()
1463 .and_then(|d| d.update_config.clone())
1464 .unwrap_or_default();
1465 let parallel = match uc.parallelism.unwrap_or(1) {
1466 0 => pending.len(),
1467 n => n as usize,
1468 };
1469 let delay = uc
1470 .delay
1471 .as_deref()
1472 .map(crate::flex::parse_duration)
1473 .transpose()
1474 .map_err(Error::invalid)?
1475 .unwrap_or_default();
1476 let monitor = uc
1477 .monitor
1478 .as_deref()
1479 .map(crate::flex::parse_duration)
1480 .transpose()
1481 .map_err(Error::invalid)?
1482 .unwrap_or(Duration::from_secs(5));
1483 let order = uc.order.unwrap_or_default();
1484 let order_name = match order {
1485 UpdateOrder::StopFirst => "stop-first",
1486 UpdateOrder::StartFirst => "start-first",
1487 };
1488 self.rollout = Some(RolloutStatus {
1489 to_rev: rev.clone(),
1490 order: order_name.into(),
1491 parallelism: parallel,
1492 done: 0,
1493 total: pending.len(),
1494 started_at: now_secs(),
1495 slots: pending
1496 .iter()
1497 .map(|(slot, old)| SlotRollout {
1498 slot: *slot,
1499 old: old.clone(),
1500 old_rev: old
1501 .as_ref()
1502 .and_then(|o| insts.iter().find(|i| i.name == *o))
1503 .map(|i| i.rev.clone()),
1504 old_state: if old.is_some() { "serving" } else { "none" }.into(),
1505 new: None,
1506 new_state: "waiting".into(),
1507 })
1508 .collect(),
1509 });
1510 let rollout_started = Instant::now();
1511 self.log(&format!(
1512 "rolling out rev {rev} to {} slot(s), {order_name}",
1513 pending.len()
1514 ));
1515 self.publish_status(def);
1516 let r = self.roll(
1517 def,
1518 &spec,
1519 &rev,
1520 &pending,
1521 parallel,
1522 delay,
1523 monitor,
1524 order,
1525 oci,
1526 probe.as_ref(),
1527 &uc,
1528 );
1529 let rollout = self.rollout.take();
1530 if let (Ok(true), Some(ro)) = (&r, rollout) {
1531 self.log(&format!(
1532 "rollout of rev {rev} complete: {}/{} slot(s) in {:.0?}",
1533 ro.done,
1534 ro.total,
1535 rollout_started.elapsed()
1536 ));
1537 }
1538 self.publish_status(def);
1539 r.map(|_| ())
1540 }
1541
1542 #[expect(clippy::too_many_arguments)]
1545 #[expect(
1546 clippy::excessive_nesting,
1547 reason = "predates the lint ratchet; split it when next changed"
1548 )]
1549 fn roll(
1550 &mut self,
1551 def: &Arc<StackDef>,
1552 spec: &SandboxSpec,
1553 rev: &String,
1554 pending: &[(u32, Option<String>)],
1555 parallel: usize,
1556 delay: Duration,
1557 monitor: Duration,
1558 order: UpdateOrder,
1559 oci: bool,
1560 probe: Option<&HealthProbe>,
1561 uc: &UpdateConfig,
1562 ) -> Result<bool> {
1563 for (n, batch) in pending.chunks(parallel).enumerate() {
1564 if n > 0 && !delay.is_zero() {
1565 std::thread::sleep(delay);
1566 }
1567 if self.superseded(def) {
1568 return Ok(false);
1569 }
1570 for (slot, old) in batch {
1571 let r = self.replace(
1572 def,
1573 spec,
1574 rev,
1575 *slot,
1576 old.as_deref(),
1577 order,
1578 oci,
1579 probe,
1580 monitor,
1581 );
1582 if let Err(e) = r {
1583 let msg = format!("slot {slot}: {e}");
1584 self.event("error", None, &msg);
1585 if old.is_none() {
1586 let prev = self.create_backoff.map(|(_, w)| w);
1588 let (wait, m) = super::failure::retry(prev, &spec.image, &msg, &e);
1589 self.create_backoff = Some((Instant::now(), wait));
1590 self.state = "failing".into();
1591 self.message = Some(m);
1592 self.publish_status(def);
1593 return Ok(false);
1594 }
1595 match uc.failure_action.unwrap_or_default() {
1596 FailureAction::Continue => continue,
1597 FailureAction::Pause => {
1598 self.event("warn", None, &format!("rollout of rev {rev} paused"));
1599 self.paused = Some((rev.clone(), format!("rollout paused: {msg}")));
1600 return Ok(false);
1601 }
1602 FailureAction::Rollback => {
1603 self.paused = Some((rev.clone(), format!("rolled back: {msg}")));
1604 let ctl = Controller {
1605 inner: self.inner.clone(),
1606 };
1607 self.event(
1608 "warn",
1609 None,
1610 &format!("rollout of rev {rev} failed; rolling back"),
1611 );
1612 if let Err(e) = ctl.rollback(&self.q) {
1613 self.event("error", None, &format!("rollback failed: {e}"));
1614 }
1615 return Ok(false);
1616 }
1617 }
1618 }
1619 }
1620 }
1621 Ok(true)
1622 }
1623
1624 fn slot_state(
1626 &mut self,
1627 def: &StackDef,
1628 slot: u32,
1629 old_state: Option<&str>,
1630 new: Option<&str>,
1631 new_state: Option<&str>,
1632 ) {
1633 if let Some(ro) = &mut self.rollout {
1634 if let Some(s) = ro.slots.iter_mut().find(|s| s.slot == slot) {
1635 if let Some(o) = old_state {
1636 s.old_state = o.into();
1637 }
1638 if let Some(n) = new {
1639 s.new = Some(n.into());
1640 }
1641 if let Some(n) = new_state {
1642 s.new_state = n.into();
1643 if n == "serving" {
1644 ro.done += 1;
1645 }
1646 }
1647 }
1648 }
1649 self.publish_status(def);
1650 }
1651
1652 fn waiting_for(&self, spec: &SandboxSpec) -> Option<String> {
1654 let st = self.inner.status.lock().unwrap();
1655 for (dep, d) in &spec.depends_on {
1656 let s = st.get(&(self.q.clone(), dep.clone()));
1657 let ok = match d.condition {
1658 DependCondition::ServiceStarted => s.is_some_and(|s| s.running > 0),
1659 DependCondition::ServiceHealthy => s.is_some_and(|s| s.healthy > 0),
1660 };
1661 if !ok {
1662 return Some(format!("waiting for {dep} ({:?})", d.condition));
1663 }
1664 }
1665 None
1666 }
1667
1668 fn handle(&self, name: &str) -> Sandbox {
1669 let (_, d) = self.template.as_ref().expect("resolved in pass");
1670 Sandbox::like(self.client(), name, d)
1671 }
1672
1673 #[expect(
1676 clippy::too_many_lines,
1677 reason = "predates the lint ratchet; split it when next changed"
1678 )]
1679 fn maintain(
1680 &mut self,
1681 def: &StackDef,
1682 i: &Inst,
1683 spec: &SandboxSpec,
1684 oci: bool,
1685 probe: Option<&HealthProbe>,
1686 ) -> Result<()> {
1687 let policy = spec.deploy.as_ref().and_then(|d| d.restart_policy.clone());
1688 let condition = policy
1689 .as_ref()
1690 .and_then(|p| p.condition)
1691 .unwrap_or_default();
1692 let sb = self.handle(&i.name);
1693 if !i.running() {
1694 self.set_rotation(&i.name, false);
1695 if condition == RestartCondition::None {
1696 return Ok(());
1697 }
1698 if self.restart_budget_spent(&i.name, policy.as_ref()) {
1699 self.message = Some(format!("{}: restart limit reached", i.name));
1700 return Ok(());
1701 }
1702 self.event(
1703 "warn",
1704 Some(&i.name),
1705 &format!("{} is {}; starting it", i.name, i.status),
1706 );
1707 self.count_restart(&i.name);
1708 if let Err(e) = sb.start() {
1709 self.log(&format!("cannot start {}: {e}", i.name));
1710 return Ok(());
1711 }
1712 }
1713 let (pid, ip) = match instance_state(self.client(), &i.name) {
1714 Ok(s) => s,
1715 Err(e) if e.is_not_found() => return Ok(()),
1716 Err(e) => return Err(e),
1717 };
1718 let rt = self.rt.entry(i.name.clone()).or_default();
1719 rt.ip = ip;
1720 if rt.pid != pid {
1721 let restarted = rt.pid != 0;
1724 rt.pid = pid;
1725 rt.started(Instant::now());
1726 rt.next_probe = None;
1727 let r = self.setup(def, &sb, spec, oci);
1728 if let Err(e) = r {
1729 if let Some(rt) = self.rt.get_mut(&i.name) {
1731 rt.pid = 0;
1732 }
1733 self.set_rotation(&i.name, false);
1734 return Err(Error::OperationFailed {
1735 step: format!("set up {}", i.name),
1736 message: e.to_string(),
1737 });
1738 }
1739 if restarted {
1741 self.mark_started(def, &i.name)?;
1742 }
1743 }
1744 let alive = self.alive(&sb, spec, oci);
1745 let healthy = match probe {
1746 None => alive,
1747 Some(p) => {
1748 let rt = self.rt.get_mut(&i.name).unwrap();
1749 if rt.since.is_none() {
1750 rt.since = Some(Instant::now());
1751 }
1752 if rt.next_probe.is_none_or(|t| Instant::now() >= t) {
1753 let r = supervise::probe(&sb, p);
1754 let rt = self.rt.get_mut(&i.name).unwrap();
1755 rt.last_probe = r.output.clone();
1756 rt.record_probe(r.ok, p, Instant::now());
1757 let wait = if rt.healthy.is_none() {
1758 p.start_interval
1759 } else {
1760 p.interval
1761 };
1762 rt.next_probe = Some(Instant::now() + wait);
1763 }
1764 self.rt[&i.name].healthy == Some(true) && alive
1765 }
1766 };
1767 self.set_rotation(&i.name, healthy);
1768 let unhealthy = self.rt[&i.name].healthy == Some(false);
1769 if unhealthy && condition != RestartCondition::None {
1770 if self.restart_budget_spent(&i.name, policy.as_ref()) {
1771 self.message = Some(format!("{}: unhealthy, restart limit reached", i.name));
1772 return Ok(());
1773 }
1774 let n = self.rt[&i.name].unhealthy_restarts;
1775 if n >= RESTARTS_BEFORE_REPLACE {
1776 self.log(&format!(
1777 "{} stayed unhealthy through {n} restarts; replacing it",
1778 i.name
1779 ));
1780 let rt = self.rt.get_mut(&i.name).unwrap();
1781 rt.unhealthy_restarts = 0;
1782 drop(sb);
1783 self.retire(&i.name)?;
1784 return Ok(());
1785 }
1786 self.event(
1787 "warn",
1788 Some(&i.name),
1789 &format!(
1790 "{} is unhealthy ({}); restarting its app",
1791 i.name,
1792 self.rt[&i.name]
1793 .last_probe
1794 .lines()
1795 .last()
1796 .unwrap_or("probe failed")
1797 ),
1798 );
1799 self.count_restart(&i.name);
1800 let rt = self.rt.get_mut(&i.name).unwrap();
1801 rt.unhealthy_restarts += 1;
1802 rt.started(Instant::now());
1803 supervise::restart_app(&sb, &self.service, oci)?;
1804 if !oci {
1805 self.mark_started(def, &i.name)?;
1808 }
1809 }
1810 Ok(())
1811 }
1812
1813 fn setup(&self, def: &StackDef, sb: &Sandbox, spec: &SandboxSpec, oci: bool) -> Result<()> {
1815 let keys = spec.secret_keys();
1817 let values = if keys.is_empty() {
1818 BTreeMap::new()
1819 } else {
1820 super::secrets::values(&self.inner.secrets, &def.org, &def.secrets, keys)?
1821 };
1822 let none = |k: &str| def.on_change(&self.service, k) == OnChange::None;
1825 let pushed = supervise::push_secrets_detailed(sb, spec, &values)?;
1828 if oci && (pushed.missing || pushed.changed.iter().any(|k| !none(k))) {
1829 supervise::restart_app(sb, &self.service, oci)?;
1830 }
1831 if spec.command.is_some() && !oci {
1832 let mut s = spec.clone();
1833 s.restart = Some(RestartMode::Always);
1834 let env = supervise::secret_env(spec, &values)?;
1835 let env_restarts =
1838 spec.env.secrets.is_empty() || spec.env.secrets.values().any(|k| !none(k));
1839 supervise::install_with(
1840 sb,
1841 &self.service,
1842 &s,
1843 spec.has_secret_files(),
1844 &env,
1845 env_restarts,
1846 )?;
1847 }
1848 Ok(())
1849 }
1850
1851 fn oci_secret_env(
1853 &self,
1854 def: &StackDef,
1855 spec: &SandboxSpec,
1856 ) -> Result<BTreeMap<String, String>> {
1857 if spec.env.secrets.is_empty() {
1858 return Ok(BTreeMap::new());
1859 }
1860 let values = super::secrets::values(
1861 &self.inner.secrets,
1862 &def.org,
1863 &def.secrets,
1864 spec.env.secrets.values().map(String::as_str),
1865 )?;
1866 supervise::secret_env(spec, &values)
1867 }
1868
1869 fn alive(&self, sb: &Sandbox, spec: &SandboxSpec, oci: bool) -> bool {
1872 if spec.command.is_some() && !oci {
1873 return supervise::unit_state(sb, &self.service).is_ok_and(|s| s == "active");
1874 }
1875 sb.info()
1876 .is_ok_and(|i| i.status.eq_ignore_ascii_case("running"))
1877 }
1878
1879 fn count_restart(&mut self, name: &str) {
1880 self.rt
1881 .entry(name.to_string())
1882 .or_default()
1883 .restarts
1884 .push_back(Instant::now());
1885 }
1886
1887 fn restart_budget_spent(
1888 &mut self,
1889 name: &str,
1890 policy: Option<&crate::spec::RestartPolicy>,
1891 ) -> bool {
1892 let Some(p) = policy else { return false };
1893 let Some(max) = p.max_attempts else {
1894 return false;
1895 };
1896 let window = p
1897 .window
1898 .as_deref()
1899 .and_then(|w| crate::flex::parse_duration(w).ok());
1900 let rt = self.rt.entry(name.to_string()).or_default();
1901 if let Some(w) = window {
1902 while rt.restarts.front().is_some_and(|t| t.elapsed() > w) {
1903 rt.restarts.pop_front();
1904 }
1905 }
1906 rt.restarts.len() as u32 >= max
1907 }
1908
1909 #[expect(clippy::too_many_arguments)]
1911 fn replace(
1912 &mut self,
1913 def: &StackDef,
1914 spec: &SandboxSpec,
1915 rev: &str,
1916 slot: u32,
1917 old: Option<&str>,
1918 order: UpdateOrder,
1919 oci: bool,
1920 probe: Option<&HealthProbe>,
1921 monitor: Duration,
1922 ) -> Result<()> {
1923 let secret_env = if oci {
1926 self.oci_secret_env(def, spec)?
1927 } else {
1928 BTreeMap::new()
1929 };
1930 let before_start = if oci && spec.has_secret_files() {
1933 let values = super::secrets::values(
1934 &self.inner.secrets,
1935 &def.org,
1936 &def.secrets,
1937 spec.secret_keys(),
1938 )?;
1939 let spec = spec.clone();
1940 Some(crate::plan::BeforeStart(Arc::new(move |c, n| {
1941 supervise::push_secret_files(c, n, &spec, &values).map(|_| ())
1942 })))
1943 } else {
1944 None
1945 };
1946 if let (Some(o), UpdateOrder::StopFirst) = (old, order) {
1947 self.log(&format!("slot {slot}: replacing {o} (stop-first)"));
1948 self.slot_state(def, slot, Some("draining"), None, None);
1949 self.retire(o)?;
1950 self.slot_state(def, slot, Some("retired"), None, None);
1951 }
1952 let name = instance_name(&self.stack, &self.service, slot, &new_id())?;
1953 self.log(&format!("slot {slot}: creating {name} (rev {rev})"));
1954 self.slot_state(def, slot, None, Some(&name), Some("creating"));
1955 let mut s = instance_spec(def, &self.service, spec, slot, rev)?;
1956 s.name = Some(name.clone());
1957 s.env.vars.extend(secret_env);
1959 let mut d = crate::sandbox::resolve(self.client(), &s, &def.file.volumes, &def.base_dir)?;
1960 d.before_start = before_start;
1961 let stack = self.q.clone();
1962 let mut report = |m: &str| eprintln!("isb serve: {stack}: {m}");
1963 let created =
1964 crate::sandbox::ensure(self.client(), &d, EnsureOptions::default(), &mut report);
1965 let result = created.and_then(|_| {
1966 let inst = Inst {
1967 name: name.clone(),
1968 slot,
1969 rev: rev.to_string(),
1970 status: "Running".into(),
1971 secrets: Some(live_versions(def, &self.service)),
1972 };
1973 self.insts.push(inst.clone());
1974 self.slot_state(def, slot, None, None, Some("probing"));
1975 self.wait_serving(def, &inst, spec, oci, probe, monitor)
1976 });
1977 if let Err(e) = result {
1978 let (e, attempt) =
1980 super::failure::explain(self.client(), &name, &self.service, oci, e, now_ms());
1981 self.inner.failures.record(&self.q, &self.service, attempt);
1982 self.event(
1983 "error",
1984 Some(&name),
1985 &format!("{name} did not come up: {e}"),
1986 );
1987 self.slot_state(def, slot, None, None, Some("failed"));
1988 let _ = self.retire(&name);
1989 return Err(e);
1990 }
1991 if let (Some(o), UpdateOrder::StartFirst) = (old, order) {
1992 self.log(&format!("slot {slot}: {name} is serving; retiring {o}"));
1993 self.slot_state(def, slot, Some("draining"), None, None);
1994 self.retire(o)?;
1995 self.slot_state(def, slot, Some("retired"), None, None);
1996 }
1997 self.slot_state(def, slot, None, None, Some("serving"));
1998 Ok(())
1999 }
2000
2001 fn wait_serving(
2005 &mut self,
2006 def: &StackDef,
2007 i: &Inst,
2008 spec: &SandboxSpec,
2009 oci: bool,
2010 probe: Option<&HealthProbe>,
2011 monitor: Duration,
2012 ) -> Result<()> {
2013 let deadline = match probe {
2016 Some(p) => {
2017 p.startup_grace.max(p.start_period)
2018 + p.interval * p.retries
2019 + Duration::from_secs(30)
2020 }
2021 None => Duration::from_secs(60),
2022 }
2023 .max(Duration::from_secs(60));
2024 let started = Instant::now();
2025 loop {
2026 self.maintain(def, i, spec, oci, probe)?;
2027 if self.rt.get(&i.name).is_some_and(|r| r.in_rotation) {
2028 break;
2029 }
2030 if self
2031 .rt
2032 .get(&i.name)
2033 .is_some_and(|r| r.healthy == Some(false))
2034 {
2035 return Err(Error::invalid(format!(
2036 "unhealthy: {}",
2037 self.rt[&i.name].last_probe
2038 )));
2039 }
2040 if started.elapsed() > deadline {
2041 let why = self
2042 .rt
2043 .get(&i.name)
2044 .map(|r| r.last_probe.clone())
2045 .filter(|s| !s.is_empty())
2046 .unwrap_or_else(|| "its app is not running".into());
2047 return Err(Error::invalid(format!(
2048 "not serving after {:?}: {why}",
2049 started.elapsed()
2050 )));
2051 }
2052 if let Some(rt) = self.rt.get_mut(&i.name) {
2054 rt.next_probe = None;
2055 }
2056 std::thread::sleep(
2057 probe
2058 .map(|p| p.start_interval)
2059 .unwrap_or(Duration::from_secs(1))
2060 .min(Duration::from_secs(2)),
2061 );
2062 }
2063 self.sync_routes();
2064 self.slot_state(def, i.slot, None, None, Some("monitoring"));
2065 let watch = Instant::now();
2066 while watch.elapsed() < monitor {
2067 std::thread::sleep(Duration::from_secs(1).min(monitor));
2068 if let Some(rt) = self.rt.get_mut(&i.name) {
2069 rt.next_probe = None;
2070 }
2071 self.maintain(def, i, spec, oci, probe)?;
2072 if !self.rt.get(&i.name).is_some_and(|r| r.in_rotation) {
2073 return Err(Error::invalid(format!(
2074 "failed within the {monitor:?} monitor period"
2075 )));
2076 }
2077 }
2078 Ok(())
2079 }
2080
2081 fn retire(&mut self, name: &str) -> Result<()> {
2084 self.drain(name);
2085 if let Ok(sb) = Sandbox::get(self.client(), name) {
2086 let _ = sb.stop(false, Duration::from_secs(10));
2087 }
2088 self.rt.remove(name);
2089 self.insts.retain(|i| i.name != name);
2090 match Sandbox::remove(self.client(), name, true) {
2091 Err(e) if !e.is_not_found() => Err(e),
2092 _ => Ok(()),
2093 }
2094 }
2095
2096 fn drain(&mut self, name: &str) {
2099 let ip = self.rt.get(name).and_then(|r| r.ip);
2100 self.set_rotation(name, false);
2101 self.sync_routes();
2102 if let Some(ip) = ip {
2103 let started = Instant::now();
2104 if let Some(o) = &self.inner.observer {
2105 o.drain(&self.q, &self.service, ip, DRAIN);
2106 }
2107 let left = DRAIN.saturating_sub(started.elapsed());
2108 for (k, p) in &self.routes {
2109 self.inner
2110 .balancer
2111 .wait_drained(k, SocketAddr::new(ip, p.target), left);
2112 }
2113 }
2114 }
2115
2116 fn set_rotation(&mut self, name: &str, on: bool) {
2117 let rt = self.rt.entry(name.to_string()).or_default();
2118 if rt.in_rotation != on {
2119 rt.in_rotation = on;
2120 self.sync_routes();
2121 }
2122 }
2123
2124 fn set_routes(&mut self, spec: &SandboxSpec) {
2126 let want: BTreeMap<String, Published> = match published(spec) {
2127 Ok(ps) => ps
2128 .into_iter()
2129 .map(|p| (format!("{}/{}/{}", self.q, self.service, p.display()), p))
2130 .collect(),
2131 Err(e) => {
2132 self.message = Some(e.to_string());
2133 BTreeMap::new()
2134 }
2135 };
2136 let stale: Vec<String> = self
2137 .routes
2138 .keys()
2139 .filter(|k| !want.contains_key(*k))
2140 .cloned()
2141 .collect();
2142 for k in stale {
2143 self.inner.balancer.remove_route(&k);
2144 self.routes.remove(&k);
2145 self.route_errors.remove(&k);
2146 }
2147 for (k, p) in want {
2148 self.routes.entry(k).or_insert(p);
2149 }
2150 self.sync_routes();
2151 }
2152
2153 fn sync_routes(&mut self) {
2154 let mut errors = BTreeMap::new();
2155 for (k, p) in self.routes.iter().filter(|(_, p)| !p.udp) {
2157 let backends: Vec<SocketAddr> = self
2158 .rt
2159 .values()
2160 .filter(|r| r.in_rotation)
2161 .filter_map(|r| r.ip)
2162 .map(|ip| SocketAddr::new(ip, p.target))
2163 .collect();
2164 if let Err(e) = self.inner.balancer.set_route(k, p.listen, backends) {
2165 errors.insert(k.clone(), e.to_string());
2166 }
2167 }
2168 for (k, e) in &errors {
2169 if self.route_errors.get(k) != Some(e) {
2170 self.log(&format!("cannot publish {k}: {e}"));
2171 }
2172 }
2173 self.route_errors = errors;
2174 self.sync_observer();
2175 self.sync_dns();
2176 }
2177
2178 fn sync_observer(&mut self) {
2180 let Some(o) = self.inner.observer.clone() else {
2181 return;
2182 };
2183 let mut ips: Vec<IpAddr> = self
2184 .rt
2185 .values()
2186 .filter(|r| r.in_rotation)
2187 .filter_map(|r| r.ip)
2188 .collect();
2189 ips.sort();
2190 if self.observed.as_ref() != Some(&ips) {
2191 o.rotation(&self.q, &self.service, &ips);
2192 self.observed = Some(ips);
2193 }
2194 }
2195
2196 #[expect(
2197 clippy::too_many_lines,
2198 reason = "predates the lint ratchet; split it when next changed"
2199 )]
2200 fn publish_status(&mut self, def: &StackDef) {
2201 self.sync_observer();
2203 self.sync_dns();
2204 let Ok(spec) = def.service(&self.service) else {
2205 return;
2206 };
2207 let rev = def.revision(&self.service).unwrap_or_default();
2208 let probe = matches!(spec.health_probe(), Ok(Some(_)));
2209 let live = live_versions(def, &self.service);
2210 let snap = self.inner.snapshot.lock().unwrap().instances.clone();
2211 let instances: Vec<InstanceStatus> = self
2212 .insts
2213 .iter()
2214 .map(|i| {
2215 let rt = self.rt.get(&i.name);
2216 let m = snap.get(&format!("{}/{}", self.oclient.project_name(), i.name));
2217 let health = match (probe, rt.and_then(|r| r.healthy)) {
2218 (false, _) => "none",
2219 (true, Some(true)) => "healthy",
2220 (true, Some(false)) => "unhealthy",
2221 (true, None) => "starting",
2222 };
2223 InstanceStatus {
2224 name: i.name.clone(),
2225 slot: i.slot,
2226 rev: i.rev.clone(),
2227 status: m
2228 .map(|m| m.status.clone())
2229 .unwrap_or_else(|| i.status.clone()),
2230 health: health.into(),
2231 ip: rt.and_then(|r| r.ip).map(|ip| ip.to_string()),
2232 in_rotation: rt.is_some_and(|r| r.in_rotation),
2233 restarts: rt.map(|r| r.restarts.len() as u32).unwrap_or(0),
2234 last_probe: rt.map(|r| r.last_probe.clone()).unwrap_or_default(),
2235 cpu_pct: m.and_then(|m| m.cpu_pct),
2236 cpu_history: m.map(|m| m.cpu_history.clone()).unwrap_or_default(),
2237 mem_bytes: m.and_then(|m| m.mem_bytes),
2238 disk_bytes: m.and_then(|m| m.disk_bytes),
2239 stale_secrets: i
2240 .secrets
2241 .as_ref()
2242 .map(|have| super::secrets::stale(have, &live))
2243 .unwrap_or_default(),
2244 }
2245 })
2246 .collect();
2247 let mut instances = instances;
2248 instances.sort_by(|a, b| (a.slot, &a.name).cmp(&(b.slot, &b.name)));
2249 let routes = self.inner.balancer.routes();
2250 let now = Instant::now();
2251 let mut ports = Vec::new();
2252 for (k, p) in &self.routes {
2253 let r = routes.iter().find(|r| r.key == *k);
2254 let accepted = r.map(|r| r.accepted).unwrap_or(0);
2255 let e = self
2256 .rates
2257 .entry(k.clone())
2258 .or_insert_with(|| (accepted, now, VecDeque::new()));
2259 let dt = now.duration_since(e.1).as_secs_f32();
2260 if dt >= 1.0 {
2263 let rate = accepted.saturating_sub(e.0) as f32 / dt;
2264 if e.2.len() == crate::metrics::HISTORY {
2265 e.2.pop_front();
2266 }
2267 e.2.push_back(rate);
2268 e.0 = accepted;
2269 e.1 = now;
2270 }
2271 ports.push(PortStatus {
2272 listen: p.display(),
2273 target: p.target,
2274 backends: match r {
2275 _ if p.udp => {
2276 let ips = self.rt.values().filter_map(|r| r.ip);
2277 ips.map(|ip| SocketAddr::new(ip, p.target).to_string())
2278 .collect()
2279 }
2280 Some(r) => r.backends.iter().map(|b| b.addr.to_string()).collect(),
2281 None => Vec::new(),
2282 },
2283 error: self.route_errors.get(k).cloned(),
2284 accepted,
2285 active: r
2286 .map(|r| {
2287 r.backends
2288 .iter()
2289 .chain(r.draining.iter())
2290 .map(|b| b.active)
2291 .sum()
2292 })
2293 .unwrap_or(0),
2294 rate_history: e.2.iter().copied().collect(),
2295 });
2296 }
2297 self.rates.retain(|k, _| self.routes.contains_key(k));
2298 let running = instances
2299 .iter()
2300 .filter(|i| i.status.eq_ignore_ascii_case("running"))
2301 .count() as u32;
2302 let healthy = instances.iter().filter(|i| i.in_rotation).count() as u32;
2303 self.watch_health(healthy, spec.replicas());
2304 let st = ServiceStatus {
2305 service: self.service.clone(),
2306 image: spec.image.clone(),
2307 rev,
2308 replicas: spec.replicas(),
2309 running,
2310 healthy,
2311 state: self.state.clone(),
2312 message: self.message.clone(),
2313 instances,
2314 ports,
2315 rollout: self.rollout.clone(),
2316 checked_at: now_secs(),
2317 domains: Vec::new(),
2318 last_failed_attempt: None,
2319 };
2320 self.inner.status.lock().unwrap().insert(self.key(), st);
2321 }
2322}
2323
2324pub fn unsettled(st: &StackStatus) -> BTreeSet<String> {
2327 st.services
2328 .iter()
2329 .filter(|s| s.state != "converged")
2330 .map(|s| s.service.clone())
2331 .collect()
2332}
2333
2334#[path = "rotation.rs"]
2335mod rotation;
2336
2337#[cfg(test)]
2338#[path = "controller_tests.rs"]
2339mod tests;