1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
|
package bridge
import (
"context"
"net/http"
"github.com/coder/websocket"
)
// WebSocketHandler returns the HTTP handler that upgrades a request to a
// WebSocket and attaches it to the bridge: replayed history first, then live
// events, while incoming messages feed HandleIncoming.
//
// Example:
//
// mux.Handle("GET /ws", b.WebSocketHandler())
func (b *Bridge) WebSocketHandler() http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
conn, err := websocket.Accept(w, r, &websocket.AcceptOptions{
// ori serves localhost setups (possibly port-forwarded); the
// SPA and this endpoint share an origin in production, and the
// Vite dev server proxies /ws, so localhost patterns suffice.
OriginPatterns: []string{"localhost:*", "127.0.0.1:*"},
})
if err != nil {
b.logger.Warn("websocket accept failed", "error", err)
return
}
defer func() { _ = conn.CloseNow() }()
b.serveClient(r.Context(), conn)
})
}
// serveClient runs one client's session: replay, then a write pump for live
// events alongside a read loop for incoming messages, until either side ends.
func (b *Bridge) serveClient(ctx context.Context, conn *websocket.Conn) {
events, replay, unsubscribe := b.Subscribe()
defer unsubscribe()
if !writeEvents(ctx, conn, replay) {
return
}
writeDone := make(chan struct{})
go func() {
defer close(writeDone)
pumpEvents(ctx, conn, events)
}()
b.readLoop(ctx, conn)
// Stop the write pump by closing its channel, then close cleanly.
unsubscribe()
<-writeDone
_ = conn.Close(websocket.StatusNormalClosure, "")
}
// writeEvents sends a batch of events; false means the connection is gone.
func writeEvents(ctx context.Context, conn *websocket.Conn, events [][]byte) bool {
for _, event := range events {
if err := conn.Write(ctx, websocket.MessageText, event); err != nil {
return false
}
}
return true
}
// pumpEvents forwards live bridge events to the socket until the
// subscription channel closes or a write fails.
func pumpEvents(ctx context.Context, conn *websocket.Conn, events <-chan []byte) {
for event := range events {
if err := conn.Write(ctx, websocket.MessageText, event); err != nil {
return
}
}
}
// readLoop feeds client messages to the bridge until the socket closes.
func (b *Bridge) readLoop(ctx context.Context, conn *websocket.Conn) {
for {
kind, data, err := conn.Read(ctx)
if err != nil {
return
}
if kind == websocket.MessageText {
b.HandleIncoming(ctx, data)
}
}
}
|