返回 DeepSeek-Reasonix
topic_removal_concurrency_test.go
根目录 / desktop / topic_removal_concurrency_test.go
1 package main
2
3 import (
4 "fmt"
5 "os"
6 "path/filepath"
7 "sync"
8 "testing"
9 "time"
10
11 "reasonix/desktop/internal/workspacestate"
12 "reasonix/internal/agent"
13 "reasonix/internal/control"
14 "reasonix/internal/provider"
15 )
16
17 func awaitRemovalTest[T any](t *testing.T, ch <-chan T) T {
18 t.Helper()
19 select {
20 case value := <-ch:
21 return value
22 case <-time.After(10 * time.Second):
23 t.Fatal("concurrent removal operation did not complete")
24 var zero T
25 return zero
26 }
27 }
28
29 // The model finishes before removal, but its commit is delayed until removal
30 // holds the title lock. Removal must already own the removal lock too, and the
31 // delayed title must not commit after archive (including legacy import).
32 func TestTopicRemovalConcurrentAIRename(t *testing.T) {
33 for _, legacy := range []bool{false, true} {
34 t.Run(fmt.Sprintf("legacy=%t", legacy), func(t *testing.T) {
35 a, _, rt, _, path := newCanonicalTitleFixture(t)
36 installNoopRuntimeEvents(a)
37 topicID := "topic-canonical"
38 appendSessionTestMessage(t, rt, "user", provider.Message{ID: "user", Role: provider.RoleUser, Origin: provider.MessageOriginUser, Content: "keep this message"})
39 var target SessionTarget
40 var err error
41 if legacy {
42 topic, createErr := a.CreateTopic("global", "", "Legacy")
43 if createErr != nil {
44 t.Fatal(createErr)
45 }
46 topicID = topic.ID
47 dir := desktopSessionDir(globalWorkspaceRoot())
48 if err := os.MkdirAll(dir, 0700); err != nil {
49 t.Fatal(err)
50 }
51 path = filepath.Join(dir, "legacy-title.jsonl")
52 writeTopicSessionWithPrompt(t, dir, filepath.Base(path), topicID, "Legacy", globalWorkspaceRoot(), "keep this message", time.Now())
53 target = SessionTarget{SessionPath: path, TopicID: topicID}
54 } else {
55 target, err = a.resolveSessionMutationTarget(SessionSelector{TopicID: topicID})
56 if err != nil {
57 t.Fatal(err)
58 }
59 }
60 req := inspectedTopicRemoval(t, a, topicID)
61 generated, commit := make(chan struct{}), make(chan struct{})
62 var release sync.Once
63 t.Cleanup(func() { release.Do(func() { close(commit) }) })
64 a.lifecycleCheckpointHook = func(phase string) {
65 switch phase {
66 case "session-title-before-commit":
67 close(generated)
68 <-commit
69 case "topic-sessions-before-archive":
70 if a.sessionRemovalMu.TryLock() {
71 a.sessionRemovalMu.Unlock()
72 t.Error("removal acquired title before removal ownership")
73 }
74 release.Do(func() { close(commit) })
75 }
76 }
77 renamed := make(chan error, 1)
78 go func() {
79 var err error
80 if legacy {
81 _, err = a.aiRenameLegacySession(t.Context(), target)
82 } else {
83 _, err = a.aiRenameCanonicalSession(t.Context(), target)
84 }
85 renamed <- err
86 }()
87 awaitRemovalTest(t, generated)
88 removed := make(chan error, 1)
89 go func() {
90 out, err := a.RemoveTopic(req)
91 if err == nil && !out.Committed {
92 err = fmt.Errorf("removal: %+v", out)
93 }
94 removed <- err
95 }()
96 if err := awaitRemovalTest(t, removed); err != nil {
97 t.Fatal(err)
98 }
99 if err := awaitRemovalTest(t, renamed); err == nil {
100 t.Fatal("stale AI title committed after archive")
101 }
102 if legacy {
103 assertLegacyLifecycle(t, a, path, "archived")
104 meta, _, err := agent.LoadBranchMeta(path)
105 if err != nil || meta.CustomTitle != "" {
106 t.Fatalf("legacy title changed: %+v %v", meta, err)
107 }
108 } else {
109 state, err := a.workspaceRegistry().Load(t.Context())
110 if err != nil || state.SessionStates[rt.Ref().SessionID].Lifecycle != workspacestate.Archived {
111 t.Fatalf("archive state: %+v %v", state, err)
112 }
113 }
114 if unlock, ok := a.tryLockRuntimeMutation("verify removal completed"); !ok {
115 t.Fatal("runtime admission leaked")
116 } else {
117 unlock()
118 }
119 })
120 }
121 }
122
123 func TestTopicRemovalJoinsAutosaveOutsideTitleLock(t *testing.T) {
124 for _, priorSave := range []bool{false, true} {
125 t.Run(fmt.Sprintf("prior-save=%t", priorSave), func(t *testing.T) {
126 a, _, rt, _, _ := newCanonicalTitleFixture(t)
127 installNoopRuntimeEvents(a)
128 tab := a.tabs["test"]
129 if priorSave {
130 a.scheduleTabSnapshot(tab.ID)
131 waitForAutosaveIdle(t, tab)
132 }
133 req := inspectedTopicRemoval(t, a, tab.TopicID)
134 titleReached := make(chan struct{})
135 a.lifecycleCheckpointHook = func(phase string) {
136 switch phase {
137 case "before-archive-commit":
138 a.scheduleTabSnapshot(tab.ID)
139 awaitRemovalTest(t, titleReached)
140 case "before-autosave-topic-title":
141 close(titleReached)
142 }
143 }
144 done := make(chan error, 1)
145 go func() {
146 out, err := a.RemoveTopic(req)
147 if err == nil && !out.Committed {
148 err = fmt.Errorf("removal: %+v", out)
149 }
150 done <- err
151 }()
152 if err := awaitRemovalTest(t, done); err != nil {
153 t.Fatal(err)
154 }
155 tab.saveMu.Lock()
156 stopped := tab.closing && !tab.saving && tab.saveCond != nil
157 tab.saveMu.Unlock()
158 if !stopped {
159 t.Fatal("removal did not join its autosave")
160 }
161 a.lifecycleCheckpointHook = nil
162 if a.maybeAutoTitleTopic(tab) {
163 t.Fatal("removed topic was auto-titled")
164 }
165 state, err := a.workspaceRegistry().Load(t.Context())
166 if err != nil || state.SessionStates[rt.Ref().SessionID].Lifecycle != workspacestate.Archived {
167 t.Fatalf("archive state: %+v %v", state, err)
168 }
169 })
170 }
171 }
172
173 func TestPlaceholderRemovalReleasesTitleLocksBeforeCleanup(t *testing.T) {
174 a, topic, _ := topicRemovalFixture(t, "global", "Placeholder")
175 tab := &WorkspaceTab{ID: "failed", TopicID: topic.ID, Scope: "global", StartupErr: "unavailable", Ready: true}
176 a.tabs[tab.ID] = tab
177 a.tabOrder = []string{tab.ID}
178 a.activeTabID = tab.ID
179 checked := false
180 a.lifecycleCheckpointHook = func(phase string) {
181 if phase != "before-topic-runtime-cleanup" {
182 return
183 }
184 checked = true
185 for name, lock := range map[string]*sync.Mutex{"title": &a.topicTitleMutationMu, "index": &topicIndexMu} {
186 if !lock.TryLock() {
187 t.Errorf("cleanup holds %s lock", name)
188 } else {
189 lock.Unlock()
190 }
191 }
192 if !a.mu.TryLock() {
193 t.Error("cleanup holds App lock")
194 } else {
195 a.mu.Unlock()
196 }
197 if a.sessionRemovalMu.TryLock() {
198 a.sessionRemovalMu.Unlock()
199 t.Error("cleanup lost removal ownership")
200 }
201 }
202 out, err := a.RemoveTopic(inspectedTopicRemoval(t, a, topic.ID))
203 if err != nil || !out.Committed || !checked {
204 t.Fatalf("removal: %+v %v cleanup=%t", out, err, checked)
205 }
206 }
207
208 type removalBlockingSnapshot struct {
209 control.SessionAPI
210 started chan struct{}
211 resume chan struct{}
212 }
213
214 func (c *removalBlockingSnapshot) Snapshot() error {
215 close(c.started)
216 <-c.resume
217 return c.SessionAPI.Snapshot()
218 }
219
220 func TestQuiesceWaitsForFirstAutosave(t *testing.T) {
221 isolateDesktopUserDirs(t)
222 a, tab := appWithTab(t, filepath.Join(t.TempDir(), "first.jsonl"))
223 c := &removalBlockingSnapshot{SessionAPI: tab.Ctrl, started: make(chan struct{}), resume: make(chan struct{})}
224 tab.Ctrl = c
225 var release sync.Once
226 defer release.Do(func() { close(c.resume) })
227 a.scheduleTabSnapshot(tab.ID)
228 awaitRemovalTest(t, c.started)
229 tab.saveMu.Lock()
230 initialized := tab.saving && tab.saveCond != nil
231 tab.saveMu.Unlock()
232 if !initialized {
233 t.Fatal("first in-flight save has no completion condition")
234 }
235 done := make(chan struct{})
236 go func() { a.quiesceTabAutosave(tab); close(done) }()
237 // Acquiring saveMu after closing is set proves the join has reached its
238 // condition wait while the actual Snapshot is still blocked.
239 deadline := time.Now().Add(10 * time.Second)
240 for {
241 tab.saveMu.Lock()
242 closing := tab.closing
243 tab.saveMu.Unlock()
244 if closing {
245 break
246 }
247 if time.Now().After(deadline) {
248 t.Fatal("quiesce did not stop save admission")
249 }
250 time.Sleep(time.Millisecond)
251 }
252 select {
253 case <-done:
254 t.Fatal("quiesce returned before the first write finished")
255 default:
256 }
257 release.Do(func() { close(c.resume) })
258 awaitRemovalTest(t, done)
259 if tab.saving {
260 t.Fatal("save still running after join")
261 }
262 }
263
263 lines GO