返回 DeepSeek-Reasonix
historical_catalog.go
根目录 / desktop / historical_catalog.go
1 package main
2
3 import (
4 "context"
5 "errors"
6 "log/slog"
7 "os"
8 "path/filepath"
9 "reflect"
10 "sort"
11 "time"
12
13 "reasonix/desktop/internal/workspacestate"
14 "reasonix/internal/historywork"
15 "reasonix/internal/session"
16 )
17
18 type historicalCatalogEntry struct {
19 scope string
20 node ProjectNode
21 sourceChanged bool
22 }
23
24 // Discovery repairs incomplete adoption before publishing metadata. Ordinary
25 // pagination never visits source directories or replays a historical log.
26 func (a *App) requestHistoricalCatalog() {
27 a.requestHistoricalCatalogWithContext(a.bootContext())
28 }
29
30 // requestHistoricalCatalogWithContext lets lifecycle owners pass the context
31 // that already governs their worker. Background catalog callbacks must not
32 // reread App.ctx while tests or the shell are replacing that field.
33 func (a *App) requestHistoricalCatalogWithContext(baseCtx context.Context) {
34 c := &a.historicalImports
35 c.mu.Lock()
36 c.initialize(baseCtx)
37 if !c.catalogEnabled || c.stopped || a.shuttingDown.Load() || c.discoveryPending || time.Since(c.catalogAt) < 5*time.Minute {
38 c.mu.Unlock()
39 return
40 }
41 c.discoveryPending = true
42 ctx := c.ctx
43 c.workers.Add(1)
44 c.mu.Unlock()
45 go func() {
46 defer c.workers.Done()
47 defer func() {
48 c.mu.Lock()
49 c.discoveryPending = false
50 c.mu.Unlock()
51 }()
52 if _, err := a.discoverHistoricalSessions(ctx, false); err != nil && !errors.Is(err, context.Canceled) {
53 slog.Warn("desktop: historical catalog discovery incomplete")
54 }
55 }()
56 }
57
58 func (a *App) listHistoricalSessions(ctx context.Context) (HistoricalImportStatus, error) {
59 return a.discoverHistoricalSessions(ctx, true)
60 }
61
62 func (a *App) discoverHistoricalSessions(ctx context.Context, includeLegacy bool) (HistoricalImportStatus, error) {
63 c := &a.historicalImports
64 c.discoveryMu.Lock()
65 defer c.discoveryMu.Unlock()
66 // Failed reads also settle the refresh interval. Otherwise a renderer read
67 // can immediately re-admit the same failed discovery.
68 defer c.finishCatalogRefresh()
69 state, err := a.workspaceRegistry().Load(ctx)
70 if err != nil {
71 return HistoricalImportStatus{Items: []HistoricalSessionView{}}, err
72 }
73 c.mu.Lock()
74 positions := make(map[string]int, len(c.catalog))
75 for i, entry := range c.catalog {
76 positions[entry.node.Key] = i
77 }
78 c.mu.Unlock()
79 publish := a.historicalCatalogPublisher(ctx, positions)
80 sources := map[string]historicalSource{}
81 add := func(path, format, scope, root, head string) {
82 key := desktopSourceKey(path, head)
83 sources[key] = historicalSource{path: path, format: format, scope: scope, root: root, head: head}
84 }
85 canonical, legacy := a.desktopHistoricalRoots()
86 var joined error
87 for _, source := range canonical {
88 joined = errors.Join(joined, scanHistoricalRoot(ctx, *source, "canonical", add, &a.historyMaintenance))
89 }
90 // Ordinary legacy discovery belongs to sessioncatalog. Only the explicit
91 // management listing enumerates it here to preserve that API's semantics.
92 if includeLegacy {
93 for _, source := range legacy {
94 joined = errors.Join(joined, scanHistoricalRoot(ctx, source, "legacy", add, &a.historyMaintenance))
95 }
96 }
97 addHistoricalRegistrySources(state, add)
98 joined = errors.Join(joined, a.addHistoricalCatalogReceipts(ctx, add))
99 joined = errors.Join(joined, addHistoricalLegacyReceiptHeads(sources, add))
100 if err := a.reconcileHistoricalLegacyCatalog(ctx, sources); err != nil {
101 return HistoricalImportStatus{Items: []HistoricalSessionView{}}, err
102 }
103 publishRecovered := func(batch []historicalCatalogEntry) {
104 if err := a.reconcileHistoricalCatalog(ctx, batch, sources); err == nil {
105 publish(batch)
106 }
107 }
108 catalog := readHistoricalCanonicalCatalog(ctx, sources, &a.historyMaintenance, publishRecovered, state)
109 if err := a.reconcileHistoricalCatalog(ctx, catalog, sources); err != nil {
110 return HistoricalImportStatus{Items: []HistoricalSessionView{}}, err
111 }
112 state, err = a.workspaceRegistry().Load(ctx)
113 if err != nil {
114 return HistoricalImportStatus{Items: []HistoricalSessionView{}}, err
115 }
116 saved, presentationErr := readHistoricalSidecar()
117 if err := ctx.Err(); err != nil {
118 return HistoricalImportStatus{Items: []HistoricalSessionView{}}, err
119 }
120 c.mu.Lock()
121 notify := false
122 defer func() {
123 c.mu.Unlock()
124 if notify {
125 a.emitProjectTreeChangedEvent()
126 }
127 }()
128 if c.stopped || ctx.Err() != nil {
129 return HistoricalImportStatus{Items: []HistoricalSessionView{}}, context.Canceled
130 }
131 c.initialize(ctx)
132 if !c.queueLoaded {
133 c.loadQueueLocked()
134 }
135 // An interrupted or inaccessible root is not evidence of deletion. Keep
136 // its previously published rows until a complete discovery can prove absence.
137 if joined != nil {
138 catalog = c.catalog
139 }
140 changed := !reflect.DeepEqual(c.catalog, catalog)
141 if changed {
142 c.catalog = catalog
143 }
144 if presentationErr == nil {
145 changed = changed || !reflect.DeepEqual(c.presentations, saved.Presentations)
146 c.presentations = saved.Presentations
147 }
148 for id, source := range sources {
149 if source.path == "" {
150 continue
151 }
152 c.sources[id] = source
153 view := historicalImportView(state, id, source, c.views[id])
154 if presentation := c.presentations[id]; presentation.Title != "" {
155 view.Title = presentation.Title
156 }
157 changed = changed || !reflect.DeepEqual(c.views[id], view)
158 c.views[id] = view
159 }
160 if changed {
161 c.catalogRevision++
162 notify = !c.stopped
163 }
164 return c.status(), joined
165 }
166
167 func historicalCatalogPlaceholder(key string, source historicalSource) historicalCatalogEntry {
168 kind := "global_topic"
169 if source.scope == "project" {
170 kind = "topic"
171 }
172 return historicalCatalogEntry{scope: source.scope, node: ProjectNode{
173 Key: "source_" + key, Kind: kind, Root: source.root, Label: filepath.Base(source.path),
174 TopicID: "historical-" + key, Historical: true, SessionPath: source.path, SortOrder: -1,
175 TurnsState: "unknown", Health: "metadata_pending", Children: []ProjectNode{},
176 Source: &SessionSourceRef{HostID: localDesktopHostID, SourceKey: key, Path: source.path},
177 }}
178 }
179
180 func readHistoricalCanonicalCatalog(ctx context.Context, sources map[string]historicalSource, maintenance *historywork.Coordinator, publish func([]historicalCatalogEntry), states ...workspacestate.State) []historicalCatalogEntry {
181 rows := []historicalCatalogEntry{}
182 ledger, ledgerErr := readDesktopMigrationLedger()
183 batchStart, count, bytes := 0, 0, int64(0)
184 var release func(int64)
185 var started time.Time
186 flush := func() {
187 if release != nil {
188 release(bytes)
189 release = nil
190 }
191 if publish != nil && batchStart < len(rows) {
192 publish(rows[batchStart:])
193 }
194 batchStart, count, bytes = len(rows), 0, 0
195 }
196 defer func() {
197 if release != nil {
198 release(bytes)
199 }
200 }()
201 ctx = maintenance.Context(ctx)
202 for key, source := range sources {
203 if ctx.Err() != nil {
204 break
205 }
206 if source.format != "canonical" || source.version != "" {
207 continue
208 }
209 // Stat reads at most three small metadata files. Charge their upper
210 // bound, including sentinel bytes used to detect oversized sidecars.
211 const metadataBytes = 3 * (historywork.ReadChunk + 1)
212 if release != nil && (count >= historywork.BatchEntries || bytes+metadataBytes > historywork.BatchBytes || time.Since(started) >= historywork.SliceDuration) {
213 flush()
214 }
215 if release == nil {
216 var err error
217 release, err = maintenance.BackgroundSlice(ctx, false)
218 if err != nil {
219 break
220 }
221 started = time.Now()
222 }
223 count++
224 bytes += metadataBytes
225 node := historicalCatalogPlaceholder(key, source).node
226 if info, err := session.NewFilesystemPersistence(filepath.Dir(source.path)).Stat(ctx, filepath.Base(source.path)); err == nil {
227 if info.Kind == session.SessionKindHeadlessRun {
228 delete(sources, key)
229 continue
230 }
231 if info.Title != "" {
232 node.Label = info.Title
233 }
234 node.Preview, node.Turns = info.Preview, info.Turns
235 node.CreatedAt, node.LastActivityAt = info.CreatedAt.UnixMilli(), info.UpdatedAt.UnixMilli()
236 if info.MetadataStatus == session.MetadataReady {
237 node.TurnsState, node.Health = "valid", "ok"
238 }
239 } else if stat, statErr := os.Stat(source.path); statErr == nil {
240 node.CreatedAt, node.LastActivityAt = stat.ModTime().UnixMilli(), stat.ModTime().UnixMilli()
241 node.Health = "degraded"
242 }
243 changed := false
244 if len(states) > 0 && ledgerErr == nil {
245 changed = historicalCanonicalSourceChanged(states[0], ledger, key, source.path)
246 }
247 rows = append(rows, historicalCatalogEntry{scope: source.scope, node: node, sourceChanged: changed})
248 }
249 // The final batch is published with completion metadata by the caller.
250 // Intermediate budget boundaries publish independently for large roots.
251 sort.Slice(rows, func(i, j int) bool { return rows[i].node.Key < rows[j].node.Key })
252 return rows
253 }
254
255 func historicalCanonicalSourceChanged(state workspacestate.State, ledger desktopMigrationLedger, key, path string) bool {
256 // Registry receipts outlive purged source directories. Missing files
257 // are not a new source revision and must not recreate a sidebar row.
258 info, err := os.Stat(path)
259 if err != nil || !info.IsDir() {
260 return false
261 }
262 if _, mapped, err := historicalMappingForSource(state, key); !mapped || err != nil {
263 return false
264 }
265 baseKey := desktopCanonicalMigrationKey(filepath.Dir(path), filepath.Base(path))
266 revision, err := desktopMigrationSourceRevision(canonicalMigrationSourceFiles(filepath.Dir(path), filepath.Base(path)))
267 if err != nil {
268 return false
269 }
270 known := false
271 for receiptKey, receipt := range ledger.Records {
272 if !historicalSourceKeyMatches(receiptKey, baseKey) || receipt.Status != "completed" {
273 continue
274 }
275 known = true
276 if receipt.SourceRevision == revision {
277 return false
278 }
279 }
280 return known
281 }
282
283 func (a *App) historicalCanonicalTopics(scope, root string, state workspacestate.State) []ProjectNode {
284 index := workspacestate.NewWorkspaceIndex(state)
285 return a.historicalCanonicalTopicsFromProjection(scope, root, state, index)
286 }
287
288 func (a *App) historicalCanonicalTopicsFromProjection(scope, root string, state workspacestate.State, index *workspacestate.WorkspaceIndex) []ProjectNode {
289 a.requestHistoricalCatalog()
290 c := &a.historicalImports
291 c.mu.Lock()
292 defer c.mu.Unlock()
293 rows := []ProjectNode{}
294 for _, entry := range c.catalog {
295 if entry.scope != scope {
296 continue
297 }
298 if scope == "project" && !index.SameRoot(entry.node.Root, root) {
299 continue
300 }
301 node := entry.node
302 if node.Health == "unavailable" {
303 continue
304 }
305 if _, adopted, err := historicalMappingForSource(state, node.Source.SourceKey); adopted || err != nil {
306 continue
307 }
308 node.PreparationStatus = "available"
309 if view, ok := c.views[node.Source.SourceKey]; ok {
310 node.PreparationStatus = view.Status
311 }
312 rows = append(rows, node)
313 }
314 return rows
315 }
316
317 func applyHistoricalPresentations(nodes []ProjectNode, saved historicalImportQueueSidecar) {
318 for i := range nodes {
319 if nodes[i].Source == nil {
320 continue
321 }
322 presentation := saved.Presentations[nodes[i].Source.SourceKey]
323 if presentation.Title != "" {
324 nodes[i].Label = presentation.Title
325 }
326 if presentation.Pinned != nil {
327 nodes[i].Pinned = *presentation.Pinned
328 }
329 }
330 }
331
332 // A shell-only read must not create workspaces or migrate organization state.
333 // Sources without canonical members still need their persisted pin overlays.
334 func (a *App) historicalPinnedShellsFromProjection(req ProjectTopicPageRequest, state workspacestate.State, index *workspacestate.WorkspaceIndex, legacy desktopProject) ([]ProjectNode, error) {
335 workspaceID, _, _ := index.Resolve(req.WorkspaceRoot)
336 if req.Scope != "project" {
337 workspaceID = workspacestate.GlobalWorkspaceID
338 }
339 adopted := map[string]bool{}
340 for _, mapping := range state.SourceMappings {
341 for _, key := range state.SourceKeys(mapping.SourceKey) {
342 adopted["source\x00local\x00"+key] = true
343 }
344 if sourceMappingHasPathAlias(mapping) {
345 adopted[sessionRuntimeKey(mapping.Path)] = true
346 }
347 }
348 page, err := a.unadoptedLegacyTopics(req, adopted, state.AdoptedTopicIDs(workspaceID))
349 if err != nil {
350 return nil, err
351 }
352 nodes := append(page.Items, a.historicalCanonicalTopicsFromProjection(req.Scope, req.WorkspaceRoot, state, index)...)
353 if saved, err := readHistoricalSidecar(); err == nil {
354 applyHistoricalPresentations(nodes, saved)
355 }
356 workspace := state.Workspaces[workspaceID]
357 org := projectedShellOrganization(workspace, state, nodes, legacy)
358 req.pinnedOnly = true
359 pins := filterWorkspaceSessionNodes(req, org, state, workspaceID, nodes)
360 sort.SliceStable(pins, func(i, j int) bool { return projectTopicLess(pins[i], pins[j], req.SortMode, org.ManualOrderEnabled) })
361 return pins, nil
362 }
363
364 func (c *historicalImportCoordinator) finishCatalogRefresh() {
365 c.mu.Lock()
366 c.catalogAt = time.Now()
367 c.mu.Unlock()
368 }
369
369 lines GO