返回 DeepSeek-Reasonix
snapshot_cache.go
根目录 / internal / transcript / snapshot_cache.go
1 package transcript
2
3 import (
4 "errors"
5 "fmt"
6 "maps"
7 "reflect"
8 )
9
10 const snapshotCacheBytes = 64 << 20
11 const snapshotCacheEntries = 3
12
13 type frozenSnapshot struct {
14 projection *Projection
15 bytes int
16 }
17
18 func validateRecordIdentities(messages []*bufferedMessage) error {
19 seen := make(map[string]int, len(messages))
20 for index, row := range messages {
21 id := row.message.RecordID
22 if id == "" {
23 return errors.New("transcript snapshot record identity is missing")
24 }
25 if previous, exists := seen[id]; exists {
26 return fmt.Errorf("transcript snapshot record identity %q is duplicated at %d and %d", id, previous, index)
27 }
28 seen[id] = index
29 }
30 return nil
31 }
32
33 // freezeLocked retains a bounded number of immutable cuts. Strings already
34 // owned by settled rows are shared; only a live accumulator is materialized.
35 // Paging/content therefore continue to describe the requested cut while the
36 // live projection advances.
37 func (p *Projection) freezeLocked(id string) (*Projection, error) {
38 if id != "" {
39 if cached, ok := p.snapshots[id]; ok {
40 return cached.projection, nil
41 }
42 return nil, nil
43 }
44 id = p.boundaryLocked().SnapshotID
45 if cached, ok := p.snapshots[id]; ok {
46 return cached.projection, nil
47 }
48 frozen := &Projection{incarnation: p.incarnation, identity: p.identity, revision: p.revision, covered: p.covered, durable: p.durable,
49 runtime: p.runtime, attempts: maps.Clone(p.attempts), prompts: maps.Clone(p.prompts)}
50 frozen.buffer.userTurns = p.buffer.userTurns
51 frozen.recordPositions = make(map[string]int, len(p.buffer.messages))
52 used := 0
53 for position, row := range p.buffer.messages {
54 message := row.materialize()
55 if message.MessageID != "" {
56 frozen.recordPositions[message.MessageID] = position
57 }
58 used += retainedBytes(reflect.ValueOf(message))
59 copy := &bufferedMessage{message: message}
60 if message.Role == "assistant" {
61 copy.content.replace(message.Content)
62 copy.reasoning.replace(message.Reasoning)
63 }
64 frozen.buffer.messages = append(frozen.buffer.messages, copy)
65 }
66 // Reject a malformed projection before it becomes a reusable frozen cut.
67 // Empty or duplicate identities would alias pagination and content reads.
68 if err := validateRecordIdentities(frozen.buffer.messages); err != nil {
69 return nil, err
70 }
71 // The turn index is derived once per cut and shares this cut's lifetime, so
72 // paging the body never shrinks navigation and repeated outline reads reuse
73 // one pass. Its previews count against the same cache budget.
74 frozen.outline = buildOutline(frozen.buffer.messages)
75 used += retainedBytes(reflect.ValueOf(frozen.recordPositions))
76 used += retainedBytes(reflect.ValueOf(frozen.outline))
77 runtime, _ := frozen.runtimeLocked()
78 used += retainedBytes(reflect.ValueOf(runtime))
79 if p.snapshots == nil {
80 p.snapshots = make(map[string]frozenSnapshot)
81 }
82 // The current cut is pinned, including unusually large sessions. Settled
83 // body strings are shared with the projection; an oversized current cut
84 // evicts every older cut instead of making the conversation unreadable.
85 for len(p.snapshotOrder) > 0 && (len(p.snapshotOrder) >= snapshotCacheEntries || p.snapshotBytes+used > snapshotCacheBytes) {
86 oldest := p.snapshotOrder[0]
87 p.snapshotOrder = p.snapshotOrder[1:]
88 p.snapshotBytes -= p.snapshots[oldest].bytes
89 delete(p.snapshots, oldest)
90 }
91 p.snapshots[id] = frozenSnapshot{projection: frozen, bytes: used}
92 p.snapshotOrder = append(p.snapshotOrder, id)
93 p.snapshotBytes += used
94 return frozen, nil
95 }
96
96 lines GO