61 lines
2.8 KiB
Go
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
|
|
}
|