959 lines
26 KiB
Go
959 lines
26 KiB
Go
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 := storeCandidateDecisionDataset(
|
|
ctx,
|
|
tx,
|
|
write,
|
|
candidate,
|
|
expired,
|
|
); err != nil {
|
|
return false, err
|
|
}
|
|
if candidate.ReadyEvent != nil && !expired {
|
|
result, err := tx.ExecContext(
|
|
ctx,
|
|
`UPDATE purchase_tasks
|
|
SET status = 'WAITING_CONFIRMATION',
|
|
version = version + 1,
|
|
updated_at = ?
|
|
WHERE id = ? AND status = 'RUNNING'
|
|
AND version = ? AND cancel_requested_at IS NULL`,
|
|
formatTimestamp(write.Now),
|
|
write.TaskID,
|
|
task.Version,
|
|
)
|
|
if err != nil {
|
|
return false, repositoryFailure(err)
|
|
}
|
|
affected, err := result.RowsAffected()
|
|
if err != nil {
|
|
return false, repositoryFailure(err)
|
|
}
|
|
if affected != 1 {
|
|
return false, usecase.ErrTaskStateConflict
|
|
}
|
|
if err := insertTaskEvent(ctx, tx, *candidate.ReadyEvent); err != nil {
|
|
return false, 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 []usecase.ExecutionCandidate
|
|
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
|
|
}
|
|
if len(candidate.EvidenceAssetIDs) != 2 {
|
|
return usecase.ErrTaskStateConflict
|
|
}
|
|
expectedHashes := []string{
|
|
candidate.DetailEvidenceSHA256,
|
|
candidate.SpecificationEvidenceSHA256,
|
|
}
|
|
for index, evidenceID := range candidate.EvidenceAssetIDs {
|
|
var actualHash string
|
|
err := tx.QueryRowContext(
|
|
ctx,
|
|
`SELECT sha256 FROM execution_evidence_assets
|
|
WHERE id = ? AND task_id = ? AND execution_id = ?`,
|
|
evidenceID,
|
|
write.TaskID,
|
|
write.ExecutionID,
|
|
).Scan(&actualHash)
|
|
if err != nil {
|
|
return repositoryFailure(err)
|
|
}
|
|
if actualHash != expectedHashes[index] {
|
|
return usecase.ErrTaskStateConflict
|
|
}
|
|
}
|
|
}
|
|
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)
|
|
}
|
|
report.DecisionDataset, err = getCandidateDecisionDataset(
|
|
ctx,
|
|
queryer,
|
|
taskID,
|
|
executionID,
|
|
)
|
|
if err != nil {
|
|
return nil, 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)
|