返回 DeepSeek-Reasonix
session_history_binding.go
根目录 / desktop / session_history_binding.go
1 package main
2
3 import (
4 "context"
5 "errors"
6 "path/filepath"
7 "sync"
8
9 "reasonix/internal/agent"
10 "reasonix/internal/session"
11 )
12
13 type SessionHistoryReadHandle struct {
14 ID string `json:"id"`
15 StorageBackend string `json:"storageBackend"`
16 SessionGeneration uint64 `json:"sessionGeneration"`
17 Capabilities []string `json:"capabilities"`
18 }
19
20 type desktopHistoryReader struct {
21 handle SessionHistoryReadHandle
22 ctx context.Context
23 release func()
24 query *session.Query
25 ref session.SessionRef
26 tab *WorkspaceTab
27 tabID, path, head string
28 identity string
29 native *nativeHistoryPreparation
30 }
31
32 type desktopHistoryReaders struct {
33 mu sync.Mutex
34 entries map[string]*desktopHistoryReader
35 native map[string]*nativeHistoryPreparation
36 workers sync.WaitGroup
37 closed bool
38 }
39
40 // BeginSessionHistoryReadForTab captures the source once. Subsequent reads
41 // never resolve whichever session happens to occupy the tab at completion.
42 func (a *App) BeginSessionHistoryReadForTab(tabID string) (SessionHistoryReadHandle, error) {
43 a.mu.RLock()
44 tab := a.tabByIDLocked(tabID)
45 if tab == nil {
46 a.mu.RUnlock()
47 return SessionHistoryReadHandle{Capabilities: []string{}}, errors.New("history tab unavailable")
48 }
49 path, head, id, generation := tab.currentSessionPath(), tab.SessionHeadID, tab.SessionID, tab.SessionGeneration
50 a.mu.RUnlock()
51 reader := &desktopHistoryReader{tab: tab, tabID: tabID, path: path, head: head, identity: id,
52 handle: SessionHistoryReadHandle{ID: newSessionRuntimeID("history"), StorageBackend: "legacy", SessionGeneration: generation, Capabilities: []string{"history-read-binding-v1"}}}
53 ctx, cancel := context.WithCancel(a.bootContext())
54 reader.ctx, reader.release = ctx, cancel
55 if id != "" || nativeStoredDirectory(path) {
56 var query *session.Query
57 var ref session.SessionRef
58 var err error
59 if id != "" {
60 query, ref, err = a.canonicalSessionQuery(tabID)
61 } else {
62 // Prototype/v3/v4 directories retain their original store. A cold
63 // query needs neither import nor a leased execution runtime.
64 var service *session.Service
65 service, err = a.historicalSessionService(filepath.Dir(path))
66 if err == nil {
67 query = service.Query()
68 ref = session.SessionRef{HostID: localDesktopHostID, SessionID: filepath.Base(path)}
69 }
70 }
71 if err != nil {
72 cancel()
73 return SessionHistoryReadHandle{Capabilities: []string{}}, err
74 }
75 shared, release, err := query.AcquireHistoryReader(ref)
76 if err != nil {
77 cancel()
78 return SessionHistoryReadHandle{Capabilities: []string{}}, err
79 }
80 cancel()
81 reader.ctx, cancel = context.WithCancel(shared)
82 reader.release = func() { cancel(); release() }
83 reader.query, reader.ref, reader.handle.StorageBackend = query, ref, "canonical"
84 }
85 var nativeGeneration, nativeSourceKey string
86 if reader.query == nil && path != "" {
87 var err error
88 nativeGeneration, err = nativeHistorySourceGeneration(path)
89 if err != nil {
90 reader.release()
91 return SessionHistoryReadHandle{Capabilities: []string{}}, err
92 }
93 nativeSourceKey = sessionRuntimeKey(path)
94 }
95 if !a.historyReaderCurrent(reader) {
96 reader.release()
97 return SessionHistoryReadHandle{Capabilities: []string{}}, errors.New("history source changed")
98 }
99 a.historyReaders.mu.Lock()
100 if a.historyReaders.entries == nil {
101 a.historyReaders.entries = make(map[string]*desktopHistoryReader)
102 }
103 if a.shuttingDown.Load() || a.historyReaders.closed {
104 a.historyReaders.mu.Unlock()
105 reader.release()
106 return SessionHistoryReadHandle{Capabilities: []string{}}, context.Canceled
107 }
108 a.historyReaders.entries[reader.handle.ID] = reader
109 if reader.query == nil && reader.path != "" {
110 reader.handle.Capabilities = append(reader.handle.Capabilities, "history-native-navigation-v1", "history-native-search-v1")
111 var releaseNative context.CancelFunc
112 reader.native, releaseNative = a.acquireNativeHistoryLocked(reader.path, reader.head, nativeSourceKey, nativeGeneration)
113 release := reader.release
114 // Source retirement cancels every read in that generation. Individual
115 // navigation cancellation still releases only this reader's reference.
116 reader.ctx, cancel = context.WithCancel(reader.native.ctx)
117 reader.release = func() { cancel(); release(); releaseNative() }
118 }
119 a.historyReaders.mu.Unlock()
120 return reader.handle, nil
121 }
122
123 func (a *App) historyReaderCurrent(reader *desktopHistoryReader) bool {
124 a.mu.RLock()
125 defer a.mu.RUnlock()
126 tab := a.tabByIDLocked(reader.tabID)
127 return tab == reader.tab && tab.SessionGeneration == reader.handle.SessionGeneration && tab.SessionID == reader.identity && tab.currentSessionPath() == reader.path && tab.SessionHeadID == reader.head
128 }
129
130 func (a *App) historyReader(id string) (*desktopHistoryReader, error) {
131 a.historyReaders.mu.Lock()
132 reader := a.historyReaders.entries[id]
133 a.historyReaders.mu.Unlock()
134 if reader == nil {
135 return nil, context.Canceled
136 }
137 if err := reader.ctx.Err(); err != nil {
138 return nil, err
139 }
140 if !a.historyReaderCurrent(reader) {
141 a.ReleaseSessionHistoryRead(id)
142 return nil, context.Canceled
143 }
144 return reader, nil
145 }
146
147 // Release is idempotent and affects only the exact handle, never a new reader.
148 func (a *App) ReleaseSessionHistoryRead(id string) {
149 a.historyReaders.mu.Lock()
150 reader := a.historyReaders.entries[id]
151 delete(a.historyReaders.entries, id)
152 a.historyReaders.mu.Unlock()
153 if reader != nil {
154 reader.release()
155 }
156 }
157
158 func (a *App) closeHistoryReaders() {
159 a.historyReaders.mu.Lock()
160 readers := a.historyReaders.entries
161 a.historyReaders.entries = nil
162 a.historyReaders.closed = true
163 // Compatibility reads also borrow preparations without a persistent RPC
164 // handle. Shutdown retires the owner, then joins every cache close barrier.
165 for _, job := range a.historyReaders.native {
166 job.cancel()
167 }
168 a.historyReaders.mu.Unlock()
169 for _, reader := range readers {
170 reader.release()
171 }
172 a.historyReaders.workers.Wait()
173 }
174
175 func (a *App) ReadSessionHistoryWindow(id string, req session.HistoryWindowRequest) (session.HistoryWindowPage, error) {
176 r, err := a.historyReader(id)
177 if err != nil {
178 return session.HistoryWindowPage{Messages: []session.PersistentMessage{}, Status: "stale_cursor"}, nil
179 }
180 if r.query == nil {
181 return session.HistoryWindowPage{Messages: []session.PersistentMessage{}, Status: "unsupported"}, nil
182 }
183 page, err := r.query.ReadHistoryWindow(r.ctx, r.ref, req)
184 if r.ctx.Err() != nil || !a.historyReaderCurrent(r) {
185 return session.HistoryWindowPage{Messages: []session.PersistentMessage{}, Status: "stale_cursor"}, nil
186 }
187 return page, err
188 }
189
190 type SessionHistoryReadSlice struct {
191 Status string `json:"status"`
192 Page HistorySlice `json:"page"`
193 }
194
195 func (a *App) ReadSessionHistorySlice(id string, req HistorySliceRequest) (SessionHistoryReadSlice, error) {
196 r, err := a.historyReader(id)
197 if err != nil {
198 return SessionHistoryReadSlice{Status: "stale_cursor", Page: emptyHistorySlice()}, nil
199 }
200 if r.query != nil || r.path == "" || nativeStoredDirectory(r.path) {
201 return SessionHistoryReadSlice{Status: "unsupported", Page: emptyHistorySlice()}, nil
202 }
203 if r.native == nil {
204 return SessionHistoryReadSlice{Status: "unsupported", Page: emptyHistorySlice()}, nil
205 }
206 pager, err := r.native.wait(r.ctx)
207 if err != nil {
208 if errors.Is(err, agent.ErrDisplayFormatUnsupported) {
209 return SessionHistoryReadSlice{Status: "unsupported", Page: emptyHistorySlice()}, nil
210 }
211 if r.ctx.Err() != nil || errors.Is(err, agent.ErrDisplaySourceChanged) {
212 return SessionHistoryReadSlice{Status: "stale_cursor", Page: emptyHistorySlice()}, nil
213 }
214 return SessionHistoryReadSlice{Status: "failed", Page: emptyHistorySlice()}, err
215 }
216 req, status, err := nativeHistorySliceAnchor(r, pager, req)
217 if r.ctx.Err() != nil || !a.historyReaderCurrent(r) || errors.Is(err, agent.ErrDisplaySourceChanged) {
218 return SessionHistoryReadSlice{Status: "stale_cursor", Page: emptyHistorySlice()}, nil
219 }
220 if err != nil || status != "ready" {
221 return SessionHistoryReadSlice{Status: status, Page: emptyHistorySlice()}, err
222 }
223 page, ready, err := a.historySliceFromPager(r.ctx, pager, filepath.Dir(r.path), r.path, normalizeHistorySliceRequest(req), r.native.key)
224 if errors.Is(err, agent.ErrDisplaySourceChanged) || r.ctx.Err() != nil || !a.historyReaderCurrent(r) {
225 return SessionHistoryReadSlice{Status: "stale_cursor", Page: emptyHistorySlice()}, nil
226 }
227 if err != nil {
228 if errors.Is(err, context.Canceled) {
229 return SessionHistoryReadSlice{Status: "stale_cursor", Page: emptyHistorySlice()}, nil
230 }
231 return SessionHistoryReadSlice{Status: "failed", Page: emptyHistorySlice()}, err
232 }
233 if !ready {
234 // No preparation was admitted for this native format. Do not advertise
235 // an indefinitely pending job; the negotiated compatibility reader can
236 // still serve this same source until its runtime projection is bound.
237 return SessionHistoryReadSlice{Status: "unsupported", Page: emptyHistorySlice()}, nil
238 }
239 status = "ready"
240 if page.Stale {
241 status = "stale_cursor"
242 }
243 for i := range page.Entries {
244 for j := range page.Entries[i].Refs {
245 page.Entries[i].Refs[j].ReadHandleID = id
246 }
247 }
248 return SessionHistoryReadSlice{Status: status, Page: page}, nil
249 }
250
251 func (a *App) ReadSessionHistoryOutline(id string, req session.HistoryOutlineRequest) (session.HistoryOutlinePage, error) {
252 r, err := a.historyReader(id)
253 if err != nil {
254 return session.HistoryOutlinePage{Entries: []session.HistoryOutlineEntry{}, Status: "stale_cursor"}, nil
255 }
256 var page session.HistoryOutlinePage
257 if r.query == nil {
258 page, err = a.readNativeHistoryOutline(r, req)
259 } else {
260 page, err = r.query.ReadHistoryOutline(r.ctx, r.ref, req)
261 }
262 if r.ctx.Err() != nil || !a.historyReaderCurrent(r) || errors.Is(err, agent.ErrDisplaySourceChanged) {
263 return session.HistoryOutlinePage{Entries: []session.HistoryOutlineEntry{}, Status: "stale_cursor"}, nil
264 }
265 return page, err
266 }
267
268 func (a *App) LocateSessionHistoryMessage(id, messageID string, snapshot uint64) (session.MessageLocation, error) {
269 r, err := a.historyReader(id)
270 if err != nil {
271 return session.MessageLocation{Status: "stale_cursor"}, nil
272 }
273 var location session.MessageLocation
274 if r.query == nil {
275 var pager *agent.DisplayPager
276 pager, err = nativeNavigationPager(r)
277 if err == nil {
278 location, err = nativeHistoryLocation(r, pager, messageID, snapshot)
279 } else {
280 location.Status = nativeNavigationStatus(r, err)
281 if location.Status != "failed" {
282 err = nil
283 }
284 }
285 } else {
286 location, err = r.query.LocateMessage(r.ctx, r.ref, messageID, snapshot)
287 }
288 if r.ctx.Err() != nil || !a.historyReaderCurrent(r) || errors.Is(err, agent.ErrDisplaySourceChanged) {
289 return session.MessageLocation{Status: "stale_cursor"}, nil
290 }
291 return location, err
292 }
293
294 func (a *App) SearchSessionHistoryRead(id, text, cursor string, limit int) (session.SearchHistoryPage, error) {
295 r, err := a.historyReader(id)
296 if err != nil {
297 return session.SearchHistoryPage{Hits: []session.SearchHistoryHit{}, Status: "stale_cursor"}, nil
298 }
299 var page session.SearchHistoryPage
300 if r.query == nil {
301 page, err = a.searchNativeHistory(r, text, cursor, limit)
302 } else {
303 page, err = r.query.SearchHistory(r.ctx, r.ref, text, cursor, limit)
304 }
305 if r.ctx.Err() != nil || !a.historyReaderCurrent(r) || errors.Is(err, agent.ErrDisplaySourceChanged) {
306 return session.SearchHistoryPage{Hits: []session.SearchHistoryHit{}, Status: "stale_cursor"}, nil
307 }
308 return page, err
309 }
310
310 lines GO