diff --git a/docker-compose.yml b/docker-compose.yml index 5ca20fd..027a07b 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,16 +1,25 @@ +networks: + shared: + name: shared + external: true + services: btcdata: build: . restart: unless-stopped ports: - - "4101:4101" + - "${DOCKER_BTCDATA_PORT}:${DOCKER_BTCDATA_PORT}" depends_on: - frankfurter + networks: + - default + - shared environment: - BTCDATA_PORT: "4101" - JDBC_URL: "jdbc:postgresql://postgres:5432/btcprod?user=postgres&password=ratata,123" - CORS_ORIGIN: "http://localhost:4041" - FRANKFURTER_URL: "http://frankfurter:8080" + BTCDATA_PORT: "${DOCKER_BTCDATA_PORT}" + JDBC_URL: "${DOCKER_JDBC_URL}" + CORS_ORIGIN: "${DOCKER_CORS_ORIGIN}" + FRANKFURTER_URL: "${DOCKER_FRANKFURTER_URL}" + STRIKE_API_KEY: "${STRIKE_API_KEY}" extra_hosts: - "postgres:host-gateway" diff --git a/resources/migrations/20260306000000-create-strike-price.down.sql b/resources/migrations/20260306000000-create-strike-price.down.sql new file mode 100644 index 0000000..af84bfa --- /dev/null +++ b/resources/migrations/20260306000000-create-strike-price.down.sql @@ -0,0 +1 @@ +DROP TABLE strike_price; diff --git a/resources/migrations/20260306000000-create-strike-price.up.sql b/resources/migrations/20260306000000-create-strike-price.up.sql new file mode 100644 index 0000000..c8a98d2 --- /dev/null +++ b/resources/migrations/20260306000000-create-strike-price.up.sql @@ -0,0 +1,5 @@ +CREATE TABLE strike_price ( + id BIGSERIAL PRIMARY KEY, + rates JSONB NOT NULL, + recorded_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); diff --git a/resources/queries.sql b/resources/queries.sql index 08e8ca2..f88487b 100644 --- a/resources/queries.sql +++ b/resources/queries.sql @@ -42,3 +42,11 @@ SET open = EXCLUDED.open, -- :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 + +-- :name insert-strike-price! :! :n +-- :doc Insert a Strike ticker snapshot +INSERT INTO strike_price (rates) VALUES (:rates) + +-- :name get-latest-strike-price :? :1 +-- :doc Get the most recent Strike ticker snapshot +SELECT rates, recorded_at FROM strike_price ORDER BY id DESC LIMIT 1 diff --git a/resources/system.edn b/resources/system.edn index 71c0d8c..57195cf 100644 --- a/resources/system.edn +++ b/resources/system.edn @@ -60,4 +60,8 @@ {:query-fn #ig/ref :db.sql/query-fn} :frankfurter/rates - {:url #or [#env FRANKFURTER_URL "http://localhost:8080"]}} + {:url #or [#env FRANKFURTER_URL "http://localhost:8080"]} + + :strike/ticker + {:query-fn #ig/ref :db.sql/query-fn + :api-key #env STRIKE_API_KEY}} diff --git a/src/clj/pmagnus/btcdata/core.clj b/src/clj/pmagnus/btcdata/core.clj index 022b73a..0d64070 100644 --- a/src/clj/pmagnus/btcdata/core.clj +++ b/src/clj/pmagnus/btcdata/core.clj @@ -18,7 +18,8 @@ ;; Pollers [pmagnus.btcdata.kraken.ohlc] [pmagnus.btcdata.kraken.ohlc-daily] - [pmagnus.btcdata.frankfurter.rates]) + [pmagnus.btcdata.frankfurter.rates] + [pmagnus.btcdata.strike.ticker]) (:gen-class)) (defonce system (atom nil)) diff --git a/src/clj/pmagnus/btcdata/frankfurter/rates.clj b/src/clj/pmagnus/btcdata/frankfurter/rates.clj index 92ef03e..870533c 100644 --- a/src/clj/pmagnus/btcdata/frankfurter/rates.clj +++ b/src/clj/pmagnus/btcdata/frankfurter/rates.clj @@ -36,10 +36,16 @@ "Fetch immediately, then poll every 60 minutes." [client base-url rates-atom running?] (future - (try - (poll! client base-url rates-atom) - (catch Exception e - (log/error e "Frankfurter initial fetch failed"))) + (loop [retries 5] + (let [ok? (try + (poll! client base-url rates-atom) + true + (catch Exception e + (log/error e (str "Frankfurter fetch failed, retrying in 5s (" retries " left)")) + false))] + (when (and (not ok?) @running? (pos? retries)) + (Thread/sleep 5000) + (recur (dec retries))))) (while @running? (Thread/sleep 3600000) (when @running? diff --git a/src/clj/pmagnus/btcdata/strike/ticker.clj b/src/clj/pmagnus/btcdata/strike/ticker.clj new file mode 100644 index 0000000..c870426 --- /dev/null +++ b/src/clj/pmagnus/btcdata/strike/ticker.clj @@ -0,0 +1,77 @@ +(ns pmagnus.btcdata.strike.ticker + (:require + [clojure.data.json :as json] + [clojure.string :as str] + [clojure.tools.logging :as log] + [integrant.core :as ig]) + (:import + [java.net URI] + [java.net.http HttpClient HttpRequest HttpResponse$BodyHandlers] + [org.postgresql.util PGobject])) + +(defn- ->jsonb + "Wrap a JSON string as a PGobject with type jsonb." + [^String s] + (doto (PGobject.) + (.setType "jsonb") + (.setValue s))) + +(defn- fetch-rates + "GET /v1/rates/ticker from Strike. Returns the raw JSON response body string." + [^HttpClient client ^String api-key] + (let [request (-> (HttpRequest/newBuilder) + (.uri (URI. "https://api.strike.me/v1/rates/ticker")) + (.header "Accept" "application/json") + (.header "Authorization" (str "Bearer " api-key)) + (.GET) + (.build)) + resp (.send client request (HttpResponse$BodyHandlers/ofString)) + body (.body resp)] + body)) + +(defn- format-rates + "Build a compact log string from the parsed rates array." + [rates] + (->> rates + (map (fn [{:keys [sourceCurrency targetCurrency amount]}] + (str sourceCurrency "/" targetCurrency " " amount))) + (str/join " "))) + +(defn- poll! + "Fetch rates and insert into the database." + [client api-key query-fn] + (let [body (fetch-rates client api-key) + rates (json/read-str body :key-fn keyword)] + (query-fn :insert-strike-price! {:rates (->jsonb body)}) + (log/info "Strike:" (format-rates rates)))) + +(defn- start-poll-loop! + "Fetch immediately, then poll every 10 seconds." + [client api-key query-fn running?] + (future + (try + (poll! client api-key query-fn) + (catch Exception e + (log/error e "Strike initial fetch failed"))) + (while @running? + (Thread/sleep 10000) + (when @running? + (try + (poll! client api-key query-fn) + (catch Exception e + (log/error e "Strike poll failed"))))))) + +(defmethod ig/init-key :strike/ticker + [_ {:keys [query-fn api-key]}] + (log/info "Starting Strike ticker poller") + (let [client (HttpClient/newHttpClient) + running? (atom true) + fut (start-poll-loop! client api-key query-fn running?)] + {:running? running? + :future fut})) + +(defmethod ig/halt-key! :strike/ticker + [_ {:keys [running? future]}] + (log/info "Stopping Strike ticker poller") + (reset! running? false) + (future-cancel future))