返回 DeepSeek-Reasonix
session_lifecycle_api.go
根目录 / desktop / session_lifecycle_api.go
1 package main
2
3 import (
4 "crypto/sha256"
5 "encoding/hex"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "sort"
10 "strings"
11 "sync"
12
13 "reasonix/desktop/internal/workspacestate"
14 "reasonix/internal/session"
15 )
16
17 type SessionLifecycleTarget struct {
18 WorkspaceID string `json:"workspaceId,omitempty"`
19 Ref *session.SessionRef `json:"ref,omitempty"`
20 RecoveryEntryID string `json:"recoveryEntryId,omitempty"`
21 }
22 type SessionLifecycleRequest struct {
23 OperationID string `json:"operationId"`
24 Action string `json:"action"`
25 Targets []SessionLifecycleTarget `json:"targets"`
26 ExpectedGeneration uint64 `json:"expectedGeneration"`
27 }
28 type SessionLifecycleItem struct {
29 Target SessionLifecycleTarget `json:"target"`
30 Ref *session.SessionRef `json:"ref,omitempty"`
31 WorkspaceID string `json:"workspaceId"`
32 Committed bool `json:"committed"`
33 ErrorCode string `json:"errorCode,omitempty"`
34 Retryable bool `json:"retryable"`
35 }
36 type SessionLifecycleResult struct {
37 OperationID string `json:"operationId"`
38 Generation uint64 `json:"generation"`
39 Committed bool `json:"committed"`
40 Items []SessionLifecycleItem `json:"items"`
41 }
42
43 // Serializes duplicate RPC deliveries; persisted receipts handle restart.
44 var desktopLifecycleCommands sync.Mutex
45
46 func (a *App) ApplySessionLifecycle(req SessionLifecycleRequest) (SessionLifecycleResult, error) {
47 desktopLifecycleCommands.Lock()
48 defer desktopLifecycleCommands.Unlock()
49 out := SessionLifecycleResult{OperationID: req.OperationID, Items: []SessionLifecycleItem{}}
50 if err := validateLifecycleRequest(req); err != nil {
51 return out, err
52 }
53 body, err := json.Marshal(req)
54 if err != nil {
55 return out, err
56 }
57 sum := sha256.Sum256(body)
58 fingerprint := hex.EncodeToString(sum[:])
59 ctx := a.bootContext()
60 store := a.workspaceRegistry()
61 state, err := store.Load(ctx)
62 if err != nil {
63 return out, err
64 }
65 key := "command-" + req.OperationID
66 if old, ok := state.PendingOperations[key]; ok {
67 if old.Kind != "command" || old.RequestFingerprint != fingerprint {
68 return out, workspacestate.ErrMutationConflict
69 }
70 if len(old.Result) > 0 {
71 if err := json.Unmarshal(old.Result, &out); err != nil {
72 return out, err
73 }
74 out.Generation = old.ResultGeneration
75 }
76 if old.Phase == "committed" {
77 return out, nil
78 }
79 } else {
80 begin := store.BeginCommand
81 if req.Action == "purge" {
82 begin = store.BeginPurgeCommand
83 }
84 if err := begin(ctx, key, fingerprint, body, req.ExpectedGeneration); err != nil {
85 return out, err
86 }
87 }
88 previous := out.Items
89 out.Items = []SessionLifecycleItem{}
90 archiveErr := a.archiveLifecycleCommand(req, key)
91 for i, target := range req.Targets {
92 if i < len(previous) && (previous[i].Committed || !previous[i].Retryable) {
93 out.Items = append(out.Items, previous[i])
94 continue
95 }
96 item, err := a.applyLifecycleTarget(req, key, i, target, state, archiveErr)
97 if err != nil {
98 return out, err
99 }
100 out.Items = append(out.Items, item)
101 }
102 state, err = store.Load(ctx)
103 if err != nil {
104 return out, err
105 }
106 out.Generation = state.Generation + 1
107 out.Committed = true
108 final := true
109 for _, item := range out.Items {
110 out.Committed = out.Committed && item.Committed
111 if !item.Committed && item.Retryable {
112 final = false
113 }
114 }
115 body, err = json.Marshal(out)
116 if err != nil {
117 return out, err
118 }
119 a.lifecycleCheckpoint("before-command-result")
120 if err := store.SaveCommandResult(ctx, key, body, final); err != nil {
121 return out, err
122 }
123 a.lifecycleCheckpoint("after-command-result")
124 state, err = store.Load(ctx)
125 if err != nil {
126 return out, err
127 }
128 out.Generation = state.PendingOperations[key].ResultGeneration
129 a.emitProjectTreeChanged()
130 return out, nil
131 }
132
133 type TrashEntry struct {
134 ID string `json:"id"`
135 Ref *session.SessionRef `json:"ref,omitempty"`
136 RecoveryEntryID string `json:"recoveryEntryId,omitempty"`
137 Title string `json:"title"`
138 WorkspaceID string `json:"workspaceId"`
139 WorkspaceTitle string `json:"workspaceTitle"`
140 ArchivedAt int64 `json:"archivedAt"`
141 Health string `json:"health"`
142 OperationPhase string `json:"operationPhase,omitempty"`
143 CanPreview bool `json:"canPreview"`
144 CanRestore bool `json:"canRestore"`
145 CanPurge bool `json:"canPurge"`
146 CleanupBatchID string `json:"cleanupBatchId,omitempty"`
147 CleanupKind string `json:"cleanupKind,omitempty"`
148 }
149 type TrashEntryPage struct {
150 Items []TrashEntry `json:"items"`
151 Generation uint64 `json:"generation"`
152 NextCursor string `json:"nextCursor,omitempty"`
153 }
154
155 func (a *App) ListTrashEntries(query, cursor string, limit int) (TrashEntryPage, error) {
156 out := TrashEntryPage{Items: []TrashEntry{}}
157 state, err := a.workspaceRegistry().Load(a.bootContext())
158 if err != nil {
159 return out, err
160 }
161 out.Generation = state.Generation
162 start, err := decodeWorkspaceSessionCursor(cursor, state.Generation)
163 if err != nil {
164 return out, err
165 }
166 rows := []TrashEntry{}
167 service := a.desktopSessionService("")
168 cleanupState, _ := a.legacyCleanup.Load(a.bootContext())
169 cleanupBySession := legacyCleanupArchivedSessions(cleanupState)
170 for id, status := range state.SessionStates {
171 op := state.PendingOperations["purge-"+id]
172 purgeState := workspacestate.ClassifyPurge(state, id)
173 pending := purgeState == workspacestate.PurgeTombstoned || purgeState == workspacestate.PurgeContentRemoved || purgeState == workspacestate.PurgeInvalid
174 if status.Lifecycle != workspacestate.Archived && !pending {
175 continue
176 }
177 ref := session.SessionRef{HostID: localDesktopHostID, SessionID: id}
178 row := TrashEntry{ID: id, Ref: &ref, Title: state.Presentation[id].Title, ArchivedAt: status.ArchivedAt, CanPurge: purgeState != workspacestate.PurgeInvalid, Health: "ready"}
179 decorateLegacyCleanupTrashEntry(&row, cleanupState.BatchID, cleanupBySession[id])
180 for wid, w := range state.Workspaces {
181 if containsDesktopString(w.SessionIDs, id) {
182 row.WorkspaceID = wid
183 row.WorkspaceTitle = w.Title
184 break
185 }
186 }
187 if info, e := service.Query().Stat(a.bootContext(), ref); e == nil {
188 if info.Title != "" {
189 row.Title = info.Title
190 }
191 row.CanPreview = !pending
192 row.CanRestore = !pending
193 } else {
194 row.Health = "unavailable"
195 }
196 if pending {
197 row.OperationPhase = op.Phase
198 row.Health = "purge_pending"
199 }
200 if row.Title == "" {
201 row.Title = id
202 }
203 if strings.Contains(strings.ToLower(row.Title+"\n"+row.WorkspaceTitle), strings.ToLower(query)) {
204 rows = append(rows, row)
205 }
206 }
207 rows = append(rows, legacyCleanupTopicTrashEntries(cleanupState, state, query)...)
208 rows = append(rows, topicRemovalTrashEntries(state, query)...)
209 sort.Slice(rows, func(i, j int) bool {
210 if rows[i].ArchivedAt != rows[j].ArchivedAt {
211 return rows[i].ArchivedAt > rows[j].ArchivedAt
212 }
213 return rows[i].ID < rows[j].ID
214 })
215 if limit <= 0 {
216 limit = 50
217 }
218 if limit > 200 {
219 limit = 200
220 }
221 if start > len(rows) {
222 start = len(rows)
223 }
224 end := min(start+limit, len(rows))
225 out.Items = append(out.Items, rows[start:end]...)
226 if end < len(rows) {
227 out.NextCursor = fmt.Sprintf("%d:%d", state.Generation, end)
228 }
229 return out, nil
230 }
231
232 func validateLifecycleRequest(req SessionLifecycleRequest) error {
233 if strings.TrimSpace(req.OperationID) == "" || len(req.OperationID) > 200 || len(req.Targets) == 0 || len(req.Targets) > 1000 {
234 return errors.New("invalid lifecycle request")
235 }
236 if req.Action != "archive" && req.Action != "restore" && req.Action != "purge" {
237 return errors.New("invalid lifecycle action")
238 }
239 seen := map[string]bool{}
240 for _, target := range req.Targets {
241 if (target.Ref == nil) == (target.RecoveryEntryID == "") {
242 return errors.New("exactly one session identity is required")
243 }
244 if target.Ref != nil {
245 if err := validateLocalSessionRef(*target.Ref); err != nil {
246 return err
247 }
248 } else if req.Action != "restore" && !(req.Action == "purge" && (strings.HasPrefix(target.RecoveryEntryID, "legacy-cleanup:") || strings.HasPrefix(target.RecoveryEntryID, "topic-removal:"))) {
249 return errors.New("historical recovery entries can only be restored")
250 }
251 body, _ := json.Marshal(target)
252 if seen[string(body)] {
253 return errors.New("duplicate lifecycle target")
254 }
255 seen[string(body)] = true
256 }
257 return nil
258 }
259
260 func (a *App) archiveLifecycleCommand(req SessionLifecycleRequest, key string) error {
261 store := a.workspaceRegistry()
262 if req.Action == "archive" {
263 child := key + "-archive"
264 state, err := store.Load(a.bootContext())
265 if err != nil {
266 return err
267 }
268 if state.PendingOperations[child].Phase != "committed" {
269 for _, target := range req.Targets {
270 if state.SessionStates[target.Ref.SessionID].Generation > req.ExpectedGeneration {
271 return workspacestate.ErrMutationConflict
272 }
273 }
274 refs := []session.SessionRef{}
275 for _, target := range req.Targets {
276 refs = append(refs, *target.Ref)
277 }
278 release := a.lockRuntimeMutation("archive lifecycle command")
279 e := a.archiveSessionRefsWithOperation(refs, child)
280 release()
281 if e != nil {
282 return e
283 }
284 }
285 }
286 return nil
287 }
288
289 func (a *App) applyLifecycleTarget(req SessionLifecycleRequest, key string, index int, target SessionLifecycleTarget, state workspacestate.State, archiveErr error) (SessionLifecycleItem, error) {
290 ctx, store := a.bootContext(), a.workspaceRegistry()
291 item := SessionLifecycleItem{Target: target, Ref: target.Ref}
292 var opErr error
293 child := fmt.Sprintf("%s-%d", key, index)
294 latest, loadErr := store.Load(ctx)
295 if loadErr != nil {
296 return item, loadErr
297 }
298 if target.Ref != nil && req.Action != "archive" && latest.PendingOperations[child].Phase != "committed" && latest.SessionStates[target.Ref.SessionID].Generation > req.ExpectedGeneration {
299 purgeState := workspacestate.ClassifyPurge(latest, target.Ref.SessionID)
300 resumingPurge := req.Action == "purge" && (purgeState == workspacestate.PurgeTombstoned || purgeState == workspacestate.PurgeContentRemoved || purgeState == workspacestate.PurgeCommitted)
301 if !resumingPurge {
302 item.ErrorCode = "state_conflict"
303 return item, nil
304 }
305 }
306 if target.Ref != nil {
307 for id, w := range state.Workspaces {
308 if containsDesktopString(w.SessionIDs, target.Ref.SessionID) {
309 item.WorkspaceID = id
310 break
311 }
312 }
313 }
314 switch req.Action {
315 case "archive":
316 opErr = archiveErr
317 case "purge":
318 if target.Ref == nil && strings.HasPrefix(target.RecoveryEntryID, "topic-removal:") {
319 opErr = a.purgeRemovedTopic(strings.TrimPrefix(target.RecoveryEntryID, "topic-removal:"), target.WorkspaceID)
320 break
321 }
322 if target.Ref == nil && strings.HasPrefix(target.RecoveryEntryID, "legacy-cleanup:") {
323 opErr = a.purgeLegacyCleanupTopic(strings.TrimPrefix(target.RecoveryEntryID, "legacy-cleanup:"), target.WorkspaceID)
324 } else {
325 release := a.lockRuntimeMutation("purge lifecycle command")
326 opErr = a.purgeCanonicalSession(ctx, *target.Ref, req.ExpectedGeneration)
327 release()
328 }
329 case "restore":
330 if target.Ref == nil && strings.HasPrefix(target.RecoveryEntryID, "topic-removal:") {
331 opErr = a.restoreRemovedTopic(strings.TrimPrefix(target.RecoveryEntryID, "topic-removal:"), target.WorkspaceID)
332 item.WorkspaceID = target.WorkspaceID
333 break
334 }
335 if target.Ref == nil && strings.HasPrefix(target.RecoveryEntryID, "legacy-cleanup:") {
336 opErr = a.restoreLegacyCleanupTopic(strings.TrimPrefix(target.RecoveryEntryID, "legacy-cleanup:"), target.WorkspaceID)
337 item.WorkspaceID = target.WorkspaceID
338 break
339 }
340 var restored SessionRestoreResult
341 if target.Ref == nil {
342 restored, opErr = a.restoreRecoveryEntryInWorkspace(target.RecoveryEntryID, child, target.WorkspaceID)
343 } else {
344 release := a.lockRuntimeMutation("restore lifecycle command")
345 saved, loadErr := store.Load(ctx)
346 if loadErr != nil {
347 opErr = loadErr
348 } else if done := saved.PendingOperations[child]; done.Phase == "committed" {
349 restored = SessionRestoreResult{Session: *target.Ref, WorkspaceID: done.WorkspaceID, Generation: done.ResultGeneration}
350 } else {
351 restored, opErr = a.restoreCanonicalSession(ctx, *target.Ref, child)
352 }
353 release()
354 }
355 if opErr == nil {
356 item.Ref = &restored.Session
357 item.WorkspaceID = restored.WorkspaceID
358 }
359 }
360 item.Committed = opErr == nil
361 if opErr != nil {
362 item.ErrorCode = "operation_failed"
363 item.Retryable = true
364 if errors.Is(opErr, workspacestate.ErrMutationConflict) {
365 item.ErrorCode = "state_conflict"
366 item.Retryable = false
367 }
368 }
369 return item, nil
370 }
371
371 lines GO