Compare commits

...
12 changed files with 1232 additions and 8 deletions
@@ -100,6 +100,7 @@ data class PurchaseAgentTask(
val ruleSnapshot: String, val ruleSnapshot: String,
val ruleSnapshotHash: String, val ruleSnapshotHash: String,
val leaseVersion: Long, val leaseVersion: Long,
val attemptNumber: Int = 0,
) )
data class CollectionHistoryItem( data class CollectionHistoryItem(
@@ -211,6 +212,8 @@ class AgentApiException(
) : Exception(message) ) : Exception(message)
class AgentApiClient(private val serverUrl: String) { class AgentApiClient(private val serverUrl: String) {
@Volatile var lastResponseServerTimeMillis: Long? = null
private set
val failureSnapshotOrigin: String get() = ServerUrlPolicy.normalize(serverUrl) val failureSnapshotOrigin: String get() = ServerUrlPolicy.normalize(serverUrl)
/** Stream the two bounded parts; never build another combined copy of the archive. */ /** Stream the two bounded parts; never build another combined copy of the archive. */
fun uploadFailureSnapshot(snapshot: cn.ilapage.goauto.agent.diagnostics.FailureSnapshot, token: String) { 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( private fun purchaseTask(data: JSONObject) = PurchaseAgentTask(
taskId = data.getLong("taskId"), taskId = data.getLong("taskId"),
taskAttemptId = data.optString("taskAttemptId"), taskAttemptId = data.optString("taskAttemptId"),
attemptNumber = data.optInt("attemptNumber"),
phase = data.optString("phase"), phase = data.optString("phase"),
executionMode = data.getString("executionMode"), executionMode = data.getString("executionMode"),
status = data.getString("status"), 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)) } connection.outputStream.use { it.write(payload.toString().toByteArray(Charsets.UTF_8)) }
} }
val status = connection.responseCode 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 if (status == HttpURLConnection.HTTP_NO_CONTENT) return null
val stream = if (status in 200..299) connection.inputStream else connection.errorStream val stream = if (status in 200..299) connection.inputStream else connection.errorStream
val body = stream?.bufferedReader(Charsets.UTF_8)?.use { it.readText() }.orEmpty() val body = stream?.bufferedReader(Charsets.UTF_8)?.use { it.readText() }.orEmpty()
@@ -85,6 +85,11 @@ class AgentForegroundService : Service() {
private val taskExecutor: ExecutorService = Executors.newSingleThreadExecutor() private val taskExecutor: ExecutorService = Executors.newSingleThreadExecutor()
private val diagnosticExecutor: ExecutorService = Executors.newSingleThreadExecutor() private val diagnosticExecutor: ExecutorService = Executors.newSingleThreadExecutor()
private val snapshotUploadExecutor: ScheduledExecutorService = Executors.newSingleThreadScheduledExecutor() 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 taskMutex = TaskExecutionMutex()
private val backfillGuard = OrderBackfillGuard(taskMutex) private val backfillGuard = OrderBackfillGuard(taskMutex)
private val runningTaskId = AtomicReference<Long?>(null) private val runningTaskId = AtomicReference<Long?>(null)
@@ -119,6 +124,7 @@ class AgentForegroundService : Service() {
identityStore = SecureDeviceStore(this) identityStore = SecureDeviceStore(this)
settingsStore = AgentSettingsStore(this) settingsStore = AgentSettingsStore(this)
stateStore = AgentStateStore(this) stateStore = AgentStateStore(this)
repurchaseState = RepurchaseState()
// Wake locks and cooldown tickets are process-local. Never restore a stale UI flag. // Wake locks and cooldown tickets are process-local. Never restore a stale UI flag.
stateStore.setKeepScreenOn(false) stateStore.setKeepScreenOn(false)
purchaseStore = PurchaseTaskStore(this) purchaseStore = PurchaseTaskStore(this)
@@ -145,9 +151,29 @@ class AgentForegroundService : Service() {
executor.scheduleWithFixedDelay(::triggerSync, 0, HEARTBEAT_SECONDS, TimeUnit.SECONDS) executor.scheduleWithFixedDelay(::triggerSync, 0, HEARTBEAT_SECONDS, TimeUnit.SECONDS)
// Independent retention/upload ticks run even when no further task is dispatched. // Independent retention/upload ticks run even when no further task is dispatched.
snapshotUploadExecutor.scheduleWithFixedDelay(::maintainFailureSnapshots, 0, 60, TimeUnit.SECONDS) 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 { 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) { if (intent?.action == ACTION_BACKFILL_STOP) {
backfillGuard.cancelled.set(true) backfillGuard.cancelled.set(true)
return START_STICKY return START_STICKY
@@ -171,6 +197,10 @@ class AgentForegroundService : Service() {
} }
override fun onDestroy() { override fun onDestroy() {
repurchaseClosed.set(true)
repurchase.requestStop(repurchase.state.roundId)
repurchaseExecutor.shutdownNow()
repurchaseState = RepurchaseState(message = "服务已停止;不会恢复整轮重购")
backfillGuard.cancelled.set(true) backfillGuard.cancelled.set(true)
runCatching { connectivityManager.unregisterNetworkCallback(networkCallback) } runCatching { connectivityManager.unregisterNetworkCallback(networkCallback) }
cancelIdleReturn("服务已停止") 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) { private fun publishBackfill(state: OrderBackfillState) {
backfillState = state backfillState = state
sendBroadcast(Intent(ACTION_BACKFILL_STATE).setPackage(packageName)) sendBroadcast(Intent(ACTION_BACKFILL_STATE).setPackage(packageName))
@@ -607,6 +741,7 @@ class AgentForegroundService : Service() {
api.startPurchaseTask(claimed.taskId, UUID.randomUUID().toString(), token) api.startPurchaseTask(claimed.taskId, UUID.randomUUID().toString(), token)
} else claimed } else claimed
check(task.status == "running" && task.taskAttemptId.isNotBlank()) { "采购任务没有有效 attempt" } check(task.status == "running" && task.taskAttemptId.isNotBlank()) { "采购任务没有有效 attempt" }
recordRepurchaseExecution(task, null)
val diagnosticDeviceId = runCatching { identityStore.credentials()?.takeIf { it.token == token }?.deviceId }.getOrNull() 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}")) && 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")) { (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 // #334: a spec probe that read zero colors/sizes must fail explicitly
// instead of being reported as a normal, empty spec_probe_completed. // instead of being reported as a normal, empty spec_probe_completed.
val outcome = PurchaseSpecProbePolicy.demote(rawOutcome) val outcome = PurchaseSpecProbePolicy.demote(rawOutcome)
recordRepurchaseExecution(task, outcome.resultType)
knownResultType = outcome.resultType knownResultType = outcome.resultType
// The live read finishes before persistence, the new failure bubble, return to Agent, or lease release. // 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") 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( private fun collectPurchaseProbe(
accessibility: GoAutoAccessibilityService, task: PurchaseAgentTask, purchaseRule: PurchaseRule, accessibility: GoAutoAccessibilityService, task: PurchaseAgentTask, purchaseRule: PurchaseRule,
diagnostic: (AgentDiagnosticEvent) -> Unit, diagnostic: (AgentDiagnosticEvent) -> Unit,
@@ -1239,6 +1381,7 @@ class AgentForegroundService : Service() {
} }
private fun evaluateIdleReturn() { private fun evaluateIdleReturn() {
if (repurchaseState.phase == RepurchasePhase.RUNNING) return
if (!idleReturn.isArmed()) return if (!idleReturn.isArmed()) return
val accessibility = GoAutoAccessibilityService.instance val accessibility = GoAutoAccessibilityService.instance
if (accessibility == null) { if (accessibility == null) {
@@ -1383,10 +1526,20 @@ class AgentForegroundService : Service() {
PendingIntent.FLAG_IMMUTABLE or PendingIntent.FLAG_UPDATE_CURRENT) PendingIntent.FLAG_IMMUTABLE or PendingIntent.FLAG_UPDATE_CURRENT)
builder.addAction(Notification.Action.Builder(null, "停止回填", stop).build()) 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 return builder
.setSmallIcon(android.R.drawable.stat_notify_sync) .setSmallIcon(android.R.drawable.stat_notify_sync)
.setContentTitle(getString(R.string.app_name)) .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) .setContentIntent(pendingIntent)
.setOngoing(true) .setOngoing(true)
.build() .build()
@@ -1433,6 +1586,12 @@ class AgentForegroundService : Service() {
} }
companion object { 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_START = "cn.ilapage.goauto.agent.BACKFILL_START"
const val ACTION_BACKFILL_STOP = "cn.ilapage.goauto.agent.BACKFILL_STOP" const val ACTION_BACKFILL_STOP = "cn.ilapage.goauto.agent.BACKFILL_STOP"
const val ACTION_BACKFILL_STATE = "cn.ilapage.goauto.agent.BACKFILL_STATE" 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 COOLDOWN_WAKE_GRACE_MILLIS = 5_000L
private const val RETURN_CONFIRM_DELAY_MILLIS = 750L private const val RETURN_CONFIRM_DELAY_MILLIS = 750L
private const val CURRENT_PAGE_RESERVATION_ID = Long.MAX_VALUE 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 PDD_PACKAGE = "com.xunmeng.pinduoduo"
private const val CURRENT_PAGE_RESULT_NOTIFICATION_MILLIS = 30_000L private const val CURRENT_PAGE_RESULT_NOTIFICATION_MILLIS = 30_000L
@@ -0,0 +1,168 @@
package cn.ilapage.goauto.agent.service
import cn.ilapage.goauto.agent.network.HistoryPage
import cn.ilapage.goauto.agent.network.PurchaseHistoryItem
import java.text.ParsePosition
import java.text.SimpleDateFormat
import java.util.Collections
import java.util.Locale
import java.util.TimeZone
data class RepurchaseHistoryQuery(val page: Int, val pageSize: Int = 50, val days: Int = 30, val status: String? = null, val taskNo: String? = null)
enum class RepurchaseBatchFailure { PAGE_FAILED, INVALID_PAGE, INVALID_TIMESTAMP, INVALID_ORDER, INCOMPLETE_BATCH, LIST_CHANGED, PAGE_LIMIT }
class RepurchaseBatchException(val reason: RepurchaseBatchFailure) : IllegalStateException(reason.name)
class RepurchaseBatch internal constructor(items: List<PurchaseHistoryItem>) {
val items: List<PurchaseHistoryItem> = immutable(items)
val candidates: List<PurchaseHistoryItem> = immutable(items.filter { it.status == "failed" && it.retryable })
val failures: List<PurchaseHistoryItem> = immutable(items.filter { it.status == "failed" })
val skips: List<PurchaseHistoryItem> = immutable(failures.filterNot { it.retryable })
val headTaskId = items.firstOrNull()?.taskId
val newestCreatedAt = items.firstOrNull()?.createdAt
val oldestCreatedAt = items.lastOrNull()?.createdAt
private fun immutable(items: List<PurchaseHistoryItem>): List<PurchaseHistoryItem> =
Collections.unmodifiableList(ArrayList(items))
}
/** Read-only discovery. The callback must use the current device's authenticated history endpoint.
*
* [nowMillis] MUST return the trusted server Date time from the response to that fetch, and throw
* when it is missing or invalid. Never fall back to the phone clock or a previous response's Date:
* the server applies the rolling 30-day cutoff, and a slow phone could accept a truncated batch.
* Tests may inject a synthetic server clock. It is read only after each successful page response.
*
* Offset overlap is tolerated; a changed head is rejected. This is not a snapshot guarantee for
* arbitrary concurrent deletion, reassignment or timestamp edits on the server.
*/
class RepurchaseBatchDiscovery(
private val fetchPage: (RepurchaseHistoryQuery) -> HistoryPage<PurchaseHistoryItem>,
private val nowMillis: () -> Long,
private val maxPages: Int = 200,
) {
init { require(maxPages > 0) }
fun discover(): RepurchaseBatch {
var latestServerTime = 0L
fun readWithServerTime(pageNumber: Int): HistoryPage<PurchaseHistoryItem> {
val page = readPage(pageNumber)
val serverTime = try { nowMillis() } catch (_: Exception) {
refuse(RepurchaseBatchFailure.INCOMPLETE_BATCH)
}
if (serverTime <= 0) refuse(RepurchaseBatchFailure.INCOMPLETE_BATCH)
latestServerTime = maxOf(latestServerTime, serverTime)
return page
}
val seen = linkedMapOf<Long, Entry>()
val batch = mutableListOf<PurchaseHistoryItem>()
var head: Long? = null
var previous: Entry? = null
var originalTotal: Long? = null
var foundBoundary = false
var reachedEnd = false
for (pageNumber in 1..maxPages) {
val page = readWithServerTime(pageNumber)
if (originalTotal == null) originalTotal = page.total
else if (page.total != originalTotal) refuse(RepurchaseBatchFailure.LIST_CHANGED)
val entries = parseOrdered(page.items)
if (pageNumber == 1) head = entries.firstOrNull()?.item?.taskId
for (entry in entries) {
val duplicate = seen[entry.item.taskId]
if (duplicate != null) {
if (duplicate.time != entry.time) refuse(RepurchaseBatchFailure.INVALID_ORDER)
continue
}
val prior = previous
if (prior != null) {
if (ENTRY_ORDER.compare(prior, entry) > 0) refuse(RepurchaseBatchFailure.INVALID_ORDER)
if (prior.time.moreThanSixtySecondsAfter(entry.time)) {
foundBoundary = true
break
}
}
seen[entry.item.taskId] = entry
batch += entry.item
previous = entry
}
if (foundBoundary) break
if (pageNumber.toLong() * PAGE_SIZE >= page.total) {
// Offset overlap without enough unique rows cannot establish a complete tail.
if (seen.size.toLong() != page.total) refuse(RepurchaseBatchFailure.INCOMPLETE_BATCH)
reachedEnd = true
break
}
}
if (!foundBoundary && !reachedEnd) refuse(RepurchaseBatchFailure.PAGE_LIMIT)
val finalPage = readWithServerTime(1)
val finalHead = parseOrdered(finalPage.items).firstOrNull()?.item?.taskId
if (finalHead != head || finalPage.total != originalTotal) refuse(RepurchaseBatchFailure.LIST_CHANGED)
if (reachedEnd && previous != null) {
// The API applies a rolling 30-day range, not a midnight-based date filter.
// HTTP Date has whole-second precision; include its possible fractional second.
// Preserve the latest observed time if a later response's Date moves backwards.
val floorMillis = latestServerTime + 1000 - DAYS * 24L * 60 * 60 * 1000
val floorSeconds = floorMillis / 1000 - if (floorMillis < 0 && floorMillis % 1000 != 0L) 1 else 0
val floor = Timestamp(floorSeconds, ((floorMillis - floorSeconds * 1000) * 1_000_000).toInt())
if (!previous.time.moreThanSixtySecondsAfter(floor)) refuse(RepurchaseBatchFailure.INCOMPLETE_BATCH)
}
return RepurchaseBatch(batch)
}
private fun readPage(pageNumber: Int): HistoryPage<PurchaseHistoryItem> {
val page = try {
fetchPage(RepurchaseHistoryQuery(page = pageNumber))
} catch (_: Exception) {
// Do not expose a network exception that might contain credentials or response data.
refuse(RepurchaseBatchFailure.PAGE_FAILED)
}
val offset = (pageNumber - 1L) * PAGE_SIZE
if (page.page != pageNumber || page.pageSize != PAGE_SIZE || page.total < 0 ||
page.items.size.toLong() != (page.total - offset).coerceIn(0, PAGE_SIZE.toLong()) ||
page.items.any { it.taskId <= 0 } || page.items.map { it.taskId }.distinct().size != page.items.size
) refuse(RepurchaseBatchFailure.INVALID_PAGE)
return page
}
private fun parseOrdered(items: List<PurchaseHistoryItem>): List<Entry> {
val entries = items.map { Entry(it, parseTime(it.createdAt)) }
val ordered = entries.sortedWith(ENTRY_ORDER)
if (entries != ordered) refuse(RepurchaseBatchFailure.INVALID_ORDER)
return ordered
}
private fun parseTime(raw: String): Timestamp {
val match = TIMESTAMP.matchEntire(raw) ?: refuse(RepurchaseBatchFailure.INVALID_TIMESTAMP)
val zone = match.groupValues[3].let { if (it == "Z") "+0000" else it.replace(":", "") }
val normalized = match.groupValues[1] + zone
val position = ParsePosition(0)
val parsed = SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ssZ", Locale.ROOT).apply {
isLenient = false
timeZone = TimeZone.getTimeZone("UTC")
}.parse(normalized, position)
if (parsed == null || position.index != normalized.length) refuse(RepurchaseBatchFailure.INVALID_TIMESTAMP)
return Timestamp(parsed.time / 1000, match.groupValues[2].padEnd(9, '0').toInt())
}
private data class Entry(val item: PurchaseHistoryItem, val time: Timestamp)
private data class Timestamp(val seconds: Long, val nanos: Int) : Comparable<Timestamp> {
override fun compareTo(other: Timestamp): Int =
seconds.compareTo(other.seconds).takeIf { it != 0 } ?: nanos.compareTo(other.nanos)
fun moreThanSixtySecondsAfter(older: Timestamp): Boolean {
val secondsApart = seconds - older.seconds
return secondsApart > 60 || (secondsApart == 60L && nanos > older.nanos)
}
}
companion object {
private const val PAGE_SIZE = 50
private const val DAYS = 30
private val TIMESTAMP = Regex("(\\d{4}-\\d{2}-\\d{2}T\\d{2}:\\d{2}:\\d{2})(?:\\.(\\d{1,9}))?(Z|[+-]\\d{2}:\\d{2})")
private val ENTRY_ORDER = compareByDescending<Entry> { it.time }.thenByDescending { it.item.taskId }
private fun refuse(reason: RepurchaseBatchFailure): Nothing = throw RepurchaseBatchException(reason)
}
}
@@ -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 val imageRequests = mutableListOf<HistoryImageRequest>()
private var currentPageReceiverRegistered = false private var currentPageReceiverRegistered = false
private var backfillPanel: LinearLayout? = null 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() { private val currentPageReceiver = object : BroadcastReceiver() {
override fun onReceive(context: Context?, intent: Intent?) { 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) { if (intent?.action == AgentForegroundService.ACTION_BACKFILL_STATE) {
renderBackfill() renderBackfill()
if (!collection && isResumed && !AgentForegroundService.backfillState.running) { if (!collection && isResumed && !AgentForegroundService.backfillState.running) {
@@ -191,6 +202,9 @@ class TaskHistoryFragment : Fragment() {
pageColumn.addView(context.screenTitle(if (collection) "采集记录" else "采购记录")) pageColumn.addView(context.screenTitle(if (collection) "采集记录" else "采购记录"))
pageColumn.addView(buildSearch()) pageColumn.addView(buildSearch())
if (!collection) { if (!collection) {
repurchasePanel = context.column(0)
pageColumn.addView(repurchasePanel)
renderRepurchase()
backfillPanel = context.column(0) backfillPanel = context.column(0)
pageColumn.addView(backfillPanel) pageColumn.addView(backfillPanel)
renderBackfill() renderBackfill()
@@ -214,12 +228,19 @@ class TaskHistoryFragment : Fragment() {
override fun onResume() { override fun onResume() {
super.onResume() super.onResume()
renderRepurchase()
detailState.taskId?.let { taskId -> detailState.taskId?.let { taskId ->
if (collection) loadCollectionDetail(taskId) else loadPurchaseDetail(taskId) if (collection) loadCollectionDetail(taskId) else loadPurchaseDetail(taskId)
} ?: load() } ?: load()
} }
override fun onDestroyView() { override fun onDestroyView() {
repurchaseDialog?.setOnDismissListener(null)
repurchaseDialog?.dismiss()
repurchaseDialog = null
repurchaseDialogRound = null
repurchasePanel = null
repurchaseButton = null
backfillPanel = null backfillPanel = null
requestGeneration++ requestGeneration++
cancelImageRequests() cancelImageRequests()
@@ -289,11 +310,16 @@ class TaskHistoryFragment : Fragment() {
}, LinearLayout.LayoutParams(ViewGroup.LayoutParams.WRAP_CONTENT, context.dp(48)).apply { marginStart = context.dp(8) }) }, LinearLayout.LayoutParams(ViewGroup.LayoutParams.WRAP_CONTENT, context.dp(48)).apply { marginStart = context.dp(8) })
} else { } else {
row.addView(MaterialButton(context).apply { row.addView(MaterialButton(context).apply {
text = "回填" repurchaseButton = this
text = "重购"
textSize = 14f textSize = 14f
minimumHeight = context.dp(48) minimumHeight = context.dp(48)
contentDescription = "回填拼多多订单号和下单时间" contentDescription = "重试本设备最后一批失败采购"
setOnClickListener { showBackfillInput() } 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) }) }, LinearLayout.LayoutParams(ViewGroup.LayoutParams.WRAP_CONTENT, context.dp(48)).apply { marginStart = context.dp(8) })
} }
addView(row, row.fullWidth()) 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() { private fun showBackfillInput() {
if (AgentForegroundService.backfillState.running) { if (AgentForegroundService.backfillState.running) {
toast("设备忙碌,请稍后操作") toast("设备忙碌,请稍后操作")
@@ -729,6 +810,7 @@ class TaskHistoryFragment : Fragment() {
if (currentPageReceiverRegistered) return if (currentPageReceiverRegistered) return
val filter = IntentFilter(AgentForegroundService.ACTION_CURRENT_PAGE_RESULT) val filter = IntentFilter(AgentForegroundService.ACTION_CURRENT_PAGE_RESULT)
filter.addAction(AgentForegroundService.ACTION_BACKFILL_STATE) filter.addAction(AgentForegroundService.ACTION_BACKFILL_STATE)
filter.addAction(AgentForegroundService.ACTION_REPURCHASE_STATE)
if (Build.VERSION.SDK_INT >= 33) { if (Build.VERSION.SDK_INT >= 33) {
requireContext().registerReceiver(currentPageReceiver, filter, Context.RECEIVER_NOT_EXPORTED) requireContext().registerReceiver(currentPageReceiver, filter, Context.RECEIVER_NOT_EXPORTED)
} else { } 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() }
}
}
}
@@ -0,0 +1,241 @@
package cn.ilapage.goauto.agent.service
import cn.ilapage.goauto.agent.network.HistoryPage
import cn.ilapage.goauto.agent.network.PurchaseHistoryItem
import org.junit.Assert.*
import org.junit.Test
import java.text.SimpleDateFormat
import java.util.Date
import java.util.Locale
import java.util.TimeZone
class RepurchaseBatchTest {
private val now = 1_791_446_400_000L
private fun at(secondsAgo: Long) = SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss'Z'", Locale.ROOT).apply {
timeZone = TimeZone.getTimeZone("UTC")
}.format(Date(now - secondsAgo * 1000))
private fun item(id: Long, secondsAgo: Long = 0, status: String = "failed", retryable: Boolean = true) = PurchaseHistoryItem(
id, status, "synthetic", "synthetic", "synthetic", "", "", "", "", 1, null, "CNY", null, null,
null, null, retryable, if (retryable) null else "不可重试", at(secondsAgo),
)
private fun discover(items: List<PurchaseHistoryItem>): RepurchaseBatch = RepurchaseBatchDiscovery({ q ->
HistoryPage(items.drop((q.page - 1) * 50).take(50), items.size.toLong(), q.page, 50)
}, { now }).discover()
private fun refused(reason: RepurchaseBatchFailure, block: () -> Unit) {
try { block(); fail("Expected $reason") } catch (e: RepurchaseBatchException) { assertEquals(reason, e.reason) }
}
@Test fun sixtySecondsJoinsTransitivelyButSixtyOneStartsOlderBatch() {
val batch = discover(listOf(item(5), item(4, 60), item(3, 120), item(2, 181), item(1, 182)))
assertEquals(listOf(5L, 4L, 3L), batch.items.map { it.taskId })
assertEquals(at(0), batch.newestCreatedAt)
assertEquals(at(120), batch.oldestCreatedAt)
assertEquals(5L, batch.headTaskId)
}
@Test fun equalInstantsUseDescendingIdsAndRespectTimezones() {
val batch = discover(listOf(item(3).copy(createdAt = "2026-10-08T08:00:00.123456789+08:00"),
item(2).copy(createdAt = "2026-10-08T00:00:00.123456789Z"),
item(1).copy(createdAt = "2026-10-07T20:00:00.123456789-04:00")))
assertEquals(listOf(3L, 2L, 1L), batch.items.map { it.taskId })
}
@Test fun subMillisecondGapAboveSixtySecondsIsBoundary() {
val batch = discover(listOf(item(2).copy(createdAt = "2026-10-08T00:01:00.000000001Z"),
item(1).copy(createdAt = "2026-10-08T00:00:00Z")))
assertEquals(listOf(2L), batch.items.map { it.taskId })
}
@Test fun readsFullStateHistoryAcrossFiftyItemBoundaryAndRechecksHead() {
val requests = mutableListOf<RepurchaseHistoryQuery>()
val data = (70L downTo 1).map { item(it, 70 - it, if (it % 2 == 0L) "failed" else "order_created") }
val result = RepurchaseBatchDiscovery({ q ->
requests += q
HistoryPage(data.drop((q.page - 1) * 50).take(50), 70, q.page, 50)
}, { now }).discover()
assertEquals(70, result.items.size)
assertEquals(35, result.candidates.size)
assertEquals(listOf(1, 2, 1), requests.map { it.page })
assertTrue(requests.all { it.days == 30 && it.pageSize == 50 && it.status == null && it.taskNo == null })
}
@Test fun overlappingPagesDeduplicateTaskIdsBeforeBoundary() {
val data = (60L downTo 1).map { item(it, if (it >= 10) 60 - it else 300) }
val result = RepurchaseBatchDiscovery({ q ->
val page = if (q.page == 1) data.take(50) else data.drop(49)
HistoryPage(page, 61, q.page, 50)
}, { now }).discover()
assertEquals(51, result.items.size)
assertEquals(51, result.items.map { it.taskId }.distinct().size)
}
@Test fun headChangeRefusesEvenWhenBoundaryAlreadyFound() {
var calls = 0
refused(RepurchaseBatchFailure.LIST_CHANGED) {
RepurchaseBatchDiscovery({ q ->
calls++
HistoryPage(if (calls == 1) listOf(item(2), item(1, 100)) else listOf(item(3), item(2), item(1, 100)),
if (calls == 1) 2 else 3, q.page, 50)
}, { now }).discover()
}
}
@Test fun invalidDatesAndMissingZonesAreRefused() {
listOf("", "2026-02-30T00:00:00Z", "2026-10-08T00:00:00", "2026-10-08T00:00:00+25:00", "2026-10-08T00:00:00Zjunk").forEach {
refused(RepurchaseBatchFailure.INVALID_TIMESTAMP) { discover(listOf(item(1).copy(createdAt = it))) }
}
}
@Test fun invertedChronologicalOrIdOrderingIsRefused() {
refused(RepurchaseBatchFailure.INVALID_ORDER) { discover(listOf(item(2, 20), item(1))) }
refused(RepurchaseBatchFailure.INVALID_ORDER) { discover(listOf(item(1), item(2))) }
}
@Test fun nearThirtyDayCutoffCannotProveBatchComplete() {
val days30 = 30L * 24 * 60 * 60
listOf(days30, days30 - 60, days30 - 61, days30 + 1).forEach { age ->
refused(RepurchaseBatchFailure.INCOMPLETE_BATCH) { discover(listOf(item(1, age))) }
}
assertEquals(1, discover(listOf(item(1, days30 - 62))).items.size)
}
@Test fun usesServerResponseClockEvenWhenPhoneIsFiveMinutesSlow() {
val phoneNow = now
val serverNow = phoneNow + 5 * 60 * 1000
val ageFromPhoneSeconds = 30L * 24 * 60 * 60 - 5 * 60 - 30
var responseTime: Long? = null
refused(RepurchaseBatchFailure.INCOMPLETE_BATCH) {
RepurchaseBatchDiscovery({ q ->
responseTime = serverNow
HistoryPage(listOf(item(1, ageFromPhoneSeconds)), 1, q.page, 50)
}, { checkNotNull(responseTime) }).discover()
}
}
@Test fun responseTimeIsNotRequiredBeforeFirstFetch() {
var responseTime: Long? = null
val batch = RepurchaseBatchDiscovery({ q ->
responseTime = now
HistoryPage(listOf(item(1)), 1, q.page, 50)
}, { checkNotNull(responseTime) }).discover()
assertEquals(1, batch.items.size)
}
@Test fun missingServerResponseClockRefusesWithoutPhoneClockFallback() {
refused(RepurchaseBatchFailure.INCOMPLETE_BATCH) {
RepurchaseBatchDiscovery({ q -> HistoryPage(listOf(item(1)), 1, q.page, 50) },
{ error("No server Date") }).discover()
}
}
@Test fun retainsLatestObservedServerClockWhenLaterResponseTimeRegresses() {
var responseTime = now
var calls = 0
refused(RepurchaseBatchFailure.INCOMPLETE_BATCH) {
RepurchaseBatchDiscovery({ q ->
responseTime = if (++calls == 1) now + 300_000 else now
HistoryPage(listOf(item(1, 30L * 24 * 60 * 60 - 330)), 1, q.page, 50)
}, { responseTime }).discover()
}
}
@Test fun exactlyFiftyRecordsCanEstablishTrustworthyEnd() {
val requested = mutableListOf<Int>()
val batch = RepurchaseBatchDiscovery({ q ->
requested += q.page
HistoryPage((50L downTo 1).map { item(it, 50 - it) }, 50, q.page, 50)
}, { now }).discover()
assertEquals(50, batch.items.size)
assertEquals(listOf(1, 1), requested)
}
@Test fun latestSuccessfulBatchDoesNotFallBackToOlderFailures() {
val result = discover(listOf(item(2, status = "order_created"), item(1, 61)))
assertEquals(1, result.items.size)
assertTrue(result.failures.isEmpty())
assertTrue(result.candidates.isEmpty())
}
@Test fun onlyFailedRetryableAreCandidatesAndSecondRoundUsesRemainingFailures() {
val items = (10L downTo 1).map { item(it, status = if (it <= 4) "failed" else "order_created") }
val first = discover(items)
assertEquals(listOf(4L, 3L, 2L, 1L), first.candidates.map { it.taskId })
val second = discover(items.map { if (it.taskId in 3L..4L) it.copy(status = "order_created") else it })
assertEquals(listOf(2L, 1L), second.candidates.map { it.taskId })
val blocked = discover(listOf(item(2, retryable = false), item(1, status = "pending")))
assertEquals(1, blocked.failures.size)
assertEquals("不可重试", blocked.skips.single().retryDisabledReason)
assertTrue(blocked.candidates.isEmpty())
}
@Test fun pageFailureCannotReturnPartOfBatch() {
refused(RepurchaseBatchFailure.PAGE_FAILED) {
RepurchaseBatchDiscovery({ q ->
if (q.page == 2) error("synthetic network failure")
HistoryPage((60L downTo 11).map { item(it, 60 - it) }, 60, q.page, 50)
}, { now }).discover()
}
}
@Test fun shortPageWithUnseenTotalAndWrongMetadataAreRefused() {
listOf(HistoryPage(listOf(item(1)), 2, 1, 50), HistoryPage(listOf(item(1)), 1, 2, 50),
HistoryPage(listOf(item(1)), 1, 1, 20), HistoryPage(emptyList(), -1, 1, 50)).forEach { page ->
refused(RepurchaseBatchFailure.INVALID_PAGE) { RepurchaseBatchDiscovery({ page }, { now }).discover() }
}
}
@Test fun maximumPagesCannotReturnTruncatedBatch() {
refused(RepurchaseBatchFailure.PAGE_LIMIT) {
RepurchaseBatchDiscovery({ q -> HistoryPage((60L downTo 11).map { item(it, 60 - it) }, 60, q.page, 50) },
{ now }, maxPages = 1).discover()
}
}
@Test fun crossPageOrderInversionAndChangedDuplicateTimestampAreRefused() {
listOf(item(1, 1), item(11, 99)).forEach { next ->
refused(RepurchaseBatchFailure.INVALID_ORDER) {
RepurchaseBatchDiscovery({ q ->
HistoryPage(if (q.page == 1) (60L downTo 11).map { item(it, 60 - it) } else listOf(next),
51, q.page, 50)
}, { now }).discover()
}
}
}
@Test fun duplicatesCannotMakeAnIncompleteTailLookComplete() {
refused(RepurchaseBatchFailure.INCOMPLETE_BATCH) {
RepurchaseBatchDiscovery({ q ->
HistoryPage(if (q.page == 1) (60L downTo 11).map { item(it, 60 - it) } else listOf(item(11, 49)),
51, q.page, 50)
}, { now }).discover()
}
}
@Test fun failedHeadRecheckRefusesTheAlreadyFoundBatch() {
var calls = 0
refused(RepurchaseBatchFailure.PAGE_FAILED) {
RepurchaseBatchDiscovery({ q ->
if (++calls > 1) error("synthetic network failure")
HistoryPage(listOf(item(1)), 1, q.page, 50)
}, { now }).discover()
}
}
@Test fun emptyHistoryIsRecheckedAndReturnedEmpty() {
var calls = 0
val batch = RepurchaseBatchDiscovery({ q -> calls++; HistoryPage(emptyList(), 0, q.page, 50) }, { now }).discover()
assertTrue(batch.items.isEmpty())
assertNull(batch.headTaskId)
assertEquals(2, calls)
}
@Test fun immutableSnapshotDoesNotExposeMutableLists() {
val batch = discover(listOf(item(1)))
listOf(batch.items, batch.candidates, batch.failures).forEach { items ->
try { (items as MutableList).clear(); fail("Mutable snapshot") } catch (_: UnsupportedOperationException) { }
}
}
}
@@ -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"))
}
}
+14 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件) generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Architecture-and-Code-Map wiki_page: Architecture-and-Code-Map
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Architecture-and-Code-Map.- wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Architecture-and-Code-Map.-
wiki_revision: 6a6a1a9a13aaf5284e5a81a79f849ad28525ceba wiki_revision: bf3ce6efd67fde1fdb6eacd3248f186aea880391
synchronized_at: 2026-10-08T07:10:44Z synchronized_at: 2026-10-08T10:15:16Z
<!-- gitea-wiki-mirror:end --> <!-- gitea-wiki-mirror:end -->
<!-- gitea-wiki-mirror:start --> <!-- gitea-wiki-mirror:start -->
@@ -652,3 +652,15 @@ Web 唯一展示位置为“采集采购 → SYB 同步记录”:列表状态
- Server `purchase/failure_snapshot*.go` 提供专用上传/安全摘要/管理员ZIP下载;`models/purchase_failure_snapshot.go` 独立私有表。追加迁移 `1791400000000_purchase_failure_snapshot.go` 保存 LONGBLOB ZIP、LONGTEXT manifest、attempt唯一约束及到期索引。 - Server `purchase/failure_snapshot*.go` 提供专用上传/安全摘要/管理员ZIP下载;`models/purchase_failure_snapshot.go` 独立私有表。追加迁移 `1791400000000_purchase_failure_snapshot.go` 保存 LONGBLOB ZIP、LONGTEXT manifest、attempt唯一约束及到期索引。
- `access.AdminAPIs` 注册两个仅管理员读取接口;普通任务DTO和Client API不增加原始诊断字段。通用操作日志精确排除三个诊断端点,私有SQL读写使用静默logger。服务启动立即清理并每小时清理过期诊断,不依赖新任务。 - `access.AdminAPIs` 注册两个仅管理员读取接口;普通任务DTO和Client API不增加原始诊断字段。通用操作日志精确排除三个诊断端点,私有SQL读写使用静默logger。服务启动立即清理并每小时清理过期诊断,不依赖新任务。
- Web `purchase-tasks/FailureSnapshotCell.vue` 是现有执行记录末尾200px管理员专属列,独立摘要查询、固定原因中文和附件下载;无原始XML在线预览。 - Web `purchase-tasks/FailureSnapshotCell.vue` 是现有执行记录末尾200px管理员专属列,独立摘要查询、固定原因中文和附件下载;无原始XML在线预览。
## Android 重购编排(#367)
实现绑定 `5430dc1`,工单分支尚未合并/安装/发布;无Server/Web、数据库或接口迁移。
- `TaskHistoryFragment` 把采购搜索栏入口改为“重购”,发送 PREPARE/CONFIRM/STOP、显示 Service 内存状态;一次确认与逐条摘要由 `RepurchaseState` 提供,不由 Fragment 驱动采购生命周期。
- `RepurchaseBatchDiscovery` 注入既有设备历史查询和服务端HTTP Date时钟,负责全状态/50条页/30天窗口、相邻60秒边界、去重、完整性及重读首页;最多200页,不使用本地历史缓存执行重购。
- `RepurchaseCoordinator` 通过小型 `RepurchasePort` 执行非阻塞状态推进;只调用原reset和只读详情,固定每任务requestId、expected attempt、10秒本地忙碌复核及15分钟单条上限。依靠原执行器的attemptNumber/UUID/phase/已知结果识别合法探测续行,拒绝外部更晚attempt冒充本轮结果。
- `AgentForegroundService` 独立单线程scheduled executor每2秒推进编排,历史网络不阻塞心跳或taskExecutor;发现使用独享API客户端,reset只短暂预留既有TaskExecutionMutex。原scheduler继续唯一负责claim/start/真实PDD操作及Outbox,编排不持整轮锁、不新建执行器路径。
- 新工作查询同时检查采集next、采购next及本设备活动状态历史,防止采购next优先返回本轮current而遮蔽其他pending。停止信号为原子roundId,配置/身份冻结,仅内存保存摘要;进程重启不恢复整轮,既有单任务恢复及失败私有诊断不变。
- `AgentApiClient` 仅补读已有采购payload的attemptNumber和标准HTTP Date响应元数据;不添加服务端字段或API。Date缺失清空样本,不沿用手机时间。
+28 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件) generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Business-Rules-and-Glossary wiki_page: Business-Rules-and-Glossary
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Business-Rules-and-Glossary.- wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Business-Rules-and-Glossary.-
wiki_revision: b9efffa2045d01399bfe02ce626cfd1bb379ec41 wiki_revision: c4870fd01a4569059fc0afc77179d14fbd286384
synchronized_at: 2026-10-08T07:10:48Z synchronized_at: 2026-10-08T10:15:21Z
<!-- gitea-wiki-mirror:end --> <!-- gitea-wiki-mirror:end -->
<!-- gitea-wiki-mirror:start --> <!-- gitea-wiki-mirror:start -->
@@ -859,3 +859,29 @@ Android 0.9.64 / versionCode 77,源码 `6550b9f`(分支实现,尚未安装
- 每attempt最多保留一份ZIP,已有ZIP不可覆盖;只有无ZIP未截取记录可升级为首次恢复ZIP。以该现场 recordedAt 起保留30天,重复上传不续期;新恢复ZIP采用真实恢复时刻。两端定期清理,手机快照逻辑总量64MiB,淘汰时同时终止上传队列。清理不修改采购结果。 - 每attempt最多保留一份ZIP,已有ZIP不可覆盖;只有无ZIP未截取记录可升级为首次恢复ZIP。以该现场 recordedAt 起保留30天,重复上传不续期;新恢复ZIP采用真实恢复时刻。两端定期清理,手机快照逻辑总量64MiB,淘汰时同时终止上传队列。清理不修改采购结果。
- 管理员采购详情可见完整/部分/未截取/暂无现场数据及安全下载入口;普通采购员、售后、客户端密钥均无原始诊断下载权限。 - 管理员采购详情可见完整/部分/未截取/暂无现场数据及安全下载入口;普通采购员、售后、客户端密钥均无原始诊断下载权限。
- 1500ms、5000节点、2MiB ZIP、8MiB展开、128窗口与深度保护是合成测试覆盖的开发边界,未经过代表性真机样本验收;单次系统Binder读取不能被这些前后检查强行中断。未知/部分必须如实显示,不宣称拿到了系统未暴露的节点。 - 1500ms、5000节点、2MiB ZIP、8MiB展开、128窗口与深度保护是合成测试覆盖的开发边界,未经过代表性真机样本验收;单次系统Binder读取不能被这些前后检查强行中断。未知/部分必须如实显示,不宣称拿到了系统未暴露的节点。
## 图搜采集部分完成的自动关联资格(#368)
实现基线:工单分支 `fix/368-image-search-partial-link` 提交 `15ccab3`(2026-10-08);尚未合并 main 或发布,生产是否启用以发布记录为准。
- 此规则仅作用于 `source=image_search` 的结果后置自动关联。正常 `completed` 沿用原行为,包括 `missing_json=[]`。
- `completed_partial` 仅在 `missing_json` 为可解析、非空字符串数组且每一项精确属于 `reviewCount`、`salesText`、`shopName` 时允许自动关联。三项均为可选描述信息。
- partial 的空值、JSON null、空数组、非法 JSON、非字符串项、空项、未知项、颜色/尺码/价格/不支持维度缺失一律不放宽。不修剪字段名、不按前缀推测。
- 保留原图搜快照、目标 PDD active、`OriginalPDDProductID` 乐观比较;人工改动后的关联不被覆盖。只有实际写入关联并返回 Linked,才进入既有规格同步和匹配流程。
- 采集任务仍为 `completed_partial`,missing 保留,界面仍显示部分完成;关联成功不等于 AI 匹配成功。匹配无结果或不可用沿用原处理,不自动创建采购。
- 不改变普通采集、手动当前页采集、手动关联、价格保护、映射算法、同步后置 hook 的耗时及重放机制。同图多结果沿用原提交顺序和乐观比较,不新增历史筛选、补关联或回退旧商品。
- 验证边界:自动化测试使用隔离 SQLite 与本地模拟 ERP 服务验证实际 SubmitResult、关联及映射落库;未以此宣称 MySQL、真实 AI、真机或线上发布验证。请求在匹配开始后取消时,匹配使用原有脱钩 context 继续;最后 Detail 读取仍使用原请求 context,可能向断开的调用端返回错误,已提交数据不因此回滚。
## Android 本设备最后一批重购(#367)
实现绑定 `5430dc1f64fc0a4d8ef124ada7356e21a975151f`(2026-10-08,工单分支 `feat/367-last-batch-repurchase`);未合并 main、未安装或发布,真机批量采购仍待独立授权验收。
- 采购页搜索栏的“回填”入口替换为“重购”。原回填弹窗、Service 动作和实现保留但不再由该按钮触发;原单条重试、继续采购不改。重购使用原任务的 `/reset`,保留 CG 编号,不使用创建替代任务的 `/retry`。
- 读取本设备近30天全部状态、不继承页面编号/状态筛选;按创建时间倒序、同刻任务ID倒序,从最新记录逐项比较相邻时间,不超过60秒归入同一推算批次。跨50条分页去重,最多200页,结束复核首页;头部或总量变化、日期/页结构异常、查询截断无法证明完整时不启动,不拿部分列表执行。批次时间边界使用独享API客户端的 HTTP Date,不使用手机时钟,缺少有效 Date 拒绝;此轻量方案不保证任意并发删除/归属变动的数据库快照。
- 最后一批没有失败或失败均不可重试时只提示,不回退更早批次。一次确认显示范围、总数、失败/可重购/跳过数量及真实待付款订单风险;系统永不付款。每轮冻结 `failed && retryable` 候选,每条只重置一次,本轮再次失败不重新入队;下次点击重新发现并确认剩余任务。
- 每条前重新检查当前资格、设备身份/Origin、网络、无障碍、新工作与本地执行/结果上传边界;不额外要求已存在完整映射,沿用服务端 unresolved 规格探测资格。明确业务资格拒绝可跳过;环境、规则、鉴权或结果未知停止后续。
- 每条固定 requestId;重置响应不明最多同ID发两次,仍不明确停止,不换UUID或把旧failed认成本次结果。本地短时忙碌且明确未发出请求时最多等待10秒复核;重放pending不是实时状态。后续 attempt 只有原执行器记录到连续探测/重匹配完成→正式采购启动的证据链才属于本轮,外部同任务另起attempt停止编排。
- Service 独立串行后台调度观察,真正领取、执行、上传仍由原执行器和租约/任务锁负责。终态且执行锁、活跃记录和Outbox均收尾后才继续。检测新采集或采购工作停止剩余,本轮当前task本身不算新工作。
- 页面和通知栏“停止重购”只停止后续,已提交/执行的任务沿用原流程;切Tab、切PDD、界面重建不中断。单条单调时钟等待上限15分钟,超时只停止编排,不取消任务、不重置或再次下单。进程死亡/服务重建不恢复整轮,原单任务恢复不改。
- 本轮摘要与最小attempt关联仅在内存,不新增地址/树采集或数据库。接口、服务端、Web、权限及 #364 诊断契约不变。自动化回归和APK构建通过不代表真机页面/批量下单已验收。