micromegas/servers/
query_audit.rs1use datafusion::physical_plan::ExecutionPlan;
11
12pub struct ScanMetrics {
14 pub output_rows: Option<u64>,
15 pub bytes_scanned: u64,
16}
17
18pub fn aggregate_scan_metrics(plan: &dyn ExecutionPlan) -> ScanMetrics {
22 fn sum_bytes(plan: &dyn ExecutionPlan) -> u64 {
23 let mut total = plan
24 .metrics()
25 .and_then(|m| m.sum_by_name("bytes_scanned"))
26 .map(|v| v.as_usize() as u64)
27 .unwrap_or(0);
28 for child in plan.children() {
29 total += sum_bytes(child.as_ref());
30 }
31 total
32 }
33 ScanMetrics {
34 output_rows: plan
35 .metrics()
36 .and_then(|m| m.output_rows())
37 .map(|r| r as u64),
38 bytes_scanned: sum_bytes(plan),
39 }
40}
41
42#[derive(serde::Serialize)]
46pub struct QueryAuditRecord {
47 pub client: String,
48 pub user: String,
49 pub email: String,
50 #[serde(skip_serializing_if = "Option::is_none")]
51 pub name: Option<String>,
52 pub service_account: bool,
53 #[serde(skip_serializing_if = "Option::is_none")]
54 pub service_account_name: Option<String>,
55 pub sql: String,
56 #[serde(skip_serializing_if = "Option::is_none")]
57 pub range_begin: Option<String>, #[serde(skip_serializing_if = "Option::is_none")]
59 pub range_end: Option<String>,
60 #[serde(skip_serializing_if = "Option::is_none")]
61 pub limit: Option<u64>,
62 pub context_init_ms: f64,
63 pub planning_ms: f64,
64 pub execution_ms: f64, pub setup_ms: f64, pub total_ms: f64, pub status: &'static str, #[serde(skip_serializing_if = "Option::is_none")]
69 pub error: Option<String>,
70 #[serde(skip_serializing_if = "Option::is_none")]
71 pub output_rows: Option<u64>,
72 pub bytes_scanned: u64,
73}