From 6d0a9dfc4b2e8a4d72fbc7880db3f2795d78a852 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Sierpi=C5=84ski?= <33436839+sierpinskid@users.noreply.github.com> Date: Wed, 19 Aug 2026 11:28:01 +0200 Subject: [PATCH] Advance _lastEventReceivedAt only by newer events. This fixes invalid /sync behavior where replaying older events would set _lastEventReceivedAt back in the past --- .../StreamChatLowLevelClient.cs | 27 ++++++-- .../StreamChatLowLevelClientTests.cs | 56 ++++++++++++++--- .../Tests/StateSync/StateSyncCatchUpTests.cs | 63 ++++++++++++++++++- 3 files changed, 132 insertions(+), 14 deletions(-) diff --git a/Assets/Plugins/StreamChat/Core/LowLevelClient/StreamChatLowLevelClient.cs b/Assets/Plugins/StreamChat/Core/LowLevelClient/StreamChatLowLevelClient.cs index 2a491b7a..45a29261 100644 --- a/Assets/Plugins/StreamChat/Core/LowLevelClient/StreamChatLowLevelClient.cs +++ b/Assets/Plugins/StreamChat/Core/LowLevelClient/StreamChatLowLevelClient.cs @@ -996,7 +996,7 @@ private void RegisterEventType(string key, #endif var eventObj = DeserializeEvent(serializedContent, out var dto); postprocess?.Invoke(dto); - _lastEventReceivedAt = eventObj.CreatedAt; + TryAdvanceLastEventReceivedAt(eventObj.CreatedAt, key); handler?.Invoke(eventObj, dto); internalHandler?.Invoke(dto); } @@ -1054,7 +1054,7 @@ private void HandleNewWebsocketMessage(string msg) if (!_eventKeyToHandler.TryGetValue(type, out var handler)) { - if (TryHandleCustomChannelEvent(msg)) + if (TryHandleCustomChannelEvent(msg, type)) { return; } @@ -1070,7 +1070,7 @@ private void HandleNewWebsocketMessage(string msg) handler(msg); } - private bool TryHandleCustomChannelEvent(string serializedContent) + private bool TryHandleCustomChannelEvent(string serializedContent, string eventType) { if (!_serializer.TryPeekValue(serializedContent, "cid", out var cid) || string.IsNullOrEmpty(cid)) @@ -1081,7 +1081,7 @@ private bool TryHandleCustomChannelEvent(string serializedContent) try { var dto = _serializer.Deserialize(serializedContent); - _lastEventReceivedAt = dto.CreatedAt; + TryAdvanceLastEventReceivedAt(dto.CreatedAt, eventType); var evt = new EventCustom(); ((ILoadableFrom)evt).LoadFromDto(dto); @@ -1141,6 +1141,25 @@ private void HandleHealthCheckEvent(EventHealthCheck healthCheckEvent, HealthChe } } + private void TryAdvanceLastEventReceivedAt(DateTimeOffset createdAt, string eventType) + { + if (createdAt == DateTimeOffset.MinValue) + { + if (_config.LogLevel.IsDebugEnabled()) + { + _logs.Warning( + $"WebSocket event `{eventType}` has no valid `created_at`; the /sync watermark was not advanced."); + } + + return; + } + + if (!_lastEventReceivedAt.HasValue || createdAt > _lastEventReceivedAt.Value) + { + _lastEventReceivedAt = createdAt; + } + } + private static bool IsUserIdValid(string userId) { var r = new Regex("^[a-zA-Z0-9@_-]+$"); diff --git a/Assets/Plugins/StreamChat/Tests/LowLevelClient/StreamChatLowLevelClientTests.cs b/Assets/Plugins/StreamChat/Tests/LowLevelClient/StreamChatLowLevelClientTests.cs index 75c56ec4..a8f35214 100644 --- a/Assets/Plugins/StreamChat/Tests/LowLevelClient/StreamChatLowLevelClientTests.cs +++ b/Assets/Plugins/StreamChat/Tests/LowLevelClient/StreamChatLowLevelClientTests.cs @@ -3,6 +3,7 @@ using System.Collections.Generic; using System.Linq; using System.Net.WebSockets; +using System.Reflection; using System.Threading; using System.Threading.Tasks; using NSubstitute; @@ -347,6 +348,27 @@ public void when_connection_state_changed_subscriber_throws_expect_remaining_sub Assert.AreNotEqual(ConnectionState.Connected, lastStateSeenByLateSubscriber); } + [Test] + public void when_event_without_created_at_expect_last_event_watermark_not_set() + { + var client = CreateConnectedClient(); + + Assert.IsNull(GetLastEventReceivedAt(client)); + } + + [Test] + public void when_event_with_created_at_expect_last_event_watermark_set() + { + var createdAt = new DateTimeOffset(2026, 8, 18, 13, 58, 59, TimeSpan.Zero); + var client = CreateClientWithMessages(logs: null, + $"{{\"connection_id\":\"fakeId\", \"type\":\"health.check\", \"created_at\":\"{createdAt:O}\"}}"); + + client.Connect(); + client.Update(deltaTime: 0.2f); + + Assert.AreEqual(createdAt, GetLastEventReceivedAt(client)); + } + private readonly List _resourcesToDispose = new List(); private IStreamChatLowLevelClient _lowLevelClient; @@ -362,6 +384,17 @@ public void when_connection_state_changed_subscriber_throws_expect_remaining_sub private IStreamClientConfig _mockStreamClientConfig; private StreamChatLowLevelClient CreateConnectedClient(ILogs logs = null) + { + var client = CreateClientWithMessages(logs, "{\"connection_id\":\"fakeId\", \"type\":\"health.check\"}"); + client.Connect(); + client.Update(deltaTime: 0.2f); + + Assert.IsTrue(client.ConnectionState == ConnectionState.Connected); + + return client; + } + + private StreamChatLowLevelClient CreateClientWithMessages(ILogs logs, params string[] websocketMessages) { var client = new StreamChatLowLevelClient(_authCredentials, _mockWebsocketClient, _mockHttpClient, new NewtonsoftJsonSerializer(), _mockTimeService, _mockNetworkMonitor, _mockApplicationInfo, @@ -370,19 +403,28 @@ private StreamChatLowLevelClient CreateConnectedClient(ILogs logs = null) _mockWebsocketClient.ConnectAsync(Arg.Any()).Returns(Task.CompletedTask); + var messages = new Queue(websocketMessages); _mockWebsocketClient.TryDequeueMessage(out Arg.Any()).Returns(arg => { - arg[0] = "{\"connection_id\":\"fakeId\", \"type\":\"health.check\"}"; - return true; - }, arg => false); - - client.Connect(); - client.Update(deltaTime: 0.2f); + if (messages.Count == 0) + { + return false; + } - Assert.IsTrue(client.ConnectionState == ConnectionState.Connected); + arg[0] = messages.Dequeue(); + return true; + }); return client; } + + private static DateTimeOffset? GetLastEventReceivedAt(StreamChatLowLevelClient client) + { + var field = typeof(StreamChatLowLevelClient).GetField("_lastEventReceivedAt", + BindingFlags.Instance | BindingFlags.NonPublic); + Assert.IsNotNull(field, "Expected _lastEventReceivedAt field to exist."); + return (DateTimeOffset?)field.GetValue(client); + } } } #endif \ No newline at end of file diff --git a/Assets/Plugins/StreamChat/Tests/StateSync/StateSyncCatchUpTests.cs b/Assets/Plugins/StreamChat/Tests/StateSync/StateSyncCatchUpTests.cs index 537ded75..0642ccfb 100644 --- a/Assets/Plugins/StreamChat/Tests/StateSync/StateSyncCatchUpTests.cs +++ b/Assets/Plugins/StreamChat/Tests/StateSync/StateSyncCatchUpTests.cs @@ -100,6 +100,44 @@ public void when_disconnection_timestamp_missing_expect_sync_not_called() AssertSyncNotCalled(); } + /// + /// /sync replays historical events through the same handlers as live events. Because the catch-up runs + /// asynchronously after reconnect, live events can already have advanced the watermark by the time the + /// replay is processed. Rewinding it here would make the next disconnect sync from a stale point. + /// + [Test] + public void when_sync_replays_events_older_than_watermark_expect_watermark_not_regressed() + { + var now = new DateTimeOffset(2026, 8, 10, 12, 0, 0, TimeSpan.Zero); + _mockTimeService.Now.Returns(now); + + SetLastEventReceivedAt(now); + SetDisconnectionLastEventReceivedAt(now.AddHours(-1)); + StubSyncResponseWithMessageEvent(now.AddHours(-1)); + + _lowLevelClient.FetchAndProcessEventsSinceLastReceivedEvent(new[] { TestChannelCid }).GetAwaiter() + .GetResult(); + + Assert.AreEqual(now, GetLastEventReceivedAt()); + } + + [Test] + public void when_sync_replays_events_newer_than_watermark_expect_watermark_advanced() + { + var now = new DateTimeOffset(2026, 8, 10, 12, 0, 0, TimeSpan.Zero); + _mockTimeService.Now.Returns(now); + + var replayedEventCreatedAt = now.AddHours(-1); + SetLastEventReceivedAt(now.AddHours(-2)); + SetDisconnectionLastEventReceivedAt(now.AddHours(-2)); + StubSyncResponseWithMessageEvent(replayedEventCreatedAt); + + _lowLevelClient.FetchAndProcessEventsSinceLastReceivedEvent(new[] { TestChannelCid }).GetAwaiter() + .GetResult(); + + Assert.AreEqual(replayedEventCreatedAt, GetLastEventReceivedAt()); + } + private const string TestChannelCid = "messaging:test-channel"; private StreamChatLowLevelClient _lowLevelClient; @@ -120,12 +158,31 @@ private void AssertSyncNotCalled() Arg.Is(uri => uri.AbsolutePath.EndsWith("/sync")), Arg.Any()); + private void StubSyncResponseWithMessageEvent(DateTimeOffset eventCreatedAt) + { + var body = + $"{{\"events\":[{{\"type\":\"message.new\",\"cid\":\"{TestChannelCid}\",\"created_at\":\"{eventCreatedAt:O}\"}}]}}"; + + _mockHttpClient + .SendHttpRequestAsync(Arg.Is(HttpMethodType.Post), Arg.Any(), Arg.Any()) + .Returns(new HttpResponse(true, 200, body, null, null)); + } + private void SetDisconnectionLastEventReceivedAt(DateTimeOffset value) + => GetPrivateField("_disconnectionLastEventReceivedAt").SetValue(_lowLevelClient, (DateTimeOffset?)value); + + private void SetLastEventReceivedAt(DateTimeOffset value) + => GetPrivateField("_lastEventReceivedAt").SetValue(_lowLevelClient, (DateTimeOffset?)value); + + private DateTimeOffset? GetLastEventReceivedAt() + => (DateTimeOffset?)GetPrivateField("_lastEventReceivedAt").GetValue(_lowLevelClient); + + private static FieldInfo GetPrivateField(string name) { - var field = typeof(StreamChatLowLevelClient).GetField("_disconnectionLastEventReceivedAt", + var field = typeof(StreamChatLowLevelClient).GetField(name, BindingFlags.Instance | BindingFlags.NonPublic); - Assert.IsNotNull(field, "Expected _disconnectionLastEventReceivedAt field to exist."); - field.SetValue(_lowLevelClient, (DateTimeOffset?)value); + Assert.IsNotNull(field, $"Expected {name} field to exist."); + return field; } } }