financial_aggregator/server/aggregator/aggregator.py
2026-05-30 20:26:01 -05:00

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...")