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..7cab5a93 --- /dev/null +++ b/Assets/Plugins/StreamChat/Core/Configs/MessageCacheWindow.cs @@ -0,0 +1,101 @@ +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 + /// . Nothing is ever removed while + /// is true - growth is + /// bounded by instead. + /// + public sealed class MessageCacheWindow + { + /// 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. + public int MaxMessages { get; } + + /// How many messages to remove per trim. Must be less than . + public int DiscardBatchSize { get; } + + /// + /// 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 MaxHistoryMessages { get; } + + public MessageCacheWindow(int maxMessages, int discardBatchSize) + : this(maxMessages, discardBatchSize, GetDefaultMaxHistoryMessages(maxMessages)) + { + } + + public MessageCacheWindow(int maxMessages, int discardBatchSize, int maxHistoryMessages) + { + 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."); + } + + if (maxHistoryMessages < maxMessages) + { + throw new ArgumentOutOfRangeException(nameof(maxHistoryMessages), maxHistoryMessages, + $"{nameof(maxHistoryMessages)} must be greater than or equal to {nameof(maxMessages)} " + + $"({maxMessages})."); + } + + MaxMessages = maxMessages; + DiscardBatchSize = discardBatchSize; + MaxHistoryMessages = maxHistoryMessages; + } + + public override string ToString() + => $"MessageCacheWindow - MaxMessages: {MaxMessages}, DiscardBatchSize: {DiscardBatchSize}, " + + $"MaxHistoryMessages: {MaxHistoryMessages}"; + + 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 GetDefaultMaxHistoryMessages(int maxMessages) + { + if (maxMessages <= 0) + { + return maxMessages; + } + + return maxMessages > int.MaxValue / DefaultMaxHistoryMessagesMultiplier + ? int.MaxValue + : maxMessages * DefaultMaxHistoryMessagesMultiplier; + } + } +} 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/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 new file mode 100644 index 00000000..c8888301 --- /dev/null +++ b/Assets/Plugins/StreamChat/Core/Helpers/ListPool.cs @@ -0,0 +1,41 @@ +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)); + } + + // 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 (!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/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/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 d9b82cf3..a837ffbf 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 /// @@ -312,10 +319,87 @@ 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 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. + /// + bool IsMessageCacheHistoryLimitReached { get; } + + /// + /// Active cache limit for this channel. null = unlimited (default). + /// 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; } + + /// + /// 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 is true. + /// + void ClearMessageCacheWindowOverride(); + + /// + /// 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; } + + /// + /// 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(); + + /// + /// 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(); + /// /// 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..a1bba234 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,26 +249,91 @@ public async Task SendNewMessageAsync(StreamSendMessageRequest s public async Task LoadOlderMessagesAsync() { - var oldestMessage = _messages.OrderBy(_ => _.CreatedAt).FirstOrDefault(); - - var request = new ChannelGetOrCreateRequestInternalDTO + // 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 (IsMessageCacheHistoryLimitReached) { - //StreamTodo: presence could be optional in config - Presence = true, - State = true, - Watch = true, - }; + WarnAboutMaxHistoryMessagesOnce(MessageCacheWindow); + return; + } + + var wasTrimmingPaused = _isMessageCacheTrimmingPaused; - if (oldestMessage != null) + // Pause before the request so live messages cannot trim the page being loaded. + PauseMessageCacheTrimming(); + + try { - request.Messages = new MessagePaginationParamsRequestInternalDTO + var oldestMessage = _messages.OrderBy(_ => _.CreatedAt).FirstOrDefault(); + + var request = new ChannelGetOrCreateRequestInternalDTO { - IdLt = oldestMessage.Id, + //StreamTodo: presence could be optional in config + Presence = true, + State = true, + Watch = true, }; + + 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; } + } - var response = await LowLevelClient.InternalChannelApi.GetOrCreateChannelAsync(Type, Id, request); - Cache.TryCreateOrUpdate(response); + public MessageCacheWindow MessageCacheWindow + => _hasMessageCacheWindowOverride ? _messageCacheWindowOverride : LowLevelClient.Config.DefaultMessageCacheWindow; + + public bool HasMessageCacheWindowOverride => _hasMessageCacheWindowOverride; + + public bool IsMessageCacheHistoryLimitReached + { + get + { + var window = MessageCacheWindow; + return window != null && _messages.Count >= window.MaxHistoryMessages; + } + } + + public void OverrideMessageCacheWindow(MessageCacheWindow window) + { + _messageCacheWindowOverride = window; + _hasMessageCacheWindowOverride = true; + _hasWarnedAboutMaxHistoryMessages = false; + TrimMessageCacheIfNeeded(); + } + + public void ClearMessageCacheWindowOverride() + { + _messageCacheWindowOverride = null; + _hasMessageCacheWindowOverride = false; + _hasWarnedAboutMaxHistoryMessages = false; + TrimMessageCacheIfNeeded(); + } + + public bool IsMessageCacheTrimmingPaused => _isMessageCacheTrimmingPaused; + + public void PauseMessageCacheTrimming() => _isMessageCacheTrimmingPaused = true; + + public void ResumeMessageCacheTrimming() + { + _isMessageCacheTrimmingPaused = false; + _hasWarnedAboutMaxHistoryMessages = false; + TrimMessageCacheIfNeeded(); } public async Task UpdateOverwriteAsync(StreamUpdateOverwriteChannelRequest updateOverwriteRequest) @@ -906,6 +976,11 @@ 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 _hasWarnedAboutMaxHistoryMessages; + private bool _muted; private bool _hidden; @@ -963,6 +1038,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 +1078,129 @@ private void InternalTruncateMessages(DateTimeOffset? deleteBeforeCreatedAt = nu Truncated?.Invoke(this); } + private void TrimMessageCacheIfNeeded() + { + var window = MessageCacheWindow; + if (window == null) + { + return; + } + + // 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); + } + + return; + } + + if (_messages.Count <= window.MaxMessages) + { + return; + } + + if (!AreMessagesSortedByCreatedAt()) + { + SortMessagesByCreatedAt(); + } + + var targetCount = window.MaxMessages - window.DiscardBatchSize; + var removeCount = _messages.Count - targetCount; + + var handler = MessagesRemovedFromCache; + + // Not pooled - this is handed to subscribers, so the SDK does not control its lifetime. + var removed = handler == null ? null : new List(removeCount); + + using (new ListPoolScope(out var tempUntrackCandidates)) + using (new HashSetPoolScope(out var tempPinnedMessageIds)) + { + for (var i = 0; i < _pinnedMessages.Count; i++) + { + tempPinnedMessageIds.Add(_pinnedMessages[i].Id); + } + + for (var i = 0; i < removeCount; i++) + { + var message = _messages[i]; + removed?.Add(message); + + if (!IsRetainedByOtherState(message, tempPinnedMessageIds)) + { + tempUntrackCandidates.Add(message); + } + } + + _messages.RemoveRange(0, removeCount); + Cache.Messages.RemoveMany(tempUntrackCandidates); + } + + // Raised after the pooled buffers are returned so subscribers can trim or send safely. + 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 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 " + + "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) + { + if (pinnedMessageIds.Contains(message.Id)) + { + 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..24f91a4f --- /dev/null +++ b/Assets/Plugins/StreamChat/Tests/StatefulClient/MessageCacheWindowTests.cs @@ -0,0 +1,555 @@ +#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); + + // 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() + => 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()}"); + + // 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.Take(4).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); + } + + /// + /// 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_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_no_removal_even_past_max_history_messages_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindowWithLowHistoryLimit); + channel.PauseMessageCacheTrimming(); + + var removedCount = 0; + channel.MessagesRemovedFromCache += (_, __) => removedCount++; + + var sent = await SendMessagesAsync(channel, 10); + + 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.IsMessageCacheHistoryLimitReached); + + var removedCount = 0; + channel.MessagesRemovedFromCache += (_, __) => removedCount++; + + // 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.GreaterOrEqual(channel.Messages.Count, window.MaxHistoryMessages); + 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"); + + 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.IsMessageCacheHistoryLimitReached); + Assert.AreEqual(window.MaxMessages - window.DiscardBatchSize, channel.Messages.Count); + Assert.AreEqual(1, removedCount); + } + + [UnityTest] + 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_expect_resume_restores_max_messages_Async() + { + var channel = await CreateUniqueTempChannelAsync(); + channel.OverrideMessageCacheWindow(SmallWindowWithLowHistoryLimit); + 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); + + 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.InternalLowLevelClient.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.InternalLowLevelClient.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 = new HashSet(sent.Take(4).Select(m => m.Id)); + 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)); + Assert.Throws(() => new MessageCacheWindow(10, 5, 9)); + Assert.DoesNotThrow(() => new MessageCacheWindow(10, 5, 10)); + } + + [Test] + public void When_max_history_messages_not_specified_expect_default_of_four_times_max_messages() + { + Assert.AreEqual(40, new MessageCacheWindow(10, 5).MaxHistoryMessages); + Assert.AreEqual(int.MaxValue, new MessageCacheWindow(int.MaxValue, 1).MaxHistoryMessages); + } + + [Test] + 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.MaxHistoryMessages); + } + + 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