Estoy tratando de mantener un flujo de estado mutable en mi clase, pero cuando le aplico algún método, se convertirá en un Flow<T> inmutable:
class MyClass : Listener<String> { private val source = Source() val flow: Flow<String?> get() = _flow // region listener override fun onUpdate(value: String?) { if (value!= null) { // emit object changes to the flow // not possible, because the builder operators on the flow below convert it to a `Flow` and it doesn't stay as a `MutableSharedFlow` :( _flow.tryEmit(value) } } // end-region @OptIn(ExperimentalCoroutinesApi::class) private val _flow by lazy { MutableStateFlow<String?>(null).onStart { emitAll( flow<String?> { val initialValue = source.getInitialValue() emit(initialValue) }.flowOn(MyDispatchers.background) ) }.onCompletion { error -> // when the flow is cancelled, stop listening to changes if (error is CancellationException) { // is was cancelled source.removeListener(this@MyClass) } }.apply { // listen to changes and send them to the flow source.addListener(this@MyClass) } } } ¿Hay alguna manera de mantener el flujo como MutableStateFlow incluso después de aplicarle los métodos onCompletion/onStart ?
Si aplica transformaciones a un flujo de estado mutable, el flujo resultante pasa a ser de solo lectura porque el flujo original actúa como su origen. Si desea emitir eventos manualmente, debe emitirlos al flujo de origen inicial.
Dicho esto, parece que lo que quiere lograr aquí es bastante simple: unir una API basada en devolución de llamada a una API de Flow . Hay una función incorporada en las rutinas de Kotlin para hacer eso, que se llama callbackFlow .
No estoy seguro de cómo su API de origen maneja la contrapresión, pero se vería así:
@OptIn(ExperimentalCoroutinesApi::class) fun Source.asFlow(): Flow<String?> = callbackFlow { send(getInitialValue()) val listener = object : Listener<String> { override fun onUpdate(value: String?) { if (value != null) { trySend(value) } } } addListener(listener) awaitClose { removeListener(listener) } } O tal vez con runBlocking { send(value) } en lugar de trySend() , dependiendo de cómo Source maneje la contrapresión y el bloqueo en su propio grupo de subprocesos.
Tenga en cuenta que flowOn podría usarse además de este flujo, pero realmente solo importaría para getInitialValue() , porque el hilo que ejecuta la devolución de llamada está controlado por la Source de todos modos.
Si agregar muchos oyentes es costoso para Source , también podría considerar compartir este flujo usando el operador shareIn() , de modo que varios suscriptores compartan la misma suscripción de oyente.