返回 DeepSeek-Reasonix
catalog_queue_test.go
根目录 / internal / session / catalog_queue_test.go
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
131 lines GO