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