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 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1 @@
|
||||
DROP TABLE binance_price;
|
||||
@@ -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()
|
||||
);
|
||||
@@ -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
|
||||
|
||||
@@ -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"}}
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 _))))
|
||||
Reference in New Issue
Block a user