update events logic

This commit is contained in:
2026-08-03 21:42:50 +03:00
parent 986078a84c
commit 7fb482f46f
25 changed files with 208 additions and 303 deletions
@@ -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)
}
}
}
@@ -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<Long, List<SessionCache.Writer>>
@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<Writer>? {
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)
}
}
@@ -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<Long, MutableSet<SessionRegistry.Writer>>
@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<Writer>? {
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)
}
}
@@ -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(
@@ -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(
@@ -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
@@ -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
}
@@ -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)
@@ -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)
@@ -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)
@@ -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<Unit?>()
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<Unit?>()
try {
outgoing.invokeOnClose { deferred.complete(null) }
deferred.await()
} finally {
sessionRegistry.remove(userId, session)
}
}
}
}
@@ -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)
}
}
@@ -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)
}
}