use serde_json::{Map, Value as JsonValue};
use sonic_rs;
use crate::control::security::identity::AuthenticatedIdentity;
use crate::control::server::response_shape::types::ShapedRows;
use crate::control::state::SharedState;
use crate::event::cdc::consume::{ConsumeError, ConsumeParams, consume_stream};
use super::super::result::{DdlError, DdlResult};
fn err(sqlstate: &str, message: impl Into<String>) -> DdlError {
DdlError {
sqlstate: sqlstate.to_string(),
message: message.into(),
}
}
pub async fn select_from_stream(
state: &SharedState,
identity: &AuthenticatedIdentity,
parts: &[&str],
) -> Result<Vec<DdlResult>, DdlError> {
let tenant_id = identity.tenant_id.as_u64();
if parts.len() < 8
|| !parts[3].eq_ignore_ascii_case("STREAM")
|| !parts[5].eq_ignore_ascii_case("CONSUMER")
|| !parts[6].eq_ignore_ascii_case("GROUP")
{
return Err(err(
"42601",
"expected SELECT * FROM STREAM <stream> CONSUMER GROUP <group> [PARTITION <p>] [LIMIT <n>]",
));
}
let stream_name = parts[4].to_lowercase();
let group_name = parts[7].to_lowercase();
let mut partition: Option<u32> = None;
let mut limit: usize = 100;
let mut i = 8;
while i < parts.len() {
if parts[i].eq_ignore_ascii_case("PARTITION") && i + 1 < parts.len() {
partition = Some(
parts[i + 1]
.parse()
.map_err(|_| err("42601", format!("invalid partition: '{}'", parts[i + 1])))?,
);
i += 2;
} else if parts[i].eq_ignore_ascii_case("LIMIT") && i + 1 < parts.len() {
limit = parts[i + 1]
.parse()
.map_err(|_| err("42601", format!("invalid limit: '{}'", parts[i + 1])))?;
i += 2;
} else {
i += 1;
}
}
let consume_params = ConsumeParams {
tenant_id,
stream_name: &stream_name,
group_name: &group_name,
partition,
limit,
};
let result = match consume_stream(state, &consume_params) {
Ok(r) => r,
Err(ConsumeError::RemotePartition { leader_node, .. }) => {
match crate::event::cdc::consume::consume_remote(state, &consume_params, leader_node)
.await
{
Ok(r) => r,
Err(e) => return Err(err("58000", e.to_string())),
}
}
Err(ConsumeError::BufferEmpty(_)) => {
let columns = result_columns();
let column_types = ShapedRows::text_types(columns.len());
return Ok(vec![DdlResult::Rows(ShapedRows {
columns,
column_types,
rows: Vec::new(),
notice: None,
})]);
}
Err(e) => {
return Err(err("42704", e.to_string()));
}
};
let columns = result_columns();
let mut rows = Vec::with_capacity(result.events.len());
for event in &result.events {
let mut row = Map::new();
row.insert(
"sequence".to_string(),
JsonValue::String(event.sequence.to_string()),
);
row.insert(
"partition".to_string(),
JsonValue::String(event.partition.to_string()),
);
row.insert(
"collection".to_string(),
JsonValue::String(event.collection.clone()),
);
row.insert(
"event_type".to_string(),
JsonValue::String(event.op.clone()),
);
row.insert(
"row_id".to_string(),
JsonValue::String(event.row_id.clone()),
);
row.insert("lsn".to_string(), JsonValue::String(event.lsn.to_string()));
row.insert(
"event_time".to_string(),
JsonValue::String(event.event_time.to_string()),
);
let new_val = event
.new_value
.as_ref()
.map(|v| sonic_rs::to_string(v).unwrap_or_default())
.unwrap_or_default();
row.insert("new_value".to_string(), JsonValue::String(new_val));
let old_val = event
.old_value
.as_ref()
.map(|v| sonic_rs::to_string(v).unwrap_or_default())
.unwrap_or_default();
row.insert("old_value".to_string(), JsonValue::String(old_val));
rows.push(row);
}
let column_types = ShapedRows::text_types(columns.len());
Ok(vec![DdlResult::Rows(ShapedRows {
columns,
column_types,
rows,
notice: None,
})])
}
fn result_columns() -> Vec<String> {
vec![
"sequence".to_string(),
"partition".to_string(),
"collection".to_string(),
"event_type".to_string(),
"row_id".to_string(),
"lsn".to_string(),
"event_time".to_string(),
"new_value".to_string(),
"old_value".to_string(),
]
}