From 5430dc1f64fc0a4d8ef124ada7356e21a975151f Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Thu, 8 Oct 2026 18:12:40 +0800 Subject: [PATCH] feat(android): repurchase last failed batch in place (#367) --- .../goauto/agent/network/AgentApiClient.kt | 7 + .../agent/service/AgentForegroundService.kt | 162 +++++++++++- .../agent/service/RepurchaseCoordinator.kt | 235 ++++++++++++++++++ .../goauto/agent/service/RepurchaseState.kt | 40 +++ .../goauto/agent/ui/TaskHistoryFragment.kt | 88 ++++++- .../goauto/agent/RepurchaseServerClockTest.kt | 37 +++ .../service/RepurchaseCoordinatorTest.kt | 182 ++++++++++++++ .../agent/service/RepurchaseStateTest.kt | 34 +++ 8 files changed, 781 insertions(+), 4 deletions(-) create mode 100644 android/app/src/main/java/cn/ilapage/goauto/agent/service/RepurchaseCoordinator.kt create mode 100644 android/app/src/main/java/cn/ilapage/goauto/agent/service/RepurchaseState.kt create mode 100644 android/app/src/test/java/cn/ilapage/goauto/agent/RepurchaseServerClockTest.kt create mode 100644 android/app/src/test/java/cn/ilapage/goauto/agent/service/RepurchaseCoordinatorTest.kt create mode 100644 android/app/src/test/java/cn/ilapage/goauto/agent/service/RepurchaseStateTest.kt diff --git a/android/app/src/main/java/cn/ilapage/goauto/agent/network/AgentApiClient.kt b/android/app/src/main/java/cn/ilapage/goauto/agent/network/AgentApiClient.kt index c0daa78..6654330 100644 --- a/android/app/src/main/java/cn/ilapage/goauto/agent/network/AgentApiClient.kt +++ b/android/app/src/main/java/cn/ilapage/goauto/agent/network/AgentApiClient.kt @@ -100,6 +100,7 @@ data class PurchaseAgentTask( val ruleSnapshot: String, val ruleSnapshotHash: String, val leaseVersion: Long, + val attemptNumber: Int = 0, ) data class CollectionHistoryItem( @@ -211,6 +212,8 @@ class AgentApiException( ) : Exception(message) class AgentApiClient(private val serverUrl: String) { + @Volatile var lastResponseServerTimeMillis: Long? = null + private set val failureSnapshotOrigin: String get() = ServerUrlPolicy.normalize(serverUrl) /** Stream the two bounded parts; never build another combined copy of the archive. */ fun uploadFailureSnapshot(snapshot: cn.ilapage.goauto.agent.diagnostics.FailureSnapshot, token: String) { @@ -666,6 +669,7 @@ class AgentApiClient(private val serverUrl: String) { private fun purchaseTask(data: JSONObject) = PurchaseAgentTask( taskId = data.getLong("taskId"), taskAttemptId = data.optString("taskAttemptId"), + attemptNumber = data.optInt("attemptNumber"), phase = data.optString("phase"), executionMode = data.getString("executionMode"), status = data.getString("status"), @@ -709,6 +713,9 @@ class AgentApiClient(private val serverUrl: String) { connection.outputStream.use { it.write(payload.toString().toByteArray(Charsets.UTF_8)) } } val status = connection.responseCode + // Metadata only: batch cutoff must use the API clock, never the phone's clock. + // A missing Date invalidates the sample instead of retaining an older response. + lastResponseServerTimeMillis = connection.getHeaderFieldDate("Date", 0L).takeIf { it > 0L } if (status == HttpURLConnection.HTTP_NO_CONTENT) return null val stream = if (status in 200..299) connection.inputStream else connection.errorStream val body = stream?.bufferedReader(Charsets.UTF_8)?.use { it.readText() }.orEmpty() 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 06830ca..a38c531 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 @@ -85,6 +85,11 @@ class AgentForegroundService : Service() { private val taskExecutor: ExecutorService = Executors.newSingleThreadExecutor() private val diagnosticExecutor: ExecutorService = Executors.newSingleThreadExecutor() private val snapshotUploadExecutor: ScheduledExecutorService = Executors.newSingleThreadScheduledExecutor() + private val repurchaseExecutor: ScheduledExecutorService = Executors.newSingleThreadScheduledExecutor() + private val repurchaseClosed = AtomicBoolean(false) + private val repurchase = RepurchaseCoordinator(SystemClock::elapsedRealtime, { UUID.randomUUID().toString() }, ::publishRepurchase) + private var repurchaseSession: RepurchaseSession? = null + private val repurchaseEvidence = java.util.concurrent.ConcurrentHashMap>() private val taskMutex = TaskExecutionMutex() private val backfillGuard = OrderBackfillGuard(taskMutex) private val runningTaskId = AtomicReference(null) @@ -119,6 +124,7 @@ class AgentForegroundService : Service() { identityStore = SecureDeviceStore(this) settingsStore = AgentSettingsStore(this) stateStore = AgentStateStore(this) + repurchaseState = RepurchaseState() // Wake locks and cooldown tickets are process-local. Never restore a stale UI flag. stateStore.setKeepScreenOn(false) purchaseStore = PurchaseTaskStore(this) @@ -145,9 +151,29 @@ class AgentForegroundService : Service() { executor.scheduleWithFixedDelay(::triggerSync, 0, HEARTBEAT_SECONDS, TimeUnit.SECONDS) // Independent retention/upload ticks run even when no further task is dispatched. snapshotUploadExecutor.scheduleWithFixedDelay(::maintainFailureSnapshots, 0, 60, TimeUnit.SECONDS) + // Separate, serialized I/O ticks keep history reads and result observation off the + // heartbeat scheduler and, especially, off the single purchase execution thread. + repurchaseExecutor.scheduleWithFixedDelay(::tickRepurchase, 2, 2, TimeUnit.SECONDS) } override fun onStartCommand(intent: Intent?, flags: Int, startId: Int): Int { + when (intent?.action) { + ACTION_REPURCHASE_PREPARE -> { + repurchaseExecutor.execute(::prepareRepurchase) + return START_STICKY + } + ACTION_REPURCHASE_CONFIRM -> { + val roundId = intent.getStringExtra("roundId").orEmpty() + repurchaseExecutor.execute { repurchase.confirm(roundId); tickRepurchase() } + return START_STICKY + } + ACTION_REPURCHASE_STOP -> { + // Atomic: a network request must not delay the user's stop signal. + repurchase.requestStop(intent.getStringExtra("roundId").orEmpty()) + repurchaseExecutor.execute(::tickRepurchase) + return START_STICKY + } + } if (intent?.action == ACTION_BACKFILL_STOP) { backfillGuard.cancelled.set(true) return START_STICKY @@ -171,6 +197,10 @@ class AgentForegroundService : Service() { } override fun onDestroy() { + repurchaseClosed.set(true) + repurchase.requestStop(repurchase.state.roundId) + repurchaseExecutor.shutdownNow() + repurchaseState = RepurchaseState(message = "服务已停止;不会恢复整轮重购") backfillGuard.cancelled.set(true) runCatching { connectivityManager.unregisterNetworkCallback(networkCallback) } cancelIdleReturn("服务已停止") @@ -333,6 +363,110 @@ class AgentForegroundService : Service() { }) } + private data class RepurchaseSession(val origin: String, val deviceId: Long, val token: String) + + private fun publishRepurchase(state: RepurchaseState) { + if (repurchaseClosed.get()) return + repurchaseState = state + sendBroadcast(Intent(ACTION_REPURCHASE_STATE).setPackage(packageName)) + updateNotification(state.message, force = true) + } + + private fun repurchaseLocalIdle(): Boolean = !working.get() && taskMutex.currentTaskId() == null && + purchaseStore.activeTaskId() == null && purchaseStore.pendingOutbox().isEmpty() && + stateStore.read().currentTaskId == null && !backfillState.running + + private fun repurchaseEnvironment(session: RepurchaseSession): String? { + if (repurchaseClosed.get()) return "服务已经停止" + val credentials = identityStore.credentials() + if (credentials?.deviceId != session.deviceId || credentials.token != session.token || + ServerUrlPolicy.normalize(settingsStore.serverUrl()) != session.origin) return "设备身份或服务器配置已变化,请重新连接" + if (!registeredThisProcess.get()) return "设备尚未完成连接校验,请先测试连接" + if (GoAutoAccessibilityService.instance == null) return "请先开启无障碍服务" + val network = connectivityManager.activeNetwork ?: return "网络不可用" + if (connectivityManager.getNetworkCapabilities(network)?.hasCapability(NetworkCapabilities.NET_CAPABILITY_INTERNET) != true) return "网络不可用" + return null + } + + private fun prepareRepurchase() { + if (repurchaseClosed.get() || !repurchase.prepare()) return + repurchaseEvidence.clear() + try { + val credentials = identityStore.credentials() ?: error("missing credentials") + val session = RepurchaseSession(ServerUrlPolicy.normalize(settingsStore.serverUrl()), credentials.deviceId, credentials.token) + repurchaseSession = session + repurchaseEnvironment(session)?.let { repurchase.fail(it); return } + val port = repurchasePort(session) + if (!repurchaseLocalIdle() || port.hasOtherWork(null)) { + repurchase.fail("设备已有未结束工作或结果上传尚未收尾,请完成后再重购") + return + } + // Discovery owns its client. No other request can overwrite its page's Date sample. + val api = AgentApiClient(session.origin) + val batch = RepurchaseBatchDiscovery({ query -> + check(!repurchase.isStopRequested() && !repurchaseClosed.get()) + check(repurchaseEnvironment(session) == null) + api.purchaseHistory(session.token, query.page, query.status, query.taskNo, query.days, query.pageSize) + }, { requireNotNull(api.lastResponseServerTimeMillis) }).discover() + repurchaseEnvironment(session)?.let { repurchase.fail(it); return } + if (!repurchaseLocalIdle() || port.hasOtherWork(null)) { + repurchase.fail("读取期间出现新的采集或采购工作,请完成后重新读取") + return + } + repurchase.prepared(batch) + } catch (error: RepurchaseBatchException) { + repurchase.fail(if (repurchase.isStopRequested()) "已停止读取,没有发起重购" else when (error.reason) { + RepurchaseBatchFailure.LIST_CHANGED -> "采购列表已变化,请重新读取" + RepurchaseBatchFailure.PAGE_FAILED -> "读取采购记录失败,请检查连接后重试" + else -> "无法完整确定最后一批采购记录,请检查时间或记录后重试(${error.reason})" + }) + } catch (_: Exception) { + repurchase.fail("设备连接或记录读取失败,请恢复后重试") + } + } + + private fun tickRepurchase() { + if (repurchaseClosed.get() || !repurchase.state.busy) return + val session = repurchaseSession ?: return + val previousTaskId = repurchase.state.currentTaskId + repurchase.tick(repurchasePort(session)) + // Reset only makes the original task pending. The original dispatcher owns all + // claim/start/automation/upload steps; the round never executes a PDD action. + if (repurchase.state.currentTaskId != null && repurchase.state.currentTaskId != previousTaskId) triggerSync() + } + + private fun repurchasePort(session: RepurchaseSession): RepurchasePort = object : RepurchasePort { + private val api = AgentApiClient(session.origin) + override fun environmentProblem() = repurchaseEnvironment(session) + override fun localIdle() = repurchaseLocalIdle() + override fun latestTaskId() = api.purchaseHistory(session.token, 1, null, null, 30, 1).items.firstOrNull()?.taskId + override fun detail(taskId: Long) = api.purchaseHistoryDetail(taskId, session.token) + override fun executionEvidence(taskId: Long) = repurchaseEvidence[taskId]?.values?.toList().orEmpty() + override fun hasOtherWork(currentTaskId: Long?): Boolean { + if (api.nextTask(session.token) != null) return true // collection IDs are a separate namespace + val next = api.nextPurchaseTask(session.token) + if (next != null && next.taskId != currentTaskId) return true + // /next returns the current running/probing purchase first and can conceal new pending work. + return listOf("pending", "running", "spec_probe_pending", "order_submit_started").any { status -> + val page = api.purchaseHistory(session.token, 1, status, null, 30, 50) + page.items.any { it.taskId != currentTaskId } || page.total > page.items.size + } + } + override fun reset(taskId: Long, requestId: String): cn.ilapage.goauto.agent.network.PurchaseResetResult { + val acquired = synchronized(taskMutex) { + !working.get() && taskMutex.tryAcquire(REPURCHASE_RESERVATION_ID) + } + if (!acquired) throw RepurchaseLocalBusy() + try { + if (environmentProblem() != null || repurchase.isStopRequested() || repurchaseClosed.get() || + purchaseStore.activeTaskId() != null || purchaseStore.pendingOutbox().isNotEmpty()) { + throw AgentApiException(409, "DEVICE_BUSY", "本地工作边界变化", false) + } + return api.resetPurchaseTask(taskId, requestId, session.token) + } finally { taskMutex.release(REPURCHASE_RESERVATION_ID) } + } + } + private fun publishBackfill(state: OrderBackfillState) { backfillState = state sendBroadcast(Intent(ACTION_BACKFILL_STATE).setPackage(packageName)) @@ -607,6 +741,7 @@ class AgentForegroundService : Service() { api.startPurchaseTask(claimed.taskId, UUID.randomUUID().toString(), token) } else claimed check(task.status == "running" && task.taskAttemptId.isNotBlank()) { "采购任务没有有效 attempt" } + recordRepurchaseExecution(task, null) val diagnosticDeviceId = runCatching { identityStore.credentials()?.takeIf { it.token == token }?.deviceId }.getOrNull() if (task.taskAttemptId.matches(Regex("[0-9a-fA-F]{8}(-[0-9a-fA-F]{4}){3}-[0-9a-fA-F]{12}")) && (diagnosticDeviceId ?: 0) > 0 && task.phase in setOf("spec_probe", "purchase")) { @@ -697,6 +832,7 @@ class AgentForegroundService : Service() { // #334: a spec probe that read zero colors/sizes must fail explicitly // instead of being reported as a normal, empty spec_probe_completed. val outcome = PurchaseSpecProbePolicy.demote(rawOutcome) + recordRepurchaseExecution(task, outcome.resultType) knownResultType = outcome.resultType // The live read finishes before persistence, the new failure bubble, return to Agent, or lease release. if (FailureSnapshotPolicy.eligible(task.phase, outcome.resultType, false)) failureSnapshot(outcome.resultType, outcome.errorCode ?: "PURCHASE_EXECUTION_FAILED") @@ -730,6 +866,12 @@ class AgentForegroundService : Service() { } } + private fun recordRepurchaseExecution(task: PurchaseAgentTask, resultType: String?) { + if (!repurchaseState.busy || repurchaseState.currentTaskId != task.taskId || task.attemptNumber <= 0) return + val taskEvidence = repurchaseEvidence.getOrPut(task.taskId) { java.util.concurrent.ConcurrentHashMap() } + taskEvidence[task.attemptNumber.toLong()] = RepurchaseExecutionEvidence(task.attemptNumber.toLong(), task.taskAttemptId, task.phase, resultType) + } + private fun collectPurchaseProbe( accessibility: GoAutoAccessibilityService, task: PurchaseAgentTask, purchaseRule: PurchaseRule, diagnostic: (AgentDiagnosticEvent) -> Unit, @@ -1239,6 +1381,7 @@ class AgentForegroundService : Service() { } private fun evaluateIdleReturn() { + if (repurchaseState.phase == RepurchasePhase.RUNNING) return if (!idleReturn.isArmed()) return val accessibility = GoAutoAccessibilityService.instance if (accessibility == null) { @@ -1383,10 +1526,20 @@ class AgentForegroundService : Service() { PendingIntent.FLAG_IMMUTABLE or PendingIntent.FLAG_UPDATE_CURRENT) builder.addAction(Notification.Action.Builder(null, "停止回填", stop).build()) } + if (repurchaseState.busy) { + val stop = PendingIntent.getService(this, 367, + Intent(this, AgentForegroundService::class.java).setAction(ACTION_REPURCHASE_STOP).putExtra("roundId", repurchaseState.roundId), + PendingIntent.FLAG_IMMUTABLE or PendingIntent.FLAG_UPDATE_CURRENT) + builder.addAction(Notification.Action.Builder(null, "停止重购", stop).build()) + } return builder .setSmallIcon(android.R.drawable.stat_notify_sync) .setContentTitle(getString(R.string.app_name)) - .setContentText(if (backfillState.running) "订单回填 · 已检查 ${backfillState.checked}" else content) + .setContentText(when { + repurchaseState.busy -> "重购 ${repurchaseState.currentTaskId?.let { "CG-$it" }.orEmpty()} · 成功${repurchaseState.successCount} 失败${repurchaseState.failedCount} 跳过${repurchaseState.skippedCount}" + backfillState.running -> "订单回填 · 已检查 ${backfillState.checked}" + else -> content + }) .setContentIntent(pendingIntent) .setOngoing(true) .build() @@ -1433,6 +1586,12 @@ class AgentForegroundService : Service() { } companion object { + const val ACTION_REPURCHASE_PREPARE = "cn.ilapage.goauto.agent.REPURCHASE_PREPARE" + const val ACTION_REPURCHASE_CONFIRM = "cn.ilapage.goauto.agent.REPURCHASE_CONFIRM" + const val ACTION_REPURCHASE_STOP = "cn.ilapage.goauto.agent.REPURCHASE_STOP" + const val ACTION_REPURCHASE_STATE = "cn.ilapage.goauto.agent.REPURCHASE_STATE" + @Volatile var repurchaseState = RepurchaseState() + private set const val ACTION_BACKFILL_START = "cn.ilapage.goauto.agent.BACKFILL_START" const val ACTION_BACKFILL_STOP = "cn.ilapage.goauto.agent.BACKFILL_STOP" const val ACTION_BACKFILL_STATE = "cn.ilapage.goauto.agent.BACKFILL_STATE" @@ -1468,6 +1627,7 @@ class AgentForegroundService : Service() { private const val COOLDOWN_WAKE_GRACE_MILLIS = 5_000L private const val RETURN_CONFIRM_DELAY_MILLIS = 750L private const val CURRENT_PAGE_RESERVATION_ID = Long.MAX_VALUE + private const val REPURCHASE_RESERVATION_ID = Long.MAX_VALUE - 1 private const val PDD_PACKAGE = "com.xunmeng.pinduoduo" private const val CURRENT_PAGE_RESULT_NOTIFICATION_MILLIS = 30_000L diff --git a/android/app/src/main/java/cn/ilapage/goauto/agent/service/RepurchaseCoordinator.kt b/android/app/src/main/java/cn/ilapage/goauto/agent/service/RepurchaseCoordinator.kt new file mode 100644 index 0000000..d5f10b6 --- /dev/null +++ b/android/app/src/main/java/cn/ilapage/goauto/agent/service/RepurchaseCoordinator.kt @@ -0,0 +1,235 @@ +package cn.ilapage.goauto.agent.service + +import cn.ilapage.goauto.agent.network.PurchaseHistoryDetail +import cn.ilapage.goauto.agent.network.PurchaseResetResult +import cn.ilapage.goauto.agent.network.AgentApiException +import java.util.concurrent.atomic.AtomicReference + +interface RepurchasePort { + fun environmentProblem(): String? + fun localIdle(): Boolean + fun hasOtherWork(currentTaskId: Long?): Boolean + fun latestTaskId(): Long? + fun detail(taskId: Long): PurchaseHistoryDetail + fun reset(taskId: Long, requestId: String): PurchaseResetResult + fun executionEvidence(taskId: Long): List = emptyList() +} + +/** Only emitted by the original executor after Start and after a known outcome. */ +data class RepurchaseExecutionEvidence(val attemptNumber: Long, val attemptId: String, val phase: String, val resultType: String?) +class RepurchaseLocalBusy : IllegalStateException() + +class RepurchaseCoordinator( + private val now: () -> Long, + private val requestId: () -> String, + private val publish: (RepurchaseState) -> Unit, +) { + @Volatile var state = RepurchaseState() + private set + private val stopRequest = AtomicReference(null) + private var position = 0 + private var expectedAttempt = 0L + private var issuedAt = 0L + private var busySince: Long? = null + private var headChecked = false + private val resetRequestIds = mutableMapOf() + + // All mutation except requestStop is confined to the Service's orchestration worker. + fun prepare(): Boolean { + if (state.busy) return false + stopRequest.set(null) + position = 0 + expectedAttempt = 0 + busySince = null + headChecked = false + resetRequestIds.clear() + emit(RepurchaseState(phase = RepurchasePhase.READING, roundId = requestId(), message = "正在读取本设备最后一批采购记录…")) + return true + } + + fun prepared(batch: RepurchaseBatch) { + if (state.phase != RepurchasePhase.READING) return + if (isStopRequested()) { finish("已停止读取,没有发起重购"); return } + val skips = batch.skips.map { RepurchaseItemResult(it.taskId, RepurchaseResult.SKIPPED, it.retryDisabledReason ?: "当前不可重试") } + val message = when { + batch.items.isEmpty() -> "近30天没有采购记录" + batch.failures.isEmpty() -> "最后一批没有失败任务,不回退更早批次" + batch.candidates.isEmpty() -> "最后一批失败任务均不可重试,请查看逐条结果" + else -> "请确认本轮重购范围" + } + emit(state.copy(phase = if (batch.candidates.isEmpty()) RepurchasePhase.FINISHED else RepurchasePhase.READY, + batch = batch, results = skips, message = message)) + } + + fun confirm(roundId: String) { + if (state.phase != RepurchasePhase.READY || roundId != state.roundId) return + if (isStopRequested()) { finish("已取消重购"); return } + emit(state.copy(phase = RepurchasePhase.RUNNING, message = "正在复核本轮任务")) + } + + fun requestStop(roundId: String) { + if (roundId.isNotBlank() && roundId == state.roundId) stopRequest.set(roundId) + } + fun isStopRequested(): Boolean = stopRequest.get() == state.roundId && state.roundId.isNotBlank() + + fun fail(message: String) { + if (!state.busy) return + if (state.currentTaskId == null) finish(message) + else if (state.phase != RepurchasePhase.STOPPING) emit(state.copy(phase = RepurchasePhase.STOPPING, message = "$message;停止后续,当前任务沿用原流程收尾")) + } + + fun tick(port: RepurchasePort) { + if (!state.busy) return + if (isStopRequested()) fail("已请求停止重购") + if (state.phase !in setOf(RepurchasePhase.RUNNING, RepurchasePhase.STOPPING)) return + if (state.currentTaskId != null && now() - issuedAt >= TASK_TIMEOUT_MILLIS) { + finish("单条等待已达15分钟,已停止编排;当前任务未被取消,请人工检查") + return + } + try { + port.environmentProblem()?.let { fail(it); return } + val current = state.currentTaskId + if (port.hasOtherWork(current)) fail("检测到新的采集或采购工作") + if (!state.busy) return + if (current != null) { + observe(port, current) + return + } + if (state.phase != RepurchasePhase.RUNNING) return + if (busySince != null && now() - requireNotNull(busySince) >= DRAIN_TIMEOUT_MILLIS) { + fail("设备忙碌或结果上传尚未收尾"); return + } + if (!port.localIdle()) { + val since = busySince ?: now().also { busySince = it } + if (now() - since >= DRAIN_TIMEOUT_MILLIS) fail("设备忙碌或结果上传尚未收尾") + return + } + val batch = requireNotNull(state.batch) + if (!headChecked) { + if (port.latestTaskId() != batch.headTaskId) { fail("采购列表已变化,请重新读取后确认"); return } + headChecked = true + } + if (position >= batch.candidates.size) { finish("本轮重购已结束"); return } + val taskId = batch.candidates[position].taskId + val detail = try { port.detail(taskId) } catch (error: AgentApiException) { + if (error.code == "PURCHASE_TASK_NOT_FOUND") { skip(taskId, "任务已过期或不再属于本设备"); return } + throw error + } + check(detail.task.taskId == taskId) + if (detail.task.status != "failed" || !detail.task.retryable) { + skip(taskId, detail.task.retryDisabledReason ?: "当前状态不可重试") + return + } + // Recheck after network reads. The reset adapter also holds a short local reservation. + if (isStopRequested()) { fail("已请求停止重购"); return } + port.environmentProblem()?.let { fail(it); return } + if (!port.localIdle()) return + if (port.hasOtherWork(null)) { fail("检测到新的采集或采购工作"); return } + if (isStopRequested()) { fail("已请求停止重购"); return } + issueReset(port, taskId, detail.attemptCount) + } catch (_: Exception) { + fail("连接、身份或服务端校验失败,请恢复后再操作") + } + } + + private fun issueReset(port: RepurchasePort, taskId: Long, priorAttempt: Long) { + val fixedRequestId = resetRequestIds.getOrPut(taskId, requestId) + issuedAt = now() + expectedAttempt = 0 + emit(state.copy(currentTaskId = taskId, message = "正在重置 CG-$taskId(原任务号)")) + var ambiguous = false + repeat(2) { attempt -> + // A stop while an ambiguous request was in flight must not start a fresh write. + if (attempt > 0 && isStopRequested()) { fail("重置响应不明确,已停止,请人工核对当前任务"); return } + try { + val result = port.reset(taskId, fixedRequestId) + check(result.taskId == taskId && result.attemptNumber > priorAttempt) + expectedAttempt = result.attemptNumber.toLong() + busySince = null + emit(state.copy(message = "等待 CG-$taskId 本次执行与结果上传")) + if (isStopRequested()) fail("已请求停止重购") + return + } catch (error: Exception) { + if (!ambiguous && error is RepurchaseLocalBusy) { + emit(state.copy(currentTaskId = null, message = "等待本地忙碌操作收尾")) + if (busySince == null) busySince = now() + return + } + if (!ambiguous && error is AgentApiException && error.status in 400..499 && error.code in BUSINESS_REFUSALS) { + emit(state.copy(currentTaskId = null)) + skip(taskId, "当前业务资格拒绝:${error.code}") + return + } + if (error is AgentApiException && error.status in 400..499) { + if (!ambiguous) emit(state.copy(currentTaskId = null)) + fail(if (ambiguous) "重置响应不明确,请人工核对" else "设备、规则或身份校验失败") + return + } + ambiguous = true + } + } + fail("重置响应仍不明确,已停止,请人工核对当前任务") + } + + private fun observe(port: RepurchasePort, taskId: Long) { + if (expectedAttempt == 0L) return // Never treat an old failed result as this reset's outcome. + val detail = port.detail(taskId) + check(detail.task.taskId == taskId) + if (detail.attemptCount < expectedAttempt) return + if (detail.attemptCount > expectedAttempt) { + val evidence = port.executionEvidence(taskId) + var linkedAttempt = expectedAttempt + // Each extra attempt needs a completed probe/rematch and the original executor's + // Start evidence for the immediately following purchase attempt. A separate reset + // after failed/cancelled/success cannot satisfy this chain. + while (linkedAttempt < detail.attemptCount) { + val from = evidence.singleOrNull { it.attemptNumber == linkedAttempt } + val to = evidence.singleOrNull { it.attemptNumber == linkedAttempt + 1 } + if (from?.resultType !in setOf("spec_probe_completed", "spec_rematch_completed") || + to?.phase != "purchase" || from?.attemptId.isNullOrBlank() || to.attemptId.isBlank() || from?.attemptId == to.attemptId) { + fail("同任务出现无法关联本轮的后续尝试,已停止,请人工核对") + if (port.localIdle()) finish(state.message) + return + } + linkedAttempt++ + } + expectedAttempt = linkedAttempt + } + val status = detail.task.status + if (status == "order_result_unknown") { + fail("订单结果待核对,请人工检查,不能继续重购") + if (port.localIdle()) finish(state.message) + return + } + if (status !in TERMINAL) return + if (!port.localIdle()) return + val failed = status == "failed" + val result = when (status) { + "order_created" -> RepurchaseItemResult(taskId, RepurchaseResult.SUCCESS, "订单已创建") + "cancelled" -> RepurchaseItemResult(taskId, RepurchaseResult.SKIPPED, "本次任务已取消") + else -> RepurchaseItemResult(taskId, RepurchaseResult.FAILED, "再次失败:${detail.task.errorCode ?: "PURCHASE_FAILED"}") + } + position++ + val stopping = state.phase == RepurchasePhase.STOPPING + emit(state.copy(currentTaskId = null, results = state.results + result)) + if (failed && detail.task.errorCode in ENVIRONMENT_FAILURES) finish("登录、验证码、风控或设备环境异常,已停止后续重购") + else if (stopping || isStopRequested()) finish(state.message) + else if (position >= requireNotNull(state.batch).candidates.size) finish("本轮重购已结束;再次失败的任务不会在本轮重复执行") + } + + private fun skip(taskId: Long, reason: String) { + position++ + emit(state.copy(results = state.results + RepurchaseItemResult(taskId, RepurchaseResult.SKIPPED, reason))) + if (position >= requireNotNull(state.batch).candidates.size) finish("本轮重购已结束") + } + private fun finish(message: String) = emit(state.copy(phase = RepurchasePhase.FINISHED, message = message)) + private fun emit(value: RepurchaseState) { state = value; publish(value) } + + companion object { + const val TASK_TIMEOUT_MILLIS = 15 * 60 * 1000L + private const val DRAIN_TIMEOUT_MILLIS = 10_000L + private val TERMINAL = setOf("order_created", "failed", "cancelled") + private val BUSINESS_REFUSALS = setOf("PURCHASE_TASK_NOT_FOUND", "PURCHASE_RETRY_NOT_ALLOWED", "PURCHASE_RETRY_UNSAFE", "PURCHASE_RETRY_STALE", "PURCHASE_SPEC_MAPPING_REQUIRED") + private val ENVIRONMENT_FAILURES = setOf("PDD_LOGIN_REQUIRED", "PDD_CAPTCHA_REQUIRED", "PDD_RISK_CONTROL", "ACCESSIBILITY_NOT_READY", "DEVICE_OFFLINE", "AGENT_RESTARTED_DURING_EXECUTION") + } +} diff --git a/android/app/src/main/java/cn/ilapage/goauto/agent/service/RepurchaseState.kt b/android/app/src/main/java/cn/ilapage/goauto/agent/service/RepurchaseState.kt new file mode 100644 index 0000000..c89794e --- /dev/null +++ b/android/app/src/main/java/cn/ilapage/goauto/agent/service/RepurchaseState.kt @@ -0,0 +1,40 @@ +package cn.ilapage.goauto.agent.service + +enum class RepurchasePhase { IDLE, READING, READY, RUNNING, STOPPING, FINISHED } +enum class RepurchaseResult { SUCCESS, FAILED, SKIPPED } +data class RepurchaseItemResult(val taskId: Long, val result: RepurchaseResult, val message: String) + +/** Process-local presentation only; never persists or resumes a round after process death. */ +data class RepurchaseState( + val phase: RepurchasePhase = RepurchasePhase.IDLE, + val roundId: String = "", + val batch: RepurchaseBatch? = null, + val currentTaskId: Long? = null, + val results: List = emptyList(), + val message: String = "", +) { + val busy get() = phase in setOf(RepurchasePhase.READING, RepurchasePhase.READY, RepurchasePhase.RUNNING, RepurchasePhase.STOPPING) + val successCount get() = results.count { it.result == RepurchaseResult.SUCCESS } + val failedCount get() = results.count { it.result == RepurchaseResult.FAILED } + val skippedCount get() = results.count { it.result == RepurchaseResult.SKIPPED } + fun confirmationText(): String = buildString { + append("本设备最后一批(近30天全部状态)\n") + append("${batch?.oldestCreatedAt.orEmpty()} ~ ${batch?.newestCreatedAt.orEmpty()}\n") + append("共 ${batch?.items?.size ?: 0} 条 · 失败 ${batch?.failures?.size ?: 0} 条\n") + append("可重购 ${batch?.candidates?.size ?: 0} 条 · 跳过 ${batch?.skips?.size ?: 0} 条\n\n") + append("相邻创建时间不超过60秒推算为同批,不受当前搜索和筛选影响。\n") + append("逐条复用原任务号重试,本轮再次失败不重复执行。单条等待上限15分钟,超时只停止后续编排。\n\n") + append("可能真实创建拼多多待付款订单,系统不会付款。停止重购不取消已经提交或正在执行的任务。") + } + + fun summaryText(): String = buildString { + if (message.isNotBlank()) append(message).append('\n') + currentTaskId?.let { append("当前 CG-$it\n") } + append("成功 $successCount · 再次失败 $failedCount · 跳过 $skippedCount") + batch?.let { + val candidateIds = it.candidates.map { task -> task.taskId }.toSet() + val handled = results.count { result -> result.taskId in candidateIds } + append("\n本轮已结束 $handled / ${it.candidates.size} · 未结束 ${(it.candidates.size - handled).coerceAtLeast(0)}") + } + } +} diff --git a/android/app/src/main/java/cn/ilapage/goauto/agent/ui/TaskHistoryFragment.kt b/android/app/src/main/java/cn/ilapage/goauto/agent/ui/TaskHistoryFragment.kt index 7880c45..f67a374 100644 --- a/android/app/src/main/java/cn/ilapage/goauto/agent/ui/TaskHistoryFragment.kt +++ b/android/app/src/main/java/cn/ilapage/goauto/agent/ui/TaskHistoryFragment.kt @@ -142,8 +142,19 @@ class TaskHistoryFragment : Fragment() { private val imageRequests = mutableListOf() private var currentPageReceiverRegistered = false private var backfillPanel: LinearLayout? = null + private var repurchasePanel: LinearLayout? = null + private var repurchaseButton: MaterialButton? = null + private var repurchaseDialog: androidx.appcompat.app.AlertDialog? = null + private var repurchaseDialogRound: String? = null private val currentPageReceiver = object : BroadcastReceiver() { override fun onReceive(context: Context?, intent: Intent?) { + if (intent?.action == AgentForegroundService.ACTION_REPURCHASE_STATE) { + renderRepurchase() + if (!collection && isResumed && AgentForegroundService.repurchaseState.phase == cn.ilapage.goauto.agent.service.RepurchasePhase.FINISHED) { + detailState.taskId?.let(::loadPurchaseDetail) ?: load() + } + return + } if (intent?.action == AgentForegroundService.ACTION_BACKFILL_STATE) { renderBackfill() if (!collection && isResumed && !AgentForegroundService.backfillState.running) { @@ -191,6 +202,9 @@ class TaskHistoryFragment : Fragment() { pageColumn.addView(context.screenTitle(if (collection) "采集记录" else "采购记录")) pageColumn.addView(buildSearch()) if (!collection) { + repurchasePanel = context.column(0) + pageColumn.addView(repurchasePanel) + renderRepurchase() backfillPanel = context.column(0) pageColumn.addView(backfillPanel) renderBackfill() @@ -214,12 +228,19 @@ class TaskHistoryFragment : Fragment() { override fun onResume() { super.onResume() + renderRepurchase() detailState.taskId?.let { taskId -> if (collection) loadCollectionDetail(taskId) else loadPurchaseDetail(taskId) } ?: load() } override fun onDestroyView() { + repurchaseDialog?.setOnDismissListener(null) + repurchaseDialog?.dismiss() + repurchaseDialog = null + repurchaseDialogRound = null + repurchasePanel = null + repurchaseButton = null backfillPanel = null requestGeneration++ cancelImageRequests() @@ -289,11 +310,16 @@ class TaskHistoryFragment : Fragment() { }, LinearLayout.LayoutParams(ViewGroup.LayoutParams.WRAP_CONTENT, context.dp(48)).apply { marginStart = context.dp(8) }) } else { row.addView(MaterialButton(context).apply { - text = "回填" + repurchaseButton = this + text = "重购" textSize = 14f minimumHeight = context.dp(48) - contentDescription = "回填拼多多订单号和下单时间" - setOnClickListener { showBackfillInput() } + contentDescription = "重试本设备最后一批失败采购" + isEnabled = !AgentForegroundService.repurchaseState.busy + setOnClickListener { + isEnabled = false + sendRepurchaseAction(AgentForegroundService.ACTION_REPURCHASE_PREPARE) + } }, LinearLayout.LayoutParams(ViewGroup.LayoutParams.WRAP_CONTENT, context.dp(48)).apply { marginStart = context.dp(8) }) } addView(row, row.fullWidth()) @@ -302,6 +328,61 @@ class TaskHistoryFragment : Fragment() { } } + private fun sendRepurchaseAction(action: String, roundId: String = AgentForegroundService.repurchaseState.roundId) { + val context = context ?: return + val intent = Intent(context, AgentForegroundService::class.java).setAction(action).putExtra("roundId", roundId) + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.O) context.startForegroundService(intent) else context.startService(intent) + } + + private fun renderRepurchase() { + val panel = repurchasePanel ?: return + val context = context ?: return + val state = AgentForegroundService.repurchaseState + val phase = state.phase + repurchaseButton?.isEnabled = !state.busy + panel.removeAllViews() + if (phase != cn.ilapage.goauto.agent.service.RepurchasePhase.READY) { + repurchaseDialog?.setOnDismissListener(null) + repurchaseDialog?.dismiss() + repurchaseDialog = null + repurchaseDialogRound = null + } + if (phase == cn.ilapage.goauto.agent.service.RepurchasePhase.IDLE) return + panel.addView(context.card(context.cardColumn().apply { + addView(context.label(if (state.busy) "本轮重购" else "重购结果", 16f)) + addView(context.label(state.summaryText(), 14f)) + if (state.busy) addView(MaterialButton(context).apply { + text = if (phase == cn.ilapage.goauto.agent.service.RepurchasePhase.STOPPING) "正在收尾…" else "停止重购" + minimumHeight = context.dp(48) + isEnabled = phase != cn.ilapage.goauto.agent.service.RepurchasePhase.STOPPING + setOnClickListener { + isEnabled = false + sendRepurchaseAction(AgentForegroundService.ACTION_REPURCHASE_STOP, state.roundId) + } + }) + if (state.results.isNotEmpty()) addView(MaterialButton(context).apply { + text = "查看逐条结果" + minimumHeight = context.dp(48) + setOnClickListener { + MaterialAlertDialogBuilder(context).setTitle("本轮重购明细") + .setMessage(state.results.joinToString("\n") { "CG-${it.taskId}:${it.message}" }) + .setPositiveButton("关闭", null).show() + } + }) + })) + if (phase == cn.ilapage.goauto.agent.service.RepurchasePhase.READY && isResumed && repurchaseDialogRound != state.roundId) { + repurchaseDialogRound = state.roundId + repurchaseDialog = MaterialAlertDialogBuilder(context).setTitle("重购最后一批") + .setMessage(state.confirmationText()) + .setNegativeButton("取消") { _, _ -> sendRepurchaseAction(AgentForegroundService.ACTION_REPURCHASE_STOP, state.roundId) } + .setPositiveButton("确认重购") { _, _ -> sendRepurchaseAction(AgentForegroundService.ACTION_REPURCHASE_CONFIRM, state.roundId) } + .create().also { dialog -> + dialog.setOnCancelListener { sendRepurchaseAction(AgentForegroundService.ACTION_REPURCHASE_STOP, state.roundId) } + dialog.show() + } + } + } + private fun showBackfillInput() { if (AgentForegroundService.backfillState.running) { toast("设备忙碌,请稍后操作") @@ -729,6 +810,7 @@ class TaskHistoryFragment : Fragment() { if (currentPageReceiverRegistered) return val filter = IntentFilter(AgentForegroundService.ACTION_CURRENT_PAGE_RESULT) filter.addAction(AgentForegroundService.ACTION_BACKFILL_STATE) + filter.addAction(AgentForegroundService.ACTION_REPURCHASE_STATE) if (Build.VERSION.SDK_INT >= 33) { requireContext().registerReceiver(currentPageReceiver, filter, Context.RECEIVER_NOT_EXPORTED) } else { diff --git a/android/app/src/test/java/cn/ilapage/goauto/agent/RepurchaseServerClockTest.kt b/android/app/src/test/java/cn/ilapage/goauto/agent/RepurchaseServerClockTest.kt new file mode 100644 index 0000000..9670cd4 --- /dev/null +++ b/android/app/src/test/java/cn/ilapage/goauto/agent/RepurchaseServerClockTest.kt @@ -0,0 +1,37 @@ +package cn.ilapage.goauto.agent + +import cn.ilapage.goauto.agent.network.AgentApiClient +import org.junit.Assert.* +import org.junit.Test +import java.net.ServerSocket +import java.util.concurrent.Executors +import java.util.concurrent.TimeUnit + +class RepurchaseServerClockTest { + @Test fun historyCapturesServerDateAndMissingDateClearsPreviousSample() { + ServerSocket(0).use { server -> + val executor = Executors.newSingleThreadExecutor() + val response = "{\"data\":{\"items\":[],\"total\":0,\"page\":1,\"pageSize\":50}}" + val serving = executor.submit { + repeat(2) { index -> + server.accept().use { socket -> + socket.soTimeout = 3000 + val reader = socket.getInputStream().bufferedReader() + while (!reader.readLine().isNullOrEmpty()) { } + val date = if (index == 0) "Date: Thu, 08 Oct 2026 09:00:00 GMT\r\n" else "" + socket.getOutputStream().write(("HTTP/1.1 200 OK\r\n" + date + + "Content-Type: application/json\r\nContent-Length: ${response.toByteArray().size}\r\nConnection: close\r\n\r\n$response").toByteArray()) + } + } + } + try { + val api = AgentApiClient("http://127.0.0.1:${server.localPort}") + api.purchaseHistory("synthetic", 1, null, null, 30, 50) + assertEquals(1791450000000L, api.lastResponseServerTimeMillis) + api.purchaseHistory("synthetic", 1, null, null, 30, 50) + assertNull(api.lastResponseServerTimeMillis) + serving.get(5, TimeUnit.SECONDS) + } finally { executor.shutdownNow() } + } + } +} diff --git a/android/app/src/test/java/cn/ilapage/goauto/agent/service/RepurchaseCoordinatorTest.kt b/android/app/src/test/java/cn/ilapage/goauto/agent/service/RepurchaseCoordinatorTest.kt new file mode 100644 index 0000000..aea7edd --- /dev/null +++ b/android/app/src/test/java/cn/ilapage/goauto/agent/service/RepurchaseCoordinatorTest.kt @@ -0,0 +1,182 @@ +package cn.ilapage.goauto.agent.service + +import cn.ilapage.goauto.agent.network.* +import org.junit.Assert.* +import org.junit.Test +import java.io.IOException + +class RepurchaseCoordinatorTest { + private fun item(id: Long, status: String = "failed", retryable: Boolean = true, error: String? = null) = PurchaseHistoryItem( + id, status, "", "", "", "", "", "", "", 1, null, "CNY", null, null, error, null, retryable, "不可重试", "2026-10-08T09:00:00Z") + private fun detail(id: Long, status: String = "failed", count: Long = 1, retryable: Boolean = true, error: String? = null) = + PurchaseHistoryDetail(item(id, status, retryable, error), count, error, null, false, null, null, null, null, false, null) + private inner class Fake : RepurchasePort { + var clock = 0L + var busy = false + var otherWork = false + var environment: String? = null + var head: Long? = 2L + var resetFailure: Exception? = null + var failOnce = false + var duringReset: (() -> Unit)? = null + val details = mutableMapOf(2L to this@RepurchaseCoordinatorTest.detail(2), 1L to this@RepurchaseCoordinatorTest.detail(1)) + val resets = mutableListOf>() + var evidence = emptyList() + override fun executionEvidence(taskId: Long) = evidence + override fun environmentProblem() = environment + override fun localIdle() = !busy + override fun hasOtherWork(currentTaskId: Long?) = otherWork + override fun latestTaskId() = head + override fun detail(taskId: Long) = details.getValue(taskId) + override fun reset(taskId: Long, requestId: String): PurchaseResetResult { + resets += taskId to requestId + duringReset?.invoke() + resetFailure?.let { if (failOnce) resetFailure = null; throw it } + return PurchaseResetResult(taskId, "CG-$taskId", 2, "pending", resets.size > 1) + } + fun runner(): RepurchaseCoordinator { + var sequence = 0 + return RepurchaseCoordinator({ clock }, { "request-${++sequence}" }, {}).also { + assertTrue(it.prepare()) + it.prepared(RepurchaseBatch(listOf(item(2), item(1)))) + it.confirm(it.state.roundId) + } + } + } + + @Test fun serialResetsWaitForNewAttemptAndLocalDrain() { + val p = Fake(); val r = p.runner() + r.tick(p); assertEquals(listOf(2L), p.resets.map { it.first }) + r.tick(p); assertEquals(0, r.state.results.size) // old failed, attempt1 + p.details[2] = detail(2, "order_created", 2); p.busy = true + r.tick(p); assertEquals(0, r.state.results.size) + p.busy = false; r.tick(p); assertEquals(1, r.state.successCount) + r.tick(p); assertEquals(listOf(2L, 1L), p.resets.map { it.first }) + p.details[1] = detail(1, count = 2); r.tick(p); r.tick(p) + assertEquals(RepurchasePhase.FINISHED, r.state.phase) + assertEquals(1, r.state.failedCount); assertEquals(2, p.resets.size) + } + @Test fun lostResetResponseUsesSameRequestIdAndReplayIsNotTerminal() { + val p = Fake(); val r = p.runner(); p.resetFailure = IOException(); p.failOnce = true + r.tick(p); r.tick(p) + assertEquals(2, p.resets.size); assertEquals(p.resets[0], p.resets[1]) + assertEquals(0, r.state.results.size) + } + @Test fun unresolvedMappedValuesDoNotPreventReset() { + val p = Fake(); val r = p.runner(); r.tick(p) + assertEquals(2L, p.resets.single().first) + } + @Test fun definiteBusinessRefusalSkipsButNetworkAmbiguityStops() { + val p = Fake(); val r = p.runner() + p.resetFailure = AgentApiException(409, "PURCHASE_RETRY_STALE", "safe", false) + r.tick(p); assertEquals(1, r.state.skippedCount) + p.resetFailure = IOException(); r.tick(p); r.tick(p) + assertTrue(r.state.phase in setOf(RepurchasePhase.STOPPING, RepurchasePhase.FINISHED)) + assertEquals(3, p.resets.size) + } + @Test fun newWorkStopsSubsequentButNotCurrent() { + val p = Fake(); val r = p.runner(); r.tick(p); p.otherWork = true + r.tick(p); assertEquals(RepurchasePhase.STOPPING, r.state.phase) + p.details[2] = detail(2, "order_created", 2); r.tick(p); r.tick(p) + assertEquals(RepurchasePhase.FINISHED, r.state.phase); assertEquals(1, p.resets.size) + } + @Test fun stopBeforeResetIssuesNothingAndStaleStopIgnored() { + val p = Fake(); val r = p.runner(); r.requestStop("old"); r.tick(p) + assertEquals(1, p.resets.size) + val p2 = Fake(); val r2 = p2.runner(); r2.requestStop(r2.state.roundId); r2.tick(p2) + assertEquals(0, p2.resets.size); assertEquals(RepurchasePhase.FINISHED, r2.state.phase) + } + @Test fun stopDuringIssuedResetAllowsOnlyCurrentToDrain() { + val p = Fake(); val r = p.runner(); p.duringReset = { r.requestStop(r.state.roundId) } + r.tick(p); p.details[2] = detail(2, "order_created", 2); r.tick(p); r.tick(p) + assertEquals(1, p.resets.size); assertEquals(RepurchasePhase.FINISHED, r.state.phase) + } + @Test fun timeoutStopsOrchestrationWithoutAnotherReset() { + val p = Fake(); val r = p.runner(); r.tick(p); p.clock = 15 * 60 * 1000L + r.tick(p); assertEquals(RepurchasePhase.FINISHED, r.state.phase); assertEquals(1, p.resets.size) + } + @Test fun unknownAndCaptchaStopTheRound() { + for ((status, error) in listOf("order_result_unknown" to null, "failed" to "PDD_CAPTCHA_REQUIRED")) { + val p = Fake(); val r = p.runner(); r.tick(p); p.details[2] = detail(2, status, 2, error = error) + r.tick(p); r.tick(p) + assertEquals(RepurchasePhase.FINISHED, r.state.phase); assertEquals(1, p.resets.size) + } + } + @Test fun headChangedOrEnvironmentUnavailableDoesNotReset() { + val p = Fake(); val r = p.runner(); p.head = 3; r.tick(p) + assertEquals(RepurchasePhase.FINISHED, r.state.phase); assertTrue(p.resets.isEmpty()) + val p2 = Fake(); val r2 = p2.runner(); p2.environment = "无障碍未开启"; r2.tick(p2) + assertEquals(RepurchasePhase.FINISHED, r2.state.phase); assertTrue(p2.resets.isEmpty()) + } + @Test fun busyHasBoundedDrainAndProcessRestartDoesNotRestoreRound() { + val p = Fake(); val r = p.runner(); p.busy = true; r.tick(p) + p.clock = 10000; r.tick(p); assertEquals(RepurchasePhase.FINISHED, r.state.phase) + assertFalse(RepurchaseCoordinator({0}, {"x"}, {}).state.busy) + } + @Test fun noRetryableCandidatesDoesNotConfirmOrFallBack() { + val r = RepurchaseCoordinator({0}, {"x"}, {}) + r.prepare(); r.prepared(RepurchaseBatch(listOf(item(3, "order_created")))) + assertEquals(RepurchasePhase.FINISHED, r.state.phase) + r.confirm(r.state.roundId); assertEquals(RepurchasePhase.FINISHED, r.state.phase) + } + + @Test fun definiteEnvironmentRefusalStopsWithoutPretendingTaskWasIssued() { + val p = Fake(); val r = p.runner() + p.resetFailure = AgentApiException(401, "DEVICE_TOKEN_INVALID", "unauthorized", false) + r.tick(p) + assertEquals(RepurchasePhase.FINISHED, r.state.phase) + assertNull(r.state.currentTaskId) + assertEquals(1, p.resets.size) + } + + @Test fun secondRoundOnlyUsesStillFailedCandidates() { + val p = Fake(); val r = p.runner(); r.tick(p) + p.details[2] = detail(2, "order_created", 2); r.tick(p); r.tick(p) + p.details[1] = detail(1, count = 2); r.tick(p) + assertTrue(r.prepare()) + r.prepared(RepurchaseBatch(listOf(item(2, "order_created"), item(1)))) + assertEquals(listOf(1L), r.state.batch!!.candidates.map { it.taskId }) + r.confirm(r.state.roundId) + // New reset must return a strictly newer attempt than the refreshed detail. + p.details[1] = detail(1) + r.tick(p) + assertEquals(listOf(2L, 1L, 1L), p.resets.map { it.first }) + } + + @Test fun unexpectedExceptionDoesNotLeakRawErrorOrSkipRest() { + val p = Fake(); val r = p.runner(); p.resetFailure = IOException("private remote body") + r.tick(p) + assertFalse(r.state.message.contains("private")) + assertEquals(RepurchasePhase.STOPPING, r.state.phase) + assertEquals(listOf(2L, 2L), p.resets.map { it.first }) + } + + @Test fun externalLaterAttemptIsNotCountedAsThisRoundsResult() { + val p = Fake(); val r = p.runner(); r.tick(p) + p.details[2] = detail(2, "order_created", 3) + p.evidence = listOf(RepurchaseExecutionEvidence(2, "attempt2", "purchase", "failed"), + RepurchaseExecutionEvidence(3, "attempt3", "purchase", "order_created")) + r.tick(p); r.tick(p) + assertEquals(RepurchasePhase.FINISHED, r.state.phase) + assertEquals(0, r.state.successCount); assertEquals(1, p.resets.size) + } + + @Test fun completedProbeCanLinkToOriginalExecutorsNextPurchaseAttempt() { + val p = Fake(); val r = p.runner(); r.tick(p) + p.details[2] = detail(2, "order_created", 3) + p.evidence = listOf(RepurchaseExecutionEvidence(2, "probe2", "spec_probe", "spec_probe_completed"), + RepurchaseExecutionEvidence(3, "purchase3", "purchase", "order_created")) + r.tick(p) + assertEquals(1, r.state.successCount) + r.tick(p); assertEquals(listOf(2L, 1L), p.resets.map { it.first }) + } + + @Test fun localReservationBusyNeverPretendsARequestWasSentAndUsesSameId() { + val p = Fake(); val r = p.runner(); p.resetFailure = RepurchaseLocalBusy(); p.failOnce = true + r.tick(p); assertNull(r.state.currentTaskId); assertEquals(RepurchasePhase.RUNNING, r.state.phase) + r.tick(p); assertEquals(p.resets[0], p.resets[1]) + val p2 = Fake(); val r2 = p2.runner(); p2.resetFailure = RepurchaseLocalBusy() + r2.tick(p2); p2.clock = 10000; r2.tick(p2) + assertEquals(RepurchasePhase.FINISHED, r2.state.phase); assertNull(r2.state.currentTaskId) + } +} diff --git a/android/app/src/test/java/cn/ilapage/goauto/agent/service/RepurchaseStateTest.kt b/android/app/src/test/java/cn/ilapage/goauto/agent/service/RepurchaseStateTest.kt new file mode 100644 index 0000000..457348a --- /dev/null +++ b/android/app/src/test/java/cn/ilapage/goauto/agent/service/RepurchaseStateTest.kt @@ -0,0 +1,34 @@ +package cn.ilapage.goauto.agent.service + +import org.junit.Assert.* +import org.junit.Test + +class RepurchaseStateTest { + @Test fun confirmationExplainsOriginalTaskAndUnpaidOrderRisk() { + val state = RepurchaseState(phase = RepurchasePhase.READY) + assertTrue(state.confirmationText().contains("不会付款")) + assertTrue(state.confirmationText().contains("原任务号")) + assertTrue(state.confirmationText().contains("15分钟")) + assertTrue(state.confirmationText().contains("60秒")) + } + + @Test fun readingReadyRunningAndStoppingPreventDuplicateRounds() { + for (phase in listOf(RepurchasePhase.READING, RepurchasePhase.READY, RepurchasePhase.RUNNING, RepurchasePhase.STOPPING)) { + assertTrue(RepurchaseState(phase = phase).busy) + } + assertFalse(RepurchaseState().busy) + assertFalse(RepurchaseState(phase = RepurchasePhase.FINISHED).busy) + } + + @Test fun countsDoNotMarkUnexecutedItemsAsFailures() { + val state = RepurchaseState(results = listOf( + RepurchaseItemResult(1, RepurchaseResult.SUCCESS, "订单已创建"), + RepurchaseItemResult(2, RepurchaseResult.FAILED, "再次失败"), + RepurchaseItemResult(3, RepurchaseResult.SKIPPED, "不可重试"), + )) + assertEquals(1, state.successCount) + assertEquals(1, state.failedCount) + assertEquals(1, state.skippedCount) + assertTrue(state.summaryText().contains("成功 1")) + } +}