Packages
electric
1.7.7
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/shape_cache/shape_status/shape_db.ex
defmodule Electric.ShapeCache.ShapeStatus.ShapeDb do
@moduledoc """
SQLite-backed persistent storage for shape metadata.
The WriteBuffer provides buffering for writes to prevent timeout cascades.
Only `handle_for_shape` and `shape_for_handle` need buffer awareness since
they are entry points for new requests. Other functions are called after
ShapeStatus has already updated its ETS cache.
"""
alias Electric.Shapes.Shape
alias Electric.ShapeCache.ShapeStatus.ShapeDb.Connection
alias Electric.ShapeCache.ShapeStatus.ShapeDb.Query
alias Electric.ShapeCache.ShapeStatus.ShapeDb.WriteBuffer
alias Electric.ShapeCache.ShapeStatus.ShapeDb.Statistics
import Electric, only: [is_stack_id: 1, is_shape_handle: 1]
import Connection,
only: [
checkout!: 3,
checkout_write!: 3,
checkout_write!: 4
]
@type shape_handle() :: Electric.shape_handle()
@type stack_id() :: Electric.stack_id()
defmodule Error do
defexception [:message]
@impl true
def exception(args) do
action = Keyword.get(args, :action, :read)
{:ok, error} = Keyword.fetch(args, :error)
%__MODULE__{message: "ShapeDb #{action} failed: #{inspect(error)}"}
end
end
def add_shape(stack_id, %Shape{} = shape, shape_handle)
when is_stack_id(stack_id) and is_shape_handle(shape_handle) do
{comparable_shape, shape_hash} = Shape.comparable_hash(shape)
relations = Shape.list_relations(shape)
with :ok <-
WriteBuffer.add_shape(
stack_id,
shape_handle,
shape,
comparable_shape,
shape_hash,
relations
) do
{:ok, shape_hash}
end
end
def remove_shape(stack_id, shape_handle) when is_stack_id(stack_id) do
if handle_exists?(stack_id, shape_handle) do
WriteBuffer.remove_shape(stack_id, shape_handle)
else
{:error, {:enoshape, shape_handle}}
end
end
def mark_snapshot_complete(stack_id, shape_handle) do
if handle_exists?(stack_id, shape_handle) do
WriteBuffer.queue_snapshot_complete(stack_id, shape_handle)
else
:error
end
end
def reset(stack_id) when is_stack_id(stack_id) do
WriteBuffer.clear(stack_id)
checkout_write!(stack_id, :reset, &Query.reset/1)
end
@doc """
Find a handle for a shape. Checks buffer first, then SQLite.
Returns :error if the handle is tombstoned (being deleted).
"""
def handle_for_shape(stack_id, %Shape{} = shape) when is_stack_id(stack_id) do
checkout_fun = &checkout!(stack_id, :handle_for_shape, &1)
handle_for_shape_inner(stack_id, shape, checkout_fun)
end
@doc """
Find a handle for a shape using the write connection to guarantee consistency.
"""
def handle_for_shape_critical(stack_id, %Shape{} = shape, timeout \\ 10_000)
when is_stack_id(stack_id) do
checkout_fun = &checkout_write!(stack_id, :handle_for_shape_critical, &1, timeout)
handle_for_shape_inner(stack_id, shape, checkout_fun)
end
defp handle_for_shape_inner(stack_id, %Shape{} = shape, checkout_fun) do
{comparable_shape, _shape_hash} = Shape.comparable_hash(shape)
case WriteBuffer.lookup_handle(stack_id, comparable_shape) do
{:ok, handle} ->
{:ok, handle}
:not_found ->
case checkout_fun.(&Query.handle_for_shape(&1, comparable_shape)) do
{:ok, handle} ->
if WriteBuffer.is_tombstoned?(stack_id, handle) do
:error
else
{:ok, handle}
end
:error ->
:error
end
end
end
@doc """
Find a shape by its handle. Checks buffer first, then SQLite.
Returns :error if the handle is tombstoned (being deleted).
"""
def shape_for_handle(stack_id, shape_handle) when is_stack_id(stack_id) do
case WriteBuffer.lookup_shape(stack_id, shape_handle) do
{:ok, shape} ->
{:ok, shape}
:not_found ->
if WriteBuffer.is_tombstoned?(stack_id, shape_handle) do
:error
else
case checkout!(stack_id, :shape_for_handle, &Query.shape_for_handle(&1, shape_handle)) do
{:ok, shape} ->
# Re-check tombstone after SQLite read to avoid race where
# shape was tombstoned between the first check and SQLite query
if WriteBuffer.is_tombstoned?(stack_id, shape_handle),
do: :error,
else: {:ok, shape}
:error ->
:error
end
end
end
end
def list_shapes(stack_id) when is_stack_id(stack_id) do
buffered = WriteBuffer.list_buffered_shapes(stack_id)
tombstones = WriteBuffer.tombstoned_handles(stack_id)
case checkout!(stack_id, :list_shapes, &Query.list_shapes/1) do
{:ok, sqlite_shapes} ->
shapes =
buffered
|> Stream.concat(
Stream.reject(sqlite_shapes, fn {handle, _shape} ->
MapSet.member?(tombstones, handle)
end)
)
# Deduplicate to handle race between buffer flush and SQLite read
|> Enum.uniq_by(fn {handle, _} -> handle end)
{:ok, shapes}
error ->
error
end
end
def list_shapes!(stack_id) when is_stack_id(stack_id) do
stack_id |> list_shapes() |> raise_on_error!(:list_shapes)
end
def shape_handles_for_relations(stack_id, relations) when is_stack_id(stack_id) do
buffered_handles = WriteBuffer.handles_for_relations(stack_id, relations)
tombstones = WriteBuffer.tombstoned_handles(stack_id)
case checkout!(
stack_id,
:shape_handles_for_relations,
&Query.shape_handles_for_relations(&1, relations)
) do
{:ok, sqlite_handles} ->
filtered_sqlite = Enum.reject(sqlite_handles, &MapSet.member?(tombstones, &1))
{:ok, Enum.uniq(buffered_handles ++ filtered_sqlite)}
error ->
error
end
end
def shape_handles_for_relations!(stack_id, relations) when is_stack_id(stack_id) do
stack_id
|> shape_handles_for_relations(relations)
|> raise_on_error!(:shape_handles_for_relations)
end
def reduce_shapes(stack_id, acc, reducer_fun) when is_function(reducer_fun, 2) do
case list_shapes(stack_id) do
{:ok, shapes} -> Enum.reduce(shapes, acc, reducer_fun)
{:error, _} = error -> error
end
end
# Only used during boot when buffer is empty
def reduce_shape_meta(stack_id, acc, reducer_fun) when is_function(reducer_fun, 2) do
checkout!(stack_id, :reduce_shape_meta, fn %Connection{} = conn ->
conn
|> Query.list_shape_meta_stream()
|> Enum.reduce(acc, reducer_fun)
end)
end
# May be slightly inaccurate during concurrent modifications due to
# the `pending_count_diff` being updated after writes are in the database
# meaning that changes may be counted twice for a (very) short period.
def count_shapes(stack_id) do
case checkout!(stack_id, :count_shapes, &Query.count_shapes/1) do
{:ok, sqlite_count} -> {:ok, sqlite_count + WriteBuffer.pending_count_diff(stack_id)}
error -> error
end
end
def count_shapes!(stack_id) do
try do
stack_id |> count_shapes() |> raise_on_error!(:count_shapes)
rescue
# the connection pool has its own registry, so attempting to checkout a
# connection will raise an ArgumentError if that registry isn't running
ArgumentError -> :error
catch
# connection pool has not started
:exit, {:noproc, {NimblePool, :checkout, _args}} -> :error
end
end
@doc false
def handle_exists?(stack_id, shape_handle) when is_stack_id(stack_id) do
case WriteBuffer.has_handle?(stack_id, shape_handle) do
true -> true
false -> false
:unknown -> checkout!(stack_id, :handle_exists?, &Query.handle_exists?(&1, shape_handle))
end
end
def validate_existing_shapes(stack_id) do
# Must flush buffer before validating so we check persisted state
WriteBuffer.flush_sync(stack_id)
with {:ok, removed_handles} <-
checkout_write!(
stack_id,
:validate_existing_shapes,
fn %Connection{} = conn ->
with {:ok, handles} <- Query.select_invalid(conn) do
Enum.each(handles, fn handle ->
:ok = Query.remove_shape(conn, handle)
end)
{:ok, handles}
end
end,
# increase timeout because we may end up doing a lot of work here
60_000
),
{:ok, count} <- count_shapes(stack_id) do
{:ok, removed_handles, count}
end
end
def explain(stack_id) do
Connection.explain(stack_id)
:ok
end
@doc "Returns the number of pending writes in the buffer"
@spec pending_buffer_size(stack_id()) :: non_neg_integer()
def pending_buffer_size(stack_id) when is_stack_id(stack_id) do
WriteBuffer.pending_operations_count(stack_id)
end
defp raise_on_error!({:ok, result}, _action), do: result
defp raise_on_error!({:error, reason}, action) do
raise Error, error: reason, action: action
end
def statistics(stack_id) do
Statistics.current(stack_id)
end
end