Skip to content

fix: pg changes committed behind a synchronous standby - #2272

Merged
leandrocp merged 12 commits into
mainfrom
repro-sync-commit-event-loss
Sep 29, 2026
Merged

leandrocp merged 12 commits into
mainfrom
repro-sync-commit-event-loss

Conversation

@leandrocp

@leandrocp leandrocp commented Sep 23, 2026 •

Copy link
Copy Markdown
Member

PG Changes can drop an RLS authorized INSERT whose COMMIT is waiting for a synchronous standby.

Reported by @kabochya (Multigres) in https://github.com/kabochya/realtime/tree/multigres-visibility-gate and built on top of his findings.

synchronous_standby_names defaults to empty in PG and that has been the config used since forever which doesn't add any wait time and thus no message loss. Any non-empty value makes COMMIT wait for an ack and messages can be lost in that window.

Multigres does it automatically and any Postgres server with non-empty value would face the same message loss problem.

Solutions considered

Option Source How it reads the slot
Only get main Reads with pg_logical_slot_get_changes. No fix. This is what runs without a synchronous standby.
peek+get https://github.com/supabase/realtime/tree/4a9c8687d13e3347749bace8e8d4e3576d35ceb7 Peeks with transaction markers, then reads up to the last settled commit.
peek+advance N/A Peeks with transaction markers, returns the settled changes, then advances the slot. Runs as one query (1).
peek+get (no markers, optimized) This branch Peeks without markers, then reads only the changes in front of the first unsettled transaction.

(1) Run as one query to avoid buffering which would use arrays to store data in-memory capped at 1GB, so a batch that decodes past that fails on every poll and the slot never advances.

Tests

main peek+get peek+advance peek+get (no markers)
1. Lost on Multigres, per run of 500 writes 40, 1, 0, 5, 2 0 0 0
2. Delivery tests, Postgres 14 to 17 and Multigres - pass pass pass
3. Parked commit deferred on OrioleDB - no no no
4. Polls to reach a change behind 1,000 unpublished transactions 1 21 21 1
5. Full poll with apply_rls, vs main - -1% to +19% -2% to +16% -1% to +17%
6. Settled read, 1,000 to 100,000 changes, vs main - +162% to +439% +93% to +220% +86% to +176%
7. Settled read, 1 and 100 changes 1.3 to 3.0 ms 2.7 to 7.2 ms 2.5 to 6.0 ms 3.1 to 6.3 ms
8. One transaction that decodes to 1.1 GB pass pass pass pass
  1. pg_changes_delivery.exs, four writers, RLS policy that reads the row.
  2. A parked commit is deferred then delivered once, a settled commit ahead of a parked one is
    delivered alone, an in-flight transaction is deferred, a denied change is consumed, 200
    concurrent writes arrive on Multigres, and commit order matches main.
  3. An OrioleDB-only transaction has no Postgres xid. Decoding reports values such as 480, 608, and
    4286579968, which the xip/xmax check reads as settled.
  4. wal2json emits B/C for every decoded transaction, and upto_nchanges counts them. At
    max_changes 100 a poll covers about 50 transactions, and an empty poll idles 500 ms
    (poll_interval_ms * 5), which caps a stream of unpublished writes near 100 txn/s.
  5. list_changes.exs, with apply_rls, which dominates.
  6. The SQL function alone, one long-lived connection per option owning a temporary slot,
    restart_lsn caught up per rep, median of 10 runs on 17.6 and five elsewhere.
  7. Same as 6, at 1 and 100 changes.
  8. 1,100 rows of 1 MB in one transaction, 14 MB of WAL. Tuplestores spill to disk.

Trade-offs

peek+get peek+advance peek+get (no markers)
Decoding passes two full one full, one fast-forward two full, or one plus a fast-forward when the peek is empty
Markers count against max_changes yes yes no
JSON parse every row every row logical messages only
Snapshot for the check after the peek the statement's, before the peek after the peek
Relies on get stops after the record that reaches upto_lsn (logicalfuncs.c), and a C marker's lsn is the end of its commit record (logical.c), so the read includes that commit. The docs only say "commit prior to the specified LSN". A C marker's lsn is the end of its commit record (logical.c), and pg_replication_slot_advance confirms exactly that LSN (logical.c, docs). get stops once returned_rows reaches upto_nchanges (17.6, unchanged since 9.4), while the docs say it stops when the count exceeds it. Logical messages are found by wal2json v2's "action":"M" and transactional fields (wal2json 2.6).

Proposal

Adopt peek+get (no markers, optimized) solution.

