| 1 | package session |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "encoding/json" |
| 6 | "fmt" |
| 7 | "os" |
| 8 | "path/filepath" |
| 9 | "testing" |
| 10 | ) |
| 11 | |
| 12 | func TestColdCatalogDrainsAfterOneListWhileSlotsWereBusy(t *testing.T) { |
| 13 | root := t.TempDir() |
| 14 | persistence := NewFilesystemPersistence(root) |
| 15 | service, err := NewService("local", persistence) |
| 16 | if err != nil { |
| 17 | t.Fatal(err) |
| 18 | } |
| 19 | for i := range 7 { |
| 20 | id := fmt.Sprintf("cold-%d", i) |
| 21 | runtime, err := service.Create(t.Context(), CreateOptions{SessionID: id}) |
| 22 | if err != nil { |
| 23 | t.Fatal(err) |
| 24 | } |
| 25 | appendRecoveryTestMessage(t, runtime.Session(), "first", "authored request") |
| 26 | if err := service.Close(t.Context(), runtime.Ref()); err != nil { |
| 27 | t.Fatal(err) |
| 28 | } |
| 29 | path := catalogMetadataPath(filepath.Join(root, ".query-cache", id)) |
| 30 | body, err := os.ReadFile(path) |
| 31 | if err != nil { |
| 32 | t.Fatal(err) |
| 33 | } |
| 34 | var old catalogMetadata |
| 35 | if err := json.Unmarshal(body, &old); err != nil { |
| 36 | t.Fatal(err) |
| 37 | } |
| 38 | old.Version = catalogMetadataVersion - 1 |
| 39 | if err := writeCatalogMetadata(filepath.Dir(path), old); err != nil { |
| 40 | t.Fatal(err) |
| 41 | } |
| 42 | } |
| 43 | if err := service.CloseAll(context.Background()); err != nil { |
| 44 | t.Fatal(err) |
| 45 | } |
| 46 | query := newQuery("local", persistence, nil) |
| 47 | t.Cleanup(query.Close) |
| 48 | for range 2 { |
| 49 | if !query.slots.tryAcquire() { |
| 50 | t.Fatal("cannot hold build slots") |
| 51 | } |
| 52 | } |
| 53 | page, err := query.List(t.Context(), "", 100) |
| 54 | if err != nil { |
| 55 | t.Fatal(err) |
| 56 | } |
| 57 | if len(page.Sessions) != 7 { |
| 58 | t.Fatalf("sessions=%d", len(page.Sessions)) |
| 59 | } |
| 60 | query.rebuildMu.Lock() |
| 61 | queued, workers := len(query.rebuilding), query.metadataWorkers |
| 62 | query.rebuildMu.Unlock() |
| 63 | query.slots.release() |
| 64 | query.slots.release() |
| 65 | if queued != 7 || workers > 2 { |
| 66 | t.Fatalf("queued=%d workers=%d", queued, workers) |
| 67 | } |
| 68 | // No second Query.List/Stat and no OpenSession: queued work must finish alone. |
| 69 | query.rebuildWG.Wait() |
| 70 | for _, row := range page.Sessions { |
| 71 | info, err := persistence.Stat(t.Context(), row.SessionID) |
| 72 | if err != nil || info.MetadataStatus != MetadataReady || info.Preview != "authored request" { |
| 73 | t.Fatalf("row=%+v err=%v", info, err) |
| 74 | } |
| 75 | } |
| 76 | } |
| 77 | |
| 78 | func TestCatalogCloseCancelsQueuedWorkers(t *testing.T) { |
| 79 | query := newQuery("local", NewFilesystemPersistence(t.TempDir()), nil) |
| 80 | query.slots.tryAcquire() |
| 81 | query.slots.tryAcquire() |
| 82 | for i := range 100 { |
| 83 | query.scheduleMetadataRebuild(fmt.Sprintf("queued-%d", i)) |
| 84 | } |
| 85 | query.Close() |
| 86 | query.slots.release() |
| 87 | query.slots.release() |
| 88 | query.rebuildMu.Lock() |
| 89 | defer query.rebuildMu.Unlock() |
| 90 | if query.metadataWorkers != 0 || len(query.metadataQueue) != 0 || len(query.rebuilding) != 0 { |
| 91 | t.Fatal("queued state survived close") |
| 92 | } |
| 93 | } |
| 94 | |
| 95 | func TestCatalogRebuildFailureIsNotAnEmptySessionOrEndlessPending(t *testing.T) { |
| 96 | s, _, dir, _ := historyBoundaryFixture(t) |
| 97 | if err := s.CloseAll(context.Background()); err != nil { |
| 98 | t.Fatal(err) |
| 99 | } |
| 100 | log, err := os.OpenFile(filepath.Join(dir, "events.frames"), os.O_WRONLY|os.O_APPEND, 0o600) |
| 101 | if err != nil { |
| 102 | t.Fatal(err) |
| 103 | } |
| 104 | _, err = log.Write([]byte("invalid frame header")) |
| 105 | _ = log.Close() |
| 106 | if err != nil { |
| 107 | t.Fatal(err) |
| 108 | } |
| 109 | q := newQuery("local", s.persistence, nil) |
| 110 | t.Cleanup(q.Close) |
| 111 | if _, err := q.List(t.Context(), "", 100); err != nil { |
| 112 | t.Fatal(err) |
| 113 | } |
| 114 | q.rebuildWG.Wait() |
| 115 | for range 2 { |
| 116 | page, err := q.List(t.Context(), "", 100) |
| 117 | if err != nil || len(page.Sessions) != 1 { |
| 118 | t.Fatalf("page=%+v err=%v", page, err) |
| 119 | } |
| 120 | info := page.Sessions[0] |
| 121 | if info.MetadataStatus != MetadataFailed || info.Error == "" { |
| 122 | t.Fatalf("failure hidden: %+v", info) |
| 123 | } |
| 124 | } |
| 125 | q.rebuildMu.Lock() |
| 126 | defer q.rebuildMu.Unlock() |
| 127 | if len(q.rebuilding) != 0 { |
| 128 | t.Fatal("failed history repeatedly requeued by list reads") |
| 129 | } |
| 130 | } |
| 131 |