1use std::collections::BTreeMap;
14use std::sync::Mutex;
15use std::time::Duration;
16
17use serde::{Deserialize, Serialize};
18
19use super::controller::{Inst, list_instances};
20use crate::client::Client;
21use crate::error::{Error, Result};
22use crate::sandbox::Sandbox;
23use crate::supervise;
24
25const READ_LINES: usize = 200;
27pub const OUTPUT_BYTES: usize = 32 * 1024;
30const MESSAGE_LINES: usize = 8;
32const MESSAGE_CHARS: usize = 1200;
34
35pub const IMAGE_RETRY: Duration = Duration::from_secs(300);
37pub const IMAGE_RETRY_MAX: Duration = Duration::from_secs(3600);
39
40pub fn image_missing(image: &str, err: &str) -> Option<String> {
45 if !err.contains("Failed getting remote image") && !err.contains("Error parsing image name") {
46 return None;
47 }
48 match crate::image_check::classify(err) {
49 crate::image_check::Probe::NotFound(why) => {
50 Some(format!("image {image} not found ({why})"))
51 }
52 crate::image_check::Probe::Denied(why) => Some(format!(
53 "image {image} cannot be pulled without credentials ({why}): it is private or does not exist"
54 )),
55 _ => None,
56 }
57}
58
59pub fn retry(prev: Option<Duration>, image: &str, msg: &str, e: &Error) -> (Duration, String) {
65 let missing = image_missing(image, &e.to_string());
66 let (first, cap) = match missing {
67 Some(_) => (IMAGE_RETRY, IMAGE_RETRY_MAX),
68 None => (Duration::from_secs(10), Duration::from_secs(300)),
69 };
70 let wait = prev.map(|w| (w * 2).clamp(first, cap)).unwrap_or(first);
71 let message = match missing {
72 Some(m) => format!(
73 "{m}: change the image and deploy again (retrying in {}m)",
74 wait.as_secs() / 60
75 ),
76 None => format!("{msg}; retrying in {wait:?}"),
77 };
78 (wait, message)
79}
80
81#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
83pub struct FailedAttempt {
84 pub instance: String,
86 pub at_ms: u64,
88 pub reason: String,
90 pub output: String,
93 #[serde(default, skip_serializing_if = "Option::is_none")]
95 pub output_note: Option<String>,
96}
97
98impl FailedAttempt {
99 pub fn summary(&self) -> serde_json::Value {
101 serde_json::json!({
102 "instance": self.instance,
103 "at_ms": self.at_ms,
104 "reason": self.reason,
105 "output_lines": self.output.lines().count(),
106 })
107 }
108}
109
110#[derive(Default)]
112pub struct Failures(Mutex<BTreeMap<(String, String), FailedAttempt>>);
113
114impl Failures {
115 pub fn record(&self, stack: &str, service: &str, a: FailedAttempt) {
116 self.0
117 .lock()
118 .unwrap_or_else(|e| e.into_inner())
119 .insert((stack.to_string(), service.to_string()), a);
120 }
121
122 pub fn last(&self, stack: &str, service: &str) -> Option<FailedAttempt> {
123 self.0
124 .lock()
125 .unwrap_or_else(|e| e.into_inner())
126 .get(&(stack.to_string(), service.to_string()))
127 .cloned()
128 }
129
130 pub fn of_stack(&self, stack: &str) -> BTreeMap<String, FailedAttempt> {
132 self.0
133 .lock()
134 .unwrap_or_else(|e| e.into_inner())
135 .iter()
136 .filter(|((s, _), _)| s == stack)
137 .map(|((_, svc), a)| (svc.clone(), a.clone()))
138 .collect()
139 }
140
141 pub fn clear(&self, stack: &str, service: &str) {
143 self.0
144 .lock()
145 .unwrap_or_else(|e| e.into_inner())
146 .remove(&(stack.to_string(), service.to_string()));
147 }
148}
149
150pub fn one_line(text: &str, n: usize, max: usize) -> String {
153 let lines: Vec<&str> = text
154 .lines()
155 .map(str::trim_end)
156 .filter(|l| !l.trim().is_empty())
157 .collect();
158 let joined = lines[lines.len().saturating_sub(n)..].join(" | ");
159 let count = joined.chars().count();
160 if count <= max {
161 return joined;
162 }
163 let tail: String = joined.chars().skip(count - max).collect();
164 format!("...{tail}")
165}
166
167const OUTPUT_WAIT: Duration = if cfg!(test) {
171 Duration::from_millis(100)
172} else {
173 Duration::from_secs(8)
174};
175
176fn poll_output(
180 read: &mut dyn FnMut() -> Result<String>,
181 within: Duration,
182 poll: Duration,
183) -> (String, Option<String>) {
184 let until = std::time::Instant::now() + within;
185 let mut last_err: Option<String>;
186 loop {
187 match read() {
188 Ok(t) if !t.trim().is_empty() => return (t, None),
189 Ok(_) => last_err = None,
190 Err(e) => last_err = Some(e.to_string()),
191 }
192 if std::time::Instant::now() >= until {
193 return (String::new(), last_err);
194 }
195 std::thread::sleep(poll);
196 }
197}
198
199fn keep_end(text: &str, max: usize) -> String {
202 if text.len() <= max {
203 return text.to_string();
204 }
205 let mut start = text.len() - max;
206 while !text.is_char_boundary(start) {
207 start += 1;
208 }
209 let tail = &text[start..];
210 match tail.find('\n') {
211 Some(i) if i + 1 < tail.len() => tail[i + 1..].to_string(),
212 _ => tail.to_string(),
213 }
214}
215
216fn read_output(client: &Client, name: &str, service: &str, oci: bool) -> (String, Option<String>) {
219 poll_output(
220 &mut || {
221 let sb = Sandbox::get(client, name)?;
222 supervise::logs(&sb, service, oci, READ_LINES)
223 },
224 OUTPUT_WAIT,
225 Duration::from_millis(500),
226 )
227}
228
229pub fn explain(
232 client: &Client,
233 name: &str,
234 service: &str,
235 oci: bool,
236 e: Error,
237 now_ms: u64,
238) -> (Error, FailedAttempt) {
239 let (output, err) = read_output(client, name, service, oci);
240 let output = keep_end(&output, OUTPUT_BYTES);
241 let output_note = output.trim().is_empty().then(|| match err {
242 Some(e) => format!("its output could not be read: {e}"),
243 None => format!(
244 "it printed nothing in the {}s its output was waited for",
245 OUTPUT_WAIT.as_secs()
246 ),
247 });
248 let reason = e.to_string();
249 let tail = one_line(&output, MESSAGE_LINES, MESSAGE_CHARS);
250 let err = if tail.is_empty() {
251 e
252 } else {
253 Error::invalid(format!("{reason}; its last output: {tail}"))
254 };
255 let attempt = FailedAttempt {
256 instance: name.to_string(),
257 at_ms: now_ms,
258 reason,
259 output,
260 output_note,
261 };
262 (err, attempt)
263}
264
265pub fn replica_logs(
269 client: &Client,
270 stack: &str,
271 service: &str,
272 oci: bool,
273 slot: Option<u32>,
274 lines: usize,
275) -> Result<BTreeMap<String, String>> {
276 let mut out = BTreeMap::new();
277 for i in list_instances(client, stack, Some(service))? {
278 if slot.is_some_and(|s| s != i.slot) {
279 continue;
280 }
281 out.insert(i.name.clone(), one_replica(client, &i, service, oci, lines));
282 }
283 Ok(out)
284}
285
286fn one_replica(client: &Client, i: &Inst, service: &str, oci: bool, lines: usize) -> String {
287 let mut last = String::new();
288 for attempt in 0..2 {
289 match Sandbox::get(client, &i.name).and_then(|sb| supervise::logs(&sb, service, oci, lines))
290 {
291 Ok(t) => return t,
292 Err(e) if e.is_not_found() => return "(replaced while reading its logs)".into(),
293 Err(e) => last = e.to_string(),
294 }
295 if attempt == 0 {
296 std::thread::sleep(Duration::from_millis(300));
297 }
298 }
299 format!("(no logs: {last})")
300}
301
302#[cfg(test)]
303mod tests {
304 use super::*;
305
306 #[test]
307 fn a_missing_image_backs_off_long_and_says_so() {
308 let e = Error::invalid(
309 "create instance x failed: Failed getting remote image info: Failed to run: skopeo inspect docker://docker.io/library/traefik:whoami: reading manifest whoami in docker.io/library/traefik: manifest unknown",
310 );
311 let (wait, m) = retry(None, "docker:traefik:whoami", "slot 1: ...", &e);
312 assert_eq!(wait, Duration::from_secs(300));
313 assert_eq!(
314 m,
315 "image docker:traefik:whoami not found (manifest unknown): change the image and deploy again (retrying in 5m)"
316 );
317 let mut w = Some(Duration::from_secs(10));
318 for _ in 0..8 {
319 w = Some(retry(w, "docker:traefik:whoami", "slot 1: ...", &e).0);
320 }
321 assert_eq!(w, Some(Duration::from_secs(3600)));
322 let (wait, m) = retry(
324 None,
325 "docker:nginx",
326 "slot 1: boom",
327 &Error::invalid("boom"),
328 );
329 assert_eq!(wait, Duration::from_secs(10));
330 assert_eq!(m, "slot 1: boom; retrying in 10s");
331 }
332
333 #[test]
334 fn a_pull_of_a_missing_image_is_said_plainly() {
335 let incus = r#"create instance web-1 failed: Failed getting remote image info: Failed to run: skopeo --insecure-policy inspect docker://docker.io/library/traefik:whoami --no-tags: exit status 2 (time="2026-10-05T04:38:08Z" level=fatal msg="Error parsing image name \"docker://docker.io/library/traefik:whoami\": reading manifest whoami in docker.io/library/traefik: manifest unknown")"#;
336 assert_eq!(
337 image_missing("docker:traefik:whoami", incus).as_deref(),
338 Some("image docker:traefik:whoami not found (manifest unknown)")
339 );
340 let offline = "create instance web-1 failed: Failed getting remote image info: Failed to run: skopeo: dial tcp: lookup registry-1.docker.io: no such host";
341 assert_eq!(image_missing("docker:nginx", offline), None);
342 assert_eq!(
343 image_missing("docker:nginx", "out of disk: manifest unknown"),
344 None
345 );
346 }
347
348 #[test]
349 fn the_message_carries_the_end_of_the_output() {
350 let out = "starting\n\nconnecting to db\nError: getaddrinfo ENOTFOUND umami-db\n at GetAddrInfoReqWrap\n";
351 assert_eq!(
352 one_line(out, 2, 500),
353 "Error: getaddrinfo ENOTFOUND umami-db | at GetAddrInfoReqWrap"
354 );
355 assert_eq!(one_line("", 8, 100), "");
356 assert_eq!(one_line("a\nb\nc\n", 8, 100), "a | b | c");
357 let long = format!("{}\nENOTFOUND", "x".repeat(300));
359 let s = one_line(&long, 8, 40);
360 assert!(s.starts_with("...") && s.ends_with("ENOTFOUND"), "{s}");
361 assert_eq!(s.chars().count(), 43);
362 }
363
364 #[test]
365 fn the_last_attempt_is_kept_per_service_until_it_serves_again() {
366 let f = Failures::default();
367 let a = FailedAttempt {
368 instance: "shop-web-1-abc".into(),
369 at_ms: 7,
370 reason: "failed within the 5s monitor period".into(),
371 output: "boom".into(),
372 output_note: None,
373 };
374 assert_eq!(f.last("shop", "web"), None);
375 f.record("shop", "web", a.clone());
376 assert_eq!(f.last("shop", "web"), Some(a.clone()));
377 assert_eq!(f.last("shop", "db"), None);
378 f.record("other", "web", a.clone());
379 assert_eq!(f.of_stack("shop").into_keys().collect::<Vec<_>>(), ["web"]);
380 f.clear("shop", "web");
381 assert_eq!(f.last("shop", "web"), None);
382 }
383
384 #[test]
385 fn an_unreadable_instance_still_explains_the_failure_without_output() {
386 let c = Client::with_socket("/nonexistent/incus.sock");
389 let (e, a) = explain(
390 &c,
391 "gone",
392 "web",
393 true,
394 Error::invalid("failed within the 5s monitor period"),
395 9,
396 );
397 assert_eq!(e.to_string(), "failed within the 5s monitor period");
398 assert_eq!(a.reason, "failed within the 5s monitor period");
399 assert!(a.output.is_empty());
400 assert!(
402 a.output_note
403 .as_deref()
404 .is_some_and(|n| n.starts_with("its output could not be read: ")),
405 "{:?}",
406 a.output_note
407 );
408 assert_eq!((a.instance.as_str(), a.at_ms), ("gone", 9));
409 }
410
411 #[test]
412 fn a_replica_deleted_while_its_logs_are_read_does_not_fail_the_call() {
413 use crate::client::fake::{Route, serve};
414 use crate::stack::{LABEL_SERVICE, LABEL_SLOT, LABEL_STACK};
415 let (_d, c) = serve(vec![Route {
417 prefix: "GET /1.0/instances?recursion=1",
418 status: 200,
419 body: serde_json::json!([{
420 "name": "web-1-aaaa",
421 "status": "Running",
422 "config": {
423 format!("user.{LABEL_STACK}"): "shop",
424 format!("user.{LABEL_SERVICE}"): "web",
425 format!("user.{LABEL_SLOT}"): "1",
426 },
427 }]),
428 }]);
429 let logs = replica_logs(&c, "shop", "web", true, None, 50).unwrap();
430 assert_eq!(logs.len(), 1);
431 assert_eq!(logs["web-1-aaaa"], "(replaced while reading its logs)");
432 assert!(
434 replica_logs(&c, "shop", "web", true, Some(2), 50)
435 .unwrap()
436 .is_empty()
437 );
438 }
439
440 #[test]
441 fn output_that_arrives_late_is_waited_for_and_silence_ends() {
442 let mut n = 0;
443 let got = poll_output(
444 &mut || {
445 n += 1;
446 match n {
447 1 => Err(Error::invalid("connection refused")),
448 2 => Ok("\n".into()),
449 _ => Ok("Error: getaddrinfo ENOTFOUND db\n".into()),
450 }
451 },
452 Duration::from_secs(5),
453 Duration::from_millis(1),
454 );
455 assert_eq!(got, ("Error: getaddrinfo ENOTFOUND db\n".to_string(), None));
456 assert_eq!(n, 3);
457 let t = std::time::Instant::now();
459 let none = poll_output(
460 &mut || Ok(String::new()),
461 Duration::from_millis(30),
462 Duration::from_millis(10),
463 );
464 assert_eq!(none, (String::new(), None));
465 assert!(t.elapsed() < Duration::from_secs(2));
466 }
467
468 #[test]
469 fn the_output_kept_is_its_end_within_the_cap() {
470 let text: String = (0..5000).map(|i| format!("line {i}\n")).collect();
471 let kept = keep_end(&text, OUTPUT_BYTES);
472 assert!(kept.len() <= OUTPUT_BYTES);
473 assert!(kept.starts_with("line ") && kept.ends_with("line 4999\n"));
474 assert_eq!(keep_end("short", 10), "short");
475 assert_eq!(keep_end("ééé", 3), "é");
477 }
478}