返回 DeepSeek-Reasonix
transcript_api.go
根目录 / internal / control / transcript_api.go
1 package control
2
3 import (
4 "context"
5 "encoding/json"
6 "errors"
7 "reasonix/internal/session"
8 "time"
9
10 "reasonix/internal/transcript"
11 "reasonix/internal/turnevent"
12 )
13
14 var ErrTranscriptProjectionUnavailable = errors.New("transcript projection is unavailable")
15
16 type TranscriptReplayRequest struct {
17 Identity transcript.Identity `json:"identity"`
18 After uint64 `json:"after"`
19 }
20
21 type TranscriptReplay struct {
22 transcript.Boundary
23 turnevent.ReplayView
24 }
25
26 type TranscriptProjectionAPI interface {
27 TranscriptSnapshot(transcript.PageRequest) (transcript.Snapshot, error)
28 TranscriptContent(transcript.ContentRequest) (transcript.ContentChunk, error)
29 TranscriptReplay(TranscriptReplayRequest) (TranscriptReplay, error)
30 }
31
32 var _ TranscriptProjectionAPI = (*Controller)(nil)
33
34 type TranscriptFollowAPI interface {
35 TranscriptFollow(context.Context, transcript.FollowRequest) (TranscriptFollowResponse, error)
36 }
37
38 type TranscriptFollowResponse struct {
39 transcript.FollowResponse
40 History *session.HistoryWindowPage `json:"history,omitempty"`
41 StorageBackend string `json:"storageBackend,omitempty"`
42 }
43
44 func (c *Controller) TranscriptFollow(ctx context.Context, req transcript.FollowRequest) (TranscriptFollowResponse, error) {
45 service, runtime, exclusive := c.v3Binding()
46 if !exclusive || runtime == nil {
47 projection, err := c.transcriptProjection()
48 if err != nil {
49 return TranscriptFollowResponse{}, err
50 }
51 view, err := projection.Follow(ctx, req)
52 return TranscriptFollowResponse{FollowResponse: view, StorageBackend: "legacy"}, err
53 }
54 view, err := runtime.FollowTranscript(ctx, req)
55 out := TranscriptFollowResponse{FollowResponse: view}
56 if err != nil || view.Snapshot == nil {
57 return out, err
58 }
59 // Make the frozen cut pageable even when accepted batches exceed the tail.
60 // The registered subscription queues concurrent commits during flush/read;
61 // no publisher lock is held.
62 cut := view.Snapshot.CoveredThroughSeq
63 if view.Snapshot.DurableSeq < cut {
64 receipt, flushErr := runtime.Session().Flush(ctx)
65 if flushErr != nil || receipt.DurableSequence < cut {
66 _, _ = runtime.Transcript().Follow(context.Background(), transcript.FollowRequest{Subscription: view.Subscription, Close: true})
67 if flushErr == nil {
68 flushErr = errors.New("transcript snapshot persistence is incomplete")
69 }
70 return out, flushErr
71 }
72 // Publish the new watermark in order, after already queued frames. Do
73 // not attach it to the older frozen view and then replay older watermarks.
74 runtime.Transcript().SetDurableSequence(receipt.DurableSequence)
75 }
76 var page session.HistoryWindowPage
77 for {
78 page, err = service.Query().ReadHistoryWindow(ctx, runtime.Ref(), session.HistoryWindowRequest{Anchor: "newest", Limit: 32, SnapshotSequence: &cut})
79 if err != nil || page.Status != "preparing" {
80 break
81 }
82 timer := time.NewTimer(10 * time.Millisecond)
83 select {
84 case <-ctx.Done():
85 timer.Stop()
86 err = ctx.Err()
87 case <-timer.C:
88 }
89 if err != nil {
90 break
91 }
92 }
93 if err != nil {
94 _, _ = runtime.Transcript().Follow(context.Background(), transcript.FollowRequest{Subscription: view.Subscription, Close: true})
95 }
96 out.History = &page
97 return out, err
98 }
99
100 // TranscriptOutlineAPI is an optional capability beside TranscriptProjectionAPI.
101 // It is deliberately separate so an existing controller implementation keeps
102 // compiling and a client can negotiate the outline independently of the body.
103 type TranscriptOutlineAPI interface {
104 TranscriptOutline(transcript.OutlineRequest) (transcript.OutlinePage, error)
105 }
106
107 var _ TranscriptOutlineAPI = (*Controller)(nil)
108
109 // SetTurnSubmissionID is called under the transport's admission boundary.
110 func (c *Controller) SetTurnSubmissionID(submissionID string) {
111 if ledger := c.turnEventLedger(); ledger != nil {
112 ledger.SetSubmissionID(submissionID)
113 }
114 }
115
116 // BindTranscriptRuntimeEpoch runs at the surface's idle runtime publication
117 // boundary. It takes only display/ledger leaf locks and invokes no callbacks.
118 func (c *Controller) BindTranscriptRuntimeEpoch(epoch string) {
119 c.turnEvents.commitMu.Lock()
120 defer c.turnEvents.commitMu.Unlock()
121 ledger := c.turnEventLedger()
122 if ledger == nil || ledger.ActiveTurnID() != "" {
123 return
124 }
125 ledger.SetRuntimeEpoch(epoch)
126 c.turnEvents.mu.RLock()
127 p := c.turnEvents.projection
128 c.turnEvents.mu.RUnlock()
129 if p != nil {
130 p.SetRuntimeEpoch(epoch)
131 }
132 }
133
134 func (c *Controller) transcriptProjection() (*transcript.Projection, error) {
135 if _, runtime, exclusive := c.v3Binding(); exclusive && runtime != nil {
136 return runtime.Transcript(), nil
137 }
138 c.turnEvents.mu.RLock()
139 defer c.turnEvents.mu.RUnlock()
140 if c.turnEvents.err != nil {
141 return nil, errors.Join(ErrTranscriptProjectionUnavailable, c.turnEvents.err)
142 }
143 if c.turnEvents.projectionErr != nil {
144 return nil, errors.Join(ErrTranscriptProjectionUnavailable, c.turnEvents.projectionErr)
145 }
146 if c.turnEvents.projection == nil {
147 return nil, ErrTranscriptProjectionUnavailable
148 }
149 return c.turnEvents.projection, nil
150 }
151
152 func (c *Controller) TranscriptSnapshot(req transcript.PageRequest) (transcript.Snapshot, error) {
153 p, err := c.transcriptProjection()
154 if err != nil {
155 return transcript.Snapshot{}, err
156 }
157 return p.Snapshot(req)
158 }
159
160 // TranscriptOutline pages the complete turn index of one snapshot. It reads the
161 // same projection the body pages do, so both describe one immutable cut.
162 func (c *Controller) TranscriptOutline(req transcript.OutlineRequest) (transcript.OutlinePage, error) {
163 p, err := c.transcriptProjection()
164 if err != nil {
165 return transcript.OutlinePage{}, err
166 }
167 return p.Outline(req)
168 }
169
170 func (c *Controller) TranscriptContent(req transcript.ContentRequest) (transcript.ContentChunk, error) {
171 p, err := c.transcriptProjection()
172 if err != nil {
173 return transcript.ContentChunk{}, err
174 }
175 return p.Content(req)
176 }
177
178 func (c *Controller) TranscriptReplay(req TranscriptReplayRequest) (TranscriptReplay, error) {
179 replay, err := c.transcriptReplay(req)
180 if err != nil {
181 return replay, err
182 }
183 // The ledger's soft budget permits an oversized first event for progress.
184 // Modern clients can obtain that data from the bounded snapshot/content
185 // protocol instead. Never acknowledge a suffix the client cannot receive.
186 encoded, err := json.Marshal(replay)
187 if err != nil {
188 return TranscriptReplay{}, err
189 }
190 if len(encoded)+1 > transcript.MaxResponseBytes {
191 replay.Events = []turnevent.Envelope{}
192 replay.ResetRequired, replay.HasMore = true, false
193 replay.NextAfterSequence = req.After
194 }
195 return replay, nil
196 }
197
198 func (c *Controller) transcriptReplay(req TranscriptReplayRequest) (TranscriptReplay, error) {
199 if _, runtime, exclusive := c.v3Binding(); exclusive && runtime != nil {
200 return TranscriptReplay{}, errors.New("transcript v2 requires Follow; legacy replay is unavailable")
201 }
202 c.turnEvents.commitMu.Lock()
203 defer c.turnEvents.commitMu.Unlock()
204 p, err := c.transcriptProjection()
205 if err != nil {
206 return TranscriptReplay{}, err
207 }
208 boundary := p.Boundary()
209 if boundary.Identity != req.Identity {
210 if boundary.Identity.SessionID != req.Identity.SessionID {
211 return TranscriptReplay{}, errors.New("transcript replay session mismatch")
212 }
213 return TranscriptReplay{Boundary: boundary, ReplayView: turnevent.ReplayView{Events: []turnevent.Envelope{}, ResetRequired: true, LatestSequence: boundary.CoveredThroughSeq}}, nil
214 }
215 view, err := c.TurnEventReplay(req.After)
216 return TranscriptReplay{Boundary: boundary, ReplayView: view}, err
217 }
218
218 lines GO