| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "encoding/json" |
| 6 | "errors" |
| 7 | "fmt" |
| 8 | "strings" |
| 9 | "testing" |
| 10 | "time" |
| 11 | ) |
| 12 | |
| 13 | func mustListProjectTree(t *testing.T, a *App) []ProjectNode { |
| 14 | t.Helper() |
| 15 | nodes, err := a.ListProjectTree() |
| 16 | if err != nil { |
| 17 | t.Fatal(err) |
| 18 | } |
| 19 | return nodes |
| 20 | } |
| 21 | |
| 22 | func TestReadSnapshotFrozenWindows(t *testing.T) { |
| 23 | for _, count := range []int{205, 405, 1000} { |
| 24 | t.Run(fmt.Sprint(count), func(t *testing.T) { |
| 25 | var store readSnapshotStore |
| 26 | defer store.close() |
| 27 | snap, err := store.build(t.Context(), "query", func(ctx context.Context, snap *readSnapshot) error { |
| 28 | for n := range count { |
| 29 | if err := store.append(ctx, snap, n); err != nil { |
| 30 | return err |
| 31 | } |
| 32 | } |
| 33 | return nil |
| 34 | }) |
| 35 | if err != nil { |
| 36 | t.Fatal(err) |
| 37 | } |
| 38 | cursor, total := "", 0 |
| 39 | for { |
| 40 | var values []int |
| 41 | next, id, _, _, err := store.page(t.Context(), "query", cursor, snap, 200, func(b []byte) error { |
| 42 | var n int |
| 43 | if err := json.Unmarshal(b, &n); err != nil { |
| 44 | return err |
| 45 | } |
| 46 | values = append(values, n) |
| 47 | return nil |
| 48 | }) |
| 49 | if err != nil || id != snap.id { |
| 50 | t.Fatalf("page: %s %v", id, err) |
| 51 | } |
| 52 | for _, n := range values { |
| 53 | if n != total { |
| 54 | t.Fatalf("row=%d want %d", n, total) |
| 55 | } |
| 56 | total++ |
| 57 | } |
| 58 | if next == "" { |
| 59 | break |
| 60 | } |
| 61 | cursor = next |
| 62 | } |
| 63 | if total != count { |
| 64 | t.Fatalf("rows=%d", total) |
| 65 | } |
| 66 | if _, _, _, _, err := store.page(t.Context(), "other", cursor, snap, 1, func([]byte) error { return nil }); err == nil { |
| 67 | t.Fatal("cross-query cursor accepted") |
| 68 | } |
| 69 | }) |
| 70 | } |
| 71 | } |
| 72 | |
| 73 | func TestReadSnapshotSpillExpiryAndRelease(t *testing.T) { |
| 74 | var store readSnapshotStore |
| 75 | defer store.close() |
| 76 | snap, err := store.build(t.Context(), "spill", func(ctx context.Context, snap *readSnapshot) error { |
| 77 | for range 600 { |
| 78 | if err := store.append(ctx, snap, strings.Repeat("x", 1024)); err != nil { |
| 79 | return err |
| 80 | } |
| 81 | } |
| 82 | return nil |
| 83 | }) |
| 84 | if err != nil { |
| 85 | t.Fatal(err) |
| 86 | } |
| 87 | if snap.data.db == nil || snap.data.memory != 0 { |
| 88 | t.Fatal("snapshot did not spill") |
| 89 | } |
| 90 | next, _, _, _, err := store.page(t.Context(), "spill", "", snap, 200, func([]byte) error { return nil }) |
| 91 | if err != nil { |
| 92 | t.Fatal(err) |
| 93 | } |
| 94 | store.mu.Lock() |
| 95 | snap.lifetime.used = time.Now().Add(-readSnapshotIdle) |
| 96 | store.mu.Unlock() |
| 97 | if _, _, _, _, err := store.page(t.Context(), "spill", next, nil, 200, func([]byte) error { return nil }); err == nil { |
| 98 | t.Fatal("expired snapshot accepted") |
| 99 | } |
| 100 | store.dispose(snap) |
| 101 | store.close() |
| 102 | if store.disk != 0 || store.memory != 0 { |
| 103 | t.Fatal("release leaked storage") |
| 104 | } |
| 105 | } |
| 106 | |
| 107 | func TestReadSnapshotCancelledWaiterLeavesBuildAlive(t *testing.T) { |
| 108 | var store readSnapshotStore |
| 109 | defer store.close() |
| 110 | entered, release := make(chan struct{}), make(chan struct{}) |
| 111 | ctx, cancel := context.WithCancel(t.Context()) |
| 112 | done := make(chan error, 1) |
| 113 | go func() { |
| 114 | _, err := store.build(ctx, "same", func(ctx context.Context, snap *readSnapshot) error { |
| 115 | close(entered) |
| 116 | <-release |
| 117 | return store.append(ctx, snap, 42) |
| 118 | }) |
| 119 | done <- err |
| 120 | }() |
| 121 | <-entered |
| 122 | cancel() |
| 123 | if err := <-done; !errors.Is(err, context.Canceled) { |
| 124 | t.Fatalf("cancel=%v", err) |
| 125 | } |
| 126 | store.mu.Lock() |
| 127 | job := store.building["same"] |
| 128 | store.mu.Unlock() |
| 129 | close(release) |
| 130 | <-job.done |
| 131 | if job.err != nil || job.snapshot == nil { |
| 132 | t.Fatalf("shared result=%v %v", job.snapshot, job.err) |
| 133 | } |
| 134 | } |
| 135 | |
| 136 | func TestReadSnapshotSharedReleaseAndShutdown(t *testing.T) { |
| 137 | app := NewApp() |
| 138 | store := &app.desktopSessions.readSnapshots |
| 139 | entered, finish := make(chan struct{}), make(chan struct{}) |
| 140 | done := make(chan *readSnapshot, 2) |
| 141 | fill := func(ctx context.Context, snap *readSnapshot) error { |
| 142 | close(entered) |
| 143 | <-finish |
| 144 | return store.append(ctx, snap, 42) |
| 145 | } |
| 146 | go func() { |
| 147 | snap, err := store.build(t.Context(), "shared", fill) |
| 148 | if err != nil { |
| 149 | t.Error(err) |
| 150 | } |
| 151 | done <- snap |
| 152 | }() |
| 153 | <-entered |
| 154 | // Done is consulted after joining the job and releasing the manager lock. |
| 155 | secondStarted := make(chan struct{}) |
| 156 | go func() { |
| 157 | snap, err := store.build(snapshotWaitContext{t.Context(), secondStarted}, "shared", func(ctx context.Context, snap *readSnapshot) error { return store.append(ctx, snap, 42) }) |
| 158 | if err != nil { |
| 159 | t.Error(err) |
| 160 | } |
| 161 | done <- snap |
| 162 | }() |
| 163 | <-secondStarted |
| 164 | close(finish) |
| 165 | a, b := <-done, <-done |
| 166 | if a == nil || b == nil || a.id == b.id || a.data != b.data { |
| 167 | t.Fatal("readers need independent handles") |
| 168 | } |
| 169 | app.ReleaseReadSnapshot(a.id) |
| 170 | app.ReleaseReadSnapshot(a.id) |
| 171 | if _, _, _, _, err := store.page(t.Context(), "shared", "", b, 1, func([]byte) error { return nil }); err != nil { |
| 172 | t.Fatalf("other reader was released: %v", err) |
| 173 | } |
| 174 | app.ReleaseReadSnapshot(b.id) |
| 175 | store.close() |
| 176 | store.close() |
| 177 | if store.memory != 0 || store.disk != 0 { |
| 178 | t.Fatal("shutdown leaked resources") |
| 179 | } |
| 180 | if _, err := store.build(t.Context(), "closed", fill); err == nil { |
| 181 | t.Fatal("closed owner restarted") |
| 182 | } |
| 183 | } |
| 184 | |
| 185 | type snapshotWaitContext struct { |
| 186 | context.Context |
| 187 | joined chan struct{} |
| 188 | } |
| 189 | |
| 190 | func (c snapshotWaitContext) Done() <-chan struct{} { |
| 191 | close(c.joined) |
| 192 | return c.Context.Done() |
| 193 | } |
| 194 | |
| 195 | func TestReadSnapshotExpiredDiskReclaimedBeforeAdmission(t *testing.T) { |
| 196 | var store readSnapshotStore |
| 197 | defer store.close() |
| 198 | fill := func(ctx context.Context, snap *readSnapshot) error { |
| 199 | return store.append(ctx, snap, strings.Repeat("x", 600*1024)) |
| 200 | } |
| 201 | for i := range 4 { |
| 202 | if _, err := store.build(t.Context(), fmt.Sprint(i), fill); err != nil { |
| 203 | t.Fatal(err) |
| 204 | } |
| 205 | } |
| 206 | store.mu.Lock() |
| 207 | reserved := store.disk |
| 208 | for _, snap := range store.entries { |
| 209 | snap.lifetime.used = time.Now().Add(-readSnapshotIdle) |
| 210 | } |
| 211 | store.mu.Unlock() |
| 212 | if reserved != readSnapshotDisk { |
| 213 | t.Fatalf("reservation=%d want %d", reserved, readSnapshotDisk) |
| 214 | } |
| 215 | if _, err := store.build(t.Context(), "after-expiry", fill); err != nil { |
| 216 | t.Fatalf("expired reservations blocked admission: %v", err) |
| 217 | } |
| 218 | store.close() |
| 219 | if store.disk != 0 || store.memory != 0 { |
| 220 | t.Fatal("expired storage leaked") |
| 221 | } |
| 222 | } |
| 223 | |
| 224 | func TestReadSnapshotShutdownWaitsForActiveAndQueuedBuilders(t *testing.T) { |
| 225 | var store readSnapshotStore |
| 226 | defer store.close() |
| 227 | entered := make(chan struct{}, 2) |
| 228 | done := make(chan error, 3) |
| 229 | fill := func(ctx context.Context, snap *readSnapshot) error { |
| 230 | if err := store.append(ctx, snap, strings.Repeat("x", 600*1024)); err != nil { |
| 231 | return err |
| 232 | } |
| 233 | entered <- struct{}{} |
| 234 | <-ctx.Done() |
| 235 | return ctx.Err() |
| 236 | } |
| 237 | for i := range 2 { |
| 238 | go func() { _, err := store.build(t.Context(), fmt.Sprint(i), fill); done <- err }() |
| 239 | } |
| 240 | <-entered |
| 241 | <-entered |
| 242 | joined := make(chan struct{}) |
| 243 | go func() { _, err := store.build(snapshotWaitContext{t.Context(), joined}, "queued", fill); done <- err }() |
| 244 | <-joined |
| 245 | store.close() |
| 246 | for range 3 { |
| 247 | if err := <-done; err == nil { |
| 248 | t.Fatal("shutdown published an unfinished snapshot") |
| 249 | } |
| 250 | } |
| 251 | if store.disk != 0 || store.memory != 0 || len(store.entries) != 0 || len(store.building) != 0 { |
| 252 | t.Fatal("shutdown returned before builder cleanup") |
| 253 | } |
| 254 | } |
| 255 | |
| 256 | func TestReadSnapshotResourceFailureDoesNotPublishPrefix(t *testing.T) { |
| 257 | var store readSnapshotStore |
| 258 | defer store.close() |
| 259 | store.disk = readSnapshotDisk |
| 260 | _, err := store.build(t.Context(), "large", func(ctx context.Context, snap *readSnapshot) error { |
| 261 | for range 600 { |
| 262 | if err := store.append(ctx, snap, strings.Repeat("x", 1024)); err != nil { |
| 263 | return err |
| 264 | } |
| 265 | } |
| 266 | return nil |
| 267 | }) |
| 268 | if err == nil || len(store.entries) != 0 || store.memory != 0 { |
| 269 | t.Fatalf("failed build leaked/published: err=%v memory=%d entries=%d", err, store.memory, len(store.entries)) |
| 270 | } |
| 271 | store.disk = 0 |
| 272 | } |
| 273 | |
| 274 | func TestReadSnapshotEvictionKeepsBoundAndRejectsOldHandle(t *testing.T) { |
| 275 | var store readSnapshotStore |
| 276 | defer store.close() |
| 277 | var first *readSnapshot |
| 278 | for i := range 70 { |
| 279 | snap, err := store.build(t.Context(), fmt.Sprint(i), func(ctx context.Context, snap *readSnapshot) error { return store.append(ctx, snap, 42) }) |
| 280 | if err != nil { |
| 281 | t.Fatal(err) |
| 282 | } |
| 283 | if i == 0 { |
| 284 | first = snap |
| 285 | } |
| 286 | if len(store.entries) > 64 { |
| 287 | t.Fatal("handle bound exceeded") |
| 288 | } |
| 289 | } |
| 290 | if _, _, _, _, err := store.page(t.Context(), "0", "", first, 1, func([]byte) error { return nil }); err == nil { |
| 291 | t.Fatal("evicted handle accepted") |
| 292 | } |
| 293 | store.close() |
| 294 | if store.memory != 0 || store.disk != 0 { |
| 295 | t.Fatal("eviction leaked storage") |
| 296 | } |
| 297 | } |
| 298 | |
| 299 | func TestReadSnapshotEvictionOrdersEqualTimestampsAndPageAccess(t *testing.T) { |
| 300 | var store readSnapshotStore |
| 301 | defer store.close() |
| 302 | var ordered []*readSnapshot |
| 303 | for i := range 32 { |
| 304 | handle, err := store.build(t.Context(), fmt.Sprint(i), func(ctx context.Context, snap *readSnapshot) error { |
| 305 | return store.append(ctx, snap, i) |
| 306 | }) |
| 307 | if err != nil { |
| 308 | t.Fatal(err) |
| 309 | } |
| 310 | ordered = append(ordered, handle.data, handle) |
| 311 | } |
| 312 | // Reading the oldest handle must refresh its position, even when every |
| 313 | // clock sample has the same value (as on a coarse platform clock). |
| 314 | first := ordered[1] |
| 315 | if _, _, _, _, err := store.page(t.Context(), "0", "", first, 1, func([]byte) error { return nil }); err != nil { |
| 316 | t.Fatal(err) |
| 317 | } |
| 318 | ordered = append(append(ordered[:1:1], ordered[2:]...), first) |
| 319 | sharedTime := time.Now() |
| 320 | store.mu.Lock() |
| 321 | for _, snap := range store.entries { |
| 322 | snap.lifetime.used = sharedTime |
| 323 | } |
| 324 | var victims []*readSnapshot |
| 325 | for range ordered { |
| 326 | victims = append(victims, store.evictOldestLocked()) |
| 327 | } |
| 328 | store.mu.Unlock() |
| 329 | for i, victim := range victims { |
| 330 | store.dispose(victim) |
| 331 | if victim != ordered[i] { |
| 332 | t.Errorf("eviction %d did not follow insertion and page-access order", i) |
| 333 | } |
| 334 | } |
| 335 | } |
| 336 |