2693 lines
77 KiB
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
|
|
}
|