Introduce realtime.list_changes_sync which is a drop-in replacement of realtime.list_changes called only when it needs to await for standbys (when synchronous_standby_names is set). How it works:

  1. Pin the range with upto := pg_current_wal_flush_lsn(), so the peek and the read see the same WAL.
  2. Peek up to upto without consuming. The peek uses the caller's own wal2json options, with no transaction markers, so max_changes counts exactly what pg_logical_slot_get_changes counts. It aggregates into one row per transaction, in commit order: the xid and the position of its first change. It skips non-transactional logical messages, which carry the xid of whatever transaction wrote them.
  3. Take a snapshot and find the first transaction still in flight: in xip, or at or past xmax.
  4. Consume with pg_logical_slot_get_changes(slot, upto, n). Here n is the number of changes in front of that transaction, or max_changes if everything settled. The read stops right after the commit that reaches n, so the in-flight transaction and everything after it stay in the slot for the next poll. If the peek returned nothing, it advances the slot to upto instead of reading again.

Call the function or not

No feature flag because this affects Multires and OrioleDB so I decided to probe the setting synchronous_standby_names and "fork" based on the DB:

  • synchronous_standby_names OFF: call current list_changes
  • synchronous_standby_names ON + OrioledDB: warning, won't work (need to follow-up)
  • synchronous_standby_names ON + Postgres ou Multigres: call list_changes_sync

Bench

dev/bench/list_changes.exs

Batch list_changes list_changes_sync delta
1 change 4.40 ms 4.50 ms +2%
100 changes 22.74 ms 23.70 ms +4%
1000 changes 182.14 ms 190.53 ms +5%
10000 changes 1.73 s 1.80 s +4%

Higher delta with less changes is because any variation in timing cause more % with less messages as a consequence.

Delivery

dev/bench/pg_changes_delivery.exs

Rows lost, out of 500.

run 1 run 2 run 3 run 4 run 5
main 121 162 3 249 265
branch 0 0 0 0 0

Note that on single-node Postgres without synchronous_standby_names there's no loss because it doesn't need to wait for a standby.

Notes

  • Our current tests on main doesn't exercise that scenario.
  • Note this only impacts Postgres Changes.

Refs

Closes REAL-1128

@blacksmith-sh

This comment has been minimized.

@leandrocp
leandrocp force-pushed the repro-sync-commit-event-loss branch 3 times, most recently from b128689 to 0cd6a8a Compare September 23, 2026 22:19
@blacksmith-sh

This comment has been minimized.

@leandrocp
leandrocp force-pushed the repro-sync-commit-event-loss branch 2 times, most recently from 78b74d3 to ed1c2dc Compare September 24, 2026 11:10
@blacksmith-sh

This comment has been minimized.

@leandrocp
leandrocp force-pushed the repro-sync-commit-event-loss branch 2 times, most recently from de4d0fe to f4d5567 Compare September 24, 2026 12:48
@github-actions

github-actions Bot commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

CRAP Score Report

Summary: files=242 functions=1333 scored=1333 worst_score=15944.76

** (Mix) CRAP threshold failed: max_score=30.00
High scores: 29
  lib/extensions/postgres_cdc_rls/subscription_manager.ex Extensions.PostgresCdcRls.SubscriptionManager.handle_info/2 score=37.16
  lib/realtime/adapters/postgres/oid_database.ex Realtime.Adapters.Postgres.OidDatabase.name_for_type_id/1 score=15944.76
  lib/realtime/application.ex Realtime.Application.setup_region_mapping/0 score=47.11
  lib/realtime/nodes.ex Realtime.Nodes.default_region_mapping/1 score=157.20
  lib/realtime/operations.ex Realtime.Operations.rebalance/0 score=42.00
  lib/realtime/operations.ex Realtime.Operations.kill_connections_to_tenant_id/2 score=90.00
  lib/realtime/tenants/connect.ex Realtime.Tenants.Connect.handle_info/2 score=41.41
  lib/realtime_web/channels/realtime_channel.ex RealtimeWeb.RealtimeChannel.handle_info/2 score=51.93
  lib/realtime_web/channels/realtime_channel.ex RealtimeWeb.RealtimeChannel.handle_in/3 score=62.32
  lib/realtime_web/dashboard/feature_flags.ex RealtimeWeb.Dashboard.FeatureFlags.handle_event/3 score=112.70
  lib/realtime_web/dashboard/node_info.ex RealtimeWeb.Dashboard.NodeInfo.fetch_node_data/2 score=42.00
  lib/realtime_web/dashboard/recon_trace.ex RealtimeWeb.Dashboard.ReconTrace.handle_event/3 score=552.00
  lib/realtime_web/dashboard/recon_trace.ex RealtimeWeb.Dashboard.ReconTrace.handle_info/2 score=210.00
  lib/realtime_web/dashboard/recon_trace.ex RealtimeWeb.Dashboard.ReconTrace.render_value/2 score=42.00
  lib/realtime_web/dashboard/recon_trace.ex RealtimeWeb.Dashboard.ReconTrace.load_module_functions/1 score=42.00
  lib/realtime_web/dashboard/recon_trace.ex RealtimeWeb.Dashboard.ReconTrace.parse_and_start/2 score=110.00
  lib/realtime_web/dashboard/recon_trace.ex RealtimeWeb.Dashboard.ReconTrace.format_value/1 score=552.00
  lib/realtime_web/dashboard/recon_trace.ex RealtimeWeb.Dashboard.ReconTrace.sort_entries/2 score=42.00
  lib/realtime_web/dashboard/sql_inspector.ex RealtimeWeb.Dashboard.SqlInspector.handle_event/3 score=72.00
  lib/realtime_web/dashboard/sql_inspector.ex RealtimeWeb.Dashboard.SqlInspector.execute_read_only/1 score=72.00
  lib/realtime_web/dashboard/sql_inspector.ex RealtimeWeb.Dashboard.SqlInspector.mask_sensitive_columns/1 score=42.00
  lib/realtime_web/dashboard/sql_inspector.ex RealtimeWeb.Dashboard.SqlInspector.compare_cells/2 score=42.00
  lib/realtime_web/dashboard/sql_inspector.ex RealtimeWeb.Dashboard.SqlInspector.format_cell/1 score=72.00
  lib/realtime_web/dashboard/tenant_migrations.ex RealtimeWeb.Dashboard.TenantMigrations.handle_info/2 score=156.00
  lib/realtime_web/live/components.ex RealtimeWeb.Components.input/1 score=90.00
  lib/realtime_web/live/inspector_live/conn_component.ex RealtimeWeb.InspectorLive.ConnComponent.handle_event/3 score=53.83
  lib/realtime_web/live/inspector_live/event_log_component.ex RealtimeWeb.InspectorLive.EventLogComponent.category_variant/1 score=35.00
  lib/realtime_web/live/inspector_live/event_log_component.ex RealtimeWeb.InspectorLive.EventLogComponent.event_label/1 score=76.13
  lib/realtime_web/live/status_live/index.ex RealtimeWeb.StatusLive.Index.handle_event/3 score=42.00

