返回 DeepSeek-Reasonix
read_snapshot.go
根目录 / desktop / internal / workspacestate / read_snapshot.go
1 package workspacestate
2
3 import (
4 "context"
5 "errors"
6 "path/filepath"
7 "slices"
8 "strings"
9
10 "reasonix/internal/pathidentity"
11 )
12
13 // ReadSnapshot is an immutable, coherent view of verified registry bytes. It
14 // owns its maps privately; retaining it never observes a later publication.
15 // Verification is O(file bytes); lookups on a retained snapshot are O(1).
16 type ReadSnapshot struct {
17 state State
18 owners map[string]string
19 conflicts map[string]bool
20 activeTopics map[string]int
21 adoptedTopics map[string]map[string]bool
22 headSources map[string]bool
23 versions *ReadVersions
24 }
25
26 // SessionMetadata deliberately excludes workspace members and organization:
27 // resolving one address must not copy a project's entire session list.
28 type SessionMetadata struct {
29 Workspace Workspace
30 State SessionState
31 Presentation Presentation
32 Registered bool
33 OwnershipConflict bool
34 SharedTopic bool
35 }
36
37 type snapshotVerification struct {
38 done chan struct{}
39 snapshot *ReadSnapshot
40 err error
41 readers int // guarded by Store.verificationMu; includes the initiating reader
42 }
43
44 func (s *Store) publishSnapshotLocked(body []byte, state State) {
45 state.sourceIdentities = newSourceIdentityIndex(state)
46 r := &ReadSnapshot{state: state, owners: make(map[string]string), conflicts: make(map[string]bool), activeTopics: make(map[string]int), headSources: make(map[string]bool), versions: NewReadVersions(state)}
47 r.adoptedTopics = adoptedTopicIndex(state)
48 for key, workspace := range state.Workspaces {
49 for _, id := range workspace.SessionIDs {
50 if _, exists := r.owners[id]; exists {
51 r.conflicts[id] = true
52 }
53 r.owners[id] = key
54 }
55 }
56 for id, p := range state.Presentation {
57 if p.TopicID != "" && state.SessionStates[id].Lifecycle == Active {
58 r.activeTopics[p.TopicID]++
59 }
60 }
61 for _, mapping := range state.SourceMappings {
62 if mapping.HeadID != "" {
63 // This is a locator hint only. Exact source/head mappings and the
64 // selected head still decide adoption; it never proves content.
65 if key, err := sourcePathKey(mapping.Path); err == nil {
66 r.headSources[key] = true
67 }
68 }
69 }
70 s.readBody = body
71 s.readSnapshot.Store(r)
72 }
73
74 // VerifySnapshot checks actual current bytes, including legacy writers that
75 // preserve generation and timestamps. Execution/mutation admission must use
76 // this boundary, never PublishedSnapshot alone.
77 func (s *Store) VerifySnapshot(ctx context.Context) (*ReadSnapshot, error) {
78 if s == nil || strings.TrimSpace(s.path) == "" || s.path == "." {
79 return nil, errors.New("workspace state path is required")
80 }
81 if err := ctx.Err(); err != nil {
82 return nil, err
83 }
84 s.verificationMu.Lock()
85 flight := s.verification
86 leader := flight == nil
87 if leader {
88 flight = &snapshotVerification{done: make(chan struct{})}
89 s.verification = flight
90 }
91 flight.readers++
92 s.verificationMu.Unlock()
93 if leader {
94 // A registry read already in flight is shared by overlapping callers.
95 // One caller's cancellation cannot poison another's validation. The
96 // file read is finite; no task or model lifetime is owned here.
97 s.mu.Lock()
98 flight.err = s.verifySnapshotLocked(context.WithoutCancel(ctx))
99 if flight.err == nil {
100 flight.snapshot = s.readSnapshot.Load()
101 }
102 s.verificationMu.Lock()
103 s.verification = nil
104 close(flight.done)
105 s.verificationMu.Unlock()
106 s.mu.Unlock()
107 }
108 select {
109 case <-ctx.Done():
110 return nil, ctx.Err()
111 case <-flight.done:
112 if err := ctx.Err(); err != nil {
113 return nil, err
114 }
115 return flight.snapshot, flight.err
116 }
117 }
118
119 // PublishedSnapshot returns the last successful verification for display only.
120 // It performs no I/O and can be stale after an external writer's change.
121 func (s *Store) PublishedSnapshot() *ReadSnapshot {
122 if s == nil {
123 return nil
124 }
125 return s.readSnapshot.Load()
126 }
127
128 func (r *ReadSnapshot) Generation() uint64 { return r.state.Generation }
129
130 func (r *ReadSnapshot) WorkspaceMetadata(id string) (Workspace, bool) {
131 w, ok := r.state.Workspaces[id]
132 w.SessionIDs, w.Organization, w.extra = nil, nil, nil
133 w.FormerRoots = slices.Clone(w.FormerRoots)
134 return w, ok
135 }
136
137 func (r *ReadSnapshot) Session(id string) SessionMetadata {
138 status, registered := r.state.SessionStates[id]
139 status.extra = nil
140 p := r.state.Presentation[id]
141 p.extra = nil
142 w, _ := r.WorkspaceMetadata(r.owners[id])
143 others := r.activeTopics[p.TopicID]
144 if status.Lifecycle == Active && p.TopicID != "" {
145 others--
146 }
147 return SessionMetadata{Workspace: w, State: status, Presentation: p, Registered: registered,
148 OwnershipConflict: r.conflicts[id], SharedTopic: p.TopicID != "" && others > 0}
149 }
150
151 func (r *ReadSnapshot) Source(key string) (SourceMapping, bool) {
152 mapping, ok, _ := r.state.ResolveSource(key)
153 return mapping, ok
154 }
155
156 func (r *ReadSnapshot) ResolveSource(key string) (SourceMapping, bool, error) {
157 return r.state.ResolveSource(key)
158 }
159
160 func sourcePathKey(path string) (string, error) {
161 if strings.TrimSpace(path) == "" {
162 return "", errors.New("source path is required")
163 }
164 abs, err := filepath.Abs(strings.TrimSpace(path))
165 if err != nil {
166 return "", err
167 }
168 identity, err := pathidentity.Resolve(abs, pathidentity.Options{FollowLeaf: true})
169 return identity.Key, err
170 }
171
172 func (r *ReadSnapshot) HasHeadSource(path string) (bool, error) {
173 if len(r.headSources) == 0 {
174 return false, nil
175 }
176 key, err := sourcePathKey(path)
177 return r.headSources[key], err
178 }
179
179 lines GO