feat(android): repurchase last failed batch in place (#367)
This commit is contained in:
@@ -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()
|
||||
|
||||
+161
-1
@@ -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<Long, java.util.concurrent.ConcurrentHashMap<Long, RepurchaseExecutionEvidence>>()
|
||||
private val taskMutex = TaskExecutionMutex()
|
||||
private val backfillGuard = OrderBackfillGuard(taskMutex)
|
||||
private val runningTaskId = AtomicReference<Long?>(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
|
||||
|
||||
|
||||
@@ -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<RepurchaseExecutionEvidence> = 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<String?>(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<Long, String>()
|
||||
|
||||
// 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")
|
||||
}
|
||||
}
|
||||
@@ -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<RepurchaseItemResult> = 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)}")
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -142,8 +142,19 @@ class TaskHistoryFragment : Fragment() {
|
||||
private val imageRequests = mutableListOf<HistoryImageRequest>()
|
||||
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 {
|
||||
|
||||
@@ -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() }
|
||||
}
|
||||
}
|
||||
}
|
||||
+182
@@ -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<Pair<Long, String>>()
|
||||
var evidence = emptyList<RepurchaseExecutionEvidence>()
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -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"))
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user