Packages
ash_sql
0.2.5
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]
def resource_to_query(resource, implementation) do
from(row in {implementation.table(resource) || "", resource}, [])
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)
|> 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} ->
case context[:data_layer][:lateral_join_source] do
{_, _} ->
data_layer_query =
data_layer_query
|> Map.update!(:__ash_bindings__, &Map.put(&1, :lateral_join?, true))
{:ok, data_layer_query}
_ ->
ash_bindings =
data_layer_query.__ash_bindings__
|> Map.put(:lateral_join?, false)
{:ok, %{data_layer_query | __ash_bindings__: ash_bindings}}
end
end
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
)
|> 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
def rewrite_nested_selects(query) do
case query.select do
%Ecto.Query.SelectExpr{
expr:
{:merge, [],
[
merge_base,
{:%{}, [], merging}
]}
} = select ->
{merging, aggregate_merges} = remap_sub_select(merging, :aggregates)
{new_sub_selects, calculation_merges} =
remap_sub_select(merging, :calculations)
new_query =
%{
query
| select: %{select | expr: {:merge, [], [merge_base, {:%{}, [], new_sub_selects}]}}
}
{calculation_merges, aggregate_merges, new_query}
%Ecto.Query.SelectExpr{expr: _other_expr} ->
{%{}, %{}, query}
end
end
# sobelow_skip ["DOS.StringToAtom"]
defp remap_sub_select(merging, sub_key) do
case Keyword.fetch(merging, sub_key) do
{:ok, {:%{}, [], nested}} ->
Enum.reduce(nested, {Keyword.delete(merging, sub_key), %{}}, fn {name, expr},
{subselect, remapping} ->
new_name = String.to_atom("__#{sub_key}__#{name}")
{Keyword.put(subselect, new_name, expr), Map.put(remapping, new_name, name)}
end)
:error ->
{merging, %{}}
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