| 1 | package agent |
| 2 | |
| 3 | import ( |
| 4 | "bufio" |
| 5 | "context" |
| 6 | "crypto/sha256" |
| 7 | "database/sql" |
| 8 | "encoding/json" |
| 9 | "errors" |
| 10 | "fmt" |
| 11 | "io" |
| 12 | "os" |
| 13 | |
| 14 | "reasonix/internal/fileops" |
| 15 | "reasonix/internal/historywork" |
| 16 | "reasonix/internal/provider" |
| 17 | "reasonix/internal/store" |
| 18 | ) |
| 19 | |
| 20 | // Schema-1 replace/append logs retain original message locations. A large |
| 21 | // replace array is decoded one message at a time; SQLite holds the live order. |
| 22 | // No migration, normalization, or repair writes are made to source files. |
| 23 | func buildEventDisplayPager(ctx context.Context, db *sql.DB, source, fingerprint string, checkpointSize int64) error { |
| 24 | return buildEventDisplayPagerObserved(ctx, db, source, fingerprint, checkpointSize, nil) |
| 25 | } |
| 26 | |
| 27 | func buildEventDisplayPagerObserved(ctx context.Context, db *sql.DB, source, fingerprint string, checkpointSize int64, observed func(string, int)) error { |
| 28 | f, err := fileops.OpenReplaceableRead(store.SessionEventLog(source)) |
| 29 | if err != nil { |
| 30 | return err |
| 31 | } |
| 32 | defer f.Close() |
| 33 | if err := validateEventPagerFile(source, f, fingerprint); err != nil { |
| 34 | return err |
| 35 | } |
| 36 | info, err := f.Stat() |
| 37 | if err != nil { |
| 38 | return err |
| 39 | } |
| 40 | progress, err := restoreEventPagerScan(ctx, db, fingerprint, info.Size()) |
| 41 | if err != nil { |
| 42 | return err |
| 43 | } |
| 44 | scanner := &eventPagerScanner{ctx: ctx, db: db, progress: progress, observed: observed} |
| 45 | if err := scanner.scan(f); err != nil { |
| 46 | return err |
| 47 | } |
| 48 | if err := finishEventDisplayPager(ctx, db, f, source, fingerprint, checkpointSize, scanner.progress.Count, observed); err != nil { |
| 49 | return err |
| 50 | } |
| 51 | return validateEventPagerFile(source, f, fingerprint) |
| 52 | } |
| 53 | |
| 54 | // Unknown event fields are skipped token by token rather than buffering the |
| 55 | // entire field/record in json.RawMessage. json.Decoder still validates nesting. |
| 56 | func skipDisplayJSONValue(ctx context.Context, decoder *json.Decoder) error { |
| 57 | depth := 0 |
| 58 | for { |
| 59 | if err := ctx.Err(); err != nil { |
| 60 | return err |
| 61 | } |
| 62 | token, err := decoder.Token() |
| 63 | if err != nil { |
| 64 | return err |
| 65 | } |
| 66 | if delim, ok := token.(json.Delim); ok { |
| 67 | switch delim { |
| 68 | case '{', '[': |
| 69 | depth++ |
| 70 | case '}', ']': |
| 71 | depth-- |
| 72 | } |
| 73 | } |
| 74 | if depth <= 0 { |
| 75 | return nil |
| 76 | } |
| 77 | } |
| 78 | } |
| 79 | |
| 80 | type displayEventLocation struct { |
| 81 | Position int |
| 82 | Offset, Length int64 |
| 83 | At int64 |
| 84 | } |
| 85 | |
| 86 | func eventDisplayLocations(ctx context.Context, db *sql.DB, lo, hi int) ([]displayEventLocation, error) { |
| 87 | rows, err := db.QueryContext(ctx, `SELECT position,offset,length,at FROM event_locations WHERE position>=? AND position<? ORDER BY position`, lo, hi) |
| 88 | if err != nil { |
| 89 | return nil, err |
| 90 | } |
| 91 | defer rows.Close() |
| 92 | locations := make([]displayEventLocation, 0, hi-lo) |
| 93 | for rows.Next() { |
| 94 | var loc displayEventLocation |
| 95 | if err := rows.Scan(&loc.Position, &loc.Offset, &loc.Length, &loc.At); err != nil { |
| 96 | return nil, err |
| 97 | } |
| 98 | locations = append(locations, loc) |
| 99 | } |
| 100 | return locations, rows.Err() |
| 101 | } |
| 102 | |
| 103 | func readDisplayEventMessage(ctx context.Context, file *os.File, loc displayEventLocation) (provider.Message, error) { |
| 104 | reader := bufio.NewReaderSize(&historywork.Reader{Context: ctx, Source: io.NewSectionReader(file, loc.Offset, loc.Length)}, historywork.ReadChunk) |
| 105 | // InputOffset may precede the comma separating adjacent array elements. |
| 106 | for { |
| 107 | prefix, err := reader.Peek(1) |
| 108 | if err != nil { |
| 109 | return provider.Message{}, err |
| 110 | } |
| 111 | if prefix[0] != ',' && prefix[0] != ' ' && prefix[0] != '\n' && prefix[0] != '\r' && prefix[0] != '\t' { |
| 112 | break |
| 113 | } |
| 114 | _, _ = reader.ReadByte() |
| 115 | } |
| 116 | var message provider.Message |
| 117 | err := json.NewDecoder(reader).Decode(&message) |
| 118 | return message, err |
| 119 | } |
| 120 | |
| 121 | func finishEventDisplayPager(ctx context.Context, db *sql.DB, file *os.File, source, fingerprint string, checkpointSize int64, count int, observed func(string, int)) error { |
| 122 | idx := SessionDisplayIndex{SchemaVersion: SessionDisplayIndexSchemaVersion, TranscriptSize: checkpointSize, ListingPreviewKnown: true} |
| 123 | hash := sha256.New() |
| 124 | users := 0 |
| 125 | if err := restoreEventProjection(ctx, db, &idx, &users, hash, count); err != nil { |
| 126 | return err |
| 127 | } |
| 128 | for lo := idx.MessageCount; lo < count; lo += historywork.BatchEntries { |
| 129 | locations, err := eventDisplayLocations(ctx, db, lo, min(lo+historywork.BatchEntries, count)) |
| 130 | if err != nil { |
| 131 | return err |
| 132 | } |
| 133 | tx, err := db.BeginTx(ctx, nil) |
| 134 | if err != nil { |
| 135 | return err |
| 136 | } |
| 137 | for _, loc := range locations { |
| 138 | message, readErr := readDisplayEventMessage(ctx, file, loc) |
| 139 | if readErr != nil { |
| 140 | err = readErr |
| 141 | break |
| 142 | } |
| 143 | var entry DisplayIndexEntry |
| 144 | entry, idx.AuthoredTurns = classifyDisplayIndexMessage(message, idx.MessageCount, loc.Offset, loc.Length, idx.AuthoredTurns) |
| 145 | body, _ := json.Marshal(entry) |
| 146 | if _, err = tx.ExecContext(ctx, `INSERT INTO entries VALUES(?,?,?,?,?,?,?)`, entry.Index, entry.Offset, entry.Length, entry.AuthoredTurn, entry.Role, users, body); err != nil { |
| 147 | break |
| 148 | } |
| 149 | body, err = json.Marshal(messageForSessionIdentity(message)) |
| 150 | if err != nil { |
| 151 | break |
| 152 | } |
| 153 | hash.Write(body) |
| 154 | hash.Write([]byte{'\n'}) |
| 155 | if entry.Role == provider.RoleUser && !entry.PinnedContextRevision { |
| 156 | users++ |
| 157 | } |
| 158 | if entry.StartsTurn && idx.ListingPreview == "" { |
| 159 | idx.ListingPreview = truncatePreview(previewProse(UserMessageText(message))) |
| 160 | } |
| 161 | idx.MessageCount++ |
| 162 | } |
| 163 | if err == nil { |
| 164 | err = saveEventProjection(ctx, tx, idx, users, hash) |
| 165 | } |
| 166 | if err != nil { |
| 167 | _ = tx.Rollback() |
| 168 | return err |
| 169 | } |
| 170 | if err := tx.Commit(); err != nil { |
| 171 | return err |
| 172 | } |
| 173 | if observed != nil { |
| 174 | observed("projection", idx.MessageCount) |
| 175 | } |
| 176 | } |
| 177 | if err := ctx.Err(); err != nil { |
| 178 | return err |
| 179 | } |
| 180 | idx.ContentDigest = fmt.Sprintf("%x", hash.Sum(nil)) |
| 181 | identity, known, err := SessionContentIdentity(source) |
| 182 | if err != nil { |
| 183 | return err |
| 184 | } |
| 185 | if known { |
| 186 | if idx.ContentDigest != identity.DigestHex { |
| 187 | return ErrDisplaySourceChanged |
| 188 | } |
| 189 | idx.Revision, idx.RevisionKnown = identity.Revision, identity.RevisionKnown |
| 190 | } |
| 191 | body, err := json.Marshal(idx) |
| 192 | if err != nil { |
| 193 | return err |
| 194 | } |
| 195 | _, err = db.ExecContext(ctx, `INSERT OR REPLACE INTO metadata VALUES('source',?),('header',?),('kind','schema1'); DROP TABLE IF EXISTS event_pending`, fingerprint, string(body)) |
| 196 | return err |
| 197 | } |
| 198 | |
| 199 | // EventMessages reads a bounded range from either native event-log format. |
| 200 | func (p *DisplayPager) EventMessages(lo, hi int) ([]provider.Message, error) { |
| 201 | if p.DAG { |
| 202 | return p.DAGMessages(lo, hi) |
| 203 | } |
| 204 | if !p.SchemaOne || lo < 0 || hi < lo || hi-lo > 500 { |
| 205 | return nil, errors.New("invalid event display window") |
| 206 | } |
| 207 | if err := p.Validate(); err != nil { |
| 208 | return nil, err |
| 209 | } |
| 210 | locations, err := eventDisplayLocations(p.ctx, p.DB, lo, hi) |
| 211 | if err != nil { |
| 212 | return nil, err |
| 213 | } |
| 214 | file, err := fileops.OpenReplaceableRead(store.SessionEventLog(p.source)) |
| 215 | if err != nil { |
| 216 | return nil, err |
| 217 | } |
| 218 | defer file.Close() |
| 219 | messages := make([]provider.Message, 0, len(locations)) |
| 220 | for _, loc := range locations { |
| 221 | message, err := readDisplayEventMessage(p.ctx, file, loc) |
| 222 | if err != nil { |
| 223 | return nil, err |
| 224 | } |
| 225 | if message.CreatedAt <= 0 && loc.At != 0 { |
| 226 | message.CreatedAt = loc.At |
| 227 | } |
| 228 | messages = append(messages, message) |
| 229 | } |
| 230 | return messages, p.Validate() |
| 231 | } |
| 232 |