Remove usage of async-uniffi as it leads to a deadlocks and memory leaks (#1381)

This commit is contained in:
ganfra 2023-09-20 12:52:57 +02:00
parent ab33d064bb
commit 87125e0ea3
9 changed files with 78 additions and 54 deletions

1
changelog.d/1381.misc Normal file
View file

@ -0,0 +1 @@
Remove usage of async-uniffi as it leads to a deadlocks and memory leaks.

View file

@ -97,8 +97,8 @@ class RustMatrixClient constructor(
private val innerRoomListService = syncService.roomListService() private val innerRoomListService = syncService.roomListService()
private val sessionDispatcher = dispatchers.io.limitedParallelism(64) private val sessionDispatcher = dispatchers.io.limitedParallelism(64)
private val sessionCoroutineScope = appCoroutineScope.childScope(dispatchers.main, "Session-${sessionId}") private val sessionCoroutineScope = appCoroutineScope.childScope(dispatchers.main, "Session-${sessionId}")
private val rustSyncService = RustSyncService(syncService, sessionCoroutineScope) private val rustSyncService = RustSyncService(syncService, dispatchers, sessionCoroutineScope)
private val verificationService = RustSessionVerificationService(rustSyncService) private val verificationService = RustSessionVerificationService(rustSyncService, dispatchers)
private val pushersService = RustPushersService( private val pushersService = RustPushersService(
client = client, client = client,
dispatchers = dispatchers, dispatchers = dispatchers,

View file

@ -53,8 +53,7 @@ class RustMatrixClientFactory @Inject constructor(
client.restoreSession(sessionData.toSession()) client.restoreSession(sessionData.toSession())
val syncService = client.syncService() val syncService = client.syncService().finishBlocking()
.finish()
RustMatrixClient( RustMatrixClient(
client = client, client = client,

View file

@ -47,15 +47,19 @@ class RustNotificationSettingsService(
notificationSettings.setDelegate(notificationSettingsDelegate) notificationSettings.setDelegate(notificationSettingsDelegate)
} }
override suspend fun getRoomNotificationSettings(roomId: RoomId, isEncrypted: Boolean, isOneToOne: Boolean): Result<RoomNotificationSettings> = override suspend fun getRoomNotificationSettings(roomId: RoomId, isEncrypted: Boolean, isOneToOne: Boolean): Result<RoomNotificationSettings> = withContext(
dispatchers.io
) {
runCatching { runCatching {
notificationSettings.getRoomNotificationSettings(roomId.value, isEncrypted, isOneToOne).let(RoomNotificationSettingsMapper::map) notificationSettings.getRoomNotificationSettingsBlocking(roomId.value, isEncrypted, isOneToOne).let(RoomNotificationSettingsMapper::map)
} }
}
override suspend fun getDefaultRoomNotificationMode(isEncrypted: Boolean, isOneToOne: Boolean): Result<RoomNotificationMode> = override suspend fun getDefaultRoomNotificationMode(isEncrypted: Boolean, isOneToOne: Boolean): Result<RoomNotificationMode> = withContext(dispatchers.io) {
runCatching { runCatching {
notificationSettings.getDefaultRoomNotificationMode(isEncrypted, isOneToOne).let(RoomNotificationSettingsMapper::mapMode) notificationSettings.getDefaultRoomNotificationModeBlocking(isEncrypted, isOneToOne).let(RoomNotificationSettingsMapper::mapMode)
} }
}
override suspend fun setDefaultRoomNotificationMode( override suspend fun setDefaultRoomNotificationMode(
isEncrypted: Boolean, isEncrypted: Boolean,
@ -63,19 +67,19 @@ class RustNotificationSettingsService(
isOneToOne: Boolean isOneToOne: Boolean
): Result<Unit> = withContext(dispatchers.io) { ): Result<Unit> = withContext(dispatchers.io) {
runCatching { runCatching {
notificationSettings.setDefaultRoomNotificationMode(isEncrypted, isOneToOne, mode.let(RoomNotificationSettingsMapper::mapMode)) notificationSettings.setDefaultRoomNotificationModeBlocking(isEncrypted, isOneToOne, mode.let(RoomNotificationSettingsMapper::mapMode))
} }
} }
override suspend fun setRoomNotificationMode(roomId: RoomId, mode: RoomNotificationMode): Result<Unit> = withContext(dispatchers.io) { override suspend fun setRoomNotificationMode(roomId: RoomId, mode: RoomNotificationMode): Result<Unit> = withContext(dispatchers.io) {
runCatching { runCatching {
notificationSettings.setRoomNotificationMode(roomId.value, mode.let(RoomNotificationSettingsMapper::mapMode)) notificationSettings.setRoomNotificationModeBlocking(roomId.value, mode.let(RoomNotificationSettingsMapper::mapMode))
} }
} }
override suspend fun restoreDefaultRoomNotificationMode(roomId: RoomId): Result<Unit> = withContext(dispatchers.io) { override suspend fun restoreDefaultRoomNotificationMode(roomId: RoomId): Result<Unit> = withContext(dispatchers.io) {
runCatching { runCatching {
notificationSettings.restoreDefaultRoomNotificationMode(roomId.value) notificationSettings.restoreDefaultRoomNotificationModeBlocking(roomId.value)
} }
} }
@ -83,31 +87,31 @@ class RustNotificationSettingsService(
override suspend fun unmuteRoom(roomId: RoomId, isEncrypted: Boolean, isOneToOne: Boolean) = withContext(dispatchers.io) { override suspend fun unmuteRoom(roomId: RoomId, isEncrypted: Boolean, isOneToOne: Boolean) = withContext(dispatchers.io) {
runCatching { runCatching {
notificationSettings.unmuteRoom(roomId.value, isEncrypted, isOneToOne) notificationSettings.unmuteRoomBlocking(roomId.value, isEncrypted, isOneToOne)
} }
} }
override suspend fun isRoomMentionEnabled(): Result<Boolean> = withContext(dispatchers.io) { override suspend fun isRoomMentionEnabled(): Result<Boolean> = withContext(dispatchers.io) {
runCatching { runCatching {
notificationSettings.isRoomMentionEnabled() notificationSettings.isRoomMentionEnabledBlocking()
} }
} }
override suspend fun setRoomMentionEnabled(enabled: Boolean): Result<Unit> = withContext(dispatchers.io) { override suspend fun setRoomMentionEnabled(enabled: Boolean): Result<Unit> = withContext(dispatchers.io) {
runCatching { runCatching {
notificationSettings.setRoomMentionEnabled(enabled) notificationSettings.setRoomMentionEnabledBlocking(enabled)
} }
} }
override suspend fun isCallEnabled(): Result<Boolean> = withContext(dispatchers.io) { override suspend fun isCallEnabled(): Result<Boolean> = withContext(dispatchers.io) {
runCatching { runCatching {
notificationSettings.isCallEnabled() notificationSettings.isCallEnabledBlocking()
} }
} }
override suspend fun setCallEnabled(enabled: Boolean): Result<Unit> = withContext(dispatchers.io) { override suspend fun setCallEnabled(enabled: Boolean): Result<Unit> = withContext(dispatchers.io) {
runCatching { runCatching {
notificationSettings.setCallEnabled(enabled) notificationSettings.setCallEnabledBlocking(enabled)
} }
} }
} }

View file

@ -55,11 +55,11 @@ import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.ExperimentalCoroutinesApi import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.cancel import kotlinx.coroutines.cancel
import kotlinx.coroutines.ensureActive
import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.withContext import kotlinx.coroutines.withContext
import kotlinx.coroutines.yield
import org.matrix.rustcomponents.sdk.RequiredState import org.matrix.rustcomponents.sdk.RequiredState
import org.matrix.rustcomponents.sdk.Room import org.matrix.rustcomponents.sdk.Room
import org.matrix.rustcomponents.sdk.RoomListItem import org.matrix.rustcomponents.sdk.RoomListItem
@ -188,12 +188,12 @@ class RustMatrixRoom(
_membersStateFlow.value = MatrixRoomMembersState.Pending(prevRoomMembers = currentMembers) _membersStateFlow.value = MatrixRoomMembersState.Pending(prevRoomMembers = currentMembers)
var rustMembers: List<RoomMember>? = null var rustMembers: List<RoomMember>? = null
try { try {
rustMembers = innerRoom.members().use { membersIterator -> rustMembers = innerRoom.membersBlocking().use { membersIterator ->
buildList { buildList {
while (true) { while (true) {
// Loading the whole membersIterator as a stop-gap measure. // Loading the whole membersIterator as a stop-gap measure.
// We should probably implement some sort of paging in the future. // We should probably implement some sort of paging in the future.
yield() ensureActive()
addAll(membersIterator.nextChunk(1000u) ?: break) addAll(membersIterator.nextChunk(1000u) ?: break)
} }
} }
@ -294,27 +294,27 @@ class RustMatrixRoom(
} }
} }
override suspend fun canUserInvite(userId: UserId): Result<Boolean> { override suspend fun canUserInvite(userId: UserId): Result<Boolean> = withContext(roomMembersDispatcher) {
return runCatching { runCatching {
innerRoom.canUserInvite(userId.value) innerRoom.canUserInviteBlocking(userId.value)
} }
} }
override suspend fun canUserRedact(userId: UserId): Result<Boolean> { override suspend fun canUserRedact(userId: UserId): Result<Boolean> = withContext(roomMembersDispatcher) {
return runCatching { runCatching {
innerRoom.canUserRedact(userId.value) innerRoom.canUserRedactBlocking(userId.value)
} }
} }
override suspend fun canUserSendState(userId: UserId, type: StateEventType): Result<Boolean> { override suspend fun canUserSendState(userId: UserId, type: StateEventType): Result<Boolean> = withContext(roomMembersDispatcher) {
return runCatching { runCatching {
innerRoom.canUserSendState(userId.value, type.map()) innerRoom.canUserSendStateBlocking(userId.value, type.map())
} }
} }
override suspend fun canUserSendMessage(userId: UserId, type: MessageEventType): Result<Boolean> { override suspend fun canUserSendMessage(userId: UserId, type: MessageEventType): Result<Boolean> = withContext(roomMembersDispatcher) {
return runCatching { runCatching {
innerRoom.canUserSendMessage(userId.value, type.map()) innerRoom.canUserSendMessageBlocking(userId.value, type.map())
} }
} }

View file

@ -16,6 +16,7 @@
package io.element.android.libraries.matrix.impl.sync package io.element.android.libraries.matrix.impl.sync
import io.element.android.libraries.core.coroutine.CoroutineDispatchers
import io.element.android.libraries.matrix.api.sync.SyncService import io.element.android.libraries.matrix.api.sync.SyncService
import io.element.android.libraries.matrix.api.sync.SyncState import io.element.android.libraries.matrix.api.sync.SyncState
import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.CoroutineScope
@ -25,27 +26,33 @@ import kotlinx.coroutines.flow.distinctUntilChanged
import kotlinx.coroutines.flow.map import kotlinx.coroutines.flow.map
import kotlinx.coroutines.flow.onEach import kotlinx.coroutines.flow.onEach
import kotlinx.coroutines.flow.stateIn import kotlinx.coroutines.flow.stateIn
import kotlinx.coroutines.withContext
import org.matrix.rustcomponents.sdk.SyncServiceInterface import org.matrix.rustcomponents.sdk.SyncServiceInterface
import org.matrix.rustcomponents.sdk.SyncServiceState import org.matrix.rustcomponents.sdk.SyncServiceState
import timber.log.Timber import timber.log.Timber
class RustSyncService( class RustSyncService(
private val innerSyncService: SyncServiceInterface, private val innerSyncService: SyncServiceInterface,
private val dispatchers: CoroutineDispatchers,
sessionCoroutineScope: CoroutineScope sessionCoroutineScope: CoroutineScope
) : SyncService { ) : SyncService {
override suspend fun startSync() = runCatching { override suspend fun startSync() = withContext(dispatchers.io) {
Timber.i("Start sync") runCatching {
innerSyncService.start() Timber.i("Start sync")
}.onFailure { innerSyncService.startBlocking()
Timber.d("Start sync failed: $it") }.onFailure {
Timber.d("Start sync failed: $it")
}
} }
override suspend fun stopSync() = runCatching { override suspend fun stopSync() = withContext(dispatchers.io){
Timber.i("Stop sync") runCatching {
innerSyncService.stop() Timber.i("Stop sync")
}.onFailure { innerSyncService.stopBlocking()
Timber.d("Stop sync failed: $it") }.onFailure {
Timber.d("Stop sync failed: $it")
}
} }
override val syncState: StateFlow<SyncState> = override val syncState: StateFlow<SyncState> =

View file

@ -44,7 +44,7 @@ internal fun Room.timelineDiffFlow(onInitialList: suspend (List<TimelineItem>) -
} }
val roomId = id() val roomId = id()
Timber.d("Open timelineDiffFlow for room $roomId") Timber.d("Open timelineDiffFlow for room $roomId")
val result = addTimelineListener(listener) val result = addTimelineListenerBlocking(listener)
try { try {
onInitialList(result.items) onInitialList(result.items)
} catch (exception: Exception) { } catch (exception: Exception) {

View file

@ -124,7 +124,7 @@ class RustMatrixTimeline(
private suspend fun fetchMembers() = withContext(dispatcher) { private suspend fun fetchMembers() = withContext(dispatcher) {
initLatch.await() initLatch.await()
innerRoom.fetchMembers() innerRoom.fetchMembersBlocking()
} }
@OptIn(ExperimentalCoroutinesApi::class) @OptIn(ExperimentalCoroutinesApi::class)

View file

@ -16,6 +16,7 @@
package io.element.android.libraries.matrix.impl.verification package io.element.android.libraries.matrix.impl.verification
import io.element.android.libraries.core.coroutine.CoroutineDispatchers
import io.element.android.libraries.core.data.tryOrNull import io.element.android.libraries.core.data.tryOrNull
import io.element.android.libraries.matrix.api.sync.SyncState import io.element.android.libraries.matrix.api.sync.SyncState
import io.element.android.libraries.matrix.api.verification.SessionVerificationService import io.element.android.libraries.matrix.api.verification.SessionVerificationService
@ -27,6 +28,7 @@ import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.combine import kotlinx.coroutines.flow.combine
import kotlinx.coroutines.withContext
import org.matrix.rustcomponents.sdk.SessionVerificationController import org.matrix.rustcomponents.sdk.SessionVerificationController
import org.matrix.rustcomponents.sdk.SessionVerificationControllerDelegate import org.matrix.rustcomponents.sdk.SessionVerificationControllerDelegate
import org.matrix.rustcomponents.sdk.SessionVerificationControllerInterface import org.matrix.rustcomponents.sdk.SessionVerificationControllerInterface
@ -35,6 +37,7 @@ import javax.inject.Inject
class RustSessionVerificationService @Inject constructor( class RustSessionVerificationService @Inject constructor(
private val syncService: RustSyncService, private val syncService: RustSyncService,
private val dispatchers: CoroutineDispatchers,
) : SessionVerificationService, SessionVerificationControllerDelegate { ) : SessionVerificationService, SessionVerificationControllerDelegate {
var verificationController: SessionVerificationControllerInterface? = null var verificationController: SessionVerificationControllerInterface? = null
@ -61,21 +64,31 @@ class RustSessionVerificationService @Inject constructor(
syncState == SyncState.Running && verificationStatus == SessionVerifiedStatus.NotVerified syncState == SyncState.Running && verificationStatus == SessionVerifiedStatus.NotVerified
} }
override suspend fun requestVerification() = tryOrFail { override suspend fun requestVerification() {
verificationController?.requestVerification() tryOrFail {
verificationController?.requestVerificationBlocking()
}
} }
override suspend fun cancelVerification() = tryOrFail { verificationController?.cancelVerification() } override suspend fun cancelVerification() {
tryOrFail { verificationController?.cancelVerificationBlocking() }
override suspend fun approveVerification() = tryOrFail { verificationController?.approveVerification() }
override suspend fun declineVerification() = tryOrFail { verificationController?.declineVerification() }
override suspend fun startVerification() = tryOrFail {
verificationController?.startSasVerification()
} }
private suspend fun tryOrFail(block: suspend () -> Unit) { override suspend fun approveVerification() {
tryOrFail { verificationController?.approveVerificationBlocking() }
}
override suspend fun declineVerification() {
tryOrFail { verificationController?.declineVerificationBlocking() }
}
override suspend fun startVerification() {
tryOrFail {
verificationController?.startSasVerificationBlocking()
}
}
private suspend fun tryOrFail(block: suspend () -> Unit) = withContext(dispatchers.io) {
runCatching { runCatching {
block() block()
}.onFailure { didFail() } }.onFailure { didFail() }