返回 DeepSeek-Reasonix
session_display_pager_dag.go
根目录 / internal / agent / session_display_pager_dag.go
1 package agent
2
3 import (
4 "context"
5 "database/sql"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "hash"
10 "io"
11 "os"
12 "time"
13
14 "reasonix/internal/fileops"
15 "reasonix/internal/historywork"
16 "reasonix/internal/provider"
17 "reasonix/internal/store"
18 )
19
20 // DAG projections retain graph edges and source locations on disk. Bodies are
21 // decoded one entry at a time, never accumulated in the runtime message array.
22 // The authoritative graph and selected head are never modified by this reader.
23 type displayDAGLocation struct {
24 Offset int64 `json:"offset"`
25 Length int64 `json:"length"`
26 ID string `json:"id"`
27 Target string `json:"target,omitempty"`
28 At time.Time `json:"at"`
29 }
30
31 func readDisplayDAGMessage(ctx context.Context, file *os.File, loc displayDAGLocation) (provider.Message, error) {
32 var entry sessionDAGEntry
33 reader := &historywork.Reader{Context: ctx, Source: io.NewSectionReader(file, loc.Offset, loc.Length)}
34 if err := json.NewDecoder(reader).Decode(&entry); err != nil {
35 return provider.Message{}, err
36 }
37 raw := entry.Msgs
38 if loc.Target != "" {
39 raw = entry.Targets[loc.Target]
40 }
41 var messages []provider.Message
42 if err := json.Unmarshal(raw, &messages); err != nil {
43 return provider.Message{}, err
44 }
45 if len(messages) != 1 {
46 return provider.Message{}, ErrSessionDisplayReadModelDamaged
47 }
48 m := messages[0]
49 if loc.ID != "" {
50 m.ID = loc.ID
51 }
52 return m, ctx.Err()
53 }
54
55 func buildDAGDisplayPager(ctx context.Context, db *sql.DB, source, fingerprint, requestedHead string, checkpointSize int64) error {
56 return buildDAGDisplayPagerObserved(ctx, db, source, fingerprint, requestedHead, checkpointSize, nil)
57 }
58
59 func buildDAGDisplayPagerObserved(ctx context.Context, db *sql.DB, source, fingerprint, requestedHead string, checkpointSize int64, observed func(string, int)) error {
60 f, err := fileops.OpenReplaceableRead(store.SessionEventLog(source))
61 if err != nil {
62 return err
63 }
64 defer f.Close()
65 if err := validateDAGPagerFile(source, f, fingerprint, requestedHead); err != nil {
66 return err
67 }
68 p, err := restoreDAGPagerProgress(ctx, db, source, fingerprint, checkpointSize)
69 if err != nil {
70 return err
71 }
72 idx, headID, err := resumeDAGDisplayPager(ctx, db, f, p, source, requestedHead, observed)
73 if err != nil {
74 return err
75 }
76 identity, known, err := SessionContentIdentity(source)
77 if err != nil {
78 return err
79 }
80 if known && requestedHead == "" {
81 if idx.ContentDigest != identity.DigestHex {
82 return fmt.Errorf("DAG selected history identity changed")
83 }
84 idx.Revision, idx.RevisionKnown = identity.Revision, identity.RevisionKnown
85 }
86 body, err := json.Marshal(idx)
87 if err != nil {
88 return err
89 }
90 if _, err = db.ExecContext(ctx, `INSERT OR REPLACE INTO metadata VALUES('source',?),('header',?),('kind','dag'),('head',?)`, fingerprint, body, headID); err != nil {
91 return err
92 }
93 return validateDAGPagerFile(source, f, fingerprint, requestedHead)
94 }
95
96 // Both ordinary reads and private rebuild transactions use this interface.
97 type displayDAGQueries interface {
98 ExecContext(context.Context, string, ...any) (sql.Result, error)
99 QueryRowContext(context.Context, string, ...any) *sql.Row
100 }
101
102 func storeDisplayDAGOverlay(ctx context.Context, db displayDAGQueries, f *os.File, kind, id string, loc displayDAGLocation) error {
103 if _, err := readDisplayDAGMessage(ctx, f, loc); err != nil {
104 return err
105 }
106 body, err := json.Marshal(loc)
107 if err != nil {
108 return err
109 }
110 _, err = db.ExecContext(ctx, `INSERT OR REPLACE INTO dag_overlays VALUES(?,?,?)`, kind, id, body)
111 return err
112 }
113
114 func displayDAGSystem(ctx context.Context, db displayDAGQueries, key string) (displayDAGLocation, error) {
115 var loc displayDAGLocation
116 var raw []byte
117 if err := db.QueryRowContext(ctx, `SELECT location FROM dag_overlays WHERE kind='system' AND id=?`, key).Scan(&raw); err != nil {
118 return loc, err
119 }
120 err := json.Unmarshal(raw, &loc)
121 return loc, err
122 }
123
124 func (p *DisplayPager) DAGMessages(lo, hi int) ([]provider.Message, error) {
125 if lo < 0 || hi < lo || hi-lo > 500 {
126 return nil, fmt.Errorf("display window exceeds page budget")
127 }
128 rows, err := p.DB.QueryContext(p.ctx, `SELECT location FROM dag_locations WHERE position>=? AND position<? ORDER BY position`, lo, hi)
129 if err != nil {
130 return nil, err
131 }
132 var locations []displayDAGLocation
133 for rows.Next() {
134 var body []byte
135 var loc displayDAGLocation
136 if err = rows.Scan(&body); err == nil {
137 err = json.Unmarshal(body, &loc)
138 }
139 if err != nil {
140 break
141 }
142 locations = append(locations, loc)
143 }
144 err = errors.Join(err, rows.Err(), rows.Close())
145 if err != nil {
146 return nil, err
147 }
148 f, err := fileops.OpenReplaceableRead(store.SessionEventLog(p.source))
149 if err != nil {
150 return nil, err
151 }
152 defer f.Close()
153 result := make([]provider.Message, 0, len(locations))
154 for _, loc := range locations {
155 m, err := readDisplayDAGMessage(p.ctx, f, loc)
156 if err != nil {
157 return nil, err
158 }
159 if m.CreatedAt <= 0 && !loc.At.IsZero() {
160 m.CreatedAt = loc.At.UnixMilli()
161 }
162 result = append(result, m)
163 }
164 return result, p.Validate()
165 }
166
167 func applyDisplayDAGLocation(ctx context.Context, db displayDAGQueries, f *os.File, state *sessionDAGState, e sessionDAGEntry, start, end int64) error {
168 loc := displayDAGLocation{Offset: start, Length: end - start, ID: e.ID, At: e.At}
169 switch e.Type {
170 case sessionDAGTypeMessage:
171 if e.ID == "" {
172 return ErrSessionDisplayReadModelDamaged
173 }
174 // Validate each record, including branches outside the selected view.
175 if _, err := readDisplayDAGMessage(ctx, f, loc); err != nil {
176 return err
177 }
178 encoded, err := json.Marshal(loc)
179 if err != nil {
180 return err
181 }
182 res, err := db.ExecContext(ctx, `INSERT OR IGNORE INTO dag_nodes VALUES(?,?,?)`, e.ID, e.Parent, encoded)
183 if err != nil {
184 return err
185 }
186 if count, _ := res.RowsAffected(); count == 0 {
187 return nil
188 }
189 h := state.headFor(e.Head, e.At)
190 h.leaf, h.lastActivity, h.lastOffset = e.ID, e.At, end
191 case sessionDAGTypePatch, sessionDAGTypeSystem, sessionDAGTypeRedact:
192 if e.Type == sessionDAGTypePatch {
193 var exists int
194 if err := db.QueryRowContext(ctx, `SELECT 1 FROM dag_nodes WHERE id=?`, e.Target).Scan(&exists); errors.Is(err, sql.ErrNoRows) {
195 return nil
196 } else if err != nil {
197 return err
198 }
199 loc.ID = e.Target
200 }
201 if e.Type == sessionDAGTypeSystem {
202 // Fork copies this immutable locator key, matching the native
203 // replay's inherited system override without retaining its body.
204 loc.ID = ""
205 key := fmt.Sprint(start)
206 state.headFor(e.Head, e.At).system = &provider.Message{ID: key}
207 if err := storeDisplayDAGOverlay(ctx, db, f, "system", key, loc); err != nil {
208 return err
209 }
210 } else if e.Type == sessionDAGTypeRedact {
211 for id := range e.Targets {
212 loc.ID, loc.Target = id, id
213 if err := storeDisplayDAGOverlay(ctx, db, f, "redact", id, loc); err != nil {
214 return err
215 }
216 }
217 } else if err := storeDisplayDAGOverlay(ctx, db, f, "patch", e.Target, loc); err != nil {
218 return err
219 }
220 case sessionDAGTypeFork, sessionDAGTypeRewind, sessionDAGTypeSelect, sessionDAGTypeRename, sessionDAGTypeRetire,
221 sessionDAGTypeTurnBegin, sessionDAGTypeTurnEnd, sessionDAGTypeCompaction:
222 if !state.applyHeadMarker(e, end) {
223 return ErrSessionDisplayReadModelDamaged
224 }
225 case sessionDAGTypeLog:
226 state.generation = e.Generation
227 case sessionDAGTypeWriter, sessionDAGTypeCheckpoint:
228 default:
229 return fmt.Errorf("unsupported DAG entry type %q", e.Type)
230 }
231 return nil
232 }
233
234 type displayDAGProjectionWriter struct {
235 ctx context.Context
236 db displayDAGQueries
237 file *os.File
238 index SessionDisplayIndex
239 hasher hash.Hash
240 users int
241 }
242
243 func (w *displayDAGProjectionWriter) append(loc displayDAGLocation) error {
244 ctx, db, f := w.ctx, w.db, w.file
245 idx, hasher := &w.index, w.hasher
246 m, err := readDisplayDAGMessage(ctx, f, loc)
247 if err != nil {
248 return err
249 }
250 entry, turn := classifyDisplayIndexMessage(m, idx.MessageCount, loc.Offset, loc.Length, idx.AuthoredTurns)
251 idx.AuthoredTurns = turn
252 entryJSON, err := json.Marshal(entry)
253 if err != nil {
254 return err
255 }
256 locationJSON, err := json.Marshal(loc)
257 if err != nil {
258 return err
259 }
260 if _, err := db.ExecContext(ctx, `INSERT INTO entries VALUES(?,?,?,?,?,?,?)`, entry.Index, entry.Offset, entry.Length, entry.AuthoredTurn, entry.Role, w.users, entryJSON); err != nil {
261 return err
262 }
263 if _, err := db.ExecContext(ctx, `INSERT INTO dag_locations VALUES(?,?)`, entry.Index, locationJSON); err != nil {
264 return err
265 }
266 body, err := json.Marshal(messageForSessionIdentity(m))
267 if err != nil {
268 return err
269 }
270 hasher.Write(body)
271 hasher.Write([]byte{'\n'})
272 if m.Role == provider.RoleUser && !IsPinnedContextRevision(m) {
273 w.users++
274 }
275 if entry.StartsTurn && idx.ListingPreview == "" {
276 idx.ListingPreview = truncatePreview(previewProse(UserMessageText(m)))
277 }
278 idx.MessageCount++
279 return nil
280 }
281
281 lines GO