feat: add deployment foundation and cross-device handoff
This commit is contained in:
@@ -0,0 +1,79 @@
|
||||
# Operation metadata store
|
||||
|
||||
Internal execution foundation; not a CLI write endpoint, application executor,
|
||||
backup engine or systemd supervisor. The existing CLI remains read-only.
|
||||
|
||||
## Contract
|
||||
|
||||
Bootstrap must supply one fixed, existing, administrator-owned local directory
|
||||
per host. Use the same canonical directory for every runner on that host.
|
||||
Untrusted users/applications must not be able to replace the directory, its
|
||||
ancestors or files. Linux deployment permissions should be 0700 for the directory
|
||||
and 0600 for metadata. Windows ACL configuration belongs to the installer.
|
||||
NFS/SMB/distributed locking is not supported.
|
||||
|
||||
`Acquire(directory, hostID)` opens an os.Root and takes a nonblocking exclusive
|
||||
OS lock on `host.lock`. Different paths/roots are not a distributed host registry.
|
||||
Hold the returned Session throughout any future application write operation.
|
||||
Never remove/replace `host.lock`, including during cleanup: its inode is the lock
|
||||
identity. Close releases it; abrupt process termination releases it at OS level.
|
||||
|
||||
- `Begin(id, planHash)` registers queued work and returns `(operation, created)`.
|
||||
Reusing ID/hash returns the original record with created=false. Another hash
|
||||
conflicts. Different IDs are blocked while an unresolved record exists.
|
||||
- `Advance(id, revision, next)` compares revision and validates the transition.
|
||||
A queued-to-running transition is the claim; a second claim fails.
|
||||
- `Get(id)` returns a value copy. Closing or poisoning a session forbids its use.
|
||||
|
||||
Allowed transitions:
|
||||
|
||||
```text
|
||||
queued → running | cancelled
|
||||
running → succeeded | failed_recovered | needs_attention | unknown
|
||||
needs_attention / unknown → succeeded | failed_recovered
|
||||
terminal states → no transitions
|
||||
```
|
||||
|
||||
Success/recovery labels are assertions by the caller, not proof. The future
|
||||
executor must verify actual effects before recording them. Unknown/running work
|
||||
survives restart without replay. Reconciliation requires inspecting actual state;
|
||||
no automatic retry, forced reset, lease expiry or stale-lock deletion is provided.
|
||||
|
||||
## Persistence
|
||||
|
||||
state.json is a versioned, canonical JSON snapshot, maximum 16 MiB, with host
|
||||
identity and SHA-256 integrity checksum. Unknown fields, duplicates, truncation,
|
||||
changed data and host mismatch are rejected. The checksum is not authentication.
|
||||
It is not an append-only audit log; event logging is a separate future component.
|
||||
|
||||
Mutation writes a unique private temporary file, syncs it, closes it, renames over
|
||||
the snapshot and (Linux) syncs the containing directory. A save error poisons the
|
||||
session because the rename may already have happened; close, reacquire and query
|
||||
before deciding anything. Temporary remnants from a killed process are ignored,
|
||||
never interpreted as successful work; automatic cleanup is not implemented.
|
||||
|
||||
Current snapshots retain all operation IDs. No pruning is provided, because
|
||||
forgetting completed IDs can re-enable an old request. At the size limit, writes
|
||||
fail closed. Admission reserves 64 bytes for the active record's later status
|
||||
and revision growth; capacity rejection does not poison read access. A retention/tombstone design is required before bounded production
|
||||
history cleanup is introduced.
|
||||
|
||||
host.lock also contains a synced initialization marker. Once work is persisted,
|
||||
a missing snapshot is rejected rather than treated as a fresh store. A crash
|
||||
between the first snapshot and marker can be repaired only from a valid snapshot.
|
||||
Protect both files. Restoring an older valid snapshot still rolls back idempotency
|
||||
history; do not restart execution without independent reconciliation. This package
|
||||
cannot detect malicious administrator edits or rollback/deletion of the entire
|
||||
state directory.
|
||||
|
||||
## Platform boundary
|
||||
|
||||
Linux uses flock and file/directory fsync. Windows uses LockFileEx for development
|
||||
tests, file sync and rename; Windows power-loss durability is not promised.
|
||||
Other OSes refuse acquisition. The local macOS panel will communicate with the
|
||||
Linux runner, not use this package as its local SQLite replacement.
|
||||
|
||||
Kernel locks are advisory on Linux; all participating writers must obey them.
|
||||
External Docker/Portainer operations are outside this lock and require drift
|
||||
checks. Multi-process tests prove process-crash behavior, not sudden power failure
|
||||
or storage-hardware reliability.
|
||||
@@ -0,0 +1,26 @@
|
||||
//go:build linux
|
||||
|
||||
package state
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"os"
|
||||
"syscall"
|
||||
)
|
||||
|
||||
func lockExclusive(f *os.File) error {
|
||||
err := syscall.Flock(int(f.Fd()), syscall.LOCK_EX|syscall.LOCK_NB)
|
||||
if errors.Is(err, syscall.EWOULDBLOCK) || errors.Is(err, syscall.EAGAIN) {
|
||||
return ErrBusy
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func syncDirectory(r *os.Root) error {
|
||||
f, err := r.Open(".")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer f.Close()
|
||||
return f.Sync()
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
//go:build !linux && !windows
|
||||
|
||||
package state
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"os"
|
||||
)
|
||||
|
||||
func lockExclusive(*os.File) error {
|
||||
return errors.New("operation store supports Linux and Windows only")
|
||||
}
|
||||
func syncDirectory(*os.Root) error { return errors.New("unsupported state durability platform") }
|
||||
@@ -0,0 +1,30 @@
|
||||
//go:build windows
|
||||
|
||||
package state
|
||||
|
||||
import (
|
||||
"os"
|
||||
"runtime"
|
||||
"syscall"
|
||||
"unsafe"
|
||||
)
|
||||
|
||||
var lockFileEx = syscall.NewLazyDLL("kernel32.dll").NewProc("LockFileEx")
|
||||
|
||||
func lockExclusive(f *os.File) error {
|
||||
var overlapped syscall.Overlapped
|
||||
// LOCKFILE_FAIL_IMMEDIATELY | LOCKFILE_EXCLUSIVE_LOCK, first byte only.
|
||||
ok, _, err := lockFileEx.Call(f.Fd(), 3, 0, 1, 0, uintptr(unsafe.Pointer(&overlapped)))
|
||||
runtime.KeepAlive(f)
|
||||
if ok != 0 {
|
||||
return nil
|
||||
}
|
||||
if err == syscall.Errno(33) {
|
||||
return ErrBusy
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
// Windows is a development platform. File.Sync is used, but Go does not provide
|
||||
// a portable directory fsync here. Do not claim Windows power-loss durability.
|
||||
func syncDirectory(*os.Root) error { return nil }
|
||||
@@ -0,0 +1,56 @@
|
||||
// Package state persists task metadata. It never executes application actions.
|
||||
package state
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"regexp"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrBusy = errors.New("host state is locked")
|
||||
ErrClosed = errors.New("session is closed")
|
||||
ErrPoisoned = errors.New("state write outcome uncertain; reopen and reconcile")
|
||||
ErrConflict = errors.New("operation identity or revision conflict")
|
||||
ErrUnresolved = errors.New("another operation requires completion or reconciliation")
|
||||
ErrTransition = errors.New("invalid operation transition")
|
||||
ErrNotFound = errors.New("operation not found")
|
||||
ErrCapacity = errors.New("operation history capacity exhausted")
|
||||
idPattern = regexp.MustCompile(`^[a-z][a-z0-9-]{0,47}$`)
|
||||
hashPattern = regexp.MustCompile(`^sha256:[a-f0-9]{64}$`)
|
||||
)
|
||||
|
||||
type Status string
|
||||
|
||||
const (
|
||||
Queued Status = "queued"
|
||||
Running Status = "running"
|
||||
Succeeded Status = "succeeded"
|
||||
FailedRecovered Status = "failed_recovered"
|
||||
NeedsAttention Status = "needs_attention"
|
||||
Unknown Status = "unknown"
|
||||
Cancelled Status = "cancelled"
|
||||
)
|
||||
|
||||
type Operation struct {
|
||||
ID string `json:"id"`
|
||||
PlanHash string `json:"planHash"`
|
||||
Status Status `json:"status"`
|
||||
Revision uint64 `json:"revision"`
|
||||
}
|
||||
|
||||
func (s Status) terminal() bool { return s == Succeeded || s == FailedRecovered || s == Cancelled }
|
||||
func (s Status) valid() bool {
|
||||
return s.terminal() || s == Queued || s == Running || s == NeedsAttention || s == Unknown
|
||||
}
|
||||
func allowed(from, to Status) bool {
|
||||
switch from {
|
||||
case Queued:
|
||||
return to == Running || to == Cancelled
|
||||
case Running:
|
||||
return to == Succeeded || to == FailedRecovered || to == NeedsAttention || to == Unknown
|
||||
case NeedsAttention, Unknown:
|
||||
return to == Succeeded || to == FailedRecovered
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,293 @@
|
||||
package state
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestIdempotencyAndTransitions(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
s, err := Acquire(dir, "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer s.Close()
|
||||
op, created, err := s.Begin("op-one", testHash)
|
||||
if err != nil || !created || op.Status != Queued || op.Revision != 1 {
|
||||
t.Fatalf("begin: %+v %v", op, err)
|
||||
}
|
||||
duplicate, created, err := s.Begin("op-one", testHash)
|
||||
if err != nil || created || duplicate != op {
|
||||
t.Fatal("duplicate mutated operation")
|
||||
}
|
||||
if _, _, err := s.Begin("op-one", "sha256:"+strings.Repeat("b", 64)); !errors.Is(err, ErrConflict) {
|
||||
t.Fatal("id rebound to another plan")
|
||||
}
|
||||
if _, _, err := s.Begin("op-two", testHash); !errors.Is(err, ErrUnresolved) {
|
||||
t.Fatal("unresolved operation bypass")
|
||||
}
|
||||
if _, err := s.Advance(op.ID, op.Revision, Succeeded); !errors.Is(err, ErrTransition) {
|
||||
t.Fatal("queued operation succeeded without running")
|
||||
}
|
||||
op, err = s.Advance(op.ID, op.Revision, Running)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := s.Advance(op.ID, 1, Succeeded); !errors.Is(err, ErrConflict) {
|
||||
t.Fatal("stale revision accepted")
|
||||
}
|
||||
if _, err := s.Advance(op.ID, op.Revision, Running); !errors.Is(err, ErrTransition) {
|
||||
t.Fatal("operation claimed twice")
|
||||
}
|
||||
op, err = s.Advance(op.ID, op.Revision, NeedsAttention)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := s.Begin("op-two", testHash); !errors.Is(err, ErrUnresolved) {
|
||||
t.Fatal("ignored manual recovery")
|
||||
}
|
||||
op, err = s.Advance(op.ID, op.Revision, FailedRecovered)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := s.Advance(op.ID, op.Revision, Running); !errors.Is(err, ErrTransition) {
|
||||
t.Fatal("terminal operation restarted")
|
||||
}
|
||||
s.Close()
|
||||
s, err = Acquire(dir, "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer s.Close()
|
||||
got, err := s.Get("op-one")
|
||||
if err != nil || got != op {
|
||||
t.Fatalf("persistence mismatch: %+v %v", got, err)
|
||||
}
|
||||
if _, created, err := s.Begin("op-two", testHash); err != nil || !created {
|
||||
t.Fatal("completed operation blocks new work")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCorruptionAndHostMismatchFailClosed(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
s, err := Acquire(dir, "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := s.Begin("op-one", testHash); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s.Close()
|
||||
if other, err := Acquire(dir, "host-two"); err == nil {
|
||||
other.Close()
|
||||
t.Fatal("wrong host accepted")
|
||||
}
|
||||
for _, data := range []string{`{`, `{}`, `null`, strings.Repeat("x", maxSnapshotBytes+1)} {
|
||||
if err := os.WriteFile(filepath.Join(dir, "state.json"), []byte(data), 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if other, err := Acquire(dir, "host-one"); err == nil {
|
||||
other.Close()
|
||||
t.Fatal("corrupt state silently reset")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestInvalidOperationsDoNotCreateSnapshot(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
s, err := Acquire(dir, "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer s.Close()
|
||||
for _, id := range []string{"", "../outside", "bad/id"} {
|
||||
if _, _, err := s.Begin(id, testHash); err == nil {
|
||||
t.Fatal("invalid operation id accepted")
|
||||
}
|
||||
}
|
||||
if _, _, err := s.Begin("op-one", "latest"); err == nil {
|
||||
t.Fatal("mutable plan binding accepted")
|
||||
}
|
||||
if _, err := s.Get("missing"); !errors.Is(err, ErrNotFound) {
|
||||
t.Fatal("missing operation not reported")
|
||||
}
|
||||
if _, err := os.Stat(filepath.Join(dir, "state.json")); !os.IsNotExist(err) {
|
||||
t.Fatal("invalid request persisted state")
|
||||
}
|
||||
}
|
||||
|
||||
func TestPersistenceFailurePoisonsSession(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
s, err := Acquire(dir, "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer s.Close()
|
||||
// A directory at the destination makes atomic replacement fail on both OSes.
|
||||
if err := os.Mkdir(filepath.Join(dir, "state.json"), 0700); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := s.Begin("op-one", testHash); err == nil {
|
||||
t.Fatal("reported failed persistence as success")
|
||||
}
|
||||
if _, _, err := s.Begin("op-two", testHash); !errors.Is(err, ErrPoisoned) {
|
||||
t.Fatalf("continued after ambiguous write: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestConcurrentClaimsHaveOneWinner(t *testing.T) {
|
||||
s, err := Acquire(t.TempDir(), "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer s.Close()
|
||||
var createdCount, claimedCount atomic.Int32
|
||||
var wg sync.WaitGroup
|
||||
for n := 0; n < 16; n++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
_, created, err := s.Begin("op-one", testHash)
|
||||
if err != nil {
|
||||
t.Error(err)
|
||||
}
|
||||
if created {
|
||||
createdCount.Add(1)
|
||||
}
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
for n := 0; n < 16; n++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
_, err := s.Advance("op-one", 1, Running)
|
||||
if err == nil {
|
||||
claimedCount.Add(1)
|
||||
} else if !errors.Is(err, ErrConflict) {
|
||||
t.Error(err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
if createdCount.Load() != 1 || claimedCount.Load() != 1 {
|
||||
t.Fatalf("duplicate winners: %d %d", createdCount.Load(), claimedCount.Load())
|
||||
}
|
||||
}
|
||||
|
||||
func TestTamperedSnapshotRejected(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
s, err := Acquire(dir, "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := s.Begin("op-one", testHash); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s.Close()
|
||||
path := filepath.Join(dir, "state.json")
|
||||
raw, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, changed := range []string{
|
||||
strings.Replace(string(raw), `"queued"`, `"succeeded"`, 1),
|
||||
strings.Replace(string(raw), `"version":1`, `"version":1,"version":1`, 1),
|
||||
string(raw) + ` {}`,
|
||||
} {
|
||||
if err := os.WriteFile(path, []byte(changed), 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if s, err := Acquire(dir, "host-one"); err == nil {
|
||||
s.Close()
|
||||
t.Fatal("tampered snapshot accepted")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAllTerminalAndReconciliationPaths(t *testing.T) {
|
||||
for _, path := range [][]Status{{Cancelled}, {Running, Succeeded}, {Running, FailedRecovered}, {Running, Unknown, FailedRecovered}, {Running, NeedsAttention, Succeeded}} {
|
||||
s, err := Acquire(t.TempDir(), "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
op, _, err := s.Begin("op-one", testHash)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, next := range path {
|
||||
op, err = s.Advance(op.ID, op.Revision, next)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if _, _, err := s.Begin("op-two", testHash); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s.Close()
|
||||
}
|
||||
}
|
||||
|
||||
func TestMissingSnapshotDoesNotResetInitializedStore(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
s, err := Acquire(dir, "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := s.Begin("op-one", testHash); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s.Close()
|
||||
if err := os.Remove(filepath.Join(dir, "state.json")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if s, err := Acquire(dir, "host-one"); err == nil {
|
||||
s.Close()
|
||||
t.Fatal("silently reset initialized store")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAdmissionReservesSpaceForCompletion(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
data := snapshot{Version: 1, HostID: "host-one", Operations: make(map[string]Operation)}
|
||||
for n := 0; n < 113357; n++ {
|
||||
id := fmt.Sprintf("op%06d", n)
|
||||
if n < 2 {
|
||||
id += strings.Repeat("a", 26)
|
||||
}
|
||||
data.Operations[id] = Operation{ID: id, PlanHash: testHash, Status: Cancelled, Revision: 2}
|
||||
}
|
||||
raw, err := encodeSnapshot(data)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.WriteFile(filepath.Join(dir, "state.json"), raw, 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Confirm this fixture reaches the actual bug boundary, not an arbitrary limit.
|
||||
data.Operations["op-new"] = Operation{ID: "op-new", PlanHash: testHash, Status: Queued, Revision: 1}
|
||||
queued, err := encodeSnapshot(data)
|
||||
if err != nil || len(queued) != maxSnapshotBytes-2 {
|
||||
t.Fatalf("fixture outside boundary: %d %v", len(queued), err)
|
||||
}
|
||||
s, err := Acquire(dir, "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer s.Close()
|
||||
if _, _, err := s.Begin("op-new", testHash); err == nil {
|
||||
t.Fatal("admitted task without room for terminal state")
|
||||
}
|
||||
if _, err := s.Get("op000002"); err != nil {
|
||||
t.Fatalf("capacity rejection poisoned read access: %v", err)
|
||||
}
|
||||
if _, err := s.Get("op-new"); !errors.Is(err, ErrNotFound) {
|
||||
t.Fatal("rejected admission persisted")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,120 @@
|
||||
package state
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"crypto/rand"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
)
|
||||
|
||||
const maxSnapshotBytes = 16 << 20
|
||||
|
||||
type snapshot struct {
|
||||
Version int `json:"version"`
|
||||
HostID string `json:"hostId"`
|
||||
Operations map[string]Operation `json:"operations"`
|
||||
Checksum string `json:"checksum"`
|
||||
}
|
||||
|
||||
func encodeSnapshot(s snapshot) ([]byte, error) {
|
||||
s.Checksum = ""
|
||||
raw, err := json.Marshal(s)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
sum := sha256.Sum256(raw)
|
||||
s.Checksum = "sha256:" + hex.EncodeToString(sum[:])
|
||||
raw, err = json.Marshal(s)
|
||||
if len(raw) > maxSnapshotBytes {
|
||||
return nil, errors.New("state capacity exceeded")
|
||||
}
|
||||
return raw, err
|
||||
}
|
||||
|
||||
func readSnapshot(r *os.Root, hostID string, initialized bool) (snapshot, error) {
|
||||
empty := snapshot{Version: 1, HostID: hostID, Operations: make(map[string]Operation)}
|
||||
if err := regularOrMissing(r, "state.json"); err != nil {
|
||||
return snapshot{}, err
|
||||
}
|
||||
f, err := r.Open("state.json")
|
||||
if os.IsNotExist(err) {
|
||||
if initialized {
|
||||
return snapshot{}, errors.New("initialized store has lost its snapshot")
|
||||
}
|
||||
return empty, nil
|
||||
}
|
||||
if err != nil {
|
||||
return snapshot{}, err
|
||||
}
|
||||
defer f.Close()
|
||||
raw, err := io.ReadAll(io.LimitReader(f, maxSnapshotBytes+1))
|
||||
if err != nil {
|
||||
return snapshot{}, err
|
||||
}
|
||||
if len(raw) > maxSnapshotBytes {
|
||||
return snapshot{}, errors.New("state exceeds size limit")
|
||||
}
|
||||
var s snapshot
|
||||
if err := json.Unmarshal(raw, &s); err != nil {
|
||||
return snapshot{}, errors.New("corrupt state")
|
||||
}
|
||||
if s.Version != 1 || s.HostID != hostID || s.Operations == nil {
|
||||
return snapshot{}, errors.New("incompatible state or host mismatch")
|
||||
}
|
||||
canonical, err := encodeSnapshot(s)
|
||||
if err != nil || !bytes.Equal(raw, canonical) {
|
||||
return snapshot{}, errors.New("state integrity check failed")
|
||||
}
|
||||
unresolved := 0
|
||||
for id, op := range s.Operations {
|
||||
if id != op.ID || !idPattern.MatchString(id) || !hashPattern.MatchString(op.PlanHash) || !op.Status.valid() || op.Revision == 0 {
|
||||
return snapshot{}, errors.New("invalid operation record")
|
||||
}
|
||||
if !op.Status.terminal() {
|
||||
unresolved++
|
||||
}
|
||||
}
|
||||
if unresolved > 1 {
|
||||
return snapshot{}, errors.New("multiple unresolved operations")
|
||||
}
|
||||
return s, nil
|
||||
}
|
||||
|
||||
func writeSnapshot(r *os.Root, s snapshot) error {
|
||||
raw, err := encodeSnapshot(s)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := regularOrMissing(r, "state.json"); err != nil {
|
||||
return err
|
||||
}
|
||||
var entropy [16]byte
|
||||
if _, err := rand.Read(entropy[:]); err != nil {
|
||||
return err
|
||||
}
|
||||
name := ".state-" + hex.EncodeToString(entropy[:]) + ".tmp"
|
||||
f, err := r.OpenFile(name, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0600)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer r.Remove(name) // Only the unique file just created; never state.json or host.lock.
|
||||
if _, err := f.Write(raw); err != nil {
|
||||
f.Close()
|
||||
return err
|
||||
}
|
||||
if err := f.Sync(); err != nil {
|
||||
f.Close()
|
||||
return err
|
||||
}
|
||||
if err := f.Close(); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := r.Rename(name, "state.json"); err != nil {
|
||||
return err
|
||||
}
|
||||
return syncDirectory(r)
|
||||
}
|
||||
@@ -0,0 +1,225 @@
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,134 @@
|
||||
package state
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
const testHash = "sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
|
||||
|
||||
func TestExclusiveSession(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
s, err := Acquire(dir, "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { s.Close() })
|
||||
if second, err := Acquire(dir, "host-one"); second != nil || !errors.Is(err, ErrBusy) {
|
||||
t.Fatalf("lock bypass: %v", err)
|
||||
}
|
||||
other, err := Acquire(t.TempDir(), "host-two")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
other.Close()
|
||||
if err := s.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := s.Begin("op-one", testHash); !errors.Is(err, ErrClosed) {
|
||||
t.Fatal("used closed session")
|
||||
}
|
||||
next, err := Acquire(dir, "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
next.Close()
|
||||
}
|
||||
|
||||
func TestInvalidHost(t *testing.T) {
|
||||
for _, host := range []string{"", "../host", "Host", strings.Repeat("a", 49)} {
|
||||
if s, err := Acquire(t.TempDir(), host); err == nil {
|
||||
s.Close()
|
||||
t.Fatal("accepted invalid host")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestStatePathsRejectSymlinks(t *testing.T) {
|
||||
for _, name := range []string{"host.lock", "state.json"} {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
dir, outside := t.TempDir(), filepath.Join(t.TempDir(), "outside")
|
||||
if err := os.WriteFile(outside, []byte("protected"), 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.Symlink(outside, filepath.Join(dir, name)); err != nil {
|
||||
t.Skipf("symlink privilege unavailable: %v", err)
|
||||
}
|
||||
if s, err := Acquire(dir, "host-one"); err == nil {
|
||||
s.Close()
|
||||
t.Fatal("accepted symlink state")
|
||||
}
|
||||
content, err := os.ReadFile(outside)
|
||||
if err != nil || string(content) != "protected" {
|
||||
t.Fatal("modified outside file")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// The subprocess exits without Close: the OS must release its lock, but its
|
||||
// running operation must remain recorded. No mock can test this boundary.
|
||||
func TestProcessDeathPreservesRunningOperation(t *testing.T) {
|
||||
if os.Getenv("DEPLOYCTL_STATE_CHILD") == "1" {
|
||||
s, err := Acquire(os.Getenv("DEPLOYCTL_STATE_DIR"), "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
op, _, err := s.Begin("op-one", testHash)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err = s.Advance(op.ID, op.Revision, Running); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
fmt.Println("READY")
|
||||
bufio.NewReader(os.Stdin).ReadByte()
|
||||
os.Exit(3)
|
||||
}
|
||||
dir := t.TempDir()
|
||||
cmd := exec.Command(os.Args[0], "-test.run=^TestProcessDeathPreservesRunningOperation$")
|
||||
cmd.Env = append(os.Environ(), "DEPLOYCTL_STATE_CHILD=1", "DEPLOYCTL_STATE_DIR="+dir)
|
||||
stdin, err := cmd.StdinPipe()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer stdin.Close()
|
||||
stdout, err := cmd.StdoutPipe()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
cmd.Stderr = os.Stderr
|
||||
if err := cmd.Start(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { cmd.Process.Kill(); cmd.Wait() })
|
||||
line, err := bufio.NewReader(stdout).ReadString('\n')
|
||||
if err != nil || strings.TrimSpace(line) != "READY" {
|
||||
t.Fatalf("child failed: %q %v", line, err)
|
||||
}
|
||||
if _, err := Acquire(dir, "host-one"); !errors.Is(err, ErrBusy) {
|
||||
t.Fatalf("cross-process lock bypass: %v", err)
|
||||
}
|
||||
if err := cmd.Process.Kill(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
cmd.Wait()
|
||||
s, err := Acquire(dir, "host-one")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer s.Close()
|
||||
op, created, err := s.Begin("op-one", testHash)
|
||||
if err != nil || created || op.Status != Running || op.Revision != 2 {
|
||||
t.Fatalf("lost interrupted task: %+v %v", op, err)
|
||||
}
|
||||
if _, _, err := s.Begin("op-two", testHash); !errors.Is(err, ErrUnresolved) {
|
||||
t.Fatal("allowed writes before reconciliation")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user