Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion gradle/libs.versions.toml
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ moshi = "1.15.2"
navigationCompose = "2.9.3"
okhttp = "4.12.0"
retrofit = "2.12.0"
streamAndroidCore = "4.0.0"
streamAndroidCore = "5.0.2"
symbolProcessingApi = "2.2.0-2.0.2"
lifecycleProcess = "2.9.1"
lifecycleViewModelCompose = "2.4.0"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,6 @@ import io.getstream.android.core.api.model.value.StreamApiKey
import io.getstream.android.core.api.model.value.StreamHttpClientInfoHeader
import io.getstream.android.core.api.model.value.StreamUserId
import io.getstream.android.core.api.model.value.StreamWsUrl
import io.getstream.android.core.api.processing.StreamBatcher
import io.getstream.android.core.api.processing.StreamRetryProcessor
import io.getstream.android.core.api.processing.StreamSerialProcessingQueue
import io.getstream.android.core.api.processing.StreamSingleFlightProcessor
Expand Down Expand Up @@ -120,13 +119,8 @@ internal fun createStreamCoreClient(
logger = logProvider.taggedLogger("SCHealthMonitor"),
scope = scope,
),
batcher =
StreamBatcher(
scope = scope,
batchSize = 10,
initialDelayMs = 100L,
maxDelayMs = 1_000L,
),
// eventAggregator is left unset: core builds one from the event parser it already
// holds, at the aggregation defaults on socketConfig.
)

return StreamClient(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import io.getstream.android.core.api.model.connection.StreamConnectedUser
import io.getstream.android.core.api.model.connection.StreamConnectionState
import io.getstream.android.core.api.model.exceptions.StreamClientException
import io.getstream.android.core.api.model.value.StreamApiKey
import io.getstream.android.core.api.processing.StreamAggregatedEvent
import io.getstream.android.core.api.socket.listeners.StreamClientListener
import io.getstream.android.core.api.subscribe.StreamSubscriptionManager
import io.getstream.feeds.android.client.api.FeedsClient
Expand Down Expand Up @@ -168,6 +169,18 @@ internal class FeedsClientImpl(
object : StreamClientListener {

override fun onEvent(event: Any) {
when (event) {
// Core aggregates events into one dispatch when the socket spikes. Arrival
// order is preserved, so replaying them one by one matches normal traffic.
is StreamAggregatedEvent<*> -> {
logger.v { "[onEvent] Received ${event.events.size} aggregated events" }
event.events.forEach(::handleEvent)
}
else -> handleEvent(event)
}
}

private fun handleEvent(event: Any?) {
if (event is WSEvent) {
logger.v { "[onEvent] Received event from core: $event" }
_events.tryEmit(event)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@ import io.getstream.android.core.api.model.connection.StreamConnectedUser
import io.getstream.android.core.api.model.connection.StreamConnectionState
import io.getstream.android.core.api.model.exceptions.StreamClientException
import io.getstream.android.core.api.model.value.StreamApiKey
import io.getstream.android.core.api.processing.StreamAggregatedEvent
import io.getstream.android.core.api.socket.listeners.StreamClientListener
import io.getstream.android.core.api.subscribe.StreamSubscriptionManager
import io.getstream.feeds.android.client.api.Moderation
import io.getstream.feeds.android.client.api.file.FeedUploader
Expand Down Expand Up @@ -62,9 +64,11 @@ import io.getstream.feeds.android.network.models.ActivityRequest
import io.getstream.feeds.android.network.models.AddActivityRequest
import io.getstream.feeds.android.network.models.DeleteActivitiesRequest
import io.getstream.feeds.android.network.models.DeleteActivitiesResponse
import io.getstream.feeds.android.network.models.WSEvent
import io.mockk.coEvery
import io.mockk.every
import io.mockk.mockk
import io.mockk.slot
import io.mockk.verify
import java.util.Date
import kotlinx.coroutines.ExperimentalCoroutinesApi
Expand Down Expand Up @@ -436,4 +440,51 @@ internal class FeedsClientImplTest {

assertEquals(query, result.query)
}

@Test
fun `on aggregated event, then dispatch every event it contains`() = runTest {
val listener = captureClientListener()

listener.onEvent(
StreamAggregatedEvent(listOf(wsEvent("activity.added"), wsEvent("activity.deleted")))
)

// One dispatch per contained event, exactly as if they had arrived separately.
verify(exactly = 2) { feedsEventsSubscriptionManager.forEach(any()) }
}

@Test
fun `on aggregated event holding a foreign payload, then still dispatch the WSEvents`() =
runTest {
val listener = captureClientListener()

// A non-WSEvent entry must not stop the rest of the batch.
listener.onEvent(
StreamAggregatedEvent(listOf("not-an-event", wsEvent("activity.added")))
)

verify(exactly = 1) { feedsEventsSubscriptionManager.forEach(any()) }
}

@Test
fun `on single event, then dispatch it once`() = runTest {
val listener = captureClientListener()

listener.onEvent(wsEvent("activity.added"))

verify(exactly = 1) { feedsEventsSubscriptionManager.forEach(any()) }
}

private suspend fun captureClientListener(): StreamClientListener {
val listener = slot<StreamClientListener>()
every { coreClient.subscribe(capture(listener)) } returns Result.success(mockk())
coEvery { coreClient.connect() } returns Result.success(mockk())
feedsClient.connect()
return listener.captured
}

private fun wsEvent(type: String): WSEvent =
object : WSEvent {
override fun getWSEventType(): String = type
}
}
Loading