From 0c3b999f8818d30b42954ff8d4fa560e72eb27c2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Sierpi=C5=84ski?= <33436839+sierpinskid@users.noreply.github.com> Date: Thu, 13 Aug 2026 12:29:59 +0200 Subject: [PATCH 1/4] Implement cache limit on messages. This helps manage memory on channels with a very high volume of received messages. The cache trimming will remove messages from the older one. It's possible to pause cache trimming if UI would be currently showing older messages so we don't clear message during user scrolling the history --- .../Core/Configs/IStreamClientConfig.cs | 8 + .../Core/Configs/MessageCacheWindow.cs | 51 ++ .../Core/Configs/MessageCacheWindow.cs.meta | 3 + .../Core/Configs/StreamClientConfig.cs | 2 + .../StreamChat/Core/Helpers/FilteredList.cs | 2 + .../StreamChat/Core/Helpers/ListPool.cs | 36 ++ .../StreamChat/Core/Helpers/ListPool.cs.meta | 3 + .../Core/StatefulModels/IStreamChannel.cs | 54 +++ .../Core/StatefulModels/StreamChannel.cs | 130 ++++++ .../StreamChat/Core/StreamChatClient.cs | 2 + .../StatefulClient/MessageCacheWindowTests.cs | 438 ++++++++++++++++++ .../MessageCacheWindowTests.cs.meta | 3 + 12 files changed, 732 insertions(+) create mode 100644 Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs create mode 100644 Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs.meta create mode 100644 Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs create mode 100644 Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs.meta create mode 100644 Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs create mode 100644 Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs.meta diff --git a/Assets/Plugins/StreamChat/Core/Configs/IStreamClientConfig.cs b/Assets/Plugins/StreamChat/Core/Configs/IStreamClientConfig.cs index a8d1628d..69887678 100644 --- a/Assets/Plugins/StreamChat/Core/Configs/IStreamClientConfig.cs +++ b/Assets/Plugins/StreamChat/Core/Configs/IStreamClientConfig.cs @@ -28,5 +28,13 @@ public interface IStreamClientConfig /// ordering matters more than instant local feedback (e.g. a shared, broadcast-ordered feed). /// bool OptimisticMessageInsert { get; set; } + + /// + /// Default local message cache limit for all channels. null (default) = unlimited. + /// Use for livestream-style channels. + /// Per-channel overrides: . + /// Does not change server history. See . + /// + MessageCacheWindow DefaultMessageCacheWindow { get; set; } } } \ No newline at end of file diff --git a/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs b/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs new file mode 100644 index 00000000..6834bbb7 --- /dev/null +++ b/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs @@ -0,0 +1,51 @@ +using System; + +namespace StreamChat.Core.Configs +{ + /// + /// Limits how many messages a channel keeps in the local cache. + /// Assign to or + /// . + /// Trimming runs in batches of once the count exceeds + /// . + /// + public sealed class MessageCacheWindow + { + /// Keep up to 500 messages; remove 100 at a time when over the limit. + public static readonly MessageCacheWindow Recommended = new MessageCacheWindow(500, 100); + + /// Trimming starts when exceeds this count. + public int MaxMessages { get; } + + /// How many messages to remove per trim. Must be less than . + public int DiscardBatchSize { get; } + + public MessageCacheWindow(int maxMessages, int discardBatchSize) + { + if (maxMessages <= 0) + { + throw new ArgumentOutOfRangeException(nameof(maxMessages), maxMessages, + $"{nameof(maxMessages)} must be greater than zero."); + } + + if (discardBatchSize <= 0) + { + throw new ArgumentOutOfRangeException(nameof(discardBatchSize), discardBatchSize, + $"{nameof(discardBatchSize)} must be greater than zero."); + } + + if (discardBatchSize >= maxMessages) + { + throw new ArgumentOutOfRangeException(nameof(discardBatchSize), discardBatchSize, + $"{nameof(discardBatchSize)} must be smaller than {nameof(maxMessages)} ({maxMessages}), " + + "otherwise a single trim would remove every message."); + } + + MaxMessages = maxMessages; + DiscardBatchSize = discardBatchSize; + } + + public override string ToString() + => $"MessageCacheWindow - MaxMessages: {MaxMessages}, DiscardBatchSize: {DiscardBatchSize}"; + } +} diff --git a/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs.meta b/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs.meta new file mode 100644 index 00000000..edd769d1 --- /dev/null +++ b/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs.meta @@ -0,0 +1,3 @@ +fileFormatVersion: 2 +guid: d0a6f5b24e7c5c9d1b8a3f6e2d4c5b7a +timeCreated: 1755079200 diff --git a/Assets/Plugins/StreamChat/Core/Configs/StreamClientConfig.cs b/Assets/Plugins/StreamChat/Core/Configs/StreamClientConfig.cs index 5e1da8bf..54b698f9 100644 --- a/Assets/Plugins/StreamChat/Core/Configs/StreamClientConfig.cs +++ b/Assets/Plugins/StreamChat/Core/Configs/StreamClientConfig.cs @@ -10,5 +10,7 @@ public class StreamClientConfig : IStreamClientConfig public StreamLogLevel LogLevel { get; set; } = StreamLogLevel.FailureOnly; public bool OptimisticMessageInsert { get; set; } = true; + + public MessageCacheWindow DefaultMessageCacheWindow { get; set; } = null; } } \ No newline at end of file diff --git a/Assets/Plugins/StreamChat/Core/Helpers/FilteredList.cs b/Assets/Plugins/StreamChat/Core/Helpers/FilteredList.cs index 43eeb84a..e4a1a31f 100644 --- a/Assets/Plugins/StreamChat/Core/Helpers/FilteredList.cs +++ b/Assets/Plugins/StreamChat/Core/Helpers/FilteredList.cs @@ -56,6 +56,8 @@ public void Insert(int index, T item) public void RemoveAt(int index) => _internalList.RemoveAt(index); + public void RemoveRange(int index, int count) => _internalList.RemoveRange(index, count); + public T this[int index] { get => _internalList[index]; diff --git a/Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs b/Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs new file mode 100644 index 00000000..3b075530 --- /dev/null +++ b/Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs @@ -0,0 +1,36 @@ +using System; +using System.Collections.Generic; + +namespace StreamChat.Core.Helpers +{ + internal static class ListPool + { + public static List Rent() + { + if (Pool.Count > 0) + { + return Pool.Pop(); + } + + return new List(); + } + + public static void Release(List list) + { + if (list == null) + { + throw new ArgumentNullException(nameof(list)); + } + + list.Clear(); + + if (Pool.Count < MaxPoolSize) + { + Pool.Push(list); + } + } + + private const int MaxPoolSize = 128; + private static readonly Stack> Pool = new Stack>(); + } +} diff --git a/Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs.meta b/Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs.meta new file mode 100644 index 00000000..3a1cc3d6 --- /dev/null +++ b/Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs.meta @@ -0,0 +1,3 @@ +fileFormatVersion: 2 +guid: c9f5e4a13d6b4b8c0a7f2e5d1c3b4a6f +timeCreated: 1755079200 diff --git a/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs b/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs index d9b82cf3..62b44422 100644 --- a/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs +++ b/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs @@ -4,6 +4,7 @@ using StreamChat.Core.Models; using StreamChat.Core.Requests; using StreamChat.Core.Responses; +using StreamChat.Core.Configs; namespace StreamChat.Core.StatefulModels { @@ -29,6 +30,12 @@ public interface IStreamChannel : IStreamStatefulModel /// event StreamMessageDeleteHandler MessageDeleted; + /// + /// Fired when old messages are removed from the local cache to save memory. This is not a server + /// delete — use for that. Remove your UI rows here. Oldest first. + /// + event StreamChannelMessagesHandler MessagesRemovedFromCache; + /// /// Event fired when a new was added to /// @@ -316,6 +323,53 @@ public interface IStreamChannel : IStreamStatefulModel /// Task LoadOlderMessagesAsync(); + /// + /// Active cache limit for this channel. null = unlimited (default). + /// When over , oldest messages are removed from + /// and the cache. Server history is unchanged — + /// can reload them. Pinned messages and open threads may + /// stay in the cache. always fires before + /// for the same message. + /// + Configs.MessageCacheWindow MessageCacheWindow { get; } + + /// + /// true if this channel has its own limit via . + /// When true, returns that override — even when it is + /// null (unlimited for this channel only). When false, the channel uses + /// . + /// + bool HasMessageCacheWindowOverride { get; } + + /// + /// Set a cache limit for this channel only. Pass null for unlimited on this channel. + /// Trims immediately unless is true. + /// + void OverrideMessageCacheWindow(Configs.MessageCacheWindow window); + + /// + /// Remove the per-channel limit and use again. + /// Trims immediately unless trimming is paused. + /// + void ClearMessageCacheWindowOverride(); + + /// + /// Whether cache trimming is paused. pauses automatically. + /// Memory is unbounded while paused. + /// + bool IsMessageCacheTrimmingPaused { get; } + + /// + /// Pause cache trimming (e.g. while the user scrolls through loaded history). + /// does this for you. + /// + void PauseMessageCacheTrimming(); + + /// + /// Resume trimming and remove excess messages now. Call when the user returns to the newest messages. + /// + void ResumeMessageCacheTrimming(); + /// /// Update channel in a complete overwrite mode. /// Important! Any data that is present on the channel and not included in a full update will be deleted. diff --git a/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs b/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs index 1cf94535..f2047742 100644 --- a/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs +++ b/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs @@ -3,6 +3,7 @@ using System.Linq; using System.Threading.Tasks; using StreamChat.Core.Helpers; +using StreamChat.Core.Configs; using StreamChat.Core.InternalDTO.Events; using StreamChat.Core.InternalDTO.Models; using StreamChat.Core.InternalDTO.Requests; @@ -21,6 +22,8 @@ namespace StreamChat.Core.StatefulModels public delegate void StreamChannelMessageHandler(IStreamChannel channel, IStreamMessage message); + public delegate void StreamChannelMessagesHandler(IStreamChannel channel, IReadOnlyList messages); + public delegate void StreamMessageDeleteHandler(IStreamChannel channel, IStreamMessage message, bool isHardDelete); public delegate void StreamChannelChangeHandler(IStreamChannel channel); @@ -50,6 +53,8 @@ internal sealed class StreamChannel : StreamStatefulModelBase, public event StreamMessageDeleteHandler MessageDeleted; + public event StreamChannelMessagesHandler MessagesRemovedFromCache; + public event StreamMessageReactionHandler ReactionAdded; public event StreamMessageReactionHandler ReactionRemoved; @@ -244,6 +249,9 @@ public async Task SendNewMessageAsync(StreamSendMessageRequest s public async Task LoadOlderMessagesAsync() { + // Pause before the request so live messages cannot trim the page being loaded. + PauseMessageCacheTrimming(); + var oldestMessage = _messages.OrderBy(_ => _.CreatedAt).FirstOrDefault(); var request = new ChannelGetOrCreateRequestInternalDTO @@ -266,6 +274,35 @@ public async Task LoadOlderMessagesAsync() Cache.TryCreateOrUpdate(response); } + public MessageCacheWindow MessageCacheWindow + => _hasMessageCacheWindowOverride ? _messageCacheWindowOverride : LowLevelClient.Config.DefaultMessageCacheWindow; + + public bool HasMessageCacheWindowOverride => _hasMessageCacheWindowOverride; + + public void OverrideMessageCacheWindow(MessageCacheWindow window) + { + _messageCacheWindowOverride = window; + _hasMessageCacheWindowOverride = true; + TrimMessageCacheIfNeeded(); + } + + public void ClearMessageCacheWindowOverride() + { + _messageCacheWindowOverride = null; + _hasMessageCacheWindowOverride = false; + TrimMessageCacheIfNeeded(); + } + + public bool IsMessageCacheTrimmingPaused => _isMessageCacheTrimmingPaused; + + public void PauseMessageCacheTrimming() => _isMessageCacheTrimmingPaused = true; + + public void ResumeMessageCacheTrimming() + { + _isMessageCacheTrimmingPaused = false; + TrimMessageCacheIfNeeded(); + } + public async Task UpdateOverwriteAsync(StreamUpdateOverwriteChannelRequest updateOverwriteRequest) { StreamAsserts.AssertNotNull(updateOverwriteRequest, nameof(updateOverwriteRequest)); @@ -906,6 +943,10 @@ protected override string InternalUniqueId private readonly List _ownCapabilities = new List(); private readonly List _pendingMessages = new List(); + private MessageCacheWindow _messageCacheWindowOverride; + private bool _hasMessageCacheWindowOverride; + private bool _isMessageCacheTrimmingPaused; + private bool _muted; private bool _hidden; @@ -963,6 +1004,9 @@ private bool InternalAppendOrUpdateMessage(MessageInternalDTO dto, out StreamMes } MessageReceived?.Invoke(this, streamMessage); + + // Trim after MessageReceived so a message is never removed from cache before it is received. + TrimMessageCacheIfNeeded(); return true; } @@ -1000,6 +1044,92 @@ private void InternalTruncateMessages(DateTimeOffset? deleteBeforeCreatedAt = nu Truncated?.Invoke(this); } + private void TrimMessageCacheIfNeeded() + { + var window = MessageCacheWindow; + if (window == null || _isMessageCacheTrimmingPaused || _messages.Count <= window.MaxMessages) + { + return; + } + + if (!AreMessagesSortedByCreatedAt()) + { + SortMessagesByCreatedAt(); + } + + var targetCount = window.MaxMessages - window.DiscardBatchSize; + var removeCount = _messages.Count - targetCount; + + var tempDeleteCandidates = ListPool.Rent(); + try + { + for (var i = 0; i < removeCount; i++) + { + tempDeleteCandidates.Add(_messages[i]); + } + + _messages.RemoveRange(0, removeCount); + + for (var i = 0; i < tempDeleteCandidates.Count; i++) + { + var message = tempDeleteCandidates[i]; + if (!IsRetainedByOtherState(message)) + { + Cache.Messages.Remove(message); + } + } + + if (MessagesRemovedFromCache != null) + { + var removed = new List(tempDeleteCandidates.Count); + for (var i = 0; i < tempDeleteCandidates.Count; i++) + { + removed.Add(tempDeleteCandidates[i]); + } + + MessagesRemovedFromCache.Invoke(this, removed); + } + } + finally + { + ListPool.Release(tempDeleteCandidates); + } + } + + // Keep in cache when pinned or referenced by a tracked thread. + private bool IsRetainedByOtherState(StreamMessage message) + { + if (_pinnedMessages.ContainsNoAlloc(message)) + { + return true; + } + + if (Cache.Threads.TryGet(message.Id, out _)) + { + return true; + } + + if (!string.IsNullOrEmpty(message.ParentId) && Cache.Threads.TryGet(message.ParentId, out _)) + { + return true; + } + + return false; + } + + private bool AreMessagesSortedByCreatedAt() + { + for (var i = 1; i < _messages.Count; i++) + { + if (_messages[i].CreatedAt < _messages[i - 1].CreatedAt) + { + return false; + } + } + + return true; + } + private void UpdateChannelFieldsFromDto(ChannelResponseInternalDTO dto, ICache cache) { #region Channel diff --git a/Assets/Plugins/StreamChat/Core/StreamChatClient.cs b/Assets/Plugins/StreamChat/Core/StreamChatClient.cs index d05928fd..cec74b01 100644 --- a/Assets/Plugins/StreamChat/Core/StreamChatClient.cs +++ b/Assets/Plugins/StreamChat/Core/StreamChatClient.cs @@ -884,6 +884,8 @@ void IStreamChatClientEventsListener.Destroy() internal StreamChatLowLevelClient InternalLowLevelClient { get; } + internal ICache InternalCache => _cache; + // We probably don't want to expose the presence, state, watch params to the public API internal async Task InternalGetOrCreateChannelWithIdAsync(ChannelType channelType, string channelId, diff --git a/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs b/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs new file mode 100644 index 00000000..8b70022a --- /dev/null +++ b/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs @@ -0,0 +1,438 @@ +#if STREAM_TESTS_ENABLED +using System; +using System.Collections; +using System.Collections.Generic; +using System.Linq; +using System.Threading.Tasks; +using NUnit.Framework; +using StreamChat.Core; +using StreamChat.Core.Configs; +using StreamChat.Core.Requests; +using StreamChat.Core.StatefulModels; +using UnityEngine.TestTools; + +namespace StreamChat.Tests.StatefulClient +{ + internal class MessageCacheWindowTests : BaseStateIntegrationTests + { + private static readonly MessageCacheWindow SmallWindow = new MessageCacheWindow(6, 3); + + [UnityTest] + public IEnumerator When_no_message_cache_window_configured_expect_no_removal() + => ConnectAndExecute(When_no_message_cache_window_configured_expect_no_removal_Async); + + private async Task When_no_message_cache_window_configured_expect_no_removal_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + var removedCount = 0; + channel.MessagesRemovedFromCache += (_, __) => removedCount++; + + await SendMessagesAsync(channel, 10); + + Assert.AreEqual(10, channel.Messages.Count); + Assert.AreEqual(0, removedCount); + } + + [UnityTest] + public IEnumerator When_message_count_exceeds_max_expect_trim_to_max_minus_discard() + => ConnectAndExecute(When_message_count_exceeds_max_expect_trim_to_max_minus_discard_Async); + + private async Task When_message_count_exceeds_max_expect_trim_to_max_minus_discard_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + + await SendMessagesAsync(channel, 7); + + Assert.AreEqual(3, channel.Messages.Count); + } + + [UnityTest] + public IEnumerator When_message_count_equals_max_expect_no_trim() + => ConnectAndExecute(When_message_count_equals_max_expect_no_trim_Async); + + private async Task When_message_count_equals_max_expect_no_trim_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + var removedCount = 0; + channel.MessagesRemovedFromCache += (_, __) => removedCount++; + + await SendMessagesAsync(channel, 6); + + Assert.AreEqual(6, channel.Messages.Count); + Assert.AreEqual(0, removedCount); + } + + [UnityTest] + public IEnumerator When_trimmed_expect_oldest_removed_and_newest_retained() + => ConnectAndExecute(When_trimmed_expect_oldest_removed_and_newest_retained_Async); + + private async Task When_trimmed_expect_oldest_removed_and_newest_retained_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + + var sent = await SendMessagesAsync(channel, 7); + var expectedIds = sent.Skip(4).Select(m => m.Id).ToList(); + var actualIds = channel.Messages.Select(m => m.Id).ToList(); + + CollectionAssert.AreEqual(expectedIds, actualIds); + } + + [UnityTest] + public IEnumerator When_trimmed_expect_removed_messages_untracked_from_cache() + => ConnectAndExecute(When_trimmed_expect_removed_messages_untracked_from_cache_Async); + + private async Task When_trimmed_expect_removed_messages_untracked_from_cache_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + + var sent = await SendMessagesAsync(channel, 7); + var removedIds = sent.Take(4).Select(m => m.Id).ToList(); + + foreach (var id in removedIds) + { + Assert.IsFalse(Client.InternalCache.Messages.TryGet(id, out _), + $"Removed message {id} should no longer be tracked in the cache."); + } + } + + [UnityTest] + public IEnumerator When_trimmed_expect_single_batched_event_with_all_removed_messages() + => ConnectAndExecute(When_trimmed_expect_single_batched_event_with_all_removed_messages_Async); + + private async Task When_trimmed_expect_single_batched_event_with_all_removed_messages_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + + var sent = await SendMessagesAsync(channel, 6); + var eventInvocations = 0; + IReadOnlyList removedBatch = null; + + channel.MessagesRemovedFromCache += (_, messages) => + { + eventInvocations++; + removedBatch = messages; + }; + + await channel.SendNewMessageAsync($"msg-trigger-{Guid.NewGuid()}"); + + Assert.AreEqual(1, eventInvocations); + Assert.NotNull(removedBatch); + Assert.AreEqual(4, removedBatch.Count); + CollectionAssert.AreEqual(sent.Select(m => m.Id).ToList(), removedBatch.Select(m => m.Id).ToList()); + } + + [UnityTest] + public IEnumerator When_message_triggers_trim_expect_MessageReceived_raised_before_removal() + => ConnectAndExecute(When_message_triggers_trim_expect_MessageReceived_raised_before_removal_Async); + + private async Task When_message_triggers_trim_expect_MessageReceived_raised_before_removal_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + + await SendMessagesAsync(channel, 6); + + var eventLog = new List(); + channel.MessageReceived += (_, msg) => eventLog.Add($"received:{msg.Id}"); + channel.MessagesRemovedFromCache += (_, messages) => + eventLog.Add($"removed:{string.Join(",", messages.Select(m => m.Id))}"); + + var trigger = await channel.SendNewMessageAsync($"msg-trigger-{Guid.NewGuid()}"); + + var receivedIndex = eventLog.FindIndex(e => e == $"received:{trigger.Id}"); + var removedIndex = eventLog.FindIndex(e => e.StartsWith("removed:")); + + Assert.Greater(receivedIndex, -1); + Assert.Greater(removedIndex, -1); + Assert.Less(receivedIndex, removedIndex); + } + + [UnityTest] + public IEnumerator When_removed_message_is_pinned_expect_it_stays_in_cache() + => ConnectAndExecute(When_removed_message_is_pinned_expect_it_stays_in_cache_Async); + + private async Task When_removed_message_is_pinned_expect_it_stays_in_cache_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + + var first = await channel.SendNewMessageAsync($"msg-0-{Guid.NewGuid()}"); + await first.PinAsync(); + await WaitWhileFalseAsync(() => channel.PinnedMessages.Any(m => m.Id == first.Id), + description: "pinned message to appear in PinnedMessages"); + + await SendMessagesAsync(channel, 6); + + Assert.IsFalse(channel.Messages.Any(m => m.Id == first.Id)); + Assert.IsTrue(Client.InternalCache.Messages.TryGet(first.Id, out _)); + } + + [UnityTest] + public IEnumerator When_pinned_message_removed_from_cache_expect_updates_still_applied() + => ConnectAndExecute(When_pinned_message_removed_from_cache_expect_updates_still_applied_Async); + + private async Task When_pinned_message_removed_from_cache_expect_updates_still_applied_Async() + { + var otherClient = await GetConnectedOtherClientAsync(); + var channel = await CreateUniqueTempChannelAsync(); + var otherChannel = await otherClient.GetOrCreateChannelWithIdAsync(channel.Type, channel.Id); + channel.OverrideMessageCacheWindow(SmallWindow); + + var first = await channel.SendNewMessageAsync($"msg-0-{Guid.NewGuid()}"); + await first.PinAsync(); + await WaitWhileFalseAsync(() => channel.PinnedMessages.Any(m => m.Id == first.Id), + description: "pinned message to appear in PinnedMessages"); + + await SendMessagesAsync(channel, 6); + + var pinned = channel.PinnedMessages.Single(m => m.Id == first.Id); + const string UpdatedText = "pinned-after-cache-removal"; + + var messageOnOther = otherChannel.Messages.Single(m => m.Id == first.Id); + await messageOnOther.UpdateAsync(new StreamUpdateMessageRequest { Text = UpdatedText }); + + await WaitWhileFalseAsync(() => pinned.Text == UpdatedText, + description: "pinned message instance to receive message.updated after cache removal"); + } + + [UnityTest] + public IEnumerator When_removed_message_is_thread_parent_expect_it_stays_in_cache() + => ConnectAndExecute(When_removed_message_is_thread_parent_expect_it_stays_in_cache_Async); + + private async Task When_removed_message_is_thread_parent_expect_it_stays_in_cache_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + + var parent = await channel.SendNewMessageAsync($"thread-parent-{Guid.NewGuid()}"); + await channel.SendNewMessageAsync(new StreamSendMessageRequest + { + ParentId = parent.Id, + ShowInChannel = false, + Text = "thread reply", + }); + + await Client.GetThreadAsync(parent.Id, replyLimit: 5, participantLimit: 5); + await SendMessagesAsync(channel, 6); + + Assert.IsFalse(channel.Messages.Any(m => m.Id == parent.Id)); + Assert.IsTrue(Client.InternalCache.Messages.TryGet(parent.Id, out _)); + } + + [UnityTest] + public IEnumerator When_LoadOlderMessagesAsync_called_expect_trimming_paused() + => ConnectAndExecute(When_LoadOlderMessagesAsync_called_expect_trimming_paused_Async); + + private async Task When_LoadOlderMessagesAsync_called_expect_trimming_paused_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + await SendMessagesAsync(channel, 3); + + await channel.LoadOlderMessagesAsync(); + + Assert.IsTrue(channel.IsMessageCacheTrimmingPaused); + } + + [UnityTest] + public IEnumerator When_trimming_paused_expect_no_removal_on_new_messages() + => ConnectAndExecute(When_trimming_paused_expect_no_removal_on_new_messages_Async); + + private async Task When_trimming_paused_expect_no_removal_on_new_messages_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + channel.PauseMessageCacheTrimming(); + + var removedCount = 0; + channel.MessagesRemovedFromCache += (_, __) => removedCount++; + + await SendMessagesAsync(channel, 10); + + Assert.AreEqual(10, channel.Messages.Count); + Assert.AreEqual(0, removedCount); + } + + [UnityTest] + public IEnumerator When_ResumeMessageCacheTrimming_called_expect_immediate_trim() + => ConnectAndExecute(When_ResumeMessageCacheTrimming_called_expect_immediate_trim_Async); + + private async Task When_ResumeMessageCacheTrimming_called_expect_immediate_trim_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + channel.PauseMessageCacheTrimming(); + + await SendMessagesAsync(channel, 10); + + var removedCount = 0; + channel.MessagesRemovedFromCache += (_, __) => removedCount++; + + channel.ResumeMessageCacheTrimming(); + + Assert.AreEqual(3, channel.Messages.Count); + Assert.AreEqual(1, removedCount); + } + + [UnityTest] + public IEnumerator When_OverrideMessageCacheWindow_called_expect_immediate_trim() + => ConnectAndExecute(When_OverrideMessageCacheWindow_called_expect_immediate_trim_Async); + + private async Task When_OverrideMessageCacheWindow_called_expect_immediate_trim_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + await SendMessagesAsync(channel, 10); + + channel.OverrideMessageCacheWindow(SmallWindow); + + Assert.AreEqual(3, channel.Messages.Count); + } + + [UnityTest] + public IEnumerator When_channel_override_is_null_expect_unlimited_despite_client_default() + => ConnectAndExecute(When_channel_override_is_null_expect_unlimited_despite_client_default_Async); + + private async Task When_channel_override_is_null_expect_unlimited_despite_client_default_Async() + { + var config = Client.LowLevelClient.Config; + var previousDefault = config.DefaultMessageCacheWindow; + try + { + config.DefaultMessageCacheWindow = SmallWindow; + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(null); + + var removedCount = 0; + channel.MessagesRemovedFromCache += (_, __) => removedCount++; + + await SendMessagesAsync(channel, 7); + + Assert.IsTrue(channel.HasMessageCacheWindowOverride); + Assert.AreEqual(7, channel.Messages.Count); + Assert.AreEqual(0, removedCount); + } + finally + { + config.DefaultMessageCacheWindow = previousDefault; + } + } + + [UnityTest] + public IEnumerator When_ClearMessageCacheWindowOverride_called_expect_client_default_reapplied() + => ConnectAndExecute(When_ClearMessageCacheWindowOverride_called_expect_client_default_reapplied_Async); + + private async Task When_ClearMessageCacheWindowOverride_called_expect_client_default_reapplied_Async() + { + var config = Client.LowLevelClient.Config; + var previousDefault = config.DefaultMessageCacheWindow; + try + { + config.DefaultMessageCacheWindow = SmallWindow; + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(null); + await SendMessagesAsync(channel, 7); + + channel.ClearMessageCacheWindowOverride(); + + Assert.IsFalse(channel.HasMessageCacheWindowOverride); + Assert.AreEqual(3, channel.Messages.Count); + } + finally + { + config.DefaultMessageCacheWindow = previousDefault; + } + } + + [UnityTest] + public IEnumerator When_messages_removed_expect_LoadOlderMessagesAsync_refetches_them() + => ConnectAndExecute(When_messages_removed_expect_LoadOlderMessagesAsync_refetches_them_Async); + + private async Task When_messages_removed_expect_LoadOlderMessagesAsync_refetches_them_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + + var sent = await SendMessagesAsync(channel, 7); + var removedIds = sent.Take(4).Select(m => m.Id).ToHashSet(); + var survivingIds = sent.Skip(4).Select(m => m.Id).ToList(); + + channel.ResumeMessageCacheTrimming(); + await channel.LoadOlderMessagesAsync(); + + var messageIds = channel.Messages.Select(m => m.Id).ToList(); + foreach (var id in removedIds) + { + Assert.Contains(id, messageIds); + } + + CollectionAssert.IsOrdered(channel.Messages.Select(m => m.CreatedAt).ToList()); + Assert.AreEqual(survivingIds.Last(), messageIds.Last()); + } + + [UnityTest] + public IEnumerator When_channel_query_returns_more_than_max_expect_no_removal_until_next_live_message() + => ConnectAndExecute(When_channel_query_returns_more_than_max_expect_no_removal_until_next_live_message_Async); + + private async Task When_channel_query_returns_more_than_max_expect_no_removal_until_next_live_message_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindow); + channel.PauseMessageCacheTrimming(); + + await SendMessagesAsync(channel, 10); + + var removedCount = 0; + channel.MessagesRemovedFromCache += (_, __) => removedCount++; + + await channel.LoadOlderMessagesAsync(); + + Assert.Greater(channel.Messages.Count, SmallWindow.MaxMessages); + Assert.AreEqual(0, removedCount); + + channel.ResumeMessageCacheTrimming(); + Assert.AreEqual(3, channel.Messages.Count); + Assert.AreEqual(1, removedCount); + + removedCount = 0; + await SendMessagesAsync(channel, 4); + + Assert.AreEqual(3, channel.Messages.Count); + Assert.AreEqual(1, removedCount); + } + + [Test] + public void MessageCacheWindow_rejects_invalid_arguments() + { + Assert.Throws(() => new MessageCacheWindow(0, 1)); + Assert.Throws(() => new MessageCacheWindow(-1, 1)); + Assert.Throws(() => new MessageCacheWindow(10, 0)); + Assert.Throws(() => new MessageCacheWindow(10, -1)); + Assert.Throws(() => new MessageCacheWindow(10, 10)); + Assert.Throws(() => new MessageCacheWindow(10, 11)); + } + + [Test] + public void MessageCacheWindow_Recommended_is_500_100() + { + Assert.AreEqual(500, MessageCacheWindow.Recommended.MaxMessages); + Assert.AreEqual(100, MessageCacheWindow.Recommended.DiscardBatchSize); + } + + private static async Task> SendMessagesAsync(IStreamChannel channel, int count) + { + var messages = new List(count); + for (var i = 0; i < count; i++) + { + messages.Add(await channel.SendNewMessageAsync($"msg-{i}-{Guid.NewGuid()}")); + } + + return messages; + } + } +} +#endif diff --git a/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs.meta b/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs.meta new file mode 100644 index 00000000..53272653 --- /dev/null +++ b/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs.meta @@ -0,0 +1,3 @@ +fileFormatVersion: 2 +guid: e1b7f6c35f8d6d0e2c9a4a7f5e3d6c8b +timeCreated: 1755079300 From 8519d3a5a37d0a7a0edddc613eadcddf3d3a08a2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Sierpi=C5=84ski?= <33436839+sierpinskid@users.noreply.github.com> Date: Thu, 13 Aug 2026 14:58:09 +0200 Subject: [PATCH 2/4] Add AbsoluteMaxMessages -> above this limit messages are cleared anyway + reduce iteration complexity --- .../Core/Configs/MessageCacheWindow.cs | 49 ++++++++- .../StreamChat/Core/Helpers/HashSetPool.cs | 41 +++++++ .../Core/Helpers/HashSetPool.cs.meta | 3 + .../Core/Helpers/HashSetPoolScope.cs | 24 ++++ .../Core/Helpers/HashSetPoolScope.cs.meta | 3 + .../StreamChat/Core/Helpers/ListPool.cs | 7 +- .../StreamChat/Core/Helpers/ListPoolScope.cs | 24 ++++ .../Core/Helpers/ListPoolScope.cs.meta | 3 + .../Core/State/Caches/CacheRepository.cs | 51 +++++++++ .../Core/State/Caches/ICacheRepository.cs | 6 + .../Core/StatefulModels/IStreamChannel.cs | 23 ++-- .../Core/StatefulModels/StreamChannel.cs | 104 ++++++++++-------- .../StatefulClient/MessageCacheWindowTests.cs | 74 ++++++++++++- 13 files changed, 351 insertions(+), 61 deletions(-) create mode 100644 Assets/Plugins/StreamChat/Core/Helpers/HashSetPool.cs create mode 100644 Assets/Plugins/StreamChat/Core/Helpers/HashSetPool.cs.meta create mode 100644 Assets/Plugins/StreamChat/Core/Helpers/HashSetPoolScope.cs create mode 100644 Assets/Plugins/StreamChat/Core/Helpers/HashSetPoolScope.cs.meta create mode 100644 Assets/Plugins/StreamChat/Core/Helpers/ListPoolScope.cs create mode 100644 Assets/Plugins/StreamChat/Core/Helpers/ListPoolScope.cs.meta diff --git a/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs b/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs index 6834bbb7..281e2fce 100644 --- a/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs +++ b/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs @@ -7,12 +7,13 @@ namespace StreamChat.Core.Configs /// Assign to or /// . /// Trimming runs in batches of once the count exceeds - /// . + /// , or once it exceeds while + /// trimming is paused. /// public sealed class MessageCacheWindow { - /// Keep up to 500 messages; remove 100 at a time when over the limit. - public static readonly MessageCacheWindow Recommended = new MessageCacheWindow(500, 100); + /// Keep up to 500 messages (2000 while paused); remove 100 at a time when over the limit. + public static readonly MessageCacheWindow Recommended = new MessageCacheWindow(500, 100, 2000); /// Trimming starts when exceeds this count. public int MaxMessages { get; } @@ -20,7 +21,22 @@ public sealed class MessageCacheWindow /// How many messages to remove per trim. Must be less than . public int DiscardBatchSize { get; } + /// + /// Upper bound that applies while + /// is true - for example while the user reads history loaded by + /// . Pausing widens the window to this + /// value instead of disabling it, so a channel stays bounded even if + /// is never called. + /// Must be greater than or equal to . Defaults to 4x . + /// + public int AbsoluteMaxMessages { get; } + public MessageCacheWindow(int maxMessages, int discardBatchSize) + : this(maxMessages, discardBatchSize, GetDefaultAbsoluteMaxMessages(maxMessages)) + { + } + + public MessageCacheWindow(int maxMessages, int discardBatchSize, int absoluteMaxMessages) { if (maxMessages <= 0) { @@ -41,11 +57,36 @@ public MessageCacheWindow(int maxMessages, int discardBatchSize) + "otherwise a single trim would remove every message."); } + if (absoluteMaxMessages < maxMessages) + { + throw new ArgumentOutOfRangeException(nameof(absoluteMaxMessages), absoluteMaxMessages, + $"{nameof(absoluteMaxMessages)} must be greater than or equal to {nameof(maxMessages)} " + + $"({maxMessages})."); + } + MaxMessages = maxMessages; DiscardBatchSize = discardBatchSize; + AbsoluteMaxMessages = absoluteMaxMessages; } public override string ToString() - => $"MessageCacheWindow - MaxMessages: {MaxMessages}, DiscardBatchSize: {DiscardBatchSize}"; + => $"MessageCacheWindow - MaxMessages: {MaxMessages}, DiscardBatchSize: {DiscardBatchSize}, " + + $"AbsoluteMaxMessages: {AbsoluteMaxMessages}"; + + private const int DefaultAbsoluteMaxMessagesMultiplier = 4; + + // Invalid values are passed through so the constructor reports the real problem instead of + // a derived one. + private static int GetDefaultAbsoluteMaxMessages(int maxMessages) + { + if (maxMessages <= 0) + { + return maxMessages; + } + + return maxMessages > int.MaxValue / DefaultAbsoluteMaxMessagesMultiplier + ? int.MaxValue + : maxMessages * DefaultAbsoluteMaxMessagesMultiplier; + } } } diff --git a/Assets/Plugins/StreamChat/Core/Helpers/HashSetPool.cs b/Assets/Plugins/StreamChat/Core/Helpers/HashSetPool.cs new file mode 100644 index 00000000..4164c37a --- /dev/null +++ b/Assets/Plugins/StreamChat/Core/Helpers/HashSetPool.cs @@ -0,0 +1,41 @@ +using System; +using System.Collections.Generic; + +namespace StreamChat.Core.Helpers +{ + internal static class HashSetPool + { + public static HashSet Rent() + { + if (Pool.Count > 0) + { + return Pool.Pop(); + } + + return new HashSet(); + } + + public static void Release(HashSet set) + { + if (set == null) + { + throw new ArgumentNullException(nameof(set)); + } + + // HashSet has no Capacity accessor, so the pre-clear count stands in for how large the + // buckets grew. See the same guard in ListPool. + var isOversized = set.Count > MaxRetainedCount; + + set.Clear(); + + if (!isOversized && Pool.Count < MaxPoolSize) + { + Pool.Push(set); + } + } + + private const int MaxPoolSize = 128; + private const int MaxRetainedCount = 4096; + private static readonly Stack> Pool = new Stack>(); + } +} diff --git a/Assets/Plugins/StreamChat/Core/Helpers/HashSetPool.cs.meta b/Assets/Plugins/StreamChat/Core/Helpers/HashSetPool.cs.meta new file mode 100644 index 00000000..1b29547f --- /dev/null +++ b/Assets/Plugins/StreamChat/Core/Helpers/HashSetPool.cs.meta @@ -0,0 +1,3 @@ +fileFormatVersion: 2 +guid: a3f8d1c65b7e2f94a8d3c6b1e5f7a2d9 +timeCreated: 1755079200 diff --git a/Assets/Plugins/StreamChat/Core/Helpers/HashSetPoolScope.cs b/Assets/Plugins/StreamChat/Core/Helpers/HashSetPoolScope.cs new file mode 100644 index 00000000..93d75e84 --- /dev/null +++ b/Assets/Plugins/StreamChat/Core/Helpers/HashSetPoolScope.cs @@ -0,0 +1,24 @@ +using System; +using System.Collections.Generic; + +namespace StreamChat.Core.Helpers +{ + internal readonly struct HashSetPoolScope : IDisposable + { + public HashSetPoolScope(out HashSet set) + { + _set = HashSetPool.Rent(); + set = _set; + } + + public void Dispose() + { + if (_set != null) + { + HashSetPool.Release(_set); + } + } + + private readonly HashSet _set; + } +} diff --git a/Assets/Plugins/StreamChat/Core/Helpers/HashSetPoolScope.cs.meta b/Assets/Plugins/StreamChat/Core/Helpers/HashSetPoolScope.cs.meta new file mode 100644 index 00000000..761b83d5 --- /dev/null +++ b/Assets/Plugins/StreamChat/Core/Helpers/HashSetPoolScope.cs.meta @@ -0,0 +1,3 @@ +fileFormatVersion: 2 +guid: e5b1a7d34c9f2e86b1d7a4c8f3e9b2d6 +timeCreated: 1755079200 diff --git a/Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs b/Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs index 3b075530..c8888301 100644 --- a/Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs +++ b/Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs @@ -22,15 +22,20 @@ public static void Release(List list) throw new ArgumentNullException(nameof(list)); } + // Clear() keeps the backing array, so a one-off bulk operation would otherwise retain an + // oversized buffer for the rest of the session. Dropping it costs one allocation later. + var isOversized = list.Capacity > MaxRetainedCapacity; + list.Clear(); - if (Pool.Count < MaxPoolSize) + if (!isOversized && Pool.Count < MaxPoolSize) { Pool.Push(list); } } private const int MaxPoolSize = 128; + private const int MaxRetainedCapacity = 4096; private static readonly Stack> Pool = new Stack>(); } } diff --git a/Assets/Plugins/StreamChat/Core/Helpers/ListPoolScope.cs b/Assets/Plugins/StreamChat/Core/Helpers/ListPoolScope.cs new file mode 100644 index 00000000..37734a9e --- /dev/null +++ b/Assets/Plugins/StreamChat/Core/Helpers/ListPoolScope.cs @@ -0,0 +1,24 @@ +using System; +using System.Collections.Generic; + +namespace StreamChat.Core.Helpers +{ + internal readonly struct ListPoolScope : IDisposable + { + public ListPoolScope(out List list) + { + _list = ListPool.Rent(); + list = _list; + } + + public void Dispose() + { + if (_list != null) + { + ListPool.Release(_list); + } + } + + private readonly List _list; + } +} diff --git a/Assets/Plugins/StreamChat/Core/Helpers/ListPoolScope.cs.meta b/Assets/Plugins/StreamChat/Core/Helpers/ListPoolScope.cs.meta new file mode 100644 index 00000000..bac9a62a --- /dev/null +++ b/Assets/Plugins/StreamChat/Core/Helpers/ListPoolScope.cs.meta @@ -0,0 +1,3 @@ +fileFormatVersion: 2 +guid: b7e2c4f81a9d4e3b8c5f7a2d9e1b3c6d +timeCreated: 1755079200 diff --git a/Assets/Plugins/StreamChat/Core/State/Caches/CacheRepository.cs b/Assets/Plugins/StreamChat/Core/State/Caches/CacheRepository.cs index e696c65a..d4304197 100644 --- a/Assets/Plugins/StreamChat/Core/State/Caches/CacheRepository.cs +++ b/Assets/Plugins/StreamChat/Core/State/Caches/CacheRepository.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using StreamChat.Core.Helpers; using StreamChat.Libs.Utils; namespace StreamChat.Core.State.Caches @@ -160,6 +161,56 @@ public void Remove(TStatefulModel trackedObject) Untracked?.Invoke(trackedObject); } + public void RemoveMany(IReadOnlyList trackedObjects) + { + if (trackedObjects == null) + { + throw new ArgumentNullException(nameof(trackedObjects)); + } + + if (trackedObjects.Count == 0) + { + return; + } + + using (new HashSetPoolScope(out var tempRemovedIds)) + using (new ListPoolScope(out var tempRemoved)) + { + for (var i = 0; i < trackedObjects.Count; i++) + { + var trackedObject = trackedObjects[i]; + if (trackedObject.UniqueId.IsNullOrEmpty()) + { + throw new ArgumentException($"{trackedObject.UniqueId} cannot be empty"); + } + + // Only untrack the exact instance this repository holds. A newer instance for the + // same id must survive, and duplicates in the input must not raise Untracked twice. + if (!_statefulModelById.TryGetValue(trackedObject.UniqueId, out var tracked) + || !ReferenceEquals(tracked, trackedObject)) + { + continue; + } + + _statefulModelById.Remove(trackedObject.UniqueId); + tempRemovedIds.Add(trackedObject.UniqueId); + tempRemoved.Add(trackedObject); + } + + if (tempRemoved.Count == 0) + { + return; + } + + _statefulModels.RemoveAll(_ => tempRemovedIds.Contains(_.UniqueId)); + + for (var i = 0; i < tempRemoved.Count; i++) + { + Untracked?.Invoke(tempRemoved[i]); + } + } + } + internal delegate TStatefulModel ConstructorHandler(string uniqueId); internal CacheRepository(ConstructorHandler constructor, ICache cache) diff --git a/Assets/Plugins/StreamChat/Core/State/Caches/ICacheRepository.cs b/Assets/Plugins/StreamChat/Core/State/Caches/ICacheRepository.cs index 9d683f1f..1324155b 100644 --- a/Assets/Plugins/StreamChat/Core/State/Caches/ICacheRepository.cs +++ b/Assets/Plugins/StreamChat/Core/State/Caches/ICacheRepository.cs @@ -64,5 +64,11 @@ TType CreateOrUpdate5(TDto dto, out bool wasCreated) where TType : class, TTrackedObject, IStreamStatefulModel, IUpdateableFrom5; void Remove(TTrackedObject trackedObject); + + /// + /// Removes multiple tracked objects in a single pass. Prefer this over calling + /// in a loop - removing N objects one by one is O(N * repository size). + /// + void RemoveMany(IReadOnlyList trackedObjects); } } \ No newline at end of file diff --git a/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs b/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs index 62b44422..51f215c5 100644 --- a/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs +++ b/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs @@ -325,7 +325,9 @@ public interface IStreamChannel : IStreamStatefulModel /// /// Active cache limit for this channel. null = unlimited (default). - /// When over , oldest messages are removed from + /// When over (or + /// while + /// is true), oldest messages are removed from /// and the cache. Server history is unchanged — /// can reload them. Pinned messages and open threads may /// stay in the cache. always fires before @@ -343,30 +345,35 @@ public interface IStreamChannel : IStreamStatefulModel /// /// Set a cache limit for this channel only. Pass null for unlimited on this channel. - /// Trims immediately unless is true. + /// Trims immediately, against when + /// is true. /// void OverrideMessageCacheWindow(Configs.MessageCacheWindow window); /// /// Remove the per-channel limit and use again. - /// Trims immediately unless trimming is paused. + /// Trims immediately, against the paused limit when trimming is paused. /// void ClearMessageCacheWindowOverride(); /// - /// Whether cache trimming is paused. pauses automatically. - /// Memory is unbounded while paused. + /// Whether the wider, paused cache limit is in effect. pauses + /// automatically so paged-in history is not removed while the user reads it. Trimming is not disabled + /// while paused - the limit becomes . /// bool IsMessageCacheTrimmingPaused { get; } /// - /// Pause cache trimming (e.g. while the user scrolls through loaded history). - /// does this for you. + /// Widen the cache limit to + /// (e.g. while the user scrolls through loaded history). + /// does this for you. /// void PauseMessageCacheTrimming(); /// - /// Resume trimming and remove excess messages now. Call when the user returns to the newest messages. + /// Restore the limit and remove excess messages now. + /// Optional - the cache stays bounded either way. Call it when the user returns to the newest messages + /// to release the memory held by paged-in history sooner. /// void ResumeMessageCacheTrimming(); diff --git a/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs b/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs index f2047742..adafab13 100644 --- a/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs +++ b/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs @@ -249,29 +249,41 @@ public async Task SendNewMessageAsync(StreamSendMessageRequest s public async Task LoadOlderMessagesAsync() { + var wasTrimmingPaused = _isMessageCacheTrimmingPaused; + // Pause before the request so live messages cannot trim the page being loaded. PauseMessageCacheTrimming(); - var oldestMessage = _messages.OrderBy(_ => _.CreatedAt).FirstOrDefault(); - - var request = new ChannelGetOrCreateRequestInternalDTO + try { - //StreamTodo: presence could be optional in config - Presence = true, - State = true, - Watch = true, - }; + var oldestMessage = _messages.OrderBy(_ => _.CreatedAt).FirstOrDefault(); - if (oldestMessage != null) - { - request.Messages = new MessagePaginationParamsRequestInternalDTO + var request = new ChannelGetOrCreateRequestInternalDTO { - IdLt = oldestMessage.Id, + //StreamTodo: presence could be optional in config + Presence = true, + State = true, + Watch = true, }; - } - var response = await LowLevelClient.InternalChannelApi.GetOrCreateChannelAsync(Type, Id, request); - Cache.TryCreateOrUpdate(response); + if (oldestMessage != null) + { + request.Messages = new MessagePaginationParamsRequestInternalDTO + { + IdLt = oldestMessage.Id, + }; + } + + var response = await LowLevelClient.InternalChannelApi.GetOrCreateChannelAsync(Type, Id, request); + Cache.TryCreateOrUpdate(response); + } + catch + { + // Nothing was paged in, so don't leave trimming suppressed for the rest of the session. + _isMessageCacheTrimmingPaused = wasTrimmingPaused; + TrimMessageCacheIfNeeded(); + throw; + } } public MessageCacheWindow MessageCacheWindow @@ -1047,7 +1059,17 @@ private void InternalTruncateMessages(DateTimeOffset? deleteBeforeCreatedAt = nu private void TrimMessageCacheIfNeeded() { var window = MessageCacheWindow; - if (window == null || _isMessageCacheTrimmingPaused || _messages.Count <= window.MaxMessages) + if (window == null) + { + return; + } + + // Pausing widens the window instead of disabling it, so a channel stays bounded even if + // ResumeMessageCacheTrimming is never called after LoadOlderMessagesAsync. + var effectiveMaxMessages + = _isMessageCacheTrimmingPaused ? window.AbsoluteMaxMessages : window.MaxMessages; + + if (_messages.Count <= effectiveMaxMessages) { return; } @@ -1057,49 +1079,45 @@ private void TrimMessageCacheIfNeeded() SortMessagesByCreatedAt(); } - var targetCount = window.MaxMessages - window.DiscardBatchSize; + var targetCount = effectiveMaxMessages - window.DiscardBatchSize; var removeCount = _messages.Count - targetCount; - var tempDeleteCandidates = ListPool.Rent(); - try - { - for (var i = 0; i < removeCount; i++) - { - tempDeleteCandidates.Add(_messages[i]); - } + var handler = MessagesRemovedFromCache; - _messages.RemoveRange(0, removeCount); + // Not pooled - this is handed to subscribers, so the SDK does not control its lifetime. + var removed = handler == null ? null : new List(removeCount); - for (var i = 0; i < tempDeleteCandidates.Count; i++) + using (new ListPoolScope(out var tempUntrackCandidates)) + using (new HashSetPoolScope(out var tempPinnedMessageIds)) + { + for (var i = 0; i < _pinnedMessages.Count; i++) { - var message = tempDeleteCandidates[i]; - if (!IsRetainedByOtherState(message)) - { - Cache.Messages.Remove(message); - } + tempPinnedMessageIds.Add(_pinnedMessages[i].Id); } - if (MessagesRemovedFromCache != null) + for (var i = 0; i < removeCount; i++) { - var removed = new List(tempDeleteCandidates.Count); - for (var i = 0; i < tempDeleteCandidates.Count; i++) + var message = _messages[i]; + removed?.Add(message); + + if (!IsRetainedByOtherState(message, tempPinnedMessageIds)) { - removed.Add(tempDeleteCandidates[i]); + tempUntrackCandidates.Add(message); } - - MessagesRemovedFromCache.Invoke(this, removed); } + + _messages.RemoveRange(0, removeCount); + Cache.Messages.RemoveMany(tempUntrackCandidates); } - finally - { - ListPool.Release(tempDeleteCandidates); - } + + // Raised after the pooled buffers are returned so subscribers can trim or send safely. + handler?.Invoke(this, removed); } // Keep in cache when pinned or referenced by a tracked thread. - private bool IsRetainedByOtherState(StreamMessage message) + private bool IsRetainedByOtherState(StreamMessage message, HashSet pinnedMessageIds) { - if (_pinnedMessages.ContainsNoAlloc(message)) + if (pinnedMessageIds.Contains(message.Id)) { return true; } diff --git a/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs b/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs index 8b70022a..d072ca77 100644 --- a/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs +++ b/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs @@ -17,6 +17,9 @@ internal class MessageCacheWindowTests : BaseStateIntegrationTests { private static readonly MessageCacheWindow SmallWindow = new MessageCacheWindow(6, 3); + // Absolute limit low enough to be reached by sending a handful of messages while paused. + private static readonly MessageCacheWindow SmallWindowWithLowAbsoluteMax = new MessageCacheWindow(4, 2, 8); + [UnityTest] public IEnumerator When_no_message_cache_window_configured_expect_no_removal() => ConnectAndExecute(When_no_message_cache_window_configured_expect_no_removal_Async); @@ -120,10 +123,13 @@ private async Task When_trimmed_expect_single_batched_event_with_all_removed_mes await channel.SendNewMessageAsync($"msg-trigger-{Guid.NewGuid()}"); + // 6 sent + 1 trigger = 7 > MaxMessages(6), trimmed down to MaxMessages - DiscardBatchSize = 3, + // so the 4 oldest are removed in a single batch, oldest first. Assert.AreEqual(1, eventInvocations); Assert.NotNull(removedBatch); Assert.AreEqual(4, removedBatch.Count); - CollectionAssert.AreEqual(sent.Select(m => m.Id).ToList(), removedBatch.Select(m => m.Id).ToList()); + CollectionAssert.AreEqual(sent.Take(4).Select(m => m.Id).ToList(), + removedBatch.Select(m => m.Id).ToList()); } [UnityTest] @@ -258,6 +264,54 @@ private async Task When_trimming_paused_expect_no_removal_on_new_messages_Async( Assert.AreEqual(0, removedCount); } + [UnityTest] + public IEnumerator When_trimming_paused_expect_trim_once_absolute_max_exceeded() + => ConnectAndExecute(When_trimming_paused_expect_trim_once_absolute_max_exceeded_Async); + + private async Task When_trimming_paused_expect_trim_once_absolute_max_exceeded_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindowWithLowAbsoluteMax); + channel.PauseMessageCacheTrimming(); + + var removedCount = 0; + channel.MessagesRemovedFromCache += (_, __) => removedCount++; + + // Above MaxMessages(4) but still within AbsoluteMaxMessages(8), so pausing holds the messages. + await SendMessagesAsync(channel, 8); + + Assert.AreEqual(8, channel.Messages.Count); + Assert.AreEqual(0, removedCount); + + // Crossing AbsoluteMaxMessages trims even though trimming is still paused, down to + // AbsoluteMaxMessages - DiscardBatchSize. This is what keeps a never-resumed channel bounded. + var sent = await SendMessagesAsync(channel, 1); + + Assert.IsTrue(channel.IsMessageCacheTrimmingPaused); + Assert.AreEqual(6, channel.Messages.Count); + Assert.AreEqual(1, removedCount); + Assert.AreEqual(sent.Single().Id, channel.Messages.Last().Id); + } + + [UnityTest] + public IEnumerator When_trimming_paused_and_bounded_expect_resume_restores_max_messages() + => ConnectAndExecute(When_trimming_paused_and_bounded_expect_resume_restores_max_messages_Async); + + private async Task When_trimming_paused_and_bounded_expect_resume_restores_max_messages_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindowWithLowAbsoluteMax); + channel.PauseMessageCacheTrimming(); + + await SendMessagesAsync(channel, 8); + Assert.AreEqual(8, channel.Messages.Count); + + channel.ResumeMessageCacheTrimming(); + + Assert.IsFalse(channel.IsMessageCacheTrimmingPaused); + Assert.AreEqual(2, channel.Messages.Count); + } + [UnityTest] public IEnumerator When_ResumeMessageCacheTrimming_called_expect_immediate_trim() => ConnectAndExecute(When_ResumeMessageCacheTrimming_called_expect_immediate_trim_Async); @@ -299,7 +353,7 @@ public IEnumerator When_channel_override_is_null_expect_unlimited_despite_client private async Task When_channel_override_is_null_expect_unlimited_despite_client_default_Async() { - var config = Client.LowLevelClient.Config; + var config = Client.InternalLowLevelClient.Config; var previousDefault = config.DefaultMessageCacheWindow; try { @@ -328,7 +382,7 @@ public IEnumerator When_ClearMessageCacheWindowOverride_called_expect_client_def private async Task When_ClearMessageCacheWindowOverride_called_expect_client_default_reapplied_Async() { - var config = Client.LowLevelClient.Config; + var config = Client.InternalLowLevelClient.Config; var previousDefault = config.DefaultMessageCacheWindow; try { @@ -358,7 +412,7 @@ private async Task When_messages_removed_expect_LoadOlderMessagesAsync_refetches channel.OverrideMessageCacheWindow(SmallWindow); var sent = await SendMessagesAsync(channel, 7); - var removedIds = sent.Take(4).Select(m => m.Id).ToHashSet(); + var removedIds = new HashSet(sent.Take(4).Select(m => m.Id)); var survivingIds = sent.Skip(4).Select(m => m.Id).ToList(); channel.ResumeMessageCacheTrimming(); @@ -414,13 +468,23 @@ public void MessageCacheWindow_rejects_invalid_arguments() Assert.Throws(() => new MessageCacheWindow(10, -1)); Assert.Throws(() => new MessageCacheWindow(10, 10)); Assert.Throws(() => new MessageCacheWindow(10, 11)); + Assert.Throws(() => new MessageCacheWindow(10, 5, 9)); + Assert.DoesNotThrow(() => new MessageCacheWindow(10, 5, 10)); + } + + [Test] + public void When_absolute_max_not_specified_expect_default_of_four_times_max_messages() + { + Assert.AreEqual(40, new MessageCacheWindow(10, 5).AbsoluteMaxMessages); + Assert.AreEqual(int.MaxValue, new MessageCacheWindow(int.MaxValue, 1).AbsoluteMaxMessages); } [Test] - public void MessageCacheWindow_Recommended_is_500_100() + public void MessageCacheWindow_Recommended_is_500_100_2000() { Assert.AreEqual(500, MessageCacheWindow.Recommended.MaxMessages); Assert.AreEqual(100, MessageCacheWindow.Recommended.DiscardBatchSize); + Assert.AreEqual(2000, MessageCacheWindow.Recommended.AbsoluteMaxMessages); } private static async Task> SendMessagesAsync(IStreamChannel channel, int count) From 60ff44c07b878789f5d71a195906ee0b083b3d26 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Sierpi=C5=84ski?= <33436839+sierpinskid@users.noreply.github.com> Date: Thu, 13 Aug 2026 15:38:12 +0200 Subject: [PATCH 3/4] Add limiting message growth when cache trimming is paused --- .../Core/Configs/MessageCacheWindow.cs | 49 ++++++---- .../Core/StatefulModels/IStreamChannel.cs | 54 +++++++---- .../Core/StatefulModels/StreamChannel.cs | 65 +++++++++++-- .../StatefulClient/MessageCacheWindowTests.cs | 97 ++++++++++++++----- 4 files changed, 200 insertions(+), 65 deletions(-) diff --git a/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs b/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs index 281e2fce..7cab5a93 100644 --- a/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs +++ b/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs @@ -7,12 +7,13 @@ namespace StreamChat.Core.Configs /// Assign to or /// . /// Trimming runs in batches of once the count exceeds - /// , or once it exceeds while - /// trimming is paused. + /// . Nothing is ever removed while + /// is true - growth is + /// bounded by instead. /// public sealed class MessageCacheWindow { - /// Keep up to 500 messages (2000 while paused); remove 100 at a time when over the limit. + /// Keep up to 500 messages; remove 100 at a time when over the limit; stop paging in history at 2000. public static readonly MessageCacheWindow Recommended = new MessageCacheWindow(500, 100, 2000); /// Trimming starts when exceeds this count. @@ -22,21 +23,29 @@ public sealed class MessageCacheWindow public int DiscardBatchSize { get; } /// - /// Upper bound that applies while - /// is true - for example while the user reads history loaded by - /// . Pausing widens the window to this - /// value instead of disabling it, so a channel stays bounded even if - /// is never called. + /// How large may grow before + /// stops paging in history. This is + /// the total message count, not a separate budget for paged-in messages. + /// It only comes into play while + /// is true - which + /// sets automatically - because + /// otherwise already keeps the channel smaller than this. + /// Reaching it never removes anything: paged-in history is exactly what a trim would delete. + /// Loading simply stops until + /// is called. Live messages are + /// still appended, so a channel that is never resumed can grow past this value; the SDK logs a warning + /// once when that happens. /// Must be greater than or equal to . Defaults to 4x . + /// Setting it equal to means "never page in history beyond the normal limit". /// - public int AbsoluteMaxMessages { get; } + public int MaxHistoryMessages { get; } public MessageCacheWindow(int maxMessages, int discardBatchSize) - : this(maxMessages, discardBatchSize, GetDefaultAbsoluteMaxMessages(maxMessages)) + : this(maxMessages, discardBatchSize, GetDefaultMaxHistoryMessages(maxMessages)) { } - public MessageCacheWindow(int maxMessages, int discardBatchSize, int absoluteMaxMessages) + public MessageCacheWindow(int maxMessages, int discardBatchSize, int maxHistoryMessages) { if (maxMessages <= 0) { @@ -57,36 +66,36 @@ public MessageCacheWindow(int maxMessages, int discardBatchSize, int absoluteMax + "otherwise a single trim would remove every message."); } - if (absoluteMaxMessages < maxMessages) + if (maxHistoryMessages < maxMessages) { - throw new ArgumentOutOfRangeException(nameof(absoluteMaxMessages), absoluteMaxMessages, - $"{nameof(absoluteMaxMessages)} must be greater than or equal to {nameof(maxMessages)} " + throw new ArgumentOutOfRangeException(nameof(maxHistoryMessages), maxHistoryMessages, + $"{nameof(maxHistoryMessages)} must be greater than or equal to {nameof(maxMessages)} " + $"({maxMessages})."); } MaxMessages = maxMessages; DiscardBatchSize = discardBatchSize; - AbsoluteMaxMessages = absoluteMaxMessages; + MaxHistoryMessages = maxHistoryMessages; } public override string ToString() => $"MessageCacheWindow - MaxMessages: {MaxMessages}, DiscardBatchSize: {DiscardBatchSize}, " - + $"AbsoluteMaxMessages: {AbsoluteMaxMessages}"; + + $"MaxHistoryMessages: {MaxHistoryMessages}"; - private const int DefaultAbsoluteMaxMessagesMultiplier = 4; + private const int DefaultMaxHistoryMessagesMultiplier = 4; // Invalid values are passed through so the constructor reports the real problem instead of // a derived one. - private static int GetDefaultAbsoluteMaxMessages(int maxMessages) + private static int GetDefaultMaxHistoryMessages(int maxMessages) { if (maxMessages <= 0) { return maxMessages; } - return maxMessages > int.MaxValue / DefaultAbsoluteMaxMessagesMultiplier + return maxMessages > int.MaxValue / DefaultMaxHistoryMessagesMultiplier ? int.MaxValue - : maxMessages * DefaultAbsoluteMaxMessagesMultiplier; + : maxMessages * DefaultMaxHistoryMessagesMultiplier; } } } diff --git a/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs b/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs index 51f215c5..fee011a0 100644 --- a/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs +++ b/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs @@ -319,19 +319,33 @@ public interface IStreamChannel : IStreamStatefulModel /// /// Load next portion of older messages. Older messages will be prepended to the list. - /// Note that loading older messages does NOT trigger the event + /// Note that loading older messages does NOT trigger the event. + /// Calling this pauses cache trimming (see ) so the + /// loaded page is not removed while the user reads it. If + /// is true this returns without loading anything - + /// stop showing your "load more" affordance and prompt the user to jump back to the newest messages, + /// which is where belongs. /// Task LoadOlderMessagesAsync(); + /// + /// true when has reached + /// and + /// will therefore load nothing. Always false when + /// is null. + /// Check this before offering "load more" so the user is not left waiting on a call that + /// cannot return anything. clears it. + /// + bool HasReachedMaxHistoryMessages { get; } + /// /// Active cache limit for this channel. null = unlimited (default). - /// When over (or - /// while - /// is true), oldest messages are removed from - /// and the cache. Server history is unchanged — + /// Once exceeds the oldest + /// messages are removed from and the cache. Server history is unchanged — /// can reload them. Pinned messages and open threads may /// stay in the cache. always fires before /// for the same message. + /// Nothing is ever removed while is true. /// Configs.MessageCacheWindow MessageCacheWindow { get; } @@ -345,35 +359,41 @@ public interface IStreamChannel : IStreamStatefulModel /// /// Set a cache limit for this channel only. Pass null for unlimited on this channel. - /// Trims immediately, against when - /// is true. + /// Trims immediately unless is true. /// void OverrideMessageCacheWindow(Configs.MessageCacheWindow window); /// /// Remove the per-channel limit and use again. - /// Trims immediately, against the paused limit when trimming is paused. + /// Trims immediately unless is true. /// void ClearMessageCacheWindowOverride(); /// - /// Whether the wider, paused cache limit is in effect. pauses - /// automatically so paged-in history is not removed while the user reads it. Trimming is not disabled - /// while paused - the limit becomes . + /// Whether cache trimming is currently suspended for this channel. + /// sets this automatically, because a trim removes the oldest + /// messages - exactly the history it just paged in. + /// While true, no message is ever removed. Growth is bounded instead: + /// stops loading once + /// is true. Incoming live messages are still + /// appended, so a channel that is never resumed keeps growing - the SDK logs a warning once when it + /// crosses that limit. /// bool IsMessageCacheTrimmingPaused { get; } /// - /// Widen the cache limit to - /// (e.g. while the user scrolls through loaded history). - /// does this for you. + /// Suspend cache trimming so no message is removed while the user reads history (e.g. while they + /// scroll back). does this for you. Always pair it with + /// ; see for what + /// bounds memory in the meantime. /// void PauseMessageCacheTrimming(); /// - /// Restore the limit and remove excess messages now. - /// Optional - the cache stays bounded either way. Call it when the user returns to the newest messages - /// to release the memory held by paged-in history sooner. + /// Resume trimming against and remove excess + /// messages now. Call this when the user returns to the newest messages. This is the only thing that + /// releases history paged in by , and the only thing that lets it + /// load more history again once is true. /// void ResumeMessageCacheTrimming(); diff --git a/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs b/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs index adafab13..a0f97683 100644 --- a/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs +++ b/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs @@ -249,6 +249,15 @@ public async Task SendNewMessageAsync(StreamSendMessageRequest s public async Task LoadOlderMessagesAsync() { + // Refusing to page in more history is the only way to bound a paused channel without + // deleting data: a trim removes the oldest messages, which is precisely the history that + // was paged in for the user to read. + if (HasReachedMaxHistoryMessages) + { + WarnAboutMaxHistoryMessagesOnce(MessageCacheWindow); + return; + } + var wasTrimmingPaused = _isMessageCacheTrimmingPaused; // Pause before the request so live messages cannot trim the page being loaded. @@ -291,10 +300,20 @@ public MessageCacheWindow MessageCacheWindow public bool HasMessageCacheWindowOverride => _hasMessageCacheWindowOverride; + public bool HasReachedMaxHistoryMessages + { + get + { + var window = MessageCacheWindow; + return window != null && _messages.Count >= window.MaxHistoryMessages; + } + } + public void OverrideMessageCacheWindow(MessageCacheWindow window) { _messageCacheWindowOverride = window; _hasMessageCacheWindowOverride = true; + _hasWarnedAboutMaxHistoryMessages = false; TrimMessageCacheIfNeeded(); } @@ -302,6 +321,7 @@ public void ClearMessageCacheWindowOverride() { _messageCacheWindowOverride = null; _hasMessageCacheWindowOverride = false; + _hasWarnedAboutMaxHistoryMessages = false; TrimMessageCacheIfNeeded(); } @@ -312,6 +332,7 @@ public void ClearMessageCacheWindowOverride() public void ResumeMessageCacheTrimming() { _isMessageCacheTrimmingPaused = false; + _hasWarnedAboutMaxHistoryMessages = false; TrimMessageCacheIfNeeded(); } @@ -958,6 +979,7 @@ protected override string InternalUniqueId private MessageCacheWindow _messageCacheWindowOverride; private bool _hasMessageCacheWindowOverride; private bool _isMessageCacheTrimmingPaused; + private bool _hasWarnedAboutMaxHistoryMessages; private bool _muted; private bool _hidden; @@ -1064,12 +1086,22 @@ private void TrimMessageCacheIfNeeded() return; } - // Pausing widens the window instead of disabling it, so a channel stays bounded even if - // ResumeMessageCacheTrimming is never called after LoadOlderMessagesAsync. - var effectiveMaxMessages - = _isMessageCacheTrimmingPaused ? window.AbsoluteMaxMessages : window.MaxMessages; + // Trimming always removes the oldest messages, which while paused is exactly the history + // LoadOlderMessagesAsync paged in for the user to read. Removing it would also move the IdLt + // anchor forward, so the next page load would re-fetch what was just removed. Paused therefore + // means "remove nothing"; growth is bounded by refusing to page in more history past + // MaxHistoryMessages instead. + if (_isMessageCacheTrimmingPaused) + { + if (_messages.Count >= window.MaxHistoryMessages) + { + WarnAboutMaxHistoryMessagesOnce(window); + } - if (_messages.Count <= effectiveMaxMessages) + return; + } + + if (_messages.Count <= window.MaxMessages) { return; } @@ -1079,7 +1111,7 @@ var effectiveMaxMessages SortMessagesByCreatedAt(); } - var targetCount = effectiveMaxMessages - window.DiscardBatchSize; + var targetCount = window.MaxMessages - window.DiscardBatchSize; var removeCount = _messages.Count - targetCount; var handler = MessagesRemovedFromCache; @@ -1114,6 +1146,27 @@ var effectiveMaxMessages handler?.Invoke(this, removed); } + // Live messages are still appended while paused, so a channel that is never resumed keeps + // growing. Only the app knows whether the user is still reading history, so the SDK reports it + // once per pause rather than guessing and discarding the user's scroll position. + private void WarnAboutMaxHistoryMessagesOnce(MessageCacheWindow window) + { + if (_hasWarnedAboutMaxHistoryMessages) + { + return; + } + + _hasWarnedAboutMaxHistoryMessages = true; + + Logs.Warning( + $"Channel `{Cid}` holds {_messages.Count} messages, reaching " + + $"{nameof(window.MaxHistoryMessages)} ({window.MaxHistoryMessages}). " + + $"{nameof(LoadOlderMessagesAsync)} will not load more history until " + + $"{nameof(ResumeMessageCacheTrimming)}() is called. Messages are never removed while cache " + + "trimming is paused, so incoming messages keep growing this channel until then - resume once " + + "the user is back at the newest messages."); + } + // Keep in cache when pinned or referenced by a tracked thread. private bool IsRetainedByOtherState(StreamMessage message, HashSet pinnedMessageIds) { diff --git a/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs b/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs index d072ca77..2d5fd842 100644 --- a/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs +++ b/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs @@ -17,8 +17,8 @@ internal class MessageCacheWindowTests : BaseStateIntegrationTests { private static readonly MessageCacheWindow SmallWindow = new MessageCacheWindow(6, 3); - // Absolute limit low enough to be reached by sending a handful of messages while paused. - private static readonly MessageCacheWindow SmallWindowWithLowAbsoluteMax = new MessageCacheWindow(4, 2, 8); + // History limit low enough to be reached by sending a handful of messages while paused. + private static readonly MessageCacheWindow SmallWindowWithLowHistoryLimit = new MessageCacheWindow(4, 2, 8); [UnityTest] public IEnumerator When_no_message_cache_window_configured_expect_no_removal() @@ -264,43 +264,96 @@ private async Task When_trimming_paused_expect_no_removal_on_new_messages_Async( Assert.AreEqual(0, removedCount); } + /// + /// Trimming removes the OLDEST messages, which while paused is exactly the history the user + /// scrolled back to. So pausing must remove nothing at all, even past MaxHistoryMessages - + /// growth is bounded by refusing to page in more history, not by deleting what is on screen. + /// [UnityTest] - public IEnumerator When_trimming_paused_expect_trim_once_absolute_max_exceeded() - => ConnectAndExecute(When_trimming_paused_expect_trim_once_absolute_max_exceeded_Async); + public IEnumerator When_trimming_paused_expect_no_removal_even_past_max_history_messages() + => ConnectAndExecute(When_trimming_paused_expect_no_removal_even_past_max_history_messages_Async); - private async Task When_trimming_paused_expect_trim_once_absolute_max_exceeded_Async() + private async Task When_trimming_paused_expect_no_removal_even_past_max_history_messages_Async() { var channel = await CreateUniqueTempChannelAsync(); - channel.OverrideMessageCacheWindow(SmallWindowWithLowAbsoluteMax); + channel.OverrideMessageCacheWindow(SmallWindowWithLowHistoryLimit); channel.PauseMessageCacheTrimming(); var removedCount = 0; channel.MessagesRemovedFromCache += (_, __) => removedCount++; - // Above MaxMessages(4) but still within AbsoluteMaxMessages(8), so pausing holds the messages. - await SendMessagesAsync(channel, 8); + var sent = await SendMessagesAsync(channel, 10); - Assert.AreEqual(8, channel.Messages.Count); + Assert.Greater(10, SmallWindowWithLowHistoryLimit.MaxHistoryMessages, + "the test must actually push the channel past MaxHistoryMessages"); + Assert.IsTrue(channel.IsMessageCacheTrimmingPaused); + Assert.AreEqual(10, channel.Messages.Count); Assert.AreEqual(0, removedCount); + CollectionAssert.AreEqual(sent.Select(m => m.Id).ToList(), + channel.Messages.Select(m => m.Id).ToList()); + } + + /// + /// The history limit stops history from being paged in rather than deleting what is already there, + /// so the oldest message stays put and the pagination anchor never moves backwards. + /// + [UnityTest] + public IEnumerator When_max_history_messages_reached_expect_LoadOlderMessagesAsync_loads_nothing() + => ConnectAndExecute(When_max_history_messages_reached_expect_LoadOlderMessagesAsync_loads_nothing_Async); + + private async Task When_max_history_messages_reached_expect_LoadOlderMessagesAsync_loads_nothing_Async() + { + var window = SmallWindowWithLowHistoryLimit; + var channel = await CreateUniqueTempChannelAsync(); + await SendMessagesAsync(channel, 12); + + channel.OverrideMessageCacheWindow(window); + Assert.AreEqual(window.MaxMessages - window.DiscardBatchSize, channel.Messages.Count); + Assert.IsFalse(channel.HasReachedMaxHistoryMessages); + + var removedCount = 0; + channel.MessagesRemovedFromCache += (_, __) => removedCount++; - // Crossing AbsoluteMaxMessages trims even though trimming is still paused, down to - // AbsoluteMaxMessages - DiscardBatchSize. This is what keeps a never-resumed channel bounded. - var sent = await SendMessagesAsync(channel, 1); + // Page back until the paused limit is reached. The channel only holds 12 messages, so this + // terminates regardless of the server's page size. + for (var i = 0; i < 5 && channel.Messages.Count < window.MaxHistoryMessages; i++) + { + await channel.LoadOlderMessagesAsync(); + } Assert.IsTrue(channel.IsMessageCacheTrimmingPaused); - Assert.AreEqual(6, channel.Messages.Count); + Assert.GreaterOrEqual(channel.Messages.Count, window.MaxHistoryMessages); + Assert.IsTrue(channel.HasReachedMaxHistoryMessages, + "the app must be able to see that loading more history is pointless"); + Assert.AreEqual(0, removedCount, "nothing may be removed while trimming is paused"); + + var oldestBefore = channel.Messages.First().Id; + var countBefore = channel.Messages.Count; + + await channel.LoadOlderMessagesAsync(); + + Assert.AreEqual(countBefore, channel.Messages.Count, + "LoadOlderMessagesAsync must not load more history once MaxHistoryMessages is reached"); + Assert.AreEqual(oldestBefore, channel.Messages.First().Id); + Assert.AreEqual(0, removedCount); + + // Resuming is what releases the paged-in history and re-enables loading. + channel.ResumeMessageCacheTrimming(); + + Assert.IsFalse(channel.IsMessageCacheTrimmingPaused); + Assert.IsFalse(channel.HasReachedMaxHistoryMessages); + Assert.AreEqual(window.MaxMessages - window.DiscardBatchSize, channel.Messages.Count); Assert.AreEqual(1, removedCount); - Assert.AreEqual(sent.Single().Id, channel.Messages.Last().Id); } [UnityTest] - public IEnumerator When_trimming_paused_and_bounded_expect_resume_restores_max_messages() - => ConnectAndExecute(When_trimming_paused_and_bounded_expect_resume_restores_max_messages_Async); + public IEnumerator When_trimming_paused_expect_resume_restores_max_messages() + => ConnectAndExecute(When_trimming_paused_expect_resume_restores_max_messages_Async); - private async Task When_trimming_paused_and_bounded_expect_resume_restores_max_messages_Async() + private async Task When_trimming_paused_expect_resume_restores_max_messages_Async() { var channel = await CreateUniqueTempChannelAsync(); - channel.OverrideMessageCacheWindow(SmallWindowWithLowAbsoluteMax); + channel.OverrideMessageCacheWindow(SmallWindowWithLowHistoryLimit); channel.PauseMessageCacheTrimming(); await SendMessagesAsync(channel, 8); @@ -473,10 +526,10 @@ public void MessageCacheWindow_rejects_invalid_arguments() } [Test] - public void When_absolute_max_not_specified_expect_default_of_four_times_max_messages() + public void When_max_history_messages_not_specified_expect_default_of_four_times_max_messages() { - Assert.AreEqual(40, new MessageCacheWindow(10, 5).AbsoluteMaxMessages); - Assert.AreEqual(int.MaxValue, new MessageCacheWindow(int.MaxValue, 1).AbsoluteMaxMessages); + Assert.AreEqual(40, new MessageCacheWindow(10, 5).MaxHistoryMessages); + Assert.AreEqual(int.MaxValue, new MessageCacheWindow(int.MaxValue, 1).MaxHistoryMessages); } [Test] @@ -484,7 +537,7 @@ public void MessageCacheWindow_Recommended_is_500_100_2000() { Assert.AreEqual(500, MessageCacheWindow.Recommended.MaxMessages); Assert.AreEqual(100, MessageCacheWindow.Recommended.DiscardBatchSize); - Assert.AreEqual(2000, MessageCacheWindow.Recommended.AbsoluteMaxMessages); + Assert.AreEqual(2000, MessageCacheWindow.Recommended.MaxHistoryMessages); } private static async Task> SendMessagesAsync(IStreamChannel channel, int count) From c4729fc17d2832599ebbb2c3c9c90f7f286c4e83 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Sierpi=C5=84ski?= <33436839+sierpinskid@users.noreply.github.com> Date: Thu, 13 Aug 2026 15:43:22 +0200 Subject: [PATCH 4/4] Rename prop --- .../Core/StatefulModels/IStreamChannel.cs | 19 +++++++++++-------- .../Core/StatefulModels/StreamChannel.cs | 6 +++--- .../StatefulClient/MessageCacheWindowTests.cs | 6 +++--- 3 files changed, 17 insertions(+), 14 deletions(-) diff --git a/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs b/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs index fee011a0..a837ffbf 100644 --- a/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs +++ b/Assets/Plugins/StreamChat/Core/StatefulModels/IStreamChannel.cs @@ -322,21 +322,24 @@ public interface IStreamChannel : IStreamStatefulModel /// Note that loading older messages does NOT trigger the event. /// Calling this pauses cache trimming (see ) so the /// loaded page is not removed while the user reads it. If - /// is true this returns without loading anything - - /// stop showing your "load more" affordance and prompt the user to jump back to the newest messages, - /// which is where belongs. + /// is true this returns without loading + /// anything - stop showing your "load more" affordance and prompt the user to jump back to the newest + /// messages, which is where belongs. /// Task LoadOlderMessagesAsync(); /// - /// true when has reached + /// true when this client's local cache for this channel has reached /// and /// will therefore load nothing. Always false when /// is null. + /// This says nothing about the server: the channel's history is intact and untouched, this + /// device has simply cached as much of it as it is allowed to. + /// clears it. /// Check this before offering "load more" so the user is not left waiting on a call that - /// cannot return anything. clears it. + /// cannot return anything. /// - bool HasReachedMaxHistoryMessages { get; } + bool IsMessageCacheHistoryLimitReached { get; } /// /// Active cache limit for this channel. null = unlimited (default). @@ -375,7 +378,7 @@ public interface IStreamChannel : IStreamStatefulModel /// messages - exactly the history it just paged in. /// While true, no message is ever removed. Growth is bounded instead: /// stops loading once - /// is true. Incoming live messages are still + /// is true. Incoming live messages are still /// appended, so a channel that is never resumed keeps growing - the SDK logs a warning once when it /// crosses that limit. /// @@ -393,7 +396,7 @@ public interface IStreamChannel : IStreamStatefulModel /// Resume trimming against and remove excess /// messages now. Call this when the user returns to the newest messages. This is the only thing that /// releases history paged in by , and the only thing that lets it - /// load more history again once is true. + /// load more history again once is true. /// void ResumeMessageCacheTrimming(); diff --git a/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs b/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs index a0f97683..a1bba234 100644 --- a/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs +++ b/Assets/Plugins/StreamChat/Core/StatefulModels/StreamChannel.cs @@ -252,7 +252,7 @@ public async Task LoadOlderMessagesAsync() // Refusing to page in more history is the only way to bound a paused channel without // deleting data: a trim removes the oldest messages, which is precisely the history that // was paged in for the user to read. - if (HasReachedMaxHistoryMessages) + if (IsMessageCacheHistoryLimitReached) { WarnAboutMaxHistoryMessagesOnce(MessageCacheWindow); return; @@ -300,7 +300,7 @@ public MessageCacheWindow MessageCacheWindow public bool HasMessageCacheWindowOverride => _hasMessageCacheWindowOverride; - public bool HasReachedMaxHistoryMessages + public bool IsMessageCacheHistoryLimitReached { get { @@ -1159,7 +1159,7 @@ private void WarnAboutMaxHistoryMessagesOnce(MessageCacheWindow window) _hasWarnedAboutMaxHistoryMessages = true; Logs.Warning( - $"Channel `{Cid}` holds {_messages.Count} messages, reaching " + $"Channel `{Cid}` holds {_messages.Count} messages in the local cache, reaching " + $"{nameof(window.MaxHistoryMessages)} ({window.MaxHistoryMessages}). " + $"{nameof(LoadOlderMessagesAsync)} will not load more history until " + $"{nameof(ResumeMessageCacheTrimming)}() is called. Messages are never removed while cache " diff --git a/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs b/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs index 2d5fd842..24f91a4f 100644 --- a/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs +++ b/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs @@ -309,7 +309,7 @@ private async Task When_max_history_messages_reached_expect_LoadOlderMessagesAsy channel.OverrideMessageCacheWindow(window); Assert.AreEqual(window.MaxMessages - window.DiscardBatchSize, channel.Messages.Count); - Assert.IsFalse(channel.HasReachedMaxHistoryMessages); + Assert.IsFalse(channel.IsMessageCacheHistoryLimitReached); var removedCount = 0; channel.MessagesRemovedFromCache += (_, __) => removedCount++; @@ -323,7 +323,7 @@ private async Task When_max_history_messages_reached_expect_LoadOlderMessagesAsy Assert.IsTrue(channel.IsMessageCacheTrimmingPaused); Assert.GreaterOrEqual(channel.Messages.Count, window.MaxHistoryMessages); - Assert.IsTrue(channel.HasReachedMaxHistoryMessages, + Assert.IsTrue(channel.IsMessageCacheHistoryLimitReached, "the app must be able to see that loading more history is pointless"); Assert.AreEqual(0, removedCount, "nothing may be removed while trimming is paused"); @@ -341,7 +341,7 @@ private async Task When_max_history_messages_reached_expect_LoadOlderMessagesAsy channel.ResumeMessageCacheTrimming(); Assert.IsFalse(channel.IsMessageCacheTrimmingPaused); - Assert.IsFalse(channel.HasReachedMaxHistoryMessages); + Assert.IsFalse(channel.IsMessageCacheHistoryLimitReached); Assert.AreEqual(window.MaxMessages - window.DiscardBatchSize, channel.Messages.Count); Assert.AreEqual(1, removedCount); }