diff --git a/android/app/src/main/java/cn/ilapage/goauto/agent/persistence/SyncDiagnosticRecorder.kt b/android/app/src/main/java/cn/ilapage/goauto/agent/persistence/SyncDiagnosticRecorder.kt index 54f2f39..c279c96 100644 --- a/android/app/src/main/java/cn/ilapage/goauto/agent/persistence/SyncDiagnosticRecorder.kt +++ b/android/app/src/main/java/cn/ilapage/goauto/agent/persistence/SyncDiagnosticRecorder.kt @@ -136,6 +136,7 @@ internal class SyncDiagnosticRecorder( private var persistenceFailed = id == null private var aborted = false private var completed = false + private var activeStep: Pair? = null fun 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 diff --git a/android/app/src/main/java/cn/ilapage/goauto/agent/service/AgentForegroundService.kt b/android/app/src/main/java/cn/ilapage/goauto/agent/service/AgentForegroundService.kt index 06a1725..a301ad2 100644 --- a/android/app/src/main/java/cn/ilapage/goauto/agent/service/AgentForegroundService.kt +++ b/android/app/src/main/java/cn/ilapage/goauto/agent/service/AgentForegroundService.kt @@ -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 } diff --git a/android/app/src/main/java/cn/ilapage/goauto/agent/service/AgentSyncCycle.kt b/android/app/src/main/java/cn/ilapage/goauto/agent/service/AgentSyncCycle.kt index a5c3fe2..81f1a66 100644 --- a/android/app/src/main/java/cn/ilapage/goauto/agent/service/AgentSyncCycle.kt +++ b/android/app/src/main/java/cn/ilapage/goauto/agent/service/AgentSyncCycle.kt @@ -72,3 +72,16 @@ private fun captureSync(action: () -> T): Result = 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) + } +} diff --git a/android/app/src/test/java/cn/ilapage/goauto/agent/persistence/HandledRecoveryDiagnosticTest.kt b/android/app/src/test/java/cn/ilapage/goauto/agent/persistence/HandledRecoveryDiagnosticTest.kt new file mode 100644 index 0000000..7ac09b7 --- /dev/null +++ b/android/app/src/test/java/cn/ilapage/goauto/agent/persistence/HandledRecoveryDiagnosticTest.kt @@ -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() + 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) + } + } +}