返回 DeepSeek-Reasonix
native_store.go
根目录 / internal / session / native_store.go
1 package session
2
3 import (
4 "bufio"
5 "bytes"
6 "context"
7 "encoding/json"
8 "errors"
9 "fmt"
10 "io"
11 "os"
12 "path/filepath"
13 "strings"
14 )
15
16 func loadStartupSessionStateForManifest(ctx context.Context, dir, eventsPath string, manifest Manifest, externalHistory bool, recovery *recoveryStore, identity storageIdentity) (*startupSessionState, int64, bool, RecoveryOpenStats, bool, error) {
17 if manifest.Codec == Codec {
18 return loadStartupSessionState(ctx, dir, eventsPath, externalHistory, recovery, identity)
19 }
20 file, err := os.Open(eventsPath)
21 if os.IsNotExist(err) {
22 projection, _ := Project(nil)
23 return &startupSessionState{projection: projection, operations: map[string]operationRecord{}}, 0, false, RecoveryOpenStats{}, false, nil
24 }
25 if err != nil {
26 return nil, 0, false, RecoveryOpenStats{}, false, err
27 }
28 defer file.Close()
29 info, err := file.Stat()
30 if err != nil {
31 return nil, 0, false, RecoveryOpenStats{}, false, err
32 }
33 projection, _ := Project(nil)
34 state := &startupSessionState{projection: projection, operations: map[string]operationRecord{}}
35 var durableEnd int64
36 var projectionErr error
37 err = scanCommitFileCodecBoundaries(ctx, file, 0, 1, manifest.Codec, nil, func(offset, end int64, commit Commit) bool {
38 if applyErr := applyProjectionCommit(&state.projection, commit); applyErr != nil {
39 projectionErr = applyErr
40 return false
41 }
42 if externalHistory {
43 for _, message := range state.projection.Messages {
44 if state.catalogPreview == "" {
45 state.catalogPreview = catalogMessagePreview(message)
46 }
47 }
48 state.projection.Messages = nil
49 }
50 state.operations[commit.OperationID] = compactOperationRecord(commit)
51 state.durable = commit.LastSequence()
52 durableEnd = end
53 state.tip = durableTip{LogOffset: end, AnchorOffset: offset, AnchorFirst: commit.FirstSequence, AnchorCommitID: commit.ID, AnchorHash: commit.OperationHash}
54 if applyErr := applyRecentCommit(&state.recentMessages, commit); applyErr != nil {
55 projectionErr = applyErr
56 return false
57 }
58 state.messageIDs.admit(commit)
59 return true
60 })
61 if err != nil || projectionErr != nil {
62 return nil, 0, false, RecoveryOpenStats{LogBytesTotal: info.Size(), LogBytesRead: durableEnd}, false, errors.Join(err, projectionErr)
63 }
64 state.projection.CommittedSequence = state.durable
65 return state, durableEnd, durableEnd < info.Size(), RecoveryOpenStats{LogBytesTotal: info.Size(), LogBytesRead: durableEnd}, false, nil
66 }
67
68 func scanCommitFileCodec(file *os.File, startOffset int64, nextSequence uint64, codec string, knownKinds map[string]bool, visit func(int64, Commit) bool) error {
69 return scanCommitFileCodecBoundaries(context.Background(), file, startOffset, nextSequence, codec, knownKinds, func(start, _ int64, commit Commit) bool {
70 return visit == nil || visit(start, commit)
71 })
72 }
73
74 // Record boundaries come from the original bytes, never re-encoded JSON: old
75 // writers may use different whitespace or field ordering.
76 func scanCommitFileCodecBoundaries(ctx context.Context, file *os.File, startOffset int64, nextSequence uint64, codec string, knownKinds map[string]bool, visit func(int64, int64, Commit) bool) error {
77 if knownKinds == nil {
78 knownKinds = ProjectionKinds
79 if codec == PrototypeCodec {
80 knownKinds = PrototypeProjectionKinds
81 }
82 }
83 if _, err := file.Seek(startOffset, io.SeekStart); err != nil {
84 return err
85 }
86 reader := bufio.NewReaderSize(file, 64<<10)
87 next := nextSequence
88 offset := startOffset
89 operations := map[string]string{}
90 for {
91 if err := ctx.Err(); err != nil {
92 return err
93 }
94 recordOffset := offset
95 line, readErr := reader.ReadBytes('\n')
96 if errors.Is(readErr, io.EOF) {
97 // Cold readers expose only the complete durable prefix. The exclusive
98 // writer path preserves and repairs this tail before accepting work.
99 break
100 }
101 if readErr != nil {
102 return readErr
103 }
104 offset += int64(len(line))
105 var commit Commit
106 if err := json.Unmarshal(bytes.TrimSuffix(line, []byte{'\n'}), &commit); err != nil {
107 return fmt.Errorf("%w: decode complete commit: %w", ErrDamagedStore, err)
108 }
109 if commit.SchemaVersion != 3 || commit.Codec != codec {
110 return fmt.Errorf("%w: event codec", ErrUnsupportedVersion)
111 }
112 if commit.RecordType != "commit" || commit.ID == "" || commit.OperationID == "" ||
113 commit.OperationHash == "" || commit.WriterGeneration == 0 || commit.FirstSequence != next ||
114 commit.EventCount != len(commit.Events) || commit.EventCount == 0 {
115 return fmt.Errorf("%w: invalid commit boundary at sequence %d", ErrDamagedStore, next)
116 }
117 if prior, ok := operations[commit.OperationID]; ok && prior != commit.OperationHash {
118 return fmt.Errorf("%w: conflicting operation %q", ErrDamagedStore, commit.OperationID)
119 }
120 operations[commit.OperationID] = commit.OperationHash
121 for i, event := range commit.Events {
122 if event.Sequence != next+uint64(i) || event.ID == "" || strings.TrimSpace(event.Kind) == "" {
123 return fmt.Errorf("%w: invalid event at sequence %d", ErrDamagedStore, next+uint64(i))
124 }
125 if !event.Optional && !knownKinds[event.Kind] {
126 return fmt.Errorf("%w: unknown required event %q", ErrUnsupportedVersion, event.Kind)
127 }
128 }
129 next = commit.LastSequence() + 1
130 normalizePrototypeHistory(&commit, codec)
131 if visit != nil && !visit(recordOffset, offset, commit) {
132 return nil
133 }
134 }
135 return nil
136 }
137
138 func (opts OpenOptions) openContext() context.Context {
139 if opts.Context != nil {
140 return opts.Context
141 }
142 return context.Background()
143 }
144 func normalizePrototypeHistory(commit *Commit, codec string) {
145 if codec != PrototypeCodec {
146 return
147 }
148 for i := range commit.Events {
149 if commit.Events[i].Kind == "context/replace" {
150 commit.Events[i].Kind = "history/replace"
151 }
152 }
153 }
154 func loadStartupCatalogPreview(dir string, manifest Manifest, externalHistory bool, startup *startupSessionState) {
155 if externalHistory && startup.catalogPreview == "" {
156 revision, revisionErr := revisionOfLog(dir)
157 cacheDir := filepath.Join(filepath.Dir(dir), ".query-cache", filepath.Base(dir))
158 if revisionErr == nil {
159 if metadata, metadataErr := readCatalogMetadata(cacheDir, manifest, revision); metadataErr == nil {
160 startup.catalogPreview = metadata.Preview
161 }
162 }
163 }
164 }
165
165 lines GO