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