I have a process that spawns different number of events when triggered (I've assumed maximum of 50) and another process that needs to be started when the first event arrives. All events need to be passed to the second process, so I use a MutableSharedFlow with replay of size 50 (it cannot be started synchronously - Android system is responsible for creating it). As soon as the second process starts I'd like to reset the replay buffer, so that the next process that subscribes won't receive the old events as well.
I have the following piece of code:
val mutableFlow = MutableSharedFlow<Event>(replay = 50)
val sharedFlow = mutableFlow.onSubscription {
log.info { "New subscription, clearing replay cache." }
val events = mutableFlow.replayCache
mutableFlow.resetReplayCache() // clears events before they arrive to subscriber
}
emitSomeValues(mutableFlow)
So this doesn't work because onSubscription is called after subscription, but before events consumption. I struggle to find any clean solution to that, so any help would be appreciated!
Still no solution for the above :(