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