返回 DeepSeek-Reasonix
session_topic_lazy_snapshot.go
根目录 / desktop / session_topic_lazy_snapshot.go
1 package main
2
3 import (
4 "encoding/json"
5 "errors"
6 "fmt"
7 "os"
8 "sort"
9 "strings"
10
11 "reasonix/desktop/internal/workspacestate"
12 "reasonix/internal/sessioncatalog"
13 )
14
15 type topicPagePosition struct {
16 cursor string
17 extra int
18 }
19
20 // Ordinary legacy history is read directly from the existing catalog WAL
21 // snapshot. Only requested pages are decoded; no all-pages loop or temporary
22 // copy of the whole result is needed to return the first fifty rows.
23 func (a *App) lazyProjectTopicSnapshot(req ProjectTopicPageRequest, reader workspaceSessionInfoReader, snap *readSnapshot, state workspacestate.State, workspaceID string, org workspacestate.Organization, versions *workspacestate.ReadVersions) (ProjectTopicPage, func() error, bool, error) {
24 catalog := a.sessionCatalog.Load()
25 if catalog == nil || !catalog.MetadataOnly() || org.ManualOrderEnabled {
26 return ProjectTopicPage{}, nil, false, nil
27 }
28 if err := applyOrganizationGroupFilter(&req, org); err != nil {
29 return ProjectTopicPage{}, nil, true, err
30 }
31 // Freeze relative filters once so canonical extras and every catalog page
32 // use the same boundary even if the cursor is resumed much later.
33 req.timeCutoff = desktopSessionTimeCutoff(req.TimeFilter)
34 lease, err := catalog.OpenReadLease(a.desktopSessions.readSnapshots.ctx)
35 if errors.Is(err, sessioncatalog.ErrReadLeaseUnavailable) {
36 return ProjectTopicPage{}, nil, false, nil
37 }
38 if err != nil {
39 return ProjectTopicPage{}, nil, true, err
40 }
41 snap.closeRead = lease.Close
42 ctx := lease.Context(a.bootContext())
43 // Decide which adapter can represent this exact WAL snapshot. A head
44 // discovered after this transaction starts belongs to the next refresh.
45 if multiple, err := catalog.HasMultipleHeads(ctx, req.Scope, req.WorkspaceRoot); err != nil || multiple {
46 lease.Close()
47 snap.closeRead = nil
48 return ProjectTopicPage{}, nil, err != nil, err
49 }
50 excluded, flat := lazyTopicSourceExclusions(state, workspaceID)
51 if !flat {
52 lease.Close()
53 snap.closeRead = nil
54 return ProjectTopicPage{}, nil, false, nil
55 }
56 workspace := state.Workspaces[workspaceID]
57 workspace.SessionIDs = admittedWorkspaceTopicMembers(req, state, workspace)
58 infos, _ := listWorkspaceSessionInfo(a.bootContext(), reader, workspace.SessionIDs)
59 sources := a.historicalCanonicalTopicsFromProjection(req.Scope, req.WorkspaceRoot, state, workspacestate.NewWorkspaceIndex(state))
60 if saved, err := readHistoricalSidecar(); err == nil {
61 applyHistoricalPresentations(sources, saved)
62 }
63 adoptedTopics := state.AdoptedTopicIDs(workspaceID)
64 // Unbacked user-created topics remain visible without inspecting a body.
65 for _, node := range a.withRemovablePlaceholderTopics(req, state, nil, adoptedTopics) {
66 found, err := catalog.HasTopicSessions(ctx, req.Scope, req.WorkspaceRoot, node.TopicID)
67 if err != nil {
68 return ProjectTopicPage{}, nil, true, err
69 }
70 if !found {
71 sources = append(sources, node)
72 }
73 }
74 extras := a.indexedWorkspaceTopics(req, state, workspace, org, infos, sources)
75 // This comparator also covers ties between flat catalog rows. Physical
76 // paths are stable identities; a background activity update cannot reorder
77 // a view whose WAL transaction has already been captured.
78 less := func(left, right ProjectNode) bool {
79 if left.Pinned == right.Pinned && projectTopicSortValue(left.CreatedAt, left.LastActivityAt, req.SortMode) == projectTopicSortValue(right.CreatedAt, right.LastActivityAt, req.SortMode) && left.TopicID == right.TopicID {
80 return topicPageIdentity(left) < topicPageIdentity(right)
81 }
82 return projectTopicLess(left, right, req.SortMode, false)
83 }
84 sort.SliceStable(extras, func(i, j int) bool { return less(extras[i], extras[j]) })
85 excludedJSON, _ := json.Marshal(excluded)
86 var deletedTopicsJSON []byte
87 if deleted := loadProjectsFile().DeletedTopics; len(deleted) > 0 {
88 deletedTopicsJSON, _ = json.Marshal(deleted)
89 }
90 query := sessioncatalog.OrdinaryPageRequest{Scope: req.Scope, WorkspaceRoot: req.WorkspaceRoot, SortMode: req.SortMode, PinnedOnly: req.pinnedOnly, ExcludePinned: req.ExcludePinned, ExcludedPathsJSON: string(excludedJSON), MinActivity: req.timeCutoff}
91 query.ExcludedTopicIDsJSON = string(deletedTopicsJSON)
92 groupJSON := ordinaryGroupSourceKeys(req)
93 // The retained page closure only needs the encoded membership predicate.
94 // Do not keep another copy of every organization's member slice alive.
95 req.groupAll, req.groupSelected = nil, nil
96 switch req.GroupFilter {
97 case "group":
98 query.IncludeSourceKeysJSON = groupJSON
99 case "ungrouped":
100 query.ExcludeSourceKeysJSON = groupJSON
101 }
102 query.ExcludeSourceKeysJSON = a.excludeUnavailableHistoricalSources(query.ExcludeSourceKeysJSON)
103 // Freeze the localized default along with the search predicate. SQLite's
104 // lower() does not implement the existing Go Unicode matching semantics.
105 recordTitle := ordinaryRecordTitle(a.localizedDefaultTopicTitle())
106 var match func(sessioncatalog.OrdinaryRecord) bool
107 if text := strings.ToLower(strings.TrimSpace(req.Query)); text != "" {
108 match = func(record sessioncatalog.OrdinaryRecord) bool {
109 key := "source\x00" + localDesktopHostID + "\x00" + record.SourceKey()
110 return strings.Contains(strings.ToLower(recordTitle(record)+"\n"+record.Preview+"\n"+key), text)
111 }
112 }
113 // Snapshot memory accounts for retained metadata and cursor checkpoints,
114 // including roots retained while a caller has not yet requested page two.
115 store := &a.desktopSessions.readSnapshots
116 fence := &readSourceFence{app: a, files: map[string]os.FileInfo{}, bindings: map[string]string{}, versions: versions, store: store, snapshot: snap, metadataOnly: true}
117 encoded, err := json.Marshal(extras)
118 if err != nil {
119 return ProjectTopicPage{}, nil, true, err
120 }
121 if err := store.reserve(snap, int64(len(encoded)+len(excludedJSON)+len(deletedTopicsJSON)+len(groupJSON)+2*len(req.Query)+1024)); err != nil {
122 return ProjectTopicPage{}, nil, true, err
123 }
124 snap.readPage = (&lazyTopicPageReader{app: a, catalog: catalog, lease: lease, req: req,
125 query: query, extras: extras, positions: map[int]topicPagePosition{0: {}},
126 match: match, recordTitle: recordTitle, less: less, fence: fence}).page
127 availability := a.catalogWorkspaceAvailability(catalog, req.Scope, req.WorkspaceRoot, ctx)
128 page := availability.decorate(ProjectTopicPage{Items: []ProjectNode{}}, catalog.Status().Revision+state.Generation)
129 validateWorkspace := a.workspaceReadFence(versions, workspace, extras)
130 return page, func() error {
131 current, err := a.workspaceRegistry().VerifySnapshot(a.bootContext())
132 if err != nil {
133 return err
134 }
135 if err := validateWorkspace(current); err != nil {
136 return err
137 }
138 return fence.validateWithCurrent(current)
139 }, true, nil
140 }
141
142 func lazyTopicSourceExclusions(state workspacestate.State, workspaceID string) ([]string, bool) {
143 excluded := []string{}
144 for _, mapping := range state.SourceMappings {
145 // Canonical directories have one source identity even if an older
146 // migration receipt retained their originating head ID.
147 if mapping.Format != "canonical" && !sourceMappingHasPathAlias(mapping) {
148 if mapping.WorkspaceID != workspaceID {
149 continue
150 }
151 // The catalog may not yet know about a new sibling, or the head
152 // index may be unavailable. Only the head-aware adapter can filter
153 // adopted identities without inventing a new path-only source.
154 return nil, false
155 }
156 // A sidecar can list a file outside the workspace that adopted it.
157 excluded = append(excluded, mapping.Path)
158 }
159 return excluded, true
160 }
161
162 // Group membership comes from the same immutable organization as the canonical
163 // rows. A group contains explicit session keys after preference import, so
164 // sessions sharing a topic must not inherit each other's membership.
165 func applyOrganizationGroupFilter(req *ProjectTopicPageRequest, org workspacestate.Organization) error {
166 req.groupInclude, req.groupExclude = nil, nil
167 req.groupIncludeJSON, req.groupExcludeJSON = "", ""
168 req.groupAll = organizationSnapshot(org, true).Groups
169 req.groupSelected = nil
170 if req.GroupFilter == "group" {
171 for i := range req.groupAll {
172 if req.groupAll[i].ID == req.GroupID {
173 req.groupSelected = &req.groupAll[i]
174 return nil
175 }
176 }
177 return fmt.Errorf("session group no longer exists")
178 }
179 return nil
180 }
181
182 func ordinaryGroupSourceKeys(req ProjectTopicPageRequest) string {
183 if req.GroupFilter != "group" && req.GroupFilter != "ungrouped" {
184 return ""
185 }
186 groups := req.groupAll
187 if req.GroupFilter == "group" {
188 groups = []desktopGroup{*req.groupSelected}
189 }
190 keys, seen := []string{}, map[string]bool{}
191 const prefix = "source\x00" + localDesktopHostID + "\x00"
192 for _, group := range groups {
193 for _, member := range group.SessionKeys {
194 if key, ok := strings.CutPrefix(strings.TrimSpace(member), prefix); ok && !seen[key] {
195 keys, seen[key] = append(keys, key), true
196 }
197 }
198 }
199 encoded, _ := json.Marshal(keys)
200 return string(encoded)
201 }
202
203 func topicPageIdentity(node ProjectNode) string {
204 if node.Session != nil {
205 return "ref\x00" + node.Session.HostID + "\x00" + node.Session.SessionID
206 }
207 if node.Source != nil {
208 return "source\x00" + node.Source.Path + "\x00" + node.Source.HeadID
209 }
210 return "node\x00" + node.Key
211 }
212
213 func ordinaryRecordTitle(defaultTitle string) func(sessioncatalog.OrdinaryRecord) string {
214 return func(record sessioncatalog.OrdinaryRecord) string {
215 if title := strings.TrimSpace(record.CustomTitle); title != "" {
216 return title
217 }
218 if strings.TrimSpace(record.TitleSource) == topicTitleSourceAuto && isDefaultTopicTitle(record.Title) {
219 return defaultTitle
220 }
221 return record.Title
222 }
223 }
224
224 lines GO