datafusion-ducklake-cli 0.2.0

Interactive and scriptable SQL shell for datafusion-ducklake-provider.
Documentation
use std::fs;
use std::io::{self, IsTerminal, Read, Write};
use std::path::Path;
use std::time::Instant;

use datafusion_ducklake_provider::{DuckLakeSessionContext, register_ducklake_object_store};

use crate::args::Cli;
use crate::error::CliError;
use crate::object_store::object_store_configs;
use crate::output::format_dataframe;
use crate::repl::run_repl;
use crate::split::split_script;

#[derive(Clone, Copy, Debug)]
pub(crate) struct RunOptions {
    pub quiet: bool,
    pub timing: bool,
    pub continue_on_error: bool,
}

pub async fn run(cli: Cli) -> Result<(), CliError> {
    let ctx = DuckLakeSessionContext::new();
    for config in object_store_configs(&cli)? {
        register_ducklake_object_store(ctx.inner(), &config)?;
    }

    let options = RunOptions {
        quiet: cli.quiet,
        timing: cli.timing,
        continue_on_error: cli.continue_on_error,
    };
    let init_options = RunOptions {
        continue_on_error: false,
        ..options
    };

    for init in &cli.init {
        run_file(&ctx, init, init_options).await?;
    }

    if !cli.execute.is_empty() {
        for sql in &cli.execute {
            run_script_source(&ctx, "command line", sql, options).await?;
        }
        return Ok(());
    }

    if let Some(file) = &cli.file {
        run_file(&ctx, file, options).await?;
        return Ok(());
    }

    if !io::stdin().is_terminal() {
        let mut sql = String::new();
        io::stdin().read_to_string(&mut sql)?;
        run_script_source(&ctx, "stdin", &sql, options).await?;
        return Ok(());
    }

    run_repl(&ctx, &cli, options).await
}

async fn run_file(
    ctx: &DuckLakeSessionContext,
    path: &Path,
    options: RunOptions,
) -> Result<(), CliError> {
    let sql = fs::read_to_string(path)?;
    run_script_source(ctx, &path.display().to_string(), &sql, options).await
}

pub(crate) async fn run_script_source(
    ctx: &DuckLakeSessionContext,
    source: &str,
    sql: &str,
    options: RunOptions,
) -> Result<(), CliError> {
    for statement in split_script(sql) {
        if let Err(err) = run_statement(ctx, &statement, options).await {
            let message = format!("{source}: {err}");
            if !options.continue_on_error {
                return Err(CliError::Message(message));
            }
            eprintln!("{message}");
        }
    }
    Ok(())
}

pub(crate) async fn run_statement(
    ctx: &DuckLakeSessionContext,
    statement: &str,
    options: RunOptions,
) -> Result<(), CliError> {
    let started = Instant::now();
    let dataframe = ctx.sql(statement).await?;
    if let Some(rendered) = format_dataframe(dataframe, options.quiet).await? {
        io::stdout().write_all(rendered.as_bytes())?;
    }
    if options.timing && !options.quiet {
        eprintln!("Time: {:.3}s", started.elapsed().as_secs_f64());
    }
    Ok(())
}