Files
rogee b6d0af1a56
management-images / build-and-publish (push) Successful in 9m15s
docs: add one-click management deployment and image workflow
2026-09-16 17:57:05 +08:00

478 lines
20 KiB
Go

package management
import (
"context"
"database/sql"
"encoding/json"
"errors"
"net/http"
"path"
"strings"
"time"
"github.com/gofiber/fiber/v3"
)
type barrier struct {
TrunkID string
OperationID string
TargetRevision int64
Reason string
State string
}
type cellApplyRequest struct {
Mode string `json:"mode"`
CellID string `json:"cell_id"`
TrunkID string `json:"trunk_id"`
Revision int64 `json:"revision"`
ExpectedLocalRevision int64 `json:"expected_local_revision"`
Config map[string]any `json:"config"`
ConfigSHA256 string `json:"config_sha256"`
}
type cellApplyAck struct {
Mode string `json:"mode"`
CellID string `json:"cell_id"`
TrunkID string `json:"trunk_id"`
Revision int64 `json:"revision"`
ConfigSHA256 string `json:"config_sha256"`
Status string `json:"status"`
Idempotent bool `json:"idempotent"`
}
func (s *Store) GetBarrier(ctx context.Context, trunkID string) (barrier, error) {
var b barrier
err := s.withRead(ctx, func(db *sql.DB) error {
return db.QueryRowContext(ctx, "SELECT trunk_id,operation_id,target_revision,reason,state FROM admission_barriers WHERE trunk_id=?", trunkID).Scan(&b.TrunkID, &b.OperationID, &b.TargetRevision, &b.Reason, &b.State)
})
return b, err
}
func (s *Store) ResumeOrStagePublication(ctx context.Context, trunkID string, revision int64, cfg TrunkConfig, cellIDs []string, operationID, reason string) error {
digest, _, err := hashJSON(cfg)
if err != nil {
return err
}
now := utcString(time.Now())
return s.withTx(ctx, func(tx *sql.Tx) error {
var existing barrier
err := tx.QueryRowContext(ctx, "SELECT trunk_id,operation_id,target_revision,reason,state FROM admission_barriers WHERE trunk_id=?", trunkID).Scan(&existing.TrunkID, &existing.OperationID, &existing.TargetRevision, &existing.Reason, &existing.State)
hasBarrier := err == nil
if err != nil && !errors.Is(err, sql.ErrNoRows) {
return err
}
if hasBarrier && existing.State == "active" && existing.OperationID != operationID {
return newAppError(409, "PUBLISH_IN_PROGRESS", "another publication currently owns the admission barrier", map[string]any{"operation_id": existing.OperationID})
}
if hasBarrier && existing.TargetRevision != revision && existing.State != "released" && existing.OperationID != operationID {
return newAppError(409, "PUBLISH_IN_PROGRESS", "an older publication must be reconciled before this resource changes", map[string]any{"target_revision": existing.TargetRevision})
}
if _, err := tx.ExecContext(ctx, "INSERT INTO admission_barriers(trunk_id,operation_id,target_revision,reason,state,created_at,updated_at) VALUES(?,?,?,?,?,?,?) ON CONFLICT(trunk_id) DO UPDATE SET operation_id=excluded.operation_id,target_revision=excluded.target_revision,reason=excluded.reason,state='active',updated_at=excluded.updated_at", trunkID, operationID, revision, reason, "active", now, now); err != nil {
return err
}
var count int
if err := tx.QueryRowContext(ctx, "SELECT COUNT(*) FROM publications WHERE trunk_id=? AND revision=?", trunkID, revision).Scan(&count); err != nil {
return err
}
if count == 0 {
if len(cellIDs) == 0 {
return newAppError(409, "NO_CAPABLE_CELL", "no compatible Cell is available", nil)
}
for _, cellID := range cellIDs {
var localRevision int64
if err := tx.QueryRowContext(ctx, "SELECT COALESCE(MAX(local_revision),0) FROM publications WHERE cell_id=? AND status='applied'", cellID).Scan(&localRevision); err != nil {
return err
}
if _, err := tx.ExecContext(ctx, "INSERT INTO publications(trunk_id,revision,cell_id,status,error_code,target_digest,local_revision,local_digest,operation_id,updated_at) VALUES(?,?,?,?,?,?,?,?,?,?)", trunkID, revision, cellID, "pending", nil, digest, localRevision, "", operationID, now); err != nil {
return err
}
}
} else {
var mismatched int
if err := tx.QueryRowContext(ctx, "SELECT COUNT(*) FROM publications WHERE trunk_id=? AND revision=? AND target_digest<>?", trunkID, revision, digest).Scan(&mismatched); err != nil {
return err
}
if mismatched > 0 {
return newAppError(409, "PUBLICATION_DIGEST_MISMATCH", "stored publication snapshot differs from the requested revision", nil)
}
if len(cellIDs) > 0 {
var expectedCount int
if err := tx.QueryRowContext(ctx, "SELECT COUNT(*) FROM publications WHERE trunk_id=? AND revision=? AND cell_id IN ("+placeholders(len(cellIDs))+")", append([]any{trunkID, revision}, stringsToAny(cellIDs)...)...).Scan(&expectedCount); err != nil {
return err
}
if expectedCount != len(cellIDs) {
return newAppError(409, "DEPLOYMENT_SET_CHANGED", "the frozen Cell deployment set changed; revalidate before retry", nil)
}
}
if _, err := tx.ExecContext(ctx, "UPDATE publications SET operation_id=?,updated_at=? WHERE trunk_id=? AND revision=?", operationID, now, trunkID, revision); err != nil {
return err
}
}
_, err = tx.ExecContext(ctx, "UPDATE trunk_versions SET state='publishing' WHERE trunk_id=? AND revision=?", trunkID, revision)
return err
})
}
func placeholders(n int) string {
if n <= 0 {
return "?"
}
return strings.TrimSuffix(strings.Repeat("?,", n), ",")
}
func stringsToAny(v []string) []any {
out := make([]any, len(v))
for i, x := range v {
out[i] = x
}
return out
}
func (s *Server) publishTrunk(c fiber.Ctx) error {
return s.executePublicationRequest(c, "publish")
}
func (s *Server) executePublicationRequest(c fiber.Ctx, action string) error {
id := c.Params("trunk_id")
scope := "sip.trunk.publish"
if action == "disable" {
scope = "sip.trunk.disable"
}
if action == "rollback" {
scope = "sip.trunk.rollback"
}
p, authErr := s.authorize(c, "admin", scope, "trunk", id)
if authErr != nil {
return s.fail(c, authErr)
}
expected, err := parseIfMatch(c)
if err != nil {
return s.fail(c, err)
}
requestID, err := requireHeader(c, "X-Request-ID")
if err != nil {
return s.fail(c, err)
}
var input struct {
TargetRevision int64 `json:"target_revision"`
}
if action == "rollback" {
if err := decodeJSON(c, &input); err != nil {
return s.fail(c, err)
}
if input.TargetRevision < 1 {
return s.fail(c, newAppError(422, "TARGET_REVISION_REQUIRED", "target_revision must be positive", nil))
}
}
ctx, cancel := s.context(c)
defer cancel()
payload := map[string]any{"if_match": expected}
if action == "rollback" {
payload["target_revision"] = input.TargetRevision
}
op, replay, beginErr := s.beginOperation(ctx, p, "trunk."+action, "trunk", id, requestID, payload, expected)
if replay || beginErr != nil {
return s.replayOrError(c, op, beginErr)
}
if action == "disable" {
return s.disableWithOperation(c, ctx, p, op, id, expected, requestID)
}
if action == "rollback" {
return s.rollbackWithOperation(c, ctx, p, op, id, expected, input.TargetRevision, requestID)
}
trunk, opErr := s.applyRevision(ctx, id, expected, op.OperationID, p.Principal, requestID, "publish", true)
if opErr != nil {
return s.finishError(ctx, op, opErr, c)
}
if err := s.finishSuccess(ctx, op, 200, trunk); err != nil {
return s.fail(c, asAppError(err))
}
return c.JSON(trunk)
}
func (s *Server) disableTrunk(c fiber.Ctx) error { return s.executePublicationRequest(c, "disable") }
func (s *Server) rollbackTrunk(c fiber.Ctx) error { return s.executePublicationRequest(c, "rollback") }
func (s *Server) disableWithOperation(c fiber.Ctx, ctx context.Context, p principal, op Operation, id string, expected int64, requestID string) error {
current, err := s.store.GetTrunk(ctx, id, s.cfg.PublicMode())
if errors.Is(err, sql.ErrNoRows) {
return s.finishError(ctx, op, newAppError(404, "TRUNK_NOT_FOUND", "trunk does not exist", nil), c)
}
if err != nil {
return s.finishError(ctx, op, asAppError(err), c)
}
if current.LatestRevision != expected {
return s.finishError(ctx, op, newAppError(409, "REVISION_CONFLICT", "trunk revision changed", map[string]any{"expected_revision": expected, "current_revision": current.LatestRevision}), c)
}
cfg, err := s.store.getTrunkConfig(ctx, id, current.LatestRevision)
if err != nil {
return s.finishError(ctx, op, asAppError(err), c)
}
if current.ActiveRevision > 0 && current.Status == "disabled" && !cfg.Enabled {
return s.finishOperationResponse(ctx, op, current, c)
}
if err := s.store.CreateAdmissionIntent(ctx, id, current.ActiveRevision, op.OperationID, "disable"); err != nil {
return s.finishError(ctx, op, asAppError(err), c)
}
cfg.Enabled = false
created, status, err := s.store.PutTrunk(ctx, id, cfg, expected, p.Principal, requestID)
_ = status
if err != nil {
_ = s.store.SetBarrierState(ctx, id, op.OperationID, "pending")
return s.finishError(ctx, op, asAppError(err), c)
}
trunk, opErr := s.applyRevision(ctx, id, created.LatestRevision, op.OperationID, p.Principal, requestID, "disable", false)
if opErr != nil {
return s.finishError(ctx, op, opErr, c)
}
if err := s.finishSuccess(ctx, op, 200, trunk); err != nil {
return s.fail(c, asAppError(err))
}
return c.JSON(trunk)
}
func (s *Server) rollbackWithOperation(c fiber.Ctx, ctx context.Context, p principal, op Operation, id string, expected, target int64, requestID string) error {
current, err := s.store.GetTrunk(ctx, id, s.cfg.PublicMode())
if errors.Is(err, sql.ErrNoRows) {
return s.finishError(ctx, op, newAppError(404, "TRUNK_NOT_FOUND", "trunk does not exist", nil), c)
}
if err != nil {
return s.finishError(ctx, op, asAppError(err), c)
}
if current.LatestRevision != expected {
return s.finishError(ctx, op, newAppError(409, "REVISION_CONFLICT", "trunk revision changed", map[string]any{"expected_revision": expected, "current_revision": current.LatestRevision}), c)
}
cfg, err := s.store.getTrunkConfig(ctx, id, target)
if errors.Is(err, sql.ErrNoRows) {
return s.finishError(ctx, op, newAppError(409, "VERSION_NOT_FOUND", "target revision does not exist", nil), c)
}
if err != nil {
return s.finishError(ctx, op, asAppError(err), c)
}
if err := s.store.CreateAdmissionIntent(ctx, id, current.ActiveRevision, op.OperationID, "rollback"); err != nil {
return s.finishError(ctx, op, asAppError(err), c)
}
created, _, err := s.store.PutTrunk(ctx, id, cfg, expected, p.Principal, requestID)
if err != nil {
_ = s.store.SetBarrierState(ctx, id, op.OperationID, "pending")
return s.finishError(ctx, op, asAppError(err), c)
}
trunk, opErr := s.applyRevision(ctx, id, created.LatestRevision, op.OperationID, p.Principal, requestID, "rollback", false)
if opErr != nil {
return s.finishError(ctx, op, opErr, c)
}
if err := s.finishSuccess(ctx, op, 200, trunk); err != nil {
return s.fail(c, asAppError(err))
}
return c.JSON(trunk)
}
func (s *Server) finishOperationResponse(ctx context.Context, op Operation, value any, c fiber.Ctx) error {
if err := s.finishSuccess(ctx, op, 200, value); err != nil {
return s.fail(c, asAppError(err))
}
return c.JSON(value)
}
func (s *Server) applyRevision(ctx context.Context, trunkID string, revision int64, operationID, actor, requestID, action string, publishOnlyEnabled bool) (AdminTrunk, *AppError) {
trunk, err := s.store.GetTrunk(ctx, trunkID, s.cfg.PublicMode())
if errors.Is(err, sql.ErrNoRows) {
return AdminTrunk{}, newAppError(404, "TRUNK_NOT_FOUND", "trunk does not exist", nil)
}
if err != nil {
return AdminTrunk{}, asAppError(err)
}
if trunk.LatestRevision != revision {
return AdminTrunk{}, newAppError(409, "REVISION_CONFLICT", "trunk revision changed", map[string]any{"current_revision": trunk.LatestRevision})
}
if action == "publish" && trunk.ActiveRevision == revision && trunk.ActiveStatus == "published" {
trunk.Mode = s.cfg.PublicMode()
return trunk, nil
}
cfg, err := s.store.getTrunkConfig(ctx, trunkID, revision)
if errors.Is(err, sql.ErrNoRows) {
return AdminTrunk{}, newAppError(409, "VERSION_NOT_FOUND", "trunk version does not exist", nil)
}
if err != nil {
return AdminTrunk{}, asAppError(err)
}
if publishOnlyEnabled && !cfg.Enabled {
return AdminTrunk{}, newAppError(409, "TRUNK_DISABLED", "disabled config cannot be published", nil)
}
if s.cfg.Mode == ModeReal {
complete, missing, verificationErr := s.store.RealVerificationComplete(ctx, trunkID, revision)
if verificationErr != nil {
return AdminTrunk{}, asAppError(verificationErr)
}
if !complete {
return AdminTrunk{}, newAppError(409, "REAL_VERIFICATION_INCOMPLETE", "real publication requires confirmed provider/line verification records", map[string]any{"missing_checks": missing})
}
}
if err := s.store.validateTrunkConfig(ctx, cfg, s.cfg.Mode == ModeReal); err != nil {
return AdminTrunk{}, asAppError(err)
}
pubs, err := s.store.GetPublicationRows(ctx, trunkID, &revision)
if err != nil {
return AdminTrunk{}, asAppError(err)
}
cellIDs := make([]string, 0, len(pubs))
for _, pub := range pubs {
cellIDs = append(cellIDs, pub.CellID)
}
if len(cellIDs) == 0 {
cellIDs, err = s.store.compatibleCellIDs(ctx, cfg)
if err != nil {
return AdminTrunk{}, asAppError(err)
}
}
if s.cfg.Mode == ModeReal {
if observationErr := s.realCellObservationError(ctx, trunkID, cellIDs); observationErr != nil {
return AdminTrunk{}, observationErr
}
}
if err := s.store.ResumeOrStagePublication(ctx, trunkID, revision, cfg, cellIDs, operationID, action); err != nil {
return AdminTrunk{}, asAppError(err)
}
pubs, err = s.store.GetPublicationRows(ctx, trunkID, &revision)
if err != nil {
return AdminTrunk{}, asAppError(err)
}
var failures []map[string]any
pending := false
for _, pub := range pubs {
if pub.Status == "applied" && pub.LocalRevision == revision && pub.LocalDigest == pub.TargetDigest {
continue
}
cell, err := s.store.GetCell(ctx, pub.CellID)
if err != nil {
failures = append(failures, map[string]any{"cell_id": pub.CellID, "code": "CELL_NOT_FOUND"})
_ = s.store.UpdatePublication(ctx, trunkID, revision, pub.CellID, "failed", "CELL_NOT_FOUND", 0, "", operationID)
continue
}
ack, applyErr := s.applyToCell(ctx, cell, trunkID, revision, pub.LocalRevision, cfg, pub.TargetDigest, requestID)
if applyErr != nil {
if applyErr.Code == "MOCK_INTENT_ONLY" {
pending = true
continue
}
failures = append(failures, map[string]any{"cell_id": pub.CellID, "code": applyErr.Code})
_ = s.store.UpdatePublication(ctx, trunkID, revision, pub.CellID, "failed", applyErr.Code, 0, "", operationID)
continue
}
if ack.Mode != s.cfg.PublicMode() || ack.CellID != pub.CellID || ack.TrunkID != trunkID || ack.Revision != revision || ack.ConfigSHA256 != pub.TargetDigest || ack.Status != "applied" {
failures = append(failures, map[string]any{"cell_id": pub.CellID, "code": "ACK_MISMATCH"})
_ = s.store.UpdatePublication(ctx, trunkID, revision, pub.CellID, "failed", "ACK_MISMATCH", ack.Revision, ack.ConfigSHA256, operationID)
continue
}
_ = s.store.UpdatePublication(ctx, trunkID, revision, pub.CellID, "applied", "", ack.Revision, ack.ConfigSHA256, operationID)
}
if pending && !containsFailure(failures) {
_ = s.store.MarkPublicationFailure(ctx, trunkID, revision, operationID, actor, requestID, action, "MOCK_INTENT_ONLY")
return AdminTrunk{}, newAppError(409, "PUBLISH_PENDING", "mock mode records publication intent but does not claim Cell application", map[string]any{"operation_id": operationID})
}
if len(failures) > 0 {
_ = s.store.MarkPublicationFailure(ctx, trunkID, revision, operationID, actor, requestID, action, "CELL_APPLY_FAILED")
return AdminTrunk{}, newAppError(409, "PUBLISH_PARTIAL", "not every target Cell acknowledged the exact snapshot; admission remains blocked", map[string]any{"operation_id": operationID, "failures": failures})
}
ok, err := s.store.FinalizePublication(ctx, trunkID, revision, cfg, operationID, actor, requestID, action)
if err != nil {
return AdminTrunk{}, asAppError(err)
}
if !ok {
return AdminTrunk{}, newAppError(409, "PUBLISH_PENDING", "publication is waiting for Cell acknowledgements", map[string]any{"operation_id": operationID})
}
trunk, err = s.store.GetTrunk(ctx, trunkID, s.cfg.PublicMode())
if err != nil {
return AdminTrunk{}, asAppError(err)
}
trunk.Mode = s.cfg.PublicMode()
return trunk, nil
}
func containsFailure(v []map[string]any) bool { return len(v) > 0 }
func (s *Server) realCellObservationError(ctx context.Context, trunkID string, cellIDs []string) *AppError {
missing := make([]map[string]any, 0)
for _, cellID := range cellIDs {
observations, err := s.store.LatestObservations(ctx, cellID, trunkID)
if err != nil {
return asAppError(err)
}
var current *ObservationInput
for index := range observations {
if observations[index].TrunkID == "" || observations[index].TrunkID == trunkID {
current = &observations[index]
break
}
}
if current == nil {
missing = append(missing, map[string]any{"cell_id": cellID, "reason": "no_current_observation"})
continue
}
age := time.Since(current.ObservedAt)
if current.ObservedAt.IsZero() || age < 0 || age > 15*time.Second {
missing = append(missing, map[string]any{"cell_id": cellID, "reason": "observation_stale_or_clock_skew", "observed_at": current.ObservedAt})
continue
}
if sources := missingSources(current.States); len(sources) > 0 {
missing = append(missing, map[string]any{"cell_id": cellID, "reason": "observation_incomplete", "missing_sources": sources})
}
}
if len(missing) > 0 {
return newAppError(409, "REAL_CELL_OBSERVATION_INCOMPLETE", "real publication requires fresh complete Cell observations", map[string]any{"cells": missing})
}
return nil
}
func cellAgentConfig(cfg TrunkConfig) map[string]any {
value := map[string]any{}
data, _ := json.Marshal(cfg)
_ = json.Unmarshal(data, &value)
if sip, ok := value["sip"].(map[string]any); ok && cfg.Sip.CredentialRef == "" {
sip["credential_ref"] = nil
}
return value
}
func (s *Server) applyToCell(ctx context.Context, cell Cell, trunkID string, revision, expectedLocalRevision int64, cfg TrunkConfig, digest, requestID string) (cellApplyAck, *AppError) {
if s.cfg.Mode == ModeMock {
if failure := s.cfg.MockCellFailures[cell.CellID]; failure != "" {
return cellApplyAck{}, newAppError(502, failure, "configured mock Cell failure", nil)
}
if !s.cfg.MockCellApply {
return cellApplyAck{}, newAppError(409, "MOCK_INTENT_ONLY", "mock publication is intent-only", nil)
}
return cellApplyAck{Mode: ModeMock, CellID: cell.CellID, TrunkID: trunkID, Revision: revision, ConfigSHA256: digest, Status: "applied", Idempotent: false}, nil
}
u, err := validateAgentURL(cell.Config.ManagementURL, s.cfg)
if err != nil {
return cellApplyAck{}, newAppError(503, "CELL_AGENT_UNAVAILABLE", err.Error(), nil)
}
client, err := s.httpClient()
if err != nil {
return cellApplyAck{}, newAppError(503, "MTLS_NOT_CONFIGURED", err.Error(), nil)
}
endpoint := *u
endpoint.Path = path.Join(strings.TrimSuffix(u.Path, "/"), "/v1/sip/trunks/", trunkID, "apply")
payload := cellApplyRequest{Mode: ModeReal, CellID: cell.CellID, TrunkID: trunkID, Revision: revision, ExpectedLocalRevision: expectedLocalRevision, Config: cellAgentConfig(cfg), ConfigSHA256: digest}
body, _ := json.Marshal(payload)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint.String(), strings.NewReader(string(body)))
if err != nil {
return cellApplyAck{}, newAppError(503, "CELL_AGENT_UNAVAILABLE", err.Error(), nil)
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("X-Request-ID", requestID)
resp, err := client.Do(req)
if err != nil {
return cellApplyAck{}, newAppError(503, "CELL_AGENT_UNAVAILABLE", "Cell Agent did not acknowledge the snapshot", nil)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return cellApplyAck{}, newAppError(502, "CELL_APPLY_FAILED", "Cell Agent rejected the snapshot", map[string]any{"status": resp.StatusCode})
}
var ack cellApplyAck
if err := json.NewDecoder(resp.Body).Decode(&ack); err != nil {
return cellApplyAck{}, newAppError(502, "INVALID_CELL_ACK", "Cell Agent acknowledgement is invalid", nil)
}
return ack, nil
}