Packages
electric
1.6.8
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/setup_effects.ex
defmodule Electric.Shapes.Consumer.SetupEffects do
# Executes ordered boot-time setup effects for consumer handler initialization.
alias Electric.Replication.ShapeLogCollector
alias Electric.Shapes.Consumer.State
alias Electric.Shapes.Filter.Indexes.SubqueryIndex
require Logger
defmodule SubscribeShape do
@moduledoc false
defstruct [:action]
end
defmodule SeedSubqueryIndex do
@moduledoc false
defstruct []
end
@type t() :: %SubscribeShape{} | %SeedSubqueryIndex{}
@spec execute([t()], State.t()) :: {:ok, State.t()} | {:error, State.t()}
def execute(effects, %State{} = state) when is_list(effects) do
Enum.reduce_while(effects, {:ok, state}, fn effect, {:ok, state} ->
case execute_effect(effect, state) do
{:ok, %State{} = state} -> {:cont, {:ok, state}}
{:error, %State{} = state} -> {:halt, {:error, state}}
end
end)
end
defp execute_effect(%SubscribeShape{action: action}, %State{} = state) do
case ShapeLogCollector.add_shape(state.stack_id, state.shape_handle, state.shape, action) do
:ok ->
{:ok, state}
{:error, error} ->
Logger.warning(
"Shape #{state.shape_handle} cannot subscribe due to #{inspect(error)} - invalidating shape"
)
{:error, state}
end
end
defp execute_effect(%SeedSubqueryIndex{}, %State{event_handler: %{views: views}} = state) do
case SubqueryIndex.for_stack(state.stack_id) do
nil ->
{:ok, state}
index ->
for {ref, view} <- views do
dep_index = ref |> List.last() |> String.to_integer()
SubqueryIndex.seed_membership(
index,
state.shape_handle,
ref,
dep_index,
view
)
end
SubqueryIndex.mark_ready(index, state.shape_handle)
{:ok, state}
end
end
defp execute_effect(%SeedSubqueryIndex{}, %State{} = state), do: {:ok, state}
end