Files
cmroubao/backend-api/internal/repository/sqlite/execution_result_repository.go
T

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)