返回 DeepSeek-Reasonix
projection.go
根目录 / internal / transcript / projection.go
1 package transcript
2
3 import (
4 "crypto/rand"
5 "encoding/json"
6 "errors"
7 "fmt"
8 "slices"
9 "sort"
10 "sync"
11
12 "reasonix/internal/event"
13 "reasonix/internal/eventwire"
14 "reasonix/internal/turnevent"
15 )
16
17 const ProtocolVersion = 1
18
19 type Identity struct {
20 SessionID string `json:"sessionId"`
21 HeadID string `json:"headId"`
22 RewriteEpoch uint64 `json:"rewriteEpoch"`
23 RuntimeEpoch string `json:"runtimeEpoch"`
24 }
25
26 type ActiveAttempt struct {
27 ID string `json:"id"`
28 MessageID string `json:"messageId"`
29 TurnID string `json:"turnId"`
30 NextIndex uint64 `json:"nextIndex"`
31 }
32
33 // Runtime is reduced from the same ordered events as the visible records.
34 // In particular it is never sampled independently after a history read.
35 type Runtime struct {
36 FinalMessageID string `json:"finalMessageId,omitempty"`
37 DurationMs int64 `json:"durationMs,omitempty"`
38 SamplingCount int `json:"samplingCount"`
39 ToolCount int `json:"toolCount"`
40 TurnID string `json:"turnId,omitempty"`
41 SubmissionID string `json:"submissionId,omitempty"`
42 Status event.TurnStatus `json:"status,omitempty"`
43 Phase string `json:"phase,omitempty"`
44 StartedAt int64 `json:"startedAt,omitempty"`
45 PendingEvents []eventwire.Event `json:"pendingEvents"`
46 CompletionSummary *eventwire.CompletionSummary `json:"completionSummary,omitempty"`
47 TurnUsage *TurnUsage `json:"turnUsage,omitempty"`
48 }
49
50 type Boundary struct {
51 ProtocolVersion int `json:"protocolVersion"`
52 SnapshotID string `json:"snapshotId"`
53 Identity Identity `json:"identity"`
54 ProjectionRevision uint64 `json:"projectionRevision"`
55 CoveredThroughSeq uint64 `json:"coveredThroughSeq"`
56 DurableSeq uint64 `json:"durableSeq"`
57 }
58
59 // Projection has one commit boundary for rows, runtime and coverage. All
60 // mutations happen after durable append and before publishing the event.
61 // It performs no callbacks or I/O under its mutex.
62 type Projection struct {
63 mu sync.Mutex
64 incarnation string
65 identity Identity
66 revision uint64
67 covered uint64
68 durable uint64
69 followers map[string]*follower
70 results map[string]uint64
71 buffer Buffer
72 runtime Runtime
73 startedTurnID string
74 attempts map[string]ActiveAttempt
75 toolCalls map[string]bool
76 prompts map[string]eventwire.Event
77 snapshots map[string]frozenSnapshot
78 snapshotOrder []string
79 snapshotBytes int
80 recordSerial uint64
81 // outline is the complete turn index of a frozen cut. It is built by
82 // freezeLocked and read only from frozen cuts, so it always describes the
83 // same revision as the records paged beside it.
84 outline []OutlineEntry
85 recordPositions map[string]int
86 }
87
88 // ensureRecordIdentity owns the last-resort identity for display-only rows.
89 // Canonical messages retain their existing m:/tool: identities; transient
90 // frames without a business sequence receive an identity scoped to this
91 // projection incarnation and keep it for every later snapshot.
92 func (p *Projection) ensureRecordIdentity(message *Message) {
93 if message.RecordID != "" {
94 return
95 }
96 switch {
97 case message.Role == "tool" && message.ToolCallID != "":
98 message.RecordID = "tool:" + message.ToolCallID
99 case message.MessageID != "":
100 message.RecordID = "m:" + message.MessageID
101 default:
102 p.recordSerial++
103 message.RecordID = fmt.Sprintf("view:%s:%d", p.incarnation, p.recordSerial)
104 }
105 }
106
107 func (p *Projection) ensureBufferRecordIdentities() {
108 for _, row := range p.buffer.messages {
109 p.ensureRecordIdentity(&row.message)
110 }
111 }
112
113 func NewProjection(identity Identity, baseline []Message, covered uint64) (*Projection, error) {
114 p := &Projection{incarnation: rand.Text(), identity: identity, covered: covered, revision: 1,
115 attempts: make(map[string]ActiveAttempt), prompts: make(map[string]eventwire.Event)}
116 // Take ownership of nested metadata as well as the slice. Callers may
117 // reuse their conversion buffers immediately after construction.
118 encoded, err := json.Marshal(baseline)
119 if err != nil {
120 return nil, newBaselineError(err, "baseline_encode_failed", len(baseline), -1, -1, Message{})
121 }
122 var owned []Message
123 if err = json.Unmarshal(encoded, &owned); err != nil {
124 return nil, newBaselineError(err, "baseline_decode_failed", len(baseline), -1, -1, Message{})
125 }
126 p.buffer.byMessageID = make(map[string]*bufferedMessage)
127 seen := make(map[string]int)
128 for index, m := range owned {
129 if m.Role == "user" {
130 p.buffer.userTurns++
131 if m.HistoryTurn == 0 {
132 m.HistoryTurn = p.buffer.userTurns
133 }
134 }
135 if m.RecordID == "" {
136 switch {
137 case m.Role == "tool" && m.ToolCallID != "":
138 m.RecordID = "tool:" + m.ToolCallID
139 case m.MessageID != "":
140 m.RecordID = "m:" + m.MessageID
141 default:
142 return nil, newBaselineError(errors.New("transcript baseline has a record without identity"), "missing_record_identity", len(owned), index, -1, m)
143 }
144 }
145 if previous, exists := seen[m.RecordID]; exists {
146 return nil, newBaselineError(fmt.Errorf("duplicate transcript record %q", m.RecordID), "duplicate_record_identity", len(owned), index, previous, m)
147 }
148 seen[m.RecordID] = index
149 row := &bufferedMessage{message: m}
150 if m.Role == "assistant" {
151 row.content.replace(m.Content)
152 row.reasoning.replace(m.Reasoning)
153 row.message.Content, row.message.Reasoning = "", ""
154 }
155 p.buffer.messages = append(p.buffer.messages, row)
156 if m.MessageID != "" && (m.Role == "assistant" || m.Role == "user") {
157 p.buffer.byMessageID[m.MessageID] = row
158 }
159 }
160 return p, nil
161 }
162
163 func (p *Projection) Apply(envelope turnevent.Envelope) error {
164 p.mu.Lock()
165 defer p.mu.Unlock()
166 if envelope.SessionID != p.identity.SessionID || (envelope.RuntimeEpoch != "" && envelope.RuntimeEpoch != p.identity.RuntimeEpoch) {
167 return errors.New("transcript event identity mismatch")
168 }
169 if envelope.Sequence <= p.covered {
170 return nil
171 }
172 if envelope.Sequence != p.covered+1 {
173 return errors.New("transcript projection sequence gap")
174 }
175 return p.applyLocked(envelope, envelope.Sequence, envelope.Sequence)
176 }
177
178 // ApplyFrame uses an independent display revision. The caller supplies the
179 // business cut; a token, phase or usage notification cannot allocate log sequence.
180 func (p *Projection) ApplyFrame(envelope turnevent.Envelope, covered uint64) error {
181 p.mu.Lock()
182 defer p.mu.Unlock()
183 if (envelope.SessionID != "" && envelope.SessionID != p.identity.SessionID) || (envelope.RuntimeEpoch != "" && envelope.RuntimeEpoch != p.identity.RuntimeEpoch) {
184 return errors.New("transcript event identity mismatch")
185 }
186 if covered < p.covered {
187 return errors.New("transcript business coverage regression")
188 }
189 if err := p.applyLocked(envelope, covered, 0); err != nil {
190 return err
191 }
192 p.trimSettledLocked()
193 return nil
194 }
195
196 // first is the business sequence the envelope itself occupies, or zero for a
197 // display frame that rides an existing cut.
198 func (p *Projection) applyLocked(envelope turnevent.Envelope, covered, first uint64) error {
199 // Detach pointer payloads before retaining them. Event publication cannot
200 // mutate a previously committed snapshot through an aliased tool slice.
201 encoded, err := json.Marshal(envelope)
202 if err != nil {
203 return err
204 }
205 var owned turnevent.Envelope
206 if err = json.Unmarshal(encoded, &owned); err != nil {
207 return err
208 }
209 if e, ok := EventFromEnvelope(owned); ok {
210 p.buffer.Apply(e)
211 p.ensureBufferRecordIdentities()
212 if e.Kind == event.TurnDone {
213 p.applyTerminalNotices(e)
214 }
215 if m := p.buffer.byMessageID[e.MessageID]; m != nil {
216 if m.message.CreatedAt == 0 {
217 m.message.CreatedAt = owned.CreatedAt
218 }
219 if owned.SubmissionID != "" {
220 m.message.SubmissionID = owned.SubmissionID
221 }
222 }
223 }
224 p.applyRuntimeLocked(owned)
225 w := owned.Event
226 p.covered = covered
227 p.revision++
228 // Legacy ledger numbering is never a chat coverage cursor.
229 w.Sequence = 0
230 if first != 0 && owned.Kind == "message" && w.MessageID != "" {
231 if p.results == nil {
232 p.results = make(map[string]uint64)
233 }
234 p.results[w.MessageID] = first
235 }
236 state := p.runtime
237 change := Change{Event: &w, Runtime: &state, FirstSeq: first}
238 if owned.Kind == "text" || owned.Kind == "reasoning" || owned.Kind == "tool_call_delta" || (owned.Kind == "tool_dispatch" && w.Tool != nil && w.Tool.Partial) {
239 if attempt, ok := p.attempts[w.AttemptID]; ok {
240 change.AttemptID, change.Index = attempt.ID, attempt.NextIndex
241 attempt.NextIndex++
242 p.attempts[attempt.ID] = attempt
243 }
244 }
245 if owned.Kind == "stream_attempt" && w.StreamAttempt != nil && w.StreamAttempt.Action == "commit" {
246 change.AttemptID, change.ResultSeq = w.StreamAttempt.ID, p.results[w.MessageID]
247 if first != 0 {
248 // A tool-only or interrupted round publishes no message event; its
249 // commit follows the durable write, so its own sequence is committed.
250 if change.ResultSeq == 0 {
251 change.ResultSeq = first
252 }
253 delete(p.results, w.MessageID)
254 }
255 change.ResultKind = "message/complete"
256 if w.StreamAttempt.Reason == "interrupted" {
257 change.ResultKind = "message/interrupted"
258 }
259 change.ResetRequired = change.ResultSeq == 0
260 }
261 p.publishChangeLocked(change)
262 return nil
263 }
264
265 func mergeTurnUsage(current *TurnUsage, usage *eventwire.Usage) *TurnUsage {
266 if usage == nil {
267 return current
268 }
269 if current == nil {
270 zero := 0
271 current = &TurnUsage{CacheReadTokens: &zero, ReasoningTokens: &zero}
272 }
273 current.UncachedInputTokens += usage.CacheMissTokens
274 if usage.CacheMissTokens == 0 && usage.CacheHitTokens == 0 {
275 current.UncachedInputTokens += usage.PromptTokens
276 }
277 current.OutputTokens += usage.CompletionTokens
278 requestTotal := usage.TotalTokens
279 if requestTotal <= 0 {
280 requestTotal = usage.PromptTokens + usage.CompletionTokens
281 }
282 current.TotalTokens += requestTotal
283 cacheRead := valueOrZero(current.CacheReadTokens) + usage.CacheHitTokens
284 current.CacheReadTokens = &cacheRead
285 reasoning := valueOrZero(current.ReasoningTokens) + usage.ReasoningTokens
286 current.ReasoningTokens = &reasoning
287 if usage.CostQuote != nil && usage.CostQuote.ModelRef != "" {
288 if !slices.Contains(current.Routes, usage.CostQuote.ModelRef) {
289 current.Routes = append(current.Routes, usage.CostQuote.ModelRef)
290 }
291 }
292 return current
293 }
294
295 func valueOrZero(value *int) int {
296 if value == nil {
297 return 0
298 }
299 return *value
300 }
301
302 // SetRuntimeEpoch is called only by the controller's idle routing boundary.
303 // Existing rows survive a runtime rebind; outstanding snapshot leases do not.
304 func (p *Projection) SetRuntimeEpoch(epoch string) {
305 p.mu.Lock()
306 defer p.mu.Unlock()
307 if p.identity.RuntimeEpoch != epoch {
308 p.identity.RuntimeEpoch = epoch
309 p.incarnation = rand.Text()
310 p.revision++
311 for _, f := range p.followers {
312 f.reset = true
313 select {
314 case f.wake <- struct{}{}:
315 default:
316 }
317 }
318 clear(p.followers)
319 }
320 }
321
322 func (p *Projection) boundaryLocked() Boundary {
323 return Boundary{ProtocolVersion: ProtocolVersion,
324 SnapshotID: fmt.Sprintf("%s:%d", p.incarnation, p.revision),
325 Identity: p.identity, ProjectionRevision: p.revision, CoveredThroughSeq: p.covered, DurableSeq: p.durable}
326 }
327
328 func (p *Projection) Boundary() Boundary {
329 p.mu.Lock()
330 defer p.mu.Unlock()
331 return p.boundaryLocked()
332 }
333
334 func (p *Projection) runtimeLocked() (Runtime, []ActiveAttempt) {
335 runtime := p.runtime
336 runtime.PendingEvents = make([]eventwire.Event, 0, len(p.prompts))
337 keys := make([]string, 0, len(p.prompts))
338 for key := range p.prompts {
339 keys = append(keys, key)
340 }
341 sort.Strings(keys)
342 for _, key := range keys {
343 runtime.PendingEvents = append(runtime.PendingEvents, p.prompts[key])
344 }
345 attempts := make([]ActiveAttempt, 0, len(p.attempts))
346 for _, attempt := range p.attempts {
347 attempts = append(attempts, attempt)
348 }
349 sort.Slice(attempts, func(i, j int) bool { return attempts[i].ID < attempts[j].ID })
350 return runtime, attempts
351 }
352
353 func (p *Projection) applyRuntimeLocked(owned turnevent.Envelope) {
354 w := owned.Event
355 if owned.TurnID != "" {
356 p.runtime.TurnID, p.runtime.Status = owned.TurnID, owned.Status
357 p.runtime.SubmissionID = owned.SubmissionID
358 }
359 switch owned.Kind {
360 case "turn_started":
361 p.runtime.FinalMessageID, p.runtime.DurationMs = "", 0
362 p.runtime.SamplingCount, p.runtime.ToolCount = 0, 0
363 p.toolCalls = make(map[string]bool)
364 p.retireRecoveryNotices()
365 if p.startedTurnID != owned.TurnID || p.runtime.StartedAt == 0 {
366 p.runtime.StartedAt = owned.CreatedAt
367 p.startedTurnID = owned.TurnID
368 }
369 p.runtime.Phase = ""
370 p.runtime.CompletionSummary = nil
371 p.runtime.TurnUsage = nil
372 p.buffer.completion = nil
373 case "usage":
374 p.runtime.TurnUsage = mergeTurnUsage(p.runtime.TurnUsage, w.Usage)
375 case "turn_phase":
376 p.runtime.Phase = w.Phase
377 case "completion_summary":
378 p.runtime.CompletionSummary = w.Completion
379 case "stream_attempt":
380 if w.StreamAttempt != nil {
381 if w.StreamAttempt.Action == "begin" {
382 p.runtime.SamplingCount++
383 p.attempts[w.StreamAttempt.ID] = ActiveAttempt{ID: w.StreamAttempt.ID, MessageID: w.MessageID, TurnID: owned.TurnID}
384 } else {
385 delete(p.attempts, w.StreamAttempt.ID)
386 }
387 }
388 case "tool_dispatch":
389 if w.Tool != nil && w.Tool.ID != "" && !p.toolCalls[w.Tool.ID] {
390 if p.toolCalls == nil {
391 p.toolCalls = make(map[string]bool)
392 }
393 p.toolCalls[w.Tool.ID] = true
394 p.runtime.ToolCount++
395 }
396 case "ask_request", "approval_request", "mcp_interaction":
397 id := w.PromptID
398 if id == "" {
399 id = owned.ItemID
400 }
401 if id != "" {
402 p.prompts[id] = w
403 }
404 case "prompt_answered":
405 delete(p.prompts, owned.ItemID)
406 case "turn_done":
407 durationMs := int64(0)
408 if p.runtime.StartedAt > 0 && owned.CreatedAt >= p.runtime.StartedAt {
409 durationMs = owned.CreatedAt - p.runtime.StartedAt
410 }
411 p.runtime.DurationMs = durationMs
412 p.buffer.attachTurnStats(owned.TurnID, p.runtime.TurnUsage, durationMs, owned.CreatedAt, p.runtime.FinalMessageID)
413 clear(p.prompts)
414 clear(p.attempts)
415 if owned.TranscriptDigest != "" {
416 p.identity.HeadID = owned.HeadID
417 p.identity.RewriteEpoch = owned.RewriteEpoch
418 }
419 }
420 }
421
421 lines GO