package sqlite_test import ( "context" "crypto/sha256" "database/sql" "encoding/hex" "errors" "fmt" "sync" "testing" "time" "cmroubao/backend-api/internal/domain" platformdatabase "cmroubao/backend-api/internal/platform/database" repository "cmroubao/backend-api/internal/repository/sqlite" "cmroubao/backend-api/internal/usecase" ) const lifecycleCreatorSubject = "local-admin" type lifecycleFixture struct { store *repository.Store db *sql.DB adminUserID string buyerOneID string buyerTwoID string deviceOneID string deviceTwoID string } func TestLifecycleRepositoryRecordsReadinessAndRequiresFreshReadyDevice( t *testing.T, ) { fixture := newLifecycleFixture(t) ctx := context.Background() now := lifecycleTime(0) task := createLifecycleTask(t, fixture, 1, now.Add(-time.Minute)) request := lifecycleClaimRequest( fixture, fixture.buyerOneID, fixture.deviceOneID, 1, now, now.Add(10*time.Minute), now.Add(-2*time.Minute), ) _, err := fixture.store.ClaimNext(ctx, request) assertLifecycleError(t, err, usecase.ErrDeviceNotReady) notReady, err := fixture.store.RecordDeviceHeartbeat( ctx, lifecycleDeviceHeartbeat( fixture.buyerOneID, fixture.deviceOneID, now, true, false, ), ) if err != nil { t.Fatalf("RecordDeviceHeartbeat(not ready) error = %v", err) } if !notReady.Device.AccessibilityEnabled || notReady.Device.PDDInstalled || notReady.Device.ReadinessAt == nil || !notReady.Device.ReadinessAt.Equal(now) { t.Fatalf("not-ready heartbeat = %+v", notReady) } _, err = fixture.store.ClaimNext(ctx, request) assertLifecycleError(t, err, usecase.ErrDeviceNotReady) readyAt := now.Add(3 * time.Minute) ready, err := fixture.store.RecordDeviceHeartbeat( ctx, lifecycleDeviceHeartbeat( fixture.buyerOneID, fixture.deviceOneID, readyAt, true, true, ), ) if err != nil { t.Fatalf("RecordDeviceHeartbeat(ready) error = %v", err) } if ready.Device.AppVersion == nil || *ready.Device.AppVersion != "0.1.0" || ready.Device.AndroidVersion == nil || *ready.Device.AndroidVersion != "16" || ready.Device.PDDVersion == nil || *ready.Device.PDDVersion != "8.17.0" || ready.Device.LastSeenAt == nil || !ready.Device.LastSeenAt.Equal(readyAt) { t.Fatalf("ready heartbeat = %+v", ready) } staleRequest := request staleRequest.Now = readyAt.Add(3 * time.Minute) staleRequest.ExpiresAt = staleRequest.Now.Add(10 * time.Minute) staleRequest.ReadinessAfter = staleRequest.Now.Add(-2 * time.Minute) _, err = fixture.store.ClaimNext(ctx, staleRequest) assertLifecycleError(t, err, usecase.ErrDeviceNotReady) freshAt := staleRequest.Now if _, err := fixture.store.RecordDeviceHeartbeat( ctx, lifecycleDeviceHeartbeat( fixture.buyerOneID, fixture.deviceOneID, freshAt, true, true, ), ); err != nil { t.Fatalf("refresh readiness: %v", err) } staleRequest.IdempotencyKey = "claim-ready" staleRequest.RequestHash = lifecycleHash("claim-ready") staleRequest.ClaimTokenHash = lifecycleHash("claim-token-ready") staleRequest.Event = lifecycleEvent( 10, "", "TASK_CLAIMED", freshAt, fixture.buyerOneID, fixture.deviceOneID, ) claimed, err := fixture.store.ClaimNext(ctx, staleRequest) if err != nil || claimed.Task == nil || claimed.Task.ID != task.ID { t.Fatalf("ClaimNext(fresh) = %+v, error = %v", claimed, err) } active, err := fixture.store.RecordDeviceHeartbeat( ctx, lifecycleDeviceHeartbeat( fixture.buyerOneID, fixture.deviceOneID, freshAt.Add(time.Second), true, true, ), ) if err != nil { t.Fatalf("RecordDeviceHeartbeat(active) error = %v", err) } if active.ActiveTaskID == nil || *active.ActiveTaskID != task.ID { t.Fatalf("active task = %+v", active.ActiveTaskID) } } func TestLifecycleRepositoryClaimIsFIFOIdempotentAndDeviceSingleActive( t *testing.T, ) { fixture := newLifecycleFixture(t) ctx := context.Background() now := lifecycleTime(0) older := createLifecycleTask(t, fixture, 1, now.Add(-2*time.Minute)) _ = createLifecycleTask(t, fixture, 2, now.Add(-time.Minute)) recordReadyHeartbeat( t, fixture, fixture.buyerOneID, fixture.deviceOneID, now, ) request := lifecycleClaimRequest( fixture, fixture.buyerOneID, fixture.deviceOneID, 20, now, now.Add(10*time.Minute), now.Add(-2*time.Minute), ) first, err := fixture.store.ClaimNext(ctx, request) if err != nil || first.Task == nil || first.Replayed { t.Fatalf("first ClaimNext() = %+v, error = %v", first, err) } if first.Task.ID != older.ID || first.Task.Status != domain.TaskStatusClaimed || first.Task.Version != older.Version+1 || first.Task.ClaimGeneration != 1 || first.Task.ClaimTokenHash == nil || *first.Task.ClaimTokenHash != request.ClaimTokenHash || first.Task.ClaimedByDeviceID == nil || *first.Task.ClaimedByDeviceID != fixture.deviceOneID { t.Fatalf("claimed task = %+v", first.Task) } replayRequest := request replayRequest.Event = lifecycleEvent( 21, "", "TASK_CLAIMED", now.Add(time.Second), fixture.buyerOneID, fixture.deviceOneID, ) replayed, err := fixture.store.ClaimNext(ctx, replayRequest) if err != nil || replayed.Task == nil || !replayed.Replayed || replayed.Task.ID != first.Task.ID || replayed.Task.ClaimGeneration != first.Task.ClaimGeneration || replayed.Task.ClaimTokenHash == nil || *replayed.Task.ClaimTokenHash != request.ClaimTokenHash { t.Fatalf("replayed ClaimNext() = %+v, error = %v", replayed, err) } assertTaskEventCount(t, fixture.db, older.ID, "TASK_CLAIMED", 1) conflict := request conflict.RequestHash = lifecycleHash("different-claim-request") _, err = fixture.store.ClaimNext(ctx, conflict) assertLifecycleError(t, err, usecase.ErrIdempotencyConflict) secondKey := request secondKey.IdempotencyKey = "claim-second-click" secondKey.RequestHash = lifecycleHash("claim-second-click") secondKey.ClaimTokenHash = lifecycleHash("claim-token-second-click") secondKey.Event = lifecycleEvent( 22, "", "TASK_CLAIMED", now.Add(2*time.Second), fixture.buyerOneID, fixture.deviceOneID, ) _, err = fixture.store.ClaimNext(ctx, secondKey) assertLifecycleError(t, err, usecase.ErrDeviceActiveTask) } func TestLifecycleRepositoryNoTaskReplayRemainsEmpty(t *testing.T) { fixture := newLifecycleFixture(t) ctx := context.Background() now := lifecycleTime(0) recordReadyHeartbeat( t, fixture, fixture.buyerOneID, fixture.deviceOneID, now, ) request := lifecycleClaimRequest( fixture, fixture.buyerOneID, fixture.deviceOneID, 30, now, now.Add(10*time.Minute), now.Add(-2*time.Minute), ) empty, err := fixture.store.ClaimNext(ctx, request) if err != nil || empty.Task != nil || empty.Replayed { t.Fatalf("empty ClaimNext() = %+v, error = %v", empty, err) } _ = createLifecycleTask(t, fixture, 1, now.Add(time.Second)) replayed, err := fixture.store.ClaimNext(ctx, request) if err != nil || replayed.Task != nil || !replayed.Replayed { t.Fatalf("empty replay = %+v, error = %v", replayed, err) } newRequest := request newRequest.IdempotencyKey = "claim-after-empty" newRequest.RequestHash = lifecycleHash("claim-after-empty") newRequest.ClaimTokenHash = lifecycleHash("claim-token-after-empty") newRequest.Event = lifecycleEvent( 31, "", "TASK_CLAIMED", now.Add(2*time.Second), fixture.buyerOneID, fixture.deviceOneID, ) claimed, err := fixture.store.ClaimNext(ctx, newRequest) if err != nil || claimed.Task == nil { t.Fatalf("new claim after empty = %+v, error = %v", claimed, err) } } func TestLifecycleRepositoryReclaimsExpiredClaimAndRejectsOldReplay( t *testing.T, ) { fixture := newLifecycleFixture(t) ctx := context.Background() now := lifecycleTime(0) task := createLifecycleTask(t, fixture, 1, now.Add(-time.Minute)) recordReadyHeartbeat( t, fixture, fixture.buyerOneID, fixture.deviceOneID, now, ) recordReadyHeartbeat( t, fixture, fixture.buyerTwoID, fixture.deviceTwoID, now, ) firstRequest := lifecycleClaimRequest( fixture, fixture.buyerOneID, fixture.deviceOneID, 40, now, now.Add(time.Minute), now.Add(-2*time.Minute), ) first, err := fixture.store.ClaimNext(ctx, firstRequest) if err != nil || first.Task == nil { t.Fatalf("first ClaimNext() = %+v, error = %v", first, err) } reclaimAt := now.Add(2 * time.Minute) secondRequest := lifecycleClaimRequest( fixture, fixture.buyerTwoID, fixture.deviceTwoID, 41, reclaimAt, reclaimAt.Add(10*time.Minute), now.Add(-time.Second), ) reclaimed, err := fixture.store.ClaimNext(ctx, secondRequest) if err != nil || reclaimed.Task == nil { t.Fatalf("reclaim ClaimNext() = %+v, error = %v", reclaimed, err) } if reclaimed.Task.ID != task.ID || reclaimed.Task.ClaimGeneration != 2 || reclaimed.Task.ClaimedByUserID == nil || *reclaimed.Task.ClaimedByUserID != fixture.buyerTwoID || reclaimed.Task.ClaimedByDeviceID == nil || *reclaimed.Task.ClaimedByDeviceID != fixture.deviceTwoID { t.Fatalf("reclaimed task = %+v", reclaimed.Task) } assertTaskEventCount(t, fixture.db, task.ID, "TASK_RECLAIMED", 1) _, err = fixture.store.ClaimNext(ctx, firstRequest) assertLifecycleError(t, err, usecase.ErrClaimReplayExpired) _, err = fixture.store.StartTask( ctx, lifecycleStartRequest( fixture, *first.Task, fixture.buyerOneID, fixture.deviceOneID, firstRequest.ClaimTokenHash, 42, reclaimAt, ), ) assertLifecycleError(t, err, usecase.ErrClaimInvalid) } func TestLifecycleRepositoryReclaimsOwnExpiredClaimBeforeFIFO( t *testing.T, ) { fixture := newLifecycleFixture(t) ctx := context.Background() now := lifecycleTime(45) ownTask := createLifecycleTask( t, fixture, 1, now.Add(-5*time.Minute), ) recordReadyHeartbeat( t, fixture, fixture.buyerOneID, fixture.deviceOneID, now.Add(-4*time.Minute), ) firstRequest := lifecycleClaimRequest( fixture, fixture.buyerOneID, fixture.deviceOneID, 45, now.Add(-4*time.Minute), now.Add(-time.Minute), now.Add(-6*time.Minute), ) first, err := fixture.store.ClaimNext(ctx, firstRequest) if err != nil || first.Task == nil || first.Task.ID != ownTask.ID { t.Fatalf("first claim = %+v, error = %v", first, err) } olderPending := createLifecycleTask( t, fixture, 2, now.Add(-10*time.Minute), ) reclaimRequest := lifecycleClaimRequest( fixture, fixture.buyerOneID, fixture.deviceOneID, 46, now, now.Add(10*time.Minute), now.Add(-10*time.Minute), ) reclaimed, err := fixture.store.ClaimNext(ctx, reclaimRequest) if err != nil || reclaimed.Task == nil || reclaimed.Task.ID != ownTask.ID { t.Fatalf("own expired reclaim = %+v, error = %v", reclaimed, err) } if reclaimed.Task.ClaimGeneration != 2 { t.Fatalf("reclaimed generation = %d", reclaimed.Task.ClaimGeneration) } pending, err := fixture.store.GetTaskDetail( ctx, lifecycleCreatorSubject, olderPending.ID, ) if err != nil || pending.Task.Status != domain.TaskStatusPending { t.Fatalf("older pending task = %+v, error = %v", pending.Task, err) } } func TestLifecycleRepositoryConcurrentClaimHasSingleWinner(t *testing.T) { fixture := newLifecycleFixture(t) secondStore := openSecondLifecycleStore(t, fixture.db) now := lifecycleTime(0) task := createLifecycleTask(t, fixture, 1, now.Add(-time.Minute)) recordReadyHeartbeat( t, fixture, fixture.buyerOneID, fixture.deviceOneID, now, ) recordReadyHeartbeat( t, fixture, fixture.buyerTwoID, fixture.deviceTwoID, now, ) requests := []usecase.ClaimNextRepositoryRequest{ lifecycleClaimRequest( fixture, fixture.buyerOneID, fixture.deviceOneID, 50, now, now.Add(10*time.Minute), now.Add(-2*time.Minute), ), lifecycleClaimRequest( fixture, fixture.buyerTwoID, fixture.deviceTwoID, 51, now, now.Add(10*time.Minute), now.Add(-2*time.Minute), ), } stores := []*repository.Store{fixture.store, secondStore} type claimOutcome struct { result usecase.ClaimNextRepositoryResult err error } start := make(chan struct{}) outcomes := make(chan claimOutcome, len(requests)) var workers sync.WaitGroup for index, request := range requests { request := request store := stores[index] workers.Add(1) go func() { defer workers.Done() <-start result, err := store.ClaimNext( context.Background(), request, ) outcomes <- claimOutcome{result: result, err: err} }() } close(start) workers.Wait() close(outcomes) claimed := 0 empty := 0 for outcome := range outcomes { if outcome.err != nil { t.Fatalf("concurrent ClaimNext() error = %v", outcome.err) } if outcome.result.Task == nil { empty++ continue } claimed++ if outcome.result.Task.ID != task.ID { t.Fatalf("concurrent claimed task = %+v", outcome.result.Task) } } if claimed != 1 || empty != 1 { t.Fatalf("concurrent outcomes: claimed=%d empty=%d", claimed, empty) } assertTaskEventCount(t, fixture.db, task.ID, "TASK_CLAIMED", 1) } func openSecondLifecycleStore( t *testing.T, first *sql.DB, ) *repository.Store { t.Helper() var sequence int var name string var path string if err := first.QueryRow(`PRAGMA database_list`).Scan( &sequence, &name, &path, ); err != nil { t.Fatalf("read lifecycle database path: %v", err) } db, err := platformdatabase.Open(context.Background(), path) if err != nil { t.Fatalf("open second lifecycle database connection: %v", err) } t.Cleanup(func() { _ = db.Close() }) store, err := repository.New(db) if err != nil { t.Fatalf("repository.New(second connection): %v", err) } return store } func TestLifecycleRepositoryStartIsIdempotentAndCreatesOneExecution( t *testing.T, ) { fixture, claimed, tokenHash, claimAt := prepareLifecycleClaim(t, 60) ctx := context.Background() startAt := claimAt.Add(time.Minute) request := lifecycleStartRequest( fixture, claimed, fixture.buyerOneID, fixture.deviceOneID, tokenHash, 61, startAt, ) started, err := fixture.store.StartTask(ctx, request) if err != nil || started.Replayed { t.Fatalf("StartTask() = %+v, error = %v", started, err) } if started.Task.Status != domain.TaskStatusRunning || started.Task.Version != claimed.Version+1 || started.Task.ClaimExpiresAt == nil || !started.Task.ClaimExpiresAt.Equal(request.ExpiresAt) || started.Execution.ID != request.Execution.ID || started.Execution.AttemptNo != 1 || started.Execution.LastHeartbeatAt == nil || !started.Execution.LastHeartbeatAt.Equal(startAt) { t.Fatalf("started task/execution = %+v / %+v", started.Task, started.Execution) } claimReplayRequest := lifecycleClaimRequest( fixture, fixture.buyerOneID, fixture.deviceOneID, 60, startAt, startAt.Add(10*time.Minute), startAt.Add(-2*time.Minute), ) claimReplay, err := fixture.store.ClaimNext(ctx, claimReplayRequest) if err != nil || !claimReplay.Replayed || claimReplay.Task == nil || claimReplay.Task.ID != started.Task.ID || claimReplay.Task.Status != domain.TaskStatusRunning { t.Fatalf( "claim replay after start = %+v, error = %v", claimReplay, err, ) } replayRequest := request replayRequest.Execution.ID = uuid(699) replayRequest.Event = lifecycleEvent( 69, claimed.ID, "TASK_STARTED", startAt.Add(time.Second), fixture.buyerOneID, fixture.deviceOneID, ) replayed, err := fixture.store.StartTask(ctx, replayRequest) if err != nil || !replayed.Replayed || replayed.Execution.ID != started.Execution.ID || replayed.Task.ID != started.Task.ID { t.Fatalf("replayed StartTask() = %+v, error = %v", replayed, err) } assertExecutionCount(t, fixture.db, claimed.ID, 1) assertTaskEventCount(t, fixture.db, claimed.ID, "TASK_STARTED", 1) conflict := request conflict.RequestHash = lifecycleHash("different-start") _, err = fixture.store.StartTask(ctx, conflict) assertLifecycleError(t, err, usecase.ErrIdempotencyConflict) } func TestLifecycleRepositoryStartRejectsVersionTokenAndExpiredLease( t *testing.T, ) { tests := []struct { name string change func(*usecase.StartTaskRepositoryRequest) expectedErr error }{ { name: "version", change: func(request *usecase.StartTaskRepositoryRequest) { request.ExpectedVersion++ }, expectedErr: usecase.ErrTaskVersionConflict, }, { name: "token", change: func(request *usecase.StartTaskRepositoryRequest) { request.ClaimTokenHash = lifecycleHash("wrong-token") }, expectedErr: usecase.ErrClaimInvalid, }, { name: "expired at boundary", change: func(request *usecase.StartTaskRepositoryRequest) { request.Now = request.ExpiresAt.Add(9 * time.Minute) request.ExpiresAt = request.Now.Add(90 * time.Second) }, expectedErr: usecase.ErrClaimExpired, }, } for index, test := range tests { t.Run(test.name, func(t *testing.T) { fixture, claimed, tokenHash, claimAt := prepareLifecycleClaim( t, 70+index*10, ) request := lifecycleStartRequest( fixture, claimed, fixture.buyerOneID, fixture.deviceOneID, tokenHash, 71+index*10, claimAt.Add(time.Minute), ) if test.name == "expired at boundary" { request.Now = *claimed.ClaimExpiresAt request.ExpiresAt = request.Now.Add(90 * time.Second) } else { test.change(&request) } _, err := fixture.store.StartTask( context.Background(), request, ) assertLifecycleError(t, err, test.expectedErr) assertExecutionCount(t, fixture.db, claimed.ID, 0) assertTaskEventCount( t, fixture.db, claimed.ID, "TASK_STARTED", 0, ) }) } } func TestLifecycleRepositoryHeartbeatExtendsLeaseWithoutEvent( t *testing.T, ) { fixture, claimed, tokenHash, claimAt := prepareLifecycleClaim(t, 100) ctx := context.Background() startRequest := lifecycleStartRequest( fixture, claimed, fixture.buyerOneID, fixture.deviceOneID, tokenHash, 101, claimAt.Add(time.Minute), ) started, err := fixture.store.StartTask(ctx, startRequest) if err != nil { t.Fatalf("StartTask() error = %v", err) } eventCountBefore := taskEventTotal(t, fixture.db, claimed.ID) heartbeatAt := startRequest.Now.Add(30 * time.Second) minimumExpiry := heartbeatAt.Add(90 * time.Second) heartbeat, err := fixture.store.HeartbeatTask( ctx, usecase.TaskHeartbeatRepositoryRequest{ UserID: fixture.buyerOneID, DeviceID: fixture.deviceOneID, TaskID: claimed.ID, ExecutionID: started.Execution.ID, ClaimGeneration: claimed.ClaimGeneration, ClaimTokenHash: tokenHash, Step: "SEARCH", Now: heartbeatAt, MinimumExpiry: minimumExpiry, }, ) if err != nil { t.Fatalf("HeartbeatTask() error = %v", err) } if heartbeat.Task.Version != started.Task.Version+1 || heartbeat.Task.ClaimExpiresAt == nil || heartbeat.Task.ClaimExpiresAt.Before(minimumExpiry) || heartbeat.Execution.CurrentStep != "SEARCH" || heartbeat.Execution.LastHeartbeatAt == nil || !heartbeat.Execution.LastHeartbeatAt.Equal(heartbeatAt) || heartbeat.CancelRequested { t.Fatalf("heartbeat = %+v", heartbeat) } if after := taskEventTotal(t, fixture.db, claimed.ID); after != eventCountBefore { t.Fatalf("heartbeat appended event: before=%d after=%d", eventCountBefore, after) } assertDeviceLastSeen(t, fixture.db, fixture.deviceOneID, heartbeatAt) expiredAt := *heartbeat.Task.ClaimExpiresAt _, err = fixture.store.HeartbeatTask( ctx, usecase.TaskHeartbeatRepositoryRequest{ UserID: fixture.buyerOneID, DeviceID: fixture.deviceOneID, TaskID: claimed.ID, ExecutionID: started.Execution.ID, ClaimGeneration: claimed.ClaimGeneration, ClaimTokenHash: tokenHash, Step: "SCAN_RESULTS", Now: expiredAt, MinimumExpiry: expiredAt.Add(90 * time.Second), }, ) assertLifecycleError(t, err, usecase.ErrClaimExpired) expiredHeartbeat, err := fixture.store.HeartbeatTask( ctx, usecase.TaskHeartbeatRepositoryRequest{ UserID: fixture.buyerOneID, DeviceID: fixture.deviceOneID, TaskID: claimed.ID, ExecutionID: started.Execution.ID, ClaimGeneration: claimed.ClaimGeneration, ClaimTokenHash: tokenHash, Step: "SAFE_STOPPED", Now: expiredAt, MinimumExpiry: expiredAt.Add(90 * time.Second), }, ) if err != nil { t.Fatalf("expired HeartbeatTask() error = %v", err) } if expiredHeartbeat.Task.ClaimExpiresAt == nil || !expiredHeartbeat.Task.ClaimExpiresAt.Equal(expiredAt) || expiredHeartbeat.Execution.CurrentStep != "SAFE_STOPPED" { t.Fatalf("expired heartbeat = %+v", expiredHeartbeat) } } func TestLifecycleRepositoryReleaseIsIdempotentAndInvalidatesClaim( t *testing.T, ) { fixture, claimed, tokenHash, claimAt := prepareLifecycleClaim(t, 120) ctx := context.Background() request := usecase.ReleaseTaskRepositoryRequest{ UserID: fixture.buyerOneID, DeviceID: fixture.deviceOneID, TaskID: claimed.ID, ClaimGeneration: claimed.ClaimGeneration, ClaimTokenHash: tokenHash, ExpectedVersion: claimed.Version, IdempotencyKey: "release-121", RequestHash: lifecycleHash("release-121"), Now: claimAt.Add(time.Minute), Event: lifecycleEvent( 121, claimed.ID, "TASK_RELEASED", claimAt.Add(time.Minute), fixture.buyerOneID, fixture.deviceOneID, ), } released, replayed, err := fixture.store.ReleaseTask(ctx, request) if err != nil || replayed { t.Fatalf("ReleaseTask() = %+v, %t, error = %v", released, replayed, err) } if released.Status != domain.TaskStatusPending || released.Version != claimed.Version+1 || released.ClaimedByUserID != nil || released.ClaimedByDeviceID != nil || released.ClaimTokenHash != nil || released.ClaimIssuedAt != nil || released.ClaimExpiresAt != nil { t.Fatalf("released task = %+v", released) } replayRequest := request replayRequest.Event = lifecycleEvent( 122, claimed.ID, "TASK_RELEASED", request.Now.Add(time.Second), fixture.buyerOneID, fixture.deviceOneID, ) replayedTask, replayed, err := fixture.store.ReleaseTask(ctx, replayRequest) if err != nil || !replayed || replayedTask.ID != claimed.ID { t.Fatalf("release replay = %+v, %t, error = %v", replayedTask, replayed, err) } assertTaskEventCount(t, fixture.db, claimed.ID, "TASK_RELEASED", 1) claimReplay := lifecycleClaimRequest( fixture, fixture.buyerOneID, fixture.deviceOneID, 120, claimAt, *claimed.ClaimExpiresAt, claimAt.Add(-2*time.Minute), ) claimReplay.ClaimTokenHash = tokenHash _, err = fixture.store.ClaimNext(ctx, claimReplay) assertLifecycleError(t, err, usecase.ErrClaimReplayExpired) } func TestLifecycleRepositoryAdminCancelClaimedTaskIsImmediate( t *testing.T, ) { fixture, claimed, _, claimAt := prepareLifecycleClaim(t, 130) adminID := fixture.adminUserID canceledAt := claimAt.Add(time.Minute) canceled, err := fixture.store.CancelTask( context.Background(), lifecycleCreatorSubject, claimed.ID, "admin canceled before start", canceledAt, domain.TaskEvent{ ID: uuid(131), TaskID: claimed.ID, ActorUserID: &adminID, Type: "TASK_CANCELED", Message: "task canceled", OccurredAt: canceledAt, }, ) if err != nil { t.Fatalf("CancelTask(CLAIMED) error = %v", err) } if canceled.Status != domain.TaskStatusCanceled || canceled.ClaimedByUserID != nil || canceled.ClaimedByDeviceID != nil || canceled.ClaimTokenHash != nil || canceled.ClaimExpiresAt != nil { t.Fatalf("canceled claimed task = %+v", canceled) } assertTaskEventCount(t, fixture.db, claimed.ID, "TASK_CANCELED", 1) } func TestLifecycleRepositoryCancelAcknowledgementEndsExecution( t *testing.T, ) { fixture, claimed, tokenHash, claimAt := prepareLifecycleClaim(t, 140) ctx := context.Background() startRequest := lifecycleStartRequest( fixture, claimed, fixture.buyerOneID, fixture.deviceOneID, tokenHash, 141, claimAt.Add(time.Minute), ) started, err := fixture.store.StartTask(ctx, startRequest) if err != nil { t.Fatalf("StartTask() error = %v", err) } expiredAt := *started.Task.ClaimExpiresAt cancelAt := startRequest.Now.Add(10 * time.Second) adminID := fixture.adminUserID cancelRequested, err := fixture.store.CancelTask( ctx, lifecycleCreatorSubject, claimed.ID, "admin requested stop", cancelAt, domain.TaskEvent{ ID: uuid(142), TaskID: claimed.ID, ActorUserID: &adminID, Type: "TASK_CANCELED", Message: "task canceled", OccurredAt: cancelAt, }, ) if err != nil { t.Fatalf("CancelTask(RUNNING) error = %v", err) } if cancelRequested.Status != domain.TaskStatusRunning || cancelRequested.CancelRequestedAt == nil || cancelRequested.CancelRequestedByUserID == nil || *cancelRequested.CancelRequestedByUserID != fixture.adminUserID || cancelRequested.CanceledAt != nil { t.Fatalf("cancel request = %+v", cancelRequested) } repeated, err := fixture.store.CancelTask( ctx, lifecycleCreatorSubject, claimed.ID, "must not replace first request", cancelAt.Add(time.Second), domain.TaskEvent{ ID: uuid(143), TaskID: claimed.ID, ActorUserID: &adminID, Type: "TASK_CANCELED", Message: "task canceled", OccurredAt: cancelAt.Add(time.Second), }, ) if err != nil || repeated.Version != cancelRequested.Version { t.Fatalf("repeat CancelTask() = %+v, error = %v", repeated, err) } assertTaskEventCount( t, fixture.db, claimed.ID, "TASK_CANCEL_REQUESTED", 1, ) detail, err := fixture.store.GetTaskDetail( ctx, lifecycleCreatorSubject, claimed.ID, ) if err != nil || detail.Execution == nil { t.Fatalf("GetTaskDetail() = %+v, error = %v", detail, err) } heartbeat, err := fixture.store.HeartbeatTask( ctx, usecase.TaskHeartbeatRepositoryRequest{ UserID: fixture.buyerOneID, DeviceID: fixture.deviceOneID, TaskID: claimed.ID, ExecutionID: started.Execution.ID, ClaimGeneration: claimed.ClaimGeneration, ClaimTokenHash: tokenHash, Step: "SAFE_STOPPED", Now: expiredAt, MinimumExpiry: expiredAt.Add(90 * time.Second), }, ) if err != nil || !heartbeat.CancelRequested { t.Fatalf("cancel heartbeat = %+v, error = %v", heartbeat, err) } ackAt := expiredAt.Add(time.Second) request := usecase.CancelAcknowledgementRepositoryRequest{ UserID: fixture.buyerOneID, DeviceID: fixture.deviceOneID, TaskID: claimed.ID, ExecutionID: started.Execution.ID, ClaimGeneration: claimed.ClaimGeneration, ClaimTokenHash: tokenHash, ExpectedVersion: heartbeat.Task.Version, IdempotencyKey: "cancel-ack-144", RequestHash: lifecycleHash("cancel-ack-144"), Now: ackAt, Event: lifecycleEvent( 144, claimed.ID, "TASK_CANCELED", ackAt, fixture.buyerOneID, fixture.deviceOneID, ), } canceled, replayed, err := fixture.store.AcknowledgeTaskCancellation( ctx, request, ) if err != nil || replayed { t.Fatalf( "AcknowledgeTaskCancellation() = %+v, %t, error = %v", canceled, replayed, err, ) } if canceled.Status != domain.TaskStatusCanceled || canceled.Version != heartbeat.Task.Version+1 || canceled.CanceledAt == nil || !canceled.CanceledAt.Equal(ackAt) || canceled.ClaimedByDeviceID != nil || canceled.ClaimTokenHash != nil { t.Fatalf("canceled task = %+v", canceled) } finished, err := fixture.store.GetTaskDetail( ctx, lifecycleCreatorSubject, claimed.ID, ) if err != nil || finished.Execution == nil || finished.Execution.FinishedAt == nil || !finished.Execution.FinishedAt.Equal(ackAt) { t.Fatalf("finished detail = %+v, error = %v", finished, err) } replayRequest := request replayRequest.Event = lifecycleEvent( 145, claimed.ID, "TASK_CANCELED", ackAt.Add(time.Second), fixture.buyerOneID, fixture.deviceOneID, ) replayedTask, replayed, err := fixture.store.AcknowledgeTaskCancellation( ctx, replayRequest, ) if err != nil || !replayed || replayedTask.ID != claimed.ID { t.Fatalf( "cancel acknowledgement replay = %+v, %t, error = %v", replayedTask, replayed, err, ) } assertTaskEventCount(t, fixture.db, claimed.ID, "TASK_CANCELED", 1) } func newLifecycleFixture(t *testing.T) *lifecycleFixture { t.Helper() db := openDatabase(t) fixture := &lifecycleFixture{ db: db, adminUserID: uuid(999), buyerOneID: uuid(901), buyerTwoID: uuid(902), deviceOneID: uuid(911), deviceTwoID: uuid(912), } createdAt := lifecycleTime(-10) seedLifecycleUser( t, db, fixture.adminUserID, "lifecycle-admin", domain.UserRoleAdmin, createdAt, ) seedLifecycleUser( t, db, fixture.buyerOneID, "lifecycle-buyer-one", domain.UserRoleBuyer, createdAt, ) seedLifecycleUser( t, db, fixture.buyerTwoID, "lifecycle-buyer-two", domain.UserRoleBuyer, createdAt, ) seedLifecycleDevice( t, db, fixture.deviceOneID, "lifecycle-device-one", fixture.buyerOneID, lifecycleHash("device-one-token"), createdAt, ) seedLifecycleDevice( t, db, fixture.deviceTwoID, "lifecycle-device-two", fixture.buyerTwoID, lifecycleHash("device-two-token"), createdAt, ) store, err := repository.New(db) if err != nil { t.Fatalf("repository.New() error = %v", err) } fixture.store = store return fixture } func seedLifecycleUser( t *testing.T, db *sql.DB, id string, username string, role domain.UserRole, createdAt time.Time, ) { t.Helper() if _, err := db.ExecContext( context.Background(), `INSERT INTO users ( id, username, password_hash, role, is_active, created_at, updated_at ) VALUES (?, ?, 'test-only-password-hash', ?, 1, ?, ?)`, id, username, role, lifecycleTimestamp(createdAt), lifecycleTimestamp(createdAt), ); err != nil { t.Fatalf("seed user %s: %v", username, err) } } func seedLifecycleDevice( t *testing.T, db *sql.DB, id string, name string, userID string, tokenHash string, createdAt time.Time, ) { t.Helper() if _, err := db.ExecContext( context.Background(), `INSERT INTO devices ( id, name, token_hash, bound_user_id, is_enabled, created_at, updated_at ) VALUES (?, ?, ?, ?, 1, ?, ?)`, id, name, tokenHash, userID, lifecycleTimestamp(createdAt), lifecycleTimestamp(createdAt), ); err != nil { t.Fatalf("seed device %s: %v", name, err) } } func createLifecycleTask( t *testing.T, fixture *lifecycleFixture, index int, createdAt time.Time, ) domain.PurchaseTask { t.Helper() asset := testAsset(index, createdAt) if _, _, err := fixture.store.CreateAssetIdempotent( context.Background(), asset, fmt.Sprintf("lifecycle-asset-%d", index), lifecycleHash(fmt.Sprintf("lifecycle-asset-%d", index)), ); err != nil { t.Fatalf("create lifecycle asset %d: %v", index, err) } task := testTask( index, asset.ID, fmt.Sprintf("lifecycle-source-%d", index), createdAt, ) created, _, err := fixture.store.CreateTaskIdempotent( context.Background(), task, lifecycleEvent( 1000+index, task.ID, "TASK_CREATED", createdAt, fixture.adminUserID, "", ), fmt.Sprintf("lifecycle-task-%d", index), lifecycleHash(fmt.Sprintf("lifecycle-task-%d", index)), ) if err != nil { t.Fatalf("create lifecycle task %d: %v", index, err) } return created } func lifecycleDeviceHeartbeat( userID string, deviceID string, reportedAt time.Time, accessibilityEnabled bool, pddInstalled bool, ) usecase.DeviceHeartbeatUpdate { return usecase.DeviceHeartbeatUpdate{ UserID: userID, DeviceID: deviceID, AppVersion: "0.1.0", AndroidVersion: "16", PDDVersion: "8.17.0", AccessibilityEnabled: accessibilityEnabled, PDDInstalled: pddInstalled, ReportedAt: reportedAt, } } func recordReadyHeartbeat( t *testing.T, fixture *lifecycleFixture, userID string, deviceID string, reportedAt time.Time, ) { t.Helper() if _, err := fixture.store.RecordDeviceHeartbeat( context.Background(), lifecycleDeviceHeartbeat( userID, deviceID, reportedAt, true, true, ), ); err != nil { t.Fatalf("record ready heartbeat: %v", err) } } func lifecycleClaimRequest( fixture *lifecycleFixture, userID string, deviceID string, index int, now time.Time, expiresAt time.Time, readinessAfter time.Time, ) usecase.ClaimNextRepositoryRequest { key := fmt.Sprintf("claim-%d", index) return usecase.ClaimNextRepositoryRequest{ CreatorSubject: lifecycleCreatorSubject, UserID: userID, DeviceID: deviceID, IdempotencyKey: key, RequestHash: lifecycleHash(key), ClaimTokenHash: lifecycleHash("token-" + key), Now: now, ExpiresAt: expiresAt, ReadinessAfter: readinessAfter, Event: lifecycleEvent( 2000+index, "", "TASK_CLAIMED", now, userID, deviceID, ), } } func lifecycleStartRequest( fixture *lifecycleFixture, task domain.PurchaseTask, userID string, deviceID string, tokenHash string, index int, now time.Time, ) usecase.StartTaskRepositoryRequest { key := fmt.Sprintf("start-%d", index) return usecase.StartTaskRepositoryRequest{ UserID: userID, DeviceID: deviceID, TaskID: task.ID, ClaimGeneration: task.ClaimGeneration, ClaimTokenHash: tokenHash, ExpectedVersion: task.Version, IdempotencyKey: key, RequestHash: lifecycleHash(key), Now: now, ExpiresAt: now.Add(90 * time.Second), Execution: domain.TaskExecution{ ID: uuid(3000 + index), TaskID: task.ID, ClaimGeneration: task.ClaimGeneration, UserID: userID, DeviceID: deviceID, CurrentStep: "PREFLIGHT", OrderSubmitted: false, StartedAt: now, }, Event: lifecycleEvent( 4000+index, task.ID, "TASK_STARTED", now, userID, deviceID, ), } } func prepareLifecycleClaim( t *testing.T, index int, ) (*lifecycleFixture, domain.PurchaseTask, string, time.Time) { t.Helper() fixture := newLifecycleFixture(t) now := lifecycleTime(index) _ = createLifecycleTask(t, fixture, 1, now.Add(-time.Minute)) recordReadyHeartbeat( t, fixture, fixture.buyerOneID, fixture.deviceOneID, now, ) request := lifecycleClaimRequest( fixture, fixture.buyerOneID, fixture.deviceOneID, index, now, now.Add(10*time.Minute), now.Add(-2*time.Minute), ) result, err := fixture.store.ClaimNext(context.Background(), request) if err != nil || result.Task == nil { t.Fatalf("prepare ClaimNext() = %+v, error = %v", result, err) } return fixture, *result.Task, request.ClaimTokenHash, now } func lifecycleEvent( index int, taskID string, eventType string, occurredAt time.Time, userID string, deviceID string, ) domain.TaskEvent { event := domain.TaskEvent{ ID: uuid(5000 + index), TaskID: taskID, Type: eventType, Message: "lifecycle test event", OccurredAt: occurredAt, } if userID != "" { event.ActorUserID = &userID } if deviceID != "" { event.ActorDeviceID = &deviceID } return event } func assertLifecycleError( t *testing.T, err error, expected error, ) { t.Helper() if !errors.Is(err, expected) { t.Fatalf("error = %v, want errors.Is(%v)", err, expected) } } func assertTaskEventCount( t *testing.T, db *sql.DB, taskID string, eventType string, expected int, ) { t.Helper() var count int if err := db.QueryRowContext( context.Background(), `SELECT COUNT(*) FROM task_events WHERE task_id = ? AND event_type = ?`, taskID, eventType, ).Scan(&count); err != nil { t.Fatalf("count %s events: %v", eventType, err) } if count != expected { t.Fatalf("%s event count = %d, want %d", eventType, count, expected) } } func taskEventTotal(t *testing.T, db *sql.DB, taskID string) int { t.Helper() var count int if err := db.QueryRowContext( context.Background(), `SELECT COUNT(*) FROM task_events WHERE task_id = ?`, taskID, ).Scan(&count); err != nil { t.Fatalf("count task events: %v", err) } return count } func assertExecutionCount( t *testing.T, db *sql.DB, taskID string, expected int, ) { t.Helper() var count int if err := db.QueryRowContext( context.Background(), `SELECT COUNT(*) FROM task_executions WHERE task_id = ?`, taskID, ).Scan(&count); err != nil { t.Fatalf("count task executions: %v", err) } if count != expected { t.Fatalf("execution count = %d, want %d", count, expected) } } func assertDeviceLastSeen( t *testing.T, db *sql.DB, deviceID string, expected time.Time, ) { t.Helper() var value string if err := db.QueryRowContext( context.Background(), `SELECT last_seen_at FROM devices WHERE id = ?`, deviceID, ).Scan(&value); err != nil { t.Fatalf("query device last_seen_at: %v", err) } parsed, err := time.Parse(time.RFC3339Nano, value) if err != nil { t.Fatalf("parse device last_seen_at %q: %v", value, err) } if !parsed.Equal(expected) { t.Fatalf("device last_seen_at = %s, want %s", parsed, expected) } } func lifecycleHash(value string) string { sum := sha256.Sum256([]byte(value)) return hex.EncodeToString(sum[:]) } func lifecycleTime(offsetMinutes int) time.Time { return time.Date(2026, 7, 26, 8, 0, 0, 0, time.UTC). Add(time.Duration(offsetMinutes) * time.Minute) } func lifecycleTimestamp(value time.Time) string { return value.UTC().Format(time.RFC3339Nano) }