Packages
electric
0.9.0
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/plug/serve_shape_plug.ex
defmodule Electric.Plug.ServeShapePlug do
use Plug.Builder, copy_opts_to_assign: :config
use Plug.ErrorHandler
# The halt/1 function is redefined further down below
import Plug.Conn, except: [halt: 1]
alias Electric.Plug.Utils
alias Electric.Shapes
alias Electric.Schema
alias Electric.Replication.LogOffset
alias Electric.Telemetry.OpenTelemetry
alias Plug.Conn
require Logger
# Aliasing for pattern matching
@before_all_offset LogOffset.before_all()
# Control messages
@up_to_date [Jason.encode!(%{headers: %{control: "up-to-date"}})]
@must_refetch Jason.encode!([%{headers: %{control: "must-refetch"}}])
@shape_definition_mismatch Jason.encode!(%{
message:
"The specified shape definition and handle do not match. " <>
"Please ensure the shape definition is correct or omit the shape handle from the request to obtain a new one."
})
defmodule Params do
use Ecto.Schema
import Ecto.Changeset
alias Electric.Replication.LogOffset
@primary_key false
embedded_schema do
field(:table, :string)
field(:offset, :string)
field(:handle, :string)
field(:live, :boolean, default: false)
field(:where, :string)
field(:columns, :string)
field(:shape_definition, :string)
field(:replica, Ecto.Enum, values: [:default, :full], default: :default)
end
def validate(params, opts) do
%__MODULE__{}
|> cast(params, __schema__(:fields) -- [:shape_definition],
message: fn _, _ -> "must be %{type}" end
)
|> validate_required([:table, :offset])
|> cast_offset()
|> cast_columns()
|> validate_handle_with_offset()
|> validate_live_with_offset()
|> cast_root_table(opts)
|> apply_action(:validate)
|> case do
{:ok, params} ->
{:ok, Map.from_struct(params)}
{:error, changeset} ->
{:error,
Ecto.Changeset.traverse_errors(changeset, fn {msg, opts} ->
Regex.replace(~r"%{(\w+)}", msg, fn _, key ->
opts |> Keyword.get(String.to_existing_atom(key), key) |> to_string()
end)
end)}
end
end
def cast_offset(%Ecto.Changeset{valid?: false} = changeset), do: changeset
def cast_offset(%Ecto.Changeset{} = changeset) do
offset = fetch_change!(changeset, :offset)
case LogOffset.from_string(offset) do
{:ok, offset} ->
put_change(changeset, :offset, offset)
{:error, message} ->
add_error(changeset, :offset, message)
end
end
def cast_columns(%Ecto.Changeset{valid?: false} = changeset), do: changeset
def cast_columns(%Ecto.Changeset{} = changeset) do
case fetch_field!(changeset, :columns) do
nil ->
changeset
columns ->
case Electric.Plug.Utils.parse_columns_param(columns) do
{:ok, parsed_cols} -> put_change(changeset, :columns, parsed_cols)
{:error, reason} -> add_error(changeset, :columns, reason)
end
end
end
def validate_handle_with_offset(%Ecto.Changeset{valid?: false} = changeset),
do: changeset
def validate_handle_with_offset(%Ecto.Changeset{} = changeset) do
offset = fetch_change!(changeset, :offset)
if offset == LogOffset.before_all() do
changeset
else
validate_required(changeset, [:handle], message: "can't be blank when offset != -1")
end
end
def validate_live_with_offset(%Ecto.Changeset{valid?: false} = changeset), do: changeset
def validate_live_with_offset(%Ecto.Changeset{} = changeset) do
offset = fetch_change!(changeset, :offset)
if offset != LogOffset.before_all() do
changeset
else
validate_exclusion(changeset, :live, [true], message: "can't be true when offset == -1")
end
end
def cast_root_table(%Ecto.Changeset{valid?: false} = changeset, _), do: changeset
def cast_root_table(%Ecto.Changeset{} = changeset, opts) do
table = fetch_change!(changeset, :table)
where = fetch_field!(changeset, :where)
columns = get_change(changeset, :columns, nil)
replica = fetch_field!(changeset, :replica)
case Shapes.Shape.new(
table,
opts ++ [where: where, columns: columns, replica: replica]
) do
{:ok, result} ->
put_change(changeset, :shape_definition, result)
{:error, {field, reasons}} ->
Enum.reduce(List.wrap(reasons), changeset, fn
{message, keys}, changeset ->
add_error(changeset, field, message, keys)
message, changeset when is_binary(message) ->
add_error(changeset, field, message)
end)
end
end
end
plug :fetch_query_params
# start_telemetry_span needs to always be the first plug after fetching query params.
plug :start_telemetry_span
plug :put_resp_content_type, "application/json"
plug :validate_query_params
plug :load_shape_info
plug :put_schema_header
# We're starting listening as soon as possible to not miss stuff that was added since we've
# asked for last offset
plug :listen_for_new_changes
plug :determine_log_chunk_offset
plug :determine_up_to_date
plug :generate_etag
plug :validate_and_put_etag
plug :put_resp_cache_headers
plug :serve_log_or_snapshot
# end_telemetry_span needs to always be the last plug here.
plug :end_telemetry_span
defp validate_query_params(%Conn{} = conn, _) do
Logger.info("Query String: #{conn.query_string}")
all_params =
Map.merge(conn.query_params, conn.path_params)
|> Map.update("live", "false", &(&1 != "false"))
case Params.validate(all_params, inspector: conn.assigns.config[:inspector]) do
{:ok, params} ->
%{conn | assigns: Map.merge(conn.assigns, params)}
{:error, error_map} ->
conn
|> send_resp(400, Jason.encode_to_iodata!(error_map))
|> halt()
end
end
defp load_shape_info(%Conn{} = conn, _) do
OpenTelemetry.with_span("shape_get.plug.load_shape_info", [], fn ->
shape_info = get_or_create_shape_handle(conn.assigns)
handle_shape_info(conn, shape_info)
end)
end
# No handle is provided so we can get the existing one for this shape
# or create a new shape if it does not yet exist
defp get_or_create_shape_handle(%{shape_definition: shape, config: config, handle: nil}) do
Shapes.get_or_create_shape_handle(config, shape)
end
# A shape handle is provided so we need to return the shape that matches the shape handle and the shape definition
defp get_or_create_shape_handle(%{shape_definition: shape, config: config}) do
Shapes.get_shape(config, shape)
end
defp handle_shape_info(
%Conn{assigns: %{shape_definition: shape, config: config, handle: shape_handle}} =
conn,
nil
) do
# There is no shape that matches the shape definition (because shape info is `nil`)
if shape_handle != nil && Shapes.has_shape?(config, shape_handle) do
# but there is a shape that matches the shape handle
# thus the shape handle does not match the shape definition
# and we return a 400 bad request status code
conn
|> send_resp(400, @shape_definition_mismatch)
|> halt()
else
# The shape handle does not exist or no longer exists
# e.g. it may have been deleted.
# Hence, create a new shape for this shape definition
# and return a 409 with a redirect to the newly created shape.
# (will be done by the recursive `handle_shape_info` call)
shape_info = Shapes.get_or_create_shape_handle(config, shape)
handle_shape_info(conn, shape_info)
end
end
defp handle_shape_info(
%Conn{assigns: %{handle: shape_handle}} = conn,
{active_shape_handle, last_offset}
)
when is_nil(shape_handle) or shape_handle == active_shape_handle do
# We found a shape that matches the shape definition
# and the shape has the same ID as the shape handle provided by the user
conn
|> assign(:active_shape_handle, active_shape_handle)
|> assign(:last_offset, last_offset)
|> put_resp_header("electric-handle", active_shape_handle)
end
defp handle_shape_info(
%Conn{assigns: %{config: config, handle: shape_handle, table: table}} = conn,
{active_shape_handle, _}
) do
if Shapes.has_shape?(config, shape_handle) do
# The shape with the provided ID exists but does not match the shape definition
# otherwise we would have found it and it would have matched the previous function clause
conn
|> send_resp(400, @shape_definition_mismatch)
|> halt()
else
# The requested shape_handle is not found, returns 409 along with a location redirect for clients to
# re-request the shape from scratch with the new shape id which acts as a consistent cache buster
# e.g. GET /v1/shape?table={root_table}&handle={new_shape_handle}&offset=-1
# TODO: discuss returning a 307 redirect rather than a 409, the client
# will have to detect this and throw out old data
conn
|> put_resp_header("electric-handle", active_shape_handle)
|> put_resp_header(
"location",
"#{conn.request_path}?table=#{table}&handle=#{active_shape_handle}&offset=-1"
)
|> send_resp(409, @must_refetch)
|> halt()
end
end
defp schema(shape) do
shape.table_info
|> Map.fetch!(shape.root_table)
|> Map.fetch!(:columns)
|> Schema.from_column_info()
|> Jason.encode!()
end
# Only adds schema header when not in live mode
defp put_schema_header(conn, _) when not conn.assigns.live do
shape = conn.assigns.shape_definition
put_resp_header(conn, "electric-schema", schema(shape))
end
defp put_schema_header(conn, _), do: conn
# If chunk offsets are available, use those instead of the latest available offset
# to optimize for cache hits and response sizes
defp determine_log_chunk_offset(%Conn{assigns: assigns} = conn, _) do
%{config: config, active_shape_handle: shape_handle, offset: offset} =
assigns
chunk_end_offset =
Shapes.get_chunk_end_log_offset(config, shape_handle, offset) ||
assigns.last_offset
conn
|> assign(:chunk_end_offset, chunk_end_offset)
|> put_resp_header("electric-offset", "#{chunk_end_offset}")
end
defp determine_up_to_date(
%Conn{
assigns: %{
offset: offset,
chunk_end_offset: chunk_end_offset,
last_offset: last_offset
}
} = conn,
_
) do
# The log can't be up to date if the last_offset is not the actual end.
# Also if client is requesting the start of the log, we don't set `up-to-date`
# here either as we want to set a long max-age on the cache-control.
if LogOffset.compare(chunk_end_offset, last_offset) == :lt or offset == @before_all_offset do
conn
|> assign(:up_to_date, [])
# header might have been added on first pass but no longer valid
# if listening to live changes and an incomplete chunk is formed
|> delete_resp_header("electric-up-to-date")
else
conn
|> assign(:up_to_date, [@up_to_date])
|> put_resp_header("electric-up-to-date", "")
end
end
defp generate_etag(%Conn{} = conn, _) do
%{
offset: offset,
active_shape_handle: active_shape_handle,
chunk_end_offset: chunk_end_offset
} = conn.assigns
conn
|> assign(
:etag,
"#{active_shape_handle}:#{offset}:#{chunk_end_offset}"
)
end
defp validate_and_put_etag(%Conn{} = conn, _) do
if_none_match =
get_req_header(conn, "if-none-match")
|> Enum.flat_map(&String.split(&1, ","))
|> Enum.map(&String.trim/1)
|> Enum.map(&String.trim(&1, ~S|"|))
cond do
conn.assigns.etag in if_none_match ->
conn
|> send_resp(304, "")
|> halt()
not conn.assigns.live ->
put_resp_header(conn, "etag", conn.assigns.etag)
true ->
conn
end
end
# If the offset is -1, set a 1 week max-age, 1 hour s-maxage (shared cache) and 1 month stale-while-revalidate
# We want private caches to cache the initial offset for a long time but for shared caches to frequently revalidate
# so they're serving a fairly fresh copy of the initials shape log.
defp put_resp_cache_headers(%Conn{assigns: %{offset: @before_all_offset}} = conn, _),
do:
conn
|> put_resp_header(
"cache-control",
"public, max-age=604800, s-maxage=3600, stale-while-revalidate=2629746"
)
# For live requests we want shorrt cache lifetimes and to update the live cursor
defp put_resp_cache_headers(%Conn{assigns: %{live: true}} = conn, _),
do:
conn
|> put_resp_header(
"cache-control",
"public, max-age=5, stale-while-revalidate=5"
)
|> put_resp_header(
"electric-cursor",
conn.assigns.config[:long_poll_timeout]
|> Utils.get_next_interval_timestamp(conn.query_params["cursor"])
|> Integer.to_string()
)
# For all other requests use the configured cache lifetimes
defp put_resp_cache_headers(%Conn{assigns: %{config: config, live: false}} = conn, _),
do:
conn
|> put_resp_header(
"cache-control",
"public, max-age=#{config[:max_age]}, stale-while-revalidate=#{config[:stale_age]}"
)
# If offset is -1, we're serving a snapshot
defp serve_log_or_snapshot(%Conn{assigns: %{offset: @before_all_offset}} = conn, _) do
OpenTelemetry.with_span("shape_get.plug.serve_snapshot", [], fn -> serve_snapshot(conn) end)
end
# Otherwise, serve log since that offset
defp serve_log_or_snapshot(conn, _) do
OpenTelemetry.with_span("shape_get.plug.serve_shape_log", [], fn -> serve_shape_log(conn) end)
end
defp serve_snapshot(
%Conn{
assigns: %{
chunk_end_offset: chunk_end_offset,
active_shape_handle: shape_handle,
up_to_date: maybe_up_to_date
}
} = conn
) do
case Shapes.get_snapshot(conn.assigns.config, shape_handle) do
{:ok, {offset, snapshot}} ->
log =
Shapes.get_log_stream(conn.assigns.config, shape_handle,
since: offset,
up_to: chunk_end_offset
)
[snapshot, log, maybe_up_to_date]
|> Stream.concat()
|> to_json_stream()
|> Stream.chunk_every(500)
|> send_stream(conn, 200)
{:error, reason} ->
error_msg = "Could not serve a snapshot because of #{inspect(reason)}"
Logger.warning(error_msg)
OpenTelemetry.record_exception(error_msg)
{status_code, message} =
if match?(%DBConnection.ConnectionError{reason: :queue_timeout}, reason),
do: {429, "Could not establish connection to database - try again later"},
else: {500, "Failed creating or fetching the snapshot"}
send_resp(
conn,
status_code,
Jason.encode_to_iodata!(%{error: message})
)
end
end
defp serve_shape_log(
%Conn{
assigns: %{
offset: offset,
chunk_end_offset: chunk_end_offset,
active_shape_handle: shape_handle,
up_to_date: maybe_up_to_date
}
} = conn
) do
log =
Shapes.get_log_stream(conn.assigns.config, shape_handle,
since: offset,
up_to: chunk_end_offset
)
if Enum.take(log, 1) == [] and conn.assigns.live do
conn
|> assign(:ot_is_immediate_response, false)
|> hold_until_change(shape_handle)
else
[log, maybe_up_to_date]
|> Stream.concat()
|> to_json_stream()
|> Stream.chunk_every(500)
|> send_stream(conn, 200)
end
end
@json_list_start "["
@json_list_end "]"
@json_item_separator ","
defp to_json_stream(items) do
Stream.concat([
[@json_list_start],
Stream.intersperse(items, @json_item_separator),
[@json_list_end]
])
end
defp send_stream(stream, conn, status) do
conn = send_chunked(conn, status)
{conn, bytes_sent} =
Enum.reduce_while(stream, {conn, 0}, fn chunk, {conn, bytes_sent} ->
chunk_size = IO.iodata_length(chunk)
OpenTelemetry.with_span("shape_get.plug.stream_chunk", [chunk_size: chunk_size], fn ->
case chunk(conn, chunk) do
{:ok, conn} ->
{:cont, {conn, bytes_sent + chunk_size}}
{:error, "closed"} ->
error_str = "Connection closed unexpectedly while streaming response"
conn = assign(conn, :error_str, error_str)
{:halt, {conn, bytes_sent}}
{:error, reason} ->
error_str = "Error while streaming response: #{inspect(reason)}"
Logger.error(error_str)
conn = assign(conn, :error_str, error_str)
{:halt, {conn, bytes_sent}}
end
end)
end)
assign(conn, :streaming_bytes_sent, bytes_sent)
end
defp listen_for_new_changes(%Conn{} = conn, _) when not conn.assigns.live, do: conn
defp listen_for_new_changes(%Conn{assigns: assigns} = conn, _) do
# Only start listening when we know there is a possibility that nothing is going to be returned
if LogOffset.compare(assigns.offset, assigns.last_offset) != :lt do
shape_handle = assigns.handle
ref = make_ref()
registry = conn.assigns.config[:registry]
Registry.register(registry, shape_handle, ref)
Logger.debug("Client #{inspect(self())} is registered for changes to #{shape_handle}")
assign(conn, :new_changes_ref, ref)
else
conn
end
end
def hold_until_change(conn, shape_handle) do
long_poll_timeout = conn.assigns.config[:long_poll_timeout]
Logger.debug("Client #{inspect(self())} is waiting for changes to #{shape_handle}")
ref = conn.assigns.new_changes_ref
receive do
{^ref, :new_changes, latest_log_offset} ->
# Stream new log since currently "held" offset
conn
|> assign(:last_offset, latest_log_offset)
|> assign(:chunk_end_offset, latest_log_offset)
# update last offset header
|> put_resp_header("electric-offset", "#{latest_log_offset}")
|> determine_up_to_date([])
|> serve_shape_log()
{^ref, :shape_rotation} ->
# We may want to notify the client better that the shape handle had changed, but just closing the response
# and letting the client handle it on reconnection is good enough.
conn
|> assign(:ot_is_shape_rotated, true)
|> assign(:ot_is_empty_response, true)
|> send_resp(200, ["[", @up_to_date, "]"])
after
# If we timeout, return an empty body and 204 as there's no response body.
long_poll_timeout ->
conn
|> assign(:ot_is_long_poll_timeout, true)
|> assign(:ot_is_empty_response, true)
|> send_resp(204, ["[", @up_to_date, "]"])
end
end
defp open_telemetry_attrs(%Conn{assigns: assigns} = conn) do
shape_handle =
if is_struct(conn.query_params, Plug.Conn.Unfetched) do
assigns[:active_shape_handle] || assigns[:shape_handle]
else
conn.query_params["handle"] || assigns[:active_shape_handle] || assigns[:shape_handle]
end
maybe_up_to_date = if up_to_date = assigns[:up_to_date], do: up_to_date != []
Electric.Plug.Utils.common_open_telemetry_attrs(conn)
|> Map.merge(%{
"shape.handle" => shape_handle,
"shape.where" => assigns[:where],
"shape.root_table" => assigns[:table],
"shape.definition" => assigns[:shape_definition],
"shape.replica" => assigns[:replica],
"shape_req.is_live" => assigns[:live],
"shape_req.offset" => assigns[:offset],
"shape_req.is_shape_rotated" => assigns[:ot_is_shape_rotated] || false,
"shape_req.is_long_poll_timeout" => assigns[:ot_is_long_poll_timeout] || false,
"shape_req.is_empty_response" => assigns[:ot_is_empty_response] || false,
"shape_req.is_immediate_response" => assigns[:ot_is_immediate_response] || true,
"shape_req.is_cached" => if(conn.status, do: conn.status == 304),
"shape_req.is_error" => if(conn.status, do: conn.status >= 400),
"shape_req.is_up_to_date" => maybe_up_to_date
})
end
#
### Telemetry
#
# Below, OpentelemetryTelemetry does the heavy lifting of setting up the span context in the
# current Elixir process to correctly attribute subsequent calls to OpenTelemetry.with_span()
# in this module as descendants of the root span, as they are all invoked in the same process
# unless a new process is spawned explicitly.
# Start the root span for the shape request, serving as an ancestor for any subsequent
# sub-span.
defp start_telemetry_span(conn, _) do
OpentelemetryTelemetry.start_telemetry_span(OpenTelemetry, "Plug_shape_get", %{}, %{})
add_span_attrs_from_conn(conn)
conn
end
# Assign root span attributes based on the latest state of Plug.Conn and end the root span.
#
# We want to have all the relevant HTTP and shape request attributes on the root span. This
# is the place to assign them because we keep this plug last in the "plug pipeline" defined
# in this module.
defp end_telemetry_span(conn, _ \\ nil) do
add_span_attrs_from_conn(conn)
OpentelemetryTelemetry.end_telemetry_span(OpenTelemetry, %{})
conn
end
defp add_span_attrs_from_conn(conn) do
conn
|> open_telemetry_attrs()
|> OpenTelemetry.add_span_attributes()
end
# This overrides Plug.Conn.halt/1 (which is deliberately "unimported" at the top of this
# module) so that we can record the response status in the OpenTelemetry span for this
# request.
defp halt(conn) do
conn
|> end_telemetry_span()
|> Plug.Conn.halt()
end
@impl Plug.ErrorHandler
def handle_errors(conn, error) do
OpenTelemetry.record_exception(error.kind, error.reason, error.stack)
error_str = Exception.format(error.kind, error.reason)
conn
|> assign(:error_str, error_str)
|> end_telemetry_span()
conn
end
end