Dodaj live sync i kolejkę zdjęć
This commit is contained in:
@@ -0,0 +1,33 @@
|
||||
package pl.firmatpp.kierowca.sync
|
||||
|
||||
import com.google.firebase.messaging.FirebaseMessagingService
|
||||
import com.google.firebase.messaging.RemoteMessage
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.launch
|
||||
import pl.firmatpp.kierowca.data.DriverRepository
|
||||
|
||||
class DriverFirebaseMessagingService : FirebaseMessagingService() {
|
||||
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
|
||||
|
||||
override fun onNewToken(token: String) {
|
||||
scope.launch {
|
||||
val repository = DriverRepository(applicationContext)
|
||||
if (repository.hasToken()) {
|
||||
runCatching { repository.storePushToken(token) }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
override fun onMessageReceived(message: RemoteMessage) {
|
||||
val data = message.data
|
||||
if (data["type"] != "driver_sync_hint") return
|
||||
|
||||
DriverSyncWorker.enqueue(
|
||||
context = applicationContext,
|
||||
date = data["date"],
|
||||
routeId = data["routeId"],
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,135 @@
|
||||
package pl.firmatpp.kierowca.sync
|
||||
|
||||
import com.google.gson.Gson
|
||||
import com.google.gson.JsonObject
|
||||
import com.google.gson.JsonParser
|
||||
import java.util.concurrent.atomic.AtomicBoolean
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.launch
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.WebSocket
|
||||
import okhttp3.WebSocketListener
|
||||
import pl.firmatpp.kierowca.BuildConfig
|
||||
import pl.firmatpp.kierowca.data.DriverRepository
|
||||
|
||||
class DriverLiveSyncClient(
|
||||
private val repository: DriverRepository,
|
||||
private val onConnected: () -> Unit,
|
||||
private val onHint: (DriverSyncHint) -> Unit,
|
||||
private val client: OkHttpClient = OkHttpClient(),
|
||||
private val gson: Gson = Gson(),
|
||||
) {
|
||||
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
|
||||
private val started = AtomicBoolean(false)
|
||||
private var webSocket: WebSocket? = null
|
||||
private var driverId: String? = null
|
||||
|
||||
fun start(driverId: String) {
|
||||
if (BuildConfig.REVERB_APP_KEY.isBlank()) return
|
||||
this.driverId = driverId
|
||||
if (!started.compareAndSet(false, true)) return
|
||||
|
||||
val wsUrl = BuildConfig.REVERB_WS_BASE_URL.trimEnd('/') +
|
||||
"/" + BuildConfig.REVERB_APP_KEY +
|
||||
"?protocol=7&client=android&version=1.0&flash=false"
|
||||
|
||||
webSocket = client.newWebSocket(
|
||||
Request.Builder().url(wsUrl).build(),
|
||||
object : WebSocketListener() {
|
||||
override fun onMessage(webSocket: WebSocket, text: String) {
|
||||
handleMessage(webSocket, text)
|
||||
}
|
||||
|
||||
override fun onClosed(webSocket: WebSocket, code: Int, reason: String) {
|
||||
started.set(false)
|
||||
scheduleReconnect()
|
||||
}
|
||||
|
||||
override fun onFailure(webSocket: WebSocket, t: Throwable, response: okhttp3.Response?) {
|
||||
started.set(false)
|
||||
scheduleReconnect()
|
||||
}
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
fun stop() {
|
||||
started.set(false)
|
||||
webSocket?.close(1000, "logout")
|
||||
webSocket = null
|
||||
driverId = null
|
||||
}
|
||||
|
||||
fun close() {
|
||||
stop()
|
||||
scope.cancel()
|
||||
}
|
||||
|
||||
private fun handleMessage(socket: WebSocket, text: String) {
|
||||
val root = runCatching { JsonParser.parseString(text).asJsonObject }.getOrNull() ?: return
|
||||
val event = root.string("event") ?: return
|
||||
|
||||
when (event) {
|
||||
"pusher:connection_established" -> {
|
||||
val socketId = root.dataObject()?.string("socket_id") ?: return
|
||||
subscribe(socket, socketId)
|
||||
onConnected()
|
||||
}
|
||||
"DriverMobileSyncHint" -> parseHint(root.dataObject())?.let(onHint)
|
||||
}
|
||||
}
|
||||
|
||||
private fun subscribe(socket: WebSocket, socketId: String) {
|
||||
val id = driverId ?: return
|
||||
val channel = "private-driver-mobile.$id"
|
||||
|
||||
scope.launch {
|
||||
runCatching {
|
||||
val auth = repository.broadcastAuth(socketId, channel).auth
|
||||
val payload = mapOf(
|
||||
"event" to "pusher:subscribe",
|
||||
"data" to mapOf(
|
||||
"channel" to channel,
|
||||
"auth" to auth,
|
||||
),
|
||||
)
|
||||
socket.send(gson.toJson(payload))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun parseHint(data: JsonObject?): DriverSyncHint? {
|
||||
if (data == null || data.string("type") != "driver_sync_hint") return null
|
||||
|
||||
return DriverSyncHint(
|
||||
scope = data.string("scope") ?: return null,
|
||||
date = data.string("date"),
|
||||
routeId = data.string("routeId"),
|
||||
version = data.string("version")?.toLongOrNull() ?: data.get("version")?.asLong ?: return null,
|
||||
checksum = data.string("checksum") ?: return null,
|
||||
)
|
||||
}
|
||||
|
||||
private fun JsonObject.dataObject(): JsonObject? {
|
||||
val data = get("data") ?: return null
|
||||
return if (data.isJsonObject) data.asJsonObject else runCatching { JsonParser.parseString(data.asString).asJsonObject }.getOrNull()
|
||||
}
|
||||
|
||||
private fun JsonObject.string(name: String): String? =
|
||||
get(name)?.takeIf { !it.isJsonNull }?.asString
|
||||
|
||||
private fun scheduleReconnect() {
|
||||
val id = driverId ?: return
|
||||
scope.launch {
|
||||
delay(5_000)
|
||||
if (!started.get() && driverId == id) {
|
||||
start(id)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
package pl.firmatpp.kierowca.sync
|
||||
|
||||
import pl.firmatpp.kierowca.data.model.SyncScopeDto
|
||||
|
||||
data class DriverSyncHint(
|
||||
val scope: String,
|
||||
val date: String?,
|
||||
val routeId: String?,
|
||||
val version: Long,
|
||||
val checksum: String,
|
||||
) {
|
||||
fun asScope(): SyncScopeDto =
|
||||
SyncScopeDto(
|
||||
scope = scope,
|
||||
date = date,
|
||||
routeId = routeId,
|
||||
version = version,
|
||||
checksum = checksum,
|
||||
computedAt = "",
|
||||
)
|
||||
}
|
||||
@@ -0,0 +1,77 @@
|
||||
package pl.firmatpp.kierowca.sync
|
||||
|
||||
import android.content.Context
|
||||
import androidx.work.BackoffPolicy
|
||||
import androidx.work.Constraints
|
||||
import androidx.work.CoroutineWorker
|
||||
import androidx.work.ExistingWorkPolicy
|
||||
import androidx.work.NetworkType
|
||||
import androidx.work.OneTimeWorkRequestBuilder
|
||||
import androidx.work.WorkManager
|
||||
import androidx.work.WorkerParameters
|
||||
import androidx.work.workDataOf
|
||||
import java.util.concurrent.TimeUnit
|
||||
import pl.firmatpp.kierowca.data.ApiErrorKind
|
||||
import pl.firmatpp.kierowca.data.ApiErrorMapper
|
||||
import pl.firmatpp.kierowca.data.sync.DriverSyncRepository
|
||||
|
||||
class DriverSyncWorker(
|
||||
appContext: Context,
|
||||
params: WorkerParameters,
|
||||
) : CoroutineWorker(appContext, params) {
|
||||
private val syncRepository = DriverSyncRepository(appContext)
|
||||
|
||||
override suspend fun doWork(): Result =
|
||||
runCatching {
|
||||
val date = inputData.getString(KEY_DATE)
|
||||
val routeId = inputData.getString(KEY_ROUTE_ID)
|
||||
val response = syncRepository.fetchSyncState(date, routeId)
|
||||
var refreshed = false
|
||||
|
||||
response.scopes.forEach { scope ->
|
||||
if (syncRepository.shouldRefresh(scope)) {
|
||||
when (scope.scope) {
|
||||
DriverSyncRepository.SCOPE_ROUTES -> {
|
||||
syncRepository.bootstrap(scope.date ?: date)
|
||||
refreshed = true
|
||||
}
|
||||
DriverSyncRepository.SCOPE_ROUTE_DETAIL -> {
|
||||
val id = scope.routeId ?: routeId
|
||||
if (!id.isNullOrBlank()) {
|
||||
syncRepository.route(id)
|
||||
refreshed = true
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (!refreshed) {
|
||||
syncRepository.saveSyncStates(response)
|
||||
}
|
||||
|
||||
Result.success()
|
||||
}.getOrElse { throwable ->
|
||||
val error = ApiErrorMapper.map(throwable)
|
||||
if (error.retryable && error.kind != ApiErrorKind.Auth) Result.retry() else Result.failure()
|
||||
}
|
||||
|
||||
companion object {
|
||||
private const val KEY_DATE = "date"
|
||||
private const val KEY_ROUTE_ID = "routeId"
|
||||
|
||||
fun enqueue(context: Context, date: String?, routeId: String?) {
|
||||
val request = OneTimeWorkRequestBuilder<DriverSyncWorker>()
|
||||
.setInputData(workDataOf(KEY_DATE to date, KEY_ROUTE_ID to routeId))
|
||||
.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,
|
||||
request,
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user