返回 DeepSeek-Reasonix
query_read_binding.go
根目录 / internal / session / query_read_binding.go
1 package session
2
3 import (
4 "context"
5 "sync"
6
7 "reasonix/internal/historywork"
8 )
9
10 type queryReadScope struct {
11 ctx context.Context
12 cancel context.CancelFunc
13 refs int
14 }
15
16 type queryReadScopeKey struct{}
17
18 // ConfigureHistoryMaintenance is called before publishing a Desktop service.
19 // All roots share the same preparation budget; passive lists never replay logs.
20 func (s *Service) ConfigureHistoryMaintenance(coordinator *historywork.Coordinator) {
21 s.query.maintenance = coordinator
22 s.query.metadataOnlyListings = true
23 }
24
25 // AcquireHistoryReader owns only query preparation, never runtime recovery or
26 // writer authority. Releasing the last reader cancels its rebuilds. A new
27 // reader receives a new context even if the old build is still unwinding.
28 func (q *Query) AcquireHistoryReader(ref SessionRef) (context.Context, func(), error) {
29 if err := ref.validate(q.hostID); err != nil {
30 return nil, nil, err
31 }
32 q.readMu.Lock()
33 if q.readScopes == nil {
34 q.readScopes = make(map[string]*queryReadScope)
35 }
36 scope := q.readScopes[ref.SessionID]
37 if scope == nil {
38 ctx, cancel := context.WithCancel(q.rebuildCtx)
39 scope = &queryReadScope{ctx: ctx, cancel: cancel}
40 q.readScopes[ref.SessionID] = scope
41 }
42 scope.refs++
43 q.readMu.Unlock()
44 var once sync.Once
45 return context.WithValue(scope.ctx, queryReadScopeKey{}, scope), func() {
46 once.Do(func() {
47 q.readMu.Lock()
48 scope.refs--
49 if scope.refs == 0 {
50 scope.cancel()
51 if q.readScopes[ref.SessionID] == scope {
52 delete(q.readScopes, ref.SessionID)
53 }
54 }
55 q.readMu.Unlock()
56 })
57 }, nil
58 }
59
60 func (q *Query) historyReadContext(_ string, callers ...context.Context) context.Context {
61 if len(callers) > 0 {
62 if scope, ok := callers[0].Value(queryReadScopeKey{}).(*queryReadScope); ok {
63 return scope.ctx
64 }
65 }
66 // A released binding must not revive a job under the service lifetime
67 // between its cancellation check and acquisition of readMu.
68 if len(callers) > 0 && callers[0].Err() != nil {
69 return callers[0]
70 }
71 // Legacy "preparing" responses end the request before its worker finishes.
72 // They cannot borrow another client's binding; new RPCs use the explicit
73 // cancellable read bindings above instead of this service lifetime.
74 return q.rebuildCtx
75 }
76
77 func (q *Query) acquireHistoryPreparation(ctx context.Context) (func(), error) {
78 if q.maintenance == nil {
79 return func() {}, ctx.Err()
80 }
81 return q.maintenance.Foreground(ctx)
82 }
83
83 lines GO