Files
worker/app/models/dead_worker_reaper_test.go
Gleb Tv 2c884c5612
Некоторые проверки не удались
CI / test (push) Successful in 2m5s
Docker / Build and publish worker image (push) Failing after 31s
refactor: adopt worker module path
2026-07-13 17:56:12 +03:00

275 строки
9.7 KiB
Go

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
}