返回 DeepSeek-Reasonix
stream_checkpoint.go
根目录 / internal / control / stream_checkpoint.go
1 package control
2
3 import (
4 "context"
5 "fmt"
6 "strings"
7 "sync"
8
9 "reasonix/internal/agent"
10 "reasonix/internal/event"
11 "reasonix/internal/provider"
12 "reasonix/internal/session"
13 )
14
15 // openStreamOutput accumulates the deltas of the stream whose message has not
16 // committed yet. The streamed text lives nowhere durable until that commit, so
17 // this is what the autosave tick can checkpoint before a kill.
18 type openStreamOutput struct {
19 mu sync.Mutex
20 messageID string
21 text strings.Builder
22 reasoning strings.Builder
23 // checkpointed is the text+reasoning length the store already holds.
24 checkpointed int
25 }
26
27 func (o *openStreamOutput) observe(e event.Event) {
28 if e.Kind == event.TurnStarted {
29 o.settle("")
30 return
31 }
32 if (e.Kind != event.Text && e.Kind != event.Reasoning) || e.MessageID == "" || e.Text == "" {
33 return
34 }
35 o.mu.Lock()
36 defer o.mu.Unlock()
37 if e.MessageID != o.messageID {
38 o.resetLocked(e.MessageID)
39 }
40 if e.Kind == event.Text {
41 o.text.WriteString(e.Text)
42 } else {
43 o.reasoning.WriteString(e.Text)
44 }
45 }
46
47 // settle forgets the stream once its message is committed; an empty id
48 // forgets whatever stream is open. Any other commit leaves the stream open but
49 // supersedes its stored checkpoint, so the next tick writes it again.
50 func (o *openStreamOutput) settle(committedID string) {
51 o.mu.Lock()
52 defer o.mu.Unlock()
53 if committedID == "" || committedID == o.messageID {
54 o.resetLocked("")
55 return
56 }
57 o.checkpointed = 0
58 }
59
60 func (o *openStreamOutput) resetLocked(messageID string) {
61 o.messageID = messageID
62 o.text.Reset()
63 o.reasoning.Reset()
64 o.checkpointed = 0
65 }
66
67 func (o *openStreamOutput) pending() (provider.Message, int, bool) {
68 o.mu.Lock()
69 defer o.mu.Unlock()
70 size := o.text.Len() + o.reasoning.Len()
71 if o.messageID == "" || size == 0 || size == o.checkpointed {
72 return provider.Message{}, 0, false
73 }
74 return agent.InterruptedStreamRecord(o.messageID, o.text.String(), o.reasoning.String()), size, true
75 }
76
77 func (o *openStreamOutput) markCheckpointed(messageID string, size int) {
78 o.mu.Lock()
79 defer o.mu.Unlock()
80 if o.messageID == messageID {
81 o.checkpointed = size
82 }
83 }
84
85 // settleOpenStreamLocked runs under commitMu with the messages just committed.
86 func (c *Controller) settleOpenStreamLocked(messages []provider.Message) {
87 for _, m := range messages {
88 c.turnEvents.openStream.settle(m.ID)
89 }
90 }
91
92 // checkpointOpenStream records the open stream's output as the local-only
93 // record its interruption would leave. It holds commitMu so the checkpoint can
94 // never land after the commit of the message it describes.
95 func (c *Controller) checkpointOpenStream(ctx context.Context) error {
96 store := c.sessionEventStore()
97 if store == nil || !c.sessionEventCommitAllowed() {
98 return nil
99 }
100 c.turnEvents.commitMu.Lock()
101 defer c.turnEvents.commitMu.Unlock()
102 if !c.messageCommitAllowedLocked(ctx, store) {
103 return nil
104 }
105 record, size, ok := c.turnEvents.openStream.pending()
106 if !ok || c.turnEvents.turnMessageIDs[record.ID] {
107 return nil
108 }
109 turnID := store.StateSnapshot().Projection.TurnID
110 if turnID == "" {
111 return nil
112 }
113 checkpoint, err := session.StreamCheckpointEvent(record)
114 if err != nil {
115 return err
116 }
117 op := fmt.Sprintf("stream-checkpoint:%s:%d", record.ID, size)
118 if _, err := c.appendSessionBatch(ctx, store, session.Batch{OperationID: op, TurnID: turnID, Events: []session.Event{checkpoint}}); err != nil {
119 return err
120 }
121 c.turnEvents.openStream.markCheckpointed(record.ID, size)
122 return nil
123 }
124
124 lines GO