defmodule QuackDB do @moduledoc """ Remote DuckDB Quack protocol client. The public API is backed by `DBConnection` so it can grow into an Ecto adapter without changing the lower-level protocol codec. """ alias QuackDB.Query alias QuackDB.Stream @type start_option :: {:uri, String.t()} | {:token, String.t()} | {:name, GenServer.name()} @type insert_row :: map() | Keyword.t() @spec start_link([start_option]) :: GenServer.on_start() def start_link(options) do QuackDB.DBConnection.start_link(options) end @spec child_spec([start_option]) :: Supervisor.child_spec() def child_spec(options) do QuackDB.DBConnection.child_spec(options) end @spec insert_rows(DBConnection.conn(), String.t() | atom(), [insert_row()], Keyword.t()) :: {:ok, QuackDB.Result.t()} | {:error, Exception.t()} def insert_rows(connection, table, rows, options \\ []) when is_list(rows) do query = %Query{statement: "APPEND #{table}", operation: {:insert_rows, table, rows, options}} case DBConnection.prepare_execute(connection, query, [], options) do {:ok, _query, result} -> {:ok, result} {:error, _error} = error -> error end end @spec insert_rows!(DBConnection.conn(), String.t() | atom(), [insert_row()], Keyword.t()) :: QuackDB.Result.t() def insert_rows!(connection, table, rows, options \\ []) do case insert_rows(connection, table, rows, options) do {:ok, result} -> result {:error, error} -> raise error end end @spec query(DBConnection.conn(), iodata(), [term()], Keyword.t()) :: {:ok, QuackDB.Result.t()} | {:error, Exception.t()} def query(connection, statement, params \\ [], options \\ []) do query = %Query{statement: statement} case DBConnection.prepare_execute(connection, query, params, options) do {:ok, _query, result} -> {:ok, result} {:error, _error} = error -> error end end @spec query!(DBConnection.conn(), iodata(), [term()], Keyword.t()) :: QuackDB.Result.t() def query!(connection, statement, params \\ [], options \\ []) do case query(connection, statement, params, options) do {:ok, result} -> result {:error, error} -> raise error end end @doc """ Runs a query and returns its result as a column-oriented map. Duplicate column names are disambiguated with suffixes such as `_2` and `_3`. Prefer `columnar/4` when you also need column order and result metadata. """ @spec columns(DBConnection.conn(), iodata(), [term()], Keyword.t()) :: {:ok, %{String.t() => [term()]}} | {:error, Exception.t()} def columns(connection, statement, params \\ [], options \\ []) do case query(connection, statement, params, options) do {:ok, result} -> {:ok, QuackDB.Result.to_columns(result)} {:error, _error} = error -> error end end @spec columns!(DBConnection.conn(), iodata(), [term()], Keyword.t()) :: %{ String.t() => [term()] } def columns!(connection, statement, params \\ [], options \\ []) do case columns(connection, statement, params, options) do {:ok, columns} -> columns {:error, error} -> raise error end end @doc """ Runs a query and returns a `QuackDB.Columns` struct. This preserves column order, original names, row count, and result metadata in addition to the column vectors. """ @spec columnar(DBConnection.conn(), iodata(), [term()], Keyword.t()) :: {:ok, QuackDB.Columns.t()} | {:error, Exception.t()} def columnar(connection, statement, params \\ [], options \\ []) do case query(connection, statement, params, options) do {:ok, result} -> {:ok, QuackDB.Result.to_columnar(result)} {:error, _error} = error -> error end end @spec columnar!(DBConnection.conn(), iodata(), [term()], Keyword.t()) :: QuackDB.Columns.t() def columnar!(connection, statement, params \\ [], options \\ []) do case columnar(connection, statement, params, options) do {:ok, columns} -> columns {:error, error} -> raise error end end @doc """ Streams query results as column-oriented batches. Each item is a map from disambiguated column names to the values in that fetch batch. This keeps large analytical results vector-shaped without materializing the whole result set. Prefer `columnar_batches/4` when you also need batch metadata. """ @spec column_batches(DBConnection.conn(), iodata(), [term()], Keyword.t()) :: Enumerable.t() def column_batches(connection, statement, params \\ [], options \\ []) do connection |> columnar_batches(statement, params, options) |> Elixir.Stream.map(& &1.columns) |> Elixir.Stream.reject(&(&1 == %{})) end @doc """ Streams query results as `QuackDB.Columns` batches. """ @spec columnar_batches(DBConnection.conn(), iodata(), [term()], Keyword.t()) :: Enumerable.t() def columnar_batches(connection, statement, params \\ [], options \\ []) do connection |> stream(statement, params, options) |> Elixir.Stream.map(&QuackDB.Result.to_columnar/1) |> Elixir.Stream.reject(&(&1.names == [])) end @spec ping(DBConnection.conn(), Keyword.t()) :: :ok | {:error, Exception.t()} def ping(connection, options \\ []) do case query(connection, "SELECT 1", [], options) do {:ok, _result} -> :ok {:error, _error} = error -> error end end @spec prepare(DBConnection.conn(), iodata(), Keyword.t()) :: {:ok, Query.t()} | {:error, Exception.t()} def prepare(connection, statement, options \\ []) do DBConnection.prepare(connection, %Query{statement: statement}, options) end @spec prepare_execute(DBConnection.conn(), iodata(), [term()], Keyword.t()) :: {:ok, Query.t(), QuackDB.Result.t()} | {:error, Exception.t()} def prepare_execute(connection, statement, params \\ [], options \\ []) do DBConnection.prepare_execute(connection, %Query{statement: statement}, params, options) end @spec stream(DBConnection.conn(), iodata(), [term()], Keyword.t()) :: Stream.t() def stream(connection, statement, params \\ [], options \\ []) do %Stream{ conn: connection, query: %Query{statement: statement}, params: params, options: options } end @spec rows(DBConnection.conn(), iodata(), [term()], Keyword.t()) :: Enumerable.t() def rows(connection, statement, params \\ [], options \\ []) do connection |> stream(statement, params, options) |> Elixir.Stream.flat_map(&(&1.rows || [])) end @spec maps(DBConnection.conn(), iodata(), [term()], Keyword.t()) :: Enumerable.t() def maps(connection, statement, params \\ [], options \\ []) do connection |> stream(statement, params, options) |> Elixir.Stream.flat_map(&result_maps/1) end defp result_maps(%QuackDB.Result{columns: columns, rows: rows}) when is_list(columns) and is_list(rows) do map_keys = QuackDB.Result.disambiguate_columns(columns) Enum.map(rows, fn row -> map_keys |> Enum.zip(row) |> Map.new() end) end defp result_maps(_result), do: [] end