mirror of
https://github.com/maziggy/bambuddy.git
synced 2026-09-30 11:12:35 +02:00
2788 lines
112 KiB
Python
2788 lines
112 KiB
Python
"""Notification service for sending push notifications via various providers."""
|
|
|
|
import asyncio
|
|
import html
|
|
import json
|
|
import logging
|
|
import re
|
|
import smtplib
|
|
from datetime import datetime, timedelta, timezone
|
|
from email.mime.image import MIMEImage
|
|
from email.mime.multipart import MIMEMultipart
|
|
from email.mime.text import MIMEText
|
|
from typing import Any
|
|
from urllib.parse import quote
|
|
|
|
import httpx
|
|
from sqlalchemy import select
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from backend.app.models.archive import PrintArchive
|
|
from backend.app.models.notification import (
|
|
NotificationDigestQueue,
|
|
NotificationLog,
|
|
NotificationProvider,
|
|
TelegramPendingVerdict,
|
|
)
|
|
from backend.app.models.notification_template import NotificationTemplate
|
|
from backend.app.services.print_confirmation import one_tap_url
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# Honest User-Agent — matches the convention used by every other outbound
|
|
# httpx client in the codebase (bambu_cloud, makerworld, firmware_check,
|
|
# inventory). Previously this client leaked python-httpx/<version>, which
|
|
# was both inconsistent with the rest of the project and a more obvious
|
|
# bot signature for upstream WAFs.
|
|
_USER_AGENT = "Bambuddy/1.0 (+https://github.com/maziggy/bambuddy)"
|
|
|
|
# Appended to a Telegram outcome prompt in reaction mode (#3046); the message
|
|
# itself carries no buttons there, so the user needs telling how to answer.
|
|
TELEGRAM_REACTION_HINT = "React with \U0001f44d or \U0001f44e to record the outcome."
|
|
|
|
|
|
def telegram_markdown_escape(message: str) -> str:
|
|
"""Escape underscores in the message body so Telegram Markdown parsing
|
|
doesn't break on job names like "A1_plate_8" or error codes like
|
|
"0300_0001". The title is already wrapped in *bold* markers, so only
|
|
escape after the first newline. Shared with the reaction poller, which
|
|
re-sends the stored body when it edits the prompt (#3046)."""
|
|
if "\n" in message:
|
|
title_part, body_part = message.split("\n", 1)
|
|
body_part = body_part.replace("_", "\\_")
|
|
return f"{title_part}\n{body_part}"
|
|
return message
|
|
|
|
|
|
def _looks_like_cloudflare_challenge(response: httpx.Response) -> bool:
|
|
"""Return True if ``response`` looks like a Cloudflare mitigation
|
|
interstitial (JS challenge / managed challenge / block page) rather
|
|
than a legitimate response passed through Cloudflare.
|
|
|
|
Self-hosted servers behind Cloudflare (Tunnel, "Bot Fight Mode", or
|
|
"Under Attack" mode) intercept non-browser clients at the edge and
|
|
return a challenge HTML page instead of forwarding to the origin —
|
|
so we never reach the user's actual ntfy / webhook backend.
|
|
Cloudflare cannot be defeated from a Python client; the user has to
|
|
add a security-skip rule on their side. We detect the shape so the
|
|
UI can tell them that, instead of dumping the raw HTML.
|
|
|
|
Detection deliberately does NOT rely on ``Server: cloudflare`` alone
|
|
— Cloudflare adds that header to every response it proxies (success
|
|
AND legitimate origin errors), so a real 401 "wrong token" from a
|
|
CF-fronted ntfy would false-positive into a misleading "your CF is
|
|
blocking" message. Reliable signals: the ``cf-mitigated`` header
|
|
(set only when CF actively mitigates) and the challenge body shape.
|
|
"""
|
|
if response.headers.get("cf-mitigated"):
|
|
return True
|
|
content_type = (response.headers.get("content-type") or "").lower()
|
|
if "html" not in content_type:
|
|
return False
|
|
body = (response.text or "")[:1024].lower()
|
|
# "Just a moment..." is Cloudflare's universal challenge-page title
|
|
# (managed challenge, JS challenge, Under Attack mode). Combined with
|
|
# an HTML content-type this is unambiguous — no legitimate ntfy or
|
|
# webhook backend returns HTML with that title. ``cf-chl-*`` and
|
|
# ``challenge-platform`` cover newer / non-default CF templates.
|
|
return "just a moment" in body or "cf-chl-bypass" in body or "cf-chl-opt" in body or "challenge-platform" in body
|
|
|
|
|
|
def _assert_safe_provider_url(url: str, *, label: str) -> str | None:
|
|
"""Validate a provider URL taken from user-supplied config.
|
|
|
|
Returns an error message on rejection, or None when the URL is
|
|
acceptable — the ``_send_*`` methods return ``tuple[bool, str]`` rather
|
|
than raising, so a message is more useful here than an exception.
|
|
|
|
Uses the LAN-service policy: self-hosting ntfy, Bark, Gotify or a webhook
|
|
receiver on the home LAN is normal and must keep working, so loopback and
|
|
RFC-1918 stay permitted. Cloud-metadata endpoints, numeric-encoded IPs and
|
|
non-HTTP schemes are rejected.
|
|
"""
|
|
from backend.app.api.routes._url_safety import assert_safe_lan_service_url
|
|
|
|
try:
|
|
assert_safe_lan_service_url(url, label=label)
|
|
except ValueError as exc:
|
|
return str(exc)
|
|
return None
|
|
|
|
|
|
def _opaque_http_failure(response: httpx.Response, *, label: str) -> str:
|
|
"""Failure message for a provider whose destination host the user supplies.
|
|
|
|
The response body is deliberately **not** returned to the caller. Provider
|
|
URLs are configurable by anyone holding ``NOTIFICATIONS_CREATE`` — which
|
|
the default Operators group carries and which does not imply
|
|
``SETTINGS_UPDATE`` — and ``POST /notifications/test-config`` accepts a URL
|
|
straight from the request body without persisting anything. Echoing the
|
|
response body there turned an intended "does my webhook work?" check into
|
|
an authenticated read primitive against any host the Bambuddy process can
|
|
reach, including services that are not exposed to the network at all.
|
|
|
|
Providers whose host Bambuddy hardcodes (Pushover, Telegram, CallMeBot)
|
|
keep returning the upstream body — there is no trust boundary to cross
|
|
when the destination cannot be influenced.
|
|
|
|
The body is logged at debug level, where it stays available to whoever
|
|
already administers the host without being handed back over the API.
|
|
"""
|
|
logger.debug(
|
|
"%s delivery failed with HTTP %s; body: %s",
|
|
label,
|
|
response.status_code,
|
|
(response.text or "")[:200],
|
|
)
|
|
return f"HTTP {response.status_code} from the configured {label} (see server logs at debug level for details)"
|
|
|
|
|
|
class NotificationService:
|
|
"""Service for sending notifications through various providers."""
|
|
|
|
def __init__(self):
|
|
self._http_client: httpx.AsyncClient | None = None
|
|
self._template_cache: dict[str, NotificationTemplate] = {}
|
|
self._digest_scheduler_task: asyncio.Task | None = None
|
|
self._last_digest_check: str = "" # "HH:MM" to avoid duplicate checks
|
|
|
|
async def _get_client(self) -> httpx.AsyncClient:
|
|
"""Get or create HTTP client.
|
|
|
|
The connect timeout is deliberately far shorter than the rest. A flat
|
|
30 s meant that when a site's internet went down, every alarm spent a
|
|
full 30 s inside ``connect`` — longer than SQLite's 15 s
|
|
``busy_timeout`` — and any other task that wanted to write during that
|
|
window failed with "database is locked" (#2770). Reaching a host either
|
|
works in a couple of seconds or is not going to; sending the body is the
|
|
part that legitimately takes time, so read/write keep the old 30 s and
|
|
an image upload on a slow uplink is unaffected.
|
|
"""
|
|
if self._http_client is None or self._http_client.is_closed:
|
|
self._http_client = httpx.AsyncClient(
|
|
timeout=httpx.Timeout(30.0, connect=5.0),
|
|
headers={"User-Agent": _USER_AGENT},
|
|
)
|
|
return self._http_client
|
|
|
|
async def close(self):
|
|
"""Close HTTP client."""
|
|
if self._http_client and not self._http_client.is_closed:
|
|
await self._http_client.aclose()
|
|
|
|
def _is_in_quiet_hours(self, provider: NotificationProvider) -> bool:
|
|
"""Check if current time is within provider's quiet hours."""
|
|
if not provider.quiet_hours_enabled:
|
|
return False
|
|
|
|
if not provider.quiet_hours_start or not provider.quiet_hours_end:
|
|
return False
|
|
|
|
try:
|
|
now = datetime.now()
|
|
current_time = now.hour * 60 + now.minute
|
|
|
|
start_parts = provider.quiet_hours_start.split(":")
|
|
end_parts = provider.quiet_hours_end.split(":")
|
|
|
|
start_minutes = int(start_parts[0]) * 60 + int(start_parts[1])
|
|
end_minutes = int(end_parts[0]) * 60 + int(end_parts[1])
|
|
|
|
# Handle overnight quiet hours (e.g., 22:00 to 07:00)
|
|
if start_minutes > end_minutes:
|
|
# Quiet hours span midnight
|
|
return current_time >= start_minutes or current_time < end_minutes
|
|
else:
|
|
# Same day quiet hours
|
|
return start_minutes <= current_time < end_minutes
|
|
except (ValueError, TypeError, AttributeError):
|
|
logger.warning("Invalid quiet hours format for provider %s", provider.name)
|
|
return False
|
|
|
|
async def _get_template(self, db: AsyncSession, event_type: str) -> NotificationTemplate | None:
|
|
"""Get a notification template by event type.
|
|
|
|
``no_autoflush`` for the same reason as ``_get_providers_for_event``:
|
|
this read runs before the provider is contacted, and must not be the
|
|
thing that opens a write transaction on the caller's session (#2770).
|
|
"""
|
|
# Check cache first
|
|
if event_type in self._template_cache:
|
|
return self._template_cache[event_type]
|
|
|
|
with db.no_autoflush:
|
|
result = await db.execute(select(NotificationTemplate).where(NotificationTemplate.event_type == event_type))
|
|
template = result.scalar_one_or_none()
|
|
|
|
if template:
|
|
self._template_cache[event_type] = template
|
|
|
|
return template
|
|
|
|
def _render_template(self, template_str: str, variables: dict[str, Any]) -> str:
|
|
"""Render a template string with variables. Missing variables become empty."""
|
|
result = template_str
|
|
for key, value in variables.items():
|
|
result = result.replace("{" + key + "}", str(value) if value is not None else "")
|
|
# Remove any remaining unreplaced placeholders
|
|
result = re.sub(r"\{[a-z_]+\}", "", result)
|
|
return result
|
|
|
|
async def _format_eta(self, seconds: int | None, db: AsyncSession) -> str:
|
|
"""Format ETA as wall-clock time, respecting user's time_format setting."""
|
|
if not seconds or seconds <= 0:
|
|
return "Unknown"
|
|
|
|
from backend.app.api.routes.settings import get_setting
|
|
|
|
eta_time = datetime.now() + timedelta(seconds=seconds)
|
|
time_format = await get_setting(db, "time_format")
|
|
|
|
if time_format == "12h":
|
|
return eta_time.strftime("%I:%M %p").lstrip("0")
|
|
# Default to 24h for "24h", "system", or unset
|
|
return eta_time.strftime("%H:%M")
|
|
|
|
def _format_duration(self, seconds: int | None) -> str:
|
|
"""Format duration in seconds to human-readable string."""
|
|
if seconds is None:
|
|
return "Unknown"
|
|
hours = seconds // 3600
|
|
minutes = (seconds % 3600) // 60
|
|
if hours > 0:
|
|
return f"{hours}h {minutes}m"
|
|
return f"{minutes}m"
|
|
|
|
def _clean_filename(self, filename: str) -> str:
|
|
"""Extract filename and remove file extensions."""
|
|
import os
|
|
|
|
# Strip path prefix (e.g., /data/Metadata/plate_5.gcode -> plate_5.gcode)
|
|
filename = os.path.basename(filename)
|
|
# Remove common extensions
|
|
if filename.endswith(".gcode.3mf"):
|
|
return filename[:-10]
|
|
elif filename.endswith(".gcode"):
|
|
return filename[:-6]
|
|
elif filename.endswith(".3mf"):
|
|
return filename[:-4]
|
|
return filename
|
|
|
|
async def _build_message_from_template(
|
|
self, db: AsyncSession, event_type: str, variables: dict[str, Any]
|
|
) -> tuple[str, str]:
|
|
"""Build notification title and body from template."""
|
|
# Add common variables
|
|
variables["timestamp"] = datetime.now().strftime("%Y-%m-%d %H:%M")
|
|
variables["app_name"] = "Bambuddy"
|
|
|
|
template = await self._get_template(db, event_type)
|
|
if not template:
|
|
# Fallback to simple message
|
|
logger.warning("Template not found for event type: %s", event_type)
|
|
return event_type.replace("_", " ").title(), str(variables)
|
|
|
|
title = self._render_template(template.title_template, variables)
|
|
body = self._render_template(template.body_template, variables)
|
|
|
|
return title, body
|
|
|
|
async def send_test_notification(
|
|
self, provider_type: str, config: dict[str, Any], db: AsyncSession | None = None
|
|
) -> tuple[bool, str]:
|
|
"""Send a test notification to verify configuration."""
|
|
if db:
|
|
title, message = await self._build_message_from_template(db, "test", {})
|
|
else:
|
|
title = "Bambuddy Test"
|
|
message = "This is a test notification. If you see this, notifications are working!"
|
|
|
|
try:
|
|
if provider_type == "callmebot":
|
|
return await self._send_callmebot(config, f"{title}\n{message}")
|
|
elif provider_type == "ntfy":
|
|
return await self._send_ntfy(config, title, message)
|
|
elif provider_type == "pushover":
|
|
return await self._send_pushover(config, title, message)
|
|
elif provider_type == "telegram":
|
|
return await self._send_telegram(config, f"*{title}*\n{message}")
|
|
elif provider_type == "email":
|
|
return await self._send_email(config, title, message)
|
|
elif provider_type == "discord":
|
|
return await self._send_discord(config, title, message)
|
|
elif provider_type == "webhook":
|
|
return await self._send_webhook(config, title, message)
|
|
elif provider_type == "homeassistant":
|
|
return await self._send_homeassistant(config, title, message, db=db)
|
|
elif provider_type == "bark":
|
|
return await self._send_bark(config, title, message)
|
|
else:
|
|
return False, f"Unknown provider type: {provider_type}"
|
|
except Exception as e:
|
|
logger.exception("Error sending test notification via %s", provider_type)
|
|
return False, str(e)
|
|
|
|
async def _send_callmebot(self, config: dict, message: str) -> tuple[bool, str]:
|
|
"""Send notification via CallMeBot (WhatsApp)."""
|
|
phone = config.get("phone", "").strip()
|
|
apikey = config.get("apikey", "").strip()
|
|
|
|
if not phone or not apikey:
|
|
return False, "Phone number and API key are required"
|
|
|
|
# URL encode the message
|
|
encoded_message = quote(message)
|
|
url = f"https://api.callmebot.com/whatsapp.php?phone={phone}&text={encoded_message}&apikey={apikey}"
|
|
|
|
client = await self._get_client()
|
|
response = await client.get(url)
|
|
|
|
if response.status_code == 200:
|
|
return True, "Message sent successfully"
|
|
else:
|
|
return False, f"HTTP {response.status_code}: {response.text[:200]}"
|
|
|
|
async def _send_bark(self, config: dict, title: str, message: str, url: str | None = None) -> tuple[bool, str]:
|
|
"""Send notification via Bark, the self-hostable iOS push service (#1495).
|
|
|
|
POSTs JSON to {server}/push. Defaults to the official api.day.app
|
|
relay; a self-hosted bark-server works by overriding the server URL.
|
|
``url`` opens on tap — the outcome confirmation (#1898) deep-links
|
|
into the archive's confirmation dialog with it.
|
|
"""
|
|
server = (config.get("server") or "https://api.day.app").strip().rstrip("/")
|
|
device_key = (config.get("device_key") or "").strip()
|
|
|
|
if not device_key:
|
|
return False, "Device key is required"
|
|
|
|
url_error = _assert_safe_provider_url(server, label="Bark server URL")
|
|
if url_error:
|
|
return False, url_error
|
|
|
|
payload: dict[str, Any] = {
|
|
"device_key": device_key,
|
|
"title": title,
|
|
"body": message,
|
|
}
|
|
group = (config.get("group") or "").strip()
|
|
if group:
|
|
payload["group"] = group
|
|
sound = (config.get("sound") or "").strip()
|
|
if sound:
|
|
payload["sound"] = sound
|
|
level = (config.get("level") or "").strip()
|
|
if level in ("active", "timeSensitive", "critical", "passive"):
|
|
payload["level"] = level
|
|
if url:
|
|
payload["url"] = url
|
|
|
|
client = await self._get_client()
|
|
response = await client.post(f"{server}/push", json=payload)
|
|
|
|
if response.status_code == 200:
|
|
# bark-server can report failures inside an HTTP 200 body
|
|
# ({"code": 400, "message": ...}), so the status alone isn't proof.
|
|
try:
|
|
body = response.json()
|
|
except ValueError:
|
|
body = None
|
|
if isinstance(body, dict) and body.get("code") not in (200, None):
|
|
# Only the numeric code is echoed. A server chosen by the caller
|
|
# controls this body too, so the free-text message is a (narrow)
|
|
# read channel of the same kind _opaque_http_failure closes.
|
|
logger.debug("Bark reported error %s: %s", body.get("code"), str(body.get("message"))[:200])
|
|
return False, f"Bark error {body.get('code')} (see server logs at debug level for details)"
|
|
return True, "Message sent successfully"
|
|
return False, _opaque_http_failure(response, label="Bark server")
|
|
|
|
async def _send_ntfy(
|
|
self,
|
|
config: dict,
|
|
title: str,
|
|
message: str,
|
|
image_data: bytes | None = None,
|
|
event_type: str | None = None,
|
|
actions: str | None = None,
|
|
) -> tuple[bool, str]:
|
|
"""Send notification via ntfy.
|
|
|
|
``actions`` is a pre-built value for ntfy's Actions header (simple
|
|
format), used by the outcome-confirmation event (#1898) to put
|
|
one-tap Good/Reject buttons directly into the push notification.
|
|
"""
|
|
server = config.get("server", "https://ntfy.sh").rstrip("/")
|
|
topic = config.get("topic", "").strip()
|
|
auth_token = config.get("auth_token", "").strip()
|
|
|
|
if not topic:
|
|
return False, "Topic is required"
|
|
|
|
url_error = _assert_safe_provider_url(server, label="ntfy server URL")
|
|
if url_error:
|
|
return False, url_error
|
|
|
|
url = f"{server}/{topic}"
|
|
# ntfy reads Title/Message from HTTP headers. httpx enforces ASCII
|
|
# for str header values, but printer names and filenames can contain
|
|
# non-ASCII characters (e.g. accented letters, CJK). Passing bytes
|
|
# bypasses the ASCII check — ntfy handles UTF-8 headers correctly.
|
|
headers: dict[str, str | bytes] = {"Title": title.encode("utf-8")}
|
|
|
|
# Per-event Priority header (#990). Only set when the user has
|
|
# explicitly mapped this event to a 1-5 value; otherwise fall through
|
|
# to the ntfy server's default so existing setups stay unchanged.
|
|
#
|
|
# The map is keyed by the provider's toggle column ("on_print_failed"),
|
|
# because that is what the dialog builds its rows from -- but every
|
|
# sender is called with the bare event name ("print_failed"), so the
|
|
# lookup used to miss for every real notification and hit only in tests
|
|
# that called this method with the prefixed name (issue #3139). Both
|
|
# spellings are accepted, which also leaves stored configs untouched.
|
|
event_priorities = config.get("event_priorities") or {}
|
|
if event_type and isinstance(event_priorities, dict):
|
|
raw = event_priorities.get(event_type)
|
|
if raw is None and not event_type.startswith("on_"):
|
|
raw = event_priorities.get(f"on_{event_type}")
|
|
try:
|
|
priority = int(raw) if raw is not None else None
|
|
except (TypeError, ValueError):
|
|
priority = None
|
|
if priority is not None and 1 <= priority <= 5:
|
|
headers["Priority"] = str(priority)
|
|
|
|
if auth_token:
|
|
headers["Authorization"] = f"Bearer {auth_token}"
|
|
|
|
if actions:
|
|
headers["Actions"] = actions
|
|
|
|
client = await self._get_client()
|
|
|
|
if image_data:
|
|
# ntfy supports image attachments via multipart form-data.
|
|
# HTTP headers cannot contain newlines, but ntfy interprets
|
|
# literal \n (backslash-n) as newlines in the Message header.
|
|
headers["Filename"] = "photo.jpg"
|
|
headers["Message"] = message.replace("\n", "\\n").encode("utf-8")
|
|
response = await client.put(url, content=image_data, headers=headers)
|
|
|
|
if response.status_code == 400 and "attachments not allowed" in response.text:
|
|
# Server has attachments disabled — retry without the image
|
|
headers.pop("Filename", None)
|
|
headers.pop("Message", None)
|
|
response = await client.post(url, content=message.encode("utf-8"), headers=headers)
|
|
else:
|
|
response = await client.post(url, content=message.encode("utf-8"), headers=headers)
|
|
|
|
if response.status_code in (200, 204):
|
|
return True, "Message sent successfully"
|
|
if _looks_like_cloudflare_challenge(response):
|
|
return False, (
|
|
f"HTTP {response.status_code} — ntfy server is behind a Cloudflare "
|
|
"challenge. Bambuddy was served the JS challenge page instead of "
|
|
"reaching ntfy. Cloudflare cannot be solved from a backend; add a "
|
|
"Cloudflare security-skip rule for this hostname, disable Bot "
|
|
"Fight Mode, or front the server with Cloudflare Access using a "
|
|
"service token. (#1534)"
|
|
)
|
|
return False, _opaque_http_failure(response, label="ntfy server")
|
|
|
|
async def _send_pushover(
|
|
self,
|
|
config: dict,
|
|
title: str,
|
|
message: str,
|
|
image_data: bytes | None = None,
|
|
url: str | None = None,
|
|
url_title: str | None = None,
|
|
) -> tuple[bool, str]:
|
|
"""Send notification via Pushover.
|
|
|
|
Args:
|
|
config: Provider configuration with user_key, app_token, priority
|
|
title: Notification title
|
|
message: Notification body
|
|
image_data: Optional JPEG image bytes to attach (max 2.5MB)
|
|
url: Optional supplementary URL shown under the message
|
|
url_title: Optional label for that URL
|
|
"""
|
|
user_key = config.get("user_key", "").strip()
|
|
app_token = config.get("app_token", "").strip()
|
|
try:
|
|
priority = int(config.get("priority", 0))
|
|
except (TypeError, ValueError):
|
|
priority = 0
|
|
|
|
if not user_key or not app_token:
|
|
return False, "User key and app token are required"
|
|
|
|
api_url = "https://api.pushover.net/1/messages.json"
|
|
data = {
|
|
"token": app_token,
|
|
"user": user_key,
|
|
"title": title,
|
|
"message": message,
|
|
"priority": priority,
|
|
}
|
|
if url:
|
|
data["url"] = url
|
|
if url_title:
|
|
data["url_title"] = url_title
|
|
|
|
# Emergency priority (2) keeps re-alerting until acknowledged, so
|
|
# Pushover *requires* retry (how often, >= 30s) and expire (when to
|
|
# give up, <= 10800s). Without them the API rejects the message. Only
|
|
# send them at priority 2 — Pushover ignores them at other priorities.
|
|
if priority == 2:
|
|
try:
|
|
retry = int(config.get("retry", 60))
|
|
except (TypeError, ValueError):
|
|
retry = 60
|
|
try:
|
|
expire = int(config.get("expire", 3600))
|
|
except (TypeError, ValueError):
|
|
expire = 3600
|
|
data["retry"] = max(30, min(retry, 10800))
|
|
data["expire"] = max(30, min(expire, 10800))
|
|
|
|
client = await self._get_client()
|
|
|
|
if image_data:
|
|
# Pushover supports image attachments via multipart form-data
|
|
files = {"attachment": ("photo.jpg", image_data, "image/jpeg")}
|
|
response = await client.post(api_url, data=data, files=files)
|
|
else:
|
|
response = await client.post(api_url, data=data)
|
|
|
|
if response.status_code == 200:
|
|
return True, "Message sent successfully"
|
|
else:
|
|
try:
|
|
error_data = response.json()
|
|
errors = error_data.get("errors", [])
|
|
return False, f"Pushover error: {', '.join(errors)}"
|
|
except Exception:
|
|
return False, f"HTTP {response.status_code}: {response.text[:200]}"
|
|
|
|
async def _send_telegram(
|
|
self,
|
|
config: dict,
|
|
message: str,
|
|
image_data: bytes | None = None,
|
|
buttons: list[dict] | None = None,
|
|
link_preview: bool = True,
|
|
) -> tuple[bool, str]:
|
|
"""Send notification via Telegram bot.
|
|
|
|
``buttons`` is one row of inline URL buttons (``{"text", "url"}``
|
|
entries), used by the outcome-confirmation event (#1898) to put
|
|
one-tap Good/Reject under the message. ``link_preview=False`` asks
|
|
Telegram not to fetch the first URL in the text for a preview card.
|
|
"""
|
|
ok, status, _ = await self._send_telegram_message(
|
|
config, message, image_data=image_data, buttons=buttons, link_preview=link_preview
|
|
)
|
|
return ok, status
|
|
|
|
async def _send_telegram_message(
|
|
self,
|
|
config: dict,
|
|
message: str,
|
|
image_data: bytes | None = None,
|
|
buttons: list[dict] | None = None,
|
|
link_preview: bool = True,
|
|
) -> tuple[bool, str, dict | None]:
|
|
"""Send via Telegram and also hand back the Bot API ``result`` (the sent Message).
|
|
|
|
The outcome confirmation in reaction mode (#3046) needs the
|
|
``message_id`` from it to recognise the reaction later; every other
|
|
caller goes through ``_send_telegram`` and ignores it.
|
|
"""
|
|
bot_token = config.get("bot_token", "").strip()
|
|
chat_id = config.get("chat_id", "").strip()
|
|
|
|
if not bot_token or not chat_id:
|
|
return False, "Bot token and chat ID are required", None
|
|
|
|
# Optional forum topic (#1518). Telegram expects message_thread_id as an
|
|
# integer in the JSON sendMessage body — a string 400s there even though
|
|
# the multipart sendPhoto call below would happily accept one. Coerce it
|
|
# once, up front, so both call sites agree and a bad value fails loudly
|
|
# instead of silently breaking only the text notifications.
|
|
thread_id_raw = str(config.get("message_thread_id") or "").strip()
|
|
message_thread_id: int | None = None
|
|
if thread_id_raw:
|
|
try:
|
|
message_thread_id = int(thread_id_raw)
|
|
except ValueError:
|
|
return False, f"Invalid message thread ID: {thread_id_raw!r} is not a number", None
|
|
|
|
message = telegram_markdown_escape(message)
|
|
|
|
client = await self._get_client()
|
|
|
|
async def _post(with_buttons: bool):
|
|
if image_data:
|
|
# Use sendPhoto to attach the thumbnail with the caption
|
|
url = f"https://api.telegram.org/bot{bot_token}/sendPhoto"
|
|
form: dict[str, Any] = {"chat_id": chat_id, "caption": message, "parse_mode": "Markdown"}
|
|
if message_thread_id is not None:
|
|
form["message_thread_id"] = message_thread_id
|
|
if with_buttons:
|
|
# Multipart form fields are strings — reply_markup goes JSON-encoded.
|
|
form["reply_markup"] = json.dumps({"inline_keyboard": [buttons]})
|
|
return await client.post(
|
|
url,
|
|
data=form,
|
|
files={"photo": ("photo.jpg", image_data, "image/jpeg")},
|
|
)
|
|
|
|
url = f"https://api.telegram.org/bot{bot_token}/sendMessage"
|
|
payload: dict[str, Any] = {
|
|
"chat_id": chat_id,
|
|
"text": message,
|
|
"parse_mode": "Markdown",
|
|
}
|
|
if not link_preview:
|
|
payload["disable_web_page_preview"] = True
|
|
if message_thread_id is not None:
|
|
payload["message_thread_id"] = message_thread_id
|
|
if with_buttons:
|
|
payload["reply_markup"] = {"inline_keyboard": [buttons]}
|
|
return await client.post(url, json=payload)
|
|
|
|
def _failure(resp) -> str | None:
|
|
"""What Telegram objected to, or None when the send went through."""
|
|
if resp.status_code != 200:
|
|
return f"HTTP {resp.status_code}: {resp.text[:200]}"
|
|
result = resp.json()
|
|
if result.get("ok"):
|
|
return None
|
|
return f"Telegram error: {result.get('description', 'Unknown error')}"
|
|
|
|
response = await _post(bool(buttons))
|
|
failure = _failure(response)
|
|
if failure and buttons:
|
|
# Telegram validates every inline-keyboard URL and refuses the whole
|
|
# send when one of them is not a URL it accepts — which is what an
|
|
# install without a public external_url produces for the #1898
|
|
# verdict links. Dropping the buttons is survivable; dropping the
|
|
# message the user is waiting for is not.
|
|
logger.warning("Telegram refused the message with inline buttons (%s); retrying without them", failure)
|
|
response = await _post(False)
|
|
failure = _failure(response)
|
|
|
|
if failure:
|
|
return False, failure, None
|
|
sent = response.json().get("result")
|
|
return True, "Message sent successfully", sent if isinstance(sent, dict) else None
|
|
|
|
async def _send_telegram_confirm_request(
|
|
self,
|
|
provider: NotificationProvider,
|
|
config: dict,
|
|
message: str,
|
|
db: AsyncSession | None,
|
|
image_data: bytes | None,
|
|
buttons: list[dict] | None,
|
|
archive_id: int | None,
|
|
) -> tuple[bool, str]:
|
|
"""Deliver the outcome prompt the way the provider's verdict mode asks (#3046).
|
|
|
|
"buttons" is the plain #1898 delivery. In "reactions" the inline
|
|
keyboard is dropped and the user answers with a thumbs-up/down on the
|
|
message itself; "both" keeps the keyboard as well. In either of those
|
|
the sent message is remembered in telegram_pending_verdicts, together
|
|
with the confirm token its links carry, so the reaction poller can map
|
|
the reaction back to the archive for as long as that token is live.
|
|
"""
|
|
mode = provider.telegram_verdict_mode or "buttons"
|
|
# Telegram's servers GET the first URL in the text to build a preview
|
|
# card. An outcome prompt whose edited body still carries {good_url}
|
|
# would have that fetch answer the question before the operator saw
|
|
# it, so the preview is off for the prompt in every mode.
|
|
if mode == "buttons":
|
|
return await self._send_telegram(
|
|
config, message, image_data=image_data, buttons=buttons, link_preview=False
|
|
)
|
|
|
|
if mode == "reactions":
|
|
buttons = None
|
|
message = f"{message}\n\n{TELEGRAM_REACTION_HINT}"
|
|
ok, status, sent = await self._send_telegram_message(
|
|
config, message, image_data=image_data, buttons=buttons, link_preview=False
|
|
)
|
|
if not ok:
|
|
return ok, status
|
|
|
|
message_id = (sent or {}).get("message_id")
|
|
if not isinstance(message_id, int) or message_id <= 0:
|
|
logger.warning("Telegram did not return a message_id for the outcome prompt; reactions cannot be matched")
|
|
return ok, status
|
|
if db is None or archive_id is None:
|
|
return ok, status
|
|
|
|
try:
|
|
# dispatch_outcome_confirmation mints a live token right before
|
|
# sending. Without one the prompt's links are dead too, and a
|
|
# reaction must not outlive them, so there is nothing to remember.
|
|
archive = await db.get(PrintArchive, archive_id)
|
|
if archive is None or not archive.confirm_token or archive.confirm_token_used_at is not None:
|
|
logger.warning("Outcome prompt for archive %s went out without a live confirm token", archive_id)
|
|
return ok, status
|
|
db.add(
|
|
TelegramPendingVerdict(
|
|
provider_id=provider.id,
|
|
chat_id=str(((sent or {}).get("chat") or {}).get("id") or config.get("chat_id", "")).strip(),
|
|
message_id=message_id,
|
|
archive_id=archive_id,
|
|
confirm_token=archive.confirm_token,
|
|
has_caption=image_data is not None,
|
|
message_text=message,
|
|
)
|
|
)
|
|
await db.commit()
|
|
except Exception as e:
|
|
# The prompt went out; losing the reaction mapping is a degraded
|
|
# outcome, not a failed notification.
|
|
logger.warning("Failed to record the Telegram outcome prompt for reactions: %s", e)
|
|
await db.rollback()
|
|
return ok, status
|
|
|
|
async def _send_email(
|
|
self,
|
|
config: dict,
|
|
subject: str,
|
|
body: str,
|
|
image_data: bytes | None = None,
|
|
finish_photo_url: str | None = None,
|
|
) -> tuple[bool, str]:
|
|
"""Send notification via email (SMTP).
|
|
|
|
Inline finish-photo embed is opt-in via the template: when the rendered
|
|
``body`` contains the substituted ``{finish_photo_url}`` value AND the
|
|
finish-photo bytes are present, the message is built as
|
|
``multipart/related`` wrapping a ``multipart/alternative`` (plain + HTML)
|
|
plus an inline ``MIMEImage`` with ``Content-ID: <bambuddy-finish-photo>``.
|
|
The HTML part replaces the URL with ``<img src="cid:...">``; the plain-
|
|
text part keeps the URL as a clickable link. When the template doesn't
|
|
reference ``{finish_photo_url}`` (or image bytes aren't available), the
|
|
original single-part text shape is used — no attachment, no surprise
|
|
inline image (#1792).
|
|
"""
|
|
smtp_server = config.get("smtp_server", "").strip()
|
|
smtp_port = int(config.get("smtp_port", 587))
|
|
username = config.get("username", "").strip()
|
|
password = config.get("password", "").strip()
|
|
from_email = config.get("from_email", "").strip()
|
|
to_email = config.get("to_email", "").strip()
|
|
# Security: "starttls" (port 587), "ssl" (port 465), "none" (port 25)
|
|
security = config.get("security", "starttls")
|
|
# Authentication: "true" or "false"
|
|
auth_enabled = config.get("auth_enabled", "true").lower() == "true"
|
|
|
|
if not all([smtp_server, from_email, to_email]):
|
|
return False, "SMTP server, from email, and to email are required"
|
|
|
|
if auth_enabled and not all([username, password]):
|
|
return False, "Username and password are required when authentication is enabled"
|
|
|
|
# Template-driven: only inline-embed when the user's template explicitly
|
|
# referenced {finish_photo_url} (so the URL appears in the rendered body)
|
|
# AND the photo bytes are available. Falls back to text-only otherwise.
|
|
inline_photo = bool(image_data and finish_photo_url and finish_photo_url in body)
|
|
|
|
try:
|
|
if inline_photo:
|
|
# multipart/related → (multipart/alternative → text, html) + inline image
|
|
msg = MIMEMultipart("related")
|
|
msg["From"] = from_email
|
|
msg["To"] = to_email
|
|
msg["Subject"] = f"[Bambuddy] {subject}"
|
|
|
|
alt = MIMEMultipart("alternative")
|
|
alt.attach(MIMEText(body, "plain"))
|
|
# Build HTML body: escape the rendered body, then swap the
|
|
# escaped URL substring for an inline <img> referencing the
|
|
# MIMEImage we attach below. Done AFTER escape so the cid: URL
|
|
# we inject isn't re-escaped.
|
|
escaped_body = html.escape(body).replace("\n", "<br>\n")
|
|
escaped_url = html.escape(finish_photo_url)
|
|
img_tag = (
|
|
'<img src="cid:bambuddy-finish-photo" '
|
|
'alt="Printer camera snapshot" '
|
|
'style="max-width:100%;height:auto;border:1px solid #ddd;border-radius:4px;">'
|
|
)
|
|
html_body = f"<html><body><p>{escaped_body.replace(escaped_url, img_tag)}</p></body></html>"
|
|
alt.attach(MIMEText(html_body, "html"))
|
|
msg.attach(alt)
|
|
|
|
img = MIMEImage(image_data, _subtype="jpeg")
|
|
# Angle-bracketed Content-ID per RFC 2392, referenced from HTML
|
|
# without the brackets via ``cid:bambuddy-finish-photo``.
|
|
img.add_header("Content-ID", "<bambuddy-finish-photo>")
|
|
img.add_header("Content-Disposition", "inline", filename="finish-photo.jpg")
|
|
msg.attach(img)
|
|
else:
|
|
msg = MIMEMultipart()
|
|
msg["From"] = from_email
|
|
msg["To"] = to_email
|
|
msg["Subject"] = f"[Bambuddy] {subject}"
|
|
msg.attach(MIMEText(body, "plain"))
|
|
|
|
# smtplib is synchronous and blocking: a wedged / greylisting /
|
|
# firewall-dropped relay leaves recv() stuck. Two problems, two
|
|
# fixes (#2572):
|
|
# 1. No timeout — smtplib defaults to the global socket timeout,
|
|
# which this app never sets, so a stuck relay blocks forever.
|
|
# Pass an explicit timeout to every connect.
|
|
# 2. Run on the event loop — a stuck (or merely slow) send freezes
|
|
# every other coroutine, including a DB session a caller is
|
|
# holding open across this notification. Offload to a worker
|
|
# thread so the loop stays live and the connection is released
|
|
# on schedule.
|
|
smtp_timeout = 30.0
|
|
msg_str = msg.as_string()
|
|
|
|
def _blocking_send() -> None:
|
|
if security == "ssl":
|
|
# Direct SSL connection (typically port 465)
|
|
server = smtplib.SMTP_SSL(smtp_server, smtp_port, timeout=smtp_timeout)
|
|
elif security == "starttls":
|
|
# STARTTLS upgrade (typically port 587)
|
|
server = smtplib.SMTP(smtp_server, smtp_port, timeout=smtp_timeout)
|
|
server.starttls()
|
|
else:
|
|
# No encryption (typically port 25) - use with caution
|
|
server = smtplib.SMTP(smtp_server, smtp_port, timeout=smtp_timeout)
|
|
try:
|
|
if auth_enabled:
|
|
server.login(username, password)
|
|
server.sendmail(from_email, to_email, msg_str)
|
|
finally:
|
|
# quit() in finally so a send error doesn't leak the socket.
|
|
try:
|
|
server.quit()
|
|
except Exception: # noqa: BLE001 — closing a broken connection is best-effort
|
|
pass
|
|
|
|
await asyncio.to_thread(_blocking_send)
|
|
|
|
return True, "Email sent successfully"
|
|
except smtplib.SMTPAuthenticationError:
|
|
return False, "SMTP authentication failed - check username/password"
|
|
except smtplib.SMTPException as e:
|
|
return False, f"SMTP error: {str(e)}"
|
|
except Exception as e:
|
|
return False, f"Email error: {str(e)}"
|
|
|
|
async def _send_discord(
|
|
self, config: dict, title: str, message: str, image_data: bytes | None = None
|
|
) -> tuple[bool, str]:
|
|
"""Send notification via Discord webhook."""
|
|
webhook_url = config.get("webhook_url", "").strip()
|
|
|
|
if not webhook_url:
|
|
return False, "Webhook URL is required"
|
|
|
|
if not (
|
|
webhook_url.startswith("https://discord.com/api/webhooks/")
|
|
or webhook_url.startswith("https://discordapp.com/api/webhooks/")
|
|
):
|
|
return False, "Invalid Discord webhook URL"
|
|
|
|
# Discord embed format for nicer messages
|
|
embed = {
|
|
"title": title,
|
|
"description": message,
|
|
"color": 0x00AE42, # Bambu green
|
|
}
|
|
|
|
client = await self._get_client()
|
|
|
|
if image_data:
|
|
# Attach image via multipart form-data and reference in embed
|
|
embed["image"] = {"url": "attachment://photo.jpg"}
|
|
payload = {"embeds": [embed]}
|
|
response = await client.post(
|
|
webhook_url,
|
|
data={"payload_json": json.dumps(payload)},
|
|
files={"files[0]": ("photo.jpg", image_data, "image/jpeg")},
|
|
)
|
|
else:
|
|
response = await client.post(webhook_url, json={"embeds": [embed]})
|
|
|
|
if response.status_code in (200, 204):
|
|
return True, "Message sent successfully"
|
|
else:
|
|
return False, f"HTTP {response.status_code}: {response.text[:200]}"
|
|
|
|
async def _send_webhook(
|
|
self,
|
|
config: dict,
|
|
title: str,
|
|
message: str,
|
|
image_data: bytes | None = None,
|
|
event_type: str | None = None,
|
|
variables: dict | None = None,
|
|
) -> tuple[bool, str]:
|
|
"""Send notification via generic webhook (POST JSON).
|
|
|
|
Supports two payload formats:
|
|
- generic: Custom field names with timestamp/source metadata + structured event data
|
|
- slack: Slack/Mattermost compatible format (just {"text": "..."})
|
|
"""
|
|
webhook_url = config.get("webhook_url", "").strip()
|
|
auth_header = config.get("auth_header", "").strip()
|
|
payload_format = config.get("payload_format", "generic").strip()
|
|
|
|
if not webhook_url:
|
|
return False, "Webhook URL is required"
|
|
|
|
url_error = _assert_safe_provider_url(webhook_url, label="Webhook URL")
|
|
if url_error:
|
|
return False, url_error
|
|
|
|
# Build payload based on format
|
|
if payload_format == "slack":
|
|
# Slack/Mattermost format - just text field
|
|
data = {"text": f"*{title}*\n{message}"}
|
|
if event_type == "print_confirm_request":
|
|
# Slack and Mattermost fetch every URL in the text to build
|
|
# preview cards, and the outcome prompt (#1898) is the one
|
|
# message whose links are single-use capabilities — that fetch
|
|
# would be a machine answering the operator's question. Off
|
|
# here for the same reason Telegram's preview is.
|
|
#
|
|
# Only here: the slack payload never attaches image bytes (the
|
|
# base64 attach below is generic-format only), so unfurling is
|
|
# the only way a {finish_photo_url} in a print_complete body
|
|
# can render as a photo in the channel. Switching it off for
|
|
# every event would quietly take that away with no setting to
|
|
# get it back.
|
|
data["unfurl_links"] = False
|
|
data["unfurl_media"] = False
|
|
else:
|
|
# Generic format with custom field names
|
|
custom_field_title = config.get("field_title", "title").strip() or "title"
|
|
custom_field_message = config.get("field_message", "message").strip() or "message"
|
|
data = {
|
|
custom_field_title: title,
|
|
custom_field_message: message,
|
|
"timestamp": datetime.now().isoformat(),
|
|
"source": "Bambuddy",
|
|
}
|
|
|
|
# For generic format, include structured event data for automation tools
|
|
if payload_format != "slack":
|
|
if event_type:
|
|
data["event"] = event_type
|
|
if variables:
|
|
for key, value in variables.items():
|
|
if key not in data: # Don't overwrite title/message/timestamp/source
|
|
data[key] = value
|
|
|
|
# Attach base64-encoded image when available (generic format only)
|
|
if image_data and payload_format != "slack":
|
|
import base64
|
|
|
|
data["image"] = base64.b64encode(image_data).decode("ascii")
|
|
|
|
headers = {"Content-Type": "application/json"}
|
|
if auth_header:
|
|
# Support "Bearer token" or just "token" format
|
|
if " " in auth_header:
|
|
headers["Authorization"] = auth_header
|
|
else:
|
|
headers["Authorization"] = f"Bearer {auth_header}"
|
|
|
|
client = await self._get_client()
|
|
try:
|
|
response = await client.post(webhook_url, json=data, headers=headers)
|
|
|
|
if response.status_code in (200, 201, 202, 204):
|
|
return True, "Webhook delivered successfully"
|
|
else:
|
|
return False, _opaque_http_failure(response, label="webhook endpoint")
|
|
except Exception as e:
|
|
return False, f"Webhook error: {str(e)}"
|
|
|
|
async def _send_homeassistant(
|
|
self, config: dict, title: str, message: str, db: AsyncSession | None = None
|
|
) -> tuple[bool, str]:
|
|
"""Send notification via Home Assistant.
|
|
|
|
Uses the globally configured HA URL/token from settings.
|
|
Defaults to persistent_notification/create, but supports
|
|
custom services via config["service"] (e.g. notify.mobile_app_myphone).
|
|
"""
|
|
# Get HA connection settings from global config
|
|
ha_url = ""
|
|
ha_token = ""
|
|
|
|
if db:
|
|
from backend.app.api.routes.settings import get_homeassistant_settings
|
|
|
|
try:
|
|
ha_settings = await get_homeassistant_settings(db)
|
|
ha_url = ha_settings.get("ha_url", "")
|
|
ha_token = ha_settings.get("ha_token", "")
|
|
except Exception as e:
|
|
logger.warning("Failed to read HA settings from database: %s", e)
|
|
else:
|
|
# Fallback: read directly from environment if no DB session
|
|
import os
|
|
|
|
ha_url = os.environ.get("HA_URL", "")
|
|
ha_token = os.environ.get("HA_TOKEN", "")
|
|
|
|
if not ha_url or not ha_token:
|
|
return False, (
|
|
"Home Assistant is not configured. Please set HA URL and token in Settings → Network → Home Assistant."
|
|
)
|
|
|
|
# Determine which HA service to call - Default: persistent_notification.create
|
|
service = (config.get("service") or "").strip()
|
|
if service:
|
|
# Allow in different forms:
|
|
# - notify.mobile_app_<device>
|
|
# - notify/mobile_app_<device>
|
|
# - api/services/notify/mobile_app_<device>
|
|
service_str = service.lstrip("/")
|
|
if service_str.startswith("api/services/"):
|
|
endpoint = service_str
|
|
elif "/" in service_str:
|
|
endpoint = f"api/services/{service_str}"
|
|
elif "." in service_str:
|
|
domain, svc = service_str.split(".", 1)
|
|
endpoint = f"api/services/{domain}/{svc}"
|
|
else:
|
|
return False, (
|
|
"Invalid Home Assistant service name. Use e.g. 'notify.mobile_app_yourdevice' or 'notify/your_service'."
|
|
)
|
|
|
|
if not re.match(r"^api/services/[a-zA-Z0-9_]+/[a-zA-Z0-9_]+$", endpoint):
|
|
return False, (
|
|
"Invalid Home Assistant service name. Domain and service must only contain letters, numbers, and underscores."
|
|
)
|
|
else:
|
|
endpoint = "api/services/persistent_notification/create"
|
|
|
|
url = f"{ha_url.rstrip('/')}/{endpoint}"
|
|
headers = {
|
|
"Authorization": f"Bearer {ha_token}",
|
|
"Content-Type": "application/json",
|
|
}
|
|
payload = {
|
|
"title": title,
|
|
"message": message,
|
|
}
|
|
|
|
# Optional custom service-data (#1441), forwarded as HA's nested "data"
|
|
# object so mobile-app push options (priority, ttl, channel, group, ...)
|
|
# reach the notify service. Only included when configured — the default
|
|
# persistent_notification.create schema rejects unknown keys.
|
|
raw_data = config.get("data")
|
|
if raw_data:
|
|
if isinstance(raw_data, str):
|
|
try:
|
|
parsed_data = json.loads(raw_data)
|
|
except json.JSONDecodeError as e:
|
|
return False, f"Invalid JSON in the Data field: {e}"
|
|
else:
|
|
parsed_data = raw_data
|
|
if not isinstance(parsed_data, dict):
|
|
return False, 'The Data field must be a JSON object, e.g. {"priority": "high", "ttl": 0}'
|
|
if parsed_data:
|
|
payload["data"] = parsed_data
|
|
|
|
client = await self._get_client()
|
|
response = await client.post(url, json=payload, headers=headers)
|
|
|
|
if response.status_code in (200, 201):
|
|
return True, "Notification sent via Home Assistant"
|
|
elif response.status_code == 401:
|
|
return False, "Home Assistant authentication failed - check your token"
|
|
else:
|
|
# ha_url comes from global settings (SETTINGS_UPDATE, admin-only), so
|
|
# this is a narrower channel than the per-request provider URLs — but
|
|
# it lands in the same NOTIFICATIONS_CREATE-gated test response, so it
|
|
# gets the same treatment.
|
|
return False, _opaque_http_failure(response, label="Home Assistant endpoint")
|
|
|
|
async def _send_to_provider(
|
|
self,
|
|
provider: NotificationProvider,
|
|
title: str,
|
|
message: str,
|
|
db: AsyncSession | None = None,
|
|
image_data: bytes | None = None,
|
|
event_type: str | None = None,
|
|
variables: dict | None = None,
|
|
) -> tuple[bool, str]:
|
|
"""Send notification to a specific provider."""
|
|
# Check quiet hours
|
|
if self._is_in_quiet_hours(provider):
|
|
logger.info("Skipping notification to %s - quiet hours active", provider.name)
|
|
return True, "Skipped - quiet hours"
|
|
|
|
config = json.loads(provider.config) if isinstance(provider.config, str) else provider.config
|
|
|
|
try:
|
|
if provider.provider_type == "callmebot":
|
|
return await self._send_callmebot(config, f"{title}\n{message}")
|
|
elif provider.provider_type == "ntfy":
|
|
# Outcome confirmation (#1898): render the verdict capability
|
|
# links as one-tap buttons on the notification itself. http so
|
|
# no browser needs to open; POST because that is the method
|
|
# that records — a GET only opens the confirmation page, which
|
|
# is what keeps unfurlers from answering the prompt.
|
|
# clear=true dismisses the notification once a button was tapped.
|
|
ntfy_actions = None
|
|
good_url = (variables or {}).get("good_url")
|
|
reject_url = (variables or {}).get("reject_url")
|
|
# Buttons need absolute URLs; without a configured external_url
|
|
# the links are relative, and the body's deep link into the
|
|
# archive has to do.
|
|
if (
|
|
event_type == "print_confirm_request"
|
|
and good_url
|
|
and reject_url
|
|
and good_url.startswith("http")
|
|
and reject_url.startswith("http")
|
|
):
|
|
ntfy_actions = (
|
|
f"http, Good, {good_url}, method=POST, clear=true; "
|
|
f"http, Reject, {reject_url}, method=POST, clear=true"
|
|
)
|
|
return await self._send_ntfy(
|
|
config, title, message, image_data=image_data, event_type=event_type, actions=ntfy_actions
|
|
)
|
|
elif provider.provider_type == "pushover":
|
|
# Outcome confirmation (#1898): Pushover has no arbitrary
|
|
# buttons, but supports one supplementary URL — deep-link into
|
|
# the archive's confirmation dialog.
|
|
supplement_url = None
|
|
supplement_url_title = None
|
|
_confirm_url = (variables or {}).get("confirm_url")
|
|
if event_type == "print_confirm_request" and _confirm_url and _confirm_url.startswith("http"):
|
|
supplement_url = _confirm_url
|
|
supplement_url_title = "Confirm print outcome"
|
|
return await self._send_pushover(
|
|
config, title, message, image_data=image_data, url=supplement_url, url_title=supplement_url_title
|
|
)
|
|
elif provider.provider_type == "telegram":
|
|
if event_type == "print_confirm_request":
|
|
# Outcome confirmation (#1898): inline URL buttons under
|
|
# the message — one tap records the verdict via the
|
|
# capability link. Same absolute-URL requirement as the
|
|
# ntfy actions. Telegram has no way to POST, so this is the
|
|
# one affordance that opens a browser, and it is the one
|
|
# that gets the one-tap marker: the page submits itself
|
|
# only for a URL that came off a button. Telegram does not
|
|
# fetch inline-keyboard URLs and nothing else can read
|
|
# them, so the marker never reaches a scanner — which is
|
|
# the difference between this and trusting the User-Agent.
|
|
tg_buttons = None
|
|
_tg_good = (variables or {}).get("good_url")
|
|
_tg_reject = (variables or {}).get("reject_url")
|
|
if _tg_good and _tg_reject and _tg_good.startswith("http") and _tg_reject.startswith("http"):
|
|
tg_buttons = [
|
|
{"text": "\U0001f44d Good", "url": one_tap_url(_tg_good)},
|
|
{"text": "\U0001f44e Reject", "url": one_tap_url(_tg_reject)},
|
|
]
|
|
# Telegram verdict mode (#3046): buttons, a reaction on
|
|
# the message, or both. The link preview is off in all of
|
|
# them (see _send_telegram_confirm_request).
|
|
_archive_id = (variables or {}).get("archive_id")
|
|
return await self._send_telegram_confirm_request(
|
|
provider,
|
|
config,
|
|
f"*{title}*\n{message}",
|
|
db,
|
|
image_data,
|
|
tg_buttons,
|
|
_archive_id if isinstance(_archive_id, int) else None,
|
|
)
|
|
# Every other event keeps Telegram's link preview, for the same
|
|
# reason as the Slack unfurl in _send_webhook: when the finish
|
|
# photo is too large to attach, the preview is how a
|
|
# {finish_photo_url} in a print_complete body still shows up as
|
|
# a photo in the chat.
|
|
return await self._send_telegram(config, f"*{title}*\n{message}", image_data=image_data)
|
|
elif provider.provider_type == "email":
|
|
# finish_photo_url is pulled from the rendered template variables
|
|
# so _send_email can detect whether the template referenced the
|
|
# URL and inline-embed the photo only in that case.
|
|
finish_photo_url = (variables or {}).get("finish_photo_url")
|
|
return await self._send_email(
|
|
config, title, message, image_data=image_data, finish_photo_url=finish_photo_url
|
|
)
|
|
elif provider.provider_type == "discord":
|
|
return await self._send_discord(config, title, message, image_data=image_data)
|
|
elif provider.provider_type == "webhook":
|
|
return await self._send_webhook(
|
|
config, title, message, image_data=image_data, event_type=event_type, variables=variables
|
|
)
|
|
elif provider.provider_type == "homeassistant":
|
|
return await self._send_homeassistant(config, title, message, db=db)
|
|
elif provider.provider_type == "bark":
|
|
# Outcome confirmation (#1898): Bark opens one URL on tap —
|
|
# deep-link into the confirmation dialog, like Pushover.
|
|
bark_url = None
|
|
_bark_confirm = (variables or {}).get("confirm_url")
|
|
if event_type == "print_confirm_request" and _bark_confirm and _bark_confirm.startswith("http"):
|
|
bark_url = _bark_confirm
|
|
return await self._send_bark(config, title, message, url=bark_url)
|
|
else:
|
|
return False, f"Unknown provider type: {provider.provider_type}"
|
|
except Exception as e:
|
|
logger.exception("Error sending notification via %s", provider.provider_type)
|
|
return False, str(e)
|
|
|
|
async def _update_provider_status(
|
|
self, db: AsyncSession, provider_id: int, success: bool, error: str | None = None
|
|
):
|
|
"""Update provider status after sending notification."""
|
|
result = await db.execute(select(NotificationProvider).where(NotificationProvider.id == provider_id))
|
|
provider = result.scalar_one_or_none()
|
|
if provider:
|
|
if success:
|
|
provider.last_success = datetime.now(timezone.utc)
|
|
else:
|
|
provider.last_error = error
|
|
provider.last_error_at = datetime.now(timezone.utc)
|
|
await db.commit()
|
|
|
|
async def _get_providers_for_event(
|
|
self,
|
|
db: AsyncSession,
|
|
event_field: str,
|
|
printer_id: int | None = None,
|
|
) -> list[NotificationProvider]:
|
|
"""Get all enabled providers that want a specific event type.
|
|
|
|
Runs under ``no_autoflush`` (#2770). Callers routinely hold pending
|
|
writes when they raise an event — the AMS sensor loop does
|
|
``db.add(history)`` and only commits after the alarms have gone out — and
|
|
without this, autoflush satisfies this SELECT by writing those rows,
|
|
which opens a write transaction on SQLite. The provider is then contacted
|
|
over the network with that transaction still open, so a site whose
|
|
internet is down holds the single SQLite writer for the whole connect
|
|
timeout and unrelated background tasks fail with "database is locked".
|
|
Deferring the flush costs nothing here: providers are committed rows, so
|
|
a pending change in the caller's session cannot be one this query wants.
|
|
"""
|
|
# Build the query dynamically based on event field
|
|
query = select(NotificationProvider).where(
|
|
NotificationProvider.enabled.is_(True),
|
|
getattr(NotificationProvider, event_field).is_(True),
|
|
)
|
|
|
|
if printer_id is not None:
|
|
query = query.where(
|
|
(NotificationProvider.printer_id.is_(None)) | (NotificationProvider.printer_id == printer_id)
|
|
)
|
|
|
|
with db.no_autoflush:
|
|
result = await db.execute(query)
|
|
return list(result.scalars().all())
|
|
|
|
async def _log_notification(
|
|
self,
|
|
db: AsyncSession,
|
|
provider_id: int,
|
|
event_type: str,
|
|
title: str,
|
|
message: str,
|
|
success: bool,
|
|
error_message: str | None = None,
|
|
printer_id: int | None = None,
|
|
printer_name: str | None = None,
|
|
):
|
|
"""Create a log entry for a sent notification."""
|
|
try:
|
|
log = NotificationLog(
|
|
provider_id=provider_id,
|
|
event_type=event_type,
|
|
title=title,
|
|
message=message,
|
|
success=success,
|
|
error_message=error_message,
|
|
printer_id=printer_id,
|
|
printer_name=printer_name,
|
|
)
|
|
db.add(log)
|
|
await db.commit()
|
|
except Exception as e:
|
|
logger.warning("Failed to log notification: %s", e)
|
|
# Don't fail the notification just because logging failed
|
|
|
|
async def _send_to_providers(
|
|
self,
|
|
providers: list[NotificationProvider],
|
|
title: str,
|
|
message: str,
|
|
db: AsyncSession,
|
|
event_type: str = "unknown",
|
|
printer_id: int | None = None,
|
|
printer_name: str | None = None,
|
|
force_immediate: bool = False,
|
|
image_data: bytes | None = None,
|
|
variables: dict | None = None,
|
|
):
|
|
"""Send notification to multiple providers and log the results.
|
|
|
|
All notifications are always sent immediately. If digest mode is enabled,
|
|
the notification is ALSO queued for the daily digest summary.
|
|
"""
|
|
for provider in providers:
|
|
try:
|
|
# Always send notification immediately
|
|
success, error = await self._send_to_provider(
|
|
provider, title, message, db, image_data=image_data, event_type=event_type, variables=variables
|
|
)
|
|
|
|
# Also queue for digest if enabled (digest is a summary, not a queue)
|
|
if provider.daily_digest_enabled and provider.daily_digest_time:
|
|
await self._queue_for_digest(
|
|
provider=provider,
|
|
event_type=event_type,
|
|
title=title,
|
|
message=message,
|
|
db=db,
|
|
printer_id=printer_id,
|
|
printer_name=printer_name,
|
|
)
|
|
await self._update_provider_status(db, provider.id, success, error if not success else None)
|
|
await self._log_notification(
|
|
db=db,
|
|
provider_id=provider.id,
|
|
event_type=event_type,
|
|
title=title,
|
|
message=message,
|
|
success=success,
|
|
error_message=error if not success else None,
|
|
printer_id=printer_id,
|
|
printer_name=printer_name,
|
|
)
|
|
if success:
|
|
logger.info("Sent notification via %s", provider.name)
|
|
else:
|
|
logger.warning("Failed to send notification via %s: %s", provider.name, error)
|
|
except Exception as e:
|
|
logger.exception("Error sending notification via %s", provider.name)
|
|
await self._update_provider_status(db, provider.id, False, str(e))
|
|
await self._log_notification(
|
|
db=db,
|
|
provider_id=provider.id,
|
|
event_type=event_type,
|
|
title=title,
|
|
message=message,
|
|
success=False,
|
|
error_message=str(e),
|
|
printer_id=printer_id,
|
|
printer_name=printer_name,
|
|
)
|
|
|
|
async def on_app_message(
|
|
self,
|
|
db: AsyncSession,
|
|
*,
|
|
sender: str,
|
|
title: str,
|
|
message: str,
|
|
url: str | None = None,
|
|
) -> int:
|
|
"""A message another application sends through Bambuddy.
|
|
|
|
Goes to every enabled channel with "Messages from connected apps" on,
|
|
through the same path as Bambuddy's own events: quiet hours, the daily
|
|
digest and the log (whose event type names the sender). The text is
|
|
the app's own; a link, when given, is appended so every channel type
|
|
carries it. Returns how many channels it was handed to.
|
|
"""
|
|
providers = await self._get_providers_for_event(db, "on_app_message")
|
|
if not providers:
|
|
return 0
|
|
body = f"{message}\n{url}" if url else message
|
|
await self._send_to_providers(providers, title, body, db, event_type=f"app:{sender}"[:50])
|
|
return len(providers)
|
|
|
|
async def on_print_start(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
data: dict,
|
|
db: AsyncSession,
|
|
archive_data: dict | None = None,
|
|
):
|
|
"""Handle print start event - send notifications to relevant providers.
|
|
|
|
Args:
|
|
printer_id: The printer ID
|
|
printer_name: The printer name
|
|
data: MQTT event data with filename, subtask_name, remaining_time, raw_data
|
|
db: Database session
|
|
archive_data: Optional archive data with print_time_seconds from 3MF parsing
|
|
"""
|
|
logger.info("on_print_start called for printer %s (%s)", printer_id, printer_name)
|
|
providers = await self._get_providers_for_event(db, "on_print_start", printer_id)
|
|
if not providers:
|
|
logger.info("No notification providers configured for print_start event on printer %s", printer_id)
|
|
return
|
|
|
|
# Use subtask_name (project name) if available, otherwise use filename
|
|
subtask_name = data.get("subtask_name")
|
|
if subtask_name:
|
|
# Replace underscores with spaces for readability
|
|
filename = subtask_name.replace("_", " ")
|
|
else:
|
|
filename = self._clean_filename(data.get("filename", "Unknown"))
|
|
|
|
# Priority for estimated_time:
|
|
# 1. Archive's print_time_seconds from 3MF parsing (most reliable)
|
|
# 2. MQTT remaining_time (may be 0 at print start)
|
|
# 3. raw_data mc_remaining_time
|
|
estimated_time = None
|
|
|
|
# Try archive data first (from 3MF parsing - most reliable)
|
|
if archive_data and archive_data.get("print_time_seconds"):
|
|
estimated_time = archive_data["print_time_seconds"]
|
|
logger.debug("Using print_time_seconds from archive: %s", estimated_time)
|
|
|
|
# Fall back to MQTT remaining_time
|
|
if estimated_time is None:
|
|
estimated_time = data.get("remaining_time")
|
|
if estimated_time:
|
|
logger.debug("Using remaining_time from MQTT: %s", estimated_time)
|
|
|
|
# Last resort: raw_data mc_remaining_time (in minutes, convert to seconds)
|
|
if estimated_time is None:
|
|
raw_time = data.get("raw_data", {}).get("mc_remaining_time")
|
|
if raw_time:
|
|
estimated_time = raw_time * 60
|
|
logger.debug("Using mc_remaining_time from raw_data: %s", estimated_time)
|
|
|
|
time_str = self._format_duration(estimated_time)
|
|
eta_str = await self._format_eta(estimated_time, db)
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"filename": filename,
|
|
"estimated_time": time_str,
|
|
"eta": eta_str,
|
|
}
|
|
|
|
# Extract image data for providers that support attachments (e.g. Pushover)
|
|
image_data = None
|
|
if archive_data:
|
|
image_data = archive_data.get("image_data")
|
|
|
|
logger.info("Found %s providers for print_start: %s", len(providers), [p.name for p in providers])
|
|
title, message = await self._build_message_from_template(db, "print_start", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"print_start",
|
|
printer_id,
|
|
printer_name,
|
|
image_data=image_data,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_print_complete(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
status: str,
|
|
data: dict,
|
|
db: AsyncSession,
|
|
archive_data: dict | None = None,
|
|
):
|
|
"""Handle print complete event - send notifications to relevant providers."""
|
|
logger.info("on_print_complete called for printer %s (%s), status=%s", printer_id, printer_name, status)
|
|
|
|
# Determine event type based on status
|
|
if status == "completed":
|
|
event_field = "on_print_complete"
|
|
event_type = "print_complete"
|
|
elif status in ("failed",):
|
|
event_field = "on_print_failed"
|
|
event_type = "print_failed"
|
|
elif status in ("aborted", "stopped", "cancelled"):
|
|
event_field = "on_print_stopped"
|
|
event_type = "print_stopped"
|
|
else:
|
|
logger.warning("Unknown print status '%s', defaulting to on_print_complete", status)
|
|
event_field = "on_print_complete"
|
|
event_type = "print_complete"
|
|
|
|
providers = await self._get_providers_for_event(db, event_field, printer_id)
|
|
if not providers:
|
|
logger.info("No notification providers configured for %s event on printer %s", event_field, printer_id)
|
|
return
|
|
|
|
# Use subtask_name (project name) if available, otherwise use filename
|
|
subtask_name = data.get("subtask_name")
|
|
if subtask_name:
|
|
filename = subtask_name.replace("_", " ")
|
|
else:
|
|
filename = self._clean_filename(data.get("filename", "Unknown"))
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"filename": filename,
|
|
"duration": "Unknown",
|
|
"filament_grams": "Unknown",
|
|
"reason": "Unknown",
|
|
}
|
|
|
|
if archive_data:
|
|
# {{duration}} on completion / failure / stopped events is the *actual*
|
|
# elapsed time (#1198). Slicer-estimated print_time_seconds is only used
|
|
# as a last-resort fallback when timestamps weren't recorded.
|
|
duration_seconds = archive_data.get("actual_time_seconds") or archive_data.get("print_time_seconds")
|
|
if duration_seconds:
|
|
variables["duration"] = self._format_duration(duration_seconds)
|
|
if archive_data.get("actual_filament_grams"):
|
|
variables["filament_grams"] = f"{archive_data['actual_filament_grams']:.1f}"
|
|
if status == "failed" and archive_data.get("failure_reason"):
|
|
variables["reason"] = archive_data["failure_reason"]
|
|
if archive_data.get("finish_photo_url"):
|
|
variables["finish_photo_url"] = archive_data["finish_photo_url"]
|
|
|
|
# Build per-slot breakdown string with AMS info when available
|
|
if archive_data.get("usage_results"):
|
|
parts = []
|
|
for u in archive_data["usage_results"]:
|
|
ams_id = u.get("ams_id", 0)
|
|
tray_id = u.get("tray_id", 0)
|
|
material = u.get("material", "Unknown") or "Unknown"
|
|
used = u.get("weight_used", 0)
|
|
if ams_id >= 128:
|
|
slot_label = "Ext"
|
|
else:
|
|
slot_label = f"AMS-{chr(65 + ams_id)} T{tray_id + 1}"
|
|
parts.append(f"{slot_label} {material}: {used:.1f}g")
|
|
variables["filament_details"] = " | ".join(parts)
|
|
elif archive_data.get("filament_slots"):
|
|
parts = []
|
|
for slot in archive_data["filament_slots"]:
|
|
ftype = slot.get("type", "Unknown") or "Unknown"
|
|
used = slot.get("used_g", 0)
|
|
parts.append(f"{ftype}: {used:.1f}g")
|
|
variables["filament_details"] = " | ".join(parts)
|
|
|
|
# Add progress for partial prints
|
|
if archive_data.get("progress") is not None:
|
|
variables["progress"] = str(archive_data["progress"])
|
|
|
|
# Extract image data for providers that support attachments (e.g. Pushover)
|
|
image_data = None
|
|
if archive_data:
|
|
image_data = archive_data.get("image_data")
|
|
|
|
logger.info("Found %s providers for %s: %s", len(providers), event_field, [p.name for p in providers])
|
|
title, message = await self._build_message_from_template(db, event_type, variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
event_type,
|
|
printer_id,
|
|
printer_name,
|
|
image_data=image_data,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_print_confirm_request(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
data: dict,
|
|
db: AsyncSession,
|
|
archive_data: dict | None = None,
|
|
good_url: str | None = None,
|
|
reject_url: str | None = None,
|
|
confirm_url: str | None = None,
|
|
archive_id: int | None = None,
|
|
):
|
|
"""Ask for a post-print outcome verdict (#1898).
|
|
|
|
Fires only for completed prints whose queue item opted in — the
|
|
provider-level toggle exists to mute a channel, not to enable the
|
|
feature. good_url / reject_url are the one-tap capability links
|
|
(rendered as ntfy action buttons), confirm_url deep-links into the
|
|
archive's confirmation dialog in the web UI. archive_id lets a Telegram
|
|
provider in reaction mode (#3046) tie the sent message to the archive.
|
|
"""
|
|
providers = await self._get_providers_for_event(db, "on_print_confirm_request", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
subtask_name = data.get("subtask_name")
|
|
if subtask_name:
|
|
filename = subtask_name.replace("_", " ")
|
|
else:
|
|
filename = self._clean_filename(data.get("filename", "Unknown"))
|
|
|
|
variables = {"printer": printer_name, "filename": filename}
|
|
if good_url:
|
|
variables["good_url"] = good_url
|
|
if reject_url:
|
|
variables["reject_url"] = reject_url
|
|
if confirm_url:
|
|
variables["confirm_url"] = confirm_url
|
|
if archive_id is not None:
|
|
variables["archive_id"] = archive_id
|
|
|
|
image_data = None
|
|
if archive_data:
|
|
if archive_data.get("finish_photo_url"):
|
|
variables["finish_photo_url"] = archive_data["finish_photo_url"]
|
|
image_data = archive_data.get("image_data")
|
|
|
|
title, message = await self._build_message_from_template(db, "print_confirm_request", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"print_confirm_request",
|
|
printer_id,
|
|
printer_name,
|
|
image_data=image_data,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_print_progress(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
filename: str,
|
|
progress: int,
|
|
db: AsyncSession,
|
|
remaining_time: int | None = None,
|
|
image_data: bytes | None = None,
|
|
):
|
|
"""Handle print progress milestone (25%, 50%, 75%)."""
|
|
providers = await self._get_providers_for_event(db, "on_print_progress", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
eta_str = await self._format_eta(remaining_time, db)
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"filename": self._clean_filename(filename),
|
|
"progress": str(progress),
|
|
"remaining_time": self._format_duration(remaining_time) if remaining_time else "Unknown",
|
|
"eta": eta_str,
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "print_progress", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"print_progress",
|
|
printer_id,
|
|
printer_name,
|
|
image_data=image_data,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_print_missing_spool_assignment(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
missing_slots: list[dict[str, str]],
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle print-start event when required trays are missing spool assignments."""
|
|
if not missing_slots:
|
|
return
|
|
|
|
providers = await self._get_providers_for_event(db, "on_print_missing_spool_assignment", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
missing_slot_names = ", ".join(slot.get("slot", "Unknown") for slot in missing_slots)
|
|
detail_lines = []
|
|
for slot in missing_slots:
|
|
slot_name = slot.get("slot", "Unknown")
|
|
profile = slot.get("profile", "Unknown")
|
|
detail_lines.append(f"- {slot_name}: {profile}")
|
|
missing_profile_details = "\n".join(detail_lines)
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"missing_slots": missing_slot_names,
|
|
"missing_slot_details": missing_profile_details,
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "print_missing_spool_assignment", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"print_missing_spool_assignment",
|
|
printer_id,
|
|
printer_name,
|
|
force_immediate=True,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_billing_charge_failed(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
filename: str,
|
|
archive_id: int | None,
|
|
error: str,
|
|
db: AsyncSession,
|
|
) -> None:
|
|
"""Notify providers that a terminal print could not be charged."""
|
|
providers = await self._get_providers_for_event(db, "on_billing_charge_failed", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"filename": self._clean_filename(filename),
|
|
"archive_id": str(archive_id) if archive_id is not None else "Unknown",
|
|
"error": error,
|
|
}
|
|
title, message = await self._build_message_from_template(db, "billing_charge_failed", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"billing_charge_failed",
|
|
printer_id,
|
|
printer_name,
|
|
force_immediate=True,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_printer_offline(self, printer_id: int, printer_name: str, db: AsyncSession):
|
|
"""Handle printer offline event."""
|
|
providers = await self._get_providers_for_event(db, "on_printer_offline", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {"printer": printer_name}
|
|
|
|
title, message = await self._build_message_from_template(db, "printer_offline", variables)
|
|
await self._send_to_providers(
|
|
providers, title, message, db, "printer_offline", printer_id, printer_name, variables=variables
|
|
)
|
|
|
|
async def on_printer_error(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
error_type: str,
|
|
db: AsyncSession,
|
|
error_detail: str | None = None,
|
|
image_data: bytes | None = None,
|
|
):
|
|
"""Handle printer error event (AMS issues, etc.)."""
|
|
providers = await self._get_providers_for_event(db, "on_printer_error", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"error_type": error_type,
|
|
"error_detail": error_detail or "No details available",
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "printer_error", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"printer_error",
|
|
printer_id,
|
|
printer_name,
|
|
image_data=image_data,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_ai_failure_detection(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
task_name: str,
|
|
confidence: float,
|
|
action: str,
|
|
db: AsyncSession,
|
|
image_data: bytes | None = None,
|
|
):
|
|
"""Handle AI failure-detection event (Obico spaghetti / print-failure ML).
|
|
|
|
Split out of on_printer_error (#1794) so a user can subscribe to AI
|
|
alerts without also being paged for every HMS hardware code.
|
|
"""
|
|
providers = await self._get_providers_for_event(db, "on_ai_failure_detection", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"task_name": task_name or "current job",
|
|
"confidence": f"{confidence:.2f}",
|
|
"action": action,
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "ai_failure_detection", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"ai_failure_detection",
|
|
printer_id,
|
|
printer_name,
|
|
image_data=image_data,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_plate_not_empty(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
db: AsyncSession,
|
|
difference_percent: float | None = None,
|
|
):
|
|
"""Handle plate not empty event - objects detected on build plate before print."""
|
|
providers = await self._get_providers_for_event(db, "on_plate_not_empty", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"difference_percent": f"{difference_percent:.1f}" if difference_percent else "N/A",
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "plate_not_empty", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"plate_not_empty",
|
|
printer_id,
|
|
printer_name,
|
|
force_immediate=True,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_plate_clear_required(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle plate-clear-required event — a print ended and the queue is gated (#2525).
|
|
|
|
Distinct from ``on_plate_not_empty``, which is the camera check *before* a
|
|
print starts. This one fires on the rising edge of the Bambuddy-side
|
|
awaiting-plate-clear flag, i.e. whenever a print reaches a terminal state
|
|
and the next queued job can't dispatch until someone confirms the bed is
|
|
free. Off by default on every provider: it lands at the same moment as the
|
|
print-complete notification, so opting in is a deliberate choice.
|
|
"""
|
|
providers = await self._get_providers_for_event(db, "on_plate_clear_required", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {"printer": printer_name}
|
|
|
|
title, message = await self._build_message_from_template(db, "plate_clear_required", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"plate_clear_required",
|
|
printer_id,
|
|
printer_name,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_filament_low(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
slot: str,
|
|
remaining_percent: int,
|
|
db: AsyncSession,
|
|
color: str | None = None,
|
|
):
|
|
"""Handle low filament event."""
|
|
providers = await self._get_providers_for_event(db, "on_filament_low", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"slot": slot,
|
|
"remaining_percent": str(remaining_percent),
|
|
"color": color or "",
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "filament_low", variables)
|
|
await self._send_to_providers(
|
|
providers, title, message, db, "filament_low", printer_id, printer_name, variables=variables
|
|
)
|
|
|
|
async def on_maintenance_due(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
maintenance_items: list[dict],
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle maintenance due event - sends notification when maintenance is due or warning."""
|
|
if not maintenance_items:
|
|
return
|
|
|
|
providers = await self._get_providers_for_event(db, "on_maintenance_due", printer_id)
|
|
if not providers:
|
|
logger.info("No notification providers configured for maintenance_due event on printer %s", printer_id)
|
|
return
|
|
|
|
# Format maintenance items list
|
|
items_list = []
|
|
for item in maintenance_items:
|
|
status = "OVERDUE" if item.get("is_due") else "Soon"
|
|
items_list.append(f"- {item['name']} ({status})")
|
|
items_str = "\n".join(items_list)
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"items": items_str,
|
|
}
|
|
|
|
logger.info("Found %s providers for maintenance_due: %s", len(providers), [p.name for p in providers])
|
|
title, message = await self._build_message_from_template(db, "maintenance_due", variables)
|
|
await self._send_to_providers(
|
|
providers, title, message, db, "maintenance_due", printer_id, printer_name, variables=variables
|
|
)
|
|
|
|
async def on_ams_humidity_high(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
ams_label: str,
|
|
humidity: float,
|
|
threshold: float,
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle AMS high humidity alarm event. Always sends immediately (bypasses digest)."""
|
|
providers = await self._get_providers_for_event(db, "on_ams_humidity_high", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"ams_label": ams_label,
|
|
"humidity": f"{humidity:.0f}",
|
|
"threshold": f"{threshold:.0f}",
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "ams_humidity_high", variables)
|
|
# Alarms always send immediately, bypassing digest mode
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"ams_humidity_high",
|
|
printer_id,
|
|
printer_name,
|
|
force_immediate=True,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_ams_temperature_high(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
ams_label: str,
|
|
temperature: float,
|
|
threshold: float,
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle AMS high temperature alarm event. Always sends immediately (bypasses digest)."""
|
|
providers = await self._get_providers_for_event(db, "on_ams_temperature_high", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"ams_label": ams_label,
|
|
"temperature": f"{temperature:.1f}",
|
|
"threshold": f"{threshold:.1f}",
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "ams_temperature_high", variables)
|
|
# Alarms always send immediately, bypassing digest mode
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"ams_temperature_high",
|
|
printer_id,
|
|
printer_name,
|
|
force_immediate=True,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_ams_drying_suspended(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
ams_label: str,
|
|
humidity: float,
|
|
threshold: float,
|
|
cycles: int,
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle automatic drying giving up on one AMS unit (#2770).
|
|
|
|
Sent immediately rather than folded into a digest: it reports that
|
|
Bambuddy has STOPPED doing something, and a report of inaction that
|
|
arrives with tomorrow's summary has already cost the user a day.
|
|
"""
|
|
providers = await self._get_providers_for_event(db, "on_ams_drying_suspended", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"ams_label": ams_label,
|
|
"humidity": f"{humidity:.0f}",
|
|
"threshold": f"{threshold:.0f}",
|
|
"cycles": str(cycles),
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "ams_drying_suspended", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"ams_drying_suspended",
|
|
printer_id,
|
|
printer_name,
|
|
force_immediate=True,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_ams_ht_humidity_high(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
ams_label: str,
|
|
humidity: float,
|
|
threshold: float,
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle AMS-HT high humidity alarm event. Always sends immediately (bypasses digest)."""
|
|
providers = await self._get_providers_for_event(db, "on_ams_ht_humidity_high", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"ams_label": ams_label,
|
|
"humidity": f"{humidity:.0f}",
|
|
"threshold": f"{threshold:.0f}",
|
|
}
|
|
|
|
# Use the same template as regular AMS (can create separate templates later if needed)
|
|
title, message = await self._build_message_from_template(db, "ams_humidity_high", variables)
|
|
# Alarms always send immediately, bypassing digest mode
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"ams_ht_humidity_high",
|
|
printer_id,
|
|
printer_name,
|
|
force_immediate=True,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_ams_ht_temperature_high(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
ams_label: str,
|
|
temperature: float,
|
|
threshold: float,
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle AMS-HT high temperature alarm event. Always sends immediately (bypasses digest)."""
|
|
providers = await self._get_providers_for_event(db, "on_ams_ht_temperature_high", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"ams_label": ams_label,
|
|
"temperature": f"{temperature:.1f}",
|
|
"threshold": f"{threshold:.1f}",
|
|
}
|
|
|
|
# Use the same template as regular AMS (can create separate templates later if needed)
|
|
title, message = await self._build_message_from_template(db, "ams_temperature_high", variables)
|
|
# Alarms always send immediately, bypassing digest mode
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"ams_ht_temperature_high",
|
|
printer_id,
|
|
printer_name,
|
|
force_immediate=True,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_bed_cooled(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
bed_temp: float,
|
|
threshold: float,
|
|
filename: str,
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle bed cooled event - bed temperature dropped below threshold after print."""
|
|
providers = await self._get_providers_for_event(db, "on_bed_cooled", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"bed_temp": f"{bed_temp:.0f}",
|
|
"threshold": f"{threshold:.0f}",
|
|
"filename": self._clean_filename(filename) if filename else "Unknown",
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "bed_cooled", variables)
|
|
await self._send_to_providers(
|
|
providers, title, message, db, "bed_cooled", printer_id, printer_name, variables=variables
|
|
)
|
|
|
|
async def on_ha_sensor_alert(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
sensor_name: str,
|
|
state: str,
|
|
db: AsyncSession,
|
|
):
|
|
"""A Home Assistant sensor bound to a printer entered its alert state (#1148).
|
|
|
|
Sent immediately rather than folded into a digest: the case this exists
|
|
for is an enclosure door left open, which is only worth telling someone
|
|
about while they can still act on it.
|
|
"""
|
|
providers = await self._get_providers_for_event(db, "on_ha_sensor_alert", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"sensor": sensor_name,
|
|
"state": state,
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "ha_sensor_alert", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"ha_sensor_alert",
|
|
printer_id,
|
|
printer_name,
|
|
force_immediate=True,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_location_ha_sensor_alert(
|
|
self,
|
|
location_name: str,
|
|
sensor_name: str,
|
|
state: str,
|
|
db: AsyncSession,
|
|
):
|
|
"""A Home Assistant sensor bound to a storage location entered its alert state (#2824).
|
|
|
|
Sent immediately rather than folded into a digest, for the same reason
|
|
as on_ha_sensor_alert above: this is the "drybox went stale" case, only
|
|
worth acting on while the humidity/temperature is still climbing.
|
|
"""
|
|
# Own column, not on_ha_sensor_alert (#2824): that one can be scoped to
|
|
# a single printer via provider.printer_id, and a location alert has no
|
|
# printer to scope by, so sharing it would leak drybox alerts to a
|
|
# provider narrowed to one printer's sensors.
|
|
providers = await self._get_providers_for_event(db, "on_location_ha_sensor_alert", None)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"location": location_name,
|
|
"sensor": sensor_name,
|
|
"state": state,
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "location_ha_sensor_alert", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"location_ha_sensor_alert",
|
|
force_immediate=True,
|
|
variables=variables,
|
|
)
|
|
|
|
async def on_first_layer_complete(
|
|
self,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
filename: str,
|
|
total_layers: int,
|
|
db: AsyncSession,
|
|
image_data: bytes | None = None,
|
|
):
|
|
"""Handle first layer complete event."""
|
|
providers = await self._get_providers_for_event(db, "on_first_layer_complete", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"printer": printer_name,
|
|
"filename": self._clean_filename(filename),
|
|
"total_layers": str(total_layers),
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "first_layer_complete", variables)
|
|
await self._send_to_providers(
|
|
providers,
|
|
title,
|
|
message,
|
|
db,
|
|
"first_layer_complete",
|
|
printer_id,
|
|
printer_name,
|
|
image_data=image_data,
|
|
variables=variables,
|
|
)
|
|
|
|
def clear_template_cache(self):
|
|
"""Clear the template cache. Call this when templates are updated."""
|
|
self._template_cache.clear()
|
|
|
|
async def send_user_print_email(
|
|
self,
|
|
event_type: str,
|
|
created_by_id: int | None,
|
|
printer_name: str,
|
|
filename: str,
|
|
db: AsyncSession,
|
|
) -> None:
|
|
"""Send a print event email notification to the user who submitted the job.
|
|
|
|
Args:
|
|
event_type: 'user_print_start', 'user_print_complete', 'user_print_failed', or 'user_print_stopped'
|
|
created_by_id: User ID who submitted the print job (from archive)
|
|
printer_name: Name of the printer
|
|
filename: Raw filename or subtask name
|
|
db: Database session
|
|
"""
|
|
if created_by_id is None:
|
|
logger.debug("[EMAIL] Skipping user print email (%s): no created_by_id", event_type)
|
|
return
|
|
|
|
try:
|
|
# Check if advanced auth is enabled - required for user email notifications
|
|
from backend.app.models.settings import Settings
|
|
|
|
result = await db.execute(select(Settings).where(Settings.key == "advanced_auth_enabled"))
|
|
setting = result.scalar_one_or_none()
|
|
if not setting or setting.value.lower() != "true":
|
|
logger.debug("[EMAIL] Skipping user print email (%s): advanced_auth not enabled", event_type)
|
|
return
|
|
|
|
# Check if user notifications are enabled (admin-controlled toggle)
|
|
notif_enabled_result = await db.execute(
|
|
select(Settings).where(Settings.key == "user_notifications_enabled")
|
|
)
|
|
notif_enabled_setting = notif_enabled_result.scalar_one_or_none()
|
|
if notif_enabled_setting and notif_enabled_setting.value.lower() == "false":
|
|
logger.debug("[EMAIL] Skipping user print email (%s): user_notifications_enabled is false", event_type)
|
|
return
|
|
|
|
# Check SMTP settings are configured - required for sending emails
|
|
from backend.app.services.email_service import get_smtp_settings, send_user_print_notification
|
|
|
|
smtp_settings = await get_smtp_settings(db)
|
|
if not smtp_settings:
|
|
logger.debug("[EMAIL] Skipping user print email (%s): SMTP settings not configured", event_type)
|
|
return
|
|
|
|
# Load user preferences
|
|
from backend.app.models.user import User
|
|
from backend.app.models.user_email_pref import UserEmailPreference
|
|
|
|
user_result = await db.execute(select(User).where(User.id == created_by_id))
|
|
user = user_result.scalar_one_or_none()
|
|
if user is None or not user.email:
|
|
logger.debug(
|
|
"[EMAIL] Skipping user print email (%s): user %s not found or has no email address",
|
|
event_type,
|
|
created_by_id,
|
|
)
|
|
return
|
|
|
|
# Load user's notification preferences
|
|
pref_result = await db.execute(
|
|
select(UserEmailPreference).where(UserEmailPreference.user_id == created_by_id)
|
|
)
|
|
pref = pref_result.scalar_one_or_none()
|
|
|
|
# Determine if this event type should be sent
|
|
should_send = False
|
|
if event_type == "user_print_start":
|
|
should_send = pref is None or pref.notify_print_start
|
|
elif event_type == "user_print_complete":
|
|
should_send = pref is None or pref.notify_print_complete
|
|
elif event_type == "user_print_failed":
|
|
should_send = pref is None or pref.notify_print_failed
|
|
elif event_type == "user_print_stopped":
|
|
should_send = pref is None or pref.notify_print_stopped
|
|
|
|
if not should_send:
|
|
logger.debug(
|
|
"[EMAIL] Skipping user print email (%s): user %s has notifications disabled for this event",
|
|
event_type,
|
|
created_by_id,
|
|
)
|
|
return
|
|
|
|
logger.info(
|
|
"[EMAIL] Sending user print email: event=%s, user=%s (%s), printer=%s, file=%s",
|
|
event_type,
|
|
user.username,
|
|
user.email,
|
|
printer_name,
|
|
filename,
|
|
)
|
|
|
|
# Build variables
|
|
variables = {
|
|
"printer": printer_name,
|
|
"filename": self._clean_filename(filename),
|
|
}
|
|
|
|
# Send the email
|
|
await send_user_print_notification(
|
|
db=db,
|
|
event_type=event_type,
|
|
user_email=user.email,
|
|
username=user.username,
|
|
variables=variables,
|
|
)
|
|
logger.info("[EMAIL] User print email sent: event=%s → %s", event_type, user.email)
|
|
except Exception as e:
|
|
logger.warning("Failed to send user print email notification: %s", e, exc_info=True)
|
|
|
|
# ==================== Queue Notifications ====================
|
|
|
|
async def on_queue_job_added(
|
|
self,
|
|
job_name: str,
|
|
target: str,
|
|
db: AsyncSession,
|
|
printer_id: int | None = None,
|
|
printer_name: str | None = None,
|
|
):
|
|
"""Handle queue job added event."""
|
|
providers = await self._get_providers_for_event(db, "on_queue_job_added", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"job_name": job_name,
|
|
"target": target, # e.g., "Printer1" or "Any X1C"
|
|
"printer": printer_name or target,
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "queue_job_added", variables)
|
|
await self._send_to_providers(
|
|
providers, title, message, db, "queue_job_added", printer_id, printer_name, variables=variables
|
|
)
|
|
|
|
async def on_queue_job_assigned(
|
|
self,
|
|
job_name: str,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
target_model: str,
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle model-based job assigned to printer event."""
|
|
providers = await self._get_providers_for_event(db, "on_queue_job_assigned", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"job_name": job_name,
|
|
"printer": printer_name,
|
|
"target_model": target_model,
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "queue_job_assigned", variables)
|
|
await self._send_to_providers(
|
|
providers, title, message, db, "queue_job_assigned", printer_id, printer_name, variables=variables
|
|
)
|
|
|
|
async def on_queue_job_started(
|
|
self,
|
|
job_name: str,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
db: AsyncSession,
|
|
estimated_time: int | None = None,
|
|
):
|
|
"""Handle queue job started printing event."""
|
|
providers = await self._get_providers_for_event(db, "on_queue_job_started", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
eta_str = await self._format_eta(estimated_time, db)
|
|
|
|
variables = {
|
|
"job_name": job_name,
|
|
"printer": printer_name,
|
|
"estimated_time": self._format_duration(estimated_time),
|
|
"eta": eta_str,
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "queue_job_started", variables)
|
|
await self._send_to_providers(
|
|
providers, title, message, db, "queue_job_started", printer_id, printer_name, variables=variables
|
|
)
|
|
|
|
async def on_queue_job_waiting(
|
|
self,
|
|
job_name: str,
|
|
target_model: str,
|
|
waiting_reason: str,
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle job waiting for filament event."""
|
|
providers = await self._get_providers_for_event(db, "on_queue_job_waiting", None)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"job_name": job_name,
|
|
"target_model": target_model,
|
|
"waiting_reason": waiting_reason,
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "queue_job_waiting", variables)
|
|
await self._send_to_providers(providers, title, message, db, "queue_job_waiting", variables=variables)
|
|
|
|
async def on_queue_job_skipped(
|
|
self,
|
|
job_name: str,
|
|
printer_id: int,
|
|
printer_name: str,
|
|
reason: str,
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle job skipped event (e.g., previous print failed)."""
|
|
providers = await self._get_providers_for_event(db, "on_queue_job_skipped", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"job_name": job_name,
|
|
"printer": printer_name,
|
|
"reason": reason,
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "queue_job_skipped", variables)
|
|
await self._send_to_providers(
|
|
providers, title, message, db, "queue_job_skipped", printer_id, printer_name, variables=variables
|
|
)
|
|
|
|
async def on_queue_job_failed(
|
|
self,
|
|
job_name: str,
|
|
printer_id: int | None,
|
|
printer_name: str | None,
|
|
reason: str,
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle job failed to start event (upload error, etc.)."""
|
|
providers = await self._get_providers_for_event(db, "on_queue_job_failed", printer_id)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"job_name": job_name,
|
|
"printer": printer_name or "Unknown",
|
|
"reason": reason,
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "queue_job_failed", variables)
|
|
await self._send_to_providers(
|
|
providers, title, message, db, "queue_job_failed", printer_id, printer_name, variables=variables
|
|
)
|
|
|
|
async def on_queue_completed(
|
|
self,
|
|
completed_count: int,
|
|
db: AsyncSession,
|
|
):
|
|
"""Handle all queue jobs completed event."""
|
|
providers = await self._get_providers_for_event(db, "on_queue_completed", None)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"completed_count": str(completed_count),
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "queue_completed", variables)
|
|
await self._send_to_providers(providers, title, message, db, "queue_completed", variables=variables)
|
|
|
|
# ==================== Inventory Stock Alerts ====================
|
|
|
|
async def on_stock_reorder_alert(
|
|
self,
|
|
material: str,
|
|
brand: str | None,
|
|
stock_g: float,
|
|
rate_g_day: float,
|
|
days_left: int,
|
|
db: AsyncSession,
|
|
):
|
|
"""Fire when an inventory SKU reaches its reorder point."""
|
|
providers = await self._get_providers_for_event(db, "on_stock_reorder_alert", None)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"material": material,
|
|
"brand": brand or "",
|
|
"stock_g": f"{stock_g:.0f}",
|
|
"rate_g_day": f"{rate_g_day:.1f}",
|
|
"days_left": str(days_left),
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "stock_reorder_alert", variables)
|
|
await self._send_to_providers(providers, title, message, db, "stock_reorder_alert", variables=variables)
|
|
|
|
async def on_stock_break_alert(
|
|
self,
|
|
material: str,
|
|
brand: str | None,
|
|
stock_g: float,
|
|
rate_g_day: float,
|
|
days_left: int,
|
|
lead_time_days: int,
|
|
db: AsyncSession,
|
|
):
|
|
"""Fire when a stock break is detected (stock runs out before lead time)."""
|
|
providers = await self._get_providers_for_event(db, "on_stock_break_alert", None)
|
|
if not providers:
|
|
return
|
|
|
|
variables = {
|
|
"material": material,
|
|
"brand": brand or "",
|
|
"stock_g": f"{stock_g:.0f}",
|
|
"rate_g_day": f"{rate_g_day:.1f}",
|
|
"days_left": str(days_left),
|
|
"lead_time_days": str(lead_time_days),
|
|
}
|
|
|
|
title, message = await self._build_message_from_template(db, "stock_break_alert", variables)
|
|
await self._send_to_providers(providers, title, message, db, "stock_break_alert", variables=variables)
|
|
|
|
async def _queue_for_digest(
|
|
self,
|
|
provider: NotificationProvider,
|
|
event_type: str,
|
|
title: str,
|
|
message: str,
|
|
db: AsyncSession,
|
|
printer_id: int | None = None,
|
|
printer_name: str | None = None,
|
|
):
|
|
"""Queue a notification for later delivery in the daily digest."""
|
|
try:
|
|
queue_entry = NotificationDigestQueue(
|
|
provider_id=provider.id,
|
|
event_type=event_type,
|
|
title=title,
|
|
message=message,
|
|
printer_id=printer_id,
|
|
printer_name=printer_name,
|
|
)
|
|
db.add(queue_entry)
|
|
await db.commit()
|
|
logger.info("Queued notification for digest: %s for provider %s", event_type, provider.name)
|
|
except Exception as e:
|
|
logger.warning("Failed to queue notification for digest: %s", e)
|
|
|
|
async def send_digest(self, provider_id: int):
|
|
"""Send all queued notifications as a single digest for a provider."""
|
|
from backend.app.core.database import async_session
|
|
|
|
async with async_session() as db:
|
|
# Get the provider
|
|
result = await db.execute(select(NotificationProvider).where(NotificationProvider.id == provider_id))
|
|
provider = result.scalar_one_or_none()
|
|
|
|
if not provider or not provider.enabled:
|
|
return
|
|
|
|
# Get all queued notifications for this provider
|
|
result = await db.execute(
|
|
select(NotificationDigestQueue)
|
|
.where(NotificationDigestQueue.provider_id == provider_id)
|
|
.order_by(NotificationDigestQueue.created_at)
|
|
)
|
|
queue_entries = list(result.scalars().all())
|
|
|
|
if not queue_entries:
|
|
logger.debug("No queued notifications for provider %s", provider.name)
|
|
return
|
|
|
|
# Build digest message
|
|
title = f"Daily Digest - {len(queue_entries)} Events"
|
|
|
|
# Group by event type
|
|
events_by_type: dict[str, list] = {}
|
|
for entry in queue_entries:
|
|
if entry.event_type not in events_by_type:
|
|
events_by_type[entry.event_type] = []
|
|
events_by_type[entry.event_type].append(entry)
|
|
|
|
# Format the digest body
|
|
body_parts = []
|
|
for event_type, entries in events_by_type.items():
|
|
event_label = event_type.replace("_", " ").title()
|
|
body_parts.append(f"== {event_label} ({len(entries)}) ==")
|
|
for entry in entries:
|
|
time_str = entry.created_at.strftime("%H:%M")
|
|
printer_info = f"[{entry.printer_name}] " if entry.printer_name else ""
|
|
body_parts.append(f" {time_str} {printer_info}{entry.title}")
|
|
body_parts.append("")
|
|
|
|
body = "\n".join(body_parts)
|
|
|
|
# Send the digest
|
|
success, error = await self._send_to_provider(provider, title, body, db)
|
|
|
|
# Log the digest
|
|
await self._log_notification(
|
|
db=db,
|
|
provider_id=provider.id,
|
|
event_type="daily_digest",
|
|
title=title,
|
|
message=body,
|
|
success=success,
|
|
error_message=error if not success else None,
|
|
)
|
|
|
|
# Clear the queue
|
|
for entry in queue_entries:
|
|
await db.delete(entry)
|
|
await db.commit()
|
|
|
|
if success:
|
|
logger.info("Sent daily digest with %s events to %s", len(queue_entries), provider.name)
|
|
else:
|
|
logger.warning("Failed to send daily digest to %s: %s", provider.name, error)
|
|
|
|
async def check_and_send_digests(self):
|
|
"""Check all providers and send digests if it's their scheduled time."""
|
|
from backend.app.core.database import async_session
|
|
|
|
current_time = datetime.now().strftime("%H:%M")
|
|
|
|
# Avoid duplicate checks within the same minute
|
|
if current_time == self._last_digest_check:
|
|
return
|
|
self._last_digest_check = current_time
|
|
|
|
async with async_session() as db:
|
|
# Find all providers with digest enabled at this time
|
|
result = await db.execute(
|
|
select(NotificationProvider).where(
|
|
NotificationProvider.enabled.is_(True),
|
|
NotificationProvider.daily_digest_enabled.is_(True),
|
|
NotificationProvider.daily_digest_time == current_time,
|
|
)
|
|
)
|
|
providers = result.scalars().all()
|
|
|
|
for provider in providers:
|
|
try:
|
|
await self.send_digest(provider.id)
|
|
except Exception as e:
|
|
logger.error("Error sending digest for provider %s: %s", provider.id, e)
|
|
|
|
def start_digest_scheduler(self):
|
|
"""Start the background scheduler for daily digest notifications."""
|
|
if self._digest_scheduler_task is None:
|
|
self._digest_scheduler_task = asyncio.create_task(self._digest_scheduler_loop())
|
|
logger.info("Notification digest scheduler started")
|
|
|
|
def stop_digest_scheduler(self):
|
|
"""Stop the background scheduler for daily digests."""
|
|
if self._digest_scheduler_task:
|
|
self._digest_scheduler_task.cancel()
|
|
self._digest_scheduler_task = None
|
|
logger.info("Notification digest scheduler stopped")
|
|
|
|
async def _digest_scheduler_loop(self):
|
|
"""Background loop that checks for scheduled digests every minute."""
|
|
while True:
|
|
try:
|
|
await self.check_and_send_digests()
|
|
except Exception as e:
|
|
logger.error("Error in digest scheduler: %s", e)
|
|
|
|
# Wait until the next minute
|
|
await asyncio.sleep(60)
|
|
|
|
|
|
# Global instance
|
|
notification_service = NotificationService()
|