返回 DeepSeek-Reasonix
source_reconciliation.go
根目录 / desktop / internal / workspacestate / source_reconciliation.go
1 package workspacestate
2
3 import (
4 "context"
5 "reflect"
6 "slices"
7 "strings"
8 )
9
10 // Archive reservations own unpublished identities. Neither an ordinary
11 // preparation nor a second source can publish an archive child active.
12 func validateArchiveImportReservation(state *State, op Operation) error {
13 if op.Kind != "import" && op.Kind != "archive-import" && !(op.Kind == "restore" && op.Mapping != nil) {
14 return nil
15 }
16 for _, other := range state.PendingOperations {
17 if other.Phase == "committed" || (other.Kind != "import" && other.Kind != "archive-import") {
18 continue
19 }
20 if op.Kind != "archive-import" && other.Kind != "archive-import" {
21 continue
22 }
23 for _, id := range op.SessionIDs {
24 if slices.Contains(other.SessionIDs, id) {
25 return ErrMutationConflict
26 }
27 }
28 }
29 return nil
30 }
31
32 func (s *Store) RecordSource(ctx context.Context, mapping SourceMapping, presentation Presentation) error {
33 return s.mutate(ctx, func(state *State) error {
34 if old, ok := state.SourceMappings[mapping.SourceKey]; ok {
35 if old.SessionID != mapping.SessionID || old.Fingerprint != mapping.Fingerprint {
36 return ErrMutationConflict
37 }
38 return nil
39 }
40 state.SourceMappings[mapping.SourceKey] = mapping
41 adoptOrganizationSource(state, mapping)
42 if _, exists := state.Presentation[mapping.SessionID]; !exists {
43 state.Presentation[mapping.SessionID] = presentation
44 }
45 return nil
46 })
47 }
48
49 // RecordRecoveredSource repairs a missing current receipt without changing the
50 // destination's content, presentation or lifecycle. The caller proves adoption
51 // from frozen content; the transaction fences competing owners and lifecycle.
52 func (s *Store) RecordRecoveredSource(ctx context.Context, mapping SourceMapping, generation uint64) error {
53 return s.mutate(ctx, func(state *State) error {
54 if state.Generation != generation || mapping.SourceKey == "" || mapping.Fingerprint == "" {
55 return ErrMutationConflict
56 }
57 if existing, found := state.SourceMappings[mapping.SourceKey]; found {
58 if existing.SessionID != mapping.SessionID || existing.Fingerprint != mapping.Fingerprint {
59 return ErrMutationConflict
60 }
61 return nil
62 }
63 lifecycle := state.SessionStates[mapping.SessionID].Lifecycle
64 if !validLifecycle(lifecycle) {
65 return ErrMutationConflict
66 }
67 if lifecycle != Deleted {
68 if owner, found := sessionOwner(*state, mapping.SessionID); !found || owner != mapping.WorkspaceID {
69 return ErrMutationConflict
70 }
71 }
72 for _, old := range state.SourceMappings {
73 if !slices.Contains(state.SourceKeys(old.SourceKey), mapping.SourceKey) {
74 continue
75 }
76 if old.SessionID == mapping.SessionID && old.WorkspaceID == mapping.WorkspaceID && old.Fingerprint == mapping.Fingerprint {
77 continue
78 } else if state.SessionStates[old.SessionID].Lifecycle != Deleted {
79 return ErrMutationConflict
80 }
81 }
82 for _, op := range state.PendingOperations {
83 if op.Phase != "committed" && op.Mapping != nil && slices.Contains(state.SourceKeys(op.Mapping.SourceKey), mapping.SourceKey) {
84 return ErrMutationConflict
85 }
86 }
87 if err := backupHistoricalVersionRegistry(s.path); err != nil {
88 return err
89 }
90 state.SourceMappings[mapping.SourceKey] = mapping
91 return nil
92 })
93 }
94
95 // RecordRetiredSource repairs a missing receipt without recreating membership
96 // or presentation. The caller proves source content while holding its read lock;
97 // the generation fence protects the observed lifecycle and competing mappings.
98 func (s *Store) RecordRetiredSource(ctx context.Context, mapping SourceMapping, generation uint64) error {
99 return s.mutate(ctx, func(state *State) error {
100 if state.Generation != generation || mapping.SourceKey == "" || mapping.Fingerprint == "" {
101 return ErrMutationConflict
102 }
103 lifecycle := state.SessionStates[mapping.SessionID].Lifecycle
104 if lifecycle != Deleted && lifecycle != Archived {
105 return ErrMutationConflict
106 }
107 if old, ok, err := state.ResolveSource(mapping.SourceKey); err != nil {
108 return err
109 } else if ok {
110 if old.SessionID != mapping.SessionID || old.Fingerprint != mapping.Fingerprint {
111 return ErrMutationConflict
112 }
113 return nil
114 }
115 for _, op := range state.PendingOperations {
116 if op.Phase != "committed" && op.Mapping != nil && op.Mapping.SourceKey == mapping.SourceKey {
117 return ErrMutationConflict
118 }
119 }
120 if err := backupHistoricalVersionRegistry(s.path); err != nil {
121 return err
122 }
123 state.SourceMappings[mapping.SourceKey] = mapping
124 return nil
125 })
126 }
127
128 // StageHistoricalArchive transfers an unpublished import to an explicit archive
129 // request. Reusing its reservation prevents recovery from publishing it active.
130 func (s *Store) StageHistoricalArchive(ctx context.Context, observed Operation, generation uint64, destinations ...string) error {
131 return s.mutate(ctx, func(state *State) error {
132 op, ok := state.PendingOperations[observed.ID]
133 if !ok || state.Generation != generation || !reflect.DeepEqual(op, observed) || op.Kind != "import" ||
134 (op.Phase != "prepared" && op.Phase != "content_ready") || op.Mapping == nil || len(op.SessionIDs) != 1 {
135 return ErrMutationConflict
136 }
137 id := op.SessionIDs[0]
138 if op.Mapping.SessionID != id || op.Mapping.WorkspaceID != op.WorkspaceID || op.RecoveryEntryID != "" || len(op.Dependencies) != 0 {
139 return ErrMutationConflict
140 }
141 if _, exists := state.SessionStates[id]; exists {
142 return ErrMutationConflict
143 }
144 if _, attached := sessionOwner(*state, id); attached {
145 return ErrMutationConflict
146 }
147 for key, other := range state.PendingOperations {
148 if key != op.ID && other.Phase != "committed" && (slices.Contains(other.SessionIDs, id) || slices.Contains(other.Dependencies, op.ID)) {
149 return ErrMutationConflict
150 }
151 }
152 if len(destinations) > 0 && destinations[0] != op.Mapping.SourceKey {
153 var err error
154 op, err = rekeyHistoricalArchive(state, op, destinations[0])
155 if err != nil {
156 return err
157 }
158 }
159 if err := backupHistoricalVersionRegistry(s.path); err != nil {
160 return err
161 }
162 op.Kind, op.Lifecycle = "archive-import", Archived
163 state.PendingOperations[op.ID] = op
164 return nil
165 })
166 }
167
168 // rekeyHistoricalArchive keeps a changed source version separate from its
169 // original receipt while preserving the unpublished target reservation.
170 func rekeyHistoricalArchive(state *State, op Operation, key string) (Operation, error) {
171 base := strings.TrimSuffix(key, ":review:"+op.Mapping.Fingerprint)
172 if base == key || !slices.Contains(state.SourceKeys(op.Mapping.SourceKey), base) {
173 return op, ErrMutationConflict
174 }
175 if _, exists, err := state.ResolveSource(key); exists || err != nil {
176 return op, ErrMutationConflict
177 }
178 for otherID, other := range state.PendingOperations {
179 if otherID != op.ID && other.Phase != "committed" && other.Mapping != nil && other.Mapping.SourceKey == key {
180 return op, ErrMutationConflict
181 }
182 }
183 mapping := *op.Mapping
184 mapping.SourceKey = key
185 op.Mapping = &mapping
186 return op, nil
187 }
188
188 lines GO