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