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