#!/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...")