iop/apps/agent/internal/projectlog/store_test.go

2693 lines
77 KiB
Go

package projectlog
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"math"
"os"
"path/filepath"
"strings"
"sync"
"testing"
"time"
"iop/packages/go/agentstate"
"iop/packages/go/agenttask"
)
func newTestStore(t *testing.T) *Store {
t.Helper()
root := t.TempDir()
statePath := filepath.Join(t.TempDir(), "manager.json")
state, err := agentstate.NewStore(statePath)
if err != nil {
t.Fatalf("agentstate.NewStore: %v", err)
}
store, err := NewStore(state, root, "proj-001", "ws-001")
if err != nil {
t.Fatalf("NewStore: %v", err)
}
return store
}
func appendTerminalRecords(t *testing.T, store *Store, n int) {
t.Helper()
j, _, _, err := store.loadJournal(context.Background())
if err != nil {
t.Fatalf("loadJournal: %v", err)
}
firstSequence := j.NextSequence
for i := 0; i < n; i++ {
rec := validRecord()
rec.Sequence = 0
rec.RecordID = fmt.Sprintf("rec-%03d", firstSequence+uint64(i))
if i == n-1 {
rec.EventType = agenttask.EventCompleted
rec.State = agenttask.WorkStateCompleted
rec.StateRevision = agenttask.StateRevision("rev-terminal")
rec.Terminal = true
}
if _, err := store.AppendRecord(context.Background(), rec); err != nil {
t.Fatalf("AppendRecord %d: %v", i, err)
}
}
}
func archiveFilesExist(t *testing.T, store *Store, ordinal uint64) {
t.Helper()
archiveDir := store.archiveDir()
for _, name := range []string{
fmt.Sprintf("%06d.intent", ordinal),
fmt.Sprintf("%06d.jsonl", ordinal),
fmt.Sprintf("%06d.timeline.jsonl", ordinal),
fmt.Sprintf("%06d.manifest.json", ordinal),
} {
path := filepath.Join(archiveDir, name)
if _, err := os.Stat(path); err != nil {
t.Errorf("missing archive file %s: %v", name, err)
}
}
}
func journalState(t *testing.T, store *Store) ([]byte, string) {
t.Helper()
payload, revision, found, err := store.state.LoadIntegrationRecord(context.Background(), store.key)
if err != nil {
t.Fatalf("LoadIntegrationRecord: %v", err)
}
if !found {
t.Fatal("expected journal state to exist")
}
return payload, revision
}
func requireJournalState(t *testing.T, store *Store, wantPayload []byte, wantRevision string) {
t.Helper()
gotPayload, gotRevision := journalState(t, store)
if gotRevision != wantRevision {
t.Fatalf("journal revision changed: got %q, want %q", gotRevision, wantRevision)
}
if !bytes.Equal(gotPayload, wantPayload) {
t.Fatal("journal payload changed")
}
}
func archiveArtifactSet(t *testing.T, store *Store, ordinal uint64) map[string][]byte {
t.Helper()
artifacts, err := readArchiveArtifactSet(store, ordinal)
if err != nil {
t.Fatal(err)
}
return artifacts
}
func readArchiveArtifactSet(store *Store, ordinal uint64) (map[string][]byte, error) {
paths := map[string]string{
"intent": store.intentPath(ordinal),
"jsonl": store.jsonlPath(ordinal),
"timeline": store.timelinePath(ordinal),
"manifest": store.manifestPath(ordinal),
}
artifacts := make(map[string][]byte, len(paths))
for name, path := range paths {
payload, err := os.ReadFile(path)
if err != nil {
return nil, fmt.Errorf("read %s artifact: %w", name, err)
}
artifacts[name] = payload
}
return artifacts, nil
}
func requireArchiveArtifactSet(t *testing.T, store *Store, ordinal uint64, want map[string][]byte) {
t.Helper()
got := archiveArtifactSet(t, store, ordinal)
for name, wantPayload := range want {
if !bytes.Equal(got[name], wantPayload) {
t.Fatalf("%s artifact changed", name)
}
}
}
func TestStoreAppendAndReplay(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 3)
records, err := store.ReplayRecords(context.Background())
if err != nil {
t.Fatalf("ReplayRecords: %v", err)
}
if len(records) != 3 {
t.Fatalf("expected 3 records, got %d", len(records))
}
for i, rec := range records {
if rec.Sequence != uint64(i+1) {
t.Fatalf("record %d sequence = %d, want %d", i, rec.Sequence, i+1)
}
}
}
func TestStoreReplayEmpty(t *testing.T) {
store := newTestStore(t)
records, err := store.ReplayRecords(context.Background())
if err != nil {
t.Fatalf("ReplayRecords: %v", err)
}
if len(records) != 0 {
t.Fatalf("expected 0 records, got %d", len(records))
}
}
func TestStoreAppendValidatesRecord(t *testing.T) {
store := newTestStore(t)
rec := validRecord()
rec.Sequence = 0
rec.RecordID = ""
_, err := store.AppendRecord(context.Background(), rec)
if !errors.Is(err, ErrInvalidIdentity) {
t.Fatalf("expected ErrInvalidIdentity, got: %v", err)
}
}
func TestStoreCASConflictRetry(t *testing.T) {
store := newTestStore(t)
const writers = 8
var wg sync.WaitGroup
errs := make(chan error, writers)
for i := 0; i < writers; i++ {
wg.Add(1)
go func(n int) {
defer wg.Done()
rec := validRecord()
rec.Sequence = 0
rec.RecordID = fmt.Sprintf("rec-%03d", n+1)
if _, err := store.AppendRecord(context.Background(), rec); err != nil {
errs <- err
}
}(i)
}
wg.Wait()
close(errs)
for err := range errs {
t.Errorf("concurrent append: %v", err)
}
records, err := store.ReplayRecords(context.Background())
if err != nil {
t.Fatalf("ReplayRecords: %v", err)
}
if len(records) != writers {
t.Fatalf("expected %d records, got %d", writers, len(records))
}
for i, rec := range records {
if rec.Sequence != uint64(i+1) {
t.Fatalf("record %d sequence = %d, want %d", i, rec.Sequence, i+1)
}
}
}
func TestStoreAppendAssignsSequence(t *testing.T) {
store := newTestStore(t)
rec1 := validRecord()
rec1.Sequence = 0
seq1, err := store.AppendRecord(context.Background(), rec1)
if err != nil {
t.Fatalf("AppendRecord 1: %v", err)
}
if seq1 != 1 {
t.Fatalf("expected sequence 1, got %d", seq1)
}
rec2 := validRecord()
rec2.Sequence = 0
rec2.RecordID = "rec-002"
seq2, err := store.AppendRecord(context.Background(), rec2)
if err != nil {
t.Fatalf("AppendRecord 2: %v", err)
}
if seq2 != 2 {
t.Fatalf("expected sequence 2, got %d", seq2)
}
records, err := store.ReplayRecords(context.Background())
if err != nil {
t.Fatalf("ReplayRecords: %v", err)
}
if len(records) != 2 || records[0].Sequence != 1 || records[1].Sequence != 2 {
t.Fatalf("unexpected replayed records: %+v", records)
}
}
func TestStoreAppendRejectsAssignedSequence(t *testing.T) {
store := newTestStore(t)
rec := validRecord()
rec.Sequence = 5
_, err := store.AppendRecord(context.Background(), rec)
if !errors.Is(err, ErrInvalidIdentity) {
t.Fatalf("expected ErrInvalidIdentity for preassigned sequence, got: %v", err)
}
}
func TestStoreRejectsJournalIdentityDrift(t *testing.T) {
t.Run("mismatched record project identity", func(t *testing.T) {
store := newTestStore(t)
rec := validRecord()
rec.Sequence = 0
rec.ProjectID = "other-proj"
rec.Locators[0].ProjectID = "other-proj"
_, err := store.AppendRecord(context.Background(), rec)
if !errors.Is(err, ErrInvalidIdentity) {
t.Fatalf("expected ErrInvalidIdentity, got: %v", err)
}
})
t.Run("mismatched record workspace identity", func(t *testing.T) {
store := newTestStore(t)
rec := validRecord()
rec.Sequence = 0
rec.WorkspaceID = "other-ws"
rec.Locators[0].WorkspaceID = "other-ws"
_, err := store.AppendRecord(context.Background(), rec)
if !errors.Is(err, ErrInvalidIdentity) {
t.Fatalf("expected ErrInvalidIdentity, got: %v", err)
}
})
t.Run("tuple key isolation", func(t *testing.T) {
sharedStatePath := filepath.Join(t.TempDir(), "shared.json")
sharedState, err := agentstate.NewStore(sharedStatePath)
if err != nil {
t.Fatalf("agentstate.NewStore: %v", err)
}
root := t.TempDir()
storeA, err := NewStore(sharedState, root, "proj:001", "ws")
if err != nil {
t.Fatalf("NewStore A: %v", err)
}
storeB, err := NewStore(sharedState, root, "proj", "001:ws")
if err != nil {
t.Fatalf("NewStore B: %v", err)
}
recA := validRecord()
recA.Sequence = 0
recA.ProjectID = "proj:001"
recA.WorkspaceID = "ws"
recA.Locators[0].ProjectID = "proj:001"
recA.Locators[0].WorkspaceID = "ws"
if _, err := storeA.AppendRecord(context.Background(), recA); err != nil {
t.Fatalf("storeA.AppendRecord: %v", err)
}
replayedB, err := storeB.ReplayRecords(context.Background())
if err != nil {
t.Fatalf("storeB.ReplayRecords: %v", err)
}
if len(replayedB) != 0 {
t.Fatalf("expected storeB to find 0 records, got %d", len(replayedB))
}
})
t.Run("persisted journal identity drift rejected", func(t *testing.T) {
sharedStatePath := filepath.Join(t.TempDir(), "shared.json")
sharedState, err := agentstate.NewStore(sharedStatePath)
if err != nil {
t.Fatalf("agentstate.NewStore: %v", err)
}
store, err := NewStore(sharedState, t.TempDir(), "proj-001", "ws-001")
if err != nil {
t.Fatalf("NewStore: %v", err)
}
corruptJournal := journal{
SchemaVersion: 1,
ProjectID: "wrong-proj",
WorkspaceID: "ws-001",
NextSequence: 1,
ArchiveOrdinal: 0,
}
payload, _ := json.Marshal(corruptJournal)
if _, err := sharedState.CompareAndSwapIntegrationRecord(context.Background(), store.key, "", payload); err != nil {
t.Fatalf("inject corrupt journal: %v", err)
}
rec := validRecord()
rec.Sequence = 0
if _, err := store.AppendRecord(context.Background(), rec); !errors.Is(err, ErrInvalidIdentity) {
t.Fatalf("expected ErrInvalidIdentity on AppendRecord with corrupt journal, got: %v", err)
}
if _, err := store.ReplayRecords(context.Background()); !errors.Is(err, ErrInvalidIdentity) {
t.Fatalf("expected ErrInvalidIdentity on ReplayRecords with corrupt journal, got: %v", err)
}
})
}
func TestStoreRejectsPersistedJournalDrift(t *testing.T) {
tests := []struct {
name string
mutate func(*journal)
wantErr error
}{
{
name: "foreign record project identity",
mutate: func(j *journal) {
j.Records[0].ProjectID = "other-proj"
j.Records[0].Locators[0].ProjectID = "other-proj"
},
wantErr: ErrInvalidIdentity,
},
{
name: "foreign record workspace identity",
mutate: func(j *journal) {
j.Records[0].WorkspaceID = "other-ws"
j.Records[0].Locators[0].WorkspaceID = "other-ws"
},
wantErr: ErrInvalidIdentity,
},
{
name: "invalid record schema",
mutate: func(j *journal) {
j.Records[0].SchemaVersion = RecordSchemaVersion + 1
},
wantErr: ErrInvalidSchemaVersion,
},
{
name: "duplicate sequence",
mutate: func(j *journal) {
second := j.Records[0]
second.RecordID = "rec-002"
j.Records = append(j.Records, second)
j.NextSequence = 2
},
wantErr: ErrInvalidIdentity,
},
{
name: "gapped sequence",
mutate: func(j *journal) {
second := j.Records[0]
second.RecordID = "rec-003"
second.Sequence = 3
j.Records = append(j.Records, second)
j.NextSequence = 4
},
wantErr: ErrInvalidIdentity,
},
{
name: "stale next sequence",
mutate: func(j *journal) {
j.NextSequence = 1
},
wantErr: ErrInvalidIdentity,
},
{
name: "zero next sequence",
mutate: func(j *journal) {
j.Records = nil
j.NextSequence = 0
},
wantErr: ErrInvalidIdentity,
},
{
name: "overflow next sequence",
mutate: func(j *journal) {
j.Records = nil
j.NextSequence = math.MaxUint64
},
wantErr: ErrInvalidIdentity,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
store := newTestStore(t)
record := validRecord()
j := newJournal(store.projectID, store.workspaceID, "")
j.Records = []ProjectLogRecord{record}
j.NextSequence = record.Sequence + 1
fingerprint, err := recordFingerprint(record)
if err != nil {
t.Fatalf("recordFingerprint: %v", err)
}
j.SeenRecords[record.RecordID] = seenRecord{
Sequence: record.Sequence,
Fingerprint: fingerprint,
}
tt.mutate(&j)
payload, err := json.Marshal(j)
if err != nil {
t.Fatalf("marshal drifted journal: %v", err)
}
if _, err := store.state.CompareAndSwapIntegrationRecord(
context.Background(),
store.key,
"",
payload,
); err != nil {
t.Fatalf("inject drifted journal: %v", err)
}
beforePayload, beforeRevision := journalState(t, store)
if _, err := store.ReplayRecords(context.Background()); !errors.Is(err, tt.wantErr) {
t.Fatalf("ReplayRecords error = %v, want %v", err, tt.wantErr)
}
requireJournalState(t, store, beforePayload, beforeRevision)
appendRecord := validRecord()
appendRecord.Sequence = 0
appendRecord.RecordID = "append-after-drift"
if _, err := store.AppendRecord(context.Background(), appendRecord); !errors.Is(err, tt.wantErr) {
t.Fatalf("AppendRecord error = %v, want %v", err, tt.wantErr)
}
requireJournalState(t, store, beforePayload, beforeRevision)
})
}
}
func TestArchiveTerminal(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("Archive: %v", err)
}
records, err := store.ReplayRecords(context.Background())
if err != nil {
t.Fatalf("ReplayRecords: %v", err)
}
if len(records) != 0 {
t.Fatalf("expected 0 records after archive, got %d", len(records))
}
archiveFilesExist(t, store, 1)
}
func TestArchiveNoTerminal(t *testing.T) {
store := newTestStore(t)
rec := validRecord()
rec.Sequence = 0
rec.RecordID = "rec-001"
if _, err := store.AppendRecord(context.Background(), rec); err != nil {
t.Fatalf("AppendRecord: %v", err)
}
err := store.Archive(context.Background())
if !errors.Is(err, ErrNoTerminalRecord) {
t.Fatalf("expected ErrNoTerminalRecord, got: %v", err)
}
}
func TestArchiveEmpty(t *testing.T) {
store := newTestStore(t)
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("Archive empty: %v", err)
}
}
func TestArchiveCrashBeforeWrite(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.failureHook = func(phase string) error {
if phase == "archive_before_write" {
return fmt.Errorf("injected: %s", phase)
}
return nil
}
if err := store.Archive(context.Background()); err == nil {
t.Fatal("expected archive to fail before write")
}
archiveDir := store.archiveDir()
if _, err := os.Stat(filepath.Join(archiveDir, "000001.intent")); err != nil {
t.Errorf("intent should exist: %v", err)
}
if _, err := os.Stat(filepath.Join(archiveDir, "000001.jsonl")); !os.IsNotExist(err) {
t.Errorf("jsonl should not exist after crash before write")
}
records, err := store.ReplayRecords(context.Background())
if err != nil {
t.Fatalf("ReplayRecords before recovery: %v", err)
}
if len(records) != 2 {
t.Fatalf("expected 2 records, got %d", len(records))
}
store.failureHook = nil
if err := store.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
records, err = store.ReplayRecords(context.Background())
if err != nil {
t.Fatalf("ReplayRecords after recovery: %v", err)
}
if len(records) != 0 {
t.Fatalf("records after recovery = %d, want 0", len(records))
}
}
func TestArchiveCrashAfterRename(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.failureHook = func(phase string) error {
if phase == "archive_after_rename" {
return fmt.Errorf("injected: %s", phase)
}
return nil
}
if err := store.Archive(context.Background()); err == nil {
t.Fatal("expected archive to fail after rename")
}
archiveDir := store.archiveDir()
for _, name := range []string{"000001.intent", "000001.jsonl", "000001.timeline.jsonl"} {
if _, err := os.Stat(filepath.Join(archiveDir, name)); err != nil {
t.Errorf("expected %s to exist: %v", name, err)
}
}
if _, err := os.Stat(filepath.Join(archiveDir, "000001.manifest.json")); !os.IsNotExist(err) {
t.Errorf("manifest should not exist after crash after rename")
}
records, _ := store.ReplayRecords(context.Background())
if len(records) != 2 {
t.Fatalf("expected 2 records, got %d", len(records))
}
store.failureHook = nil
if err := store.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
records, _ = store.ReplayRecords(context.Background())
if len(records) != 0 {
t.Fatalf("expected 0 records after reconcile, got %d", len(records))
}
if _, err := os.Stat(filepath.Join(archiveDir, "000001.manifest.json")); err != nil {
t.Errorf("manifest should exist after reconcile: %v", err)
}
}
func TestArchiveCrashBeforeManifest(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.failureHook = func(phase string) error {
if phase == "archive_before_manifest" {
return fmt.Errorf("injected: %s", phase)
}
return nil
}
if err := store.Archive(context.Background()); err == nil {
t.Fatal("expected archive to fail before manifest")
}
archiveDir := store.archiveDir()
for _, name := range []string{"000001.intent", "000001.jsonl", "000001.timeline.jsonl"} {
if _, err := os.Stat(filepath.Join(archiveDir, name)); err != nil {
t.Errorf("expected %s to exist: %v", name, err)
}
}
if _, err := os.Stat(filepath.Join(archiveDir, "000001.manifest.json")); !os.IsNotExist(err) {
t.Errorf("final manifest should not exist after crash before manifest")
}
records, _ := store.ReplayRecords(context.Background())
if len(records) != 2 {
t.Fatalf("expected 2 records, got %d", len(records))
}
store.failureHook = nil
if err := store.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
records, _ = store.ReplayRecords(context.Background())
if len(records) != 0 {
t.Fatalf("expected 0 records after reconcile, got %d", len(records))
}
}
func TestArchiveCrashBeforeCleanup(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.failureHook = func(phase string) error {
if phase == "archive_before_cleanup" {
return fmt.Errorf("injected: %s", phase)
}
return nil
}
if err := store.Archive(context.Background()); err == nil {
t.Fatal("expected archive to fail before cleanup")
}
archiveFilesExist(t, store, 1)
records, _ := store.ReplayRecords(context.Background())
if len(records) != 2 {
t.Fatalf("expected 2 records before reconcile, got %d", len(records))
}
store.failureHook = nil
if err := store.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
records, _ = store.ReplayRecords(context.Background())
if len(records) != 0 {
t.Fatalf("expected 0 records after reconcile, got %d", len(records))
}
}
func TestArchiveExactlyOnce(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("first Archive: %v", err)
}
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("second Archive: %v", err)
}
archiveDir := store.archiveDir()
entries, err := os.ReadDir(archiveDir)
if err != nil {
t.Fatalf("ReadDir: %v", err)
}
intentCount := 0
for _, e := range entries {
if filepath.Ext(e.Name()) == ".intent" {
intentCount++
}
}
if intentCount != 1 {
t.Fatalf("expected 1 intent file, got %d", intentCount)
}
}
func TestArchiveConcurrentCallsShareImmutableOrdinal(t *testing.T) {
root := t.TempDir()
statePath := filepath.Join(t.TempDir(), "manager.json")
state, err := agentstate.NewStore(statePath)
if err != nil {
t.Fatalf("agentstate.NewStore: %v", err)
}
storeA, err := NewStore(state, root, "proj-001", "ws-001")
if err != nil {
t.Fatalf("NewStore A: %v", err)
}
storeB, err := NewStore(state, root, "proj-001", "ws-001")
if err != nil {
t.Fatalf("NewStore B: %v", err)
}
appendTerminalRecords(t, storeA, 2)
firstAtWrite := make(chan struct{})
releaseFirst := make(chan struct{})
secondBeforeLock := make(chan struct{})
var committedArtifacts map[string][]byte
storeA.WithFailureHook(func(phase string) error {
switch phase {
case "archive_before_write":
close(firstAtWrite)
<-releaseFirst
case "archive_before_cleanup":
var err error
committedArtifacts, err = readArchiveArtifactSet(storeA, 1)
return err
}
return nil
})
storeB.WithFailureHook(func(phase string) error {
if phase == "archive_before_lock" {
close(secondBeforeLock)
}
return nil
})
results := make(chan error, 2)
go func() {
results <- storeA.Archive(context.Background())
}()
<-firstAtWrite
go func() {
results <- storeB.Archive(context.Background())
}()
<-secondBeforeLock
close(releaseFirst)
for range 2 {
if err := <-results; err != nil {
t.Fatalf("concurrent Archive: %v", err)
}
}
if committedArtifacts == nil {
t.Fatal("first Archive did not capture its committed artifact set")
}
requireArchiveArtifactSet(t, storeA, 1, committedArtifacts)
records, err := storeA.ReplayRecords(context.Background())
if err != nil {
t.Fatalf("ReplayRecords: %v", err)
}
if len(records) != 0 {
t.Fatalf("expected journal to be pruned exactly once, got %d records", len(records))
}
if err := storeB.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
requireArchiveArtifactSet(t, storeB, 1, committedArtifacts)
archiveFilesExist(t, storeB, 1)
}
func TestArchiveRejectsPreCleanupArtifactDrift(t *testing.T) {
malformedIntent := []byte("{\"schema_version\":")
tests := []struct {
name string
mutate func(*Store) error
verify func(*testing.T, *Store)
}{
{
name: "intent",
mutate: func(store *Store) error {
payload, err := os.ReadFile(store.intentPath(1))
if err != nil {
return err
}
var intent archiveIntent
if err := json.Unmarshal(payload, &intent); err != nil {
return err
}
intent.SchemaVersion++
payload, err = json.Marshal(intent)
if err != nil {
return err
}
return os.WriteFile(store.intentPath(1), append(payload, '\n'), 0o600)
},
},
{
name: "malformed intent",
mutate: func(store *Store) error {
return os.WriteFile(store.intentPath(1), malformedIntent, 0o600)
},
verify: func(t *testing.T, store *Store) {
t.Helper()
got, err := os.ReadFile(store.intentPath(1))
if err != nil {
t.Fatalf("read malformed intent: %v", err)
}
if !bytes.Equal(got, malformedIntent) {
t.Fatal("malformed intent bytes changed")
}
},
},
{
name: "jsonl",
mutate: func(store *Store) error {
return os.WriteFile(store.jsonlPath(1), []byte("tampered jsonl\n"), 0o600)
},
},
{
name: "manifest",
mutate: func(store *Store) error {
payload, err := os.ReadFile(store.manifestPath(1))
if err != nil {
return err
}
var manifest ArchiveManifest
if err := json.Unmarshal(payload, &manifest); err != nil {
return err
}
manifest.ProjectID = "other-proj"
payload, err = json.Marshal(manifest)
if err != nil {
return err
}
return os.WriteFile(store.manifestPath(1), append(payload, '\n'), 0o600)
},
},
{
name: "timeline",
mutate: func(store *Store) error {
return os.WriteFile(store.timelinePath(1), []byte("tampered timeline\n"), 0o600)
},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
beforePayload, beforeRevision := journalState(t, store)
store.WithFailureHook(func(phase string) error {
if phase == "archive_before_cleanup" {
return tt.mutate(store)
}
return nil
})
err := store.Archive(context.Background())
if !errors.Is(err, ErrArchiveConflict) {
t.Fatalf("Archive error = %v, want ErrArchiveConflict", err)
}
requireJournalState(t, store, beforePayload, beforeRevision)
if tt.verify != nil {
tt.verify(t, store)
}
})
}
}
func TestArchiveRetryReusesCommittedArtifacts(t *testing.T) {
t.Run("valid retry reuses bytes", func(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.WithFailureHook(func(phase string) error {
if phase == "archive_before_cleanup" {
return fmt.Errorf("injected: archive_before_cleanup")
}
return nil
})
if err := store.Archive(context.Background()); err == nil {
t.Fatal("expected initial Archive to stop before cleanup")
}
beforeArtifacts := archiveArtifactSet(t, store, 1)
store.WithFailureHook(nil)
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("retry Archive: %v", err)
}
requireArchiveArtifactSet(t, store, 1, beforeArtifacts)
records, err := store.ReplayRecords(context.Background())
if err != nil {
t.Fatalf("ReplayRecords: %v", err)
}
if len(records) != 0 {
t.Fatalf("expected retry to prune verified records, got %d", len(records))
}
})
t.Run("conflicting retry preserves bytes and journal", func(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.WithFailureHook(func(phase string) error {
if phase == "archive_before_cleanup" {
return fmt.Errorf("injected: archive_before_cleanup")
}
return nil
})
if err := store.Archive(context.Background()); err == nil {
t.Fatal("expected initial Archive to stop before cleanup")
}
store.WithFailureHook(nil)
if err := os.WriteFile(store.timelinePath(1), []byte("retry conflict\n"), 0o600); err != nil {
t.Fatalf("tamper timeline: %v", err)
}
beforeArtifacts := archiveArtifactSet(t, store, 1)
beforePayload, beforeRevision := journalState(t, store)
err := store.Archive(context.Background())
if !errors.Is(err, ErrArchiveConflict) {
t.Fatalf("retry Archive error = %v, want ErrArchiveConflict", err)
}
requireArchiveArtifactSet(t, store, 1, beforeArtifacts)
requireJournalState(t, store, beforePayload, beforeRevision)
})
}
func TestReconcileIdempotent(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("Archive: %v", err)
}
if err := store.Reconcile(context.Background()); err != nil {
t.Fatalf("first Reconcile: %v", err)
}
if err := store.Reconcile(context.Background()); err != nil {
t.Fatalf("second Reconcile: %v", err)
}
records, _ := store.ReplayRecords(context.Background())
if len(records) != 0 {
t.Fatalf("expected 0 records, got %d", len(records))
}
}
func TestReconcileConflictFailsClosed(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("Archive: %v", err)
}
archiveDir := store.archiveDir()
manifestPath := filepath.Join(archiveDir, "000001.manifest.json")
data, err := os.ReadFile(manifestPath)
if err != nil {
t.Fatalf("ReadFile manifest: %v", err)
}
var m map[string]any
if err := json.Unmarshal(data, &m); err != nil {
t.Fatalf("Unmarshal manifest: %v", err)
}
m["checksum"] = "sha256:tampered"
tampered, err := json.Marshal(m)
if err != nil {
t.Fatalf("Marshal tampered: %v", err)
}
if err := os.WriteFile(manifestPath, tampered, 0o600); err != nil {
t.Fatalf("WriteFile tampered: %v", err)
}
err = store.Reconcile(context.Background())
if !errors.Is(err, ErrArchiveConflict) {
t.Fatalf("expected ErrArchiveConflict, got: %v", err)
}
}
func TestReconcileRejectsCommittedArtifactDrift(t *testing.T) {
t.Run("missing jsonl returns incomplete", func(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.WithFailureHook(func(phase string) error {
if phase == "archive_before_write" {
return fmt.Errorf("injected")
}
return nil
})
_ = store.Archive(context.Background())
store.WithFailureHook(nil)
if err := store.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile missing JSONL: %v", err)
}
})
t.Run("tampered jsonl returns conflict", func(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.WithFailureHook(func(phase string) error {
if phase == "archive_after_rename" {
return fmt.Errorf("injected")
}
return nil
})
_ = store.Archive(context.Background())
store.WithFailureHook(nil)
archiveDir := store.archiveDir()
jsonlPath := filepath.Join(archiveDir, "000001.jsonl")
_ = os.WriteFile(jsonlPath, []byte("tampered jsonl content\n"), 0o600)
err := store.Reconcile(context.Background())
if !errors.Is(err, ErrArchiveConflict) {
t.Fatalf("expected ErrArchiveConflict, got: %v", err)
}
})
t.Run("manifest field drift with unchanged checksum returns conflict", func(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
_ = store.Archive(context.Background())
archiveDir := store.archiveDir()
manifestPath := filepath.Join(archiveDir, "000001.manifest.json")
data, _ := os.ReadFile(manifestPath)
var m ArchiveManifest
_ = json.Unmarshal(data, &m)
m.ProjectID = "tampered-proj"
tamperedBytes, _ := json.Marshal(m)
_ = os.WriteFile(manifestPath, tamperedBytes, 0o600)
err := store.Reconcile(context.Background())
if !errors.Is(err, ErrArchiveConflict) {
t.Fatalf("expected ErrArchiveConflict for manifest field drift, got: %v", err)
}
})
t.Run("tampered intent returns conflict", func(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
_ = store.Archive(context.Background())
archiveDir := store.archiveDir()
intentPath := filepath.Join(archiveDir, "000001.intent")
data, _ := os.ReadFile(intentPath)
var intent archiveIntent
_ = json.Unmarshal(data, &intent)
intent.SchemaVersion = 99
tamperedBytes, _ := json.Marshal(intent)
_ = os.WriteFile(intentPath, tamperedBytes, 0o600)
err := store.Reconcile(context.Background())
if !errors.Is(err, ErrArchiveConflict) {
t.Fatalf("expected ErrArchiveConflict for tampered intent, got: %v", err)
}
})
t.Run("malformed intent returns conflict without mutation", func(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("Archive: %v", err)
}
malformedIntent := []byte("{\"schema_version\":")
if err := os.WriteFile(store.intentPath(1), malformedIntent, 0o600); err != nil {
t.Fatalf("write malformed intent: %v", err)
}
beforeArtifacts := archiveArtifactSet(t, store, 1)
beforePayload, beforeRevision := journalState(t, store)
err := store.Reconcile(context.Background())
if !errors.Is(err, ErrArchiveConflict) {
t.Fatalf("expected ErrArchiveConflict for malformed intent, got: %v", err)
}
requireArchiveArtifactSet(t, store, 1, beforeArtifacts)
requireJournalState(t, store, beforePayload, beforeRevision)
})
t.Run("tampered timeline returns conflict", func(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
_ = store.Archive(context.Background())
archiveDir := store.archiveDir()
timelinePath := filepath.Join(archiveDir, "000001.timeline.jsonl")
_ = os.WriteFile(timelinePath, []byte("tampered timeline\n"), 0o600)
err := store.Reconcile(context.Background())
if !errors.Is(err, ErrArchiveConflict) {
t.Fatalf("expected ErrArchiveConflict for tampered timeline, got: %v", err)
}
})
}
func TestReconcileReconstructsMissingArtifact(t *testing.T) {
t.Run("reconstructs missing manifest", func(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.WithFailureHook(func(phase string) error {
if phase == "archive_after_rename" {
return fmt.Errorf("injected")
}
return nil
})
_ = store.Archive(context.Background())
store.WithFailureHook(nil)
archiveDir := store.archiveDir()
manifestPath := filepath.Join(archiveDir, "000001.manifest.json")
if _, err := os.Stat(manifestPath); !os.IsNotExist(err) {
t.Fatal("manifest should not exist before reconcile")
}
if err := store.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
if _, err := os.Stat(manifestPath); err != nil {
t.Fatalf("manifest should exist after reconcile: %v", err)
}
})
t.Run("reconstructs missing timeline", func(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
_ = store.Archive(context.Background())
archiveDir := store.archiveDir()
timelinePath := filepath.Join(archiveDir, "000001.timeline.jsonl")
_ = os.Remove(timelinePath)
if err := store.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
if _, err := os.Stat(timelinePath); err != nil {
t.Fatalf("timeline should be reconstructed after reconcile: %v", err)
}
})
}
func TestReconcilePreservesConcurrentAppend(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.WithFailureHook(func(phase string) error {
if phase == "archive_before_cleanup" {
return fmt.Errorf("injected: archive_before_cleanup")
}
return nil
})
if err := store.Archive(context.Background()); err == nil {
t.Fatal("expected archive to fail before cleanup")
}
store.WithFailureHook(nil)
rec3 := validRecord()
rec3.Sequence = 0
rec3.RecordID = "rec-003"
seq3, err := store.AppendRecord(context.Background(), rec3)
if err != nil {
t.Fatalf("AppendRecord 3: %v", err)
}
if seq3 != 3 {
t.Fatalf("expected sequence 3, got %d", seq3)
}
if err := store.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
records, err := store.ReplayRecords(context.Background())
if err != nil {
t.Fatalf("ReplayRecords: %v", err)
}
if len(records) != 1 {
t.Fatalf("expected 1 record remaining in journal, got %d", len(records))
}
if records[0].Sequence != 3 || records[0].RecordID != "rec-003" {
t.Fatalf("expected remaining record seq 3 rec-003, got %+v", records[0])
}
rec4 := validRecord()
rec4.Sequence = 0
rec4.RecordID = "rec-004"
rec4.EventType = agenttask.EventCompleted
rec4.State = agenttask.WorkStateCompleted
rec4.StateRevision = "rev-terminal-2"
rec4.Terminal = true
if _, err := store.AppendRecord(context.Background(), rec4); err != nil {
t.Fatalf("AppendRecord 4: %v", err)
}
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("second Archive: %v", err)
}
records, err = store.ReplayRecords(context.Background())
if err != nil {
t.Fatalf("ReplayRecords after second archive: %v", err)
}
if len(records) != 0 {
t.Fatalf("expected 0 records after second archive, got %d", len(records))
}
archiveFilesExist(t, store, 1)
archiveFilesExist(t, store, 2)
}
func TestReconcileRejectsDivergentPrefix(t *testing.T) {
sharedStatePath := filepath.Join(t.TempDir(), "shared.json")
sharedState, err := agentstate.NewStore(sharedStatePath)
if err != nil {
t.Fatalf("agentstate.NewStore: %v", err)
}
store, err := NewStore(sharedState, t.TempDir(), "proj-001", "ws-001")
if err != nil {
t.Fatalf("NewStore: %v", err)
}
appendTerminalRecords(t, store, 2)
store.WithFailureHook(func(phase string) error {
if phase == "archive_before_cleanup" {
return fmt.Errorf("injected")
}
return nil
})
_ = store.Archive(context.Background())
store.WithFailureHook(nil)
payload, revision, found, err := sharedState.LoadIntegrationRecord(context.Background(), store.key)
if err != nil || !found {
t.Fatalf("LoadIntegrationRecord: %v, found=%v", err, found)
}
var j journal
_ = json.Unmarshal(payload, &j)
j.Records[0].Message = "divergent message"
mutatedPayload, _ := json.Marshal(j)
if _, err := sharedState.CompareAndSwapIntegrationRecord(context.Background(), store.key, revision, mutatedPayload); err != nil {
t.Fatalf("CompareAndSwapIntegrationRecord: %v", err)
}
err = store.Reconcile(context.Background())
if !errors.Is(err, ErrInvalidIdentity) {
t.Fatalf("expected ErrInvalidIdentity for divergent replay index, got: %v", err)
}
}
func TestReconcileRestartConverges(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.failureHook = func(phase string) error {
if phase == "archive_before_cleanup" {
return fmt.Errorf("injected: %s", phase)
}
return nil
}
if err := store.Archive(context.Background()); err == nil {
t.Fatal("expected archive to fail before cleanup")
}
store.failureHook = nil
if err := store.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
records, _ := store.ReplayRecords(context.Background())
if len(records) != 0 {
t.Fatalf("expected 0 records after reconcile, got %d", len(records))
}
archiveFilesExist(t, store, 1)
}
func TestArchiveManifestChecksumMatchesJSONL(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("Archive: %v", err)
}
archiveDir := store.archiveDir()
jsonlData, err := os.ReadFile(filepath.Join(archiveDir, "000001.jsonl"))
if err != nil {
t.Fatalf("ReadFile jsonl: %v", err)
}
manifestData, err := os.ReadFile(filepath.Join(archiveDir, "000001.manifest.json"))
if err != nil {
t.Fatalf("ReadFile manifest: %v", err)
}
var m ArchiveManifest
if err := json.Unmarshal(manifestData, &m); err != nil {
t.Fatalf("Unmarshal manifest: %v", err)
}
if computeJSONLChecksum(jsonlData) != m.Checksum {
t.Fatal("JSONL checksum does not match manifest checksum")
}
}
func TestArchiveTimelineRedacted(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("Archive: %v", err)
}
archiveDir := store.archiveDir()
timelineData, err := os.ReadFile(filepath.Join(archiveDir, "000001.timeline.jsonl"))
if err != nil {
t.Fatalf("ReadFile timeline: %v", err)
}
for _, line := range bytes.Split(timelineData, []byte("\n")) {
line = bytes.TrimSpace(line)
if len(line) == 0 {
continue
}
var entry WorkLogEntry
if err := json.Unmarshal(line, &entry); err != nil {
t.Fatalf("decode timeline entry: %v", err)
}
if entry.Sequence == 0 {
t.Fatal("timeline entry missing sequence")
}
}
}
func TestArchiveMultipleOrdinals(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("first Archive: %v", err)
}
appendTerminalRecords(t, store, 2)
if err := store.Archive(context.Background()); err != nil {
t.Fatalf("second Archive: %v", err)
}
archiveDir := store.archiveDir()
for _, name := range []string{"000001.intent", "000002.intent"} {
if _, err := os.Stat(filepath.Join(archiveDir, name)); err != nil {
t.Errorf("missing %s: %v", name, err)
}
}
}
func TestArchiveRejectsMalformedJSONLWithMatchingIntentChecksum(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.WithFailureHook(func(phase string) error {
if phase == "archive_before_cleanup" {
return fmt.Errorf("injected: archive_before_cleanup")
}
return nil
})
if err := store.Archive(context.Background()); err == nil {
t.Fatal("expected archive to fail before cleanup")
}
store.WithFailureHook(nil)
jsonlPath := store.jsonlPath(1)
malformedJSONL := []byte("invalid jsonl payload\n")
if err := os.WriteFile(jsonlPath, malformedJSONL, 0644); err != nil {
t.Fatalf("WriteFile jsonl: %v", err)
}
intentPath := store.intentPath(1)
intentData, err := os.ReadFile(intentPath)
if err != nil {
t.Fatalf("ReadFile intent: %v", err)
}
var intent archiveIntent
if err := json.Unmarshal(intentData, &intent); err != nil {
t.Fatalf("Unmarshal intent: %v", err)
}
intent.Checksum = computeJSONLChecksum(malformedJSONL)
updatedIntentData, err := json.Marshal(intent)
if err != nil {
t.Fatalf("Marshal intent: %v", err)
}
if err := os.WriteFile(intentPath, updatedIntentData, 0644); err != nil {
t.Fatalf("WriteFile intent: %v", err)
}
wantArtifacts := archiveArtifactSet(t, store, 1)
wantPayload, wantRevision := journalState(t, store)
err = store.Archive(context.Background())
if !errors.Is(err, ErrArchiveConflict) {
t.Fatalf("expected ErrArchiveConflict, got: %v", err)
}
requireArchiveArtifactSet(t, store, 1, wantArtifacts)
requireJournalState(t, store, wantPayload, wantRevision)
}
func TestReconcileRejectsMalformedJSONLWithMatchingIntentChecksum(t *testing.T) {
store := newTestStore(t)
appendTerminalRecords(t, store, 2)
store.WithFailureHook(func(phase string) error {
if phase == "archive_before_cleanup" {
return fmt.Errorf("injected: archive_before_cleanup")
}
return nil
})
if err := store.Archive(context.Background()); err == nil {
t.Fatal("expected archive to fail before cleanup")
}
store.WithFailureHook(nil)
jsonlPath := store.jsonlPath(1)
malformedJSONL := []byte("invalid jsonl payload\n")
if err := os.WriteFile(jsonlPath, malformedJSONL, 0644); err != nil {
t.Fatalf("WriteFile jsonl: %v", err)
}
intentPath := store.intentPath(1)
intentData, err := os.ReadFile(intentPath)
if err != nil {
t.Fatalf("ReadFile intent: %v", err)
}
var intent archiveIntent
if err := json.Unmarshal(intentData, &intent); err != nil {
t.Fatalf("Unmarshal intent: %v", err)
}
intent.Checksum = computeJSONLChecksum(malformedJSONL)
updatedIntentData, err := json.Marshal(intent)
if err != nil {
t.Fatalf("Marshal intent: %v", err)
}
if err := os.WriteFile(intentPath, updatedIntentData, 0644); err != nil {
t.Fatalf("WriteFile intent: %v", err)
}
wantArtifacts := archiveArtifactSet(t, store, 1)
wantPayload, wantRevision := journalState(t, store)
err = store.Reconcile(context.Background())
if !errors.Is(err, ErrArchiveConflict) {
t.Fatalf("expected ErrArchiveConflict, got: %v", err)
}
requireArchiveArtifactSet(t, store, 1, wantArtifacts)
requireJournalState(t, store, wantPayload, wantRevision)
}
func TestStoreReplayDeduplicatesBeforeAndAfterArchive(t *testing.T) {
ctx := context.Background()
root := t.TempDir()
statePath := filepath.Join(root, "projectlog-state.json")
archiveRoot := filepath.Join(root, "archives")
state, err := agentstate.NewStore(statePath)
if err != nil {
t.Fatalf("agentstate.NewStore: %v", err)
}
store, err := NewStore(state, archiveRoot, "proj-replay", "ws-replay")
if err != nil {
t.Fatalf("NewStore: %v", err)
}
scoped, err := store.ForWorkUnit("work-replay")
if err != nil {
t.Fatalf("ForWorkUnit: %v", err)
}
record := scopedTerminalRecord(
"proj-replay",
"ws-replay",
"work-replay",
"attempt-replay",
"evt-sha256-"+strings.Repeat("a", 64),
)
first, err := scoped.AppendRecord(ctx, record)
if err != nil {
t.Fatalf("first AppendRecord: %v", err)
}
replayed, err := scoped.AppendRecord(ctx, record)
if err != nil {
t.Fatalf("pre-archive replay: %v", err)
}
if replayed != first {
t.Fatalf("pre-archive replay sequence = %d, want %d", replayed, first)
}
records, err := scoped.ReplayRecords(ctx)
if err != nil || len(records) != 1 {
t.Fatalf("ReplayRecords before archive = %d, %v", len(records), err)
}
if err := scoped.Archive(ctx); err != nil {
t.Fatalf("Archive: %v", err)
}
reopenedState, err := agentstate.NewStore(statePath)
if err != nil {
t.Fatalf("reopen state: %v", err)
}
reopened, err := NewStore(reopenedState, archiveRoot, "proj-replay", "ws-replay")
if err != nil {
t.Fatalf("reopen store: %v", err)
}
reopenedScoped, err := reopened.ForWorkUnit("work-replay")
if err != nil {
t.Fatalf("reopen scope: %v", err)
}
replayed, err = reopenedScoped.AppendRecord(ctx, record)
if err != nil {
t.Fatalf("post-archive replay: %v", err)
}
if replayed != first {
t.Fatalf("post-archive replay sequence = %d, want %d", replayed, first)
}
records, err = reopenedScoped.ReplayRecords(ctx)
if err != nil || len(records) != 0 {
t.Fatalf("ReplayRecords after replay = %d, %v", len(records), err)
}
if err := reopenedScoped.Archive(ctx); err != nil {
t.Fatalf("second Archive: %v", err)
}
if _, err := os.Stat(reopenedScoped.manifestPath(2)); !os.IsNotExist(err) {
t.Fatalf("second archive was created: %v", err)
}
}
func TestStoreRejectsConflictingRecordReplay(t *testing.T) {
ctx := context.Background()
store := newTestStore(t)
record := scopedTerminalRecord(
"proj-001",
"ws-001",
"work-001",
"attempt-001",
"evt-sha256-"+strings.Repeat("b", 64),
)
if _, err := store.AppendRecord(ctx, record); err != nil {
t.Fatalf("AppendRecord: %v", err)
}
conflict := record
conflict.Message = "different safe content"
if _, err := store.AppendRecord(ctx, conflict); !errors.Is(err, ErrRecordReplayConflict) {
t.Fatalf("conflicting replay error = %v, want ErrRecordReplayConflict", err)
}
records, err := store.ReplayRecords(ctx)
if err != nil || len(records) != 1 || records[0].Message != record.Message {
t.Fatalf("conflicting replay mutated journal: %+v, %v", records, err)
}
if err := store.Archive(ctx); err != nil {
t.Fatalf("Archive: %v", err)
}
if _, err := store.AppendRecord(ctx, conflict); !errors.Is(err, ErrRecordReplayConflict) {
t.Fatalf("post-archive conflict error = %v, want ErrRecordReplayConflict", err)
}
}
func TestStoreScopesParallelWorkArchives(t *testing.T) {
ctx := context.Background()
root := t.TempDir()
state, err := agentstate.NewStore(filepath.Join(root, "state.json"))
if err != nil {
t.Fatalf("agentstate.NewStore: %v", err)
}
store, err := NewStore(state, filepath.Join(root, "archives"), "proj-scoped", "ws-scoped")
if err != nil {
t.Fatalf("NewStore: %v", err)
}
workA, err := store.ForWorkUnit("work-A")
if err != nil {
t.Fatalf("ForWorkUnit A: %v", err)
}
workB, err := store.ForWorkUnit("work-B")
if err != nil {
t.Fatalf("ForWorkUnit B: %v", err)
}
if workA.key == workB.key || workA.archiveDir() == workB.archiveDir() {
t.Fatal("work-unit scopes share a journal key or archive root")
}
recordA := scopedTerminalRecord(
"proj-scoped",
"ws-scoped",
"work-A",
"attempt-A",
"evt-sha256-"+strings.Repeat("c", 64),
)
recordB := scopedTerminalRecord(
"proj-scoped",
"ws-scoped",
"work-B",
"attempt-B",
"evt-sha256-"+strings.Repeat("d", 64),
)
if _, err := workA.AppendRecord(ctx, recordA); err != nil {
t.Fatalf("AppendRecord A: %v", err)
}
if _, err := workB.AppendRecord(ctx, recordB); err != nil {
t.Fatalf("AppendRecord B: %v", err)
}
if _, err := workA.AppendRecord(ctx, recordB); !errors.Is(err, ErrInvalidIdentity) {
t.Fatalf("cross-scope append error = %v, want ErrInvalidIdentity", err)
}
if err := workA.Archive(ctx); err != nil {
t.Fatalf("Archive A: %v", err)
}
if err := workB.Archive(ctx); err != nil {
t.Fatalf("Archive B: %v", err)
}
for name, scoped := range map[string]*Store{"A": workA, "B": workB} {
manifestData, err := os.ReadFile(scoped.manifestPath(1))
if err != nil {
t.Fatalf("read manifest %s: %v", name, err)
}
var manifest ArchiveManifest
if err := json.Unmarshal(manifestData, &manifest); err != nil {
t.Fatalf("decode manifest %s: %v", name, err)
}
if manifest.ArchiveOrdinal != 1 ||
len(manifest.WorkUnitIDs) != 1 ||
manifest.WorkUnitIDs[0] != scoped.workUnitID {
t.Fatalf("manifest %s scope = %+v", name, manifest)
}
}
}
func TestStoreEventReplayFingerprintSurvivesPruneAndRestart(t *testing.T) {
ctx := context.Background()
root := t.TempDir()
statePath := filepath.Join(root, "journal-state.json")
archiveRoot := filepath.Join(root, "archives")
state, err := agentstate.NewStore(statePath)
if err != nil {
t.Fatalf("agentstate.NewStore: %v", err)
}
store, err := NewStore(state, archiveRoot, "proj-event-replay", "ws-event-replay")
if err != nil {
t.Fatalf("NewStore: %v", err)
}
event := agenttask.Event{
EventID: "event-replay-after-prune",
Type: agenttask.EventCompleted,
ProjectID: "proj-event-replay",
WorkspaceID: "ws-event-replay",
WorkUnitID: "work-event-replay",
AttemptID: "attempt-event-replay",
Ordinal: 4,
State: agenttask.WorkStateCompleted,
Detail: "logical completion",
Timestamp: time.Date(2026, 7, 30, 14, 0, 0, 0, time.UTC),
}
recordID, err := opaqueRecordID(event.EventID)
if err != nil {
t.Fatalf("opaqueRecordID: %v", err)
}
fingerprint, err := eventFingerprint(event)
if err != nil {
t.Fatalf("eventFingerprint: %v", err)
}
record := scopedTerminalRecord(
event.ProjectID,
event.WorkspaceID,
event.WorkUnitID,
event.AttemptID,
recordID,
)
record.DispatchOrdinal = event.Ordinal
record.Timestamp = event.Timestamp
record.Message = event.Detail
sequence, err := store.AppendEventRecord(ctx, record, fingerprint)
if err != nil {
t.Fatalf("AppendEventRecord: %v", err)
}
scoped, _ := store.ForWorkUnit(event.WorkUnitID)
if err := scoped.Archive(ctx); err != nil {
t.Fatalf("Archive: %v", err)
}
state, err = agentstate.NewStore(statePath)
if err != nil {
t.Fatalf("reopen agentstate.NewStore: %v", err)
}
store, err = NewStore(state, archiveRoot, event.ProjectID, event.WorkspaceID)
if err != nil {
t.Fatalf("reopen NewStore: %v", err)
}
replayed, err := store.CheckEventReplay(
ctx,
event.WorkUnitID,
recordID,
fingerprint,
)
if err != nil || !replayed {
t.Fatalf("CheckEventReplay = %t, %v", replayed, err)
}
index, _, found, err := store.loadEventReplayIndex(ctx)
if err != nil || !found {
t.Fatalf("load retained event replay index: found=%t, err=%v", found, err)
}
entry := index.Entries[recordID]
if entry.WorkUnitID != event.WorkUnitID ||
entry.Sequence != sequence ||
entry.EventFingerprint != fingerprint {
t.Fatalf("retained replay entry = %+v", entry)
}
volatileProjection := record
volatileProjection.Timestamp = volatileProjection.Timestamp.Add(12 * time.Hour)
volatileProjection.StateRevision = "manager-state-revision-999"
replayedSequence, err := store.AppendEventRecord(ctx, volatileProjection, fingerprint)
if err != nil {
t.Fatalf("volatile AppendEventRecord replay: %v", err)
}
if replayedSequence != sequence {
t.Fatalf("replay sequence = %d, want %d", replayedSequence, sequence)
}
scoped, _ = store.ForWorkUnit(event.WorkUnitID)
records, err := scoped.ReplayRecords(ctx)
if err != nil || len(records) != 0 {
t.Fatalf("records after pruned replay = %d, %v", len(records), err)
}
if err := scoped.Archive(ctx); err != nil {
t.Fatalf("idempotent Archive: %v", err)
}
if _, err := os.Stat(scoped.manifestPath(2)); !os.IsNotExist(err) {
t.Fatalf("unexpected archive ordinal 2: %v", err)
}
}
func TestStoreRejectsLogicalEventIDReuseAcrossScopes(t *testing.T) {
tests := []struct {
name string
from agenttask.WorkUnitID
to agenttask.WorkUnitID
}{
{name: "work A to work B", from: "work-A", to: "work-B"},
{name: "project only to work", to: "work-B"},
{name: "work to project only", from: "work-A"},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
ctx := context.Background()
store := newTestStore(t)
original := agenttask.Event{
EventID: "event-logical-conflict",
Type: agenttask.EventDispatchStarted,
ProjectID: store.projectID,
WorkspaceID: store.workspaceID,
WorkUnitID: test.from,
AttemptID: attemptForScope(test.from, "original"),
Ordinal: 3,
State: agenttask.WorkStateDispatching,
Detail: "original detail",
Timestamp: time.Date(2026, 7, 30, 15, 0, 0, 0, time.UTC),
}
recordID, err := opaqueRecordID(original.EventID)
if err != nil {
t.Fatalf("opaqueRecordID: %v", err)
}
fingerprint, err := eventFingerprint(original)
if err != nil {
t.Fatalf("eventFingerprint: %v", err)
}
originalRecord := projectLogRecordForEvent(original, recordID)
originalSequence, err := store.AppendEventRecord(ctx, originalRecord, fingerprint)
if err != nil {
t.Fatalf("AppendEventRecord: %v", err)
}
conflicting := original
conflicting.WorkUnitID = test.to
conflicting.AttemptID = attemptForScope(test.to, "changed")
conflicting.Detail = "changed logical detail"
conflictingFingerprint, err := eventFingerprint(conflicting)
if err != nil {
t.Fatalf("conflicting eventFingerprint: %v", err)
}
replayed, err := store.CheckEventReplay(
ctx,
test.to,
recordID,
conflictingFingerprint,
)
if replayed || !errors.Is(err, ErrRecordReplayConflict) {
t.Fatalf("conflicting CheckEventReplay = %t, %v", replayed, err)
}
if _, err := store.AppendEventRecord(
ctx,
projectLogRecordForEvent(conflicting, recordID),
conflictingFingerprint,
); !errors.Is(err, ErrRecordReplayConflict) {
t.Fatalf("conflicting AppendEventRecord error = %v", err)
}
originalScope, err := store.scopeForWorkUnit(test.from)
if err != nil {
t.Fatalf("original scope: %v", err)
}
originalRecords, err := originalScope.ReplayRecords(ctx)
if err != nil || len(originalRecords) != 1 ||
originalRecords[0].Sequence != originalSequence {
t.Fatalf("original records = %+v, %v", originalRecords, err)
}
changedScope, err := store.scopeForWorkUnit(test.to)
if err != nil {
t.Fatalf("changed scope: %v", err)
}
changedRecords, err := changedScope.ReplayRecords(ctx)
if err != nil || len(changedRecords) != 0 {
t.Fatalf("changed scope records = %+v, %v", changedRecords, err)
}
})
}
}
func TestStoreEventReplayIndexSerializesCrossScopeCAS(t *testing.T) {
ctx := context.Background()
root := t.TempDir()
statePath := filepath.Join(root, "projectlog-state.json")
archiveRoot := filepath.Join(root, "archives")
stores := make([]*Store, 2)
for index := range stores {
state, err := agentstate.NewStore(statePath)
if err != nil {
t.Fatalf("agentstate.NewStore %d: %v", index, err)
}
stores[index], err = NewStore(
state,
archiveRoot,
"proj-concurrent-replay",
"ws-concurrent-replay",
)
if err != nil {
t.Fatalf("NewStore %d: %v", index, err)
}
}
events := []agenttask.Event{
{
EventID: "event-concurrent-scope-claim",
Type: agenttask.EventDispatchStarted,
ProjectID: stores[0].projectID,
WorkspaceID: stores[0].workspaceID,
WorkUnitID: "work-A",
AttemptID: "attempt-A",
Ordinal: 1,
State: agenttask.WorkStateDispatching,
Detail: "claim A",
Timestamp: time.Date(2026, 7, 30, 15, 30, 0, 0, time.UTC),
},
{
EventID: "event-concurrent-scope-claim",
Type: agenttask.EventDispatchStarted,
ProjectID: stores[0].projectID,
WorkspaceID: stores[0].workspaceID,
WorkUnitID: "work-B",
AttemptID: "attempt-B",
Ordinal: 2,
State: agenttask.WorkStateDispatching,
Detail: "claim B",
Timestamp: time.Date(2026, 7, 30, 15, 30, 1, 0, time.UTC),
},
}
recordID, err := opaqueRecordID(events[0].EventID)
if err != nil {
t.Fatalf("opaqueRecordID: %v", err)
}
start := make(chan struct{})
results := make(chan error, len(events))
for index := range events {
index := index
go func() {
fingerprint, fingerprintErr := eventFingerprint(events[index])
if fingerprintErr != nil {
results <- fingerprintErr
return
}
<-start
_, appendErr := stores[index].AppendEventRecord(
ctx,
projectLogRecordForEvent(events[index], recordID),
fingerprint,
)
results <- appendErr
}()
}
close(start)
var successes, conflicts int
for range events {
err := <-results
switch {
case err == nil:
successes++
case errors.Is(err, ErrRecordReplayConflict):
conflicts++
default:
t.Fatalf("concurrent append error: %v", err)
}
}
if successes != 1 || conflicts != 1 {
t.Fatalf("successes/conflicts = %d/%d, want 1/1", successes, conflicts)
}
index, _, found, err := stores[0].loadEventReplayIndex(ctx)
if err != nil || !found || len(index.Entries) != 1 {
t.Fatalf("event replay index = %+v, found=%t, err=%v", index, found, err)
}
totalRecords := 0
for _, event := range events {
scope, err := stores[0].scopeForWorkUnit(event.WorkUnitID)
if err != nil {
t.Fatalf("scope %q: %v", event.WorkUnitID, err)
}
records, err := scope.ReplayRecords(ctx)
if err != nil {
t.Fatalf("ReplayRecords %q: %v", event.WorkUnitID, err)
}
totalRecords += len(records)
}
if totalRecords != 1 {
t.Fatalf("journal records = %d, want 1", totalRecords)
}
}
func TestStoreEventReplayIndexRecoversLegacyScopedEntry(t *testing.T) {
ctx := context.Background()
root := t.TempDir()
statePath := filepath.Join(root, "projectlog-state.json")
state, err := agentstate.NewStore(statePath)
if err != nil {
t.Fatalf("agentstate.NewStore: %v", err)
}
store, err := NewStore(
state,
filepath.Join(root, "archives"),
"proj-legacy-replay",
"ws-legacy-replay",
)
if err != nil {
t.Fatalf("NewStore: %v", err)
}
event := agenttask.Event{
EventID: "event-legacy-replay",
Type: agenttask.EventDispatchStarted,
ProjectID: store.projectID,
WorkspaceID: store.workspaceID,
WorkUnitID: "work-legacy-A",
AttemptID: "attempt-legacy-A",
Ordinal: 4,
State: agenttask.WorkStateDispatching,
Detail: "legacy retained event",
Timestamp: time.Date(2026, 7, 30, 16, 0, 0, 0, time.UTC),
}
recordID, err := opaqueRecordID(event.EventID)
if err != nil {
t.Fatalf("opaqueRecordID: %v", err)
}
fingerprint, err := eventFingerprint(event)
if err != nil {
t.Fatalf("eventFingerprint: %v", err)
}
legacyScope, err := store.ForWorkUnit(event.WorkUnitID)
if err != nil {
t.Fatalf("ForWorkUnit: %v", err)
}
sequence, err := legacyScope.appendRecord(
ctx,
projectLogRecordForEvent(event, recordID),
fingerprint,
)
if err != nil {
t.Fatalf("legacy appendRecord: %v", err)
}
if _, _, found, err := state.LoadIntegrationRecord(ctx, store.replayKey); err != nil || found {
t.Fatalf("global replay index exists before migration: found=%t, err=%v", found, err)
}
reopenedState, err := agentstate.NewStore(statePath)
if err != nil {
t.Fatalf("reopen agentstate.NewStore: %v", err)
}
reopened, err := NewStore(
reopenedState,
store.root,
store.projectID,
store.workspaceID,
)
if err != nil {
t.Fatalf("reopen NewStore: %v", err)
}
replayed, err := reopened.CheckEventReplay(
ctx,
"work-legacy-B",
recordID,
fingerprint,
)
if replayed || !errors.Is(err, ErrRecordReplayConflict) {
t.Fatalf("cross-scope legacy replay = %t, %v", replayed, err)
}
secondState, err := agentstate.NewStore(statePath)
if err != nil {
t.Fatalf("second reopen agentstate.NewStore: %v", err)
}
second, err := NewStore(secondState, store.root, store.projectID, store.workspaceID)
if err != nil {
t.Fatalf("second reopen NewStore: %v", err)
}
replayed, err = second.CheckEventReplay(
ctx,
event.WorkUnitID,
recordID,
fingerprint,
)
if err != nil || !replayed {
t.Fatalf("exact recovered replay = %t, %v", replayed, err)
}
index, _, found, err := second.loadEventReplayIndex(ctx)
if err != nil || !found {
t.Fatalf("load recovered index: found=%t, err=%v", found, err)
}
entry := index.Entries[recordID]
if entry.WorkUnitID != event.WorkUnitID ||
entry.Sequence != sequence ||
entry.EventFingerprint != fingerprint {
t.Fatalf("recovered entry = %+v", entry)
}
}
func TestStoreEventReplayIndexRejectsMalformedAndConflictingState(t *testing.T) {
ctx := context.Background()
t.Run("malformed retained index", func(t *testing.T) {
store := newTestStore(t)
payload, err := json.Marshal(eventReplayIndex{
SchemaVersion: eventReplayIndexSchemaVersion + 1,
ProjectID: store.projectID,
WorkspaceID: store.workspaceID,
Entries: map[string]eventReplayEntry{},
})
if err != nil {
t.Fatalf("json.Marshal: %v", err)
}
if _, err := store.state.CompareAndSwapIntegrationRecord(
ctx,
store.replayKey,
"",
payload,
); err != nil {
t.Fatalf("seed malformed index: %v", err)
}
_, err = store.CheckEventReplay(
ctx,
"",
"evt-sha256-"+strings.Repeat("e", 64),
"sha256:"+strings.Repeat("f", 64),
)
if !errors.Is(err, ErrInvalidSchemaVersion) {
t.Fatalf("malformed index error = %v", err)
}
})
t.Run("conflicting legacy scopes", func(t *testing.T) {
store := newTestStore(t)
recordID := "evt-sha256-" + strings.Repeat("a", 64)
fingerprints := []string{
"sha256:" + strings.Repeat("b", 64),
"sha256:" + strings.Repeat("c", 64),
}
for index, workUnitID := range []agenttask.WorkUnitID{"work-A", "work-B"} {
scope, err := store.ForWorkUnit(workUnitID)
if err != nil {
t.Fatalf("ForWorkUnit: %v", err)
}
record := scopedTerminalRecord(
store.projectID,
store.workspaceID,
workUnitID,
agenttask.AttemptID(fmt.Sprintf("attempt-%d", index)),
recordID,
)
if _, err := scope.appendRecord(ctx, record, fingerprints[index]); err != nil {
t.Fatalf("legacy append %d: %v", index, err)
}
}
replayed, err := store.CheckEventReplay(ctx, "work-A", recordID, fingerprints[0])
if replayed || !errors.Is(err, ErrRecordReplayConflict) {
t.Fatalf("conflicting legacy replay = %t, %v", replayed, err)
}
if _, _, found, err := store.state.LoadIntegrationRecord(
ctx,
store.replayKey,
); err != nil || found {
t.Fatalf("conflicting recovery persisted an index: found=%t, err=%v", found, err)
}
})
}
func TestS12LoopParallelArchiveMatrix(t *testing.T) {
phases := []string{
"archive_before_lock",
"archive_before_write",
"archive_after_rename",
"archive_before_manifest",
"archive_before_cleanup",
}
for _, phase := range phases {
t.Run(phase, func(t *testing.T) {
runS12ArchivePhase(t, phase)
})
}
}
func runS12ArchivePhase(t *testing.T, failurePhase string) {
t.Helper()
ctx := context.Background()
root := t.TempDir()
managerPath := filepath.Join(root, "manager-state.json")
journalPath := filepath.Join(root, "projectlog-state.json")
archiveRoot := filepath.Join(root, "archives")
managerState, err := agentstate.NewStore(managerPath)
if err != nil {
t.Fatalf("manager state: %v", err)
}
journalState, err := agentstate.NewStore(journalPath)
if err != nil {
t.Fatalf("journal state: %v", err)
}
store, sink := newS12StoreAndSink(
t,
managerState,
journalState,
archiveRoot,
)
workA := newS12Work("work-A", 1, 1, agenttask.WorkStateDispatching)
workB := newS12Work("work-B", 1, 2, agenttask.WorkStateDispatching)
putS12Work(t, managerState, workA)
putS12Work(t, managerState, workB)
start := time.Date(2026, 7, 30, 16, 0, 0, 0, time.UTC)
eventAStart := s12Event(
"event-A-start",
agenttask.EventDispatchStarted,
workA,
"Dispatch token sk-sensitive-value",
start,
)
if err := sink.Emit(ctx, eventAStart); err != nil {
t.Fatalf("Emit A start: %v", err)
}
eventBStart := s12Event(
"event-B-start",
agenttask.EventDispatchStarted,
workB,
"Task B starts on the same project frontier",
start.Add(time.Second),
)
if err := sink.Emit(ctx, eventBStart); err != nil {
t.Fatalf("Emit B start: %v", err)
}
scopeA, _ := store.ForWorkUnit(workA.Unit.ID)
if err := scopeA.Archive(ctx); !errors.Is(err, ErrNoTerminalRecord) {
t.Fatalf("non-terminal Archive error = %v, want ErrNoTerminalRecord", err)
}
var eventATerminal agenttask.Event
var lastFollowupEvent agenttask.Event
for pair := 1; pair <= 11; pair++ {
workA.State = agenttask.WorkStateReviewing
workA.Review = &agenttask.ReviewResult{
ProjectID: "proj-loop",
WorkUnitID: workA.Unit.ID,
AttemptID: workA.AttemptID,
ArtifactID: agenttask.ArtifactID(fmt.Sprintf("artifact-A-%02d", pair)),
Verdict: agenttask.ReviewVerdictFail,
Message: "review requested a follow-up",
Rework: true,
}
putS12Work(t, managerState, workA)
reviewEvent := s12Event(
fmt.Sprintf("event-A-review-%02d", pair),
agenttask.EventReviewResult,
workA,
fmt.Sprintf("review failure %02d", pair),
start.Add(time.Duration(pair*2)*time.Minute),
)
if err := sink.Emit(ctx, reviewEvent); err != nil {
t.Fatalf("Emit A review %d: %v", pair, err)
}
workA = newS12Work(
workA.Unit.ID,
uint32(pair+1),
workA.DispatchOrdinal,
agenttask.WorkStateDispatching,
)
putS12Work(t, managerState, workA)
followupEvent := s12Event(
fmt.Sprintf("event-A-followup-%02d", pair),
agenttask.EventFollowup,
workA,
fmt.Sprintf("follow-up %02d", pair),
start.Add(time.Duration(pair*2+1)*time.Minute),
)
if err := sink.Emit(ctx, followupEvent); err != nil {
t.Fatalf("Emit A follow-up %d: %v", pair, err)
}
lastFollowupEvent = followupEvent
if pair == 3 {
workB.State = agenttask.WorkStateCompleted
workB.ChangeSet = &agenttask.ChangeSetIdentity{
ID: "change-B",
Revision: "change-B-rev",
ArtifactID: "artifact-B",
}
workB.IntegrationAttempt = 1
workB.Integration = &agenttask.IntegrationResult{
ProjectID: "proj-loop",
WorkUnitID: workB.Unit.ID,
ChangeSet: *workB.ChangeSet,
Ordinal: workB.DispatchOrdinal,
Attempt: 1,
Outcome: agenttask.IntegrationOutcomeIntegrated,
BeforeRevision: "base-B-before",
AfterRevision: "base-B-after",
}
putS12Work(t, managerState, workB)
eventBTerminal := s12Event(
"event-B-completed",
agenttask.EventCompleted,
workB,
"Task B completed while Task A remained active",
start.Add(8*time.Minute),
)
if err := sink.Emit(ctx, eventBTerminal); err != nil {
t.Fatalf("Emit B completed: %v", err)
}
scopeB, _ := store.ForWorkUnit(workB.Unit.ID)
if err := scopeB.Archive(ctx); err != nil {
t.Fatalf("Archive B: %v", err)
}
}
if pair == 6 {
managerState, err = agentstate.NewStore(managerPath)
if err != nil {
t.Fatalf("reopen manager state: %v", err)
}
journalState, err = agentstate.NewStore(journalPath)
if err != nil {
t.Fatalf("reopen journal state: %v", err)
}
store, sink = newS12StoreAndSink(
t,
managerState,
journalState,
archiveRoot,
)
advanceS12ManagerRevision(t, managerState)
replayedFollowup := lastFollowupEvent
replayedFollowup.Timestamp = replayedFollowup.Timestamp.Add(4 * time.Hour)
if err := sink.Emit(ctx, replayedFollowup); err != nil {
t.Fatalf("stable follow-up replay after reopen: %v", err)
}
scopeA, _ := store.ForWorkUnit(workA.Unit.ID)
records, err := scopeA.ReplayRecords(ctx)
if err != nil || len(records) != 13 {
t.Fatalf("A records after stable reopen replay = %d, %v", len(records), err)
}
}
}
workA.State = agenttask.WorkStateCompleted
workA.ChangeSet = &agenttask.ChangeSetIdentity{
ID: "change-A",
Revision: "change-A-rev",
ArtifactID: "artifact-A-terminal",
}
workA.IntegrationAttempt = 1
workA.Integration = &agenttask.IntegrationResult{
ProjectID: "proj-loop",
WorkUnitID: workA.Unit.ID,
ChangeSet: *workA.ChangeSet,
Ordinal: workA.DispatchOrdinal,
Attempt: 1,
Outcome: agenttask.IntegrationOutcomeIntegrated,
BeforeRevision: "base-A-before",
AfterRevision: "base-A-after",
}
completion := agenttask.LocatorRecord{
Kind: agenttask.LocatorCompletion,
Opaque: "completion-A-12",
Revision: "completion-A-12-rev",
ProjectID: "proj-loop",
WorkspaceID: "ws-loop",
WorkUnitID: workA.Unit.ID,
AttemptID: workA.AttemptID,
}
workA.Locators[agenttask.LocatorCompletion] = completion
workA.Integration.CompletionLocator = &completion
putS12Work(t, managerState, workA)
eventATerminal = s12Event(
"event-A-completed",
agenttask.EventCompleted,
workA,
"Task A completed after 11 review/follow-up pairs",
start.Add(24*time.Minute),
)
if err := sink.Emit(ctx, eventATerminal); err != nil {
t.Fatalf("Emit A completed: %v", err)
}
scopeA, _ = store.ForWorkUnit(workA.Unit.ID)
scopeA.WithFailureHook(func(phase string) error {
if phase == failurePhase {
return fmt.Errorf("injected failure at %s", phase)
}
return nil
})
if err := scopeA.Archive(ctx); err == nil {
t.Fatalf("Archive succeeded despite injected %s failure", failurePhase)
}
managerState, err = agentstate.NewStore(managerPath)
if err != nil {
t.Fatalf("second manager reopen: %v", err)
}
journalState, err = agentstate.NewStore(journalPath)
if err != nil {
t.Fatalf("second journal reopen: %v", err)
}
store, sink = newS12StoreAndSink(t, managerState, journalState, archiveRoot)
scopeA, _ = store.ForWorkUnit(workA.Unit.ID)
scopeB, _ := store.ForWorkUnit(workB.Unit.ID)
if err := scopeA.Reconcile(ctx); err != nil {
t.Fatalf("Reconcile A after %s: %v", failurePhase, err)
}
if err := scopeB.Reconcile(ctx); err != nil {
t.Fatalf("Reconcile B after %s: %v", failurePhase, err)
}
advanceS12ManagerRevision(t, managerState)
replayedTerminal := eventATerminal
replayedTerminal.Timestamp = replayedTerminal.Timestamp.Add(8 * time.Hour)
if err := sink.Emit(ctx, replayedTerminal); err != nil {
t.Fatalf("post-archive terminal replay: %v", err)
}
records, err := scopeA.ReplayRecords(ctx)
if err != nil || len(records) != 0 {
t.Fatalf("A retained records after replay = %d, %v", len(records), err)
}
if err := scopeA.Archive(ctx); err != nil {
t.Fatalf("idempotent Archive A: %v", err)
}
if _, err := os.Stat(scopeA.manifestPath(2)); !os.IsNotExist(err) {
t.Fatalf("unexpected A archive ordinal 2: %v", err)
}
manifestA := readArchiveManifest(t, scopeA, 1)
manifestB := readArchiveManifest(t, scopeB, 1)
if manifestA.ArchiveOrdinal != 1 ||
manifestA.RecordCount != 24 ||
len(manifestA.WorkUnitIDs) != 1 ||
manifestA.WorkUnitIDs[0] != workA.Unit.ID {
t.Fatalf("A manifest = %+v", manifestA)
}
if manifestB.ArchiveOrdinal != 1 ||
manifestB.RecordCount != 2 ||
len(manifestB.WorkUnitIDs) != 1 ||
manifestB.WorkUnitIDs[0] != workB.Unit.ID {
t.Fatalf("B manifest = %+v", manifestB)
}
entriesA := readWorkLogEntries(t, scopeA, 1)
if len(entriesA) != 24 {
t.Fatalf("A WORK_LOG entries = %d, want 24", len(entriesA))
}
reviews := 0
followups := 0
redacted := false
for _, entry := range entriesA {
if entry.WorkUnitID != workA.Unit.ID || entry.DispatchOrdinal != 1 {
t.Fatalf("cross-task or dispatch drift in A WORK_LOG: %+v", entry)
}
for _, locator := range entry.Locators {
if locator.WorkUnitID != entry.WorkUnitID || locator.AttemptID != entry.AttemptID {
t.Fatalf("locator drift in A WORK_LOG: %+v", entry)
}
}
switch entry.EventType {
case agenttask.EventReviewResult:
reviews++
if entry.RoleStage != "review" || entry.Result != string(agenttask.ReviewVerdictFail) {
t.Fatalf("review WORK_LOG evidence = %+v", entry)
}
case agenttask.EventFollowup:
followups++
if entry.RoleStage != "followup" {
t.Fatalf("follow-up WORK_LOG evidence = %+v", entry)
}
}
if entry.Message == "[REDACTED]" {
redacted = true
}
}
if reviews != 11 || followups != 11 || !redacted {
t.Fatalf("A WORK_LOG reviews=%d followups=%d redacted=%t", reviews, followups, redacted)
}
last := entriesA[len(entriesA)-1]
if !last.Terminal ||
last.LoopOrdinal != 12 ||
last.AttemptID != "attempt-A-12" ||
last.Result != string(agenttask.IntegrationOutcomeIntegrated) ||
len(last.Locators) != 3 {
t.Fatalf("terminal A WORK_LOG evidence = %+v", last)
}
entriesB := readWorkLogEntries(t, scopeB, 1)
if len(entriesB) != 2 ||
!entriesB[1].Terminal ||
entriesB[1].DispatchOrdinal != 2 {
t.Fatalf("B WORK_LOG evidence = %+v", entriesB)
}
}
func advanceS12ManagerRevision(
t *testing.T,
stateStore *agentstate.Store,
) {
t.Helper()
ctx := context.Background()
state, revision, err := stateStore.Load(ctx)
if err != nil {
t.Fatalf("advance manager revision Load: %v", err)
}
if state.Commands == nil {
state.Commands = make(map[agenttask.CommandID]agenttask.CommandRecord)
}
state.Commands["unrelated-command"] = agenttask.CommandRecord{
Intent: agenttask.StartIntent{
CommandID: "unrelated-command",
ProjectID: "unrelated-project",
WorkspaceID: "unrelated-workspace",
MilestoneID: "unrelated-milestone",
WorkflowRevision: "unrelated-workflow",
ConfigRevision: "unrelated-config",
GrantRevision: "unrelated-grant",
StartedAt: time.Date(2026, 7, 30, 20, 0, 0, 0, time.UTC),
},
}
if _, err := stateStore.CompareAndSwap(ctx, revision, state); err != nil {
t.Fatalf("advance manager revision CompareAndSwap: %v", err)
}
}
func scopedTerminalRecord(
projectID agenttask.ProjectID,
workspaceID agenttask.WorkspaceID,
workUnitID agenttask.WorkUnitID,
attemptID agenttask.AttemptID,
recordID string,
) ProjectLogRecord {
record := validRecord()
record.Sequence = 0
record.RecordID = recordID
record.ProjectID = projectID
record.WorkspaceID = workspaceID
record.WorkUnitID = workUnitID
record.AttemptID = attemptID
record.EventType = agenttask.EventCompleted
record.State = agenttask.WorkStateCompleted
record.StateRevision = "state-terminal"
record.Terminal = true
record.Locators[0].ProjectID = projectID
record.Locators[0].WorkspaceID = workspaceID
record.Locators[0].WorkUnitID = workUnitID
record.Locators[0].AttemptID = attemptID
return record
}
func attemptForScope(
workUnitID agenttask.WorkUnitID,
suffix string,
) agenttask.AttemptID {
if workUnitID == "" {
return ""
}
return agenttask.AttemptID("attempt-" + string(workUnitID) + "-" + suffix)
}
func projectLogRecordForEvent(
event agenttask.Event,
recordID string,
) ProjectLogRecord {
record := validRecord()
record.Sequence = 0
record.RecordID = recordID
record.ProjectID = event.ProjectID
record.WorkspaceID = event.WorkspaceID
record.WorkUnitID = event.WorkUnitID
record.AttemptID = event.AttemptID
record.CommandID = event.CommandID
record.EventType = event.Type
record.State = event.State
if event.State != "" {
record.StateRevision = "manager-state-revision"
} else {
record.StateRevision = ""
}
record.DispatchOrdinal = event.Ordinal
record.RouteStatus = nil
record.QuotaObservation = nil
record.Locators = nil
record.Message = event.Detail
record.Metadata = nil
record.Timestamp = event.Timestamp
record.Terminal = event.State.Terminal()
return record
}
func newS12StoreAndSink(
t *testing.T,
managerState *agentstate.Store,
journalState *agentstate.Store,
archiveRoot string,
) (*Store, *Sink) {
t.Helper()
store, err := NewStore(journalState, archiveRoot, "proj-loop", "ws-loop")
if err != nil {
t.Fatalf("NewStore: %v", err)
}
resolver, err := NewStateStoreEvidenceResolver(managerState)
if err != nil {
t.Fatalf("NewStateStoreEvidenceResolver: %v", err)
}
sink, err := NewSink(store, resolver)
if err != nil {
t.Fatalf("NewSink: %v", err)
}
return store, sink
}
func newS12Work(
workUnitID agenttask.WorkUnitID,
attempt uint32,
dispatch agenttask.DispatchOrdinal,
state agenttask.WorkState,
) agenttask.WorkRecord {
attemptID := agenttask.AttemptID(fmt.Sprintf("attempt-%s-%02d", workUnitID[len("work-"):], attempt))
locators := make(map[agenttask.LocatorKind]agenttask.LocatorRecord)
for _, kind := range []agenttask.LocatorKind{
agenttask.LocatorProcess,
agenttask.LocatorSession,
} {
locators[kind] = agenttask.LocatorRecord{
Kind: kind,
Opaque: fmt.Sprintf("%s-%s-%02d", kind, workUnitID, attempt),
Revision: fmt.Sprintf("%s-rev-%02d", kind, attempt),
ProjectID: "proj-loop",
WorkspaceID: "ws-loop",
WorkUnitID: workUnitID,
AttemptID: attemptID,
}
}
return agenttask.WorkRecord{
Unit: agenttask.WorkUnit{
ID: workUnitID,
MilestoneID: "iop-agent-cli-runtime",
WriteSetKind: agenttask.WriteSetDisjoint,
IsolationMode: "overlay",
},
State: state,
Attempt: attempt,
AttemptID: attemptID,
DispatchOrdinal: dispatch,
Target: &agenttask.ExecutionTarget{
ProviderID: "provider-s12",
ModelID: "model-s12",
ProfileID: "profile-s12",
ProfileRevision: "profile-s12-rev",
ConfigRevision: "config-s12-rev",
Capacity: 2,
},
Locators: locators,
}
}
func putS12Work(
t *testing.T,
stateStore *agentstate.Store,
work agenttask.WorkRecord,
) agenttask.StateRevision {
t.Helper()
ctx := context.Background()
state, revision, err := stateStore.Load(ctx)
if err != nil {
t.Fatalf("manager Load: %v", err)
}
if state.Projects == nil {
state.Projects = make(map[agenttask.ProjectID]agenttask.ProjectRecord)
}
project := state.Projects["proj-loop"]
if project.ProjectID == "" {
project = agenttask.ProjectRecord{
ProjectID: "proj-loop",
WorkspaceID: "ws-loop",
Status: agenttask.ProjectStatusRunning,
Intent: &agenttask.StartIntent{
CommandID: "cmd-loop",
ProjectID: "proj-loop",
WorkspaceID: "ws-loop",
MilestoneID: "iop-agent-cli-runtime",
WorkflowRevision: "workflow-loop-rev",
ConfigRevision: "config-s12-rev",
GrantRevision: "grant-s12-rev",
StartedAt: time.Date(2026, 7, 30, 15, 0, 0, 0, time.UTC),
},
Works: make(map[agenttask.WorkUnitID]agenttask.WorkRecord),
}
}
if project.Works == nil {
project.Works = make(map[agenttask.WorkUnitID]agenttask.WorkRecord)
}
project.Works[work.Unit.ID] = work
state.Projects[project.ProjectID] = project
next, err := stateStore.CompareAndSwap(ctx, revision, state)
if err != nil {
t.Fatalf("manager CompareAndSwap: %v", err)
}
return next
}
func s12Event(
eventID string,
eventType agenttask.EventType,
work agenttask.WorkRecord,
detail string,
timestamp time.Time,
) agenttask.Event {
event := agenttask.Event{
EventID: eventID,
Type: eventType,
ProjectID: "proj-loop",
WorkspaceID: "ws-loop",
WorkUnitID: work.Unit.ID,
CommandID: "cmd-loop",
AttemptID: work.AttemptID,
Ordinal: work.DispatchOrdinal,
State: work.State,
ProviderID: work.Target.ProviderID,
ProfileID: work.Target.ProfileID,
Detail: detail,
Timestamp: timestamp,
}
if work.ChangeSet != nil {
event.ChangeSetID = work.ChangeSet.ID
event.ChangeSetRevision = work.ChangeSet.Revision
}
event.IntegrationAttempt = work.IntegrationAttempt
return event
}
func readArchiveManifest(t *testing.T, store *Store, ordinal uint64) ArchiveManifest {
t.Helper()
data, err := os.ReadFile(store.manifestPath(ordinal))
if err != nil {
t.Fatalf("read manifest: %v", err)
}
var manifest ArchiveManifest
if err := json.Unmarshal(data, &manifest); err != nil {
t.Fatalf("decode manifest: %v", err)
}
return manifest
}
func readWorkLogEntries(t *testing.T, store *Store, ordinal uint64) []WorkLogEntry {
t.Helper()
data, err := os.ReadFile(store.timelinePath(ordinal))
if err != nil {
t.Fatalf("read WORK_LOG timeline: %v", err)
}
var entries []WorkLogEntry
for _, line := range bytes.Split(data, []byte("\n")) {
line = bytes.TrimSpace(line)
if len(line) == 0 {
continue
}
var entry WorkLogEntry
if err := json.Unmarshal(line, &entry); err != nil {
t.Fatalf("decode WORK_LOG entry: %v", err)
}
entries = append(entries, entry)
}
return entries
}