This repository has been archived by the owner on Apr 26, 2024. It is now read-only.
-
-
Notifications
You must be signed in to change notification settings - Fork 2.1k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Ensure stream position never goes backwards, and add test
- Loading branch information
1 parent
56e7fb3
commit fddfade
Showing
3 changed files
with
213 additions
and
14 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,133 @@ | ||
import json | ||
|
||
from parameterized import parameterized | ||
|
||
from synapse.rest import admin | ||
from synapse.rest.client import login, room, sync | ||
from synapse.types import JsonDict | ||
|
||
from tests import unittest | ||
from tests.test_utils.event_injection import create_event | ||
|
||
|
||
class TestLateArrivingMessages(unittest.HomeserverTestCase): | ||
servlets = [ | ||
admin.register_servlets_for_client_rest_resource, | ||
room.register_servlets, | ||
login.register_servlets, | ||
sync.register_servlets, | ||
] | ||
|
||
@parameterized.expand([(1,), (5,), (10,), (20,)]) | ||
def test_late_messages_have_received_ts(self, sync_timeline_limit: int) -> None: | ||
user1 = self.register_user("user1", "pass") | ||
tok1 = self.login("user1", "pass") | ||
|
||
room_id = self.helper.create_room_as(user1, is_public=True, tok=tok1) | ||
|
||
self.helper.send(room_id, "Hi!", tok=tok1) | ||
|
||
prev_event_ids = None | ||
event_and_contexts = [] | ||
depth = None | ||
for i in range(100): | ||
event, unpersisted_context = self.get_success( | ||
create_event( | ||
self.hs, | ||
prev_event_ids=prev_event_ids, | ||
room_id=room_id, | ||
sender=user1, | ||
type="m.room.message", | ||
content={"body": f"late message {i}", "msgtype": "m.text"}, | ||
depth=depth, | ||
) | ||
) | ||
if not prev_event_ids: | ||
depth = event.depth | ||
depth += 1 | ||
context = self.get_success(unpersisted_context.persist(event)) | ||
|
||
event_and_contexts.append((event, context)) | ||
prev_event_ids = [event.event_id] | ||
|
||
self.reactor.advance(60 * 60) | ||
|
||
for e, _ in event_and_contexts[-5:]: | ||
self.get_success(self.hs.get_datastores().main.received_event(e)) | ||
|
||
self.get_success( | ||
self.hs.get_storage_controllers().persistence.persist_events( | ||
event_and_contexts[-5:], backfilled=False | ||
) | ||
) | ||
|
||
self.reactor.advance(60 * 60) | ||
|
||
for e, _ in event_and_contexts[-50:-5]: | ||
self.get_success(self.hs.get_datastores().main.received_event(e)) | ||
|
||
self.get_success( | ||
self.hs.get_storage_controllers().persistence.persist_events( | ||
event_and_contexts[-50:-5], backfilled=True | ||
) | ||
) | ||
|
||
sync_filter = json.dumps({"room": {"timeline": {"limit": sync_timeline_limit}}}) | ||
|
||
channel = self.make_request( | ||
"GET", f"/sync?filter={sync_filter}", access_token=tok1 | ||
) | ||
self.assertEqual(channel.code, 200) | ||
|
||
room_result = channel.json_body["rooms"]["join"][room_id] | ||
|
||
group_id = None | ||
|
||
def check_is_late_event(event_json: JsonDict) -> None: | ||
nonlocal group_id | ||
|
||
self.assertSubstring("late message", event_json["content"]["body"]) | ||
self.assertIn("io.element.late_event", event_json["unsigned"]) | ||
late_metadata = event_json["unsigned"]["io.element.late_event"] | ||
self.assertIn("received_ts", late_metadata) | ||
|
||
self.assertGreater( | ||
late_metadata["received_ts"] - event_json["origin_server_ts"], 60 * 60 | ||
) | ||
|
||
self.assertIn("group_id", late_metadata) | ||
if group_id: | ||
self.assertEqual(late_metadata["group_id"], group_id) | ||
else: | ||
group_id = late_metadata["group_id"] | ||
|
||
for event_json in room_result["timeline"]["events"]: | ||
check_is_late_event(event_json) | ||
|
||
prev_batch = room_result["timeline"]["prev_batch"] | ||
|
||
channel = self.make_request( | ||
"GET", | ||
f"/rooms/{room_id}/messages?from={prev_batch}&dir=b", | ||
access_token=tok1, | ||
) | ||
self.assertEqual(channel.code, 200) | ||
|
||
events = channel.json_body["chunk"] | ||
next_batch = channel.json_body["end"] | ||
|
||
for event_json in events: | ||
check_is_late_event(event_json) | ||
|
||
channel = self.make_request( | ||
"GET", | ||
f"/rooms/{room_id}/messages?from={next_batch}&dir=b", | ||
access_token=tok1, | ||
) | ||
self.assertEqual(channel.code, 200) | ||
|
||
events = channel.json_body["chunk"] | ||
next_batch = channel.json_body["end"] | ||
|
||
for event_json in events: | ||
check_is_late_event(event_json) |