change(room members): makes sure to subscribe to timeline items changes
This commit is contained in:
parent
5f2453b128
commit
a3c81d5f25
10 changed files with 162 additions and 106 deletions
|
|
@ -54,6 +54,7 @@ interface Timeline : AutoCloseable {
|
||||||
|
|
||||||
val mode: Mode
|
val mode: Mode
|
||||||
val membershipChangeEventReceived: Flow<Unit>
|
val membershipChangeEventReceived: Flow<Unit>
|
||||||
|
val onSyncedEventReceived: Flow<Unit>
|
||||||
suspend fun sendReadReceipt(eventId: EventId, receiptType: ReceiptType): Result<Unit>
|
suspend fun sendReadReceipt(eventId: EventId, receiptType: ReceiptType): Result<Unit>
|
||||||
suspend fun markAsRead(receiptType: ReceiptType): Result<Unit>
|
suspend fun markAsRead(receiptType: ReceiptType): Result<Unit>
|
||||||
suspend fun paginate(direction: PaginationDirection): Result<Boolean>
|
suspend fun paginate(direction: PaginationDirection): Result<Boolean>
|
||||||
|
|
@ -233,4 +234,5 @@ interface Timeline : AutoCloseable {
|
||||||
* Get the latest event id of the timeline.
|
* Get the latest event id of the timeline.
|
||||||
*/
|
*/
|
||||||
suspend fun getLatestEventId(): Result<EventId?>
|
suspend fun getLatestEventId(): Result<EventId?>
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -50,6 +50,7 @@ import io.element.android.libraries.matrix.impl.widget.generateWidgetWebViewUrl
|
||||||
import io.element.android.services.toolbox.api.systemclock.SystemClock
|
import io.element.android.services.toolbox.api.systemclock.SystemClock
|
||||||
import kotlinx.coroutines.flow.Flow
|
import kotlinx.coroutines.flow.Flow
|
||||||
import kotlinx.coroutines.flow.MutableStateFlow
|
import kotlinx.coroutines.flow.MutableStateFlow
|
||||||
|
import kotlinx.coroutines.flow.SharingStarted.Companion.WhileSubscribed
|
||||||
import kotlinx.coroutines.flow.combine
|
import kotlinx.coroutines.flow.combine
|
||||||
import kotlinx.coroutines.flow.distinctUntilChanged
|
import kotlinx.coroutines.flow.distinctUntilChanged
|
||||||
import kotlinx.coroutines.flow.drop
|
import kotlinx.coroutines.flow.drop
|
||||||
|
|
@ -57,6 +58,7 @@ import kotlinx.coroutines.flow.launchIn
|
||||||
import kotlinx.coroutines.flow.map
|
import kotlinx.coroutines.flow.map
|
||||||
import kotlinx.coroutines.flow.onEach
|
import kotlinx.coroutines.flow.onEach
|
||||||
import kotlinx.coroutines.flow.onStart
|
import kotlinx.coroutines.flow.onStart
|
||||||
|
import kotlinx.coroutines.flow.stateIn
|
||||||
import kotlinx.coroutines.withContext
|
import kotlinx.coroutines.withContext
|
||||||
import org.matrix.rustcomponents.sdk.DateDividerMode
|
import org.matrix.rustcomponents.sdk.DateDividerMode
|
||||||
import org.matrix.rustcomponents.sdk.IdentityStatusChangeListener
|
import org.matrix.rustcomponents.sdk.IdentityStatusChangeListener
|
||||||
|
|
@ -91,8 +93,6 @@ class JoinedRustRoom(
|
||||||
private val roomDispatcher = coroutineDispatchers.io.limitedParallelism(32)
|
private val roomDispatcher = coroutineDispatchers.io.limitedParallelism(32)
|
||||||
private val innerRoom = baseRoom.innerRoom
|
private val innerRoom = baseRoom.innerRoom
|
||||||
|
|
||||||
override val syncUpdateFlow = MutableStateFlow(0L)
|
|
||||||
|
|
||||||
override val roomTypingMembersFlow: Flow<List<UserId>> = mxCallbackFlow {
|
override val roomTypingMembersFlow: Flow<List<UserId>> = mxCallbackFlow {
|
||||||
val initial = emptyList<UserId>()
|
val initial = emptyList<UserId>()
|
||||||
channel.trySend(initial)
|
channel.trySend(initial)
|
||||||
|
|
@ -135,11 +135,21 @@ class JoinedRustRoom(
|
||||||
|
|
||||||
override val roomNotificationSettingsStateFlow = MutableStateFlow<RoomNotificationSettingsState>(RoomNotificationSettingsState.Unknown)
|
override val roomNotificationSettingsStateFlow = MutableStateFlow<RoomNotificationSettingsState>(RoomNotificationSettingsState.Unknown)
|
||||||
|
|
||||||
override val liveTimeline = liveInnerTimeline.map(mode = Timeline.Mode.Live) {
|
override val liveTimeline = liveInnerTimeline.map(mode = Timeline.Mode.Live)
|
||||||
syncUpdateFlow.value = systemClock.epochMillis()
|
|
||||||
}
|
override val syncUpdateFlow = liveTimeline
|
||||||
|
.onSyncedEventReceived.map { systemClock.epochMillis() }
|
||||||
|
.stateIn(
|
||||||
|
scope = roomCoroutineScope,
|
||||||
|
started = WhileSubscribed(),
|
||||||
|
initialValue = systemClock.epochMillis(),
|
||||||
|
)
|
||||||
|
|
||||||
init {
|
init {
|
||||||
|
subscribeToRoomMembersChange()
|
||||||
|
}
|
||||||
|
|
||||||
|
private fun subscribeToRoomMembersChange() {
|
||||||
val powerLevelChanges = roomInfoFlow.map { it.roomPowerLevels }.distinctUntilChanged()
|
val powerLevelChanges = roomInfoFlow.map { it.roomPowerLevels }.distinctUntilChanged()
|
||||||
val membershipChanges = liveTimeline.membershipChangeEventReceived.onStart { emit(Unit) }
|
val membershipChanges = liveTimeline.membershipChangeEventReceived.onStart { emit(Unit) }
|
||||||
combine(membershipChanges, powerLevelChanges) { _, _ -> }
|
combine(membershipChanges, powerLevelChanges) { _, _ -> }
|
||||||
|
|
@ -478,7 +488,6 @@ class JoinedRustRoom(
|
||||||
|
|
||||||
private fun InnerTimeline.map(
|
private fun InnerTimeline.map(
|
||||||
mode: Timeline.Mode,
|
mode: Timeline.Mode,
|
||||||
onNewSyncedEvent: () -> Unit = {},
|
|
||||||
): Timeline {
|
): Timeline {
|
||||||
val timelineCoroutineScope = roomCoroutineScope.childScope(coroutineDispatchers.main, "TimelineScope-$roomId-$this")
|
val timelineCoroutineScope = roomCoroutineScope.childScope(coroutineDispatchers.main, "TimelineScope-$roomId-$this")
|
||||||
return RustTimeline(
|
return RustTimeline(
|
||||||
|
|
@ -489,7 +498,6 @@ class JoinedRustRoom(
|
||||||
coroutineScope = timelineCoroutineScope,
|
coroutineScope = timelineCoroutineScope,
|
||||||
dispatcher = roomDispatcher,
|
dispatcher = roomDispatcher,
|
||||||
roomContentForwarder = roomContentForwarder,
|
roomContentForwarder = roomContentForwarder,
|
||||||
onNewSyncedEvent = onNewSyncedEvent,
|
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -7,9 +7,10 @@
|
||||||
|
|
||||||
package io.element.android.libraries.matrix.impl.timeline
|
package io.element.android.libraries.matrix.impl.timeline
|
||||||
|
|
||||||
|
import androidx.compose.ui.util.fastForEach
|
||||||
import io.element.android.libraries.matrix.api.timeline.MatrixTimelineItem
|
import io.element.android.libraries.matrix.api.timeline.MatrixTimelineItem
|
||||||
import io.element.android.libraries.matrix.api.timeline.item.event.RoomMembershipContent
|
import io.element.android.libraries.matrix.api.timeline.item.event.RoomMembershipContent
|
||||||
import kotlinx.coroutines.flow.Flow
|
import io.element.android.libraries.matrix.api.timeline.item.event.TimelineItemEventOrigin
|
||||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||||
import kotlinx.coroutines.flow.first
|
import kotlinx.coroutines.flow.first
|
||||||
import kotlinx.coroutines.sync.Mutex
|
import kotlinx.coroutines.sync.Mutex
|
||||||
|
|
@ -20,58 +21,60 @@ import timber.log.Timber
|
||||||
|
|
||||||
internal class MatrixTimelineDiffProcessor(
|
internal class MatrixTimelineDiffProcessor(
|
||||||
private val timelineItems: MutableSharedFlow<List<MatrixTimelineItem>>,
|
private val timelineItems: MutableSharedFlow<List<MatrixTimelineItem>>,
|
||||||
private val timelineItemFactory: MatrixTimelineItemMapper,
|
private val membershipChangeEventReceivedFlow: MutableSharedFlow<Unit>,
|
||||||
|
private val syncedEventReceivedFlow: MutableSharedFlow<Unit>,
|
||||||
|
private val timelineItemMapper: MatrixTimelineItemMapper,
|
||||||
) {
|
) {
|
||||||
private val mutex = Mutex()
|
private val mutex = Mutex()
|
||||||
|
|
||||||
private val _membershipChangeEventReceived = MutableSharedFlow<Unit>(extraBufferCapacity = 1)
|
|
||||||
val membershipChangeEventReceived: Flow<Unit> = _membershipChangeEventReceived
|
|
||||||
|
|
||||||
suspend fun postDiffs(diffs: List<TimelineDiff>) {
|
suspend fun postDiffs(diffs: List<TimelineDiff>) {
|
||||||
updateTimelineItems {
|
mutex.withLock {
|
||||||
Timber.v("Update timeline items from postDiffs (with ${diffs.size} items) on ${Thread.currentThread()}")
|
Timber.v("Update timeline items from postDiffs (with ${diffs.size} items) on ${Thread.currentThread()}")
|
||||||
diffs.forEach { diff ->
|
val result = processDiffs(diffs)
|
||||||
applyDiff(diff)
|
timelineItems.emit(result.items())
|
||||||
|
if (result.hasNewEventsFromSync()) {
|
||||||
|
syncedEventReceivedFlow.emit(Unit)
|
||||||
|
}
|
||||||
|
if (result.hasMembershipChangeEventFromSync()) {
|
||||||
|
membershipChangeEventReceivedFlow.emit(Unit)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private suspend fun updateTimelineItems(block: MutableList<MatrixTimelineItem>.() -> Unit) =
|
private suspend fun processDiffs(diffs: List<TimelineDiff>): DiffingResult {
|
||||||
mutex.withLock {
|
val mutableTimelineItems = if (timelineItems.replayCache.isNotEmpty()) {
|
||||||
val mutableTimelineItems = if (timelineItems.replayCache.isNotEmpty()) {
|
timelineItems.first().toMutableList()
|
||||||
timelineItems.first().toMutableList()
|
} else {
|
||||||
} else {
|
mutableListOf()
|
||||||
mutableListOf()
|
|
||||||
}
|
|
||||||
block(mutableTimelineItems)
|
|
||||||
timelineItems.tryEmit(mutableTimelineItems)
|
|
||||||
}
|
}
|
||||||
|
val result = DiffingResult(items = mutableTimelineItems)
|
||||||
|
diffs.forEach { diff ->
|
||||||
|
result.applyDiff(diff)
|
||||||
|
}
|
||||||
|
return result
|
||||||
|
}
|
||||||
|
|
||||||
private fun MutableList<MatrixTimelineItem>.applyDiff(diff: TimelineDiff) {
|
private fun DiffingResult.applyDiff(diff: TimelineDiff) {
|
||||||
when (diff) {
|
when (diff) {
|
||||||
is TimelineDiff.Append -> {
|
is TimelineDiff.Append -> {
|
||||||
val items = diff.values.map { it.asMatrixTimelineItem() }
|
diff.values.fastForEach { item ->
|
||||||
addAll(items)
|
add(item.map())
|
||||||
|
}
|
||||||
}
|
}
|
||||||
is TimelineDiff.PushBack -> {
|
is TimelineDiff.PushBack -> {
|
||||||
val item = diff.value.asMatrixTimelineItem()
|
val item = diff.value.map()
|
||||||
if (item is MatrixTimelineItem.Event && item.event.content is RoomMembershipContent) {
|
|
||||||
// TODO - This is a temporary solution to notify the room screen about membership changes
|
|
||||||
// Ideally, this should be implemented by the Rust SDK
|
|
||||||
_membershipChangeEventReceived.tryEmit(Unit)
|
|
||||||
}
|
|
||||||
add(item)
|
add(item)
|
||||||
}
|
}
|
||||||
is TimelineDiff.PushFront -> {
|
is TimelineDiff.PushFront -> {
|
||||||
val item = diff.value.asMatrixTimelineItem()
|
val item = diff.value.map()
|
||||||
add(0, item)
|
add(0, item)
|
||||||
}
|
}
|
||||||
is TimelineDiff.Set -> {
|
is TimelineDiff.Set -> {
|
||||||
val item = diff.value.asMatrixTimelineItem()
|
val item = diff.value.map()
|
||||||
set(diff.index.toInt(), item)
|
set(diff.index.toInt(), item)
|
||||||
}
|
}
|
||||||
is TimelineDiff.Insert -> {
|
is TimelineDiff.Insert -> {
|
||||||
val item = diff.value.asMatrixTimelineItem()
|
val item = diff.value.map()
|
||||||
add(diff.index.toInt(), item)
|
add(diff.index.toInt(), item)
|
||||||
}
|
}
|
||||||
is TimelineDiff.Remove -> {
|
is TimelineDiff.Remove -> {
|
||||||
|
|
@ -79,25 +82,92 @@ internal class MatrixTimelineDiffProcessor(
|
||||||
}
|
}
|
||||||
is TimelineDiff.Reset -> {
|
is TimelineDiff.Reset -> {
|
||||||
clear()
|
clear()
|
||||||
val items = diff.values.map { it.asMatrixTimelineItem() }
|
diff.values.fastForEach { item ->
|
||||||
addAll(items)
|
add(item.map())
|
||||||
|
}
|
||||||
}
|
}
|
||||||
TimelineDiff.PopFront -> {
|
TimelineDiff.PopFront -> {
|
||||||
removeFirstOrNull()
|
removeFirst()
|
||||||
}
|
}
|
||||||
TimelineDiff.PopBack -> {
|
TimelineDiff.PopBack -> {
|
||||||
removeLastOrNull()
|
removeLast()
|
||||||
}
|
}
|
||||||
TimelineDiff.Clear -> {
|
TimelineDiff.Clear -> {
|
||||||
clear()
|
clear()
|
||||||
}
|
}
|
||||||
is TimelineDiff.Truncate -> {
|
is TimelineDiff.Truncate -> {
|
||||||
subList(diff.length.toInt(), size).clear()
|
truncate(diff.length.toInt())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun TimelineItem.asMatrixTimelineItem(): MatrixTimelineItem {
|
private fun TimelineItem.map(): MatrixTimelineItem {
|
||||||
return timelineItemFactory.map(this)
|
return timelineItemMapper.map(this)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private class DiffingResult(
|
||||||
|
private val items: MutableList<MatrixTimelineItem>,
|
||||||
|
private var hasNewEventsFromSync: Boolean = false,
|
||||||
|
private var hasMembershipChangeEventFromSync: Boolean = false,
|
||||||
|
) {
|
||||||
|
|
||||||
|
fun items(): List<MatrixTimelineItem> = items
|
||||||
|
fun hasNewEventsFromSync(): Boolean = hasNewEventsFromSync
|
||||||
|
fun hasMembershipChangeEventFromSync(): Boolean = hasMembershipChangeEventFromSync
|
||||||
|
|
||||||
|
fun add(item: MatrixTimelineItem) {
|
||||||
|
processItem(item)
|
||||||
|
items.add(item)
|
||||||
|
}
|
||||||
|
|
||||||
|
fun add(index: Int, item: MatrixTimelineItem) {
|
||||||
|
processItem(item)
|
||||||
|
items.add(index, item)
|
||||||
|
}
|
||||||
|
|
||||||
|
fun set(index: Int, item: MatrixTimelineItem) {
|
||||||
|
processItem(item)
|
||||||
|
items[index] = item
|
||||||
|
}
|
||||||
|
|
||||||
|
fun removeAt(index: Int) {
|
||||||
|
items.removeAt(index)
|
||||||
|
}
|
||||||
|
|
||||||
|
fun removeFirst() {
|
||||||
|
items.removeFirstOrNull()
|
||||||
|
}
|
||||||
|
|
||||||
|
fun removeLast() {
|
||||||
|
items.removeLastOrNull()
|
||||||
|
}
|
||||||
|
|
||||||
|
fun truncate(length: Int) {
|
||||||
|
items.subList(length, items.size).clear()
|
||||||
|
}
|
||||||
|
|
||||||
|
fun clear() {
|
||||||
|
items.clear()
|
||||||
|
}
|
||||||
|
|
||||||
|
private fun processItem(item: MatrixTimelineItem) {
|
||||||
|
if (skipProcessing()) return
|
||||||
|
when (item) {
|
||||||
|
is MatrixTimelineItem.Event -> {
|
||||||
|
if (item.event.origin == TimelineItemEventOrigin.SYNC) {
|
||||||
|
hasNewEventsFromSync = true
|
||||||
|
when (item.event.content) {
|
||||||
|
is RoomMembershipContent -> hasMembershipChangeEventFromSync = true
|
||||||
|
else -> Unit
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
else -> Unit
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private fun skipProcessing(): Boolean {
|
||||||
|
return hasNewEventsFromSync && hasMembershipChangeEventFromSync
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -80,16 +80,18 @@ private const val PAGINATION_SIZE = 50
|
||||||
class RustTimeline(
|
class RustTimeline(
|
||||||
private val inner: InnerTimeline,
|
private val inner: InnerTimeline,
|
||||||
override val mode: Timeline.Mode,
|
override val mode: Timeline.Mode,
|
||||||
systemClock: SystemClock,
|
private val systemClock: SystemClock,
|
||||||
private val joinedRoom: JoinedRoom,
|
private val joinedRoom: JoinedRoom,
|
||||||
private val coroutineScope: CoroutineScope,
|
private val coroutineScope: CoroutineScope,
|
||||||
private val dispatcher: CoroutineDispatcher,
|
private val dispatcher: CoroutineDispatcher,
|
||||||
private val roomContentForwarder: RoomContentForwarder,
|
private val roomContentForwarder: RoomContentForwarder,
|
||||||
onNewSyncedEvent: () -> Unit,
|
|
||||||
) : Timeline {
|
) : Timeline {
|
||||||
private val _timelineItems: MutableSharedFlow<List<MatrixTimelineItem>> =
|
private val _timelineItems: MutableSharedFlow<List<MatrixTimelineItem>> =
|
||||||
MutableSharedFlow(replay = 1, extraBufferCapacity = Int.MAX_VALUE)
|
MutableSharedFlow(replay = 1, extraBufferCapacity = Int.MAX_VALUE)
|
||||||
|
|
||||||
|
private val _membershipChangeEventReceived = MutableSharedFlow<Unit>(extraBufferCapacity = 1)
|
||||||
|
private val _onSyncedEventReceived: MutableSharedFlow<Unit> = MutableSharedFlow(extraBufferCapacity = 1)
|
||||||
|
|
||||||
private val timelineEventContentMapper = TimelineEventContentMapper()
|
private val timelineEventContentMapper = TimelineEventContentMapper()
|
||||||
private val inReplyToMapper = InReplyToMapper(timelineEventContentMapper)
|
private val inReplyToMapper = InReplyToMapper(timelineEventContentMapper)
|
||||||
private val timelineItemMapper = MatrixTimelineItemMapper(
|
private val timelineItemMapper = MatrixTimelineItemMapper(
|
||||||
|
|
@ -98,18 +100,19 @@ class RustTimeline(
|
||||||
virtualTimelineItemMapper = VirtualTimelineItemMapper(),
|
virtualTimelineItemMapper = VirtualTimelineItemMapper(),
|
||||||
eventTimelineItemMapper = EventTimelineItemMapper(
|
eventTimelineItemMapper = EventTimelineItemMapper(
|
||||||
contentMapper = timelineEventContentMapper
|
contentMapper = timelineEventContentMapper
|
||||||
)
|
),
|
||||||
)
|
)
|
||||||
private val timelineDiffProcessor = MatrixTimelineDiffProcessor(
|
private val timelineDiffProcessor = MatrixTimelineDiffProcessor(
|
||||||
timelineItems = _timelineItems,
|
timelineItems = _timelineItems,
|
||||||
timelineItemFactory = timelineItemMapper,
|
membershipChangeEventReceivedFlow = _membershipChangeEventReceived,
|
||||||
|
syncedEventReceivedFlow = _onSyncedEventReceived,
|
||||||
|
timelineItemMapper = timelineItemMapper,
|
||||||
)
|
)
|
||||||
private val timelineItemsSubscriber = TimelineItemsSubscriber(
|
private val timelineItemsSubscriber = TimelineItemsSubscriber(
|
||||||
timeline = inner,
|
timeline = inner,
|
||||||
timelineCoroutineScope = coroutineScope,
|
timelineCoroutineScope = coroutineScope,
|
||||||
timelineDiffProcessor = timelineDiffProcessor,
|
timelineDiffProcessor = timelineDiffProcessor,
|
||||||
dispatcher = dispatcher,
|
dispatcher = dispatcher,
|
||||||
onNewSyncedEvent = onNewSyncedEvent,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
private val roomBeginningPostProcessor = RoomBeginningPostProcessor(mode)
|
private val roomBeginningPostProcessor = RoomBeginningPostProcessor(mode)
|
||||||
|
|
@ -151,7 +154,13 @@ class RustTimeline(
|
||||||
.launchIn(this)
|
.launchIn(this)
|
||||||
}
|
}
|
||||||
|
|
||||||
override val membershipChangeEventReceived: Flow<Unit> = timelineDiffProcessor.membershipChangeEventReceived
|
override val membershipChangeEventReceived: Flow<Unit> = _membershipChangeEventReceived
|
||||||
|
.onStart { timelineItemsSubscriber.subscribeIfNeeded() }
|
||||||
|
.onCompletion { timelineItemsSubscriber.unsubscribeIfNeeded() }
|
||||||
|
|
||||||
|
override val onSyncedEventReceived: Flow<Unit> = _onSyncedEventReceived
|
||||||
|
.onStart { timelineItemsSubscriber.subscribeIfNeeded() }
|
||||||
|
.onCompletion { timelineItemsSubscriber.unsubscribeIfNeeded() }
|
||||||
|
|
||||||
override suspend fun sendReadReceipt(eventId: EventId, receiptType: ReceiptType): Result<Unit> = withContext(dispatcher) {
|
override suspend fun sendReadReceipt(eventId: EventId, receiptType: ReceiptType): Result<Unit> = withContext(dispatcher) {
|
||||||
runCatchingExceptions {
|
runCatchingExceptions {
|
||||||
|
|
|
||||||
|
|
@ -1,32 +0,0 @@
|
||||||
/*
|
|
||||||
* Copyright 2023, 2024 New Vector Ltd.
|
|
||||||
*
|
|
||||||
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
|
|
||||||
* Please see LICENSE files in the repository root for full details.
|
|
||||||
*/
|
|
||||||
|
|
||||||
package io.element.android.libraries.matrix.impl.timeline
|
|
||||||
|
|
||||||
import org.matrix.rustcomponents.sdk.TimelineDiff
|
|
||||||
import org.matrix.rustcomponents.sdk.TimelineItem
|
|
||||||
import uniffi.matrix_sdk_ui.EventItemOrigin
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Tries to get an event origin from the TimelineDiff.
|
|
||||||
* If there is multiple events in the diff, uses the first one as it should be a good indicator.
|
|
||||||
*/
|
|
||||||
internal fun TimelineDiff.eventOrigin(): EventItemOrigin? {
|
|
||||||
return when (this) {
|
|
||||||
is TimelineDiff.Append -> values.firstOrNull()?.eventOrigin()
|
|
||||||
is TimelineDiff.PushBack -> value.eventOrigin()
|
|
||||||
is TimelineDiff.PushFront -> value.eventOrigin()
|
|
||||||
is TimelineDiff.Set -> value.eventOrigin()
|
|
||||||
is TimelineDiff.Insert -> value.eventOrigin()
|
|
||||||
is TimelineDiff.Reset -> values.firstOrNull()?.eventOrigin()
|
|
||||||
else -> null
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private fun TimelineItem.eventOrigin(): EventItemOrigin? {
|
|
||||||
return asEvent()?.origin
|
|
||||||
}
|
|
||||||
|
|
@ -11,13 +11,11 @@ import io.element.android.libraries.core.coroutine.childScope
|
||||||
import kotlinx.coroutines.CoroutineDispatcher
|
import kotlinx.coroutines.CoroutineDispatcher
|
||||||
import kotlinx.coroutines.CoroutineScope
|
import kotlinx.coroutines.CoroutineScope
|
||||||
import kotlinx.coroutines.cancelChildren
|
import kotlinx.coroutines.cancelChildren
|
||||||
import kotlinx.coroutines.coroutineScope
|
|
||||||
import kotlinx.coroutines.flow.launchIn
|
import kotlinx.coroutines.flow.launchIn
|
||||||
import kotlinx.coroutines.flow.onEach
|
import kotlinx.coroutines.flow.onEach
|
||||||
import kotlinx.coroutines.sync.Mutex
|
import kotlinx.coroutines.sync.Mutex
|
||||||
import kotlinx.coroutines.sync.withLock
|
import kotlinx.coroutines.sync.withLock
|
||||||
import org.matrix.rustcomponents.sdk.Timeline
|
import org.matrix.rustcomponents.sdk.Timeline
|
||||||
import uniffi.matrix_sdk_ui.EventItemOrigin
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* This class is responsible for subscribing to a timeline and post the items/diffs to the timelineDiffProcessor.
|
* This class is responsible for subscribing to a timeline and post the items/diffs to the timelineDiffProcessor.
|
||||||
|
|
@ -28,7 +26,6 @@ internal class TimelineItemsSubscriber(
|
||||||
dispatcher: CoroutineDispatcher,
|
dispatcher: CoroutineDispatcher,
|
||||||
private val timeline: Timeline,
|
private val timeline: Timeline,
|
||||||
private val timelineDiffProcessor: MatrixTimelineDiffProcessor,
|
private val timelineDiffProcessor: MatrixTimelineDiffProcessor,
|
||||||
private val onNewSyncedEvent: () -> Unit,
|
|
||||||
) {
|
) {
|
||||||
private var subscriptionCount = 0
|
private var subscriptionCount = 0
|
||||||
private val mutex = Mutex()
|
private val mutex = Mutex()
|
||||||
|
|
@ -43,9 +40,6 @@ internal class TimelineItemsSubscriber(
|
||||||
if (subscriptionCount == 0) {
|
if (subscriptionCount == 0) {
|
||||||
timeline.timelineDiffFlow()
|
timeline.timelineDiffFlow()
|
||||||
.onEach { diffs ->
|
.onEach { diffs ->
|
||||||
if (diffs.any { diff -> diff.eventOrigin() == EventItemOrigin.SYNC }) {
|
|
||||||
onNewSyncedEvent()
|
|
||||||
}
|
|
||||||
timelineDiffProcessor.postDiffs(diffs)
|
timelineDiffProcessor.postDiffs(diffs)
|
||||||
}
|
}
|
||||||
.launchIn(coroutineScope)
|
.launchIn(coroutineScope)
|
||||||
|
|
|
||||||
|
|
@ -168,10 +168,12 @@ class MatrixTimelineDiffProcessorTest {
|
||||||
}
|
}
|
||||||
|
|
||||||
internal fun TestScope.createMatrixTimelineDiffProcessor(
|
internal fun TestScope.createMatrixTimelineDiffProcessor(
|
||||||
timelineItems: MutableSharedFlow<List<MatrixTimelineItem>>,
|
timelineItems: MutableSharedFlow<List<MatrixTimelineItem>> = MutableSharedFlow(),
|
||||||
): MatrixTimelineDiffProcessor {
|
membershipChangeEventReceivedFlow: MutableSharedFlow<Unit> = MutableSharedFlow(),
|
||||||
|
syncedEventReceivedFlow: MutableSharedFlow<Unit> = MutableSharedFlow(),
|
||||||
|
): MatrixTimelineDiffProcessor {
|
||||||
val timelineEventContentMapper = TimelineEventContentMapper()
|
val timelineEventContentMapper = TimelineEventContentMapper()
|
||||||
val timelineItemMapper = MatrixTimelineItemMapper(
|
val timelineItemFactory = MatrixTimelineItemMapper(
|
||||||
fetchDetailsForEvent = { _ -> Result.success(Unit) },
|
fetchDetailsForEvent = { _ -> Result.success(Unit) },
|
||||||
coroutineScope = this,
|
coroutineScope = this,
|
||||||
virtualTimelineItemMapper = VirtualTimelineItemMapper(),
|
virtualTimelineItemMapper = VirtualTimelineItemMapper(),
|
||||||
|
|
@ -181,6 +183,8 @@ internal fun TestScope.createMatrixTimelineDiffProcessor(
|
||||||
)
|
)
|
||||||
return MatrixTimelineDiffProcessor(
|
return MatrixTimelineDiffProcessor(
|
||||||
timelineItems = timelineItems,
|
timelineItems = timelineItems,
|
||||||
timelineItemFactory = timelineItemMapper,
|
membershipChangeEventReceivedFlow = membershipChangeEventReceivedFlow,
|
||||||
|
syncedEventReceivedFlow = syncedEventReceivedFlow,
|
||||||
|
timelineItemMapper = timelineItemFactory,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -98,7 +98,6 @@ private fun TestScope.createRustTimeline(
|
||||||
coroutineScope: CoroutineScope = backgroundScope,
|
coroutineScope: CoroutineScope = backgroundScope,
|
||||||
dispatcher: CoroutineDispatcher = testCoroutineDispatchers().io,
|
dispatcher: CoroutineDispatcher = testCoroutineDispatchers().io,
|
||||||
roomContentForwarder: RoomContentForwarder = RoomContentForwarder(FakeFfiRoomListService()),
|
roomContentForwarder: RoomContentForwarder = RoomContentForwarder(FakeFfiRoomListService()),
|
||||||
onNewSyncedEvent: () -> Unit = {},
|
|
||||||
): RustTimeline {
|
): RustTimeline {
|
||||||
return RustTimeline(
|
return RustTimeline(
|
||||||
inner = inner,
|
inner = inner,
|
||||||
|
|
@ -108,6 +107,5 @@ private fun TestScope.createRustTimeline(
|
||||||
coroutineScope = coroutineScope,
|
coroutineScope = coroutineScope,
|
||||||
dispatcher = dispatcher,
|
dispatcher = dispatcher,
|
||||||
roomContentForwarder = roomContentForwarder,
|
roomContentForwarder = roomContentForwarder,
|
||||||
onNewSyncedEvent = onNewSyncedEvent,
|
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -13,8 +13,6 @@ import io.element.android.libraries.matrix.api.timeline.MatrixTimelineItem
|
||||||
import io.element.android.libraries.matrix.impl.fixtures.factories.aRustEventTimelineItem
|
import io.element.android.libraries.matrix.impl.fixtures.factories.aRustEventTimelineItem
|
||||||
import io.element.android.libraries.matrix.impl.fixtures.fakes.FakeFfiTimeline
|
import io.element.android.libraries.matrix.impl.fixtures.fakes.FakeFfiTimeline
|
||||||
import io.element.android.libraries.matrix.impl.fixtures.fakes.FakeFfiTimelineItem
|
import io.element.android.libraries.matrix.impl.fixtures.fakes.FakeFfiTimelineItem
|
||||||
import io.element.android.tests.testutils.lambda.lambdaError
|
|
||||||
import io.element.android.tests.testutils.lambda.lambdaRecorder
|
|
||||||
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||||
import kotlinx.coroutines.test.StandardTestDispatcher
|
import kotlinx.coroutines.test.StandardTestDispatcher
|
||||||
|
|
@ -35,9 +33,12 @@ class TimelineItemsSubscriberTest {
|
||||||
val timelineItems: MutableSharedFlow<List<MatrixTimelineItem>> =
|
val timelineItems: MutableSharedFlow<List<MatrixTimelineItem>> =
|
||||||
MutableSharedFlow(replay = 1, extraBufferCapacity = Int.MAX_VALUE)
|
MutableSharedFlow(replay = 1, extraBufferCapacity = Int.MAX_VALUE)
|
||||||
val timeline = FakeFfiTimeline()
|
val timeline = FakeFfiTimeline()
|
||||||
|
val diffProcessor = createMatrixTimelineDiffProcessor(
|
||||||
|
timelineItems = timelineItems,
|
||||||
|
)
|
||||||
val timelineItemsSubscriber = createTimelineItemsSubscriber(
|
val timelineItemsSubscriber = createTimelineItemsSubscriber(
|
||||||
timeline = timeline,
|
timeline = timeline,
|
||||||
timelineItems = timelineItems,
|
timelineDiffProcessor = diffProcessor,
|
||||||
)
|
)
|
||||||
timelineItems.test {
|
timelineItems.test {
|
||||||
timelineItemsSubscriber.subscribeIfNeeded()
|
timelineItemsSubscriber.subscribeIfNeeded()
|
||||||
|
|
@ -56,9 +57,12 @@ class TimelineItemsSubscriberTest {
|
||||||
val timelineItems: MutableSharedFlow<List<MatrixTimelineItem>> =
|
val timelineItems: MutableSharedFlow<List<MatrixTimelineItem>> =
|
||||||
MutableSharedFlow(replay = 1, extraBufferCapacity = Int.MAX_VALUE)
|
MutableSharedFlow(replay = 1, extraBufferCapacity = Int.MAX_VALUE)
|
||||||
val timeline = FakeFfiTimeline()
|
val timeline = FakeFfiTimeline()
|
||||||
|
val diffProcessor = createMatrixTimelineDiffProcessor(
|
||||||
|
timelineItems = timelineItems,
|
||||||
|
)
|
||||||
val timelineItemsSubscriber = createTimelineItemsSubscriber(
|
val timelineItemsSubscriber = createTimelineItemsSubscriber(
|
||||||
timeline = timeline,
|
timeline = timeline,
|
||||||
timelineItems = timelineItems,
|
timelineDiffProcessor = diffProcessor,
|
||||||
)
|
)
|
||||||
timelineItems.test {
|
timelineItems.test {
|
||||||
timelineItemsSubscriber.subscribeIfNeeded()
|
timelineItemsSubscriber.subscribeIfNeeded()
|
||||||
|
|
@ -73,15 +77,16 @@ class TimelineItemsSubscriberTest {
|
||||||
|
|
||||||
@Ignore("JNA direct mapping has broken unit tests with FFI fakes")
|
@Ignore("JNA direct mapping has broken unit tests with FFI fakes")
|
||||||
@Test
|
@Test
|
||||||
fun `when timeline emits an item with SYNC origin, the callback onNewSyncedEvent is invoked`() = runTest {
|
fun `when timeline emits an item with SYNC origin`() = runTest {
|
||||||
val timelineItems: MutableSharedFlow<List<MatrixTimelineItem>> =
|
val timelineItems: MutableSharedFlow<List<MatrixTimelineItem>> =
|
||||||
MutableSharedFlow(replay = 1, extraBufferCapacity = Int.MAX_VALUE)
|
MutableSharedFlow(replay = 1, extraBufferCapacity = Int.MAX_VALUE)
|
||||||
val timeline = FakeFfiTimeline()
|
val timeline = FakeFfiTimeline()
|
||||||
val onNewSyncedEventRecorder = lambdaRecorder<Unit> { }
|
val diffProcessor = createMatrixTimelineDiffProcessor(
|
||||||
|
timelineItems = timelineItems,
|
||||||
|
)
|
||||||
val timelineItemsSubscriber = createTimelineItemsSubscriber(
|
val timelineItemsSubscriber = createTimelineItemsSubscriber(
|
||||||
timeline = timeline,
|
timeline = timeline,
|
||||||
timelineItems = timelineItems,
|
timelineDiffProcessor = diffProcessor,
|
||||||
onNewSyncedEvent = onNewSyncedEventRecorder,
|
|
||||||
)
|
)
|
||||||
timelineItems.test {
|
timelineItems.test {
|
||||||
timelineItemsSubscriber.subscribeIfNeeded()
|
timelineItemsSubscriber.subscribeIfNeeded()
|
||||||
|
|
@ -100,7 +105,6 @@ class TimelineItemsSubscriberTest {
|
||||||
assertThat(final).isNotEmpty()
|
assertThat(final).isNotEmpty()
|
||||||
timelineItemsSubscriber.unsubscribeIfNeeded()
|
timelineItemsSubscriber.unsubscribeIfNeeded()
|
||||||
}
|
}
|
||||||
onNewSyncedEventRecorder.assertions().isCalledOnce()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Ignore("JNA direct mapping has broken unit tests with FFI fakes")
|
@Ignore("JNA direct mapping has broken unit tests with FFI fakes")
|
||||||
|
|
@ -116,14 +120,12 @@ class TimelineItemsSubscriberTest {
|
||||||
|
|
||||||
private fun TestScope.createTimelineItemsSubscriber(
|
private fun TestScope.createTimelineItemsSubscriber(
|
||||||
timeline: Timeline = FakeFfiTimeline(),
|
timeline: Timeline = FakeFfiTimeline(),
|
||||||
timelineItems: MutableSharedFlow<List<MatrixTimelineItem>> = MutableSharedFlow(replay = 1, extraBufferCapacity = Int.MAX_VALUE),
|
timelineDiffProcessor: MatrixTimelineDiffProcessor = createMatrixTimelineDiffProcessor(),
|
||||||
onNewSyncedEvent: () -> Unit = { lambdaError() },
|
|
||||||
): TimelineItemsSubscriber {
|
): TimelineItemsSubscriber {
|
||||||
return TimelineItemsSubscriber(
|
return TimelineItemsSubscriber(
|
||||||
timelineCoroutineScope = backgroundScope,
|
timelineCoroutineScope = backgroundScope,
|
||||||
dispatcher = StandardTestDispatcher(testScheduler),
|
dispatcher = StandardTestDispatcher(testScheduler),
|
||||||
timeline = timeline,
|
timeline = timeline,
|
||||||
timelineDiffProcessor = createMatrixTimelineDiffProcessor(timelineItems),
|
timelineDiffProcessor = timelineDiffProcessor,
|
||||||
onNewSyncedEvent = onNewSyncedEvent,
|
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -47,6 +47,7 @@ class FakeTimeline(
|
||||||
)
|
)
|
||||||
),
|
),
|
||||||
override val membershipChangeEventReceived: Flow<Unit> = MutableSharedFlow(),
|
override val membershipChangeEventReceived: Flow<Unit> = MutableSharedFlow(),
|
||||||
|
override val onSyncedEventReceived: Flow<Unit> = MutableSharedFlow(),
|
||||||
private val cancelSendResult: (TransactionId) -> Result<Unit> = { lambdaError() },
|
private val cancelSendResult: (TransactionId) -> Result<Unit> = { lambdaError() },
|
||||||
override val mode: Timeline.Mode = Timeline.Mode.Live,
|
override val mode: Timeline.Mode = Timeline.Mode.Live,
|
||||||
private val markAsReadResult: (ReceiptType) -> Result<Unit> = { lambdaError() },
|
private val markAsReadResult: (ReceiptType) -> Result<Unit> = { lambdaError() },
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue