返回 DeepSeek-Reasonix
broadcaster.go
根目录 / internal / serve / broadcaster.go
1 package serve
2
3 import (
4 "encoding/json"
5 "strings"
6 "sync"
7 "time"
8
9 "reasonix/internal/agent"
10 "reasonix/internal/billing"
11 "reasonix/internal/event"
12 "reasonix/internal/eventwire"
13 )
14
15 type subscription struct {
16 all bool
17 }
18
19 const (
20 subscriberBufferSize = 128
21 subscriberPriorityReserve = 32
22 )
23
24 // Broadcaster is the event.Sink the controllers emit to in server mode. It
25 // marshals each event once and fans it out to every connected SSE subscriber.
26 // A slow subscriber's buffer is allowed to drop rather than back-pressure the
27 // agent goroutine — a browser that can't keep up loses intermediate frames, not
28 // the whole session (it can refetch /history).
29 type Broadcaster struct {
30 mu sync.Mutex
31 subs map[chan []byte]subscription
32 ledgers map[string]*billing.Ledger
33 current string
34 displayCurrency string
35 modelApplicationChanged func()
36 }
37
38 // NewBroadcaster returns an empty Broadcaster ready to accept subscribers.
39 func NewBroadcaster() *Broadcaster {
40 return &Broadcaster{
41 subs: map[chan []byte]subscription{},
42 ledgers: map[string]*billing.Ledger{},
43 }
44 }
45
46 // SetDisplayCurrency rebinds the session ledger to a stored valuation. Empty
47 // keeps automatic mode: a single original currency is selected and mixed
48 // currencies remain buckets.
49 func (b *Broadcaster) SetDisplayCurrency(currency string) {
50 if b == nil {
51 return
52 }
53 b.mu.Lock()
54 b.displayCurrency = billing.NormalizeCurrency(currency)
55 b.mu.Unlock()
56 }
57
58 // sessionRouteKey is the one normalization rule for session references the
59 // broadcaster keys on: final-format identity routes ("session-id:<id>") are not
60 // filesystem paths and stay verbatim, legacy transcript paths take their
61 // canonical form. Every emit, ledger and current-session update goes through
62 // it so an identity route can never be rewritten into a cwd-relative pseudo
63 // path that no subscriber or registry matches.
64 func sessionRouteKey(path string) string {
65 path = strings.TrimSpace(path)
66 if strings.HasPrefix(path, remoteSessionIDQueryPrefix) {
67 return path
68 }
69 return agent.CanonicalSessionPath(path)
70 }
71
72 // hiddenFromCurrentOnly reports whether a routed frame is filtered out for
73 // subscribers that watch only the current session. The filter compares path
74 // keys; identity routes are not path-keyed and are routed by the client from
75 // the frame itself, exactly like identity-tagged live frames (which carry no
76 // path), so they are never dropped here.
77 func hiddenFromCurrentOnly(sessionPath, current string) bool {
78 if sessionPath == "" || strings.HasPrefix(sessionPath, remoteSessionIDQueryPrefix) {
79 return false
80 }
81 return sessionPath != current
82 }
83
84 // SetCurrentSession records the controller shown by current-only subscribers.
85 // Untagged events remain compatible and are attributed to this session.
86 func (b *Broadcaster) SetCurrentSession(path string) {
87 if b == nil {
88 return
89 }
90 path = sessionRouteKey(path)
91 b.mu.Lock()
92 b.current = path
93 b.mu.Unlock()
94 }
95
96 // CurrentSession reports the session currently selected by Serve.
97 func (b *Broadcaster) CurrentSession() string {
98 if b == nil {
99 return ""
100 }
101 b.mu.Lock()
102 defer b.mu.Unlock()
103 return b.current
104 }
105
106 // ResetSession clears the current session usage ledger for legacy callers.
107 func (b *Broadcaster) ResetSession() {
108 b.ResetSessionPath("")
109 }
110
111 // ResetSessionPath clears one session ledger without affecting detached
112 // sessions. Empty selects the current session.
113 func (b *Broadcaster) ResetSessionPath(path string) {
114 if b == nil {
115 return
116 }
117 b.mu.Lock()
118 if path == "" {
119 path = b.current
120 } else {
121 path = sessionRouteKey(path)
122 }
123 delete(b.ledgers, path)
124 b.mu.Unlock()
125 }
126
127 // SessionCostQuote returns the current aggregate quote without repricing.
128 func (b *Broadcaster) SessionCostQuote() billing.CostQuote {
129 return b.SessionCostQuoteFor("")
130 }
131
132 // SessionCostQuoteFor returns one session's aggregate quote. Empty selects
133 // the current session so existing single-session callers keep their contract.
134 func (b *Broadcaster) SessionCostQuoteFor(path string) billing.CostQuote {
135 if b == nil {
136 return billing.AggregateQuotes(nil, "")
137 }
138 b.mu.Lock()
139 defer b.mu.Unlock()
140 return b.ledgerLocked(path).Total(b.displayCurrency)
141 }
142
143 func (b *Broadcaster) ledgerLocked(path string) *billing.Ledger {
144 if path == "" {
145 path = b.current
146 } else {
147 path = sessionRouteKey(path)
148 }
149 ledger := b.ledgers[path]
150 if ledger == nil {
151 ledger = billing.NewLedger()
152 b.ledgers[path] = ledger
153 }
154 return ledger
155 }
156
157 // Emit marshals the event to JSON and delivers it to every subscriber. Drops to
158 // a subscriber whose buffer is full rather than blocking. A marshal failure is
159 // dropped silently — one bad event shouldn't stall the stream.
160 func (b *Broadcaster) Emit(e event.Event) {
161 if e.SessionPath != "" {
162 e.SessionPath = sessionRouteKey(e.SessionPath)
163 }
164 wired := eventwire.ToWire(e)
165 b.mu.Lock()
166 observedCurrent := b.current
167 b.mu.Unlock()
168 wired.SessionCurrent = e.SessionPath != "" && e.SessionPath == observedCurrent
169 data, err := json.Marshal(wired)
170 if err != nil {
171 return
172 }
173 b.mu.Lock()
174 defer b.mu.Unlock()
175 if b.current != observedCurrent {
176 wired.SessionCurrent = e.SessionPath != "" && e.SessionPath == b.current
177 data, err = json.Marshal(wired)
178 if err != nil {
179 return
180 }
181 }
182 if e.Kind == event.Usage && e.Usage != nil && e.CostQuote != nil {
183 b.ledgerLocked(e.SessionPath).Add(*e.CostQuote, billing.UsageTokens{
184 PromptTokens: e.Usage.PromptTokens, CompletionTokens: e.Usage.CompletionTokens,
185 CacheHitTokens: e.Usage.CacheHitTokens, CacheMissTokens: e.Usage.CacheMissTokens,
186 CacheWriteTokens: e.Usage.CacheWriteTokens, CacheWriteBilledTokens: e.Usage.CacheWriteBilledTokens,
187 Estimated: e.Usage.Estimated,
188 }, time.Now().UTC())
189 }
190 for ch, sub := range b.subs {
191 if !sub.all && hiddenFromCurrentOnly(e.SessionPath, b.current) {
192 continue
193 }
194 enqueueSubscriberFrame(ch, data, e.Kind)
195 }
196 }
197
198 // EmitTo delivers an event only to the supplied subscriber. It is used for
199 // connection-local recovery frames, such as replaying a prompt to a browser
200 // that attached after the original event was emitted. Normal runtime events
201 // should continue to use Emit so every subscriber receives them.
202 func (b *Broadcaster) EmitTo(target <-chan []byte, e event.Event) {
203 if e.SessionPath != "" {
204 e.SessionPath = sessionRouteKey(e.SessionPath)
205 }
206 wired := eventwire.ToWire(e)
207 b.mu.Lock()
208 observedCurrent := b.current
209 b.mu.Unlock()
210 wired.SessionCurrent = e.SessionPath != "" && e.SessionPath == observedCurrent
211 data, err := json.Marshal(wired)
212 if err != nil {
213 return
214 }
215 b.mu.Lock()
216 defer b.mu.Unlock()
217 if b.current != observedCurrent {
218 wired.SessionCurrent = e.SessionPath != "" && e.SessionPath == b.current
219 data, err = json.Marshal(wired)
220 if err != nil {
221 return
222 }
223 }
224 for ch, sub := range b.subs {
225 if (<-chan []byte)(ch) != target {
226 continue
227 }
228 if !sub.all && hiddenFromCurrentOnly(e.SessionPath, b.current) {
229 return
230 }
231 enqueueSubscriberFrame(ch, data, e.Kind)
232 return
233 }
234 }
235
236 // eventIsPriority keeps lifecycle and terminal frames out of the high-volume
237 // delta budget. Slow subscribers may recover text/history over HTTP, but they
238 // must still learn that a turn, prompt, or foreground-session transition ended.
239 func eventIsPriority(kind event.Kind) bool {
240 switch kind {
241 case event.Reasoning, event.Text, event.ToolProgress, event.StreamAttempt:
242 return false
243 default:
244 return true
245 }
246 }
247
248 // eventMustReachSubscriber identifies lifecycle truth that cannot be recovered
249 // by refetching history. The reserved queue budget protects these frames from
250 // deltas; if other priority traffic also exhausts that budget, the newest
251 // terminal/routing frame evicts a recoverable queued frame instead of vanishing.
252 func eventMustReachSubscriber(kind event.Kind, data []byte) bool {
253 switch kind {
254 case event.TurnDone, event.SessionChanged:
255 return true
256 case event.Notice:
257 return wireFrameMustReachSubscriber(data)
258 default:
259 return false
260 }
261 }
262
263 func enqueueSubscriberFrame(ch chan []byte, data []byte, kind event.Kind) {
264 priority := eventIsPriority(kind)
265 if !priority && len(ch) >= cap(ch)-subscriberPriorityReserve {
266 return
267 }
268 select {
269 case ch <- data:
270 return
271 default:
272 // A slow subscriber exhausted even the priority reserve. Ordinary
273 // priority events remain lossy, but lifecycle truth gets one slot by
274 // evicting an older frame while Broadcaster.mu serializes producers.
275 if !eventMustReachSubscriber(kind, data) || !evictRecoverableSubscriberFrame(ch) {
276 return
277 }
278 }
279 // Broadcaster.mu excludes every producer, and eviction leaves at least one
280 // slot. A concurrent consumer can only create more capacity, so this send is
281 // bounded while guaranteeing the terminal frame is retained.
282 ch <- data
283 }
284
285 func evictRecoverableSubscriberFrame(ch chan []byte) bool {
286 queued := len(ch)
287 if queued < cap(ch) {
288 return true
289 }
290 frames := make([][]byte, 0, queued)
291 drain:
292 for range queued {
293 select {
294 case frame := <-ch:
295 frames = append(frames, frame)
296 default:
297 break drain
298 }
299 }
300 if len(frames) == 0 {
301 return true
302 }
303 evict := -1
304 for i, frame := range frames {
305 if !wireFrameMustReachSubscriber(frame) {
306 evict = i
307 break
308 }
309 }
310 if evict < 0 {
311 // A bounded queue cannot retain an unbounded run of terminal frames.
312 // Prefer the latest lifecycle truth over an older one in that degenerate
313 // case; normal saturation always finds a recoverable delta/status frame.
314 evict = 0
315 }
316 for i, frame := range frames {
317 if i != evict {
318 ch <- frame
319 }
320 }
321 return true
322 }
323
324 func wireFrameMustReachSubscriber(data []byte) bool {
325 return eventwire.FrameMustReachMirror(data)
326 }
327
328 // EmitWire publishes an externally authored wire frame (same JSON contract as
329 // locally emitted events) to every subscriber. The local-takeover mirror uses
330 // it: a desktop writer that owns the session file pushes its frames through
331 // Serve so the remote tab keeps rendering the conversation live without any
332 // change to its own pipeline.
333 func (b *Broadcaster) EmitWire(wired eventwire.Event) {
334 if wired.SessionPath != "" {
335 wired.SessionPath = sessionRouteKey(wired.SessionPath)
336 }
337 b.mu.Lock()
338 observedCurrent := b.current
339 b.mu.Unlock()
340 wired.SessionCurrent = wired.SessionPath != "" && wired.SessionPath == observedCurrent
341 data, err := json.Marshal(wired)
342 if err != nil {
343 return
344 }
345 b.mu.Lock()
346 defer b.mu.Unlock()
347 if b.current != observedCurrent {
348 wired.SessionCurrent = wired.SessionPath != "" && wired.SessionPath == b.current
349 data, err = json.Marshal(wired)
350 if err != nil {
351 return
352 }
353 }
354 for ch := range b.subs {
355 enqueueSubscriberWireFrame(ch, data, wired.Kind)
356 }
357 }
358
359 // wireKindIsPriority mirrors eventIsPriority for externally supplied frames,
360 // which arrive as wire kind strings rather than typed event.Kind values.
361 func wireKindIsPriority(kind string) bool {
362 return !eventwire.WireKindIsRecoverable(kind)
363 }
364
365 func enqueueSubscriberWireFrame(ch chan []byte, data []byte, kind string) {
366 priority := wireKindIsPriority(kind)
367 if !priority && len(ch) >= cap(ch)-subscriberPriorityReserve {
368 return
369 }
370 select {
371 case ch <- data:
372 return
373 default:
374 if !wireFrameMustReachSubscriber(data) || !evictRecoverableSubscriberFrame(ch) {
375 return
376 }
377 }
378 ch <- data
379 }
380
381 // Subscribe registers a new SSE client and returns its channel plus an
382 // unsubscribe func the handler must call (defer) when the client disconnects.
383 func (b *Broadcaster) Subscribe() (<-chan []byte, func()) {
384 return b.subscribe(false)
385 }
386
387 // SubscribeAll receives tagged frames from current and detached sessions.
388 // Desktop uses it to maintain per-session runtime state; browser clients keep
389 // using Subscribe and see only the selected session.
390 func (b *Broadcaster) SubscribeAll() (<-chan []byte, func()) {
391 return b.subscribe(true)
392 }
393
394 func (b *Broadcaster) subscribe(all bool) (<-chan []byte, func()) {
395 ch := make(chan []byte, subscriberBufferSize)
396 b.mu.Lock()
397 b.subs[ch] = subscription{all: all}
398 b.mu.Unlock()
399 return ch, func() {
400 b.mu.Lock()
401 if _, ok := b.subs[ch]; ok {
402 delete(b.subs, ch)
403 close(ch)
404 }
405 b.mu.Unlock()
406 }
407 }
408
409 // Subscribers reports the current connection count (for diagnostics/tests).
410 func (b *Broadcaster) Subscribers() int {
411 b.mu.Lock()
412 defer b.mu.Unlock()
413 return len(b.subs)
414 }
415
415 lines GO