Skip to main content

alopex_cli/commands/
sql.rs

1//! SQL Command - SQL query execution
2//!
3//! Supports: query execution, file-based queries
4
5use std::collections::HashSet;
6use std::fs;
7use std::io::{self, Read, Write};
8
9use alopex_embedded::Database;
10
11use crate::batch::{BatchMode, DistributedReadOutcome};
12use crate::cli::{OutputFormat, RoutingReportFormat, SqlCommand, SqlReadMode};
13use crate::client::http::{ClientError, HttpClient};
14use crate::error::{CliError, Result};
15use crate::models::{Column, DataType, Row, Value};
16use crate::output::formatter::{create_formatter, Formatter};
17use crate::profile::config::ResolvedSqlReadMode;
18use crate::streaming::timeout::parse_deadline;
19use crate::streaming::{
20    write_distributed_read_routing_report, CancelSignal, Deadline, DistributedReadRoutingReport,
21    StreamingWriter, WriteStatus,
22};
23use crate::tui::{is_tty, TuiApp};
24use crate::ui::mode::UiMode;
25use futures_util::StreamExt;
26use reqwest::Response;
27
28#[doc(hidden)]
29pub struct SqlExecutionOptions<'a> {
30    pub limit: Option<usize>,
31    pub quiet: bool,
32    pub cancel: &'a CancelSignal,
33    pub deadline: &'a Deadline,
34    pub admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
35}
36
37/// Source of formatters for SQL result output.
38///
39/// The `sql` command emits one result block per statement, so a single
40/// invocation may need more than one formatter instance.
41enum SqlOutput<'a> {
42    /// Create formatters from an output format. `OutputFormat::Json` also
43    /// wraps the output in an array of per-statement result sets.
44    Format(OutputFormat),
45    /// Obtain formatters from a caller-supplied factory (e.g. the admin TUI
46    /// capture formatter). No statement-array wrapping is applied.
47    Custom(&'a mut dyn FnMut() -> Box<dyn Formatter>),
48}
49
50impl SqlOutput<'_> {
51    /// Create a formatter for the next result block.
52    fn create_block_formatter(&mut self) -> Box<dyn Formatter> {
53        match self {
54            SqlOutput::Format(format) => create_formatter(*format),
55            SqlOutput::Custom(factory) => factory(),
56        }
57    }
58
59    /// Whether output must be wrapped in an array of per-statement result sets.
60    fn statement_array(&self) -> bool {
61        matches!(self, SqlOutput::Format(OutputFormat::Json))
62    }
63
64    /// Whether the produced formatters support streaming output.
65    fn supports_streaming(&mut self) -> bool {
66        match self {
67            SqlOutput::Format(format) => format.supports_streaming(),
68            SqlOutput::Custom(factory) => factory().supports_streaming(),
69        }
70    }
71}
72
73/// Write adapter that records whether any bytes have been written.
74///
75/// Used by the JSON statement-array wrapping to decide, on error paths, how
76/// many opened arrays must be closed to keep stdout well-formed JSON.
77struct CountingWriter<W> {
78    inner: W,
79    bytes: u64,
80}
81
82impl<W: Write> CountingWriter<W> {
83    fn new(inner: W) -> Self {
84        Self { inner, bytes: 0 }
85    }
86
87    fn bytes_written(&self) -> u64 {
88        self.bytes
89    }
90}
91
92impl<W: Write> Write for CountingWriter<W> {
93    fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
94        let written = self.inner.write(buf)?;
95        self.bytes += written as u64;
96        Ok(written)
97    }
98
99    fn flush(&mut self) -> io::Result<()> {
100        self.inner.flush()
101    }
102}
103
104/// Execute a SQL command with dynamic column detection.
105///
106/// This function creates formatters internally based on the query result type,
107/// ensuring that SELECT queries use the correct column headers.
108///
109/// # Output contract
110///
111/// - `--output json` emits an array of per-statement result sets; a single
112///   statement yields a 1-element array. With `--quiet`, DDL/DML status
113///   result sets are omitted, so input consisting only of such statements
114///   yields an empty array.
115/// - Other formats emit one result block per statement, in statement order.
116///
117/// # Arguments
118///
119/// * `db` - The database instance.
120/// * `cmd` - The SQL command to execute.
121/// * `writer` - The output writer.
122/// * `output_format` - The output format to use.
123/// * `limit` - Optional row limit.
124/// * `quiet` - Whether to suppress warnings and status result sets.
125#[allow(clippy::too_many_arguments)]
126pub fn execute_with_formatter<'a, W: Write>(
127    db: &Database,
128    cmd: SqlCommand,
129    batch_mode: &BatchMode,
130    ui_mode: UiMode,
131    writer: &mut W,
132    output_format: OutputFormat,
133    admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
134    limit: Option<usize>,
135    quiet: bool,
136) -> Result<()> {
137    let deadline = Deadline::new(parse_deadline(cmd.deadline.as_deref())?);
138    let cancel = CancelSignal::new();
139
140    execute_with_formatter_control(
141        db,
142        cmd,
143        batch_mode,
144        ui_mode,
145        writer,
146        output_format,
147        SqlExecutionOptions {
148            limit,
149            quiet,
150            cancel: &cancel,
151            deadline: &deadline,
152            admin_launcher,
153        },
154    )
155}
156
157/// Execute a SQL command using formatters supplied by `make_formatter`.
158///
159/// One formatter is created per result block (a multi-statement input emits
160/// one block per statement). No JSON statement-array wrapping is applied;
161/// this entry point is intended for callers that capture results with a
162/// custom formatter (e.g. the admin TUI).
163#[allow(clippy::too_many_arguments)]
164pub fn execute_with_formatter_factory<'a, W: Write>(
165    db: &Database,
166    cmd: SqlCommand,
167    batch_mode: &BatchMode,
168    ui_mode: UiMode,
169    writer: &mut W,
170    make_formatter: &mut dyn FnMut() -> Box<dyn Formatter>,
171    admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
172    limit: Option<usize>,
173    quiet: bool,
174) -> Result<()> {
175    let deadline = Deadline::new(parse_deadline(cmd.deadline.as_deref())?);
176    let cancel = CancelSignal::new();
177    let mut output = SqlOutput::Custom(make_formatter);
178
179    execute_with_output_control(
180        db,
181        cmd,
182        batch_mode,
183        ui_mode,
184        writer,
185        &mut output,
186        SqlExecutionOptions {
187            limit,
188            quiet,
189            cancel: &cancel,
190            deadline: &deadline,
191            admin_launcher,
192        },
193    )
194}
195
196#[doc(hidden)]
197pub fn execute_with_formatter_control<W: Write>(
198    db: &Database,
199    cmd: SqlCommand,
200    batch_mode: &BatchMode,
201    ui_mode: UiMode,
202    writer: &mut W,
203    output_format: OutputFormat,
204    options: SqlExecutionOptions<'_>,
205) -> Result<()> {
206    let mut output = SqlOutput::Format(output_format);
207    execute_with_output_control(db, cmd, batch_mode, ui_mode, writer, &mut output, options)
208}
209
210fn execute_with_output_control<W: Write>(
211    db: &Database,
212    cmd: SqlCommand,
213    batch_mode: &BatchMode,
214    ui_mode: UiMode,
215    writer: &mut W,
216    output: &mut SqlOutput<'_>,
217    mut options: SqlExecutionOptions<'_>,
218) -> Result<()> {
219    let sql = cmd.resolve_query(batch_mode)?;
220    let effective_limit = merge_limit(options.limit, cmd.max_rows);
221    options.limit = effective_limit;
222
223    if ui_mode == UiMode::Tui {
224        return execute_tui_local_or_fallback(db, &sql, writer, output, options);
225    }
226
227    execute_sql_with_formatter(db, &sql, writer, output, &options)
228}
229
230fn is_select_query(sql: &str) -> Result<bool> {
231    use alopex_sql::{AlopexDialect, Parser};
232
233    let dialect = AlopexDialect;
234    let stmts = Parser::parse_sql(&dialect, sql).map_err(|e| CliError::Parse(format!("{}", e)))?;
235    Ok(stmts.len() == 1
236        && stmts
237            .first()
238            .is_some_and(|statement| statement.kind.is_query()))
239}
240
241/// Execute a SQL command against a remote server using HttpClient.
242///
243/// `--output json` emits an array of result sets, matching the local and
244/// server-side per-statement output contract.
245#[allow(clippy::too_many_arguments)]
246pub async fn execute_remote_with_formatter<'a, W: Write>(
247    client: &HttpClient,
248    cmd: &SqlCommand,
249    batch_mode: &BatchMode,
250    ui_mode: UiMode,
251    writer: &mut W,
252    output_format: OutputFormat,
253    admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
254    limit: Option<usize>,
255    quiet: bool,
256) -> Result<()> {
257    let effective_limit = merge_limit(limit, cmd.max_rows);
258    let deadline = Deadline::new(parse_deadline(cmd.deadline.as_deref())?);
259    let cancel = CancelSignal::new();
260    let options = SqlExecutionOptions {
261        limit: effective_limit,
262        quiet,
263        cancel: &cancel,
264        deadline: &deadline,
265        admin_launcher,
266    };
267
268    execute_remote_with_formatter_control(
269        client,
270        cmd,
271        batch_mode,
272        ui_mode,
273        writer,
274        output_format,
275        options,
276    )
277    .await
278}
279
280/// Dispatch a server SQL command after profile-side routing resolution.
281///
282/// Cluster candidates deliberately use the versioned distributed-read route
283/// instead of the legacy local SQL endpoint. Until a server supplies a
284/// prepared-result response, any non-success remains a classified client
285/// error; it is never retried through local execution.
286#[allow(clippy::too_many_arguments)]
287pub async fn execute_remote_with_routing<'a, W: Write>(
288    client: &HttpClient,
289    cmd: &SqlCommand,
290    read_mode: ResolvedSqlReadMode,
291    batch_mode: &BatchMode,
292    ui_mode: UiMode,
293    writer: &mut W,
294    output_format: OutputFormat,
295    admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
296    limit: Option<usize>,
297    quiet: bool,
298) -> Result<()> {
299    let requested_mode = requested_read_mode_name(cmd.read_mode, read_mode);
300    let (effective_mode, decision) = match read_mode {
301        ResolvedSqlReadMode::Local => (Some("local".to_string()), "local"),
302        ResolvedSqlReadMode::Cluster(mode) => (Some(read_mode_name(mode).to_string()), "cluster"),
303    };
304    let result = match read_mode {
305        ResolvedSqlReadMode::Local => {
306            execute_remote_with_formatter(
307                client,
308                cmd,
309                batch_mode,
310                ui_mode,
311                writer,
312                output_format,
313                admin_launcher,
314                limit,
315                quiet,
316            )
317            .await
318        }
319        ResolvedSqlReadMode::Cluster(mode) => {
320            execute_remote_cluster_read(client, cmd, batch_mode, ui_mode, mode).await
321        }
322    };
323    if let Some(format) = cmd.routing_report {
324        emit_routing_report(format, requested_mode, effective_mode, decision, &result)?;
325    }
326    result
327}
328
329async fn execute_remote_cluster_read(
330    client: &HttpClient,
331    cmd: &SqlCommand,
332    batch_mode: &BatchMode,
333    ui_mode: UiMode,
334    mode: SqlReadMode,
335) -> Result<()> {
336    if ui_mode == UiMode::Tui {
337        return Err(distributed_read_error(
338            DistributedReadOutcome::Unsupported,
339            "cluster SQL read does not support TUI preview",
340        ));
341    }
342    let sql = cmd.resolve_query(batch_mode)?;
343    if !is_select_query(&sql)? {
344        return Err(distributed_read_error(
345            DistributedReadOutcome::Unsupported,
346            "remote DDL, DML, and transactions are not supported for cluster reads",
347        ));
348    }
349    let request = RemoteDistributedReadRequest {
350        sql,
351        read_mode: read_mode_name(mode),
352    };
353    // P2.12 currently returns a typed capability failure until a fenced
354    // coordinator is installed. If a future server unexpectedly returns a
355    // success envelope without a prepared-result adapter, fail closed rather
356    // than treating it as a local SQL success.
357    let _: serde_json::Value = client
358        .post_json("v1/sql/reads", &request)
359        .await
360        .map_err(map_distributed_client_error)?;
361    Err(distributed_read_error(
362        DistributedReadOutcome::TerminalFailure,
363        "distributed-read server response requires a prepared-result stream adapter",
364    ))
365}
366
367fn read_mode_name(mode: SqlReadMode) -> &'static str {
368    match mode {
369        SqlReadMode::Local => "local",
370        SqlReadMode::Inherit => "inherit",
371        SqlReadMode::Strong => "strong",
372        SqlReadMode::Stale => "stale",
373    }
374}
375
376fn requested_read_mode_name(
377    command_mode: Option<SqlReadMode>,
378    resolved_mode: ResolvedSqlReadMode,
379) -> &'static str {
380    let default = match resolved_mode {
381        ResolvedSqlReadMode::Local => SqlReadMode::Local,
382        ResolvedSqlReadMode::Cluster(_) => SqlReadMode::Inherit,
383    };
384    read_mode_name(command_mode.unwrap_or(default))
385}
386
387fn emit_routing_report(
388    format: RoutingReportFormat,
389    requested_mode: &str,
390    effective_mode: Option<String>,
391    decision: &str,
392    result: &Result<()>,
393) -> Result<()> {
394    let (outcome, reason) = match result {
395        Ok(()) => (DistributedReadOutcome::Success.as_str().to_string(), None),
396        Err(CliError::DistributedReadOutcome {
397            outcome, reason, ..
398        }) => (outcome.clone(), Some(reason.clone())),
399        Err(CliError::Cancelled) => (
400            DistributedReadOutcome::Cancelled.as_str().to_string(),
401            Some("client cancellation".into()),
402        ),
403        Err(CliError::Timeout(reason)) => (
404            DistributedReadOutcome::RetryableFailure
405                .as_str()
406                .to_string(),
407            Some(reason.clone()),
408        ),
409        Err(CliError::ServerUnsupported(reason)) => (
410            DistributedReadOutcome::Unsupported.as_str().to_string(),
411            Some(reason.clone()),
412        ),
413        Err(error) => (
414            DistributedReadOutcome::TerminalFailure.as_str().to_string(),
415            Some(error.to_string()),
416        ),
417    };
418    let report = DistributedReadRoutingReport::new(
419        requested_mode,
420        effective_mode,
421        decision,
422        outcome,
423        reason,
424    );
425    let stderr = std::io::stderr();
426    let mut stderr = stderr.lock();
427    write_distributed_read_routing_report(&mut stderr, format, &report)
428}
429
430fn distributed_read_error(outcome: DistributedReadOutcome, reason: impl Into<String>) -> CliError {
431    let reason = reason.into();
432    CliError::DistributedReadOutcome {
433        outcome: outcome.as_str().to_string(),
434        reason,
435        exit_code: outcome.exit_code(),
436    }
437}
438
439fn map_distributed_client_error(error: ClientError) -> CliError {
440    match error {
441        ClientError::Request { source, .. } => distributed_read_error(
442            DistributedReadOutcome::RetryableFailure,
443            format!("request failed: {source}"),
444        ),
445        ClientError::Auth(error) => distributed_read_error(
446            DistributedReadOutcome::AuthorizationFailure,
447            error.to_string(),
448        ),
449        ClientError::HttpStatus { status, body } => {
450            let code = serde_json::from_str::<serde_json::Value>(&body)
451                .ok()
452                .and_then(|value| {
453                    value
454                        .get("error")
455                        .and_then(|error| error.get("code"))
456                        .and_then(|code| code.as_str())
457                        .map(str::to_owned)
458                });
459            let outcome = match code.as_deref() {
460                Some("UNAUTHORIZED") | Some("AUTHORIZATION_FAILURE") => {
461                    DistributedReadOutcome::AuthorizationFailure
462                }
463                Some("UNSUPPORTED_REMOTE")
464                | Some("NOT_IMPLEMENTED")
465                | Some("FUTURE_DISTRIBUTED_EXECUTION_REQUIRED") => {
466                    DistributedReadOutcome::Unsupported
467                }
468                Some("CAPABILITY_UNAVAILABLE") => DistributedReadOutcome::TerminalFailure,
469                _ if status.is_server_error() || status.as_u16() == 408 => {
470                    DistributedReadOutcome::RetryableFailure
471                }
472                _ => DistributedReadOutcome::TerminalFailure,
473            };
474            distributed_read_error(outcome, format!("HTTP {}: {body}", status.as_u16()))
475        }
476        ClientError::InvalidUrl(message) | ClientError::Build(message) => {
477            distributed_read_error(DistributedReadOutcome::TerminalFailure, message)
478        }
479    }
480}
481
482fn sql_context_message(sql: &str) -> String {
483    let condensed = sql.split_whitespace().collect::<Vec<_>>().join(" ");
484    let max_len = 200;
485    if condensed.chars().count() > max_len {
486        let truncated: String = condensed.chars().take(max_len).collect();
487        format!("SQL: {truncated}...")
488    } else {
489        format!("SQL: {condensed}")
490    }
491}
492
493/// Execute a remote SQL command using formatters supplied by `make_formatter`.
494///
495/// See [`execute_with_formatter_factory`] for the output contract.
496#[allow(clippy::too_many_arguments)]
497pub async fn execute_remote_with_formatter_factory<'a, W: Write>(
498    client: &HttpClient,
499    cmd: &SqlCommand,
500    batch_mode: &BatchMode,
501    ui_mode: UiMode,
502    writer: &mut W,
503    make_formatter: &mut dyn FnMut() -> Box<dyn Formatter>,
504    admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
505    limit: Option<usize>,
506    quiet: bool,
507) -> Result<()> {
508    let effective_limit = merge_limit(limit, cmd.max_rows);
509    let deadline = Deadline::new(parse_deadline(cmd.deadline.as_deref())?);
510    let cancel = CancelSignal::new();
511    let options = SqlExecutionOptions {
512        limit: effective_limit,
513        quiet,
514        cancel: &cancel,
515        deadline: &deadline,
516        admin_launcher,
517    };
518    let mut output = SqlOutput::Custom(make_formatter);
519    execute_remote_with_output_control(
520        client,
521        cmd,
522        batch_mode,
523        ui_mode,
524        writer,
525        &mut output,
526        options,
527    )
528    .await
529}
530
531#[doc(hidden)]
532pub async fn execute_remote_with_formatter_control<W: Write>(
533    client: &HttpClient,
534    cmd: &SqlCommand,
535    batch_mode: &BatchMode,
536    ui_mode: UiMode,
537    writer: &mut W,
538    output_format: OutputFormat,
539    options: SqlExecutionOptions<'_>,
540) -> Result<()> {
541    let mut output = SqlOutput::Format(output_format);
542    execute_remote_with_output_control(
543        client,
544        cmd,
545        batch_mode,
546        ui_mode,
547        writer,
548        &mut output,
549        options,
550    )
551    .await
552}
553
554async fn execute_remote_with_output_control<W: Write>(
555    client: &HttpClient,
556    cmd: &SqlCommand,
557    batch_mode: &BatchMode,
558    ui_mode: UiMode,
559    writer: &mut W,
560    output: &mut SqlOutput<'_>,
561    options: SqlExecutionOptions<'_>,
562) -> Result<()> {
563    let sql = cmd.resolve_query(batch_mode)?;
564    if ui_mode == UiMode::Tui {
565        return execute_tui_remote_or_fallback(client, &sql, cmd, writer, output, options).await;
566    }
567    execute_remote_with_formatter_impl(client, &sql, cmd, writer, output, &options).await
568}
569
570async fn execute_remote_with_formatter_impl<W: Write>(
571    client: &HttpClient,
572    sql: &str,
573    cmd: &SqlCommand,
574    writer: &mut W,
575    output: &mut SqlOutput<'_>,
576    options: &SqlExecutionOptions<'_>,
577) -> Result<()> {
578    if is_select_query(sql)? && output.supports_streaming() {
579        return execute_remote_streaming(
580            client,
581            sql,
582            writer,
583            output,
584            options,
585            cmd.fetch_size,
586            cmd.max_rows,
587        )
588        .await;
589    }
590
591    let request = RemoteSqlRequest {
592        sql: sql.to_string(),
593        streaming: false,
594        fetch_size: cmd.fetch_size,
595        max_rows: cmd.max_rows,
596    };
597    let response: RemoteSqlResponse = tokio::select! {
598        result = tokio::time::timeout(options.deadline.remaining(), client.post_json("api/sql/query", &request)) => {
599            match result {
600                Ok(value) => value.map_err(map_client_error)?,
601                Err(_) => {
602                    let _ = send_cancel_request(client).await;
603                    return Err(CliError::Timeout(format!(
604                        "deadline exceeded after {}",
605                        humantime::format_duration(options.deadline.duration())
606                    )));
607                }
608            }
609        }
610        _ = options.cancel.wait() => {
611            let _ = send_cancel_request(client).await;
612            return Err(CliError::Cancelled);
613        }
614    };
615
616    // JSON output is an array of result sets, matching the server response.
617    let json_array = output.statement_array();
618    if json_array {
619        writeln!(writer, "[")?;
620    }
621    let mut emitted = false;
622    for response in response.into_results() {
623        if response.columns.is_empty() && options.quiet {
624            continue;
625        }
626        if json_array && emitted {
627            writeln!(writer, ",")?;
628        }
629        emitted = true;
630
631        if response.columns.is_empty() {
632            let message = match response.affected_rows {
633                Some(count) => format!("{count} row(s) affected"),
634                None => "Operation completed successfully".to_string(),
635            };
636            if json_array {
637                writeln!(writer, "[")?;
638            }
639            let columns = sql_status_columns();
640            {
641                let mut streaming_writer = StreamingWriter::new(
642                    &mut *writer,
643                    output.create_block_formatter(),
644                    columns,
645                    options.limit,
646                )
647                .with_quiet(options.quiet);
648                streaming_writer.prepare(Some(1))?;
649                let row = Row::new(vec![Value::Text("OK".to_string()), Value::Text(message)]);
650                streaming_writer.write_row(row)?;
651                streaming_writer.finish()?;
652            }
653            if json_array {
654                writeln!(writer, "]")?;
655            }
656            continue;
657        }
658
659        let columns: Vec<Column> = response
660            .columns
661            .iter()
662            .map(|col| Column::new(&col.name, data_type_from_string(&col.data_type)))
663            .collect();
664        if json_array {
665            writeln!(writer, "[")?;
666        }
667        {
668            let mut streaming_writer = StreamingWriter::new(
669                &mut *writer,
670                output.create_block_formatter(),
671                columns,
672                options.limit,
673            )
674            .with_quiet(options.quiet);
675            streaming_writer.prepare(Some(response.rows.len()))?;
676            for row in response.rows {
677                if options.cancel.is_cancelled() {
678                    let _ = send_cancel_request(client).await;
679                    return Err(CliError::Cancelled);
680                }
681                options.deadline.check()?;
682                let values = row.into_iter().map(remote_value_to_value).collect();
683                match streaming_writer.write_row(Row::new(values))? {
684                    WriteStatus::LimitReached => break,
685                    WriteStatus::Continue => {}
686                }
687            }
688            streaming_writer.finish()?;
689        }
690        if json_array {
691            writeln!(writer, "]")?;
692        }
693    }
694    if json_array {
695        writeln!(writer, "]")?;
696    }
697    Ok(())
698}
699
700fn execute_tui_local_or_fallback<'a, W: Write>(
701    db: &Database,
702    sql: &str,
703    writer: &mut W,
704    output: &mut SqlOutput<'_>,
705    mut options: SqlExecutionOptions<'a>,
706) -> Result<()> {
707    if !is_tty() {
708        if !options.quiet {
709            eprintln!("Warning: --tui requires a TTY, falling back to batch output.");
710        }
711        return execute_sql_with_formatter(db, sql, writer, output, &options);
712    }
713
714    let admin_launcher = options.admin_launcher.take();
715    match execute_tui_local(db, sql, &options, admin_launcher) {
716        Ok(()) => Ok(()),
717        Err(err) => {
718            if !options.quiet {
719                eprintln!("Warning: TUI failed ({err}); falling back to batch output.");
720            }
721            execute_sql_with_formatter(db, sql, writer, output, &options)
722        }
723    }
724}
725
726fn execute_tui_local<'a>(
727    db: &Database,
728    sql: &str,
729    options: &SqlExecutionOptions<'a>,
730    admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
731) -> Result<()> {
732    use alopex_sql::ExecutionResult;
733
734    options.deadline.check()?;
735    let result = db.execute_sql(sql)?;
736    options.deadline.check()?;
737
738    let (columns, rows) = match result {
739        ExecutionResult::Success => {
740            let columns = sql_status_columns();
741            let row = Row::new(vec![
742                Value::Text("OK".to_string()),
743                Value::Text("Operation completed successfully".to_string()),
744            ]);
745            (columns, vec![row])
746        }
747        ExecutionResult::RowsAffected(count) => {
748            let columns = sql_status_columns();
749            let row = Row::new(vec![
750                Value::Text("OK".to_string()),
751                Value::Text(format!("{count} row(s) affected")),
752            ]);
753            (columns, vec![row])
754        }
755        ExecutionResult::Query(query_result) => {
756            let columns = query_result
757                .columns
758                .iter()
759                .map(|col| Column::new(&col.name, DataType::Text))
760                .collect::<Vec<_>>();
761            let mut rows = Vec::with_capacity(query_result.rows.len());
762            for sql_row in query_result.rows {
763                let values = sql_row.into_iter().map(sql_value_to_value).collect();
764                rows.push(Row::new(values));
765            }
766            if let Some(limit) = options.limit {
767                rows.truncate(limit);
768            }
769            (columns, rows)
770        }
771    };
772
773    let app = TuiApp::new(columns, rows, "local", false)
774        .with_context_message(Some(sql_context_message(sql)))
775        .with_admin_launcher(admin_launcher);
776    app.run()
777}
778
779async fn execute_tui_remote_or_fallback<'a, W: Write>(
780    client: &HttpClient,
781    sql: &str,
782    cmd: &SqlCommand,
783    writer: &mut W,
784    output: &mut SqlOutput<'_>,
785    mut options: SqlExecutionOptions<'a>,
786) -> Result<()> {
787    if !is_tty() {
788        if !options.quiet {
789            eprintln!("Warning: --tui requires a TTY, falling back to batch output.");
790        }
791        return execute_remote_with_formatter_impl(client, sql, cmd, writer, output, &options)
792            .await;
793    }
794
795    let admin_launcher = options.admin_launcher.take();
796    match execute_tui_remote(client, sql, cmd, &options, admin_launcher).await {
797        Ok(()) => Ok(()),
798        Err(err) => {
799            if !options.quiet {
800                eprintln!("Warning: TUI failed ({err}); falling back to batch output.");
801            }
802            execute_remote_with_formatter_impl(client, sql, cmd, writer, output, &options).await
803        }
804    }
805}
806
807async fn execute_tui_remote<'a>(
808    client: &HttpClient,
809    sql: &str,
810    cmd: &SqlCommand,
811    options: &SqlExecutionOptions<'a>,
812    admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
813) -> Result<()> {
814    if is_select_query(sql)? {
815        let (columns, rows) =
816            collect_remote_streaming_rows(client, sql, options, cmd.fetch_size, cmd.max_rows)
817                .await?;
818        let app = TuiApp::new(columns, rows, "server", false)
819            .with_context_message(Some(sql_context_message(sql)))
820            .with_admin_launcher(admin_launcher);
821        return app.run();
822    }
823
824    let request = RemoteSqlRequest {
825        sql: sql.to_string(),
826        streaming: false,
827        fetch_size: cmd.fetch_size,
828        max_rows: cmd.max_rows,
829    };
830
831    let response: RemoteSqlResponse = tokio::select! {
832        result = tokio::time::timeout(options.deadline.remaining(), client.post_json("api/sql/query", &request)) => {
833            match result {
834                Ok(value) => value.map_err(map_client_error)?,
835                Err(_) => {
836                    let _ = send_cancel_request(client).await;
837                    return Err(CliError::Timeout(format!(
838                        "deadline exceeded after {}",
839                        humantime::format_duration(options.deadline.duration())
840                    )));
841                }
842            }
843        }
844        _ = options.cancel.wait() => {
845            let _ = send_cancel_request(client).await;
846            return Err(CliError::Cancelled);
847        }
848    };
849
850    let response = response
851        .into_results()
852        .into_iter()
853        .last()
854        .unwrap_or_default();
855    let (columns, rows) = if response.columns.is_empty() {
856        let columns = sql_status_columns();
857        let message = match response.affected_rows {
858            Some(count) => format!("{count} row(s) affected"),
859            None => "Operation completed successfully".to_string(),
860        };
861        let row = Row::new(vec![Value::Text("OK".to_string()), Value::Text(message)]);
862        (columns, vec![row])
863    } else {
864        let columns: Vec<Column> = response
865            .columns
866            .iter()
867            .map(|col| Column::new(&col.name, data_type_from_string(&col.data_type)))
868            .collect();
869        let mut rows = response
870            .rows
871            .into_iter()
872            .map(|row| Row::new(row.into_iter().map(remote_value_to_value).collect()))
873            .collect::<Vec<_>>();
874        if let Some(limit) = options.limit {
875            rows.truncate(limit);
876        }
877        (columns, rows)
878    };
879
880    let app = TuiApp::new(columns, rows, "server", false)
881        .with_context_message(Some(sql_context_message(sql)))
882        .with_admin_launcher(admin_launcher);
883    app.run()
884}
885
886async fn collect_remote_streaming_rows(
887    client: &HttpClient,
888    sql: &str,
889    options: &SqlExecutionOptions<'_>,
890    fetch_size: Option<usize>,
891    max_rows: Option<usize>,
892) -> Result<(Vec<Column>, Vec<Row>)> {
893    let request = RemoteSqlRequest {
894        sql: sql.to_string(),
895        streaming: true,
896        fetch_size,
897        max_rows,
898    };
899
900    let response = tokio::select! {
901        result = tokio::time::timeout(options.deadline.remaining(), client.post_json_stream("api/sql/query", &request)) => {
902            match result {
903                Ok(value) => value.map_err(map_client_error)?,
904                Err(_) => {
905                    let _ = send_cancel_request(client).await;
906                    return Err(CliError::Timeout(format!(
907                        "deadline exceeded after {}",
908                        humantime::format_duration(options.deadline.duration())
909                    )));
910                }
911            }
912        }
913        _ = options.cancel.wait() => {
914            let _ = send_cancel_request(client).await;
915            return Err(CliError::Cancelled);
916        }
917    };
918
919    if let Some(content_type) = response
920        .headers()
921        .get(reqwest::header::CONTENT_TYPE)
922        .and_then(|value| value.to_str().ok())
923    {
924        if content_type.starts_with("application/jsonl") {
925            return collect_remote_jsonl_rows(client, response, options).await;
926        }
927    }
928
929    let mut stream = response.bytes_stream();
930    let mut buffer: Vec<u8> = Vec::new();
931    let mut pos: usize = 0;
932    let mut done = false;
933    let mut saw_array_start = false;
934    let mut columns: Option<Vec<String>> = None;
935    let mut column_set: Option<HashSet<String>> = None;
936    let mut rows: Vec<Row> = Vec::new();
937
938    while !done {
939        if options.cancel.is_cancelled() {
940            let _ = send_cancel_request(client).await;
941            return Err(CliError::Cancelled);
942        }
943        if let Err(err) = options.deadline.check() {
944            let _ = send_cancel_request(client).await;
945            return Err(err);
946        }
947
948        let next = tokio::select! {
949            _ = options.cancel.wait() => {
950                let _ = send_cancel_request(client).await;
951                return Err(CliError::Cancelled);
952            }
953            result = tokio::time::timeout(options.deadline.remaining(), stream.next()) => {
954                match result {
955                    Ok(value) => value,
956                    Err(_) => {
957                        let _ = send_cancel_request(client).await;
958                        return Err(CliError::Timeout(format!(
959                            "deadline exceeded after {}",
960                            humantime::format_duration(options.deadline.duration())
961                        )));
962                    }
963                }
964            }
965        };
966
967        let chunk = match next {
968            Some(chunk) => chunk,
969            None => break,
970        };
971
972        let bytes = match chunk {
973            Ok(bytes) => bytes,
974            Err(err) => return Err(CliError::ServerConnection(format!("request failed: {err}"))),
975        };
976
977        buffer.extend_from_slice(&bytes);
978
979        loop {
980            skip_whitespace(&buffer, &mut pos);
981            if pos >= buffer.len() {
982                break;
983            }
984
985            if !saw_array_start {
986                if buffer[pos] != b'[' {
987                    return Err(CliError::InvalidArgument(
988                        "Invalid streaming response: expected JSON array".into(),
989                    ));
990                }
991                pos += 1;
992                saw_array_start = true;
993                continue;
994            }
995
996            skip_whitespace(&buffer, &mut pos);
997            if pos >= buffer.len() {
998                break;
999            }
1000
1001            if buffer[pos] == b']' {
1002                pos += 1;
1003                done = true;
1004                break;
1005            }
1006
1007            let slice = &buffer[pos..];
1008            let mut stream =
1009                serde_json::Deserializer::from_slice(slice).into_iter::<serde_json::Value>();
1010            let value = match stream.next() {
1011                Some(Ok(value)) => value,
1012                Some(Err(err)) if err.is_eof() => break,
1013                Some(Err(err)) => return Err(CliError::Json(err)),
1014                None => break,
1015            };
1016            pos = pos.saturating_add(stream.byte_offset());
1017
1018            let object = value.as_object().ok_or_else(|| {
1019                CliError::InvalidArgument("Invalid streaming row: expected JSON object".into())
1020            })?;
1021
1022            if columns.is_none() {
1023                let names: Vec<String> = object.keys().cloned().collect();
1024                let set: HashSet<String> = names.iter().cloned().collect();
1025                if names.is_empty() {
1026                    return Err(CliError::InvalidArgument(
1027                        "Invalid streaming row: empty object".into(),
1028                    ));
1029                }
1030                columns = Some(names);
1031                column_set = Some(set);
1032            }
1033
1034            let names = columns
1035                .as_ref()
1036                .ok_or_else(|| CliError::InvalidArgument("Missing columns".into()))?;
1037            let set = column_set
1038                .as_ref()
1039                .ok_or_else(|| CliError::InvalidArgument("Missing column set".into()))?;
1040
1041            if object.len() != names.len() || !object.keys().all(|key| set.contains(key)) {
1042                return Err(CliError::InvalidArgument(
1043                    "Invalid streaming row: column mismatch".into(),
1044                ));
1045            }
1046
1047            let values = names
1048                .iter()
1049                .map(|name| {
1050                    object.get(name).ok_or_else(|| {
1051                        CliError::InvalidArgument(format!(
1052                            "Invalid streaming row: missing column '{name}'"
1053                        ))
1054                    })
1055                })
1056                .map(|value| value.and_then(json_value_to_value))
1057                .collect::<Result<Vec<_>>>()?;
1058            rows.push(Row::new(values));
1059
1060            if let Some(limit) = options.limit {
1061                if rows.len() >= limit {
1062                    let _ = send_cancel_request(client).await;
1063                    done = true;
1064                    break;
1065                }
1066            }
1067
1068            skip_whitespace(&buffer, &mut pos);
1069            if pos >= buffer.len() {
1070                break;
1071            }
1072            match buffer[pos] {
1073                b',' => {
1074                    pos += 1;
1075                }
1076                b']' => {
1077                    pos += 1;
1078                    done = true;
1079                    break;
1080                }
1081                _ => {
1082                    return Err(CliError::InvalidArgument(
1083                        "Invalid streaming response: expected ',' or ']'".into(),
1084                    ))
1085                }
1086            }
1087        }
1088
1089        if pos > 0 {
1090            buffer.drain(..pos);
1091            pos = 0;
1092        }
1093    }
1094
1095    if done {
1096        if has_non_whitespace(&buffer) {
1097            return Err(CliError::InvalidArgument(
1098                "Invalid streaming response: unexpected trailing data".into(),
1099            ));
1100        }
1101        buffer.clear();
1102        loop {
1103            let next = tokio::select! {
1104                _ = options.cancel.wait() => {
1105                    let _ = send_cancel_request(client).await;
1106                    return Err(CliError::Cancelled);
1107                }
1108                result = tokio::time::timeout(options.deadline.remaining(), stream.next()) => {
1109                    match result {
1110                        Ok(value) => value,
1111                        Err(_) => {
1112                            let _ = send_cancel_request(client).await;
1113                            return Err(CliError::Timeout(format!(
1114                                "deadline exceeded after {}",
1115                                humantime::format_duration(options.deadline.duration())
1116                            )));
1117                        }
1118                    }
1119                }
1120            };
1121
1122            let chunk = match next {
1123                Some(chunk) => chunk,
1124                None => break,
1125            };
1126
1127            let bytes = match chunk {
1128                Ok(bytes) => bytes,
1129                Err(err) => {
1130                    return Err(CliError::ServerConnection(format!("request failed: {err}")))
1131                }
1132            };
1133
1134            if has_non_whitespace(&bytes) {
1135                return Err(CliError::InvalidArgument(
1136                    "Invalid streaming response: unexpected trailing data".into(),
1137                ));
1138            }
1139        }
1140    } else {
1141        skip_whitespace(&buffer, &mut pos);
1142        if pos < buffer.len() {
1143            return Err(CliError::InvalidArgument(
1144                "Invalid streaming response: unexpected trailing data".into(),
1145            ));
1146        }
1147        return Err(CliError::InvalidArgument(
1148            "Invalid streaming response: unexpected end of stream".into(),
1149        ));
1150    }
1151
1152    let columns = columns
1153        .unwrap_or_default()
1154        .into_iter()
1155        .map(|name| Column::new(name, DataType::Text))
1156        .collect();
1157    Ok((columns, rows))
1158}
1159
1160#[allow(dead_code)]
1161async fn collect_remote_non_streaming_rows(
1162    client: &HttpClient,
1163    sql: &str,
1164    options: &SqlExecutionOptions<'_>,
1165    fetch_size: Option<usize>,
1166    max_rows: Option<usize>,
1167) -> Result<(Vec<Column>, Vec<Row>)> {
1168    let request = RemoteSqlRequest {
1169        sql: sql.to_string(),
1170        streaming: false,
1171        fetch_size,
1172        max_rows,
1173    };
1174    let response: RemoteSqlResponse = tokio::select! {
1175        result = tokio::time::timeout(options.deadline.remaining(), client.post_json("api/sql/query", &request)) => {
1176            match result {
1177                Ok(value) => value.map_err(map_client_error)?,
1178                Err(_) => {
1179                    let _ = send_cancel_request(client).await;
1180                    return Err(CliError::Timeout(format!(
1181                        "deadline exceeded after {}",
1182                        humantime::format_duration(options.deadline.duration())
1183                    )));
1184                }
1185            }
1186        }
1187        _ = options.cancel.wait() => {
1188            let _ = send_cancel_request(client).await;
1189            return Err(CliError::Cancelled);
1190        }
1191    };
1192
1193    let response = response
1194        .into_results()
1195        .into_iter()
1196        .last()
1197        .unwrap_or_default();
1198    if response.columns.is_empty() {
1199        let message = match response.affected_rows {
1200            Some(count) => format!("{count} row(s) affected"),
1201            None => "Operation completed successfully".to_string(),
1202        };
1203        let columns = sql_status_columns();
1204        let row = Row::new(vec![Value::Text("OK".to_string()), Value::Text(message)]);
1205        return Ok((columns, vec![row]));
1206    }
1207
1208    let columns = response
1209        .columns
1210        .iter()
1211        .map(|col| Column::new(&col.name, data_type_from_string(&col.data_type)))
1212        .collect::<Vec<_>>();
1213    let mut rows = response
1214        .rows
1215        .into_iter()
1216        .map(|row| Row::new(row.into_iter().map(remote_value_to_value).collect()))
1217        .collect::<Vec<_>>();
1218    if let Some(limit) = options.limit {
1219        rows.truncate(limit);
1220    }
1221    Ok((columns, rows))
1222}
1223
1224async fn execute_remote_streaming<W: Write>(
1225    client: &HttpClient,
1226    sql: &str,
1227    writer: &mut W,
1228    output: &mut SqlOutput<'_>,
1229    options: &SqlExecutionOptions<'_>,
1230    fetch_size: Option<usize>,
1231    max_rows: Option<usize>,
1232) -> Result<()> {
1233    let request = RemoteSqlRequest {
1234        sql: sql.to_string(),
1235        streaming: true,
1236        fetch_size,
1237        max_rows,
1238    };
1239
1240    let response = tokio::select! {
1241        result = tokio::time::timeout(options.deadline.remaining(), client.post_json_stream("api/sql/query", &request)) => {
1242            match result {
1243                Ok(value) => value.map_err(map_client_error)?,
1244                Err(_) => {
1245                    let _ = send_cancel_request(client).await;
1246                    return Err(CliError::Timeout(format!(
1247                        "deadline exceeded after {}",
1248                        humantime::format_duration(options.deadline.duration())
1249                    )));
1250                }
1251            }
1252        }
1253        _ = options.cancel.wait() => {
1254            let _ = send_cancel_request(client).await;
1255            return Err(CliError::Cancelled);
1256        }
1257    };
1258
1259    // The streaming API is used only for a single SELECT; wrap that one
1260    // result set in the same outer array as the non-streaming path.
1261    if !output.statement_array() {
1262        return stream_remote_result_set(
1263            client,
1264            response,
1265            writer,
1266            output.create_block_formatter(),
1267            options,
1268        )
1269        .await;
1270    }
1271
1272    writeln!(writer, "[")?;
1273    let mut result_set_writer = CountingWriter::new(&mut *writer);
1274    let result = stream_remote_result_set(
1275        client,
1276        response,
1277        &mut result_set_writer,
1278        output.create_block_formatter(),
1279        options,
1280    )
1281    .await;
1282    let result_set_started = result_set_writer.bytes_written() > 0;
1283    match result {
1284        Ok(()) => {
1285            writeln!(writer, "]")?;
1286            Ok(())
1287        }
1288        Err(err) => {
1289            // Keep stdout well-formed JSON on stream errors: rows are always
1290            // written as complete objects, so closing the result-set array
1291            // (when it was opened) and the statement array leaves valid JSON
1292            // alongside the non-zero exit code.
1293            if result_set_started {
1294                let _ = writeln!(writer, "]");
1295            }
1296            let _ = writeln!(writer, "]");
1297            Err(err)
1298        }
1299    }
1300}
1301
1302/// Stream a single remote result set (JSON array or JSONL body) to `writer`.
1303async fn stream_remote_result_set<W: Write>(
1304    client: &HttpClient,
1305    response: Response,
1306    writer: &mut W,
1307    formatter: Box<dyn Formatter>,
1308    options: &SqlExecutionOptions<'_>,
1309) -> Result<()> {
1310    if let Some(content_type) = response
1311        .headers()
1312        .get(reqwest::header::CONTENT_TYPE)
1313        .and_then(|value| value.to_str().ok())
1314    {
1315        if content_type.starts_with("application/jsonl") {
1316            return execute_remote_jsonl_streaming(client, response, writer, formatter, options)
1317                .await;
1318        }
1319    }
1320
1321    let mut stream = response.bytes_stream();
1322    let mut buffer: Vec<u8> = Vec::new();
1323    let mut pos: usize = 0;
1324    let mut streaming_writer: Option<StreamingWriter<&mut W>> = None;
1325    let mut formatter = Some(formatter);
1326    let mut columns: Option<Vec<String>> = None;
1327    let mut column_set: Option<HashSet<String>> = None;
1328    let mut done = false;
1329    let mut saw_array_start = false;
1330
1331    while !done {
1332        if options.cancel.is_cancelled() {
1333            let _ = send_cancel_request(client).await;
1334            return Err(CliError::Cancelled);
1335        }
1336        if let Err(err) = options.deadline.check() {
1337            let _ = send_cancel_request(client).await;
1338            return Err(err);
1339        }
1340
1341        let next = tokio::select! {
1342            _ = options.cancel.wait() => {
1343                let _ = send_cancel_request(client).await;
1344                return Err(CliError::Cancelled);
1345            }
1346            result = tokio::time::timeout(options.deadline.remaining(), stream.next()) => {
1347                match result {
1348                    Ok(value) => value,
1349                    Err(_) => {
1350                        let _ = send_cancel_request(client).await;
1351                        return Err(CliError::Timeout(format!(
1352                            "deadline exceeded after {}",
1353                            humantime::format_duration(options.deadline.duration())
1354                        )));
1355                    }
1356                }
1357            }
1358        };
1359
1360        let chunk = match next {
1361            Some(chunk) => chunk,
1362            None => break,
1363        };
1364
1365        let bytes = match chunk {
1366            Ok(bytes) => bytes,
1367            Err(err) => return Err(CliError::ServerConnection(format!("request failed: {err}"))),
1368        };
1369
1370        buffer.extend_from_slice(&bytes);
1371
1372        loop {
1373            skip_whitespace(&buffer, &mut pos);
1374            if pos >= buffer.len() {
1375                break;
1376            }
1377
1378            if !saw_array_start {
1379                if buffer[pos] != b'[' {
1380                    return Err(CliError::InvalidArgument(
1381                        "Invalid streaming response: expected JSON array".into(),
1382                    ));
1383                }
1384                pos += 1;
1385                saw_array_start = true;
1386                continue;
1387            }
1388
1389            skip_whitespace(&buffer, &mut pos);
1390            if pos >= buffer.len() {
1391                break;
1392            }
1393
1394            if buffer[pos] == b']' {
1395                pos += 1;
1396                done = true;
1397                break;
1398            }
1399
1400            let slice = &buffer[pos..];
1401            let mut stream =
1402                serde_json::Deserializer::from_slice(slice).into_iter::<serde_json::Value>();
1403            let value = match stream.next() {
1404                Some(Ok(value)) => value,
1405                Some(Err(err)) if err.is_eof() => break,
1406                Some(Err(err)) => return Err(CliError::Json(err)),
1407                None => break,
1408            };
1409            pos = pos.saturating_add(stream.byte_offset());
1410
1411            let object = value.as_object().ok_or_else(|| {
1412                CliError::InvalidArgument("Invalid streaming row: expected JSON object".into())
1413            })?;
1414
1415            if columns.is_none() {
1416                let names: Vec<String> = object.keys().cloned().collect();
1417                let set: HashSet<String> = names.iter().cloned().collect();
1418                if names.is_empty() {
1419                    return Err(CliError::InvalidArgument(
1420                        "Invalid streaming row: empty object".into(),
1421                    ));
1422                }
1423                let cols = names
1424                    .iter()
1425                    .map(|name| Column::new(name, DataType::Text))
1426                    .collect::<Vec<_>>();
1427                let formatter = formatter
1428                    .take()
1429                    .ok_or_else(|| CliError::InvalidArgument("Missing formatter".into()))?;
1430                let mut writer = StreamingWriter::new(&mut *writer, formatter, cols, options.limit)
1431                    .with_quiet(options.quiet);
1432                writer.prepare(None)?;
1433                streaming_writer = Some(writer);
1434                columns = Some(names);
1435                column_set = Some(set);
1436            }
1437
1438            let names = columns
1439                .as_ref()
1440                .ok_or_else(|| CliError::InvalidArgument("Missing columns".into()))?;
1441            let set = column_set
1442                .as_ref()
1443                .ok_or_else(|| CliError::InvalidArgument("Missing column set".into()))?;
1444
1445            if object.len() != names.len() || !object.keys().all(|key| set.contains(key)) {
1446                return Err(CliError::InvalidArgument(
1447                    "Invalid streaming row: column mismatch".into(),
1448                ));
1449            }
1450
1451            let values = names
1452                .iter()
1453                .map(|name| {
1454                    object.get(name).ok_or_else(|| {
1455                        CliError::InvalidArgument(format!(
1456                            "Invalid streaming row: missing column '{name}'"
1457                        ))
1458                    })
1459                })
1460                .map(|value| value.and_then(json_value_to_value))
1461                .collect::<Result<Vec<_>>>()?;
1462
1463            if let Some(writer) = streaming_writer.as_mut() {
1464                match writer.write_row(Row::new(values))? {
1465                    WriteStatus::LimitReached => {
1466                        let _ = send_cancel_request(client).await;
1467                        return writer.finish();
1468                    }
1469                    WriteStatus::Continue => {}
1470                }
1471            }
1472
1473            skip_whitespace(&buffer, &mut pos);
1474            if pos >= buffer.len() {
1475                break;
1476            }
1477            match buffer[pos] {
1478                b',' => {
1479                    pos += 1;
1480                }
1481                b']' => {
1482                    pos += 1;
1483                    done = true;
1484                    break;
1485                }
1486                _ => {
1487                    return Err(CliError::InvalidArgument(
1488                        "Invalid streaming response: expected ',' or ']'".into(),
1489                    ))
1490                }
1491            }
1492        }
1493
1494        if pos > 0 {
1495            buffer.drain(..pos);
1496            pos = 0;
1497        }
1498    }
1499
1500    if done {
1501        if has_non_whitespace(&buffer) {
1502            return Err(CliError::InvalidArgument(
1503                "Invalid streaming response: unexpected trailing data".into(),
1504            ));
1505        }
1506        buffer.clear();
1507        loop {
1508            let next = tokio::select! {
1509                _ = options.cancel.wait() => {
1510                    let _ = send_cancel_request(client).await;
1511                    return Err(CliError::Cancelled);
1512                }
1513                result = tokio::time::timeout(options.deadline.remaining(), stream.next()) => {
1514                    match result {
1515                        Ok(value) => value,
1516                        Err(_) => {
1517                            let _ = send_cancel_request(client).await;
1518                            return Err(CliError::Timeout(format!(
1519                                "deadline exceeded after {}",
1520                                humantime::format_duration(options.deadline.duration())
1521                            )));
1522                        }
1523                    }
1524                }
1525            };
1526
1527            let chunk = match next {
1528                Some(chunk) => chunk,
1529                None => break,
1530            };
1531
1532            let bytes = match chunk {
1533                Ok(bytes) => bytes,
1534                Err(err) => {
1535                    return Err(CliError::ServerConnection(format!("request failed: {err}")))
1536                }
1537            };
1538
1539            if has_non_whitespace(&bytes) {
1540                return Err(CliError::InvalidArgument(
1541                    "Invalid streaming response: unexpected trailing data".into(),
1542                ));
1543            }
1544        }
1545    } else {
1546        skip_whitespace(&buffer, &mut pos);
1547        if pos < buffer.len() {
1548            return Err(CliError::InvalidArgument(
1549                "Invalid streaming response: unexpected trailing data".into(),
1550            ));
1551        }
1552        return Err(CliError::InvalidArgument(
1553            "Invalid streaming response: unexpected end of stream".into(),
1554        ));
1555    }
1556
1557    if let Some(mut writer) = streaming_writer {
1558        return writer.finish();
1559    }
1560
1561    if done && saw_array_start {
1562        if let Some(formatter) = formatter.take() {
1563            let mut writer =
1564                StreamingWriter::new(&mut *writer, formatter, Vec::new(), options.limit)
1565                    .with_quiet(options.quiet);
1566            writer.prepare(None)?;
1567            return writer.finish();
1568        }
1569    }
1570
1571    Ok(())
1572}
1573
1574async fn collect_remote_jsonl_rows(
1575    client: &HttpClient,
1576    response: Response,
1577    options: &SqlExecutionOptions<'_>,
1578) -> Result<(Vec<Column>, Vec<Row>)> {
1579    let mut stream = response.bytes_stream();
1580    let mut buffer: Vec<u8> = Vec::new();
1581    let mut columns: Option<Vec<Column>> = None;
1582    let mut rows: Vec<Row> = Vec::new();
1583    let mut done = false;
1584
1585    while !done {
1586        if options.cancel.is_cancelled() {
1587            let _ = send_cancel_request(client).await;
1588            return Err(CliError::Cancelled);
1589        }
1590        if let Err(err) = options.deadline.check() {
1591            let _ = send_cancel_request(client).await;
1592            return Err(err);
1593        }
1594
1595        let next = tokio::select! {
1596            _ = options.cancel.wait() => {
1597                let _ = send_cancel_request(client).await;
1598                return Err(CliError::Cancelled);
1599            }
1600            result = tokio::time::timeout(options.deadline.remaining(), stream.next()) => {
1601                match result {
1602                    Ok(value) => value,
1603                    Err(_) => {
1604                        let _ = send_cancel_request(client).await;
1605                        return Err(CliError::Timeout(format!(
1606                            "deadline exceeded after {}",
1607                            humantime::format_duration(options.deadline.duration())
1608                        )));
1609                    }
1610                }
1611            }
1612        };
1613
1614        let chunk = match next {
1615            Some(chunk) => chunk,
1616            None => break,
1617        };
1618
1619        let bytes = match chunk {
1620            Ok(bytes) => bytes,
1621            Err(err) => return Err(CliError::ServerConnection(format!("request failed: {err}"))),
1622        };
1623
1624        buffer.extend_from_slice(&bytes);
1625
1626        while let Some(newline) = buffer.iter().position(|&b| b == b'\n') {
1627            let line = buffer.drain(..=newline).collect::<Vec<u8>>();
1628            if let Some(item) = parse_jsonl_line(&line)? {
1629                if let Some(error) = item.error {
1630                    return Err(CliError::InvalidArgument(format!(
1631                        "Server error: {}",
1632                        error.message
1633                    )));
1634                }
1635                if item.done {
1636                    done = true;
1637                    break;
1638                }
1639                let row = item.row.ok_or_else(|| {
1640                    CliError::InvalidArgument("Invalid streaming response: missing row".into())
1641                })?;
1642                if columns.is_none() {
1643                    columns = Some(default_stream_columns(row.len()));
1644                }
1645                rows.push(Row::new(
1646                    row.into_iter().map(remote_value_to_value).collect(),
1647                ));
1648                if let Some(limit) = options.limit {
1649                    if rows.len() >= limit {
1650                        let _ = send_cancel_request(client).await;
1651                        done = true;
1652                        break;
1653                    }
1654                }
1655            }
1656        }
1657    }
1658
1659    if !done {
1660        if let Some(item) = parse_jsonl_line(&buffer)? {
1661            if let Some(error) = item.error {
1662                return Err(CliError::InvalidArgument(format!(
1663                    "Server error: {}",
1664                    error.message
1665                )));
1666            }
1667            if item.done {
1668                done = true;
1669            } else if let Some(row) = item.row {
1670                if columns.is_none() {
1671                    columns = Some(default_stream_columns(row.len()));
1672                }
1673                rows.push(Row::new(
1674                    row.into_iter().map(remote_value_to_value).collect(),
1675                ));
1676                if let Some(limit) = options.limit {
1677                    if rows.len() >= limit {
1678                        let _ = send_cancel_request(client).await;
1679                        done = true;
1680                    }
1681                }
1682            } else {
1683                return Err(CliError::InvalidArgument(
1684                    "Invalid streaming response: missing row".into(),
1685                ));
1686            }
1687        }
1688    }
1689
1690    if !done {
1691        return Err(CliError::InvalidArgument(
1692            "Invalid streaming response: unexpected end of stream".into(),
1693        ));
1694    }
1695
1696    Ok((columns.unwrap_or_default(), rows))
1697}
1698
1699async fn execute_remote_jsonl_streaming<W: Write>(
1700    client: &HttpClient,
1701    response: Response,
1702    writer: &mut W,
1703    formatter: Box<dyn Formatter>,
1704    options: &SqlExecutionOptions<'_>,
1705) -> Result<()> {
1706    let mut stream = response.bytes_stream();
1707    let mut buffer: Vec<u8> = Vec::new();
1708    let mut streaming_writer: Option<StreamingWriter<&mut W>> = None;
1709    let mut formatter = Some(formatter);
1710    let mut done = false;
1711
1712    while !done {
1713        if options.cancel.is_cancelled() {
1714            let _ = send_cancel_request(client).await;
1715            return Err(CliError::Cancelled);
1716        }
1717        if let Err(err) = options.deadline.check() {
1718            let _ = send_cancel_request(client).await;
1719            return Err(err);
1720        }
1721
1722        let next = tokio::select! {
1723            _ = options.cancel.wait() => {
1724                let _ = send_cancel_request(client).await;
1725                return Err(CliError::Cancelled);
1726            }
1727            result = tokio::time::timeout(options.deadline.remaining(), stream.next()) => {
1728                match result {
1729                    Ok(value) => value,
1730                    Err(_) => {
1731                        let _ = send_cancel_request(client).await;
1732                        return Err(CliError::Timeout(format!(
1733                            "deadline exceeded after {}",
1734                            humantime::format_duration(options.deadline.duration())
1735                        )));
1736                    }
1737                }
1738            }
1739        };
1740
1741        let chunk = match next {
1742            Some(chunk) => chunk,
1743            None => break,
1744        };
1745
1746        let bytes = match chunk {
1747            Ok(bytes) => bytes,
1748            Err(err) => return Err(CliError::ServerConnection(format!("request failed: {err}"))),
1749        };
1750
1751        buffer.extend_from_slice(&bytes);
1752
1753        while let Some(newline) = buffer.iter().position(|&b| b == b'\n') {
1754            let line = buffer.drain(..=newline).collect::<Vec<u8>>();
1755            if let Some(item) = parse_jsonl_line(&line)? {
1756                if let Some(error) = item.error {
1757                    return Err(CliError::InvalidArgument(format!(
1758                        "Server error: {}",
1759                        error.message
1760                    )));
1761                }
1762                if item.done {
1763                    done = true;
1764                    break;
1765                }
1766                let row = item.row.ok_or_else(|| {
1767                    CliError::InvalidArgument("Invalid streaming response: missing row".into())
1768                })?;
1769                if streaming_writer.is_none() {
1770                    let columns = default_stream_columns(row.len());
1771                    let formatter = formatter
1772                        .take()
1773                        .ok_or_else(|| CliError::InvalidArgument("Missing formatter".into()))?;
1774                    let mut writer =
1775                        StreamingWriter::new(&mut *writer, formatter, columns, options.limit)
1776                            .with_quiet(options.quiet);
1777                    writer.prepare(None)?;
1778                    streaming_writer = Some(writer);
1779                }
1780                let values = row.into_iter().map(remote_value_to_value).collect();
1781                if let Some(writer) = streaming_writer.as_mut() {
1782                    match writer.write_row(Row::new(values))? {
1783                        WriteStatus::LimitReached => {
1784                            let _ = send_cancel_request(client).await;
1785                            done = true;
1786                            break;
1787                        }
1788                        WriteStatus::Continue => {}
1789                    }
1790                }
1791            }
1792        }
1793    }
1794
1795    if !done {
1796        if let Some(item) = parse_jsonl_line(&buffer)? {
1797            if let Some(error) = item.error {
1798                return Err(CliError::InvalidArgument(format!(
1799                    "Server error: {}",
1800                    error.message
1801                )));
1802            }
1803            if item.done {
1804                done = true;
1805            } else if let Some(row) = item.row {
1806                if streaming_writer.is_none() {
1807                    let columns = default_stream_columns(row.len());
1808                    let formatter = formatter
1809                        .take()
1810                        .ok_or_else(|| CliError::InvalidArgument("Missing formatter".into()))?;
1811                    let mut writer =
1812                        StreamingWriter::new(&mut *writer, formatter, columns, options.limit)
1813                            .with_quiet(options.quiet);
1814                    writer.prepare(None)?;
1815                    streaming_writer = Some(writer);
1816                }
1817                let values = row.into_iter().map(remote_value_to_value).collect();
1818                if let Some(writer) = streaming_writer.as_mut() {
1819                    match writer.write_row(Row::new(values))? {
1820                        WriteStatus::LimitReached => {
1821                            let _ = send_cancel_request(client).await;
1822                            done = true;
1823                        }
1824                        WriteStatus::Continue => {}
1825                    }
1826                }
1827            } else {
1828                return Err(CliError::InvalidArgument(
1829                    "Invalid streaming response: missing row".into(),
1830                ));
1831            }
1832        }
1833    }
1834
1835    if !done {
1836        return Err(CliError::InvalidArgument(
1837            "Invalid streaming response: unexpected end of stream".into(),
1838        ));
1839    }
1840
1841    if let Some(mut writer) = streaming_writer {
1842        return writer.finish();
1843    }
1844
1845    if let Some(formatter) = formatter.take() {
1846        let mut writer = StreamingWriter::new(&mut *writer, formatter, Vec::new(), options.limit)
1847            .with_quiet(options.quiet);
1848        writer.prepare(None)?;
1849        return writer.finish();
1850    }
1851
1852    Ok(())
1853}
1854
1855fn default_stream_columns(count: usize) -> Vec<Column> {
1856    (0..count)
1857        .map(|idx| Column::new(format!("col{}", idx + 1), DataType::Text))
1858        .collect()
1859}
1860
1861fn parse_jsonl_line(line: &[u8]) -> Result<Option<RemoteStreamItem>> {
1862    if line
1863        .iter()
1864        .all(|b| matches!(b, b' ' | b'\n' | b'\r' | b'\t'))
1865    {
1866        return Ok(None);
1867    }
1868    let text = std::str::from_utf8(line)
1869        .map_err(|err| CliError::InvalidArgument(format!("Invalid UTF-8: {err}")))?;
1870    let trimmed = text.trim();
1871    if trimmed.is_empty() {
1872        return Ok(None);
1873    }
1874    let item = serde_json::from_str::<RemoteStreamItem>(trimmed)?;
1875    Ok(Some(item))
1876}
1877
1878fn skip_whitespace(buffer: &[u8], pos: &mut usize) {
1879    while *pos < buffer.len() {
1880        match buffer[*pos] {
1881            b' ' | b'\n' | b'\r' | b'\t' => *pos += 1,
1882            _ => break,
1883        }
1884    }
1885}
1886
1887fn has_non_whitespace(buffer: &[u8]) -> bool {
1888    buffer
1889        .iter()
1890        .any(|byte| !matches!(byte, b' ' | b'\n' | b'\r' | b'\t'))
1891}
1892
1893fn json_value_to_value(value: &serde_json::Value) -> Result<Value> {
1894    match value {
1895        serde_json::Value::Null => Ok(Value::Null),
1896        serde_json::Value::Bool(value) => Ok(Value::Bool(*value)),
1897        serde_json::Value::Number(value) => {
1898            if let Some(value) = value.as_i64() {
1899                Ok(Value::Int(value))
1900            } else if let Some(value) = value.as_f64() {
1901                Ok(Value::Float(value))
1902            } else {
1903                Err(CliError::InvalidArgument(
1904                    "Invalid numeric value in streaming row".into(),
1905                ))
1906            }
1907        }
1908        serde_json::Value::String(value) => Ok(Value::Text(value.clone())),
1909        serde_json::Value::Array(values) => {
1910            let mut vector = Vec::with_capacity(values.len());
1911            for entry in values {
1912                let number = entry.as_f64().ok_or_else(|| {
1913                    CliError::InvalidArgument("Invalid vector value in streaming row".into())
1914                })?;
1915                vector.push(number as f32);
1916            }
1917            Ok(Value::Vector(vector))
1918        }
1919        serde_json::Value::Object(_) => Err(CliError::InvalidArgument(
1920            "Invalid streaming row: nested objects are not supported".into(),
1921        )),
1922    }
1923}
1924
1925/// Legacy execute function for backward compatibility with tests.
1926#[allow(dead_code)]
1927pub fn execute<W: Write>(
1928    db: &Database,
1929    cmd: SqlCommand,
1930    batch_mode: &BatchMode,
1931    writer: &mut StreamingWriter<W>,
1932) -> Result<()> {
1933    let sql = cmd.resolve_query(batch_mode)?;
1934
1935    execute_sql(db, &sql, writer)
1936}
1937
1938impl SqlCommand {
1939    /// Resolve the SQL query source (argument, file, or stdin).
1940    pub fn resolve_query(&self, batch_mode: &BatchMode) -> Result<String> {
1941        match (&self.query, &self.file) {
1942            (Some(query), None) => Ok(query.clone()),
1943            (None, Some(file)) => fs::read_to_string(file).map_err(|e| {
1944                CliError::InvalidArgument(format!("Failed to read SQL file '{}': {}", file, e))
1945            }),
1946            (None, None) if !batch_mode.is_tty => {
1947                let mut buf = String::new();
1948                io::stdin().read_to_string(&mut buf)?;
1949                Ok(buf)
1950            }
1951            (None, None) => Err(CliError::NoQueryProvided),
1952            (Some(_), Some(_)) => Err(CliError::InvalidArgument(
1953                "Cannot specify both query and file".to_string(),
1954            )),
1955        }
1956    }
1957}
1958
1959/// Execute SQL and write results.
1960fn execute_sql<W: Write>(db: &Database, sql: &str, writer: &mut StreamingWriter<W>) -> Result<()> {
1961    use alopex_embedded::SqlResult;
1962
1963    let result = db.execute_sql(sql)?;
1964
1965    match result {
1966        SqlResult::Success => {
1967            // DDL success - output simple status
1968            writer.prepare(Some(1))?;
1969            let row = Row::new(vec![
1970                Value::Text("OK".to_string()),
1971                Value::Text("Operation completed successfully".to_string()),
1972            ]);
1973            writer.write_row(row)?;
1974            writer.finish()?;
1975        }
1976        SqlResult::RowsAffected(count) => {
1977            // DML success - output affected rows count
1978            writer.prepare(Some(1))?;
1979            let row = Row::new(vec![
1980                Value::Text("OK".to_string()),
1981                Value::Text(format!("{} row(s) affected", count)),
1982            ]);
1983            writer.write_row(row)?;
1984            writer.finish()?;
1985        }
1986        SqlResult::Query(query_result) => {
1987            // SELECT result - output rows
1988            let row_count = query_result.rows.len();
1989            writer.prepare(Some(row_count))?;
1990
1991            for sql_row in query_result.rows {
1992                let values: Vec<Value> = sql_row.into_iter().map(sql_value_to_value).collect();
1993                let row = Row::new(values);
1994
1995                match writer.write_row(row)? {
1996                    WriteStatus::LimitReached => break,
1997                    WriteStatus::Continue => {}
1998                }
1999            }
2000
2001            writer.finish()?;
2002        }
2003    }
2004
2005    Ok(())
2006}
2007
2008/// Execute SQL locally, emitting one result block per statement.
2009///
2010/// - A single query statement uses the streaming path (FR-7).
2011/// - Everything else (DDL/DML and multi-statement input) executes all
2012///   statements in one transaction via [`Database::execute_sql_multi`] and
2013///   emits each statement's result in order.
2014/// - For `--output json` the output is always an array of per-statement
2015///   result sets (a single statement yields a 1-element array).
2016///
2017/// FR-7 Compliance: Uses the SQL parser to detect query statements instead of a heuristic.
2018/// This properly handles:
2019/// - VALUES query bodies
2020/// - WITH clauses (CTEs)
2021/// - Leading comments
2022/// - Complex query structures
2023fn execute_sql_with_formatter<W: Write>(
2024    db: &Database,
2025    sql: &str,
2026    writer: &mut W,
2027    output: &mut SqlOutput<'_>,
2028    options: &SqlExecutionOptions<'_>,
2029) -> Result<()> {
2030    use alopex_sql::{AlopexDialect, Parser};
2031
2032    // FR-7: Parse the statement instead of relying on a SELECT-prefix heuristic.
2033    // This handles VALUES, WITH clauses, leading comments, and complex query structures.
2034    let dialect = AlopexDialect;
2035    let stmts = Parser::parse_sql(&dialect, sql).map_err(|e| CliError::Parse(format!("{}", e)))?;
2036
2037    let is_single_query = stmts.len() == 1
2038        && stmts
2039            .first()
2040            .is_some_and(|statement| statement.kind.is_query());
2041
2042    if is_single_query {
2043        // A single query statement uses the streaming path (FR-7).
2044        let formatter = output.create_block_formatter();
2045        if !output.statement_array() {
2046            return execute_sql_select_streaming(db, sql, writer, formatter, options, false);
2047        }
2048        let mut result_set_writer = CountingWriter::new(&mut *writer);
2049        let result =
2050            execute_sql_select_streaming(db, sql, &mut result_set_writer, formatter, options, true);
2051        let output_started = result_set_writer.bytes_written() > 0;
2052        match result {
2053            Ok(()) => {
2054                // Close the statement-result array opened by the streaming callback.
2055                writeln!(writer, "]")?;
2056                Ok(())
2057            }
2058            Err(err) => {
2059                // Keep stdout well-formed JSON on mid-stream errors: the
2060                // callback opens the statement array and the result-set array
2061                // together before writing rows (always complete objects), so
2062                // closing both leaves valid JSON alongside the non-zero exit
2063                // code. If nothing was written (e.g. planning failed), there
2064                // is nothing to close.
2065                if output_started {
2066                    let _ = writeln!(writer, "]");
2067                    let _ = writeln!(writer, "]");
2068                }
2069                Err(err)
2070            }
2071        }
2072    } else {
2073        // DDL/DML and multi-statement input: emit one result block per statement.
2074        execute_sql_statements(db, sql, writer, output, options)
2075    }
2076}
2077
2078/// Execute SELECT query with streaming callback (FR-7).
2079///
2080/// This function uses `execute_sql_with_rows` for true streaming output.
2081/// The callback receives rows one at a time from the iterator, and the
2082/// transaction is kept alive during streaming.
2083///
2084/// When `open_statement_array` is true (JSON output), the callback writes the
2085/// opening `[` of the statement-result array before streaming; the caller is
2086/// responsible for writing the closing `]` after this function returns.
2087fn execute_sql_select_streaming<W: Write>(
2088    db: &Database,
2089    sql: &str,
2090    writer: &mut W,
2091    formatter: Box<dyn Formatter>,
2092    options: &SqlExecutionOptions<'_>,
2093    open_statement_array: bool,
2094) -> Result<()> {
2095    use alopex_embedded::StreamingQueryResult;
2096    use std::sync::atomic::{AtomicBool, Ordering};
2097    use std::sync::Arc;
2098
2099    // Helper to convert CliError to alopex_embedded::Error for callback
2100    fn cli_err_to_embedded(e: crate::error::CliError) -> alopex_embedded::Error {
2101        alopex_embedded::Error::Sql(alopex_sql::SqlError::Execution {
2102            message: e.to_string(),
2103            code: "ALOPEX-C001",
2104        })
2105    }
2106
2107    let cancelled = Arc::new(AtomicBool::new(false));
2108    let timed_out = Arc::new(AtomicBool::new(false));
2109    let cancel_flag = cancelled.clone();
2110    let timeout_flag = timed_out.clone();
2111
2112    let result = db.execute_sql_with_rows(sql, |mut rows| {
2113        // FR-7: SELECT result - stream rows directly from iterator while transaction is alive
2114        let columns = columns_from_streaming_rows(&rows);
2115        // The prefix is written here (not by the caller) so that planning
2116        // errors do not leave a dangling `[` on stdout.
2117        if open_statement_array {
2118            writeln!(writer, "[").map_err(|e| cli_err_to_embedded(CliError::Io(e)))?;
2119        }
2120        let mut streaming_writer = StreamingWriter::new(writer, formatter, columns, options.limit)
2121            .with_quiet(options.quiet);
2122
2123        // FR-7: Use None for row count hint to support true streaming output
2124        streaming_writer
2125            .prepare(None)
2126            .map_err(cli_err_to_embedded)?;
2127
2128        if let Err(err) = options.deadline.check() {
2129            timeout_flag.store(true, Ordering::SeqCst);
2130            return Err(cli_err_to_embedded(err));
2131        }
2132
2133        // Consume iterator row by row for true streaming.
2134        // Errors from `next_row` must propagate: swallowing them here would
2135        // silently truncate the result set (GitHub issue #23).
2136        while let Some(sql_row) = rows.next_row()? {
2137            if options.cancel.is_cancelled() {
2138                cancel_flag.store(true, Ordering::SeqCst);
2139                return Err(cli_err_to_embedded(CliError::Cancelled));
2140            }
2141            if let Err(err) = options.deadline.check() {
2142                timeout_flag.store(true, Ordering::SeqCst);
2143                return Err(cli_err_to_embedded(err));
2144            }
2145            let values: Vec<Value> = sql_row.into_iter().map(sql_value_to_value).collect();
2146            let row = Row::new(values);
2147
2148            match streaming_writer
2149                .write_row(row)
2150                .map_err(cli_err_to_embedded)?
2151            {
2152                WriteStatus::LimitReached => break,
2153                WriteStatus::Continue => {}
2154            }
2155        }
2156
2157        streaming_writer.finish().map_err(cli_err_to_embedded)?;
2158        Ok(())
2159    });
2160
2161    let result = match result {
2162        Ok(value) => value,
2163        Err(err) => {
2164            if cancelled.load(Ordering::SeqCst) {
2165                return Err(CliError::Cancelled);
2166            }
2167            if timed_out.load(Ordering::SeqCst) {
2168                return Err(CliError::Timeout(format!(
2169                    "deadline exceeded after {}",
2170                    humantime::format_duration(options.deadline.duration())
2171                )));
2172            }
2173            return Err(CliError::Database(err));
2174        }
2175    };
2176
2177    match result {
2178        StreamingQueryResult::QueryProcessed(()) => Ok(()),
2179        StreamingQueryResult::Success | StreamingQueryResult::RowsAffected(_) => {
2180            // Unexpected: SELECT should not return these
2181            Ok(())
2182        }
2183    }
2184}
2185
2186#[derive(serde::Serialize)]
2187struct RemoteSqlRequest {
2188    sql: String,
2189    #[serde(default)]
2190    streaming: bool,
2191    #[serde(skip_serializing_if = "Option::is_none")]
2192    fetch_size: Option<usize>,
2193    #[serde(skip_serializing_if = "Option::is_none")]
2194    max_rows: Option<usize>,
2195}
2196
2197#[derive(serde::Serialize)]
2198struct RemoteDistributedReadRequest {
2199    sql: String,
2200    read_mode: &'static str,
2201}
2202
2203#[derive(serde::Deserialize)]
2204struct RemoteColumnInfo {
2205    name: String,
2206    data_type: String,
2207}
2208
2209#[derive(serde::Deserialize, Default)]
2210struct RemoteSqlResult {
2211    #[serde(default)]
2212    columns: Vec<RemoteColumnInfo>,
2213    #[serde(default)]
2214    rows: Vec<Vec<alopex_sql::storage::SqlValue>>,
2215    #[serde(default)]
2216    affected_rows: Option<u64>,
2217}
2218
2219#[derive(serde::Deserialize)]
2220struct RemoteSqlResponse {
2221    #[serde(default)]
2222    results: Vec<RemoteSqlResult>,
2223    #[serde(flatten)]
2224    legacy: RemoteSqlResult,
2225}
2226
2227impl RemoteSqlResponse {
2228    fn into_results(self) -> Vec<RemoteSqlResult> {
2229        if self.results.is_empty() {
2230            vec![self.legacy]
2231        } else {
2232            self.results
2233        }
2234    }
2235}
2236
2237#[derive(serde::Deserialize)]
2238struct RemoteStreamItem {
2239    row: Option<Vec<alopex_sql::storage::SqlValue>>,
2240    error: Option<RemoteStreamError>,
2241    #[serde(default)]
2242    done: bool,
2243}
2244
2245#[derive(serde::Deserialize)]
2246struct RemoteStreamError {
2247    message: String,
2248}
2249
2250fn map_client_error(err: ClientError) -> CliError {
2251    match err {
2252        ClientError::Request { source, .. } => {
2253            CliError::ServerConnection(format!("request failed: {source}"))
2254        }
2255        ClientError::InvalidUrl(message) => CliError::InvalidArgument(message),
2256        ClientError::Build(message) => CliError::InvalidArgument(message),
2257        ClientError::Auth(err) => CliError::InvalidArgument(err.to_string()),
2258        ClientError::HttpStatus { status, body } => {
2259            CliError::InvalidArgument(format!("Server error: HTTP {} - {}", status.as_u16(), body))
2260        }
2261    }
2262}
2263
2264async fn send_cancel_request(client: &HttpClient) -> Result<()> {
2265    #[derive(serde::Serialize)]
2266    struct CancelRequest {}
2267
2268    let request = CancelRequest {};
2269    let _: serde_json::Value = client
2270        .post_json("api/sql/cancel", &request)
2271        .await
2272        .map_err(map_client_error)?;
2273    Ok(())
2274}
2275
2276fn merge_limit(limit: Option<usize>, max_rows: Option<usize>) -> Option<usize> {
2277    match (limit, max_rows) {
2278        (Some(a), Some(b)) => Some(a.min(b)),
2279        (Some(value), None) | (None, Some(value)) => Some(value),
2280        (None, None) => None,
2281    }
2282}
2283
2284/// Execute DDL/DML/multi-statement SQL and emit one result block per statement.
2285///
2286/// All statements run in a single auto-commit transaction
2287/// ([`Database::execute_sql_multi`]); if any statement fails, the whole batch
2288/// is rolled back and an error is returned (non-zero exit code).
2289///
2290/// Output contract:
2291/// - `--output json`: an array of per-statement result sets. DDL/DML
2292///   statements contribute a status result set (`status`/`message` columns);
2293///   `--quiet` omits status result sets.
2294/// - Other formats: one result block per statement, in statement order
2295///   (status blocks are suppressed by `--quiet`).
2296fn execute_sql_statements<W: Write>(
2297    db: &Database,
2298    sql: &str,
2299    writer: &mut W,
2300    output: &mut SqlOutput<'_>,
2301    options: &SqlExecutionOptions<'_>,
2302) -> Result<()> {
2303    use alopex_sql::ExecutionResult;
2304
2305    options.deadline.check()?;
2306    let results = db.execute_sql_multi(sql)?;
2307    options.deadline.check()?;
2308
2309    let json_array = output.statement_array();
2310    if json_array {
2311        writeln!(writer, "[")?;
2312    }
2313    let mut first = true;
2314    for result in results {
2315        let is_status = !matches!(result, ExecutionResult::Query(_));
2316        if is_status && options.quiet {
2317            continue;
2318        }
2319        if json_array && !first {
2320            writeln!(writer, ",")?;
2321        }
2322        first = false;
2323        emit_execution_result(result, writer, output, options)?;
2324    }
2325    if json_array {
2326        writeln!(writer, "]")?;
2327    }
2328    Ok(())
2329}
2330
2331/// Emit a single statement's execution result as one result block.
2332fn emit_execution_result<W: Write>(
2333    result: alopex_sql::ExecutionResult,
2334    writer: &mut W,
2335    output: &mut SqlOutput<'_>,
2336    options: &SqlExecutionOptions<'_>,
2337) -> Result<()> {
2338    use alopex_sql::ExecutionResult;
2339
2340    match result {
2341        ExecutionResult::Success => emit_status_result(
2342            writer,
2343            output,
2344            options,
2345            "Operation completed successfully".to_string(),
2346        ),
2347        ExecutionResult::RowsAffected(count) => {
2348            emit_status_result(writer, output, options, format!("{count} row(s) affected"))
2349        }
2350        ExecutionResult::Query(query_result) => {
2351            let columns = columns_from_query_result(&query_result);
2352            let mut streaming_writer = StreamingWriter::new(
2353                &mut *writer,
2354                output.create_block_formatter(),
2355                columns,
2356                options.limit,
2357            )
2358            .with_quiet(options.quiet);
2359            streaming_writer.prepare(Some(query_result.rows.len()))?;
2360            for sql_row in query_result.rows {
2361                let values: Vec<Value> = sql_row.into_iter().map(sql_value_to_value).collect();
2362                match streaming_writer.write_row(Row::new(values))? {
2363                    WriteStatus::LimitReached => break,
2364                    WriteStatus::Continue => {}
2365                }
2366            }
2367            streaming_writer.finish()
2368        }
2369    }
2370}
2371
2372/// Emit an `OK` status result set (`status`/`message` columns).
2373fn emit_status_result<W: Write>(
2374    writer: &mut W,
2375    output: &mut SqlOutput<'_>,
2376    options: &SqlExecutionOptions<'_>,
2377    message: String,
2378) -> Result<()> {
2379    let columns = sql_status_columns();
2380    let mut streaming_writer = StreamingWriter::new(
2381        &mut *writer,
2382        output.create_block_formatter(),
2383        columns,
2384        options.limit,
2385    )
2386    .with_quiet(options.quiet);
2387    streaming_writer.prepare(Some(1))?;
2388    streaming_writer.write_row(Row::new(vec![
2389        Value::Text("OK".to_string()),
2390        Value::Text(message),
2391    ]))?;
2392    streaming_writer.finish()
2393}
2394
2395/// Convert alopex_sql::SqlValue to our Value type.
2396fn sql_value_to_value(sql_value: alopex_sql::SqlValue) -> Value {
2397    use alopex_sql::SqlValue;
2398
2399    match sql_value {
2400        SqlValue::Null => Value::Null,
2401        SqlValue::Integer(i) => Value::Int(i as i64),
2402        SqlValue::BigInt(i) => Value::Int(i),
2403        SqlValue::Float(f) => Value::Float(f as f64),
2404        SqlValue::Double(f) => Value::Float(f),
2405        SqlValue::Text(s) => Value::Text(s),
2406        SqlValue::Blob(b) => Value::Bytes(b),
2407        SqlValue::Boolean(b) => Value::Bool(b),
2408        SqlValue::Timestamp(ts) => {
2409            // Format timestamp as ISO 8601 string
2410            Value::Text(format!("{}", ts))
2411        }
2412        SqlValue::Vector(v) => Value::Vector(v),
2413    }
2414}
2415
2416fn remote_value_to_value(sql_value: alopex_sql::storage::SqlValue) -> Value {
2417    use alopex_sql::storage::SqlValue;
2418
2419    match sql_value {
2420        SqlValue::Null => Value::Null,
2421        SqlValue::Integer(i) => Value::Int(i as i64),
2422        SqlValue::BigInt(i) => Value::Int(i),
2423        SqlValue::Float(f) => Value::Float(f as f64),
2424        SqlValue::Double(f) => Value::Float(f),
2425        SqlValue::Text(s) => Value::Text(s),
2426        SqlValue::Blob(b) => Value::Bytes(b),
2427        SqlValue::Boolean(b) => Value::Bool(b),
2428        SqlValue::Timestamp(ts) => Value::Text(ts.to_string()),
2429        SqlValue::Vector(v) => Value::Vector(v),
2430    }
2431}
2432
2433fn data_type_from_string(value: &str) -> DataType {
2434    let upper = value.to_ascii_uppercase();
2435    if upper.starts_with("INT") || upper.starts_with("BIGINT") {
2436        DataType::Int
2437    } else if upper.starts_with("FLOAT") || upper.starts_with("DOUBLE") {
2438        DataType::Float
2439    } else if upper.starts_with("BLOB") {
2440        DataType::Bytes
2441    } else if upper.starts_with("BOOLEAN") {
2442        DataType::Bool
2443    } else if upper.starts_with("VECTOR") {
2444        DataType::Vector
2445    } else {
2446        DataType::Text
2447    }
2448}
2449
2450/// Convert alopex_sql::executor::ColumnInfo to our Column type.
2451fn sql_column_to_column(col: &alopex_sql::executor::ColumnInfo) -> Column {
2452    use alopex_sql::planner::ResolvedType;
2453
2454    let data_type = match &col.data_type {
2455        ResolvedType::Integer | ResolvedType::BigInt => DataType::Int,
2456        ResolvedType::Float | ResolvedType::Double => DataType::Float,
2457        ResolvedType::Text => DataType::Text,
2458        ResolvedType::Blob => DataType::Bytes,
2459        ResolvedType::Boolean => DataType::Bool,
2460        ResolvedType::Timestamp => DataType::Text, // Display as text
2461        ResolvedType::Vector { .. } => DataType::Vector,
2462        ResolvedType::Null => DataType::Text, // Fallback
2463    };
2464
2465    Column::new(&col.name, data_type)
2466}
2467
2468/// Create columns from SQL query result.
2469fn columns_from_query_result(query_result: &alopex_sql::executor::QueryResult) -> Vec<Column> {
2470    query_result
2471        .columns
2472        .iter()
2473        .map(sql_column_to_column)
2474        .collect()
2475}
2476
2477/// Create columns from streaming query result iterator (FR-7).
2478#[allow(dead_code)] // Kept for potential future use with old streaming API
2479fn columns_from_streaming_result(
2480    query_iter: &alopex_embedded::QueryRowIterator<'_>,
2481) -> Vec<Column> {
2482    query_iter
2483        .columns()
2484        .iter()
2485        .map(sql_column_to_column)
2486        .collect()
2487}
2488
2489/// Create columns from callback-based streaming rows (FR-7).
2490fn columns_from_streaming_rows(rows: &alopex_embedded::StreamingRows<'_>) -> Vec<Column> {
2491    rows.columns().iter().map(sql_column_to_column).collect()
2492}
2493
2494/// Create columns for status output.
2495pub fn sql_status_columns() -> Vec<Column> {
2496    vec![
2497        Column::new("status", DataType::Text),
2498        Column::new("message", DataType::Text),
2499    ]
2500}
2501
2502#[cfg(test)]
2503mod tests {
2504    use super::*;
2505    use crate::batch::BatchModeSource;
2506    use crate::output::jsonl::JsonlFormatter;
2507
2508    fn create_test_db() -> Database {
2509        Database::open_in_memory().unwrap()
2510    }
2511
2512    fn create_status_writer(output: &mut Vec<u8>) -> StreamingWriter<&mut Vec<u8>> {
2513        let formatter = Box::new(JsonlFormatter::new());
2514        let columns = sql_status_columns();
2515        StreamingWriter::new(output, formatter, columns, None)
2516    }
2517
2518    fn create_query_writer(
2519        output: &mut Vec<u8>,
2520        columns: Vec<Column>,
2521    ) -> StreamingWriter<&mut Vec<u8>> {
2522        let formatter = Box::new(JsonlFormatter::new());
2523        StreamingWriter::new(output, formatter, columns, None)
2524    }
2525
2526    fn default_batch_mode() -> BatchMode {
2527        BatchMode {
2528            is_batch: false,
2529            is_tty: true,
2530            source: BatchModeSource::Default,
2531        }
2532    }
2533
2534    #[test]
2535    fn test_create_table() {
2536        let db = create_test_db();
2537
2538        let mut output = Vec::new();
2539        {
2540            let mut writer = create_status_writer(&mut output);
2541            execute_sql(
2542                &db,
2543                "CREATE TABLE users (id INTEGER PRIMARY KEY, name TEXT);",
2544                &mut writer,
2545            )
2546            .unwrap();
2547        }
2548
2549        let result = String::from_utf8(output).unwrap();
2550        assert!(result.contains("OK"));
2551    }
2552
2553    #[test]
2554    fn test_insert_and_select() {
2555        let db = create_test_db();
2556
2557        // Create table
2558        {
2559            let mut output = Vec::new();
2560            let mut writer = create_status_writer(&mut output);
2561            execute_sql(
2562                &db,
2563                "CREATE TABLE users (id INTEGER PRIMARY KEY, name TEXT);",
2564                &mut writer,
2565            )
2566            .unwrap();
2567        }
2568
2569        // Insert
2570        {
2571            let mut output = Vec::new();
2572            let mut writer = create_status_writer(&mut output);
2573            execute_sql(
2574                &db,
2575                "INSERT INTO users (id, name) VALUES (1, 'Alice');",
2576                &mut writer,
2577            )
2578            .unwrap();
2579            let result = String::from_utf8(output).unwrap();
2580            assert!(result.contains("row(s) affected"));
2581        }
2582
2583        // Select - we need columns from the query result
2584        {
2585            let mut output = Vec::new();
2586            let columns = vec![
2587                Column::new("id", DataType::Int),
2588                Column::new("name", DataType::Text),
2589            ];
2590            let mut writer = create_query_writer(&mut output, columns);
2591            execute_sql(&db, "SELECT id, name FROM users;", &mut writer).unwrap();
2592            let result = String::from_utf8(output).unwrap();
2593            assert!(result.contains("Alice"));
2594        }
2595    }
2596
2597    #[test]
2598    fn test_syntax_error() {
2599        let db = create_test_db();
2600
2601        let mut output = Vec::new();
2602        let mut writer = create_status_writer(&mut output);
2603        let result = execute_sql(&db, "CREATE TABEL invalid_syntax;", &mut writer);
2604        assert!(result.is_err());
2605    }
2606
2607    #[test]
2608    fn test_multiple_statements() {
2609        let db = create_test_db();
2610
2611        let mut output = Vec::new();
2612        {
2613            let mut writer = create_status_writer(&mut output);
2614            execute_sql(
2615                &db,
2616                "CREATE TABLE t (id INTEGER PRIMARY KEY); INSERT INTO t (id) VALUES (1);",
2617                &mut writer,
2618            )
2619            .unwrap();
2620        }
2621
2622        // Verify the table exists and has data
2623        {
2624            let mut output = Vec::new();
2625            let columns = vec![Column::new("id", DataType::Int)];
2626            let mut writer = create_query_writer(&mut output, columns);
2627            execute_sql(&db, "SELECT id FROM t;", &mut writer).unwrap();
2628            let result = String::from_utf8(output).unwrap();
2629            assert!(result.contains("1"));
2630        }
2631    }
2632
2633    /// Regression test for issue #23: errors surfaced by `next_row` during a
2634    /// streaming SELECT must propagate as `Err`, not be swallowed into a
2635    /// silently-empty result set with exit code 0 (the old
2636    /// `while let Ok(Some(...))` pattern did exactly that).
2637    #[test]
2638    fn streaming_select_row_error_propagates() {
2639        let db = create_test_db();
2640        db.execute_sql("CREATE TABLE t (id INTEGER PRIMARY KEY);")
2641            .unwrap();
2642        db.execute_sql("INSERT INTO t (id) VALUES (1);").unwrap();
2643
2644        let deadline = Deadline::new(parse_deadline(None).unwrap());
2645        let cancel = CancelSignal::new();
2646        let options = SqlExecutionOptions {
2647            limit: None,
2648            quiet: false,
2649            cancel: &cancel,
2650            deadline: &deadline,
2651            admin_launcher: None,
2652        };
2653
2654        let mut output = Vec::new();
2655        let formatter = Box::new(JsonlFormatter::new());
2656        // `id / 0` fails during row evaluation (division by zero), i.e. inside
2657        // the `next_row` loop of `execute_sql_select_streaming`.
2658        let result = execute_sql_select_streaming(
2659            &db,
2660            "SELECT id / 0 AS x FROM t;",
2661            &mut output,
2662            formatter,
2663            &options,
2664            false,
2665        );
2666        assert!(
2667            result.is_err(),
2668            "row evaluation errors must propagate instead of yielding an empty result"
2669        );
2670    }
2671
2672    #[test]
2673    fn test_sql_value_conversion() {
2674        use alopex_sql::SqlValue;
2675
2676        assert!(matches!(sql_value_to_value(SqlValue::Null), Value::Null));
2677        assert!(matches!(
2678            sql_value_to_value(SqlValue::Integer(42)),
2679            Value::Int(42)
2680        ));
2681        assert!(matches!(
2682            sql_value_to_value(SqlValue::BigInt(100)),
2683            Value::Int(100)
2684        ));
2685        assert!(matches!(
2686            sql_value_to_value(SqlValue::Boolean(true)),
2687            Value::Bool(true)
2688        ));
2689        assert!(
2690            matches!(sql_value_to_value(SqlValue::Text("hello".to_string())), Value::Text(s) if s == "hello")
2691        );
2692    }
2693
2694    #[test]
2695    fn resolve_query_from_argument() {
2696        let cmd = SqlCommand {
2697            query: Some("SELECT 1".to_string()),
2698            file: None,
2699            fetch_size: None,
2700            max_rows: None,
2701            deadline: None,
2702            read_mode: None,
2703            routing_report: None,
2704            tui: false,
2705        };
2706
2707        let sql = cmd.resolve_query(&default_batch_mode()).unwrap();
2708        assert_eq!(sql, "SELECT 1");
2709    }
2710
2711    #[test]
2712    fn resolve_query_from_file() {
2713        let mut file = tempfile::NamedTempFile::new().unwrap();
2714        writeln!(file, "SELECT * FROM users").unwrap();
2715
2716        let cmd = SqlCommand {
2717            query: None,
2718            file: Some(file.path().display().to_string()),
2719            fetch_size: None,
2720            max_rows: None,
2721            deadline: None,
2722            read_mode: None,
2723            routing_report: None,
2724            tui: false,
2725        };
2726
2727        let sql = cmd.resolve_query(&default_batch_mode()).unwrap();
2728        assert_eq!(sql, "SELECT * FROM users\n");
2729    }
2730
2731    #[test]
2732    fn resolve_query_returns_no_query_error() {
2733        let cmd = SqlCommand {
2734            query: None,
2735            file: None,
2736            fetch_size: None,
2737            max_rows: None,
2738            deadline: None,
2739            read_mode: None,
2740            routing_report: None,
2741            tui: false,
2742        };
2743
2744        let err = cmd.resolve_query(&default_batch_mode()).unwrap_err();
2745        assert!(matches!(err, CliError::NoQueryProvided));
2746    }
2747
2748    #[test]
2749    fn resolve_query_rejects_query_and_file() {
2750        let cmd = SqlCommand {
2751            query: Some("SELECT 1".to_string()),
2752            file: Some("query.sql".into()),
2753            fetch_size: None,
2754            max_rows: None,
2755            deadline: None,
2756            read_mode: None,
2757            routing_report: None,
2758            tui: false,
2759        };
2760
2761        let err = cmd.resolve_query(&default_batch_mode()).unwrap_err();
2762        assert!(matches!(
2763            err,
2764            CliError::InvalidArgument(msg) if msg == "Cannot specify both query and file"
2765        ));
2766    }
2767}