返回 DeepSeek-Reasonix
history_index.go
根目录 / internal / session / history_index.go
1 package session
2
3 import (
4 "context"
5 "database/sql"
6 "encoding/base64"
7 "encoding/json"
8 "errors"
9 "fmt"
10 "os"
11 "path/filepath"
12 "strings"
13
14 "reasonix/internal/agent"
15 "reasonix/internal/projectiondb"
16 "reasonix/internal/provider"
17 "reasonix/internal/sessioncontent"
18 )
19
20 func (q *Query) prepareHistoryIndex(ctx context.Context, ref SessionRef) (*FilesystemPersistence, string, error) {
21 if q == nil {
22 return nil, "", errors.New("session: nil query")
23 }
24 if err := ref.validate(q.hostID); err != nil {
25 return nil, "", err
26 }
27 filesystem, ok := q.persistence.(*FilesystemPersistence)
28 if !ok {
29 return nil, "", errors.New("session: history index requires filesystem persistence")
30 }
31 path := historyIndexPath(filesystem.Root, ref.SessionID)
32 lock := q.projectionLock("history", ref.SessionID)
33 lock.Lock()
34 err := ensureHistoryIndex(ctx, filesystem, ref.SessionID, path)
35 lock.Unlock()
36 if err != nil {
37 return nil, "", err
38 }
39 return filesystem, path, nil
40 }
41
42 func (q *Query) ReadContent(ctx context.Context, ref SessionRef, contentRef sessioncontent.Ref, offset, length int64) ([]byte, error) {
43 if q == nil {
44 return nil, errors.New("session: nil query")
45 }
46 if err := ref.validate(q.hostID); err != nil {
47 return nil, err
48 }
49 filesystem, ok := q.persistence.(*FilesystemPersistence)
50 if !ok {
51 return nil, errors.New("session: content reads require filesystem persistence")
52 }
53 if offset < 0 || length < 0 || length > 1<<20 || offset > contentRef.Bytes || length > contentRef.Bytes-offset {
54 return nil, errors.New("session: invalid or oversized content range")
55 }
56 if !q.contentAuthorized(ref.SessionID, contentRef.Digest, contentRef.Bytes, contentRef.IndexDigest) {
57 return nil, errors.New("session: content reference is not authorized for this session")
58 }
59 return contentStoreForSessionDir(filepath.Join(filesystem.Root, ref.SessionID)).ReadRange(ctx, contentRef, offset, length)
60 }
61
62 func ensureHistoryIndex(ctx context.Context, persistence *FilesystemPersistence, sessionID, path string) error {
63 dir := filepath.Join(persistence.Root, sessionID)
64 revision, err := revisionOfLog(dir)
65 if err != nil {
66 return err
67 }
68 if historyIndexCurrent(ctx, dir, path, sessionID, revision) {
69 return nil
70 }
71 if updated, err := incrementHistoryIndex(ctx, dir, path, sessionID, revision); updated || (err != nil && !errors.Is(err, ErrDamagedStore)) {
72 return err
73 }
74 // A bad derived checkpoint is not proof that the authoritative log is bad.
75 // Rebuild validates the complete log before replacing the old index; genuine
76 // corruption still fails and leaves the previous index intact.
77 return rebuildHistoryIndex(ctx, dir, path, sessionID, revision)
78 }
79
80 type historyIndexMetadata struct {
81 sessionID string
82 logSize int64
83 storageRevision int
84 projection int
85 durableSequence uint64
86 viewSequence uint64
87 generation string
88 }
89
90 func readHistoryIndexMetadata(ctx context.Context, db *sql.DB) (historyIndexMetadata, error) {
91 values := map[string]string{}
92 rows, err := db.QueryContext(ctx, `SELECT key,value FROM metadata WHERE key IN ('session_id','log_size','storage_revision','projection_version','durable_sequence','history_view_sequence','generation')`)
93 if err != nil {
94 return historyIndexMetadata{}, err
95 }
96 defer rows.Close()
97 for rows.Next() {
98 var key, value string
99 if err := rows.Scan(&key, &value); err != nil {
100 return historyIndexMetadata{}, err
101 }
102 values[key] = value
103 }
104 if err := rows.Err(); err != nil {
105 return historyIndexMetadata{}, err
106 }
107 metadata := historyIndexMetadata{sessionID: values["session_id"], generation: values["generation"]}
108 if _, err := fmt.Sscan(values["log_size"], &metadata.logSize); err != nil {
109 return historyIndexMetadata{}, err
110 }
111 if _, err := fmt.Sscan(values["storage_revision"], &metadata.storageRevision); err != nil {
112 return historyIndexMetadata{}, err
113 }
114 if _, err := fmt.Sscan(values["projection_version"], &metadata.projection); err != nil {
115 return historyIndexMetadata{}, err
116 }
117 if _, err := fmt.Sscan(values["durable_sequence"], &metadata.durableSequence); err != nil {
118 return historyIndexMetadata{}, err
119 }
120 if values["history_view_sequence"] != "" {
121 if _, err := fmt.Sscan(values["history_view_sequence"], &metadata.viewSequence); err != nil {
122 return historyIndexMetadata{}, err
123 }
124 }
125 return metadata, nil
126 }
127
128 // incrementHistoryIndex advances only the complete transactions appended after
129 // the published coverage watermark. It returns updated=false when the existing
130 // file cannot be trusted as a base and must be rebuilt atomically.
131 func incrementHistoryIndex(ctx context.Context, dir, path, sessionID string, revision logRevision) (updated bool, result error) {
132 if _, err := os.Stat(path); err != nil {
133 return false, nil
134 }
135 handle, err := projectiondb.Open(ctx, projectiondb.OpenOptions{Path: path, Migrations: historyMigrations, RequireDisk: true, MaxOpenConns: 1})
136 if err != nil {
137 return false, nil
138 }
139 defer handle.DB.Close()
140 metadata, err := readHistoryIndexMetadata(ctx, handle.DB)
141 if err != nil {
142 return false, nil
143 }
144 generation, err := historyProjectionGeneration(dir, 0)
145 if err != nil {
146 return false, err
147 }
148 if !metadata.canIncrement(sessionID, revision, generation) {
149 return false, nil
150 }
151
152 manifest, err := readManifest(filepath.Join(dir, "manifest.json"))
153 if err != nil || manifest.Codec != Codec {
154 return false, nil
155 }
156 log, err := os.Open(logPathForManifest(dir, manifest))
157 if err != nil {
158 return false, err
159 }
160 defer log.Close()
161
162 state, err := loadCurrentHistoryState(ctx, handle.DB)
163 if err != nil {
164 return false, nil
165 }
166
167 tx, err := handle.DB.BeginTx(ctx, nil)
168 if err != nil {
169 return false, err
170 }
171 state.tx = tx
172 state.statements, err = prepareHistoryBuildStatements(ctx, tx)
173 if err != nil {
174 _ = tx.Rollback()
175 return false, err
176 }
177 defer func() {
178 state.statements.close()
179 if tx != nil {
180 _ = tx.Rollback()
181 }
182 }()
183 content := contentStoreForSessionDir(dir)
184 viewSequence := metadata.viewSequence
185 var buildErr error
186 progress, err := scanHistoryLog(ctx, log, metadata.logSize, metadata.durableSequence+1, revision.Size, content, func(commit Commit) bool {
187 state.commitTurn, state.commitTime = commit.TurnID, commit.CreatedAt.UnixMilli()
188 state.transactions = append(state.transactions, []any{commit.ID, commit.FirstSequence, commit.LastSequence(), commit.OperationID, commit.OperationHash, commit.TurnID, commit.CreatedAt.UTC().Format("2006-01-02T15:04:05.999999999Z07:00")})
189 for _, event := range commit.Events {
190 if event.Kind == "history/replace" {
191 viewSequence = event.Sequence
192 }
193 digest := ""
194 var contentBytes int64
195 if event.PayloadRef != nil {
196 digest, contentBytes = event.PayloadRef.Digest, event.PayloadRef.Bytes
197 if err := insertContentRef(ctx, &state, *event.PayloadRef); err != nil {
198 buildErr = err
199 return false
200 }
201 }
202 state.events = append(state.events, []any{event.Sequence, commit.ID, event.ID, event.Kind, digest, contentBytes})
203 if err := indexMessageEvent(ctx, content, &state, event); err != nil {
204 buildErr = err
205 return false
206 }
207 }
208 return true
209 })
210 if err != nil || buildErr != nil {
211 return false, errors.Join(err, buildErr)
212 }
213 if err := validateHistoryLog(ctx, dir, log, generation); err != nil {
214 return false, err
215 }
216 if progress.end == metadata.logSize {
217 // An incomplete physical tail remains unpublished. A writer will preserve
218 // and repair it before the next append.
219 return true, nil
220 }
221 if err := flushHistoryBuildRows(ctx, tx, &state); err != nil {
222 return false, err
223 }
224 generation = strings.TrimSuffix(generation, ":0") + fmt.Sprintf(":%d", viewSequence)
225 values := map[string]string{
226 "log_size": fmt.Sprint(progress.end), "log_mtime_ns": fmt.Sprint(progress.modTimeNS),
227 "durable_sequence": fmt.Sprint(progress.sequence), "history_view_sequence": fmt.Sprint(viewSequence), "generation": generation,
228 }
229 for key, value := range values {
230 if _, err := tx.ExecContext(ctx, `INSERT INTO metadata(key,value) VALUES(?,?) ON CONFLICT(key) DO UPDATE SET value=excluded.value`, key, value); err != nil {
231 return false, err
232 }
233 }
234 state.statements.close()
235 state.statements = nil
236 if err := tx.Commit(); err != nil {
237 return false, err
238 }
239 tx = nil
240 _, _ = handle.DB.ExecContext(ctx, `PRAGMA shrink_memory`)
241 return true, nil
242 }
243
244 func loadCurrentHistoryState(ctx context.Context, db *sql.DB) (historyBuildState, error) {
245 state := historyBuildState{positions: map[string]int64{}, turns: map[string]int{}, versions: map[string]int{}}
246 // Versions belong to the stable message identity, including retired rows.
247 // A later rewrite can restore a removed message; its next version must not
248 // collide with the versions retained for older snapshots.
249 rows, err := db.QueryContext(ctx, `SELECT message_id,
250 COALESCE(MAX(CASE WHEN current=1 THEN position END),0),
251 COALESCE(MAX(CASE WHEN current=1 THEN visible_turn END),0),MAX(version)
252 FROM messages GROUP BY message_id`)
253 if err != nil {
254 return historyBuildState{}, err
255 }
256 for rows.Next() {
257 var id string
258 var position int64
259 var visibleTurn, version int
260 if err := rows.Scan(&id, &position, &visibleTurn, &version); err != nil {
261 _ = rows.Close()
262 return historyBuildState{}, err
263 }
264 state.versions[id] = version
265 if position == 0 {
266 continue
267 }
268 state.positions[id], state.turns[id] = position, visibleTurn
269 state.nextPosition = max(state.nextPosition, position)
270 state.visibleTurn = max(state.visibleTurn, visibleTurn)
271 }
272 return state, errors.Join(rows.Err(), rows.Close())
273 }
274
275 func historyIndexCurrent(ctx context.Context, dir, path, sessionID string, revision logRevision) bool {
276 if _, err := os.Stat(path); err != nil {
277 return false
278 }
279 handle, err := projectiondb.Open(ctx, projectiondb.OpenOptions{Path: path, Migrations: historyMigrations, RequireDisk: true, MaxOpenConns: 1})
280 if err != nil {
281 return false
282 }
283 defer handle.DB.Close()
284 values := map[string]string{}
285 rows, err := handle.DB.QueryContext(ctx, `SELECT key,value FROM metadata WHERE key IN ('session_id','log_size','log_mtime_ns','storage_revision','projection_version','generation')`)
286 if err != nil {
287 return false
288 }
289 defer rows.Close()
290 for rows.Next() {
291 var key, value string
292 if rows.Scan(&key, &value) != nil {
293 return false
294 }
295 values[key] = value
296 }
297 manifest, err := readStoredManifest(filepath.Join(dir, "manifest.json"))
298 if err != nil {
299 return false
300 }
301 identity, err := readStorageIdentity(dir, manifest)
302 if err != nil || !strings.HasPrefix(values["generation"], identity.Generation+":") {
303 return false
304 }
305 return values["session_id"] == sessionID && values["log_size"] == fmt.Sprint(revision.Size) && values["log_mtime_ns"] == fmt.Sprint(revision.ModTimeNS) && values["storage_revision"] == fmt.Sprint(StorageRevision) && values["projection_version"] == fmt.Sprint(historyIndexVersion)
306 }
307
308 func historyProjectionGeneration(dir string, viewSequence uint64) (string, error) {
309 manifest, err := readStoredManifest(filepath.Join(dir, "manifest.json"))
310 if err != nil {
311 return "", err
312 }
313 identity, err := readStorageIdentity(dir, manifest)
314 if err != nil {
315 return "", err
316 }
317 return fmt.Sprintf("%s:%d", identity.Generation, viewSequence), nil
318 }
319
320 func rebuildHistoryIndex(ctx context.Context, dir, path, sessionID string, revision logRevision) error {
321 manifest, err := readManifest(filepath.Join(dir, "manifest.json"))
322 if err != nil {
323 return err
324 }
325 log, err := os.Open(logPathForManifest(dir, manifest))
326 if err != nil {
327 return err
328 }
329 defer log.Close()
330 content := contentStoreForSessionDir(dir)
331 return projectiondb.Rebuild(ctx, projectiondb.OpenOptions{Path: path, Migrations: historyMigrations, RequireDisk: true, MaxOpenConns: 1, QuickCheck: true}, func(ctx context.Context, db *sql.DB) error {
332 return populateHistoryIndex(ctx, db, log, content, dir, sessionID, revision)
333 })
334 }
335
336 func populateHistoryIndex(ctx context.Context, db *sql.DB, log *os.File, content *sessioncontent.Store, dir, sessionID string, revision logRevision) error {
337 generation, err := historyProjectionGeneration(dir, 0)
338 if err != nil {
339 return err
340 }
341 // Keep SQLite's derived-data working set explicit. The history database
342 // may be many GiB, but neither its page cache nor temporary sort state
343 // belongs in the runtime's cumulative memory footprint.
344 // Rebuild writes an unpublished, disposable replacement beside the live
345 // index. Avoid WAL and durability work for that private file; Rebuild
346 // validates it before one atomic publish, and the event log remains the
347 // durable source if a crash leaves or corrupts the temporary database.
348 if err := configureHistoryRebuild(ctx, db); err != nil {
349 return err
350 }
351 var tx *sql.Tx
352 state := historyBuildState{positions: map[string]int64{}, turns: map[string]int{}, versions: map[string]int{}}
353 defer func() {
354 state.statements.close()
355 if tx != nil {
356 _ = tx.Rollback()
357 }
358 }()
359 beginChunk := func() error {
360 var err error
361 tx, err = db.BeginTx(ctx, nil)
362 if err != nil {
363 return err
364 }
365 state.tx = tx
366 state.statements, err = prepareHistoryBuildStatements(ctx, tx)
367 return err
368 }
369 if err := beginChunk(); err != nil {
370 return err
371 }
372 var viewSequence uint64
373 var buildErr error
374 chunkEvents := 0
375 var chunkBytes int64
376 commitChunk := func() error {
377 if err := flushHistoryBuildRows(ctx, tx, &state); err != nil {
378 return err
379 }
380 state.statements.close()
381 state.statements = nil
382 if err := tx.Commit(); err != nil {
383 return err
384 }
385 // modernc SQLite allocates its page cache on the Go heap. Release dirty
386 // pages after each bounded transaction so a multi-GiB derived index does
387 // not retain every completed chunk until the database closes.
388 if _, err := db.ExecContext(ctx, `PRAGMA shrink_memory`); err != nil {
389 return err
390 }
391 chunkEvents, chunkBytes = 0, 0
392 return beginChunk()
393 }
394 progress, err := scanHistoryLog(ctx, log, 0, 1, revision.Size, content, func(commit Commit) bool {
395 state.commitTurn, state.commitTime = commit.TurnID, commit.CreatedAt.UnixMilli()
396 state.transactions = append(state.transactions, []any{commit.ID, commit.FirstSequence, commit.LastSequence(), commit.OperationID, commit.OperationHash, commit.TurnID, commit.CreatedAt.UTC().Format("2006-01-02T15:04:05.999999999Z07:00")})
397 for _, event := range commit.Events {
398 if event.Kind == "history/replace" {
399 viewSequence = event.Sequence
400 }
401 digest := ""
402 var bytes int64
403 if event.PayloadRef != nil {
404 digest, bytes = event.PayloadRef.Digest, event.PayloadRef.Bytes
405 if err := insertContentRef(ctx, &state, *event.PayloadRef); err != nil {
406 buildErr = err
407 return false
408 }
409 }
410 state.events = append(state.events, []any{event.Sequence, commit.ID, event.ID, event.Kind, digest, bytes})
411 if err := indexMessageEvent(ctx, content, &state, event); err != nil {
412 buildErr = err
413 return false
414 }
415 chunkEvents++
416 chunkBytes += int64(len(event.Payload))
417 if event.PayloadRef != nil {
418 chunkBytes += min(event.PayloadRef.Bytes, int64(historyIndexTxnBytes))
419 }
420 if chunkEvents >= historyIndexTxnEvents || chunkBytes >= historyIndexTxnBytes {
421 if err := commitChunk(); err != nil {
422 buildErr = err
423 return false
424 }
425 }
426 }
427 return true
428 })
429 if err != nil {
430 return err
431 }
432 if buildErr != nil {
433 return buildErr
434 }
435 if err := flushHistoryBuildRows(ctx, tx, &state); err != nil {
436 return err
437 }
438 if err := validateHistoryLog(ctx, dir, log, generation); err != nil {
439 return err
440 }
441 generation = strings.TrimSuffix(generation, ":0") + fmt.Sprintf(":%d", viewSequence)
442 metadata := map[string]string{"session_id": sessionID, "log_size": fmt.Sprint(progress.end), "log_mtime_ns": fmt.Sprint(progress.modTimeNS), "storage_revision": fmt.Sprint(StorageRevision), "projection_version": fmt.Sprint(historyIndexVersion), "durable_sequence": fmt.Sprint(progress.sequence), "history_view_sequence": fmt.Sprint(viewSequence), "generation": generation}
443 for key, value := range metadata {
444 if _, err := tx.ExecContext(ctx, `INSERT INTO metadata(key,value) VALUES(?,?)`, key, value); err != nil {
445 return err
446 }
447 }
448 state.statements.close()
449 state.statements = nil
450 err = tx.Commit()
451 tx = nil
452 return err
453 }
454
455 func flushHistoryBuildRows(ctx context.Context, tx *sql.Tx, state *historyBuildState) error {
456 if err := insertHistoryRows(ctx, tx, `INSERT INTO transactions(commit_id,first_sequence,last_sequence,operation_id,operation_hash,turn_id,created_at) VALUES `, 7, state.transactions); err != nil {
457 return err
458 }
459 if err := insertHistoryRows(ctx, tx, `INSERT INTO events(sequence,commit_id,event_id,kind,payload_digest,payload_bytes) VALUES `, 6, state.events); err != nil {
460 return err
461 }
462 if err := insertHistoryRows(ctx, tx, `INSERT OR IGNORE INTO content_refs(digest,bytes,index_digest) VALUES `, 3, state.contentRefs); err != nil {
463 return err
464 }
465 if err := insertHistoryRows(ctx, tx, `INSERT INTO messages(message_id,version,position,event_sequence,valid_to,role,preview,inline,content_digest,content_bytes,content_index_digest,current,search_text,visible_turn,visible_user) VALUES `, 15, state.messages); err != nil {
466 return err
467 }
468 state.transactions = state.transactions[:0]
469 state.events = state.events[:0]
470 state.contentRefs = state.contentRefs[:0]
471 state.messages = state.messages[:0]
472 return nil
473 }
474
475 func insertHistoryRows(ctx context.Context, tx *sql.Tx, prefix string, columns int, rows [][]any) error {
476 const rowsPerStatement = 128
477 for start := 0; start < len(rows); start += rowsPerStatement {
478 end := min(start+rowsPerStatement, len(rows))
479 var query strings.Builder
480 query.WriteString(prefix)
481 args := make([]any, 0, (end-start)*columns)
482 for rowIndex, row := range rows[start:end] {
483 if len(row) != columns {
484 return errors.New("session: invalid history index row width")
485 }
486 if rowIndex > 0 {
487 query.WriteByte(',')
488 }
489 query.WriteByte('(')
490 for column := range columns {
491 if column > 0 {
492 query.WriteByte(',')
493 }
494 query.WriteByte('?')
495 }
496 query.WriteByte(')')
497 args = append(args, row...)
498 }
499 if _, err := tx.ExecContext(ctx, query.String(), args...); err != nil {
500 return err
501 }
502 }
503 return nil
504 }
505
506 func indexMessageEvent(ctx context.Context, content *sessioncontent.Store, state *historyBuildState, event Event) error {
507 if event.Kind == "diagnostic" {
508 return indexDisplayNotice(ctx, content, state, event)
509 }
510 if event.Kind == "submission/accepted" {
511 payload := event.Payload
512 if event.PayloadRef != nil {
513 var err error
514 payload, err = resolveContentPayload(ctx, content, *event.PayloadRef)
515 if err != nil {
516 return err
517 }
518 }
519 var receipt SubmissionReceipt
520 if err := json.Unmarshal(payload, &receipt); err != nil {
521 return err
522 }
523 _, err := state.tx.ExecContext(ctx, `INSERT INTO submissions(session_id,submission_id,message_id,sequence) VALUES(?,?,?,?) ON CONFLICT(session_id,submission_id) DO NOTHING`, receipt.SessionID, receipt.SubmissionID, receipt.MessageID, event.Sequence)
524 return err
525 }
526 if err := indexTurnEvent(ctx, content, state, event); err != nil {
527 return err
528 }
529 if event.Kind != "message/complete" && event.Kind != "message/upsert" && event.Kind != "message/retract" && event.Kind != "history/replace" && event.Kind != "legacy/import" {
530 return nil
531 }
532 payload := event.Payload
533 if event.PayloadRef != nil {
534 var err error
535 payload, err = resolveContentPayload(ctx, content, *event.PayloadRef)
536 if err != nil {
537 return err
538 }
539 }
540 switch event.Kind {
541 case "message/retract":
542 ids, err := retractedMessageIDs(event, payload)
543 if err != nil {
544 return err
545 }
546 if err := flushHistoryBuildRows(ctx, state.tx, state); err != nil {
547 return err
548 }
549 for _, id := range ids {
550 if _, err := state.statements.expire.ExecContext(ctx, event.Sequence, id); err != nil {
551 return err
552 }
553 delete(state.positions, id)
554 delete(state.turns, id)
555 }
556 return renumberVisibleTurns(ctx, state, event.Sequence)
557 case "message/complete", "message/upsert":
558 var body struct {
559 Message *provider.Message `json:"message"`
560 }
561 if err := strictPayload(payload, &body); err != nil || body.Message == nil {
562 return damagedPayload(event, err)
563 }
564 message := body.Message
565 if state.commitTurn != "" && message.Role == provider.RoleAssistant && !message.LocalOnly && (strings.TrimSpace(message.Content) != "" || strings.TrimSpace(message.RawContent) != "") {
566 if _, err := state.tx.ExecContext(ctx, `UPDATE turn_summaries SET final_message_id=?,ended_at=MAX(ended_at,started_at+?) WHERE turn_id=?`, message.ID, message.WorkDurationMs, state.commitTurn); err != nil {
567 return err
568 }
569 }
570 return indexOneMessage(ctx, content, state, *body.Message, event.Sequence, event.Kind == "message/upsert")
571 case "history/replace", "legacy/import":
572 messages, err := replacementEventMessages(event, payload)
573 if err != nil {
574 return err
575 }
576 return replaceIndexedMessages(ctx, content, state, messages, event.Sequence)
577 }
578 return nil
579 }
580
581 func replaceIndexedMessages(ctx context.Context, content *sessioncontent.Store, state *historyBuildState, messages []provider.Message, sequence uint64) error {
582 if err := flushHistoryBuildRows(ctx, state.tx, state); err != nil {
583 return err
584 }
585 if _, err := state.statements.clear.ExecContext(ctx, sequence); err != nil {
586 return err
587 }
588 state.nextPosition = 0
589 state.visibleTurn = 0
590 state.positions = map[string]int64{}
591 state.turns = map[string]int{}
592 // Keep each identity's version watermark when replacing its visible row.
593 // Retired versions remain in SQLite for fixed-snapshot readers.
594 legacyTurn := 0
595 var final *provider.Message
596 flushLegacyTurn := func() error {
597 if final == nil {
598 return nil
599 }
600 _, err := state.tx.ExecContext(ctx, `INSERT OR REPLACE INTO turn_summaries(turn_id,start_sequence,end_sequence,started_at,ended_at,final_message_id) VALUES(?,?,?,?,?,?)`, fmt.Sprintf("legacy:%d:%d", sequence, legacyTurn), sequence, sequence, 0, final.WorkDurationMs, final.ID)
601 return err
602 }
603 for _, message := range messages {
604 if message.Role == provider.RoleUser {
605 if err := flushLegacyTurn(); err != nil {
606 return err
607 }
608 legacyTurn++
609 final = nil
610 }
611 if message.Role == provider.RoleAssistant && !message.LocalOnly && (strings.TrimSpace(message.Content) != "" || strings.TrimSpace(message.RawContent) != "") {
612 copy := message
613 final = &copy
614 }
615 if err := indexOneMessage(ctx, content, state, message, sequence, false); err != nil {
616 return err
617 }
618 }
619 return flushLegacyTurn()
620 }
621
622 func renumberVisibleTurns(ctx context.Context, state *historyBuildState, sequence uint64) error {
623 rows, err := state.tx.QueryContext(ctx, `SELECT message_id,ordinal FROM (SELECT message_id,visible_turn,SUM(visible_user) OVER (ORDER BY position) AS ordinal FROM messages WHERE current=1) WHERE visible_turn<>ordinal`)
624 if err != nil {
625 return err
626 }
627 type changedTurn struct {
628 id string
629 turn int
630 }
631 var changed []changedTurn
632 for rows.Next() {
633 var item changedTurn
634 if err := rows.Scan(&item.id, &item.turn); err != nil {
635 _ = rows.Close()
636 return err
637 }
638 changed = append(changed, item)
639 }
640 if err := errors.Join(rows.Err(), rows.Close()); err != nil {
641 return err
642 }
643 for _, item := range changed {
644 version := state.versions[item.id] + 1
645 if _, err := state.statements.expire.ExecContext(ctx, sequence, item.id); err != nil {
646 return err
647 }
648 if _, err := state.tx.ExecContext(ctx, `INSERT INTO messages(message_id,version,position,event_sequence,valid_to,role,preview,inline,content_digest,content_bytes,content_index_digest,current,search_text,visible_turn,visible_user) SELECT message_id,?,position,?,0,role,preview,inline,content_digest,content_bytes,content_index_digest,1,search_text,?,visible_user FROM messages WHERE message_id=? ORDER BY version DESC LIMIT 1`, version, sequence, item.turn, item.id); err != nil {
649 return err
650 }
651 state.versions[item.id], state.turns[item.id] = version, item.turn
652 }
653 state.visibleTurn = 0
654 for _, turn := range state.turns {
655 state.visibleTurn = max(state.visibleTurn, turn)
656 }
657 return nil
658 }
659
660 func messageSearchText(message provider.Message) string {
661 parts := []string{message.Content, message.RawContent, message.ReasoningContent}
662 return strings.Join(parts, "\n")
663 }
664
665 func insertContentRef(ctx context.Context, state *historyBuildState, ref sessioncontent.Ref) error {
666 if err := ctx.Err(); err != nil {
667 return err
668 }
669 state.contentRefs = append(state.contentRefs, []any{ref.Digest, ref.Bytes, ref.IndexDigest})
670 return nil
671 }
672
673 func messagePreview(message provider.Message) string {
674 // Reference-only history rows have no origin field until hydrated. Never
675 // publish host protocol text in that temporary user-message preview.
676 if agent.IsHostGeneratedUserMessage(message) {
677 return ""
678 }
679 preview := strings.TrimSpace(message.Content)
680 if message.Role == provider.RoleUser {
681 // Derive display text before truncating; a truncated injected block can
682 // no longer be separated from the user's request.
683 preview = agent.UserMessageText(message)
684 if message.Origin == "" && strings.TrimSpace(message.RawContent) == "" {
685 if body, ok := strings.CutPrefix(preview, `<session-context version="1">`); ok && strings.HasPrefix(strings.TrimSpace(body), "This host-generated snapshot supersedes every earlier session-context snapshot.") {
686 return ""
687 }
688 }
689 } else if preview == "" {
690 preview = strings.TrimSpace(message.RawContent)
691 }
692 runes := []rune(preview)
693 if len(runes) > 240 {
694 preview = string(runes[:240])
695 }
696 return preview
697 }
698
699 func encodeHistoryCursor(cursor historyCursor) (string, error) {
700 data, err := json.Marshal(cursor)
701 if err != nil {
702 return "", err
703 }
704 return base64.RawURLEncoding.EncodeToString(data), nil
705 }
706
707 func decodeHistoryCursor(value string) (historyCursor, error) {
708 data, err := base64.RawURLEncoding.DecodeString(value)
709 if err != nil {
710 return historyCursor{}, errors.New("session: invalid history cursor")
711 }
712 var cursor historyCursor
713 if err := json.Unmarshal(data, &cursor); err != nil {
714 return historyCursor{}, errors.New("session: invalid history cursor")
715 }
716 return cursor, nil
717 }
718
719 func encodeSearchHistoryCursor(cursor searchHistoryCursor) (string, error) {
720 data, err := json.Marshal(cursor)
721 if err != nil {
722 return "", err
723 }
724 return base64.RawURLEncoding.EncodeToString(data), nil
725 }
726
727 func decodeSearchHistoryCursor(value string) (searchHistoryCursor, error) {
728 data, err := base64.RawURLEncoding.DecodeString(value)
729 if err != nil {
730 return searchHistoryCursor{}, errors.New("session: invalid history search cursor")
731 }
732 var cursor searchHistoryCursor
733 if err := json.Unmarshal(data, &cursor); err != nil {
734 return searchHistoryCursor{}, errors.New("session: invalid history search cursor")
735 }
736 return cursor, nil
737 }
738
739 func scanMetadataUint(row *sql.Row, target *uint64) error {
740 var value string
741 if err := row.Scan(&value); err != nil {
742 return err
743 }
744 _, err := fmt.Sscan(value, target)
745 return err
746 }
747
747 lines GO