From 7cec9ca2fab78cd30adb578b5b967598e10f6dfd Mon Sep 17 00:00:00 2001 From: Matthew Hodgson Date: Sat, 8 Aug 2026 22:53:41 +0300 Subject: [PATCH] Address review of the hierarchy changes * The prefetch window iterated the queue in the opposite order to the pops. A room can appear in the queue more than once (linked from two spaces; the queue is only deduped at pop time), so the entry whose summary got computed could be the wrong one - summarising with another parent's via/depth. Iterate in pop order. * Bound the prefetch window by the number of rooms still needed for the page, so a small ?limit= no longer speculatively issues (and then discards) up to ten federation requests. * Cap how many children are summarised at once when answering a federation hierarchy request: that endpoint is reachable by any federating server, and each summary now fans out internally, so an uncapped gather over 50 children could monopolise the DB pool. * get_room_with_stats caches None for unknown rooms, and the federation hierarchy endpoint summarises whatever room id it is asked about, so a negative entry could outlive the room's creation (the stats writer only invalidates once it catches up) and trip _build_room_entry's assert. Invalidate on the room-creation paths, and omit the room from the summary rather than asserting if it has no room entry. * Lower get_room_hierarchy_state's max_entries: values hold every m.space.child event of the room, so entries are not uniformly small. --- synapse/handlers/room_summary.py | 55 +++++++++++++++++++------ synapse/storage/databases/main/room.py | 12 ++++++ synapse/storage/databases/main/state.py | 4 +- 3 files changed, 58 insertions(+), 13 deletions(-) diff --git a/synapse/handlers/room_summary.py b/synapse/handlers/room_summary.py index 47733b17a0..40c19ee699 100644 --- a/synapse/handlers/room_summary.py +++ b/synapse/handlers/room_summary.py @@ -51,7 +51,7 @@ from synapse.storage.databases.main.state import RoomHierarchyState from synapse.types import JsonDict, JsonMapping, Requester, StateMap, StrCollection from synapse.types.state import StateFilter from synapse.util import unwrapFirstError -from synapse.util.async_helpers import gather_results, yieldable_gather_results +from synapse.util.async_helpers import concurrently_execute, gather_results from synapse.util.caches.response_cache import ResponseCache if TYPE_CHECKING: @@ -72,6 +72,10 @@ MAX_SERVERS_PER_SPACE = 3 # requests) are fetched concurrently, ahead of the strictly-ordered traversal. PREFETCH_SUMMARIES = 10 +# how many of a space's children are summarised at once when responding to a +# federation hierarchy request. +CHILD_SUMMARY_CONCURRENCY = 10 + @attr.s(slots=True, frozen=True, auto_attribs=True) class _PaginationKey: @@ -331,8 +335,17 @@ class RoomSummaryHandler: # Iterate through the queue until we reach the limit or run out of # rooms to include. while room_queue and len(rooms_result) < limit: - # Kick off summary fetches for the next few rooms in traversal order. - for upcoming in room_queue[-PREFETCH_SUMMARIES:]: + # Kick off summary fetches for the next few rooms in traversal + # order, but never for more rooms than are still needed to fill the + # page: each one may be a federation request, and any we do not + # consume is wasted (and would be re-issued for the next page). + window = min(PREFETCH_SUMMARIES, limit - len(rooms_result)) + # NB reversed(): the queue is a stack, so the LAST entries are + # processed first. The same room can appear in the queue more than + # once (a room linked from two spaces; the queue is only deduped at + # pop time), and the entry carries the via/depth used to summarise + # it — so the entry registered here must be the one popped first. + for upcoming in reversed(room_queue[-window:]): if ( upcoming.room_id not in processed_rooms and upcoming.room_id not in prefetched @@ -518,19 +531,31 @@ class RoomSummaryHandler: assert isinstance(room_id, str) child_ids.append(room_id) - async def summarize_child( - room_id: str, - ) -> tuple[str, Optional["_RoomEntry"]] | None: + summaries: list[tuple[str, Optional["_RoomEntry"]] | None] = [None] * len( + child_ids + ) + + async def summarize_child(indexed_room_id: tuple[int, str]) -> None: + index, room_id = indexed_room_id # If the room is unknown, skip it. if not await self._store.is_host_joined(room_id, self._server_name): - return None - return room_id, await self._summarize_local_room( - None, origin, room_id, suggested_only, include_children=False + return + summaries[index] = ( + room_id, + await self._summarize_local_room( + None, origin, room_id, suggested_only, include_children=False + ), ) # Summarise each child (but not its children) concurrently, adding them - # to the response in their original order. - for result in await yieldable_gather_results(summarize_child, child_ids): + # to the response in their original order. This endpoint is reachable by + # any federating server and each summary itself fans out, so cap how + # many run at once rather than issuing all MAX_ROOMS_PER_SPACE at once. + await concurrently_execute( + summarize_child, list(enumerate(child_ids)), CHILD_SUMMARY_CONCURRENCY + ) + + for result in summaries: if result is None: continue room_id, room_entry = result @@ -600,7 +625,7 @@ class RoomSummaryHandler: except (NotFoundError, UnsupportedRoomVersionError): pass - is_partial_state, hierarchy_state, _, _ = await make_deferred_yieldable( + is_partial_state, hierarchy_state, stats, _ = await make_deferred_yieldable( gather_results( ( # NB lambdas as run_in_background's typing rejects @cached methods. @@ -632,6 +657,12 @@ class RoomSummaryHandler: admin_skip_room_visibility_check, ) + if stats is None: + # We have state for the room but no row in the rooms table, so there + # is nothing to summarise. _build_room_entry would assert. + logger.info("room %s has no room entry, omitting from summary", room_id) + return None + if ( not admin_skip_room_visibility_check and not await self._is_local_room_accessible_from_summary( diff --git a/synapse/storage/databases/main/room.py b/synapse/storage/databases/main/room.py index 8b7893a2a6..0f7cf76f72 100644 --- a/synapse/storage/databases/main/room.py +++ b/synapse/storage/databases/main/room.py @@ -404,6 +404,11 @@ class RoomWorkerStore(CacheInvalidationWorkerStore): logger.error("store_room with room_id=%s failed: %s", room_id, e) raise StoreError(500, "Problem creating room.") + # The room may have been looked up before we had a row for it (e.g. a + # remote server asking us to summarise it), caching a None; the stats + # writer only invalidates once it catches up. + await self.invalidate_cache_and_stream("get_room_with_stats", (room_id,)) + async def get_room(self, room_id: str) -> tuple[bool, bool] | None: """Retrieve a room. @@ -2489,6 +2494,9 @@ class RoomWorkerStore(CacheInvalidationWorkerStore): "has_auth_chain_index": has_auth_chain_index, }, ) + # As in store_room: drop any negative entry cached before the room row + # existed. + await self.invalidate_cache_and_stream("get_room_with_stats", (room_id,)) class _BackgroundUpdates: @@ -2996,6 +3004,10 @@ class RoomStore(RoomBackgroundUpdateStore, RoomWorkerStore): "has_auth_chain_index": has_auth_chain_index, }, ) + # The room may have been looked up before we had a row for it (e.g. a + # remote server asking us to summarise it), caching a None; the stats + # writer only invalidates once it catches up. + await self.invalidate_cache_and_stream("get_room_with_stats", (room_id,)) async def store_partial_state_room( self, diff --git a/synapse/storage/databases/main/state.py b/synapse/storage/databases/main/state.py index c23f0a497a..3d31d080fa 100644 --- a/synapse/storage/databases/main/state.py +++ b/synapse/storage/databases/main/state.py @@ -611,7 +611,9 @@ class StateGroupWorkerStore(EventsWorkerStore, SQLBaseStore): "get_filtered_current_state_ids", _get_filtered_current_state_ids_txn ) - @cached(max_entries=50000) + # NB max_entries is lower than the sibling state caches: each value holds + # every m.space.child event of the room, so entries are not uniformly small. + @cached(max_entries=10000) async def get_room_hierarchy_state( self, room_id: str ) -> Optional["RoomHierarchyState"]: