返回 DeepSeek-Reasonix
snapshot.go
根目录 / internal / transcript / snapshot.go
1 package transcript
2
3 import (
4 "encoding/json"
5 "errors"
6 "reflect"
7 "unicode/utf8"
8 )
9
10 const (
11 // MaxResponseBytes includes the complete encoded response and its newline.
12 MaxResponseBytes = 2 << 20
13 defaultPageRecords = 120
14 maxPageRecords = 1000
15 defaultPageBytes = 512 << 10
16 maxPageBytes = MaxResponseBytes
17 inlineFieldBytes = 64 << 10
18 previewBytes = 4 << 10
19 contentChunkBytes = 64 << 10
20 )
21
22 type PageRequest struct {
23 SnapshotID string `json:"snapshotId"`
24 MessageID string `json:"messageId,omitempty"`
25 Before int `json:"before"`
26 Records int `json:"records"`
27 Bytes int `json:"bytes"`
28 }
29
30 type ContentRef struct {
31 SnapshotID string `json:"snapshotId"`
32 RecordID string `json:"recordId"`
33 Path []string `json:"path"`
34 Bytes int `json:"bytes"`
35 }
36
37 type Record struct {
38 ID string `json:"id"`
39 Order int `json:"order"`
40 Message Message `json:"message"`
41 Refs []ContentRef `json:"refs"`
42 }
43
44 type Snapshot struct {
45 Boundary
46 Records []Record `json:"records"`
47 Runtime Runtime `json:"runtime"`
48 ActiveAttempts []ActiveAttempt `json:"activeAttempts"`
49 ActiveRecords []Record `json:"activeRecords"`
50 Before int `json:"before"`
51 HasOlder bool `json:"hasOlder"`
52 TotalRecords int `json:"totalRecords"`
53 TotalTurns int `json:"totalTurns"`
54 Stale bool `json:"stale"`
55 NotFound bool `json:"notFound,omitempty"`
56 }
57
58 type ContentRequest struct {
59 ContentRef
60 Offset int `json:"offset"`
61 }
62
63 type ContentChunk struct {
64 Data string `json:"data"`
65 NextOffset int `json:"nextOffset"`
66 Done bool `json:"done"`
67 Stale bool `json:"stale"`
68 }
69
70 // Snapshot returns a bounded immutable window. Pages and content refs are
71 // bound to this exact revision. A bounded cache retains immutable cuts while
72 // streaming continues; an evicted cut returns Stale rather than mixed history.
73 func (p *Projection) Snapshot(req PageRequest) (Snapshot, error) {
74 p.mu.Lock()
75 defer p.mu.Unlock()
76 frozen, err := p.freezeLocked(req.SnapshotID)
77 if err != nil {
78 return Snapshot{}, err
79 }
80 if frozen == nil {
81 return Snapshot{Boundary: p.boundaryLocked(), Stale: true}, nil
82 }
83 return frozen.snapshotCurrent(req)
84 }
85
86 func (p *Projection) snapshotCurrent(req PageRequest) (Snapshot, error) {
87 out := Snapshot{Boundary: p.boundaryLocked(), Records: []Record{}, ActiveAttempts: []ActiveAttempt{}, ActiveRecords: []Record{}}
88 if req.SnapshotID != "" && req.SnapshotID != out.SnapshotID {
89 out.Stale = true
90 return out, nil
91 }
92 limit := req.Records
93 if limit <= 0 {
94 limit = defaultPageRecords
95 }
96 limit = min(limit, maxPageRecords)
97 budget := req.Bytes
98 if budget <= 0 {
99 budget = defaultPageBytes
100 }
101 budget = min(budget, maxPageBytes)
102 end, found := p.snapshotEnd(req)
103 if !found {
104 out.NotFound = true
105 return out, nil
106 }
107 out.TotalRecords = len(p.buffer.messages)
108 out.TotalTurns = p.buffer.userTurns
109 used := 0
110 for i := end - 1; i >= 0 && len(out.Records) < limit; i-- {
111 row, err := boundedRecord(p.buffer.messages[i].materialize(), out.SnapshotID)
112 if err != nil {
113 return Snapshot{}, err
114 }
115 row.Order = i
116 encoded, err := json.Marshal(row)
117 if err != nil {
118 return Snapshot{}, err
119 }
120 if len(encoded) > maxPageBytes {
121 return Snapshot{}, errors.New("transcript record metadata exceeds the page limit")
122 }
123 if len(out.Records) > 0 && used+len(encoded) > budget {
124 break
125 }
126 out.Records = append(out.Records, row)
127 used += len(encoded)
128 }
129 for i, j := 0, len(out.Records)-1; i < j; i, j = i+1, j-1 {
130 out.Records[i], out.Records[j] = out.Records[j], out.Records[i]
131 }
132 out.Before = end - len(out.Records)
133 out.HasOlder = out.Before > 0
134 runtime, attempts := p.runtimeLocked()
135 // Detach nested prompt pointers before releasing the mutex.
136 encoded, err := json.Marshal(runtime)
137 if err != nil {
138 return Snapshot{}, err
139 }
140 if err = json.Unmarshal(encoded, &out.Runtime); err != nil {
141 return Snapshot{}, err
142 }
143 out.ActiveAttempts = attempts
144 // A running assistant can precede a large batch of tool result rows. Keep
145 // every still-mutable owner in the same snapshot cut so a later delta can
146 // never be appended to an unloaded prefix.
147 {
148 present := make(map[string]struct{}, len(out.Records))
149 for _, record := range out.Records {
150 present[record.ID] = struct{}{}
151 }
152 for _, i := range activeRecordIndexes(p.buffer.messages, out.Before, runtime) {
153 m := p.buffer.messages[i].materialize()
154 active := (!runtime.Status.Terminal() && m.Pending) || (m.Role == "user" && m.TurnID == runtime.TurnID && runtime.TurnID != "")
155 for _, call := range m.ToolCalls {
156 active = active || (!runtime.Status.Terminal() && call.Pending)
157 }
158 if !active {
159 continue
160 }
161 row, err := boundedRecord(m, out.SnapshotID)
162 if err != nil {
163 return Snapshot{}, err
164 }
165 row.Order = i
166 if _, alreadyPresent := present[row.ID]; alreadyPresent {
167 continue
168 }
169 present[row.ID] = struct{}{}
170 out.ActiveRecords = append(out.ActiveRecords, row)
171 }
172 }
173 encoded, err = json.Marshal(out)
174 if err != nil {
175 return Snapshot{}, err
176 }
177 if len(encoded)+1 > maxPageBytes {
178 return Snapshot{}, errors.New("transcript snapshot metadata exceeds the page limit")
179 }
180 return out, nil
181 }
182
183 // activeRecordIndexes walks only the current mutable turn. Older settled
184 // history cannot gain a streamed suffix, so scanning it on every page request
185 // only extends the projection mutex hold time with the total transcript size.
186 // A missing runtime turn ID keeps the conservative full-prefix behavior for
187 // legacy baselines that do not carry turn identity.
188 func activeRecordIndexes(messages []*bufferedMessage, before int, runtime Runtime) []int {
189 if runtime.Status.Terminal() {
190 return nil
191 }
192 if runtime.TurnID == "" {
193 indexes := make([]int, 0, before)
194 for i := range before {
195 indexes = append(indexes, i)
196 }
197 return indexes
198 }
199 indexes := make([]int, 0, 8)
200 seenCurrentTurn := false
201 for i := before - 1; i >= 0; i-- {
202 m := messages[i].materialize()
203 isCurrentTurn := m.TurnID == runtime.TurnID
204 active := (!runtime.Status.Terminal() && m.Pending) || (m.Role == "user" && isCurrentTurn)
205 for _, call := range m.ToolCalls {
206 active = active || (!runtime.Status.Terminal() && call.Pending)
207 }
208 if active {
209 indexes = append(indexes, i)
210 }
211 if isCurrentTurn {
212 seenCurrentTurn = true
213 continue
214 }
215 if seenCurrentTurn {
216 break
217 }
218 }
219 return indexes
220 }
221
222 func (p *Projection) Content(req ContentRequest) (ContentChunk, error) {
223 p.mu.Lock()
224 defer p.mu.Unlock()
225 if req.SnapshotID == "" {
226 return ContentChunk{}, errors.New("snapshotId is required")
227 }
228 frozen, err := p.freezeLocked(req.SnapshotID)
229 if err != nil {
230 return ContentChunk{}, err
231 }
232 if frozen == nil {
233 return ContentChunk{Stale: true}, nil
234 }
235 return frozen.contentCurrent(req)
236 }
237
238 func (p *Projection) contentCurrent(req ContentRequest) (ContentChunk, error) {
239 if req.SnapshotID != p.boundaryLocked().SnapshotID {
240 return ContentChunk{Stale: true}, nil
241 }
242 if req.Offset < 0 || len(req.Path) == 0 || len(req.Path) > 16 {
243 return ContentChunk{}, errors.New("invalid transcript content request")
244 }
245 for _, row := range p.buffer.messages {
246 if row.message.RecordID != req.RecordID {
247 continue
248 }
249 text, ok := contentStringAt(row.materialize(), req.Path)
250 if !ok || req.Offset > len(text) || (req.Offset < len(text) && !utf8.RuneStart(text[req.Offset])) {
251 return ContentChunk{}, errors.New("invalid transcript content offset")
252 }
253 end := runeBoundary(text, min(req.Offset+contentChunkBytes, len(text)))
254 return ContentChunk{Data: text[req.Offset:end], NextOffset: end, Done: end == len(text)}, nil
255 }
256 return ContentChunk{}, errors.New("transcript record not found")
257 }
258
259 func boundedRecord(message Message, snapshotID string) (Record, error) {
260 if message.RecordID == "" {
261 return Record{}, errors.New("transcript snapshot record identity is missing")
262 }
263 out := Record{ID: message.RecordID, Refs: []ContentRef{}}
264 out.Message = mapContentStrings(reflect.ValueOf(message), nil, func(text string, path []string) string {
265 if len(text) <= inlineFieldBytes {
266 return text
267 }
268 out.Refs = append(out.Refs, ContentRef{SnapshotID: snapshotID, RecordID: out.ID, Path: path, Bytes: len(text)})
269 return text[:runeBoundary(text, previewBytes)]
270 }).Interface().(Message)
271 return out, nil
272 }
273
274 func (p *Projection) snapshotEnd(req PageRequest) (int, bool) {
275 if req.MessageID != "" {
276 position, found := p.recordPositions[req.MessageID]
277 return position + 1, found
278 }
279 end := len(p.buffer.messages)
280 if req.SnapshotID != "" {
281 end = min(max(req.Before, 0), end)
282 }
283 return end, true
284 }
285
286 func runeBoundary(text string, offset int) int {
287 for offset > 0 && offset < len(text) && !utf8.RuneStart(text[offset]) {
288 offset--
289 }
290 return offset
291 }
292
292 lines GO