defmodule LiveSync.Replication do @moduledoc """ Most of this code was sourced from https://github.com/josevalim/sync. """ use Postgrex.ReplicationConnection require Logger # TODO: Allow the publications to be passed as parameters def start_link(opts) do name = Keyword.get(opts, :name) if not is_atom(name) or name == nil do raise ArgumentError, "an atom :name is required when starting #{inspect(__MODULE__)}" end opts = Keyword.put_new(opts, :auto_reconnect, true) Postgrex.ReplicationConnection.start_link(__MODULE__, [], opts) end def subscribe(name) do Registry.register(LiveSync.Registry, name, []) end @doc """ Wait for connection. This is typically used by boot to make sure the replication is running to avoid unnecessary syncs. It accepts a maximum timeout. This function will exit if the server is not running. It returns `:ok` or `:timeout` otherwise. """ def wait_for_connection!(name, timeout \\ 5000) do ref = :erlang.monitor(:process, name, alias: :reply_demonitor) send(name, {:wait_for_connection, ref}) receive do :ok -> :ok {:DOWN, ^ref, _, _, reason} -> exit({reason, {__MODULE__, :wait_for_connection!, [name, timeout]}}) after timeout -> :timeout end end @doc """ Use to emulate a disconnection from the database. """ def disconnect(name) do send(name, :disconnect) end ## Callbacks @impl true def init(_opts) do # TODO: make devhub dynamic path = Application.app_dir(:devhub, "ebin") schemas = LiveSync.Sync |> Protocol.extract_impls([path]) |> Map.new(fn schema -> opts = LiveSync.opts(schema) fields = :fields |> schema.__schema__() |> Map.new(fn field -> {to_string(field), schema.__schema__(:type, field)} end) {opts[:table] || schema.__schema__(:source), %{module: schema, fields: fields}} end) state = %{ schemas: schemas, relations: %{}, # {:disconnected, []} | :connected | [operation] replication: {:disconnected, []} } {:ok, state} end @impl true def handle_disconnect(%{replication: replication} = state) do waiting = case replication do {:disconnected, waiting} -> waiting _ -> [] end {:noreply, %{state | replication: {:disconnected, waiting}}} end @impl true def handle_connect(state) do slot = random_slot_name() query = "CREATE_REPLICATION_SLOT #{slot} TEMPORARY LOGICAL pgoutput NOEXPORT_SNAPSHOT" {:query, query, state} end @impl true def handle_info({:wait_for_connection, ref}, state) do case state.replication do {:disconnected, waiting} -> {:noreply, %{state | replication: {:disconnected, [ref | waiting]}}} _ -> send(ref, :ok) {:noreply, state} end end def handle_info(:disconnect, _state) do {:disconnect, "user requested"} end def handle_info(_message, state) do {:noreply, state} end @impl true def handle_result([result], %{replication: {:disconnected, waiting}} = state) do %Postgrex.Result{ command: :create, columns: ["slot_name", "consistent_point", "snapshot_name", "output_plugin"], rows: [[slot, _lsn, nil, "pgoutput"]] } = result for ref <- waiting do send(ref, :ok) end query = "START_REPLICATION SLOT #{slot} LOGICAL 0/0 (proto_version '2', publication_names 'live_sync')" {:stream, query, [], %{state | replication: :connected}} end def handle_result(%Postgrex.Error{} = error, _state) do raise Exception.message(error) end @impl true # https://www.postgresql.org/docs/14/protocol-replication.html def handle_data(<>, state) do case rest do <> when state.replication == :connected -> handle_begin(state) <> when is_list(state.replication) -> handle_commit(lsn, state) <> when is_list(state.replication) -> handle_tuple_data(:insert, oid, count, tuple_data, state) <> when is_list(state.replication) -> handle_tuple_data(:update, oid, count, tuple_data, state) <> when is_list(state.replication) -> %{^oid => {schema, table, _columns}} = state.relation Logger.error( "A primary key of a row has been changed or its replica identity has been set to full, " <> "those operations are not currently supported by sync on #{schema}.#{table}" ) {:noreply, state} <> -> handle_relation(oid, rest, state) <> when is_list(state.replication) -> handle_tuple_data(:delete, oid, count, tuple_data, state) msg -> IO.inspect(msg, limit: :infinity) {:noreply, state} end end def handle_data(<>, state) do messages = case reply do 1 -> [<>] 0 -> [] end {:noreply, messages, state} end ## Decoding messages defp handle_begin(state) do {:noreply, %{state | replication: []}} end defp handle_relation(oid, rest, state) do [schema, rest] = :binary.split(rest, <<0>>) schema = if schema == "", do: "pg_catalog", else: schema [table, <<_replica_identity::8, count::16, rest::binary>>] = :binary.split(rest, <<0>>) columns = parse_columns(count, rest) state = put_in(state.relations[oid], {schema, table, columns}) {:noreply, state} end defp handle_tuple_data(kind, oid, count, tuple_data, state) do {schema, table, columns} = Map.fetch!(state.relations, oid) data = parse_tuple_data(count, columns, tuple_data) operation = %{schema: schema, table: table, op: kind, data: Map.new(data)} {:noreply, update_in(state.replication, &[operation | &1])} end defp handle_commit(_lsn, state) do data = state.replication |> Enum.filter(fn data -> Map.has_key?(state.schemas, data.table) end) |> Enum.map(fn %{data: data, table: table, op: op} -> %{module: module, fields: fields} = state.schemas[table] data = data |> Map.take(Map.keys(fields)) |> Map.new(fn {k, v} -> field = String.to_existing_atom(k) value = case fields[k] do :boolean -> v == "t" # TODO: handle errors type -> Ecto.Type.cast!(type, v) end {field, value} end) struct = Ecto.Schema.Loader.load_struct(module, nil, table) {op, Map.merge(struct, data)} end) Registry.dispatch(LiveSync.Registry, "live_sync", fn subscribed -> for {pid, _opts} <- subscribed, do: send(pid, {:live_sync, data}) end) {:noreply, %{state | replication: :connected}} end # TODO: if an entry has been soft-deleted, we could emit a special delete # instruction instead of sending the whole update. defp parse_tuple_data(0, [], <<>>), do: [] defp parse_tuple_data(count, [{name, _oid, _modifier} | columns], data) do case data do <> -> [{name, nil} | parse_tuple_data(count - 1, columns, rest)] # TODO: We are using text for convenience, we must set binary on the protocol <> -> [{name, value} | parse_tuple_data(count - 1, columns, rest)] <> -> raise "binary values not supported by sync" <> -> parse_tuple_data(count - 1, columns, rest) end end defp parse_columns(0, <<>>), do: [] defp parse_columns(count, <<_flags, rest::binary>>) do [name, <>] = :binary.split(rest, <<0>>) [{name, oid, modifier} | parse_columns(count - 1, rest)] end ## Helpers @epoch DateTime.to_unix(~U[2000-01-01 00:00:00Z], :microsecond) defp current_time, do: System.os_time(:microsecond) - @epoch defp random_slot_name do "phx_sync_" <> Base.encode32(:crypto.strong_rand_bytes(5), case: :lower) end end