feat: harden offline sync and add diagnostics

This commit is contained in:
admin
2026-07-10 23:58:07 +02:00
parent 0564101e9e
commit 89235fbb46
34 changed files with 1828 additions and 197 deletions
@@ -7,6 +7,7 @@ import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.launch
import pl.firmatpp.kierowca.data.DriverRepository
import pl.firmatpp.kierowca.data.sync.NetworkMonitor
class DriverFirebaseMessagingService : FirebaseMessagingService() {
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
@@ -14,7 +15,7 @@ class DriverFirebaseMessagingService : FirebaseMessagingService() {
override fun onNewToken(token: String) {
scope.launch {
val repository = DriverRepository(applicationContext)
if (repository.hasToken()) {
if (repository.hasToken() && NetworkMonitor(applicationContext).isCurrentlyValidated()) {
runCatching { repository.storePushToken(token) }
}
}
@@ -38,20 +39,26 @@ class DriverFirebaseMessagingService : FirebaseMessagingService() {
title = data["title"],
body = data["body"],
)
DriverSyncWorker.enqueue(
context = applicationContext,
date = data["date"],
routeId = data["routeId"],
)
if (!DriverRuntimeSyncState.websocketOwnsForegroundHints()) {
DriverSyncWorker.enqueue(
context = applicationContext,
date = data["date"],
routeId = data["routeId"],
)
}
return
}
if (data["type"] != "driver_sync_hint") return
if (DriverRuntimeSyncState.websocketOwnsForegroundHints()) return
DriverSyncWorker.enqueue(
context = applicationContext,
date = data["date"],
routeId = data["routeId"],
scope = data["scope"],
version = data["version"]?.toLongOrNull(),
checksum = data["checksum"],
)
}
}
@@ -11,6 +11,9 @@ import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import okhttp3.OkHttpClient
import okhttp3.Request
import okhttp3.WebSocket
@@ -19,6 +22,8 @@ import pl.firmatpp.kierowca.data.DriverRepository
import pl.firmatpp.kierowca.data.model.BroadcastAuthResponse
import pl.firmatpp.kierowca.data.model.RealtimeConfigDto
import pl.firmatpp.kierowca.diagnostics.AppDiagnostics
import pl.firmatpp.kierowca.data.ApiErrorKind
import pl.firmatpp.kierowca.data.ApiErrorMapper
interface DriverLiveSyncGateway {
suspend fun broadcastAuth(socketId: String, channelName: String): BroadcastAuthResponse
@@ -47,7 +52,7 @@ class OkHttpLiveWebSocketFactory(
client.newWebSocket(Request.Builder().url(url).build(), listener)
}
private enum class LiveSyncConnectionState {
enum class LiveSyncConnectionState {
Stopped,
WaitingForNetwork,
Disconnected,
@@ -56,10 +61,24 @@ private enum class LiveSyncConnectionState {
Connected,
}
data class LiveSyncDiagnostics(
val state: LiveSyncConnectionState,
val desiredActive: Boolean,
val networkAvailable: Boolean,
val foregroundActive: Boolean,
val configured: Boolean,
val socketId: String?,
val reconnectAttempt: Int,
val messageVersion: Long,
val lastMessageAtEpochMillis: Long?,
val lastDisconnectReason: String?,
)
class DriverLiveSyncClient(
private val gateway: DriverLiveSyncGateway,
private val onConnected: () -> Unit,
private val onHint: (DriverSyncHint) -> Unit,
private val onAuthenticationFailed: () -> Unit = {},
private val webSocketFactory: LiveWebSocketFactory = OkHttpLiveWebSocketFactory(defaultOkHttpClient()),
private val gson: Gson = Gson(),
private val scope: CoroutineScope = CoroutineScope(SupervisorJob() + Dispatchers.IO),
@@ -70,12 +89,14 @@ class DriverLiveSyncClient(
repository: DriverRepository,
onConnected: () -> Unit,
onHint: (DriverSyncHint) -> Unit,
onAuthenticationFailed: () -> Unit = {},
client: OkHttpClient = defaultOkHttpClient(),
gson: Gson = Gson(),
) : this(
gateway = DriverRepositoryLiveSyncGateway(repository),
onConnected = onConnected,
onHint = onHint,
onAuthenticationFailed = onAuthenticationFailed,
webSocketFactory = OkHttpLiveWebSocketFactory(client),
gson = gson,
)
@@ -104,6 +125,31 @@ class DriverLiveSyncClient(
private var staleWatchdogJob: Job? = null
private var reconnectAttempt: Int = 0
private var messageVersion: Long = 0
private var lastMessageAtEpochMillis: Long? = null
private var lastDisconnectReason: String? = null
private var foregroundActive: Boolean = true
private val _connectionState = MutableStateFlow(LiveSyncConnectionState.Stopped)
val connectionState: StateFlow<LiveSyncConnectionState> = _connectionState.asStateFlow()
fun diagnostics(): LiveSyncDiagnostics = synchronized(lock) {
LiveSyncDiagnostics(
state = state,
desiredActive = desiredActive,
networkAvailable = networkAvailable,
foregroundActive = foregroundActive,
configured = realtimeConfig != null,
socketId = socketId,
reconnectAttempt = reconnectAttempt,
messageVersion = messageVersion,
lastMessageAtEpochMillis = lastMessageAtEpochMillis,
lastDisconnectReason = lastDisconnectReason,
)
}
private fun updateState(next: LiveSyncConnectionState) {
state = next
_connectionState.value = next
}
fun start(driverId: String, config: RealtimeConfigDto?) {
if (!isConfigUsable(config)) return
@@ -119,6 +165,7 @@ class DriverLiveSyncClient(
fun ensureConnected() {
val shouldConnect = synchronized(lock) {
desiredActive &&
foregroundActive &&
networkAvailable &&
state !in setOf(
LiveSyncConnectionState.Connecting,
@@ -136,22 +183,32 @@ class DriverLiveSyncClient(
reconnectJob?.cancel()
reconnectJob = null
closeSocketLocked("network_lost")
state = if (desiredActive) LiveSyncConnectionState.WaitingForNetwork else LiveSyncConnectionState.Stopped
updateState(if (desiredActive) LiveSyncConnectionState.WaitingForNetwork else LiveSyncConnectionState.Stopped)
false
} else {
desiredActive && state == LiveSyncConnectionState.WaitingForNetwork
}
}
if (available) {
reportRealtimeStatus("reconnecting", "network_available")
} else {
reportRealtimeStatus("disconnected", "network_lost")
}
if (shouldReconnect) connectNow(resetAttempt = true)
}
fun setForeground(active: Boolean) {
val shouldConnect = synchronized(lock) {
foregroundActive = active
if (!active) {
reconnectJob?.cancel()
reconnectJob = null
closeSocketLocked("app_background")
updateState(LiveSyncConnectionState.Stopped)
false
} else {
desiredActive && networkAvailable
}
}
if (shouldConnect) connectNow(resetAttempt = true)
}
fun stop() {
synchronized(lock) {
desiredActive = false
@@ -162,7 +219,7 @@ class DriverLiveSyncClient(
realtimeConfig = null
socketId = null
reconnectAttempt = 0
state = LiveSyncConnectionState.Stopped
updateState(LiveSyncConnectionState.Stopped)
}
reportRealtimeStatus("disconnected", "client_stop")
}
@@ -176,14 +233,14 @@ class DriverLiveSyncClient(
val wsUrl = synchronized(lock) {
val config = realtimeConfig ?: return
val id = driverId ?: return
if (!desiredActive || !networkAvailable || !isConfigUsable(config)) return
if (!desiredActive || !foregroundActive || !networkAvailable || !isConfigUsable(config)) return
if (state == LiveSyncConnectionState.Connecting || state == LiveSyncConnectionState.Subscribing || state == LiveSyncConnectionState.Connected) return
if (resetAttempt) reconnectAttempt = 0
reconnectJob?.cancel()
reconnectJob = null
socketId = null
state = LiveSyncConnectionState.Connecting
updateState(LiveSyncConnectionState.Connecting)
val appKey = config.reverbAppKey.orEmpty()
val wsBaseUrl = config.reverbWsBaseUrl.orEmpty()
@@ -227,6 +284,7 @@ class DriverLiveSyncClient(
false
} else {
messageVersion += 1
lastMessageAtEpochMillis = System.currentTimeMillis()
true
}
}
@@ -244,7 +302,7 @@ class DriverLiveSyncClient(
val nextSocketId = root.dataObject()?.string("socket_id") ?: return
synchronized(lock) {
socketId = nextSocketId
state = LiveSyncConnectionState.Subscribing
updateState(LiveSyncConnectionState.Subscribing)
}
subscribe(socket, nextSocketId)
}
@@ -254,7 +312,7 @@ class DriverLiveSyncClient(
if (channel == "private-driver-mobile.$id") {
synchronized(lock) {
reconnectAttempt = 0
state = LiveSyncConnectionState.Connected
updateState(LiveSyncConnectionState.Connected)
}
reportRealtimeStatus("connected")
startHeartbeat()
@@ -290,6 +348,10 @@ class DriverLiveSyncClient(
}
}.onFailure { throwable ->
AppDiagnostics.log("realtime_subscription_error: ${throwable.message ?: throwable::class.java.simpleName}")
if (ApiErrorMapper.map(throwable).kind == ApiErrorKind.Auth) {
handleAuthenticationFailure(socket)
return@onFailure
}
socket.close(1000, "subscription_failed")
handleDisconnect("error", "subscription_failed", socket)
}
@@ -304,15 +366,16 @@ class DriverLiveSyncClient(
stopStaleWatchdogLocked()
webSocket = null
socketId = null
lastDisconnectReason = reason
if (!desiredActive) {
state = LiveSyncConnectionState.Stopped
updateState(LiveSyncConnectionState.Stopped)
false
} else if (!networkAvailable) {
state = LiveSyncConnectionState.WaitingForNetwork
updateState(LiveSyncConnectionState.WaitingForNetwork)
false
} else {
state = LiveSyncConnectionState.Disconnected
updateState(LiveSyncConnectionState.Disconnected)
true
}
}
@@ -323,7 +386,7 @@ class DriverLiveSyncClient(
private fun scheduleReconnect() {
val delayMs = synchronized(lock) {
if (!desiredActive || !networkAvailable || state == LiveSyncConnectionState.Stopped) return
if (!desiredActive || !foregroundActive || !networkAvailable || state == LiveSyncConnectionState.Stopped) return
if (reconnectJob?.isActive == true) return
val delay = reconnectDelaysMs.getOrElse(reconnectAttempt) { reconnectDelaysMs.last() }
reconnectAttempt += 1
@@ -401,10 +464,26 @@ class DriverLiveSyncClient(
scope.launch {
runCatching {
gateway.storeRealtimeStatus(status, socketId, error)
}.onFailure { throwable ->
if (ApiErrorMapper.map(throwable).kind == ApiErrorKind.Auth) handleAuthenticationFailure()
}
}
}
private fun handleAuthenticationFailure(socket: WebSocket? = null) {
val notify = synchronized(lock) {
if (socket != null && webSocket !== socket) return
val wasActive = desiredActive
desiredActive = false
reconnectJob?.cancel()
reconnectJob = null
closeSocketLocked("authentication_failed")
updateState(LiveSyncConnectionState.Stopped)
wasActive
}
if (notify) onAuthenticationFailed()
}
private fun isConfigUsable(config: RealtimeConfigDto?): Boolean =
config?.reverbEnabled == true &&
!config.reverbAppKey.isNullOrBlank() &&
@@ -0,0 +1,8 @@
package pl.firmatpp.kierowca.sync
object DriverRuntimeSyncState {
@Volatile var foreground: Boolean = false
@Volatile var realtimeConnected: Boolean = false
fun websocketOwnsForegroundHints(): Boolean = foreground && realtimeConnected
}
@@ -14,6 +14,8 @@ import java.util.concurrent.TimeUnit
import pl.firmatpp.kierowca.data.ApiErrorKind
import pl.firmatpp.kierowca.data.ApiErrorMapper
import pl.firmatpp.kierowca.data.sync.DriverSyncRepository
import pl.firmatpp.kierowca.data.sync.NetworkMonitor
import pl.firmatpp.kierowca.data.model.SyncScopeDto
import pl.firmatpp.kierowca.diagnostics.AppDiagnostics
class DriverSyncWorker(
@@ -23,9 +25,20 @@ class DriverSyncWorker(
private val syncRepository = DriverSyncRepository(appContext)
override suspend fun doWork(): Result =
if (!NetworkMonitor(applicationContext).isCurrentlyValidated()) Result.retry() else
runCatching {
val date = inputData.getString(KEY_DATE)
val routeId = inputData.getString(KEY_ROUTE_ID)
val hintedScope = inputData.getString(KEY_SCOPE)?.let { scope ->
val version = inputData.getLong(KEY_VERSION, -1L)
val checksum = inputData.getString(KEY_CHECKSUM)
if (version >= 0L && !checksum.isNullOrBlank()) {
SyncScopeDto(scope, date, routeId, version, checksum, computedAt = "")
} else null
}
if (hintedScope != null && !syncRepository.shouldRefresh(hintedScope)) {
return@runCatching Result.success()
}
val response = syncRepository.fetchSyncState(date, routeId)
var refreshed = false
@@ -85,17 +98,33 @@ class DriverSyncWorker(
companion object {
private const val KEY_DATE = "date"
private const val KEY_ROUTE_ID = "routeId"
private const val KEY_SCOPE = "scope"
private const val KEY_VERSION = "version"
private const val KEY_CHECKSUM = "checksum"
fun enqueue(context: Context, date: String?, routeId: String?) {
fun enqueue(
context: Context,
date: String?,
routeId: String?,
scope: String? = null,
version: Long? = null,
checksum: String? = null,
) {
val request = OneTimeWorkRequestBuilder<DriverSyncWorker>()
.setInputData(workDataOf(KEY_DATE to date, KEY_ROUTE_ID to routeId))
.setInputData(workDataOf(
KEY_DATE to date,
KEY_ROUTE_ID to routeId,
KEY_SCOPE to scope,
KEY_VERSION to version,
KEY_CHECKSUM to checksum,
))
.setConstraints(Constraints.Builder().setRequiredNetworkType(NetworkType.CONNECTED).build())
.setBackoffCriteria(BackoffPolicy.EXPONENTIAL, 30, TimeUnit.SECONDS)
.build()
WorkManager.getInstance(context).enqueueUniqueWork(
listOf("driver-sync", date ?: "-", routeId ?: "-").joinToString("-"),
ExistingWorkPolicy.REPLACE,
listOf("driver-sync", scope ?: "all", date ?: "-", routeId ?: "-").joinToString("-"),
ExistingWorkPolicy.KEEP,
request,
)
}
@@ -26,6 +26,7 @@ import pl.firmatpp.kierowca.R
import pl.firmatpp.kierowca.data.ApiErrorKind
import pl.firmatpp.kierowca.data.ApiErrorMapper
import pl.firmatpp.kierowca.data.DriverRepository
import pl.firmatpp.kierowca.data.sync.NetworkMonitor
import pl.firmatpp.kierowca.diagnostics.AppDiagnostics
import retrofit2.HttpException
@@ -36,6 +37,7 @@ class NewRouteNotificationWorker(
private val repository = DriverRepository(appContext)
override suspend fun doWork(): Result {
if (!NetworkMonitor(applicationContext).isCurrentlyValidated()) return Result.retry()
val routeId = inputData.getString(KEY_ROUTE_ID)?.takeIf(String::isNotBlank) ?: return Result.success()
return runCatching {