fix(android): reconcile rejected runtime state and sync admission #377
This commit is contained in:
@@ -278,6 +278,12 @@ class PurchaseTaskStore(context: Context) : SQLiteOpenHelper(context, DATABASE_N
|
||||
"SELECT COUNT(*) FROM purchase_outbox WHERE upload_status='rejected'", null,
|
||||
).use { it.moveToFirst(); it.getInt(0) }
|
||||
|
||||
@Synchronized
|
||||
fun isCurrentAttemptRejected(taskId: Long, attemptId: String): Boolean = readableDatabase.rawQuery(
|
||||
"SELECT 1 FROM purchase_task WHERE task_id=? AND attempt_id=? AND upload_status='rejected'",
|
||||
arrayOf(taskId.toString(), attemptId),
|
||||
).use { it.moveToFirst() }
|
||||
|
||||
@Synchronized
|
||||
fun acknowledgeRejection(outboxId: Long) {
|
||||
writableDatabase.execSQL(PurchaseRejectionSql.acknowledge, arrayOf(System.currentTimeMillis(), outboxId))
|
||||
|
||||
@@ -204,6 +204,7 @@ class AgentForegroundService : Service() {
|
||||
}
|
||||
|
||||
override fun onDestroy() {
|
||||
stateStore.stop()
|
||||
probeHandoff.invalidate()
|
||||
if (diagnosticInstance === this) diagnosticInstance = null
|
||||
repurchaseClosed.set(true)
|
||||
@@ -621,7 +622,7 @@ class AgentForegroundService : Service() {
|
||||
publishCurrentPageResult(0L, "采集间隔中,还需 $remaining 秒。", replacementOriginType, replacementOriginTaskId)
|
||||
return
|
||||
}
|
||||
if (!taskMutex.tryAcquire(CURRENT_PAGE_RESERVATION_ID)) {
|
||||
if (!tryAcquireCurrentPage(taskMutex, working, CURRENT_PAGE_RESERVATION_ID)) {
|
||||
publishCurrentPageResult(0L, "设备正在执行任务,请稍后再试。", replacementOriginType, replacementOriginTaskId)
|
||||
return
|
||||
}
|
||||
@@ -1085,6 +1086,13 @@ class AgentForegroundService : Service() {
|
||||
markRejected = { item, code ->
|
||||
probeHandoff.invalidate()
|
||||
purchaseStore.markRejected(item, code)
|
||||
synchronized(taskMutex) {
|
||||
val runtime = stateStore.read()
|
||||
if (shouldClearRejectedPurchase(item.taskId, runtime.currentTaskId, runtime.currentTaskType,
|
||||
taskMutex.currentTaskId(), purchaseStore.isCurrentAttemptRejected(item.taskId, item.attemptId))) {
|
||||
stateStore.clearActiveTask(item.taskId)
|
||||
}
|
||||
}
|
||||
sendBroadcast(Intent(ACTION_PURCHASE_REJECTIONS_CHANGED).setPackage(packageName))
|
||||
},
|
||||
).flush()
|
||||
|
||||
@@ -191,8 +191,9 @@ class AgentSettingsStore(context: Context) {
|
||||
|
||||
class AgentStateStore(context: Context) {
|
||||
private val preferences = context.getSharedPreferences(PREFERENCES, Context.MODE_PRIVATE)
|
||||
private val publicationGate = AgentStatePublicationGate()
|
||||
|
||||
fun update(code: String, message: String, deviceId: Long? = null, tokenStored: Boolean? = null, heartbeat: Boolean = false, connection: Boolean = false) {
|
||||
fun update(code: String, message: String, deviceId: Long? = null, tokenStored: Boolean? = null, heartbeat: Boolean = false, connection: Boolean = false) = publicationGate.publish {
|
||||
val editor = preferences.edit()
|
||||
.putString(STATE_CODE, code)
|
||||
.putString(STATE_MESSAGE, message)
|
||||
@@ -203,6 +204,12 @@ class AgentStateStore(context: Context) {
|
||||
editor.apply()
|
||||
}
|
||||
|
||||
/** This service instance cannot publish a late heartbeat over STOPPED or a restarted service. */
|
||||
fun stop() = publicationGate.stop {
|
||||
preferences.edit().putString(STATE_CODE, "STOPPED").putString(CONNECTION_CODE, "STOPPED")
|
||||
.putString(STATE_MESSAGE, "前台服务已停止").apply()
|
||||
}
|
||||
|
||||
fun connectionCode(): String = preferences.getString(CONNECTION_CODE, "STOPPED") ?: "STOPPED"
|
||||
|
||||
fun read(): AgentState = AgentState(
|
||||
|
||||
+16
@@ -0,0 +1,16 @@
|
||||
package cn.ilapage.goauto.agent.service
|
||||
|
||||
internal fun shouldClearRejectedPurchase(
|
||||
rejectedTaskId: Long, runtimeTaskId: Long?, runtimeTaskType: String?,
|
||||
executingTaskId: Long?, currentAttemptRejected: Boolean,
|
||||
): Boolean = executingTaskId == null && runtimeTaskId == rejectedTaskId &&
|
||||
runtimeTaskType == "purchase" && currentAttemptRejected
|
||||
|
||||
internal class AgentStatePublicationGate {
|
||||
private var stopped = false
|
||||
@Synchronized fun publish(write: () -> Unit) { if (!stopped) write() }
|
||||
@Synchronized fun stop(write: () -> Unit) { stopped = true; write() }
|
||||
}
|
||||
|
||||
internal fun tryAcquireCurrentPage(mutex: TaskExecutionMutex, syncing: java.util.concurrent.atomic.AtomicBoolean, reservation: Long): Boolean =
|
||||
synchronized(mutex) { !syncing.get() && mutex.tryAcquire(reservation) }
|
||||
+59
@@ -0,0 +1,59 @@
|
||||
package cn.ilapage.goauto.agent.service
|
||||
|
||||
import org.junit.Assert.*
|
||||
import org.junit.Test
|
||||
|
||||
class PurchaseRuntimeReconciliationTest {
|
||||
@Test fun `manual current page cannot enter while sync owns recovery and flush`() {
|
||||
val mutex = TaskExecutionMutex()
|
||||
val syncing = java.util.concurrent.atomic.AtomicBoolean(false)
|
||||
val syncEntered = java.util.concurrent.CountDownLatch(1)
|
||||
val finishSync = java.util.concurrent.CountDownLatch(1)
|
||||
val worker = java.util.concurrent.Executors.newSingleThreadExecutor()
|
||||
val future = worker.submit {
|
||||
synchronized(mutex) { assertTrue(syncing.compareAndSet(false, true)) }
|
||||
syncEntered.countDown()
|
||||
check(finishSync.await(3, java.util.concurrent.TimeUnit.SECONDS))
|
||||
syncing.set(false)
|
||||
}
|
||||
try {
|
||||
assertTrue(syncEntered.await(3, java.util.concurrent.TimeUnit.SECONDS))
|
||||
assertFalse(tryAcquireCurrentPage(mutex, syncing, Long.MAX_VALUE))
|
||||
assertNull(mutex.currentTaskId())
|
||||
finishSync.countDown()
|
||||
future.get(3, java.util.concurrent.TimeUnit.SECONDS)
|
||||
assertTrue(tryAcquireCurrentPage(mutex, syncing, Long.MAX_VALUE))
|
||||
} finally { finishSync.countDown(); worker.shutdownNow() }
|
||||
}
|
||||
|
||||
@Test fun `reboot restored rejected task clears runtime so repurchase becomes idle`() {
|
||||
var runtimeTaskId: Long? = 42
|
||||
val currentDatabaseAttemptRejected = true
|
||||
if (shouldClearRejectedPurchase(42, runtimeTaskId, "purchase", null, currentDatabaseAttemptRejected)) runtimeTaskId = null
|
||||
assertNull(runtimeTaskId)
|
||||
assertTrue(runtimeTaskId == null) // final runtime predicate in repurchaseLocalIdle
|
||||
}
|
||||
|
||||
@Test fun `another task collection active executor or newer attempt cannot be cleared`() {
|
||||
assertFalse(shouldClearRejectedPurchase(42, 43, "purchase", null, true))
|
||||
assertFalse(shouldClearRejectedPurchase(42, 42, "collection", null, true))
|
||||
assertFalse(shouldClearRejectedPurchase(42, 42, "purchase", 42, true))
|
||||
assertFalse(shouldClearRejectedPurchase(42, 42, "purchase", 43, true))
|
||||
assertFalse(shouldClearRejectedPurchase(42, 42, "purchase", null, false))
|
||||
}
|
||||
|
||||
@Test fun `stopped service remains stopped after an in flight heartbeat completes`() {
|
||||
val gate = AgentStatePublicationGate()
|
||||
var connection = "CONNECTING"
|
||||
gate.publish { connection = "ONLINE" }
|
||||
assertEquals("ONLINE", connection)
|
||||
gate.stop { connection = "STOPPED" }
|
||||
gate.publish { connection = "ONLINE" }
|
||||
assertEquals("STOPPED", connection)
|
||||
val restartedService = AgentStatePublicationGate()
|
||||
restartedService.publish { connection = "CONNECTING" }
|
||||
assertEquals("CONNECTING", connection)
|
||||
gate.publish { connection = "ONLINE" }
|
||||
assertEquals("CONNECTING", connection)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user