Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
83 changes: 40 additions & 43 deletions lib/efsql.ex
Original file line number Diff line number Diff line change
Expand Up @@ -5,33 +5,33 @@ defmodule Efsql do
A statement flows through the textbook pipeline:

SQL text
|> Efsql.Parser.to_logical() # parse tree -> Efsql.Logical.Select
|> Efsql.Parser.sql_to_logical() # parse tree -> Efsql.Logical.Select
|> resolve tenant
|> Efsql.Rewrite.normalize() # rewrite passes
|> Efsql.Planner.plan() # access-path selection -> Efsql.Physical.Plan
|> Efsql.Executor.run() # adapter pull + operator pipeline
"""

import Ecto.Query

def hello() do
tenant = EctoFoundationDB.Tenant.open!(Efsql.Repo, "localhost")

query = from(s in "secrets", select: [id: s.id, iv: s.iv])
r1 = Efsql.Repo.all(query, prefix: tenant)
A query across tenants (`*.table`) is planned by `Efsql.Fanout` instead,
which plans each tenant's read with `Efsql.Planner`.

r2 = all("select id, iv from localhost.secrets;")

{r1, r2}
end
The CLI and the TUI run statements in an `Efsql.Session` (open tenants,
settings, the active tenant) and show the `Efsql.Result` it returns.
`qall/3` and `all/2` run one outside any session.
"""

def all(sql, options \\ []) do
{_, result, _tenants} = qall(sql, options)
result
end

@doc """
Runs one statement outside any session, reusing and returning the
`tenants` cache: `{plan, rows, tenants}`. See `Efsql.Session` for the
CLI's and TUI's way in.
"""
def qall(sql, options \\ [], tenants \\ %{}) do
sql |> Efsql.Parser.sql_to_logical() |> run_logical(options, tenants)
{result, session} = Efsql.Session.run(%Efsql.Session{tenants: tenants}, sql, options)
{result.plan, result.rows, session.tenants}
end

