返回 DeepSeek-Reasonix
history_replace_duplicate_id_test.go
根目录 / internal / session / history_replace_duplicate_id_test.go
1 package session
2
3 import (
4 "context"
5 "encoding/json"
6 "errors"
7 "os"
8 "path/filepath"
9 "testing"
10
11 "reasonix/internal/provider"
12 )
13
14 func duplicateReplacementPayload(t *testing.T, kind string) json.RawMessage {
15 t.Helper()
16 messages := []provider.Message{
17 {ID: "dup", Role: provider.RoleUser, Content: "first dup"},
18 {ID: "answer", Role: provider.RoleAssistant, Content: "plan"},
19 {ID: "dup", Role: provider.RoleUser, Content: "second dup"},
20 {ID: "tail", Role: provider.RoleAssistant, Content: "done"},
21 }
22 body := map[string]any{"messages": messages, "reason": "fresh-resume"}
23 if kind == "legacy/import" {
24 body = map[string]any{"messages": messages, "source": Source{Path: "legacy.jsonl"}}
25 }
26 payload, err := json.Marshal(body)
27 if err != nil {
28 t.Fatal(err)
29 }
30 return payload
31 }
32
33 func TestReplacementRefusesDuplicateMessageID(t *testing.T) {
34 for _, kind := range []string{"history/replace", "legacy/import"} {
35 t.Run(kind, func(t *testing.T) {
36 service, err := NewService("local", NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions")))
37 if err != nil {
38 t.Fatal(err)
39 }
40 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
41 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "dup-write"})
42 if err != nil {
43 t.Fatal(err)
44 }
45 _, err = runtime.Session().Append(t.Context(), Batch{OperationID: "dup", Events: []Event{{Kind: kind, Payload: duplicateReplacementPayload(t, kind)}}})
46 if !errors.Is(err, ErrDuplicateMessageID) {
47 t.Fatalf("append duplicate %s: %v, want ErrDuplicateMessageID", kind, err)
48 }
49 })
50 }
51 }
52
53 // commitDuplicateReplacement writes a history/replace the writer now refuses,
54 // standing in for a log committed before that check existed.
55 func commitDuplicateReplacement(t *testing.T, s *Session) {
56 t.Helper()
57 valid, err := json.Marshal(map[string]any{"messages": []provider.Message{{ID: "placeholder", Role: provider.RoleUser, Content: "x"}}})
58 if err != nil {
59 t.Fatal(err)
60 }
61 prepared, err := s.PrepareBatchContext(t.Context(), "dup-on-disk", Batch{Events: []Event{{Kind: "history/replace", Payload: valid}}})
62 if err != nil {
63 t.Fatal(err)
64 }
65 payload := duplicateReplacementPayload(t, "history/replace")
66 prepared.events[0].Payload = payload
67 prepared.storedEvents[0].Payload, prepared.storedEvents[0].PayloadRef = payload, nil
68 if _, err := s.CommitPrepared(prepared); err != nil {
69 t.Fatal(err)
70 }
71 if _, err := s.Flush(t.Context()); err != nil {
72 t.Fatal(err)
73 }
74 }
75
76 func TestDuplicateMessageIDOnDiskStillHydrates(t *testing.T) {
77 root := filepath.Join(t.TempDir(), "sessions")
78 service, err := NewService("local", NewFilesystemPersistence(root))
79 if err != nil {
80 t.Fatal(err)
81 }
82 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
83 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "dup-read"})
84 if err != nil {
85 t.Fatal(err)
86 }
87 commitDuplicateReplacement(t, runtime.Session())
88 want := []string{"dup", "answer", "tail"}
89 if got := messageIDs(runtime.Session().Snapshot().Projection.Messages); !idsEqual(got, want) {
90 t.Errorf("projection ids = %v, want %v", got, want)
91 }
92 ref := runtime.Ref()
93 if err := service.CloseAll(t.Context()); err != nil {
94 t.Fatal(err)
95 }
96 for _, removeCache := range []bool{false, true} {
97 if removeCache {
98 if err := os.RemoveAll(filepath.Join(root, ".query-cache")); err != nil {
99 t.Fatal(err)
100 }
101 }
102 reopened, err := NewService("local", NewFilesystemPersistence(root))
103 if err != nil {
104 t.Fatal(err)
105 }
106 query := reopened.Query()
107 if _, _, err := query.prepareHistoryIndex(t.Context(), ref); err != nil {
108 t.Fatalf("prepare history index (cache removed=%v): %v", removeCache, err)
109 }
110 page, err := query.ReadHistoryWindow(t.Context(), ref, HistoryWindowRequest{Anchor: "newest"})
111 if err != nil || page.Status != "ready" {
112 t.Fatalf("read history (cache removed=%v): %+v, %v", removeCache, page, err)
113 }
114 if got := windowIDs(t, page); !idsEqual(got, want) {
115 t.Fatalf("history ids (cache removed=%v) = %v, want %v", removeCache, got, want)
116 }
117 if search := searchHistoryReady(t, query, ref, "dup", "", 10); len(search.Hits) != 1 {
118 t.Fatalf("search (cache removed=%v): %+v", removeCache, search)
119 }
120 opened, err := reopened.Open(t.Context(), ref)
121 if err != nil {
122 t.Fatalf("reopen (cache removed=%v): %v", removeCache, err)
123 }
124 if got := messageIDs(opened.Runtime().Session().Snapshot().Projection.Messages); !idsEqual(got, want) {
125 t.Fatalf("reopened projection ids = %v, want %v", got, want)
126 }
127 if err := opened.Release(t.Context()); err != nil {
128 t.Fatal(err)
129 }
130 if err := reopened.CloseAll(t.Context()); err != nil {
131 t.Fatal(err)
132 }
133 }
134 }
135
136 func messageIDs(messages []provider.Message) []string {
137 ids := make([]string, 0, len(messages))
138 for _, message := range messages {
139 ids = append(ids, message.ID)
140 }
141 return ids
142 }
143
143 lines GO