返回 DeepSeek-Reasonix
session_display_pager_events.go
根目录 / internal / agent / session_display_pager_events.go
1 package agent
2
3 import (
4 "bufio"
5 "context"
6 "crypto/sha256"
7 "database/sql"
8 "encoding/json"
9 "errors"
10 "fmt"
11 "io"
12 "os"
13
14 "reasonix/internal/fileops"
15 "reasonix/internal/historywork"
16 "reasonix/internal/provider"
17 "reasonix/internal/store"
18 )
19
20 // Schema-1 replace/append logs retain original message locations. A large
21 // replace array is decoded one message at a time; SQLite holds the live order.
22 // No migration, normalization, or repair writes are made to source files.
23 func buildEventDisplayPager(ctx context.Context, db *sql.DB, source, fingerprint string, checkpointSize int64) error {
24 return buildEventDisplayPagerObserved(ctx, db, source, fingerprint, checkpointSize, nil)
25 }
26
27 func buildEventDisplayPagerObserved(ctx context.Context, db *sql.DB, source, fingerprint string, checkpointSize int64, observed func(string, int)) error {
28 f, err := fileops.OpenReplaceableRead(store.SessionEventLog(source))
29 if err != nil {
30 return err
31 }
32 defer f.Close()
33 if err := validateEventPagerFile(source, f, fingerprint); err != nil {
34 return err
35 }
36 info, err := f.Stat()
37 if err != nil {
38 return err
39 }
40 progress, err := restoreEventPagerScan(ctx, db, fingerprint, info.Size())
41 if err != nil {
42 return err
43 }
44 scanner := &eventPagerScanner{ctx: ctx, db: db, progress: progress, observed: observed}
45 if err := scanner.scan(f); err != nil {
46 return err
47 }
48 if err := finishEventDisplayPager(ctx, db, f, source, fingerprint, checkpointSize, scanner.progress.Count, observed); err != nil {
49 return err
50 }
51 return validateEventPagerFile(source, f, fingerprint)
52 }
53
54 // Unknown event fields are skipped token by token rather than buffering the
55 // entire field/record in json.RawMessage. json.Decoder still validates nesting.
56 func skipDisplayJSONValue(ctx context.Context, decoder *json.Decoder) error {
57 depth := 0
58 for {
59 if err := ctx.Err(); err != nil {
60 return err
61 }
62 token, err := decoder.Token()
63 if err != nil {
64 return err
65 }
66 if delim, ok := token.(json.Delim); ok {
67 switch delim {
68 case '{', '[':
69 depth++
70 case '}', ']':
71 depth--
72 }
73 }
74 if depth <= 0 {
75 return nil
76 }
77 }
78 }
79
80 type displayEventLocation struct {
81 Position int
82 Offset, Length int64
83 At int64
84 }
85
86 func eventDisplayLocations(ctx context.Context, db *sql.DB, lo, hi int) ([]displayEventLocation, error) {
87 rows, err := db.QueryContext(ctx, `SELECT position,offset,length,at FROM event_locations WHERE position>=? AND position<? ORDER BY position`, lo, hi)
88 if err != nil {
89 return nil, err
90 }
91 defer rows.Close()
92 locations := make([]displayEventLocation, 0, hi-lo)
93 for rows.Next() {
94 var loc displayEventLocation
95 if err := rows.Scan(&loc.Position, &loc.Offset, &loc.Length, &loc.At); err != nil {
96 return nil, err
97 }
98 locations = append(locations, loc)
99 }
100 return locations, rows.Err()
101 }
102
103 func readDisplayEventMessage(ctx context.Context, file *os.File, loc displayEventLocation) (provider.Message, error) {
104 reader := bufio.NewReaderSize(&historywork.Reader{Context: ctx, Source: io.NewSectionReader(file, loc.Offset, loc.Length)}, historywork.ReadChunk)
105 // InputOffset may precede the comma separating adjacent array elements.
106 for {
107 prefix, err := reader.Peek(1)
108 if err != nil {
109 return provider.Message{}, err
110 }
111 if prefix[0] != ',' && prefix[0] != ' ' && prefix[0] != '\n' && prefix[0] != '\r' && prefix[0] != '\t' {
112 break
113 }
114 _, _ = reader.ReadByte()
115 }
116 var message provider.Message
117 err := json.NewDecoder(reader).Decode(&message)
118 return message, err
119 }
120
121 func finishEventDisplayPager(ctx context.Context, db *sql.DB, file *os.File, source, fingerprint string, checkpointSize int64, count int, observed func(string, int)) error {
122 idx := SessionDisplayIndex{SchemaVersion: SessionDisplayIndexSchemaVersion, TranscriptSize: checkpointSize, ListingPreviewKnown: true}
123 hash := sha256.New()
124 users := 0
125 if err := restoreEventProjection(ctx, db, &idx, &users, hash, count); err != nil {
126 return err
127 }
128 for lo := idx.MessageCount; lo < count; lo += historywork.BatchEntries {
129 locations, err := eventDisplayLocations(ctx, db, lo, min(lo+historywork.BatchEntries, count))
130 if err != nil {
131 return err
132 }
133 tx, err := db.BeginTx(ctx, nil)
134 if err != nil {
135 return err
136 }
137 for _, loc := range locations {
138 message, readErr := readDisplayEventMessage(ctx, file, loc)
139 if readErr != nil {
140 err = readErr
141 break
142 }
143 var entry DisplayIndexEntry
144 entry, idx.AuthoredTurns = classifyDisplayIndexMessage(message, idx.MessageCount, loc.Offset, loc.Length, idx.AuthoredTurns)
145 body, _ := json.Marshal(entry)
146 if _, err = tx.ExecContext(ctx, `INSERT INTO entries VALUES(?,?,?,?,?,?,?)`, entry.Index, entry.Offset, entry.Length, entry.AuthoredTurn, entry.Role, users, body); err != nil {
147 break
148 }
149 body, err = json.Marshal(messageForSessionIdentity(message))
150 if err != nil {
151 break
152 }
153 hash.Write(body)
154 hash.Write([]byte{'\n'})
155 if entry.Role == provider.RoleUser && !entry.PinnedContextRevision {
156 users++
157 }
158 if entry.StartsTurn && idx.ListingPreview == "" {
159 idx.ListingPreview = truncatePreview(previewProse(UserMessageText(message)))
160 }
161 idx.MessageCount++
162 }
163 if err == nil {
164 err = saveEventProjection(ctx, tx, idx, users, hash)
165 }
166 if err != nil {
167 _ = tx.Rollback()
168 return err
169 }
170 if err := tx.Commit(); err != nil {
171 return err
172 }
173 if observed != nil {
174 observed("projection", idx.MessageCount)
175 }
176 }
177 if err := ctx.Err(); err != nil {
178 return err
179 }
180 idx.ContentDigest = fmt.Sprintf("%x", hash.Sum(nil))
181 identity, known, err := SessionContentIdentity(source)
182 if err != nil {
183 return err
184 }
185 if known {
186 if idx.ContentDigest != identity.DigestHex {
187 return ErrDisplaySourceChanged
188 }
189 idx.Revision, idx.RevisionKnown = identity.Revision, identity.RevisionKnown
190 }
191 body, err := json.Marshal(idx)
192 if err != nil {
193 return err
194 }
195 _, err = db.ExecContext(ctx, `INSERT OR REPLACE INTO metadata VALUES('source',?),('header',?),('kind','schema1'); DROP TABLE IF EXISTS event_pending`, fingerprint, string(body))
196 return err
197 }
198
199 // EventMessages reads a bounded range from either native event-log format.
200 func (p *DisplayPager) EventMessages(lo, hi int) ([]provider.Message, error) {
201 if p.DAG {
202 return p.DAGMessages(lo, hi)
203 }
204 if !p.SchemaOne || lo < 0 || hi < lo || hi-lo > 500 {
205 return nil, errors.New("invalid event display window")
206 }
207 if err := p.Validate(); err != nil {
208 return nil, err
209 }
210 locations, err := eventDisplayLocations(p.ctx, p.DB, lo, hi)
211 if err != nil {
212 return nil, err
213 }
214 file, err := fileops.OpenReplaceableRead(store.SessionEventLog(p.source))
215 if err != nil {
216 return nil, err
217 }
218 defer file.Close()
219 messages := make([]provider.Message, 0, len(locations))
220 for _, loc := range locations {
221 message, err := readDisplayEventMessage(p.ctx, file, loc)
222 if err != nil {
223 return nil, err
224 }
225 if message.CreatedAt <= 0 && loc.At != 0 {
226 message.CreatedAt = loc.At
227 }
228 messages = append(messages, message)
229 }
230 return messages, p.Validate()
231 }
232
232 lines GO