| 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 |