返回 DeepSeek-Reasonix
topic_removal_recovery.go
根目录 / desktop / topic_removal_recovery.go
1 package main
2
3 import (
4 "context"
5 "encoding/json"
6 "errors"
7 "slices"
8 "strings"
9
10 "reasonix/desktop/internal/legacycleanup"
11 "reasonix/desktop/internal/workspacestate"
12 "reasonix/internal/topicstate"
13 )
14
15 func topicRemovalTrashEntries(state workspacestate.State, query string) []TrashEntry {
16 rows := []TrashEntry{}
17 for _, operation := range state.TopicRemovals {
18 if operation.Disposition != "archive_placeholder" || (operation.Phase != "committed" && operation.Phase != "prepared" && operation.Phase != "restoring") {
19 continue
20 }
21 var item legacycleanup.Candidate
22 if json.Unmarshal(operation.Snapshot, &item) != nil || item.Topic == nil {
23 continue
24 }
25 workspaceTitle := state.Workspaces[item.WorkspaceID].Title
26 if workspaceTitle == "" {
27 workspaceTitle = globalProjectTitle()
28 if item.Topic.Scope == "project" {
29 workspaceTitle = workspaceName(item.Topic.WorkspaceRoot)
30 }
31 }
32 if !strings.Contains(strings.ToLower(item.Title+"\n"+workspaceTitle), strings.ToLower(query)) {
33 continue
34 }
35 ready := operation.Phase == "committed"
36 row := TrashEntry{ID: "topic-removal:" + operation.ID, RecoveryEntryID: "topic-removal:" + operation.ID, Title: item.Title,
37 WorkspaceID: item.WorkspaceID, WorkspaceTitle: workspaceTitle, ArchivedAt: operation.ArchivedAt,
38 CanRestore: ready, CanPurge: ready, Health: "ready"}
39 if !ready {
40 row.Health, row.OperationPhase = "operation_pending", operation.Phase
41 }
42 rows = append(rows, row)
43 }
44 return rows
45 }
46
47 // Only explicit user removals are replayed here, never a new cleanup batch.
48 func (a *App) reconcileTopicRemovals(state workspacestate.State) error {
49 var joined error
50 for _, operation := range state.TopicRemovals {
51 switch operation.Phase {
52 case "prepared":
53 out, err := a.RemoveTopic(TopicRemovalRequest{OperationID: operation.ID, Target: TopicRemovalTarget{WorkspaceID: operation.WorkspaceID, TopicID: operation.TopicID}, ExpectedToken: operation.Token})
54 if err == nil {
55 err = topicRemovalError(out)
56 }
57 joined = errors.Join(joined, err)
58 case "restoring":
59 joined = errors.Join(joined, a.restoreRemovedTopic(operation.ID, operation.WorkspaceID))
60 }
61 }
62 return joined
63 }
64
65 func (a *App) purgeRemovedTopic(id, workspaceID string) error {
66 return a.workspaceRegistry().TransitionTopicRemoval(a.bootContext(), id, func(_ workspacestate.State, operation *workspacestate.TopicRemoval, _ func() error) error {
67 if operation.ID != id || operation.WorkspaceID != workspaceID || operation.Disposition != "archive_placeholder" {
68 return workspacestate.ErrMutationConflict
69 }
70 if operation.Phase == "purged" {
71 return nil
72 }
73 if operation.Phase != "committed" {
74 return workspacestate.ErrMutationConflict
75 }
76 operation.Phase, operation.Snapshot = "purged", nil
77 operation.Metadata = nil
78 return nil
79 })
80 }
81
82 func (a *App) restoreRemovedTopic(id, workspaceID string) error {
83 release, ok := a.tryLockRuntimeMutation("restore topic placeholder")
84 if !ok {
85 return errTopicArchiveBusy
86 }
87 defer release()
88 a.topicTitleMutationMu.Lock()
89 defer a.topicTitleMutationMu.Unlock()
90 topicIndexMu.Lock()
91 defer topicIndexMu.Unlock()
92 return a.workspaceRegistry().TransitionTopicRemoval(a.bootContext(), id, func(state workspacestate.State, operation *workspacestate.TopicRemoval, checkpoint func() error) error {
93 if operation.ID != id || operation.WorkspaceID != workspaceID || operation.Disposition != "archive_placeholder" {
94 return workspacestate.ErrMutationConflict
95 }
96 if operation.Phase == "restored" {
97 return nil
98 }
99 if operation.Phase != "committed" && operation.Phase != "restoring" {
100 return workspacestate.ErrMutationConflict
101 }
102 var item legacycleanup.Candidate
103 if err := json.Unmarshal(operation.Snapshot, &item); err != nil {
104 return err
105 }
106 if item.Topic == nil {
107 return workspacestate.ErrMutationConflict
108 }
109 var metadata topicstate.Record
110 if err := json.Unmarshal(operation.Metadata, &metadata); err != nil {
111 return err
112 }
113 if err := a.checkRemovedTopicSources(state, item); err != nil {
114 return err
115 }
116 desktopProjectsFileMu.Lock()
117 defer desktopProjectsFileMu.Unlock()
118 release, err := acquireDesktopProjectsFileLock()
119 if err != nil {
120 return err
121 }
122 defer release()
123 file, err := readTopicRemovalProjects()
124 if err != nil {
125 return err
126 }
127 if topicRemovalHasOtherOwner(file, item) {
128 return workspacestate.ErrMutationConflict
129 }
130 if item.Topic.Scope == "project" && projectIndexByRoot(file.Projects, item.Topic.WorkspaceRoot) < 0 {
131 return workspacestate.ErrWorkspaceNotFound
132 }
133 return desktopTopicState.withExclusiveScope(topicTitleRoot(item.Topic.Scope, item.Topic.WorkspaceRoot), func(ctx context.Context, store *topicstate.Store) error {
134 return a.restoreRemovedTopicMetadataLocked(ctx, store, state, file, item, metadata, operation, checkpoint)
135 })
136 })
137 }
138
139 func restoreRemovedTopicIndex(file *desktopProjectFile, item legacycleanup.Candidate) {
140 snapshot := item.Topic
141 file.DeletedTopics = removeString(file.DeletedTopics, item.TopicID)
142 if snapshot.Scope == "global" {
143 file.GlobalTopics = insertLegacyTopicAt(file.GlobalTopics, item.TopicID, snapshot.Order)
144 if snapshot.Pinned && !slices.Contains(file.GlobalPinnedTopics, item.TopicID) {
145 file.GlobalPinnedTopics = append(file.GlobalPinnedTopics, item.TopicID)
146 }
147 if groups, changed := restoreLegacyTopicGroup(file.GlobalGroups, snapshot, item.TopicID); changed {
148 file.GlobalGroups = groups
149 file.GlobalGroupsRevision++
150 }
151 return
152 }
153 i := projectIndexByRoot(file.Projects, snapshot.WorkspaceRoot)
154 project := &file.Projects[i]
155 project.Topics = insertLegacyTopicAt(project.Topics, item.TopicID, snapshot.Order)
156 if snapshot.Pinned && !slices.Contains(project.PinnedTopics, item.TopicID) {
157 project.PinnedTopics = append(project.PinnedTopics, item.TopicID)
158 }
159 if groups, changed := restoreLegacyTopicGroup(project.Groups, snapshot, item.TopicID); changed {
160 project.Groups = groups
161 project.GroupsRevision++
162 }
163 }
164
165 // Requires registry, projects and topic-store writer ownership.
166 func (a *App) restoreRemovedTopicMetadataLocked(ctx context.Context, store *topicstate.Store, state workspacestate.State, file desktopProjectFile, item legacycleanup.Candidate, metadata topicstate.Record, operation *workspacestate.TopicRemoval, checkpoint func() error) error {
167 snapshot, err := store.Snapshot(ctx)
168 if err != nil {
169 return err
170 }
171 record, exists := snapshot.Records[item.TopicID]
172 indexed := topicIndexedInProjectsSnapshot(file, item.Topic.Scope, item.Topic.WorkspaceRoot, item.TopicID)
173 if operation.Phase == "committed" {
174 if exists || indexed || !slices.Contains(file.DeletedTopics, item.TopicID) {
175 return workspacestate.ErrMutationConflict
176 }
177 operation.Phase, operation.RestoreRevision = "restoring", snapshot.State.Revision+1
178 if err := checkpoint(); err != nil {
179 return err
180 }
181 } else if exists {
182 if record.RowRevision != operation.RestoreRevision || record.Title != item.Title || record.TitleSource != item.Topic.TitleSource || record.CreatedAtMS != item.Topic.CreatedAt {
183 return workspacestate.ErrMutationConflict
184 }
185 } else if indexed {
186 return workspacestate.ErrMutationConflict
187 }
188 if !exists {
189 // Another topic may have advanced the scope revision while this restore
190 // was interrupted before its first write. No target row exists yet, so
191 // reserve the new revision durably; an existing row still uses the CAS above.
192 if operation.RestoreRevision != snapshot.State.Revision+1 {
193 operation.RestoreRevision = snapshot.State.Revision + 1
194 if err := checkpoint(); err != nil {
195 return err
196 }
197 }
198 a.lifecycleCheckpoint("topic-restore-before-metadata")
199 if _, err := store.Update(ctx, item.TopicID, func(record *topicstate.Record) {
200 *record = metadata
201 }); err != nil {
202 return err
203 }
204 }
205 if indexed {
206 // Do not overwrite organization edits made after an interrupted restore.
207 current, err := topicRemovalCandidate(state, file, TopicRemovalTarget{WorkspaceID: item.WorkspaceID, TopicID: item.TopicID})
208 if err != nil {
209 return err
210 }
211 if current.Topic.Pinned != item.Topic.Pinned {
212 return workspacestate.ErrMutationConflict
213 }
214 } else {
215 a.lifecycleCheckpoint("topic-restore-before-index")
216 restoreRemovedTopicIndex(&file, item)
217 if err := saveProjectsFile(file); err != nil {
218 return err
219 }
220 }
221 a.lifecycleCheckpoint("topic-restore-before-commit")
222 operation.Phase = "restored"
223 return nil
224 }
225
225 lines GO