Packages
electric
1.4.6
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.ShapeCache.ShapeStatus
alias Electric.Shapes.Shape
import Electric, only: [is_stack_id: 1, is_shape_handle: 1]
@type shape_handle :: Electric.shape_handle()
@type stack_id :: Electric.stack_id()
@doc """
Get the snapshot followed by the log.
"""
def get_merged_log_stream(stack_id, shape_handle, opts)
when is_shape_handle(shape_handle) and is_stack_id(stack_id) do
offset = Access.get(opts, :since, LogOffset.before_all())
max_offset = Access.get(opts, :up_to, LogOffset.last())
if ShapeCache.has_shape?(shape_handle, stack_id) do
with :started <- ShapeCache.await_snapshot_start(shape_handle, stack_id) do
storage = shape_storage(stack_id, shape_handle)
{:ok, Storage.get_log_stream(offset, max_offset, storage)}
end
else
# If we have a shape handle, but no shape, it means the shape was deleted. Send a 409
# and expect the client to retry - if the state of the world allows, it'll get a new handle.
{:error, Electric.Shapes.Api.Error.must_refetch()}
end
end
@doc """
Get the shape handle that corresponds to this shape definition and return it
"""
@spec fetch_handle_by_shape(stack_id(), Shape.t()) :: {:ok, shape_handle()} | :error
def fetch_handle_by_shape(stack_id, %Shape{} = shape_def) when is_stack_id(stack_id) do
ShapeCache.fetch_handle_by_shape(shape_def, stack_id)
end
@spec fetch_shape_by_handle(stack_id(), shape_handle()) :: Shape.t() | :error
def fetch_shape_by_handle(stack_id, shape_handle)
when is_shape_handle(shape_handle) and is_stack_id(stack_id) do
ShapeCache.fetch_shape_by_handle(shape_handle, stack_id)
end
@doc """
Cheaply validate that a shape handle matches the shape definition.
"""
def resolve_shape_handle(stack_id, shape_handle, %Shape{} = shape) do
ShapeCache.resolve_shape_handle(shape_handle, shape, stack_id)
end
@doc """
Get or create a shape handle and return it along with the latest offset of the shape
"""
@spec get_or_create_shape_handle(stack_id(), Shape.t()) :: {shape_handle(), LogOffset.t()}
def get_or_create_shape_handle(stack_id, shape_def) when is_stack_id(stack_id) do
ShapeCache.get_or_create_shape_handle(
shape_def,
stack_id,
otel_ctx: :otel_ctx.get_current()
)
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(stack_id(), shape_handle(), LogOffset.t()) :: LogOffset.t() | nil
def get_chunk_end_log_offset(stack_id, shape_handle, offset) do
storage = shape_storage(stack_id, shape_handle)
Storage.get_chunk_end_log_offset(offset, storage)
end
@doc """
Check whether the log has an entry for a given shape handle
"""
@spec has_shape?(stack_id(), shape_handle()) :: boolean()
def has_shape?(stack_id, shape_handle) do
ShapeCache.has_shape?(shape_handle, stack_id)
end
@doc """
Remove and clean up all data (meta data and shape log + snapshot) associated with
the given shape handle
"""
@spec clean_shape(stack_id(), shape_handle()) :: :ok
def clean_shape(stack_id, shape_handle) do
ShapeCache.clean_shape(shape_handle, stack_id)
:ok
end
@spec clean_shapes(stack_id(), [shape_handle()]) :: :ok
def clean_shapes(stack_id, shape_handles) do
for shape_handle <- shape_handles do
ShapeCache.clean_shape(shape_handle, stack_id)
end
:ok
end
@doc """
Wrap the writing of a snapshot to some Storage backend with the required
ShapeStatus update calls.
"""
@spec make_new_snapshot!(
Electric.Shapes.Querying.json_result_stream(),
Storage.shape_storage(),
stack_id(),
shape_handle()
) :: :ok | {:error, term()}
def make_new_snapshot!(stream, storage, stack_id, shape_handle) do
with :ok <- Storage.make_new_snapshot!(stream, storage),
:ok <- ShapeStatus.mark_snapshot_complete(stack_id, shape_handle) do
:ok
end
end
@spec mark_snapshot_started(Storage.shape_storage(), stack_id(), shape_handle()) ::
:ok | {:error, term()}
def mark_snapshot_started(storage, stack_id, shape_handle) do
with :ok <- Storage.mark_snapshot_as_started(storage) do
ShapeStatus.mark_snapshot_started(stack_id, shape_handle)
end
end
defp shape_storage(stack_id, shape_handle) do
Storage.for_shape(shape_handle, Storage.for_stack(stack_id))
end
def query_subset(handle, shape, subset, opts) do
Electric.Shapes.PartialModes.query_subset(handle, shape, subset, opts)
end
end