返回 DeepSeek-Reasonix
session_head_ops.go
根目录 / internal / agent / session_head_ops.go
1 package agent
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "log/slog"
8 "strings"
9 "time"
10
11 "reasonix/internal/store"
12 )
13
14 // ErrSessionHeadUnknown reports a head id that the log does not contain.
15 var ErrSessionHeadUnknown = errors.New("session head not found")
16
17 // ErrSessionNotDAG reports a head operation on a schema-1 session.
18 var ErrSessionNotDAG = errors.New("session log is not schema 2")
19
20 // ForkHead starts a new head at fromMessageID (empty means the root) and moves
21 // this session onto it: the in-memory transcript becomes that prefix and the
22 // next save appends behind it. kind is HeadKindFork or HeadKindRewind; the
23 // previous head keeps its full chain as a version. It returns the new head id.
24 func (s *Session) ForkHead(path, fromMessageID, kind, name string) (string, error) {
25 if kind == "" {
26 kind = HeadKindFork
27 }
28 newHead := NewHeadID()
29 var st *sessionDAGState
30 err := s.withSessionSaveLocks(path, func() error {
31 var err error
32 st, err = s.dagStateForHeadOp(path)
33 if err != nil {
34 return err
35 }
36 ref, _ := s.Head()
37 current := ref.HeadID
38 if current == "" || st.heads[current] == nil {
39 current = st.selectedHead()
40 }
41 if fromMessageID != "" && st.nodes[fromMessageID] == nil {
42 return fmt.Errorf("fork head: message %s: %w", fromMessageID, ErrSessionHeadUnknown)
43 }
44 now := time.Now().UTC()
45 entries := []sessionDAGEntry{
46 {Type: sessionDAGTypeFork, Head: current, NewHead: newHead, From: fromMessageID, Kind: kind, Name: strings.TrimSpace(name), At: now},
47 {Type: sessionDAGTypeSelect, Head: newHead, Reason: kind, At: now},
48 }
49 return appendHeadEntries(path, st, entries)
50 })
51 if err != nil {
52 return "", err
53 }
54 s.adoptHead(st, newHead, path)
55 return newHead, nil
56 }
57
58 // SwitchHead moves this session onto an existing head and records the choice
59 // with a select marker, so a later open lands on it too.
60 func (s *Session) SwitchHead(path, headID string) error {
61 var st *sessionDAGState
62 err := s.withSessionSaveLocks(path, func() error {
63 var err error
64 st, err = s.dagStateForHeadOp(path)
65 if err != nil {
66 return err
67 }
68 h := st.heads[headID]
69 if h == nil {
70 return fmt.Errorf("switch head %s: %w", headID, ErrSessionHeadUnknown)
71 }
72 if h.retired {
73 return fmt.Errorf("switch head %s: head is retired", headID)
74 }
75 if st.selected == headID {
76 return nil
77 }
78 return appendHeadEntries(path, st, []sessionDAGEntry{{Type: sessionDAGTypeSelect, Head: headID, Reason: "switch", At: time.Now().UTC()}})
79 })
80 if err != nil {
81 return err
82 }
83 s.adoptHead(st, headID, path)
84 return nil
85 }
86
87 // LoadSessionHeadReadOnly materializes one legacy schema-2 head without
88 // appending a select marker or changing the source log's default head. New
89 // runtimes use it to migrate historical heads into independent v3 sessions.
90 func LoadSessionHeadReadOnly(path, headID string) (*Session, error) {
91 return loadSessionHeadReadOnlyWithLimits(context.Background(), path, headID, defaultSessionReplayLimits)
92 }
93
94 func LoadSessionHeadReadOnlyContext(ctx context.Context, path, headID string) (*Session, error) {
95 return loadSessionHeadReadOnlyWithLimits(ctx, path, headID, defaultSessionReplayLimits)
96 }
97
98 // LoadSessionHeadForMigration materializes one head from a frozen legacy DAG
99 // without applying cumulative interactive-history replay budgets.
100 func LoadSessionHeadForMigration(ctx context.Context, path, headID string) (*Session, error) {
101 if ctx == nil {
102 ctx = context.Background()
103 }
104 return loadSessionHeadReadOnlyWithLimits(ctx, path, headID, migrationSessionReplayLimits())
105 }
106
107 func loadSessionHeadReadOnlyWithLimits(ctx context.Context, path, headID string, limits sessionReplayLimits) (*Session, error) {
108 st, err := replayDAGForHeadOpReadOnlyWithLimits(ctx, path, limits)
109 if err != nil {
110 return nil, err
111 }
112 head := st.heads[headID]
113 if head == nil || head.retired {
114 return nil, fmt.Errorf("load head %s: %w", headID, ErrSessionHeadUnknown)
115 }
116 s := NewSession("")
117 s.adoptHead(st, headID, path)
118 return s, nil
119 }
120
121 func replayDAGForHeadOpReadOnlyWithLimits(ctx context.Context, path string, limits sessionReplayLimits) (*sessionDAGState, error) {
122 probe, err := probeSessionEventLog(path)
123 if err != nil {
124 return nil, err
125 }
126 if !probe.dag {
127 return nil, ErrSessionNotDAG
128 }
129 st, err := replaySessionDAG(ctx, store.SessionEventLog(path), limits)
130 if err != nil {
131 return nil, err
132 }
133 if st.damaged {
134 return nil, fmt.Errorf("legacy session has an incomplete tail and is read-only")
135 }
136 return st, nil
137 }
138
139 // SelectSessionHead records the default head of a session that is not open
140 // in this process (the versions UI acting on a closed conversation).
141 func SelectSessionHead(path, headID string) error {
142 return appendSessionHeadMarker(path, headID, func(st *sessionDAGState) (sessionDAGEntry, error) {
143 if st.heads[headID].retired {
144 return sessionDAGEntry{}, fmt.Errorf("select head %s: head is retired", headID)
145 }
146 return sessionDAGEntry{Type: sessionDAGTypeSelect, Head: headID, Reason: "select"}, nil
147 })
148 }
149
150 // RetireSessionHead marks a head as cleaned up. Its exclusive entries are
151 // dropped at the next single-writer rotation; until then it is hidden.
152 func RetireSessionHead(path, headID string) error {
153 return appendSessionHeadMarker(path, headID, func(st *sessionDAGState) (sessionDAGEntry, error) {
154 if st.selectedHead() == headID {
155 return sessionDAGEntry{}, fmt.Errorf("retire head %s: it is the selected head", headID)
156 }
157 return sessionDAGEntry{Type: sessionDAGTypeRetire, Head: headID}, nil
158 })
159 }
160
161 // RenameSessionHead sets the display name of a head.
162 func RenameSessionHead(path, headID, name string) error {
163 return appendSessionHeadMarker(path, headID, func(*sessionDAGState) (sessionDAGEntry, error) {
164 return sessionDAGEntry{Type: sessionDAGTypeRename, Head: headID, Name: strings.TrimSpace(name)}, nil
165 })
166 }
167
168 func appendSessionHeadMarker(path, headID string, build func(*sessionDAGState) (sessionDAGEntry, error)) error {
169 if strings.TrimSpace(path) == "" || headID == "" {
170 return fmt.Errorf("session head marker: path and head id are required")
171 }
172 unlock := lockSessionSavePath(path)
173 defer unlock()
174 unlockFile, err := lockSessionFile(path)
175 if err != nil {
176 return fmt.Errorf("lock session file: %w", err)
177 }
178 defer unlockFile()
179 st, err := replayDAGForHeadOp(path)
180 if err != nil {
181 return err
182 }
183 if st.heads[headID] == nil {
184 return fmt.Errorf("head %s: %w", headID, ErrSessionHeadUnknown)
185 }
186 entry, err := build(st)
187 if err != nil {
188 return err
189 }
190 entry.At = time.Now().UTC()
191 return appendHeadEntries(path, st, []sessionDAGEntry{entry})
192 }
193
194 // dagStateForHeadOp returns the replayed graph for a head operation on an
195 // open session, reusing the cached state like a save does.
196 func (s *Session) dagStateForHeadOp(path string) (*sessionDAGState, error) {
197 if _, ok := s.Head(); !ok {
198 return nil, ErrSessionNotDAG
199 }
200 return s.dagStateForSave(context.Background(), path, time.Now().UTC())
201 }
202
203 func replayDAGForHeadOp(path string) (*sessionDAGState, error) {
204 probe, err := probeSessionEventLog(path)
205 if err != nil {
206 return nil, err
207 }
208 if !probe.dag {
209 return nil, ErrSessionNotDAG
210 }
211 st, err := replaySessionDAG(context.Background(), store.SessionEventLog(path), defaultSessionReplayLimits)
212 if err != nil {
213 return nil, err
214 }
215 if st.damaged {
216 if err := settleDAGTail(context.Background(), path, st, time.Now().UTC()); err != nil {
217 return nil, err
218 }
219 }
220 return st, nil
221 }
222
223 // appendHeadEntries lands marker entries, folds them into st, and refreshes
224 // the head index and meta mirror so listings see the change at once.
225 func appendHeadEntries(path string, st *sessionDAGState, entries []sessionDAGEntry) error {
226 tail := st.lastGoodEnd
227 if _, err := appendSessionDAGEntries(path, entries, true); err != nil {
228 return err
229 }
230 if err := st.replayFrom(context.Background(), tail, defaultSessionReplayLimits); err != nil {
231 return err
232 }
233 if st.damaged {
234 return fmt.Errorf("session log %s: head markers did not replay", path)
235 }
236 if err := writeSessionDAGIndex(context.Background(), path, st); err != nil {
237 slog.Warn("session: head index write after head marker failed", "path", path, "err", err)
238 }
239 selected := st.selectedHead()
240 if err := UpdateBranchMeta(path, false, func(meta *BranchMeta) error {
241 meta.HeadID = selected
242 meta.HeadCount = len(st.heads)
243 meta.LogSchema = sessionDAGSchemaVersion
244 meta.LogGeneration = st.generation
245 return nil
246 }); err != nil {
247 slog.Warn("session: head metadata update after head marker failed", "path", path, "err", err)
248 }
249 return nil
250 }
251
252 // adoptHead points the live session at headID: the transcript becomes the
253 // head's materialized chain and the persisted baseline is that chain, so the
254 // next save diffs against it instead of the previous head.
255 func (s *Session) adoptHead(st *sessionDAGState, headID, path string) {
256 msgs, _ := st.materialize(headID)
257 digest, err := digestSessionMessages(msgs)
258 s.mu.Lock()
259 s.Messages = msgs
260 s.version++
261 s.rewriteVersion++
262 version, rewriteVersion := s.version, s.rewriteVersion
263 s.head.ref = HeadRef{HeadID: headID, LeafID: st.heads[headID].leaf, LogGeneration: st.generation, LogOffset: st.lastGoodEnd}
264 s.head.dag = true
265 s.head.state = st
266 s.head.headCount = len(st.heads)
267 s.head.openTurn = st.heads[headID].openTurn
268 s.mu.Unlock()
269 if err == nil {
270 revision, _, _ := sessionContentRevision(path)
271 s.setPersistedBaseline(path, digest, version, revision, true, true, rewriteVersion, msgs)
272 }
273 }
274
274 lines GO