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, 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#[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#[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 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#[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 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 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 if result_set_started {
1295 let _ = writeln!(writer, "]");
1296 }
1297 let _ = writeln!(writer, "]");
1298 Err(err)
1299 }
1300 }
1301}
1302
1303async 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#[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 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
1960fn 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 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 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 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
2009fn 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 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 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 writeln!(writer, "]")?;
2057 Ok(())
2058 }
2059 Err(err) => {
2060 if output_started {
2067 let _ = writeln!(writer, "]");
2068 let _ = writeln!(writer, "]");
2069 }
2070 Err(err)
2071 }
2072 }
2073 } else {
2074 execute_sql_statements(db, sql, writer, output, options)
2076 }
2077}
2078
2079fn 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 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 let columns = columns_from_streaming_rows(&rows);
2116 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 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 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 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
2285fn 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
2332fn 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
2373fn 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
2396fn 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 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
2451fn 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, ResolvedType::Vector { .. } => DataType::Vector,
2463 ResolvedType::Null => DataType::Text, };
2465
2466 Column::new(&col.name, data_type)
2467}
2468
2469fn 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#[allow(dead_code)] fn 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
2490fn columns_from_streaming_rows(rows: &alopex_embedded::StreamingRows<'_>) -> Vec<Column> {
2492 rows.columns().iter().map(sql_column_to_column).collect()
2493}
2494
2495pub 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 {
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 {
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 {
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 {
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 #[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 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}