diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/BackendRemoteEventHandler.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/BackendRemoteEventHandler.kt index 453c86a..d1da013 100644 --- a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/BackendRemoteEventHandler.kt +++ b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/BackendRemoteEventHandler.kt @@ -1,7 +1,7 @@ package ru.shadowsparky.vbox.backend.data import kotlinx.serialization.json.Json -import org.koin.core.annotation.Factory +import org.koin.core.annotation.Single import org.slf4j.Logger import org.slf4j.LoggerFactory import ru.shadowsparky.vbox.shared.domain.RemoteEvent @@ -9,34 +9,25 @@ import ru.shadowsparky.vbox.shared.domain.RemoteEventHandler val eventLogger: Logger = LoggerFactory.getLogger("events") -@Factory +@Single class BackendRemoteEventHandler( private val json: Json, - private val sessionCache: SessionCache + private val registry: SessionRegistry ) : RemoteEventHandler { - - suspend fun put(userId: Long, socketSession: SessionCache.Writer) { - sessionCache.put(userId, socketSession) - } - - suspend fun remove(userId: Long, socketSession: SessionCache.Writer) { - sessionCache.remove(userId, socketSession) - } - override suspend fun notify(eventInfo: RemoteEvent) { - val json = json.encodeToString(eventInfo) - val writers = sessionCache.get(eventInfo.userId) + val jsonText = json.encodeToString(eventInfo) + val writers = registry.get(eventInfo.userId) if (writers.isNullOrEmpty()) { - eventLogger.info("unable to notify $eventInfo. sessions not found, cache $sessionCache") - } else { - writers.forEach { - eventLogger.info("notify[$eventInfo]. session $it", RuntimeException("called")) - try { - it.writeText(json) - } catch (_: Exception) { - eventLogger.error("unable to notify $it. delete session") - remove(eventInfo.userId, it) - } + eventLogger.info("unable to notify {}. sessions not found, cache {}", eventInfo, registry) + return + } + writers.toList().forEach { writer -> + eventLogger.info("notify[{}]. session {}", eventInfo, writer) + try { + writer.writeText(jsonText) + } catch (e: Exception) { + eventLogger.error("unable to notify {}. delete session. Reason: {}", writer, e.message) + registry.remove(eventInfo.userId, writer) } } } diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/SessionCache.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/SessionCache.kt deleted file mode 100644 index 83aa08d..0000000 --- a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/SessionCache.kt +++ /dev/null @@ -1,42 +0,0 @@ -package ru.shadowsparky.vbox.backend.data - -import kotlinx.coroutines.sync.Mutex -import kotlinx.coroutines.sync.withLock -import org.koin.core.annotation.Factory -import java.util.Collections - -typealias SessionMap = MutableMap> - -@Factory -class SessionCache { - private val sessionMap: SessionMap = Collections.synchronizedMap(hashMapOf()) - - private val mutex = Mutex() - - suspend fun put(userId: Long, socketSession: Writer) { - eventLogger.info("put[$userId]=$socketSession") - mutex.withLock { - sessionMap[userId] = (sessionMap[userId] ?: mutableListOf()) + listOf(socketSession) - } - } - - suspend fun get(userId: Long): List? { - mutex.withLock { - return sessionMap[userId] - } - } - - suspend fun remove(userId: Long, session: Writer) { - eventLogger.info("remove[$userId]=$session") - mutex.withLock { - val sessionFromMap = sessionMap[userId]?.toMutableList() - val item = sessionFromMap?.firstOrNull { it == session } ?: return - sessionFromMap.remove(item) - sessionMap[userId] = sessionFromMap - } - } - - fun interface Writer { - suspend fun writeText(text: String) - } -} diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/SessionRegistry.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/SessionRegistry.kt new file mode 100644 index 0000000..88c2fcd --- /dev/null +++ b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/SessionRegistry.kt @@ -0,0 +1,33 @@ +package ru.shadowsparky.vbox.backend.data + +import org.koin.core.annotation.Single +import java.util.concurrent.ConcurrentHashMap + +private typealias SessionMap = ConcurrentHashMap> + +@Single +class SessionRegistry { + private val sessionMap: SessionMap = ConcurrentHashMap() + + fun put(userId: Long, socketSession: Writer) { + eventLogger.info("put[$userId]=$socketSession") + sessionMap.computeIfAbsent(userId) { ConcurrentHashMap.newKeySet() } + .add(socketSession) + } + + fun get(userId: Long): Set? { + return sessionMap[userId] + } + + fun remove(userId: Long, session: Writer) { + eventLogger.info("remove[$userId]=$session") + sessionMap.computeIfPresent(userId) { _, set -> + set.remove(session) + if (set.isEmpty()) null else set + } + } + + fun interface Writer { + suspend fun writeText(text: String) + } +} diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/tags/BackendMovieTagRepository.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/tags/BackendMovieTagRepository.kt index 919e9a9..dc5c2f1 100644 --- a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/tags/BackendMovieTagRepository.kt +++ b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/tags/BackendMovieTagRepository.kt @@ -8,11 +8,9 @@ import kotlinx.coroutines.flow.flow import kotlinx.coroutines.flow.map import kotlinx.coroutines.withContext import org.koin.core.annotation.Factory -import org.koin.core.annotation.Named import ru.shadowsparky.domain.DispatcherProvider import ru.shadowsparky.http.domain.BadRequestException import ru.shadowsparky.vbox.backend.AppDatabase -import ru.shadowsparky.vbox.shared.domain.EventType import ru.shadowsparky.vbox.shared.domain.RemoteEvent import ru.shadowsparky.vbox.shared.domain.RemoteEventHandler import ru.shadowsparky.vbox.shared.domain.tag.MovieTagRepository @@ -22,7 +20,7 @@ import ru.shadowsparky.vbox.shared.domain.tag.TagInfo class MovieTagRepositoryFactory( private val appDatabase: AppDatabase, private val dispatcherProvider: DispatcherProvider, - @Named(EventType.MOVIE_TAG) private val movieEventHandler: RemoteEventHandler + private val movieEventHandler: RemoteEventHandler ) { fun create(userId: Long): MovieTagRepository { return BackendMovieTagRepository( diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/tags/BackendUserTagRepository.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/tags/BackendUserTagRepository.kt index aac4bf9..837b415 100644 --- a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/tags/BackendUserTagRepository.kt +++ b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/data/tags/BackendUserTagRepository.kt @@ -5,10 +5,8 @@ import app.cash.sqldelight.coroutines.mapToList import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.map import org.koin.core.annotation.Factory -import org.koin.core.annotation.Named import ru.shadowsparky.vbox.backend.AppDatabase import ru.shadowsparky.vbox.shared.di.factory.DispatcherProvider -import ru.shadowsparky.vbox.shared.domain.EventType import ru.shadowsparky.vbox.shared.domain.RemoteEvent import ru.shadowsparky.vbox.shared.domain.RemoteEventHandler import ru.shadowsparky.vbox.shared.domain.tag.TagInfo @@ -18,7 +16,7 @@ import ru.shadowsparky.vbox.shared.domain.tag.UserTagRepository class UserTagRepositoryFactory( private val db: AppDatabase, private val dispatcherProvider: DispatcherProvider, - @Named(EventType.USER_TAG) private val userTagEventHandler: RemoteEventHandler + private val userTagEventHandler: RemoteEventHandler ) { fun create(userId: Long): UserTagRepository { return BackendUserTagRepository( diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/EntryPoint.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/EntryPoint.kt index cdda560..5adafcb 100644 --- a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/EntryPoint.kt +++ b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/EntryPoint.kt @@ -1,13 +1,12 @@ package ru.shadowsparky.vbox.backend.di -import org.koin.core.annotation.Named import org.koin.core.annotation.Single import ru.shadowsparky.backend.data.TokenVerifier import ru.shadowsparky.backend.domain.JwtInfo import ru.shadowsparky.backend.domain.LoginVerifier import ru.shadowsparky.chat.backend.domain.ProcessUserMessageUseCase -import ru.shadowsparky.vbox.backend.data.BackendRemoteEventHandler import ru.shadowsparky.vbox.backend.data.BackendUpdateFetcherFactory +import ru.shadowsparky.vbox.backend.data.SessionRegistry import ru.shadowsparky.vbox.backend.data.auth.AuthTokenRepositoryFactory import ru.shadowsparky.vbox.backend.data.chat.BackendChatRepositoryFactory import ru.shadowsparky.vbox.backend.data.tags.MovieTagRepositoryFactory @@ -15,7 +14,7 @@ import ru.shadowsparky.vbox.backend.data.tags.UserTagRepositoryFactory import ru.shadowsparky.vbox.backend.di.factory.RecentlyWatchedRepositoryFactory import ru.shadowsparky.vbox.backend.di.factory.SavedMovieRepositoryFactory import ru.shadowsparky.vbox.backend.di.factory.SearchRepositoryFactory -import ru.shadowsparky.vbox.shared.domain.EventType +import ru.shadowsparky.vbox.shared.domain.RemoteEventHandler import ru.shadowsparky.vbox.shared.domain.VideoApi @Single @@ -28,17 +27,14 @@ class RoutingEntryPoint( val userTagFactory: UserTagRepositoryFactory, val updateFetcherFactory: BackendUpdateFetcherFactory, val chatRepositoryFactory: BackendChatRepositoryFactory, - val processUserMessageUseCase: ProcessUserMessageUseCase + val processUserMessageUseCase: ProcessUserMessageUseCase, + val remoteEventHandler: RemoteEventHandler ) @Single class WebSocketEntryPoint( - @Named(EventType.SEARCH) val searchEventHandler: BackendRemoteEventHandler, - @Named(EventType.RECENT) val recentlyEventHandler: BackendRemoteEventHandler, - @Named(EventType.SAVED) val savedEventHandler: BackendRemoteEventHandler, - @Named(EventType.USER_TAG) val userTagEventHandler: BackendRemoteEventHandler, - @Named(EventType.MOVIE_TAG) val movieTagEventHandler: BackendRemoteEventHandler, - val tokenVerifier: TokenVerifier + val tokenVerifier: TokenVerifier, + val sessionRegistry: SessionRegistry ) @Single diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/EventModule.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/EventModule.kt index 5f5c361..41c8cba 100644 --- a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/EventModule.kt +++ b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/EventModule.kt @@ -1,32 +1,9 @@ package ru.shadowsparky.vbox.backend.di import org.koin.core.annotation.Module -import org.koin.core.annotation.Named -import org.koin.core.annotation.Single -import ru.shadowsparky.vbox.backend.data.BackendRemoteEventHandler -import ru.shadowsparky.vbox.shared.domain.EventType -import ru.shadowsparky.vbox.shared.domain.RemoteEventHandler @Module class EventModule { - @Single(binds = [BackendRemoteEventHandler::class]) - @Named(EventType.SEARCH) - fun provideSearch(impl: BackendRemoteEventHandler): RemoteEventHandler = impl - @Single(binds = [BackendRemoteEventHandler::class]) - @Named(EventType.RECENT) - fun provideRecent(impl: BackendRemoteEventHandler): RemoteEventHandler = impl - - @Single(binds = [BackendRemoteEventHandler::class]) - @Named(EventType.SAVED) - fun provideSaved(impl: BackendRemoteEventHandler): RemoteEventHandler = impl - - @Single(binds = [BackendRemoteEventHandler::class]) - @Named(EventType.USER_TAG) - fun provideUserTag(impl: BackendRemoteEventHandler): RemoteEventHandler = impl - - @Single(binds = [BackendRemoteEventHandler::class]) - @Named(EventType.MOVIE_TAG) - fun provideMovieTag(impl: BackendRemoteEventHandler): RemoteEventHandler = impl } diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/factory/RecentlyWatchedRepositoryFactory.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/factory/RecentlyWatchedRepositoryFactory.kt index 9fcbb7b..cc09ae4 100644 --- a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/factory/RecentlyWatchedRepositoryFactory.kt +++ b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/factory/RecentlyWatchedRepositoryFactory.kt @@ -1,18 +1,16 @@ package ru.shadowsparky.vbox.backend.di.factory import org.koin.core.annotation.Factory -import org.koin.core.annotation.Named import ru.shadowsparky.vbox.backend.AppDatabase import ru.shadowsparky.vbox.backend.data.BackendRecentlyWatchedRepository import ru.shadowsparky.vbox.shared.di.factory.DispatcherProvider -import ru.shadowsparky.vbox.shared.domain.EventType import ru.shadowsparky.vbox.shared.domain.RemoteEventHandler @Factory class RecentlyWatchedRepositoryFactory( private val db: AppDatabase, private val dispatcherProvider: DispatcherProvider, - @Named(EventType.RECENT) private val remoteEventHandler: RemoteEventHandler + private val remoteEventHandler: RemoteEventHandler ) { fun create(userId: Long): BackendRecentlyWatchedRepository { return BackendRecentlyWatchedRepository(db, userId, dispatcherProvider, remoteEventHandler) diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/factory/SavedMovieRepositoryFactory.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/factory/SavedMovieRepositoryFactory.kt index f60b076..66794b0 100644 --- a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/factory/SavedMovieRepositoryFactory.kt +++ b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/factory/SavedMovieRepositoryFactory.kt @@ -1,13 +1,11 @@ package ru.shadowsparky.vbox.backend.di.factory import org.koin.core.annotation.Factory -import org.koin.core.annotation.Named import ru.shadowsparky.backend.data.Logger import ru.shadowsparky.vbox.backend.AppDatabase import ru.shadowsparky.vbox.backend.data.BackendRemoteEventHandler import ru.shadowsparky.vbox.backend.data.BackendSavedMovieRepository import ru.shadowsparky.vbox.shared.di.factory.DispatcherProvider -import ru.shadowsparky.vbox.shared.domain.EventType import ru.shadowsparky.vbox.shared.domain.SavedMovieRepository @Factory @@ -15,7 +13,7 @@ class SavedMovieRepositoryFactory( private val db: AppDatabase, private val dispatcherProvider: DispatcherProvider, private val logger: Logger, - @Named(EventType.SAVED) private val eventHandler: BackendRemoteEventHandler + private val eventHandler: BackendRemoteEventHandler ) { fun create(userId: Long): SavedMovieRepository { return BackendSavedMovieRepository(db, userId, dispatcherProvider, logger, eventHandler) diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/factory/SearchRepositoryFactory.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/factory/SearchRepositoryFactory.kt index 97336e5..82d22da 100644 --- a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/factory/SearchRepositoryFactory.kt +++ b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/di/factory/SearchRepositoryFactory.kt @@ -1,11 +1,9 @@ package ru.shadowsparky.vbox.backend.di.factory import org.koin.core.annotation.Factory -import org.koin.core.annotation.Named import ru.shadowsparky.vbox.backend.AppDatabase import ru.shadowsparky.vbox.backend.data.BackendSearchRepository import ru.shadowsparky.vbox.shared.di.factory.DispatcherProvider -import ru.shadowsparky.vbox.shared.domain.EventType import ru.shadowsparky.vbox.shared.domain.RemoteEventHandler import ru.shadowsparky.vbox.shared.domain.SearchRepository @@ -13,7 +11,7 @@ import ru.shadowsparky.vbox.shared.domain.SearchRepository class SearchRepositoryFactory( private val db: AppDatabase, private val dispatcherProvider: DispatcherProvider, - @Named(EventType.SEARCH) private val eventHandler: RemoteEventHandler + private val eventHandler: RemoteEventHandler ) { fun create(userId: Long): SearchRepository { return BackendSearchRepository(db, userId, dispatcherProvider, eventHandler) diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/presentation/WebSocket.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/presentation/WebSocket.kt index 078b42a..275689f 100644 --- a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/presentation/WebSocket.kt +++ b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/presentation/WebSocket.kt @@ -2,7 +2,6 @@ package ru.shadowsparky.vbox.backend.presentation import io.ktor.server.application.Application import io.ktor.server.application.install -import io.ktor.server.routing.Route import io.ktor.server.routing.routing import io.ktor.server.websocket.WebSockets import io.ktor.server.websocket.pingPeriod @@ -11,57 +10,47 @@ import io.ktor.server.websocket.webSocket import io.ktor.websocket.Frame import io.ktor.websocket.readText import kotlinx.coroutines.CompletableDeferred -import ru.shadowsparky.backend.data.TokenVerifier -import ru.shadowsparky.backend.domain.USER_ID_ARG -import ru.shadowsparky.vbox.backend.data.BackendRemoteEventHandler -import ru.shadowsparky.vbox.backend.data.SessionCache +import kotlinx.coroutines.awaitCancellation +import ru.shadowsparky.vbox.backend.data.SessionRegistry import ru.shadowsparky.vbox.backend.di.WebSocketEntryPoint -import ru.shadowsparky.vbox.shared.domain.EventType +import ru.shadowsparky.vbox.shared.domain.RecentlyWatchedRepository import ru.shadowsparky.vbox.shared.domain.RemoteEventHandler +import ru.shadowsparky.vbox.shared.domain.SavedMovieRepository +import ru.shadowsparky.vbox.shared.domain.SearchRepository +import ru.shadowsparky.vbox.shared.domain.USER_ID_ARG +import ru.shadowsparky.vbox.shared.domain.tag.MovieTagRepository +import ru.shadowsparky.vbox.shared.domain.tag.UserTagRepository import kotlin.time.Duration.Companion.seconds fun Application.configureWebSocket(socketEntryPoint: WebSocketEntryPoint) = with(socketEntryPoint) { install(WebSockets) { - pingPeriod = (15).seconds - timeout = (15).seconds + pingPeriod = (30).seconds + timeout = (30).seconds maxFrameSize = Long.MAX_VALUE masking = false } routing { - hashMapOf( - EventType.SEARCH to searchEventHandler, - EventType.RECENT to recentlyEventHandler, - EventType.SAVED to savedEventHandler, - EventType.USER_TAG to userTagEventHandler, - EventType.MOVIE_TAG to movieTagEventHandler - ).forEach { (key, handler) -> webSocket(key, handler, tokenVerifier) } - } -} - -private fun Route.webSocket( - eventType: String, - handler: BackendRemoteEventHandler, - tokenVerifier: TokenVerifier -) { - val path = when (eventType) { - EventType.SEARCH -> RemoteEventHandler.SEARCH_CHANGED - EventType.RECENT -> RemoteEventHandler.RECENTLY_CHANGED - EventType.SAVED -> RemoteEventHandler.SAVED_CHANGED - EventType.USER_TAG -> RemoteEventHandler.USER_TAG_CHANGED - EventType.MOVIE_TAG -> RemoteEventHandler.MOVIE_TAG_CHANGED - else -> error("Unsupported type $eventType") - } - webSocket(path) { - val frame = (incoming.receive() as Frame.Text).readText() - val userId = tokenVerifier.verify(frame).getClaim(USER_ID_ARG).asLong() - val session = SessionCache.Writer { text -> outgoing.trySend(Frame.Text(text)) } - handler.put(userId, session) - val deferred = CompletableDeferred() - try { - outgoing.invokeOnClose { deferred.complete(null) } - deferred.await() - } finally { - handler.remove(userId, session) + listOf( + "${SearchRepository.PREFIX}/onChange", + "${RecentlyWatchedRepository.PREFIX}/onChange", + "${SavedMovieRepository.PREFIX}/onChange", + "${UserTagRepository.PREFIX}/onChange", + "${MovieTagRepository.PREFIX}/onChange" + ).forEach { + webSocket(it) { awaitCancellation() } + } + webSocket(RemoteEventHandler.ON_EVENT) { + val frame = (incoming.receive() as Frame.Text).readText() + val userId = tokenVerifier.verify(frame).getClaim(USER_ID_ARG).asLong() + val session = SessionRegistry.Writer { text -> outgoing.trySend(Frame.Text(text)) } + sessionRegistry.put(userId, session) + val deferred = CompletableDeferred() + try { + outgoing.invokeOnClose { deferred.complete(null) } + deferred.await() + } finally { + sessionRegistry.remove(userId, session) + } } } } diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/presentation/routing/AuthRouting.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/presentation/routing/AuthRouting.kt index c155ef5..9c4f6af 100644 --- a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/presentation/routing/AuthRouting.kt +++ b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/presentation/routing/AuthRouting.kt @@ -33,7 +33,6 @@ fun Routing.setupAuthMethods( call.respond(HttpStatusCode.OK) } authenticate(AUTH_JWT_NAME) { - get("/chech") { call.respond(HttpStatusCode.OK) } get(HeathCheck.PATH) { call.respond(HttpStatusCode.OK) } post(AuthTokenRepository.CHANGE_PASS_PATH) { authEntryPoint.authTokenRepositoryFactory.create(call.obtainUserId()) @@ -45,6 +44,6 @@ fun Routing.setupAuthMethods( setupSavedMovie(savedMovieFactory) setupTagsRouting(userTagFactory, movieTagFactory) setupUpdates(updateFetcherFactory) - setupChat(chatRepositoryFactory, processUserMessageUseCase) + setupChat(chatRepositoryFactory, processUserMessageUseCase, remoteEventHandler) } } diff --git a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/presentation/routing/ChatRouting.kt b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/presentation/routing/ChatRouting.kt index 5da29bd..fdae83f 100644 --- a/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/presentation/routing/ChatRouting.kt +++ b/apps/vbox/backend/src/main/kotlin/ru/shadowsparky/vbox/backend/presentation/routing/ChatRouting.kt @@ -14,10 +14,13 @@ import ru.shadowsparky.http.domain.BadRequestException import ru.shadowsparky.vbox.backend.data.chat.BackendChatRepositoryFactory import ru.shadowsparky.vbox.backend.presentation.obtainUserId import ru.shadowsparky.vbox.shared.domain.DYNAMIC_PREFIX +import ru.shadowsparky.vbox.shared.domain.RemoteEvent +import ru.shadowsparky.vbox.shared.domain.RemoteEventHandler fun Route.setupChat( chatRepositoryFactory: BackendChatRepositoryFactory, - processUserMessageUseCase: ProcessUserMessageUseCase + processUserMessageUseCase: ProcessUserMessageUseCase, + remoteEventHandler: RemoteEventHandler ) { get("$DYNAMIC_PREFIX/chat/query") { val repository = chatRepositoryFactory.create(call.obtainUserId()) @@ -30,6 +33,8 @@ fun Route.setupChat( if (info.role != ChatRoles.USER) throw BadRequestException("invalid role") val repository = chatRepositoryFactory.create(userId) repository.send(info) - call.respond(processUserMessageUseCase.execute(userId)) + val response = processUserMessageUseCase.execute(userId) + remoteEventHandler.notify(RemoteEvent.OnChatUpdate(userId)) + call.respond(response) } } diff --git a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/RemoteEventListenerImpl.kt b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/RemoteEventListenerImpl.kt index b3c7d63..b57b1a9 100644 --- a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/RemoteEventListenerImpl.kt +++ b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/RemoteEventListenerImpl.kt @@ -7,79 +7,68 @@ import io.ktor.http.HttpMethod import io.ktor.websocket.Frame import io.ktor.websocket.readText import kotlinx.coroutines.CancellationException -import kotlinx.coroutines.currentCoroutineContext +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.delay import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.SharingStarted import kotlinx.coroutines.flow.first import kotlinx.coroutines.flow.flow import kotlinx.coroutines.flow.retryWhen -import kotlinx.coroutines.isActive +import kotlinx.coroutines.flow.shareIn import kotlinx.serialization.json.Json -import org.koin.core.annotation.Factory +import org.koin.core.annotation.Single import ru.shadowsparky.domain.Log import ru.shadowsparky.http.domain.TokenStorage import ru.shadowsparky.vbox.shared.domain.RemoteEvent +import ru.shadowsparky.vbox.shared.domain.RemoteEventHandler import ru.shadowsparky.vbox.shared.domain.RemoteEventListener import ru.shadowsparky.vbox.shared.domain.ServerConfigurationRepository import ru.shadowsparky.vbox.shared.domain.getServerConfiguration import kotlin.math.pow +import kotlin.time.Duration.Companion.milliseconds -@Factory -class RemoteEventListenerFactory( +@Single +class RemoteEventListenerImpl( private val httpClient: HttpClient, private val serverConfigurationRepository: ServerConfigurationRepository, private val authTokenCache: TokenStorage, private val json: Json, private val log: Log -) { - fun create(suffix: String): RemoteEventListener { - return RemoteEventListenerImpl( - httpClient, - serverConfigurationRepository, - authTokenCache, - json, - log, - suffix - ) - } -} - -private class RemoteEventListenerImpl( - private val httpClient: HttpClient, - private val serverConfigurationRepository: ServerConfigurationRepository, - private val authTokenCache: TokenStorage, - private val json: Json, - private val log: Log, - private val suffix: String ) : RemoteEventListener { - private val eventFlow = flow { + private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default) + + override val event = flow { val config = serverConfigurationRepository.getServerConfiguration() val token = authTokenCache.token.first() ?: return@flow - log.d(TAG, "listener $suffix started") + log.d(TAG, "listener started") try { httpClient.webSocket( method = HttpMethod.Get, - request = { url(config.asWebSocketStr() + "/$suffix") } + request = { url(config.asWebSocketStr() + "/${RemoteEventHandler.ON_EVENT}") } ) { outgoing.send(Frame.Text(token.token)) - while (currentCoroutineContext().isActive) { - val text = (incoming.receive() as? Frame.Text)?.readText() - if (text != null) { - val decoded = json.decodeFromString(text) - log.d(TAG, "receive $text $decoded") - emit(decoded) + for (frame in incoming) { + if (frame is Frame.Text) { + val text = frame.readText() + try { + val decoded = json.decodeFromString(text) + log.d(TAG, "receive $text $decoded") + emit(decoded) + } catch (e: Exception) { + log.e(TAG, e, "failed to decode message: $text") + } } } } } catch (e: Exception) { - log.e(TAG, e, "websocket closed") + log.e(TAG, e, "websocket closed with error") throw e } finally { - log.d(TAG, "listener $suffix finished") + log.d(TAG, "listener finished") } - } - - override val event: Flow = eventFlow.retryExponential { true } + }.retryExponential { true }.shareIn(scope, SharingStarted.WhileSubscribed(500), 0) private fun Flow.retryExponential( maxRetries: Int = Int.MAX_VALUE, @@ -88,7 +77,7 @@ private class RemoteEventListenerImpl( factor: Double = 2.0, shouldRetry: (Throwable) -> Boolean = { true } ): Flow = retryWhen { cause, attempt -> - log.d(TAG, "error occurred $suffix", cause) + log.d(TAG, "error occurred, attempt=$attempt", cause) if (!shouldRetry(cause) || attempt >= maxRetries || cause is CancellationException) { false } else { @@ -96,7 +85,7 @@ private class RemoteEventListenerImpl( .toLong() .coerceAtMost(maxDelay) log.d(TAG, "retry after delay $delayTime", cause) - delay(delayTime) + delay(delayTime.milliseconds) true } } diff --git a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/recent/RemoteRecentlyWatchedRepository.kt b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/recent/RemoteRecentlyWatchedRepository.kt index cf3aad0..cf78f21 100644 --- a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/recent/RemoteRecentlyWatchedRepository.kt +++ b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/recent/RemoteRecentlyWatchedRepository.kt @@ -14,7 +14,6 @@ import ru.shadowsparky.vbox.shared.data.apply import ru.shadowsparky.vbox.shared.data.get import ru.shadowsparky.vbox.shared.data.setJsonBody import ru.shadowsparky.vbox.shared.di.DbType -import ru.shadowsparky.vbox.shared.domain.EventType import ru.shadowsparky.vbox.shared.domain.RecentlyWatchedRepository import ru.shadowsparky.vbox.shared.domain.RecentlyWatchedRequest import ru.shadowsparky.vbox.shared.domain.RecentlyWatchedResponse @@ -30,7 +29,7 @@ class RemoteRecentlyWatchedRepository( private val httpClient: HttpClient, private val serverConfigurationRepository: ServerConfigurationRepository, private val authTokenCache: TokenStorage, - @Named(EventType.RECENT) private val remoteEventListener: RemoteEventListener, + private val remoteEventListener: RemoteEventListener ) : RecentlyWatchedRepository { private suspend fun obtainConfig() = serverConfigurationRepository.getServerConfiguration() diff --git a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/saved/RemoteSavedMovieRepository.kt b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/saved/RemoteSavedMovieRepository.kt index dea4807..04bbe4f 100644 --- a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/saved/RemoteSavedMovieRepository.kt +++ b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/saved/RemoteSavedMovieRepository.kt @@ -4,13 +4,12 @@ import io.ktor.client.HttpClient import io.ktor.client.request.post import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.filter +import kotlinx.coroutines.flow.filterIsInstance import kotlinx.coroutines.flow.map import kotlinx.coroutines.flow.onStart -import org.koin.core.annotation.Named import ru.shadowsparky.vbox.shared.data.apply import ru.shadowsparky.vbox.shared.data.get import ru.shadowsparky.vbox.shared.data.setJsonBody -import ru.shadowsparky.vbox.shared.domain.EventType import ru.shadowsparky.vbox.shared.domain.IsSavedResponse import ru.shadowsparky.vbox.shared.domain.RemoteEvent import ru.shadowsparky.vbox.shared.domain.RemoteEventListener @@ -24,11 +23,11 @@ import ru.shadowsparky.vbox.shared.domain.model.VideoDetails class RemoteSavedMovieRepository( private val httpClient: HttpClient, private val serverConfigurationRepository: ServerConfigurationRepository, - @Named(EventType.SAVED) private val eventListener: RemoteEventListener + private val eventListener: RemoteEventListener ) : SavedMovieRepository { override fun getAll(): Flow> { - return eventListener.event.filter { it is RemoteEvent.OnSaved } + return eventListener.event.filterIsInstance() .map { getAllSingle() } .onStart { emit(getAllSingle()) } } diff --git a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/search/RemoteSearchRepository.kt b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/search/RemoteSearchRepository.kt index f3348e5..9997bf5 100644 --- a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/search/RemoteSearchRepository.kt +++ b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/search/RemoteSearchRepository.kt @@ -4,6 +4,7 @@ import io.ktor.client.HttpClient import io.ktor.client.request.delete import io.ktor.client.request.post import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.filterIsInstance import kotlinx.coroutines.flow.map import kotlinx.coroutines.flow.onStart import org.koin.core.annotation.Factory @@ -12,7 +13,7 @@ import ru.shadowsparky.vbox.shared.data.apply import ru.shadowsparky.vbox.shared.data.get import ru.shadowsparky.vbox.shared.data.setJsonBody import ru.shadowsparky.vbox.shared.di.DbType -import ru.shadowsparky.vbox.shared.domain.EventType +import ru.shadowsparky.vbox.shared.domain.RemoteEvent import ru.shadowsparky.vbox.shared.domain.RemoteEventListener import ru.shadowsparky.vbox.shared.domain.SearchDeleteRequest import ru.shadowsparky.vbox.shared.domain.SearchQueryRequest @@ -27,11 +28,11 @@ import ru.shadowsparky.vbox.shared.domain.getServerConfiguration class RemoteSearchRepository( private val httpClient: HttpClient, private val serverConfigurationProvider: ServerConfigurationRepository, - @Named(EventType.SEARCH) private val removeListener: RemoteEventListener + private val removeListener: RemoteEventListener ) : SearchRepository { override fun search(query: String): Flow> { - return removeListener.event + return removeListener.event.filterIsInstance() .map { searchSingle(query) } .onStart { emit(searchSingle(query)) } } diff --git a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/tag/RemoteMovieTagRepository.kt b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/tag/RemoteMovieTagRepository.kt index 01d089e..0633fc1 100644 --- a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/tag/RemoteMovieTagRepository.kt +++ b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/tag/RemoteMovieTagRepository.kt @@ -16,7 +16,6 @@ import ru.shadowsparky.vbox.shared.data.apply import ru.shadowsparky.vbox.shared.data.get import ru.shadowsparky.vbox.shared.data.setJsonBody import ru.shadowsparky.vbox.shared.di.DbType -import ru.shadowsparky.vbox.shared.domain.EventType import ru.shadowsparky.vbox.shared.domain.RemoteEvent import ru.shadowsparky.vbox.shared.domain.RemoteEventListener import ru.shadowsparky.vbox.shared.domain.ServerConfigurationRepository @@ -32,7 +31,7 @@ import ru.shadowsparky.vbox.shared.domain.tag.TagResponse class RemoteMovieTagRepository( private val httpClient: HttpClient, private val serverConfigurationProvider: ServerConfigurationRepository, - @Named(EventType.MOVIE_TAG) private val movieTagListener: RemoteEventListener, + private val movieTagListener: RemoteEventListener, private val dispatcherProvider: DispatcherProvider ) : MovieTagRepository { diff --git a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/tag/RemoteUserTagRepository.kt b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/tag/RemoteUserTagRepository.kt index bce7f86..9c99af1 100644 --- a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/tag/RemoteUserTagRepository.kt +++ b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/data/tag/RemoteUserTagRepository.kt @@ -15,7 +15,6 @@ import ru.shadowsparky.vbox.shared.data.apply import ru.shadowsparky.vbox.shared.data.get import ru.shadowsparky.vbox.shared.data.setJsonBody import ru.shadowsparky.vbox.shared.di.DbType -import ru.shadowsparky.vbox.shared.domain.EventType import ru.shadowsparky.vbox.shared.domain.RemoteEvent import ru.shadowsparky.vbox.shared.domain.RemoteEventListener import ru.shadowsparky.vbox.shared.domain.ServerConfigurationRepository @@ -32,7 +31,7 @@ import ru.shadowsparky.vbox.shared.domain.tag.UserTagRepository class RemoteUserTagRepository( private val httpClient: HttpClient, private val serverConfigurationProvider: ServerConfigurationRepository, - @Named(EventType.USER_TAG) private val userTagListener: RemoteEventListener, + private val userTagListener: RemoteEventListener, private val dispatcherProvider: DispatcherProvider ) : UserTagRepository { diff --git a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/di/EventsModule.kt b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/di/EventsModule.kt index e2c63b7..422aff8 100644 --- a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/di/EventsModule.kt +++ b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/di/EventsModule.kt @@ -1,53 +1,7 @@ package ru.shadowsparky.vbox.shared.di -import org.koin.core.annotation.Factory import org.koin.core.annotation.Module -import org.koin.core.annotation.Named -import ru.shadowsparky.vbox.shared.data.RemoteEventListenerFactory -import ru.shadowsparky.vbox.shared.domain.EventType -import ru.shadowsparky.vbox.shared.domain.RemoteEventHandler -import ru.shadowsparky.vbox.shared.domain.RemoteEventListener @Module class EventsModule { - - @Factory - @Named(EventType.RECENT) - fun provideRecentRemoteEventListener( - factory: RemoteEventListenerFactory - ): RemoteEventListener { - return factory.create(RemoteEventHandler.RECENTLY_CHANGED) - } - - @Factory - @Named(EventType.SEARCH) - fun provideSearchRemoteEventListener( - factory: RemoteEventListenerFactory - ): RemoteEventListener { - return factory.create(RemoteEventHandler.SEARCH_CHANGED) - } - - @Factory - @Named(EventType.SAVED) - fun provideSavedRemoteEventListener( - factory: RemoteEventListenerFactory - ): RemoteEventListener { - return factory.create(RemoteEventHandler.SAVED_CHANGED) - } - - @Factory - @Named(EventType.MOVIE_TAG) - fun provideMovieTagRemoteEventListener( - factory: RemoteEventListenerFactory - ): RemoteEventListener { - return factory.create(RemoteEventHandler.MOVIE_TAG_CHANGED) - } - - @Factory - @Named(EventType.USER_TAG) - fun provideUserTagRemoteEventListener( - factory: RemoteEventListenerFactory - ): RemoteEventListener { - return factory.create(RemoteEventHandler.USER_TAG_CHANGED) - } } diff --git a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/di/SavedModule.kt b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/di/SavedModule.kt index 9bf06ea..870d563 100644 --- a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/di/SavedModule.kt +++ b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/di/SavedModule.kt @@ -9,7 +9,6 @@ import ru.shadowsparky.vbox.shared.data.AppDatabaseProvider import ru.shadowsparky.vbox.shared.data.saved.PosterWithoutSuffixSavedMovieRepository import ru.shadowsparky.vbox.shared.data.saved.RemoteSavedMovieRepository import ru.shadowsparky.vbox.shared.data.saved.SqlSavedMovieRepository -import ru.shadowsparky.vbox.shared.domain.EventType import ru.shadowsparky.vbox.shared.domain.RemoteEventListener import ru.shadowsparky.vbox.shared.domain.SavedMovieRepository import ru.shadowsparky.vbox.shared.domain.ServerConfigurationRepository @@ -32,7 +31,7 @@ class SavedModule { fun provideSavedRemoteRepository( httpClient: HttpClient, serverConfigurationRepository: ServerConfigurationRepository, - @Named(EventType.SAVED) eventListener: RemoteEventListener + eventListener: RemoteEventListener ): SavedMovieRepository { val impl = RemoteSavedMovieRepository(httpClient, serverConfigurationRepository, eventListener) diff --git a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/presentation/VBoxRootComponent.kt b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/presentation/VBoxRootComponent.kt index e090292..0e4a56f 100644 --- a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/presentation/VBoxRootComponent.kt +++ b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/presentation/VBoxRootComponent.kt @@ -6,16 +6,20 @@ import com.arkivanov.decompose.childContext import com.arkivanov.decompose.router.stack.childStackWebNavigation import com.arkivanov.decompose.router.webhistory.WebNavigation import com.arkivanov.decompose.router.webhistory.WebNavigationOwner +import com.arkivanov.essenty.lifecycle.coroutines.repeatOnLifecycle +import kotlinx.coroutines.flow.collect import kotlinx.coroutines.flow.map import kotlinx.coroutines.launch import kotlinx.serialization.KSerializer import org.koin.core.Koin import org.koin.core.annotation.Factory import ru.shadowsparky.http.domain.TokenStorage +import ru.shadowsparky.http.domain.hasAccount import ru.shadowsparky.ui.nav.RootComponent import ru.shadowsparky.updater.presentation.UpdateComponentFactory import ru.shadowsparky.vbox.shared.di.factory.ImageLoaderProvider import ru.shadowsparky.vbox.shared.domain.HeathCheck +import ru.shadowsparky.vbox.shared.domain.RemoteEventListener import ru.shadowsparky.vbox.shared.presentation.nav.Route @OptIn(ExperimentalDecomposeApi::class) @@ -26,6 +30,7 @@ class VBoxRootComponent( private val healthCheck: HeathCheck, val imageLoaderProvider: ImageLoaderProvider, updateComponentFactory: UpdateComponentFactory, + private val remoteEventListener: RemoteEventListener, override val serializer: KSerializer = Route.serializer(), override val koin: Koin = ru.shadowsparky.koin ) : RootComponent(componentContext), WebNavigationOwner { @@ -41,7 +46,20 @@ class VBoxRootComponent( } init { - scope.launch { updateComponent.checkUpdates(false) } + scope.launch { + repeatOnLifecycle { + if (tokenStorage.hasAccount()) { + remoteEventListener.event.collect() + } + } + } + scope.launch { + tokenStorage.token.collect { + if (it != null) { + updateComponent.checkUpdates(false) + } + } + } } override val webNavigation: WebNavigation<*> = @@ -63,7 +81,8 @@ class RootFactory( private val tokenStorage: TokenStorage, private val heathCheck: HeathCheck, private val imageLoaderProvider: ImageLoaderProvider, - private val updateComponentFactory: UpdateComponentFactory + private val updateComponentFactory: UpdateComponentFactory, + private val remoteEventListener: RemoteEventListener ) { fun create( context: ComponentContext, @@ -82,7 +101,8 @@ class RootFactory( listOf(Route.Loading(next = initStack.toSet())), heathCheck, imageLoaderProvider, - updateComponentFactory + updateComponentFactory, + remoteEventListener ) } } diff --git a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/presentation/chat/ChatComponentFactory.kt b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/presentation/chat/ChatComponentFactory.kt index f289cbf..60ddfbf 100644 --- a/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/presentation/chat/ChatComponentFactory.kt +++ b/apps/vbox/client/shared/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/presentation/chat/ChatComponentFactory.kt @@ -1,6 +1,8 @@ package ru.shadowsparky.vbox.shared.presentation.chat import com.arkivanov.decompose.ComponentContext +import kotlinx.coroutines.flow.filterIsInstance +import kotlinx.coroutines.flow.map import org.koin.core.annotation.Factory import org.koin.core.annotation.Named import ru.shadowsparky.chat.domain.ChatRepository @@ -8,6 +10,8 @@ import ru.shadowsparky.chat.domain.ChatTextParser import ru.shadowsparky.chat.presentation.ChatComponent import ru.shadowsparky.ui.nav.ComponentFactory import ru.shadowsparky.ui.nav.RouteNavigator +import ru.shadowsparky.vbox.shared.domain.RemoteEvent +import ru.shadowsparky.vbox.shared.domain.RemoteEventListener import ru.shadowsparky.vbox.shared.presentation.nav.Child import ru.shadowsparky.vbox.shared.presentation.nav.Route import ru.shadowsparky.vbox.shared.presentation.nav.RouteIds @@ -16,7 +20,8 @@ import ru.shadowsparky.vbox.shared.presentation.nav.RouteIds @Named(RouteIds.CHAT) class ChatComponentFactory( private val chatRepository: ChatRepository, - private val textParser: ChatTextParser + private val textParser: ChatTextParser, + private val remoteEventListener: RemoteEventListener ) : ComponentFactory { override fun create( route: Route.Chat, @@ -27,7 +32,8 @@ class ChatComponentFactory( child, chatRepository, textParser, - onLinkClicked = { root.nav(Route.Videos(it)) } + onLinkClicked = { root.nav(Route.Videos(it)) }, + refreshFlow = remoteEventListener.event.filterIsInstance().map {} ) return Child.ChatChild(comp) } diff --git a/apps/vbox/common/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/domain/RemoteEventHandler.kt b/apps/vbox/common/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/domain/RemoteEventHandler.kt index d42833b..ed59633 100644 --- a/apps/vbox/common/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/domain/RemoteEventHandler.kt +++ b/apps/vbox/common/src/commonMain/kotlin/ru/shadowsparky/vbox/shared/domain/RemoteEventHandler.kt @@ -1,19 +1,14 @@ package ru.shadowsparky.vbox.shared.domain import kotlinx.coroutines.flow.Flow +import kotlinx.serialization.SerialName import kotlinx.serialization.Serializable -import ru.shadowsparky.vbox.shared.domain.tag.MovieTagRepository -import ru.shadowsparky.vbox.shared.domain.tag.UserTagRepository interface RemoteEventHandler { suspend fun notify(eventInfo: RemoteEvent) companion object { - const val SEARCH_CHANGED = "${SearchRepository.PREFIX}/onChange" - const val RECENTLY_CHANGED = "${RecentlyWatchedRepository.PREFIX}/onChange" - const val SAVED_CHANGED = "${SavedMovieRepository.PREFIX}/onChange" - const val USER_TAG_CHANGED = "${UserTagRepository.PREFIX}/onChange" - const val MOVIE_TAG_CHANGED = "${MovieTagRepository.PREFIX}/onChange" + const val ON_EVENT = "$DYNAMIC_PREFIX/onEvent" } } @@ -26,26 +21,27 @@ sealed interface RemoteEvent { val userId: Long @Serializable + @SerialName("OnSearch") data class OnSearch(val millis: Long, override val userId: Long) : RemoteEvent @Serializable + @SerialName("OnRecent") data class OnRecent(val millis: Long, override val userId: Long, val movieId: Long) : RemoteEvent @Serializable + @SerialName("OnSaved") data class OnSaved(override val userId: Long, val movieId: Long) : RemoteEvent @Serializable + @SerialName("OnUserTag") data class OnUserTag(override val userId: Long) : RemoteEvent @Serializable + @SerialName("OnMovieTag") data class OnMovieTag(override val userId: Long, val movieId: Long) : RemoteEvent -} -object EventType { - const val SEARCH = "search" - const val RECENT = "recent" - const val SAVED = "saved" - const val USER_TAG = "user_tag" - const val MOVIE_TAG = "movie_tag" + @Serializable + @SerialName("OnChatUpdate") + data class OnChatUpdate(override val userId: Long) : RemoteEvent } diff --git a/feature/chat/chat-client/src/commonMain/kotlin/ru/shadowsparky/chat/presentation/ChatComponent.kt b/feature/chat/chat-client/src/commonMain/kotlin/ru/shadowsparky/chat/presentation/ChatComponent.kt index 16b6618..9e01e4e 100644 --- a/feature/chat/chat-client/src/commonMain/kotlin/ru/shadowsparky/chat/presentation/ChatComponent.kt +++ b/feature/chat/chat-client/src/commonMain/kotlin/ru/shadowsparky/chat/presentation/ChatComponent.kt @@ -3,11 +3,14 @@ package ru.shadowsparky.chat.presentation import com.arkivanov.decompose.ComponentContext import com.arkivanov.decompose.value.MutableValue import com.arkivanov.decompose.value.Value +import com.arkivanov.essenty.lifecycle.coroutines.repeatOnLifecycle import com.arkivanov.essenty.lifecycle.doOnDestroy +import kotlinx.coroutines.CoroutineExceptionHandler import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel +import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.launch import ru.shadowsparky.chat.domain.ChatRepository import ru.shadowsparky.chat.domain.ChatRoles @@ -31,10 +34,14 @@ class ChatComponent( componentContext: ComponentContext, private val repository: ChatRepository, private val textParser: ChatTextParser?, - val onLinkClicked: (String) -> Unit + val onLinkClicked: (String) -> Unit, + val refreshFlow: Flow ) : ComponentContext by componentContext { + private val errorHandler = CoroutineExceptionHandler { _, throwable -> + _state.value = ChatState.Error(throwable.message ?: throwable.toString()) + } - private val scope = CoroutineScope(Dispatchers.Main.immediate + SupervisorJob()).also { ctx -> + private val scope = CoroutineScope(SupervisorJob() + errorHandler + Dispatchers.Main.immediate).also { ctx -> lifecycle.doOnDestroy { ctx.cancel() } } @@ -43,6 +50,9 @@ class ChatComponent( init { loadInitialMessages() + scope.launch { + repeatOnLifecycle { refreshFlow.collect { refreshList() } } + } } fun parseMessage(content: String): List { @@ -52,12 +62,8 @@ class ChatComponent( fun loadInitialMessages() { _state.value = ChatState.Loading scope.launch { - try { - val dbList = repository.query(afterId = null) - _state.value = ChatState.Data(messages = dbList) - } catch (e: Exception) { - _state.value = ChatState.Error(e.message ?: e.toString()) - } + val dbList = repository.query(afterId = null) + _state.value = ChatState.Data(messages = dbList) } } @@ -66,20 +72,20 @@ class ChatComponent( if (content.isBlank() || currentState.isSending) return _state.value = currentState.copy(isSending = true, sendError = null) scope.launch { - try { - repository.send(Message(role = ChatRoles.USER, content = content, timestamp = 0L)) - val latestMessageId = currentState.messages.lastOrNull()?.id - val freshMessages = repository.query(afterId = latestMessageId) - _state.value = currentState.copy( - messages = currentState.messages + freshMessages, - isSending = false - ) - } catch (e: Exception) { - _state.value = currentState.copy( - isSending = false, - sendError = e.message ?: e.toString() - ) - } + repository.send(Message(role = ChatRoles.USER, content = content, timestamp = 0L)) + } + } + + private fun refreshList() { + val currentState = _state.value as? ChatState.Data ?: return + scope.launch { + _state.value = currentState.copy(isSending = true) + val latestMessageId = currentState.messages.lastOrNull()?.id + val freshMessages = repository.query(afterId = latestMessageId) + _state.value = currentState.copy( + messages = currentState.messages + freshMessages, + isSending = false + ) } }