echarts to display the datafusion SQL
This commit is contained in:
195
graphs/src/main.rs
Normal file
195
graphs/src/main.rs
Normal file
@@ -0,0 +1,195 @@
|
||||
use std::{env, net::SocketAddr};
|
||||
|
||||
use axum::{
|
||||
Json, Router,
|
||||
extract::State,
|
||||
http::StatusCode,
|
||||
response::Html,
|
||||
routing::{get, post},
|
||||
};
|
||||
mod analytics {
|
||||
include!(concat!(
|
||||
env!("CARGO_MANIFEST_DIR"),
|
||||
"/../common/src/proto/komp_ac.analytics.rs"
|
||||
));
|
||||
}
|
||||
|
||||
use analytics::{
|
||||
AnalyticsResultColumn, AnalyticsValue, ExecuteAnalyticsQueryRequest,
|
||||
analytics_service_client::AnalyticsServiceClient, analytics_value,
|
||||
};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use serde_json::{Value, json};
|
||||
use tonic::{Request, metadata::MetadataValue, transport::Channel};
|
||||
|
||||
const INDEX_HTML: &str = include_str!("../static/index.html");
|
||||
|
||||
#[derive(Clone)]
|
||||
struct AppState {
|
||||
analytics: AnalyticsServiceClient<Channel>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct QueryInput {
|
||||
profile_name: String,
|
||||
sql: String,
|
||||
#[serde(default)]
|
||||
token: String,
|
||||
#[serde(default = "default_max_rows")]
|
||||
max_rows: u32,
|
||||
}
|
||||
|
||||
fn default_max_rows() -> u32 {
|
||||
1_000
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
struct QueryOutput {
|
||||
columns: Vec<ColumnOutput>,
|
||||
rows: Vec<Vec<Value>>,
|
||||
row_count: u64,
|
||||
elapsed_ms: u64,
|
||||
truncated: bool,
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
struct ColumnOutput {
|
||||
name: String,
|
||||
data_type: String,
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
struct ErrorOutput {
|
||||
error: String,
|
||||
}
|
||||
|
||||
type ApiResult<T> = Result<Json<T>, (StatusCode, Json<ErrorOutput>)>;
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
let grpc_endpoint =
|
||||
env::var("ANALYTICS_GRPC_ENDPOINT").unwrap_or_else(|_| "http://[::1]:50051".into());
|
||||
let listen_address = env::var("LISTEN_ADDRESS")
|
||||
.unwrap_or_else(|_| "127.0.0.1:3000".into())
|
||||
.parse::<SocketAddr>()?;
|
||||
|
||||
let channel = Channel::from_shared(grpc_endpoint.clone())?.connect_lazy();
|
||||
let state = AppState {
|
||||
analytics: AnalyticsServiceClient::new(channel),
|
||||
};
|
||||
|
||||
let app = Router::new()
|
||||
.route("/", get(index))
|
||||
.route("/api/query", post(run_query))
|
||||
.with_state(state);
|
||||
|
||||
let listener = tokio::net::TcpListener::bind(listen_address).await?;
|
||||
println!("Analytics graphs: http://{listen_address}");
|
||||
println!("Analytics gRPC endpoint: {grpc_endpoint}");
|
||||
axum::serve(listener, app).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn index() -> Html<&'static str> {
|
||||
Html(INDEX_HTML)
|
||||
}
|
||||
|
||||
async fn run_query(State(state): State<AppState>, Json(input): Json<QueryInput>) -> ApiResult<QueryOutput> {
|
||||
if input.profile_name.trim().is_empty() || input.sql.trim().is_empty() {
|
||||
return Err(api_error(
|
||||
StatusCode::BAD_REQUEST,
|
||||
"Profile name and SQL are required",
|
||||
));
|
||||
}
|
||||
|
||||
let mut request = Request::new(ExecuteAnalyticsQueryRequest {
|
||||
profile_name: input.profile_name,
|
||||
sql: input.sql,
|
||||
max_rows: input.max_rows,
|
||||
});
|
||||
|
||||
if !input.token.trim().is_empty() {
|
||||
let value = MetadataValue::try_from(format!("Bearer {}", input.token.trim()))
|
||||
.map_err(|_| api_error(StatusCode::BAD_REQUEST, "The bearer token is invalid"))?;
|
||||
request.metadata_mut().insert("authorization", value);
|
||||
}
|
||||
|
||||
let mut client = state.analytics;
|
||||
let mut stream = client
|
||||
.execute_analytics_query(request)
|
||||
.await
|
||||
.map_err(grpc_error)?
|
||||
.into_inner();
|
||||
|
||||
let mut columns = Vec::new();
|
||||
let mut rows = Vec::new();
|
||||
let mut row_count = 0;
|
||||
let mut elapsed_ms = 0;
|
||||
let mut truncated = false;
|
||||
|
||||
while let Some(batch) = stream.message().await.map_err(grpc_error)? {
|
||||
if !batch.columns.is_empty() {
|
||||
columns = batch.columns.into_iter().map(column_output).collect();
|
||||
}
|
||||
rows.extend(batch.rows.into_iter().map(|row| {
|
||||
row.values.into_iter().map(value_output).collect::<Vec<_>>()
|
||||
}));
|
||||
|
||||
if batch.is_final {
|
||||
row_count = batch.row_count;
|
||||
elapsed_ms = batch.elapsed_ms;
|
||||
truncated = batch.truncated;
|
||||
}
|
||||
}
|
||||
|
||||
Ok(Json(QueryOutput {
|
||||
columns,
|
||||
rows,
|
||||
row_count,
|
||||
elapsed_ms,
|
||||
truncated,
|
||||
}))
|
||||
}
|
||||
|
||||
fn column_output(column: AnalyticsResultColumn) -> ColumnOutput {
|
||||
ColumnOutput {
|
||||
name: column.name,
|
||||
data_type: column.data_type,
|
||||
}
|
||||
}
|
||||
|
||||
fn value_output(value: AnalyticsValue) -> Value {
|
||||
match value.kind {
|
||||
None | Some(analytics_value::Kind::NullValue(_)) => Value::Null,
|
||||
Some(analytics_value::Kind::BoolValue(value)) => json!(value),
|
||||
Some(analytics_value::Kind::Int64Value(value)) => json!(value),
|
||||
Some(analytics_value::Kind::Uint64Value(value)) => json!(value),
|
||||
Some(analytics_value::Kind::DoubleValue(value)) => json!(value),
|
||||
Some(analytics_value::Kind::StringValue(value)) => json!(value),
|
||||
Some(analytics_value::Kind::BytesValue(value)) => {
|
||||
Value::String(value.iter().map(|byte| format!("{byte:02x}")).collect())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn grpc_error(error: tonic::Status) -> (StatusCode, Json<ErrorOutput>) {
|
||||
let status = match error.code() {
|
||||
tonic::Code::InvalidArgument => StatusCode::BAD_REQUEST,
|
||||
tonic::Code::Unauthenticated => StatusCode::UNAUTHORIZED,
|
||||
tonic::Code::PermissionDenied => StatusCode::FORBIDDEN,
|
||||
tonic::Code::NotFound => StatusCode::NOT_FOUND,
|
||||
tonic::Code::ResourceExhausted => StatusCode::TOO_MANY_REQUESTS,
|
||||
tonic::Code::Unavailable => StatusCode::SERVICE_UNAVAILABLE,
|
||||
_ => StatusCode::BAD_GATEWAY,
|
||||
};
|
||||
api_error(status, error.message())
|
||||
}
|
||||
|
||||
fn api_error(status: StatusCode, message: impl Into<String>) -> (StatusCode, Json<ErrorOutput>) {
|
||||
(
|
||||
status,
|
||||
Json(ErrorOutput {
|
||||
error: message.into(),
|
||||
}),
|
||||
)
|
||||
}
|
||||
Reference in New Issue
Block a user