Files
server-26/drb-c2-core/app/internal/node_sweeper.py
T

90 lines
3.4 KiB
Python

import asyncio
from datetime import datetime, timezone, timedelta
from app.config import settings
from app.internal.logger import logger
from app.internal import firestore as fstore
SWEEP_INTERVAL = 90 # seconds — matches node_offline_threshold; no gain in checking faster
async def sweeper_loop():
"""
Periodically check for nodes that haven't checked in recently
and mark them offline in Firestore.
"""
logger.info("Node sweeper started.")
while True:
await asyncio.sleep(SWEEP_INTERVAL)
try:
await _sweep()
except Exception as e:
logger.error(f"Sweeper error: {e}")
async def _sweep():
threshold = datetime.now(timezone.utc) - timedelta(seconds=settings.node_offline_threshold)
def _query():
from app.internal.firestore import db
return [
doc.to_dict()
for doc in db.collection("nodes").stream()
]
nodes = await asyncio.to_thread(_query)
for node in nodes:
status = node.get("status", "offline")
if status == "offline":
continue
last_seen_raw = node.get("last_seen")
if not last_seen_raw:
continue
# last_seen may be a Firestore Timestamp, a datetime, or an ISO string
if isinstance(last_seen_raw, str):
last_seen = datetime.fromisoformat(last_seen_raw)
else:
last_seen = last_seen_raw
if last_seen.tzinfo is None:
last_seen = last_seen.replace(tzinfo=timezone.utc)
if last_seen < threshold:
node_id = node.get("node_id")
await fstore.doc_update("nodes", node_id, {"status": "offline"})
logger.info(f"Node {node_id} marked offline (last seen: {last_seen.isoformat()})")
from app.routers.tokens import release_token
await release_token(node_id)
continue
# Check for expired system overrides (only for fixed nodes with timeout enforced)
override_timeout_raw = node.get("override_timeout_at")
enforce_timeout = node.get("enforce_override_timeout", True)
node_type = node.get("node_type", "fixed")
if override_timeout_raw and enforce_timeout and node_type != "portable":
if isinstance(override_timeout_raw, str):
override_timeout = datetime.fromisoformat(override_timeout_raw)
else:
override_timeout = override_timeout_raw
if override_timeout.tzinfo is None:
override_timeout = override_timeout.replace(tzinfo=timezone.utc)
if datetime.now(timezone.utc) > override_timeout:
node_id = node.get("node_id")
assigned_system_id = node.get("assigned_system_id")
logger.info(f"Node {node_id} override has expired. Reverting to system {assigned_system_id}.")
# Push the original assigned config if it exists
if assigned_system_id:
system_doc = await fstore.doc_get("systems", assigned_system_id)
if system_doc:
from app.internal.mqtt_handler import mqtt_handler
mqtt_handler.push_config(node_id, system_doc)
await fstore.doc_update("nodes", node_id, {
"is_overridden": False,
"override_system_id": None,
"override_timeout_at": None,
})