mirror of
https://github.com/agessaman/meshcore-bot.git
synced 2026-10-05 22:07:17 +00:00
Closes #279. A MeshCore client with no region configured sends every channel message as a plain FLOOD, which every repeater on the mesh rebroadcasts. This adds two things: free observation of how much of that the bot hears, and an opt-in warning to the senders. Observation classifies each channel message as scoped, global or unknown and tallies it per channel per local day. It costs one upsert and no airtime, and it is what lets an operator see the size of the problem before deciding to spend airtime on it. The classification runs ahead of the flood_scopes allowlist, because an unscoped message is exactly what that allowlist drops. Warnings only ever fire on positive RF evidence of an unscoped FLOOD. Absence of correlation is not proof that a sender omitted a region, so it classifies as unknown and stays quiet. They are off by default, start in dry run, and are fenced by min_unscoped_messages, a per-sender cooldown, a mesh-wide cooldown and a daily cap on attempts. Dry run consumes the same budget it previews, so the log is what going live would put on the air, not an upper bound. All three limits read from the database, so a restart cannot release a burst. Channel-delivered warnings go out at global scope on purpose: the recipient is by definition outside any region the bot replies under. New Settings -> Region Warnings page in the web viewer covers all of it. Its dark-mode striped rows exposed a pre-existing base.html bug where Bootstrap's light-theme text color survived on a dark row background (~1.3:1); fixed there for every table in the app.
754 lines
28 KiB
Python
754 lines
28 KiB
Python
#!/usr/bin/env python3
|
|
"""Regional flood-scope observation and the optional "set a region code" warning.
|
|
|
|
MeshCore puts a channel message on the air in one of two ways. An ordinary
|
|
``FLOOD`` is rebroadcast by every repeater that hears it, anywhere on the mesh.
|
|
A ``TC_FLOOD`` carries a transport code derived from a region key — the "region
|
|
code" — and only repeaters holding that key pass it on. A client with no region
|
|
configured therefore floods the entire mesh with every message it sends, which
|
|
is the traffic problem this module is about (issue #279).
|
|
|
|
Two things live here:
|
|
|
|
* **Observation.** Every channel message the bot hears is classified as
|
|
``scoped``, ``global`` or ``unknown`` and tallied into ``region_scope_daily``.
|
|
This costs one upsert per message and no airtime, and it is what lets an
|
|
operator see how much unscoped traffic their mesh actually carries *before*
|
|
deciding whether telling anyone about it is worth the airtime.
|
|
|
|
* **Warning.** When explicitly enabled, a sender whose messages are confirmed
|
|
unscoped can be sent a short note asking them to set a region code. This
|
|
spends airtime automatically, so it is off by default, starts in dry-run, and
|
|
is fenced by a per-sender cooldown, a mesh-wide cooldown and a daily cap.
|
|
|
|
The classification only ever warns on *positive* evidence of a global flood.
|
|
Absence of evidence classifies as ``unknown`` and is never warned about — see
|
|
``MessageHandler._classify_channel_flood_scope``, which produces the verdict
|
|
this module consumes.
|
|
|
|
The pure ``load_settings`` / summary helpers take a ``configparser`` and a
|
|
``DBManager`` rather than the bot, because the web viewer runs in a separate
|
|
process and has no bot object.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import time
|
|
from dataclasses import dataclass, field
|
|
from datetime import date, datetime, timedelta
|
|
from typing import Any, Optional
|
|
|
|
CONFIG_SECTION = "Region_Warnings"
|
|
|
|
# Verdicts produced by MessageHandler._classify_channel_flood_scope.
|
|
VERDICT_SCOPED = "scoped"
|
|
VERDICT_GLOBAL = "global"
|
|
VERDICT_UNKNOWN = "unknown"
|
|
VERDICTS = (VERDICT_SCOPED, VERDICT_GLOBAL, VERDICT_UNKNOWN)
|
|
|
|
_VERDICT_COLUMN = {
|
|
VERDICT_SCOPED: "scoped_count",
|
|
VERDICT_GLOBAL: "global_count",
|
|
VERDICT_UNKNOWN: "unknown_count",
|
|
}
|
|
|
|
DELIVERY_DM = "dm"
|
|
DELIVERY_CHANNEL = "channel"
|
|
|
|
ACTION_SENT = "sent"
|
|
ACTION_DRY_RUN = "dry_run"
|
|
ACTION_FAILED = "failed"
|
|
|
|
DEFAULT_MESSAGE = (
|
|
"Heads up: your messages have no region code, so they flood the whole mesh. "
|
|
"Setting one in your MeshCore app keeps things quiet. Thanks!"
|
|
)
|
|
|
|
# Body budget for a DM, matching CommandManager.get_max_message_length.
|
|
DM_BODY_LIMIT = 158
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class RegionWarningSettings:
|
|
"""Resolved ``[Region_Warnings]`` configuration."""
|
|
|
|
enabled: bool = False
|
|
dry_run: bool = True
|
|
delivery: str = DELIVERY_DM
|
|
channels: tuple[str, ...] = ()
|
|
message: str = DEFAULT_MESSAGE
|
|
min_unscoped_messages: int = 3
|
|
per_sender_cooldown_hours: float = 168.0
|
|
mesh_cooldown_minutes: float = 30.0
|
|
max_warnings_per_day: int = 6
|
|
track_traffic: bool = True
|
|
|
|
def monitors_channel(self, channel: Optional[str]) -> bool:
|
|
"""Whether ``channel`` is in the allowlist (an empty allowlist means all)."""
|
|
if not self.channels:
|
|
return True
|
|
return normalize_channel(channel) in self.channels
|
|
|
|
|
|
def normalize_channel(channel: Optional[str]) -> str:
|
|
"""Case-fold a channel name and drop a leading ``#`` so config matches the wire."""
|
|
return (channel or "").strip().lstrip("#").lower()
|
|
|
|
|
|
def _get(config: Any, key: str, fallback: str = "") -> str:
|
|
try:
|
|
if config is not None and config.has_option(CONFIG_SECTION, key):
|
|
return (config.get(CONFIG_SECTION, key) or "").strip()
|
|
except Exception:
|
|
pass
|
|
return fallback
|
|
|
|
|
|
def _get_bool(config: Any, key: str, fallback: bool) -> bool:
|
|
raw = _get(config, key).lower()
|
|
if raw in ("1", "true", "yes", "on"):
|
|
return True
|
|
if raw in ("0", "false", "no", "off"):
|
|
return False
|
|
return fallback
|
|
|
|
|
|
def _get_number(config: Any, key: str, fallback: float, *, minimum: float = 0.0) -> float:
|
|
raw = _get(config, key)
|
|
if not raw:
|
|
return fallback
|
|
try:
|
|
value = float(raw)
|
|
except ValueError:
|
|
return fallback
|
|
return max(minimum, value)
|
|
|
|
|
|
def load_settings(config: Any) -> RegionWarningSettings:
|
|
"""Read ``[Region_Warnings]`` into a settings object, falling back to defaults.
|
|
|
|
Unparseable values fall back rather than raising: this runs on every config
|
|
reload in the message path, and a typo should not stop the bot handling mail.
|
|
"""
|
|
delivery = _get(config, "delivery", DELIVERY_DM).lower()
|
|
if delivery not in (DELIVERY_DM, DELIVERY_CHANNEL):
|
|
delivery = DELIVERY_DM
|
|
|
|
channels = tuple(
|
|
normalize_channel(part)
|
|
for part in _get(config, "channels").split(",")
|
|
if normalize_channel(part)
|
|
)
|
|
|
|
message = _get(config, "message") or DEFAULT_MESSAGE
|
|
|
|
return RegionWarningSettings(
|
|
enabled=_get_bool(config, "enabled", False),
|
|
dry_run=_get_bool(config, "dry_run", True),
|
|
delivery=delivery,
|
|
channels=channels,
|
|
message=message,
|
|
min_unscoped_messages=int(_get_number(config, "min_unscoped_messages", 3, minimum=1)),
|
|
per_sender_cooldown_hours=_get_number(config, "per_sender_cooldown_hours", 168.0),
|
|
mesh_cooldown_minutes=_get_number(config, "mesh_cooldown_minutes", 30.0),
|
|
max_warnings_per_day=int(_get_number(config, "max_warnings_per_day", 6)),
|
|
track_traffic=_get_bool(config, "track_traffic", True),
|
|
)
|
|
|
|
|
|
def settings_to_config_values(settings: RegionWarningSettings) -> dict[str, str]:
|
|
"""Render settings back to INI strings for :mod:`modules.settings_store`."""
|
|
def _num(value: float) -> str:
|
|
return str(int(value)) if float(value).is_integer() else str(value)
|
|
|
|
return {
|
|
"enabled": "true" if settings.enabled else "false",
|
|
"dry_run": "true" if settings.dry_run else "false",
|
|
"delivery": settings.delivery,
|
|
"channels": ", ".join(settings.channels),
|
|
"message": settings.message,
|
|
"min_unscoped_messages": str(settings.min_unscoped_messages),
|
|
"per_sender_cooldown_hours": _num(settings.per_sender_cooldown_hours),
|
|
"mesh_cooldown_minutes": _num(settings.mesh_cooldown_minutes),
|
|
"max_warnings_per_day": str(settings.max_warnings_per_day),
|
|
"track_traffic": "true" if settings.track_traffic else "false",
|
|
}
|
|
|
|
|
|
def render_message(template: str, sender: Optional[str], channel: Optional[str]) -> str:
|
|
"""Substitute ``{sender}`` / ``{channel}``, leaving unknown braces untouched."""
|
|
text = template or DEFAULT_MESSAGE
|
|
replacements = {
|
|
"{sender}": sender or "",
|
|
"{channel}": channel or "",
|
|
}
|
|
for token, value in replacements.items():
|
|
text = text.replace(token, value)
|
|
return text.strip()
|
|
|
|
|
|
def truncate_to_bytes(text: str, limit: int) -> str:
|
|
"""Trim ``text`` to ``limit`` UTF-8 bytes without splitting a character."""
|
|
encoded = text.encode("utf-8")
|
|
if len(encoded) <= limit:
|
|
return text
|
|
return encoded[:limit].decode("utf-8", errors="ignore")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Read helpers — shared with the web viewer, which has no bot object
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def local_now(config: Any, logger: Any = None) -> datetime:
|
|
"""Now in the configured ``[Bot] timezone``, as a naive datetime.
|
|
|
|
Everything this module stores and compares — tallies, the daily cap, both
|
|
cooldowns — goes through here, so the whole feature agrees on which day it
|
|
is even when ``[Bot] timezone`` differs from the host's. Naive, because the
|
|
values are written to SQLite and read back by ``date(created_at)``, which
|
|
has no notion of offsets.
|
|
"""
|
|
try:
|
|
from modules.utils import get_config_timezone
|
|
|
|
tz, _name = get_config_timezone(config, logger)
|
|
return datetime.now(tz).replace(tzinfo=None)
|
|
except Exception:
|
|
return datetime.now()
|
|
|
|
|
|
def _local_today(config: Any, logger: Any = None) -> date:
|
|
"""Today's date in the configured ``[Bot] timezone``."""
|
|
return local_now(config, logger).date()
|
|
|
|
|
|
def traffic_summary(
|
|
db_manager: Any,
|
|
config: Any = None,
|
|
days: int = 14,
|
|
logger: Any = None,
|
|
) -> dict[str, Any]:
|
|
"""Per-channel scoped/global/unknown tallies over the last ``days`` local days."""
|
|
days = max(1, int(days))
|
|
since = (_local_today(config, logger) - timedelta(days=days - 1)).isoformat()
|
|
try:
|
|
rows = db_manager.execute_query(
|
|
"SELECT channel, "
|
|
"SUM(scoped_count) AS scoped, "
|
|
"SUM(global_count) AS global_, "
|
|
"SUM(unknown_count) AS unknown "
|
|
"FROM region_scope_daily WHERE date >= ? "
|
|
"GROUP BY channel ORDER BY (SUM(global_count) + SUM(scoped_count) + SUM(unknown_count)) DESC",
|
|
(since,),
|
|
)
|
|
except Exception:
|
|
rows = []
|
|
|
|
channels: list[dict[str, Any]] = []
|
|
totals: dict[str, Any] = {"scoped": 0, "global": 0, "unknown": 0}
|
|
for row in rows:
|
|
scoped = int(row.get("scoped") or 0)
|
|
globally = int(row.get("global_") or 0)
|
|
unknown = int(row.get("unknown") or 0)
|
|
total = scoped + globally + unknown
|
|
if total <= 0:
|
|
continue
|
|
channels.append({
|
|
"channel": row.get("channel") or "",
|
|
"scoped": scoped,
|
|
"global": globally,
|
|
"unknown": unknown,
|
|
"total": total,
|
|
# Share of *classified* traffic that was unscoped. Unknown messages
|
|
# are excluded from the denominator rather than counted as scoped,
|
|
# so a mesh the bot cannot classify reads as "no data", not "clean".
|
|
"unscoped_pct": round(100.0 * globally / (scoped + globally), 1) if (scoped + globally) else None,
|
|
})
|
|
totals["scoped"] += scoped
|
|
totals["global"] += globally
|
|
totals["unknown"] += unknown
|
|
|
|
classified = totals["scoped"] + totals["global"]
|
|
totals["total"] = classified + totals["unknown"]
|
|
totals["unscoped_pct"] = round(100.0 * totals["global"] / classified, 1) if classified else None
|
|
return {"days": days, "since": since, "channels": channels, "totals": totals}
|
|
|
|
|
|
def daily_series(
|
|
db_manager: Any,
|
|
config: Any = None,
|
|
days: int = 14,
|
|
logger: Any = None,
|
|
) -> list[dict[str, Any]]:
|
|
"""One row per local date over the window, with zero-filled gaps."""
|
|
days = max(1, int(days))
|
|
today = _local_today(config, logger)
|
|
since = (today - timedelta(days=days - 1)).isoformat()
|
|
try:
|
|
rows = db_manager.execute_query(
|
|
"SELECT date, SUM(scoped_count) AS scoped, SUM(global_count) AS global_, "
|
|
"SUM(unknown_count) AS unknown FROM region_scope_daily "
|
|
"WHERE date >= ? GROUP BY date",
|
|
(since,),
|
|
)
|
|
except Exception:
|
|
rows = []
|
|
by_date = {
|
|
row.get("date"): (
|
|
int(row.get("scoped") or 0),
|
|
int(row.get("global_") or 0),
|
|
int(row.get("unknown") or 0),
|
|
)
|
|
for row in rows
|
|
}
|
|
series = []
|
|
for offset in range(days - 1, -1, -1):
|
|
key = (today - timedelta(days=offset)).isoformat()
|
|
scoped, globally, unknown = by_date.get(key, (0, 0, 0))
|
|
series.append({
|
|
"date": key,
|
|
"scoped": scoped,
|
|
"global": globally,
|
|
"unknown": unknown,
|
|
})
|
|
return series
|
|
|
|
|
|
def recent_events(db_manager: Any, limit: int = 50) -> list[dict[str, Any]]:
|
|
"""Most recent warning decisions, newest first."""
|
|
limit = max(1, min(int(limit), 500))
|
|
try:
|
|
return db_manager.execute_query(
|
|
"SELECT id, created_at, sender_id, sender_pubkey, channel, delivery, action, detail "
|
|
"FROM region_warning_events ORDER BY created_at DESC, id DESC LIMIT ?",
|
|
(limit,),
|
|
)
|
|
except Exception:
|
|
return []
|
|
|
|
|
|
def warning_budget(
|
|
db_manager: Any,
|
|
settings: RegionWarningSettings,
|
|
config: Any = None,
|
|
logger: Any = None,
|
|
) -> dict[str, Any]:
|
|
"""How much of the daily cap is spent and when the mesh cooldown lifts.
|
|
|
|
The cap counts *attempts*, failures included. Its job is to bound how much
|
|
unprompted activity this feature can produce in a day, and a send that
|
|
reported failure may still have put something on the air before it did.
|
|
"""
|
|
today = _local_today(config, logger).isoformat()
|
|
used = 0
|
|
last_at: Optional[str] = None
|
|
try:
|
|
rows = db_manager.execute_query(
|
|
"SELECT COUNT(*) AS used FROM region_warning_events "
|
|
"WHERE date(created_at) = ?",
|
|
(today,),
|
|
)
|
|
if rows:
|
|
used = int(rows[0].get("used") or 0)
|
|
rows = db_manager.execute_query(
|
|
"SELECT MAX(created_at) AS last_at FROM region_warning_events WHERE action IN (?, ?)",
|
|
(ACTION_SENT, ACTION_DRY_RUN),
|
|
)
|
|
if rows:
|
|
last_at = rows[0].get("last_at")
|
|
except Exception:
|
|
pass
|
|
|
|
cap = settings.max_warnings_per_day
|
|
return {
|
|
"date": today,
|
|
"used_today": used,
|
|
"cap": cap,
|
|
"remaining": None if cap <= 0 else max(0, cap - used),
|
|
"unlimited": cap <= 0,
|
|
"last_warning_at": last_at,
|
|
}
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Monitor — bot-side
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@dataclass
|
|
class _SenderState:
|
|
"""In-memory per-sender counter, rebuilt after a restart."""
|
|
|
|
unscoped_seen: int = 0
|
|
counted_at: float = field(default_factory=time.monotonic)
|
|
|
|
|
|
class RegionWarningMonitor:
|
|
"""Tally flood scopes and, when enabled, warn senders who never set one."""
|
|
|
|
#: How long a sender's unscoped-message run survives without new traffic.
|
|
#: A sender who floods three times a year should not accumulate their way
|
|
#: to a warning, so the run resets after a quiet spell.
|
|
RUN_TTL_SECONDS = 24 * 3600
|
|
|
|
#: Ceiling on tracked senders. A busy mesh sees thousands of names over a
|
|
#: long uptime, and an unbounded dict here would be a slow leak in the
|
|
#: channel-message path. Expired runs are swept first; if that is not
|
|
#: enough, the least recently seen are dropped, which only delays a
|
|
#: warning for a sender who had stopped talking anyway.
|
|
MAX_TRACKED_SENDERS = 2000
|
|
|
|
def __init__(self, bot: Any) -> None:
|
|
self.bot = bot
|
|
self.logger = bot.logger
|
|
self.settings = load_settings(getattr(bot, "config", None))
|
|
self._senders: dict[str, _SenderState] = {}
|
|
# Written-through from the DB on first use so a restart cannot reset the
|
|
# mesh-wide cooldown and let a burst of warnings out.
|
|
self._last_warning_monotonic: Optional[float] = None
|
|
self._last_warning_wall: Optional[datetime] = None
|
|
self._budget_loaded = False
|
|
if self.settings.enabled:
|
|
self.logger.info(
|
|
"Region warnings enabled (%s, delivery=%s, cap=%s/day)",
|
|
"dry run" if self.settings.dry_run else "live",
|
|
self.settings.delivery,
|
|
self.settings.max_warnings_per_day or "unlimited",
|
|
)
|
|
|
|
# -- configuration -----------------------------------------------------
|
|
|
|
def reload_config(self) -> None:
|
|
"""Re-read ``[Region_Warnings]`` after a hot config reload."""
|
|
self.settings = load_settings(getattr(self.bot, "config", None))
|
|
|
|
# -- observation -------------------------------------------------------
|
|
|
|
async def observe(
|
|
self,
|
|
*,
|
|
verdict: str,
|
|
sender_id: Optional[str],
|
|
sender_pubkey: Optional[str],
|
|
channel: Optional[str],
|
|
) -> None:
|
|
"""Record one channel message's scope verdict and warn if it earns one.
|
|
|
|
Never raises: this sits on the channel-message path and a bookkeeping
|
|
failure must not cost the message.
|
|
"""
|
|
try:
|
|
if verdict not in _VERDICT_COLUMN:
|
|
return
|
|
if self.settings.track_traffic:
|
|
self._record_verdict(verdict, channel)
|
|
if verdict != VERDICT_GLOBAL:
|
|
if verdict == VERDICT_SCOPED and sender_id:
|
|
# A scoped message proves the sender has a region set now,
|
|
# so an earlier unscoped run is stale evidence.
|
|
self._senders.pop(sender_id, None)
|
|
return
|
|
await self._consider_warning(sender_id, sender_pubkey, channel)
|
|
except Exception:
|
|
self.logger.exception("Region warning monitor failed on a channel message")
|
|
|
|
def _record_verdict(self, verdict: str, channel: Optional[str]) -> None:
|
|
column = _VERDICT_COLUMN[verdict]
|
|
today = _local_today(getattr(self.bot, "config", None), self.logger).isoformat()
|
|
name = (channel or "").strip() or "(unknown)"
|
|
db_manager = getattr(self.bot, "db_manager", None)
|
|
if not db_manager:
|
|
return
|
|
with db_manager.connection() as conn:
|
|
# Column name comes from the _VERDICT_COLUMN table, never from input.
|
|
conn.execute(
|
|
f"INSERT INTO region_scope_daily (date, channel, {column}) VALUES (?, ?, 1) "
|
|
f"ON CONFLICT(date, channel) DO UPDATE SET {column} = {column} + 1",
|
|
(today, name),
|
|
)
|
|
conn.commit()
|
|
|
|
# -- warning decision --------------------------------------------------
|
|
|
|
async def _consider_warning(
|
|
self,
|
|
sender_id: Optional[str],
|
|
sender_pubkey: Optional[str],
|
|
channel: Optional[str],
|
|
) -> None:
|
|
settings = self.settings
|
|
if not settings.enabled or not sender_id:
|
|
return
|
|
if not settings.monitors_channel(channel):
|
|
return
|
|
|
|
cmd_mgr = getattr(self.bot, "command_manager", None)
|
|
if cmd_mgr is None:
|
|
return
|
|
if cmd_mgr.is_user_banned(sender_id):
|
|
return
|
|
if self._is_self(sender_id, sender_pubkey):
|
|
return
|
|
# channelpause silences the bot on channels; a warning is a bot response
|
|
# and has no business being the one thing that keeps talking.
|
|
if not getattr(self.bot, "channel_responses_enabled", True):
|
|
return
|
|
|
|
state = self._senders.get(sender_id)
|
|
now = time.monotonic()
|
|
if state is None or (now - state.counted_at) > self.RUN_TTL_SECONDS:
|
|
state = _SenderState()
|
|
self._senders[sender_id] = state
|
|
self._evict_stale_senders(now)
|
|
state.unscoped_seen += 1
|
|
state.counted_at = now
|
|
if state.unscoped_seen < settings.min_unscoped_messages:
|
|
return
|
|
|
|
self._load_budget_once()
|
|
|
|
if self._mesh_cooldown_active(now):
|
|
return
|
|
if self._sender_cooldown_active(sender_id):
|
|
return
|
|
if self._daily_cap_reached():
|
|
self.logger.debug(
|
|
"Region warning for %s withheld: daily cap of %d reached",
|
|
sender_id, settings.max_warnings_per_day,
|
|
)
|
|
return
|
|
|
|
# The run has been spent whether or not the send itself succeeds; not
|
|
# resetting it would retry on the sender's very next message.
|
|
state.unscoped_seen = 0
|
|
|
|
text = render_message(settings.message, sender_id, channel)
|
|
if not text:
|
|
return
|
|
|
|
if settings.dry_run:
|
|
self._record_event(sender_id, sender_pubkey, channel, ACTION_DRY_RUN, text)
|
|
self._mark_warning_sent(now)
|
|
self.logger.info(
|
|
"Region warning (dry run) for %s on %s: %s", sender_id, channel or "?", text
|
|
)
|
|
return
|
|
|
|
sent, detail = await self._send_warning(sender_id, channel, text)
|
|
self._record_event(
|
|
sender_id, sender_pubkey, channel,
|
|
ACTION_SENT if sent else ACTION_FAILED,
|
|
detail,
|
|
)
|
|
if sent:
|
|
# Only a successful send starts the mesh cooldown; a failed one
|
|
# spent no airtime and should not silence a sender who would.
|
|
self._mark_warning_sent(now)
|
|
|
|
async def _send_warning(
|
|
self, sender_id: str, channel: Optional[str], text: str
|
|
) -> tuple[bool, str]:
|
|
cmd_mgr = self.bot.command_manager
|
|
if self.settings.delivery == DELIVERY_CHANNEL:
|
|
if not channel:
|
|
return False, "no channel to reply on"
|
|
body = truncate_to_bytes(text, self._channel_body_limit())
|
|
# Deliberately global scope: the recipient is by definition not
|
|
# inside any region the bot replies under, so a scoped reply would
|
|
# never reach them.
|
|
ok = await cmd_mgr.send_channel_message(
|
|
channel, body, skip_user_rate_limit=True, scope="*"
|
|
)
|
|
return ok, body if ok else f"channel send failed: {body}"
|
|
|
|
body = truncate_to_bytes(text, DM_BODY_LIMIT)
|
|
ok = await cmd_mgr.send_dm(sender_id, body, skip_user_rate_limit=True)
|
|
return ok, body if ok else f"DM send failed (contact unknown or radio busy): {body}"
|
|
|
|
def _channel_body_limit(self) -> int:
|
|
"""Channel body budget for a global-scope send from this node."""
|
|
try:
|
|
from modules.models import MeshMessage
|
|
|
|
probe = MeshMessage(content="", channel="", is_dm=False, reply_scope="")
|
|
return int(self.bot.command_manager.get_max_message_length(probe))
|
|
except Exception:
|
|
return 130
|
|
|
|
# -- gating ------------------------------------------------------------
|
|
|
|
def _evict_stale_senders(self, now: float) -> None:
|
|
"""Keep ``_senders`` bounded (see ``MAX_TRACKED_SENDERS``)."""
|
|
if len(self._senders) <= self.MAX_TRACKED_SENDERS:
|
|
return
|
|
for name, state in list(self._senders.items()):
|
|
if (now - state.counted_at) > self.RUN_TTL_SECONDS:
|
|
del self._senders[name]
|
|
overflow = len(self._senders) - self.MAX_TRACKED_SENDERS
|
|
if overflow <= 0:
|
|
return
|
|
oldest = sorted(self._senders.items(), key=lambda kv: kv[1].counted_at)
|
|
for name, _state in oldest[:overflow]:
|
|
del self._senders[name]
|
|
|
|
def _now(self) -> datetime:
|
|
return local_now(getattr(self.bot, "config", None), self.logger)
|
|
|
|
def _is_self(self, sender_id: str, sender_pubkey: Optional[str]) -> bool:
|
|
"""Whether this message came from the bot's own node.
|
|
|
|
Checked by public key first: the configured name can lag the device's,
|
|
and another node is free to call itself whatever it likes.
|
|
"""
|
|
own_key = self._own_public_key()
|
|
prefix = (sender_pubkey or "").strip().lower()
|
|
if own_key and prefix:
|
|
return own_key.startswith(prefix) or prefix.startswith(own_key)
|
|
try:
|
|
bot_name = self.bot.config.get("Bot", "bot_name", fallback="") or ""
|
|
except Exception:
|
|
bot_name = ""
|
|
return bool(bot_name) and sender_id.strip().lower() == bot_name.strip().lower()
|
|
|
|
def _own_public_key(self) -> str:
|
|
"""This node's public key as lowercase hex, or "" when unavailable."""
|
|
try:
|
|
self_info = getattr(getattr(self.bot, "meshcore", None), "self_info", None)
|
|
if self_info is None:
|
|
return ""
|
|
if isinstance(self_info, dict):
|
|
key = self_info.get("public_key", "")
|
|
else:
|
|
key = getattr(self_info, "public_key", "")
|
|
if isinstance(key, (bytes, bytearray)):
|
|
return bytes(key).hex()
|
|
return str(key or "").strip().lower()
|
|
except Exception:
|
|
return ""
|
|
|
|
def _load_budget_once(self) -> None:
|
|
"""Seed the mesh cooldown from the DB so a restart cannot bypass it."""
|
|
if self._budget_loaded:
|
|
return
|
|
self._budget_loaded = True
|
|
db_manager = getattr(self.bot, "db_manager", None)
|
|
if not db_manager:
|
|
return
|
|
try:
|
|
rows = db_manager.execute_query(
|
|
"SELECT MAX(created_at) AS last_at FROM region_warning_events "
|
|
"WHERE action IN (?, ?)",
|
|
(ACTION_SENT, ACTION_DRY_RUN),
|
|
)
|
|
except Exception:
|
|
return
|
|
raw = rows[0].get("last_at") if rows else None
|
|
if not raw:
|
|
return
|
|
parsed = _parse_timestamp(raw)
|
|
if parsed is None:
|
|
return
|
|
self._last_warning_wall = parsed
|
|
elapsed = (self._now() - parsed).total_seconds()
|
|
if elapsed >= 0:
|
|
self._last_warning_monotonic = time.monotonic() - elapsed
|
|
|
|
def _mark_warning_sent(self, now: float) -> None:
|
|
self._last_warning_monotonic = now
|
|
self._last_warning_wall = self._now()
|
|
|
|
def _mesh_cooldown_active(self, now: float) -> bool:
|
|
window = self.settings.mesh_cooldown_minutes * 60.0
|
|
if window <= 0 or self._last_warning_monotonic is None:
|
|
return False
|
|
return (now - self._last_warning_monotonic) < window
|
|
|
|
def _sender_cooldown_active(self, sender_id: str) -> bool:
|
|
hours = self.settings.per_sender_cooldown_hours
|
|
if hours <= 0:
|
|
return False
|
|
db_manager = getattr(self.bot, "db_manager", None)
|
|
if not db_manager:
|
|
return False
|
|
try:
|
|
rows = db_manager.execute_query(
|
|
"SELECT MAX(created_at) AS last_at FROM region_warning_events "
|
|
"WHERE sender_id = ? AND action IN (?, ?)",
|
|
(sender_id, ACTION_SENT, ACTION_DRY_RUN),
|
|
)
|
|
except Exception:
|
|
return False
|
|
raw = rows[0].get("last_at") if rows else None
|
|
parsed = _parse_timestamp(raw) if raw else None
|
|
if parsed is None:
|
|
return False
|
|
return (self._now() - parsed) < timedelta(hours=hours)
|
|
|
|
def _daily_cap_reached(self) -> bool:
|
|
cap = self.settings.max_warnings_per_day
|
|
if cap <= 0:
|
|
return False
|
|
db_manager = getattr(self.bot, "db_manager", None)
|
|
if not db_manager:
|
|
return False
|
|
today = self._now().date().isoformat()
|
|
try:
|
|
# Attempts, not successes: see warning_budget.
|
|
rows = db_manager.execute_query(
|
|
"SELECT COUNT(*) AS used FROM region_warning_events "
|
|
"WHERE date(created_at) = ?",
|
|
(today,),
|
|
)
|
|
except Exception:
|
|
return False
|
|
used = int(rows[0].get("used") or 0) if rows else 0
|
|
return used >= cap
|
|
|
|
# -- persistence -------------------------------------------------------
|
|
|
|
def _record_event(
|
|
self,
|
|
sender_id: str,
|
|
sender_pubkey: Optional[str],
|
|
channel: Optional[str],
|
|
action: str,
|
|
detail: str,
|
|
) -> None:
|
|
db_manager = getattr(self.bot, "db_manager", None)
|
|
if not db_manager:
|
|
return
|
|
try:
|
|
with db_manager.connection() as conn:
|
|
conn.execute(
|
|
"INSERT INTO region_warning_events "
|
|
"(created_at, sender_id, sender_pubkey, channel, delivery, action, detail) "
|
|
"VALUES (?, ?, ?, ?, ?, ?, ?)",
|
|
(
|
|
local_now(getattr(self.bot, "config", None), self.logger)
|
|
.isoformat(sep=" ", timespec="seconds"),
|
|
sender_id,
|
|
(sender_pubkey or "")[:64] or None,
|
|
channel,
|
|
self.settings.delivery,
|
|
action,
|
|
detail[:400] if detail else None,
|
|
),
|
|
)
|
|
conn.commit()
|
|
except Exception:
|
|
self.logger.exception("Failed to record region warning event")
|
|
|
|
|
|
def _parse_timestamp(raw: Any) -> Optional[datetime]:
|
|
"""Parse a stored event timestamp, tolerating ``T`` or space separators."""
|
|
if not raw:
|
|
return None
|
|
text = str(raw).strip().replace("T", " ")
|
|
for fmt in ("%Y-%m-%d %H:%M:%S.%f", "%Y-%m-%d %H:%M:%S", "%Y-%m-%d %H:%M"):
|
|
try:
|
|
return datetime.strptime(text, fmt)
|
|
except ValueError:
|
|
continue
|
|
return None
|