1use crate::*;
2use futures::stream;
3
4async fn history() -> Result<Value> {
5 host::call(json!({"operation":"deploy.history"})).await
6}
7async fn current() -> Result<Option<String>> {
8 Ok(serde_json::from_value(
9 host::call(json!({"operation":"deploy.current"})).await?,
10 )?)
11}
12async fn stages(app: Arc<App>) -> Result<Value> {
13 let mut values = host::call(json!({"operation":"deploy.stages"})).await?;
14 let jobs = core::scan(app).await.ok();
15 for stage in values.as_array_mut().unwrap() {
16 let job = jobs.as_ref().map(|jobs| &jobs.value[string(&stage["id"])]);
17 stage["hostname"] = json!(
18 job.and_then(|job| job["job"]["Meta"]["studio_hostname"].as_str())
19 .filter(|hostname| !hostname.is_empty())
20 );
21 stage["health"] = json!(job.map(|job| job["health"].as_str().unwrap_or("down")));
22 }
23 Ok(values)
24}
25async fn tree(app: Arc<App>, id: &str) -> Result<Arc<Document>> {
26 let id = id.to_owned();
27 app.cache
28 .get(
29 format!("release:{id}"),
30 Duration::from_secs(365 * 86400),
31 move || async move {
32 host::call(json!({"operation":"deploy.release","release":id})).await
33 },
34 )
35 .await
36}
37async fn changed(
38 app: Arc<App>,
39 from: Option<&str>,
40 to: Option<&str>,
41) -> Result<Option<Vec<String>>> {
42 let (Some(from), Some(to)) = (from, to) else {
43 return Ok(None);
44 };
45 let (a, b) = tokio::try_join!(tree(app.clone(), from), tree(app.clone(), to))?;
46 let (Some(a), Some(b)) = (a.value.as_object(), b.value.as_object()) else {
47 return Ok(None);
48 };
49 let ids: std::collections::BTreeSet<_> = a.keys().chain(b.keys()).collect();
50 Ok(Some(
51 ids.into_iter()
52 .filter(|id| a.get(*id) != b.get(*id))
53 .cloned()
54 .collect(),
55 ))
56}
57fn diff(before: &str, after: &str) -> Vec<Value> {
58 let a: Vec<_> = if before.is_empty() {
59 Vec::new()
60 } else {
61 before.split('\n').collect()
62 };
63 let b: Vec<_> = if after.is_empty() {
64 Vec::new()
65 } else {
66 after.split('\n').collect()
67 };
68 let mut common = vec![vec![0usize; b.len() + 1]; a.len() + 1];
69 for i in (0..a.len()).rev() {
70 for j in (0..b.len()).rev() {
71 common[i][j] = if a[i] == b[j] {
72 common[i + 1][j + 1] + 1
73 } else {
74 common[i + 1][j].max(common[i][j + 1])
75 };
76 }
77 }
78 let (mut i, mut j) = (0, 0);
79 let mut lines = Vec::new();
80 while i < a.len() || j < b.len() {
81 if i < a.len() && j < b.len() && a[i] == b[j] {
82 lines.push(json!([" ", a[i]]));
83 i += 1;
84 j += 1;
85 } else if i < a.len() && (j == b.len() || common[i + 1][j] >= common[i][j + 1]) {
86 lines.push(json!(["-", a[i]]));
87 i += 1;
88 } else {
89 lines.push(json!(["+", b[j]]));
90 j += 1;
91 }
92 }
93 lines
94}
95async fn stage(app: Arc<App>, id: &str) -> Result<Value> {
96 if !core::valid_id(id) {
97 return Err(Error::new(
98 400,
99 "Stage IDs are lowercase letters, digits and dashes.",
100 ));
101 }
102 array(&stages(app).await?)
103 .iter()
104 .find(|s| s["id"] == id)
105 .cloned()
106 .ok_or_else(|| {
107 Error::new(
108 404,
109 format!("No stage named {id}. It may have been destroyed; go back to deploys."),
110 )
111 })
112}
113fn entry(history: &Value, n: &str) -> Result<(usize, Value)> {
114 let n = n
115 .parse::<usize>()
116 .ok()
117 .filter(|n| *n > 0)
118 .ok_or_else(|| Error::new(400, "Invalid deploy number."))?;
119 array(history)
120 .get(n - 1)
121 .cloned()
122 .map(|v| (n, v))
123 .ok_or_else(|| {
124 Error::new(
125 404,
126 format!("No deploy number {n}. Go back to deploys to see the history."),
127 )
128 })
129}
130struct Scope {
131 jobs: Vec<String>,
132 allocations: Vec<Value>,
133 from: f64,
134 to: Option<f64>,
135}
136async fn scope(app: Arc<App>, kind: &str, target: &str) -> Result<Scope> {
137 let mut after = f64::NEG_INFINITY;
138 let mut until = f64::INFINITY;
139 let mut to = None;
140 let jobs = if kind == "stages" {
141 vec![string(&stage(app.clone(), target).await?["id"]).to_owned()]
142 } else {
143 let history = history().await?;
144 let (n, entry) = entry(&history, target)?;
145 let previous = if n >= 2 {
146 &history[n - 2]
147 } else {
148 &Value::Null
149 };
150 after = previous["time"].as_f64().unwrap_or(f64::NEG_INFINITY);
151 until = number(&entry["time"]);
152 to = history[n]["time"].as_f64();
153 match changed(
154 app.clone(),
155 previous["release"].as_str(),
156 entry["release"].as_str(),
157 )
158 .await?
159 {
160 Some(jobs) => jobs,
161 None => tree(app.clone(), string(&entry["release"]))
162 .await?
163 .value
164 .as_object()
165 .into_iter()
166 .flat_map(|t| t.keys().cloned())
167 .collect(),
168 }
169 };
170 let scanned = core::scan(app).await?;
171 let mut allocations = Vec::new();
172 let mut created = Vec::new();
173 for id in &jobs {
174 let found = &scanned.value[id];
175 let mine=array(&found["allocs"]).iter().filter(|a| number(&a["CreateTime"])/1e9>after && number(&a["CreateTime"])/1e9<=until).map(|a| {
176 let tasks:Vec<_>=a["TaskStates"].as_object().into_iter().flat_map(|t| t.iter()).map(|(name,task)| json!({"name":name,"restarts":task["Restarts"],"lastRestart":core::last_restart(task)})).collect();let finished:Vec<_>=a["TaskStates"].as_object().into_iter().flat_map(|t| t.values()).map(|t| core::seconds(&t["FinishedAt"]).unwrap_or(0.0)).collect();let ended=if ["complete","failed","lost"].contains(&string(&a["ClientStatus"])) && !finished.is_empty(){finished.into_iter().max_by(f64::total_cmp)}else{None};let created_at=number(&a["CreateTime"])/1e9;created.push(created_at);
177 json!({"id":a["ID"],"state":a["ClientStatus"],"healthy":a["DeploymentStatus"]["Healthy"],"created":created_at,"ended":ended,"tasks":tasks,"checks":if a["DesiredStatus"]=="run" && ["running","pending"].contains(&string(&a["ClientStatus"])) {core::checks_of(&a["home_checks"])}else{json!([])}})
178 }).collect::<Vec<_>>();
179 allocations.push(json!(mine));
180 }
181 Ok(Scope {
182 jobs,
183 allocations,
184 from: created
185 .into_iter()
186 .min_by(f64::total_cmp)
187 .unwrap_or(after.max(0.0)),
188 to,
189 })
190}
191async fn runtime(app: Arc<App>, scope: Scope) -> Result<Value> {
192 let end = scope.to.unwrap_or((now() / 2.0).floor() * 2.0);
193 let mut jobs = Vec::new();
194 for (i, id) in scope.jobs.iter().enumerate() {
195 let query = HashMap::from([("service".into(), id.clone())]);
196 let usage = tokio::try_join!(
197 telemetry::metrics(
198 app.clone(),
199 "service.cpu",
200 &query,
201 None,
202 Some((scope.from, end))
203 ),
204 telemetry::metrics(
205 app.clone(),
206 "service.memory",
207 &query,
208 None,
209 Some((scope.from, end))
210 )
211 );
212 let (cpu, memory) = match usage {
213 Ok((cpu, memory)) => (cpu.value[0].clone(), memory.value[0].clone()),
214 Err(e) if e.status == 501 => (Value::Null, Value::Null),
215 Err(e) => return Err(e),
216 };
217 jobs.push(json!({"id":id,"allocations":scope.allocations[i],"cpu":cpu,"memory":memory}));
218 }
219 Ok(json!({"from":scope.from,"to":scope.to,"jobs":jobs}))
220}
221fn display_run(run: Value, history: &Value) -> Result<Value> {
222 if run.is_null() {
223 return Ok(run);
224 }
225 let target = string(&run["target"]);
226 let title = match string(&run["action"]) {
227 "secret-set" => format!("set {target}/{}", string(&run["key"])),
228 "secret-rotate" => format!("rotate {target}/{}", string(&run["key"])),
229 action @ ("start" | "stop" | "restart") => format!("{action} {target}"),
230 "deploy" => "deploy main".to_owned(),
231 "promote" => format!("promote {target}"),
232 "destroy" => format!("destroy {target}"),
233 "rollback" => {
234 let (_, entry) = entry(history, target)?;
235 format!(
236 "roll back to {}",
237 string(&entry["release"]).get(..8).unwrap_or_default()
238 )
239 }
240 _ => {
241 return Err(Error::new(
242 502,
243 "The deployment action couldn't be read. Reload deploys.",
244 ));
245 }
246 };
247 Ok(json!({"id":run["id"],"title":title,"target":target,"code":run["code"]}))
248}
249async fn start(app: Arc<App>, action: &str, target: &str) -> Result<Value> {
250 let run =
251 host::call(json!({"operation":"deploy.start","action":action,"target":target})).await?;
252 app.cache.invalidate("deploys");
253 display_run(run, &history().await?)
254}
255pub async fn run_stream(app: Arc<App>, id: &str) -> Result<Response> {
256 host::call(json!({"operation":"deploy.run","id":id})).await?;
257 let id = id.to_owned();
258 let stream = stream::unfold(
259 (app, id, 0usize, false),
260 move |(app, id, mut sent, ended)| async move {
261 if ended {
262 return None;
263 }
264 let run = id.clone();
265 let response =
266 app.cache
267 .get(
268 format!("run:{id}"),
269 Duration::from_secs(1),
270 move || async move {
271 host::call(json!({"operation":"deploy.run","id":run})).await
272 },
273 )
274 .await;
275 match response {
276 Ok(response) => {
277 let mut chunk = String::new();
278 let lines = array(&response.value["lines"]);
279 while sent < lines.len() {
280 for line in string(&lines[sent]).lines() {
281 chunk.push_str(&format!("data: {line}\n"));
282 }
283 chunk.push('\n');
284 sent += 1;
285 }
286 let ended = !response.value["code"].is_null();
287 if ended {
288 chunk.push_str(&format!(
289 "event: exit\ndata: {}\n\n",
290 number(&response.value["code"])
291 ));
292 }
293 if chunk.is_empty() {
294 chunk.push_str(": keepalive\n\n");
295 }
296 tokio::time::sleep(Duration::from_secs(1)).await;
297 Some((
298 Ok::<_, std::io::Error>(Bytes::from(chunk)),
299 (app, id, sent, ended),
300 ))
301 }
302 Err(error) => Some((
303 Ok(Bytes::from(format!(
304 "event: exit\ndata: 1\n\n: {}\n\n",
305 error.message.replace('\n', " ")
306 ))),
307 (app, id, sent, true),
308 )),
309 }
310 },
311 );
312 Ok((
313 [
314 ("content-type", "text/event-stream"),
315 ("cache-control", "no-cache"),
316 ],
317 axum::body::Body::from_stream(stream),
318 )
319 .into_response())
320}
321
322pub async fn route(
323 app: Arc<App>,
324 method: &Method,
325 parts: &[&str],
326 query: &HashMap<String, String>,
327) -> Result<Response> {
328 let value = match parts {
329 [] if method == Method::GET => {
330 let state = app.clone();
331 let document = app
332 .cache
333 .get(
334 "deploys".into(),
335 Duration::from_secs(2),
336 move || async move {
337 let history = history().await?;
338 let current = current().await?;
339 let main = host::call(json!({"operation":"deploy.main"})).await?;
340 let mut stages = stages(state.clone()).await?;
341 let mut entries = Vec::new();
342 for (i, item) in array(&history).iter().enumerate() {
343 let mut item = item.clone();
344 item["n"] = json!(i + 1);
345 item["changed"] = json!(
346 changed(
347 state.clone(),
348 i.checked_sub(1)
349 .and_then(|n| history[n]["release"].as_str()),
350 item["release"].as_str()
351 )
352 .await?
353 );
354 item["redeploys"] = json!(
355 changed(state.clone(), current.as_deref(), item["release"].as_str()).await?
356 );
357 entries.push(item);
358 }
359 entries.reverse();
360 for stage in stages.as_array_mut().unwrap() {
361 let base = array(&history).iter().rposition(|e| {
362 e["source"] == stage["id"] || number(&e["time"]) <= number(&stage["created"])
363 });
364 stage["files"] = if stage["clone"].is_null() {
365 Value::Null
366 } else {
367 apps::file_link(string(&stage["mount"]), false)
368 };
369 stage["base"] = json!(base.map(|n| n + 1));
370 stage["changed"] = json!(
371 changed(state.clone(), current.as_deref(), stage["release"].as_str()).await?
372 );
373 stage["behind"] = json!(
374 changed(
375 state.clone(),
376 base.and_then(|n| history[n]["release"].as_str()),
377 current.as_deref()
378 )
379 .await?
380 .unwrap_or_default()
381 );
382 }
383 let run = display_run(
384 host::call(json!({"operation":"deploy.last"})).await?,
385 &history,
386 )?;
387 Ok(json!({
388 "current": current,
389 "main": main,
390 "recorded": array(&history).last().and_then(|e| e["release"].as_str()) == current.as_deref(),
391 "history": entries, "stages": stages, "run": run
392 }))
393 },
394 )
395 .await?;
396 return Ok(document.response());
397 }
398 ["changes"] if method == Method::GET => {
399 let to = query
400 .get("to")
401 .ok_or_else(|| Error::new(400, "Pick a release."))?;
402 let a = if let Some(from) = query.get("from") {
403 tree(app.clone(), from).await?
404 } else {
405 Arc::new(Document::new(json!({})))
406 };
407 let b = tree(app.clone(), to).await?;
408 let (Some(a), Some(b)) = (a.value.as_object(), b.value.as_object()) else {
409 return Err(Error::new(
410 410,
411 "That release is no longer on the host, so its changes can't be shown.",
412 ));
413 };
414 let ids: std::collections::BTreeSet<_> = a.keys().chain(b.keys()).collect();
415 json!(ids.into_iter().filter(|id| a.get(*id)!=b.get(*id)).map(|id| json!({"id":id,"lines":diff(a.get(id).map(string).unwrap_or_default(),b.get(id).map(string).unwrap_or_default())})).collect::<Vec<_>>())
416 }
417 [kind, target, "runtime"]
418 if method == Method::GET && ["stages", "history"].contains(kind) =>
419 {
420 runtime(app.clone(), scope(app, kind, target).await?).await?
421 }
422 [kind, target, "logs"] if method == Method::GET && ["stages", "history"].contains(kind) => {
423 let scope = scope(app.clone(), kind, target).await?;
424 let job = query
425 .get("job")
426 .or(scope.jobs.first())
427 .ok_or_else(|| Error::new(404, "No job ran here, so there are no logs to show."))?;
428 if !scope.jobs.contains(job) {
429 return Err(Error::new(
430 404,
431 format!("{job} didn't run here, so there are no logs to show."),
432 ));
433 }
434 job_logs(app, job, query, scope.from, scope.to).await?
435 }
436 [kind, target, "output"]
437 if method == Method::GET && ["stages", "history"].contains(kind) =>
438 {
439 json!({"lines":host::call(json!({"operation":"deploy.output","kind":kind,"target":target})).await?})
440 }
441 ["main", release, "deploy"] if method == Method::POST => {
442 start(app, "deploy", release).await?
443 }
444 ["stages", id, "destroy"] if method == Method::POST => start(app, "destroy", id).await?,
445 ["history", n, "rollback"] if method == Method::POST => start(app, "rollback", n).await?,
446 _ => return Err(Error::new(404, "Not Found")),
447 };
448 Ok(Document::new(value).response())
449}
450
451async fn job_logs(
452 app: Arc<App>,
453 job: &str,
454 query: &HashMap<String, String>,
455 from: f64,
456 to: Option<f64>,
457) -> Result<Value> {
458 let jobs = core::scan(app.clone()).await?;
459 let found = &jobs.value[job];
460 let limit = telemetry::query_number(query, "limit", 300.0, 1.0, 1000.0)? as usize;
461 let after =
462 telemetry::query_number(query, "after", from, f64::NEG_INFINITY, f64::INFINITY)?.max(from);
463 let before = telemetry::query_number(
464 query,
465 "before",
466 to.unwrap_or(f64::MAX),
467 f64::NEG_INFINITY,
468 f64::INFINITY,
469 )?
470 .min(to.unwrap_or(f64::MAX));
471 let mut lines = Vec::new();
472 let ansi = regex::Regex::new(r"\x1b\[[0-?]*[ -/]*[@-~]|\r").unwrap();
473 let error =
474 regex::Regex::new(r"(?i)\b(ERROR|ERR|FATAL|CRITICAL|PANIC)\b|level=(error|fatal)").unwrap();
475 let warn = regex::Regex::new(r"(?i)\b(WARN|WARNING|WRN)\b|level=warn").unwrap();
476 for alloc in array(&found["allocs"]) {
477 for task in alloc["TaskStates"]
478 .as_object()
479 .into_iter()
480 .flat_map(|t| t.keys())
481 {
482 let path = format!(
483 "/v1/client/fs/logs/{}?{}",
484 string(&alloc["ID"]),
485 params(&[
486 ("task", task.clone()),
487 ("type", "stdout".into()),
488 ("origin", "end".into()),
489 ("offset", (-512 * 1024).to_string()),
490 ("plain", "true".into())
491 ])
492 );
493 let Some(bytes) = core::nomad_raw(&app, &path).await? else {
494 continue;
495 };
496 let text = String::from_utf8_lossy(&bytes);
497 let mut partial = String::new();
498 for line in text.lines() {
499 let fields: Vec<_> = line.splitn(4, ' ').collect();
500 if fields.len() != 4
501 || !["stdout", "stderr"].contains(&fields[1])
502 || !["F", "P"].contains(&fields[2])
503 {
504 continue;
505 }
506 let Some(t) = core::seconds(&json!(fields[0])) else {
507 continue;
508 };
509 partial.push_str(fields[3]);
510 if fields[2] == "P" {
511 continue;
512 }
513 let text = std::mem::take(&mut partial);
514 if t <= after || t >= before {
515 continue;
516 }
517 let plain = ansi.replace_all(&text, "");
518 let level = if error.is_match(&plain) {
519 Some("error")
520 } else if warn.is_match(&plain) {
521 Some("warn")
522 } else {
523 None
524 };
525 lines.push(
526 json!({"t":t,"container":task,"stream":fields[1],"level":level,"text":text}),
527 );
528 }
529 }
530 }
531 lines.sort_by(|a, b| number(&b["t"]).total_cmp(&number(&a["t"])));
532 lines.truncate(limit);
533 Ok(json!(lines))
534}