1use crate::*;
2use futures::{StreamExt, stream};
3use sha1::{Digest, Sha1};
4use std::path::Path;
5
6pub async fn nomad(app: &App, path: &str) -> Result<Value> {
7 let value = nomad_raw(app, path).await?;
8 Ok(value
9 .map(|bytes| serde_json::from_slice(&bytes))
10 .transpose()?
11 .unwrap_or(Value::Null))
12}
13pub async fn nomad_raw(app: &App, path: &str) -> Result<Option<Bytes>> {
14 let request = async {
15 let token = tokio::fs::read_to_string(env(
16 "STUDIO_NOMAD_TOKEN_FILE",
17 "/var/lib/studio/dashboard.token",
18 ))
19 .await
20 .map_err(|e| Error::new(501, format!("Nomad isn't connected: {e}")))?;
21 let _slot = app.nomad_slots.acquire().await?;
22 let base = match &app.internal {
23 Some((base, _)) => format!("{}nomad", base.as_str()),
24 None => env("NOMAD_ADDR", "http://127.0.0.1:4646"),
25 };
26 let response = app
27 .request(Method::GET, &format!("{base}{path}"))?
28 .header("X-Nomad-Token", token.trim())
29 .send()
30 .await
31 .map_err(|_| {
32 Error::new(502, "Service status is taking too long. Retry in a moment.")
33 })?;
34 if response.status() == 404 {
35 return Ok(None);
36 }
37 let status = response.status();
38 let bytes = response.bytes().await?;
39 if !status.is_success() {
40 return Err(Error::new(
41 502,
42 format!(
43 "Nomad refused {} ({status}): {}",
44 path.split('?').next().unwrap(),
45 String::from_utf8_lossy(&bytes).trim()
46 ),
47 ));
48 }
49 Ok(Some(bytes))
50 };
51 tokio::time::timeout(Duration::from_secs(15), request)
52 .await
53 .map_err(|_| Error::new(502, "Service status is taking too long. Retry in a moment."))?
54}
55
56fn live_alloc(alloc: &Value) -> bool {
57 alloc["DesiredStatus"] == "run"
58 && matches!(string(&alloc["ClientStatus"]), "pending" | "running")
59}
60pub fn health(job: &Value, allocs: &[Value], checks: &[Value]) -> &'static str {
61 if job["Stop"] == true {
62 return "stopped";
63 }
64 let running: Vec<_> = allocs.iter().filter(|a| live_alloc(a)).collect();
65 if running.is_empty() {
66 return "down";
67 }
68 if running
69 .iter()
70 .any(|a| a["JobVersion"] != running[0]["JobVersion"])
71 {
72 return "deploying";
73 }
74 if running.iter().any(|a| a["ClientStatus"] == "pending") {
75 return "starting";
76 }
77 if running.iter().any(|a| {
78 a["TaskStates"]
79 .as_object()
80 .into_iter()
81 .flat_map(|s| s.values())
82 .any(|s| s["State"] == "pending")
83 }) {
84 return "restarting";
85 }
86 if checks.iter().any(|c| c["Status"] == "failure")
87 || running
88 .iter()
89 .any(|a| a["DeploymentStatus"]["Healthy"] == false)
90 {
91 return "degraded";
92 }
93 "healthy"
94}
95
96pub async fn scan(app: Arc<App>) -> Result<Arc<Document>> {
97 let state = app.clone();
98 app.cache
99 .get(
100 "nomad-scan".into(),
101 Duration::from_secs(10),
102 move || async move {
103 let (stubs, allocations, nodes) = tokio::try_join!(
104 nomad(&state, "/v1/jobs"),
105 nomad(&state, "/v1/allocations"),
106 nomad(&state, "/v1/nodes?resources=true")
107 )?;
108 let cpu = &nodes[0]["NodeResources"]["Cpu"];
109 let mhz = number(&cpu["CpuShares"]) / number(&cpu["TotalCpuCores"]);
110 if !stubs.is_array() || !allocations.is_array() || !mhz.is_finite() || mhz <= 0.0 {
111 return Err(Error::new(
112 502,
113 "Nomad answered without its jobs, allocations or node.",
114 ));
115 }
116 let jobs = stream::iter(array(&stubs).iter().cloned())
117 .map(|stub| {
118 let state = state.clone();
119 async move {
120 let id = string(&stub["ID"]).to_owned();
121 let cached = state
122 .specs
123 .lock()
124 .unwrap()
125 .get(&id)
126 .filter(|job| job["JobModifyIndex"] == stub["JobModifyIndex"])
127 .cloned();
128 let job = match cached {
129 Some(job) => job,
130 None => nomad(&state, &format!("/v1/job/{}", encoded(&id))).await?,
131 };
132 state.specs.lock().unwrap().insert(id.clone(), job.clone());
133 Ok::<_, Error>((id, job))
134 }
135 })
136 .buffer_unordered(4)
137 .collect::<Vec<_>>()
138 .await
139 .into_iter()
140 .collect::<Result<HashMap<_, _>>>()?;
141 state
142 .specs
143 .lock()
144 .unwrap()
145 .retain(|id, _| jobs.contains_key(id));
146 let mut allocs = array(&allocations).to_vec();
147 allocs
148 .sort_by(|a, b| number(&b["CreateTime"]).total_cmp(&number(&a["CreateTime"])));
149 let checked = stream::iter(allocs)
150 .map(|alloc| {
151 let state = state.clone();
152 async move {
153 let checks = if live_alloc(&alloc) {
154 nomad(
155 &state,
156 &format!(
157 "/v1/client/allocation/{}/checks",
158 string(&alloc["ID"])
159 ),
160 )
161 .await?
162 .as_object()
163 .map(|v| v.values().cloned().collect::<Vec<_>>())
164 .unwrap_or_default()
165 } else {
166 Vec::new()
167 };
168 let mut alloc = alloc;
169 alloc["home_checks"] = json!(checks);
170 Ok::<_, Error>(alloc)
171 }
172 })
173 .buffered(4)
174 .collect::<Vec<_>>()
175 .await
176 .into_iter()
177 .collect::<Result<Vec<_>>>()?;
178 let mut found = serde_json::Map::new();
179 for (id, job) in jobs {
180 let mine: Vec<_> = checked.iter().filter(|a| a["JobID"] == id).collect();
181 let allocs: Vec<_> = mine.iter().map(|a| (*a).clone()).collect();
182 let checks: Vec<_> = mine
183 .iter()
184 .flat_map(|a| array(&a["home_checks"]).iter().cloned())
185 .collect();
186 let health = health(&job, &allocs, &checks);
187 found.insert(
188 id,
189 json!({"job":job,"allocs":allocs,"mhz":mhz,"health":health}),
190 );
191 }
192 Ok(Value::Object(found))
193 },
194 )
195 .await
196}
197
198pub async fn managed() -> Result<Vec<String>> {
199 Ok(host::call(json!({"operation":"deploy.managed"}))
200 .await?
201 .as_array()
202 .unwrap()
203 .iter()
204 .map(|id| string(id).to_owned())
205 .collect())
206}
207pub async fn read_json(file: &std::path::Path, missing: Value) -> Result<Value> {
208 match tokio::fs::read(file).await {
209 Ok(bytes) => Ok(serde_json::from_slice(&bytes)?),
210 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(missing),
211 Err(e) => Err(e.into()),
212 }
213}
214pub async fn write_json(file: &std::path::Path, value: &Value) -> Result<()> {
215 tokio::fs::create_dir_all(file.parent().unwrap()).await?;
216 let temporary = file.with_extension(format!("{}.tmp", rand::random::<u64>()));
217 tokio::fs::write(&temporary, serde_json::to_vec(value)?).await?;
218 tokio::fs::rename(&temporary, file).await?;
219 Ok(())
220}
221pub async fn service_file(app: &App, id: &str) -> Result<PathBuf> {
222 if !valid_id(id) {
223 return Err(Error::new(400, "Invalid service ID."));
224 }
225 let excluded = excluded_services(&app.repo).await?;
226 if excluded.iter().any(|service| service == id) {
227 return Err(Error::new(
228 404,
229 format!("No service is named {id}. Pick one from the sidebar."),
230 ));
231 }
232 let direct = app.repo.join("service").join(id).join("service.pkl");
233 if tokio::fs::try_exists(&direct).await? {
234 return Ok(direct);
235 }
236 let mut dirs = tokio::fs::read_dir(app.repo.join("service")).await?;
237 while let Some(dir) = dirs.next_entry().await? {
238 if excluded
239 .iter()
240 .any(|service| dir.file_name() == service.as_str())
241 {
242 continue;
243 }
244 let file = dir.path().join(format!("{id}.pkl"));
245 if tokio::fs::try_exists(&file).await? {
246 return Ok(file);
247 }
248 }
249 Ok(direct)
250}
251pub fn valid_id(id: &str) -> bool {
252 !id.is_empty()
253 && id.as_bytes()[0].is_ascii_lowercase()
254 && id
255 .bytes()
256 .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-')
257}
258
259async fn icons(app: Arc<App>, id: String) -> Result<Arc<Document>> {
260 let state = app.clone();
261 app.cache.get(format!("icon:{id}"),Duration::from_secs(30),move || async move {
262 let file = service_file(&state,&id).await?;
263 let dir = if file.file_name().unwrap() == "service.pkl" { file.parent().unwrap().to_path_buf() } else { file.parent().unwrap().join(&id) };
264 let mut found = serde_json::Map::new();
265 for file in ["icon.svg","icon-light.svg","icon-dark.svg"] {
266 let candidate = dir.join(file);
267 let fallback = PathBuf::from(env("STUDIO_ICON_DIR","server/icons")).join(&id).join(file);
268 let path = if tokio::fs::try_exists(&candidate).await? { candidate } else { fallback };
269 if let Ok(bytes) = tokio::fs::read(&path).await {
270 found.insert(file.to_owned(),json!({"path":path,"hash":format!("{:x}",Sha1::digest(&bytes))[..12].to_owned()}));
271 }
272 }
273 Ok(Value::Object(found))
274 }).await
275}
276async fn icon_of(app: Arc<App>, id: &str) -> Value {
277 let Ok(found) = icons(app, id.into()).await else {
278 return Value::Null;
279 };
280 let url = |mode| {
281 let file = if !found.value[mode].is_null() {
282 mode
283 } else {
284 "icon.svg"
285 };
286 found.value[file]["hash"]
287 .as_str()
288 .map(|hash| format!("/api/icons/{id}/{file}?v={hash}"))
289 };
290 match (url("icon-light.svg"), url("icon-dark.svg")) {
291 (Some(light), Some(dark)) => json!({"light":light,"dark":dark}),
292 _ => Value::Null,
293 }
294}
295pub async fn icon(
296 app: &Arc<App>,
297 parts: &[&str],
298 query: &HashMap<String, String>,
299) -> Result<Response> {
300 if parts.len() != 3
301 || !valid_id(parts[1])
302 || !["icon.svg", "icon-light.svg", "icon-dark.svg"].contains(&parts[2])
303 {
304 return Err(Error::new(404, "Not Found"));
305 }
306 let found = icons(app.clone(), parts[1].into()).await?;
307 let path = found.value[parts[2]]["path"]
308 .as_str()
309 .ok_or_else(|| Error::new(404, "Not Found"))?;
310 Ok((
311 [
312 ("content-type", "image/svg+xml"),
313 (
314 "cache-control",
315 if query.contains_key("v") {
316 "public, max-age=31536000, immutable"
317 } else {
318 "no-cache"
319 },
320 ),
321 ],
322 tokio::fs::read(path).await?,
323 )
324 .into_response())
325}
326
327async fn summarize(app: Arc<App>, id: &str, found: &Value) -> Value {
328 let meta = &found["job"]["Meta"];
329 let tasks: Vec<_> = array(&found["job"]["TaskGroups"])
330 .iter()
331 .flat_map(|g| array(&g["Tasks"]))
332 .collect();
333 let usage = app
334 .usage
335 .lock()
336 .unwrap()
337 .get(id)
338 .cloned()
339 .unwrap_or(Value::Null);
340 let mhz = number(&found["mhz"]);
341 let cpu: f64 = tasks.iter().map(|t| number(&t["Resources"]["CPU"])).sum();
342 let memory: f64 = tasks
343 .iter()
344 .map(|t| {
345 let max = number(&t["Resources"]["MemoryMaxMB"]);
346 if max > 0.0 {
347 max
348 } else {
349 number(&t["Resources"]["MemoryMB"])
350 }
351 })
352 .sum();
353 json!({"id":id,"name":meta["studio_name"].as_str().unwrap_or(id),"health":found["health"].as_str().unwrap_or("down"),"url":meta["studio_hostname"].as_str().filter(|h| !h.is_empty() && meta["studio_launcher"] != "false").map(|h| format!("https://{h}")),"icon":icon_of(app.clone(),id).await,"access":meta["studio_auth_role"],"tagline":meta["studio_tagline"],"cpu":usage["cpu"],"cpuLimit":if mhz > 0.0 { cpu/mhz } else {0.0},"memory":usage["memory"],"memoryLimit":memory*1048576.0})
354}
355async fn summaries(app: Arc<App>) -> Result<Arc<Document>> {
356 let state = app.clone();
357 app.cache
358 .get(
359 "services".into(),
360 Duration::from_secs(2),
361 move || async move {
362 let (ids, jobs) = tokio::try_join!(managed(), scan(state.clone()))?;
363 let values = stream::iter(ids)
364 .map(|id| {
365 let state = state.clone();
366 let jobs = jobs.clone();
367 async move { summarize(state, &id, &jobs.value[&id]).await }
368 })
369 .buffered(4)
370 .collect::<Vec<_>>()
371 .await;
372 Ok(json!(values))
373 },
374 )
375 .await
376}
377
378pub fn seconds(value: &Value) -> Option<f64> {
379 chrono::DateTime::parse_from_rfc3339(value.as_str()?)
380 .ok()
381 .map(|t| t.timestamp() as f64 + t.timestamp_subsec_nanos() as f64 / 1e9)
382 .filter(|s| *s > 0.0)
383}
384pub fn last_restart(task: &Value) -> Value {
385 let Some(t) = seconds(&task["LastRestart"]) else {
386 return Value::Null;
387 };
388 let events = array(&task["Events"]);
389 let end = events
390 .iter()
391 .rposition(|e| e["Type"] == "Restarting")
392 .map(|i| i + 1)
393 .unwrap_or(0);
394 let reason = events[..end]
395 .iter()
396 .rev()
397 .find(|e| {
398 ["Terminated", "Restart Signaled", "Driver Failure", "Killed"]
399 .contains(&string(&e["Type"]))
400 })
401 .map(|e| string(&e["DisplayMessage"]))
402 .filter(|s| !s.is_empty())
403 .unwrap_or("restarted");
404 json!({"t":t,"reason":reason})
405}
406pub fn checks_of(checks: &Value) -> Value {
407 json!(array(checks).iter().map(|c| json!({"name":c["Check"],"passing":c["Status"] != "failure","output":c["Output"]})).collect::<Vec<_>>())
408}
409fn containers(found: &Value) -> Value {
410 let allocs = array(&found["allocs"]);
411 let alloc = allocs
412 .iter()
413 .find(|a| live_alloc(a))
414 .or_else(|| allocs.first())
415 .unwrap_or(&Value::Null);
416 let tasks: HashMap<_, _> = array(&found["job"]["TaskGroups"])
417 .iter()
418 .flat_map(|g| array(&g["Tasks"]))
419 .map(|t| (string(&t["Name"]), t))
420 .collect();
421 json!(alloc["TaskStates"].as_object().into_iter().flat_map(|s| s.iter()).map(|(name,state)| {
422 let spec = tasks.get(name.as_str()).copied().unwrap_or(&Value::Null);
423 json!({"name":name,"image":spec["Config"]["image"],"hook":spec["Lifecycle"]["Hook"],"state":state["State"],"restarts":state["Restarts"],"lastRestart":last_restart(state),"startedAt":seconds(&state["StartedAt"])})
424 }).collect::<Vec<_>>())
425}
426async fn issues(app: Arc<App>) -> Result<Value> {
427 let summaries = app
428 .cache
429 .peek("services", Duration::from_secs(2))
430 .ok_or_else(|| Error::new(503, "Service status is unavailable."))?;
431 let jobs = app
432 .cache
433 .peek("nomad-scan", Duration::from_secs(10))
434 .ok_or_else(|| Error::new(503, "Service status is unavailable."))?;
435 let mut issues = Vec::new();
436 for service in array(&summaries.value) {
437 let id = string(&service["id"]);
438 let health = string(&service["health"]);
439 let found = &jobs.value[id];
440 if matches!(health, "down" | "degraded") {
441 issues.push(json!({"service":{"id":id,"name":service["name"],"icon":service["icon"],"url":service["url"]},"health":health,"failing":array(&found["allocs"]).iter().flat_map(|a| array(&a["home_checks"])).find(|c| c["Status"] == "failure").map(|c| c["Output"].clone())}));
442 }
443 }
444 Ok(json!(issues))
445}
446
447async fn excluded_services(repo: &Path) -> Result<Vec<String>> {
448 Ok(serde_json::from_slice(
449 &tokio::fs::read(repo.join("config/excluded-services.json")).await?,
450 )?)
451}
452
453async fn launcher_module(repo: &Path) -> Result<String> {
454 let excluded = excluded_services(repo).await?;
455 let mut files = Vec::new();
456 let mut directories = tokio::fs::read_dir(repo.join("service")).await?;
457 while let Some(directory) = directories.next_entry().await? {
458 let directory_name = directory.file_name().to_string_lossy().into_owned();
459 if excluded.contains(&directory_name) || !directory.file_type().await?.is_dir() {
460 continue;
461 }
462 let mut entries = tokio::fs::read_dir(directory.path()).await?;
463 while let Some(entry) = entries.next_entry().await? {
464 let path = entry.path();
465 if !path.extension().is_some_and(|extension| extension == "pkl") {
466 continue;
467 }
468 let name = if entry.file_name() == "service.pkl" {
469 directory_name.clone()
470 } else {
471 path.file_stem().unwrap().to_string_lossy().into_owned()
472 };
473 if !excluded.contains(&name) {
474 files.push(path);
475 }
476 }
477 }
478 files.sort();
479 let mut module = String::new();
480 let mut services = Vec::new();
481 for (index, path) in files.iter().enumerate() {
482 let uri = serde_json::to_string(&url::Url::from_file_path(path).unwrap().as_str())?;
483 module.push_str(&format!("import {uri} as service{index}\n"));
484 services.push(format!("{uri}, service{index}"));
485 }
486 module.push_str(&format!("local services = Map({})\n", services.join(", ")));
487 module.push_str(r#"output { renderer = new JsonRenderer { omitNullProperties = false }
488value = services.filter((_, s) -> s.enabled && s.meta.launcher).mapValues((_, s) ->
489let (tasks = s.containers?.toMap() ?? Map("app", s.container))
490let (routes = tasks.filter((_, task) -> task?.http?.hostname != null).mapValues((_, task) -> task.http))
491let (route = routes.values.firstOrNull)
492new Dynamic { name = s.meta.name; tagline = s.meta.tagline; hostname = route?.hostname; access = s.meta.access ?? route?.authRole; urls = routes.mapValues((_, http) -> "https://\(http.hostname)") }) }"#);
493 Ok(module)
494}
495
496async fn launcher(app: Arc<App>) -> Result<Arc<Document>> {
497 let real = tokio::fs::canonicalize(&app.repo).await?;
498 let state = app.clone();
499 app.cache.get(format!("launcher:{}",real.display()),Duration::from_secs(365*86400),move || async move {
500 let module = launcher_module(&real).await?;
501 let value: Value = serde_json::from_slice(&command("pkl",&["eval","-"],Some(module.as_bytes())).await?)?;
502 let mut apps = Vec::new();
503 for (file, value) in value.as_object().into_iter().flat_map(|v| v.iter()) {
504 if let Some(host) = value["hostname"].as_str() {
505 let path = std::path::Path::new(file);
506 let id = if path.file_name().unwrap() == "service.pkl" { path.parent().unwrap().file_name().unwrap() } else { path.file_stem().unwrap() }.to_string_lossy();
507 apps.push(json!({"id":id,"name":value["name"],"tagline":value["tagline"].as_str().filter(|s| !s.is_empty()),"url":format!("https://{host}"),"urls":value["urls"],"access":value["access"],"icon":icon_of(state.clone(),&id).await}));
508 }
509 }
510 Ok(json!(apps))
511 }).await
512}
513
514async fn service(app: Arc<App>, id: &str) -> Result<Arc<Document>> {
515 let state = app.clone();
516 let id = id.to_owned();
517 app.cache.get(format!("service:{id}"), Duration::from_secs(2), move || async move {
518 let app = state;
519 let id = id.as_str();
520 let ids = managed().await?;
521 let jobs = scan(app.clone()).await?;
522 let found = &jobs.value[id];
523 if !ids.iter().any(|s| s == id) || found.is_null() {
524 return Err(Error::new(
525 404,
526 format!("No service is named {id}. Pick one from the sidebar."),
527 ));
528 }
529 let all = summaries(app.clone()).await?;
530 let link = |other: &str, kind: &Value| {
531 array(&all.value)
532 .iter()
533 .find(|s| s["id"] == other)
534 .map(|s| json!({"id":other,"name":s["name"],"kind":kind,"health":s["health"]}))
535 };
536 let needs = |other: &str| {
537 serde_json::from_str::<Value>(string(&jobs.value[other]["job"]["Meta"]["studio_requires"]))
538 .ok()
539 };
540 let requirements = needs(id).map(|n| {
541 n.as_object()
542 .into_iter()
543 .flat_map(|v| v.iter())
544 .filter_map(|(id, kind)| link(id, kind))
545 .collect::<Vec<_>>()
546 });
547 let dependents = if ids.iter().all(|id| needs(id).is_some()) {
548 Some(
549 ids.iter()
550 .filter_map(|other| {
551 needs(other).and_then(|n| n.get(id).and_then(|kind| link(other, kind)))
552 })
553 .collect::<Vec<_>>(),
554 )
555 } else {
556 None
557 };
558 let mut summary = summarize(app.clone(), id, found).await;
559 let meta = &found["job"]["Meta"];
560 let application_traces = if let Some(name) = meta["studio_trace_service"]
561 .as_str()
562 .filter(|s| !s.is_empty())
563 {
564 telemetry::has_traces(app.clone(), name)
565 .await
566 .unwrap_or(false)
567 } else {
568 false
569 };
570 let metrics =
571 if array(&found["job"]["TaskGroups"]).iter().flat_map(|g| array(&g["Services"])).any(|s| array(&s["Tags"]).iter().any(|tag| string(tag).starts_with("studio-metrics-path="))) || meta["studio_metrics_pushed"] == "true" {
572 telemetry::metric_names(app.clone(), id)
573 .await
574 .unwrap_or(json!([]))
575 } else {
576 json!([])
577 };
578 let release = host::call(json!({"operation":"deploy.current"})).await?;
579 let secrets = serde_json::from_str::<Value>(string(&meta["studio_secrets"]))
580 .ok()
581 .map(|v| {
582 array(&v)
583 .iter()
584 .map(|s| json!({"name":s["name"],"generated":s["generated"]}))
585 .collect::<Vec<_>>()
586 });
587 let ack = read_json(&app.data.join("restarts-acknowledged.json"), json!({})).await?;
588 let datasets = app.cache.get("zfs-datasets".into(),Duration::from_secs(30),move || async move {
589 let rows = host::call(json!({"operation":"storage.datasets"})).await?;
590 Ok(json!(array(&rows).iter().map(|d| json!({"name":d["name"],"mountpoint":d["mountpoint"],"used":d["usedbydataset"],"snapshots":d["usedbysnapshots"]})).collect::<Vec<_>>()))
591 }).await?;
592 let logs = array(&all.value).iter().find(|s| s["id"] == "victoria" && s["url"].is_string()).map(|s| json!({"app":{"id":s["id"],"name":s["name"],"icon":s["icon"]},"url":format!("{}/select/vmui/#/?query={}",string(&s["url"]),encoded(&format!("{{job=\"{id}\"}}")))}));
593 let fields = json!({"applicationTraces":application_traces,"applicationMetrics":metrics,"release":release,"deployedAt":number(&found["job"]["SubmitTime"])/1e9,"rollout":if array(&found["job"]["TaskGroups"]).iter().any(|g| number(&g["Update"]["Canary"]) > 0.0) {"overlapped"} else {"simple"},"containers":containers(found),"checks":checks_of(&json!(array(&found["allocs"]).iter().flat_map(|a| array(&a["home_checks"])).collect::<Vec<_>>())),"requirements":requirements,"dependents":dependents,"secrets":secrets,"datasets":array(&datasets.value).iter().filter(|d| string(&d["name"]).ends_with(&format!("/prod/{id}"))).collect::<Vec<_>>(),"restartsAcknowledged":ack[id],"logs":logs});
594 summary
595 .as_object_mut()
596 .unwrap()
597 .extend(fields.as_object().unwrap().clone());
598 Ok(summary)
599 }).await
600}
601
602async fn change(app: Arc<App>, request: Value) -> Result<()> {
603 let run = host::call(request).await?;
604 app.cache.invalidate("deploys");
605 let outcome = tokio::time::timeout(Duration::from_secs(120), async {
606 loop {
607 let status = host::call(json!({"operation":"deploy.run","id":run["id"]})).await?;
608 if let Some(code) = status["code"].as_i64() {
609 return if code == 0 {
610 Ok(())
611 } else {
612 Err(Error::new(
613 502,
614 "The change failed. Check its run in deploys before retrying.",
615 ))
616 };
617 }
618 tokio::time::sleep(Duration::from_millis(500)).await;
619 }
620 })
621 .await
622 .map_err(|_| {
623 Error::new(
624 504,
625 "The change is still running. Check its run in deploys before retrying.",
626 )
627 });
628 app.cache.invalidate("deploys");
629 app.cache.invalidate("nomad-scan");
630 app.cache.invalidate("services");
631 app.cache
632 .invalidate(&format!("service:{}", string(&run["target"])));
633 outcome??;
634 Ok(())
635}
636
637pub async fn route(
638 app: Arc<App>,
639 method: &Method,
640 parts: &[&str],
641 query: &HashMap<String, String>,
642 me: &Value,
643 body: Value,
644) -> Result<Response> {
645 let empty = || StatusCode::NO_CONTENT.into_response();
646 let value = match parts {
647 ["me"] if method == Method::GET => me.clone(),
648 ["host"] if method == Method::GET => {
649 let sample = host::sample(app).await?;
650 json!({"cores":sample.value["cores"],"memory":sample.value["memory"]["total"],"bootedAt":sample.value["bootedAt"],"outages":null})
651 }
652 ["services"] if method == Method::GET => return Ok(summaries(app).await?.response()),
653 ["launcher"] if method == Method::GET => {
654 let found = launcher(app.clone()).await?;
655 let jobs = app.cache.peek("nomad-scan", Duration::from_secs(10));
656 let groups: Vec<_> = array(&me["groups"]).iter().map(string).collect();
657 let mut apps = Vec::new();
658 for value in array(&found.value)
659 .iter()
660 .filter(|v| can_open(&groups, v["access"].as_str()))
661 {
662 let mut value = value.clone();
663 value["health"] = jobs
664 .as_ref()
665 .map(|j| j.value[string(&value["id"])]["health"].clone())
666 .unwrap_or(Value::Null);
667 value.as_object_mut().unwrap().remove("access");
668 apps.push(value);
669 }
670 json!(apps)
671 }
672 ["status"] if method == Method::GET => {
673 let issues = issues(app.clone()).await.ok();
674 let apps = launcher(app.clone()).await?;
675 let groups: Vec<_> = array(&me["groups"]).iter().map(string).collect();
676 let admin = array(&me["sections"]).iter().any(|s| s == "admin");
677 let issues = issues.map(|v| {
678 array(&v)
679 .iter()
680 .filter(|i| {
681 admin
682 || array(&apps.value).iter().any(|a| {
683 a["id"] == i["service"]["id"]
684 && can_open(&groups, a["access"].as_str())
685 })
686 })
687 .map(|issue| {
688 let mut issue = issue.clone();
689 if !admin {
690 issue["failing"] = Value::Null;
691 }
692 issue
693 })
694 .collect::<Vec<_>>()
695 });
696 let sample = host::sample(app).await?;
697 json!({"load":(number(&sample.value["load"])/number(&sample.value["cores"])*100.0).min(100.0),"issues":issues})
698 }
699 ["services", id] if method == Method::GET => return Ok(service(app, id).await?.response()),
700 ["services", id, "definition"] if method == Method::GET => definition(app, id).await?,
701 ["services", id, "restarts", "acknowledge"] if method == Method::POST => {
702 let file = app.data.join("restarts-acknowledged.json");
703 let mut value = read_json(&file, json!({})).await?;
704 value[*id] = json!(now());
705 write_json(&file, &value).await?;
706 app.cache.invalidate(&format!("service:{id}"));
707 return Ok(empty());
708 }
709 ["services", id, action] if method == Method::POST => {
710 if !["start", "stop", "restart"].contains(action) {
711 return Err(Error::new(400, "Choose start, stop, or restart."));
712 }
713 change(
714 app,
715 json!({"operation":"deploy.start","action":action,"target":id}),
716 )
717 .await?;
718 return Ok(empty());
719 }
720 ["services", id, "secrets", name] if method == Method::GET => {
721 json!({"value":host::call(json!({"operation":"deploy.secret.get","service":id,"key":name})).await?})
722 }
723 ["services", id, "secrets", name] | ["services", id, "secrets", name, "rotate"]
724 if (parts.len() == 4 && method == Method::PUT)
725 || (parts.len() == 5 && method == Method::POST) =>
726 {
727 let mut request = json!({"operation":if parts.len()==5 {"deploy.secret.rotate"} else {"deploy.secret.set"},"service":id,"key":name});
728 if parts.len() == 4 {
729 request["value"] = body["value"].clone();
730 }
731 change(app, request).await?;
732 return Ok(empty());
733 }
734 ["metrics", metric] if method == Method::GET => {
735 return Ok(telemetry::metrics(app, metric, query, None, None)
736 .await?
737 .response());
738 }
739 ["services", id, "metrics", name] if method == Method::GET => {
740 let detail = service(app.clone(), id).await?;
741 if !array(&detail.value["applicationMetrics"])
742 .iter()
743 .any(|v| v == *name)
744 {
745 return Err(Error::new(404, "Metric unavailable for this service."));
746 }
747 return Ok(telemetry::metrics(app, name, query, Some(id), None)
748 .await?
749 .response());
750 }
751 ["services", id, "logs"] if method == Method::GET => {
752 telemetry::logs(app, id, query).await?
753 }
754 ["services", id, "traces"] if method == Method::GET => {
755 telemetry::traces(app, id, query, |_| true).await?
756 }
757 ["traces", id] if method == Method::GET => telemetry::trace(app, id).await?,
758 _ => return Err(Error::new(404, "Not Found")),
759 };
760 Ok(Document::new(value).response())
761}
762
763pub fn start(app: Arc<App>) {
764 let state = app.clone();
765 tokio::spawn(async move {
766 loop {
767 if let Err(e) = summaries(state.clone()).await {
768 eprintln!("status: {}", e.message);
769 }
770 tokio::time::sleep(Duration::from_secs(2)).await;
771 }
772 });
773 tokio::spawn(async move {
774 let mut previous = (0.0, 0.0);
775 let mut interval = tokio::time::interval(Duration::from_secs(2));
776 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
777 loop {
778 interval.tick().await;
779 let sample = async {
780 let sample = host::sample(app.clone()).await?;
781 let sample = &sample.value;
782 let total = number(&sample["cpu"]["total"]); let busy = number(&sample["cpu"]["busy"]);
783 let cpu = if previous.1 > 0.0 && total > previous.1 { ((busy-previous.0)/(total-previous.1)*100.0).clamp(0.0,100.0) } else {0.0}; previous = (busy,total);
784 let jobs = app.cache.peek("services",Duration::from_secs(2));
785 let services: serde_json::Map<_,_> = jobs.as_ref().into_iter().flat_map(|j| array(&j.value)).map(|s| (string(&s["id"]).to_owned(),json!({"cpu":s["cpu"],"memory":s["memory"],"health":s["health"]}))).collect();
786 Ok::<_,Error>(Document::new(json!({"t":now(),"host":{"cpu":cpu,"memory":sample["memory"]["used"],"arc":sample["arc"],"temperature":sample["temperature"]},"services":services,"ups":null})))
787 }.await;
788 match sample {
789 Ok(value) => {
790 app.live.send_replace(value.bytes);
791 }
792 Err(e) => eprintln!("host sample: {}", e.message),
793 }
794 }
795 });
796}
797
798async fn definition(app: Arc<App>, id: &str) -> Result<Value> {
799 if !managed().await?.iter().any(|s| s == id) {
800 return Err(Error::new(
801 404,
802 format!("No service is named {id}. Pick one from the sidebar."),
803 ));
804 }
805 let jobs = scan(app.clone()).await?;
806 let found = &jobs.value[id];
807 if found.is_null() {
808 return Err(Error::new(
809 404,
810 format!("No service is named {id}. Pick one from the sidebar."),
811 ));
812 }
813 let mhz = number(&found["mhz"]);
814 let groups=array(&found["job"]["TaskGroups"]).iter().map(|group| {
815 let tags=|service: &Value,key: &str| array(&service["Tags"]).iter().filter_map(|s| string(s).strip_prefix(&format!("{key}=")).map(str::to_owned)).collect::<Vec<_>>();
816 let services=array(&group["Services"]).iter().map(|service| {
817 let check=&service["Checks"][0];let restart=&check["CheckRestart"];
818 json!({"name":service["Name"],"port":service["PortLabel"],"hostnames":tags(service,"caddy-host"),"authRole":tags(service,"caddy-auth-role").first(),"check":if check.is_null() {Value::Null} else {json!({"task":check["TaskName"],"type":check["Type"],"path":check["Path"],"interval":number(&check["Interval"])/1e9,"timeout":number(&check["Timeout"])/1e9,"restartAfter":if number(&restart["Limit"])>0.0 {json!({"failures":restart["Limit"],"grace":number(&restart["Grace"])/1e9})} else {Value::Null}})}})
819 }).collect::<Vec<_>>();
820 let tasks=array(&group["Tasks"]).iter().map(|task| {
821 let mut ports=HashMap::new();for network in array(&group["Networks"]) {for (key,dynamic) in [("DynamicPorts",true),("ReservedPorts",false)] {for port in array(&network[key]) {ports.insert(string(&port["Label"]),(port,dynamic));}}}
822 let config=&task["Config"];let mut env:Vec<Value>=task["Env"].as_object().into_iter().flat_map(|v| v.keys()).map(|name| json!({"name":name,"from":"job"})).collect();
823 let assignments=regex::Regex::new(r"(?m)(?:^|\}\})([A-Za-z_]\w*)=").unwrap();
824 for template in array(&task["Templates"]).iter().filter(|t| t["Envvars"]==true) {let text=string(&template["EmbeddedTmpl"]);for capture in assignments.captures_iter(text) {env.push(json!({"name":&capture[1],"from":if text.contains("nomadVar") {"secret"} else {"template"}}));}}
825 json!({"name":task["Name"],"hook":task["Lifecycle"]["Hook"],"sidecar":task["Lifecycle"]["Sidecar"].as_bool().unwrap_or(false),"image":config["image"],"user":task["User"].as_str().filter(|s| !s.is_empty()),"cpu":number(&task["Resources"]["CPU"])/mhz,"memory":number(&task["Resources"]["MemoryMB"])*1048576.0,"memoryMax":if number(&task["Resources"]["MemoryMaxMB"])>0.0 {json!(number(&task["Resources"]["MemoryMaxMB"])*1048576.0)} else {Value::Null},"ports":array(&config["ports"]).iter().map(|label| {let port=ports.get(string(label));json!({"label":label,"container":port.and_then(|(p,_)| p["To"].as_u64()).filter(|n| *n>0),"host":port.filter(|(_,dynamic)| !dynamic).map(|(p,_)| p["Value"].clone()),"network":port.map(|(p,_)| p["HostNetwork"].as_str().unwrap_or("default")).unwrap_or("default")})}).collect::<Vec<_>>(),"mounts":array(&config["volumes"]).iter().map(|volume| {let fields:Vec<_>=string(volume).split(':').collect();json!({"source":fields[0],"target":fields.get(1).unwrap_or(&fields[0]),"readOnly":fields.get(2).is_some_and(|s| s.split(',').any(|s| s=="ro"))})}).collect::<Vec<_>>(),"tmpfs":array(&config["tmpfs"]),"devices":array(&config["devices"]),"capabilities":array(&config["cap_add"]),"hostNetwork":config["network_mode"]=="host","extraHosts":array(&config["extra_hosts"]),"env":env})
826 }).collect::<Vec<_>>();
827 let restart=&group["RestartPolicy"];
828 json!({"name":group["Name"],"count":group["Count"],"restart":{"attempts":restart["Attempts"],"interval":number(&restart["Interval"])/1e9,"delay":number(&restart["Delay"])/1e9,"fail":restart["Mode"]=="fail"},"services":services,"tasks":tasks})
829 }).collect::<Vec<_>>();
830 let source = tokio::fs::read_to_string(service_file(&app, id).await?)
831 .await
832 .ok();
833 Ok(json!({"groups":groups,"source":source}))
834}
835
836#[cfg(test)]
837mod tests {
838 use super::*;
839 #[tokio::test]
840 async fn launcher_excludes_directories_and_grouped_definitions_before_import() {
841 let root = std::env::temp_dir().join(format!("studio-launcher-{}", uuid::Uuid::new_v4()));
842 tokio::fs::create_dir_all(root.join("config")).await.unwrap();
843 for directory in ["retired", "personal"] {
844 tokio::fs::create_dir_all(root.join("service").join(directory))
845 .await
846 .unwrap();
847 }
848 tokio::fs::write(root.join("config/excluded-services.json"), r#"["retired"]"#)
849 .await
850 .unwrap();
851 for file in [
852 "retired/service.pkl",
853 "personal/retired.pkl",
854 "personal/service.pkl",
855 "personal/fixture.pkl",
856 ] {
857 tokio::fs::write(root.join("service").join(file), "unread fixture")
858 .await
859 .unwrap();
860 }
861 let module = launcher_module(&root).await.unwrap();
862 assert!(!module.contains("retired"));
863 assert!(!module.contains("import*"));
864 assert!(module.contains("/personal/service.pkl"));
865 assert!(module.contains("/personal/fixture.pkl"));
866 assert_eq!(
867 module.lines().filter(|line| line.starts_with("import ")).count(),
868 2,
869 );
870 tokio::fs::remove_dir_all(root).await.unwrap();
871 }
872 #[test]
873 fn health_precedence() {
874 let job = json!({"Stop":false});
875 let alloc = json!({"DesiredStatus":"run","ClientStatus":"running","JobVersion":1,"TaskStates":{"app":{"State":"running"}}});
876 assert_eq!(health(&job, &[], &[]), "down");
877 assert_eq!(health(&job, &[alloc.clone()], &[]), "healthy");
878 assert_eq!(
879 health(&job, &[alloc.clone()], &[json!({"Status":"failure"})]),
880 "degraded"
881 );
882 let mut pending = alloc.clone();
883 pending["TaskStates"]["app"]["State"] = json!("pending");
884 assert_eq!(
885 health(&job, &[pending], &[json!({"Status":"failure"})]),
886 "restarting"
887 );
888 let mut newer = alloc.clone();
889 newer["JobVersion"] = json!(2);
890 assert_eq!(health(&job, &[alloc, newer], &[]), "deploying");
891 assert_eq!(health(&json!({"Stop":true}), &[], &[]), "stopped");
892 }
893}