diff --git a/app/ingestor.py b/app/ingestor.py index bd7043a..56909b5 100644 --- a/app/ingestor.py +++ b/app/ingestor.py @@ -148,22 +148,26 @@ async def start_nats_consumer(): nc = await nats.connect(NATS_URLS) js = nc.jetstream() - # Create stream if not exists + # Create stream if not exists; if it exists, UPDATE its subject list so + # newly added event types (e.g. events.camera) are actually covered. + from nats.js.api import StreamConfig + stream_cfg = StreamConfig( + name=NATS_STREAM, + subjects=[ + "events.gdelt", "events.rss", "events.social", + "events.earthquake", "events.disaster", "events.weather", + "events.fire", "events.satellite", "events.new", "events.alert", + "events.camera", + ], + retention=nats.js.api.RetentionPolicy.LIMITS, + max_msgs=1_000_000, + ) try: - await js.add_stream( - name=NATS_STREAM, - subjects=[ - "events.gdelt", "events.rss", "events.social", - "events.earthquake", "events.disaster", "events.weather", - "events.fire", "events.satellite", "events.new", "events.alert", - "events.camera", - ], - retention=nats.js.api.RetentionPolicy.LIMITS, - max_msgs=1_000_000, - ) + await js.add_stream(stream_cfg) logger.info("Created NATS stream %s", NATS_STREAM) except Exception: - logger.debug("Stream %s already exists", NATS_STREAM) + await js.update_stream(stream_cfg) + logger.info("Updated NATS stream %s subjects", NATS_STREAM) # Create durable consumer sub = await js.pull_subscribe(