Harden lease lifecycle durability

This commit is contained in:
kami
2026-07-30 14:34:29 +04:00
parent 1ff0af2e69
commit f6ee0e3060
40 changed files with 2108 additions and 590 deletions
+148 -43
View File
@@ -32,27 +32,10 @@ func Open(dir string) (*Store, error) {
if err := os.MkdirAll(s.cas, 0755); err != nil {
return nil, err
}
var snapshotSeq uint64
if b, readErr := os.ReadFile(s.snapshot); readErr == nil {
var snap struct {
Seq uint64 `json:"seq"`
Tasks []domain.Task `json:"tasks"`
}
if json.Unmarshal(b, &snap) != nil {
return nil, fmt.Errorf("invalid snapshot")
}
for _, t := range snap.Tasks {
s.tasks[t.ID] = t
s.external[t.Source+"\x00"+t.ExternalID] = t.ID
}
snapshotSeq = snap.Seq
// Continue event numbering after the snapshot. Without restoring this
// cursor, the first append after a restart reused sequence 1 and made
// the append-only log unreplayable.
s.seq = snapshotSeq
} else if !errors.Is(readErr, os.ErrNotExist) {
return nil, readErr
}
// A snapshot is a disposable read cache, never recovery authority. Loading
// it before the log let a partially-written snapshot become a different
// history than events.jsonl after a crash. Rebuild every projection from
// the append-only, fsynced log instead.
f, err := os.Open(s.path)
if os.IsNotExist(err) {
return s, nil
@@ -62,16 +45,13 @@ func Open(dir string) (*Store, error) {
}
defer f.Close()
sc := bufio.NewScanner(f)
var expected uint64 = snapshotSeq + 1
var expected uint64 = 1
for sc.Scan() {
var e domain.Event
if err := json.Unmarshal(sc.Bytes(), &e); err == nil {
if err := domain.ValidateEvent(e); err != nil {
return nil, err
}
if e.Seq < expected {
continue
}
if e.Seq != expected {
return nil, fmt.Errorf("event sequence gap: got %d, want %d", e.Seq, expected)
}
@@ -151,9 +131,20 @@ func (s *Store) apply(e domain.Event) error {
s.external[t.Source+"\x00"+t.ExternalID] = t.ID
case "TaskLeased":
t.State = domain.StateLeased
t.Lease = &domain.Lease{HarnessID: p["harness_id"].(string), Until: time.Unix(0, int64(p["until_ns"].(float64)))}
epoch, _ := p["lease_epoch"].(string)
if epoch == "" {
// A pre-fencing event cannot safely be renewed by an old worker.
// Deriving a stable token from the durable event identity makes the
// recovered lease observable but non-renewable until it expires.
epoch = "legacy:" + e.ID
}
t.Lease = &domain.Lease{HarnessID: p["harness_id"].(string), Epoch: epoch, Until: time.Unix(0, int64(p["until_ns"].(float64)))}
case "TaskLeaseRenewed":
t.Lease = &domain.Lease{HarnessID: p["harness_id"].(string), Until: time.Unix(0, int64(p["until_ns"].(float64)))}
epoch, _ := p["lease_epoch"].(string)
if epoch == "" && t.Lease != nil {
epoch = t.Lease.Epoch
}
t.Lease = &domain.Lease{HarnessID: p["harness_id"].(string), Epoch: epoch, Until: time.Unix(0, int64(p["until_ns"].(float64)))}
case "TaskReleased":
t.State = domain.StateQueued
t.Lease = nil
@@ -292,6 +283,9 @@ func (s *Store) Append(e domain.Event) error {
return domain.ErrConflict
}
}
if err := s.validateTransition(e, t, taskExists, contract); err != nil {
return err
}
if e.Type == "TaskLeased" {
var p struct {
ExpectedVersion *int `json:"expected_version"`
@@ -333,9 +327,6 @@ func (s *Store) Append(e domain.Event) error {
}
}
}
if err := s.apply(e); err != nil {
return err
}
f, err := os.OpenFile(s.path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0644)
if err != nil {
return err
@@ -348,10 +339,60 @@ func (s *Store) Append(e domain.Event) error {
if err = f.Sync(); err != nil {
return err
}
// The event is the commit record. Do not expose a projection that cannot
// be recovered from it after a power loss.
if err := s.apply(e); err != nil {
return err
}
s.events = append(s.events, e)
s.seq = e.Seq
if err := s.writeSnapshot(); err != nil {
return err
// Snapshot failure does not roll back a committed event. Open always
// rebuilds from the log, so leaving a stale cache is safe.
_ = s.writeSnapshot()
return nil
}
// validateTransition keeps lifecycle authority at the durable append
// boundary. A task may be completed/failed while queued by an external
// provider, but once a lease exists its owner and fencing epoch are required
// for every lifecycle mutation.
func (s *Store) validateTransition(e domain.Event, t domain.Task, exists bool, p map[string]any) error {
if !exists {
if e.Type != "TaskCreated" && e.Type != "QuotaReported" && e.Type != "StandupAdvisory" && e.Type != "ApprovalGranted" && e.Type != "ApprovalDenied" {
return domain.ErrNotFound
}
return nil
}
if e.Type == "TaskLeased" && t.State != domain.StateQueued {
return domain.ErrConflict
}
if t.State != domain.StateLeased || t.Lease == nil {
return nil
}
switch e.Type {
case "TaskLeaseRenewed", "TaskReleased", "TaskPickupValidated", "TaskCompleted", "TaskBlocked", "TaskFailed":
owner, _ := p["harness_id"].(string)
epoch, _ := p["lease_epoch"].(string)
// Expiry is the one coordinator-owned relinquish path. It still binds
// the exact epoch that was observed when the timer fired.
if e.Type == "TaskReleased" {
if reason, _ := p["reason"].(string); (reason == "lease_expired" || reason == "pane_exited") && owner == t.Lease.HarnessID && epoch == t.Lease.Epoch {
return nil
}
}
if owner != t.Lease.HarnessID || epoch == "" || epoch != t.Lease.Epoch {
return domain.ErrConflict
}
case "TaskCorrected":
// Corrections may repair metadata while a task is leased, but cannot
// smuggle in a lifecycle transition around the current fenced owner.
if _, changesState := p["state"]; changesState {
owner, _ := p["harness_id"].(string)
epoch, _ := p["lease_epoch"].(string)
if owner != t.Lease.HarnessID || epoch == "" || epoch != t.Lease.Epoch {
return domain.ErrConflict
}
}
}
return nil
}
@@ -368,10 +409,28 @@ func (s *Store) writeSnapshot() error {
return err
}
tmp := s.snapshot + ".tmp"
if err = os.WriteFile(tmp, b, 0644); err != nil {
f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0644)
if err != nil {
return err
}
return os.Rename(tmp, s.snapshot)
if _, err = f.Write(b); err == nil {
err = f.Sync()
}
if closeErr := f.Close(); err == nil {
err = closeErr
}
if err != nil {
return err
}
if err = os.Rename(tmp, s.snapshot); err != nil {
return err
}
dir, err := os.Open(filepath.Dir(s.snapshot))
if err != nil {
return err
}
defer dir.Close()
return dir.Sync()
}
func (s *Store) Tasks() []domain.Task {
s.mu.Lock()
@@ -397,9 +456,39 @@ func (s *Store) PutArtifact(b []byte) (string, error) {
h := domain.Hash(b)
p := filepath.Join(s.cas, h)
if _, err := os.Stat(p); errors.Is(err, os.ErrNotExist) {
if err = os.WriteFile(p, b, 0644); err != nil {
tmp := p + ".tmp-" + domain.NewID()
f, openErr := os.OpenFile(tmp, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0644)
if openErr != nil {
return "", openErr
}
if _, err = f.Write(b); err == nil {
err = f.Sync()
}
if closeErr := f.Close(); err == nil {
err = closeErr
}
if err != nil {
_ = os.Remove(tmp)
return "", err
}
if err = os.Rename(tmp, p); err != nil && !errors.Is(err, os.ErrExist) {
_ = os.Remove(tmp)
return "", err
}
if !errors.Is(err, os.ErrExist) {
dir, openErr := os.Open(s.cas)
if openErr != nil {
return "", openErr
}
syncErr := dir.Sync()
closeErr := dir.Close()
if syncErr != nil {
return "", syncErr
}
if closeErr != nil {
return "", closeErr
}
}
}
return h, nil
}
@@ -450,7 +539,7 @@ func (s *Store) Lease(id, harness string, ttl time.Duration) (domain.Event, erro
if t.State != domain.StateQueued {
return domain.Event{}, domain.ErrConflict
}
payload := map[string]any{"harness_id": harness, "ttl": ttl.Seconds(), "until_ns": time.Now().Add(ttl).UnixNano(), "expected_version": t.Version}
payload := map[string]any{"harness_id": harness, "lease_epoch": domain.NewID(), "ttl": ttl.Seconds(), "until_ns": time.Now().Add(ttl).UnixNano(), "expected_version": t.Version}
if t.HandoffRef != "" {
payload["handoff_ref"] = t.HandoffRef
payload["transaction_id"] = t.ReleaseTransaction
@@ -464,7 +553,7 @@ func (s *Store) Lease(id, harness string, ttl time.Duration) (domain.Event, erro
// RenewLease atomically extends the current owner's lease. The observed task
// version is part of the request so an old worker can never renew a lease
// after release/reassignment.
func (s *Store) RenewLease(id, harness string, expectedVersion int, ttl time.Duration) (domain.Event, error) {
func (s *Store) RenewLease(id, harness, epoch string, expectedVersion int, ttl time.Duration) (domain.Event, error) {
if ttl <= 0 {
return domain.Event{}, fmt.Errorf("%w: ttl must be positive", domain.ErrInvalid)
}
@@ -472,10 +561,10 @@ func (s *Store) RenewLease(id, harness string, expectedVersion int, ttl time.Dur
if !ok {
return domain.Event{}, domain.ErrNotFound
}
if t.State != domain.StateLeased || t.Lease == nil || t.Lease.HarnessID != harness || t.Version != expectedVersion {
if t.State != domain.StateLeased || t.Lease == nil || t.Lease.HarnessID != harness || t.Lease.Epoch != epoch || t.Version != expectedVersion {
return domain.Event{}, domain.ErrConflict
}
p, _ := json.Marshal(map[string]any{"harness_id": harness, "until_ns": time.Now().Add(ttl).UnixNano(), "expected_version": expectedVersion})
p, _ := json.Marshal(map[string]any{"harness_id": harness, "lease_epoch": epoch, "until_ns": time.Now().Add(ttl).UnixNano(), "expected_version": expectedVersion})
e := domain.Event{ID: domain.NewID(), Type: "TaskLeaseRenewed", TaskID: id, Version: t.Version + 1, Payload: p, Surface: string(authz.System)}
return e, s.Append(e)
}
@@ -483,14 +572,30 @@ func (s *Store) RenewLease(id, harness string, expectedVersion int, ttl time.Dur
func (s *Store) ExpireLeases(now time.Time) ([]domain.Event, error) {
var out []domain.Event
for _, t := range s.Tasks() {
if t.State == domain.StateLeased && t.Lease != nil && !t.Lease.Until.After(now) {
p, _ := json.Marshal(map[string]any{"reason": "lease_expired", "harness_id": t.Lease.HarnessID})
e := domain.Event{ID: domain.NewID(), Type: "TaskReleased", TaskID: t.ID, Version: t.Version + 1, Payload: p, Surface: string(authz.System)}
if err := s.Append(e); err != nil {
if e, err := s.ExpireLease(t.ID, now); err != nil {
if !errors.Is(err, domain.ErrConflict) {
return out, err
}
} else if e.ID != "" {
out = append(out, e)
}
}
return out, nil
}
// ExpireLease releases exactly the observed lease if, and only if, its TTL
// has elapsed. Coordinators use this one-task form to stop their local pane
// before publishing the release event; the batch helper remains for
// deployments without a local coordinator.
func (s *Store) ExpireLease(id string, now time.Time) (domain.Event, error) {
t, ok := s.Task(id)
if !ok {
return domain.Event{}, domain.ErrNotFound
}
if t.State != domain.StateLeased || t.Lease == nil || t.Lease.Until.After(now) {
return domain.Event{}, domain.ErrConflict
}
p, _ := json.Marshal(map[string]any{"reason": "lease_expired", "harness_id": t.Lease.HarnessID, "lease_epoch": t.Lease.Epoch, "expected_version": t.Version})
e := domain.Event{ID: domain.NewID(), Type: "TaskReleased", TaskID: t.ID, Version: t.Version + 1, Payload: p, Surface: string(authz.System)}
return e, s.Append(e)
}
+105 -7
View File
@@ -3,6 +3,7 @@ package store
import (
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"strings"
@@ -35,7 +36,7 @@ func TestLeaseCarriesReleasedHandoffRef(t *testing.T) {
t.Fatal(err)
}
task, _ := s.Task("t")
p, _ := json.Marshal(map[string]string{"handoff_ref": ref, "anchor_sha": "0123456789012345678901234567890123456789"})
p, _ := json.Marshal(map[string]any{"handoff_ref": ref, "anchor_sha": "0123456789012345678901234567890123456789", "harness_id": task.Lease.HarnessID, "lease_epoch": task.Lease.Epoch, "expected_version": task.Version})
if err := s.Append(domain.Event{ID: "release", Type: "TaskReleased", TaskID: "t", Version: task.Version + 1, Payload: p, Surface: string(authz.System)}); err != nil {
t.Fatal(err)
}
@@ -64,13 +65,13 @@ func TestRenewLeaseRequiresCurrentOwnerAndVersion(t *testing.T) {
t.Fatal(err)
}
before, _ := s.Task("t")
if _, err := s.RenewLease("t", "worker-b", before.Version, time.Hour); !errors.Is(err, domain.ErrConflict) {
if _, err := s.RenewLease("t", "worker-b", before.Lease.Epoch, before.Version, time.Hour); !errors.Is(err, domain.ErrConflict) {
t.Fatalf("other worker renewal = %v, want conflict", err)
}
if _, err := s.RenewLease("t", "worker-a", before.Version-1, time.Hour); !errors.Is(err, domain.ErrConflict) {
if _, err := s.RenewLease("t", "worker-a", before.Lease.Epoch, before.Version-1, time.Hour); !errors.Is(err, domain.ErrConflict) {
t.Fatalf("stale renewal = %v, want conflict", err)
}
e, err := s.RenewLease("t", "worker-a", before.Version, time.Hour)
e, err := s.RenewLease("t", "worker-a", before.Lease.Epoch, before.Version, time.Hour)
if err != nil {
t.Fatal(err)
}
@@ -78,11 +79,107 @@ func TestRenewLeaseRequiresCurrentOwnerAndVersion(t *testing.T) {
if e.Type != "TaskLeaseRenewed" || after.Version != before.Version+1 || after.Lease == nil || !after.Lease.Until.After(before.Lease.Until) {
t.Fatalf("renewal was not projected: before=%+v after=%+v event=%+v", before, after, e)
}
if _, err := s.RenewLease("t", "worker-a", before.Version, time.Hour); !errors.Is(err, domain.ErrConflict) {
if _, err := s.RenewLease("t", "worker-a", before.Lease.Epoch, before.Version, time.Hour); !errors.Is(err, domain.ErrConflict) {
t.Fatalf("replayed renewal = %v, want conflict", err)
}
}
func TestLeaseEpochFencesStaleOwnerLifecycleWrites(t *testing.T) {
s, err := Open(t.TempDir())
if err != nil {
t.Fatal(err)
}
if err := s.Append(domain.Event{ID: "create", Type: "TaskCreated", TaskID: "t", Version: 1, Payload: []byte(`{"source":"s","external_id":"epoch","project":"p"}`), Surface: string(authz.System)}); err != nil {
t.Fatal(err)
}
if _, err := s.Lease("t", "worker", time.Minute); err != nil {
t.Fatal(err)
}
first, _ := s.Task("t")
ref, err := s.PutArtifact([]byte("handoff"))
if err != nil {
t.Fatal(err)
}
staleRelease, _ := json.Marshal(map[string]any{"handoff_ref": ref, "anchor_sha": strings.Repeat("a", 40), "harness_id": "worker", "lease_epoch": "stale", "expected_version": first.Version})
if err := s.Append(domain.Event{ID: "stale-release", Type: "TaskReleased", TaskID: "t", Version: first.Version + 1, Payload: staleRelease, Surface: string(authz.System)}); !errors.Is(err, domain.ErrConflict) {
t.Fatalf("stale release = %v, want conflict", err)
}
release, _ := json.Marshal(map[string]any{"handoff_ref": ref, "anchor_sha": strings.Repeat("a", 40), "harness_id": "worker", "lease_epoch": first.Lease.Epoch, "expected_version": first.Version})
if err := s.Append(domain.Event{ID: "release", Type: "TaskReleased", TaskID: "t", Version: first.Version + 1, Payload: release, Surface: string(authz.System)}); err != nil {
t.Fatal(err)
}
if _, err := s.Lease("t", "worker", time.Minute); err != nil {
t.Fatal(err)
}
second, _ := s.Task("t")
if second.Lease.Epoch == first.Lease.Epoch {
t.Fatal("re-lease reused fencing epoch")
}
report, err := s.PutArtifact([]byte("report"))
if err != nil {
t.Fatal(err)
}
staleComplete, _ := json.Marshal(map[string]any{"report_ref": report, "receipt": map[string]any{"consumed": 1}, "harness_id": "worker", "lease_epoch": first.Lease.Epoch, "expected_version": second.Version})
if err := s.Append(domain.Event{ID: "stale-complete", Type: "TaskCompleted", TaskID: "t", Version: second.Version + 1, Payload: staleComplete, Surface: string(authz.System)}); !errors.Is(err, domain.ErrConflict) {
t.Fatalf("stale completion = %v, want conflict", err)
}
}
func TestOpenRebuildsOnlyFromLogAndIgnoresCorruptSnapshot(t *testing.T) {
dir := t.TempDir()
s, err := Open(dir)
if err != nil {
t.Fatal(err)
}
if err := s.Append(created("create")); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(dir, "snapshot.json"), []byte(`not json`), 0600); err != nil {
t.Fatal(err)
}
restarted, err := Open(dir)
if err != nil {
t.Fatalf("corrupt disposable snapshot prevented log recovery: %v", err)
}
if task, ok := restarted.Task("task-1"); !ok || task.State != domain.StateQueued {
t.Fatalf("log projection = %#v, present=%v", task, ok)
}
}
func TestReplayLegacyLeaseDerivesNonRenewableFence(t *testing.T) {
dir := t.TempDir()
until := time.Now().Add(time.Hour).UnixNano()
events := []domain.Event{
{SchemaVersion: 2, Seq: 1, ID: "create", Type: "TaskCreated", TaskID: "t", Version: 1, Payload: []byte(`{"source":"s","external_id":"legacy","project":"p"}`), Surface: string(authz.System)},
{SchemaVersion: 2, Seq: 2, ID: "lease", Type: "TaskLeased", TaskID: "t", Version: 2, Payload: []byte(`{"harness_id":"worker","until_ns":` + fmt.Sprint(until) + `,"expected_version":1}`), Surface: string(authz.System)},
}
f, err := os.Create(filepath.Join(dir, "events.jsonl"))
if err != nil {
t.Fatal(err)
}
for _, e := range events {
b, _ := json.Marshal(e)
if _, err := f.Write(append(b, '\n')); err != nil {
_ = f.Close()
t.Fatal(err)
}
}
if err := f.Close(); err != nil {
t.Fatal(err)
}
s, err := Open(dir)
if err != nil {
t.Fatal(err)
}
task, _ := s.Task("t")
if task.Lease == nil || task.Lease.Epoch != "legacy:lease" {
t.Fatalf("legacy lease fence = %#v", task.Lease)
}
if _, err := s.RenewLease("t", "worker", "", task.Version, time.Hour); !errors.Is(err, domain.ErrConflict) {
t.Fatalf("legacy lease renewed without derived fence: %v", err)
}
}
func TestBlockedTaskProjectsStructuredDiagnosisAndLegacyFallback(t *testing.T) {
s, err := Open(t.TempDir())
if err != nil {
@@ -407,7 +504,8 @@ func TestReleaseTransactionSurvivesReLeaseUntilMatchingPickup(t *testing.T) {
t.Fatal(err)
}
anchor := strings.Repeat("a", 40)
p, _ := json.Marshal(map[string]any{"handoff_ref": ref, "anchor_sha": anchor, "transaction_id": "tx-1", "expected_version": leased.Version})
first, _ := s.Task("task-1")
p, _ := json.Marshal(map[string]any{"handoff_ref": ref, "anchor_sha": anchor, "transaction_id": "tx-1", "harness_id": first.Lease.HarnessID, "lease_epoch": first.Lease.Epoch, "expected_version": leased.Version})
if err := s.Append(domain.Event{ID: "release", Type: "TaskReleased", TaskID: "task-1", Version: leased.Version + 1, Payload: p, Surface: string(authz.System)}); err != nil {
t.Fatal(err)
}
@@ -418,7 +516,7 @@ func TestReleaseTransactionSurvivesReLeaseUntilMatchingPickup(t *testing.T) {
if task.ReleaseTransaction != "tx-1" || task.ReleaseAnchor != anchor || task.HandoffRef != ref {
t.Fatalf("re-lease lost transaction: %+v", task)
}
p, _ = json.Marshal(map[string]any{"transaction_id": "tx-1", "handoff_ref": ref, "anchor_sha": anchor, "harness_id": "successor", "lease_version": task.Version, "expected_version": task.Version})
p, _ = json.Marshal(map[string]any{"transaction_id": "tx-1", "handoff_ref": ref, "anchor_sha": anchor, "harness_id": "successor", "lease_epoch": task.Lease.Epoch, "lease_version": task.Version, "expected_version": task.Version})
if err := s.Append(domain.Event{ID: "pickup", Type: "TaskPickupValidated", TaskID: "task-1", Version: task.Version + 1, Payload: p, Surface: string(authz.System)}); err != nil {
t.Fatal(err)
}