返回 DeepSeek-Reasonix
service.go
根目录 / internal / session / service.go
1 package session
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "log/slog"
8 "os"
9 "path/filepath"
10 "sync"
11 "sync/atomic"
12 "time"
13
14 "reasonix/internal/agent"
15 "reasonix/internal/event"
16 "reasonix/internal/transcript"
17 )
18
19 // SessionRef is the only execution identity used by the linear session
20 // service. Paths and UI selection are deliberately absent.
21 type SessionRef struct {
22 HostID string `json:"hostId"`
23 SessionID string `json:"sessionId"`
24 }
25
26 func (r SessionRef) validate(hostID string) error {
27 if r.HostID == "" || r.HostID != hostID {
28 return fmt.Errorf("session: session host %q does not match service host %q", r.HostID, hostID)
29 }
30 return validateSessionID(r.SessionID)
31 }
32
33 var (
34 ErrSessionNotRunning = errors.New("session runtime is not attached")
35 ErrRuntimeBusy = errors.New("session runtime already has an activity")
36 ErrRuntimeBound = errors.New("session runtime still has client bindings")
37 ErrRuntimeRetiring = errors.New("session runtime is retiring")
38 ErrRecoveryRequired = errors.New("session runtime requires recovery")
39 ErrStaleActivity = errors.New("session activity no longer owns commit authority")
40 ErrStaleExecution = errors.New("session execution generation no longer owns commit authority")
41 )
42
43 type RuntimePhase string
44
45 const (
46 RuntimeIdle RuntimePhase = "idle"
47 RuntimeRunning RuntimePhase = "running"
48 RuntimeCancelling RuntimePhase = "cancelling"
49 RuntimeFinalizing RuntimePhase = "finalizing"
50 RuntimeRecoveryRequired RuntimePhase = "recovery_required"
51 RuntimeClosed RuntimePhase = "closed"
52 )
53
54 type RuntimeSnapshot struct {
55 Ref SessionRef `json:"session"`
56 Epoch string `json:"runtimeEpoch"`
57 ActivityRevision uint64 `json:"activityRevision"`
58 Phase RuntimePhase `json:"phase"`
59 Activity string `json:"activity,omitempty"`
60 Session Snapshot `json:"sessionSnapshot"`
61 }
62
63 type CancelReceipt struct {
64 Ref SessionRef `json:"session"`
65 Accepted bool `json:"accepted"`
66 RuntimeEpoch string `json:"runtimeEpoch,omitempty"`
67 ActivityRevision uint64 `json:"activityRevision,omitempty"`
68 Phase RuntimePhase `json:"phase"`
69 }
70
71 // Runtime is the sole owner of a live Session and its write handle. Execution
72 // lifecycle lives in the bound turn-loop; persisted running events never
73 // create a Runtime after process restart.
74 type Runtime struct {
75 transcript *transcript.Projection
76 ref SessionRef
77 epoch string
78 session *Session
79 owner *Service
80 // instance stamps the publish grant so a delayed owner can prove it still
81 // refers to the exact instance it published.
82 instance string
83
84 mu sync.Mutex
85 phase RuntimePhase
86 activity string
87 revision atomic.Uint64
88 // execution is the generation-scoped turn-loop. Cancel loads it without
89 // taking mu so Stop never waits on a commit or persistence lock.
90 execution atomic.Pointer[executionBinding]
91 bindGen atomic.Uint64
92 canceling atomic.Bool
93 closeDone chan struct{}
94 closeErr error
95 }
96
97 func newRuntime(ref SessionRef, session *Session) (*Runtime, error) {
98 runtime, err := initializeRuntime(ref, session)
99 // Logging can perform I/O; release the session lock before emitting.
100 var diagnostic *TranscriptInitializationError
101 if errors.As(err, &diagnostic) {
102 slog.Error("session transcript initialization failed", "diagnostic", diagnostic)
103 }
104 return runtime, err
105 }
106
107 func initializeRuntime(ref SessionRef, session *Session) (*Runtime, error) {
108 session.mu.Lock()
109 defer session.mu.Unlock()
110 runtime := &Runtime{ref: ref, epoch: randomID(), session: session, phase: RuntimeIdle}
111 baseline := session.recentMessages
112 if len(baseline) == 0 {
113 baseline = session.projection.Messages
114 }
115 totalMessages := len(baseline)
116 if len(baseline) > 96 {
117 baseline = baseline[len(baseline)-96:]
118 }
119 projection, err := transcript.NewProjection(transcript.Identity{SessionID: ref.SessionID, RuntimeEpoch: runtime.epoch}, session.transcriptRows(baseline), session.next-1)
120 if err != nil {
121 return nil, &TranscriptInitializationError{sessionID: ref.SessionID, covered: session.next - 1,
122 messageCount: len(baseline), totalMessages: totalMessages, cause: err}
123 }
124 runtime.transcript = projection
125 durable := uint64(0)
126 if session.binding != nil {
127 durable, _, _ = session.binding.progress()
128 }
129 restored := transcript.Runtime{TurnID: session.projection.TurnID, Status: session.projection.TurnStatus, FinalMessageID: session.projection.CurrentTurnMessageID}
130 if receipt, ok := session.projection.Submissions.byTurn[session.id+"\x00"+restored.TurnID]; ok {
131 restored.SubmissionID = receipt.SubmissionID
132 }
133 restored.SamplingCount, restored.ToolCount = len(session.projection.CurrentAttempts), len(session.projection.CurrentCalls)
134 if restored.TurnID != "" && !restored.Status.Terminal() {
135 // A persisted open turn is recovery evidence, not a running model.
136 restored.Status = event.TurnRecoveryRequired
137 }
138 if restored.TurnID == "" && len(session.projection.Turns) > 0 {
139 last := session.projection.Turns[len(session.projection.Turns)-1]
140 restored.TurnID, restored.FinalMessageID = last.TurnID, last.MessageID
141 restored.DurationMs = last.DurationMs
142 restored.SamplingCount, restored.ToolCount = last.SamplingCount, last.ToolCount
143 }
144 for _, message := range baseline {
145 if message.ID == restored.FinalMessageID {
146 restored.DurationMs = max(restored.DurationMs, message.WorkDurationMs)
147 }
148 }
149 runtime.transcript.RestoreRuntime(restored, durable)
150 session.transcript = runtime.transcript
151 runtime.revision.Store(1)
152 return runtime, nil
153 }
154
155 func (r *Runtime) Ref() SessionRef { return r.ref }
156
157 func (r *Runtime) Session() *Session { return r.session }
158
159 func (r *Runtime) Snapshot() RuntimeSnapshot {
160 state := r.activitySnapshot()
161 state.Session = r.session.Snapshot()
162 return state
163 }
164
165 func (r *Runtime) StateSnapshot() RuntimeSnapshot {
166 state := r.activitySnapshot()
167 state.Session = r.session.StateSnapshot()
168 return state
169 }
170
171 // ExecutionSnapshot returns the provider projection and lightweight turn
172 // boundaries without reconstructing the durable UI transcript.
173 func (r *Runtime) ExecutionSnapshot() RuntimeSnapshot {
174 state := r.activitySnapshot()
175 state.Session = r.session.ExecutionSnapshot()
176 return state
177 }
178
179 func (r *Runtime) activitySnapshot() RuntimeSnapshot {
180 r.mu.Lock()
181 phase := r.phase
182 if r.canceling.Load() && phase == RuntimeRunning {
183 phase = RuntimeCancelling
184 }
185 state := RuntimeSnapshot{Ref: r.ref, Epoch: r.epoch, ActivityRevision: r.revision.Load(), Phase: phase, Activity: r.activity}
186 r.mu.Unlock()
187 return state
188 }
189
190 // Cancel forwards Stop to the bound turn-loop without taking the runtime
191 // mutex. An unbound runtime is already idle.
192 func (r *Runtime) Cancel() bool {
193 var binding *executionBinding
194 var revision uint64
195 for {
196 exec := r.loadExecution()
197 if exec == nil || exec.control == nil {
198 return false
199 }
200 revision = r.revision.Load()
201 if !exec.control.Cancel() {
202 // A host cutover may linearize while Cancel is inside the outgoing
203 // loop. Retry only when ownership actually changed; a stable owner
204 // rejecting Cancel remains a normal idle result.
205 if r.loadExecution() != exec {
206 continue
207 }
208 return false
209 }
210 binding = exec
211 break
212 }
213 // The callback may finish its worker and start queued work before it
214 // returns. Never apply its acknowledgement to that successor's activity.
215 update := func() {
216 defer r.mu.Unlock()
217 if r.execution.Load() == binding && r.revision.Load() == revision && r.phase == RuntimeRunning {
218 r.canceling.Store(true)
219 r.phase = RuntimeCancelling
220 if r.activity != MaintenanceActivity {
221 r.activity = "cancelling"
222 }
223 r.revision.Add(1)
224 }
225 }
226 if r.mu.TryLock() {
227 update()
228 } else {
229 go func() { r.mu.Lock(); update() }()
230 }
231 return true
232 }
233
234 func (r *Runtime) RequireRecovery(activity string) {
235 r.mu.Lock()
236 defer r.mu.Unlock()
237 if r.phase == RuntimeClosed {
238 return
239 }
240 r.phase = RuntimeRecoveryRequired
241 r.activity = activity
242 r.revision.Add(1)
243 }
244
245 func (r *Runtime) close(ctx context.Context) error {
246 r.mu.Lock()
247 if r.closeDone != nil {
248 done := r.closeDone
249 r.mu.Unlock()
250 <-done
251 return r.closeErr
252 }
253 if r.phase.busy() && r.loadExecution() != nil {
254 r.mu.Unlock()
255 return ErrRuntimeBusy
256 }
257 // Seal admission in the same critical section as the idle check. The
258 // irreversible close has one uncancellable result for every caller.
259 r.closeDone = make(chan struct{})
260 r.phase = RuntimeClosed
261 r.activity = ""
262 r.revision.Add(1)
263 r.mu.Unlock()
264 r.transcript.CloseFollowers()
265 r.closeErr = r.session.close(context.Background())
266 close(r.closeDone)
267 return r.closeErr
268 }
269
270 func osClosedError() error { return errors.New("session runtime is closed") }
271
272 // Service applies DSH's prepare/publish/exact-detach rule. Candidate handles
273 // are opened outside the registry lock; only the exact published Runtime can
274 // later unregister itself.
275 type Service struct {
276 hostID string
277 persistence SessionPersistence
278
279 mu sync.Mutex
280 active map[SessionRef]*Runtime
281 closed map[SessionRef]error
282 preparing map[SessionRef]*prepareRuntime
283 bindings map[*Runtime]int
284 retiring map[*Runtime]chan struct{}
285 retireIdle map[*Runtime]bool
286 idleTimers map[*Runtime]*time.Timer
287 idleWeight map[*Runtime]int64
288 idleOrder map[*Runtime]uint64
289 idleUsed int64
290 idleClock uint64
291 idleBudget int64
292 idlePool *IdlePool
293 idleTTL time.Duration
294 query *Query
295 revision atomic.Uint64
296 }
297
298 func (s *Service) Cancel(ref SessionRef) (RuntimeSnapshot, error) {
299 runtime, ok := s.Runtime(ref)
300 if !ok {
301 return RuntimeSnapshot{}, ErrSessionNotRunning
302 }
303 runtime.Cancel()
304 return runtime.StateSnapshot(), nil
305 }
306
307 // CancelSession is the public session-scoped Stop contract. A missing runtime
308 // is already idle and therefore succeeds idempotently; no caller-supplied turn
309 // id participates in routing or authorization.
310 func (s *Service) CancelSession(ref SessionRef) (CancelReceipt, error) {
311 if err := ref.validate(s.hostID); err != nil {
312 return CancelReceipt{}, err
313 }
314 runtime, ok := s.Runtime(ref)
315 if !ok {
316 return CancelReceipt{Ref: ref, Accepted: true, Phase: RuntimeIdle}, nil
317 }
318 if runtime.Cancel() {
319 return CancelReceipt{
320 Ref: ref,
321 Accepted: true,
322 RuntimeEpoch: runtime.epoch,
323 ActivityRevision: runtime.revision.Load(),
324 Phase: RuntimeCancelling,
325 }, nil
326 }
327 snapshot := runtime.activitySnapshot()
328 return CancelReceipt{Ref: ref, Accepted: true, RuntimeEpoch: snapshot.Epoch, ActivityRevision: snapshot.ActivityRevision, Phase: snapshot.Phase}, nil
329 }
330
331 func (s *Service) Flush(ctx context.Context, ref SessionRef) (DurableReceipt, error) {
332 runtime, ok := s.Runtime(ref)
333 if !ok {
334 return DurableReceipt{}, ErrSessionNotRunning
335 }
336 return runtime.session.Flush(ctx)
337 }
338
339 // ContinueLegacy freezes one legacy head, publishes its deterministic final
340 // session, then attaches that exact session. It never writes the source and it
341 // does not accept the caller's pending submission; hosts enqueue the unchanged
342 // submission only after this method returns the new immutable identity.
343 func (s *Service) ContinueLegacy(ctx context.Context, sourcePath, headID string) (*Runtime, MigrationResult, error) {
344 filesystem, ok := s.persistence.(*FilesystemPersistence)
345 if !ok {
346 return nil, MigrationResult{}, errors.New("session: persistence does not support legacy migration")
347 }
348 result, err := migrateLegacyHeadForHost(ctx, sourcePath, filesystem.Root, headID)
349 if err != nil {
350 return nil, result, err
351 }
352 runtime, err := s.openRuntime(ctx, SessionRef{HostID: s.hostID, SessionID: result.TargetID})
353 return runtime, result, err
354 }
355
356 // ContinueImported resolves the paired legacy transcript and retired event
357 // sidecar as one frozen migration decision. It refuses divergent histories
358 // instead of letting a caller accidentally resume whichever source it opened
359 // first.
360 func (s *Service) ContinueImported(ctx context.Context, sourcePath, headID string) (*Runtime, ImportResult, error) {
361 return s.ContinueImportedWithHeader(ctx, sourcePath, headID, CreateOptions{})
362 }
363
364 // ContinueImportedWithHeader publishes Desktop ownership in the same atomic
365 // directory publication as the imported history.
366 func (s *Service) ContinueImportedWithHeader(ctx context.Context, sourcePath, headID string, options CreateOptions) (*Runtime, ImportResult, error) {
367 filesystem, ok := s.persistence.(*FilesystemPersistence)
368 if !ok {
369 return nil, ImportResult{}, errors.New("session: persistence does not support imported sessions")
370 }
371 result, err := importSourceForLegacyWithHeader(ctx, sourcePath, filesystem.Root, headID, options)
372 if err != nil {
373 return nil, result, err
374 }
375 runtime, err := s.openRuntime(ctx, SessionRef{HostID: s.hostID, SessionID: result.TargetID})
376 return runtime, result, err
377 }
378
379 // ContinuePrototype is the explicit, fail-closed bridge for the retired
380 // sidecar codec. Unknown required events or conflicting tails remain read-only.
381 func (s *Service) ContinuePrototype(ctx context.Context, sourceDir string) (*Runtime, PrototypeImportResult, error) {
382 filesystem, ok := s.persistence.(*FilesystemPersistence)
383 if !ok {
384 return nil, PrototypeImportResult{}, errors.New("session: persistence does not support prototype import")
385 }
386 result, err := ImportPrototype(ctx, sourceDir, filesystem.Root)
387 if err != nil {
388 return nil, result, err
389 }
390 runtime, err := s.openRuntime(ctx, SessionRef{HostID: s.hostID, SessionID: result.TargetID})
391 return runtime, result, err
392 }
393
394 // ContinueImportedFrom resolves legacy history against its original paired
395 // store while publishing only into this service's separate staging root.
396 func (s *Service) ContinueImportedFrom(ctx context.Context, sourcePath, sourceRoot, headID string) (*Runtime, ImportResult, error) {
397 return s.ContinueImportedSource(ctx, sourcePath, filepath.Join(sourceRoot, agent.BranchID(sourcePath)), headID)
398 }
399
400 // ContinueImportedSource accepts a provenance-linked directory whose identity
401 // may have changed when a legacy head was previously converted.
402 func (s *Service) ContinueImportedSource(ctx context.Context, sourcePath, sourceDir, headID string) (*Runtime, ImportResult, error) {
403 filesystem, ok := s.persistence.(*FilesystemPersistence)
404 if !ok {
405 return nil, ImportResult{}, errors.New("session: persistence does not support imported sessions")
406 }
407 result, err := importSourceForLegacyAt(ctx, sourcePath, sourceDir, filesystem.Root, headID, CreateOptions{})
408 if err != nil {
409 return nil, result, err
410 }
411 sourceRoot := filepath.Dir(sourceDir)
412 if result.Kind == "final" && filepath.Clean(sourceRoot) != filepath.Clean(filesystem.Root) {
413 tmp, err := os.MkdirTemp("", "reasonix-canonical-stage-")
414 if err != nil {
415 return nil, result, err
416 }
417 defer os.RemoveAll(tmp)
418 bundle := filepath.Join(tmp, "bundle")
419 if err := NewFilesystemPersistence(sourceRoot).exportCold(ctx, result.TargetID, bundle); err != nil {
420 return nil, result, err
421 }
422 if _, err := s.Import(ctx, bundle); err != nil {
423 return nil, result, err
424 }
425 }
426 runtime, err := s.openRuntime(ctx, SessionRef{HostID: s.hostID, SessionID: result.TargetID})
427 return runtime, result, err
428 }
429
430 // ContinueStoredPreview upgrades a pre-ownership linear store selected by its
431 // former session id. The old directory remains read-only; execution resumes on
432 // the deterministic final-codec identity returned here.
433 func (s *Service) ContinueStoredPreview(ctx context.Context, sessionID string) (*Runtime, PrototypeImportResult, error) {
434 filesystem, ok := s.persistence.(*FilesystemPersistence)
435 if !ok {
436 return nil, PrototypeImportResult{}, errors.New("session: persistence does not support preview import")
437 }
438 if err := validateSessionID(sessionID); err != nil {
439 return nil, PrototypeImportResult{}, err
440 }
441 sourceDir, err := filesystem.sessionDir(sessionID, true)
442 if err != nil {
443 return nil, PrototypeImportResult{}, err
444 }
445 result, err := ImportStoredPreview(ctx, sourceDir, filesystem.Root)
446 if err != nil {
447 return nil, result, err
448 }
449 runtime, err := s.openRuntime(ctx, SessionRef{HostID: s.hostID, SessionID: result.TargetID})
450 return runtime, result, err
451 }
452
453 // Fork creates an independent child at the exact end event of a completed
454 // turn. No message-count inference is involved.
455 func (s *Service) Fork(ctx context.Context, ref SessionRef, afterTurnID, childID string) (*Runtime, error) {
456 runtime, ok := s.Runtime(ref)
457 if !ok {
458 return nil, ErrSessionNotRunning
459 }
460 turn, ok := completedTurn(runtime.session.Snapshot().Projection.Turns, afterTurnID)
461 if !ok {
462 return nil, fmt.Errorf("session: completed turn %q not found", afterTurnID)
463 }
464 return s.forkAt(ctx, runtime, turn.EndSequence, childID)
465 }
466
467 // Rewind creates a child from the event immediately before beforeTurnID.
468 func (s *Service) Rewind(ctx context.Context, ref SessionRef, beforeTurnID, childID string) (*Runtime, error) {
469 runtime, ok := s.Runtime(ref)
470 if !ok {
471 return nil, ErrSessionNotRunning
472 }
473 turn, ok := completedTurn(runtime.session.Snapshot().Projection.Turns, beforeTurnID)
474 if !ok {
475 return nil, fmt.Errorf("session: completed turn %q not found", beforeTurnID)
476 }
477 return s.forkAt(ctx, runtime, turn.StartSequence-1, childID)
478 }
479
480 func completedTurn(turns []TurnBoundary, id string) (TurnBoundary, bool) {
481 for _, turn := range turns {
482 if turn.TurnID == id {
483 return turn, true
484 }
485 }
486 return TurnBoundary{}, false
487 }
488
489 func (s *Service) forkAt(ctx context.Context, parent *Runtime, sequence uint64, childID string) (*Runtime, error) {
490 filesystem, ok := s.persistence.(*FilesystemPersistence)
491 if !ok {
492 return nil, errors.New("session: persistence does not support filesystem fork")
493 }
494 if childID == "" {
495 childID = randomID()
496 }
497 if err := validateSessionID(childID); err != nil {
498 return nil, err
499 }
500 childDir := filepath.Join(filesystem.Root, childID)
501 if _, err := parent.session.Fork(ctx, childDir, childID, sequence); err != nil {
502 return nil, err
503 }
504 return s.openRuntime(ctx, SessionRef{HostID: s.hostID, SessionID: childID})
505 }
506
507 func (s *Service) Close(ctx context.Context, ref SessionRef) error {
508 if err := ref.validate(s.hostID); err != nil {
509 return err
510 }
511 s.mu.Lock()
512 runtime := s.active[ref]
513 closedErr, closed := s.closed[ref]
514 s.mu.Unlock()
515 if runtime == nil {
516 if closed {
517 return closedErr
518 }
519 return ErrSessionNotRunning
520 }
521 return s.closeOwned(ctx, runtime, "")
522 }
523
524 // closeOwned is the teardown entry point for a RuntimeOwner holding one exact
525 // instance grant. A delayed old disposer must never close its same-ID
526 // successor, and a client-bound runtime is never torn down underneath it.
527 func (s *Service) closeOwned(ctx context.Context, runtime *Runtime, instance string) error {
528 return s.closeRuntime(ctx, runtime, instance, false)
529 }
530
531 // closeRuntime adds the terminal variant Shutdown needs. Refusing a bound
532 // runtime is right while the process keeps running, but at shutdown it would
533 // strand the writer lease and recovery handles for the process lifetime, so
534 // the final teardown releases them and still reports the leaked binding.
535 func (s *Service) closeRuntime(ctx context.Context, runtime *Runtime, instance string, terminal bool) error {
536 if runtime == nil {
537 return ErrSessionNotRunning
538 }
539 if err := runtime.ref.validate(s.hostID); err != nil {
540 return err
541 }
542 if instance != "" && runtime.instance != instance {
543 return ErrSessionNotRunning
544 }
545 var leaked error
546 s.mu.Lock()
547 if timer := s.idleTimers[runtime]; timer != nil {
548 s.removeIdleCacheLocked(runtime, true)
549 }
550 if s.bindings[runtime] != 0 {
551 if !terminal {
552 s.mu.Unlock()
553 return ErrRuntimeBound
554 }
555 delete(s.bindings, runtime)
556 leaked = fmt.Errorf("%w: %s", ErrRuntimeBound, runtime.ref.SessionID)
557 }
558 if s.active[runtime.ref] != runtime {
559 s.mu.Unlock()
560 return errors.Join(leaked, runtime.close(ctx))
561 }
562 if done := s.retiring[runtime]; done != nil {
563 s.mu.Unlock()
564 select {
565 case <-done:
566 return errors.Join(leaked, runtime.close(ctx))
567 case <-ctx.Done():
568 return errors.Join(leaked, ctx.Err())
569 }
570 }
571 done := make(chan struct{})
572 s.retiring[runtime] = done
573 s.mu.Unlock()
574 err := runtime.close(ctx)
575 s.mu.Lock()
576 delete(s.retiring, runtime)
577 if !errors.Is(err, ErrRuntimeBusy) && s.active[runtime.ref] == runtime {
578 delete(s.active, runtime.ref)
579 delete(s.retireIdle, runtime)
580 s.removeIdleCacheLocked(runtime, true)
581 s.closed[runtime.ref] = err
582 s.revision.Add(1)
583 }
584 close(done)
585 s.mu.Unlock()
586 return errors.Join(leaked, err)
587 }
588
589 // Detach removes a runtime only if it is still the exact published instance.
590 // It is used by host callbacks that may arrive after a replacement.
591 func (s *Service) Detach(runtime *Runtime) bool {
592 if runtime == nil {
593 return false
594 }
595 s.mu.Lock()
596 defer s.mu.Unlock()
597 if s.active[runtime.ref] != runtime {
598 return false
599 }
600 if timer := s.idleTimers[runtime]; timer != nil {
601 s.removeIdleCacheLocked(runtime, true)
602 }
603 delete(s.active, runtime.ref)
604 s.revision.Add(1)
605 return true
606 }
607
608 type ObserveResult struct {
609 Runtime *RuntimeSnapshot `json:"runtime,omitempty"`
610 Events EventPage `json:"events"`
611 }
612
613 // SessionDir resolves the on-disk directory of a final-format identity
614 // without opening it. Hosts use it to probe writer occupancy for takeover
615 // flows; the writer lease itself is never taken here.
616 func (s *Service) SessionDir(ctx context.Context, ref SessionRef) (string, error) {
617 if err := ref.validate(s.hostID); err != nil {
618 return "", err
619 }
620 info, err := s.persistence.Stat(ctx, ref.SessionID)
621 if err != nil {
622 return "", err
623 }
624 return info.Path, nil
625 }
626
627 func (s *Service) Observe(ctx context.Context, ref SessionRef, cursor uint64, limit int) (ObserveResult, error) {
628 if err := ref.validate(s.hostID); err != nil {
629 return ObserveResult{}, err
630 }
631 if runtime, ok := s.Runtime(ref); ok {
632 // Observe reports runtime state plus an explicitly paged event tail. It
633 // must not duplicate the provider model workset into every poll.
634 snapshot := runtime.StateSnapshot()
635 page, err := runtime.session.AcceptedPage(ctx, cursor, limit)
636 return ObserveResult{Runtime: &snapshot, Events: page}, err
637 }
638 handle, err := s.persistence.Open(ref.SessionID, ReadOnly)
639 if err != nil {
640 return ObserveResult{}, err
641 }
642 defer handle.Close(context.Background())
643 page, err := handle.Read(ctx, cursor, limit)
644 return ObserveResult{Events: page}, err
645 }
646
646 lines GO