diff --git a/synapse/storage/databases/main/events.py b/synapse/storage/databases/main/events.py index 62ae371ca3..909792095c 100644 --- a/synapse/storage/databases/main/events.py +++ b/synapse/storage/databases/main/events.py @@ -1720,10 +1720,10 @@ class PersistEventsStore: txn.execute_batch( f""" INSERT INTO sliding_sync_membership_snapshots - (room_id, user_id, membership_event_id, membership, event_stream_ordering + (room_id, user_id, sender, membership_event_id, membership, event_stream_ordering {("," + ", ".join(insert_keys)) if insert_keys else ""}) VALUES ( - ?, ?, ?, ?, ? + ?, ?, ?, ?, ?, ? {("," + ", ".join("?" for _ in insert_values)) if insert_values else ""} ) ON CONFLICT (room_id, user_id) @@ -1737,6 +1737,7 @@ class PersistEventsStore: [ room_id, membership_info.user_id, + membership_info.sender, membership_info.membership_event_id, membership_info.membership, membership_info.membership_event_stream_ordering, @@ -2693,6 +2694,7 @@ class PersistEventsStore: raw_stripped_state_events = knock_room_state insert_values = { + "sender": event.sender, "membership_event_id": event.event_id, "membership": event.membership, "event_stream_ordering": event.internal_metadata.stream_ordering, diff --git a/synapse/storage/databases/main/events_bg_updates.py b/synapse/storage/databases/main/events_bg_updates.py index 5d04955bf1..a2bfaabafd 100644 --- a/synapse/storage/databases/main/events_bg_updates.py +++ b/synapse/storage/databases/main/events_bg_updates.py @@ -1631,6 +1631,8 @@ class EventsBackgroundUpdatesStore(StreamWorkerStore, StateDeltasStore, SQLBaseS ) def _backfill_table_txn(txn: LoggingTransaction) -> None: + # Handle updating the `sliding_sync_joined_rooms` table + # last_successful_room_id: Optional[str] = None for room_id, insert_map in joined_room_updates.items(): ( @@ -1965,9 +1967,12 @@ class EventsBackgroundUpdatesStore(StreamWorkerStore, StateDeltasStore, SQLBaseS ) def _backfill_table_txn(txn: LoggingTransaction) -> None: + # Handle updating the `sliding_sync_membership_snapshots` table + # for key, insert_map in to_insert_membership_snapshots.items(): room_id, user_id = key membership_info = to_insert_membership_infos[key] + sender = membership_info.sender membership_event_id = membership_info.membership_event_id membership = membership_info.membership membership_event_stream_ordering = ( @@ -1983,10 +1988,10 @@ class EventsBackgroundUpdatesStore(StreamWorkerStore, StateDeltasStore, SQLBaseS txn.execute( f""" INSERT INTO sliding_sync_membership_snapshots - (room_id, user_id, membership_event_id, membership, event_stream_ordering + (room_id, user_id, sender, membership_event_id, membership, event_stream_ordering {("," + ", ".join(insert_keys)) if insert_keys else ""}) VALUES ( - ?, ?, ?, ?, ? + ?, ?, ?, ?, ?, ? {("," + ", ".join("?" for _ in insert_values)) if insert_values else ""} ) ON CONFLICT (room_id, user_id) @@ -1995,6 +2000,7 @@ class EventsBackgroundUpdatesStore(StreamWorkerStore, StateDeltasStore, SQLBaseS [ room_id, user_id, + sender, membership_event_id, membership, membership_event_stream_ordering, diff --git a/synapse/storage/schema/main/delta/87/01_sliding_sync_memberships.sql b/synapse/storage/schema/main/delta/87/01_sliding_sync_memberships.sql index 16b3f84c3d..5fac6af619 100644 --- a/synapse/storage/schema/main/delta/87/01_sliding_sync_memberships.sql +++ b/synapse/storage/schema/main/delta/87/01_sliding_sync_memberships.sql @@ -60,6 +60,8 @@ CREATE TABLE IF NOT EXISTS sliding_sync_joined_rooms( CREATE TABLE IF NOT EXISTS sliding_sync_membership_snapshots( room_id TEXT NOT NULL REFERENCES rooms(room_id), user_id TEXT NOT NULL, + -- Useful to be able to tell leaves from kicks (where the `user_id` is different from the `sender`) + sender TEXT NOT NULL, membership_event_id TEXT NOT NULL REFERENCES events(event_id), membership TEXT NOT NULL, -- `stream_ordering` of the `membership_event_id`