Packages
ash_sql
0.2.44
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.5
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.16
0.3.15
0.3.14
0.3.13
0.3.12
0.3.11
0.3.10
0.3.9
0.3.8
0.3.7
0.3.6
0.3.5
0.3.4
0.3.3
0.3.2
0.3.1
0.3.0
0.2.93
0.2.92
0.2.91
0.2.90
0.2.89
0.2.88
0.2.87
0.2.86
0.2.85
0.2.84
0.2.83
0.2.82
0.2.81
0.2.80
0.2.79
0.2.78
0.2.77
0.2.76
0.2.75
0.2.74
0.2.73
0.2.72
0.2.71
0.2.70
0.2.69
0.2.68
0.2.67
0.2.66
0.2.65
0.2.64
0.2.63
0.2.62
0.2.61
0.2.60
0.2.59
0.2.58
0.2.57
0.2.56
0.2.55
0.2.54
0.2.53
0.2.52
0.2.51
0.2.50
0.2.49
0.2.48
0.2.47
0.2.46
0.2.45
0.2.44
0.2.43
0.2.42
0.2.41
0.2.40
0.2.39
0.2.38
0.2.37
0.2.36
0.2.35
0.2.34
0.2.33
0.2.32
0.2.31
0.2.30
0.2.29
0.2.28
0.2.27
0.2.26
0.2.25
0.2.24
0.2.23
0.2.22
0.2.21
0.2.20
0.2.19
0.2.18
0.2.17
0.2.16
0.2.15
0.2.14
0.2.13
0.2.12
0.2.11
0.2.10
0.2.9
0.2.8
0.2.7
0.2.6
0.2.5
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.3
0.1.2
0.1.1-rc.20
0.1.1-rc.19
0.1.1-rc.18
0.1.1-rc.17
0.1.1-rc.16
0.1.1-rc.15
0.1.1-rc.14
0.1.1-rc.13
0.1.1-rc.12
0.1.1-rc.11
0.1.1-rc.10
0.1.1-rc.9
0.1.1-rc.8
0.1.1-rc.7
0.1.1-rc.6
0.1.1-rc.5
0.1.1-rc.4
0.1.1-rc.3
0.1.1-rc.2
0.1.1-rc.1
0.1.1-rc.0
Shared utilities for ecto-based sql data layers.
Current section
Files
Jump to
Current section
Files
lib/query.ex
defmodule AshSql.Query do
@moduledoc false
import Ecto.Query, only: [subquery: 1, from: 2]
require Ash.Expr
def resource_to_query(resource, implementation, domain \\ nil) do
from(row in {implementation.table(resource) || "", resource}, [])
|> Map.put(:__ash_domain__, domain)
end
def set_context(resource, data_layer_query, sql_behaviour, context) do
start_bindings = context[:data_layer][:start_bindings_at] || 0
data_layer_query = from(row in data_layer_query, as: ^start_bindings)
data_layer_query =
if context[:data_layer][:table] do
%{
data_layer_query
| from: %{data_layer_query.from | source: {context[:data_layer][:table], resource}}
}
else
data_layer_query
end
data_layer_query =
if context[:data_layer][:schema] do
Ecto.Query.put_query_prefix(data_layer_query, to_string(context[:data_layer][:schema]))
else
data_layer_query
end
data_layer_query =
data_layer_query
|> AshSql.Bindings.default_bindings(resource, sql_behaviour, context)
|> AshSql.Bindings.add_parent_bindings(context)
data_layer_query =
case context[:data_layer][:lateral_join_source] do
{data, path} ->
lateral_join_source_query = path |> List.first() |> elem(0)
lateral_join_source_query.resource
|> Ash.Query.set_context(%{
:data_layer => lateral_join_source_query.context[:data_layer]
})
|> Ash.Query.set_tenant(lateral_join_source_query.tenant)
|> set_lateral_join_prefix(data_layer_query)
|> filter_for_records(data)
|> case do
%{valid?: true} = query ->
Ash.Query.data_layer_query(query)
query ->
{:error, query}
end
|> case do
{:ok, lateral_join_source_query} ->
lateral_join_source_query =
if Enum.count(path) == 2 do
Map.update!(lateral_join_source_query, :__ash_bindings__, fn bindings ->
bindings
|> Map.put(:lateral_join_bindings, [start_bindings + 1])
|> Map.update!(:bindings, fn bindings ->
Map.put(
bindings,
start_bindings + 1,
%{
source: path |> Enum.at(1) |> elem(3) |> Map.get(:source),
path: [path |> Enum.at(1) |> elem(3) |> Map.get(:name)],
type: :inner
}
)
end)
end)
else
lateral_join_source_query
end
{:ok,
Map.update!(data_layer_query, :__ash_bindings__, fn bindings ->
Map.put(
bindings,
:lateral_join_source_query,
lateral_join_source_query
)
|> Map.update!(:current, &(&1 + 1))
end)}
{:error, error} ->
{:error, error}
end
_ ->
{:ok, data_layer_query}
end
case data_layer_query do
{:error, error} ->
{:error, error}
{:ok, data_layer_query} ->
{domain, data_layer_query} = Map.pop(data_layer_query, :__ash_domain__)
case context[:data_layer][:lateral_join_source] do
{_, _} ->
data_layer_query =
data_layer_query
|> Map.update!(:__ash_bindings__, &Map.put(&1, :lateral_join?, true))
|> Map.update!(:__ash_bindings__, &Map.put(&1, :domain, domain))
{:ok, data_layer_query}
_ ->
ash_bindings =
data_layer_query.__ash_bindings__
|> Map.put(:lateral_join?, false)
|> Map.put(:domain, domain)
{:ok, %{data_layer_query | __ash_bindings__: ash_bindings}}
end
end
end
defp filter_for_records(query, records) do
keys =
case Ash.Resource.Info.primary_key(query.resource) do
[] ->
case Ash.Resource.Info.identities(query.resource) do
%{keys: keys} -> keys
_ -> []
end
pkey ->
pkey
end
expr =
case keys do
[] ->
raise "Cannot use lateral joins with a resource that has no primary key and no identities"
keys ->
Enum.reduce(records, Ash.Expr.expr(false), fn record, filter_expr ->
all_keys_match_expr =
Enum.reduce(keys, Ash.Expr.expr(true), fn key, key_expr ->
Ash.Expr.expr(^key_expr and ^Ash.Expr.ref(key) == ^Map.get(record, key))
end)
Ash.Expr.expr(^filter_expr or ^all_keys_match_expr)
end)
end
Ash.Query.do_filter(query, expr)
end
def return_query(%{__ash_bindings__: %{lateral_join?: true}} = query, resource) do
query =
AshSql.Bindings.default_bindings(query, resource, query.__ash_bindings__.sql_behaviour)
if query.__ash_bindings__[:sort_applied?] do
{:ok, query}
else
AshSql.Sort.apply_sort(
query,
query.__ash_bindings__[:sort],
query.__ash_bindings__.resource
)
end
end
def return_query(query, resource) do
query =
AshSql.Bindings.default_bindings(query, resource, query.__ash_bindings__.sql_behaviour)
with_sort_applied =
if query.__ash_bindings__[:sort_applied?] do
{:ok, query}
else
AshSql.Sort.apply_sort(query, query.__ash_bindings__[:sort], resource)
end
case with_sort_applied do
{:error, error} ->
{:error, error}
{:ok, query} ->
if query.__ash_bindings__[:__order__?] && query.windows[:order] do
if query.distinct do
{calculations_require_rewrite, aggregates_require_rewrite, query} =
rewrite_nested_selects(query)
query_with_order =
from(row in query, select_merge: %{__order__: over(row_number(), :order)})
query_without_limit_and_offset =
query_with_order
|> Ecto.Query.exclude(:limit)
|> Ecto.Query.exclude(:offset)
{:ok,
from(row in subquery(query_without_limit_and_offset),
select: row,
order_by: row.__order__
)
|> Map.put(:limit, query.limit)
|> Map.put(:offset, query.offset)
|> AshSql.Bindings.default_bindings(
resource,
query.__ash_bindings__.sql_behaviour,
query.__ash_bindings__.context
)
|> Map.update!(:__ash_bindings__, fn bindings ->
Map.merge(
bindings,
%{
calculations_require_rewrite: calculations_require_rewrite,
aggregates_require_rewrite: aggregates_require_rewrite
},
fn _, v1, v2 -> Map.merge(v1, v2) end
)
end)}
else
order_by = %{query.windows[:order] | expr: query.windows[:order].expr[:order_by]}
{:ok,
%{
query
| windows: Keyword.delete(query.windows, :order),
order_bys: [order_by]
}}
end
else
{:ok, %{query | windows: Keyword.delete(query.windows, :order)}}
end
end
end
defp set_lateral_join_prefix(ash_query, query) do
if Ash.Resource.Info.multitenancy_strategy(ash_query.resource) == :context do
Ash.Query.set_tenant(ash_query, query.prefix)
else
ash_query
end
end
# sobelow_skip ["DOS.StringToAtom"]
def rewrite_nested_selects(query) do
case query.select do
%Ecto.Query.SelectExpr{
expr:
{:merge, [],
[
merge_base,
{:%{}, [], current_merging}
]}
} = select ->
# as we flatten these, they must all remain in the same relative order
# I'm actually not sure why this is required by ecto, but it is :)
merging =
Enum.flat_map(current_merging, fn
{type, {:%{}, _, type_exprs}} when type in [:calculations, :aggregates] ->
Enum.map(type_exprs, fn {name, expr} ->
{String.to_atom("__#{type}__#{name}"), expr}
end)
{type, other} ->
[{type, other}]
end)
aggregate_merges =
current_merging
|> Keyword.get(:aggregates, {:%{}, [], []})
|> elem(2)
|> Map.new(fn {name, _expr} ->
{String.to_existing_atom("__aggregates__#{name}"), name}
end)
calculation_merges =
current_merging
|> Keyword.get(:calculations, {:%{}, [], []})
|> elem(2)
|> Map.new(fn {name, _expr} ->
{String.to_existing_atom("__calculations__#{name}"), name}
end)
new_query = %{
query
| select: %{select | expr: {:merge, [], [merge_base, {:%{}, [], merging}]}}
}
{calculation_merges, aggregate_merges, new_query}
%Ecto.Query.SelectExpr{expr: _other_expr} ->
{%{}, %{}, query}
end
end
def remap_mapped_fields(
results,
query,
calculations_require_rewrite \\ %{},
aggregates_require_rewrite \\ %{}
) do
calculation_names = query.__ash_bindings__.calculation_names
aggregate_names = query.__ash_bindings__.aggregate_names
calculations_require_rewrite =
Map.merge(
query.__ash_bindings__[:calculations_require_rewrite] || %{},
calculations_require_rewrite
)
aggregates_require_rewrite =
Map.merge(
query.__ash_bindings__[:aggregates_require_rewrite] || %{},
aggregates_require_rewrite
)
if Enum.empty?(calculation_names) and Enum.empty?(aggregate_names) and
Enum.empty?(calculations_require_rewrite) and Enum.empty?(aggregates_require_rewrite) do
results
else
Enum.map(results, fn result ->
result
|> remap_to_nested(:calculations, calculations_require_rewrite)
|> remap_to_nested(:aggregates, aggregates_require_rewrite)
|> remap(:calculations, calculation_names)
|> remap(:aggregates, aggregate_names)
end)
end
end
defp remap_to_nested(record, _subfield, mapping) when mapping == %{} do
record
end
defp remap_to_nested(record, subfield, mapping) do
Map.update!(record, subfield, fn subfield_values ->
Enum.reduce(mapping, subfield_values, fn {source, dest}, subfield_values ->
subfield_values
|> Map.put(dest, Map.get(record, source))
|> Map.delete(source)
end)
end)
end
defp remap(record, _subfield, mapping) when mapping == %{} do
record
end
defp remap(record, subfield, mapping) do
Map.update!(record, subfield, fn subfield_values ->
Enum.reduce(mapping, subfield_values, fn {dest, source}, subfield_values ->
subfield_values
|> Map.put(dest, Map.get(subfield_values, source))
|> Map.delete(source)
end)
end)
end
end