mirror of
https://github.com/agessaman/meshcore-bot.git
synced 2026-10-06 06:17:17 +00:00
Region scopes had to be hand-edited in config.ini, and the radio's own default
region could not be set from the bot at all. Both now live on the Radio page.
Node Settings gains a Default Region Scope field beside the path hash size,
which is the firmware's own setting — NodePrefs.default_scope_name and
default_scope_key, the region the radio falls back to for any send the bot does
not scope itself. Firmware that has no such command is reported as not having
answered rather than shown as an empty field, and a stored key that is not the
stored name's hash is flagged, because the radio routes by the key and the name
beside it is only a label. A build-flag default stores the bare name while
hashing the '#' form, so the comparison normalizes first.
Clearing it does not go through the meshcore library, because none of its reset
paths can: set_default_flood_scope(None) raises on len(None), "" earns
ILLEGAL_ARG from the firmware, and "*" only works because the frame is padded by
character count rather than by the name it wrote, landing one byte short of the
length the firmware reads as "a scope follows". The firmware's own contract for
clearing is a bare CMD_SET_DEFAULT_FLOOD_SCOPE, so that is what goes out. The
same padding bug is why a device scope name must be ASCII: a multi-byte
character displaces the transport key out of its field.
A separate Region Scopes card covers the bot's side — [Channels] flood_scopes as
an explicit choice between "reply whatever the scope" and an allowlist, and
outgoing_flood_scope_override. It writes config.ini, queues a hot reload, then
polls that reload and reports what the bot did with it, including saying plainly
when nothing picked it up. Per-channel flood_scope.<channel> entries are listed
read-only, since they are the reason a channel can ignore the default.
The two settings interact in a way worth stating on the page: once the bot sends
a scoped message it restores with set_flood_scope("*"), which leaves the radio in
forced-unscoped mode, so the device default stops applying until the bot scopes
another send.
Also from building it:
- flood_scopes left in [Bot] is surfaced rather than hidden. CommandManager
still honours it when [Channels] is empty, so a page that ignored it would
show "replies to every scope" while the bot enforced an allowlist — and
saving "off" hands the allowlist straight back to [Bot], which the banner now
says instead of claiming the old key stops mattering.
- outgoing_flood_scope_override = none is read as global flood on the send path,
as it already was everywhere else. send_channel_message tested the raw value
against a fixed tuple, so the lowercase spelling became the region "#none".
- Scope normalization moved to modules/flood_scope.py so the viewer, a separate
process, can validate a typed name without importing the bot's command
machinery. CommandManager delegates to it.
- Dark mode named only .border, so a .border-start divider drew Bootstrap's
light #dee2e6 onto a dark card. All four directional utilities now match.
2126 lines
97 KiB
Python
2126 lines
97 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
Message scheduler functionality for the MeshCore Bot
|
|
Handles scheduled messages and timing
|
|
"""
|
|
|
|
import asyncio
|
|
import datetime
|
|
import hashlib
|
|
import json
|
|
import os
|
|
import re
|
|
import socket
|
|
import sqlite3
|
|
import threading
|
|
import time
|
|
import uuid
|
|
from pathlib import Path
|
|
from typing import Any, Optional
|
|
|
|
from apscheduler.schedulers.background import BackgroundScheduler
|
|
from apscheduler.triggers.date import DateTrigger
|
|
from meshcore.events import EventType
|
|
|
|
from .flood_scope import scope_key_hex
|
|
from .maintenance import MaintenanceRunner
|
|
from .models import CHANNEL_REGIONAL_FLOOD_SCOPE_BODY_OVERHEAD
|
|
from .scheduled_message_cron import (
|
|
is_valid_legacy_hhmm,
|
|
parse_schedule_key,
|
|
parse_scheduled_message_value,
|
|
)
|
|
from .security_utils import validate_external_url
|
|
from .utils import (
|
|
decode_escape_sequences,
|
|
format_keyword_response_with_placeholders,
|
|
get_config_timezone,
|
|
)
|
|
|
|
_CHANNEL_OPERATION_TYPES = ('add', 'remove')
|
|
_RADIO_OPERATION_TYPES = (
|
|
'radio_reboot',
|
|
'radio_connect',
|
|
'radio_disconnect',
|
|
'firmware_read',
|
|
'firmware_write',
|
|
'radio_params_read',
|
|
'radio_params_write',
|
|
'radio_advert',
|
|
)
|
|
_CONFIG_OPERATION_TYPES = ('config_reload',)
|
|
|
|
# Let short host stalls or CPU pressure delay a scheduled transmission without
|
|
# dropping it. Coalescing prevents a long suspension from replaying multiple
|
|
# stale occurrences onto the mesh when the process resumes.
|
|
SCHEDULE_MISFIRE_GRACE_SECONDS = 300
|
|
|
|
|
|
class MessageScheduler:
|
|
"""Manages scheduled messages and timing"""
|
|
|
|
def __init__(self, bot):
|
|
self.bot = bot
|
|
self.logger = bot.logger
|
|
self._claim_owner_host = socket.gethostname()
|
|
self._claim_owner_pid = os.getpid()
|
|
# Unique to this scheduler/process start; retained with the PID to make
|
|
# ownership auditable and prevent one scheduler instance finalizing a
|
|
# claim made by another instance in the same process.
|
|
self._claim_owner_boot_id = uuid.uuid4().hex
|
|
self.scheduled_messages = {}
|
|
self.scheduler_thread = None
|
|
self._apscheduler: Optional[BackgroundScheduler] = None
|
|
self.last_channel_ops_check_time = 0
|
|
self.last_message_queue_check_time = 0
|
|
self.last_radio_ops_check_time = 0
|
|
self._data_retention_interval_seconds = 86400 # 24 hours
|
|
self._data_retention_startup_delay_seconds = 60
|
|
# Enforce retention shortly after startup, independently of the nightly
|
|
# email timer. After that first run, the scheduler resets this timestamp
|
|
# and retains the normal daily cadence.
|
|
self.last_data_retention_run = (
|
|
time.time()
|
|
- self._data_retention_interval_seconds
|
|
+ self._data_retention_startup_delay_seconds
|
|
)
|
|
self.last_nightly_email_time = time.time() # don't send immediately on startup
|
|
self.last_db_backup_run = 0
|
|
self.last_log_rotation_check_time = 0
|
|
self.maintenance = MaintenanceRunner(bot, get_current_time=self.get_current_time)
|
|
db_manager = getattr(bot, 'db_manager', None)
|
|
db_path = getattr(db_manager, 'db_path', None)
|
|
if db_manager is not None and isinstance(db_path, (str, os.PathLike)):
|
|
self._recover_interrupted_operations()
|
|
|
|
def get_current_time(self):
|
|
"""Get current time in configured timezone"""
|
|
tz, _ = get_config_timezone(self.bot.config, self.logger)
|
|
return datetime.datetime.now(tz)
|
|
|
|
def _shutdown_apscheduler_if_running(self) -> None:
|
|
"""Stop APScheduler if it exists and is running (idempotent, no spurious errors)."""
|
|
if self._apscheduler is None:
|
|
return
|
|
if not getattr(self._apscheduler, "running", False):
|
|
return
|
|
try:
|
|
self._apscheduler.shutdown(wait=False)
|
|
except Exception as e:
|
|
self.logger.debug("Error shutting down scheduler: %s", e)
|
|
|
|
def setup_scheduled_messages(self):
|
|
"""Setup scheduled messages from config using APScheduler."""
|
|
# Stop and recreate the APScheduler to avoid duplicate jobs on reload
|
|
self._shutdown_apscheduler_if_running()
|
|
tz, _ = get_config_timezone(self.bot.config, self.logger)
|
|
self._apscheduler = BackgroundScheduler(
|
|
timezone=tz,
|
|
job_defaults={
|
|
'misfire_grace_time': SCHEDULE_MISFIRE_GRACE_SECONDS,
|
|
'coalesce': True,
|
|
},
|
|
)
|
|
self.scheduled_messages.clear()
|
|
|
|
if self.bot.config.has_section('Scheduled_Messages'):
|
|
self.logger.info("Found Scheduled_Messages section")
|
|
for schedule_key, message_info in self.bot.config.items('Scheduled_Messages'):
|
|
self.logger.info(f"Processing scheduled message: '{schedule_key}' -> '{message_info}'")
|
|
try:
|
|
parsed = parse_schedule_key(schedule_key, tz)
|
|
if parsed.trigger is None:
|
|
self.logger.warning(
|
|
f"Invalid schedule '{schedule_key}' for scheduled message: {message_info}"
|
|
)
|
|
continue
|
|
|
|
if parsed.is_deprecated_hhmm:
|
|
hh = int(schedule_key[:2])
|
|
mm = int(schedule_key[2:])
|
|
cron_suggestion = f"{mm} {hh} * * *"
|
|
self.logger.warning(
|
|
"Scheduled_Messages key %r uses deprecated HHMM daily format; "
|
|
"migrate to 5-field cron (minute hour dom mon dow), e.g. %r. "
|
|
"HHMM support will be removed in a future release.",
|
|
schedule_key,
|
|
cron_suggestion,
|
|
)
|
|
|
|
channel, message, scope = parse_scheduled_message_value(message_info)
|
|
message = decode_escape_sequences(message)
|
|
|
|
if self._has_command_placeholders(message):
|
|
interval = self._min_fire_interval_seconds(parsed.trigger, tz)
|
|
floor = self.MIN_COMMAND_PLACEHOLDER_INTERVAL_SECONDS
|
|
if interval is not None and interval < floor:
|
|
self.logger.error(
|
|
"Scheduled message %r uses a {cmd:...} placeholder but fires "
|
|
"every %.0fs; the minimum is %ds because each firing spends "
|
|
"airtime. Not scheduled: %s",
|
|
schedule_key, interval, floor, message,
|
|
)
|
|
continue
|
|
|
|
job_id = "schedmsg_" + hashlib.sha256(
|
|
f"{schedule_key}\0{channel}\0{scope or ''}\0{message}".encode()
|
|
).hexdigest()[:24]
|
|
|
|
self._apscheduler.add_job(
|
|
self.send_scheduled_message,
|
|
parsed.trigger,
|
|
args=[channel, message],
|
|
kwargs={"schedule_key": schedule_key, "scope": scope},
|
|
id=job_id,
|
|
replace_existing=True,
|
|
)
|
|
self.scheduled_messages[schedule_key] = (
|
|
channel,
|
|
message,
|
|
parsed.display_label,
|
|
scope,
|
|
)
|
|
scope_note = f" scope={scope}" if scope else ""
|
|
self.logger.info(
|
|
f"Scheduled message: {parsed.display_label} -> {channel}{scope_note}: {message}"
|
|
)
|
|
except ValueError:
|
|
self.logger.warning(f"Invalid scheduled message format: {message_info}")
|
|
except Exception as e:
|
|
self.logger.warning(f"Error setting up scheduled message '{schedule_key}': {e}")
|
|
|
|
self._apscheduler.start()
|
|
self.logger.info(f"APScheduler started with {len(self.scheduled_messages)} scheduled message(s)")
|
|
|
|
self._setup_device_mode_scheduler_jobs()
|
|
|
|
# Setup interval-based advertising
|
|
self.setup_interval_advertising()
|
|
|
|
def setup_interval_advertising(self):
|
|
"""Setup interval-based advertising from config"""
|
|
try:
|
|
advert_interval_hours = self.bot.config.getint('Bot', 'advert_interval_hours', fallback=0)
|
|
if advert_interval_hours > 0:
|
|
self.logger.info(f"Setting up interval-based advertising every {advert_interval_hours} hours")
|
|
# Initialize bot's last advert time to now to prevent immediate advert if not already set
|
|
if not hasattr(self.bot, 'last_advert_time') or self.bot.last_advert_time is None:
|
|
self.bot.last_advert_time = time.time()
|
|
else:
|
|
self.logger.info("Interval-based advertising disabled (advert_interval_hours = 0)")
|
|
except Exception as e:
|
|
self.logger.warning(f"Error setting up interval advertising: {e}")
|
|
|
|
def _setup_device_mode_scheduler_jobs(self) -> None:
|
|
"""One-shot jobs for auto_manage_contacts=device: firmware autoadd + favourite hygiene."""
|
|
if self._apscheduler is None:
|
|
return
|
|
if self.bot.config.get('Bot', 'auto_manage_contacts', fallback='device').lower() != 'device':
|
|
return
|
|
try:
|
|
delay_fw = max(0, self.bot.config.getint('Bot', 'device_mode_firmware_delay_seconds', fallback=30))
|
|
delay_p1 = max(0, self.bot.config.getint('Bot', 'device_mode_favourite_pass1_delay_seconds', fallback=90))
|
|
delay_p2 = max(0, self.bot.config.getint('Bot', 'device_mode_favourite_pass2_delay_seconds', fallback=180))
|
|
base = self.get_current_time()
|
|
self._apscheduler.add_job(
|
|
self._device_mode_firmware_job_sync,
|
|
trigger=DateTrigger(run_date=base + datetime.timedelta(seconds=delay_fw)),
|
|
id='device_mode_firmware_autoadd',
|
|
replace_existing=True,
|
|
)
|
|
self._apscheduler.add_job(
|
|
self._device_mode_favourite_pass1_job_sync,
|
|
trigger=DateTrigger(run_date=base + datetime.timedelta(seconds=delay_p1)),
|
|
id='device_mode_favourite_pass1',
|
|
replace_existing=True,
|
|
)
|
|
self._apscheduler.add_job(
|
|
self._device_mode_favourite_pass2_job_sync,
|
|
trigger=DateTrigger(run_date=base + datetime.timedelta(seconds=delay_p2)),
|
|
id='device_mode_favourite_pass2',
|
|
replace_existing=True,
|
|
)
|
|
self.logger.info(
|
|
'Scheduled device-mode jobs: firmware +%ss, favourite pass1 +%ss, pass2 +%ss',
|
|
delay_fw,
|
|
delay_p1,
|
|
delay_p2,
|
|
)
|
|
except Exception as e:
|
|
self.logger.warning('Could not schedule device-mode contact jobs: %s', e)
|
|
|
|
def _run_async_on_main_loop(self, coro: Any, timeout: float = 300.0) -> None:
|
|
"""Run async coroutine on bot main loop from APScheduler thread (same pattern as send_scheduled_message)."""
|
|
import asyncio
|
|
|
|
loop = getattr(self.bot, 'main_event_loop', None)
|
|
if not loop or not loop.is_running():
|
|
# Close the coroutine we were handed. Dropping it unawaited leaks it and
|
|
# emits "coroutine ... was never awaited" RuntimeWarning.
|
|
coro.close()
|
|
self.logger.warning('No running main_event_loop — skipping device-mode scheduled job')
|
|
return
|
|
|
|
try:
|
|
future = asyncio.run_coroutine_threadsafe(coro, loop)
|
|
except RuntimeError as e:
|
|
# This job runs on an APScheduler thread, so the loop can stop between
|
|
# the is_running() check above and this submit. The coroutine never got
|
|
# scheduled, so close it here too.
|
|
coro.close()
|
|
self.logger.warning('Event loop gone during device-mode job: %s', e)
|
|
return
|
|
|
|
try:
|
|
future.result(timeout=timeout)
|
|
except RuntimeError as e:
|
|
self.logger.warning('Event loop gone during device-mode job: %s', e)
|
|
except Exception as e:
|
|
self.logger.error('Device-mode scheduled job failed: %s', e)
|
|
|
|
async def _device_mode_firmware_coro(self) -> None:
|
|
await self.bot.repeater_manager.apply_device_mode_firmware_preferences()
|
|
|
|
async def _device_mode_favourite_pass1_coro(self) -> None:
|
|
await self.bot.repeater_manager.sync_device_mode_favourites_pass1()
|
|
|
|
async def _device_mode_favourite_pass2_coro(self) -> None:
|
|
await self.bot.repeater_manager.sync_device_mode_favourites_pass2()
|
|
|
|
def _device_mode_firmware_job_sync(self) -> None:
|
|
if self.bot.config.get('Bot', 'auto_manage_contacts', fallback='device').lower() != 'device':
|
|
self.logger.debug('Skipping device_mode_firmware job — not device mode')
|
|
return
|
|
self._run_async_on_main_loop(self._device_mode_firmware_coro(), timeout=120.0)
|
|
|
|
def _device_mode_favourite_pass1_job_sync(self) -> None:
|
|
if self.bot.config.get('Bot', 'auto_manage_contacts', fallback='device').lower() != 'device':
|
|
return
|
|
self._run_async_on_main_loop(self._device_mode_favourite_pass1_coro(), timeout=600.0)
|
|
|
|
def _device_mode_favourite_pass2_job_sync(self) -> None:
|
|
if self.bot.config.get('Bot', 'auto_manage_contacts', fallback='device').lower() != 'device':
|
|
return
|
|
self._run_async_on_main_loop(self._device_mode_favourite_pass2_coro(), timeout=600.0)
|
|
|
|
def _is_valid_time_format(self, time_str: str) -> bool:
|
|
"""Validate deprecated legacy time format (HHMM). Prefer cron in config keys."""
|
|
return is_valid_legacy_hhmm(time_str)
|
|
|
|
def _scheduled_message_stagger_seconds(self, schedule_key: str) -> float:
|
|
"""Deterministic delay in [0, max) so simultaneous cron jobs do not stack on the radio."""
|
|
max_s = self.bot.config.getfloat(
|
|
"Bot", "scheduled_message_max_stagger_seconds", fallback=1.5
|
|
)
|
|
if max_s <= 0 or not (schedule_key or "").strip():
|
|
return 0.0
|
|
digest = hashlib.sha256(schedule_key.encode("utf-8")).digest()
|
|
slot = int.from_bytes(digest[:4], "big") / (2**32)
|
|
return float(slot * max_s)
|
|
|
|
def send_scheduled_message(
|
|
self,
|
|
channel: str,
|
|
message: str,
|
|
schedule_key: str = "",
|
|
scope: str | None = None,
|
|
):
|
|
"""Send a scheduled message (synchronous wrapper for schedule library)"""
|
|
if self.bot.is_radio_zombie:
|
|
self.logger.warning("send_scheduled_message suppressed — radio is in zombie state")
|
|
return
|
|
if self.bot.is_radio_offline:
|
|
self.logger.warning("send_scheduled_message suppressed — radio is offline (repeated send timeouts)")
|
|
return
|
|
|
|
current_time = self.get_current_time()
|
|
scope_note = f" [{scope}]" if scope else ""
|
|
self.logger.info(
|
|
f"📅 Sending scheduled message at {current_time.strftime('%H:%M:%S')} "
|
|
f"to {channel}{scope_note}: {message}"
|
|
)
|
|
|
|
import asyncio
|
|
|
|
# Use the main event loop if available, otherwise create a new one
|
|
# This prevents deadlock when the main loop is already running
|
|
if hasattr(self.bot, 'main_event_loop') and self.bot.main_event_loop and self.bot.main_event_loop.is_running():
|
|
# Schedule coroutine in the running main event loop
|
|
future = asyncio.run_coroutine_threadsafe(
|
|
self._send_scheduled_message_async(
|
|
channel, message, schedule_key=schedule_key, scope=scope
|
|
),
|
|
self.bot.main_event_loop,
|
|
)
|
|
# Wait for completion (with timeout to prevent indefinite blocking)
|
|
try:
|
|
future.result(timeout=60) # 60 second timeout
|
|
self.bot._record_send_success()
|
|
except RuntimeError as e:
|
|
self.logger.warning("Event loop gone during scheduled message: %s", e)
|
|
except Exception as e:
|
|
self.logger.error(f"Error sending scheduled message: {type(e).__name__}: {e}")
|
|
self.bot._record_send_failure(scheduler=self)
|
|
else:
|
|
# Fallback: create a temporary event loop and close it when done
|
|
loop = asyncio.new_event_loop()
|
|
try:
|
|
loop.run_until_complete(
|
|
self._send_scheduled_message_async(
|
|
channel, message, schedule_key=schedule_key, scope=scope
|
|
)
|
|
)
|
|
finally:
|
|
loop.close()
|
|
|
|
async def _get_mesh_info(self) -> dict[str, Any]:
|
|
"""Get mesh network information for scheduled messages"""
|
|
info = {
|
|
'total_contacts': 0,
|
|
'total_repeaters': 0,
|
|
'total_companions': 0,
|
|
'total_roomservers': 0,
|
|
'total_sensors': 0,
|
|
'recent_activity_24h': 0,
|
|
'new_companions_7d': 0,
|
|
'new_repeaters_7d': 0,
|
|
'new_roomservers_7d': 0,
|
|
'new_sensors_7d': 0,
|
|
'total_contacts_30d': 0,
|
|
'total_repeaters_30d': 0,
|
|
'total_companions_30d': 0,
|
|
'total_roomservers_30d': 0,
|
|
'total_sensors_30d': 0
|
|
}
|
|
|
|
try:
|
|
# Get contact statistics from repeater manager if available
|
|
if hasattr(self.bot, 'repeater_manager'):
|
|
try:
|
|
stats = await self.bot.repeater_manager.get_contact_statistics()
|
|
if stats:
|
|
info['total_contacts'] = stats.get('total_heard', 0)
|
|
by_role = stats.get('by_role', {})
|
|
info['total_repeaters'] = by_role.get('repeater', 0)
|
|
info['total_companions'] = by_role.get('companion', 0)
|
|
info['total_roomservers'] = by_role.get('roomserver', 0)
|
|
info['total_sensors'] = by_role.get('sensor', 0)
|
|
info['recent_activity_24h'] = stats.get('recent_activity', 0)
|
|
except Exception as e:
|
|
self.logger.debug(f"Error getting stats from repeater_manager: {e}")
|
|
|
|
# Fallback to device contacts if repeater manager stats not available
|
|
if info['total_contacts'] == 0 and hasattr(self.bot, 'meshcore') and hasattr(self.bot.meshcore, 'contacts'):
|
|
info['total_contacts'] = len(self.bot.meshcore.contacts)
|
|
|
|
# Count repeaters and companions
|
|
if hasattr(self.bot, 'repeater_manager'):
|
|
for contact_data in self.bot.meshcore.contacts.values():
|
|
if self.bot.repeater_manager._is_repeater_device(contact_data):
|
|
info['total_repeaters'] += 1
|
|
else:
|
|
info['total_companions'] += 1
|
|
|
|
# Get recent activity from message_stats if available
|
|
if info['recent_activity_24h'] == 0:
|
|
try:
|
|
with self.bot.db_manager.connection() as conn:
|
|
cursor = conn.cursor()
|
|
# Check if message_stats table exists
|
|
cursor.execute('''
|
|
SELECT name FROM sqlite_master
|
|
WHERE type='table' AND name='message_stats'
|
|
''')
|
|
if cursor.fetchone():
|
|
cutoff_time = int(time.time()) - (24 * 60 * 60)
|
|
cursor.execute('''
|
|
SELECT COUNT(DISTINCT sender_id)
|
|
FROM message_stats
|
|
WHERE timestamp >= ? AND is_dm = 0
|
|
''', (cutoff_time,))
|
|
result = cursor.fetchone()
|
|
if result:
|
|
info['recent_activity_24h'] = result[0]
|
|
except Exception as e:
|
|
self.logger.debug("Error querying message_stats: %s", e)
|
|
|
|
# Calculate new devices in last 7 days (matching web viewer logic)
|
|
# Query devices first heard in the last 7 days, grouped by role
|
|
# Also calculate devices active in last 30 days (last_heard)
|
|
try:
|
|
with self.bot.db_manager.connection() as conn:
|
|
cursor = conn.cursor()
|
|
# Check if complete_contact_tracking table exists
|
|
cursor.execute('''
|
|
SELECT name FROM sqlite_master
|
|
WHERE type='table' AND name='complete_contact_tracking'
|
|
''')
|
|
if cursor.fetchone():
|
|
# Get new devices by role (first_heard in last 7 days)
|
|
# Use role field for matching (more reliable than device_type)
|
|
cursor.execute('''
|
|
SELECT role, COUNT(DISTINCT public_key) as count
|
|
FROM complete_contact_tracking
|
|
WHERE first_heard >= datetime('now', '-7 days')
|
|
AND role IS NOT NULL AND role != ''
|
|
GROUP BY role
|
|
''')
|
|
for row in cursor.fetchall():
|
|
role = (row[0] or '').lower()
|
|
count = row[1] or 0
|
|
|
|
if role == 'companion':
|
|
info['new_companions_7d'] = count
|
|
elif role == 'repeater':
|
|
info['new_repeaters_7d'] = count
|
|
elif role == 'roomserver':
|
|
info['new_roomservers_7d'] = count
|
|
elif role == 'sensor':
|
|
info['new_sensors_7d'] = count
|
|
|
|
# Get total contacts active in last 30 days (last_heard)
|
|
cursor.execute('''
|
|
SELECT COUNT(DISTINCT public_key) as count
|
|
FROM complete_contact_tracking
|
|
WHERE last_heard >= datetime('now', '-30 days')
|
|
''')
|
|
result = cursor.fetchone()
|
|
if result:
|
|
info['total_contacts_30d'] = result[0] or 0
|
|
|
|
# Get devices active in last 30 days by role (last_heard)
|
|
cursor.execute('''
|
|
SELECT role, COUNT(DISTINCT public_key) as count
|
|
FROM complete_contact_tracking
|
|
WHERE last_heard >= datetime('now', '-30 days')
|
|
AND role IS NOT NULL AND role != ''
|
|
GROUP BY role
|
|
''')
|
|
for row in cursor.fetchall():
|
|
role = (row[0] or '').lower()
|
|
count = row[1] or 0
|
|
|
|
if role == 'companion':
|
|
info['total_companions_30d'] = count
|
|
elif role == 'repeater':
|
|
info['total_repeaters_30d'] = count
|
|
elif role == 'roomserver':
|
|
info['total_roomservers_30d'] = count
|
|
elif role == 'sensor':
|
|
info['total_sensors_30d'] = count
|
|
except Exception as e:
|
|
self.logger.debug(f"Error getting new device counts or 30-day activity: {e}")
|
|
|
|
except Exception as e:
|
|
self.logger.debug(f"Error getting mesh info: {e}")
|
|
|
|
return info
|
|
|
|
# {cmd:<command> [args]} — run a command and substitute its reply text.
|
|
# Non-greedy and brace-free inside, matching the placeholder limits elsewhere.
|
|
_COMMAND_PLACEHOLDER_RE = re.compile(r"\{cmd:([^{}]+)\}")
|
|
|
|
# Floor on how often a schedule containing {cmd:...} may fire. Every firing is a
|
|
# transmission on a shared medium, and a command placeholder makes it trivial to
|
|
# write a cron that airs several times an hour. Deliberately not configurable.
|
|
MIN_COMMAND_PLACEHOLDER_INTERVAL_SECONDS = 900
|
|
|
|
@staticmethod
|
|
def _min_fire_interval_seconds(trigger, tz, samples: int = 12) -> Optional[float]:
|
|
"""Smallest gap between consecutive firings of *trigger*, in seconds.
|
|
|
|
Sampled rather than derived, so uneven crons are measured by their tightest
|
|
gap: ``0,1 * * * *`` is a 60-second schedule, not a half-hourly one.
|
|
|
|
Returns None when the trigger has no future firings to compare.
|
|
"""
|
|
now = datetime.datetime.now(tz)
|
|
previous = trigger.get_next_fire_time(None, now)
|
|
if previous is None:
|
|
return None
|
|
|
|
smallest = None
|
|
for _ in range(samples):
|
|
nxt = trigger.get_next_fire_time(
|
|
previous, previous + datetime.timedelta(microseconds=1)
|
|
)
|
|
if nxt is None:
|
|
break
|
|
gap = (nxt - previous).total_seconds()
|
|
if gap > 0 and (smallest is None or gap < smallest):
|
|
smallest = gap
|
|
previous = nxt
|
|
return smallest
|
|
|
|
def _has_command_placeholders(self, message: str) -> bool:
|
|
return bool(self._COMMAND_PLACEHOLDER_RE.search(message))
|
|
|
|
async def _expand_command_placeholders(self, message: str, channel: str) -> str:
|
|
"""Replace each {cmd:...} with the command's output.
|
|
|
|
A placeholder whose command is unknown, disabled, admin-only or failing
|
|
expands to an empty string rather than leaving the raw ``{cmd:...}`` text
|
|
on the air. Command output is never re-scanned, so a reply that happens to
|
|
contain ``{cmd:...}`` cannot cause recursion.
|
|
"""
|
|
timeout = self.bot.config.getfloat(
|
|
'Bot', 'scheduled_command_timeout_seconds', fallback=30.0
|
|
)
|
|
|
|
out = []
|
|
last = 0
|
|
for match in self._COMMAND_PLACEHOLDER_RE.finditer(message):
|
|
out.append(message[last:match.start()])
|
|
spec = match.group(1).strip()
|
|
try:
|
|
rendered = await self.bot.command_manager.render_command_output(
|
|
spec, channel=channel, timeout=timeout
|
|
)
|
|
except Exception as e:
|
|
self.logger.warning(
|
|
"Error rendering scheduled command placeholder %r: %s", spec, e
|
|
)
|
|
rendered = None
|
|
if rendered:
|
|
self.logger.info("Scheduled message rendered {cmd:%s}", spec)
|
|
else:
|
|
self.logger.warning(
|
|
"Scheduled message placeholder {cmd:%s} produced nothing; omitted", spec
|
|
)
|
|
out.append(rendered or "")
|
|
last = match.end()
|
|
out.append(message[last:])
|
|
return "".join(out)
|
|
|
|
def _has_mesh_info_placeholders(self, message: str) -> bool:
|
|
"""Check if message contains mesh info placeholders"""
|
|
placeholders = [
|
|
'{total_contacts}', '{total_repeaters}', '{total_companions}',
|
|
'{total_roomservers}', '{total_sensors}', '{recent_activity_24h}',
|
|
'{new_companions_7d}', '{new_repeaters_7d}', '{new_roomservers_7d}', '{new_sensors_7d}',
|
|
'{total_contacts_30d}', '{total_repeaters_30d}', '{total_companions_30d}',
|
|
'{total_roomservers_30d}', '{total_sensors_30d}',
|
|
# Legacy placeholders for backward compatibility
|
|
'{repeaters}', '{companions}'
|
|
]
|
|
return any(placeholder in message for placeholder in placeholders)
|
|
|
|
def _channel_body_budget(self, scope: str | None) -> int:
|
|
"""UTF-8 byte budget for one channel message body.
|
|
|
|
Mirrors BaseCommand.get_max_message_length: channel sends are framed as
|
|
"<username>: <body>", and a regional flood scope costs extra header bytes.
|
|
"""
|
|
username = ""
|
|
try:
|
|
self_info = getattr(getattr(self.bot, "meshcore", None), "self_info", None)
|
|
if isinstance(self_info, dict):
|
|
username = self_info.get("name") or self_info.get("user_name") or ""
|
|
elif self_info is not None:
|
|
username = getattr(self_info, "name", "") or getattr(self_info, "user_name", "")
|
|
except Exception: # noqa: BLE001 - budget must never break a send
|
|
username = ""
|
|
if not isinstance(username, str) or not username:
|
|
try:
|
|
username = self.bot.config.get("Bot", "bot_name", fallback="") or ""
|
|
except Exception: # noqa: BLE001 - budget must never break a send
|
|
username = ""
|
|
# A stubbed or misconfigured source can hand back a non-string; fall back to
|
|
# the most conservative budget rather than raising inside the send path.
|
|
if not isinstance(username, str):
|
|
username = ""
|
|
|
|
budget = 160 - len(username.encode("utf-8")) - 2
|
|
if (scope or "").strip():
|
|
budget -= CHANNEL_REGIONAL_FLOOD_SCOPE_BODY_OVERHEAD
|
|
return max(budget, 32)
|
|
|
|
def _effective_send_scope(self, channel: str, scope: str | None) -> str | None:
|
|
"""The scope the send will actually use, not just the one on the schedule.
|
|
|
|
send_channel_message resolves an unset scope from ``flood_scope.<channel>`` and
|
|
then ``outgoing_flood_scope_override``. Budgeting on the raw schedule scope
|
|
alone would size chunks for a global send and overshoot by the regional header
|
|
once the sender adds it.
|
|
"""
|
|
if (scope or "").strip():
|
|
return scope
|
|
try:
|
|
resolved = self.bot.command_manager.resolve_channel_send_scope(
|
|
scope=None, channel=channel
|
|
)
|
|
if (resolved or "").strip():
|
|
return resolved
|
|
override = self.bot.config.get(
|
|
"Channels", "outgoing_flood_scope_override", fallback=""
|
|
)
|
|
return override if (override or "").strip() else None
|
|
except Exception: # noqa: BLE001 - budgeting must never break a send
|
|
# Unknown means assume regional, which only ever makes chunks smaller.
|
|
return "#unknown"
|
|
|
|
@staticmethod
|
|
def _split_to_budget(text: str, budget: int) -> list[str]:
|
|
"""Split *text* into chunks of at most *budget* UTF-8 bytes, on line breaks
|
|
where possible so a rendered command's lines are not cut mid-sentence."""
|
|
if len(text.encode("utf-8")) <= budget:
|
|
return [text]
|
|
|
|
chunks: list[str] = []
|
|
current = ""
|
|
for line in text.split("\n"):
|
|
candidate = f"{current}\n{line}" if current else line
|
|
if len(candidate.encode("utf-8")) <= budget:
|
|
current = candidate
|
|
continue
|
|
if current:
|
|
chunks.append(current)
|
|
current = ""
|
|
# A single line over budget still has to go out; cut it on a character
|
|
# boundary that keeps the encoded length within the limit.
|
|
while len(line.encode("utf-8")) > budget:
|
|
cut = budget
|
|
while cut > 0 and len(line[:cut].encode("utf-8")) > budget:
|
|
cut -= 1
|
|
chunks.append(line[:cut])
|
|
line = line[cut:]
|
|
current = line
|
|
if current:
|
|
chunks.append(current)
|
|
return [c for c in chunks if c]
|
|
|
|
async def _send_scheduled_message_async(
|
|
self,
|
|
channel: str,
|
|
message: str,
|
|
*,
|
|
schedule_key: str = "",
|
|
scope: str | None = None,
|
|
):
|
|
"""Send a scheduled message (async implementation)"""
|
|
stagger = self._scheduled_message_stagger_seconds(schedule_key)
|
|
if stagger > 0:
|
|
self.logger.debug(
|
|
"Scheduled message stagger %.2fs (schedule_key=%r)", stagger, schedule_key
|
|
)
|
|
await asyncio.sleep(stagger)
|
|
|
|
# Command placeholders first: their output may itself contain mesh info
|
|
# placeholders, which the pass below then resolves.
|
|
if self._has_command_placeholders(message):
|
|
message = await self._expand_command_placeholders(message, channel)
|
|
if not message.strip():
|
|
self.logger.warning(
|
|
"Scheduled message for %s is empty after expanding command "
|
|
"placeholders; nothing sent", channel
|
|
)
|
|
return
|
|
|
|
# Check if message contains mesh info placeholders
|
|
if self._has_mesh_info_placeholders(message):
|
|
try:
|
|
mesh_info = await self._get_mesh_info()
|
|
# Use shared formatting function (message=None for scheduled messages)
|
|
try:
|
|
message = format_keyword_response_with_placeholders(
|
|
message,
|
|
message=None, # No message object for scheduled messages
|
|
bot=self.bot,
|
|
mesh_info=mesh_info
|
|
)
|
|
self.logger.debug("Replaced mesh info placeholders in scheduled message")
|
|
except (KeyError, ValueError) as e:
|
|
self.logger.warning(f"Error replacing placeholders in scheduled message: {e}. Sending message as-is.")
|
|
except Exception as e:
|
|
self.logger.warning(f"Error fetching mesh info for scheduled message: {e}. Sending message as-is.")
|
|
|
|
import asyncio as _asyncio
|
|
send_timeout = self.bot.config.getint('Bot', 'send_timeout_seconds', fallback=30)
|
|
|
|
# A {cmd:...} placeholder can expand to more than one message's worth of text,
|
|
# and send_channel_message does not split. Chunk to the RF body budget so a
|
|
# long rendered reply airs as several messages instead of failing at the device.
|
|
effective_scope = self._effective_send_scope(channel, scope)
|
|
chunks = self._split_to_budget(message, self._channel_body_budget(effective_scope))
|
|
if len(chunks) > 1:
|
|
self.logger.info(
|
|
"Scheduled message for %s split into %d chunks to fit the RF budget",
|
|
channel, len(chunks),
|
|
)
|
|
await _asyncio.wait_for(
|
|
self.bot.command_manager.send_channel_messages_chunked(
|
|
channel, chunks, skip_user_rate_limit=True, scope=scope
|
|
),
|
|
timeout=send_timeout * len(chunks),
|
|
)
|
|
return
|
|
|
|
await _asyncio.wait_for(
|
|
self.bot.command_manager.send_channel_message(
|
|
channel, message, skip_user_rate_limit=True, scope=scope
|
|
),
|
|
timeout=send_timeout,
|
|
)
|
|
|
|
def start(self):
|
|
"""Start the scheduler in a separate thread"""
|
|
self.scheduler_thread = threading.Thread(target=self.run_scheduler, daemon=True)
|
|
self.scheduler_thread.start()
|
|
|
|
def join(self, timeout: float = 5.0) -> None:
|
|
"""Wait for the scheduler thread to finish and stop APScheduler (e.g. during shutdown)."""
|
|
self._shutdown_apscheduler_if_running()
|
|
if self.scheduler_thread and self.scheduler_thread.is_alive():
|
|
self.scheduler_thread.join(timeout=timeout)
|
|
if self.scheduler_thread.is_alive():
|
|
self.logger.debug("Scheduler thread did not finish within %s s", timeout)
|
|
|
|
def run_scheduler(self):
|
|
"""Run the scheduler in a separate thread"""
|
|
self.logger.info("Scheduler thread started")
|
|
last_log_time = 0
|
|
last_feed_poll_time = 0
|
|
last_job_count = 0
|
|
last_job_log_time = 0
|
|
|
|
while self.bot.keep_running:
|
|
current_time = self.get_current_time()
|
|
|
|
# Log current time every 5 minutes for debugging
|
|
if time.time() - last_log_time > 300: # 5 minutes
|
|
self.logger.info(f"Scheduler running - Current time: {current_time.strftime('%Y-%m-%d %H:%M:%S %Z')}")
|
|
last_log_time = time.time()
|
|
|
|
# Log APScheduler job count when it changes (max once per 30 seconds)
|
|
if self._apscheduler is not None:
|
|
current_job_count = len(self._apscheduler.get_jobs())
|
|
current_time_sec = time.time()
|
|
if current_job_count != last_job_count and (current_time_sec - last_job_log_time) >= 30:
|
|
if current_job_count > 0:
|
|
self.logger.debug(f"Found {current_job_count} scheduled jobs")
|
|
last_job_count = current_job_count
|
|
last_job_log_time = current_time_sec
|
|
|
|
# Check for interval-based advertising
|
|
self.check_interval_advertising()
|
|
|
|
# Poll feeds every minute (but feeds themselves control their check intervals)
|
|
if time.time() - last_feed_poll_time >= 60: # Every 60 seconds
|
|
if (hasattr(self.bot, 'feed_manager') and self.bot.feed_manager and
|
|
hasattr(self.bot.feed_manager, 'enabled') and self.bot.feed_manager.enabled and
|
|
hasattr(self.bot, 'connected') and self.bot.connected):
|
|
# Run feed polling in async context
|
|
import asyncio
|
|
if hasattr(self.bot, 'main_event_loop') and self.bot.main_event_loop and self.bot.main_event_loop.is_running():
|
|
# Schedule coroutine in the running main event loop
|
|
future = asyncio.run_coroutine_threadsafe(
|
|
self.bot.feed_manager.poll_all_feeds(),
|
|
self.bot.main_event_loop
|
|
)
|
|
future.add_done_callback(
|
|
lambda f: self.logger.error("Error in feed polling cycle: %s", f.exception())
|
|
if not f.cancelled() and f.exception() else None
|
|
)
|
|
else:
|
|
# Fallback: create a temporary event loop and close it when done
|
|
loop = asyncio.new_event_loop()
|
|
try:
|
|
loop.run_until_complete(self.bot.feed_manager.poll_all_feeds())
|
|
self.logger.debug("Feed polling cycle completed")
|
|
except Exception as e:
|
|
self.logger.error(f"Error in feed polling cycle: {e}")
|
|
finally:
|
|
loop.close()
|
|
last_feed_poll_time = time.time()
|
|
|
|
# Channels are fetched once on launch only - no periodic refresh
|
|
# This prevents losing channels during incomplete updates
|
|
|
|
# Process pending channel operations from web viewer (every 5 seconds)
|
|
if time.time() - self.last_channel_ops_check_time >= 5: # Every 5 seconds
|
|
if (hasattr(self.bot, 'channel_manager') and self.bot.channel_manager and
|
|
hasattr(self.bot, 'connected') and self.bot.connected):
|
|
import asyncio
|
|
if hasattr(self.bot, 'main_event_loop') and self.bot.main_event_loop and self.bot.main_event_loop.is_running():
|
|
# Schedule coroutine in the running main event loop
|
|
future = asyncio.run_coroutine_threadsafe(
|
|
self._process_channel_operations(),
|
|
self.bot.main_event_loop
|
|
)
|
|
future.add_done_callback(
|
|
lambda f: self.logger.exception("Error processing channel operations: %s", f.exception())
|
|
if not f.cancelled() and f.exception() else None
|
|
)
|
|
else:
|
|
# Fallback: create new event loop if main loop not available
|
|
try:
|
|
loop = asyncio.get_event_loop()
|
|
except RuntimeError:
|
|
loop = asyncio.new_event_loop()
|
|
asyncio.set_event_loop(loop)
|
|
|
|
loop.run_until_complete(self._process_channel_operations())
|
|
self.last_channel_ops_check_time = time.time()
|
|
|
|
# Process pending radio operations from web viewer (every 5 seconds)
|
|
if time.time() - self.last_radio_ops_check_time >= 5:
|
|
if hasattr(self.bot, 'main_event_loop') and self.bot.main_event_loop and self.bot.main_event_loop.is_running():
|
|
import asyncio
|
|
future = asyncio.run_coroutine_threadsafe(
|
|
self._process_radio_operations(),
|
|
self.bot.main_event_loop
|
|
)
|
|
future.add_done_callback(
|
|
lambda f: self.logger.exception("Error processing radio operations: %s", f.exception())
|
|
if not f.cancelled() and f.exception() else None
|
|
)
|
|
# Config-reload requests from the web viewer use the same table.
|
|
config_future = asyncio.run_coroutine_threadsafe(
|
|
self._process_config_operations(),
|
|
self.bot.main_event_loop
|
|
)
|
|
config_future.add_done_callback(
|
|
lambda f: self.logger.exception("Error processing config operations: %s", f.exception())
|
|
if not f.cancelled() and f.exception() else None
|
|
)
|
|
self.last_radio_ops_check_time = time.time()
|
|
|
|
# Process feed message queue (every 2 seconds, fire-and-forget)
|
|
# process_message_queue() returns immediately if a run is already in progress,
|
|
# so we never block this thread waiting for per-feed send intervals.
|
|
if time.time() - self.last_message_queue_check_time >= 2:
|
|
if (hasattr(self.bot, 'feed_manager') and self.bot.feed_manager and
|
|
hasattr(self.bot, 'connected') and self.bot.connected):
|
|
import asyncio
|
|
if hasattr(self.bot, 'main_event_loop') and self.bot.main_event_loop and self.bot.main_event_loop.is_running():
|
|
future = asyncio.run_coroutine_threadsafe(
|
|
self.bot.feed_manager.process_message_queue(),
|
|
self.bot.main_event_loop
|
|
)
|
|
future.add_done_callback(
|
|
lambda f: self.logger.exception("Error processing message queue: %s", f.exception())
|
|
if not f.cancelled() and f.exception() else None
|
|
)
|
|
self.last_message_queue_check_time = time.time()
|
|
|
|
# Data retention: run daily (packet_stream, repeater tables, stats, caches, mesh_connections)
|
|
if time.time() - self.last_data_retention_run >= self._data_retention_interval_seconds:
|
|
self.maintenance.run_data_retention()
|
|
self.last_data_retention_run = time.time()
|
|
|
|
# Nightly maintenance email (24 h interval, after retention so stats are fresh)
|
|
if time.time() - self.last_nightly_email_time >= self._data_retention_interval_seconds:
|
|
self.maintenance.send_nightly_email()
|
|
self.last_nightly_email_time = time.time()
|
|
|
|
# Log rotation live-apply: check bot_metadata for config changes every 60 s
|
|
if time.time() - self.last_log_rotation_check_time >= 60:
|
|
self.maintenance.apply_log_rotation_config()
|
|
self.last_log_rotation_check_time = time.time()
|
|
|
|
# DB backup: evaluate schedule every 5 minutes
|
|
if time.time() - self.last_db_backup_run >= 300:
|
|
self.maintenance.maybe_run_db_backup()
|
|
self.last_db_backup_run = time.time()
|
|
|
|
time.sleep(1)
|
|
|
|
self.logger.info("Scheduler thread stopped")
|
|
|
|
def _run_data_retention(self):
|
|
"""Run data retention cleanup: packet_stream, repeater tables, stats, caches, mesh_connections."""
|
|
import asyncio
|
|
|
|
def get_retention_days(section: str, key: str, default: int) -> int:
|
|
try:
|
|
if self.bot.config.has_section(section) and self.bot.config.has_option(section, key):
|
|
return self.bot.config.getint(section, key)
|
|
except Exception:
|
|
pass
|
|
return default
|
|
|
|
packet_stream_days = get_retention_days('Data_Retention', 'packet_stream_retention_days', 3)
|
|
purging_log_days = get_retention_days('Data_Retention', 'purging_log_retention_days', 90)
|
|
daily_stats_days = get_retention_days('Data_Retention', 'daily_stats_retention_days', 90)
|
|
observed_paths_days = get_retention_days('Data_Retention', 'observed_paths_retention_days', 90)
|
|
mesh_connections_days = get_retention_days('Data_Retention', 'mesh_connections_retention_days', 7)
|
|
stats_days = get_retention_days('Stats_Command', 'data_retention_days', 7)
|
|
|
|
try:
|
|
# Packet stream (web viewer integration)
|
|
if hasattr(self.bot, 'web_viewer_integration') and self.bot.web_viewer_integration:
|
|
bi = getattr(self.bot.web_viewer_integration, 'bot_integration', None)
|
|
if bi and hasattr(bi, 'cleanup_old_data'):
|
|
bi.cleanup_old_data(packet_stream_days)
|
|
|
|
# Repeater manager: purging_log and optional daily_stats / unique_advert / observed_paths
|
|
if hasattr(self.bot, 'repeater_manager') and self.bot.repeater_manager:
|
|
if hasattr(self.bot, 'main_event_loop') and self.bot.main_event_loop and self.bot.main_event_loop.is_running():
|
|
future = asyncio.run_coroutine_threadsafe(
|
|
self.bot.repeater_manager.cleanup_database(purging_log_days),
|
|
self.bot.main_event_loop
|
|
)
|
|
try:
|
|
future.result(timeout=60)
|
|
except RuntimeError as e:
|
|
self.logger.warning("Event loop gone during cleanup_database: %s", e)
|
|
except Exception as e:
|
|
self.logger.error(f"Error in repeater_manager.cleanup_database: {type(e).__name__}: {e}")
|
|
else:
|
|
try:
|
|
loop = asyncio.get_event_loop()
|
|
except RuntimeError:
|
|
loop = asyncio.new_event_loop()
|
|
asyncio.set_event_loop(loop)
|
|
loop.run_until_complete(self.bot.repeater_manager.cleanup_database(purging_log_days))
|
|
if hasattr(self.bot.repeater_manager, 'cleanup_repeater_retention'):
|
|
self.bot.repeater_manager.cleanup_repeater_retention(
|
|
daily_stats_days=daily_stats_days,
|
|
observed_paths_days=observed_paths_days
|
|
)
|
|
|
|
# Stats tables (message_stats, command_stats, path_stats)
|
|
if hasattr(self.bot, 'command_manager') and self.bot.command_manager:
|
|
stats_cmd = self.bot.command_manager.commands.get('stats') if getattr(self.bot.command_manager, 'commands', None) else None
|
|
if stats_cmd and hasattr(stats_cmd, 'cleanup_old_stats'):
|
|
stats_cmd.cleanup_old_stats(stats_days)
|
|
|
|
# Expired caches (geocoding_cache, generic_cache)
|
|
if hasattr(self.bot, 'db_manager') and self.bot.db_manager and hasattr(self.bot.db_manager, 'cleanup_expired_cache'):
|
|
self.bot.db_manager.cleanup_expired_cache()
|
|
|
|
# Mesh connections (DB prune to match in-memory expiration)
|
|
if hasattr(self.bot, 'mesh_graph') and self.bot.mesh_graph and hasattr(self.bot.mesh_graph, 'delete_expired_edges_from_db'):
|
|
self.bot.mesh_graph.delete_expired_edges_from_db(mesh_connections_days)
|
|
|
|
ran_at = datetime.datetime.now(datetime.timezone.utc).isoformat()
|
|
self._last_retention_stats['ran_at'] = ran_at
|
|
try:
|
|
self.bot.db_manager.set_metadata('maint.status.data_retention_ran_at', ran_at)
|
|
self.bot.db_manager.set_metadata('maint.status.data_retention_outcome', 'ok')
|
|
except Exception:
|
|
pass
|
|
|
|
except Exception as e:
|
|
self.logger.exception(f"Error during data retention cleanup: {e}")
|
|
self._last_retention_stats['error'] = str(e)
|
|
try:
|
|
ran_at = datetime.datetime.now(datetime.timezone.utc).isoformat()
|
|
self.bot.db_manager.set_metadata('maint.status.data_retention_ran_at', ran_at)
|
|
self.bot.db_manager.set_metadata('maint.status.data_retention_outcome', f'error: {e}')
|
|
except Exception:
|
|
pass
|
|
|
|
def check_interval_advertising(self):
|
|
"""Check if it's time to send an interval-based advert"""
|
|
try:
|
|
advert_interval_hours = self.bot.config.getint('Bot', 'advert_interval_hours', fallback=0)
|
|
if advert_interval_hours <= 0:
|
|
return # Interval advertising disabled
|
|
|
|
current_time = time.time()
|
|
|
|
# Check if enough time has passed since last advert
|
|
if not hasattr(self.bot, 'last_advert_time') or self.bot.last_advert_time is None:
|
|
# First time, set the timer
|
|
self.bot.last_advert_time = current_time
|
|
return
|
|
|
|
time_since_last_advert = current_time - self.bot.last_advert_time
|
|
interval_seconds = advert_interval_hours * 3600 # Convert hours to seconds
|
|
|
|
if time_since_last_advert >= interval_seconds:
|
|
self.logger.info(f"Time for interval-based advert (every {advert_interval_hours} hours)")
|
|
self.send_interval_advert()
|
|
self.bot.last_advert_time = current_time
|
|
|
|
except Exception as e:
|
|
self.logger.error(f"Error checking interval advertising: {e}")
|
|
|
|
def send_interval_advert(self):
|
|
"""Send an interval-based advert (synchronous wrapper)"""
|
|
if self.bot.is_radio_zombie:
|
|
self.logger.warning("send_interval_advert suppressed — radio is in zombie state")
|
|
return
|
|
if self.bot.is_radio_offline:
|
|
self.logger.warning("send_interval_advert suppressed — radio is offline (repeated send timeouts)")
|
|
return
|
|
|
|
current_time = self.get_current_time()
|
|
self.logger.info(f"📢 Sending interval-based flood advert at {current_time.strftime('%H:%M:%S')}")
|
|
|
|
import asyncio
|
|
|
|
# Use the main event loop if available, otherwise create a new one
|
|
# This prevents deadlock when the main loop is already running
|
|
if hasattr(self.bot, 'main_event_loop') and self.bot.main_event_loop and self.bot.main_event_loop.is_running():
|
|
# Schedule coroutine in the running main event loop
|
|
future = asyncio.run_coroutine_threadsafe(
|
|
self._send_interval_advert_async(),
|
|
self.bot.main_event_loop
|
|
)
|
|
# Wait for completion (with timeout to prevent indefinite blocking)
|
|
try:
|
|
future.result(timeout=60) # 60 second timeout
|
|
self.bot._record_send_success()
|
|
except Exception as e:
|
|
self.logger.error(f"Error sending interval advert: {type(e).__name__}: {e}")
|
|
self.bot._record_send_failure(scheduler=self)
|
|
else:
|
|
# Fallback: create new event loop if main loop not available
|
|
try:
|
|
loop = asyncio.get_event_loop()
|
|
except RuntimeError:
|
|
loop = asyncio.new_event_loop()
|
|
asyncio.set_event_loop(loop)
|
|
|
|
# Run the async function in the event loop
|
|
loop.run_until_complete(self._send_interval_advert_async())
|
|
|
|
async def _send_interval_advert_async(self):
|
|
"""Send an interval-based advert (async implementation)"""
|
|
import asyncio
|
|
|
|
from meshcore.events import EventType
|
|
try:
|
|
result = await asyncio.wait_for(
|
|
self.bot.meshcore.commands.send_advert(flood=True),
|
|
timeout=30.0,
|
|
)
|
|
except asyncio.TimeoutError:
|
|
# Feed interval advert timeouts into existing zombie-detection heuristics.
|
|
self.bot._radio_fail_count = getattr(self.bot, "_radio_fail_count", 0) + 1
|
|
self.logger.warning(
|
|
"send_interval_advert timed out after 30s; "
|
|
"_radio_fail_count=%d",
|
|
self.bot._radio_fail_count,
|
|
)
|
|
raise
|
|
if hasattr(result, 'type') and result.type == EventType.ERROR:
|
|
reason = (result.payload or {}).get('reason', 'unknown') if hasattr(result, 'payload') else 'unknown'
|
|
raise RuntimeError(f"send_advert failed: {reason}")
|
|
self.logger.info("Interval-based flood advert sent successfully")
|
|
|
|
def _claim_operation(self, operation_types: tuple[str, ...]) -> Optional[dict[str, Any]]:
|
|
"""Atomically claim the oldest pending operation in one serialized group.
|
|
|
|
``BEGIN IMMEDIATE`` prevents two scheduler ticks (or two bot processes)
|
|
from selecting the same pending row. A group permits only one
|
|
``processing`` operation at a time, preserving device-operation order
|
|
when an earlier command is slow.
|
|
|
|
Processing rows are never auto-requeued. If a previous same-host owner
|
|
is provably dead at startup, its claim becomes ``interrupted`` so later
|
|
work can proceed without replaying the ambiguous action. Live or
|
|
unprovable owners continue to block the group conservatively.
|
|
"""
|
|
if not operation_types:
|
|
raise ValueError("operation_types must not be empty")
|
|
|
|
placeholders = ', '.join('?' for _ in operation_types)
|
|
with self.bot.db_manager.connection() as conn:
|
|
conn.row_factory = sqlite3.Row
|
|
cursor = conn.cursor()
|
|
cursor.execute('BEGIN IMMEDIATE')
|
|
|
|
cursor.execute(
|
|
f'''SELECT id
|
|
FROM channel_operations
|
|
WHERE status = 'processing'
|
|
AND operation_type IN ({placeholders})
|
|
LIMIT 1''',
|
|
operation_types,
|
|
)
|
|
if cursor.fetchone() is not None:
|
|
conn.commit()
|
|
return None
|
|
|
|
cursor.execute(
|
|
f'''SELECT *
|
|
FROM channel_operations
|
|
WHERE status = 'pending'
|
|
AND operation_type IN ({placeholders})
|
|
ORDER BY created_at ASC, id ASC
|
|
LIMIT 1''',
|
|
operation_types,
|
|
)
|
|
row = cursor.fetchone()
|
|
if row is None:
|
|
conn.commit()
|
|
return None
|
|
|
|
cursor.execute(
|
|
'''UPDATE channel_operations
|
|
SET status = 'processing',
|
|
claimed_at = CURRENT_TIMESTAMP,
|
|
claim_owner_host = ?,
|
|
claim_owner_pid = ?,
|
|
claim_owner_boot_id = ?,
|
|
processed_at = NULL,
|
|
error_message = NULL,
|
|
result_data = NULL
|
|
WHERE id = ? AND status = 'pending' ''',
|
|
(
|
|
self._claim_owner_host,
|
|
self._claim_owner_pid,
|
|
self._claim_owner_boot_id,
|
|
row['id'],
|
|
),
|
|
)
|
|
if cursor.rowcount != 1:
|
|
conn.rollback()
|
|
return None
|
|
|
|
conn.commit()
|
|
claimed = dict(row)
|
|
claimed['status'] = 'processing'
|
|
return claimed
|
|
|
|
@staticmethod
|
|
def _is_local_pid_alive(pid: int) -> Optional[bool]:
|
|
"""Return PID liveness, or ``None`` when it cannot be proven either way."""
|
|
if not isinstance(pid, int) or pid <= 0:
|
|
return None
|
|
try:
|
|
os.kill(pid, 0)
|
|
except ProcessLookupError:
|
|
return False
|
|
except PermissionError:
|
|
# The kernel found the process but this user cannot signal it.
|
|
return True
|
|
except OSError:
|
|
return None
|
|
return True
|
|
|
|
def _recover_interrupted_operations(self) -> int:
|
|
"""Resolve only startup claims provably abandoned by a local process.
|
|
|
|
Legacy ownerless rows may be interrupted. An owned row is interrupted
|
|
only when it belongs to this host and its PID is provably dead. Live
|
|
same-host owners, other hosts, and incomplete/unknown identities stay
|
|
blocked conservatively. This method is only invoked during scheduler
|
|
construction, never by polling or a timer.
|
|
"""
|
|
explanation = (
|
|
"Bot restarted while this operation was processing; the action may "
|
|
"already have reached the device. Automatic retry is disabled. "
|
|
"Verify device state before submitting another operation."
|
|
)
|
|
try:
|
|
with self.bot.db_manager.connection() as conn:
|
|
conn.row_factory = sqlite3.Row
|
|
cursor = conn.cursor()
|
|
cursor.execute('BEGIN IMMEDIATE')
|
|
cursor.execute(
|
|
'''SELECT id, claim_owner_host, claim_owner_pid,
|
|
claim_owner_boot_id
|
|
FROM channel_operations
|
|
WHERE status = 'processing' ''',
|
|
)
|
|
rows = cursor.fetchall()
|
|
recovered = 0
|
|
blocked = 0
|
|
for row in rows:
|
|
owner_host = row['claim_owner_host']
|
|
owner_pid = row['claim_owner_pid']
|
|
owner_boot_id = row['claim_owner_boot_id']
|
|
legacy_ownerless = (
|
|
owner_host is None
|
|
and owner_pid is None
|
|
and owner_boot_id is None
|
|
)
|
|
same_host_dead = (
|
|
owner_host == self._claim_owner_host
|
|
and self._is_local_pid_alive(owner_pid) is False
|
|
)
|
|
if not legacy_ownerless and not same_host_dead:
|
|
blocked += 1
|
|
continue
|
|
|
|
cursor.execute(
|
|
'''UPDATE channel_operations
|
|
SET status = 'interrupted',
|
|
processed_at = CURRENT_TIMESTAMP,
|
|
error_message = ?
|
|
WHERE id = ? AND status = 'processing' ''',
|
|
(explanation, row['id']),
|
|
)
|
|
recovered += cursor.rowcount
|
|
conn.commit()
|
|
except sqlite3.Error as exc:
|
|
self.logger.exception(
|
|
"Could not recover interrupted channel/radio operations at startup: %s",
|
|
exc,
|
|
)
|
|
return 0
|
|
|
|
if recovered:
|
|
self.logger.warning(
|
|
"Marked %s operation(s) interrupted after restart; verify device state before retrying",
|
|
recovered,
|
|
)
|
|
if blocked:
|
|
self.logger.warning(
|
|
"Left %s processing operation(s) blocked because their owner is live or cannot be proven dead",
|
|
blocked,
|
|
)
|
|
return recovered
|
|
|
|
def _finish_claimed_operation(
|
|
self,
|
|
op_id: int,
|
|
*,
|
|
success: bool,
|
|
result_payload: Optional[dict[str, Any]] = None,
|
|
error_message: Optional[str] = None,
|
|
) -> None:
|
|
"""Finalize a claim without overwriting externally resolved state."""
|
|
with self.bot.db_manager.connection() as conn:
|
|
cursor = conn.cursor()
|
|
if success:
|
|
cursor.execute(
|
|
'''UPDATE channel_operations
|
|
SET status = 'completed',
|
|
processed_at = CURRENT_TIMESTAMP,
|
|
result_data = ?,
|
|
error_message = NULL
|
|
WHERE id = ? AND status = 'processing'
|
|
AND claim_owner_host = ?
|
|
AND claim_owner_pid = ?
|
|
AND claim_owner_boot_id = ? ''',
|
|
(
|
|
json.dumps(result_payload or {'success': True}),
|
|
op_id,
|
|
self._claim_owner_host,
|
|
self._claim_owner_pid,
|
|
self._claim_owner_boot_id,
|
|
),
|
|
)
|
|
else:
|
|
cursor.execute(
|
|
'''UPDATE channel_operations
|
|
SET status = 'failed',
|
|
processed_at = CURRENT_TIMESTAMP,
|
|
error_message = ?
|
|
WHERE id = ? AND status = 'processing'
|
|
AND claim_owner_host = ?
|
|
AND claim_owner_pid = ?
|
|
AND claim_owner_boot_id = ? ''',
|
|
(
|
|
error_message or 'Unknown error',
|
|
op_id,
|
|
self._claim_owner_host,
|
|
self._claim_owner_pid,
|
|
self._claim_owner_boot_id,
|
|
),
|
|
)
|
|
if cursor.rowcount != 1:
|
|
self.logger.warning(
|
|
"Operation %s was no longer processing when finalization was attempted",
|
|
op_id,
|
|
)
|
|
conn.commit()
|
|
|
|
async def _process_channel_operations(self):
|
|
"""Process pending channel operations from the web viewer"""
|
|
try:
|
|
# Preserve the previous per-tick batch capacity, but claim only one
|
|
# row immediately before executing it. The next claim cannot
|
|
# succeed until the previous row has reached a terminal state.
|
|
for _ in range(10):
|
|
op = self._claim_operation(_CHANNEL_OPERATION_TYPES)
|
|
if not op:
|
|
return
|
|
await self._execute_claimed_channel_operation(op)
|
|
|
|
except Exception as e:
|
|
db_path = getattr(self.bot.db_manager, 'db_path', 'unknown')
|
|
db_path_str = str(db_path) if db_path != 'unknown' else 'unknown'
|
|
self.logger.exception(f"Error in _process_channel_operations: {e}")
|
|
if db_path_str != 'unknown':
|
|
path_obj = Path(db_path_str)
|
|
self.logger.error(f"Database path: {db_path_str} (exists: {path_obj.exists()}, readable: {os.access(db_path_str, os.R_OK) if path_obj.exists() else False}, writable: {os.access(db_path_str, os.W_OK) if path_obj.exists() else False})")
|
|
# Check parent directory permissions
|
|
if path_obj.exists():
|
|
parent = path_obj.parent
|
|
self.logger.error(f"Parent directory: {parent} (exists: {parent.exists()}, writable: {os.access(str(parent), os.W_OK) if parent.exists() else False})")
|
|
else:
|
|
self.logger.error(f"Database path: {db_path_str}")
|
|
|
|
async def _execute_claimed_channel_operation(self, op: dict[str, Any]) -> None:
|
|
"""Execute and finalize one already-claimed channel operation."""
|
|
op_id = op['id']
|
|
op_type = op['operation_type']
|
|
channel_idx = op['channel_idx']
|
|
channel_name = op['channel_name']
|
|
channel_key_hex = op['channel_key_hex']
|
|
self.logger.info("Processing claimed channel operation %s: %s", op_id, op_type)
|
|
|
|
try:
|
|
success = False
|
|
error_msg = None
|
|
|
|
if op_type == 'add':
|
|
if channel_key_hex:
|
|
channel_secret = bytes.fromhex(channel_key_hex)
|
|
success = await self.bot.channel_manager.add_channel(
|
|
channel_idx, channel_name, channel_secret=channel_secret
|
|
)
|
|
else:
|
|
success = await self.bot.channel_manager.add_channel(
|
|
channel_idx, channel_name
|
|
)
|
|
|
|
if success:
|
|
self.logger.info(f"Successfully processed channel add operation: {channel_name} at index {channel_idx}")
|
|
else:
|
|
error_msg = "Failed to add channel"
|
|
|
|
elif op_type == 'remove':
|
|
success = await self.bot.channel_manager.remove_channel(channel_idx)
|
|
|
|
if success:
|
|
self.logger.info(f"Successfully processed channel remove operation: index {channel_idx}")
|
|
else:
|
|
error_msg = "Failed to remove channel"
|
|
|
|
self._finish_claimed_operation(
|
|
op_id,
|
|
success=success,
|
|
error_message=error_msg,
|
|
)
|
|
|
|
except Exception as e:
|
|
self.logger.error(f"Error processing channel operation {op_id}: {e}")
|
|
try:
|
|
self._finish_claimed_operation(op_id, success=False, error_message=str(e))
|
|
except Exception as update_error:
|
|
self.logger.error(f"Error updating operation status: {update_error}")
|
|
|
|
async def _process_radio_operations(self):
|
|
"""Process pending radio connect/disconnect/reboot/firmware operations from the web viewer."""
|
|
try:
|
|
op = self._claim_operation(_RADIO_OPERATION_TYPES)
|
|
if not op:
|
|
return
|
|
|
|
op_id = op['id']
|
|
op_type = op['operation_type']
|
|
self.logger.info(f"Processing radio operation {op_id}: {op_type}")
|
|
|
|
try:
|
|
result_payload = {'success': True}
|
|
if op_type == 'radio_reboot':
|
|
success = await self.bot.reboot_radio()
|
|
elif op_type == 'radio_connect':
|
|
success = await self.bot.reconnect_radio()
|
|
elif op_type == 'radio_disconnect':
|
|
success = await self.bot.disconnect_radio()
|
|
elif op_type == 'firmware_read':
|
|
success, result_payload = await self._firmware_read_op()
|
|
elif op_type == 'firmware_write':
|
|
payload = json.loads(op['payload_data'] or '{}')
|
|
success, result_payload = await self._firmware_write_op(payload)
|
|
elif op_type == 'radio_params_read':
|
|
success, result_payload = await self._radio_params_read_op()
|
|
elif op_type == 'radio_params_write':
|
|
payload = json.loads(op['payload_data'] or '{}')
|
|
success, result_payload = await self._radio_params_write_op(payload)
|
|
elif op_type == 'radio_advert':
|
|
payload = json.loads(op['payload_data'] or '{}')
|
|
success, result_payload = await self._radio_advert_op(payload)
|
|
else:
|
|
success = False
|
|
|
|
error_msg = result_payload.get('error', 'Radio operation returned False') \
|
|
if isinstance(result_payload, dict) else 'Radio operation returned False'
|
|
self._finish_claimed_operation(
|
|
op_id,
|
|
success=success,
|
|
result_payload=result_payload,
|
|
error_message=error_msg,
|
|
)
|
|
|
|
except Exception as e:
|
|
self.logger.error(f"Error executing radio operation {op_id}: {e}")
|
|
try:
|
|
self._finish_claimed_operation(op_id, success=False, error_message=str(e))
|
|
except Exception as update_error:
|
|
self.logger.error(f"Error updating radio operation status: {update_error}")
|
|
|
|
except Exception as e:
|
|
self.logger.exception(f"Error in _process_radio_operations: {e}")
|
|
|
|
async def _process_config_operations(self):
|
|
"""Process pending config_reload requests queued by the web viewer.
|
|
|
|
The web viewer (a separate process) inserts a ``config_reload`` row into
|
|
``channel_operations``; here we call ``bot.reload_config()`` so command
|
|
settings saved in the UI take effect without a restart. Service plugin
|
|
start/stop is not handled by reload_config and still needs a restart.
|
|
"""
|
|
try:
|
|
op = self._claim_operation(_CONFIG_OPERATION_TYPES)
|
|
if not op:
|
|
return
|
|
|
|
op_id = op['id']
|
|
self.logger.info(f"Processing config reload operation {op_id}")
|
|
|
|
try:
|
|
success, message = self.bot.reload_config()
|
|
except Exception as e: # noqa: BLE001 - never let reload crash the scheduler
|
|
success, message = False, str(e)
|
|
self.logger.exception("Error during config reload")
|
|
|
|
self._finish_claimed_operation(
|
|
op_id,
|
|
success=success,
|
|
result_payload={'success': True, 'message': message},
|
|
error_message=message,
|
|
)
|
|
|
|
self.logger.info("Config reload %s: %s", 'succeeded' if success else 'failed', message)
|
|
|
|
except Exception as e:
|
|
self.logger.exception(f"Error in _process_config_operations: {e}")
|
|
|
|
async def _firmware_read_op(self):
|
|
"""Read the path hash mode and default flood scope from radio firmware."""
|
|
import asyncio
|
|
try:
|
|
meshcore = getattr(self.bot, 'meshcore', None)
|
|
if not meshcore or not getattr(meshcore, 'is_connected', False):
|
|
return False, {'error': 'Radio not connected'}
|
|
|
|
path_hash_mode = await asyncio.wait_for(
|
|
meshcore.commands.get_path_hash_mode(), timeout=10
|
|
)
|
|
|
|
result: dict[str, Any] = {'path_hash_mode': path_hash_mode}
|
|
result.update(await self._read_default_flood_scope(meshcore))
|
|
return True, result
|
|
except Exception as e:
|
|
self.logger.error(f"Firmware read failed: {e}")
|
|
return False, {'error': str(e)}
|
|
|
|
async def _read_default_flood_scope(self, meshcore) -> dict[str, Any]:
|
|
"""Read the radio's own default flood scope, if it has one to report.
|
|
|
|
Never raises: firmware without CMD_GET_DEFAULT_FLOOD_SCOPE answers with
|
|
an error or not at all, and that must not cost the caller the path hash
|
|
mode it asked for in the same read. The absence of the keys is how the
|
|
page knows the radio did not answer, which is not the same as the radio
|
|
answering "no scope set".
|
|
"""
|
|
getter = getattr(meshcore.commands, 'get_default_flood_scope', None)
|
|
if getter is None:
|
|
self.logger.debug(
|
|
"meshcore library has no get_default_flood_scope; skipping"
|
|
)
|
|
return {}
|
|
try:
|
|
event = await asyncio.wait_for(getter(), timeout=10)
|
|
except Exception as e: # noqa: BLE001 - optional read, never fatal
|
|
self.logger.info("Default flood scope read unavailable: %s", e)
|
|
return {}
|
|
|
|
if event is None or getattr(event, 'type', None) == EventType.ERROR:
|
|
self.logger.info(
|
|
"Radio did not report a default flood scope (result=%s)", event
|
|
)
|
|
return {}
|
|
|
|
payload = getattr(event, 'payload', None) or {}
|
|
name = payload.get('scope_name', '') or ''
|
|
key = payload.get('scope_key', '') or ''
|
|
# A radio with no default scope sends the one-byte sentinel, which the
|
|
# library turns into an empty payload. Report that as the cleared
|
|
# state, distinct from the radio never having answered.
|
|
result: dict[str, Any] = {
|
|
'default_scope_name': name,
|
|
'default_scope_key': key,
|
|
'default_scope_supported': True,
|
|
}
|
|
if name and key:
|
|
# The name is a label stored beside the key; the radio routes by
|
|
# the key. A build-flag default stores the name without its '#'
|
|
# while hashing the '#' form, so compare on the normalized name.
|
|
result['default_scope_key_matches'] = (
|
|
scope_key_hex(name).lower() == key.lower()
|
|
)
|
|
return result
|
|
|
|
async def _firmware_write_op(self, payload: dict):
|
|
"""Write the path hash mode to radio firmware."""
|
|
import asyncio
|
|
|
|
from meshcore.events import EventType
|
|
try:
|
|
meshcore = getattr(self.bot, 'meshcore', None)
|
|
if not meshcore or not getattr(meshcore, 'is_connected', False):
|
|
return False, {'error': 'Radio not connected'}
|
|
|
|
results = {}
|
|
errors = []
|
|
|
|
if 'path_hash_mode' in payload:
|
|
mode = int(payload['path_hash_mode'])
|
|
result = await asyncio.wait_for(
|
|
meshcore.commands.set_path_hash_mode(mode), timeout=10
|
|
)
|
|
ok = getattr(result, 'type', None) == EventType.OK
|
|
results['path_hash_mode'] = ok
|
|
if not ok:
|
|
errors.append(f"set_path_hash_mode({mode}) failed: {result}")
|
|
|
|
if 'default_flood_scope' in payload:
|
|
scope = payload['default_flood_scope']
|
|
ok, error = await self._write_default_flood_scope(meshcore, scope)
|
|
results['default_flood_scope'] = ok
|
|
if error:
|
|
errors.append(error)
|
|
|
|
success = len(errors) == 0
|
|
response: dict[str, Any] = {'results': results}
|
|
if errors:
|
|
response['errors'] = errors
|
|
return success, response
|
|
except Exception as e:
|
|
self.logger.error(f"Firmware write failed: {e}")
|
|
return False, {'error': str(e)}
|
|
|
|
async def _write_default_flood_scope(self, meshcore, scope) -> tuple[bool, str]:
|
|
"""Set or clear the radio's default flood scope.
|
|
|
|
Clearing does not go through the library. ``set_default_flood_scope``
|
|
cannot do it: ``None`` raises on ``len(None)``, ``""`` makes the
|
|
firmware answer ILLEGAL_ARG, and ``"*"`` only works because the frame
|
|
is padded by character count rather than by the name it wrote, leaving
|
|
it one byte short of the length the firmware reads as "a scope
|
|
follows". The firmware's own contract for clearing is a bare
|
|
CMD_SET_DEFAULT_FLOOD_SCOPE with nothing after it, so send that.
|
|
"""
|
|
from meshcore.packets import CommandType
|
|
|
|
setter = getattr(meshcore.commands, 'set_default_flood_scope', None)
|
|
if setter is None:
|
|
return False, (
|
|
'This meshcore library version cannot set the default flood scope'
|
|
)
|
|
|
|
if scope:
|
|
result = await asyncio.wait_for(setter(scope), timeout=10)
|
|
label = f'set_default_flood_scope({scope})'
|
|
else:
|
|
result = await asyncio.wait_for(
|
|
meshcore.commands.send(
|
|
bytearray([CommandType.SET_DEFAULT_FLOOD_SCOPE.value]),
|
|
[EventType.OK, EventType.ERROR],
|
|
),
|
|
timeout=10,
|
|
)
|
|
label = 'clear default_flood_scope'
|
|
|
|
ok = getattr(result, 'type', None) == EventType.OK
|
|
if ok:
|
|
return True, ''
|
|
return False, f'{label} failed: {result}'
|
|
|
|
async def _radio_params_read_op(self):
|
|
"""Read current radio and node parameters via SELF_INFO (appstart)."""
|
|
try:
|
|
meshcore = getattr(self.bot, 'meshcore', None)
|
|
if not meshcore or not getattr(meshcore, 'is_connected', False):
|
|
return False, {'error': 'Radio not connected'}
|
|
|
|
event = await asyncio.wait_for(
|
|
meshcore.commands.send_appstart(), timeout=10
|
|
)
|
|
if event is None or event.type == EventType.ERROR:
|
|
return False, {'error': 'Failed to read radio parameters'}
|
|
|
|
p = event.payload or {}
|
|
return True, {
|
|
'freq': p.get('radio_freq'),
|
|
'bw': p.get('radio_bw'),
|
|
'sf': p.get('radio_sf'),
|
|
'cr': p.get('radio_cr'),
|
|
'tx_power': p.get('tx_power'),
|
|
'max_tx_power': p.get('max_tx_power'),
|
|
'name': p.get('name'),
|
|
'adv_lat': p.get('adv_lat'),
|
|
'adv_lon': p.get('adv_lon'),
|
|
'adv_loc_policy': p.get('adv_loc_policy'),
|
|
'manual_add_contacts': p.get('manual_add_contacts'),
|
|
'multi_acks': p.get('multi_acks'),
|
|
'telemetry_mode_base': p.get('telemetry_mode_base'),
|
|
'telemetry_mode_loc': p.get('telemetry_mode_loc'),
|
|
'telemetry_mode_env': p.get('telemetry_mode_env'),
|
|
}
|
|
except Exception as e:
|
|
self.logger.error(f"Radio params read failed: {e}")
|
|
return False, {'error': str(e)}
|
|
|
|
# Fields applied through a single CMD_SET_OTHER_PARAMS frame (read-modify-write).
|
|
OTHER_PARAMS_FIELDS = (
|
|
'manual_add_contacts', 'multi_acks', 'adv_loc_policy',
|
|
'telemetry_mode_base', 'telemetry_mode_loc', 'telemetry_mode_env',
|
|
)
|
|
|
|
async def _radio_params_write_op(self, payload: dict):
|
|
"""Write radio and node parameters to the device.
|
|
|
|
Accepts any mix of: freq/bw/sf/cr (together), tx_power, name, lat/lon
|
|
(together), rx_delay/airtime_factor (together), and the
|
|
CMD_SET_OTHER_PARAMS fields (manual_add_contacts, multi_acks,
|
|
adv_loc_policy, telemetry_mode_*).
|
|
"""
|
|
try:
|
|
meshcore = getattr(self.bot, 'meshcore', None)
|
|
if not meshcore or not getattr(meshcore, 'is_connected', False):
|
|
return False, {'error': 'Radio not connected'}
|
|
|
|
results = {}
|
|
errors = []
|
|
|
|
if any(k in payload for k in ('freq', 'bw', 'sf', 'cr')):
|
|
freq = float(payload['freq'])
|
|
bw = float(payload['bw'])
|
|
sf = int(payload['sf'])
|
|
cr = int(payload['cr'])
|
|
result = await asyncio.wait_for(
|
|
meshcore.commands.set_radio(freq, bw, sf, cr), timeout=10
|
|
)
|
|
ok = getattr(result, 'type', None) == EventType.OK
|
|
results['radio'] = ok
|
|
if not ok:
|
|
errors.append(f"set_radio failed: {result}")
|
|
|
|
if 'tx_power' in payload:
|
|
result = await asyncio.wait_for(
|
|
meshcore.commands.set_tx_power(int(payload['tx_power'])), timeout=10
|
|
)
|
|
ok = getattr(result, 'type', None) == EventType.OK
|
|
results['tx_power'] = ok
|
|
if not ok:
|
|
errors.append(f"set_tx_power failed: {result}")
|
|
|
|
if 'name' in payload:
|
|
result = await asyncio.wait_for(
|
|
meshcore.commands.set_name(str(payload['name'])), timeout=10
|
|
)
|
|
ok = getattr(result, 'type', None) == EventType.OK
|
|
results['name'] = ok
|
|
if not ok:
|
|
errors.append(f"set_name failed: {result}")
|
|
|
|
if 'lat' in payload and 'lon' in payload:
|
|
result = await asyncio.wait_for(
|
|
meshcore.commands.set_coords(
|
|
float(payload['lat']), float(payload['lon'])
|
|
), timeout=10
|
|
)
|
|
ok = getattr(result, 'type', None) == EventType.OK
|
|
results['coords'] = ok
|
|
if not ok:
|
|
errors.append(f"set_coords failed: {result}")
|
|
|
|
if any(k in payload for k in self.OTHER_PARAMS_FIELDS):
|
|
# CMD_SET_OTHER_PARAMS writes all of these at once, so read the
|
|
# current values first and overlay only the requested changes.
|
|
infos_event = await asyncio.wait_for(
|
|
meshcore.commands.send_appstart(), timeout=10
|
|
)
|
|
if infos_event is None or infos_event.type == EventType.ERROR:
|
|
errors.append('Failed to read current device settings before update')
|
|
else:
|
|
infos = dict(infos_event.payload or {})
|
|
if 'manual_add_contacts' in payload:
|
|
infos['manual_add_contacts'] = bool(payload['manual_add_contacts'])
|
|
for key in ('multi_acks', 'adv_loc_policy',
|
|
'telemetry_mode_base', 'telemetry_mode_loc',
|
|
'telemetry_mode_env'):
|
|
if key in payload:
|
|
infos[key] = int(payload[key])
|
|
result = await asyncio.wait_for(
|
|
meshcore.commands.set_other_params_from_infos(infos), timeout=10
|
|
)
|
|
ok = getattr(result, 'type', None) == EventType.OK
|
|
results['other_params'] = ok
|
|
if not ok:
|
|
errors.append(f"set_other_params failed: {result}")
|
|
|
|
if 'rx_delay' in payload and 'airtime_factor' in payload:
|
|
# Firmware stores these as floats; the wire format is value x1000.
|
|
rx_ms = int(round(float(payload['rx_delay']) * 1000))
|
|
af_ms = int(round(float(payload['airtime_factor']) * 1000))
|
|
result = await asyncio.wait_for(
|
|
meshcore.commands.set_tuning(rx_ms, af_ms), timeout=10
|
|
)
|
|
ok = getattr(result, 'type', None) == EventType.OK
|
|
results['tuning'] = ok
|
|
if not ok:
|
|
errors.append(f"set_tuning failed: {result}")
|
|
|
|
success = len(errors) == 0
|
|
response: dict = {'results': results}
|
|
if errors:
|
|
response['errors'] = errors
|
|
return success, response
|
|
except Exception as e:
|
|
self.logger.error(f"Radio params write failed: {e}")
|
|
return False, {'error': str(e)}
|
|
|
|
async def _radio_advert_op(self, payload: dict):
|
|
"""Send a self-advertisement (optionally flooded) from the device."""
|
|
try:
|
|
meshcore = getattr(self.bot, 'meshcore', None)
|
|
if not meshcore or not getattr(meshcore, 'is_connected', False):
|
|
return False, {'error': 'Radio not connected'}
|
|
|
|
flood = bool(payload.get('flood', False))
|
|
result = await asyncio.wait_for(
|
|
meshcore.commands.send_advert(flood=flood), timeout=10
|
|
)
|
|
ok = getattr(result, 'type', None) == EventType.OK
|
|
if not ok:
|
|
return False, {'error': f"send_advert failed: {result}"}
|
|
return True, {'flood': flood}
|
|
except Exception as e:
|
|
self.logger.error(f"Radio advert failed: {e}")
|
|
return False, {'error': str(e)}
|
|
|
|
# ── Maintenance (delegates to MaintenanceRunner) ─────────────────────────
|
|
|
|
@property
|
|
def _last_retention_stats(self) -> dict[str, Any]:
|
|
return self.maintenance._last_retention_stats
|
|
|
|
@_last_retention_stats.setter
|
|
def _last_retention_stats(self, value: dict[str, Any]) -> None:
|
|
self.maintenance._last_retention_stats.clear()
|
|
self.maintenance._last_retention_stats.update(value)
|
|
|
|
@property
|
|
def _last_db_backup_stats(self) -> dict[str, Any]:
|
|
return self.maintenance._last_db_backup_stats
|
|
|
|
@_last_db_backup_stats.setter
|
|
def _last_db_backup_stats(self, value: dict[str, Any]) -> None:
|
|
self.maintenance._last_db_backup_stats.clear()
|
|
self.maintenance._last_db_backup_stats.update(value)
|
|
|
|
@property
|
|
def _last_log_rotation_applied(self) -> dict[str, str]:
|
|
return self.maintenance._last_log_rotation_applied
|
|
|
|
@_last_log_rotation_applied.setter
|
|
def _last_log_rotation_applied(self, value: dict[str, str]) -> None:
|
|
self.maintenance._last_log_rotation_applied.clear()
|
|
self.maintenance._last_log_rotation_applied.update(value)
|
|
|
|
def run_db_backup(self) -> None:
|
|
"""Run a DB backup immediately (manual / HTTP)."""
|
|
self.maintenance.run_db_backup()
|
|
|
|
def _run_data_retention(self) -> None:
|
|
self.maintenance.run_data_retention()
|
|
|
|
def _get_notif(self, key: str) -> str:
|
|
return self.maintenance.get_notif(key)
|
|
|
|
def _collect_email_stats(self) -> dict[str, Any]:
|
|
return self.maintenance.collect_email_stats()
|
|
|
|
def _format_email_body(self, stats: dict[str, Any], period_start: str, period_end: str) -> str:
|
|
return self.maintenance.format_email_body(stats, period_start, period_end)
|
|
|
|
# ── Zombie radio alert email ─────────────────────────────────────────────
|
|
|
|
def send_zombie_alert_email(self, fail_count: int, threshold: int, interval: int) -> None:
|
|
"""Send an immediate alert email when a zombie radio is detected.
|
|
|
|
Uses the same SMTP settings as the nightly digest. Recipients are taken
|
|
from the ``radio_zombie_alert_email`` config key; if that key is empty the
|
|
nightly maintenance recipients are used as a fallback.
|
|
|
|
This method is intentionally synchronous so it can be run in a thread
|
|
executor from the async event loop without blocking it.
|
|
"""
|
|
import smtplib
|
|
import ssl as _ssl
|
|
from email.message import EmailMessage
|
|
|
|
zombie_alert_enabled = self.bot.config.getboolean(
|
|
'Connection',
|
|
'radio_zombie_alert_enabled',
|
|
fallback=self.bot.config.getboolean('Bot', 'radio_zombie_alert_enabled', fallback=False),
|
|
)
|
|
zombie_alert_email_cfg: str | None = None
|
|
db_manager = getattr(self.bot, 'db_manager', None)
|
|
if db_manager is not None and hasattr(db_manager, 'get_metadata'):
|
|
try:
|
|
meta_enabled = db_manager.get_metadata('zombie.alert_enabled')
|
|
if isinstance(meta_enabled, str) and meta_enabled.strip():
|
|
zombie_alert_enabled = meta_enabled.strip().lower() in {
|
|
'1', 'true', 'yes', 'on',
|
|
}
|
|
meta_email = db_manager.get_metadata('zombie.alert_email')
|
|
if isinstance(meta_email, str) and meta_email.strip():
|
|
zombie_alert_email_cfg = meta_email.strip()
|
|
except Exception:
|
|
pass
|
|
if not zombie_alert_enabled:
|
|
return
|
|
|
|
smtp_host = self._get_notif('smtp_host')
|
|
smtp_security = self._get_notif('smtp_security') or 'starttls'
|
|
smtp_user = self._get_notif('smtp_user')
|
|
smtp_password = self._get_notif('smtp_password')
|
|
from_name = self._get_notif('from_name') or 'MeshCore Bot'
|
|
from_email = self._get_notif('from_email')
|
|
|
|
# Alert recipients: dedicated config key, falls back to nightly recipients
|
|
alert_email_cfg = zombie_alert_email_cfg or self.bot.config.get(
|
|
'Connection',
|
|
'radio_zombie_alert_email',
|
|
fallback=self.bot.config.get('Bot', 'radio_zombie_alert_email', fallback=''),
|
|
).strip()
|
|
if alert_email_cfg:
|
|
recipients = [r.strip() for r in alert_email_cfg.split(',') if r.strip()]
|
|
else:
|
|
recipients = [r.strip() for r in self._get_notif('recipients').split(',') if r.strip()]
|
|
|
|
if not smtp_host or not from_email or not recipients:
|
|
self.bot.logger.warning(
|
|
"Zombie alert email enabled but SMTP settings incomplete "
|
|
f"(host={smtp_host!r}, from={from_email!r}, recipients={recipients}) "
|
|
"— alert email not sent"
|
|
)
|
|
return
|
|
|
|
allow_local = self._get_notif('allow_local_smtp').lower() == 'true'
|
|
if not validate_external_url(f'http://{smtp_host}', allow_private=allow_local):
|
|
self.bot.logger.error(
|
|
"Zombie alert email aborted: SMTP host %r resolves to a private or reserved address",
|
|
smtp_host,
|
|
)
|
|
return
|
|
|
|
try:
|
|
smtp_port = int(self._get_notif('smtp_port') or (465 if smtp_security == 'ssl' else 587))
|
|
except ValueError:
|
|
smtp_port = 587
|
|
|
|
now_utc = datetime.datetime.now(datetime.timezone.utc)
|
|
connection_type = self.bot.config.get('Connection', 'connection_type', fallback='unknown')
|
|
serial_port = self.bot.config.get('Connection', 'serial_port', fallback='n/a')
|
|
interval_min = interval // 60
|
|
|
|
subject = (
|
|
f'ALERT: MeshCore Bot — Zombie Radio Detected '
|
|
f'[{now_utc.strftime("%Y-%m-%d %H:%M UTC")}]'
|
|
)
|
|
body = '\n'.join([
|
|
'MeshCore Bot — Zombie Radio Alert',
|
|
'=' * 44,
|
|
f'Time : {now_utc.strftime("%Y-%m-%d %H:%M:%S UTC")}',
|
|
'',
|
|
'RADIO STATUS',
|
|
'─' * 30,
|
|
f' Connection : {connection_type}',
|
|
f' Port / Device : {serial_port}',
|
|
f' Failed probes : {fail_count} of {threshold} (threshold)',
|
|
f' Probe interval: {interval}s ({interval_min} min)',
|
|
'',
|
|
'ACTION REQUIRED',
|
|
'─' * 30,
|
|
' The radio firmware is unresponsive (zombie state).',
|
|
' A physical POWER CYCLE of the radio is required.',
|
|
' Disconnect/reconnect of the serial/BLE transport will NOT fix this.',
|
|
'',
|
|
' Steps to recover:',
|
|
' 1. Power off the radio hardware',
|
|
' 2. Wait 10 seconds',
|
|
' 3. Power on the radio hardware',
|
|
' 4. The bot will reconnect and resume normal operation automatically',
|
|
'',
|
|
'─' * 44,
|
|
'Probe monitoring has been suspended to avoid log spam.',
|
|
'It will resume automatically after the next successful reconnect.',
|
|
])
|
|
|
|
try:
|
|
msg = EmailMessage()
|
|
msg['Subject'] = subject
|
|
msg['From'] = f'{from_name} <{from_email}>'
|
|
msg['To'] = ', '.join(recipients)
|
|
msg.set_content(body)
|
|
|
|
context = _ssl.create_default_context()
|
|
_smtp_timeout = 30
|
|
|
|
if smtp_security == 'ssl':
|
|
with smtplib.SMTP_SSL(smtp_host, smtp_port, context=context, timeout=_smtp_timeout) as s:
|
|
if smtp_user and smtp_password:
|
|
s.login(smtp_user, smtp_password)
|
|
s.send_message(msg)
|
|
else:
|
|
with smtplib.SMTP(smtp_host, smtp_port, timeout=_smtp_timeout) as s:
|
|
if smtp_security == 'starttls':
|
|
s.ehlo()
|
|
s.starttls(context=context)
|
|
s.ehlo()
|
|
if smtp_user and smtp_password:
|
|
s.login(smtp_user, smtp_password)
|
|
s.send_message(msg)
|
|
|
|
self.bot.logger.info(
|
|
f"Zombie radio alert email sent to {recipients}"
|
|
)
|
|
except Exception as e:
|
|
self.bot.logger.error(f"Failed to send zombie radio alert email: {e}")
|
|
|
|
# ── Maintenance helpers ──────────────────────────────────────────────────
|
|
|
|
def _get_maint(self, key: str) -> str:
|
|
return self.maintenance.get_maint(key)
|
|
|
|
def _apply_log_rotation_config(self) -> None:
|
|
self.maintenance.apply_log_rotation_config()
|
|
|
|
def _maybe_run_db_backup(self) -> None:
|
|
self.maintenance.maybe_run_db_backup()
|
|
|
|
def _run_db_backup(self) -> None:
|
|
self.maintenance.run_db_backup()
|
|
|
|
# ── Radio offline alert email ────────────────────────────────────────────
|
|
|
|
def send_radio_offline_alert_email(self, fail_count: int, threshold: int) -> None:
|
|
"""Send an immediate alert email when the radio-offline state is entered.
|
|
|
|
Uses the same SMTP settings as the nightly digest. Recipients are taken
|
|
from the ``radio_offline_alert_email`` config key; if that key is empty the
|
|
nightly maintenance recipients are used as a fallback.
|
|
|
|
Intentionally synchronous — intended to be run in a daemon thread.
|
|
"""
|
|
import smtplib
|
|
import ssl as _ssl
|
|
from email.message import EmailMessage
|
|
|
|
alert_enabled = self.bot.config.getboolean(
|
|
'Connection',
|
|
'radio_offline_alert_enabled',
|
|
fallback=self.bot.config.getboolean('Bot', 'radio_offline_alert_enabled', fallback=False),
|
|
)
|
|
if not alert_enabled:
|
|
return
|
|
|
|
smtp_host = self._get_notif('smtp_host')
|
|
smtp_security = self._get_notif('smtp_security') or 'starttls'
|
|
smtp_user = self._get_notif('smtp_user')
|
|
smtp_password = self._get_notif('smtp_password')
|
|
from_name = self._get_notif('from_name') or 'MeshCore Bot'
|
|
from_email = self._get_notif('from_email')
|
|
|
|
alert_email_cfg = self.bot.config.get(
|
|
'Connection',
|
|
'radio_offline_alert_email',
|
|
fallback=self.bot.config.get('Bot', 'radio_offline_alert_email', fallback=''),
|
|
).strip()
|
|
if alert_email_cfg:
|
|
recipients = [r.strip() for r in alert_email_cfg.split(',') if r.strip()]
|
|
else:
|
|
recipients = [r.strip() for r in self._get_notif('recipients').split(',') if r.strip()]
|
|
|
|
if not smtp_host or not from_email or not recipients:
|
|
self.bot.logger.warning(
|
|
"Radio-offline alert email enabled but SMTP settings incomplete "
|
|
f"(host={smtp_host!r}, from={from_email!r}, recipients={recipients}) "
|
|
"— alert email not sent"
|
|
)
|
|
return
|
|
|
|
allow_local = self._get_notif('allow_local_smtp').lower() == 'true'
|
|
if not validate_external_url(f'http://{smtp_host}', allow_private=allow_local):
|
|
self.bot.logger.error(
|
|
"Radio-offline alert email aborted: SMTP host %r resolves to a private or reserved address",
|
|
smtp_host,
|
|
)
|
|
return
|
|
|
|
try:
|
|
smtp_port = int(self._get_notif('smtp_port') or (465 if smtp_security == 'ssl' else 587))
|
|
except ValueError:
|
|
smtp_port = 587
|
|
|
|
now_utc = datetime.datetime.now(datetime.timezone.utc)
|
|
connection_type = self.bot.config.get('Connection', 'connection_type', fallback='unknown')
|
|
serial_port = self.bot.config.get('Connection', 'serial_port', fallback='n/a')
|
|
|
|
subject = (
|
|
f'ALERT: MeshCore Bot — Radio Offline '
|
|
f'[{now_utc.strftime("%Y-%m-%d %H:%M UTC")}]'
|
|
)
|
|
body = '\n'.join([
|
|
'MeshCore Bot — Radio Offline Alert',
|
|
'=' * 44,
|
|
f'Time : {now_utc.strftime("%Y-%m-%d %H:%M:%S UTC")}',
|
|
'',
|
|
'RADIO STATUS',
|
|
'─' * 30,
|
|
f' Connection : {connection_type}',
|
|
f' Port / Device : {serial_port}',
|
|
f' Failed sends : {fail_count} of {threshold} (threshold)',
|
|
'',
|
|
'WHAT THIS MEANS',
|
|
'─' * 30,
|
|
' The bot can no longer send outbound messages to the mesh.',
|
|
' Inbound packets from the radio may still be arriving normally.',
|
|
' This is NOT a zombie (firmware lock-up) — the radio is responsive',
|
|
' but outbound sends are timing out.',
|
|
'',
|
|
'ACTION REQUIRED',
|
|
'─' * 30,
|
|
' Check the radio power supply and physical connection.',
|
|
' Use the dashboard "Clear Offline Flag" button once the issue',
|
|
' is resolved, or restart the bot to auto-probe.',
|
|
'',
|
|
'─' * 44,
|
|
'Outbound sends will be suppressed until the offline flag is cleared.',
|
|
])
|
|
|
|
try:
|
|
msg = EmailMessage()
|
|
msg['Subject'] = subject
|
|
msg['From'] = f'{from_name} <{from_email}>'
|
|
msg['To'] = ', '.join(recipients)
|
|
msg.set_content(body)
|
|
|
|
context = _ssl.create_default_context()
|
|
_smtp_timeout = 30
|
|
|
|
if smtp_security == 'ssl':
|
|
with smtplib.SMTP_SSL(smtp_host, smtp_port, context=context, timeout=_smtp_timeout) as s:
|
|
if smtp_user and smtp_password:
|
|
s.login(smtp_user, smtp_password)
|
|
s.send_message(msg)
|
|
else:
|
|
with smtplib.SMTP(smtp_host, smtp_port, timeout=_smtp_timeout) as s:
|
|
if smtp_security == 'starttls':
|
|
s.ehlo()
|
|
s.starttls(context=context)
|
|
s.ehlo()
|
|
if smtp_user and smtp_password:
|
|
s.login(smtp_user, smtp_password)
|
|
s.send_message(msg)
|
|
|
|
self.bot.logger.info(f"Radio-offline alert email sent to {recipients}")
|
|
except Exception as e:
|
|
self.bot.logger.error(f"Failed to send radio-offline alert email: {e}")
|
|
|
|
# ── Maintenance helpers ──────────────────────────────────────────────────
|