返回 DeepSeek-Reasonix
session_topic_lazy_page.go
根目录 / desktop / session_topic_lazy_page.go
1 package main
2
3 import (
4 "context"
5 "encoding/json"
6 "reasonix/internal/sessioncatalog"
7 )
8
9 // The snapshot lease guards both fixed query inputs and cursor checkpoints.
10 type lazyTopicPageReader struct {
11 app *App
12 catalog *sessioncatalog.Catalog
13 lease *sessioncatalog.ReadLease
14 req ProjectTopicPageRequest
15 query sessioncatalog.OrdinaryPageRequest
16 extras []ProjectNode
17 positions map[int]topicPagePosition
18 match func(sessioncatalog.OrdinaryRecord) bool
19 recordTitle func(sessioncatalog.OrdinaryRecord) string
20 less func(ProjectNode, ProjectNode) bool
21 fence *readSourceFence
22 }
23
24 func (r *lazyTopicPageReader) page(readCtx context.Context, offset, limit int) (result [][]byte, hasMore bool, resultErr error) {
25 select {
26 case <-r.catalog.Invalidated():
27 return nil, false, snapshotStale("catalog_replaced")
28 default:
29 }
30 defer func() {
31 select {
32 case <-r.catalog.Invalidated():
33 result, hasMore, resultErr = nil, false, snapshotStale("catalog_replaced")
34 default:
35 }
36 }()
37 a, catalog, lease, req := r.app, r.catalog, r.lease, r.req
38 query, extras, positions := r.query, r.extras, r.positions
39 match, recordTitle, less := r.match, r.recordTitle, r.less
40 fence, store, snap := r.fence, r.fence.store, r.fence.snapshot
41 position, ok := positions[offset]
42 if !ok {
43 return nil, false, snapshotStale("invalid_cursor")
44 }
45 request := query
46 request.Cursor, request.Limit = position.cursor, min(limit+1, sessioncatalog.MaxLimit)
47 records, err := catalog.ListMatchingOrdinarySessions(lease.Context(readCtx), request, match)
48 if err != nil {
49 return nil, false, err
50 }
51 _, runtime := a.catalogRuntimeOverlays()
52 legacy := make([]ProjectNode, 0, len(records))
53 for _, record := range records {
54 overlay := runtime[sessionRuntimeKey(record.Path)]
55 kind := "topic"
56 if req.Scope == "global" {
57 kind = "global_topic"
58 }
59 title := recordTitle(record)
60 legacy = append(legacy, ProjectNode{Key: projectSessionNodeKey(req.Scope, record.Path), Kind: kind, Label: title, Root: req.WorkspaceRoot, TopicID: record.TopicID, SessionPath: record.Path,
61 Historical: true, Source: &SessionSourceRef{HostID: localDesktopHostID, Path: record.Path, SourceKey: record.SourceKey()},
62 Preview: record.Preview, Turns: record.Turns, TurnsState: string(record.TurnsState), Health: string(record.Health), CreatedAt: record.CreatedAt, LastActivityAt: record.LastActivityAt, Pinned: record.Pinned, SortOrder: -1,
63 Recovered: record.Recovered, RecoveryReason: record.RecoveryReason, RecoveryDigest: record.RecoveryDigest, RecoveryParentID: record.ParentID, Open: overlay.open, Running: overlay.running, Status: overlay.status, Children: []ProjectNode{}})
64 }
65 rows := [][]byte{}
66 index := 0
67 for len(rows) < limit && (index < len(legacy) || position.extra < len(extras)) {
68 var node ProjectNode
69 if index < len(legacy) && (position.extra >= len(extras) || less(legacy[index], extras[position.extra])) {
70 node = legacy[index]
71 position.cursor = records[index].Cursor
72 index++
73 } else {
74 node = extras[position.extra]
75 position.extra++
76 }
77 b, err := json.Marshal(node)
78 if err != nil {
79 return nil, false, err
80 }
81 rows = append(rows, b)
82 if node.Session == nil && node.SessionPath != "" {
83 if err := fence.add(lease.Context(readCtx), node.SessionPath); err != nil {
84 return nil, false, err
85 }
86 }
87 }
88 more := index < len(legacy) || position.extra < len(extras) || len(records) == request.Limit
89 if more {
90 if _, exists := positions[offset+len(rows)]; !exists {
91 if err := store.reserve(snap, int64(len(position.cursor)+64)); err != nil {
92 return nil, false, err
93 }
94 positions[offset+len(rows)] = position
95 }
96 }
97 return rows, more, nil
98 }
99
99 lines GO