1use super::*;
12
13impl Controller {
14 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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, ¤t, spec, oci, probe.as_ref(), monitor)
594 }
595
596 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}