返回 DeepSeek-Reasonix
lifecycle.go
根目录 / desktop / internal / workspacestate / lifecycle.go
1 package workspacestate
2
3 import (
4 "bytes"
5 "context"
6 "crypto/sha256"
7 "encoding/hex"
8 "encoding/json"
9 "errors"
10 "fmt"
11 "os"
12 "path/filepath"
13 "reflect"
14 "slices"
15 "strings"
16 "time"
17
18 "reasonix/internal/fileutil"
19 )
20
21 const (
22 Active = "active"
23 Archived = "archived"
24 Deleted = "deleted"
25 )
26
27 type SessionState struct {
28 Lifecycle string `json:"lifecycle"`
29 Generation uint64 `json:"generation"`
30 ArchivedAt int64 `json:"archivedAt,omitempty"`
31 extra map[string]json.RawMessage
32 }
33 type SourceMapping struct {
34 RetainedArtifacts []string `json:"retainedArtifacts,omitempty"`
35 SourceKey string `json:"sourceKey"`
36 Path string `json:"path"`
37 HeadID string `json:"headId,omitempty"`
38 Format string `json:"format"`
39 Fingerprint string `json:"fingerprint"`
40 SessionID string `json:"sessionId"`
41 WorkspaceID string `json:"workspaceId"`
42 extra map[string]json.RawMessage
43 }
44 type Presentation struct {
45 TopicID string `json:"topicId,omitempty"`
46 Title string `json:"title,omitempty"`
47 Pinned bool `json:"pinned,omitempty"`
48 SortOrder int `json:"sortOrder"`
49 extra map[string]json.RawMessage
50 }
51 type RecoveryEntry struct {
52 ID string `json:"id"`
53 SourceKey string `json:"sourceKey"`
54 Path string `json:"path,omitempty"`
55 HeadID string `json:"headId,omitempty"`
56 SessionID string `json:"sessionId,omitempty"`
57 WorkspaceID string `json:"workspaceId,omitempty"`
58 Scope string `json:"scope,omitempty"`
59 WorkspaceRoot string `json:"workspaceRoot,omitempty"`
60 Format string `json:"format"`
61 Reason string `json:"reason"`
62 Status string `json:"status"`
63 Fingerprint string `json:"fingerprint,omitempty"`
64 extra map[string]json.RawMessage
65 }
66 type Operation struct {
67 RequestFingerprint string `json:"requestFingerprint,omitempty"`
68 Request json.RawMessage `json:"request,omitempty"`
69 Result json.RawMessage `json:"result,omitempty"`
70 ID string `json:"id"`
71 Kind string `json:"kind"`
72 Phase string `json:"phase"`
73 SessionIDs []string `json:"sessionIds"`
74 WorkspaceID string `json:"workspaceId,omitempty"`
75 Lifecycle string `json:"lifecycle"`
76 ExpectedGeneration uint64 `json:"expectedGeneration"`
77 ResultGeneration uint64 `json:"resultGeneration,omitempty"`
78 RecoveryEntryID string `json:"recoveryEntryId,omitempty"`
79 Mapping *SourceMapping `json:"mapping,omitempty"`
80 Presentation *Presentation `json:"presentation,omitempty"`
81 Dependencies []string `json:"dependencies,omitempty"`
82 extra map[string]json.RawMessage
83 }
84
85 type PurgeState uint8
86
87 const (
88 PurgeAbsent PurgeState = iota
89 PurgePrepared
90 PurgePreparedStale
91 PurgeTombstoned
92 PurgeContentRemoved
93 PurgeCommitted
94 PurgeInvalid
95 )
96
97 func validLifecycle(value string) bool {
98 return value == Active || value == Archived || value == Deleted
99 }
100 func setLifecycle(state *State, id, lifecycle string) {
101 current := state.SessionStates[id]
102 if lifecycle == Archived && current.Lifecycle != Archived {
103 current.ArchivedAt = time.Now().UnixMilli()
104 }
105 current.Lifecycle, current.Generation = lifecycle, state.Generation+1
106 state.SessionStates[id] = current
107 }
108 func validateLifecycleState(state State) error {
109 for id, item := range state.SessionStates {
110 if strings.TrimSpace(id) == "" || !validLifecycle(item.Lifecycle) {
111 return fmt.Errorf("%w: session lifecycle", ErrUnsupportedVersion)
112 }
113 }
114 for key, item := range state.SourceMappings {
115 if key == "" || item.SourceKey != key || item.SessionID == "" || item.Fingerprint == "" {
116 return errors.New("invalid source mapping")
117 }
118 }
119 for key, item := range state.PendingOperations {
120 if key == "" || item.ID != key || !validLifecycle(item.Lifecycle) {
121 return errors.New("invalid session operation")
122 }
123 switch item.Kind {
124 case "import", "archive-import", "archive", "purge", "restore", "command":
125 default:
126 return fmt.Errorf("%w: operation kind", ErrUnsupportedVersion)
127 }
128 switch item.Phase {
129 case "prepared", "content_ready", "committed":
130 case "tombstoned", "content_removed":
131 if item.Kind != "purge" {
132 return ErrUnsupportedVersion
133 }
134 default:
135 return fmt.Errorf("%w: operation phase", ErrUnsupportedVersion)
136 }
137 }
138 for key, item := range state.RecoveryEntries {
139 if key == "" || item.ID != key || item.SourceKey == "" {
140 return errors.New("invalid recovery entry")
141 }
142 switch item.Status {
143 case "pending", "failed", "restored":
144 default:
145 return fmt.Errorf("%w: recovery status", ErrUnsupportedVersion)
146 }
147 }
148 return nil
149 }
150
151 func (s *Store) SetLifecycle(ctx context.Context, ids []string, lifecycle string) error {
152 return s.mutate(ctx, func(state *State) error {
153 if !validLifecycle(lifecycle) {
154 return ErrUnsupportedVersion
155 }
156 for _, id := range ids {
157 if state.SessionStates[id].Lifecycle == Deleted {
158 return ErrMutationConflict
159 }
160 if _, ok := sessionOwner(*state, id); !ok {
161 return ErrSessionNotFound
162 }
163 if err := validateLifecycleSupersedesPreparedPurge(*state, id); err != nil {
164 return err
165 }
166 }
167 for _, id := range ids {
168 cancelPreparedPurge(state, id)
169 setLifecycle(state, id, lifecycle)
170 if lifecycle == Active {
171 owner, _ := sessionOwner(*state, id)
172 workspace := state.Workspaces[owner]
173 workspace.Visible = true
174 state.Workspaces[owner] = workspace
175 }
176 }
177 return nil
178 })
179 }
180
181 func (s *Store) BeginOperation(ctx context.Context, op Operation) error {
182 return s.mutate(ctx, func(state *State) error {
183 for _, id := range op.SessionIDs {
184 if state.SessionStates[id].Lifecycle == Deleted {
185 return ErrMutationConflict
186 }
187 }
188 if op.ID == "" || !validLifecycle(op.Lifecycle) {
189 return ErrMutationConflict
190 }
191 if old, ok := state.PendingOperations[op.ID]; ok {
192 if old.Kind != op.Kind || old.RecoveryEntryID != op.RecoveryEntryID || old.Lifecycle != op.Lifecycle || old.WorkspaceID != op.WorkspaceID || !slices.Equal(old.Dependencies, op.Dependencies) || (len(op.SessionIDs) != 0 && !slices.Equal(old.SessionIDs, op.SessionIDs)) {
193 return ErrMutationConflict
194 }
195 return nil
196 }
197 if op.ExpectedGeneration != 0 && op.ExpectedGeneration != state.Generation {
198 return ErrMutationConflict
199 }
200 if err := validateArchiveImportReservation(state, op); err != nil {
201 return err
202 }
203 // Bind child admission to the original command's observed state,
204 // not to a newer snapshot taken after content/ownership validation.
205 if split := strings.LastIndex(op.ID, "-"); split > 0 {
206 parent := state.PendingOperations[op.ID[:split]]
207 for _, id := range op.SessionIDs {
208 if parent.Kind == "command" && state.SessionStates[id].Generation > parent.ExpectedGeneration {
209 return ErrMutationConflict
210 }
211 }
212 }
213 op.ExpectedGeneration = state.Generation
214 op.Phase = "prepared"
215 if op.SessionIDs == nil {
216 op.SessionIDs = []string{}
217 }
218 state.PendingOperations[op.ID] = op
219 return nil
220 })
221 }
222
223 func (s *Store) ReserveOperationTargets(ctx context.Context, id string, ids []string, sources ...*SourceMapping) error {
224 return s.mutate(ctx, func(state *State) error {
225 op, ok := state.PendingOperations[id]
226 if !ok || len(ids) == 0 {
227 return ErrMutationConflict
228 }
229 if len(op.SessionIDs) != 0 && !slices.Equal(op.SessionIDs, ids) {
230 return ErrMutationConflict
231 }
232 if len(sources) > 0 && sources[0] != nil {
233 mapping := sources[0]
234 if op.Mapping != nil && (op.Mapping.SourceKey != mapping.SourceKey || op.Mapping.Fingerprint != mapping.Fingerprint || op.Mapping.SessionID != mapping.SessionID) {
235 return ErrMutationConflict
236 }
237 if op.Phase == "prepared" {
238 op.Mapping = mapping
239 }
240 }
241 op.SessionIDs = append([]string{}, ids...)
242 state.PendingOperations[id] = op
243 return nil
244 })
245 }
246
247 func (s *Store) PrepareOperationContent(ctx context.Context, id string, ids []string, mapping *SourceMapping, presentation *Presentation) error {
248 return s.mutate(ctx, func(state *State) error {
249 op, ok := state.PendingOperations[id]
250 if !ok {
251 return ErrMutationConflict
252 }
253 if op.Phase != "prepared" {
254 if !slices.Equal(op.SessionIDs, ids) || !reflect.DeepEqual(op.Mapping, mapping) || !reflect.DeepEqual(op.Presentation, presentation) {
255 return ErrMutationConflict
256 }
257 return nil
258 }
259 if len(op.SessionIDs) != 0 && !slices.Equal(op.SessionIDs, ids) {
260 return ErrMutationConflict
261 }
262 op.Phase, op.SessionIDs, op.Mapping, op.Presentation = "content_ready", append([]string{}, ids...), mapping, presentation
263 state.PendingOperations[id] = op
264 return nil
265 })
266 }
267
268 // CommitOperation publishes membership, lifecycle and provenance in one durable
269 // registry replacement. File publication must have been validated beforehand.
270 func (s *Store) CommitOperation(ctx context.Context, id string) error {
271 return s.mutate(ctx, func(state *State) error {
272 if op, ok := state.PendingOperations[id]; ok && op.Kind == "archive-import" {
273 return ErrMutationConflict
274 }
275 if err := validateTopicRemovalArchive(*state, state.PendingOperations[id]); err != nil {
276 return err
277 }
278 return commitOperation(state, id, map[string]bool{})
279 })
280 }
281
282 // CommitHistoricalArchive publishes a standalone legacy trash import without
283 // admitting arbitrary archive-import children intended for an atomic batch.
284 func (s *Store) CommitHistoricalArchive(ctx context.Context, id string, archivedAt ...int64) error {
285 return s.mutate(ctx, func(state *State) error {
286 op, ok := state.PendingOperations[id]
287 if !ok || op.Kind != "archive-import" || op.Lifecycle != Archived || op.Mapping == nil {
288 return ErrMutationConflict
289 }
290 if op.Phase == "committed" {
291 return nil
292 }
293 if err := commitOperation(state, id, map[string]bool{}); err != nil {
294 return err
295 }
296 // No proven archive timestamp is available for these old records.
297 for _, sessionID := range op.SessionIDs {
298 status := state.SessionStates[sessionID]
299 status.ArchivedAt = 0
300 if len(archivedAt) > 0 && archivedAt[0] > 0 {
301 status.ArchivedAt = archivedAt[0]
302 }
303 state.SessionStates[sessionID] = status
304 }
305 return nil
306 })
307 }
308
309 func commitOperation(state *State, id string, visiting map[string]bool) error {
310 if visiting[id] {
311 return ErrMutationConflict
312 }
313 visiting[id] = true
314 defer delete(visiting, id)
315 op, ok := state.PendingOperations[id]
316 if !ok {
317 return ErrMutationConflict
318 }
319 if op.Phase == "committed" {
320 return nil
321 }
322 if op.Phase != "content_ready" || len(op.SessionIDs) == 0 {
323 return ErrMutationConflict
324 }
325 for _, dependency := range op.Dependencies {
326 child, ok := state.PendingOperations[dependency]
327 if !ok || child.Kind != "archive-import" || child.Phase != "content_ready" {
328 return ErrMutationConflict
329 }
330 for _, target := range child.SessionIDs {
331 if !slices.Contains(op.SessionIDs, target) {
332 return ErrMutationConflict
333 }
334 }
335 if err := commitOperation(state, dependency, visiting); err != nil {
336 return err
337 }
338 }
339 for _, sessionID := range op.SessionIDs {
340 if err := validateLifecycleSupersedesPreparedPurge(*state, sessionID); err != nil {
341 return err
342 }
343 if state.SessionStates[sessionID].Lifecycle == Deleted {
344 return ErrMutationConflict
345 }
346 if current, exists := state.SessionStates[sessionID]; exists && current.Generation > op.ExpectedGeneration && current.Generation != state.Generation+1 {
347 return ErrMutationConflict
348 }
349 owner, attached := sessionOwner(*state, sessionID)
350 if attached && op.WorkspaceID != "" && owner != op.WorkspaceID {
351 return ErrMutationConflict
352 }
353 if !attached {
354 workspace, exists := state.Workspaces[op.WorkspaceID]
355 if !exists {
356 return ErrWorkspaceNotFound
357 }
358 workspace.SessionIDs = insertBefore(workspace.SessionIDs, sessionID, "")
359 // Source adoption replaces the imported source slot below. Publishing
360 // a canonical default first would incorrectly override that choice.
361 if op.Mapping == nil {
362 attachOrganizationSession(&workspace, sessionID, "")
363 mirrorOrganizationOrder(&workspace)
364 }
365 workspace.UpdatedAt = time.Now().UTC()
366 state.Workspaces[workspace.ID] = workspace
367 owner = workspace.ID
368 }
369 cancelPreparedPurge(state, sessionID)
370 setLifecycle(state, sessionID, op.Lifecycle)
371 if op.Lifecycle == Active {
372 workspace := state.Workspaces[owner]
373 workspace.Visible = true
374 state.Workspaces[owner] = workspace
375 }
376 if op.Presentation != nil {
377 state.Presentation[sessionID] = *op.Presentation
378 }
379 }
380 if err := commitSourceMapping(state, op.Mapping); err != nil {
381 return err
382 }
383 if op.RecoveryEntryID != "" {
384 entry, exists := state.RecoveryEntries[op.RecoveryEntryID]
385 if !exists {
386 return ErrMutationConflict
387 }
388 entry.Status, entry.SessionID = "restored", op.SessionIDs[0]
389 state.RecoveryEntries[entry.ID] = entry
390 settleRecoveryVersion(state, entry)
391 }
392 op.Phase, op.ResultGeneration = "committed", state.Generation+1
393 state.PendingOperations[id] = op
394 return nil
395 }
396
397 func ClassifyPurge(state State, id string) PurgeState {
398 key := "purge-" + id
399 op, exists := state.PendingOperations[key]
400 if !exists {
401 return PurgeAbsent
402 }
403 status, known := state.SessionStates[id]
404 if !known || op.ID != key || op.Kind != "purge" || op.Lifecycle != Deleted || len(op.SessionIDs) != 1 || op.SessionIDs[0] != id {
405 return PurgeInvalid
406 }
407 switch op.Phase {
408 case "prepared":
409 if status.Lifecycle == Deleted {
410 return PurgeInvalid
411 }
412 if status.Lifecycle != Archived || status.Generation > op.ExpectedGeneration {
413 return PurgePreparedStale
414 }
415 return PurgePrepared
416 case "tombstoned":
417 if status.Lifecycle == Deleted {
418 return PurgeTombstoned
419 }
420 case "content_removed":
421 if status.Lifecycle == Deleted {
422 return PurgeContentRemoved
423 }
424 case "committed":
425 if status.Lifecycle == Deleted {
426 return PurgeCommitted
427 }
428 }
429 return PurgeInvalid
430 }
431
432 func validateLifecycleSupersedesPreparedPurge(state State, id string) error {
433 switch ClassifyPurge(state, id) {
434 case PurgeAbsent, PurgePrepared, PurgePreparedStale:
435 return nil
436 default:
437 return ErrMutationConflict
438 }
439 }
440
441 func cancelPreparedPurge(state *State, id string) {
442 key := "purge-" + id
443 if op, exists := state.PendingOperations[key]; exists && op.Kind == "purge" && op.Phase == "prepared" {
444 delete(state.PendingOperations, key)
445 }
446 }
447
448 func samePurgeIdentity(left, right Operation) bool {
449 return left.ID == right.ID && left.Kind == "purge" && right.Kind == "purge" && left.Lifecycle == right.Lifecycle &&
450 left.ExpectedGeneration == right.ExpectedGeneration && slices.Equal(left.SessionIDs, right.SessionIDs) && bytes.Equal(left.Request, right.Request)
451 }
452
453 // BeginPurge atomically validates the archived generation, publishes the
454 // deletion tombstone and records the resumable purge operation.
455 func (s *Store) BeginPurge(ctx context.Context, id string, expected uint64) error {
456 return s.beginOrResumePurge(ctx, id, expected, nil, false)
457 }
458
459 // ResumePurge continues only the observed operation. It cannot recreate a
460 // deletion intent after a restore superseded that operation.
461 func (s *Store) ResumePurge(ctx context.Context, id string, observed Operation) error {
462 return s.beginOrResumePurge(ctx, id, observed.ExpectedGeneration, &observed, false)
463 }
464
465 // ResumePurgeForRequest keeps both the request snapshot and the observed
466 // transaction identity. Neither may be refreshed while waiting for locks.
467 func (s *Store) ResumePurgeForRequest(ctx context.Context, id string, expected uint64, observed Operation) error {
468 return s.beginOrResumePurge(ctx, id, expected, &observed, false)
469 }
470
471 func (s *Store) beginOrResumePurge(ctx context.Context, id string, expected uint64, observed *Operation, cleanupSources bool) error {
472 stale := false
473 err := s.mutate(ctx, func(state *State) error {
474 key := "purge-" + id
475 current, exists := state.PendingOperations[key]
476 if observed != nil && (!exists || !samePurgeIdentity(current, *observed)) {
477 return fmt.Errorf("%w: observed purge was removed or replaced", ErrMutationConflict)
478 }
479 switch ClassifyPurge(*state, id) {
480 case PurgePreparedStale:
481 delete(state.PendingOperations, key)
482 stale = true
483 return nil
484 case PurgePrepared:
485 if observed == nil {
486 return ErrMutationConflict
487 }
488 status := state.SessionStates[id]
489 if status.Generation > expected {
490 return fmt.Errorf("%w: purge generation %d exceeds request %d", ErrMutationConflict, status.Generation, expected)
491 }
492 setLifecycle(state, id, Deleted)
493 current.Phase = "tombstoned"
494 state.PendingOperations[key] = current
495 return nil
496 case PurgeTombstoned, PurgeContentRemoved, PurgeCommitted:
497 return nil
498 case PurgeInvalid:
499 return fmt.Errorf("%w: inconsistent purge state", ErrMutationConflict)
500 case PurgeAbsent:
501 if observed != nil {
502 return ErrMutationConflict
503 }
504 status, known := state.SessionStates[id]
505 if !known || status.Lifecycle != Archived || status.Generation > expected {
506 return ErrMutationConflict
507 }
508 op := Operation{ID: key, Kind: "purge", Phase: "tombstoned", Lifecycle: Deleted, SessionIDs: []string{id}, ExpectedGeneration: status.Generation}
509 if cleanupSources {
510 var err error
511 op.Request, err = purgeSourceCleanupRequest(*state, id)
512 if err != nil {
513 return err
514 }
515 }
516 state.PendingOperations[key] = op
517 setLifecycle(state, id, Deleted)
518 return nil
519 default:
520 return ErrMutationConflict
521 }
522 })
523 if err != nil {
524 return err
525 }
526 if stale {
527 return fmt.Errorf("%w: stale purge preparation removed", ErrMutationConflict)
528 }
529 return nil
530 }
531
532 func (s *Store) AdvancePurge(ctx context.Context, id, phase string) error {
533 return s.mutate(ctx, func(state *State) error {
534 key := "purge-" + id
535 op, ok := state.PendingOperations[key]
536 if !ok || op.Kind != "purge" {
537 return ErrMutationConflict
538 }
539 if ClassifyPurge(*state, id) == PurgeInvalid {
540 return ErrMutationConflict
541 }
542 if op.Phase == "committed" || op.Phase == phase || op.Phase == "content_removed" {
543 return nil
544 }
545 if phase != "content_removed" || op.Phase != "tombstoned" {
546 return ErrMutationConflict
547 }
548 op.Phase = phase
549 state.PendingOperations[key] = op
550 return nil
551 })
552 }
553
554 func (s *Store) CompletePurge(ctx context.Context, id string) error {
555 return s.mutate(ctx, func(state *State) error {
556 key := "purge-" + id
557 op := state.PendingOperations[key]
558 switch ClassifyPurge(*state, id) {
559 case PurgeCommitted:
560 return nil
561 case PurgeContentRemoved:
562 default:
563 return ErrMutationConflict
564 }
565 // Keep only the consumed topic identity in the existing purge receipt.
566 // Canonical sessions and runtime-adopted sources may have no import
567 // journal from which a future display reader could recover that identity.
568 if topicID := state.Presentation[id].TopicID; topicID != "" {
569 op.WorkspaceID, _ = sessionOwner(*state, id)
570 if op.Presentation == nil {
571 op.Presentation = &Presentation{TopicID: topicID}
572 }
573 }
574 for key, workspace := range state.Workspaces {
575 if slices.Contains(workspace.SessionIDs, id) {
576 op.WorkspaceID = workspace.ID
577 }
578 workspace.SessionIDs = remove(workspace.SessionIDs, id)
579 state.Workspaces[key] = workspace
580 }
581 // The content and presentation can go, but their topic ownership must
582 // survive: otherwise a residual metadata row looks like a new topic.
583 if topicID := state.Presentation[id].TopicID; topicID != "" {
584 if op.Presentation == nil {
585 op.Presentation = &Presentation{}
586 }
587 op.Presentation.TopicID = topicID
588 }
589 delete(state.Presentation, id)
590 op.Phase, op.ResultGeneration = "committed", state.Generation+1
591 state.PendingOperations[key] = op
592 return nil
593 })
594 }
595
596 func (s *Store) UpdatePresentation(ctx context.Context, ids []string, title *string, pinned *bool) error {
597 return s.mutate(ctx, func(state *State) error {
598 for _, id := range ids {
599 if _, ok := sessionOwner(*state, id); !ok {
600 return ErrSessionNotFound
601 }
602 value, exists := state.Presentation[id]
603 if !exists {
604 value.SortOrder = -1
605 }
606 if title != nil {
607 value.Title = *title
608 }
609 if pinned != nil {
610 value.Pinned = *pinned
611 }
612 state.Presentation[id] = value
613 }
614 return nil
615 })
616 }
617
618 // EnsureSessionTopic publishes the initial display group without letting
619 // later tab rebuilds overwrite a user's persisted presentation.
620 func (s *Store) EnsureSessionTopic(ctx context.Context, id, topicID, title string) error {
621 if id == "" || topicID == "" {
622 return nil
623 }
624 return s.mutate(ctx, func(state *State) error {
625 if _, ok := sessionOwner(*state, id); !ok {
626 return ErrSessionNotFound
627 }
628 value := state.Presentation[id]
629 if value.TopicID != "" {
630 return nil
631 }
632 value.TopicID = topicID
633 if value.Title == "" {
634 value.Title = title
635 }
636 state.Presentation[id] = value
637 return nil
638 })
639 }
640
641 func (s *Store) RecordRecovery(ctx context.Context, entry RecoveryEntry) error {
642 return s.mutate(ctx, func(state *State) error {
643 if _, restored := restoredRecoveryVersion(*state, entry.SourceKey, entry.Fingerprint); restored {
644 return nil
645 }
646 if old, ok := state.RecoveryEntries[entry.ID]; ok {
647 if old.Status == "restored" && old.Fingerprint == entry.Fingerprint {
648 return nil
649 }
650 if strings.Contains(old.Reason, "conflict") && !strings.Contains(entry.Reason, "conflict") {
651 entry.Reason = old.Reason
652 }
653 entry.extra = old.extra
654 }
655 if entry.Status == "" {
656 entry.Status = "pending"
657 }
658 state.RecoveryEntries[entry.ID] = entry
659 return nil
660 })
661 }
662
663 // ReconcileDiscoveredSession rechecks ownership under the writer lock. A scan
664 // snapshot can predate a concurrent create, import, archive, or restore.
665 // Discovery must never classify those published/reserved IDs as orphans.
666 func (s *Store) ReconcileDiscoveredSession(ctx context.Context, entry RecoveryEntry, workspace *Workspace) error {
667 return s.mutate(ctx, func(state *State) error {
668 id := entry.SessionID
669 if state.SessionStates[id].Lifecycle == Deleted {
670 return nil
671 }
672 if id == "" {
673 return errors.New("discovery requires a session id")
674 }
675 if _, ok := sessionOwner(*state, id); ok {
676 return nil
677 }
678 if _, ok := state.PendingCreates[id]; ok {
679 return nil
680 }
681 for _, op := range state.PendingOperations {
682 if op.Phase != "committed" && slices.Contains(op.SessionIDs, id) {
683 return nil
684 }
685 }
686 if lifecycle, ok := state.SessionStates[id]; ok && lifecycle.Lifecycle != Active {
687 workspace = nil
688 entry.Reason = "historical_state_unknown"
689 }
690 if workspace == nil {
691 if old, ok := state.RecoveryEntries[entry.ID]; ok {
692 if old.Status == "restored" {
693 return nil
694 }
695 entry.extra = old.extra
696 }
697 state.RecoveryEntries[entry.ID] = entry
698 return nil
699 }
700 // A discovery snapshot can predate another process registering this
701 // directory. Resolve its owner under the writer lock, just like session
702 // ownership above, and preserve the owner's title and visibility.
703 workspaceID, found, err := ResolveWorkspaceID(*state, workspace.Root)
704 if err != nil {
705 return err
706 }
707 if !found {
708 workspaceID = workspace.ID
709 if existing, exists := state.Workspaces[workspaceID]; exists && existing.Root != workspace.Root {
710 if workspaceID != GlobalWorkspaceID || strings.TrimSpace(workspace.Root) == "" {
711 return ErrMutationConflict
712 }
713 rebindGlobalRoot(state, workspace.Root, time.Now().UTC())
714 }
715 }
716 value, ok := state.Workspaces[workspaceID]
717 if !ok {
718 value = *workspace
719 state.WorkspaceIDs = append(state.WorkspaceIDs, value.ID)
720 }
721 value.SessionIDs = append(value.SessionIDs, id)
722 value.UpdatedAt = time.Now().UTC()
723 state.Workspaces[value.ID] = value
724 return nil
725 })
726 }
727
728 // backupV1 runs under the same cross-process lock as the schema publication.
729 // Content-addressed backups are immutable: a repeat never overwrites evidence.
730 func backupV1(path string) error {
731 body, err := os.ReadFile(path)
732 if os.IsNotExist(err) {
733 return nil
734 }
735 if err != nil {
736 return err
737 }
738 var header struct {
739 Version int `json:"version"`
740 }
741 if err := json.Unmarshal(body, &header); err != nil {
742 return err
743 }
744 if header.Version != 1 {
745 return nil
746 }
747 digest := sha256.Sum256(body)
748 backup := filepath.Join(filepath.Dir(path), "upgrade-backups", "workspace-v1-"+hex.EncodeToString(digest[:])+".json")
749 if old, err := os.ReadFile(backup); err == nil {
750 if sha256.Sum256(old) != digest {
751 return errors.New("workspace upgrade backup is corrupt")
752 }
753 return nil
754 } else if !os.IsNotExist(err) {
755 return err
756 }
757 if err := os.MkdirAll(filepath.Dir(backup), 0700); err != nil {
758 return err
759 }
760 return fileutil.AtomicWriteFileStrict(backup, body, 0600)
761 }
762
763 // unknownFields retains future, user-owned metadata during read/modify/write.
764 func unknownFields(body []byte, names ...string) (map[string]json.RawMessage, error) {
765 var fields map[string]json.RawMessage
766 if err := json.Unmarshal(body, &fields); err != nil {
767 return nil, err
768 }
769 for _, name := range names {
770 delete(fields, name)
771 }
772 return fields, nil
773 }
774
775 func commitSourceMapping(state *State, mapping *SourceMapping) error {
776 if mapping != nil {
777 mapping := *mapping
778 if old, exists := state.SourceMappings[mapping.SourceKey]; exists && (old.SessionID != mapping.SessionID || old.Fingerprint != mapping.Fingerprint) {
779 return ErrMutationConflict
780 }
781 state.SourceMappings[mapping.SourceKey] = mapping
782 adoptOrganizationSource(state, mapping)
783 }
784 return nil
785 }
786
786 lines GO