| 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 |