返回 DeepSeek-Reasonix
topic_archive.go
根目录 / desktop / topic_archive.go
1 package main
2
3 import (
4 "errors"
5 "fmt"
6 "log/slog"
7 "os"
8 "path/filepath"
9 "strings"
10 "time"
11
12 "reasonix/internal/agent"
13 "reasonix/internal/control"
14 "reasonix/internal/store"
15 )
16
17 var (
18 errTopicHasActiveWork = errors.New("wait for the session to finish, answer pending prompts, and stop background jobs before archiving this topic")
19 errTopicArchiveBusy = errors.New("Reasonix is finishing another session change — wait a moment and retry archiving")
20 )
21
22 var topicArchiveCleanupHookForTest func() error
23
24 type topicArchiveTrace struct {
25 phase string
26 targetCount int
27 runtimeCount int
28 }
29
30 func (a *App) TrashTopic(topicID string) error {
31 if key, historical := strings.CutPrefix(topicID, "historical-"); historical {
32 if _, err := a.ListHistoricalSessions(); err != nil {
33 return err
34 }
35 c := &a.historicalImports
36 c.mu.Lock()
37 source, found := c.sources[key]
38 c.mu.Unlock()
39 if !found {
40 return newSessionOperationError(sessionOperationTargetNotFound, "The historical source is unavailable.")
41 }
42 _, err := a.ArchiveSessionTarget(SessionSelector{Source: &SessionSourceRef{Path: source.path, SourceKey: key, HeadID: source.head}})
43 return err
44 }
45 return friendlySessionFileError(a.archiveCompatibleTopic(topicID))
46 }
47
48 func (a *App) topicHasActiveRuntimeWork(topicID string) bool {
49 a.mu.RLock()
50 defer a.mu.RUnlock()
51 for _, tabs := range []map[string]*WorkspaceTab{a.tabs, a.detachedSessions} {
52 for _, tab := range tabs {
53 if tab != nil && tab.TopicID == topicID && tab.hasActiveRuntimeWork() {
54 return true
55 }
56 }
57 }
58 return false
59 }
60
61 func (a *App) trashTopic(topicID string) (retErr error) {
62 topicID = strings.TrimSpace(topicID)
63 if topicID == "" {
64 return fmt.Errorf("topicID is required")
65 }
66 started := time.Now()
67 trace := topicArchiveTrace{phase: "start"}
68 defer func() {
69 outcome := "ok"
70 if retErr != nil {
71 outcome = "failed"
72 }
73 slog.Debug("desktop: topic archive timing", "outcome", outcome, "phase", trace.phase,
74 "total_ms", time.Since(started).Milliseconds(), "target_count", trace.targetCount, "runtime_count", trace.runtimeCount)
75 }()
76 fallback, changedDirs, err := a.commitTopicArchive(topicID, &trace)
77 if err != nil {
78 return err
79 }
80 if fallback.workspaceRoot != "" {
81 changedDirs = append(changedDirs, desktopSessionDir(fallback.workspaceRoot))
82 } else if fallback.needs {
83 changedDirs = append(changedDirs, desktopSessionDir(globalWorkspaceRoot()))
84 }
85 // The last visible topic leaves no replacement runtime (the frontend lands
86 // on the workspace draft). Abandoned transient blanks still go, or
87 // reconciliation could promote a default-titled blank into the registry.
88 keepPath := ""
89 a.mu.RLock()
90 if tab := a.tabs[a.activeTabID]; tab != nil {
91 keepPath = tab.SessionPath
92 }
93 a.mu.RUnlock()
94 trace.phase = "discard_unused_blanks"
95 a.discardUnusedTransientBlankSessions(changedDirs, keepPath)
96 trace.phase = "notify"
97 if len(changedDirs) > 0 {
98 a.emitProjectTreeChangedForSessionDirs(changedDirs...)
99 } else {
100 a.emitProjectTreeMetadataChanged()
101 }
102 return nil
103 }
104
105 func (a *App) commitTopicArchive(topicID string, trace *topicArchiveTrace) (fallbackRuntimeTarget, []string, error) {
106 trace.phase = "runtime_lock"
107 releaseRuntime, ok := a.tryLockRuntimeMutationBounded("trash-topic")
108 if !ok {
109 return fallbackRuntimeTarget{}, nil, errTopicArchiveBusy
110 }
111 defer releaseRuntime()
112 trace.phase = "removal_lock"
113 if !a.sessionRemovalMu.TryLock() {
114 return fallbackRuntimeTarget{}, nil, errTopicArchiveBusy
115 }
116 defer a.sessionRemovalMu.Unlock()
117 trace.phase = "active_work_check"
118 if a.topicHasActiveRuntimeWork(topicID) {
119 return fallbackRuntimeTarget{}, nil, errTopicHasActiveWork
120 }
121 trace.phase = "target_scan"
122 targets, err := a.topicTrashTargets(topicID)
123 if err != nil {
124 return fallbackRuntimeTarget{}, nil, err
125 }
126 trace.targetCount = len(targets)
127 changedDirs := make([]string, 0, len(targets))
128 for _, target := range targets {
129 changedDirs = append(changedDirs, target.dir)
130 }
131 trace.phase = "snapshot"
132 removed := a.captureTopicRuntimeBindings(topicID)
133 trace.runtimeCount = len(removed)
134 if err := a.snapshotTopicRuntimeBindings(removed); err != nil {
135 return fallbackRuntimeTarget{}, nil, err
136 }
137 trace.phase = "acquire_removal_ownership"
138 ownership, err := acquireTopicArchiveOwnership(targets, removed)
139 if err != nil {
140 return fallbackRuntimeTarget{}, nil, err
141 }
142 defer ownership.release()
143 trace.phase = "mark_cleanup_pending"
144 rollbackMarkers, err := markTopicArchiveCleanupPending(topicID, targets)
145 if err != nil {
146 ownership.rollback()
147 return fallbackRuntimeTarget{}, nil, err
148 }
149 trace.phase = "detach_runtimes"
150 fallback, unchanged := a.removeTopicRuntimeBindingsIfUnchanged(topicID, removed)
151 if !unchanged {
152 rollbackMarkers()
153 ownership.rollback()
154 return fallbackRuntimeTarget{}, nil, errTopicArchiveBusy
155 }
156 a.finalizeRemovedTopicRuntimes(removed)
157 destroyBegun := false
158 closedRemoved := map[control.SessionAPI]bool{}
159 defer func() {
160 if destroyBegun {
161 a.closeRemainingRemovedSessionRuntimesAfterDestroyAdmissionHeld(removed, closedRemoved)
162 } else {
163 a.closeRemainingRemovedSessionRuntimesAdmissionHeld(removed, closedRemoved)
164 }
165 }()
166 trace.phase = "teardown"
167 destroyBatches := make([][]control.SessionDestroyHandle, len(targets))
168 for i, target := range targets {
169 destroys := a.destroyHandlesForSession(target.dir, target.sessionPath, removed)
170 destroyBatches[i] = destroys
171 destroyBegun = destroyBegun || len(destroys) > 0
172 }
173 timedOutTargets := waitDestroyHandleBatches(destroyBatches)
174 for i, target := range targets {
175 destroys := destroyBatches[i]
176 a.closeRemovedSessionRuntimesForSessionAfterDestroyAdmissionHeld(removed, target.dir, target.sessionPath, closedRemoved)
177 a.removeSessionCatalogPath(target.sessionPath, "topic_archived")
178 if timedOutTargets[i] {
179 guard := ownership.take(target.sessionPath)
180 go delayedDesktopTopicTrash(target.dir, target.sessionPath, target.key, guard, destroys)
181 continue
182 }
183 trace.phase = "move_artifacts"
184 var err error
185 if hook := topicArchiveCleanupHookForTest; hook != nil {
186 err = hook()
187 }
188 if err == nil {
189 err = trashSessionArtifactsWithGuard(target.dir, target.sessionPath, target.key, ownership.take(target.sessionPath))
190 }
191 finishDestroyHandles(destroys)
192 if err != nil {
193 // Cleanup-pending is the durable commit point. Once bindings have
194 // detached, report the archive as accepted and let startup
195 // reconciliation finish any filesystem operation that could not.
196 slog.Warn("desktop: topic archive cleanup remains pending")
197 }
198 }
199 trace.phase = "delete_topic_metadata"
200 if err := a.deleteTopic(topicID); err != nil {
201 slog.Warn("desktop: topic archive metadata cleanup remains pending")
202 } else if err := clearTopicArchiveMetadataPending(topicID); err != nil {
203 slog.Warn("desktop: topic archive metadata marker cleanup remains pending")
204 }
205 return fallback, changedDirs, nil
206 }
207
208 func markTopicArchiveCleanupPending(topicID string, targets []topicTrashTarget) (func(), error) {
209 if err := markTopicArchiveMetadataPending(topicID, targets); err != nil {
210 return nil, err
211 }
212 marked := make([]string, 0, len(targets))
213 rollback := func() {
214 for _, path := range marked {
215 if err := agent.ClearCleanupPending(path); err != nil {
216 slog.Warn("desktop: rollback topic archive marker failed")
217 }
218 }
219 if err := clearTopicArchiveMetadataPending(topicID); err != nil {
220 slog.Warn("desktop: rollback topic archive metadata marker failed")
221 }
222 }
223 for _, target := range targets {
224 if err := agent.MarkCleanupPending(target.sessionPath, "delete"); err != nil {
225 rollback()
226 return rollback, err
227 }
228 marked = append(marked, target.sessionPath)
229 }
230 return rollback, nil
231 }
232
233 type topicTrashTarget struct {
234 dir string
235 sessionPath string
236 key string
237 }
238
239 func (a *App) topicTrashTargets(topicID string) ([]topicTrashTarget, error) {
240 topicID = strings.TrimSpace(topicID)
241 var targets []topicTrashTarget
242 seen := map[string]bool{}
243 addTarget := func(dir, path string) error {
244 sessionPath, key, err := validateSessionPath(dir, path)
245 if err != nil {
246 return err
247 }
248 id := dir + "\x00" + sessionPath
249 if seen[id] {
250 return nil
251 }
252 seen[id] = true
253 if err := validateSessionTrashTarget(dir, sessionPath, key); err != nil {
254 return err
255 }
256 targets = append(targets, topicTrashTarget{dir: dir, sessionPath: sessionPath, key: key})
257 return nil
258 }
259 for _, dir := range a.knownSessionDirs() {
260 index, err := topicSessionIndexForDir(dir)
261 if err != nil {
262 return nil, err
263 }
264 for _, match := range index.byTopic[topicID] {
265 if !agent.IsCleanupPending(match.path) {
266 if err := addTarget(dir, match.path); err != nil {
267 return nil, err
268 }
269 }
270 }
271 }
272 a.mu.RLock()
273 var runtimeTargets []struct{ dir, path string }
274 for _, tab := range a.runtimeTabsLocked() {
275 if tab == nil || tab.TopicID != topicID {
276 continue
277 }
278 if path := canonicalTabSessionPath(tab.currentSessionPath()); path != "" {
279 dir := tabSessionDir(tab)
280 if filepath.IsAbs(path) {
281 dir = filepath.Dir(path)
282 }
283 runtimeTargets = append(runtimeTargets, struct{ dir, path string }{dir: dir, path: path})
284 }
285 }
286 a.mu.RUnlock()
287 for _, target := range runtimeTargets {
288 if err := addTarget(target.dir, target.path); err != nil {
289 return nil, err
290 }
291 }
292 return targets, nil
293 }
294
295 func topicIndexedInRegistry(scope, workspaceRoot, topicID string) bool {
296 topicID = strings.TrimSpace(topicID)
297 if topicID == "" {
298 return false
299 }
300 if strings.TrimSpace(loadTopicTitles(topicTitleRoot(scope, workspaceRoot))[topicID]) != "" {
301 return true
302 }
303 f := loadProjectsFile()
304 if scope != "project" {
305 return containsDesktopString(f.GlobalTopics, topicID)
306 }
307 if i := projectIndexByRoot(f.Projects, workspaceRoot); i >= 0 {
308 return containsDesktopString(f.Projects[i].Topics, topicID)
309 }
310 return false
311 }
312
313 func (a *App) discardUnusedTransientBlankSessions(dirs []string, keepPath string) {
314 // Forced legacy repair promotes indexed sidecars under this same lock. Hold
315 // it through classification, deletion, and scoped registry cleanup so a
316 // concurrent reconcile cannot reinsert a zero-byte ghost after we sweep it.
317 legacyMigrationMu.Lock()
318 defer legacyMigrationMu.Unlock()
319
320 keepPath = canonicalTabSessionPath(strings.TrimSpace(keepPath))
321 kept := map[string]bool{}
322 if keepPath != "" {
323 kept[keepPath] = true
324 }
325 if a != nil {
326 a.mu.RLock()
327 for _, tabs := range []map[string]*WorkspaceTab{a.tabs, a.detachedSessions} {
328 for _, tab := range tabs {
329 if tab == nil {
330 continue
331 }
332 if path := canonicalTabSessionPath(strings.TrimSpace(tab.SessionPath)); path != "" {
333 kept[path] = true
334 }
335 }
336 }
337 a.mu.RUnlock()
338 }
339 siblingDirs := append([]string(nil), dirs...)
340 if a != nil {
341 siblingDirs = append(siblingDirs, a.knownSessionDirs()...)
342 }
343 seen := map[string]bool{}
344 for _, dir := range dirs {
345 dir = strings.TrimSpace(dir)
346 if dir == "" || seen[dir] {
347 continue
348 }
349 seen[dir] = true
350 entries, err := os.ReadDir(dir)
351 if err != nil {
352 continue
353 }
354 for _, entry := range entries {
355 name := entry.Name()
356 if entry.IsDir() || !store.IsSessionTranscriptName(name) {
357 continue
358 }
359 path := filepath.Join(dir, name)
360 if kept[canonicalTabSessionPath(path)] {
361 continue
362 }
363 if !unusedTransientBlankSession(dir, path) {
364 continue
365 }
366 meta, hasMeta, _ := agent.LoadBranchMeta(path)
367 removed := discardTransientBlankSessionArtifacts(path)
368 if removed && hasMeta && !transientTopicHasSibling(siblingDirs, path, meta) {
369 cleanupTransientBlankTopicRegistration(meta)
370 }
371 if removed && a != nil {
372 a.removeSessionCatalogPath(path, "transient_blank_discarded")
373 }
374 }
375 }
376 }
377
378 func unusedTransientBlankSession(dir, path string) bool {
379 resolved, ok := pinnedTabSessionPath(dir, path)
380 if !ok {
381 return false
382 }
383 info, err := os.Stat(resolved)
384 if err != nil || info.IsDir() || info.Size() != 0 {
385 return false
386 }
387 meta, ok, err := agent.LoadBranchMeta(resolved)
388 if err != nil || !ok {
389 return true
390 }
391 topicID := strings.TrimSpace(meta.TopicID)
392 if topicID == "" {
393 return true
394 }
395 if !isDefaultTopicTitle(meta.TopicTitle) && strings.TrimSpace(meta.TopicTitle) != "" {
396 return false
397 }
398 // A zero-byte default sidecar stays transient after registry projection.
399 // The caller's keep set protects every visible or detached runtime; registry
400 // presence alone cannot prove user content.
401 return true
402 }
403
404 func transientTopicHasSibling(dirs []string, excludedPath string, target agent.BranchMeta) bool {
405 seen := make(map[string]bool, len(dirs))
406 for _, dir := range dirs {
407 dir = strings.TrimSpace(dir)
408 key := projectRootKey(dir)
409 if dir == "" || seen[key] {
410 continue
411 }
412 seen[key] = true
413 entries, err := os.ReadDir(dir)
414 if err != nil {
415 continue
416 }
417 for _, entry := range entries {
418 if entry.IsDir() || !store.IsSessionTranscriptName(entry.Name()) {
419 continue
420 }
421 path := filepath.Join(dir, entry.Name())
422 if sameDesktopPath(path, excludedPath) {
423 continue
424 }
425 meta, ok, err := agent.LoadBranchMeta(path)
426 if err != nil || !ok || strings.TrimSpace(meta.TopicID) != strings.TrimSpace(target.TopicID) {
427 continue
428 }
429 sameRoot := meta.DefaultScope() != "project" || sameProjectRoot(meta.WorkspaceRoot, target.WorkspaceRoot)
430 if meta.DefaultScope() == target.DefaultScope() && sameRoot {
431 return true
432 }
433 }
434 }
435 return false
436 }
437
438 func cleanupTransientBlankTopicRegistration(meta agent.BranchMeta) {
439 topicID := strings.TrimSpace(meta.TopicID)
440 if topicID == "" {
441 return
442 }
443 scope, root := meta.DefaultScope(), normalizeProjectRoot(meta.WorkspaceRoot)
444 _ = updateProjectsFile(func(f *desktopProjectFile) (bool, error) {
445 changed := false
446 if scope != "project" {
447 if next := removeString(f.GlobalTopics, topicID); !sameStringList(next, f.GlobalTopics) {
448 f.GlobalTopics, changed = next, true
449 }
450 if next := removeString(f.GlobalPinnedTopics, topicID); !sameStringList(next, f.GlobalPinnedTopics) {
451 f.GlobalPinnedTopics, changed = next, true
452 }
453 if next, removed := groupsWithoutTopic(f.GlobalGroups, topicID); removed {
454 f.GlobalGroups, f.GlobalGroupsRevision, changed = next, f.GlobalGroupsRevision+1, true
455 }
456 return changed, nil
457 }
458 if index := projectIndexByRoot(f.Projects, root); index >= 0 {
459 project := &f.Projects[index]
460 if next := removeString(project.Topics, topicID); !sameStringList(next, project.Topics) {
461 project.Topics, changed = next, true
462 }
463 if next := removeString(project.PinnedTopics, topicID); !sameStringList(next, project.PinnedTopics) {
464 project.PinnedTopics, changed = next, true
465 }
466 if next, removed := groupsWithoutTopic(project.Groups, topicID); removed {
467 project.Groups, project.GroupsRevision, changed = next, project.GroupsRevision+1, true
468 }
469 }
470 return changed, nil
471 })
472 titleRoot := topicTitleRoot(scope, root)
473 _ = deleteTopicState(titleRoot, topicID)
474 }
475
475 lines GO