use super::*;
use nodedb_types::Value;
use nodedb_types::protocol::opcodes::ResponseStatus;
use nodedb_types::protocol::{MAX_FRAME_SIZE, NativeResponse};
#[test]
fn chunk_large_response_splits_rows() {
let columns = vec!["id".to_string(), "data".to_string()];
let rows: Vec<Vec<Value>> = (0..100)
.map(|i| {
vec![
Value::Integer(i),
Value::String(format!("row-data-{i}-padding-{}", "x".repeat(150))),
]
})
.collect();
let response = NativeResponse {
seq: 1,
status: ResponseStatus::Ok,
columns: Some(columns),
rows: Some(rows),
rows_affected: None,
watermark_lsn: 42,
error: None,
auth: None,
warnings: Vec::new(),
};
let frames = chunk_large_response(response, codec::FrameFormat::MessagePack).unwrap();
assert!(!frames.is_empty());
for (i, frame) in frames.iter().enumerate() {
let resp: NativeResponse = zerompk::from_msgpack(frame).unwrap();
assert!(resp.rows.is_some());
if i < frames.len() - 1 {
assert_eq!(resp.status, ResponseStatus::Partial);
} else {
assert_eq!(resp.status, ResponseStatus::Ok);
}
}
}
#[test]
fn chunk_large_response_no_rows_passthrough() {
let response = NativeResponse {
seq: 1,
status: ResponseStatus::Ok,
columns: None,
rows: None,
rows_affected: Some(5),
watermark_lsn: 42,
error: None,
auth: None,
warnings: Vec::new(),
};
let frames = chunk_large_response(response, codec::FrameFormat::MessagePack).unwrap();
assert_eq!(
frames.len(),
1,
"no-rows response should pass through as-is"
);
}
#[test]
fn chunk_large_response_preserves_all_rows() {
let columns = vec!["id".to_string(), "value".to_string()];
let row_count = 100_000;
let rows: Vec<Vec<Value>> = (0..row_count)
.map(|i| {
vec![
Value::Integer(i),
Value::String(format!("v{i}-{}", "p".repeat(150))),
]
})
.collect();
let response = NativeResponse {
seq: 42,
status: ResponseStatus::Ok,
columns: Some(columns.clone()),
rows: Some(rows),
rows_affected: None,
watermark_lsn: 99,
error: None,
auth: None,
warnings: Vec::new(),
};
let frames = chunk_large_response(response, codec::FrameFormat::MessagePack).unwrap();
assert!(frames.len() > 1, "should produce multiple frames");
let mut total_rows: Vec<Vec<Value>> = Vec::new();
for frame in &frames {
let resp: NativeResponse = zerompk::from_msgpack(frame).unwrap();
if let Some(rows) = resp.rows {
total_rows.extend(rows);
}
}
assert_eq!(total_rows.len(), row_count as usize);
let first: NativeResponse = zerompk::from_msgpack(&frames[0]).unwrap();
assert_eq!(first.columns, Some(columns));
assert_eq!(first.status, ResponseStatus::Partial);
let last: NativeResponse = zerompk::from_msgpack(frames.last().unwrap()).unwrap();
assert_eq!(last.status, ResponseStatus::Ok);
for frame in &frames {
assert!(
frame.len() <= MAX_FRAME_SIZE as usize,
"frame size {} exceeds MAX_FRAME_SIZE {}",
frame.len(),
MAX_FRAME_SIZE
);
}
}