Update rust sdk to 0.1.31: new app service
This commit is contained in:
parent
d3a86bffee
commit
ec04250a9b
4 changed files with 61 additions and 19 deletions
|
|
@ -47,6 +47,7 @@ import io.element.android.libraries.matrix.impl.room.RoomContentForwarder
|
||||||
import io.element.android.libraries.matrix.impl.room.RustMatrixRoom
|
import io.element.android.libraries.matrix.impl.room.RustMatrixRoom
|
||||||
import io.element.android.libraries.matrix.impl.room.RustRoomSummaryDataSource
|
import io.element.android.libraries.matrix.impl.room.RustRoomSummaryDataSource
|
||||||
import io.element.android.libraries.matrix.impl.room.roomOrNull
|
import io.element.android.libraries.matrix.impl.room.roomOrNull
|
||||||
|
import io.element.android.libraries.matrix.impl.room.stateFlow
|
||||||
import io.element.android.libraries.matrix.impl.sync.RustSyncService
|
import io.element.android.libraries.matrix.impl.sync.RustSyncService
|
||||||
import io.element.android.libraries.matrix.impl.usersearch.UserProfileMapper
|
import io.element.android.libraries.matrix.impl.usersearch.UserProfileMapper
|
||||||
import io.element.android.libraries.matrix.impl.usersearch.UserSearchResultMapper
|
import io.element.android.libraries.matrix.impl.usersearch.UserSearchResultMapper
|
||||||
|
|
@ -85,11 +86,14 @@ class RustMatrixClient constructor(
|
||||||
) : MatrixClient {
|
) : MatrixClient {
|
||||||
|
|
||||||
override val sessionId: UserId = UserId(client.userId())
|
override val sessionId: UserId = UserId(client.userId())
|
||||||
private val roomListService = client.roomListServiceWithEncryption()
|
private val app = client.app().use { builder ->
|
||||||
|
builder.finish()
|
||||||
|
}
|
||||||
|
private val roomListService = app.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 verificationService = RustSessionVerificationService()
|
private val verificationService = RustSessionVerificationService()
|
||||||
private val syncService = RustSyncService(roomListService, sessionCoroutineScope)
|
private val syncService = RustSyncService(app, roomListService.stateFlow(), sessionCoroutineScope)
|
||||||
private val pushersService = RustPushersService(
|
private val pushersService = RustPushersService(
|
||||||
client = client,
|
client = client,
|
||||||
dispatchers = dispatchers,
|
dispatchers = dispatchers,
|
||||||
|
|
@ -252,6 +256,7 @@ class RustMatrixClient constructor(
|
||||||
sessionCoroutineScope.cancel()
|
sessionCoroutineScope.cancel()
|
||||||
client.setDelegate(null)
|
client.setDelegate(null)
|
||||||
verificationService.destroy()
|
verificationService.destroy()
|
||||||
|
app.destroy()
|
||||||
roomListService.destroy()
|
roomListService.destroy()
|
||||||
notificationClient.destroy()
|
notificationClient.destroy()
|
||||||
client.destroy()
|
client.destroy()
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,36 @@
|
||||||
|
/*
|
||||||
|
* Copyright (c) 2023 New Vector Ltd
|
||||||
|
*
|
||||||
|
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
* you may not use this file except in compliance with the License.
|
||||||
|
* You may obtain a copy of the License at
|
||||||
|
*
|
||||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
*
|
||||||
|
* Unless required by applicable law or agreed to in writing, software
|
||||||
|
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
* See the License for the specific language governing permissions and
|
||||||
|
* limitations under the License.
|
||||||
|
*/
|
||||||
|
|
||||||
|
package io.element.android.libraries.matrix.impl.sync
|
||||||
|
|
||||||
|
import io.element.android.libraries.matrix.impl.util.mxCallbackFlow
|
||||||
|
import kotlinx.coroutines.channels.Channel
|
||||||
|
import kotlinx.coroutines.channels.trySendBlocking
|
||||||
|
import kotlinx.coroutines.flow.Flow
|
||||||
|
import kotlinx.coroutines.flow.buffer
|
||||||
|
import org.matrix.rustcomponents.sdk.App
|
||||||
|
import org.matrix.rustcomponents.sdk.AppState
|
||||||
|
import org.matrix.rustcomponents.sdk.AppStateObserver
|
||||||
|
|
||||||
|
fun App.stateFlow(): Flow<AppState> =
|
||||||
|
mxCallbackFlow {
|
||||||
|
val listener = object : AppStateObserver {
|
||||||
|
override fun onUpdate(state: AppState) {
|
||||||
|
trySendBlocking(state)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
state(listener)
|
||||||
|
}.buffer(Channel.UNLIMITED)
|
||||||
|
|
@ -17,6 +17,7 @@
|
||||||
package io.element.android.libraries.matrix.impl.sync
|
package io.element.android.libraries.matrix.impl.sync
|
||||||
|
|
||||||
import io.element.android.libraries.matrix.api.sync.SyncState
|
import io.element.android.libraries.matrix.api.sync.SyncState
|
||||||
|
import org.matrix.rustcomponents.sdk.AppState
|
||||||
import org.matrix.rustcomponents.sdk.RoomListServiceState
|
import org.matrix.rustcomponents.sdk.RoomListServiceState
|
||||||
|
|
||||||
internal fun RoomListServiceState.toSyncState(): SyncState {
|
internal fun RoomListServiceState.toSyncState(): SyncState {
|
||||||
|
|
@ -28,3 +29,11 @@ internal fun RoomListServiceState.toSyncState(): SyncState {
|
||||||
RoomListServiceState.TERMINATED -> SyncState.Terminated
|
RoomListServiceState.TERMINATED -> SyncState.Terminated
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
internal fun AppState.toSyncState(): SyncState {
|
||||||
|
return when (this) {
|
||||||
|
AppState.RUNNING -> SyncState.Syncing
|
||||||
|
AppState.TERMINATED -> SyncState.Terminated
|
||||||
|
AppState.ERROR -> SyncState.InError
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -18,47 +18,39 @@ package io.element.android.libraries.matrix.impl.sync
|
||||||
|
|
||||||
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 io.element.android.libraries.matrix.impl.room.stateFlow
|
|
||||||
import kotlinx.coroutines.CoroutineScope
|
import kotlinx.coroutines.CoroutineScope
|
||||||
|
import kotlinx.coroutines.flow.Flow
|
||||||
import kotlinx.coroutines.flow.SharingStarted
|
import kotlinx.coroutines.flow.SharingStarted
|
||||||
import kotlinx.coroutines.flow.StateFlow
|
import kotlinx.coroutines.flow.StateFlow
|
||||||
import kotlinx.coroutines.flow.distinctUntilChanged
|
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 org.matrix.rustcomponents.sdk.RoomListService
|
import org.matrix.rustcomponents.sdk.App
|
||||||
import org.matrix.rustcomponents.sdk.RoomListServiceState
|
import org.matrix.rustcomponents.sdk.RoomListServiceState
|
||||||
import timber.log.Timber
|
import timber.log.Timber
|
||||||
import java.util.concurrent.atomic.AtomicBoolean
|
|
||||||
|
|
||||||
class RustSyncService(
|
class RustSyncService(
|
||||||
private val roomListService: RoomListService,
|
private val app: App,
|
||||||
|
roomListStateFlow: Flow<RoomListServiceState>,
|
||||||
sessionCoroutineScope: CoroutineScope
|
sessionCoroutineScope: CoroutineScope
|
||||||
) : SyncService {
|
) : SyncService {
|
||||||
|
|
||||||
private val isSyncing = AtomicBoolean(false)
|
|
||||||
|
|
||||||
override fun startSync() = runCatching {
|
override fun startSync() = runCatching {
|
||||||
if (isSyncing.compareAndSet(false, true)) {
|
Timber.v("Start sync")
|
||||||
Timber.v("Start sync")
|
app.start()
|
||||||
roomListService.sync()
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun stopSync() = runCatching {
|
override fun stopSync() = runCatching {
|
||||||
if (isSyncing.compareAndSet(true, false)) {
|
Timber.v("Stop sync")
|
||||||
Timber.v("Stop sync")
|
app.pause()
|
||||||
roomListService.stopSync()
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
override val syncState: StateFlow<SyncState> =
|
override val syncState: StateFlow<SyncState> =
|
||||||
roomListService
|
roomListStateFlow
|
||||||
.stateFlow()
|
|
||||||
.map(RoomListServiceState::toSyncState)
|
.map(RoomListServiceState::toSyncState)
|
||||||
.onEach { state ->
|
.onEach { state ->
|
||||||
Timber.v("Sync state=$state")
|
Timber.v("Sync state=$state")
|
||||||
isSyncing.set(state == SyncState.Syncing)
|
|
||||||
}
|
}
|
||||||
.distinctUntilChanged()
|
.distinctUntilChanged()
|
||||||
.stateIn(sessionCoroutineScope, SharingStarted.Eagerly, SyncState.Idle)
|
.stateIn(sessionCoroutineScope, SharingStarted.Eagerly, SyncState.Idle)
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue