Packages
electric
1.4.12
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/querying.ex
defmodule Electric.Shapes.Querying do
alias Electric.ShapeCache.LogChunker
alias Electric.Utils
alias Electric.Shapes.Shape
alias Electric.Shapes.Shape.SubqueryMoves
alias Electric.Telemetry.OpenTelemetry
@value_prefix SubqueryMoves.value_prefix()
@null_sentinel SubqueryMoves.null_sentinel()
def query_move_in(conn, stack_id, shape_handle, shape, {where, params}) do
table = Utils.relation_to_sql(shape.root_table)
{json_like_select, _} =
json_like_select(shape, %{"is_move_in" => true}, stack_id, shape_handle)
key_select = key_select(shape)
tag_select = make_tags(shape, stack_id, shape_handle) |> Enum.join(", ")
query =
Postgrex.prepare!(
conn,
table,
~s|SELECT #{key_select}, ARRAY[#{tag_select}]::text[], #{json_like_select} FROM #{table} WHERE #{where}|
)
Postgrex.stream(conn, query, params)
|> Stream.flat_map(& &1.rows)
end
def query_subset(conn, stack_id, shape_handle, shape, subset, headers \\ []) do
# When querying a subset, we select same columns as the base shape
table = Utils.relation_to_sql(shape.root_table)
where =
case {shape.where, subset.where} do
{nil, nil} ->
""
{nil, %{query: where}} ->
" WHERE " <> where
{%{query: where}, nil} ->
" WHERE " <> where
{%{query: base_where}, %{query: where}} ->
" WHERE " <> base_where <> " AND (" <> where <> ")"
end
order_by = if order_by = subset.order_by, do: " ORDER BY " <> order_by, else: ""
limit = if limit = subset.limit, do: " LIMIT #{limit}", else: ""
offset = if offset = subset.offset, do: " OFFSET #{offset}", else: ""
{json_like_select, params} = json_like_select(shape, headers, stack_id, shape_handle)
query =
Postgrex.prepare!(
conn,
table,
~s|SELECT #{json_like_select} FROM #{table} #{where} #{order_by} #{limit} #{offset}|
)
Postgrex.stream(conn, query, params)
|> Stream.flat_map(& &1.rows)
rescue
e in Postgrex.Error ->
case e.postgres do
# invalid_text_representation - e.g. invalid enum value
%{code: :invalid_text_representation, message: message} ->
# This is a type of error we expect, because we allow enums in subset where clauses
# even though we can't validate them fully.
raise __MODULE__.QueryError, message: message
_ ->
reraise e, __STACKTRACE__
end
end
defmodule QueryError do
defexception [:message]
end
@doc """
Streams the initial data for a shape. Query results are returned as a stream of JSON strings, as prepared on PostgreSQL.
"""
@type json_iodata :: iodata()
@type json_result_stream :: Enumerable.t(json_iodata())
@spec stream_initial_data(
DBConnection.t(),
String.t(),
String.t(),
Shape.t(),
non_neg_integer()
) ::
json_result_stream()
def stream_initial_data(
conn,
stack_id,
shape_handle,
shape,
chunk_bytes_threshold \\ LogChunker.default_chunk_size_threshold()
)
def stream_initial_data(_, _, _, %Shape{log_mode: :changes_only}, _chunk_bytes_threshold) do
[]
end
def stream_initial_data(
conn,
stack_id,
shape_handle,
%Shape{root_table: root_table} = shape,
chunk_bytes_threshold
) do
OpenTelemetry.with_span("shape_read.stream_initial_data", [], stack_id, fn ->
table = Utils.relation_to_sql(root_table)
where =
if not is_nil(shape.where), do: " WHERE " <> shape.where.query, else: ""
{json_like_select, params} = json_like_select(shape, [], stack_id, shape_handle)
query =
Postgrex.prepare!(conn, table, ~s|SELECT #{json_like_select} FROM #{table} #{where}|)
Postgrex.stream(conn, query, params)
|> Stream.flat_map(& &1.rows)
|> Stream.transform(0, fn [line], chunk_size ->
# Reason to add 1 byte to expected length is to account for `\n` breaks when the data is written.
case LogChunker.fit_into_chunk(
IO.iodata_length(line) + 1,
chunk_size,
chunk_bytes_threshold
) do
{:ok, new_chunk_size} ->
{[line], new_chunk_size}
{:threshold_exceeded, new_chunk_size} ->
{[line, :chunk_boundary], new_chunk_size}
end
end)
end)
end
defp key_select(%Shape{root_table: root_table, root_pk: pk_cols}) do
~s['#{escape_relation(root_table)}' || '/' || #{join_primary_keys(pk_cols)}]
end
# Converts a tag structure to something PG select can fill, but returns a list of separate strings for each tag
# - it's up to the caller to interpolate them into the query correctly
defp make_tags(%Shape{tag_structure: tag_structure}, stack_id, shape_handle) do
Enum.map(tag_structure, fn pattern ->
Enum.map(pattern, fn
column_name when is_binary(column_name) ->
col = pg_cast_column_to_text(column_name)
namespaced = pg_namespace_value_sql(col)
~s[md5('#{stack_id}#{shape_handle}' || #{namespaced})]
{:hash_together, columns} ->
column_parts =
Enum.map(columns, fn col_name ->
col = pg_cast_column_to_text(col_name)
~s['#{col_name}:' || #{pg_namespace_value_sql(col)}]
end)
~s[md5('#{stack_id}#{shape_handle}' || #{Enum.join(column_parts, " || ")})]
end)
|> Enum.join("|| '/' ||")
end)
end
defp json_like_select(
%Shape{
root_table: root_table,
selected_columns: columns
} = shape,
additional_headers,
stack_id,
shape_handle
) do
tags = make_tags(shape, stack_id, shape_handle)
key_part = build_key_part(shape)
value_part = build_value_part(columns)
headers_part = build_headers_part(root_table, additional_headers, tags)
# We're building a JSON string that looks like this:
#
# {
# "key": "\"public\".\"test_table\"/\"1\"",
# "value": {
# "id": "1",
# "name": "John Doe",
# "email": "john.doe@example.com",
# "nullable": null
# },
# "headers": {"operation": "insert", "relation": ["public", "test_table"]}
# }
query =
~s['{' || #{key_part} || ',' || #{value_part} || ',' || #{headers_part} || '}']
{query, []}
end
defp build_headers_part(rel, headers, tags) when is_list(headers),
do: build_headers_part(rel, Map.new(headers), tags)
defp build_headers_part({relation, table}, additional_headers, tags) do
headers = %{operation: "insert", relation: [relation, table]}
headers =
headers
|> Map.merge(additional_headers)
|> Jason.encode!()
|> Utils.escape_quotes(?')
headers =
if tags != [] do
"{" <> json = headers
tags = Enum.join(tags, ~s[ || '","' || ])
~s/{"tags":["' || #{tags} || '"],/ <> json
else
headers
end
~s['"headers":#{headers}']
end
defp build_key_part(shape) do
# Because relation part of the key is known at query building time, we can use $1 to inject escaped version of the relation
~s['"key":' || ] <> pg_escape_string_for_json(key_select(shape))
end
# This is a bespoke derivation of the record from its contents for Postgres but it must
# exactly match the algorithm implemented in `Electric.Replication.Changes.build_key/3`.
defp join_primary_keys(pk_cols) do
pk_cols
|> Enum.map(&pg_cast_column_to_text/1)
|> Enum.map(&~s['"' || replace(#{&1}, '/', '//') || '"'])
# NULL values are not allowed in PKs, but they are possible on pk-less tables where we consider all columns to be PKs
|> Enum.map(&~s[coalesce(#{&1}, '_')])
|> Enum.join(~s[ || '/' || ])
end
defp build_value_part(columns) do
column_parts = Enum.map(columns, &build_column_part/1)
~s['"value":{' || #{Enum.join(column_parts, " || ',' || ")} || '}']
end
defp build_column_part(column) do
escaped_name = escape_sql_json_interpolation(column)
escaped_value = escape_column_value(column)
# Since `||` returns NULL if any of the arguments is NULL, we need to use `coalesce` to handle NULL values
~s['"#{escaped_name}":' || #{pg_coalesce_json_string(escaped_value)}]
end
defp escape_sql_json_interpolation(str) do
str
|> String.replace(~S|"|, ~S|\"|)
|> String.replace(~S|'|, ~S|''|)
end
defp escape_relation(relation) do
relation |> Utils.relation_to_sql(true) |> String.replace(~S|'|, ~S|''|)
end
defp escape_column_value(column) do
column
|> pg_cast_column_to_text()
|> pg_escape_string_for_json()
|> pg_coalesce_json_string()
end
defp pg_cast_column_to_text(column), do: ~s["#{Utils.escape_quotes(column)}"::text]
defp pg_escape_string_for_json(str), do: ~s[to_json(#{str})::text]
defp pg_coalesce_json_string(str), do: ~s[coalesce(#{str} , 'null')]
# Generates SQL to namespace a value for tag hashing.
# This MUST produce identical output to SubqueryMoves.namespace_value/1 for
# the same input values, or Elixir-side and SQL-side tag computation will diverge.
defp pg_namespace_value_sql(col_sql) do
~s[CASE WHEN #{col_sql} IS NULL THEN '#{@null_sentinel}' ELSE '#{@value_prefix}' || #{col_sql} END]
end
end