| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "crypto/sha256" |
| 6 | "fmt" |
| 7 | "os" |
| 8 | "path/filepath" |
| 9 | "sync" |
| 10 | |
| 11 | "reasonix/internal/agent" |
| 12 | "reasonix/internal/config" |
| 13 | "reasonix/internal/fileops" |
| 14 | "reasonix/internal/store" |
| 15 | ) |
| 16 | |
| 17 | // One preparation owns a source/head generation. Readers own references, |
| 18 | // rather than independently rebuilding the same derived SQLite index. |
| 19 | type nativeHistoryPreparation struct { |
| 20 | key string |
| 21 | cacheKey string |
| 22 | path string |
| 23 | refs int // guarded by desktopHistoryReaders.mu |
| 24 | done chan struct{} |
| 25 | closed chan struct{} |
| 26 | cancel context.CancelFunc |
| 27 | ctx context.Context |
| 28 | pager *agent.DisplayPager |
| 29 | err error |
| 30 | searchMu sync.Mutex |
| 31 | search *nativeHistorySearch |
| 32 | } |
| 33 | |
| 34 | // Stat identities are task admission tokens, not content proof. OpenDisplayPager |
| 35 | // independently validates the selected branch and authoritative content. |
| 36 | func nativeHistorySourceGeneration(path string) (string, error) { |
| 37 | hash := sha256.New() |
| 38 | for i, candidate := range []string{path, store.SessionEventLog(path), store.SessionDisplayIndex(path)} { |
| 39 | info, err := os.Stat(candidate) |
| 40 | if i > 0 && os.IsNotExist(err) { |
| 41 | fmt.Fprintf(hash, "%d:absent\x00", i) |
| 42 | continue |
| 43 | } |
| 44 | if err != nil { |
| 45 | return "", err |
| 46 | } |
| 47 | target, version := fileops.DiskSnapshot(candidate, info) |
| 48 | fmt.Fprintf(hash, "%d:%s:%s\x00", i, target.Key, version) |
| 49 | } |
| 50 | return fmt.Sprintf("%x", hash.Sum(nil)), nil |
| 51 | } |
| 52 | |
| 53 | func (p *nativeHistoryPreparation) wait(ctx context.Context) (*agent.DisplayPager, error) { |
| 54 | select { |
| 55 | case <-ctx.Done(): |
| 56 | return nil, ctx.Err() |
| 57 | case <-p.ctx.Done(): |
| 58 | return nil, p.ctx.Err() |
| 59 | case <-p.done: |
| 60 | if err := ctx.Err(); err != nil { |
| 61 | return nil, err |
| 62 | } |
| 63 | // Completion and retirement may both be ready. Retirement wins even |
| 64 | // when select chose completion; the pager no longer admits readers. |
| 65 | if err := p.ctx.Err(); err != nil { |
| 66 | return nil, err |
| 67 | } |
| 68 | return p.pager, p.err |
| 69 | } |
| 70 | } |
| 71 | |
| 72 | // Caller holds historyReaders.mu. All disk work runs outside both this mutex |
| 73 | // and App.mu. Its lifetime never owns model execution or persistence. |
| 74 | func (a *App) acquireNativeHistoryLocked(path, head, sourceKey, generation string) (*nativeHistoryPreparation, context.CancelFunc) { |
| 75 | manager := &a.historyReaders |
| 76 | if manager.native == nil { |
| 77 | manager.native = map[string]*nativeHistoryPreparation{} |
| 78 | } |
| 79 | sum := sha256.Sum256([]byte(sourceKey + "\x00" + head)) |
| 80 | cacheKey := fmt.Sprintf("%x", sum) |
| 81 | key := cacheKey + ":" + generation |
| 82 | job := manager.native[key] |
| 83 | if job == nil || job.ctx.Err() != nil { |
| 84 | // A replaced source invalidates its old readers. Retire those derived |
| 85 | // handles before replacing the shared cache, including on Windows. |
| 86 | var retired []<-chan struct{} |
| 87 | for _, previous := range manager.native { |
| 88 | if previous.cacheKey == cacheKey { |
| 89 | previous.cancel() |
| 90 | retired = append(retired, previous.closed) |
| 91 | } |
| 92 | } |
| 93 | ctx, cancel := context.WithCancel(a.bootContext()) |
| 94 | ctx = a.historyMaintenance.Context(ctx) |
| 95 | job = &nativeHistoryPreparation{key: key, cacheKey: cacheKey, path: path, done: make(chan struct{}), closed: make(chan struct{}), cancel: cancel, ctx: ctx} |
| 96 | manager.native[key] = job |
| 97 | manager.workers.Go(func() { |
| 98 | defer func() { |
| 99 | job.closeSearch() |
| 100 | if job.pager != nil { |
| 101 | _ = job.pager.Close() |
| 102 | } |
| 103 | manager.mu.Lock() |
| 104 | if manager.native[key] == job { |
| 105 | delete(manager.native, key) |
| 106 | } |
| 107 | close(job.closed) |
| 108 | manager.mu.Unlock() |
| 109 | }() |
| 110 | for _, closed := range retired { |
| 111 | // Preserve the cache retirement chain even if this successor is |
| 112 | // canceled before admission. A third reader must not rebuild the |
| 113 | // same cache while a predecessor still holds its database open. |
| 114 | <-closed |
| 115 | } |
| 116 | release, err := a.historyMaintenance.Foreground(ctx) |
| 117 | if err == nil { |
| 118 | root := config.CacheDir() |
| 119 | current, verifyErr := nativeHistorySourceGeneration(path) |
| 120 | if verifyErr != nil { |
| 121 | err = verifyErr |
| 122 | } else if current != generation { |
| 123 | err = agent.ErrDisplaySourceChanged |
| 124 | } else if root == "" { |
| 125 | err = fmt.Errorf("history cache unavailable") |
| 126 | } else { |
| 127 | job.pager, err = agent.OpenDisplayPager(ctx, path, filepath.Join(root, "history-display-v1", cacheKey+".sqlite"), head) |
| 128 | if err == nil { |
| 129 | current, verifyErr = nativeHistorySourceGeneration(path) |
| 130 | if verifyErr != nil { |
| 131 | err = verifyErr |
| 132 | } else if current != generation { |
| 133 | err = agent.ErrDisplaySourceChanged |
| 134 | } |
| 135 | } |
| 136 | } |
| 137 | release() |
| 138 | } |
| 139 | job.err = err |
| 140 | close(job.done) |
| 141 | <-ctx.Done() |
| 142 | }) |
| 143 | } |
| 144 | job.refs++ |
| 145 | var once sync.Once |
| 146 | return job, func() { |
| 147 | once.Do(func() { |
| 148 | manager.mu.Lock() |
| 149 | job.refs-- |
| 150 | if job.refs == 0 { |
| 151 | // Keep the retiring owner discoverable until its database closes. |
| 152 | // An immediate reopen must join the close barrier, not reuse it. |
| 153 | job.cancel() |
| 154 | } |
| 155 | manager.mu.Unlock() |
| 156 | }) |
| 157 | } |
| 158 | } |
| 159 |