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