Files
wangjia 3d5bac66b4 feat(provision): 弹性节点基建 Terraform + 一键更换 (tsk_6u0FxmbC7Yeq)
IaC 面 (infra/) 与控制面 (server/internal/provision) 双产出,落地 doc/04 §4
「节点是牲口」弹性拓扑与 make-before-break 一键更换。

server/internal/provision:
- CloudAdapter 适配层 + Registry(首发 vultr 消耗品池 / hetzner 精品池各一);
  厂商凭证仅从 PROVISION_<VENDOR>_* env 注入,不入库不入 git。
- ProvisionService:CreateNode(幂等键重放不重复开机)、DestroyNode(幂等)、
  RotateIP(换 IP 不换机 + version bump)、ListProviders。
- Replace 一键更换:先建后拆,新机 up 先于旧机 draining(容量不下降),
  replacement_uuid 幂等键 + replacements 表分步记录,崩溃可续跑不重复。
- RotatePool:池内滚动轮换,并发度 1–2。
- cmd/nodectl CLI:create/destroy/rotate-ip/replace/rotate-pool/providers。
- 单测(mock 厂商 API + 内存 Store):幂等重放、make-before-break 时序断言、
  开机失败/探活超时→destroyed+失败计数+告警钩子、崩溃续跑、RotatePool。

infra/:
- terraform/:探针机 + 控制面基线模块化(probe / control-plane)+ README,
  低频基线进 state,节点不进 Terraform。
