返回 DeepSeek-Reasonix
layout_v3.go
根目录 / internal / checkpoint / layout_v3.go
1 package checkpoint
2
3 import (
4 "encoding/json"
5 "fmt"
6 "os"
7 "path/filepath"
8 "strconv"
9
10 "reasonix/internal/fileutil"
11 fileenc "reasonix/internal/fileutil/encoding"
12 )
13
14 func (s *Store) turnsDir() string {
15 return filepath.Join(s.dir, "turns")
16 }
17
18 func (s *Store) turnDir(turn int) string {
19 return filepath.Join(s.turnsDir(), strconv.Itoa(turn))
20 }
21
22 func (s *Store) v3MetaPath(turn int) string {
23 return filepath.Join(s.turnDir(turn), "meta.json")
24 }
25
26 func (s *Store) v3BeforePath(turn, index int) string {
27 return filepath.Join(s.turnDir(turn), "files", fmt.Sprintf("%04d.before", index))
28 }
29
30 func v3PayloadBytes(f FileSnap) ([]byte, error) {
31 if f.rawContent != nil {
32 return f.rawContent, nil
33 }
34 if f.Content == nil {
35 return nil, nil
36 }
37 if f.Encoding != nil {
38 return fileenc.Encode(*f.Content, *f.Encoding)
39 }
40 return []byte(*f.Content), nil
41 }
42
43 func (s *Store) persistV3(c *Checkpoint) error {
44 // Previous builds only inspect turn-N.json. Keep a payload-free v2 marker
45 // so their NextTurn remains monotonic across a downgrade; write it first so
46 // a crash can leave reduced rewind visibility, never an invisible turn.
47 marker := *c
48 marker.Result = nil
49 marker.SchemaVersion = SchemaV2
50 marker.Files = []FileSnap{}
51 marker.Coverage = CoverageNone
52 marker.CoverageGaps = nil
53 marker.ActiveWriters = nil
54 marker.ExpiredFilePayload = true
55 markerBytes, err := json.Marshal(&marker)
56 if err != nil {
57 return err
58 }
59 if err := os.MkdirAll(s.dir, 0o755); err != nil {
60 return err
61 }
62 markerPath := filepath.Join(s.dir, fmt.Sprintf("turn-%d.json", c.Turn))
63 if err := fileutil.AtomicWriteFileStrict(markerPath, markerBytes, 0o644); err != nil {
64 return err
65 }
66
67 turnDir := s.turnDir(c.Turn)
68 if err := os.MkdirAll(filepath.Join(turnDir, "files"), 0o755); err != nil {
69 return err
70 }
71 wire := *c
72 wire.SchemaVersion = SchemaV3
73 wire.Files = make([]FileSnap, len(c.Files))
74 for i, f := range c.Files {
75 snap := f
76 snap.BlobRef = ""
77 payloadPath := s.v3BeforePath(c.Turn, i)
78 if f.Content != nil && !f.PayloadExpired {
79 payload, err := v3PayloadBytes(f)
80 if err != nil {
81 return err
82 }
83 if err := fileutil.AtomicWriteFile(payloadPath, payload, 0o644); err != nil {
84 return err
85 }
86 snap.Content = nil
87 snap.Encoding = nil
88 } else if err := os.Remove(payloadPath); err != nil && !os.IsNotExist(err) {
89 return err
90 }
91 wire.Files[i] = snap
92 }
93 b, err := json.Marshal(&wire)
94 if err != nil {
95 return err
96 }
97 return fileutil.AtomicWriteFileStrict(s.v3MetaPath(c.Turn), b, 0o644)
98 }
99
100 func (s *Store) removeTurnArtifacts(turns map[int]bool) error {
101 if s.dir == "" {
102 return nil
103 }
104 for turn := range turns {
105 // The compatibility marker is the cross-version liveness record. Remove
106 // payloads first so a crash can leave a visible payload-free turn, never
107 // a markerless directory that a newer reader resurrects after downgrade.
108 if err := os.RemoveAll(s.turnDir(turn)); err != nil {
109 return fmt.Errorf("remove checkpoint turn %d: %w", turn, err)
110 }
111 for _, path := range []string{
112 filepath.Join(s.dir, fmt.Sprintf("turn-%d.json", turn)),
113 filepath.Join(s.expiredDir(), fmt.Sprintf("turn-%d.json", turn)),
114 } {
115 if err := os.Remove(path); err != nil && !os.IsNotExist(err) {
116 return fmt.Errorf("remove checkpoint turn %d: %w", turn, err)
117 }
118 }
119 }
120 return nil
121 }
122
123 func (s *Store) loadV3Turns() []*Checkpoint {
124 ents, err := os.ReadDir(s.turnsDir())
125 if err != nil {
126 return nil
127 }
128 var turns []*Checkpoint
129 for _, e := range ents {
130 if !e.IsDir() {
131 continue
132 }
133 turn, err := strconv.Atoi(e.Name())
134 if err != nil {
135 continue
136 }
137 // Old readers truncate only turn-N.json. Treat marker absence as a
138 // tombstone so reopening with a newer build cannot revive those turns.
139 if _, err := os.Stat(filepath.Join(s.dir, fmt.Sprintf("turn-%d.json", turn))); err != nil {
140 continue
141 }
142 b, err := fileenc.ReadFileUTF8(s.v3MetaPath(turn))
143 if err != nil {
144 continue
145 }
146 var c Checkpoint
147 if json.Unmarshal(b, &c) != nil {
148 continue
149 }
150 c.Turn = turn
151 c.SchemaVersion = SchemaV3
152 for i := range c.Files {
153 if c.Files[i].PayloadExpired {
154 continue
155 }
156 raw, err := os.ReadFile(s.v3BeforePath(turn, i))
157 if err != nil {
158 continue
159 }
160 if want := c.Files[i].SHA256; want != "" && Digest(raw) != want {
161 continue
162 }
163 enc, detected := fileenc.Detect(raw)
164 text := string(fileenc.Decode(detected, enc))
165 c.Files[i].Content = &text
166 c.Files[i].Encoding = &enc
167 c.Files[i].BlobRef = ""
168 c.Files[i].rawContent = append([]byte(nil), raw...)
169 }
170 turns = append(turns, &c)
171 }
172 return turns
173 }
174
175 func (s *Store) v3PayloadSize(turn int) (int64, error) {
176 // Metadata also holds the bounded frozen result patches.
177 root := s.turnDir(turn)
178 var total int64
179 err := filepath.WalkDir(root, func(_ string, d os.DirEntry, err error) error {
180 if err != nil {
181 if os.IsNotExist(err) {
182 return nil
183 }
184 return err
185 }
186 if d.IsDir() {
187 return nil
188 }
189 info, err := d.Info()
190 if err != nil {
191 return err
192 }
193 if info.Mode().IsRegular() {
194 total += info.Size()
195 }
196 return nil
197 })
198 if os.IsNotExist(err) {
199 return 0, nil
200 }
201 return total, err
202 }
203
204 // pruneV3TurnsLocked applies both count retention and the legacy 1 GiB soft
205 // payload budget to complete v3 turn directories. The current and protected
206 // turns may temporarily exceed the budget; the next unprotected turn prunes
207 // whole oldest directories. Older v1/v2 metadata remains on its legacy path.
208 func (s *Store) pruneV3TurnsLocked() {
209 if s.dir == "" || s.retainN <= 0 {
210 return
211 }
212 var turns []*Checkpoint
213 for _, c := range s.all() {
214 if c.SchemaVersion >= SchemaV3 {
215 turns = append(turns, c)
216 }
217 }
218 excess := len(turns) - s.retainN
219 sizes := make(map[*Checkpoint]int64, len(turns))
220 var totalSize int64
221 sizeKnown := true
222 for _, c := range turns {
223 size, err := s.v3PayloadSize(c.Turn)
224 if err != nil {
225 sizeKnown = false
226 break
227 }
228 sizes[c] = size
229 totalSize += size
230 }
231 quotaExceeded := func() bool {
232 return sizeKnown && s.blobQuota > 0 && totalSize > s.blobQuota
233 }
234 if excess <= 0 && !quotaExceeded() {
235 return
236 }
237 removed := make(map[*Checkpoint]bool)
238 for _, c := range turns {
239 if excess <= 0 && !quotaExceeded() {
240 break
241 }
242 if c == s.cur || s.protectTurns[c.Turn] {
243 continue
244 }
245 if err := s.removeTurnArtifacts(map[int]bool{c.Turn: true}); err != nil {
246 continue
247 }
248 removed[c] = true
249 if excess > 0 {
250 excess--
251 }
252 totalSize -= sizes[c]
253 }
254 if len(removed) == 0 {
255 return
256 }
257 kept := s.done[:0]
258 for _, c := range s.done {
259 if !removed[c] {
260 kept = append(kept, c)
261 }
262 }
263 s.done = kept
264 // Pre-v3 builds briefly wrote both a turn directory and a blob. Once such a
265 // turn ages out, the legacy mark-and-sweep can reclaim its orphaned blob.
266 s.pruneBlobsLocked()
267 }
268
268 lines GO