diff --git a/src/servers/src/http/prom_store.rs b/src/servers/src/http/prom_store.rs index cb7b6307c4b..60c000403fd 100644 --- a/src/servers/src/http/prom_store.rs +++ b/src/servers/src/http/prom_store.rs @@ -175,12 +175,20 @@ async fn remote_write_v1( processor.set_pipeline(pipeline_handler, query_ctx.clone(), pipeline_def); } - let mut req = decode_remote_write_request(is_zstd, body, prom_validation_mode, &mut processor)?; + let mut decoded = + decode_remote_write_request(is_zstd, body, prom_validation_mode, &mut processor)?; + // Parsing borrows the decode buffer, but row building copies out of it: tag + // values through `decode_string`, column names through `to_owned`, and the + // borrowing `col_indexes` dies inside `as_insert_requests`. Nothing below + // references the buffer, so it need not span the write. let req = if processor.use_pipeline { + drop(decoded); processor.exec_pipeline().await? } else { - req.as_insert_requests() + let req = decoded.as_insert_requests(); + drop(decoded); + req }; let batches = into_prom_write_batches(req, query_ctx); diff --git a/src/servers/src/http/result/json_result.rs b/src/servers/src/http/result/json_result.rs index 537ea66cd02..74a75f26309 100644 --- a/src/servers/src/http/result/json_result.rs +++ b/src/servers/src/http/result/json_result.rs @@ -90,17 +90,20 @@ impl IntoResponse for JsonResponse { .to_string(), Some(GreptimeQueryOutput::Records(records)) => { - let schema = records.schema(); + // Borrow the schema field directly so the rows can be moved out. + let schema = &records.schema; let data: Vec> = records .rows - .iter() - .map(|row| { + .into_iter() + .map(|mut row| { + // Slicing keeps the out-of-bounds panic for short rows; + // a plain zip would silently truncate them. schema .column_schemas .iter() - .enumerate() - .map(|(i, col)| (col.name.clone(), row[i].clone())) + .zip(row[..schema.column_schemas.len()].iter_mut()) + .map(|(col, value)| (col.name.clone(), value.take())) .collect::>() }) .collect(); @@ -133,3 +136,47 @@ impl IntoResponse for JsonResponse { .into_response() } } + +#[cfg(test)] +mod tests { + use axum::body::to_bytes; + + use super::*; + + #[tokio::test] + async fn test_records_response_preserves_values_and_duplicate_columns() { + let response: JsonResponse = serde_json::from_value(json!({ + "output": [{"records": { + "schema": {"column_schemas": [ + {"name": "duplicate", "data_type": "String"}, + {"name": "nested", "data_type": "Json"}, + {"name": "duplicate", "data_type": "String"}, + {"name": "escaped\"column", "data_type": "String"} + ]}, + "rows": [ + ["discarded", {"array": [null, true, "中文"]}, "last", "line\n\\\""], + ["discarded", [1, {"key": "value"}], null, ""] + ] + }}], + "execution_time_ms": 7 + })) + .unwrap(); + + let response = response.into_response(); + assert_eq!(response.status(), axum::http::StatusCode::OK); + assert_eq!(response.headers()[header::CONTENT_TYPE], "application/json"); + assert_eq!(response.headers()[&GREPTIME_DB_HEADER_FORMAT], "json"); + assert_eq!(response.headers()[&GREPTIME_DB_HEADER_EXECUTION_TIME], "7"); + let body = to_bytes(response.into_body(), usize::MAX).await.unwrap(); + assert_eq!( + serde_json::from_slice::(&body).unwrap(), + json!({ + "data": [ + {"duplicate": "last", "nested": {"array": [null, true, "中文"]}, "escaped\"column": "line\n\\\""}, + {"duplicate": null, "nested": [1, {"key": "value"}], "escaped\"column": ""} + ], + "execution_time_ms": 7 + }) + ); + } +} diff --git a/src/servers/src/prom_remote_write/mod.rs b/src/servers/src/prom_remote_write/mod.rs index 73f968b6084..c9aabfb2745 100644 --- a/src/servers/src/prom_remote_write/mod.rs +++ b/src/servers/src/prom_remote_write/mod.rs @@ -77,6 +77,8 @@ pub fn decode_remote_write_request( // fallback to the other compression method try_decompress(!is_zstd, &body[..])? }; + // Decompression copied the payload out, so the compressed body is no longer needed. + drop(body); let mut request = PROM_WRITE_REQUEST_POOL.pull(PromWriteRequest::default); diff --git a/src/servers/src/prom_remote_write/v2.rs b/src/servers/src/prom_remote_write/v2.rs index b388071106f..2f1f4d410c4 100644 --- a/src/servers/src/prom_remote_write/v2.rs +++ b/src/servers/src/prom_remote_write/v2.rs @@ -143,6 +143,8 @@ pub(crate) fn decode_remote_write_v2( } else { try_decompress(!is_zstd, &body[..])? }; + // Decompression copied the payload out, so the compressed body is no longer needed. + drop(body); let request = BorrowedRequest::decode(&buf).context(error::DecodePromRemoteRequestSnafu)?; drop(decode_timer);