Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
94 changes: 78 additions & 16 deletions src/mapper/row_compact.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,24 +17,31 @@ pub struct QueryResult {
pub columns: Vec<ColumnSchema>,
pub rows: Vec<Vec<Value>>,
pub statistics: QueryStatistics,
#[serde(skip_serializing_if = "Option::is_none")]
pub query_id: Option<String>,
}

/// Parses the output of ClickHouse `FORMAT JSONCompactEachRowWithNamesAndTypes`.
/// Line 1: JSON array of column names `["id", "amount"]`
/// Line 2: JSON array of ClickHouse data types `["UInt64", "Decimal(18, 4)"]`
/// Lines 3+: JSON arrays of row values `[18446744073709551615, 123.4500]`
/// Automatically normalizes row values according to `ColumnSchema::mapped_type` (e.g. converting 64-bit numbers and Decimals into JSON strings to prevent 53-bit float overflow in JS/Flutter).
///
/// `limit`, when set, stops parsing (and allocating) further data rows once
/// that many have been read, instead of parsing the entire result set and
/// discarding the excess — this bounds peak memory for large result sets
/// (see issue #47).
pub fn parse_compact_output(
output_text: &str,
elapsed_ms: u64,
limit: Option<usize>,
) -> Result<QueryResult, DriverError> {
let lines: Vec<&str> = output_text
let mut lines = output_text
.lines()
.map(|l| l.trim())
.filter(|l| !l.is_empty())
.collect();
.filter(|l| !l.is_empty());

if lines.is_empty() {
let Some(names_line) = lines.next() else {
return Ok(QueryResult {
columns: vec![],
rows: vec![],
Expand All @@ -43,23 +50,24 @@ pub fn parse_compact_output(
bytes_read: output_text.len(),
elapsed_ms,
},
query_id: None,
});
}
};

if lines.len() < 2 {
let Some(types_line) = lines.next() else {
return Err(DriverError::Client(
"Malformed JSONCompactEachRowWithNamesAndTypes output: missing names or types row"
.to_string(),
));
}
};

let names: Vec<String> = serde_json::from_str(lines[0]).map_err(|e| {
let names: Vec<String> = serde_json::from_str(names_line).map_err(|e| {
DriverError::Client(format!(
"Failed to parse column names from ClickHouse output: {}",
e
))
})?;
let types: Vec<String> = serde_json::from_str(lines[1]).map_err(|e| {
let types: Vec<String> = serde_json::from_str(types_line).map_err(|e| {
DriverError::Client(format!(
"Failed to parse column types from ClickHouse output: {}",
e
Expand All @@ -79,8 +87,12 @@ pub fn parse_compact_output(
columns.push(ColumnSchema::new(name, ch_type));
}

let mut rows = Vec::with_capacity(lines.len().saturating_sub(2));
for line in &lines[2..] {
let mut rows = Vec::with_capacity(limit.unwrap_or(16).min(1024));
for line in lines {
if limit.is_some_and(|limit| rows.len() >= limit) {
break;
}

let mut raw_row: Vec<Value> = serde_json::from_str(line).map_err(|e| {
DriverError::Client(format!(
"Failed to parse data row JSON array from ClickHouse: {}",
Expand Down Expand Up @@ -124,6 +136,7 @@ pub fn parse_compact_output(
bytes_read: output_text.len(),
elapsed_ms,
},
query_id: None,
})
}

Expand All @@ -139,7 +152,7 @@ mod tests {
[18446744073709551615, "Alice", 1234567.8901, true]
[102, null, 0.0000, false]"#;

let res = parse_compact_output(raw_output, 15).unwrap();
let res = parse_compact_output(raw_output, 15, None).unwrap();
assert_eq!(res.columns.len(), 4);
assert_eq!(res.columns[0].mapped_type, "string");
assert_eq!(res.columns[1].mapped_type, "string");
Expand Down Expand Up @@ -170,32 +183,81 @@ mod tests {
let version_output = r#"["version()"]
["String"]
["24.3.1.2452"]"#;
let res = parse_compact_output(version_output, 0).unwrap();
let res = parse_compact_output(version_output, 0, None).unwrap();
assert_eq!(res.rows.len(), 1);
assert_eq!(res.rows[0][0], json!("24.3.1.2452"));

let uptime_output = r#"["uptime()"]
["UInt32"]
[123456]"#;
let res = parse_compact_output(uptime_output, 0).unwrap();
let res = parse_compact_output(uptime_output, 0, None).unwrap();
assert_eq!(res.rows[0][0].as_u64(), Some(123456));
}

#[test]
fn test_parse_compact_output_empty() {
let res = parse_compact_output("", 5).unwrap();
let res = parse_compact_output("", 5, None).unwrap();
assert!(res.columns.is_empty());
assert!(res.rows.is_empty());
assert_eq!(res.statistics.rows_read, 0);
}

#[test]
fn test_parse_compact_output_enforces_limit() {
// Regression for issue #47: `limit` must stop row parsing early instead
// of parsing the whole result set and discarding the excess, since a
// multi-GB result would otherwise be fully materialized in memory first.
let raw_output = r#"["id"]
["UInt64"]
[1]
[2]
[3]
[4]
[5]"#;

let unlimited = parse_compact_output(raw_output, 0, None).unwrap();
assert_eq!(unlimited.rows.len(), 5);
assert_eq!(unlimited.statistics.rows_read, 5);

let limited = parse_compact_output(raw_output, 0, Some(2)).unwrap();
assert_eq!(limited.rows.len(), 2);
// UInt64 is normalized to a JSON string to protect 53-bit JS precision.
assert_eq!(limited.rows[0][0], json!("1"));
assert_eq!(limited.rows[1][0], json!("2"));
assert_eq!(limited.statistics.rows_read, 2);

// A limit larger than the actual row count is a no-op.
let generous_limit = parse_compact_output(raw_output, 0, Some(100)).unwrap();
assert_eq!(generous_limit.rows.len(), 5);

// A zero limit returns no rows at all, without erroring.
let zero_limit = parse_compact_output(raw_output, 0, Some(0)).unwrap();
assert!(zero_limit.rows.is_empty());
}

#[test]
fn test_parse_compact_output_sets_query_id_field() {
// Regression for issue #47: query_id lives directly on QueryResult so
// callers don't need a second serde_json::to_value pass just to splice
// a queryId key into the already-serialized response.
let raw_output = r#"["id"]
["UInt64"]
[1]"#;
let mut res = parse_compact_output(raw_output, 0, None).unwrap();
assert_eq!(res.query_id, None);
res.query_id = Some("querya-job-1-2-3".to_string());

let serialized = serde_json::to_value(&res).unwrap();
assert_eq!(serialized["queryId"], json!("querya-job-1-2-3"));
}

#[test]
fn test_parse_compact_output_complex_types() {
let raw_output = r#"["arr", "tup", "dt", "big_arr"]
["Array(Int32)", "Tuple(Int32, String)", "DateTime64(3)", "Array(UInt64)"]
[[10, 20, 30], [100, "foo"], "2026-07-11 12:34:56.789", [18446744073709551615, 42]]"#;

let res = parse_compact_output(raw_output, 8).unwrap();
let res = parse_compact_output(raw_output, 8, None).unwrap();
assert_eq!(res.columns.len(), 4);
assert_eq!(res.columns[0].mapped_type, "array");
assert_eq!(res.columns[1].mapped_type, "json");
Expand Down
45 changes: 34 additions & 11 deletions src/rpc/handlers/query.rs
Original file line number Diff line number Diff line change
Expand Up @@ -341,14 +341,13 @@ pub async fn handle_query(params: Option<Value>) -> Result<Value, DriverError> {
["UInt64", "String", "Nullable(UInt64)"]
[18446744073709551615, "page_view", 42]
[100, "click", null]"#;
let mut parsed_val = serde_json::to_value(parse_compact_output(
let mut result = parse_compact_output(
mock_output,
start_time.elapsed().as_millis() as u64,
)?)?;
if let Some(obj) = parsed_val.as_object_mut() {
obj.insert("queryId".to_string(), json!(actual_query_id));
}
return Ok(parsed_val);
query_params.limit,
)?;
result.query_id = Some(actual_query_id);
return Ok(serde_json::to_value(result)?);
} else {
return Ok(build_non_tabular_result(
&upper_sql,
Expand All @@ -370,11 +369,9 @@ pub async fn handle_query(params: Option<Value>) -> Result<Value, DriverError> {
let elapsed = start_time.elapsed().as_millis() as u64;

if is_tabular_query {
let mut parsed_val = serde_json::to_value(parse_compact_output(&text, elapsed)?)?;
if let Some(obj) = parsed_val.as_object_mut() {
obj.insert("queryId".to_string(), json!(actual_query_id));
}
Ok(parsed_val)
let mut result = parse_compact_output(&text, elapsed, query_params.limit)?;
result.query_id = Some(actual_query_id);
Ok(serde_json::to_value(result)?)
} else {
Ok(build_non_tabular_result(
&upper_sql,
Expand Down Expand Up @@ -770,6 +767,32 @@ mod tests {
ConnectionPool::global().remove(111);
}

#[tokio::test]
async fn test_handle_query_enforces_limit() {
// Regression for issue #47: a `limit` in the request must actually
// truncate the parsed rows instead of being silently ignored.
let _guard = crate::utils::test_lock::GLOBAL_TEST_LOCK.lock().await;
let client = ClickHouseClient::from_params(ConnectParams {
connection_id: 113,
connection_string: Some("mock://localhost:8123/default".to_string()),
..Default::default()
})
.unwrap();
ConnectionPool::global().insert(client);

let query_params = json!({
"connectionId": 113,
"sql": "SELECT id, event_name, user_id FROM events",
"limit": 1
});

let res = handle_query(Some(query_params)).await.unwrap();
assert_eq!(res["rows"].as_array().unwrap().len(), 1);
assert_eq!(res["statistics"]["rowsRead"], 1);

ConnectionPool::global().remove(113);
}

#[tokio::test]
async fn test_handle_query_tabular_detection_edge_cases() {
// Regression for issue #58: leading comments, CTE `WITH` queries and a
Expand Down
12 changes: 6 additions & 6 deletions src/rpc/handlers/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -277,7 +277,7 @@ pub async fn handle_context_actions(params: Option<Value>) -> Result<Value, Driv
crate::utils::sql_escape::escape_sql_string_literal(col_name)
);
let text = run_introspection_query(p.connection_id, &sql).await?;
crate::mapper::row_compact::parse_compact_output(&text, 0)
crate::mapper::row_compact::parse_compact_output(&text, 0, None)
.ok()
.and_then(|result| result.rows.into_iter().next())
.and_then(|row| row.into_iter().next())
Expand Down Expand Up @@ -386,15 +386,15 @@ pub async fn handle_get_server_stats(params: Option<Value>) -> Result<Value, Dri
});

let mut version_str = "ClickHouse".to_string();
if let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&version_text, 0)
if let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&version_text, 0, None)
&& let Some(row) = parsed.rows.first()
&& let Some(v) = row.first().and_then(|x| x.as_str())
{
version_str = format!("ClickHouse {}", v);
}

let mut uptime_sec = 0;
if let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&uptime_text, 0)
if let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&uptime_text, 0, None)
&& let Some(row) = parsed.rows.first()
&& let Some(v) = row.first().and_then(|x| x.as_u64())
{
Expand All @@ -409,7 +409,7 @@ pub async fn handle_get_server_stats(params: Option<Value>) -> Result<Value, Dri
.await
.unwrap_or_default();
let mut db_sizes = serde_json::Map::new();
if let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&db_sizes_text, 0) {
if let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&db_sizes_text, 0, None) {
for row in parsed.rows {
if let (Some(db), Some(size)) = (
row.first().and_then(|x| x.as_str()),
Expand Down Expand Up @@ -482,7 +482,7 @@ pub async fn handle_get_object_metadata(params: Option<Value>) -> Result<Value,
);
let mut ddl_str = String::new();
if let Ok(text) = client.post_sql(&ddl_sql, |_| {}).await
&& let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&text, 0)
&& let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&text, 0, None)
&& let Some(row) = parsed.rows.first()
&& let Some(v) = row.first().and_then(|x| x.as_str())
{
Expand All @@ -496,7 +496,7 @@ pub async fn handle_get_object_metadata(params: Option<Value>) -> Result<Value,
);
let mut columns = Vec::new();
if let Ok(text) = client.post_sql(&cols_sql, |_| {}).await
&& let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&text, 0)
&& let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&text, 0, None)
{
for row in parsed.rows {
let name = row.first().and_then(|x| x.as_str()).unwrap_or("unknown");
Expand Down
10 changes: 5 additions & 5 deletions src/sdui/tree.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ pub fn build_root_databases_nodes(
let mut nodes = Vec::new();

if let Some(output) = compact_output {
let parsed = parse_compact_output(output, 0)?;
let parsed = parse_compact_output(output, 0, None)?;
for row in parsed.rows {
let name = row.first().and_then(|v| v.as_str()).unwrap_or("unknown");
let engine = row.get(1).and_then(|v| v.as_str()).unwrap_or("");
Expand Down Expand Up @@ -146,7 +146,7 @@ pub fn parse_tables_nodes(
compact_output: &str,
filter_view: bool,
) -> Result<Vec<SduiTreeNode>, DriverError> {
let parsed = parse_compact_output(compact_output, 0)?;
let parsed = parse_compact_output(compact_output, 0, None)?;
let mut nodes = Vec::new();

for row in parsed.rows {
Expand Down Expand Up @@ -195,7 +195,7 @@ pub fn parse_dictionaries_nodes(
db_name: &str,
compact_output: &str,
) -> Result<Vec<SduiTreeNode>, DriverError> {
let parsed = parse_compact_output(compact_output, 0)?;
let parsed = parse_compact_output(compact_output, 0, None)?;
let mut nodes = Vec::new();

for row in parsed.rows {
Expand Down Expand Up @@ -233,7 +233,7 @@ pub fn parse_columns_nodes(
table_name: &str,
compact_output: &str,
) -> Result<Vec<SduiTreeNode>, DriverError> {
let parsed = parse_compact_output(compact_output, 0)?;
let parsed = parse_compact_output(compact_output, 0, None)?;
let mut nodes = Vec::new();

for row in parsed.rows {
Expand Down Expand Up @@ -269,7 +269,7 @@ pub fn parse_partitions_nodes(
table_name: &str,
compact_output: &str,
) -> Result<Vec<SduiTreeNode>, DriverError> {
let parsed = parse_compact_output(compact_output, 0)?;
let parsed = parse_compact_output(compact_output, 0, None)?;
let mut nodes = Vec::new();

for row in parsed.rows {
Expand Down
Loading