fix(android): record handled purchase recovery failures in sync diagnostics (#378)
This commit is contained in:
+19
-1
@@ -136,6 +136,7 @@ internal class SyncDiagnosticRecorder(
|
||||
private var persistenceFailed = id == null
|
||||
private var aborted = false
|
||||
private var completed = false
|
||||
private var activeStep: Pair<JSONObject, Long>? = null
|
||||
|
||||
fun <T> step(name: String, action: () -> T): T {
|
||||
require(name in setOf("register", "heartbeat", "recover", "flush", "claim"))
|
||||
@@ -148,10 +149,11 @@ internal class SyncDiagnosticRecorder(
|
||||
.put("startedOffsetMs", (begin - start).coerceAtLeast(0))
|
||||
.put("status", "running")
|
||||
entries.put(entry)
|
||||
activeStep = entry to begin
|
||||
persistSteps()
|
||||
try {
|
||||
val value = action()
|
||||
entry.put("status", "success")
|
||||
if (entry.getString("status") != "failure") entry.put("status", "success")
|
||||
return value
|
||||
} catch (e: CancellationException) {
|
||||
aborted = true
|
||||
@@ -171,9 +173,25 @@ internal class SyncDiagnosticRecorder(
|
||||
.put("durationMs", (elapsed() - begin).coerceAtLeast(0))
|
||||
persistSteps()
|
||||
}
|
||||
activeStep = null
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Observe a business-handled exception without changing that operation's return/throw
|
||||
* policy.
|
||||
*/
|
||||
fun handledFailure(error: Exception) {
|
||||
if (error is CancellationException) {
|
||||
aborted = true
|
||||
throw error
|
||||
}
|
||||
val (entry, begin) = activeStep ?: return failed(error)
|
||||
failure = syncFailure(error, entry.getString("step"), wall(), elapsed() - begin)
|
||||
entry.put("status", "failure").put("category", failure!!.category)
|
||||
persistSteps()
|
||||
}
|
||||
|
||||
fun failed(error: Exception) {
|
||||
if (error is CancellationException) {
|
||||
aborted = true
|
||||
|
||||
+11
-7
@@ -373,7 +373,11 @@ class AgentForegroundService : Service() {
|
||||
else purchaseStore.activeTaskId().also(runningTaskId::set)
|
||||
},
|
||||
executionActive = { taskMutex.currentTaskId() != null },
|
||||
recover = { round.step("recover") { recoverInterruptedPurchases(api, activeCredentials.token) } },
|
||||
recover = {
|
||||
round.step("recover") {
|
||||
recoverInterruptedPurchases(api, activeCredentials.token, round::handledFailure)
|
||||
}
|
||||
},
|
||||
flush = { round.step("flush") { flushPurchaseOutbox(api, activeCredentials.token) } },
|
||||
hasPending = { purchaseStore.pendingOutbox().isNotEmpty() },
|
||||
).run()
|
||||
@@ -1091,12 +1095,16 @@ class AgentForegroundService : Service() {
|
||||
.toString()
|
||||
}
|
||||
|
||||
private fun recoverInterruptedPurchases(api: AgentApiClient, token: String) {
|
||||
private fun recoverInterruptedPurchases(
|
||||
api: AgentApiClient,
|
||||
token: String,
|
||||
onHandledFailure: (AgentApiException) -> Unit,
|
||||
) {
|
||||
purchaseStore.interruptedAttempts().forEach { interrupted ->
|
||||
probeHandoff.beginExecution()
|
||||
if (interrupted.status == PurchaseTaskStore.STATUS_ORDER_SUBMIT_STARTED) {
|
||||
val boundaryRequestId = interrupted.orderSubmitRequestId ?: return@forEach
|
||||
try {
|
||||
recoverPurchaseBoundary(onHandledFailure) {
|
||||
api.markPurchaseOrderSubmitStarted(interrupted.taskId, boundaryRequestId, token)
|
||||
val automation = GoAutoAccessibilityService.instance?.let(::PurchaseLiveAutomation)
|
||||
val evidence = try { automation?.readOrderResult() } catch (error: Exception) {
|
||||
@@ -1122,10 +1130,6 @@ class AgentForegroundService : Service() {
|
||||
}
|
||||
val requestId = UUID.randomUUID().toString()
|
||||
purchaseStore.completeAndEnqueue(interrupted.taskId, interrupted.attemptId, requestId, purchaseResultPayload(requestId, interrupted.attemptId, outcome))
|
||||
} catch (error: AgentApiException) {
|
||||
// Keep the local irreversible marker. The same server request ID
|
||||
// is replayed after connectivity returns; the order click is never repeated.
|
||||
if (isAgentAuthenticationError(error)) throw error
|
||||
}
|
||||
return@forEach
|
||||
}
|
||||
|
||||
@@ -72,3 +72,16 @@ private fun <T> captureSync(action: () -> T): Result<T> = try {
|
||||
} catch (error: Exception) {
|
||||
Result.failure(error)
|
||||
}
|
||||
|
||||
/** Preserve the local irreversible marker when replay fails; the original request ID is reused. */
|
||||
internal fun recoverPurchaseBoundary(
|
||||
onHandledFailure: (AgentApiException) -> Unit,
|
||||
recover: () -> Unit,
|
||||
) {
|
||||
try {
|
||||
recover()
|
||||
} catch (error: AgentApiException) {
|
||||
if (isAgentAuthenticationError(error)) throw error
|
||||
onHandledFailure(error)
|
||||
}
|
||||
}
|
||||
|
||||
+106
@@ -0,0 +1,106 @@
|
||||
package cn.ilapage.goauto.agent.persistence
|
||||
|
||||
import cn.ilapage.goauto.agent.network.AgentApiException
|
||||
import cn.ilapage.goauto.agent.network.HeartbeatResult
|
||||
import cn.ilapage.goauto.agent.service.AgentSyncCycle
|
||||
import cn.ilapage.goauto.agent.service.recoverPurchaseBoundary
|
||||
import java.sql.DriverManager
|
||||
import java.util.concurrent.CancellationException
|
||||
import org.json.JSONArray
|
||||
import org.junit.Assert.*
|
||||
import org.junit.Test
|
||||
|
||||
class HandledRecoveryDiagnosticTest {
|
||||
@Test
|
||||
fun handledBoundaryApiFailureIsPersistedButDoesNotChangeRecoveryOrClaimPolicy() {
|
||||
for (status in listOf(500, 409)) {
|
||||
for (pending in listOf(false, true)) {
|
||||
DriverManager.getConnection("jdbc:sqlite::memory:").use { db ->
|
||||
val sql = JdbcSyncSql(db)
|
||||
AgentDiagnosticSchema.syncStatements.forEach { sql.execute(it) }
|
||||
val store = SyncDiagnosticStore(sql)
|
||||
var elapsed = 10L
|
||||
val round =
|
||||
SyncDiagnosticRecorder(store, "p", { 1000 }, { elapsed })
|
||||
.begin(SyncContext())
|
||||
val events = mutableListOf<String>()
|
||||
var originalMarkerPresent = true
|
||||
val replayBoundary: () -> Unit = {
|
||||
events += "recover"
|
||||
elapsed = 45
|
||||
throw AgentApiException(
|
||||
status,
|
||||
"PRIVATE_CODE",
|
||||
"private URL and payload",
|
||||
false,
|
||||
)
|
||||
}
|
||||
val cycle =
|
||||
AgentSyncCycle(
|
||||
heartbeat = {
|
||||
round.step("heartbeat") {
|
||||
events += "heartbeat"
|
||||
HeartbeatResult(1, true, false, 15)
|
||||
}
|
||||
},
|
||||
activeTaskId = { null },
|
||||
executionActive = { false },
|
||||
recover = {
|
||||
round.step("recover") {
|
||||
recoverPurchaseBoundary(round::handledFailure) {
|
||||
replayBoundary()
|
||||
originalMarkerPresent = false
|
||||
}
|
||||
}
|
||||
},
|
||||
flush = { round.step("flush") { events += "flush" } },
|
||||
hasPending = { pending },
|
||||
)
|
||||
.run()
|
||||
round.finish()
|
||||
assertEquals(listOf("heartbeat", "recover", "flush"), events)
|
||||
assertTrue(originalMarkerPresent)
|
||||
assertNull(cycle.failureCode)
|
||||
assertEquals(!pending, cycle.canClaim)
|
||||
val row = sql.query("SELECT * FROM sync_round").single()
|
||||
assertEquals("failure", row["status"])
|
||||
assertEquals("recover", row["failure_step"])
|
||||
assertEquals(
|
||||
if (status == 500) "api_server" else "api_conflict",
|
||||
row["failure_category"],
|
||||
)
|
||||
assertEquals("35", row["failure_duration_ms"])
|
||||
assertEquals(status.toString(), row["http_status"])
|
||||
assertNull(row["api_code"])
|
||||
val steps = JSONArray(row["steps"])
|
||||
assertEquals("failure", steps.getJSONObject(2).getString("status"))
|
||||
assertEquals(35, steps.getJSONObject(2).getLong("durationMs"))
|
||||
assertEquals("success", steps.getJSONObject(3).getString("status"))
|
||||
assertFalse(row.toString().contains("private"))
|
||||
assertEquals(1, store.snapshot("p", 1100)!!.report.getInt("failedRounds"))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun boundaryPolicyStillRethrowsAuthenticationCancellationAndFatalErrors() {
|
||||
for (error in
|
||||
listOf(
|
||||
AgentApiException(401, "OTHER", "private", false),
|
||||
AgentApiException(409, "DEVICE_DISABLED", "private", false),
|
||||
CancellationException(),
|
||||
OutOfMemoryError(),
|
||||
)) {
|
||||
val thrown =
|
||||
assertThrows(error.javaClass) {
|
||||
recoverPurchaseBoundary({
|
||||
fail("only handled nonauthentication API errors are observed")
|
||||
}) {
|
||||
throw error
|
||||
}
|
||||
}
|
||||
assertSame(error, thrown)
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user