diff --git a/CLAUDE.md b/CLAUDE.md index 90dbe00..950bc2e 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -1,15 +1,14 @@ # btcprice -Bitcoin price tracker - Clojure Kit web application. +Bitcoin price tracker UI — Clojure Kit web application. Serves the HTML/HTMX/Tailwind UI. Connects to btcdata for live price data via SSE. ## Build & Development Commands -- `make run` — Start dev server on port 3000 -- `make repl` — Start nREPL on port 7888 +- `make run` — Start dev server on port 4000 +- `make repl` — Start nREPL - `make test` — Run tests - `make tailwind` — Watch Tailwind CSS for changes - `make uberjar` — Build production JAR -- `docker compose up -d` — Start PostgreSQL ## REPL Commands @@ -18,10 +17,6 @@ Bitcoin price tracker - Clojure Kit web application. (go) ;; Start system (reset) ;; Reload code & restart (halt) ;; Stop system -(reset-db) ;; Drop & re-migrate database -(migrate) ;; Run pending migrations -(rollback) ;; Rollback last migration -(query-fn) ;; Get database query function ``` ## Architecture @@ -29,9 +24,10 @@ Bitcoin price tracker - Clojure Kit web application. - **Framework:** Kit (Integrant-based) - **Server:** Undertow - **Routing:** Reitit -- **Templates:** Hiccup + HTMX +- **Templates:** Hiccup + vanilla JS EventSource - **Styling:** Tailwind CSS v4 -- **Database:** PostgreSQL + conman + Migratus +- **Data backend:** btcdata (separate service at BTCDATA_URL, default http://localhost:4100) +- **No database** — all data comes from btcdata ## Source Layout @@ -50,6 +46,6 @@ src/clj/pmagnus/btcprice/ │ └── formats.clj # Content negotiation └── routes/ ├── api.clj # /api routes (JSON) - ├── ui.clj # UI routes (HTML) + ├── ui.clj # UI routes (HTML + EventSource) └── utils.clj # Route utilities ``` diff --git a/Dockerfile b/Dockerfile index 5019456..48b8a56 100644 --- a/Dockerfile +++ b/Dockerfile @@ -23,6 +23,6 @@ COPY --from=build /build/target/btcprice-standalone.jar /btcprice/btcprice-stand EXPOSE 4040 ENV PORT=4040 -ENV JDBC_URL=jdbc:postgresql://postgres:5432/btcprod +ENV BTCDATA_URL=http://btcdata:4100 CMD ["java", "-jar", "/btcprice/btcprice-standalone.jar"] diff --git a/deps.edn b/deps.edn index 85fde67..22ab3a2 100644 --- a/deps.edn +++ b/deps.edn @@ -1,7 +1,5 @@ {:paths ["src/clj" "resources"] :deps {org.clojure/clojure {:mvn/version "1.12.3"} - org.clojure/data.json {:mvn/version "2.5.1"} - ;; Kit io.github.kit-clj/kit-core {:mvn/version "1.0.6"} @@ -27,12 +25,7 @@ ;; Serialization metosin/muuntaja {:mvn/version "0.6.11"} - luminus-transit/luminus-transit {:mvn/version "0.1.6"} - - ;; Database - io.github.kit-clj/kit-postgres {:mvn/version "1.0.7"} - io.github.kit-clj/kit-sql-conman {:mvn/version "1.10.5"} - io.github.kit-clj/kit-sql-migratus {:mvn/version "1.0.5"}} + luminus-transit/luminus-transit {:mvn/version "0.1.6"}} :aliases {:build {:deps {io.github.clojure/tools.build {:mvn/version "0.10.9"}} diff --git a/docker-compose.yml b/docker-compose.yml index a997b19..4622ccb 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,10 +1,22 @@ services: + btcdata: + build: ../btcdata + restart: unless-stopped + ports: + - "4100:4100" + environment: + BTCDATA_PORT: "4100" + JDBC_URL: "jdbc:postgresql://postgres:5432/btcprod?user=postgres&password=ratata,123" + CORS_ORIGIN: "http://localhost:4040" + extra_hosts: + - "postgres:host-gateway" + app: build: . restart: unless-stopped ports: - "4040:4040" environment: - JDBC_URL: "jdbc:postgresql://postgres:5432/btcprod?user=postgres&password=ratata,123" - extra_hosts: - - "postgres:host-gateway" + BTCDATA_URL: "http://btcdata:4100" + depends_on: + - btcdata diff --git a/env/dev/clj/user.clj b/env/dev/clj/user.clj index 41e4600..4e8531c 100644 --- a/env/dev/clj/user.clj +++ b/env/dev/clj/user.clj @@ -1,6 +1,5 @@ (ns user (:require - [clojure.tools.logging :as log] [integrant.core :as ig] [integrant.repl :refer [go halt reset reset-all set-prep!]] [integrant.repl.state :as state] @@ -19,30 +18,6 @@ (-> (config/system-config {:profile :test}) (ig/expand))))) -(defn reset-db [] - (let [sys (or @pmagnus.btcprice.core/system - integrant.repl.state/system)] - (when-let [mig (:db.sql/migrations sys)] - (migratus.core/reset mig) - (log/info "Database reset complete")))) - -(defn rollback [] - (let [sys (or @pmagnus.btcprice.core/system - integrant.repl.state/system)] - (when-let [mig (:db.sql/migrations sys)] - (migratus.core/rollback mig)))) - -(defn migrate [] - (let [sys (or @pmagnus.btcprice.core/system - integrant.repl.state/system)] - (when-let [mig (:db.sql/migrations sys)] - (migratus.core/migrate mig)))) - -(defn query-fn [] - (let [sys (or @pmagnus.btcprice.core/system - integrant.repl.state/system)] - (:db.sql/query-fn sys))) - (comment (dev-prep!) (go) diff --git a/resources/migrations/20260217000000-create-binance-price.down.sql b/resources/migrations/20260217000000-create-binance-price.down.sql deleted file mode 100644 index 24d8300..0000000 --- a/resources/migrations/20260217000000-create-binance-price.down.sql +++ /dev/null @@ -1 +0,0 @@ -DROP TABLE binance_price; diff --git a/resources/migrations/20260217000000-create-binance-price.up.sql b/resources/migrations/20260217000000-create-binance-price.up.sql deleted file mode 100644 index 3f7c41b..0000000 --- a/resources/migrations/20260217000000-create-binance-price.up.sql +++ /dev/null @@ -1,5 +0,0 @@ -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/migrations/20260218000000-create-kraken-hour.down.sql b/resources/migrations/20260218000000-create-kraken-hour.down.sql deleted file mode 100644 index 073213e..0000000 --- a/resources/migrations/20260218000000-create-kraken-hour.down.sql +++ /dev/null @@ -1 +0,0 @@ -DROP TABLE kraken_hour; diff --git a/resources/migrations/20260218000000-create-kraken-hour.up.sql b/resources/migrations/20260218000000-create-kraken-hour.up.sql deleted file mode 100644 index eb1ad69..0000000 --- a/resources/migrations/20260218000000-create-kraken-hour.up.sql +++ /dev/null @@ -1,11 +0,0 @@ -CREATE TABLE kraken_hour ( - ts BIGINT NOT NULL PRIMARY KEY, - open NUMERIC(18,8) NOT NULL, - high NUMERIC(18,8) NOT NULL, - low NUMERIC(18,8) NOT NULL, - close NUMERIC(18,8) NOT NULL, - vwap NUMERIC(18,8) NOT NULL, - volume NUMERIC(24,8) NOT NULL, - trade_count INTEGER NOT NULL, - created_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); diff --git a/resources/migrations/20260218100000-create-kraken-day.down.sql b/resources/migrations/20260218100000-create-kraken-day.down.sql deleted file mode 100644 index a8f05eb..0000000 --- a/resources/migrations/20260218100000-create-kraken-day.down.sql +++ /dev/null @@ -1 +0,0 @@ -DROP TABLE kraken_day; diff --git a/resources/migrations/20260218100000-create-kraken-day.up.sql b/resources/migrations/20260218100000-create-kraken-day.up.sql deleted file mode 100644 index b694557..0000000 --- a/resources/migrations/20260218100000-create-kraken-day.up.sql +++ /dev/null @@ -1,11 +0,0 @@ -CREATE TABLE kraken_day ( - ts BIGINT NOT NULL PRIMARY KEY, - open NUMERIC(18,8) NOT NULL, - high NUMERIC(18,8) NOT NULL, - low NUMERIC(18,8) NOT NULL, - close NUMERIC(18,8) NOT NULL, - vwap NUMERIC(18,8) NOT NULL, - volume NUMERIC(24,8) NOT NULL, - trade_count INTEGER NOT NULL, - created_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); diff --git a/resources/migrations/placeholder.txt b/resources/migrations/placeholder.txt deleted file mode 100644 index 2586130..0000000 --- a/resources/migrations/placeholder.txt +++ /dev/null @@ -1 +0,0 @@ -Migration files go here. diff --git a/resources/queries.sql b/resources/queries.sql deleted file mode 100644 index 7c7a9ba..0000000 --- a/resources/queries.sql +++ /dev/null @@ -1,44 +0,0 @@ --- queries for btcprice - --- :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 - --- :name upsert-kraken-hour! :! :n --- :doc Upsert a Kraken hourly OHLC candle -INSERT INTO kraken_hour (ts, open, high, low, close, vwap, volume, trade_count) -VALUES (:ts, :open, :high, :low, :close, :vwap, :volume, :trade-count) -ON CONFLICT (ts) DO UPDATE -SET open = EXCLUDED.open, - high = EXCLUDED.high, - low = EXCLUDED.low, - close = EXCLUDED.close, - vwap = EXCLUDED.vwap, - volume = EXCLUDED.volume, - trade_count = EXCLUDED.trade_count - --- :name get-latest-kraken-hour :? :1 --- :doc Get the most recent Kraken hourly candle by timestamp -SELECT ts, open, high, low, close, vwap, volume, trade_count, created_at -FROM kraken_hour ORDER BY ts DESC LIMIT 1 - --- :name upsert-kraken-day! :! :n --- :doc Upsert a Kraken daily OHLC candle -INSERT INTO kraken_day (ts, open, high, low, close, vwap, volume, trade_count) -VALUES (:ts, :open, :high, :low, :close, :vwap, :volume, :trade-count) -ON CONFLICT (ts) DO UPDATE -SET open = EXCLUDED.open, - high = EXCLUDED.high, - low = EXCLUDED.low, - close = EXCLUDED.close, - vwap = EXCLUDED.vwap, - volume = EXCLUDED.volume, - trade_count = EXCLUDED.trade_count - --- :name get-latest-kraken-day :? :1 --- :doc Get the most recent Kraken daily candle by timestamp -SELECT ts FROM kraken_day ORDER BY ts DESC LIMIT 1 diff --git a/resources/system.edn b/resources/system.edn index a662a2c..ebc4640 100644 --- a/resources/system.edn +++ b/resources/system.edn @@ -32,34 +32,9 @@ {:routes #ig/ref :router/routes :env #ig/ref :system/env} - :db.sql/connection - {:jdbc-url #env JDBC_URL} - - :db.sql/query-fn - {:conn #ig/ref :db.sql/connection - :options {} - :filename "queries.sql"} - - :db.sql/migrations - {:store :database - :db {:datasource #ig/ref :db.sql/connection} - :migrate-on-init? true} - :reitit.routes/api - {:base-path "/api" - :query-fn #ig/ref :db.sql/query-fn} + {:base-path "/api"} :reitit.routes/ui - {:base-path "" - :query-fn #ig/ref :db.sql/query-fn - :binance #ig/ref :ws/binance} - - :ws/binance - {:query-fn #ig/ref :db.sql/query-fn - :uri "wss://stream.binance.com:9443/ws/btcusdt@trade"} - - :kraken/ohlc - {:query-fn #ig/ref :db.sql/query-fn} - - :kraken/ohlc-day - {:query-fn #ig/ref :db.sql/query-fn}} + {:base-path "" + :btcdata-url #or [#env BTCDATA_URL "http://localhost:4100"]}} diff --git a/src/clj/pmagnus/btcprice/core.clj b/src/clj/pmagnus/btcprice/core.clj index 6c3a6bb..f38ccb0 100644 --- a/src/clj/pmagnus/btcprice/core.clj +++ b/src/clj/pmagnus/btcprice/core.clj @@ -6,19 +6,11 @@ [pmagnus.btcprice.env :refer [defaults]] ;; Edges - [kit.edge.db.postgres] - [kit.edge.db.sql.conman] - [kit.edge.db.sql.migratus] [kit.edge.server.undertow] [pmagnus.btcprice.web.handler] ;; Routes [pmagnus.btcprice.web.routes.api] - [pmagnus.btcprice.web.routes.ui] - ;; WebSocket clients - [pmagnus.btcprice.ws.binance] - ;; Pollers - [pmagnus.btcprice.kraken.ohlc] - [pmagnus.btcprice.kraken.ohlc-daily]) + [pmagnus.btcprice.web.routes.ui]) (:gen-class)) (defonce system (atom nil)) diff --git a/src/clj/pmagnus/btcprice/kraken/ohlc.clj b/src/clj/pmagnus/btcprice/kraken/ohlc.clj deleted file mode 100644 index df2b36d..0000000 --- a/src/clj/pmagnus/btcprice/kraken/ohlc.clj +++ /dev/null @@ -1,124 +0,0 @@ -(ns pmagnus.btcprice.kraken.ohlc - (:require - [clojure.data.json :as json] - [clojure.tools.logging :as log] - [integrant.core :as ig]) - (:import - [java.net URI] - [java.net.http HttpClient HttpRequest HttpResponse$BodyHandlers])) - -(def ^:private kraken-ohlc-url - "https://api.kraken.com/0/public/OHLC?pair=XBTUSD&interval=60") - -(defn- fetch-ohlc - "HTTP GET to Kraken OHLC endpoint. Returns parsed JSON result map. - When `since` is provided, appends &since= to fetch only newer candles." - [^HttpClient client since] - (let [url (if since - (str kraken-ohlc-url "&since=" since) - kraken-ohlc-url) - request (-> (HttpRequest/newBuilder) - (.uri (URI. url)) - (.header "Accept" "application/json") - (.GET) - (.build)) - resp (.send client request (HttpResponse$BodyHandlers/ofString)) - body (json/read-str (.body resp) :key-fn keyword)] - (when-let [errors (seq (:error body))] - (throw (ex-info "Kraken API error" {:errors errors}))) - (:result body))) - -(defn- parse-candle - "Convert a Kraken OHLC array [ts, open, high, low, close, vwap, volume, count] - to a map with bigdec values." - [[ts open high low close vwap volume count]] - {:ts (long ts) - :open (bigdec open) - :high (bigdec high) - :low (bigdec low) - :close (bigdec close) - :vwap (bigdec vwap) - :volume (bigdec volume) - :trade-count (int count)}) - -(defn- save-candles! - "Upsert each candle into the kraken_hour table." - [query-fn candles] - (doseq [candle candles] - (query-fn :upsert-kraken-hour! candle))) - -(defn- seed-since-from-db - "Query DB for the latest candle timestamp. Returns it or nil." - [query-fn] - (some-> (query-fn :get-latest-kraken-hour {}) - :ts)) - -(defn- poll! - "Fetch OHLC data, drop the last (in-progress) candle, save completed ones. - Returns the count of saved candles." - [client query-fn since-atom] - (let [result (fetch-ohlc client @since-atom) - ;; Kraken returns a map with the pair key and a "last" key - last-ts (:last result) - pair-key (first (remove #{:last} (keys result))) - raw (get result pair-key) - candles (map parse-candle (butlast raw))] - (when (seq candles) - (save-candles! query-fn candles) - (when last-ts - (reset! since-atom last-ts)) - (log/info "Fetched" (count candles) "completed Kraken hourly candles")) - (count candles))) - -(defn- ms-until-next-poll - "Milliseconds from now until next HH:01:00 UTC." - [] - (let [now (java.time.ZonedDateTime/now java.time.ZoneOffset/UTC) - next (-> now - (.truncatedTo java.time.temporal.ChronoUnit/HOURS) - (.plusMinutes 1)) - target (if (.isAfter now next) - (.plusHours next 1) - next)] - (.toMillis (java.time.Duration/between now target)))) - -(defn- start-poll-loop! - "Start a background future that polls Kraken aligned to HH:01:00 UTC. - When `has-data?` is false, fetches immediately to backfill." - [client query-fn since-atom running? has-data?] - (future - (when-not has-data? - (try - (poll! client query-fn since-atom) - (catch Exception e - (log/error e "Kraken OHLC initial poll failed")))) - (while @running? - (let [wait (ms-until-next-poll)] - (log/info "Next Kraken OHLC poll in" (quot wait 60000) "minutes") - (Thread/sleep wait) - (when @running? - (try - (poll! client query-fn since-atom) - (catch Exception e - (log/error e "Kraken OHLC poll failed")))))))) - -(defmethod ig/init-key :kraken/ohlc - [_ {:keys [query-fn]}] - (let [client (HttpClient/newHttpClient) - db-since (seed-since-from-db query-fn) - since (atom db-since) - running? (atom true) - has-data? (some? db-since) - fut (start-poll-loop! client query-fn since running? has-data?)] - (if has-data? - (log/info "Starting Kraken OHLC poller, next fetch in" (quot (ms-until-next-poll) 60000) "minutes") - (log/info "Starting Kraken OHLC poller, fetching immediately (no data)")) - {:running? running? - :future fut - :since since})) - -(defmethod ig/halt-key! :kraken/ohlc - [_ {:keys [running? future]}] - (log/info "Stopping Kraken OHLC poller") - (reset! running? false) - (future-cancel future)) diff --git a/src/clj/pmagnus/btcprice/kraken/ohlc_daily.clj b/src/clj/pmagnus/btcprice/kraken/ohlc_daily.clj deleted file mode 100644 index 632d05a..0000000 --- a/src/clj/pmagnus/btcprice/kraken/ohlc_daily.clj +++ /dev/null @@ -1,123 +0,0 @@ -(ns pmagnus.btcprice.kraken.ohlc-daily - (:require - [clojure.data.json :as json] - [clojure.tools.logging :as log] - [integrant.core :as ig]) - (:import - [java.net URI] - [java.net.http HttpClient HttpRequest HttpResponse$BodyHandlers])) - -(def ^:private kraken-ohlc-url - "https://api.kraken.com/0/public/OHLC?pair=XBTUSD&interval=1440") - -(defn- fetch-ohlc - "HTTP GET to Kraken OHLC endpoint (daily). Returns parsed JSON result map. - When `since` is provided, appends &since= to fetch only newer candles." - [^HttpClient client since] - (let [url (if since - (str kraken-ohlc-url "&since=" since) - kraken-ohlc-url) - request (-> (HttpRequest/newBuilder) - (.uri (URI. url)) - (.header "Accept" "application/json") - (.GET) - (.build)) - resp (.send client request (HttpResponse$BodyHandlers/ofString)) - body (json/read-str (.body resp) :key-fn keyword)] - (when-let [errors (seq (:error body))] - (throw (ex-info "Kraken API error" {:errors errors}))) - (:result body))) - -(defn- parse-candle - "Convert a Kraken OHLC array [ts, open, high, low, close, vwap, volume, count] - to a map with bigdec values." - [[ts open high low close vwap volume count]] - {:ts (long ts) - :open (bigdec open) - :high (bigdec high) - :low (bigdec low) - :close (bigdec close) - :vwap (bigdec vwap) - :volume (bigdec volume) - :trade-count (int count)}) - -(defn- save-candles! - "Upsert each candle into the kraken_day table." - [query-fn candles] - (doseq [candle candles] - (query-fn :upsert-kraken-day! candle))) - -(defn- seed-since-from-db - "Query DB for the latest daily candle timestamp. Returns it or nil." - [query-fn] - (some-> (query-fn :get-latest-kraken-day {}) - :ts)) - -(defn- poll! - "Fetch daily OHLC data, drop the last (in-progress) candle, save completed ones. - Returns the count of saved candles." - [client query-fn since-atom] - (let [result (fetch-ohlc client @since-atom) - last-ts (:last result) - pair-key (first (remove #{:last} (keys result))) - raw (get result pair-key) - candles (map parse-candle (butlast raw))] - (when (seq candles) - (save-candles! query-fn candles) - (when last-ts - (reset! since-atom last-ts)) - (log/info "Fetched" (count candles) "completed Kraken daily candles")) - (count candles))) - -(defn- ms-until-next-daily-poll - "Milliseconds from now until next 00:01:00 UTC." - [] - (let [now (java.time.ZonedDateTime/now java.time.ZoneOffset/UTC) - next (-> now - (.truncatedTo java.time.temporal.ChronoUnit/DAYS) - (.plusMinutes 1)) - target (if (.isAfter now next) - (.plusDays next 1) - next)] - (.toMillis (java.time.Duration/between now target)))) - -(defn- start-poll-loop! - "Start a background future that polls Kraken daily OHLC aligned to 00:01:00 UTC. - When `has-data?` is false, fetches immediately to backfill." - [client query-fn since-atom running? has-data?] - (future - (when-not has-data? - (try - (poll! client query-fn since-atom) - (catch Exception e - (log/error e "Kraken daily OHLC initial poll failed")))) - (while @running? - (let [wait (ms-until-next-daily-poll)] - (log/info "Next Kraken daily OHLC poll in" (quot wait 60000) "minutes") - (Thread/sleep wait) - (when @running? - (try - (poll! client query-fn since-atom) - (catch Exception e - (log/error e "Kraken daily OHLC poll failed")))))))) - -(defmethod ig/init-key :kraken/ohlc-day - [_ {:keys [query-fn]}] - (let [client (HttpClient/newHttpClient) - db-since (seed-since-from-db query-fn) - since (atom db-since) - running? (atom true) - has-data? (some? db-since) - fut (start-poll-loop! client query-fn since running? has-data?)] - (if has-data? - (log/info "Starting Kraken daily OHLC poller, next fetch in" (quot (ms-until-next-daily-poll) 60000) "minutes") - (log/info "Starting Kraken daily OHLC poller, fetching immediately (no data)")) - {:running? running? - :future fut - :since since})) - -(defmethod ig/halt-key! :kraken/ohlc-day - [_ {:keys [running? future]}] - (log/info "Stopping Kraken daily OHLC poller") - (reset! running? false) - (future-cancel future)) diff --git a/src/clj/pmagnus/btcprice/web/htmx.clj b/src/clj/pmagnus/btcprice/web/htmx.clj index 329d2f7..d5019c0 100644 --- a/src/clj/pmagnus/btcprice/web/htmx.clj +++ b/src/clj/pmagnus/btcprice/web/htmx.clj @@ -17,9 +17,7 @@ [:meta {:name "apple-mobile-web-app-capable" :content "yes"}] [: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-ext-sse@2.2.2/sse.js"}]] + [:link {:rel "stylesheet" :href "/css/output.css"}]] [:body [:div.mx-auto.max-w-lg.px-4.py-6 ~@content]]]))})) diff --git a/src/clj/pmagnus/btcprice/web/routes/api.clj b/src/clj/pmagnus/btcprice/web/routes/api.clj index be1d90b..dbb8bc3 100644 --- a/src/clj/pmagnus/btcprice/web/routes/api.clj +++ b/src/clj/pmagnus/btcprice/web/routes/api.clj @@ -8,8 +8,7 @@ [reitit.ring.coercion :as coercion] [reitit.ring.middleware.muuntaja :as muuntaja] [reitit.ring.middleware.parameters :as parameters] - [reitit.swagger :as swagger] - [reitit.swagger-ui :as swagger-ui])) + [reitit.swagger :as swagger])) (defn- api-routes [_opts] [["/swagger.json" diff --git a/src/clj/pmagnus/btcprice/web/routes/ui.clj b/src/clj/pmagnus/btcprice/web/routes/ui.clj index 108260a..f02e187 100644 --- a/src/clj/pmagnus/btcprice/web/routes/ui.clj +++ b/src/clj/pmagnus/btcprice/web/routes/ui.clj @@ -1,96 +1,60 @@ (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]] [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]) - (:import - [io.undertow.server HttpServerExchange] - [io.undertow.util HttpString] - [java.io BufferedWriter OutputStreamWriter] - [java.time ZoneId] - [java.time.format DateTimeFormatter] - [java.util.concurrent LinkedBlockingQueue TimeUnit])) + [reitit.ring.middleware.parameters :as parameters])) -(def ^:private time-fmt - (-> (DateTimeFormatter/ofPattern "dd-MM HH:mm:ss") - (.withZone (ZoneId/systemDefault)))) - -(defn home-page [_req] +(defn- home-page [{:keys [btcdata-url]} _req] (page {:title "BTC Price"} [:header.text-center.mb-8 [:h1.text-3xl.font-bold.text-gray-900 "BTC Price"] [:p.text-sm.text-gray-500.mt-1 "Bitcoin price tracker"]] - [:div#price-panel {:hx-ext "sse" :sse-connect "/price/stream"} - [:div {:sse-swap "price-update"} + [:div#price-panel + [:div#price-display [:div.bg-white.rounded-2xl.shadow-sm.border.border-gray-200.p-6 - [:p.text-center.text-gray-400.text-sm "Connecting..."]]]])) + [:p.text-center.text-gray-400.text-sm "Connecting..."]]]] -;; Tailwind classes referenced for price color: -;; text-green-600 text-red-600 text-gray-900 -(defn- render-price-html [{:keys [price prev-price recorded-at]}] - (let [color (cond - (nil? prev-price) "text-gray-900" - (> price prev-price) "text-green-600" - (< price prev-price) "text-red-600" - :else "text-gray-900")] - (str - (h/html - [:div.bg-white.rounded-2xl.shadow-sm.border.border-gray-200.p-6 - [:div.text-center - [:p {:class (str "text-4xl font-bold " color)} - (str "$" (format "%,.2f" (double price)))] - [:p.text-xs.text-gray-400.mt-2 - (str "Updated " (.format time-fmt 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)) + [:script + (hiccup2.core/raw + (str + "document.addEventListener('DOMContentLoaded', function() {\n" + " var btcdataUrl = '" btcdata-url "';\n" + " var display = document.getElementById('price-display');\n" + " var source = new EventSource(btcdataUrl + '/api/price/stream');\n" + "\n" + " source.addEventListener('price-update', function(e) {\n" + " var d = JSON.parse(e.data);\n" + " var price = parseFloat(d.price);\n" + " var prevPrice = d.prev_price ? parseFloat(d.prev_price) : null;\n" + " var color = 'text-gray-900';\n" + " if (prevPrice !== null) {\n" + " if (price > prevPrice) color = 'text-green-600';\n" + " else if (price < prevPrice) color = 'text-red-600';\n" + " }\n" + " var fmt = '$' + price.toLocaleString('en-US', {minimumFractionDigits: 2, maximumFractionDigits: 2});\n" + " var ts = new Date(d.recorded_at);\n" + " var pad = function(n) { return n < 10 ? '0' + n : n; };\n" + " var timeStr = pad(ts.getDate()) + '-' + pad(ts.getMonth()+1) + ' ' + pad(ts.getHours()) + ':' + pad(ts.getMinutes()) + ':' + pad(ts.getSeconds());\n" + " display.innerHTML = '
' +\n" + " '
' +\n" + " '

' + fmt + '

' +\n" + " '

Updated ' + timeStr + '

' +\n" + " '
';\n" + " });\n" + "\n" + " source.onerror = function() {\n" + " display.innerHTML = '
' +\n" + " '

Connection lost. Reconnecting...

';\n" + " };\n" + "});\n"))])) (defn- ui-routes [opts] [["/" - {:get home-page}] - ["/price/stream" - {:get (fn [req] (price-stream-handler opts req)) - :middleware []}]]) + {:get (fn [req] (home-page 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 deleted file mode 100644 index 529866f..0000000 --- a/src/clj/pmagnus/btcprice/ws/binance.clj +++ /dev/null @@ -1,93 +0,0 @@ -(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. - `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] - (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) - (let [prev-price (:price @latest-price)] - (reset! latest-price - {:price price - :prev-price prev-price - :recorded-at (java.time.Instant/now)})))))) - (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 latest-price)] - (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 [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 - :latest-price latest-price})) - -(defmethod ig/halt-key! :ws/binance - [_ {:keys [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 _))))