返回 DeepSeek-Reasonix
store.go
根目录 / internal / session / store.go
1 // Package session owns Reasonix's canonical session service.
2 //
3 // Runtime owns live business state and activity authority, PersistenceBinding
4 // owns accepted work and durability progress, Store owns framed bytes and the
5 // writer lease, and Query owns rebuildable disk projections. A commit is
6 // accepted before it is durable; Flush establishes a semantic checkpoint.
7 package session
8
9 import (
10 "context"
11 "crypto/rand"
12 "crypto/sha256"
13 "encoding/hex"
14 "encoding/json"
15 "errors"
16 "fmt"
17 "io"
18 "maps"
19 "os"
20 "path/filepath"
21 "strings"
22 "sync"
23 "time"
24
25 "reasonix/internal/filelock"
26 "reasonix/internal/fileutil"
27 "reasonix/internal/provider"
28 "reasonix/internal/sessioncontent"
29 )
30
31 const (
32 SchemaVersion = V4SchemaVersion
33 // StorageRevision distinguishes the final v4 layout from unpublished v4
34 // drafts. Physical layout changes are migration boundaries even when the
35 // logical codec remains v4.
36 StorageRevision = 3
37 // Codec identifies the current framed linear session format. Earlier linear
38 // and prototype stores are immutable migration inputs.
39 Codec = V4Codec
40 FinalV31Codec = "reasonix.session.linear/v3.1"
41 LegacyLinearCodec = "reasonix.session.linear/v3"
42 PrototypeCodec = "reasonix.session.events/v3"
43 LiveBatchDelay = 200 * time.Millisecond
44 currentLogName = "events.frames"
45 legacyLogName = "events.jsonl"
46 )
47
48 var (
49 ErrUnsupportedVersion = errors.New("unsupported session storage version")
50 ErrDamagedStore = errors.New("damaged session store")
51 ErrStaleGeneration = errors.New("stale session writer generation")
52 ErrOperationConflict = errors.New("session operation id conflicts with an earlier batch")
53 ErrPersistenceUncertain = errors.New("session persistence result is uncertain")
54 ErrSessionNotFound = errors.New("session not found")
55 ErrSessionExists = errors.New("session already exists")
56 ErrWriterOwned = errors.New("session writer is owned by another runtime")
57 ErrReadOnly = errors.New("session handle is read-only")
58 )
59
60 type Manifest struct {
61 SchemaVersion int `json:"schemaVersion"`
62 Codec string `json:"codec"`
63 StorageRevision int `json:"storageRevision,omitempty"`
64 ContentRoot string `json:"contentRoot,omitempty"`
65 SessionID string `json:"sessionId"`
66 CreatedAt time.Time `json:"createdAt"`
67 WriterGeneration uint64 `json:"writerGeneration"`
68 InheritedEvents uint64 `json:"inheritedEventCount,omitempty"`
69 Source *Source `json:"source,omitempty"`
70 Kind SessionKind `json:"kind,omitempty"`
71 }
72
73 const sharedContentRoot = "../.content-v1"
74
75 type Source struct {
76 Path string `json:"path"`
77 Size int64 `json:"size"`
78 SHA256 string `json:"sha256"`
79 Version string `json:"version,omitempty"`
80 LegacyHeadID string `json:"legacyHeadId,omitempty"`
81 }
82
83 type Event struct {
84 ID string `json:"id"`
85 Sequence uint64 `json:"seq"`
86 Kind string `json:"kind"`
87 Optional bool `json:"optional,omitempty"`
88 Required bool `json:"required,omitempty"`
89 Payload json.RawMessage `json:"payload,omitempty"`
90 PayloadRef *sessioncontent.Ref `json:"payloadRef,omitempty"`
91 }
92
93 type Commit struct {
94 SchemaVersion int `json:"schemaVersion"`
95 Codec string `json:"codec"`
96 RecordType string `json:"recordType"`
97 ID string `json:"commitId"`
98 OperationID string `json:"operationId"`
99 OperationHash string `json:"operationHash"`
100 FirstSequence uint64 `json:"firstSeq"`
101 EventCount int `json:"eventCount"`
102 TurnID string `json:"turnId,omitempty"`
103 WriterGeneration uint64 `json:"writerGeneration"`
104 CreatedAt time.Time `json:"createdAt"`
105 Events []Event `json:"events"`
106 }
107
108 func (c Commit) LastSequence() uint64 {
109 if c.EventCount == 0 {
110 return c.FirstSequence
111 }
112 return c.FirstSequence + uint64(c.EventCount) - 1
113 }
114
115 type Batch struct {
116 OperationID string
117 TurnID string
118 Events []Event
119 }
120
121 type DurableReceipt struct {
122 DurableSequence uint64 `json:"durableSequence"`
123 }
124
125 type PersistenceStatus string
126
127 const (
128 PersistenceReady PersistenceStatus = "ready"
129 PersistencePending PersistenceStatus = "pending"
130 PersistenceFailed PersistenceStatus = "failed"
131 PersistenceUncertain PersistenceStatus = "uncertain"
132 )
133
134 type Snapshot struct {
135 EventSequence uint64
136 DurableSequence uint64
137 PersistenceStatus PersistenceStatus
138 PersistenceError string
139 Projection Projection
140 }
141
142 type timerHandle interface{ Stop() bool }
143
144 type OpenOptions struct {
145 Context context.Context
146 AfterFunc func(time.Duration, func()) timerHandle
147 Write func(context.Context, io.Writer, []byte) error
148 Sync func(*os.File) error
149 // ExternalHistory keeps durable UI messages out of the live projection.
150 // Production SessionService enables it; low-level compatibility callers
151 // retain the historical full-projection behavior unless requested.
152 ExternalHistory bool
153 // ObserveRecovery receives bounded open-path I/O counters after recovery.
154 // Capacity tests use it to distinguish checkpoint recovery from prefix replay.
155 ObserveRecovery func(RecoveryOpenStats)
156 // DisableRecoveryPublish is a fault-injection hook used to verify that a
157 // durable log tail is replayed from the previous checkpoint.
158 DisableRecoveryPublish bool
159 }
160
161 type operationRecord struct {
162 hash string
163 commit Commit
164 }
165
166 func compactOperationRecord(commit Commit) operationRecord {
167 metadata := commit
168 metadata.Events = nil
169 return operationRecord{hash: commit.OperationHash, commit: metadata}
170 }
171
172 type uncertainWrite struct {
173 start int64
174 stagedPath string
175 stagedBytes int64
176 commitCount int
177 }
178
179 type uncertainAppendError struct {
180 cause error
181 write uncertainWrite
182 }
183
184 func (e *uncertainAppendError) Error() string {
185 return fmt.Sprintf("%v: append at offset %d may have changed the log: %v", ErrPersistenceUncertain, e.write.start, e.cause)
186 }
187
188 func (e *uncertainAppendError) Unwrap() error { return ErrPersistenceUncertain }
189
190 // Store is the physical framed-log handle for one session directory. It owns the
191 // writer lease, the open file, and the rebuildable sparse offset index. It
192 // deliberately holds no projection, operation table, or accepted commit list:
193 // those belong to Session.
194 type Store struct {
195 mu sync.Mutex
196 indexMu sync.Mutex
197 closeOnce sync.Once
198 closeErr error
199 dir string
200 manifest Manifest
201 file *os.File
202 releaseLease func()
203 closed bool
204 index sparseIndex
205 content *sessioncontent.Store
206 startup *startupSessionState
207 recovery *recoveryStore
208 identity storageIdentity
209 tip durableTip
210
211 writeFn func(context.Context, io.Writer, []byte) error
212 syncFn func(*os.File) error
213 }
214
215 type startupSessionState struct {
216 projection Projection
217 operations map[string]operationRecord
218 durable uint64
219 catalogPreview string
220 recentMessages []provider.Message
221 messageIDs messageIdentities
222 tip durableTip
223 }
224
225 type durableTip struct {
226 LogOffset int64
227 AnchorOffset int64
228 AnchorFirst uint64
229 AnchorCommitID string
230 AnchorHash string
231 }
232
233 // ID returns the immutable session identity of the physical store.
234 func (s *Store) ID() string {
235 return s.SessionID()
236 }
237
238 func (s *Store) SessionID() string {
239 if s == nil {
240 return ""
241 }
242 s.mu.Lock()
243 defer s.mu.Unlock()
244 return s.manifest.SessionID
245 }
246
247 func (s *Store) Manifest() Manifest {
248 if s == nil {
249 return Manifest{}
250 }
251 s.mu.Lock()
252 defer s.mu.Unlock()
253 manifest := s.manifest
254 if manifest.Source != nil {
255 source := *manifest.Source
256 manifest.Source = &source
257 }
258 return manifest
259 }
260
261 // Dir reports the confined directory that backs this handle.
262 func (s *Store) Dir() string {
263 if s == nil {
264 return ""
265 }
266 s.mu.Lock()
267 defer s.mu.Unlock()
268 return s.dir
269 }
270
271 // writableFile returns the open append file for durability repair. The caller
272 // must already hold the writer lease, which the handle acquired at open.
273 func (s *Store) writableFile() (*os.File, error) {
274 if s == nil {
275 return nil, os.ErrClosed
276 }
277 s.mu.Lock()
278 defer s.mu.Unlock()
279 if s.closed || s.file == nil {
280 return nil, os.ErrClosed
281 }
282 return s.file, nil
283 }
284
285 // Open opens an existing legacy-compatible path, creating it when absent. It
286 // returns a live Session: the in-memory log is the caller's business state.
287 func Open(dir, sessionID string) (*Session, error) {
288 return OpenWithOptions(dir, sessionID, OpenOptions{})
289 }
290
291 // OpenWithOptions is the low-level test/import constructor retained while the
292 // old controller adapter is removed. Production callers use
293 // FilesystemPersistence, whose Open is strict and never creates a session.
294 func OpenWithOptions(dir, sessionID string, opts OpenOptions) (*Session, error) {
295 if _, err := os.Stat(filepath.Clean(strings.TrimSpace(dir))); os.IsNotExist(err) {
296 return CreateWithOptions(dir, sessionID, opts)
297 }
298 handle, err := openExistingHandle(dir, sessionID, opts)
299 if err != nil {
300 return nil, err
301 }
302 return bindSession(handle, opts)
303 }
304
305 func CreateStore(dir, sessionID string) (*Session, error) {
306 return CreateWithOptions(dir, sessionID, OpenOptions{})
307 }
308
309 func CreateWithOptions(dir, sessionID string, opts OpenOptions) (*Session, error) {
310 return createWithOptions(dir, sessionID, opts, nil, "")
311 }
312
313 func createWithOptions(dir, sessionID string, opts OpenOptions, header *SessionHeader, kind SessionKind) (*Session, error) {
314 dir = filepath.Clean(strings.TrimSpace(dir))
315 sessionID = strings.TrimSpace(sessionID)
316 if dir == "." || sessionID == "" {
317 return nil, fmt.Errorf("session: directory and session id are required")
318 }
319 if err := validateSessionID(sessionID); err != nil {
320 return nil, err
321 }
322 if err := os.MkdirAll(filepath.Dir(dir), 0o700); err != nil {
323 return nil, err
324 }
325 if err := os.Mkdir(dir, 0o700); err != nil {
326 if os.IsExist(err) {
327 return nil, fmt.Errorf("%w: %s", ErrSessionExists, sessionID)
328 }
329 return nil, err
330 }
331 created := true
332 defer func() {
333 if created {
334 _ = os.RemoveAll(dir)
335 }
336 }()
337 createdAt := time.Now().UTC()
338 if header != nil {
339 header.SessionID = sessionID
340 header.CreatedAt = createdAt
341 }
342 manifest := Manifest{SchemaVersion: SchemaVersion, Codec: Codec, StorageRevision: StorageRevision, ContentRoot: sharedContentRoot, SessionID: sessionID, CreatedAt: createdAt, Kind: kind}
343 if err := writeManifestFile(filepath.Join(dir, "manifest.json"), manifest); err != nil {
344 return nil, err
345 }
346 if header != nil {
347 if err := writeSessionHeader(dir, *header); err != nil {
348 return nil, err
349 }
350 }
351 if err := fileutil.AtomicWriteFileStrict(filepath.Join(dir, currentLogName), nil, 0o600); err != nil {
352 return nil, err
353 }
354 created = false
355 return OpenWithOptions(dir, sessionID, opts)
356 }
357
358 func openExistingHandle(dir, sessionID string, opts OpenOptions) (*Store, error) {
359 ctx := opts.openContext()
360 if err := ctx.Err(); err != nil {
361 return nil, err
362 }
363 dir = filepath.Clean(strings.TrimSpace(dir))
364 sessionID = strings.TrimSpace(sessionID)
365 if dir == "." || sessionID == "" {
366 return nil, fmt.Errorf("session: directory and session id are required")
367 }
368 if err := validateSessionID(sessionID); err != nil {
369 return nil, err
370 }
371 info, err := os.Stat(dir)
372 if os.IsNotExist(err) {
373 return nil, fmt.Errorf("%w: %s", ErrSessionNotFound, sessionID)
374 }
375 if err != nil {
376 return nil, err
377 }
378 if !info.IsDir() {
379 return nil, fmt.Errorf("session: session path is not a directory: %s", dir)
380 }
381 releaseLease, err := acquireSessionWriter(dir)
382 if err != nil {
383 if errors.Is(err, filelock.ErrHeld) {
384 return nil, fmt.Errorf("%w: %s", ErrWriterOwned, sessionID)
385 }
386 return nil, err
387 }
388 fail := func(err error) (*Store, error) {
389 releaseLease()
390 return nil, err
391 }
392 manifestPath := filepath.Join(dir, "manifest.json")
393 manifest, err := readStoredManifest(manifestPath)
394 if err != nil {
395 return fail(err)
396 }
397 if manifest.SessionID != sessionID {
398 return fail(fmt.Errorf("session: manifest belongs to %q", manifest.SessionID))
399 }
400 if !supportedStoredManifest(manifest) {
401 return fail(fmt.Errorf("%w: manifest schema or codec", ErrUnsupportedVersion))
402 }
403 eventsPath := logPathForManifest(dir, manifest)
404 identity, err := ensureStorageIdentity(dir, manifest)
405 if err != nil {
406 return fail(err)
407 }
408 recovery, err := openRecoveryStoreRepair(dir, identity)
409 if err != nil {
410 return fail(err)
411 }
412 failRecovery := func(err error) (*Store, error) {
413 _ = recovery.close()
414 return fail(err)
415 }
416 // Build runtime state in one streaming validation pass. A newer required
417 // event or a damaged complete batch leaves the original tail untouched.
418 startup, durableEnd, torn, stats, usedCheckpoint, err := loadStartupSessionStateForManifest(ctx, dir, eventsPath, manifest, opts.ExternalHistory, recovery, identity)
419 if opts.ObserveRecovery != nil {
420 opts.ObserveRecovery(stats)
421 }
422 if err != nil {
423 return failRecovery(err)
424 }
425 if err := ctx.Err(); err != nil {
426 return failRecovery(err)
427 }
428 loadStartupCatalogPreview(dir, manifest, opts.ExternalHistory, startup)
429 if torn {
430 // Cold readers deliberately stop at the last complete record. A writer
431 // may repair that tail only after acquiring the exclusive lease above:
432 // preserve the original bytes first, then truncate back to the durable
433 // commit boundary. This never invents or partially replays an event.
434 if _, err := preserveAndTruncateTail(eventsPath, durableEnd, "torn"); err != nil {
435 return failRecovery(fmt.Errorf("recover torn v3 tail: %w", err))
436 }
437 }
438 if !usedCheckpoint {
439 checkpoint := checkpointFromStartup(manifest, identity, startup)
440 if err := recovery.publish(context.Background(), checkpoint, startup.operations); err == nil {
441 startup.operations = map[string]operationRecord{}
442 }
443 }
444 // Keep the source codec stable. Legacy stored sessions are writable through
445 // their original line codec; explicit conversion is the only codec change.
446 if manifest.Codec == Codec && manifest.StorageRevision > 0 {
447 manifest.StorageRevision = StorageRevision
448 }
449 manifest.WriterGeneration++
450 if err := writeManifestFile(manifestPath, manifest); err != nil {
451 return failRecovery(err)
452 }
453 f, err := os.OpenFile(eventsPath, os.O_CREATE|os.O_RDWR|os.O_APPEND, 0o600)
454 if err != nil {
455 return failRecovery(err)
456 }
457 writeFn := opts.Write
458 if writeFn == nil {
459 writeFn = writeAllContext
460 }
461 syncFn := opts.Sync
462 if syncFn == nil {
463 syncFn = func(file *os.File) error { return file.Sync() }
464 }
465 // Opening a runtime must not rebuild a historical seek index. The writer
466 // needs only the durable sequence and byte end; readers maintain their own
467 // disposable locator in the query cache.
468 index := sparseIndex{Codec: sparseIndexCodec, LogSize: durableEnd, LastSequence: startup.durable, partial: startup.durable > 0}
469 return &Store{
470 dir: dir, manifest: manifest, file: f, releaseLease: releaseLease,
471 index: index, content: contentStoreForSessionDir(dir), startup: startup,
472 recovery: recovery, identity: identity, tip: startup.tip,
473 writeFn: writeFn, syncFn: syncFn,
474 }, nil
475 }
476
477 func loadStartupSessionState(ctx context.Context, dir, eventsPath string, externalHistory bool, recovery *recoveryStore, identity storageIdentity) (*startupSessionState, int64, bool, RecoveryOpenStats, bool, error) {
478 file, err := os.Open(eventsPath)
479 if os.IsNotExist(err) {
480 projection, _ := Project(nil)
481 return &startupSessionState{projection: projection, operations: map[string]operationRecord{}}, 0, false, RecoveryOpenStats{}, false, nil
482 }
483 if err != nil {
484 return nil, 0, false, RecoveryOpenStats{}, false, err
485 }
486 defer file.Close()
487 info, err := file.Stat()
488 if err != nil {
489 return nil, 0, false, RecoveryOpenStats{}, false, err
490 }
491 if state, end, torn, stats, ok := loadRecoveryStartupState(ctx, dir, file, info, recovery, identity); ok {
492 return state, end, torn, stats, true, nil
493 }
494 stats := RecoveryOpenStats{LogBytesTotal: info.Size()}
495 if externalHistory {
496 state, end, torn, err := loadBoundedStartupSessionState(ctx, dir, file, info)
497 stats.LogBytesRead = end
498 return state, end, torn, stats, false, err
499 }
500 projection, _ := Project(nil)
501 state := &startupSessionState{projection: projection, operations: map[string]operationRecord{}}
502 var durableEnd int64
503 var projectionErr error
504 err = scanV4CommitFile(ctx, file, 0, 1, contentStoreForSessionDir(dir), nil, func(offset int64, commit Commit) bool {
505 if applyErr := applyProjectionCommit(&state.projection, commit); applyErr != nil {
506 projectionErr = applyErr
507 return false
508 }
509 if externalHistory {
510 for _, message := range state.projection.Messages {
511 if state.catalogPreview == "" {
512 state.catalogPreview = catalogMessagePreview(message)
513 }
514 }
515 state.projection.Messages = nil
516 }
517 state.operations[commit.OperationID] = compactOperationRecord(commit)
518 state.durable = commit.LastSequence()
519 durableEnd, _ = file.Seek(0, io.SeekCurrent)
520 state.tip = durableTip{LogOffset: durableEnd, AnchorOffset: offset, AnchorFirst: commit.FirstSequence, AnchorCommitID: commit.ID, AnchorHash: commit.OperationHash}
521 if err := applyRecentCommit(&state.recentMessages, commit); err != nil {
522 projectionErr = err
523 return false
524 }
525 state.messageIDs.admit(commit)
526 return true
527 })
528 if err != nil {
529 return nil, 0, false, stats, false, err
530 }
531 if projectionErr != nil {
532 return nil, 0, false, stats, false, projectionErr
533 }
534 stats.LogBytesRead = durableEnd
535 return state, durableEnd, durableEnd < info.Size(), stats, false, nil
536 }
537
538 // loadBoundedStartupSessionState separates lightweight business recovery from
539 // model-context recovery. The first pass validates every transaction without
540 // resolving historical message bodies and locates the newest context reset.
541 // The second pass materializes only the provider workset after that reset.
542 func loadBoundedStartupSessionState(ctx context.Context, dir string, file *os.File, info os.FileInfo) (*startupSessionState, int64, bool, error) {
543 projection, _ := Project(nil)
544 state := &startupSessionState{projection: projection, operations: map[string]operationRecord{}}
545 modelOffset, modelSequence := int64(0), uint64(1)
546 var durableEnd int64
547 var projectionErr error
548 sawModelEvent := false
549 repeated := map[uint64]bool{} // both passes keep an id's first message
550 content := contentStoreForSessionDir(dir)
551 err := scanV4CommitFileRefs(ctx, file, 0, 1, content, nil, func(offset int64, commit Commit) bool {
552 business := commit
553 business.Events = nil
554 for _, event := range commit.Events {
555 if modelProjectionEvent(event.Kind) {
556 sawModelEvent = true
557 supersedeStreamCheckpoint(&state.projection, event.Kind)
558 if modelProjectionReset(event.Kind) {
559 modelOffset, modelSequence = offset, commit.FirstSequence
560 }
561 continue
562 }
563 resolved, err := resolveProjectionEvent(ctx, content, event)
564 if err != nil {
565 projectionErr = err
566 return false
567 }
568 business.Events = append(business.Events, resolved)
569 }
570 if err := applyProjectionCommit(&state.projection, business); err != nil {
571 projectionErr = err
572 return false
573 }
574 recent := commit
575 recent.Events = nil
576 for _, event := range commit.Events {
577 switch event.Kind {
578 case "message/complete", "message/upsert", "message/retract", "history/replace", "legacy/import":
579 resolved, err := resolveProjectionEvent(ctx, content, event)
580 if err != nil {
581 projectionErr = err
582 return false
583 }
584 if !state.messageIDs.admitEvent(resolved, repeated) {
585 continue
586 }
587 recent.Events = append(recent.Events, resolved)
588 applyTranscriptMetadata(&state.projection, commit, resolved)
589 }
590 }
591 if err := applyRecentCommit(&state.recentMessages, recent); err != nil {
592 projectionErr = err
593 return false
594 }
595 state.operations[commit.OperationID] = compactOperationRecord(commit)
596 state.durable = commit.LastSequence()
597 durableEnd, _ = file.Seek(0, io.SeekCurrent)
598 state.tip = durableTip{LogOffset: durableEnd, AnchorOffset: offset, AnchorFirst: commit.FirstSequence, AnchorCommitID: commit.ID, AnchorHash: commit.OperationHash}
599 return true
600 })
601 if err != nil {
602 return nil, 0, false, err
603 }
604 if projectionErr != nil {
605 return nil, 0, false, projectionErr
606 }
607 if sawModelEvent {
608 inputs, hidden, retracted := state.projection.TranscriptInputs, state.projection.HiddenTurns, state.projection.RetractedInputs
609 // The first pass saw canonical results even before the latest context
610 // reset. Replaying only that reset's wire view cannot recover their state.
611 rejected := maps.Clone(state.projection.RejectedToolResults)
612 checkpoint := state.projection.StreamCheckpoint
613 if err := loadCurrentModelProjection(ctx, file, content, state, modelOffset, modelSequence, repeated); err != nil {
614 return nil, 0, false, err
615 }
616 state.projection.TranscriptInputs, state.projection.HiddenTurns = inputs, hidden
617 state.projection.RetractedInputs = retracted
618 state.projection.RejectedToolResults = rejected
619 state.projection.StreamCheckpoint = checkpoint
620 }
621 state.projection.Messages = nil
622 state.projection.CommittedSequence = state.durable
623 return state, durableEnd, durableEnd < info.Size(), nil
624 }
625
626 func loadCurrentModelProjection(ctx context.Context, file *os.File, content *sessioncontent.Store, state *startupSessionState, offset int64, sequence uint64, repeated map[uint64]bool) error {
627 var projectionErr error
628 err := scanV4CommitFileRefs(ctx, file, offset, sequence, content, nil, func(_ int64, commit Commit) bool {
629 model := commit
630 model.Events = nil
631 for _, event := range commit.Events {
632 if !modelProjectionEvent(event.Kind) || repeated[event.Sequence] {
633 continue
634 }
635 resolved, err := resolveProjectionEvent(ctx, content, event)
636 if err != nil {
637 projectionErr = err
638 return false
639 }
640 model.Events = append(model.Events, resolved)
641 }
642 if err := applyProjectionCommit(&state.projection, model); err != nil {
643 projectionErr = err
644 return false
645 }
646 return true
647 })
648 if err != nil {
649 return err
650 }
651 return projectionErr
652 }
653
654 func resolveProjectionEvent(ctx context.Context, content *sessioncontent.Store, event Event) (Event, error) {
655 if event.PayloadRef == nil {
656 return event, nil
657 }
658 payload, err := resolveContentPayload(ctx, content, *event.PayloadRef)
659 if err != nil {
660 return Event{}, fmt.Errorf("%w: read v4 event %s payload: %w", ErrDamagedStore, event.ID, err)
661 }
662 event.Payload, event.PayloadRef = payload, nil
663 return event, nil
664 }
665
666 func modelProjectionEvent(kind string) bool {
667 switch kind {
668 case "message/complete", "message/upsert", "message/retract", "history/replace", "model/context-replace", "compaction", "legacy/import":
669 return true
670 default:
671 return false
672 }
673 }
674
675 func modelProjectionReset(kind string) bool {
676 switch kind {
677 case "history/replace", "model/context-replace", "compaction", "legacy/import":
678 return true
679 default:
680 return false
681 }
682 }
683
684 // bindSession replays the durable prefix and constructs the in-memory Session
685 // over a binding for the exact handle.
686 func bindSession(handle *Store, opts OpenOptions) (*Session, error) {
687 dir := handle.Dir()
688 state := handle.startup
689 if state == nil {
690 projection, _ := Project(nil)
691 state = &startupSessionState{projection: projection, operations: map[string]operationRecord{}}
692 }
693 handle.startup = nil
694 binding := newPersistenceBinding(handle, dir, state.durable, opts)
695 session := newSession(handle.Manifest().SessionID, handle.Manifest(), nil, state.projection, binding)
696 session.next = state.durable + 1
697 session.operations = state.operations
698 session.externalHistory = opts.ExternalHistory
699 session.catalogPreview = state.catalogPreview
700 session.recentMessages = detachMessages(state.recentMessages)
701 session.durableRecent = detachMessages(state.recentMessages)
702 session.messageIDs = state.messageIDs
703 session.storageGeneration = handle.identity.Generation
704 session.recovery = handle.recovery
705 binding.metadataSource = session.metadataForDurable
706 binding.recoverySource = session.recoveryForDurable
707 binding.recoveryPublished = session.recoveryPublished
708 binding.disableRecoveryPublish = opts.DisableRecoveryPublish
709 return session, nil
710 }
711
712 func writeManifestFile(path string, manifest Manifest) error {
713 b, err := json.MarshalIndent(manifest, "", " ")
714 if err != nil {
715 return err
716 }
717 return fileutil.AtomicWriteFileStrict(path, append(b, '\n'), 0o600)
718 }
719
720 func contentStoreForSessionDir(dir string) *sessioncontent.Store {
721 root := filepath.Join(filepath.Dir(dir), ".content-v1")
722 if data, err := os.ReadFile(filepath.Join(dir, "manifest.json")); err == nil {
723 var location struct {
724 ContentRoot string `json:"contentRoot"`
725 }
726 if json.Unmarshal(data, &location) == nil && strings.TrimSpace(location.ContentRoot) != "" {
727 candidate := filepath.Clean(filepath.Join(dir, location.ContentRoot))
728 relative, relErr := filepath.Rel(dir, candidate)
729 if relErr == nil && (relative == ".content-v1" || relative == sharedContentRoot) {
730 root = candidate
731 }
732 }
733 }
734 return sessioncontent.New(root)
735 }
736
737 func logPathForManifest(dir string, manifest Manifest) string {
738 if manifest.Codec == Codec {
739 return filepath.Join(dir, currentLogName)
740 }
741 return filepath.Join(dir, legacyLogName)
742 }
743
744 func supportedStoredManifest(manifest Manifest) bool {
745 if currentStoredManifest(manifest) {
746 return true
747 }
748 if manifest.SchemaVersion == SchemaVersion && manifest.Codec == Codec && manifest.StorageRevision == 0 {
749 return true
750 }
751 return manifest.SchemaVersion == 3 &&
752 (manifest.Codec == FinalV31Codec || manifest.Codec == LegacyLinearCodec || manifest.Codec == PrototypeCodec)
753 }
754
755 func currentStoredManifest(manifest Manifest) bool {
756 return manifest.SchemaVersion == SchemaVersion && manifest.Codec == Codec &&
757 manifest.StorageRevision >= 1 && manifest.StorageRevision <= StorageRevision
758 }
759
760 func readStoredManifest(path string) (Manifest, error) {
761 b, err := os.ReadFile(path)
762 if err != nil {
763 return Manifest{}, err
764 }
765 var manifest Manifest
766 if err := json.Unmarshal(b, &manifest); err != nil {
767 return Manifest{}, err
768 }
769 if !supportedStoredManifest(manifest) {
770 return Manifest{}, fmt.Errorf("%w: manifest schema or codec", ErrUnsupportedVersion)
771 }
772 return manifest, nil
773 }
774
775 // Append writes committed batches in order. Session owns sequence allocation,
776 // validation, and idempotency; this method reports physical write uncertainty.
777 func (s *Store) Append(ctx context.Context, commits []Commit) error {
778 if s == nil {
779 return fmt.Errorf("session: nil store")
780 }
781 if err := ctx.Err(); err != nil {
782 return err
783 }
784 if len(commits) == 0 {
785 return nil
786 }
787 s.mu.Lock()
788 defer s.mu.Unlock()
789 if s.closed || s.file == nil {
790 return os.ErrClosed
791 }
792 s.indexMu.Lock()
793 next := s.index.LastSequence + 1
794 s.indexMu.Unlock()
795 wantSchema, wantCodec := SchemaVersion, Codec
796 if s.manifest.Codec != Codec {
797 wantSchema, wantCodec = 3, s.manifest.Codec
798 }
799 for i, commit := range commits {
800 if commit.SchemaVersion != wantSchema || commit.Codec != wantCodec || commit.RecordType != "commit" ||
801 commit.ID == "" || commit.OperationID == "" || commit.OperationHash == "" ||
802 commit.WriterGeneration != s.manifest.WriterGeneration || commit.FirstSequence != next ||
803 commit.EventCount == 0 || commit.EventCount != len(commit.Events) {
804 return fmt.Errorf("%w: invalid physical commit %d at sequence %d", ErrDamagedStore, i, next)
805 }
806 for eventIndex, event := range commit.Events {
807 want := commit.FirstSequence + uint64(eventIndex)
808 if event.Sequence != want || event.ID == "" || strings.TrimSpace(event.Kind) == "" {
809 return fmt.Errorf("%w: invalid physical event at sequence %d", ErrDamagedStore, want)
810 }
811 }
812 next = commit.LastSequence() + 1
813 }
814 return s.persist(ctx, s.file, commits)
815 }
816
817 func (s *Store) persist(ctx context.Context, file *os.File, commits []Commit) error {
818 staged, err := os.CreateTemp(s.dir, ".append-*.staged")
819 if err != nil {
820 return fmt.Errorf("stage v4 append: %w", err)
821 }
822 stagedPath := staged.Name()
823 keepStaged := false
824 defer func() {
825 _ = staged.Close()
826 if !keepStaged {
827 _ = os.Remove(stagedPath)
828 }
829 }()
830 var lengths []int64
831 if s.manifest.Codec == Codec {
832 lengths, err = encodeV4Commits(ctx, staged, s.content, commits)
833 } else {
834 lengths, err = encodeLegacyCommits(ctx, staged, commits)
835 }
836 if err != nil {
837 return err
838 }
839 if err := staged.Sync(); err != nil {
840 return fmt.Errorf("fsync staged v4 append: %w", err)
841 }
842 stagedInfo, err := staged.Stat()
843 if err != nil {
844 return err
845 }
846 if _, err := staged.Seek(0, io.SeekStart); err != nil {
847 return err
848 }
849 start, err := file.Seek(0, io.SeekEnd)
850 if err != nil {
851 return err
852 }
853 s.adoptPersistedIndex(file, start)
854 if err := copyStagedAppend(ctx, staged, file, s.writeFn); err != nil {
855 end, statErr := file.Seek(0, io.SeekEnd)
856 if statErr == nil && end == start {
857 return err
858 }
859 if statErr == nil && end == start+stagedInfo.Size() {
860 if syncErr := s.syncFn(file); syncErr == nil {
861 s.recordPersistedIndex(file, start, commits, lengths)
862 return nil
863 }
864 }
865 keepStaged = true
866 return &uncertainAppendError{cause: err, write: uncertainWrite{start: start, stagedPath: stagedPath, stagedBytes: stagedInfo.Size(), commitCount: len(commits)}}
867 }
868 if err := s.syncFn(file); err != nil {
869 keepStaged = true
870 return &uncertainAppendError{cause: fmt.Errorf("fsync: %w", err), write: uncertainWrite{start: start, stagedPath: stagedPath, stagedBytes: stagedInfo.Size(), commitCount: len(commits)}}
871 }
872 s.recordPersistedIndex(file, start, commits, lengths)
873 return nil
874 }
875
876 func encodeLegacyCommits(ctx context.Context, dst io.Writer, commits []Commit) ([]int64, error) {
877 lengths := make([]int64, 0, len(commits))
878 for _, commit := range commits {
879 if commit.Codec == PrototypeCodec {
880 commit = cloneCommit(commit)
881 for i := range commit.Events {
882 if commit.Events[i].Kind == "history/replace" {
883 commit.Events[i].Kind = "context/replace"
884 }
885 }
886 }
887 if err := ctx.Err(); err != nil {
888 return nil, err
889 }
890 encoded, err := json.Marshal(commit)
891 if err != nil {
892 return nil, err
893 }
894 encoded = append(encoded, '\n')
895 if _, err := dst.Write(encoded); err != nil {
896 return nil, err
897 }
898 lengths = append(lengths, int64(len(encoded)))
899 }
900 return lengths, nil
901 }
902
903 func copyStagedAppend(ctx context.Context, source io.Reader, destination io.Writer, writeFn func(context.Context, io.Writer, []byte) error) error {
904 buffer := make([]byte, 1<<20)
905 for {
906 if err := ctx.Err(); err != nil {
907 return err
908 }
909 n, readErr := source.Read(buffer)
910 if n > 0 {
911 if err := writeFn(ctx, destination, buffer[:n]); err != nil {
912 return err
913 }
914 }
915 if readErr != nil {
916 if errors.Is(readErr, io.EOF) {
917 return nil
918 }
919 return readErr
920 }
921 }
922 }
923
924 // Sync fsyncs the physical log and reports the durable sequence observed on
925 // disk. It is the handle half of a semantic checkpoint; PersistenceBinding
926 // pairs it with queue drain.
927 func (s *Store) Sync(ctx context.Context) (DurableReceipt, error) {
928 if s == nil {
929 return DurableReceipt{}, os.ErrClosed
930 }
931 if err := ctx.Err(); err != nil {
932 return DurableReceipt{}, err
933 }
934 s.mu.Lock()
935 defer s.mu.Unlock()
936 if s.closed || s.file == nil {
937 return DurableReceipt{}, os.ErrClosed
938 }
939 if err := s.syncFn(s.file); err != nil {
940 return DurableReceipt{}, err
941 }
942 s.indexMu.Lock()
943 sequence := s.index.LastSequence
944 s.indexMu.Unlock()
945 return DurableReceipt{DurableSequence: sequence}, nil
946 }
947
948 func (s *Store) Close(_ context.Context) error {
949 if s == nil {
950 return nil
951 }
952 s.closeOnce.Do(func() {
953 s.mu.Lock()
954 file := s.file
955 recovery := s.recovery
956 s.file = nil
957 s.recovery = nil
958 s.closed = true
959 releaseLease := s.releaseLease
960 s.releaseLease = nil
961 s.mu.Unlock()
962 var closeErr error
963 if file != nil {
964 closeErr = file.Close()
965 }
966 if recovery != nil {
967 closeErr = errors.Join(closeErr, recovery.close())
968 }
969 if releaseLease != nil {
970 releaseLease()
971 }
972 s.closeErr = closeErr
973 })
974 return s.closeErr
975 }
976
977 func writeAllContext(ctx context.Context, w io.Writer, data []byte) error {
978 for len(data) > 0 {
979 if err := ctx.Err(); err != nil {
980 return err
981 }
982 n, err := w.Write(data)
983 if err != nil {
984 return err
985 }
986 if n == 0 {
987 return io.ErrShortWrite
988 }
989 data = data[n:]
990 }
991 return nil
992 }
993
994 // Replay returns the complete durable prefix. It ignores only an unterminated
995 // final record; a later write owner preserves and repairs that tail after it
996 // acquires the exclusive lease.
997 func Replay(dir string, knownKinds map[string]bool) ([]Commit, error) {
998 commits := []Commit{}
999 err := scanDurableCommits(dir, knownKinds, func(commit Commit) bool {
1000 commits = append(commits, commit)
1001 return true
1002 })
1003 return commits, err
1004 }
1005
1006 // scanDurableCommits validates records in sequence and lets paged readers stop
1007 // without materializing the rest of a large log. The next page resumes from a
1008 // sequence cursor; a rebuildable offset index can optimize seeking without
1009 // changing this validation contract.
1010 func scanDurableCommits(dir string, knownKinds map[string]bool, visit func(Commit) bool) error {
1011 manifest, err := readStoredManifest(filepath.Join(dir, "manifest.json"))
1012 if err != nil {
1013 return err
1014 }
1015 file, err := os.Open(logPathForManifest(dir, manifest))
1016 if os.IsNotExist(err) {
1017 return nil
1018 }
1019 if err != nil {
1020 return err
1021 }
1022 defer file.Close()
1023 adapter := func(_ int64, commit Commit) bool {
1024 if visit == nil {
1025 return true
1026 }
1027 return visit(commit)
1028 }
1029 if manifest.Codec == Codec {
1030 return scanV4CommitFile(context.Background(), file, 0, 1, contentStoreForSessionDir(dir), knownKinds, adapter)
1031 }
1032 return scanCommitFileCodec(file, 0, 1, manifest.Codec, knownKinds, adapter)
1033 }
1034
1035 func preserveAndTruncateTail(path string, cut int64, label string) (string, error) {
1036 input, err := os.Open(path)
1037 if err != nil {
1038 return "", err
1039 }
1040 info, err := input.Stat()
1041 if err != nil {
1042 _ = input.Close()
1043 return "", err
1044 }
1045 if cut < 0 || cut > info.Size() {
1046 _ = input.Close()
1047 return "", fmt.Errorf("invalid durable tail offset %d for %d-byte log", cut, info.Size())
1048 }
1049 backup := filepath.Join(filepath.Dir(path), fmt.Sprintf("events.%s-%d.tail", label, time.Now().UTC().UnixNano()))
1050 out, err := os.OpenFile(backup, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o600)
1051 if err != nil {
1052 _ = input.Close()
1053 return "", fmt.Errorf("preserve original tail: %w", err)
1054 }
1055 _, seekErr := input.Seek(cut, io.SeekStart)
1056 _, copyErr := io.CopyBuffer(out, input, make([]byte, 1<<20))
1057 syncErr := out.Sync()
1058 closeOutErr := out.Close()
1059 closeInErr := input.Close()
1060 if err := errors.Join(seekErr, copyErr, syncErr, closeOutErr, closeInErr); err != nil {
1061 _ = os.Remove(backup)
1062 return "", fmt.Errorf("preserve original tail: %w", err)
1063 }
1064 writable, err := os.OpenFile(path, os.O_RDWR, 0o600)
1065 if err != nil {
1066 return "", err
1067 }
1068 truncateErr := writable.Truncate(cut)
1069 syncErr = writable.Sync()
1070 closeErr := writable.Close()
1071 if err := errors.Join(truncateErr, syncErr, closeErr); err != nil {
1072 return "", fmt.Errorf("truncate to durable prefix: %w", err)
1073 }
1074 return backup, nil
1075 }
1076
1077 func hashOperation(sessionID, turnID string, events []Event) (string, error) {
1078 type inputEvent struct {
1079 ID string `json:"id,omitempty"`
1080 Kind string `json:"kind"`
1081 Optional bool `json:"optional,omitempty"`
1082 Payload json.RawMessage `json:"payload,omitempty"`
1083 }
1084 inputs := make([]inputEvent, len(events))
1085 for i, event := range events {
1086 inputs[i] = inputEvent{ID: event.ID, Kind: strings.TrimSpace(event.Kind), Optional: event.Optional, Payload: event.Payload}
1087 }
1088 b, err := json.Marshal(struct {
1089 SessionID string `json:"sessionId"`
1090 TurnID string `json:"turnId"`
1091 Events []inputEvent `json:"events"`
1092 }{SessionID: sessionID, TurnID: turnID, Events: inputs})
1093 if err != nil {
1094 return "", err
1095 }
1096 sum := sha256.Sum256(b)
1097 return hex.EncodeToString(sum[:]), nil
1098 }
1099
1100 func deterministicID(seed string) string {
1101 sum := sha256.Sum256([]byte(seed))
1102 return hex.EncodeToString(sum[:16])
1103 }
1104
1105 func batchContains(events []Event, kind string) bool {
1106 for _, event := range events {
1107 if strings.TrimSpace(event.Kind) == kind {
1108 return true
1109 }
1110 }
1111 return false
1112 }
1113
1114 func sameCommitPrefix(all, prefix []Commit) bool {
1115 for i := range prefix {
1116 if all[i].ID != prefix[i].ID {
1117 return false
1118 }
1119 }
1120 return true
1121 }
1122
1123 func cloneEvents(events []Event) []Event {
1124 out := make([]Event, len(events))
1125 copy(out, events)
1126 for i := range out {
1127 out[i].Payload = append(json.RawMessage(nil), out[i].Payload...)
1128 if out[i].PayloadRef != nil {
1129 ref := *out[i].PayloadRef
1130 out[i].PayloadRef = &ref
1131 }
1132 }
1133 return out
1134 }
1135
1136 func cloneCommit(commit Commit) Commit {
1137 commit.Events = cloneEvents(commit.Events)
1138 return commit
1139 }
1140
1141 func cloneCommits(commits []Commit) []Commit {
1142 out := make([]Commit, len(commits))
1143 for i := range commits {
1144 out[i] = cloneCommit(commits[i])
1145 }
1146 return out
1147 }
1148
1149 func randomID() string {
1150 var b [16]byte
1151 if _, err := rand.Read(b[:]); err != nil {
1152 return fmt.Sprintf("%d", time.Now().UnixNano())
1153 }
1154 return hex.EncodeToString(b[:])
1155 }
1156
1157 func errorString(err error) string {
1158 if err == nil {
1159 return ""
1160 }
1161 return err.Error()
1162 }
1163
1163 lines GO