| 1 | package agent |
| 2 | |
| 3 | import ( |
| 4 | "bufio" |
| 5 | "context" |
| 6 | "crypto/sha256" |
| 7 | "database/sql" |
| 8 | "encoding" |
| 9 | "encoding/json" |
| 10 | "errors" |
| 11 | "fmt" |
| 12 | "io" |
| 13 | "os" |
| 14 | |
| 15 | "reasonix/internal/fileops" |
| 16 | "reasonix/internal/historywork" |
| 17 | "reasonix/internal/provider" |
| 18 | ) |
| 19 | |
| 20 | // Build an offset-only projection for a checkpoint without publishing or |
| 21 | // repairing any session sidecar. At most one message and one SQLite batch are |
| 22 | // retained. Digest validation happens before the disposable index is published. |
| 23 | func buildCheckpointDisplayPager(ctx context.Context, db *sql.DB, source, fingerprint string) error { |
| 24 | return buildCheckpointDisplayPagerObserved(ctx, db, source, fingerprint, nil) |
| 25 | } |
| 26 | |
| 27 | // The observer is used by deterministic interruption tests after durable batch |
| 28 | // publication. Production callers do not install one. |
| 29 | func buildCheckpointDisplayPagerObserved(ctx context.Context, db *sql.DB, source, fingerprint string, committed func(int)) error { |
| 30 | f, err := fileops.OpenReplaceableRead(source) |
| 31 | if err != nil { |
| 32 | return err |
| 33 | } |
| 34 | defer f.Close() |
| 35 | if err := validateCheckpointPagerFile(source, f, fingerprint); err != nil { |
| 36 | return err |
| 37 | } |
| 38 | hash := sha256.New() |
| 39 | idx := SessionDisplayIndex{SchemaVersion: SessionDisplayIndexSchemaVersion, ListingPreviewKnown: true} |
| 40 | users := 0 |
| 41 | if err := restoreCheckpointPagerProgress(ctx, db, &idx, &users, hash); err != nil { |
| 42 | return err |
| 43 | } |
| 44 | if _, err := f.Seek(idx.TranscriptSize, io.SeekStart); err != nil { |
| 45 | return err |
| 46 | } |
| 47 | reader := bufio.NewReaderSize(&historywork.Reader{Context: ctx, Source: f}, historywork.ReadChunk) |
| 48 | done := false |
| 49 | for !done { |
| 50 | tx, err := db.BeginTx(ctx, nil) |
| 51 | if err != nil { |
| 52 | return err |
| 53 | } |
| 54 | for range historywork.BatchEntries { |
| 55 | line, readErr := readSessionDisplayIndexLine(reader) |
| 56 | if readErr != nil && !errors.Is(readErr, io.EOF) { |
| 57 | err = readErr |
| 58 | break |
| 59 | } |
| 60 | if len(line) > 0 { |
| 61 | var message provider.Message |
| 62 | if err = json.Unmarshal(line, &message); err != nil { |
| 63 | break |
| 64 | } |
| 65 | if message.Role == "" { |
| 66 | err = errors.New("checkpoint message has no role") |
| 67 | // Ancient event rows use kind/type instead of provider roles. |
| 68 | // Classify only the first record as an unsupported format; |
| 69 | // a foreign record after valid messages is damaged content. |
| 70 | if idx.MessageCount == 0 { |
| 71 | var event struct { |
| 72 | Kind string `json:"kind"` |
| 73 | Type string `json:"type"` |
| 74 | } |
| 75 | if json.Unmarshal(line, &event) == nil && (event.Kind != "" || event.Type != "") { |
| 76 | err = ErrDisplayFormatUnsupported |
| 77 | } |
| 78 | } |
| 79 | break |
| 80 | } |
| 81 | var entry DisplayIndexEntry |
| 82 | entry, idx.AuthoredTurns = classifyDisplayIndexMessage(message, idx.MessageCount, idx.TranscriptSize, int64(len(line)), idx.AuthoredTurns) |
| 83 | var encoded []byte |
| 84 | encoded, err = json.Marshal(entry) |
| 85 | if err != nil { |
| 86 | break |
| 87 | } |
| 88 | _, err = tx.ExecContext(ctx, `INSERT INTO entries VALUES(?,?,?,?,?,?,?)`, entry.Index, entry.Offset, entry.Length, entry.AuthoredTurn, entry.Role, users, encoded) |
| 89 | if err != nil { |
| 90 | break |
| 91 | } |
| 92 | encoded, err = json.Marshal(messageForSessionIdentity(message)) |
| 93 | if err != nil { |
| 94 | break |
| 95 | } |
| 96 | _, _ = hash.Write(encoded) |
| 97 | _, _ = hash.Write([]byte{'\n'}) |
| 98 | if entry.Role == provider.RoleUser && !entry.PinnedContextRevision { |
| 99 | users++ |
| 100 | } |
| 101 | if entry.StartsTurn && idx.ListingPreview == "" { |
| 102 | idx.ListingPreview = truncatePreview(previewProse(UserMessageText(message))) |
| 103 | } |
| 104 | idx.MessageCount++ |
| 105 | idx.TranscriptSize += int64(len(line)) |
| 106 | } |
| 107 | if errors.Is(readErr, io.EOF) { |
| 108 | done = true |
| 109 | break |
| 110 | } |
| 111 | } |
| 112 | if err != nil { |
| 113 | _ = tx.Rollback() |
| 114 | return fmt.Errorf("build checkpoint display index: %w", err) |
| 115 | } |
| 116 | state, marshalErr := hash.(encoding.BinaryMarshaler).MarshalBinary() |
| 117 | if marshalErr == nil { |
| 118 | marshalErr = saveCheckpointPagerProgress(ctx, tx, idx, users, state) |
| 119 | } |
| 120 | if marshalErr != nil { |
| 121 | _ = tx.Rollback() |
| 122 | return marshalErr |
| 123 | } |
| 124 | if err := tx.Commit(); err != nil { |
| 125 | return err |
| 126 | } |
| 127 | if committed != nil { |
| 128 | committed(idx.MessageCount) |
| 129 | } |
| 130 | } |
| 131 | return publishCheckpointDisplayPager(ctx, db, source, f, fingerprint, idx, hash) |
| 132 | } |
| 133 | |
| 134 | // Fence both the handle that supplied the bytes and its current pathname |
| 135 | // before publishing. A replacement or a newly authoritative event log cannot |
| 136 | // turn a resumed checkpoint prefix into a complete generation. |
| 137 | func validateCheckpointPagerFile(source string, f *os.File, fingerprint string) error { |
| 138 | info, err := f.Stat() |
| 139 | if err != nil { |
| 140 | return err |
| 141 | } |
| 142 | target, version := fileops.DiskHandleSnapshot(source, f, info) |
| 143 | if fmt.Sprintf("%s:%s:checkpoint", target.Key, version) != fingerprint { |
| 144 | return ErrDisplaySourceChanged |
| 145 | } |
| 146 | current, err := os.Stat(source) |
| 147 | if err != nil { |
| 148 | return err |
| 149 | } |
| 150 | target, version = fileops.DiskSnapshot(source, current) |
| 151 | if !os.SameFile(info, current) || fmt.Sprintf("%s:%s:checkpoint", target.Key, version) != fingerprint { |
| 152 | return ErrDisplaySourceChanged |
| 153 | } |
| 154 | return validateCheckpointDisplaySource(source) |
| 155 | } |
| 156 |