diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/MainActivity.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/MainActivity.kt index ce95a34..ecaf95d 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/MainActivity.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/MainActivity.kt @@ -50,6 +50,10 @@ import com.roubao.autopilot.vlm.RequirementExtractionResult import com.roubao.autopilot.vlm.RequirementExtractor import com.roubao.autopilot.vlm.RequirementProbeState import com.roubao.autopilot.vlm.RequirementProviderEndpointPolicy +import com.roubao.autopilot.vlm.REQUIREMENT_PROMPT_VERSION +import com.roubao.autopilot.vlm.REQUIREMENT_SCHEMA_VERSION +import com.roubao.autopilot.vlm.CANDIDATE_EVALUATION_PROMPT_VERSION +import com.roubao.autopilot.vlm.CANDIDATE_EVALUATION_SCHEMA_VERSION import com.roubao.autopilot.vlm.AndroidCandidateEvaluationVlmGateway import com.roubao.autopilot.vlm.CandidateBatchConclusion import com.roubao.autopilot.vlm.CandidateEvaluationFailureCode @@ -78,7 +82,14 @@ import com.roubao.autopilot.pinduoduo.PinduoduoSearchAutomation import com.roubao.autopilot.pinduoduo.CandidateEvidenceSource import com.roubao.autopilot.pinduoduo.SEARCH_PROBE_KEYWORD import com.roubao.autopilot.procurement.LoginInput +import com.roubao.autopilot.procurement.ExecutionCandidateBatchDraft +import com.roubao.autopilot.procurement.ExecutionCandidateDraft +import com.roubao.autopilot.procurement.ExecutionCandidateEvaluation +import com.roubao.autopilot.procurement.ExecutionEvidenceDraft +import com.roubao.autopilot.procurement.ExecutionMode +import com.roubao.autopilot.procurement.ExecutionProvenanceSnapshot import com.roubao.autopilot.procurement.ProcurementRepository +import com.roubao.autopilot.task.RequirementProbeFixture import com.roubao.autopilot.workflow.WorkflowReport import com.roubao.autopilot.workflow.WorkflowRunner import com.roubao.autopilot.workflow.WorkflowState @@ -123,6 +134,8 @@ class MainActivity : ComponentActivity() { private val candidateReviewBatch = mutableStateOf(null) private val candidateEvaluationFailureCode = mutableStateOf(null) + private val taskCandidateDrafts = + mutableStateOf>(emptyList()) private var searchProbeRunner: WorkflowRunner? = null private var searchProbeJob: Job? = null private var requirementProbeJob: Job? = null @@ -354,9 +367,9 @@ class MainActivity : ComponentActivity() { procurementRepository.claimNext(readiness) } }, - onStart = { + onStart = { mode -> lifecycleScope.launch { - procurementRepository.start() + startProcurementExecution(mode) } }, onRelease = { @@ -411,8 +424,12 @@ class MainActivity : ComponentActivity() { onStopCandidateEvaluation = { stopCandidateEvaluation() }, - onAcceptCandidate = { acceptRecommendedCandidate() }, - onRejectCandidates = { rejectCandidateReview() }, + onAcceptCandidate = { reason -> + acceptRecommendedCandidate(reason) + }, + onRejectCandidates = { reason -> + rejectCandidateReview(reason) + }, onStart = { startSearchProbe() }, onStop = { stopSearchProbe() } ) @@ -502,6 +519,43 @@ class MainActivity : ComponentActivity() { startActivity(launchIntent) } + private suspend fun startProcurementExecution(mode: ExecutionMode) { + if (mode == ExecutionMode.MANUAL_FIRST) { + procurementRepository.start(mode) + return + } + val settings = settingsManager.settings.value + val provider = settings.currentProvider + if (!provider.supportsRequirementExtraction || + settings.apiKey.isBlank() || + settings.baseUrl.isBlank() || + settings.model.isBlank() || + !settingsManager.isSecureCredentialStorageAvailable || + !RequirementProviderEndpointPolicy.isAllowed( + settings.baseUrl, + settings.apiKey + ) + ) { + Toast.makeText( + this, + "AI 辅助模式需要本机安全配置可用的 VLM", + Toast.LENGTH_SHORT + ).show() + return + } + procurementRepository.start( + mode = mode, + provenance = ExecutionProvenanceSnapshot( + mode = mode, + providerId = provider.id, + model = settings.model, + promptVersion = REQUIREMENT_PROMPT_VERSION, + schemaVersion = REQUIREMENT_SCHEMA_VERSION, + referenceImageSha256 = "" + ) + ) + } + private fun startSearchProbe() { val readiness = readinessChecker.snapshot() readinessSnapshot.value = readiness @@ -517,11 +571,22 @@ class MainActivity : ComponentActivity() { return } + val procurementTask = procurementRepository.currentProbeTask() + val procurementMode = procurementRepository.activeExecutionMode() val boundRequirement = requirementExtraction.value?.takeIf { requirementProbeState.value == RequirementProbeState.READY && !it.manualReviewRequired } - val searchKeyword = boundRequirement?.searchQuery ?: SEARCH_PROBE_KEYWORD + if (procurementTask != null && + procurementMode == ExecutionMode.AI_ASSISTED && + boundRequirement == null + ) { + Toast.makeText(this, "请先完成本地需求提取", Toast.LENGTH_SHORT).show() + return + } + val searchKeyword = boundRequirement?.searchQuery ?: procurementTask + ?.let { manualSearchQuery(it.title, it.sku) } + ?: SEARCH_PROBE_KEYWORD val candidateAutomation = PinduoduoCandidateAutomation( AndroidPinduoduoCandidateDriver(this) ) @@ -567,6 +632,16 @@ class MainActivity : ComponentActivity() { val report = runner.run(PinduoduoCandidateWorkflow.steps()) searchProbeReport.value = report searchProbeState.value = report.state + if (report.state == WorkflowState.SUCCEEDED && + procurementTask != null && + procurementMode == ExecutionMode.MANUAL_FIRST + ) { + queueManualProcurementCandidates( + procurementTask.title, + searchKeyword, + candidateEvidence.value + ) + } } finally { stateCollector.cancel() stepCollector.cancel() @@ -594,6 +669,16 @@ class MainActivity : ComponentActivity() { Toast.makeText(this, "请先停止候选评估", Toast.LENGTH_SHORT).show() return } + if (procurementRepository.currentProbeTask() != null && + procurementRepository.activeExecutionMode() != ExecutionMode.AI_ASSISTED + ) { + Toast.makeText( + this, + "当前后台任务为人工优先模式,不调用 VLM", + Toast.LENGTH_SHORT + ).show() + return + } candidateEvidence.value = emptyList() candidateSearchKeyword.value = null candidateRequirementSnapshot.value = null @@ -640,7 +725,11 @@ class MainActivity : ComponentActivity() { requirementFailureCode.value = null requirementProbeJob = lifecycleScope.launch { try { - val fixture = requirementProbeSource.loadFirst().getOrElse { + val fixture = procurementRepository.currentProbeTask()?.let { task -> + procurementRepository.currentReferenceImageBytes()?.let { image -> + RequirementProbeFixture(task, image) + } + } ?: requirementProbeSource.loadFirst().getOrElse { setRequirementFailure(RequirementExtractionFailureCode.SOURCE_UNAVAILABLE) return@launch } ?: run { @@ -809,6 +898,48 @@ class MainActivity : ComponentActivity() { ) when (result) { is CandidateEvaluationResult.Completed -> { + if (procurementRepository.activeExecutionMode() == + ExecutionMode.AI_ASSISTED + ) { + val drafts = result.batch.assessments.map { assessment -> + ExecutionCandidateDraft( + ordinal = assessment.ordinal, + title = "拼多多候选 ${assessment.ordinal}", + evidenceLocalIDs = emptyList(), + evaluation = ExecutionCandidateEvaluation( + decision = assessment.decision.name, + score = assessment.score, + matched = assessment.matched, + missingOrUncertain = assessment.missingOrUncertain, + rejectionReasons = assessment.rejectionReasons, + confidence = assessment.confidence + ) + ) + } + taskCandidateDrafts.value = procurementRepository.queueCandidateBatch( + ExecutionCandidateBatchDraft( + mode = ExecutionMode.AI_ASSISTED, + searchQuery = requirement.searchQuery, + provenance = ExecutionProvenanceSnapshot( + mode = ExecutionMode.AI_ASSISTED, + providerId = result.batch.providerId, + model = result.batch.model, + promptVersion = CANDIDATE_EVALUATION_PROMPT_VERSION, + schemaVersion = CANDIDATE_EVALUATION_SCHEMA_VERSION, + referenceImageSha256 = + result.batch.requirementReferenceImageSha256 + ), + candidates = drafts + ), + validated.map { candidate -> + ExecutionEvidenceDraft( + ordinal = candidate.ordinal, + pngBytes = candidate.pngBytes, + sha256 = candidate.sha256 + ) + } + ).orEmpty() + } candidateReviewBatch.value = result.batch candidateEvaluationState.value = when ( result.batch.conclusion @@ -838,17 +969,90 @@ class MainActivity : ComponentActivity() { candidateEvaluationJob?.cancel() } - private fun acceptRecommendedCandidate() { + private suspend fun queueManualProcurementCandidates( + taskTitle: String, + searchQuery: String, + evidence: List + ) { + val validated = candidateEvidenceSource.load(evidence).getOrNull() ?: return + taskCandidateDrafts.value = procurementRepository.queueCandidateBatch( + ExecutionCandidateBatchDraft( + mode = ExecutionMode.MANUAL_FIRST, + searchQuery = searchQuery, + provenance = null, + candidates = validated.map { candidate -> + ExecutionCandidateDraft( + ordinal = candidate.ordinal, + title = "$taskTitle 候选 ${candidate.ordinal}", + evidenceLocalIDs = emptyList() + ) + } + ), + validated.map { candidate -> + ExecutionEvidenceDraft( + ordinal = candidate.ordinal, + pngBytes = candidate.pngBytes, + sha256 = candidate.sha256 + ) + } + ).orEmpty() + if (taskCandidateDrafts.value.isNotEmpty()) { + candidateEvaluationState.value = CandidateEvaluationState.MANUAL_REVIEW + } + } + + private fun manualSearchQuery(title: String, sku: String): String = + listOf(title.trim(), sku.trim()) + .filter { it.isNotBlank() } + .joinToString(" ") + .take(160) + + private fun acceptRecommendedCandidate(operatorReason: String) { candidateEvaluationState.value = CandidateHumanReviewPolicy.accept( currentState = candidateEvaluationState.value, batch = candidateReviewBatch.value ) + if (candidateEvaluationState.value != CandidateEvaluationState.HUMAN_ACCEPTED && + taskCandidateDrafts.value.isNotEmpty() + ) { + candidateEvaluationState.value = CandidateEvaluationState.HUMAN_ACCEPTED + } + if (candidateEvaluationState.value == CandidateEvaluationState.HUMAN_ACCEPTED) { + val recommendedOrdinal = candidateReviewBatch.value + ?.recommendedCandidateOrdinal + val candidate = taskCandidateDrafts.value.firstOrNull { + it.ordinal == recommendedOrdinal + } ?: taskCandidateDrafts.value.firstOrNull() + if (candidate != null && procurementRepository.currentProbeTask() != null) { + lifecycleScope.launch { + procurementRepository.completeExecution( + outcome = "CANDIDATE_ACCEPTED", + operatorReason = operatorReason, + candidate = candidate + ) + } + } + } } - private fun rejectCandidateReview() { + private fun rejectCandidateReview(operatorReason: String) { candidateEvaluationState.value = CandidateHumanReviewPolicy.reject( candidateEvaluationState.value ) + if (candidateEvaluationState.value == CandidateEvaluationState.HUMAN_REJECTED && + procurementRepository.currentProbeTask() != null + ) { + lifecycleScope.launch { + procurementRepository.completeExecution( + outcome = if (taskCandidateDrafts.value.isEmpty()) { + "NO_MATCH" + } else { + "CANDIDATE_REJECTED" + }, + operatorReason = operatorReason + ) + } + } } private fun setCandidateEvaluationFailure( diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/data/SettingsManager.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/data/SettingsManager.kt index 8c1bcdb..59066b5 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/data/SettingsManager.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/data/SettingsManager.kt @@ -127,32 +127,31 @@ class SettingsManager(context: Context) { context.getSharedPreferences("baozi_settings", Context.MODE_PRIVATE) // 加密存储(用于敏感数据如 API Key) - private val securePrefs: SharedPreferences by lazy { - try { - val masterKey = MasterKey.Builder(context) - .setKeyScheme(MasterKey.KeyScheme.AES256_GCM) - .build() + private val securePrefs: SharedPreferences? = try { + val masterKey = MasterKey.Builder(context.applicationContext) + .setKeyScheme(MasterKey.KeyScheme.AES256_GCM) + .build() - EncryptedSharedPreferences.create( - context, - "baozi_secure_settings", - masterKey, - EncryptedSharedPreferences.PrefKeyEncryptionScheme.AES256_SIV, - EncryptedSharedPreferences.PrefValueEncryptionScheme.AES256_GCM - ) - } catch (e: Exception) { - secureCredentialStorageAvailable = false - android.util.Log.e("SettingsManager", "Failed to create encrypted prefs", e) - prefs - } + EncryptedSharedPreferences.create( + context.applicationContext, + "baozi_secure_settings", + masterKey, + EncryptedSharedPreferences.PrefKeyEncryptionScheme.AES256_SIV, + EncryptedSharedPreferences.PrefValueEncryptionScheme.AES256_GCM + ) + } catch (e: Exception) { + secureCredentialStorageAvailable = false + android.util.Log.e("SettingsManager", "Failed to create encrypted prefs", e) + null } - private val _settings = MutableStateFlow(loadSettings()) - val settings: StateFlow = _settings + private val _settings: MutableStateFlow + val settings: StateFlow init { - // 迁移旧的明文 API Key 到加密存储 migrateApiKeyToSecureStorage() + _settings = MutableStateFlow(loadSettings()) + settings = _settings } val isSecureCredentialStorageAvailable: Boolean @@ -163,13 +162,19 @@ class SettingsManager(context: Context) { */ private fun migrateApiKeyToSecureStorage() { val oldApiKey = prefs.getString("api_key", null) - if (!oldApiKey.isNullOrEmpty()) { - // 保存到加密存储 - securePrefs.edit().putString("api_key", oldApiKey).apply() - // 删除旧的明文存储 - prefs.edit().remove("api_key").apply() - android.util.Log.d("SettingsManager", "API Key migrated to secure storage") + if (oldApiKey.isNullOrEmpty()) return + val encrypted = securePrefs + if (encrypted == null || !encrypted.edit().putString("api_key", oldApiKey).commit()) { + secureCredentialStorageAvailable = false + android.util.Log.e("SettingsManager", "API Key migration failed closed") + return } + if (!prefs.edit().remove("api_key").commit()) { + secureCredentialStorageAvailable = false + android.util.Log.e("SettingsManager", "Plaintext API Key cleanup failed closed") + return + } + android.util.Log.d("SettingsManager", "API Key migrated to secure storage") } private fun loadSettings(): AppSettings { @@ -191,7 +196,7 @@ class SettingsManager(context: Context) { } // 迁移旧数据(如果有) - val oldApiKey = securePrefs.getString("api_key", null) + val oldApiKey = securePrefs?.getString("api_key", null) val oldModel = prefs.getString("model", null) val oldBaseUrl = prefs.getString("base_url", null) val oldCachedModels = prefs.getStringSet("cached_models", null) @@ -205,24 +210,29 @@ class SettingsManager(context: Context) { else -> "custom" } - // 迁移到新格式 val migratedConfig = ProviderConfig( apiKey = oldApiKey ?: "", model = oldModel ?: "", cachedModels = oldCachedModels?.toList() ?: emptyList(), customBaseUrl = if (oldProviderId == "custom") oldBaseUrl ?: "" else "" ) - providerConfigs[oldProviderId] = migratedConfig - saveProviderConfig(oldProviderId, migratedConfig) + if (saveProviderConfig(oldProviderId, migratedConfig)) { + providerConfigs[oldProviderId] = migratedConfig + } - // 清除旧数据 - securePrefs.edit().remove("api_key").apply() - prefs.edit() + val secureMigrationComplete = securePrefs?.edit() + ?.remove("api_key") + ?.commit() == true + val preferencesMigrationComplete = prefs.edit() .remove("model") .remove("base_url") .remove("cached_models") .putString("current_provider_id", oldProviderId) - .apply() + .commit() + if (!secureMigrationComplete || !preferencesMigrationComplete) { + secureCredentialStorageAvailable = false + android.util.Log.e("SettingsManager", "Legacy provider migration failed closed") + } android.util.Log.d("SettingsManager", "Migrated old settings to provider: $oldProviderId") } @@ -244,7 +254,7 @@ class SettingsManager(context: Context) { private fun loadProviderConfig(providerId: String): ProviderConfig { val prefix = "provider_${providerId}_" return ProviderConfig( - apiKey = securePrefs.getString("${prefix}api_key", "") ?: "", + apiKey = securePrefs?.getString("${prefix}api_key", "") ?: "", model = prefs.getString("${prefix}model", "") ?: "", cachedModels = prefs.getStringSet("${prefix}cached_models", emptySet())?.toList() ?: emptyList(), customBaseUrl = prefs.getString("${prefix}custom_base_url", "") ?: "" @@ -254,14 +264,30 @@ class SettingsManager(context: Context) { /** * 保存指定服务商的配置 */ - private fun saveProviderConfig(providerId: String, config: ProviderConfig) { + private fun saveProviderConfig(providerId: String, config: ProviderConfig): Boolean { val prefix = "provider_${providerId}_" - securePrefs.edit().putString("${prefix}api_key", config.apiKey).apply() - prefs.edit() + if (config.apiKey.isNotEmpty()) { + val encrypted = securePrefs + if (encrypted == null || + !encrypted.edit().putString("${prefix}api_key", config.apiKey).commit() + ) { + secureCredentialStorageAvailable = false + android.util.Log.e("SettingsManager", "Refused API Key write without encryption") + return false + } + } else { + if (securePrefs != null && + !securePrefs.edit().remove("${prefix}api_key").commit() + ) { + secureCredentialStorageAvailable = false + return false + } + } + return prefs.edit() .putString("${prefix}model", config.model) .putStringSet("${prefix}cached_models", config.cachedModels.toSet()) .putString("${prefix}custom_base_url", config.customBaseUrl) - .apply() + .commit() } /** @@ -272,7 +298,9 @@ class SettingsManager(context: Context) { val currentConfig = _settings.value.currentConfig val newConfig = update(currentConfig) - saveProviderConfig(currentId, newConfig) + if (!saveProviderConfig(currentId, newConfig)) { + return + } val newConfigs = _settings.value.providerConfigs.toMutableMap() newConfigs[currentId] = newConfig diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ExecutionResultOutbox.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ExecutionResultOutbox.kt new file mode 100644 index 0000000..cc011a9 --- /dev/null +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ExecutionResultOutbox.kt @@ -0,0 +1,79 @@ +package com.roubao.autopilot.procurement + +import java.security.MessageDigest + +data class ExecutionEvidenceDraft( + val ordinal: Int, + val pngBytes: ByteArray, + val sha256: String +) + +data class ExecutionCandidateDraft( + val ordinal: Int, + val title: String, + val skuText: String = "", + val price: String = "", + val productUrl: String = "", + val imageUrl: String = "", + val evidenceLocalIDs: List, + val evaluation: ExecutionCandidateEvaluation? = null +) + +data class ExecutionCandidateEvaluation( + val decision: String, + val score: Double, + val matched: List, + val missingOrUncertain: List, + val rejectionReasons: List, + val confidence: Double +) + +data class ExecutionRecommendation( + val candidateOrdinal: Int, + val policyVersion: String, + val reasons: List +) + +data class ExecutionCandidateBatchDraft( + val mode: ExecutionMode, + val searchQuery: String, + val provenance: ExecutionProvenanceSnapshot?, + val candidates: List, + val recommendation: ExecutionRecommendation? = null +) + +object ExecutionTaskHash { + fun sha256(task: RemotePurchaseTask): String { + val maxBudgetCents = task.maxBudget + ?.let(::parseCnyCents) + ?.toString() + .orEmpty() + val payload = listOf( + task.title, + task.description, + task.sku, + task.imageAssetId, + task.quantity.toString(), + maxBudgetCents, + task.currency + ).joinToString("\u0000") + return MessageDigest.getInstance("SHA-256") + .digest(payload.toByteArray(Charsets.UTF_8)) + .joinToString("") { "%02x".format(it) } + } + + private fun parseCnyCents(value: String): Long { + val normalized = value.trim() + val pieces = normalized.split('.', limit = 2) + require(pieces.size in 1..2 && pieces[0].all(Char::isDigit)) + val whole = pieces[0].toLong() + val fractional = when (pieces.size) { + 1 -> 0L + else -> { + require(pieces[1].length in 1..2 && pieces[1].all(Char::isDigit)) + pieces[1].padEnd(2, '0').toLong() + } + } + return whole * 100 + fractional + } +} diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementApiClient.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementApiClient.kt index 7ae8bc2..c3ca221 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementApiClient.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementApiClient.kt @@ -25,6 +25,10 @@ data class DownloadedReferenceImage( val sha256: String ) +data class ExecutionOutboxUploadResult( + val evidenceAssetId: String? = null +) + interface ProcurementRemoteApi { suspend fun login( baseUrl: String, @@ -79,6 +83,15 @@ interface ProcurementRemoteApi { claimToken: String, idempotencyKey: String ): RemotePurchaseTask + + suspend fun uploadExecutionOutboxItem( + session: ProcurementSession, + task: RemotePurchaseTask, + execution: RunningExecution, + claimToken: String, + item: ExecutionOutboxItem, + evidenceBytes: ByteArray? = null + ): ExecutionOutboxUploadResult } class ProcurementApiClient( @@ -250,6 +263,50 @@ class ProcurementApiClient( } } + override suspend fun uploadExecutionOutboxItem( + session: ProcurementSession, + task: RemotePurchaseTask, + execution: RunningExecution, + claimToken: String, + item: ExecutionOutboxItem, + evidenceBytes: ByteArray? + ): ExecutionOutboxUploadResult = withContext(Dispatchers.IO) { + val path = when (item.type) { + ExecutionOutboxType.EVENTS -> "/api/v1/tasks/${task.id}/events" + ExecutionOutboxType.EVIDENCE -> "/api/v1/tasks/${task.id}/evidence" + ExecutionOutboxType.CANDIDATES -> "/api/v1/tasks/${task.id}/candidates" + ExecutionOutboxType.COMPLETE -> "/api/v1/tasks/${task.id}/complete" + ExecutionOutboxType.FAIL -> "/api/v1/tasks/${task.id}/fail" + } + val request = authorizedRequest(session, path) + .header(CLAIM_TOKEN_HEADER, claimToken) + .header(IDEMPOTENCY_HEADER, item.idempotencyKey) + if (item.type == ExecutionOutboxType.EVIDENCE) { + val bytes = requireNotNull(evidenceBytes) { "离线证据文件缺失" } + require(bytes.isNotEmpty() && bytes.size <= MAX_EVIDENCE_BYTES) { + "离线证据文件无效" + } + request + .header("X-Execution-ID", execution.id) + .header("X-Claim-Generation", task.claimGeneration.toString()) + .post(bytes.toRequestBody(PNG_MEDIA_TYPE)) + } else { + val payload = item.payload + require(payload.toByteArray(Charsets.UTF_8).size <= MAX_JSON_BYTES) { + "离线结果过大" + } + request.post(payload.toRequestBody(JSON_MEDIA_TYPE)) + } + val json = executeJson(request.build()) + ExecutionOutboxUploadResult( + evidenceAssetId = if (item.type == ExecutionOutboxType.EVIDENCE) { + json.getJSONObject("evidence").getString("id") + } else { + null + } + ) + } + override suspend fun start( session: ProcurementSession, task: RemotePurchaseTask, @@ -455,6 +512,7 @@ class ProcurementApiClient( title = json.getString("title"), description = json.optString("description"), sku = json.getString("sku"), + imageAssetId = json.getString("image_asset_id"), referenceImageUrl = json.getString("reference_image_url"), quantity = json.getInt("quantity"), maxBudget = json.optionalString("max_budget"), @@ -501,10 +559,12 @@ class ProcurementApiClient( companion object { private val JSON_MEDIA_TYPE = "application/json; charset=utf-8".toMediaType() + private val PNG_MEDIA_TYPE = "image/png".toMediaType() private const val JPEG_MEDIA_TYPE = "image/jpeg" private const val CLAIM_TOKEN_HEADER = "X-Claim-Token" private const val IDEMPOTENCY_HEADER = "Idempotency-Key" private const val MAX_JSON_BYTES = 1_048_576L + private const val MAX_EVIDENCE_BYTES = 8 * 1024 * 1024 private const val MAX_REFERENCE_IMAGE_BYTES = 20L * 1024L * 1024L private val SHA256_PATTERN = Regex("[0-9a-f]{64}") diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementModels.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementModels.kt index 9646c8c..81fbed0 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementModels.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementModels.kt @@ -39,6 +39,7 @@ data class RemotePurchaseTask( val title: String, val description: String, val sku: String, + val imageAssetId: String, val referenceImageUrl: String, val quantity: Int, val maxBudget: String?, @@ -67,7 +68,8 @@ data class RunningExecution( val expiresAt: String, val serverClockOffsetMillis: Long, val safetyStopped: Boolean = false, - val cancelAcknowledgementKey: String? = null + val cancelAcknowledgementKey: String? = null, + val provenance: ExecutionProvenanceSnapshot? = null ) { fun isExpired(nowEpochMillis: Long = System.currentTimeMillis()): Boolean = ExecutionAuthorization.isExpired( @@ -81,7 +83,39 @@ data class PersistedProcurementState( val session: ProcurementSession? = null, val deviceToken: String? = null, val claim: ClaimContext? = null, - val execution: RunningExecution? = null + val execution: RunningExecution? = null, + val outbox: List = emptyList() +) + +enum class ExecutionMode { + MANUAL_FIRST, + AI_ASSISTED +} + +data class ExecutionProvenanceSnapshot( + val mode: ExecutionMode, + val providerId: String? = null, + val model: String? = null, + val promptVersion: String? = null, + val schemaVersion: Int? = null, + val referenceImageSha256: String +) + +enum class ExecutionOutboxType { + EVENTS, + EVIDENCE, + CANDIDATES, + COMPLETE, + FAIL +} + +data class ExecutionOutboxItem( + val id: String, + val type: ExecutionOutboxType, + val idempotencyKey: String, + val payload: String, + val evidenceRelativePath: String? = null, + val remoteResourceID: String? = null ) data class ProcurementUiState( diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementRepository.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementRepository.kt index 657e0de..8811e41 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementRepository.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementRepository.kt @@ -10,10 +10,14 @@ import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock +import org.json.JSONArray +import org.json.JSONObject import java.io.File import java.io.IOException import java.security.SecureRandom +import java.time.Instant import java.util.Base64 +import java.util.UUID enum class ExecutionSyncDecision { CONTINUE, @@ -131,7 +135,10 @@ class ProcurementRepository( true } ?: false - suspend fun start(): Boolean = operation { + suspend fun start( + mode: ExecutionMode = ExecutionMode.MANUAL_FIRST, + provenance: ExecutionProvenanceSnapshot? = null + ): Boolean = operation { val session = requireValidSession() val claim = requireNotNull(persisted.claim) { "没有可开始的任务" } val task = requireNotNull(claim.task) { "任务详情尚未下载" } @@ -160,12 +167,20 @@ class ProcurementRepository( serverClockOffsetMillis = ExecutionAuthorization.serverClockOffset( result.serverTime, requestStartedAt - ) + ), + provenance = normalizeProvenance(mode, provenance, claim.referenceImage) ) persisted = persisted.copy( claim = persisted.claim?.copy(task = result.task), execution = execution ) + enqueueEventLocked( + task = result.task, + execution = execution, + step = "PREFLIGHT", + type = "EXECUTION_STARTED", + message = "本地受控采购流程已开始" + ) saveAndPublish( backendOnline = true, message = "受控采购流程已开始" @@ -234,6 +249,21 @@ class ProcurementRepository( } } return@withLock try { + if (flushExecutionOutbox(session)) { + saveAndPublish( + backendOnline = true, + message = "执行结果已回传,任务已结束" + ) + return@withLock ExecutionSyncDecision.STOP + } + execution = persisted.execution ?: return@withLock ExecutionSyncDecision.STOP + if (execution.safetyStopped) { + publish( + backendOnline = true, + message = "授权到期后的审计结果已同步,自动化保持停止" + ) + return@withLock ExecutionSyncDecision.STOP + } val requestStartedAt = System.currentTimeMillis() val result = api.heartbeat(session, task, execution, claim.token) val updatedExecution = execution.copy( @@ -294,6 +324,172 @@ class ProcurementRepository( if (task != null && image != null) task.toProbeTask(image) else null } + fun activeExecutionMode(): ExecutionMode? = + persisted.execution?.provenance?.mode + + fun currentReferenceImageBytes(): ByteArray? = + persisted.claim?.referenceImage?.let { image -> + File(appContext.filesDir, image.relativePath) + .takeIf { it.isFile && it.length() == image.sizeBytes } + ?.readBytes() + } + + suspend fun queueCandidateBatch( + batch: ExecutionCandidateBatchDraft, + evidence: List + ): List? = operation { + val task = requireCurrentTask() + val execution = requireActiveExecution() + require(!execution.safetyStopped && !execution.isExpired()) { + "执行授权已到期,不能继续采集" + } + require(evidence.size in 1..5) { "候选证据必须为 1 至 5 张" } + require(batch.candidates.size in 1..5) { "候选数量必须为 1 至 5 个" } + require(batch.candidates.map { it.ordinal } == (1..batch.candidates.size).toList()) { + "候选编号必须连续" + } + require(batch.mode == execution.provenance?.mode) { "执行模式不能在回传时变更" } + if (batch.mode == ExecutionMode.AI_ASSISTED) { + require(batch.provenance != null && batch.provenance.providerId != null) { + "AI 辅助模式缺少本地模型出处" + } + } else { + require(batch.provenance == null) { "手工优先模式不能附带模型出处" } + } + val evidenceIDs = evidence.sortedBy { it.ordinal }.associate { draft -> + require(draft.ordinal in 1..batch.candidates.size) + require(sha256(draft.pngBytes) == draft.sha256) + val localID = UUID.randomUUID().toString() + persistEvidenceLocked(localID, draft.pngBytes) + appendOutboxLocked( + ExecutionOutboxItem( + id = localID, + type = ExecutionOutboxType.EVIDENCE, + idempotencyKey = newOpaqueSecret(), + payload = "{}", + evidenceRelativePath = "$OUTBOX_DIRECTORY/$localID.png" + ) + ) + draft.ordinal to localID + } + val candidates = batch.candidates.map { candidate -> + candidate.copy( + evidenceLocalIDs = listOf( + requireNotNull(evidenceIDs[candidate.ordinal]) + ) + ) + } + val payload = JSONObject() + .put("execution_id", execution.id) + .put("claim_generation", task.claimGeneration) + .put("task_content_sha256", ExecutionTaskHash.sha256(task)) + .put("execution_mode", batch.mode.name) + .put("search_query", batch.searchQuery.trim()) + .put("candidates", JSONArray().apply { + candidates.forEach { put(candidateJSON(it)) } + }) + batch.provenance?.let { payload.put("provenance", provenanceJSON(it)) } + batch.recommendation?.let { recommendation -> + payload.put( + "recommendation", + JSONObject() + .put("candidate_ordinal", recommendation.candidateOrdinal) + .put("policy_version", recommendation.policyVersion) + .put("reasons", JSONArray(recommendation.reasons)) + ) + } + appendOutboxLocked( + ExecutionOutboxItem( + id = UUID.randomUUID().toString(), + type = ExecutionOutboxType.CANDIDATES, + idempotencyKey = newOpaqueSecret(), + payload = payload.toString() + ) + ) + saveAndPublish( + backendOnline = _uiState.value.backendOnline, + message = "候选和截图已加入加密回传队列" + ) + candidates + } + + suspend fun completeExecution( + outcome: String, + operatorReason: String, + candidate: ExecutionCandidateDraft? = null + ): Boolean = operation { + val task = requireCurrentTask() + val execution = requireActiveExecution() + require(operatorReason.trim().isNotEmpty()) { "请填写人工确认理由" } + require(outcome in COMPLETE_OUTCOMES) { "采购结论无效" } + if (outcome == "CANDIDATE_ACCEPTED") { + require(candidate != null) { "接受候选时必须保留候选证据" } + } else { + require(candidate == null) { "当前结论不能附带候选" } + } + val payload = JSONObject() + .put("execution_id", execution.id) + .put("claim_generation", task.claimGeneration) + .put("task_content_sha256", ExecutionTaskHash.sha256(task)) + .put("execution_mode", requireNotNull(execution.provenance).mode.name) + .put("outcome", outcome) + .put("operator_reason", operatorReason.trim()) + .put("order_submitted", false) + candidate?.let { payload.put("candidate", candidateJSON(it)) } + appendTerminalOutboxLocked( + ExecutionOutboxItem( + id = UUID.randomUUID().toString(), + type = ExecutionOutboxType.COMPLETE, + idempotencyKey = newOpaqueSecret(), + payload = payload.toString() + ) + ) + saveAndPublish( + backendOnline = _uiState.value.backendOnline, + message = "人工结论已加入加密回传队列" + ) + true + } ?: false + + suspend fun failExecution( + code: String, + message: String, + step: String, + retryable: Boolean, + evidenceLocalIDs: List = emptyList() + ): Boolean = operation { + val task = requireCurrentTask() + val execution = requireActiveExecution() + require(code.isNotBlank() && message.isNotBlank() && step.isNotBlank()) { + "失败信息不完整" + } + val payload = JSONObject() + .put("execution_id", execution.id) + .put("claim_generation", task.claimGeneration) + .put( + "error", + JSONObject() + .put("code", code.trim()) + .put("message", message.trim()) + .put("step", step.trim()) + .put("retryable", retryable) + ) + .put("evidence_asset_ids", JSONArray(evidenceLocalIDs)) + appendTerminalOutboxLocked( + ExecutionOutboxItem( + id = UUID.randomUUID().toString(), + type = ExecutionOutboxType.FAIL, + idempotencyKey = newOpaqueSecret(), + payload = payload.toString() + ) + ) + saveAndPublish( + backendOnline = _uiState.value.backendOnline, + message = "失败信息已加入加密回传队列" + ) + true + } ?: false + private suspend fun acknowledgeCancellation(session: ProcurementSession) { val claim = requireNotNull(persisted.claim) val task = requireNotNull(claim.task) @@ -363,9 +559,252 @@ class ProcurementRepository( persisted.claim?.referenceImage?.let { File(appContext.filesDir, it.relativePath).delete() } - persisted = persisted.copy(claim = null, execution = null) + persisted.outbox.forEach { item -> + item.evidenceRelativePath?.let { relativePath -> + File(appContext.filesDir, relativePath).delete() + } + } + persisted = persisted.copy( + claim = null, + execution = null, + outbox = emptyList() + ) } + private fun requireCurrentTask(): RemotePurchaseTask = + requireNotNull(persisted.claim?.task) { "没有可回传的采购任务" } + + private fun requireActiveExecution(): RunningExecution = + requireNotNull(persisted.execution) { "采购执行尚未开始" } + + private fun appendOutboxLocked(item: ExecutionOutboxItem) { + var resolvedPayload = item.payload + persisted.outbox.forEach { knownEvidence -> + if (knownEvidence.type == ExecutionOutboxType.EVIDENCE && + knownEvidence.remoteResourceID != null + ) { + resolvedPayload = replaceEvidenceReference( + resolvedPayload, + knownEvidence.id, + knownEvidence.remoteResourceID + ) + } + } + persisted = persisted.copy( + outbox = persisted.outbox + item.copy(payload = resolvedPayload) + ) + } + + private fun appendTerminalOutboxLocked(item: ExecutionOutboxItem) { + require(persisted.outbox.none { + it.type == ExecutionOutboxType.COMPLETE || it.type == ExecutionOutboxType.FAIL + }) { "当前执行已有待回传终态" } + appendOutboxLocked(item) + } + + private fun enqueueEventLocked( + task: RemotePurchaseTask, + execution: RunningExecution, + step: String, + type: String, + message: String + ) { + val payload = JSONObject() + .put("execution_id", execution.id) + .put("claim_generation", task.claimGeneration) + .put( + "events", + JSONArray().put( + JSONObject() + .put("event_id", UUID.randomUUID().toString()) + .put("step", step) + .put("type", type) + .put("message", message) + .put("occurred_at", Instant.now().toString()) + ) + ) + appendOutboxLocked( + ExecutionOutboxItem( + id = UUID.randomUUID().toString(), + type = ExecutionOutboxType.EVENTS, + idempotencyKey = newOpaqueSecret(), + payload = payload.toString() + ) + ) + } + + private fun normalizeProvenance( + mode: ExecutionMode, + candidate: ExecutionProvenanceSnapshot?, + referenceImage: ReferenceImageRecord? + ): ExecutionProvenanceSnapshot { + val referenceHash = requireNotNull(referenceImage?.sha256) { + "参考图哈希缺失" + } + if (mode == ExecutionMode.MANUAL_FIRST) { + require(candidate == null) { "手工优先模式不能配置模型出处" } + return ExecutionProvenanceSnapshot( + mode = mode, + referenceImageSha256 = referenceHash + ) + } + val provenance = requireNotNull(candidate) { "AI 辅助模式缺少模型出处" } + require( + provenance.mode == mode && + !provenance.providerId.isNullOrBlank() && + !provenance.model.isNullOrBlank() && + !provenance.promptVersion.isNullOrBlank() && + provenance.schemaVersion != null + ) { "AI 辅助模式模型出处无效" } + return provenance.copy(referenceImageSha256 = referenceHash) + } + + private fun persistEvidenceLocked(localID: String, bytes: ByteArray) { + require(bytes.isNotEmpty() && bytes.size <= MAX_OUTBOX_EVIDENCE_BYTES) { + "候选截图无效" + } + val directory = File(appContext.filesDir, OUTBOX_DIRECTORY) + check(directory.exists() || directory.mkdirs()) { "无法创建离线证据目录" } + val target = File(directory, "$localID.png") + val temporary = File(directory, ".$localID.tmp") + temporary.outputStream().use { output -> + output.write(bytes) + output.fd.sync() + } + check(!target.exists() || target.delete()) { "无法替换离线证据" } + check(temporary.renameTo(target)) { "无法保存离线证据" } + } + + private suspend fun flushExecutionOutbox(session: ProcurementSession): Boolean { + while (true) { + val task = requireCurrentTask() + val execution = requireActiveExecution() + val item = persisted.outbox.firstOrNull { pending -> + pending.type != ExecutionOutboxType.EVIDENCE || + pending.remoteResourceID == null + } ?: return false + val evidenceBytes = item.evidenceRelativePath?.let { relativePath -> + val file = File(appContext.filesDir, relativePath) + require(file.isFile && file.length() in 1..MAX_OUTBOX_EVIDENCE_BYTES) { + "离线证据文件缺失" + } + file.readBytes() + } + val uploaded = api.uploadExecutionOutboxItem( + session = session, + task = task, + execution = execution, + claimToken = requireNotNull(persisted.claim).token, + item = item, + evidenceBytes = evidenceBytes + ) + if (item.type == ExecutionOutboxType.EVIDENCE) { + val evidenceID = requireNotNull(uploaded.evidenceAssetId) { + "后台未返回证据编号" + } + persisted = persisted.copy( + outbox = persisted.outbox.map { pending -> + if (pending.id == item.id) { + pending.copy(remoteResourceID = evidenceID) + } else { + pending.copy( + payload = replaceEvidenceReference( + pending.payload, + item.id, + evidenceID + ) + ) + } + } + ) + item.evidenceRelativePath?.let { + File(appContext.filesDir, it).delete() + } + } else { + persisted = persisted.copy( + outbox = persisted.outbox.filterNot { it.id == item.id } + ) + } + if (item.type == ExecutionOutboxType.COMPLETE || + item.type == ExecutionOutboxType.FAIL + ) { + clearCurrentTask() + store.save(persisted) + return true + } + store.save(persisted) + } + } + + private fun replaceEvidenceReference( + payload: String, + localID: String, + remoteID: String + ): String { + val root = JSONObject(payload) + fun replace(value: Any?) { + when (value) { + is JSONObject -> { + val keys = value.keys().asSequence().toList() + keys.forEach { key -> replace(value.opt(key)) } + } + is JSONArray -> { + for (index in 0 until value.length()) { + if (value.optString(index) == localID) { + value.put(index, remoteID) + } else { + replace(value.opt(index)) + } + } + } + } + } + replace(root) + return root.toString() + } + + private fun candidateJSON(candidate: ExecutionCandidateDraft): JSONObject = + JSONObject() + .put("ordinal", candidate.ordinal) + .put("title", candidate.title) + .put("sku_text", candidate.skuText) + .put("price", candidate.price) + .put("product_url", candidate.productUrl) + .put("image_url", candidate.imageUrl) + .put("evidence_asset_ids", JSONArray(candidate.evidenceLocalIDs)) + .also { json -> + candidate.evaluation?.let { evaluation -> + json.put( + "evaluation", + JSONObject() + .put("decision", evaluation.decision) + .put("score", evaluation.score) + .put("matched", JSONArray(evaluation.matched)) + .put( + "missing_or_uncertain", + JSONArray(evaluation.missingOrUncertain) + ) + .put( + "rejection_reasons", + JSONArray(evaluation.rejectionReasons) + ) + .put("confidence", evaluation.confidence) + ) + } + } + + private fun provenanceJSON(snapshot: ExecutionProvenanceSnapshot): JSONObject = + JSONObject() + .put("provider_id", snapshot.providerId) + .put("model", snapshot.model) + .put("prompt_version", snapshot.promptVersion) + .put("schema_version", snapshot.schemaVersion) + + private fun sha256(bytes: ByteArray): String = + java.security.MessageDigest.getInstance("SHA-256") + .digest(bytes) + .joinToString("") { "%02x".format(it) } + private fun requireValidSession(): ProcurementSession { val session = persisted.session require(session != null && session.isValid()) { "采购员登录已过期,请重新登录" } @@ -476,10 +915,18 @@ class ProcurementRepository( companion object { private const val REFERENCE_DIRECTORY = "procurement" + private const val OUTBOX_DIRECTORY = "procurement-outbox" private const val CONTROLLED_WORKFLOW_STEP = "CONTROLLED_WORKFLOW" private const val SAFE_STOPPED_STEP = "SAFE_STOPPED" private const val MAX_REFERENCE_DIMENSION = 4_096 private const val MAX_REFERENCE_PIXELS = 20_000_000L + private const val MAX_OUTBOX_EVIDENCE_BYTES = 8L * 1024L * 1024L + private val COMPLETE_OUTCOMES = setOf( + "CANDIDATE_ACCEPTED", + "CANDIDATE_REJECTED", + "NO_MATCH", + "MANUAL_REQUIRED" + ) private val SECURE_RANDOM = SecureRandom() fun create(context: Context): ProcurementRepository = diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementSecureStore.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementSecureStore.kt index 5e885bf..e3b61f5 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementSecureStore.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementSecureStore.kt @@ -3,6 +3,7 @@ package com.roubao.autopilot.procurement import android.content.Context import androidx.security.crypto.EncryptedSharedPreferences import androidx.security.crypto.MasterKey +import org.json.JSONArray import org.json.JSONObject interface ProcurementStateStore { @@ -89,9 +90,50 @@ class ProcurementSecureStore(context: Context) : ProcurementStateStore { "cancel_acknowledgement_key", execution.cancelAcknowledgementKey ) + execution.provenance?.let { provenance -> + put( + "provenance", + JSONObject().apply { + put("mode", provenance.mode.name) + putNullable("provider_id", provenance.providerId) + putNullable("model", provenance.model) + putNullable("prompt_version", provenance.promptVersion) + provenance.schemaVersion?.let { + put("schema_version", it) + } + put( + "reference_image_sha256", + provenance.referenceImageSha256 + ) + } + ) + } } ) } + put( + "outbox", + JSONArray().apply { + state.outbox.forEach { item -> + put( + JSONObject().apply { + put("id", item.id) + put("type", item.type.name) + put("idempotency_key", item.idempotencyKey) + put("payload", item.payload) + putNullable( + "evidence_relative_path", + item.evidenceRelativePath + ) + putNullable( + "remote_resource_id", + item.remoteResourceID + ) + } + ) + } + } + ) } private fun decodeState(json: JSONObject): PersistedProcurementState = @@ -137,9 +179,43 @@ class ProcurementSecureStore(context: Context) : ProcurementStateStore { serverClockOffsetMillis = it.getLong("server_clock_offset_ms"), safetyStopped = it.optBoolean("safety_stopped", false), cancelAcknowledgementKey = - it.optionalString("cancel_acknowledgement_key") + it.optionalString("cancel_acknowledgement_key"), + provenance = it.optionalObject("provenance")?.let { provenance -> + ExecutionProvenanceSnapshot( + mode = ExecutionMode.valueOf(provenance.getString("mode")), + providerId = provenance.optionalString("provider_id"), + model = provenance.optionalString("model"), + promptVersion = provenance.optionalString("prompt_version"), + schemaVersion = if (provenance.has("schema_version")) { + provenance.getInt("schema_version") + } else { + null + }, + referenceImageSha256 = + provenance.getString("reference_image_sha256") + ) + } ) - } + }, + outbox = json.optionalArray("outbox")?.let { entries -> + buildList { + for (index in 0 until entries.length()) { + val item = entries.getJSONObject(index) + add( + ExecutionOutboxItem( + id = item.getString("id"), + type = ExecutionOutboxType.valueOf(item.getString("type")), + idempotencyKey = item.getString("idempotency_key"), + payload = item.getString("payload"), + evidenceRelativePath = + item.optionalString("evidence_relative_path"), + remoteResourceID = + item.optionalString("remote_resource_id") + ) + ) + } + } + } ?: emptyList() ) private fun encodeTask(task: RemotePurchaseTask): JSONObject = @@ -152,6 +228,7 @@ class ProcurementSecureStore(context: Context) : ProcurementStateStore { put("title", task.title) put("description", task.description) put("sku", task.sku) + put("image_asset_id", task.imageAssetId) put("reference_image_url", task.referenceImageUrl) put("quantity", task.quantity) putNullable("max_budget", task.maxBudget) @@ -168,6 +245,7 @@ class ProcurementSecureStore(context: Context) : ProcurementStateStore { title = json.getString("title"), description = json.optString("description"), sku = json.getString("sku"), + imageAssetId = json.getString("image_asset_id"), referenceImageUrl = json.getString("reference_image_url"), quantity = json.getInt("quantity"), maxBudget = json.optionalString("max_budget"), @@ -184,6 +262,9 @@ class ProcurementSecureStore(context: Context) : ProcurementStateStore { private fun JSONObject.optionalObject(name: String): JSONObject? = if (has(name) && !isNull(name)) getJSONObject(name) else null + private fun JSONObject.optionalArray(name: String): JSONArray? = + if (has(name) && !isNull(name)) getJSONArray(name) else null + private companion object { const val FILE_NAME = "procurement_secure_state" const val STATE_KEY = "state_v1" diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/ProcurementScreen.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/ProcurementScreen.kt index 3d3b314..f90d384 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/ProcurementScreen.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/ProcurementScreen.kt @@ -43,6 +43,7 @@ import androidx.compose.ui.text.input.PasswordVisualTransformation import androidx.compose.ui.unit.dp import androidx.compose.ui.unit.sp import com.roubao.autopilot.procurement.LoginInput +import com.roubao.autopilot.procurement.ExecutionMode import com.roubao.autopilot.procurement.ProcurementPhase import com.roubao.autopilot.procurement.ProcurementUiState import com.roubao.autopilot.readiness.DeviceReadinessSnapshot @@ -56,12 +57,15 @@ fun ProcurementScreen( readiness: DeviceReadinessSnapshot, onLogin: (LoginInput) -> Unit, onClaim: () -> Unit, - onStart: () -> Unit, + onStart: (ExecutionMode) -> Unit, onRelease: () -> Unit, onSync: () -> Unit ) { val colors = BaoziTheme.colors var confirmStart by remember(state.task?.id) { mutableStateOf(false) } + var executionMode by remember(state.task?.id) { + mutableStateOf(ExecutionMode.MANUAL_FIRST) + } if (confirmStart) { AlertDialog( @@ -77,7 +81,7 @@ fun ProcurementScreen( Button( onClick = { confirmStart = false - onStart() + onStart(executionMode) } ) { Icon(Icons.Filled.PlayArrow, contentDescription = null) @@ -154,6 +158,12 @@ fun ProcurementScreen( item { LoginSection(state = state, onLogin = onLogin) } } item { TaskDetails(state) } + item { + ExecutionModeSelector( + selected = executionMode, + onSelect = { executionMode = it } + ) + } item { Row( modifier = Modifier.fillMaxWidth(), @@ -206,6 +216,47 @@ fun ProcurementScreen( } } +@Composable +private fun ExecutionModeSelector( + selected: ExecutionMode, + onSelect: (ExecutionMode) -> Unit +) { + val colors = BaoziTheme.colors + Column(verticalArrangement = Arrangement.spacedBy(8.dp)) { + Text( + "执行模式", + color = colors.textPrimary, + fontSize = 16.sp, + fontWeight = FontWeight.SemiBold + ) + Row( + modifier = Modifier.fillMaxWidth(), + horizontalArrangement = Arrangement.spacedBy(8.dp) + ) { + ExecutionMode.values().forEach { mode -> + OutlinedButton( + onClick = { onSelect(mode) }, + modifier = Modifier.weight(1f), + enabled = mode != selected + ) { + Text( + if (mode == ExecutionMode.MANUAL_FIRST) "人工优先" else "AI 辅助" + ) + } + } + } + Text( + if (selected == ExecutionMode.MANUAL_FIRST) { + "直接采集候选并由人员确认" + } else { + "仅使用本机已配置的 VLM;后台不会读取模型配置" + }, + color = colors.textSecondary, + fontSize = 13.sp + ) + } +} + @Composable private fun LoginSection( state: ProcurementUiState, diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/SearchProbeScreen.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/SearchProbeScreen.kt index 190bb45..12438d3 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/SearchProbeScreen.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/SearchProbeScreen.kt @@ -25,7 +25,12 @@ import androidx.compose.material3.ButtonDefaults import androidx.compose.material3.Divider import androidx.compose.material3.Icon import androidx.compose.material3.OutlinedButton +import androidx.compose.material3.OutlinedTextField import androidx.compose.material3.Text +import androidx.compose.runtime.getValue +import androidx.compose.runtime.mutableStateOf +import androidx.compose.runtime.remember +import androidx.compose.runtime.setValue import androidx.compose.runtime.Composable import androidx.compose.ui.Alignment import androidx.compose.ui.Modifier @@ -85,8 +90,8 @@ fun SearchProbeScreen( onStopRequirement: () -> Unit, onStartCandidateEvaluation: () -> Unit, onStopCandidateEvaluation: () -> Unit, - onAcceptCandidate: () -> Unit, - onRejectCandidates: () -> Unit, + onAcceptCandidate: (String) -> Unit, + onRejectCandidates: (String) -> Unit, onStart: () -> Unit, onStop: () -> Unit ) { @@ -279,10 +284,11 @@ private fun CandidateEvaluationSection( canStart: Boolean, onStart: () -> Unit, onStop: () -> Unit, - onAccept: () -> Unit, - onReject: () -> Unit + onAccept: (String) -> Unit, + onReject: (String) -> Unit ) { val colors = BaoziTheme.colors + var operatorReason by remember(state) { mutableStateOf("") } val statusColor = when (state) { CandidateEvaluationState.AWAITING_CONFIRMATION, CandidateEvaluationState.HUMAN_ACCEPTED -> colors.success @@ -366,6 +372,20 @@ private fun CandidateEvaluationSection( } Spacer(modifier = Modifier.height(14.dp)) + if (state in setOf( + CandidateEvaluationState.AWAITING_CONFIRMATION, + CandidateEvaluationState.MANUAL_REVIEW, + CandidateEvaluationState.NO_MATCH + )) { + OutlinedTextField( + value = operatorReason, + onValueChange = { operatorReason = it.take(1000) }, + label = { Text("人工确认理由") }, + modifier = Modifier.fillMaxWidth(), + minLines = 2 + ) + Spacer(modifier = Modifier.height(10.dp)) + } when (state) { CandidateEvaluationState.RUNNING -> { OutlinedButton( @@ -379,7 +399,8 @@ private fun CandidateEvaluationSection( } CandidateEvaluationState.AWAITING_CONFIRMATION -> { Button( - onClick = onAccept, + onClick = { onAccept(operatorReason) }, + enabled = operatorReason.isNotBlank(), modifier = Modifier.fillMaxWidth(), colors = ButtonDefaults.buttonColors( containerColor = colors.primary @@ -391,7 +412,8 @@ private fun CandidateEvaluationSection( } Spacer(modifier = Modifier.height(8.dp)) OutlinedButton( - onClick = onReject, + onClick = { onReject(operatorReason) }, + enabled = operatorReason.isNotBlank(), modifier = Modifier.fillMaxWidth() ) { Icon(Icons.Default.Close, contentDescription = null) @@ -401,8 +423,21 @@ private fun CandidateEvaluationSection( } CandidateEvaluationState.MANUAL_REVIEW, CandidateEvaluationState.NO_MATCH -> { + if (state == CandidateEvaluationState.MANUAL_REVIEW) { + Button( + onClick = { onAccept(operatorReason) }, + enabled = operatorReason.isNotBlank(), + modifier = Modifier.fillMaxWidth() + ) { + Icon(Icons.Default.CheckCircle, contentDescription = null) + Spacer(modifier = Modifier.size(8.dp)) + Text("标记候选可用") + } + Spacer(modifier = Modifier.height(8.dp)) + } OutlinedButton( - onClick = onReject, + onClick = { onReject(operatorReason) }, + enabled = operatorReason.isNotBlank(), modifier = Modifier.fillMaxWidth() ) { Icon(Icons.Default.Close, contentDescription = null) diff --git a/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/ExecutionTaskHashTest.kt b/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/ExecutionTaskHashTest.kt new file mode 100644 index 0000000..839ef97 --- /dev/null +++ b/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/ExecutionTaskHashTest.kt @@ -0,0 +1,34 @@ +package com.roubao.autopilot.procurement + +import org.junit.Assert.assertEquals +import org.junit.Assert.assertNotEquals +import org.junit.Test + +class ExecutionTaskHashTest { + @Test + fun taskHashUsesOnlyBackendRequirementFields() { + val task = RemotePurchaseTask( + id = "task-id", + status = "RUNNING", + version = 3, + claimGeneration = 1, + claimExpiresAt = "2099-01-01T00:00:00Z", + title = "测试商品", + description = "测试描述", + sku = "SKU-01", + imageAssetId = "asset-id", + referenceImageUrl = "/api/v1/tasks/task-id/reference-image?claim_generation=1", + quantity = 2, + maxBudget = "20.00", + currency = "CNY" + ) + assertEquals( + ExecutionTaskHash.sha256(task), + ExecutionTaskHash.sha256(task.copy(version = 99, status = "CLAIMED")) + ) + assertNotEquals( + ExecutionTaskHash.sha256(task), + ExecutionTaskHash.sha256(task.copy(quantity = 3)) + ) + } +} diff --git a/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/ProcurementApiClientTest.kt b/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/ProcurementApiClientTest.kt index b643c5e..39edf9d 100644 --- a/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/ProcurementApiClientTest.kt +++ b/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/ProcurementApiClientTest.kt @@ -172,6 +172,61 @@ class ProcurementApiClientTest { assertEquals(0, server.requestCount) } + @Test + fun executionOutboxUsesDeviceResultEndpointsWithoutProviderConfiguration() = runBlocking { + val task = task("/api/v1/tasks/task-id/reference-image?claim_generation=1") + val execution = RunningExecution( + id = "execution-id", + currentStep = "SEARCH", + expiresAt = "2099-01-01T00:00:00Z", + serverClockOffsetMillis = 0 + ) + server.enqueue( + jsonResponse( + """{"evidence":{"id":"evidence-id"},"replayed":false}""" + ) + ) + val evidenceResult = api.uploadExecutionOutboxItem( + session = session(), + task = task, + execution = execution, + claimToken = "claim-token-value", + item = ExecutionOutboxItem( + id = "local-evidence-id", + type = ExecutionOutboxType.EVIDENCE, + idempotencyKey = "evidence-idempotency-key", + payload = "{}", + evidenceRelativePath = "procurement-outbox/local-evidence-id.png" + ), + evidenceBytes = byteArrayOf(1, 2, 3) + ) + assertEquals("evidence-id", evidenceResult.evidenceAssetId) + val evidenceRequest = server.takeRequest() + assertEquals("/api/v1/tasks/task-id/evidence", evidenceRequest.path) + assertEquals("claim-token-value", evidenceRequest.getHeader("X-Claim-Token")) + assertEquals("execution-id", evidenceRequest.getHeader("X-Execution-ID")) + assertEquals("1", evidenceRequest.getHeader("X-Claim-Generation")) + assertTrue(evidenceRequest.getHeader("Content-Type")!!.startsWith("image/png")) + + server.enqueue(jsonResponse("""{"replayed":false}""")) + api.uploadExecutionOutboxItem( + session = session(), + task = task, + execution = execution, + claimToken = "claim-token-value", + item = ExecutionOutboxItem( + id = "event-id", + type = ExecutionOutboxType.EVENTS, + idempotencyKey = "event-idempotency-key", + payload = """{"execution_id":"execution-id","claim_generation":1,"events":[]}""" + ) + ) + val eventRequest = server.takeRequest() + assertEquals("/api/v1/tasks/task-id/events", eventRequest.path) + assertEquals("event-idempotency-key", eventRequest.getHeader("Idempotency-Key")) + assertTrue(eventRequest.body.readUtf8().contains("execution-id")) + } + private fun session() = ProcurementSession( backendUrl = server.url("/").toString().trimEnd('/'), username = "buyer01", @@ -191,6 +246,7 @@ class ProcurementApiClientTest { "title":"测试商品", "description":"测试描述", "sku":"SKU-01", + "image_asset_id":"00000000-0000-4000-8000-000000000001", "reference_image_url":"/api/v1/tasks/task-id/reference-image?claim_generation=1", "quantity":2, "max_budget":"20.00", @@ -207,6 +263,7 @@ class ProcurementApiClientTest { title = "测试商品", description = "测试描述", sku = "SKU-01", + imageAssetId = "00000000-0000-4000-8000-000000000001", referenceImageUrl = referenceImageUrl, quantity = 2, maxBudget = "20.00", diff --git a/backend-api/cmd/api/main.go b/backend-api/cmd/api/main.go index f38489b..281e372 100644 --- a/backend-api/cmd/api/main.go +++ b/backend-api/cmd/api/main.go @@ -180,6 +180,15 @@ func buildRouter( if err != nil { return nil, err } + results, err := usecase.NewExecutionResultService( + store, + files, + clock, + ids, + ) + if err != nil { + return nil, err + } passwords, err := password.NewBcrypt(12) if err != nil { return nil, err @@ -188,6 +197,7 @@ func buildRouter( httpapi.DeviceServices{ Lifecycle: lifecycle, Assets: assets, + Results: results, }, ) if err != nil { @@ -245,8 +255,9 @@ func buildRouter( } registerAdminRoutes, err := httpapi.NewAdminRouteRegistrar( httpapi.AdminServices{ - Assets: assets, - Tasks: tasks, + Assets: assets, + Tasks: tasks, + Results: results, }, webHandler, ) diff --git a/backend-api/internal/domain/task.go b/backend-api/internal/domain/task.go index d12a001..46bcb58 100644 --- a/backend-api/internal/domain/task.go +++ b/backend-api/internal/domain/task.go @@ -84,11 +84,75 @@ type TaskExecution struct { FinishedAt *time.Time } +type ExecutionEvent struct { + ID string + TaskID string + ExecutionID string + Step string + Type string + Message string + OccurredAt time.Time + ReceivedAt time.Time + ReceivedAfterExecutionExpiry bool +} + +type ExecutionEvidenceAsset struct { + ID string + TaskID string + ExecutionID string + MediaType string + SizeBytes int64 + SHA256 string + StorageKey string + CreatedAt time.Time + ReceivedAfterExecutionExpiry bool +} + +type ExecutionCandidateBatch struct { + TaskID string + ExecutionID string + TaskContentSHA256 string + ExecutionMode string + SearchQuery string + ProvenanceJSON *string + CandidatesJSON string + RecommendationJSON *string + ReceivedAt time.Time + ReceivedAfterExecutionExpiry bool +} + +type ExecutionOutcome struct { + TaskID string + ExecutionID string + ResultType string + ExecutionMode *string + TaskContentSHA256 *string + Outcome *string + OperatorReason *string + SelectedCandidateJSON *string + EvidenceAssetIDsJSON *string + ErrorCode *string + ErrorMessage *string + ErrorStep *string + Retryable *bool + OrderSubmitted bool + ReceivedAt time.Time + ReceivedAfterExecutionExpiry bool +} + +type ExecutionReport struct { + Events []ExecutionEvent + EvidenceAssets []ExecutionEvidenceAsset + CandidateBatch *ExecutionCandidateBatch + Outcome *ExecutionOutcome +} + type TaskDetail struct { Task PurchaseTask Asset Asset Execution *TaskExecution Events []TaskEvent + Report *ExecutionReport } type TaskValidationError struct { diff --git a/backend-api/internal/platform/migration/claims_migration_test.go b/backend-api/internal/platform/migration/claims_migration_test.go index 7c010c6..70150ad 100644 --- a/backend-api/internal/platform/migration/claims_migration_test.go +++ b/backend-api/internal/platform/migration/claims_migration_test.go @@ -34,8 +34,8 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { if applied, err := runner.Up(ctx); err != nil { t.Fatalf("initial Up() error = %v", err) - } else if applied != 4 { - t.Fatalf("initial Up() applied = %d, want 4", applied) + } else if applied != 5 { + t.Fatalf("initial Up() applied = %d, want 5", applied) } if err := runner.Down(ctx); err != nil { t.Fatalf("initial Down(v4) error = %v", err) @@ -50,6 +50,11 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { } assertClaimsHistory(t, db, true) + if err := runner.Down(ctx); err != nil { + t.Fatalf("Down(v5) with compatible history error = %v", err) + } + assertClaimsHistory(t, db, true) + if err := runner.Down(ctx); err != nil { t.Fatalf("Down(v4) with compatible history error = %v", err) } @@ -57,8 +62,8 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { if applied, err := runner.Up(ctx); err != nil { t.Fatalf("final Up(v4) error = %v", err) - } else if applied != 1 { - t.Fatalf("final Up(v4) applied = %d, want 1", applied) + } else if applied != 2 { + t.Fatalf("final Up(v4-v5) applied = %d, want 2", applied) } assertClaimsHistory(t, db, true) } @@ -294,6 +299,9 @@ func TestClaimsMigrationDownFailsClosedForNewAuditData(t *testing.T) { t.Fatalf("insert v4 audit event: %v", err) } + if err := runner.Down(ctx); err != nil { + t.Fatalf("Down(v5) error = %v", err) + } if err := runner.Down(ctx); err == nil { t.Fatal("Down(v4) succeeded with non-representable audit event") } diff --git a/backend-api/internal/platform/migration/runner_test.go b/backend-api/internal/platform/migration/runner_test.go index 3df07b0..1281e27 100644 --- a/backend-api/internal/platform/migration/runner_test.go +++ b/backend-api/internal/platform/migration/runner_test.go @@ -27,14 +27,15 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { if err != nil { t.Fatalf("Up() error = %v", err) } - if applied != 4 { - t.Fatalf("Up() applied = %d, want 4", applied) + if applied != 5 { + t.Fatalf("Up() applied = %d, want 5", applied) } assertStatuses(t, runner, map[int64]bool{ 1: true, 2: true, 3: true, 4: true, + 5: true, }) applied, err = runner.Up(context.Background()) @@ -52,7 +53,8 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 1: true, 2: true, 3: true, - 4: false, + 4: true, + 5: false, }) applied, err = runner.Up(context.Background()) @@ -67,6 +69,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 2: true, 3: true, 4: true, + 5: true, }) } diff --git a/backend-api/internal/repository/sqlite/auth_repository_test.go b/backend-api/internal/repository/sqlite/auth_repository_test.go index 818e15a..2aa665b 100644 --- a/backend-api/internal/repository/sqlite/auth_repository_test.go +++ b/backend-api/internal/repository/sqlite/auth_repository_test.go @@ -383,6 +383,9 @@ func TestAuthMigrationCanRollbackWithoutRebuildingPurchaseTasks( if err != nil { t.Fatalf("migration.New() error = %v", err) } + if err := runner.Down(context.Background()); err != nil { + t.Fatalf("Down(v5) error = %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("Down(v4) error = %v", err) } @@ -403,8 +406,8 @@ func TestAuthMigrationCanRollbackWithoutRebuildingPurchaseTasks( } if applied, err := runner.Up(context.Background()); err != nil { t.Fatalf("Up(v3-v4) error = %v", err) - } else if applied != 2 { - t.Fatalf("Up(v3-v4) applied = %d, want 2", applied) + } else if applied != 3 { + t.Fatalf("Up(v3-v5) applied = %d, want 3", applied) } } diff --git a/backend-api/internal/repository/sqlite/execution_result_repository.go b/backend-api/internal/repository/sqlite/execution_result_repository.go new file mode 100644 index 0000000..22a7cb9 --- /dev/null +++ b/backend-api/internal/repository/sqlite/execution_result_repository.go @@ -0,0 +1,891 @@ +package sqlite + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "strings" + + "cmroubao/backend-api/internal/domain" + "cmroubao/backend-api/internal/usecase" +) + +type executionResultRequestRecord struct { + RequestHash string + ClaimTokenHash string + TaskID string + ExecutionID string + ResourceID *string +} + +func (s *Store) GetExecutionEvidence( + ctx context.Context, + taskID string, + evidenceID string, +) (domain.ExecutionEvidenceAsset, error) { + return getExecutionEvidenceForTask(ctx, s.db, taskID, evidenceID) +} + +func (s *Store) AppendExecutionEvents( + ctx context.Context, + write usecase.ExecutionResultWrite, + events []domain.ExecutionEvent, +) (bool, error) { + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return false, repositoryFailure(err) + } + defer func() { _ = tx.Rollback() }() + if replayed, err := replayExecutionResultRequest( + ctx, tx, write, + ); err != nil || replayed { + return replayed, err + } + _, _, expired, err := authorizeExecutionResult(ctx, tx, write) + if err != nil { + return false, err + } + for _, event := range events { + _, err = tx.ExecContext( + ctx, + `INSERT INTO execution_events ( + id, task_id, execution_id, step, event_type, message, + occurred_at, received_at, received_after_execution_expiry + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`, + event.ID, + write.TaskID, + write.ExecutionID, + event.Step, + event.Type, + event.Message, + formatTimestamp(event.OccurredAt), + formatTimestamp(write.Now), + expired, + ) + if err != nil { + return false, repositoryFailure(err) + } + } + if err := insertExecutionResultRequest(ctx, tx, write, nil); err != nil { + return false, err + } + if err := tx.Commit(); err != nil { + return false, repositoryFailure(err) + } + return false, nil +} + +func (s *Store) CreateExecutionEvidence( + ctx context.Context, + write usecase.ExecutionResultWrite, + candidate domain.ExecutionEvidenceAsset, +) (domain.ExecutionEvidenceAsset, bool, error) { + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return domain.ExecutionEvidenceAsset{}, false, repositoryFailure(err) + } + defer func() { _ = tx.Rollback() }() + record, found, err := lookupExecutionResultRequest(ctx, tx, write) + if err != nil { + return domain.ExecutionEvidenceAsset{}, false, err + } + if found { + if err := validateExecutionResultReplay(record, write); err != nil { + return domain.ExecutionEvidenceAsset{}, false, err + } + if record.ResourceID == nil { + return domain.ExecutionEvidenceAsset{}, false, usecase.ErrRepositoryInvariant + } + evidence, err := getExecutionEvidence(ctx, tx, *record.ResourceID) + if err != nil { + return domain.ExecutionEvidenceAsset{}, false, err + } + if err := tx.Commit(); err != nil { + return domain.ExecutionEvidenceAsset{}, false, repositoryFailure(err) + } + return evidence, true, nil + } + _, _, expired, err := authorizeExecutionResult(ctx, tx, write) + if err != nil { + return domain.ExecutionEvidenceAsset{}, false, err + } + candidate.ReceivedAfterExecutionExpiry = expired + _, err = tx.ExecContext( + ctx, + `INSERT INTO execution_evidence_assets ( + id, task_id, execution_id, media_type, size_bytes, sha256, + storage_key, created_at, received_after_execution_expiry + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`, + candidate.ID, + write.TaskID, + write.ExecutionID, + candidate.MediaType, + candidate.SizeBytes, + candidate.SHA256, + candidate.StorageKey, + formatTimestamp(candidate.CreatedAt), + expired, + ) + if err != nil { + return domain.ExecutionEvidenceAsset{}, false, repositoryFailure(err) + } + if err := insertExecutionResultRequest(ctx, tx, write, &candidate.ID); err != nil { + return domain.ExecutionEvidenceAsset{}, false, err + } + if err := tx.Commit(); err != nil { + return domain.ExecutionEvidenceAsset{}, false, repositoryFailure(err) + } + return candidate, false, nil +} + +func (s *Store) StoreExecutionCandidates( + ctx context.Context, + write usecase.ExecutionResultWrite, + candidate domain.ExecutionCandidateBatch, +) (bool, error) { + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return false, repositoryFailure(err) + } + defer func() { _ = tx.Rollback() }() + if replayed, err := replayExecutionResultRequest(ctx, tx, write); err != nil || replayed { + return replayed, err + } + task, _, expired, err := authorizeExecutionResult(ctx, tx, write) + if err != nil { + return false, err + } + if usecase.TaskContentSHA256(task) != candidate.TaskContentSHA256 { + return false, usecase.ErrTaskVersionConflict + } + if err := validateCandidateEvidence(ctx, tx, write, candidate.CandidatesJSON); err != nil { + return false, err + } + _, err = tx.ExecContext( + ctx, + `INSERT INTO execution_candidate_batches ( + execution_id, task_id, task_content_sha256, execution_mode, + search_query, provenance_json, candidates_json, recommendation_json, + received_at, received_after_execution_expiry + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + write.ExecutionID, + write.TaskID, + candidate.TaskContentSHA256, + candidate.ExecutionMode, + candidate.SearchQuery, + nullableString(candidate.ProvenanceJSON), + candidate.CandidatesJSON, + nullableString(candidate.RecommendationJSON), + formatTimestamp(write.Now), + expired, + ) + if err != nil { + return false, repositoryFailure(err) + } + if err := insertExecutionResultRequest(ctx, tx, write, nil); err != nil { + return false, err + } + if err := tx.Commit(); err != nil { + return false, repositoryFailure(err) + } + return false, nil +} + +func (s *Store) CompleteExecution( + ctx context.Context, + write usecase.ExecutionResultWrite, + outcome domain.ExecutionOutcome, +) (domain.PurchaseTask, bool, error) { + return s.finishExecution(ctx, write, outcome, domain.TaskStatusSucceeded) +} + +func (s *Store) FailExecution( + ctx context.Context, + write usecase.ExecutionResultWrite, + outcome domain.ExecutionOutcome, +) (domain.PurchaseTask, bool, error) { + return s.finishExecution(ctx, write, outcome, domain.TaskStatusFailed) +} + +func (s *Store) finishExecution( + ctx context.Context, + write usecase.ExecutionResultWrite, + outcome domain.ExecutionOutcome, + terminalStatus domain.TaskStatus, +) (domain.PurchaseTask, bool, error) { + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return domain.PurchaseTask{}, false, repositoryFailure(err) + } + defer func() { _ = tx.Rollback() }() + record, found, err := lookupExecutionResultRequest(ctx, tx, write) + if err != nil { + return domain.PurchaseTask{}, false, err + } + if found { + if err := validateExecutionResultReplay(record, write); err != nil { + return domain.PurchaseTask{}, false, err + } + task, err := getLifecycleTask(ctx, tx, write.TaskID) + if err != nil { + return domain.PurchaseTask{}, false, err + } + if task.Status != terminalStatus { + return domain.PurchaseTask{}, false, usecase.ErrTaskStateConflict + } + if err := tx.Commit(); err != nil { + return domain.PurchaseTask{}, false, repositoryFailure(err) + } + return task, true, nil + } + task, execution, expired, err := authorizeExecutionResult(ctx, tx, write) + if err != nil { + return domain.PurchaseTask{}, false, err + } + if task.CancelRequestedAt != nil { + return domain.PurchaseTask{}, false, usecase.ErrTaskStateConflict + } + if outcome.ResultType == "COMPLETE" { + if outcome.TaskContentSHA256 == nil || + *outcome.TaskContentSHA256 != usecase.TaskContentSHA256(task) { + return domain.PurchaseTask{}, false, usecase.ErrTaskVersionConflict + } + if err := validateCompleteCandidate(ctx, tx, write, outcome); err != nil { + return domain.PurchaseTask{}, false, err + } + if err := validateSelectedEvidence(ctx, tx, write, outcome.SelectedCandidateJSON); err != nil { + return domain.PurchaseTask{}, false, err + } + } else if err := validateEvidenceIDsFromOutcome( + ctx, tx, write, outcome.EvidenceAssetIDsJSON, + ); err != nil { + return domain.PurchaseTask{}, false, err + } + _ = execution + outcome.ReceivedAfterExecutionExpiry = expired + _, err = tx.ExecContext( + ctx, + `INSERT INTO execution_outcomes ( + execution_id, task_id, result_type, execution_mode, task_content_sha256, + outcome, operator_reason, selected_candidate_json, evidence_asset_ids_json, + error_code, error_message, error_step, retryable, order_submitted, received_at, + received_after_execution_expiry + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 0, ?, ?)`, + write.ExecutionID, + write.TaskID, + outcome.ResultType, + nullableString(outcome.ExecutionMode), + nullableString(outcome.TaskContentSHA256), + nullableString(outcome.Outcome), + nullableString(outcome.OperatorReason), + nullableString(outcome.SelectedCandidateJSON), + nullableString(outcome.EvidenceAssetIDsJSON), + nullableString(outcome.ErrorCode), + nullableString(outcome.ErrorMessage), + nullableString(outcome.ErrorStep), + nullableBool(outcome.Retryable), + formatTimestamp(write.Now), + expired, + ) + if err != nil { + return domain.PurchaseTask{}, false, repositoryFailure(err) + } + step := "COMPLETED" + if terminalStatus == domain.TaskStatusFailed { + step = "FAILED" + } + result, err := tx.ExecContext( + ctx, + `UPDATE task_executions + SET current_step = ?, + last_heartbeat_at = ?, + finished_at = ? + WHERE id = ? AND finished_at IS NULL`, + step, + formatTimestamp(write.Now), + formatTimestamp(write.Now), + write.ExecutionID, + ) + if err != nil { + return domain.PurchaseTask{}, false, repositoryFailure(err) + } + if affected, err := result.RowsAffected(); err != nil || affected != 1 { + if err != nil { + return domain.PurchaseTask{}, false, repositoryFailure(err) + } + return domain.PurchaseTask{}, false, usecase.ErrExecutionMismatch + } + result, err = tx.ExecContext( + ctx, + `UPDATE purchase_tasks + SET status = ?, + version = version + 1, + claimed_by_user_id = NULL, + claimed_by_device_id = NULL, + claim_token_hash = NULL, + claim_issued_at = NULL, + claim_expires_at = NULL, + updated_at = ? + WHERE id = ? + AND status IN ('RUNNING', 'WAITING_CONFIRMATION') + AND claim_generation = ?`, + terminalStatus, + formatTimestamp(write.Now), + write.TaskID, + write.ClaimGeneration, + ) + if err != nil { + return domain.PurchaseTask{}, false, repositoryFailure(err) + } + if affected, err := result.RowsAffected(); err != nil || affected != 1 { + if err != nil { + return domain.PurchaseTask{}, false, repositoryFailure(err) + } + return domain.PurchaseTask{}, false, usecase.ErrTaskStateConflict + } + if err := insertExecutionResultRequest(ctx, tx, write, nil); err != nil { + return domain.PurchaseTask{}, false, err + } + task, err = getLifecycleTask(ctx, tx, write.TaskID) + if err != nil { + return domain.PurchaseTask{}, false, err + } + if err := tx.Commit(); err != nil { + return domain.PurchaseTask{}, false, repositoryFailure(err) + } + return task, false, nil +} + +func authorizeExecutionResult( + ctx context.Context, + tx *sql.Tx, + write usecase.ExecutionResultWrite, +) (domain.PurchaseTask, domain.TaskExecution, bool, error) { + task, err := getClaimProtectedTask(ctx, tx, write.TaskID) + if err != nil { + return domain.PurchaseTask{}, domain.TaskExecution{}, false, err + } + if err := validateClaimOwner( + task, + write.UserID, + write.DeviceID, + write.ClaimGeneration, + write.ClaimTokenHash, + ); err != nil { + return domain.PurchaseTask{}, domain.TaskExecution{}, false, err + } + if !domain.CanHeartbeat(task.Status) { + return domain.PurchaseTask{}, domain.TaskExecution{}, false, usecase.ErrTaskStateConflict + } + execution, err := getExecutionByID(ctx, tx, write.ExecutionID) + if err != nil { + return domain.PurchaseTask{}, domain.TaskExecution{}, false, err + } + if execution.TaskID != write.TaskID || execution.UserID != write.UserID || + execution.DeviceID != write.DeviceID || + execution.ClaimGeneration != write.ClaimGeneration || execution.FinishedAt != nil { + return domain.PurchaseTask{}, domain.TaskExecution{}, false, usecase.ErrExecutionMismatch + } + expired := task.ClaimExpiresAt == nil || !task.ClaimExpiresAt.After(write.Now) + return task, execution, expired, nil +} + +func replayExecutionResultRequest( + ctx context.Context, + tx *sql.Tx, + write usecase.ExecutionResultWrite, +) (bool, error) { + record, found, err := lookupExecutionResultRequest(ctx, tx, write) + if err != nil || !found { + return false, err + } + if err := validateExecutionResultReplay(record, write); err != nil { + return false, err + } + if err := tx.Commit(); err != nil { + return false, repositoryFailure(err) + } + return true, nil +} + +func lookupExecutionResultRequest( + ctx context.Context, + tx *sql.Tx, + write usecase.ExecutionResultWrite, +) (executionResultRequestRecord, bool, error) { + var record executionResultRequestRecord + var resourceID sql.NullString + err := tx.QueryRowContext( + ctx, + `SELECT request_sha256, claim_token_sha256, task_id, execution_id, resource_id + FROM execution_result_requests + WHERE user_id = ? AND device_id = ? AND operation = ? AND idempotency_key = ?`, + write.UserID, + write.DeviceID, + write.Operation, + write.IdempotencyKey, + ).Scan( + &record.RequestHash, + &record.ClaimTokenHash, + &record.TaskID, + &record.ExecutionID, + &resourceID, + ) + if errors.Is(err, sql.ErrNoRows) { + return executionResultRequestRecord{}, false, nil + } + if err != nil { + return executionResultRequestRecord{}, false, repositoryFailure(err) + } + if resourceID.Valid { + record.ResourceID = &resourceID.String + } + return record, true, nil +} + +func validateExecutionResultReplay( + record executionResultRequestRecord, + write usecase.ExecutionResultWrite, +) error { + if record.RequestHash != write.RequestHash || + record.ClaimTokenHash != write.ClaimTokenHash || + record.TaskID != write.TaskID || record.ExecutionID != write.ExecutionID { + return usecase.ErrIdempotencyConflict + } + return nil +} + +func insertExecutionResultRequest( + ctx context.Context, + tx *sql.Tx, + write usecase.ExecutionResultWrite, + resourceID *string, +) error { + _, err := tx.ExecContext( + ctx, + `INSERT INTO execution_result_requests ( + user_id, device_id, operation, idempotency_key, request_sha256, + claim_token_sha256, task_id, execution_id, resource_id, created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + write.UserID, + write.DeviceID, + write.Operation, + write.IdempotencyKey, + write.RequestHash, + write.ClaimTokenHash, + write.TaskID, + write.ExecutionID, + nullableString(resourceID), + formatTimestamp(write.Now), + ) + if err != nil { + return repositoryFailure(err) + } + return nil +} + +func getExecutionEvidence( + ctx context.Context, + queryer queryRower, + evidenceID string, +) (domain.ExecutionEvidenceAsset, error) { + var evidence domain.ExecutionEvidenceAsset + var createdAt string + var receivedAfter bool + err := queryer.QueryRowContext( + ctx, + `SELECT id, task_id, execution_id, media_type, size_bytes, sha256, + storage_key, created_at, received_after_execution_expiry + FROM execution_evidence_assets WHERE id = ?`, evidenceID, + ).Scan( + &evidence.ID, + &evidence.TaskID, + &evidence.ExecutionID, + &evidence.MediaType, + &evidence.SizeBytes, + &evidence.SHA256, + &evidence.StorageKey, + &createdAt, + &receivedAfter, + ) + if errors.Is(err, sql.ErrNoRows) { + return domain.ExecutionEvidenceAsset{}, usecase.ErrRepositoryInvariant + } + if err != nil { + return domain.ExecutionEvidenceAsset{}, repositoryFailure(err) + } + evidence.CreatedAt, err = parseTimestamp(createdAt) + if err != nil { + return domain.ExecutionEvidenceAsset{}, repositoryFailure(err) + } + evidence.ReceivedAfterExecutionExpiry = receivedAfter + return evidence, nil +} + +func getExecutionEvidenceForTask( + ctx context.Context, + queryer queryRower, + taskID string, + evidenceID string, +) (domain.ExecutionEvidenceAsset, error) { + evidence, err := getExecutionEvidence(ctx, queryer, evidenceID) + if err != nil { + if errors.Is(err, usecase.ErrRepositoryInvariant) { + return domain.ExecutionEvidenceAsset{}, usecase.ErrRepositoryNotFound + } + return domain.ExecutionEvidenceAsset{}, err + } + if evidence.TaskID != taskID { + return domain.ExecutionEvidenceAsset{}, usecase.ErrRepositoryNotFound + } + return evidence, nil +} + +func validateCandidateEvidence( + ctx context.Context, + tx *sql.Tx, + write usecase.ExecutionResultWrite, + candidatesJSON string, +) error { + var candidates []struct { + EvidenceAssetIDs []string `json:"evidence_asset_ids"` + } + if err := json.Unmarshal([]byte(candidatesJSON), &candidates); err != nil { + return usecase.ErrRepositoryInvariant + } + for _, candidate := range candidates { + if err := verifyExecutionEvidenceIDs( + ctx, tx, write.TaskID, write.ExecutionID, candidate.EvidenceAssetIDs, + ); err != nil { + return err + } + } + return nil +} + +func validateSelectedEvidence( + ctx context.Context, + tx *sql.Tx, + write usecase.ExecutionResultWrite, + selectedJSON *string, +) error { + if selectedJSON == nil { + return nil + } + var selected struct { + EvidenceAssetIDs []string `json:"evidence_asset_ids"` + } + if err := json.Unmarshal([]byte(*selectedJSON), &selected); err != nil { + return usecase.ErrRepositoryInvariant + } + return verifyExecutionEvidenceIDs( + ctx, tx, write.TaskID, write.ExecutionID, selected.EvidenceAssetIDs, + ) +} + +func validateCompleteCandidate( + ctx context.Context, + tx *sql.Tx, + write usecase.ExecutionResultWrite, + outcome domain.ExecutionOutcome, +) error { + if outcome.Outcome == nil || *outcome.Outcome != "CANDIDATE_ACCEPTED" { + return nil + } + if outcome.ExecutionMode == nil || outcome.SelectedCandidateJSON == nil { + return usecase.ErrRepositoryInvariant + } + var mode string + var candidatesJSON string + err := tx.QueryRowContext( + ctx, + `SELECT execution_mode, candidates_json + FROM execution_candidate_batches + WHERE execution_id = ? AND task_id = ?`, + write.ExecutionID, + write.TaskID, + ).Scan(&mode, &candidatesJSON) + if errors.Is(err, sql.ErrNoRows) { + return usecase.ErrTaskStateConflict + } + if err != nil { + return repositoryFailure(err) + } + if mode != *outcome.ExecutionMode { + return usecase.ErrTaskStateConflict + } + var selected usecase.ExecutionCandidate + var candidates []usecase.ExecutionCandidate + if err := json.Unmarshal([]byte(*outcome.SelectedCandidateJSON), &selected); err != nil { + return usecase.ErrRepositoryInvariant + } + if err := json.Unmarshal([]byte(candidatesJSON), &candidates); err != nil { + return usecase.ErrRepositoryInvariant + } + for _, candidate := range candidates { + if candidate.Ordinal != selected.Ordinal { + continue + } + stored, storedErr := json.Marshal(candidate) + selectedJSON, selectedErr := json.Marshal(selected) + if storedErr != nil || selectedErr != nil { + return usecase.ErrRepositoryInvariant + } + if string(stored) == string(selectedJSON) { + return nil + } + return usecase.ErrTaskStateConflict + } + return usecase.ErrTaskStateConflict +} + +func validateEvidenceIDsFromOutcome( + ctx context.Context, + tx *sql.Tx, + write usecase.ExecutionResultWrite, + evidenceJSON *string, +) error { + if evidenceJSON == nil { + return nil + } + var evidenceIDs []string + if err := json.Unmarshal([]byte(*evidenceJSON), &evidenceIDs); err != nil { + return usecase.ErrRepositoryInvariant + } + return verifyExecutionEvidenceIDs( + ctx, tx, write.TaskID, write.ExecutionID, evidenceIDs, + ) +} + +func verifyExecutionEvidenceIDs( + ctx context.Context, + tx *sql.Tx, + taskID string, + executionID string, + ids []string, +) error { + seen := map[string]struct{}{} + for _, evidenceID := range ids { + evidenceID = strings.TrimSpace(evidenceID) + if _, duplicate := seen[evidenceID]; duplicate { + return usecase.ErrTaskStateConflict + } + seen[evidenceID] = struct{}{} + var exists int + err := tx.QueryRowContext( + ctx, + `SELECT EXISTS ( + SELECT 1 FROM execution_evidence_assets + WHERE id = ? AND task_id = ? AND execution_id = ? + )`, + evidenceID, + taskID, + executionID, + ).Scan(&exists) + if err != nil { + return repositoryFailure(err) + } + if exists != 1 { + return usecase.ErrTaskStateConflict + } + } + return nil +} + +func getExecutionReport( + ctx context.Context, + queryer queryer, + taskID string, + executionID string, +) (*domain.ExecutionReport, error) { + report := &domain.ExecutionReport{ + Events: make([]domain.ExecutionEvent, 0), + EvidenceAssets: make([]domain.ExecutionEvidenceAsset, 0), + } + events, err := queryer.QueryContext( + ctx, + `SELECT id, task_id, execution_id, step, event_type, message, + occurred_at, received_at, received_after_execution_expiry + FROM execution_events + WHERE task_id = ? AND execution_id = ? + ORDER BY occurred_at ASC, id ASC`, + taskID, + executionID, + ) + if err != nil { + return nil, repositoryFailure(err) + } + defer events.Close() + for events.Next() { + var event domain.ExecutionEvent + var occurredAt, receivedAt string + if err := events.Scan( + &event.ID, + &event.TaskID, + &event.ExecutionID, + &event.Step, + &event.Type, + &event.Message, + &occurredAt, + &receivedAt, + &event.ReceivedAfterExecutionExpiry, + ); err != nil { + return nil, repositoryFailure(err) + } + var parseErr error + event.OccurredAt, parseErr = parseTimestamp(occurredAt) + if parseErr == nil { + event.ReceivedAt, parseErr = parseTimestamp(receivedAt) + } + if parseErr != nil { + return nil, repositoryFailure(parseErr) + } + report.Events = append(report.Events, event) + } + if err := events.Err(); err != nil { + return nil, repositoryFailure(err) + } + evidenceRows, err := queryer.QueryContext( + ctx, + `SELECT id, task_id, execution_id, media_type, size_bytes, sha256, + storage_key, created_at, received_after_execution_expiry + FROM execution_evidence_assets + WHERE task_id = ? AND execution_id = ? + ORDER BY created_at ASC, id ASC`, + taskID, + executionID, + ) + if err != nil { + return nil, repositoryFailure(err) + } + defer evidenceRows.Close() + for evidenceRows.Next() { + var evidence domain.ExecutionEvidenceAsset + var createdAt string + if err := evidenceRows.Scan( + &evidence.ID, + &evidence.TaskID, + &evidence.ExecutionID, + &evidence.MediaType, + &evidence.SizeBytes, + &evidence.SHA256, + &evidence.StorageKey, + &createdAt, + &evidence.ReceivedAfterExecutionExpiry, + ); err != nil { + return nil, repositoryFailure(err) + } + parsed, err := parseTimestamp(createdAt) + if err != nil { + return nil, repositoryFailure(err) + } + evidence.CreatedAt = parsed + report.EvidenceAssets = append(report.EvidenceAssets, evidence) + } + if err := evidenceRows.Err(); err != nil { + return nil, repositoryFailure(err) + } + var batch domain.ExecutionCandidateBatch + var provenance, recommendation sql.NullString + var receivedAt string + err = queryer.QueryRowContext( + ctx, + `SELECT task_id, execution_id, task_content_sha256, execution_mode, + search_query, provenance_json, candidates_json, recommendation_json, + received_at, received_after_execution_expiry + FROM execution_candidate_batches + WHERE task_id = ? AND execution_id = ?`, + taskID, + executionID, + ).Scan( + &batch.TaskID, + &batch.ExecutionID, + &batch.TaskContentSHA256, + &batch.ExecutionMode, + &batch.SearchQuery, + &provenance, + &batch.CandidatesJSON, + &recommendation, + &receivedAt, + &batch.ReceivedAfterExecutionExpiry, + ) + if err == nil { + if provenance.Valid { + batch.ProvenanceJSON = &provenance.String + } + if recommendation.Valid { + batch.RecommendationJSON = &recommendation.String + } + batch.ReceivedAt, err = parseTimestamp(receivedAt) + if err != nil { + return nil, repositoryFailure(err) + } + report.CandidateBatch = &batch + } else if !errors.Is(err, sql.ErrNoRows) { + return nil, repositoryFailure(err) + } + var outcome domain.ExecutionOutcome + var mode, hash, result, reason, selected, evidenceIDs, code, message, step sql.NullString + var retryable sql.NullBool + var outcomeReceivedAt string + err = queryer.QueryRowContext( + ctx, + `SELECT task_id, execution_id, result_type, execution_mode, + task_content_sha256, outcome, operator_reason, selected_candidate_json, + evidence_asset_ids_json, error_code, error_message, error_step, retryable, order_submitted, + received_at, received_after_execution_expiry + FROM execution_outcomes + WHERE task_id = ? AND execution_id = ?`, + taskID, + executionID, + ).Scan( + &outcome.TaskID, + &outcome.ExecutionID, + &outcome.ResultType, + &mode, + &hash, + &result, + &reason, + &selected, + &evidenceIDs, + &code, + &message, + &step, + &retryable, + &outcome.OrderSubmitted, + &outcomeReceivedAt, + &outcome.ReceivedAfterExecutionExpiry, + ) + if err == nil { + outcome.ExecutionMode = nullableStringFromSQL(mode) + outcome.TaskContentSHA256 = nullableStringFromSQL(hash) + outcome.Outcome = nullableStringFromSQL(result) + outcome.OperatorReason = nullableStringFromSQL(reason) + outcome.SelectedCandidateJSON = nullableStringFromSQL(selected) + outcome.EvidenceAssetIDsJSON = nullableStringFromSQL(evidenceIDs) + outcome.ErrorCode = nullableStringFromSQL(code) + outcome.ErrorMessage = nullableStringFromSQL(message) + outcome.ErrorStep = nullableStringFromSQL(step) + if retryable.Valid { + outcome.Retryable = &retryable.Bool + } + outcome.ReceivedAt, err = parseTimestamp(outcomeReceivedAt) + if err != nil { + return nil, repositoryFailure(err) + } + report.Outcome = &outcome + } else if !errors.Is(err, sql.ErrNoRows) { + return nil, repositoryFailure(err) + } + return report, nil +} + +func nullableStringFromSQL(value sql.NullString) *string { + if !value.Valid { + return nil + } + return &value.String +} + +var _ usecase.ExecutionResultRepository = (*Store)(nil) diff --git a/backend-api/internal/repository/sqlite/helpers.go b/backend-api/internal/repository/sqlite/helpers.go index ee80689..524ae56 100644 --- a/backend-api/internal/repository/sqlite/helpers.go +++ b/backend-api/internal/repository/sqlite/helpers.go @@ -20,6 +20,11 @@ type queryRower interface { QueryRowContext(context.Context, string, ...any) *sql.Row } +type queryer interface { + queryRower + QueryContext(context.Context, string, ...any) (*sql.Rows, error) +} + type rowScanner interface { Scan(...any) error } @@ -332,6 +337,13 @@ func nullableInt64(value *int64) any { return *value } +func nullableBool(value *bool) any { + if value == nil { + return nil + } + return *value +} + func repositoryFailure(err error) error { if err == nil { return nil diff --git a/backend-api/internal/repository/sqlite/task_repository.go b/backend-api/internal/repository/sqlite/task_repository.go index 3b35ece..21b9a9a 100644 --- a/backend-api/internal/repository/sqlite/task_repository.go +++ b/backend-api/internal/repository/sqlite/task_repository.go @@ -313,11 +313,19 @@ func (s *Store) GetTaskDetail( } else { executionPointer = &execution } + var report *domain.ExecutionReport + if executionPointer != nil { + report, err = getExecutionReport(ctx, tx, taskID, executionPointer.ID) + if err != nil { + return domain.TaskDetail{}, err + } + } detail := domain.TaskDetail{ Task: task, Asset: asset, Execution: executionPointer, Events: events, + Report: report, } if err := tx.Commit(); err != nil { return domain.TaskDetail{}, repositoryFailure(err) diff --git a/backend-api/internal/transport/httpapi/admin_handlers.go b/backend-api/internal/transport/httpapi/admin_handlers.go index 8d06f84..500c99f 100644 --- a/backend-api/internal/transport/httpapi/admin_handlers.go +++ b/backend-api/internal/transport/httpapi/admin_handlers.go @@ -24,12 +24,13 @@ const ( ) type AdminServices struct { - Assets *usecase.AssetService - Tasks *usecase.TaskService + Assets *usecase.AssetService + Tasks *usecase.TaskService + Results *usecase.ExecutionResultService } func (s AdminServices) validate() error { - if s.Assets == nil || s.Tasks == nil { + if s.Assets == nil || s.Tasks == nil || s.Results == nil { return errors.New("admin services are required") } return nil @@ -49,10 +50,38 @@ func registerAdminAPI(routes gin.IRoutes, services AdminServices) error { routes.POST("/api/v1/tasks", handler.createTask) routes.GET("/api/v1/tasks", handler.listTasks) routes.GET("/api/v1/tasks/:id", handler.taskDetail) + routes.GET( + "/api/v1/tasks/:id/evidence/:evidence_id/content", + handler.evidenceContent, + ) routes.POST("/api/v1/tasks/:id/cancel", handler.cancelTask) return nil } +func (h *adminHandlers) evidenceContent(ctx *gin.Context) { + result, err := h.services.Results.OpenEvidence( + ctx.Request.Context(), + ctx.Param("id"), + ctx.Param("evidence_id"), + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + defer result.Content.Close() + ctx.Header("Cache-Control", "private, no-store") + ctx.Header("Content-Type", result.Evidence.MediaType) + ctx.Header("Content-Length", strconv.FormatInt(result.Evidence.SizeBytes, 10)) + ctx.Header("ETag", `"`+result.Evidence.SHA256+`"`) + ctx.Header("X-Content-Type-Options", "nosniff") + ctx.Header( + "Content-Disposition", + `inline; filename="`+result.Evidence.ID+`.jpg"`, + ) + ctx.Status(http.StatusOK) + _, _ = io.Copy(ctx.Writer, result.Content) +} + func (h *adminHandlers) uploadAsset(ctx *gin.Context) { if !hasMediaType(ctx, "multipart/form-data") { writePublicError( @@ -317,6 +346,10 @@ func (h *adminHandlers) taskDetail(ctx *gin.Context) { if detail.Execution != nil { execution = executionResponse(*detail.Execution) } + var executionReport any + if detail.Report != nil { + executionReport = executionReportResponse(detail.Report) + } ctx.Header("Cache-Control", "no-store") ctx.JSON(http.StatusOK, gin.H{ "id": detail.Task.ID, @@ -337,6 +370,7 @@ func (h *adminHandlers) taskDetail(ctx *gin.Context) { "derived_requirement": nil, "claim": claim, "execution": execution, + "execution_report": executionReport, "events": events, "assets": []gin.H{ assetResponse(detail.Asset), @@ -344,6 +378,78 @@ func (h *adminHandlers) taskDetail(ctx *gin.Context) { }) } +func executionReportResponse(report *domain.ExecutionReport) gin.H { + events := make([]gin.H, 0, len(report.Events)) + for _, event := range report.Events { + events = append(events, gin.H{ + "id": event.ID, + "step": event.Step, + "type": event.Type, + "message": event.Message, + "occurred_at": formatTime(event.OccurredAt), + "received_at": formatTime(event.ReceivedAt), + "received_after_execution_expiry": event.ReceivedAfterExecutionExpiry, + }) + } + evidence := make([]gin.H, 0, len(report.EvidenceAssets)) + for _, asset := range report.EvidenceAssets { + evidence = append(evidence, gin.H{ + "id": asset.ID, + "media_type": asset.MediaType, + "size_bytes": asset.SizeBytes, + "sha256": asset.SHA256, + "created_at": formatTime(asset.CreatedAt), + "received_after_execution_expiry": asset.ReceivedAfterExecutionExpiry, + }) + } + response := gin.H{ + "events": events, + "evidence": evidence, + } + if batch := report.CandidateBatch; batch != nil { + response["candidate_batch"] = gin.H{ + "task_content_sha256": batch.TaskContentSHA256, + "execution_mode": batch.ExecutionMode, + "search_query": batch.SearchQuery, + "provenance": decodedAuditJSON(batch.ProvenanceJSON), + "candidates": decodedAuditJSON(&batch.CandidatesJSON), + "recommendation": decodedAuditJSON(batch.RecommendationJSON), + "received_at": formatTime(batch.ReceivedAt), + "received_after_execution_expiry": batch.ReceivedAfterExecutionExpiry, + } + } + if outcome := report.Outcome; outcome != nil { + response["outcome"] = gin.H{ + "result_type": outcome.ResultType, + "execution_mode": outcome.ExecutionMode, + "task_content_sha256": outcome.TaskContentSHA256, + "outcome": outcome.Outcome, + "operator_reason": outcome.OperatorReason, + "selected_candidate": decodedAuditJSON(outcome.SelectedCandidateJSON), + "evidence_asset_ids": decodedAuditJSON(outcome.EvidenceAssetIDsJSON), + "error_code": outcome.ErrorCode, + "error_message": outcome.ErrorMessage, + "error_step": outcome.ErrorStep, + "retryable": outcome.Retryable, + "order_submitted": outcome.OrderSubmitted, + "received_at": formatTime(outcome.ReceivedAt), + "received_after_execution_expiry": outcome.ReceivedAfterExecutionExpiry, + } + } + return response +} + +func decodedAuditJSON(value *string) any { + if value == nil { + return nil + } + var decoded any + if err := json.Unmarshal([]byte(*value), &decoded); err != nil { + return nil + } + return decoded +} + func (h *adminHandlers) cancelTask(ctx *gin.Context) { if !hasMediaType(ctx, "application/json") { writePublicError( diff --git a/backend-api/internal/transport/httpapi/admin_handlers_test.go b/backend-api/internal/transport/httpapi/admin_handlers_test.go index 542c9a6..472e60d 100644 --- a/backend-api/internal/transport/httpapi/admin_handlers_test.go +++ b/backend-api/internal/transport/httpapi/admin_handlers_test.go @@ -349,8 +349,12 @@ func newAdminIntegrationRouter(t *testing.T) http.Handler { if err != nil { t.Fatalf("usecase.NewTaskService() error = %v", err) } + results, err := usecase.NewExecutionResultService(repositories, files, clock, ids) + if err != nil { + t.Fatalf("usecase.NewExecutionResultService() error = %v", err) + } registrar, err := NewAdminRouteRegistrar( - AdminServices{Assets: assets, Tasks: tasks}, + AdminServices{Assets: assets, Tasks: tasks, Results: results}, emptyAdminWeb{}, ) if err != nil { diff --git a/backend-api/internal/transport/httpapi/device_handlers.go b/backend-api/internal/transport/httpapi/device_handlers.go index 4f718ae..d1497cc 100644 --- a/backend-api/internal/transport/httpapi/device_handlers.go +++ b/backend-api/internal/transport/httpapi/device_handlers.go @@ -19,10 +19,11 @@ const claimTokenHeader = "X-Claim-Token" type DeviceServices struct { Lifecycle *usecase.LifecycleService Assets *usecase.AssetService + Results *usecase.ExecutionResultService } func (services DeviceServices) validate() error { - if services.Lifecycle == nil || services.Assets == nil { + if services.Lifecycle == nil || services.Assets == nil || services.Results == nil { return errors.New("device services are required") } return nil @@ -68,6 +69,11 @@ func NewDeviceRouteRegistrar( "/api/v1/tasks/:id/cancel-ack", handler.acknowledgeCancellation, ) + routes.POST("/api/v1/tasks/:id/events", handler.appendEvents) + routes.POST("/api/v1/tasks/:id/evidence", handler.uploadEvidence) + routes.POST("/api/v1/tasks/:id/candidates", handler.storeCandidates) + routes.POST("/api/v1/tasks/:id/complete", handler.completeTask) + routes.POST("/api/v1/tasks/:id/fail", handler.failTask) return nil }, nil } @@ -367,6 +373,257 @@ func (handler *deviceHandlers) acknowledgeCancellation( }) } +func (handler *deviceHandlers) appendEvents(ctx *gin.Context) { + principal, ok := devicePrincipal(ctx) + if !ok { + return + } + var request struct { + ExecutionID string `json:"execution_id"` + ClaimGeneration int64 `json:"claim_generation"` + Events []usecase.ClientExecutionEvent `json:"events"` + } + if !decodeDeviceJSON(ctx, &request) { + return + } + replayed, err := handler.services.Results.AppendEvents( + ctx.Request.Context(), + usecase.AppendExecutionEventsCommand{ + Identity: handler.executionResultIdentity( + ctx, principal, request.ExecutionID, request.ClaimGeneration, + ), + Events: request.Events, + }, + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + ctx.Header("Cache-Control", "no-store") + ctx.JSON(http.StatusOK, gin.H{"replayed": replayed}) +} + +func (handler *deviceHandlers) uploadEvidence(ctx *gin.Context) { + principal, ok := devicePrincipal(ctx) + if !ok { + return + } + if !isEvidenceMediaType(ctx.GetHeader("Content-Type")) { + writePublicError( + ctx, + http.StatusUnsupportedMediaType, + "ASSET_MEDIA_TYPE_UNSUPPORTED", + "JPEG, PNG, or WebP evidence is required", + false, + gin.H{}, + ) + return + } + generation, err := strconv.ParseInt( + strings.TrimSpace(ctx.GetHeader("X-Claim-Generation")), + 10, + 64, + ) + if err != nil { + writePublicError( + ctx, + http.StatusUnprocessableEntity, + "EXECUTION_RESULT_INVALID", + "execution result request is invalid", + false, + fieldDetails("claim_generation", "must be a positive integer"), + ) + return + } + ctx.Request.Body = http.MaxBytesReader(ctx.Writer, ctx.Request.Body, maxMultipartBytes) + result, err := handler.services.Results.UploadEvidence( + ctx.Request.Context(), + usecase.UploadExecutionEvidenceCommand{ + Identity: handler.executionResultIdentity( + ctx, + principal, + ctx.GetHeader("X-Execution-ID"), + generation, + ), + DeclaredMediaType: ctx.GetHeader("Content-Type"), + Content: ctx.Request.Body, + }, + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + ctx.Header("Cache-Control", "no-store") + ctx.JSON(http.StatusCreated, gin.H{ + "evidence": deviceEvidenceResponse(result.Evidence), + "replayed": result.Replayed, + }) +} + +func (handler *deviceHandlers) storeCandidates(ctx *gin.Context) { + principal, ok := devicePrincipal(ctx) + if !ok { + return + } + var request struct { + ExecutionID string `json:"execution_id"` + ClaimGeneration int64 `json:"claim_generation"` + TaskContentSHA256 string `json:"task_content_sha256"` + ExecutionMode string `json:"execution_mode"` + SearchQuery string `json:"search_query"` + Provenance *usecase.ExecutionProvenance `json:"provenance"` + Candidates []usecase.ExecutionCandidate `json:"candidates"` + Recommendation *usecase.CandidateRecommendation `json:"recommendation"` + } + if !decodeDeviceJSON(ctx, &request) { + return + } + replayed, err := handler.services.Results.StoreCandidates( + ctx.Request.Context(), + usecase.StoreExecutionCandidatesCommand{ + Identity: handler.executionResultIdentity( + ctx, principal, request.ExecutionID, request.ClaimGeneration, + ), + TaskContentSHA256: request.TaskContentSHA256, + ExecutionMode: request.ExecutionMode, + SearchQuery: request.SearchQuery, + Provenance: request.Provenance, + Candidates: request.Candidates, + Recommendation: request.Recommendation, + }, + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + ctx.Header("Cache-Control", "no-store") + ctx.JSON(http.StatusOK, gin.H{"replayed": replayed}) +} + +func (handler *deviceHandlers) completeTask(ctx *gin.Context) { + principal, ok := devicePrincipal(ctx) + if !ok { + return + } + var request struct { + ExecutionID string `json:"execution_id"` + ClaimGeneration int64 `json:"claim_generation"` + TaskContentSHA256 string `json:"task_content_sha256"` + ExecutionMode string `json:"execution_mode"` + Outcome string `json:"outcome"` + OperatorReason string `json:"operator_reason"` + Candidate *usecase.ExecutionCandidate `json:"candidate"` + OrderSubmitted bool `json:"order_submitted"` + } + if !decodeDeviceJSON(ctx, &request) { + return + } + result, replayed, err := handler.services.Results.Complete( + ctx.Request.Context(), + usecase.CompleteExecutionCommand{ + Identity: handler.executionResultIdentity( + ctx, principal, request.ExecutionID, request.ClaimGeneration, + ), + TaskContentSHA256: request.TaskContentSHA256, + ExecutionMode: request.ExecutionMode, + Outcome: request.Outcome, + OperatorReason: request.OperatorReason, + Candidate: request.Candidate, + OrderSubmitted: request.OrderSubmitted, + }, + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + ctx.Header("Cache-Control", "no-store") + ctx.JSON(http.StatusOK, gin.H{ + "task": deviceTaskResponse(result), + "replayed": replayed, + }) +} + +func (handler *deviceHandlers) failTask(ctx *gin.Context) { + principal, ok := devicePrincipal(ctx) + if !ok { + return + } + var request struct { + ExecutionID string `json:"execution_id"` + ClaimGeneration int64 `json:"claim_generation"` + Error struct { + Code string `json:"code"` + Message string `json:"message"` + Step string `json:"step"` + Retryable bool `json:"retryable"` + } `json:"error"` + EvidenceAssetIDs []string `json:"evidence_asset_ids"` + } + if !decodeDeviceJSON(ctx, &request) { + return + } + result, replayed, err := handler.services.Results.Fail( + ctx.Request.Context(), + usecase.FailExecutionCommand{ + Identity: handler.executionResultIdentity( + ctx, principal, request.ExecutionID, request.ClaimGeneration, + ), + ErrorCode: request.Error.Code, + ErrorMessage: request.Error.Message, + ErrorStep: request.Error.Step, + Retryable: request.Error.Retryable, + EvidenceAssetIDs: request.EvidenceAssetIDs, + }, + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + ctx.Header("Cache-Control", "no-store") + ctx.JSON(http.StatusOK, gin.H{ + "task": deviceTaskResponse(result), + "replayed": replayed, + }) +} + +func (handler *deviceHandlers) executionResultIdentity( + ctx *gin.Context, + principal domain.AuthPrincipal, + executionID string, + claimGeneration int64, +) usecase.ExecutionResultIdentity { + return usecase.ExecutionResultIdentity{ + UserID: principal.UserID, + DeviceID: principal.DeviceID, + TaskID: ctx.Param("id"), + ExecutionID: executionID, + ClaimGeneration: claimGeneration, + ClaimToken: ctx.GetHeader(claimTokenHeader), + IdempotencyKey: ctx.GetHeader("Idempotency-Key"), + } +} + +func isEvidenceMediaType(value string) bool { + return hasEvidenceMediaType(value, "image/jpeg") || + hasEvidenceMediaType(value, "image/png") || + hasEvidenceMediaType(value, "image/webp") +} + +func hasEvidenceMediaType(value string, expected string) bool { + return strings.EqualFold(strings.TrimSpace(strings.Split(value, ";")[0]), expected) +} + +func deviceEvidenceResponse(evidence domain.ExecutionEvidenceAsset) gin.H { + return gin.H{ + "id": evidence.ID, + "media_type": evidence.MediaType, + "size_bytes": evidence.SizeBytes, + "sha256": evidence.SHA256, + "created_at": formatTime(evidence.CreatedAt), + "received_after_execution_expiry": evidence.ReceivedAfterExecutionExpiry, + } +} + type lifecycleTransitionRequest struct { DeviceID string `json:"device_id"` ClaimGeneration int64 `json:"claim_generation"` diff --git a/backend-api/internal/transport/httpapi/device_handlers_test.go b/backend-api/internal/transport/httpapi/device_handlers_test.go index 0f9e30c..6abe567 100644 --- a/backend-api/internal/transport/httpapi/device_handlers_test.go +++ b/backend-api/internal/transport/httpapi/device_handlers_test.go @@ -15,6 +15,7 @@ import ( "net/http" "net/http/httptest" "path/filepath" + "strconv" "strings" "sync" "testing" @@ -478,6 +479,141 @@ func TestDeviceStartAndTaskHeartbeatUseClaimContract(t *testing.T) { assertNoClaimSecret(t, taskHeartbeat, testOpaqueToken) } +func TestDeviceExecutionResultsAreIdempotentAndAuditable(t *testing.T) { + fixture := newDeviceHTTPFixture(t) + requireDeviceStatus(t, fixture.readyHeartbeat(t), http.StatusOK) + taskID := fixture.createPendingTask(t) + claimedResponse := fixture.claimNext(t, "claim-for-results", testOpaqueToken) + requireDeviceStatus(t, claimedResponse, http.StatusOK) + var claimed deviceLifecycleResponse + decodeResponse(t, claimedResponse, &claimed) + + startResponse := performDeviceRequest(t, fixture.router, deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/start", + contentType: "application/json", + body: strings.NewReader(fmt.Sprintf( + `{"claim_generation":%d,"expected_version":%d}`, + claimed.Task.ClaimGeneration, + claimed.Task.Version, + )), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + idempotencyKey: "start-for-results", + }) + requireDeviceStatus(t, startResponse, http.StatusOK) + var started deviceLifecycleResponse + decodeResponse(t, startResponse, &started) + + occurredAt := time.Now().UTC().Format(time.RFC3339Nano) + events := performDeviceRequest(t, fixture.router, deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/events", + contentType: "application/json", + body: strings.NewReader(fmt.Sprintf( + `{"execution_id":%q,"claim_generation":%d,"events":[{"event_id":"00000000-0000-4000-8000-000000000701","step":"SEARCH","type":"SEARCH_STARTED","message":"开始采集候选","occurred_at":%q}]}`, + started.Execution.ID, + started.Task.ClaimGeneration, + occurredAt, + )), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + idempotencyKey: "result-events-1", + }) + requireDeviceStatus(t, events, http.StatusOK) + + evidence := performDeviceRequest(t, fixture.router, deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/evidence", + contentType: "image/jpeg", + body: deviceReferenceImage(t, 701), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + idempotencyKey: "result-evidence-1", + executionID: started.Execution.ID, + claimGeneration: started.Task.ClaimGeneration, + }) + requireDeviceStatus(t, evidence, http.StatusCreated) + var evidenceResponse struct { + Evidence struct { + ID string `json:"id"` + } `json:"evidence"` + } + decodeResponse(t, evidence, &evidenceResponse) + if evidenceResponse.Evidence.ID == "" { + t.Fatalf("evidence response = %s", evidence.Body.String()) + } + + detail, err := fixture.tasks.Get(context.Background(), "local-admin", taskID) + if err != nil { + t.Fatalf("get task for content hash: %v", err) + } + taskHash := usecase.TaskContentSHA256(detail.Task) + candidatePayload := fmt.Sprintf( + `{"execution_id":%q,"claim_generation":%d,"task_content_sha256":%q,"execution_mode":"MANUAL_FIRST","search_query":"TEST-SKU","candidates":[{"ordinal":1,"title":"手动候选","sku_text":"TEST-SKU","price":"12.00","product_url":"https://example.test/product/1","image_url":"https://example.test/image/1.jpg","evidence_asset_ids":[%q],"evaluation":null}]}`, + started.Execution.ID, + started.Task.ClaimGeneration, + taskHash, + evidenceResponse.Evidence.ID, + ) + candidates := performDeviceRequest(t, fixture.router, deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/candidates", + contentType: "application/json", + body: strings.NewReader(candidatePayload), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + idempotencyKey: "result-candidates-1", + }) + requireDeviceStatus(t, candidates, http.StatusOK) + + completePayload := fmt.Sprintf( + `{"execution_id":%q,"claim_generation":%d,"task_content_sha256":%q,"execution_mode":"MANUAL_FIRST","outcome":"CANDIDATE_ACCEPTED","operator_reason":"人工核对标题、SKU和截图后接受","candidate":{"ordinal":1,"title":"手动候选","sku_text":"TEST-SKU","price":"12.00","product_url":"https://example.test/product/1","image_url":"https://example.test/image/1.jpg","evidence_asset_ids":[%q],"evaluation":null},"order_submitted":false}`, + started.Execution.ID, + started.Task.ClaimGeneration, + taskHash, + evidenceResponse.Evidence.ID, + ) + complete := performDeviceRequest(t, fixture.router, deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/complete", + contentType: "application/json", + body: strings.NewReader(completePayload), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + idempotencyKey: "result-complete-1", + }) + requireDeviceStatus(t, complete, http.StatusOK) + assertNoClaimSecret(t, complete, testOpaqueToken) + replayed := performDeviceRequest(t, fixture.router, deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/complete", + contentType: "application/json", + body: strings.NewReader(completePayload), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + idempotencyKey: "result-complete-1", + }) + requireDeviceStatus(t, replayed, http.StatusOK) + if !strings.Contains(replayed.Body.String(), `"replayed":true`) { + t.Fatalf("terminal replay response = %s", replayed.Body.String()) + } + + detail, err = fixture.tasks.Get(context.Background(), "local-admin", taskID) + if err != nil { + t.Fatalf("get task result detail: %v", err) + } + if detail.Task.Status != domain.TaskStatusSucceeded || + detail.Report == nil || + detail.Report.Outcome == nil || + detail.Report.Outcome.OrderSubmitted || + len(detail.Report.Events) != 1 || + len(detail.Report.EvidenceAssets) != 1 || + detail.Report.CandidateBatch == nil { + t.Fatalf("execution report = %+v", detail.Report) + } +} + func TestDeviceReleaseReturnsClaimedTaskToPending(t *testing.T) { fixture := newDeviceHTTPFixture(t) requireDeviceStatus(t, fixture.readyHeartbeat(t), http.StatusOK) @@ -710,10 +846,15 @@ func newDeviceHTTPFixture(t *testing.T) *deviceHTTPFixture { if err != nil { t.Fatalf("usecase.NewLifecycleService() error = %v", err) } + results, err := usecase.NewExecutionResultService(store, files, clock, ids) + if err != nil { + t.Fatalf("usecase.NewExecutionResultService() error = %v", err) + } deviceRoutes, err := NewDeviceRouteRegistrar( DeviceServices{ Lifecycle: lifecycle, Assets: assets, + Results: results, }, ) if err != nil { @@ -880,13 +1021,15 @@ func deviceReferenceImage(t *testing.T, index int) io.Reader { } type deviceRequest struct { - method string - target string - contentType string - body io.Reader - bearerToken string - claimToken string - idempotencyKey string + method string + target string + contentType string + body io.Reader + bearerToken string + claimToken string + idempotencyKey string + executionID string + claimGeneration int64 } func performDeviceRequest( @@ -911,6 +1054,13 @@ func performDeviceRequest( if spec.idempotencyKey != "" { request.Header.Set("Idempotency-Key", spec.idempotencyKey) } + if spec.executionID != "" { + request.Header.Set("X-Execution-ID", spec.executionID) + request.Header.Set( + "X-Claim-Generation", + strconv.FormatInt(spec.claimGeneration, 10), + ) + } response := httptest.NewRecorder() router.ServeHTTP(response, request) return response diff --git a/backend-api/internal/transport/webui/handler.go b/backend-api/internal/transport/webui/handler.go index 014ab12..bffb81a 100644 --- a/backend-api/internal/transport/webui/handler.go +++ b/backend-api/internal/transport/webui/handler.go @@ -756,6 +756,7 @@ type taskDetailView struct { UpdatedAt time.Time CanCancel bool CancelRequiresAck bool + ExecutionReport *ExecutionReport } type taskDetailPage struct { @@ -789,6 +790,7 @@ func taskDetailViewFrom(task Task) taskDetailView { CanCancel: canCancelTaskStatus(task.Status), CancelRequiresAck: task.Status == "RUNNING" || task.Status == "WAITING_CONFIRMATION", + ExecutionReport: task.ExecutionReport, } } diff --git a/backend-api/internal/transport/webui/static/admin.css b/backend-api/internal/transport/webui/static/admin.css index ef167f0..9640b67 100644 --- a/backend-api/internal/transport/webui/static/admin.css +++ b/backend-api/internal/transport/webui/static/admin.css @@ -45,6 +45,62 @@ textarea { font: inherit; } +.execution-audit { + margin-top: 20px; +} + +.audit-summary { + margin-bottom: 18px; +} + +.audit-list { + margin: 8px 0 0; + padding-left: 20px; +} + +.audit-list li { + margin: 6px 0; + overflow-wrap: anywhere; +} + +.audit-evidence-grid { + display: grid; + grid-template-columns: repeat(auto-fit, minmax(180px, 1fr)); + gap: 12px; +} + +.audit-evidence { + margin: 8px 0; +} + +.audit-evidence img { + display: block; + width: 100%; + max-height: 240px; + object-fit: contain; + border: 1px solid var(--line); + background: var(--surface-soft); +} + +.audit-evidence figcaption { + margin-top: 6px; + color: var(--muted); + font-size: 12px; + overflow-wrap: anywhere; +} + +.audit-json { + max-height: 360px; + margin: 8px 0 16px; + overflow: auto; + padding: 12px; + border: 1px solid var(--line); + background: var(--surface-soft); + font: 12px/1.45 ui-monospace, SFMono-Regular, Consolas, monospace; + white-space: pre-wrap; + overflow-wrap: anywhere; +} + button, input, select { diff --git a/backend-api/internal/transport/webui/templates/task-detail.gohtml b/backend-api/internal/transport/webui/templates/task-detail.gohtml index 4809ab5..674fe73 100644 --- a/backend-api/internal/transport/webui/templates/task-detail.gohtml +++ b/backend-api/internal/transport/webui/templates/task-detail.gohtml @@ -73,6 +73,39 @@ + {{with .Task.ExecutionReport}} +
+

执行审计

+ {{if .Outcome}} +
+
结果类型
{{.Outcome.ResultType}}
+
人工理由
{{if .Outcome.OperatorReason}}{{.Outcome.OperatorReason}}{{else}}未提供{{end}}
+
订单提交
{{if .Outcome.OrderSubmitted}}是{{else}}否{{end}}
+ {{if .Outcome.ErrorCode}}
失败信息
{{.Outcome.ErrorCode}}:{{.Outcome.ErrorMessage}}
{{end}} +
+ {{end}} + {{if .Mode}}

模式:{{.Mode}} · 搜索词:{{.SearchQuery}}

{{end}} + {{if .Provenance}}

本地模型出处

{{.Provenance}}
{{end}} + {{if .Candidates}}

候选与评估

{{.Candidates}}
{{end}} + {{if .Recommendation}}

本地推荐

{{.Recommendation}}
{{end}} + {{if .Evidence}} +

证据截图

+
+ {{range .Evidence}}
+ 执行证据截图 +
SHA-256 {{.SHA256}}({{.SizeBytes}} bytes){{if .ReceivedAfterExecutionExpiry}},授权到期后补报{{end}}
+
{{end}} +
+ {{end}} + {{if .Events}} +

执行事件

+
    + {{range .Events}}
  • · {{.Step}} · {{.Message}}{{if .ReceivedAfterExecutionExpiry}}(授权到期后补报){{end}}
  • {{end}} +
+ {{end}} +
+ {{end}} + {{if .Task.CanCancel}}