返回 DeepSeek-Reasonix
topic_removal.go
根目录 / desktop / topic_removal.go
1 package main
2
3 import (
4 "context"
5 "crypto/sha256"
6 "encoding/hex"
7 "encoding/json"
8 "errors"
9 "fmt"
10 "os"
11 "path/filepath"
12 "slices"
13 "strings"
14 "time"
15
16 "reasonix/desktop/internal/legacycleanup"
17 "reasonix/desktop/internal/workspacestate"
18 "reasonix/internal/topicstate"
19 )
20
21 type TopicRemovalTarget struct {
22 WorkspaceID string `json:"workspaceId"`
23 TopicID string `json:"topicId"`
24 }
25 type TopicRemovalInspection struct {
26 Target TopicRemovalTarget `json:"target"`
27 Disposition string `json:"disposition"`
28 Allowed bool `json:"allowed"`
29 Reason string `json:"reason,omitempty"`
30 Token string `json:"token"`
31 }
32 type TopicRemovalRequest struct {
33 OperationID string `json:"operationId"`
34 Target TopicRemovalTarget `json:"target"`
35 ExpectedToken string `json:"expectedToken"`
36 }
37 type TopicRemovalResult struct {
38 Committed bool `json:"committed"`
39 Disposition string `json:"disposition"`
40 RecoveryEntryID string `json:"recoveryEntryId,omitempty"`
41 ErrorCode string `json:"errorCode,omitempty"`
42 ErrorMessage string `json:"errorMessage,omitempty"`
43 Retryable bool `json:"retryable"`
44 }
45
46 // Removal must not inherit the display reader's best-effort empty fallback.
47 func readTopicRemovalProjects() (desktopProjectFile, error) {
48 var file desktopProjectFile
49 body, err := readFileUTF8(filepath.Join(desktopConfigDir(), desktopProjectsFile))
50 if err != nil && !os.IsNotExist(err) {
51 return file, err
52 }
53 if err == nil {
54 if err = json.Unmarshal(body, &file); err != nil {
55 return file, err
56 }
57 }
58 file = normalizeProjectsFile(file)
59 body, err = readFileUTF8(filepath.Join(desktopConfigDir(), desktopProjectOrganizationFile))
60 if os.IsNotExist(err) {
61 return file, nil
62 }
63 if err != nil {
64 return file, err
65 }
66 var organization desktopProjectOrganizationFileData
67 if err := json.Unmarshal(body, &organization); err != nil {
68 return file, err
69 }
70 if organization.Version < 1 || organization.Version > desktopProjectOrganizationVersion {
71 return file, errors.New("unsupported topic organization version")
72 }
73 return applyProjectOrganization(file, organization), nil
74 }
75
76 func topicRemovalWorkspaceID(state workspacestate.State, scope, root string) string {
77 for id, workspace := range state.Workspaces {
78 if sameDesktopPath(workspace.Root, desktopWorkspaceRoot(scope, root)) {
79 return id
80 }
81 }
82 return desktopWorkspaceID(scope, root)
83 }
84
85 func topicRemovalCandidate(state workspacestate.State, file desktopProjectFile, target TopicRemovalTarget) (legacycleanup.Candidate, error) {
86 var candidates []legacycleanup.Candidate
87 add := func(scope, root string, ids, pinned []string, groups []desktopGroup) {
88 order := slices.Index(ids, target.TopicID)
89 if order < 0 {
90 return
91 }
92 wid := topicRemovalWorkspaceID(state, scope, root)
93 if target.WorkspaceID != "" && target.WorkspaceID != wid {
94 return
95 }
96 candidates = append(candidates, legacycleanup.Candidate{Kind: "topic", WorkspaceID: wid, TopicID: target.TopicID,
97 Topic: legacyCleanupTopicSnapshot(scope, root, target.TopicID, "", 0, 0, order, pinned, groups)})
98 }
99 add("global", "", file.GlobalTopics, file.GlobalPinnedTopics, file.GlobalGroups)
100 for _, project := range file.Projects {
101 add("project", project.Root, project.Topics, project.PinnedTopics, project.Groups)
102 }
103 if len(candidates) != 1 {
104 return legacycleanup.Candidate{}, workspacestate.ErrMutationConflict
105 }
106 // A topic ID must have one owner even when a caller supplied a workspace.
107 owners := 0
108 if slices.Contains(file.GlobalTopics, target.TopicID) {
109 owners++
110 }
111 for _, p := range file.Projects {
112 if slices.Contains(p.Topics, target.TopicID) {
113 owners++
114 }
115 }
116 if owners != 1 {
117 return legacycleanup.Candidate{}, workspacestate.ErrMutationConflict
118 }
119 return candidates[0], nil
120 }
121
122 func topicRemovalToken(item legacycleanup.Candidate, associations any) string {
123 body, _ := json.Marshal([]any{item.WorkspaceID, item.TopicID, item.Title, item.Topic, associations})
124 sum := sha256.Sum256(body)
125 return hex.EncodeToString(sum[:])
126 }
127
128 func (a *App) inspectTopicRemovalSnapshot(state workspacestate.State, file desktopProjectFile, snapshot topicstate.Snapshot, item legacycleanup.Candidate) (TopicRemovalInspection, legacycleanup.Candidate, error) {
129 out := TopicRemovalInspection{Target: TopicRemovalTarget{WorkspaceID: item.WorkspaceID, TopicID: item.TopicID}}
130 record, exists := snapshot.Records[item.TopicID]
131 if !exists || record.RowRevision == 0 {
132 return out, item, workspacestate.ErrMutationConflict
133 }
134 item.Title = record.Title
135 item.Topic.TitleSource, item.Topic.CreatedAt, item.Topic.RowRevision = record.TitleSource, record.CreatedAtMS, record.RowRevision
136 out.Disposition = "archive_placeholder"
137 if isDefaultTopicTitle(record.Title) && record.TitleSource == topicTitleSourceAuto && !item.Topic.Pinned && item.Topic.GroupID == "" {
138 out.Disposition = "discard_placeholder"
139 }
140 associations := map[string]any{}
141 for id, presentation := range state.Presentation {
142 if presentation.TopicID == item.TopicID && state.SessionStates[id].Lifecycle != workspacestate.Deleted {
143 out.Disposition = "archive_sessions"
144 associations[id] = state.SessionStates[id]
145 }
146 }
147 targets, err := a.topicTrashTargets(item.TopicID)
148 if err != nil {
149 return out, item, err
150 }
151 for _, target := range targets {
152 associations[target.sessionPath] = target.key
153 out.Disposition = "archive_sessions"
154 }
155 out.Token = topicRemovalToken(item, associations)
156 if a.topicHasActiveRuntimeWork(item.TopicID) {
157 out.Reason = "busy"
158 return out, item, nil
159 }
160 if out.Disposition != "archive_sessions" {
161 if topicRemovalHasPendingCreate(state, item.TopicID) {
162 out.Reason = "busy"
163 return out, item, nil
164 }
165 if classification, reason := classifyLegacyCleanupTopicSources(item, a.knownSessionDirs()); classification != "empty" {
166 out.Reason = reason
167 return out, item, nil
168 }
169 a.mu.RLock()
170 for _, tab := range a.runtimeTabsLocked() {
171 if tab != nil && tab.TopicID == item.TopicID && (tab.Ctrl != nil || tab.SessionID != "" || tab.StartupErr == "") {
172 out.Reason = "busy"
173 break
174 }
175 }
176 a.mu.RUnlock()
177 if out.Reason != "" {
178 return out, item, nil
179 }
180 }
181 out.Allowed = true
182 return out, item, nil
183 }
184
185 func (a *App) InspectTopicRemoval(target TopicRemovalTarget) (TopicRemovalInspection, error) {
186 out, _, err := a.inspectTopicRemoval(target)
187 return out, sessionOperationErrorForTarget(err, target.TopicID, "")
188 }
189
190 func (a *App) inspectTopicRemoval(target TopicRemovalTarget) (TopicRemovalInspection, legacycleanup.Candidate, error) {
191 var out TopicRemovalInspection
192 var item legacycleanup.Candidate
193 if strings.TrimSpace(target.TopicID) == "" {
194 return out, item, workspacestate.ErrMutationConflict
195 }
196 state, err := a.workspaceRegistry().Load(a.bootContext())
197 if err != nil {
198 return out, item, err
199 }
200 if canonical, found, err := a.inspectCanonicalTopicRemoval(state, target); found || err != nil {
201 return canonical, item, err
202 }
203 file, err := readTopicRemovalProjects()
204 if err != nil {
205 return out, item, err
206 }
207 item, err = topicRemovalCandidate(state, file, target)
208 if err != nil {
209 // Indexed ambiguity must not be bypassed by a source fallback.
210 for _, project := range file.Projects {
211 if slices.Contains(project.Topics, target.TopicID) {
212 return out, item, err
213 }
214 }
215 if slices.Contains(file.GlobalTopics, target.TopicID) {
216 return out, item, err
217 }
218 out, err = a.inspectLegacyTopicRemoval(state, target)
219 return out, item, err
220 }
221 snapshot, err := desktopTopicState.snapshot(topicTitleRoot(item.Topic.Scope, item.Topic.WorkspaceRoot))
222 if err != nil {
223 return out, item, err
224 }
225 return a.inspectTopicRemovalSnapshot(state, file, snapshot, item)
226 }
227
228 func (a *App) inspectCanonicalTopicRemoval(state workspacestate.State, target TopicRemovalTarget) (TopicRemovalInspection, bool, error) {
229 out := TopicRemovalInspection{Target: target, Disposition: "archive_sessions", Allowed: true}
230 token, owner, err := workspacestate.TopicSessionRemovalToken(state, target.WorkspaceID, target.TopicID)
231 if err != nil {
232 return out, false, err
233 }
234 if token == "" {
235 return out, false, nil
236 }
237 out.Target.WorkspaceID = owner
238 out.Token = token
239 if a.topicHasActiveRuntimeWork(target.TopicID) {
240 out.Allowed, out.Reason = false, "busy"
241 }
242 return out, true, nil
243 }
244
245 func (a *App) removeCompatiblePlaceholderAdmissionHeld(topicID string) error {
246 inspection, _, err := a.inspectTopicRemoval(TopicRemovalTarget{TopicID: topicID})
247 if err != nil {
248 return err
249 }
250 _, err = a.removePlaceholderAdmissionHeld(TopicRemovalRequest{OperationID: "topic-" + newTabID(), Target: inspection.Target, ExpectedToken: inspection.Token})
251 return err
252 }
253
254 func (a *App) RemoveTopic(req TopicRemovalRequest) (TopicRemovalResult, error) {
255 out := TopicRemovalResult{}
256 if req.OperationID == "" || len(req.OperationID) > 200 || req.ExpectedToken == "" || req.Target.WorkspaceID == "" {
257 return out, errors.New("invalid topic removal request")
258 }
259 release, ok := a.tryLockRuntimeMutationBounded("remove topic")
260 if !ok {
261 out.ErrorCode, out.Retryable = "busy", true
262 return out, nil
263 }
264 defer release()
265 state, err := a.workspaceRegistry().Load(a.bootContext())
266 if err == nil {
267 if old, exists := state.TopicRemovals[req.OperationID]; exists && old.WorkspaceID == req.Target.WorkspaceID && old.TopicID == req.Target.TopicID && old.Token == req.ExpectedToken && (old.Phase == "committed" || old.Phase == "restored" || old.Phase == "purged") {
268 return topicRemovalResult(old), nil
269 }
270 }
271 if err == nil {
272 // Placeholder requests have durable retry state; inspect occurs again in
273 // the writer transaction. Session-backed topics retain canonical archive.
274 if _, pending := state.TopicRemovals[req.OperationID]; !pending {
275 var inspection TopicRemovalInspection
276 inspection, _, err = a.inspectTopicRemoval(req.Target)
277 if err == nil && (!inspection.Allowed || inspection.Token != req.ExpectedToken) {
278 err = workspacestate.ErrMutationConflict
279 }
280 if err == nil && inspection.Disposition == "archive_sessions" {
281 out, err = a.removeTopicSessionsAdmissionHeld(req)
282 if err == nil {
283 return out, nil
284 }
285 }
286 } else if state.TopicRemovals[req.OperationID].Disposition == "archive_sessions" {
287 out, err = a.removeTopicSessionsAdmissionHeld(req)
288 if err == nil {
289 return out, nil
290 }
291 }
292 }
293 if err == nil {
294 out, err = a.removePlaceholderAdmissionHeld(req)
295 }
296 if err != nil {
297 out.ErrorCode, out.Retryable = "operation_failed", true
298 out.ErrorMessage = err.Error()
299 if errors.Is(err, workspacestate.ErrMutationConflict) {
300 out.ErrorCode, out.Retryable = "state_conflict", false
301 }
302 if errors.Is(err, errTopicArchiveBusy) || errors.Is(err, errTopicHasActiveWork) {
303 out.ErrorCode = "busy"
304 }
305 a.emitProjectTreeChanged()
306 }
307 return out, nil
308 }
309
310 func topicRemovalResult(item workspacestate.TopicRemoval) TopicRemovalResult {
311 out := TopicRemovalResult{Committed: true, Disposition: item.Disposition}
312 if item.Disposition == "archive_placeholder" {
313 out.RecoveryEntryID = "topic-removal:" + item.ID
314 }
315 return out
316 }
317
318 func (a *App) removeTopicSessionsAdmissionHeld(req TopicRemovalRequest) (TopicRemovalResult, error) {
319 // Match AI/manual rename: removal -> title -> index. Cleanup must run
320 // after title/index unlock, while removal and runtime admission remain held.
321 a.sessionRemovalMu.Lock()
322 defer a.sessionRemovalMu.Unlock()
323 cleanup := archivedRuntimeCleanup{app: a}
324 defer cleanup.finish()
325 // Serialize title/organization changes from this process across admission.
326 a.topicTitleMutationMu.Lock()
327 defer a.topicTitleMutationMu.Unlock()
328 topicIndexMu.Lock()
329 defer topicIndexMu.Unlock()
330 state, err := a.workspaceRegistry().Load(a.bootContext())
331 if err != nil {
332 return TopicRemovalResult{}, err
333 }
334 archiveID := "topic-sessions-" + req.OperationID
335 finished := state.PendingOperations[archiveID].Phase == "committed"
336 if !finished {
337 inspection, _, err := a.inspectTopicRemoval(req.Target)
338 if err != nil {
339 return TopicRemovalResult{}, err
340 }
341 if !inspection.Allowed || inspection.Disposition != "archive_sessions" || inspection.Token != req.ExpectedToken {
342 return TopicRemovalResult{}, workspacestate.ErrMutationConflict
343 }
344 }
345 err = a.workspaceRegistry().TransitionTopicRemoval(a.bootContext(), req.OperationID, func(state workspacestate.State, item *workspacestate.TopicRemoval, _ func() error) error {
346 if item.ID != "" {
347 if item.WorkspaceID != req.Target.WorkspaceID || item.TopicID != req.Target.TopicID || item.Token != req.ExpectedToken || item.Disposition != "archive_sessions" {
348 return workspacestate.ErrMutationConflict
349 }
350 return nil
351 }
352 token, _, err := workspacestate.TopicSessionRemovalToken(state, req.Target.WorkspaceID, req.Target.TopicID)
353 if err != nil {
354 return err
355 }
356 if token != "" && token != req.ExpectedToken {
357 return workspacestate.ErrMutationConflict
358 }
359 *item = workspacestate.TopicRemoval{ID: req.OperationID, WorkspaceID: req.Target.WorkspaceID, TopicID: req.Target.TopicID, Token: req.ExpectedToken, SessionToken: token, Disposition: "archive_sessions", Phase: "prepared"}
360 return nil
361 })
362 if err != nil {
363 return TopicRemovalResult{}, err
364 }
365 if !finished {
366 a.lifecycleCheckpoint("topic-sessions-before-archive")
367 if err := a.archiveCompatibleTopicWithCleanupAdmissionHeld(req.Target.TopicID, archiveID, &cleanup); err != nil {
368 return TopicRemovalResult{}, err
369 }
370 }
371 err = a.workspaceRegistry().TransitionTopicRemoval(a.bootContext(), req.OperationID, func(_ workspacestate.State, item *workspacestate.TopicRemoval, _ func() error) error {
372 item.Phase = "committed"
373 return nil
374 })
375 return TopicRemovalResult{Committed: err == nil, Disposition: "archive_sessions"}, err
376 }
377
378 func (a *App) removePlaceholderAdmissionHeld(req TopicRemovalRequest) (TopicRemovalResult, error) {
379 var result workspacestate.TopicRemoval
380 if !a.sessionRemovalMu.TryLock() {
381 return TopicRemovalResult{}, errTopicArchiveBusy
382 }
383 defer a.sessionRemovalMu.Unlock()
384 var removed []removedSessionRuntime
385 defer func() { a.finalizeRemovedTopicRuntimes(removed) }()
386 a.topicTitleMutationMu.Lock()
387 defer a.topicTitleMutationMu.Unlock()
388 topicIndexMu.Lock()
389 defer topicIndexMu.Unlock()
390 err := a.workspaceRegistry().TransitionTopicRemoval(a.bootContext(), req.OperationID, func(state workspacestate.State, operation *workspacestate.TopicRemoval, checkpoint func() error) error {
391 if operation.ID != "" {
392 if operation.Token != req.ExpectedToken || operation.TopicID != req.Target.TopicID || operation.WorkspaceID != req.Target.WorkspaceID {
393 return workspacestate.ErrMutationConflict
394 }
395 if operation.Phase == "committed" || operation.Phase == "restored" || operation.Phase == "purged" {
396 result = *operation
397 return nil
398 }
399 if operation.Phase != "prepared" {
400 return workspacestate.ErrMutationConflict
401 }
402 }
403 desktopProjectsFileMu.Lock()
404 defer desktopProjectsFileMu.Unlock()
405 release, err := acquireDesktopProjectsFileLock()
406 if err != nil {
407 return err
408 }
409 defer release()
410 file, err := readTopicRemovalProjects()
411 if err != nil {
412 return err
413 }
414 var item legacycleanup.Candidate
415 if operation.ID != "" {
416 if err := json.Unmarshal(operation.Snapshot, &item); err != nil {
417 return err
418 }
419 } else {
420 item, err = topicRemovalCandidate(state, file, req.Target)
421 if err != nil {
422 return err
423 }
424 }
425 if item.Topic == nil {
426 return workspacestate.ErrMutationConflict
427 }
428 if topicRemovalHasOtherOwner(file, item) {
429 return workspacestate.ErrMutationConflict
430 }
431 err = desktopTopicState.withExclusiveScope(topicTitleRoot(item.Topic.Scope, item.Topic.WorkspaceRoot), func(ctx context.Context, store *topicstate.Store) error {
432 return a.removePlaceholderMetadataLocked(ctx, store, state, file, item, req, operation, checkpoint)
433 })
434 result = *operation
435 return err
436 })
437 if err != nil {
438 return TopicRemovalResult{}, err
439 }
440 captured := a.captureTopicRuntimeBindings(req.Target.TopicID)
441 _, unchanged := a.removeTopicRuntimeBindingsIfUnchanged(req.Target.TopicID, captured)
442 if unchanged {
443 removed = captured
444 }
445 a.emitProjectTreeChanged()
446 return topicRemovalResult(result), nil
447 }
448
449 func (a *App) checkRemovedTopicSources(state workspacestate.State, item legacycleanup.Candidate) error {
450 if topicRemovalHasPendingCreate(state, item.TopicID) {
451 return errTopicHasActiveWork
452 }
453 for id, presentation := range state.Presentation {
454 if presentation.TopicID == item.TopicID && state.SessionStates[id].Lifecycle != workspacestate.Deleted {
455 return workspacestate.ErrMutationConflict
456 }
457 }
458 if classification, _ := classifyLegacyCleanupTopicSources(item, a.knownSessionDirs()); classification != "empty" {
459 return workspacestate.ErrMutationConflict
460 }
461 if a.topicHasActiveRuntimeWork(item.TopicID) {
462 return errTopicHasActiveWork
463 }
464 a.mu.RLock()
465 defer a.mu.RUnlock()
466 for _, tab := range a.runtimeTabsLocked() {
467 if tab != nil && tab.TopicID == item.TopicID && (tab.Ctrl != nil || tab.SessionID != "" || tab.StartupErr == "") {
468 return errTopicHasActiveWork
469 }
470 }
471 return nil
472 }
473
474 func topicRemovalHasPendingCreate(state workspacestate.State, topicID string) bool {
475 for _, create := range state.PendingCreates {
476 if create.Presentation != nil && create.Presentation.TopicID == topicID {
477 return true
478 }
479 }
480 for _, operation := range state.PendingOperations {
481 if operation.Phase != "committed" && operation.Presentation != nil && operation.Presentation.TopicID == topicID {
482 return true
483 }
484 }
485 return false
486 }
487
488 func topicRemovalError(result TopicRemovalResult) error {
489 if result.Committed {
490 return nil
491 }
492 return fmt.Errorf("topic removal incomplete: %s", result.ErrorCode)
493 }
494
495 // Requires registry, projects and topic-store writer ownership.
496 func (a *App) removePlaceholderMetadataLocked(ctx context.Context, store *topicstate.Store, state workspacestate.State, file desktopProjectFile, item legacycleanup.Candidate, req TopicRemovalRequest, operation *workspacestate.TopicRemoval, checkpoint func() error) error {
497 snapshot, err := store.Snapshot(ctx)
498 if err != nil {
499 return err
500 }
501 indexed := topicIndexedInProjectsSnapshot(file, item.Topic.Scope, item.Topic.WorkspaceRoot, item.TopicID)
502 if indexed {
503 current, err := topicRemovalCandidate(state, file, req.Target)
504 if err != nil {
505 return err
506 }
507 inspection, captured, err := a.inspectTopicRemovalSnapshot(state, file, snapshot, current)
508 if err != nil {
509 return err
510 }
511 if !inspection.Allowed || inspection.Token != req.ExpectedToken || inspection.Disposition == "archive_sessions" {
512 return workspacestate.ErrMutationConflict
513 }
514 item = captured
515 if operation.ID == "" {
516 body, err := json.Marshal(item)
517 if err != nil {
518 return err
519 }
520 *operation = workspacestate.TopicRemoval{ID: req.OperationID, WorkspaceID: item.WorkspaceID, TopicID: item.TopicID, Token: req.ExpectedToken, Disposition: inspection.Disposition, Phase: "prepared", ArchivedAt: time.Now().UnixMilli(), Snapshot: body}
521 operation.Metadata, err = json.Marshal(snapshot.Records[item.TopicID])
522 if err != nil {
523 return err
524 }
525 if err := checkpoint(); err != nil {
526 return err
527 }
528 }
529 } else {
530 if operation.ID == "" || !slices.Contains(file.DeletedTopics, item.TopicID) {
531 return workspacestate.ErrMutationConflict
532 }
533 if r, exists := snapshot.Records[item.TopicID]; exists && r.RowRevision != item.Topic.RowRevision {
534 return workspacestate.ErrMutationConflict
535 }
536 if err := a.checkRemovedTopicSources(state, item); err != nil {
537 return err
538 }
539 }
540 a.lifecycleCheckpoint("topic-removal-before-index")
541 if err := removeTopicFromProjectsFileCrossProcessLocked(item.TopicID); err != nil {
542 return err
543 }
544 a.lifecycleCheckpoint("topic-removal-before-metadata")
545 if _, err := store.Delete(ctx, item.TopicID); err != nil {
546 return err
547 }
548 a.lifecycleCheckpoint("topic-removal-before-commit")
549 operation.Phase = "committed"
550 if operation.Disposition == "discard_placeholder" {
551 operation.Snapshot = nil
552 operation.Metadata = nil
553 }
554 return nil
555 }
556
556 lines GO