返回 DeepSeek-Reasonix
legacy_empty_session_cleanup.go
根目录 / desktop / legacy_empty_session_cleanup.go
1 package main
2
3 import (
4 "context"
5 "crypto/sha256"
6 "encoding/hex"
7 "encoding/json"
8 "errors"
9 "fmt"
10 "io"
11 "log/slog"
12 "os"
13 "path/filepath"
14 "slices"
15 "sync"
16 "time"
17
18 "reasonix/desktop/internal/legacycleanup"
19 "reasonix/desktop/internal/workspacestate"
20 "reasonix/internal/agent"
21 "reasonix/internal/filelock"
22 "reasonix/internal/session"
23 "reasonix/internal/store"
24 )
25
26 const legacyEmptySessionCleanupEvent = "legacy-empty-session-cleanup:changed"
27
28 var errLegacyCleanupStateChanged = errors.New("legacy cleanup candidate changed after inspection")
29
30 type LegacyEmptySessionCleanupItem struct {
31 ID string `json:"id"`
32 Kind string `json:"kind"`
33 WorkspaceID string `json:"workspaceId,omitempty"`
34 SessionID string `json:"sessionId,omitempty"`
35 TopicID string `json:"topicId,omitempty"`
36 Title string `json:"title,omitempty"`
37 Phase string `json:"phase"`
38 Classification string `json:"classification,omitempty"`
39 Reason string `json:"reason,omitempty"`
40 }
41
42 type LegacyEmptySessionCleanupStatus struct {
43 Version int `json:"version"`
44 BatchID string `json:"batchId,omitempty"`
45 State string `json:"state"`
46 Removed int `json:"removed"`
47 Pending int `json:"pending"`
48 Busy int `json:"busy"`
49 Unknown int `json:"unknown"`
50 Protected int `json:"protected"`
51 HasContent int `json:"hasContent"`
52 Items []LegacyEmptySessionCleanupItem `json:"items"`
53 }
54
55 type legacyCleanupDecision struct {
56 classification string
57 reason string
58 info session.SessionInfo
59 snapshot session.Snapshot
60 }
61
62 type legacyCleanupWorkerState struct {
63 mu sync.Mutex
64 running bool
65 // beforeArchive is a deterministic race-test hook. Set before cleanup and
66 // never mutate it concurrently.
67 beforeArchive func()
68 }
69
70 func legacyCleanupBatchID(state workspacestate.State) string {
71 ids := make([]string, 0, len(state.SessionStates))
72 for id := range state.SessionStates {
73 ids = append(ids, id)
74 }
75 slices.Sort(ids)
76 body, _ := json.Marshal(struct {
77 Generation uint64
78 IDs []string
79 At int64
80 }{state.Generation, ids, time.Now().UTC().UnixNano()})
81 sum := sha256.Sum256(body)
82 return hex.EncodeToString(sum[:12])
83 }
84
85 func legacyCleanupSourcePaths(sessionPath string) ([]string, error) {
86 paths := make([]string, 0, 24)
87 for _, artifact := range sessionTrashArtifacts(sessionPath, filepath.Base(sessionPath)) {
88 paths = append(paths, artifact.src)
89 }
90 subagents, err := agent.ListSubagentsByParent(filepath.Dir(sessionPath), agent.BranchID(sessionPath))
91 if err != nil {
92 return nil, err
93 }
94 for _, artifact := range subagents {
95 paths = append(paths, artifact.SessionPath, artifact.MetaPath)
96 paths = append(paths, store.SessionSidecarFiles(artifact.SessionPath)...)
97 paths = append(paths,
98 store.SessionCheckpointDir(artifact.SessionPath),
99 store.SessionJobsDir(artifact.SessionPath),
100 store.SessionInboxDir(artifact.SessionPath),
101 )
102 }
103 paths = uniqueStrings(paths)
104 slices.Sort(paths)
105 return paths, nil
106 }
107
108 // legacyCleanupSourceFingerprint covers every durable artifact that can make a
109 // legacy session recoverable. It deliberately excludes lease and lock files.
110 func legacyCleanupSourceFingerprint(sessionPath string) (string, error) {
111 paths, err := legacyCleanupSourcePaths(sessionPath)
112 if err != nil {
113 return "", err
114 }
115 h := sha256.New()
116 found := false
117 for _, path := range paths {
118 info, err := os.Lstat(path)
119 if os.IsNotExist(err) {
120 continue
121 }
122 if err != nil {
123 return "", err
124 }
125 if info.Mode()&os.ModeSymlink != 0 {
126 return "", errors.New("legacy cleanup source contains a symbolic link")
127 }
128 found = true
129 root := filepath.Dir(sessionPath)
130 if info.IsDir() {
131 err = filepath.WalkDir(path, func(child string, entry os.DirEntry, walkErr error) error {
132 if walkErr != nil {
133 return walkErr
134 }
135 childInfo, infoErr := entry.Info()
136 if infoErr != nil {
137 return infoErr
138 }
139 if childInfo.Mode()&os.ModeSymlink != 0 {
140 return errors.New("legacy cleanup source contains a symbolic link")
141 }
142 rel, relErr := filepath.Rel(root, child)
143 if relErr != nil {
144 return relErr
145 }
146 fmt.Fprintf(h, "%s\x00%d\x00", filepath.ToSlash(rel), childInfo.Size())
147 if entry.IsDir() {
148 return nil
149 }
150 if !childInfo.Mode().IsRegular() {
151 return errors.New("legacy cleanup source contains a non-regular file")
152 }
153 file, openErr := os.Open(child)
154 if openErr != nil {
155 return openErr
156 }
157 _, copyErr := io.Copy(h, file)
158 closeErr := file.Close()
159 return errors.Join(copyErr, closeErr)
160 })
161 if err != nil {
162 return "", err
163 }
164 continue
165 }
166 if !info.Mode().IsRegular() {
167 return "", errors.New("legacy cleanup source contains a non-regular file")
168 }
169 rel, err := filepath.Rel(root, path)
170 if err != nil {
171 return "", err
172 }
173 fmt.Fprintf(h, "%s\x00%d\x00", filepath.ToSlash(rel), info.Size())
174 file, err := os.Open(path)
175 if err != nil {
176 return "", err
177 }
178 _, copyErr := io.Copy(h, file)
179 closeErr := file.Close()
180 if err := errors.Join(copyErr, closeErr); err != nil {
181 return "", err
182 }
183 }
184 if !found {
185 return "", os.ErrNotExist
186 }
187 return hex.EncodeToString(h.Sum(nil)), nil
188 }
189
190 // initializeLegacyEmptySessionCleanupBatch freezes identities before the
191 // renderer can create a new draft-backed session. Content inspection remains a
192 // background operation after migration and draft recovery.
193 func (a *App) initializeLegacyEmptySessionCleanupBatch() error {
194 if a == nil || a.legacyCleanup == nil {
195 return errors.New("legacy cleanup store is unavailable")
196 }
197 if _, err := a.legacyCleanup.Load(a.bootContext()); err == nil {
198 return nil
199 } else if !errors.Is(err, legacycleanup.ErrNotInitialized) {
200 return err
201 }
202 state, err := a.workspaceRegistry().Load(a.bootContext())
203 if err != nil {
204 return err
205 }
206 builder := newLegacyCleanupBatchBuilder(a, state)
207 builder.registerCanonicalSessions(a.bootContext())
208 projects := loadProjectsFile()
209 builder.registerTopics("global", "", workspacestate.GlobalWorkspaceID, projects.GlobalTopics, projects.GlobalPinnedTopics, projects.GlobalGroups)
210 for _, project := range projects.Projects {
211 workspaceID := builder.workspaceIDForRoot(project.Root)
212 if workspaceID != "" {
213 builder.registerTopics("project", project.Root, workspaceID, project.Topics, project.PinnedTopics, project.Groups)
214 }
215 }
216 _, _, err = a.legacyCleanup.Initialize(a.bootContext(), legacycleanup.State{
217 BatchID: legacyCleanupBatchID(state),
218 Items: builder.items,
219 })
220 return err
221 }
222
223 func (a *App) GetLegacyEmptySessionCleanupStatus() (LegacyEmptySessionCleanupStatus, error) {
224 state, err := a.legacyCleanup.Load(a.bootContext())
225 if errors.Is(err, legacycleanup.ErrNotInitialized) {
226 return LegacyEmptySessionCleanupStatus{Version: 1, State: "not_initialized", Items: []LegacyEmptySessionCleanupItem{}}, nil
227 }
228 if err != nil {
229 return LegacyEmptySessionCleanupStatus{}, err
230 }
231 return legacyCleanupStatus(state), nil
232 }
233
234 func legacyCleanupStatus(state legacycleanup.State) LegacyEmptySessionCleanupStatus {
235 out := LegacyEmptySessionCleanupStatus{Version: state.Version, BatchID: state.BatchID, State: "complete", Items: []LegacyEmptySessionCleanupItem{}}
236 for _, item := range legacycleanup.SortedItems(state) {
237 out.Items = append(out.Items, LegacyEmptySessionCleanupItem{ID: item.ID, Kind: item.Kind, WorkspaceID: item.WorkspaceID, SessionID: item.SessionID, TopicID: item.TopicID, Title: item.Title, Phase: item.Phase, Classification: item.Classification, Reason: item.Reason})
238 switch item.Phase {
239 case "archived":
240 out.Removed++
241 case "busy":
242 out.Busy++
243 out.Pending++
244 case "unknown", "registered", "verified", "archive_pending":
245 out.Unknown++
246 out.Pending++
247 case "protected", "restored":
248 out.Protected++
249 case "has_content":
250 out.HasContent++
251 }
252 }
253 if out.Pending > 0 {
254 out.State = "pending"
255 }
256 return out
257 }
258
259 func (a *App) RetryLegacyEmptySessionCleanup() (LegacyEmptySessionCleanupStatus, error) {
260 status, err := a.GetLegacyEmptySessionCleanupStatus()
261 if err != nil {
262 return status, err
263 }
264 return status, errors.New("automatic empty conversation cleanup has been retired; existing trash entries can still be restored")
265 }
266
267 func (a *App) runLegacyEmptySessionCleanup(includeUnknown bool) {
268 if a == nil || a.legacyCleanup == nil {
269 return
270 }
271 a.legacyCleanupWorker.mu.Lock()
272 if a.legacyCleanupWorker.running {
273 a.legacyCleanupWorker.mu.Unlock()
274 return
275 }
276 a.legacyCleanupWorker.running = true
277 a.legacyCleanupWorker.mu.Unlock()
278 defer func() {
279 a.legacyCleanupWorker.mu.Lock()
280 a.legacyCleanupWorker.running = false
281 a.legacyCleanupWorker.mu.Unlock()
282 }()
283 releaseWorker, err := a.legacyCleanup.TryAcquireWorker()
284 if err != nil {
285 if !errors.Is(err, filelock.ErrHeld) {
286 slog.Warn("desktop: legacy empty session cleanup worker unavailable", "err", err)
287 }
288 return
289 }
290 defer releaseWorker()
291 if a.desktopMigrationDone != nil {
292 select {
293 case <-a.desktopMigrationDone:
294 case <-a.bootContext().Done():
295 return
296 }
297 }
298 a.mu.RLock()
299 tabsRestored := a.tabsRestored
300 a.mu.RUnlock()
301 if tabsRestored != nil {
302 select {
303 case <-tabsRestored:
304 case <-a.bootContext().Done():
305 return
306 }
307 }
308 state, err := a.legacyCleanup.Load(a.bootContext())
309 if err != nil {
310 if !errors.Is(err, legacycleanup.ErrNotInitialized) {
311 slog.Warn("desktop: legacy empty session cleanup disabled", "err", err)
312 }
313 return
314 }
315 before := legacyCleanupStatus(state).Removed
316 for _, item := range legacycleanup.SortedItems(state) {
317 if item.Restored || item.Phase == "archived" || item.Phase == "has_content" || item.Phase == "protected" {
318 continue
319 }
320 if item.Phase == "unknown" && !includeUnknown {
321 continue
322 }
323 switch item.Kind {
324 case "session":
325 a.processLegacyCleanupSession(item)
326 case "legacy":
327 a.processLegacyCleanupSource(item)
328 case "topic":
329 a.processLegacyCleanupTopic(item)
330 }
331 }
332 after, err := a.GetLegacyEmptySessionCleanupStatus()
333 if err == nil {
334 a.emitRuntimeEvent(legacyEmptySessionCleanupEvent, after)
335 if after.Removed > before {
336 a.emitProjectTreeChanged()
337 }
338 }
339 }
340
341 func (a *App) reconcileLegacyCleanupArchivedOperation(item legacycleanup.Candidate, sessionID string) bool {
342 state, err := a.workspaceRegistry().Load(a.bootContext())
343 if err != nil {
344 return false
345 }
346 op, ok := state.PendingOperations[item.OperationID]
347 if !ok || op.Kind != "archive" || op.Phase != "committed" || !slices.Contains(op.SessionIDs, sessionID) ||
348 state.SessionStates[sessionID].Lifecycle != workspacestate.Archived {
349 return false
350 }
351 a.updateLegacyCleanupItem(item.ID, func(next *legacycleanup.Candidate) {
352 next.SessionID = sessionID
353 next.Phase, next.Classification, next.Reason = "archived", "empty", ""
354 if next.ArchivedAt == 0 {
355 next.ArchivedAt = state.SessionStates[sessionID].ArchivedAt
356 if next.ArchivedAt == 0 {
357 next.ArchivedAt = time.Now().UTC().UnixMilli()
358 }
359 }
360 })
361 return true
362 }
363
364 func (a *App) updateLegacyCleanupItem(id string, update func(*legacycleanup.Candidate)) {
365 _, err := a.legacyCleanup.Update(a.bootContext(), func(state *legacycleanup.State) error {
366 item, ok := state.Items[id]
367 if !ok || item.Restored {
368 return errLegacyCleanupStateChanged
369 }
370 update(&item)
371 state.Items[id] = item
372 return nil
373 })
374 if err != nil && !errors.Is(err, errLegacyCleanupStateChanged) {
375 slog.Warn("desktop: legacy cleanup state update failed", "err", err)
376 }
377 }
378
379 func (a *App) processLegacyCleanupSession(item legacycleanup.Candidate) {
380 if a.reconcileLegacyCleanupArchivedOperation(item, item.SessionID) {
381 return
382 }
383 ref := session.SessionRef{HostID: localDesktopHostID, SessionID: item.SessionID}
384 decision := a.classifyLegacyCleanupSession(a.bootContext(), ref, item)
385 if decision.classification != "empty" {
386 phase := decision.classification
387 if phase == "empty" || phase == "" {
388 phase = "unknown"
389 }
390 a.updateLegacyCleanupItem(item.ID, func(next *legacycleanup.Candidate) {
391 next.Phase, next.Classification, next.Reason = phase, decision.classification, decision.reason
392 })
393 return
394 }
395 if a.legacyCleanupWorker.beforeArchive != nil {
396 a.legacyCleanupWorker.beforeArchive()
397 }
398 release, ok := a.tryLockRuntimeMutation("legacy empty session cleanup")
399 if !ok {
400 a.updateLegacyCleanupItem(item.ID, func(next *legacycleanup.Candidate) {
401 next.Phase, next.Classification, next.Reason = "busy", "busy", "runtime_mutation"
402 })
403 return
404 }
405 defer release()
406 verify := func(ctx context.Context, latest workspacestate.State) error {
407 fresh := a.classifyLegacyCleanupSession(ctx, ref, item)
408 if fresh.classification != "empty" {
409 return fmt.Errorf("%w: %s", errLegacyCleanupStateChanged, fresh.classification)
410 }
411 return nil
412 }
413 err := a.archiveSessionRefsWithOperationConditional([]session.SessionRef{ref}, item.OperationID, verify)
414 if err != nil {
415 classification, reason := "unknown", "archive_failed"
416 if errors.Is(err, errTopicHasActiveWork) || errors.Is(err, errTopicArchiveBusy) {
417 classification, reason = "busy", "runtime_active"
418 } else if errors.Is(err, errLegacyCleanupStateChanged) || errors.Is(err, workspacestate.ErrMutationConflict) {
419 classification, reason = "protected", "state_changed"
420 }
421 a.updateLegacyCleanupItem(item.ID, func(next *legacycleanup.Candidate) {
422 next.Phase, next.Classification, next.Reason = classification, classification, reason
423 })
424 return
425 }
426 a.updateLegacyCleanupItem(item.ID, func(next *legacycleanup.Candidate) {
427 next.Phase, next.Classification, next.Reason, next.ArchivedAt = "archived", "empty", "", time.Now().UTC().UnixMilli()
428 })
429 }
430
430 lines GO