Skip to content

Postgres Changes can deliver frames encoded for another socket protocol #2309

Description

@hsusul

Describe the bug

Postgres Changes can send a frame encoded for another subscriber's negotiated socket protocol. When V1 and V2 subscribers have the same channel topic and change parameters, both receive the first subscriber's wire format. The endpoint explicitly supports both protocols.

To Reproduce

On upstream main 1921720d883bb13feb2a3bf6cae17c4dce42c095, save this as a local .exs file and run mix run --no-start /path/to/repro.exs after compiling dependencies. This exercises the actual CDC dispatcher and supported socket serializers without a database or timing dependency:

alias Extensions.PostgresCdcRls.MessageDispatcher

topic = "realtime:topic"
serializers = [Phoenix.Socket.V1.JSONSerializer, RealtimeWeb.Socket.V2Serializer]

subscriptions =
  Enum.map(serializers, fn serializer ->
    {self(), {:subscriber_fastlane, self(), serializer, [{"subscription", 1}], topic, true}}
  end)

payload = Jason.encode!(%{"schema" => "public", "table" => "todos", "type" => "INSERT", "record" => %{"id" => 1}})
:ok = MessageDispatcher.dispatch(subscriptions, self(), {"INSERT", payload, MapSet.new(["subscription"])})

for serializer <- serializers do
  receive do
    {:socket_push, :text, frame} ->
      IO.inspect({serializer, Jason.decode!(IO.iodata_to_binary(frame))})
  after
    1_000 -> raise "missing frame"
  end
end

Both decoded frames are JSON objects. Reverse the serializers and both are JSON arrays. The same failure occurs with the legacy API (true changed to false).

In production, connect two authorized clients using vsn=1.0.0 and vsn=2.0.0, join the same channel topic with identical Postgres Changes parameters, then insert a matching row. postgres_cdc_subscribe/2 retains each negotiated serializer and computes the same outgoing subscription ID from identical parameters, making the logical outgoing messages equal despite separate client subscriptions.

Expected behavior

Each transport receives a frame encoded by its own serializer: a JSON object for V1, and [join_ref, ref, topic, event, payload] for V2. Delivery must not depend on subscriber enumeration order.

Actual behavior

MessageDispatcher.broadcast_message/4 caches encoded frames using only the outgoing Phoenix.Socket.Broadcast as its key. An equal message for a different serializer hits that cache and bypasses the second serializer. The V2 subscriber receives a V1 object, or the V1 subscriber receives a V2 array.

System information

  • Realtime v2.140.3 / upstream main 1921720d883bb13feb2a3bf6cae17c4dce42c095
  • Reproduced locally with Elixir 1.20.4 / OTP 29
  • Supported serializers configured in lib/realtime_web/endpoint.ex

Additional context

This affects mixed-protocol subscribers, including client migrations: a matching database change reaches the transport but is not in the negotiated frame format, so clients cannot reliably decode their change events. No malformed input or race is required.

The cache should be scoped by both serializer and message, retaining reuse among subscribers with the same serializer. A regression using both real serializers fails before this change and passes after it, in both subscriber orders and both API modes. The existing cache-reuse test also verifies that identical same-serializer messages are encoded only once.

I searched open/closed issues and PRs, discussions, upstream branches and commit history for the dispatcher, serializer caching, mixed V1/V2 protocols and wire-format symptoms; I found no equivalent fix. #1615 concerns error caching in the separate Broadcast dispatcher, not this CDC cache key.

No activity

Activity on this issue will appear here.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions