返回 DeepSeek-Reasonix
recovery_store.go
根目录 / internal / session / recovery_store.go
1 package session
2
3 import (
4 "bytes"
5 "context"
6 "crypto/sha256"
7 "encoding/json"
8 "errors"
9 "fmt"
10 "io"
11 "os"
12 "path/filepath"
13 "slices"
14 "sort"
15 "strings"
16 "time"
17
18 "github.com/klauspost/compress/zstd"
19 bolt "go.etcd.io/bbolt"
20
21 "reasonix/internal/agent"
22 "reasonix/internal/fileutil"
23 "reasonix/internal/provider"
24 "reasonix/internal/sessioncontent"
25 )
26
27 // Version 8 carries the live message ids the writer checks new completions
28 // against. Older projections are disposable and rebuild from the unchanged log.
29 const recoveryProjectionVersion = 8
30
31 const (
32 recoveryFormatVersion = 1
33 recoveryDBName = "recovery-v1.bolt"
34 storageIdentityName = "storage.identity.json"
35 RecentMessageLimit = 100
36 recentInlineBytes = 32 << 10
37 recentResponseBytes = 512 << 10
38 recentSnapshotName = "recent-v1.json"
39 )
40
41 var (
42 recoveryMetaBucket = []byte("meta")
43 recoveryCheckpointBucket = []byte("checkpoints")
44 recoveryOperationBucket = []byte("operations")
45 recoveryCurrentKey = []byte("current")
46 recoveryPreviousKey = []byte("previous")
47 )
48
49 // RecoveryOpenStats reports work performed by the normal open path. It is
50 // intentionally small and stable enough for capacity tests and host telemetry.
51 type RecoveryOpenStats struct {
52 UsedCheckpoint bool `json:"usedCheckpoint"`
53 LogBytesRead int64 `json:"logBytesRead"`
54 LogBytesTotal int64 `json:"logBytesTotal"`
55 TailCommits int `json:"tailCommits"`
56 }
57
58 // RecentSnapshot is the bounded, read-only baseline used before history or
59 // search projections are available.
60 type RecentSnapshot struct {
61 Version int `json:"version"`
62 SessionID string `json:"sessionId"`
63 StorageGeneration string `json:"storageGeneration"`
64 DurableSequence uint64 `json:"durableSequence"`
65 TotalTurns int `json:"totalTurns"`
66 Entries []PersistentMessage `json:"entries"`
67 Title string `json:"title,omitempty"`
68 ModelRef string `json:"modelRef,omitempty"`
69 ModelIdentity string `json:"modelIdentity,omitempty"`
70 }
71
72 type storageIdentity struct {
73 Version int `json:"version"`
74 SessionID string `json:"sessionId"`
75 Generation string `json:"generation"`
76 LogPrefixDigest string `json:"logPrefixDigest,omitempty"`
77 CreatedAt time.Time `json:"createdAt"`
78 }
79
80 type recoveryCheckpoint struct {
81 Version int `json:"version"`
82 SessionID string `json:"sessionId"`
83 StorageGeneration string `json:"storageGeneration"`
84 StorageRevision int `json:"storageRevision"`
85 DurableSequence uint64 `json:"durableSequence"`
86 LogOffset int64 `json:"logOffset"`
87 AnchorOffset int64 `json:"anchorOffset,omitempty"`
88 AnchorFirst uint64 `json:"anchorFirst,omitempty"`
89 AnchorCommitID string `json:"anchorCommitId,omitempty"`
90 AnchorHash string `json:"anchorHash,omitempty"`
91 ProjectionVersion int `json:"projectionVersion"`
92 Projection Projection `json:"projection"`
93 RecentMessages []provider.Message `json:"recentMessages,omitempty"`
94 CatalogPreview string `json:"catalogPreview,omitempty"`
95 MessageIDs []string `json:"messageIds,omitempty"`
96 CreatedAt time.Time `json:"createdAt"`
97 }
98
99 type recoveryOperation struct {
100 Hash string `json:"hash"`
101 CommitID string `json:"commitId"`
102 FirstSequence uint64 `json:"firstSequence"`
103 EventCount int `json:"eventCount"`
104 TurnID string `json:"turnId,omitempty"`
105 OperationID string `json:"operationId"`
106 OperationHash string `json:"operationHash"`
107 WriterGeneration uint64 `json:"writerGeneration"`
108 CreatedAt time.Time `json:"createdAt"`
109 }
110
111 func operationForRecovery(record operationRecord) recoveryOperation {
112 commit := record.commit
113 return recoveryOperation{
114 Hash: record.hash, CommitID: commit.ID, FirstSequence: commit.FirstSequence,
115 EventCount: commit.EventCount, TurnID: commit.TurnID, OperationID: commit.OperationID,
116 OperationHash: commit.OperationHash, WriterGeneration: commit.WriterGeneration,
117 CreatedAt: commit.CreatedAt,
118 }
119 }
120
121 func (o recoveryOperation) record() operationRecord {
122 return operationRecord{hash: o.Hash, commit: Commit{
123 SchemaVersion: SchemaVersion, Codec: Codec, RecordType: "commit", ID: o.CommitID,
124 OperationID: o.OperationID, OperationHash: o.OperationHash,
125 FirstSequence: o.FirstSequence, EventCount: o.EventCount, TurnID: o.TurnID,
126 WriterGeneration: o.WriterGeneration, CreatedAt: o.CreatedAt,
127 }}
128 }
129
130 type recoveryStore struct {
131 db *bolt.DB
132 path string
133 recent string
134 sessionDir string
135 identity storageIdentity
136 }
137
138 func recoveryCacheDir(sessionDir string) string {
139 return filepath.Join(filepath.Dir(sessionDir), ".recovery-cache", filepath.Base(sessionDir))
140 }
141
142 func ensureStorageIdentity(sessionDir string, manifest Manifest) (storageIdentity, error) {
143 path := filepath.Join(sessionDir, storageIdentityName)
144 prefix, prefixErr := storageLogPrefix(sessionDir, manifest)
145 if prefixErr != nil && !os.IsNotExist(prefixErr) {
146 return storageIdentity{}, prefixErr
147 }
148 if data, err := os.ReadFile(path); err == nil {
149 var identity storageIdentity
150 if json.Unmarshal(data, &identity) == nil && identity.Version == recoveryFormatVersion && identity.SessionID == manifest.SessionID && strings.TrimSpace(identity.Generation) != "" {
151 if identity.LogPrefixDigest == "" && prefix != "" {
152 identity.LogPrefixDigest = prefix
153 encoded, marshalErr := json.Marshal(identity)
154 if marshalErr != nil {
155 return storageIdentity{}, marshalErr
156 }
157 if writeErr := fileutil.AtomicWriteFileStrict(path, append(encoded, '\n'), 0o600); writeErr != nil {
158 return storageIdentity{}, writeErr
159 }
160 return identity, nil
161 }
162 if prefix == "" || identity.LogPrefixDigest == prefix {
163 return identity, nil
164 }
165 }
166 }
167 identity := storageIdentity{Version: recoveryFormatVersion, SessionID: manifest.SessionID, Generation: randomID(), LogPrefixDigest: prefix, CreatedAt: time.Now().UTC()}
168 data, err := json.Marshal(identity)
169 if err != nil {
170 return storageIdentity{}, err
171 }
172 if err := fileutil.AtomicWriteFileStrict(path, append(data, '\n'), 0o600); err != nil {
173 return storageIdentity{}, err
174 }
175 return identity, nil
176 }
177
178 func storageLogPrefix(sessionDir string, manifest Manifest) (string, error) {
179 file, err := os.Open(logPathForManifest(sessionDir, manifest))
180 if err != nil {
181 return "", err
182 }
183 defer file.Close()
184 // The framed transaction header makes the first 64 bytes immutable after
185 // the first commit. Hashing a larger short-file prefix would change merely
186 // because a normal append extended a log shorter than that prefix.
187 buffer := make([]byte, 64)
188 n, err := file.Read(buffer)
189 if err != nil && !errors.Is(err, io.EOF) {
190 return "", err
191 }
192 if n == 0 {
193 return "", nil
194 }
195 digest := sha256.Sum256(buffer[:n])
196 return fmt.Sprintf("%x", digest[:]), nil
197 }
198
199 func readStorageIdentity(sessionDir string, manifest Manifest) (storageIdentity, error) {
200 data, err := os.ReadFile(filepath.Join(sessionDir, storageIdentityName))
201 if err != nil {
202 return storageIdentity{}, err
203 }
204 var identity storageIdentity
205 if err := json.Unmarshal(data, &identity); err != nil {
206 return storageIdentity{}, err
207 }
208 if identity.Version != recoveryFormatVersion || identity.SessionID != manifest.SessionID || strings.TrimSpace(identity.Generation) == "" {
209 return storageIdentity{}, ErrStaleGeneration
210 }
211 prefix, err := storageLogPrefix(sessionDir, manifest)
212 if err != nil && !os.IsNotExist(err) {
213 return storageIdentity{}, err
214 }
215 if identity.LogPrefixDigest != "" && prefix != identity.LogPrefixDigest {
216 return storageIdentity{}, ErrStaleGeneration
217 }
218 return identity, nil
219 }
220
221 func openRecoveryStore(sessionDir string, identity storageIdentity) (*recoveryStore, error) {
222 dir := recoveryCacheDir(sessionDir)
223 if err := os.MkdirAll(dir, 0o700); err != nil {
224 return nil, err
225 }
226 path := filepath.Join(dir, recoveryDBName)
227 db, err := bolt.Open(path, 0o600, &bolt.Options{Timeout: 250 * time.Millisecond, NoFreelistSync: true})
228 if err != nil {
229 return nil, err
230 }
231 store := &recoveryStore{db: db, path: path, recent: filepath.Join(dir, recentSnapshotName), sessionDir: sessionDir, identity: identity}
232 err = db.Update(func(tx *bolt.Tx) error {
233 meta, err := tx.CreateBucketIfNotExists(recoveryMetaBucket)
234 if err != nil {
235 return err
236 }
237 checkpoints, err := tx.CreateBucketIfNotExists(recoveryCheckpointBucket)
238 if err != nil {
239 return err
240 }
241 operations, err := tx.CreateBucketIfNotExists(recoveryOperationBucket)
242 if err != nil {
243 return err
244 }
245 _ = checkpoints
246 _ = operations
247 stored := string(meta.Get([]byte("storage_generation")))
248 if stored != "" && stored != identity.Generation {
249 if err := tx.DeleteBucket(recoveryCheckpointBucket); err != nil {
250 return err
251 }
252 if err := tx.DeleteBucket(recoveryOperationBucket); err != nil {
253 return err
254 }
255 if _, err := tx.CreateBucket(recoveryCheckpointBucket); err != nil {
256 return err
257 }
258 if _, err := tx.CreateBucket(recoveryOperationBucket); err != nil {
259 return err
260 }
261 }
262 return meta.Put([]byte("storage_generation"), []byte(identity.Generation))
263 })
264 if err != nil {
265 _ = db.Close()
266 return nil, err
267 }
268 return store, nil
269 }
270
271 func openRecoveryStoreRepair(sessionDir string, identity storageIdentity) (*recoveryStore, error) {
272 store, err := openRecoveryStore(sessionDir, identity)
273 if err == nil {
274 return store, nil
275 }
276 path := filepath.Join(recoveryCacheDir(sessionDir), recoveryDBName)
277 if _, statErr := os.Stat(path); statErr == nil {
278 _ = os.Rename(path, path+fmt.Sprintf(".corrupt-%d", time.Now().UTC().UnixNano()))
279 }
280 return openRecoveryStore(sessionDir, identity)
281 }
282
283 func (s *recoveryStore) close() error {
284 if s == nil || s.db == nil {
285 return nil
286 }
287 return s.db.Close()
288 }
289
290 func encodeRecoveryValue(value any) ([]byte, error) {
291 raw, err := json.Marshal(value)
292 if err != nil {
293 return nil, err
294 }
295 encoder, err := zstd.NewWriter(nil, zstd.WithEncoderConcurrency(1), zstd.WithEncoderLevel(zstd.SpeedFastest))
296 if err != nil {
297 return nil, err
298 }
299 defer encoder.Close()
300 return encoder.EncodeAll(raw, nil), nil
301 }
302
303 func decodeRecoveryValue(data []byte, value any) error {
304 decoder, err := zstd.NewReader(nil, zstd.WithDecoderConcurrency(1), zstd.WithDecoderMaxMemory(64<<20))
305 if err != nil {
306 return err
307 }
308 defer decoder.Close()
309 raw, err := decoder.DecodeAll(data, nil)
310 if err != nil {
311 return err
312 }
313 return json.Unmarshal(raw, value)
314 }
315
316 func (s *recoveryStore) loadCheckpoints() ([]recoveryCheckpoint, error) {
317 if s == nil || s.db == nil {
318 return nil, os.ErrNotExist
319 }
320 var encoded [][]byte
321 err := s.db.View(func(tx *bolt.Tx) error {
322 bucket := tx.Bucket(recoveryCheckpointBucket)
323 if bucket == nil {
324 return os.ErrNotExist
325 }
326 for _, key := range [][]byte{recoveryCurrentKey, recoveryPreviousKey} {
327 if value := bucket.Get(key); value != nil {
328 encoded = append(encoded, append([]byte(nil), value...))
329 }
330 }
331 if len(encoded) == 0 {
332 return os.ErrNotExist
333 }
334 return nil
335 })
336 if err != nil {
337 return nil, err
338 }
339 checkpoints := make([]recoveryCheckpoint, 0, len(encoded))
340 for _, value := range encoded {
341 var checkpoint recoveryCheckpoint
342 if decodeRecoveryValue(value, &checkpoint) == nil {
343 checkpoints = append(checkpoints, checkpoint)
344 }
345 }
346 if len(checkpoints) == 0 {
347 return nil, ErrDamagedStore
348 }
349 return checkpoints, nil
350 }
351
352 func (s *recoveryStore) lookupOperation(operationID string) (operationRecord, bool, error) {
353 if s == nil || s.db == nil {
354 return operationRecord{}, false, nil
355 }
356 var encoded []byte
357 err := s.db.View(func(tx *bolt.Tx) error {
358 bucket := tx.Bucket(recoveryOperationBucket)
359 if bucket == nil {
360 return nil
361 }
362 encoded = append(encoded, bucket.Get([]byte(operationID))...)
363 return nil
364 })
365 if err != nil || len(encoded) == 0 {
366 return operationRecord{}, false, err
367 }
368 var operation recoveryOperation
369 if err := json.Unmarshal(encoded, &operation); err != nil {
370 return operationRecord{}, false, err
371 }
372 return operation.record(), true, nil
373 }
374
375 func (s *recoveryStore) publish(ctx context.Context, checkpoint recoveryCheckpoint, operations map[string]operationRecord) error {
376 if s == nil || s.db == nil {
377 return errors.New("session: recovery store unavailable")
378 }
379 if err := ctx.Err(); err != nil {
380 return err
381 }
382 if s.identity.LogPrefixDigest == "" {
383 manifest, err := readStoredManifest(filepath.Join(s.sessionDir, "manifest.json"))
384 if err != nil {
385 return err
386 }
387 prefix, err := storageLogPrefix(s.sessionDir, manifest)
388 if err != nil {
389 return err
390 }
391 if prefix != "" {
392 s.identity.LogPrefixDigest = prefix
393 encodedIdentity, err := json.Marshal(s.identity)
394 if err != nil {
395 return err
396 }
397 if err := fileutil.AtomicWriteFileStrict(filepath.Join(s.sessionDir, storageIdentityName), append(encodedIdentity, '\n'), 0o600); err != nil {
398 return err
399 }
400 }
401 }
402 checkpoint.Version = recoveryFormatVersion
403 checkpoint.StorageGeneration = s.identity.Generation
404 checkpoint.CreatedAt = time.Now().UTC()
405 encoded, err := encodeRecoveryValue(checkpoint)
406 if err != nil {
407 return err
408 }
409 keys := make([]string, 0, len(operations))
410 for key := range operations {
411 keys = append(keys, key)
412 }
413 sort.Strings(keys)
414 err = s.db.Update(func(tx *bolt.Tx) error {
415 checkpoints := tx.Bucket(recoveryCheckpointBucket)
416 operationBucket := tx.Bucket(recoveryOperationBucket)
417 if checkpoints == nil || operationBucket == nil {
418 return errors.New("session: recovery buckets unavailable")
419 }
420 if current := checkpoints.Get(recoveryCurrentKey); current != nil {
421 if err := checkpoints.Put(recoveryPreviousKey, current); err != nil {
422 return err
423 }
424 }
425 for _, key := range keys {
426 operation := operationForRecovery(operations[key])
427 value, err := json.Marshal(operation)
428 if err != nil {
429 return err
430 }
431 if err := operationBucket.Put([]byte(key), value); err != nil {
432 return err
433 }
434 }
435 if err := checkpoints.Put(recoveryCurrentKey, encoded); err != nil {
436 return err
437 }
438 return tx.Bucket(recoveryMetaBucket).Put([]byte("coverage_sequence"), fmt.Append(nil, checkpoint.DurableSequence))
439 })
440 if err != nil {
441 return err
442 }
443 recent := RecentSnapshot{
444 Version: recoveryFormatVersion, SessionID: checkpoint.SessionID,
445 StorageGeneration: checkpoint.StorageGeneration, DurableSequence: checkpoint.DurableSequence,
446 Title: checkpoint.Projection.Title,
447 ModelRef: checkpoint.Projection.ModelRef, ModelIdentity: checkpoint.Projection.ModelIdentity,
448 TotalTurns: visibleBoundaryCount(checkpoint.Projection, false),
449 }
450 recent.Entries, err = buildRecentEntries(ctx, s.sessionDir, checkpoint.RecentMessages, checkpoint.DurableSequence, recent.TotalTurns)
451 attachSubmissionEntries(checkpoint.Projection.Submissions, checkpoint.SessionID, recent.Entries)
452 if err != nil {
453 return err
454 }
455 data, err := json.Marshal(recent)
456 if err != nil {
457 return err
458 }
459 return fileutil.AtomicWriteFileStrict(s.recent, append(data, '\n'), 0o600)
460 }
461
462 func readRecentSnapshot(sessionDir string, identity storageIdentity) (RecentSnapshot, error) {
463 data, err := os.ReadFile(filepath.Join(recoveryCacheDir(sessionDir), recentSnapshotName))
464 if err != nil {
465 return RecentSnapshot{}, err
466 }
467 var snapshot RecentSnapshot
468 if err := json.Unmarshal(data, &snapshot); err != nil {
469 return RecentSnapshot{}, err
470 }
471 if snapshot.Version != recoveryFormatVersion || snapshot.SessionID != identity.SessionID || snapshot.StorageGeneration != identity.Generation {
472 return RecentSnapshot{}, ErrStaleGeneration
473 }
474 if len(snapshot.Entries) > RecentMessageLimit {
475 return RecentSnapshot{}, ErrDamagedStore
476 }
477 return snapshot, nil
478 }
479
480 // buildRecentEntries produces the bounded public baseline. Large canonical
481 // messages are stored once in ContentStore and represented by a preview plus a
482 // range-readable reference, so recent-v1.json cannot grow with tool output.
483 func buildRecentEntries(ctx context.Context, sessionDir string, messages []provider.Message, sequence uint64, totalTurns int) ([]PersistentMessage, error) {
484 if len(messages) > RecentMessageLimit {
485 messages = messages[len(messages)-RecentMessageLimit:]
486 }
487 entries := make([]PersistentMessage, 0, len(messages))
488 content := contentStoreForSessionDir(sessionDir)
489 inlineBytes := 0
490 visibleTurn := totalTurns
491 for _, message := range messages {
492 if agent.IsUserAuthoredTurnMessage(message) {
493 visibleTurn--
494 }
495 }
496 visibleTurn = max(visibleTurn, 0)
497 for position, message := range messages {
498 if err := ctx.Err(); err != nil {
499 return nil, err
500 }
501 if agent.IsUserAuthoredTurnMessage(message) {
502 visibleTurn++
503 }
504 body, err := json.Marshal(message)
505 if err != nil {
506 return nil, err
507 }
508 entry := PersistentMessage{
509 MessageID: message.ID, Position: int64(position + 1), Version: 1,
510 Role: string(message.Role), Preview: messagePreview(message),
511 EventSequence: sequence, VisibleTurn: visibleTurn,
512 }
513 if len(body) <= recentInlineBytes && inlineBytes+len(body) <= recentResponseBytes {
514 entry.Inline = body
515 inlineBytes += len(body)
516 } else {
517 ref, err := content.Put(ctx, bytes.NewReader(body), sessioncontent.Metadata{MediaType: "application/json"})
518 if err != nil {
519 return nil, err
520 }
521 entry.ContentRef = &ref
522 previewBody, err := recentDisplayMessage(message)
523 if err != nil {
524 return nil, err
525 }
526 if inlineBytes+len(previewBody) <= recentResponseBytes {
527 entry.Inline = previewBody
528 inlineBytes += len(previewBody)
529 }
530 }
531 entries = append(entries, entry)
532 }
533 return entries, nil
534 }
535
536 func recentDisplayMessage(message provider.Message) (json.RawMessage, error) {
537 preview := detachMessages([]provider.Message{message})[0]
538 preview.Content = messagePreview(message)
539 preview.RawContent = ""
540 preview.ProviderContent = ""
541 preview.Images = nil
542 preview.ImageInputs = nil
543 preview.ResponsesItems = nil
544 preview.ThinkingBlocks = nil
545 if runes := []rune(preview.ReasoningContent); len(runes) > 4096 {
546 preview.ReasoningContent = string(runes[:4096])
547 }
548 for i := range preview.ToolCalls {
549 if runes := []rune(preview.ToolCalls[i].Arguments); len(runes) > 2048 {
550 preview.ToolCalls[i].Arguments = string(runes[:2048])
551 }
552 }
553 body, err := json.Marshal(preview)
554 if err != nil {
555 return nil, err
556 }
557 if len(body) <= recentInlineBytes {
558 return body, nil
559 }
560 // Preserve the fields required to place the row even when optional display
561 // metadata alone exceeds the per-message preview budget.
562 return json.Marshal(provider.Message{
563 ID: preview.ID, Role: preview.Role, Content: preview.Content,
564 ToolCallID: preview.ToolCallID, Name: preview.Name,
565 CreatedAt: preview.CreatedAt, WorkDurationMs: preview.WorkDurationMs,
566 })
567 }
568
569 func checkpointFromStartup(manifest Manifest, identity storageIdentity, state *startupSessionState) recoveryCheckpoint {
570 projection, _ := Project(nil)
571 if state != nil {
572 projection = cloneProjection(state.projection)
573 projection.Messages = nil
574 }
575 checkpoint := recoveryCheckpoint{
576 Version: recoveryFormatVersion, SessionID: manifest.SessionID,
577 StorageGeneration: identity.Generation, StorageRevision: StorageRevision,
578 ProjectionVersion: recoveryProjectionVersion, Projection: projection,
579 }
580 if state != nil {
581 checkpoint.DurableSequence = state.durable
582 checkpoint.LogOffset = state.tip.LogOffset
583 checkpoint.AnchorOffset = state.tip.AnchorOffset
584 checkpoint.AnchorFirst = state.tip.AnchorFirst
585 checkpoint.AnchorCommitID = state.tip.AnchorCommitID
586 checkpoint.AnchorHash = state.tip.AnchorHash
587 checkpoint.RecentMessages = detachMessages(state.recentMessages)
588 checkpoint.CatalogPreview = state.catalogPreview
589 checkpoint.MessageIDs = state.messageIDs.list()
590 }
591 return checkpoint
592 }
593
594 func loadRecoveryStartupState(ctx context.Context, dir string, file *os.File, info os.FileInfo, recovery *recoveryStore, identity storageIdentity) (*startupSessionState, int64, bool, RecoveryOpenStats, bool) {
595 stats := RecoveryOpenStats{LogBytesTotal: info.Size()}
596 checkpoints, err := recovery.loadCheckpoints()
597 if err != nil {
598 return nil, 0, false, stats, false
599 }
600 for _, checkpoint := range checkpoints {
601 if state, end, torn, attempt, ok := tryRecoveryCheckpoint(ctx, dir, file, info, identity, checkpoint); ok {
602 return state, end, torn, attempt, true
603 }
604 }
605 return nil, 0, false, stats, false
606 }
607
608 func tryRecoveryCheckpoint(ctx context.Context, dir string, file *os.File, info os.FileInfo, identity storageIdentity, checkpoint recoveryCheckpoint) (*startupSessionState, int64, bool, RecoveryOpenStats, bool) {
609 stats := RecoveryOpenStats{LogBytesTotal: info.Size()}
610 if checkpoint.Version != recoveryFormatVersion || checkpoint.ProjectionVersion != recoveryProjectionVersion ||
611 checkpoint.SessionID != identity.SessionID || checkpoint.StorageGeneration != identity.Generation ||
612 checkpoint.StorageRevision != StorageRevision || checkpoint.LogOffset < 0 || checkpoint.LogOffset > info.Size() ||
613 checkpoint.Projection.CommittedSequence != checkpoint.DurableSequence {
614 return nil, 0, false, stats, false
615 }
616 if checkpoint.DurableSequence == 0 {
617 if checkpoint.LogOffset != 0 {
618 return nil, 0, false, stats, false
619 }
620 } else {
621 if checkpoint.AnchorOffset < 0 || checkpoint.AnchorOffset >= checkpoint.LogOffset || checkpoint.AnchorFirst == 0 || checkpoint.AnchorCommitID == "" {
622 return nil, 0, false, stats, false
623 }
624 var anchor Commit
625 var anchorEnd int64
626 err := scanV4CommitFileRefs(ctx, file, checkpoint.AnchorOffset, checkpoint.AnchorFirst, contentStoreForSessionDir(dir), nil, func(_ int64, commit Commit) bool {
627 anchor = commit
628 anchorEnd, _ = file.Seek(0, 1)
629 return false
630 })
631 if err != nil || anchor.ID != checkpoint.AnchorCommitID || anchor.OperationHash != checkpoint.AnchorHash ||
632 anchor.LastSequence() != checkpoint.DurableSequence || anchorEnd != checkpoint.LogOffset {
633 return nil, 0, false, stats, false
634 }
635 stats.LogBytesRead += anchorEnd - checkpoint.AnchorOffset
636 }
637
638 state := &startupSessionState{
639 projection: cloneProjection(checkpoint.Projection), operations: map[string]operationRecord{},
640 durable: checkpoint.DurableSequence, catalogPreview: checkpoint.CatalogPreview,
641 recentMessages: detachMessages(checkpoint.RecentMessages),
642 messageIDs: identitiesOf(checkpoint.MessageIDs),
643 tip: durableTip{LogOffset: checkpoint.LogOffset, AnchorOffset: checkpoint.AnchorOffset,
644 AnchorFirst: checkpoint.AnchorFirst, AnchorCommitID: checkpoint.AnchorCommitID, AnchorHash: checkpoint.AnchorHash},
645 }
646 var projectionErr error
647 content := contentStoreForSessionDir(dir)
648 repeated := map[uint64]bool{}
649 err := scanV4CommitFile(ctx, file, checkpoint.LogOffset, checkpoint.DurableSequence+1, content, nil, func(offset int64, commit Commit) bool {
650 kept := commit
651 kept.Events = slices.DeleteFunc(slices.Clone(commit.Events), func(event Event) bool { return !state.messageIDs.admitEvent(event, repeated) })
652 if err := applyProjectionCommit(&state.projection, kept); err != nil {
653 projectionErr = err
654 return false
655 }
656 if err := applyRecentCommit(&state.recentMessages, kept); err != nil {
657 projectionErr = err
658 return false
659 }
660 state.operations[commit.OperationID] = compactOperationRecord(commit)
661 state.durable = commit.LastSequence()
662 state.tip.AnchorOffset = offset
663 state.tip.AnchorFirst = commit.FirstSequence
664 state.tip.AnchorCommitID = commit.ID
665 state.tip.AnchorHash = commit.OperationHash
666 state.tip.LogOffset, _ = file.Seek(0, 1)
667 stats.TailCommits++
668 return true
669 })
670 if err != nil || projectionErr != nil {
671 return nil, 0, false, stats, false
672 }
673 state.projection.Messages = nil
674 state.projection.CommittedSequence = state.durable
675 stats.UsedCheckpoint = true
676 stats.LogBytesRead += max(state.tip.LogOffset-checkpoint.LogOffset, 0)
677 return state, state.tip.LogOffset, state.tip.LogOffset < info.Size(), stats, true
678 }
679
680 func applyRecentCommit(messages *[]provider.Message, commit Commit) error {
681 projection, _ := Project(nil)
682 projection.Messages = detachMessages(*messages)
683 recent := commit
684 recent.Events = nil
685 for _, event := range commit.Events {
686 switch event.Kind {
687 case "message/complete", "message/upsert", "message/retract", "history/replace", "legacy/import":
688 recent.Events = append(recent.Events, event)
689 }
690 }
691 if len(recent.Events) == 0 {
692 return nil
693 }
694 if err := applyProjectionCommit(&projection, recent); err != nil {
695 return err
696 }
697 if len(projection.Messages) > RecentMessageLimit {
698 projection.Messages = append([]provider.Message(nil), projection.Messages[len(projection.Messages)-RecentMessageLimit:]...)
699 }
700 *messages = detachMessages(projection.Messages)
701 return nil
702 }
703
703 lines GO