返回 DeepSeek-Reasonix
index.go
根目录 / internal / session / index.go
1 package session
2
3 import (
4 "context"
5 "crypto/sha256"
6 "encoding/hex"
7 "encoding/json"
8 "fmt"
9 "io"
10 "os"
11 "path/filepath"
12 "sort"
13
14 "reasonix/internal/fileutil"
15 )
16
17 const (
18 sparseIndexCodec = "reasonix.session.offset-index/v1"
19 sparseIndexInterval = 256
20 sparseIdentityBytes = 4096
21 )
22
23 type sparseIndex struct {
24 Codec string `json:"codec"`
25 LogSize int64 `json:"logSize"`
26 LogModTimeNS int64 `json:"logModTimeNs"`
27 LogIdentity string `json:"logIdentity"`
28 LastSequence uint64 `json:"lastSequence"`
29 CommitCount uint64 `json:"commitCount"`
30 Entries []sparseIndexEntry `json:"entries"`
31 partial bool
32 }
33
34 type sparseIndexEntry struct {
35 FirstSequence uint64 `json:"firstSequence"`
36 Offset int64 `json:"offset"`
37 }
38
39 func sparseIndexPath(cacheDir string) string {
40 return filepath.Join(cacheDir, "events.offset-index.json")
41 }
42
43 // loadOrBuildSparseIndex treats the index as an expendable cache. A missing or
44 // corrupt cache is rebuilt by validating the complete durable log. Failure to
45 // write the rebuilt cache does not make an otherwise readable session fail.
46 func loadOrBuildSparseIndex(ctx context.Context, dir, cacheDir string) (sparseIndex, error) {
47 if err := ctx.Err(); err != nil {
48 return sparseIndex{}, err
49 }
50 manifest, err := readStoredManifest(filepath.Join(dir, "manifest.json"))
51 if err != nil {
52 return sparseIndex{}, err
53 }
54 logPath := logPathForManifest(dir, manifest)
55 file, err := os.Open(logPath)
56 if os.IsNotExist(err) {
57 return sparseIndex{Codec: sparseIndexCodec, Entries: []sparseIndexEntry{}}, nil
58 }
59 if err != nil {
60 return sparseIndex{}, err
61 }
62 defer file.Close()
63 info, err := file.Stat()
64 if err != nil {
65 return sparseIndex{}, err
66 }
67 identity, err := sparseLogIdentity(file, info)
68 if err != nil {
69 return sparseIndex{}, err
70 }
71 if cached, ok := readValidSparseIndex(cacheDir, info, identity); ok {
72 return cached, nil
73 }
74
75 rebuilt := sparseIndex{
76 Codec: sparseIndexCodec,
77 LogSize: info.Size(),
78 LogModTimeNS: info.ModTime().UnixNano(),
79 LogIdentity: identity,
80 Entries: []sparseIndexEntry{},
81 }
82 commitIndex := 0
83 visit := func(offset int64, commit Commit) bool {
84 if ctx.Err() != nil {
85 return false
86 }
87 if commitIndex%sparseIndexInterval == 0 {
88 rebuilt.Entries = append(rebuilt.Entries, sparseIndexEntry{FirstSequence: commit.FirstSequence, Offset: offset})
89 }
90 commitIndex++
91 rebuilt.CommitCount++
92 rebuilt.LastSequence = commit.LastSequence()
93 return true
94 }
95 if manifest.Codec == Codec {
96 err = scanV4CommitFileRefs(ctx, file, 0, 1, contentStoreForSessionDir(dir), nil, visit)
97 } else {
98 err = scanCommitFileCodec(file, 0, 1, manifest.Codec, nil, visit)
99 }
100 if err != nil {
101 return sparseIndex{}, err
102 }
103 if err := ctx.Err(); err != nil {
104 return sparseIndex{}, err
105 }
106 writeSparseIndex(cacheDir, rebuilt)
107 return rebuilt, nil
108 }
109
110 func readValidSparseIndex(cacheDir string, info os.FileInfo, identity string) (sparseIndex, bool) {
111 data, err := os.ReadFile(sparseIndexPath(cacheDir))
112 if err != nil {
113 return sparseIndex{}, false
114 }
115 var cached sparseIndex
116 if json.Unmarshal(data, &cached) != nil || !cached.validFor(info, identity) {
117 return sparseIndex{}, false
118 }
119 return cached, true
120 }
121
122 func (idx sparseIndex) validFor(info os.FileInfo, identity string) bool {
123 if idx.Codec != sparseIndexCodec || idx.LogSize != info.Size() || idx.LogModTimeNS != info.ModTime().UnixNano() || idx.LogIdentity != identity {
124 return false
125 }
126 var previousSequence uint64
127 var previousOffset int64 = -1
128 for i, entry := range idx.Entries {
129 if entry.FirstSequence == 0 || entry.Offset < 0 || entry.Offset >= idx.LogSize || entry.FirstSequence <= previousSequence || entry.Offset <= previousOffset {
130 return false
131 }
132 if i == 0 && (entry.FirstSequence != 1 || entry.Offset != 0) {
133 return false
134 }
135 previousSequence, previousOffset = entry.FirstSequence, entry.Offset
136 }
137 return (idx.LastSequence == 0) == (len(idx.Entries) == 0) && (idx.CommitCount == 0) == (idx.LastSequence == 0)
138 }
139
140 func writeSparseIndex(cacheDir string, index sparseIndex) {
141 data, err := json.Marshal(index)
142 if err != nil || os.MkdirAll(cacheDir, 0o700) != nil {
143 return
144 }
145 _ = fileutil.AtomicWriteFileStrict(sparseIndexPath(cacheDir), append(data, '\n'), 0o600)
146 }
147
148 func (s *Store) recordPersistedIndex(file *os.File, start int64, commits []Commit, lengths []int64) {
149 if s == nil || file == nil || len(commits) == 0 || len(commits) != len(lengths) {
150 return
151 }
152 info, err := file.Stat()
153 if err != nil {
154 return
155 }
156 offset := start
157 for i, commit := range commits {
158 if i == len(commits)-1 {
159 s.tip = durableTip{LogOffset: offset + int64(lengths[i]), AnchorOffset: offset, AnchorFirst: commit.FirstSequence, AnchorCommitID: commit.ID, AnchorHash: commit.OperationHash}
160 }
161 offset += int64(lengths[i])
162 }
163 identity, err := sparseLogIdentity(file, info)
164 if err != nil {
165 return
166 }
167 s.indexMu.Lock()
168 index := s.index
169 if index.partial {
170 index.LogSize = info.Size()
171 index.LogModTimeNS = info.ModTime().UnixNano()
172 index.LogIdentity = identity
173 index.LastSequence = commits[len(commits)-1].LastSequence()
174 s.index = index
175 s.indexMu.Unlock()
176 return
177 }
178 if index.Codec != sparseIndexCodec || index.LogSize != start {
179 s.indexMu.Unlock()
180 _ = s.rebuildWriterIndex(file)
181 return
182 }
183 offset = start
184 for i, commit := range commits {
185 if index.CommitCount%sparseIndexInterval == 0 {
186 index.Entries = append(index.Entries, sparseIndexEntry{FirstSequence: commit.FirstSequence, Offset: offset})
187 }
188 index.CommitCount++
189 index.LastSequence = commit.LastSequence()
190 offset += int64(lengths[i])
191 }
192 index.LogSize = info.Size()
193 index.LogModTimeNS = info.ModTime().UnixNano()
194 index.LogIdentity = identity
195 s.index = index
196 s.indexMu.Unlock()
197 writeSparseIndex(s.dir, index)
198 }
199
200 // A reopened writer knows the durable sequence but not the older checkpoints.
201 // Adopt an already validated cache before appending, so the new commits can
202 // extend it without scanning history or leaving the disk index stale.
203 func (s *Store) adoptPersistedIndex(file *os.File, start int64) {
204 s.indexMu.Lock()
205 partial, lastSequence := s.index.partial, s.index.LastSequence
206 s.indexMu.Unlock()
207 if !partial {
208 return
209 }
210 info, err := file.Stat()
211 if err != nil || info.Size() != start {
212 return
213 }
214 identity, err := sparseLogIdentity(file, info)
215 if err != nil {
216 return
217 }
218 cached, ok := readValidSparseIndex(s.dir, info, identity)
219 if !ok || cached.LastSequence != lastSequence {
220 return
221 }
222 s.indexMu.Lock()
223 if s.index.partial && s.index.LastSequence == lastSequence {
224 s.index = cached
225 }
226 s.indexMu.Unlock()
227 }
228
229 func (s *Store) rebuildWriterIndex(_ *os.File) error {
230 index, err := loadOrBuildSparseIndex(context.Background(), s.dir, s.dir)
231 if err != nil {
232 return err
233 }
234 s.indexMu.Lock()
235 s.index = index
236 s.indexMu.Unlock()
237 return nil
238 }
239
240 func (idx sparseIndex) checkpoint(offset uint64) sparseIndexEntry {
241 if len(idx.Entries) == 0 {
242 return sparseIndexEntry{FirstSequence: 1}
243 }
244 target := offset + 1
245 position := max(sort.Search(len(idx.Entries), func(i int) bool { return idx.Entries[i].FirstSequence > target })-1, 0)
246 return idx.Entries[position]
247 }
248
249 func sparseLogIdentity(file *os.File, info os.FileInfo) (string, error) {
250 hash := sha256.New()
251 _, _ = fmt.Fprintf(hash, "%d:%d:", info.Size(), info.ModTime().UnixNano())
252 readChunk := func(offset, length int64) error {
253 if length <= 0 {
254 return nil
255 }
256 buf := make([]byte, length)
257 n, err := file.ReadAt(buf, offset)
258 if err != nil && err != io.EOF {
259 return err
260 }
261 _, _ = hash.Write(buf[:n])
262 return nil
263 }
264 first := min(info.Size(), sparseIdentityBytes)
265 if err := readChunk(0, first); err != nil {
266 return "", err
267 }
268 if info.Size() > first {
269 last := min(info.Size()-first, sparseIdentityBytes)
270 if err := readChunk(info.Size()-last, last); err != nil {
271 return "", err
272 }
273 }
274 return hex.EncodeToString(hash.Sum(nil)), nil
275 }
276
276 lines GO