notification_service.py 58 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478
  1. """Notification service for sending push notifications via various providers."""
  2. import asyncio
  3. import json
  4. import logging
  5. import re
  6. import smtplib
  7. from datetime import datetime, timedelta, timezone
  8. from email.mime.multipart import MIMEMultipart
  9. from email.mime.text import MIMEText
  10. from typing import Any
  11. from urllib.parse import quote
  12. import httpx
  13. from sqlalchemy import select
  14. from sqlalchemy.ext.asyncio import AsyncSession
  15. from backend.app.models.notification import NotificationDigestQueue, NotificationLog, NotificationProvider
  16. from backend.app.models.notification_template import NotificationTemplate
  17. logger = logging.getLogger(__name__)
  18. class NotificationService:
  19. """Service for sending notifications through various providers."""
  20. def __init__(self):
  21. self._http_client: httpx.AsyncClient | None = None
  22. self._template_cache: dict[str, NotificationTemplate] = {}
  23. self._digest_scheduler_task: asyncio.Task | None = None
  24. self._last_digest_check: str = "" # "HH:MM" to avoid duplicate checks
  25. async def _get_client(self) -> httpx.AsyncClient:
  26. """Get or create HTTP client."""
  27. if self._http_client is None or self._http_client.is_closed:
  28. self._http_client = httpx.AsyncClient(timeout=30.0)
  29. return self._http_client
  30. async def close(self):
  31. """Close HTTP client."""
  32. if self._http_client and not self._http_client.is_closed:
  33. await self._http_client.aclose()
  34. def _is_in_quiet_hours(self, provider: NotificationProvider) -> bool:
  35. """Check if current time is within provider's quiet hours."""
  36. if not provider.quiet_hours_enabled:
  37. return False
  38. if not provider.quiet_hours_start or not provider.quiet_hours_end:
  39. return False
  40. try:
  41. now = datetime.now()
  42. current_time = now.hour * 60 + now.minute
  43. start_parts = provider.quiet_hours_start.split(":")
  44. end_parts = provider.quiet_hours_end.split(":")
  45. start_minutes = int(start_parts[0]) * 60 + int(start_parts[1])
  46. end_minutes = int(end_parts[0]) * 60 + int(end_parts[1])
  47. # Handle overnight quiet hours (e.g., 22:00 to 07:00)
  48. if start_minutes > end_minutes:
  49. # Quiet hours span midnight
  50. return current_time >= start_minutes or current_time < end_minutes
  51. else:
  52. # Same day quiet hours
  53. return start_minutes <= current_time < end_minutes
  54. except (ValueError, TypeError, AttributeError):
  55. logger.warning("Invalid quiet hours format for provider %s", provider.name)
  56. return False
  57. async def _get_template(self, db: AsyncSession, event_type: str) -> NotificationTemplate | None:
  58. """Get a notification template by event type."""
  59. # Check cache first
  60. if event_type in self._template_cache:
  61. return self._template_cache[event_type]
  62. result = await db.execute(select(NotificationTemplate).where(NotificationTemplate.event_type == event_type))
  63. template = result.scalar_one_or_none()
  64. if template:
  65. self._template_cache[event_type] = template
  66. return template
  67. def _render_template(self, template_str: str, variables: dict[str, Any]) -> str:
  68. """Render a template string with variables. Missing variables become empty."""
  69. result = template_str
  70. for key, value in variables.items():
  71. result = result.replace("{" + key + "}", str(value) if value is not None else "")
  72. # Remove any remaining unreplaced placeholders
  73. result = re.sub(r"\{[a-z_]+\}", "", result)
  74. return result
  75. async def _format_eta(self, seconds: int | None, db: AsyncSession) -> str:
  76. """Format ETA as wall-clock time, respecting user's time_format setting."""
  77. if not seconds or seconds <= 0:
  78. return "Unknown"
  79. from backend.app.api.routes.settings import get_setting
  80. eta_time = datetime.now() + timedelta(seconds=seconds)
  81. time_format = await get_setting(db, "time_format")
  82. if time_format == "12h":
  83. return eta_time.strftime("%I:%M %p").lstrip("0")
  84. # Default to 24h for "24h", "system", or unset
  85. return eta_time.strftime("%H:%M")
  86. def _format_duration(self, seconds: int | None) -> str:
  87. """Format duration in seconds to human-readable string."""
  88. if seconds is None:
  89. return "Unknown"
  90. hours = seconds // 3600
  91. minutes = (seconds % 3600) // 60
  92. if hours > 0:
  93. return f"{hours}h {minutes}m"
  94. return f"{minutes}m"
  95. def _clean_filename(self, filename: str) -> str:
  96. """Extract filename and remove file extensions."""
  97. import os
  98. # Strip path prefix (e.g., /data/Metadata/plate_5.gcode -> plate_5.gcode)
  99. filename = os.path.basename(filename)
  100. # Remove common extensions
  101. if filename.endswith(".gcode.3mf"):
  102. return filename[:-10]
  103. elif filename.endswith(".gcode"):
  104. return filename[:-6]
  105. elif filename.endswith(".3mf"):
  106. return filename[:-4]
  107. return filename
  108. async def _build_message_from_template(
  109. self, db: AsyncSession, event_type: str, variables: dict[str, Any]
  110. ) -> tuple[str, str]:
  111. """Build notification title and body from template."""
  112. # Add common variables
  113. variables["timestamp"] = datetime.now().strftime("%Y-%m-%d %H:%M")
  114. variables["app_name"] = "Bambuddy"
  115. template = await self._get_template(db, event_type)
  116. if not template:
  117. # Fallback to simple message
  118. logger.warning("Template not found for event type: %s", event_type)
  119. return event_type.replace("_", " ").title(), str(variables)
  120. title = self._render_template(template.title_template, variables)
  121. body = self._render_template(template.body_template, variables)
  122. return title, body
  123. async def send_test_notification(
  124. self, provider_type: str, config: dict[str, Any], db: AsyncSession | None = None
  125. ) -> tuple[bool, str]:
  126. """Send a test notification to verify configuration."""
  127. if db:
  128. title, message = await self._build_message_from_template(db, "test", {})
  129. else:
  130. title = "Bambuddy Test"
  131. message = "This is a test notification. If you see this, notifications are working!"
  132. try:
  133. if provider_type == "callmebot":
  134. return await self._send_callmebot(config, f"{title}\n{message}")
  135. elif provider_type == "ntfy":
  136. return await self._send_ntfy(config, title, message)
  137. elif provider_type == "pushover":
  138. return await self._send_pushover(config, title, message)
  139. elif provider_type == "telegram":
  140. return await self._send_telegram(config, f"*{title}*\n{message}")
  141. elif provider_type == "email":
  142. return await self._send_email(config, title, message)
  143. elif provider_type == "discord":
  144. return await self._send_discord(config, title, message)
  145. elif provider_type == "webhook":
  146. return await self._send_webhook(config, title, message)
  147. elif provider_type == "homeassistant":
  148. return await self._send_homeassistant(config, title, message, db=db)
  149. else:
  150. return False, f"Unknown provider type: {provider_type}"
  151. except Exception as e:
  152. logger.exception("Error sending test notification via %s", provider_type)
  153. return False, str(e)
  154. async def _send_callmebot(self, config: dict, message: str) -> tuple[bool, str]:
  155. """Send notification via CallMeBot (WhatsApp)."""
  156. phone = config.get("phone", "").strip()
  157. apikey = config.get("apikey", "").strip()
  158. if not phone or not apikey:
  159. return False, "Phone number and API key are required"
  160. # URL encode the message
  161. encoded_message = quote(message)
  162. url = f"https://api.callmebot.com/whatsapp.php?phone={phone}&text={encoded_message}&apikey={apikey}"
  163. client = await self._get_client()
  164. response = await client.get(url)
  165. if response.status_code == 200:
  166. return True, "Message sent successfully"
  167. else:
  168. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  169. async def _send_ntfy(
  170. self, config: dict, title: str, message: str, image_data: bytes | None = None
  171. ) -> tuple[bool, str]:
  172. """Send notification via ntfy."""
  173. server = config.get("server", "https://ntfy.sh").rstrip("/")
  174. topic = config.get("topic", "").strip()
  175. auth_token = config.get("auth_token", "").strip()
  176. if not topic:
  177. return False, "Topic is required"
  178. url = f"{server}/{topic}"
  179. headers = {"Title": title}
  180. if auth_token:
  181. headers["Authorization"] = f"Bearer {auth_token}"
  182. client = await self._get_client()
  183. if image_data:
  184. # ntfy supports image attachments via multipart form-data.
  185. # HTTP headers cannot contain newlines, but ntfy interprets
  186. # literal \n (backslash-n) as newlines in the Message header.
  187. headers["Filename"] = "photo.jpg"
  188. headers["Message"] = message.replace("\n", "\\n")
  189. response = await client.put(url, content=image_data, headers=headers)
  190. if response.status_code == 400 and "attachments not allowed" in response.text:
  191. # Server has attachments disabled — retry without the image
  192. headers.pop("Filename", None)
  193. headers.pop("Message", None)
  194. response = await client.post(url, content=message, headers=headers)
  195. else:
  196. response = await client.post(url, content=message, headers=headers)
  197. if response.status_code in (200, 204):
  198. return True, "Message sent successfully"
  199. else:
  200. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  201. async def _send_pushover(
  202. self, config: dict, title: str, message: str, image_data: bytes | None = None
  203. ) -> tuple[bool, str]:
  204. """Send notification via Pushover.
  205. Args:
  206. config: Provider configuration with user_key, app_token, priority
  207. title: Notification title
  208. message: Notification body
  209. image_data: Optional JPEG image bytes to attach (max 2.5MB)
  210. """
  211. user_key = config.get("user_key", "").strip()
  212. app_token = config.get("app_token", "").strip()
  213. priority = config.get("priority", 0)
  214. if not user_key or not app_token:
  215. return False, "User key and app token are required"
  216. url = "https://api.pushover.net/1/messages.json"
  217. data = {
  218. "token": app_token,
  219. "user": user_key,
  220. "title": title,
  221. "message": message,
  222. "priority": priority,
  223. }
  224. client = await self._get_client()
  225. if image_data:
  226. # Pushover supports image attachments via multipart form-data
  227. files = {"attachment": ("photo.jpg", image_data, "image/jpeg")}
  228. response = await client.post(url, data=data, files=files)
  229. else:
  230. response = await client.post(url, data=data)
  231. if response.status_code == 200:
  232. return True, "Message sent successfully"
  233. else:
  234. try:
  235. error_data = response.json()
  236. errors = error_data.get("errors", [])
  237. return False, f"Pushover error: {', '.join(errors)}"
  238. except Exception:
  239. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  240. async def _send_telegram(self, config: dict, message: str, image_data: bytes | None = None) -> tuple[bool, str]:
  241. """Send notification via Telegram bot."""
  242. bot_token = config.get("bot_token", "").strip()
  243. chat_id = config.get("chat_id", "").strip()
  244. if not bot_token or not chat_id:
  245. return False, "Bot token and chat ID are required"
  246. # Escape underscores in the message body so Telegram Markdown
  247. # parsing doesn't break on job names like "A1_plate_8" or error
  248. # codes like "0300_0001". The title is already wrapped in *bold*
  249. # markers, so only escape after the first newline.
  250. if "\n" in message:
  251. title_part, body_part = message.split("\n", 1)
  252. body_part = body_part.replace("_", "\\_")
  253. message = f"{title_part}\n{body_part}"
  254. client = await self._get_client()
  255. if image_data:
  256. # Use sendPhoto to attach the thumbnail with the caption
  257. url = f"https://api.telegram.org/bot{bot_token}/sendPhoto"
  258. response = await client.post(
  259. url,
  260. data={"chat_id": chat_id, "caption": message, "parse_mode": "Markdown"},
  261. files={"photo": ("photo.jpg", image_data, "image/jpeg")},
  262. )
  263. else:
  264. url = f"https://api.telegram.org/bot{bot_token}/sendMessage"
  265. data = {
  266. "chat_id": chat_id,
  267. "text": message,
  268. "parse_mode": "Markdown",
  269. }
  270. response = await client.post(url, json=data)
  271. if response.status_code == 200:
  272. result = response.json()
  273. if result.get("ok"):
  274. return True, "Message sent successfully"
  275. else:
  276. return False, f"Telegram error: {result.get('description', 'Unknown error')}"
  277. else:
  278. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  279. async def _send_email(self, config: dict, subject: str, body: str) -> tuple[bool, str]:
  280. """Send notification via email (SMTP)."""
  281. smtp_server = config.get("smtp_server", "").strip()
  282. smtp_port = int(config.get("smtp_port", 587))
  283. username = config.get("username", "").strip()
  284. password = config.get("password", "").strip()
  285. from_email = config.get("from_email", "").strip()
  286. to_email = config.get("to_email", "").strip()
  287. # Security: "starttls" (port 587), "ssl" (port 465), "none" (port 25)
  288. security = config.get("security", "starttls")
  289. # Authentication: "true" or "false"
  290. auth_enabled = config.get("auth_enabled", "true").lower() == "true"
  291. if not all([smtp_server, from_email, to_email]):
  292. return False, "SMTP server, from email, and to email are required"
  293. if auth_enabled and not all([username, password]):
  294. return False, "Username and password are required when authentication is enabled"
  295. try:
  296. msg = MIMEMultipart()
  297. msg["From"] = from_email
  298. msg["To"] = to_email
  299. msg["Subject"] = f"[Bambuddy] {subject}"
  300. msg.attach(MIMEText(body, "plain"))
  301. if security == "ssl":
  302. # Direct SSL connection (typically port 465)
  303. server = smtplib.SMTP_SSL(smtp_server, smtp_port)
  304. elif security == "starttls":
  305. # STARTTLS upgrade (typically port 587)
  306. server = smtplib.SMTP(smtp_server, smtp_port)
  307. server.starttls()
  308. else:
  309. # No encryption (typically port 25) - use with caution
  310. server = smtplib.SMTP(smtp_server, smtp_port)
  311. if auth_enabled:
  312. server.login(username, password)
  313. server.sendmail(from_email, to_email, msg.as_string())
  314. server.quit()
  315. return True, "Email sent successfully"
  316. except smtplib.SMTPAuthenticationError:
  317. return False, "SMTP authentication failed - check username/password"
  318. except smtplib.SMTPException as e:
  319. return False, f"SMTP error: {str(e)}"
  320. except Exception as e:
  321. return False, f"Email error: {str(e)}"
  322. async def _send_discord(
  323. self, config: dict, title: str, message: str, image_data: bytes | None = None
  324. ) -> tuple[bool, str]:
  325. """Send notification via Discord webhook."""
  326. webhook_url = config.get("webhook_url", "").strip()
  327. if not webhook_url:
  328. return False, "Webhook URL is required"
  329. if not webhook_url.startswith("https://discord.com/api/webhooks/"):
  330. return False, "Invalid Discord webhook URL"
  331. # Discord embed format for nicer messages
  332. embed = {
  333. "title": title,
  334. "description": message,
  335. "color": 0x00AE42, # Bambu green
  336. }
  337. client = await self._get_client()
  338. if image_data:
  339. # Attach image via multipart form-data and reference in embed
  340. embed["image"] = {"url": "attachment://photo.jpg"}
  341. payload = {"embeds": [embed]}
  342. response = await client.post(
  343. webhook_url,
  344. data={"payload_json": json.dumps(payload)},
  345. files={"files[0]": ("photo.jpg", image_data, "image/jpeg")},
  346. )
  347. else:
  348. response = await client.post(webhook_url, json={"embeds": [embed]})
  349. if response.status_code in (200, 204):
  350. return True, "Message sent successfully"
  351. else:
  352. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  353. async def _send_webhook(self, config: dict, title: str, message: str) -> tuple[bool, str]:
  354. """Send notification via generic webhook (POST JSON).
  355. Supports two payload formats:
  356. - generic: Custom field names with timestamp/source metadata
  357. - slack: Slack/Mattermost compatible format (just {"text": "..."})
  358. """
  359. webhook_url = config.get("webhook_url", "").strip()
  360. auth_header = config.get("auth_header", "").strip()
  361. payload_format = config.get("payload_format", "generic").strip()
  362. if not webhook_url:
  363. return False, "Webhook URL is required"
  364. # Build payload based on format
  365. if payload_format == "slack":
  366. # Slack/Mattermost format - just text field
  367. data = {"text": f"*{title}*\n{message}"}
  368. else:
  369. # Generic format with custom field names
  370. custom_field_title = config.get("field_title", "title").strip() or "title"
  371. custom_field_message = config.get("field_message", "message").strip() or "message"
  372. data = {
  373. custom_field_title: title,
  374. custom_field_message: message,
  375. "timestamp": datetime.now().isoformat(),
  376. "source": "Bambuddy",
  377. }
  378. headers = {"Content-Type": "application/json"}
  379. if auth_header:
  380. # Support "Bearer token" or just "token" format
  381. if " " in auth_header:
  382. headers["Authorization"] = auth_header
  383. else:
  384. headers["Authorization"] = f"Bearer {auth_header}"
  385. client = await self._get_client()
  386. try:
  387. response = await client.post(webhook_url, json=data, headers=headers)
  388. if response.status_code in (200, 201, 202, 204):
  389. return True, "Webhook delivered successfully"
  390. else:
  391. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  392. except Exception as e:
  393. return False, f"Webhook error: {str(e)}"
  394. async def _send_homeassistant(
  395. self, config: dict, title: str, message: str, db: AsyncSession | None = None
  396. ) -> tuple[bool, str]:
  397. """Send notification via Home Assistant persistent notifications.
  398. Uses the globally configured HA URL/token from settings,
  399. and calls POST /api/services/persistent_notification/create.
  400. """
  401. # Get HA connection settings from global config
  402. ha_url = ""
  403. ha_token = ""
  404. if db:
  405. from backend.app.api.routes.settings import get_homeassistant_settings
  406. try:
  407. ha_settings = await get_homeassistant_settings(db)
  408. ha_url = ha_settings.get("ha_url", "")
  409. ha_token = ha_settings.get("ha_token", "")
  410. except Exception as e:
  411. logger.warning("Failed to read HA settings from database: %s", e)
  412. else:
  413. # Fallback: read directly from environment if no DB session
  414. import os
  415. ha_url = os.environ.get("HA_URL", "")
  416. ha_token = os.environ.get("HA_TOKEN", "")
  417. if not ha_url or not ha_token:
  418. return False, (
  419. "Home Assistant is not configured. Please set HA URL and token in Settings → Network → Home Assistant."
  420. )
  421. url = f"{ha_url.rstrip('/')}/api/services/persistent_notification/create"
  422. headers = {
  423. "Authorization": f"Bearer {ha_token}",
  424. "Content-Type": "application/json",
  425. }
  426. payload = {
  427. "title": title,
  428. "message": message,
  429. }
  430. client = await self._get_client()
  431. response = await client.post(url, json=payload, headers=headers)
  432. if response.status_code in (200, 201):
  433. return True, "Notification sent via Home Assistant"
  434. elif response.status_code == 401:
  435. return False, "Home Assistant authentication failed - check your token"
  436. else:
  437. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  438. async def _send_to_provider(
  439. self,
  440. provider: NotificationProvider,
  441. title: str,
  442. message: str,
  443. db: AsyncSession | None = None,
  444. image_data: bytes | None = None,
  445. ) -> tuple[bool, str]:
  446. """Send notification to a specific provider."""
  447. # Check quiet hours
  448. if self._is_in_quiet_hours(provider):
  449. logger.info("Skipping notification to %s - quiet hours active", provider.name)
  450. return True, "Skipped - quiet hours"
  451. config = json.loads(provider.config) if isinstance(provider.config, str) else provider.config
  452. try:
  453. if provider.provider_type == "callmebot":
  454. return await self._send_callmebot(config, f"{title}\n{message}")
  455. elif provider.provider_type == "ntfy":
  456. return await self._send_ntfy(config, title, message, image_data=image_data)
  457. elif provider.provider_type == "pushover":
  458. return await self._send_pushover(config, title, message, image_data=image_data)
  459. elif provider.provider_type == "telegram":
  460. return await self._send_telegram(config, f"*{title}*\n{message}", image_data=image_data)
  461. elif provider.provider_type == "email":
  462. return await self._send_email(config, title, message)
  463. elif provider.provider_type == "discord":
  464. return await self._send_discord(config, title, message, image_data=image_data)
  465. elif provider.provider_type == "webhook":
  466. return await self._send_webhook(config, title, message)
  467. elif provider.provider_type == "homeassistant":
  468. return await self._send_homeassistant(config, title, message, db=db)
  469. else:
  470. return False, f"Unknown provider type: {provider.provider_type}"
  471. except Exception as e:
  472. logger.exception("Error sending notification via %s", provider.provider_type)
  473. return False, str(e)
  474. async def _update_provider_status(
  475. self, db: AsyncSession, provider_id: int, success: bool, error: str | None = None
  476. ):
  477. """Update provider status after sending notification."""
  478. result = await db.execute(select(NotificationProvider).where(NotificationProvider.id == provider_id))
  479. provider = result.scalar_one_or_none()
  480. if provider:
  481. if success:
  482. provider.last_success = datetime.now(timezone.utc)
  483. else:
  484. provider.last_error = error
  485. provider.last_error_at = datetime.now(timezone.utc)
  486. await db.commit()
  487. async def _get_providers_for_event(
  488. self,
  489. db: AsyncSession,
  490. event_field: str,
  491. printer_id: int | None = None,
  492. ) -> list[NotificationProvider]:
  493. """Get all enabled providers that want a specific event type."""
  494. # Build the query dynamically based on event field
  495. query = select(NotificationProvider).where(
  496. NotificationProvider.enabled.is_(True),
  497. getattr(NotificationProvider, event_field).is_(True),
  498. )
  499. if printer_id is not None:
  500. query = query.where(
  501. (NotificationProvider.printer_id.is_(None)) | (NotificationProvider.printer_id == printer_id)
  502. )
  503. result = await db.execute(query)
  504. return list(result.scalars().all())
  505. async def _log_notification(
  506. self,
  507. db: AsyncSession,
  508. provider_id: int,
  509. event_type: str,
  510. title: str,
  511. message: str,
  512. success: bool,
  513. error_message: str | None = None,
  514. printer_id: int | None = None,
  515. printer_name: str | None = None,
  516. ):
  517. """Create a log entry for a sent notification."""
  518. try:
  519. log = NotificationLog(
  520. provider_id=provider_id,
  521. event_type=event_type,
  522. title=title,
  523. message=message,
  524. success=success,
  525. error_message=error_message,
  526. printer_id=printer_id,
  527. printer_name=printer_name,
  528. )
  529. db.add(log)
  530. await db.commit()
  531. except Exception as e:
  532. logger.warning("Failed to log notification: %s", e)
  533. # Don't fail the notification just because logging failed
  534. async def _send_to_providers(
  535. self,
  536. providers: list[NotificationProvider],
  537. title: str,
  538. message: str,
  539. db: AsyncSession,
  540. event_type: str = "unknown",
  541. printer_id: int | None = None,
  542. printer_name: str | None = None,
  543. force_immediate: bool = False,
  544. image_data: bytes | None = None,
  545. ):
  546. """Send notification to multiple providers and log the results.
  547. All notifications are always sent immediately. If digest mode is enabled,
  548. the notification is ALSO queued for the daily digest summary.
  549. """
  550. for provider in providers:
  551. try:
  552. # Always send notification immediately
  553. success, error = await self._send_to_provider(provider, title, message, db, image_data=image_data)
  554. # Also queue for digest if enabled (digest is a summary, not a queue)
  555. if provider.daily_digest_enabled and provider.daily_digest_time:
  556. await self._queue_for_digest(
  557. provider=provider,
  558. event_type=event_type,
  559. title=title,
  560. message=message,
  561. db=db,
  562. printer_id=printer_id,
  563. printer_name=printer_name,
  564. )
  565. await self._update_provider_status(db, provider.id, success, error if not success else None)
  566. await self._log_notification(
  567. db=db,
  568. provider_id=provider.id,
  569. event_type=event_type,
  570. title=title,
  571. message=message,
  572. success=success,
  573. error_message=error if not success else None,
  574. printer_id=printer_id,
  575. printer_name=printer_name,
  576. )
  577. if success:
  578. logger.info("Sent notification via %s", provider.name)
  579. else:
  580. logger.warning("Failed to send notification via %s: %s", provider.name, error)
  581. except Exception as e:
  582. logger.exception("Error sending notification via %s", provider.name)
  583. await self._update_provider_status(db, provider.id, False, str(e))
  584. await self._log_notification(
  585. db=db,
  586. provider_id=provider.id,
  587. event_type=event_type,
  588. title=title,
  589. message=message,
  590. success=False,
  591. error_message=str(e),
  592. printer_id=printer_id,
  593. printer_name=printer_name,
  594. )
  595. async def on_print_start(
  596. self,
  597. printer_id: int,
  598. printer_name: str,
  599. data: dict,
  600. db: AsyncSession,
  601. archive_data: dict | None = None,
  602. ):
  603. """Handle print start event - send notifications to relevant providers.
  604. Args:
  605. printer_id: The printer ID
  606. printer_name: The printer name
  607. data: MQTT event data with filename, subtask_name, remaining_time, raw_data
  608. db: Database session
  609. archive_data: Optional archive data with print_time_seconds from 3MF parsing
  610. """
  611. logger.info("on_print_start called for printer %s (%s)", printer_id, printer_name)
  612. providers = await self._get_providers_for_event(db, "on_print_start", printer_id)
  613. if not providers:
  614. logger.info("No notification providers configured for print_start event on printer %s", printer_id)
  615. return
  616. # Use subtask_name (project name) if available, otherwise use filename
  617. subtask_name = data.get("subtask_name")
  618. if subtask_name:
  619. # Replace underscores with spaces for readability
  620. filename = subtask_name.replace("_", " ")
  621. else:
  622. filename = self._clean_filename(data.get("filename", "Unknown"))
  623. # Priority for estimated_time:
  624. # 1. Archive's print_time_seconds from 3MF parsing (most reliable)
  625. # 2. MQTT remaining_time (may be 0 at print start)
  626. # 3. raw_data mc_remaining_time
  627. estimated_time = None
  628. # Try archive data first (from 3MF parsing - most reliable)
  629. if archive_data and archive_data.get("print_time_seconds"):
  630. estimated_time = archive_data["print_time_seconds"]
  631. logger.debug("Using print_time_seconds from archive: %s", estimated_time)
  632. # Fall back to MQTT remaining_time
  633. if estimated_time is None:
  634. estimated_time = data.get("remaining_time")
  635. if estimated_time:
  636. logger.debug("Using remaining_time from MQTT: %s", estimated_time)
  637. # Last resort: raw_data mc_remaining_time (in minutes, convert to seconds)
  638. if estimated_time is None:
  639. raw_time = data.get("raw_data", {}).get("mc_remaining_time")
  640. if raw_time:
  641. estimated_time = raw_time * 60
  642. logger.debug("Using mc_remaining_time from raw_data: %s", estimated_time)
  643. time_str = self._format_duration(estimated_time)
  644. eta_str = await self._format_eta(estimated_time, db)
  645. variables = {
  646. "printer": printer_name,
  647. "filename": filename,
  648. "estimated_time": time_str,
  649. "eta": eta_str,
  650. }
  651. # Extract image data for providers that support attachments (e.g. Pushover)
  652. image_data = None
  653. if archive_data:
  654. image_data = archive_data.get("image_data")
  655. logger.info("Found %s providers for print_start: %s", len(providers), [p.name for p in providers])
  656. title, message = await self._build_message_from_template(db, "print_start", variables)
  657. await self._send_to_providers(
  658. providers, title, message, db, "print_start", printer_id, printer_name, image_data=image_data
  659. )
  660. async def on_print_complete(
  661. self,
  662. printer_id: int,
  663. printer_name: str,
  664. status: str,
  665. data: dict,
  666. db: AsyncSession,
  667. archive_data: dict | None = None,
  668. ):
  669. """Handle print complete event - send notifications to relevant providers."""
  670. logger.info("on_print_complete called for printer %s (%s), status=%s", printer_id, printer_name, status)
  671. # Determine event type based on status
  672. if status == "completed":
  673. event_field = "on_print_complete"
  674. event_type = "print_complete"
  675. elif status in ("failed",):
  676. event_field = "on_print_failed"
  677. event_type = "print_failed"
  678. elif status in ("aborted", "stopped", "cancelled"):
  679. event_field = "on_print_stopped"
  680. event_type = "print_stopped"
  681. else:
  682. logger.warning("Unknown print status '%s', defaulting to on_print_complete", status)
  683. event_field = "on_print_complete"
  684. event_type = "print_complete"
  685. providers = await self._get_providers_for_event(db, event_field, printer_id)
  686. if not providers:
  687. logger.info("No notification providers configured for %s event on printer %s", event_field, printer_id)
  688. return
  689. # Use subtask_name (project name) if available, otherwise use filename
  690. subtask_name = data.get("subtask_name")
  691. if subtask_name:
  692. filename = subtask_name.replace("_", " ")
  693. else:
  694. filename = self._clean_filename(data.get("filename", "Unknown"))
  695. variables = {
  696. "printer": printer_name,
  697. "filename": filename,
  698. "duration": "Unknown",
  699. "filament_grams": "Unknown",
  700. "reason": "Unknown",
  701. }
  702. if archive_data:
  703. if archive_data.get("print_time_seconds"):
  704. variables["duration"] = self._format_duration(archive_data["print_time_seconds"])
  705. if archive_data.get("actual_filament_grams"):
  706. variables["filament_grams"] = f"{archive_data['actual_filament_grams']:.1f}"
  707. if status == "failed" and archive_data.get("failure_reason"):
  708. variables["reason"] = archive_data["failure_reason"]
  709. if archive_data.get("finish_photo_url"):
  710. variables["finish_photo_url"] = archive_data["finish_photo_url"]
  711. # Build per-slot breakdown string with AMS info when available
  712. if archive_data.get("usage_results"):
  713. parts = []
  714. for u in archive_data["usage_results"]:
  715. ams_id = u.get("ams_id", 0)
  716. tray_id = u.get("tray_id", 0)
  717. material = u.get("material", "Unknown") or "Unknown"
  718. used = u.get("weight_used", 0)
  719. if ams_id >= 128:
  720. slot_label = "Ext"
  721. else:
  722. slot_label = f"AMS-{chr(65 + ams_id)} T{tray_id + 1}"
  723. parts.append(f"{slot_label} {material}: {used:.1f}g")
  724. variables["filament_details"] = " | ".join(parts)
  725. elif archive_data.get("filament_slots"):
  726. parts = []
  727. for slot in archive_data["filament_slots"]:
  728. ftype = slot.get("type", "Unknown") or "Unknown"
  729. used = slot.get("used_g", 0)
  730. parts.append(f"{ftype}: {used:.1f}g")
  731. variables["filament_details"] = " | ".join(parts)
  732. # Add progress for partial prints
  733. if archive_data.get("progress") is not None:
  734. variables["progress"] = str(archive_data["progress"])
  735. # Extract image data for providers that support attachments (e.g. Pushover)
  736. image_data = None
  737. if archive_data:
  738. image_data = archive_data.get("image_data")
  739. logger.info("Found %s providers for %s: %s", len(providers), event_field, [p.name for p in providers])
  740. title, message = await self._build_message_from_template(db, event_type, variables)
  741. await self._send_to_providers(
  742. providers, title, message, db, event_type, printer_id, printer_name, image_data=image_data
  743. )
  744. async def on_print_progress(
  745. self,
  746. printer_id: int,
  747. printer_name: str,
  748. filename: str,
  749. progress: int,
  750. db: AsyncSession,
  751. remaining_time: int | None = None,
  752. image_data: bytes | None = None,
  753. ):
  754. """Handle print progress milestone (25%, 50%, 75%)."""
  755. providers = await self._get_providers_for_event(db, "on_print_progress", printer_id)
  756. if not providers:
  757. return
  758. eta_str = await self._format_eta(remaining_time, db)
  759. variables = {
  760. "printer": printer_name,
  761. "filename": self._clean_filename(filename),
  762. "progress": str(progress),
  763. "remaining_time": self._format_duration(remaining_time) if remaining_time else "Unknown",
  764. "eta": eta_str,
  765. }
  766. title, message = await self._build_message_from_template(db, "print_progress", variables)
  767. await self._send_to_providers(
  768. providers, title, message, db, "print_progress", printer_id, printer_name, image_data=image_data
  769. )
  770. async def on_printer_offline(self, printer_id: int, printer_name: str, db: AsyncSession):
  771. """Handle printer offline event."""
  772. providers = await self._get_providers_for_event(db, "on_printer_offline", printer_id)
  773. if not providers:
  774. return
  775. variables = {"printer": printer_name}
  776. title, message = await self._build_message_from_template(db, "printer_offline", variables)
  777. await self._send_to_providers(providers, title, message, db, "printer_offline", printer_id, printer_name)
  778. async def on_printer_error(
  779. self,
  780. printer_id: int,
  781. printer_name: str,
  782. error_type: str,
  783. db: AsyncSession,
  784. error_detail: str | None = None,
  785. image_data: bytes | None = None,
  786. ):
  787. """Handle printer error event (AMS issues, etc.)."""
  788. providers = await self._get_providers_for_event(db, "on_printer_error", printer_id)
  789. if not providers:
  790. return
  791. variables = {
  792. "printer": printer_name,
  793. "error_type": error_type,
  794. "error_detail": error_detail or "No details available",
  795. }
  796. title, message = await self._build_message_from_template(db, "printer_error", variables)
  797. await self._send_to_providers(
  798. providers, title, message, db, "printer_error", printer_id, printer_name, image_data=image_data
  799. )
  800. async def on_plate_not_empty(
  801. self,
  802. printer_id: int,
  803. printer_name: str,
  804. db: AsyncSession,
  805. difference_percent: float | None = None,
  806. ):
  807. """Handle plate not empty event - objects detected on build plate before print."""
  808. providers = await self._get_providers_for_event(db, "on_plate_not_empty", printer_id)
  809. if not providers:
  810. return
  811. variables = {
  812. "printer": printer_name,
  813. "difference_percent": f"{difference_percent:.1f}" if difference_percent else "N/A",
  814. }
  815. title, message = await self._build_message_from_template(db, "plate_not_empty", variables)
  816. await self._send_to_providers(
  817. providers, title, message, db, "plate_not_empty", printer_id, printer_name, force_immediate=True
  818. )
  819. async def on_filament_low(
  820. self,
  821. printer_id: int,
  822. printer_name: str,
  823. slot: int,
  824. remaining_percent: int,
  825. db: AsyncSession,
  826. color: str | None = None,
  827. ):
  828. """Handle low filament event."""
  829. providers = await self._get_providers_for_event(db, "on_filament_low", printer_id)
  830. if not providers:
  831. return
  832. variables = {
  833. "printer": printer_name,
  834. "slot": str(slot),
  835. "remaining_percent": str(remaining_percent),
  836. "color": color or "",
  837. }
  838. title, message = await self._build_message_from_template(db, "filament_low", variables)
  839. await self._send_to_providers(providers, title, message, db, "filament_low", printer_id, printer_name)
  840. async def on_maintenance_due(
  841. self,
  842. printer_id: int,
  843. printer_name: str,
  844. maintenance_items: list[dict],
  845. db: AsyncSession,
  846. ):
  847. """Handle maintenance due event - sends notification when maintenance is due or warning."""
  848. if not maintenance_items:
  849. return
  850. providers = await self._get_providers_for_event(db, "on_maintenance_due", printer_id)
  851. if not providers:
  852. logger.info("No notification providers configured for maintenance_due event on printer %s", printer_id)
  853. return
  854. # Format maintenance items list
  855. items_list = []
  856. for item in maintenance_items:
  857. status = "OVERDUE" if item.get("is_due") else "Soon"
  858. items_list.append(f"- {item['name']} ({status})")
  859. items_str = "\n".join(items_list)
  860. variables = {
  861. "printer": printer_name,
  862. "items": items_str,
  863. }
  864. logger.info("Found %s providers for maintenance_due: %s", len(providers), [p.name for p in providers])
  865. title, message = await self._build_message_from_template(db, "maintenance_due", variables)
  866. await self._send_to_providers(providers, title, message, db, "maintenance_due", printer_id, printer_name)
  867. async def on_ams_humidity_high(
  868. self,
  869. printer_id: int,
  870. printer_name: str,
  871. ams_label: str,
  872. humidity: float,
  873. threshold: float,
  874. db: AsyncSession,
  875. ):
  876. """Handle AMS high humidity alarm event. Always sends immediately (bypasses digest)."""
  877. providers = await self._get_providers_for_event(db, "on_ams_humidity_high", printer_id)
  878. if not providers:
  879. return
  880. variables = {
  881. "printer": printer_name,
  882. "ams_label": ams_label,
  883. "humidity": f"{humidity:.0f}",
  884. "threshold": f"{threshold:.0f}",
  885. }
  886. title, message = await self._build_message_from_template(db, "ams_humidity_high", variables)
  887. # Alarms always send immediately, bypassing digest mode
  888. await self._send_to_providers(
  889. providers, title, message, db, "ams_humidity_high", printer_id, printer_name, force_immediate=True
  890. )
  891. async def on_ams_temperature_high(
  892. self,
  893. printer_id: int,
  894. printer_name: str,
  895. ams_label: str,
  896. temperature: float,
  897. threshold: float,
  898. db: AsyncSession,
  899. ):
  900. """Handle AMS high temperature alarm event. Always sends immediately (bypasses digest)."""
  901. providers = await self._get_providers_for_event(db, "on_ams_temperature_high", printer_id)
  902. if not providers:
  903. return
  904. variables = {
  905. "printer": printer_name,
  906. "ams_label": ams_label,
  907. "temperature": f"{temperature:.1f}",
  908. "threshold": f"{threshold:.1f}",
  909. }
  910. title, message = await self._build_message_from_template(db, "ams_temperature_high", variables)
  911. # Alarms always send immediately, bypassing digest mode
  912. await self._send_to_providers(
  913. providers, title, message, db, "ams_temperature_high", printer_id, printer_name, force_immediate=True
  914. )
  915. async def on_ams_ht_humidity_high(
  916. self,
  917. printer_id: int,
  918. printer_name: str,
  919. ams_label: str,
  920. humidity: float,
  921. threshold: float,
  922. db: AsyncSession,
  923. ):
  924. """Handle AMS-HT high humidity alarm event. Always sends immediately (bypasses digest)."""
  925. providers = await self._get_providers_for_event(db, "on_ams_ht_humidity_high", printer_id)
  926. if not providers:
  927. return
  928. variables = {
  929. "printer": printer_name,
  930. "ams_label": ams_label,
  931. "humidity": f"{humidity:.0f}",
  932. "threshold": f"{threshold:.0f}",
  933. }
  934. # Use the same template as regular AMS (can create separate templates later if needed)
  935. title, message = await self._build_message_from_template(db, "ams_humidity_high", variables)
  936. # Alarms always send immediately, bypassing digest mode
  937. await self._send_to_providers(
  938. providers, title, message, db, "ams_ht_humidity_high", printer_id, printer_name, force_immediate=True
  939. )
  940. async def on_ams_ht_temperature_high(
  941. self,
  942. printer_id: int,
  943. printer_name: str,
  944. ams_label: str,
  945. temperature: float,
  946. threshold: float,
  947. db: AsyncSession,
  948. ):
  949. """Handle AMS-HT high temperature alarm event. Always sends immediately (bypasses digest)."""
  950. providers = await self._get_providers_for_event(db, "on_ams_ht_temperature_high", printer_id)
  951. if not providers:
  952. return
  953. variables = {
  954. "printer": printer_name,
  955. "ams_label": ams_label,
  956. "temperature": f"{temperature:.1f}",
  957. "threshold": f"{threshold:.1f}",
  958. }
  959. # Use the same template as regular AMS (can create separate templates later if needed)
  960. title, message = await self._build_message_from_template(db, "ams_temperature_high", variables)
  961. # Alarms always send immediately, bypassing digest mode
  962. await self._send_to_providers(
  963. providers, title, message, db, "ams_ht_temperature_high", printer_id, printer_name, force_immediate=True
  964. )
  965. async def on_bed_cooled(
  966. self,
  967. printer_id: int,
  968. printer_name: str,
  969. bed_temp: float,
  970. threshold: float,
  971. filename: str,
  972. db: AsyncSession,
  973. ):
  974. """Handle bed cooled event - bed temperature dropped below threshold after print."""
  975. providers = await self._get_providers_for_event(db, "on_bed_cooled", printer_id)
  976. if not providers:
  977. return
  978. variables = {
  979. "printer": printer_name,
  980. "bed_temp": f"{bed_temp:.0f}",
  981. "threshold": f"{threshold:.0f}",
  982. "filename": self._clean_filename(filename) if filename else "Unknown",
  983. }
  984. title, message = await self._build_message_from_template(db, "bed_cooled", variables)
  985. await self._send_to_providers(providers, title, message, db, "bed_cooled", printer_id, printer_name)
  986. async def on_first_layer_complete(
  987. self,
  988. printer_id: int,
  989. printer_name: str,
  990. filename: str,
  991. total_layers: int,
  992. db: AsyncSession,
  993. image_data: bytes | None = None,
  994. ):
  995. """Handle first layer complete event."""
  996. providers = await self._get_providers_for_event(db, "on_first_layer_complete", printer_id)
  997. if not providers:
  998. return
  999. variables = {
  1000. "printer": printer_name,
  1001. "filename": self._clean_filename(filename),
  1002. "total_layers": str(total_layers),
  1003. }
  1004. title, message = await self._build_message_from_template(db, "first_layer_complete", variables)
  1005. await self._send_to_providers(
  1006. providers, title, message, db, "first_layer_complete", printer_id, printer_name, image_data=image_data
  1007. )
  1008. def clear_template_cache(self):
  1009. """Clear the template cache. Call this when templates are updated."""
  1010. self._template_cache.clear()
  1011. # ==================== Queue Notifications ====================
  1012. async def on_queue_job_added(
  1013. self,
  1014. job_name: str,
  1015. target: str,
  1016. db: AsyncSession,
  1017. printer_id: int | None = None,
  1018. printer_name: str | None = None,
  1019. ):
  1020. """Handle queue job added event."""
  1021. providers = await self._get_providers_for_event(db, "on_queue_job_added", printer_id)
  1022. if not providers:
  1023. return
  1024. variables = {
  1025. "job_name": job_name,
  1026. "target": target, # e.g., "Printer1" or "Any X1C"
  1027. "printer": printer_name or target,
  1028. }
  1029. title, message = await self._build_message_from_template(db, "queue_job_added", variables)
  1030. await self._send_to_providers(providers, title, message, db, "queue_job_added", printer_id, printer_name)
  1031. async def on_queue_job_assigned(
  1032. self,
  1033. job_name: str,
  1034. printer_id: int,
  1035. printer_name: str,
  1036. target_model: str,
  1037. db: AsyncSession,
  1038. ):
  1039. """Handle model-based job assigned to printer event."""
  1040. providers = await self._get_providers_for_event(db, "on_queue_job_assigned", printer_id)
  1041. if not providers:
  1042. return
  1043. variables = {
  1044. "job_name": job_name,
  1045. "printer": printer_name,
  1046. "target_model": target_model,
  1047. }
  1048. title, message = await self._build_message_from_template(db, "queue_job_assigned", variables)
  1049. await self._send_to_providers(providers, title, message, db, "queue_job_assigned", printer_id, printer_name)
  1050. async def on_queue_job_started(
  1051. self,
  1052. job_name: str,
  1053. printer_id: int,
  1054. printer_name: str,
  1055. db: AsyncSession,
  1056. estimated_time: int | None = None,
  1057. ):
  1058. """Handle queue job started printing event."""
  1059. providers = await self._get_providers_for_event(db, "on_queue_job_started", printer_id)
  1060. if not providers:
  1061. return
  1062. eta_str = await self._format_eta(estimated_time, db)
  1063. variables = {
  1064. "job_name": job_name,
  1065. "printer": printer_name,
  1066. "estimated_time": self._format_duration(estimated_time),
  1067. "eta": eta_str,
  1068. }
  1069. title, message = await self._build_message_from_template(db, "queue_job_started", variables)
  1070. await self._send_to_providers(providers, title, message, db, "queue_job_started", printer_id, printer_name)
  1071. async def on_queue_job_waiting(
  1072. self,
  1073. job_name: str,
  1074. target_model: str,
  1075. waiting_reason: str,
  1076. db: AsyncSession,
  1077. ):
  1078. """Handle job waiting for filament event."""
  1079. providers = await self._get_providers_for_event(db, "on_queue_job_waiting", None)
  1080. if not providers:
  1081. return
  1082. variables = {
  1083. "job_name": job_name,
  1084. "target_model": target_model,
  1085. "waiting_reason": waiting_reason,
  1086. }
  1087. title, message = await self._build_message_from_template(db, "queue_job_waiting", variables)
  1088. await self._send_to_providers(providers, title, message, db, "queue_job_waiting")
  1089. async def on_queue_job_skipped(
  1090. self,
  1091. job_name: str,
  1092. printer_id: int,
  1093. printer_name: str,
  1094. reason: str,
  1095. db: AsyncSession,
  1096. ):
  1097. """Handle job skipped event (e.g., previous print failed)."""
  1098. providers = await self._get_providers_for_event(db, "on_queue_job_skipped", printer_id)
  1099. if not providers:
  1100. return
  1101. variables = {
  1102. "job_name": job_name,
  1103. "printer": printer_name,
  1104. "reason": reason,
  1105. }
  1106. title, message = await self._build_message_from_template(db, "queue_job_skipped", variables)
  1107. await self._send_to_providers(providers, title, message, db, "queue_job_skipped", printer_id, printer_name)
  1108. async def on_queue_job_failed(
  1109. self,
  1110. job_name: str,
  1111. printer_id: int | None,
  1112. printer_name: str | None,
  1113. reason: str,
  1114. db: AsyncSession,
  1115. ):
  1116. """Handle job failed to start event (upload error, etc.)."""
  1117. providers = await self._get_providers_for_event(db, "on_queue_job_failed", printer_id)
  1118. if not providers:
  1119. return
  1120. variables = {
  1121. "job_name": job_name,
  1122. "printer": printer_name or "Unknown",
  1123. "reason": reason,
  1124. }
  1125. title, message = await self._build_message_from_template(db, "queue_job_failed", variables)
  1126. await self._send_to_providers(providers, title, message, db, "queue_job_failed", printer_id, printer_name)
  1127. async def on_queue_completed(
  1128. self,
  1129. completed_count: int,
  1130. db: AsyncSession,
  1131. ):
  1132. """Handle all queue jobs completed event."""
  1133. providers = await self._get_providers_for_event(db, "on_queue_completed", None)
  1134. if not providers:
  1135. return
  1136. variables = {
  1137. "completed_count": str(completed_count),
  1138. }
  1139. title, message = await self._build_message_from_template(db, "queue_completed", variables)
  1140. await self._send_to_providers(providers, title, message, db, "queue_completed")
  1141. async def _queue_for_digest(
  1142. self,
  1143. provider: NotificationProvider,
  1144. event_type: str,
  1145. title: str,
  1146. message: str,
  1147. db: AsyncSession,
  1148. printer_id: int | None = None,
  1149. printer_name: str | None = None,
  1150. ):
  1151. """Queue a notification for later delivery in the daily digest."""
  1152. try:
  1153. queue_entry = NotificationDigestQueue(
  1154. provider_id=provider.id,
  1155. event_type=event_type,
  1156. title=title,
  1157. message=message,
  1158. printer_id=printer_id,
  1159. printer_name=printer_name,
  1160. )
  1161. db.add(queue_entry)
  1162. await db.commit()
  1163. logger.info("Queued notification for digest: %s for provider %s", event_type, provider.name)
  1164. except Exception as e:
  1165. logger.warning("Failed to queue notification for digest: %s", e)
  1166. async def send_digest(self, provider_id: int):
  1167. """Send all queued notifications as a single digest for a provider."""
  1168. from backend.app.core.database import async_session
  1169. async with async_session() as db:
  1170. # Get the provider
  1171. result = await db.execute(select(NotificationProvider).where(NotificationProvider.id == provider_id))
  1172. provider = result.scalar_one_or_none()
  1173. if not provider or not provider.enabled:
  1174. return
  1175. # Get all queued notifications for this provider
  1176. result = await db.execute(
  1177. select(NotificationDigestQueue)
  1178. .where(NotificationDigestQueue.provider_id == provider_id)
  1179. .order_by(NotificationDigestQueue.created_at)
  1180. )
  1181. queue_entries = list(result.scalars().all())
  1182. if not queue_entries:
  1183. logger.debug("No queued notifications for provider %s", provider.name)
  1184. return
  1185. # Build digest message
  1186. title = f"Daily Digest - {len(queue_entries)} Events"
  1187. # Group by event type
  1188. events_by_type: dict[str, list] = {}
  1189. for entry in queue_entries:
  1190. if entry.event_type not in events_by_type:
  1191. events_by_type[entry.event_type] = []
  1192. events_by_type[entry.event_type].append(entry)
  1193. # Format the digest body
  1194. body_parts = []
  1195. for event_type, entries in events_by_type.items():
  1196. event_label = event_type.replace("_", " ").title()
  1197. body_parts.append(f"== {event_label} ({len(entries)}) ==")
  1198. for entry in entries:
  1199. time_str = entry.created_at.strftime("%H:%M")
  1200. printer_info = f"[{entry.printer_name}] " if entry.printer_name else ""
  1201. body_parts.append(f" {time_str} {printer_info}{entry.title}")
  1202. body_parts.append("")
  1203. body = "\n".join(body_parts)
  1204. # Send the digest
  1205. success, error = await self._send_to_provider(provider, title, body, db)
  1206. # Log the digest
  1207. await self._log_notification(
  1208. db=db,
  1209. provider_id=provider.id,
  1210. event_type="daily_digest",
  1211. title=title,
  1212. message=body,
  1213. success=success,
  1214. error_message=error if not success else None,
  1215. )
  1216. # Clear the queue
  1217. for entry in queue_entries:
  1218. await db.delete(entry)
  1219. await db.commit()
  1220. if success:
  1221. logger.info("Sent daily digest with %s events to %s", len(queue_entries), provider.name)
  1222. else:
  1223. logger.warning("Failed to send daily digest to %s: %s", provider.name, error)
  1224. async def check_and_send_digests(self):
  1225. """Check all providers and send digests if it's their scheduled time."""
  1226. from backend.app.core.database import async_session
  1227. current_time = datetime.now().strftime("%H:%M")
  1228. # Avoid duplicate checks within the same minute
  1229. if current_time == self._last_digest_check:
  1230. return
  1231. self._last_digest_check = current_time
  1232. async with async_session() as db:
  1233. # Find all providers with digest enabled at this time
  1234. result = await db.execute(
  1235. select(NotificationProvider).where(
  1236. NotificationProvider.enabled.is_(True),
  1237. NotificationProvider.daily_digest_enabled.is_(True),
  1238. NotificationProvider.daily_digest_time == current_time,
  1239. )
  1240. )
  1241. providers = result.scalars().all()
  1242. for provider in providers:
  1243. try:
  1244. await self.send_digest(provider.id)
  1245. except Exception as e:
  1246. logger.error("Error sending digest for provider %s: %s", provider.id, e)
  1247. def start_digest_scheduler(self):
  1248. """Start the background scheduler for daily digest notifications."""
  1249. if self._digest_scheduler_task is None:
  1250. self._digest_scheduler_task = asyncio.create_task(self._digest_scheduler_loop())
  1251. logger.info("Notification digest scheduler started")
  1252. def stop_digest_scheduler(self):
  1253. """Stop the background scheduler for daily digests."""
  1254. if self._digest_scheduler_task:
  1255. self._digest_scheduler_task.cancel()
  1256. self._digest_scheduler_task = None
  1257. logger.info("Notification digest scheduler stopped")
  1258. async def _digest_scheduler_loop(self):
  1259. """Background loop that checks for scheduled digests every minute."""
  1260. while True:
  1261. try:
  1262. await self.check_and_send_digests()
  1263. except Exception as e:
  1264. logger.error("Error in digest scheduler: %s", e)
  1265. # Wait until the next minute
  1266. await asyncio.sleep(60)
  1267. # Global instance
  1268. notification_service = NotificationService()