Apply PR review: advertise org.matrix.msc4525 in /versions, tolerate malformed pos, protect cold-start drain, fix total_rooms

- /versions now advertises org.matrix.msc4525, honouring both the global
  flag and the per-user experimental feature (as Erik requested inline).
- An unparsable pos no longer 400s: it is treated like an unrecognised
  one and the connection starts afresh (MSC4525 has no client error
  path).
- The aging lane now also reserves page slots for never-sent rooms, so
  continuous traffic in already-delivered rooms can't stall the initial
  drain indefinitely.
- total_rooms is counted before partial-state filtering so cold-start
  progress reporting isn't understated.
This commit is contained in:
Matthew Hodgson
2026-08-07 13:38:08 +03:00
parent 318029506c
commit 20feb1ef15
7 changed files with 137 additions and 14 deletions
+1
View File
@@ -64,6 +64,7 @@ pub struct ExperimentalConfig {
pub msc4108_enabled: bool,
pub msc4108_delegation_endpoint: Option<String>,
pub msc3575_enabled: bool,
pub msc4525_enabled: bool,
pub msc4133_enabled: bool,
pub msc4155_enabled: bool,
pub msc4306_enabled: bool,
+22
View File
@@ -120,9 +120,27 @@ async fn build_versions_response(
None => global_unstable_feature_map.msc3575,
};
let msc4525_enabled = match user_id {
Some(user_id) => {
// Don't both looking anything up if it's enabled for everyone
if global_unstable_feature_map.msc4525 {
true
} else {
// Look up whether it's explicitly enabled/disabled for this user
store
.is_feature_enabled_for_user(user_id, PerUserExperimentalFeature::MSC4525)
.await?
// Default to false if there is no entry for this user
.unwrap_or(false)
}
}
None => global_unstable_feature_map.msc4525,
};
let unstable_feature_map = UnstableFeatureMap {
msc3575: msc3575_enabled,
msc3881: msc3881_enabled,
msc4525: msc4525_enabled,
..*global_unstable_feature_map
};
@@ -236,6 +254,9 @@ pub struct UnstableFeatureMap {
/// Simplified sliding sync
#[serde(rename = "org.matrix.simplified_msc3575")]
msc3575: bool,
/// MSC4525: Paginated sync
#[serde(rename = "org.matrix.msc4525")]
msc4525: bool,
/// Arbitrary key-value profile fields.
#[serde(rename = "uk.tcpip.msc4133")]
msc4133: bool,
@@ -309,6 +330,7 @@ pub fn synapse_config_to_global_unstable_feature_map(
|| (config.experimental.msc4108_delegation_endpoint.is_some()),
msc4140: config.server.msc4140_enabled,
msc3575: config.experimental.msc3575_enabled,
msc4525: config.experimental.msc4525_enabled,
msc4133: config.experimental.msc4133_enabled,
msc4133_stable: true,
msc4155: config.experimental.msc4155_enabled,
+2
View File
@@ -29,6 +29,8 @@ pub enum PerUserExperimentalFeature {
MSC3575,
#[serde(rename = "msc4222")]
MSC4222,
#[serde(rename = "msc4525")]
MSC4525,
}
impl std::fmt::Display for PerUserExperimentalFeature {
+37 -11
View File
@@ -230,6 +230,9 @@ class MSC4525PaginatedSyncHandler(SlidingSyncHandler):
)
sync_room_map = dict(interested_rooms.room_membership_for_user_map)
# For `total_rooms` (cold-start progress): the account's room count,
# before any partial-state filtering below.
total_rooms = len(sync_room_map)
newly_joined_rooms = interested_rooms.newly_joined_rooms
newly_left_rooms = interested_rooms.newly_left_rooms
dm_room_ids = interested_rooms.dm_room_ids
@@ -350,23 +353,46 @@ class MSC4525PaginatedSyncHandler(SlidingSyncHandler):
# Page: most recently active rooms first. When the page overflows, a
# slice of it is reserved for the longest-deferred rooms so that
# nothing is starved by busier rooms perpetually sorting first.
# nothing - neither a deferred update nor the never-sent cold-start
# backlog - is starved by busier rooms perpetually sorting first.
page_room_ids: list[str] = []
if candidates:
aged_room_ids: list[str] = []
if len(candidates) > page_size and previously_room_ids:
if (
from_token is not None
and len(candidates) > page_size
and (previously_room_ids or never_room_ids)
):
aging_lane_size = max(1, page_size // AGING_LANE_FRACTION)
def last_sent_stream_pos(room_id: str) -> int:
room_status = previous_connection_state.rooms.have_sent_room(
room_id
# Never-sent rooms first (deferred since the connection
# started, and invisible to the client until delivered), most
# recently active first to match the page's ordering; then the
# longest-deferred previously-sent rooms.
if never_room_ids:
never_room_infos = await self.room_lists.sort_rooms(
{room_id: sync_room_map[room_id] for room_id in never_room_ids},
to_token,
limit=aging_lane_size,
)
aged_room_ids.extend(
room_info.room_id for room_info in never_room_infos
)
assert room_status.last_token is not None
return room_status.last_token.stream
aged_room_ids = sorted(previously_room_ids, key=last_sent_stream_pos)[
:aging_lane_size
]
if len(aged_room_ids) < aging_lane_size and previously_room_ids:
def last_sent_stream_pos(room_id: str) -> int:
room_status = previous_connection_state.rooms.have_sent_room(
room_id
)
assert room_status.last_token is not None
return room_status.last_token.stream
aged_room_ids.extend(
sorted(previously_room_ids, key=last_sent_stream_pos)[
: aging_lane_size - len(aged_room_ids)
]
)
sorted_room_infos = await self.room_lists.sort_rooms(
{room_id: sync_room_map[room_id] for room_id in candidates},
@@ -466,7 +492,7 @@ class MSC4525PaginatedSyncHandler(SlidingSyncHandler):
rooms=rooms,
extensions=extensions,
pending=pending,
total_rooms=len(sync_room_map),
total_rooms=total_rooms,
)
set_tag("paginated_sync.result", bool(result))
+13 -3
View File
@@ -23,6 +23,7 @@ sync servlet; only the request parsing and the top-level response differ.
import logging
from typing import TYPE_CHECKING
from synapse.api.errors import SynapseError
from synapse.http.server import HttpServer
from synapse.http.servlet import (
parse_and_validate_json_object_from_request,
@@ -80,9 +81,18 @@ class MSC4525PaginatedSyncRestServlet(SlidingSyncRestServlet):
from_token = None
if from_token_string is not None:
from_token = await SlidingSyncStreamToken.from_string(
self.store, from_token_string
)
try:
from_token = await SlidingSyncStreamToken.from_string(
self.store, from_token_string
)
except SynapseError:
# There is no client error path in this API: an unparsable
# `pos` is treated the same as an unrecognised one - as
# absent, so the connection starts afresh.
logger.info(
"Unparsable paginated sync pos for %s; starting the connection afresh",
user,
)
body = parse_and_validate_json_object_from_request(
request, MSC4525PaginatedSyncBody
@@ -339,6 +339,35 @@ class MSC4525PaginatedSyncTestCase(unittest.HomeserverTestCase):
for room_id, room_response in response["rooms"].items():
self.assertTrue(room_response.get("initial"), room_id)
# An entirely unparsable pos is treated the same way, not a 400.
response = self._sync(body, pos="not a token at all")
self.assertEqual(set(response["rooms"].keys()), set(room_ids))
def test_cold_start_backlog_not_starved_by_live_traffic(self) -> None:
"""Rooms never sent on the connection get a reserved slice of every
page, so continuous traffic in already-sent rooms can't stall the
initial drain indefinitely."""
room_ids = self._create_rooms(6)
body = {"page_size": 2, "limit": 5, "history": 1}
response = self._sync(body)
seen_rooms = dict(response["rooms"])
pos = response["pos"]
self.assertEqual(len(seen_rooms), 2)
# Keep the already-delivered rooms busy while draining: without the
# aging lane every page would fill with the busy (most recently
# active) rooms and the never-sent backlog would starve.
for _ in range(20):
if not response.get("pending"):
break
for room_id in seen_rooms:
self.helper.send(room_id, body="chatter", tok=self.tok)
response = self._sync(body, pos=pos)
seen_rooms.update(response["rooms"])
pos = response["pos"]
self.assertEqual(set(seen_rooms.keys()), set(room_ids))
def test_extensions_apply_without_scoping(self) -> None:
"""Extensions have no lists/rooms scoping: enabling one is enough for
it to apply to the rooms in the response."""
+33
View File
@@ -142,6 +142,39 @@ class VersionsTestCase(unittest.HomeserverTestCase):
channel.json_body,
)
def test_msc4525_advertised_per_user(self) -> None:
user1_id = self.register_user("user1", "pass")
user1_tok = self.login(user1_id, "pass")
# Disabled by default.
channel = self.make_request(
"GET",
"/_matrix/client/versions",
access_token=user1_tok,
)
self.assertEqual(channel.code, 200, channel.result)
self.assertEqual(
channel.json_body["unstable_features"]["org.matrix.msc4525"],
False,
channel.json_body,
)
# Advertised once enabled for the user.
self._enable_experimental_feature_for_user(
target_user_id=user1_id, features={"msc4525": True}
)
channel = self.make_request(
"GET",
"/_matrix/client/versions",
access_token=user1_tok,
)
self.assertEqual(channel.code, 200, channel.result)
self.assertEqual(
channel.json_body["unstable_features"]["org.matrix.msc4525"],
True,
channel.json_body,
)
def test_msc4446_false_by_default(self) -> None:
channel = self.make_request("GET", "/_matrix/client/versions")
self.assertEqual(channel.code, 200, channel.result)