Files
pangolin/server/internal/scheduler/orchestrate/grayscale.go
wangjia 7627c8a195 feat(scheduler): 自动更换编排 + 养机灰度 [tsk_mxEFSvKHX64G]
实现 15E 任务:Replacer 状态机(pending→creating→probing→activating→draining_old→done/failed)
+ Grayscale 权重梯度(10→25→50→75→100,每步 6h)。

关键设计:
- Redis JSON 记录(sched:replace:{uuid})+ SetNX 崩溃安全,防止重复 CreateNode
- 幂等 key:attempt 0 = uuid,retry N = uuid:retry:N
- probePass:国内 ≥2/3 ISP 通过 + 海外 L1 全通,连续 2 Tick 才晋级
- 探活超时 15min → failProbeAttempt → 轮换 provider;3 次失败 → failed + NotifyFault
- StubBreaker(15F 未就绪);NewReplacer(Config{}) 依赖注入
- 6 个集成测试全部通过(happy path / max-retry / crash recovery / queue idempotency / grayscale)

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-15 21:26:07 +08:00

184 lines
6.4 KiB
Go

package orchestrate
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"strings"
"time"
"github.com/redis/go-redis/v9"
)
// ─────────────────────────────────────────────────────────────────────────────
// Redis key namespace for grayscale records
// ─────────────────────────────────────────────────────────────────────────────
const (
// grayKeyPrefix is the namespace for per-node grayscale warm-up records.
// Full key: sched:gray:{nodeID}
grayKeyPrefix = "sched:gray:"
)
// ─────────────────────────────────────────────────────────────────────────────
// GrayscaleRecord
// ─────────────────────────────────────────────────────────────────────────────
// GrayscaleRecord is the Redis-persisted state of a node in weight warm-up.
// Key: sched:gray:{nodeID} Value: JSON-encoded GrayscaleRecord.
type GrayscaleRecord struct {
NodeID string `json:"nodeId"`
CurrentWeight int `json:"currentWeight"`
StartedAt time.Time `json:"startedAt"`
LastAdvancedAt time.Time `json:"lastAdvancedAt"`
}
// startGrayscale creates a new GrayscaleRecord for nodeID at weight 10 (the
// initial weight). It is called from stepActivating immediately after the node
// is transitioned to "up".
func startGrayscale(ctx context.Context, rdb *redis.Client, clock Clock, nodeID string) error {
now := clock.Now()
rec := GrayscaleRecord{
NodeID: nodeID,
CurrentWeight: GrayscaleWeights[0], // 10
StartedAt: now,
LastAdvancedAt: now,
}
data, err := json.Marshal(rec)
if err != nil {
return err
}
return rdb.Set(ctx, grayKeyPrefix+nodeID, data, 0).Err()
}
// ─────────────────────────────────────────────────────────────────────────────
// Grayscale
// ─────────────────────────────────────────────────────────────────────────────
// Grayscale advances the weight ramp for all nodes currently in warm-up.
// It is called by the CapacityLoop (task 15H); the implementation lives here
// (task 15E) while wiring is done in 15H.
//
// Weight ladder: 10 → 25 → 50 → 75 → 100, one step every 6 h.
// When weight 100 is reached, the grayscale record is deleted.
type Grayscale struct {
rdb *redis.Client
lc LifecycleService
clock Clock
}
// NewGrayscale creates a Grayscale with the given dependencies.
// If clock is nil, RealClock is used.
func NewGrayscale(rdb *redis.Client, lc LifecycleService, clock Clock) *Grayscale {
if clock == nil {
clock = RealClock{}
}
return &Grayscale{rdb: rdb, lc: lc, clock: clock}
}
// Advance scans all sched:gray:* keys and advances any node whose
// LastAdvancedAt is at least GrayscaleInterval (6 h) in the past.
func (g *Grayscale) Advance(ctx context.Context) error {
keys, err := g.scanGrayKeys(ctx)
if err != nil {
return fmt.Errorf("grayscale: scan keys: %w", err)
}
for _, key := range keys {
if advErr := g.advanceNode(ctx, key); advErr != nil {
slog.Error("grayscale: advance node", "key", key, "error", advErr)
// continue; one error must not halt other nodes
}
}
return nil
}
// advanceNode loads the record at key, advances the weight if the interval has
// elapsed, calls SetWeight, writes an audit log, and deletes the record when
// the final weight (100) is reached.
func (g *Grayscale) advanceNode(ctx context.Context, key string) error {
val, err := g.rdb.Get(ctx, key).Result()
if err == redis.Nil {
return nil // already deleted
}
if err != nil {
return fmt.Errorf("grayscale: get %s: %w", key, err)
}
var rec GrayscaleRecord
if err := json.Unmarshal([]byte(val), &rec); err != nil {
return fmt.Errorf("grayscale: unmarshal %s: %w", key, err)
}
now := g.clock.Now()
if now.Sub(rec.LastAdvancedAt) < GrayscaleInterval {
return nil // interval not elapsed yet
}
// Find the next weight in the ladder.
nextWeight, ok := nextGrayscaleWeight(rec.CurrentWeight)
if !ok {
// Already at maximum; remove the record.
_ = g.rdb.Del(ctx, key).Err()
return nil
}
// Apply the new weight.
if err := g.lc.SetWeight(ctx, rec.NodeID, nextWeight); err != nil {
return fmt.Errorf("grayscale: set weight %d for %s: %w", nextWeight, rec.NodeID, err)
}
// Audit trail.
meta := fmt.Sprintf(`{"from":%d,"to":%d,"node":%q}`, rec.CurrentWeight, nextWeight, rec.NodeID)
_ = g.lc.WriteAuditLog(ctx, "grayscale", "weight_advanced", "node:"+rec.NodeID, meta)
slog.Info("grayscale: weight advanced",
"node", rec.NodeID,
"from", rec.CurrentWeight,
"to", nextWeight,
)
if nextWeight >= GrayscaleWeights[len(GrayscaleWeights)-1] {
// Final weight reached: remove the grayscale record.
_ = g.rdb.Del(ctx, key).Err()
return nil
}
// Persist updated record.
rec.CurrentWeight = nextWeight
rec.LastAdvancedAt = now
data, err := json.Marshal(rec)
if err != nil {
return err
}
return g.rdb.Set(ctx, key, data, 0).Err()
}
// nextGrayscaleWeight returns the next weight after current in GrayscaleWeights.
// Returns (0, false) if current is already at or beyond the final weight.
func nextGrayscaleWeight(current int) (int, bool) {
for i, w := range GrayscaleWeights {
if w == current && i+1 < len(GrayscaleWeights) {
return GrayscaleWeights[i+1], true
}
}
return 0, false
}
// scanGrayKeys returns all Redis keys matching sched:gray:*.
func (g *Grayscale) scanGrayKeys(ctx context.Context) ([]string, error) {
var keys []string
iter := g.rdb.Scan(ctx, 0, grayKeyPrefix+"*", 0).Iterator()
for iter.Next(ctx) {
k := iter.Val()
// Skip any keys that don't have the expected format (e.g. the index key).
if strings.HasPrefix(k, grayKeyPrefix) {
keys = append(keys, k)
}
}
if err := iter.Err(); err != nil {
return nil, err
}
return keys, nil
}