commit f64aa98b0494fa68027a43087a24c07c90585af6 Author: Per Magnus Petersen Date: Mon Feb 23 20:07:30 2026 +0100 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 diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..65cc69f --- /dev/null +++ b/.dockerignore @@ -0,0 +1,4 @@ +target/ +.env +.cpcache/ +.nrepl-port diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..cdb6118 --- /dev/null +++ b/.gitignore @@ -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 diff --git a/CLAUDE.md b/CLAUDE.md new file mode 100644 index 0000000..369bff0 --- /dev/null +++ b/CLAUDE.md @@ -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 +``` diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..37232ce --- /dev/null +++ b/Dockerfile @@ -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"] diff --git a/Makefile b/Makefile new file mode 100644 index 0000000..3ff33ac --- /dev/null +++ b/Makefile @@ -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 diff --git a/build.clj b/build.clj new file mode 100644 index 0000000..996e053 --- /dev/null +++ b/build.clj @@ -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)) diff --git a/deps.edn b/deps.edn new file mode 100644 index 0000000..bf9097a --- /dev/null +++ b/deps.edn @@ -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]"]}}} diff --git a/env/dev/clj/pmagnus/btcdata/dev_middleware.clj b/env/dev/clj/pmagnus/btcdata/dev_middleware.clj new file mode 100644 index 0000000..72098df --- /dev/null +++ b/env/dev/clj/pmagnus/btcdata/dev_middleware.clj @@ -0,0 +1,4 @@ +(ns pmagnus.btcdata.dev-middleware) + +(defn wrap-dev [handler _opts] + handler) diff --git a/env/dev/clj/pmagnus/btcdata/env.clj b/env/dev/clj/pmagnus/btcdata/env.clj new file mode 100644 index 0000000..060c518 --- /dev/null +++ b/env/dev/clj/pmagnus/btcdata/env.clj @@ -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}}) diff --git a/env/dev/clj/user.clj b/env/dev/clj/user.clj new file mode 100644 index 0000000..625e2a4 --- /dev/null +++ b/env/dev/clj/user.clj @@ -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)) diff --git a/env/dev/resources/logback.xml b/env/dev/resources/logback.xml new file mode 100644 index 0000000..1440e27 --- /dev/null +++ b/env/dev/resources/logback.xml @@ -0,0 +1,15 @@ + + + + + %d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n + + + + + + + + + + diff --git a/env/prod/clj/pmagnus/btcdata/env.clj b/env/prod/clj/pmagnus/btcdata/env.clj new file mode 100644 index 0000000..8ec5ab5 --- /dev/null +++ b/env/prod/clj/pmagnus/btcdata/env.clj @@ -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}}) diff --git a/env/prod/resources/logback.xml b/env/prod/resources/logback.xml new file mode 100644 index 0000000..5194993 --- /dev/null +++ b/env/prod/resources/logback.xml @@ -0,0 +1,14 @@ + + + + + %d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n + + + + + + + + + diff --git a/env/test/resources/logback.xml b/env/test/resources/logback.xml new file mode 100644 index 0000000..5194993 --- /dev/null +++ b/env/test/resources/logback.xml @@ -0,0 +1,14 @@ + + + + + %d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n + + + + + + + + + diff --git a/kit.edn b/kit.edn new file mode 100644 index 0000000..bcd627b --- /dev/null +++ b/kit.edn @@ -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"}]}} 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..96aabc2 --- /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..22e092c --- /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/migrations/20260218000000-create-kraken-hour.down.sql b/resources/migrations/20260218000000-create-kraken-hour.down.sql new file mode 100644 index 0000000..0d2879b --- /dev/null +++ b/resources/migrations/20260218000000-create-kraken-hour.down.sql @@ -0,0 +1 @@ +DROP TABLE kraken_hour; diff --git a/resources/migrations/20260218000000-create-kraken-hour.up.sql b/resources/migrations/20260218000000-create-kraken-hour.up.sql new file mode 100644 index 0000000..b2c3510 --- /dev/null +++ b/resources/migrations/20260218000000-create-kraken-hour.up.sql @@ -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() +); diff --git a/resources/migrations/20260218100000-create-kraken-day.down.sql b/resources/migrations/20260218100000-create-kraken-day.down.sql new file mode 100644 index 0000000..7a4717f --- /dev/null +++ b/resources/migrations/20260218100000-create-kraken-day.down.sql @@ -0,0 +1 @@ +DROP TABLE kraken_day; diff --git a/resources/migrations/20260218100000-create-kraken-day.up.sql b/resources/migrations/20260218100000-create-kraken-day.up.sql new file mode 100644 index 0000000..7215589 --- /dev/null +++ b/resources/migrations/20260218100000-create-kraken-day.up.sql @@ -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() +); diff --git a/resources/queries.sql b/resources/queries.sql new file mode 100644 index 0000000..08e8ca2 --- /dev/null +++ b/resources/queries.sql @@ -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 diff --git a/resources/system.edn b/resources/system.edn new file mode 100644 index 0000000..bdb8e96 --- /dev/null +++ b/resources/system.edn @@ -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}} diff --git a/src/clj/pmagnus/btcdata/config.clj b/src/clj/pmagnus/btcdata/config.clj new file mode 100644 index 0000000..1ebc616 --- /dev/null +++ b/src/clj/pmagnus/btcdata/config.clj @@ -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)) diff --git a/src/clj/pmagnus/btcdata/core.clj b/src/clj/pmagnus/btcdata/core.clj new file mode 100644 index 0000000..d3e96c6 --- /dev/null +++ b/src/clj/pmagnus/btcdata/core.clj @@ -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))) diff --git a/src/clj/pmagnus/btcdata/kraken/ohlc.clj b/src/clj/pmagnus/btcdata/kraken/ohlc.clj new file mode 100644 index 0000000..e9c1c9f --- /dev/null +++ b/src/clj/pmagnus/btcdata/kraken/ohlc.clj @@ -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)) diff --git a/src/clj/pmagnus/btcdata/kraken/ohlc_daily.clj b/src/clj/pmagnus/btcdata/kraken/ohlc_daily.clj new file mode 100644 index 0000000..659dc19 --- /dev/null +++ b/src/clj/pmagnus/btcdata/kraken/ohlc_daily.clj @@ -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)) diff --git a/src/clj/pmagnus/btcdata/web/controllers/health.clj b/src/clj/pmagnus/btcdata/web/controllers/health.clj new file mode 100644 index 0000000..fe59f94 --- /dev/null +++ b/src/clj/pmagnus/btcdata/web/controllers/health.clj @@ -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"}})) diff --git a/src/clj/pmagnus/btcdata/web/handler.clj b/src/clj/pmagnus/btcdata/web/handler.clj new file mode 100644 index 0000000..6a58441 --- /dev/null +++ b/src/clj/pmagnus/btcdata/web/handler.clj @@ -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)))) diff --git a/src/clj/pmagnus/btcdata/web/middleware/core.clj b/src/clj/pmagnus/btcdata/web/middleware/core.clj new file mode 100644 index 0000000..dbde6bb --- /dev/null +++ b/src/clj/pmagnus/btcdata/web/middleware/core.clj @@ -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))))) diff --git a/src/clj/pmagnus/btcdata/web/middleware/exception.clj b/src/clj/pmagnus/btcdata/web/middleware/exception.clj new file mode 100644 index 0000000..08289f3 --- /dev/null +++ b/src/clj/pmagnus/btcdata/web/middleware/exception.clj @@ -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))}))) diff --git a/src/clj/pmagnus/btcdata/web/middleware/formats.clj b/src/clj/pmagnus/btcdata/web/middleware/formats.clj new file mode 100644 index 0000000..e84f374 --- /dev/null +++ b/src/clj/pmagnus/btcdata/web/middleware/formats.clj @@ -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))))) diff --git a/src/clj/pmagnus/btcdata/web/routes/api.clj b/src/clj/pmagnus/btcdata/web/routes/api.clj new file mode 100644 index 0000000..09cb11c --- /dev/null +++ b/src/clj/pmagnus/btcdata/web/routes/api.clj @@ -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)]) diff --git a/src/clj/pmagnus/btcdata/web/routes/utils.clj b/src/clj/pmagnus/btcdata/web/routes/utils.clj new file mode 100644 index 0000000..733776e --- /dev/null +++ b/src/clj/pmagnus/btcdata/web/routes/utils.clj @@ -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)) diff --git a/src/clj/pmagnus/btcdata/ws/binance.clj b/src/clj/pmagnus/btcdata/ws/binance.clj new file mode 100644 index 0000000..a7239d7 --- /dev/null +++ b/src/clj/pmagnus/btcdata/ws/binance.clj @@ -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 _))))