@doc """
Expand All @@ -52,58 +52,55 @@ defmodule Efsql do
{plan, Efsql.Executor.run(plan), tenants}
end

def stream(sql) do
{logical, _tenants} = sql_to_logical(sql)
query = logical |> Efsql.Rewrite.normalize() |> Efsql.Planner.to_ecto_query()
{query, Efsql.Repo.stream(query)}
end

def sql_to_logical(sql, tenants \\ %{}) do
logical = %Efsql.Logical.Select{} = Efsql.Parser.sql_to_logical(sql)
resolve_tenant(logical, tenants)
end

def resolve_tenant(%Efsql.Logical.Select{prefix: {:all_tenants, _}}, _tenants) do
raise Efsql.Exception.Unsupported, "a query across tenants can't be streamed"
end

# Already resolved, e.g. to the TUI session's active tenant.
def resolve_tenant(%Efsql.Logical.Select{prefix: nil, tenant: tenant} = logical, tenants)
when tenant != nil,
do: {logical, tenants}
defp resolve_tenant(%Efsql.Logical.Select{prefix: nil, tenant: tenant} = logical, tenants)
when tenant != nil,
do: {logical, tenants}

def resolve_tenant(%Efsql.Logical.Select{} = logical, tenants) do
defp resolve_tenant(%Efsql.Logical.Select{} = logical, tenants) do
{tenant_name, storage_id} =
case logical.prefix do
{storage_id, tenant_name} -> {tenant_name, storage_id}
nil -> raise "Tenant required"
tenant_name -> {tenant_name, nil}
end

if not Map.has_key?(tenants, {tenant_name, storage_id}) and
not EctoFoundationDB.Tenant.exists?(Efsql.Repo, tenant_name) do
raise Efsql.Exception.Unsupported, "Tenant '#{tenant_name}' does not exist"
end

{tenant, tenants} = open_tenant(tenants, tenant_name, storage_id)
{tenant, tenants} = open_tenant(tenants, tenant_name, storage_id, check_exists: true)
{%Efsql.Logical.Select{logical | tenant: tenant}, tenants}
end

@doc """
Opens a tenant through the session's cache of open tenants, keyed by
name and storage id (nil for the Repo's own).

Read-only: the tenant is opened with `migrate: false`, and
`Tenant.open/3` only opens (never creates), so efsql never writes to the
database it is exploring. A storage id other than the Repo's gets its
tenant cache started first. With `check_exists: true`, a tenant that
doesn't exist in that storage id is an `Efsql.Exception.Unsupported`
error; tenants just listed from the storage id skip the check.
"""
def open_tenant(tenants, tenant_name, storage_id) do
def open_tenant(tenants, tenant_name, storage_id, options \\ []) do
case Map.fetch(tenants, {tenant_name, storage_id}) do
{:ok, tenant} ->
{tenant, tenants}

:error ->
if storage_id, do: Efsql.Discover.ensure_storage_cache(storage_id)
config = Efsql.Repo.config()
config = if storage_id, do: Keyword.put(config, :storage_id, storage_id), else: config

if Keyword.get(options, :check_exists, false) and
not EctoFoundationDB.Tenant.Backend.exists?(
Ecto.Adapters.FoundationDB.db(Efsql.Repo),
tenant_name,
config
) do
raise Efsql.Exception.Unsupported, "Tenant '#{tenant_name}' does not exist"
end

open_opts = if storage_id, do: [storage_id: storage_id], else: []

# migrate: false keeps this read-only. Tenant.open/3 already requires
# the tenant to exist (only open!/3 creates), so with the migration
# step skipped efsql never writes to the database it is exploring.
tenant =
EctoFoundationDB.Tenant.open(Efsql.Repo, tenant_name, [migrate: false] ++ open_opts)

Expand Down
35 changes: 2 additions & 33 deletions lib/efsql/aggregate.ex
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ defmodule Efsql.Aggregate do
@spec run([map], [atom], [aggregate]) :: [map]
def run(rows, group_by, aggregates) do
rows
|> Enum.group_by(fn row -> Enum.map(group_by, &key(Map.get(row, &1))) end)
|> Enum.group_by(fn row -> Enum.map(group_by, &Types.equality_key(Map.get(row, &1))) end)
|> ensure_one_group(group_by)
|> Enum.map(fn {_key, rows} ->
first = List.first(rows, %{})
Expand All @@ -35,20 +35,13 @@ defmodule Efsql.Aggregate do
Map.put(acc, name, compute(function, arg, rows))
end)
end)
|> Enum.sort(&(compare_groups(&1, &2, group_by) != :gt))
|> Efsql.Executor.sort(Enum.map(group_by, &{:asc, &1}))
end

# A whole-table aggregate answers even when no row matched.
defp ensure_one_group(groups, []) when groups == %{}, do: [{[], []}]
defp ensure_one_group(groups, _group_by), do: Enum.to_list(groups)

# A term that is equal for values Types.compare/2 calls equal.
defp key(%Decimal{} = d), do: Decimal.normalize(d)
defp key(%NaiveDateTime{microsecond: {us, _}} = t), do: %{t | microsecond: {us, 6}}
defp key(%DateTime{microsecond: {us, _}} = t), do: %{t | microsecond: {us, 6}}
defp key(%Time{microsecond: {us, _}} = t), do: %{t | microsecond: {us, 6}}
defp key(value), do: value

# -- aggregates --

defp compute(:count, :star, rows), do: length(rows)
Expand Down Expand Up @@ -105,28 +98,4 @@ defmodule Efsql.Aggregate do
defp decimal(%Decimal{} = d), do: d
defp decimal(n) when is_integer(n), do: Decimal.new(n)
defp decimal(n) when is_float(n), do: Decimal.from_float(n)

# -- group order --

# By key, NULLs last, like an ascending ORDER BY.
defp compare_groups(_a, _b, []), do: :eq

defp compare_groups(a, b, [field | rest]) do
case {Map.get(a, field), Map.get(b, field)} do
{nil, nil} ->
compare_groups(a, b, rest)

{nil, _} ->
:gt

{_, nil} ->
:lt

{va, vb} ->
case Types.compare(va, vb) do
:eq -> compare_groups(a, b, rest)
cmp -> cmp
end
end
end
end
77 changes: 31 additions & 46 deletions lib/efsql/cli.ex
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
defmodule Efsql.Cli do
defstruct args: [], history: [], debug: false, tenants: %{}, limit: 15, tenant_batch: nil
defstruct args: [], history: [], debug: false, session: %Efsql.Session{}

use GenServer

Expand Down Expand Up @@ -92,9 +92,7 @@ defmodule Efsql.Cli do
"""
Meta-commands:
\\tenants [storage_id] list tenants (optionally for a specific storage id)
\\set limit N set the default row limit (currently #{state.limit})
\\set tenant_batch N|off read *.table queries N tenants per transaction
(currently #{state.tenant_batch || "off"})
#{settings_help(state.session.settings)}
\\? show this help
""",
:light_black
Expand Down Expand Up @@ -140,57 +138,41 @@ defmodule Efsql.Cli do
state
end

defp handle_input("\\set limit " <> rest, state = %__MODULE__{}) do
case Integer.parse(String.trim(rest)) do
{n, ""} when n > 0 ->
Owl.IO.puts(Owl.Data.tag("limit set to #{n}", :light_black))
%__MODULE__{state | limit: n}
defp handle_input("\\set " <> rest, state = %__MODULE__{}) do
case Efsql.Settings.set(state.session.settings, rest) do
{:ok, settings, message} ->
Owl.IO.puts(Owl.Data.tag(message, :light_black))
%__MODULE__{state | session: %{state.session | settings: settings}}

_ ->
print_error("Usage: \\set limit <positive integer>")
state
end
end

defp handle_input("\\set tenant_batch " <> rest, state = %__MODULE__{}) do
case {String.trim(rest), Integer.parse(String.trim(rest))} do
{"off", _} ->
Owl.IO.puts(Owl.Data.tag("tenant_batch off: one transaction", :light_black))
%__MODULE__{state | tenant_batch: nil}

{_, {n, ""}} when n > 0 ->
Owl.IO.puts(Owl.Data.tag("tenant_batch set to #{n}", :light_black))
%__MODULE__{state | tenant_batch: n}

_ ->
print_error("Usage: \\set tenant_batch <positive integer> | off")
{:error, usage} ->
print_error(usage)
state
end
end

defp handle_input(data, state = %__MODULE__{}) do
limit_sql = "limit #{state.limit + 1}"
limit = state.session.settings.limit
limit_sql = "limit #{limit + 1}"

{tenants} =
session =
try do
{sql, display_limit} =
if String.match?(data, ~r/\blimit\b/i),
do: {data, :all},
else: {String.replace(data, ~r/;\s*$/, " #{limit_sql};"), state.limit}

options = if state.tenant_batch, do: [tenant_batch: state.tenant_batch], else: []
{call, rows, tenants} = Efsql.qall(sql, options, state.tenants)
if state.debug, do: print_debug(call)
print_table(rows, display_limit, call.columns)
print_snapshots(Efsql.Fanout.transactions(call))
{tenants}
else: {String.replace(data, ~r/;\s*$/, " #{limit_sql};"), limit}

{result, session} = Efsql.Session.run(state.session, sql)
if state.debug, do: print_debug(result.plan)
print_table(result.rows, display_limit, result.columns)
print_transactions(result.transactions)
session
rescue
e ->
print_error(e)
{state.tenants}
state.session
end

%__MODULE__{state | history: [data | state.history], tenants: tenants}
%__MODULE__{state | history: [data | state.history], session: session}
end

def init_ecto_foundationdb!(args) do
Expand Down Expand Up @@ -258,12 +240,10 @@ defmodule Efsql.Cli do
print_rows(display_rows, more?, columns)
end

# Columns in select-list order; `select *` (no columns) sorts them.
defp print_rows(rows, more?, columns) do
rows
|> Enum.map(fn row ->
row = if columns, do: Map.new(columns, &{&1, Map.get(row, &1)}), else: row
Map.new(row, fn {k, v} -> {to_string(k), format_value(v)} end)
Map.new(columns, &{to_string(&1), format_value(Map.get(row, &1))})
end)
|> Owl.Table.new(
border_style: :solid_rounded,
Expand All @@ -277,14 +257,19 @@ defmodule Efsql.Cli do
Owl.IO.puts(Owl.Data.tag(label, :light_black))
end

defp print_snapshots(1), do: :ok
defp settings_help(settings) do
Enum.map_join(Efsql.Settings.describe(settings), "\n", fn {command, what} ->
" " <> String.pad_trailing(command, 22) <> " " <> what
end)
end

defp print_snapshots(n) do
defp print_transactions(1), do: :ok

defp print_transactions(n) do
Owl.IO.puts(Owl.Data.tag("(read in #{n} transactions)", :yellow))
end

defp column_sorter(nil), do: :asc

# Owl sorts a table's columns; keep them in the result's order.
defp column_sorter(columns) do
position = columns |> Enum.map(&to_string/1) |> Enum.with_index() |> Map.new()
&(Map.fetch!(position, &1) <= Map.fetch!(position, &2))
Expand Down
4 changes: 1 addition & 3 deletions lib/efsql/discover.ex
Original file line number Diff line number Diff line change
Expand Up @@ -171,9 +171,7 @@ defmodule Efsql.Discover do
end

def indexes(tenant, source) do
EctoFoundationDB.Layer.Metadata.transactional(tenant, source, fn _tx, metadata ->
for idx <- metadata.indexes, do: %{name: idx[:id], fields: idx[:fields]}
end)
for idx <- Planner.indexes(tenant, source), do: %{name: idx[:id], fields: idx[:fields]}
end

# -- schema cache --
Expand Down
Loading
Loading