feat(t218): reconcile single order submissions

This commit is contained in:
QiuSW
2026-07-28 17:48:54 +08:00
parent 682cc7327c
commit a77d8be06d
32 changed files with 4306 additions and 79 deletions
+33
View File
@@ -355,6 +355,39 @@ type OrderDryRun struct {
ReadyAt *time.Time
}
type OrderSubmissionStatus string
const (
OrderSubmissionFenced OrderSubmissionStatus = "FENCED"
OrderSubmissionReconciled OrderSubmissionStatus = "RECONCILED"
OrderSubmissionManualReview OrderSubmissionStatus = "MANUAL_REVIEW"
)
type OrderSubmission struct {
ID string
AuthorizationID string
DryRunID string
TaskID string
ExecutionID string
CommandSHA256 string
DryRunEvidenceSHA256 string
Status OrderSubmissionStatus
ExpectedTitle string
ExpectedSKU string
ExpectedQuantity int
ExpectedUnitPriceCents int64
ExpectedTotalPriceCents int64
PlatformOrderNo *string
PlatformOrderedAt *time.Time
PlatformOrderStatus *string
ReconciliationEvidenceAssetID *string
ReconciliationEvidenceSHA256 *string
ManualReasonCode *string
FencedAt time.Time
ReconciledAt *time.Time
ManualReviewAt *time.Time
}
type ExecutionReport struct {
Events []ExecutionEvent
EvidenceAssets []ExecutionEvidenceAsset
@@ -34,8 +34,11 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) {
if applied, err := runner.Up(ctx); err != nil {
t.Fatalf("initial Up() error = %v", err)
} else if applied != 10 {
t.Fatalf("initial Up() applied = %d, want 10", applied)
} else if applied != 11 {
t.Fatalf("initial Up() applied = %d, want 11", applied)
}
if err := runner.Down(ctx); err != nil {
t.Fatalf("initial Down(v11) error = %v", err)
}
if err := runner.Down(ctx); err != nil {
t.Fatalf("initial Down(v10) error = %v", err)
@@ -59,9 +62,14 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) {
seedClaimsHistoricalFixture(t, db)
if applied, err := runner.Up(ctx); err != nil {
t.Fatalf("Up(v5-v10) over historical data error = %v", err)
} else if applied != 6 {
t.Fatalf("Up(v5-v10) applied = %d, want 6", applied)
t.Fatalf("Up(v5-v11) over historical data error = %v", err)
} else if applied != 7 {
t.Fatalf("Up(v5-v11) applied = %d, want 7", applied)
}
assertClaimsHistory(t, db, true)
if err := runner.Down(ctx); err != nil {
t.Fatalf("Down(v11) with compatible history error = %v", err)
}
assertClaimsHistory(t, db, true)
@@ -101,9 +109,9 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) {
assertClaimsHistory(t, db, false)
if applied, err := runner.Up(ctx); err != nil {
t.Fatalf("final Up(v4-v10) error = %v", err)
} else if applied != 7 {
t.Fatalf("final Up(v4-v10) applied = %d, want 7", applied)
t.Fatalf("final Up(v4-v11) error = %v", err)
} else if applied != 8 {
t.Fatalf("final Up(v4-v11) applied = %d, want 8", applied)
}
assertClaimsHistory(t, db, true)
}
@@ -339,6 +347,9 @@ func TestClaimsMigrationDownFailsClosedForNewAuditData(t *testing.T) {
t.Fatalf("insert v4 audit event: %v", err)
}
if err := runner.Down(ctx); err != nil {
t.Fatalf("Down(v11) error = %v", err)
}
if err := runner.Down(ctx); err != nil {
t.Fatalf("Down(v10) error = %v", err)
}
@@ -27,8 +27,8 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) {
if err != nil {
t.Fatalf("Up() error = %v", err)
}
if applied != 10 {
t.Fatalf("Up() applied = %d, want 10", applied)
if applied != 11 {
t.Fatalf("Up() applied = %d, want 11", applied)
}
assertStatuses(t, runner, map[int64]bool{
1: true,
@@ -41,6 +41,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) {
8: true,
9: true,
10: true,
11: true,
})
applied, err = runner.Up(context.Background())
@@ -64,7 +65,8 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) {
7: true,
8: true,
9: true,
10: false,
10: true,
11: false,
})
applied, err = runner.Up(context.Background())
@@ -85,6 +87,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) {
8: true,
9: true,
10: true,
11: true,
})
}
@@ -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(v11) error = %v", err)
}
if err := runner.Down(context.Background()); err != nil {
t.Fatalf("Down(v10) error = %v", err)
}
@@ -420,9 +423,9 @@ func TestAuthMigrationCanRollbackWithoutRebuildingPurchaseTasks(
t.Fatal("purchase_tasks was lost during auth migration rollback")
}
if applied, err := runner.Up(context.Background()); err != nil {
t.Fatalf("Up(v3-v10) error = %v", err)
} else if applied != 8 {
t.Fatalf("Up(v3-v10) applied = %d, want 8", applied)
t.Fatalf("Up(v3-v11) error = %v", err)
} else if applied != 9 {
t.Fatalf("Up(v3-v11) applied = %d, want 9", applied)
}
}
@@ -0,0 +1,880 @@
package sqlite
import (
"context"
"database/sql"
"encoding/json"
"errors"
"strings"
"time"
"unicode"
"cmroubao/backend-api/internal/domain"
"cmroubao/backend-api/internal/usecase"
)
const orderReconciliationClockWindow = 5 * time.Minute
func (s *Store) StartOrderSubmission(
ctx context.Context,
write usecase.StartOrderSubmissionWrite,
) (domain.OrderSubmission, bool, error) {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
defer func() { _ = tx.Rollback() }()
if result, found, err := replayOrderSubmissionRequest(
ctx,
tx,
write.DeviceID,
"START",
write.IdempotencyKey,
write.RequestSHA256,
write.TaskID,
); err != nil || found {
return commitOrderSubmissionReplay(tx, result, found, err)
}
if err := validateDeviceOrderCommandClaim(
ctx,
tx,
write.UserID,
write.DeviceID,
write.TaskID,
write.ExecutionID,
write.ClaimGeneration,
write.ClaimTokenHash,
write.Now,
); err != nil {
return domain.OrderSubmission{}, false, err
}
authorization, err := getOrderAuthorization(
ctx,
tx,
write.AuthorizationID,
)
if err != nil {
return domain.OrderSubmission{}, false, err
}
if err := validateSubmissionAuthorization(
authorization,
write.UserID,
write.DeviceID,
write.TaskID,
write.ExecutionID,
write.ClaimGeneration,
write.CommandSHA256,
); err != nil {
return domain.OrderSubmission{}, false, err
}
if authorization.Status != domain.OrderAuthorizationExecuting {
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
dryRun, found, err := getOrderDryRunByAuthorization(
ctx,
tx,
authorization.ID,
)
if err != nil {
return domain.OrderSubmission{}, false, err
}
if !found ||
dryRun.ID != write.DryRunID ||
dryRun.Status != domain.OrderDryRunReady ||
dryRun.CommandSHA256 != write.CommandSHA256 ||
dryRun.EvidenceSHA256 == nil ||
*dryRun.EvidenceSHA256 != write.DryRunEvidenceSHA256 ||
dryRun.ObservedTitle == nil ||
normalizeSubmissionText(*dryRun.ObservedTitle) !=
normalizeSubmissionText(write.ObservedTitle) ||
dryRun.SelectedSKU == nil ||
strings.TrimSpace(*dryRun.SelectedSKU) != write.SelectedSKU ||
dryRun.Quantity == nil ||
*dryRun.Quantity != write.Quantity ||
dryRun.UnitPriceCents == nil ||
*dryRun.UnitPriceCents != write.UnitPriceCents ||
dryRun.TotalPriceCents == nil ||
*dryRun.TotalPriceCents != write.TotalPriceCents {
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
if _, found, err := getOrderSubmissionByAuthorization(
ctx,
tx,
authorization.ID,
); err != nil {
return domain.OrderSubmission{}, false, err
} else if found {
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
_, err = tx.ExecContext(
ctx,
`INSERT INTO order_submissions (
id, authorization_id, dry_run_id, task_id, execution_id,
command_sha256, dry_run_evidence_sha256, status,
expected_title, expected_sku, expected_quantity,
expected_unit_price_cents, expected_total_price_cents, fenced_at
) VALUES (?, ?, ?, ?, ?, ?, ?, 'FENCED', ?, ?, ?, ?, ?, ?)`,
write.SubmissionID,
authorization.ID,
dryRun.ID,
write.TaskID,
write.ExecutionID,
write.CommandSHA256,
write.DryRunEvidenceSHA256,
write.ObservedTitle,
write.SelectedSKU,
write.Quantity,
write.UnitPriceCents,
write.TotalPriceCents,
formatTimestamp(write.Now),
)
if err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
if err := insertTaskEvent(ctx, tx, write.Event); err != nil {
return domain.OrderSubmission{}, false, err
}
if err := insertOrderSubmissionRequest(
ctx,
tx,
write.DeviceID,
"START",
write.IdempotencyKey,
write.RequestSHA256,
write.TaskID,
write.SubmissionID,
write.Now,
); err != nil {
return domain.OrderSubmission{}, false, err
}
submission, err := getOrderSubmissionByID(
ctx,
tx,
write.SubmissionID,
)
if err != nil {
return domain.OrderSubmission{}, false, err
}
if err := tx.Commit(); err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
return submission, false, nil
}
func (s *Store) ReconcileOrderSubmission(
ctx context.Context,
write usecase.ReconcileOrderSubmissionWrite,
) (domain.OrderSubmission, bool, error) {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
defer func() { _ = tx.Rollback() }()
if result, found, err := replayOrderSubmissionRequest(
ctx,
tx,
write.DeviceID,
"RECONCILE",
write.IdempotencyKey,
write.RequestSHA256,
write.TaskID,
); err != nil || found {
return commitOrderSubmissionReplay(tx, result, found, err)
}
if err := validateFencedSubmissionOwner(
ctx,
tx,
write.UserID,
write.DeviceID,
write.TaskID,
write.ExecutionID,
write.ClaimGeneration,
write.ClaimTokenHash,
); err != nil {
return domain.OrderSubmission{}, false, err
}
authorization, err := getOrderAuthorization(
ctx,
tx,
write.AuthorizationID,
)
if err != nil {
return domain.OrderSubmission{}, false, err
}
if err := validateSubmissionAuthorization(
authorization,
write.UserID,
write.DeviceID,
write.TaskID,
write.ExecutionID,
write.ClaimGeneration,
write.CommandSHA256,
); err != nil {
return domain.OrderSubmission{}, false, err
}
if authorization.Status != domain.OrderAuthorizationExecuting {
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
submission, err := getOrderSubmissionByID(
ctx,
tx,
write.SubmissionID,
)
if err != nil {
return domain.OrderSubmission{}, false, err
}
if err := validateSubmissionRecord(
submission,
authorization.ID,
write.TaskID,
write.ExecutionID,
write.CommandSHA256,
); err != nil {
return domain.OrderSubmission{}, false, err
}
if submission.Status != domain.OrderSubmissionFenced &&
submission.Status != domain.OrderSubmissionManualReview {
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
if normalizeSubmissionText(write.ObservedTitle) !=
normalizeSubmissionText(submission.ExpectedTitle) ||
write.SelectedSKU != submission.ExpectedSKU ||
write.Quantity != submission.ExpectedQuantity ||
write.TotalPriceCents != submission.ExpectedTotalPriceCents ||
write.PlatformOrderStatus != "PENDING_PAYMENT" {
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
if write.ParsedPlatformOrderedAt.Before(
submission.FencedAt.Add(-orderReconciliationClockWindow),
) || write.ParsedPlatformOrderedAt.After(
write.Now.Add(orderReconciliationClockWindow),
) {
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
evidence, err := getExecutionEvidence(
ctx,
tx,
write.EvidenceAssetID,
)
if err != nil {
return domain.OrderSubmission{}, false, err
}
if evidence.TaskID != write.TaskID ||
evidence.ExecutionID != write.ExecutionID ||
evidence.MediaType != "image/jpeg" ||
evidence.SHA256 != write.EvidenceSHA256 {
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
result, err := tx.ExecContext(
ctx,
`UPDATE order_submissions
SET status = 'RECONCILED',
platform_order_no = ?,
platform_ordered_at = ?,
platform_order_status = 'PENDING_PAYMENT',
reconciliation_evidence_asset_id = ?,
reconciliation_evidence_sha256 = ?,
manual_reason_code = NULL,
reconciled_at = ?,
manual_review_at = NULL
WHERE id = ? AND status IN ('FENCED', 'MANUAL_REVIEW')`,
write.PlatformOrderNo,
formatTimestamp(write.ParsedPlatformOrderedAt),
write.EvidenceAssetID,
write.EvidenceSHA256,
formatTimestamp(write.Now),
submission.ID,
)
if err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
if affected, err := result.RowsAffected(); err != nil || affected != 1 {
if err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
result, err = tx.ExecContext(
ctx,
`UPDATE order_authorizations
SET status = 'CONSUMED', consumed_at = ?
WHERE id = ? AND status = 'EXECUTING'`,
formatTimestamp(write.Now),
authorization.ID,
)
if err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
if affected, err := result.RowsAffected(); err != nil || affected != 1 {
if err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
if err := finishReconciledOrderExecution(
ctx,
tx,
write,
authorization,
); err != nil {
return domain.OrderSubmission{}, false, err
}
if err := insertTaskEvent(ctx, tx, write.Event); err != nil {
return domain.OrderSubmission{}, false, err
}
if err := insertOrderSubmissionRequest(
ctx,
tx,
write.DeviceID,
"RECONCILE",
write.IdempotencyKey,
write.RequestSHA256,
write.TaskID,
submission.ID,
write.Now,
); err != nil {
return domain.OrderSubmission{}, false, err
}
submission, err = getOrderSubmissionByID(ctx, tx, submission.ID)
if err != nil {
return domain.OrderSubmission{}, false, err
}
if err := tx.Commit(); err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
return submission, false, nil
}
func (s *Store) ManualReviewOrderSubmission(
ctx context.Context,
write usecase.ManualReviewOrderSubmissionWrite,
) (domain.OrderSubmission, bool, error) {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
defer func() { _ = tx.Rollback() }()
if result, found, err := replayOrderSubmissionRequest(
ctx,
tx,
write.DeviceID,
"MANUAL_REVIEW",
write.IdempotencyKey,
write.RequestSHA256,
write.TaskID,
); err != nil || found {
return commitOrderSubmissionReplay(tx, result, found, err)
}
if err := validateFencedSubmissionOwner(
ctx,
tx,
write.UserID,
write.DeviceID,
write.TaskID,
write.ExecutionID,
write.ClaimGeneration,
write.ClaimTokenHash,
); err != nil {
return domain.OrderSubmission{}, false, err
}
authorization, err := getOrderAuthorization(
ctx,
tx,
write.AuthorizationID,
)
if err != nil {
return domain.OrderSubmission{}, false, err
}
if err := validateSubmissionAuthorization(
authorization,
write.UserID,
write.DeviceID,
write.TaskID,
write.ExecutionID,
write.ClaimGeneration,
write.CommandSHA256,
); err != nil {
return domain.OrderSubmission{}, false, err
}
if authorization.Status != domain.OrderAuthorizationExecuting {
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
submission, err := getOrderSubmissionByID(
ctx,
tx,
write.SubmissionID,
)
if err != nil {
return domain.OrderSubmission{}, false, err
}
if err := validateSubmissionRecord(
submission,
authorization.ID,
write.TaskID,
write.ExecutionID,
write.CommandSHA256,
); err != nil {
return domain.OrderSubmission{}, false, err
}
if submission.Status != domain.OrderSubmissionFenced {
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
var evidenceID, evidenceSHA any
if write.EvidenceAssetID != "" {
evidence, err := getExecutionEvidence(
ctx,
tx,
write.EvidenceAssetID,
)
if err != nil {
return domain.OrderSubmission{}, false, err
}
if evidence.TaskID != write.TaskID ||
evidence.ExecutionID != write.ExecutionID ||
evidence.MediaType != "image/jpeg" ||
evidence.SHA256 != write.EvidenceSHA256 {
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
evidenceID = write.EvidenceAssetID
evidenceSHA = write.EvidenceSHA256
}
result, err := tx.ExecContext(
ctx,
`UPDATE order_submissions
SET status = 'MANUAL_REVIEW',
reconciliation_evidence_asset_id = ?,
reconciliation_evidence_sha256 = ?,
manual_reason_code = ?,
manual_review_at = ?
WHERE id = ? AND status = 'FENCED'`,
evidenceID,
evidenceSHA,
write.ReasonCode,
formatTimestamp(write.Now),
submission.ID,
)
if err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
if affected, err := result.RowsAffected(); err != nil || affected != 1 {
if err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
return domain.OrderSubmission{}, false, usecase.ErrTaskStateConflict
}
if err := insertTaskEvent(ctx, tx, write.Event); err != nil {
return domain.OrderSubmission{}, false, err
}
if err := insertOrderSubmissionRequest(
ctx,
tx,
write.DeviceID,
"MANUAL_REVIEW",
write.IdempotencyKey,
write.RequestSHA256,
write.TaskID,
submission.ID,
write.Now,
); err != nil {
return domain.OrderSubmission{}, false, err
}
submission, err = getOrderSubmissionByID(ctx, tx, submission.ID)
if err != nil {
return domain.OrderSubmission{}, false, err
}
if err := tx.Commit(); err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
return submission, false, nil
}
func validateSubmissionAuthorization(
authorization domain.OrderAuthorization,
userID, deviceID, taskID, executionID string,
claimGeneration int64,
commandSHA256 string,
) error {
if authorization.UserID != userID ||
authorization.DeviceID != deviceID ||
authorization.TaskID != taskID ||
authorization.ExecutionID != executionID ||
authorization.ClaimGeneration != claimGeneration ||
authorization.CommandSHA256 == nil ||
*authorization.CommandSHA256 != commandSHA256 {
return usecase.ErrExecutionMismatch
}
return nil
}
func validateFencedSubmissionOwner(
ctx context.Context,
tx *sql.Tx,
userID, deviceID, taskID, executionID string,
claimGeneration int64,
claimTokenHash string,
) error {
task, err := getClaimProtectedTask(ctx, tx, taskID)
if err != nil {
return err
}
if err := validateClaimOwner(
task,
userID,
deviceID,
claimGeneration,
claimTokenHash,
); err != nil {
return err
}
if task.Status != domain.TaskStatusWaitingConfirmation ||
task.CancelRequestedAt != nil {
return usecase.ErrTaskStateConflict
}
execution, err := getExecutionByID(ctx, tx, executionID)
if err != nil {
return err
}
if execution.TaskID != taskID ||
execution.UserID != userID ||
execution.DeviceID != deviceID ||
execution.ClaimGeneration != claimGeneration ||
execution.FinishedAt != nil {
return usecase.ErrExecutionMismatch
}
return nil
}
func validateSubmissionRecord(
submission domain.OrderSubmission,
authorizationID, taskID, executionID, commandSHA256 string,
) error {
if submission.AuthorizationID != authorizationID ||
submission.TaskID != taskID ||
submission.ExecutionID != executionID ||
submission.CommandSHA256 != commandSHA256 {
return usecase.ErrExecutionMismatch
}
return nil
}
func finishReconciledOrderExecution(
ctx context.Context,
tx *sql.Tx,
write usecase.ReconcileOrderSubmissionWrite,
authorization domain.OrderAuthorization,
) error {
var executionMode string
err := tx.QueryRowContext(
ctx,
`SELECT COALESCE(
(SELECT execution_mode
FROM execution_candidate_batches
WHERE execution_id = ? AND task_id = ?),
(SELECT execution_mode
FROM candidate_search_runs
WHERE execution_id = ? AND task_id = ?)
)`,
write.ExecutionID,
write.TaskID,
write.ExecutionID,
write.TaskID,
).Scan(&executionMode)
if errors.Is(err, sql.ErrNoRows) {
return usecase.ErrRepositoryInvariant
}
if err != nil {
return repositoryFailure(err)
}
evidenceIDs, err := json.Marshal([]string{write.EvidenceAssetID})
if err != nil {
return repositoryFailure(err)
}
expired := false
var claimExpiresAt sql.NullString
if err := tx.QueryRowContext(
ctx,
`SELECT claim_expires_at FROM purchase_tasks WHERE id = ?`,
write.TaskID,
).Scan(&claimExpiresAt); err != nil {
return repositoryFailure(err)
}
if !claimExpiresAt.Valid {
expired = true
} else {
expiresAt, err := parseTimestamp(claimExpiresAt.String)
if err != nil {
return repositoryFailure(err)
}
expired = !expiresAt.After(write.Now)
}
_, 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 (?, ?, 'COMPLETE', ?, ?, 'CANDIDATE_ACCEPTED',
?, NULL, ?, NULL, NULL, NULL, NULL, 1, ?, ?)`,
write.ExecutionID,
write.TaskID,
executionMode,
authorization.TaskContentSHA256,
"pending-payment order uniquely reconciled",
string(evidenceIDs),
formatTimestamp(write.Now),
expired,
)
if err != nil {
return repositoryFailure(err)
}
result, err := tx.ExecContext(
ctx,
`UPDATE task_executions
SET current_step = 'ORDER_RECONCILED',
last_heartbeat_at = ?,
finished_at = ?
WHERE id = ? AND finished_at IS NULL`,
formatTimestamp(write.Now),
formatTimestamp(write.Now),
write.ExecutionID,
)
if err != nil {
return repositoryFailure(err)
}
if affected, err := result.RowsAffected(); err != nil || affected != 1 {
if err != nil {
return repositoryFailure(err)
}
return usecase.ErrTaskStateConflict
}
result, err = tx.ExecContext(
ctx,
`UPDATE purchase_tasks
SET status = 'SUCCEEDED',
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 = 'WAITING_CONFIRMATION'
AND claim_generation = ?`,
formatTimestamp(write.Now),
write.TaskID,
write.ClaimGeneration,
)
if err != nil {
return repositoryFailure(err)
}
if affected, err := result.RowsAffected(); err != nil || affected != 1 {
if err != nil {
return repositoryFailure(err)
}
return usecase.ErrTaskStateConflict
}
return nil
}
func getOrderSubmissionByAuthorization(
ctx context.Context,
queryer queryRower,
authorizationID string,
) (domain.OrderSubmission, bool, error) {
submission, err := scanOrderSubmission(queryer.QueryRowContext(
ctx,
orderSubmissionSelect+` WHERE authorization_id = ?`,
authorizationID,
))
if errors.Is(err, sql.ErrNoRows) {
return domain.OrderSubmission{}, false, nil
}
if err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
return submission, true, nil
}
func getOrderSubmissionByID(
ctx context.Context,
queryer queryRower,
submissionID string,
) (domain.OrderSubmission, error) {
submission, err := scanOrderSubmission(queryer.QueryRowContext(
ctx,
orderSubmissionSelect+` WHERE id = ?`,
submissionID,
))
if errors.Is(err, sql.ErrNoRows) {
return domain.OrderSubmission{}, usecase.ErrRepositoryNotFound
}
if err != nil {
return domain.OrderSubmission{}, repositoryFailure(err)
}
return submission, nil
}
const orderSubmissionSelect = `SELECT
id, authorization_id, dry_run_id, task_id, execution_id,
command_sha256, dry_run_evidence_sha256, status,
expected_title, expected_sku, expected_quantity,
expected_unit_price_cents, expected_total_price_cents,
platform_order_no, platform_ordered_at, platform_order_status,
reconciliation_evidence_asset_id, reconciliation_evidence_sha256,
manual_reason_code, fenced_at, reconciled_at, manual_review_at
FROM order_submissions`
func scanOrderSubmission(
scanner rowScanner,
) (domain.OrderSubmission, error) {
var submission domain.OrderSubmission
var orderNo, orderedAt, orderStatus sql.NullString
var evidenceID, evidenceSHA, manualReason sql.NullString
var fencedAt string
var reconciledAt, manualReviewAt sql.NullString
err := scanner.Scan(
&submission.ID,
&submission.AuthorizationID,
&submission.DryRunID,
&submission.TaskID,
&submission.ExecutionID,
&submission.CommandSHA256,
&submission.DryRunEvidenceSHA256,
&submission.Status,
&submission.ExpectedTitle,
&submission.ExpectedSKU,
&submission.ExpectedQuantity,
&submission.ExpectedUnitPriceCents,
&submission.ExpectedTotalPriceCents,
&orderNo,
&orderedAt,
&orderStatus,
&evidenceID,
&evidenceSHA,
&manualReason,
&fencedAt,
&reconciledAt,
&manualReviewAt,
)
if err != nil {
return domain.OrderSubmission{}, err
}
submission.PlatformOrderNo = optionalSQLString(orderNo)
submission.PlatformOrderStatus = optionalSQLString(orderStatus)
submission.ReconciliationEvidenceAssetID = optionalSQLString(evidenceID)
submission.ReconciliationEvidenceSHA256 = optionalSQLString(evidenceSHA)
submission.ManualReasonCode = optionalSQLString(manualReason)
submission.FencedAt, err = parseTimestamp(fencedAt)
if err != nil {
return domain.OrderSubmission{}, err
}
if orderedAt.Valid {
parsed, parseErr := parseTimestamp(orderedAt.String)
if parseErr != nil {
return domain.OrderSubmission{}, parseErr
}
submission.PlatformOrderedAt = &parsed
}
if submission.ReconciledAt, err = parseNullableTimestamp(reconciledAt); err != nil {
return domain.OrderSubmission{}, err
}
if submission.ManualReviewAt, err = parseNullableTimestamp(manualReviewAt); err != nil {
return domain.OrderSubmission{}, err
}
return submission, nil
}
func optionalSQLString(value sql.NullString) *string {
if !value.Valid {
return nil
}
return &value.String
}
func replayOrderSubmissionRequest(
ctx context.Context,
tx *sql.Tx,
deviceID, operation, idempotencyKey, requestSHA256, taskID string,
) (domain.OrderSubmission, bool, error) {
var knownHash, knownTask, submissionID string
err := tx.QueryRowContext(
ctx,
`SELECT request_sha256, task_id, submission_id
FROM device_order_submission_requests
WHERE device_id = ? AND operation = ? AND idempotency_key = ?`,
deviceID,
operation,
idempotencyKey,
).Scan(&knownHash, &knownTask, &submissionID)
if errors.Is(err, sql.ErrNoRows) {
return domain.OrderSubmission{}, false, nil
}
if err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
if knownHash != requestSHA256 || knownTask != taskID {
return domain.OrderSubmission{}, false, usecase.ErrIdempotencyConflict
}
result, err := getOrderSubmissionByID(ctx, tx, submissionID)
if err != nil {
return domain.OrderSubmission{}, false, err
}
return result, true, nil
}
func commitOrderSubmissionReplay(
tx *sql.Tx,
result domain.OrderSubmission,
found bool,
err error,
) (domain.OrderSubmission, bool, error) {
if err != nil || !found {
return domain.OrderSubmission{}, false, err
}
if err := tx.Commit(); err != nil {
return domain.OrderSubmission{}, false, repositoryFailure(err)
}
return result, true, nil
}
func insertOrderSubmissionRequest(
ctx context.Context,
tx *sql.Tx,
deviceID, operation, idempotencyKey, requestSHA256, taskID,
submissionID string,
now time.Time,
) error {
_, err := tx.ExecContext(
ctx,
`INSERT INTO device_order_submission_requests (
device_id, operation, idempotency_key, request_sha256,
task_id, submission_id, created_at
) VALUES (?, ?, ?, ?, ?, ?, ?)`,
deviceID,
operation,
idempotencyKey,
requestSHA256,
taskID,
submissionID,
formatTimestamp(now.UTC()),
)
if err != nil {
return repositoryFailure(err)
}
return nil
}
func normalizeSubmissionText(value string) string {
return strings.Map(func(character rune) rune {
if unicode.IsSpace(character) || unicode.IsPunct(character) {
return -1
}
return unicode.ToLower(character)
}, value)
}
var _ usecase.OrderSubmissionRepository = (*Store)(nil)
@@ -410,6 +410,9 @@ func TestAdminOrderAuthorizationIsIdempotentAndRevisioned(t *testing.T) {
if err != nil {
t.Fatalf("migration.New() error = %v", err)
}
if err := runner.Down(context.Background()); err != nil {
t.Fatalf("order submission migration down: %v", err)
}
if err := runner.Down(context.Background()); err != nil {
t.Fatalf("order dry-run migration down: %v", err)
}
@@ -17,17 +17,18 @@ import (
const claimTokenHeader = "X-Claim-Token"
type DeviceServices struct {
Lifecycle *usecase.LifecycleService
Assets *usecase.AssetService
Results *usecase.ExecutionResultService
Commands *usecase.DeviceOrderCommandService
DryRuns *usecase.OrderDryRunService
Lifecycle *usecase.LifecycleService
Assets *usecase.AssetService
Results *usecase.ExecutionResultService
Commands *usecase.DeviceOrderCommandService
DryRuns *usecase.OrderDryRunService
Submissions *usecase.OrderSubmissionService
}
func (services DeviceServices) validate() error {
if services.Lifecycle == nil || services.Assets == nil ||
services.Results == nil || services.Commands == nil ||
services.DryRuns == nil {
services.DryRuns == nil || services.Submissions == nil {
return errors.New("device services are required")
}
return nil
@@ -93,12 +94,215 @@ func NewDeviceRouteRegistrar(
"/api/v1/tasks/:id/order-dry-runs/:command_id/ready",
handler.readyOrderDryRun,
)
routes.POST(
"/api/v1/tasks/:id/order-submissions/start",
handler.startOrderSubmission,
)
routes.POST(
"/api/v1/tasks/:id/order-submissions/:submission_id/reconcile",
handler.reconcileOrderSubmission,
)
routes.POST(
"/api/v1/tasks/:id/order-submissions/:submission_id/manual-review",
handler.manualReviewOrderSubmission,
)
routes.POST("/api/v1/tasks/:id/complete", handler.completeTask)
routes.POST("/api/v1/tasks/:id/fail", handler.failTask)
return nil
}, nil
}
func (handler *deviceHandlers) startOrderSubmission(ctx *gin.Context) {
principal, ok := devicePrincipal(ctx)
if !ok {
return
}
var request struct {
DeviceID string `json:"device_id"`
ExecutionID string `json:"execution_id"`
ClaimGeneration int64 `json:"claim_generation"`
CommandID string `json:"command_id"`
CommandSHA256 string `json:"command_sha256"`
DryRunID string `json:"dry_run_id"`
DryRunEvidenceSHA256 string `json:"dry_run_evidence_sha256"`
ObservedTitle string `json:"observed_title"`
SelectedSKU string `json:"selected_sku"`
Quantity int `json:"quantity"`
UnitPriceCents int64 `json:"unit_price_cents"`
TotalPriceCents int64 `json:"total_price_cents"`
}
if !decodeDeviceJSON(ctx, &request) ||
!deviceIDMatches(ctx, request.DeviceID, principal.DeviceID) {
return
}
result, err := handler.services.Submissions.Start(
ctx.Request.Context(),
usecase.StartOrderSubmissionCommand{
UserID: principal.UserID,
DeviceID: principal.DeviceID,
TaskID: ctx.Param("id"),
ExecutionID: request.ExecutionID,
AuthorizationID: request.CommandID,
ClaimGeneration: request.ClaimGeneration,
ClaimToken: ctx.GetHeader(claimTokenHeader),
CommandSHA256: request.CommandSHA256,
DryRunID: request.DryRunID,
DryRunEvidenceSHA256: request.DryRunEvidenceSHA256,
ObservedTitle: request.ObservedTitle,
SelectedSKU: request.SelectedSKU,
Quantity: request.Quantity,
UnitPriceCents: request.UnitPriceCents,
TotalPriceCents: request.TotalPriceCents,
IdempotencyKey: ctx.GetHeader("Idempotency-Key"),
},
)
if err != nil {
writeUsecaseError(ctx, err)
return
}
writeOrderSubmission(ctx, result)
}
func (handler *deviceHandlers) reconcileOrderSubmission(ctx *gin.Context) {
principal, ok := devicePrincipal(ctx)
if !ok {
return
}
var request struct {
DeviceID string `json:"device_id"`
ExecutionID string `json:"execution_id"`
ClaimGeneration int64 `json:"claim_generation"`
CommandID string `json:"command_id"`
CommandSHA256 string `json:"command_sha256"`
PlatformOrderNo string `json:"platform_order_no"`
PlatformOrderedAt string `json:"platform_ordered_at"`
PlatformOrderStatus string `json:"platform_order_status"`
ObservedTitle string `json:"observed_title"`
SelectedSKU string `json:"selected_sku"`
Quantity int `json:"quantity"`
TotalPriceCents int64 `json:"total_price_cents"`
EvidenceAssetID string `json:"evidence_asset_id"`
EvidenceSHA256 string `json:"evidence_sha256"`
}
if !decodeDeviceJSON(ctx, &request) ||
!deviceIDMatches(ctx, request.DeviceID, principal.DeviceID) {
return
}
result, err := handler.services.Submissions.Reconcile(
ctx.Request.Context(),
usecase.ReconcileOrderSubmissionCommand{
UserID: principal.UserID,
DeviceID: principal.DeviceID,
TaskID: ctx.Param("id"),
ExecutionID: request.ExecutionID,
AuthorizationID: request.CommandID,
ClaimGeneration: request.ClaimGeneration,
ClaimToken: ctx.GetHeader(claimTokenHeader),
CommandSHA256: request.CommandSHA256,
SubmissionID: ctx.Param("submission_id"),
PlatformOrderNo: request.PlatformOrderNo,
PlatformOrderedAt: request.PlatformOrderedAt,
PlatformOrderStatus: request.PlatformOrderStatus,
ObservedTitle: request.ObservedTitle,
SelectedSKU: request.SelectedSKU,
Quantity: request.Quantity,
TotalPriceCents: request.TotalPriceCents,
EvidenceAssetID: request.EvidenceAssetID,
EvidenceSHA256: request.EvidenceSHA256,
IdempotencyKey: ctx.GetHeader("Idempotency-Key"),
},
)
if err != nil {
writeUsecaseError(ctx, err)
return
}
writeOrderSubmission(ctx, result)
}
func (handler *deviceHandlers) manualReviewOrderSubmission(ctx *gin.Context) {
principal, ok := devicePrincipal(ctx)
if !ok {
return
}
var request struct {
DeviceID string `json:"device_id"`
ExecutionID string `json:"execution_id"`
ClaimGeneration int64 `json:"claim_generation"`
CommandID string `json:"command_id"`
CommandSHA256 string `json:"command_sha256"`
ReasonCode string `json:"reason_code"`
EvidenceAssetID string `json:"evidence_asset_id"`
EvidenceSHA256 string `json:"evidence_sha256"`
}
if !decodeDeviceJSON(ctx, &request) ||
!deviceIDMatches(ctx, request.DeviceID, principal.DeviceID) {
return
}
result, err := handler.services.Submissions.ManualReview(
ctx.Request.Context(),
usecase.ManualReviewOrderSubmissionCommand{
UserID: principal.UserID,
DeviceID: principal.DeviceID,
TaskID: ctx.Param("id"),
ExecutionID: request.ExecutionID,
AuthorizationID: request.CommandID,
ClaimGeneration: request.ClaimGeneration,
ClaimToken: ctx.GetHeader(claimTokenHeader),
CommandSHA256: request.CommandSHA256,
SubmissionID: ctx.Param("submission_id"),
ReasonCode: request.ReasonCode,
EvidenceAssetID: request.EvidenceAssetID,
EvidenceSHA256: request.EvidenceSHA256,
IdempotencyKey: ctx.GetHeader("Idempotency-Key"),
},
)
if err != nil {
writeUsecaseError(ctx, err)
return
}
writeOrderSubmission(ctx, result)
}
func writeOrderSubmission(
ctx *gin.Context,
result usecase.OrderSubmissionResult,
) {
ctx.Header("Cache-Control", "no-store")
ctx.JSON(http.StatusOK, gin.H{
"submission": orderSubmissionResponse(result.Submission),
"replayed": result.Replayed,
})
}
func orderSubmissionResponse(submission domain.OrderSubmission) gin.H {
return gin.H{
"id": submission.ID,
"command_id": submission.AuthorizationID,
"dry_run_id": submission.DryRunID,
"task_id": submission.TaskID,
"execution_id": submission.ExecutionID,
"command_sha256": submission.CommandSHA256,
"dry_run_evidence_sha256": submission.DryRunEvidenceSHA256,
"status": submission.Status,
"expected_title": submission.ExpectedTitle,
"expected_sku": submission.ExpectedSKU,
"expected_quantity": submission.ExpectedQuantity,
"expected_unit_price_cents": submission.ExpectedUnitPriceCents,
"expected_total_price_cents": submission.ExpectedTotalPriceCents,
"platform_order_no": submission.PlatformOrderNo,
"platform_ordered_at": formatOptionalTime(
submission.PlatformOrderedAt,
),
"platform_order_status": submission.PlatformOrderStatus,
"reconciliation_evidence_asset_id": submission.ReconciliationEvidenceAssetID,
"reconciliation_evidence_sha256": submission.ReconciliationEvidenceSHA256,
"manual_reason_code": submission.ManualReasonCode,
"fenced_at": formatTime(submission.FencedAt),
"reconciled_at": formatOptionalTime(submission.ReconciledAt),
"manual_review_at": formatOptionalTime(submission.ManualReviewAt),
}
}
func (handler *deviceHandlers) startOrderDryRun(ctx *gin.Context) {
principal, ok := devicePrincipal(ctx)
if !ok {
@@ -713,6 +713,9 @@ func TestDeviceExecutionResultsAreIdempotentAndAuditable(t *testing.T) {
if err != nil {
t.Fatalf("migration.New() after review error = %v", err)
}
if err := runner.Down(context.Background()); err != nil {
t.Fatalf("order submission migration down: %v", err)
}
if err := runner.Down(context.Background()); err != nil {
t.Fatalf("order dry-run migration down: %v", err)
}
@@ -724,8 +727,8 @@ func TestDeviceExecutionResultsAreIdempotentAndAuditable(t *testing.T) {
}
if applied, err := runner.Up(context.Background()); err != nil {
t.Fatalf("restore device command migration: %v", err)
} else if applied != 2 {
t.Fatalf("restored migrations = %d, want 2", applied)
} else if applied != 3 {
t.Fatalf("restored migrations = %d, want 3", applied)
}
completePayload := fmt.Sprintf(
@@ -1125,6 +1128,12 @@ func TestDeviceOrderCommandDeliveryAndAcknowledgementAreRecoverable(
if !strings.Contains(readyDryRun.Body.String(), `"status":"READY"`) {
t.Fatalf("dry-run ready = %s", readyDryRun.Body.String())
}
var readyDryRunResponse struct {
DryRun struct {
ID string `json:"id"`
} `json:"dry_run"`
}
decodeResponse(t, readyDryRun, &readyDryRunResponse)
readyDryRunReplay := performDeviceRequest(t, fixture.router, deviceRequest{
method: http.MethodPost,
target: "/api/v1/tasks/" + taskID +
@@ -1139,12 +1148,182 @@ func TestDeviceOrderCommandDeliveryAndAcknowledgementAreRecoverable(
if !strings.Contains(readyDryRunReplay.Body.String(), `"replayed":true`) {
t.Fatalf("dry-run ready replay = %s", readyDryRunReplay.Body.String())
}
var deliveredEvents, acknowledgedEvents, dryRunStartedEvents, dryRunReadyEvents int
startSubmissionPayload := fmt.Sprintf(
`{"device_id":%q,"execution_id":%q,"claim_generation":%d,"command_id":%q,"command_sha256":%q,"dry_run_id":%q,"dry_run_evidence_sha256":%q,"observed_title":%q,"selected_sku":%q,"quantity":2,"unit_price_cents":2150,"total_price_cents":4300}`,
deviceTestDeviceID,
started.Execution.ID,
started.Task.ClaimGeneration,
command.ID,
command.CommandSHA256,
readyDryRunResponse.DryRun.ID,
dryRunEvidenceSHA,
command.Candidate.Title,
command.OriginalSKU,
)
startSubmission := performDeviceRequest(t, fixture.router, deviceRequest{
method: http.MethodPost,
target: "/api/v1/tasks/" + taskID + "/order-submissions/start",
contentType: "application/json",
body: strings.NewReader(startSubmissionPayload),
bearerToken: testOpaqueToken,
claimToken: testOpaqueToken,
idempotencyKey: "order-submission-start",
})
requireDeviceStatus(t, startSubmission, http.StatusOK)
var submissionResponse struct {
Submission struct {
ID string `json:"id"`
} `json:"submission"`
}
decodeResponse(t, startSubmission, &submissionResponse)
if submissionResponse.Submission.ID == "" ||
!strings.Contains(startSubmission.Body.String(), `"status":"FENCED"`) {
t.Fatalf("submission start = %s", startSubmission.Body.String())
}
startSubmissionReplay := performDeviceRequest(t, fixture.router, deviceRequest{
method: http.MethodPost,
target: "/api/v1/tasks/" + taskID + "/order-submissions/start",
contentType: "application/json",
body: strings.NewReader(startSubmissionPayload),
bearerToken: testOpaqueToken,
claimToken: testOpaqueToken,
idempotencyKey: "order-submission-start",
})
requireDeviceStatus(t, startSubmissionReplay, http.StatusOK)
if !strings.Contains(startSubmissionReplay.Body.String(), `"replayed":true`) {
t.Fatalf("submission start replay = %s", startSubmissionReplay.Body.String())
}
secondSubmission := performDeviceRequest(t, fixture.router, deviceRequest{
method: http.MethodPost,
target: "/api/v1/tasks/" + taskID + "/order-submissions/start",
contentType: "application/json",
body: strings.NewReader(startSubmissionPayload),
bearerToken: testOpaqueToken,
claimToken: testOpaqueToken,
idempotencyKey: "order-submission-start-second",
})
requireDeviceStatus(t, secondSubmission, http.StatusConflict)
manualReviewPayload := fmt.Sprintf(
`{"device_id":%q,"execution_id":%q,"claim_generation":%d,"command_id":%q,"command_sha256":%q,"reason_code":"ORDER_PAGE_UNKNOWN"}`,
deviceTestDeviceID,
started.Execution.ID,
started.Task.ClaimGeneration,
command.ID,
command.CommandSHA256,
)
manualReview := performDeviceRequest(t, fixture.router, deviceRequest{
method: http.MethodPost,
target: "/api/v1/tasks/" + taskID + "/order-submissions/" +
submissionResponse.Submission.ID + "/manual-review",
contentType: "application/json",
body: strings.NewReader(manualReviewPayload),
bearerToken: testOpaqueToken,
claimToken: testOpaqueToken,
idempotencyKey: "order-submission-manual",
})
requireDeviceStatus(t, manualReview, http.StatusOK)
if !strings.Contains(manualReview.Body.String(), `"status":"MANUAL_REVIEW"`) {
t.Fatalf("submission manual review = %s", manualReview.Body.String())
}
const reconciliationEvidenceID = "00000000-0000-4000-8000-000000000078"
reconciliationEvidenceSHA := strings.Repeat("d", 64)
if _, err := fixture.db.Exec(
`INSERT INTO execution_evidence_assets (
id, task_id, execution_id, media_type, size_bytes, sha256,
storage_key, created_at, received_after_execution_expiry
) VALUES (?, ?, ?, 'image/jpeg', 10, ?, ?, ?, 0)`,
reconciliationEvidenceID,
taskID,
started.Execution.ID,
reconciliationEvidenceSHA,
"orders/reconciliation.jpg",
time.Now().UTC().Format(time.RFC3339Nano),
); err != nil {
t.Fatalf("seed reconciliation evidence: %v", err)
}
reconcilePayload := fmt.Sprintf(
`{"device_id":%q,"execution_id":%q,"claim_generation":%d,"command_id":%q,"command_sha256":%q,"platform_order_no":"12345678901234567890","platform_ordered_at":%q,"platform_order_status":"PENDING_PAYMENT","observed_title":%q,"selected_sku":%q,"quantity":2,"total_price_cents":4300,"evidence_asset_id":%q,"evidence_sha256":%q}`,
deviceTestDeviceID,
started.Execution.ID,
started.Task.ClaimGeneration,
command.ID,
command.CommandSHA256,
time.Now().UTC().Format(time.RFC3339),
command.Candidate.Title,
command.OriginalSKU,
reconciliationEvidenceID,
reconciliationEvidenceSHA,
)
reconciled := performDeviceRequest(t, fixture.router, deviceRequest{
method: http.MethodPost,
target: "/api/v1/tasks/" + taskID + "/order-submissions/" +
submissionResponse.Submission.ID + "/reconcile",
contentType: "application/json",
body: strings.NewReader(reconcilePayload),
bearerToken: testOpaqueToken,
claimToken: testOpaqueToken,
idempotencyKey: "order-submission-reconcile",
})
requireDeviceStatus(t, reconciled, http.StatusOK)
if !strings.Contains(reconciled.Body.String(), `"status":"RECONCILED"`) ||
!strings.Contains(reconciled.Body.String(), `"platform_order_status":"PENDING_PAYMENT"`) {
t.Fatalf("submission reconciliation = %s", reconciled.Body.String())
}
reconciledReplay := performDeviceRequest(t, fixture.router, deviceRequest{
method: http.MethodPost,
target: "/api/v1/tasks/" + taskID + "/order-submissions/" +
submissionResponse.Submission.ID + "/reconcile",
contentType: "application/json",
body: strings.NewReader(reconcilePayload),
bearerToken: testOpaqueToken,
claimToken: testOpaqueToken,
idempotencyKey: "order-submission-reconcile",
})
requireDeviceStatus(t, reconciledReplay, http.StatusOK)
if !strings.Contains(reconciledReplay.Body.String(), `"replayed":true`) {
t.Fatalf("submission reconcile replay = %s", reconciledReplay.Body.String())
}
var outcomeSubmitted bool
if err := fixture.db.QueryRow(
`SELECT order_submitted FROM execution_outcomes
WHERE execution_id = ?`,
started.Execution.ID,
).Scan(&outcomeSubmitted); err != nil {
t.Fatalf("query reconciled outcome: %v", err)
}
if !outcomeSubmitted {
t.Fatal("reconciled execution outcome did not record order_submitted")
}
var taskStatus, authorizationStatus string
if err := fixture.db.QueryRow(
`SELECT status FROM purchase_tasks WHERE id = ?`,
taskID,
).Scan(&taskStatus); err != nil {
t.Fatalf("query reconciled task: %v", err)
}
if err := fixture.db.QueryRow(
`SELECT status FROM order_authorizations WHERE id = ?`,
command.ID,
).Scan(&authorizationStatus); err != nil {
t.Fatalf("query consumed authorization: %v", err)
}
if taskStatus != "SUCCEEDED" || authorizationStatus != "CONSUMED" {
t.Fatalf(
"reconciled task/authorization = %s/%s",
taskStatus,
authorizationStatus,
)
}
var deliveredEvents, acknowledgedEvents, dryRunStartedEvents,
dryRunReadyEvents, fencedEvents, manualEvents, reconciledEvents int
for eventType, target := range map[string]*int{
"ORDER_AUTHORIZATION_DELIVERED": &deliveredEvents,
"ORDER_AUTHORIZATION_ACKNOWLEDGED": &acknowledgedEvents,
"ORDER_DRY_RUN_STARTED": &dryRunStartedEvents,
"ORDER_DRY_RUN_READY": &dryRunReadyEvents,
"ORDER_SUBMISSION_FENCED": &fencedEvents,
"ORDER_SUBMISSION_MANUAL_REVIEW": &manualEvents,
"ORDER_SUBMISSION_RECONCILED": &reconciledEvents,
} {
if err := fixture.db.QueryRow(
`SELECT COUNT(*) FROM task_events
@@ -1156,13 +1335,17 @@ func TestDeviceOrderCommandDeliveryAndAcknowledgementAreRecoverable(
}
}
if deliveredEvents != 1 || acknowledgedEvents != 1 ||
dryRunStartedEvents != 1 || dryRunReadyEvents != 1 {
dryRunStartedEvents != 1 || dryRunReadyEvents != 1 ||
fencedEvents != 1 || manualEvents != 1 || reconciledEvents != 1 {
t.Fatalf(
"delivery/ack/dry-run events = %d/%d/%d/%d",
"order workflow events = %d/%d/%d/%d/%d/%d/%d",
deliveredEvents,
acknowledgedEvents,
dryRunStartedEvents,
dryRunReadyEvents,
fencedEvents,
manualEvents,
reconciledEvents,
)
}
runner, err := migration.New(fixture.db)
@@ -1170,7 +1353,7 @@ func TestDeviceOrderCommandDeliveryAndAcknowledgementAreRecoverable(
t.Fatalf("migration.New() error = %v", err)
}
if err := runner.Down(context.Background()); err == nil {
t.Fatal("order dry-run migration down succeeded with dry-run data")
t.Fatal("order submission migration down succeeded with retained data")
}
}
@@ -1503,13 +1686,18 @@ func newDeviceHTTPFixture(t *testing.T) *deviceHTTPFixture {
if err != nil {
t.Fatalf("usecase.NewOrderDryRunService() error = %v", err)
}
submissions, err := usecase.NewOrderSubmissionService(store, clock, ids)
if err != nil {
t.Fatalf("usecase.NewOrderSubmissionService() error = %v", err)
}
deviceRoutes, err := NewDeviceRouteRegistrar(
DeviceServices{
Lifecycle: lifecycle,
Assets: assets,
Results: results,
Commands: commands,
DryRuns: dryRuns,
Lifecycle: lifecycle,
Assets: assets,
Results: results,
Commands: commands,
DryRuns: dryRuns,
Submissions: submissions,
},
)
if err != nil {
@@ -0,0 +1,526 @@
package usecase
import (
"context"
"errors"
"math"
"regexp"
"strings"
"time"
"cmroubao/backend-api/internal/domain"
)
const (
orderSubmissionStartOperation = "START"
orderSubmissionReconcileOperation = "RECONCILE"
orderSubmissionManualReviewOperation = "MANUAL_REVIEW"
)
type StartOrderSubmissionCommand struct {
UserID string
DeviceID string
TaskID string
ExecutionID string
AuthorizationID string
ClaimGeneration int64
ClaimToken string
CommandSHA256 string
DryRunID string
DryRunEvidenceSHA256 string
ObservedTitle string
SelectedSKU string
Quantity int
UnitPriceCents int64
TotalPriceCents int64
IdempotencyKey string
}
type StartOrderSubmissionWrite struct {
StartOrderSubmissionCommand
SubmissionID string
ClaimTokenHash string
RequestSHA256 string
Now time.Time
Event domain.TaskEvent
}
type ReconcileOrderSubmissionCommand struct {
UserID string
DeviceID string
TaskID string
ExecutionID string
AuthorizationID string
ClaimGeneration int64
ClaimToken string
CommandSHA256 string
SubmissionID string
PlatformOrderNo string
PlatformOrderedAt string
PlatformOrderStatus string
ObservedTitle string
SelectedSKU string
Quantity int
TotalPriceCents int64
EvidenceAssetID string
EvidenceSHA256 string
IdempotencyKey string
}
type ReconcileOrderSubmissionWrite struct {
ReconcileOrderSubmissionCommand
ParsedPlatformOrderedAt time.Time
ClaimTokenHash string
RequestSHA256 string
Now time.Time
Event domain.TaskEvent
}
type ManualReviewOrderSubmissionCommand struct {
UserID string
DeviceID string
TaskID string
ExecutionID string
AuthorizationID string
ClaimGeneration int64
ClaimToken string
CommandSHA256 string
SubmissionID string
ReasonCode string
EvidenceAssetID string
EvidenceSHA256 string
IdempotencyKey string
}
type ManualReviewOrderSubmissionWrite struct {
ManualReviewOrderSubmissionCommand
ClaimTokenHash string
RequestSHA256 string
Now time.Time
Event domain.TaskEvent
}
type OrderSubmissionResult struct {
Submission domain.OrderSubmission
Replayed bool
}
type OrderSubmissionRepository interface {
StartOrderSubmission(
context.Context,
StartOrderSubmissionWrite,
) (domain.OrderSubmission, bool, error)
ReconcileOrderSubmission(
context.Context,
ReconcileOrderSubmissionWrite,
) (domain.OrderSubmission, bool, error)
ManualReviewOrderSubmission(
context.Context,
ManualReviewOrderSubmissionWrite,
) (domain.OrderSubmission, bool, error)
}
type OrderSubmissionService struct {
repository OrderSubmissionRepository
clock Clock
ids IDGenerator
}
func NewOrderSubmissionService(
repository OrderSubmissionRepository,
clock Clock,
ids IDGenerator,
) (*OrderSubmissionService, error) {
if repository == nil || clock == nil || ids == nil {
return nil, errors.New("order submission service dependencies are required")
}
return &OrderSubmissionService{
repository: repository,
clock: clock,
ids: ids,
}, nil
}
func (service *OrderSubmissionService) Start(
ctx context.Context,
command StartOrderSubmissionCommand,
) (OrderSubmissionResult, error) {
command = normalizeStartOrderSubmission(command)
fields := orderSubmissionIdentityFields(
command.UserID,
command.DeviceID,
command.TaskID,
command.ExecutionID,
command.AuthorizationID,
command.ClaimGeneration,
command.ClaimToken,
command.CommandSHA256,
command.IdempotencyKey,
)
if !isUUID(command.DryRunID) {
fields["dry_run_id"] = "must be a UUID"
}
if !sha256Pattern.MatchString(command.DryRunEvidenceSHA256) {
fields["dry_run_evidence_sha256"] = "must be lowercase SHA-256"
}
validateOrderSnapshot(
fields,
command.ObservedTitle,
command.SelectedSKU,
command.Quantity,
command.UnitPriceCents,
command.TotalPriceCents,
)
if command.Quantity > 0 &&
command.UnitPriceCents <= math.MaxInt64/int64(command.Quantity) &&
command.TotalPriceCents !=
command.UnitPriceCents*int64(command.Quantity) {
fields["total_price_cents"] = "must equal unit price times quantity"
}
if len(fields) > 0 {
return OrderSubmissionResult{}, invalidError(
"ORDER_SUBMISSION_START_INVALID",
"order submission start request is invalid",
fields,
)
}
requestHash, err := lifecycleRequestHash(command)
if err != nil {
return OrderSubmissionResult{}, internalLifecycleFailure(err)
}
submissionID, err := service.ids.NewID()
if err != nil {
return OrderSubmissionResult{}, internalLifecycleFailure(err)
}
eventID, err := service.ids.NewID()
if err != nil {
return OrderSubmissionResult{}, internalLifecycleFailure(err)
}
now := service.clock.Now().UTC()
userID, deviceID := command.UserID, command.DeviceID
submission, replayed, err := service.repository.StartOrderSubmission(
ctx,
StartOrderSubmissionWrite{
StartOrderSubmissionCommand: command,
SubmissionID: submissionID,
ClaimTokenHash: hashSecret(command.ClaimToken),
RequestSHA256: requestHash,
Now: now,
Event: domain.TaskEvent{
ID: eventID,
TaskID: command.TaskID,
ActorUserID: &userID,
ActorDeviceID: &deviceID,
Type: "ORDER_SUBMISSION_FENCED",
Message: "single order submission fenced",
OccurredAt: now,
},
},
)
if err != nil {
return OrderSubmissionResult{}, wrapLifecycleRepositoryError(err)
}
return OrderSubmissionResult{
Submission: submission,
Replayed: replayed,
}, nil
}
func (service *OrderSubmissionService) Reconcile(
ctx context.Context,
command ReconcileOrderSubmissionCommand,
) (OrderSubmissionResult, error) {
command = normalizeReconcileOrderSubmission(command)
fields := orderSubmissionIdentityFields(
command.UserID,
command.DeviceID,
command.TaskID,
command.ExecutionID,
command.AuthorizationID,
command.ClaimGeneration,
command.ClaimToken,
command.CommandSHA256,
command.IdempotencyKey,
)
if !isUUID(command.SubmissionID) {
fields["submission_id"] = "must be a UUID"
}
if !platformOrderNumberPattern.MatchString(command.PlatformOrderNo) {
fields["platform_order_no"] = "must be 8 to 40 digits"
}
orderedAt, err := time.Parse(time.RFC3339, command.PlatformOrderedAt)
if err != nil {
fields["platform_ordered_at"] = "must be RFC3339"
}
if command.PlatformOrderStatus != "PENDING_PAYMENT" {
fields["platform_order_status"] = "must be PENDING_PAYMENT"
}
validateReconciledOrderSnapshot(
fields,
command.ObservedTitle,
command.SelectedSKU,
command.Quantity,
command.TotalPriceCents,
)
if !isUUID(command.EvidenceAssetID) {
fields["evidence_asset_id"] = "must be a UUID"
}
if !sha256Pattern.MatchString(command.EvidenceSHA256) {
fields["evidence_sha256"] = "must be lowercase SHA-256"
}
if len(fields) > 0 {
return OrderSubmissionResult{}, invalidError(
"ORDER_SUBMISSION_RECONCILE_INVALID",
"order submission reconciliation request is invalid",
fields,
)
}
requestHash, err := lifecycleRequestHash(command)
if err != nil {
return OrderSubmissionResult{}, internalLifecycleFailure(err)
}
eventID, err := service.ids.NewID()
if err != nil {
return OrderSubmissionResult{}, internalLifecycleFailure(err)
}
now := service.clock.Now().UTC()
userID, deviceID := command.UserID, command.DeviceID
submission, replayed, err := service.repository.ReconcileOrderSubmission(
ctx,
ReconcileOrderSubmissionWrite{
ReconcileOrderSubmissionCommand: command,
ParsedPlatformOrderedAt: orderedAt.UTC(),
ClaimTokenHash: hashSecret(command.ClaimToken),
RequestSHA256: requestHash,
Now: now,
Event: domain.TaskEvent{
ID: eventID,
TaskID: command.TaskID,
ActorUserID: &userID,
ActorDeviceID: &deviceID,
Type: "ORDER_SUBMISSION_RECONCILED",
Message: "pending-payment order uniquely reconciled",
OccurredAt: now,
},
},
)
if err != nil {
return OrderSubmissionResult{}, wrapLifecycleRepositoryError(err)
}
return OrderSubmissionResult{
Submission: submission,
Replayed: replayed,
}, nil
}
func (service *OrderSubmissionService) ManualReview(
ctx context.Context,
command ManualReviewOrderSubmissionCommand,
) (OrderSubmissionResult, error) {
command = normalizeManualReviewOrderSubmission(command)
fields := orderSubmissionIdentityFields(
command.UserID,
command.DeviceID,
command.TaskID,
command.ExecutionID,
command.AuthorizationID,
command.ClaimGeneration,
command.ClaimToken,
command.CommandSHA256,
command.IdempotencyKey,
)
if !isUUID(command.SubmissionID) {
fields["submission_id"] = "must be a UUID"
}
if _, ok := orderSubmissionManualReasons[command.ReasonCode]; !ok {
fields["reason_code"] = "must be an allowed reason"
}
if (command.EvidenceAssetID == "") != (command.EvidenceSHA256 == "") {
fields["evidence"] = "asset id and SHA-256 must be provided together"
} else if command.EvidenceAssetID != "" {
if !isUUID(command.EvidenceAssetID) {
fields["evidence_asset_id"] = "must be a UUID"
}
if !sha256Pattern.MatchString(command.EvidenceSHA256) {
fields["evidence_sha256"] = "must be lowercase SHA-256"
}
}
if len(fields) > 0 {
return OrderSubmissionResult{}, invalidError(
"ORDER_SUBMISSION_MANUAL_REVIEW_INVALID",
"order submission manual-review request is invalid",
fields,
)
}
requestHash, err := lifecycleRequestHash(command)
if err != nil {
return OrderSubmissionResult{}, internalLifecycleFailure(err)
}
eventID, err := service.ids.NewID()
if err != nil {
return OrderSubmissionResult{}, internalLifecycleFailure(err)
}
now := service.clock.Now().UTC()
userID, deviceID := command.UserID, command.DeviceID
submission, replayed, err := service.repository.ManualReviewOrderSubmission(
ctx,
ManualReviewOrderSubmissionWrite{
ManualReviewOrderSubmissionCommand: command,
ClaimTokenHash: hashSecret(command.ClaimToken),
RequestSHA256: requestHash,
Now: now,
Event: domain.TaskEvent{
ID: eventID,
TaskID: command.TaskID,
ActorUserID: &userID,
ActorDeviceID: &deviceID,
Type: "ORDER_SUBMISSION_MANUAL_REVIEW",
Message: "order submission requires manual reconciliation",
OccurredAt: now,
},
},
)
if err != nil {
return OrderSubmissionResult{}, wrapLifecycleRepositoryError(err)
}
return OrderSubmissionResult{
Submission: submission,
Replayed: replayed,
}, nil
}
func orderSubmissionIdentityFields(
userID, deviceID, taskID, executionID, authorizationID string,
claimGeneration int64,
claimToken, commandSHA256, idempotencyKey string,
) map[string]string {
return dryRunIdentityFields(
userID,
deviceID,
taskID,
executionID,
authorizationID,
claimGeneration,
claimToken,
commandSHA256,
idempotencyKey,
)
}
func validateOrderSnapshot(
fields map[string]string,
title, sku string,
quantity int,
unitPriceCents, totalPriceCents int64,
) {
if title == "" || len([]byte(title)) > 1024 {
fields["observed_title"] = "must be 1 to 1024 UTF-8 bytes"
}
if sku == "" || len([]byte(sku)) > 512 {
fields["selected_sku"] = "must be 1 to 512 UTF-8 bytes"
}
if quantity < 1 || quantity > 99 {
fields["quantity"] = "must be 1 to 99"
}
if unitPriceCents < 1 {
fields["unit_price_cents"] = "must be positive"
}
if totalPriceCents < 1 {
fields["total_price_cents"] = "must be positive"
}
if quantity > 0 &&
unitPriceCents > math.MaxInt64/int64(quantity) {
fields["total_price_cents"] = "price multiplication overflows"
}
}
func validateReconciledOrderSnapshot(
fields map[string]string,
title, sku string,
quantity int,
totalPriceCents int64,
) {
if title == "" || len([]byte(title)) > 1024 {
fields["observed_title"] = "must be 1 to 1024 UTF-8 bytes"
}
if sku == "" || len([]byte(sku)) > 512 {
fields["selected_sku"] = "must be 1 to 512 UTF-8 bytes"
}
if quantity < 1 || quantity > 99 {
fields["quantity"] = "must be 1 to 99"
}
if totalPriceCents < 1 {
fields["total_price_cents"] = "must be positive"
}
}
func normalizeStartOrderSubmission(
command StartOrderSubmissionCommand,
) StartOrderSubmissionCommand {
command.UserID = strings.TrimSpace(command.UserID)
command.DeviceID = strings.TrimSpace(command.DeviceID)
command.TaskID = strings.TrimSpace(command.TaskID)
command.ExecutionID = strings.TrimSpace(command.ExecutionID)
command.AuthorizationID = strings.TrimSpace(command.AuthorizationID)
command.ClaimToken = strings.TrimSpace(command.ClaimToken)
command.CommandSHA256 = strings.TrimSpace(command.CommandSHA256)
command.DryRunID = strings.TrimSpace(command.DryRunID)
command.DryRunEvidenceSHA256 = strings.TrimSpace(
command.DryRunEvidenceSHA256,
)
command.ObservedTitle = strings.TrimSpace(command.ObservedTitle)
command.SelectedSKU = strings.TrimSpace(command.SelectedSKU)
command.IdempotencyKey = strings.TrimSpace(command.IdempotencyKey)
return command
}
func normalizeReconcileOrderSubmission(
command ReconcileOrderSubmissionCommand,
) ReconcileOrderSubmissionCommand {
command.UserID = strings.TrimSpace(command.UserID)
command.DeviceID = strings.TrimSpace(command.DeviceID)
command.TaskID = strings.TrimSpace(command.TaskID)
command.ExecutionID = strings.TrimSpace(command.ExecutionID)
command.AuthorizationID = strings.TrimSpace(command.AuthorizationID)
command.ClaimToken = strings.TrimSpace(command.ClaimToken)
command.CommandSHA256 = strings.TrimSpace(command.CommandSHA256)
command.SubmissionID = strings.TrimSpace(command.SubmissionID)
command.PlatformOrderNo = strings.TrimSpace(command.PlatformOrderNo)
command.PlatformOrderedAt = strings.TrimSpace(command.PlatformOrderedAt)
command.PlatformOrderStatus = strings.TrimSpace(command.PlatformOrderStatus)
command.ObservedTitle = strings.TrimSpace(command.ObservedTitle)
command.SelectedSKU = strings.TrimSpace(command.SelectedSKU)
command.EvidenceAssetID = strings.TrimSpace(command.EvidenceAssetID)
command.EvidenceSHA256 = strings.TrimSpace(command.EvidenceSHA256)
command.IdempotencyKey = strings.TrimSpace(command.IdempotencyKey)
return command
}
func normalizeManualReviewOrderSubmission(
command ManualReviewOrderSubmissionCommand,
) ManualReviewOrderSubmissionCommand {
command.UserID = strings.TrimSpace(command.UserID)
command.DeviceID = strings.TrimSpace(command.DeviceID)
command.TaskID = strings.TrimSpace(command.TaskID)
command.ExecutionID = strings.TrimSpace(command.ExecutionID)
command.AuthorizationID = strings.TrimSpace(command.AuthorizationID)
command.ClaimToken = strings.TrimSpace(command.ClaimToken)
command.CommandSHA256 = strings.TrimSpace(command.CommandSHA256)
command.SubmissionID = strings.TrimSpace(command.SubmissionID)
command.ReasonCode = strings.TrimSpace(command.ReasonCode)
command.EvidenceAssetID = strings.TrimSpace(command.EvidenceAssetID)
command.EvidenceSHA256 = strings.TrimSpace(command.EvidenceSHA256)
command.IdempotencyKey = strings.TrimSpace(command.IdempotencyKey)
return command
}
var platformOrderNumberPattern = regexp.MustCompile(`^[0-9]{8,40}$`)
var orderSubmissionManualReasons = map[string]struct{}{
"ORDER_NOT_FOUND": {},
"ORDER_AMBIGUOUS": {},
"ORDER_FIELDS_INCOMPLETE": {},
"ORDER_PAGE_UNKNOWN": {},
"RISK_OR_PAYMENT_BOUNDARY": {},
"EVIDENCE_UNAVAILABLE": {},
}