Files
kubernetes/scripts/setup_jetstream.py
2026-07-18 12:00:33 +05:30

167 lines
6.0 KiB
Python

#!/usr/bin/env python3
"""
Setup JetStream stream and consumer for NATS
Run this after NATS is deployed and running
"""
import asyncio
import nats
import os
import sys
async def setup_jetstream():
"""Configure JetStream with two streams (DELIVERIES, RIDER) and per-subject consumers."""
nats_url = os.getenv("NATS_URL", "nats://nats.workolik.com:4222")
nats_user = os.getenv("NATS_USER", "admin")
nats_password = os.getenv("NATS_PASSWORD", "package@321#")
# Stream definitions
streams = {
"DELIVERIES": {
"subjects": [
"api.v1.deliveries.createdeliveries",
"api.v1.deliveries.updatedelivery",
"api.v2.deliveries.createdeliverylog",
],
},
"RIDER": {
"subjects": [
"api.v2.partners.createriderlog",
"api.v2.partners.createbreaklog",
"api.v2.partners.updatebreaklog",
],
},
"ORDERS": {
"subjects": [
"api.v1.mob.orders.createorder",
],
},
"PRODUCTS": {
"subjects": [
"api.v1.web.products.create",
],
},
"CUSTOMERS": {
"subjects": [
"api.v1.mob.customers.login",
"api.v1.mob.customers.create",
],
"retention": "work" # Special handling for Login queue: delete immediately after ack
},
}
# Per-subject durable consumers
consumers = {
"DELIVERIES": {
"api.v1.deliveries.createdeliveries": "deliveries_createdeliveries",
"api.v1.deliveries.updatedelivery": "deliveries_updatedelivery",
"api.v2.deliveries.createdeliverylog": "deliveries_createdeliverylog",
},
"RIDER": {
"api.v2.partners.createriderlog": "rider_createriderlog",
"api.v2.partners.createbreaklog": "rider_createbreaklog",
"api.v2.partners.updatebreaklog": "rider_updatebreaklog",
},
"ORDERS": {
"api.v1.mob.orders.createorder": "orders_createorder",
},
"PRODUCTS": {
"api.v1.web.products.create": "products_create",
},
"CUSTOMERS": {
"api.v1.mob.customers.login": "customers_login",
"api.v1.mob.customers.create": "customers_create",
},
}
try:
print(f"Connecting to NATS at {nats_url}...")
nc = await nats.connect(
servers=[nats_url],
user=nats_user,
password=nats_password,
)
print("✅ Connected to NATS")
js = nc.jetstream()
# Create / recreate streams
for stream_name, cfg in streams.items():
try:
info = await js.stream_info(stream_name)
print(f"⚠️ Stream '{stream_name}' already exists with subjects={info.config.subjects}, updating...")
# Determine retention policy
retention_policy = cfg.get("retention", "limits")
await js.update_stream(
name=stream_name,
subjects=cfg["subjects"],
storage="memory",
retention=retention_policy,
max_age=24 * 60 * 60,
max_msgs=50000,
max_bytes=512 * 1024 * 1024,
)
print(f"✅ Stream '{stream_name}' updated")
except Exception as e:
if "not found" in str(e).lower() or "404" in str(e).lower():
print(f"Creating stream '{stream_name}'...")
# Determine retention policy (use 'limits' by default, 'work' for queues)
retention_policy = cfg.get("retention", "limits")
await js.add_stream(
name=stream_name,
subjects=cfg["subjects"],
storage="memory",
retention=retention_policy,
max_age=24 * 60 * 60,
max_msgs=50000,
max_bytes=512 * 1024 * 1024,
)
print(f"✅ Stream '{stream_name}' created")
else:
print(f"⚠️ Could not inspect stream '{stream_name}': {e}")
# Create durable consumers per subject
print("\nConfiguring consumers...")
for stream_name, subject_map in consumers.items():
for subject, durable in subject_map.items():
try:
print(f"Creating consumer '{durable}' on stream '{stream_name}' for subject '{subject}'...")
await js.add_consumer(
stream_name,
durable_name=durable,
filter_subject=subject,
ack_policy="explicit",
deliver_policy="all",
max_deliver=5,
ack_wait=30,
)
print(f"✅ Consumer '{durable}' created")
except Exception as e:
if "already in use" in str(e).lower() or "already exists" in str(e).lower():
print(f"⚠️ Consumer '{durable}' already exists, skipping...")
else:
raise
print("\n✅ JetStream setup complete!")
print(" Streams:")
for name, cfg in streams.items():
print(f" - {name}: {', '.join(cfg['subjects'])}")
print(" Consumers:")
for stream_name, subject_map in consumers.items():
for subject, durable in subject_map.items():
print(f" - {durable}: stream={stream_name}, subject={subject}")
await nc.close()
sys.exit(0)
except Exception as e:
print(f"❌ Error: {e}")
sys.exit(1)
if __name__ == "__main__":
asyncio.run(setup_jetstream())