| 1 | use crate::*; |
| 2 | use futures::{StreamExt, stream}; |
| 3 | |
| 4 | pub async fn endpoint(app: Arc<App>, service: &str) -> Result<String> { |
| 5 | if let Some((base, _)) = &app.internal { |
| 6 | return Ok(format!("{}services/{}", base.as_str(), encoded(service))); |
| 7 | } |
| 8 | let id = service.to_owned(); |
| 9 | let state = app.clone(); |
| 10 | let found = app |
| 11 | .cache |
| 12 | .get( |
| 13 | format!("endpoint:{service}"), |
| 14 | Duration::from_secs(10), |
| 15 | move || async move { |
| 16 | let values = core::nomad(&state, &format!("/v1/service/{}", encoded(&id))).await?; |
| 17 | let instance = array(&values) |
| 18 | .iter() |
| 19 | .find(|v| v["Address"] == "127.0.0.1") |
| 20 | .or_else(|| array(&values).first()) |
| 21 | .ok_or_else(|| { |
| 22 | Error::new(501, format!("{id} isn't running on this host yet.")) |
| 23 | })?; |
| 24 | let address = string(&instance["Address"]); |
| 25 | Ok(json!(format!( |
| 26 | "http://{}:{}", |
| 27 | if address.contains(':') { |
| 28 | format!("[{address}]") |
| 29 | } else { |
| 30 | address.into() |
| 31 | }, |
| 32 | number(&instance["Port"]) as u16 |
| 33 | ))) |
| 34 | }, |
| 35 | ) |
| 36 | .await?; |
| 37 | Ok(string(&found.value).to_owned()) |
| 38 | } |
| 39 | |
| 40 | pub async fn read( |
| 41 | app: Arc<App>, |
| 42 | url: String, |
| 43 | form: Option<Vec<(String, String)>>, |
| 44 | ) -> Result<Arc<Document>> { |
| 45 | let state = app.clone(); |
| 46 | let key = format!("read:{url}:{}", serde_json::to_string(&form)?); |
| 47 | app.cache |
| 48 | .get(key, Duration::from_secs(2), move || async move { |
| 49 | let request = if let Some(form) = &form { |
| 50 | state.request(Method::POST, &url)?.form(form) |
| 51 | } else { |
| 52 | state.request(Method::GET, &url)? |
| 53 | }; |
| 54 | let response = request |
| 55 | .timeout(Duration::from_secs(8)) |
| 56 | .send() |
| 57 | .await |
| 58 | .map_err(|_| Error::new(502, "The telemetry store isn't answering."))?; |
| 59 | if !response.status().is_success() { |
| 60 | return Err(Error::new( |
| 61 | 502, |
| 62 | format!("The telemetry store answered {}.", response.status()), |
| 63 | )); |
| 64 | } |
| 65 | if form.is_some() { |
| 66 | Ok(json!(response.text().await?)) |
| 67 | } else { |
| 68 | Ok(response.json().await?) |
| 69 | } |
| 70 | }) |
| 71 | .await |
| 72 | } |
| 73 | |
| 74 | const METRICS: &[(&str, &str)] = &[ |
| 75 | ("host.cpu", "studio_host_cpu_percent"), |
| 76 | ("host.memory", "studio_host_memory_bytes"), |
| 77 | ("host.temperature", "studio_host_temperature_celsius"), |
| 78 | ("host.power", "studio_host_power_watts"), |
| 79 | ("host.gpu", "studio_host_gpu_percent"), |
| 80 | ("host.network", "studio_host_network_bytes_per_second"), |
| 81 | ("service.cpu", "studio_service_cpu_cores"), |
| 82 | ("service.memory", "studio_service_memory_bytes"), |
| 83 | ("vm.cpu", "studio_vm_cpu_percent"), |
| 84 | ("vm.memory", "studio_vm_memory_bytes"), |
| 85 | ]; |
| 86 | |
| 87 | fn series(result: &Value, application: bool) -> Value { |
| 88 | json!( |
| 89 | array(&result["data"]["result"]) |
| 90 | .iter() |
| 91 | .take(if application { 12 } else { usize::MAX }) |
| 92 | .map(|row| { |
| 93 | let metric = &row["metric"]; |
| 94 | let name = if application { |
| 95 | let labels = metric |
| 96 | .as_object() |
| 97 | .into_iter() |
| 98 | .flat_map(|v| v.iter()) |
| 99 | .filter(|(k, _)| { |
| 100 | !["__name__", "service", "service_name"].contains(&k.as_str()) |
| 101 | }) |
| 102 | .map(|(k, v)| format!("{k}={}", string(v))) |
| 103 | .collect::<Vec<_>>() |
| 104 | .join(" · "); |
| 105 | if labels.is_empty() { |
| 106 | "total".into() |
| 107 | } else { |
| 108 | labels |
| 109 | } |
| 110 | } else { |
| 111 | metric["service"] |
| 112 | .as_str() |
| 113 | .or(metric["vm"].as_str()) |
| 114 | .or(metric["direction"].as_str()) |
| 115 | .unwrap_or("total") |
| 116 | .into() |
| 117 | }; |
| 118 | let t: Vec<_> = array(&row["values"]).iter().map(|v| v[0].clone()).collect(); |
| 119 | let v: Vec<_> = array(&row["values"]) |
| 120 | .iter() |
| 121 | .map(|v| { |
| 122 | let n = string(&v[1]).parse::<f64>().unwrap_or(f64::NAN); |
| 123 | if n.is_finite() { json!(n) } else { Value::Null } |
| 124 | }) |
| 125 | .collect(); |
| 126 | json!({"name":name,"t":t,"v":v}) |
| 127 | }) |
| 128 | .collect::<Vec<_>>() |
| 129 | ) |
| 130 | } |
| 131 | |
| 132 | pub fn query_number( |
| 133 | query: &HashMap<String, String>, |
| 134 | key: &str, |
| 135 | default: f64, |
| 136 | min: f64, |
| 137 | max: f64, |
| 138 | ) -> Result<f64> { |
| 139 | let value = query |
| 140 | .get(key) |
| 141 | .map(|v| v.parse::<f64>()) |
| 142 | .transpose() |
| 143 | .map_err(|_| Error::new(400, format!("Invalid {key}.")))? |
| 144 | .unwrap_or(default); |
| 145 | if !value.is_finite() || value < min || value > max { |
| 146 | return Err(Error::new(400, format!("Invalid {key}."))); |
| 147 | } |
| 148 | Ok(value) |
| 149 | } |
| 150 | pub async fn metrics( |
| 151 | app: Arc<App>, |
| 152 | name: &str, |
| 153 | query: &HashMap<String, String>, |
| 154 | application: Option<&str>, |
| 155 | window: Option<(f64, f64)>, |
| 156 | ) -> Result<Arc<Document>> { |
| 157 | let range = query_number(query, "range", 3600.0, 60.0, 30.0 * 86400.0)?; |
| 158 | let to = window |
| 159 | .map(|(_, to)| to) |
| 160 | .unwrap_or_else(|| (now() / 2.0).floor() * 2.0); |
| 161 | let from = window.map(|(from, _)| from).unwrap_or(to - range); |
| 162 | let service = query.get("service"); |
| 163 | let quote = |s: &str| serde_json::to_string(s).unwrap(); |
| 164 | let selector = if let Some(id) = application { |
| 165 | format!( |
| 166 | "{name}{{service={}}} or {name}{{service_name={}}}", |
| 167 | quote(id), |
| 168 | quote(id) |
| 169 | ) |
| 170 | } else { |
| 171 | let metric = METRICS |
| 172 | .iter() |
| 173 | .find(|(key, _)| *key == name) |
| 174 | .ok_or_else(|| Error::new(400, "Invalid metric."))? |
| 175 | .1; |
| 176 | format!( |
| 177 | "{metric}{}", |
| 178 | service |
| 179 | .map(|id| format!( |
| 180 | "{{{}={}}}", |
| 181 | if name.starts_with("vm.") { |
| 182 | "vm" |
| 183 | } else { |
| 184 | "service" |
| 185 | }, |
| 186 | quote(id) |
| 187 | )) |
| 188 | .unwrap_or_default() |
| 189 | ) |
| 190 | }; |
| 191 | let key = if window.is_some() { |
| 192 | format!("series:{selector}:fixed:{from}:{to}") |
| 193 | } else { |
| 194 | format!("series:{selector}:live:{range}") |
| 195 | }; |
| 196 | let state = app.clone(); |
| 197 | let network = name == "host.network"; |
| 198 | let application = application.is_some(); |
| 199 | app.cache |
| 200 | .get(key, Duration::from_secs(2), move || async move { |
| 201 | let base = endpoint(state.clone(), "victoria-metrics").await?; |
| 202 | let query = params(&[ |
| 203 | ("query", selector), |
| 204 | ("start", from.to_string()), |
| 205 | ("end", to.to_string()), |
| 206 | ("step", ((to - from) / 300.0).round().max(2.0).to_string()), |
| 207 | ]); |
| 208 | let response = read(state, format!("{base}/api/v1/query_range?{query}"), None).await?; |
| 209 | if response.value["status"] != "success" { |
| 210 | return Err(Error::new(502, "The metrics query failed.")); |
| 211 | } |
| 212 | let mut values = series(&response.value, application); |
| 213 | if network { |
| 214 | values |
| 215 | .as_array_mut() |
| 216 | .unwrap() |
| 217 | .sort_by_key(|s| if s["name"] == "down" { 0 } else { 1 }); |
| 218 | } |
| 219 | Ok(values) |
| 220 | }) |
| 221 | .await |
| 222 | } |
| 223 | pub async fn metric_names(app: Arc<App>, id: &str) -> Result<Value> { |
| 224 | let query = params(&["service", "service_name"].map(|label| { |
| 225 | ( |
| 226 | "match[]", |
| 227 | format!( |
| 228 | "{{{label}={},__name__!~\"studio_.*\"}}", |
| 229 | serde_json::to_string(id).unwrap() |
| 230 | ), |
| 231 | ) |
| 232 | })); |
| 233 | let url = format!( |
| 234 | "{}/api/v1/label/__name__/values?{query}", |
| 235 | endpoint(app.clone(), "victoria-metrics").await? |
| 236 | ); |
| 237 | let response = read(app, url, None).await?; |
| 238 | let valid = regex::Regex::new(r"^[a-zA-Z_:][a-zA-Z0-9_:]*$").unwrap(); |
| 239 | let mut values: Vec<_> = array(&response.value["data"]) |
| 240 | .iter() |
| 241 | .filter(|v| valid.is_match(string(v))) |
| 242 | .cloned() |
| 243 | .collect(); |
| 244 | values.sort_by(|a, b| string(a).cmp(string(b))); |
| 245 | Ok(json!(values)) |
| 246 | } |
| 247 | |
| 248 | #[derive(Debug, PartialEq)] |
| 249 | struct Term { |
| 250 | field: Option<String>, |
| 251 | op: String, |
| 252 | value: String, |
| 253 | exclude: bool, |
| 254 | } |
| 255 | fn terms(query: &str, trace: bool) -> Result<Vec<Term>> { |
| 256 | let tokens = |
| 257 | regex::Regex::new(r#"(-?)(?:([\w.]+):(>=|<=|>|<)?)?(?:"([^"]*)"?|(\S+))"#).unwrap(); |
| 258 | let wildcard = regex::Regex::new(r"(?i)^(\d+)x+$").unwrap(); |
| 259 | let mut found = Vec::new(); |
| 260 | for captures in tokens.captures_iter(query) { |
| 261 | let field = captures.get(2).map(|m| m.as_str()); |
| 262 | let recognized = field.filter(|s| { |
| 263 | if trace { |
| 264 | ["status", "method", "path", "host", "client", "duration"].contains(s) |
| 265 | || s.contains('.') |
| 266 | } else { |
| 267 | ["container", "stream"].contains(s) |
| 268 | } |
| 269 | }); |
| 270 | let value = if field.is_some() && recognized.is_none() { |
| 271 | captures[0].trim_start_matches('-').to_owned() |
| 272 | } else { |
| 273 | captures |
| 274 | .get(4) |
| 275 | .or(captures.get(5)) |
| 276 | .map(|m| m.as_str()) |
| 277 | .unwrap_or_default() |
| 278 | .to_owned() |
| 279 | }; |
| 280 | let value = wildcard.replace(&value, "$1").into_owned(); |
| 281 | if value.is_empty() { |
| 282 | continue; |
| 283 | } |
| 284 | found.push(Term { |
| 285 | field: recognized.map(str::to_owned), |
| 286 | op: if recognized.is_some() { |
| 287 | captures.get(3).map(|m| m.as_str()).unwrap_or(":") |
| 288 | } else { |
| 289 | ":" |
| 290 | } |
| 291 | .into(), |
| 292 | value, |
| 293 | exclude: &captures[1] == "-", |
| 294 | }); |
| 295 | } |
| 296 | Ok(found) |
| 297 | } |
| 298 | pub async fn logs(app: Arc<App>, id: &str, query: &HashMap<String, String>) -> Result<Value> { |
| 299 | let limit = query_number(query, "limit", 200.0, 1.0, 5000.0)?; |
| 300 | let mut filters = vec![ |
| 301 | "source:=nomad".to_owned(), |
| 302 | format!("job:={}", serde_json::to_string(id)?), |
| 303 | ]; |
| 304 | if let Some(level) = query.get("level") { |
| 305 | filters.push( |
| 306 | match level.as_str() { |
| 307 | "error" => "level:=error", |
| 308 | "warn" => "(level:=warn OR level:=error)", |
| 309 | _ => return Err(Error::new(400, "Invalid log level.")), |
| 310 | } |
| 311 | .into(), |
| 312 | ); |
| 313 | } |
| 314 | for term in terms( |
| 315 | query.get("q").map(String::as_str).unwrap_or_default(), |
| 316 | false, |
| 317 | )? { |
| 318 | let field = match term.field.as_deref() { |
| 319 | Some("container") => "task", |
| 320 | Some("stream") => "stream", |
| 321 | _ => "_msg", |
| 322 | }; |
| 323 | filters.push(format!( |
| 324 | "{}{field}:~{}", |
| 325 | if term.exclude { "!" } else { "" }, |
| 326 | serde_json::to_string(&format!( |
| 327 | "(?i){}{}", |
| 328 | if term.field.is_some() { "^" } else { "" }, |
| 329 | regex::escape(&term.value) |
| 330 | ))? |
| 331 | )); |
| 332 | } |
| 333 | let mut form = vec![ |
| 334 | ("query".into(), filters.join(" ")), |
| 335 | ("limit".into(), limit.to_string()), |
| 336 | ]; |
| 337 | for (key, field, add) in [("after", "start", 0.000001), ("before", "end", 0.0)] { |
| 338 | if query.contains_key(key) { |
| 339 | form.push(( |
| 340 | field.into(), |
| 341 | (query_number(query, key, 0.0, f64::NEG_INFINITY, f64::INFINITY)? + add) |
| 342 | .to_string(), |
| 343 | )); |
| 344 | } |
| 345 | } |
| 346 | let url = format!( |
| 347 | "{}/select/logsql/query", |
| 348 | endpoint(app.clone(), "victoria-logs").await? |
| 349 | ); |
| 350 | let response = read(app, url, Some(form)).await?; |
| 351 | let mut rows = string(&response.value).lines().filter(|l| !l.trim().is_empty()).map(|line| { |
| 352 | let row: Value = serde_json::from_str(line)?; |
| 353 | Ok::<_,Error>(json!({"t":core::seconds(&row["_time"]),"container":row["task"].as_str().unwrap_or("app"),"stream":if row["stream"] == "stderr" {"stderr"} else {"stdout"},"level":if row["level"] == "error" || row["level"] == "warn" {row["level"].clone()} else {Value::Null},"text":row["_msg"].as_str().unwrap_or_default()})) |
| 354 | }).collect::<Result<Vec<_>>>()?; |
| 355 | rows.sort_by(|a, b| number(&b["t"]).total_cmp(&number(&a["t"]))); |
| 356 | Ok(json!(rows)) |
| 357 | } |
| 358 | |
| 359 | fn span(row: &Value) -> Value { |
| 360 | let mut attributes: serde_json::Map<String, Value> = row |
| 361 | .as_object() |
| 362 | .into_iter() |
| 363 | .flat_map(|v| v.iter()) |
| 364 | .filter_map(|(key, value)| { |
| 365 | key.strip_prefix("span_attr:") |
| 366 | .map(|key| (key.into(), value.clone())) |
| 367 | }) |
| 368 | .collect(); |
| 369 | for (from, to) in [ |
| 370 | ("http.request.method", "http.request.method"), |
| 371 | ("http.response.status_code", "http.response.status_code"), |
| 372 | ("http.route", "url.path"), |
| 373 | ] { |
| 374 | if let Some(value) = row.get(format!("resource_attr:{from}")) { |
| 375 | attributes |
| 376 | .entry(to.to_owned()) |
| 377 | .or_insert_with(|| value.clone()); |
| 378 | } |
| 379 | } |
| 380 | let status = attributes |
| 381 | .get("http.response.status_code") |
| 382 | .or(attributes.get("http.status_code")) |
| 383 | .cloned() |
| 384 | .unwrap_or(Value::Null); |
| 385 | let service = attributes |
| 386 | .get("studio.service") |
| 387 | .cloned() |
| 388 | .unwrap_or_else(|| { |
| 389 | row["resource_attr:service.name"] |
| 390 | .as_str() |
| 391 | .map(|s| json!(s)) |
| 392 | .unwrap_or(json!("unknown")) |
| 393 | }); |
| 394 | let error = if row["status_code"] == "2" || number(&status) >= 500.0 { |
| 395 | if !status.is_null() { |
| 396 | status |
| 397 | } else { |
| 398 | json!( |
| 399 | row["status_message"] |
| 400 | .as_str() |
| 401 | .filter(|s| !s.is_empty()) |
| 402 | .unwrap_or("request failed") |
| 403 | ) |
| 404 | } |
| 405 | } else { |
| 406 | Value::Null |
| 407 | }; |
| 408 | let name = if attributes.get("studio.kind").is_some_and(|v| v == "edge") { |
| 409 | format!("edge: {}", row["name"].as_str().unwrap_or("request")) |
| 410 | } else { |
| 411 | string(&row["name"]).into() |
| 412 | }; |
| 413 | json!({"id":row["span_id"],"parent":row["parent_span_id"].as_str().filter(|s| !s.is_empty()),"service":service,"name":name,"start":number(&row["start_time_unix_nano"])/1e9,"duration":number(&row["duration"])/1e9,"error":error,"attributes":attributes}) |
| 414 | } |
| 415 | async fn spans(app: Arc<App>, query: String) -> Result<Vec<Value>> { |
| 416 | let url = format!( |
| 417 | "{}/select/logsql/query", |
| 418 | endpoint(app.clone(), "victoria-traces").await? |
| 419 | ); |
| 420 | let rows = read(app, url, Some(vec![("query".into(), query)])).await?; |
| 421 | let mut traces: HashMap<String, Vec<Value>> = HashMap::new(); |
| 422 | for line in string(&rows.value).lines().filter(|l| !l.trim().is_empty()) { |
| 423 | let row: Value = serde_json::from_str(line)?; |
| 424 | if string(&row["trace_id"]).is_empty() || string(&row["span_id"]).is_empty() { |
| 425 | continue; |
| 426 | } |
| 427 | traces |
| 428 | .entry(string(&row["trace_id"]).to_owned()) |
| 429 | .or_default() |
| 430 | .push(span(&row)); |
| 431 | } |
| 432 | Ok(traces |
| 433 | .into_iter() |
| 434 | .map(|(id, mut spans)| { |
| 435 | let parents: HashMap<_, _> = spans |
| 436 | .iter() |
| 437 | .map(|s| (string(&s["id"]).to_owned(), string(&s["parent"]).to_owned())) |
| 438 | .collect(); |
| 439 | let depth = |span: &Value| { |
| 440 | let mut parent = string(&span["parent"]); |
| 441 | let mut seen = std::collections::HashSet::new(); |
| 442 | let mut n = 0; |
| 443 | while let Some(next) = parents.get(parent) { |
| 444 | if !seen.insert(parent) { |
| 445 | break; |
| 446 | } |
| 447 | n += 1; |
| 448 | parent = next; |
| 449 | } |
| 450 | n |
| 451 | }; |
| 452 | spans.sort_by(|a, b| { |
| 453 | depth(a) |
| 454 | .cmp(&depth(b)) |
| 455 | .then_with(|| number(&a["start"]).total_cmp(&number(&b["start"]))) |
| 456 | }); |
| 457 | json!({"id":id,"spans":spans}) |
| 458 | }) |
| 459 | .collect()) |
| 460 | } |
| 461 | pub async fn trace(app: Arc<App>, id: &str) -> Result<Value> { |
| 462 | spans(app, format!("trace_id:={}", serde_json::to_string(id)?)) |
| 463 | .await? |
| 464 | .into_iter() |
| 465 | .next() |
| 466 | .ok_or_else(|| { |
| 467 | Error::new( |
| 468 | 404, |
| 469 | "That trace has aged out of the trace store. Pick a newer one.", |
| 470 | ) |
| 471 | }) |
| 472 | } |
| 473 | pub async fn has_traces(app: Arc<App>, service: &str) -> Result<bool> { |
| 474 | let url = format!( |
| 475 | "{}/select/logsql/query", |
| 476 | endpoint(app.clone(), "victoria-traces").await? |
| 477 | ); |
| 478 | Ok(!string( |
| 479 | &read( |
| 480 | app, |
| 481 | url, |
| 482 | Some(vec![( |
| 483 | "query".into(), |
| 484 | format!( |
| 485 | "_time:7d \"resource_attr:service.name\":={} | limit 1", |
| 486 | serde_json::to_string(service)? |
| 487 | ), |
| 488 | )]), |
| 489 | ) |
| 490 | .await? |
| 491 | .value, |
| 492 | ) |
| 493 | .trim() |
| 494 | .is_empty()) |
| 495 | } |
| 496 | pub async fn traces( |
| 497 | app: Arc<App>, |
| 498 | id: &str, |
| 499 | query: &HashMap<String, String>, |
| 500 | visible: impl Fn(&Value) -> bool, |
| 501 | ) -> Result<Value> { |
| 502 | let jobs = core::scan(app.clone()).await?; |
| 503 | let Some(service) = jobs.value[id]["job"]["Meta"]["studio_trace_service"] |
| 504 | .as_str() |
| 505 | .filter(|s| !s.is_empty()) |
| 506 | else { |
| 507 | return Ok(json!([])); |
| 508 | }; |
| 509 | let sort = query.get("sort").map(String::as_str).unwrap_or("newest"); |
| 510 | let (key, descending) = match sort { |
| 511 | "newest" => ("_time", true), |
| 512 | "oldest" => ("_time", false), |
| 513 | "slowest" => ("duration", true), |
| 514 | "fastest" => ("duration", false), |
| 515 | _ => return Err(Error::new(400, "Invalid trace order.")), |
| 516 | }; |
| 517 | let limit = query_number(query, "limit", 50.0, 1.0, 500.0)?; |
| 518 | let quote = |s: &str| serde_json::to_string(s).unwrap(); |
| 519 | let like = |field: &str, pattern: &str| { |
| 520 | format!("{}:~{}", quote(field), quote(&format!("(?i){pattern}"))) |
| 521 | }; |
| 522 | let mut filters = vec![ |
| 523 | "_time:7d".into(), |
| 524 | format!("\"resource_attr:service.name\":={}", quote(service)), |
| 525 | ]; |
| 526 | if query.get("errors").is_some_and(|v| v == "1") { |
| 527 | filters.push("(status_code:=2 OR \"span_attr:http.response.status_code\":>=500)".into()); |
| 528 | } |
| 529 | let duration = regex::Regex::new(r"^(\d+(?:\.\d+)?)(us|µs|ms|s|m)?$").unwrap(); |
| 530 | for term in terms(query.get("q").map(String::as_str).unwrap_or_default(), true)? { |
| 531 | let escaped = regex::escape(&term.value); |
| 532 | let filter = match term.field.as_deref() { |
| 533 | Some("duration") => { |
| 534 | let captures = duration.captures(&term.value).ok_or_else(|| { |
| 535 | Error::new( |
| 536 | 400, |
| 537 | format!("Write a duration like 500ms or 2s, not {}.", term.value), |
| 538 | ) |
| 539 | })?; |
| 540 | let nanos = match captures.get(2).map(|m| m.as_str()).unwrap_or("ms") { |
| 541 | "us" | "µs" => 1e3, |
| 542 | "s" => 1e9, |
| 543 | "m" => 60e9, |
| 544 | _ => 1e6, |
| 545 | }; |
| 546 | format!( |
| 547 | "duration:{}{}", |
| 548 | if term.op == ":" { ">=" } else { &term.op }, |
| 549 | (captures[1].parse::<f64>().unwrap() * nanos).round() |
| 550 | ) |
| 551 | } |
| 552 | None => format!( |
| 553 | "({})", |
| 554 | [ |
| 555 | "name", |
| 556 | "span_attr:url.path", |
| 557 | "span_attr:http.request.method", |
| 558 | "span_attr:http.response.status_code", |
| 559 | "span_attr:server.address" |
| 560 | ] |
| 561 | .map(|f| like(f, &escaped)) |
| 562 | .join(" OR ") |
| 563 | ), |
| 564 | Some(field) => { |
| 565 | let alias = match field { |
| 566 | "status" => "http.response.status_code", |
| 567 | "method" => "http.request.method", |
| 568 | "path" => "url.path", |
| 569 | "host" => "server.address", |
| 570 | "client" => "client.address", |
| 571 | _ => field, |
| 572 | }; |
| 573 | if term.op == ":" { |
| 574 | like(&format!("span_attr:{alias}"), &format!("^{escaped}")) |
| 575 | } else { |
| 576 | let value = term |
| 577 | .value |
| 578 | .parse::<f64>() |
| 579 | .ok() |
| 580 | .filter(|n| n.is_finite()) |
| 581 | .ok_or_else(|| { |
| 582 | Error::new(400, format!("Put a number after {field}:{}.", term.op)) |
| 583 | })?; |
| 584 | format!( |
| 585 | "{}:{}{value}", |
| 586 | quote(&format!("span_attr:{alias}")), |
| 587 | term.op |
| 588 | ) |
| 589 | } |
| 590 | } |
| 591 | }; |
| 592 | filters.push(if term.exclude { |
| 593 | format!("!({filter})") |
| 594 | } else { |
| 595 | filter |
| 596 | }); |
| 597 | } |
| 598 | let mut found = spans( |
| 599 | app, |
| 600 | format!( |
| 601 | "trace_id:in({} | sort by ({key}{}) | limit {limit} | fields trace_id)", |
| 602 | filters.join(" "), |
| 603 | if descending { " desc" } else { "" } |
| 604 | ), |
| 605 | ) |
| 606 | .await?; |
| 607 | for trace in &mut found { |
| 608 | trace["spans"].as_array_mut().unwrap().retain(&visible); |
| 609 | } |
| 610 | found.retain(|trace| !array(&trace["spans"]).is_empty()); |
| 611 | let value = |trace: &Value| { |
| 612 | array(&trace["spans"]) |
| 613 | .iter() |
| 614 | .find(|s| s["service"] == service && s["attributes"]["studio.kind"] != "edge") |
| 615 | .map(|s| number(&s[if key == "_time" { "start" } else { "duration" }])) |
| 616 | .unwrap_or(0.0) |
| 617 | }; |
| 618 | found.sort_by(|a, b| { |
| 619 | if descending { |
| 620 | value(b).total_cmp(&value(a)) |
| 621 | } else { |
| 622 | value(a).total_cmp(&value(b)) |
| 623 | } |
| 624 | }); |
| 625 | Ok(json!(found.into_iter().map(|t| { let spans = array(&t["spans"]); let mut services: Vec<_> = spans.iter().map(|s| s["service"].clone()).collect(); services.sort_by(|a,b| string(a).cmp(string(b))); services.dedup(); json!({"id":t["id"],"root":spans.first(),"spans":spans.len(),"services":services,"error":spans.iter().find(|s| !s["error"].is_null()).map(|s| s["error"].clone())}) }).collect::<Vec<_>>())) |
| 626 | } |
| 627 | |
| 628 | #[derive(Default)] |
| 629 | struct Recorder { |
| 630 | signature: Value, |
| 631 | cpu: HashMap<String, (f64, f64)>, |
| 632 | vms: HashMap<String, (f64, f64)>, |
| 633 | network: Option<(f64, f64, f64)>, |
| 634 | } |
| 635 | impl Recorder { |
| 636 | async fn collect(&mut self, app: Arc<App>) -> Result<Vec<String>> { |
| 637 | let jobs = core::scan(app.clone()).await?; |
| 638 | let mut signature = Vec::new(); |
| 639 | let mut owners = HashMap::new(); |
| 640 | for (id, found) in jobs.value.as_object().unwrap() { |
| 641 | for alloc in array(&found["allocs"]) { |
| 642 | let alloc_id = string(&alloc["ID"]).to_owned(); |
| 643 | owners.insert(alloc_id.clone(), id.clone()); |
| 644 | for (name, task) in alloc["TaskStates"] |
| 645 | .as_object() |
| 646 | .into_iter() |
| 647 | .flat_map(|v| v.iter()) |
| 648 | { |
| 649 | signature.push((alloc_id.clone(), name.clone(), task["StartedAt"].clone())); |
| 650 | } |
| 651 | } |
| 652 | } |
| 653 | signature.sort_by(|a, b| (&a.0, &a.1).cmp(&(&b.0, &b.1))); |
| 654 | let signature = json!(signature); |
| 655 | let counters = |
| 656 | host::call(json!({"operation":"host.usage","refresh":signature != self.signature})) |
| 657 | .await?; |
| 658 | self.signature = signature; |
| 659 | let at = number(&counters["at"]); |
| 660 | let timestamp = (now() * 1000.0) as i64; |
| 661 | let mut next = HashMap::new(); |
| 662 | let mut usage: HashMap<String, (f64, f64)> = HashMap::new(); |
| 663 | for container in array(&counters["containers"]) { |
| 664 | let name = string(&container["name"]); |
| 665 | let alloc = name |
| 666 | .get(name.len().saturating_sub(36)..) |
| 667 | .unwrap_or_default(); |
| 668 | let Some(service) = owners.get(alloc) else { |
| 669 | continue; |
| 670 | }; |
| 671 | let usec = number(&container["cpu"]); |
| 672 | let memory = number(&container["memory"]); |
| 673 | let id = string(&container["id"]); |
| 674 | let cpu = self |
| 675 | .cpu |
| 676 | .get(id) |
| 677 | .filter(|(_, t)| at > *t) |
| 678 | .map(|(old, t)| ((usec - old) / ((at - t) * 1e6)).max(0.0)) |
| 679 | .unwrap_or(0.0); |
| 680 | next.insert(id.to_owned(), (usec, at)); |
| 681 | let item = usage.entry(service.clone()).or_default(); |
| 682 | item.0 += cpu; |
| 683 | item.1 += memory; |
| 684 | } |
| 685 | self.cpu = next; |
| 686 | *app.usage.lock().unwrap() = usage |
| 687 | .iter() |
| 688 | .map(|(id, (cpu, memory))| (id.clone(), json!({"cpu":cpu,"memory":memory}))) |
| 689 | .collect(); |
| 690 | let bytes = app.live.borrow().clone(); |
| 691 | let live: Value = if bytes.is_empty() { |
| 692 | json!({}) |
| 693 | } else { |
| 694 | serde_json::from_slice(&bytes)? |
| 695 | }; |
| 696 | let mut lines = vec![ |
| 697 | format!( |
| 698 | "studio_host_cpu_percent {} {timestamp}", |
| 699 | number(&live["host"]["cpu"]) |
| 700 | ), |
| 701 | format!( |
| 702 | "studio_host_memory_bytes {} {timestamp}", |
| 703 | number(&live["host"]["memory"]) |
| 704 | ), |
| 705 | ]; |
| 706 | for (id, (cpu, memory)) in usage { |
| 707 | let id = serde_json::to_string(&id)?; |
| 708 | lines.push(format!( |
| 709 | "studio_service_cpu_cores{{service={id}}} {cpu} {timestamp}" |
| 710 | )); |
| 711 | lines.push(format!( |
| 712 | "studio_service_memory_bytes{{service={id}}} {memory} {timestamp}" |
| 713 | )); |
| 714 | } |
| 715 | let sample = host::sample(app.clone()).await?; |
| 716 | let sample = &sample.value; |
| 717 | let rx = number(&sample["network"]["rx"]); |
| 718 | let tx = number(&sample["network"]["tx"]); |
| 719 | let network_at = number(&sample["at"]); |
| 720 | if let Some((old_rx, old_tx, t)) = self.network.filter(|(_, _, t)| network_at > *t) { |
| 721 | lines.push(format!( |
| 722 | "studio_host_network_bytes_per_second{{direction=\"down\"}} {} {timestamp}", |
| 723 | ((rx - old_rx) / (network_at - t)).max(0.0) |
| 724 | )); |
| 725 | lines.push(format!( |
| 726 | "studio_host_network_bytes_per_second{{direction=\"up\"}} {} {timestamp}", |
| 727 | ((tx - old_tx) / (network_at - t)).max(0.0) |
| 728 | )); |
| 729 | } |
| 730 | self.network = Some((rx, tx, network_at)); |
| 731 | if let Some(value) = live["host"]["temperature"].as_f64() { |
| 732 | lines.push(format!( |
| 733 | "studio_host_temperature_celsius {value} {timestamp}" |
| 734 | )); |
| 735 | } |
| 736 | if let Some(value) = sample["gpu"].as_f64() { |
| 737 | lines.push(format!("studio_host_gpu_percent {value} {timestamp}")); |
| 738 | } |
| 739 | if let Ok(Ok(stats)) = tokio::time::timeout( |
| 740 | Duration::from_secs(8), |
| 741 | host::call(json!({"operation":"vm.stats"})), |
| 742 | ) |
| 743 | .await |
| 744 | { |
| 745 | let mut next = HashMap::new(); |
| 746 | let mut usage = HashMap::new(); |
| 747 | for (id, stat) in stats.as_object().into_iter().flat_map(|v| v.iter()) { |
| 748 | let cpu = self |
| 749 | .vms |
| 750 | .get(id) |
| 751 | .filter(|(_, t)| at > *t) |
| 752 | .map(|(old, t)| { |
| 753 | ((number(&stat["cpu"]) - old) / (at - t) / number(&stat["vcpus"]) * 100.0) |
| 754 | .max(0.0) |
| 755 | }) |
| 756 | .unwrap_or(0.0); |
| 757 | usage.insert(id.clone(), json!({"cpu":cpu,"memory":stat["memory"]})); |
| 758 | next.insert(id.clone(), (number(&stat["cpu"]), at)); |
| 759 | lines.push(format!( |
| 760 | "studio_vm_cpu_percent{{vm={}}} {cpu} {timestamp}", |
| 761 | serde_json::to_string(id)? |
| 762 | )); |
| 763 | lines.push(format!( |
| 764 | "studio_vm_memory_bytes{{vm={}}} {} {timestamp}", |
| 765 | serde_json::to_string(id)?, |
| 766 | number(&stat["memory"]) |
| 767 | )); |
| 768 | } |
| 769 | self.vms = next; |
| 770 | *app.vm_usage.lock().unwrap() = usage; |
| 771 | } |
| 772 | Ok(lines) |
| 773 | } |
| 774 | } |
| 775 | pub fn start(app: Arc<App>) { |
| 776 | tokio::spawn(async move { |
| 777 | let mut recorder = Recorder::default(); |
| 778 | loop { |
| 779 | let result = async { |
| 780 | let base = endpoint(app.clone(), "victoria-metrics").await?; |
| 781 | if env("STUDIO_METRICS", "") == "follow" { |
| 782 | let names: Vec<_> = METRICS |
| 783 | .iter() |
| 784 | .filter(|(k, _)| k.starts_with("service.") || k.starts_with("vm.")) |
| 785 | .map(|(_, v)| *v) |
| 786 | .collect(); |
| 787 | let url = format!( |
| 788 | "{base}/api/v1/query?{}", |
| 789 | params(&[("query", format!("{{__name__=~\"{}\"}}", names.join("|")))]) |
| 790 | ); |
| 791 | let result = read(app.clone(), url, None).await?; |
| 792 | let mut services = HashMap::<String, Value>::new(); |
| 793 | let mut vms = HashMap::<String, Value>::new(); |
| 794 | for row in array(&result.value["data"]["result"]) { |
| 795 | let metric = &row["metric"]; |
| 796 | let (map, id) = if metric["service"].is_string() { |
| 797 | (&mut services, string(&metric["service"])) |
| 798 | } else { |
| 799 | (&mut vms, string(&metric["vm"])) |
| 800 | }; |
| 801 | let value = map |
| 802 | .entry(id.into()) |
| 803 | .or_insert_with(|| json!({"cpu":0,"memory":0})); |
| 804 | value[if string(&metric["__name__"]).contains("_cpu_") { |
| 805 | "cpu" |
| 806 | } else { |
| 807 | "memory" |
| 808 | }] = json!(number(&row["value"][1])); |
| 809 | } |
| 810 | *app.usage.lock().unwrap() = services; |
| 811 | *app.vm_usage.lock().unwrap() = vms; |
| 812 | } else { |
| 813 | let lines = recorder.collect(app.clone()).await?; |
| 814 | app.request(Method::POST, &format!("{base}/api/v1/import/prometheus"))? |
| 815 | .body(lines.join("\n") + "\n") |
| 816 | .send() |
| 817 | .await? |
| 818 | .error_for_status()?; |
| 819 | let jobs = core::scan(app.clone()).await?; |
| 820 | let scrapes: Vec<_> = jobs |
| 821 | .value |
| 822 | .as_object() |
| 823 | .unwrap() |
| 824 | .iter() |
| 825 | .flat_map(|(id, job)| { |
| 826 | array(&job["job"]["TaskGroups"]) |
| 827 | .iter() |
| 828 | .flat_map(|group| array(&group["Services"])) |
| 829 | .flat_map(move |service| { |
| 830 | array(&service["Tags"]).iter().filter_map(move |tag| { |
| 831 | string(tag).strip_prefix("studio-metrics-path=").map( |
| 832 | |path| { |
| 833 | ( |
| 834 | id.clone(), |
| 835 | string(&service["Name"]).to_owned(), |
| 836 | path.to_owned(), |
| 837 | ) |
| 838 | }, |
| 839 | ) |
| 840 | }) |
| 841 | }) |
| 842 | }) |
| 843 | .collect(); |
| 844 | stream::iter(scrapes) |
| 845 | .for_each_concurrent(4, |(id, service, path)| { |
| 846 | let app = app.clone(); |
| 847 | let base = base.clone(); |
| 848 | async move { |
| 849 | let result = async { |
| 850 | let response = app |
| 851 | .request( |
| 852 | Method::GET, |
| 853 | &format!( |
| 854 | "{}{path}", |
| 855 | endpoint(app.clone(), &service).await? |
| 856 | ), |
| 857 | )? |
| 858 | .timeout(Duration::from_secs(5)) |
| 859 | .send() |
| 860 | .await? |
| 861 | .error_for_status()?; |
| 862 | app.request( |
| 863 | Method::POST, |
| 864 | &format!( |
| 865 | "{base}/api/v1/import/prometheus?{}", |
| 866 | params(&[ |
| 867 | ("extra_label", format!("service={id}")), |
| 868 | ("extra_label", format!("instance={service}")) |
| 869 | ]) |
| 870 | ), |
| 871 | )? |
| 872 | .body(response.text().await?) |
| 873 | .send() |
| 874 | .await? |
| 875 | .error_for_status()?; |
| 876 | Ok::<_, Error>(()) |
| 877 | } |
| 878 | .await; |
| 879 | if let Err(e) = result { |
| 880 | eprintln!("metrics {id}: {}", e.message); |
| 881 | } |
| 882 | } |
| 883 | }) |
| 884 | .await; |
| 885 | } |
| 886 | Ok::<_, Error>(()) |
| 887 | } |
| 888 | .await; |
| 889 | if let Err(e) = result |
| 890 | && e.status != 501 |
| 891 | { |
| 892 | eprintln!("metrics: {}", e.message); |
| 893 | } |
| 894 | tokio::time::sleep(Duration::from_secs(15)).await; |
| 895 | } |
| 896 | }); |
| 897 | } |
| 898 | |
| 899 | #[cfg(test)] |
| 900 | mod tests { |
| 901 | use super::*; |
| 902 | #[test] |
| 903 | fn search_fields_comparisons_exclusions_and_unknown_fields() { |
| 904 | let parsed = terms( |
| 905 | r#"status:5xx -path:/health duration:>1s "GET /x" user:42 -noise"#, |
| 906 | true, |
| 907 | ) |
| 908 | .unwrap(); |
| 909 | assert_eq!( |
| 910 | parsed |
| 911 | .iter() |
| 912 | .map(|t| ( |
| 913 | t.field.as_deref(), |
| 914 | t.op.as_str(), |
| 915 | t.value.as_str(), |
| 916 | t.exclude |
| 917 | )) |
| 918 | .collect::<Vec<_>>(), |
| 919 | vec![ |
| 920 | (Some("status"), ":", "5", false), |
| 921 | (Some("path"), ":", "/health", true), |
| 922 | (Some("duration"), ">", "1s", false), |
| 923 | (None, ":", "GET /x", false), |
| 924 | (None, ":", "user:42", false), |
| 925 | (None, ":", "noise", true) |
| 926 | ] |
| 927 | ); |
| 928 | assert_eq!(terms(r#" "" - "#, false).unwrap().len(), 1); |
| 929 | } |
| 930 | #[test] |
| 931 | fn non_finite_points_are_null() { |
| 932 | assert_eq!( |
| 933 | series( |
| 934 | &json!({"data":{"result":[{"metric":{"service":"a"},"values":[[1,"2"],[2,"NaN"],[3,"+Inf"]]}]}}), |
| 935 | false |
| 936 | ), |
| 937 | json!([{"name":"a","t":[1,2,3],"v":[2.0,null,null]}]) |
| 938 | ); |
| 939 | } |
| 940 | #[test] |
| 941 | fn span_error_status_is_preserved() { |
| 942 | let s = span( |
| 943 | &json!({"span_id":"a","parent_span_id":"a","start_time_unix_nano":"1000000000","duration":"20","status_code":"2","status_message":"broken"}), |
| 944 | ); |
| 945 | assert_eq!(s["error"], "broken"); |
| 946 | } |
| 947 | } |