返回 DeepSeek-Reasonix
business.go
根目录 / internal / transcript / business.go
1 package transcript
2
3 import (
4 "slices"
5
6 "reasonix/internal/event"
7 "reasonix/internal/eventwire"
8 )
9
10 func (p *Projection) PersistenceFailed() {
11 p.mu.Lock()
12 defer p.mu.Unlock()
13 p.runtime.Status = event.TurnRecoveryRequired
14 p.revision++
15 p.publishChangeLocked(Change{Event: &eventwire.Event{Kind: "turn_status", Status: string(event.TurnRecoveryRequired)}})
16 p.revision++
17 p.publishChangeLocked(Change{Event: &eventwire.Event{Kind: "notice", Level: "warn", Code: "transcript_save_failed", Text: "Session output could not be saved. Displayed output is retained; recovery is required."}})
18 }
19
20 func (p *Projection) RestoreRuntime(runtime Runtime, durable uint64) {
21 p.mu.Lock()
22 defer p.mu.Unlock()
23 p.runtime, p.durable = runtime, durable
24 }
25
26 // AcceptBusiness advances coverage for every accepted batch, including batches
27 // that have no visible records. Rows come from the canonical message reader,
28 // never from a wire-event ledger. Mutable stream owners are retained separately
29 // from the bounded settled tail.
30 func (p *Projection) AcceptBusiness(rows []Message, covered uint64, turnID string, rewrite bool, finalMessageID ...string) {
31 p.acceptBusiness(rows, nil, covered, turnID, rewrite, finalMessageID...)
32 }
33
34 // AcceptRetractions invalidates the reading cut while retaining unrelated
35 // active output. The canonical reader supplies the new visible history.
36 func (p *Projection) AcceptRetractions(rows []Message, removed []string, covered uint64, turnID, finalMessageID string) {
37 p.acceptBusiness(rows, removed, covered, turnID, false, finalMessageID)
38 }
39
40 func (p *Projection) acceptBusiness(rows []Message, removed []string, covered uint64, turnID string, rewrite bool, finalMessageID ...string) {
41 p.mu.Lock()
42 defer p.mu.Unlock()
43 if covered <= p.covered {
44 return
45 }
46 if rewrite {
47 p.buffer.Reset()
48 }
49 for _, id := range removed {
50 for index, row := range slices.Backward(p.buffer.messages) {
51 if row.message.MessageID == id {
52 if row.message.Role == "user" {
53 p.buffer.userTurns--
54 }
55 p.buffer.messages = append(p.buffer.messages[:index], p.buffer.messages[index+1:]...)
56 }
57 }
58 delete(p.buffer.byMessageID, id)
59 delete(p.results, id)
60 for key, attempt := range p.attempts {
61 if attempt.MessageID == id {
62 delete(p.attempts, key)
63 }
64 }
65 }
66 reset := rewrite || len(removed) > 0
67 if reset {
68 p.identity.RewriteEpoch++
69 clear(p.snapshots)
70 p.snapshotOrder = nil
71 p.snapshotBytes = 0
72 }
73 if p.buffer.byMessageID == nil {
74 p.buffer.byMessageID = make(map[string]*bufferedMessage)
75 }
76 if p.results == nil {
77 p.results = make(map[string]uint64)
78 }
79 published := make([]Message, 0, len(rows))
80 for _, message := range rows {
81 p.ensureRecordIdentity(&message)
82 // The published change is a copy. Fill turn identity before that copy,
83 // or the follower installs the user row with no turn and leaves live
84 // output above it. The buffer below must see the same value.
85 if message.TurnID == "" {
86 message.TurnID = turnID
87 }
88 published = append(published, message)
89 var row *bufferedMessage
90 for _, existing := range p.buffer.messages {
91 if existing.message.RecordID == message.RecordID {
92 row = existing
93 break
94 }
95 }
96 if row == nil {
97 row = &bufferedMessage{}
98 p.buffer.messages = append(p.buffer.messages, row)
99 if message.Role == "user" {
100 p.buffer.userTurns++
101 }
102 }
103 row.message = message
104 if message.Role == "assistant" {
105 if len(finalMessageID) == 0 {
106 p.runtime.FinalMessageID = message.MessageID
107 }
108 row.content.replace(message.Content)
109 row.reasoning.replace(message.Reasoning)
110 row.message.Content, row.message.Reasoning = "", ""
111 }
112 if message.MessageID != "" && (message.Role == "assistant" || message.Role == "user") {
113 p.buffer.byMessageID[message.MessageID] = row
114 p.results[message.MessageID] = covered
115 }
116 }
117 first := p.covered + 1
118 if len(finalMessageID) > 0 {
119 p.runtime.FinalMessageID = finalMessageID[0]
120 }
121 p.covered = covered
122 p.revision++
123 p.trimSettledLocked()
124 p.publishChangeLocked(Change{FirstSeq: first, Records: published, ResetRequired: reset})
125 }
126
127 func (p *Projection) trimSettledLocked() {
128 const retainedRecords = 96
129 if len(p.buffer.messages) <= retainedRecords {
130 return
131 }
132 keep := make([]*bufferedMessage, 0, retainedRecords+len(p.attempts))
133 for i, row := range p.buffer.messages {
134 active := row.message.Pending
135 for _, attempt := range p.attempts {
136 active = active || attempt.MessageID == row.message.MessageID
137 }
138 if active || i >= len(p.buffer.messages)-retainedRecords {
139 keep = append(keep, row)
140 } else if p.buffer.byMessageID[row.message.MessageID] == row {
141 delete(p.buffer.byMessageID, row.message.MessageID)
142 delete(p.results, row.message.MessageID)
143 }
144 }
145 p.buffer.messages = keep
146 }
147
147 lines GO