| 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 |