Files
server-deploy/internal/state/store.go

226 lines
5.7 KiB
Go

package state
import (
"errors"
"os"
"path/filepath"
"sync"
)
// Session holds an exclusive host lock for its entire lifetime. The caller must
// retain it throughout any future write operation, not only metadata updates.
// Never unlink host.lock: replacing its inode defeats OS locking.
type Session struct {
mu sync.Mutex
root *os.Root
lock *os.File
data snapshot
closed bool
poisoned bool
initialized bool
}
const initializedMarker = "server-deploy-state-v1\n"
// Acquire requires an existing administrator-owned LOCAL directory. Bootstrap
// owns directory creation/permissions; this function never creates a new root.
func Acquire(directory, hostID string) (*Session, error) {
if !idPattern.MatchString(hostID) || !filepath.IsAbs(directory) {
return nil, errors.New("invalid host or state directory")
}
r, err := os.OpenRoot(directory)
if err != nil {
return nil, err
}
fail := func(err error) (*Session, error) { r.Close(); return nil, err }
if err := regularOrMissing(r, "host.lock"); err != nil {
return fail(err)
}
lock, err := r.OpenFile("host.lock", os.O_CREATE|os.O_RDWR, 0600)
if err != nil {
return fail(err)
}
if err := lockExclusive(lock); err != nil {
lock.Close()
return fail(err)
}
s := &Session{root: r, lock: lock}
info, err := lock.Stat()
if err != nil {
s.Close()
return nil, err
}
if info.Size() != 0 {
if info.Size() != int64(len(initializedMarker)) {
s.Close()
return nil, errors.New("invalid state initialization marker")
}
marker := make([]byte, len(initializedMarker))
if _, err := lock.ReadAt(marker, 0); err != nil || string(marker) != initializedMarker {
s.Close()
return nil, errors.New("corrupt state initialization marker")
}
s.initialized = true
}
s.data, err = readSnapshot(r, hostID, s.initialized)
if err != nil {
s.Close()
return nil, err
}
// A crash may happen after the first snapshot rename but before its marker.
// Repair only from a fully validated existing snapshot, never from absence.
if len(s.data.Operations) > 0 {
if err := s.markInitialized(); err != nil {
s.Close()
return nil, err
}
}
return s, nil
}
func (s *Session) markInitialized() error {
if s.initialized {
return nil
}
if _, err := s.lock.WriteAt([]byte(initializedMarker), 0); err != nil {
return err
}
if err := s.lock.Sync(); err != nil {
return err
}
s.initialized = true
return nil
}
func regularOrMissing(r *os.Root, name string) error {
info, err := r.Lstat(name)
if os.IsNotExist(err) {
return nil
}
if err != nil {
return err
}
if !info.Mode().IsRegular() {
return errors.New("state path must be a regular file, not a link")
}
return nil
}
func (s *Session) ready() error {
if s.closed {
return ErrClosed
}
if s.poisoned {
return ErrPoisoned
}
return nil
}
func (s *Session) Close() error {
s.mu.Lock()
defer s.mu.Unlock()
if s.closed {
return nil
}
s.closed = true
return errors.Join(s.lock.Close(), s.root.Close())
}
func (s *Session) Get(id string) (Operation, error) {
s.mu.Lock()
defer s.mu.Unlock()
if err := s.ready(); err != nil {
return Operation{}, err
}
op, ok := s.data.Operations[id]
if !ok {
return Operation{}, ErrNotFound
}
return op, nil
}
// Begin is idempotent, not an execution claim. A repeated ID cannot be rebound
// to a different plan, and any unresolved operation blocks unrelated work.
func (s *Session) Begin(id, planHash string) (Operation, bool, error) {
s.mu.Lock()
defer s.mu.Unlock()
if err := s.ready(); err != nil {
return Operation{}, false, err
}
if !idPattern.MatchString(id) || !hashPattern.MatchString(planHash) {
return Operation{}, false, errors.New("invalid operation identity")
}
if op, ok := s.data.Operations[id]; ok {
if op.PlanHash != planHash {
return Operation{}, false, ErrConflict
}
return op, false, nil
}
for _, op := range s.data.Operations {
if !op.Status.terminal() {
return Operation{}, false, ErrUnresolved
}
}
op := Operation{ID: id, PlanHash: planHash, Status: Queued, Revision: 1}
if err := s.save(op); err != nil {
return Operation{}, false, err
}
return op, true, nil
}
// Advance uses a revision check so only one caller can claim queued work.
// Reconciliation to a terminal status requires external evidence; this metadata
// layer cannot prove that a deployment or restore actually succeeded.
func (s *Session) Advance(id string, revision uint64, next Status) (Operation, error) {
s.mu.Lock()
defer s.mu.Unlock()
if err := s.ready(); err != nil {
return Operation{}, err
}
op, ok := s.data.Operations[id]
if !ok {
return Operation{}, ErrNotFound
}
if op.Revision != revision {
return Operation{}, ErrConflict
}
if !allowed(op.Status, next) {
return Operation{}, ErrTransition
}
if op.Revision == ^uint64(0) {
return Operation{}, ErrConflict
}
op.Status = next
op.Revision++
if err := s.save(op); err != nil {
return Operation{}, err
}
return op, nil
}
func (s *Session) save(op Operation) error {
updated := snapshot{Version: 1, HostID: s.data.HostID, Operations: make(map[string]Operation, len(s.data.Operations)+1)}
for id, old := range s.data.Operations {
updated.Operations[id] = old
}
updated.Operations[op.ID] = op
if op.Revision == 1 {
// Only one unresolved operation is admitted. Reserve more than the
// longest status growth plus uint64 revision growth (at most 29 bytes).
raw, err := encodeSnapshot(updated)
if err != nil || len(raw) > maxSnapshotBytes-64 {
return ErrCapacity
}
}
if err := writeSnapshot(s.root, updated); err != nil {
s.poisoned = true
return errors.Join(ErrPoisoned, err)
}
if err := s.markInitialized(); err != nil {
s.poisoned = true
return errors.Join(ErrPoisoned, err)
}
s.data = updated
return nil
}