返回 DeepSeek-Reasonix
topic_list.go
根目录 / internal / sessioncatalog / topic_list.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "strings"
6 )
7
8 func (c *Catalog) ListTopics(ctx context.Context, req TopicPageRequest) (TopicPage, error) {
9 out := TopicPage{Items: []TopicRecord{}, Revision: c.readRevision(ctx)}
10 req.Scope, req.WorkspaceRoot = normalizeScope(req.Scope, req.WorkspaceRoot)
11 if req.Limit <= 0 {
12 req.Limit = DefaultLimit
13 }
14 if req.Limit > MaxLimit {
15 req.Limit = MaxLimit
16 }
17 cursor, err := decodeCursor(req.Cursor)
18 if err != nil {
19 return out, err
20 }
21 if cursor != nil && cursor.ManualOrder != req.ManualOrder {
22 return out, errCursorSortModeChanged
23 }
24 if cursor != nil && cursor.Binding != req.CursorBinding {
25 return out, errCursorSortModeChanged
26 }
27 rootKey := c.workspaceRootKey(req.Scope, req.WorkspaceRoot)
28 args := []any{req.Scope, rootKey}
29 where := `scope=? AND workspace_root_key=?`
30 if req.PinnedOnly {
31 where += ` AND pinned=1`
32 }
33 if query := strings.TrimSpace(req.Query); query != "" {
34 where += ` AND lower(title) LIKE ?`
35 args = append(args, "%"+strings.ToLower(query)+"%")
36 }
37 if cutoff := timeFilterCutoff(req.TimeFilter, c.readTime(ctx)); cutoff > 0 {
38 where += ` AND last_activity_at>=?`
39 args = append(args, cutoff)
40 }
41 where += ` AND (?='' OR topic_id IN (SELECT value FROM json_each(COALESCE(NULLIF(?,''),'[]'))))` +
42 ` AND (?='' OR topic_id NOT IN (SELECT value FROM json_each(COALESCE(NULLIF(?,''),'[]'))))` +
43 ` AND (?=0 OR pinned=0)`
44 args = append(args, req.IncludeTopicIDsJSON, req.IncludeTopicIDsJSON,
45 req.ExcludeTopicIDsJSON, req.ExcludeTopicIDsJSON, req.ExcludePinned)
46 scanCursor := cursor
47 scanLimit := max(req.Limit+1, 64)
48 for len(out.Items) <= req.Limit {
49 query, pageArgs := topicPageQuery(req, where, args, scanCursor, scanLimit)
50 rows, queryErr := c.readDB(ctx).QueryContext(ctx, query, pageArgs...)
51 if queryErr != nil {
52 return out, queryErr
53 }
54 scanned, scanErr := scanTopicRows(rows, scanLimit)
55 if scanErr != nil {
56 return out, scanErr
57 }
58 rawCount := len(scanned)
59 overflow := false
60 for _, item := range scanned {
61 sessions, listErr := c.listTopicSessionsByRootKey(ctx, TopicKey{
62 Scope: item.Scope, WorkspaceRoot: item.WorkspaceRoot, TopicID: item.TopicID,
63 }, rootKey)
64 if listErr != nil {
65 return TopicPage{Items: []TopicRecord{}, Revision: out.Revision}, listErr
66 }
67 if len(sessions) == 0 {
68 continue
69 }
70 // Skip recovery shells that lost their ordinary representative while
71 // lineage is re-anchored, unless the topic is explicitly pinned.
72 hasOrdinary := false
73 for _, session := range sessions {
74 if session.OrdinaryVisible || (!session.Recovered && !session.RecoveryCopy) {
75 hasOrdinary = true
76 break
77 }
78 }
79 if !hasOrdinary && !item.Pinned {
80 continue
81 }
82 item.Sessions = sessions
83 hydrateTopicDisplay(&item)
84 out.Items = append(out.Items, item)
85 if len(out.Items) > req.Limit {
86 overflow = true
87 break
88 }
89 }
90 if overflow || rawCount < scanLimit || rawCount == 0 {
91 break
92 }
93 scanCursor = cursorForTopic(scanned[rawCount-1], req)
94 }
95 more := len(out.Items) > req.Limit
96 if more {
97 out.Items = out.Items[:req.Limit]
98 }
99 if more && len(out.Items) > 0 {
100 out.NextCursor = encodeCursor(*cursorForTopic(out.Items[len(out.Items)-1], req))
101 }
102 return out, nil
103 }
104
105 func topicPageQuery(req TopicPageRequest, where string, args []any, cursor *pageCursor, limit int) (string, []any) {
106 sortExpression := topicPageSortExpression(req.SortMode)
107 manualSortExpression := topicPageManualSortExpression()
108 pageArgs := append([]any(nil), args...)
109 if cursor != nil && req.ManualOrder {
110 where += ` AND (pinned<? OR (pinned=? AND ` + manualSortExpression + `>?) OR ` +
111 `(pinned=? AND ` + manualSortExpression + `=? AND ` + sortExpression + `<?) OR ` +
112 `(pinned=? AND ` + manualSortExpression + `=? AND ` + sortExpression + `=? AND topic_id>?))`
113 pageArgs = append(pageArgs,
114 cursor.Pinned,
115 cursor.Pinned, cursor.SortOrder,
116 cursor.Pinned, cursor.SortOrder, cursor.Activity,
117 cursor.Pinned, cursor.SortOrder, cursor.Activity, cursor.TopicID,
118 )
119 } else if cursor != nil {
120 where += ` AND (pinned<? OR (pinned=? AND ` + sortExpression + `<?) OR (pinned=? AND ` + sortExpression + `=? AND topic_id>?))`
121 pageArgs = append(pageArgs, cursor.Pinned, cursor.Pinned, cursor.Activity,
122 cursor.Pinned, cursor.Activity, cursor.TopicID)
123 }
124 orderBy := `pinned DESC,` + sortExpression + ` DESC,topic_id ASC`
125 if req.ManualOrder {
126 orderBy = `pinned DESC,` + manualSortExpression + ` ASC,` + sortExpression + ` DESC,topic_id ASC`
127 }
128 pageArgs = append(pageArgs, limit)
129 return `SELECT scope,workspace_root,topic_id,title,title_source,pinned,
130 CASE WHEN metadata_present=1 THEN sort_order ELSE -1 END,
131 turns,turns_state,created_at,last_activity_at,recovery_state,recovery_branch_count,
132 recovery_unresolved_count,recovery_cleanup_eligible_count,health
133 FROM catalog_topics WHERE ` + where + ` ORDER BY ` + orderBy + ` LIMIT ?`, pageArgs
134 }
135
136 func scanTopicRows(rows interface {
137 Next() bool
138 Scan(...any) error
139 Err() error
140 Close() error
141 }, capacity int) ([]TopicRecord, error) {
142 defer rows.Close()
143 // Drain before hydrating sessions: the nested read needs another connection
144 // and an open cursor deadlocks when the in-memory pool is saturated.
145 scanned := make([]TopicRecord, 0, capacity)
146 for rows.Next() {
147 var item TopicRecord
148 if err := rows.Scan(&item.Scope, &item.WorkspaceRoot, &item.TopicID, &item.Title,
149 &item.TitleSource, &item.Pinned, &item.SortOrder, &item.Turns, &item.TurnsState,
150 &item.CreatedAt, &item.LastActivityAt, &item.RecoveryState, &item.RecoveryBranchCount,
151 &item.RecoveryUnresolvedCount, &item.RecoveryCleanupEligibleCount, &item.Health); err != nil {
152 return nil, err
153 }
154 scanned = append(scanned, item)
155 }
156 if err := rows.Err(); err != nil {
157 return nil, err
158 }
159 return scanned, nil
160 }
161
162 func cursorForTopic(topic TopicRecord, req TopicPageRequest) *pageCursor {
163 pinned := 0
164 if topic.Pinned {
165 pinned = 1
166 }
167 return &pageCursor{
168 Pinned: pinned, ManualOrder: req.ManualOrder,
169 SortOrder: topicPageManualSortValue(topic),
170 Activity: topicPageSortValue(topic, req.SortMode), TopicID: topic.TopicID,
171 Binding: req.CursorBinding,
172 }
173 }
174
174 lines GO