Compare commits

..
55 changed files with 4750 additions and 538 deletions
@@ -0,0 +1,126 @@
package cn.ilapage.goauto.agent.automation
/** Current-capture evidence only. No clicks, historical selection proof, or whole-screen text fallback. */
internal object PurchaseFinalSpecVerifier {
private fun normalize(value: String) = SpecValueNormalizer.normalizeFinalConfirmation(value)
fun verify(snapshot: UiSnapshot, input: PurchaseExecutionInput, diagnostic: (String) -> Unit) {
val config = PurchaseRehearsalExecutor.DEFAULT_COLLECTOR
val byPath = snapshot.nodes.associateBy(SnapshotNode::path)
fun within(node: SnapshotNode, container: SnapshotNode): Boolean {
val visited = mutableSetOf<String>()
var current: SnapshotNode? = node
while (current != null && visited.add(current.path)) {
if (current.path == container.path) return true
current = current.parentPath?.let(byPath::get)
}
return false
}
fun contained(node: SnapshotNode, container: SnapshotNode) =
node.bounds.width > 0 && node.bounds.height > 0 &&
node.bounds.left >= container.bounds.left && node.bounds.right <= container.bounds.right &&
node.bounds.top >= container.bounds.top && node.bounds.bottom <= container.bounds.bottom
val original = PddScreenParser.parse(snapshot, config, input.goodsId, null, forFinalConfirmation = true)
val fullPanel = boundedQuantityPanel(snapshot).takeIf { original.specPanelOpen }
// #331/#335 sheet ancestry, as represented by SpecPanelFixtures.sheet and taskOptionDedupSheet.
// The inner scroll list is not the full sheet: summary and quantity may be its siblings.
val scoped = fullPanel?.let { panel ->
val nodes = snapshot.nodes.filter { node ->
(within(node, panel) && contained(node, panel)) ||
(node.parentPath == null && node.label.isBlank() && !node.clickable)
}
PddScreenParser.parse(snapshot.copy(nodes = nodes), config, input.goodsId, null, forFinalConfirmation = true)
.takeIf { it.specPanelOpen }
}
val optionScope = if (scoped != null) fullPanel else original.specPanelContainer
val dimensions = (scoped ?: original).dimensions.groupBy(VisibleDimension::key).mapValues { (_, groups) ->
val options = groups.flatMap(VisibleDimension::values).filter { option ->
optionScope != null && within(option.node, optionScope) && contained(option.node, optionScope)
}
// A blank clickable wrapper and its single text child may survive optionBlock separately.
// The #331 size fixture exposes exactly the same rectangle for both. Aggregate only a
// same-value ancestor chain with that physical rectangle, never same-label siblings.
options.groupBy { option ->
options.filter { other -> other.node.bounds == option.node.bounds &&
normalize(other.text) == normalize(option.text) && within(option.node, other.node)
}.minBy { it.node.path.length }.node.path
}.values.map { members ->
members.first().let { first -> first.copy(
available = members.all(VisibleSpecValue::available),
node = first.node.copy(
selected = members.any { it.node.selected },
checked = members.any { it.node.checked },
),
) }
}
}
val summaries = scoped?.finalSelectionSummaries.orEmpty().map { raw ->
// Only the anchored recognized prefix and its optional colon are removed.
val trimmed = raw.trim()
val prefix = config.textAliases.selection.selectedPrefixes.firstOrNull(trimmed::startsWith)
val value = if (prefix == null) trimmed else trimmed.removePrefix(prefix).trimStart().let {
if (it.startsWith(':') || it.startsWith(':')) it.drop(1) else it
}
normalize(value)
}.distinct()
val targets = linkedMapOf("color" to normalize(input.mappedColor), "size" to normalize(input.mappedSize))
.filterValues(String::isNotBlank)
fun reject(dimension: String, stage: String, selected: Boolean = false, conflict: Boolean = false): Nothing {
val detail = "finalSpec;dimension=$dimension;stage=$stage;summaryPresent=${if (summaries.isNotEmpty()) 1 else 0};" +
"selectedFound=${if (selected) 1 else 0};conflict=${if (conflict) 1 else 0}"
diagnostic(detail)
throw PurchaseLiveException("PURCHASE_SPEC_NOT_MATCHED", "创建订单前规格复核失败;$detail")
}
if (targets.isEmpty()) reject("none", "missing_target")
val selectedProof = mutableSetOf<String>()
targets.forEach { (dimension, target) ->
val options = dimensions[dimension].orEmpty()
val selected = options.filter { it.node.selected || it.node.checked }
// Parsed text already follows the established UI trailing-price separation contract.
// The final-only normalizer removes no additional content from values, input or summary.
val exact = options.filter { normalize(it.text) == target }
if (selected.any { normalize(it.text) != target }) reject(dimension, "selected_conflict", selected = true, conflict = true)
if (exact.size > 1) reject(dimension, "candidate_collision", selected.isNotEmpty(), conflict = true)
if (exact.any { !it.available }) reject(dimension, "unavailable", selected.isNotEmpty(), conflict = true)
if (exact.any { it.node.selected || it.node.checked }) selectedProof += dimension
}
if (summaries.size > 1) reject(targets.keys.first(), "summary_conflict", selectedProof.isNotEmpty(), conflict = true)
val summary = summaries.singleOrNull()
// Evidence only supports color then size with a space. Do not guess reverse order or comma separators.
val expected = targets.values.joinToString(" ")
if (summary != null && targets.size == 2) {
val colors = (dimensions["color"].orEmpty().map { normalize(it.text) } + targets.getValue("color")).distinct()
val sizes = (dimensions["size"].orEmpty().map { normalize(it.text) } + targets.getValue("size")).distinct()
if (colors.any { color -> sizes.any { size ->
"$color $size" == summary && (color != targets["color"] || size != targets["size"])
} }) reject("color", "composition_ambiguity", selectedProof.isNotEmpty(), conflict = true)
}
targets.keys.forEach { dimension ->
if (dimension !in selectedProof && summary != expected) reject(dimension, "missing_proof")
}
}
/** Existing #335 quantity-ancestor form, retaining identity as well as bounds for summary ownership. */
private fun boundedQuantityPanel(snapshot: UiSnapshot): SnapshotNode? {
val visible = snapshot.nodes.filter(SnapshotNode::visible)
val screenWidth = visible.maxOfOrNull { it.bounds.right } ?: return null
val screenHeight = visible.maxOfOrNull { it.bounds.bottom } ?: return null
val screenArea = screenWidth.toLong() * screenHeight
val byPath = visible.associateBy(SnapshotNode::path)
var current = visible.singleOrNull {
it.enabled && it.className == "android.widget.EditText" && it.label.toIntOrNull()?.let { value -> value > 0 } == true
} ?: return null
var panel: SnapshotNode? = null
val visited = mutableSetOf(current.path)
while (true) {
val parent = current.parentPath?.let(byPath::get) ?: break
if (!visited.add(parent.path)) return null
val area = parent.bounds.width.toLong() * parent.bounds.height
if (area <= 0 || area >= screenArea || parent.parentPath == null) break
panel = parent
current = parent
}
return panel
}
}
@@ -3,7 +3,6 @@ package cn.ilapage.goauto.agent.automation
import java.text.SimpleDateFormat
import java.util.Locale
import java.util.TimeZone
import org.json.JSONObject
data class ShippingAddressProof internal constructor(
internal val expectedAddress: String,
@@ -18,28 +17,7 @@ data class FinalConfirmationEvidence(
val unitPriceCent: Long,
val addressSuffix: String,
val activityName: String,
val observedQuantity: Long? = null,
val selectedSummaryPresent: Boolean = false,
val summaryPromptPresent: Boolean = false,
) {
/** Local boundary evidence only; never used in the uploaded purchase result. */
fun toLocalJson(): String = JSONObject()
.put("goodsId", goodsId)
.put("mappedColor", mappedColor)
.put("mappedSize", mappedSize)
.put("targetQuantity", quantity)
.put("observedQuantity", observedQuantity ?: JSONObject.NULL)
.put("selectedSummaryPresent", selectedSummaryPresent)
.put("summaryPromptPresent", summaryPromptPresent)
.put("unitPriceCent", unitPriceCent)
.put("addressSuffix", addressSuffix)
.put("activityName", activityName)
.toString()
}
internal fun observedPurchaseQuantity(snapshot: UiSnapshot): Long? = snapshot.nodes.filter {
it.visible && it.enabled && it.className?.endsWith("EditText") == true && it.label.toLongOrNull() != null
}.singleOrNull()?.label?.toLongOrNull()
)
data class PurchaseOrderEvidence(val orderNo: String, val submittedAt: String, val pddOrderAmountCent: Long? = null)
data class PurchaseOrderReadFailure(
@@ -337,20 +315,18 @@ class PurchaseLiveAutomation(
val snapshot = driver.capture()
pageProblem(snapshot)
if (snapshot.packageName != PDD_PACKAGE) fail("PURCHASE_CONFIRMATION_LOST", "最终提交前页面已经变化,禁止创建订单")
if (!hasSavedAddressEvidence(snapshot, address.expectedAddress, address.suffix)) {
if (!hasFinalSavedAddressEvidence(snapshot, address.expectedAddress, address.suffix)) {
fail("PURCHASE_ADDRESS_UPDATE_FAILED", "收货地址保存后复核失败,未创建订单")
}
val submit = finalSubmitTargets(snapshot)
if (submit.size != 1) fail("PURCHASE_SUBMIT_TARGET_AMBIGUOUS", "创建订单按钮不是唯一目标,禁止创建订单")
val screen = PddScreenParser.parse(snapshot, PurchaseRehearsalExecutor.DEFAULT_COLLECTOR, input.goodsId, null)
rejectSubmitSelectionPrompt(submit.single())
PurchaseFinalSpecVerifier.verify(snapshot, input, panelDiagnostic)
val price = screen.priceCent ?: fail("RULE_NOT_MATCHED", "创建订单前没有读取到商品单价")
if (price !in input.minUnitPriceCent..input.maxUnitPriceCent) fail("PURCHASE_PRICE_OUT_OF_RANGE", "当前商品单价超出允许范围", price)
val selection = PurchaseSelectionEvidence.observe(snapshot)
return FinalConfirmationEvidence(
input.goodsId, input.mappedColor, input.mappedSize, input.quantity, price, address.suffix, snapshot.activityName.orEmpty(),
observedPurchaseQuantity(snapshot), selection.selectedSummaryPresent, selection.summaryPromptPresent,
)
val quantities = snapshot.nodes.filter { it.visible && it.enabled && it.className?.endsWith("EditText") == true }.mapNotNull { it.label.toLongOrNull() }
if (quantities.singleOrNull() != input.quantity) fail("PURCHASE_QUANTITY_MISMATCH", "创建订单前数量复核失败")
return FinalConfirmationEvidence(input.goodsId, input.mappedColor, input.mappedSize, input.quantity, price, address.suffix, snapshot.activityName.orEmpty())
}
fun submitOrderOnce() {
@@ -361,9 +337,8 @@ class PurchaseLiveAutomation(
if (snapshot.packageName != PDD_PACKAGE || targets.size != 1) {
fail("PURCHASE_SUBMIT_TARGET_AMBIGUOUS", "最终页面或创建订单按钮已变化,禁止创建订单")
}
rejectSubmitSelectionPrompt(targets.single())
submitAttempted = true
when (driver.clickFresh(targets.single().chosen)) {
when (driver.clickFresh(targets.single())) {
FreshActionResult.SUCCESS -> Unit
FreshActionResult.AMBIGUOUS -> fail("PURCHASE_SUBMIT_TARGET_AMBIGUOUS", "创建订单按钮不唯一")
else -> fail("PURCHASE_ORDER_RESULT_UNKNOWN", "创建订单点击结果不明确,请人工检查")
@@ -904,24 +879,7 @@ class PurchaseLiveAutomation(
* (Samsung: only a zero-size text node), fail explicitly (`bottom_row_unlabelled`) —
* never fall back to a higher row.
*/
private data class SubmitTarget(val row: SnapshotNode, val chosen: SnapshotNode, val text: String)
internal fun selectionPromptPresent(snapshot: UiSnapshot): Boolean =
PurchaseSelectionEvidence.observe(snapshot).summaryPromptPresent ||
finalSubmitTargets(snapshot).singleOrNull()?.let { isSelectionPrompt(it.text) } == true
private fun isSelectionPrompt(text: String): Boolean =
Regex("^(选择|選擇).+(后|後)[,,](提交订单|提交訂單)$").matches(text.filterNot(Char::isWhitespace))
private fun rejectSubmitSelectionPrompt(target: SubmitTarget) {
if (isSelectionPrompt(target.text)) {
val detail = "finalSpec;stage=selection_prompt"
panelDiagnostic(detail)
fail("PURCHASE_SPEC_NOT_MATCHED", "创建订单前规格复核失败;$detail")
}
}
private fun finalSubmitTargets(snapshot: UiSnapshot): List<SubmitTarget> {
private fun finalSubmitTargets(snapshot: UiSnapshot): List<SnapshotNode> {
fun hasArea(node: SnapshotNode) = node.bounds.width > 0 && node.bounds.height > 0
val byPath = snapshot.nodes.associateBy { it.path }
fun clickableAncestor(node: SnapshotNode): SnapshotNode? {
@@ -990,14 +948,7 @@ class PurchaseLiveAutomation(
"leafClass=${chosen.className};rowClass=${row.className};w=${row.bounds.width};h=${row.bounds.height};" +
"price=${hasPrice.diagFlag()};legacyMarker=${legacyMarker.diagFlag()};type=${screen.specPanelType}",
)
// Prompt evidence uses each visible leaf's own label exactly once. A nested
// clickable button is a separate action and cannot supply text for this row.
val text = (listOf(row) + subtreeDescendants(row, snapshot)).filter { node ->
node.visible && node.label.isNotBlank() && clickableAncestor(node)?.path == row.path &&
subtreeDescendants(node, snapshot).none { it.label.isNotBlank() }
}.sortedWith(compareBy<SnapshotNode> { it.bounds.top }.thenBy { it.bounds.left })
.joinToString("") { it.label }
return listOf(SubmitTarget(row, chosen, text))
return listOf(chosen)
}
private fun subtreeDescendants(node: SnapshotNode, snapshot: UiSnapshot): List<SnapshotNode> {
@@ -772,11 +772,6 @@ class PurchaseRehearsalExecutor(
val summaryTokenMatched = !token.isNullOrBlank() && SpecValueNormalizer.summaryHasExactToken(screen.selectedSummary, token)
if (!summaryTokenMatched) {
if (screen.selectedSummary == null && proofUsable) {
if (candidates.any { it.available && it.node.visible &&
SpecValueNormalizer.normalizeFinalConfirmation(it.text) == SpecValueNormalizer.normalizeFinalConfirmation(target) &&
!it.node.selected && !it.node.checked }) {
return FinalSpecVerification(false, "target_visible_unselected", false, candidates.size, true, true, matchingNodes.size)
}
return FinalSpecVerification(true, "attempt_selection_state", false, candidates.size, true, true, matchingNodes.size)
}
return FinalSpecVerification(false, "summary_token_missing", false, candidates.size, proofRecorded, proofUsable, matchingNodes.size)
@@ -808,15 +803,12 @@ class PurchaseRehearsalExecutor(
target: String,
proof: ExactSpecSelectionProof?,
): FinalSpecVerification {
var selectionPromptPresent = false
fun captureScreen() = currentScreen(input) { snapshot -> selectionPromptPresent = hasSelectionPrompt(snapshot) }
var screen = captureScreen()
var screen = currentScreen(input)
screen.problem?.let {
return FinalSpecVerification(false, it.code, false, 0, proof != null, false)
}
if (selectionPromptPresent) return FinalSpecVerification(false, "selection_prompt", false, 0, proof != null, false)
var verification = verifyExactSpecSelection(screen, dimension, target, proof)
if (verification.confirmed || verification.targetMatchCount > 0 || verification.reason in setOf("visible_selected_conflict", "selection_prompt", "target_visible_unselected")) return verification
if (verification.confirmed || verification.targetMatchCount > 0 || verification.reason == "visible_selected_conflict") return verification
if (!screen.specPanelOpen) return verification.copy(reason = "relocation_panel_closed", relocationAttempted = true)
var signature = screen.dimensions.joinToString("|") { item ->
@@ -840,7 +832,7 @@ class PurchaseRehearsalExecutor(
}
swipes++
pause(300)
screen = captureScreen()
screen = currentScreen(input)
screen.problem?.let {
return verification.copy(
reason = it.code,
@@ -848,9 +840,6 @@ class PurchaseRehearsalExecutor(
relocationSwipes = swipes,
)
}
if (selectionPromptPresent) {
return verification.copy(confirmed = false, reason = "selection_prompt", relocationAttempted = true, relocationSwipes = swipes)
}
if (!screen.specPanelOpen) {
return verification.copy(
reason = "relocation_panel_closed",
@@ -862,7 +851,7 @@ class PurchaseRehearsalExecutor(
relocationAttempted = true,
relocationSwipes = swipes,
)
if (verification.confirmed || verification.targetMatchCount > 0 || verification.reason in setOf("visible_selected_conflict", "selection_prompt", "target_visible_unselected")) return verification
if (verification.confirmed || verification.targetMatchCount > 0 || verification.reason == "visible_selected_conflict") return verification
val refreshedSignature = screen.dimensions.joinToString("|") { item ->
"${item.key}:${item.values.joinToString(",") { value -> "${value.text}:${value.available}" }}"
}
@@ -1140,9 +1129,6 @@ class PurchaseRehearsalExecutor(
dimension to (normalizedTarget(dimension, rawTarget)
?: return failure(SPEC_SAFE_TARGET_MISSING, "下发规格无法安全规范化"))
}
if (selected.isEmpty() && hasSelectionPrompt(driver.capture())) {
return failure(SPEC_SELECTION_UNCONFIRMED, "最终规格复核未能确认精确选中状态 [reason=selection_prompt]")
}
selected.forEach { (dimension, target) ->
val verification = verifyExactSpecSelectionWithRelocation(input, dimension, target, selectionProofs[dimension])
if (!verification.confirmed) {
@@ -1166,9 +1152,6 @@ class PurchaseRehearsalExecutor(
return null
}
private fun hasSelectionPrompt(snapshot: UiSnapshot): Boolean =
PurchaseLiveAutomation(driver, pause, panelDiagnostic).selectionPromptPresent(snapshot)
private fun hasAllRecordedSelectionProofs(
input: PurchaseExecutionInput,
proofs: Map<String, ExactSpecSelectionProof>,
@@ -1236,9 +1219,8 @@ class PurchaseRehearsalExecutor(
return normalized.takeIf(SpecValueNormalizer::isSafeSize)
}
private fun currentScreen(input: PurchaseExecutionInput, snapshotObserved: (UiSnapshot) -> Unit = {}): ParsedPddScreen {
private fun currentScreen(input: PurchaseExecutionInput): ParsedPddScreen {
val snapshot = driver.capture()
snapshotObserved(snapshot)
val screen = PddScreenParser.parse(snapshot, DEFAULT_COLLECTOR, input.goodsId, null, purchasePanelContext)
purchasePanelContext = if (screen.isPddPackage && screen.problem == null && screen.specPanelOpen) {
screen.specPanelContainer?.let {
@@ -1,68 +0,0 @@
package cn.ilapage.goauto.agent.automation
/** Boolean-only summary observations within the #374 bounded quantity panel. */
internal data class PurchaseSelectionEvidence(
val selectedSummaryPresent: Boolean = false,
val summaryPromptPresent: Boolean = false,
) {
companion object {
fun observe(snapshot: UiSnapshot): PurchaseSelectionEvidence {
val config = PurchaseRehearsalExecutor.DEFAULT_COLLECTOR
val original = PddScreenParser.parse(snapshot, config, "", null)
val panel = boundedQuantityPanel(snapshot).takeIf { original.specPanelOpen }
?: return PurchaseSelectionEvidence()
val byPath = snapshot.nodes.associateBy(SnapshotNode::path)
fun within(node: SnapshotNode): Boolean {
val visited = mutableSetOf<String>()
var current: SnapshotNode? = node
while (current != null && visited.add(current.path)) {
if (current.path == panel.path) return true
current = current.parentPath?.let(byPath::get)
}
return false
}
fun contained(node: SnapshotNode) = node.bounds.width > 0 && node.bounds.height > 0 &&
node.bounds.left >= panel.bounds.left && node.bounds.right <= panel.bounds.right &&
node.bounds.top >= panel.bounds.top && node.bounds.bottom <= panel.bounds.bottom
val nodes = snapshot.nodes.filter { (within(it) && contained(it)) ||
(it.parentPath == null && it.label.isBlank() && !it.clickable) }
val scoped = PddScreenParser.parse(snapshot.copy(nodes = nodes), config, "", null, forFinalConfirmation = true)
if (!scoped.specPanelOpen) return PurchaseSelectionEvidence()
val quantity = nodes.singleOrNull { it.visible && it.enabled && it.className == "android.widget.EditText" && it.label.toIntOrNull()?.let { value -> value > 0 } == true }
val firstOptionTop = scoped.dimensions.flatMap(VisibleDimension::values).minOfOrNull { it.node.bounds.top } ?: Int.MAX_VALUE
val prompt = quantity != null && nodes.any { node ->
node.visible && !node.clickable && node.className?.endsWith("TextView") == true &&
node.bounds.bottom <= minOf(quantity.bounds.top, firstOptionTop) &&
node.label.filterNot(Char::isWhitespace).let { it.startsWith("请选择") || it.startsWith("請選擇") }
}
val selected = scoped.finalSelectionSummaries.any { raw ->
val compact = raw.filterNot(Char::isWhitespace)
config.textAliases.selection.selectedPrefixes.any(compact::startsWith)
}
return PurchaseSelectionEvidence(selected, prompt)
}
/** Retains #374's real ancestry and bounded full sheet, not the inner option list. */
private fun boundedQuantityPanel(snapshot: UiSnapshot): SnapshotNode? {
val visible = snapshot.nodes.filter(SnapshotNode::visible)
val screenWidth = visible.maxOfOrNull { it.bounds.right } ?: return null
val screenHeight = visible.maxOfOrNull { it.bounds.bottom } ?: return null
val screenArea = screenWidth.toLong() * screenHeight
val byPath = visible.associateBy(SnapshotNode::path)
var current = visible.singleOrNull {
it.enabled && it.className == "android.widget.EditText" && it.label.toIntOrNull()?.let { value -> value > 0 } == true
} ?: return null
var panel: SnapshotNode? = null
val visited = mutableSetOf(current.path)
while (true) {
val parent = current.parentPath?.let(byPath::get) ?: break
if (!visited.add(parent.path)) return null
val area = parent.bounds.width.toLong() * parent.bounds.height
if (area <= 0 || area >= screenArea || parent.parentPath == null) break
panel = parent
current = parent
}
return panel
}
}
}
@@ -35,6 +35,7 @@ data class HeartbeatResult(
val online: Boolean,
val busy: Boolean,
val heartbeatIntervalSeconds: Int,
val acceptsClientReport: Boolean = false,
)
data class AgentAppRelease(
@@ -221,7 +222,10 @@ class AgentApiException(
val retryable: Boolean,
) : Exception(message)
class AgentApiClient(private val serverUrl: String) {
class AgentApiClient(
internal val serverUrl: String,
private val connectionFactory: (URL) -> HttpURLConnection = { it.openConnection() as HttpURLConnection },
) {
@Volatile var lastResponseServerTimeMillis: Long? = null
private set
val failureSnapshotOrigin: String get() = ServerUrlPolicy.normalize(serverUrl)
@@ -309,17 +313,19 @@ class AgentApiClient(private val serverUrl: String) {
)
}
fun heartbeat(token: String, currentTaskId: Long?, capabilities: List<String> = emptyList()): HeartbeatResult {
fun heartbeat(token: String, currentTaskId: Long?, capabilities: List<String> = emptyList(), clientReport: JSONObject? = null, requestId: String = UUID.randomUUID().toString()): HeartbeatResult {
val payload = JSONObject()
.put("requestId", UUID.randomUUID().toString())
.put("requestId", requestId)
.put("currentTaskId", currentTaskId ?: JSONObject.NULL)
.put("capabilities", JSONArray(capabilities))
if (clientReport != null) payload.put("clientReport", clientReport)
val data = post("/api/agent/v1/heartbeat", payload, token).getJSONObject("data")
return HeartbeatResult(
deviceId = data.getLong("deviceId"),
online = data.getBoolean("online"),
busy = data.getBoolean("busy"),
heartbeatIntervalSeconds = data.optInt("heartbeatIntervalSeconds", 15),
acceptsClientReport = data.opt("acceptsClientReport") == true,
)
}
@@ -707,7 +713,7 @@ class AgentApiClient(private val serverUrl: String) {
}
private fun request(method: String, path: String, payload: JSONObject?, token: String?, recoveryCode: String? = null): JSONObject? {
val connection = (URL(serverUrl + path).openConnection() as HttpURLConnection).apply {
val connection = connectionFactory(URL(serverUrl + path)).apply {
requestMethod = method
connectTimeout = 10_000
readTimeout = 15_000
@@ -719,10 +725,14 @@ class AgentApiClient(private val serverUrl: String) {
if (!token.isNullOrBlank()) setRequestProperty("Authorization", "Bearer $token")
if (!recoveryCode.isNullOrBlank()) setRequestProperty("X-GoAuto-Device-Recovery-Code", recoveryCode)
}
var stage = RequestStage.CONNECT
try {
connection.connect()
stage = RequestStage.UNKNOWN
if (payload != null) {
connection.outputStream.use { it.write(payload.toString().toByteArray(Charsets.UTF_8)) }
}
stage = RequestStage.RESPONSE
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.
@@ -740,6 +750,8 @@ class AgentApiClient(private val serverUrl: String) {
)
}
return json
} catch (error: java.net.SocketTimeoutException) {
throw AgentTransportTimeout(stage, error)
} finally {
connection.disconnect()
}
@@ -0,0 +1,65 @@
package cn.ilapage.goauto.agent.network
import java.util.UUID
import org.json.JSONObject
/** Service lifetime, never persisted: a new process must negotiate again. */
internal class HeartbeatReportNegotiation {
var acceptsReport = false
private set
private var origin: String? = null
private var token: String? = null
fun heartbeat(
api: AgentApiClient,
token: String,
taskId: Long?,
capabilities: List<String>,
report: () -> JSONObject?,
ack: () -> Unit,
validate: (HeartbeatResult) -> Unit = {},
): HeartbeatResult {
if (origin != api.serverUrl || this.token != token) {
acceptsReport = false
origin = api.serverUrl
this.token = token
}
val snapshot = if (acceptsReport) report() else null
val requestId = UUID.randomUUID().toString()
val result =
try {
api.heartbeat(token, taskId, capabilities, snapshot, requestId)
} catch (error: AgentApiException) {
if (
snapshot == null ||
error.status != 422 ||
error.code != "INVALID_REQUEST" ||
error.message != "请求 JSON 无效"
)
throw error
acceptsReport = false
// A fallback is never an acknowledgement or a capability grant.
return api.heartbeat(token, taskId, capabilities, null, requestId).also(validate)
}
validate(result)
acceptsReport = result.acceptsClientReport
if (snapshot != null) ack()
return result
}
}
internal enum class RequestStage {
CONNECT,
RESPONSE,
UNKNOWN,
}
internal class AgentTransportTimeout(
val stage: RequestStage,
cause: java.net.SocketTimeoutException,
) : java.net.SocketTimeoutException("agent_transport_timeout") {
init {
initCause(cause)
}
}
@@ -1,7 +1,22 @@
package cn.ilapage.goauto.agent.persistence
internal object AgentDiagnosticSchema {
const val VERSION = 5
const val VERSION = 6
val syncStatements = listOf(
"""CREATE TABLE IF NOT EXISTS sync_round (
id INTEGER PRIMARY KEY AUTOINCREMENT, process_id TEXT NOT NULL,
started_at INTEGER NOT NULL, ended_at INTEGER, duration_ms INTEGER, start_gap_ms INTEGER,
device_id INTEGER, task_id INTEGER, attempt_id TEXT, phase TEXT,
validated_network INTEGER, transport TEXT NOT NULL, foreground INTEGER NOT NULL,
status TEXT NOT NULL, steps TEXT NOT NULL DEFAULT '[]',
failure_at INTEGER, failure_step TEXT, failure_category TEXT, failure_duration_ms INTEGER,
http_status INTEGER, api_code TEXT, exception_class TEXT,
completion_seq INTEGER, reported INTEGER NOT NULL DEFAULT 0
)""".trimIndent(),
"CREATE TABLE IF NOT EXISTS sync_meta (id INTEGER PRIMARY KEY CHECK(id=1), sequence INTEGER NOT NULL)",
"INSERT OR IGNORE INTO sync_meta(id,sequence) VALUES(1,0)",
)
val failureSnapshotStatements = listOf(
"""CREATE TABLE IF NOT EXISTS purchase_failure_snapshot (
@@ -99,5 +114,6 @@ internal object AgentDiagnosticSchema {
(if (oldVersion < 4 && newVersion >= 4) failureSnapshotStatements else emptyList()) +
(if (oldVersion < 5 && newVersion >= 5) colorDiscoveryOutcomeColumns.mapNotNull { (name, definition) ->
if (name in existingColumns) null else "ALTER TABLE agent_diagnostic ADD COLUMN $name $definition"
} else emptyList())
} else emptyList()) +
(if (oldVersion < 6 && newVersion >= 6) syncStatements else emptyList())
}
@@ -150,6 +150,7 @@ class AgentDiagnosticStore(context: Context) : SQLiteOpenHelper(context, DATABAS
db.execSQL(AgentDiagnosticSchema.createTableSql)
db.execSQL("CREATE INDEX idx_agent_diagnostic_task ON agent_diagnostic(task_id, id)")
AgentDiagnosticSchema.failureSnapshotStatements.forEach(db::execSQL)
AgentDiagnosticSchema.syncStatements.forEach(db::execSQL)
}
override fun onUpgrade(db: SQLiteDatabase, oldVersion: Int, newVersion: Int) {
@@ -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,53 @@ 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 isCurrentAttemptRejected(taskId: Long, attemptId: String): Boolean = readableDatabase.rawQuery(
"SELECT 1 FROM purchase_task WHERE task_id=? AND attempt_id=? AND upload_status='rejected'",
arrayOf(taskId.toString(), attemptId),
).use { it.moveToFirst() }
@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"
@@ -0,0 +1,234 @@
package cn.ilapage.goauto.agent.persistence
import cn.ilapage.goauto.agent.network.AgentApiException
import cn.ilapage.goauto.agent.network.AgentTransportTimeout
import cn.ilapage.goauto.agent.network.RequestStage
import java.util.concurrent.CancellationException
import org.json.JSONArray
import org.json.JSONObject
/** Collection and purchase IDs belong to separate domains and can have the same numeric value. */
internal fun syncTaskIdentity(
taskId: Long?,
taskType: String?,
executionActive: Boolean,
live: SyncContext?,
recoveredAttempt: () -> String?,
): SyncContext {
if (taskType == "purchase" && live?.taskId == taskId && live != null) return live
return SyncContext(
taskId = taskId,
attemptId = if (!executionActive && taskId != null) recoveredAttempt() else null,
)
}
/** Only fixed classifications enter this store; exception text, URLs and payloads never do. */
internal fun syncFailure(error: Exception, step: String, at: Long, duration: Long): SyncFailure {
val api = error as? AgentApiException
val timeout = error as? AgentTransportTimeout
val category =
when {
api != null ->
when {
cn.ilapage.goauto.agent.service.isAgentAuthenticationError(api) -> "api_auth"
api.status == 409 -> "api_conflict"
api.status >= 500 -> "api_server"
else -> "api_other"
}
timeout != null ->
when (timeout.stage) {
RequestStage.CONNECT -> "timeout_connect"
RequestStage.RESPONSE -> "timeout_response"
RequestStage.UNKNOWN -> "timeout_unknown"
}
error is java.net.SocketTimeoutException -> "timeout_unknown"
error is java.io.IOException -> "network"
else -> "exception"
}
val code =
api?.code?.takeIf {
it in
setOf(
"DEVICE_TOKEN_INVALID",
"DEVICE_DISABLED",
"DEVICE_INSTALL_ID_CONFLICT",
"DEVICE_TASK_MISMATCH",
"INVALID_REQUEST",
"PURCHASE_STATE_CONFLICT",
"PURCHASE_LEASE_EXPIRED",
"PURCHASE_RESULT_CONFLICT",
)
}
val type =
when (error) {
is java.net.SocketTimeoutException -> "SocketTimeoutException"
is java.net.UnknownHostException -> "UnknownHostException"
is java.net.ConnectException -> "ConnectException"
is javax.net.ssl.SSLException -> "SSLException"
is java.net.SocketException -> "SocketException"
is java.io.EOFException -> "EOFException"
else -> null
}
return SyncFailure(
step,
category,
at.coerceAtLeast(1),
duration.coerceAtLeast(0),
api?.status,
code,
type,
)
}
internal class SyncDiagnosticRecorder(
private val store: SyncDiagnosticStore,
private val process: String,
private val wall: () -> Long,
private val elapsed: () -> Long,
private val warning: () -> Unit = {},
) {
private var previousStart: Long? = null
fun begin(context: SyncContext): Round {
val start = elapsed()
val at = wall()
val gap = previousStart?.let { (start - it).coerceAtLeast(0) }
previousStart = start
val id = safe { store.start(process, at, gap, context) }
val round = Round(id, at, start)
round.persistSteps()
safe { store.cleanup(at) }
return round
}
fun snapshot(): SyncReportSnapshot? = safe { store.snapshot(process, wall()) }
fun ack(snapshot: SyncReportSnapshot) {
safe { store.ack(snapshot) }
}
private fun <T> safe(action: () -> T): T? =
try {
action()
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
try {
warning()
} catch (cancel: CancellationException) {
throw cancel
} catch (ignored: Exception) {}
null
}
inner class Round(private val id: Long?, at: Long, private val start: Long) {
private val entries =
JSONArray()
.put(
JSONObject()
.put("ordinal", 0)
.put("step", "sync")
.put("startedAt", at)
.put("startedOffsetMs", 0)
.put("status", "running")
)
private var failure: SyncFailure? = null
private var persistenceFailed = id == null
private var aborted = false
private var completed = false
private var activeStep: Pair<JSONObject, Long>? = null
fun <T> step(name: String, action: () -> T): T {
require(name in setOf("register", "heartbeat", "recover", "flush", "claim"))
val begin = elapsed()
val entry =
JSONObject()
.put("ordinal", entries.length())
.put("step", name)
.put("startedAt", wall())
.put("startedOffsetMs", (begin - start).coerceAtLeast(0))
.put("status", "running")
entries.put(entry)
activeStep = entry to begin
persistSteps()
try {
val value = action()
if (entry.getString("status") != "failure") entry.put("status", "success")
return value
} catch (e: CancellationException) {
aborted = true
throw e
} catch (e: Exception) {
failure = syncFailure(e, name, wall(), elapsed() - begin)
entry.put("status", "failure").put("category", failure!!.category)
throw e
} catch (e: Error) {
aborted = true
throw e
} finally {
if (!aborted) {
entry
.put("endedAt", wall())
.put("endedOffsetMs", (elapsed() - start).coerceAtLeast(0))
.put("durationMs", (elapsed() - begin).coerceAtLeast(0))
persistSteps()
}
activeStep = null
}
}
/**
* Observe a business-handled exception without changing that operation's return/throw
* policy.
*/
fun handledFailure(error: Exception) {
if (error is CancellationException) {
aborted = true
throw error
}
val (entry, begin) = activeStep ?: return failed(error)
failure = syncFailure(error, entry.getString("step"), wall(), elapsed() - begin)
entry.put("status", "failure").put("category", failure!!.category)
persistSteps()
}
fun failed(error: Exception) {
if (error is CancellationException) {
aborted = true
throw error
}
if (failure == null) failure = syncFailure(error, "sync", wall(), elapsed() - start)
}
fun abort() {
aborted = true
}
fun finish() {
if (completed || aborted) return
completed = true
entries
.getJSONObject(0)
.put("endedAt", wall())
.put("endedOffsetMs", (elapsed() - start).coerceAtLeast(0))
.put("durationMs", (elapsed() - start).coerceAtLeast(0))
.put("status", if (failure == null) "success" else "failure")
persistSteps()
// A failed start/step write must never turn an incomplete durable trace into a fake
// success.
if (!persistenceFailed && id != null)
safe { store.finish(id, wall(), (elapsed() - start).coerceAtLeast(0), failure) }
}
internal fun persistSteps() {
if (
id != null &&
safe {
store.steps(id, entries.toString())
true
} != true
)
persistenceFailed = true
}
}
}
@@ -0,0 +1,218 @@
package cn.ilapage.goauto.agent.persistence
import android.database.sqlite.SQLiteDatabase
import java.text.SimpleDateFormat
import java.util.Date
import java.util.Locale
import java.util.TimeZone
import org.json.JSONObject
/** Shared SQL is exercised against SQLite JDBC; Android only supplies the driver. */
internal interface SyncSql {
fun execute(sql: String, args: List<Any?> = emptyList())
fun query(sql: String, args: List<Any?> = emptyList()): List<Map<String, String?>>
fun <T> transaction(block: () -> T): T
}
internal class AndroidSyncSql(private val open: () -> SQLiteDatabase) : SyncSql {
override fun execute(sql: String, args: List<Any?>) = open().execSQL(sql, args.toTypedArray())
override fun query(sql: String, args: List<Any?>): List<Map<String, String?>> =
open().rawQuery(sql, args.map { it?.toString() }.toTypedArray()).use { c ->
buildList {
while (c.moveToNext()) add(
c.columnNames.associateWith { name ->
val i = c.getColumnIndexOrThrow(name)
if (c.isNull(i)) null else c.getString(i)
}
)
}
}
override fun <T> transaction(block: () -> T): T {
val db = open()
db.beginTransaction()
try {
val result = block()
db.setTransactionSuccessful()
return result
} finally {
db.endTransaction()
}
}
}
internal data class SyncContext(
val deviceId: Long? = null,
val taskId: Long? = null,
val attemptId: String? = null,
val phase: String? = null,
val validatedNetwork: Boolean? = null,
val transport: String = "unknown",
val foreground: Boolean = true,
)
internal data class SyncFailure(
val step: String,
val category: String,
val at: Long,
val durationMs: Long,
val httpStatus: Int? = null,
val apiCode: String? = null,
val exceptionClass: String? = null,
)
internal data class SyncReportSnapshot(val report: JSONObject, val boundary: Long)
internal class SyncDiagnosticStore(private val db: SyncSql) {
@Synchronized
fun start(process: String, at: Long, gap: Long?, context: SyncContext): Long =
db.transaction {
db.execute(
"INSERT INTO sync_round(process_id,started_at,start_gap_ms,device_id,task_id,attempt_id,phase,validated_network,transport,foreground,status) VALUES(?,?,?,?,?,?,?,?,?,?,'running')",
listOf(
process,
at,
gap,
context.deviceId?.takeIf { it > 0 },
context.taskId?.takeIf { it > 0 },
context.attemptId?.takeIf { ATTEMPT_UUID.matches(it) },
context.phase?.takeIf { it in setOf("spec_probe", "purchase") },
context.validatedNetwork?.let { if (it) 1 else 0 },
context.transport.takeIf { it in setOf("wifi", "cellular", "none", "unknown") }
?: "unknown",
if (context.foreground) 1 else 0,
),
)
db.query("SELECT last_insert_rowid() AS id").single().getValue("id")!!.toLong()
}
@Synchronized
fun steps(id: Long, json: String) {
db.execute(
"UPDATE sync_round SET steps=? WHERE id=? AND status='running'",
listOf(json, id),
)
}
@Synchronized
fun finish(id: Long, at: Long, duration: Long, failure: SyncFailure?) =
db.transaction {
val sequence = nextSequence()
db.execute(
"UPDATE sync_round SET ended_at=?,duration_ms=?,status=?,failure_at=?,failure_step=?,failure_category=?,failure_duration_ms=?,http_status=?,api_code=?,exception_class=?,completion_seq=? WHERE id=? AND status='running'",
listOf(
at,
duration,
if (failure == null) "success" else "failure",
failure?.at,
failure?.step,
failure?.category,
failure?.durationMs,
failure?.httpStatus,
failure?.apiCode,
failure?.exceptionClass,
sequence,
id,
),
)
}
@Synchronized
fun cleanup(now: Long) =
db.transaction {
db.execute("DELETE FROM sync_round WHERE started_at < ?", listOf(now - RETENTION_MS))
db.execute(
"DELETE FROM sync_round WHERE id NOT IN (SELECT id FROM sync_round ORDER BY id DESC LIMIT 6000)"
)
}
@Synchronized
fun snapshot(process: String, now: Long): SyncReportSnapshot? =
db.transaction {
val cutoff = now - RETENTION_MS
val boundary =
db.query("SELECT sequence FROM sync_meta WHERE id=1")
.single()
.getValue("sequence")!!
.toLong()
val failures =
db.query(
"SELECT * FROM sync_round WHERE status='failure' AND reported=0 AND completion_seq<=? AND started_at>=? ORDER BY completion_seq DESC",
listOf(boundary, cutoff),
)
val previous =
db.query(
"SELECT duration_ms FROM sync_round WHERE status!='running' AND started_at>=? ORDER BY completion_seq DESC LIMIT 1",
listOf(cutoff),
)
.firstOrNull()
val interrupted =
db.query(
"SELECT MAX(id) AS id FROM sync_round WHERE status='running' AND process_id!=? AND started_at>=?",
listOf(process, cutoff),
)
.single()["id"]
?.toLong() ?: 0L
val sequence = nextSequence()
val latest = failures.firstOrNull()
val report =
JSONObject()
.put("snapshotSeq", sequence)
.put("failedRounds", failures.size)
.put(
"latestFailureAt",
latest?.get("failure_at")?.toLong()?.let(::isoTime) ?: JSONObject.NULL,
)
.put("latestFailureStep", latest?.get("failure_step") ?: JSONObject.NULL)
.put(
"latestFailureCategory",
latest?.get("failure_category") ?: JSONObject.NULL,
)
.put(
"latestFailureDurationMs",
latest?.get("failure_duration_ms")?.toLong()?.coerceIn(0, RETENTION_MS)
?: JSONObject.NULL,
)
.put(
"previousSyncDurationMs",
previous?.get("duration_ms")?.toLong()?.coerceIn(0, RETENTION_MS)
?: JSONObject.NULL,
)
.put("hasInterruptedRound", interrupted > 0)
check(
sequence in 1..9_007_199_254_740_991L &&
report.toString().toByteArray(Charsets.UTF_8).size <= 2048
)
SyncReportSnapshot(report, boundary)
}
@Synchronized
fun ack(snapshot: SyncReportSnapshot) =
db.transaction {
db.execute(
"UPDATE sync_round SET reported=1 WHERE status!='running' AND completion_seq<=?",
listOf(snapshot.boundary),
)
}
private fun nextSequence(): Long {
db.execute("UPDATE sync_meta SET sequence=sequence+1 WHERE id=1")
return db.query("SELECT sequence FROM sync_meta WHERE id=1")
.single()
.getValue("sequence")!!
.toLong()
}
companion object {
const val RETENTION_MS = 86_400_000L
private val ATTEMPT_UUID = Regex("[0-9a-fA-F]{8}(-[0-9a-fA-F]{4}){3}-[0-9a-fA-F]{12}")
private fun isoTime(time: Long) =
SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss.SSS'Z'", Locale.US)
.apply { timeZone = TimeZone.getTimeZone("UTC") }
.format(Date(time))
}
}
@@ -57,6 +57,13 @@ import cn.ilapage.goauto.agent.automation.RuleValidationException
import cn.ilapage.goauto.agent.identity.SecureDeviceStore
import cn.ilapage.goauto.agent.network.AgentApiClient
import cn.ilapage.goauto.agent.network.AgentApiException
import cn.ilapage.goauto.agent.network.HeartbeatReportNegotiation
import cn.ilapage.goauto.agent.persistence.AndroidSyncSql
import cn.ilapage.goauto.agent.persistence.SyncContext
import cn.ilapage.goauto.agent.persistence.SyncDiagnosticRecorder
import cn.ilapage.goauto.agent.persistence.SyncDiagnosticStore
import cn.ilapage.goauto.agent.persistence.SyncReportSnapshot
import cn.ilapage.goauto.agent.persistence.syncTaskIdentity
import cn.ilapage.goauto.agent.network.DeviceInfo
import cn.ilapage.goauto.agent.network.PurchaseAgentTask
import cn.ilapage.goauto.agent.network.ServerUrlPolicy
@@ -115,6 +122,10 @@ class AgentForegroundService : Service() {
private lateinit var stateStore: AgentStateStore
private lateinit var purchaseStore: PurchaseTaskStore
private lateinit var diagnosticStore: AgentDiagnosticStore
private lateinit var syncDiagnostics: SyncDiagnosticRecorder
private val heartbeatReports = HeartbeatReportNegotiation()
private val syncTaskContext = AtomicReference<SyncContext?>(null)
@Volatile private var foregroundStarted = false
private lateinit var diagnosticRecorder: SafeAgentDiagnosticRecorder
private lateinit var failureSnapshotStore: PurchaseFailureSnapshotStore
private lateinit var connectivityManager: ConnectivityManager
@@ -135,6 +146,10 @@ class AgentForegroundService : Service() {
stateStore.setKeepScreenOn(false)
purchaseStore = PurchaseTaskStore(this)
diagnosticStore = AgentDiagnosticStore(this)
syncDiagnostics = SyncDiagnosticRecorder(
SyncDiagnosticStore(AndroidSyncSql { diagnosticStore.writableDatabase }), UUID.randomUUID().toString(),
System::currentTimeMillis, SystemClock::elapsedRealtime,
) { Log.w("GoAutoDiagnostic", "sync_diagnostic_persistence_failed") }
failureSnapshotStore = PurchaseFailureSnapshotStore(diagnosticStore)
diagnosticRecorder = SafeAgentDiagnosticRecorder(
persist = { event ->
@@ -152,6 +167,7 @@ class AgentForegroundService : Service() {
connectivityManager = getSystemService(ConnectivityManager::class.java)
createNotificationChannel()
startForeground(NOTIFICATION_ID, notification("正在启动"))
foregroundStarted = true
registerNetworkCallback()
resumeCollectionCooldown()
executor.scheduleWithFixedDelay(::triggerSync, 0, HEARTBEAT_SECONDS, TimeUnit.SECONDS)
@@ -204,6 +220,8 @@ class AgentForegroundService : Service() {
}
override fun onDestroy() {
foregroundStarted = false
stateStore.stop()
probeHandoff.invalidate()
if (diagnosticInstance === this) diagnosticInstance = null
repurchaseClosed.set(true)
@@ -268,15 +286,63 @@ class AgentForegroundService : Service() {
}
private fun synchronizeAgent(): String {
val round = syncDiagnostics.begin(syncDiagnosticContext())
try {
return synchronizeAgent(round)
} catch (error: Error) {
round.abort()
throw error
} catch (error: java.util.concurrent.CancellationException) {
round.abort()
throw error
} finally {
round.finish()
}
}
private fun syncDiagnosticContext(): SyncContext = try {
val taskId = runningTaskId.get() ?: purchaseStore.activeTaskId()
val state = stateStore.read()
val taskContext = syncTaskIdentity(
taskId, state.currentTaskType.takeIf { state.currentTaskId == taskId },
taskMutex.currentTaskId() != null, syncTaskContext.get(),
) {
purchaseStore.interruptedAttempts().firstOrNull { it.taskId == taskId }?.attemptId
}
val network = connectivityManager.activeNetwork
val capabilities = network?.let(connectivityManager::getNetworkCapabilities)
SyncContext(
deviceId = identityStore.credentials()?.deviceId,
taskId = taskId,
attemptId = taskContext.attemptId,
phase = taskContext.phase,
validatedNetwork = capabilities?.hasCapability(NetworkCapabilities.NET_CAPABILITY_VALIDATED),
transport = when {
network == null -> "none"
capabilities?.hasTransport(NetworkCapabilities.TRANSPORT_WIFI) == true -> "wifi"
capabilities?.hasTransport(NetworkCapabilities.TRANSPORT_CELLULAR) == true -> "cellular"
else -> "unknown"
},
foreground = foregroundStarted,
)
} catch (error: java.util.concurrent.CancellationException) {
throw error
} catch (error: Exception) {
Log.w("GoAutoDiagnostic", "sync_diagnostic_context_failed")
SyncContext(foreground = foregroundStarted)
}
private fun synchronizeAgent(round: SyncDiagnosticRecorder.Round): 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)
val registration = round.step("register") { api.register(deviceInfo(), credentials?.token) }
if (registration.deviceToken != null) {
val issuedToken = registration.deviceToken
identityStore.saveCredentials(registration.deviceId, issuedToken)
@@ -290,17 +356,41 @@ 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 ->
round.step("heartbeat") {
var snapshot: SyncReportSnapshot? = null
heartbeatReports.heartbeat(
api, activeCredentials.token, currentTaskId, AgentCapabilities.supported,
report = { syncDiagnostics.snapshot().also { snapshot = it }?.report },
ack = { snapshot?.let(syncDiagnostics::ack) },
validate = { check(it.deviceId == activeCredentials.deviceId) { "心跳设备身份不一致" } },
)
}
},
activeTaskId = {
if (taskMutex.currentTaskId() != null) runningTaskId.get()
else purchaseStore.activeTaskId().also(runningTaskId::set)
},
executionActive = { taskMutex.currentTaskId() != null },
recover = {
round.step("recover") {
recoverInterruptedPurchases(api, activeCredentials.token, round::handledFailure)
}
},
flush = { round.step("flush") { flushPurchaseOutbox(api, activeCredentials.token) } },
hasPending = { purchaseStore.pendingOutbox().isNotEmpty() },
).run()
sync.failure?.let(round::failed)
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 +399,34 @@ class AgentForegroundService : Service() {
heartbeat = true,
)
updateNotification(if (heartbeat.busy) "在线 · 执行中" else "在线 · 空闲")
scheduleTask(api, activeCredentials.token)
if (sync.canClaim) round.step("claim") { scheduleTask(api, activeCredentials.token) } else MANUAL_BUSY
} catch (error: AgentApiException) {
round.failed(error)
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) {
round.failed(error)
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 +435,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)) {
@@ -607,7 +699,7 @@ class AgentForegroundService : Service() {
publishCurrentPageResult(0L, "采集间隔中,还需 $remaining 秒。", replacementOriginType, replacementOriginTaskId)
return
}
if (!taskMutex.tryAcquire(CURRENT_PAGE_RESERVATION_ID)) {
if (!tryAcquireCurrentPage(taskMutex, working, CURRENT_PAGE_RESERVATION_ID)) {
publishCurrentPageResult(0L, "设备正在执行任务,请稍后再试。", replacementOriginType, replacementOriginTaskId)
return
}
@@ -765,6 +857,7 @@ class AgentForegroundService : Service() {
var resultSafelyStored = false
var knownResultType: String? = null
var snapshotContext: SnapshotContext? = null
var observedSyncTaskContext: SyncContext? = null
var executionEntered = false
var snapshotAttempted = false
var probeDiagnosticEvents = 0
@@ -782,6 +875,8 @@ class AgentForegroundService : Service() {
api.startPurchaseTask(claimed.taskId, UUID.randomUUID().toString(), token)
} else claimed
check(task.status == "running" && task.taskAttemptId.isNotBlank()) { "采购任务没有有效 attempt" }
observedSyncTaskContext = SyncContext(taskId = task.taskId, attemptId = task.taskAttemptId, phase = task.phase)
syncTaskContext.set(observedSyncTaskContext)
recordRepurchaseExecution(task, null)
val diagnosticDeviceId = runCatching { identityStore.credentials()?.takeIf { it.token == token }?.deviceId }.getOrNull()
if (task.taskAttemptId.matches(Regex("[0-9a-fA-F]{8}(-[0-9a-fA-F]{4}){3}-[0-9a-fA-F]{12}")) &&
@@ -864,7 +959,15 @@ class AgentForegroundService : Service() {
panelDiagnostic = { evidence -> Log.i("GoAutoPurchasePanel", "task=${task.taskId};$evidence") },
beforeOrderSubmit = { evidence ->
val boundaryRequestId = UUID.randomUUID().toString()
val finalEvidence = evidence.toLocalJson()
val finalEvidence = JSONObject()
.put("goodsId", evidence.goodsId)
.put("mappedColor", evidence.mappedColor)
.put("mappedSize", evidence.mappedSize)
.put("quantity", evidence.quantity)
.put("unitPriceCent", evidence.unitPriceCent)
.put("addressSuffix", evidence.addressSuffix)
.put("activityName", evidence.activityName)
.toString()
purchaseStore.markOrderSubmitStarted(task.taskId, task.taskAttemptId, boundaryRequestId, finalEvidence)
api.markPurchaseOrderSubmitStarted(task.taskId, boundaryRequestId, token)
},
@@ -929,9 +1032,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)
@@ -940,6 +1044,7 @@ class AgentForegroundService : Service() {
stateStore.update("TASK_ERROR", error.message ?: "采购演练执行异常", tokenStored = true)
} finally {
if (snapshotContext?.phase == "spec_probe" && probeDiagnosticEvents == 0) Log.i("GoAutoDiagnostic", "purchase_diagnostic_no_event")
observedSyncTaskContext?.let { syncTaskContext.compareAndSet(it, null) }
if (!resultSafelyStored) cancelIdleReturn("采购结果未安全保存")
releaseTaskWakeLock()
}
@@ -990,12 +1095,16 @@ class AgentForegroundService : Service() {
.toString()
}
private fun recoverInterruptedPurchases(api: AgentApiClient, token: String) {
private fun recoverInterruptedPurchases(
api: AgentApiClient,
token: String,
onHandledFailure: (AgentApiException) -> Unit,
) {
purchaseStore.interruptedAttempts().forEach { interrupted ->
probeHandoff.beginExecution()
if (interrupted.status == PurchaseTaskStore.STATUS_ORDER_SUBMIT_STARTED) {
val boundaryRequestId = interrupted.orderSubmitRequestId ?: return@forEach
try {
recoverPurchaseBoundary(onHandledFailure) {
api.markPurchaseOrderSubmitStarted(interrupted.taskId, boundaryRequestId, token)
val automation = GoAutoAccessibilityService.instance?.let(::PurchaseLiveAutomation)
val evidence = try { automation?.readOrderResult() } catch (error: Exception) {
@@ -1021,9 +1130,6 @@ class AgentForegroundService : Service() {
}
val requestId = UUID.randomUUID().toString()
purchaseStore.completeAndEnqueue(interrupted.taskId, interrupted.attemptId, requestId, purchaseResultPayload(requestId, interrupted.attemptId, outcome))
} catch (_: AgentApiException) {
// Keep the local irreversible marker. The same server request ID
// is replayed after connectivity returns; the order click is never repeated.
}
return@forEach
}
@@ -1058,6 +1164,18 @@ class AgentForegroundService : Service() {
runCatching { onFirstSubmitted(item, requireNotNull(acknowledgement)) }
.onFailure { probeHandoff.invalidate(); logDiagnosticPersistenceFailure(it) }
},
markRejected = { item, code ->
probeHandoff.invalidate()
purchaseStore.markRejected(item, code)
synchronized(taskMutex) {
val runtime = stateStore.read()
if (shouldClearRejectedPurchase(item.taskId, runtime.currentTaskId, runtime.currentTaskType,
taskMutex.currentTaskId(), purchaseStore.isCurrentAttemptRejected(item.taskId, item.attemptId))) {
stateStore.clearActiveTask(item.taskId)
}
}
sendBroadcast(Intent(ACTION_PURCHASE_REJECTIONS_CHANGED).setPackage(packageName))
},
).flush()
runningTaskId.set(purchaseStore.activeTaskId())
}
@@ -1685,6 +1803,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"
@@ -191,17 +191,27 @@ class AgentSettingsStore(context: Context) {
class AgentStateStore(context: Context) {
private val preferences = context.getSharedPreferences(PREFERENCES, Context.MODE_PRIVATE)
private val publicationGate = AgentStatePublicationGate()
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) = publicationGate.publish {
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()
}
/** This service instance cannot publish a late heartbeat over STOPPED or a restarted service. */
fun stop() = publicationGate.stop {
preferences.edit().putString(STATE_CODE, "STOPPED").putString(CONNECTION_CODE, "STOPPED")
.putString(STATE_MESSAGE, "前台服务已停止").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 +289,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,87 @@
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,
val failure: Exception? = null,
)
/** 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 = captureSync { heartbeat(activeTaskId()) }
var authFailure = isAgentAuthenticationError(pulse.exceptionOrNull())
var recovered = false
var flushed = false
if (!authFailure && !executionActive()) {
val recovery = captureSync(recover)
recovered = recovery.isSuccess
authFailure = isAgentAuthenticationError(recovery.exceptionOrNull())
if (!authFailure) {
val upload = captureSync(flush)
flushed = upload.isSuccess
authFailure = isAgentAuthenticationError(upload.exceptionOrNull())
}
if (!authFailure && (pulse.exceptionOrNull() as? AgentApiException)?.code == "DEVICE_TASK_MISMATCH") {
pulse = captureSync { 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,
failure as? Exception,
)
}
}
private fun <T> captureSync(action: () -> T): Result<T> = try {
Result.success(action())
} catch (error: java.util.concurrent.CancellationException) {
throw error
} catch (error: Exception) {
Result.failure(error)
}
/** Preserve the local irreversible marker when replay fails; the original request ID is reused. */
internal fun recoverPurchaseBoundary(
onHandledFailure: (AgentApiException) -> Unit,
recover: () -> Unit,
) {
try {
recover()
} catch (error: AgentApiException) {
if (isAgentAuthenticationError(error)) throw error
onHandledFailure(error)
}
}
@@ -0,0 +1,16 @@
package cn.ilapage.goauto.agent.service
internal fun shouldClearRejectedPurchase(
rejectedTaskId: Long, runtimeTaskId: Long?, runtimeTaskType: String?,
executingTaskId: Long?, currentAttemptRejected: Boolean,
): Boolean = executingTaskId == null && runtimeTaskId == rejectedTaskId &&
runtimeTaskType == "purchase" && currentAttemptRejected
internal class AgentStatePublicationGate {
private var stopped = false
@Synchronized fun publish(write: () -> Unit) { if (!stopped) write() }
@Synchronized fun stop(write: () -> Unit) { stopped = true; write() }
}
internal fun tryAcquireCurrentPage(mutex: TaskExecutionMutex, syncing: java.util.concurrent.atomic.AtomicBoolean, reservation: Long): Boolean =
synchronized(mutex) { !syncing.get() && mutex.tryAcquire(reservation) }
@@ -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")
@@ -0,0 +1,76 @@
package cn.ilapage.goauto.agent
import cn.ilapage.goauto.agent.network.*
import cn.ilapage.goauto.agent.persistence.syncFailure
import java.net.HttpURLConnection
import java.net.SocketTimeoutException
import java.net.URL
import org.junit.Assert.*
import org.junit.Test
class AgentTransportTimeoutTest {
@Test
fun loopbackPeerThatWithholdsResponseUsesExistingReadTimeoutAndResponseCategory() {
val socket = java.net.ServerSocket(0, 1, java.net.InetAddress.getLoopbackAddress())
val release = java.util.concurrent.CountDownLatch(1)
val accepted = java.util.concurrent.CountDownLatch(1)
val worker =
kotlin.concurrent.thread(isDaemon = true) {
socket.accept().use {
accepted.countDown()
release.await(20, java.util.concurrent.TimeUnit.SECONDS)
}
}
try {
val error =
assertThrows(AgentTransportTimeout::class.java) {
AgentApiClient("http://127.0.0.1:${socket.localPort}")
.heartbeat("synthetic", null)
}
assertEquals(0, accepted.count)
assertEquals(RequestStage.RESPONSE, error.stage)
assertEquals("timeout_response", syncFailure(error, "heartbeat", 1000, 1).category)
} finally {
release.countDown()
socket.close()
worker.join(1000)
}
}
@Test
fun timeoutCategoryComesFromActualConnectionStageNotExceptionMessage() {
for (stage in RequestStage.values()) {
val connection =
object : HttpURLConnection(URL("http://127.0.0.1")) {
override fun connect() {
if (stage == RequestStage.CONNECT)
throw SocketTimeoutException("Read timed out")
}
override fun disconnect() {}
override fun usingProxy() = false
override fun getOutputStream(): java.io.OutputStream {
if (stage == RequestStage.UNKNOWN)
throw SocketTimeoutException("connect timed out")
return java.io.ByteArrayOutputStream()
}
override fun getResponseCode(): Int {
throw SocketTimeoutException("connect timed out")
}
}
val api = AgentApiClient("http://127.0.0.1", connectionFactory = { connection })
val error =
assertThrows(AgentTransportTimeout::class.java) { api.heartbeat("synthetic", null) }
assertEquals(stage, error.stage)
assertEquals(
"timeout_${stage.name.lowercase()}",
syncFailure(error, "heartbeat", 1000, 1).category,
)
assertEquals(10_000, connection.connectTimeout)
assertEquals(15_000, connection.readTimeout)
}
}
}
@@ -0,0 +1,315 @@
package cn.ilapage.goauto.agent
import cn.ilapage.goauto.agent.network.*
import cn.ilapage.goauto.agent.persistence.*
import java.net.InetAddress
import java.net.ServerSocket
import java.sql.DriverManager
import java.util.concurrent.CopyOnWriteArrayList
import kotlin.concurrent.thread
import org.json.JSONArray
import org.json.JSONObject
import org.junit.Assert.*
import org.junit.Test
class HeartbeatReportCompatibilityTest {
@Test
fun trueFalseMissingFreshProcessAndTokenChangeControlNextRequest() {
val responses =
listOf(
success(true),
success(false),
success(true),
success(null),
success(true),
success(true),
success(true),
)
HttpFixture(responses).use { http ->
val negotiation = HeartbeatReportNegotiation()
var ack = 0
repeat(5) {
negotiation.heartbeat(
AgentApiClient(http.origin),
"token",
null,
emptyList(),
{ report() },
{ ack++ },
)
}
negotiation.heartbeat(
AgentApiClient(http.origin),
"changed",
null,
emptyList(),
{ report() },
{ ack++ },
)
HeartbeatReportNegotiation()
.heartbeat(
AgentApiClient(http.origin),
"changed",
null,
emptyList(),
{ report() },
{ ack++ },
)
assertEquals(
listOf(false, true, false, true, false, false, false),
http.requests.map { it.has("clientReport") },
)
assertEquals(2, ack)
}
}
@Test
fun originChangeRenegotiatesEvenWhenTokenUnchanged() {
HttpFixture(listOf(success(true))).use { first ->
HttpFixture(listOf(success(true))).use { second ->
val negotiation = HeartbeatReportNegotiation()
for (http in listOf(first, second)) negotiation.heartbeat(
AgentApiClient(http.origin),
"token",
null,
emptyList(),
{ report() },
{ fail("new origin cannot acknowledge") },
)
assertFalse(second.requests.single().has("clientReport"))
}
}
}
@Test
fun onlyExactOldServerRejectionCanDowngradeAndFallbackIsAtMostOnce() {
for ((status, code, message) in
listOf(
401 to "INVALID_REQUEST" to "请求 JSON 无效",
403 to "INVALID_REQUEST" to "请求 JSON 无效",
409 to "DEVICE_TASK_MISMATCH" to "请求 JSON 无效",
500 to "INVALID_REQUEST" to "请求 JSON 无效",
422 to "INVALID_REQUEST" to "different",
422 to "OTHER" to "请求 JSON 无效",
)
.map { Triple(it.first.first, it.first.second, it.second) }) {
HttpFixture(listOf(success(true), failure(status, code, message))).use { http ->
val negotiation = HeartbeatReportNegotiation()
val api = AgentApiClient(http.origin)
negotiation.heartbeat(
api,
"token",
null,
emptyList(),
{ report() },
{ fail("ack") },
)
assertThrows(AgentApiException::class.java) {
negotiation.heartbeat(
api,
"token",
null,
emptyList(),
{ report() },
{ fail("ack") },
)
}
assertEquals(2, http.requests.size)
}
}
HttpFixture(
listOf(
success(true),
failure(422, "INVALID_REQUEST", "请求 JSON 无效"),
failure(422, "INVALID_REQUEST", "请求 JSON 无效"),
)
)
.use { http ->
val negotiation = HeartbeatReportNegotiation()
val api = AgentApiClient(http.origin)
negotiation.heartbeat(
api,
"token",
null,
emptyList(),
{ report() },
{ fail("ack") },
)
assertThrows(AgentApiException::class.java) {
negotiation.heartbeat(
api,
"token",
null,
emptyList(),
{ report() },
{ fail("ack") },
)
}
assertEquals(3, http.requests.size)
}
}
@Test
fun missingReportDoesNotFallbackAndInvalidIdentityCannotAcknowledge() {
HttpFixture(listOf(failure(422, "INVALID_REQUEST", "请求 JSON 无效"))).use { http ->
assertThrows(AgentApiException::class.java) {
HeartbeatReportNegotiation()
.heartbeat(
AgentApiClient(http.origin),
"t",
null,
emptyList(),
{ report() },
{ fail("ack") },
)
}
assertEquals(1, http.requests.size)
}
HttpFixture(listOf(success(true), success(false))).use { http ->
val negotiation = HeartbeatReportNegotiation()
val api = AgentApiClient(http.origin)
negotiation.heartbeat(api, "token", null, emptyList(), { report() }, { fail("ack") })
assertThrows(IllegalStateException::class.java) {
negotiation.heartbeat(
api,
"token",
null,
emptyList(),
{ report() },
{ fail("ack") },
) {
check(it.deviceId == 2L)
}
}
}
}
@Test
fun lostResponseDoesNotAcknowledgeOrResend() {
HttpFixture(listOf(success(true), null)).use { http ->
val negotiation = HeartbeatReportNegotiation()
val api = AgentApiClient(http.origin)
negotiation.heartbeat(api, "t", null, emptyList(), { report() }, { fail("ack") })
assertThrows(java.io.IOException::class.java) {
negotiation.heartbeat(api, "t", null, emptyList(), { report() }, { fail("ack") })
}
assertEquals(2, http.requests.size)
}
}
@Test
fun actualSerializedReportsFromSqliteReachHttpAndExportCrossLanguageFixtures() {
DriverManager.getConnection("jdbc:sqlite::memory:").use { db ->
AgentDiagnosticSchema.syncStatements.forEach {
db.createStatement().use { s -> s.execute(it) }
}
val store = SyncDiagnosticStore(JdbcSyncSql(db))
val at = 1_700_000_000_000L
val zero = store.snapshot("p", at)!!
val id = store.start("p", at, null, SyncContext())
store.finish(id, at + 100, 100, SyncFailure("heartbeat", "network", at + 100, 100))
val failed = store.snapshot("p", at + 200)!!
store.ack(failed)
store.start("old", at + 300, null, SyncContext())
val interrupted = store.snapshot("p", at + 400)!!
HttpFixture(List(4) { success(true) }).use { http ->
val negotiation = HeartbeatReportNegotiation()
val api = AgentApiClient(http.origin)
negotiation.heartbeat(api, "synthetic", null, emptyList(), { null }, {})
for (snapshot in listOf(zero, failed, interrupted)) negotiation.heartbeat(
api,
"synthetic",
null,
emptyList(),
{ snapshot.report },
{},
)
val reports =
JSONArray(http.requests.drop(1).map { it.getJSONObject("clientReport") })
assertEquals(3, reports.length())
for (i in 0..2) {
assertEquals(8, reports.getJSONObject(i).length())
assertTrue(reports.getJSONObject(i).toString().toByteArray().size <= 2048)
}
val file = java.io.File("../build/diagnostics378-client-reports.json")
requireNotNull(file.parentFile).mkdirs()
file.writeText(reports.toString(2), Charsets.UTF_8)
}
}
}
private fun report() =
JSONObject()
.put("snapshotSeq", 1)
.put("failedRounds", 0)
.put("latestFailureAt", JSONObject.NULL)
.put("latestFailureStep", JSONObject.NULL)
.put("latestFailureCategory", JSONObject.NULL)
.put("latestFailureDurationMs", JSONObject.NULL)
.put("previousSyncDurationMs", JSONObject.NULL)
.put("hasInterruptedRound", false)
private fun success(capability: Boolean?): Pair<Int, String> =
200 to
JSONObject()
.put(
"data",
JSONObject().put("deviceId", 1).put("online", true).put("busy", false).apply {
capability?.let { put("acceptsClientReport", it) }
},
)
.toString()
private fun failure(status: Int, code: String, message: String) =
status to JSONObject().put("code", code).put("message", message).toString()
}
private class HttpFixture(responses: List<Pair<Int, String>?>) : AutoCloseable {
private val socket =
ServerSocket(0, 10, InetAddress.getLoopbackAddress()).apply { soTimeout = 5000 }
val origin = "http://127.0.0.1:${socket.localPort}"
val requests = CopyOnWriteArrayList<JSONObject>()
private var error: Throwable? = null
private val worker =
thread(isDaemon = true) {
try {
for (response in responses) socket.accept().use { client ->
client.soTimeout = 5000
val input = client.getInputStream().bufferedReader()
var length = 0
while (true) {
val line = input.readLine() ?: error("missing headers")
if (line.isEmpty()) break
if (line.startsWith("Content-Length:", true))
length = line.substringAfter(':').trim().toInt()
}
check(length in 1..4096)
val body = CharArray(length)
var read = 0
while (read < length) {
val count = input.read(body, read, length - read)
check(count > 0)
read += count
}
requests.add(JSONObject(String(body)))
if (response != null) {
val bytes = response.second.toByteArray(Charsets.UTF_8)
client
.getOutputStream()
.write(
("HTTP/1.1 ${response.first} OK\r\nContent-Length: ${bytes.size}\r\nConnection: close\r\n\r\n")
.toByteArray() + bytes
)
}
}
} catch (e: Throwable) {
if (!socket.isClosed) error = e
}
}
override fun close() {
socket.close()
worker.join(1000)
error?.let { throw AssertionError("HTTP fixture failed", it) }
}
}
@@ -0,0 +1,86 @@
package cn.ilapage.goauto.agent
import cn.ilapage.goauto.agent.network.*
import java.net.ServerSocket
import kotlin.concurrent.thread
import org.json.JSONObject
import org.junit.Assert.*
import org.junit.Test
class HeartbeatReportProtocolTest {
@Test
fun negotiatesThenDowngradesOnceWithSameRequestIdAndNoAck() {
val requests = mutableListOf<JSONObject>()
val server =
ServerSocket(0, 3, java.net.InetAddress.getLoopbackAddress()).apply { soTimeout = 5000 }
val worker =
thread(isDaemon = true) {
repeat(3) {
server.accept().use { socket ->
socket.soTimeout = 5000
val input = socket.getInputStream().bufferedReader()
var length = 0
while (true) {
val line = input.readLine() ?: error("missing header")
if (line.isEmpty()) break
if (line.startsWith("Content-Length:", true))
length = line.substringAfter(':').trim().toInt()
}
check(length in 1..4096)
val chars = CharArray(length)
var read = 0
while (read < length) {
val count = input.read(chars, read, length - read)
check(count > 0)
read += count
}
val request = JSONObject(String(chars))
requests.add(request)
val reject = request.has("clientReport")
val response =
if (reject) """{"code":"INVALID_REQUEST","message":"请求 JSON 无效"}"""
else
"""{"data":{"deviceId":1,"online":true,"busy":false,"acceptsClientReport":true}}"""
val bytes = response.toByteArray()
socket
.getOutputStream()
.write(
("HTTP/1.1 ${if(reject) 422 else 200} OK\r\nContent-Length: ${bytes.size}\r\nConnection: close\r\n\r\n")
.toByteArray() + bytes
)
}
}
}
try {
val negotiation = HeartbeatReportNegotiation()
var ack = 0
val api = AgentApiClient("http://127.0.0.1:${server.localPort}")
negotiation.heartbeat(
api,
"token",
null,
emptyList(),
{ JSONObject().put("snapshotSeq", 1) },
{ ack++ },
)
negotiation.heartbeat(
api,
"token",
null,
emptyList(),
{ JSONObject().put("snapshotSeq", 1) },
{ ack++ },
)
assertEquals(3, requests.size)
assertFalse(requests[0].has("clientReport"))
assertTrue(requests[1].has("clientReport"))
assertFalse(requests[2].has("clientReport"))
assertEquals(requests[1].getString("requestId"), requests[2].getString("requestId"))
assertEquals(0, ack)
assertFalse(negotiation.acceptsReport)
} finally {
server.close()
worker.join(1000)
}
}
}
@@ -3,171 +3,224 @@ package cn.ilapage.goauto.agent
import cn.ilapage.goauto.agent.automation.*
import org.junit.Assert.*
import org.junit.Test
import org.json.JSONObject
/** #376 replaces #374's final value-proof contract. All fixtures are synthetic. */
/** Final capture only; derived from the existing #331/#332 bounded purchase-sheet fixtures. */
class PurchaseFinalSpecConfirmationTest {
@Test fun `local final evidence distinguishes target observed quantity and scoped summary flags`() {
val json = JSONObject(check(sheet("已选 黑色 XL"), input().copy(quantity = 2)).toLocalJson())
assertFalse(json.has("quantity"))
assertEquals(2, json.getLong("targetQuantity"))
assertEquals(1, json.getLong("observedQuantity"))
assertTrue(json.getBoolean("selectedSummaryPresent"))
assertFalse(json.getBoolean("summaryPromptPresent"))
val prompt = check(sheet("请 选择:颜色"), input())
assertFalse(prompt.selectedSummaryPresent)
assertTrue(prompt.summaryPromptPresent)
val external = sheet(null).addText("r/outside", "已选 黑色 XL", NodeBounds(20, 400, 800, 460))
assertFalse(check(external, input()).selectedSummaryPresent)
assertFalse(check(sheet("#黑色 XL#"), input()).selectedSummaryPresent)
@Test fun `bracket display variants and offscreen color are proven by full panel summary`() {
val snapshot = sheet("已选 [香槟]RW圆领 L码[105-115斤]", colors = listOf("其他颜色"), sizes = listOf("L码[105-115斤]"))
.selected("size/o0/t")
accepts(snapshot, "【香槟】RW圆领", "L码【105-115斤】")
}
@Test fun `observed quantity is nullable for missing or multiple numeric editors`() {
val single = sheet(null)
assertEquals(1L, observedPurchaseQuantity(single))
assertNull(observedPurchaseQuantity(single.mapNodes { if (it.className == "android.widget.EditText") it.copy(visible = false) else it }))
val quantity = single.nodes.single { it.className == "android.widget.EditText" }
assertNull(observedPurchaseQuantity(single.copy(nodes = single.nodes + quantity.copy(path = "r/second"))))
val evidence = check(single, input()).copy(observedQuantity = null)
val json = JSONObject(evidence.toLocalJson())
assertTrue(json.has("observedQuantity"))
assertTrue(json.isNull("observedQuantity"))
@Test fun `summary above scroll area is accepted for scrollable and non scrollable sheets`() {
for (scrollable in listOf(true, false)) {
accepts(sheet("已选 黑色 XL", scrollable = scrollable))
}
}
@Test fun `summary prompt is scoped to visible nonclickable header rather than option or external text`() {
val normal = sheet(null)
val driver = object : PurchaseUiDriver {
override fun capture() = normal
override fun swipePurchaseIn(target: SnapshotNode, direction: SwipeDirection, durationMs: Long): Boolean = error("unexpected swipe")
override fun clickFresh(target: SnapshotNode): FreshActionResult = error("unexpected click")
override fun tapPurchaseFresh(target: SnapshotNode): FreshActionResult = error("unexpected tap")
override fun inputFresh(target: SnapshotNode, value: String): FreshActionResult = error("unexpected input")
override fun swipePurchase(direction: SwipeDirection, durationMs: Long): Boolean = error("unexpected swipe")
override fun backPurchase(): Boolean = error("unexpected back")
}
val live = PurchaseLiveAutomation(driver, pause = {})
for (prompt in listOf("请选择颜色", "请 选择:颜色", "請選擇: 顏色")) {
assertTrue(live.selectionPromptPresent(sheet(prompt)))
}
val summaryBounds = NodeBounds(384, 791, 1068, 852)
for ((index, snapshot) in listOf(
normal.addText("r/external", "请选择颜色", summaryBounds),
normal.addText("r/sheet/fake", "请选择颜色", summaryBounds).mapNodes { if (it.path.endsWith("fake")) it.copy(parentPath = "r") else it },
sheet("请选择颜色").mapNodes { if (it.path.endsWith("body/summary")) it.copy(visible = false) else it },
sheet("请选择颜色").mapNodes { if (it.path.endsWith("body/summary")) it.copy(clickable = true) else it },
normal.addText("r/sheet/body/list/color/o0/prompt", "请选择颜色", NodeBounds(40, 1600, 300, 1640)),
normal.mapNodes { if (it.path.startsWith("r/sheet/body/qty")) it.copy(bounds = it.bounds.copy(top = 1990, bottom = 2020)) else it }
.addText("r/sheet/body/list/color/o0/prompt", "请选择颜色", NodeBounds(40, 1600, 300, 1640)),
sheet("请选择颜色").mapNodes { if (it.path in listOf("r/sheet", "r/sheet/body")) it.copy(bounds = NodeBounds(0, 0, 1080, 2216)) else it },
).withIndex()) {
assertFalse("variant=$index", live.selectionPromptPresent(snapshot))
assertFalse(PurchaseSelectionEvidence.observe(snapshot).summaryPromptPresent)
}
// A missing final button never adds a pre-address failure by itself.
assertFalse(live.selectionPromptPresent(normal.mapNodes { if (it.path.startsWith("r/sheet/submit")) it.copy(visible = false) else it }))
}
@Test fun `normal button permits missing heading marketing summary and offscreen size`() {
accepts(sheet("#示例活动#" ).mapNodes { if (it.path.endsWith("color/h")) it.copy(visible = false) else it })
accepts(sheet("#示例活动#").mapNodes { if (it.path.contains("list/size")) it.copy(visible = false) else it })
@Test fun `checked only exact option and physical parent child duplicates provide selection proof`() {
accepts(sheet(null).selected("color/o0/img", checked = true).selected("size/o0/t", checked = true))
}
@Test fun `normal button permits wrong selected values unavailable targets and collisions`() {
accepts(sheet(null).mapNodes {
when {
it.path.endsWith("color/o0") -> it.copy(selected = true, contentDescription = "白色")
it.path.endsWith("size/o0/t") -> it.copy(selected = true, text = "L")
@Test fun `selected other color or size overrides matching summary and positive target`() {
rejects(sheet("已选 黑色 XL", colors = listOf("黑色", "白色")).selected("color/o0").selected("color/o1"))
rejects(sheet("已选 黑色 XL", sizes = listOf("XL", "L")).selected("size/o1/t"))
}
@Test fun `selected bracket variants are not a conflict`() {
accepts(sheet("已选 [白色] XL", colors = listOf("[白色]")).selected("color/o0").selected("size/o0/t"), "【白色】")
}
@Test fun `distinct physical identical or normalized target options are collisions`() {
for (colors in listOf(listOf("黑色", "黑色"), listOf("【白色】", "[白色]"))) {
rejects(sheet("已选 ${colors[0]} XL", colors = colors).selected("color/o0"), colors[0])
}
}
@Test fun `explicit unavailable target cannot be rescued by summary`() {
rejects(sheet("已选 黑色 XL").mapNodes { if (it.path.contains("color/o0")) it.copy(enabled = false) else it })
}
@Test fun `all required selected proofs take priority over a single non proving summary`() {
accepts(sheet("已选 其他颜色 其他尺码").selected("color/o0/img").selected("size/o0/t"))
rejects(sheet("已选 白色 L", colors = listOf("黑色", "白色"), sizes = listOf("XL", "L"))
.selected("color/o0/img").selected("size/o0/t")) // Known alternative composition is a conflict.
}
@Test fun `selected options retain established parser separation of trailing display price`() {
accepts(sheet("已选 黑色 XL", colors = listOf("黑色 ¥19.60"), sizes = listOf("XL ¥19.60"))
.selected("color/o0/img").selected("size/o0/t"))
rejects(sheet("已选 黑色 XL", colors = listOf("黑色加绒 ¥19.60")).selected("color/o0/img"))
rejects(sheet("已选 黑色 XL", colors = listOf("黑色 ¥19.60", "黑色 ¥20.00")).selected("color/o0/img"))
}
@Test fun `unselected visible options alone and unrelated title cannot prove selection`() {
rejects(sheet(null).addText("r/title", "示例黑色 XL 商品"))
}
@Test fun `panel external prefixed summary is ignored even at same coordinates`() {
rejects(sheet(null).addText("r/fake", "已选 黑色 XL"))
accepts(sheet("已选 黑色 XL").addText("r/fake", "已选 白色 L"))
}
@Test fun `structural ancestry overrides textual path prefix`() {
rejects(sheet(null).addText("r/sheet/fake", "已选 黑色 XL", parent = "r"))
}
@Test fun `whole summary rejects substring sizes colors extra dimensions quantity and unsupported separators`() {
for (summary in listOf("已选 黑色 XL", "已选 黑色 加绒 L", "已选 黑色 L 2件", "已选 黑色,L", "已选 L 黑色")) {
rejects(sheet(summary, sizes = listOf("L")), size = "L")
}
rejects(sheet("已选 [白色+黑色]超值两件装 XL", colors = emptyList()), "【白色】")
}
@Test fun `known candidates making another full composition reject ambiguous summary`() {
rejects(sheet("已选 黑色 加绒 XL", colors = listOf("黑色", "黑色 加绒"), sizes = listOf("加绒 XL", "XL")), size = "加绒 XL")
rejects(sheet("已选 黑色 加绒 XL", colors = listOf("黑色", "黑色 加绒"), sizes = listOf("加绒 XL", "XL")), color = "黑色 加绒")
}
@Test fun `summary permits invisible color despite another unselected visible color`() {
accepts(sheet("已选 黑色 XL", colors = listOf("白色")))
}
@Test fun `required single dimension must equal whole summary and empty pair is never proof`() {
accepts(sheet("已选 黑色"), size = "")
accepts(sheet("已选 XL"), color = "")
rejects(sheet("已选 黑色 XL"), size = "")
rejects(sheet("已选 黑色 XL"), color = "")
rejects(sheet("已选 黑色 XL"), color = " ", size = "")
}
@Test fun `anchored selected prefix variants preserve internal colons commas and whitespace`() {
for (prefix in listOf("已选", "已选 ", "已选:", "已选: ", " 已选: ")) {
accepts(sheet("$prefix 黑色:加绒,款 XL", colors = emptyList()), color = "黑色:加绒,款")
}
accepts(sheet("已选 黑色 加绒 XL", colors = emptyList()), color = "黑色 加绒")
rejects(sheet("已选 黑色 加绒 XL", colors = emptyList()), color = "黑色加绒")
rejects(sheet("已选 黑色:加绒 XL", colors = emptyList()), color = "黑色加绒")
rejects(sheet("已选::黑色 XL"))
accepts(sheet("已选::黑色 XL", colors = emptyList()), color = ":黑色")
}
@Test fun `conflicting scoped summaries reject even both exact options selected`() {
rejects(sheet("已选 黑色 XL").selected("color/o0").selected("size/o0/t").addText("r/sheet/second", "已选 白色 L"))
accepts(sheet("已选 黑色 XL").addText("r/sheet/second", "已选:黑色 XL"))
}
@Test fun `display normalization is narrow and content preserving`() {
accepts(sheet("已选 [白色](加绒) RW12 XL", colors = emptyList()), color = "〔白色〕(加绒) RW12")
rejects(sheet("已选 白色 RW12 XL", colors = emptyList()), color = "白色 rw12")
rejects(sheet("已选 白色 1 XL", colors = emptyList()), color = "白色 ①")
rejects(sheet("已选 白色 XL", colors = emptyList()), color = "白色 ¥19.6")
}
@Test fun `existing structurally recognized unprefixed summary is usable`() {
accepts(prefixlessSheet())
rejects(sheet("黑色 XL")) // No quantity/header relationship in this fixture.
rejects(sheet(null).addText("r/title", "黑色 XL"))
rejects(prefixlessSheet().addText("r/sheet/body/info/second", "黑色 L"))
}
@Test fun `same dimension groups retain earlier selected conflicts and physical candidates`() {
val base = sheet("已选 黑色 XL").addText("r/sheet/body/list/extra/h", "颜色分类")
.addText("r/sheet/body/list/extra/o", "白色")
.mapNodes { when (it.path) {
"r/sheet/body/list/extra/h" -> it.copy(bounds = NodeBounds(20, 1000, 500, 1040), parentPath = "r/sheet/body/list")
"r/sheet/body/list/extra/o" -> it.copy(bounds = NodeBounds(20, 1045, 500, 1100), parentPath = "r/sheet/body/list", clickable = true, selected = true)
"r/sheet/body/list" -> it.copy(bounds = NodeBounds(0, 980, 1080, 2036))
else -> it
}
})
accepts(sheet("已选 白色 L").mapNodes { if (it.path.contains("color/o0")) it.copy(enabled = false) else it })
accepts(sheet(null).let { it.copy(nodes = it.nodes + it.nodes.filter { n -> n.path.contains("color/o0") }.map { n ->
n.copy(path = n.path.replace("o0", "o1"), parentPath = n.parentPath?.replace("o0", "o1"), bounds = n.bounds.copy(left = 400, right = 700))
}) })
} }
rejects(base)
rejects(base.mapNodes { if (it.path.endsWith("extra/o")) it.copy(text = "黑色") else it })
}
@Test fun `normal button permits mismatched quantity and summary prompt`() {
accepts(sheet("请选择:颜色分类 尺码"), input().copy(quantity = 2))
accepts(sheet("請選擇: 顏色 尺碼"))
}
@Test fun `special button rejects simplified traditional whitespace and split prompt`() {
for (prompt in listOf("选择颜色分类及尺码后,提交订单", "選擇顏色及尺碼後,提交訂單", " 选 择 颜色 后 , 提交 订单 ")) {
rejects(sheet("已选 黑色 XL").button(prompt), "PURCHASE_SPEC_NOT_MATCHED")
}
rejects(splitButton(), "PURCHASE_SPEC_NOT_MATCHED")
rejects(splitButton().mapNodes { if (it.path == "r/sheet/submit/l") it.copy(text = "选择颜色后,提交订单") else it }, "PURCHASE_SPEC_NOT_MATCHED")
}
@Test fun `whole match ignores prefix suffix hidden external and independent child fragments`() {
for (label in listOf("提示选择颜色后,提交订单", "选择颜色后,提交订单优惠", "选择颜色提交订单", "提交订单")) accepts(sheet(null).button(label))
accepts(sheet(null).addText("r/outside", "选择颜色后,提交订单", NodeBounds(20, 10, 800, 70)))
accepts(sheet(null).addText("r/sheet/submit/hidden", "选择颜色后,提交订单", NodeBounds(20, 2140, 800, 2170)).mapNodes {
if (it.path.endsWith("hidden")) it.copy(visible = false) else it
})
accepts(sheet(null).mapNodes { when (it.path) {
"r/sheet/submit" -> it.copy(text = "选择颜色后,提交订单")
"r/sheet/submit/l/t" -> it.copy(text = "选择颜色后,提交订单", visible = false)
else -> it
} }) // Do not use a parent's aggregate label when its labelled child is hidden.
accepts(splitButton().mapNodes { if (it.path.endsWith("/first")) it.copy(clickable = true) else it })
accepts(splitButton().mapNodes { if (it.path.endsWith("/first")) it.copy(path = "r/outside", parentPath = "r") else it })
}
@Test fun `address price and missing button still reject without actions`() {
rejects(sheet(null), "PURCHASE_PRICE_OUT_OF_RANGE", input().copy(maxUnitPriceCent = 1))
rejects(sheet(null).mapNodes { if (it.path.endsWith("addr/a/detail")) it.copy(text = "示例地址") else it }, "PURCHASE_ADDRESS_UPDATE_FAILED")
rejects(sheet(null).mapNodes { if (it.path.startsWith("r/sheet/submit")) it.copy(visible = false) else it }, "PURCHASE_SUBMIT_TARGET_AMBIGUOUS")
}
@Test fun `collector opt in retains independent candidates without changing regular parsing`() {
val snapshot = SpecPanelFixtures.sheet(SpecPanelFixtures.Sheet(colorLabels = listOf("黑色", "黑色")))
@Test fun `collector default still collapses same raw labels and leaves summary behavior unchanged`() {
val snapshot = sheet("已选 黑色 XL", colors = listOf("黑色", "黑色"))
val config = PurchaseRehearsalExecutor.DEFAULT_COLLECTOR
val regular = PddScreenParser.parse(snapshot, config, "376", null)
val final = PddScreenParser.parse(snapshot, config, "376", null, forFinalConfirmation = true)
val regular = PddScreenParser.parse(snapshot, config, "374", null)
val final = PddScreenParser.parse(snapshot, config, "374", null, forFinalConfirmation = true)
assertEquals(1, regular.dimensions.single { it.key == "color" }.values.size)
assertEquals(2, final.dimensions.single { it.key == "color" }.values.size)
assertEquals(regular.selectedSummary, final.selectedSummary)
assertTrue(regular.finalSelectionSummaries.isEmpty())
}
private fun splitButton(): UiSnapshot = sheet(null).button("提交订单").mapNodes {
if (it.path.endsWith("submit/l/t")) it.copy(bounds = NodeBounds(600, 2180, 1000, 2216)) else it
}.addText("r/sheet/submit/l/first", "选择颜色后,", NodeBounds(20, 2180, 590, 2216))
@Test fun `independently bounded selected options survive absent full panel boundary`() {
val snapshot = sheet(null).selected("color/o0/img").selected("size/o0/t")
.mapNodes { if (it.path == "r/sheet") it.copy(bounds = NodeBounds(0, 0, 1080, 2216)) else it }
PurchaseFinalSpecVerifier.verify(snapshot, input()) { fail(it) }
val noSelections = snapshot.mapNodes { it.copy(selected = false, checked = false) }
.addText("r/sheet/fake", "已选 黑色 XL")
try {
PurchaseFinalSpecVerifier.verify(noSelections, input()) {}
fail("unproven full panel must not supply summary")
} catch (failure: PurchaseLiveException) { assertEquals("PURCHASE_SPEC_NOT_MATCHED", failure.code) }
}
private fun sheet(summary: String?): UiSnapshot = SpecPanelFixtures.sheet(SpecPanelFixtures.Sheet(colorLabels = listOf("黑色"), sizeLabels = listOf("XL")))
@Test fun `price quantity and address checks remain enforced`() {
rejects(sheet("已选 黑色 XL"), expectedCode = "PURCHASE_PRICE_OUT_OF_RANGE", input = input().copy(maxUnitPriceCent = 1))
rejects(sheet("已选 黑色 XL"), expectedCode = "PURCHASE_QUANTITY_MISMATCH", input = input().copy(quantity = 2))
rejects(sheet("已选 黑色 XL").mapNodes { if (it.path == "r/sheet/body/addr/a/detail") it.copy(text = "示例地址") else it }, expectedCode = "PURCHASE_ADDRESS_UPDATE_FAILED")
}
private fun sheet(
summary: String?, colors: List<String> = listOf("黑色"), sizes: List<String> = listOf("XL"), scrollable: Boolean = true,
): UiSnapshot = SpecPanelFixtures.sheet(SpecPanelFixtures.Sheet(colorLabels = colors, sizeLabels = sizes, listScrollable = scrollable))
.mapNodes { when (it.path) {
"r/sheet/body/summary" -> it.copy(text = summary, visible = summary != null)
"r/sheet/body/addr/a/detail" -> it.copy(text = ADDRESS)
"r/sheet/submit/l/t" -> it.copy(text = "提交订单")
else -> it
} }
/** Reuses the info-above-list shape from PddProductDetailCollectorTest.prefixlessPanel. */
private fun prefixlessSheet(): UiSnapshot {
val base = sheet("黑色 XL")
val info = base.nodes.single { it.path == "r/sheet/body/qty" }.copy(
path = "r/sheet/body/info", bounds = NodeBounds(0, 770, 1080, 1000), parentPath = "r/sheet/body",
)
return base.mapNodes { when {
it.path == "r/sheet/body/summary" -> it.copy(path = "r/sheet/body/info/summary", parentPath = info.path)
it.path == "r/sheet/body/qty" -> it.copy(path = "r/sheet/body/info/qty", parentPath = info.path)
it.path.startsWith("r/sheet/body/qty/") -> it.copy(path = it.path.replace("body/qty/", "body/info/qty/"), parentPath = "r/sheet/body/info/qty")
else -> it
} }.let { it.copy(nodes = it.nodes + info) }
}
private fun UiSnapshot.mapNodes(transform: (SnapshotNode) -> SnapshotNode) = copy(nodes = nodes.map(transform))
private fun UiSnapshot.button(label: String) = mapNodes { if (it.path.endsWith("submit/l/t")) it.copy(text = label) else it }
private fun UiSnapshot.addText(path: String, text: String, bounds: NodeBounds) = copy(nodes = nodes + SnapshotNode(
path, path.substringBeforeLast('/'), text, null, null, "android.widget.TextView", bounds,
private fun UiSnapshot.selected(suffix: String, checked: Boolean = false) = mapNodes {
if (it.path.endsWith(suffix)) it.copy(selected = !checked, checked = checked) else it
}
private fun UiSnapshot.addText(path: String, text: String, parent: String = path.substringBeforeLast('/')) = copy(nodes = nodes + SnapshotNode(
path, parent, text, null, null, "android.widget.TextView", NodeBounds(384, 791, 1068, 852),
false, false, false, false, true, true,
))
private fun input() = PurchaseExecutionInput(
taskId = 376, executionMode = "live", phase = "purchase", url = "https://mobile.yangkeduo.com/goods.html?goods_id=376",
goodsId = "376", mappedColor = "黑色", mappedSize = "XL", quantity = 1,
minUnitPriceCent = 1, maxUnitPriceCent = 3000, addressSuffix = "_cg376",
private fun input(color: String = "黑色", size: String = "XL") = PurchaseExecutionInput(
taskId = 374, executionMode = "live", phase = "purchase", url = "https://mobile.yangkeduo.com/goods.html?goods_id=374",
goodsId = "374", mappedColor = color, mappedSize = size, quantity = 1,
minUnitPriceCent = 1, maxUnitPriceCent = 3000, addressSuffix = "_cg374",
)
private fun check(snapshot: UiSnapshot, input: PurchaseExecutionInput, diagnostics: MutableList<String> = mutableListOf()): FinalConfirmationEvidence {
val driver = object : PurchaseUiDriver {
override fun capture() = snapshot
override fun backPurchase(): Boolean = error("unexpected back")
override fun clickFresh(target: SnapshotNode): FreshActionResult = error("final confirmation must be read only")
override fun tapPurchaseFresh(target: SnapshotNode): FreshActionResult = error("unexpected tap")
override fun clickFresh(target: SnapshotNode): FreshActionResult = error("unexpected click")
override fun inputFresh(target: SnapshotNode, value: String): FreshActionResult = error("unexpected input")
override fun swipePurchase(direction: SwipeDirection, durationMs: Long): Boolean = error("unexpected swipe")
override fun swipePurchaseIn(target: SnapshotNode, direction: SwipeDirection, durationMs: Long): Boolean = error("unexpected scoped swipe")
override fun backPurchase(): Boolean = error("unexpected back")
}
return PurchaseLiveAutomation(driver, pause = { error("unexpected pause") }, panelDiagnostic = diagnostics::add)
.finalConfirmation(input, ShippingAddressProof(ADDRESS, "_cg376"))
.finalConfirmation(input, ShippingAddressProof(ADDRESS, "_cg374"))
}
private fun accepts(snapshot: UiSnapshot, input: PurchaseExecutionInput = input()) {
assertEquals(input.mappedColor, check(snapshot, input).mappedColor)
private fun accepts(snapshot: UiSnapshot, color: String = "黑色", size: String = "XL") {
val diagnostics = mutableListOf<String>()
val result = try { check(snapshot, input(color, size), diagnostics) } catch (failure: PurchaseLiveException) {
throw AssertionError("${failure.code}: ${diagnostics.lastOrNull()}", failure)
}
assertEquals(color, result.mappedColor)
}
private fun rejects(snapshot: UiSnapshot, expectedCode: String, input: PurchaseExecutionInput = input()) {
private fun rejects(snapshot: UiSnapshot, color: String = "黑色", size: String = "XL", expectedCode: String = "PURCHASE_SPEC_NOT_MATCHED", input: PurchaseExecutionInput = input(color, size)) {
val diagnostics = mutableListOf<String>()
try {
check(snapshot, input, diagnostics)
@@ -175,11 +228,14 @@ class PurchaseFinalSpecConfirmationTest {
} catch (failure: PurchaseLiveException) {
assertEquals(expectedCode, failure.code)
if (expectedCode == "PURCHASE_SPEC_NOT_MATCHED") {
assertEquals("finalSpec;stage=selection_prompt", diagnostics.last())
assertTrue(failure.message.orEmpty().contains(diagnostics.last()))
val detail = diagnostics.last()
assertTrue(detail, detail.matches(Regex("finalSpec;dimension=(color|size|none);stage=[a-z_]+;summaryPresent=[01];selectedFound=[01];conflict=[01]")))
assertTrue(failure.message.orEmpty(), failure.message.orEmpty().contains(detail))
assertFalse(failure.message.orEmpty().contains(ADDRESS))
assertFalse(failure.message.orEmpty().contains("黑色"))
assertFalse(failure.message.orEmpty().contains("XL"))
}
}
}
private companion object { const val ADDRESS = "示例区示例路376号_cg376" }
private companion object { const val ADDRESS = "示例区示例路374号_cg374" }
}
@@ -8,9 +8,6 @@ import cn.ilapage.goauto.agent.automation.PurchaseExecutionInput
import cn.ilapage.goauto.agent.automation.PurchaseLiveAutomation
import cn.ilapage.goauto.agent.automation.PurchaseLiveException
import cn.ilapage.goauto.agent.automation.PurchaseRuleParser
import cn.ilapage.goauto.agent.automation.PurchaseRehearsalExecutor
import cn.ilapage.goauto.agent.automation.PurchaseAgentCapabilities
import cn.ilapage.goauto.agent.automation.PurchaseExecutionOutcome
import cn.ilapage.goauto.agent.automation.PurchaseUiDriver
import cn.ilapage.goauto.agent.automation.SnapshotNode
import cn.ilapage.goauto.agent.automation.SwipeDirection
@@ -21,105 +18,6 @@ import org.junit.Assert.assertTrue
import org.junit.Test
class PurchaseLiveAutomationTest {
@Test
fun `pre address summary or button selection prompt stops executor before any address action`() {
for (variant in listOf("summary", "button", "second-dimension", "after-swipe")) {
val result = runPromptBoundary(prePrompt = variant)
assertEquals(variant, "failed", result.outcome.resultType)
assertEquals(variant, "PURCHASE_SPEC_SELECTION_UNCONFIRMED", result.outcome.errorCode)
assertTrue(result.outcome.message, result.outcome.message.contains("reason=selection_prompt"))
assertEquals(variant, 0, result.boundaries)
assertEquals(variant, 0, result.driver.addressPathClicks)
assertEquals(variant, 0, result.driver.inputCount)
assertEquals(variant, 0, result.driver.submitClicks)
assertEquals(variant, if (variant == "after-swipe") 1 else 0, result.preSwipes)
}
}
@Test
fun `final selection prompt is ordinary failure before boundary and unknown after boundary with zero clicks`() {
val before = runPromptBoundary(finalPrompt = true)
assertEquals("failed", before.outcome.resultType)
assertEquals("PURCHASE_SPEC_NOT_MATCHED", before.outcome.errorCode)
assertTrue(before.outcome.message.contains("finalSpec;stage=selection_prompt"))
assertEquals(0, before.boundaries)
assertEquals(0, before.driver.submitClicks)
assertEquals(1, before.driver.inputCount)
val after = runPromptBoundary(boundaryPrompt = true)
assertEquals("order_result_unknown", after.outcome.resultType)
assertEquals("PURCHASE_ORDER_RESULT_UNKNOWN", after.outcome.errorCode)
assertEquals(1, after.boundaries)
assertEquals(0, after.driver.submitClicks)
assertEquals(0, after.driver.scopedSwipes + after.driver.genericSwipes)
}
@Test
fun `post address wrong spec and quantity normal button complete executor once`() {
val result = runPromptBoundary(wrongFinalValues = true)
assertEquals("order_created", result.outcome.resultType)
assertEquals(1, result.boundaries)
assertEquals(1, result.driver.submitClicks)
}
private data class PromptBoundaryResult(val outcome: PurchaseExecutionOutcome, val driver: LiveDriver, val boundaries: Int, val preSwipes: Int)
private fun runPromptBoundary(
prePrompt: String? = null, finalPrompt: Boolean = false, boundaryPrompt: Boolean = false, wrongFinalValues: Boolean = false,
): PromptBoundaryResult {
val addressDriver = LiveDriver()
var step = ""
var boundaries = 0
var verifyCaptures = 0
var preSwipes = 0
val selectedSheet = SpecPanelFixtures.sheet(SpecPanelFixtures.Sheet(colorLabels = listOf("黑色"), sizeLabels = listOf("XL"))).let { sheet ->
sheet.copy(nodes = sheet.nodes.map { when {
it.path.endsWith("body/summary") -> it.copy(text = "已选 黑色 XL")
it.path.endsWith("qty/input") -> it.copy(text = "2")
it.path.endsWith("submit/l/t") -> it.copy(text = "提交订单")
it.path.endsWith("color/o0") || it.path.endsWith("size/o0/t") -> it.copy(selected = true)
else -> it
} })
}
val driver = object : PurchaseUiDriver by addressDriver {
override fun capture(): UiSnapshot {
if (step in listOf("updateShippingAddress", "createOrder", "readOrderResult")) {
val snapshot = addressDriver.capture()
if (step != "createOrder") return snapshot
return snapshot.copy(nodes = snapshot.nodes.map { node -> when {
node.path == "panel/submit" && (finalPrompt || (boundaryPrompt && boundaries > 0)) ->
node.copy(text = "选择颜色分类及尺码后,提交订单")
wrongFinalValues && node.path == "panel/selected" -> node.copy(text = "已选 白色 L")
wrongFinalValues && node.path == "panel/quantity" -> node.copy(text = "9")
else -> node
} })
}
if (step != "verifyOrderSummary") return selectedSheet
verifyCaptures++
return selectedSheet.copy(nodes = selectedSheet.nodes.map { node -> when {
prePrompt == "after-swipe" && node.path.contains("list/color") -> node.copy(visible = false)
prePrompt == "after-swipe" && preSwipes == 0 && node.path.endsWith("body/summary") -> node.copy(text = "已选 白色 XL")
(prePrompt == "summary" || (prePrompt == "after-swipe" && preSwipes > 0)) && node.path.endsWith("body/summary") -> node.copy(text = "请选择:颜色分类")
(prePrompt == "button" || (prePrompt == "second-dimension" && verifyCaptures > 1)) && node.path.endsWith("submit/l/t") ->
node.copy(text = "选择颜色分类及尺码后,提交订单")
else -> node
} })
}
override fun swipePurchaseIn(target: SnapshotNode, direction: SwipeDirection, durationMs: Long): Boolean {
if (step == "verifyOrderSummary") { preSwipes++; return true }
return addressDriver.swipePurchaseIn(target, direction, durationMs)
}
}
val rule = PurchaseRuleParser.parse("""{
"schemaVersion":1,"ruleType":"pddPurchase",
"requiredCapabilities":["purchase.live.v1","purchase.address-update.v1","purchase.order-create.v1"],
"actions":[{"type":"openProduct"},{"type":"verifyProduct"},{"type":"openSpecPanel"},{"type":"selectSpec"},{"type":"setQuantity"},{"type":"verifyUnitPrice"},{"type":"verifyOrderSummary"},{"type":"updateShippingAddress"},{"type":"createOrder"},{"type":"readOrderResult"}]
}""")
val outcome = PurchaseRehearsalExecutor(driver, { error("unexpected open") }, { null }, pause = {},
stepChanged = { step = it }, beforeOrderSubmit = { boundaries++ },
).execute(input().copy(reuseProbeProduct = true), rule, PurchaseAgentCapabilities.supported)
return PromptBoundaryResult(outcome, addressDriver, boundaries, preSwipes)
}
@Test
fun `hint delayed address update preserves fresh payment back budget for late order evidence`() {
val observationDriver = OrderObservationDriver { sample ->
@@ -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() }
}
}
}
}
@@ -723,31 +723,6 @@ class PurchaseRehearsalExecutorTest {
assertEquals(1, driver.clicked.count { it == "黑色" })
}
@Test
fun `visible unselected target without summary cannot reuse historical proof`() {
val driver = FakePurchaseDriver(sizes = listOf("L"), hideSelectedSummaryAfterQuantitySet = true, hideSizeSelectedState = true)
val outcome = executor(driver, { driver.browser = true; true }, { null }, pause = {})
.execute(input().copy(mappedSize = "L"), PurchaseRuleParser.parse(rule()), PurchaseAgentCapabilities.supported)
assertEquals("failed", outcome.resultType)
assertEquals("PURCHASE_SPEC_SELECTION_UNCONFIRMED", outcome.errorCode)
assertTrue(outcome.message, outcome.message.contains("reason=target_visible_unselected"))
assertTrue(outcome.message, outcome.message.contains("proofRecorded=true"))
}
@Test
fun `visible unselected normalized target blocks historical proof without relocation`() {
val target = "L【建议105-115斤】"
val driver = FakePurchaseDriver(
sizes = listOf(target), hideSelectedSummaryAfterQuantitySet = true,
finalSizesAfterQuantitySet = listOf("L[建议105-115斤]"),
)
val outcome = executor(driver, { driver.browser = true; true }, { null }, pause = {})
.execute(input().copy(mappedSize = target), PurchaseRuleParser.parse(rule()), PurchaseAgentCapabilities.supported)
assertEquals("PURCHASE_SPEC_SELECTION_UNCONFIRMED", outcome.errorCode)
assertTrue(outcome.message, outcome.message.contains("reason=target_visible_unselected"))
assertFalse(driver.swipeInPaths.contains("scroll"))
}
@Test
fun `final verification keeps current attempt proof when selected panel becomes unclassified`() {
val driver = FakePurchaseDriver(panelBecomesUnknownAfterSizeProof = true)
@@ -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))
}
}
@@ -502,7 +502,7 @@ class SpecPanelRecognitionTest {
val base = SpecPanelFixtures.sheet(Sheet(submit = "none"))
val snapshot = base.copy(
nodes = base.nodes + listOf(
node("r/real-submit", "提交订单", NodeBounds(0, 2135, 1080, 2200), clickable = true, parent = "r"),
node("r/real-submit", "选择颜色分类及尺码后,提交订单", NodeBounds(0, 2135, 1080, 2200), clickable = true, parent = "r"),
// Zero-*area* (right==left) and lower on screen: must never win despite
// having its own non-blank label and a clickable ancestor.
node("r/decoy", "¥0.0", NodeBounds(500, 2200, 500, 2216), clickable = true, parent = "r"),
@@ -512,7 +512,7 @@ class SpecPanelRecognitionTest {
PurchaseLiveAutomation(driver, pause = {}).submitOrderOnce()
assertEquals(listOf("提交订单"), driver.clicked)
assertEquals(listOf("选择颜色分类及尺码后,提交订单"), driver.clicked)
}
@Test
@@ -196,13 +196,7 @@ class TruncatedSpecCardTest {
var sizeSelected = false
val clicks = mutableListOf<String>()
override fun capture() = if (opened) snapshot.copy(nodes = snapshot.nodes.map { n ->
when {
n.path.contains("/size/o") -> n.copy(selected = sizeSelected && n.path.startsWith("r/sheet/body/list/size/o2"))
// #376: a completed fake selection must no longer expose PDD's incomplete-selection prompts.
sizeSelected && n.label.startsWith("请选择") -> n.copy(text = "已选 黑色示例长裤【有抽绳】 有口袋不起球 2XL建议130-150斤")
sizeSelected && n.label == "选择颜色分类及尺码后,提交订单" -> n.copy(text = "提交订单")
else -> n
}
if (n.path.contains("/size/o")) n.copy(selected = sizeSelected && n.path.startsWith("r/sheet/body/list/size/o2")) else n
}) else SpecPanelFixtures.productDetailPage()
override fun clickFresh(target: SnapshotNode): FreshActionResult {
clicks += target.label
@@ -248,12 +242,7 @@ class TruncatedSpecCardTest {
}
@Test fun `existing final confirmation target first ordering is documented not changed`() {
val snapshot = sheet(otherSelected = true).let { sheet -> sheet.copy(nodes = sheet.nodes.map { node -> when {
node.label.startsWith("请选择") -> node.copy(text = "已选 $full $size")
node.label == "选择颜色分类及尺码后,提交订单" -> node.copy(text = "提交订单")
else -> node
} }) }
val screen = parse(snapshot)
val screen = parse(sheet(otherSelected = true))
val executor = executor(Driver(sheet()))
val immediate = PurchaseRehearsalExecutor::class.java.declaredMethods.single { it.name == "isExactSpecSelected" }
immediate.isAccessible = true
@@ -8,10 +8,32 @@ import org.junit.Assert.assertTrue
import org.junit.Test
class AgentDiagnosticStoreMigrationTest {
@Test fun v1AndV5UpgradeToV6PreservesEveryExistingTableAndReopens() {
for (version in listOf(1,5)) {
val file=java.io.File.createTempFile("sync-v6-upgrade-",".db")
try {
val old: Map<String,List<List<String?>>>
DriverManager.getConnection("jdbc:sqlite:${file.absolutePath}").use { db ->
createOldDatabase(db,version); insertOldDiagnostic(db); old=snapshotTables(db)
migrate(db,version,6); assertOldDiagnosticPreserved(db)
old.forEach { (name,rows)-> assertEquals(rows,tableRows(db,name)) }
val store=SyncDiagnosticStore(JdbcSyncSql(db))
val id=store.start("first",1000,null,SyncContext())
store.finish(id,1100,100,SyncFailure("heartbeat","network",1100,100))
store.ack(store.snapshot("first",1200)!!)
}
DriverManager.getConnection("jdbc:sqlite:${file.absolutePath}").use { db ->
assertOldDiagnosticPreserved(db); old.forEach { (name,rows)-> assertEquals(rows,tableRows(db,name)) }
val snapshot=SyncDiagnosticStore(JdbcSyncSql(db)).snapshot("reopened",1300)!!
assertEquals(0,snapshot.report.getInt("failedRounds")); assertTrue(snapshot.report.getLong("snapshotSeq")>2)
}
} finally { check(file.delete()) }
}
}
@Test
fun freshV5DatabaseHasNullableColorDiscoveryOutcomeColumns() = withDatabase { db ->
db.createStatement().use { it.execute(AgentDiagnosticSchema.createTableSql) }
assertEquals(5, AgentDiagnosticSchema.VERSION)
assertEquals(6, AgentDiagnosticSchema.VERSION)
assertV5Columns(db)
insertOldDiagnostic(db)
assertOldDiagnosticPreserved(db)
@@ -0,0 +1,106 @@
package cn.ilapage.goauto.agent.persistence
import cn.ilapage.goauto.agent.network.AgentApiException
import cn.ilapage.goauto.agent.network.HeartbeatResult
import cn.ilapage.goauto.agent.service.AgentSyncCycle
import cn.ilapage.goauto.agent.service.recoverPurchaseBoundary
import java.sql.DriverManager
import java.util.concurrent.CancellationException
import org.json.JSONArray
import org.junit.Assert.*
import org.junit.Test
class HandledRecoveryDiagnosticTest {
@Test
fun handledBoundaryApiFailureIsPersistedButDoesNotChangeRecoveryOrClaimPolicy() {
for (status in listOf(500, 409)) {
for (pending in listOf(false, true)) {
DriverManager.getConnection("jdbc:sqlite::memory:").use { db ->
val sql = JdbcSyncSql(db)
AgentDiagnosticSchema.syncStatements.forEach { sql.execute(it) }
val store = SyncDiagnosticStore(sql)
var elapsed = 10L
val round =
SyncDiagnosticRecorder(store, "p", { 1000 }, { elapsed })
.begin(SyncContext())
val events = mutableListOf<String>()
var originalMarkerPresent = true
val replayBoundary: () -> Unit = {
events += "recover"
elapsed = 45
throw AgentApiException(
status,
"PRIVATE_CODE",
"private URL and payload",
false,
)
}
val cycle =
AgentSyncCycle(
heartbeat = {
round.step("heartbeat") {
events += "heartbeat"
HeartbeatResult(1, true, false, 15)
}
},
activeTaskId = { null },
executionActive = { false },
recover = {
round.step("recover") {
recoverPurchaseBoundary(round::handledFailure) {
replayBoundary()
originalMarkerPresent = false
}
}
},
flush = { round.step("flush") { events += "flush" } },
hasPending = { pending },
)
.run()
round.finish()
assertEquals(listOf("heartbeat", "recover", "flush"), events)
assertTrue(originalMarkerPresent)
assertNull(cycle.failureCode)
assertEquals(!pending, cycle.canClaim)
val row = sql.query("SELECT * FROM sync_round").single()
assertEquals("failure", row["status"])
assertEquals("recover", row["failure_step"])
assertEquals(
if (status == 500) "api_server" else "api_conflict",
row["failure_category"],
)
assertEquals("35", row["failure_duration_ms"])
assertEquals(status.toString(), row["http_status"])
assertNull(row["api_code"])
val steps = JSONArray(row["steps"])
assertEquals("failure", steps.getJSONObject(2).getString("status"))
assertEquals(35, steps.getJSONObject(2).getLong("durationMs"))
assertEquals("success", steps.getJSONObject(3).getString("status"))
assertFalse(row.toString().contains("private"))
assertEquals(1, store.snapshot("p", 1100)!!.report.getInt("failedRounds"))
}
}
}
}
@Test
fun boundaryPolicyStillRethrowsAuthenticationCancellationAndFatalErrors() {
for (error in
listOf(
AgentApiException(401, "OTHER", "private", false),
AgentApiException(409, "DEVICE_DISABLED", "private", false),
CancellationException(),
OutOfMemoryError(),
)) {
val thrown =
assertThrows(error.javaClass) {
recoverPurchaseBoundary({
fail("only handled nonauthentication API errors are observed")
}) {
throw error
}
}
assertSame(error, thrown)
}
}
}
@@ -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,225 @@
package cn.ilapage.goauto.agent.persistence
import java.sql.DriverManager
import java.util.concurrent.CancellationException
import org.junit.Assert.*
import org.junit.Test
class SyncDiagnosticRecorderTest {
@Test
fun collectionWithSameNumericIdCannotInheritPurchaseAttempt() {
val old = SyncContext(taskId = 7, attemptId = "old-purchase", phase = "purchase")
val collection =
syncTaskIdentity(7, "collection", true, old) {
error("must not read purchase fallback")
}
assertEquals(7L, collection.taskId)
assertNull(collection.attemptId)
assertNull(collection.phase)
val purchase = syncTaskIdentity(7, "purchase", true, old) { null }
assertEquals("old-purchase", purchase.attemptId)
val recovered = syncTaskIdentity(8, null, false, old) { "recovered" }
assertEquals("recovered", recovered.attemptId)
assertNull(recovered.phase)
}
@Test
fun failureOfEachPersistenceStageIsIsolatedAndNeverCreatesFakeSuccess() {
for (target in
listOf(
"open",
"INSERT INTO sync_round",
"UPDATE sync_round SET steps",
"UPDATE sync_round SET ended_at",
"SELECT sequence",
"DELETE FROM sync_round",
)) {
DriverManager.getConnection("jdbc:sqlite::memory:").use { db ->
val real = JdbcSyncSql(db)
AgentDiagnosticSchema.syncStatements.forEach { real.execute(it) }
val broken =
object : SyncSql {
override fun execute(sql: String, args: List<Any?>) {
if (sql.startsWith(target)) error("private secret")
real.execute(sql, args)
}
override fun query(
sql: String,
args: List<Any?>,
): List<Map<String, String?>> {
if (sql.startsWith(target)) error("private secret")
return real.query(sql, args)
}
override fun <T> transaction(block: () -> T): T {
if (target == "open") error("private secret")
return real.transaction(block)
}
}
var warnings = 0
val recorder =
SyncDiagnosticRecorder(SyncDiagnosticStore(broken), "p", { 1000 }, { 10 }) {
warnings++
}
var business = 0
val round = recorder.begin(SyncContext())
round.step("heartbeat") { business++ }
round.finish()
assertEquals(target, 1, business)
assertTrue(target, warnings > 0)
val rows = real.query("SELECT * FROM sync_round")
if (target != "DELETE FROM sync_round")
assertTrue(target, rows.none { it["status"] == "success" })
assertFalse(rows.toString().contains("private secret"))
}
}
}
@Test
fun readAckCleanupFailuresKeepUnacknowledgedFailuresAndCancellationIsNotSwallowed() {
DriverManager.getConnection("jdbc:sqlite::memory:").use { db ->
val real = JdbcSyncSql(db)
AgentDiagnosticSchema.syncStatements.forEach { real.execute(it) }
val store = SyncDiagnosticStore(real)
val id = store.start("p", 1000, null, SyncContext())
store.finish(id, 1100, 100, SyncFailure("flush", "exception", 1100, 100))
val snapshot = store.snapshot("p", 1200)!!
var failure: Exception = IllegalStateException("secret")
val broken =
object : SyncSql {
override fun execute(sql: String, args: List<Any?>) {
throw failure
}
override fun query(sql: String, args: List<Any?>): List<Map<String, String?>> {
throw failure
}
override fun <T> transaction(block: () -> T): T = real.transaction(block)
}
val recorder =
SyncDiagnosticRecorder(SyncDiagnosticStore(broken), "p", { 1300 }, { 10 })
recorder.ack(snapshot)
assertNull(recorder.snapshot())
assertEquals(1, store.snapshot("p", 1400)!!.report.getInt("failedRounds"))
failure = CancellationException()
assertThrows(CancellationException::class.java) { recorder.snapshot() }
assertThrows(CancellationException::class.java) { recorder.begin(SyncContext()) }
}
}
@Test
fun freshRecorderDoesNotReusePreviousElapsedClockAndAbortedRoundStaysRunning() {
DriverManager.getConnection("jdbc:sqlite::memory:").use { db ->
val sql = JdbcSyncSql(db)
AgentDiagnosticSchema.syncStatements.forEach { sql.execute(it) }
val store = SyncDiagnosticStore(sql)
SyncDiagnosticRecorder(store, "old", { 1000 }, { 9000 }).begin(SyncContext()).finish()
val round = SyncDiagnosticRecorder(store, "new", { 2000 }, { 5 }).begin(SyncContext())
round.abort()
round.finish()
val row = sql.query("SELECT * FROM sync_round ORDER BY id DESC LIMIT 1").single()
assertNull(row["start_gap_ms"])
assertEquals("running", row["status"])
}
}
@Test
fun wallClockJumpDoesNotChangeElapsedDurationAndRepeatedStepsHaveDistinctOrdinals() {
DriverManager.getConnection("jdbc:sqlite::memory:").use { db ->
AgentDiagnosticSchema.syncStatements.forEach {
db.createStatement().use { s -> s.execute(it) }
}
val sql = JdbcSyncSql(db)
var wall = 1000L
var elapsed = 10L
val recorder =
SyncDiagnosticRecorder(SyncDiagnosticStore(sql), "process", { wall }, { elapsed })
val round = recorder.begin(SyncContext())
round.step("heartbeat") {
wall = -500
elapsed = 30
}
round.step("heartbeat") { elapsed = 40 }
round.finish()
val row = sql.query("SELECT * FROM sync_round").single()
assertEquals("30", row["duration_ms"])
val steps = org.json.JSONArray(row["steps"])
assertEquals(3, steps.length())
assertEquals(1, steps.getJSONObject(1).getInt("ordinal"))
assertEquals(2, steps.getJSONObject(2).getInt("ordinal"))
assertEquals(0, steps.getJSONObject(1).getLong("startedOffsetMs"))
assertEquals(20, steps.getJSONObject(1).getLong("endedOffsetMs"))
assertEquals(20, steps.getJSONObject(2).getLong("startedOffsetMs"))
assertEquals(30, steps.getJSONObject(2).getLong("endedOffsetMs"))
elapsed = 110
recorder.begin(SyncContext()).finish()
assertEquals(
"100",
sql.query("SELECT start_gap_ms FROM sync_round ORDER BY id DESC LIMIT 1")
.single()["start_gap_ms"],
)
}
}
@Test
fun authenticationCodesAndUninstrumentedTimeoutAreSafelyClassified() {
assertEquals(
"api_auth",
syncFailure(
cn.ilapage.goauto.agent.network.AgentApiException(
409,
"DEVICE_DISABLED",
"private",
false,
),
"heartbeat",
1000,
10,
)
.category,
)
val failure =
syncFailure(
java.net.SocketTimeoutException("connect timed out url secret"),
"register",
1000,
10,
)
assertEquals("timeout_unknown", failure.category)
assertEquals("SocketTimeoutException", failure.exceptionClass)
assertNull(
syncFailure(IllegalArgumentException("private"), "flush", 1000, 10).exceptionClass
)
}
@Test
fun storageFailuresNeverBlockBusinessButCancellationAndFatalErrorsPropagate() {
val sql =
object : SyncSql {
override fun execute(sql: String, args: List<Any?>) {
throw IllegalStateException("private message")
}
override fun query(sql: String, args: List<Any?>): List<Map<String, String?>> {
throw IllegalStateException("private message")
}
override fun <T> transaction(block: () -> T): T = block()
}
val recorder = SyncDiagnosticRecorder(SyncDiagnosticStore(sql), "p", { 1000 }, { 10 })
var count = 0
val round = recorder.begin(SyncContext())
round.step("heartbeat") { count++ }
round.finish()
assertEquals(1, count)
assertNull(recorder.snapshot())
assertThrows(CancellationException::class.java) {
round.step("flush") { throw CancellationException() }
}
assertThrows(OutOfMemoryError::class.java) {
round.step("flush") { throw OutOfMemoryError() }
}
}
}
@@ -0,0 +1,168 @@
package cn.ilapage.goauto.agent.persistence
import java.sql.DriverManager
import org.junit.Assert.*
import org.junit.Test
class SyncDiagnosticStoreTest {
@Test
fun contextRejectsArbitraryTextInsteadOfPersistingIt() {
DriverManager.getConnection("jdbc:sqlite::memory:").use { db ->
val sql = JdbcSyncSql(db)
AgentDiagnosticSchema.syncStatements.forEach { sql.execute(it) }
SyncDiagnosticStore(sql)
.start(
"p",
1000,
null,
SyncContext(
deviceId = -1,
taskId = -1,
attemptId = "secret URL",
phase = "secret URL",
transport = "secret URL",
),
)
val row = sql.query("SELECT * FROM sync_round").single()
assertNull(row["device_id"])
assertNull(row["task_id"])
assertNull(row["attempt_id"])
assertNull(row["phase"])
assertEquals("unknown", row["transport"])
assertFalse(row.toString().contains("secret URL"))
}
}
@Test
fun retentionBoundsSyncHistoryOnlyAndSequenceDoesNotReset() {
DriverManager.getConnection("jdbc:sqlite::memory:").use { db ->
val sql = JdbcSyncSql(db)
AgentDiagnosticSchema.syncStatements.forEach { sql.execute(it) }
sql.execute(AgentDiagnosticSchema.createTableSql)
sql.execute(
"INSERT INTO agent_diagnostic(task_id,stage,reason,attempt,elapsed_ms,agent_version,created_at) VALUES(1,'x','x',0,0,'old',1)"
)
val store = SyncDiagnosticStore(sql)
val expired = store.start("p", 1, null, SyncContext())
store.finish(expired, 2, 1, SyncFailure("sync", "exception", 2, 1))
val before = store.snapshot("p", 10)!!.report.getLong("snapshotSeq")
val now = SyncDiagnosticStore.RETENTION_MS + 1000
store.start("p", 999, null, SyncContext())
store.start("p", 1000, null, SyncContext())
store.cleanup(now)
assertEquals("1", sql.query("SELECT COUNT(*) AS n FROM sync_round").single()["n"])
repeat(6001) { store.start("p", now, null, SyncContext()) }
store.cleanup(now)
assertEquals("6000", sql.query("SELECT COUNT(*) AS n FROM sync_round").single()["n"])
assertEquals("1", sql.query("SELECT COUNT(*) AS n FROM agent_diagnostic").single()["n"])
assertTrue(store.snapshot("p", now)!!.report.getLong("snapshotSeq") > before)
}
}
@Test
fun ackCannotEraseRunningRoundThatFailsAfterSnapshotAndRestartKeepsSequence() {
DriverManager.getConnection("jdbc:sqlite::memory:").use { db ->
AgentDiagnosticSchema.syncStatements.forEach { sql ->
db.createStatement().use { it.execute(sql) }
}
val sql = JdbcSyncSql(db)
val store = SyncDiagnosticStore(sql)
val first = store.start("p1", 1000, null, SyncContext())
store.finish(first, 1100, 100, SyncFailure("heartbeat", "network", 1100, 100))
val running = store.start("p1", 1200, 200, SyncContext())
val snapshot = store.snapshot("p1", 1250)!!
assertEquals(1, snapshot.report.getInt("failedRounds"))
store.finish(running, 1300, 100, SyncFailure("flush", "exception", 1300, 100))
store.ack(snapshot)
val restart = SyncDiagnosticStore(sql).snapshot("p2", 1400)!!
assertEquals(1, restart.report.getInt("failedRounds"))
assertEquals("flush", restart.report.getString("latestFailureStep"))
assertTrue(
restart.report.getLong("snapshotSeq") > snapshot.report.getLong("snapshotSeq")
)
assertEquals(1, store.snapshot("p2", 1450)!!.report.getInt("failedRounds"))
}
}
@Test
fun oldRunningIsInterruptedAndNeverInventsFailure() {
DriverManager.getConnection("jdbc:sqlite::memory:").use { db ->
AgentDiagnosticSchema.syncStatements.forEach { sql ->
db.createStatement().use { it.execute(sql) }
}
val store = SyncDiagnosticStore(JdbcSyncSql(db))
store.start("old", 1000, null, SyncContext())
val snapshot = store.snapshot("new", 1200)!!
assertTrue(snapshot.report.getBoolean("hasInterruptedRound"))
assertEquals(0, snapshot.report.getInt("failedRounds"))
assertTrue(snapshot.report.isNull("latestFailureAt"))
store.ack(snapshot)
assertTrue(store.snapshot("new", 1300)!!.report.getBoolean("hasInterruptedRound"))
store.cleanup(1001 + SyncDiagnosticStore.RETENTION_MS)
assertFalse(
store
.snapshot("new", 1001 + SyncDiagnosticStore.RETENTION_MS)!!
.report
.getBoolean("hasInterruptedRound")
)
}
}
@Test
fun v6AddsSyncHistoryWithoutChangingLegacyRetention() {
assertEquals(6, AgentDiagnosticSchema.VERSION)
assertEquals(50, AgentDiagnosticRetentionPolicy.MAX_RECORDS)
DriverManager.getConnection("jdbc:sqlite::memory:").use { db ->
db.createStatement().use { it.execute(AgentDiagnosticSchema.createTableSql) }
AgentDiagnosticSchema.failureSnapshotStatements.forEach { sql ->
db.createStatement().use { it.execute(sql) }
}
AgentDiagnosticSchema.migrationStatements(5, 6, emptySet()).forEach { sql ->
db.createStatement().use { it.execute(sql) }
}
db.createStatement().use { s ->
s.executeQuery("SELECT COUNT(*) FROM sync_round").use {
assertTrue(it.next())
assertEquals(0, it.getInt(1))
}
}
}
}
}
internal class JdbcSyncSql(private val db: java.sql.Connection) : SyncSql {
override fun execute(sql: String, args: List<Any?>) {
db.prepareStatement(sql).use { s ->
args.forEachIndexed { i, v -> s.setObject(i + 1, v) }
s.executeUpdate()
}
}
override fun query(sql: String, args: List<Any?>): List<Map<String, String?>> =
db.prepareStatement(sql).use { s ->
args.forEachIndexed { i, v -> s.setObject(i + 1, v) }
s.executeQuery().use { r ->
buildList {
while (r.next()) add(
(1..r.metaData.columnCount).associate {
r.metaData.getColumnLabel(it) to r.getString(it)
}
)
}
}
}
override fun <T> transaction(block: () -> T): T {
db.autoCommit = false
try {
val value = block()
db.commit()
return value
} catch (e: Throwable) {
db.rollback()
throw e
} finally {
db.autoCommit = true
}
}
}
@@ -0,0 +1,128 @@
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 {
@Test fun taskIdentityReadFailureRemainsObservableWithoutChangingCycleOrder() {
val error=IllegalStateException("private")
val result=AgentSyncCycle({ online }, { throw error }, { false }, {}, {}, { false }).run()
assertSame(error,result.failure)
assertFalse(result.canClaim)
}
@Test fun cancellationAndFatalErrorsAreNeverConvertedToBusinessFailures() {
for(error in listOf(java.util.concurrent.CancellationException(), OutOfMemoryError())) {
assertThrows(error.javaClass) {
AgentSyncCycle({ throw error }, { null }, { false }, {}, {}, { false }).run()
}
}
}
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))
}
}
@@ -0,0 +1,59 @@
package cn.ilapage.goauto.agent.service
import org.junit.Assert.*
import org.junit.Test
class PurchaseRuntimeReconciliationTest {
@Test fun `manual current page cannot enter while sync owns recovery and flush`() {
val mutex = TaskExecutionMutex()
val syncing = java.util.concurrent.atomic.AtomicBoolean(false)
val syncEntered = java.util.concurrent.CountDownLatch(1)
val finishSync = java.util.concurrent.CountDownLatch(1)
val worker = java.util.concurrent.Executors.newSingleThreadExecutor()
val future = worker.submit {
synchronized(mutex) { assertTrue(syncing.compareAndSet(false, true)) }
syncEntered.countDown()
check(finishSync.await(3, java.util.concurrent.TimeUnit.SECONDS))
syncing.set(false)
}
try {
assertTrue(syncEntered.await(3, java.util.concurrent.TimeUnit.SECONDS))
assertFalse(tryAcquireCurrentPage(mutex, syncing, Long.MAX_VALUE))
assertNull(mutex.currentTaskId())
finishSync.countDown()
future.get(3, java.util.concurrent.TimeUnit.SECONDS)
assertTrue(tryAcquireCurrentPage(mutex, syncing, Long.MAX_VALUE))
} finally { finishSync.countDown(); worker.shutdownNow() }
}
@Test fun `reboot restored rejected task clears runtime so repurchase becomes idle`() {
var runtimeTaskId: Long? = 42
val currentDatabaseAttemptRejected = true
if (shouldClearRejectedPurchase(42, runtimeTaskId, "purchase", null, currentDatabaseAttemptRejected)) runtimeTaskId = null
assertNull(runtimeTaskId)
assertTrue(runtimeTaskId == null) // final runtime predicate in repurchaseLocalIdle
}
@Test fun `another task collection active executor or newer attempt cannot be cleared`() {
assertFalse(shouldClearRejectedPurchase(42, 43, "purchase", null, true))
assertFalse(shouldClearRejectedPurchase(42, 42, "collection", null, true))
assertFalse(shouldClearRejectedPurchase(42, 42, "purchase", 42, true))
assertFalse(shouldClearRejectedPurchase(42, 42, "purchase", 43, true))
assertFalse(shouldClearRejectedPurchase(42, 42, "purchase", null, false))
}
@Test fun `stopped service remains stopped after an in flight heartbeat completes`() {
val gate = AgentStatePublicationGate()
var connection = "CONNECTING"
gate.publish { connection = "ONLINE" }
assertEquals("ONLINE", connection)
gate.stop { connection = "STOPPED" }
gate.publish { connection = "ONLINE" }
assertEquals("STOPPED", connection)
val restartedService = AgentStatePublicationGate()
restartedService.publish { connection = "CONNECTING" }
assertEquals("CONNECTING", connection)
gate.publish { connection = "ONLINE" }
assertEquals("CONNECTING", connection)
}
}
+27 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Architecture-and-Code-Map
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Architecture-and-Code-Map.-
wiki_revision: 76854d900b2372ddedd46a2b82b5677d8bf5cb31
synchronized_at: 2026-10-09T07:45:10Z
wiki_revision: 5508fe99e1040dff61d3bc6dbdf08ce0a1fd02ad
synchronized_at: 2026-10-10T03:10:31Z
<!-- gitea-wiki-mirror:end -->
<!-- gitea-wiki-mirror:start -->
@@ -673,3 +673,28 @@ Web 唯一展示位置为“采集采购 → SYB 同步记录”:列表状态
- `PddProductDetailCollector.collectColors` 只基于既有解析帧记录颜色阶段 UP 派发次数、空发现轮次和横向探索三态;不增加抓取、手势或等待。标签及签名只在内存使用,持久化仍为安全聚合。
- 不新增 Server/Web API 字段、服务端迁移或权限;普通结果 payload/missing 与采购探测失败策略未改变。完整遍历与不完整候选拦截仍属 #370 后续待证据实施范围,不能把本地诊断值作为已实现的业务门禁。
- 读取与三态限制见 Troubleshooting;旧版不保证降级打开 v5。JDBC 测试与实际 Android 安装升级分别验收,测试不代替真实竖向颜色采集验证。
## Admin 列表按当前筛选导出全部页(#372,待验收)
实现绑定 `a4a793f`(分支 `feat/372-list-export`,未合并、未发布)。
- 采购管理与 SYB 商品页各增「导出」按钮:`web/src/utils/list-export.js` 先按当前筛选读第 1 页取总数,超过 5000 条不导出,否则逐页读取(采购每页 100、SYB 每页 500)再用 `web/src/vendor/Export2Excel.js`(经 `utils/save-excel.js` 按需加载)生成 xlsx;任一页或批次失败不生成文件;筛选条件以点击时复制的参数为准。无服务端导出接口、无异步导出任务。
- `GET /api/admin/v1/purchase-tasks` 的 `AdminListRequest`(`server/app/goauto/purchase/admin_query.go`)新增可选 `orderSubmittedFrom` / `orderSubmittedTo`(`YYYY-MM-DD`,含首尾两天,按 `order_submitted_at`,UTC+8 与 SYB 创建时间筛选一致);格式错误或起止倒置返回 `INVALID_REQUEST`;无下单时间的任务在设置该筛选时不出现。采购页筛选栏用日期范围选择器。
- `GET /api/admin/v1/return-matches` 的 `returnmatch.ListItem` 新增只读 `yeekeQuantity`(`server/app/goauto/returnmatch/service.go`,来自既有 `yeeke_return_item` join,不参与匹配);该接口原本就支持一次传多个 `sybProductId`,SYB 导出按每批 200 个商品 ID 调用,不逐行请求,且只采用 `matched`/`confirmed` 的匹配。
- 无数据库迁移、无新接口、无新权限;导出按钮不单独鉴权,能打开列表的用户(含采购员)都可导出。
## 设备心跳诊断与逐设备离线事务(#378,实现基线)
Server 实现绑定 `d84728d`,Android 绑定 `dbe7e98`,分支 `feat/378-heartbeat-diagnostics`;main 合并记录见 #378;未执行线上迁移、发布或装机。以下是分支实现,不代表线上或现有手机已经更新。
- `server/app/goauto/device/heartbeat.go`:成功且非 requestId 重放的心跳在核心事务提交后、HTTP 返回前尽力保存历史;设备 offline→online、超时离线和已有停用入口在状态变化事务内保存事件。没有新增启用入口,身份恢复/重新注册不补造配对事件。
- `device/diagnostics.go`:严格校验报告、独立 200ms 写历史 context、按生命周期运行清理;历史失败不改变在线判断。清理与每 15 秒离线扫描分离。
- `agent_heartbeat_log` 保存 device_id、received_at、request_id、上报/服务端任务 ID、busy、agent_version、校验后的报告 JSON 和状态。索引 `ix_heartbeat_device_received(device_id,received_at)`、`ix_heartbeat_received(received_at)`。
- `agent_device_status_event` 保存 device_id、occurred_at、from_status/to_status、固定 reason、last_heartbeat_at、failed_task_count、order_result_unknown_count;索引 `ix_device_event_occurred(occurred_at)`。失败与订单结果未知分别计数。
- 离线候选列表在事务外读取;每台设备单独事务,设备主键及相关任务 NOWAIT 锁定当前读,复核心跳和状态后使用已经锁定的任务 ID/状态更新任务与 attempt、设备和事件。禁止用跨设备旧快照重选任务。
- MySQL 3572 只回滚并跳过当前设备,下一轮再试;不持有前面设备的锁等待后面设备,不本轮循环重试。其他错误保留错误语义及此前真实提交计数。事件日志只在提交后输出;没有全系统锁顺序重构。
- Android `AgentForegroundService` 在实际同步开始时创建记录,包装注册、心跳、恢复、补传与领任务;`SyncDiagnosticRecorder` 保存步骤序号、单调时钟起止偏移与总耗时。恢复函数内部已处理的非认证 API 错误也被观察,原业务吞错、占位与后续流程不变。
- `SyncDiagnosticStore` 在 `goauto_diagnostics.db` v6 追加 `sync_round`,并以单行 `sync_meta` 保存跨清理、重启的递增序号。既有采集诊断、三个采购失败快照表及其保留策略不变。同步表每轮尽力清理到最近24小时且最多6000轮。
- `HeartbeatReportNegotiation` 由服务实例持有,不因每轮重建 `AgentApiClient` 而丢失协商状态;只确认实际发出的已完成报告边界。HTTP请求阶段由 `AgentApiClient.request` 标注,不依靠异常消息区分连接/响应超时。
- 当前任务上下文只采用实际可核验的设备/任务ID、attempt UUID及固定phase;采集和采购即使数字ID相同也不复用旧采购上下文。foreground表示Android前台服务已启动,不表示Agent Activity当前占据屏幕。
+39 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Business-Rules-and-Glossary
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Business-Rules-and-Glossary.-
wiki_revision: f8bb75d1fc0b5e8822186c740ac2034fbd562057
synchronized_at: 2026-10-09T09:25:31Z
wiki_revision: 84e5a638893e7136b70bfdf7b0ce0a5186cd4ee4
synchronized_at: 2026-10-10T03:10:34Z
<!-- gitea-wiki-mirror:end -->
<!-- gitea-wiki-mirror:start -->
@@ -933,3 +933,40 @@ Android 0.9.64 / versionCode 77,源码 `6550b9f`(分支实现,尚未安装
- 不新增分享、OCR/VLM、滑动取证、自动重新选规格或付款;合成测试与 APK 构建不代表真实订单验证。安装、真实采购、合并和发布须按相应授权执行。
验证:19 个相关套件共418项通过,0失败/错误/跳过;Debug APK与release单测源码编译通过。执行器测试模拟提交边界回调,并验证本地JSON序列化;未执行真实SQLite/服务端边界端到端或真实采购验证。不宣称全量Android测试通过(既有图搜返回忙循环用例不在本次回归范围)。
## Android 采购结果明确拒收恢复(#377)
实现位于工单分支 fix/377-rejected-purchase-outbox(实现提交 8b75a95,运行占位与互斥补齐 309dce3);仅 Android,未合并 main、未安装或发布,真机恢复尚待独立授权验证。
- 每轮先发心跳,再在无任务执行锁时恢复中断任务及补传。心跳任务不一致不能挡住拒收收敛;收尾后只额外重发一次心跳。认证错误停止补传,只有最终心跳成功、没有未完成本地任务/待上传结果且收尾正常才领取新任务。不放宽服务端任务身份校验,不修改心跳间隔或离线阈值。
- 只有 HTTP 409 且错误码精确属于 PURCHASE_STATE_CONFLICT、PURCHASE_LEASE_EXPIRED、PURCHASE_RESULT_CONFLICT 才视为明确拒收。未知 409、网络异常、5xx 及认证错误不按拒收处理;缺少 retryable 不构成拒收证据。
- 本地拒收使用 rejected,不伪装 sent。按 outbox 与 task/attempt 归属事务保存,保留原 payload;不再自动补传该条,继续处理后续结果。拒收不触发上传成功或规格探测交接,交接凭证失效。拒收记录不再充当待上传占位;不改变服务端采购结果和原重购停止边界。
- order_created、order_result_unknown 或带订单号的拒收结果在采购页持续提示人工核对。“已核对”只持久隐藏提醒,不删除结果、不调用服务端、不允许再次下单;非订单拒收仍可在设置及任务详情查看。
- goauto_purchase.db v2→v3 只追加拒收时间、白名单错误码及人工核对时间三个可空列,旧结果与记录保持不变;v1 升级仍经过原 v2 追加列。旧版 SQLiteOpenHelper 不保证能打开 v3,不能靠卸载清数据回退,优先使用保留 v3 schema 的修复版本。
- 设置页区分同步中、真实心跳失败、已连接但有拒收;新说明采用固定文案,不展示任意服务端错误消息。本单只修复明确拒收后的持续阻塞,不代表最初心跳中断的原因已查明,历史诊断仍属于 #378。
- 拒收后只在执行锁空闲、当前运行标记属于同一采购任务且本地当前 attempt 已拒收时清理恢复占位,不清除其他采集/采购或新 attempt;手动当前页采集与同步共用既有 working/任务锁门禁。服务停止写 STOPPED,原服务实例的迟到状态更新不覆盖它。
## Admin 列表导出 Excel(#372,待验收)
实现绑定 `a4a793f`,分支 `feat/372-list-export`,未合并、未发布。
- 范围:采购管理和 SYB 商品页的「导出」导出**当前筛选条件下的全部页**,不是只导出当前页;筛选以点击导出时的值为准。总数超过 **5000** 条时提示缩小筛选范围并不导出;任一页读取失败不生成文件;无数据提示且不生成文件。分页读取期间数据变化可能有少量重复或遗漏,以导出时刻近似结果为准。
- 格式:时间 `YYYY-MM-DD HH:mm:ss`(浏览器本地时区);金额单位元,写成 Excel **数字单元格**(按分换算,例如 1230 分 → 12.3,可直接求和)并显示为 `0.00`;空值留空,0 写成 0;文件名 `采购导出_YYYYMMDD-HHmmss.xlsx`、`SYB商品导出_YYYYMMDD-HHmmss.xlsx`。不导出收货人、地址、账号等个人信息。
- 采购导出列(按顺序):SYB订单号、虾皮商品ID、颜色、尺码(均为 SYB 目标规格快照)、数量、状态(页面同款中文)、PDD单号、下单时间、价格(元)(数字单元格)。价格取 `pddOrderAmountCent / 100`,即 Agent 在 PDD 待付款页读到的应付金额;旧版 Agent 或人工补录前为空则留空。要导出成功采购,先把状态筛为「订单已创建」并按需选择下单日期。
- 采购管理新增「下单日期」筛选:含首尾两天,按 `order_submitted_at`,日期边界 UTC+8,与 SYB 商品页创建时间筛选一致;没有下单时间的任务在设置该筛选时不出现。
- SYB 商品导出列(按顺序):SYB订单、虾皮商品ID、**SYB售价(TWD)**、SYB颜色、SYB尺码、SYB数量、yeeke颜色、yeeke尺码、yeeke数量、匹配时间。yeeke 颜色/尺码来自该商品有效(`matched` 或 `confirmed`)退货匹配的 `variationName`,按第一个英文或中文逗号拆分,没有逗号时整段放 yeeke 颜色;yeeke 数量为对应退货商品数量;SYB售价(TWD)取商品行 `unitPriceCent / 100`,与页面「售价」列同源(SYB `detail/listByStock` 的 `productPrice`,单件成交单价,元/TWD),数字单元格,无值留空;没有有效匹配时 yeeke 列与匹配时间留空。要导出已用退货,先把处理阶段筛为「已用退货」。
- 权限:不新增按钮权限,能打开列表的用户(含采购员)都可导出。
## 设备离线诊断的事务边界(#378,实现基线)
Server实现 `d84728d`、Android实现 `dbe7e98`,main 合并记录见 #378;未发布或装机;本节不代表线上/手机行为已经更新。
默认心跳 15 秒、离线阈值 45 秒、离线扫描 15 秒不变。扫描逐设备独立事务:锁定最新设备/任务状态后,将设备离线、相关任务和 attempt 收敛、对应状态事件一起提交。锁争用时本设备不作任何部分修改,跳过后继续其他设备,下轮再检查;不自动重试采购或换机。
运行中采集和未提交订单采购仍转失败;已 `order_submit_started` 的采购仍转 `order_result_unknown`,人工核对、禁止自动重派。两个数量分开记录。诊断历史是尽力保存,不作为设备在线或任务业务状态的事实来源;事件失败回滚当前设备业务变化,其他已提交设备保持不变。
Android同步诊断不改变#377的心跳→恢复→补传→必要的一次任务不一致心跳→符合条件才领任务的顺序。诊断failure是本轮观察到的同步错误,不等同于业务任务失败;恢复内部暂缓处理的API错误也记录,但不改变原恢复占位、后续flush/claim规则,不增加采购重试或付款动作。
可恢复诊断存储异常不阻断同步,取消和致命错误不记成假成功。同步记录独立保留24小时/6000轮,不占采集50条额度;旧进行中记录只表示中断/不完整,不伪造原现场、失败原因或成功。报告确认不会清除发送期间新增失败;本地确认失败或响应丢失允许重复,报告不能用于改变设备在线或采购状态。
+64 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Troubleshooting
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Troubleshooting
wiki_revision: 6b370764320d1f5148fbc31f0ef6ce999370173d
synchronized_at: 2026-10-09T02:06:57Z
wiki_revision: d06fa26d85339641e3d67e5ee6314205173c24c5
synchronized_at: 2026-10-10T03:10:51Z
<!-- gitea-wiki-mirror:end -->
# 故障排查
@@ -158,3 +158,65 @@ adb -s <serial> shell run-as cn.ilapage.goauto.agent sqlite3 -readonly databases
```
排查须绑定同一设备、Agent/PDD 版本、任务与规则快照,并对照同次 App 布局。网页竖排与 App 横排不是同一复现场景;当前横排只能做横向兼容性回归。真机测试结果及待验证项以 #370 工单为准。
## 采购结果拒收后设备持续离线(#377)
修复绑定工单分支 fix/377-rejected-purchase-outbox 的 8b75a95、309dce3;尚未合并 main、安装或发布。自动化覆盖不等于真实设备验收,不能假定线上手机已具备本节恢复行为。
旧版同步在心跳之前补传,一条结果被服务端永久拒收就能让每轮跳过心跳及领任务。离线之后本地仍携带旧任务 ID,又可能被心跳的 DEVICE_TASK_MISMATCH 拒绝;仅把心跳前置不足以恢复。#377 先发心跳,仍进行有界拒收收尾,再在任务不一致时重发一次心跳;认证失败不会继续收尾,执行中的任务不会被恢复流程并发接管。
排查步骤:
1. 核实 APK 对应提交,不仅看版本名(本次未改 versionCode/versionName)。设置页“正在同步”不代表断线;真实心跳网络、服务端或认证失败使用固定原因说明;任务处理错误不再覆盖此前成功的连接状态。
2. “已连接 · 有 N 条采购结果被服务端拒收”表示本地已有 rejected 结果,不等于服务端接收成功。可在设置页“查看拒收结果”和对应采购详情查看;有订单证据的结果还在采购页顶部持续提醒。
3. 人工核对 PDD 实际订单及后台任务。需要补录订单时沿用 Admin 原有人工流程;Agent 的“已核对”只隐藏这条本机提醒,不能代替后台保存、恢复回填、解锁重试或完成任务。原结果/订单号与核对时间继续保存。
4. 未知 409、网络失败或 5xx 保持 pending;如果仍卡住,核对错误分类及原任务状态,不自动把 pending 改成 sent。不反复重试真实采购来试探是否已下单。
5. 早期现场曾经另获授权,通过同时修改 outbox 与本地任务 upload_status=sent 恢复设备,但这不代表服务端已收件,也会把订单号留在手机。它只是历史应急处置,不是新版运维步骤,不应复制到脚本或作为常规修复。
数据库升级仅限手机私有 goauto_purchase.db:v3 在 purchase_outbox 添加 rejected_at、rejection_error_code、acknowledged_at;保留原始 payload 和任务数据。不得把整库、订单/地址内容上传到工单、Wiki、普通日志或 SynapBus。旧 APK 不保证能够降级打开 v3;回退优先使用保留 v3 的兼容修复构建,不卸载清数据、不擅自降低库版本。
已验证 SQLite JDBC 执行同一生产迁移/拒收/核对 SQL、事务回滚及文件重开;尚未执行 Android SQLiteOpenHelper 仪器测试、真机界面或真实采购。最初心跳和结果同时中断的原因仍未查明,#378 的诊断计划不属于本修复。
## 心跳历史、离线事件与锁争用(#378,实现基线)
Server适用 `d84728d`、Android适用 `dbe7e98`,main 合并记录见 #378;未发布或装机;先确认目标实例/手机确实包含实现并完成对应迁移。旧版本没有下列表,不能直接执行这些查询。
经授权只读排查时用参数化查询指定设备和时间,不输出 Token、任务正文或其他个人数据:
```sql
SELECT device_id, received_at, request_id, current_task_id, running_task_id,
busy, agent_version, client_report_status
FROM agent_heartbeat_log
WHERE device_id = ? AND received_at >= ? AND received_at < ?
ORDER BY received_at, id;
SELECT device_id, occurred_at, from_status, to_status, reason,
last_heartbeat_at, failed_task_count, order_result_unknown_count
FROM agent_device_status_event
WHERE device_id = ? AND occurred_at >= ? AND occurred_at < ?
ORDER BY occurred_at, id;
```
- 心跳记录相邻间隔超过 45 秒只说明历史有缺口;200ms 历史写入超时或失败也会缺记录,不能单凭此宣称设备离线。
- 实际离线次数/时长以同设备 timeout/resumed 事件配对为准;仅有一端、跨时间窗、停用或身份恢复造成的不配对单列,不把缺少事件伪造为完整区间。
- `heartbeat_timeout` 中失败数不包含 `order_result_unknown_count`,后者表示订单已提交后的未知结果,不能自动重下单。
- 离线扫描遇到 MySQL 3572 只记固定锁争用分类及 deviceId 的 warn;它表示该设备当前被操作占用,未发生离线状态变更,等下一轮扫描。其他设备继续处理。
- 非 3572 数据库错误仍报告;事件写入失败会回滚该设备的状态与任务,已经成功处理的其他设备不回滚。不得通过删除任务或清空诊断来掩盖错误。
- 需要核对锁范围时对实际锁定 SQL 做 EXPLAIN;小型测试数据的执行计划不等于线上计划,不未经授权修改生产索引。
手机上经授权只读查询 `goauto_diagnostics.db` 的同步表,参数时间为epoch毫秒;不要为心跳排错公开整库,库内其他表可能含敏感失败现场:
```sql
SELECT id, process_id, started_at, ended_at, start_gap_ms, duration_ms,
device_id, task_id, attempt_id, phase, validated_network, transport,
foreground, status, failure_step, failure_category, failure_duration_ms,
http_status, api_code, exception_class, steps, completion_seq, reported
FROM sync_round
WHERE started_at >= ? AND started_at < ?
ORDER BY id;
```
- started_at/ended_at是墙钟对照;start_gap_ms、duration_ms、steps里的startedOffsetMs/endedOffsetMs来自单调时钟,不用墙钟相减解释耗时。新服务实例不延用上个实例的单调起点。
- process_id是诊断服务实例的随机标识,不是Android系统PID。
- running表示诊断没有完整结束,可能是卡住、进程退出、取消或诊断写入失败,不能仅凭它断言崩溃原因。旧实例running在保留期内一直触发hasInterruptedRound,报告确认不会把它改成成功/失败或隐藏。
- 一轮中观察到步骤失败(包括恢复内部暂缓处理的500/409),诊断可记failure,即使业务随后恢复连接;这不是采集/采购任务失败,也不能据此重下单。已正常收敛的明确拒收仍按#377处理。
- timeout_connect/timeout_response/timeout_unknown根据请求实际阶段分类,未知阶段不猜。诊断故障只输出sync_diagnostic_persistence_failed或sync_diagnostic_context_failed固定提示;不附错误正文、URL或SQL参数。
- reported只表示Agent已收到该报告请求的成功响应并在本地确认,不保证Server历史写入成功;响应丢失或确认失败可能重复补报。
+33 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Android-Agent-API-Contract
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Android-Agent-API-Contract.-
wiki_revision: 2102252a3f0ecb646b49941cc603423abb76ca2c
synchronized_at: 2026-10-08T07:11:01Z
wiki_revision: 1dc3c7550c0c4da4291e08b840170d861042e97a
synchronized_at: 2026-10-10T03:10:58Z
<!-- gitea-wiki-mirror:end -->
<!-- gitea-wiki-mirror:start -->
@@ -1530,3 +1530,34 @@ Server追加迁移 `1791400000000_purchase_failure_snapshot.go` 建立purchase_f
- 有ZIP的新上传只接受相符的非成功执行 attempt;不能向 pending/成功 attempt 新增失败诊断,已有有效载荷重放仍幂等。manifest 与 multipart metadata 语义相同;ZIP拒绝未知条目、重复路径、加密/未知压缩、CRC错误,XML拒绝DTD和处理指令,除限定XML声明。
- 本地队列绑定 Server Origin 和载荷 fingerprint,切换服务器不向新服务器转交旧ZIP;旧上传响应不删除后来首次恢复ZIP。每60秒维护,网络/408/429/5xx按60秒至1小时退避,410与其他永久4xx停止该队列项。结果Outbox优先,不因设备忙而永久饥饿;上传线程不持有任务互斥、不操作设备。
- 诊断接口不进入通用请求正文日志;私有SQL记录器不输出manifest/BLOB。管理员摘要不包含窗口和节点内容。默认HTTP例外只沿用既有Agent设置,不新增开关;传输私有原始内容时同样存在明文风险,部署须确认现有传输策略。
## 心跳诊断报告契约(#378,Server/Android 实现基线)
Server绑定 `d84728d`、Android绑定 `dbe7e98`;main 合并记录见 #378;未在线迁移/发布或装机。本节契约已在实现基线并作本地跨端验证,不代表现有手机已经补报。
心跳成功响应追加 `acceptsClientReport: true`;原 requestId 幂等及 currentTaskId 一致性校验不变。请求可以省略 `clientReport`(记录状态 none);显式 null、非法类型/字段/值只使报告为 client_report_invalid,历史 JSON 为 NULL,不让合法基础心跳失败。主 JSON 非法仍按原错误处理,整个请求上限仍为 64KiB。
报告 JSON 原始片段最多 2048 bytes,必须包含以下 8 个精确 camelCase 键(大小写别名不接受),不允许额外键:
| 字段 | 约束 |
|---|---|
| snapshotSeq | 整数 1..9007199254740991 |
| failedRounds | 整数 0..1000000 |
| latestFailureAt | null 或有效非零 RFC3339 时间 |
| latestFailureStep | null 或 sync/register/heartbeat/recover/flush/claim |
| latestFailureCategory | null 或 api_auth/api_conflict/api_server/api_other/timeout_connect/timeout_response/timeout_unknown/network/exception |
| latestFailureDurationMs | null 或整数 0..86400000 |
| previousSyncDurationMs | null 或整数 0..86400000 |
| hasInterruptedRound | 非 null 布尔 |
failedRounds 为 0 时四个 latestFailure* 必须全为 null;大于 0 时四个均须有效且非 null。前三个非 nullable 标量为 snapshotSeq、failedRounds、hasInterruptedRound。只存经校验并重新编码的白名单,不存 URL、异常消息、正文或页面内容。无有效正数快照序号时应省略报告。
服务端仅对成功非 requestId 重放的心跳写历史,核心事务提交后使用独立 200ms context,写失败不改变响应;acceptsClientReport 表示协议能力,不是历史持久化保证。
Android 已实现的协商边界:上次成功心跳明确 true 才在下一次携带报告;缺少标志/false、进程重启、地址或认证上下文变化撤销缓存。带报告请求遇旧版明确 HTTP422/INVALID_REQUEST/“请求 JSON 无效”时,可移除报告并沿用 requestId 最多重发一次基础心跳;401/403、任务冲突、网络或5xx不触发此兼容重发。基础回退成功不能确认报告。只确认实际发送快照及之前的本地失败,发送期间新失败保留;丢响应/确认失败允许重复,不代表服务端按 snapshotSeq 去重。
Agent在响应设备身份校验通过后才更新协商和确认报告。兼容回退响应不确认报告,也不直接授予补报能力,下轮基础心跳重新协商。HTTP连接/响应超时仍为10秒/15秒,不因诊断改变。
snapshotSeq来自持久递增序号;本地同时捕获“已完成记录”的序号边界,成功响应只确认这条边界以前的失败。当前正在运行的一轮即使在发送期间后来失败,也不会被较早快照清除。hasInterruptedRound表示保留期内是否存在旧服务实例的running记录,确认报告不清除该事实。
诊断数据库打开、写入、读取、确认或清理的可恢复异常只损失诊断证据,正常同步继续;不吞取消或fatal error,不将不完整持久化伪装为成功。报告和本地同步表均不保存异常message、URL、请求/响应正文或PDD页面内容。
+39 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Deployment-and-Operations
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Deployment-and-Operations.-
wiki_revision: 3f69f2e39675752dfd43dc38de41cb7defb74a42
synchronized_at: 2026-10-09T07:58:00Z
wiki_revision: b50250e0d61eff0b33f468e7cdf74139fe76bb65
synchronized_at: 2026-10-10T03:10:43Z
<!-- gitea-wiki-mirror:end -->
<!-- gitea-wiki-mirror:start -->
@@ -372,3 +372,40 @@ Provider 故障日志只允许记录调用关联 ID、操作类型、耗时、
- 集成验证25套件478项通过,0失败/错误/跳过;compileReleaseUnitTestKotlin与assembleDebug通过。不是全量单测结论;#374记录的未修改图搜返回测试忙循环未在本次重跑。未进行真机采购。
- 集成APK位于本机 `D:/OPC/goauto-worktrees/release-373-374/android/app/build/outputs/apk/debug/app-debug.apk`,版本仍为 `0.9.69-373-fix1 / code82`,**已同时包含#374,不能只靠版本文字区分前一#373单独包**。SHA256 `6e2d0e0ed1eb943b3f168bb16d628d2c2bb4ce3b7c7c8b9e567dcdf8259150f8`;apksigner校验通过,沿用原签名证书 `bf86d7465c6092be74ce9c4187eb30c9a7d045c3b91891a547e2598429766dec`。
- 本次仅合并、验证和推送,未覆盖安装、未执行真实采购/付款、未发布或重启服务;不需要业务库迁移。手机仍使用此前安装版本。后续装机须另行授权,保留原签名及用户数据。
## #372 Admin 列表导出发布(2026-10-10)
- 用户授权合并 main 并发布。#372 以 `--no-ff` 合并为 main `72a15f8`(同时包含此前已在 main 的 #377 Android 提交,本次不构建或安装 Android)。合并时 `docs/03` 镜像冲突,选用包含 #377 与 #372 两段的较新导出版本(Business-Rules-and-Glossary@1c57a86),`sync --check` 通过。
- 发布目录 `/home/goauto/releases/20261010-72a15f8-372`,上一目录 `/home/goauto/releases/20261009-e26743c-integrated` 保留供回滚。自上次发布 `e26743c` 以来 Server/Web 的源码变化只有 #372:采购列表新增可选下单日期参数、`return-matches` 列表新增只读 `yeekeQuantity`、两页导出按钮。没有数据库迁移、权限或菜单变更,未执行 migrate。
- 切换前复查:采集、采购及 attempt 都没有执行中的任务;有 1 个 pending 采购任务保持原状,没有取消、重置或发起业务。config 从上一发布逐字复制(diff 一致);static/temp/var 沿用同一真实目录;上一发布的 js/css 哈希资源以不覆盖方式保留,便于已打开页面加载旧分块。发布目录 0755、`www` 用户可读 index 已在切换前确认。
- Server 为 Linux amd64、CGO_ENABLED=0、go1.26.5 构建,SHA256 `ab3943a1261f438eaad4d4acf006910e8a996cf97afc18a399ad5bee6c8f3a1f`,运行进程二进制一致;Web 入口 index.html SHA256 `c3465960bd6f450d7b0d51994475bfdc126c7b6a47068d92bad9b4ac751d0529`。current 原子切换,goauto.service 重启后 active;Nginx 配置未变,未 reload。
- 按内容验收:公网 `/` 与 `/login` 返回同一 index(含 `id="app"`,无欢迎页),10 项入口 JS/CSS 200,`/api/v1/health` 200,验证码 code=200;未登录采购列表(带下单日期参数)与 return-matches 业务码 401;服务日志无 panic/fatal/1146/1054。
- 未验证:登录后采购管理、SYB 商品多页导出,以及采购员账号导出,待业务验收。两台采购手机在发布前(10-09 晚)已离线,与本次发布无关。
- 回滚:Server/Web 成套把 current 切回 `/home/goauto/releases/20261009-e26743c-integrated` 并重启;无迁移,不需要处理数据库。
## #372 v2 导出售价列的纯 Web 发布(2026-10-10)
- 用户授权合并 #372 v2 并更新线上。`feat/372-export-price` 以 `--no-ff` 合并为 main `4ad59bd`;合并后 Web 单元测试 93/93 通过,`build:prod` 通过。与上一发布 `72a15f8` 相比,Server 源码没有变化。
- 发布目录 `/home/goauto/releases/20261010-4ad59bd-372v2-web`,只替换 dist。goauto-server 从上一发布目录逐字复制(cmp 一致),config 一致,static/temp/var 沿用同一真实目录;上一版 js/css 哈希资源以不覆盖方式保留。index.html SHA256 `83c2dfa845c29073234bbdaaabac0690c46b51d32c85609aaf2743b854915a57`。
- 已确认 Nginx root 与运行进程 `GOAUTO_WEB_DIST` 都指向 `/home/goauto/current/dist`,所以按 #363 的纯 Web 方式只原子切换 current,不重启 GoAuto 或 Nginx;切换前后进程 PID 4786 不变,服务 active。
- 按内容验收:`/` 与 `/login` 返回新 index(含 `id="app"`,无欢迎页),10 项入口 JS/CSS 200,导出模块包含新表头「SYB售价(TWD)」,health 200,验证码 code=200,未登录业务接口 401。
- 未验证:登录后实际导出,以及 Excel 打开后售价与页面一致、能否求和,待业务验收。
- 回滚:把 current 原子切回 `/home/goauto/releases/20261010-72a15f8-372` 即可,不需要重启或处理数据库。
## 设备心跳诊断迁移与回退(#378,实现基线)
Server实现 `d84728d`、Android实现 `dbe7e98`,main 合并记录见 #378;未执行线上迁移、发布、重启或装机。以下为后续经授权部署时的边界,不是本次线上操作记录。
- 追加迁移 `1791600000000_agent_diagnostics.go` 已注册,新增 `agent_heartbeat_log`、`agent_device_status_event` 及索引;不改业务表结构。本轮并发修正不另加迁移或索引。
- 本地验证只能使用明确隔离且初始为空的 MySQL 8 库 `goauto_378_test`,凭据由当前进程注入 `GOAUTO_378_TEST_MYSQL_DSN`,执行 `go test ./app/goauto/device -run TestMySQLAgentDiagnostics -count=1 -v`(cwd server)。不得把线上或本地业务库 DSN 用于该测试;测试会创建并清理合成表。
- 现有迁移框架只有 up;本地删除新诊断表/版本记录后 re-up 的测试不等于存在生产 down 命令。MySQL DDL 隐式提交,不能承诺事务回滚整批 DDL。
- 常规运行回退停止诊断写入和能力广告,保留表及数据;不自动 DROP 表。删除生产数据或执行迁移仍需单独授权。
- API 生命周期中独立清理:心跳历史每小时清理严格超过 7 天的数据,状态事件每天清理严格超过30天的数据,每次至多5000条,剩余等下轮;不占离线15秒扫描,不修改Admin定时任务开关。
- 离线扫描使用 MySQL8 的 NOWAIT,争用只跳过本设备。测试执行计划仅是测试库证据,不代表线上数据分布;不据此擅自新增索引。
- Android诊断库v5→v6仅追加sync_round和单行sync_meta;v1旧库按既有升级链追加,保留原采集及采购失败快照。旧APK的SQLiteOpenHelper不保证能打开v6,回退应使用保留v6结构的修复构建,不卸载/清数据、不自动降库。
- 本地Android验证(cwd仓库根,JDK17及现有Android SDK):`./android/gradlew.bat -p android :app:testDebugUnitTest --tests '*HeartbeatReport*Test' --tests '*AgentTransportTimeoutTest' --tests '*SyncDiagnostic*Test' --tests '*HandledRecoveryDiagnosticTest' --tests '*AgentDiagnosticStoreMigrationTest' --tests '*AgentSyncCycleTest' :app:assembleDebug :app:compileReleaseUnitTestKotlin`。发布前还应执行#377受影响的outbox、运行占位、互斥、重购等回归;SQLite JDBC验证不等同于真机升级。
- Android兼容性测试在`android/build/diagnostics378-client-reports.json`生成纯合成HTTP捕获报告;在server目录将该绝对路径注入`GOAUTO_378_ANDROID_REPORTS`,执行`go test ./app/goauto/device -run TestAndroidGeneratedHeartbeatReports -count=1 -v`,使用真实validator/handler验证三类报告与重放。未设置路径时该项明确skip,不得据此宣称跨端通过;产物不提交Git。
- 线上迁移、发布、重启、装机仍需单独授权;本地构建/夹具不代表真机或线上30分钟心跳验收。
@@ -0,0 +1,89 @@
package device
import (
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"os"
"strings"
"testing"
"github.com/gin-gonic/gin"
"github.com/google/uuid"
)
// Feed reports captured by Android's real HTTP client into the real Server
// handler. This optional cross-build check uses synthetic test data only.
func TestAndroidGeneratedHeartbeatReports(t *testing.T) {
path := os.Getenv("GOAUTO_378_ANDROID_REPORTS")
if path == "" {
t.Skip("set GOAUTO_378_ANDROID_REPORTS to the Android test-generated JSON array")
}
data, err := os.ReadFile(path)
if err != nil {
t.Fatal(err)
}
var reports []json.RawMessage
if err := json.Unmarshal(data, &reports); err != nil {
t.Fatal("Android report fixture is not a JSON array")
}
if len(reports) < 3 {
t.Fatal("expected zero-failure, failed and interrupted reports")
}
db := openTestDatabase(t)
service := newTestService(t, db)
_, registered := registerHeartbeatDevice(t, service)
gin.SetMode(gin.TestMode)
router := gin.New()
router.POST("/heartbeat", Handler{DB: db}.Heartbeat)
var sawZero, sawFailure, sawInterrupted bool
for i, raw := range reports {
t.Run(fmt.Sprintf("android_report_%d", i), func(t *testing.T) {
canonical, status := validateHeartbeatClientReport(raw)
if status != "accepted" || canonical == nil {
t.Fatalf("Android-generated report rejected: classification=%s", status)
}
var report heartbeatClientReport
if err := json.Unmarshal(raw, &report); err != nil {
t.Fatal("report decode failed")
}
sawZero = sawZero || report.FailedRounds == 0
sawFailure = sawFailure || report.FailedRounds > 0
sawInterrupted = sawInterrupted || report.HasInterruptedRound
id := uuid.NewString()
body := fmt.Sprintf(`{"requestId":%q,"currentTaskId":null,"clientReport":%s}`, id, raw)
for replay := 0; replay < 2; replay++ {
req := httptest.NewRequest(http.MethodPost, "/heartbeat", strings.NewReader(body))
req.Header.Set("Authorization", "Bearer "+testDeviceToken)
res := httptest.NewRecorder()
router.ServeHTTP(res, req)
if res.Code != http.StatusOK {
t.Fatalf("heartbeat status=%d", res.Code)
}
var response struct {
Data struct {
AcceptsClientReport bool `json:"acceptsClientReport"`
} `json:"data"`
}
if json.Unmarshal(res.Body.Bytes(), &response) != nil || !response.Data.AcceptsClientReport {
t.Fatal("successful response did not advertise report capability")
}
}
var rows []struct {
DeviceID uint64
ClientReportStatus string
ClientReportJSON *string
}
if err := db.Table("agent_heartbeat_log").Where("request_id = ?", id).Find(&rows).Error; err != nil {
t.Fatal(err)
}
if len(rows) != 1 || rows[0].DeviceID != registered.DeviceID || rows[0].ClientReportStatus != "accepted" || rows[0].ClientReportJSON == nil || *rows[0].ClientReportJSON != *canonical {
t.Fatal("accepted Android report was not stored once with the correct device and canonical payload")
}
})
}
if !sawZero || !sawFailure || !sawInterrupted {
t.Fatal("fixture must cover zero-failure, failed and interrupted reports")
}
}
+165
View File
@@ -0,0 +1,165 @@
package device
import (
"bytes"
"context"
"encoding/json"
"errors"
"time"
"go-admin/app/goauto/models"
log "github.com/go-admin-team/go-admin-core/logger"
"gorm.io/gorm"
"gorm.io/gorm/logger"
)
const heartbeatHistoryBudget = 200 * time.Millisecond
type heartbeatClientReport struct {
SnapshotSeq int64 `json:"snapshotSeq"`
FailedRounds int64 `json:"failedRounds"`
LatestFailureAt *time.Time `json:"latestFailureAt"`
LatestFailureStep *string `json:"latestFailureStep"`
LatestFailureCategory *string `json:"latestFailureCategory"`
LatestFailureDurationMs *int64 `json:"latestFailureDurationMs"`
PreviousSyncDurationMs *int64 `json:"previousSyncDurationMs"`
HasInterruptedRound bool `json:"hasInterruptedRound"`
}
func validateHeartbeatClientReport(raw json.RawMessage) (*string, string) {
if len(raw) == 0 {
return nil, "none"
}
invalid := func() (*string, string) { return nil, "client_report_invalid" }
if len(raw) > 2048 {
return invalid()
}
var fields map[string]json.RawMessage
if json.Unmarshal(raw, &fields) != nil || len(fields) != 8 {
return invalid()
}
// encoding/json matches struct keys case-insensitively. Require every
// exact contract key before decoding so aliases cannot hide a missing key.
for _, name := range []string{"snapshotSeq", "failedRounds", "latestFailureAt", "latestFailureStep", "latestFailureCategory", "latestFailureDurationMs", "previousSyncDurationMs", "hasInterruptedRound"} {
if _, ok := fields[name]; !ok {
return invalid()
}
}
for _, name := range []string{"snapshotSeq", "failedRounds", "hasInterruptedRound"} {
if v, ok := fields[name]; !ok || bytes.Equal(bytes.TrimSpace(v), []byte("null")) {
return invalid()
}
}
var report heartbeatClientReport
decoder := json.NewDecoder(bytes.NewReader(raw))
decoder.DisallowUnknownFields()
if decoder.Decode(&report) != nil || report.SnapshotSeq < 1 || report.SnapshotSeq > 9007199254740991 || report.FailedRounds < 0 || report.FailedRounds > 1000000 {
return invalid()
}
for _, v := range []*int64{report.LatestFailureDurationMs, report.PreviousSyncDurationMs} {
if v != nil && (*v < 0 || *v > 86400000) {
return invalid()
}
}
if report.FailedRounds == 0 {
if report.LatestFailureAt != nil || report.LatestFailureStep != nil || report.LatestFailureCategory != nil || report.LatestFailureDurationMs != nil {
return invalid()
}
} else {
if report.LatestFailureAt == nil || report.LatestFailureAt.IsZero() || report.LatestFailureStep == nil || report.LatestFailureCategory == nil || report.LatestFailureDurationMs == nil {
return invalid()
}
if !oneOf(*report.LatestFailureStep, "sync", "register", "heartbeat", "recover", "flush", "claim") || !oneOf(*report.LatestFailureCategory, "api_auth", "api_conflict", "api_server", "api_other", "timeout_connect", "timeout_response", "timeout_unknown", "network", "exception") {
return invalid()
}
}
canonical, _ := json.Marshal(report)
value := string(canonical)
return &value, "accepted"
}
func oneOf(value string, allowed ...string) bool {
for _, v := range allowed {
if value == v {
return true
}
}
return false
}
func (s *Service) saveHeartbeatHistory(ctx context.Context, row models.AgentHeartbeatLog) {
// Client cancellation after the core commit must not suppress history. There
// is deliberately no goroutine: the caller waits at most this DB deadline.
writeCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), heartbeatHistoryBudget)
defer cancel()
err := s.DB.WithContext(writeCtx).Session(&gorm.Session{Logger: logger.Default.LogMode(logger.Silent), SkipDefaultTransaction: true}).Create(&row).Error
if err != nil {
classification := "database_error"
if errors.Is(err, context.DeadlineExceeded) || writeCtx.Err() != nil {
classification = "deadline_exceeded"
}
log.Warnf("agent heartbeat history skipped: device_id=%d classification=%s", row.DeviceID, classification)
}
}
func insertDeviceStatusEvent(tx *gorm.DB, event *models.AgentDeviceStatusEvent) error {
// Never issue a diagnostic SELECT or log the SQL error/body.
if err := tx.Session(&gorm.Session{Logger: logger.Default.LogMode(logger.Silent)}).Create(event).Error; err != nil {
return internalError(errors.New("device status event write failed"))
}
return nil
}
func logDeviceStatusEvent(event models.AgentDeviceStatusEvent) {
log.Infof("agent device status changed: device_id=%d reason=%s failed_task_count=%d order_result_unknown_count=%d", event.DeviceID, event.Reason, event.FailedTaskCount, event.OrderResultUnknownCount)
}
const diagnosticsCleanupBatch = 5000
// RunDiagnosticsCleanup shares the API server lifecycle but never occupies the
// offline monitor's 15-second scan loop.
func RunDiagnosticsCleanup(ctx context.Context, s *Service, onError func(error)) {
heartbeats := time.NewTicker(time.Hour)
defer heartbeats.Stop()
events := time.NewTicker(24 * time.Hour)
defer events.Stop()
runDiagnosticsCleanup(ctx, s, heartbeats.C, events.C, onError)
}
func runDiagnosticsCleanup(ctx context.Context, s *Service, heartbeats, events <-chan time.Time, onError func(error)) {
for {
select {
case <-ctx.Done():
return
case <-heartbeats:
if _, err := s.CleanupHeartbeatHistory(ctx); err != nil && ctx.Err() == nil && onError != nil {
onError(err)
}
case <-events:
if _, err := s.CleanupDeviceStatusEvents(ctx); err != nil && ctx.Err() == nil && onError != nil {
onError(err)
}
}
}
}
func (s *Service) CleanupHeartbeatHistory(ctx context.Context) (int64, error) {
return cleanupDiagnostics(s.DB.WithContext(ctx), &models.AgentHeartbeatLog{}, "received_at", s.Now().Add(-7*24*time.Hour))
}
func (s *Service) CleanupDeviceStatusEvents(ctx context.Context) (int64, error) {
return cleanupDiagnostics(s.DB.WithContext(ctx), &models.AgentDeviceStatusEvent{}, "occurred_at", s.Now().Add(-30*24*time.Hour))
}
func cleanupDiagnostics(db *gorm.DB, model any, timestamp string, cutoff time.Time) (int64, error) {
// A fixed-size ID batch works on both MySQL and SQLite. Recheck expiration
// in DELETE, and never loop through an unbounded backlog in one run.
var ids []uint64
if err := db.Model(model).Where(timestamp+" < ?", cutoff).Order(timestamp+", id").Limit(diagnosticsCleanupBatch).Pluck("id", &ids).Error; err != nil {
return 0, errors.New("diagnostics cleanup read failed")
}
if len(ids) == 0 {
return 0, nil
}
result := db.Where("id IN ? AND "+timestamp+" < ?", ids, cutoff).Delete(model)
if result.Error != nil {
return 0, errors.New("diagnostics cleanup delete failed")
}
return result.RowsAffected, nil
}
@@ -0,0 +1,372 @@
package device_test
import (
"context"
"errors"
"net"
"os"
"strings"
"sync"
"testing"
"time"
. "go-admin/app/goauto/device"
"go-admin/app/goauto/models"
versionlocal "go-admin/cmd/migrate/migration/version-local"
common "go-admin/common/models"
driver "github.com/go-sql-driver/mysql"
"github.com/google/uuid"
"gorm.io/driver/mysql"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"gorm.io/gorm/logger"
)
// Explicit opt-in only; refuses production/remote DSNs and any nonempty schema.
// Credentials are injected by the caller and never included in test errors.
func TestMySQLAgentDiagnostics(t *testing.T) {
raw := os.Getenv("GOAUTO_378_TEST_MYSQL_DSN")
if raw == "" {
t.Skip("isolated local MySQL DSN not provided")
}
cfg, err := driver.ParseDSN(raw)
if err != nil {
t.Fatal("invalid test DSN")
}
host, _, err := net.SplitHostPort(cfg.Addr)
if err != nil || cfg.Net != "tcp" || (host != "127.0.0.1" && host != "localhost" && host != "::1") || cfg.DBName != "goauto_378_test" {
t.Fatal("requires local tcp database goauto_378_test")
}
cfg.ParseTime = true
db, err := gorm.Open(mysql.Open(cfg.FormatDSN()), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
if err != nil {
t.Fatal("test MySQL connection failed")
}
sqlDB, err := db.DB()
if err != nil {
t.Fatal("test database unavailable")
}
t.Cleanup(func() { _ = sqlDB.Close() })
sqlDB.SetMaxOpenConns(8)
var version string
if db.Raw("SELECT VERSION()").Scan(&version).Error != nil || !strings.HasPrefix(version, "8.") {
t.Fatal("requires MySQL 8")
}
tables, err := db.Migrator().GetTables()
if err != nil || len(tables) != 0 {
t.Fatal("test schema must be empty")
}
owned := []any{&common.Migration{}, &models.AgentDevice{}, &models.PDDProduct{}, &models.CollectionRule{}, &models.CollectionTask{}, &models.PurchaseTask{}, &models.PurchaseTaskAttempt{}, &models.PDDAccount{}, &models.ShopeeProduct{}, &models.SYBProduct{}}
t.Cleanup(func() {
for _, model := range append([]any{&models.AgentHeartbeatLog{}, &models.AgentDeviceStatusEvent{}}, owned...) {
if err := db.Migrator().DropTable(model); err != nil {
t.Errorf("test table cleanup failed for %T", model)
}
}
})
if err := db.AutoMigrate(owned...); err != nil {
t.Fatal("test base migration failed")
}
const migrationVersion = "1791600000000_agent_diagnostics"
for i := 0; i < 2; i++ {
if versionlocal.MigrateAgentDiagnostics(db, migrationVersion) != nil {
t.Fatal("diagnostic migration failed")
}
}
for _, idx := range []struct {
model any
name string
}{{&models.AgentHeartbeatLog{}, "ix_heartbeat_device_received"}, {&models.AgentHeartbeatLog{}, "ix_heartbeat_received"}, {&models.AgentDeviceStatusEvent{}, "ix_device_event_occurred"}} {
if !db.Migrator().HasIndex(idx.model, idx.name) {
t.Fatalf("missing MySQL index %s", idx.name)
}
}
// Exercise reviewed rollback only on fresh diagnostic tables in this schema.
if db.Migrator().DropTable(&models.AgentHeartbeatLog{}, &models.AgentDeviceStatusEvent{}) != nil {
t.Fatal("diagnostic test rollback failed")
}
if db.Where("version = ?", migrationVersion).Delete(&common.Migration{}).Error != nil {
t.Fatal("test version rollback failed")
}
if versionlocal.MigrateAgentDiagnostics(db, migrationVersion) != nil {
t.Fatal("diagnostic reapply failed")
}
seed := func(t *testing.T) (*Service, RegisterResponse, string) {
t.Helper()
s := NewService(db)
s.Now = func() time.Time { return time.Date(2026, 10, 10, 0, 0, 0, 0, time.UTC) }
token := "isolated-mysql-" + uuid.NewString()
s.GenerateToken = func() (string, error) { return token, nil }
d, err := s.Register(context.Background(), RegisterRequest{RequestID: uuid.NewString(), InstallID: uuid.NewString(), Name: "test", Manufacturer: "test", Model: "test", AndroidVersion: "14", AgentVersion: "test", PDDVersion: "test"}, "")
if err != nil {
t.Fatal("synthetic device registration failed")
}
if db.Model(&models.AgentDevice{}).Where("id = ?", d.DeviceID).Update("last_heartbeat_at", s.Now().Add(-time.Minute)).Error != nil {
t.Fatal("test seed failed")
}
t.Cleanup(func() {
if db.Model(&models.AgentDevice{}).Where("id = ?", d.DeviceID).Update("status", models.DeviceStatusDisabled).Error != nil {
t.Error("synthetic device cleanup failed")
}
})
return s, d, token
}
t.Run("v4", func(t *testing.T) { runOfflineV4MySQLTests(t, db, seed) })
t.Run("scan_revalidates_real_heartbeat_after_candidate_read", func(t *testing.T) {
s, d, token := seed(t)
type scanMarker struct{}
callback := "v4_heartbeat_before_scan_lock"
if db.Callback().Query().Before("gorm:query").Register(callback, func(tx *gorm.DB) {
if tx.Statement.Context.Value(scanMarker{}) != true || tx.Statement.Table != "agent_device" {
return
}
if _, ok := tx.Statement.Clauses["FOR"]; !ok {
return
}
if _, err := s.Heartbeat(context.Background(), HeartbeatRequest{RequestID: uuid.NewString()}, token); err != nil {
tx.AddError(errors.New("concurrent heartbeat failed"))
}
}) != nil {
t.Fatal("register heartbeat barrier failed")
}
defer db.Callback().Query().Remove(callback)
n, err := s.MarkStaleDevicesOffline(context.WithValue(context.Background(), scanMarker{}, true), DefaultOfflineThreshold)
if err != nil || n != 0 {
t.Fatal("fresh heartbeat falsely marked offline")
}
requireOfflineDevice(t, db, d.DeviceID, models.DeviceStatusOnline, 0)
})
t.Run("concurrent_scans_and_resume_are_unique", func(t *testing.T) {
s, d, testDeviceToken := seed(t)
var wg sync.WaitGroup
errs := make(chan error, 2)
for i := 0; i < 2; i++ {
wg.Add(1)
go func() {
defer wg.Done()
_, err := s.MarkStaleDevicesOffline(context.Background(), DefaultOfflineThreshold)
errs <- err
}()
}
wg.Wait()
for i := 0; i < 2; i++ {
if err := <-errs; err != nil {
var dbErr *driver.MySQLError
if errors.As(err, &dbErr) {
t.Fatalf("concurrent scan failed mysql_errno=%d", dbErr.Number)
}
t.Fatalf("concurrent scan failed classification=%T", err)
}
}
req := HeartbeatRequest{RequestID: uuid.NewString()}
for i := 0; i < 2; i++ {
wg.Add(1)
go func() {
defer wg.Done()
_, err := s.Heartbeat(context.Background(), req, testDeviceToken)
errs <- err
}()
}
wg.Wait()
for i := 0; i < 2; i++ {
if <-errs != nil {
t.Fatal("concurrent heartbeat failed")
}
}
var events []models.AgentDeviceStatusEvent
db.Where("device_id = ?", d.DeviceID).Order("id").Find(&events)
if len(events) != 2 || events[0].Reason != "heartbeat_timeout" || events[1].Reason != "heartbeat_resumed" {
t.Fatalf("transition count=%d", len(events))
}
var n int64
db.Model(&models.AgentHeartbeatLog{}).Where("device_id = ?", d.DeviceID).Count(&n)
if n != 1 {
t.Fatalf("replay history count=%d", n)
}
})
t.Run("event_failure_rolls_back_real_mysql", func(t *testing.T) {
s, d, _ := seed(t)
product := models.PDDProduct{GoodsID: "37801", URL: "https://example.invalid/test"}
rule := models.CollectionRule{Name: "test", ContentJSON: `{}`}
if db.Create(&product).Error != nil || db.Create(&rule).Error != nil {
t.Fatal("synthetic task input failed")
}
task := models.CollectionTask{PDDProductID: &product.ID, RuleID: rule.ID, DeviceID: &d.DeviceID, Status: models.TaskStatusRunning, URLSnapshot: product.URL, GoodsIDSnapshot: product.GoodsID, RuleSnapshot: rule.ContentJSON}
if db.Create(&task).Error != nil {
t.Fatal("synthetic task creation failed")
}
if db.Exec("CREATE TRIGGER diagnostic_test_reject BEFORE INSERT ON agent_device_status_event FOR EACH ROW SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = 'synthetic_failure'").Error != nil {
t.Fatal("test trigger creation failed")
}
defer db.Exec("DROP TRIGGER IF EXISTS diagnostic_test_reject")
n, err := s.MarkStaleDevicesOffline(context.Background(), DefaultOfflineThreshold)
if err == nil || n != 0 {
t.Fatal("failed event did not abort")
}
var device models.AgentDevice
db.First(&device, d.DeviceID)
if device.Status != models.DeviceStatusOnline {
t.Fatal("state escaped rollback")
}
var storedTask models.CollectionTask
if db.First(&storedTask, task.ID).Error != nil || storedTask.Status != models.TaskStatusRunning || storedTask.DeviceRunSlot == nil {
t.Fatal("task escaped rollback")
}
if db.Exec("DROP TRIGGER diagnostic_test_reject").Error != nil {
t.Fatal("test trigger removal failed")
}
n, err = s.MarkStaleDevicesOffline(context.Background(), DefaultOfflineThreshold)
if err != nil || n != 1 {
t.Fatal("retry failed")
}
var event models.AgentDeviceStatusEvent
if db.Where("device_id = ?", d.DeviceID).First(&event).Error != nil || event.FailedTaskCount != 1 || event.OrderResultUnknownCount != 0 {
t.Fatal("per-device counts incorrect")
}
})
for _, domain := range []string{"collection", "purchase"} {
t.Run("scan_does_not_deadlock_task_then_device_reset_"+domain, func(t *testing.T) {
s, d, _ := seed(t)
goodsID := "37802"
if domain == "purchase" {
goodsID = "37803"
}
product := models.PDDProduct{GoodsID: goodsID, URL: "https://example.invalid/test"}
rule := models.CollectionRule{Name: "test", ContentJSON: `{}`}
if db.Create(&product).Error != nil || db.Create(&rule).Error != nil {
t.Fatal("synthetic inputs failed")
}
var taskID uint64
var taskModel any
if domain == "collection" {
task := models.CollectionTask{PDDProductID: &product.ID, RuleID: rule.ID, DeviceID: &d.DeviceID, Status: models.TaskStatusRunning, URLSnapshot: product.URL, GoodsIDSnapshot: product.GoodsID, RuleSnapshot: rule.ContentJSON}
if db.Create(&task).Error != nil {
t.Fatal("synthetic task failed")
}
taskID = task.ID
taskModel = &models.CollectionTask{}
} else {
task := models.PurchaseTask{PDDProductID: product.ID, DeviceID: &d.DeviceID, ExecutionMode: models.PurchaseExecutionModeRehearsal, Status: models.PurchaseTaskStatusRunning, PDDURLSnapshot: product.URL, PDDGoodsIDSnapshot: product.GoodsID, Quantity: 1, ReferenceUnitPriceCent: 100, MinUnitPriceCent: 20, MaxUnitPriceCent: 150, Currency: "CNY", RuleType: "pddPurchase", RuleSchemaVersion: 1, RequiredCapabilitiesJSON: `[]`, RuleSnapshot: `{}`, CreateRequestID: uuid.NewString()}
if db.Create(&task).Error != nil {
t.Fatal("synthetic purchase task failed")
}
taskID = task.ID
taskModel = &models.PurchaseTask{}
}
resetTx := db.Begin()
defer resetTx.Rollback()
if resetTx.Clauses(clause.Locking{Strength: "UPDATE"}).First(taskModel, taskID).Error != nil {
t.Fatal("reset task lock failed")
}
type scanMarker struct{}
deviceLocked, releaseScan := make(chan struct{}), make(chan struct{})
callback := "diagnostic_test_scan_lock_barrier"
if db.Callback().Query().After("gorm:query").Register(callback, func(tx *gorm.DB) {
if tx.Statement.Context.Value(scanMarker{}) == true && tx.Statement.Table == "agent_device" {
if _, ok := tx.Statement.Clauses["FOR"]; ok {
close(deviceLocked)
<-releaseScan
}
}
}) != nil {
t.Fatal("register scan barrier failed")
}
defer db.Callback().Query().Remove(callback)
scanDone, resetDone := make(chan error, 1), make(chan error, 1)
go func() {
_, err := s.MarkStaleDevicesOffline(context.WithValue(context.Background(), scanMarker{}, true), DefaultOfflineThreshold)
scanDone <- err
}()
select {
case <-deviceLocked:
case <-time.After(5 * time.Second):
close(releaseScan)
t.Fatal("scan device lock not reached")
}
go func() {
var device models.AgentDevice
err := resetTx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&device, d.DeviceID).Error
resetTx.Rollback()
resetDone <- err
}()
close(releaseScan)
scanErr, resetErr := <-scanDone, <-resetDone
if resetErr != nil {
var mysqlErr *driver.MySQLError
if errors.As(resetErr, &mysqlErr) {
t.Fatalf("reset lock failed mysql_errno=%d", mysqlErr.Number)
}
t.Fatal("reset lock failed")
}
if scanErr != nil {
t.Fatal("scan must skip busy task without returning contention")
}
var eventCount int64
db.Model(&models.AgentDeviceStatusEvent{}).Where("device_id = ?", d.DeviceID).Count(&eventCount)
if eventCount != 0 {
t.Fatal("busy task created event")
}
var device models.AgentDevice
db.First(&device, d.DeviceID)
if device.Status != models.DeviceStatusOnline {
t.Fatal("busy task changed status")
}
db.Callback().Query().Remove(callback)
n, err := s.MarkStaleDevicesOffline(context.Background(), DefaultOfflineThreshold)
if err != nil || n != 1 {
t.Fatal("scan did not retry after reset released locks")
}
})
}
t.Run("history_timeout_is_bounded_on_real_lock", func(t *testing.T) {
s, _, testDeviceToken := seed(t)
// Hold the diagnostic table's insert gap; liveness updates remain free.
lock := db.Begin()
defer lock.Rollback()
if lock.Exec("SELECT id FROM agent_heartbeat_log WHERE id > 0 FOR UPDATE").Error != nil {
t.Fatal("history lock failed")
}
started := time.Now()
res, err := s.Heartbeat(context.Background(), HeartbeatRequest{RequestID: uuid.NewString()}, testDeviceToken)
t.Logf("heartbeat with locked history table elapsed=%s", time.Since(started))
if err != nil || !res.Online || time.Since(started) > time.Second {
t.Fatal("history lock changed heartbeat outcome or exceeded budget")
}
})
t.Run("strict_cleanup", func(t *testing.T) {
s, d, _ := seed(t)
boundary := models.AgentHeartbeatLog{DeviceID: d.DeviceID, ReceivedAt: s.Now().Add(-7 * 24 * time.Hour), RequestID: uuid.NewString(), AgentVersion: "test", ClientReportStatus: "none"}
expired := boundary
expired.ReceivedAt = expired.ReceivedAt.Add(-time.Second)
db.Create(&boundary)
db.Create(&expired)
n, err := s.CleanupHeartbeatHistory(context.Background())
if err != nil || n != 1 {
t.Fatalf("cleanup=%d failed=%v", n, err != nil)
}
history := make([]models.AgentHeartbeatLog, 5001)
events := make([]models.AgentDeviceStatusEvent, 5001)
for i := range history {
history[i] = models.AgentHeartbeatLog{DeviceID: d.DeviceID, ReceivedAt: s.Now().Add(-31 * 24 * time.Hour), ClientReportStatus: "none"}
events[i] = models.AgentDeviceStatusEvent{DeviceID: d.DeviceID, OccurredAt: s.Now().Add(-31 * 24 * time.Hour)}
}
if db.CreateInBatches(&history, 500).Error != nil || db.CreateInBatches(&events, 500).Error != nil {
t.Fatal("synthetic cleanup batch seed failed")
}
n, err = s.CleanupHeartbeatHistory(context.Background())
if err != nil || n != 5000 {
t.Fatalf("history cleanup batch=%d failed=%v", n, err != nil)
}
n, err = s.CleanupDeviceStatusEvents(context.Background())
if err != nil || n != 5000 {
t.Fatalf("event cleanup batch=%d failed=%v", n, err != nil)
}
var remaining int64
db.Model(&models.AgentHeartbeatLog{}).Where("device_id = ?", d.DeviceID).Count(&remaining)
if remaining != 2 {
t.Fatalf("history boundary/remaining=%d", remaining)
}
})
}
@@ -0,0 +1,434 @@
package device
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"go-admin/app/goauto/models"
"github.com/gin-gonic/gin"
"github.com/google/uuid"
"gorm.io/gorm"
)
func TestHeartbeatDiagnosticsHistoryAndReportIsolation(t *testing.T) {
db := openTestDatabase(t)
service := newTestService(t, db)
_, registered := registerHeartbeatDevice(t, service)
gin.SetMode(gin.TestMode)
router := gin.New()
router.POST("/heartbeat", Handler{DB: db}.Heartbeat)
for _, tc := range []struct{ name, report, status string }{
{"missing", "", "none"},
{"accepted", `,"clientReport":{"snapshotSeq":1,"failedRounds":0,"latestFailureAt":null,"latestFailureStep":null,"latestFailureCategory":null,"latestFailureDurationMs":null,"previousSyncDurationMs":25,"hasInterruptedRound":false}`, "accepted"},
{"wrong_type", `,"clientReport":"secret-not-stored"`, "client_report_invalid"},
{"unknown_field", `,"clientReport":{"url":"secret-not-stored"}`, "client_report_invalid"},
} {
t.Run(tc.name, func(t *testing.T) {
id := uuid.NewString()
body := fmt.Sprintf(`{"requestId":%q,"currentTaskId":null%s}`, id, tc.report)
for i := 0; i < 2; i++ {
req := httptest.NewRequest(http.MethodPost, "/heartbeat", strings.NewReader(body))
req.Header.Set("Authorization", "Bearer "+testDeviceToken)
res := httptest.NewRecorder()
router.ServeHTTP(res, req)
if res.Code != 200 {
t.Fatalf("heartbeat status=%d body=%s", res.Code, res.Body.String())
}
var envelope struct {
Data map[string]any `json:"data"`
}
_ = json.Unmarshal(res.Body.Bytes(), &envelope)
if envelope.Data["acceptsClientReport"] != true {
t.Fatal("capability missing")
}
}
var rows []struct {
DeviceID uint64
ClientReportStatus string
ClientReportJSON *string
}
if err := db.Table("agent_heartbeat_log").Where("request_id = ?", id).Find(&rows).Error; err != nil {
t.Fatal(err)
}
if len(rows) != 1 || rows[0].DeviceID != registered.DeviceID || rows[0].ClientReportStatus != tc.status || (rows[0].ClientReportJSON != nil) != (tc.status == "accepted") {
t.Fatalf("history=%+v", rows)
}
})
}
}
func TestDisableWritesOneAtomicStatusEvent(t *testing.T) {
db := openTestDatabase(t)
s := newTestService(t, db)
_, d := registerHeartbeatDevice(t, s)
for i := 0; i < 2; i++ {
if err := s.Disable(context.Background(), d.DeviceID); err != nil {
t.Fatal(err)
}
}
var events []models.AgentDeviceStatusEvent
if err := db.Find(&events).Error; err != nil {
t.Fatal(err)
}
if len(events) != 1 || events[0].Reason != "disabled" || events[0].FromStatus != "online" || events[0].ToStatus != "disabled" {
t.Fatalf("events=%+v", events)
}
}
func TestStatusEventFailureRollsBackAndCanRetry(t *testing.T) {
for _, action := range []string{"timeout", "resumed", "disabled"} {
t.Run(action, func(t *testing.T) {
db := openTestDatabase(t)
s := newTestService(t, db)
_, d := registerHeartbeatDevice(t, s)
status := models.DeviceStatusOnline
if action == "resumed" {
status = models.DeviceStatusOffline
}
old := s.Now().Add(-time.Minute)
if err := db.Model(&models.AgentDevice{}).Where("id = ?", d.DeviceID).Updates(map[string]any{"status": status, "last_heartbeat_at": old}).Error; err != nil {
t.Fatal(err)
}
var task models.CollectionTask
if action == "timeout" || action == "disabled" {
product := models.PDDProduct{GoodsID: "1001", URL: "https://example.invalid/test"}
rule := models.CollectionRule{Name: "test", ContentJSON: `{}`}
if err := db.Create(&product).Error; err != nil {
t.Fatal(err)
}
if err := db.Create(&rule).Error; err != nil {
t.Fatal(err)
}
task = newHeartbeatTask(product, rule, d.DeviceID)
if err := db.Create(&task).Error; err != nil {
t.Fatal(err)
}
}
operation := func() error {
switch action {
case "timeout":
_, err := s.MarkStaleDevicesOffline(context.Background(), DefaultOfflineThreshold)
return err
case "disabled":
return s.Disable(context.Background(), d.DeviceID)
default:
_, err := s.Heartbeat(context.Background(), HeartbeatRequest{RequestID: uuid.NewString()}, testDeviceToken)
return err
}
}
callback := "diagnostic_event_failure"
if err := db.Callback().Create().Before("gorm:create").Register(callback, func(tx *gorm.DB) {
if tx.Statement.Table == "agent_device_status_event" {
tx.AddError(errors.New("injected private error"))
}
}); err != nil {
t.Fatal(err)
}
if err := operation(); err == nil {
t.Fatal("event failure must abort transaction")
}
var stored models.AgentDevice
if err := db.First(&stored, d.DeviceID).Error; err != nil {
t.Fatal(err)
}
if stored.Status != status || !stored.LastHeartbeatAt.Equal(old) {
t.Fatalf("state escaped rollback=%+v", stored.Status)
}
if task.ID != 0 {
var storedTask models.CollectionTask
if err := db.First(&storedTask, task.ID).Error; err != nil {
t.Fatal(err)
}
if storedTask.Status != models.TaskStatusRunning || storedTask.DeviceRunSlot == nil {
t.Fatal("task changes escaped rollback")
}
}
_ = db.Callback().Create().Remove(callback)
if err := operation(); err != nil {
t.Fatal(err)
}
var n int64
db.Model(&models.AgentDeviceStatusEvent{}).Count(&n)
if n != 1 {
t.Fatalf("events=%d", n)
}
})
}
}
func TestOfflineEventCountsPerDeviceIncludingUnknownBoundary(t *testing.T) {
db := openTestDatabase(t)
s := newTestService(t, db)
ids := make([]uint64, 3)
for i := range ids {
s.GenerateToken = func() (string, error) { return fmt.Sprintf("test-token-%d", i), nil }
_, registered := registerHeartbeatDevice(t, s)
ids[i] = registered.DeviceID
if err := db.Model(&models.AgentDevice{}).Where("id = ?", ids[i]).Update("last_heartbeat_at", s.Now().Add(-time.Minute)).Error; err != nil {
t.Fatal(err)
}
}
product := models.PDDProduct{GoodsID: "1002", URL: "https://example.invalid/test"}
rule := models.CollectionRule{Name: "test", ContentJSON: `{}`}
if err := db.Create(&product).Error; err != nil {
t.Fatal(err)
}
if err := db.Create(&rule).Error; err != nil {
t.Fatal(err)
}
collection := newHeartbeatTask(product, rule, ids[0])
if err := db.Create(&collection).Error; err != nil {
t.Fatal(err)
}
for i, status := range []string{models.PurchaseTaskStatusRunning, models.PurchaseTaskStatusOrderSubmitStarted} {
task := newHeartbeatPurchaseTask(product, ids[i], status)
if err := db.Create(&task).Error; err != nil {
t.Fatal(err)
}
attempt := newRunningPurchaseAttempt(task.ID, ids[i])
if err := db.Create(&attempt).Error; err != nil {
t.Fatal(err)
}
}
if n, err := s.MarkStaleDevicesOffline(context.Background(), DefaultOfflineThreshold); err != nil || n != 3 {
t.Fatalf("scan=%d %v", n, err)
}
var events []models.AgentDeviceStatusEvent
if err := db.Order("device_id").Find(&events).Error; err != nil {
t.Fatal(err)
}
if len(events) != 3 || events[0].FailedTaskCount != 2 || events[0].OrderResultUnknownCount != 0 || events[1].FailedTaskCount != 0 || events[1].OrderResultUnknownCount != 1 || events[2].FailedTaskCount != 0 || events[2].OrderResultUnknownCount != 0 {
t.Fatalf("events=%+v", events)
}
}
func TestHeartbeatHistoryOnlyForSuccessfulMainRequest(t *testing.T) {
db := openTestDatabase(t)
s := newTestService(t, db)
registerHeartbeatDevice(t, s)
if _, err := s.Heartbeat(context.Background(), HeartbeatRequest{RequestID: uuid.NewString()}, "wrong-token"); err == nil {
t.Fatal("bad token accepted")
}
if _, err := s.Heartbeat(context.Background(), HeartbeatRequest{RequestID: "invalid"}, testDeviceToken); err == nil {
t.Fatal("bad request accepted")
}
id := uint64(999)
if _, err := s.Heartbeat(context.Background(), HeartbeatRequest{RequestID: uuid.NewString(), CurrentTaskID: &id}, testDeviceToken); err == nil {
t.Fatal("task mismatch accepted")
}
var count int64
if err := db.Model(&models.AgentHeartbeatLog{}).Count(&count).Error; err != nil {
t.Fatal(err)
}
if count != 0 {
t.Fatal("failed request stored history")
}
}
func TestHeartbeatHistoryBoundedIndependentAndBestEffort(t *testing.T) {
db := openTestDatabase(t)
s := newTestService(t, db)
_, d := registerHeartbeatDevice(t, s)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
callback := "history_wait_for_deadline"
if err := db.Callback().Create().Before("gorm:create").Register(callback, func(tx *gorm.DB) {
if tx.Statement.Table != "agent_heartbeat_log" {
return
}
// We are already outside the core transaction; cancellation must not
// cancel the independent history context.
cancel()
if tx.Statement.Context.Err() != nil {
t.Error("history inherited request cancellation")
}
deadline, ok := tx.Statement.Context.Deadline()
if !ok || time.Until(deadline) > 201*time.Millisecond {
t.Error("missing bounded deadline")
}
var current models.AgentDevice
if err := db.First(&current, d.DeviceID).Error; err != nil || current.LastHeartbeatAt == nil {
t.Error("core was not committed before history")
}
<-tx.Statement.Context.Done()
tx.AddError(tx.Statement.Context.Err())
}); err != nil {
t.Fatal(err)
}
started := time.Now()
res, err := s.Heartbeat(ctx, HeartbeatRequest{RequestID: uuid.NewString()}, testDeviceToken)
if err != nil || !res.Online || time.Since(started) > time.Second {
t.Fatalf("best effort changed result: %v", err)
}
_ = db.Callback().Create().Remove(callback)
if err := db.Migrator().DropTable(&models.AgentHeartbeatLog{}); err != nil {
t.Fatal(err)
}
if _, err := s.Heartbeat(context.Background(), HeartbeatRequest{RequestID: uuid.NewString()}, testDeviceToken); err != nil {
t.Fatal(err)
}
}
func TestClientReportWhitelistBounds(t *testing.T) {
valid := `{"snapshotSeq":1,"failedRounds":1,"latestFailureAt":"2026-10-10T00:00:00Z","latestFailureStep":"heartbeat","latestFailureCategory":"network","latestFailureDurationMs":12,"previousSyncDurationMs":null,"hasInterruptedRound":false}`
for _, tc := range []struct{ name, raw, status string }{
{"accepted", valid, "accepted"}, {"none", "", "none"}, {"null", "null", "client_report_invalid"}, {"array", "[]", "client_report_invalid"},
{"sequence", strings.Replace(valid, `"snapshotSeq":1`, `"snapshotSeq":0`, 1), "client_report_invalid"},
{"duration", strings.Replace(valid, `"latestFailureDurationMs":12`, `"latestFailureDurationMs":86400001`, 1), "client_report_invalid"},
{"unknown_step", strings.Replace(valid, `"heartbeat"`, `"http://private"`, 1), "client_report_invalid"},
{"inconsistent", strings.Replace(valid, `"failedRounds":1`, `"failedRounds":0`, 1), "client_report_invalid"},
{"unknown", strings.Replace(valid, `"snapshotSeq"`, `"privateText"`, 1), "client_report_invalid"},
{"null_bool", strings.Replace(valid, `"hasInterruptedRound":false`, `"hasInterruptedRound":null`, 1), "client_report_invalid"},
{"case_alias_cannot_replace_missing_required", `{"snapshotSeq":1,"failedRounds":0,"latestFailureAt":null,"latestFailureStep":null,"latestFailureCategory":null,"SNAPSHOTSEQ":2,"previousSyncDurationMs":null,"hasInterruptedRound":false}`, "client_report_invalid"},
{"oversize", strings.Repeat(" ", 2049) + valid, "client_report_invalid"},
} {
t.Run(tc.name, func(t *testing.T) {
value, status := validateHeartbeatClientReport(json.RawMessage(tc.raw))
if status != tc.status || (status != "accepted" && value != nil) {
t.Fatalf("status=%s valuePresent=%v", status, value != nil)
}
})
}
}
func TestDiagnosticsCleanupIsBoundedAndStrictlyExpired(t *testing.T) {
db := openTestDatabase(t)
s := newTestService(t, db)
now := s.Now()
old := now.Add(-31 * 24 * time.Hour)
heartbeats := make([]models.AgentHeartbeatLog, 5001)
events := make([]models.AgentDeviceStatusEvent, 5001)
for i := range heartbeats {
heartbeats[i] = models.AgentHeartbeatLog{DeviceID: 1, ReceivedAt: old, RequestID: uuid.NewString(), AgentVersion: "test", ClientReportStatus: "none"}
events[i] = models.AgentDeviceStatusEvent{DeviceID: 1, OccurredAt: old, FromStatus: "online", ToStatus: "offline", Reason: "heartbeat_timeout"}
}
if err := db.CreateInBatches(&heartbeats, 100).Error; err != nil {
t.Fatal(err)
}
if err := db.CreateInBatches(&events, 100).Error; err != nil {
t.Fatal(err)
}
boundaryHeartbeat := models.AgentHeartbeatLog{DeviceID: 1, ReceivedAt: now.Add(-7 * 24 * time.Hour), RequestID: uuid.NewString(), AgentVersion: "test", ClientReportStatus: "none"}
boundaryEvent := models.AgentDeviceStatusEvent{DeviceID: 1, OccurredAt: now.Add(-30 * 24 * time.Hour), FromStatus: "online", ToStatus: "offline", Reason: "heartbeat_timeout"}
db.Create(&boundaryHeartbeat)
db.Create(&boundaryEvent)
n, err := s.CleanupHeartbeatHistory(context.Background())
if err != nil || n != 5000 {
t.Fatalf("heartbeat cleanup=%d %v", n, err)
}
n, err = s.CleanupDeviceStatusEvents(context.Background())
if err != nil || n != 5000 {
t.Fatalf("event cleanup=%d %v", n, err)
}
var count int64
db.Model(&models.AgentHeartbeatLog{}).Count(&count)
if count != 2 {
t.Fatalf("remaining=%d", count)
}
db.Model(&models.AgentDeviceStatusEvent{}).Count(&count)
if count != 2 {
t.Fatalf("remaining=%d", count)
}
n, err = s.CleanupHeartbeatHistory(context.Background())
if err != nil || n != 1 {
t.Fatalf("strict heartbeat cleanup=%d %v", n, err)
}
n, err = s.CleanupDeviceStatusEvents(context.Background())
if err != nil || n != 1 {
t.Fatalf("strict event cleanup=%d %v", n, err)
}
}
func TestDiagnosticsCleanupRunsOnSeparateSchedulesAndStops(t *testing.T) {
db := openTestDatabase(t)
s := newTestService(t, db)
old := s.Now().Add(-31 * 24 * time.Hour)
history := models.AgentHeartbeatLog{ReceivedAt: old, ClientReportStatus: "none"}
event := models.AgentDeviceStatusEvent{OccurredAt: old}
if err := db.Create(&history).Error; err != nil {
t.Fatal(err)
}
if err := db.Create(&event).Error; err != nil {
t.Fatal(err)
}
heartbeats, events := make(chan time.Time), make(chan time.Time)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
done := make(chan struct{})
go func() {
defer close(done)
runDiagnosticsCleanup(ctx, s, heartbeats, events, func(err error) { t.Error(err) })
}()
heartbeats <- s.Now()
// A second tick synchronizes completion of the first cleanup.
heartbeats <- s.Now()
var n int64
db.Model(&models.AgentHeartbeatLog{}).Count(&n)
if n != 0 {
t.Fatal("hourly history cleanup did not run")
}
db.Model(&models.AgentDeviceStatusEvent{}).Count(&n)
if n != 1 {
t.Fatal("hourly tick cleaned daily events")
}
events <- s.Now()
events <- s.Now()
cancel()
select {
case <-done:
case <-time.After(time.Second):
t.Fatal("cleanup did not stop")
}
db.Model(&models.AgentDeviceStatusEvent{}).Count(&n)
if n != 0 {
t.Fatal("daily event cleanup did not run")
}
}
func TestOfflineDiagnosticsAtomicAndResumed(t *testing.T) {
db := openTestDatabase(t)
s := newTestService(t, db)
_, d := registerHeartbeatDevice(t, s)
old := s.Now().Add(-time.Minute)
if err := db.Model(&models.AgentDevice{}).Where("id = ?", d.DeviceID).Update("last_heartbeat_at", old).Error; err != nil {
t.Fatal(err)
}
for i := 0; i < 2; i++ {
n, err := s.MarkStaleDevicesOffline(context.Background(), DefaultOfflineThreshold)
if err != nil || n != int64(1-i) {
t.Fatalf("scan=%d %v", n, err)
}
}
var events []struct {
Reason, FromStatus, ToStatus string
LastHeartbeatAt *time.Time
FailedTaskCount, OrderResultUnknownCount int64
}
if err := db.Table("agent_device_status_event").Find(&events).Error; err != nil {
t.Fatal(err)
}
if len(events) != 1 || events[0].Reason != "heartbeat_timeout" || events[0].FromStatus != "online" || events[0].ToStatus != "offline" || events[0].LastHeartbeatAt == nil || !events[0].LastHeartbeatAt.Equal(old) {
t.Fatalf("events=%+v", events)
}
req := HeartbeatRequest{RequestID: uuid.NewString()}
for i := 0; i < 2; i++ {
if _, err := s.Heartbeat(context.Background(), req, testDeviceToken); err != nil {
t.Fatal(err)
}
}
events = nil
if err := db.Table("agent_device_status_event").Order("id").Find(&events).Error; err != nil {
t.Fatal(err)
}
if len(events) != 2 || events[1].Reason != "heartbeat_resumed" {
t.Fatalf("events=%+v", events)
}
}
+166 -57
View File
@@ -2,14 +2,19 @@ package device
import (
"context"
"encoding/json"
"errors"
"fmt"
"time"
"go-admin/app/goauto/models"
log "github.com/go-admin-team/go-admin-core/logger"
mysqlDriver "github.com/go-sql-driver/mysql"
"github.com/google/uuid"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"gorm.io/gorm/logger"
)
const (
@@ -18,9 +23,10 @@ const (
)
type HeartbeatRequest struct {
RequestID string `json:"requestId"`
CurrentTaskID *uint64 `json:"currentTaskId"`
Capabilities []string `json:"capabilities,omitempty"`
RequestID string `json:"requestId"`
CurrentTaskID *uint64 `json:"currentTaskId"`
Capabilities []string `json:"capabilities,omitempty"`
ClientReport json.RawMessage `json:"clientReport,omitempty"`
}
type HeartbeatResponse struct {
@@ -32,6 +38,7 @@ type HeartbeatResponse struct {
ServerTime string `json:"serverTime"`
HeartbeatIntervalSeconds int `json:"heartbeatIntervalSeconds"`
Replayed bool `json:"replayed,omitempty"`
AcceptsClientReport bool `json:"acceptsClientReport"`
}
func (service *Service) Heartbeat(ctx context.Context, request HeartbeatRequest, presentedToken string) (HeartbeatResponse, error) {
@@ -50,10 +57,12 @@ func (service *Service) Heartbeat(ctx context.Context, request HeartbeatRequest,
interval = DefaultHeartbeatIntervalSeconds
}
var response HeartbeatResponse
var history models.AgentHeartbeatLog
var resumed *models.AgentDeviceStatusEvent
err = service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var device models.AgentDevice
digest := tokenDigest(presentedToken)
if err := tx.Where("token_digest = ? AND token_revoked_at IS NULL", digest).First(&device).Error; err != nil {
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("token_digest = ? AND token_revoked_at IS NULL", digest).First(&device).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return tokenInvalidError()
}
@@ -93,14 +102,30 @@ func (service *Service) Heartbeat(ctx context.Context, request HeartbeatRequest,
if result.RowsAffected == 0 {
return tokenInvalidError()
}
if device.Status == models.DeviceStatusOffline {
resumed = &models.AgentDeviceStatusEvent{DeviceID: device.ID, OccurredAt: now, FromStatus: device.Status, ToStatus: models.DeviceStatusOnline, Reason: "heartbeat_resumed", LastHeartbeatAt: device.LastHeartbeatAt}
if err := insertDeviceStatusEvent(tx, resumed); err != nil {
return err
}
}
}
response = HeartbeatResponse{
RequestID: request.RequestID, DeviceID: device.ID, CurrentTaskID: runningTaskID,
Online: true, Busy: runningTaskID != nil, ServerTime: now.UTC().Format(time.RFC3339Nano),
HeartbeatIntervalSeconds: interval, Replayed: replayed,
HeartbeatIntervalSeconds: interval, Replayed: replayed, AcceptsClientReport: true,
}
history = models.AgentHeartbeatLog{DeviceID: device.ID, ReceivedAt: now, RequestID: request.RequestID, CurrentTaskID: request.CurrentTaskID, RunningTaskID: runningTaskID, Busy: runningTaskID != nil, AgentVersion: device.AgentVersion}
return nil
})
if err == nil {
if resumed != nil {
logDeviceStatusEvent(*resumed)
}
if !response.Replayed {
history.ClientReportJSON, history.ClientReportStatus = validateHeartbeatClientReport(request.ClientReport)
service.saveHeartbeatHistory(ctx, history)
}
}
return response, err
}
@@ -115,8 +140,8 @@ func tokenInvalidError() error {
return &ServiceError{Code: CodeTokenInvalid, Message: "设备 Token 无效", Retryable: false}
}
// MarkStaleDevicesOffline atomically marks stale devices offline and fails any
// running task held by those devices. It returns the number of devices changed.
// MarkStaleDevicesOffline commits each device, its running tasks and its event
// independently. The count always describes committed device transitions.
func (service *Service) MarkStaleDevicesOffline(ctx context.Context, threshold time.Duration) (int64, error) {
if threshold <= 0 {
return 0, fmt.Errorf("offline threshold must be positive")
@@ -124,69 +149,152 @@ func (service *Service) MarkStaleDevicesOffline(ctx context.Context, threshold t
now := service.Now()
cutoff := now.Add(-threshold)
var changed int64
err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var deviceIDs []uint64
if err := tx.Model(&models.AgentDevice{}).
Where("status = ? AND COALESCE(last_heartbeat_at, created_at) < ?", models.DeviceStatusOnline, cutoff).
Pluck("id", &deviceIDs).Error; err != nil {
return internalError(err)
var deviceIDs []uint64
if err := service.DB.WithContext(ctx).Model(&models.AgentDevice{}).
Where("status = ? AND COALESCE(last_heartbeat_at, created_at) < ?", models.DeviceStatusOnline, cutoff).
Order("id").Pluck("id", &deviceIDs).Error; err != nil {
return 0, internalError(err)
}
if len(deviceIDs) == 0 {
return 0, nil
}
for _, deviceID := range deviceIDs {
if err := ctx.Err(); err != nil {
return changed, err
}
if len(deviceIDs) == 0 {
return nil
var event *models.AgentDeviceStatusEvent
// Expected NOWAIT contention must not be printed by GORM with SQL or
// parameters. The fixed warning below is the only contention log.
err := service.DB.WithContext(ctx).Session(&gorm.Session{Logger: logger.Default.LogMode(logger.Silent)}).Transaction(func(tx *gorm.DB) error {
var device models.AgentDevice
// The unlocked candidate list is only a hint. Lock the exact primary
// key without waiting, then revalidate its current state and heartbeat.
if err := tx.Clauses(clause.Locking{Strength: "UPDATE", Options: "NOWAIT"}).First(&device, deviceID).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return nil
}
return internalError(err)
}
last := device.CreatedAt
if device.LastHeartbeatAt != nil {
last = *device.LastHeartbeatAt
}
if device.Status != models.DeviceStatusOnline || !last.Before(cutoff) {
return nil
}
// Existing Admin reset paths lock task -> device, while collection
// Claim locks device -> task. Never wait for a task holding this device:
// contention rolls back only this device; a later scan retries it.
tasks, err := lockOfflineTasksNowait(tx, device.ID)
if err != nil {
return err
}
var collectionFailed int64
if len(tasks.CollectionIDs) > 0 {
result := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.CollectionTask{}).
Where("id IN ? AND status = ?", tasks.CollectionIDs, models.TaskStatusRunning).
Updates(map[string]any{
"status": models.TaskStatusFailed, "active_slot": gorm.Expr("NULL"), "device_run_slot": gorm.Expr("NULL"),
"lease_expires_at": gorm.Expr("NULL"), "error_code": "DEVICE_OFFLINE",
"error_message": "设备心跳超时,任务失败且不会自动重试", "finished_at": now,
})
if result.Error != nil {
return internalError(result.Error)
}
collectionFailed = result.RowsAffected
}
failed, unknown, err := markOfflinePurchaseTasks(tx, tasks, now)
if err != nil {
return err
}
event = &models.AgentDeviceStatusEvent{DeviceID: device.ID, OccurredAt: now, FromStatus: device.Status, ToStatus: models.DeviceStatusOffline, Reason: "heartbeat_timeout", LastHeartbeatAt: device.LastHeartbeatAt, FailedTaskCount: collectionFailed + failed, OrderResultUnknownCount: unknown}
result := tx.Model(&models.AgentDevice{}).Where("id = ? AND status = ?", device.ID, models.DeviceStatusOnline).
Update("status", models.DeviceStatusOffline)
if result.Error != nil {
return internalError(result.Error)
}
return insertDeviceStatusEvent(tx, event)
})
if err != nil {
if ctxErr := ctx.Err(); ctxErr != nil {
return changed, ctxErr
}
var mysqlErr *mysqlDriver.MySQLError
if errors.As(err, &mysqlErr) && mysqlErr.Number == 3572 {
log.Warnf("agent offline scan skipped: device_id=%d classification=lock_busy", deviceID)
continue
}
return changed, err
}
if err := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.CollectionTask{}).
Where("device_id IN ? AND status = ?", deviceIDs, models.TaskStatusRunning).
Updates(map[string]any{
"status": models.TaskStatusFailed, "active_slot": gorm.Expr("NULL"), "device_run_slot": gorm.Expr("NULL"),
"lease_expires_at": gorm.Expr("NULL"), "error_code": "DEVICE_OFFLINE",
"error_message": "设备心跳超时,任务失败且不会自动重试", "finished_at": now,
}).Error; err != nil {
return internalError(err)
if event != nil {
changed++
logDeviceStatusEvent(*event)
}
if err := markOfflinePurchaseTasks(tx, deviceIDs, now); err != nil {
return err
}
result := tx.Model(&models.AgentDevice{}).Where("id IN ? AND status = ?", deviceIDs, models.DeviceStatusOnline).
Update("status", models.DeviceStatusOffline)
if result.Error != nil {
return internalError(result.Error)
}
changed = result.RowsAffected
return nil
})
return changed, err
}
return changed, nil
}
func markOfflinePurchaseTasks(tx *gorm.DB, deviceIDs []uint64, now time.Time) error {
if err := markOfflinePurchaseStatus(
tx, deviceIDs, models.PurchaseTaskStatusRunning, models.PurchaseTaskStatusFailed,
"failed", "设备心跳超时,采购任务失败且不会自动重试", now,
); err != nil {
return err
type offlineTasks struct {
CollectionIDs []uint64
PurchaseRunningIDs []uint64
PurchaseSubmitStartedIDs []uint64
}
func lockOfflineTasksNowait(tx *gorm.DB, deviceID uint64) (offlineTasks, error) {
var tasks offlineTasks
locking := clause.Locking{Strength: "UPDATE", Options: "NOWAIT"}
if err := tx.Model(&models.CollectionTask{}).Clauses(locking).
Where("device_id = ? AND status = ?", deviceID, models.TaskStatusRunning).
Order("id").Pluck("id", &tasks.CollectionIDs).Error; err != nil {
return tasks, internalError(err)
}
return markOfflinePurchaseStatus(
tx, deviceIDs, models.PurchaseTaskStatusOrderSubmitStarted, models.PurchaseTaskStatusOrderResultUnknown,
var purchases []struct {
ID uint64
Status string
}
if err := tx.Model(&models.PurchaseTask{}).Select("id, status").Clauses(locking).
Where("device_id = ? AND status IN ?", deviceID, []string{models.PurchaseTaskStatusRunning, models.PurchaseTaskStatusOrderSubmitStarted}).
Order("id").Find(&purchases).Error; err != nil {
return tasks, internalError(err)
}
for _, task := range purchases {
if task.Status == models.PurchaseTaskStatusRunning {
tasks.PurchaseRunningIDs = append(tasks.PurchaseRunningIDs, task.ID)
} else {
tasks.PurchaseSubmitStartedIDs = append(tasks.PurchaseSubmitStartedIDs, task.ID)
}
}
return tasks, nil
}
func markOfflinePurchaseTasks(tx *gorm.DB, tasks offlineTasks, now time.Time) (int64, int64, error) {
failed, err := markOfflinePurchaseStatus(
tx, tasks.PurchaseRunningIDs, models.PurchaseTaskStatusRunning, models.PurchaseTaskStatusFailed,
"failed", "设备心跳超时,采购任务失败且不会自动重试", now,
)
if err != nil {
return 0, 0, err
}
unknown, err := markOfflinePurchaseStatus(
tx, tasks.PurchaseSubmitStartedIDs, models.PurchaseTaskStatusOrderSubmitStarted, models.PurchaseTaskStatusOrderResultUnknown,
"order_result_unknown", "设备心跳超时,订单结果未知,请人工核对,禁止自动重试", now,
)
return failed, unknown, err
}
func markOfflinePurchaseStatus(
tx *gorm.DB,
deviceIDs []uint64,
taskIDs []uint64,
fromStatus string,
toStatus string,
resultType string,
errorMessage string,
now time.Time,
) error {
var taskIDs []uint64
if err := tx.Model(&models.PurchaseTask{}).
Where("device_id IN ? AND status = ?", deviceIDs, fromStatus).
Pluck("id", &taskIDs).Error; err != nil {
return internalError(err)
}
) (int64, error) {
// Reuse the IDs and source states from the locking current read. An
// ordinary SELECT here could see an older REPEATABLE READ snapshot.
if len(taskIDs) == 0 {
return nil
return 0, nil
}
errorCode := "DEVICE_OFFLINE"
if err := tx.Model(&models.PurchaseTaskAttempt{}).
@@ -195,7 +303,7 @@ func markOfflinePurchaseStatus(
"status": models.PurchaseAttemptStatusFailed, "result_type": resultType,
"error_code": errorCode, "error_message": errorMessage, "finished_at": now,
}).Error; err != nil {
return internalError(err)
return 0, internalError(err)
}
updates := map[string]any{
"status": toStatus, "device_run_slot": gorm.Expr("NULL"), "account_run_slot": gorm.Expr("NULL"),
@@ -205,12 +313,13 @@ func markOfflinePurchaseStatus(
if toStatus == models.PurchaseTaskStatusFailed {
updates["active_slot"] = gorm.Expr("NULL")
}
if err := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).
result := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).
Where("id IN ? AND status = ?", taskIDs, fromStatus).
Updates(updates).Error; err != nil {
return internalError(err)
Updates(updates)
if result.Error != nil {
return 0, internalError(result.Error)
}
return nil
return result.RowsAffected, nil
}
func RunOfflineMonitor(ctx context.Context, service *Service, scanInterval, threshold time.Duration, onError func(error)) {
@@ -0,0 +1,432 @@
package device_test
import (
"context"
"errors"
"fmt"
"strings"
"sync"
"testing"
"time"
"go-admin/app/goauto/device"
"go-admin/app/goauto/models"
"go-admin/app/goauto/purchase"
"go-admin/app/goauto/task"
driver "github.com/go-sql-driver/mysql"
"github.com/google/uuid"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
type offlineSeed func(*testing.T) (*device.Service, device.RegisterResponse, string)
func offlineCollection(t *testing.T, db *gorm.DB, deviceID uint64, status string) models.CollectionTask {
t.Helper()
product := models.PDDProduct{GoodsID: strings.ReplaceAll(uuid.NewString(), "-", ""), URL: "https://example.invalid/test"}
rule := models.CollectionRule{Name: "test", ContentJSON: `{"steps":[]}`}
if err := db.Create(&product).Error; err != nil {
t.Fatal("synthetic product failed")
}
if err := db.Create(&rule).Error; err != nil {
t.Fatal("synthetic rule failed")
}
task := models.CollectionTask{PDDProductID: &product.ID, RuleID: rule.ID, DeviceID: &deviceID, Status: status, URLSnapshot: product.URL, GoodsIDSnapshot: product.GoodsID, RuleSnapshot: rule.ContentJSON}
if err := db.Create(&task).Error; err != nil {
t.Fatal("synthetic collection task failed")
}
return task
}
func offlinePurchase(t *testing.T, db *gorm.DB, deviceID uint64, status string) models.PurchaseTask {
t.Helper()
product := models.PDDProduct{GoodsID: strings.ReplaceAll(uuid.NewString(), "-", ""), URL: "https://example.invalid/test"}
if err := db.Create(&product).Error; err != nil {
t.Fatal("synthetic product failed")
}
task := models.PurchaseTask{PDDProductID: product.ID, DeviceID: &deviceID, ExecutionMode: models.PurchaseExecutionModeRehearsal, Status: status, PDDURLSnapshot: product.URL, PDDGoodsIDSnapshot: product.GoodsID, Quantity: 1, ReferenceUnitPriceCent: 100, MinUnitPriceCent: 20, MaxUnitPriceCent: 150, Currency: "CNY", RuleType: "pddPurchase", RuleSchemaVersion: 1, RequiredCapabilitiesJSON: `[]`, RuleSnapshot: `{}`, CreateRequestID: uuid.NewString()}
if err := db.Create(&task).Error; err != nil {
t.Fatal("synthetic purchase task failed")
}
return task
}
func requireOfflineDevice(t *testing.T, db *gorm.DB, id uint64, status string, events int64) {
t.Helper()
var row models.AgentDevice
if db.First(&row, id).Error != nil || row.Status != status {
t.Fatalf("device %d status=%s want=%s", id, row.Status, status)
}
var count int64
if db.Model(&models.AgentDeviceStatusEvent{}).Where("device_id = ?", id).Count(&count).Error != nil || count != events {
t.Fatalf("device %d events=%d want=%d", id, count, events)
}
}
func runOfflineV4MySQLTests(t *testing.T, db *gorm.DB, seed offlineSeed) {
for _, busyKind := range []string{"device", "collection", "purchase"} {
t.Run("busy_"+busyKind+"_does_not_block_other_devices", func(t *testing.T) {
s, first, _ := seed(t)
_, busy, _ := seed(t)
_, last, _ := seed(t)
var lockedModel any
var lockedID uint64
switch busyKind {
case "device":
lockedModel = &models.AgentDevice{}
lockedID = busy.DeviceID
case "collection":
task := offlineCollection(t, db, busy.DeviceID, models.TaskStatusRunning)
lockedModel = &models.CollectionTask{}
lockedID = task.ID
case "purchase":
task := offlinePurchase(t, db, busy.DeviceID, models.PurchaseTaskStatusRunning)
lockedModel = &models.PurchaseTask{}
lockedID = task.ID
}
holder := db.Begin()
if holder.Error != nil {
t.Fatal("begin failed")
}
defer holder.Rollback()
if holder.Clauses(clause.Locking{Strength: "UPDATE"}).First(lockedModel, lockedID).Error != nil {
t.Fatal("synthetic lock failed")
}
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
n, err := s.MarkStaleDevicesOffline(ctx, device.DefaultOfflineThreshold)
if err != nil || n != 2 {
t.Fatalf("busy %s prevented unrelated devices: changed=%d failed=%v", busyKind, n, err != nil)
}
requireOfflineDevice(t, db, first.DeviceID, models.DeviceStatusOffline, 1)
requireOfflineDevice(t, db, busy.DeviceID, models.DeviceStatusOnline, 0)
requireOfflineDevice(t, db, last.DeviceID, models.DeviceStatusOffline, 1)
if holder.Rollback().Error != nil {
t.Fatal("release failed")
}
n, err = s.MarkStaleDevicesOffline(context.Background(), device.DefaultOfflineThreshold)
if err != nil || n != 1 {
t.Fatal("busy device not retried on next scan")
}
})
}
for _, target := range []string{models.PurchaseTaskStatusRunning, models.PurchaseTaskStatusOrderSubmitStarted} {
t.Run("second_device_current_read_"+target, func(t *testing.T) {
s, first, _ := seed(t)
_, second, _ := seed(t)
offlinePurchase(t, db, first.DeviceID, models.PurchaseTaskStatusRunning)
task := offlinePurchase(t, db, second.DeviceID, models.PurchaseTaskStatusPending)
attempt := models.PurchaseTaskAttempt{TaskID: task.ID, AttemptID: uuid.NewString(), AttemptNumber: 1, Phase: models.PurchaseAttemptPhasePurchase, Status: models.PurchaseAttemptStatusRunning, DeviceID: &second.DeviceID, RuleSnapshotHash: "test", SpecDecisionSnapshot: `{}`}
if db.Create(&attempt).Error != nil {
t.Fatal("synthetic attempt failed")
}
callback := "v4_advance_second_device"
if db.Callback().Create().After("gorm:create").Register(callback, func(tx *gorm.DB) {
event, ok := tx.Statement.Dest.(*models.AgentDeviceStatusEvent)
if !ok || event.DeviceID != first.DeviceID {
return
}
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
if err := db.WithContext(ctx).Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("id = ?", task.ID).Updates(map[string]any{"status": target, "device_run_slot": 1}).Error; err != nil {
tx.AddError(errors.New("synthetic second device advancement failed"))
}
}) != nil {
t.Fatal("register advancement failed")
}
defer db.Callback().Create().Remove(callback)
n, err := s.MarkStaleDevicesOffline(context.Background(), device.DefaultOfflineThreshold)
if err != nil || n != 2 {
t.Fatalf("scan changed=%d failed=%v", n, err != nil)
}
var stored models.PurchaseTask
var storedAttempt models.PurchaseTaskAttempt
var event models.AgentDeviceStatusEvent
if db.First(&stored, task.ID).Error != nil || db.First(&storedAttempt, attempt.ID).Error != nil || db.Where("device_id = ?", second.DeviceID).First(&event).Error != nil {
t.Fatal("read outcome failed")
}
expected := models.PurchaseTaskStatusFailed
failed, unknown := int64(1), int64(0)
if target == models.PurchaseTaskStatusOrderSubmitStarted {
expected = models.PurchaseTaskStatusOrderResultUnknown
failed, unknown = 0, 1
}
if stored.Status != expected || storedAttempt.Status != models.PurchaseAttemptStatusFailed || event.FailedTaskCount != failed || event.OrderResultUnknownCount != unknown {
t.Fatalf("stale task snapshot: status=%s attempt=%s failed=%d unknown=%d", stored.Status, storedAttempt.Status, event.FailedTaskCount, event.OrderResultUnknownCount)
}
})
}
t.Run("event_failure_preserves_prior_device_commit", func(t *testing.T) {
s, first, _ := seed(t)
_, second, _ := seed(t)
_, last, _ := seed(t)
task := offlineCollection(t, db, second.DeviceID, models.TaskStatusRunning)
statement := fmt.Sprintf("CREATE TRIGGER v4_reject_event BEFORE INSERT ON agent_device_status_event FOR EACH ROW BEGIN IF NEW.device_id = %d THEN SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = 'synthetic_failure'; END IF; END", second.DeviceID)
if db.Exec(statement).Error != nil {
t.Fatal("test trigger failed")
}
defer db.Exec("DROP TRIGGER IF EXISTS v4_reject_event")
n, err := s.MarkStaleDevicesOffline(context.Background(), device.DefaultOfflineThreshold)
if err == nil || n != 1 {
t.Fatalf("must return earlier committed count on event failure: changed=%d failed=%v", n, err != nil)
}
requireOfflineDevice(t, db, first.DeviceID, models.DeviceStatusOffline, 1)
requireOfflineDevice(t, db, second.DeviceID, models.DeviceStatusOnline, 0)
requireOfflineDevice(t, db, last.DeviceID, models.DeviceStatusOnline, 0)
var stored models.CollectionTask
if db.First(&stored, task.ID).Error != nil || stored.Status != models.TaskStatusRunning {
t.Fatal("current task escaped rollback")
}
if db.Exec("DROP TRIGGER v4_reject_event").Error != nil {
t.Fatal("remove trigger failed")
}
n, err = s.MarkStaleDevicesOffline(context.Background(), device.DefaultOfflineThreshold)
if err != nil || n != 2 {
t.Fatal("next scan retry failed")
}
})
for _, kind := range []string{"1213", "1205", "text_3572", "cancel", "cancel_with_3572"} {
t.Run("non_contention_error_"+kind, func(t *testing.T) {
s, first, _ := seed(t)
_, second, _ := seed(t)
_, last, _ := seed(t)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
callback := "v4_inject_non_contention"
if db.Callback().Query().After("gorm:query").Register(callback, func(tx *gorm.DB) {
row, ok := tx.Statement.Dest.(*models.AgentDevice)
if !ok || row.ID != second.DeviceID {
return
}
if _, locking := tx.Statement.Clauses["FOR"]; !locking {
return
}
switch kind {
case "1213":
tx.AddError(&driver.MySQLError{Number: 1213, Message: "synthetic"})
case "1205":
tx.AddError(&driver.MySQLError{Number: 1205, Message: "synthetic"})
case "text_3572":
tx.AddError(errors.New("synthetic 3572"))
case "cancel":
cancel()
tx.AddError(context.Canceled)
case "cancel_with_3572":
cancel()
tx.AddError(&driver.MySQLError{Number: 3572, Message: "synthetic"})
}
}) != nil {
t.Fatal("register injection failed")
}
defer db.Callback().Query().Remove(callback)
n, err := s.MarkStaleDevicesOffline(ctx, device.DefaultOfflineThreshold)
if err == nil || n != 1 {
t.Fatalf("non-contention error swallowed or lost prior commit: changed=%d failed=%v", n, err != nil)
}
if (kind == "cancel" || kind == "cancel_with_3572") && !errors.Is(err, context.Canceled) {
t.Fatal("cancellation not preserved")
}
requireOfflineDevice(t, db, first.DeviceID, models.DeviceStatusOffline, 1)
requireOfflineDevice(t, db, second.DeviceID, models.DeviceStatusOnline, 0)
requireOfflineDevice(t, db, last.DeviceID, models.DeviceStatusOnline, 0)
})
}
for _, domain := range []string{"collection_claim", "purchase_replayed_claim"} {
t.Run("actual_"+domain+"_does_not_deadlock_scan", func(t *testing.T) {
s, d, token := seed(t)
_, other, _ := seed(t)
type claimMarker struct{}
locked, release := make(chan struct{}), make(chan struct{})
var once sync.Once
callback := "v4_actual_claim_barrier"
if db.Callback().Query().After("gorm:query").Register(callback, func(tx *gorm.DB) {
if tx.Statement.Context.Value(claimMarker{}) != true {
return
}
lockTable := "agent_device"
if domain == "purchase_replayed_claim" {
lockTable = "purchase_task"
}
if tx.Statement.Table != lockTable {
return
}
if _, ok := tx.Statement.Clauses["FOR"]; !ok {
return
}
once.Do(func() { close(locked); <-release })
}) != nil {
t.Fatal("register actual claim barrier failed")
}
defer db.Callback().Query().Remove(callback)
ctx, cancel := context.WithTimeout(context.WithValue(context.Background(), claimMarker{}, true), 5*time.Second)
defer cancel()
done := make(chan error, 1)
var taskID uint64
if domain == "collection_claim" {
record := offlineCollection(t, db, d.DeviceID, models.TaskStatusPending)
taskID = record.ID
service := task.NewService(db)
service.Now = s.Now
go func() {
_, err := service.Claim(ctx, record.ID, task.ActionRequest{RequestID: uuid.NewString()}, token)
done <- err
}()
} else {
record := offlinePurchase(t, db, d.DeviceID, models.PurchaseTaskStatusRunning)
taskID = record.ID
requestID := uuid.NewString()
if db.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("id = ?", record.ID).Update("claim_request_id", requestID).Error != nil {
t.Fatal("replay fixture failed")
}
service := purchase.NewService(db)
service.Now = s.Now
go func() {
_, err := service.Claim(ctx, record.ID, purchase.ActionRequest{RequestID: requestID}, token)
done <- err
}()
}
select {
case <-locked:
case <-ctx.Done():
close(release)
t.Fatal("actual claim lock not reached")
}
scanCtx, stopScan := context.WithTimeout(context.Background(), 2*time.Second)
n, err := s.MarkStaleDevicesOffline(scanCtx, device.DefaultOfflineThreshold)
stopScan()
close(release)
claimErr := <-done
if err != nil || claimErr != nil || n != 1 {
t.Fatalf("claim/scan conflict: changed=%d scanFailed=%v claimFailed=%v", n, err != nil, claimErr != nil)
}
requireOfflineDevice(t, db, d.DeviceID, models.DeviceStatusOnline, 0)
requireOfflineDevice(t, db, other.DeviceID, models.DeviceStatusOffline, 1)
if domain == "collection_claim" {
var stored models.CollectionTask
if db.First(&stored, taskID).Error != nil || stored.Status != models.TaskStatusPending || stored.LeaseExpiresAt == nil || stored.ClaimRequestID == nil {
t.Fatal("claim partially persisted")
}
}
})
}
t.Run("real_heartbeat_started_during_scan_resumes_once", func(t *testing.T) {
s, d, token := seed(t)
type scanMarker struct{}
type heartbeatMarker struct{}
locked, release, heartbeatStarted := make(chan struct{}), make(chan struct{}), make(chan struct{})
scanCallback, heartbeatCallback := "v4_scan_heartbeat_barrier", "v4_heartbeat_scan_barrier"
if db.Callback().Query().After("gorm:query").Register(scanCallback, func(tx *gorm.DB) {
if tx.Statement.Context.Value(scanMarker{}) == true && tx.Statement.Table == "agent_device" {
if _, ok := tx.Statement.Clauses["FOR"]; ok {
close(locked)
<-release
}
}
}) != nil {
t.Fatal("register scan barrier failed")
}
defer db.Callback().Query().Remove(scanCallback)
if db.Callback().Query().Before("gorm:query").Register(heartbeatCallback, func(tx *gorm.DB) {
if tx.Statement.Context.Value(heartbeatMarker{}) == true && tx.Statement.Table == "agent_device" {
if _, ok := tx.Statement.Clauses["FOR"]; ok {
close(heartbeatStarted)
}
}
}) != nil {
t.Fatal("register heartbeat barrier failed")
}
defer db.Callback().Query().Remove(heartbeatCallback)
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
scanDone, heartbeatDone := make(chan error, 1), make(chan error, 1)
go func() {
n, err := s.MarkStaleDevicesOffline(context.WithValue(ctx, scanMarker{}, true), device.DefaultOfflineThreshold)
if err == nil && n != 1 {
err = errors.New("wrong scan count")
}
scanDone <- err
}()
select {
case <-locked:
case <-ctx.Done():
close(release)
t.Fatal("scan lock not reached")
}
go func() {
_, err := s.Heartbeat(context.WithValue(ctx, heartbeatMarker{}, true), device.HeartbeatRequest{RequestID: uuid.NewString()}, token)
heartbeatDone <- err
}()
select {
case <-heartbeatStarted:
case <-ctx.Done():
close(release)
t.Fatal("heartbeat not started")
}
close(release)
if err := <-scanDone; err != nil {
t.Fatal("scan failed")
}
if err := <-heartbeatDone; err != nil {
t.Fatal("heartbeat failed")
}
requireOfflineDevice(t, db, d.DeviceID, models.DeviceStatusOnline, 2)
})
t.Run("explain_actual_lock_queries", func(t *testing.T) {
s, d, _ := seed(t)
offlineCollection(t, db, d.DeviceID, models.TaskStatusRunning)
offlinePurchase(t, db, d.DeviceID, models.PurchaseTaskStatusRunning)
type captured struct {
sql string
vars []any
}
queries := map[string]captured{}
callback := "v4_capture_lock_sql"
if db.Callback().Query().After("gorm:query").Register(callback, func(tx *gorm.DB) {
if _, ok := tx.Statement.Clauses["FOR"]; ok {
queries[tx.Statement.Table] = captured{tx.Statement.SQL.String(), append([]any(nil), tx.Statement.Vars...)}
}
}) != nil {
t.Fatal("capture SQL failed")
}
n, err := s.MarkStaleDevicesOffline(context.Background(), device.DefaultOfflineThreshold)
db.Callback().Query().Remove(callback)
if err != nil || n != 1 {
t.Fatal("capture scan failed")
}
for _, table := range []string{"agent_device", "collection_task", "purchase_task"} {
query, ok := queries[table]
if !ok {
t.Fatalf("missing production lock query %s", table)
}
var plan []struct {
Type string
Key *string
PossibleKeys *string
Rows int64
Extra string
}
if db.Raw("EXPLAIN "+query.sql, query.vars...).Scan(&plan).Error != nil {
t.Fatal("EXPLAIN failed")
}
t.Logf("synthetic fixture only table=%s sql=%s", table, query.sql)
for _, row := range plan {
key := "NULL"
if row.Key != nil {
key = *row.Key
}
possible := "NULL"
if row.PossibleKeys != nil {
possible = *row.PossibleKeys
}
t.Logf("type=%s key=%s possible_keys=%s estimated_rows=%d extra=%s", row.Type, key, possible, row.Rows, row.Extra)
}
}
})
}
+22 -4
View File
@@ -16,6 +16,7 @@ import (
"github.com/google/uuid"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
const DefaultHeartbeatIntervalSeconds = 15
@@ -305,7 +306,15 @@ func (service *Service) RevokeToken(ctx context.Context, deviceID uint64) error
func (service *Service) deactivate(ctx context.Context, deviceID uint64, revoke bool) error {
now := service.Now()
return service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var event *models.AgentDeviceStatusEvent
err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var device models.AgentDevice
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&device, deviceID).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return &ServiceError{Code: CodeDeviceNotFound, Message: "设备不存在", Retryable: false}
}
return internalError(err)
}
updates := map[string]any{"status": models.DeviceStatusDisabled}
if revoke {
updates["token_revoked_at"] = now
@@ -324,17 +333,26 @@ func (service *Service) deactivate(ctx context.Context, deviceID uint64, revoke
}
}
errorCode, errorMessage := CodeDeviceDisabled, "设备已由管理员停用"
if err := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.CollectionTask{}).
result = tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.CollectionTask{}).
Where("device_id = ? AND status = ?", deviceID, models.TaskStatusRunning).
Updates(map[string]any{
"status": models.TaskStatusFailed, "active_slot": gorm.Expr("NULL"),
"device_run_slot": gorm.Expr("NULL"), "lease_expires_at": gorm.Expr("NULL"),
"error_code": errorCode, "error_message": errorMessage, "finished_at": now,
}).Error; err != nil {
return internalError(err)
})
if result.Error != nil {
return internalError(result.Error)
}
if device.Status != models.DeviceStatusDisabled {
event = &models.AgentDeviceStatusEvent{DeviceID: device.ID, OccurredAt: now, FromStatus: device.Status, ToStatus: models.DeviceStatusDisabled, Reason: "disabled", LastHeartbeatAt: device.LastHeartbeatAt, FailedTaskCount: result.RowsAffected}
return insertDeviceStatusEvent(tx, event)
}
return nil
})
if err == nil && event != nil {
logDeviceStatusEvent(*event)
}
return err
}
func normalizeRegisterRequest(request RegisterRequest) RegisterRequest {
+2
View File
@@ -30,6 +30,8 @@ LIMIT 1`
func MigratedModels() []any {
return []any{
&models.AgentDevice{},
&models.AgentHeartbeatLog{},
&models.AgentDeviceStatusEvent{},
&models.AgentAppRelease{},
&models.AgentAppReleaseSetting{},
&models.PDDProduct{},
@@ -0,0 +1,33 @@
package models
import "time"
// AgentHeartbeatLog is diagnostic history, never a device liveness authority.
type AgentHeartbeatLog struct {
ID uint64 `gorm:"primaryKey;autoIncrement"`
DeviceID uint64 `gorm:"not null;index:ix_heartbeat_device_received,priority:1"`
ReceivedAt time.Time `gorm:"not null;index:ix_heartbeat_device_received,priority:2;index:ix_heartbeat_received"`
RequestID string `gorm:"size:36;not null"`
CurrentTaskID *uint64
RunningTaskID *uint64
Busy bool `gorm:"not null"`
AgentVersion string `gorm:"size:32;not null"`
ClientReportJSON *string `gorm:"type:text"`
ClientReportStatus string `gorm:"size:32;not null;check:ck_heartbeat_report_status,client_report_status IN ('none','accepted','client_report_invalid')"`
}
func (AgentHeartbeatLog) TableName() string { return "agent_heartbeat_log" }
type AgentDeviceStatusEvent struct {
ID uint64 `gorm:"primaryKey;autoIncrement"`
DeviceID uint64 `gorm:"not null"`
OccurredAt time.Time `gorm:"not null;index:ix_device_event_occurred"`
FromStatus string `gorm:"size:16;not null"`
ToStatus string `gorm:"size:16;not null"`
Reason string `gorm:"size:32;not null"`
LastHeartbeatAt *time.Time
FailedTaskCount int64 `gorm:"not null"`
OrderResultUnknownCount int64 `gorm:"not null"`
}
func (AgentDeviceStatusEvent) TableName() string { return "agent_device_status_event" }
+3
View File
@@ -129,6 +129,9 @@ func run() error {
defer stopOfflineMonitors()
for _, db := range sdk.Runtime.GetDb() {
service := goautodevice.NewService(db)
go goautodevice.RunDiagnosticsCleanup(offlineMonitorContext, service,
func(err error) { log.Error("device diagnostics cleanup failed") },
)
go goautopurchase.RunFailureSnapshotCleanup(
offlineMonitorContext, goautopurchase.NewService(db), time.Hour,
func(err error) { log.Error("purchase failure snapshot cleanup failed") },
@@ -0,0 +1,29 @@
package version_local
import (
"runtime"
"go-admin/app/goauto/models"
"go-admin/cmd/migrate/migration"
common "go-admin/common/models"
"gorm.io/gorm"
)
func init() {
_, file, _, _ := runtime.Caller(0)
migration.Migrate.SetVersion(migration.GetFilename(file), MigrateAgentDiagnostics)
}
// MigrateAgentDiagnostics only adds #378 diagnostic tables. MySQL DDL commits
// implicitly: idempotent table creation allows recovery after a partial up.
func MigrateAgentDiagnostics(db *gorm.DB, version string) error {
for _, model := range []any{&models.AgentHeartbeatLog{}, &models.AgentDeviceStatusEvent{}} {
if !db.Migrator().HasTable(model) {
if err := db.Migrator().CreateTable(model); err != nil {
return err
}
}
}
return db.Where("version = ?", version).FirstOrCreate(&common.Migration{Version: version}).Error
}
@@ -0,0 +1,48 @@
package version_local
import (
"go-admin/app/goauto/models"
common "go-admin/common/models"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"testing"
)
func TestAgentDiagnosticsMigrationAppendAndRollback(t *testing.T) {
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err := db.AutoMigrate(&common.Migration{}); err != nil {
t.Fatal(err)
}
for i := 0; i < 2; i++ {
if err := MigrateAgentDiagnostics(db, "test378"); err != nil {
t.Fatal(err)
}
}
for _, model := range []any{&models.AgentHeartbeatLog{}, &models.AgentDeviceStatusEvent{}} {
if !db.Migrator().HasTable(model) {
t.Fatalf("missing table %T", model)
}
}
for _, idx := range []struct {
model any
name string
}{{&models.AgentHeartbeatLog{}, "ix_heartbeat_device_received"}, {&models.AgentHeartbeatLog{}, "ix_heartbeat_received"}, {&models.AgentDeviceStatusEvent{}, "ix_device_event_occurred"}} {
if !db.Migrator().HasIndex(idx.model, idx.name) {
t.Fatalf("missing index %s", idx.name)
}
}
// Existing migration runner only supports up. Rollback is an explicit,
// operator-authorized removal of these two new diagnostic tables/version.
if err := db.Migrator().DropTable(&models.AgentHeartbeatLog{}, &models.AgentDeviceStatusEvent{}); err != nil {
t.Fatal(err)
}
if err := db.Where("version = ?", "test378").Delete(&common.Migration{}).Error; err != nil {
t.Fatal(err)
}
if err := MigrateAgentDiagnostics(db, "test378"); err != nil {
t.Fatal(err)
}
}