|
|
@@ -1,7 +1,9 @@
|
|
|
"""Notification service for sending push notifications via various providers."""
|
|
|
|
|
|
+import asyncio
|
|
|
import json
|
|
|
import logging
|
|
|
+import re
|
|
|
import smtplib
|
|
|
from datetime import datetime
|
|
|
from email.mime.multipart import MIMEMultipart
|
|
|
@@ -13,9 +15,8 @@ import httpx
|
|
|
from sqlalchemy import select
|
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
|
|
-from backend.app.models.notification import NotificationProvider
|
|
|
-from backend.app.models.settings import Settings
|
|
|
-from backend.app.i18n import Translator
|
|
|
+from backend.app.models.notification import NotificationLog, NotificationProvider, NotificationDigestQueue
|
|
|
+from backend.app.models.notification_template import NotificationTemplate
|
|
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
@@ -25,6 +26,9 @@ class NotificationService:
|
|
|
|
|
|
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."""
|
|
|
@@ -66,136 +70,77 @@ class NotificationService:
|
|
|
logger.warning(f"Invalid quiet hours format for provider {provider.name}")
|
|
|
return False
|
|
|
|
|
|
- async def _get_notification_language(self, db: AsyncSession) -> str:
|
|
|
- """Get the notification language from settings."""
|
|
|
+ async def _get_template(self, db: AsyncSession, event_type: str) -> NotificationTemplate | None:
|
|
|
+ """Get a notification template by event type."""
|
|
|
+ # Check cache first
|
|
|
+ if event_type in self._template_cache:
|
|
|
+ return self._template_cache[event_type]
|
|
|
+
|
|
|
result = await db.execute(
|
|
|
- select(Settings).where(Settings.key == "notification_language")
|
|
|
+ select(NotificationTemplate).where(NotificationTemplate.event_type == event_type)
|
|
|
)
|
|
|
- setting = result.scalar_one_or_none()
|
|
|
- return setting.value if setting else "en"
|
|
|
+ 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
|
|
|
|
|
|
- def _format_duration(self, seconds: int | None, translator: Translator) -> str:
|
|
|
+ def _format_duration(self, seconds: int | None) -> str:
|
|
|
"""Format duration in seconds to human-readable string."""
|
|
|
if seconds is None:
|
|
|
- return translator.t("notification.unknown")
|
|
|
+ return "Unknown"
|
|
|
hours = seconds // 3600
|
|
|
minutes = (seconds % 3600) // 60
|
|
|
if hours > 0:
|
|
|
return f"{hours}h {minutes}m"
|
|
|
return f"{minutes}m"
|
|
|
|
|
|
- def _build_print_start_message(self, printer_name: str, data: dict, translator: Translator) -> tuple[str, str]:
|
|
|
- """Build notification message for print start event."""
|
|
|
- filename = data.get("filename", translator.t("notification.unknown"))
|
|
|
- # Clean up filename
|
|
|
- if filename.endswith(".gcode.3mf"):
|
|
|
- filename = filename[:-10]
|
|
|
- elif filename.endswith(".3mf"):
|
|
|
- filename = filename[:-4]
|
|
|
-
|
|
|
- title = translator.t("notification.print_started")
|
|
|
-
|
|
|
- estimated_time = data.get("raw_data", {}).get("print", {}).get("mc_remaining_time")
|
|
|
- time_str = self._format_duration(estimated_time * 60 if estimated_time else None, translator)
|
|
|
-
|
|
|
- message = f"{printer_name}: {filename}\n{translator.t('notification.estimated')}: {time_str}"
|
|
|
- return title, message
|
|
|
-
|
|
|
- def _build_print_complete_message(
|
|
|
- self, printer_name: str, status: str, data: dict, translator: Translator, archive_data: dict | None = None
|
|
|
- ) -> tuple[str, str]:
|
|
|
- """Build notification message for print complete event."""
|
|
|
- filename = data.get("filename", translator.t("notification.unknown"))
|
|
|
+ def _clean_filename(self, filename: str) -> str:
|
|
|
+ """Remove file extensions from filename."""
|
|
|
if filename.endswith(".gcode.3mf"):
|
|
|
- filename = filename[:-10]
|
|
|
+ return filename[:-10]
|
|
|
elif filename.endswith(".3mf"):
|
|
|
- filename = filename[:-4]
|
|
|
-
|
|
|
- if status == "completed":
|
|
|
- title = translator.t("notification.print_completed")
|
|
|
- elif status == "failed":
|
|
|
- title = translator.t("notification.print_failed")
|
|
|
- elif status in ("aborted", "stopped", "cancelled"):
|
|
|
- title = translator.t("notification.print_stopped")
|
|
|
- else:
|
|
|
- title = translator.t("notification.print_ended")
|
|
|
+ return filename[:-4]
|
|
|
+ return filename
|
|
|
|
|
|
- lines = [f"{printer_name}: {filename}"]
|
|
|
-
|
|
|
- if archive_data:
|
|
|
- # Add print time if available
|
|
|
- if archive_data.get("print_time_seconds"):
|
|
|
- lines.append(f"{translator.t('notification.time')}: {self._format_duration(archive_data['print_time_seconds'], translator)}")
|
|
|
- # Add filament used if available
|
|
|
- if archive_data.get("actual_filament_grams"):
|
|
|
- lines.append(f"{translator.t('notification.filament')}: {archive_data['actual_filament_grams']:.1f}g")
|
|
|
- # Add failure reason if failed
|
|
|
- if status == "failed" and archive_data.get("failure_reason"):
|
|
|
- lines.append(f"{translator.t('notification.reason')}: {archive_data['failure_reason']}")
|
|
|
-
|
|
|
- message = "\n".join(lines)
|
|
|
- return title, message
|
|
|
-
|
|
|
- def _build_progress_message(
|
|
|
- self, printer_name: str, filename: str, progress: int, translator: Translator
|
|
|
+ async def _build_message_from_template(
|
|
|
+ self, db: AsyncSession, event_type: str, variables: dict[str, Any]
|
|
|
) -> tuple[str, str]:
|
|
|
- """Build notification message for print progress milestone."""
|
|
|
- if filename.endswith(".gcode.3mf"):
|
|
|
- filename = filename[:-10]
|
|
|
- elif filename.endswith(".3mf"):
|
|
|
- filename = filename[:-4]
|
|
|
-
|
|
|
- title = translator.t("notification.print_progress", progress=progress)
|
|
|
- message = f"{printer_name}: {filename}"
|
|
|
- return title, message
|
|
|
+ """Build notification title and body from template."""
|
|
|
+ # Add common variables
|
|
|
+ variables["timestamp"] = datetime.now().strftime("%Y-%m-%d %H:%M")
|
|
|
+ variables["app_name"] = "BambuTrack"
|
|
|
|
|
|
- def _build_printer_offline_message(self, printer_name: str, translator: Translator) -> tuple[str, str]:
|
|
|
- """Build notification message for printer offline event."""
|
|
|
- title = translator.t("notification.printer_offline")
|
|
|
- message = translator.t("notification.printer_disconnected", printer=printer_name)
|
|
|
- return title, message
|
|
|
+ template = await self._get_template(db, event_type)
|
|
|
+ if not template:
|
|
|
+ # Fallback to simple message
|
|
|
+ logger.warning(f"Template not found for event type: {event_type}")
|
|
|
+ return event_type.replace("_", " ").title(), str(variables)
|
|
|
|
|
|
- def _build_printer_error_message(
|
|
|
- self, printer_name: str, error_type: str, translator: Translator, error_detail: str | None = None
|
|
|
- ) -> tuple[str, str]:
|
|
|
- """Build notification message for printer error event."""
|
|
|
- title = translator.t("notification.printer_error", error_type=error_type)
|
|
|
- message = f"{printer_name}"
|
|
|
- if error_detail:
|
|
|
- message += f"\n{error_detail}"
|
|
|
- return title, message
|
|
|
-
|
|
|
- def _build_filament_low_message(
|
|
|
- self, printer_name: str, slot: int, remaining_percent: int, translator: Translator
|
|
|
- ) -> tuple[str, str]:
|
|
|
- """Build notification message for low filament event."""
|
|
|
- title = translator.t("notification.filament_low")
|
|
|
- message = translator.t("notification.slot_at_percent", printer=printer_name, slot=slot, percent=remaining_percent)
|
|
|
- return title, message
|
|
|
+ title = self._render_template(template.title_template, variables)
|
|
|
+ body = self._render_template(template.body_template, variables)
|
|
|
|
|
|
- def _build_maintenance_due_message(
|
|
|
- self, printer_name: str, maintenance_items: list[dict], translator: Translator
|
|
|
- ) -> tuple[str, str]:
|
|
|
- """Build notification message for maintenance due event."""
|
|
|
- title = translator.t("notification.maintenance_due")
|
|
|
- lines = [f"{printer_name}:"]
|
|
|
- for item in maintenance_items:
|
|
|
- status = translator.t("notification.overdue") if item.get("is_due") else translator.t("notification.soon")
|
|
|
- lines.append(f"• {item['name']} ({status})")
|
|
|
- message = "\n".join(lines)
|
|
|
- return title, message
|
|
|
+ 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."""
|
|
|
- lang = "en"
|
|
|
if db:
|
|
|
- lang = await self._get_notification_language(db)
|
|
|
- translator = Translator(lang)
|
|
|
-
|
|
|
- title = translator.t("notification.test_title")
|
|
|
- message = translator.t("notification.test_message")
|
|
|
+ title, message = await self._build_message_from_template(db, "test", {})
|
|
|
+ else:
|
|
|
+ title = "BambuTrack Test"
|
|
|
+ message = "This is a test notification. If you see this, notifications are working!"
|
|
|
|
|
|
try:
|
|
|
if provider_type == "callmebot":
|
|
|
@@ -208,6 +153,10 @@ class NotificationService:
|
|
|
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)
|
|
|
else:
|
|
|
return False, f"Unknown provider type: {provider_type}"
|
|
|
except Exception as e:
|
|
|
@@ -366,6 +315,70 @@ class NotificationService:
|
|
|
except Exception as e:
|
|
|
return False, f"Email error: {str(e)}"
|
|
|
|
|
|
+ async def _send_discord(self, config: dict, title: str, message: str) -> 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/"):
|
|
|
+ return False, "Invalid Discord webhook URL"
|
|
|
+
|
|
|
+ # Discord embed format for nicer messages
|
|
|
+ data = {
|
|
|
+ "embeds": [{
|
|
|
+ "title": title,
|
|
|
+ "description": message,
|
|
|
+ "color": 0x00AE42, # Bambu green
|
|
|
+ }]
|
|
|
+ }
|
|
|
+
|
|
|
+ client = await self._get_client()
|
|
|
+ response = await client.post(webhook_url, json=data)
|
|
|
+
|
|
|
+ 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) -> tuple[bool, str]:
|
|
|
+ """Send notification via generic webhook (POST JSON)."""
|
|
|
+ webhook_url = config.get("webhook_url", "").strip()
|
|
|
+ auth_header = config.get("auth_header", "").strip()
|
|
|
+ custom_field_title = config.get("field_title", "title").strip() or "title"
|
|
|
+ custom_field_message = config.get("field_message", "message").strip() or "message"
|
|
|
+
|
|
|
+ if not webhook_url:
|
|
|
+ return False, "Webhook URL is required"
|
|
|
+
|
|
|
+ # Build payload with custom field names
|
|
|
+ data = {
|
|
|
+ custom_field_title: title,
|
|
|
+ custom_field_message: message,
|
|
|
+ "timestamp": datetime.now().isoformat(),
|
|
|
+ "source": "BambuTrack",
|
|
|
+ }
|
|
|
+
|
|
|
+ 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, f"HTTP {response.status_code}: {response.text[:200]}"
|
|
|
+ except Exception as e:
|
|
|
+ return False, f"Webhook error: {str(e)}"
|
|
|
+
|
|
|
async def _send_to_provider(
|
|
|
self, provider: NotificationProvider, title: str, message: str
|
|
|
) -> tuple[bool, str]:
|
|
|
@@ -388,6 +401,10 @@ class NotificationService:
|
|
|
return await self._send_telegram(config, f"*{title}*\n{message}")
|
|
|
elif provider.provider_type == "email":
|
|
|
return await self._send_email(config, title, message)
|
|
|
+ elif provider.provider_type == "discord":
|
|
|
+ return await self._send_discord(config, title, message)
|
|
|
+ elif provider.provider_type == "webhook":
|
|
|
+ return await self._send_webhook(config, title, message)
|
|
|
else:
|
|
|
return False, f"Unknown provider type: {provider.provider_type}"
|
|
|
except Exception as e:
|
|
|
@@ -431,18 +448,75 @@ class NotificationService:
|
|
|
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(f"Failed to log notification: {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,
|
|
|
):
|
|
|
- """Send notification to multiple providers."""
|
|
|
+ """Send notification to multiple providers and log the results."""
|
|
|
for provider in providers:
|
|
|
try:
|
|
|
+ # Check if provider wants digest mode
|
|
|
+ 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,
|
|
|
+ )
|
|
|
+ continue
|
|
|
+
|
|
|
success, error = await self._send_to_provider(provider, title, message)
|
|
|
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(f"Sent notification via {provider.name}")
|
|
|
else:
|
|
|
@@ -450,6 +524,17 @@ class NotificationService:
|
|
|
except Exception as e:
|
|
|
logger.exception(f"Error sending notification via {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_print_start(
|
|
|
self, printer_id: int, printer_name: str, data: dict, db: AsyncSession
|
|
|
@@ -461,11 +546,19 @@ class NotificationService:
|
|
|
logger.info(f"No notification providers configured for print_start event on printer {printer_id}")
|
|
|
return
|
|
|
|
|
|
- lang = await self._get_notification_language(db)
|
|
|
- translator = Translator(lang)
|
|
|
+ filename = self._clean_filename(data.get("filename", "Unknown"))
|
|
|
+ estimated_time = data.get("raw_data", {}).get("print", {}).get("mc_remaining_time")
|
|
|
+ time_str = self._format_duration(estimated_time * 60 if estimated_time else None)
|
|
|
+
|
|
|
+ variables = {
|
|
|
+ "printer": printer_name,
|
|
|
+ "filename": filename,
|
|
|
+ "estimated_time": time_str,
|
|
|
+ }
|
|
|
+
|
|
|
logger.info(f"Found {len(providers)} providers for print_start: {[p.name for p in providers]}")
|
|
|
- title, message = self._build_print_start_message(printer_name, data, translator)
|
|
|
- await self._send_to_providers(providers, title, message, db)
|
|
|
+ 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)
|
|
|
|
|
|
async def on_print_complete(
|
|
|
self,
|
|
|
@@ -478,28 +571,48 @@ class NotificationService:
|
|
|
):
|
|
|
"""Handle print complete event - send notifications to relevant providers."""
|
|
|
logger.info(f"on_print_complete called for printer {printer_id} ({printer_name}), status={status}")
|
|
|
- # Determine which event type this is
|
|
|
+
|
|
|
+ # 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:
|
|
|
- # Unknown status, default to on_print_complete
|
|
|
logger.warning(f"Unknown print status '{status}', defaulting to on_print_complete")
|
|
|
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(f"No notification providers configured for {event_field} event on printer {printer_id}")
|
|
|
return
|
|
|
|
|
|
- lang = await self._get_notification_language(db)
|
|
|
- translator = Translator(lang)
|
|
|
+ filename = self._clean_filename(data.get("filename", "Unknown"))
|
|
|
+
|
|
|
+ variables = {
|
|
|
+ "printer": printer_name,
|
|
|
+ "filename": filename,
|
|
|
+ "duration": "",
|
|
|
+ "filament_grams": "",
|
|
|
+ "reason": "",
|
|
|
+ }
|
|
|
+
|
|
|
+ if archive_data:
|
|
|
+ if archive_data.get("print_time_seconds"):
|
|
|
+ variables["duration"] = self._format_duration(archive_data["print_time_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"]
|
|
|
+
|
|
|
logger.info(f"Found {len(providers)} providers for {event_field}: {[p.name for p in providers]}")
|
|
|
- title, message = self._build_print_complete_message(printer_name, status, data, translator, archive_data)
|
|
|
- await self._send_to_providers(providers, title, message, db)
|
|
|
+ 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)
|
|
|
|
|
|
async def on_print_progress(
|
|
|
self,
|
|
|
@@ -508,16 +621,22 @@ class NotificationService:
|
|
|
filename: str,
|
|
|
progress: int,
|
|
|
db: AsyncSession,
|
|
|
+ remaining_time: int | 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
|
|
|
|
|
|
- lang = await self._get_notification_language(db)
|
|
|
- translator = Translator(lang)
|
|
|
- title, message = self._build_progress_message(printer_name, filename, progress, translator)
|
|
|
- await self._send_to_providers(providers, title, message, db)
|
|
|
+ variables = {
|
|
|
+ "printer": printer_name,
|
|
|
+ "filename": self._clean_filename(filename),
|
|
|
+ "progress": str(progress),
|
|
|
+ "remaining_time": self._format_duration(remaining_time) if remaining_time else "",
|
|
|
+ }
|
|
|
+
|
|
|
+ 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)
|
|
|
|
|
|
async def on_printer_offline(
|
|
|
self, printer_id: int, printer_name: str, db: AsyncSession
|
|
|
@@ -527,10 +646,10 @@ class NotificationService:
|
|
|
if not providers:
|
|
|
return
|
|
|
|
|
|
- lang = await self._get_notification_language(db)
|
|
|
- translator = Translator(lang)
|
|
|
- title, message = self._build_printer_offline_message(printer_name, translator)
|
|
|
- await self._send_to_providers(providers, title, message, db)
|
|
|
+ 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)
|
|
|
|
|
|
async def on_printer_error(
|
|
|
self,
|
|
|
@@ -545,10 +664,14 @@ class NotificationService:
|
|
|
if not providers:
|
|
|
return
|
|
|
|
|
|
- lang = await self._get_notification_language(db)
|
|
|
- translator = Translator(lang)
|
|
|
- title, message = self._build_printer_error_message(printer_name, error_type, translator, error_detail)
|
|
|
- await self._send_to_providers(providers, title, message, db)
|
|
|
+ variables = {
|
|
|
+ "printer": printer_name,
|
|
|
+ "error_type": error_type,
|
|
|
+ "error_detail": error_detail or "",
|
|
|
+ }
|
|
|
+
|
|
|
+ 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)
|
|
|
|
|
|
async def on_filament_low(
|
|
|
self,
|
|
|
@@ -557,16 +680,22 @@ class NotificationService:
|
|
|
slot: int,
|
|
|
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
|
|
|
|
|
|
- lang = await self._get_notification_language(db)
|
|
|
- translator = Translator(lang)
|
|
|
- title, message = self._build_filament_low_message(printer_name, slot, remaining_percent, translator)
|
|
|
- await self._send_to_providers(providers, title, message, db)
|
|
|
+ variables = {
|
|
|
+ "printer": printer_name,
|
|
|
+ "slot": str(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)
|
|
|
|
|
|
async def on_maintenance_due(
|
|
|
self,
|
|
|
@@ -584,11 +713,176 @@ class NotificationService:
|
|
|
logger.info(f"No notification providers configured for maintenance_due event on printer {printer_id}")
|
|
|
return
|
|
|
|
|
|
- lang = await self._get_notification_language(db)
|
|
|
- translator = Translator(lang)
|
|
|
+ # 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(f"Found {len(providers)} providers for maintenance_due: {[p.name for p in providers]}")
|
|
|
- title, message = self._build_maintenance_due_message(printer_name, maintenance_items, translator)
|
|
|
- await self._send_to_providers(providers, title, message, db)
|
|
|
+ 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)
|
|
|
+
|
|
|
+ def clear_template_cache(self):
|
|
|
+ """Clear the template cache. Call this when templates are updated."""
|
|
|
+ self._template_cache.clear()
|
|
|
+
|
|
|
+ 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(f"Queued notification for digest: {event_type} for provider {provider.name}")
|
|
|
+ except Exception as e:
|
|
|
+ logger.warning(f"Failed to queue notification for digest: {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(f"No queued notifications for provider {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)
|
|
|
+
|
|
|
+ # 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(f"Sent daily digest with {len(queue_entries)} events to {provider.name}")
|
|
|
+ else:
|
|
|
+ logger.warning(f"Failed to send daily digest to {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 == True,
|
|
|
+ NotificationProvider.daily_digest_enabled == 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(f"Error sending digest for provider {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(f"Error in digest scheduler: {e}")
|
|
|
+
|
|
|
+ # Wait until the next minute
|
|
|
+ await asyncio.sleep(60)
|
|
|
|
|
|
|
|
|
# Global instance
|