返回 DeepSeek-Reasonix
session_recovery_publish.go
根目录 / internal / session / session_recovery_publish.go
1 package session
2
3 type recoveryPublishState struct {
4 checkpoint recoveryCheckpoint
5 operations map[string]operationRecord
6 }
7
8 func (s *Session) recoveryForDurable(durable uint64) (recoveryPublishState, bool) {
9 if s == nil || s.recovery == nil {
10 return recoveryPublishState{}, false
11 }
12 s.mu.Lock()
13 defer s.mu.Unlock()
14 if durable+1 != s.next {
15 return recoveryPublishState{}, false
16 }
17 store, _ := s.binding.handle.(*Store)
18 if store == nil {
19 return recoveryPublishState{}, false
20 }
21 store.mu.Lock()
22 tip := store.tip
23 store.mu.Unlock()
24 projection := cloneProjection(s.projection)
25 projection.Messages = nil
26 checkpoint := recoveryCheckpoint{
27 Version: recoveryFormatVersion, SessionID: s.id, StorageGeneration: s.storageGeneration,
28 StorageRevision: StorageRevision, DurableSequence: durable, LogOffset: tip.LogOffset,
29 AnchorOffset: tip.AnchorOffset, AnchorFirst: tip.AnchorFirst,
30 AnchorCommitID: tip.AnchorCommitID, AnchorHash: tip.AnchorHash,
31 ProjectionVersion: recoveryProjectionVersion, Projection: projection,
32 RecentMessages: detachMessages(s.recentMessages), CatalogPreview: s.catalogPreview,
33 MessageIDs: s.messageIDs.list(),
34 }
35 operations := make(map[string]operationRecord, len(s.commits))
36 for _, commit := range s.commits {
37 if commit.LastSequence() <= durable {
38 operations[commit.OperationID] = compactOperationRecord(commit)
39 }
40 }
41 return recoveryPublishState{checkpoint: checkpoint, operations: operations}, true
42 }
43
44 func (s *Session) recoveryPublished(durable uint64) {
45 if s == nil {
46 return
47 }
48 s.mu.Lock()
49 defer s.mu.Unlock()
50 if durable+1 == s.next {
51 s.durableRecent = detachMessages(s.recentMessages)
52 }
53 cut := 0
54 for cut < len(s.commits) && s.commits[cut].LastSequence() <= durable {
55 delete(s.operations, s.commits[cut].OperationID)
56 cut++
57 }
58 if cut > 0 {
59 s.commits = append([]Commit(nil), s.commits[cut:]...)
60 }
61 if s.externalHistory {
62 s.projection.Messages = nil
63 }
64 }
65
65 lines GO