fix(android): isolate rejected purchase results and restore sync #377
This commit is contained in:
+16
-1
@@ -1,15 +1,30 @@
|
||||
package cn.ilapage.goauto.agent.persistence
|
||||
|
||||
import cn.ilapage.goauto.agent.network.AgentApiException
|
||||
|
||||
object PurchaseRejectionPolicy {
|
||||
val codes = setOf("PURCHASE_STATE_CONFLICT", "PURCHASE_LEASE_EXPIRED", "PURCHASE_RESULT_CONFLICT")
|
||||
fun terminalCode(error: Exception): String? = (error as? AgentApiException)
|
||||
?.takeIf { it.status == 409 && it.code in codes }?.code
|
||||
}
|
||||
|
||||
/** Uploads already-persisted results only; it never invokes device automation. */
|
||||
class PurchaseOutboxUploader(
|
||||
private val pending: () -> List<PendingPurchaseOutbox>,
|
||||
private val submit: (PendingPurchaseOutbox) -> Unit,
|
||||
private val markUploaded: (PendingPurchaseOutbox) -> Unit,
|
||||
private val afterUploaded: (PendingPurchaseOutbox) -> Unit = {},
|
||||
private val markRejected: (PendingPurchaseOutbox, String) -> Unit = { _, _ -> error("Rejection persistence not configured") },
|
||||
) {
|
||||
fun flush() {
|
||||
pending().forEach { item ->
|
||||
submit(item)
|
||||
try {
|
||||
submit(item)
|
||||
} catch (error: Exception) {
|
||||
val code = PurchaseRejectionPolicy.terminalCode(error) ?: throw error
|
||||
markRejected(item, code)
|
||||
return@forEach
|
||||
}
|
||||
markUploaded(item)
|
||||
afterUploaded(item)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,27 @@
|
||||
package cn.ilapage.goauto.agent.persistence
|
||||
|
||||
data class RejectedPurchaseResult(
|
||||
val id: Long, val taskId: Long, val attemptId: String, val payloadJson: String,
|
||||
val rejectedAt: Long, val errorCode: String, val acknowledgedAt: Long?,
|
||||
)
|
||||
|
||||
internal object PurchaseRejectionSql {
|
||||
val migration = listOf(
|
||||
"ALTER TABLE purchase_outbox ADD COLUMN rejected_at INTEGER",
|
||||
"ALTER TABLE purchase_outbox ADD COLUMN rejection_error_code TEXT",
|
||||
"ALTER TABLE purchase_outbox ADD COLUMN acknowledged_at INTEGER",
|
||||
)
|
||||
const val activeTask = "SELECT task_id FROM purchase_task WHERE upload_status != 'rejected' AND (status IN ('running','order_submit_started') OR upload_status='pending') ORDER BY updated_at ASC LIMIT 1"
|
||||
const val acknowledge = "UPDATE purchase_outbox SET acknowledged_at=? WHERE id=? AND upload_status='rejected' AND acknowledged_at IS NULL"
|
||||
|
||||
/** The caller must wrap both writes in one transaction. Zero task rows means a newer attempt replaced it. */
|
||||
fun reject(item: PendingPurchaseOutbox, code: String, now: Long, update: (String, List<Any>) -> Int) {
|
||||
require(code in PurchaseRejectionPolicy.codes)
|
||||
check(update(
|
||||
"UPDATE purchase_outbox SET upload_status='rejected',rejected_at=?,rejection_error_code=?,updated_at=? WHERE id=? AND task_id=? AND attempt_id=? AND upload_status='pending'",
|
||||
listOf(now, code, now, item.id, item.taskId, item.attemptId),
|
||||
) == 1)
|
||||
update("UPDATE purchase_task SET upload_status='rejected',updated_at=? WHERE task_id=? AND attempt_id=?",
|
||||
listOf(now, item.taskId, item.attemptId))
|
||||
}
|
||||
}
|
||||
@@ -66,6 +66,7 @@ class PurchaseTaskStore(context: Context) : SQLiteOpenHelper(context, DATABASE_N
|
||||
)
|
||||
db.execSQL("CREATE INDEX idx_purchase_outbox_pending ON purchase_outbox(upload_status, id)")
|
||||
db.execSQL("CREATE INDEX idx_purchase_attempt_task ON purchase_attempt(task_id, created_at)")
|
||||
PurchaseRejectionSql.migration.forEach(db::execSQL)
|
||||
}
|
||||
|
||||
override fun onUpgrade(db: SQLiteDatabase, oldVersion: Int, newVersion: Int) {
|
||||
@@ -74,6 +75,7 @@ class PurchaseTaskStore(context: Context) : SQLiteOpenHelper(context, DATABASE_N
|
||||
db.execSQL("ALTER TABLE purchase_task ADD COLUMN final_confirmation_json TEXT")
|
||||
db.execSQL("ALTER TABLE purchase_task ADD COLUMN irreversible_at INTEGER")
|
||||
}
|
||||
if (oldVersion < 3 && newVersion >= 3) PurchaseRejectionSql.migration.forEach(db::execSQL)
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
@@ -225,8 +227,8 @@ class PurchaseTaskStore(context: Context) : SQLiteOpenHelper(context, DATABASE_N
|
||||
|
||||
@Synchronized
|
||||
fun activeTaskId(): Long? = readableDatabase.rawQuery(
|
||||
"SELECT task_id FROM purchase_task WHERE status IN (?,?) OR upload_status=? ORDER BY updated_at ASC LIMIT 1",
|
||||
arrayOf(STATUS_RUNNING, STATUS_ORDER_SUBMIT_STARTED, UPLOAD_PENDING),
|
||||
PurchaseRejectionSql.activeTask,
|
||||
null,
|
||||
).use { cursor -> if (cursor.moveToFirst()) cursor.getLong(0) else null }
|
||||
|
||||
@Synchronized
|
||||
@@ -243,9 +245,47 @@ class PurchaseTaskStore(context: Context) : SQLiteOpenHelper(context, DATABASE_N
|
||||
}
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
fun markRejected(item: PendingPurchaseOutbox, code: String) {
|
||||
val db = writableDatabase
|
||||
db.beginTransaction()
|
||||
try {
|
||||
PurchaseRejectionSql.reject(item, code, System.currentTimeMillis()) { sql, args ->
|
||||
db.compileStatement(sql).use { statement ->
|
||||
args.forEachIndexed { index, value ->
|
||||
if (value is Long) statement.bindLong(index + 1, value)
|
||||
else statement.bindString(index + 1, value.toString())
|
||||
}
|
||||
statement.executeUpdateDelete()
|
||||
}
|
||||
}
|
||||
db.setTransactionSuccessful()
|
||||
} finally { db.endTransaction() }
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
fun rejectedResults(): List<RejectedPurchaseResult> = readableDatabase.rawQuery(
|
||||
"SELECT id,task_id,attempt_id,payload_json,rejected_at,rejection_error_code,acknowledged_at FROM purchase_outbox WHERE upload_status='rejected' ORDER BY id DESC", null,
|
||||
).use { cursor -> buildList {
|
||||
while (cursor.moveToNext()) add(RejectedPurchaseResult(
|
||||
cursor.getLong(0), cursor.getLong(1), cursor.getString(2), cursor.getString(3),
|
||||
cursor.getLong(4), cursor.getString(5), if (cursor.isNull(6)) null else cursor.getLong(6),
|
||||
))
|
||||
} }
|
||||
|
||||
@Synchronized
|
||||
fun rejectedCount(): Int = readableDatabase.rawQuery(
|
||||
"SELECT COUNT(*) FROM purchase_outbox WHERE upload_status='rejected'", null,
|
||||
).use { it.moveToFirst(); it.getInt(0) }
|
||||
|
||||
@Synchronized
|
||||
fun acknowledgeRejection(outboxId: Long) {
|
||||
writableDatabase.execSQL(PurchaseRejectionSql.acknowledge, arrayOf(System.currentTimeMillis(), outboxId))
|
||||
}
|
||||
|
||||
companion object {
|
||||
private const val DATABASE_NAME = "goauto_purchase.db"
|
||||
private const val DATABASE_VERSION = 2
|
||||
private const val DATABASE_VERSION = 3
|
||||
private const val STATUS_RUNNING = "running"
|
||||
const val STATUS_ORDER_SUBMIT_STARTED = "order_submit_started"
|
||||
private const val STATUS_COMPLETED = "completed"
|
||||
|
||||
+48
-26
@@ -268,12 +268,13 @@ class AgentForegroundService : Service() {
|
||||
}
|
||||
|
||||
private fun synchronizeAgent(): String {
|
||||
var heartbeatSucceeded = false
|
||||
return try {
|
||||
val serverUrl = ServerUrlPolicy.normalize(settingsStore.serverUrl())
|
||||
val api = AgentApiClient(serverUrl)
|
||||
var credentials = identityStore.credentials()
|
||||
stateStore.update("CONNECTING", "正在连接并校验设备身份", credentials?.deviceId ?: 0L, credentials != null)
|
||||
updateNotification("正在连接")
|
||||
stateStore.update("CONNECTING", "正在同步", credentials?.deviceId ?: 0L, credentials != null, connection = true)
|
||||
updateNotification("正在同步")
|
||||
|
||||
if (!registeredThisProcess.get()) {
|
||||
val registration = api.register(deviceInfo(), credentials?.token)
|
||||
@@ -290,17 +291,30 @@ class AgentForegroundService : Service() {
|
||||
}
|
||||
|
||||
val activeCredentials = credentials ?: error("设备尚未取得认证凭据")
|
||||
if (taskMutex.currentTaskId() == null) {
|
||||
recoverInterruptedPurchases(api, activeCredentials.token)
|
||||
flushPurchaseOutbox(api, activeCredentials.token)
|
||||
runningTaskId.set(purchaseStore.activeTaskId())
|
||||
val sync = AgentSyncCycle(
|
||||
heartbeat = { currentTaskId ->
|
||||
api.heartbeat(activeCredentials.token, currentTaskId, AgentCapabilities.supported).also {
|
||||
check(it.deviceId == activeCredentials.deviceId) { "心跳设备身份不一致" }
|
||||
}
|
||||
},
|
||||
activeTaskId = {
|
||||
if (taskMutex.currentTaskId() != null) runningTaskId.get()
|
||||
else purchaseStore.activeTaskId().also(runningTaskId::set)
|
||||
},
|
||||
executionActive = { taskMutex.currentTaskId() != null },
|
||||
recover = { recoverInterruptedPurchases(api, activeCredentials.token) },
|
||||
flush = { flushPurchaseOutbox(api, activeCredentials.token) },
|
||||
hasPending = { purchaseStore.pendingOutbox().isNotEmpty() },
|
||||
).run()
|
||||
if (sync.failureCode != null) {
|
||||
cancelIdleReturn("同步未完成")
|
||||
val message = cn.ilapage.goauto.agent.ui.PurchaseRejectionPresentation.connection(sync.failureCode, 0)
|
||||
stateStore.update(sync.failureCode, message, activeCredentials.deviceId, true, connection = true)
|
||||
updateNotification(message)
|
||||
return if (sync.failureCode == "AUTH_ERROR") MANUAL_AUTH_ERROR else MANUAL_NETWORK_ERROR
|
||||
}
|
||||
val heartbeat = api.heartbeat(
|
||||
activeCredentials.token,
|
||||
currentTaskId = runningTaskId.get(),
|
||||
capabilities = AgentCapabilities.supported,
|
||||
)
|
||||
check(heartbeat.deviceId == activeCredentials.deviceId) { "心跳返回了不同的设备身份" }
|
||||
val heartbeat = requireNotNull(sync.heartbeat)
|
||||
heartbeatSucceeded = true
|
||||
stateStore.update(
|
||||
code = if (heartbeat.busy) "BUSY" else "ONLINE",
|
||||
message = if (heartbeat.busy) "设备在线,正在执行任务" else "设备在线空闲",
|
||||
@@ -309,30 +323,32 @@ class AgentForegroundService : Service() {
|
||||
heartbeat = true,
|
||||
)
|
||||
updateNotification(if (heartbeat.busy) "在线 · 执行中" else "在线 · 空闲")
|
||||
scheduleTask(api, activeCredentials.token)
|
||||
if (sync.canClaim) scheduleTask(api, activeCredentials.token) else MANUAL_BUSY
|
||||
} catch (error: AgentApiException) {
|
||||
cancelIdleReturn("网络或服务端请求失败")
|
||||
val authenticationError = error.code in setOf(
|
||||
"DEVICE_TOKEN_INVALID", "DEVICE_INSTALL_ID_CONFLICT", "DEVICE_DISABLED",
|
||||
)
|
||||
val authenticationError = isAgentAuthenticationError(error)
|
||||
val connectionFailure = syncConnectionFailureCode(error, heartbeatSucceeded)
|
||||
val code = connectionFailure ?: "TASK_ERROR"
|
||||
val storedCredentials = runCatching { identityStore.credentials() }.getOrNull()
|
||||
stateStore.update(
|
||||
code = if (authenticationError) "AUTH_ERROR" else "NETWORK_ERROR",
|
||||
message = "${error.code}:${error.message}",
|
||||
code = code,
|
||||
message = if (connectionFailure == null) "任务同步暂未完成,请稍后重试" else cn.ilapage.goauto.agent.ui.PurchaseRejectionPresentation.connection(code, 0),
|
||||
deviceId = storedCredentials?.deviceId ?: 0L,
|
||||
tokenStored = storedCredentials != null,
|
||||
connection = connectionFailure != null,
|
||||
)
|
||||
updateNotification(if (authenticationError) "设备认证失败" else "连接失败,等待网络恢复")
|
||||
updateNotification(if (authenticationError) "设备认证失败" else if (heartbeatSucceeded) "任务同步暂未完成" else "连接失败,等待恢复")
|
||||
if (authenticationError) MANUAL_AUTH_ERROR else MANUAL_NETWORK_ERROR
|
||||
} catch (error: Exception) {
|
||||
cancelIdleReturn("Agent 运行异常")
|
||||
val configured = settingsStore.serverUrl().isNotBlank()
|
||||
val storedCredentials = runCatching { identityStore.credentials() }.getOrNull()
|
||||
stateStore.update(
|
||||
code = if (configured) "ERROR" else "CONFIG_REQUIRED",
|
||||
message = error.message ?: "Agent 运行失败",
|
||||
code = if (heartbeatSucceeded) "TASK_ERROR" else if (configured) "ERROR" else "CONFIG_REQUIRED",
|
||||
message = if (configured) "同步暂未完成,请稍后重试" else "请配置服务端",
|
||||
deviceId = storedCredentials?.deviceId ?: 0L,
|
||||
tokenStored = storedCredentials != null,
|
||||
connection = !heartbeatSucceeded,
|
||||
)
|
||||
updateNotification(if (configured) "运行异常" else "等待配置服务端")
|
||||
if (configured) MANUAL_ERROR else MANUAL_CONFIG_REQUIRED
|
||||
@@ -341,8 +357,6 @@ class AgentForegroundService : Service() {
|
||||
|
||||
private fun scheduleTask(api: AgentApiClient, token: String): String {
|
||||
if (taskMutex.currentTaskId() != null) return MANUAL_BUSY
|
||||
recoverInterruptedPurchases(api, token)
|
||||
flushPurchaseOutbox(api, token)
|
||||
val collectionCooldown = activeCollectionCooldown()
|
||||
val purchaseTask = api.nextPurchaseTask(token)
|
||||
when (TaskDispatchPolicy.decide(purchaseTask?.status, collectionCooldown != null)) {
|
||||
@@ -937,9 +951,10 @@ class AgentForegroundService : Service() {
|
||||
)
|
||||
}
|
||||
}
|
||||
val message = if (outcome.resultType == "failed") "${outcome.errorCode}:${outcome.message}" else outcome.message
|
||||
val rejected = purchaseStore.rejectedResults().any { it.taskId == task.taskId && it.attemptId == task.taskAttemptId }
|
||||
val message = if (rejected) "采购结果被服务端拒收,请在采购页或设置中查看" else if (outcome.resultType == "failed") "${outcome.errorCode}:${outcome.message}" else outcome.message
|
||||
stateStore.update(if (outcome.resultType == "failed") "TASK_ERROR" else "ONLINE", message, tokenStored = true)
|
||||
updateNotification(if (outcome.resultType == "failed") "$taskLabel #${task.taskId} 失败" else "$taskLabel #${task.taskId} 已提交")
|
||||
updateNotification(if (rejected) "$taskLabel #${task.taskId} 结果被拒收" else if (outcome.resultType == "failed") "$taskLabel #${task.taskId} 失败" else "$taskLabel #${task.taskId} 已提交")
|
||||
} catch (error: AgentApiException) {
|
||||
failureSnapshot("failed", error.code.takeIf { it.matches(Regex("[A-Z][A-Z0-9_]{0,63}")) } ?: "PURCHASE_API_FAILED")
|
||||
stateStore.update("TASK_ERROR", "${error.code}:${error.message}", tokenStored = true)
|
||||
@@ -1029,9 +1044,10 @@ class AgentForegroundService : Service() {
|
||||
}
|
||||
val requestId = UUID.randomUUID().toString()
|
||||
purchaseStore.completeAndEnqueue(interrupted.taskId, interrupted.attemptId, requestId, purchaseResultPayload(requestId, interrupted.attemptId, outcome))
|
||||
} catch (_: AgentApiException) {
|
||||
} catch (error: AgentApiException) {
|
||||
// Keep the local irreversible marker. The same server request ID
|
||||
// is replayed after connectivity returns; the order click is never repeated.
|
||||
if (isAgentAuthenticationError(error)) throw error
|
||||
}
|
||||
return@forEach
|
||||
}
|
||||
@@ -1066,6 +1082,11 @@ class AgentForegroundService : Service() {
|
||||
runCatching { onFirstSubmitted(item, requireNotNull(acknowledgement)) }
|
||||
.onFailure { probeHandoff.invalidate(); logDiagnosticPersistenceFailure(it) }
|
||||
},
|
||||
markRejected = { item, code ->
|
||||
probeHandoff.invalidate()
|
||||
purchaseStore.markRejected(item, code)
|
||||
sendBroadcast(Intent(ACTION_PURCHASE_REJECTIONS_CHANGED).setPackage(packageName))
|
||||
},
|
||||
).flush()
|
||||
runningTaskId.set(purchaseStore.activeTaskId())
|
||||
}
|
||||
@@ -1693,6 +1714,7 @@ class AgentForegroundService : Service() {
|
||||
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"
|
||||
const val ACTION_PURCHASE_REJECTIONS_CHANGED = "cn.ilapage.goauto.agent.PURCHASE_REJECTIONS_CHANGED"
|
||||
@Volatile var repurchaseState = RepurchaseState()
|
||||
private set
|
||||
const val ACTION_BACKFILL_START = "cn.ilapage.goauto.agent.BACKFILL_START"
|
||||
|
||||
@@ -192,16 +192,19 @@ class AgentSettingsStore(context: Context) {
|
||||
class AgentStateStore(context: Context) {
|
||||
private val preferences = context.getSharedPreferences(PREFERENCES, Context.MODE_PRIVATE)
|
||||
|
||||
fun update(code: String, message: String, deviceId: Long? = null, tokenStored: Boolean? = null, heartbeat: Boolean = false) {
|
||||
fun update(code: String, message: String, deviceId: Long? = null, tokenStored: Boolean? = null, heartbeat: Boolean = false, connection: Boolean = false) {
|
||||
val editor = preferences.edit()
|
||||
.putString(STATE_CODE, code)
|
||||
.putString(STATE_MESSAGE, message)
|
||||
deviceId?.let { editor.putLong(DEVICE_ID, it) }
|
||||
tokenStored?.let { editor.putBoolean(TOKEN_STORED, it) }
|
||||
if (heartbeat) editor.putLong(LAST_HEARTBEAT, System.currentTimeMillis())
|
||||
if (heartbeat || connection) editor.putString(CONNECTION_CODE, code)
|
||||
editor.apply()
|
||||
}
|
||||
|
||||
fun connectionCode(): String = preferences.getString(CONNECTION_CODE, "STOPPED") ?: "STOPPED"
|
||||
|
||||
fun read(): AgentState = AgentState(
|
||||
code = preferences.getString(STATE_CODE, "STOPPED") ?: "STOPPED",
|
||||
message = preferences.getString(STATE_MESSAGE, "前台服务尚未启动") ?: "前台服务尚未启动",
|
||||
@@ -279,6 +282,7 @@ class AgentStateStore(context: Context) {
|
||||
private companion object {
|
||||
const val PREFERENCES = "goauto_agent_runtime"
|
||||
const val STATE_CODE = "state_code"
|
||||
const val CONNECTION_CODE = "connection_code"
|
||||
const val STATE_MESSAGE = "state_message"
|
||||
const val DEVICE_ID = "device_id"
|
||||
const val LAST_HEARTBEAT = "last_heartbeat"
|
||||
|
||||
@@ -0,0 +1,64 @@
|
||||
package cn.ilapage.goauto.agent.service
|
||||
|
||||
import cn.ilapage.goauto.agent.network.AgentApiException
|
||||
import cn.ilapage.goauto.agent.network.HeartbeatResult
|
||||
|
||||
internal fun isAgentAuthenticationError(error: Throwable?): Boolean = error is AgentApiException &&
|
||||
(error.status in setOf(401, 403) || error.code in setOf("DEVICE_TOKEN_INVALID", "DEVICE_INSTALL_ID_CONFLICT", "DEVICE_DISABLED"))
|
||||
|
||||
internal fun syncConnectionFailureCode(error: Throwable, heartbeatSucceeded: Boolean): String? = when {
|
||||
isAgentAuthenticationError(error) -> "AUTH_ERROR"
|
||||
heartbeatSucceeded -> null
|
||||
error is AgentApiException && error.code == "DEVICE_TASK_MISMATCH" -> "TASK_MISMATCH"
|
||||
error is AgentApiException && error.status >= 500 -> "SERVER_ERROR"
|
||||
else -> "NETWORK_ERROR"
|
||||
}
|
||||
|
||||
internal data class AgentSyncResult(
|
||||
val heartbeat: HeartbeatResult?,
|
||||
val failureCode: String?,
|
||||
val canClaim: Boolean,
|
||||
)
|
||||
|
||||
/** Run under the service's existing sync exclusion; device execution keeps its task mutex. */
|
||||
internal class AgentSyncCycle(
|
||||
private val heartbeat: (Long?) -> HeartbeatResult,
|
||||
private val activeTaskId: () -> Long?,
|
||||
private val executionActive: () -> Boolean,
|
||||
private val recover: () -> Unit,
|
||||
private val flush: () -> Unit,
|
||||
private val hasPending: () -> Boolean,
|
||||
) {
|
||||
fun run(): AgentSyncResult {
|
||||
var pulse = runCatching { heartbeat(activeTaskId()) }
|
||||
var authFailure = isAgentAuthenticationError(pulse.exceptionOrNull())
|
||||
var recovered = false
|
||||
var flushed = false
|
||||
if (!authFailure && !executionActive()) {
|
||||
val recovery = runCatching(recover)
|
||||
recovered = recovery.isSuccess
|
||||
authFailure = isAgentAuthenticationError(recovery.exceptionOrNull())
|
||||
if (!authFailure) {
|
||||
val upload = runCatching(flush)
|
||||
flushed = upload.isSuccess
|
||||
authFailure = isAgentAuthenticationError(upload.exceptionOrNull())
|
||||
}
|
||||
if (!authFailure && (pulse.exceptionOrNull() as? AgentApiException)?.code == "DEVICE_TASK_MISMATCH") {
|
||||
pulse = runCatching { heartbeat(activeTaskId()) }
|
||||
}
|
||||
}
|
||||
val failure = pulse.exceptionOrNull()
|
||||
return AgentSyncResult(
|
||||
pulse.getOrNull(),
|
||||
when {
|
||||
authFailure -> "AUTH_ERROR"
|
||||
failure == null -> null
|
||||
isAgentAuthenticationError(failure) -> "AUTH_ERROR"
|
||||
(failure as? AgentApiException)?.code == "DEVICE_TASK_MISMATCH" -> "TASK_MISMATCH"
|
||||
failure is AgentApiException && failure.status >= 500 -> "SERVER_ERROR"
|
||||
else -> "NETWORK_ERROR"
|
||||
},
|
||||
pulse.isSuccess && !authFailure && recovered && flushed && !executionActive() && !hasPending() && activeTaskId() == null,
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -24,6 +24,7 @@ import cn.ilapage.goauto.agent.network.CollectionHistoryItem
|
||||
import cn.ilapage.goauto.agent.network.PurchaseHistoryItem
|
||||
import cn.ilapage.goauto.agent.network.ServerUrlPolicy
|
||||
import cn.ilapage.goauto.agent.persistence.TaskHistoryCache
|
||||
import cn.ilapage.goauto.agent.persistence.PurchaseTaskStore
|
||||
import cn.ilapage.goauto.agent.service.AgentForegroundService
|
||||
import cn.ilapage.goauto.agent.service.AgentSettingsStore
|
||||
import cn.ilapage.goauto.agent.service.AgentStateStore
|
||||
@@ -58,6 +59,7 @@ class AgentSettingsFragment : Fragment() {
|
||||
private lateinit var connectionFeedback: TextView
|
||||
private lateinit var disabledReason: TextView
|
||||
private lateinit var diagnostics: TextView
|
||||
private lateinit var rejectedResultsButton: MaterialButton
|
||||
private lateinit var accessibilityText: TextView
|
||||
private lateinit var collectionIntervalStartLayout: TextInputLayout
|
||||
private lateinit var collectionIntervalStartInput: TextInputEditText
|
||||
@@ -326,6 +328,20 @@ class AgentSettingsFragment : Fragment() {
|
||||
diagnostics = context.label("—", 14f, context.getColor(R.color.agent_text))
|
||||
diagnostics.setPadding(0, context.dp(10), 0, 0)
|
||||
addView(diagnostics)
|
||||
rejectedResultsButton = MaterialButton(context).apply {
|
||||
text = "查看拒收结果"
|
||||
minHeight = context.dp(48)
|
||||
setOnClickListener {
|
||||
val records = PurchaseTaskStore(context).use { it.rejectedResults() }
|
||||
MaterialAlertDialogBuilder(context).setTitle("采购结果被服务端拒收")
|
||||
.setMessage(records.joinToString("\n\n") {
|
||||
PurchaseRejectionPresentation.notice(it) + "\n${it.errorCode}" +
|
||||
if (it.acknowledgedAt != null) "\n已核对(仅本机记录)" else ""
|
||||
}.ifBlank { "暂无拒收结果" })
|
||||
.setPositiveButton("知道了", null).show()
|
||||
}
|
||||
}
|
||||
addView(rejectedResultsButton, fullWidth(8))
|
||||
}))
|
||||
}
|
||||
return context.page(content)
|
||||
@@ -669,12 +685,14 @@ class AgentSettingsFragment : Fragment() {
|
||||
disabledReason.text = if (busy) "任务执行中,暂时不能修改服务器或设备名称。" else ""
|
||||
|
||||
val installId = runCatching { identityStore.installId() }.getOrElse { "读取失败" }
|
||||
val rejectedCount = PurchaseTaskStore(context).use { it.rejectedCount() }
|
||||
rejectedResultsButton.visibility = if (rejectedCount > 0) View.VISIBLE else View.GONE
|
||||
diagnostics.text = buildString {
|
||||
append("installId:$installId\n")
|
||||
append("Device Token:${if (state.tokenStored) "已配置" else "未配置"}\n")
|
||||
append("注册状态:${if (state.deviceId > 0) "已注册(设备 ${state.deviceId})" else "未注册"}\n")
|
||||
append("Agent 版本:${BuildConfig.VERSION_NAME}\n")
|
||||
append("服务端连接:${if (state.code in setOf("ONLINE", "BUSY", "COLLECTION_COOLDOWN")) "已连接" else "未连接"}\n")
|
||||
append("服务端连接:${PurchaseRejectionPresentation.connection(stateStore.connectionCode(), rejectedCount)}\n")
|
||||
append("保持屏幕常亮:${if (state.keepScreenOn) "已开启" else "仅在任务执行或采集间隔时开启"}")
|
||||
}
|
||||
accessibilityText.text = when (AccessibilityReadinessDetector.current(context)) {
|
||||
|
||||
@@ -0,0 +1,28 @@
|
||||
package cn.ilapage.goauto.agent.ui
|
||||
|
||||
import cn.ilapage.goauto.agent.persistence.RejectedPurchaseResult
|
||||
import org.json.JSONObject
|
||||
|
||||
internal object PurchaseRejectionPresentation {
|
||||
private fun payload(item: RejectedPurchaseResult) = runCatching { JSONObject(item.payloadJson) }.getOrNull()
|
||||
private fun orderNo(item: RejectedPurchaseResult): String? = payload(item)?.optString("pddOrderNo")
|
||||
?.takeIf { it.matches(Regex("[0-9A-Za-z-]{1,80}")) }
|
||||
fun showBanner(item: RejectedPurchaseResult): Boolean = item.acknowledgedAt == null &&
|
||||
(payload(item)?.optString("resultType") in setOf("order_created", "order_result_unknown") ||
|
||||
!payload(item)?.optString("pddOrderNo").isNullOrBlank())
|
||||
fun notice(item: RejectedPurchaseResult): String = buildString {
|
||||
append("采购任务 #${item.taskId} 的结果被服务端拒收")
|
||||
orderNo(item)?.let { append("\n拼多多订单号:$it") }
|
||||
append("\n请人工核对拼多多订单及后台任务,避免重复采购")
|
||||
}
|
||||
fun connection(code: String, rejected: Int): String = when (code) {
|
||||
"CONNECTING" -> "正在同步"
|
||||
"ONLINE", "BUSY", "COLLECTION_COOLDOWN", "TASK_ERROR" -> if (rejected > 0) "已连接 · 有 $rejected 条采购结果被服务端拒收" else "已连接"
|
||||
"AUTH_ERROR" -> "未连接 · 设备认证失败"
|
||||
"TASK_MISMATCH" -> "未连接 · 设备任务状态不一致"
|
||||
"NETWORK_ERROR" -> "未连接 · 网络暂不可用"
|
||||
"SERVER_ERROR" -> "未连接 · 服务端暂不可用"
|
||||
"CONFIG_REQUIRED" -> "未连接 · 请配置服务端"
|
||||
else -> "未连接 · 同步暂未完成"
|
||||
}
|
||||
}
|
||||
@@ -37,6 +37,8 @@ import cn.ilapage.goauto.agent.network.HistoryColorImage
|
||||
import cn.ilapage.goauto.agent.network.PurchaseHistoryDetail
|
||||
import cn.ilapage.goauto.agent.network.PurchaseHistoryItem
|
||||
import cn.ilapage.goauto.agent.persistence.TaskHistoryCache
|
||||
import cn.ilapage.goauto.agent.persistence.PurchaseTaskStore
|
||||
import cn.ilapage.goauto.agent.persistence.RejectedPurchaseResult
|
||||
import cn.ilapage.goauto.agent.service.AgentForegroundService
|
||||
import cn.ilapage.goauto.agent.service.AgentSettingsStore
|
||||
import cn.ilapage.goauto.agent.service.AgentStateStore
|
||||
@@ -146,8 +148,14 @@ class TaskHistoryFragment : Fragment() {
|
||||
private var repurchaseButton: MaterialButton? = null
|
||||
private var repurchaseDialog: androidx.appcompat.app.AlertDialog? = null
|
||||
private var repurchaseDialogRound: String? = null
|
||||
private var rejectionPanel: LinearLayout? = null
|
||||
private var displayedRejections: List<RejectedPurchaseResult>? = null
|
||||
private val currentPageReceiver = object : BroadcastReceiver() {
|
||||
override fun onReceive(context: Context?, intent: Intent?) {
|
||||
if (intent?.action == AgentForegroundService.ACTION_PURCHASE_REJECTIONS_CHANGED) {
|
||||
renderRejections()
|
||||
return
|
||||
}
|
||||
if (intent?.action == AgentForegroundService.ACTION_REPURCHASE_STATE) {
|
||||
renderRepurchase()
|
||||
if (!collection && isResumed && AgentForegroundService.repurchaseState.phase == cn.ilapage.goauto.agent.service.RepurchasePhase.FINISHED) {
|
||||
@@ -200,6 +208,10 @@ class TaskHistoryFragment : Fragment() {
|
||||
val context = requireContext()
|
||||
pageColumn = context.column()
|
||||
pageColumn.addView(context.screenTitle(if (collection) "采集记录" else "采购记录"))
|
||||
if (!collection) {
|
||||
rejectionPanel = context.column(0)
|
||||
pageColumn.addView(rejectionPanel)
|
||||
}
|
||||
pageColumn.addView(buildSearch())
|
||||
if (!collection) {
|
||||
repurchasePanel = context.column(0)
|
||||
@@ -228,6 +240,7 @@ class TaskHistoryFragment : Fragment() {
|
||||
|
||||
override fun onResume() {
|
||||
super.onResume()
|
||||
renderRejections()
|
||||
renderRepurchase()
|
||||
detailState.taskId?.let { taskId ->
|
||||
if (collection) loadCollectionDetail(taskId) else loadPurchaseDetail(taskId)
|
||||
@@ -235,6 +248,8 @@ class TaskHistoryFragment : Fragment() {
|
||||
}
|
||||
|
||||
override fun onDestroyView() {
|
||||
rejectionPanel = null
|
||||
displayedRejections = null
|
||||
repurchaseDialog?.setOnDismissListener(null)
|
||||
repurchaseDialog?.dismiss()
|
||||
repurchaseDialog = null
|
||||
@@ -247,6 +262,29 @@ class TaskHistoryFragment : Fragment() {
|
||||
super.onDestroyView()
|
||||
}
|
||||
|
||||
private fun renderRejections() {
|
||||
val panel = rejectionPanel ?: return
|
||||
val context = context ?: return
|
||||
val records = PurchaseTaskStore(context).use { it.rejectedResults() }.filter(PurchaseRejectionPresentation::showBanner)
|
||||
if (records == displayedRejections) return
|
||||
displayedRejections = records
|
||||
panel.removeAllViews()
|
||||
records.forEach { record ->
|
||||
panel.addView(context.card(context.cardColumn().apply {
|
||||
addView(context.label(PurchaseRejectionPresentation.notice(record), 14f, context.getColor(R.color.agent_warning), true))
|
||||
addView(MaterialButton(context).apply {
|
||||
text = "已核对"
|
||||
minHeight = context.dp(48)
|
||||
contentDescription = "采购任务 #${record.taskId} 已核对,仅隐藏本机提醒"
|
||||
setOnClickListener {
|
||||
val saved = runCatching { PurchaseTaskStore(context).use { it.acknowledgeRejection(record.id) } }.isSuccess
|
||||
if (saved) renderRejections() else toast("暂时无法保存核对记录,请稍后重试")
|
||||
}
|
||||
}, fullWidth(8))
|
||||
}))
|
||||
}
|
||||
}
|
||||
|
||||
override fun onDestroy() {
|
||||
unregisterCurrentPageReceiver()
|
||||
imageLoader.close()
|
||||
@@ -811,6 +849,7 @@ class TaskHistoryFragment : Fragment() {
|
||||
val filter = IntentFilter(AgentForegroundService.ACTION_CURRENT_PAGE_RESULT)
|
||||
filter.addAction(AgentForegroundService.ACTION_BACKFILL_STATE)
|
||||
filter.addAction(AgentForegroundService.ACTION_REPURCHASE_STATE)
|
||||
filter.addAction(AgentForegroundService.ACTION_PURCHASE_REJECTIONS_CHANGED)
|
||||
if (Build.VERSION.SDK_INT >= 33) {
|
||||
requireContext().registerReceiver(currentPageReceiver, filter, Context.RECEIVER_NOT_EXPORTED)
|
||||
} else {
|
||||
@@ -893,6 +932,11 @@ class TaskHistoryFragment : Fragment() {
|
||||
private fun renderPurchaseDetail(detail: PurchaseHistoryDetail) {
|
||||
val context = requireContext()
|
||||
val task = detail.task
|
||||
PurchaseTaskStore(context).use { it.rejectedResults() }.filter { it.taskId == task.taskId }.forEach { record ->
|
||||
resultColumn.addView(context.card(context.cardColumn().apply {
|
||||
addView(context.label(PurchaseRejectionPresentation.notice(record) + "\n${record.errorCode}", 14f, context.getColor(R.color.agent_warning)))
|
||||
}))
|
||||
}
|
||||
val info = buildString {
|
||||
append("蝦皮订单号:${task.shopeeOrderNo.ifBlank { "—" }}\n")
|
||||
append("PDD 商品:${task.pddGoodsId}\n")
|
||||
|
||||
@@ -8,6 +8,22 @@ import org.junit.Test
|
||||
class PurchaseOutboxHandoffTest {
|
||||
private val item = PendingPurchaseOutbox(1, 42, "probe", "request", "{}")
|
||||
|
||||
@Test fun `terminal conflict does not block later pending result`() {
|
||||
val submitted = mutableListOf<Long>()
|
||||
val uploaded = mutableListOf<Long>()
|
||||
val result = runCatching {
|
||||
PurchaseOutboxUploader(
|
||||
{ listOf(item, item.copy(id = 2, taskId = 43)) },
|
||||
{ submitted += it.id; if (it.id == 1L) throw cn.ilapage.goauto.agent.network.AgentApiException(409, "PURCHASE_STATE_CONFLICT", "untrusted", false) },
|
||||
{ uploaded += it.id },
|
||||
markRejected = { _, _ -> },
|
||||
).flush()
|
||||
}
|
||||
assertTrue(result.isSuccess)
|
||||
assertEquals(listOf(1L, 2L), submitted)
|
||||
assertEquals(listOf(2L), uploaded)
|
||||
}
|
||||
|
||||
@Test fun `handoff callback runs only after successful submit and local upload mark`() {
|
||||
for (failure in listOf("submit", "mark", "none")) {
|
||||
val events = mutableListOf<String>()
|
||||
@@ -31,4 +47,37 @@ class PurchaseOutboxHandoffTest {
|
||||
PurchaseOutboxUploader({ listOf(item) }, { submitted++ }, {}).flush()
|
||||
assertEquals(1, submitted)
|
||||
}
|
||||
|
||||
@Test fun `only whitelisted 409 is terminal independent of retryable`() {
|
||||
for (code in listOf("PURCHASE_STATE_CONFLICT", "PURCHASE_LEASE_EXPIRED", "PURCHASE_RESULT_CONFLICT", "OTHER")) {
|
||||
for (status in listOf(409, 401, 403, 500)) for (retryable in listOf(false, true)) {
|
||||
val events = mutableListOf<String>()
|
||||
val result = runCatching {
|
||||
PurchaseOutboxUploader({ listOf(item) },
|
||||
{ throw cn.ilapage.goauto.agent.network.AgentApiException(status, code, "secret URL", retryable) },
|
||||
{ events += "uploaded" }, { events += "handoff" },
|
||||
{ _, fixedCode -> events += fixedCode },
|
||||
).flush()
|
||||
}
|
||||
val rejected = status == 409 && code != "OTHER"
|
||||
assertEquals(rejected, result.isSuccess)
|
||||
assertEquals(if (rejected) listOf(code) else emptyList<String>(), events)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Test fun `network and rejected persistence failures retain pending and stop flush`() {
|
||||
for (diskFailure in listOf(false, true)) {
|
||||
val events = mutableListOf<String>()
|
||||
val result = runCatching {
|
||||
PurchaseOutboxUploader({ listOf(item, item.copy(id = 2)) },
|
||||
{ events += "submit"; if (diskFailure) throw cn.ilapage.goauto.agent.network.AgentApiException(409, "PURCHASE_STATE_CONFLICT", "raw", false) else error("network") },
|
||||
{ events += "uploaded" }, { events += "handoff" },
|
||||
{ _, _ -> events += "reject"; error("disk") },
|
||||
).flush()
|
||||
}
|
||||
assertTrue(result.isFailure)
|
||||
assertEquals(if (diskFailure) listOf("submit", "reject") else listOf("submit"), events)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,54 @@
|
||||
package cn.ilapage.goauto.agent
|
||||
|
||||
import cn.ilapage.goauto.agent.network.AgentApiClient
|
||||
import cn.ilapage.goauto.agent.network.AgentApiException
|
||||
import cn.ilapage.goauto.agent.persistence.PendingPurchaseOutbox
|
||||
import cn.ilapage.goauto.agent.persistence.PurchaseOutboxUploader
|
||||
import java.net.ServerSocket
|
||||
import java.util.concurrent.Executors
|
||||
import java.util.concurrent.TimeUnit
|
||||
import org.junit.Assert.*
|
||||
import org.junit.Test
|
||||
|
||||
class PurchaseOutboxHttpTest {
|
||||
@Test fun `missing retryable is not evidence that an unknown 409 is terminal`() {
|
||||
for (code in listOf("PURCHASE_STATE_CONFLICT", "PURCHASE_LEASE_EXPIRED", "PURCHASE_RESULT_CONFLICT", "UNKNOWN_CONFLICT")) {
|
||||
ServerSocket(0).use { server ->
|
||||
server.soTimeout = 3000
|
||||
val executor = Executors.newSingleThreadExecutor()
|
||||
val serving = executor.submit {
|
||||
server.accept().use { socket ->
|
||||
socket.soTimeout = 3000
|
||||
val reader = socket.getInputStream().bufferedReader()
|
||||
var contentLength = 0
|
||||
while (true) {
|
||||
val line = reader.readLine() ?: break
|
||||
if (line.isEmpty()) break
|
||||
if (line.startsWith("Content-Length:", true)) contentLength = line.substringAfter(':').trim().toInt()
|
||||
}
|
||||
repeat(contentLength) { check(reader.read() >= 0) }
|
||||
val response = """{"code":"$code","message":"untrusted diagnostic"}"""
|
||||
socket.getOutputStream().write(("HTTP/1.1 409 Conflict\r\nContent-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}")
|
||||
var pending = true
|
||||
val result = runCatching {
|
||||
PurchaseOutboxUploader(
|
||||
{ listOf(PendingPurchaseOutbox(1,42,"probe","request","{}")) },
|
||||
{ api.submitPurchaseResult(it.taskId,it.payloadJson,"synthetic") },
|
||||
{ fail("409 is never uploaded") },
|
||||
{ fail("409 is never a handoff") },
|
||||
{ _, _ -> pending = false },
|
||||
).flush()
|
||||
}
|
||||
assertEquals(code == "UNKNOWN_CONFLICT", pending)
|
||||
assertEquals(code != "UNKNOWN_CONFLICT", result.isSuccess)
|
||||
if (pending) assertFalse((result.exceptionOrNull() as AgentApiException).retryable)
|
||||
serving.get(5, TimeUnit.SECONDS)
|
||||
} finally { executor.shutdownNow() }
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
+32
@@ -0,0 +1,32 @@
|
||||
package cn.ilapage.goauto.agent
|
||||
|
||||
import cn.ilapage.goauto.agent.persistence.RejectedPurchaseResult
|
||||
import cn.ilapage.goauto.agent.ui.PurchaseRejectionPresentation
|
||||
import org.junit.Assert.*
|
||||
import org.junit.Test
|
||||
|
||||
class PurchaseRejectionPresentationTest {
|
||||
private fun item(payload: String, acknowledged: Long? = null) = RejectedPurchaseResult(1,42,"attempt",payload,100,"PURCHASE_STATE_CONFLICT",acknowledged)
|
||||
@Test fun `only unacknowledged order evidence appears on procurement banner`() {
|
||||
for (payload in listOf("""{"resultType":"order_created"}""", """{"resultType":"order_result_unknown"}""", """{"resultType":"failed","pddOrderNo":"123-456"}""")) {
|
||||
assertTrue(PurchaseRejectionPresentation.showBanner(item(payload)))
|
||||
assertFalse(PurchaseRejectionPresentation.showBanner(item(payload, 200)))
|
||||
}
|
||||
assertFalse(PurchaseRejectionPresentation.showBanner(item("""{"resultType":"spec_probe_completed"}""")))
|
||||
assertFalse(PurchaseRejectionPresentation.showBanner(item("""{"resultType":"spec_probe_completed","pddOrderNo":null}""")))
|
||||
}
|
||||
@Test fun `notice includes task optional order and fixed review instruction without raw errors`() {
|
||||
val message = PurchaseRejectionPresentation.notice(item("""{"resultType":"order_created","pddOrderNo":"123-456","message":"https://secret"}"""))
|
||||
assertTrue(message.contains("#42")); assertTrue(message.contains("123-456"))
|
||||
assertTrue(message.contains("请人工核对拼多多订单及后台任务,避免重复采购"))
|
||||
assertFalse(message.contains("secret"))
|
||||
}
|
||||
@Test fun `connection classifies sync heartbeat and rejected separately`() {
|
||||
assertEquals("正在同步", PurchaseRejectionPresentation.connection("CONNECTING", 2))
|
||||
assertEquals("已连接 · 有 2 条采购结果被服务端拒收", PurchaseRejectionPresentation.connection("ONLINE", 2))
|
||||
assertEquals("未连接 · 设备认证失败", PurchaseRejectionPresentation.connection("AUTH_ERROR", 2))
|
||||
assertEquals("未连接 · 设备任务状态不一致", PurchaseRejectionPresentation.connection("TASK_MISMATCH", 2))
|
||||
assertEquals("未连接 · 网络暂不可用", PurchaseRejectionPresentation.connection("NETWORK_ERROR", 2))
|
||||
assertEquals("未连接 · 服务端暂不可用", PurchaseRejectionPresentation.connection("SERVER_ERROR", 2))
|
||||
}
|
||||
}
|
||||
+100
@@ -0,0 +1,100 @@
|
||||
package cn.ilapage.goauto.agent.persistence
|
||||
|
||||
import java.sql.Connection
|
||||
import java.sql.DriverManager
|
||||
import org.junit.Assert.*
|
||||
import org.junit.Test
|
||||
|
||||
class PurchaseRejectionStoreTest {
|
||||
@Test fun `v2 upgrade adds exactly three nullable columns and preserves all old data`() = database { db ->
|
||||
val before = rows(db, "SELECT * FROM purchase_outbox")
|
||||
val taskBefore = rows(db, "SELECT * FROM purchase_task")
|
||||
migrate(db)
|
||||
assertEquals(3, PurchaseRejectionSql.migration.size)
|
||||
assertEquals(taskBefore, rows(db, "SELECT * FROM purchase_task"))
|
||||
assertEquals(before, rows(db, "SELECT id, task_id, attempt_id, request_id, payload_json, upload_status, created_at, updated_at FROM purchase_outbox"))
|
||||
assertEquals(listOf(listOf(null, null, null)), rows(db, "SELECT rejected_at,rejection_error_code,acknowledged_at FROM purchase_outbox"))
|
||||
val columns = rows(db, "PRAGMA table_info(purchase_outbox)").takeLast(3)
|
||||
assertTrue(columns.all { it[3] == "0" && it[4] == null })
|
||||
}
|
||||
|
||||
@Test fun `reject is atomic preserves payload excludes pending and targets current attempt only`() = database { db ->
|
||||
migrate(db)
|
||||
reject(db)
|
||||
assertEquals(listOf(listOf("rejected", "payload", "100", "PURCHASE_STATE_CONFLICT")), rows(db, "SELECT upload_status,payload_json,rejected_at,rejection_error_code FROM purchase_outbox"))
|
||||
assertEquals(listOf(listOf("completed", "rejected", "payload")), rows(db, "SELECT status,upload_status,result_json FROM purchase_task"))
|
||||
assertTrue(rows(db, PurchaseRejectionSql.activeTask).isEmpty())
|
||||
assertTrue(rows(db, "SELECT id FROM purchase_outbox WHERE upload_status='pending'").isEmpty())
|
||||
}
|
||||
|
||||
@Test fun `old outbox rejection cannot overwrite a newer task attempt`() = database { db ->
|
||||
migrate(db)
|
||||
db.createStatement().use { it.execute("UPDATE purchase_task SET attempt_id='new',status='running',upload_status='none'") }
|
||||
reject(db)
|
||||
assertEquals(listOf(listOf("new", "running", "none")), rows(db, "SELECT attempt_id,status,upload_status FROM purchase_task"))
|
||||
}
|
||||
|
||||
@Test fun `rejected record never counts active even if legacy local status says running`() = database { db ->
|
||||
migrate(db); reject(db)
|
||||
db.createStatement().use { it.execute("UPDATE purchase_task SET status='running'") }
|
||||
assertTrue(rows(db, PurchaseRejectionSql.activeTask).isEmpty())
|
||||
}
|
||||
|
||||
@Test fun `invalid whitelist and mismatched outbox identity cannot change evidence`() = database { db ->
|
||||
migrate(db)
|
||||
for ((item, code) in listOf(
|
||||
PendingPurchaseOutbox(1,42,"old","request","payload") to "UNKNOWN",
|
||||
PendingPurchaseOutbox(1,42,"new","request","payload") to "PURCHASE_STATE_CONFLICT",
|
||||
)) {
|
||||
assertTrue(runCatching { PurchaseRejectionSql.reject(item, code, 100) { sql, args -> update(db, sql, args) } }.isFailure)
|
||||
}
|
||||
assertEquals(listOf(listOf("pending", null, null)), rows(db, "SELECT upload_status,rejected_at,rejection_error_code FROM purchase_outbox"))
|
||||
}
|
||||
|
||||
@Test fun `task update failure rolls back rejection evidence too`() = database { db ->
|
||||
migrate(db)
|
||||
db.createStatement().use { it.execute("CREATE TRIGGER fail_task BEFORE UPDATE ON purchase_task BEGIN SELECT RAISE(ABORT, 'test'); END") }
|
||||
assertTrue(runCatching { reject(db) }.isFailure)
|
||||
assertEquals(listOf(listOf("pending", null, null)), rows(db, "SELECT upload_status,rejected_at,rejection_error_code FROM purchase_outbox"))
|
||||
assertEquals(listOf(listOf("pending")), rows(db, "SELECT upload_status FROM purchase_task"))
|
||||
}
|
||||
|
||||
@Test fun `ack survives database reopen and changes no payload or task state`() {
|
||||
val file = java.io.File.createTempFile("purchase-rejection-", ".db")
|
||||
try {
|
||||
connect(file.absolutePath).use { db ->
|
||||
seed(db); migrate(db); reject(db)
|
||||
update(db, PurchaseRejectionSql.acknowledge, listOf(200L, 1L))
|
||||
}
|
||||
connect(file.absolutePath).use { db ->
|
||||
assertEquals(listOf(listOf("200", "payload", "rejected")), rows(db, "SELECT acknowledged_at,payload_json,upload_status FROM purchase_outbox"))
|
||||
assertEquals(listOf(listOf("completed", "rejected")), rows(db, "SELECT status,upload_status FROM purchase_task"))
|
||||
}
|
||||
} finally { check(file.delete()) }
|
||||
}
|
||||
|
||||
private fun reject(db: Connection) {
|
||||
db.autoCommit = false
|
||||
try {
|
||||
PurchaseRejectionSql.reject(PendingPurchaseOutbox(1,42,"old","request","payload"), "PURCHASE_STATE_CONFLICT", 100) { sql, args -> update(db, sql, args) }
|
||||
db.commit()
|
||||
} catch (error: Exception) { db.rollback(); throw error }
|
||||
finally { db.autoCommit = true }
|
||||
}
|
||||
private fun update(db: Connection, sql: String, args: List<Any>): Int = db.prepareStatement(sql).use { statement ->
|
||||
args.forEachIndexed { index, arg -> statement.setObject(index + 1, arg) }
|
||||
statement.executeUpdate()
|
||||
}
|
||||
private fun migrate(db: Connection) = PurchaseRejectionSql.migration.forEach { sql -> db.createStatement().use { it.execute(sql) } }
|
||||
private fun connect(path: String): Connection { Class.forName("org.sqlite.JDBC"); return DriverManager.getConnection("jdbc:sqlite:$path") }
|
||||
private fun database(block: (Connection) -> Unit) = connect(":memory:").use { seed(it); block(it) }
|
||||
private fun rows(db: Connection, sql: String): List<List<String?>> = db.createStatement().use { statement ->
|
||||
statement.executeQuery(sql).use { r -> buildList { while(r.next()) add((1..r.metaData.columnCount).map { r.getString(it) }) } }
|
||||
}
|
||||
private fun seed(db: Connection) = db.createStatement().use {
|
||||
it.execute("CREATE TABLE purchase_task (task_id INTEGER PRIMARY KEY,attempt_id TEXT NOT NULL,rule_snapshot_hash TEXT NOT NULL,current_step TEXT NOT NULL,status TEXT NOT NULL,result_json TEXT,upload_status TEXT NOT NULL,created_at INTEGER NOT NULL,updated_at INTEGER NOT NULL,order_submit_request_id TEXT,final_confirmation_json TEXT,irreversible_at INTEGER)")
|
||||
it.execute("CREATE TABLE purchase_outbox (id INTEGER PRIMARY KEY AUTOINCREMENT,task_id INTEGER NOT NULL,attempt_id TEXT NOT NULL,request_id TEXT NOT NULL UNIQUE,payload_json TEXT NOT NULL,upload_status TEXT NOT NULL,created_at INTEGER NOT NULL,updated_at INTEGER NOT NULL)")
|
||||
it.execute("INSERT INTO purchase_task VALUES (42,'old','hash','submit_result','completed','payload','pending',1,2,NULL,NULL,NULL)")
|
||||
it.execute("INSERT INTO purchase_outbox VALUES (1,42,'old','request','payload','pending',1,2)")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,115 @@
|
||||
package cn.ilapage.goauto.agent.service
|
||||
|
||||
import cn.ilapage.goauto.agent.network.AgentApiException
|
||||
import cn.ilapage.goauto.agent.network.HeartbeatResult
|
||||
import cn.ilapage.goauto.agent.persistence.PendingPurchaseOutbox
|
||||
import cn.ilapage.goauto.agent.persistence.PurchaseOutboxUploader
|
||||
import org.junit.Assert.*
|
||||
import org.junit.Test
|
||||
|
||||
class AgentSyncCycleTest {
|
||||
private val online = HeartbeatResult(7, true, false, 15)
|
||||
private fun api(status: Int, code: String) = AgentApiException(status, code, "untrusted URL", false)
|
||||
|
||||
@Test fun `stale completed probe heartbeat mismatch rejected then refreshed heartbeat permits next claim`() {
|
||||
val events = mutableListOf<String>()
|
||||
var active: Long? = 42
|
||||
var pending = true
|
||||
val result = AgentSyncCycle(
|
||||
heartbeat = { id -> events += "heartbeat:$id"; if (id != null) throw api(409, "DEVICE_TASK_MISMATCH") else online },
|
||||
activeTaskId = { active }, executionActive = { false },
|
||||
recover = { events += "recover" },
|
||||
flush = {
|
||||
PurchaseOutboxUploader(
|
||||
{ listOf(PendingPurchaseOutbox(1,42,"probe","request","{}")) },
|
||||
{ throw api(409,"PURCHASE_STATE_CONFLICT") },
|
||||
{ fail("rejected probe must not be marked uploaded") },
|
||||
{ fail("rejected probe must not establish handoff") },
|
||||
{ _, _ -> events += "reject-probe"; pending = false; active = null },
|
||||
).flush()
|
||||
},
|
||||
hasPending = { pending },
|
||||
).run()
|
||||
if (result.canClaim) events += "claim"
|
||||
assertEquals(listOf("heartbeat:42", "recover", "reject-probe", "heartbeat:null", "claim"), events)
|
||||
assertEquals(online, result.heartbeat)
|
||||
}
|
||||
|
||||
@Test fun `authentication errors skip recovery and flush including unrecognized HTTP 401 and 403`() {
|
||||
for ((status, code) in listOf(401 to "OTHER", 403 to "OTHER", 409 to "DEVICE_TOKEN_INVALID", 409 to "DEVICE_INSTALL_ID_CONFLICT", 409 to "DEVICE_DISABLED")) {
|
||||
val result = AgentSyncCycle({ throw api(status, code) }, { 42 }, { false },
|
||||
{ fail("recovery after auth failure") }, { fail("flush after auth failure") }, { true }).run()
|
||||
assertFalse(result.canClaim)
|
||||
assertEquals("AUTH_ERROR", result.failureCode)
|
||||
}
|
||||
}
|
||||
|
||||
@Test fun `heartbeat survives recovery or flush failure but claiming stops`() {
|
||||
for (recoverFails in listOf(false, true)) {
|
||||
var flushed = false
|
||||
val result = AgentSyncCycle({ online }, { null }, { false },
|
||||
{ if (recoverFails) error("disk") },
|
||||
{ flushed = true; if (!recoverFails) throw api(500, "INTERNAL") }, { false }).run()
|
||||
assertTrue(flushed)
|
||||
assertEquals(online, result.heartbeat)
|
||||
assertFalse(result.canClaim)
|
||||
}
|
||||
}
|
||||
|
||||
@Test fun `failed heartbeat still flushes but never claims and mismatch retries only once`() {
|
||||
for (code in listOf("DEVICE_TASK_MISMATCH", "INTERNAL")) {
|
||||
var heartbeats = 0
|
||||
var flushed = false
|
||||
val result = AgentSyncCycle({ heartbeats++; throw api(409, code) }, { null }, { false }, {},
|
||||
{ flushed = true }, { false }).run()
|
||||
assertTrue(flushed)
|
||||
assertEquals(if (code == "DEVICE_TASK_MISMATCH") 2 else 1, heartbeats)
|
||||
assertFalse(result.canClaim)
|
||||
}
|
||||
}
|
||||
|
||||
@Test fun `network failure leaves pending even if heartbeat succeeds`() {
|
||||
var pending = true
|
||||
val result = AgentSyncCycle({ online }, { 42 }, { false }, {}, { throw java.io.IOException("network") }, { pending }).run()
|
||||
assertTrue(pending)
|
||||
assertEquals(online, result.heartbeat)
|
||||
assertFalse(result.canClaim)
|
||||
}
|
||||
|
||||
@Test fun `active execution sends heartbeat without recovery or flush`() {
|
||||
val result = AgentSyncCycle({ online }, { 42 }, { true }, { fail("recover") }, { fail("flush") }, { false }).run()
|
||||
assertEquals(online, result.heartbeat)
|
||||
assertFalse(result.canClaim)
|
||||
}
|
||||
|
||||
@Test fun `retained irreversible local task prevents claim even without outbox`() {
|
||||
val result = AgentSyncCycle({ online }, { 42 }, { false }, {}, {}, { false }).run()
|
||||
assertFalse(result.canClaim)
|
||||
}
|
||||
|
||||
@Test fun `authentication failure during recovery or flush stops round without heartbeat retry`() {
|
||||
for (duringRecovery in listOf(true, false)) {
|
||||
var heartbeats = 0
|
||||
var flushes = 0
|
||||
val result = AgentSyncCycle({ heartbeats++; throw api(409,"DEVICE_TASK_MISMATCH") }, { 42 }, { false },
|
||||
{ if (duringRecovery) throw api(401,"OTHER") },
|
||||
{ flushes++; throw api(403,"OTHER") }, { true }).run()
|
||||
assertEquals(1, heartbeats)
|
||||
assertEquals(if (duringRecovery) 0 else 1, flushes)
|
||||
assertEquals("AUTH_ERROR", result.failureCode)
|
||||
assertFalse(result.canClaim)
|
||||
}
|
||||
}
|
||||
|
||||
@Test fun `server failure is classified separately from transport failure`() {
|
||||
val result = AgentSyncCycle({ throw api(503,"UNAVAILABLE") }, { null }, { false }, {}, {}, { false }).run()
|
||||
assertEquals("SERVER_ERROR", result.failureCode)
|
||||
}
|
||||
|
||||
@Test fun `claim failure cannot erase successful heartbeat except authentication`() {
|
||||
assertNull(syncConnectionFailureCode(api(409,"PURCHASE_STATE_CONFLICT"), true))
|
||||
assertNull(syncConnectionFailureCode(IllegalStateException("disk"), true))
|
||||
assertNull(syncConnectionFailureCode(api(503,"INTERNAL"), true))
|
||||
assertEquals("AUTH_ERROR", syncConnectionFailureCode(api(401,"OTHER"), true))
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user