Compare commits

...
Author SHA1 Message Date
QiuSWandClaude Opus 5.5 d8b4f3a9d5 fix(android): track horizontal color rows frame to frame (#370)
Follow each horizontal color row by overlap with its previous frame (with a
position fallback for whole-page moves on multi-row grids), skip reverse
sweeps when rows share one container, and keep the earliest abnormal
COLOR_DISCOVERY reason.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NTDbDcwbDw1TSAcE6wfh2F
2026-10-10 15:04:06 +08:00
QiuSWandClaude Opus 5.5 23dafcf73f docs: record #372 v3 web release (#372)
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NTDbDcwbDw1TSAcE6wfh2F
2026-10-10 11:58:24 +08:00
QiuSWandClaude Opus 5.5 3d5e9646f6 merge: shop column in SYB product export (#372)
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NTDbDcwbDw1TSAcE6wfh2F
2026-10-10 11:52:01 +08:00
QiuSWandClaude Sonnet 5.5 e351a69912 feat(web): add shop column to SYB product export (#372)
Insert a text "店铺" column (row.shopName) right after SYB订单 and move
the 0.00 number format to the shifted price column. Sync the wiki mirror.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NTDbDcwbDw1TSAcE6wfh2F
2026-10-10 11:50:18 +08:00
QiuSW db64a49e69 merge: heartbeat history and Android sync diagnostics (#378) 2026-10-10 11:15:38 +08:00
QiuSW b79e281b97 docs: complete Android heartbeat diagnostic contracts (#378) 2026-10-10 11:12:07 +08:00
QiuSW dbe7e98365 fix(android): record handled purchase recovery failures in sync diagnostics (#378) 2026-10-10 11:02:46 +08:00
QiuSW 9f489fdf62 test: validate Android heartbeat reports through Server handler (#378) 2026-10-10 10:58:35 +08:00
QiuSW a28b869701 feat(android): persist bounded sync diagnostics and negotiate heartbeat reports (#378) 2026-10-10 10:57:33 +08:00
QiuSW c18bf78e0e docs: record Server heartbeat diagnostics and contention boundaries (#378) 2026-10-10 10:26:58 +08:00
QiuSW f5edd889f7 test(server): clarify heartbeat concurrency evidence (#378) 2026-10-10 10:24:00 +08:00
QiuSW d84728dde9 fix(server): isolate offline scan contention per device (#378) 2026-10-10 10:19:44 +08:00
QiuSWandClaude Opus 5.5 575d5f34cc docs: record #372 v2 web release (#372)
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NTDbDcwbDw1TSAcE6wfh2F
2026-10-10 09:54:04 +08:00
QiuSWandClaude Opus 5.5 4ad59bd231 merge: SYB price column and numeric amounts in exports (#372)
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NTDbDcwbDw1TSAcE6wfh2F
2026-10-10 09:48:42 +08:00
QiuSW 010e9b90ae feat: record bounded heartbeat and device transition diagnostics (#378) 2026-10-10 09:26:45 +08:00
37 changed files with 3883 additions and 129 deletions
@@ -1412,6 +1412,14 @@ class PddProductDetailCollector(
// These signatures stay in memory; only the aggregate tri-state is saved.
val diagnosticHorizontalRows = mutableMapOf<Set<String>, Boolean?>()
var diagnosticTermination = AgentDiagnosticReason.COLOR_FOUND
// Keep the earliest abnormal reason; later normal ends must not hide it.
fun terminate(reason: AgentDiagnosticReason) {
if (diagnosticTermination == AgentDiagnosticReason.COLOR_FOUND ||
diagnosticTermination == AgentDiagnosticReason.COLOR_EDGE_REACHED
) {
diagnosticTermination = reason
}
}
val imageAttempts = mutableSetOf<String>()
fun collectVisibleImages(values: List<VisibleSpecValue>) {
@@ -1561,7 +1569,18 @@ class PddProductDetailCollector(
// established left/right paging behavior.
if (rows.any { it.size > 1 }) {
val rowSeeds = rows.map { row -> row.map(VisibleSpecValue::text).toSet() }
// Frame-to-frame row tracking: rows are followed by overlap with
// their previous frame, not by the pass's starting frame.
val tracks = rows.map(::colorRowTrack).toMutableList()
val coveredRows = mutableSetOf<Int>()
var coveredIncomplete: Boolean? = null
rowSeeds.forEachIndexed { rowIndex, seed ->
if (rowIndex in coveredRows) {
// A shared container already swept this row's values.
diagnosticHorizontalRows[seed] = coveredIncomplete
?: diagnosticHorizontalRows[seed]?.takeIf { it }
return@forEachIndexed
}
val moveRight = rowIndex % 2 == 0
val horizontalSignatureReads = mutableMapOf<List<String>, Int>()
var diagnosticPreviousSignature: List<String>? = null
@@ -1569,11 +1588,13 @@ class PddProductDetailCollector(
var diagnosticMoved = false
var diagnosticRowUncertain = false
var diagnosticRowIncomplete: Boolean? = null
var sharedMoves = 0
var sharedBroken = tracks.size < 2
for (horizontalPass in 0..config.limits.getValue("specHorizontalSwipes")) {
var diagnosticBeforeClicks: List<String>? = null
collectVisibleColors { observedRows ->
diagnosticBeforeClicks = observedRows.filter { row -> row.any { it.text in seed } }
.singleOrNull()?.sortedBy { it.node.bounds.left }
diagnosticBeforeClicks = matchColorRow(tracks[rowIndex], observedRows, tracks.size, requireUniqueOverlap = true)
?.sortedBy { it.node.bounds.left }
?.let { if (moveRight) it else it.reversed() }?.let(::optionSignature)
}?.let { return it }
screen = parse(goodsId, config, evidence)
@@ -1581,11 +1602,25 @@ class PddProductDetailCollector(
if (!screen.pageEvidenceMatched) return failure("RULE_NOT_MATCHED", "采集期间离开 PDD 商品详情页")
rows = colorRows(screen)
observeColorDiscovery(screen, rows)
val matchedRow = rows.maxByOrNull { row -> row.count { it.text in seed } }
?.takeIf { row -> row.any { it.text in seed } }
val matches = tracks.map { matchColorRow(it, rows, tracks.size) }
val matchedRow = matches[rowIndex]
if (matchedRow != null &&
matchedRow.map(VisibleSpecValue::text).toSet() != tracks[rowIndex].texts
) {
sharedMoves++
val everyOtherRowMoved = tracks.indices.all { other ->
other == rowIndex || matches[other]?.let { match ->
match.map(VisibleSpecValue::text).toSet() != tracks[other].texts
} == true
}
if (!everyOtherRowMoved) sharedBroken = true
}
matches.forEachIndexed { index, match ->
match?.let { tracks[index] = colorRowTrack(it) }
}
if (matchedRow == null) {
diagnosticHorizontalUnknown = true
diagnosticTermination = AgentDiagnosticReason.COLOR_ROW_REFLOWED
terminate(AgentDiagnosticReason.COLOR_ROW_REFLOWED)
break
}
val currentRow = matchedRow.sortedBy { it.node.bounds.left }
@@ -1608,6 +1643,13 @@ class PddProductDetailCollector(
if (signatureReads > config.limits.getValue("stableEdgeReads")) {
if (!diagnosticRowUncertain && diagnosticStableReads >= config.limits.getValue("stableEdgeReads")) {
diagnosticRowIncomplete = false
// Every other row changed with every move of this
// row and this row is stable at its edge: they share
// one container and were swept along with it.
if (!sharedBroken && sharedMoves > 0) {
coveredRows += (rowIndex + 1 until tracks.size)
coveredIncomplete = false
}
}
break
}
@@ -1640,7 +1682,7 @@ class PddProductDetailCollector(
diagnosticTrailingEmptyReads = if (rows.isEmpty()) diagnosticTrailingEmptyReads + 1 else 0
if (rows.isEmpty()) diagnosticHorizontalUnknown = true
if (!discoverVertically) {
diagnosticTermination = AgentDiagnosticReason.COLOR_EDGE_REACHED
terminate(AgentDiagnosticReason.COLOR_EDGE_REACHED)
break
}
@@ -1652,21 +1694,21 @@ class PddProductDetailCollector(
}
previousVerticalSignature = verticalSignature
if (verticalStable >= config.limits.getValue("stableEdgeReads")) {
diagnosticTermination = AgentDiagnosticReason.COLOR_EDGE_REACHED
terminate(AgentDiagnosticReason.COLOR_EDGE_REACHED)
break
}
if (verticalPass == config.limits.getValue("specVerticalSwipes")) {
diagnosticTermination = AgentDiagnosticReason.COLOR_SCAN_LIMIT
terminate(AgentDiagnosticReason.COLOR_SCAN_LIMIT)
break
}
val anchor = rows.flatten().firstOrNull()?.node ?: specPanelContainer
if (anchor == null) {
diagnosticHorizontalUnknown = true
diagnosticTermination = AgentDiagnosticReason.COLOR_CONTAINER_UNAVAILABLE
terminate(AgentDiagnosticReason.COLOR_CONTAINER_UNAVAILABLE)
break
}
if (!driver.swipeSpec(SwipeDirection.UP, anchor)) {
diagnosticTermination = AgentDiagnosticReason.COLOR_SWIPE_FAILED
terminate(AgentDiagnosticReason.COLOR_SWIPE_FAILED)
break
}
diagnosticVerticalSwipes++
@@ -2219,6 +2261,39 @@ class PddProductDetailCollector(
const val COLOR_IMAGE_MAX_TOTAL_BYTES = 8 * 1024 * 1024
}
/** Last observed frame of one color row, used to follow it across horizontal swipes. */
private class ColorRowTrack(val texts: Set<String>, val centerY: Int, val tolerance: Int)
private fun colorRowTrack(row: List<VisibleSpecValue>) = ColorRowTrack(
row.map(VisibleSpecValue::text).toSet(),
row.map { it.node.bounds.centerY }.average().toInt(),
(row.first().node.bounds.height / 2).coerceIn(24, 80),
)
/**
* Finds the current-frame row for a tracked row. Overlap with the previous
* frame wins. With no overlap (a whole-page move) a multi-row grid whose
* row count is unchanged may match by vertical position, but only when
* exactly one current row sits within the colorRows tolerance. Anything
* else is not the same row and stays a reflow.
*/
private fun matchColorRow(
track: ColorRowTrack,
rows: List<List<VisibleSpecValue>>,
trackedRowCount: Int,
requireUniqueOverlap: Boolean = false,
): List<VisibleSpecValue>? {
val overlapping = rows.filter { row -> row.any { it.text in track.texts } }
if (overlapping.isNotEmpty()) {
if (requireUniqueOverlap && overlapping.size > 1) return null
return overlapping.maxByOrNull { row -> row.count { it.text in track.texts } }
}
if (trackedRowCount < 2 || rows.size != trackedRowCount) return null
return rows.filter { row ->
kotlin.math.abs(row.map { it.node.bounds.centerY }.average().toInt() - track.centerY) <= track.tolerance
}.singleOrNull()
}
private fun colorRows(screen: ParsedPddScreen): List<List<VisibleSpecValue>> {
val values = screen.dimensions.filter { it.key == "color" }.flatMap { it.values }
.sortedWith(compareBy({ it.node.bounds.centerY }, { it.node.bounds.left }))
@@ -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) {
@@ -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,7 @@ class AgentForegroundService : Service() {
}
override fun onDestroy() {
foregroundStarted = false
stateStore.stop()
probeHandoff.invalidate()
if (diagnosticInstance === this) diagnosticInstance = null
@@ -269,6 +286,53 @@ 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())
@@ -278,7 +342,7 @@ class AgentForegroundService : Service() {
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)
@@ -294,8 +358,14 @@ class AgentForegroundService : Service() {
val activeCredentials = credentials ?: error("设备尚未取得认证凭据")
val sync = AgentSyncCycle(
heartbeat = { currentTaskId ->
api.heartbeat(activeCredentials.token, currentTaskId, AgentCapabilities.supported).also {
check(it.deviceId == activeCredentials.deviceId) { "心跳设备身份不一致" }
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 = {
@@ -303,10 +373,15 @@ class AgentForegroundService : Service() {
else purchaseStore.activeTaskId().also(runningTaskId::set)
},
executionActive = { taskMutex.currentTaskId() != null },
recover = { recoverInterruptedPurchases(api, activeCredentials.token) },
flush = { flushPurchaseOutbox(api, activeCredentials.token) },
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)
@@ -324,8 +399,9 @@ class AgentForegroundService : Service() {
heartbeat = true,
)
updateNotification(if (heartbeat.busy) "在线 · 执行中" else "在线 · 空闲")
if (sync.canClaim) scheduleTask(api, activeCredentials.token) else MANUAL_BUSY
if (sync.canClaim) round.step("claim") { scheduleTask(api, activeCredentials.token) } else MANUAL_BUSY
} catch (error: AgentApiException) {
round.failed(error)
cancelIdleReturn("网络或服务端请求失败")
val authenticationError = isAgentAuthenticationError(error)
val connectionFailure = syncConnectionFailureCode(error, heartbeatSucceeded)
@@ -341,6 +417,7 @@ class AgentForegroundService : Service() {
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()
@@ -780,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
@@ -797,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}")) &&
@@ -964,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()
}
@@ -1014,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) {
@@ -1045,10 +1130,6 @@ class AgentForegroundService : Service() {
}
val requestId = UUID.randomUUID().toString()
purchaseStore.completeAndEnqueue(interrupted.taskId, interrupted.attemptId, requestId, purchaseResultPayload(requestId, interrupted.attemptId, outcome))
} catch (error: AgentApiException) {
// Keep the local irreversible marker. The same server request ID
// is replayed after connectivity returns; the order click is never repeated.
if (isAgentAuthenticationError(error)) throw error
}
return@forEach
}
@@ -18,6 +18,7 @@ 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. */
@@ -30,21 +31,21 @@ internal class AgentSyncCycle(
private val hasPending: () -> Boolean,
) {
fun run(): AgentSyncResult {
var pulse = runCatching { heartbeat(activeTaskId()) }
var pulse = captureSync { heartbeat(activeTaskId()) }
var authFailure = isAgentAuthenticationError(pulse.exceptionOrNull())
var recovered = false
var flushed = false
if (!authFailure && !executionActive()) {
val recovery = runCatching(recover)
val recovery = captureSync(recover)
recovered = recovery.isSuccess
authFailure = isAgentAuthenticationError(recovery.exceptionOrNull())
if (!authFailure) {
val upload = runCatching(flush)
val upload = captureSync(flush)
flushed = upload.isSuccess
authFailure = isAgentAuthenticationError(upload.exceptionOrNull())
}
if (!authFailure && (pulse.exceptionOrNull() as? AgentApiException)?.code == "DEVICE_TASK_MISMATCH") {
pulse = runCatching { heartbeat(activeTaskId()) }
pulse = captureSync { heartbeat(activeTaskId()) }
}
}
val failure = pulse.exceptionOrNull()
@@ -59,6 +60,28 @@ internal class AgentSyncCycle(
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,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)
}
}
}
@@ -2092,6 +2092,122 @@ class PddProductDetailCollectorTest {
}
}
private fun gridColors(count: Int) = (1..count).map { "款%02d色".format(it) }
private fun gridConfig() = config().copy(
timeoutsMs = config().timeoutsMs + ("overall" to 600_000),
limits = config().limits + mapOf("specHorizontalSwipes" to 12, "specVerticalSwipes" to 3, "stableEdgeReads" to 2),
)
private fun sharedGridCollects47Colors(step: Int) {
// 47 colors in 2 rows x 24 columns sharing one horizontal container.
val colors = gridColors(47)
val grid = listOf(colors.filterIndexed { i, _ -> i % 2 == 0 }, colors.filterIndexed { i, _ -> i % 2 == 1 })
assertEquals(listOf(24, 23), grid.map { it.size })
val prices = colors.mapIndexed { i, c -> c to 1000L + i * 10L }.toMap()
val events = mutableListOf<AgentDiagnosticEvent>()
val driver = FakeCollectorDriver(
colors = colors,
prices = prices,
horizontalGrid = grid,
gridVisibleColumns = 4,
gridStep = step,
sizePages = listOf(listOf("S"), listOf("M")),
hideColorHeadingAfterFirstVerticalPage = true,
)
var clock = 0L
val result = PddProductDetailCollector(driver, { clock }, { clock += it }, taskId = 370, diagnostic = events::add)
.collect(GOODS_ID, rule(gridConfig()))
assertTrue(result.successful)
val payload = requireNotNull(result.payload)
assertEquals("completed", payload.status)
assertEquals(colors.toSet(), payload.dimensions.first { it.key == "color" }.values.toSet())
assertEquals(47, payload.colorPrices.size)
assertEquals(prices, payload.colorPrices.associate { it.color to it.priceCent })
// 47 colors x 2 sizes, each SKU carrying its own color's price.
assertEquals(94, payload.skus.size)
assertEquals(setOf("S", "M"), payload.skus.map { it.specs.getValue("size") }.toSet())
payload.skus.forEach { assertEquals(prices.getValue(it.specs.getValue("color")), it.priceCent) }
assertTrue(payload.missing.isEmpty())
// The shared container is swept once: no reverse sweep away from an edge.
assertTrue(driver.gridSwipeLog.none { it.first == SwipeDirection.RIGHT && it.second > 0 })
assertTrue(driver.gridSwipeLog.count { it.first == SwipeDirection.LEFT } <= 12)
val event = events.single { it.stage == AgentDiagnosticStage.COLOR_DISCOVERY }
assertEquals(47, event.clickableColorCount)
}
@Test
fun `shared container two row grid collects all colors when swipes overlap`() = sharedGridCollects47Colors(step = 3)
@Test
fun `shared container two row grid collects all colors when swipes move whole pages`() = sharedGridCollects47Colors(step = 4)
@Test
fun `independent row grid keeps the snake sweep and does not claim a shared container`() {
val colors = gridColors(20)
val grid = listOf(colors.take(10), colors.drop(10))
val driver = FakeCollectorDriver(
colors = colors,
horizontalGrid = grid,
gridVisibleColumns = 4,
gridStep = 3,
gridSharedContainer = false,
sizePages = listOf(listOf("S"), listOf("M")),
hideColorHeadingAfterFirstVerticalPage = true,
)
var clock = 0L
val result = PddProductDetailCollector(driver, { clock }, { clock += it }).collect(GOODS_ID, rule(gridConfig()))
assertTrue(result.successful)
val collected = requireNotNull(result.payload).dimensions.first { it.key == "color" }.values
// Row 0 is swept to its end; row 1 only exposes its first window because
// it starts at its own left edge, exactly as before.
assertEquals((colors.take(10) + colors.drop(10).take(4)).toSet(), collected.toSet())
assertEquals(14, collected.size)
assertTrue(driver.gridSwipeLog.any { it.first == SwipeDirection.RIGHT && it.second > 0 })
}
@Test
fun `row disappearing keeps reflow reason after later empty vertical reads`() {
val colors = listOf("A色", "B色", "C色", "D色", "E色", "F色")
val events = mutableListOf<AgentDiagnosticEvent>()
val driver = FakeCollectorDriver(
colors = colors,
colorPages = listOf(colors.take(4), colors.takeLast(2)),
rowSize = 2,
sizePages = listOf(listOf("S"), listOf("M")),
hideColorHeadingAfterFirstVerticalPage = true,
)
var clock = 0L
val result = PddProductDetailCollector(driver, { clock }, { clock += it }, taskId = 370, diagnostic = events::add)
.collect(GOODS_ID, rule(gridConfig()))
assertTrue(result.successful)
val event = events.single { it.stage == AgentDiagnosticStage.COLOR_DISCOVERY }
assertEquals(AgentDiagnosticReason.COLOR_ROW_REFLOWED, event.reason)
assertTrue(requireNotNull(event.colorTrailingEmptyReadCount) > 0)
assertEquals(null, event.colorHorizontalIncomplete)
}
@Test
fun `single row disjoint reflow keeps reflow reason`() {
val colors = listOf("A色", "B色", "C色", "D色")
val events = mutableListOf<AgentDiagnosticEvent>()
var clock = 0L
val result = PddProductDetailCollector(
FakeCollectorDriver(colors = colors, colorPages = listOf(colors.take(2), colors.takeLast(2))),
{ clock }, { clock += it }, taskId = 370, diagnostic = events::add,
).collect(GOODS_ID, rule())
assertTrue(result.successful)
assertEquals(AgentDiagnosticReason.COLOR_ROW_REFLOWED, events.single { it.stage == AgentDiagnosticStage.COLOR_DISCOVERY }.reason)
}
private class FakeCollectorDriver(
colors: List<String> = listOf("红色"),
sizes: List<String> = listOf("S"),
@@ -2146,6 +2262,14 @@ class PddProductDetailCollectorTest {
// restored to top. Heading count drops but color/size values are
// untouched, so this must NOT be treated as a collapsed panel.
private val specPanelDropsExtraDimensionAfterTopSwipe: Boolean = false,
// #370: horizontal multi-row color grid. Each swipe moves gridStep
// columns; gridSharedContainer moves all rows together, otherwise only
// the anchored row. Left/right swipes are logged with the largest
// column offset before the gesture.
private val horizontalGrid: List<List<String>>? = null,
private val gridVisibleColumns: Int = 4,
private val gridStep: Int = 3,
private val gridSharedContainer: Boolean = true,
) : PddCollectorDriver {
var captureCount = 0
var clickCount = 0
@@ -2155,6 +2279,8 @@ class PddProductDetailCollectorTest {
var backCount = 0
var entryClickCount = 0
var restoreGestures = 0
val gridSwipeLog = mutableListOf<Pair<SwipeDirection, Int>>()
private val gridOffsets = IntArray(horizontalGrid?.size ?: 0)
private var selected: String? = initialSelectedColor
private var previousSelected: String? = null
private var horizontalPage = 0
@@ -2224,19 +2350,22 @@ class PddProductDetailCollectorTest {
val continuationPage = hideDimensionHeadingsAfterFirstVerticalPage && verticalPage > 0
val panelCollapsedNow = specPanelCollapsesAfterTopSwipe && downSwipeCount >= 1
val hideColorNow = specPanelHidesColorInitially && downSwipeCount == 0
val visibleColors = colorVerticalPages?.get(verticalPage.coerceAtMost(colorVerticalPages.lastIndex))
val gridPlaced = horizontalGrid?.flatMapIndexed { r, row ->
row.drop(gridOffsets[r]).take(gridVisibleColumns).mapIndexed { c, color -> Triple(color, r, c) }
}
val visibleColors = gridPlaced?.map { it.first }
?: colorVerticalPages?.get(verticalPage.coerceAtMost(colorVerticalPages.lastIndex))
?: colorPages[horizontalPage.coerceAtMost(colorPages.lastIndex)]
val colorRowCount = (visibleColors.size + rowSize - 1) / rowSize
val colorRowCount = horizontalGrid?.size ?: ((visibleColors.size + rowSize - 1) / rowSize)
val placedColors = gridPlaced ?: visibleColors
.filterNot { hideSelectedColorOption && it == selected }
.mapIndexed { index, color -> Triple(color, index / rowSize, index % rowSize) }
val sizeHeadingTop = maxOf(700, 470 + colorRowCount * 90 + 20)
if (!continuationPage && !panelCollapsedNow && !(hideColorHeadingAfterFirstVerticalPage && verticalPage > 0)) {
nodes += node("scroll/color-heading", "颜色分类", 20, 400, 300, 450, parentPath = "scroll")
}
if (!continuationPage && !hideColorNow && !panelCollapsedNow) {
visibleColors
.filterNot { hideSelectedColorOption && it == selected }
.forEachIndexed { index, color ->
val row = index / rowSize
val column = index % rowSize
placedColors.forEach { (color, row, column) ->
val path = "scroll/color-$color-$captureCount"
val left = 30 + column * 230
val top = 470 + row * 90
@@ -2357,6 +2486,18 @@ class PddProductDetailCollectorTest {
override fun swipeSpec(direction: SwipeDirection, anchor: SnapshotNode?): Boolean {
swipes += direction to anchor
if (horizontalGrid != null && (direction == SwipeDirection.LEFT || direction == SwipeDirection.RIGHT)) {
gridSwipeLog += direction to (gridOffsets.maxOrNull() ?: 0)
val delta = if (direction == SwipeDirection.LEFT) gridStep else -gridStep
val moved = if (gridSharedContainer) horizontalGrid.indices.toList()
else listOf((((anchor?.bounds?.top ?: 470) - 470) / 90).coerceIn(0, horizontalGrid.lastIndex))
val sharedMax = horizontalGrid.maxOf { it.size } - gridVisibleColumns
moved.forEach { r ->
val max = if (gridSharedContainer) sharedMax else horizontalGrid[r].size - gridVisibleColumns
gridOffsets[r] = (gridOffsets[r] + delta).coerceIn(0, maxOf(0, max))
}
return true
}
if ((direction == SwipeDirection.UP || direction == SwipeDirection.DOWN) && !verticalSwipeSucceeds) return false
when (direction) {
SwipeDirection.LEFT -> horizontalPage = (horizontalPage + 1).coerceAtMost(colorPages.lastIndex)
@@ -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,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
}
}
}
@@ -8,6 +8,19 @@ 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)
+18 -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: 08da88a4f106004cb5f02448260d2fc22a83e9cc
synchronized_at: 2026-10-10T00:59:09Z
wiki_revision: 5508fe99e1040dff61d3bc6dbdf08ce0a1fd02ad
synchronized_at: 2026-10-10T03:10:31Z
<!-- gitea-wiki-mirror:end -->
<!-- gitea-wiki-mirror:start -->
@@ -682,3 +682,19 @@ Web 唯一展示位置为“采集采购 → SYB 同步记录”:列表状态
- `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当前占据屏幕。
+16 -3
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: 35c42f6136426bd3c295fe912d491c3fd980fcb4
synchronized_at: 2026-10-10T01:41:05Z
wiki_revision: d427a46283561df71714929c13b7156ee42debea
synchronized_at: 2026-10-10T03:47:23Z
<!-- gitea-wiki-mirror:end -->
<!-- gitea-wiki-mirror:start -->
@@ -955,5 +955,18 @@ Android 0.9.64 / versionCode 77,源码 `6550b9f`(分支实现,尚未安装
- 格式:时间 `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 列与匹配时间留空。要导出已用退货,先把处理阶段筛为「已用退货」。
- SYB 商品导出列(按顺序):SYB订单、**店铺**、虾皮商品ID、**SYB售价(TWD)**、SYB颜色、SYB尺码、SYB数量、yeeke颜色、yeeke尺码、yeeke数量、匹配时间。店铺取商品行 `shopName`(SYB 订单的店铺,与页面「店铺」列同值,文本单元格,无值留空,不使用店铺筛选下拉的显示名);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条额度;旧进行中记录只表示中断/不完整,不伪造原现场、失败原因或成功。报告确认不会清除发送期间新增失败;本地确认失败或响应丢失允许重复,报告不能用于改变设备在线或采购状态。
+46 -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: af209183f85200ce53f79f82ab682e17bb0aa942
synchronized_at: 2026-10-09T10:28:29Z
wiki_revision: d06fa26d85339641e3d67e5ee6314205173c24c5
synchronized_at: 2026-10-10T03:10:51Z
<!-- gitea-wiki-mirror:end -->
# 故障排查
@@ -176,3 +176,47 @@ adb -s <serial> shell run-as cn.ilapage.goauto.agent sqlite3 -readonly databases
数据库升级仅限手机私有 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页面内容。
+38 -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: ad212f537b7d153e6337b3f2851e661680bfba20
synchronized_at: 2026-10-10T01:17:11Z
wiki_revision: d08a0dfbb8356f9cd7212d495e4f11beca59e293
synchronized_at: 2026-10-10T03:56:35Z
<!-- gitea-wiki-mirror:end -->
<!-- gitea-wiki-mirror:start -->
@@ -383,3 +383,39 @@ Provider 故障日志只允许记录调用关联 ID、操作类型、耗时、
- 按内容验收:公网 `/` 与 `/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分钟心跳验收。
## #372 v3 导出店铺列的纯 Web 发布(2026-10-10)
- 用户授权合并 #372 v3 并部署线上。`feat/372-export-shop` 以 `--no-ff` 合并为 main `3d5e964`;Web 单元测试 93/93 通过,`build:prod` 通过。相对上一线上 Web(`4ad59bd`),Web 只有 v3 的导出改动。
- **main 已含 #378 的 Server 代码和迁移,本次没有部署**:线上 Server 仍是 `72a15f8` 构建的二进制。#378 的 Server 发布和迁移需另行授权。
- 发布目录 `/home/goauto/releases/20261010-3d5e964-372v3-web`,只替换 dist。goauto-server 与 config 从上一发布逐字复制(cmp/diff 一致),static/temp/var 沿用原目录,旧 js/css 哈希资源保留。index.html SHA256 `b900a40ec4691f0fb2c3d89e1e20826e7ca3205db311a33c24087767a515979a`。
- 只原子切换 current,不重启 GoAuto/Nginx;切换前后 PID 4786 不变,服务 active。
- 按内容验收:`/` 与 `/login` 返回新 index,10 项 JS/CSS 200,线上导出模块的 SYB 表头含「店铺」,health 200,验证码 code=200,未登录业务接口 401。登录后的实际导出待业务验收。
- 回滚:current 原子切回 `/home/goauto/releases/20261010-4ad59bd-372v2-web`,不需要重启。
@@ -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)
}
}
+3 -3
View File
@@ -51,7 +51,7 @@ const blank = value => (value === null || value === undefined ? '' : value)
// Column indexes written as numbers; the saver shows them as 0.00.
export const PURCHASE_EXPORT_NUMBER_COLS = [8]
export const SYB_PRODUCT_EXPORT_NUMBER_COLS = [2]
export const SYB_PRODUCT_EXPORT_NUMBER_COLS = [3]
export const PURCHASE_EXPORT_HEADER = ['SYB订单号', '虾皮商品ID', '颜色', '尺码', '数量', '状态', 'PDD单号', '下单时间', '价格(元)']
@@ -64,7 +64,7 @@ export function purchaseExportRow(task, statusLabel = value => value) {
]
}
export const SYB_PRODUCT_EXPORT_HEADER = ['SYB订单', '虾皮商品ID', 'SYB售价(TWD)', 'SYB颜色', 'SYB尺码', 'SYB数量', 'yeeke颜色', 'yeeke尺码', 'yeeke数量', '匹配时间']
export const SYB_PRODUCT_EXPORT_HEADER = ['SYB订单', '店铺', '虾皮商品ID', 'SYB售价(TWD)', 'SYB颜色', 'SYB尺码', 'SYB数量', 'yeeke颜色', 'yeeke尺码', 'yeeke数量', '匹配时间']
export const isActiveReturnMatch = match => Boolean(match) && (match.status === 'matched' || match.status === 'confirmed')
@@ -73,7 +73,7 @@ export function sybProductExportRow(product, match) {
const active = isActiveReturnMatch(match) ? match : null
const yeeke = active ? splitYeekeVariation(active.variationName) : { color: '', size: '' }
return [
blank(product.orderCode), blank(product.shopeeItemId), centToYuanNumber(product.unitPriceCent), blank(product.targetColor), blank(product.targetSize), blank(product.quantity),
blank(product.orderCode), blank(product.shopName), blank(product.shopeeItemId), centToYuanNumber(product.unitPriceCent), blank(product.targetColor), blank(product.targetSize), blank(product.quantity),
yeeke.color, yeeke.size, active ? blank(active.yeekeQuantity) : '', active ? formatExportTime(active.matchedAt) : ''
]
}
+17 -14
View File
@@ -58,17 +58,18 @@ test('cent to yuan number is numeric, zero stays zero, empty stays empty', () =>
})
test('syb product row uses only an active match and ignores cancelled ones', () => {
const product = { orderCode: 'O1', shopeeItemId: '99', unitPriceCent: 1230, targetColor: '白色', targetSize: 'L', quantity: 1 }
const product = { orderCode: 'O1', shopName: 'shop-a', shopeeItemId: '99', unitPriceCent: 1230, targetColor: '白色', targetSize: 'L', quantity: 1 }
const matchedAt = new Date(2026, 9, 9, 9, 0, 0).toISOString()
assert.deepEqual(lib.SYB_PRODUCT_EXPORT_HEADER, ['SYB订单', '虾皮商品ID', 'SYB售价(TWD)', 'SYB颜色', 'SYB尺码', 'SYB数量', 'yeeke颜色', 'yeeke尺码', 'yeeke数量', '匹配时间'])
assert.deepEqual(lib.SYB_PRODUCT_EXPORT_NUMBER_COLS, [2])
assert.deepEqual(lib.SYB_PRODUCT_EXPORT_HEADER, ['SYB订单', '店铺', '虾皮商品ID', 'SYB售价(TWD)', 'SYB颜色', 'SYB尺码', 'SYB数量', 'yeeke颜色', 'yeeke尺码', 'yeeke数量', '匹配时间'])
assert.deepEqual(lib.SYB_PRODUCT_EXPORT_NUMBER_COLS, [3])
assert.deepEqual(lib.sybProductExportRow(product, { status: 'confirmed', variationName: '白色,L', yeekeQuantity: 3, matchedAt }),
['O1', '99', 12.3, '白色', 'L', 1, '白色', 'L', 3, '2026-10-09 09:00:00'])
assert.deepEqual(lib.sybProductExportRow(product, { status: 'matched', variationName: '均碼', yeekeQuantity: 1, matchedAt }).slice(6, 9), ['均碼', '', 1])
assert.deepEqual(lib.sybProductExportRow(product, { status: 'cancelled', variationName: 'x', yeekeQuantity: 1, matchedAt }), ['O1', '99', 12.3, '白色', 'L', 1, '', '', '', ''])
assert.deepEqual(lib.sybProductExportRow(product, null), ['O1', '99', 12.3, '白色', 'L', 1, '', '', '', ''])
assert.equal(lib.sybProductExportRow({ ...product, unitPriceCent: 0 }, null)[2], 0)
for (const value of [null, undefined]) assert.equal(lib.sybProductExportRow({ ...product, unitPriceCent: value }, null)[2], '')
['O1', 'shop-a', '99', 12.3, '白色', 'L', 1, '白色', 'L', 3, '2026-10-09 09:00:00'])
assert.deepEqual(lib.sybProductExportRow(product, { status: 'matched', variationName: '均碼', yeekeQuantity: 1, matchedAt }).slice(7, 10), ['均碼', '', 1])
assert.deepEqual(lib.sybProductExportRow(product, { status: 'cancelled', variationName: 'x', yeekeQuantity: 1, matchedAt }), ['O1', 'shop-a', '99', 12.3, '白色', 'L', 1, '', '', '', ''])
assert.deepEqual(lib.sybProductExportRow(product, null), ['O1', 'shop-a', '99', 12.3, '白色', 'L', 1, '', '', '', ''])
assert.equal(lib.sybProductExportRow({ ...product, unitPriceCent: 0 }, null)[3], 0)
for (const value of [null, undefined]) assert.equal(lib.sybProductExportRow({ ...product, unitPriceCent: value }, null)[3], '')
for (const value of [null, undefined]) assert.equal(lib.sybProductExportRow({ ...product, shopName: value }, null)[1], '')
})
// Runs the real Export2Excel writer with file-saver stubbed and reads the xlsx back.
@@ -96,15 +97,17 @@ test('excel writer makes numeric cells, 0.00 format only for requested columns',
test('real export row types and format end up in the sheet', () => {
const header = lib.SYB_PRODUCT_EXPORT_HEADER
const rows = [
lib.sybProductExportRow({ orderCode: 'O1', unitPriceCent: 1230 }, null),
lib.sybProductExportRow({ orderCode: 'O1', shopName: 'shop-a', unitPriceCent: 1230 }, null),
lib.sybProductExportRow({ orderCode: 'O2', unitPriceCent: 0 }, null),
lib.sybProductExportRow({ orderCode: 'O3', unitPriceCent: null }, null)
]
const sheet = writeWorkbook({ header, data: rows, filename: 't', numberFormatCols: lib.SYB_PRODUCT_EXPORT_NUMBER_COLS })
assert.equal(sheet.C1.v, 'SYB售价(TWD)')
assert.equal(sheet.C2.t, 'n'); assert.equal(sheet.C2.v, 12.3); assert.equal(sheet.C2.z, '0.00')
assert.equal(sheet.C3.t, 'n'); assert.equal(sheet.C3.v, 0)
assert.equal(sheet.C4.v, ''); assert.notEqual(sheet.C4.t, 'n')
assert.equal(sheet.B1.v, '店铺')
assert.equal(sheet.B2.t, 's'); assert.equal(sheet.B2.v, 'shop-a'); assert.notEqual(sheet.B2.z, '0.00')
assert.equal(sheet.D1.v, 'SYB售价(TWD)')
assert.equal(sheet.D2.t, 'n'); assert.equal(sheet.D2.v, 12.3); assert.equal(sheet.D2.z, '0.00')
assert.equal(sheet.D3.t, 'n'); assert.equal(sheet.D3.v, 0)
assert.equal(sheet.D4.v, ''); assert.notEqual(sheet.D4.t, 'n')
})
function harness(total, options = {}) {