211 lines
7.1 KiB
Python
211 lines
7.1 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
Financial Aggregator Service
|
|
Polls Yahoo Finance and Finnhub APIs, publishes to MQTT broker.
|
|
Runs every 3 minutes to match ESP32 refresh interval.
|
|
"""
|
|
|
|
import asyncio
|
|
import json
|
|
import logging
|
|
import os
|
|
import socket
|
|
import sys
|
|
import time
|
|
from typing import Any
|
|
|
|
import aiohttp
|
|
import paho.mqtt.client as mqtt
|
|
|
|
logging.basicConfig(
|
|
level=logging.INFO,
|
|
format="%(asctime)s - %(levelname)s - %(message)s",
|
|
handlers=[logging.StreamHandler(sys.stdout)],
|
|
)
|
|
logger = logging.getLogger(__name__)
|
|
|
|
MQTT_HOST = os.getenv("MQTT_HOST", "mosquitto")
|
|
MQTT_PORT = int(os.getenv("MQTT_PORT", "1883"))
|
|
MQTT_USER = os.getenv("MQTT_USER", "esp32")
|
|
MQTT_PASS = os.getenv("MQTT_PASS", "your_password")
|
|
MQTT_TOPIC = os.getenv("MQTT_TOPIC", "finance/prices")
|
|
|
|
FINNHUB_API_KEY = os.getenv("FINNHUB_API_KEY", "d74t3bpr01qg1eo6pln0")
|
|
|
|
POLL_INTERVAL = 180
|
|
|
|
SYMBOLS = {
|
|
"btc": {"source": "finnhub", "symbol": "BINANCE:BTCUSDT"},
|
|
"gold": {"source": "yahoo", "symbol": "GC=F"},
|
|
"silver": {"source": "yahoo", "symbol": "SI=F"},
|
|
}
|
|
|
|
YAHOO_USER_AGENT = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36"
|
|
|
|
|
|
class MQTTClient:
|
|
def __init__(self):
|
|
self.client = mqtt.Client(client_id=f"aggregator-{socket.gethostname()}")
|
|
self.client.username_pw_set(MQTT_USER, MQTT_PASS)
|
|
self.client.on_connect = self._on_connect
|
|
self.client.on_disconnect = self._on_disconnect
|
|
self.connected = False
|
|
self._connect()
|
|
|
|
def _connect(self):
|
|
retries = 0
|
|
max_retries = 10
|
|
while retries < max_retries:
|
|
try:
|
|
self.client.connect(MQTT_HOST, MQTT_PORT, keepalive=60)
|
|
self.client.loop_start()
|
|
logger.info(f"Connected to MQTT broker at {MQTT_HOST}:{MQTT_PORT}")
|
|
return
|
|
except Exception as e:
|
|
retries += 1
|
|
logger.warning(
|
|
f"MQTT connection attempt {retries}/{max_retries} failed: {e}"
|
|
)
|
|
time.sleep(5)
|
|
logger.error("Failed to connect to MQTT broker after max retries")
|
|
|
|
def _on_connect(self, client, userdata, flags, rc):
|
|
if rc == 0:
|
|
self.connected = True
|
|
logger.info("MQTT connection established")
|
|
else:
|
|
logger.error(f"MQTT connection failed with code: {rc}")
|
|
|
|
def _on_disconnect(self, client, userdata, rc):
|
|
self.connected = False
|
|
logger.warning(f"MQTT disconnected with code: {rc}")
|
|
if rc != 0:
|
|
logger.info("Attempting to reconnect...")
|
|
time.sleep(5)
|
|
self._connect()
|
|
|
|
def publish(self, topic: str, payload: dict) -> bool:
|
|
if not self.connected:
|
|
logger.warning("MQTT not connected, attempting reconnect")
|
|
self._connect()
|
|
if not self.connected:
|
|
return False
|
|
|
|
result = self.client.publish(topic, json.dumps(payload), qos=1, retain=True)
|
|
if result.rc == mqtt.MQTT_ERR_SUCCESS:
|
|
logger.debug(f"Published to {topic}: {payload}")
|
|
return True
|
|
else:
|
|
logger.error(f"Failed to publish: {result}")
|
|
return False
|
|
|
|
|
|
async def fetch_yahoo_price(
|
|
session: aiohttp.ClientSession, symbol: str
|
|
) -> float | None:
|
|
url = f"https://query1.finance.yahoo.com/v8/finance/chart/{symbol}?interval=1d&range=1d"
|
|
headers = {"User-Agent": YAHOO_USER_AGENT}
|
|
try:
|
|
async with session.get(
|
|
url, headers=headers, timeout=aiohttp.ClientTimeout(total=15)
|
|
) as resp:
|
|
if resp.status != 200:
|
|
logger.warning(f"Yahoo API returned {resp.status} for {symbol}")
|
|
return None
|
|
data = await resp.json()
|
|
quote = data.get("chart", {}).get("result", [{}])[0].get("meta", {})
|
|
price = quote.get("regularMarketPrice")
|
|
if price:
|
|
logger.debug(f"Yahoo {symbol}: ${price}")
|
|
return float(price)
|
|
logger.warning(f"No price in Yahoo response for {symbol}")
|
|
return None
|
|
except asyncio.TimeoutError:
|
|
logger.warning(f"Yahoo API timeout for {symbol}")
|
|
return None
|
|
except Exception as e:
|
|
logger.error(f"Yahoo API error for {symbol}: {e}")
|
|
return None
|
|
|
|
|
|
async def fetch_finnhub_price(
|
|
session: aiohttp.ClientSession, symbol: str
|
|
) -> float | None:
|
|
url = f"https://finnhub.io/api/v1/quote?symbol={symbol}&token={FINNHUB_API_KEY}"
|
|
try:
|
|
async with session.get(url, timeout=aiohttp.ClientTimeout(total=15)) as resp:
|
|
if resp.status != 200:
|
|
logger.warning(f"Finnhub API returned {resp.status} for {symbol}")
|
|
return None
|
|
data = await resp.json()
|
|
price = data.get("c")
|
|
if price and price > 0:
|
|
logger.debug(f"Finnhub {symbol}: ${price}")
|
|
return float(price)
|
|
logger.warning(f"No valid price in Finnhub response for {symbol}")
|
|
return None
|
|
except asyncio.TimeoutError:
|
|
logger.warning(f"Finnhub API timeout for {symbol}")
|
|
return None
|
|
except Exception as e:
|
|
logger.error(f"Finnhub API error for {symbol}: {e}")
|
|
return None
|
|
|
|
|
|
async def fetch_all_prices() -> dict[str, float]:
|
|
prices = {}
|
|
async with aiohttp.ClientSession() as session:
|
|
tasks = []
|
|
for key, config in SYMBOLS.items():
|
|
if config["source"] == "yahoo":
|
|
tasks.append((key, fetch_yahoo_price(session, config["symbol"])))
|
|
else:
|
|
tasks.append((key, fetch_finnhub_price(session, config["symbol"])))
|
|
|
|
results = await asyncio.gather(*[t[1] for t in tasks], return_exceptions=True)
|
|
for (key, _), result in zip(tasks, results):
|
|
if isinstance(result, Exception):
|
|
logger.error(f"Exception fetching {key}: {result}")
|
|
elif result is not None:
|
|
prices[key] = result
|
|
else:
|
|
logger.warning(f"Failed to fetch price for {key}")
|
|
|
|
return prices
|
|
|
|
|
|
async def main_loop():
|
|
mqtt_client = MQTTClient()
|
|
logger.info("Starting financial aggregator service")
|
|
logger.info(f"Poll interval: {POLL_INTERVAL} seconds")
|
|
logger.info(f"MQTT topic: {MQTT_TOPIC}")
|
|
|
|
while True:
|
|
try:
|
|
logger.info("Fetching prices...")
|
|
prices = await fetch_all_prices()
|
|
|
|
if prices:
|
|
payload = {
|
|
"timestamp": int(time.time()),
|
|
**{k: v for k, v in prices.items()},
|
|
}
|
|
|
|
if mqtt_client.publish(MQTT_TOPIC, payload):
|
|
logger.info(f"Published {len(prices)} prices: {prices}")
|
|
else:
|
|
logger.error("Failed to publish prices")
|
|
else:
|
|
logger.warning("No prices fetched, skipping publish")
|
|
|
|
except Exception as e:
|
|
logger.error(f"Error in main loop: {e}", exc_info=True)
|
|
|
|
await asyncio.sleep(POLL_INTERVAL)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
try:
|
|
asyncio.run(main_loop())
|
|
except KeyboardInterrupt:
|
|
logger.info("Shutting down...")
|