Packages
electric
1.0.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.Telemetry.OpenTelemetry
require Logger
@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(), Shape.t(), non_neg_integer()) ::
json_result_stream()
def stream_initial_data(
conn,
stack_id,
%Shape{root_table: root_table} = shape,
chunk_bytes_threshold \\ LogChunker.default_chunk_size_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)
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 json_like_select(%Shape{
root_table: root_table,
selected_columns: columns,
root_pk: pk_cols
}) do
key_part = build_key_part(root_table, pk_cols)
value_part = build_value_part(columns)
headers_part = build_headers_part(root_table)
# 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(root_table) do
~s['"headers":{"operation":"insert","relation":#{build_relation_header(root_table)}}']
end
defp build_relation_header({schema, table}) do
~s'["#{escape_sql_json_interpolation(schema)}","#{escape_sql_json_interpolation(table)}"]'
end
defp build_key_part(root_table, pk_cols) do
pk_part = join_primary_keys(pk_cols)
# 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(~s['#{escape_relation(root_table)}' || '/"' || #{pk_part} || '"'])
end
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')]
end