From a50256a0426b763a6a4dd72f9753db1de2a6a167 Mon Sep 17 00:00:00 2001 From: Per Magnus Petersen Date: Tue, 17 Feb 2026 15:51:10 +0100 Subject: [PATCH] Add real-time BTC price feed via Binance WebSocket Connect to Binance trade stream, persist prices to binance_price table (throttled to every 5s), and live-update the home page via HTMX polling. Co-Authored-By: Claude Opus 4.6 --- ...260217000000-create-binance-price.down.sql | 1 + ...20260217000000-create-binance-price.up.sql | 5 ++ resources/queries.sql | 9 +- resources/system.edn | 6 +- src/clj/pmagnus/btcprice/core.clj | 4 +- src/clj/pmagnus/btcprice/web/routes/ui.clj | 21 ++++- src/clj/pmagnus/btcprice/ws/binance.clj | 85 +++++++++++++++++++ 7 files changed, 125 insertions(+), 6 deletions(-) create mode 100644 resources/migrations/20260217000000-create-binance-price.down.sql create mode 100644 resources/migrations/20260217000000-create-binance-price.up.sql create mode 100644 src/clj/pmagnus/btcprice/ws/binance.clj diff --git a/resources/migrations/20260217000000-create-binance-price.down.sql b/resources/migrations/20260217000000-create-binance-price.down.sql new file mode 100644 index 0000000..24d8300 --- /dev/null +++ b/resources/migrations/20260217000000-create-binance-price.down.sql @@ -0,0 +1 @@ +DROP TABLE binance_price; diff --git a/resources/migrations/20260217000000-create-binance-price.up.sql b/resources/migrations/20260217000000-create-binance-price.up.sql new file mode 100644 index 0000000..3f7c41b --- /dev/null +++ b/resources/migrations/20260217000000-create-binance-price.up.sql @@ -0,0 +1,5 @@ +CREATE TABLE binance_price ( + id BIGSERIAL PRIMARY KEY, + price NUMERIC(18,8) NOT NULL, + recorded_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); diff --git a/resources/queries.sql b/resources/queries.sql index 710dc18..98297fd 100644 --- a/resources/queries.sql +++ b/resources/queries.sql @@ -1,2 +1,9 @@ -- queries for btcprice --- Add HugSQL queries here as tables are created + +-- :name insert-binance-price! :! :n +-- :doc Insert a Binance BTC price +INSERT INTO binance_price (price) VALUES (:price) + +-- :name get-latest-binance-price :? :1 +-- :doc Get the most recent Binance BTC price +SELECT price, recorded_at FROM binance_price ORDER BY id DESC LIMIT 1 diff --git a/resources/system.edn b/resources/system.edn index 3603376..277c40b 100644 --- a/resources/system.edn +++ b/resources/system.edn @@ -51,4 +51,8 @@ :reitit.routes/ui {:base-path "" - :query-fn #ig/ref :db.sql/query-fn}} + :query-fn #ig/ref :db.sql/query-fn} + + :ws/binance + {:query-fn #ig/ref :db.sql/query-fn + :uri "wss://stream.binance.com:9443/ws/btcusdt@trade"}} diff --git a/src/clj/pmagnus/btcprice/core.clj b/src/clj/pmagnus/btcprice/core.clj index c184c8b..e312ec4 100644 --- a/src/clj/pmagnus/btcprice/core.clj +++ b/src/clj/pmagnus/btcprice/core.clj @@ -13,7 +13,9 @@ [pmagnus.btcprice.web.handler] ;; Routes [pmagnus.btcprice.web.routes.api] - [pmagnus.btcprice.web.routes.ui]) + [pmagnus.btcprice.web.routes.ui] + ;; WebSocket clients + [pmagnus.btcprice.ws.binance]) (:gen-class)) (defonce system (atom nil)) diff --git a/src/clj/pmagnus/btcprice/web/routes/ui.clj b/src/clj/pmagnus/btcprice/web/routes/ui.clj index a4bb894..7c4b777 100644 --- a/src/clj/pmagnus/btcprice/web/routes/ui.clj +++ b/src/clj/pmagnus/btcprice/web/routes/ui.clj @@ -1,7 +1,7 @@ (ns pmagnus.btcprice.web.routes.ui (:require [integrant.core :as ig] - [pmagnus.btcprice.web.htmx :refer [page]] + [pmagnus.btcprice.web.htmx :refer [page fragment]] [pmagnus.btcprice.web.middleware.exception :as exception] [pmagnus.btcprice.web.middleware.formats :as formats] [reitit.ring.middleware.muuntaja :as muuntaja] @@ -14,12 +14,27 @@ [: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- ui-routes [_opts] +(defn price-fragment [{:keys [query-fn]}] + (let [row (query-fn :get-latest-binance-price {})] + (fragment {} + [: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..."])]))) + +(defn- ui-routes [opts] [["/" - {:get home-page}]]) + {:get home-page}] + ["/price" + {:get (fn [_req] (price-fragment opts))}]]) (defn route-data [opts] (merge diff --git a/src/clj/pmagnus/btcprice/ws/binance.clj b/src/clj/pmagnus/btcprice/ws/binance.clj new file mode 100644 index 0000000..32045dd --- /dev/null +++ b/src/clj/pmagnus/btcprice/ws/binance.clj @@ -0,0 +1,85 @@ +(ns pmagnus.btcprice.ws.binance + (:require + [clojure.data.json :as json] + [clojure.tools.logging :as log] + [integrant.core :as ig]) + (:import + [java.net URI] + [java.net.http HttpClient WebSocket WebSocket$Listener] + [java.util.concurrent CompletableFuture CompletionStage])) + +(defn- save-price! [query-fn price] + (try + (query-fn :insert-binance-price! {:price price}) + (catch Exception e + (log/error e "Failed to save BTC price")))) + +(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] + (let [client (HttpClient/newHttpClient) + listener (reify WebSocket$Listener + (onOpen [_ ws] + (log/info "Binance WebSocket connected") + (.request ws 1)) + + (onText [_ ws data last?] + (let [buf (:buffer @state)] + (.append buf data) + (when last? + (let [text (str buf)] + (.setLength buf 0) + (try + (let [msg (json/read-str text :key-fn keyword) + price (some-> (:p msg) bigdec)] + (when price + (let [now (System/currentTimeMillis)] + (when (> (- now (:last-write @state)) 5000) + (swap! state assoc :last-write now) + (log/info "BTC price:" (str price)) + (save-price! query-fn price))))) + (catch Exception e + (log/error e "Failed to parse Binance message"))))) + (let [^CompletionStage cf (CompletableFuture/completedFuture nil)] + (.request ws 1) + cf))) + + (onClose [_ _ws code reason] + (log/warn "Binance WebSocket closed:" code reason) + (when (:running? @state) + (future + (Thread/sleep 3000) + (when (:running? @state) + (log/info "Reconnecting to Binance...") + (try + (let [ws (connect! uri query-fn state)] + (swap! state assoc :ws ws)) + (catch Exception e + (log/error e "Binance reconnect failed"))))))) + + (onError [_ _ws error] + (log/error error "Binance WebSocket error")))] + (-> (.newWebSocketBuilder client) + (.buildAsync (URI. uri) listener) + (.join)))) + +(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)] + (swap! state assoc :ws ws) + state)) + +(defmethod ig/halt-key! :ws/binance + [_ state] + (log/info "Stopping Binance WebSocket listener") + (swap! state assoc :running? false) + (when-let [ws (:ws @state)] + (try + (.sendClose ws WebSocket/NORMAL_CLOSURE "shutting down") + (catch Exception _))))