Ingestor: update JetStream stream subjects on startup so events.camera gets covered (add_stream silently no-ops when stream exists)
All checks were successful
build-and-deploy / build (push) Successful in 1m52s

This commit is contained in:
Sirius DevOps 2026-08-24 21:19:20 -04:00
parent 456712fc32
commit b1c611d4c4

View file

@ -148,22 +148,26 @@ async def start_nats_consumer():
nc = await nats.connect(NATS_URLS) nc = await nats.connect(NATS_URLS)
js = nc.jetstream() 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: try:
await js.add_stream( await js.add_stream(stream_cfg)
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,
)
logger.info("Created NATS stream %s", NATS_STREAM) logger.info("Created NATS stream %s", NATS_STREAM)
except Exception: 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 # Create durable consumer
sub = await js.pull_subscribe( sub = await js.pull_subscribe(