|
26 | 26 |
|
27 | 27 | from synapse.logging.opentracing import trace
|
28 | 28 | from synapse.storage._base import SQLBaseStore
|
29 |
| -from synapse.storage.database import LoggingTransaction |
| 29 | +from synapse.storage.database import LoggingTransaction, make_in_list_sql_clause |
30 | 30 | from synapse.storage.databases.main.stream import _filter_results_by_stream
|
31 |
| -from synapse.types import RoomStreamToken |
| 31 | +from synapse.types import RoomStreamToken, StrCollection |
32 | 32 | from synapse.util.caches.stream_change_cache import StreamChangeCache
|
| 33 | +from synapse.util.iterutils import batch_iter |
33 | 34 |
|
34 | 35 | logger = logging.getLogger(__name__)
|
35 | 36 |
|
@@ -200,3 +201,62 @@ def get_current_state_deltas_for_room_txn(
|
200 | 201 | return await self.db_pool.runInteraction(
|
201 | 202 | "get_current_state_deltas_for_room", get_current_state_deltas_for_room_txn
|
202 | 203 | )
|
| 204 | + |
| 205 | + @trace |
| 206 | + async def get_current_state_deltas_for_rooms( |
| 207 | + self, |
| 208 | + room_ids: StrCollection, |
| 209 | + from_token: RoomStreamToken, |
| 210 | + to_token: RoomStreamToken, |
| 211 | + ) -> List[StateDelta]: |
| 212 | + """Get the state deltas between two tokens for the set of rooms.""" |
| 213 | + |
| 214 | + room_ids = self._curr_state_delta_stream_cache.get_entities_changed( |
| 215 | + room_ids, from_token.stream |
| 216 | + ) |
| 217 | + if not room_ids: |
| 218 | + return [] |
| 219 | + |
| 220 | + def get_current_state_deltas_for_rooms_txn( |
| 221 | + txn: LoggingTransaction, |
| 222 | + room_ids: StrCollection, |
| 223 | + ) -> List[StateDelta]: |
| 224 | + clause, args = make_in_list_sql_clause( |
| 225 | + self.database_engine, "room_id", room_ids |
| 226 | + ) |
| 227 | + |
| 228 | + sql = f""" |
| 229 | + SELECT instance_name, stream_id, room_id, type, state_key, event_id, prev_event_id |
| 230 | + FROM current_state_delta_stream |
| 231 | + WHERE {clause} AND ? < stream_id AND stream_id <= ? |
| 232 | + ORDER BY stream_id ASC |
| 233 | + """ |
| 234 | + args.append(from_token.stream) |
| 235 | + args.append(to_token.get_max_stream_pos()) |
| 236 | + |
| 237 | + txn.execute(sql, args) |
| 238 | + |
| 239 | + return [ |
| 240 | + StateDelta( |
| 241 | + stream_id=row[1], |
| 242 | + room_id=row[2], |
| 243 | + event_type=row[3], |
| 244 | + state_key=row[4], |
| 245 | + event_id=row[5], |
| 246 | + prev_event_id=row[6], |
| 247 | + ) |
| 248 | + for row in txn |
| 249 | + if _filter_results_by_stream(from_token, to_token, row[0], row[1]) |
| 250 | + ] |
| 251 | + |
| 252 | + results = [] |
| 253 | + for batch in batch_iter(room_ids, 1000): |
| 254 | + deltas = await self.db_pool.runInteraction( |
| 255 | + "get_current_state_deltas_for_rooms", |
| 256 | + get_current_state_deltas_for_rooms_txn, |
| 257 | + batch, |
| 258 | + ) |
| 259 | + |
| 260 | + results.extend(deltas) |
| 261 | + |
| 262 | + return results |
0 commit comments