diff --git a/resources/system.edn b/resources/system.edn index 277c40b..7475742 100644 --- a/resources/system.edn +++ b/resources/system.edn @@ -51,7 +51,8 @@ :reitit.routes/ui {:base-path "" - :query-fn #ig/ref :db.sql/query-fn} + :query-fn #ig/ref :db.sql/query-fn + :binance #ig/ref :ws/binance} :ws/binance {:query-fn #ig/ref :db.sql/query-fn diff --git a/src/clj/pmagnus/btcprice/web/htmx.clj b/src/clj/pmagnus/btcprice/web/htmx.clj index 69f8976..329d2f7 100644 --- a/src/clj/pmagnus/btcprice/web/htmx.clj +++ b/src/clj/pmagnus/btcprice/web/htmx.clj @@ -18,7 +18,8 @@ [:meta {:name "apple-mobile-web-app-status-bar-style" :content "black-translucent"}] [:title (or (:title opts#) "BTC Price")] [:link {:rel "stylesheet" :href "/css/output.css"}] - [:script {:src "https://unpkg.com/htmx.org@2.0.8"}]] + [:script {:src "https://unpkg.com/htmx.org@2.0.8"}] + [:script {:src "https://unpkg.com/htmx-ext-sse@2.2.2/sse.js"}]] [:body [:div.mx-auto.max-w-lg.px-4.py-6 ~@content]]]))})) diff --git a/src/clj/pmagnus/btcprice/web/routes/ui.clj b/src/clj/pmagnus/btcprice/web/routes/ui.clj index 7c4b777..2138542 100644 --- a/src/clj/pmagnus/btcprice/web/routes/ui.clj +++ b/src/clj/pmagnus/btcprice/web/routes/ui.clj @@ -1,11 +1,18 @@ (ns pmagnus.btcprice.web.routes.ui (:require + [clojure.string :as str] + [hiccup2.core :as h] [integrant.core :as ig] - [pmagnus.btcprice.web.htmx :refer [page fragment]] + [pmagnus.btcprice.web.htmx :refer [page]] [pmagnus.btcprice.web.middleware.exception :as exception] [pmagnus.btcprice.web.middleware.formats :as formats] [reitit.ring.middleware.muuntaja :as muuntaja] - [reitit.ring.middleware.parameters :as parameters])) + [reitit.ring.middleware.parameters :as parameters]) + (:import + [io.undertow.server HttpServerExchange] + [io.undertow.util HttpString] + [java.io BufferedWriter OutputStreamWriter] + [java.util.concurrent LinkedBlockingQueue TimeUnit])) (defn home-page [_req] (page {:title "BTC Price"} @@ -13,28 +20,63 @@ [:h1.text-3xl.font-bold.text-gray-900 "BTC Price"] [:p.text-sm.text-gray-500.mt-1 "Bitcoin price tracker"]] - [:div#price-panel.space-y-4 - {:hx-get "/price" :hx-trigger "every 2s" :hx-swap "innerHTML"} - [:div.bg-white.rounded-2xl.shadow-sm.border.border-gray-200.p-6 - [:p.text-center.text-gray-400.text-sm "Loading..."]]])) - -(defn price-fragment [{:keys [query-fn]}] - (let [row (query-fn :get-latest-binance-price {})] - (fragment {} + [:div#price-panel {:hx-ext "sse" :sse-connect "/price/stream"} + [:div {:sse-swap "price-update"} [:div.bg-white.rounded-2xl.shadow-sm.border.border-gray-200.p-6 - (if row - [:div.text-center - [:p.text-4xl.font-bold.text-gray-900 - (str "$" (format "%,.2f" (double (:price row))))] - [:p.text-xs.text-gray-400.mt-2 - (str "Updated " (:recorded_at row))]] - [:p.text-center.text-gray-400.text-sm "Waiting for data..."])]))) + [:p.text-center.text-gray-400.text-sm "Connecting..."]]]])) + +(defn- render-price-html [{:keys [price recorded-at]}] + (str + (h/html + [:div.bg-white.rounded-2xl.shadow-sm.border.border-gray-200.p-6 + [:div.text-center + [:p.text-4xl.font-bold.text-gray-900 + (str "$" (format "%,.2f" (double price)))] + [:p.text-xs.text-gray-400.mt-2 + (str "Updated " recorded-at)]]]))) + +(defn- format-sse [event data] + (str "event: " event "\n" + (str/join "\n" (map #(str "data: " %) (str/split-lines data))) + "\n\n")) + +(defn- price-stream-handler [{:keys [binance]} req] + (let [^HttpServerExchange exchange (:server-exchange req) + price-atom (:latest-price binance) + headers (.getResponseHeaders exchange)] + (.setStatusCode exchange 200) + (.put headers (HttpString. "Content-Type") "text/event-stream") + (.put headers (HttpString. "Cache-Control") "no-cache") + (.put headers (HttpString. "X-Accel-Buffering") "no") + (when-not (.isBlocking exchange) + (.startBlocking exchange)) + (let [out (.getOutputStream exchange) + writer (BufferedWriter. (OutputStreamWriter. out "UTF-8")) + queue (LinkedBlockingQueue.) + wkey (keyword (gensym "sse-"))] + (add-watch price-atom wkey + (fn [_ _ _ v] (.offer queue v))) + (when-let [v @price-atom] + (.offer queue v)) + (try + (loop [] + (if-let [v (.poll queue 15 TimeUnit/SECONDS)] + (do (.write writer (format-sse "price-update" (render-price-html v))) + (.flush writer)) + (do (.write writer ": heartbeat\n\n") + (.flush writer))) + (recur)) + (catch Exception _) + (finally + (remove-watch price-atom wkey) + (try (.close writer) (catch Exception _))))) + nil)) (defn- ui-routes [opts] [["/" {:get home-page}] - ["/price" - {:get (fn [_req] (price-fragment opts))}]]) + ["/price/stream" + {:get (fn [req] (price-stream-handler opts req))}]]) (defn route-data [opts] (merge diff --git a/src/clj/pmagnus/btcprice/ws/binance.clj b/src/clj/pmagnus/btcprice/ws/binance.clj index 32045dd..e5beded 100644 --- a/src/clj/pmagnus/btcprice/ws/binance.clj +++ b/src/clj/pmagnus/btcprice/ws/binance.clj @@ -16,8 +16,9 @@ (defn- connect! "Opens a WebSocket to Binance trade stream and returns the WebSocket instance. - `state` is an atom with keys :running?, :last-write, :buffer." - [uri query-fn state] + `state` is an atom with keys :running?, :last-write, :buffer. + `latest-price` is an atom reset with each throttled price for SSE consumers." + [uri query-fn state latest-price] (let [client (HttpClient/newHttpClient) listener (reify WebSocket$Listener (onOpen [_ ws] @@ -38,7 +39,10 @@ (when (> (- now (:last-write @state)) 5000) (swap! state assoc :last-write now) (log/info "BTC price:" (str price)) - (save-price! query-fn price))))) + (save-price! query-fn price) + (reset! latest-price + {:price price + :recorded-at (java.time.Instant/now)}))))) (catch Exception e (log/error e "Failed to parse Binance message"))))) (let [^CompletionStage cf (CompletableFuture/completedFuture nil)] @@ -53,7 +57,7 @@ (when (:running? @state) (log/info "Reconnecting to Binance...") (try - (let [ws (connect! uri query-fn state)] + (let [ws (connect! uri query-fn state latest-price)] (swap! state assoc :ws ws)) (catch Exception e (log/error e "Binance reconnect failed"))))))) @@ -67,16 +71,18 @@ (defmethod ig/init-key :ws/binance [_ {:keys [query-fn uri]}] (log/info "Starting Binance WebSocket listener:" uri) - (let [state (atom {:running? true - :last-write 0 - :buffer (StringBuilder.) - :ws nil}) - ws (connect! uri query-fn state)] + (let [latest-price (atom nil) + state (atom {:running? true + :last-write 0 + :buffer (StringBuilder.) + :ws nil}) + ws (connect! uri query-fn state latest-price)] (swap! state assoc :ws ws) - state)) + {:state state + :latest-price latest-price})) (defmethod ig/halt-key! :ws/binance - [_ state] + [_ {:keys [state]}] (log/info "Stopping Binance WebSocket listener") (swap! state assoc :running? false) (when-let [ws (:ws @state)]