75396963ef
The clearing rule sat below the reducer switch, so it ran for every event rather than for the correction it was written for. TaskReleased found the task unblocked and erased the contradiction, which is the rotation the reopen itself causes: the planning session launched one lease later and was told nothing again. Scoped to TaskCorrected, and the store test walks the real sequence (mismatch, reopen, release, lease) rather than reading the projection at the moment it is written. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01CVbaKucEYBjMqVeUgJUsc1
795 lines
31 KiB
Go
795 lines
31 KiB
Go
package store
|
|
|
|
import (
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"path/filepath"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
|
|
"orchestra/internal/authz"
|
|
"orchestra/internal/domain"
|
|
)
|
|
|
|
func created(id string) domain.Event {
|
|
b, _ := json.Marshal(map[string]any{"source": "jsonl", "external_id": "42", "project": "demo", "capability": []string{"mechanical"}})
|
|
return domain.Event{ID: id, Type: "TaskCreated", TaskID: "task-1", Version: 1, Payload: b, Surface: string(authz.System)}
|
|
}
|
|
|
|
func TestLeaseCarriesReleasedHandoffRef(t *testing.T) {
|
|
s, err := Open(t.TempDir())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
b := []byte(`{"source":"s","external_id":"x","project":"p"}`)
|
|
if err := s.Append(domain.Event{ID: "create", Type: "TaskCreated", TaskID: "t", Version: 1, Payload: b, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := s.Lease("t", "h", time.Minute); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
ref, err := s.PutArtifact([]byte("handoff"))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
task, _ := s.Task("t")
|
|
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)
|
|
}
|
|
e, err := s.Lease("t", "h", time.Minute)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var got map[string]any
|
|
if err := json.Unmarshal(e.Payload, &got); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if got["handoff_ref"] != ref {
|
|
t.Fatalf("handoff_ref=%v want %s", got["handoff_ref"], ref)
|
|
}
|
|
}
|
|
|
|
func TestRenewLeaseRequiresCurrentOwnerAndVersion(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":"renew","project":"p"}`), Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := s.Lease("t", "worker-a", time.Minute); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
before, _ := s.Task("t")
|
|
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.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.Lease.Epoch, before.Version, time.Hour)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
after, _ := s.Task("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.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 TestNeedsAttentionRetainsFencedLeaseForLateCompletion(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":"attention","project":"p"}`), Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := s.Lease("t", "worker", time.Hour); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
leased, _ := s.Task("t")
|
|
attention, _ := json.Marshal(map[string]any{"blocker": "prompt response uncertain", "block_reason": "lease_failure", "harness_id": leased.Lease.HarnessID, "lease_epoch": leased.Lease.Epoch, "expected_version": leased.Version})
|
|
if err := s.Append(domain.Event{ID: "attention", Type: "TaskNeedsAttention", TaskID: "t", Version: leased.Version + 1, Payload: attention, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
current, _ := s.Task("t")
|
|
if current.State != domain.StateNeedsAttention || current.Lease == nil || current.Lease.Epoch != leased.Lease.Epoch {
|
|
t.Fatalf("attention revoked or replaced lease: %+v", current)
|
|
}
|
|
report, err := s.PutArtifact([]byte("late completion"))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
completion, _ := json.Marshal(map[string]any{"report_ref": report, "receipt": map[string]any{"consumed": 1}, "harness_id": current.Lease.HarnessID, "lease_epoch": current.Lease.Epoch, "expected_version": current.Version})
|
|
if err := s.Append(domain.Event{ID: "late-complete", Type: "TaskCompleted", TaskID: "t", Version: current.Version + 1, Payload: completion, Surface: string(authz.System)}); err != nil {
|
|
t.Fatalf("late completion from retained owner: %v", err)
|
|
}
|
|
completed, _ := s.Task("t")
|
|
if completed.State != domain.StateCompleted || completed.Lease != nil {
|
|
t.Fatalf("late completion did not settle task: %+v", completed)
|
|
}
|
|
}
|
|
|
|
func TestReclaimPersistsAttemptAndBackoffAcrossReopen(t *testing.T) {
|
|
dir := t.TempDir()
|
|
s, err := Open(dir)
|
|
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":"retry","project":"p"}`), Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := s.Lease("t", "worker", time.Hour); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
task, _ := s.Task("t")
|
|
at := time.Date(2026, 7, 30, 12, 0, 0, 0, time.UTC)
|
|
p, _ := json.Marshal(map[string]any{"reason": "lease_expired", "failure_class": "worker_lost", "harness_id": task.Lease.HarnessID, "lease_epoch": task.Lease.Epoch, "expected_version": task.Version})
|
|
if err := s.Append(domain.Event{ID: "reclaim", Type: "TaskReleased", TaskID: "t", Version: task.Version + 1, At: at, Payload: p, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
got, _ := s.Task("t")
|
|
if got.Attempt != 1 || got.FailureClass != "worker_lost" || !got.NextRetryAt.Equal(at.Add(time.Minute)) {
|
|
t.Fatalf("reclaim projection=%+v", got)
|
|
}
|
|
reopened, err := Open(dir)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
got, _ = reopened.Task("t")
|
|
if got.Attempt != 1 || !got.NextRetryAt.Equal(at.Add(time.Minute)) {
|
|
t.Fatalf("reopen lost retry state: %+v", got)
|
|
}
|
|
}
|
|
|
|
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 TestSchedulingAndQuotaIndexesReplayFromEvents(t *testing.T) {
|
|
dir := t.TempDir()
|
|
s, err := Open(dir)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(created("indexed")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := s.Lease("task-1", "h1", time.Hour); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
snapshot := s.SchedulingSnapshot()
|
|
if got := snapshot.ActiveLeases["h1"]; got != 1 {
|
|
t.Fatalf("active lease index=%d, want 1", got)
|
|
}
|
|
now := time.Now().UTC()
|
|
report := func(at time.Time, consumed float64, known bool) {
|
|
payload, _ := json.Marshal(map[string]any{"harness_id": "h1", "consumed": consumed, "known": known})
|
|
if err := s.Append(domain.Event{ID: domain.NewID(), Type: "QuotaReported", TaskID: "quota", Version: 1, At: at, Payload: payload, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
// Deliberately append out of timestamp order: the index must answer the
|
|
// rolling query by receipt time, not event sequence.
|
|
report(now, 4, true)
|
|
report(now.Add(-2*time.Hour), 3, true)
|
|
report(now.Add(-time.Hour), 1, false)
|
|
if used, known := s.QuotaSince("h1", now.Add(-90*time.Minute)); used != 5 || known {
|
|
t.Fatalf("indexed quota=(%v,%v), want (5,false)", used, known)
|
|
}
|
|
if used, known := s.QuotaSince("h1", now.Add(-30*time.Minute)); used != 4 || !known {
|
|
t.Fatalf("recent indexed quota=(%v,%v), want (4,true)", used, known)
|
|
}
|
|
reopened, err := Open(dir)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if got := reopened.SchedulingSnapshot().ActiveLeases["h1"]; got != 1 {
|
|
t.Fatalf("replayed active lease index=%d, want 1", got)
|
|
}
|
|
if used, known := reopened.QuotaSince("h1", now.Add(-90*time.Minute)); used != 5 || known {
|
|
t.Fatalf("replayed indexed quota=(%v,%v), want (5,false)", used, known)
|
|
}
|
|
}
|
|
|
|
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 {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(domain.Event{ID: "create", Type: "TaskCreated", TaskID: "t", Version: 1, Payload: []byte(`{"source":"s","external_id":"blocked","project":"p"}`), Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
payload := []byte(`{"blocker":"worker is offline","block_reason":"worker_offline","pane_id":"p1","harness_id":"w1","pane_state":"unreachable"}`)
|
|
if err := s.Append(domain.Event{ID: "blocked", Type: "TaskBlocked", TaskID: "t", Version: 2, Payload: payload, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
task, _ := s.Task("t")
|
|
if task.BlockReason != domain.BlockReasonWorkerOffline || task.LastPaneID != "p1" || task.PaneState != "unreachable" {
|
|
t.Fatalf("blocked diagnosis was not projected: %+v", task)
|
|
}
|
|
if got := domain.InferBlockReason("handoff validation failed"); got != domain.BlockReasonHandoffValidation {
|
|
t.Fatalf("legacy fallback=%q", got)
|
|
}
|
|
}
|
|
|
|
func TestTerminalSessionEvidenceSurvivesSessionCleanup(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":"evidence","project":"p"}`), Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
ref, err := s.PutArtifact([]byte("report"))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
payload := []byte(`{"report_ref":"` + ref + `","receipt":{"source":"worker"},"session_evidence":{"pane_id":"w:p1","harness_id":"worker-1","pane_state":"open","source":"worker","captured_at":"2026-07-29T12:00:00Z","checked_at":"2026-07-29T12:00:01Z"}}`)
|
|
if err := s.Append(domain.Event{ID: "complete", Type: "TaskCompleted", TaskID: "t", Version: 2, Payload: payload, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
task, _ := s.Task("t")
|
|
if task.LastSession.PaneID != "w:p1" || task.LastSession.Source != "worker" || task.LastSession.CapturedAt.IsZero() || task.LastHarness != "worker-1" {
|
|
t.Fatalf("terminal session evidence was lost: %+v", task)
|
|
}
|
|
}
|
|
|
|
func TestAppendReplayAndDeduplicate(t *testing.T) {
|
|
dir := t.TempDir()
|
|
s, err := Open(dir)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(created("e1")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(created("e2")); !errors.Is(err, domain.ErrDuplicate) {
|
|
t.Fatalf("expected ErrDuplicate, got %v", err)
|
|
}
|
|
if got := len(s.Events(0)); got != 1 {
|
|
t.Fatalf("duplicate ingest appended %d events", got)
|
|
}
|
|
ref, err := s.PutArtifact([]byte("report"))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
completion, _ := json.Marshal(map[string]any{"report_ref": ref, "receipt": map[string]any{"harness_id": "h", "consumed": 1}})
|
|
if err := s.Append(domain.Event{Type: "TaskCompleted", TaskID: "task-1", Version: 2, Payload: completion, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(domain.Event{Type: "TaskReleased", TaskID: "task-1", Version: 2, Surface: string(authz.System), Payload: json.RawMessage(`{"handoff_ref":"` + ref + `","anchor_sha":"0123456789012345678901234567890123456789"}`)}); err != domain.ErrConflict {
|
|
t.Fatalf("expected conflict, got %v", err)
|
|
}
|
|
s2, err := Open(dir)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if got := s2.Tasks()[0].State; got != domain.StateCompleted {
|
|
t.Fatalf("replay state = %s", got)
|
|
}
|
|
if _, err := os.Stat(filepath.Join(dir, "snapshot.json")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
// TestTaskAmendedAppliesAllFields guards S7: the projection previously only
|
|
// applied "title" from a TaskAmended payload, silently discarding due,
|
|
// description, and inherent_priority amendments even though they were
|
|
// accepted and logged.
|
|
func TestTaskAmendedAppliesAllFields(t *testing.T) {
|
|
dir := t.TempDir()
|
|
s, err := Open(dir)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(created("e1")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
amend, _ := json.Marshal(map[string]any{
|
|
"title": "new title",
|
|
"description": "new description",
|
|
"inherent_priority": 5.0,
|
|
"due": "2026-08-01T00:00:00Z",
|
|
})
|
|
if err := s.Append(domain.Event{Type: "TaskAmended", TaskID: "task-1", Version: 2, Payload: amend, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
tk := s.Tasks()[0]
|
|
if tk.Title != "new title" || tk.Description != "new description" || tk.InherentPriority != 5 {
|
|
t.Fatalf("unexpected task after amendment: %+v", tk)
|
|
}
|
|
if tk.Due == nil || !tk.Due.Equal(time.Date(2026, 8, 1, 0, 0, 0, 0, time.UTC)) {
|
|
t.Fatalf("unexpected due after amendment: %+v", tk.Due)
|
|
}
|
|
}
|
|
|
|
// TestTaskCorrected guards S8: the spec (§3.1) requires corrections to be
|
|
// compensating events appended on top of a wrong one, never an edit of the
|
|
// log — a wrongly-emitted terminal state (e.g. a mistaken TaskFailed) must be
|
|
// repairable by appending a new event that references the one it corrects,
|
|
// with both surviving in the replayed log.
|
|
func TestTaskCorrected(t *testing.T) {
|
|
dir := t.TempDir()
|
|
s, err := Open(dir)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(created("e1")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
failPayload, _ := json.Marshal(map[string]any{"reason": "mistaken failure"})
|
|
failEvt := domain.Event{ID: "e2", Type: "TaskFailed", TaskID: "task-1", Version: 2, Payload: failPayload, Surface: string(authz.System)}
|
|
if err := s.Append(failEvt); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if tk := s.Tasks()[0]; tk.State != domain.StateFailed {
|
|
t.Fatalf("expected failed, got %s", tk.State)
|
|
}
|
|
|
|
// Referencing an unknown event is rejected.
|
|
badCorrection, _ := json.Marshal(map[string]any{"corrects": "does-not-exist", "state": "queued"})
|
|
if err := s.Append(domain.Event{Type: "TaskCorrected", TaskID: "task-1", Version: 3, Payload: badCorrection, Surface: string(authz.System)}); !errors.Is(err, domain.ErrInvalid) {
|
|
t.Fatalf("expected ErrInvalid for unknown corrects target, got %v", err)
|
|
}
|
|
|
|
correction, _ := json.Marshal(map[string]any{"corrects": "e2", "state": "queued"})
|
|
if err := s.Append(domain.Event{ID: "e3", Type: "TaskCorrected", TaskID: "task-1", Version: 3, Payload: correction, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
tk := s.Tasks()[0]
|
|
if tk.State != domain.StateQueued {
|
|
t.Fatalf("expected correction to restore queued state, got %s", tk.State)
|
|
}
|
|
|
|
// Replay from disk still applies both the wrong event and its correction
|
|
// (the snapshot only elides already-applied events from Events(), never
|
|
// from the durable log itself).
|
|
s2, err := Open(dir)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if tk2 := s2.Tasks()[0]; tk2.State != domain.StateQueued {
|
|
t.Fatalf("expected replayed state queued, got %s", tk2.State)
|
|
}
|
|
}
|
|
|
|
// TestLeaseAndExpireEventIDsAreUnique guards S5: Event.ID was set to the
|
|
// task id in both Lease and ExpireLeases, so every lease of the same task
|
|
// produced a TaskLeased/TaskReleased event with a colliding ID — unsound
|
|
// for ApplyAdvisory or any future ID-based lookup.
|
|
func TestLeaseAndExpireEventIDsAreUnique(t *testing.T) {
|
|
s, err := Open(t.TempDir())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(created("e1")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
task := s.Tasks()[0]
|
|
leaseEvt, err := s.Lease(task.ID, "h1", time.Millisecond)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if leaseEvt.ID == task.ID || leaseEvt.ID == "" {
|
|
t.Fatalf("lease event ID %q collides with task ID %q", leaseEvt.ID, task.ID)
|
|
}
|
|
time.Sleep(2 * time.Millisecond)
|
|
expired, err := s.ExpireLeases(time.Now())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(expired) != 1 {
|
|
t.Fatalf("expected 1 expiry, got %d", len(expired))
|
|
}
|
|
if expired[0].ID == task.ID || expired[0].ID == leaseEvt.ID || expired[0].ID == "" {
|
|
t.Fatalf("expiry event ID %q collides", expired[0].ID)
|
|
}
|
|
}
|
|
|
|
// TestTaskBySourceResolvesDuplicate covers the S6 fix: a caller that gets
|
|
// ErrDuplicate from Append must be able to look up the already-ingested
|
|
// task by its dedup key instead of guessing at "the last event in the log".
|
|
func TestTaskBySourceResolvesDuplicate(t *testing.T) {
|
|
s, err := Open(t.TempDir())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(created("e1")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
want := s.Tasks()[0]
|
|
if err := s.Append(created("e2")); !errors.Is(err, domain.ErrDuplicate) {
|
|
t.Fatalf("expected ErrDuplicate, got %v", err)
|
|
}
|
|
got, ok := s.TaskBySource("jsonl", "42")
|
|
if !ok || got.ID != want.ID {
|
|
t.Fatalf("TaskBySource = %+v, ok=%v, want %+v", got, ok, want)
|
|
}
|
|
if _, ok := s.TaskBySource("jsonl", "does-not-exist"); ok {
|
|
t.Fatal("expected not found")
|
|
}
|
|
}
|
|
|
|
func TestArtifactIsContentAddressed(t *testing.T) {
|
|
s, err := Open(t.TempDir())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
h1, err := s.PutArtifact([]byte("proof"))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
h2, err := s.PutArtifact([]byte("proof"))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if h1 != h2 {
|
|
t.Fatal("same artifact received different hashes")
|
|
}
|
|
if _, err := os.Stat(filepath.Join(s.cas, h1)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
func TestLifecycleEventsRequireEvidence(t *testing.T) {
|
|
cases := []struct {
|
|
name string
|
|
typ string
|
|
body string
|
|
}{
|
|
{"release", "TaskReleased", `{}`},
|
|
{"complete", "TaskCompleted", `{}`},
|
|
{"block", "TaskBlocked", `{}`},
|
|
}
|
|
for _, tc := range cases {
|
|
t.Run(tc.name, func(t *testing.T) {
|
|
err := domain.ValidateEvent(domain.Event{Type: tc.typ, TaskID: "task-1", Version: 1, Payload: json.RawMessage(tc.body), Surface: string(authz.System)})
|
|
if err == nil {
|
|
t.Fatal("expected lifecycle evidence validation error")
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
// TestAppendEnforcesAuthorizationAtTheBus proves authorization is checked
|
|
// once at the append boundary (spec §7.1/§1.4), not only in HTTP handlers:
|
|
// a caller writing to the store directly with a notify-only surface, or with
|
|
// no declared surface at all, is rejected exactly like an HTTP request would
|
|
// be — there is no in-process bypass for the router, coordinator, or a
|
|
// provider adapter that forgets to declare who it's acting as.
|
|
func TestAppendEnforcesAuthorizationAtTheBus(t *testing.T) {
|
|
s, err := Open(t.TempDir())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
b, _ := json.Marshal(map[string]any{"source": "jsonl", "external_id": "1", "project": "demo"})
|
|
|
|
// A notify-only surface (e.g. Telegram) must never be able to create a
|
|
// task by calling the store directly, even though it bypasses HTTP.
|
|
if err := s.Append(domain.Event{ID: "e1", Type: "TaskCreated", TaskID: "task-1", Version: 1, Payload: b, Surface: string(authz.Telegram)}); err == nil {
|
|
t.Fatal("notify-only surface created a task via direct store access")
|
|
}
|
|
|
|
// An internal producer that forgets to declare a surface is rejected,
|
|
// not silently trusted as the plane.
|
|
if err := s.Append(domain.Event{ID: "e2", Type: "TaskCreated", TaskID: "task-1", Version: 1, Payload: b}); err == nil {
|
|
t.Fatal("event with no declared surface was accepted")
|
|
}
|
|
|
|
// The plane (router/coordinator/provider) authorizes as System and
|
|
// succeeds.
|
|
if err := s.Append(domain.Event{ID: "e3", Type: "TaskCreated", TaskID: "task-1", Version: 1, Payload: b, Surface: string(authz.System)}); err != nil {
|
|
t.Fatalf("system surface rejected: %v", err)
|
|
}
|
|
}
|
|
|
|
func TestExpectedVersionIsCheckedForEveryWriter(t *testing.T) {
|
|
s, err := Open(t.TempDir())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(created("create")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
p := json.RawMessage(`{"reason":"rotate","expected_version":0}`)
|
|
err = s.Append(domain.Event{Type: "TaskReleased", TaskID: "task-1", Version: 2, Payload: p, Surface: string(authz.System)})
|
|
if err != domain.ErrConflict {
|
|
t.Fatalf("expected CAS conflict, got %v", err)
|
|
}
|
|
}
|
|
|
|
func TestReleaseTransactionSurvivesReLeaseUntilMatchingPickup(t *testing.T) {
|
|
s, err := Open(t.TempDir())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(created("create")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
leased, err := s.Lease("task-1", "predecessor", time.Minute)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
ref, err := s.PutArtifact([]byte("handoff"))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
anchor := strings.Repeat("a", 40)
|
|
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)
|
|
}
|
|
if _, err := s.Lease("task-1", "successor", time.Minute); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
task, _ := s.Task("task-1")
|
|
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_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)
|
|
}
|
|
task, _ = s.Task("task-1")
|
|
if task.PickupTransaction != "tx-1" || task.PickupLeaseVersion != 4 {
|
|
t.Fatalf("pickup not bound to transaction/epoch: %+v", task)
|
|
}
|
|
}
|
|
|
|
// An empty window is observable zero consumption, not missing data. Reporting
|
|
// it as unknown made the router refuse every harness that had a configured
|
|
// quota limit and no receipt history, which no first lease could ever produce.
|
|
func TestQuotaSinceReportsEmptyWindowAsKnownZero(t *testing.T) {
|
|
s, err := Open(t.TempDir())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
now := time.Now().UTC()
|
|
if used, known := s.QuotaSince("fresh", now.Add(-fiveHours)); used != 0 || !known {
|
|
t.Fatalf("never-reported harness quota=(%v,%v), want (0,true)", used, known)
|
|
}
|
|
payload, _ := json.Marshal(map[string]any{"harness_id": "h1", "consumed": 7.0, "known": true})
|
|
if err := s.Append(domain.Event{ID: domain.NewID(), Type: "QuotaReported", TaskID: "quota", Version: 1, At: now.Add(-24 * time.Hour), Payload: payload, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if used, known := s.QuotaSince("h1", now.Add(-fiveHours)); used != 0 || !known {
|
|
t.Fatalf("receipts only outside the window quota=(%v,%v), want (0,true)", used, known)
|
|
}
|
|
if used, known := s.QuotaSince("h1", now.Add(-48*time.Hour)); used != 7 || !known {
|
|
t.Fatalf("receipts inside the window quota=(%v,%v), want (7,true)", used, known)
|
|
}
|
|
}
|
|
|
|
const fiveHours = 5 * time.Hour
|
|
|
|
// TestExpiryRetainsLeaseEpoch covers the late-handoff fence: a worker that
|
|
// pushed its release anchor and then lost the lease can only commit if the
|
|
// projection still knows which epoch just ended.
|
|
func TestExpiryRetainsLeaseEpoch(t *testing.T) {
|
|
s, err := Open(t.TempDir())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(created("e1")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
id := s.Tasks()[0].ID
|
|
if _, err := s.Lease(id, "h1", time.Millisecond); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
leased, _ := s.Task(id)
|
|
epoch := leased.Lease.Epoch
|
|
time.Sleep(2 * time.Millisecond)
|
|
if _, err := s.ExpireLease(id, time.Now()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
after, _ := s.Task(id)
|
|
if after.Lease != nil {
|
|
t.Fatal("expired task still holds a lease")
|
|
}
|
|
if after.LastLeaseEpoch != epoch || epoch == "" {
|
|
t.Fatalf("last lease epoch %q, want %q", after.LastLeaseEpoch, epoch)
|
|
}
|
|
}
|
|
|
|
// A worker can hold a pushed anchor whose commit was refused when an operator
|
|
// blocks the task. A reopen returns it to the queue, and the late-handoff path
|
|
// can only accept that exact owner if the epoch that ended is still recorded.
|
|
func TestBlockRetainsLeaseEpochForALaterReopen(t *testing.T) {
|
|
s, err := Open(t.TempDir())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Append(created("e1")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
id := s.Tasks()[0].ID
|
|
if _, err := s.Lease(id, "h1", time.Minute); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
leased, _ := s.Task(id)
|
|
epoch := leased.Lease.Epoch
|
|
p, _ := json.Marshal(map[string]any{"blocker": "parked by the operator", "harness_id": "h1", "lease_epoch": epoch})
|
|
if err := s.Append(domain.Event{ID: domain.NewID(), Type: "TaskBlocked", TaskID: id, Version: leased.Version + 1, Payload: p, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
after, _ := s.Task(id)
|
|
if after.State != domain.StateBlocked || after.Lease != nil {
|
|
t.Fatalf("expected a blocked unleased task, got %s lease=%v", after.State, after.Lease)
|
|
}
|
|
if after.LastLeaseEpoch != epoch || epoch == "" {
|
|
t.Fatalf("last lease epoch %q, want %q", after.LastLeaseEpoch, epoch)
|
|
}
|
|
}
|
|
|
|
// F66's projection has to survive the rotation the reopen causes. The session
|
|
// that reported the contradiction hands off, a successor leases, and only then
|
|
// is the planning context rendered. Live on run 21, the field was gone by
|
|
// then: the reopen recorded it and the rotation lost it.
|
|
func TestTheContradictionSurvivesTheRotationItCauses(t *testing.T) {
|
|
s, err := Open(t.TempDir())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
b := []byte(`{"source":"s","external_id":"x","project":"p"}`)
|
|
if err := s.Append(domain.Event{ID: "create", Type: "TaskCreated", TaskID: "t", Version: 1, Payload: b, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := s.Lease("t", "h", time.Minute); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
// Frame to implement, the shortest legal route to the phase a mismatch
|
|
// can be reported from.
|
|
task, _ := s.Task("t")
|
|
toImplement, _ := json.Marshal(map[string]any{"phase": "implement", "from": "frame"})
|
|
if err := s.Append(domain.Event{ID: "impl", Type: domain.EventWorkPhaseChanged, TaskID: "t", Version: task.Version + 1, Payload: toImplement, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
task, _ = s.Task("t")
|
|
m, _ := json.Marshal(map[string]any{
|
|
"plan_ref": "plan-a", "phase_id": "phase-2",
|
|
"at_sha": "0123456789012345678901234567890123456789",
|
|
"observed": "the body is assembled inline",
|
|
"contradicts": "the plan says one helper returns it",
|
|
"harness_id": task.Lease.HarnessID,
|
|
"lease_epoch": task.Lease.Epoch,
|
|
"requested_action": "replan",
|
|
})
|
|
if err := s.Append(domain.Event{ID: "mismatch", Type: domain.EventPlanMismatchRecorded, TaskID: "t", Version: task.Version + 1, Payload: m, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if got, _ := s.Task("t"); got.PlanMismatch == nil {
|
|
t.Fatal("the contradiction was not projected at all")
|
|
}
|
|
|
|
// The reopen, then the rotation it causes.
|
|
task, _ = s.Task("t")
|
|
ph, _ := json.Marshal(map[string]any{"phase": "plan", "from": "implement", "reopen": domain.EventPlanMismatchRecorded, "reopen_phase_id": "phase-2"})
|
|
if err := s.Append(domain.Event{ID: "reopen", Type: domain.EventWorkPhaseChanged, TaskID: "t", Version: task.Version + 1, Payload: ph, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
ref, err := s.PutArtifact([]byte("handoff"))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
task, _ = s.Task("t")
|
|
rel, _ := 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: rel, Surface: string(authz.System)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := s.Lease("t", "h", time.Minute); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
// This is the moment the planning context is rendered.
|
|
got, _ := s.Task("t")
|
|
if got.PlanMismatch == nil {
|
|
t.Fatal("the planner is convened to settle a contradiction it is no longer told about")
|
|
}
|
|
if got.PlanMismatch.PhaseID != "phase-2" {
|
|
t.Fatalf("the contradiction changed across the rotation: %+v", got.PlanMismatch)
|
|
}
|
|
}
|