From 05f5ce231b81c645d88a52fabf30e90e40f47aab Mon Sep 17 00:00:00 2001 From: Kegan Dougal <7190048+kegsay@users.noreply.github.com> Date: Wed, 22 Apr 2026 16:58:37 +0100 Subject: [PATCH 1/3] fedclient: add support for MSC4242 fields Pending unit tests --- synapse/federation/federation_client.py | 177 +++++++++++++++--------- synapse/federation/transport/client.py | 32 +++-- synapse/handlers/federation.py | 7 + tests/handlers/test_device.py | 1 + tests/handlers/test_federation.py | 1 + tests/handlers/test_room_member.py | 1 + tests/storage/test_stream.py | 1 + 7 files changed, 150 insertions(+), 70 deletions(-) diff --git a/synapse/federation/federation_client.py b/synapse/federation/federation_client.py index 78a1900c731..1b444c31795 100644 --- a/synapse/federation/federation_client.py +++ b/synapse/federation/federation_client.py @@ -126,6 +126,9 @@ class SendJoinResult: # Always contains the server we joined off. servers_in_room: AbstractSet[str] + # Only valid for state DAG rooms (MSC4242) + state_dag: list[EventBase] | None + class FederationClient(FederationBase): def __init__(self, hs: "HomeServer"): @@ -1108,11 +1111,12 @@ async def send_join( SynapseError: if the chosen remote server returns a 300/400 code, or no servers successfully handle the request. """ - # See related restriction in /createRoom requests in handlers/room.py - if room_version.msc4242_state_dags: - raise UnsupportedRoomVersionError( - "Homeserver does not support this room version over federation" - ) + + def find_create_event(events: list[EventBase]) -> EventBase | None: + for e in events: + if (e.type, e.state_key) == (EventTypes.Create, ""): + return e + return None async def send_request(destination: str) -> SendJoinResult: response = await self._do_send_join( @@ -1142,13 +1146,16 @@ async def send_request(destination: str) -> SendJoinResult: state = response.state auth_chain = response.auth_events - - create_event = None - for e in state: - if (e.type, e.state_key) == (EventTypes.Create, ""): - create_event = e - break - + state_dag: list[EventBase] = [] + if room_version.msc4242_state_dags: + if not response.state_dag: + raise InvalidResponseError("No state_dag returned") + state_dag = response.state_dag + + # Validate the create event and room version are what we expect to see. + create_event = find_create_event( + state_dag if room_version.msc4242_state_dags else state + ) if create_event is None: # If the state doesn't have a create event then the room is # invalid, and it would fail auth checks anyway. @@ -1166,8 +1173,31 @@ async def send_request(destination: str) -> SendJoinResult: % (create_room_version,) ) + # Validate and set faster room joins fields + servers_in_room = None + if response.servers_in_room is not None: + servers_in_room = set(response.servers_in_room) + + if response.members_omitted: + if not servers_in_room: + raise InvalidResponseError( + "members_omitted was set, but no servers were listed in the room" + ) + + if not partial_state: + raise InvalidResponseError( + "members_omitted was set, but we asked for full state" + ) + + # `servers_in_room` is supposed to be a complete list. + # Fix things up in case the remote homeserver is badly behaved. + servers_in_room.add(destination) + logger.info( - "Processing from send_join %d events", len(state) + len(auth_chain) + "Processing from send_join %d events", + len(state_dag) + if room_version.msc4242_state_dags + else (len(state) + len(auth_chain)), ) # We now go and check the signatures and hashes for the event. Note @@ -1185,62 +1215,81 @@ async def _execute(pdu: EventBase) -> None: if valid_pdu: valid_pdus_map[valid_pdu.event_id] = valid_pdu - await concurrently_execute( - _execute, itertools.chain(state, auth_chain), 10000 - ) - - # NB: We *need* to copy to ensure that we don't have multiple - # references being passed on, as that causes... issues. - signed_state = [ - copy.copy(valid_pdus_map[p.event_id]) - for p in state - if p.event_id in valid_pdus_map - ] - - signed_auth = [ - valid_pdus_map[p.event_id] - for p in auth_chain - if p.event_id in valid_pdus_map - ] - - # NB: We *need* to copy to ensure that we don't have multiple - # references being passed on, as that causes... issues. - for s in signed_state: - s.internal_metadata = s.internal_metadata.copy() - - # double-check that the auth chain doesn't include a different create event - auth_chain_create_events = [ - e.event_id - for e in signed_auth - if (e.type, e.state_key) == (EventTypes.Create, "") - ] - if auth_chain_create_events and auth_chain_create_events != [ - create_event.event_id - ]: - raise InvalidResponseError( - "Unexpected create event(s) in auth chain: %s" - % (auth_chain_create_events,) + # Verify signatures/hashes on events, and make sure they all refer to the same room. + if room_version.msc4242_state_dags: + if state or auth_chain: + raise InvalidResponseError( + "State DAG rooms must not set state or auth_chain fields" + ) + await concurrently_execute(_execute, itertools.chain(state_dag), 10000) + # Copy valid PDUs along with internal metadata. + # It's unclear why this is needed but the code seems to expect it. + signed_state_dag = [ + copy.copy(valid_pdus_map[p.event_id]) + for p in state_dag + if p.event_id in valid_pdus_map + ] + for s in signed_state_dag: + s.internal_metadata = s.internal_metadata.copy() + + # Verify each event is for this room (and thus has the same create event as it is v12+) + for state_event in signed_state_dag: + if state_event.room_id != pdu.room_id: + raise InvalidResponseError( + "%s in state_dag belongs to room %s, not %s which we are joining" + % (state_event.event_id, state_event.room_id, pdu.room_id) + ) + return SendJoinResult( + event=event, + state=[], + auth_chain=[], + state_dag=signed_state_dag, + origin=destination, + partial_state=response.members_omitted, + servers_in_room=servers_in_room or frozenset(), ) - - servers_in_room = None - if response.servers_in_room is not None: - servers_in_room = set(response.servers_in_room) - - if response.members_omitted: - if not servers_in_room: + else: + if state_dag: raise InvalidResponseError( - "members_omitted was set, but no servers were listed in the room" + "Room does not support state DAGs but set state_dag field" ) + await concurrently_execute( + _execute, itertools.chain(state, auth_chain), 10000 + ) - if not partial_state: + # NB: We *need* to copy to ensure that we don't have multiple + # references being passed on, as that causes... issues. + signed_state = [ + copy.copy(valid_pdus_map[p.event_id]) + for p in state + if p.event_id in valid_pdus_map + ] + + signed_auth = [ + valid_pdus_map[p.event_id] + for p in auth_chain + if p.event_id in valid_pdus_map + ] + + # NB: We *need* to copy to ensure that we don't have multiple + # references being passed on, as that causes... issues. + for s in signed_state: + s.internal_metadata = s.internal_metadata.copy() + + # double-check that the auth chain doesn't include a different create event + auth_chain_create_events = [ + e.event_id + for e in signed_auth + if (e.type, e.state_key) == (EventTypes.Create, "") + ] + if auth_chain_create_events and auth_chain_create_events != [ + create_event.event_id + ]: raise InvalidResponseError( - "members_omitted was set, but we asked for full state" + "Unexpected create event(s) in auth chain: %s" + % (auth_chain_create_events,) ) - # `servers_in_room` is supposed to be a complete list. - # Fix things up in case the remote homeserver is badly behaved. - servers_in_room.add(destination) - return SendJoinResult( event=event, state=signed_state, @@ -1248,6 +1297,7 @@ async def _execute(pdu: EventBase) -> None: origin=destination, partial_state=response.members_omitted, servers_in_room=servers_in_room or frozenset(), + state_dag=None, ) # MSC3083 defines additional error codes for room joins. @@ -1548,6 +1598,7 @@ async def get_missing_events( limit: int, min_depth: int, timeout: int, + state_dag: bool = False, ) -> list[EventBase]: """Tries to fetch events we are missing. This is called when we receive an event without having received all of its ancestors. @@ -1563,6 +1614,7 @@ async def get_missing_events( limit: Maximum number of events to return. min_depth: Minimum depth of events to return. timeout: Max time to wait in ms + state_dag: True to walk the state DAG (MSC4242 rooms) """ try: content = await self.transport_layer.get_missing_events( @@ -1573,6 +1625,7 @@ async def get_missing_events( limit=limit, min_depth=min_depth, timeout=timeout, + state_dag=state_dag, ) room_version = await self.store.get_room_version(room_id) diff --git a/synapse/federation/transport/client.py b/synapse/federation/transport/client.py index 5d5212ef96c..e4747016c40 100644 --- a/synapse/federation/transport/client.py +++ b/synapse/federation/transport/client.py @@ -776,18 +776,21 @@ async def get_missing_events( limit: int, min_depth: int, timeout: int, + state_dag: bool, ) -> JsonDict: path = _create_v1_path("/get_missing_events/%s", room_id) - + request_body = { + "limit": int(limit), + "min_depth": int(min_depth), + "earliest_events": earliest_events, + "latest_events": latest_events, + } + if state_dag: + request_body["org.matrix.msc4242.state_dag"] = True return await self.client.post_json( destination=destination, path=path, - data={ - "limit": int(limit), - "min_depth": int(min_depth), - "earliest_events": earliest_events, - "latest_events": latest_events, - }, + data=request_body, timeout=timeout, ) @@ -986,6 +989,10 @@ class SendJoinResponse: # "event" is not included in the response. event: EventBase | None = None + # MSC4242: State DAGs. Always included for state dag rooms, else None. + # Replaces auth_events. + state_dag: list[EventBase] | None = None + # The room state is incomplete members_omitted: bool = False @@ -1068,7 +1075,7 @@ class SendJoinParser(ByteParser[SendJoinResponse]): MAX_RESPONSE_SIZE = 500 * 1024 * 1024 def __init__(self, room_version: RoomVersion, v1_api: bool): - self._response = SendJoinResponse([], [], event_dict={}) + self._response = SendJoinResponse([], [], event_dict={}, state_dag=[]) self._room_version = room_version self._coros: list[Generator[None, bytes, None]] = [] @@ -1112,6 +1119,15 @@ def __init__(self, room_version: RoomVersion, v1_api: bool): ) ) + if room_version.msc4242_state_dags: + self._coros.append( + ijson.items_coro( + _event_list_parser(room_version, self._response.state_dag), + prefix + "state_dag.item", + use_float=True, + ) + ) + def write(self, data: bytes) -> int: for c in self._coros: c.send(data) diff --git a/synapse/handlers/federation.py b/synapse/handlers/federation.py index b3444dd2ef0..c7f12310a7b 100644 --- a/synapse/handlers/federation.py +++ b/synapse/handlers/federation.py @@ -53,6 +53,7 @@ PartialStateConflictError, RequestSendFailed, SynapseError, + UnsupportedRoomVersionError, ) from synapse.api.room_versions import KNOWN_ROOM_VERSIONS, RoomVersion from synapse.crypto.event_signing import compute_event_signature @@ -646,6 +647,12 @@ async def do_invite_join( room_id ) + # See related restriction in /createRoom requests in handlers/room.py + if room_version_obj.msc4242_state_dags: + raise UnsupportedRoomVersionError( + "Homeserver does not support this room version over federation" + ) + ret = await self.federation_client.send_join( host_list, event, diff --git a/tests/handlers/test_device.py b/tests/handlers/test_device.py index 9e44b1dc1e0..9026e1697f9 100644 --- a/tests/handlers/test_device.py +++ b/tests/handlers/test_device.py @@ -822,6 +822,7 @@ def test_local_device_changes_sent_to_new_servers_on_un_partial_state( partial_state=True, # Only REMOTE1_SERVER_NAME is known at join time. servers_in_room={self.REMOTE1_SERVER_NAME}, + state_dag=None, ) ) diff --git a/tests/handlers/test_federation.py b/tests/handlers/test_federation.py index e4a41cf1ae4..efb800ba06b 100644 --- a/tests/handlers/test_federation.py +++ b/tests/handlers/test_federation.py @@ -652,6 +652,7 @@ def test_failed_partial_join_is_clean(self) -> None: ], partial_state=True, servers_in_room={"example.com"}, + state_dag=None, ) ) diff --git a/tests/handlers/test_room_member.py b/tests/handlers/test_room_member.py index 3890abdbc83..c84ad5e9e7a 100644 --- a/tests/handlers/test_room_member.py +++ b/tests/handlers/test_room_member.py @@ -172,6 +172,7 @@ def test_remote_joins_contribute_to_rate_limit(self) -> None: auth_chain=[create_event], partial_state=False, servers_in_room=frozenset(), + state_dag=None, ) ) diff --git a/tests/storage/test_stream.py b/tests/storage/test_stream.py index d51fa1f8bad..8fce7c03e9b 100644 --- a/tests/storage/test_stream.py +++ b/tests/storage/test_stream.py @@ -1455,6 +1455,7 @@ def test_remote_join(self) -> None: auth_chain=[create_event, creator_join_event], partial_state=False, servers_in_room=frozenset(), + state_dag=None, ) ) From 541e395920c653876f13eb4b925b4c3da0b5d0eb Mon Sep 17 00:00:00 2001 From: Kegan Dougal <7190048+kegsay@users.noreply.github.com> Date: Wed, 19 Aug 2026 13:59:12 +0100 Subject: [PATCH 2/3] Add changelog --- changelog.d/20127.feature | 1 + 1 file changed, 1 insertion(+) create mode 100644 changelog.d/20127.feature diff --git a/changelog.d/20127.feature b/changelog.d/20127.feature new file mode 100644 index 00000000000..5e6738e0359 --- /dev/null +++ b/changelog.d/20127.feature @@ -0,0 +1 @@ +Add experimental federation client support for [MSC4242](https://github.com/matrix-org/matrix-spec-proposals/pull/4242): State DAGs. From 4d9d1da80e43ed89c1354a01f237feba3967bd11 Mon Sep 17 00:00:00 2001 From: Kegan Dougal <7190048+kegsay@users.noreply.github.com> Date: Tue, 1 Sep 2026 11:09:43 +0100 Subject: [PATCH 3/3] Review comments --- synapse/federation/federation_client.py | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/synapse/federation/federation_client.py b/synapse/federation/federation_client.py index b860c097f34..fd3f9ab9122 100644 --- a/synapse/federation/federation_client.py +++ b/synapse/federation/federation_client.py @@ -1216,13 +1216,13 @@ async def _execute(pdu: EventBase) -> None: # Verify signatures/hashes on events, and make sure they all refer to the same room. if room_version.msc4242_state_dags: - if state or auth_chain: + if state or auth_chain or servers_in_room: raise InvalidResponseError( - "State DAG rooms must not set state or auth_chain fields" + "State DAG rooms must not set servers_in_room, state or auth_chain fields" ) await concurrently_execute(_execute, itertools.chain(state_dag), 10000) - # Copy valid PDUs along with internal metadata. - # It's unclear why this is needed but the code seems to expect it. + # NB: We *need* to copy to ensure that we don't have multiple + # references being passed on, as that causes... issues. signed_state_dag = [ valid_pdus_map[p.event_id].deep_copy() for p in state_dag @@ -1242,8 +1242,10 @@ async def _execute(pdu: EventBase) -> None: auth_chain=[], state_dag=signed_state_dag, origin=destination, - partial_state=response.members_omitted, - servers_in_room=servers_in_room or frozenset(), + # The current Synapse implementation of MSC4242 does not support + # faster remote room joins, so always set partial_state=False. + partial_state=False, + servers_in_room=frozenset(), ) else: if state_dag: