返回 DeepSeek-Reasonix
query.go
根目录 / internal / session / query.go
1 package session
2
3 import (
4 "context"
5 "fmt"
6 "path/filepath"
7 "strings"
8 "sync"
9 "time"
10
11 "reasonix/internal/historywork"
12 "reasonix/internal/provider"
13 )
14
15 // Query is the read model shared by desktop, CLI, Serve, ACP and bots. It uses
16 // an attached runtime when one exists and otherwise opens only a read handle;
17 // querying cold history never constructs an Agent or acquires writer ownership.
18 type Query struct {
19 maintenance *historywork.Coordinator
20 readMu sync.Mutex
21 readScopes map[string]*queryReadScope
22 metadataOnlyListings bool
23 hostID string
24 persistence SessionPersistence
25 service *Service
26 rebuildMu sync.Mutex
27 rebuilding map[string]struct{}
28 metadataQueue []metadataRebuildTask
29 metadataWorkers int
30 metadataFailures map[string]error
31 generation map[string]uint64
32 rebuildCtx context.Context
33 rebuildStop context.CancelFunc
34 rebuildWG sync.WaitGroup
35 closed bool
36 slots *rebuildSlots
37 indexMu sync.Mutex
38 indexLocks map[string]*sync.Mutex
39 contentMu sync.Mutex
40 contentGrants map[string]time.Time
41 searchMu sync.Mutex
42 searchBuilds map[string]*searchPreparation
43 historyMu sync.Mutex
44 historyBuilds map[string]*historyPreparation
45 }
46
47 func (s *Service) Query() *Query {
48 if s == nil {
49 return nil
50 }
51 return s.query
52 }
53
54 func newQuery(hostID string, persistence SessionPersistence, service *Service) *Query {
55 rebuildCtx, rebuildStop := context.WithCancel(context.Background())
56 query := &Query{
57 hostID: hostID, persistence: persistence, service: service,
58 rebuilding: map[string]struct{}{}, generation: map[string]uint64{}, rebuildCtx: rebuildCtx,
59 metadataFailures: map[string]error{},
60 rebuildStop: rebuildStop, slots: newRebuildSlots(2),
61 indexLocks: map[string]*sync.Mutex{},
62 contentGrants: map[string]time.Time{},
63 searchBuilds: map[string]*searchPreparation{},
64 historyBuilds: map[string]*historyPreparation{},
65 }
66 return query
67 }
68
69 func contentGrantKey(sessionID, storageGeneration, digest string, bytes int64, indexDigest string) string {
70 return sessionID + "\x00" + storageGeneration + "\x00" + digest + "\x00" + fmt.Sprint(bytes) + "\x00" + indexDigest
71 }
72
73 func (q *Query) storageGeneration(sessionID string) string {
74 filesystem, ok := q.persistence.(*FilesystemPersistence)
75 if !ok {
76 return ""
77 }
78 dir := filepath.Join(filesystem.Root, sessionID)
79 manifest, err := readStoredManifest(filepath.Join(dir, "manifest.json"))
80 if err != nil {
81 return ""
82 }
83 identity, err := readStorageIdentity(dir, manifest)
84 if err != nil {
85 return ""
86 }
87 return identity.Generation
88 }
89
90 func (q *Query) authorizeContentForGeneration(sessionID, generation, digest string, bytes int64, indexDigest string) {
91 if generation == "" {
92 return
93 }
94 q.contentMu.Lock()
95 defer q.contentMu.Unlock()
96 now := time.Now()
97 for key, expiry := range q.contentGrants {
98 if !expiry.After(now) {
99 delete(q.contentGrants, key)
100 }
101 }
102 q.contentGrants[contentGrantKey(sessionID, generation, digest, bytes, indexDigest)] = now.Add(15 * time.Minute)
103 }
104
105 func (q *Query) contentAuthorized(sessionID, digest string, bytes int64, indexDigest string) bool {
106 generation := q.storageGeneration(sessionID)
107 if generation == "" {
108 return false
109 }
110 q.contentMu.Lock()
111 defer q.contentMu.Unlock()
112 key := contentGrantKey(sessionID, generation, digest, bytes, indexDigest)
113 expiry, ok := q.contentGrants[key]
114 if !ok || !expiry.After(time.Now()) {
115 delete(q.contentGrants, key)
116 return false
117 }
118 return true
119 }
120
121 func (q *Query) projectionLock(kind, sessionID string) *sync.Mutex {
122 key := kind + "\x00" + sessionID
123 q.indexMu.Lock()
124 defer q.indexMu.Unlock()
125 lock := q.indexLocks[key]
126 if lock == nil {
127 lock = &sync.Mutex{}
128 q.indexLocks[key] = lock
129 }
130 return lock
131 }
132
133 // Close stops catalog work owned by this query. Individual List callers do not
134 // own shared rebuilds, so cancelling one request never cancels work another
135 // caller may use; the host query lifetime is the cancellation boundary.
136 func (q *Query) Close() {
137 if q == nil {
138 return
139 }
140 q.rebuildMu.Lock()
141 if !q.closed {
142 q.closed = true
143 if q.rebuildStop != nil {
144 q.rebuildStop()
145 }
146 }
147 q.rebuildMu.Unlock()
148 q.rebuildWG.Wait()
149 }
150
151 func (q *Query) Snapshot(ctx context.Context, ref SessionRef) (Snapshot, error) {
152 if q == nil || q.persistence == nil {
153 return Snapshot{}, fmt.Errorf("session: nil session query")
154 }
155 if err := ref.validate(q.hostID); err != nil {
156 return Snapshot{}, err
157 }
158 if q.service != nil {
159 if runtime, ok := q.service.Runtime(ref); ok {
160 return runtime.Session().Snapshot(), nil
161 }
162 }
163 handle, err := q.persistence.Open(ref.SessionID, ReadOnly)
164 if err != nil {
165 return Snapshot{}, err
166 }
167 defer handle.Close(context.WithoutCancel(ctx))
168 projection := Projection{}
169 var cursor uint64
170 for {
171 page, readErr := handle.Read(ctx, cursor, 1000)
172 if readErr != nil {
173 return Snapshot{}, readErr
174 }
175 for _, commit := range page.Commits {
176 if err := applyProjectionCommit(&projection, commit); err != nil {
177 return Snapshot{}, err
178 }
179 }
180 if !page.Truncated {
181 break
182 }
183 if page.Next <= cursor {
184 return Snapshot{}, fmt.Errorf("%w: cold history cursor did not advance", ErrDamagedStore)
185 }
186 cursor = page.Next
187 }
188 sequence := projection.CommittedSequence
189 return Snapshot{EventSequence: sequence, DurableSequence: sequence, PersistenceStatus: PersistenceReady, Projection: projection}, nil
190 }
191
192 func (q *Query) History(ctx context.Context, ref SessionRef) ([]provider.Message, error) {
193 snapshot, err := q.Snapshot(ctx, ref)
194 if err != nil {
195 return nil, err
196 }
197 return append([]provider.Message(nil), snapshot.Projection.Messages...), nil
198 }
199
200 // Stat returns one header-backed metadata observation without opening event
201 // bodies. Live projection state overlays the disposable cache, matching List.
202 func (q *Query) Stat(ctx context.Context, ref SessionRef) (SessionInfo, error) {
203 if q == nil || q.persistence == nil {
204 return SessionInfo{}, fmt.Errorf("session: nil session query")
205 }
206 if err := ref.validate(q.hostID); err != nil {
207 return SessionInfo{}, err
208 }
209 info, err := q.persistence.Stat(ctx, ref.SessionID)
210 if err != nil {
211 return SessionInfo{}, err
212 }
213 q.enrichInfo(&info)
214 return info, nil
215 }
216
217 // ResolveSessionID turns a caller-selected opaque ID into a server-observed
218 // SessionRef. The requested value is used only for equality; every identity
219 // returned to filesystem-backed readers comes from the persistence catalog.
220 func (q *Query) ResolveSessionID(ctx context.Context, requested string) (SessionRef, error) {
221 if q == nil || q.persistence == nil {
222 return SessionRef{}, fmt.Errorf("session: nil session query")
223 }
224 candidate := strings.TrimSpace(requested)
225 if err := validateSessionID(candidate); err != nil {
226 return SessionRef{}, err
227 }
228 cursor := ""
229 for {
230 page, err := q.persistence.List(ctx, cursor, 100)
231 if err != nil {
232 return SessionRef{}, err
233 }
234 for _, info := range page.Sessions {
235 if info.SessionID == candidate {
236 return SessionRef{HostID: q.hostID, SessionID: info.SessionID}, nil
237 }
238 }
239 if page.NextCursor == "" {
240 return SessionRef{}, fmt.Errorf("%w: %s", ErrSessionNotFound, candidate)
241 }
242 if page.NextCursor <= cursor {
243 return SessionRef{}, fmt.Errorf("%w: catalog cursor did not advance", ErrDamagedStore)
244 }
245 cursor = page.NextCursor
246 }
247 }
248
249 func (q *Query) List(ctx context.Context, cursor string, limit int) (SessionPage, error) {
250 if q == nil || q.persistence == nil {
251 return SessionPage{}, fmt.Errorf("session: nil session query")
252 }
253 page, err := q.persistence.List(ctx, cursor, limit)
254 if err != nil {
255 return SessionPage{}, err
256 }
257 for i := range page.Sessions {
258 q.enrichInfo(&page.Sessions[i])
259 }
260 return page, nil
261 }
262
263 func (q *Query) enrichInfo(info *SessionInfo) {
264 info.Ref = SessionRef{HostID: q.hostID, SessionID: info.SessionID}
265 if info.Error != "" {
266 return
267 }
268 if q.service != nil {
269 if runtime, ok := q.service.Runtime(info.Ref); ok {
270 applyCatalogMetadata(info, runtime.Session().CatalogMetadata())
271 return
272 }
273 }
274 if !q.metadataOnlyListings && info.Codec == Codec && info.MetadataStatus != MetadataReady {
275 q.rebuildMu.Lock()
276 failure := q.metadataFailures[info.SessionID]
277 q.rebuildMu.Unlock()
278 if failure != nil {
279 info.MetadataStatus, info.Error = MetadataFailed, failure.Error()
280 return
281 }
282 q.scheduleMetadataRebuild(info.SessionID)
283 }
284 }
285
286 func applyCatalogMetadata(info *SessionInfo, metadata catalogMetadata) {
287 info.Title, info.TitleSequence = metadata.Title, metadata.TitleSequence
288 info.ModelRef, info.ModelIdentity = metadata.ModelRef, metadata.ModelIdentity
289 info.Turns, info.Preview, info.MetadataStatus = metadata.Turns, metadata.Preview, MetadataReady
290 info.EventSequence, info.ResultSequence = metadata.Sequence, metadata.ResultSequence
291 }
292
293 func (q *Query) scheduleMetadataRebuild(sessionID string) {
294 if _, ok := q.persistence.(*FilesystemPersistence); !ok {
295 return
296 }
297 q.rebuildMu.Lock()
298 if q.closed {
299 q.rebuildMu.Unlock()
300 return
301 }
302 if _, exists := q.rebuilding[sessionID]; exists {
303 q.rebuildMu.Unlock()
304 return
305 }
306 q.rebuilding[sessionID] = struct{}{}
307 delete(q.metadataFailures, sessionID)
308 generation := q.generation[sessionID]
309 q.metadataQueue = append(q.metadataQueue, metadataRebuildTask{sessionID, generation})
310 if q.metadataWorkers < 2 {
311 q.metadataWorkers++
312 q.rebuildWG.Add(1)
313 go q.runMetadataRebuilds()
314 }
315 q.rebuildMu.Unlock()
316 }
317
318 type metadataRebuildTask struct {
319 sessionID string
320 generation uint64
321 }
322
323 // Keep a deduplicated ID queue, not one goroutine per session. Workers wait
324 // behind user history/recovery work and continue without another List call.
325 func (q *Query) runMetadataRebuilds() {
326 defer q.rebuildWG.Done()
327 for {
328 q.rebuildMu.Lock()
329 if q.closed || len(q.metadataQueue) == 0 {
330 q.metadataWorkers--
331 if q.closed {
332 q.metadataQueue = nil
333 clear(q.rebuilding)
334 }
335 q.rebuildMu.Unlock()
336 return
337 }
338 task := q.metadataQueue[0]
339 q.metadataQueue[0] = metadataRebuildTask{}
340 q.metadataQueue = q.metadataQueue[1:]
341 q.rebuildMu.Unlock()
342 if err := q.slots.acquire(q.rebuildCtx, rebuildPriorityPrefetch); err != nil {
343 continue
344 }
345 err := q.rebuildCatalogMetadata(task.sessionID, task.generation)
346 q.slots.release()
347 q.rebuildMu.Lock()
348 if err != nil && q.generation[task.sessionID] == task.generation {
349 q.metadataFailures[task.sessionID] = err
350 }
351 delete(q.rebuilding, task.sessionID)
352 q.rebuildMu.Unlock()
353 }
354 }
355
356 // invalidateCatalog fences every metadata task scheduled for an older session
357 // incarnation. Service calls it before deleting the directory, so a delayed
358 // task cannot publish a cache entry that makes the deleted session reappear.
359 func (q *Query) invalidateCatalog(sessionID string) {
360 if q == nil {
361 return
362 }
363 q.rebuildMu.Lock()
364 q.generation[sessionID]++
365 delete(q.metadataFailures, sessionID)
366 q.rebuildMu.Unlock()
367 }
368
369 func (q *Query) rebuildCatalogMetadata(sessionID string, generation uint64, callers ...context.Context) error {
370 ctx := q.rebuildCtx
371 if len(callers) > 0 {
372 ctx = callers[0]
373 }
374 filesystem, ok := q.persistence.(*FilesystemPersistence)
375 if !ok {
376 return nil
377 }
378 handle, err := q.persistence.Open(sessionID, ReadOnly)
379 if err != nil {
380 return err
381 }
382 defer handle.Close(context.WithoutCancel(ctx))
383 sessionDir := filepath.Join(filesystem.Root, sessionID)
384 cacheDir := filepath.Join(filesystem.Root, ".query-cache", filepath.Base(sessionID))
385 manifest, err := readManifest(filepath.Join(sessionDir, "manifest.json"))
386 if err != nil {
387 return err
388 }
389 metadata, err := reduceCatalogMetadata(ctx, handle, manifest)
390 if err != nil {
391 return err
392 }
393 q.rebuildMu.Lock()
394 defer q.rebuildMu.Unlock()
395 if ctx.Err() != nil || q.rebuildCtx.Err() != nil || q.generation[sessionID] != generation {
396 return nil
397 }
398 // Hold the generation boundary through publication. Deletion invalidates
399 // before it moves the directory, so it either wins first or waits until this
400 // exact-incarnation cache is completely written and then removes it.
401 return writeCatalogMetadataForSession(cacheDir, sessionDir, metadata)
402 }
403
404 // RefreshMetadata is an explicit management-operation fence. Passive list
405 // reads never call it; imports and archive candidates need an exact committed
406 // sequence before they can freeze their optimistic mutation preconditions.
407 func (q *Query) RefreshMetadata(ctx context.Context, ref SessionRef) (SessionInfo, error) {
408 if err := ref.validate(q.hostID); err != nil {
409 return SessionInfo{}, err
410 }
411 q.rebuildMu.Lock()
412 generation := q.generation[ref.SessionID]
413 q.rebuildMu.Unlock()
414 if err := q.rebuildCatalogMetadata(ref.SessionID, generation, ctx); err != nil {
415 return SessionInfo{}, err
416 }
417 return q.Stat(ctx, ref)
418 }
419
419 lines GO