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