| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "errors" |
| 6 | "path/filepath" |
| 7 | "strings" |
| 8 | |
| 9 | "reasonix/internal/agent" |
| 10 | "reasonix/internal/provider" |
| 11 | ) |
| 12 | |
| 13 | func (a *App) pagedColdHistorySlice(ctx context.Context, sessionDir, path string, req HistorySliceRequest, heads ...string) (HistorySlice, bool, error) { |
| 14 | head := "" |
| 15 | if len(heads) > 0 { |
| 16 | head = heads[0] |
| 17 | } |
| 18 | page, handled := emptyHistorySlice(), false |
| 19 | err := a.withNativeHistoryPager(ctx, path, head, func(ctx context.Context, pager *agent.DisplayPager, sourceID string) error { |
| 20 | handled = true |
| 21 | var err error |
| 22 | page, _, err = a.historySliceFromPager(ctx, pager, sessionDir, path, req, sourceID) |
| 23 | return err |
| 24 | }) |
| 25 | if handled || nativeHistoryPreparationFailure(err) { |
| 26 | if err != nil { |
| 27 | return emptyHistorySlice(), true, err |
| 28 | } |
| 29 | return page, true, nil |
| 30 | } |
| 31 | return HistorySlice{}, false, nil |
| 32 | } |
| 33 | |
| 34 | func nativeHistoryPreparationFailure(err error) bool { |
| 35 | // Only a positively identified unsupported format may use the legacy |
| 36 | // adapter. I/O, parse and cache failures must not silently trigger a full |
| 37 | // replay after the bounded reader has failed. |
| 38 | return err != nil && !errors.Is(err, agent.ErrDisplayFormatUnsupported) |
| 39 | } |
| 40 | |
| 41 | // Compatibility paging and field reads borrow the same preparation and cache |
| 42 | // owner as bound readers. Retirement cancels both the read and its decoding; |
| 43 | // only the last borrower retires the projection, after the callback returns. |
| 44 | func (a *App) withNativeHistoryPager(ctx context.Context, path, head string, read func(context.Context, *agent.DisplayPager, string) error) error { |
| 45 | if err := ctx.Err(); err != nil { |
| 46 | return err |
| 47 | } |
| 48 | generation, err := nativeHistorySourceGeneration(path) |
| 49 | if err != nil { |
| 50 | return errors.Join(agent.ErrDisplaySourceChanged, err) |
| 51 | } |
| 52 | sourceKey := sessionRuntimeKey(path) |
| 53 | a.historyReaders.mu.Lock() |
| 54 | if a.historyReaders.closed || a.shuttingDown.Load() { |
| 55 | a.historyReaders.mu.Unlock() |
| 56 | return context.Canceled |
| 57 | } |
| 58 | job, release := a.acquireNativeHistoryLocked(path, head, sourceKey, generation) |
| 59 | a.historyReaders.mu.Unlock() |
| 60 | defer release() |
| 61 | ctx, cancel := context.WithCancel(ctx) |
| 62 | stop := context.AfterFunc(job.ctx, cancel) |
| 63 | defer func() { stop(); cancel() }() |
| 64 | pager, err := job.wait(ctx) |
| 65 | if err != nil { |
| 66 | return err |
| 67 | } |
| 68 | pager = pager.WithContext(ctx) |
| 69 | if err := pager.Validate(); err != nil { |
| 70 | return err |
| 71 | } |
| 72 | err = read(ctx, pager, job.key) |
| 73 | if ctx.Err() != nil || job.ctx.Err() != nil { |
| 74 | return context.Canceled |
| 75 | } |
| 76 | if err != nil { |
| 77 | return err |
| 78 | } |
| 79 | return pager.Validate() |
| 80 | } |
| 81 | |
| 82 | func (a *App) historySliceFromPager(ctx context.Context, pager *agent.DisplayPager, sessionDir, path string, req HistorySliceRequest, sourceID string) (HistorySlice, bool, error) { |
| 83 | src := historySourceFromPager(ctx, pager, path, sourceID) |
| 84 | page, err := a.pageHistorySliceSource(src, req, sessionDisplayResolver(sessionDir, path), sessionPlannerDisplayTurns(sessionDir, path), nil, path) |
| 85 | if src.readErr != nil { |
| 86 | return emptyHistorySlice(), true, src.readErr |
| 87 | } |
| 88 | page.Source = "index" |
| 89 | if pager.Built { |
| 90 | page.Source = "scan" |
| 91 | } |
| 92 | if pager.DAG || pager.SchemaOne { |
| 93 | page.Source = "event-log" |
| 94 | } |
| 95 | return page, true, err |
| 96 | } |
| 97 | |
| 98 | func historySourceFromPager(ctx context.Context, pager *agent.DisplayPager, path, sourceID string) *historySliceSource { |
| 99 | pager = pager.WithContext(ctx) |
| 100 | idx := &pager.Header |
| 101 | src := &historySliceSource{sessionID: strings.TrimSuffix(filepath.Base(path), ".jsonl"), total: idx.MessageCount, |
| 102 | totalTurns: idx.AuthoredTurns, revision: idx.Revision, revKnown: idx.RevisionKnown, digest: idx.ContentDigest, |
| 103 | position: pager.Entry, usersBefore: pager.UsersBefore, maxFetch: 500, windowOnly: true} |
| 104 | src.sourceID = sourceID |
| 105 | if idx.RevisionKnown { |
| 106 | src.epoch = int(idx.Revision) |
| 107 | } else { |
| 108 | src.revision = 0 |
| 109 | } |
| 110 | src.fetch = func(lo, hi int) ([]provider.Message, error) { |
| 111 | if err := pager.Validate(); err != nil { |
| 112 | return nil, err |
| 113 | } |
| 114 | if pager.DAG || pager.SchemaOne { |
| 115 | return pager.EventMessages(lo, hi) |
| 116 | } |
| 117 | entries, err := pager.Entries(lo, hi) |
| 118 | if err != nil { |
| 119 | return nil, err |
| 120 | } |
| 121 | messages, err := readSessionMessagesAtOffsetsContext(ctx, path, entries) |
| 122 | if err != nil { |
| 123 | return nil, err |
| 124 | } |
| 125 | return messages, pager.Validate() |
| 126 | } |
| 127 | src.windowBytes = func(lo, hi int) int64 { |
| 128 | if hi <= lo { |
| 129 | return 0 |
| 130 | } |
| 131 | if pager.DAG || pager.SchemaOne { |
| 132 | entries, err := pager.Entries(lo, hi) |
| 133 | if err != nil { |
| 134 | src.readErr = err |
| 135 | return 0 |
| 136 | } |
| 137 | var total int64 |
| 138 | for _, entry := range entries { |
| 139 | total += entry.Length |
| 140 | } |
| 141 | return total |
| 142 | } |
| 143 | first, err := pager.Entry(lo) |
| 144 | if err != nil { |
| 145 | src.readErr = err |
| 146 | return 0 |
| 147 | } |
| 148 | last, err := pager.Entry(hi - 1) |
| 149 | if err != nil { |
| 150 | src.readErr = err |
| 151 | return 0 |
| 152 | } |
| 153 | return last.Offset + last.Length - first.Offset |
| 154 | } |
| 155 | return src |
| 156 | } |
| 157 |