102 lines
3.2 KiB
Go
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)
|
|
}
|