返回 DeepSeek-Reasonix
session_display_pager_checkpoint.go
根目录 / internal / agent / session_display_pager_checkpoint.go
1 package agent
2
3 import (
4 "bufio"
5 "context"
6 "crypto/sha256"
7 "database/sql"
8 "encoding"
9 "encoding/json"
10 "errors"
11 "fmt"
12 "io"
13 "os"
14
15 "reasonix/internal/fileops"
16 "reasonix/internal/historywork"
17 "reasonix/internal/provider"
18 )
19
20 // Build an offset-only projection for a checkpoint without publishing or
21 // repairing any session sidecar. At most one message and one SQLite batch are
22 // retained. Digest validation happens before the disposable index is published.
23 func buildCheckpointDisplayPager(ctx context.Context, db *sql.DB, source, fingerprint string) error {
24 return buildCheckpointDisplayPagerObserved(ctx, db, source, fingerprint, nil)
25 }
26
27 // The observer is used by deterministic interruption tests after durable batch
28 // publication. Production callers do not install one.
29 func buildCheckpointDisplayPagerObserved(ctx context.Context, db *sql.DB, source, fingerprint string, committed func(int)) error {
30 f, err := fileops.OpenReplaceableRead(source)
31 if err != nil {
32 return err
33 }
34 defer f.Close()
35 if err := validateCheckpointPagerFile(source, f, fingerprint); err != nil {
36 return err
37 }
38 hash := sha256.New()
39 idx := SessionDisplayIndex{SchemaVersion: SessionDisplayIndexSchemaVersion, ListingPreviewKnown: true}
40 users := 0
41 if err := restoreCheckpointPagerProgress(ctx, db, &idx, &users, hash); err != nil {
42 return err
43 }
44 if _, err := f.Seek(idx.TranscriptSize, io.SeekStart); err != nil {
45 return err
46 }
47 reader := bufio.NewReaderSize(&historywork.Reader{Context: ctx, Source: f}, historywork.ReadChunk)
48 done := false
49 for !done {
50 tx, err := db.BeginTx(ctx, nil)
51 if err != nil {
52 return err
53 }
54 for range historywork.BatchEntries {
55 line, readErr := readSessionDisplayIndexLine(reader)
56 if readErr != nil && !errors.Is(readErr, io.EOF) {
57 err = readErr
58 break
59 }
60 if len(line) > 0 {
61 var message provider.Message
62 if err = json.Unmarshal(line, &message); err != nil {
63 break
64 }
65 if message.Role == "" {
66 err = errors.New("checkpoint message has no role")
67 // Ancient event rows use kind/type instead of provider roles.
68 // Classify only the first record as an unsupported format;
69 // a foreign record after valid messages is damaged content.
70 if idx.MessageCount == 0 {
71 var event struct {
72 Kind string `json:"kind"`
73 Type string `json:"type"`
74 }
75 if json.Unmarshal(line, &event) == nil && (event.Kind != "" || event.Type != "") {
76 err = ErrDisplayFormatUnsupported
77 }
78 }
79 break
80 }
81 var entry DisplayIndexEntry
82 entry, idx.AuthoredTurns = classifyDisplayIndexMessage(message, idx.MessageCount, idx.TranscriptSize, int64(len(line)), idx.AuthoredTurns)
83 var encoded []byte
84 encoded, err = json.Marshal(entry)
85 if err != nil {
86 break
87 }
88 _, err = tx.ExecContext(ctx, `INSERT INTO entries VALUES(?,?,?,?,?,?,?)`, entry.Index, entry.Offset, entry.Length, entry.AuthoredTurn, entry.Role, users, encoded)
89 if err != nil {
90 break
91 }
92 encoded, err = json.Marshal(messageForSessionIdentity(message))
93 if err != nil {
94 break
95 }
96 _, _ = hash.Write(encoded)
97 _, _ = hash.Write([]byte{'\n'})
98 if entry.Role == provider.RoleUser && !entry.PinnedContextRevision {
99 users++
100 }
101 if entry.StartsTurn && idx.ListingPreview == "" {
102 idx.ListingPreview = truncatePreview(previewProse(UserMessageText(message)))
103 }
104 idx.MessageCount++
105 idx.TranscriptSize += int64(len(line))
106 }
107 if errors.Is(readErr, io.EOF) {
108 done = true
109 break
110 }
111 }
112 if err != nil {
113 _ = tx.Rollback()
114 return fmt.Errorf("build checkpoint display index: %w", err)
115 }
116 state, marshalErr := hash.(encoding.BinaryMarshaler).MarshalBinary()
117 if marshalErr == nil {
118 marshalErr = saveCheckpointPagerProgress(ctx, tx, idx, users, state)
119 }
120 if marshalErr != nil {
121 _ = tx.Rollback()
122 return marshalErr
123 }
124 if err := tx.Commit(); err != nil {
125 return err
126 }
127 if committed != nil {
128 committed(idx.MessageCount)
129 }
130 }
131 return publishCheckpointDisplayPager(ctx, db, source, f, fingerprint, idx, hash)
132 }
133
134 // Fence both the handle that supplied the bytes and its current pathname
135 // before publishing. A replacement or a newly authoritative event log cannot
136 // turn a resumed checkpoint prefix into a complete generation.
137 func validateCheckpointPagerFile(source string, f *os.File, fingerprint string) error {
138 info, err := f.Stat()
139 if err != nil {
140 return err
141 }
142 target, version := fileops.DiskHandleSnapshot(source, f, info)
143 if fmt.Sprintf("%s:%s:checkpoint", target.Key, version) != fingerprint {
144 return ErrDisplaySourceChanged
145 }
146 current, err := os.Stat(source)
147 if err != nil {
148 return err
149 }
150 target, version = fileops.DiskSnapshot(source, current)
151 if !os.SameFile(info, current) || fmt.Sprintf("%s:%s:checkpoint", target.Key, version) != fingerprint {
152 return ErrDisplaySourceChanged
153 }
154 return validateCheckpointDisplaySource(source)
155 }
156
156 lines GO