feat(t207): audit local procurement results
This commit is contained in:
@@ -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)
|
||||
Reference in New Issue
Block a user