返回 DeepSeek-Reasonix
session_head.go
根目录 / internal / agent / session_head.go
1 package agent
2
3 import (
4 "context"
5 "fmt"
6 "slices"
7 "time"
8
9 "reasonix/internal/store"
10 )
11
12 // sessionHeadState is what a Session knows about its schema-2 log position;
13 // it stays zero for schema-1 sessions. state caches the replayed graph so a
14 // save only reads the bytes appended since the last observed tail.
15 type sessionHeadState struct {
16 ref HeadRef
17 dag bool
18 headCount int
19 state *sessionDAGState
20 events []HeadEvent
21 pending []sessionDAGEntry
22 openTurn *sessionDAGTurn
23 }
24
25 // clone retains execution identity without sharing the mutable replay cache.
26 // Pending entries are turn markers containing only scalar fields.
27 func (h sessionHeadState) clone() sessionHeadState {
28 h.state = nil
29 h.events = slices.Clone(h.events)
30 h.pending = slices.Clone(h.pending)
31 if h.openTurn != nil {
32 turn := *h.openTurn
33 h.openTurn = &turn
34 }
35 return h
36 }
37
38 // HeadEvent reports a head-level fact a save discovered; the controller turns
39 // it into a user-facing notice.
40 type HeadEvent struct {
41 Kind string
42 HeadID string
43 OtherWriter string
44 }
45
46 const (
47 // HeadEventForkedConcurrent: this writer's head diverged from another
48 // writer's appends and continued on a fresh head.
49 HeadEventForkedConcurrent = "forked_concurrent"
50 // HeadEventMultipleRecentHeads: the session opened on one of several heads
51 // that were active within recentHeadWindow; the others remain versions.
52 HeadEventMultipleRecentHeads = "multiple_recent_heads"
53 )
54
55 // recentHeadWindow bounds how old a competing head may be before opening a
56 // conversation stops mentioning it.
57 const recentHeadWindow = 24 * time.Hour
58
59 // loadHeadEvents derives the events a fresh load should surface.
60 func loadHeadEvents(st *sessionDAGState, selected string) []HeadEvent {
61 if st == nil {
62 return nil
63 }
64 h := st.heads[selected]
65 if h == nil {
66 return nil
67 }
68 for _, id := range st.liveHeads() {
69 if id == selected {
70 continue
71 }
72 if other := st.heads[id]; other != nil && h.lastActivity.Sub(other.lastActivity) < recentHeadWindow {
73 return []HeadEvent{{Kind: HeadEventMultipleRecentHeads, HeadID: selected, OtherWriter: other.writer}}
74 }
75 }
76 return nil
77 }
78
79 // Head reports the head this session was loaded from or last saved to. ok is
80 // false for sessions that still live in a schema-1 log or a bare .jsonl.
81 func (s *Session) Head() (HeadRef, bool) {
82 if s == nil {
83 return HeadRef{}, false
84 }
85 s.mu.RLock()
86 defer s.mu.RUnlock()
87 return s.head.ref, s.head.dag
88 }
89
90 // DrainHeadEvents returns and clears the head events recorded by saves.
91 func (s *Session) DrainHeadEvents() []HeadEvent {
92 if s == nil {
93 return nil
94 }
95 s.mu.Lock()
96 defer s.mu.Unlock()
97 events := s.head.events
98 s.head.events = nil
99 return events
100 }
101
102 // selectedHead picks the head a plain open lands on: the last explicit select
103 // while that head is alive, otherwise the most recently active live head, with
104 // the larger log offset breaking ties so every reader of the same bytes agrees.
105 func (st *sessionDAGState) selectedHead() string {
106 if h := st.heads[st.selected]; h != nil && !h.retired {
107 return h.id
108 }
109 pick := func(includeRetired bool) string {
110 best := ""
111 for _, id := range st.headOrder {
112 h := st.heads[id]
113 if h == nil || (h.retired && !includeRetired) {
114 continue
115 }
116 if best == "" {
117 best = id
118 continue
119 }
120 b := st.heads[best]
121 if h.lastActivity.After(b.lastActivity) || (h.lastActivity.Equal(b.lastActivity) && h.lastOffset > b.lastOffset) {
122 best = id
123 }
124 }
125 return best
126 }
127 if best := pick(false); best != "" {
128 return best
129 }
130 if best := pick(true); best != "" {
131 return best
132 }
133 return SessionMainHead
134 }
135
136 func (st *sessionDAGState) headRecord(id string, selected string) SessionHead {
137 h := st.heads[id]
138 msgs, _ := st.materialize(id)
139 preview, turns := SessionPreviewFromMessages(msgs)
140 return SessionHead{
141 ID: h.id,
142 Kind: h.kind,
143 Name: h.name,
144 ParentHead: h.parentHead,
145 ForkFrom: h.forkFrom,
146 Writer: h.writer,
147 LeafID: h.leaf,
148 CreatedAt: h.createdAt,
149 LastActivity: h.lastActivity,
150 Retired: h.retired,
151 Selected: h.id == selected,
152 MessageCount: len(msgs),
153 Turns: turns,
154 Preview: preview,
155 }
156 }
157
158 // headList lists every head in declaration order with the selected one marked
159 // and covered heads flagged: retiring one of those loses no message.
160 func (st *sessionDAGState) headList() []SessionHead {
161 selected := st.selectedHead()
162 onChain := map[string]struct{}{}
163 for _, id := range st.chainIDs(selected) {
164 onChain[id] = struct{}{}
165 }
166 out := make([]SessionHead, 0, len(st.headOrder))
167 for _, id := range st.headOrder {
168 h := st.heads[id]
169 if h == nil {
170 continue
171 }
172 rec := st.headRecord(id, selected)
173 if id != selected && !h.retired {
174 _, rec.Covered = onChain[h.leaf]
175 rec.Covered = rec.Covered || h.leaf == ""
176 }
177 out = append(out, rec)
178 }
179 return out
180 }
181
182 // liveHeads returns the ids of heads that have not been retired.
183 func (st *sessionDAGState) liveHeads() []string {
184 out := make([]string, 0, len(st.headOrder))
185 for _, id := range st.headOrder {
186 if h := st.heads[id]; h != nil && !h.retired {
187 out = append(out, id)
188 }
189 }
190 return out
191 }
192
193 // reachable returns every node id on the chain of any live head.
194 func (st *sessionDAGState) reachable() map[string]struct{} {
195 keep := map[string]struct{}{}
196 for _, id := range st.liveHeads() {
197 for _, mid := range st.chainIDs(id) {
198 keep[mid] = struct{}{}
199 }
200 }
201 return keep
202 }
203
204 // ListSessionHeads replays a schema-2 log and returns its heads. A schema-1
205 // session has no heads and returns nil, nil.
206 func ListSessionHeads(path string) ([]SessionHead, error) {
207 return listSessionHeads(context.Background(), path, defaultSessionReplayLimits, false)
208 }
209
210 // ListSessionHeadsForMigration enumerates a frozen source without the cumulative
211 // interactive replay budget. The migration owner must freeze the source first.
212 func ListSessionHeadsForMigration(ctx context.Context, path string) ([]SessionHead, error) {
213 return listSessionHeads(ctx, path, migrationSessionReplayLimits(), true)
214 }
215
216 func listSessionHeads(ctx context.Context, path string, limits sessionReplayLimits, strict bool) ([]SessionHead, error) {
217 probe, err := probeSessionEventLogWithLimits(path, limits)
218 if err != nil {
219 return nil, err
220 }
221 if !probe.dag {
222 return nil, nil
223 }
224 st, err := replaySessionDAG(ctx, store.SessionEventLog(path), limits)
225 if err != nil {
226 return nil, fmt.Errorf("list session heads: %w", err)
227 }
228 if strict && st.damaged {
229 return nil, fmt.Errorf("list session heads: incomplete legacy DAG")
230 }
231 heads := st.headList()
232 slices.SortStableFunc(heads, func(a, b SessionHead) int {
233 return a.CreatedAt.Compare(b.CreatedAt)
234 })
235 return heads, nil
236 }
237
237 lines GO