| 1 | use crate::*; |
| 2 | use futures::{StreamExt, stream}; |
| 3 | use sha1::{Digest, Sha1}; |
| 4 | use std::path::Path; |
| 5 | |
| 6 | pub 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 | } |
| 13 | pub 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 | |
| 56 | fn live_alloc(alloc: &Value) -> bool { |
| 57 | alloc["DesiredStatus"] == "run" |
| 58 | && matches!(string(&alloc["ClientStatus"]), "pending" | "running") |
| 59 | } |
| 60 | pub 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 | |
| 96 | pub 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 | |
| 198 | pub 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 | } |
| 207 | pub 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 | } |
| 214 | pub 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 | } |
| 221 | pub 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 | } |
| 251 | pub 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 | |
| 259 | async 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 | } |
| 276 | async 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 | } |
| 295 | pub 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 | |
| 327 | async 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 | } |
| 355 | async 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 | |
| 378 | pub 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 | } |
| 384 | pub 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 | } |
| 406 | pub 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 | } |
| 409 | fn 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 | } |
| 426 | async 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 | |
| 447 | async 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 | |
| 453 | async 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 } |
| 488 | value = services.filter((_, s) -> s.enabled && s.meta.launcher).mapValues((_, s) -> |
| 489 | let (tasks = s.containers?.toMap() ?? Map("app", s.container)) |
| 490 | let (routes = tasks.filter((_, task) -> task?.http?.hostname != null).mapValues((_, task) -> task.http)) |
| 491 | let (route = routes.values.firstOrNull) |
| 492 | new 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 | |
| 496 | async 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 | |
| 514 | async 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 | |
| 602 | async 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 | |
| 637 | pub 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 | |
| 763 | pub 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 | |
| 798 | async 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)] |
| 837 | mod 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 | } |