1use std::collections::BTreeMap;
12use std::sync::Mutex;
13use std::time::Duration;
14
15use serde::Serialize;
16
17use super::controller::{Inst, list_instances};
18use crate::client::Client;
19use crate::error::{Error, Result};
20use crate::sandbox::Sandbox;
21use crate::supervise;
22
23const READ_LINES: usize = 60;
25const MESSAGE_LINES: usize = 8;
27const MESSAGE_CHARS: usize = 1200;
29
30#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
32pub struct FailedAttempt {
33 pub instance: String,
35 pub at_ms: u64,
37 pub reason: String,
39 pub output: String,
41}
42
43#[derive(Default)]
45pub struct Failures(Mutex<BTreeMap<(String, String), FailedAttempt>>);
46
47impl Failures {
48 pub fn record(&self, stack: &str, service: &str, a: FailedAttempt) {
49 self.0
50 .lock()
51 .unwrap_or_else(|e| e.into_inner())
52 .insert((stack.to_string(), service.to_string()), a);
53 }
54
55 pub fn last(&self, stack: &str, service: &str) -> Option<FailedAttempt> {
56 self.0
57 .lock()
58 .unwrap_or_else(|e| e.into_inner())
59 .get(&(stack.to_string(), service.to_string()))
60 .cloned()
61 }
62
63 pub fn clear(&self, stack: &str, service: &str) {
65 self.0
66 .lock()
67 .unwrap_or_else(|e| e.into_inner())
68 .remove(&(stack.to_string(), service.to_string()));
69 }
70}
71
72pub fn one_line(text: &str, n: usize, max: usize) -> String {
75 let lines: Vec<&str> = text
76 .lines()
77 .map(str::trim_end)
78 .filter(|l| !l.trim().is_empty())
79 .collect();
80 let joined = lines[lines.len().saturating_sub(n)..].join(" | ");
81 let count = joined.chars().count();
82 if count <= max {
83 return joined;
84 }
85 let tail: String = joined.chars().skip(count - max).collect();
86 format!("...{tail}")
87}
88
89const OUTPUT_WAIT: Duration = if cfg!(test) {
93 Duration::from_millis(100)
94} else {
95 Duration::from_secs(8)
96};
97
98fn poll_output(
102 read: &mut dyn FnMut() -> Result<String>,
103 within: Duration,
104 poll: Duration,
105) -> String {
106 let until = std::time::Instant::now() + within;
107 loop {
108 if let Ok(t) = read() {
109 if !t.trim().is_empty() {
110 return t;
111 }
112 }
113 if std::time::Instant::now() >= until {
114 return String::new();
115 }
116 std::thread::sleep(poll);
117 }
118}
119
120fn read_output(client: &Client, name: &str, service: &str, oci: bool) -> String {
122 poll_output(
123 &mut || {
124 let sb = Sandbox::get(client, name)?;
125 supervise::logs(&sb, service, oci, READ_LINES)
126 },
127 OUTPUT_WAIT,
128 Duration::from_millis(500),
129 )
130}
131
132pub fn explain(
135 client: &Client,
136 name: &str,
137 service: &str,
138 oci: bool,
139 e: Error,
140 now_ms: u64,
141) -> (Error, FailedAttempt) {
142 let output = read_output(client, name, service, oci);
143 let reason = e.to_string();
144 let tail = one_line(&output, MESSAGE_LINES, MESSAGE_CHARS);
145 let err = if tail.is_empty() {
146 e
147 } else {
148 Error::invalid(format!("{reason}; its last output: {tail}"))
149 };
150 let attempt = FailedAttempt {
151 instance: name.to_string(),
152 at_ms: now_ms,
153 reason,
154 output,
155 };
156 (err, attempt)
157}
158
159pub fn replica_logs(
163 client: &Client,
164 stack: &str,
165 service: &str,
166 oci: bool,
167 slot: Option<u32>,
168 lines: usize,
169) -> Result<BTreeMap<String, String>> {
170 let mut out = BTreeMap::new();
171 for i in list_instances(client, stack, Some(service))? {
172 if slot.is_some_and(|s| s != i.slot) {
173 continue;
174 }
175 out.insert(i.name.clone(), one_replica(client, &i, service, oci, lines));
176 }
177 Ok(out)
178}
179
180fn one_replica(client: &Client, i: &Inst, service: &str, oci: bool, lines: usize) -> String {
181 let mut last = String::new();
182 for attempt in 0..2 {
183 match Sandbox::get(client, &i.name).and_then(|sb| supervise::logs(&sb, service, oci, lines))
184 {
185 Ok(t) => return t,
186 Err(e) if e.is_not_found() => return "(replaced while reading its logs)".into(),
187 Err(e) => last = e.to_string(),
188 }
189 if attempt == 0 {
190 std::thread::sleep(Duration::from_millis(300));
191 }
192 }
193 format!("(no logs: {last})")
194}
195
196#[cfg(test)]
197mod tests {
198 use super::*;
199
200 #[test]
201 fn the_message_carries_the_end_of_the_output() {
202 let out = "starting\n\nconnecting to db\nError: getaddrinfo ENOTFOUND umami-db\n at GetAddrInfoReqWrap\n";
203 assert_eq!(
204 one_line(out, 2, 500),
205 "Error: getaddrinfo ENOTFOUND umami-db | at GetAddrInfoReqWrap"
206 );
207 assert_eq!(one_line("", 8, 100), "");
208 assert_eq!(one_line("a\nb\nc\n", 8, 100), "a | b | c");
209 let long = format!("{}\nENOTFOUND", "x".repeat(300));
211 let s = one_line(&long, 8, 40);
212 assert!(s.starts_with("...") && s.ends_with("ENOTFOUND"), "{s}");
213 assert_eq!(s.chars().count(), 43);
214 }
215
216 #[test]
217 fn the_last_attempt_is_kept_per_service_until_it_serves_again() {
218 let f = Failures::default();
219 let a = FailedAttempt {
220 instance: "shop-web-1-abc".into(),
221 at_ms: 7,
222 reason: "failed within the 5s monitor period".into(),
223 output: "boom".into(),
224 };
225 assert_eq!(f.last("shop", "web"), None);
226 f.record("shop", "web", a.clone());
227 assert_eq!(f.last("shop", "web"), Some(a));
228 assert_eq!(f.last("shop", "db"), None);
229 f.clear("shop", "web");
230 assert_eq!(f.last("shop", "web"), None);
231 }
232
233 #[test]
234 fn an_unreadable_instance_still_explains_the_failure_without_output() {
235 let c = Client::with_socket("/nonexistent/incus.sock");
238 let (e, a) = explain(
239 &c,
240 "gone",
241 "web",
242 true,
243 Error::invalid("failed within the 5s monitor period"),
244 9,
245 );
246 assert_eq!(e.to_string(), "failed within the 5s monitor period");
247 assert_eq!(a.reason, "failed within the 5s monitor period");
248 assert!(a.output.is_empty());
249 assert_eq!((a.instance.as_str(), a.at_ms), ("gone", 9));
250 }
251
252 #[test]
253 fn a_replica_deleted_while_its_logs_are_read_does_not_fail_the_call() {
254 use crate::client::fake::{Route, serve};
255 use crate::stack::{LABEL_SERVICE, LABEL_SLOT, LABEL_STACK};
256 let (_d, c) = serve(vec![Route {
258 prefix: "GET /1.0/instances?recursion=1",
259 status: 200,
260 body: serde_json::json!([{
261 "name": "web-1-aaaa",
262 "status": "Running",
263 "config": {
264 format!("user.{LABEL_STACK}"): "shop",
265 format!("user.{LABEL_SERVICE}"): "web",
266 format!("user.{LABEL_SLOT}"): "1",
267 },
268 }]),
269 }]);
270 let logs = replica_logs(&c, "shop", "web", true, None, 50).unwrap();
271 assert_eq!(logs.len(), 1);
272 assert_eq!(logs["web-1-aaaa"], "(replaced while reading its logs)");
273 assert!(
275 replica_logs(&c, "shop", "web", true, Some(2), 50)
276 .unwrap()
277 .is_empty()
278 );
279 }
280
281 #[test]
282 fn output_that_arrives_late_is_waited_for_and_silence_ends() {
283 let mut n = 0;
284 let got = poll_output(
285 &mut || {
286 n += 1;
287 match n {
288 1 => Err(Error::invalid("connection refused")),
289 2 => Ok("\n".into()),
290 _ => Ok("Error: getaddrinfo ENOTFOUND db\n".into()),
291 }
292 },
293 Duration::from_secs(5),
294 Duration::from_millis(1),
295 );
296 assert_eq!(got, "Error: getaddrinfo ENOTFOUND db\n");
297 assert_eq!(n, 3);
298 let t = std::time::Instant::now();
300 let none = poll_output(
301 &mut || Ok(String::new()),
302 Duration::from_millis(30),
303 Duration::from_millis(10),
304 );
305 assert_eq!(none, "");
306 assert!(t.elapsed() < Duration::from_secs(2));
307 }
308}