返回 DeepSeek-Reasonix
service_ownership.go
根目录 / internal / session / service_ownership.go
1 package session
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "sync"
8 "time"
9 )
10
11 // HostID returns the immutable host namespace used to validate SessionRef.
12 func (s *Service) HostID() string {
13 if s == nil {
14 return ""
15 }
16 return s.hostID
17 }
18
19 type prepareRuntime struct{ done chan struct{} }
20
21 // ClientBinding grants a caller access to one exact published runtime without
22 // granting authority to close the host-owned writer. Release only detaches this
23 // client; the host retires an idle runtime after its last binding disappears.
24 type ClientBinding struct {
25 service *Service
26 runtime *Runtime
27 once sync.Once
28 err error
29 }
30
31 func (b *ClientBinding) Runtime() *Runtime {
32 if b == nil {
33 return nil
34 }
35 return b.runtime
36 }
37
38 func (b *ClientBinding) Release(ctx context.Context) error {
39 if b == nil || b.service == nil || b.runtime == nil {
40 return nil
41 }
42 b.once.Do(func() { b.err = b.service.releaseBinding(ctx, b.runtime) })
43 return b.err
44 }
45
46 // PreparedRuntime owns a write handle that has been fully opened but is not
47 // yet visible through the host registry. Callers may build projections, seed
48 // initial events, and flush before atomically publishing the exact instance.
49 // A candidate must be either published or discarded.
50 type PreparedRuntime struct {
51 service *Service
52 runtime *Runtime
53 instance string
54
55 mu sync.Mutex
56 published bool
57 discarded bool
58 }
59
60 func (p *PreparedRuntime) Runtime() *Runtime {
61 if p == nil {
62 return nil
63 }
64 return p.runtime
65 }
66
67 func NewService(hostID string, persistence SessionPersistence) (*Service, error) {
68 if hostID == "" || persistence == nil {
69 return nil, errors.New("session: host id and persistence are required")
70 }
71 service := &Service{
72 hostID: hostID, persistence: persistence,
73 active: map[SessionRef]*Runtime{}, closed: map[SessionRef]error{}, preparing: map[SessionRef]*prepareRuntime{},
74 bindings: map[*Runtime]int{}, retiring: map[*Runtime]chan struct{}{}, retireIdle: map[*Runtime]bool{},
75 idleTimers: map[*Runtime]*time.Timer{}, idleWeight: map[*Runtime]int64{}, idleOrder: map[*Runtime]uint64{},
76 idleBudget: 256 << 20, idleTTL: 60 * time.Second,
77 }
78 service.query = newQuery(hostID, persistence, service)
79 return service, nil
80 }
81
82 func (s *Service) Create(ctx context.Context, options CreateOptions) (*Runtime, error) {
83 prepared, err := s.PrepareCreate(ctx, options)
84 if err != nil {
85 return nil, err
86 }
87 owner, err := s.Publish(prepared)
88 if err != nil {
89 _ = s.Discard(context.Background(), prepared)
90 return nil, err
91 }
92 return owner.Runtime(), nil
93 }
94
95 // PrepareCreate reserves the immutable session identity and its writer lease
96 // without publishing an attachable runtime. This is the DSH prepare phase:
97 // host/controller state remains untouched until Publish succeeds.
98 func (s *Service) PrepareCreate(ctx context.Context, options CreateOptions) (*PreparedRuntime, error) {
99 if err := ctx.Err(); err != nil {
100 return nil, err
101 }
102 session, err := s.persistence.Create(options)
103 if err != nil {
104 return nil, err
105 }
106 session.externalizeDurableHistory()
107 ref := SessionRef{HostID: s.hostID, SessionID: session.ID()}
108 candidate, err := newRuntime(ref, session)
109 if err != nil {
110 return nil, errors.Join(err, session.Close(context.Background()))
111 }
112 candidate.owner = s
113 return &PreparedRuntime{service: s, runtime: candidate, instance: randomID()}, nil
114 }
115
116 // Publish makes the exact prepared runtime visible and returns the host
117 // authority over that instance. It never replaces an existing instance with the
118 // same identity; the caller must resolve that ownership conflict explicitly.
119 //
120 // Only a RuntimeOwner can terminate a published runtime. Clients attach through
121 // ClientBinding and can only detach themselves.
122 func (s *Service) Publish(prepared *PreparedRuntime) (*RuntimeOwner, error) {
123 if prepared == nil || prepared.service != s || prepared.runtime == nil {
124 return nil, errors.New("session: invalid prepared runtime")
125 }
126 prepared.mu.Lock()
127 defer prepared.mu.Unlock()
128 if prepared.discarded {
129 return nil, errors.New("session: prepared runtime was discarded")
130 }
131 candidate := prepared.runtime
132 if prepared.published {
133 return &RuntimeOwner{service: s, runtime: candidate, instance: prepared.instance}, nil
134 }
135 ref := candidate.ref
136 s.mu.Lock()
137 if current := s.active[ref]; current != nil {
138 s.mu.Unlock()
139 return nil, fmt.Errorf("%w: %s", ErrSessionExists, ref.SessionID)
140 }
141 delete(s.closed, ref)
142 candidate.instance = prepared.instance
143 s.active[ref] = candidate
144 s.revision.Add(1)
145 s.mu.Unlock()
146 prepared.published = true
147 return &RuntimeOwner{service: s, runtime: candidate, instance: prepared.instance}, nil
148 }
149
150 // RuntimeOwner is the host-side authority over one exact published instance.
151 // Controller and client code must never hold one: they use ClientBinding so a
152 // failed attach can only undo its own bind, never dispose a shared runtime.
153 type RuntimeOwner struct {
154 service *Service
155 runtime *Runtime
156 instance string
157 }
158
159 // Runtime exposes the owned instance for host preparation work.
160 func (o *RuntimeOwner) Runtime() *Runtime {
161 if o == nil {
162 return nil
163 }
164 return o.runtime
165 }
166
167 // Bind attaches a client to this instance without transferring ownership.
168 func (o *RuntimeOwner) Bind() (*ClientBinding, error) {
169 if o == nil || o.service == nil || o.runtime == nil {
170 return nil, ErrSessionNotRunning
171 }
172 return o.service.Bind(o.runtime)
173 }
174
175 // Close terminates this exact instance. It refuses while any client is still
176 // bound, and it can never affect a same-ID successor published later.
177 func (o *RuntimeOwner) Close(ctx context.Context) error {
178 if o == nil || o.service == nil || o.runtime == nil {
179 return ErrSessionNotRunning
180 }
181 return o.service.closeOwned(ctx, o.runtime, o.instance)
182 }
183
184 // Owner returns the host authority for the exact active instance. It fails for
185 // an unknown or superseded instance, so a delayed caller can never acquire
186 // authority over a same-ID successor. It is used by the flows that publish a
187 // brand-new identity in the same call; attaching to an existing session must go
188 // through Open and its ClientBinding instead.
189 func (s *Service) Owner(runtime *Runtime) (*RuntimeOwner, error) {
190 if runtime == nil || runtime.owner != s {
191 return nil, ErrSessionNotRunning
192 }
193 s.mu.Lock()
194 defer s.mu.Unlock()
195 if runtime.instance == "" || s.active[runtime.ref] != runtime {
196 return nil, ErrSessionNotRunning
197 }
198 return &RuntimeOwner{service: s, runtime: runtime, instance: runtime.instance}, nil
199 }
200
201 // Discard closes an unpublished candidate and releases its writer lease.
202 // Published runtimes must be closed through Service.Close so exact-instance
203 // unregistering cannot be bypassed.
204 func (s *Service) Discard(ctx context.Context, prepared *PreparedRuntime) error {
205 if prepared == nil || prepared.service != s || prepared.runtime == nil {
206 return nil
207 }
208 prepared.mu.Lock()
209 defer prepared.mu.Unlock()
210 if prepared.published {
211 return errors.New("session: published runtime cannot be discarded")
212 }
213 if prepared.discarded {
214 return prepared.runtime.close(ctx)
215 }
216 prepared.discarded = true
217 return prepared.runtime.close(ctx)
218 }
219
220 // Open attaches a client to an existing session and returns a binding that can
221 // only detach this client. It never hands out authority to close a runtime that
222 // another client may already be using.
223 func (s *Service) Open(ctx context.Context, ref SessionRef) (*ClientBinding, error) {
224 return s.openBinding(ctx, ref)
225 }
226
227 // openRuntime resolves the exact published instance. It is internal so only the
228 // service's own prepare/publish flows can reach a runtime without a grant.
229 func (s *Service) openRuntime(ctx context.Context, ref SessionRef) (*Runtime, error) {
230 if err := ref.validate(s.hostID); err != nil {
231 return nil, err
232 }
233 for {
234 s.mu.Lock()
235 if current := s.active[ref]; current != nil {
236 s.mu.Unlock()
237 return current, nil
238 }
239 if pending := s.preparing[ref]; pending != nil {
240 done := pending.done
241 s.mu.Unlock()
242 select {
243 case <-done:
244 continue
245 case <-ctx.Done():
246 return nil, ctx.Err()
247 }
248 }
249 pending := &prepareRuntime{done: make(chan struct{})}
250 s.preparing[ref] = pending
251 s.mu.Unlock()
252 break
253 }
254
255 var session *Session
256 var err error
257 if contextual, ok := s.persistence.(interface {
258 OpenContext(context.Context, string, AccessMode) (*Session, error)
259 }); ok {
260 session, err = contextual.OpenContext(ctx, ref.SessionID, ReadWrite)
261 } else {
262 session, err = s.persistence.Open(ref.SessionID, ReadWrite)
263 }
264 if err != nil {
265 s.finishPrepare(ref)
266 return nil, err
267 }
268 if _, _, recoverErr := session.RecoverInterrupted(ctx); recoverErr != nil {
269 _ = session.Close(context.Background())
270 s.finishPrepare(ref)
271 return nil, fmt.Errorf("session: close interrupted runtime: %w", recoverErr)
272 }
273 session.externalizeDurableHistory()
274 candidate, err := newRuntime(ref, session)
275 if err != nil {
276 closeErr := session.Close(context.Background())
277 s.finishPrepare(ref)
278 return nil, errors.Join(err, closeErr)
279 }
280 candidate.owner = s
281 s.mu.Lock()
282 pending := s.preparing[ref]
283 delete(s.preparing, ref)
284 if current := s.active[ref]; current != nil {
285 if pending != nil {
286 close(pending.done)
287 }
288 s.mu.Unlock()
289 _ = candidate.close(context.Background())
290 return current, nil
291 }
292 delete(s.closed, ref)
293 // Stamp the publish grant so Owner can hand the host authority over exactly
294 // this instance and no same-ID successor.
295 candidate.instance = randomID()
296 s.active[ref] = candidate
297 s.revision.Add(1)
298 if pending != nil {
299 close(pending.done)
300 }
301 s.mu.Unlock()
302 return candidate, nil
303 }
304
305 func (s *Service) finishPrepare(ref SessionRef) {
306 s.mu.Lock()
307 if pending := s.preparing[ref]; pending != nil {
308 delete(s.preparing, ref)
309 close(pending.done)
310 }
311 s.mu.Unlock()
312 }
313
314 func (s *Service) Runtime(ref SessionRef) (*Runtime, bool) {
315 if ref.validate(s.hostID) != nil {
316 return nil, false
317 }
318 s.mu.Lock()
319 runtime := s.active[ref]
320 s.mu.Unlock()
321 return runtime, runtime != nil
322 }
323
324 // Bind attaches a client to an exact published runtime. The returned binding
325 // can only release its own reference; it cannot dispose the shared runtime.
326 func (s *Service) Bind(runtime *Runtime) (*ClientBinding, error) {
327 if runtime == nil || runtime.owner != s {
328 return nil, ErrSessionNotRunning
329 }
330 s.mu.Lock()
331 defer s.mu.Unlock()
332 if s.retiring[runtime] != nil {
333 return nil, ErrRuntimeRetiring
334 }
335 if s.active[runtime.ref] != runtime {
336 return nil, ErrSessionNotRunning
337 }
338 s.bindings[runtime]++
339 delete(s.retireIdle, runtime)
340 if timer := s.idleTimers[runtime]; timer != nil {
341 s.removeIdleCacheLocked(runtime, true)
342 }
343 return &ClientBinding{service: s, runtime: runtime}, nil
344 }
345
346 // OpenBinding opens or reuses a runtime and attaches a client capability.
347 func (s *Service) openBinding(ctx context.Context, ref SessionRef) (*ClientBinding, error) {
348 for {
349 runtime, err := s.openRuntime(ctx, ref)
350 if err != nil {
351 return nil, err
352 }
353 binding, err := s.Bind(runtime)
354 if err == nil {
355 return binding, err
356 }
357 if errors.Is(err, ErrRuntimeRetiring) {
358 s.mu.Lock()
359 done := s.retiring[runtime]
360 s.mu.Unlock()
361 if done != nil {
362 select {
363 case <-done:
364 continue
365 case <-ctx.Done():
366 return nil, ctx.Err()
367 }
368 }
369 continue
370 }
371 return nil, err
372 }
373 }
374
375 func (s *Service) releaseBinding(ctx context.Context, runtime *Runtime) error {
376 s.mu.Lock()
377 count := s.bindings[runtime]
378 if count <= 1 {
379 delete(s.bindings, runtime)
380 s.retireIdle[runtime] = true
381 } else {
382 s.bindings[runtime] = count - 1
383 }
384 s.mu.Unlock()
385 if count <= 1 {
386 var flushErr error
387 if !runtime.executionBusy() {
388 _, flushErr = runtime.session.Flush(ctx)
389 }
390 s.scheduleIdleRetirement(runtime)
391 return flushErr
392 }
393 return nil
394 }
395
396 func (s *Service) closeIfUnbound(ctx context.Context, runtime *Runtime) error {
397 if runtime == nil {
398 return nil
399 }
400 s.scheduleIdleRetirement(runtime)
401 return nil
402 }
403
404 func (s *Service) scheduleIdleRetirement(runtime *Runtime) {
405 if runtime == nil {
406 return
407 }
408 weight := runtime.session.cacheWeight()
409 s.mu.Lock()
410 if s.active[runtime.ref] != runtime || s.bindings[runtime] != 0 || !s.retireIdle[runtime] {
411 s.mu.Unlock()
412 return
413 }
414 if runtime.executionBusy() {
415 s.mu.Unlock()
416 return
417 }
418 if s.idleTimers[runtime] != nil {
419 s.mu.Unlock()
420 return
421 }
422 ttl := s.idleTTL
423 s.idleClock++
424 s.idleOrder[runtime] = s.idleClock
425 idleEpoch := s.idleClock
426 s.idleWeight[runtime] = weight
427 s.idleUsed += weight
428 if ttl > 0 {
429 s.idleTimers[runtime] = time.AfterFunc(ttl, func() {
430 _ = s.retireIfUnboundAt(context.Background(), runtime, idleEpoch)
431 })
432 }
433 var victims []*Runtime
434 var sharedVictims []idlePoolEntry
435 if s.idlePool != nil {
436 sharedVictims = s.idlePool.add(s, runtime, weight)
437 }
438 for s.idlePool == nil && s.idleBudget >= 0 && s.idleUsed > s.idleBudget && len(s.idleOrder) > 0 {
439 var oldest *Runtime
440 var order uint64
441 for candidate, candidateOrder := range s.idleOrder {
442 if oldest == nil || candidateOrder < order {
443 oldest, order = candidate, candidateOrder
444 }
445 }
446 if oldest == nil {
447 break
448 }
449 s.removeIdleCacheLocked(oldest, true)
450 victims = append(victims, oldest)
451 }
452 s.mu.Unlock()
453 for _, victim := range victims {
454 _ = s.retireIfUnbound(context.Background(), victim)
455 }
456 for _, victim := range sharedVictims {
457 _ = victim.service.retireIfUnboundAt(context.Background(), victim.runtime, victim.epoch)
458 }
459 if ttl <= 0 && len(victims) == 0 {
460 _ = s.retireIfUnbound(context.Background(), runtime)
461 }
462 }
463
464 func (s *Service) removeIdleCacheLocked(runtime *Runtime, stop bool) {
465 s.idlePool.remove(runtime)
466 if timer := s.idleTimers[runtime]; timer != nil && stop {
467 timer.Stop()
468 }
469 delete(s.idleTimers, runtime)
470 s.idleUsed -= s.idleWeight[runtime]
471 if s.idleUsed < 0 {
472 s.idleUsed = 0
473 }
474 delete(s.idleWeight, runtime)
475 delete(s.idleOrder, runtime)
476 }
477
478 func (s *Service) retireIfUnbound(ctx context.Context, runtime *Runtime) error {
479 return s.retireIfUnboundAt(ctx, runtime, 0)
480 }
481
482 func (s *Service) retireIfUnboundAt(ctx context.Context, runtime *Runtime, epoch uint64) error {
483 s.mu.Lock()
484 if epoch != 0 && s.idleOrder[runtime] != epoch {
485 s.mu.Unlock()
486 return nil
487 }
488 s.removeIdleCacheLocked(runtime, false)
489 if s.active[runtime.ref] != runtime || s.bindings[runtime] != 0 || !s.retireIdle[runtime] || runtime.executionBusy() {
490 s.mu.Unlock()
491 return nil
492 }
493 if done := s.retiring[runtime]; done != nil {
494 s.mu.Unlock()
495 select {
496 case <-done:
497 return nil
498 case <-ctx.Done():
499 return ctx.Err()
500 }
501 }
502 done := make(chan struct{})
503 s.retiring[runtime] = done
504 s.mu.Unlock()
505
506 err := runtime.close(ctx)
507 s.mu.Lock()
508 delete(s.retiring, runtime)
509 if !errors.Is(err, ErrRuntimeBusy) && s.active[runtime.ref] == runtime && s.bindings[runtime] == 0 {
510 delete(s.active, runtime.ref)
511 delete(s.retireIdle, runtime)
512 s.removeIdleCacheLocked(runtime, true)
513 s.closed[runtime.ref] = err
514 s.revision.Add(1)
515 }
516 close(done)
517 s.mu.Unlock()
518 return err
519 }
520
520 lines GO