diff --git a/README.md b/README.md index 14adf20..1098286 100644 --- a/README.md +++ b/README.md @@ -7,6 +7,10 @@ auto-reconnecting jetstream proxy with a default pool that should work NOTES: - this should run as close to your infrastructure as possible. you then would connect to `ws://localhost:6969/subscribe` +- no cursor support. since there will be multiple jetstream upstreams that + run at separate cursor timelines, it would be pretty hard to rewrite cursors + in such a way that everything works. if your application relies on cursors, + then it's probably best for your application to deal with multi-upstream support. ``` go build diff --git a/main.go b/main.go index 5744e20..69af466 100644 --- a/main.go +++ b/main.go @@ -7,6 +7,7 @@ import ( "os" "strings" "sync" + "sync/atomic" "time" "github.com/gorilla/websocket" @@ -26,6 +27,7 @@ var DEFAULT_POOL = []string{ type Broadcaster struct { listeners []chan []byte mu sync.Mutex + connected atomic.Bool } // Subscribe returns a new channel that will receive Jetstream events @@ -84,6 +86,8 @@ func measureLatency(url string) (time.Duration, error) { } start := time.Now() + // jetstream instances return the "Welcome to jetstream!" banner on / which + // should be useful enough for latency resp, err := client.Get(url) if err != nil { return 0, err @@ -99,6 +103,19 @@ var upgrader = websocket.Upgrader{ }, } +// handleHealth returns 200 if connected to upstream +func handleHealth(broadcaster *Broadcaster) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + if broadcaster.connected.Load() { + w.WriteHeader(http.StatusOK) + w.Write([]byte("Welcome to jetstream!")) + } else { + w.WriteHeader(http.StatusServiceUnavailable) + w.Write([]byte("Not connected to upstream")) + } + } +} + // handleSubscribe upgrades HTTP connection to websocket and streams events func handleSubscribe(broadcaster *Broadcaster) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { @@ -130,8 +147,8 @@ func handleSubscribe(broadcaster *Broadcaster) http.HandlerFunc { // connectToUpstream maintains a connection to the upstream websocket and broadcasts messages func connectToUpstream(pool []string, broadcaster *Broadcaster) { - backoff := time.Second - maxBackoff := time.Minute + backoff := 50 * time.Millisecond + maxBackoff := 20 * time.Second var currentUpstream string for { @@ -157,6 +174,7 @@ func connectToUpstream(pool []string, broadcaster *Broadcaster) { conn, _, err := websocket.DefaultDialer.Dial(currentUpstream+"/subscribe", nil) if err != nil { slog.Error("Failed to connect to upstream", slog.String("url", currentUpstream), slog.Any("error", err)) + broadcaster.connected.Store(false) time.Sleep(backoff) backoff *= 2 if backoff > maxBackoff { @@ -166,6 +184,7 @@ func connectToUpstream(pool []string, broadcaster *Broadcaster) { } slog.Info("Connected to upstream", slog.String("url", currentUpstream)) + broadcaster.connected.Store(true) backoff = time.Second // Reset backoff on successful connection // Read messages from upstream and broadcast them @@ -173,6 +192,7 @@ func connectToUpstream(pool []string, broadcaster *Broadcaster) { messageType, message, err := conn.ReadMessage() if err != nil { slog.Error("Error reading from upstream", slog.Any("error", err)) + broadcaster.connected.Store(false) conn.Close() break } @@ -271,6 +291,7 @@ func main() { go connectToUpstream(pool, broadcaster) // Setup HTTP server + http.HandleFunc("/", handleHealth(broadcaster)) http.HandleFunc("/subscribe", handleSubscribe(broadcaster)) slog.Info("Starting proxy server", slog.String("bind", bindAddr))