Skip to main content

isb_core/stack/
rotation.rs

1//! What a new secret version does to the stack services using it.
2//!
3//! The controller notices a new version (`isb secret set`, a forced
4//! refresh, or its polling round over driver-backed bindings, which asks
5//! each driver once per org for all its names), saves it in the stack's
6//! bindings and says per service what happens, by the service's
7//! `on_change`: `roll` is a new revision (the ordinary rolling update),
8//! `restart` and `none` are the worker's to do on the replicas it has
9//! ([`Worker::cycle_in_place`]).
10
11use super::*;
12
13impl Controller {
14    /// A stored secret got a new value (`isb secret set`): every stack in
15    /// the org bound to it moves to the new version, and each service using
16    /// it rolls, restarts in place or is left stale, per its `on_change`.
17    /// Returns what happened, per service.
18    pub fn secret_changed(&self, org: &OrgId, name: &str) -> Vec<Cycle> {
19        let mut out = Vec::new();
20        for def in self.definitions() {
21            if def.org != *org {
22                continue;
23            }
24            let keys: Vec<String> = def
25                .secrets
26                .iter()
27                .filter(|(_, b)| b.name == name)
28                .map(|(k, _)| k.clone())
29                .collect();
30            if !keys.is_empty() {
31                out.extend(self.check_bindings(&def.qualified(), &keys, false, None));
32            }
33        }
34        out
35    }
36
37    /// Re-read every binding to `name` in the org from its driver now
38    /// (`isb secret refresh`): once per driver, however many stacks use it.
39    /// Returns each driver and version found, and what the new version did.
40    pub fn refresh_secret(
41        &self,
42        org: &OrgId,
43        name: &str,
44    ) -> Result<crate::stack::secrets::Refreshed> {
45        let mut found: Vec<(String, u64)> = Vec::new();
46        let mut cycles = Vec::new();
47        let mut read = crate::stack::secrets::Polled::new();
48        for def in self.definitions() {
49            if def.org != *org {
50                continue;
51            }
52            let q = def.qualified();
53            let keys: Vec<String> = def
54                .secrets
55                .iter()
56                .filter(|(_, b)| b.name == name)
57                .map(|(k, _)| k.clone())
58                .collect();
59            for k in &keys {
60                let b = &def.secrets[k];
61                let id = (org.clone(), b.driver.clone(), name.to_string());
62                if let std::collections::btree_map::Entry::Vacant(e) = read.entry(id) {
63                    let v = self.inner.secrets.refresh_in(&b.driver, org, name)?;
64                    found.push((b.driver.clone(), v));
65                    e.insert(Ok(v));
66                }
67                let every = def
68                    .file
69                    .secrets
70                    .get(k)
71                    .map(crate::spec::SecretDef::refresh_interval)
72                    .unwrap_or(crate::spec::DEFAULT_SECRET_REFRESH);
73                self.inner
74                    .refresh
75                    .lock()
76                    .unwrap()
77                    .reset(&q, k, every, Instant::now());
78            }
79            if !keys.is_empty() {
80                cycles.extend(self.check_bindings(&q, &keys, true, Some(&read)));
81            }
82        }
83        Ok((found, cycles))
84    }
85
86    /// The driver-backed bindings whose refresh interval is up.
87    pub(super) fn check_due_secrets(&self) {
88        let defs: Vec<(String, Arc<StackDef>)> = self
89            .inner
90            .stacks
91            .lock()
92            .unwrap()
93            .iter()
94            .map(|(q, d)| (q.clone(), d.clone()))
95            .collect();
96        let due = self
97            .inner
98            .refresh
99            .lock()
100            .unwrap()
101            .due(defs.iter().map(|(q, d)| (q.as_str(), &**d)), Instant::now());
102        self.poll(due);
103    }
104
105    /// Check the given `(stack, key)` bindings against their drivers in one
106    /// round: every driver is asked once per org for all its names (one
107    /// 1Password lookup per item, whatever the number of fields and stacks
108    /// using it), then each stack moves to what changed.
109    pub(super) fn poll(&self, due: Vec<(String, String)>) {
110        if due.is_empty() {
111            return;
112        }
113        let mut by_stack: BTreeMap<String, Vec<String>> = BTreeMap::new();
114        let mut refs = Vec::new();
115        for (q, k) in due {
116            if let Ok(def) = self.get_def(&q) {
117                if let Some(b) = def.secrets.get(&k) {
118                    refs.push((def.org.clone(), b.driver.clone(), b.name.clone()));
119                }
120            }
121            by_stack.entry(q).or_default().push(k);
122        }
123        let versions = crate::stack::secrets::poll_versions(&self.inner.secrets, refs);
124        for (q, keys) in by_stack {
125            self.check_bindings(&q, &keys, false, Some(&versions));
126        }
127    }
128
129    /// Compare the given bindings of a stack with the store's current
130    /// versions (from `known` when given, else asked now); on any change,
131    /// save the new versions, hand the definition to the workers, and say
132    /// what each service using a moved secret does about it. `quiet` skips
133    /// logging lookups that fail (the caller reports them).
134    pub(super) fn check_bindings(
135        &self,
136        q: &str,
137        keys: &[String],
138        quiet: bool,
139        known: Option<&crate::stack::secrets::Polled>,
140    ) -> Vec<Cycle> {
141        let _g = self.inner.edit.lock().unwrap();
142        let Ok(cur) = self.get_def(q) else {
143            return Vec::new();
144        };
145        let mut def = (*cur).clone();
146        let mut moved: Vec<(String, String, u64, u64)> = Vec::new();
147        for k in keys {
148            let Some(b) = def.secrets.get_mut(k) else {
149                continue;
150            };
151            let id = (def.org.clone(), b.driver.clone(), b.name.clone());
152            let v = match known.and_then(|m| m.get(&id)) {
153                Some(r) => r.clone().map_err(Error::invalid),
154                None => self.inner.secrets.version_in(&b.driver, &def.org, &b.name),
155            };
156            match v {
157                Ok(v) if v != b.version => {
158                    moved.push((k.clone(), b.name.clone(), b.version, v));
159                    b.version = v;
160                }
161                Ok(_) => {}
162                // The services keep the value they have.
163                Err(e) if !quiet => self.note(
164                    "warn",
165                    q,
166                    format!("secret {}: cannot check its version: {e}", b.name),
167                ),
168                Err(_) => {}
169            }
170        }
171        if moved.is_empty() {
172            return Vec::new();
173        }
174        // A new version from a driver takes effect where the old one is
175        // stored first (`rotate`); `isb secret set` did that before storing.
176        let mut failed: BTreeMap<String, String> = BTreeMap::new();
177        for (key, _, from, _) in &moved {
178            let b = &def.secrets[key];
179            let has_rotate = def
180                .file
181                .secrets
182                .get(key)
183                .is_some_and(|d| d.rotate.is_some());
184            if !b.is_driver_backed() || !has_rotate {
185                continue;
186            }
187            let r = b
188                .read(&self.inner.secrets, &def.org)
189                .and_then(|v| self.run_rotate(&def, key, &v));
190            if let Err(e) = r {
191                failed.insert(key.clone(), e.to_string());
192                if let Some(b) = def.secrets.get_mut(key) {
193                    b.version = *from;
194                }
195            }
196        }
197        let mut cycles = Vec::new();
198        for (key, secret, from, to) in &moved {
199            for service in def.services_using(key) {
200                cycles.push(Cycle {
201                    stack: def.name.clone(),
202                    action: def.on_change(&service, key),
203                    service,
204                    key: key.clone(),
205                    secret: secret.clone(),
206                    from: *from,
207                    to: *to,
208                    error: failed.get(key).cloned(),
209                });
210            }
211        }
212        if failed.len() < moved.len() {
213            if let Err(e) = self.inner.store.save(&def) {
214                self.note("error", q, format!("cannot save new secret versions: {e}"));
215                return Vec::new();
216            }
217            self.apply(Arc::new(def));
218        }
219        self.report_cycles(q, &cycles);
220        cycles
221    }
222
223    /// `isb secret set NAME`, before the value is stored: run the `rotate`
224    /// command of every stack secret bound to `name` in the org, with the
225    /// new value, so it takes effect where the old one is kept (a database
226    /// user's password) before any replica is given it. The first failure
227    /// is the answer, and the caller stores nothing.
228    pub fn rotate_before_set(
229        &self,
230        org: &OrgId,
231        name: &str,
232        value: &[u8],
233    ) -> Result<Vec<crate::stack::secrets::Applied>> {
234        let mut out = Vec::new();
235        for def in self.definitions() {
236            if def.org != *org {
237                continue;
238            }
239            for (key, b) in &def.secrets {
240                let rotates = def
241                    .file
242                    .secrets
243                    .get(key)
244                    .is_some_and(|d| d.rotate.is_some());
245                if b.name == name && rotates {
246                    out.extend(self.run_rotate(&def, key, value).map_err(|e| {
247                        Error::invalid(format!(
248                            "secret {name} not changed: stack {}: {e}",
249                            def.name
250                        ))
251                    })?);
252                }
253            }
254        }
255        Ok(out)
256    }
257
258    /// Run the top-level secret `key`'s `rotate` command, with `value` on
259    /// stdin, in one running replica of each service using it.
260    fn run_rotate(
261        &self,
262        def: &StackDef,
263        key: &str,
264        value: &[u8],
265    ) -> Result<Vec<crate::stack::secrets::Applied>> {
266        let Some(argv) = def.file.secrets.get(key).and_then(|d| d.rotate.clone()) else {
267            return Ok(Vec::new());
268        };
269        let client = crate::org::client(&self.inner.client, &def.org);
270        let mut out = Vec::new();
271        for service in def.services_using(key) {
272            let rev = def.revision(&service)?;
273            let mut insts: Vec<Inst> = list_instances(&client, &def.name, Some(&service))?
274                .into_iter()
275                .filter(Inst::running)
276                .collect();
277            // The current revision's first, then the lowest slot.
278            insts.sort_by_key(|i| (i.rev != rev, i.slot));
279            let Some(i) = insts.first() else {
280                return Err(Error::invalid(format!(
281                    "rotate: no running replica of {service} to run it in"
282                )));
283            };
284            let o = Sandbox::get(&client, &i.name)?.exec_with(
285                argv.clone(),
286                crate::exec::ExecOptions::default()
287                    .stdin(crate::exec::Stdin::Bytes(value.to_vec()))
288                    .timeout(Duration::from_secs(120)),
289            )?;
290            if !o.success() {
291                let err = o.stderr_text();
292                let out_text = o.stdout_text();
293                let why = [err.trim(), out_text.trim()]
294                    .into_iter()
295                    .find(|s| !s.is_empty())
296                    .unwrap_or("no output");
297                let tail: String = why.lines().rev().take(5).collect::<Vec<_>>().join(" | ");
298                return Err(Error::invalid(format!(
299                    "rotate in {} exited with {:?}: {tail}",
300                    i.name, o.exit_code
301                )));
302            }
303            self.event(
304                "secret.rotated",
305                "info",
306                &def.qualified(),
307                &service,
308                format!(
309                    "secret {key}: the new value took effect in {} (rotate)",
310                    i.name
311                ),
312            );
313            out.push(crate::stack::secrets::Applied {
314                stack: def.name.clone(),
315                service,
316                instance: i.name.clone(),
317            });
318        }
319        Ok(out)
320    }
321
322    /// One `secret.rotated` event per service a new secret version reached,
323    /// saying what it does about it: the history keeps them, and
324    /// notification channels can report them.
325    fn report_cycles(&self, q: &str, cycles: &[Cycle]) {
326        let mut by_service: BTreeMap<&str, Vec<&Cycle>> = BTreeMap::new();
327        for c in cycles {
328            by_service.entry(&c.service).or_default().push(c);
329        }
330        for (service, cs) in by_service {
331            let what: Vec<String> = cs
332                .iter()
333                .map(|c| format!("{} v{} -> v{}", c.secret, c.from, c.to))
334                .collect();
335            let action = cs.iter().map(|c| c.action).max().unwrap_or_default();
336            if let Some(e) = cs.iter().find_map(|c| c.error.as_deref()) {
337                self.event(
338                    "secret.rotated",
339                    "error",
340                    q,
341                    service,
342                    format!(
343                        "new secret version ({}) not taken up: {e}; the service keeps the version it has",
344                        what.join(", ")
345                    ),
346                );
347                continue;
348            }
349            let (level, how) = match action {
350                OnChange::Roll => ("info", "rolling its replicas".to_string()),
351                OnChange::Restart => ("info", "restarting its replicas in place".to_string()),
352                OnChange::None => (
353                    "warn",
354                    format!(
355                        "not cycled (on_change: none): files have the new value, the apps keep v{} until they next start",
356                        cs.iter().map(|c| c.from).min().unwrap_or(0)
357                    ),
358                ),
359            };
360            self.event(
361                "secret.rotated",
362                level,
363                q,
364                service,
365                format!("new secret version ({}): {how}", what.join(", ")),
366            );
367        }
368    }
369}
370
371impl Worker {
372    /// Record on the instance that its app started with the versions bound
373    /// now of the secrets the service takes in place ([`LABEL_SECRETS`]).
374    pub(super) fn mark_started(&mut self, def: &StackDef, name: &str) -> Result<()> {
375        let want = live_versions(def, &self.service);
376        let Some(i) = self.insts.iter().find(|i| i.name == name) else {
377            return Ok(());
378        };
379        if want.is_empty() || i.secrets.as_ref() == Some(&want) {
380            return Ok(());
381        }
382        self.set_secrets_label(name, &want)
383    }
384
385    fn set_secrets_label(&mut self, name: &str, v: &BTreeMap<String, u64>) -> Result<()> {
386        self.client().mutate(
387            "PATCH",
388            &format!("/1.0/instances/{}", encode_segment(name)),
389            Some(&serde_json::json!({"config": {
390                format!("user.{LABEL_SECRETS}"): crate::stack::secrets::versions_label(v),
391            }})),
392            "label secret versions",
393            Duration::from_secs(60),
394        )?;
395        if let Some(i) = self.insts.iter_mut().find(|i| i.name == name) {
396            i.secrets = Some(v.clone());
397        }
398        Ok(())
399    }
400
401    /// Bring the replicas to the versions bound now of the secrets the
402    /// service takes in place: `on_change: restart` restarts each app with
403    /// the new value, in batches of `update_config.parallelism`, each
404    /// drained first and waited on until it serves again; `on_change: none`
405    /// only delivers the files (and the variables a later start reads), and
406    /// the replica stays stale until it next starts.
407    pub(super) fn cycle_in_place(
408        &mut self,
409        def: &Arc<StackDef>,
410        spec: &SandboxSpec,
411        uc: &UpdateConfig,
412    ) -> Result<()> {
413        let live = def.live_secrets(&self.service);
414        if live.is_empty() {
415            self.restart_failed = None;
416            return Ok(());
417        }
418        let want: BTreeMap<String, u64> = live.iter().map(|(k, (_, v))| (k.clone(), *v)).collect();
419        if self
420            .restart_failed
421            .as_ref()
422            .is_some_and(|(w, _)| *w != want)
423        {
424            self.restart_failed = None;
425        }
426        let mut restart = Vec::new();
427        for i in self.insts.clone() {
428            // An instance from before the label, or a secret new to it, is
429            // taken to run what is bound: there is nothing to compare with.
430            let mut have = i.secrets.clone().unwrap_or_default();
431            let mut adopted = i.secrets.is_none();
432            for (k, v) in &want {
433                if !have.contains_key(k) {
434                    have.insert(k.clone(), *v);
435                    adopted = true;
436                }
437            }
438            have.retain(|k, _| want.contains_key(k));
439            if adopted {
440                self.set_secrets_label(&i.name, &have)?;
441            }
442            let stale = crate::stack::secrets::stale(&have, &want);
443            if stale.is_empty() {
444                continue;
445            }
446            if stale.iter().any(|s| live[&s.key].0 == OnChange::Restart) {
447                restart.push(i);
448                continue;
449            }
450            // Only `none` secrets moved: deliver once per version.
451            let delivered = self.rt.get(&i.name).map(|r| &r.delivered);
452            if delivered != Some(&want) {
453                self.deliver_live(def, spec, &i.name)?;
454                let keys: Vec<String> = stale
455                    .iter()
456                    .map(|s| format!("{} v{} -> v{}", s.key, s.running, s.current))
457                    .collect();
458                self.event(
459                    "info",
460                    Some(&i.name),
461                    &format!(
462                        "{}: delivered new secret files ({}); its app keeps the old value until it next starts (on_change: none)",
463                        i.name,
464                        keys.join(", ")
465                    ),
466                );
467                self.rt.entry(i.name.clone()).or_default().delivered = want.clone();
468            }
469        }
470        if restart.is_empty() || self.restart_failed.is_some() {
471            return Ok(());
472        }
473        self.restart_batches(def, spec, uc, &restart, &want)
474    }
475
476    /// Restart `restart` in place, in batches of `update_config.parallelism`
477    /// with its `delay`, each replica watched for `monitor`. The first that
478    /// fails stops the rest until the versions move again.
479    fn restart_batches(
480        &mut self,
481        def: &Arc<StackDef>,
482        spec: &SandboxSpec,
483        uc: &UpdateConfig,
484        restart: &[Inst],
485        want: &BTreeMap<String, u64>,
486    ) -> Result<()> {
487        let parallel = match uc.parallelism.unwrap_or(1) {
488            0 => restart.len(),
489            n => n as usize,
490        };
491        let delay = uc
492            .delay
493            .as_deref()
494            .map(crate::flex::parse_duration)
495            .transpose()
496            .map_err(Error::invalid)?
497            .unwrap_or_default();
498        let monitor = uc
499            .monitor
500            .as_deref()
501            .map(crate::flex::parse_duration)
502            .transpose()
503            .map_err(Error::invalid)?
504            .unwrap_or(Duration::from_secs(5));
505        self.state = "updating".into();
506        self.message = Some(format!(
507            "restarting {} replica(s) in place for new secret versions",
508            restart.len()
509        ));
510        self.publish_status(def);
511        let started = Instant::now();
512        for (n, batch) in restart.chunks(parallel).enumerate() {
513            if n > 0 && !delay.is_zero() {
514                std::thread::sleep(delay);
515            }
516            if self.superseded(def) {
517                return Ok(());
518            }
519            for i in batch {
520                if let Err(e) = self.restart_in_place(def, spec, monitor, i, want) {
521                    let msg = format!(
522                        "{}: restart for new secret versions failed: {e}; the other replicas keep the old value until the next version or deploy",
523                        i.name
524                    );
525                    self.event("error", Some(&i.name), &msg);
526                    self.restart_failed = Some((want.clone(), msg.clone()));
527                    self.message = Some(msg);
528                    self.publish_status(def);
529                    return Ok(());
530                }
531            }
532        }
533        self.event(
534            "info",
535            None,
536            &format!(
537                "restarted {} replica(s) in place for new secret versions in {:.0?}",
538                restart.len(),
539                started.elapsed()
540            ),
541        );
542        self.message = None;
543        Ok(())
544    }
545
546    /// Drain one replica, give it the values bound now (files, variables),
547    /// restart its app, and wait until it serves again.
548    fn restart_in_place(
549        &mut self,
550        def: &StackDef,
551        spec: &SandboxSpec,
552        monitor: Duration,
553        i: &Inst,
554        want: &BTreeMap<String, u64>,
555    ) -> Result<()> {
556        let oci = crate::plan::ImageSource::parse(&spec.image)?.is_oci();
557        let probe = spec.health_probe().map_err(Error::invalid)?;
558        // Read first: a secret that cannot be read leaves it serving.
559        let keys = spec.secret_keys();
560        let values =
561            crate::stack::secrets::values(&self.inner.secrets, &def.org, &def.secrets, keys)?;
562        self.log(&format!(
563            "{}: restarting in place for new secret versions",
564            i.name
565        ));
566        self.drain(&i.name);
567        let sb = self.handle(&i.name);
568        supervise::push_secrets(&sb, spec, &values)?;
569        let env = supervise::secret_env(spec, &values)?;
570        if oci {
571            supervise::set_oci_env(&sb, &env)?;
572            supervise::restart_app(&sb, &self.service, oci)?;
573        } else if spec.command.is_some() {
574            let mut s = spec.clone();
575            s.restart = Some(RestartMode::Always);
576            if !supervise::install(&sb, &self.service, &s, spec.has_secret_files(), &env)? {
577                supervise::restart_app(&sb, &self.service, oci)?;
578            }
579        }
580        self.set_secrets_label(&i.name, want)?;
581        let rt = self.rt.entry(i.name.clone()).or_default();
582        rt.failures = 0;
583        rt.healthy = None;
584        rt.next_probe = None;
585        rt.since = Some(Instant::now());
586        rt.delivered = want.clone();
587        let current = self
588            .insts
589            .iter()
590            .find(|x| x.name == i.name)
591            .cloned()
592            .unwrap_or_else(|| i.clone());
593        self.wait_serving(def, &current, spec, oci, probe.as_ref(), monitor)
594    }
595
596    /// Put the values bound now where a running replica can take them
597    /// without a restart: its secret files, and the variables its next
598    /// start reads (the unit's environment file, or an OCI instance's
599    /// config).
600    fn deliver_live(&self, def: &StackDef, spec: &SandboxSpec, name: &str) -> Result<()> {
601        let oci = crate::plan::ImageSource::parse(&spec.image)?.is_oci();
602        let keys = spec.secret_keys();
603        let values =
604            crate::stack::secrets::values(&self.inner.secrets, &def.org, &def.secrets, keys)?;
605        let sb = self.handle(name);
606        supervise::push_secrets(&sb, spec, &values)?;
607        let env = supervise::secret_env(spec, &values)?;
608        if oci {
609            supervise::set_oci_env(&sb, &env)?;
610        } else if spec.command.is_some() {
611            let mut s = spec.clone();
612            s.restart = Some(RestartMode::Always);
613            supervise::install_with(&sb, &self.service, &s, spec.has_secret_files(), &env, false)?;
614        }
615        Ok(())
616    }
617}