Packages
electric
1.6.1
1.7.8
1.7.7
1.7.6
1.7.5
1.7.4
1.7.3
1.7.2
1.7.1
1.7.0
1.6.10
1.6.9
1.6.8
1.6.7
1.6.6
1.6.5
1.6.4
1.6.3
1.6.2
1.6.1
1.6.0
1.5.1
1.5.0
1.4.16
1.4.16-beta-1
1.4.15
1.4.14
1.4.13
1.4.12
1.4.11
1.4.10
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.4
1.3.3
1.3.2
1.2.4
1.2.3
1.2.2
1.2.1
1.2.0
1.1.14
1.1.13
1.1.12
1.1.11
1.1.10
1.1.9
1.1.8
1.1.7
1.1.6
retired
1.1.5
retired
1.1.4
retired
1.1.3
retired
1.1.2
1.1.1
1.1.0
1.0.24
1.0.23
1.0.22
1.0.21
1.0.20
1.0.19
1.0.18
1.0.17
1.0.15
1.0.13
1.0.12
1.0.11
1.0.10
1.0.9
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
1.0.0-beta.23
1.0.0-beta.22
1.0.0-beta.20
1.0.0-beta.19
1.0.0-beta.18
1.0.0-beta.17
1.0.0-beta.16
1.0.0-beta.15
1.0.0-beta.14
1.0.0-beta.13
1.0.0-beta.12
1.0.0-beta.11
1.0.0-beta.10
1.0.0-beta.9
1.0.0-beta.8
1.0.0-beta.7
1.0.0-beta.6
1.0.0-beta.5
1.0.0-beta.4
1.0.0-beta.3
1.0.0-beta.2
1.0.0-beta.1
0.9.5
0.9.4
0.9.3
0.9.2
0.9.1
0.9.0
0.8.1
0.8.0
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.3
0.6.2
0.6.1
0.5.2
0.4.4
Postgres sync engine. Sync little subsets of your Postgres data into local apps and services.
Current section
Files
Jump to
Current section
Files
lib/electric/shapes/consumer/event_handler_builder.ex
defmodule Electric.Shapes.Consumer.EventHandlerBuilder do
# Builds the initial event handler and ordered setup effects for a consumer shape.
alias Electric.Shapes.Consumer.EventHandler
alias Electric.Shapes.Consumer.Materializer
alias Electric.Shapes.Consumer.SetupEffects
alias Electric.Shapes.Consumer.State
alias Electric.Shapes.DnfPlan
alias Electric.Shapes.Shape
@spec build(State.t(), :create | :restore) ::
{:ok, EventHandler.t(), [SetupEffects.t()]}
def build(%State{shape: %Shape{shape_dependencies_handles: dep_handles}} = state, action)
when dep_handles != [] do
{:ok, dnf_plan} = DnfPlan.compile(state.shape)
dependency_move_policy = dependency_move_policy(state.stack_id, state.shape)
{views, dep_handle_to_ref, dep_index_to_ref} =
dep_handles
|> Enum.with_index()
|> Enum.reduce({%{}, %{}, %{}}, fn {handle, index},
{views, handle_mapping, index_mapping} ->
materializer_opts = %{stack_id: state.stack_id, shape_handle: handle}
:ok = Materializer.wait_until_ready(materializer_opts)
view = Materializer.get_link_values(materializer_opts)
ref = ["$sublink", Integer.to_string(index)]
{Map.put(views, ref, view), Map.put(handle_mapping, handle, {index, ref}),
Map.put(index_mapping, index, ref)}
end)
buffer_max_transactions =
Electric.StackConfig.lookup(
state.stack_id,
:subquery_buffer_max_transactions,
Electric.Config.default(:subquery_buffer_max_transactions)
)
handler = %EventHandler.Subqueries.Steady{
shape_info: %Electric.Shapes.Consumer.Subqueries.ShapeInfo{
shape: state.shape,
stack_id: state.stack_id,
shape_handle: state.shape_handle,
dnf_plan: dnf_plan,
ref_resolver:
Electric.Shapes.Consumer.Subqueries.RefResolver.new(dep_handle_to_ref, dep_index_to_ref),
buffer_max_transactions: buffer_max_transactions,
dependency_move_policy: dependency_move_policy
},
views: views
}
{:ok, handler,
[%SetupEffects.SubscribeShape{action: action}, %SetupEffects.SeedSubqueryIndex{}]}
end
def build(%State{} = state, action) do
handler = %EventHandler.Default{
shape: state.shape,
stack_id: state.stack_id,
shape_handle: state.shape_handle
}
{:ok, handler, [%SetupEffects.SubscribeShape{action: action}]}
end
defp dependency_move_policy(stack_id, _shape) do
feature_flags = Electric.StackConfig.lookup(stack_id, :feature_flags, [])
if "tagged_subqueries" not in feature_flags do
:invalidate_on_dependency_move
else
:stream_dependency_moves
end
end
end