Packages
electric
0.7.7
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.ex
defmodule Electric.Shapes do
alias Electric.Replication.LogOffset
alias Electric.ShapeCache.Storage
alias Electric.ShapeCache
alias Electric.Shapes.Shape
require Logger
@type shape_id :: Electric.ShapeCacheBehaviour.shape_id()
@doc """
Get snapshot for the shape ID
"""
def get_snapshot(config, shape_id) do
{shape_cache, opts} = Access.get(config, :shape_cache, {ShapeCache, []})
storage = shape_storage(config, shape_id)
if shape_cache.has_shape?(shape_id, opts) do
with :started <- shape_cache.await_snapshot_start(shape_id, opts) do
{:ok, Storage.get_snapshot(storage)}
end
else
{:error, "invalid shape_id #{inspect(shape_id)}"}
end
end
@doc """
Get stream of the log since a given offset
"""
def get_log_stream(config, shape_id, opts) do
{shape_cache, shape_cache_opts} = Access.get(config, :shape_cache, {ShapeCache, []})
offset = Access.get(opts, :since, LogOffset.before_all())
max_offset = Access.get(opts, :up_to, LogOffset.last())
storage = shape_storage(config, shape_id)
if shape_cache.has_shape?(shape_id, shape_cache_opts) do
Storage.get_log_stream(offset, max_offset, storage)
else
raise "Unknown shape: #{shape_id}"
end
end
@doc """
Get the shape that corresponds to this shape definition and return it along with the latest offset of the shape
"""
@spec get_shape(keyword(), Shape.t()) :: {shape_id(), LogOffset.t()}
def get_shape(config, shape_def) do
{shape_cache, opts} = Access.get(config, :shape_cache, {ShapeCache, []})
shape_cache.get_shape(shape_def, opts)
end
@doc """
Get or create a shape ID and return it along with the latest offset of the shape
"""
@spec get_or_create_shape_id(keyword(), Shape.t()) :: {shape_id(), LogOffset.t()}
def get_or_create_shape_id(config, shape_def) do
{shape_cache, opts} = Access.get(config, :shape_cache, {ShapeCache, []})
shape_cache.get_or_create_shape_id(shape_def, opts)
end
@doc """
Get the last exclusive offset of the chunk starting from the given offset
If `nil` is returned, chunk is not complete and the shape's latest offset should be used
"""
@spec get_chunk_end_log_offset(keyword(), shape_id(), LogOffset.t()) ::
LogOffset.t() | nil
def get_chunk_end_log_offset(config, shape_id, offset) do
storage = shape_storage(config, shape_id)
Storage.get_chunk_end_log_offset(offset, storage)
end
@doc """
Check whether the log has an entry for a given shape ID
"""
@spec has_shape?(keyword(), shape_id()) :: boolean()
def has_shape?(config, shape_id) do
{shape_cache, opts} = Access.get(config, :shape_cache, {ShapeCache, []})
shape_cache.has_shape?(shape_id, opts)
end
@doc """
Clean up all data (meta data and shape log + snapshot) associated with the given shape ID
"""
@spec clean_shape(shape_id(), keyword()) :: :ok
def clean_shape(shape_id, opts \\ []) do
{shape_cache, opts} = Access.get(opts, :shape_cache, {ShapeCache, []})
shape_cache.clean_shape(shape_id, opts)
:ok
end
@spec clean_shapes([shape_id()], keyword()) :: :ok
def clean_shapes(shape_ids, opts \\ []) do
{shape_cache, opts} = Access.get(opts, :shape_cache, {ShapeCache, []})
for shape_id <- shape_ids do
shape_cache.clean_shape(shape_id, opts)
end
:ok
end
defp shape_storage(config, shape_id) do
Storage.for_shape(shape_id, Access.fetch!(config, :storage))
end
end