返回 DeepSeek-Reasonix
session_display_pager.go
根目录 / internal / agent / session_display_pager.go
1 package agent
2
3 import (
4 "context"
5 "database/sql"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "io"
10 "os"
11
12 "reasonix/internal/fileops"
13 "reasonix/internal/historywork"
14 "reasonix/internal/projectiondb"
15 "reasonix/internal/provider"
16 "reasonix/internal/store"
17 )
18
19 // DisplayPager is a disposable disk projection of exact legacy byte offsets.
20 // It never converts or rewrites the source transcript. Header contains no
21 // Entries; individual positions are queried through SQLite's bounded cache.
22 type DisplayPager struct {
23 DB *sql.DB
24 Header SessionDisplayIndex
25 Built bool
26 DAG bool
27 SchemaOne bool
28 ctx context.Context
29 source string
30 sourceInfo os.FileInfo
31 eventInfo os.FileInfo
32 sourceVersion fileops.Version
33 eventVersion fileops.Version
34 }
35
36 var ErrDisplaySourceChanged = errors.New("display source changed")
37 var ErrDisplayFormatUnsupported = errors.New("display format requires compatibility reader")
38
39 var displayPagerMigrations = []projectiondb.Migration{{Version: 1, Apply: func(ctx context.Context, tx *sql.Tx) error {
40 _, err := tx.ExecContext(ctx, `CREATE TABLE metadata(key TEXT PRIMARY KEY,value TEXT NOT NULL);
41 CREATE TABLE entries(position INTEGER PRIMARY KEY,offset INTEGER NOT NULL,length INTEGER NOT NULL,turn INTEGER NOT NULL,role TEXT NOT NULL,user_before INTEGER NOT NULL,entry BLOB NOT NULL);
42 CREATE INDEX entries_turn ON entries(turn,position);`)
43 return err
44 }}}
45
46 func OpenDisplayPager(ctx context.Context, source, cachePath string, heads ...string) (*DisplayPager, error) {
47 head := ""
48 if len(heads) > 0 {
49 head = heads[0]
50 }
51 return openDisplayPager(ctx, source, cachePath, head, false)
52 }
53
54 func openDisplayPager(ctx context.Context, source, cachePath, head string, forceSource bool) (*DisplayPager, error) {
55 if err := ctx.Err(); err != nil {
56 return nil, err
57 }
58 info, err := os.Stat(source)
59 if err != nil {
60 return nil, err
61 }
62 identity, known, err := SessionContentIdentity(source)
63 if err != nil {
64 return nil, err
65 }
66 indexPath := store.SessionDisplayIndex(source)
67 indexInfo, err := os.Stat(indexPath)
68 plain := os.IsNotExist(err)
69 if err != nil && !plain {
70 return nil, err
71 }
72 // Equal timestamps are ambiguous. Leave them to the explicit authoritative
73 // preparation path, never certify stale offsets from size alone.
74 if !plain && (forceSource || !known || !indexInfo.ModTime().After(SessionContentModTime(source))) {
75 plain = true
76 }
77 sourceTarget, sourceVersion := fileops.DiskSnapshot(source, info)
78 fingerprint := fmt.Sprintf("%s:%s", sourceTarget.Key, sourceVersion)
79 opts := projectiondb.OpenOptions{Path: cachePath, Migrations: displayPagerMigrations, RequireDisk: true, MaxOpenConns: 1}
80 handle, err := projectiondb.Open(ctx, opts)
81 if err != nil {
82 return nil, err
83 }
84 defer func() {
85 if handle != nil {
86 _ = handle.DB.Close()
87 }
88 }()
89 var stored string
90 storedErr := handle.DB.QueryRowContext(ctx, `SELECT value FROM metadata WHERE key='source'`).Scan(&stored)
91 observed, err := observeDisplayPagerEvents(ctx, source, head, forceSource, plain, indexInfo, fingerprint, stored, storedErr)
92 if err != nil {
93 return nil, err
94 }
95 fingerprint, plain = observed.fingerprint, observed.plain
96 eventInfo, eventVersion, dag, schemaOne := observed.eventInfo, observed.eventVersion, observed.dag, observed.schemaOne
97 built := storedErr != nil || stored != fingerprint
98 if built {
99 _ = handle.DB.Close()
100 err = observed.rebuild(ctx, opts, source, indexPath, head, info.Size())
101 if err != nil {
102 if !dag && !schemaOne && !plain && ctx.Err() == nil {
103 return openDisplayPager(ctx, source, cachePath, head, true)
104 }
105 return nil, err
106 }
107 handle, err = projectiondb.Open(ctx, opts)
108 if err != nil {
109 return nil, err
110 }
111 }
112 p := &DisplayPager{DB: handle.DB, ctx: ctx, source: source, sourceInfo: info, eventInfo: eventInfo, sourceVersion: sourceVersion, eventVersion: eventVersion, DAG: dag, SchemaOne: schemaOne, Built: built && (plain || dag || schemaOne)}
113 matching, err := p.loadMatchingHeader(ctx, info, identity, known, head)
114 if err != nil {
115 p.Close()
116 return nil, err
117 }
118 if !matching {
119 p.Close()
120 if !dag && !schemaOne && !plain {
121 return openDisplayPager(ctx, source, cachePath, head, true)
122 }
123 return nil, ErrDisplaySourceChanged
124 }
125
126 if err := p.prepareQueries(ctx); err != nil {
127 p.Close()
128 return nil, err
129 }
130 handle = nil // The returned pager owns the open database from here.
131 return p, nil
132 }
133
134 func (p *DisplayPager) prepareQueries(ctx context.Context) error {
135 if err := p.Validate(); err != nil {
136 return err
137 }
138 // Optional accelerator over disposable metadata. Existing cache readers
139 // can ignore it; no authoritative format or schema version changes.
140 _, err := p.DB.ExecContext(ctx, `CREATE INDEX IF NOT EXISTS entries_authored_turn ON entries(turn,position) WHERE json_extract(CAST(entry AS TEXT),'$.starts_turn')=1`)
141 return err
142 }
143
144 func (p *DisplayPager) Close() error { return p.DB.Close() }
145
146 // WithContext borrows the projection for a single caller. Only its preparation
147 // owner closes the database; canceling a read cannot cancel other readers.
148 func (p *DisplayPager) WithContext(ctx context.Context) *DisplayPager {
149 copy := *p
150 copy.ctx = ctx
151 return &copy
152 }
153
154 func (p *DisplayPager) Validate() error {
155 if err := p.ctx.Err(); err != nil {
156 return err
157 }
158 st, err := os.Stat(p.source)
159 if err != nil {
160 return err
161 }
162 if !os.SameFile(st, p.sourceInfo) || st.Size() != p.sourceInfo.Size() || !st.ModTime().Equal(p.sourceInfo.ModTime()) {
163 return ErrDisplaySourceChanged
164 }
165 if _, version := fileops.DiskSnapshot(p.source, st); version != p.sourceVersion {
166 return ErrDisplaySourceChanged
167 }
168 if p.eventInfo != nil {
169 current, err := os.Stat(store.SessionEventLog(p.source))
170 if err != nil {
171 return err
172 }
173 if !os.SameFile(current, p.eventInfo) || current.Size() != p.eventInfo.Size() || !current.ModTime().Equal(p.eventInfo.ModTime()) {
174 return ErrDisplaySourceChanged
175 }
176 if _, version := fileops.DiskSnapshot(store.SessionEventLog(p.source), current); version != p.eventVersion {
177 return ErrDisplaySourceChanged
178 }
179 } else {
180 current, err := os.Stat(store.SessionEventLog(p.source))
181 if err == nil && current.Size() > 0 {
182 return ErrDisplaySourceChanged
183 }
184 if err != nil && !os.IsNotExist(err) {
185 return err
186 }
187 }
188 return nil
189 }
190
191 func (p *DisplayPager) Entry(position int) (DisplayIndexEntry, error) {
192 var entry DisplayIndexEntry
193 var body []byte
194 if err := p.DB.QueryRowContext(p.ctx, `SELECT entry FROM entries WHERE position=?`, position).Scan(&body); err != nil {
195 return entry, err
196 }
197 err := json.Unmarshal(body, &entry)
198 return entry, err
199 }
200
201 func (p *DisplayPager) Entries(lo, hi int) ([]DisplayIndexEntry, error) {
202 if lo < 0 || hi < lo || hi-lo > 500 {
203 return nil, errors.New("display window exceeds page budget")
204 }
205 rows, err := p.DB.QueryContext(p.ctx, `SELECT entry FROM entries WHERE position>=? AND position<? ORDER BY position`, lo, hi)
206 if err != nil {
207 return nil, err
208 }
209 defer rows.Close()
210 out := make([]DisplayIndexEntry, 0, hi-lo)
211 for rows.Next() {
212 var b []byte
213 if err := rows.Scan(&b); err != nil {
214 return nil, err
215 }
216 var e DisplayIndexEntry
217 if err := json.Unmarshal(b, &e); err != nil {
218 return nil, err
219 }
220 out = append(out, e)
221 }
222 return out, rows.Err()
223 }
224
225 func (p *DisplayPager) UsersBefore(position int) (int, error) {
226 var count int
227 err := p.DB.QueryRowContext(p.ctx, `SELECT user_before FROM entries WHERE position=?`, position).Scan(&count)
228 return count, err
229 }
230
231 // TurnEntries returns only authored user positions in the selected branch.
232 // Hidden control messages and long tool runs do not require prefix decoding.
233 func (p *DisplayPager) TurnEntries(start, limit int) ([]DisplayIndexEntry, error) {
234 if limit <= 0 || limit > 1000 {
235 return nil, errors.New("display outline exceeds page budget")
236 }
237 rows, err := p.DB.QueryContext(p.ctx, `SELECT entry FROM entries WHERE turn>=? AND json_extract(CAST(entry AS TEXT),'$.starts_turn')=1 ORDER BY turn,position LIMIT ?`, max(start, 1), limit)
238 if err != nil {
239 return nil, err
240 }
241 defer rows.Close()
242 entries := make([]DisplayIndexEntry, 0, limit)
243 for rows.Next() {
244 var body []byte
245 var entry DisplayIndexEntry
246 if err := rows.Scan(&body); err != nil {
247 return nil, err
248 }
249 if err := json.Unmarshal(body, &entry); err != nil {
250 return nil, err
251 }
252 entries = append(entries, entry)
253 }
254 return entries, rows.Err()
255 }
256
257 // Import the existing JSON index one entry at a time. A corrupt or cancelled
258 // build never replaces the last published projection.
259 func importDisplayPager(ctx context.Context, db *sql.DB, indexPath, fingerprint string) error {
260 return importDisplayPagerObserved(ctx, db, indexPath, fingerprint, nil)
261 }
262
263 func importDisplayPagerObserved(ctx context.Context, db *sql.DB, indexPath, fingerprint string, committed func(int)) error {
264 f, err := os.Open(indexPath)
265 if err != nil {
266 return err
267 }
268 defer f.Close()
269 key, err := displayImportFileKey(indexPath, f)
270 if err != nil {
271 return err
272 }
273 info, err := f.Stat()
274 if err != nil {
275 return err
276 }
277 progress, err := restoreDisplayImportProgress(ctx, db, key, info.Size())
278 if err != nil {
279 return err
280 }
281 decoder, base, err := displayImportDecoder(ctx, f, progress.sourceOffset)
282 if err != nil {
283 return err
284 }
285 progress.decoderBase, progress.committed = base, committed
286 token, err := decoder.Token()
287 if err != nil || token != json.Delim('{') {
288 return errors.New("invalid display index object")
289 }
290 header := progress.header
291 sawEntries := false
292 for decoder.More() {
293 key, err := decoder.Token()
294 if err != nil {
295 return err
296 }
297 name, ok := key.(string)
298 if !ok {
299 return errors.New("invalid display index field")
300 }
301 if name != "entries" {
302 var raw json.RawMessage
303 if err := decoder.Decode(&raw); err != nil {
304 return err
305 }
306 header[name] = raw
307 continue
308 }
309 if sawEntries {
310 return errors.New("duplicate display entries")
311 }
312 sawEntries = true
313 token, err := decoder.Token()
314 if err != nil || token != json.Delim('[') {
315 return errors.New("invalid display entries")
316 }
317 if err := progress.readEntries(ctx, db, decoder); err != nil {
318 return err
319 }
320 if _, err := decoder.Token(); err != nil {
321 return err
322 }
323 }
324 if _, err := decoder.Token(); err != nil {
325 return err
326 }
327 var extra any
328 if err := decoder.Decode(&extra); !errors.Is(err, io.EOF) {
329 return errors.Join(errors.New("trailing display index data"), err)
330 }
331 body, err := json.Marshal(header)
332 if err != nil {
333 return err
334 }
335 var idx SessionDisplayIndex
336 if err := json.Unmarshal(body, &idx); err != nil {
337 return err
338 }
339 if !sawEntries || idx.SchemaVersion != SessionDisplayIndexSchemaVersion || idx.MessageCount != progress.count || idx.TranscriptSize != progress.offset {
340 return errors.New("display index header does not cover entries")
341 }
342 current, err := displayImportFileKey(indexPath, f)
343 if err != nil {
344 return err
345 }
346 if current != key {
347 return ErrDisplaySourceChanged
348 }
349 _, err = db.ExecContext(ctx, `INSERT OR REPLACE INTO metadata VALUES('source',?),('header',?)`, fingerprint, string(body))
350 return err
351 }
352
353 type displayIndexImportProgress struct {
354 count, users int
355 offset int64
356 sourceOffset int64
357 decoderBase int64
358 sourceKey string
359 header map[string]json.RawMessage
360 committed func(int)
361 }
362
363 func (progress *displayIndexImportProgress) readEntries(ctx context.Context, db *sql.DB, decoder *json.Decoder) error {
364 for decoder.More() {
365 tx, err := db.BeginTx(ctx, nil)
366 if err != nil {
367 return err
368 }
369 for batch := 0; batch < historywork.BatchEntries && decoder.More(); batch++ {
370 var e DisplayIndexEntry
371 if err = decoder.Decode(&e); err != nil {
372 break
373 }
374 if e.Index != progress.count || e.Offset != progress.offset || e.Length <= 0 {
375 err = errors.New("invalid display offset chain")
376 break
377 }
378 body, marshalErr := json.Marshal(e)
379 if marshalErr != nil {
380 err = marshalErr
381 break
382 }
383 _, err = tx.ExecContext(ctx, `INSERT INTO entries VALUES(?,?,?,?,?,?,?)`, progress.count, e.Offset, e.Length, e.AuthoredTurn, e.Role, progress.users, body)
384 if err != nil {
385 break
386 }
387 if e.Role == provider.RoleUser && !e.PinnedContextRevision {
388 progress.users++
389 }
390 progress.count++
391 progress.offset += e.Length
392 }
393 if err != nil {
394 _ = tx.Rollback()
395 return err
396 }
397 progress.sourceOffset = decoder.InputOffset() + progress.decoderBase
398 if err := progress.save(ctx, tx); err != nil {
399 _ = tx.Rollback()
400 return err
401 }
402 if err := tx.Commit(); err != nil {
403 return err
404 }
405 if progress.committed != nil {
406 progress.committed(progress.count)
407 }
408 }
409 return nil
410 }
411
411 lines GO