feat(android): persist bounded sync diagnostics and negotiate heartbeat reports (#378)

This commit is contained in:
QiuSW
2026-10-10 10:57:33 +08:00
parent c18bf78e0e
commit a28b869701
15 changed files with 1537 additions and 17 deletions
@@ -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,216 @@
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
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)
persistSteps()
try {
val value = action()
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()
}
}
}
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,11 @@ 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) } },
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 +395,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 +413,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 +853,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 +871,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 +1040,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()
}
@@ -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,15 @@ 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)
}
@@ -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)
}
}
}
@@ -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,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)