@coveralls

coveralls commented Sep 24, 2026 •

Copy link
Copy Markdown

Coverage Status

Coverage is 91.829% — repro-sync-commit-event-loss into main. No base build found for main.

@leandrocp
leandrocp force-pushed the repro-sync-commit-event-loss branch from c43dee0 to 687b2ca Compare September 24, 2026 13:14
@leandrocp
leandrocp force-pushed the repro-sync-commit-event-loss branch 3 times, most recently from fb7c9d0 to e85dd62 Compare September 24, 2026 15:32
exit;
end if;

boundary := commit_lsns[i];

@kabochya kabochya Sep 24, 2026 •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If we are concerned about decoding twice, we can save the extra pg_logical_slot_get_changes below by buffering all visible changes in the for loop, return them, and do a pg_replication_slot_advance to the last visible txn's LSN

@leandrocp leandrocp Sep 24, 2026 •

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hey @kabochya let me measure, it might work.

/edit the gains are minimal and adds more complexity so I'm not sure it's worth it 🤔

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That's fine too; mostly just raising this approach for completeness.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Oh yeah ideas are welcome, we should measure and test all 👍🏻

@leandrocp
leandrocp force-pushed the repro-sync-commit-event-loss branch 2 times, most recently from 829c884 to 1b1959c Compare September 24, 2026 19:08
@leandrocp
leandrocp force-pushed the repro-sync-commit-event-loss branch from 1b1959c to b1d9e88 Compare September 24, 2026 19:09
-- A drop-in replacement for pg_logical_slot_get_changes, taking and forwarding the same
-- plugin options, that holds back a change whose transaction has not settled yet. It works
-- the boundary out for itself, so it takes no upto_lsn.
CREATE OR REPLACE FUNCTION realtime.settled_changes(

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Note that I deliberately didn't add a guard to run this function only when sync standby is detected otherwise we'd create kind of a fork in one of the most used functions which would complicate debug, support, etc.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Update: added a feature flag to preserve the original functions while we test the new approach.

# Retried because a pooler can drop the connection mid-statement and OrioleDB can block on
# OTablesMetaTranche. The schema is dropped before it is recreated, so a lost attempt leaves the
# database with no realtime schema and every migration then fails with 3F000.
defp reset_realtime_schema!(settings) do

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Unfortunately adding Multigres is creating some flaky tests and this change was an attempt to fix no 'realtime' schema errors. An attempt because I still see such errors but less than before so it's a win still. But it will need a bit more work to eliminate all.

@blacksmith-sh

This comment has been minimized.

@leandrocp
leandrocp force-pushed the repro-sync-commit-event-loss branch from 8c2f0b2 to a5e6d94 Compare September 25, 2026 18:21
1. uses xid as a type
2. caches the running xids and age once instead of per-change lookup
3. removes the group by since the COMMIT messages are guaranteed to show
up
Comment thread lib/extensions/postgres_cdc_rls/replications.ex Outdated
@leandrocp
leandrocp merged commit 6884715 into main Sep 29, 2026
50 checks passed
@leandrocp
leandrocp deleted the repro-sync-commit-event-loss branch September 29, 2026 17:01
@realtime-release-bot

Copy link
Copy Markdown

🎉 This PR is included in version 2.140.1 🎉

The release is available on GitHub release

Your semantic-release bot 📦🚀

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

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants