1use std::io::Write;
6
7use serde::Serialize;
8
9use crate::cli::RoutingReportFormat;
10use crate::error::{CliError, Result};
11use crate::models::{Column, Row};
12use crate::output::formatter::Formatter;
13pub const DEFAULT_BUFFER_LIMIT: usize = 10 * 1024 * 1024;
15
16#[derive(Debug, Clone, Serialize)]
20pub struct DistributedReadRoutingReport {
21 pub schema_version: u32,
22 pub requested_mode: String,
23 pub effective_mode: Option<String>,
24 pub decision: String,
25 pub ranges: Vec<String>,
26 pub freshness: Option<String>,
27 pub retry_count: u32,
28 pub failover_count: u32,
29 pub outcome: String,
30 #[serde(skip_serializing_if = "Option::is_none")]
31 pub reason: Option<String>,
32}
33
34impl DistributedReadRoutingReport {
35 pub fn new(
36 requested_mode: impl Into<String>,
37 effective_mode: Option<String>,
38 decision: impl Into<String>,
39 outcome: impl Into<String>,
40 reason: Option<String>,
41 ) -> Self {
42 Self {
43 schema_version: 1,
44 requested_mode: requested_mode.into(),
45 effective_mode,
46 decision: decision.into(),
47 ranges: Vec::new(),
48 freshness: None,
49 retry_count: 0,
50 failover_count: 0,
51 outcome: outcome.into(),
52 reason,
53 }
54 }
55}
56
57pub fn write_distributed_read_routing_report<W: Write>(
61 writer: &mut W,
62 format: RoutingReportFormat,
63 report: &DistributedReadRoutingReport,
64) -> Result<()> {
65 match format {
66 RoutingReportFormat::Json => {
67 serde_json::to_writer(&mut *writer, report)?;
68 writeln!(writer)?;
69 }
70 RoutingReportFormat::Human => {
71 writeln!(writer, "distributed read routing report")?;
72 writeln!(writer, " requested_mode: {}", report.requested_mode)?;
73 if let Some(mode) = &report.effective_mode {
74 writeln!(writer, " effective_mode: {mode}")?;
75 }
76 writeln!(writer, " decision: {}", report.decision)?;
77 writeln!(writer, " ranges: {}", report.ranges.join(","))?;
78 if let Some(freshness) = &report.freshness {
79 writeln!(writer, " freshness: {freshness}")?;
80 }
81 writeln!(writer, " retry_count: {}", report.retry_count)?;
82 writeln!(writer, " failover_count: {}", report.failover_count)?;
83 writeln!(writer, " outcome: {}", report.outcome)?;
84 if let Some(reason) = &report.reason {
85 writeln!(writer, " reason: {reason}")?;
86 }
87 }
88 }
89 writer.flush()?;
90 Ok(())
91}
92
93#[derive(Debug, Clone, Copy, PartialEq, Eq)]
95pub enum WriteStatus {
96 Continue,
98 LimitReached,
100}
101
102pub struct StreamingWriter<W> {
113 writer: W,
115 formatter: Box<dyn Formatter>,
117 columns: Vec<Column>,
119 limit: Option<usize>,
121 buffer_limit: usize,
123 buffer: Vec<Row>,
125 buffer_bytes: usize,
127 written_count: usize,
129 output_started: bool,
131 quiet: bool,
133}
134
135impl<W: Write> StreamingWriter<W> {
136 pub fn new(
145 writer: W,
146 formatter: Box<dyn Formatter>,
147 columns: Vec<Column>,
148 limit: Option<usize>,
149 ) -> Self {
150 Self {
151 writer,
152 formatter,
153 columns,
154 limit,
155 buffer_limit: DEFAULT_BUFFER_LIMIT,
156 buffer: Vec::new(),
157 buffer_bytes: 0,
158 written_count: 0,
159 output_started: false,
160 quiet: false,
161 }
162 }
163
164 #[allow(dead_code)]
166 pub fn with_buffer_limit(mut self, buffer_limit: usize) -> Self {
167 self.buffer_limit = buffer_limit;
168 self
169 }
170
171 pub fn with_quiet(mut self, quiet: bool) -> Self {
173 self.quiet = quiet;
174 self
175 }
176
177 pub fn is_quiet(&self) -> bool {
181 self.quiet
182 }
183
184 pub fn prepare(&mut self, row_count_hint: Option<usize>) -> Result<()> {
193 let _ = row_count_hint;
194
195 if self.formatter.supports_streaming() {
197 self.formatter
198 .write_header(&mut self.writer, &self.columns)?;
199 self.output_started = true;
200 }
201
202 Ok(())
203 }
204
205 pub fn write_row(&mut self, row: Row) -> Result<WriteStatus> {
219 if let Some(limit) = self.limit {
221 if self.written_count >= limit {
222 return Ok(WriteStatus::LimitReached);
223 }
224 }
225
226 if self.formatter.supports_streaming() {
227 self.formatter.write_row(&mut self.writer, &row)?;
229 self.written_count += 1;
230 } else {
231 let row_bytes = estimate_row_bytes(&row);
233 self.buffer_bytes = self.buffer_bytes.saturating_add(row_bytes);
234 if self.buffer_bytes > self.buffer_limit {
235 return Err(CliError::InvalidArgument(
236 "Buffer limit exceeded (~10MB). \
237 Use --output json|csv|tsv or --limit to reduce results."
238 .into(),
239 ));
240 }
241 self.buffer.push(row);
242 self.written_count += 1;
243 }
244
245 Ok(WriteStatus::Continue)
246 }
247
248 pub fn finish(&mut self) -> Result<()> {
250 if !self.output_started {
252 self.formatter
253 .write_header(&mut self.writer, &self.columns)?;
254 self.output_started = true;
255
256 for row in self.buffer.drain(..) {
258 self.formatter.write_row(&mut self.writer, &row)?;
259 }
260 }
261
262 self.formatter.write_footer(&mut self.writer)?;
264
265 Ok(())
266 }
267
268 #[allow(dead_code)]
270 pub fn written_count(&self) -> usize {
271 self.written_count
272 }
273
274 #[allow(dead_code)]
276 pub fn output_started(&self) -> bool {
277 self.output_started
278 }
279
280 #[allow(dead_code)]
282 pub fn buffered_bytes(&self) -> usize {
283 self.buffer_bytes
284 }
285}
286
287fn estimate_row_bytes(row: &Row) -> usize {
288 row.columns
289 .iter()
290 .map(estimate_value_bytes)
291 .sum::<usize>()
292 .saturating_add(row.columns.len() * 8)
293}
294
295fn estimate_value_bytes(value: &crate::models::Value) -> usize {
296 match value {
297 crate::models::Value::Null => 4,
298 crate::models::Value::Bool(_) => 1,
299 crate::models::Value::Int(_) => 8,
300 crate::models::Value::Float(_) => 8,
301 crate::models::Value::Text(text) => text.len(),
302 crate::models::Value::Bytes(bytes) => bytes.len(),
303 crate::models::Value::Vector(values) => values.len() * 4,
304 }
305}
306
307#[cfg(test)]
308mod tests {
309 use super::*;
310 use crate::error::CliError;
311 use crate::models::{DataType, Value};
312 use crate::output::csv::CsvFormatter;
313 use crate::output::json::JsonFormatter;
314 use crate::output::jsonl::JsonlFormatter;
315 use crate::output::table::TableFormatter;
316
317 fn test_columns() -> Vec<Column> {
318 vec![
319 Column::new("id", DataType::Int),
320 Column::new("name", DataType::Text),
321 ]
322 }
323
324 fn test_row(id: i64, name: &str) -> Row {
325 Row::new(vec![Value::Int(id), Value::Text(name.to_string())])
326 }
327
328 #[test]
329 fn test_streaming_format_immediate_output() {
330 let mut output = Vec::new();
331 let formatter = Box::new(JsonlFormatter::new());
332 let columns = test_columns();
333
334 let mut writer = StreamingWriter::new(&mut output, formatter, columns, None);
335
336 writer.prepare(None).unwrap();
337 assert!(writer.output_started());
338
339 let status = writer.write_row(test_row(1, "Alice")).unwrap();
340 assert_eq!(status, WriteStatus::Continue);
341 assert_eq!(writer.written_count(), 1);
342
343 writer.finish().unwrap();
344
345 let result = String::from_utf8(output).unwrap();
346 assert!(result.contains("\"id\":1"));
347 assert!(result.contains("\"name\":\"Alice\""));
348 }
349
350 #[test]
351 fn test_non_streaming_format_buffered_output() {
352 let mut output = Vec::new();
353 let formatter = Box::new(TableFormatter::new());
354 let columns = test_columns();
355
356 let mut writer = StreamingWriter::new(&mut output, formatter, columns, None);
357
358 writer.prepare(None).unwrap();
359 assert!(!writer.output_started()); let status = writer.write_row(test_row(1, "Alice")).unwrap();
362 assert_eq!(status, WriteStatus::Continue);
363 assert!(!writer.output_started()); writer.finish().unwrap();
366 assert!(writer.output_started()); let result = String::from_utf8(output).unwrap();
369 assert!(result.contains("id"));
370 assert!(result.contains("Alice"));
371 }
372
373 #[test]
374 fn test_limit_enforcement() {
375 let mut output = Vec::new();
376 let formatter = Box::new(CsvFormatter::new());
377 let columns = test_columns();
378
379 let mut writer = StreamingWriter::new(&mut output, formatter, columns, Some(2));
380
381 writer.prepare(None).unwrap();
382
383 assert_eq!(
384 writer.write_row(test_row(1, "Alice")).unwrap(),
385 WriteStatus::Continue
386 );
387 assert_eq!(
388 writer.write_row(test_row(2, "Bob")).unwrap(),
389 WriteStatus::Continue
390 );
391 assert_eq!(
392 writer.write_row(test_row(3, "Charlie")).unwrap(),
393 WriteStatus::LimitReached
394 );
395
396 assert_eq!(writer.written_count(), 2);
397
398 writer.finish().unwrap();
399 }
400
401 #[test]
402 fn test_buffer_overflow_errors() {
403 let mut output = Vec::new();
404 let formatter = Box::new(TableFormatter::new());
405 let columns = test_columns();
406
407 let mut writer =
408 StreamingWriter::new(&mut output, formatter, columns, None).with_buffer_limit(40); writer.prepare(None).unwrap();
411 assert!(!writer.output_started());
412
413 assert_eq!(
415 writer.write_row(test_row(1, "Alice")).unwrap(),
416 WriteStatus::Continue
417 );
418 let err = writer.write_row(test_row(2, "Bob")).unwrap_err();
419 assert!(matches!(err, CliError::InvalidArgument(_)));
420 }
421
422 #[test]
423 fn test_empty_output() {
424 let mut output = Vec::new();
425 let formatter = Box::new(JsonFormatter::new());
426 let columns = test_columns();
427
428 let mut writer = StreamingWriter::new(&mut output, formatter, columns, None);
429
430 writer.prepare(None).unwrap();
431 writer.finish().unwrap();
432
433 let result = String::from_utf8(output).unwrap();
434 assert!(result.contains('['));
436 assert!(result.contains(']'));
437 }
438
439 #[test]
440 fn test_csv_streaming() {
441 let mut output = Vec::new();
442 let formatter = Box::new(CsvFormatter::new());
443 let columns = test_columns();
444
445 let mut writer = StreamingWriter::new(&mut output, formatter, columns, None);
446
447 writer.prepare(None).unwrap();
448 assert!(writer.output_started()); writer.write_row(test_row(1, "Alice")).unwrap();
451 writer.write_row(test_row(2, "Bob")).unwrap();
452 writer.finish().unwrap();
453
454 let result = String::from_utf8(output).unwrap();
455 assert_eq!(result, "id,name\n1,Alice\n2,Bob\n");
456 }
457
458 #[test]
459 fn test_written_count() {
460 let mut output = Vec::new();
461 let formatter = Box::new(CsvFormatter::new());
462 let columns = test_columns();
463
464 let mut writer = StreamingWriter::new(&mut output, formatter, columns, None);
465
466 writer.prepare(None).unwrap();
467
468 assert_eq!(writer.written_count(), 0);
469 writer.write_row(test_row(1, "Alice")).unwrap();
470 assert_eq!(writer.written_count(), 1);
471 writer.write_row(test_row(2, "Bob")).unwrap();
472 assert_eq!(writer.written_count(), 2);
473
474 writer.finish().unwrap();
475 }
476
477 #[test]
478 fn test_table_with_small_data() {
479 let mut output = Vec::new();
480 let formatter = Box::new(TableFormatter::new());
481 let columns = test_columns();
482
483 let mut writer = StreamingWriter::new(&mut output, formatter, columns, None);
484
485 writer.prepare(None).unwrap();
486 assert!(!writer.output_started()); writer.write_row(test_row(1, "Alice")).unwrap();
489 writer.write_row(test_row(2, "Bob")).unwrap();
490 writer.finish().unwrap();
491
492 let result = String::from_utf8(output).unwrap();
493 assert!(result.contains("id"));
495 assert!(result.contains("name"));
496 assert!(result.contains("Alice"));
497 assert!(result.contains("Bob"));
498 }
499
500 #[test]
501 fn test_streaming_large_row_count_does_not_buffer() {
502 let formatter = Box::new(JsonlFormatter::new());
503 let columns = test_columns();
504
505 let mut writer = StreamingWriter::new(std::io::sink(), formatter, columns, None);
506
507 writer.prepare(None).unwrap();
508 for i in 0..12_000 {
509 writer.write_row(test_row(i, "row")).unwrap();
510 }
511
512 assert_eq!(writer.written_count(), 12_000);
513 assert!(writer.buffer.is_empty());
514
515 writer.finish().unwrap();
516 }
517
518 #[test]
519 fn routing_report_is_serialized_to_its_own_writer_without_row_output() {
520 let report = DistributedReadRoutingReport::new(
521 "strong",
522 Some("strong".into()),
523 "cluster",
524 "retryable_failure",
525 Some("read point expired".into()),
526 );
527 let mut stderr = Vec::new();
528 write_distributed_read_routing_report(&mut stderr, RoutingReportFormat::Json, &report)
529 .unwrap();
530 let value: serde_json::Value = serde_json::from_slice(&stderr).unwrap();
531 assert_eq!(value["requested_mode"], "strong");
532 assert_eq!(value["outcome"], "retryable_failure");
533 assert_eq!(value["reason"], "read point expired");
534 }
535}