| 1 | package session |
| 2 | |
| 3 | import ( |
| 4 | "bytes" |
| 5 | "context" |
| 6 | "errors" |
| 7 | "os" |
| 8 | "path/filepath" |
| 9 | "testing" |
| 10 | ) |
| 11 | |
| 12 | // Seed only revision-1 events, then label the fixture with its original |
| 13 | // revision. Do not use this helper to downgrade real session data. |
| 14 | func revisionOneFixture(t *testing.T) (string, string, string, []byte, []byte) { |
| 15 | t.Helper() |
| 16 | root, id := filepath.Join(t.TempDir(), "sessions"), "revision-one" |
| 17 | persistence := NewFilesystemPersistence(root) |
| 18 | writer, err := persistence.Create(CreateOptions{SessionID: id}) |
| 19 | if err != nil { |
| 20 | t.Fatal(err) |
| 21 | } |
| 22 | if _, err := writer.Append(t.Context(), Batch{OperationID: "old-message", Events: []Event{{Kind: "message/complete", Payload: []byte(`{"message":{"id":"old","role":"user","content":"preserve original bytes"}}`)}}}); err != nil { |
| 23 | t.Fatal(err) |
| 24 | } |
| 25 | if _, err := writer.Flush(t.Context()); err != nil { |
| 26 | t.Fatal(err) |
| 27 | } |
| 28 | if err := writer.Close(t.Context()); err != nil { |
| 29 | t.Fatal(err) |
| 30 | } |
| 31 | dir := filepath.Join(root, id) |
| 32 | manifestPath := filepath.Join(dir, "manifest.json") |
| 33 | manifest, err := readStoredManifest(manifestPath) |
| 34 | if err != nil { |
| 35 | t.Fatal(err) |
| 36 | } |
| 37 | manifest.StorageRevision = 1 |
| 38 | if err := writeManifestFile(manifestPath, manifest); err != nil { |
| 39 | t.Fatal(err) |
| 40 | } |
| 41 | manifestBytes, err := os.ReadFile(manifestPath) |
| 42 | if err != nil { |
| 43 | t.Fatal(err) |
| 44 | } |
| 45 | logBytes, err := os.ReadFile(filepath.Join(dir, currentLogName)) |
| 46 | if err != nil { |
| 47 | t.Fatal(err) |
| 48 | } |
| 49 | return root, id, dir, manifestBytes, logBytes |
| 50 | } |
| 51 | |
| 52 | func assertRevisionFixtureBytes(t *testing.T, dir string, manifestBytes, logBytes []byte) { |
| 53 | t.Helper() |
| 54 | for name, want := range map[string][]byte{"manifest.json": manifestBytes, currentLogName: logBytes} { |
| 55 | got, err := os.ReadFile(filepath.Join(dir, name)) |
| 56 | if err != nil { |
| 57 | t.Fatal(err) |
| 58 | } |
| 59 | if !bytes.Equal(got, want) { |
| 60 | t.Fatalf("%s bytes changed", name) |
| 61 | } |
| 62 | } |
| 63 | } |
| 64 | |
| 65 | func TestStorageRevisionReadOnlyDoesNotUpgrade(t *testing.T) { |
| 66 | root, id, dir, manifestBytes, logBytes := revisionOneFixture(t) |
| 67 | persistence := NewFilesystemPersistence(root) |
| 68 | reader, err := persistence.Open(id, ReadOnly) |
| 69 | if err != nil { |
| 70 | t.Fatal(err) |
| 71 | } |
| 72 | if _, err := reader.Read(t.Context(), 0, 10); err != nil { |
| 73 | t.Fatal(err) |
| 74 | } |
| 75 | if err := reader.Close(t.Context()); err != nil { |
| 76 | t.Fatal(err) |
| 77 | } |
| 78 | assertRevisionFixtureBytes(t, dir, manifestBytes, logBytes) |
| 79 | } |
| 80 | |
| 81 | func TestStorageRevisionWriterUpgradesWithoutRewritingLog(t *testing.T) { |
| 82 | root, id, dir, _, logBytes := revisionOneFixture(t) |
| 83 | writer, err := NewFilesystemPersistence(root).Open(id, ReadWrite) |
| 84 | if err != nil { |
| 85 | t.Fatal(err) |
| 86 | } |
| 87 | if err := writer.Close(t.Context()); err != nil { |
| 88 | t.Fatal(err) |
| 89 | } |
| 90 | manifest, err := readStoredManifest(filepath.Join(dir, "manifest.json")) |
| 91 | if err != nil { |
| 92 | t.Fatal(err) |
| 93 | } |
| 94 | if manifest.StorageRevision != StorageRevision { |
| 95 | t.Fatalf("revision=%d, want=%d", manifest.StorageRevision, StorageRevision) |
| 96 | } |
| 97 | got, err := os.ReadFile(filepath.Join(dir, currentLogName)) |
| 98 | if err != nil { |
| 99 | t.Fatal(err) |
| 100 | } |
| 101 | if !bytes.Equal(got, logBytes) { |
| 102 | t.Fatal("upgrade rewrote original event bytes") |
| 103 | } |
| 104 | } |
| 105 | |
| 106 | func TestStorageRevisionDamagedLogDoesNotUpgrade(t *testing.T) { |
| 107 | root, id, dir, manifestBytes, logBytes := revisionOneFixture(t) |
| 108 | logBytes[0] ^= 0xff |
| 109 | if err := os.WriteFile(filepath.Join(dir, currentLogName), logBytes, 0o600); err != nil { |
| 110 | t.Fatal(err) |
| 111 | } |
| 112 | writer, err := NewFilesystemPersistence(root).Open(id, ReadWrite) |
| 113 | if err == nil { |
| 114 | _ = writer.Close(context.Background()) |
| 115 | t.Fatal("damaged log was opened for writing") |
| 116 | } |
| 117 | assertRevisionFixtureBytes(t, dir, manifestBytes, logBytes) |
| 118 | } |
| 119 | |
| 120 | func TestStorageRevisionExternalLeasePreventsUpgrade(t *testing.T) { |
| 121 | root, id, dir, manifestBytes, logBytes := revisionOneFixture(t) |
| 122 | release, err := acquireSessionWriter(dir) |
| 123 | if err != nil { |
| 124 | t.Fatal(err) |
| 125 | } |
| 126 | defer release() |
| 127 | writer, err := NewFilesystemPersistence(root).Open(id, ReadWrite) |
| 128 | if writer != nil { |
| 129 | _ = writer.Close(context.Background()) |
| 130 | } |
| 131 | if !errors.Is(err, ErrWriterOwned) { |
| 132 | t.Fatalf("open while owned: %v", err) |
| 133 | } |
| 134 | assertRevisionFixtureBytes(t, dir, manifestBytes, logBytes) |
| 135 | } |
| 136 | |
| 137 | func TestPreviousReaderRejectsStorageRevisionThree(t *testing.T) { |
| 138 | const previousMaxStorageRevision = 2 |
| 139 | if StorageRevision <= previousMaxStorageRevision { |
| 140 | t.Fatalf("StorageRevision=%d is not newer than previous reader max %d", StorageRevision, previousMaxStorageRevision) |
| 141 | } |
| 142 | root, id, dir, _, _ := revisionOneFixture(t) |
| 143 | writer, err := NewFilesystemPersistence(root).Open(id, ReadWrite) |
| 144 | if err != nil { |
| 145 | t.Fatal(err) |
| 146 | } |
| 147 | if err := writer.Close(t.Context()); err != nil { |
| 148 | t.Fatal(err) |
| 149 | } |
| 150 | manifest, err := readStoredManifest(filepath.Join(dir, "manifest.json")) |
| 151 | if err != nil { |
| 152 | t.Fatal(err) |
| 153 | } |
| 154 | if manifest.StorageRevision != StorageRevision { |
| 155 | t.Fatalf("revision=%d, want=%d", manifest.StorageRevision, StorageRevision) |
| 156 | } |
| 157 | previousAccepts := manifest.SchemaVersion == SchemaVersion && manifest.Codec == Codec && |
| 158 | manifest.StorageRevision >= 1 && manifest.StorageRevision <= previousMaxStorageRevision |
| 159 | if previousAccepts { |
| 160 | t.Fatal("previous reader predicate accepted StorageRevision 3") |
| 161 | } |
| 162 | } |
| 163 | |
| 164 | func TestStorageRevisionUnknownRejectedByReadersAndWriters(t *testing.T) { |
| 165 | for _, mode := range []AccessMode{ReadOnly, ReadWrite} { |
| 166 | root, id, dir, _, logBytes := revisionOneFixture(t) |
| 167 | manifestPath := filepath.Join(dir, "manifest.json") |
| 168 | manifest, err := readStoredManifest(manifestPath) |
| 169 | if err != nil { |
| 170 | t.Fatal(err) |
| 171 | } |
| 172 | manifest.StorageRevision = 999 |
| 173 | if err := writeManifestFile(manifestPath, manifest); err != nil { |
| 174 | t.Fatal(err) |
| 175 | } |
| 176 | manifestBytes, err := os.ReadFile(manifestPath) |
| 177 | if err != nil { |
| 178 | t.Fatal(err) |
| 179 | } |
| 180 | handle, err := NewFilesystemPersistence(root).Open(id, mode) |
| 181 | if handle != nil { |
| 182 | _ = handle.Close(context.Background()) |
| 183 | } |
| 184 | if !errors.Is(err, ErrUnsupportedVersion) { |
| 185 | t.Fatalf("mode=%v unknown revision: %v", mode, err) |
| 186 | } |
| 187 | assertRevisionFixtureBytes(t, dir, manifestBytes, logBytes) |
| 188 | } |
| 189 | } |
| 190 | |
| 191 | func TestStorageRevisionManifestPublishFailureReleasesWriter(t *testing.T) { |
| 192 | _, id, dir, manifestBytes, logBytes := revisionOneFixture(t) |
| 193 | manifestPath := filepath.Join(dir, "manifest.json") |
| 194 | backupPath := filepath.Join(dir, "manifest.before-failure.json") |
| 195 | // ObserveRecovery is called after validation and before manifest publication. |
| 196 | // A directory at the rename destination deterministically rejects the atomic |
| 197 | // publish on every platform, without timing or permission assumptions. |
| 198 | writer, err := OpenWithOptions(dir, id, OpenOptions{ObserveRecovery: func(RecoveryOpenStats) { |
| 199 | if err := os.Rename(manifestPath, backupPath); err != nil { |
| 200 | t.Fatal(err) |
| 201 | } |
| 202 | if err := os.Mkdir(manifestPath, 0o700); err != nil { |
| 203 | t.Fatal(err) |
| 204 | } |
| 205 | }}) |
| 206 | if writer != nil { |
| 207 | _ = writer.Close(context.Background()) |
| 208 | } |
| 209 | if err == nil { |
| 210 | t.Fatal("manifest publish unexpectedly succeeded") |
| 211 | } |
| 212 | got, err := os.ReadFile(backupPath) |
| 213 | if err != nil { |
| 214 | t.Fatal(err) |
| 215 | } |
| 216 | if !bytes.Equal(got, manifestBytes) { |
| 217 | t.Fatal("failed upgrade changed old manifest bytes") |
| 218 | } |
| 219 | got, err = os.ReadFile(filepath.Join(dir, currentLogName)) |
| 220 | if err != nil { |
| 221 | t.Fatal(err) |
| 222 | } |
| 223 | if !bytes.Equal(got, logBytes) { |
| 224 | t.Fatal("failed upgrade changed log bytes") |
| 225 | } |
| 226 | if err := os.Remove(manifestPath); err != nil { |
| 227 | t.Fatal(err) |
| 228 | } |
| 229 | if err := os.Rename(backupPath, manifestPath); err != nil { |
| 230 | t.Fatal(err) |
| 231 | } |
| 232 | // The failed open must release its lease, so a retry can finish the upgrade. |
| 233 | retry, err := Open(dir, id) |
| 234 | if err != nil { |
| 235 | t.Fatal(err) |
| 236 | } |
| 237 | if err := retry.Close(t.Context()); err != nil { |
| 238 | t.Fatal(err) |
| 239 | } |
| 240 | manifest, err := readStoredManifest(manifestPath) |
| 241 | if err != nil { |
| 242 | t.Fatal(err) |
| 243 | } |
| 244 | if manifest.StorageRevision != StorageRevision { |
| 245 | t.Fatalf("retry revision=%d, want=%d", manifest.StorageRevision, StorageRevision) |
| 246 | } |
| 247 | } |
| 248 |