| 1 | package session |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "database/sql" |
| 6 | "errors" |
| 7 | "strings" |
| 8 | |
| 9 | "reasonix/internal/projectiondb" |
| 10 | "reasonix/internal/textutil" |
| 11 | ) |
| 12 | |
| 13 | func ensureHistoryOutlineIndexes(ctx context.Context, db *sql.DB) error { |
| 14 | tx, err := db.BeginTx(ctx, nil) |
| 15 | if err != nil { |
| 16 | return err |
| 17 | } |
| 18 | defer func() { _ = tx.Rollback() }() |
| 19 | if _, err := tx.ExecContext(ctx, `CREATE INDEX IF NOT EXISTS messages_outline_users ON messages(visible_user,visible_turn,position,event_sequence,valid_to); CREATE INDEX IF NOT EXISTS messages_outline_answers ON messages(visible_turn,role,position DESC,event_sequence,valid_to)`); err != nil { |
| 20 | return err |
| 21 | } |
| 22 | return tx.Commit() |
| 23 | } |
| 24 | |
| 25 | // HistoryOutlineRequest addresses visible turns in a fixed durable cut. |
| 26 | type HistoryOutlineRequest struct { |
| 27 | Generation string `json:"generation,omitempty"` |
| 28 | SnapshotSequence *uint64 `json:"snapshotSequence,omitempty"` |
| 29 | StartTurn int `json:"startTurn,omitempty"` |
| 30 | Limit int `json:"limit,omitempty"` |
| 31 | } |
| 32 | |
| 33 | // HistoryOutlineEntry contains display metadata, never message bodies. |
| 34 | type HistoryOutlineEntry struct { |
| 35 | MessageID string `json:"messageId"` |
| 36 | Turn int `json:"turn"` |
| 37 | Position int64 `json:"position"` |
| 38 | Prompt string `json:"prompt"` |
| 39 | Answer string `json:"answer,omitempty"` |
| 40 | } |
| 41 | |
| 42 | // HistoryOutlinePage shares its identity with canonical history windows. |
| 43 | type HistoryOutlinePage struct { |
| 44 | Status string `json:"status"` |
| 45 | Generation string `json:"generation"` |
| 46 | SnapshotSequence uint64 `json:"snapshotSequence"` |
| 47 | CoverageSequence uint64 `json:"coverageSequence"` |
| 48 | TotalTurns int `json:"totalTurns"` |
| 49 | Entries []HistoryOutlineEntry `json:"entries"` |
| 50 | NextTurn int `json:"nextTurn"` |
| 51 | Done bool `json:"done"` |
| 52 | } |
| 53 | |
| 54 | // ReadHistoryOutline reads indexed previews independently of the bounded live |
| 55 | // projection. The fixed version interval is identical to ReadHistoryWindow. |
| 56 | func (q *Query) ReadHistoryOutline(ctx context.Context, ref SessionRef, req HistoryOutlineRequest) (HistoryOutlinePage, error) { |
| 57 | page := HistoryOutlinePage{Entries: []HistoryOutlineEntry{}} |
| 58 | if q == nil { |
| 59 | return page, errors.New("session: nil query") |
| 60 | } |
| 61 | if err := ref.validate(q.hostID); err != nil { |
| 62 | return page, err |
| 63 | } |
| 64 | filesystem, ok := q.persistence.(*FilesystemPersistence) |
| 65 | if !ok { |
| 66 | return page, errors.New("session: history outline requires filesystem persistence") |
| 67 | } |
| 68 | path := historyIndexPath(filesystem.Root, ref.SessionID) |
| 69 | ready, err := q.historyLocatorReady(ctx, filesystem, ref.SessionID, path) |
| 70 | if err != nil { |
| 71 | return page, err |
| 72 | } |
| 73 | if !ready { |
| 74 | page.Status = "preparing" |
| 75 | return page, nil |
| 76 | } |
| 77 | handle, err := projectiondb.Open(ctx, projectiondb.OpenOptions{Path: path, Migrations: historyMigrations, RequireDisk: true, MaxOpenConns: 1}) |
| 78 | if err != nil { |
| 79 | return page, err |
| 80 | } |
| 81 | defer handle.DB.Close() |
| 82 | // These are optional accelerators, not a data-format migration. Raising |
| 83 | // schema_migrations would make previous readers reject an unchanged format. |
| 84 | if err := ensureHistoryOutlineIndexes(ctx, handle.DB); err != nil { |
| 85 | return page, err |
| 86 | } |
| 87 | // Metadata and both queries belong to one SQLite read transaction. A |
| 88 | // concurrent index update must not mix ordinal counts with another cut. |
| 89 | tx, err := handle.DB.BeginTx(ctx, nil) |
| 90 | if err != nil { |
| 91 | return page, err |
| 92 | } |
| 93 | defer func() { _ = tx.Rollback() }() |
| 94 | var generation string |
| 95 | var durable uint64 |
| 96 | if err := tx.QueryRowContext(ctx, `SELECT (SELECT value FROM metadata WHERE key='generation'),(SELECT value FROM metadata WHERE key='durable_sequence')`).Scan(&generation, &durable); err != nil { |
| 97 | return page, err |
| 98 | } |
| 99 | page.Generation, page.CoverageSequence, page.SnapshotSequence = generation, durable, durable |
| 100 | if req.SnapshotSequence != nil { |
| 101 | page.SnapshotSequence = *req.SnapshotSequence |
| 102 | } |
| 103 | if (req.Generation != "" && req.Generation != generation) || page.SnapshotSequence > durable { |
| 104 | page.Status = "stale_cursor" |
| 105 | return page, nil |
| 106 | } |
| 107 | cut := page.SnapshotSequence |
| 108 | if err := tx.QueryRowContext(ctx, `SELECT COALESCE(MAX(visible_turn),0) FROM messages WHERE visible_user=1 AND event_sequence<=? AND (valid_to=0 OR valid_to>?)`, cut, cut).Scan(&page.TotalTurns); err != nil { |
| 109 | return page, err |
| 110 | } |
| 111 | limit := req.Limit |
| 112 | if limit <= 0 { |
| 113 | limit = 128 |
| 114 | } |
| 115 | limit = min(limit, 1000) |
| 116 | rows, err := tx.QueryContext(ctx, `SELECT u.message_id,u.visible_turn,u.position,u.preview, |
| 117 | COALESCE((SELECT a.preview FROM messages a WHERE a.visible_turn=u.visible_turn AND a.role='assistant' AND TRIM(a.preview)<>'' AND a.event_sequence<=? AND (a.valid_to=0 OR a.valid_to>?) ORDER BY a.position DESC LIMIT 1),'') |
| 118 | FROM messages u WHERE u.visible_user=1 AND u.visible_turn>=? AND u.event_sequence<=? AND (u.valid_to=0 OR u.valid_to>?) ORDER BY u.visible_turn,u.position LIMIT ?`, cut, cut, max(1, req.StartTurn), cut, cut, limit) |
| 119 | if err != nil { |
| 120 | return page, err |
| 121 | } |
| 122 | defer rows.Close() |
| 123 | for rows.Next() { |
| 124 | var entry HistoryOutlineEntry |
| 125 | if err := rows.Scan(&entry.MessageID, &entry.Turn, &entry.Position, &entry.Prompt, &entry.Answer); err != nil { |
| 126 | return page, err |
| 127 | } |
| 128 | entry.Prompt = textutil.ClipGraphemes(strings.Join(strings.Fields(entry.Prompt), " "), 50, "…") |
| 129 | entry.Answer = textutil.ClipGraphemes(strings.Join(strings.Fields(entry.Answer), " "), 120, "…") |
| 130 | page.Entries = append(page.Entries, entry) |
| 131 | } |
| 132 | if err := rows.Err(); err != nil { |
| 133 | return page, err |
| 134 | } |
| 135 | page.NextTurn = max(1, req.StartTurn) |
| 136 | if len(page.Entries) > 0 { |
| 137 | page.NextTurn = page.Entries[len(page.Entries)-1].Turn + 1 |
| 138 | } |
| 139 | page.Done = page.NextTurn > page.TotalTurns |
| 140 | page.Status = "ready" |
| 141 | return page, nil |
| 142 | } |
| 143 |