Kotlin Flow and Websockets with clean architecture on Android

Viewed 1406

Recently, our team tried to implement websockets. We easily thought to use Rx when listening on events but I wondered how to do without it. So, we tried the famous Kotlin Flow but I don't know if our implementation is correct.

Our app's architecture is splitted in four layers:

  • Service - emits and receives events from socket,
  • Repository - filters, maps, transforms, etc,
  • ViewModel - populates the LiveDatas
  • Activity - observes the changes and updates UI.

Therefore, we listen the events received into the Service as follows:

fun listenMessages(): Flow<List<Message>> = channelFlow {
    socket.on("NewMessage") { args ->
        val message = gson.fromJson(args[0].toString(), ...)
        trySend(message)
    }
    awaitClose()
}

We use channelFlow's coroutine to send to the consumer when an event is received with trySend and we keep this Flow alive by using awaitClose.

The Repository does some logic after catching the Flow and sends it back to the ViewModel:

fun getMessages(): Flow<List<Message>> {
    return service.listenMessages()
            .filter { ... }
            .map { ... }
}

Then, the ViewModel launches the coroutines and updates the LiveData when collecting the Flow:

fun getMessages() {
    viewModelScope.launch(context = Dispatcher.IO) {
        repository.getMessages()
                .collect {
                    messagesLiveData.postValue(it)
                }
    }
}

This works well however this raises some questions:

  • Is this the correct implementation?
  • When we need to constantly listening, does channelFlow is the right choice?
  • In this case, should we use classic Channels instead of Flow (hot vs cold)?

Thanks in advance for your advices.

1 Answers

The concept of channels is meant as a primitive to mainly communicate between coroutines 1. As far as I understand it based on the channelFlow docs, it uses a Channel under the hood and translates it to a Flow. By using a Channel under the hood, there is something important to realize:

every value that is sent to the channel is received once. You cannot use channels to distribute events or state updates in a way that allows multiple subscribers to independently receive and react upon them. 1

Depending on your architecture this may not matter, but this example may illustrate something important:

fun ws(): Flow<List<Message>> = channelFlow {
    println("channelFlow block called")
    trySend(listOf(Message(0), Message(1)))
    delay(2_000)
    trySend(listOf(Message(2)))
}

fun main() = runBlocking {
    val source = ws()

    launch {
        source.collect {
            println("First collect got $it")
        }
    }
    launch {
        source.collect {
            println("Second collect got $it")
        }
    }
}

Which generates the output:

channelFlow block called
channelFlow block called
First collect got [Message(id=0), Message(id=1)]
Second collect got [Message(id=0), Message(id=1)]
First collect got [Message(id=2)]
Second collect got [Message(id=2)]

Since a Channel cannot be shared, each time .collect is called on source, it triggers the channelFlow body to be called! What may be more suprising is that the same Flow is referenced via source, there is not a new call to ws() going on.

Looking more closely at the channelFlow docs you can see that

The resulting flow is cold, which means that block is called every time a terminal operator is applied to the resulting flow.2

Since .collect() is a terminal operator, it triggers a new call to the channelFlow body.

Now, does this matter for your use case? I'm not sure; you will have to figure out if you are subscribing to the flow produced by listenMessages in multiple locations, or if the body of that channelFlow is expensive to perform more than once.

As a more general recommendation, I would suggest to be more explicit about the behavior of your service. Should listenMessages emit to everyone that is subscribed? Should it only emit to the first one that is available (see also fan-out behavior)? If you prefer the former, a SharedFlow would be recommended, if the later I would be explicit about exposing a Channel directly so you don't confuse downstream consumers about the issue I showed above.

Related