返回 DeepSeek-Reasonix
save_dag.go
根目录 / internal / agent / save_dag.go
1 package agent
2
3 import (
4 "context"
5 "crypto/sha256"
6 "fmt"
7 "log/slog"
8 "os"
9 "time"
10
11 "reasonix/internal/provider"
12 "reasonix/internal/store"
13 )
14
15 type dagRoute int
16
17 const (
18 dagRouteSchemaOne dagRoute = iota
19 dagRouteNative
20 dagRouteUpgrade
21 )
22
23 // SessionLogSchemaEnv is the emergency switch back to the schema-1 writer for
24 // sessions that have not been upgraded yet; upgraded logs stay schema 2.
25 const SessionLogSchemaEnv = "REASONIX_SESSION_LOG"
26
27 func dagWriterEnabled() bool {
28 return os.Getenv(SessionLogSchemaEnv) != "v1"
29 }
30
31 // probeLogForSave classifies the event log for a save, refusing a schema this
32 // build cannot own and healing a torn schema-1 tail before anything appends
33 // behind it.
34 func probeLogForSave(path string) (sessionEventLogProbe, error) {
35 probe, err := probeSessionEventLog(path)
36 if err != nil {
37 return probe, err
38 }
39 if probe.futureSchema {
40 return probe, fmt.Errorf("session event log for %s uses schema %d; this build supports up to %d", path, probe.schemaVersion, sessionDAGSchemaVersion)
41 }
42 if probe.native && probe.size > 0 {
43 if err := repairSessionEventLogTail(path); err != nil {
44 return probe, fmt.Errorf("repair session event log: %w", err)
45 }
46 }
47 return probe, nil
48 }
49
50 // Loaded bare checkpoints remain checkpoint-only on continuation.
51 func (s *Session) probeNativeLogForSave(path string) (sessionEventLogProbe, error) {
52 probe, err := probeLogForSave(path)
53 if s.persistFormat == sessionPersistLegacy && probe.size == 0 {
54 probe.native = false
55 }
56 return probe, err
57 }
58
59 // Loaded sessions retain their schema; fresh sessions use normal DAG admission.
60 func (s *Session) dagSaveRoute(path string, probe sessionEventLogProbe) dagRoute {
61 if probe.dag {
62 return dagRouteNative
63 }
64 if s != nil && s.persistFormat == sessionPersistLegacy {
65 return dagRouteSchemaOne
66 }
67 if !dagWriterEnabled() || !probe.native {
68 return dagRouteSchemaOne
69 }
70 if probe.size > 0 {
71 if !SessionLeaseHeldByCurrentRuntime(path) {
72 return dagRouteSchemaOne
73 }
74 return dagRouteUpgrade
75 }
76 if SessionLeaseHeldByOtherRuntime(path) {
77 return dagRouteSchemaOne
78 }
79 return dagRouteUpgrade
80 }
81
82 // saveDAGLocked persists msgs to a schema-2 log: it replays (or extends) the
83 // cached graph, diffs the in-memory transcript against this session's head,
84 // appends the resulting entries in one batch, and refreshes the derived
85 // files. Concurrent writers never conflict; a writer that only fell behind
86 // its own head reports a stale-prefix conflict so the caller adopts disk.
87 func (s *Session) saveDAGLocked(path string, mode sessionSaveMode, route dagRoute, msgs []provider.Message, version uint64, rewriteVersion int, digest [sha256.Size]byte) error {
88 ctx := context.Background()
89 now := time.Now().UTC()
90 baseRevision, _, err := sessionContentRevision(path)
91 ledgerUnreadable := err != nil
92 if ledgerUnreadable {
93 // A persisted session with an unreadable ledger fails closed, as in
94 // schema 1; a brand-new transcript still lands first and the record
95 // step below reports the sidecar.
96 if sessionArtifactsHaveContent(path) {
97 return err
98 }
99 baseRevision = 0
100 }
101 var st *sessionDAGState
102 if route == dagRouteUpgrade {
103 st, err = s.upgradeLogForSave(ctx, path, now)
104 } else {
105 st, err = s.dagStateForSave(ctx, path, now)
106 }
107 if err != nil {
108 return err
109 }
110 s.ensureMessageIDsForSave(msgs)
111 plan, err := s.planDAGWrite(path, st, msgs, mode, now)
112 if err != nil {
113 return err
114 }
115 deferProjection := mode.defersProjection()
116 pending := s.takePendingMarkers()
117 for i := range pending {
118 pending[i].Head = plan.head
119 }
120 plan.entries = append(plan.entries, pending...)
121 if len(plan.entries) == 0 {
122 s.adoptDAGPosition(st, plan)
123 s.republishDAGDerivedIfPending(ctx, path, st, plan, msgs, digest, baseRevision)
124 s.maintainDAGLog(ctx, path, st, plan, mode, now)
125 s.markCheckpointPersisted(path, digest, version, baseRevision, rewriteVersion, msgs, deferProjection)
126 return nil
127 }
128 reserved, err := invalidateSessionListingProjection(path)
129 if err != nil {
130 if !ledgerUnreadable {
131 return fmt.Errorf("invalidate session listing projection: %w", err)
132 }
133 reserved = 0
134 }
135 tail := st.lastGoodEnd
136 if _, err := appendSessionDAGEntries(path, plan.entries, true); err != nil {
137 s.requeuePendingMarkers(pending)
138 return err
139 }
140 if err := st.replayFrom(ctx, tail, defaultSessionReplayLimits); err != nil {
141 return err
142 }
143 if st.damaged {
144 return fmt.Errorf("session log %s: appended entries did not replay", path)
145 }
146 plan.applyIDRenames(s)
147 s.adoptDAGPosition(st, plan)
148 // The compatibility checkpoint lands before the ledger, as in schema 1, so
149 // a metadata failure never leaves the anchor behind the log.
150 selected := st.selectedHead()
151 displayCurrent := selected == plan.head && writeDAGCheckpointCache(path, plan, msgs, baseRevision)
152 revision, err := recordSessionContentRevision(path, digest, baseRevision, reserved)
153 if err != nil {
154 return err
155 }
156 s.publishDAGDerived(ctx, path, st, plan, msgs, digest, revision, selected, displayCurrent, deferProjection)
157 s.maintainDAGLog(ctx, path, st, plan, mode, now)
158 s.markCheckpointPersisted(path, digest, version, revision, rewriteVersion, msgs, deferProjection)
159 return nil
160 }
161
162 // dagStateForSave returns the replayed graph, extending the cached state by
163 // the bytes appended since its last observed tail when the generation and
164 // size still allow it, and settles a torn tail before any append.
165 func (s *Session) dagStateForSave(ctx context.Context, path string, now time.Time) (*sessionDAGState, error) {
166 logPath := store.SessionEventLog(path)
167 s.mu.RLock()
168 cached := s.head.state
169 s.mu.RUnlock()
170 header, ok, err := readSessionDAGHeader(path)
171 if err != nil {
172 return nil, err
173 }
174 var st *sessionDAGState
175 if cached != nil && ok && cached.path == logPath && header.generation == cached.generation && !cached.damaged {
176 if info, err := os.Stat(logPath); err == nil && info.Size() >= cached.lastGoodEnd {
177 st = cached
178 if err := st.replayFrom(ctx, st.lastGoodEnd, defaultSessionReplayLimits); err != nil {
179 return nil, err
180 }
181 }
182 }
183 if st == nil {
184 st, err = replaySessionDAG(ctx, logPath, defaultSessionReplayLimits)
185 if err != nil {
186 return nil, err
187 }
188 }
189 if st.damaged {
190 if err := settleDAGTail(ctx, path, st, now); err != nil {
191 return nil, err
192 }
193 }
194 return st, nil
195 }
196
197 // settleDAGTail waits out the quiet window a torn tail might still be
198 // finishing, then repairs it; appending after an unrepaired partial line
199 // would bury this writer's own entry inside it.
200 func settleDAGTail(ctx context.Context, path string, st *sessionDAGState, now time.Time) error {
201 logPath := store.SessionEventLog(path)
202 info, err := os.Stat(logPath)
203 if err != nil {
204 return err
205 }
206 if age := now.Sub(info.ModTime()); age < sessionDAGTailRepairMinAge {
207 time.Sleep(sessionDAGTailRepairMinAge - age)
208 st.damaged = false
209 if err := st.replayFrom(ctx, st.lastGoodEnd, defaultSessionReplayLimits); err != nil {
210 return err
211 }
212 if !st.damaged {
213 return nil
214 }
215 }
216 repaired, err := repairSessionDAGTail(path, st, time.Now().UTC())
217 if err != nil {
218 return err
219 }
220 if !repaired {
221 return fmt.Errorf("session log %s has a torn tail that is still being written", path)
222 }
223 return nil
224 }
225
226 // upgradeLogForSave replaces the schema-1 log (or bare checkpoint) with
227 // generation 1 of a schema-2 log built from the transcript on disk. The
228 // in-memory delta is then written by the ordinary diff, exactly like any
229 // other save. Callers hold the file lock and satisfied dagSaveRoute.
230 func (s *Session) upgradeLogForSave(ctx context.Context, path string, now time.Time) (*sessionDAGState, error) {
231 disk, err := loadSessionTranscript(ctx, path, defaultSessionReplayLimits, nil)
232 if err != nil && !os.IsNotExist(err) {
233 return nil, err
234 }
235 msgs := migrateLegacyProviderContent(NormalizeSession(disk.msgs))
236 assignLegacyMessageIDs(path, msgs)
237 var inFlight *InFlightTurnMeta
238 if meta, ok, err := LoadBranchMeta(path); err == nil && ok {
239 inFlight = meta.InFlightTurn
240 }
241 if err := upgradeSessionLogToDAG(path, msgs, disk.times, inFlight, now); err != nil {
242 return nil, err
243 }
244 st, err := replaySessionDAG(ctx, store.SessionEventLog(path), defaultSessionReplayLimits)
245 if err != nil {
246 return nil, err
247 }
248 if len(msgs) > 0 {
249 slog.Info("session: upgraded event log to schema 2", "path", path, "messages", len(msgs))
250 }
251 return st, nil
252 }
253
254 // ensureMessageIDsForSave mints ids for messages that reached the session
255 // without one and writes them back by position so the next save sees the
256 // same ids the log now holds.
257 func (s *Session) ensureMessageIDsForSave(msgs []provider.Message) {
258 minted := false
259 for i := range msgs {
260 if msgs[i].ID == "" {
261 msgs[i].ID = NewMessageID()
262 minted = true
263 }
264 }
265 if !minted {
266 return
267 }
268 s.mu.Lock()
269 defer s.mu.Unlock()
270 for i := range msgs {
271 if i < len(s.Messages) && s.Messages[i].ID == "" && messagesEqualForStorage(s.Messages[i], msgs[i]) {
272 s.Messages[i].ID = msgs[i].ID
273 }
274 }
275 }
276
277 func (s *Session) adoptDAGPosition(st *sessionDAGState, plan *dagWritePlan) {
278 h := st.heads[plan.head]
279 leaf := ""
280 if h != nil {
281 leaf = h.leaf
282 }
283 s.mu.Lock()
284 defer s.mu.Unlock()
285 s.head.ref = HeadRef{HeadID: plan.head, LeafID: leaf, LogGeneration: st.generation, LogOffset: st.lastGoodEnd}
286 s.head.dag = true
287 s.head.state = st
288 s.head.headCount = len(st.heads)
289 if plan.forked {
290 s.head.events = append(s.head.events, HeadEvent{Kind: HeadEventForkedConcurrent, HeadID: plan.head, OtherWriter: plan.otherWriter})
291 }
292 }
293
294 // writeDAGCheckpointCache refreshes the .jsonl random-read model for the
295 // selected head, extending it in place for a pure append. It reports whether
296 // the cache now matches msgs; failures are logged, never fatal.
297 func writeDAGCheckpointCache(path string, plan *dagWritePlan, msgs []provider.Message, baseRevision int64) bool {
298 if plan.pureAppend {
299 current, err := appendSessionDisplayReadModel(path, msgs, plan.appendFrom, baseRevision)
300 if err != nil {
301 slog.Warn("session: keeping save after display read-model append failure", "path", path, "err", err)
302 }
303 if current {
304 return true
305 }
306 }
307 if err := writeSessionMessages(path, msgs); err != nil {
308 slog.Warn("session: keeping save after display read-model write failure", "path", path, "err", err)
309 return false
310 }
311 return true
312 }
313
314 // publishDAGDerived refreshes the display index when the .jsonl cache is
315 // current for this head, and the head index plus meta mirror on every save.
316 func (s *Session) publishDAGDerived(ctx context.Context, path string, st *sessionDAGState, plan *dagWritePlan, msgs []provider.Message, digest [sha256.Size]byte, revision int64, selected string, displayCurrent, deferProjection bool) {
317 if displayCurrent {
318 appendFrom := -1
319 if plan.pureAppend {
320 appendFrom = plan.appendFrom
321 }
322 if err := refreshCheckpointDisplayIndex(path, msgs, digest, revision, appendFrom, deferProjection); err != nil {
323 slog.Warn("session: keeping save after display index write failure", "path", path, "err", err)
324 }
325 }
326 if err := writeSessionDAGIndex(ctx, path, st); err != nil {
327 slog.Warn("session: keeping save after head index write failure", "path", path, "err", err)
328 }
329 if err := UpdateBranchMeta(path, false, func(meta *BranchMeta) error {
330 meta.HeadID = selected
331 meta.HeadCount = len(st.heads)
332 meta.LogSchema = sessionDAGSchemaVersion
333 meta.LogGeneration = st.generation
334 return nil
335 }); err != nil {
336 slog.Warn("session: head metadata update deferred", "path", path, "err", err)
337 }
338 }
339
340 // maintainDAGLog rotates the log when it has outgrown its live chains or a
341 // redaction needs its bytes physically erased, but only under the
342 // single-writer proof; otherwise the log simply keeps growing for now.
343 func (s *Session) maintainDAGLog(ctx context.Context, path string, st *sessionDAGState, plan *dagWritePlan, mode sessionSaveMode, now time.Time) {
344 if mode != sessionSaveRewriteCompact && st.holes == 0 && !sessionDAGLogOversized(st) {
345 return
346 }
347 if err := sessionDAGSingleWriterProof(path, st, now); err != nil {
348 slog.Info("session: log rotation deferred", "path", path, "reason", err)
349 return
350 }
351 if err := rotateSessionDAG(path, st, now); err != nil {
352 slog.Warn("session: log rotation failed", "path", path, "err", err)
353 return
354 }
355 fresh, err := replaySessionDAG(ctx, store.SessionEventLog(path), defaultSessionReplayLimits)
356 if err != nil {
357 slog.Warn("session: replay after rotation failed", "path", path, "err", err)
358 return
359 }
360 s.adoptDAGPosition(fresh, &dagWritePlan{head: plan.head})
361 if err := writeSessionDAGIndex(ctx, path, fresh); err != nil {
362 slog.Warn("session: head index write after rotation failed", "path", path, "err", err)
363 }
364 }
365
365 lines GO