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 }