- cloud-init/node.yaml.tmpl:节点引导模板(注入一次性 bootstrap token,task #5)。
- identity-isolation.md:身份隔离登记表(doc/06 §2 红线,无任何凭证)。

migrations/000008:nodes 增 provider_instance_id/elastic_ip_id、node_events
增 ip_rotated、provision_idempotency / replacements 表(附加式,不动现网)。

红线:仅面向新厂商池,绝不纳管现网生产 EC2(deploy/ marzban)。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-13 14:23:39 +08:00

313 lines
9.2 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package provision
import (
"context"
"fmt"
"sync"
"time"
)
// stepRank orders ReplaceStep so resume logic can compare progress.
var stepRank = map[ReplaceStep]int{
StepOpenNew: 0,
StepProbing: 1,
StepNewUp: 2,
StepDraining: 3,
StepDestroyOld: 4,
StepDone: 5,
}
// ReplaceResult reports the outcome of one Replace.
type ReplaceResult struct {
ReplacementUUID string
OldNodeID int64
NewNodeID int64
NewNode *Node
}
// Replace performs a make-before-break one-click replacement of nodeID
// (doc/04 §4.1). A fresh node is brought UP before the old one is drained and
// destroyed, so capacity never dips.
//
// replacementUUID is the idempotency key for the whole orchestration. Pass ""
// to start a new replacement; pass an existing UUID to RESUME after a crash —
// progress is persisted in the replacements table and CreateNode's own
// idempotency guarantees no duplicate boot (validation: "崩溃重启续跑不重复").
func (s *Service) Replace(ctx context.Context, nodeID int64, replacementUUID string) (*ReplaceResult, error) {
old, err := s.store.GetNode(ctx, nodeID)
if err != nil {
return nil, err
}
if old == nil {
return nil, fmt.Errorf("provision: node %d not found", nodeID)
}
pool := poolForTier(old.Tier)
// Load or create the orchestration record.
var rec *Replacement
if replacementUUID != "" {
rec, err = s.store.GetReplacement(ctx, replacementUUID)
if err != nil {
return nil, err
}
}
if rec == nil {
if replacementUUID == "" {
if replacementUUID, err = newUUID(); err != nil {
return nil, err
}
}
rec = &Replacement{
UUID: replacementUUID,
OldNodeID: nodeID,
Pool: pool,
Status: ReplaceRunning,
Step: StepOpenNew,
}
if err := s.store.CreateReplacement(ctx, rec); err != nil {
return nil, err
}
}
if rec.Status == ReplaceDone {
// Already finished; return the recorded result.
newNode, _ := s.store.GetNode(ctx, rec.NewNodeID)
return &ReplaceResult{ReplacementUUID: rec.UUID, OldNodeID: rec.OldNodeID, NewNodeID: rec.NewNodeID, NewNode: newNode}, nil
}
return s.runReplace(ctx, rec, old, pool)
}
// runReplace drives the replacement state machine forward from rec.Step,
// persisting after each transition so a crash resumes cleanly.
func (s *Service) runReplace(ctx context.Context, rec *Replacement, old *Node, pool Pool) (*ReplaceResult, error) {
atLeast := func(step ReplaceStep) bool { return stepRank[rec.Step] >= stepRank[step] }
advance := func(step ReplaceStep) error {
rec.Step = step
return s.store.UpdateReplacement(ctx, rec)
}
var newNode *Node
// Step 1: open the new node (same region, pool may switch vendor).
if !atLeast(StepProbing) {
spec, err := s.replacementSpec(ctx, old, pool)
if err != nil {
return s.failReplace(ctx, rec, old, pool, err)
}
nn, err := s.CreateNode(ctx, spec, rec.UUID+":create")
if err != nil {
return s.failReplace(ctx, rec, old, pool, fmt.Errorf("open new node: %w", err))
}
newNode = nn
rec.NewNodeID = nn.ID
if err := advance(StepProbing); err != nil {
return nil, err
}
}
if newNode == nil && rec.NewNodeID != 0 {
if newNode, _ = s.store.GetNode(ctx, rec.NewNodeID); newNode == nil {
return s.failReplace(ctx, rec, old, pool, fmt.Errorf("new node %d vanished", rec.NewNodeID))
}
}
// Step 2: wait for agent self-register + simplified probing.
if !atLeast(StepNewUp) {
if s.prober != nil {
pctx, cancel := context.WithTimeout(ctx, s.probeTimeout)
err := s.prober.WaitReady(pctx, newNode)
cancel()
if err != nil {
_ = s.store.WriteNodeEvent(ctx, newNode.ID, EventProbeFail, jsonDetail(map[string]any{"error": err.Error()}))
// Probe timeout: destroy the failed new node, count + alert.
_ = s.DestroyNode(ctx, newNode.ID)
count := s.bumpFail(pool)
s.alert.Fire(ctx, Alert{Kind: AlertProbeTimeout, NodeUUID: newNode.UUID, Pool: pool, Message: err.Error(), FailCount: count})
rec.Status = ReplaceFailed
_ = s.store.UpdateReplacement(ctx, rec)
return nil, fmt.Errorf("provision: replace probe failed: %w", err)
}
}
_ = s.store.WriteNodeEvent(ctx, newNode.ID, EventProbePass, "")
if err := advance(StepNewUp); err != nil {
return nil, err
}
}
// Step 3: promote the new node UP and publish (capacity is now restored
// BEFORE the old node leaves the directory — make-before-break).
if !atLeast(StepDraining) {
if err := s.store.UpdateNodeStatus(ctx, newNode.ID, StatusUp); err != nil {
return nil, err
}
_ = s.store.WriteNodeEvent(ctx, newNode.ID, EventMarkedUp, "")
if _, err := s.store.BumpDirectoryVersion(ctx); err != nil {
return nil, err
}
s.resetFail(pool)
if err := advance(StepDraining); err != nil {
return nil, err
}
}
// Step 4: drain the old node. Free/consumable nodes hard-cut.
if !atLeast(StepDestroyOld) {
if err := s.store.UpdateNodeStatus(ctx, old.ID, StatusDraining); err != nil {
return nil, err
}
_ = s.store.WriteNodeEvent(ctx, old.ID, EventDraining, "")
if _, err := s.store.BumpDirectoryVersion(ctx); err != nil {
return nil, err
}
if pool != PoolConsumable {
if err := s.clock.Sleep(ctx, s.drainTimeout); err != nil {
return nil, err
}
}
if err := advance(StepDestroyOld); err != nil {
return nil, err
}
}
// Step 5: destroy the old node + release IP.
if !atLeast(StepDone) {
if err := s.DestroyNode(ctx, old.ID); err != nil {
return s.failReplace(ctx, rec, old, pool, fmt.Errorf("destroy old node: %w", err))
}
_ = s.store.WriteNodeEvent(ctx, old.ID, EventReplaced, jsonDetail(map[string]any{
"replacement_uuid": rec.UUID,
"new_node_id": rec.NewNodeID,
}))
if err := advance(StepDone); err != nil {
return nil, err
}
}
rec.Status = ReplaceDone
if err := s.store.UpdateReplacement(ctx, rec); err != nil {
return nil, err
}
_ = s.store.WriteAuditLog(ctx, "provision", "replace", "node:"+old.UUID,
jsonDetail(map[string]any{"replacement_uuid": rec.UUID, "new_node_id": rec.NewNodeID}))
if newNode == nil && rec.NewNodeID != 0 {
newNode, _ = s.store.GetNode(ctx, rec.NewNodeID)
}
return &ReplaceResult{
ReplacementUUID: rec.UUID,
OldNodeID: rec.OldNodeID,
NewNodeID: rec.NewNodeID,
NewNode: newNode,
}, nil
}
// failReplace records a fatal replacement failure (alert + audit).
func (s *Service) failReplace(ctx context.Context, rec *Replacement, old *Node, pool Pool, cause error) (*ReplaceResult, error) {
rec.Status = ReplaceFailed
_ = s.store.UpdateReplacement(ctx, rec)
count := s.bumpFail(pool)
_ = s.store.WriteAuditLog(ctx, "provision", "replace_failed", "node:"+old.UUID,
jsonDetail(map[string]any{"replacement_uuid": rec.UUID, "error": cause.Error(), "fail_count": count}))
s.alert.Fire(ctx, Alert{Kind: AlertReplaceFail, NodeUUID: old.UUID, Pool: pool, Message: cause.Error(), FailCount: count})
return nil, fmt.Errorf("provision: replace failed: %w", cause)
}
// replacementSpec derives the spec for the new node from the old one, keeping
// the same region/tier/role but allowing a different vendor in the same pool
// (doc/04 §4.1: "同 region · 厂商池内可换家").
func (s *Service) replacementSpec(ctx context.Context, old *Node, pool Pool) (NodeSpec, error) {
providerID := old.ProviderID
providers, err := s.store.ListProviders(ctx, pool)
if err != nil {
return NodeSpec{}, err
}
// Prefer a different enabled provider in the pool to spread exposure.
for _, p := range providers {
if p.ID != old.ProviderID && providerSupportsRegion(p, old.Region) {
providerID = p.ID
break
}
}
return NodeSpec{
Region: old.Region,
Role: old.Role,
Tier: old.Tier,
ProviderID: providerID,
NameZH: old.NameZH,
NameEn: old.NameEn,
RealityPBK: old.RealityPBK,
RealitySNI: old.RealitySNI,
HY2Port: old.HY2Port,
Weight: old.Weight,
Tags: old.Tags,
}, nil
}
func providerSupportsRegion(p *Provider, region string) bool {
if len(p.Regions) == 0 {
return true // unconstrained
}
for _, r := range p.Regions {
if r == region {
return true
}
}
return false
}
// RotatePool rolls Replace across every up node in a pool with bounded
// concurrency (doc/04 §4: 并发度 12). Used for routine rotation or large-scale
// event rebuilds.
func (s *Service) RotatePool(ctx context.Context, pool Pool, concurrency int) ([]ReplaceResult, error) {
if concurrency < 1 {
concurrency = 1
}
if concurrency > 2 {
concurrency = 2
}
nodes, err := s.store.ListNodesByPool(ctx, pool, StatusUp)
if err != nil {
return nil, err
}
var (
mu sync.Mutex
results []ReplaceResult
firstEr error
wg sync.WaitGroup
)
sem := make(chan struct{}, concurrency)
for _, n := range nodes {
n := n
if firstErrSet(&mu, &firstEr) {
break
}
wg.Add(1)
sem <- struct{}{}
go func() {
defer wg.Done()
defer func() { <-sem }()
res, err := s.Replace(ctx, n.ID, "")
mu.Lock()
defer mu.Unlock()
if err != nil {
if firstEr == nil {
firstEr = err
}
return
}
results = append(results, *res)
}()
}
wg.Wait()
return results, firstEr
}
func firstErrSet(mu *sync.Mutex, e *error) bool {
mu.Lock()
defer mu.Unlock()
return *e != nil
}
// drainDeadline is exposed for tests/monitoring of the configured drain window.
func (s *Service) drainDeadline(start time.Time) time.Time { return start.Add(s.drainTimeout) }