add some stuff to backend base
This commit is contained in:
+16
-80
@@ -4,114 +4,50 @@ import io.ktor.client.HttpClient
|
||||
import io.ktor.client.plugins.websocket.webSocket
|
||||
import io.ktor.client.request.url
|
||||
import io.ktor.http.HttpMethod
|
||||
import io.ktor.websocket.Frame
|
||||
import io.ktor.websocket.readText
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.channels.ReceiveChannel
|
||||
import kotlinx.coroutines.channels.SendChannel
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.SharingStarted
|
||||
import kotlinx.coroutines.flow.emitAll
|
||||
import kotlinx.coroutines.flow.filterNotNull
|
||||
import kotlinx.coroutines.flow.retryWhen
|
||||
import kotlinx.coroutines.flow.map
|
||||
import kotlinx.coroutines.flow.shareIn
|
||||
import kotlinx.coroutines.flow.transformLatest
|
||||
import kotlinx.serialization.json.Json
|
||||
import org.koin.core.annotation.Single
|
||||
import ru.shadowsparky.domain.Log
|
||||
import ru.shadowsparky.http.data.RemoveEventProcessor
|
||||
import ru.shadowsparky.http.data.retryExponential
|
||||
import ru.shadowsparky.http.domain.TokenStorage
|
||||
import ru.shadowsparky.vbox.shared.domain.AuthRequest
|
||||
import ru.shadowsparky.vbox.shared.domain.AuthResponse
|
||||
import ru.shadowsparky.vbox.shared.domain.HealthCheck
|
||||
import ru.shadowsparky.vbox.shared.domain.ON_EVENT
|
||||
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
|
||||
|
||||
@Single
|
||||
class RemoteEventListenerImpl(
|
||||
private val httpClient: HttpClient,
|
||||
private val serverConfigurationRepository: ServerConfigurationRepository,
|
||||
private val authTokenCache: TokenStorage,
|
||||
authTokenCache: TokenStorage,
|
||||
private val healthCheck: HealthCheck,
|
||||
private val json: Json,
|
||||
private val log: Log
|
||||
private val processor: RemoveEventProcessor,
|
||||
private val json: Json
|
||||
) : RemoteEventListener {
|
||||
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
|
||||
|
||||
@OptIn(ExperimentalCoroutinesApi::class)
|
||||
override val event = authTokenCache.token.filterNotNull().transformLatest { token ->
|
||||
override val event: Flow<RemoteEvent> = authTokenCache.token.filterNotNull().transformLatest { token ->
|
||||
val config = serverConfigurationRepository.getServerConfiguration()
|
||||
log.d(TAG, "listener started")
|
||||
try {
|
||||
httpClient.webSocket(
|
||||
method = HttpMethod.Get,
|
||||
request = { url(config.asWebSocketStr() + "/${RemoteEventHandler.ON_EVENT}") }
|
||||
) {
|
||||
authV2(token.token, incoming, outgoing)
|
||||
for (frame in incoming) {
|
||||
if (frame is Frame.Text) {
|
||||
val text = frame.readText()
|
||||
try {
|
||||
val decoded = json.decodeFromString<RemoteEvent>(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 with error")
|
||||
throw e
|
||||
} finally {
|
||||
log.d(TAG, "listener finished")
|
||||
httpClient.webSocket(
|
||||
method = HttpMethod.Get,
|
||||
request = { url(config.asWebSocketStr() + "/${ON_EVENT}") }
|
||||
) {
|
||||
val flow = processor.process(token.token, this) { healthCheck.check() }
|
||||
.map { json.decodeFromString<RemoteEvent>(it) }
|
||||
emitAll(flow)
|
||||
}
|
||||
}.retryExponential { true }.shareIn(scope, SharingStarted.WhileSubscribed(500), 0)
|
||||
|
||||
private suspend fun authV2(
|
||||
token: String,
|
||||
input: ReceiveChannel<Frame>,
|
||||
output: SendChannel<Frame>
|
||||
) {
|
||||
output.send(Frame.Text(json.encodeToString(AuthRequest(token))))
|
||||
val rawFrame = (input.receive() as Frame.Text).readText()
|
||||
val response = json.decodeFromString<AuthResponse>(rawFrame)
|
||||
if (!response.ok) {
|
||||
healthCheck.check()
|
||||
error("Authentication failed")
|
||||
}
|
||||
}
|
||||
|
||||
private fun <T> Flow<T>.retryExponential(
|
||||
maxRetries: Int = Int.MAX_VALUE,
|
||||
initialDelay: Long = 5000L,
|
||||
maxDelay: Long = 300_000L,
|
||||
factor: Double = 2.0,
|
||||
shouldRetry: (Throwable) -> Boolean = { true }
|
||||
): Flow<T> = retryWhen { cause, attempt ->
|
||||
log.d(TAG, "error occurred, attempt=$attempt", cause)
|
||||
if (!shouldRetry(cause) || attempt >= maxRetries || cause is CancellationException) {
|
||||
false
|
||||
} else {
|
||||
val delayTime = (initialDelay * factor.pow(attempt.toDouble()))
|
||||
.toLong()
|
||||
.coerceAtMost(maxDelay)
|
||||
log.d(TAG, "retry after delay $delayTime", cause)
|
||||
delay(delayTime.milliseconds)
|
||||
true
|
||||
}
|
||||
}
|
||||
|
||||
private companion object {
|
||||
const val TAG = "RemoteSearchEventListener"
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user