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

102 lines
3.2 KiB
Go

package management
import (
"context"
"database/sql"
"encoding/json"
"strings"
)
func (s *Store) ListUnfinishedOperations(ctx context.Context) ([]Operation, error) {
var out []Operation
err := s.withRead(ctx, func(db *sql.DB) error {
rows, err := db.QueryContext(ctx, "SELECT operation_id,actor,action,resource_type,resource_id,request_id,request_digest,expected_revision,state,http_status,response_json,error_json,created_at,updated_at,expires_at FROM operations WHERE state IN ('in_progress','pending','needs_reconciliation') ORDER BY created_at ASC")
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var op Operation
var response, failure sql.NullString
if err := rows.Scan(&op.OperationID, &op.Actor, &op.Action, &op.ResourceType, &op.ResourceID, &op.RequestID, &op.RequestDigest, &op.ExpectedRevision, &op.State, &op.HTTPStatus, &response, &failure, &op.CreatedAt, &op.UpdatedAt, &op.ExpiresAt); err != nil {
return err
}
if response.Valid {
op.ResponseJSON = response.String
}
if failure.Valid {
op.ErrorJSON = failure.String
}
out = append(out, op)
}
return rows.Err()
})
return out, err
}
// Recover finalizes only durable, already-acknowledged snapshots. It never redelivers an uncertain Cell request.
func (s *Server) Recover(ctx context.Context) error {
ops, err := s.store.ListUnfinishedOperations(ctx)
if err != nil {
return err
}
for _, op := range ops {
if !strings.HasPrefix(op.Action, "trunk.") {
continue
}
pubs, err := s.store.GetPublicationRows(ctx, op.ResourceID, nil)
if err != nil {
return err
}
var revision int64
allApplied := false
for _, pub := range pubs {
if pub.OperationID != op.OperationID {
continue
}
if revision == 0 {
revision = pub.Revision
allApplied = true
}
if pub.Revision != revision || pub.Status != "applied" {
allApplied = false
}
}
if revision == 0 {
_ = s.store.FinishOperation(ctx, op.OperationID, "needs_reconciliation", 409, "", recoveryError(op))
continue
}
if !allApplied {
_ = s.store.FinishOperation(ctx, op.OperationID, "needs_reconciliation", 409, "", recoveryError(op))
continue
}
cfg, err := s.store.getTrunkConfig(ctx, op.ResourceID, revision)
if err != nil {
return err
}
ok, err := s.store.FinalizePublication(ctx, op.ResourceID, revision, cfg, op.OperationID, op.Actor, op.RequestID, strings.TrimPrefix(op.Action, "trunk."))
if err != nil {
return err
}
if !ok {
_ = s.store.FinishOperation(ctx, op.OperationID, "needs_reconciliation", 409, "", recoveryError(op))
continue
}
trunk, err := s.store.GetTrunk(ctx, op.ResourceID, s.cfg.PublicMode())
if err != nil {
return err
}
trunk.Mode = s.cfg.PublicMode()
body, _ := json.Marshal(trunk)
if err := s.store.FinishOperation(ctx, op.OperationID, "succeeded", 200, string(body), ""); err != nil {
return err
}
}
return nil
}
func recoveryError(op Operation) string {
b, _ := json.Marshal(errorResponse{Error: ErrorBody{Code: "OPERATION_RECOVERY_REQUIRED", Message: "the operation stopped before a durable all-Cell acknowledgement; no request was redelivered", Fields: map[string]any{"operation_id": op.OperationID}}})
return string(b)
}