返回 DeepSeek-Reasonix
session_history_preparation.go
根目录 / desktop / session_history_preparation.go
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
159 lines GO