Skip to main content

alopex_cli/streaming/
writer.rs

1//! StreamingWriter - Streaming output controller
2//!
3//! Manages streaming output with buffer limits for non-streaming formats.
4
5use 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;
13/// Default buffer limit for non-streaming formats (table).
14pub const DEFAULT_BUFFER_LIMIT: usize = 10 * 1024 * 1024;
15
16/// A routing/status report kept separate from row output. Fields unavailable
17/// before a server coordinator responds remain null rather than being filled
18/// with local guesses.
19#[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
57/// Write the report to a caller-provided stderr writer. No `StreamingWriter`
58/// or SQL row formatter is involved, so table/JSON/CSV/TSV/JSONL stdout stays
59/// byte-for-byte independent from routing diagnostics.
60pub 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/// Status returned by write operations.
94#[derive(Debug, Clone, Copy, PartialEq, Eq)]
95pub enum WriteStatus {
96    /// Row was written successfully, continue writing.
97    Continue,
98    /// Limit reached, no more rows will be written.
99    LimitReached,
100}
101
102/// Streaming writer for output.
103///
104/// Controls streaming output with the following behaviors:
105/// - **Streaming formats** (json, jsonl, csv, tsv): Output rows immediately.
106/// - **Non-streaming formats** (table): Buffer rows up to `buffer_limit`.
107///
108/// # Output Boundary
109///
110/// - **Output started**: When `output_started == true` (header has been written).
111/// - **Buffer overflow**: Returns an error prompting `--output json|csv|tsv` or `--limit`.
112pub struct StreamingWriter<W> {
113    /// The underlying writer.
114    writer: W,
115    /// The formatter to use for output.
116    formatter: Box<dyn Formatter>,
117    /// Column definitions (schema).
118    columns: Vec<Column>,
119    /// Optional row limit.
120    limit: Option<usize>,
121    /// Buffer limit for non-streaming formats.
122    buffer_limit: usize,
123    /// Buffer for non-streaming formats (table).
124    buffer: Vec<Row>,
125    /// Approximate buffered size in bytes.
126    buffer_bytes: usize,
127    /// Number of rows written (for limit checking).
128    written_count: usize,
129    /// Whether header has been output (output started).
130    output_started: bool,
131    /// Whether quiet mode is enabled (suppress warnings).
132    quiet: bool,
133}
134
135impl<W: Write> StreamingWriter<W> {
136    /// Create a new StreamingWriter.
137    ///
138    /// # Arguments
139    ///
140    /// * `writer` - The output writer (e.g., stdout).
141    /// * `formatter` - The formatter to use for output.
142    /// * `columns` - Column definitions for the output.
143    /// * `limit` - Optional row limit.
144    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    /// Create a new StreamingWriter with a custom buffer limit.
165    #[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    /// Enable quiet mode (suppress warnings).
172    pub fn with_quiet(mut self, quiet: bool) -> Self {
173        self.quiet = quiet;
174        self
175    }
176
177    /// Check if quiet mode is enabled.
178    ///
179    /// When quiet mode is enabled, status-only output (OK messages) should be suppressed.
180    pub fn is_quiet(&self) -> bool {
181        self.quiet
182    }
183
184    /// Prepare the writer for output.
185    ///
186    /// For streaming formats (json, jsonl, csv, tsv), immediately outputs the header.
187    /// For non-streaming formats (table), defers header output.
188    ///
189    /// # Arguments
190    ///
191    /// * `row_count_hint` - Optional estimated row count (unused for buffer sizing).
192    pub fn prepare(&mut self, row_count_hint: Option<usize>) -> Result<()> {
193        let _ = row_count_hint;
194
195        // For streaming formats, output header immediately
196        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    /// Write a row to the output.
206    ///
207    /// For streaming formats, the row is output immediately.
208    /// For non-streaming formats, the row is buffered.
209    ///
210    /// # Returns
211    ///
212    /// * `WriteStatus::Continue` - Row was written, continue writing.
213    /// * `WriteStatus::LimitReached` - Row limit reached, stop writing.
214    ///
215    /// # Note
216    ///
217    /// The row is taken by ownership to avoid unnecessary cloning when buffering.
218    pub fn write_row(&mut self, row: Row) -> Result<WriteStatus> {
219        // Check limit
220        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            // Streaming format: output immediately
228            self.formatter.write_row(&mut self.writer, &row)?;
229            self.written_count += 1;
230        } else {
231            // Non-streaming format: buffer the row
232            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    /// Finish output, flushing any buffered rows and writing the footer.
249    pub fn finish(&mut self) -> Result<()> {
250        // For non-streaming formats, output header if not yet done
251        if !self.output_started {
252            self.formatter
253                .write_header(&mut self.writer, &self.columns)?;
254            self.output_started = true;
255
256            // Flush any buffered rows
257            for row in self.buffer.drain(..) {
258                self.formatter.write_row(&mut self.writer, &row)?;
259            }
260        }
261
262        // Write footer
263        self.formatter.write_footer(&mut self.writer)?;
264
265        Ok(())
266    }
267
268    /// Returns the number of rows written.
269    #[allow(dead_code)]
270    pub fn written_count(&self) -> usize {
271        self.written_count
272    }
273
274    /// Returns whether output has started.
275    #[allow(dead_code)]
276    pub fn output_started(&self) -> bool {
277        self.output_started
278    }
279
280    /// Returns approximate buffered bytes.
281    #[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()); // Header not output yet
360
361        let status = writer.write_row(test_row(1, "Alice")).unwrap();
362        assert_eq!(status, WriteStatus::Continue);
363        assert!(!writer.output_started()); // Still buffering
364
365        writer.finish().unwrap();
366        assert!(writer.output_started()); // Now output started
367
368        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); // Very small buffer
409
410        writer.prepare(None).unwrap();
411        assert!(!writer.output_started());
412
413        // Add rows to buffer
414        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        // JSON array format: should be valid empty array
435        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()); // CSV is streaming
449
450        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()); // Table is non-streaming
487
488        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        // Table output should contain the data
494        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}