1use crate::*;
2use futures::{StreamExt, stream};
3
4pub 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
40pub 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
74const 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
87fn 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
132pub 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}
150pub 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}
223pub 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)]
249struct Term {
250 field: Option<String>,
251 op: String,
252 value: String,
253 exclude: bool,
254}
255fn 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}
298pub 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
359fn 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}
415async 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}
461pub 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}
473pub 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}
496pub 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)]
629struct Recorder {
630 signature: Value,
631 cpu: HashMap<String, (f64, f64)>,
632 vms: HashMap<String, (f64, f64)>,
633 network: Option<(f64, f64, f64)>,
634}
635impl 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}
775pub 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)]
900mod 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}