Current section

Files

Jump to
ash_neo4j lib data_layer.ex
Raw

lib/data_layer.ex

defmodule AshNeo4j.DataLayer do
@moduledoc "Ash DataLayer for Neo4j"
@behaviour Ash.DataLayer
alias Ash.Actions.Sort
alias AshNeo4j.DataLayer.Info
alias AshNeo4j.QueryHelper
@filter_stream_size 100
@impl true
def can?(_, :read), do: true
#def can?(_, :create), do: true
#def can?(_, :update), do: true
#def can?(_, :upsert), do: true
#def can?(_, :destroy), do: true
def can?(_, :sort), do: true
def can?(_, :filter), do: true
def can?(_, :limit), do: true
# def can?(_, :bulk_create), do: true
def can?(_, :offset), do: true
def can?(_, :boolean_filter), do: true
#def can?(_, :transact), do: true
def can?(_, {:filter_expr, _}), do: true
def can?(_, :nested_expressions), do: true
def can?(_, :expression_calculation_sort), do: true
def can?(_, {:sort, _}), do: true
def can?(_, _), do: false
@neo4j %Spark.Dsl.Section{
name: :neo4j,
examples: [
"""
neo4j do
label :Comment
store [:title]
translate id: :uuid
end
"""
],
schema: [
label: [
type: :atom,
doc: "The node label",
required: true
],
store: [
type: {:list, :atom},
doc: "The attributes to be stored as node properties, without translation",
required: true
],
translate: [
type: :keyword_list,
doc: "Optional attribute to node property translations"
]
]
}
@impl true
def limit(query, offset, _), do: {:ok, %{query | limit: offset}}
@impl true
def offset(query, offset, _), do: {:ok, %{query | offset: offset}}
@impl true
def filter(query, filter, _resource) do
{:ok, %{query | filter: filter}}
end
@impl true
def sort(query, sort, _resource) do
{:ok, %{query | sort: sort}}
end
@doc false
def store_opt(attributes) do
if Enum.all?(attributes, &is_atom/1) do
{:ok, attributes}
else
{:error, "Expected all attribute names to be atoms"}
end
end
@sections [@neo4j]
use Spark.Dsl.Extension,
sections: @sections,
persisters: [AshNeo4j.DataLayer.Transformer]
defmodule Query do
@moduledoc false
defstruct [:resource, :sort, :filter, :limit, :offset, :domain]
end
@impl true
def run_query(query, resource) do
label = Info.label(resource)
nodes = QueryHelper.query_nodes(label, query)
#|> IO.inspect(label: "AshNeo4j.DataLayer.run_query match_nodes result")
results =
convert_nodes_to_resources(nodes, label)
|> filter_stream(query.domain, query.filter)
|> sort_stream(resource, query.domain, query.sort)
|> offset_stream(query.offset)
|> limit_stream(query.limit)
#|> IO.inspect(label: "AshNeo4j.DataLayer.run_query result")
{:ok, results}
end
@impl true
def resource_to_query(resource, domain) do
%Query{resource: resource, domain: domain}
end
@impl true
def transaction(resource, fun, _timeout, _) do
label = Info.label(resource)
:global.trans(
{{:neo4j, label}, System.unique_integer()},
fn ->
try do
Process.put({:neo4j_in_transaction, label}, true)
{:res, fun.()}
catch
{{:neo4j_rollback, ^label}, value} ->
{:error, value}
end
end,
[node() | :erlang.nodes()],
0
)
|> case do
{:res, result} -> {:ok, result}
{:error, error} -> {:error, error}
:aborted -> {:error, "transaction failed"}
end
end
@impl true
def rollback(resource, error) do
throw({{:neo4j_rollback, Info.label(resource)}, error})
end
@impl true
def in_transaction?(resource) do
Process.get({:neo4j_in_transaction, Info.label(resource)}, false) == true
end
def filter_matches(records, nil, _domain), do: records
def filter_matches(records, filter, domain) do
{:ok, records} = Ash.Filter.Runtime.filter_matches(domain, records, filter)
records
end
# converts nodes to resources, where the input is a list of related node groups
# the output of each group is a single resource, enriched with attributes linking related nodes
defp convert_nodes_to_resources(groups, label) when is_list(groups) and is_atom(label) do
#IO.inspect(label, label: "AshNeo4j.DataLayer.convert_nodes_to_resources label")
resource = Info.resource(label)
groups #|> IO.inspect(label: "AshNeo4j.DataLayer.convert_nodes_to_resources groups")
|> Stream.map(fn related_nodes ->
source_node = Map.get(related_nodes, "s")
dest_node = Map.get(related_nodes, "d")
if dest_node != nil do
dest_label = List.first(dest_node.labels)
dest_resource = convert_node_to_resource(Info.resource(String.to_atom(dest_label)), dest_node, [])
relationship = Ash.Resource.Info.relationship(resource, String.downcase(dest_label))
enrichment = {relationship.source_attribute, Map.get(dest_resource, relationship.destination_attribute)}
convert_node_to_resource(resource, source_node, [enrichment])
else
convert_node_to_resource(resource, source_node, [])
end
#|> IO.inspect(label: "AshNeo4j.DataLayer.convert_nodes_to_resources resource")
end)
end
defp convert_node_to_resource(resource, node, enrichments) when is_atom(resource) and is_map(node) do
#IO.inspect(node, label: "AshNeo4j.DataLayer.convert_node_to_resource node")
store = Info.store(resource)
translate = Info.translate(resource)
stored = Enum.into(store, %{}, fn field ->
{field, Map.get(node.properties, to_string(field))}
end)
translated = Enum.into(translate, stored, fn {resource_field, node_field} ->
{resource_field, Map.get(node.properties, to_string(node_field))}
end)
Enum.into(enrichments, translated, fn {field, value} ->
{field, value}
end) # |> IO.inspect(label: "AshNeo4j.DataLayer.convert_node_to_resource enriched")
|> Map.put(:__struct__, resource)
|> Map.put(:__data_layer__, __MODULE__)
# TODO metadata should be a struct including neo4j node id?
|> Map.put(:__metadata__, %{})
|> Map.put(:aggregates, %{})
|> Map.put(:calculations, %{})
#|> IO.inspect(label: "AshNeo4j.DataLayer.convert_node_to_resource result")
end
defp sort_stream(stream, _resource, _domain, sort) when sort in [nil, []] do
stream
end
defp sort_stream(stream, resource, domain, sort) do
Sort.runtime_sort(stream, sort, domain: domain, resource: resource)
end
defp filter_stream(stream, _domain, nil), do: stream
defp filter_stream(stream, domain, filter) do
stream
|> Stream.chunk_every(@filter_stream_size)
|> Stream.flat_map(fn chunk ->
filter_matches(chunk, filter, domain)
end)
end
defp offset_stream(stream, offset) when offset in [0, nil], do: stream
defp offset_stream(stream, offset), do: Stream.drop(stream, offset)
defp limit_stream(stream, nil), do: stream
defp limit_stream(stream, limit), do: Stream.take(stream, limit)
end