返回 DeepSeek-Reasonix
transcript_restore.go
根目录 / internal / control / transcript_restore.go
1 package control
2
3 import (
4 "crypto/rand"
5 "errors"
6 "path/filepath"
7
8 "reasonix/internal/agent"
9 "reasonix/internal/provider"
10 "reasonix/internal/transcript"
11 "reasonix/internal/turnevent"
12 )
13
14 func (c *Controller) restoreTranscriptProjection(sessionPath string, ledger *turnevent.Ledger) (*transcript.Projection, error) {
15 identity := transcript.Identity{SessionID: agent.BranchID(sessionPath), RuntimeEpoch: rand.Text()}
16 ledger.SetRuntimeEpoch(identity.RuntimeEpoch)
17 var messages []provider.Message
18 if c.executor != nil && c.executor.Session() != nil {
19 messages, identity.HeadID, identity.RewriteEpoch = c.executor.Session().DisplayBaseline()
20 }
21 latest, _ := ledger.ProjectionCursor()
22 p, restored, err := restoreFromTranscriptCheckpoint(sessionPath, ledger, identity, messages, latest)
23 if err != nil || restored {
24 return p, err
25 }
26 dir := c.sessionDir
27 if dir == "" && sessionPath != "" {
28 dir = filepath.Dir(sessionPath)
29 }
30 legacy, err := transcript.LoadLegacyDisplays(dir, sessionPath)
31 if err != nil {
32 return nil, err
33 }
34 pending := ledger.PendingProjections()
35 users := make([]provider.Message, 0)
36 for _, m := range messages {
37 if agent.IsUserAuthoredTurnMessage(m) {
38 users = append(users, m)
39 }
40 }
41 installedTurns := make(map[string]bool)
42 for _, turn := range legacy.Turns {
43 if turn.TurnID != "" {
44 installedTurns[turn.TurnID] = true
45 }
46 }
47 covered, prefixEnd := latest, len(messages)
48 for i, group := range pending {
49 userID := ""
50 for _, envelope := range group.Events {
51 if envelope.Kind == "user_message" && (envelope.Source == "" || envelope.Source == "executor") {
52 userID = envelope.Event.MessageID
53 break
54 }
55 }
56 if userID != "" {
57 // The identified user event is the start of the not-yet-projected
58 // suffix. Do not seed its already-autosaved assistant deltas twice.
59 for index, m := range messages {
60 if m.ID == userID {
61 prefixEnd = min(prefixEnd, index)
62 break
63 }
64 }
65 if len(group.Events) > 0 {
66 covered = min(covered, group.Events[0].Sequence-1)
67 }
68 continue
69 }
70 if installedTurns[group.TurnID] {
71 continue
72 }
73 rows := transcript.PendingDisplayMessages(group, transcript.Formatter{})
74 if len(rows) == 0 {
75 continue
76 }
77 userIndex := len(users) - len(pending) + i
78 if userIndex < 0 || userIndex >= len(users) {
79 return nil, errors.New("legacy display recovery has no matching user boundary")
80 }
81 legacy.Turns = append(legacy.Turns, transcript.LegacyDisplayTurn{TurnID: group.TurnID, UserMessageID: users[userIndex].ID, Messages: rows})
82 }
83 rows := transcript.History(messages[:prefixEnd], transcript.HistoryOptions{LegacyTurns: legacy.Turns, CheckpointTurns: c.CheckpointTurnsByMessageIndex(), SubmitContent: func(m provider.Message) string {
84 return StripReferencedContextPrefix(StripComposePrefixes(m.Content))
85 }, UserContent: func(m provider.Message) string {
86 if display := legacy.Users[transcript.LegacyDisplayKey(m.Content)]; display != "" {
87 return display
88 }
89 if m.RawContent != "" {
90 return m.RawContent
91 }
92 return StripReferencedContextPrefix(StripComposePrefixes(m.Content))
93 }})
94 p, err = transcript.NewProjection(identity, rows, covered)
95 if err != nil {
96 return nil, err
97 }
98 if err := replayTranscriptSuffix(p, ledger, covered, latest, identity.RuntimeEpoch); err != nil {
99 return nil, err
100 }
101 return p, nil
102 }
103
104 func restoreFromTranscriptCheckpoint(sessionPath string, ledger *turnevent.Ledger, identity transcript.Identity, messages []provider.Message, latest uint64) (*transcript.Projection, bool, error) {
105 checkpoint, exists, err := transcript.LoadCheckpoint(sessionPath)
106 if err != nil {
107 return nil, false, err
108 }
109 if !exists || checkpoint.CoveredThroughSeq > latest || checkpoint.ProviderCount < 0 || checkpoint.ProviderCount > len(messages) {
110 return nil, false, nil
111 }
112 prefixDigest, err := agent.ContentDigestForMessages(messages[:checkpoint.ProviderCount])
113 if err != nil {
114 return nil, false, err
115 }
116 if prefixDigest != checkpoint.TranscriptDigest || checkpoint.Identity.HeadID != identity.HeadID || checkpoint.Identity.RewriteEpoch != identity.RewriteEpoch {
117 return nil, false, nil
118 }
119 p, err := transcript.RestoreCheckpoint(checkpoint, identity)
120 if err != nil {
121 return nil, false, err
122 }
123 if err := replayTranscriptSuffix(p, ledger, checkpoint.CoveredThroughSeq, latest, identity.RuntimeEpoch); err != nil {
124 return nil, false, err
125 }
126 return p, true, nil
127 }
128
129 func replayTranscriptSuffix(p *transcript.Projection, ledger *turnevent.Ledger, after, latest uint64, runtimeEpoch string) error {
130 if after == latest {
131 return nil
132 }
133 events, err := ledger.EventsAfter(after)
134 if err != nil {
135 return err
136 }
137 for _, envelope := range events {
138 // Recovered turns are terminal. Adopt their display data into this
139 // runtime, without replaying any provider/tool/prompt side effect.
140 envelope.RuntimeEpoch = runtimeEpoch
141 if err := p.Apply(envelope); err != nil {
142 return err
143 }
144 }
145 if p.Boundary().CoveredThroughSeq != latest {
146 return errors.New("transcript recovery suffix is no longer retained")
147 }
148 return nil
149 }
150
150 lines GO