返回 DeepSeek-Reasonix
history_slice_paged.go
根目录 / desktop / history_slice_paged.go
1 package main
2
3 import (
4 "context"
5 "errors"
6 "path/filepath"
7 "strings"
8
9 "reasonix/internal/agent"
10 "reasonix/internal/provider"
11 )
12
13 func (a *App) pagedColdHistorySlice(ctx context.Context, sessionDir, path string, req HistorySliceRequest, heads ...string) (HistorySlice, bool, error) {
14 head := ""
15 if len(heads) > 0 {
16 head = heads[0]
17 }
18 page, handled := emptyHistorySlice(), false
19 err := a.withNativeHistoryPager(ctx, path, head, func(ctx context.Context, pager *agent.DisplayPager, sourceID string) error {
20 handled = true
21 var err error
22 page, _, err = a.historySliceFromPager(ctx, pager, sessionDir, path, req, sourceID)
23 return err
24 })
25 if handled || nativeHistoryPreparationFailure(err) {
26 if err != nil {
27 return emptyHistorySlice(), true, err
28 }
29 return page, true, nil
30 }
31 return HistorySlice{}, false, nil
32 }
33
34 func nativeHistoryPreparationFailure(err error) bool {
35 // Only a positively identified unsupported format may use the legacy
36 // adapter. I/O, parse and cache failures must not silently trigger a full
37 // replay after the bounded reader has failed.
38 return err != nil && !errors.Is(err, agent.ErrDisplayFormatUnsupported)
39 }
40
41 // Compatibility paging and field reads borrow the same preparation and cache
42 // owner as bound readers. Retirement cancels both the read and its decoding;
43 // only the last borrower retires the projection, after the callback returns.
44 func (a *App) withNativeHistoryPager(ctx context.Context, path, head string, read func(context.Context, *agent.DisplayPager, string) error) error {
45 if err := ctx.Err(); err != nil {
46 return err
47 }
48 generation, err := nativeHistorySourceGeneration(path)
49 if err != nil {
50 return errors.Join(agent.ErrDisplaySourceChanged, err)
51 }
52 sourceKey := sessionRuntimeKey(path)
53 a.historyReaders.mu.Lock()
54 if a.historyReaders.closed || a.shuttingDown.Load() {
55 a.historyReaders.mu.Unlock()
56 return context.Canceled
57 }
58 job, release := a.acquireNativeHistoryLocked(path, head, sourceKey, generation)
59 a.historyReaders.mu.Unlock()
60 defer release()
61 ctx, cancel := context.WithCancel(ctx)
62 stop := context.AfterFunc(job.ctx, cancel)
63 defer func() { stop(); cancel() }()
64 pager, err := job.wait(ctx)
65 if err != nil {
66 return err
67 }
68 pager = pager.WithContext(ctx)
69 if err := pager.Validate(); err != nil {
70 return err
71 }
72 err = read(ctx, pager, job.key)
73 if ctx.Err() != nil || job.ctx.Err() != nil {
74 return context.Canceled
75 }
76 if err != nil {
77 return err
78 }
79 return pager.Validate()
80 }
81
82 func (a *App) historySliceFromPager(ctx context.Context, pager *agent.DisplayPager, sessionDir, path string, req HistorySliceRequest, sourceID string) (HistorySlice, bool, error) {
83 src := historySourceFromPager(ctx, pager, path, sourceID)
84 page, err := a.pageHistorySliceSource(src, req, sessionDisplayResolver(sessionDir, path), sessionPlannerDisplayTurns(sessionDir, path), nil, path)
85 if src.readErr != nil {
86 return emptyHistorySlice(), true, src.readErr
87 }
88 page.Source = "index"
89 if pager.Built {
90 page.Source = "scan"
91 }
92 if pager.DAG || pager.SchemaOne {
93 page.Source = "event-log"
94 }
95 return page, true, err
96 }
97
98 func historySourceFromPager(ctx context.Context, pager *agent.DisplayPager, path, sourceID string) *historySliceSource {
99 pager = pager.WithContext(ctx)
100 idx := &pager.Header
101 src := &historySliceSource{sessionID: strings.TrimSuffix(filepath.Base(path), ".jsonl"), total: idx.MessageCount,
102 totalTurns: idx.AuthoredTurns, revision: idx.Revision, revKnown: idx.RevisionKnown, digest: idx.ContentDigest,
103 position: pager.Entry, usersBefore: pager.UsersBefore, maxFetch: 500, windowOnly: true}
104 src.sourceID = sourceID
105 if idx.RevisionKnown {
106 src.epoch = int(idx.Revision)
107 } else {
108 src.revision = 0
109 }
110 src.fetch = func(lo, hi int) ([]provider.Message, error) {
111 if err := pager.Validate(); err != nil {
112 return nil, err
113 }
114 if pager.DAG || pager.SchemaOne {
115 return pager.EventMessages(lo, hi)
116 }
117 entries, err := pager.Entries(lo, hi)
118 if err != nil {
119 return nil, err
120 }
121 messages, err := readSessionMessagesAtOffsetsContext(ctx, path, entries)
122 if err != nil {
123 return nil, err
124 }
125 return messages, pager.Validate()
126 }
127 src.windowBytes = func(lo, hi int) int64 {
128 if hi <= lo {
129 return 0
130 }
131 if pager.DAG || pager.SchemaOne {
132 entries, err := pager.Entries(lo, hi)
133 if err != nil {
134 src.readErr = err
135 return 0
136 }
137 var total int64
138 for _, entry := range entries {
139 total += entry.Length
140 }
141 return total
142 }
143 first, err := pager.Entry(lo)
144 if err != nil {
145 src.readErr = err
146 return 0
147 }
148 last, err := pager.Entry(hi - 1)
149 if err != nil {
150 src.readErr = err
151 return 0
152 }
153 return last.Offset + last.Length - first.Offset
154 }
155 return src
156 }
157
157 lines GO