"""Long-running ingester: pulls OSINT sources into NATS JetStream and consumes NATS messages into PostgreSQL. Runs forever (one process, two concurrent tasks): * producer loop — fetch RSS / GDELT / USGS on an interval, publish to NATS * consumer loop — pull from the NATS `events.>` stream, write to Postgres Env (all optional, 12-factor): RSS_URL comma-separated feed URLs to poll (default: none) INGEST_INTERVAL seconds between producer cycles (default: 300) GDELT_QUERY GDELT search term (default: "") INGEST_EARTHQUAKES "1"/"true" to enable USGS feed (default: "1") """ from __future__ import annotations import asyncio import logging import os from pathlib import Path sys_path = str(Path(__file__).parent) import sys sys.path.insert(0, sys_path) from config import NATS_URL, FIRMS_INTERVAL, FIRMS_DATASET, AISSTREAM_IN_INGEST # noqa: E402 from sources import ingest_rss_feed, ingest_gdelt, ingest_earthquakes # noqa: E402 from fire_sources import ingest_fires # noqa: E402 from ingestor import ingest_event, start_nats_consumer # noqa: E402 logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s: %(message)s") logger = logging.getLogger("osint.ingester") RSS_URLS = [u.strip() for u in os.getenv("RSS_URL", "").split(",") if u.strip()] INTERVAL = int(os.getenv("INGEST_INTERVAL", "300")) GDELT_QUERY = os.getenv("GDELT_QUERY", "") ENABLE_QUAKES = os.getenv("INGEST_EARTHQUAKES", "1").lower() in ("1", "true", "yes") ENABLE_FIRES = os.getenv("INGEST_FIRES", "1").lower() in ("1", "true", "yes") NATS_STREAM = "events" async def producer_loop() -> None: """Periodically fetch external sources and publish to NATS.""" while True: try: for url in RSS_URLS: try: n = await ingest_rss_feed(url) logger.info("RSS %s -> %d items", url, n) except Exception: # noqa: BLE001 logger.exception("RSS fetch failed: %s", url) try: g = await ingest_gdelt(query=GDELT_QUERY, max_articles=50) logger.info("GDELT -> %d articles", g) except Exception: # noqa: BLE001 logger.exception("GDELT fetch failed") if ENABLE_QUAKES: try: q = await ingest_earthquakes() logger.info("USGS -> %d events", q) except Exception: # noqa: BLE001 logger.exception("USGS fetch failed") except Exception: # noqa: BLE001 logger.exception("producer cycle error") await asyncio.sleep(INTERVAL) async def consumer_loop() -> None: """Continuously pull NATS messages and persist to Postgres.""" _, sub = await start_nats_consumer() while True: try: msgs = await sub.fetch(100, timeout=5) except Exception: # noqa: BLE001 (TimeoutError / no messages) continue for msg in msgs: try: import json data = json.loads(msg.data) await ingest_event(data) await msg.ack() except Exception: # noqa: BLE001 logger.exception("failed to process message") async def fire_loop() -> None: """Poll NASA FIRMS active fires on the ~15-minute cadence. Runs as its own task so its slower cadence (FIRMS_INTERVAL, default 900s) is independent of the RSS/GDELT/USGS producer cycle (INGEST_INTERVAL). """ logger.info("fire loop starting (interval=%ss, dataset=%s)", FIRMS_INTERVAL, FIRMS_DATASET) while True: try: n = await ingest_fires() logger.info("FIRMS -> %d hotspots", n) except Exception: # noqa: BLE001 — keep the loop alive across transient failures logger.exception("FIRMS fetch failed") await asyncio.sleep(FIRMS_INTERVAL) async def main() -> None: logger.info( "ingester starting (rss=%d feeds, gdelt_q=%r, quakes=%s, fires=%s, interval=%ss)", len(RSS_URLS), GDELT_QUERY, ENABLE_QUAKES, ENABLE_FIRES, INTERVAL, ) tasks: list[asyncio.Task] = [] if ENABLE_FIRES: # Fire ingest only starts once FIRMS_MAP_KEY is set (ingest_fires logs # and idles otherwise). tasks.append(asyncio.create_task(fire_loop())) if AISSTREAM_IN_INGEST: from ais_stream import run_ais_worker # noqa: E402 tasks.append(asyncio.create_task(run_ais_worker())) await asyncio.gather(producer_loop(), consumer_loop(), *tasks) if __name__ == "__main__": asyncio.run(main())