Files
pangolin/server/internal/provision/providers/registry.go
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

122 lines
3.5 KiB
Go

package providers
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"os"
"time"
"github.com/wangjia/pangolin/server/internal/provision"
)
// Factory builds a CloudAdapter, reading credentials from independent secrets
// (env). It returns provision.ErrNoCredentials when the required secret is unset.
type Factory func() (provision.CloudAdapter, error)
// builtins maps api_kind → Factory. Extend here when onboarding a vendor.
var builtins = map[string]Factory{
"vultr": newVultrFromEnv,
"hetzner": newHetznerFromEnv,
}
// Registry resolves provision.Provider rows to live adapters. It satisfies
// provision.AdapterFactory and caches one adapter per api_kind.
type Registry struct {
factories map[string]Factory
cache map[string]provision.CloudAdapter
}
// NewRegistry returns a Registry backed by the built-in vendor factories.
func NewRegistry() *Registry {
fs := make(map[string]Factory, len(builtins))
for k, v := range builtins {
fs[k] = v
}
return &Registry{factories: fs, cache: map[string]provision.CloudAdapter{}}
}
// Register adds or overrides a factory for api_kind (used in tests / extension).
func (r *Registry) Register(apiKind string, f Factory) { r.factories[apiKind] = f }
// For resolves the adapter for a provider row.
func (r *Registry) For(p *provision.Provider) (provision.CloudAdapter, error) {
if a, ok := r.cache[p.APIKind]; ok {
return a, nil
}
f, ok := r.factories[p.APIKind]
if !ok {
return nil, fmt.Errorf("providers: no adapter registered for api_kind %q", p.APIKind)
}
a, err := f()
if err != nil {
return nil, err
}
r.cache[p.APIKind] = a
return a, nil
}
var _ provision.AdapterFactory = (*Registry)(nil)
// --- shared HTTP helper ---
// httpClient is a small JSON REST helper shared by the adapters.
type httpClient struct {
base string
bearer string
hc *http.Client
}
func newHTTPClient(base, bearer string) *httpClient {
return &httpClient{base: base, bearer: bearer, hc: &http.Client{Timeout: 30 * time.Second}}
}
// do issues an authenticated JSON request and decodes the response into out
// (out may be nil). It returns an error on any non-2xx status.
func (c *httpClient) do(ctx context.Context, method, path string, body, out any) error {
var rdr io.Reader
if body != nil {
buf, err := json.Marshal(body)
if err != nil {
return fmt.Errorf("providers: marshal body: %w", err)
}
rdr = bytes.NewReader(buf)
}
req, err := http.NewRequestWithContext(ctx, method, c.base+path, rdr)
if err != nil {
return fmt.Errorf("providers: build request: %w", err)
}
req.Header.Set("Authorization", "Bearer "+c.bearer)
if body != nil {
req.Header.Set("Content-Type", "application/json")
}
resp, err := c.hc.Do(req)
if err != nil {
return fmt.Errorf("providers: %s %s: %w", method, path, err)
}
defer resp.Body.Close()
data, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
// Never echo credentials; only status + vendor body (no auth header).
return fmt.Errorf("providers: %s %s: status %d: %s", method, path, resp.StatusCode, string(data))
}
if out != nil && len(data) > 0 {
if err := json.Unmarshal(data, out); err != nil {
return fmt.Errorf("providers: decode response: %w", err)
}
}
return nil
}
// secret reads an env-injected credential, returning ErrNoCredentials if unset.
func secret(env string) (string, error) {
v := os.Getenv(env)
if v == "" {
return "", fmt.Errorf("%w (env %s)", provision.ErrNoCredentials, env)
}
return v, nil
}