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

61 lines
2.8 KiB
Go

package management
import (
"context"
"database/sql"
"errors"
"time"
)
// CreateAdmissionIntent writes the fail-closed boundary before a disable or rollback creates a new snapshot.
func (s *Store) CreateAdmissionIntent(ctx context.Context, trunkID string, revision int64, operationID, reason string) error {
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)
if err == nil && existing.State != "released" && existing.OperationID != operationID {
return newAppError(409, "PUBLISH_IN_PROGRESS", "another publication currently owns the admission barrier", map[string]any{"operation_id": existing.OperationID})
}
if err != nil && !errors.Is(err, sql.ErrNoRows) {
return err
}
_, 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)
return err
})
}
func (s *Store) SetBarrierState(ctx context.Context, trunkID, operationID, state string) error {
if state != "active" && state != "pending" && state != "released" {
return errors.New("invalid barrier state")
}
return s.withTx(ctx, func(tx *sql.Tx) error {
_, err := tx.ExecContext(ctx, "UPDATE admission_barriers SET state=?,updated_at=? WHERE trunk_id=? AND operation_id=?", state, utcString(time.Now()), trunkID, operationID)
return err
})
}
// AdmissionAllowed is the single persisted gate an external scheduler can consult before creating a new call.
func (s *Store) AdmissionAllowed(ctx context.Context, trunkID string) (bool, string, error) {
var activeStatus string
var barrierState sql.NullString
err := s.withRead(ctx, func(db *sql.DB) error {
return db.QueryRowContext(ctx, "SELECT active_status FROM trunks WHERE trunk_id=?", trunkID).Scan(&activeStatus)
})
if err != nil {
return false, "trunk_not_found", err
}
err = s.withRead(ctx, func(db *sql.DB) error {
return db.QueryRowContext(ctx, "SELECT state FROM admission_barriers WHERE trunk_id=?", trunkID).Scan(&barrierState)
})
if err == nil && barrierState.Valid && barrierState.String != "released" {
return false, "admission_barrier", nil
}
if err != nil && !errors.Is(err, sql.ErrNoRows) {
return false, "unknown", err
}
if activeStatus != "published" {
return false, "trunk_not_published", nil
}
return true, "ok", nil
}