返回 DeepSeek-Reasonix
metadata_restart_test.go
根目录 / internal / sessioncatalog / metadata_restart_test.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "fmt"
6 "os"
7 "path/filepath"
8 "sync"
9 "testing"
10 "time"
11
12 "reasonix/internal/historywork"
13 )
14
15 func TestMetadataQueueSupersedesInvalidScanBeforeReadingItsRemainder(t *testing.T) {
16 dir := t.TempDir()
17 for i := range 3 * historywork.BatchEntries {
18 if err := os.WriteFile(filepath.Join(dir, fmt.Sprintf("%04d.jsonl", i)), []byte("unread body\n"), 0600); err != nil {
19 t.Fatal(err)
20 }
21 }
22 c, err := Open(t.Context(), Options{InMemory: true, MetadataOnly: true, StartPaused: true})
23 if err != nil {
24 t.Fatal(err)
25 }
26 entered, release := make(chan struct{}), make(chan struct{})
27 var releaseOnce, firstSlice sync.Once
28 t.Cleanup(func() { releaseOnce.Do(func() { close(release) }); _ = c.Close(context.Background()) })
29 starts, firstRows := 0, 0
30 c.testReconcileStartHook = func(DirectoryTarget) { starts++ }
31 c.testReconcileBatchHook = func(count int) {
32 if starts == 1 {
33 firstRows += count
34 firstSlice.Do(func() {
35 close(entered)
36 <-release
37 })
38 }
39 }
40 target := DirectoryTarget{Path: dir, Scope: "global"}
41 done, accepted := c.ScheduleReconcile(target)
42 if !accepted {
43 t.Fatal("initial root rejected")
44 }
45 c.ResumeDiscovery()
46 <-entered
47 // The scan is held before publishing its first slice, so the invalidation
48 // deterministically arrives before it can read a second directory batch.
49 joined, accepted := c.ScheduleReconcile(target)
50 if !accepted || joined != done {
51 t.Fatal("replacement lost shared completion")
52 }
53 releaseOnce.Do(func() { close(release) })
54 select {
55 case <-done:
56 case <-t.Context().Done():
57 t.Fatal("replacement did not settle")
58 }
59 if count, err := c.CountDirectorySessions(t.Context(), dir); err != nil || count != 3*historywork.BatchEntries {
60 t.Fatalf("replacement lost a source: count=%d err=%v", count, err)
61 }
62 if err := c.Close(t.Context()); err != nil {
63 t.Fatal(err)
64 }
65 if starts != 2 || firstRows > historywork.BatchEntries {
66 t.Fatalf("obsolete iterator consumed the remainder: starts=%d firstRows=%d", starts, firstRows)
67 }
68 }
69
70 func TestMetadataRestartDebounceCannotExtendItsAdmissionDeadline(t *testing.T) {
71 now := time.Unix(100, 0)
72 job := &metadataQueueJob{}
73 deferMetadataRestart(job, now)
74 deadline := job.restartDeadline
75 if !job.ready.Equal(now.Add(historywork.PauseDuration)) {
76 t.Fatal("restart did not yield a bounded slice")
77 }
78 for i := 1; i <= 20; i++ {
79 deferMetadataRestart(job, now.Add(time.Duration(i)*50*time.Millisecond))
80 if job.restartDeadline != deadline || job.ready.After(deadline) {
81 t.Fatal("continued notifications extended the admission deadline")
82 }
83 }
84 if job.ready.After(now.Add(time.Second)) {
85 t.Fatal("continuous updates starved discovery admission")
86 }
87 }
88
88 lines GO