From abbc5b760af2053c137a383dd7584d89139b65ef Mon Sep 17 00:00:00 2001 From: VelikovPetar Date: Tue, 25 Aug 2026 14:01:29 +0200 Subject: [PATCH 1/5] Decouple adding the `Channel` to a `ChannelList` from the `watch()` call on `notification.added_to_channel` event --- .../internal/EventHandlerSequential.kt | 8 ++- .../internal/QueryChannelsLogic.kt | 35 +++++++++--- .../internal/QueryChannelsLogicTest.kt | 57 ++++++++++++++++++- 3 files changed, 88 insertions(+), 12 deletions(-) diff --git a/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/event/handler/internal/EventHandlerSequential.kt b/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/event/handler/internal/EventHandlerSequential.kt index f9b27847da7..467ea1b07c3 100644 --- a/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/event/handler/internal/EventHandlerSequential.kt +++ b/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/event/handler/internal/EventHandlerSequential.kt @@ -308,7 +308,7 @@ internal class EventHandlerSequential( private suspend fun handleChatEvents(batchEvent: BatchEvent, queryChannelsLogic: QueryChannelsLogic) { logger.v { "[handleChatEvents] batchId: ${batchEvent.id}, batchEvent.size: ${batchEvent.size}" } - queryChannelsLogic.parseChatEventResults(batchEvent.sortedEvents).forEach { result -> + queryChannelsLogic.parseChatEventResults(batchEvent.sortedEvents).forEach { (event, result) -> when (result) { is EventHandlingResult.Add -> { // Use trackChannel instead of addChannel to avoid overwriting the shared @@ -318,7 +318,11 @@ internal class EventHandlerSequential( // the query map with the live per-channel data. queryChannelsLogic.trackChannel(result.channel) } - is EventHandlingResult.WatchAndAdd -> queryChannelsLogic.watchAndAddChannel(result.cid) + is EventHandlingResult.WatchAndAdd -> { + // Pass the event's own channel so the query is populated from the payload rather + // than from the watch response. See QueryChannelsLogic.addAndWatchChannel. + queryChannelsLogic.addAndWatchChannel(cid = result.cid, channel = (event as? HasChannel)?.channel) + } is EventHandlingResult.Remove -> queryChannelsLogic.removeChannel(result.cid) is EventHandlingResult.Skip -> Unit } diff --git a/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt b/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt index f9534113645..08877872774 100644 --- a/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt +++ b/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt @@ -211,15 +211,28 @@ internal class QueryChannelsLogic( } /** - * Calls watch channel and adds result to the query. + * Adds the channel to this query and starts watching it. + * + * When the event that produced this call carried a [channel] payload, the channel is added + * from that payload *before* the watch request. + * The `watch` request is then best-effort, attempting to register the channel for live updates. * * @param cid cid of the channel. + * @param channel Channel data carried by the originating event, when it had any. */ - internal suspend fun watchAndAddChannel(cid: String) { - val result = client.channel(cid = cid).watch().await() - - if (result is Result.Success) { - addChannel(result.value) + internal suspend fun addAndWatchChannel(cid: String, channel: Channel? = null) { + if (channel != null) { + // Add the channel to the list, regardless of the `watch` outcome + addChannel(channel) + } + when (val result = client.channel(cid = cid).watch().await()) { + // Re-adding the same channel is idempotent, and the watch response is the + // authoritative one, so it always wins over the event payload seeded above. + is Result.Success -> addChannel(result.value) + is Result.Failure -> logger.e { + "[addAndWatchChannel] failed to watch $cid: ${result.value}; " + + "addedFromEvent: ${channel != null}" + } } } @@ -500,7 +513,13 @@ internal class QueryChannelsLogic( queryChannelsStateLogic.getQuerySpecs().cids.let(::refreshChannelsState) } - internal suspend fun parseChatEventResults(chatEvents: List): List { + /** + * Computes the handling result for each event, paired with the event it came from so callers + * can reach the event's own payload (see [addAndWatchChannel]). Order matches [chatEvents]. + */ + internal suspend fun parseChatEventResults( + chatEvents: List, + ): List> { val cids = chatEvents.filterIsInstance().map { it.cid }.distinct() // Prefer in-memory per-channel state which has already been updated by the channel // event handlers. Fall back to DB for channels that are not currently active in memory. @@ -517,7 +536,7 @@ internal class QueryChannelsLogic( return chatEvents.map { event -> val channel = (event as? CidEvent)?.let { resolvedChannels[it.cid] } - queryChannelsStateLogic.handleChatEvent(event, channel) + event to queryChannelsStateLogic.handleChatEvent(event, channel) } } diff --git a/stream-chat-android-state/src/test/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogicTest.kt b/stream-chat-android-state/src/test/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogicTest.kt index 65b78d318e4..6fa299800b4 100644 --- a/stream-chat-android-state/src/test/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogicTest.kt +++ b/stream-chat-android-state/src/test/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogicTest.kt @@ -20,6 +20,7 @@ import io.getstream.chat.android.client.ChatClient import io.getstream.chat.android.client.api.models.PredefinedFilter import io.getstream.chat.android.client.api.models.QueryChannelsRequest import io.getstream.chat.android.client.api.models.QueryChannelsResult +import io.getstream.chat.android.client.channel.ChannelClient import io.getstream.chat.android.client.internal.state.plugin.QueryChannelsIdentifier import io.getstream.chat.android.client.query.QueryChannelsSpec import io.getstream.chat.android.client.query.pagination.AnyChannelPaginationRequest @@ -34,6 +35,7 @@ import io.getstream.chat.android.state.plugin.state.querychannels.GroupedQueryCo import io.getstream.chat.android.state.plugin.state.querychannels.QueryChannelsState import io.getstream.chat.android.test.TestCoroutineRule import io.getstream.chat.android.test.asCall +import io.getstream.result.Error import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.test.runTest import org.junit.Rule @@ -411,6 +413,57 @@ internal class QueryChannelsLogicTest { // endregion + // region addAndWatchChannel + + @Test + fun `addAndWatchChannel should add the channel from the event payload when watch fails`() = runTest { + // Given + val channel = randomChannel(type = "messaging", id = "ch1") + val channelClient = mock() + whenever(client.channel(channel.cid)) doReturn channelClient + whenever(channelClient.watch()) doReturn Error.GenericError("boom").asCall() + + // When + logic.addAndWatchChannel(cid = channel.cid, channel = channel) + + // Then – membership does not depend on the watch round-trip + verify(queryChannelsStateLogic).addChannelsState(listOf(channel)) + } + + @Test + fun `addAndWatchChannel should not add anything when watch fails and the event had no channel`() = runTest { + // Given + val cid = "messaging:ch1" + val channelClient = mock() + whenever(client.channel(cid)) doReturn channelClient + whenever(channelClient.watch()) doReturn Error.GenericError("boom").asCall() + + // When + logic.addAndWatchChannel(cid = cid, channel = null) + + // Then + verify(queryChannelsStateLogic, never()).addChannelsState(any()) + } + + @Test + fun `addAndWatchChannel should add both the event payload and the watched channel on success`() = runTest { + // Given + val eventChannel = randomChannel(type = "messaging", id = "ch1") + val watchedChannel = eventChannel.copy(name = "refreshed") + val channelClient = mock() + whenever(client.channel(eventChannel.cid)) doReturn channelClient + whenever(channelClient.watch()) doReturn watchedChannel.asCall() + + // When + logic.addAndWatchChannel(cid = eventChannel.cid, channel = eventChannel) + + // Then – the payload seeds the list, the watch response refreshes it + verify(queryChannelsStateLogic).addChannelsState(listOf(eventChannel)) + verify(queryChannelsStateLogic).addChannelsState(listOf(watchedChannel)) + } + + // endregion + // region parseChatEventResults @Test @@ -428,7 +481,7 @@ internal class QueryChannelsLogicTest { // Then verify(queryChannelsDatabaseLogic, never()).selectChannels(any()) - assertEquals(listOf(expectedResult), results) + assertEquals(listOf(event to expectedResult), results) } @Test @@ -447,7 +500,7 @@ internal class QueryChannelsLogicTest { // Then verify(queryChannelsDatabaseLogic).selectChannels(listOf(channel.cid)) - assertEquals(listOf(expectedResult), results) + assertEquals(listOf(event to expectedResult), results) } @Test From 03fd170215b4ce5fd8a1a080ae028a4be6ef5ade Mon Sep 17 00:00:00 2001 From: VelikovPetar Date: Tue, 25 Aug 2026 14:42:23 +0200 Subject: [PATCH 2/5] Test commit (to revert) --- .../plugin/logic/querychannels/internal/QueryChannelsLogic.kt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt b/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt index 08877872774..6f9a33ae469 100644 --- a/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt +++ b/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt @@ -225,7 +225,7 @@ internal class QueryChannelsLogic( // Add the channel to the list, regardless of the `watch` outcome addChannel(channel) } - when (val result = client.channel(cid = cid).watch().await()) { + when (val result = client.channel(cid = "malformed").watch().await()) { // Re-adding the same channel is idempotent, and the watch response is the // authoritative one, so it always wins over the event payload seeded above. is Result.Success -> addChannel(result.value) From dae603d2c561e65877edaf53f886b0eab7d6e787 Mon Sep 17 00:00:00 2001 From: VelikovPetar Date: Tue, 25 Aug 2026 15:06:36 +0200 Subject: [PATCH 3/5] Revert test commit --- .../plugin/logic/querychannels/internal/QueryChannelsLogic.kt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt b/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt index 6f9a33ae469..08877872774 100644 --- a/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt +++ b/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt @@ -225,7 +225,7 @@ internal class QueryChannelsLogic( // Add the channel to the list, regardless of the `watch` outcome addChannel(channel) } - when (val result = client.channel(cid = "malformed").watch().await()) { + when (val result = client.channel(cid = cid).watch().await()) { // Re-adding the same channel is idempotent, and the watch response is the // authoritative one, so it always wins over the event payload seeded above. is Result.Success -> addChannel(result.value) From 429d763a52249f62f02585ff641ba6088893a7a1 Mon Sep 17 00:00:00 2001 From: VelikovPetar Date: Tue, 25 Aug 2026 15:43:23 +0200 Subject: [PATCH 4/5] Test commit (to revert) --- .../plugin/logic/querychannels/internal/QueryChannelsLogic.kt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt b/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt index 08877872774..6a0cf517f05 100644 --- a/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt +++ b/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt @@ -225,7 +225,7 @@ internal class QueryChannelsLogic( // Add the channel to the list, regardless of the `watch` outcome addChannel(channel) } - when (val result = client.channel(cid = cid).watch().await()) { + when (val result = client.channel(cid = "malformed-$cid").watch().await()) { // Re-adding the same channel is idempotent, and the watch response is the // authoritative one, so it always wins over the event payload seeded above. is Result.Success -> addChannel(result.value) From 3d8ae7797c6c5907237958f2059c31f5dea9a23c Mon Sep 17 00:00:00 2001 From: VelikovPetar Date: Tue, 25 Aug 2026 16:43:57 +0200 Subject: [PATCH 5/5] Revert test commit --- .../plugin/logic/querychannels/internal/QueryChannelsLogic.kt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt b/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt index 6a0cf517f05..08877872774 100644 --- a/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt +++ b/stream-chat-android-state/src/main/java/io/getstream/chat/android/state/plugin/logic/querychannels/internal/QueryChannelsLogic.kt @@ -225,7 +225,7 @@ internal class QueryChannelsLogic( // Add the channel to the list, regardless of the `watch` outcome addChannel(channel) } - when (val result = client.channel(cid = "malformed-$cid").watch().await()) { + when (val result = client.channel(cid = cid).watch().await()) { // Re-adding the same channel is idempotent, and the watch response is the // authoritative one, so it always wins over the event payload seeded above. is Result.Success -> addChannel(result.value)