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 pub fn remaining_limit(&self) -> Option<usize> {
276 self.limit
277 .map(|limit| limit.saturating_sub(self.written_count))
278 }
279
280 #[allow(dead_code)]
282 pub fn output_started(&self) -> bool {
283 self.output_started
284 }
285
286 #[allow(dead_code)]
288 pub fn buffered_bytes(&self) -> usize {
289 self.buffer_bytes
290 }
291}
292
293fn estimate_row_bytes(row: &Row) -> usize {
294 row.columns
295 .iter()
296 .map(estimate_value_bytes)
297 .sum::<usize>()
298 .saturating_add(row.columns.len() * 8)
299}
300
301fn estimate_value_bytes(value: &crate::models::Value) -> usize {
302 match value {
303 crate::models::Value::Null => 4,
304 crate::models::Value::Bool(_) => 1,
305 crate::models::Value::Int(_) => 8,
306 crate::models::Value::Float(_) => 8,
307 crate::models::Value::Text(text) => text.len(),
308 crate::models::Value::Bytes(bytes) => bytes.len(),
309 crate::models::Value::Vector(values) => values.len() * 4,
310 }
311}
312
313#[cfg(test)]
314mod tests {
315 use super::*;
316 use crate::error::CliError;
317 use crate::models::{DataType, Value};
318 use crate::output::csv::CsvFormatter;
319 use crate::output::json::JsonFormatter;
320 use crate::output::jsonl::JsonlFormatter;
321 use crate::output::table::TableFormatter;
322
323 fn test_columns() -> Vec<Column> {
324 vec![
325 Column::new("id", DataType::Int),
326 Column::new("name", DataType::Text),
327 ]
328 }
329
330 fn test_row(id: i64, name: &str) -> Row {
331 Row::new(vec![Value::Int(id), Value::Text(name.to_string())])
332 }
333
334 #[test]
335 fn test_streaming_format_immediate_output() {
336 let mut output = Vec::new();
337 let formatter = Box::new(JsonlFormatter::new());
338 let columns = test_columns();
339
340 let mut writer = StreamingWriter::new(&mut output, formatter, columns, None);
341
342 writer.prepare(None).unwrap();
343 assert!(writer.output_started());
344
345 let status = writer.write_row(test_row(1, "Alice")).unwrap();
346 assert_eq!(status, WriteStatus::Continue);
347 assert_eq!(writer.written_count(), 1);
348
349 writer.finish().unwrap();
350
351 let result = String::from_utf8(output).unwrap();
352 assert!(result.contains("\"id\":1"));
353 assert!(result.contains("\"name\":\"Alice\""));
354 }
355
356 #[test]
357 fn test_non_streaming_format_buffered_output() {
358 let mut output = Vec::new();
359 let formatter = Box::new(TableFormatter::new());
360 let columns = test_columns();
361
362 let mut writer = StreamingWriter::new(&mut output, formatter, columns, None);
363
364 writer.prepare(None).unwrap();
365 assert!(!writer.output_started()); let status = writer.write_row(test_row(1, "Alice")).unwrap();
368 assert_eq!(status, WriteStatus::Continue);
369 assert!(!writer.output_started()); writer.finish().unwrap();
372 assert!(writer.output_started()); let result = String::from_utf8(output).unwrap();
375 assert!(result.contains("id"));
376 assert!(result.contains("Alice"));
377 }
378
379 #[test]
380 fn test_limit_enforcement() {
381 let mut output = Vec::new();
382 let formatter = Box::new(CsvFormatter::new());
383 let columns = test_columns();
384
385 let mut writer = StreamingWriter::new(&mut output, formatter, columns, Some(2));
386
387 writer.prepare(None).unwrap();
388
389 assert_eq!(
390 writer.write_row(test_row(1, "Alice")).unwrap(),
391 WriteStatus::Continue
392 );
393 assert_eq!(
394 writer.write_row(test_row(2, "Bob")).unwrap(),
395 WriteStatus::Continue
396 );
397 assert_eq!(
398 writer.write_row(test_row(3, "Charlie")).unwrap(),
399 WriteStatus::LimitReached
400 );
401
402 assert_eq!(writer.written_count(), 2);
403
404 writer.finish().unwrap();
405 }
406
407 #[test]
408 fn test_buffer_overflow_errors() {
409 let mut output = Vec::new();
410 let formatter = Box::new(TableFormatter::new());
411 let columns = test_columns();
412
413 let mut writer =
414 StreamingWriter::new(&mut output, formatter, columns, None).with_buffer_limit(40); writer.prepare(None).unwrap();
417 assert!(!writer.output_started());
418
419 assert_eq!(
421 writer.write_row(test_row(1, "Alice")).unwrap(),
422 WriteStatus::Continue
423 );
424 let err = writer.write_row(test_row(2, "Bob")).unwrap_err();
425 assert!(matches!(err, CliError::InvalidArgument(_)));
426 }
427
428 #[test]
429 fn test_empty_output() {
430 let mut output = Vec::new();
431 let formatter = Box::new(JsonFormatter::new());
432 let columns = test_columns();
433
434 let mut writer = StreamingWriter::new(&mut output, formatter, columns, None);
435
436 writer.prepare(None).unwrap();
437 writer.finish().unwrap();
438
439 let result = String::from_utf8(output).unwrap();
440 assert!(result.contains('['));
442 assert!(result.contains(']'));
443 }
444
445 #[test]
446 fn test_csv_streaming() {
447 let mut output = Vec::new();
448 let formatter = Box::new(CsvFormatter::new());
449 let columns = test_columns();
450
451 let mut writer = StreamingWriter::new(&mut output, formatter, columns, None);
452
453 writer.prepare(None).unwrap();
454 assert!(writer.output_started()); writer.write_row(test_row(1, "Alice")).unwrap();
457 writer.write_row(test_row(2, "Bob")).unwrap();
458 writer.finish().unwrap();
459
460 let result = String::from_utf8(output).unwrap();
461 assert_eq!(result, "id,name\n1,Alice\n2,Bob\n");
462 }
463
464 #[test]
465 fn test_written_count() {
466 let mut output = Vec::new();
467 let formatter = Box::new(CsvFormatter::new());
468 let columns = test_columns();
469
470 let mut writer = StreamingWriter::new(&mut output, formatter, columns, None);
471
472 writer.prepare(None).unwrap();
473
474 assert_eq!(writer.written_count(), 0);
475 writer.write_row(test_row(1, "Alice")).unwrap();
476 assert_eq!(writer.written_count(), 1);
477 writer.write_row(test_row(2, "Bob")).unwrap();
478 assert_eq!(writer.written_count(), 2);
479
480 writer.finish().unwrap();
481 }
482
483 #[test]
484 fn test_table_with_small_data() {
485 let mut output = Vec::new();
486 let formatter = Box::new(TableFormatter::new());
487 let columns = test_columns();
488
489 let mut writer = StreamingWriter::new(&mut output, formatter, columns, None);
490
491 writer.prepare(None).unwrap();
492 assert!(!writer.output_started()); writer.write_row(test_row(1, "Alice")).unwrap();
495 writer.write_row(test_row(2, "Bob")).unwrap();
496 writer.finish().unwrap();
497
498 let result = String::from_utf8(output).unwrap();
499 assert!(result.contains("id"));
501 assert!(result.contains("name"));
502 assert!(result.contains("Alice"));
503 assert!(result.contains("Bob"));
504 }
505
506 #[test]
507 fn test_streaming_large_row_count_does_not_buffer() {
508 let formatter = Box::new(JsonlFormatter::new());
509 let columns = test_columns();
510
511 let mut writer = StreamingWriter::new(std::io::sink(), formatter, columns, None);
512
513 writer.prepare(None).unwrap();
514 for i in 0..12_000 {
515 writer.write_row(test_row(i, "row")).unwrap();
516 }
517
518 assert_eq!(writer.written_count(), 12_000);
519 assert!(writer.buffer.is_empty());
520
521 writer.finish().unwrap();
522 }
523
524 #[test]
525 fn routing_report_is_serialized_to_its_own_writer_without_row_output() {
526 let report = DistributedReadRoutingReport::new(
527 "strong",
528 Some("strong".into()),
529 "cluster",
530 "retryable_failure",
531 Some("read point expired".into()),
532 );
533 let mut stderr = Vec::new();
534 write_distributed_read_routing_report(&mut stderr, RoutingReportFormat::Json, &report)
535 .unwrap();
536 let value: serde_json::Value = serde_json::from_slice(&stderr).unwrap();
537 assert_eq!(value["requested_mode"], "strong");
538 assert_eq!(value["outcome"], "retryable_failure");
539 assert_eq!(value["reason"], "read point expired");
540 }
541}