diff --git a/backend/internal/handler/gateway_handler_warmup_intercept_unit_test.go b/backend/internal/handler/gateway_handler_warmup_intercept_unit_test.go index 685231fc91..8f3362f919 100644 --- a/backend/internal/handler/gateway_handler_warmup_intercept_unit_test.go +++ b/backend/internal/handler/gateway_handler_warmup_intercept_unit_test.go @@ -43,6 +43,12 @@ func (f *fakeSchedulerCache) RetireBucket(_ context.Context, _ service.Scheduler func (f *fakeSchedulerCache) ReopenBucket(_ context.Context, bucket service.SchedulerBucket) (service.SchedulerBucketWriteToken, error) { return service.SchedulerBucketWriteToken{Bucket: bucket, Epoch: 1}, nil } +func (f *fakeSchedulerCache) TryAcquireGroupLifecycleLease(_ context.Context, _ int64, _ time.Duration) (service.SchedulerGroupLifecycleLease, bool, error) { + return service.SchedulerGroupLifecycleLease{}, false, nil +} +func (f *fakeSchedulerCache) ReleaseGroupLifecycleLease(_ context.Context, _ service.SchedulerGroupLifecycleLease) error { + return nil +} func (f *fakeSchedulerCache) GetAccount(_ context.Context, id int64) (*service.Account, error) { for _, account := range f.accounts { if account != nil && account.ID == id { diff --git a/backend/internal/repository/account_repo_integration_test.go b/backend/internal/repository/account_repo_integration_test.go index 7c928072e4..f7e227e424 100644 --- a/backend/internal/repository/account_repo_integration_test.go +++ b/backend/internal/repository/account_repo_integration_test.go @@ -47,6 +47,14 @@ func (s *schedulerCacheRecorder) ReopenBucket(ctx context.Context, bucket servic return service.SchedulerBucketWriteToken{Bucket: bucket, Epoch: 1}, nil } +func (s *schedulerCacheRecorder) TryAcquireGroupLifecycleLease(_ context.Context, groupID int64, _ time.Duration) (service.SchedulerGroupLifecycleLease, bool, error) { + return service.SchedulerGroupLifecycleLease{GroupID: groupID, OwnerToken: "scheduler-cache-recorder"}, true, nil +} + +func (s *schedulerCacheRecorder) ReleaseGroupLifecycleLease(context.Context, service.SchedulerGroupLifecycleLease) error { + return nil +} + func (s *schedulerCacheRecorder) GetAccount(ctx context.Context, accountID int64) (*service.Account, error) { if s.accounts == nil { return nil, nil diff --git a/backend/internal/repository/scheduler_cache.go b/backend/internal/repository/scheduler_cache.go index 72d12d799b..9c098d0648 100644 --- a/backend/internal/repository/scheduler_cache.go +++ b/backend/internal/repository/scheduler_cache.go @@ -2,6 +2,8 @@ package repository import ( "context" + "crypto/rand" + "encoding/hex" "encoding/json" "fmt" "log/slog" @@ -33,6 +35,11 @@ const ( snapshotGraceTTLSeconds = 60 ) +const ( + schedulerGroupLifecycleLockPrefix = "sched:group:lifecycle-lock:" + schedulerGroupLifecycleOwnerTokenBytes = 16 +) + var ( captureBucketWriteTokenScript = redis.NewScript(` if redis.call('EXISTS', KEYS[2]) == 1 then @@ -124,6 +131,13 @@ if currentActive ~= false then end redis.call('DEL', KEYS[4], KEYS[5]) return currentEpoch +`) + + releaseGroupLifecycleLeaseScript = redis.NewScript(` +if redis.call('GET', KEYS[1]) == ARGV[1] then + return redis.call('DEL', KEYS[1]) +end +return 0 `) // activateSnapshotScript 原子 CAS 切换快照版本。 @@ -310,6 +324,57 @@ func (c *schedulerCache) ReopenBucket(ctx context.Context, bucket service.Schedu return service.SchedulerBucketWriteToken{Bucket: bucket, Epoch: result}, nil } +func (c *schedulerCache) TryAcquireGroupLifecycleLease(ctx context.Context, groupID int64, ttl time.Duration) (service.SchedulerGroupLifecycleLease, bool, error) { + if groupID <= 0 { + return service.SchedulerGroupLifecycleLease{}, false, fmt.Errorf("%w: group id must be positive", service.ErrSchedulerGroupLifecycleLeaseInvalid) + } + if ttl <= 0 { + return service.SchedulerGroupLifecycleLease{}, false, fmt.Errorf("%w: ttl must be positive", service.ErrSchedulerGroupLifecycleLeaseInvalid) + } + ownerToken, err := newSchedulerGroupLifecycleOwnerToken() + if err != nil { + return service.SchedulerGroupLifecycleLease{}, false, err + } + acquired, err := c.rdb.SetNX(ctx, schedulerGroupLifecycleLockKey(groupID), ownerToken, ttl).Result() + if err != nil { + return service.SchedulerGroupLifecycleLease{}, false, err + } + if !acquired { + return service.SchedulerGroupLifecycleLease{}, false, nil + } + return service.SchedulerGroupLifecycleLease{GroupID: groupID, OwnerToken: ownerToken}, true, nil +} + +func (c *schedulerCache) ReleaseGroupLifecycleLease(ctx context.Context, lease service.SchedulerGroupLifecycleLease) error { + if !lease.ValidFor(lease.GroupID) { + return service.ErrSchedulerGroupLifecycleLeaseInvalid + } + result, err := releaseGroupLifecycleLeaseScript.Run( + ctx, + c.rdb, + []string{schedulerGroupLifecycleLockKey(lease.GroupID)}, + lease.OwnerToken, + ).Int64() + if err != nil { + return err + } + if result == 0 { + return fmt.Errorf("%w: group=%d", service.ErrSchedulerGroupLifecycleLeaseLost, lease.GroupID) + } + if result != 1 { + return fmt.Errorf("release scheduler group lifecycle lease returned %d", result) + } + return nil +} + +func newSchedulerGroupLifecycleOwnerToken() (string, error) { + raw := make([]byte, schedulerGroupLifecycleOwnerTokenBytes) + if _, err := rand.Read(raw); err != nil { + return "", fmt.Errorf("generate scheduler group lifecycle owner token: %w", err) + } + return hex.EncodeToString(raw), nil +} + func (c *schedulerCache) SetSnapshot(ctx context.Context, bucket service.SchedulerBucket, token service.SchedulerBucketWriteToken, accounts []service.Account) error { if !token.ValidFor(bucket) { return fmt.Errorf("%w: bucket=%s", service.ErrSchedulerBucketWriteFenced, bucket.String()) @@ -535,6 +600,10 @@ func schedulerBucketKey(prefix string, bucket service.SchedulerBucket) string { return fmt.Sprintf("%s%d:%s:%s", prefix, bucket.GroupID, bucket.Platform, bucket.Mode) } +func schedulerGroupLifecycleLockKey(groupID int64) string { + return schedulerGroupLifecycleLockPrefix + strconv.FormatInt(groupID, 10) +} + func schedulerSnapshotKey(bucket service.SchedulerBucket, version string) string { return fmt.Sprintf("%s%d:%s:%s:v%s", schedulerSnapshotPrefix, bucket.GroupID, bucket.Platform, bucket.Mode, version) } diff --git a/backend/internal/repository/scheduler_cache_integration_test.go b/backend/internal/repository/scheduler_cache_integration_test.go index 00966d95b6..52673315c8 100644 --- a/backend/internal/repository/scheduler_cache_integration_test.go +++ b/backend/internal/repository/scheduler_cache_integration_test.go @@ -137,3 +137,40 @@ func TestSchedulerCacheRetireAndReopenFencesOldEpochIntegration(t *testing.T) { require.Len(t, snapshot, 1) require.Equal(t, account.ID, snapshot[0].ID) } + +func TestSchedulerCacheGroupLifecycleLeaseOwnerAndTTLIntegration(t *testing.T) { + ctx := context.Background() + rdb := testRedis(t) + cache := NewSchedulerCache(rdb) + const groupID int64 = 78 + const ttl = 500 * time.Millisecond + + first, acquired, err := cache.TryAcquireGroupLifecycleLease(ctx, groupID, ttl) + require.NoError(t, err) + require.True(t, acquired) + pttl, err := rdb.PTTL(ctx, schedulerGroupLifecycleLockKey(groupID)).Result() + require.NoError(t, err) + require.Positive(t, pttl) + require.LessOrEqual(t, pttl, ttl) + + var second service.SchedulerGroupLifecycleLease + require.Eventually(t, func() bool { + var acquireErr error + second, acquired, acquireErr = cache.TryAcquireGroupLifecycleLease(ctx, groupID, time.Minute) + return acquireErr == nil && acquired + }, 5*time.Second, 20*time.Millisecond) + require.NotEqual(t, first.OwnerToken, second.OwnerToken) + + require.ErrorIs(t, cache.ReleaseGroupLifecycleLease(ctx, first), service.ErrSchedulerGroupLifecycleLeaseLost) + _, acquired, err = cache.TryAcquireGroupLifecycleLease(ctx, groupID, time.Minute) + require.NoError(t, err) + require.False(t, acquired, "a stale release must not delete the successor lease") + + require.NoError(t, cache.ReleaseGroupLifecycleLease(ctx, second)) + require.ErrorIs(t, cache.ReleaseGroupLifecycleLease(ctx, second), service.ErrSchedulerGroupLifecycleLeaseLost) + third, acquired, err := cache.TryAcquireGroupLifecycleLease(ctx, groupID, time.Minute) + require.NoError(t, err) + require.True(t, acquired) + require.True(t, third.ValidFor(groupID)) + require.NoError(t, cache.ReleaseGroupLifecycleLease(ctx, third)) +} diff --git a/backend/internal/repository/scheduler_cache_unit_test.go b/backend/internal/repository/scheduler_cache_unit_test.go index bcba575145..f933603c47 100644 --- a/backend/internal/repository/scheduler_cache_unit_test.go +++ b/backend/internal/repository/scheduler_cache_unit_test.go @@ -4,6 +4,8 @@ package repository import ( "context" + "encoding/hex" + "strings" "testing" "time" @@ -438,3 +440,174 @@ func TestSchedulerCacheReopenExpiresPreviousActiveSnapshot(t *testing.T) { require.Zero(t, exists) require.NoError(t, cache.SetSnapshot(ctx, bucket, newToken, []service.Account{account})) } + +func TestSchedulerCacheGroupLifecycleLeaseConcurrentAcquireSingleOwner(t *testing.T) { + ctx := context.Background() + cache := newSchedulerCacheUnit(t) + const groupID int64 = 71 + + type result struct { + lease service.SchedulerGroupLifecycleLease + acquired bool + err error + } + start := make(chan struct{}) + results := make(chan result, 32) + for range 32 { + go func() { + <-start + lease, acquired, err := cache.TryAcquireGroupLifecycleLease(ctx, groupID, time.Minute) + results <- result{lease: lease, acquired: acquired, err: err} + }() + } + close(start) + + var owner service.SchedulerGroupLifecycleLease + acquiredCount := 0 + for range 32 { + got := <-results + require.NoError(t, got.err) + if got.acquired { + acquiredCount++ + owner = got.lease + require.True(t, got.lease.ValidFor(groupID)) + } else { + require.Equal(t, service.SchedulerGroupLifecycleLease{}, got.lease) + } + } + require.Equal(t, 1, acquiredCount) + require.Len(t, owner.OwnerToken, schedulerGroupLifecycleOwnerTokenBytes*2) + require.Equal(t, strings.ToLower(owner.OwnerToken), owner.OwnerToken) + decodedOwner, err := hex.DecodeString(owner.OwnerToken) + require.NoError(t, err) + require.Len(t, decodedOwner, schedulerGroupLifecycleOwnerTokenBytes) + + require.NoError(t, cache.ReleaseGroupLifecycleLease(ctx, owner)) + next, acquired, err := cache.TryAcquireGroupLifecycleLease(ctx, groupID, time.Minute) + require.NoError(t, err) + require.True(t, acquired) + require.True(t, next.ValidFor(groupID)) + require.NotEqual(t, owner.OwnerToken, next.OwnerToken) + require.NoError(t, cache.ReleaseGroupLifecycleLease(ctx, next)) +} + +func TestSchedulerCacheGroupLifecycleLeaseStaleReleaseCannotDeleteSuccessor(t *testing.T) { + ctx := context.Background() + cache, mr := newSchedulerCacheUnitWithRedis(t) + const groupID int64 = 72 + const ttl = time.Minute + + first, acquired, err := cache.TryAcquireGroupLifecycleLease(ctx, groupID, ttl) + require.NoError(t, err) + require.True(t, acquired) + + mr.FastForward(ttl + time.Second) + second, acquired, err := cache.TryAcquireGroupLifecycleLease(ctx, groupID, ttl) + require.NoError(t, err) + require.True(t, acquired) + require.NotEqual(t, first.OwnerToken, second.OwnerToken) + + require.ErrorIs(t, cache.ReleaseGroupLifecycleLease(ctx, first), service.ErrSchedulerGroupLifecycleLeaseLost) + owner, err := cache.rdb.Get(ctx, schedulerGroupLifecycleLockKey(groupID)).Result() + require.NoError(t, err) + require.Equal(t, second.OwnerToken, owner) + + _, acquired, err = cache.TryAcquireGroupLifecycleLease(ctx, groupID, ttl) + require.NoError(t, err) + require.False(t, acquired) + require.NoError(t, cache.ReleaseGroupLifecycleLease(ctx, second)) + require.ErrorIs(t, cache.ReleaseGroupLifecycleLease(ctx, second), service.ErrSchedulerGroupLifecycleLeaseLost) +} + +func TestSchedulerCacheGroupLifecycleLeaseExpiredReleaseIsLost(t *testing.T) { + ctx := context.Background() + cache, mr := newSchedulerCacheUnitWithRedis(t) + const groupID int64 = 73 + const ttl = time.Minute + + lease, acquired, err := cache.TryAcquireGroupLifecycleLease(ctx, groupID, ttl) + require.NoError(t, err) + require.True(t, acquired) + mr.FastForward(ttl + time.Second) + + require.ErrorIs(t, cache.ReleaseGroupLifecycleLease(ctx, lease), service.ErrSchedulerGroupLifecycleLeaseLost) +} + +func TestSchedulerCacheGroupLifecycleLeaseWrongOwnerAndCrossGroupAreLost(t *testing.T) { + ctx := context.Background() + cache := newSchedulerCacheUnit(t) + const firstGroupID int64 = 74 + const secondGroupID int64 = 75 + + first, acquired, err := cache.TryAcquireGroupLifecycleLease(ctx, firstGroupID, time.Minute) + require.NoError(t, err) + require.True(t, acquired) + second, acquired, err := cache.TryAcquireGroupLifecycleLease(ctx, secondGroupID, time.Minute) + require.NoError(t, err) + require.True(t, acquired, "different groups must acquire independently") + require.NotEqual(t, first.OwnerToken, second.OwnerToken) + + wrongOwner := first + wrongOwner.OwnerToken = strings.Repeat("0", schedulerGroupLifecycleOwnerTokenBytes*2) + if wrongOwner.OwnerToken == first.OwnerToken { + wrongOwner.OwnerToken = strings.Repeat("1", schedulerGroupLifecycleOwnerTokenBytes*2) + } + require.ErrorIs(t, cache.ReleaseGroupLifecycleLease(ctx, wrongOwner), service.ErrSchedulerGroupLifecycleLeaseLost) + + crossGroup := first + crossGroup.GroupID = secondGroupID + require.ErrorIs(t, cache.ReleaseGroupLifecycleLease(ctx, crossGroup), service.ErrSchedulerGroupLifecycleLeaseLost) + + firstOwner, err := cache.rdb.Get(ctx, schedulerGroupLifecycleLockKey(firstGroupID)).Result() + require.NoError(t, err) + require.Equal(t, first.OwnerToken, firstOwner) + secondOwner, err := cache.rdb.Get(ctx, schedulerGroupLifecycleLockKey(secondGroupID)).Result() + require.NoError(t, err) + require.Equal(t, second.OwnerToken, secondOwner) + + require.NoError(t, cache.ReleaseGroupLifecycleLease(ctx, first)) + require.NoError(t, cache.ReleaseGroupLifecycleLease(ctx, second)) +} + +func TestSchedulerCacheGroupLifecycleLeaseCanceledContextFailsClosed(t *testing.T) { + cache := newSchedulerCacheUnit(t) + canceledCtx, cancel := context.WithCancel(context.Background()) + cancel() + + lease, acquired, err := cache.TryAcquireGroupLifecycleLease(canceledCtx, 76, time.Minute) + require.ErrorIs(t, err, context.Canceled) + require.False(t, acquired) + require.Equal(t, service.SchedulerGroupLifecycleLease{}, lease) + + ctx := context.Background() + lease, acquired, err = cache.TryAcquireGroupLifecycleLease(ctx, 76, time.Minute) + require.NoError(t, err) + require.True(t, acquired) + require.ErrorIs(t, cache.ReleaseGroupLifecycleLease(canceledCtx, lease), context.Canceled) + owner, err := cache.rdb.Get(ctx, schedulerGroupLifecycleLockKey(lease.GroupID)).Result() + require.NoError(t, err) + require.Equal(t, lease.OwnerToken, owner) + require.NoError(t, cache.ReleaseGroupLifecycleLease(ctx, lease)) +} + +func TestSchedulerCacheGroupLifecycleLeaseRejectsInvalidInput(t *testing.T) { + ctx := context.Background() + cache := newSchedulerCacheUnit(t) + + lease, acquired, err := cache.TryAcquireGroupLifecycleLease(ctx, 0, time.Minute) + require.ErrorIs(t, err, service.ErrSchedulerGroupLifecycleLeaseInvalid) + require.False(t, acquired) + require.Equal(t, service.SchedulerGroupLifecycleLease{}, lease) + + lease, acquired, err = cache.TryAcquireGroupLifecycleLease(ctx, 73, 0) + require.ErrorIs(t, err, service.ErrSchedulerGroupLifecycleLeaseInvalid) + require.False(t, acquired) + require.Equal(t, service.SchedulerGroupLifecycleLease{}, lease) + + canceledCtx, cancel := context.WithCancel(ctx) + cancel() + require.ErrorIs(t, cache.ReleaseGroupLifecycleLease(canceledCtx, service.SchedulerGroupLifecycleLease{}), service.ErrSchedulerGroupLifecycleLeaseInvalid) + keys, err := cache.rdb.DBSize(ctx).Result() + require.NoError(t, err) + require.Zero(t, keys) +} diff --git a/backend/internal/service/scheduler_cache.go b/backend/internal/service/scheduler_cache.go index b117c6194f..d507d6ce2a 100644 --- a/backend/internal/service/scheduler_cache.go +++ b/backend/internal/service/scheduler_cache.go @@ -16,8 +16,10 @@ const ( ) var ( - ErrSchedulerBucketRetired = errors.New("scheduler bucket retired") - ErrSchedulerBucketWriteFenced = errors.New("scheduler bucket write fenced") + ErrSchedulerBucketRetired = errors.New("scheduler bucket retired") + ErrSchedulerBucketWriteFenced = errors.New("scheduler bucket write fenced") + ErrSchedulerGroupLifecycleLeaseInvalid = errors.New("scheduler group lifecycle lease invalid") + ErrSchedulerGroupLifecycleLeaseLost = errors.New("scheduler group lifecycle lease lost") ) // SchedulerBucketWriteToken fences a snapshot writer to one bucket epoch. @@ -31,6 +33,17 @@ func (t SchedulerBucketWriteToken) ValidFor(bucket SchedulerBucket) bool { return t.Epoch > 0 && t.Bucket == bucket } +// SchedulerGroupLifecycleLease identifies one owner of a group's short-lived +// retirement/reopen critical section. +type SchedulerGroupLifecycleLease struct { + GroupID int64 + OwnerToken string +} + +func (l SchedulerGroupLifecycleLease) ValidFor(groupID int64) bool { + return groupID > 0 && l.GroupID == groupID && l.OwnerToken != "" +} + type SchedulerBucket struct { GroupID int64 Platform string @@ -76,9 +89,17 @@ type SchedulerCache interface { // ReopenBucket is the only operation allowed to clear a tombstone. It returns // the retirement generation established by RetireBucket; repeated calls for // the same generation are idempotent. Callers must serialize a fresh authority - // check through ReopenBucket with RetireBucket under the same bucket lifecycle - // lock; ordinary rebuild paths never call ReopenBucket. + // check through ReopenBucket with RetireBucket under the same group lifecycle + // lease; ordinary rebuild paths never call ReopenBucket. ReopenBucket(ctx context.Context, bucket SchedulerBucket) (SchedulerBucketWriteToken, error) + // TryAcquireGroupLifecycleLease serializes authoritative retirement/reopen + // decisions for one non-zero group across instances. + TryAcquireGroupLifecycleLease(ctx context.Context, groupID int64, ttl time.Duration) (SchedulerGroupLifecycleLease, bool, error) + // ReleaseGroupLifecycleLease releases the lease only if its owner token still + // matches, so an expired holder cannot delete a successor's lease. Missing, + // expired, mismatched, and already released leases return + // ErrSchedulerGroupLifecycleLeaseLost. + ReleaseGroupLifecycleLease(ctx context.Context, lease SchedulerGroupLifecycleLease) error // GetAccount 获取单账号快照。 GetAccount(ctx context.Context, accountID int64) (*Account, error) // SetAccount 写入单账号快照(包含不可调度状态)。 diff --git a/backend/internal/service/scheduler_snapshot_hydration_test.go b/backend/internal/service/scheduler_snapshot_hydration_test.go index 5238529c63..e87aac180d 100644 --- a/backend/internal/service/scheduler_snapshot_hydration_test.go +++ b/backend/internal/service/scheduler_snapshot_hydration_test.go @@ -35,6 +35,14 @@ func (c *snapshotHydrationCache) ReopenBucket(ctx context.Context, bucket Schedu return SchedulerBucketWriteToken{Bucket: bucket, Epoch: 1}, nil } +func (c *snapshotHydrationCache) TryAcquireGroupLifecycleLease(context.Context, int64, time.Duration) (SchedulerGroupLifecycleLease, bool, error) { + return SchedulerGroupLifecycleLease{}, false, nil +} + +func (c *snapshotHydrationCache) ReleaseGroupLifecycleLease(context.Context, SchedulerGroupLifecycleLease) error { + return nil +} + func (c *snapshotHydrationCache) GetAccount(ctx context.Context, accountID int64) (*Account, error) { if c.accounts == nil { return nil, nil diff --git a/backend/internal/service/scheduler_snapshot_outbox_cleanup_test.go b/backend/internal/service/scheduler_snapshot_outbox_cleanup_test.go index acbdcdb259..9a7d360e2e 100644 --- a/backend/internal/service/scheduler_snapshot_outbox_cleanup_test.go +++ b/backend/internal/service/scheduler_snapshot_outbox_cleanup_test.go @@ -37,6 +37,14 @@ func (c *outboxCleanupCache) ReopenBucket(ctx context.Context, bucket SchedulerB return SchedulerBucketWriteToken{Bucket: bucket, Epoch: 1}, nil } +func (c *outboxCleanupCache) TryAcquireGroupLifecycleLease(context.Context, int64, time.Duration) (SchedulerGroupLifecycleLease, bool, error) { + return SchedulerGroupLifecycleLease{}, false, nil +} + +func (c *outboxCleanupCache) ReleaseGroupLifecycleLease(context.Context, SchedulerGroupLifecycleLease) error { + return nil +} + func (c *outboxCleanupCache) GetAccount(ctx context.Context, accountID int64) (*Account, error) { return nil, nil }