Initial commit: btcdata — BTC data backend service

Split from btcprice monolith. Owns DB, migrations, Binance WebSocket,
Kraken OHLC pollers, and exposes a JSON API with SSE price stream.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
2026-02-23 20:07:30 +01:00
co-authored by Claude Opus 4.6
commit f64aa98b04
35 changed files with 1079 additions and 0 deletions
+4
View File
@@ -0,0 +1,4 @@
target/
.env
.cpcache/
.nrepl-port
+17
View File
@@ -0,0 +1,17 @@
/target
/classes
/checkouts
profiles.clj
pom.xml
pom.xml.asc
*.jar
*.class
/.lein-*
/.nrepl-port
/.cpcache
/.clj-kondo
/log
.DS_Store
*.log
/.env
/token
+62
View File
@@ -0,0 +1,62 @@
# btcdata
Bitcoin data backend — Clojure Kit application. Fetches BTC prices from Binance WebSocket and Kraken HTTP, stores in PostgreSQL, exposes a JSON API with SSE price stream.
## Build & Development Commands
- `make run` — Start dev server on port 4100
- `make repl` — Start nREPL
- `make test` — Run tests
- `make uberjar` — Build production JAR
## REPL Commands
```clojure
(dev-prep!) ;; Prepare dev system
(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
- **Framework:** Kit (Integrant-based)
- **Server:** Undertow
- **Routing:** Reitit
- **Database:** PostgreSQL + conman + Migratus
- **Port:** 4100 (env: BTCDATA_PORT)
- **CORS:** Configured via CORS_ORIGIN env var
## API Endpoints
- `GET /api/health` — Health check
- `GET /api/price/stream` — SSE stream with JSON price updates
- `GET /api/price/latest` — Latest price as JSON
## Source Layout
```
src/clj/pmagnus/btcdata/
├── core.clj # App entry point
├── config.clj # System config loader
├── ws/
│ └── binance.clj # Binance WebSocket client
├── kraken/
│ ├── ohlc.clj # Hourly OHLC poller
│ └── ohlc_daily.clj # Daily OHLC poller
└── web/
├── handler.clj # Ring handler + routes
├── controllers/
│ └── health.clj # Health check
├── middleware/
│ ├── core.clj # Base middleware + CORS
│ ├── exception.clj # Exception handling
│ └── formats.clj # Content negotiation
└── routes/
├── api.clj # /api routes (JSON + SSE)
└── utils.clj # Route utilities
```
+18
View File
@@ -0,0 +1,18 @@
FROM clojure:temurin-21-tools-deps-alpine AS build
WORKDIR /build
COPY deps.edn build.clj ./
RUN clj -Sforce -P
COPY . .
RUN clj -T:build all
FROM eclipse-temurin:21-jre-alpine
COPY --from=build /build/target/btcdata-standalone.jar /btcdata/btcdata-standalone.jar
EXPOSE 4100
ENV BTCDATA_PORT=4100
ENV CORS_ORIGIN=http://localhost:4000
ENV JDBC_URL=jdbc:postgresql://postgres:5432/btcprod
CMD ["java", "-jar", "/btcdata/btcdata-standalone.jar"]
+19
View File
@@ -0,0 +1,19 @@
include .env
export
.PHONY: clean run repl test uberjar
clean:
rm -rf target
run:
clj -M:dev -e "(dev-prep!) (go)" -r
repl:
clj -M:dev:nrepl
test:
clj -M:test
uberjar:
clj -T:build all
+28
View File
@@ -0,0 +1,28 @@
(ns build
(:require [clojure.tools.build.api :as b]))
(def lib 'pmagnus/btcdata)
(def version "0.0.1-SNAPSHOT")
(def target-dir "target")
(def class-dir (str target-dir "/classes"))
(def basis (delay (b/create-basis {:project "deps.edn"})))
(def uber-file (str target-dir "/btcdata-standalone.jar"))
(defn clean [_]
(b/delete {:path target-dir}))
(defn uber [_]
(clean nil)
(b/copy-dir {:src-dirs ["src/clj" "resources" "env/prod/resources" "env/prod/clj"]
:target-dir class-dir})
(b/compile-clj {:basis @basis
:ns-compile '[pmagnus.btcdata.core]
:class-dir class-dir})
(b/uber {:class-dir class-dir
:uber-file uber-file
:basis @basis
:main 'pmagnus.btcdata.core}))
(defn all [_]
(uber nil))
+60
View File
@@ -0,0 +1,60 @@
{: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"}
;; Logging
ch.qos.logback/logback-classic {:mvn/version "1.5.20"}
;; HTTP Server
io.github.kit-clj/kit-undertow {:mvn/version "1.0.10"}
;; Routing
metosin/reitit {:mvn/version "0.9.2"}
metosin/reitit-ring {:mvn/version "0.9.2"}
metosin/reitit-middleware {:mvn/version "0.9.2"}
metosin/reitit-swagger {:mvn/version "0.9.2"}
metosin/reitit-swagger-ui {:mvn/version "0.9.2"}
metosin/reitit-malli {:mvn/version "0.9.2"}
;; Web
ring/ring-core {:mvn/version "1.15.3"}
ring/ring-defaults {:mvn/version "0.7.0"}
metosin/ring-http-response {:mvn/version "0.9.5"}
;; 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"}}
:aliases
{:build {:deps {io.github.clojure/tools.build {:mvn/version "0.10.9"}}
:ns-default build}
:dev {:extra-paths ["env/dev/clj" "env/dev/resources" "test/clj"]
:extra-deps {integrant/repl {:mvn/version "0.5.0"}
criterium/criterium {:mvn/version "0.4.6"}
expound/expound {:mvn/version "0.9.0"}
com.lambdaisland/classpath {:mvn/version "0.4.44"}
ring/ring-devel {:mvn/version "1.15.3"}}}
:test {:extra-paths ["env/test/resources"]
:extra-deps {io.github.cognitect-labs/test-runner {:git/tag "v0.5.1" :git/sha "dfb30dd"}
peridot/peridot {:mvn/version "0.5.4"}
clj-commons/byte-streams {:mvn/version "0.3.4"}}
:main-opts ["-m" "cognitect.test-runner"]
:exec-fn cognitect.test-runner.api/test}
:nrepl {:extra-deps {nrepl/nrepl {:mvn/version "1.3.1"}}
:main-opts ["-m" "nrepl.cmdline"
"--middleware" "[integrant.repl/middleware]"]}
:cider {:extra-deps {cider/cider-nrepl {:mvn/version "0.55.7"}}
:main-opts ["-m" "nrepl.cmdline"
"--middleware" "[cider.nrepl/cider-middleware]"]}}}
+4
View File
@@ -0,0 +1,4 @@
(ns pmagnus.btcdata.dev-middleware)
(defn wrap-dev [handler _opts]
handler)
+11
View File
@@ -0,0 +1,11 @@
(ns pmagnus.btcdata.env
(:require
[clojure.tools.logging :as log]
[pmagnus.btcdata.dev-middleware :refer [wrap-dev]]))
(def defaults
{:init (fn [] (log/info "\n-=[btcdata starting...dev/test profile]=-"))
:start (fn [] (log/info "\n-=[btcdata started...dev/test profile]=-"))
:stop (fn [] (log/info "\n-=[btcdata has shut down]=-"))
:middleware wrap-dev
:opts {:profile :dev}})
+50
View File
@@ -0,0 +1,50 @@
(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]
[pmagnus.btcdata.config :as config]
[pmagnus.btcdata.core]))
(defn dev-prep! []
(set-prep!
(fn []
(-> (config/system-config {:profile :dev})
(ig/expand)))))
(defn test-prep! []
(set-prep!
(fn []
(-> (config/system-config {:profile :test})
(ig/expand)))))
(defn reset-db []
(let [sys (or @pmagnus.btcdata.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.btcdata.core/system
integrant.repl.state/system)]
(when-let [mig (:db.sql/migrations sys)]
(migratus.core/rollback mig))))
(defn migrate []
(let [sys (or @pmagnus.btcdata.core/system
integrant.repl.state/system)]
(when-let [mig (:db.sql/migrations sys)]
(migratus.core/migrate mig))))
(defn query-fn []
(let [sys (or @pmagnus.btcdata.core/system
integrant.repl.state/system)]
(:db.sql/query-fn sys)))
(comment
(dev-prep!)
(go)
(reset)
(halt))
+15
View File
@@ -0,0 +1,15 @@
<configuration>
<statusListener class="ch.qos.logback.core.status.NopStatusListener" />
<appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n</pattern>
</encoder>
</appender>
<root level="INFO">
<appender-ref ref="STDOUT" />
</root>
<logger name="pmagnus.btcdata" level="DEBUG" />
<logger name="org.eclipse.jetty" level="WARN" />
<logger name="io.undertow" level="WARN" />
<logger name="org.xnio" level="WARN" />
</configuration>
+10
View File
@@ -0,0 +1,10 @@
(ns pmagnus.btcdata.env
(:require
[clojure.tools.logging :as log]))
(def defaults
{:init (fn [] (log/info "\n-=[btcdata starting]=-"))
:start (fn [] (log/info "\n-=[btcdata started successfully]=-"))
:stop (fn [] (log/info "\n-=[btcdata has shut down successfully]=-"))
:middleware (fn [handler _] handler)
:opts {:profile :prod}})
+14
View File
@@ -0,0 +1,14 @@
<configuration>
<statusListener class="ch.qos.logback.core.status.NopStatusListener" />
<appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n</pattern>
</encoder>
</appender>
<root level="INFO">
<appender-ref ref="STDOUT" />
</root>
<logger name="org.eclipse.jetty" level="WARN" />
<logger name="io.undertow" level="WARN" />
<logger name="org.xnio" level="WARN" />
</configuration>
+14
View File
@@ -0,0 +1,14 @@
<configuration>
<statusListener class="ch.qos.logback.core.status.NopStatusListener" />
<appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n</pattern>
</encoder>
</appender>
<root level="INFO">
<appender-ref ref="STDOUT" />
</root>
<logger name="org.eclipse.jetty" level="WARN" />
<logger name="io.undertow" level="WARN" />
<logger name="org.xnio" level="WARN" />
</configuration>
+8
View File
@@ -0,0 +1,8 @@
{:full-name "pmagnus/btcdata"
:ns-name "pmagnus.btcdata"
:sanitized "pmagnus/btcdata"
:name "btcdata"
:modules {:root "modules"
:repositories [{:url "https://github.com/kit-clj/modules.git"
:tag "master"
:name "kit-modules"}]}}
@@ -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()
);
@@ -0,0 +1 @@
DROP TABLE kraken_hour;
@@ -0,0 +1,11 @@
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()
);
@@ -0,0 +1 @@
DROP TABLE kraken_day;
@@ -0,0 +1,11 @@
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()
);
+44
View File
@@ -0,0 +1,44 @@
-- queries for btcdata
-- :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
+59
View File
@@ -0,0 +1,59 @@
{:system/env #profile {:dev :dev :test :test :prod :prod}
:server/http
{:port #long #or [#env BTCDATA_PORT 4100]
:host #or [#env HTTP_HOST "0.0.0.0"]
:handler #ig/ref :handler/ring}
:handler/ring
{:router #ig/ref :router/core
:api-path "/api"
:cors-origin #or [#env CORS_ORIGIN "*"]
:site-defaults-config
{:params {:keywordize true
:multipart true
:nested true}
:cookies false
:session false
:security {:anti-forgery false
:xss-protection {:enable? true :mode :block}
:frame-options :sameorigin}
:responses {:not-modified-responses true
:absolute-redirects true
:content-types true
:default-charset "utf-8"}}}
:router/routes
{:routes #ig/refset :reitit/routes}
:router/core
{: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
: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}}
+8
View File
@@ -0,0 +1,8 @@
(ns pmagnus.btcdata.config
(:require
[kit.config :as config]))
(def ^:const system-filename "system.edn")
(defn system-config [options]
(config/read-config system-filename options))
+44
View File
@@ -0,0 +1,44 @@
(ns pmagnus.btcdata.core
(:require
[clojure.tools.logging :as log]
[integrant.core :as ig]
[pmagnus.btcdata.config :as config]
[pmagnus.btcdata.env :refer [defaults]]
;; Edges
[kit.edge.db.postgres]
[kit.edge.db.sql.conman]
[kit.edge.db.sql.migratus]
[kit.edge.server.undertow]
[pmagnus.btcdata.web.handler]
;; Routes
[pmagnus.btcdata.web.routes.api]
;; WebSocket clients
[pmagnus.btcdata.ws.binance]
;; Pollers
[pmagnus.btcdata.kraken.ohlc]
[pmagnus.btcdata.kraken.ohlc-daily])
(:gen-class))
(defonce system (atom nil))
(defn stop-app []
((or (:stop defaults) (fn [])))
(some-> (deref system) (ig/halt!))
(shutdown-agents))
(defn start-app [& [params]]
((or (:init defaults) (fn [])))
(->> (config/system-config (or params {}))
(ig/expand)
(ig/init)
(reset! system))
((or (:start defaults) (fn []))))
(defn -main [& _]
(Thread/setDefaultUncaughtExceptionHandler
(reify Thread$UncaughtExceptionHandler
(uncaughtException [_ _thread ex]
(log/error ex "Uncaught exception"))))
(start-app)
(.addShutdownHook (Runtime/getRuntime) (Thread. stop-app)))
+124
View File
@@ -0,0 +1,124 @@
(ns pmagnus.btcdata.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))
@@ -0,0 +1,123 @@
(ns pmagnus.btcdata.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))
@@ -0,0 +1,11 @@
(ns pmagnus.btcdata.web.controllers.health
(:require
[ring.util.http-response :as response])
(:import
[java.lang.management ManagementFactory]))
(defn healthcheck! [_req]
(response/ok
{:time (str (java.time.Instant/now))
:up-time (.. ManagementFactory getRuntimeMXBean getUptime)
:app {:status "up"}}))
+30
View File
@@ -0,0 +1,30 @@
(ns pmagnus.btcdata.web.handler
(:require
[integrant.core :as ig]
[reitit.ring :as ring]
[pmagnus.btcdata.web.middleware.core :as middleware]))
(defmethod ig/init-key :handler/ring
[_ {:keys [router api-path] :as opts}]
(ring/ring-handler
(router)
(ring/routes
(when (some? api-path)
(reitit.ring/create-default-handler))
(ring/create-default-handler
{:not-found (constantly {:status 404 :body "Page not found"})
:method-not-allowed (constantly {:status 405 :body "Not allowed"})
:not-acceptable (constantly {:status 406 :body "Not acceptable"})}))
{:middleware [(middleware/wrap-base opts)]}))
(defmethod ig/init-key :router/routes
[_ {:keys [routes]}]
(mapv (fn [route]
(if (fn? route) (route) route))
routes))
(defmethod ig/init-key :router/core
[_ {:keys [routes env]}]
(if (= env :dev)
(fn [] (ring/router routes))
(constantly (ring/router routes))))
@@ -0,0 +1,27 @@
(ns pmagnus.btcdata.web.middleware.core
(:require
[ring.middleware.defaults :as defaults]))
(defn- wrap-cors
"Add CORS headers to every response. `origin` is the allowed origin string."
[handler origin]
(fn [request]
(if (= :options (:request-method request))
{:status 204
:headers {"Access-Control-Allow-Origin" origin
"Access-Control-Allow-Methods" "GET, OPTIONS"
"Access-Control-Allow-Headers" "Content-Type"
"Access-Control-Max-Age" "86400"}}
(when-let [resp (handler request)]
(update resp :headers assoc
"Access-Control-Allow-Origin" origin)))))
(defn wrap-base [{:keys [site-defaults-config cors-origin]}]
(let [origin (or cors-origin
(System/getenv "CORS_ORIGIN")
"*")]
(fn [handler]
(-> handler
(defaults/wrap-defaults
(or site-defaults-config defaults/api-defaults))
(wrap-cors origin)))))
@@ -0,0 +1,21 @@
(ns pmagnus.btcdata.web.middleware.exception
(:require
[clojure.tools.logging :as log]
[reitit.ring.middleware.exception :as exception]))
(defn- handler [message exception request]
(let [uri (:uri request)]
(log/error exception (str message " at " uri))
{:status 500
:body {:message message
:exception (.getClass exception)
:data (ex-data exception)
:uri uri}}))
(def exception-middleware
(exception/create-exception-middleware
(merge
exception/default-handlers
{::exception/default (partial handler "Internal error")
::exception/wrap (fn [handler e request]
(handler e request))})))
@@ -0,0 +1,12 @@
(ns pmagnus.btcdata.web.middleware.formats
(:require
[luminus-transit.time :as time]
[muuntaja.core :as m]))
(def instance
(m/create
(-> m/default-options
(update-in [:formats "application/transit+json" :decoder-opts]
(partial merge time/time-deserialization-handlers))
(update-in [:formats "application/transit+json" :encoder-opts]
(partial merge time/time-serialization-handlers)))))
+110
View File
@@ -0,0 +1,110 @@
(ns pmagnus.btcdata.web.routes.api
(:require
[clojure.data.json :as json]
[clojure.string :as str]
[integrant.core :as ig]
[pmagnus.btcdata.web.controllers.health :as health]
[pmagnus.btcdata.web.middleware.exception :as exception]
[pmagnus.btcdata.web.middleware.formats :as formats]
[reitit.coercion.malli :as malli]
[reitit.ring.coercion :as coercion]
[reitit.ring.middleware.muuntaja :as muuntaja]
[reitit.ring.middleware.parameters :as parameters]
[reitit.swagger :as swagger])
(:import
[io.undertow.server HttpServerExchange]
[io.undertow.util HttpString]
[java.io BufferedWriter OutputStreamWriter]
[java.util.concurrent LinkedBlockingQueue TimeUnit]))
(defn- format-sse [event data]
(str "event: " event "\n"
(str/join "\n" (map #(str "data: " %) (str/split-lines data)))
"\n\n"))
(defn- price->json [{:keys [price prev-price recorded-at]}]
(json/write-str
{:price (str price)
:prev_price (when prev-price (str prev-price))
:recorded_at (str recorded-at)}))
(defn- price-stream-handler [{:keys [binance]} req]
(let [^HttpServerExchange exchange (:server-exchange req)
price-atom (:latest-price binance)
headers (.getResponseHeaders exchange)
cors-origin (or (System/getenv "CORS_ORIGIN") "*")]
(.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")
(.put headers (HttpString. "Access-Control-Allow-Origin") cors-origin)
(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" (price->json 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- latest-price-handler [{:keys [binance]} _req]
(let [price-atom (:latest-price binance)
v @price-atom]
(if v
{:status 200
:headers {"Content-Type" "application/json"}
:body (price->json v)}
{:status 503
:headers {"Content-Type" "application/json"}
:body (json/write-str {:error "No price data yet"})})))
(defn- api-routes [opts]
[["/swagger.json"
{:get {:no-doc true
:swagger {:info {:title "btcdata API"}}
:handler (swagger/create-swagger-handler)}}]
["/health"
{:get health/healthcheck!}]
["/price/stream"
{:get (fn [req] (price-stream-handler opts req))
:no-doc true
:middleware []}]
["/price/latest"
{:get (fn [req] (latest-price-handler opts req))}]])
(defn route-data [opts]
(merge
opts
{:coercion malli/coercion
:muuntaja formats/instance
:swagger {:id ::api}
:middleware [parameters/parameters-middleware
muuntaja/format-negotiate-middleware
muuntaja/format-response-middleware
exception/exception-middleware
muuntaja/format-request-middleware
coercion/coerce-request-middleware
coercion/coerce-response-middleware]}))
(derive :reitit.routes/api :reitit/routes)
(defmethod ig/init-key :reitit.routes/api
[_ {:keys [base-path]
:or {base-path ""}
:as opts}]
[base-path (route-data opts) (api-routes opts)])
@@ -0,0 +1,9 @@
(ns pmagnus.btcdata.web.routes.utils)
(def route-data-path [:reitit.core/match :data])
(defn route-data [req]
(get-in req route-data-path))
(defn route-data-key [req k]
(get (route-data req) k))
+93
View File
@@ -0,0 +1,93 @@
(ns pmagnus.btcdata.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 _))))