diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..8fce603 --- /dev/null +++ b/.gitignore @@ -0,0 +1 @@ +data/ diff --git a/files/nats_registration_listener.py b/files/nats_registration_listener.py index 97250a3..e850ddd 100644 --- a/files/nats_registration_listener.py +++ b/files/nats_registration_listener.py @@ -1,12 +1,13 @@ #!/usr/bin/env python3 """ -NATS Registration Listener (fox100 + NetBox hostname/upgrade_cmd lookup + timing + problem counter) -------------------------------------------------------------------------------------------- -- One device GET (status + tags + custom_fields.upgrade_cmd), no duplicate fetch -- If fox100 and upgrade_cmd present -> prepend 3x ASCII BEL to stdout line -- Publish logic for fox100: - * upgrade_cmd empty/absent -> publish with task_name="sot-updater" - * upgrade_cmd present -> publish with task_name="sot-updater-upgradecmd" (+ 3 BELs) +NATS Registration Listener (fox100 + NetBox hostname/action_next lookup + timing + problem counter) +----------------------------------------------------------------------------------------------- +- One device GET (status + tags + custom_fields.action_next), no duplicate fetch +- If fox100 and action_next present -> prepend 3x ASCII BEL to stdout line +- RabbitMQ publish logic is DISABLED for now: + * We still PREPARE the RabbitMQ message body + * We clearly log: "ok, here i will execute " + * But we DO NOT actually POST it to RabbitMQ (publish call commented out) - Keeps: nb_problems counter, timings, iface_id diagnostics, same formatting """ @@ -38,7 +39,7 @@ NB_TOKEN = "7648e4f5ee370cda7834682e61b47c2ee8e95623" # keep as provided NB_TIMEOUT = 3.0 # seconds per HTTP GET # ========================= -# RabbitMQ hardcoded config (immediate publish like rmq-ikeja-pub3.sh without delay) +# RabbitMQ hardcoded config (publish currently disabled) # ========================= RMQ_HOST = "10.210.12.2" RMQ_PORT = 15672 @@ -129,19 +130,29 @@ def extract_fields(obj: Dict[str, Any]): # NetBox lookup (urllib) # ========================= def normalize_mac(mac: str) -> Optional[str]: + """ + Normalize MAC to lowercase colon-separated hex pairs: aa:bb:cc:dd:ee:ff + Accepts common formats: aa-bb-cc-dd-ee-ff, aabbccddeeff, aa:bb:... + """ if not mac or not isinstance(mac, str): return None - s = mac.strip().lower().replace("-", ":") + s = mac.strip().lower() + + # Strip all non-hex chars to handle aabb.ccdd.eeff, aa-bb-..., etc. hex_only = "".join(ch for ch in s if ch in "0123456789abcdef") if len(hex_only) == 12: - return ":join".replace(":", "").join([":".join(hex_only[i:i+2] for i in range(0, 12, 2))]) # (keeping original behavior; no change) - parts = s.split(":") + return ":".join(hex_only[i:i + 2] for i in range(0, 12, 2)) + + # If it's already colon-separated, validate it strictly. + parts = s.replace("-", ":").split(":") if len(parts) == 6 and all(len(p) == 2 and all(c in "0123456789abcdef" for c in p) for p in parts): - return s + return ":".join(parts) + return None -def http_get_json(url: str, params: Optional[Dict[str, str]] = None, headers: Optional[Dict[str, str]] = None, timeout: float = NB_TIMEOUT): +def http_get_json(url: str, params: Optional[Dict[str, str]] = None, headers: Optional[Dict[str, str]] = None, + timeout: float = NB_TIMEOUT): if params: url = f"{url}?{urlencode(params)}" req = Request(url, headers=headers or {}, method="GET") @@ -160,7 +171,9 @@ def http_get_json(url: str, params: Optional[Dict[str, str]] = None, headers: Op # RabbitMQ management API POST helper (basic auth; JSON in/out) -def http_post_json(url: str, payload_obj: Dict[str, Any], user: Optional[str] = None, password: Optional[str] = None, timeout: float = RMQ_TIMEOUT): +# NOTE: currently not used because publishing is disabled, but left in place for later. +def http_post_json(url: str, payload_obj: Dict[str, Any], user: Optional[str] = None, password: Optional[str] = None, + timeout: float = RMQ_TIMEOUT): body = json.dumps(payload_obj).encode("utf-8") headers = {"Content-Type": "application/json"} if user and password: @@ -190,11 +203,13 @@ async def nb_problem(log_status, msg: str): await log_status(msg) -def nb_lookup_device_by_mac(mac: str, log_status) -> Tuple[Optional[str], Optional[int], Optional[int], Optional[str], Optional[Set[str]], Optional[Any]]: +def nb_lookup_device_by_mac(mac: str, log_status) -> Tuple[ + Optional[str], Optional[int], Optional[int], Optional[str], Optional[Set[str]], Optional[Any] +]: """ - Resolve MAC -> (hostname, iface_id, device_id, device_status_value, tag_slugs_set, upgrade_cmd) + Resolve MAC -> (hostname, iface_id, device_id, device_status_value, tag_slugs_set, action_next) - Logs problems for anomalies (mac not found, unassigned, wrong type, iface fetch fail). - - If device detail fetch fails, returns host/id with status/tags/upgrade_cmd as None. + - If device detail fetch fails, returns host/id with status/tags/action_next as None. """ mac_norm = normalize_mac(mac) if not mac_norm: @@ -208,7 +223,8 @@ def nb_lookup_device_by_mac(mac: str, log_status) -> Tuple[Optional[str], Option } # Step 1: MAC lookup - data, code = http_get_json(f"{base}/api/dcim/mac-addresses/", params={"mac_address": mac_norm, "limit": "2"}, headers=h) + data, code = http_get_json(f"{base}/api/dcim/mac-addresses/", params={"mac_address": mac_norm, "limit": "2"}, + headers=h) if code != 200 or not data: asyncio.create_task(nb_problem(log_status, f"[{ts()}] nb: mac query http={code} mac={mac_norm}")) return None, None, None, None, None, None @@ -223,7 +239,9 @@ def nb_lookup_device_by_mac(mac: str, log_status) -> Tuple[Optional[str], Option aoid = rec.get("assigned_object_id") if len(results) > 1: - asyncio.create_task(nb_problem(log_status, f"[{ts()}] nb: multiple mac records mac={mac_norm} iface_id={aoid if aoid is not None else '-'}")) + asyncio.create_task( + nb_problem(log_status, f"[{ts()}] nb: multiple mac records mac={mac_norm} iface_id={aoid if aoid is not None else '-'}") + ) if not aot or aoid is None: asyncio.create_task(nb_problem(log_status, f"[{ts()}] nb: mac unassigned mac={mac_norm}")) @@ -245,7 +263,7 @@ def nb_lookup_device_by_mac(mac: str, log_status) -> Tuple[Optional[str], Option asyncio.create_task(nb_problem(log_status, f"[{ts()}] nb: iface has no device iface_id={aoid}")) return None, aoid, None, None, None, None - # Step 3: Device detail (single fetch for status, tags, custom_fields.upgrade_cmd) + # Step 3: Device detail (single fetch for status, tags, custom_fields.action_next) device, code3 = http_get_json(f"{base}/api/dcim/devices/{dev_id}/", headers=h) if code3 != 200 or not device: # treat as "no extra info" @@ -260,9 +278,9 @@ def nb_lookup_device_by_mac(mac: str, log_status) -> Tuple[Optional[str], Option tag_slugs.add(slug) cf = device.get("custom_fields") or {} - upgrade_cmd = cf.get("upgrade_cmd") + action_next = cf.get("action_next") - return host, aoid, dev_id, (status_val if isinstance(status_val, str) else None), tag_slugs, upgrade_cmd + return host, aoid, dev_id, (status_val if isinstance(status_val, str) else None), tag_slugs, action_next # ========================= @@ -337,64 +355,70 @@ async def main(): host_suffix = "" action_suffix = "" # kept; no longer used for tag action - bell_prefix = "" # ASCII BEL when upgrade_cmd present (3x) + bell_prefix = "" # ASCII BEL when action_next present (3x) netbox_time_ms = 0.0 - # NOTE: no tag check anymore; behavior depends on upgrade_cmd only + # Behavior depends on action_next presence only (for fox100) if product == "fox100": nb_start = time.perf_counter() try: - host, iface_id, dev_id, status_val, tag_slugs, upgrade_cmd = nb_lookup_device_by_mac(mac=mac, log_status=log_status) + host, iface_id, dev_id, status_val, tag_slugs, action_next = nb_lookup_device_by_mac(mac=mac, log_status=log_status) if host: host_suffix = f" host={host}" if iface_id and not host: host_suffix += f" iface_id={iface_id}" - # NEW LOGIC: - # - If host known: - # * upgrade_cmd empty/absent -> publish with task_name="sot-updater" - # * upgrade_cmd present -> publish with task_name="sot-updater-upgradecmd" and ring 3 BELs + # If host known, prepare RabbitMQ message and log it, but DO NOT publish. if host: try: rmq_url = f"http://{RMQ_HOST}:{RMQ_PORT}/api/exchanges/{RMQ_VHOST}/{RMQ_EXCHANGE_WORK}/publish" - has_upgrade_cmd = False + has_action_next = False try: - # consider non-empty string or any truthy value as "present" - if isinstance(upgrade_cmd, str): - has_upgrade_cmd = len(upgrade_cmd.strip()) > 0 + if isinstance(action_next, str): + has_action_next = len(action_next.strip()) > 0 else: - has_upgrade_cmd = bool(upgrade_cmd) + has_action_next = bool(action_next) except Exception: - has_upgrade_cmd = False + has_action_next = False - task_name = "sot-updater-upgradecmd" if has_upgrade_cmd else "sot-updater" + # Task naming: keep it simple and explicit. + task_name = "sot-updater-actionnext" if has_action_next else "sot-updater" + + rmq_payload_obj = { + "inscope_device": host, + "task_name": task_name, + # IMPORTANT: we are not "executing" it here; we only carry it for later scripts. + "action_next": action_next, + } rmq_body = { "properties": { "content_type": "application/json" }, "routing_key": RMQ_ROUTING_KEY, - "payload": json.dumps({ - "inscope_device": host, - "task_name": task_name, - }), + "payload": json.dumps(rmq_payload_obj), "payload_encoding": "string", } - resp, code = http_post_json(rmq_url, rmq_body, user=RMQ_USER, password=RMQ_PASS, timeout=RMQ_TIMEOUT) - if code != 200: - await log_status(f"[{ts()}] rmq: publish immediate http={code} host={host}") - else: - routed = False - try: - routed = bool((resp or {}).get("routed", False)) - except Exception: - routed = False - if not routed: - await log_status(f"[{ts()}] rmq: publish immediate routed=false host={host}") - # Bell behavior: ring 3x BEL only if upgrade_cmd present - if has_upgrade_cmd: + # (b) Explicit log of what would be executed + await log_status(f"[{ts()}] ok, here i will execute {json.dumps(rmq_body, ensure_ascii=False)} url={rmq_url}") + + # (c) Publishing is disabled for now (leave code in place, commented out) + # resp, code = http_post_json(rmq_url, rmq_body, user=RMQ_USER, password=RMQ_PASS, timeout=RMQ_TIMEOUT) + # if code != 200: + # await log_status(f"[{ts()}] rmq: publish immediate http={code} host={host}") + # else: + # routed = False + # try: + # routed = bool((resp or {}).get("routed", False)) + # except Exception: + # routed = False + # if not routed: + # await log_status(f"[{ts()}] rmq: publish immediate routed=false host={host}") + + # Bell behavior: ring 3x BEL only if action_next present + if has_action_next: bell_prefix = "\a" * 3 except Exception as e: @@ -416,7 +440,7 @@ async def main(): line += f" subject={msg.subject}" async with print_lock: - # Prepend BEL only when we had upgrade_cmd present (3x) + # Prepend BEL only when we had action_next present (3x) sys.stdout.write(bell_prefix + line + "\n") sys.stdout.flush() @@ -452,3 +476,4 @@ if __name__ == "__main__": asyncio.run(main()) except KeyboardInterrupt: pass +