mirror of
https://github.com/element-hq/synapse.git
synced 2024-12-20 10:55:09 +03:00
Work on snapshots background update
This commit is contained in:
parent
ce7e1b4e67
commit
27caa47c50
2 changed files with 101 additions and 54 deletions
|
@ -137,7 +137,7 @@ class SlidingSyncStateInsertValues(TypedDict, total=False):
|
||||||
"""
|
"""
|
||||||
|
|
||||||
room_type: Optional[str]
|
room_type: Optional[str]
|
||||||
is_encrypted: Optional[bool]
|
is_encrypted: bool
|
||||||
room_name: Optional[str]
|
room_name: Optional[str]
|
||||||
tombstone_successor_room_id: Optional[str]
|
tombstone_successor_room_id: Optional[str]
|
||||||
|
|
||||||
|
@ -150,7 +150,7 @@ class SlidingSyncMembershipSnapshotSharedInsertValues(
|
||||||
multiple memberships
|
multiple memberships
|
||||||
"""
|
"""
|
||||||
|
|
||||||
has_known_state: Optional[bool]
|
has_known_state: bool
|
||||||
|
|
||||||
|
|
||||||
@attr.s(slots=True, auto_attribs=True)
|
@attr.s(slots=True, auto_attribs=True)
|
||||||
|
|
|
@ -1881,10 +1881,9 @@ class EventsBackgroundUpdatesStore(StreamWorkerStore, StateDeltasStore, SQLBaseS
|
||||||
keyvalues={"room_id": room_id},
|
keyvalues={"room_id": room_id},
|
||||||
values={},
|
values={},
|
||||||
insertion_values={
|
insertion_values={
|
||||||
|
# The reason we're only *inserting* (not *updating*) is because if
|
||||||
|
# they are present, that means they are already up-to-date.
|
||||||
**update_map,
|
**update_map,
|
||||||
# The reason we're only *inserting* (not *updating*) `event_stream_ordering`
|
|
||||||
# and `bump_stamp` is because if they are present, that means they are already
|
|
||||||
# up-to-date.
|
|
||||||
"event_stream_ordering": event_stream_ordering,
|
"event_stream_ordering": event_stream_ordering,
|
||||||
"bump_stamp": bump_stamp,
|
"bump_stamp": bump_stamp,
|
||||||
},
|
},
|
||||||
|
@ -2122,6 +2121,7 @@ class EventsBackgroundUpdatesStore(StreamWorkerStore, StateDeltasStore, SQLBaseS
|
||||||
last_to_insert_membership_infos_by_room_id: Dict[
|
last_to_insert_membership_infos_by_room_id: Dict[
|
||||||
str, SlidingSyncMembershipInfoWithEventPos
|
str, SlidingSyncMembershipInfoWithEventPos
|
||||||
] = {}
|
] = {}
|
||||||
|
# TODO: `concurrently_execute` based on buckets of room_ids
|
||||||
for (
|
for (
|
||||||
room_id,
|
room_id,
|
||||||
room_id_from_rooms_table,
|
room_id_from_rooms_table,
|
||||||
|
@ -2226,6 +2226,9 @@ class EventsBackgroundUpdatesStore(StreamWorkerStore, StateDeltasStore, SQLBaseS
|
||||||
# we only want some non-membership state
|
# we only want some non-membership state
|
||||||
await_full_state=False,
|
await_full_state=False,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# TODO: Read the state from the `room_stats_state` table if we can
|
||||||
|
|
||||||
# We're iterating over rooms that we are joined to so they should
|
# We're iterating over rooms that we are joined to so they should
|
||||||
# have `current_state_events` and we should have some current state
|
# have `current_state_events` and we should have some current state
|
||||||
# for each room
|
# for each room
|
||||||
|
@ -2404,59 +2407,103 @@ class EventsBackgroundUpdatesStore(StreamWorkerStore, StateDeltasStore, SQLBaseS
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# Assemble data so it's ready for the batch queries in the transaction
|
||||||
|
key_names = ("room_id", "user_id")
|
||||||
|
key_values: List[Tuple[str, str]] = []
|
||||||
|
insertion_value_names = (
|
||||||
|
"has_known_state",
|
||||||
|
"room_type",
|
||||||
|
"is_encrypted",
|
||||||
|
"room_name",
|
||||||
|
"tombstone_successor_room_id",
|
||||||
|
"sender",
|
||||||
|
"membership_event_id",
|
||||||
|
"membership",
|
||||||
|
"event_stream_ordering",
|
||||||
|
"event_instance_name",
|
||||||
|
)
|
||||||
|
insertion_value_values: List[
|
||||||
|
Tuple[
|
||||||
|
bool,
|
||||||
|
Optional[str],
|
||||||
|
bool,
|
||||||
|
Optional[str],
|
||||||
|
Optional[str],
|
||||||
|
str,
|
||||||
|
str,
|
||||||
|
str,
|
||||||
|
int,
|
||||||
|
str,
|
||||||
|
]
|
||||||
|
] = []
|
||||||
|
forgotten_update_query_args: List[Tuple[str, str, str]] = []
|
||||||
|
for key, insert_map in to_insert_membership_snapshots.items():
|
||||||
|
room_id, user_id = key
|
||||||
|
membership_info = to_insert_membership_infos[(room_id, user_id)]
|
||||||
|
sender = membership_info.sender
|
||||||
|
membership_event_id = membership_info.membership_event_id
|
||||||
|
membership = membership_info.membership
|
||||||
|
membership_event_stream_ordering = (
|
||||||
|
membership_info.membership_event_stream_ordering
|
||||||
|
)
|
||||||
|
membership_event_instance_name = (
|
||||||
|
membership_info.membership_event_instance_name
|
||||||
|
)
|
||||||
|
|
||||||
|
key_values.append((room_id, user_id))
|
||||||
|
insertion_value_values.append(
|
||||||
|
(
|
||||||
|
# `has_known_state` should be set to *some* True/False value
|
||||||
|
insert_map["has_known_state"],
|
||||||
|
insert_map.get("room_type"),
|
||||||
|
insert_map.get("is_encrypted", False),
|
||||||
|
insert_map.get("room_name"),
|
||||||
|
insert_map.get("tombstone_successor_room_id"),
|
||||||
|
sender,
|
||||||
|
membership_event_id,
|
||||||
|
membership,
|
||||||
|
membership_event_stream_ordering,
|
||||||
|
membership_event_instance_name,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
forgotten_update_query_args.append(
|
||||||
|
(
|
||||||
|
membership_event_id,
|
||||||
|
room_id,
|
||||||
|
user_id,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
def _fill_table_txn(txn: LoggingTransaction) -> None:
|
def _fill_table_txn(txn: LoggingTransaction) -> None:
|
||||||
# Handle updating the `sliding_sync_membership_snapshots` table
|
# Handle updating the `sliding_sync_membership_snapshots` table
|
||||||
#
|
#
|
||||||
for key, insert_map in to_insert_membership_snapshots.items():
|
# We don't need to update in the upsert the state because we never partially
|
||||||
room_id, user_id = key
|
# insert/update the snapshots and anything already there is up-to-date
|
||||||
membership_info = to_insert_membership_infos[(room_id, user_id)]
|
# EXCEPT for the `forgotten` field since that is updated out-of-band from
|
||||||
sender = membership_info.sender
|
# the membership changes. We're using an upsert to avoid unique
|
||||||
membership_event_id = membership_info.membership_event_id
|
# violation errors that would happen from directly inserting.
|
||||||
membership = membership_info.membership
|
self.db_pool.simple_upsert_many_txn(
|
||||||
membership_event_stream_ordering = (
|
txn,
|
||||||
membership_info.membership_event_stream_ordering
|
table="sliding_sync_membership_snapshots",
|
||||||
)
|
key_names=key_names,
|
||||||
membership_event_instance_name = (
|
key_values=key_values,
|
||||||
membership_info.membership_event_instance_name
|
# TODO: Implement these
|
||||||
)
|
insertion_value_names=insertion_value_names,
|
||||||
|
insertion_value_values=insertion_value_values,
|
||||||
|
)
|
||||||
|
|
||||||
# We don't need to upsert the state because we never partially
|
# We need to find the `forgotten` value during the transaction because
|
||||||
# insert/update the snapshots and anything already there is up-to-date
|
# we can't risk inserting stale data.
|
||||||
# EXCEPT for the `forgotten` field since that is updated out-of-band
|
txn.execute_batch(
|
||||||
# from the membership changes.
|
"""
|
||||||
#
|
UPDATE sliding_sync_membership_snapshots
|
||||||
# Even though we're only doing insertions, we're using
|
SET
|
||||||
# `simple_upsert_txn()` here to avoid unique violation errors that would
|
forgotten = (SELECT forgotten FROM room_memberships WHERE event_id = ?)
|
||||||
# happen from `simple_insert_txn()`
|
WHERE room_id = ? and user_id = ?
|
||||||
self.db_pool.simple_upsert_txn(
|
""",
|
||||||
txn,
|
forgotten_update_query_args,
|
||||||
table="sliding_sync_membership_snapshots",
|
)
|
||||||
keyvalues={"room_id": room_id, "user_id": user_id},
|
|
||||||
values={},
|
|
||||||
insertion_values={
|
|
||||||
**insert_map,
|
|
||||||
"sender": sender,
|
|
||||||
"membership_event_id": membership_event_id,
|
|
||||||
"membership": membership,
|
|
||||||
"event_stream_ordering": membership_event_stream_ordering,
|
|
||||||
"event_instance_name": membership_event_instance_name,
|
|
||||||
},
|
|
||||||
)
|
|
||||||
# We need to find the `forgotten` value during the transaction because
|
|
||||||
# we can't risk inserting stale data.
|
|
||||||
txn.execute(
|
|
||||||
"""
|
|
||||||
UPDATE sliding_sync_membership_snapshots
|
|
||||||
SET
|
|
||||||
forgotten = (SELECT forgotten FROM room_memberships WHERE event_id = ?)
|
|
||||||
WHERE room_id = ? and user_id = ?
|
|
||||||
""",
|
|
||||||
(
|
|
||||||
membership_event_id,
|
|
||||||
room_id,
|
|
||||||
user_id,
|
|
||||||
),
|
|
||||||
)
|
|
||||||
|
|
||||||
await self.db_pool.runInteraction(
|
await self.db_pool.runInteraction(
|
||||||
"sliding_sync_membership_snapshots_bg_update", _fill_table_txn
|
"sliding_sync_membership_snapshots_bg_update", _fill_table_txn
|
||||||
|
|
Loading…
Reference in a new issue