| 1 | package session |
| 2 | |
| 3 | import ( |
| 4 | "bytes" |
| 5 | "context" |
| 6 | "crypto/sha256" |
| 7 | "encoding/json" |
| 8 | "errors" |
| 9 | "fmt" |
| 10 | "io" |
| 11 | "os" |
| 12 | "path/filepath" |
| 13 | "slices" |
| 14 | "sort" |
| 15 | "strings" |
| 16 | "time" |
| 17 | |
| 18 | "github.com/klauspost/compress/zstd" |
| 19 | bolt "go.etcd.io/bbolt" |
| 20 | |
| 21 | "reasonix/internal/agent" |
| 22 | "reasonix/internal/fileutil" |
| 23 | "reasonix/internal/provider" |
| 24 | "reasonix/internal/sessioncontent" |
| 25 | ) |
| 26 | |
| 27 | // Version 8 carries the live message ids the writer checks new completions |
| 28 | // against. Older projections are disposable and rebuild from the unchanged log. |
| 29 | const recoveryProjectionVersion = 8 |
| 30 | |
| 31 | const ( |
| 32 | recoveryFormatVersion = 1 |
| 33 | recoveryDBName = "recovery-v1.bolt" |
| 34 | storageIdentityName = "storage.identity.json" |
| 35 | RecentMessageLimit = 100 |
| 36 | recentInlineBytes = 32 << 10 |
| 37 | recentResponseBytes = 512 << 10 |
| 38 | recentSnapshotName = "recent-v1.json" |
| 39 | ) |
| 40 | |
| 41 | var ( |
| 42 | recoveryMetaBucket = []byte("meta") |
| 43 | recoveryCheckpointBucket = []byte("checkpoints") |
| 44 | recoveryOperationBucket = []byte("operations") |
| 45 | recoveryCurrentKey = []byte("current") |
| 46 | recoveryPreviousKey = []byte("previous") |
| 47 | ) |
| 48 | |
| 49 | // RecoveryOpenStats reports work performed by the normal open path. It is |
| 50 | // intentionally small and stable enough for capacity tests and host telemetry. |
| 51 | type RecoveryOpenStats struct { |
| 52 | UsedCheckpoint bool `json:"usedCheckpoint"` |
| 53 | LogBytesRead int64 `json:"logBytesRead"` |
| 54 | LogBytesTotal int64 `json:"logBytesTotal"` |
| 55 | TailCommits int `json:"tailCommits"` |
| 56 | } |
| 57 | |
| 58 | // RecentSnapshot is the bounded, read-only baseline used before history or |
| 59 | // search projections are available. |
| 60 | type RecentSnapshot struct { |
| 61 | Version int `json:"version"` |
| 62 | SessionID string `json:"sessionId"` |
| 63 | StorageGeneration string `json:"storageGeneration"` |
| 64 | DurableSequence uint64 `json:"durableSequence"` |
| 65 | TotalTurns int `json:"totalTurns"` |
| 66 | Entries []PersistentMessage `json:"entries"` |
| 67 | Title string `json:"title,omitempty"` |
| 68 | ModelRef string `json:"modelRef,omitempty"` |
| 69 | ModelIdentity string `json:"modelIdentity,omitempty"` |
| 70 | } |
| 71 | |
| 72 | type storageIdentity struct { |
| 73 | Version int `json:"version"` |
| 74 | SessionID string `json:"sessionId"` |
| 75 | Generation string `json:"generation"` |
| 76 | LogPrefixDigest string `json:"logPrefixDigest,omitempty"` |
| 77 | CreatedAt time.Time `json:"createdAt"` |
| 78 | } |
| 79 | |
| 80 | type recoveryCheckpoint struct { |
| 81 | Version int `json:"version"` |
| 82 | SessionID string `json:"sessionId"` |
| 83 | StorageGeneration string `json:"storageGeneration"` |
| 84 | StorageRevision int `json:"storageRevision"` |
| 85 | DurableSequence uint64 `json:"durableSequence"` |
| 86 | LogOffset int64 `json:"logOffset"` |
| 87 | AnchorOffset int64 `json:"anchorOffset,omitempty"` |
| 88 | AnchorFirst uint64 `json:"anchorFirst,omitempty"` |
| 89 | AnchorCommitID string `json:"anchorCommitId,omitempty"` |
| 90 | AnchorHash string `json:"anchorHash,omitempty"` |
| 91 | ProjectionVersion int `json:"projectionVersion"` |
| 92 | Projection Projection `json:"projection"` |
| 93 | RecentMessages []provider.Message `json:"recentMessages,omitempty"` |
| 94 | CatalogPreview string `json:"catalogPreview,omitempty"` |
| 95 | MessageIDs []string `json:"messageIds,omitempty"` |
| 96 | CreatedAt time.Time `json:"createdAt"` |
| 97 | } |
| 98 | |
| 99 | type recoveryOperation struct { |
| 100 | Hash string `json:"hash"` |
| 101 | CommitID string `json:"commitId"` |
| 102 | FirstSequence uint64 `json:"firstSequence"` |
| 103 | EventCount int `json:"eventCount"` |
| 104 | TurnID string `json:"turnId,omitempty"` |
| 105 | OperationID string `json:"operationId"` |
| 106 | OperationHash string `json:"operationHash"` |
| 107 | WriterGeneration uint64 `json:"writerGeneration"` |
| 108 | CreatedAt time.Time `json:"createdAt"` |
| 109 | } |
| 110 | |
| 111 | func operationForRecovery(record operationRecord) recoveryOperation { |
| 112 | commit := record.commit |
| 113 | return recoveryOperation{ |
| 114 | Hash: record.hash, CommitID: commit.ID, FirstSequence: commit.FirstSequence, |
| 115 | EventCount: commit.EventCount, TurnID: commit.TurnID, OperationID: commit.OperationID, |
| 116 | OperationHash: commit.OperationHash, WriterGeneration: commit.WriterGeneration, |
| 117 | CreatedAt: commit.CreatedAt, |
| 118 | } |
| 119 | } |
| 120 | |
| 121 | func (o recoveryOperation) record() operationRecord { |
| 122 | return operationRecord{hash: o.Hash, commit: Commit{ |
| 123 | SchemaVersion: SchemaVersion, Codec: Codec, RecordType: "commit", ID: o.CommitID, |
| 124 | OperationID: o.OperationID, OperationHash: o.OperationHash, |
| 125 | FirstSequence: o.FirstSequence, EventCount: o.EventCount, TurnID: o.TurnID, |
| 126 | WriterGeneration: o.WriterGeneration, CreatedAt: o.CreatedAt, |
| 127 | }} |
| 128 | } |
| 129 | |
| 130 | type recoveryStore struct { |
| 131 | db *bolt.DB |
| 132 | path string |
| 133 | recent string |
| 134 | sessionDir string |
| 135 | identity storageIdentity |
| 136 | } |
| 137 | |
| 138 | func recoveryCacheDir(sessionDir string) string { |
| 139 | return filepath.Join(filepath.Dir(sessionDir), ".recovery-cache", filepath.Base(sessionDir)) |
| 140 | } |
| 141 | |
| 142 | func ensureStorageIdentity(sessionDir string, manifest Manifest) (storageIdentity, error) { |
| 143 | path := filepath.Join(sessionDir, storageIdentityName) |
| 144 | prefix, prefixErr := storageLogPrefix(sessionDir, manifest) |
| 145 | if prefixErr != nil && !os.IsNotExist(prefixErr) { |
| 146 | return storageIdentity{}, prefixErr |
| 147 | } |
| 148 | if data, err := os.ReadFile(path); err == nil { |
| 149 | var identity storageIdentity |
| 150 | if json.Unmarshal(data, &identity) == nil && identity.Version == recoveryFormatVersion && identity.SessionID == manifest.SessionID && strings.TrimSpace(identity.Generation) != "" { |
| 151 | if identity.LogPrefixDigest == "" && prefix != "" { |
| 152 | identity.LogPrefixDigest = prefix |
| 153 | encoded, marshalErr := json.Marshal(identity) |
| 154 | if marshalErr != nil { |
| 155 | return storageIdentity{}, marshalErr |
| 156 | } |
| 157 | if writeErr := fileutil.AtomicWriteFileStrict(path, append(encoded, '\n'), 0o600); writeErr != nil { |
| 158 | return storageIdentity{}, writeErr |
| 159 | } |
| 160 | return identity, nil |
| 161 | } |
| 162 | if prefix == "" || identity.LogPrefixDigest == prefix { |
| 163 | return identity, nil |
| 164 | } |
| 165 | } |
| 166 | } |
| 167 | identity := storageIdentity{Version: recoveryFormatVersion, SessionID: manifest.SessionID, Generation: randomID(), LogPrefixDigest: prefix, CreatedAt: time.Now().UTC()} |
| 168 | data, err := json.Marshal(identity) |
| 169 | if err != nil { |
| 170 | return storageIdentity{}, err |
| 171 | } |
| 172 | if err := fileutil.AtomicWriteFileStrict(path, append(data, '\n'), 0o600); err != nil { |
| 173 | return storageIdentity{}, err |
| 174 | } |
| 175 | return identity, nil |
| 176 | } |
| 177 | |
| 178 | func storageLogPrefix(sessionDir string, manifest Manifest) (string, error) { |
| 179 | file, err := os.Open(logPathForManifest(sessionDir, manifest)) |
| 180 | if err != nil { |
| 181 | return "", err |
| 182 | } |
| 183 | defer file.Close() |
| 184 | // The framed transaction header makes the first 64 bytes immutable after |
| 185 | // the first commit. Hashing a larger short-file prefix would change merely |
| 186 | // because a normal append extended a log shorter than that prefix. |
| 187 | buffer := make([]byte, 64) |
| 188 | n, err := file.Read(buffer) |
| 189 | if err != nil && !errors.Is(err, io.EOF) { |
| 190 | return "", err |
| 191 | } |
| 192 | if n == 0 { |
| 193 | return "", nil |
| 194 | } |
| 195 | digest := sha256.Sum256(buffer[:n]) |
| 196 | return fmt.Sprintf("%x", digest[:]), nil |
| 197 | } |
| 198 | |
| 199 | func readStorageIdentity(sessionDir string, manifest Manifest) (storageIdentity, error) { |
| 200 | data, err := os.ReadFile(filepath.Join(sessionDir, storageIdentityName)) |
| 201 | if err != nil { |
| 202 | return storageIdentity{}, err |
| 203 | } |
| 204 | var identity storageIdentity |
| 205 | if err := json.Unmarshal(data, &identity); err != nil { |
| 206 | return storageIdentity{}, err |
| 207 | } |
| 208 | if identity.Version != recoveryFormatVersion || identity.SessionID != manifest.SessionID || strings.TrimSpace(identity.Generation) == "" { |
| 209 | return storageIdentity{}, ErrStaleGeneration |
| 210 | } |
| 211 | prefix, err := storageLogPrefix(sessionDir, manifest) |
| 212 | if err != nil && !os.IsNotExist(err) { |
| 213 | return storageIdentity{}, err |
| 214 | } |
| 215 | if identity.LogPrefixDigest != "" && prefix != identity.LogPrefixDigest { |
| 216 | return storageIdentity{}, ErrStaleGeneration |
| 217 | } |
| 218 | return identity, nil |
| 219 | } |
| 220 | |
| 221 | func openRecoveryStore(sessionDir string, identity storageIdentity) (*recoveryStore, error) { |
| 222 | dir := recoveryCacheDir(sessionDir) |
| 223 | if err := os.MkdirAll(dir, 0o700); err != nil { |
| 224 | return nil, err |
| 225 | } |
| 226 | path := filepath.Join(dir, recoveryDBName) |
| 227 | db, err := bolt.Open(path, 0o600, &bolt.Options{Timeout: 250 * time.Millisecond, NoFreelistSync: true}) |
| 228 | if err != nil { |
| 229 | return nil, err |
| 230 | } |
| 231 | store := &recoveryStore{db: db, path: path, recent: filepath.Join(dir, recentSnapshotName), sessionDir: sessionDir, identity: identity} |
| 232 | err = db.Update(func(tx *bolt.Tx) error { |
| 233 | meta, err := tx.CreateBucketIfNotExists(recoveryMetaBucket) |
| 234 | if err != nil { |
| 235 | return err |
| 236 | } |
| 237 | checkpoints, err := tx.CreateBucketIfNotExists(recoveryCheckpointBucket) |
| 238 | if err != nil { |
| 239 | return err |
| 240 | } |
| 241 | operations, err := tx.CreateBucketIfNotExists(recoveryOperationBucket) |
| 242 | if err != nil { |
| 243 | return err |
| 244 | } |
| 245 | _ = checkpoints |
| 246 | _ = operations |
| 247 | stored := string(meta.Get([]byte("storage_generation"))) |
| 248 | if stored != "" && stored != identity.Generation { |
| 249 | if err := tx.DeleteBucket(recoveryCheckpointBucket); err != nil { |
| 250 | return err |
| 251 | } |
| 252 | if err := tx.DeleteBucket(recoveryOperationBucket); err != nil { |
| 253 | return err |
| 254 | } |
| 255 | if _, err := tx.CreateBucket(recoveryCheckpointBucket); err != nil { |
| 256 | return err |
| 257 | } |
| 258 | if _, err := tx.CreateBucket(recoveryOperationBucket); err != nil { |
| 259 | return err |
| 260 | } |
| 261 | } |
| 262 | return meta.Put([]byte("storage_generation"), []byte(identity.Generation)) |
| 263 | }) |
| 264 | if err != nil { |
| 265 | _ = db.Close() |
| 266 | return nil, err |
| 267 | } |
| 268 | return store, nil |
| 269 | } |
| 270 | |
| 271 | func openRecoveryStoreRepair(sessionDir string, identity storageIdentity) (*recoveryStore, error) { |
| 272 | store, err := openRecoveryStore(sessionDir, identity) |
| 273 | if err == nil { |
| 274 | return store, nil |
| 275 | } |
| 276 | path := filepath.Join(recoveryCacheDir(sessionDir), recoveryDBName) |
| 277 | if _, statErr := os.Stat(path); statErr == nil { |
| 278 | _ = os.Rename(path, path+fmt.Sprintf(".corrupt-%d", time.Now().UTC().UnixNano())) |
| 279 | } |
| 280 | return openRecoveryStore(sessionDir, identity) |
| 281 | } |
| 282 | |
| 283 | func (s *recoveryStore) close() error { |
| 284 | if s == nil || s.db == nil { |
| 285 | return nil |
| 286 | } |
| 287 | return s.db.Close() |
| 288 | } |
| 289 | |
| 290 | func encodeRecoveryValue(value any) ([]byte, error) { |
| 291 | raw, err := json.Marshal(value) |
| 292 | if err != nil { |
| 293 | return nil, err |
| 294 | } |
| 295 | encoder, err := zstd.NewWriter(nil, zstd.WithEncoderConcurrency(1), zstd.WithEncoderLevel(zstd.SpeedFastest)) |
| 296 | if err != nil { |
| 297 | return nil, err |
| 298 | } |
| 299 | defer encoder.Close() |
| 300 | return encoder.EncodeAll(raw, nil), nil |
| 301 | } |
| 302 | |
| 303 | func decodeRecoveryValue(data []byte, value any) error { |
| 304 | decoder, err := zstd.NewReader(nil, zstd.WithDecoderConcurrency(1), zstd.WithDecoderMaxMemory(64<<20)) |
| 305 | if err != nil { |
| 306 | return err |
| 307 | } |
| 308 | defer decoder.Close() |
| 309 | raw, err := decoder.DecodeAll(data, nil) |
| 310 | if err != nil { |
| 311 | return err |
| 312 | } |
| 313 | return json.Unmarshal(raw, value) |
| 314 | } |
| 315 | |
| 316 | func (s *recoveryStore) loadCheckpoints() ([]recoveryCheckpoint, error) { |
| 317 | if s == nil || s.db == nil { |
| 318 | return nil, os.ErrNotExist |
| 319 | } |
| 320 | var encoded [][]byte |
| 321 | err := s.db.View(func(tx *bolt.Tx) error { |
| 322 | bucket := tx.Bucket(recoveryCheckpointBucket) |
| 323 | if bucket == nil { |
| 324 | return os.ErrNotExist |
| 325 | } |
| 326 | for _, key := range [][]byte{recoveryCurrentKey, recoveryPreviousKey} { |
| 327 | if value := bucket.Get(key); value != nil { |
| 328 | encoded = append(encoded, append([]byte(nil), value...)) |
| 329 | } |
| 330 | } |
| 331 | if len(encoded) == 0 { |
| 332 | return os.ErrNotExist |
| 333 | } |
| 334 | return nil |
| 335 | }) |
| 336 | if err != nil { |
| 337 | return nil, err |
| 338 | } |
| 339 | checkpoints := make([]recoveryCheckpoint, 0, len(encoded)) |
| 340 | for _, value := range encoded { |
| 341 | var checkpoint recoveryCheckpoint |
| 342 | if decodeRecoveryValue(value, &checkpoint) == nil { |
| 343 | checkpoints = append(checkpoints, checkpoint) |
| 344 | } |
| 345 | } |
| 346 | if len(checkpoints) == 0 { |
| 347 | return nil, ErrDamagedStore |
| 348 | } |
| 349 | return checkpoints, nil |
| 350 | } |
| 351 | |
| 352 | func (s *recoveryStore) lookupOperation(operationID string) (operationRecord, bool, error) { |
| 353 | if s == nil || s.db == nil { |
| 354 | return operationRecord{}, false, nil |
| 355 | } |
| 356 | var encoded []byte |
| 357 | err := s.db.View(func(tx *bolt.Tx) error { |
| 358 | bucket := tx.Bucket(recoveryOperationBucket) |
| 359 | if bucket == nil { |
| 360 | return nil |
| 361 | } |
| 362 | encoded = append(encoded, bucket.Get([]byte(operationID))...) |
| 363 | return nil |
| 364 | }) |
| 365 | if err != nil || len(encoded) == 0 { |
| 366 | return operationRecord{}, false, err |
| 367 | } |
| 368 | var operation recoveryOperation |
| 369 | if err := json.Unmarshal(encoded, &operation); err != nil { |
| 370 | return operationRecord{}, false, err |
| 371 | } |
| 372 | return operation.record(), true, nil |
| 373 | } |
| 374 | |
| 375 | func (s *recoveryStore) publish(ctx context.Context, checkpoint recoveryCheckpoint, operations map[string]operationRecord) error { |
| 376 | if s == nil || s.db == nil { |
| 377 | return errors.New("session: recovery store unavailable") |
| 378 | } |
| 379 | if err := ctx.Err(); err != nil { |
| 380 | return err |
| 381 | } |
| 382 | if s.identity.LogPrefixDigest == "" { |
| 383 | manifest, err := readStoredManifest(filepath.Join(s.sessionDir, "manifest.json")) |
| 384 | if err != nil { |
| 385 | return err |
| 386 | } |
| 387 | prefix, err := storageLogPrefix(s.sessionDir, manifest) |
| 388 | if err != nil { |
| 389 | return err |
| 390 | } |
| 391 | if prefix != "" { |
| 392 | s.identity.LogPrefixDigest = prefix |
| 393 | encodedIdentity, err := json.Marshal(s.identity) |
| 394 | if err != nil { |
| 395 | return err |
| 396 | } |
| 397 | if err := fileutil.AtomicWriteFileStrict(filepath.Join(s.sessionDir, storageIdentityName), append(encodedIdentity, '\n'), 0o600); err != nil { |
| 398 | return err |
| 399 | } |
| 400 | } |
| 401 | } |
| 402 | checkpoint.Version = recoveryFormatVersion |
| 403 | checkpoint.StorageGeneration = s.identity.Generation |
| 404 | checkpoint.CreatedAt = time.Now().UTC() |
| 405 | encoded, err := encodeRecoveryValue(checkpoint) |
| 406 | if err != nil { |
| 407 | return err |
| 408 | } |
| 409 | keys := make([]string, 0, len(operations)) |
| 410 | for key := range operations { |
| 411 | keys = append(keys, key) |
| 412 | } |
| 413 | sort.Strings(keys) |
| 414 | err = s.db.Update(func(tx *bolt.Tx) error { |
| 415 | checkpoints := tx.Bucket(recoveryCheckpointBucket) |
| 416 | operationBucket := tx.Bucket(recoveryOperationBucket) |
| 417 | if checkpoints == nil || operationBucket == nil { |
| 418 | return errors.New("session: recovery buckets unavailable") |
| 419 | } |
| 420 | if current := checkpoints.Get(recoveryCurrentKey); current != nil { |
| 421 | if err := checkpoints.Put(recoveryPreviousKey, current); err != nil { |
| 422 | return err |
| 423 | } |
| 424 | } |
| 425 | for _, key := range keys { |
| 426 | operation := operationForRecovery(operations[key]) |
| 427 | value, err := json.Marshal(operation) |
| 428 | if err != nil { |
| 429 | return err |
| 430 | } |
| 431 | if err := operationBucket.Put([]byte(key), value); err != nil { |
| 432 | return err |
| 433 | } |
| 434 | } |
| 435 | if err := checkpoints.Put(recoveryCurrentKey, encoded); err != nil { |
| 436 | return err |
| 437 | } |
| 438 | return tx.Bucket(recoveryMetaBucket).Put([]byte("coverage_sequence"), fmt.Append(nil, checkpoint.DurableSequence)) |
| 439 | }) |
| 440 | if err != nil { |
| 441 | return err |
| 442 | } |
| 443 | recent := RecentSnapshot{ |
| 444 | Version: recoveryFormatVersion, SessionID: checkpoint.SessionID, |
| 445 | StorageGeneration: checkpoint.StorageGeneration, DurableSequence: checkpoint.DurableSequence, |
| 446 | Title: checkpoint.Projection.Title, |
| 447 | ModelRef: checkpoint.Projection.ModelRef, ModelIdentity: checkpoint.Projection.ModelIdentity, |
| 448 | TotalTurns: visibleBoundaryCount(checkpoint.Projection, false), |
| 449 | } |
| 450 | recent.Entries, err = buildRecentEntries(ctx, s.sessionDir, checkpoint.RecentMessages, checkpoint.DurableSequence, recent.TotalTurns) |
| 451 | attachSubmissionEntries(checkpoint.Projection.Submissions, checkpoint.SessionID, recent.Entries) |
| 452 | if err != nil { |
| 453 | return err |
| 454 | } |
| 455 | data, err := json.Marshal(recent) |
| 456 | if err != nil { |
| 457 | return err |
| 458 | } |
| 459 | return fileutil.AtomicWriteFileStrict(s.recent, append(data, '\n'), 0o600) |
| 460 | } |
| 461 | |
| 462 | func readRecentSnapshot(sessionDir string, identity storageIdentity) (RecentSnapshot, error) { |
| 463 | data, err := os.ReadFile(filepath.Join(recoveryCacheDir(sessionDir), recentSnapshotName)) |
| 464 | if err != nil { |
| 465 | return RecentSnapshot{}, err |
| 466 | } |
| 467 | var snapshot RecentSnapshot |
| 468 | if err := json.Unmarshal(data, &snapshot); err != nil { |
| 469 | return RecentSnapshot{}, err |
| 470 | } |
| 471 | if snapshot.Version != recoveryFormatVersion || snapshot.SessionID != identity.SessionID || snapshot.StorageGeneration != identity.Generation { |
| 472 | return RecentSnapshot{}, ErrStaleGeneration |
| 473 | } |
| 474 | if len(snapshot.Entries) > RecentMessageLimit { |
| 475 | return RecentSnapshot{}, ErrDamagedStore |
| 476 | } |
| 477 | return snapshot, nil |
| 478 | } |
| 479 | |
| 480 | // buildRecentEntries produces the bounded public baseline. Large canonical |
| 481 | // messages are stored once in ContentStore and represented by a preview plus a |
| 482 | // range-readable reference, so recent-v1.json cannot grow with tool output. |
| 483 | func buildRecentEntries(ctx context.Context, sessionDir string, messages []provider.Message, sequence uint64, totalTurns int) ([]PersistentMessage, error) { |
| 484 | if len(messages) > RecentMessageLimit { |
| 485 | messages = messages[len(messages)-RecentMessageLimit:] |
| 486 | } |
| 487 | entries := make([]PersistentMessage, 0, len(messages)) |
| 488 | content := contentStoreForSessionDir(sessionDir) |
| 489 | inlineBytes := 0 |
| 490 | visibleTurn := totalTurns |
| 491 | for _, message := range messages { |
| 492 | if agent.IsUserAuthoredTurnMessage(message) { |
| 493 | visibleTurn-- |
| 494 | } |
| 495 | } |
| 496 | visibleTurn = max(visibleTurn, 0) |
| 497 | for position, message := range messages { |
| 498 | if err := ctx.Err(); err != nil { |
| 499 | return nil, err |
| 500 | } |
| 501 | if agent.IsUserAuthoredTurnMessage(message) { |
| 502 | visibleTurn++ |
| 503 | } |
| 504 | body, err := json.Marshal(message) |
| 505 | if err != nil { |
| 506 | return nil, err |
| 507 | } |
| 508 | entry := PersistentMessage{ |
| 509 | MessageID: message.ID, Position: int64(position + 1), Version: 1, |
| 510 | Role: string(message.Role), Preview: messagePreview(message), |
| 511 | EventSequence: sequence, VisibleTurn: visibleTurn, |
| 512 | } |
| 513 | if len(body) <= recentInlineBytes && inlineBytes+len(body) <= recentResponseBytes { |
| 514 | entry.Inline = body |
| 515 | inlineBytes += len(body) |
| 516 | } else { |
| 517 | ref, err := content.Put(ctx, bytes.NewReader(body), sessioncontent.Metadata{MediaType: "application/json"}) |
| 518 | if err != nil { |
| 519 | return nil, err |
| 520 | } |
| 521 | entry.ContentRef = &ref |
| 522 | previewBody, err := recentDisplayMessage(message) |
| 523 | if err != nil { |
| 524 | return nil, err |
| 525 | } |
| 526 | if inlineBytes+len(previewBody) <= recentResponseBytes { |
| 527 | entry.Inline = previewBody |
| 528 | inlineBytes += len(previewBody) |
| 529 | } |
| 530 | } |
| 531 | entries = append(entries, entry) |
| 532 | } |
| 533 | return entries, nil |
| 534 | } |
| 535 | |
| 536 | func recentDisplayMessage(message provider.Message) (json.RawMessage, error) { |
| 537 | preview := detachMessages([]provider.Message{message})[0] |
| 538 | preview.Content = messagePreview(message) |
| 539 | preview.RawContent = "" |
| 540 | preview.ProviderContent = "" |
| 541 | preview.Images = nil |
| 542 | preview.ImageInputs = nil |
| 543 | preview.ResponsesItems = nil |
| 544 | preview.ThinkingBlocks = nil |
| 545 | if runes := []rune(preview.ReasoningContent); len(runes) > 4096 { |
| 546 | preview.ReasoningContent = string(runes[:4096]) |
| 547 | } |
| 548 | for i := range preview.ToolCalls { |
| 549 | if runes := []rune(preview.ToolCalls[i].Arguments); len(runes) > 2048 { |
| 550 | preview.ToolCalls[i].Arguments = string(runes[:2048]) |
| 551 | } |
| 552 | } |
| 553 | body, err := json.Marshal(preview) |
| 554 | if err != nil { |
| 555 | return nil, err |
| 556 | } |
| 557 | if len(body) <= recentInlineBytes { |
| 558 | return body, nil |
| 559 | } |
| 560 | // Preserve the fields required to place the row even when optional display |
| 561 | // metadata alone exceeds the per-message preview budget. |
| 562 | return json.Marshal(provider.Message{ |
| 563 | ID: preview.ID, Role: preview.Role, Content: preview.Content, |
| 564 | ToolCallID: preview.ToolCallID, Name: preview.Name, |
| 565 | CreatedAt: preview.CreatedAt, WorkDurationMs: preview.WorkDurationMs, |
| 566 | }) |
| 567 | } |
| 568 | |
| 569 | func checkpointFromStartup(manifest Manifest, identity storageIdentity, state *startupSessionState) recoveryCheckpoint { |
| 570 | projection, _ := Project(nil) |
| 571 | if state != nil { |
| 572 | projection = cloneProjection(state.projection) |
| 573 | projection.Messages = nil |
| 574 | } |
| 575 | checkpoint := recoveryCheckpoint{ |
| 576 | Version: recoveryFormatVersion, SessionID: manifest.SessionID, |
| 577 | StorageGeneration: identity.Generation, StorageRevision: StorageRevision, |
| 578 | ProjectionVersion: recoveryProjectionVersion, Projection: projection, |
| 579 | } |
| 580 | if state != nil { |
| 581 | checkpoint.DurableSequence = state.durable |
| 582 | checkpoint.LogOffset = state.tip.LogOffset |
| 583 | checkpoint.AnchorOffset = state.tip.AnchorOffset |
| 584 | checkpoint.AnchorFirst = state.tip.AnchorFirst |
| 585 | checkpoint.AnchorCommitID = state.tip.AnchorCommitID |
| 586 | checkpoint.AnchorHash = state.tip.AnchorHash |
| 587 | checkpoint.RecentMessages = detachMessages(state.recentMessages) |
| 588 | checkpoint.CatalogPreview = state.catalogPreview |
| 589 | checkpoint.MessageIDs = state.messageIDs.list() |
| 590 | } |
| 591 | return checkpoint |
| 592 | } |
| 593 | |
| 594 | func loadRecoveryStartupState(ctx context.Context, dir string, file *os.File, info os.FileInfo, recovery *recoveryStore, identity storageIdentity) (*startupSessionState, int64, bool, RecoveryOpenStats, bool) { |
| 595 | stats := RecoveryOpenStats{LogBytesTotal: info.Size()} |
| 596 | checkpoints, err := recovery.loadCheckpoints() |
| 597 | if err != nil { |
| 598 | return nil, 0, false, stats, false |
| 599 | } |
| 600 | for _, checkpoint := range checkpoints { |
| 601 | if state, end, torn, attempt, ok := tryRecoveryCheckpoint(ctx, dir, file, info, identity, checkpoint); ok { |
| 602 | return state, end, torn, attempt, true |
| 603 | } |
| 604 | } |
| 605 | return nil, 0, false, stats, false |
| 606 | } |
| 607 | |
| 608 | func tryRecoveryCheckpoint(ctx context.Context, dir string, file *os.File, info os.FileInfo, identity storageIdentity, checkpoint recoveryCheckpoint) (*startupSessionState, int64, bool, RecoveryOpenStats, bool) { |
| 609 | stats := RecoveryOpenStats{LogBytesTotal: info.Size()} |
| 610 | if checkpoint.Version != recoveryFormatVersion || checkpoint.ProjectionVersion != recoveryProjectionVersion || |
| 611 | checkpoint.SessionID != identity.SessionID || checkpoint.StorageGeneration != identity.Generation || |
| 612 | checkpoint.StorageRevision != StorageRevision || checkpoint.LogOffset < 0 || checkpoint.LogOffset > info.Size() || |
| 613 | checkpoint.Projection.CommittedSequence != checkpoint.DurableSequence { |
| 614 | return nil, 0, false, stats, false |
| 615 | } |
| 616 | if checkpoint.DurableSequence == 0 { |
| 617 | if checkpoint.LogOffset != 0 { |
| 618 | return nil, 0, false, stats, false |
| 619 | } |
| 620 | } else { |
| 621 | if checkpoint.AnchorOffset < 0 || checkpoint.AnchorOffset >= checkpoint.LogOffset || checkpoint.AnchorFirst == 0 || checkpoint.AnchorCommitID == "" { |
| 622 | return nil, 0, false, stats, false |
| 623 | } |
| 624 | var anchor Commit |
| 625 | var anchorEnd int64 |
| 626 | err := scanV4CommitFileRefs(ctx, file, checkpoint.AnchorOffset, checkpoint.AnchorFirst, contentStoreForSessionDir(dir), nil, func(_ int64, commit Commit) bool { |
| 627 | anchor = commit |
| 628 | anchorEnd, _ = file.Seek(0, 1) |
| 629 | return false |
| 630 | }) |
| 631 | if err != nil || anchor.ID != checkpoint.AnchorCommitID || anchor.OperationHash != checkpoint.AnchorHash || |
| 632 | anchor.LastSequence() != checkpoint.DurableSequence || anchorEnd != checkpoint.LogOffset { |
| 633 | return nil, 0, false, stats, false |
| 634 | } |
| 635 | stats.LogBytesRead += anchorEnd - checkpoint.AnchorOffset |
| 636 | } |
| 637 | |
| 638 | state := &startupSessionState{ |
| 639 | projection: cloneProjection(checkpoint.Projection), operations: map[string]operationRecord{}, |
| 640 | durable: checkpoint.DurableSequence, catalogPreview: checkpoint.CatalogPreview, |
| 641 | recentMessages: detachMessages(checkpoint.RecentMessages), |
| 642 | messageIDs: identitiesOf(checkpoint.MessageIDs), |
| 643 | tip: durableTip{LogOffset: checkpoint.LogOffset, AnchorOffset: checkpoint.AnchorOffset, |
| 644 | AnchorFirst: checkpoint.AnchorFirst, AnchorCommitID: checkpoint.AnchorCommitID, AnchorHash: checkpoint.AnchorHash}, |
| 645 | } |
| 646 | var projectionErr error |
| 647 | content := contentStoreForSessionDir(dir) |
| 648 | repeated := map[uint64]bool{} |
| 649 | err := scanV4CommitFile(ctx, file, checkpoint.LogOffset, checkpoint.DurableSequence+1, content, nil, func(offset int64, commit Commit) bool { |
| 650 | kept := commit |
| 651 | kept.Events = slices.DeleteFunc(slices.Clone(commit.Events), func(event Event) bool { return !state.messageIDs.admitEvent(event, repeated) }) |
| 652 | if err := applyProjectionCommit(&state.projection, kept); err != nil { |
| 653 | projectionErr = err |
| 654 | return false |
| 655 | } |
| 656 | if err := applyRecentCommit(&state.recentMessages, kept); err != nil { |
| 657 | projectionErr = err |
| 658 | return false |
| 659 | } |
| 660 | state.operations[commit.OperationID] = compactOperationRecord(commit) |
| 661 | state.durable = commit.LastSequence() |
| 662 | state.tip.AnchorOffset = offset |
| 663 | state.tip.AnchorFirst = commit.FirstSequence |
| 664 | state.tip.AnchorCommitID = commit.ID |
| 665 | state.tip.AnchorHash = commit.OperationHash |
| 666 | state.tip.LogOffset, _ = file.Seek(0, 1) |
| 667 | stats.TailCommits++ |
| 668 | return true |
| 669 | }) |
| 670 | if err != nil || projectionErr != nil { |
| 671 | return nil, 0, false, stats, false |
| 672 | } |
| 673 | state.projection.Messages = nil |
| 674 | state.projection.CommittedSequence = state.durable |
| 675 | stats.UsedCheckpoint = true |
| 676 | stats.LogBytesRead += max(state.tip.LogOffset-checkpoint.LogOffset, 0) |
| 677 | return state, state.tip.LogOffset, state.tip.LogOffset < info.Size(), stats, true |
| 678 | } |
| 679 | |
| 680 | func applyRecentCommit(messages *[]provider.Message, commit Commit) error { |
| 681 | projection, _ := Project(nil) |
| 682 | projection.Messages = detachMessages(*messages) |
| 683 | recent := commit |
| 684 | recent.Events = nil |
| 685 | for _, event := range commit.Events { |
| 686 | switch event.Kind { |
| 687 | case "message/complete", "message/upsert", "message/retract", "history/replace", "legacy/import": |
| 688 | recent.Events = append(recent.Events, event) |
| 689 | } |
| 690 | } |
| 691 | if len(recent.Events) == 0 { |
| 692 | return nil |
| 693 | } |
| 694 | if err := applyProjectionCommit(&projection, recent); err != nil { |
| 695 | return err |
| 696 | } |
| 697 | if len(projection.Messages) > RecentMessageLimit { |
| 698 | projection.Messages = append([]provider.Message(nil), projection.Messages[len(projection.Messages)-RecentMessageLimit:]...) |
| 699 | } |
| 700 | *messages = detachMessages(projection.Messages) |
| 701 | return nil |
| 702 | } |
| 703 |