package models_test import ( "context" "fmt" "testing" "time" "github.com/google/uuid" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "gorm.io/datatypes" "rocketgit.ru/rsmon/worker/app/models" ) // TestReapDeadWorkers_SkipsHealthyAndKillsSilent mirrors the // reaper-shape test for ReapExpiredTasks in task_test.go: seed three // workers (heartbeat fresh / heartbeat stale / already dead) plus two // leased tasks owned by the stale worker and one leased task on a // healthy worker (which must NOT be touched). Then assert: // // - the fresh worker is left active // - the silent worker flips to "dead" // - the already-dead worker is left dead (idempotent) // - the silent worker's leased tasks are reset to queued // - the healthy worker's leased task is untouched func TestReapDeadWorkers_SkipsHealthyAndKillsSilent(t *testing.T) { models.Drop() models.Migrate() seedRegion(t, "test") healthy := &models.WorkerNode{ WorkerID: "w-healthy", RegionCode: "test", Status: "active", AuthToken: "tok-healthy", LastSeen: timePtr(time.Now()), } stale := &models.WorkerNode{ WorkerID: "w-stale", RegionCode: "test", Status: "active", AuthToken: "tok-stale", LastSeen: timePtr(time.Now().Add(-models.DeadWorkerHeartbeatTimeout - time.Minute)), } alreadyDead := &models.WorkerNode{ WorkerID: "w-dead", RegionCode: "test", Status: "dead", AuthToken: "tok-dead", LastSeen: timePtr(time.Now().Add(-time.Hour)), } require.NoError(t, models.DB().Create(healthy).Error) require.NoError(t, models.DB().Create(stale).Error) require.NoError(t, models.DB().Create(alreadyDead).Error) // Two leased tasks on the stale worker — both must come back. task1 := mustLeaseTask(t, stale.WorkerID, "test-acct") task2 := mustLeaseTask(t, stale.WorkerID, "test-acct") // One leased task on the healthy worker — must stay leased. healthyTask := mustLeaseTask(t, healthy.WorkerID, "test-acct") // Already-dead worker with a leased task — not part of THIS reaper's // reaping set (it would only be touched by a fresh reaper pass), so // leave it leased to demonstrate that we don't accidentally reassign // from prior-dead leases too. deadPriorTask := mustLeaseTask(t, alreadyDead.WorkerID, "test-acct") reaped, reassigned, err := models.ReapDeadWorkers() require.NoError(t, err) assert.Equal(t, 1, reaped, "only the stale worker should flip (already-dead is left untouched)") // 3 tasks come back: the 2 leased by the stale worker (just flipped // to dead) + the 1 leased by the prior-dead worker, which had never // been cleaned up because no previous reaper ran. The reaper matches // dead workers by status, so any leased task on a dead worker is a // stranded lease and must be returned to the pool regardless of when // the worker flipped. assert.Equal(t, 3, reassigned) var healthyRow, staleRow, deadRow models.WorkerNode require.NoError(t, models.DB().First(&healthyRow, healthy.ID).Error) assert.Equal(t, "active", healthyRow.Status, "healthy worker must stay active") require.NoError(t, models.DB().First(&staleRow, stale.ID).Error) assert.Equal(t, "dead", staleRow.Status, "stale worker should flip to dead") require.NoError(t, models.DB().First(&deadRow, alreadyDead.ID).Error) assert.Equal(t, "dead", deadRow.Status, "already-dead worker should not be touched") got := func(id int64) models.Task { var row models.Task require.NoError(t, models.DB().First(&row, id).Error) return row } t1 := got(task1.ID) assert.Equal(t, models.TaskStateQueued, t1.State, "stale worker task must come back to queued") assert.Empty(t, t1.LeaseOwner) assert.Nil(t, t1.LeaseExpiresAt) t2 := got(task2.ID) assert.Equal(t, models.TaskStateQueued, t2.State) assert.Empty(t, t2.LeaseOwner) assert.Nil(t, t2.LeaseExpiresAt) ht := got(healthyTask.ID) assert.Equal(t, models.TaskStateLeased, ht.State, "healthy worker's lease must be untouched") assert.Equal(t, healthy.WorkerID, ht.LeaseOwner) dt := got(deadPriorTask.ID) assert.Equal(t, models.TaskStateQueued, dt.State, "prior-dead task must also be returned to the pool — any leased task on a dead worker is a stranded lease") assert.Empty(t, dt.LeaseOwner) } // TestReapDeadWorkers_NoOpOnHealthyFleet verifies the cheap path: when // no workers are stale the reaper returns (0, 0, nil) without doing any // work. Mirrors the "reaped = int(r.RowsAffected)" branch in // ReapExpiredTasks. func TestReapDeadWorkers_NoOpOnHealthyFleet(t *testing.T) { models.Drop() models.Migrate() seedRegion(t, "test") w := &models.WorkerNode{ WorkerID: "w-only-healthy", RegionCode: "test", Status: "active", AuthToken: uuid.NewString(), LastSeen: timePtr(time.Now()), } require.NoError(t, models.DB().Create(w).Error) reaped, reassigned, err := models.ReapDeadWorkers() require.NoError(t, err) assert.Equal(t, 0, reaped) assert.Equal(t, 0, reassigned) var row models.WorkerNode require.NoError(t, models.DB().First(&row, w.ID).Error) assert.Equal(t, "active", row.Status) } // TestReapDeadWorkers_OnlyReassignsLeasedNotOthers ensures that the // reaper does not touch queued/succeeded/failed_retry tasks on the dead // worker — only leased ones need to be returned to the queue. Tasks in // other states either belong to no one (queued) or are terminal/semi- // terminal and have their own audit trail. func TestReapDeadWorkers_OnlyReassignsLeasedNotOthers(t *testing.T) { models.Drop() models.Migrate() seedRegion(t, "test") now := time.Now().Add(-models.DeadWorkerHeartbeatTimeout - time.Minute) w := &models.WorkerNode{ WorkerID: "w-mix", RegionCode: "test", Status: "active", AuthToken: uuid.NewString(), LastSeen: &now, } require.NoError(t, models.DB().Create(w).Error) leased := mustLeaseTask(t, w.WorkerID, "acct-mix") succeeded := mustInsertTask(t, w.WorkerID, "acct-mix", models.TaskStateSucceeded) failedRetry := mustInsertTask(t, w.WorkerID, "acct-mix", models.TaskStateFailedRetry) failedPerm := mustInsertTask(t, w.WorkerID, "acct-mix", models.TaskStateFailedPerm) otherWorker := mustLeaseTask(t, "w-other", "acct-mix") reaped, reassigned, err := models.ReapDeadWorkers() require.NoError(t, err) assert.Equal(t, 1, reaped) assert.Equal(t, 1, reassigned, "exactly the one leased task on the dead worker") got := func(id int64) string { var row models.Task require.NoError(t, models.DB().First(&row, id).Error) return row.State } assert.Equal(t, models.TaskStateQueued, got(leased.ID)) assert.Equal(t, models.TaskStateSucceeded, got(succeeded.ID), "succeeded must not move") assert.Equal(t, models.TaskStateFailedRetry, got(failedRetry.ID), "failed_retry must not move") assert.Equal(t, models.TaskStateFailedPerm, got(failedPerm.ID), "failed_perm must not move") assert.Equal(t, models.TaskStateLeased, got(otherWorker.ID), "tasks leased by another worker must not move") } // TestStartDeadWorkerReaper_TickerFiresOnce is a smoke test for the // background helper: spin up the reaper with a tight 10ms ticker and // a cancellable context, wait for the first tick, then cancel the // context so the goroutine exits cleanly without leaking. Mirrors the // shape of how main.init() uses StartTaskReaper (it can't be torn // down, but for tests we always pass a cancellable context). func TestStartDeadWorkerReaper_TickerFiresOnce(t *testing.T) { models.Drop() models.Migrate() seedRegion(t, "test") w := &models.WorkerNode{ WorkerID: "w-ticker", RegionCode: "test", Status: "active", AuthToken: uuid.NewString(), LastSeen: timePtr(time.Now().Add(-2 * models.DeadWorkerHeartbeatTimeout)), } require.NoError(t, models.DB().Create(w).Error) ctx, cancel := context.WithCancel(context.Background()) defer cancel() models.StartDeadWorkerReaper(ctx, 10*time.Millisecond) deadline := time.Now().Add(2 * time.Second) var got models.WorkerNode for time.Now().Before(deadline) { require.NoError(t, models.DB().First(&got, w.ID).Error) if got.Status == "dead" { return } time.Sleep(20 * time.Millisecond) } t.Fatalf("reaper goroutine did not flip worker to dead within 2s; last status=%q", got.Status) } // --------------------------------------------------------------------------- // test helpers // --------------------------------------------------------------------------- func timePtr(t time.Time) *time.Time { return &t } // mustLeaseTask inserts a leased Task row owned by workerID. The state // is the only field that matters for reaper tests; payload/idempotency // are stubs. func mustLeaseTask(t *testing.T, workerID, accountLabel string) models.Task { t.Helper() return mustInsertTask(t, workerID, accountLabel, models.TaskStateLeased) } func mustInsertTask(t *testing.T, workerID, accountLabel string, state string) models.Task { t.Helper() leaseUntil := time.Now().Add(time.Hour) stamp := time.Now().UnixNano() jobID := fmt.Sprintf("reap-%s-%s-%d", state, workerID, stamp) // Cap at 64 chars (tasks.job_id VARCHAR(64)). The components above // already stay under the limit because mustLeaseTask keeps workerID // short ("w-mix", "w-stale", …) and state is bounded. if len(jobID) > 64 { jobID = jobID[:64] } idemp := fmt.Sprintf("reap:%s:%d", workerID, stamp) if len(idemp) > 255 { idemp = idemp[:255] } task := models.Task{ JobID: jobID, Kind: models.TaskKindNotification, State: state, AccountID: 1, Payload: datatypes.JSON([]byte(`{"method":"email"}`)), NotBefore: time.Now().Add(-time.Minute), LeaseOwner: workerID, LeaseExpiresAt: &leaseUntil, Attempts: 1, MaxAttempts: 5, IdempotencyKey: idemp, } if state != models.TaskStateLeased { task.LeaseOwner = "" task.LeaseExpiresAt = nil } require.NoError(t, models.DB().Create(&task).Error) require.NotZero(t, task.ID) return task }