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