something
This commit is contained in:
249
komp-app/src/analytics.rs
Normal file
249
komp-app/src/analytics.rs
Normal file
@@ -0,0 +1,249 @@
|
||||
use anyhow::{Result, bail, ensure};
|
||||
use common::proto::komp_ac::analytics::{
|
||||
AnalyticsResultBatch, AnalyticsResultColumn, AnalyticsValue, analytics_value::Kind,
|
||||
};
|
||||
use serde::Serialize;
|
||||
|
||||
const MAX_ROWS: usize = 10_000;
|
||||
const MAX_RESULT_BYTES: usize = 32 * 1024 * 1024;
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize)]
|
||||
#[serde(tag = "kind", content = "value", rename_all = "snake_case")]
|
||||
pub enum Cell {
|
||||
Null,
|
||||
Text(String),
|
||||
Integer(String),
|
||||
Unsigned(String),
|
||||
Float(f64),
|
||||
Boolean(bool),
|
||||
Bytes(String),
|
||||
}
|
||||
|
||||
impl TryFrom<AnalyticsValue> for Cell {
|
||||
type Error = anyhow::Error;
|
||||
|
||||
fn try_from(value: AnalyticsValue) -> Result<Self> {
|
||||
Ok(match value.kind {
|
||||
Some(Kind::NullValue(0)) => Self::Null,
|
||||
Some(Kind::StringValue(value)) => Self::Text(value),
|
||||
Some(Kind::Int64Value(value)) => Self::Integer(value.to_string()),
|
||||
Some(Kind::Uint64Value(value)) => Self::Unsigned(value.to_string()),
|
||||
Some(Kind::DoubleValue(value)) if value.is_finite() => Self::Float(value),
|
||||
Some(Kind::BoolValue(value)) => Self::Boolean(value),
|
||||
Some(Kind::BytesValue(value)) => {
|
||||
Self::Bytes(value.iter().map(|byte| format!("{byte:02x}")).collect())
|
||||
}
|
||||
_ => bail!("Analytics returned an invalid or non-finite cell"),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl Cell {
|
||||
fn size(&self) -> usize {
|
||||
match self {
|
||||
Self::Text(value)
|
||||
| Self::Integer(value)
|
||||
| Self::Unsigned(value)
|
||||
| Self::Bytes(value) => value.len().saturating_add(32),
|
||||
_ => 32,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Default, Serialize)]
|
||||
pub struct QueryResult {
|
||||
pub columns: Vec<AnalyticsResultColumn>,
|
||||
pub rows: Vec<Vec<Cell>>,
|
||||
pub row_count: u64,
|
||||
pub elapsed_ms: u64,
|
||||
pub truncated: bool,
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct Collector {
|
||||
result: QueryResult,
|
||||
bytes: usize,
|
||||
finished: bool,
|
||||
}
|
||||
|
||||
impl Collector {
|
||||
fn push(&mut self, batch: AnalyticsResultBatch) -> Result<()> {
|
||||
ensure!(!self.finished, "Analytics sent data after completion");
|
||||
if !batch.columns.is_empty() {
|
||||
if self.result.columns.is_empty() {
|
||||
self.result.columns = batch.columns;
|
||||
} else {
|
||||
ensure!(
|
||||
self.result.columns == batch.columns,
|
||||
"Analytics changed its result schema"
|
||||
);
|
||||
}
|
||||
}
|
||||
ensure!(
|
||||
self.result.rows.len().saturating_add(batch.rows.len()) <= MAX_ROWS,
|
||||
"Analytics result exceeds the client row limit"
|
||||
);
|
||||
for row in batch.rows {
|
||||
ensure!(
|
||||
row.values.len() == self.result.columns.len(),
|
||||
"Analytics row does not match its schema"
|
||||
);
|
||||
let row: Vec<_> = row
|
||||
.values
|
||||
.into_iter()
|
||||
.map(Cell::try_from)
|
||||
.collect::<Result<_>>()?;
|
||||
self.bytes = self
|
||||
.bytes
|
||||
.saturating_add(row.iter().map(Cell::size).sum::<usize>());
|
||||
ensure!(
|
||||
self.bytes <= MAX_RESULT_BYTES,
|
||||
"Analytics result exceeds 32 MiB; narrow the query or filters"
|
||||
);
|
||||
self.result.rows.push(row);
|
||||
}
|
||||
if batch.is_final {
|
||||
ensure!(
|
||||
batch.row_count == self.result.rows.len() as u64,
|
||||
"Analytics result is incomplete"
|
||||
);
|
||||
self.finished = true;
|
||||
self.result.row_count = batch.row_count;
|
||||
self.result.elapsed_ms = batch.elapsed_ms;
|
||||
self.result.truncated = batch.truncated;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn finish(self) -> Result<QueryResult> {
|
||||
ensure!(self.finished, "Analytics stream ended before completion");
|
||||
Ok(self.result)
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn collect_result(
|
||||
mut stream: tonic::Streaming<AnalyticsResultBatch>,
|
||||
) -> Result<QueryResult> {
|
||||
let mut collector = Collector::default();
|
||||
while let Some(batch) = stream.message().await? {
|
||||
collector.push(batch)?;
|
||||
}
|
||||
collector.finish()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use common::proto::komp_ac::analytics::AnalyticsResultRow;
|
||||
|
||||
#[test]
|
||||
fn ipc_preserves_exact_values_and_distinguishes_null_from_empty_text() {
|
||||
assert_eq!(
|
||||
Cell::try_from(AnalyticsValue {
|
||||
kind: Some(Kind::Int64Value(i64::MAX))
|
||||
})
|
||||
.unwrap(),
|
||||
Cell::Integer("9223372036854775807".into())
|
||||
);
|
||||
assert_eq!(
|
||||
Cell::try_from(AnalyticsValue {
|
||||
kind: Some(Kind::Uint64Value(u64::MAX))
|
||||
})
|
||||
.unwrap(),
|
||||
Cell::Unsigned("18446744073709551615".into())
|
||||
);
|
||||
assert_eq!(
|
||||
Cell::try_from(AnalyticsValue {
|
||||
kind: Some(Kind::StringValue("999999999999999999.99".into()))
|
||||
})
|
||||
.unwrap(),
|
||||
Cell::Text("999999999999999999.99".into())
|
||||
);
|
||||
assert_ne!(
|
||||
Cell::try_from(AnalyticsValue {
|
||||
kind: Some(Kind::NullValue(0))
|
||||
})
|
||||
.unwrap(),
|
||||
Cell::try_from(AnalyticsValue {
|
||||
kind: Some(Kind::StringValue(String::new()))
|
||||
})
|
||||
.unwrap()
|
||||
);
|
||||
assert!(Cell::try_from(AnalyticsValue { kind: None }).is_err());
|
||||
assert!(
|
||||
Cell::try_from(AnalyticsValue {
|
||||
kind: Some(Kind::DoubleValue(f64::NAN))
|
||||
})
|
||||
.is_err()
|
||||
);
|
||||
}
|
||||
|
||||
fn data() -> AnalyticsResultBatch {
|
||||
AnalyticsResultBatch {
|
||||
columns: vec![AnalyticsResultColumn {
|
||||
name: "amount".into(),
|
||||
data_type: "Decimal128(20, 2)".into(),
|
||||
}],
|
||||
rows: vec![AnalyticsResultRow {
|
||||
values: vec![AnalyticsValue {
|
||||
kind: Some(Kind::StringValue("123.45".into())),
|
||||
}],
|
||||
}],
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn interrupted_stream_is_never_returned_as_a_successful_partial_report() {
|
||||
let mut collector = Collector::default();
|
||||
collector.push(data()).unwrap();
|
||||
assert!(collector.finish().is_err());
|
||||
let mut collector = Collector::default();
|
||||
collector.push(data()).unwrap();
|
||||
assert!(
|
||||
collector
|
||||
.push(AnalyticsResultBatch {
|
||||
is_final: true,
|
||||
row_count: 2,
|
||||
..Default::default()
|
||||
})
|
||||
.is_err()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn completion_preserves_schema_and_truncation() {
|
||||
let mut collector = Collector::default();
|
||||
collector.push(data()).unwrap();
|
||||
collector
|
||||
.push(AnalyticsResultBatch {
|
||||
is_final: true,
|
||||
row_count: 1,
|
||||
truncated: true,
|
||||
elapsed_ms: 7,
|
||||
..Default::default()
|
||||
})
|
||||
.unwrap();
|
||||
let result = collector.finish().unwrap();
|
||||
assert_eq!(result.columns[0].data_type, "Decimal128(20, 2)");
|
||||
assert_eq!(result.row_count, 1);
|
||||
assert_eq!(result.elapsed_ms, 7);
|
||||
assert!(result.truncated);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn mismatched_schema_and_post_completion_batches_are_rejected() {
|
||||
let mut collector = Collector::default();
|
||||
let mut invalid = data();
|
||||
invalid.rows[0].values.clear();
|
||||
assert!(collector.push(invalid).is_err());
|
||||
let mut collector = Collector::default();
|
||||
collector
|
||||
.push(AnalyticsResultBatch {
|
||||
is_final: true,
|
||||
..Default::default()
|
||||
})
|
||||
.unwrap();
|
||||
assert!(collector.push(data()).is_err());
|
||||
}
|
||||
}
|
||||
@@ -1,4 +1,5 @@
|
||||
use crate::search::SearchGrpc;
|
||||
mod reporting;
|
||||
use anyhow::{Context, Result, anyhow};
|
||||
use crate::transport::{
|
||||
DEFAULT_GRPC_ENDPOINT, authenticated_request as request_with_auth_token, connect_channel,
|
||||
|
||||
134
komp-app/src/grpc/reporting.rs
Normal file
134
komp-app/src/grpc/reporting.rs
Normal file
@@ -0,0 +1,134 @@
|
||||
use super::GrpcClient;
|
||||
use anyhow::{Context, Result};
|
||||
use common::proto::komp_ac::analytics::{reporting_service_client::ReportingServiceClient, *};
|
||||
|
||||
impl GrpcClient {
|
||||
pub async fn list_report_assets(
|
||||
&mut self,
|
||||
request: ListReportAssetsRequest,
|
||||
) -> Result<ListReportAssetsResponse> {
|
||||
let request = self.authenticated_request(request)?;
|
||||
let response = ReportingServiceClient::new(self.channel.clone())
|
||||
.list_assets(request)
|
||||
.await
|
||||
.context("gRPC ReportingService list_assets call failed")?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
|
||||
pub async fn get_report_asset(
|
||||
&mut self,
|
||||
request: GetReportAssetRequest,
|
||||
) -> Result<ReportAsset> {
|
||||
let request = self.authenticated_request(request)?;
|
||||
let response = ReportingServiceClient::new(self.channel.clone())
|
||||
.get_asset(request)
|
||||
.await
|
||||
.context("gRPC ReportingService get_asset call failed")?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
|
||||
pub async fn save_report_draft(
|
||||
&mut self,
|
||||
request: SaveReportDraftRequest,
|
||||
) -> Result<ReportAsset> {
|
||||
let request = self.authenticated_request(request)?;
|
||||
let response = ReportingServiceClient::new(self.channel.clone())
|
||||
.save_draft(request)
|
||||
.await
|
||||
.context("gRPC ReportingService save_draft call failed")?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
|
||||
pub async fn publish_report(&mut self, request: PublishReportRequest) -> Result<ReportAsset> {
|
||||
let request = self.authenticated_request(request)?;
|
||||
let response = ReportingServiceClient::new(self.channel.clone())
|
||||
.publish(request)
|
||||
.await
|
||||
.context("gRPC ReportingService publish call failed")?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
|
||||
pub async fn list_report_versions(
|
||||
&mut self,
|
||||
request: ReportAssetRef,
|
||||
) -> Result<ListReportVersionsResponse> {
|
||||
let request = self.authenticated_request(request)?;
|
||||
let response = ReportingServiceClient::new(self.channel.clone())
|
||||
.list_versions(request)
|
||||
.await
|
||||
.context("gRPC ReportingService list_versions call failed")?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
|
||||
pub async fn restore_report_draft(
|
||||
&mut self,
|
||||
request: RestoreReportDraftRequest,
|
||||
) -> Result<ReportAsset> {
|
||||
let request = self.authenticated_request(request)?;
|
||||
let response = ReportingServiceClient::new(self.channel.clone())
|
||||
.restore_draft(request)
|
||||
.await
|
||||
.context("gRPC ReportingService restore_draft call failed")?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
|
||||
pub async fn set_report_archived(
|
||||
&mut self,
|
||||
request: SetReportArchivedRequest,
|
||||
) -> Result<ReportAsset> {
|
||||
let request = self.authenticated_request(request)?;
|
||||
let response = ReportingServiceClient::new(self.channel.clone())
|
||||
.set_archived(request)
|
||||
.await
|
||||
.context("gRPC ReportingService set_archived call failed")?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
|
||||
pub async fn get_report_personal_views(
|
||||
&mut self,
|
||||
request: ReportAssetRef,
|
||||
) -> Result<GetReportPersonalViewsResponse> {
|
||||
let request = self.authenticated_request(request)?;
|
||||
let response = ReportingServiceClient::new(self.channel.clone())
|
||||
.get_personal_views(request)
|
||||
.await
|
||||
.context("gRPC ReportingService get_personal_views call failed")?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
|
||||
pub async fn save_report_personal_view(
|
||||
&mut self,
|
||||
request: SaveReportPersonalViewRequest,
|
||||
) -> Result<ReportPersonalView> {
|
||||
let request = self.authenticated_request(request)?;
|
||||
let response = ReportingServiceClient::new(self.channel.clone())
|
||||
.save_personal_view(request)
|
||||
.await
|
||||
.context("gRPC ReportingService save_personal_view call failed")?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
|
||||
pub async fn delete_report_personal_view(
|
||||
&mut self,
|
||||
request: DeleteReportPersonalViewRequest,
|
||||
) -> Result<ReportMutationResult> {
|
||||
let request = self.authenticated_request(request)?;
|
||||
let response = ReportingServiceClient::new(self.channel.clone())
|
||||
.delete_personal_view(request)
|
||||
.await
|
||||
.context("gRPC ReportingService delete_personal_view call failed")?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
|
||||
pub async fn execute_report_dataset(
|
||||
&mut self,
|
||||
request: ExecuteReportDatasetRequest,
|
||||
) -> Result<tonic::Streaming<AnalyticsResultBatch>> {
|
||||
let request = self.authenticated_request(request)?;
|
||||
let response = ReportingServiceClient::new(self.channel.clone())
|
||||
.execute_dataset(request)
|
||||
.await
|
||||
.context("gRPC ReportingService execute_dataset call failed")?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
}
|
||||
@@ -5,6 +5,7 @@
|
||||
//! owns client behavior that must not drift between them.
|
||||
|
||||
pub mod auth;
|
||||
pub mod analytics;
|
||||
pub mod csv;
|
||||
pub mod grpc;
|
||||
pub mod import_export;
|
||||
|
||||
Reference in New Issue
Block a user