Files
orchestra/internal/store/store_test.go
T
kami 3c7cf95d8c Settle a release transaction deterministically in every case
F60, and the general rule F58 and F59 were reaching for one case at a
time: a transaction must settle or be abandoned deterministically, and
must never spin on an answer that cannot change.

Terminal now means failed or completed. Both drop the transaction and
free the session; nothing will ever lease either task again.

Blocked keeps the transaction, because a reopen returns the task to the
queue and that exact owner can still commit. TaskBlocked therefore
retains the ending epoch the way TaskReleased already did, or the
late-handoff path would have nothing to fence against after the reopen.

A refusal parks the commit instead of retrying every five seconds. It
is the coordinator's answer about who owns the task, so it stays true
until an event about that task arrives, and any such event un-parks it.
A reopen arrives as TaskCorrected, so the rule cannot be a list of
event types. Backoff runs 30s to a 5 minute cap.

A transport failure is not an answer and keeps retrying at once. That
distinction is the whole reason the park keys on a 4xx StatusError
rather than on any error at all.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CVbaKucEYBjMqVeUgJUsc1
2026-08-28 23:52:18 +04:00

726 lines
28 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)
}
}