notification_service.py 81 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067
  1. """Notification service for sending push notifications via various providers."""
  2. import asyncio
  3. import html
  4. import json
  5. import logging
  6. import re
  7. import smtplib
  8. from datetime import datetime, timedelta, timezone
  9. from email.mime.image import MIMEImage
  10. from email.mime.multipart import MIMEMultipart
  11. from email.mime.text import MIMEText
  12. from typing import Any
  13. from urllib.parse import quote
  14. import httpx
  15. from sqlalchemy import select
  16. from sqlalchemy.ext.asyncio import AsyncSession
  17. from backend.app.models.notification import NotificationDigestQueue, NotificationLog, NotificationProvider
  18. from backend.app.models.notification_template import NotificationTemplate
  19. logger = logging.getLogger(__name__)
  20. # Honest User-Agent — matches the convention used by every other outbound
  21. # httpx client in the codebase (bambu_cloud, makerworld, firmware_check,
  22. # inventory). Previously this client leaked python-httpx/<version>, which
  23. # was both inconsistent with the rest of the project and a more obvious
  24. # bot signature for upstream WAFs.
  25. _USER_AGENT = "Bambuddy/1.0 (+https://github.com/maziggy/bambuddy)"
  26. def _looks_like_cloudflare_challenge(response: httpx.Response) -> bool:
  27. """Return True if ``response`` looks like a Cloudflare mitigation
  28. interstitial (JS challenge / managed challenge / block page) rather
  29. than a legitimate response passed through Cloudflare.
  30. Self-hosted servers behind Cloudflare (Tunnel, "Bot Fight Mode", or
  31. "Under Attack" mode) intercept non-browser clients at the edge and
  32. return a challenge HTML page instead of forwarding to the origin —
  33. so we never reach the user's actual ntfy / webhook backend.
  34. Cloudflare cannot be defeated from a Python client; the user has to
  35. add a security-skip rule on their side. We detect the shape so the
  36. UI can tell them that, instead of dumping the raw HTML.
  37. Detection deliberately does NOT rely on ``Server: cloudflare`` alone
  38. — Cloudflare adds that header to every response it proxies (success
  39. AND legitimate origin errors), so a real 401 "wrong token" from a
  40. CF-fronted ntfy would false-positive into a misleading "your CF is
  41. blocking" message. Reliable signals: the ``cf-mitigated`` header
  42. (set only when CF actively mitigates) and the challenge body shape.
  43. """
  44. if response.headers.get("cf-mitigated"):
  45. return True
  46. content_type = (response.headers.get("content-type") or "").lower()
  47. if "html" not in content_type:
  48. return False
  49. body = (response.text or "")[:1024].lower()
  50. # "Just a moment..." is Cloudflare's universal challenge-page title
  51. # (managed challenge, JS challenge, Under Attack mode). Combined with
  52. # an HTML content-type this is unambiguous — no legitimate ntfy or
  53. # webhook backend returns HTML with that title. ``cf-chl-*`` and
  54. # ``challenge-platform`` cover newer / non-default CF templates.
  55. return "just a moment" in body or "cf-chl-bypass" in body or "cf-chl-opt" in body or "challenge-platform" in body
  56. class NotificationService:
  57. """Service for sending notifications through various providers."""
  58. def __init__(self):
  59. self._http_client: httpx.AsyncClient | None = None
  60. self._template_cache: dict[str, NotificationTemplate] = {}
  61. self._digest_scheduler_task: asyncio.Task | None = None
  62. self._last_digest_check: str = "" # "HH:MM" to avoid duplicate checks
  63. async def _get_client(self) -> httpx.AsyncClient:
  64. """Get or create HTTP client."""
  65. if self._http_client is None or self._http_client.is_closed:
  66. self._http_client = httpx.AsyncClient(
  67. timeout=30.0,
  68. headers={"User-Agent": _USER_AGENT},
  69. )
  70. return self._http_client
  71. async def close(self):
  72. """Close HTTP client."""
  73. if self._http_client and not self._http_client.is_closed:
  74. await self._http_client.aclose()
  75. def _is_in_quiet_hours(self, provider: NotificationProvider) -> bool:
  76. """Check if current time is within provider's quiet hours."""
  77. if not provider.quiet_hours_enabled:
  78. return False
  79. if not provider.quiet_hours_start or not provider.quiet_hours_end:
  80. return False
  81. try:
  82. now = datetime.now()
  83. current_time = now.hour * 60 + now.minute
  84. start_parts = provider.quiet_hours_start.split(":")
  85. end_parts = provider.quiet_hours_end.split(":")
  86. start_minutes = int(start_parts[0]) * 60 + int(start_parts[1])
  87. end_minutes = int(end_parts[0]) * 60 + int(end_parts[1])
  88. # Handle overnight quiet hours (e.g., 22:00 to 07:00)
  89. if start_minutes > end_minutes:
  90. # Quiet hours span midnight
  91. return current_time >= start_minutes or current_time < end_minutes
  92. else:
  93. # Same day quiet hours
  94. return start_minutes <= current_time < end_minutes
  95. except (ValueError, TypeError, AttributeError):
  96. logger.warning("Invalid quiet hours format for provider %s", provider.name)
  97. return False
  98. async def _get_template(self, db: AsyncSession, event_type: str) -> NotificationTemplate | None:
  99. """Get a notification template by event type."""
  100. # Check cache first
  101. if event_type in self._template_cache:
  102. return self._template_cache[event_type]
  103. result = await db.execute(select(NotificationTemplate).where(NotificationTemplate.event_type == event_type))
  104. template = result.scalar_one_or_none()
  105. if template:
  106. self._template_cache[event_type] = template
  107. return template
  108. def _render_template(self, template_str: str, variables: dict[str, Any]) -> str:
  109. """Render a template string with variables. Missing variables become empty."""
  110. result = template_str
  111. for key, value in variables.items():
  112. result = result.replace("{" + key + "}", str(value) if value is not None else "")
  113. # Remove any remaining unreplaced placeholders
  114. result = re.sub(r"\{[a-z_]+\}", "", result)
  115. return result
  116. async def _format_eta(self, seconds: int | None, db: AsyncSession) -> str:
  117. """Format ETA as wall-clock time, respecting user's time_format setting."""
  118. if not seconds or seconds <= 0:
  119. return "Unknown"
  120. from backend.app.api.routes.settings import get_setting
  121. eta_time = datetime.now() + timedelta(seconds=seconds)
  122. time_format = await get_setting(db, "time_format")
  123. if time_format == "12h":
  124. return eta_time.strftime("%I:%M %p").lstrip("0")
  125. # Default to 24h for "24h", "system", or unset
  126. return eta_time.strftime("%H:%M")
  127. def _format_duration(self, seconds: int | None) -> str:
  128. """Format duration in seconds to human-readable string."""
  129. if seconds is None:
  130. return "Unknown"
  131. hours = seconds // 3600
  132. minutes = (seconds % 3600) // 60
  133. if hours > 0:
  134. return f"{hours}h {minutes}m"
  135. return f"{minutes}m"
  136. def _clean_filename(self, filename: str) -> str:
  137. """Extract filename and remove file extensions."""
  138. import os
  139. # Strip path prefix (e.g., /data/Metadata/plate_5.gcode -> plate_5.gcode)
  140. filename = os.path.basename(filename)
  141. # Remove common extensions
  142. if filename.endswith(".gcode.3mf"):
  143. return filename[:-10]
  144. elif filename.endswith(".gcode"):
  145. return filename[:-6]
  146. elif filename.endswith(".3mf"):
  147. return filename[:-4]
  148. return filename
  149. async def _build_message_from_template(
  150. self, db: AsyncSession, event_type: str, variables: dict[str, Any]
  151. ) -> tuple[str, str]:
  152. """Build notification title and body from template."""
  153. # Add common variables
  154. variables["timestamp"] = datetime.now().strftime("%Y-%m-%d %H:%M")
  155. variables["app_name"] = "Bambuddy"
  156. template = await self._get_template(db, event_type)
  157. if not template:
  158. # Fallback to simple message
  159. logger.warning("Template not found for event type: %s", event_type)
  160. return event_type.replace("_", " ").title(), str(variables)
  161. title = self._render_template(template.title_template, variables)
  162. body = self._render_template(template.body_template, variables)
  163. return title, body
  164. async def send_test_notification(
  165. self, provider_type: str, config: dict[str, Any], db: AsyncSession | None = None
  166. ) -> tuple[bool, str]:
  167. """Send a test notification to verify configuration."""
  168. if db:
  169. title, message = await self._build_message_from_template(db, "test", {})
  170. else:
  171. title = "Bambuddy Test"
  172. message = "This is a test notification. If you see this, notifications are working!"
  173. try:
  174. if provider_type == "callmebot":
  175. return await self._send_callmebot(config, f"{title}\n{message}")
  176. elif provider_type == "ntfy":
  177. return await self._send_ntfy(config, title, message)
  178. elif provider_type == "pushover":
  179. return await self._send_pushover(config, title, message)
  180. elif provider_type == "telegram":
  181. return await self._send_telegram(config, f"*{title}*\n{message}")
  182. elif provider_type == "email":
  183. return await self._send_email(config, title, message)
  184. elif provider_type == "discord":
  185. return await self._send_discord(config, title, message)
  186. elif provider_type == "webhook":
  187. return await self._send_webhook(config, title, message)
  188. elif provider_type == "homeassistant":
  189. return await self._send_homeassistant(config, title, message, db=db)
  190. else:
  191. return False, f"Unknown provider type: {provider_type}"
  192. except Exception as e:
  193. logger.exception("Error sending test notification via %s", provider_type)
  194. return False, str(e)
  195. async def _send_callmebot(self, config: dict, message: str) -> tuple[bool, str]:
  196. """Send notification via CallMeBot (WhatsApp)."""
  197. phone = config.get("phone", "").strip()
  198. apikey = config.get("apikey", "").strip()
  199. if not phone or not apikey:
  200. return False, "Phone number and API key are required"
  201. # URL encode the message
  202. encoded_message = quote(message)
  203. url = f"https://api.callmebot.com/whatsapp.php?phone={phone}&text={encoded_message}&apikey={apikey}"
  204. client = await self._get_client()
  205. response = await client.get(url)
  206. if response.status_code == 200:
  207. return True, "Message sent successfully"
  208. else:
  209. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  210. async def _send_ntfy(
  211. self,
  212. config: dict,
  213. title: str,
  214. message: str,
  215. image_data: bytes | None = None,
  216. event_type: str | None = None,
  217. ) -> tuple[bool, str]:
  218. """Send notification via ntfy."""
  219. server = config.get("server", "https://ntfy.sh").rstrip("/")
  220. topic = config.get("topic", "").strip()
  221. auth_token = config.get("auth_token", "").strip()
  222. if not topic:
  223. return False, "Topic is required"
  224. url = f"{server}/{topic}"
  225. # ntfy reads Title/Message from HTTP headers. httpx enforces ASCII
  226. # for str header values, but printer names and filenames can contain
  227. # non-ASCII characters (e.g. accented letters, CJK). Passing bytes
  228. # bypasses the ASCII check — ntfy handles UTF-8 headers correctly.
  229. headers: dict[str, str | bytes] = {"Title": title.encode("utf-8")}
  230. # Per-event Priority header (#990). Only set when the user has
  231. # explicitly mapped this event to a 1-5 value; otherwise fall through
  232. # to the ntfy server's default so existing setups stay unchanged.
  233. event_priorities = config.get("event_priorities") or {}
  234. if event_type and isinstance(event_priorities, dict):
  235. raw = event_priorities.get(event_type)
  236. try:
  237. priority = int(raw) if raw is not None else None
  238. except (TypeError, ValueError):
  239. priority = None
  240. if priority is not None and 1 <= priority <= 5:
  241. headers["Priority"] = str(priority)
  242. if auth_token:
  243. headers["Authorization"] = f"Bearer {auth_token}"
  244. client = await self._get_client()
  245. if image_data:
  246. # ntfy supports image attachments via multipart form-data.
  247. # HTTP headers cannot contain newlines, but ntfy interprets
  248. # literal \n (backslash-n) as newlines in the Message header.
  249. headers["Filename"] = "photo.jpg"
  250. headers["Message"] = message.replace("\n", "\\n").encode("utf-8")
  251. response = await client.put(url, content=image_data, headers=headers)
  252. if response.status_code == 400 and "attachments not allowed" in response.text:
  253. # Server has attachments disabled — retry without the image
  254. headers.pop("Filename", None)
  255. headers.pop("Message", None)
  256. response = await client.post(url, content=message.encode("utf-8"), headers=headers)
  257. else:
  258. response = await client.post(url, content=message.encode("utf-8"), headers=headers)
  259. if response.status_code in (200, 204):
  260. return True, "Message sent successfully"
  261. if _looks_like_cloudflare_challenge(response):
  262. return False, (
  263. f"HTTP {response.status_code} — ntfy server is behind a Cloudflare "
  264. "challenge. Bambuddy was served the JS challenge page instead of "
  265. "reaching ntfy. Cloudflare cannot be solved from a backend; add a "
  266. "Cloudflare security-skip rule for this hostname, disable Bot "
  267. "Fight Mode, or front the server with Cloudflare Access using a "
  268. "service token. (#1534)"
  269. )
  270. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  271. async def _send_pushover(
  272. self, config: dict, title: str, message: str, image_data: bytes | None = None
  273. ) -> tuple[bool, str]:
  274. """Send notification via Pushover.
  275. Args:
  276. config: Provider configuration with user_key, app_token, priority
  277. title: Notification title
  278. message: Notification body
  279. image_data: Optional JPEG image bytes to attach (max 2.5MB)
  280. """
  281. user_key = config.get("user_key", "").strip()
  282. app_token = config.get("app_token", "").strip()
  283. try:
  284. priority = int(config.get("priority", 0))
  285. except (TypeError, ValueError):
  286. priority = 0
  287. if not user_key or not app_token:
  288. return False, "User key and app token are required"
  289. url = "https://api.pushover.net/1/messages.json"
  290. data = {
  291. "token": app_token,
  292. "user": user_key,
  293. "title": title,
  294. "message": message,
  295. "priority": priority,
  296. }
  297. # Emergency priority (2) keeps re-alerting until acknowledged, so
  298. # Pushover *requires* retry (how often, >= 30s) and expire (when to
  299. # give up, <= 10800s). Without them the API rejects the message. Only
  300. # send them at priority 2 — Pushover ignores them at other priorities.
  301. if priority == 2:
  302. try:
  303. retry = int(config.get("retry", 60))
  304. except (TypeError, ValueError):
  305. retry = 60
  306. try:
  307. expire = int(config.get("expire", 3600))
  308. except (TypeError, ValueError):
  309. expire = 3600
  310. data["retry"] = max(30, min(retry, 10800))
  311. data["expire"] = max(30, min(expire, 10800))
  312. client = await self._get_client()
  313. if image_data:
  314. # Pushover supports image attachments via multipart form-data
  315. files = {"attachment": ("photo.jpg", image_data, "image/jpeg")}
  316. response = await client.post(url, data=data, files=files)
  317. else:
  318. response = await client.post(url, data=data)
  319. if response.status_code == 200:
  320. return True, "Message sent successfully"
  321. else:
  322. try:
  323. error_data = response.json()
  324. errors = error_data.get("errors", [])
  325. return False, f"Pushover error: {', '.join(errors)}"
  326. except Exception:
  327. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  328. async def _send_telegram(self, config: dict, message: str, image_data: bytes | None = None) -> tuple[bool, str]:
  329. """Send notification via Telegram bot."""
  330. bot_token = config.get("bot_token", "").strip()
  331. chat_id = config.get("chat_id", "").strip()
  332. if not bot_token or not chat_id:
  333. return False, "Bot token and chat ID are required"
  334. # Escape underscores in the message body so Telegram Markdown
  335. # parsing doesn't break on job names like "A1_plate_8" or error
  336. # codes like "0300_0001". The title is already wrapped in *bold*
  337. # markers, so only escape after the first newline.
  338. if "\n" in message:
  339. title_part, body_part = message.split("\n", 1)
  340. body_part = body_part.replace("_", "\\_")
  341. message = f"{title_part}\n{body_part}"
  342. client = await self._get_client()
  343. if image_data:
  344. # Use sendPhoto to attach the thumbnail with the caption
  345. url = f"https://api.telegram.org/bot{bot_token}/sendPhoto"
  346. response = await client.post(
  347. url,
  348. data={"chat_id": chat_id, "caption": message, "parse_mode": "Markdown"},
  349. files={"photo": ("photo.jpg", image_data, "image/jpeg")},
  350. )
  351. else:
  352. url = f"https://api.telegram.org/bot{bot_token}/sendMessage"
  353. data = {
  354. "chat_id": chat_id,
  355. "text": message,
  356. "parse_mode": "Markdown",
  357. }
  358. response = await client.post(url, json=data)
  359. if response.status_code == 200:
  360. result = response.json()
  361. if result.get("ok"):
  362. return True, "Message sent successfully"
  363. else:
  364. return False, f"Telegram error: {result.get('description', 'Unknown error')}"
  365. else:
  366. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  367. async def _send_email(
  368. self,
  369. config: dict,
  370. subject: str,
  371. body: str,
  372. image_data: bytes | None = None,
  373. finish_photo_url: str | None = None,
  374. ) -> tuple[bool, str]:
  375. """Send notification via email (SMTP).
  376. Inline finish-photo embed is opt-in via the template: when the rendered
  377. ``body`` contains the substituted ``{finish_photo_url}`` value AND the
  378. finish-photo bytes are present, the message is built as
  379. ``multipart/related`` wrapping a ``multipart/alternative`` (plain + HTML)
  380. plus an inline ``MIMEImage`` with ``Content-ID: <bambuddy-finish-photo>``.
  381. The HTML part replaces the URL with ``<img src="cid:...">``; the plain-
  382. text part keeps the URL as a clickable link. When the template doesn't
  383. reference ``{finish_photo_url}`` (or image bytes aren't available), the
  384. original single-part text shape is used — no attachment, no surprise
  385. inline image (#1792).
  386. """
  387. smtp_server = config.get("smtp_server", "").strip()
  388. smtp_port = int(config.get("smtp_port", 587))
  389. username = config.get("username", "").strip()
  390. password = config.get("password", "").strip()
  391. from_email = config.get("from_email", "").strip()
  392. to_email = config.get("to_email", "").strip()
  393. # Security: "starttls" (port 587), "ssl" (port 465), "none" (port 25)
  394. security = config.get("security", "starttls")
  395. # Authentication: "true" or "false"
  396. auth_enabled = config.get("auth_enabled", "true").lower() == "true"
  397. if not all([smtp_server, from_email, to_email]):
  398. return False, "SMTP server, from email, and to email are required"
  399. if auth_enabled and not all([username, password]):
  400. return False, "Username and password are required when authentication is enabled"
  401. # Template-driven: only inline-embed when the user's template explicitly
  402. # referenced {finish_photo_url} (so the URL appears in the rendered body)
  403. # AND the photo bytes are available. Falls back to text-only otherwise.
  404. inline_photo = bool(image_data and finish_photo_url and finish_photo_url in body)
  405. try:
  406. if inline_photo:
  407. # multipart/related → (multipart/alternative → text, html) + inline image
  408. msg = MIMEMultipart("related")
  409. msg["From"] = from_email
  410. msg["To"] = to_email
  411. msg["Subject"] = f"[Bambuddy] {subject}"
  412. alt = MIMEMultipart("alternative")
  413. alt.attach(MIMEText(body, "plain"))
  414. # Build HTML body: escape the rendered body, then swap the
  415. # escaped URL substring for an inline <img> referencing the
  416. # MIMEImage we attach below. Done AFTER escape so the cid: URL
  417. # we inject isn't re-escaped.
  418. escaped_body = html.escape(body).replace("\n", "<br>\n")
  419. escaped_url = html.escape(finish_photo_url)
  420. img_tag = (
  421. '<img src="cid:bambuddy-finish-photo" '
  422. 'alt="Printer camera snapshot" '
  423. 'style="max-width:100%;height:auto;border:1px solid #ddd;border-radius:4px;">'
  424. )
  425. html_body = f"<html><body><p>{escaped_body.replace(escaped_url, img_tag)}</p></body></html>"
  426. alt.attach(MIMEText(html_body, "html"))
  427. msg.attach(alt)
  428. img = MIMEImage(image_data, _subtype="jpeg")
  429. # Angle-bracketed Content-ID per RFC 2392, referenced from HTML
  430. # without the brackets via ``cid:bambuddy-finish-photo``.
  431. img.add_header("Content-ID", "<bambuddy-finish-photo>")
  432. img.add_header("Content-Disposition", "inline", filename="finish-photo.jpg")
  433. msg.attach(img)
  434. else:
  435. msg = MIMEMultipart()
  436. msg["From"] = from_email
  437. msg["To"] = to_email
  438. msg["Subject"] = f"[Bambuddy] {subject}"
  439. msg.attach(MIMEText(body, "plain"))
  440. # smtplib is synchronous and blocking: a wedged / greylisting /
  441. # firewall-dropped relay leaves recv() stuck. Two problems, two
  442. # fixes (#2572):
  443. # 1. No timeout — smtplib defaults to the global socket timeout,
  444. # which this app never sets, so a stuck relay blocks forever.
  445. # Pass an explicit timeout to every connect.
  446. # 2. Run on the event loop — a stuck (or merely slow) send freezes
  447. # every other coroutine, including a DB session a caller is
  448. # holding open across this notification. Offload to a worker
  449. # thread so the loop stays live and the connection is released
  450. # on schedule.
  451. smtp_timeout = 30.0
  452. msg_str = msg.as_string()
  453. def _blocking_send() -> None:
  454. if security == "ssl":
  455. # Direct SSL connection (typically port 465)
  456. server = smtplib.SMTP_SSL(smtp_server, smtp_port, timeout=smtp_timeout)
  457. elif security == "starttls":
  458. # STARTTLS upgrade (typically port 587)
  459. server = smtplib.SMTP(smtp_server, smtp_port, timeout=smtp_timeout)
  460. server.starttls()
  461. else:
  462. # No encryption (typically port 25) - use with caution
  463. server = smtplib.SMTP(smtp_server, smtp_port, timeout=smtp_timeout)
  464. try:
  465. if auth_enabled:
  466. server.login(username, password)
  467. server.sendmail(from_email, to_email, msg_str)
  468. finally:
  469. # quit() in finally so a send error doesn't leak the socket.
  470. try:
  471. server.quit()
  472. except Exception: # noqa: BLE001 — closing a broken connection is best-effort
  473. pass
  474. await asyncio.to_thread(_blocking_send)
  475. return True, "Email sent successfully"
  476. except smtplib.SMTPAuthenticationError:
  477. return False, "SMTP authentication failed - check username/password"
  478. except smtplib.SMTPException as e:
  479. return False, f"SMTP error: {str(e)}"
  480. except Exception as e:
  481. return False, f"Email error: {str(e)}"
  482. async def _send_discord(
  483. self, config: dict, title: str, message: str, image_data: bytes | None = None
  484. ) -> tuple[bool, str]:
  485. """Send notification via Discord webhook."""
  486. webhook_url = config.get("webhook_url", "").strip()
  487. if not webhook_url:
  488. return False, "Webhook URL is required"
  489. if not (
  490. webhook_url.startswith("https://discord.com/api/webhooks/")
  491. or webhook_url.startswith("https://discordapp.com/api/webhooks/")
  492. ):
  493. return False, "Invalid Discord webhook URL"
  494. # Discord embed format for nicer messages
  495. embed = {
  496. "title": title,
  497. "description": message,
  498. "color": 0x00AE42, # Bambu green
  499. }
  500. client = await self._get_client()
  501. if image_data:
  502. # Attach image via multipart form-data and reference in embed
  503. embed["image"] = {"url": "attachment://photo.jpg"}
  504. payload = {"embeds": [embed]}
  505. response = await client.post(
  506. webhook_url,
  507. data={"payload_json": json.dumps(payload)},
  508. files={"files[0]": ("photo.jpg", image_data, "image/jpeg")},
  509. )
  510. else:
  511. response = await client.post(webhook_url, json={"embeds": [embed]})
  512. if response.status_code in (200, 204):
  513. return True, "Message sent successfully"
  514. else:
  515. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  516. async def _send_webhook(
  517. self,
  518. config: dict,
  519. title: str,
  520. message: str,
  521. image_data: bytes | None = None,
  522. event_type: str | None = None,
  523. variables: dict | None = None,
  524. ) -> tuple[bool, str]:
  525. """Send notification via generic webhook (POST JSON).
  526. Supports two payload formats:
  527. - generic: Custom field names with timestamp/source metadata + structured event data
  528. - slack: Slack/Mattermost compatible format (just {"text": "..."})
  529. """
  530. webhook_url = config.get("webhook_url", "").strip()
  531. auth_header = config.get("auth_header", "").strip()
  532. payload_format = config.get("payload_format", "generic").strip()
  533. if not webhook_url:
  534. return False, "Webhook URL is required"
  535. # Build payload based on format
  536. if payload_format == "slack":
  537. # Slack/Mattermost format - just text field
  538. data = {"text": f"*{title}*\n{message}"}
  539. else:
  540. # Generic format with custom field names
  541. custom_field_title = config.get("field_title", "title").strip() or "title"
  542. custom_field_message = config.get("field_message", "message").strip() or "message"
  543. data = {
  544. custom_field_title: title,
  545. custom_field_message: message,
  546. "timestamp": datetime.now().isoformat(),
  547. "source": "Bambuddy",
  548. }
  549. # For generic format, include structured event data for automation tools
  550. if payload_format != "slack":
  551. if event_type:
  552. data["event"] = event_type
  553. if variables:
  554. for key, value in variables.items():
  555. if key not in data: # Don't overwrite title/message/timestamp/source
  556. data[key] = value
  557. # Attach base64-encoded image when available (generic format only)
  558. if image_data and payload_format != "slack":
  559. import base64
  560. data["image"] = base64.b64encode(image_data).decode("ascii")
  561. headers = {"Content-Type": "application/json"}
  562. if auth_header:
  563. # Support "Bearer token" or just "token" format
  564. if " " in auth_header:
  565. headers["Authorization"] = auth_header
  566. else:
  567. headers["Authorization"] = f"Bearer {auth_header}"
  568. client = await self._get_client()
  569. try:
  570. response = await client.post(webhook_url, json=data, headers=headers)
  571. if response.status_code in (200, 201, 202, 204):
  572. return True, "Webhook delivered successfully"
  573. else:
  574. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  575. except Exception as e:
  576. return False, f"Webhook error: {str(e)}"
  577. async def _send_homeassistant(
  578. self, config: dict, title: str, message: str, db: AsyncSession | None = None
  579. ) -> tuple[bool, str]:
  580. """Send notification via Home Assistant.
  581. Uses the globally configured HA URL/token from settings.
  582. Defaults to persistent_notification/create, but supports
  583. custom services via config["service"] (e.g. notify.mobile_app_myphone).
  584. """
  585. # Get HA connection settings from global config
  586. ha_url = ""
  587. ha_token = ""
  588. if db:
  589. from backend.app.api.routes.settings import get_homeassistant_settings
  590. try:
  591. ha_settings = await get_homeassistant_settings(db)
  592. ha_url = ha_settings.get("ha_url", "")
  593. ha_token = ha_settings.get("ha_token", "")
  594. except Exception as e:
  595. logger.warning("Failed to read HA settings from database: %s", e)
  596. else:
  597. # Fallback: read directly from environment if no DB session
  598. import os
  599. ha_url = os.environ.get("HA_URL", "")
  600. ha_token = os.environ.get("HA_TOKEN", "")
  601. if not ha_url or not ha_token:
  602. return False, (
  603. "Home Assistant is not configured. Please set HA URL and token in Settings → Network → Home Assistant."
  604. )
  605. # Determine which HA service to call - Default: persistent_notification.create
  606. service = (config.get("service") or "").strip()
  607. if service:
  608. # Allow in different forms:
  609. # - notify.mobile_app_<device>
  610. # - notify/mobile_app_<device>
  611. # - api/services/notify/mobile_app_<device>
  612. service_str = service.lstrip("/")
  613. if service_str.startswith("api/services/"):
  614. endpoint = service_str
  615. elif "/" in service_str:
  616. endpoint = f"api/services/{service_str}"
  617. elif "." in service_str:
  618. domain, svc = service_str.split(".", 1)
  619. endpoint = f"api/services/{domain}/{svc}"
  620. else:
  621. return False, (
  622. "Invalid Home Assistant service name. Use e.g. 'notify.mobile_app_yourdevice' or 'notify/your_service'."
  623. )
  624. if not re.match(r"^api/services/[a-zA-Z0-9_]+/[a-zA-Z0-9_]+$", endpoint):
  625. return False, (
  626. "Invalid Home Assistant service name. Domain and service must only contain letters, numbers, and underscores."
  627. )
  628. else:
  629. endpoint = "api/services/persistent_notification/create"
  630. url = f"{ha_url.rstrip('/')}/{endpoint}"
  631. headers = {
  632. "Authorization": f"Bearer {ha_token}",
  633. "Content-Type": "application/json",
  634. }
  635. payload = {
  636. "title": title,
  637. "message": message,
  638. }
  639. client = await self._get_client()
  640. response = await client.post(url, json=payload, headers=headers)
  641. if response.status_code in (200, 201):
  642. return True, "Notification sent via Home Assistant"
  643. elif response.status_code == 401:
  644. return False, "Home Assistant authentication failed - check your token"
  645. else:
  646. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  647. async def _send_to_provider(
  648. self,
  649. provider: NotificationProvider,
  650. title: str,
  651. message: str,
  652. db: AsyncSession | None = None,
  653. image_data: bytes | None = None,
  654. event_type: str | None = None,
  655. variables: dict | None = None,
  656. ) -> tuple[bool, str]:
  657. """Send notification to a specific provider."""
  658. # Check quiet hours
  659. if self._is_in_quiet_hours(provider):
  660. logger.info("Skipping notification to %s - quiet hours active", provider.name)
  661. return True, "Skipped - quiet hours"
  662. config = json.loads(provider.config) if isinstance(provider.config, str) else provider.config
  663. try:
  664. if provider.provider_type == "callmebot":
  665. return await self._send_callmebot(config, f"{title}\n{message}")
  666. elif provider.provider_type == "ntfy":
  667. return await self._send_ntfy(config, title, message, image_data=image_data, event_type=event_type)
  668. elif provider.provider_type == "pushover":
  669. return await self._send_pushover(config, title, message, image_data=image_data)
  670. elif provider.provider_type == "telegram":
  671. return await self._send_telegram(config, f"*{title}*\n{message}", image_data=image_data)
  672. elif provider.provider_type == "email":
  673. # finish_photo_url is pulled from the rendered template variables
  674. # so _send_email can detect whether the template referenced the
  675. # URL and inline-embed the photo only in that case.
  676. finish_photo_url = (variables or {}).get("finish_photo_url")
  677. return await self._send_email(
  678. config, title, message, image_data=image_data, finish_photo_url=finish_photo_url
  679. )
  680. elif provider.provider_type == "discord":
  681. return await self._send_discord(config, title, message, image_data=image_data)
  682. elif provider.provider_type == "webhook":
  683. return await self._send_webhook(
  684. config, title, message, image_data=image_data, event_type=event_type, variables=variables
  685. )
  686. elif provider.provider_type == "homeassistant":
  687. return await self._send_homeassistant(config, title, message, db=db)
  688. else:
  689. return False, f"Unknown provider type: {provider.provider_type}"
  690. except Exception as e:
  691. logger.exception("Error sending notification via %s", provider.provider_type)
  692. return False, str(e)
  693. async def _update_provider_status(
  694. self, db: AsyncSession, provider_id: int, success: bool, error: str | None = None
  695. ):
  696. """Update provider status after sending notification."""
  697. result = await db.execute(select(NotificationProvider).where(NotificationProvider.id == provider_id))
  698. provider = result.scalar_one_or_none()
  699. if provider:
  700. if success:
  701. provider.last_success = datetime.now(timezone.utc)
  702. else:
  703. provider.last_error = error
  704. provider.last_error_at = datetime.now(timezone.utc)
  705. await db.commit()
  706. async def _get_providers_for_event(
  707. self,
  708. db: AsyncSession,
  709. event_field: str,
  710. printer_id: int | None = None,
  711. ) -> list[NotificationProvider]:
  712. """Get all enabled providers that want a specific event type."""
  713. # Build the query dynamically based on event field
  714. query = select(NotificationProvider).where(
  715. NotificationProvider.enabled.is_(True),
  716. getattr(NotificationProvider, event_field).is_(True),
  717. )
  718. if printer_id is not None:
  719. query = query.where(
  720. (NotificationProvider.printer_id.is_(None)) | (NotificationProvider.printer_id == printer_id)
  721. )
  722. result = await db.execute(query)
  723. return list(result.scalars().all())
  724. async def _log_notification(
  725. self,
  726. db: AsyncSession,
  727. provider_id: int,
  728. event_type: str,
  729. title: str,
  730. message: str,
  731. success: bool,
  732. error_message: str | None = None,
  733. printer_id: int | None = None,
  734. printer_name: str | None = None,
  735. ):
  736. """Create a log entry for a sent notification."""
  737. try:
  738. log = NotificationLog(
  739. provider_id=provider_id,
  740. event_type=event_type,
  741. title=title,
  742. message=message,
  743. success=success,
  744. error_message=error_message,
  745. printer_id=printer_id,
  746. printer_name=printer_name,
  747. )
  748. db.add(log)
  749. await db.commit()
  750. except Exception as e:
  751. logger.warning("Failed to log notification: %s", e)
  752. # Don't fail the notification just because logging failed
  753. async def _send_to_providers(
  754. self,
  755. providers: list[NotificationProvider],
  756. title: str,
  757. message: str,
  758. db: AsyncSession,
  759. event_type: str = "unknown",
  760. printer_id: int | None = None,
  761. printer_name: str | None = None,
  762. force_immediate: bool = False,
  763. image_data: bytes | None = None,
  764. variables: dict | None = None,
  765. ):
  766. """Send notification to multiple providers and log the results.
  767. All notifications are always sent immediately. If digest mode is enabled,
  768. the notification is ALSO queued for the daily digest summary.
  769. """
  770. for provider in providers:
  771. try:
  772. # Always send notification immediately
  773. success, error = await self._send_to_provider(
  774. provider, title, message, db, image_data=image_data, event_type=event_type, variables=variables
  775. )
  776. # Also queue for digest if enabled (digest is a summary, not a queue)
  777. if provider.daily_digest_enabled and provider.daily_digest_time:
  778. await self._queue_for_digest(
  779. provider=provider,
  780. event_type=event_type,
  781. title=title,
  782. message=message,
  783. db=db,
  784. printer_id=printer_id,
  785. printer_name=printer_name,
  786. )
  787. await self._update_provider_status(db, provider.id, success, error if not success else None)
  788. await self._log_notification(
  789. db=db,
  790. provider_id=provider.id,
  791. event_type=event_type,
  792. title=title,
  793. message=message,
  794. success=success,
  795. error_message=error if not success else None,
  796. printer_id=printer_id,
  797. printer_name=printer_name,
  798. )
  799. if success:
  800. logger.info("Sent notification via %s", provider.name)
  801. else:
  802. logger.warning("Failed to send notification via %s: %s", provider.name, error)
  803. except Exception as e:
  804. logger.exception("Error sending notification via %s", provider.name)
  805. await self._update_provider_status(db, provider.id, False, str(e))
  806. await self._log_notification(
  807. db=db,
  808. provider_id=provider.id,
  809. event_type=event_type,
  810. title=title,
  811. message=message,
  812. success=False,
  813. error_message=str(e),
  814. printer_id=printer_id,
  815. printer_name=printer_name,
  816. )
  817. async def on_print_start(
  818. self,
  819. printer_id: int,
  820. printer_name: str,
  821. data: dict,
  822. db: AsyncSession,
  823. archive_data: dict | None = None,
  824. ):
  825. """Handle print start event - send notifications to relevant providers.
  826. Args:
  827. printer_id: The printer ID
  828. printer_name: The printer name
  829. data: MQTT event data with filename, subtask_name, remaining_time, raw_data
  830. db: Database session
  831. archive_data: Optional archive data with print_time_seconds from 3MF parsing
  832. """
  833. logger.info("on_print_start called for printer %s (%s)", printer_id, printer_name)
  834. providers = await self._get_providers_for_event(db, "on_print_start", printer_id)
  835. if not providers:
  836. logger.info("No notification providers configured for print_start event on printer %s", printer_id)
  837. return
  838. # Use subtask_name (project name) if available, otherwise use filename
  839. subtask_name = data.get("subtask_name")
  840. if subtask_name:
  841. # Replace underscores with spaces for readability
  842. filename = subtask_name.replace("_", " ")
  843. else:
  844. filename = self._clean_filename(data.get("filename", "Unknown"))
  845. # Priority for estimated_time:
  846. # 1. Archive's print_time_seconds from 3MF parsing (most reliable)
  847. # 2. MQTT remaining_time (may be 0 at print start)
  848. # 3. raw_data mc_remaining_time
  849. estimated_time = None
  850. # Try archive data first (from 3MF parsing - most reliable)
  851. if archive_data and archive_data.get("print_time_seconds"):
  852. estimated_time = archive_data["print_time_seconds"]
  853. logger.debug("Using print_time_seconds from archive: %s", estimated_time)
  854. # Fall back to MQTT remaining_time
  855. if estimated_time is None:
  856. estimated_time = data.get("remaining_time")
  857. if estimated_time:
  858. logger.debug("Using remaining_time from MQTT: %s", estimated_time)
  859. # Last resort: raw_data mc_remaining_time (in minutes, convert to seconds)
  860. if estimated_time is None:
  861. raw_time = data.get("raw_data", {}).get("mc_remaining_time")
  862. if raw_time:
  863. estimated_time = raw_time * 60
  864. logger.debug("Using mc_remaining_time from raw_data: %s", estimated_time)
  865. time_str = self._format_duration(estimated_time)
  866. eta_str = await self._format_eta(estimated_time, db)
  867. variables = {
  868. "printer": printer_name,
  869. "filename": filename,
  870. "estimated_time": time_str,
  871. "eta": eta_str,
  872. }
  873. # Extract image data for providers that support attachments (e.g. Pushover)
  874. image_data = None
  875. if archive_data:
  876. image_data = archive_data.get("image_data")
  877. logger.info("Found %s providers for print_start: %s", len(providers), [p.name for p in providers])
  878. title, message = await self._build_message_from_template(db, "print_start", variables)
  879. await self._send_to_providers(
  880. providers,
  881. title,
  882. message,
  883. db,
  884. "print_start",
  885. printer_id,
  886. printer_name,
  887. image_data=image_data,
  888. variables=variables,
  889. )
  890. async def on_print_complete(
  891. self,
  892. printer_id: int,
  893. printer_name: str,
  894. status: str,
  895. data: dict,
  896. db: AsyncSession,
  897. archive_data: dict | None = None,
  898. ):
  899. """Handle print complete event - send notifications to relevant providers."""
  900. logger.info("on_print_complete called for printer %s (%s), status=%s", printer_id, printer_name, status)
  901. # Determine event type based on status
  902. if status == "completed":
  903. event_field = "on_print_complete"
  904. event_type = "print_complete"
  905. elif status in ("failed",):
  906. event_field = "on_print_failed"
  907. event_type = "print_failed"
  908. elif status in ("aborted", "stopped", "cancelled"):
  909. event_field = "on_print_stopped"
  910. event_type = "print_stopped"
  911. else:
  912. logger.warning("Unknown print status '%s', defaulting to on_print_complete", status)
  913. event_field = "on_print_complete"
  914. event_type = "print_complete"
  915. providers = await self._get_providers_for_event(db, event_field, printer_id)
  916. if not providers:
  917. logger.info("No notification providers configured for %s event on printer %s", event_field, printer_id)
  918. return
  919. # Use subtask_name (project name) if available, otherwise use filename
  920. subtask_name = data.get("subtask_name")
  921. if subtask_name:
  922. filename = subtask_name.replace("_", " ")
  923. else:
  924. filename = self._clean_filename(data.get("filename", "Unknown"))
  925. variables = {
  926. "printer": printer_name,
  927. "filename": filename,
  928. "duration": "Unknown",
  929. "filament_grams": "Unknown",
  930. "reason": "Unknown",
  931. }
  932. if archive_data:
  933. # {{duration}} on completion / failure / stopped events is the *actual*
  934. # elapsed time (#1198). Slicer-estimated print_time_seconds is only used
  935. # as a last-resort fallback when timestamps weren't recorded.
  936. duration_seconds = archive_data.get("actual_time_seconds") or archive_data.get("print_time_seconds")
  937. if duration_seconds:
  938. variables["duration"] = self._format_duration(duration_seconds)
  939. if archive_data.get("actual_filament_grams"):
  940. variables["filament_grams"] = f"{archive_data['actual_filament_grams']:.1f}"
  941. if status == "failed" and archive_data.get("failure_reason"):
  942. variables["reason"] = archive_data["failure_reason"]
  943. if archive_data.get("finish_photo_url"):
  944. variables["finish_photo_url"] = archive_data["finish_photo_url"]
  945. # Build per-slot breakdown string with AMS info when available
  946. if archive_data.get("usage_results"):
  947. parts = []
  948. for u in archive_data["usage_results"]:
  949. ams_id = u.get("ams_id", 0)
  950. tray_id = u.get("tray_id", 0)
  951. material = u.get("material", "Unknown") or "Unknown"
  952. used = u.get("weight_used", 0)
  953. if ams_id >= 128:
  954. slot_label = "Ext"
  955. else:
  956. slot_label = f"AMS-{chr(65 + ams_id)} T{tray_id + 1}"
  957. parts.append(f"{slot_label} {material}: {used:.1f}g")
  958. variables["filament_details"] = " | ".join(parts)
  959. elif archive_data.get("filament_slots"):
  960. parts = []
  961. for slot in archive_data["filament_slots"]:
  962. ftype = slot.get("type", "Unknown") or "Unknown"
  963. used = slot.get("used_g", 0)
  964. parts.append(f"{ftype}: {used:.1f}g")
  965. variables["filament_details"] = " | ".join(parts)
  966. # Add progress for partial prints
  967. if archive_data.get("progress") is not None:
  968. variables["progress"] = str(archive_data["progress"])
  969. # Extract image data for providers that support attachments (e.g. Pushover)
  970. image_data = None
  971. if archive_data:
  972. image_data = archive_data.get("image_data")
  973. logger.info("Found %s providers for %s: %s", len(providers), event_field, [p.name for p in providers])
  974. title, message = await self._build_message_from_template(db, event_type, variables)
  975. await self._send_to_providers(
  976. providers,
  977. title,
  978. message,
  979. db,
  980. event_type,
  981. printer_id,
  982. printer_name,
  983. image_data=image_data,
  984. variables=variables,
  985. )
  986. async def on_print_progress(
  987. self,
  988. printer_id: int,
  989. printer_name: str,
  990. filename: str,
  991. progress: int,
  992. db: AsyncSession,
  993. remaining_time: int | None = None,
  994. image_data: bytes | None = None,
  995. ):
  996. """Handle print progress milestone (25%, 50%, 75%)."""
  997. providers = await self._get_providers_for_event(db, "on_print_progress", printer_id)
  998. if not providers:
  999. return
  1000. eta_str = await self._format_eta(remaining_time, db)
  1001. variables = {
  1002. "printer": printer_name,
  1003. "filename": self._clean_filename(filename),
  1004. "progress": str(progress),
  1005. "remaining_time": self._format_duration(remaining_time) if remaining_time else "Unknown",
  1006. "eta": eta_str,
  1007. }
  1008. title, message = await self._build_message_from_template(db, "print_progress", variables)
  1009. await self._send_to_providers(
  1010. providers,
  1011. title,
  1012. message,
  1013. db,
  1014. "print_progress",
  1015. printer_id,
  1016. printer_name,
  1017. image_data=image_data,
  1018. variables=variables,
  1019. )
  1020. async def on_print_missing_spool_assignment(
  1021. self,
  1022. printer_id: int,
  1023. printer_name: str,
  1024. missing_slots: list[dict[str, str]],
  1025. db: AsyncSession,
  1026. ):
  1027. """Handle print-start event when required trays are missing spool assignments."""
  1028. if not missing_slots:
  1029. return
  1030. providers = await self._get_providers_for_event(db, "on_print_missing_spool_assignment", printer_id)
  1031. if not providers:
  1032. return
  1033. missing_slot_names = ", ".join(slot.get("slot", "Unknown") for slot in missing_slots)
  1034. detail_lines = []
  1035. for slot in missing_slots:
  1036. slot_name = slot.get("slot", "Unknown")
  1037. profile = slot.get("profile", "Unknown")
  1038. detail_lines.append(f"- {slot_name}: {profile}")
  1039. missing_profile_details = "\n".join(detail_lines)
  1040. variables = {
  1041. "printer": printer_name,
  1042. "missing_slots": missing_slot_names,
  1043. "missing_slot_details": missing_profile_details,
  1044. }
  1045. title, message = await self._build_message_from_template(db, "print_missing_spool_assignment", variables)
  1046. await self._send_to_providers(
  1047. providers,
  1048. title,
  1049. message,
  1050. db,
  1051. "print_missing_spool_assignment",
  1052. printer_id,
  1053. printer_name,
  1054. force_immediate=True,
  1055. variables=variables,
  1056. )
  1057. async def on_printer_offline(self, printer_id: int, printer_name: str, db: AsyncSession):
  1058. """Handle printer offline event."""
  1059. providers = await self._get_providers_for_event(db, "on_printer_offline", printer_id)
  1060. if not providers:
  1061. return
  1062. variables = {"printer": printer_name}
  1063. title, message = await self._build_message_from_template(db, "printer_offline", variables)
  1064. await self._send_to_providers(
  1065. providers, title, message, db, "printer_offline", printer_id, printer_name, variables=variables
  1066. )
  1067. async def on_printer_error(
  1068. self,
  1069. printer_id: int,
  1070. printer_name: str,
  1071. error_type: str,
  1072. db: AsyncSession,
  1073. error_detail: str | None = None,
  1074. image_data: bytes | None = None,
  1075. ):
  1076. """Handle printer error event (AMS issues, etc.)."""
  1077. providers = await self._get_providers_for_event(db, "on_printer_error", printer_id)
  1078. if not providers:
  1079. return
  1080. variables = {
  1081. "printer": printer_name,
  1082. "error_type": error_type,
  1083. "error_detail": error_detail or "No details available",
  1084. }
  1085. title, message = await self._build_message_from_template(db, "printer_error", variables)
  1086. await self._send_to_providers(
  1087. providers,
  1088. title,
  1089. message,
  1090. db,
  1091. "printer_error",
  1092. printer_id,
  1093. printer_name,
  1094. image_data=image_data,
  1095. variables=variables,
  1096. )
  1097. async def on_ai_failure_detection(
  1098. self,
  1099. printer_id: int,
  1100. printer_name: str,
  1101. task_name: str,
  1102. confidence: float,
  1103. action: str,
  1104. db: AsyncSession,
  1105. image_data: bytes | None = None,
  1106. ):
  1107. """Handle AI failure-detection event (Obico spaghetti / print-failure ML).
  1108. Split out of on_printer_error (#1794) so a user can subscribe to AI
  1109. alerts without also being paged for every HMS hardware code.
  1110. """
  1111. providers = await self._get_providers_for_event(db, "on_ai_failure_detection", printer_id)
  1112. if not providers:
  1113. return
  1114. variables = {
  1115. "printer": printer_name,
  1116. "task_name": task_name or "current job",
  1117. "confidence": f"{confidence:.2f}",
  1118. "action": action,
  1119. }
  1120. title, message = await self._build_message_from_template(db, "ai_failure_detection", variables)
  1121. await self._send_to_providers(
  1122. providers,
  1123. title,
  1124. message,
  1125. db,
  1126. "ai_failure_detection",
  1127. printer_id,
  1128. printer_name,
  1129. image_data=image_data,
  1130. variables=variables,
  1131. )
  1132. async def on_plate_not_empty(
  1133. self,
  1134. printer_id: int,
  1135. printer_name: str,
  1136. db: AsyncSession,
  1137. difference_percent: float | None = None,
  1138. ):
  1139. """Handle plate not empty event - objects detected on build plate before print."""
  1140. providers = await self._get_providers_for_event(db, "on_plate_not_empty", printer_id)
  1141. if not providers:
  1142. return
  1143. variables = {
  1144. "printer": printer_name,
  1145. "difference_percent": f"{difference_percent:.1f}" if difference_percent else "N/A",
  1146. }
  1147. title, message = await self._build_message_from_template(db, "plate_not_empty", variables)
  1148. await self._send_to_providers(
  1149. providers,
  1150. title,
  1151. message,
  1152. db,
  1153. "plate_not_empty",
  1154. printer_id,
  1155. printer_name,
  1156. force_immediate=True,
  1157. variables=variables,
  1158. )
  1159. async def on_filament_low(
  1160. self,
  1161. printer_id: int,
  1162. printer_name: str,
  1163. slot: int,
  1164. remaining_percent: int,
  1165. db: AsyncSession,
  1166. color: str | None = None,
  1167. ):
  1168. """Handle low filament event."""
  1169. providers = await self._get_providers_for_event(db, "on_filament_low", printer_id)
  1170. if not providers:
  1171. return
  1172. variables = {
  1173. "printer": printer_name,
  1174. "slot": str(slot),
  1175. "remaining_percent": str(remaining_percent),
  1176. "color": color or "",
  1177. }
  1178. title, message = await self._build_message_from_template(db, "filament_low", variables)
  1179. await self._send_to_providers(
  1180. providers, title, message, db, "filament_low", printer_id, printer_name, variables=variables
  1181. )
  1182. async def on_maintenance_due(
  1183. self,
  1184. printer_id: int,
  1185. printer_name: str,
  1186. maintenance_items: list[dict],
  1187. db: AsyncSession,
  1188. ):
  1189. """Handle maintenance due event - sends notification when maintenance is due or warning."""
  1190. if not maintenance_items:
  1191. return
  1192. providers = await self._get_providers_for_event(db, "on_maintenance_due", printer_id)
  1193. if not providers:
  1194. logger.info("No notification providers configured for maintenance_due event on printer %s", printer_id)
  1195. return
  1196. # Format maintenance items list
  1197. items_list = []
  1198. for item in maintenance_items:
  1199. status = "OVERDUE" if item.get("is_due") else "Soon"
  1200. items_list.append(f"- {item['name']} ({status})")
  1201. items_str = "\n".join(items_list)
  1202. variables = {
  1203. "printer": printer_name,
  1204. "items": items_str,
  1205. }
  1206. logger.info("Found %s providers for maintenance_due: %s", len(providers), [p.name for p in providers])
  1207. title, message = await self._build_message_from_template(db, "maintenance_due", variables)
  1208. await self._send_to_providers(
  1209. providers, title, message, db, "maintenance_due", printer_id, printer_name, variables=variables
  1210. )
  1211. async def on_ams_humidity_high(
  1212. self,
  1213. printer_id: int,
  1214. printer_name: str,
  1215. ams_label: str,
  1216. humidity: float,
  1217. threshold: float,
  1218. db: AsyncSession,
  1219. ):
  1220. """Handle AMS high humidity alarm event. Always sends immediately (bypasses digest)."""
  1221. providers = await self._get_providers_for_event(db, "on_ams_humidity_high", printer_id)
  1222. if not providers:
  1223. return
  1224. variables = {
  1225. "printer": printer_name,
  1226. "ams_label": ams_label,
  1227. "humidity": f"{humidity:.0f}",
  1228. "threshold": f"{threshold:.0f}",
  1229. }
  1230. title, message = await self._build_message_from_template(db, "ams_humidity_high", variables)
  1231. # Alarms always send immediately, bypassing digest mode
  1232. await self._send_to_providers(
  1233. providers,
  1234. title,
  1235. message,
  1236. db,
  1237. "ams_humidity_high",
  1238. printer_id,
  1239. printer_name,
  1240. force_immediate=True,
  1241. variables=variables,
  1242. )
  1243. async def on_ams_temperature_high(
  1244. self,
  1245. printer_id: int,
  1246. printer_name: str,
  1247. ams_label: str,
  1248. temperature: float,
  1249. threshold: float,
  1250. db: AsyncSession,
  1251. ):
  1252. """Handle AMS high temperature alarm event. Always sends immediately (bypasses digest)."""
  1253. providers = await self._get_providers_for_event(db, "on_ams_temperature_high", printer_id)
  1254. if not providers:
  1255. return
  1256. variables = {
  1257. "printer": printer_name,
  1258. "ams_label": ams_label,
  1259. "temperature": f"{temperature:.1f}",
  1260. "threshold": f"{threshold:.1f}",
  1261. }
  1262. title, message = await self._build_message_from_template(db, "ams_temperature_high", variables)
  1263. # Alarms always send immediately, bypassing digest mode
  1264. await self._send_to_providers(
  1265. providers,
  1266. title,
  1267. message,
  1268. db,
  1269. "ams_temperature_high",
  1270. printer_id,
  1271. printer_name,
  1272. force_immediate=True,
  1273. variables=variables,
  1274. )
  1275. async def on_ams_ht_humidity_high(
  1276. self,
  1277. printer_id: int,
  1278. printer_name: str,
  1279. ams_label: str,
  1280. humidity: float,
  1281. threshold: float,
  1282. db: AsyncSession,
  1283. ):
  1284. """Handle AMS-HT high humidity alarm event. Always sends immediately (bypasses digest)."""
  1285. providers = await self._get_providers_for_event(db, "on_ams_ht_humidity_high", printer_id)
  1286. if not providers:
  1287. return
  1288. variables = {
  1289. "printer": printer_name,
  1290. "ams_label": ams_label,
  1291. "humidity": f"{humidity:.0f}",
  1292. "threshold": f"{threshold:.0f}",
  1293. }
  1294. # Use the same template as regular AMS (can create separate templates later if needed)
  1295. title, message = await self._build_message_from_template(db, "ams_humidity_high", variables)
  1296. # Alarms always send immediately, bypassing digest mode
  1297. await self._send_to_providers(
  1298. providers,
  1299. title,
  1300. message,
  1301. db,
  1302. "ams_ht_humidity_high",
  1303. printer_id,
  1304. printer_name,
  1305. force_immediate=True,
  1306. variables=variables,
  1307. )
  1308. async def on_ams_ht_temperature_high(
  1309. self,
  1310. printer_id: int,
  1311. printer_name: str,
  1312. ams_label: str,
  1313. temperature: float,
  1314. threshold: float,
  1315. db: AsyncSession,
  1316. ):
  1317. """Handle AMS-HT high temperature alarm event. Always sends immediately (bypasses digest)."""
  1318. providers = await self._get_providers_for_event(db, "on_ams_ht_temperature_high", printer_id)
  1319. if not providers:
  1320. return
  1321. variables = {
  1322. "printer": printer_name,
  1323. "ams_label": ams_label,
  1324. "temperature": f"{temperature:.1f}",
  1325. "threshold": f"{threshold:.1f}",
  1326. }
  1327. # Use the same template as regular AMS (can create separate templates later if needed)
  1328. title, message = await self._build_message_from_template(db, "ams_temperature_high", variables)
  1329. # Alarms always send immediately, bypassing digest mode
  1330. await self._send_to_providers(
  1331. providers,
  1332. title,
  1333. message,
  1334. db,
  1335. "ams_ht_temperature_high",
  1336. printer_id,
  1337. printer_name,
  1338. force_immediate=True,
  1339. variables=variables,
  1340. )
  1341. async def on_bed_cooled(
  1342. self,
  1343. printer_id: int,
  1344. printer_name: str,
  1345. bed_temp: float,
  1346. threshold: float,
  1347. filename: str,
  1348. db: AsyncSession,
  1349. ):
  1350. """Handle bed cooled event - bed temperature dropped below threshold after print."""
  1351. providers = await self._get_providers_for_event(db, "on_bed_cooled", printer_id)
  1352. if not providers:
  1353. return
  1354. variables = {
  1355. "printer": printer_name,
  1356. "bed_temp": f"{bed_temp:.0f}",
  1357. "threshold": f"{threshold:.0f}",
  1358. "filename": self._clean_filename(filename) if filename else "Unknown",
  1359. }
  1360. title, message = await self._build_message_from_template(db, "bed_cooled", variables)
  1361. await self._send_to_providers(
  1362. providers, title, message, db, "bed_cooled", printer_id, printer_name, variables=variables
  1363. )
  1364. async def on_first_layer_complete(
  1365. self,
  1366. printer_id: int,
  1367. printer_name: str,
  1368. filename: str,
  1369. total_layers: int,
  1370. db: AsyncSession,
  1371. image_data: bytes | None = None,
  1372. ):
  1373. """Handle first layer complete event."""
  1374. providers = await self._get_providers_for_event(db, "on_first_layer_complete", printer_id)
  1375. if not providers:
  1376. return
  1377. variables = {
  1378. "printer": printer_name,
  1379. "filename": self._clean_filename(filename),
  1380. "total_layers": str(total_layers),
  1381. }
  1382. title, message = await self._build_message_from_template(db, "first_layer_complete", variables)
  1383. await self._send_to_providers(
  1384. providers,
  1385. title,
  1386. message,
  1387. db,
  1388. "first_layer_complete",
  1389. printer_id,
  1390. printer_name,
  1391. image_data=image_data,
  1392. variables=variables,
  1393. )
  1394. def clear_template_cache(self):
  1395. """Clear the template cache. Call this when templates are updated."""
  1396. self._template_cache.clear()
  1397. async def send_user_print_email(
  1398. self,
  1399. event_type: str,
  1400. created_by_id: int | None,
  1401. printer_name: str,
  1402. filename: str,
  1403. db: AsyncSession,
  1404. ) -> None:
  1405. """Send a print event email notification to the user who submitted the job.
  1406. Args:
  1407. event_type: 'user_print_start', 'user_print_complete', 'user_print_failed', or 'user_print_stopped'
  1408. created_by_id: User ID who submitted the print job (from archive)
  1409. printer_name: Name of the printer
  1410. filename: Raw filename or subtask name
  1411. db: Database session
  1412. """
  1413. if created_by_id is None:
  1414. logger.debug("[EMAIL] Skipping user print email (%s): no created_by_id", event_type)
  1415. return
  1416. try:
  1417. # Check if advanced auth is enabled - required for user email notifications
  1418. from backend.app.models.settings import Settings
  1419. result = await db.execute(select(Settings).where(Settings.key == "advanced_auth_enabled"))
  1420. setting = result.scalar_one_or_none()
  1421. if not setting or setting.value.lower() != "true":
  1422. logger.debug("[EMAIL] Skipping user print email (%s): advanced_auth not enabled", event_type)
  1423. return
  1424. # Check if user notifications are enabled (admin-controlled toggle)
  1425. notif_enabled_result = await db.execute(
  1426. select(Settings).where(Settings.key == "user_notifications_enabled")
  1427. )
  1428. notif_enabled_setting = notif_enabled_result.scalar_one_or_none()
  1429. if notif_enabled_setting and notif_enabled_setting.value.lower() == "false":
  1430. logger.debug("[EMAIL] Skipping user print email (%s): user_notifications_enabled is false", event_type)
  1431. return
  1432. # Check SMTP settings are configured - required for sending emails
  1433. from backend.app.services.email_service import get_smtp_settings, send_user_print_notification
  1434. smtp_settings = await get_smtp_settings(db)
  1435. if not smtp_settings:
  1436. logger.debug("[EMAIL] Skipping user print email (%s): SMTP settings not configured", event_type)
  1437. return
  1438. # Load user preferences
  1439. from backend.app.models.user import User
  1440. from backend.app.models.user_email_pref import UserEmailPreference
  1441. user_result = await db.execute(select(User).where(User.id == created_by_id))
  1442. user = user_result.scalar_one_or_none()
  1443. if user is None or not user.email:
  1444. logger.debug(
  1445. "[EMAIL] Skipping user print email (%s): user %s not found or has no email address",
  1446. event_type,
  1447. created_by_id,
  1448. )
  1449. return
  1450. # Load user's notification preferences
  1451. pref_result = await db.execute(
  1452. select(UserEmailPreference).where(UserEmailPreference.user_id == created_by_id)
  1453. )
  1454. pref = pref_result.scalar_one_or_none()
  1455. # Determine if this event type should be sent
  1456. should_send = False
  1457. if event_type == "user_print_start":
  1458. should_send = pref is None or pref.notify_print_start
  1459. elif event_type == "user_print_complete":
  1460. should_send = pref is None or pref.notify_print_complete
  1461. elif event_type == "user_print_failed":
  1462. should_send = pref is None or pref.notify_print_failed
  1463. elif event_type == "user_print_stopped":
  1464. should_send = pref is None or pref.notify_print_stopped
  1465. if not should_send:
  1466. logger.debug(
  1467. "[EMAIL] Skipping user print email (%s): user %s has notifications disabled for this event",
  1468. event_type,
  1469. created_by_id,
  1470. )
  1471. return
  1472. logger.info(
  1473. "[EMAIL] Sending user print email: event=%s, user=%s (%s), printer=%s, file=%s",
  1474. event_type,
  1475. user.username,
  1476. user.email,
  1477. printer_name,
  1478. filename,
  1479. )
  1480. # Build variables
  1481. variables = {
  1482. "printer": printer_name,
  1483. "filename": self._clean_filename(filename),
  1484. }
  1485. # Send the email
  1486. await send_user_print_notification(
  1487. db=db,
  1488. event_type=event_type,
  1489. user_email=user.email,
  1490. username=user.username,
  1491. variables=variables,
  1492. )
  1493. logger.info("[EMAIL] User print email sent: event=%s → %s", event_type, user.email)
  1494. except Exception as e:
  1495. logger.warning("Failed to send user print email notification: %s", e, exc_info=True)
  1496. # ==================== Queue Notifications ====================
  1497. async def on_queue_job_added(
  1498. self,
  1499. job_name: str,
  1500. target: str,
  1501. db: AsyncSession,
  1502. printer_id: int | None = None,
  1503. printer_name: str | None = None,
  1504. ):
  1505. """Handle queue job added event."""
  1506. providers = await self._get_providers_for_event(db, "on_queue_job_added", printer_id)
  1507. if not providers:
  1508. return
  1509. variables = {
  1510. "job_name": job_name,
  1511. "target": target, # e.g., "Printer1" or "Any X1C"
  1512. "printer": printer_name or target,
  1513. }
  1514. title, message = await self._build_message_from_template(db, "queue_job_added", variables)
  1515. await self._send_to_providers(
  1516. providers, title, message, db, "queue_job_added", printer_id, printer_name, variables=variables
  1517. )
  1518. async def on_queue_job_assigned(
  1519. self,
  1520. job_name: str,
  1521. printer_id: int,
  1522. printer_name: str,
  1523. target_model: str,
  1524. db: AsyncSession,
  1525. ):
  1526. """Handle model-based job assigned to printer event."""
  1527. providers = await self._get_providers_for_event(db, "on_queue_job_assigned", printer_id)
  1528. if not providers:
  1529. return
  1530. variables = {
  1531. "job_name": job_name,
  1532. "printer": printer_name,
  1533. "target_model": target_model,
  1534. }
  1535. title, message = await self._build_message_from_template(db, "queue_job_assigned", variables)
  1536. await self._send_to_providers(
  1537. providers, title, message, db, "queue_job_assigned", printer_id, printer_name, variables=variables
  1538. )
  1539. async def on_queue_job_started(
  1540. self,
  1541. job_name: str,
  1542. printer_id: int,
  1543. printer_name: str,
  1544. db: AsyncSession,
  1545. estimated_time: int | None = None,
  1546. ):
  1547. """Handle queue job started printing event."""
  1548. providers = await self._get_providers_for_event(db, "on_queue_job_started", printer_id)
  1549. if not providers:
  1550. return
  1551. eta_str = await self._format_eta(estimated_time, db)
  1552. variables = {
  1553. "job_name": job_name,
  1554. "printer": printer_name,
  1555. "estimated_time": self._format_duration(estimated_time),
  1556. "eta": eta_str,
  1557. }
  1558. title, message = await self._build_message_from_template(db, "queue_job_started", variables)
  1559. await self._send_to_providers(
  1560. providers, title, message, db, "queue_job_started", printer_id, printer_name, variables=variables
  1561. )
  1562. async def on_queue_job_waiting(
  1563. self,
  1564. job_name: str,
  1565. target_model: str,
  1566. waiting_reason: str,
  1567. db: AsyncSession,
  1568. ):
  1569. """Handle job waiting for filament event."""
  1570. providers = await self._get_providers_for_event(db, "on_queue_job_waiting", None)
  1571. if not providers:
  1572. return
  1573. variables = {
  1574. "job_name": job_name,
  1575. "target_model": target_model,
  1576. "waiting_reason": waiting_reason,
  1577. }
  1578. title, message = await self._build_message_from_template(db, "queue_job_waiting", variables)
  1579. await self._send_to_providers(providers, title, message, db, "queue_job_waiting", variables=variables)
  1580. async def on_queue_job_skipped(
  1581. self,
  1582. job_name: str,
  1583. printer_id: int,
  1584. printer_name: str,
  1585. reason: str,
  1586. db: AsyncSession,
  1587. ):
  1588. """Handle job skipped event (e.g., previous print failed)."""
  1589. providers = await self._get_providers_for_event(db, "on_queue_job_skipped", printer_id)
  1590. if not providers:
  1591. return
  1592. variables = {
  1593. "job_name": job_name,
  1594. "printer": printer_name,
  1595. "reason": reason,
  1596. }
  1597. title, message = await self._build_message_from_template(db, "queue_job_skipped", variables)
  1598. await self._send_to_providers(
  1599. providers, title, message, db, "queue_job_skipped", printer_id, printer_name, variables=variables
  1600. )
  1601. async def on_queue_job_failed(
  1602. self,
  1603. job_name: str,
  1604. printer_id: int | None,
  1605. printer_name: str | None,
  1606. reason: str,
  1607. db: AsyncSession,
  1608. ):
  1609. """Handle job failed to start event (upload error, etc.)."""
  1610. providers = await self._get_providers_for_event(db, "on_queue_job_failed", printer_id)
  1611. if not providers:
  1612. return
  1613. variables = {
  1614. "job_name": job_name,
  1615. "printer": printer_name or "Unknown",
  1616. "reason": reason,
  1617. }
  1618. title, message = await self._build_message_from_template(db, "queue_job_failed", variables)
  1619. await self._send_to_providers(
  1620. providers, title, message, db, "queue_job_failed", printer_id, printer_name, variables=variables
  1621. )
  1622. async def on_queue_completed(
  1623. self,
  1624. completed_count: int,
  1625. db: AsyncSession,
  1626. ):
  1627. """Handle all queue jobs completed event."""
  1628. providers = await self._get_providers_for_event(db, "on_queue_completed", None)
  1629. if not providers:
  1630. return
  1631. variables = {
  1632. "completed_count": str(completed_count),
  1633. }
  1634. title, message = await self._build_message_from_template(db, "queue_completed", variables)
  1635. await self._send_to_providers(providers, title, message, db, "queue_completed", variables=variables)
  1636. # ==================== Inventory Stock Alerts ====================
  1637. async def on_stock_reorder_alert(
  1638. self,
  1639. material: str,
  1640. brand: str | None,
  1641. stock_g: float,
  1642. rate_g_day: float,
  1643. days_left: int,
  1644. db: AsyncSession,
  1645. ):
  1646. """Fire when an inventory SKU reaches its reorder point."""
  1647. providers = await self._get_providers_for_event(db, "on_stock_reorder_alert", None)
  1648. if not providers:
  1649. return
  1650. variables = {
  1651. "material": material,
  1652. "brand": brand or "",
  1653. "stock_g": f"{stock_g:.0f}",
  1654. "rate_g_day": f"{rate_g_day:.1f}",
  1655. "days_left": str(days_left),
  1656. }
  1657. title, message = await self._build_message_from_template(db, "stock_reorder_alert", variables)
  1658. await self._send_to_providers(providers, title, message, db, "stock_reorder_alert", variables=variables)
  1659. async def on_stock_break_alert(
  1660. self,
  1661. material: str,
  1662. brand: str | None,
  1663. stock_g: float,
  1664. rate_g_day: float,
  1665. days_left: int,
  1666. lead_time_days: int,
  1667. db: AsyncSession,
  1668. ):
  1669. """Fire when a stock break is detected (stock runs out before lead time)."""
  1670. providers = await self._get_providers_for_event(db, "on_stock_break_alert", None)
  1671. if not providers:
  1672. return
  1673. variables = {
  1674. "material": material,
  1675. "brand": brand or "",
  1676. "stock_g": f"{stock_g:.0f}",
  1677. "rate_g_day": f"{rate_g_day:.1f}",
  1678. "days_left": str(days_left),
  1679. "lead_time_days": str(lead_time_days),
  1680. }
  1681. title, message = await self._build_message_from_template(db, "stock_break_alert", variables)
  1682. await self._send_to_providers(providers, title, message, db, "stock_break_alert", variables=variables)
  1683. async def _queue_for_digest(
  1684. self,
  1685. provider: NotificationProvider,
  1686. event_type: str,
  1687. title: str,
  1688. message: str,
  1689. db: AsyncSession,
  1690. printer_id: int | None = None,
  1691. printer_name: str | None = None,
  1692. ):
  1693. """Queue a notification for later delivery in the daily digest."""
  1694. try:
  1695. queue_entry = NotificationDigestQueue(
  1696. provider_id=provider.id,
  1697. event_type=event_type,
  1698. title=title,
  1699. message=message,
  1700. printer_id=printer_id,
  1701. printer_name=printer_name,
  1702. )
  1703. db.add(queue_entry)
  1704. await db.commit()
  1705. logger.info("Queued notification for digest: %s for provider %s", event_type, provider.name)
  1706. except Exception as e:
  1707. logger.warning("Failed to queue notification for digest: %s", e)
  1708. async def send_digest(self, provider_id: int):
  1709. """Send all queued notifications as a single digest for a provider."""
  1710. from backend.app.core.database import async_session
  1711. async with async_session() as db:
  1712. # Get the provider
  1713. result = await db.execute(select(NotificationProvider).where(NotificationProvider.id == provider_id))
  1714. provider = result.scalar_one_or_none()
  1715. if not provider or not provider.enabled:
  1716. return
  1717. # Get all queued notifications for this provider
  1718. result = await db.execute(
  1719. select(NotificationDigestQueue)
  1720. .where(NotificationDigestQueue.provider_id == provider_id)
  1721. .order_by(NotificationDigestQueue.created_at)
  1722. )
  1723. queue_entries = list(result.scalars().all())
  1724. if not queue_entries:
  1725. logger.debug("No queued notifications for provider %s", provider.name)
  1726. return
  1727. # Build digest message
  1728. title = f"Daily Digest - {len(queue_entries)} Events"
  1729. # Group by event type
  1730. events_by_type: dict[str, list] = {}
  1731. for entry in queue_entries:
  1732. if entry.event_type not in events_by_type:
  1733. events_by_type[entry.event_type] = []
  1734. events_by_type[entry.event_type].append(entry)
  1735. # Format the digest body
  1736. body_parts = []
  1737. for event_type, entries in events_by_type.items():
  1738. event_label = event_type.replace("_", " ").title()
  1739. body_parts.append(f"== {event_label} ({len(entries)}) ==")
  1740. for entry in entries:
  1741. time_str = entry.created_at.strftime("%H:%M")
  1742. printer_info = f"[{entry.printer_name}] " if entry.printer_name else ""
  1743. body_parts.append(f" {time_str} {printer_info}{entry.title}")
  1744. body_parts.append("")
  1745. body = "\n".join(body_parts)
  1746. # Send the digest
  1747. success, error = await self._send_to_provider(provider, title, body, db)
  1748. # Log the digest
  1749. await self._log_notification(
  1750. db=db,
  1751. provider_id=provider.id,
  1752. event_type="daily_digest",
  1753. title=title,
  1754. message=body,
  1755. success=success,
  1756. error_message=error if not success else None,
  1757. )
  1758. # Clear the queue
  1759. for entry in queue_entries:
  1760. await db.delete(entry)
  1761. await db.commit()
  1762. if success:
  1763. logger.info("Sent daily digest with %s events to %s", len(queue_entries), provider.name)
  1764. else:
  1765. logger.warning("Failed to send daily digest to %s: %s", provider.name, error)
  1766. async def check_and_send_digests(self):
  1767. """Check all providers and send digests if it's their scheduled time."""
  1768. from backend.app.core.database import async_session
  1769. current_time = datetime.now().strftime("%H:%M")
  1770. # Avoid duplicate checks within the same minute
  1771. if current_time == self._last_digest_check:
  1772. return
  1773. self._last_digest_check = current_time
  1774. async with async_session() as db:
  1775. # Find all providers with digest enabled at this time
  1776. result = await db.execute(
  1777. select(NotificationProvider).where(
  1778. NotificationProvider.enabled.is_(True),
  1779. NotificationProvider.daily_digest_enabled.is_(True),
  1780. NotificationProvider.daily_digest_time == current_time,
  1781. )
  1782. )
  1783. providers = result.scalars().all()
  1784. for provider in providers:
  1785. try:
  1786. await self.send_digest(provider.id)
  1787. except Exception as e:
  1788. logger.error("Error sending digest for provider %s: %s", provider.id, e)
  1789. def start_digest_scheduler(self):
  1790. """Start the background scheduler for daily digest notifications."""
  1791. if self._digest_scheduler_task is None:
  1792. self._digest_scheduler_task = asyncio.create_task(self._digest_scheduler_loop())
  1793. logger.info("Notification digest scheduler started")
  1794. def stop_digest_scheduler(self):
  1795. """Stop the background scheduler for daily digests."""
  1796. if self._digest_scheduler_task:
  1797. self._digest_scheduler_task.cancel()
  1798. self._digest_scheduler_task = None
  1799. logger.info("Notification digest scheduler stopped")
  1800. async def _digest_scheduler_loop(self):
  1801. """Background loop that checks for scheduled digests every minute."""
  1802. while True:
  1803. try:
  1804. await self.check_and_send_digests()
  1805. except Exception as e:
  1806. logger.error("Error in digest scheduler: %s", e)
  1807. # Wait until the next minute
  1808. await asyncio.sleep(60)
  1809. # Global instance
  1810. notification_service = NotificationService()