| 1 | package session |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "encoding/json" |
| 6 | "os" |
| 7 | "path/filepath" |
| 8 | "strconv" |
| 9 | "strings" |
| 10 | "testing" |
| 11 | "time" |
| 12 | |
| 13 | bolt "go.etcd.io/bbolt" |
| 14 | "reasonix/internal/provider" |
| 15 | ) |
| 16 | |
| 17 | func TestRecoveryCheckpointRestoresTailAndIdempotencyWithoutPrefixReplay(t *testing.T) { |
| 18 | dir := filepath.Join(t.TempDir(), "sessions-v4", "fast") |
| 19 | s, err := CreateWithOptions(dir, "fast", OpenOptions{ExternalHistory: true}) |
| 20 | if err != nil { |
| 21 | t.Fatal(err) |
| 22 | } |
| 23 | |
| 24 | for i := range 80 { |
| 25 | message := provider.Message{ID: "message-" + strings.Repeat("x", 8) + strconv.Itoa(i), Role: provider.RoleUser, Content: strings.Repeat("history", 256)} |
| 26 | payload, marshalErr := json.Marshal(map[string]any{"message": message}) |
| 27 | if marshalErr != nil { |
| 28 | t.Fatal(marshalErr) |
| 29 | } |
| 30 | if _, err := s.Append(context.Background(), Batch{OperationID: "operation-" + string(rune(0x100+i)), Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil { |
| 31 | t.Fatal(err) |
| 32 | } |
| 33 | } |
| 34 | if _, err := s.Flush(t.Context()); err != nil { |
| 35 | t.Fatal(err) |
| 36 | } |
| 37 | wantModel, err := json.Marshal(s.DeriveMessages()) |
| 38 | if err != nil { |
| 39 | t.Fatal(err) |
| 40 | } |
| 41 | if err := s.Close(t.Context()); err != nil { |
| 42 | t.Fatal(err) |
| 43 | } |
| 44 | |
| 45 | // Every query projection is disposable. Removing it must not disable the |
| 46 | // recovery fast path or recent-message access. |
| 47 | if err := os.RemoveAll(filepath.Join(filepath.Dir(dir), ".query-cache", "fast")); err != nil { |
| 48 | t.Fatal(err) |
| 49 | } |
| 50 | var stats RecoveryOpenStats |
| 51 | reopened, err := OpenWithOptions(dir, "fast", OpenOptions{ |
| 52 | ExternalHistory: true, |
| 53 | ObserveRecovery: func(got RecoveryOpenStats) { stats = got }, |
| 54 | }) |
| 55 | if err != nil { |
| 56 | t.Fatal(err) |
| 57 | } |
| 58 | t.Cleanup(func() { _ = reopened.Close(context.Background()) }) |
| 59 | if !stats.UsedCheckpoint { |
| 60 | t.Fatalf("recovery stats = %+v, want checkpoint", stats) |
| 61 | } |
| 62 | if stats.LogBytesRead >= stats.LogBytesTotal { |
| 63 | t.Fatalf("recovery read %d of %d bytes, want bounded tail", stats.LogBytesRead, stats.LogBytesTotal) |
| 64 | } |
| 65 | gotModel, err := json.Marshal(reopened.DeriveMessages()) |
| 66 | if err != nil { |
| 67 | t.Fatal(err) |
| 68 | } |
| 69 | if string(gotModel) != string(wantModel) { |
| 70 | t.Fatal("provider-visible model projection changed across checkpoint recovery") |
| 71 | } |
| 72 | recent := reopened.RecentSnapshot() |
| 73 | if recent.StorageGeneration == "" || recent.DurableSequence != reopened.EventSequence() || len(recent.Entries) == 0 || len(recent.Entries) > RecentMessageLimit { |
| 74 | t.Fatalf("recent snapshot = %+v", recent) |
| 75 | } |
| 76 | |
| 77 | // A durable operation remains globally idempotent even though old operation |
| 78 | // records are no longer resident in Session.operations. |
| 79 | firstPayload, _ := json.Marshal(map[string]any{"message": provider.Message{ID: "message-" + strings.Repeat("x", 8) + "0", Role: provider.RoleUser, Content: strings.Repeat("history", 256)}}) |
| 80 | first, err := reopened.Append(t.Context(), Batch{OperationID: "operation-" + string(rune(0x100)), Events: []Event{{Kind: "message/complete", Payload: firstPayload}}}) |
| 81 | if err != nil { |
| 82 | t.Fatal(err) |
| 83 | } |
| 84 | if first.FirstSequence != 1 { |
| 85 | t.Fatalf("idempotent commit starts at %d, want 1", first.FirstSequence) |
| 86 | } |
| 87 | conflictPayload, _ := json.Marshal(map[string]any{"message": provider.Message{ID: "conflict", Role: provider.RoleUser, Content: "different"}}) |
| 88 | if _, err := reopened.Append(t.Context(), Batch{OperationID: "operation-" + string(rune(0x100)), Events: []Event{{Kind: "message/complete", Payload: conflictPayload}}}); err == nil { |
| 89 | t.Fatal("same operation key with different content unexpectedly succeeded") |
| 90 | } |
| 91 | } |
| 92 | |
| 93 | func TestRecentEntriesKeepAbsoluteVisibleTurns(t *testing.T) { |
| 94 | dir := filepath.Join(t.TempDir(), "session") |
| 95 | if err := os.MkdirAll(dir, 0o700); err != nil { |
| 96 | t.Fatal(err) |
| 97 | } |
| 98 | messages := []provider.Message{ |
| 99 | {ID: "eight", Role: provider.RoleUser, Content: "eight"}, |
| 100 | {ID: "assistant", Role: provider.RoleAssistant, Content: "answer"}, |
| 101 | {ID: "nine", Role: provider.RoleUser, Content: "nine"}, |
| 102 | {ID: "ten", Role: provider.RoleUser, Content: "ten"}, |
| 103 | } |
| 104 | entries, err := buildRecentEntries(t.Context(), dir, messages, 10, 10) |
| 105 | if err != nil { |
| 106 | t.Fatal(err) |
| 107 | } |
| 108 | want := []int{8, 8, 9, 10} |
| 109 | for index := range entries { |
| 110 | if entries[index].VisibleTurn != want[index] { |
| 111 | t.Fatalf("entry %d visible turn=%d, want %d", index, entries[index].VisibleTurn, want[index]) |
| 112 | } |
| 113 | } |
| 114 | } |
| 115 | |
| 116 | func TestRecoveryCheckpointReplaysOnlyDurableTail(t *testing.T) { |
| 117 | root := filepath.Join(t.TempDir(), "sessions-v4") |
| 118 | dir := filepath.Join(root, "tail") |
| 119 | first, err := CreateWithOptions(dir, "tail", OpenOptions{ExternalHistory: true}) |
| 120 | if err != nil { |
| 121 | t.Fatal(err) |
| 122 | } |
| 123 | appendRecoveryTestMessage(t, first, "one", "first") |
| 124 | if _, err := first.Flush(t.Context()); err != nil { |
| 125 | t.Fatal(err) |
| 126 | } |
| 127 | if err := first.Close(t.Context()); err != nil { |
| 128 | t.Fatal(err) |
| 129 | } |
| 130 | |
| 131 | // Suppress checkpoint publication for one append to simulate a log that is |
| 132 | // ahead of its last valid recovery point. |
| 133 | second, err := OpenWithOptions(dir, "tail", OpenOptions{ExternalHistory: true, DisableRecoveryPublish: true}) |
| 134 | if err != nil { |
| 135 | t.Fatal(err) |
| 136 | } |
| 137 | appendRecoveryTestMessage(t, second, "two", "second") |
| 138 | if _, err := second.Flush(t.Context()); err != nil { |
| 139 | t.Fatal(err) |
| 140 | } |
| 141 | if err := second.Close(t.Context()); err != nil { |
| 142 | t.Fatal(err) |
| 143 | } |
| 144 | |
| 145 | var stats RecoveryOpenStats |
| 146 | third, err := OpenWithOptions(dir, "tail", OpenOptions{ExternalHistory: true, ObserveRecovery: func(got RecoveryOpenStats) { stats = got }}) |
| 147 | if err != nil { |
| 148 | t.Fatal(err) |
| 149 | } |
| 150 | defer third.Close(context.Background()) |
| 151 | if !stats.UsedCheckpoint || stats.TailCommits != 1 { |
| 152 | t.Fatalf("recovery stats = %+v, want one tail commit", stats) |
| 153 | } |
| 154 | if got := third.DeriveMessages(); len(got) != 2 || got[0].Content != "first" || got[1].Content != "second" { |
| 155 | t.Fatalf("tail projection = %#v", got) |
| 156 | } |
| 157 | } |
| 158 | |
| 159 | func TestRecoveryCheckpointFallsBackToPreviousGeneration(t *testing.T) { |
| 160 | root := filepath.Join(t.TempDir(), "sessions-v4") |
| 161 | dir := filepath.Join(root, "previous") |
| 162 | store, err := CreateWithOptions(dir, "previous", OpenOptions{ExternalHistory: true}) |
| 163 | if err != nil { |
| 164 | t.Fatal(err) |
| 165 | } |
| 166 | appendRecoveryTestMessage(t, store, "one", "first") |
| 167 | if _, err := store.Flush(t.Context()); err != nil { |
| 168 | t.Fatal(err) |
| 169 | } |
| 170 | appendRecoveryTestMessage(t, store, "two", "second") |
| 171 | if _, err := store.Flush(t.Context()); err != nil { |
| 172 | t.Fatal(err) |
| 173 | } |
| 174 | if err := store.Close(t.Context()); err != nil { |
| 175 | t.Fatal(err) |
| 176 | } |
| 177 | |
| 178 | path := filepath.Join(recoveryCacheDir(dir), recoveryDBName) |
| 179 | db, err := bolt.Open(path, 0o600, &bolt.Options{Timeout: time.Second}) |
| 180 | if err != nil { |
| 181 | t.Fatal(err) |
| 182 | } |
| 183 | if err := db.Update(func(tx *bolt.Tx) error { |
| 184 | return tx.Bucket(recoveryCheckpointBucket).Put(recoveryCurrentKey, []byte("damaged")) |
| 185 | }); err != nil { |
| 186 | _ = db.Close() |
| 187 | t.Fatal(err) |
| 188 | } |
| 189 | if err := db.Close(); err != nil { |
| 190 | t.Fatal(err) |
| 191 | } |
| 192 | |
| 193 | var stats RecoveryOpenStats |
| 194 | reopened, err := OpenWithOptions(dir, "previous", OpenOptions{ExternalHistory: true, ObserveRecovery: func(got RecoveryOpenStats) { stats = got }}) |
| 195 | if err != nil { |
| 196 | t.Fatal(err) |
| 197 | } |
| 198 | defer reopened.Close(context.Background()) |
| 199 | if !stats.UsedCheckpoint { |
| 200 | t.Fatalf("fallback recovery stats = %+v", stats) |
| 201 | } |
| 202 | if got := reopened.DeriveMessages(); len(got) != 2 || got[0].ID != "one" || got[1].ID != "two" { |
| 203 | t.Fatalf("fallback messages = %+v", got) |
| 204 | } |
| 205 | } |
| 206 | |
| 207 | func appendRecoveryTestMessage(t *testing.T, s *Session, id, content string) { |
| 208 | t.Helper() |
| 209 | payload, err := json.Marshal(map[string]any{"message": provider.Message{ID: id, Role: provider.RoleUser, Content: content}}) |
| 210 | if err != nil { |
| 211 | t.Fatal(err) |
| 212 | } |
| 213 | if _, err := s.Append(t.Context(), Batch{OperationID: "append-" + id, Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil { |
| 214 | t.Fatal(err) |
| 215 | } |
| 216 | } |
| 217 |