Coroutine: introduce scoped dispatcher with limitedParalellism
This commit is contained in:
parent
c2c81d3747
commit
d02f7fb871
5 changed files with 68 additions and 55 deletions
|
|
@ -54,7 +54,7 @@ import io.element.android.libraries.matrix.impl.verification.RustSessionVerifica
|
||||||
import io.element.android.libraries.sessionstorage.api.SessionStore
|
import io.element.android.libraries.sessionstorage.api.SessionStore
|
||||||
import io.element.android.services.toolbox.api.systemclock.SystemClock
|
import io.element.android.services.toolbox.api.systemclock.SystemClock
|
||||||
import kotlinx.coroutines.CoroutineScope
|
import kotlinx.coroutines.CoroutineScope
|
||||||
import kotlinx.coroutines.Dispatchers
|
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||||
import kotlinx.coroutines.cancel
|
import kotlinx.coroutines.cancel
|
||||||
import kotlinx.coroutines.flow.filter
|
import kotlinx.coroutines.flow.filter
|
||||||
import kotlinx.coroutines.flow.first
|
import kotlinx.coroutines.flow.first
|
||||||
|
|
@ -73,6 +73,7 @@ import org.matrix.rustcomponents.sdk.CreateRoomParameters as RustCreateRoomParam
|
||||||
import org.matrix.rustcomponents.sdk.RoomPreset as RustRoomPreset
|
import org.matrix.rustcomponents.sdk.RoomPreset as RustRoomPreset
|
||||||
import org.matrix.rustcomponents.sdk.RoomVisibility as RustRoomVisibility
|
import org.matrix.rustcomponents.sdk.RoomVisibility as RustRoomVisibility
|
||||||
|
|
||||||
|
@OptIn(ExperimentalCoroutinesApi::class)
|
||||||
class RustMatrixClient constructor(
|
class RustMatrixClient constructor(
|
||||||
private val client: Client,
|
private val client: Client,
|
||||||
private val sessionStore: SessionStore,
|
private val sessionStore: SessionStore,
|
||||||
|
|
@ -85,6 +86,7 @@ class RustMatrixClient constructor(
|
||||||
|
|
||||||
override val sessionId: UserId = UserId(client.userId())
|
override val sessionId: UserId = UserId(client.userId())
|
||||||
private val roomListService = client.roomListServiceWithEncryption()
|
private val roomListService = client.roomListServiceWithEncryption()
|
||||||
|
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 verificationService = RustSessionVerificationService()
|
private val verificationService = RustSessionVerificationService()
|
||||||
private val syncService = RustSyncService(roomListService, sessionCoroutineScope)
|
private val syncService = RustSyncService(roomListService, sessionCoroutineScope)
|
||||||
|
|
@ -92,6 +94,7 @@ class RustMatrixClient constructor(
|
||||||
client = client,
|
client = client,
|
||||||
dispatchers = dispatchers,
|
dispatchers = dispatchers,
|
||||||
)
|
)
|
||||||
|
|
||||||
private val notificationService = RustNotificationService(client)
|
private val notificationService = RustNotificationService(client)
|
||||||
|
|
||||||
private val clientDelegate = object : ClientDelegate {
|
private val clientDelegate = object : ClientDelegate {
|
||||||
|
|
@ -105,7 +108,7 @@ class RustMatrixClient constructor(
|
||||||
RustRoomSummaryDataSource(
|
RustRoomSummaryDataSource(
|
||||||
roomListService = roomListService,
|
roomListService = roomListService,
|
||||||
sessionCoroutineScope = sessionCoroutineScope,
|
sessionCoroutineScope = sessionCoroutineScope,
|
||||||
coroutineDispatchers = dispatchers,
|
dispatcher = sessionDispatcher,
|
||||||
)
|
)
|
||||||
|
|
||||||
override val roomSummaryDataSource: RoomSummaryDataSource
|
override val roomSummaryDataSource: RoomSummaryDataSource
|
||||||
|
|
@ -150,7 +153,7 @@ class RustMatrixClient constructor(
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
private suspend fun pairOfRoom(roomId: RoomId): Pair<RoomListItem, Room>? = withContext(dispatchers.io) {
|
private suspend fun pairOfRoom(roomId: RoomId): Pair<RoomListItem, Room>? = withContext(sessionDispatcher) {
|
||||||
val cachedRoomListItem = roomListService.roomOrNull(roomId.value)
|
val cachedRoomListItem = roomListService.roomOrNull(roomId.value)
|
||||||
val fullRoom = cachedRoomListItem?.fullRoom()
|
val fullRoom = cachedRoomListItem?.fullRoom()
|
||||||
if (cachedRoomListItem == null || fullRoom == null) {
|
if (cachedRoomListItem == null || fullRoom == null) {
|
||||||
|
|
@ -165,19 +168,19 @@ class RustMatrixClient constructor(
|
||||||
return roomId?.let { getRoom(it) }
|
return roomId?.let { getRoom(it) }
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun ignoreUser(userId: UserId): Result<Unit> = withContext(dispatchers.io) {
|
override suspend fun ignoreUser(userId: UserId): Result<Unit> = withContext(sessionDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
client.ignoreUser(userId.value)
|
client.ignoreUser(userId.value)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun unignoreUser(userId: UserId): Result<Unit> = withContext(dispatchers.io) {
|
override suspend fun unignoreUser(userId: UserId): Result<Unit> = withContext(sessionDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
client.unignoreUser(userId.value)
|
client.unignoreUser(userId.value)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun createRoom(createRoomParams: CreateRoomParameters): Result<RoomId> = withContext(dispatchers.io) {
|
override suspend fun createRoom(createRoomParams: CreateRoomParameters): Result<RoomId> = withContext(sessionDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
val rustParams = RustCreateRoomParameters(
|
val rustParams = RustCreateRoomParameters(
|
||||||
name = createRoomParams.name,
|
name = createRoomParams.name,
|
||||||
|
|
@ -221,14 +224,14 @@ class RustMatrixClient constructor(
|
||||||
return createRoom(createRoomParams)
|
return createRoom(createRoomParams)
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun getProfile(userId: UserId): Result<MatrixUser> = withContext(Dispatchers.IO) {
|
override suspend fun getProfile(userId: UserId): Result<MatrixUser> = withContext(sessionDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
client.getProfile(userId.value).let(UserProfileMapper::map)
|
client.getProfile(userId.value).let(UserProfileMapper::map)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun searchUsers(searchTerm: String, limit: Long): Result<MatrixSearchUserResults> =
|
override suspend fun searchUsers(searchTerm: String, limit: Long): Result<MatrixSearchUserResults> =
|
||||||
withContext(dispatchers.io) {
|
withContext(sessionDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
client.searchUsers(searchTerm, limit.toULong()).let(UserSearchResultMapper::map)
|
client.searchUsers(searchTerm, limit.toULong()).let(UserSearchResultMapper::map)
|
||||||
}
|
}
|
||||||
|
|
@ -260,7 +263,7 @@ class RustMatrixClient constructor(
|
||||||
baseDirectory.deleteSessionDirectory(userID = sessionId.value, deleteCryptoDb = false)
|
baseDirectory.deleteSessionDirectory(userID = sessionId.value, deleteCryptoDb = false)
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun logout() = withContext(dispatchers.io) {
|
override suspend fun logout() = withContext(sessionDispatcher) {
|
||||||
try {
|
try {
|
||||||
client.logout()
|
client.logout()
|
||||||
} catch (failure: Throwable) {
|
} catch (failure: Throwable) {
|
||||||
|
|
@ -271,20 +274,20 @@ class RustMatrixClient constructor(
|
||||||
sessionStore.removeSession(sessionId.value)
|
sessionStore.removeSession(sessionId.value)
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun loadUserDisplayName(): Result<String> = withContext(dispatchers.io) {
|
override suspend fun loadUserDisplayName(): Result<String> = withContext(sessionDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
client.displayName()
|
client.displayName()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun loadUserAvatarURLString(): Result<String?> = withContext(dispatchers.io) {
|
override suspend fun loadUserAvatarURLString(): Result<String?> = withContext(sessionDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
client.avatarUrl()
|
client.avatarUrl()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@OptIn(ExperimentalUnsignedTypes::class)
|
@OptIn(ExperimentalUnsignedTypes::class)
|
||||||
override suspend fun uploadMedia(mimeType: String, data: ByteArray, progressCallback: ProgressCallback?): Result<String> = withContext(dispatchers.io) {
|
override suspend fun uploadMedia(mimeType: String, data: ByteArray, progressCallback: ProgressCallback?): Result<String> = withContext(sessionDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
client.uploadMedia(mimeType, data.toUByteArray().toList(), progressCallback?.toProgressWatcher())
|
client.uploadMedia(mimeType, data.toUByteArray().toList(), progressCallback?.toProgressWatcher())
|
||||||
}
|
}
|
||||||
|
|
@ -305,7 +308,7 @@ class RustMatrixClient constructor(
|
||||||
private suspend fun File.getCacheSize(
|
private suspend fun File.getCacheSize(
|
||||||
userID: String,
|
userID: String,
|
||||||
includeCryptoDb: Boolean = false,
|
includeCryptoDb: Boolean = false,
|
||||||
): Long = withContext(dispatchers.io) {
|
): Long = withContext(sessionDispatcher) {
|
||||||
// Rust sanitises the user ID replacing invalid characters with an _
|
// Rust sanitises the user ID replacing invalid characters with an _
|
||||||
val sanitisedUserID = userID.replace(":", "_")
|
val sanitisedUserID = userID.replace(":", "_")
|
||||||
val sessionDirectory = File(this@getCacheSize, sanitisedUserID)
|
val sessionDirectory = File(this@getCacheSize, sanitisedUserID)
|
||||||
|
|
@ -327,7 +330,7 @@ class RustMatrixClient constructor(
|
||||||
private suspend fun File.deleteSessionDirectory(
|
private suspend fun File.deleteSessionDirectory(
|
||||||
userID: String,
|
userID: String,
|
||||||
deleteCryptoDb: Boolean = false,
|
deleteCryptoDb: Boolean = false,
|
||||||
): Boolean = withContext(dispatchers.io) {
|
): Boolean = withContext(sessionDispatcher) {
|
||||||
// Rust sanitises the user ID replacing invalid characters with an _
|
// Rust sanitises the user ID replacing invalid characters with an _
|
||||||
val sanitisedUserID = userID.replace(":", "_")
|
val sanitisedUserID = userID.replace(":", "_")
|
||||||
val sessionDirectory = File(this@deleteSessionDirectory, sanitisedUserID)
|
val sessionDirectory = File(this@deleteSessionDirectory, sanitisedUserID)
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,7 @@ import io.element.android.libraries.core.coroutine.CoroutineDispatchers
|
||||||
import io.element.android.libraries.matrix.api.media.MatrixMediaLoader
|
import io.element.android.libraries.matrix.api.media.MatrixMediaLoader
|
||||||
import io.element.android.libraries.matrix.api.media.MediaFile
|
import io.element.android.libraries.matrix.api.media.MediaFile
|
||||||
import io.element.android.libraries.matrix.api.media.MediaSource
|
import io.element.android.libraries.matrix.api.media.MediaSource
|
||||||
|
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||||
import kotlinx.coroutines.withContext
|
import kotlinx.coroutines.withContext
|
||||||
import org.matrix.rustcomponents.sdk.Client
|
import org.matrix.rustcomponents.sdk.Client
|
||||||
import org.matrix.rustcomponents.sdk.mediaSourceFromUrl
|
import org.matrix.rustcomponents.sdk.mediaSourceFromUrl
|
||||||
|
|
@ -29,10 +30,12 @@ import org.matrix.rustcomponents.sdk.MediaSource as RustMediaSource
|
||||||
|
|
||||||
class RustMediaLoader(
|
class RustMediaLoader(
|
||||||
baseCacheDirectory: File,
|
baseCacheDirectory: File,
|
||||||
private val dispatchers: CoroutineDispatchers,
|
dispatchers: CoroutineDispatchers,
|
||||||
private val innerClient: Client,
|
private val innerClient: Client,
|
||||||
) : MatrixMediaLoader {
|
) : MatrixMediaLoader {
|
||||||
|
|
||||||
|
@OptIn(ExperimentalCoroutinesApi::class)
|
||||||
|
private val mediaDispatcher = dispatchers.io.limitedParallelism(32)
|
||||||
private val cacheDirectory = File(baseCacheDirectory, "temp/media").apply {
|
private val cacheDirectory = File(baseCacheDirectory, "temp/media").apply {
|
||||||
if (!exists()) {
|
if (!exists()) {
|
||||||
mkdirs()
|
mkdirs()
|
||||||
|
|
@ -41,7 +44,7 @@ class RustMediaLoader(
|
||||||
|
|
||||||
@OptIn(ExperimentalUnsignedTypes::class)
|
@OptIn(ExperimentalUnsignedTypes::class)
|
||||||
override suspend fun loadMediaContent(source: MediaSource): Result<ByteArray> =
|
override suspend fun loadMediaContent(source: MediaSource): Result<ByteArray> =
|
||||||
withContext(dispatchers.io) {
|
withContext(mediaDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
source.toRustMediaSource().use { source ->
|
source.toRustMediaSource().use { source ->
|
||||||
innerClient.getMediaContent(source).toUByteArray().toByteArray()
|
innerClient.getMediaContent(source).toUByteArray().toByteArray()
|
||||||
|
|
@ -55,7 +58,7 @@ class RustMediaLoader(
|
||||||
width: Long,
|
width: Long,
|
||||||
height: Long
|
height: Long
|
||||||
): Result<ByteArray> =
|
): Result<ByteArray> =
|
||||||
withContext(dispatchers.io) {
|
withContext(mediaDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
source.toRustMediaSource().use { mediaSource ->
|
source.toRustMediaSource().use { mediaSource ->
|
||||||
innerClient.getMediaThumbnail(
|
innerClient.getMediaThumbnail(
|
||||||
|
|
@ -68,7 +71,7 @@ class RustMediaLoader(
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun downloadMediaFile(source: MediaSource, mimeType: String?, body: String?): Result<MediaFile> =
|
override suspend fun downloadMediaFile(source: MediaSource, mimeType: String?, body: String?): Result<MediaFile> =
|
||||||
withContext(dispatchers.io) {
|
withContext(mediaDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
source.toRustMediaSource().use { mediaSource ->
|
source.toRustMediaSource().use { mediaSource ->
|
||||||
val mediaFile = innerClient.getMediaFile(
|
val mediaFile = innerClient.getMediaFile(
|
||||||
|
|
|
||||||
|
|
@ -23,7 +23,6 @@ import io.element.android.libraries.matrix.api.core.ProgressCallback
|
||||||
import io.element.android.libraries.matrix.api.core.RoomId
|
import io.element.android.libraries.matrix.api.core.RoomId
|
||||||
import io.element.android.libraries.matrix.api.core.SessionId
|
import io.element.android.libraries.matrix.api.core.SessionId
|
||||||
import io.element.android.libraries.matrix.api.core.UserId
|
import io.element.android.libraries.matrix.api.core.UserId
|
||||||
import io.element.android.libraries.matrix.api.room.location.AssetType
|
|
||||||
import io.element.android.libraries.matrix.api.media.AudioInfo
|
import io.element.android.libraries.matrix.api.media.AudioInfo
|
||||||
import io.element.android.libraries.matrix.api.media.FileInfo
|
import io.element.android.libraries.matrix.api.media.FileInfo
|
||||||
import io.element.android.libraries.matrix.api.media.ImageInfo
|
import io.element.android.libraries.matrix.api.media.ImageInfo
|
||||||
|
|
@ -32,17 +31,19 @@ import io.element.android.libraries.matrix.api.room.MatrixRoom
|
||||||
import io.element.android.libraries.matrix.api.room.MatrixRoomMembersState
|
import io.element.android.libraries.matrix.api.room.MatrixRoomMembersState
|
||||||
import io.element.android.libraries.matrix.api.room.MessageEventType
|
import io.element.android.libraries.matrix.api.room.MessageEventType
|
||||||
import io.element.android.libraries.matrix.api.room.StateEventType
|
import io.element.android.libraries.matrix.api.room.StateEventType
|
||||||
|
import io.element.android.libraries.matrix.api.room.location.AssetType
|
||||||
import io.element.android.libraries.matrix.api.room.roomMembers
|
import io.element.android.libraries.matrix.api.room.roomMembers
|
||||||
import io.element.android.libraries.matrix.api.timeline.MatrixTimeline
|
import io.element.android.libraries.matrix.api.timeline.MatrixTimeline
|
||||||
import io.element.android.libraries.matrix.api.timeline.item.event.EventType
|
import io.element.android.libraries.matrix.api.timeline.item.event.EventType
|
||||||
import io.element.android.libraries.matrix.impl.core.toProgressWatcher
|
import io.element.android.libraries.matrix.impl.core.toProgressWatcher
|
||||||
import io.element.android.libraries.matrix.impl.room.location.toInner
|
|
||||||
import io.element.android.libraries.matrix.impl.media.map
|
import io.element.android.libraries.matrix.impl.media.map
|
||||||
|
import io.element.android.libraries.matrix.impl.room.location.toInner
|
||||||
import io.element.android.libraries.matrix.impl.timeline.RustMatrixTimeline
|
import io.element.android.libraries.matrix.impl.timeline.RustMatrixTimeline
|
||||||
import io.element.android.libraries.matrix.impl.timeline.backPaginationStatusFlow
|
import io.element.android.libraries.matrix.impl.timeline.backPaginationStatusFlow
|
||||||
import io.element.android.libraries.matrix.impl.timeline.timelineDiffFlow
|
import io.element.android.libraries.matrix.impl.timeline.timelineDiffFlow
|
||||||
import io.element.android.services.toolbox.api.systemclock.SystemClock
|
import io.element.android.services.toolbox.api.systemclock.SystemClock
|
||||||
import kotlinx.coroutines.CoroutineScope
|
import kotlinx.coroutines.CoroutineScope
|
||||||
|
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||||
import kotlinx.coroutines.cancel
|
import kotlinx.coroutines.cancel
|
||||||
import kotlinx.coroutines.flow.MutableStateFlow
|
import kotlinx.coroutines.flow.MutableStateFlow
|
||||||
import kotlinx.coroutines.flow.StateFlow
|
import kotlinx.coroutines.flow.StateFlow
|
||||||
|
|
@ -62,6 +63,7 @@ import org.matrix.rustcomponents.sdk.messageEventContentFromMarkdown
|
||||||
import timber.log.Timber
|
import timber.log.Timber
|
||||||
import java.io.File
|
import java.io.File
|
||||||
|
|
||||||
|
@OptIn(ExperimentalCoroutinesApi::class)
|
||||||
class RustMatrixRoom(
|
class RustMatrixRoom(
|
||||||
override val sessionId: SessionId,
|
override val sessionId: SessionId,
|
||||||
private val roomListItem: RoomListItem,
|
private val roomListItem: RoomListItem,
|
||||||
|
|
@ -74,6 +76,11 @@ class RustMatrixRoom(
|
||||||
|
|
||||||
override val roomId = RoomId(innerRoom.id())
|
override val roomId = RoomId(innerRoom.id())
|
||||||
|
|
||||||
|
// Create a dispatcher for all room methods...
|
||||||
|
private val roomDispatcher = coroutineDispatchers.io.limitedParallelism(32)
|
||||||
|
//...except getMember methods as it could quickly fill the roomDispatcher...
|
||||||
|
private val roomMembersDispatcher = coroutineDispatchers.io.limitedParallelism(8)
|
||||||
|
|
||||||
private val roomCoroutineScope = sessionCoroutineScope.childScope(coroutineDispatchers.main, "RoomScope-$roomId")
|
private val roomCoroutineScope = sessionCoroutineScope.childScope(coroutineDispatchers.main, "RoomScope-$roomId")
|
||||||
private val _membersStateFlow = MutableStateFlow<MatrixRoomMembersState>(MatrixRoomMembersState.Unknown)
|
private val _membersStateFlow = MutableStateFlow<MatrixRoomMembersState>(MatrixRoomMembersState.Unknown)
|
||||||
private val isInit = MutableStateFlow(false)
|
private val isInit = MutableStateFlow(false)
|
||||||
|
|
@ -83,7 +90,7 @@ class RustMatrixRoom(
|
||||||
matrixRoom = this,
|
matrixRoom = this,
|
||||||
innerRoom = innerRoom,
|
innerRoom = innerRoom,
|
||||||
roomCoroutineScope = roomCoroutineScope,
|
roomCoroutineScope = roomCoroutineScope,
|
||||||
coroutineDispatchers = coroutineDispatchers
|
dispatcher = roomDispatcher
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -105,7 +112,7 @@ class RustMatrixRoom(
|
||||||
timelineLimit = null
|
timelineLimit = null
|
||||||
)
|
)
|
||||||
roomListItem.subscribe(settings)
|
roomListItem.subscribe(settings)
|
||||||
roomCoroutineScope.launch(coroutineDispatchers.computation) {
|
roomCoroutineScope.launch(roomDispatcher) {
|
||||||
innerRoom.timelineDiffFlow { initialList ->
|
innerRoom.timelineDiffFlow { initialList ->
|
||||||
_timeline.postItems(initialList)
|
_timeline.postItems(initialList)
|
||||||
}.onEach {
|
}.onEach {
|
||||||
|
|
@ -175,7 +182,7 @@ class RustMatrixRoom(
|
||||||
override val activeMemberCount: Long
|
override val activeMemberCount: Long
|
||||||
get() = innerRoom.activeMembersCount().toLong()
|
get() = innerRoom.activeMembersCount().toLong()
|
||||||
|
|
||||||
override suspend fun updateMembers(): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun updateMembers(): Result<Unit> = withContext(roomMembersDispatcher) {
|
||||||
val currentState = _membersStateFlow.value
|
val currentState = _membersStateFlow.value
|
||||||
val currentMembers = currentState.roomMembers()
|
val currentMembers = currentState.roomMembers()
|
||||||
_membersStateFlow.value = MatrixRoomMembersState.Pending(prevRoomMembers = currentMembers)
|
_membersStateFlow.value = MatrixRoomMembersState.Pending(prevRoomMembers = currentMembers)
|
||||||
|
|
@ -189,20 +196,20 @@ class RustMatrixRoom(
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun userDisplayName(userId: UserId): Result<String?> =
|
override suspend fun userDisplayName(userId: UserId): Result<String?> =
|
||||||
withContext(coroutineDispatchers.io) {
|
withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.memberDisplayName(userId.value)
|
innerRoom.memberDisplayName(userId.value)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun userAvatarUrl(userId: UserId): Result<String?> =
|
override suspend fun userAvatarUrl(userId: UserId): Result<String?> =
|
||||||
withContext(coroutineDispatchers.io) {
|
withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.memberAvatarUrl(userId.value)
|
innerRoom.memberAvatarUrl(userId.value)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun sendMessage(message: String): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun sendMessage(message: String): Result<Unit> = withContext(roomDispatcher) {
|
||||||
val transactionId = genTransactionId()
|
val transactionId = genTransactionId()
|
||||||
messageEventContentFromMarkdown(message).use { content ->
|
messageEventContentFromMarkdown(message).use { content ->
|
||||||
runCatching {
|
runCatching {
|
||||||
|
|
@ -211,7 +218,7 @@ class RustMatrixRoom(
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun editMessage(originalEventId: EventId?, transactionId: String?, message: String): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun editMessage(originalEventId: EventId?, transactionId: String?, message: String): Result<Unit> = withContext(roomDispatcher) {
|
||||||
if (originalEventId != null) {
|
if (originalEventId != null) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.edit(/* TODO use content */ message, originalEventId.value, transactionId)
|
innerRoom.edit(/* TODO use content */ message, originalEventId.value, transactionId)
|
||||||
|
|
@ -224,7 +231,7 @@ class RustMatrixRoom(
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun replyMessage(eventId: EventId, message: String): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun replyMessage(eventId: EventId, message: String): Result<Unit> = withContext(roomDispatcher) {
|
||||||
val transactionId = genTransactionId()
|
val transactionId = genTransactionId()
|
||||||
// val content = messageEventContentFromMarkdown(message)
|
// val content = messageEventContentFromMarkdown(message)
|
||||||
runCatching {
|
runCatching {
|
||||||
|
|
@ -232,50 +239,50 @@ class RustMatrixRoom(
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun redactEvent(eventId: EventId, reason: String?) = withContext(coroutineDispatchers.io) {
|
override suspend fun redactEvent(eventId: EventId, reason: String?) = withContext(roomDispatcher) {
|
||||||
val transactionId = genTransactionId()
|
val transactionId = genTransactionId()
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.redact(eventId.value, reason, transactionId)
|
innerRoom.redact(eventId.value, reason, transactionId)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun leave(): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun leave(): Result<Unit> = withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.leave()
|
innerRoom.leave()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun acceptInvitation(): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun acceptInvitation(): Result<Unit> = withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.acceptInvitation()
|
innerRoom.acceptInvitation()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun rejectInvitation(): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun rejectInvitation(): Result<Unit> = withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.rejectInvitation()
|
innerRoom.rejectInvitation()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun inviteUserById(id: UserId): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun inviteUserById(id: UserId): Result<Unit> = withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.inviteUserById(id.value)
|
innerRoom.inviteUserById(id.value)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun canInvite(): Result<Boolean> = withContext(coroutineDispatchers.io) {
|
override suspend fun canInvite(): Result<Boolean> = withContext(roomMembersDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.member(sessionId.value).use(RoomMember::canInvite)
|
innerRoom.member(sessionId.value).use(RoomMember::canInvite)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun canSendStateEvent(type: StateEventType): Result<Boolean> = withContext(coroutineDispatchers.io) {
|
override suspend fun canSendStateEvent(type: StateEventType): Result<Boolean> = withContext(roomMembersDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.member(sessionId.value).use { it.canSendState(type.map()) }
|
innerRoom.member(sessionId.value).use { it.canSendState(type.map()) }
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun canSendEvent(type: MessageEventType): Result<Boolean> = withContext(coroutineDispatchers.io) {
|
override suspend fun canSendEvent(type: MessageEventType): Result<Boolean> = withContext(roomMembersDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.member(sessionId.value).use { it.canSendMessage(type.map()) }
|
innerRoom.member(sessionId.value).use { it.canSendMessage(type.map()) }
|
||||||
}
|
}
|
||||||
|
|
@ -305,13 +312,13 @@ class RustMatrixRoom(
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun toggleReaction(emoji: String, eventId: EventId): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun toggleReaction(emoji: String, eventId: EventId): Result<Unit> = withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.toggleReaction(key = emoji, eventId = eventId.value)
|
innerRoom.toggleReaction(key = emoji, eventId = eventId.value)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun forwardEvent(eventId: EventId, roomIds: List<RoomId>): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun forwardEvent(eventId: EventId, roomIds: List<RoomId>): Result<Unit> = withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
roomContentForwarder.forward(fromRoom = innerRoom, eventId = eventId, toRoomIds = roomIds)
|
roomContentForwarder.forward(fromRoom = innerRoom, eventId = eventId, toRoomIds = roomIds)
|
||||||
}.onFailure {
|
}.onFailure {
|
||||||
|
|
@ -320,14 +327,14 @@ class RustMatrixRoom(
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun retrySendMessage(transactionId: String): Result<Unit> =
|
override suspend fun retrySendMessage(transactionId: String): Result<Unit> =
|
||||||
withContext(coroutineDispatchers.io) {
|
withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.retrySend(transactionId)
|
innerRoom.retrySend(transactionId)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun cancelSend(transactionId: String): Result<Unit> =
|
override suspend fun cancelSend(transactionId: String): Result<Unit> =
|
||||||
withContext(coroutineDispatchers.io) {
|
withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.cancelSend(transactionId)
|
innerRoom.cancelSend(transactionId)
|
||||||
}
|
}
|
||||||
|
|
@ -335,40 +342,40 @@ class RustMatrixRoom(
|
||||||
|
|
||||||
@OptIn(ExperimentalUnsignedTypes::class)
|
@OptIn(ExperimentalUnsignedTypes::class)
|
||||||
override suspend fun updateAvatar(mimeType: String, data: ByteArray): Result<Unit> =
|
override suspend fun updateAvatar(mimeType: String, data: ByteArray): Result<Unit> =
|
||||||
withContext(coroutineDispatchers.io) {
|
withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.uploadAvatar(mimeType, data.toUByteArray().toList())
|
innerRoom.uploadAvatar(mimeType, data.toUByteArray().toList())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun removeAvatar(): Result<Unit> =
|
override suspend fun removeAvatar(): Result<Unit> =
|
||||||
withContext(coroutineDispatchers.io) {
|
withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.removeAvatar()
|
innerRoom.removeAvatar()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun setName(name: String): Result<Unit> =
|
override suspend fun setName(name: String): Result<Unit> =
|
||||||
withContext(coroutineDispatchers.io) {
|
withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.setName(name)
|
innerRoom.setName(name)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun setTopic(topic: String): Result<Unit> =
|
override suspend fun setTopic(topic: String): Result<Unit> =
|
||||||
withContext(coroutineDispatchers.io) {
|
withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.setTopic(topic)
|
innerRoom.setTopic(topic)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private suspend fun fetchMembers() = withContext(coroutineDispatchers.io) {
|
private suspend fun fetchMembers() = withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.fetchMembers()
|
innerRoom.fetchMembers()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun reportContent(eventId: EventId, reason: String, blockUserId: UserId?): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun reportContent(eventId: EventId, reason: String, blockUserId: UserId?): Result<Unit> = withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.reportContent(eventId = eventId.value, score = null, reason = reason)
|
innerRoom.reportContent(eventId = eventId.value, score = null, reason = reason)
|
||||||
if (blockUserId != null) {
|
if (blockUserId != null) {
|
||||||
|
|
@ -383,7 +390,7 @@ class RustMatrixRoom(
|
||||||
description: String?,
|
description: String?,
|
||||||
zoomLevel: Int?,
|
zoomLevel: Int?,
|
||||||
assetType: AssetType?,
|
assetType: AssetType?,
|
||||||
): Result<Unit> = withContext(coroutineDispatchers.io) {
|
): Result<Unit> = withContext(roomDispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.sendLocation(
|
innerRoom.sendLocation(
|
||||||
body = body,
|
body = body,
|
||||||
|
|
|
||||||
|
|
@ -16,9 +16,9 @@
|
||||||
|
|
||||||
package io.element.android.libraries.matrix.impl.room
|
package io.element.android.libraries.matrix.impl.room
|
||||||
|
|
||||||
import io.element.android.libraries.core.coroutine.CoroutineDispatchers
|
|
||||||
import io.element.android.libraries.matrix.api.room.RoomSummary
|
import io.element.android.libraries.matrix.api.room.RoomSummary
|
||||||
import io.element.android.libraries.matrix.api.room.RoomSummaryDataSource
|
import io.element.android.libraries.matrix.api.room.RoomSummaryDataSource
|
||||||
|
import kotlinx.coroutines.CoroutineDispatcher
|
||||||
import kotlinx.coroutines.CoroutineScope
|
import kotlinx.coroutines.CoroutineScope
|
||||||
import kotlinx.coroutines.flow.Flow
|
import kotlinx.coroutines.flow.Flow
|
||||||
import kotlinx.coroutines.flow.MutableStateFlow
|
import kotlinx.coroutines.flow.MutableStateFlow
|
||||||
|
|
@ -41,7 +41,7 @@ import timber.log.Timber
|
||||||
internal class RustRoomSummaryDataSource(
|
internal class RustRoomSummaryDataSource(
|
||||||
private val roomListService: RoomListService,
|
private val roomListService: RoomListService,
|
||||||
private val sessionCoroutineScope: CoroutineScope,
|
private val sessionCoroutineScope: CoroutineScope,
|
||||||
coroutineDispatchers: CoroutineDispatchers,
|
dispatcher: CoroutineDispatcher,
|
||||||
roomSummaryDetailsFactory: RoomSummaryDetailsFactory = RoomSummaryDetailsFactory(),
|
roomSummaryDetailsFactory: RoomSummaryDetailsFactory = RoomSummaryDetailsFactory(),
|
||||||
) : RoomSummaryDataSource {
|
) : RoomSummaryDataSource {
|
||||||
|
|
||||||
|
|
@ -53,7 +53,7 @@ internal class RustRoomSummaryDataSource(
|
||||||
private val inviteRoomsListProcessor = RoomSummaryListProcessor(inviteRooms, roomListService, roomSummaryDetailsFactory, shouldFetchFullRoom = true)
|
private val inviteRoomsListProcessor = RoomSummaryListProcessor(inviteRooms, roomListService, roomSummaryDetailsFactory, shouldFetchFullRoom = true)
|
||||||
|
|
||||||
init {
|
init {
|
||||||
sessionCoroutineScope.launch(coroutineDispatchers.computation) {
|
sessionCoroutineScope.launch(dispatcher) {
|
||||||
val allRooms = roomListService.allRooms()
|
val allRooms = roomListService.allRooms()
|
||||||
allRooms
|
allRooms
|
||||||
.observeEntriesWithProcessor(allRoomsListProcessor)
|
.observeEntriesWithProcessor(allRoomsListProcessor)
|
||||||
|
|
|
||||||
|
|
@ -16,7 +16,6 @@
|
||||||
|
|
||||||
package io.element.android.libraries.matrix.impl.timeline
|
package io.element.android.libraries.matrix.impl.timeline
|
||||||
|
|
||||||
import io.element.android.libraries.core.coroutine.CoroutineDispatchers
|
|
||||||
import io.element.android.libraries.matrix.api.core.EventId
|
import io.element.android.libraries.matrix.api.core.EventId
|
||||||
import io.element.android.libraries.matrix.api.room.MatrixRoom
|
import io.element.android.libraries.matrix.api.room.MatrixRoom
|
||||||
import io.element.android.libraries.matrix.api.timeline.MatrixTimeline
|
import io.element.android.libraries.matrix.api.timeline.MatrixTimeline
|
||||||
|
|
@ -25,6 +24,7 @@ import io.element.android.libraries.matrix.impl.timeline.item.event.EventMessage
|
||||||
import io.element.android.libraries.matrix.impl.timeline.item.event.EventTimelineItemMapper
|
import io.element.android.libraries.matrix.impl.timeline.item.event.EventTimelineItemMapper
|
||||||
import io.element.android.libraries.matrix.impl.timeline.item.event.TimelineEventContentMapper
|
import io.element.android.libraries.matrix.impl.timeline.item.event.TimelineEventContentMapper
|
||||||
import io.element.android.libraries.matrix.impl.timeline.item.virtual.VirtualTimelineItemMapper
|
import io.element.android.libraries.matrix.impl.timeline.item.virtual.VirtualTimelineItemMapper
|
||||||
|
import kotlinx.coroutines.CoroutineDispatcher
|
||||||
import kotlinx.coroutines.CoroutineScope
|
import kotlinx.coroutines.CoroutineScope
|
||||||
import kotlinx.coroutines.FlowPreview
|
import kotlinx.coroutines.FlowPreview
|
||||||
import kotlinx.coroutines.flow.Flow
|
import kotlinx.coroutines.flow.Flow
|
||||||
|
|
@ -45,7 +45,7 @@ class RustMatrixTimeline(
|
||||||
roomCoroutineScope: CoroutineScope,
|
roomCoroutineScope: CoroutineScope,
|
||||||
private val matrixRoom: MatrixRoom,
|
private val matrixRoom: MatrixRoom,
|
||||||
private val innerRoom: Room,
|
private val innerRoom: Room,
|
||||||
private val coroutineDispatchers: CoroutineDispatchers,
|
private val dispatcher: CoroutineDispatcher,
|
||||||
) : MatrixTimeline {
|
) : MatrixTimeline {
|
||||||
|
|
||||||
private val _timelineItems: MutableStateFlow<List<MatrixTimelineItem>> =
|
private val _timelineItems: MutableStateFlow<List<MatrixTimelineItem>> =
|
||||||
|
|
@ -109,13 +109,13 @@ class RustMatrixTimeline(
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun fetchDetailsForEvent(eventId: EventId): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun fetchDetailsForEvent(eventId: EventId): Result<Unit> = withContext(dispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.fetchDetailsForEvent(eventId.value)
|
innerRoom.fetchDetailsForEvent(eventId.value)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun paginateBackwards(requestSize: Int, untilNumberOfItems: Int): Result<Unit> = withContext(coroutineDispatchers.io) {
|
override suspend fun paginateBackwards(requestSize: Int, untilNumberOfItems: Int): Result<Unit> = withContext(dispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
Timber.v("Start back paginating for room ${matrixRoom.roomId} ")
|
Timber.v("Start back paginating for room ${matrixRoom.roomId} ")
|
||||||
val paginationOptions = PaginationOptions.UntilNumItems(
|
val paginationOptions = PaginationOptions.UntilNumItems(
|
||||||
|
|
@ -131,7 +131,7 @@ class RustMatrixTimeline(
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override suspend fun sendReadReceipt(eventId: EventId) = withContext(coroutineDispatchers.io) {
|
override suspend fun sendReadReceipt(eventId: EventId) = withContext(dispatcher) {
|
||||||
runCatching {
|
runCatching {
|
||||||
innerRoom.sendReadReceipt(eventId = eventId.value)
|
innerRoom.sendReadReceipt(eventId = eventId.value)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue