From aafc3c9e3e104d50b8914808296dada4aee03e71 Mon Sep 17 00:00:00 2001 From: Mario Nachbaur Date: Tue, 18 Aug 2026 19:03:47 +0200 Subject: [PATCH] read from dotenv ot getenv --- packages/ingestor/src/ingestor/__init__.py | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/packages/ingestor/src/ingestor/__init__.py b/packages/ingestor/src/ingestor/__init__.py index 034e8a3..79b2758 100644 --- a/packages/ingestor/src/ingestor/__init__.py +++ b/packages/ingestor/src/ingestor/__init__.py @@ -2,6 +2,7 @@ import asyncio import json import logging import sqlite3 +from os import getenv import dotenv from atproto_jetstream import Jetstream, JetstreamCommitEvent, JetstreamOptions @@ -12,7 +13,7 @@ logger = logging.getLogger(__name__) async def ingest_jetstream(config: dict[str, str | None]): endpoint = "wss://jetstream1.us-east.bsky.network/subscribe" options = JetstreamOptions( - endpoint=config.get("JETSTREAM_URL") or endpoint, + endpoint=read_config(config, key="JETSTREAM_URL", default=endpoint), wanted_collections=["at.ligo.*"], compress=True, ) @@ -56,7 +57,7 @@ def handle_commit( (prefix, did), ) else: - logger.debug(f"creating or updating {prefix} for {did}") + print(f"creating or updating {prefix} for {did}") if commit.record["$type"] != type: return content = json.dumps(commit.record) @@ -71,10 +72,14 @@ def handle_commit( def get_database(config: dict[str, str | None]) -> sqlite3.Connection | None: - database_name = config.get("FLASK_KEYVAL_DB_URL") or "keyval.db" + database_name = read_config(config, key="FLASK_KEYVAL_DB_URL", default="keyval.db") return sqlite3.connect(database_name) +def read_config(config: dict[str, str | None], key: str, default: str) -> str: + return config.get(key) or getenv(key) or default + + async def async_main(config: dict[str, str | None]): try: await ingest_jetstream(config) -- 2.51.2