返回 DeepSeek-Reasonix
session.go
根目录 / internal / session / session.go
1 package session
2
3 import (
4 "bytes"
5 "context"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "os"
10 "strings"
11 "sync"
12 "time"
13
14 "reasonix/internal/provider"
15 "reasonix/internal/sessioncontent"
16 "reasonix/internal/transcript"
17 )
18
19 // Session owns the live business state for one session identity: sequence
20 // allocation, the bounded accepted tail, compact operation identities, and the
21 // current projection. Durable history bodies and searchable operation details
22 // live behind Query; the physical handle owns bytes and the writer lease.
23 //
24 // A commit becomes observable here before it is durable. That is deliberate and
25 // mirrors DSH: accepting an event updates the projection and the UI, while the
26 // binding batches the write-behind and Flush marks a semantic checkpoint.
27 type Session struct {
28 transcript *transcript.Projection
29 mu sync.Mutex
30 id string
31 manifest Manifest
32 next uint64
33 commits []Commit
34 operations map[string]operationRecord
35 projection Projection
36 binding *PersistenceBinding
37 sealed bool
38 sealedError error
39 readOnly bool
40 // externalHistory means durable UI messages live in HistoryQuery rather
41 // than this runtime projection. Messages then contains only the accepted,
42 // not-yet-durable tail; ModelMessages remains the exact provider workset.
43 externalHistory bool
44 catalogPreview string
45 recentMessages []provider.Message
46 durableRecent []provider.Message
47 messageIDs messageIdentities
48 storageGeneration string
49 recovery *recoveryStore
50 // coldHandle backs a read-only session, which has no binding because it
51 // never enqueues or drains anything.
52 coldHandle SessionHandle
53 }
54
55 // PreparedBatch is a validated, self-contained commit payload. Every expensive
56 // or fallible step — payload copying, schema validation, turn identity
57 // derivation, and the operation hash — happens here, before the caller takes
58 // the activity commit gate. CommitPrepared then only assigns identity and
59 // extends the log.
60 type PreparedBatch struct {
61 sessionID string
62 writerGeneration uint64
63 operationID string
64 turnID string
65 events []Event
66 storedEvents []Event
67 hash string
68 reservation *queueReservation
69 prior *operationRecord
70 }
71
72 // OperationID reports the stable idempotency key of the prepared batch.
73 func (p PreparedBatch) OperationID() string { return p.operationID }
74
75 // Empty reports whether the batch carries no committable event.
76 func (p PreparedBatch) Empty() bool { return len(p.events) == 0 }
77
78 // Release returns queue capacity when a prepared batch loses its activity or
79 // CAS race before acceptance. It is safe after CommitPrepared consumes it.
80 func (p PreparedBatch) Release() { p.reservation.release() }
81
82 // newSession builds the live in-memory session over an already-open binding.
83 // commits is the durable prefix replayed by the handle; the projection is
84 // rebuilt from it so no business state is inherited from the physical layer.
85 func newSession(id string, manifest Manifest, commits []Commit, projection Projection, binding *PersistenceBinding) *Session {
86 operations := make(map[string]operationRecord, len(commits))
87 next := uint64(1)
88 for _, commit := range commits {
89 operations[commit.OperationID] = compactOperationRecord(commit)
90 next = commit.LastSequence() + 1
91 }
92 ids := make([]string, 0, len(projection.Messages))
93 for _, message := range projection.Messages {
94 ids = append(ids, message.ID)
95 }
96 return &Session{
97 id: id, manifest: manifest, next: next,
98 operations: operations, projection: projection, binding: binding,
99 messageIDs: identitiesOf(ids),
100 }
101 }
102
103 // ID returns the immutable session identity.
104 func (s *Session) ID() string {
105 if s == nil {
106 return ""
107 }
108 s.mu.Lock()
109 defer s.mu.Unlock()
110 return s.id
111 }
112
113 // Manifest returns the persistence manifest that describes this session.
114 func (s *Session) Manifest() Manifest {
115 if s == nil {
116 return Manifest{}
117 }
118 s.mu.Lock()
119 defer s.mu.Unlock()
120 manifest := s.manifest
121 if manifest.Source != nil {
122 source := *manifest.Source
123 manifest.Source = &source
124 }
125 return manifest
126 }
127
128 // ExplicitPermissionPreset excludes a fork's inherited parent choice. Only a
129 // preset event written under this session's own identity is its user choice.
130 func (s *Session) ExplicitPermissionPreset() string {
131 if s == nil {
132 return ""
133 }
134 s.mu.Lock()
135 defer s.mu.Unlock()
136 sequence := s.projection.PermissionPresetSequence
137 if sequence == 0 || (s.manifest.Source != nil && sequence <= s.manifest.InheritedEvents) {
138 return ""
139 }
140 return s.projection.PermissionPreset
141 }
142
143 // EventSequence reports the last accepted sequence. Accepted events may still
144 // be waiting in the binding's write-behind queue.
145 func (s *Session) EventSequence() uint64 {
146 if s == nil {
147 return 0
148 }
149 s.mu.Lock()
150 defer s.mu.Unlock()
151 return s.next - 1
152 }
153
154 // PrepareBatch validates the logical batch and computes its idempotency digest
155 // without touching the commit lock. Callers must treat the result as immutable.
156 func (s *Session) PrepareBatch(operationID string, batch Batch) (PreparedBatch, error) {
157 return s.PrepareBatchContext(context.Background(), operationID, batch)
158 }
159
160 // PrepareBatchContext publishes large immutable payloads before the batch can
161 // enter the accepted sequence. It performs all disk I/O outside the session
162 // commit lock; CommitPrepared only revalidates identity and generation.
163 func (s *Session) PrepareBatchContext(ctx context.Context, operationID string, batch Batch) (PreparedBatch, error) {
164 if s == nil {
165 return PreparedBatch{}, fmt.Errorf("session: nil session")
166 }
167 operationID = strings.TrimSpace(operationID)
168 if operationID == "" {
169 operationID = strings.TrimSpace(batch.OperationID)
170 }
171 if operationID == "" || len(batch.Events) == 0 {
172 return PreparedBatch{}, fmt.Errorf("session: operation id and events are required")
173 }
174 turnID := strings.TrimSpace(batch.TurnID)
175 if turnID == "" && batchContains(batch.Events, "turn/start") {
176 turnID = deterministicID("turn\x00" + s.id + "\x00" + operationID)
177 }
178 events := cloneEvents(batch.Events)
179 for i := range events {
180 event := &events[i]
181 event.Kind = strings.TrimSpace(event.Kind)
182 if event.Kind == "message/retract" {
183 if event.Optional {
184 return PreparedBatch{}, fmt.Errorf("message/retract must be required")
185 }
186 event.Required = true
187 }
188 if event.Kind == "" {
189 return PreparedBatch{}, fmt.Errorf("session: events[%d].kind is required", i)
190 }
191 if !event.Optional && !ProjectionKinds[event.Kind] {
192 return PreparedBatch{}, fmt.Errorf("%w: unknown required event %q", ErrUnsupportedVersion, event.Kind)
193 }
194 if err := checkReplacementIdentities(*event); err != nil {
195 return PreparedBatch{}, err
196 }
197 // The durable sequence is assigned at commit time so a rejected batch
198 // never consumes one.
199 event.Sequence = 0
200 }
201 // The digest covers exactly what the caller supplied. Generated event ids
202 // are deliberately excluded so retrying the same logical batch is
203 // idempotent instead of reporting a spurious operation conflict.
204 hash, err := hashOperation(s.id, turnID, events)
205 if err != nil {
206 return PreparedBatch{}, err
207 }
208 prior, found, err := s.lookupOperation(operationID)
209 if err != nil {
210 return PreparedBatch{}, fmt.Errorf("session: lookup operation %q: %w", operationID, err)
211 }
212 if found && prior.hash != hash {
213 return PreparedBatch{}, fmt.Errorf("%w: %q", ErrOperationConflict, operationID)
214 }
215 for i := range events {
216 if events[i].ID == "" {
217 events[i].ID = randomID()
218 }
219 }
220 s.mu.Lock()
221 sessionID, writerGeneration, binding, manifestCodec := s.id, s.manifest.WriterGeneration, s.binding, s.manifest.Codec
222 s.mu.Unlock()
223 if binding == nil {
224 return PreparedBatch{}, ErrReadOnly
225 }
226 storedEvents := cloneEvents(events)
227 if manifestCodec == Codec {
228 content := s.contentStore()
229 for i := range storedEvents {
230 if len(storedEvents[i].Payload) <= v4InlinePayloadBytes {
231 continue
232 }
233 if content == nil {
234 return PreparedBatch{}, errors.New("session: content store unavailable for large event payload")
235 }
236 ref, err := content.Put(ctx, bytes.NewReader(storedEvents[i].Payload), sessioncontent.Metadata{MediaType: "application/json"})
237 if err != nil {
238 return PreparedBatch{}, fmt.Errorf("prepare event %s content: %w", storedEvents[i].ID, err)
239 }
240 storedEvents[i].Payload = nil
241 storedEvents[i].PayloadRef = &ref
242 }
243 }
244 reservation, err := binding.reserve(ctx, commitHotBytes(Commit{Events: storedEvents}))
245 if err != nil {
246 return PreparedBatch{}, err
247 }
248 prepared := PreparedBatch{sessionID: sessionID, writerGeneration: writerGeneration, operationID: operationID, turnID: turnID, events: events, storedEvents: storedEvents, hash: hash, reservation: reservation}
249 if found {
250 copy := prior
251 prepared.prior = &copy
252 }
253 return prepared, nil
254 }
255
256 func (s *Session) lookupOperation(operationID string) (operationRecord, bool, error) {
257 if s == nil {
258 return operationRecord{}, false, nil
259 }
260 s.mu.Lock()
261 if record, ok := s.operations[operationID]; ok {
262 s.mu.Unlock()
263 return record, true, nil
264 }
265 recovery := s.recovery
266 s.mu.Unlock()
267 return recovery.lookupOperation(operationID)
268 }
269
270 func (s *Session) contentStore() *sessioncontent.Store {
271 if s == nil || s.binding == nil {
272 return nil
273 }
274 store, _ := s.binding.handle.(*Store)
275 if store == nil {
276 return nil
277 }
278 return store.content
279 }
280
281 func (s *Session) ContentStore() *sessioncontent.Store { return s.contentStore() }
282
283 func (s *Session) StorageGeneration() string {
284 if s == nil {
285 return ""
286 }
287 s.mu.Lock()
288 defer s.mu.Unlock()
289 return s.storageGeneration
290 }
291
292 // CommitPrepared appends an already validated batch under one short memory
293 // lock. The persistence binding only receives an immutable batch into its
294 // write-behind queue, so this never performs file I/O and never blocks on a
295 // subscriber.
296 func (s *Session) CommitPrepared(prepared PreparedBatch) (Commit, error) {
297 return s.commitPrepared(prepared, nil)
298 }
299
300 func (s *Session) commitPrepared(prepared PreparedBatch, expectedTitleSequence *uint64) (Commit, error) {
301 defer prepared.Release()
302 if s == nil {
303 return Commit{}, fmt.Errorf("session: nil session")
304 }
305 if prepared.Empty() {
306 return Commit{}, fmt.Errorf("session: operation id and events are required")
307 }
308 s.mu.Lock()
309 if expectedTitleSequence != nil && s.projection.TitleSequence != *expectedTitleSequence {
310 s.mu.Unlock()
311 return Commit{}, ErrSessionTitleChanged
312 }
313 if prepared.sessionID != s.id || prepared.writerGeneration != s.manifest.WriterGeneration {
314 s.mu.Unlock()
315 return Commit{}, ErrStaleGeneration
316 }
317 if s.readOnly {
318 s.mu.Unlock()
319 return Commit{}, ErrReadOnly
320 }
321 if s.sealed {
322 err := s.sealedError
323 s.mu.Unlock()
324 if err == nil {
325 err = osClosedError()
326 }
327 return Commit{}, err
328 }
329 if prior, ok := s.operations[prepared.operationID]; ok {
330 if prior.hash != prepared.hash {
331 s.mu.Unlock()
332 return Commit{}, fmt.Errorf("%w: %q", ErrOperationConflict, prepared.operationID)
333 }
334 commit := prior.commit
335 commit.Events = cloneEvents(prepared.events)
336 for i := range commit.Events {
337 commit.Events[i].Sequence = commit.FirstSequence + uint64(i)
338 }
339 s.mu.Unlock()
340 return commit, nil
341 }
342 if prepared.prior != nil {
343 commit := prepared.prior.commit
344 commit.Events = cloneEvents(prepared.events)
345 for i := range commit.Events {
346 commit.Events[i].Sequence = commit.FirstSequence + uint64(i)
347 }
348 s.mu.Unlock()
349 return commit, nil
350 }
351 commitSchema, commitCodec := SchemaVersion, Codec
352 if s.manifest.Codec != Codec {
353 commitSchema, commitCodec = 3, s.manifest.Codec
354 }
355 commit := Commit{
356 SchemaVersion: commitSchema, Codec: commitCodec, RecordType: "commit", ID: randomID(),
357 OperationID: prepared.operationID, OperationHash: prepared.hash, FirstSequence: s.next,
358 EventCount: len(prepared.events), TurnID: prepared.turnID,
359 WriterGeneration: s.manifest.WriterGeneration, CreatedAt: time.Now().UTC(),
360 // PreparedBatch already owns a private clone. Transfer it into the
361 // immutable accepted commit instead of copying every payload again.
362 Events: prepared.events,
363 }
364 for i := range commit.Events {
365 commit.Events[i].Sequence = commit.FirstSequence + uint64(i)
366 }
367 storedCommit := commit
368 storedCommit.Events = prepared.storedEvents
369 for i := range storedCommit.Events {
370 storedCommit.Events[i].Sequence = storedCommit.FirstSequence + uint64(i)
371 }
372 identities := s.messageIDs.changeFor(commit)
373 if identities.duplicate != "" {
374 s.mu.Unlock()
375 return Commit{}, duplicateMessageError(identities.duplicate)
376 }
377 projection := cloneProjection(s.projection)
378 if err := applyProjectionCommit(&projection, commit); err != nil {
379 s.mu.Unlock()
380 return Commit{}, err
381 }
382 binding := s.binding
383 if binding == nil {
384 s.mu.Unlock()
385 return Commit{}, ErrReadOnly
386 }
387 // Session and the persistence queue form one in-memory acceptance boundary.
388 // accept performs no I/O or callbacks; after projection validation there are
389 // no remaining fallible state changes in the closure.
390 err := binding.accept(storedCommit, prepared.reservation, func() {
391 s.commits = append(s.commits, commit)
392 s.projection = projection
393 s.messageIDs.apply(identities)
394 _ = applyRecentCommit(&s.recentMessages, commit)
395 s.next = commit.LastSequence() + 1
396 s.operations[prepared.operationID] = compactOperationRecord(commit)
397 s.acceptTranscriptCommit(commit)
398 })
399 s.mu.Unlock()
400 if err != nil {
401 return Commit{}, err
402 }
403 return cloneCommit(commit), nil
404 }
405
406 // AppendBatch is the convenience form used by callers that do not need to
407 // separate validation from the commit.
408 func (s *Session) AppendBatch(ctx context.Context, operationID string, events []Event) (Commit, error) {
409 return s.Append(ctx, Batch{OperationID: operationID, Events: events})
410 }
411
412 // Append validates and commits a logical batch in one step.
413 func (s *Session) Append(ctx context.Context, batch Batch) (Commit, error) {
414 if err := ctx.Err(); err != nil {
415 return Commit{}, err
416 }
417 prepared, err := s.PrepareBatchContext(ctx, batch.OperationID, batch)
418 if err != nil {
419 return Commit{}, err
420 }
421 return s.CommitPrepared(prepared)
422 }
423
424 // Snapshot returns the full live view, including message and turn history.
425 func (s *Session) Snapshot() Snapshot { return s.snapshot(true, true) }
426
427 // StateSnapshot omits history so progress notifications do not copy every
428 // message and completed turn on each activity update.
429 func (s *Session) StateSnapshot() Snapshot { return s.snapshot(false, false) }
430
431 // ExecutionSnapshot exposes the current model workset and business state,
432 // retaining local refusal evidence until outbound request normalization. It
433 // never materializes durable UI history; UI history is obtained from Query.
434 func (s *Session) ExecutionSnapshot() Snapshot {
435 snapshot := s.snapshot(false, true)
436 restoreRejectedToolResults(&snapshot.Projection)
437 return snapshot
438 }
439
440 func (s *Session) snapshot(includeHistory, includeModel bool) Snapshot {
441 if s == nil {
442 return Snapshot{PersistenceStatus: PersistenceFailed, PersistenceError: "nil session"}
443 }
444 s.mu.Lock()
445 projection := s.projection
446 sequence := s.next - 1
447 externalHistory := s.externalHistory
448 accepted := cloneCommits(s.commits)
449 var history eventPageReader
450 if s.binding != nil {
451 history = s.binding.handle
452 } else if s.coldHandle != nil {
453 history = s.coldHandle
454 }
455 s.mu.Unlock()
456 if includeHistory && externalHistory {
457 // Snapshot is the explicit full-history compatibility boundary. Service
458 // progress and Goal paths use StateSnapshot; paged clients use Query.
459 // Reconstructing here preserves existing callers without keeping a second
460 // durable UI transcript resident in every runtime.
461 if messages, err := materializeSnapshotMessages(history, accepted, sequence); err == nil {
462 projection.Messages = messages
463 }
464 }
465 if !includeHistory {
466 projection.Messages = nil
467 }
468 if !includeModel {
469 projection.ModelMessages, projection.Turns = nil, nil
470 }
471 snapshot := Snapshot{EventSequence: sequence, Projection: cloneProjection(projection)}
472 if s.binding != nil {
473 durable, status, detail := s.binding.progress()
474 snapshot.DurableSequence, snapshot.PersistenceStatus, snapshot.PersistenceError = durable, status, detail
475 } else {
476 snapshot.PersistenceStatus = PersistenceReady
477 }
478 // Nested provider metadata is immutable internally but Go cannot freeze
479 // returned slices. Detach it outside the commit lock before exposing it.
480 snapshot.Projection.Messages = detachMessages(snapshot.Projection.Messages)
481 snapshot.Projection.ModelMessages = detachMessages(snapshot.Projection.ModelMessages)
482 return snapshot
483 }
484
485 // CatalogMetadata returns the rebuildable list projection of this session.
486 func (s *Session) CatalogMetadata() catalogMetadata {
487 if s == nil {
488 return catalogMetadata{Version: catalogMetadataVersion, Codec: Codec}
489 }
490 s.mu.Lock()
491 defer s.mu.Unlock()
492 metadata := metadataFromProjection(s.manifest, s.next-1, s.projection)
493 return metadata
494 }
495
496 // DeriveMessages returns the model history projection.
497 func (s *Session) DeriveMessages() []provider.Message {
498 if s == nil {
499 return nil
500 }
501 s.mu.Lock()
502 messages := detachMessages(s.projection.ModelMessages)
503 s.mu.Unlock()
504 return messages
505 }
506
507 func (s *Session) cacheWeight() int64 {
508 if s == nil {
509 return 0
510 }
511 s.mu.Lock()
512 defer s.mu.Unlock()
513 weight := int64(64 << 10)
514 for id, result := range s.projection.RejectedToolResults {
515 weight += int64(128 + len(id) + len(result.ToolCallID) + len(result.Name) + len(result.State))
516 }
517 for _, message := range s.projection.ModelMessages {
518 weight += int64(len(message.ID) + len(message.Content) + len(message.RawContent) + len(message.ProviderContent) + len(message.ReasoningContent) + len(message.ReasoningSignature) + len(message.Original))
519 for _, image := range message.Images {
520 weight += int64(len(image))
521 }
522 for _, call := range message.ToolCalls {
523 weight += int64(len(call.ID) + len(call.Name) + len(call.Arguments) + len(call.Diff))
524 }
525 for _, item := range message.ResponsesItems {
526 weight += int64(len(item))
527 }
528 for _, block := range message.ThinkingBlocks {
529 encoded, _ := json.Marshal(block)
530 weight += int64(len(encoded))
531 }
532 }
533 weight += int64(len(s.projection.PlanState) + len(s.projection.GoalState))
534 return weight
535 }
536
537 // RecentSnapshot returns the bounded chat baseline without consulting the
538 // history locator or search index.
539 func (s *Session) RecentSnapshot() RecentSnapshot {
540 if s == nil {
541 return RecentSnapshot{}
542 }
543 durable := uint64(0)
544 if s.binding != nil {
545 durable = s.binding.durableSequence()
546 }
547 s.mu.Lock()
548 messages := detachMessages(s.durableRecent)
549 sessionDir := ""
550 if s.binding != nil {
551 sessionDir = s.binding.dir
552 }
553 snapshot := RecentSnapshot{
554 Version: recoveryFormatVersion, SessionID: s.id, StorageGeneration: s.storageGeneration,
555 DurableSequence: durable,
556 Title: s.projection.Title, ModelRef: s.projection.ModelRef, ModelIdentity: s.projection.ModelIdentity,
557 TotalTurns: visibleBoundaryCount(s.projection, false),
558 }
559 s.mu.Unlock()
560 if sessionDir != "" {
561 snapshot.Entries, _ = buildRecentEntries(context.Background(), sessionDir, messages, durable, snapshot.TotalTurns)
562 }
563 return snapshot
564 }
565
566 func materializeSnapshotMessages(history eventPageReader, accepted []Commit, acceptedSequence uint64) ([]provider.Message, error) {
567 projection, _ := Project(nil)
568 var cursor uint64
569 if history != nil {
570 for {
571 startCursor := cursor
572 page, err := history.Read(context.Background(), cursor, 1000)
573 if err != nil {
574 return nil, err
575 }
576 for _, commit := range page.Commits {
577 if commit.LastSequence() > acceptedSequence {
578 break
579 }
580 if err := applyProjectionCommit(&projection, commit); err != nil {
581 return nil, err
582 }
583 // Only UI messages are requested at this compatibility boundary.
584 // Clearing the provider projection after each commit prevents a
585 // second cumulative model-history allocation during reconstruction.
586 projection.ModelMessages = nil
587 cursor = commit.LastSequence()
588 }
589 if !page.Truncated || cursor >= acceptedSequence {
590 break
591 }
592 if cursor <= startCursor {
593 return nil, fmt.Errorf("%w: full snapshot cursor did not advance", ErrDamagedStore)
594 }
595 }
596 }
597 for _, commit := range accepted {
598 if commit.LastSequence() <= cursor || commit.FirstSequence > acceptedSequence {
599 continue
600 }
601 if err := applyProjectionCommit(&projection, commit); err != nil {
602 return nil, err
603 }
604 projection.ModelMessages = nil
605 cursor = commit.LastSequence()
606 }
607 return projection.Messages, nil
608 }
609
610 // externalizeDurableHistory switches a Service-owned runtime to the bounded
611 // history model. It is intentionally not used by the low-level Store API,
612 // whose compatibility callers still request a complete projection.
613 func (s *Session) externalizeDurableHistory() {
614 if s == nil {
615 return
616 }
617 s.mu.Lock()
618 defer s.mu.Unlock()
619 if s.externalHistory {
620 return
621 }
622 for _, message := range s.projection.Messages {
623 if s.catalogPreview == "" {
624 s.catalogPreview = catalogMessagePreview(message)
625 }
626 }
627 s.projection.Messages = nil
628 s.externalHistory = true
629 }
630
631 // AcceptedPage returns the live accepted prefix, including events that have not
632 // crossed a durability checkpoint yet. SessionHandle.Read on the physical layer
633 // continues to expose only durable records.
634 func (s *Session) AcceptedPage(ctx context.Context, offset uint64, limit int) (EventPage, error) {
635 if err := ctx.Err(); err != nil {
636 return EventPage{}, err
637 }
638 if limit == 0 {
639 limit = 100
640 }
641 if limit < 1 || limit > 1000 {
642 return EventPage{}, fmt.Errorf("session: read limit must be 1..1000 commits")
643 }
644 s.mu.Lock()
645 tail := cloneCommits(s.commits)
646 var handle SessionHandle
647 if s.binding != nil {
648 handle = s.binding.handle
649 } else {
650 handle = s.coldHandle
651 }
652 s.mu.Unlock()
653 page := EventPage{Commits: []Commit{}}
654 if handle != nil {
655 var err error
656 page, err = handle.Read(ctx, offset, limit)
657 if err != nil {
658 return EventPage{}, err
659 }
660 if page.Truncated || len(page.Commits) == limit {
661 return page, nil
662 }
663 }
664 for _, commit := range tail {
665 if commit.LastSequence() <= offset {
666 continue
667 }
668 if len(page.Commits) > 0 && commit.LastSequence() <= page.Commits[len(page.Commits)-1].LastSequence() {
669 continue
670 }
671 if len(page.Commits) == limit {
672 page.Truncated = true
673 break
674 }
675 page.Commits = append(page.Commits, cloneCommit(commit))
676 page.Next = commit.LastSequence()
677 }
678 return page, nil
679 }
680
681 // Flush drains the write-behind queue and returns the durable sequence.
682 func (s *Session) Flush(ctx context.Context) (DurableReceipt, error) {
683 if s == nil || s.binding == nil {
684 return DurableReceipt{}, ErrReadOnly
685 }
686 return s.binding.Flush(ctx)
687 }
688
689 // FlushThrough waits only for the captured accepted prefix.
690 func (s *Session) FlushThrough(ctx context.Context, through uint64) (DurableReceipt, error) {
691 if s == nil || s.binding == nil {
692 return DurableReceipt{}, ErrReadOnly
693 }
694 if through > s.EventSequence() {
695 return DurableReceipt{}, fmt.Errorf("session: watermark exceeds accepted sequence")
696 }
697 return s.binding.FlushThrough(ctx, through)
698 }
699
700 // Read exposes the durable prefix through a paged read. Events accepted but not
701 // yet checkpointed are visible through AcceptedPage instead.
702 func (s *Session) Read(ctx context.Context, offset uint64, limit int) (EventPage, error) {
703 handle := s.Handle()
704 if handle == nil {
705 return EventPage{}, os.ErrClosed
706 }
707 return handle.Read(ctx, offset, limit)
708 }
709
710 // Close releases persistence ownership for this session. It is idempotent and
711 // uncancellable in effect: every caller observes the same close result.
712 func (s *Session) Close(ctx context.Context) error { return s.close(ctx) }
713
714 // Sync forces the physical log to stable storage without draining new work.
715 func (s *Session) Sync(ctx context.Context) (DurableReceipt, error) {
716 handle := s.Handle()
717 if handle == nil {
718 return DurableReceipt{}, ErrReadOnly
719 }
720 return handle.Sync(ctx)
721 }
722
723 // newReadSession builds a cold session over a read-only handle. It replays
724 // nothing: cold callers consume the durable prefix through paged Read, which is
725 // what keeps catalog and history queries independent of log length.
726 func newReadSession(handle SessionHandle) *Session {
727 manifest := handle.Manifest()
728 return &Session{
729 id: manifest.SessionID, manifest: manifest, next: 1,
730 operations: map[string]operationRecord{}, readOnly: true,
731 coldHandle: handle,
732 }
733 }
734
735 // Handle exposes the physical persistence handle. Callers must not treat it as
736 // a business-state owner: reads see only the durable prefix.
737 func (s *Session) Handle() SessionHandle {
738 if s == nil {
739 return nil
740 }
741 s.mu.Lock()
742 defer s.mu.Unlock()
743 if s.binding == nil {
744 return s.coldHandle
745 }
746 return s.binding.handle
747 }
748
749 // WritableHandle exposes the leased physical handle for fork, export, and
750 // recovery operations that are defined in terms of durable bytes.
751 func (s *Session) WritableHandle() WritableSessionHandle {
752 if s == nil {
753 return nil
754 }
755 s.mu.Lock()
756 defer s.mu.Unlock()
757 if s.binding == nil {
758 return nil
759 }
760 handle, _ := s.binding.handle.(WritableSessionHandle)
761 return handle
762 }
763
764 // seal stops accepting new commits. It is the admission boundary that Runtime
765 // close establishes before the binding is drained.
766 func (s *Session) seal(err error) {
767 if s == nil {
768 return
769 }
770 s.mu.Lock()
771 s.sealed = true
772 s.sealedError = err
773 binding := s.binding
774 s.mu.Unlock()
775 if binding != nil {
776 binding.stopAccepting()
777 }
778 }
779
780 // close seals the session, drains the binding, and closes the physical handle.
781 // It is uncancellable and idempotent: every caller observes the same result.
782 func (s *Session) close(ctx context.Context) error {
783 if s == nil {
784 return nil
785 }
786 s.seal(osClosedError())
787 s.mu.Lock()
788 binding, cold := s.binding, s.coldHandle
789 s.mu.Unlock()
790 if binding != nil {
791 return binding.Close(ctx)
792 }
793 if cold != nil {
794 // A cold session holds no binding, so its handle is released directly.
795 // Leaving it open would leak the read handle for every catalog scan.
796 return cold.Close(ctx)
797 }
798 return nil
799 }
800
800 lines GO