notification_service.py 113 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079208020812082208320842085208620872088208920902091209220932094209520962097209820992100210121022103210421052106210721082109211021112112211321142115211621172118211921202121212221232124212521262127212821292130213121322133213421352136213721382139214021412142214321442145214621472148214921502151215221532154215521562157215821592160216121622163216421652166216721682169217021712172217321742175217621772178217921802181218221832184218521862187218821892190219121922193219421952196219721982199220022012202220322042205220622072208220922102211221222132214221522162217221822192220222122222223222422252226222722282229223022312232223322342235223622372238223922402241224222432244224522462247224822492250225122522253225422552256225722582259226022612262226322642265226622672268226922702271227222732274227522762277227822792280228122822283228422852286228722882289229022912292229322942295229622972298229923002301230223032304230523062307230823092310231123122313231423152316231723182319232023212322232323242325232623272328232923302331233223332334233523362337233823392340234123422343234423452346234723482349235023512352235323542355235623572358235923602361236223632364236523662367236823692370237123722373237423752376237723782379238023812382238323842385238623872388238923902391239223932394239523962397239823992400240124022403240424052406240724082409241024112412241324142415241624172418241924202421242224232424242524262427242824292430243124322433243424352436243724382439244024412442244324442445244624472448244924502451245224532454245524562457245824592460246124622463246424652466246724682469247024712472247324742475247624772478247924802481248224832484248524862487248824892490249124922493249424952496249724982499250025012502250325042505250625072508250925102511251225132514251525162517251825192520252125222523252425252526252725282529253025312532253325342535253625372538253925402541254225432544254525462547254825492550255125522553255425552556255725582559256025612562256325642565256625672568256925702571257225732574257525762577257825792580258125822583258425852586258725882589259025912592259325942595259625972598259926002601260226032604260526062607260826092610261126122613261426152616261726182619262026212622262326242625262626272628262926302631263226332634263526362637263826392640264126422643264426452646264726482649265026512652265326542655265626572658265926602661266226632664266526662667266826692670267126722673267426752676267726782679268026812682268326842685268626872688268926902691269226932694269526962697269826992700270127022703270427052706270727082709271027112712271327142715271627172718271927202721272227232724272527262727272827292730273127322733273427352736273727382739274027412742274327442745274627472748274927502751275227532754275527562757275827592760276127622763276427652766276727682769277027712772277327742775277627772778277927802781278227832784278527862787278827892790279127922793279427952796279727982799280028012802280328042805280628072808280928102811
  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.archive import PrintArchive
  18. from backend.app.models.notification import (
  19. NotificationDigestQueue,
  20. NotificationLog,
  21. NotificationProvider,
  22. TelegramPendingVerdict,
  23. )
  24. from backend.app.models.notification_template import NotificationTemplate
  25. from backend.app.services.print_confirmation import one_tap_url
  26. logger = logging.getLogger(__name__)
  27. # Honest User-Agent — matches the convention used by every other outbound
  28. # httpx client in the codebase (bambu_cloud, makerworld, firmware_check,
  29. # inventory). Previously this client leaked python-httpx/<version>, which
  30. # was both inconsistent with the rest of the project and a more obvious
  31. # bot signature for upstream WAFs.
  32. _USER_AGENT = "Bambuddy/1.0 (+https://github.com/maziggy/bambuddy)"
  33. # Appended to a Telegram outcome prompt in reaction mode (#3046); the message
  34. # itself carries no buttons there, so the user needs telling how to answer.
  35. TELEGRAM_REACTION_HINT = "React with \U0001f44d or \U0001f44e to record the outcome."
  36. def telegram_markdown_escape(message: str) -> str:
  37. """Escape underscores in the message body so Telegram Markdown parsing
  38. doesn't break on job names like "A1_plate_8" or error codes like
  39. "0300_0001". The title is already wrapped in *bold* markers, so only
  40. escape after the first newline. Shared with the reaction poller, which
  41. re-sends the stored body when it edits the prompt (#3046)."""
  42. if "\n" in message:
  43. title_part, body_part = message.split("\n", 1)
  44. body_part = body_part.replace("_", "\\_")
  45. return f"{title_part}\n{body_part}"
  46. return message
  47. def _looks_like_cloudflare_challenge(response: httpx.Response) -> bool:
  48. """Return True if ``response`` looks like a Cloudflare mitigation
  49. interstitial (JS challenge / managed challenge / block page) rather
  50. than a legitimate response passed through Cloudflare.
  51. Self-hosted servers behind Cloudflare (Tunnel, "Bot Fight Mode", or
  52. "Under Attack" mode) intercept non-browser clients at the edge and
  53. return a challenge HTML page instead of forwarding to the origin —
  54. so we never reach the user's actual ntfy / webhook backend.
  55. Cloudflare cannot be defeated from a Python client; the user has to
  56. add a security-skip rule on their side. We detect the shape so the
  57. UI can tell them that, instead of dumping the raw HTML.
  58. Detection deliberately does NOT rely on ``Server: cloudflare`` alone
  59. — Cloudflare adds that header to every response it proxies (success
  60. AND legitimate origin errors), so a real 401 "wrong token" from a
  61. CF-fronted ntfy would false-positive into a misleading "your CF is
  62. blocking" message. Reliable signals: the ``cf-mitigated`` header
  63. (set only when CF actively mitigates) and the challenge body shape.
  64. """
  65. if response.headers.get("cf-mitigated"):
  66. return True
  67. content_type = (response.headers.get("content-type") or "").lower()
  68. if "html" not in content_type:
  69. return False
  70. body = (response.text or "")[:1024].lower()
  71. # "Just a moment..." is Cloudflare's universal challenge-page title
  72. # (managed challenge, JS challenge, Under Attack mode). Combined with
  73. # an HTML content-type this is unambiguous — no legitimate ntfy or
  74. # webhook backend returns HTML with that title. ``cf-chl-*`` and
  75. # ``challenge-platform`` cover newer / non-default CF templates.
  76. return "just a moment" in body or "cf-chl-bypass" in body or "cf-chl-opt" in body or "challenge-platform" in body
  77. def _assert_safe_provider_url(url: str, *, label: str) -> str | None:
  78. """Validate a provider URL taken from user-supplied config.
  79. Returns an error message on rejection, or None when the URL is
  80. acceptable — the ``_send_*`` methods return ``tuple[bool, str]`` rather
  81. than raising, so a message is more useful here than an exception.
  82. Uses the LAN-service policy: self-hosting ntfy, Bark, Gotify or a webhook
  83. receiver on the home LAN is normal and must keep working, so loopback and
  84. RFC-1918 stay permitted. Cloud-metadata endpoints, numeric-encoded IPs and
  85. non-HTTP schemes are rejected.
  86. """
  87. from backend.app.api.routes._url_safety import assert_safe_lan_service_url
  88. try:
  89. assert_safe_lan_service_url(url, label=label)
  90. except ValueError as exc:
  91. return str(exc)
  92. return None
  93. def _opaque_http_failure(response: httpx.Response, *, label: str) -> str:
  94. """Failure message for a provider whose destination host the user supplies.
  95. The response body is deliberately **not** returned to the caller. Provider
  96. URLs are configurable by anyone holding ``NOTIFICATIONS_CREATE`` — which
  97. the default Operators group carries and which does not imply
  98. ``SETTINGS_UPDATE`` — and ``POST /notifications/test-config`` accepts a URL
  99. straight from the request body without persisting anything. Echoing the
  100. response body there turned an intended "does my webhook work?" check into
  101. an authenticated read primitive against any host the Bambuddy process can
  102. reach, including services that are not exposed to the network at all.
  103. Providers whose host Bambuddy hardcodes (Pushover, Telegram, CallMeBot)
  104. keep returning the upstream body — there is no trust boundary to cross
  105. when the destination cannot be influenced.
  106. The body is logged at debug level, where it stays available to whoever
  107. already administers the host without being handed back over the API.
  108. """
  109. logger.debug(
  110. "%s delivery failed with HTTP %s; body: %s",
  111. label,
  112. response.status_code,
  113. (response.text or "")[:200],
  114. )
  115. return f"HTTP {response.status_code} from the configured {label} (see server logs at debug level for details)"
  116. class NotificationService:
  117. """Service for sending notifications through various providers."""
  118. def __init__(self):
  119. self._http_client: httpx.AsyncClient | None = None
  120. self._template_cache: dict[str, NotificationTemplate] = {}
  121. self._digest_scheduler_task: asyncio.Task | None = None
  122. self._last_digest_check: str = "" # "HH:MM" to avoid duplicate checks
  123. async def _get_client(self) -> httpx.AsyncClient:
  124. """Get or create HTTP client.
  125. The connect timeout is deliberately far shorter than the rest. A flat
  126. 30 s meant that when a site's internet went down, every alarm spent a
  127. full 30 s inside ``connect`` — longer than SQLite's 15 s
  128. ``busy_timeout`` — and any other task that wanted to write during that
  129. window failed with "database is locked" (#2770). Reaching a host either
  130. works in a couple of seconds or is not going to; sending the body is the
  131. part that legitimately takes time, so read/write keep the old 30 s and
  132. an image upload on a slow uplink is unaffected.
  133. """
  134. if self._http_client is None or self._http_client.is_closed:
  135. self._http_client = httpx.AsyncClient(
  136. timeout=httpx.Timeout(30.0, connect=5.0),
  137. headers={"User-Agent": _USER_AGENT},
  138. )
  139. return self._http_client
  140. async def close(self):
  141. """Close HTTP client."""
  142. if self._http_client and not self._http_client.is_closed:
  143. await self._http_client.aclose()
  144. def _is_in_quiet_hours(self, provider: NotificationProvider) -> bool:
  145. """Check if current time is within provider's quiet hours."""
  146. if not provider.quiet_hours_enabled:
  147. return False
  148. if not provider.quiet_hours_start or not provider.quiet_hours_end:
  149. return False
  150. try:
  151. now = datetime.now()
  152. current_time = now.hour * 60 + now.minute
  153. start_parts = provider.quiet_hours_start.split(":")
  154. end_parts = provider.quiet_hours_end.split(":")
  155. start_minutes = int(start_parts[0]) * 60 + int(start_parts[1])
  156. end_minutes = int(end_parts[0]) * 60 + int(end_parts[1])
  157. # Handle overnight quiet hours (e.g., 22:00 to 07:00)
  158. if start_minutes > end_minutes:
  159. # Quiet hours span midnight
  160. return current_time >= start_minutes or current_time < end_minutes
  161. else:
  162. # Same day quiet hours
  163. return start_minutes <= current_time < end_minutes
  164. except (ValueError, TypeError, AttributeError):
  165. logger.warning("Invalid quiet hours format for provider %s", provider.name)
  166. return False
  167. async def _get_template(self, db: AsyncSession, event_type: str) -> NotificationTemplate | None:
  168. """Get a notification template by event type.
  169. ``no_autoflush`` for the same reason as ``_get_providers_for_event``:
  170. this read runs before the provider is contacted, and must not be the
  171. thing that opens a write transaction on the caller's session (#2770).
  172. """
  173. # Check cache first
  174. if event_type in self._template_cache:
  175. return self._template_cache[event_type]
  176. with db.no_autoflush:
  177. result = await db.execute(select(NotificationTemplate).where(NotificationTemplate.event_type == event_type))
  178. template = result.scalar_one_or_none()
  179. if template:
  180. self._template_cache[event_type] = template
  181. return template
  182. def _render_template(self, template_str: str, variables: dict[str, Any]) -> str:
  183. """Render a template string with variables. Missing variables become empty."""
  184. result = template_str
  185. for key, value in variables.items():
  186. result = result.replace("{" + key + "}", str(value) if value is not None else "")
  187. # Remove any remaining unreplaced placeholders
  188. result = re.sub(r"\{[a-z_]+\}", "", result)
  189. return result
  190. async def _format_eta(self, seconds: int | None, db: AsyncSession) -> str:
  191. """Format ETA as wall-clock time, respecting user's time_format setting."""
  192. if not seconds or seconds <= 0:
  193. return "Unknown"
  194. from backend.app.api.routes.settings import get_setting
  195. eta_time = datetime.now() + timedelta(seconds=seconds)
  196. time_format = await get_setting(db, "time_format")
  197. if time_format == "12h":
  198. return eta_time.strftime("%I:%M %p").lstrip("0")
  199. # Default to 24h for "24h", "system", or unset
  200. return eta_time.strftime("%H:%M")
  201. def _format_duration(self, seconds: int | None) -> str:
  202. """Format duration in seconds to human-readable string."""
  203. if seconds is None:
  204. return "Unknown"
  205. hours = seconds // 3600
  206. minutes = (seconds % 3600) // 60
  207. if hours > 0:
  208. return f"{hours}h {minutes}m"
  209. return f"{minutes}m"
  210. def _clean_filename(self, filename: str) -> str:
  211. """Extract filename and remove file extensions."""
  212. import os
  213. # Strip path prefix (e.g., /data/Metadata/plate_5.gcode -> plate_5.gcode)
  214. filename = os.path.basename(filename)
  215. # Remove common extensions
  216. if filename.endswith(".gcode.3mf"):
  217. return filename[:-10]
  218. elif filename.endswith(".gcode"):
  219. return filename[:-6]
  220. elif filename.endswith(".3mf"):
  221. return filename[:-4]
  222. return filename
  223. async def _build_message_from_template(
  224. self, db: AsyncSession, event_type: str, variables: dict[str, Any]
  225. ) -> tuple[str, str]:
  226. """Build notification title and body from template."""
  227. # Add common variables
  228. variables["timestamp"] = datetime.now().strftime("%Y-%m-%d %H:%M")
  229. variables["app_name"] = "Bambuddy"
  230. template = await self._get_template(db, event_type)
  231. if not template:
  232. # Fallback to simple message
  233. logger.warning("Template not found for event type: %s", event_type)
  234. return event_type.replace("_", " ").title(), str(variables)
  235. title = self._render_template(template.title_template, variables)
  236. body = self._render_template(template.body_template, variables)
  237. return title, body
  238. async def send_test_notification(
  239. self, provider_type: str, config: dict[str, Any], db: AsyncSession | None = None
  240. ) -> tuple[bool, str]:
  241. """Send a test notification to verify configuration."""
  242. if db:
  243. title, message = await self._build_message_from_template(db, "test", {})
  244. else:
  245. title = "Bambuddy Test"
  246. message = "This is a test notification. If you see this, notifications are working!"
  247. try:
  248. if provider_type == "callmebot":
  249. return await self._send_callmebot(config, f"{title}\n{message}")
  250. elif provider_type == "ntfy":
  251. return await self._send_ntfy(config, title, message)
  252. elif provider_type == "pushover":
  253. return await self._send_pushover(config, title, message)
  254. elif provider_type == "telegram":
  255. return await self._send_telegram(config, f"*{title}*\n{message}")
  256. elif provider_type == "email":
  257. return await self._send_email(config, title, message)
  258. elif provider_type == "discord":
  259. return await self._send_discord(config, title, message)
  260. elif provider_type == "webhook":
  261. return await self._send_webhook(config, title, message)
  262. elif provider_type == "homeassistant":
  263. return await self._send_homeassistant(config, title, message, db=db)
  264. elif provider_type == "bark":
  265. return await self._send_bark(config, title, message)
  266. else:
  267. return False, f"Unknown provider type: {provider_type}"
  268. except Exception as e:
  269. logger.exception("Error sending test notification via %s", provider_type)
  270. return False, str(e)
  271. async def _send_callmebot(self, config: dict, message: str) -> tuple[bool, str]:
  272. """Send notification via CallMeBot (WhatsApp)."""
  273. phone = config.get("phone", "").strip()
  274. apikey = config.get("apikey", "").strip()
  275. if not phone or not apikey:
  276. return False, "Phone number and API key are required"
  277. # URL encode the message
  278. encoded_message = quote(message)
  279. url = f"https://api.callmebot.com/whatsapp.php?phone={phone}&text={encoded_message}&apikey={apikey}"
  280. client = await self._get_client()
  281. response = await client.get(url)
  282. if response.status_code == 200:
  283. return True, "Message sent successfully"
  284. else:
  285. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  286. async def _send_bark(self, config: dict, title: str, message: str, url: str | None = None) -> tuple[bool, str]:
  287. """Send notification via Bark, the self-hostable iOS push service (#1495).
  288. POSTs JSON to {server}/push. Defaults to the official api.day.app
  289. relay; a self-hosted bark-server works by overriding the server URL.
  290. ``url`` opens on tap — the outcome confirmation (#1898) deep-links
  291. into the archive's confirmation dialog with it.
  292. """
  293. server = (config.get("server") or "https://api.day.app").strip().rstrip("/")
  294. device_key = (config.get("device_key") or "").strip()
  295. if not device_key:
  296. return False, "Device key is required"
  297. url_error = _assert_safe_provider_url(server, label="Bark server URL")
  298. if url_error:
  299. return False, url_error
  300. payload: dict[str, Any] = {
  301. "device_key": device_key,
  302. "title": title,
  303. "body": message,
  304. }
  305. group = (config.get("group") or "").strip()
  306. if group:
  307. payload["group"] = group
  308. sound = (config.get("sound") or "").strip()
  309. if sound:
  310. payload["sound"] = sound
  311. level = (config.get("level") or "").strip()
  312. if level in ("active", "timeSensitive", "critical", "passive"):
  313. payload["level"] = level
  314. if url:
  315. payload["url"] = url
  316. client = await self._get_client()
  317. response = await client.post(f"{server}/push", json=payload)
  318. if response.status_code == 200:
  319. # bark-server can report failures inside an HTTP 200 body
  320. # ({"code": 400, "message": ...}), so the status alone isn't proof.
  321. try:
  322. body = response.json()
  323. except ValueError:
  324. body = None
  325. if isinstance(body, dict) and body.get("code") not in (200, None):
  326. # Only the numeric code is echoed. A server chosen by the caller
  327. # controls this body too, so the free-text message is a (narrow)
  328. # read channel of the same kind _opaque_http_failure closes.
  329. logger.debug("Bark reported error %s: %s", body.get("code"), str(body.get("message"))[:200])
  330. return False, f"Bark error {body.get('code')} (see server logs at debug level for details)"
  331. return True, "Message sent successfully"
  332. return False, _opaque_http_failure(response, label="Bark server")
  333. async def _send_ntfy(
  334. self,
  335. config: dict,
  336. title: str,
  337. message: str,
  338. image_data: bytes | None = None,
  339. event_type: str | None = None,
  340. actions: str | None = None,
  341. ) -> tuple[bool, str]:
  342. """Send notification via ntfy.
  343. ``actions`` is a pre-built value for ntfy's Actions header (simple
  344. format), used by the outcome-confirmation event (#1898) to put
  345. one-tap Good/Reject buttons directly into the push notification.
  346. """
  347. server = config.get("server", "https://ntfy.sh").rstrip("/")
  348. topic = config.get("topic", "").strip()
  349. auth_token = config.get("auth_token", "").strip()
  350. if not topic:
  351. return False, "Topic is required"
  352. url_error = _assert_safe_provider_url(server, label="ntfy server URL")
  353. if url_error:
  354. return False, url_error
  355. url = f"{server}/{topic}"
  356. # ntfy reads Title/Message from HTTP headers. httpx enforces ASCII
  357. # for str header values, but printer names and filenames can contain
  358. # non-ASCII characters (e.g. accented letters, CJK). Passing bytes
  359. # bypasses the ASCII check — ntfy handles UTF-8 headers correctly.
  360. headers: dict[str, str | bytes] = {"Title": title.encode("utf-8")}
  361. # Per-event Priority header (#990). Only set when the user has
  362. # explicitly mapped this event to a 1-5 value; otherwise fall through
  363. # to the ntfy server's default so existing setups stay unchanged.
  364. #
  365. # The map is keyed by the provider's toggle column ("on_print_failed"),
  366. # because that is what the dialog builds its rows from -- but every
  367. # sender is called with the bare event name ("print_failed"), so the
  368. # lookup used to miss for every real notification and hit only in tests
  369. # that called this method with the prefixed name (issue #3139). Both
  370. # spellings are accepted, which also leaves stored configs untouched.
  371. event_priorities = config.get("event_priorities") or {}
  372. if event_type and isinstance(event_priorities, dict):
  373. raw = event_priorities.get(event_type)
  374. if raw is None and not event_type.startswith("on_"):
  375. raw = event_priorities.get(f"on_{event_type}")
  376. try:
  377. priority = int(raw) if raw is not None else None
  378. except (TypeError, ValueError):
  379. priority = None
  380. if priority is not None and 1 <= priority <= 5:
  381. headers["Priority"] = str(priority)
  382. if auth_token:
  383. headers["Authorization"] = f"Bearer {auth_token}"
  384. if actions:
  385. headers["Actions"] = actions
  386. client = await self._get_client()
  387. if image_data:
  388. # ntfy supports image attachments via multipart form-data.
  389. # HTTP headers cannot contain newlines, but ntfy interprets
  390. # literal \n (backslash-n) as newlines in the Message header.
  391. headers["Filename"] = "photo.jpg"
  392. headers["Message"] = message.replace("\n", "\\n").encode("utf-8")
  393. response = await client.put(url, content=image_data, headers=headers)
  394. if response.status_code == 400 and "attachments not allowed" in response.text:
  395. # Server has attachments disabled — retry without the image
  396. headers.pop("Filename", None)
  397. headers.pop("Message", None)
  398. response = await client.post(url, content=message.encode("utf-8"), headers=headers)
  399. else:
  400. response = await client.post(url, content=message.encode("utf-8"), headers=headers)
  401. if response.status_code in (200, 204):
  402. return True, "Message sent successfully"
  403. if _looks_like_cloudflare_challenge(response):
  404. return False, (
  405. f"HTTP {response.status_code} — ntfy server is behind a Cloudflare "
  406. "challenge. Bambuddy was served the JS challenge page instead of "
  407. "reaching ntfy. Cloudflare cannot be solved from a backend; add a "
  408. "Cloudflare security-skip rule for this hostname, disable Bot "
  409. "Fight Mode, or front the server with Cloudflare Access using a "
  410. "service token. (#1534)"
  411. )
  412. return False, _opaque_http_failure(response, label="ntfy server")
  413. async def _send_pushover(
  414. self,
  415. config: dict,
  416. title: str,
  417. message: str,
  418. image_data: bytes | None = None,
  419. url: str | None = None,
  420. url_title: str | None = None,
  421. ) -> tuple[bool, str]:
  422. """Send notification via Pushover.
  423. Args:
  424. config: Provider configuration with user_key, app_token, priority
  425. title: Notification title
  426. message: Notification body
  427. image_data: Optional JPEG image bytes to attach (max 2.5MB)
  428. url: Optional supplementary URL shown under the message
  429. url_title: Optional label for that URL
  430. """
  431. user_key = config.get("user_key", "").strip()
  432. app_token = config.get("app_token", "").strip()
  433. try:
  434. priority = int(config.get("priority", 0))
  435. except (TypeError, ValueError):
  436. priority = 0
  437. if not user_key or not app_token:
  438. return False, "User key and app token are required"
  439. api_url = "https://api.pushover.net/1/messages.json"
  440. data = {
  441. "token": app_token,
  442. "user": user_key,
  443. "title": title,
  444. "message": message,
  445. "priority": priority,
  446. }
  447. if url:
  448. data["url"] = url
  449. if url_title:
  450. data["url_title"] = url_title
  451. # Emergency priority (2) keeps re-alerting until acknowledged, so
  452. # Pushover *requires* retry (how often, >= 30s) and expire (when to
  453. # give up, <= 10800s). Without them the API rejects the message. Only
  454. # send them at priority 2 — Pushover ignores them at other priorities.
  455. if priority == 2:
  456. try:
  457. retry = int(config.get("retry", 60))
  458. except (TypeError, ValueError):
  459. retry = 60
  460. try:
  461. expire = int(config.get("expire", 3600))
  462. except (TypeError, ValueError):
  463. expire = 3600
  464. data["retry"] = max(30, min(retry, 10800))
  465. data["expire"] = max(30, min(expire, 10800))
  466. client = await self._get_client()
  467. if image_data:
  468. # Pushover supports image attachments via multipart form-data
  469. files = {"attachment": ("photo.jpg", image_data, "image/jpeg")}
  470. response = await client.post(api_url, data=data, files=files)
  471. else:
  472. response = await client.post(api_url, data=data)
  473. if response.status_code == 200:
  474. return True, "Message sent successfully"
  475. else:
  476. try:
  477. error_data = response.json()
  478. errors = error_data.get("errors", [])
  479. return False, f"Pushover error: {', '.join(errors)}"
  480. except Exception:
  481. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  482. async def _send_telegram(
  483. self,
  484. config: dict,
  485. message: str,
  486. image_data: bytes | None = None,
  487. buttons: list[dict] | None = None,
  488. link_preview: bool = True,
  489. ) -> tuple[bool, str]:
  490. """Send notification via Telegram bot.
  491. ``buttons`` is one row of inline URL buttons (``{"text", "url"}``
  492. entries), used by the outcome-confirmation event (#1898) to put
  493. one-tap Good/Reject under the message. ``link_preview=False`` asks
  494. Telegram not to fetch the first URL in the text for a preview card.
  495. """
  496. ok, status, _ = await self._send_telegram_message(
  497. config, message, image_data=image_data, buttons=buttons, link_preview=link_preview
  498. )
  499. return ok, status
  500. async def _send_telegram_message(
  501. self,
  502. config: dict,
  503. message: str,
  504. image_data: bytes | None = None,
  505. buttons: list[dict] | None = None,
  506. link_preview: bool = True,
  507. ) -> tuple[bool, str, dict | None]:
  508. """Send via Telegram and also hand back the Bot API ``result`` (the sent Message).
  509. The outcome confirmation in reaction mode (#3046) needs the
  510. ``message_id`` from it to recognise the reaction later; every other
  511. caller goes through ``_send_telegram`` and ignores it.
  512. """
  513. bot_token = config.get("bot_token", "").strip()
  514. chat_id = config.get("chat_id", "").strip()
  515. if not bot_token or not chat_id:
  516. return False, "Bot token and chat ID are required", None
  517. # Optional forum topic (#1518). Telegram expects message_thread_id as an
  518. # integer in the JSON sendMessage body — a string 400s there even though
  519. # the multipart sendPhoto call below would happily accept one. Coerce it
  520. # once, up front, so both call sites agree and a bad value fails loudly
  521. # instead of silently breaking only the text notifications.
  522. thread_id_raw = str(config.get("message_thread_id") or "").strip()
  523. message_thread_id: int | None = None
  524. if thread_id_raw:
  525. try:
  526. message_thread_id = int(thread_id_raw)
  527. except ValueError:
  528. return False, f"Invalid message thread ID: {thread_id_raw!r} is not a number", None
  529. message = telegram_markdown_escape(message)
  530. client = await self._get_client()
  531. async def _post(with_buttons: bool):
  532. if image_data:
  533. # Use sendPhoto to attach the thumbnail with the caption
  534. url = f"https://api.telegram.org/bot{bot_token}/sendPhoto"
  535. form: dict[str, Any] = {"chat_id": chat_id, "caption": message, "parse_mode": "Markdown"}
  536. if message_thread_id is not None:
  537. form["message_thread_id"] = message_thread_id
  538. if with_buttons:
  539. # Multipart form fields are strings — reply_markup goes JSON-encoded.
  540. form["reply_markup"] = json.dumps({"inline_keyboard": [buttons]})
  541. return await client.post(
  542. url,
  543. data=form,
  544. files={"photo": ("photo.jpg", image_data, "image/jpeg")},
  545. )
  546. url = f"https://api.telegram.org/bot{bot_token}/sendMessage"
  547. payload: dict[str, Any] = {
  548. "chat_id": chat_id,
  549. "text": message,
  550. "parse_mode": "Markdown",
  551. }
  552. if not link_preview:
  553. payload["disable_web_page_preview"] = True
  554. if message_thread_id is not None:
  555. payload["message_thread_id"] = message_thread_id
  556. if with_buttons:
  557. payload["reply_markup"] = {"inline_keyboard": [buttons]}
  558. return await client.post(url, json=payload)
  559. def _failure(resp) -> str | None:
  560. """What Telegram objected to, or None when the send went through."""
  561. if resp.status_code != 200:
  562. return f"HTTP {resp.status_code}: {resp.text[:200]}"
  563. result = resp.json()
  564. if result.get("ok"):
  565. return None
  566. return f"Telegram error: {result.get('description', 'Unknown error')}"
  567. response = await _post(bool(buttons))
  568. failure = _failure(response)
  569. if failure and buttons:
  570. # Telegram validates every inline-keyboard URL and refuses the whole
  571. # send when one of them is not a URL it accepts — which is what an
  572. # install without a public external_url produces for the #1898
  573. # verdict links. Dropping the buttons is survivable; dropping the
  574. # message the user is waiting for is not.
  575. logger.warning("Telegram refused the message with inline buttons (%s); retrying without them", failure)
  576. response = await _post(False)
  577. failure = _failure(response)
  578. if failure:
  579. return False, failure, None
  580. sent = response.json().get("result")
  581. return True, "Message sent successfully", sent if isinstance(sent, dict) else None
  582. async def _send_telegram_confirm_request(
  583. self,
  584. provider: NotificationProvider,
  585. config: dict,
  586. message: str,
  587. db: AsyncSession | None,
  588. image_data: bytes | None,
  589. buttons: list[dict] | None,
  590. archive_id: int | None,
  591. ) -> tuple[bool, str]:
  592. """Deliver the outcome prompt the way the provider's verdict mode asks (#3046).
  593. "buttons" is the plain #1898 delivery. In "reactions" the inline
  594. keyboard is dropped and the user answers with a thumbs-up/down on the
  595. message itself; "both" keeps the keyboard as well. In either of those
  596. the sent message is remembered in telegram_pending_verdicts, together
  597. with the confirm token its links carry, so the reaction poller can map
  598. the reaction back to the archive for as long as that token is live.
  599. """
  600. mode = provider.telegram_verdict_mode or "buttons"
  601. # Telegram's servers GET the first URL in the text to build a preview
  602. # card. An outcome prompt whose edited body still carries {good_url}
  603. # would have that fetch answer the question before the operator saw
  604. # it, so the preview is off for the prompt in every mode.
  605. if mode == "buttons":
  606. return await self._send_telegram(
  607. config, message, image_data=image_data, buttons=buttons, link_preview=False
  608. )
  609. if mode == "reactions":
  610. buttons = None
  611. message = f"{message}\n\n{TELEGRAM_REACTION_HINT}"
  612. ok, status, sent = await self._send_telegram_message(
  613. config, message, image_data=image_data, buttons=buttons, link_preview=False
  614. )
  615. if not ok:
  616. return ok, status
  617. message_id = (sent or {}).get("message_id")
  618. if not isinstance(message_id, int) or message_id <= 0:
  619. logger.warning("Telegram did not return a message_id for the outcome prompt; reactions cannot be matched")
  620. return ok, status
  621. if db is None or archive_id is None:
  622. return ok, status
  623. try:
  624. # dispatch_outcome_confirmation mints a live token right before
  625. # sending. Without one the prompt's links are dead too, and a
  626. # reaction must not outlive them, so there is nothing to remember.
  627. archive = await db.get(PrintArchive, archive_id)
  628. if archive is None or not archive.confirm_token or archive.confirm_token_used_at is not None:
  629. logger.warning("Outcome prompt for archive %s went out without a live confirm token", archive_id)
  630. return ok, status
  631. db.add(
  632. TelegramPendingVerdict(
  633. provider_id=provider.id,
  634. chat_id=str(((sent or {}).get("chat") or {}).get("id") or config.get("chat_id", "")).strip(),
  635. message_id=message_id,
  636. archive_id=archive_id,
  637. confirm_token=archive.confirm_token,
  638. has_caption=image_data is not None,
  639. message_text=message,
  640. )
  641. )
  642. await db.commit()
  643. except Exception as e:
  644. # The prompt went out; losing the reaction mapping is a degraded
  645. # outcome, not a failed notification.
  646. logger.warning("Failed to record the Telegram outcome prompt for reactions: %s", e)
  647. await db.rollback()
  648. return ok, status
  649. async def _send_email(
  650. self,
  651. config: dict,
  652. subject: str,
  653. body: str,
  654. image_data: bytes | None = None,
  655. finish_photo_url: str | None = None,
  656. ) -> tuple[bool, str]:
  657. """Send notification via email (SMTP).
  658. Inline finish-photo embed is opt-in via the template: when the rendered
  659. ``body`` contains the substituted ``{finish_photo_url}`` value AND the
  660. finish-photo bytes are present, the message is built as
  661. ``multipart/related`` wrapping a ``multipart/alternative`` (plain + HTML)
  662. plus an inline ``MIMEImage`` with ``Content-ID: <bambuddy-finish-photo>``.
  663. The HTML part replaces the URL with ``<img src="cid:...">``; the plain-
  664. text part keeps the URL as a clickable link. When the template doesn't
  665. reference ``{finish_photo_url}`` (or image bytes aren't available), the
  666. original single-part text shape is used — no attachment, no surprise
  667. inline image (#1792).
  668. """
  669. smtp_server = config.get("smtp_server", "").strip()
  670. smtp_port = int(config.get("smtp_port", 587))
  671. username = config.get("username", "").strip()
  672. password = config.get("password", "").strip()
  673. from_email = config.get("from_email", "").strip()
  674. to_email = config.get("to_email", "").strip()
  675. # Security: "starttls" (port 587), "ssl" (port 465), "none" (port 25)
  676. security = config.get("security", "starttls")
  677. # Authentication: "true" or "false"
  678. auth_enabled = config.get("auth_enabled", "true").lower() == "true"
  679. if not all([smtp_server, from_email, to_email]):
  680. return False, "SMTP server, from email, and to email are required"
  681. if auth_enabled and not all([username, password]):
  682. return False, "Username and password are required when authentication is enabled"
  683. # Template-driven: only inline-embed when the user's template explicitly
  684. # referenced {finish_photo_url} (so the URL appears in the rendered body)
  685. # AND the photo bytes are available. Falls back to text-only otherwise.
  686. inline_photo = bool(image_data and finish_photo_url and finish_photo_url in body)
  687. try:
  688. if inline_photo:
  689. # multipart/related → (multipart/alternative → text, html) + inline image
  690. msg = MIMEMultipart("related")
  691. msg["From"] = from_email
  692. msg["To"] = to_email
  693. msg["Subject"] = f"[Bambuddy] {subject}"
  694. alt = MIMEMultipart("alternative")
  695. alt.attach(MIMEText(body, "plain"))
  696. # Build HTML body: escape the rendered body, then swap the
  697. # escaped URL substring for an inline <img> referencing the
  698. # MIMEImage we attach below. Done AFTER escape so the cid: URL
  699. # we inject isn't re-escaped.
  700. escaped_body = html.escape(body).replace("\n", "<br>\n")
  701. escaped_url = html.escape(finish_photo_url)
  702. img_tag = (
  703. '<img src="cid:bambuddy-finish-photo" '
  704. 'alt="Printer camera snapshot" '
  705. 'style="max-width:100%;height:auto;border:1px solid #ddd;border-radius:4px;">'
  706. )
  707. html_body = f"<html><body><p>{escaped_body.replace(escaped_url, img_tag)}</p></body></html>"
  708. alt.attach(MIMEText(html_body, "html"))
  709. msg.attach(alt)
  710. img = MIMEImage(image_data, _subtype="jpeg")
  711. # Angle-bracketed Content-ID per RFC 2392, referenced from HTML
  712. # without the brackets via ``cid:bambuddy-finish-photo``.
  713. img.add_header("Content-ID", "<bambuddy-finish-photo>")
  714. img.add_header("Content-Disposition", "inline", filename="finish-photo.jpg")
  715. msg.attach(img)
  716. else:
  717. msg = MIMEMultipart()
  718. msg["From"] = from_email
  719. msg["To"] = to_email
  720. msg["Subject"] = f"[Bambuddy] {subject}"
  721. msg.attach(MIMEText(body, "plain"))
  722. # smtplib is synchronous and blocking: a wedged / greylisting /
  723. # firewall-dropped relay leaves recv() stuck. Two problems, two
  724. # fixes (#2572):
  725. # 1. No timeout — smtplib defaults to the global socket timeout,
  726. # which this app never sets, so a stuck relay blocks forever.
  727. # Pass an explicit timeout to every connect.
  728. # 2. Run on the event loop — a stuck (or merely slow) send freezes
  729. # every other coroutine, including a DB session a caller is
  730. # holding open across this notification. Offload to a worker
  731. # thread so the loop stays live and the connection is released
  732. # on schedule.
  733. smtp_timeout = 30.0
  734. msg_str = msg.as_string()
  735. def _blocking_send() -> None:
  736. if security == "ssl":
  737. # Direct SSL connection (typically port 465)
  738. server = smtplib.SMTP_SSL(smtp_server, smtp_port, timeout=smtp_timeout)
  739. elif security == "starttls":
  740. # STARTTLS upgrade (typically port 587)
  741. server = smtplib.SMTP(smtp_server, smtp_port, timeout=smtp_timeout)
  742. server.starttls()
  743. else:
  744. # No encryption (typically port 25) - use with caution
  745. server = smtplib.SMTP(smtp_server, smtp_port, timeout=smtp_timeout)
  746. try:
  747. if auth_enabled:
  748. server.login(username, password)
  749. server.sendmail(from_email, to_email, msg_str)
  750. finally:
  751. # quit() in finally so a send error doesn't leak the socket.
  752. try:
  753. server.quit()
  754. except Exception: # noqa: BLE001 — closing a broken connection is best-effort
  755. pass
  756. await asyncio.to_thread(_blocking_send)
  757. return True, "Email sent successfully"
  758. except smtplib.SMTPAuthenticationError:
  759. return False, "SMTP authentication failed - check username/password"
  760. except smtplib.SMTPException as e:
  761. return False, f"SMTP error: {str(e)}"
  762. except Exception as e:
  763. return False, f"Email error: {str(e)}"
  764. async def _send_discord(
  765. self, config: dict, title: str, message: str, image_data: bytes | None = None
  766. ) -> tuple[bool, str]:
  767. """Send notification via Discord webhook."""
  768. webhook_url = config.get("webhook_url", "").strip()
  769. if not webhook_url:
  770. return False, "Webhook URL is required"
  771. if not (
  772. webhook_url.startswith("https://discord.com/api/webhooks/")
  773. or webhook_url.startswith("https://discordapp.com/api/webhooks/")
  774. ):
  775. return False, "Invalid Discord webhook URL"
  776. # Discord embed format for nicer messages
  777. embed = {
  778. "title": title,
  779. "description": message,
  780. "color": 0x00AE42, # Bambu green
  781. }
  782. client = await self._get_client()
  783. if image_data:
  784. # Attach image via multipart form-data and reference in embed
  785. embed["image"] = {"url": "attachment://photo.jpg"}
  786. payload = {"embeds": [embed]}
  787. response = await client.post(
  788. webhook_url,
  789. data={"payload_json": json.dumps(payload)},
  790. files={"files[0]": ("photo.jpg", image_data, "image/jpeg")},
  791. )
  792. else:
  793. response = await client.post(webhook_url, json={"embeds": [embed]})
  794. if response.status_code in (200, 204):
  795. return True, "Message sent successfully"
  796. else:
  797. return False, f"HTTP {response.status_code}: {response.text[:200]}"
  798. async def _send_webhook(
  799. self,
  800. config: dict,
  801. title: str,
  802. message: str,
  803. image_data: bytes | None = None,
  804. event_type: str | None = None,
  805. variables: dict | None = None,
  806. ) -> tuple[bool, str]:
  807. """Send notification via generic webhook (POST JSON).
  808. Supports two payload formats:
  809. - generic: Custom field names with timestamp/source metadata + structured event data
  810. - slack: Slack/Mattermost compatible format (just {"text": "..."})
  811. """
  812. webhook_url = config.get("webhook_url", "").strip()
  813. auth_header = config.get("auth_header", "").strip()
  814. payload_format = config.get("payload_format", "generic").strip()
  815. if not webhook_url:
  816. return False, "Webhook URL is required"
  817. url_error = _assert_safe_provider_url(webhook_url, label="Webhook URL")
  818. if url_error:
  819. return False, url_error
  820. # Build payload based on format
  821. if payload_format == "slack":
  822. # Slack/Mattermost format - just text field
  823. data = {"text": f"*{title}*\n{message}"}
  824. if event_type == "print_confirm_request":
  825. # Slack and Mattermost fetch every URL in the text to build
  826. # preview cards, and the outcome prompt (#1898) is the one
  827. # message whose links are single-use capabilities — that fetch
  828. # would be a machine answering the operator's question. Off
  829. # here for the same reason Telegram's preview is.
  830. #
  831. # Only here: the slack payload never attaches image bytes (the
  832. # base64 attach below is generic-format only), so unfurling is
  833. # the only way a {finish_photo_url} in a print_complete body
  834. # can render as a photo in the channel. Switching it off for
  835. # every event would quietly take that away with no setting to
  836. # get it back.
  837. data["unfurl_links"] = False
  838. data["unfurl_media"] = False
  839. else:
  840. # Generic format with custom field names
  841. custom_field_title = config.get("field_title", "title").strip() or "title"
  842. custom_field_message = config.get("field_message", "message").strip() or "message"
  843. data = {
  844. custom_field_title: title,
  845. custom_field_message: message,
  846. "timestamp": datetime.now().isoformat(),
  847. "source": "Bambuddy",
  848. }
  849. # For generic format, include structured event data for automation tools
  850. if payload_format != "slack":
  851. if event_type:
  852. data["event"] = event_type
  853. if variables:
  854. for key, value in variables.items():
  855. if key not in data: # Don't overwrite title/message/timestamp/source
  856. data[key] = value
  857. # Attach base64-encoded image when available (generic format only)
  858. if image_data and payload_format != "slack":
  859. import base64
  860. data["image"] = base64.b64encode(image_data).decode("ascii")
  861. headers = {"Content-Type": "application/json"}
  862. if auth_header:
  863. # Support "Bearer token" or just "token" format
  864. if " " in auth_header:
  865. headers["Authorization"] = auth_header
  866. else:
  867. headers["Authorization"] = f"Bearer {auth_header}"
  868. client = await self._get_client()
  869. try:
  870. response = await client.post(webhook_url, json=data, headers=headers)
  871. if response.status_code in (200, 201, 202, 204):
  872. return True, "Webhook delivered successfully"
  873. else:
  874. return False, _opaque_http_failure(response, label="webhook endpoint")
  875. except Exception as e:
  876. return False, f"Webhook error: {str(e)}"
  877. async def _send_homeassistant(
  878. self, config: dict, title: str, message: str, db: AsyncSession | None = None
  879. ) -> tuple[bool, str]:
  880. """Send notification via Home Assistant.
  881. Uses the globally configured HA URL/token from settings.
  882. Defaults to persistent_notification/create, but supports
  883. custom services via config["service"] (e.g. notify.mobile_app_myphone).
  884. """
  885. # Get HA connection settings from global config
  886. ha_url = ""
  887. ha_token = ""
  888. if db:
  889. from backend.app.api.routes.settings import get_homeassistant_settings
  890. try:
  891. ha_settings = await get_homeassistant_settings(db)
  892. ha_url = ha_settings.get("ha_url", "")
  893. ha_token = ha_settings.get("ha_token", "")
  894. except Exception as e:
  895. logger.warning("Failed to read HA settings from database: %s", e)
  896. else:
  897. # Fallback: read directly from environment if no DB session
  898. import os
  899. ha_url = os.environ.get("HA_URL", "")
  900. ha_token = os.environ.get("HA_TOKEN", "")
  901. if not ha_url or not ha_token:
  902. return False, (
  903. "Home Assistant is not configured. Please set HA URL and token in Settings → Network → Home Assistant."
  904. )
  905. # Determine which HA service to call - Default: persistent_notification.create
  906. service = (config.get("service") or "").strip()
  907. if service:
  908. # Allow in different forms:
  909. # - notify.mobile_app_<device>
  910. # - notify/mobile_app_<device>
  911. # - api/services/notify/mobile_app_<device>
  912. service_str = service.lstrip("/")
  913. if service_str.startswith("api/services/"):
  914. endpoint = service_str
  915. elif "/" in service_str:
  916. endpoint = f"api/services/{service_str}"
  917. elif "." in service_str:
  918. domain, svc = service_str.split(".", 1)
  919. endpoint = f"api/services/{domain}/{svc}"
  920. else:
  921. return False, (
  922. "Invalid Home Assistant service name. Use e.g. 'notify.mobile_app_yourdevice' or 'notify/your_service'."
  923. )
  924. if not re.match(r"^api/services/[a-zA-Z0-9_]+/[a-zA-Z0-9_]+$", endpoint):
  925. return False, (
  926. "Invalid Home Assistant service name. Domain and service must only contain letters, numbers, and underscores."
  927. )
  928. else:
  929. endpoint = "api/services/persistent_notification/create"
  930. url = f"{ha_url.rstrip('/')}/{endpoint}"
  931. headers = {
  932. "Authorization": f"Bearer {ha_token}",
  933. "Content-Type": "application/json",
  934. }
  935. payload = {
  936. "title": title,
  937. "message": message,
  938. }
  939. # Optional custom service-data (#1441), forwarded as HA's nested "data"
  940. # object so mobile-app push options (priority, ttl, channel, group, ...)
  941. # reach the notify service. Only included when configured — the default
  942. # persistent_notification.create schema rejects unknown keys.
  943. raw_data = config.get("data")
  944. if raw_data:
  945. if isinstance(raw_data, str):
  946. try:
  947. parsed_data = json.loads(raw_data)
  948. except json.JSONDecodeError as e:
  949. return False, f"Invalid JSON in the Data field: {e}"
  950. else:
  951. parsed_data = raw_data
  952. if not isinstance(parsed_data, dict):
  953. return False, 'The Data field must be a JSON object, e.g. {"priority": "high", "ttl": 0}'
  954. if parsed_data:
  955. payload["data"] = parsed_data
  956. client = await self._get_client()
  957. response = await client.post(url, json=payload, headers=headers)
  958. if response.status_code in (200, 201):
  959. return True, "Notification sent via Home Assistant"
  960. elif response.status_code == 401:
  961. return False, "Home Assistant authentication failed - check your token"
  962. else:
  963. # ha_url comes from global settings (SETTINGS_UPDATE, admin-only), so
  964. # this is a narrower channel than the per-request provider URLs — but
  965. # it lands in the same NOTIFICATIONS_CREATE-gated test response, so it
  966. # gets the same treatment.
  967. return False, _opaque_http_failure(response, label="Home Assistant endpoint")
  968. async def _send_to_provider(
  969. self,
  970. provider: NotificationProvider,
  971. title: str,
  972. message: str,
  973. db: AsyncSession | None = None,
  974. image_data: bytes | None = None,
  975. event_type: str | None = None,
  976. variables: dict | None = None,
  977. ) -> tuple[bool, str]:
  978. """Send notification to a specific provider."""
  979. # Check quiet hours
  980. if self._is_in_quiet_hours(provider):
  981. logger.info("Skipping notification to %s - quiet hours active", provider.name)
  982. return True, "Skipped - quiet hours"
  983. config = json.loads(provider.config) if isinstance(provider.config, str) else provider.config
  984. try:
  985. if provider.provider_type == "callmebot":
  986. return await self._send_callmebot(config, f"{title}\n{message}")
  987. elif provider.provider_type == "ntfy":
  988. # Outcome confirmation (#1898): render the verdict capability
  989. # links as one-tap buttons on the notification itself. http so
  990. # no browser needs to open; POST because that is the method
  991. # that records — a GET only opens the confirmation page, which
  992. # is what keeps unfurlers from answering the prompt.
  993. # clear=true dismisses the notification once a button was tapped.
  994. ntfy_actions = None
  995. good_url = (variables or {}).get("good_url")
  996. reject_url = (variables or {}).get("reject_url")
  997. # Buttons need absolute URLs; without a configured external_url
  998. # the links are relative, and the body's deep link into the
  999. # archive has to do.
  1000. if (
  1001. event_type == "print_confirm_request"
  1002. and good_url
  1003. and reject_url
  1004. and good_url.startswith("http")
  1005. and reject_url.startswith("http")
  1006. ):
  1007. ntfy_actions = (
  1008. f"http, Good, {good_url}, method=POST, clear=true; "
  1009. f"http, Reject, {reject_url}, method=POST, clear=true"
  1010. )
  1011. return await self._send_ntfy(
  1012. config, title, message, image_data=image_data, event_type=event_type, actions=ntfy_actions
  1013. )
  1014. elif provider.provider_type == "pushover":
  1015. # Outcome confirmation (#1898): Pushover has no arbitrary
  1016. # buttons, but supports one supplementary URL — deep-link into
  1017. # the archive's confirmation dialog.
  1018. supplement_url = None
  1019. supplement_url_title = None
  1020. _confirm_url = (variables or {}).get("confirm_url")
  1021. if event_type == "print_confirm_request" and _confirm_url and _confirm_url.startswith("http"):
  1022. supplement_url = _confirm_url
  1023. supplement_url_title = "Confirm print outcome"
  1024. return await self._send_pushover(
  1025. config, title, message, image_data=image_data, url=supplement_url, url_title=supplement_url_title
  1026. )
  1027. elif provider.provider_type == "telegram":
  1028. if event_type == "print_confirm_request":
  1029. # Outcome confirmation (#1898): inline URL buttons under
  1030. # the message — one tap records the verdict via the
  1031. # capability link. Same absolute-URL requirement as the
  1032. # ntfy actions. Telegram has no way to POST, so this is the
  1033. # one affordance that opens a browser, and it is the one
  1034. # that gets the one-tap marker: the page submits itself
  1035. # only for a URL that came off a button. Telegram does not
  1036. # fetch inline-keyboard URLs and nothing else can read
  1037. # them, so the marker never reaches a scanner — which is
  1038. # the difference between this and trusting the User-Agent.
  1039. tg_buttons = None
  1040. _tg_good = (variables or {}).get("good_url")
  1041. _tg_reject = (variables or {}).get("reject_url")
  1042. if _tg_good and _tg_reject and _tg_good.startswith("http") and _tg_reject.startswith("http"):
  1043. tg_buttons = [
  1044. {"text": "\U0001f44d Good", "url": one_tap_url(_tg_good)},
  1045. {"text": "\U0001f44e Reject", "url": one_tap_url(_tg_reject)},
  1046. ]
  1047. # Telegram verdict mode (#3046): buttons, a reaction on
  1048. # the message, or both. The link preview is off in all of
  1049. # them (see _send_telegram_confirm_request).
  1050. _archive_id = (variables or {}).get("archive_id")
  1051. return await self._send_telegram_confirm_request(
  1052. provider,
  1053. config,
  1054. f"*{title}*\n{message}",
  1055. db,
  1056. image_data,
  1057. tg_buttons,
  1058. _archive_id if isinstance(_archive_id, int) else None,
  1059. )
  1060. # Every other event keeps Telegram's link preview, for the same
  1061. # reason as the Slack unfurl in _send_webhook: when the finish
  1062. # photo is too large to attach, the preview is how a
  1063. # {finish_photo_url} in a print_complete body still shows up as
  1064. # a photo in the chat.
  1065. return await self._send_telegram(config, f"*{title}*\n{message}", image_data=image_data)
  1066. elif provider.provider_type == "email":
  1067. # finish_photo_url is pulled from the rendered template variables
  1068. # so _send_email can detect whether the template referenced the
  1069. # URL and inline-embed the photo only in that case.
  1070. finish_photo_url = (variables or {}).get("finish_photo_url")
  1071. return await self._send_email(
  1072. config, title, message, image_data=image_data, finish_photo_url=finish_photo_url
  1073. )
  1074. elif provider.provider_type == "discord":
  1075. return await self._send_discord(config, title, message, image_data=image_data)
  1076. elif provider.provider_type == "webhook":
  1077. return await self._send_webhook(
  1078. config, title, message, image_data=image_data, event_type=event_type, variables=variables
  1079. )
  1080. elif provider.provider_type == "homeassistant":
  1081. return await self._send_homeassistant(config, title, message, db=db)
  1082. elif provider.provider_type == "bark":
  1083. # Outcome confirmation (#1898): Bark opens one URL on tap —
  1084. # deep-link into the confirmation dialog, like Pushover.
  1085. bark_url = None
  1086. _bark_confirm = (variables or {}).get("confirm_url")
  1087. if event_type == "print_confirm_request" and _bark_confirm and _bark_confirm.startswith("http"):
  1088. bark_url = _bark_confirm
  1089. return await self._send_bark(config, title, message, url=bark_url)
  1090. else:
  1091. return False, f"Unknown provider type: {provider.provider_type}"
  1092. except Exception as e:
  1093. logger.exception("Error sending notification via %s", provider.provider_type)
  1094. return False, str(e)
  1095. async def _update_provider_status(
  1096. self, db: AsyncSession, provider_id: int, success: bool, error: str | None = None
  1097. ):
  1098. """Update provider status after sending notification."""
  1099. result = await db.execute(select(NotificationProvider).where(NotificationProvider.id == provider_id))
  1100. provider = result.scalar_one_or_none()
  1101. if provider:
  1102. if success:
  1103. provider.last_success = datetime.now(timezone.utc)
  1104. else:
  1105. provider.last_error = error
  1106. provider.last_error_at = datetime.now(timezone.utc)
  1107. await db.commit()
  1108. async def _get_providers_for_event(
  1109. self,
  1110. db: AsyncSession,
  1111. event_field: str,
  1112. printer_id: int | None = None,
  1113. ) -> list[NotificationProvider]:
  1114. """Get all enabled providers that want a specific event type.
  1115. Runs under ``no_autoflush`` (#2770). Callers routinely hold pending
  1116. writes when they raise an event — the AMS sensor loop does
  1117. ``db.add(history)`` and only commits after the alarms have gone out — and
  1118. without this, autoflush satisfies this SELECT by writing those rows,
  1119. which opens a write transaction on SQLite. The provider is then contacted
  1120. over the network with that transaction still open, so a site whose
  1121. internet is down holds the single SQLite writer for the whole connect
  1122. timeout and unrelated background tasks fail with "database is locked".
  1123. Deferring the flush costs nothing here: providers are committed rows, so
  1124. a pending change in the caller's session cannot be one this query wants.
  1125. """
  1126. # Build the query dynamically based on event field
  1127. query = select(NotificationProvider).where(
  1128. NotificationProvider.enabled.is_(True),
  1129. getattr(NotificationProvider, event_field).is_(True),
  1130. )
  1131. if printer_id is not None:
  1132. query = query.where(
  1133. (NotificationProvider.printer_id.is_(None)) | (NotificationProvider.printer_id == printer_id)
  1134. )
  1135. with db.no_autoflush:
  1136. result = await db.execute(query)
  1137. return list(result.scalars().all())
  1138. async def _log_notification(
  1139. self,
  1140. db: AsyncSession,
  1141. provider_id: int,
  1142. event_type: str,
  1143. title: str,
  1144. message: str,
  1145. success: bool,
  1146. error_message: str | None = None,
  1147. printer_id: int | None = None,
  1148. printer_name: str | None = None,
  1149. ):
  1150. """Create a log entry for a sent notification."""
  1151. try:
  1152. log = NotificationLog(
  1153. provider_id=provider_id,
  1154. event_type=event_type,
  1155. title=title,
  1156. message=message,
  1157. success=success,
  1158. error_message=error_message,
  1159. printer_id=printer_id,
  1160. printer_name=printer_name,
  1161. )
  1162. db.add(log)
  1163. await db.commit()
  1164. except Exception as e:
  1165. logger.warning("Failed to log notification: %s", e)
  1166. # Don't fail the notification just because logging failed
  1167. async def _send_to_providers(
  1168. self,
  1169. providers: list[NotificationProvider],
  1170. title: str,
  1171. message: str,
  1172. db: AsyncSession,
  1173. event_type: str = "unknown",
  1174. printer_id: int | None = None,
  1175. printer_name: str | None = None,
  1176. force_immediate: bool = False,
  1177. image_data: bytes | None = None,
  1178. variables: dict | None = None,
  1179. ):
  1180. """Send notification to multiple providers and log the results.
  1181. All notifications are always sent immediately. If digest mode is enabled,
  1182. the notification is ALSO queued for the daily digest summary.
  1183. """
  1184. for provider in providers:
  1185. try:
  1186. # Always send notification immediately
  1187. success, error = await self._send_to_provider(
  1188. provider, title, message, db, image_data=image_data, event_type=event_type, variables=variables
  1189. )
  1190. # Also queue for digest if enabled (digest is a summary, not a queue)
  1191. if provider.daily_digest_enabled and provider.daily_digest_time:
  1192. await self._queue_for_digest(
  1193. provider=provider,
  1194. event_type=event_type,
  1195. title=title,
  1196. message=message,
  1197. db=db,
  1198. printer_id=printer_id,
  1199. printer_name=printer_name,
  1200. )
  1201. await self._update_provider_status(db, provider.id, success, error if not success else None)
  1202. await self._log_notification(
  1203. db=db,
  1204. provider_id=provider.id,
  1205. event_type=event_type,
  1206. title=title,
  1207. message=message,
  1208. success=success,
  1209. error_message=error if not success else None,
  1210. printer_id=printer_id,
  1211. printer_name=printer_name,
  1212. )
  1213. if success:
  1214. logger.info("Sent notification via %s", provider.name)
  1215. else:
  1216. logger.warning("Failed to send notification via %s: %s", provider.name, error)
  1217. except Exception as e:
  1218. logger.exception("Error sending notification via %s", provider.name)
  1219. await self._update_provider_status(db, provider.id, False, str(e))
  1220. await self._log_notification(
  1221. db=db,
  1222. provider_id=provider.id,
  1223. event_type=event_type,
  1224. title=title,
  1225. message=message,
  1226. success=False,
  1227. error_message=str(e),
  1228. printer_id=printer_id,
  1229. printer_name=printer_name,
  1230. )
  1231. async def on_app_message(
  1232. self,
  1233. db: AsyncSession,
  1234. *,
  1235. sender: str,
  1236. title: str,
  1237. message: str,
  1238. url: str | None = None,
  1239. ) -> int:
  1240. """A message another application sends through Bambuddy.
  1241. Goes to every enabled channel with "Messages from connected apps" on,
  1242. through the same path as Bambuddy's own events: quiet hours, the daily
  1243. digest and the log (whose event type names the sender). The text is
  1244. the app's own; a link, when given, is appended so every channel type
  1245. carries it. Returns how many channels it was handed to.
  1246. """
  1247. providers = await self._get_providers_for_event(db, "on_app_message")
  1248. if not providers:
  1249. return 0
  1250. body = f"{message}\n{url}" if url else message
  1251. await self._send_to_providers(providers, title, body, db, event_type=f"app:{sender}"[:50])
  1252. return len(providers)
  1253. async def on_print_start(
  1254. self,
  1255. printer_id: int,
  1256. printer_name: str,
  1257. data: dict,
  1258. db: AsyncSession,
  1259. archive_data: dict | None = None,
  1260. ):
  1261. """Handle print start event - send notifications to relevant providers.
  1262. Args:
  1263. printer_id: The printer ID
  1264. printer_name: The printer name
  1265. data: MQTT event data with filename, subtask_name, remaining_time, raw_data
  1266. db: Database session
  1267. archive_data: Optional archive data with print_time_seconds from 3MF parsing
  1268. """
  1269. logger.info("on_print_start called for printer %s (%s)", printer_id, printer_name)
  1270. providers = await self._get_providers_for_event(db, "on_print_start", printer_id)
  1271. if not providers:
  1272. logger.info("No notification providers configured for print_start event on printer %s", printer_id)
  1273. return
  1274. # Use subtask_name (project name) if available, otherwise use filename
  1275. subtask_name = data.get("subtask_name")
  1276. if subtask_name:
  1277. # Replace underscores with spaces for readability
  1278. filename = subtask_name.replace("_", " ")
  1279. else:
  1280. filename = self._clean_filename(data.get("filename", "Unknown"))
  1281. # Priority for estimated_time:
  1282. # 1. Archive's print_time_seconds from 3MF parsing (most reliable)
  1283. # 2. MQTT remaining_time (may be 0 at print start)
  1284. # 3. raw_data mc_remaining_time
  1285. estimated_time = None
  1286. # Try archive data first (from 3MF parsing - most reliable)
  1287. if archive_data and archive_data.get("print_time_seconds"):
  1288. estimated_time = archive_data["print_time_seconds"]
  1289. logger.debug("Using print_time_seconds from archive: %s", estimated_time)
  1290. # Fall back to MQTT remaining_time
  1291. if estimated_time is None:
  1292. estimated_time = data.get("remaining_time")
  1293. if estimated_time:
  1294. logger.debug("Using remaining_time from MQTT: %s", estimated_time)
  1295. # Last resort: raw_data mc_remaining_time (in minutes, convert to seconds)
  1296. if estimated_time is None:
  1297. raw_time = data.get("raw_data", {}).get("mc_remaining_time")
  1298. if raw_time:
  1299. estimated_time = raw_time * 60
  1300. logger.debug("Using mc_remaining_time from raw_data: %s", estimated_time)
  1301. time_str = self._format_duration(estimated_time)
  1302. eta_str = await self._format_eta(estimated_time, db)
  1303. variables = {
  1304. "printer": printer_name,
  1305. "filename": filename,
  1306. "estimated_time": time_str,
  1307. "eta": eta_str,
  1308. }
  1309. # Extract image data for providers that support attachments (e.g. Pushover)
  1310. image_data = None
  1311. if archive_data:
  1312. image_data = archive_data.get("image_data")
  1313. logger.info("Found %s providers for print_start: %s", len(providers), [p.name for p in providers])
  1314. title, message = await self._build_message_from_template(db, "print_start", variables)
  1315. await self._send_to_providers(
  1316. providers,
  1317. title,
  1318. message,
  1319. db,
  1320. "print_start",
  1321. printer_id,
  1322. printer_name,
  1323. image_data=image_data,
  1324. variables=variables,
  1325. )
  1326. async def on_print_complete(
  1327. self,
  1328. printer_id: int,
  1329. printer_name: str,
  1330. status: str,
  1331. data: dict,
  1332. db: AsyncSession,
  1333. archive_data: dict | None = None,
  1334. ):
  1335. """Handle print complete event - send notifications to relevant providers."""
  1336. logger.info("on_print_complete called for printer %s (%s), status=%s", printer_id, printer_name, status)
  1337. # Determine event type based on status
  1338. if status == "completed":
  1339. event_field = "on_print_complete"
  1340. event_type = "print_complete"
  1341. elif status in ("failed",):
  1342. event_field = "on_print_failed"
  1343. event_type = "print_failed"
  1344. elif status in ("aborted", "stopped", "cancelled"):
  1345. event_field = "on_print_stopped"
  1346. event_type = "print_stopped"
  1347. else:
  1348. logger.warning("Unknown print status '%s', defaulting to on_print_complete", status)
  1349. event_field = "on_print_complete"
  1350. event_type = "print_complete"
  1351. providers = await self._get_providers_for_event(db, event_field, printer_id)
  1352. if not providers:
  1353. logger.info("No notification providers configured for %s event on printer %s", event_field, printer_id)
  1354. return
  1355. # Use subtask_name (project name) if available, otherwise use filename
  1356. subtask_name = data.get("subtask_name")
  1357. if subtask_name:
  1358. filename = subtask_name.replace("_", " ")
  1359. else:
  1360. filename = self._clean_filename(data.get("filename", "Unknown"))
  1361. variables = {
  1362. "printer": printer_name,
  1363. "filename": filename,
  1364. "duration": "Unknown",
  1365. "filament_grams": "Unknown",
  1366. "reason": "Unknown",
  1367. }
  1368. if archive_data:
  1369. # {{duration}} on completion / failure / stopped events is the *actual*
  1370. # elapsed time (#1198). Slicer-estimated print_time_seconds is only used
  1371. # as a last-resort fallback when timestamps weren't recorded.
  1372. duration_seconds = archive_data.get("actual_time_seconds") or archive_data.get("print_time_seconds")
  1373. if duration_seconds:
  1374. variables["duration"] = self._format_duration(duration_seconds)
  1375. if archive_data.get("actual_filament_grams"):
  1376. variables["filament_grams"] = f"{archive_data['actual_filament_grams']:.1f}"
  1377. if status == "failed" and archive_data.get("failure_reason"):
  1378. variables["reason"] = archive_data["failure_reason"]
  1379. if archive_data.get("finish_photo_url"):
  1380. variables["finish_photo_url"] = archive_data["finish_photo_url"]
  1381. # Build per-slot breakdown string with AMS info when available
  1382. if archive_data.get("usage_results"):
  1383. parts = []
  1384. for u in archive_data["usage_results"]:
  1385. ams_id = u.get("ams_id", 0)
  1386. tray_id = u.get("tray_id", 0)
  1387. material = u.get("material", "Unknown") or "Unknown"
  1388. used = u.get("weight_used", 0)
  1389. if ams_id >= 128:
  1390. slot_label = "Ext"
  1391. else:
  1392. slot_label = f"AMS-{chr(65 + ams_id)} T{tray_id + 1}"
  1393. parts.append(f"{slot_label} {material}: {used:.1f}g")
  1394. variables["filament_details"] = " | ".join(parts)
  1395. elif archive_data.get("filament_slots"):
  1396. parts = []
  1397. for slot in archive_data["filament_slots"]:
  1398. ftype = slot.get("type", "Unknown") or "Unknown"
  1399. used = slot.get("used_g", 0)
  1400. parts.append(f"{ftype}: {used:.1f}g")
  1401. variables["filament_details"] = " | ".join(parts)
  1402. # Add progress for partial prints
  1403. if archive_data.get("progress") is not None:
  1404. variables["progress"] = str(archive_data["progress"])
  1405. # Extract image data for providers that support attachments (e.g. Pushover)
  1406. image_data = None
  1407. if archive_data:
  1408. image_data = archive_data.get("image_data")
  1409. logger.info("Found %s providers for %s: %s", len(providers), event_field, [p.name for p in providers])
  1410. title, message = await self._build_message_from_template(db, event_type, variables)
  1411. await self._send_to_providers(
  1412. providers,
  1413. title,
  1414. message,
  1415. db,
  1416. event_type,
  1417. printer_id,
  1418. printer_name,
  1419. image_data=image_data,
  1420. variables=variables,
  1421. )
  1422. async def on_print_confirm_request(
  1423. self,
  1424. printer_id: int,
  1425. printer_name: str,
  1426. data: dict,
  1427. db: AsyncSession,
  1428. archive_data: dict | None = None,
  1429. good_url: str | None = None,
  1430. reject_url: str | None = None,
  1431. confirm_url: str | None = None,
  1432. archive_id: int | None = None,
  1433. ):
  1434. """Ask for a post-print outcome verdict (#1898).
  1435. Fires only for completed prints whose queue item opted in — the
  1436. provider-level toggle exists to mute a channel, not to enable the
  1437. feature. good_url / reject_url are the one-tap capability links
  1438. (rendered as ntfy action buttons), confirm_url deep-links into the
  1439. archive's confirmation dialog in the web UI. archive_id lets a Telegram
  1440. provider in reaction mode (#3046) tie the sent message to the archive.
  1441. """
  1442. providers = await self._get_providers_for_event(db, "on_print_confirm_request", printer_id)
  1443. if not providers:
  1444. return
  1445. subtask_name = data.get("subtask_name")
  1446. if subtask_name:
  1447. filename = subtask_name.replace("_", " ")
  1448. else:
  1449. filename = self._clean_filename(data.get("filename", "Unknown"))
  1450. variables = {"printer": printer_name, "filename": filename}
  1451. if good_url:
  1452. variables["good_url"] = good_url
  1453. if reject_url:
  1454. variables["reject_url"] = reject_url
  1455. if confirm_url:
  1456. variables["confirm_url"] = confirm_url
  1457. if archive_id is not None:
  1458. variables["archive_id"] = archive_id
  1459. image_data = None
  1460. if archive_data:
  1461. if archive_data.get("finish_photo_url"):
  1462. variables["finish_photo_url"] = archive_data["finish_photo_url"]
  1463. image_data = archive_data.get("image_data")
  1464. title, message = await self._build_message_from_template(db, "print_confirm_request", variables)
  1465. await self._send_to_providers(
  1466. providers,
  1467. title,
  1468. message,
  1469. db,
  1470. "print_confirm_request",
  1471. printer_id,
  1472. printer_name,
  1473. image_data=image_data,
  1474. variables=variables,
  1475. )
  1476. async def on_print_progress(
  1477. self,
  1478. printer_id: int,
  1479. printer_name: str,
  1480. filename: str,
  1481. progress: int,
  1482. db: AsyncSession,
  1483. remaining_time: int | None = None,
  1484. image_data: bytes | None = None,
  1485. ):
  1486. """Handle print progress milestone (25%, 50%, 75%)."""
  1487. providers = await self._get_providers_for_event(db, "on_print_progress", printer_id)
  1488. if not providers:
  1489. return
  1490. eta_str = await self._format_eta(remaining_time, db)
  1491. variables = {
  1492. "printer": printer_name,
  1493. "filename": self._clean_filename(filename),
  1494. "progress": str(progress),
  1495. "remaining_time": self._format_duration(remaining_time) if remaining_time else "Unknown",
  1496. "eta": eta_str,
  1497. }
  1498. title, message = await self._build_message_from_template(db, "print_progress", variables)
  1499. await self._send_to_providers(
  1500. providers,
  1501. title,
  1502. message,
  1503. db,
  1504. "print_progress",
  1505. printer_id,
  1506. printer_name,
  1507. image_data=image_data,
  1508. variables=variables,
  1509. )
  1510. async def on_print_missing_spool_assignment(
  1511. self,
  1512. printer_id: int,
  1513. printer_name: str,
  1514. missing_slots: list[dict[str, str]],
  1515. db: AsyncSession,
  1516. ):
  1517. """Handle print-start event when required trays are missing spool assignments."""
  1518. if not missing_slots:
  1519. return
  1520. providers = await self._get_providers_for_event(db, "on_print_missing_spool_assignment", printer_id)
  1521. if not providers:
  1522. return
  1523. missing_slot_names = ", ".join(slot.get("slot", "Unknown") for slot in missing_slots)
  1524. detail_lines = []
  1525. for slot in missing_slots:
  1526. slot_name = slot.get("slot", "Unknown")
  1527. profile = slot.get("profile", "Unknown")
  1528. detail_lines.append(f"- {slot_name}: {profile}")
  1529. missing_profile_details = "\n".join(detail_lines)
  1530. variables = {
  1531. "printer": printer_name,
  1532. "missing_slots": missing_slot_names,
  1533. "missing_slot_details": missing_profile_details,
  1534. }
  1535. title, message = await self._build_message_from_template(db, "print_missing_spool_assignment", variables)
  1536. await self._send_to_providers(
  1537. providers,
  1538. title,
  1539. message,
  1540. db,
  1541. "print_missing_spool_assignment",
  1542. printer_id,
  1543. printer_name,
  1544. force_immediate=True,
  1545. variables=variables,
  1546. )
  1547. async def on_billing_charge_failed(
  1548. self,
  1549. printer_id: int,
  1550. printer_name: str,
  1551. filename: str,
  1552. archive_id: int | None,
  1553. error: str,
  1554. db: AsyncSession,
  1555. ) -> None:
  1556. """Notify providers that a terminal print could not be charged."""
  1557. providers = await self._get_providers_for_event(db, "on_billing_charge_failed", printer_id)
  1558. if not providers:
  1559. return
  1560. variables = {
  1561. "printer": printer_name,
  1562. "filename": self._clean_filename(filename),
  1563. "archive_id": str(archive_id) if archive_id is not None else "Unknown",
  1564. "error": error,
  1565. }
  1566. title, message = await self._build_message_from_template(db, "billing_charge_failed", variables)
  1567. await self._send_to_providers(
  1568. providers,
  1569. title,
  1570. message,
  1571. db,
  1572. "billing_charge_failed",
  1573. printer_id,
  1574. printer_name,
  1575. force_immediate=True,
  1576. variables=variables,
  1577. )
  1578. async def on_printer_offline(self, printer_id: int, printer_name: str, db: AsyncSession):
  1579. """Handle printer offline event."""
  1580. providers = await self._get_providers_for_event(db, "on_printer_offline", printer_id)
  1581. if not providers:
  1582. return
  1583. variables = {"printer": printer_name}
  1584. title, message = await self._build_message_from_template(db, "printer_offline", variables)
  1585. await self._send_to_providers(
  1586. providers, title, message, db, "printer_offline", printer_id, printer_name, variables=variables
  1587. )
  1588. async def on_printer_error(
  1589. self,
  1590. printer_id: int,
  1591. printer_name: str,
  1592. error_type: str,
  1593. db: AsyncSession,
  1594. error_detail: str | None = None,
  1595. image_data: bytes | None = None,
  1596. ):
  1597. """Handle printer error event (AMS issues, etc.)."""
  1598. providers = await self._get_providers_for_event(db, "on_printer_error", printer_id)
  1599. if not providers:
  1600. return
  1601. variables = {
  1602. "printer": printer_name,
  1603. "error_type": error_type,
  1604. "error_detail": error_detail or "No details available",
  1605. }
  1606. title, message = await self._build_message_from_template(db, "printer_error", variables)
  1607. await self._send_to_providers(
  1608. providers,
  1609. title,
  1610. message,
  1611. db,
  1612. "printer_error",
  1613. printer_id,
  1614. printer_name,
  1615. image_data=image_data,
  1616. variables=variables,
  1617. )
  1618. async def on_ai_failure_detection(
  1619. self,
  1620. printer_id: int,
  1621. printer_name: str,
  1622. task_name: str,
  1623. confidence: float,
  1624. action: str,
  1625. db: AsyncSession,
  1626. image_data: bytes | None = None,
  1627. ):
  1628. """Handle AI failure-detection event (Obico spaghetti / print-failure ML).
  1629. Split out of on_printer_error (#1794) so a user can subscribe to AI
  1630. alerts without also being paged for every HMS hardware code.
  1631. """
  1632. providers = await self._get_providers_for_event(db, "on_ai_failure_detection", printer_id)
  1633. if not providers:
  1634. return
  1635. variables = {
  1636. "printer": printer_name,
  1637. "task_name": task_name or "current job",
  1638. "confidence": f"{confidence:.2f}",
  1639. "action": action,
  1640. }
  1641. title, message = await self._build_message_from_template(db, "ai_failure_detection", variables)
  1642. await self._send_to_providers(
  1643. providers,
  1644. title,
  1645. message,
  1646. db,
  1647. "ai_failure_detection",
  1648. printer_id,
  1649. printer_name,
  1650. image_data=image_data,
  1651. variables=variables,
  1652. )
  1653. async def on_plate_not_empty(
  1654. self,
  1655. printer_id: int,
  1656. printer_name: str,
  1657. db: AsyncSession,
  1658. difference_percent: float | None = None,
  1659. ):
  1660. """Handle plate not empty event - objects detected on build plate before print."""
  1661. providers = await self._get_providers_for_event(db, "on_plate_not_empty", printer_id)
  1662. if not providers:
  1663. return
  1664. variables = {
  1665. "printer": printer_name,
  1666. "difference_percent": f"{difference_percent:.1f}" if difference_percent else "N/A",
  1667. }
  1668. title, message = await self._build_message_from_template(db, "plate_not_empty", variables)
  1669. await self._send_to_providers(
  1670. providers,
  1671. title,
  1672. message,
  1673. db,
  1674. "plate_not_empty",
  1675. printer_id,
  1676. printer_name,
  1677. force_immediate=True,
  1678. variables=variables,
  1679. )
  1680. async def on_plate_clear_required(
  1681. self,
  1682. printer_id: int,
  1683. printer_name: str,
  1684. db: AsyncSession,
  1685. ):
  1686. """Handle plate-clear-required event — a print ended and the queue is gated (#2525).
  1687. Distinct from ``on_plate_not_empty``, which is the camera check *before* a
  1688. print starts. This one fires on the rising edge of the Bambuddy-side
  1689. awaiting-plate-clear flag, i.e. whenever a print reaches a terminal state
  1690. and the next queued job can't dispatch until someone confirms the bed is
  1691. free. Off by default on every provider: it lands at the same moment as the
  1692. print-complete notification, so opting in is a deliberate choice.
  1693. """
  1694. providers = await self._get_providers_for_event(db, "on_plate_clear_required", printer_id)
  1695. if not providers:
  1696. return
  1697. variables = {"printer": printer_name}
  1698. title, message = await self._build_message_from_template(db, "plate_clear_required", variables)
  1699. await self._send_to_providers(
  1700. providers,
  1701. title,
  1702. message,
  1703. db,
  1704. "plate_clear_required",
  1705. printer_id,
  1706. printer_name,
  1707. variables=variables,
  1708. )
  1709. async def on_filament_low(
  1710. self,
  1711. printer_id: int,
  1712. printer_name: str,
  1713. slot: str,
  1714. remaining_percent: int,
  1715. db: AsyncSession,
  1716. color: str | None = None,
  1717. ):
  1718. """Handle low filament event."""
  1719. providers = await self._get_providers_for_event(db, "on_filament_low", printer_id)
  1720. if not providers:
  1721. return
  1722. variables = {
  1723. "printer": printer_name,
  1724. "slot": slot,
  1725. "remaining_percent": str(remaining_percent),
  1726. "color": color or "",
  1727. }
  1728. title, message = await self._build_message_from_template(db, "filament_low", variables)
  1729. await self._send_to_providers(
  1730. providers, title, message, db, "filament_low", printer_id, printer_name, variables=variables
  1731. )
  1732. async def on_maintenance_due(
  1733. self,
  1734. printer_id: int,
  1735. printer_name: str,
  1736. maintenance_items: list[dict],
  1737. db: AsyncSession,
  1738. ):
  1739. """Handle maintenance due event - sends notification when maintenance is due or warning."""
  1740. if not maintenance_items:
  1741. return
  1742. providers = await self._get_providers_for_event(db, "on_maintenance_due", printer_id)
  1743. if not providers:
  1744. logger.info("No notification providers configured for maintenance_due event on printer %s", printer_id)
  1745. return
  1746. # Format maintenance items list
  1747. items_list = []
  1748. for item in maintenance_items:
  1749. status = "OVERDUE" if item.get("is_due") else "Soon"
  1750. items_list.append(f"- {item['name']} ({status})")
  1751. items_str = "\n".join(items_list)
  1752. variables = {
  1753. "printer": printer_name,
  1754. "items": items_str,
  1755. }
  1756. logger.info("Found %s providers for maintenance_due: %s", len(providers), [p.name for p in providers])
  1757. title, message = await self._build_message_from_template(db, "maintenance_due", variables)
  1758. await self._send_to_providers(
  1759. providers, title, message, db, "maintenance_due", printer_id, printer_name, variables=variables
  1760. )
  1761. async def on_ams_humidity_high(
  1762. self,
  1763. printer_id: int,
  1764. printer_name: str,
  1765. ams_label: str,
  1766. humidity: float,
  1767. threshold: float,
  1768. db: AsyncSession,
  1769. ):
  1770. """Handle AMS high humidity alarm event. Always sends immediately (bypasses digest)."""
  1771. providers = await self._get_providers_for_event(db, "on_ams_humidity_high", printer_id)
  1772. if not providers:
  1773. return
  1774. variables = {
  1775. "printer": printer_name,
  1776. "ams_label": ams_label,
  1777. "humidity": f"{humidity:.0f}",
  1778. "threshold": f"{threshold:.0f}",
  1779. }
  1780. title, message = await self._build_message_from_template(db, "ams_humidity_high", variables)
  1781. # Alarms always send immediately, bypassing digest mode
  1782. await self._send_to_providers(
  1783. providers,
  1784. title,
  1785. message,
  1786. db,
  1787. "ams_humidity_high",
  1788. printer_id,
  1789. printer_name,
  1790. force_immediate=True,
  1791. variables=variables,
  1792. )
  1793. async def on_ams_temperature_high(
  1794. self,
  1795. printer_id: int,
  1796. printer_name: str,
  1797. ams_label: str,
  1798. temperature: float,
  1799. threshold: float,
  1800. db: AsyncSession,
  1801. ):
  1802. """Handle AMS high temperature alarm event. Always sends immediately (bypasses digest)."""
  1803. providers = await self._get_providers_for_event(db, "on_ams_temperature_high", printer_id)
  1804. if not providers:
  1805. return
  1806. variables = {
  1807. "printer": printer_name,
  1808. "ams_label": ams_label,
  1809. "temperature": f"{temperature:.1f}",
  1810. "threshold": f"{threshold:.1f}",
  1811. }
  1812. title, message = await self._build_message_from_template(db, "ams_temperature_high", variables)
  1813. # Alarms always send immediately, bypassing digest mode
  1814. await self._send_to_providers(
  1815. providers,
  1816. title,
  1817. message,
  1818. db,
  1819. "ams_temperature_high",
  1820. printer_id,
  1821. printer_name,
  1822. force_immediate=True,
  1823. variables=variables,
  1824. )
  1825. async def on_ams_drying_suspended(
  1826. self,
  1827. printer_id: int,
  1828. printer_name: str,
  1829. ams_label: str,
  1830. humidity: float,
  1831. threshold: float,
  1832. cycles: int,
  1833. db: AsyncSession,
  1834. ):
  1835. """Handle automatic drying giving up on one AMS unit (#2770).
  1836. Sent immediately rather than folded into a digest: it reports that
  1837. Bambuddy has STOPPED doing something, and a report of inaction that
  1838. arrives with tomorrow's summary has already cost the user a day.
  1839. """
  1840. providers = await self._get_providers_for_event(db, "on_ams_drying_suspended", printer_id)
  1841. if not providers:
  1842. return
  1843. variables = {
  1844. "printer": printer_name,
  1845. "ams_label": ams_label,
  1846. "humidity": f"{humidity:.0f}",
  1847. "threshold": f"{threshold:.0f}",
  1848. "cycles": str(cycles),
  1849. }
  1850. title, message = await self._build_message_from_template(db, "ams_drying_suspended", variables)
  1851. await self._send_to_providers(
  1852. providers,
  1853. title,
  1854. message,
  1855. db,
  1856. "ams_drying_suspended",
  1857. printer_id,
  1858. printer_name,
  1859. force_immediate=True,
  1860. variables=variables,
  1861. )
  1862. async def on_ams_ht_humidity_high(
  1863. self,
  1864. printer_id: int,
  1865. printer_name: str,
  1866. ams_label: str,
  1867. humidity: float,
  1868. threshold: float,
  1869. db: AsyncSession,
  1870. ):
  1871. """Handle AMS-HT high humidity alarm event. Always sends immediately (bypasses digest)."""
  1872. providers = await self._get_providers_for_event(db, "on_ams_ht_humidity_high", printer_id)
  1873. if not providers:
  1874. return
  1875. variables = {
  1876. "printer": printer_name,
  1877. "ams_label": ams_label,
  1878. "humidity": f"{humidity:.0f}",
  1879. "threshold": f"{threshold:.0f}",
  1880. }
  1881. # Use the same template as regular AMS (can create separate templates later if needed)
  1882. title, message = await self._build_message_from_template(db, "ams_humidity_high", variables)
  1883. # Alarms always send immediately, bypassing digest mode
  1884. await self._send_to_providers(
  1885. providers,
  1886. title,
  1887. message,
  1888. db,
  1889. "ams_ht_humidity_high",
  1890. printer_id,
  1891. printer_name,
  1892. force_immediate=True,
  1893. variables=variables,
  1894. )
  1895. async def on_ams_ht_temperature_high(
  1896. self,
  1897. printer_id: int,
  1898. printer_name: str,
  1899. ams_label: str,
  1900. temperature: float,
  1901. threshold: float,
  1902. db: AsyncSession,
  1903. ):
  1904. """Handle AMS-HT high temperature alarm event. Always sends immediately (bypasses digest)."""
  1905. providers = await self._get_providers_for_event(db, "on_ams_ht_temperature_high", printer_id)
  1906. if not providers:
  1907. return
  1908. variables = {
  1909. "printer": printer_name,
  1910. "ams_label": ams_label,
  1911. "temperature": f"{temperature:.1f}",
  1912. "threshold": f"{threshold:.1f}",
  1913. }
  1914. # Use the same template as regular AMS (can create separate templates later if needed)
  1915. title, message = await self._build_message_from_template(db, "ams_temperature_high", variables)
  1916. # Alarms always send immediately, bypassing digest mode
  1917. await self._send_to_providers(
  1918. providers,
  1919. title,
  1920. message,
  1921. db,
  1922. "ams_ht_temperature_high",
  1923. printer_id,
  1924. printer_name,
  1925. force_immediate=True,
  1926. variables=variables,
  1927. )
  1928. async def on_bed_cooled(
  1929. self,
  1930. printer_id: int,
  1931. printer_name: str,
  1932. bed_temp: float,
  1933. threshold: float,
  1934. filename: str,
  1935. db: AsyncSession,
  1936. ):
  1937. """Handle bed cooled event - bed temperature dropped below threshold after print."""
  1938. providers = await self._get_providers_for_event(db, "on_bed_cooled", printer_id)
  1939. if not providers:
  1940. return
  1941. variables = {
  1942. "printer": printer_name,
  1943. "bed_temp": f"{bed_temp:.0f}",
  1944. "threshold": f"{threshold:.0f}",
  1945. "filename": self._clean_filename(filename) if filename else "Unknown",
  1946. }
  1947. title, message = await self._build_message_from_template(db, "bed_cooled", variables)
  1948. await self._send_to_providers(
  1949. providers, title, message, db, "bed_cooled", printer_id, printer_name, variables=variables
  1950. )
  1951. async def on_ha_sensor_alert(
  1952. self,
  1953. printer_id: int,
  1954. printer_name: str,
  1955. sensor_name: str,
  1956. state: str,
  1957. db: AsyncSession,
  1958. ):
  1959. """A Home Assistant sensor bound to a printer entered its alert state (#1148).
  1960. Sent immediately rather than folded into a digest: the case this exists
  1961. for is an enclosure door left open, which is only worth telling someone
  1962. about while they can still act on it.
  1963. """
  1964. providers = await self._get_providers_for_event(db, "on_ha_sensor_alert", printer_id)
  1965. if not providers:
  1966. return
  1967. variables = {
  1968. "printer": printer_name,
  1969. "sensor": sensor_name,
  1970. "state": state,
  1971. }
  1972. title, message = await self._build_message_from_template(db, "ha_sensor_alert", variables)
  1973. await self._send_to_providers(
  1974. providers,
  1975. title,
  1976. message,
  1977. db,
  1978. "ha_sensor_alert",
  1979. printer_id,
  1980. printer_name,
  1981. force_immediate=True,
  1982. variables=variables,
  1983. )
  1984. async def on_location_ha_sensor_alert(
  1985. self,
  1986. location_name: str,
  1987. sensor_name: str,
  1988. state: str,
  1989. db: AsyncSession,
  1990. ):
  1991. """A Home Assistant sensor bound to a storage location entered its alert state (#2824).
  1992. Sent immediately rather than folded into a digest, for the same reason
  1993. as on_ha_sensor_alert above: this is the "drybox went stale" case, only
  1994. worth acting on while the humidity/temperature is still climbing.
  1995. """
  1996. # Own column, not on_ha_sensor_alert (#2824): that one can be scoped to
  1997. # a single printer via provider.printer_id, and a location alert has no
  1998. # printer to scope by, so sharing it would leak drybox alerts to a
  1999. # provider narrowed to one printer's sensors.
  2000. providers = await self._get_providers_for_event(db, "on_location_ha_sensor_alert", None)
  2001. if not providers:
  2002. return
  2003. variables = {
  2004. "location": location_name,
  2005. "sensor": sensor_name,
  2006. "state": state,
  2007. }
  2008. title, message = await self._build_message_from_template(db, "location_ha_sensor_alert", variables)
  2009. await self._send_to_providers(
  2010. providers,
  2011. title,
  2012. message,
  2013. db,
  2014. "location_ha_sensor_alert",
  2015. force_immediate=True,
  2016. variables=variables,
  2017. )
  2018. async def on_first_layer_complete(
  2019. self,
  2020. printer_id: int,
  2021. printer_name: str,
  2022. filename: str,
  2023. total_layers: int,
  2024. db: AsyncSession,
  2025. image_data: bytes | None = None,
  2026. ):
  2027. """Handle first layer complete event."""
  2028. providers = await self._get_providers_for_event(db, "on_first_layer_complete", printer_id)
  2029. if not providers:
  2030. return
  2031. variables = {
  2032. "printer": printer_name,
  2033. "filename": self._clean_filename(filename),
  2034. "total_layers": str(total_layers),
  2035. }
  2036. title, message = await self._build_message_from_template(db, "first_layer_complete", variables)
  2037. await self._send_to_providers(
  2038. providers,
  2039. title,
  2040. message,
  2041. db,
  2042. "first_layer_complete",
  2043. printer_id,
  2044. printer_name,
  2045. image_data=image_data,
  2046. variables=variables,
  2047. )
  2048. def clear_template_cache(self):
  2049. """Clear the template cache. Call this when templates are updated."""
  2050. self._template_cache.clear()
  2051. async def send_user_print_email(
  2052. self,
  2053. event_type: str,
  2054. created_by_id: int | None,
  2055. printer_name: str,
  2056. filename: str,
  2057. db: AsyncSession,
  2058. ) -> None:
  2059. """Send a print event email notification to the user who submitted the job.
  2060. Args:
  2061. event_type: 'user_print_start', 'user_print_complete', 'user_print_failed', or 'user_print_stopped'
  2062. created_by_id: User ID who submitted the print job (from archive)
  2063. printer_name: Name of the printer
  2064. filename: Raw filename or subtask name
  2065. db: Database session
  2066. """
  2067. if created_by_id is None:
  2068. logger.debug("[EMAIL] Skipping user print email (%s): no created_by_id", event_type)
  2069. return
  2070. try:
  2071. # Check if advanced auth is enabled - required for user email notifications
  2072. from backend.app.models.settings import Settings
  2073. result = await db.execute(select(Settings).where(Settings.key == "advanced_auth_enabled"))
  2074. setting = result.scalar_one_or_none()
  2075. if not setting or setting.value.lower() != "true":
  2076. logger.debug("[EMAIL] Skipping user print email (%s): advanced_auth not enabled", event_type)
  2077. return
  2078. # Check if user notifications are enabled (admin-controlled toggle)
  2079. notif_enabled_result = await db.execute(
  2080. select(Settings).where(Settings.key == "user_notifications_enabled")
  2081. )
  2082. notif_enabled_setting = notif_enabled_result.scalar_one_or_none()
  2083. if notif_enabled_setting and notif_enabled_setting.value.lower() == "false":
  2084. logger.debug("[EMAIL] Skipping user print email (%s): user_notifications_enabled is false", event_type)
  2085. return
  2086. # Check SMTP settings are configured - required for sending emails
  2087. from backend.app.services.email_service import get_smtp_settings, send_user_print_notification
  2088. smtp_settings = await get_smtp_settings(db)
  2089. if not smtp_settings:
  2090. logger.debug("[EMAIL] Skipping user print email (%s): SMTP settings not configured", event_type)
  2091. return
  2092. # Load user preferences
  2093. from backend.app.models.user import User
  2094. from backend.app.models.user_email_pref import UserEmailPreference
  2095. user_result = await db.execute(select(User).where(User.id == created_by_id))
  2096. user = user_result.scalar_one_or_none()
  2097. if user is None or not user.email:
  2098. logger.debug(
  2099. "[EMAIL] Skipping user print email (%s): user %s not found or has no email address",
  2100. event_type,
  2101. created_by_id,
  2102. )
  2103. return
  2104. # Load user's notification preferences
  2105. pref_result = await db.execute(
  2106. select(UserEmailPreference).where(UserEmailPreference.user_id == created_by_id)
  2107. )
  2108. pref = pref_result.scalar_one_or_none()
  2109. # Determine if this event type should be sent
  2110. should_send = False
  2111. if event_type == "user_print_start":
  2112. should_send = pref is None or pref.notify_print_start
  2113. elif event_type == "user_print_complete":
  2114. should_send = pref is None or pref.notify_print_complete
  2115. elif event_type == "user_print_failed":
  2116. should_send = pref is None or pref.notify_print_failed
  2117. elif event_type == "user_print_stopped":
  2118. should_send = pref is None or pref.notify_print_stopped
  2119. if not should_send:
  2120. logger.debug(
  2121. "[EMAIL] Skipping user print email (%s): user %s has notifications disabled for this event",
  2122. event_type,
  2123. created_by_id,
  2124. )
  2125. return
  2126. logger.info(
  2127. "[EMAIL] Sending user print email: event=%s, user=%s (%s), printer=%s, file=%s",
  2128. event_type,
  2129. user.username,
  2130. user.email,
  2131. printer_name,
  2132. filename,
  2133. )
  2134. # Build variables
  2135. variables = {
  2136. "printer": printer_name,
  2137. "filename": self._clean_filename(filename),
  2138. }
  2139. # Send the email
  2140. await send_user_print_notification(
  2141. db=db,
  2142. event_type=event_type,
  2143. user_email=user.email,
  2144. username=user.username,
  2145. variables=variables,
  2146. )
  2147. logger.info("[EMAIL] User print email sent: event=%s → %s", event_type, user.email)
  2148. except Exception as e:
  2149. logger.warning("Failed to send user print email notification: %s", e, exc_info=True)
  2150. # ==================== Queue Notifications ====================
  2151. async def on_queue_job_added(
  2152. self,
  2153. job_name: str,
  2154. target: str,
  2155. db: AsyncSession,
  2156. printer_id: int | None = None,
  2157. printer_name: str | None = None,
  2158. ):
  2159. """Handle queue job added event."""
  2160. providers = await self._get_providers_for_event(db, "on_queue_job_added", printer_id)
  2161. if not providers:
  2162. return
  2163. variables = {
  2164. "job_name": job_name,
  2165. "target": target, # e.g., "Printer1" or "Any X1C"
  2166. "printer": printer_name or target,
  2167. }
  2168. title, message = await self._build_message_from_template(db, "queue_job_added", variables)
  2169. await self._send_to_providers(
  2170. providers, title, message, db, "queue_job_added", printer_id, printer_name, variables=variables
  2171. )
  2172. async def on_queue_job_assigned(
  2173. self,
  2174. job_name: str,
  2175. printer_id: int,
  2176. printer_name: str,
  2177. target_model: str,
  2178. db: AsyncSession,
  2179. ):
  2180. """Handle model-based job assigned to printer event."""
  2181. providers = await self._get_providers_for_event(db, "on_queue_job_assigned", printer_id)
  2182. if not providers:
  2183. return
  2184. variables = {
  2185. "job_name": job_name,
  2186. "printer": printer_name,
  2187. "target_model": target_model,
  2188. }
  2189. title, message = await self._build_message_from_template(db, "queue_job_assigned", variables)
  2190. await self._send_to_providers(
  2191. providers, title, message, db, "queue_job_assigned", printer_id, printer_name, variables=variables
  2192. )
  2193. async def on_queue_job_started(
  2194. self,
  2195. job_name: str,
  2196. printer_id: int,
  2197. printer_name: str,
  2198. db: AsyncSession,
  2199. estimated_time: int | None = None,
  2200. ):
  2201. """Handle queue job started printing event."""
  2202. providers = await self._get_providers_for_event(db, "on_queue_job_started", printer_id)
  2203. if not providers:
  2204. return
  2205. eta_str = await self._format_eta(estimated_time, db)
  2206. variables = {
  2207. "job_name": job_name,
  2208. "printer": printer_name,
  2209. "estimated_time": self._format_duration(estimated_time),
  2210. "eta": eta_str,
  2211. }
  2212. title, message = await self._build_message_from_template(db, "queue_job_started", variables)
  2213. await self._send_to_providers(
  2214. providers, title, message, db, "queue_job_started", printer_id, printer_name, variables=variables
  2215. )
  2216. async def on_queue_job_waiting(
  2217. self,
  2218. job_name: str,
  2219. target_model: str,
  2220. waiting_reason: str,
  2221. db: AsyncSession,
  2222. ):
  2223. """Handle job waiting for filament event."""
  2224. providers = await self._get_providers_for_event(db, "on_queue_job_waiting", None)
  2225. if not providers:
  2226. return
  2227. variables = {
  2228. "job_name": job_name,
  2229. "target_model": target_model,
  2230. "waiting_reason": waiting_reason,
  2231. }
  2232. title, message = await self._build_message_from_template(db, "queue_job_waiting", variables)
  2233. await self._send_to_providers(providers, title, message, db, "queue_job_waiting", variables=variables)
  2234. async def on_queue_job_skipped(
  2235. self,
  2236. job_name: str,
  2237. printer_id: int,
  2238. printer_name: str,
  2239. reason: str,
  2240. db: AsyncSession,
  2241. ):
  2242. """Handle job skipped event (e.g., previous print failed)."""
  2243. providers = await self._get_providers_for_event(db, "on_queue_job_skipped", printer_id)
  2244. if not providers:
  2245. return
  2246. variables = {
  2247. "job_name": job_name,
  2248. "printer": printer_name,
  2249. "reason": reason,
  2250. }
  2251. title, message = await self._build_message_from_template(db, "queue_job_skipped", variables)
  2252. await self._send_to_providers(
  2253. providers, title, message, db, "queue_job_skipped", printer_id, printer_name, variables=variables
  2254. )
  2255. async def on_queue_job_failed(
  2256. self,
  2257. job_name: str,
  2258. printer_id: int | None,
  2259. printer_name: str | None,
  2260. reason: str,
  2261. db: AsyncSession,
  2262. ):
  2263. """Handle job failed to start event (upload error, etc.)."""
  2264. providers = await self._get_providers_for_event(db, "on_queue_job_failed", printer_id)
  2265. if not providers:
  2266. return
  2267. variables = {
  2268. "job_name": job_name,
  2269. "printer": printer_name or "Unknown",
  2270. "reason": reason,
  2271. }
  2272. title, message = await self._build_message_from_template(db, "queue_job_failed", variables)
  2273. await self._send_to_providers(
  2274. providers, title, message, db, "queue_job_failed", printer_id, printer_name, variables=variables
  2275. )
  2276. async def on_queue_completed(
  2277. self,
  2278. completed_count: int,
  2279. db: AsyncSession,
  2280. ):
  2281. """Handle all queue jobs completed event."""
  2282. providers = await self._get_providers_for_event(db, "on_queue_completed", None)
  2283. if not providers:
  2284. return
  2285. variables = {
  2286. "completed_count": str(completed_count),
  2287. }
  2288. title, message = await self._build_message_from_template(db, "queue_completed", variables)
  2289. await self._send_to_providers(providers, title, message, db, "queue_completed", variables=variables)
  2290. # ==================== Inventory Stock Alerts ====================
  2291. async def on_stock_reorder_alert(
  2292. self,
  2293. material: str,
  2294. brand: str | None,
  2295. stock_g: float,
  2296. rate_g_day: float,
  2297. days_left: int,
  2298. db: AsyncSession,
  2299. *,
  2300. subtype: str | None = None,
  2301. color: str | None = None,
  2302. skip_break_subscribers: bool = False,
  2303. ):
  2304. """Fire when an inventory SKU reaches its reorder point.
  2305. ``subtype`` and ``color`` are what tell apart the messages for two colours of one
  2306. product, which the forecast (grouped by colour as well) reports separately.
  2307. A SKU in stock break has also reached its reorder point. ``skip_break_subscribers``
  2308. leaves out the providers that have the break alert on, so the producer can send
  2309. the reorder event for a break SKU without telling those providers twice.
  2310. """
  2311. providers = await self._get_providers_for_event(db, "on_stock_reorder_alert", None)
  2312. if skip_break_subscribers:
  2313. providers = [p for p in providers if not p.on_stock_break_alert]
  2314. if not providers:
  2315. return
  2316. variables = {
  2317. "material": material,
  2318. "subtype": subtype or "",
  2319. "brand": brand or "",
  2320. "color": color or "",
  2321. "stock_g": f"{stock_g:.0f}",
  2322. "rate_g_day": f"{rate_g_day:.1f}",
  2323. "days_left": str(days_left),
  2324. }
  2325. title, message = await self._build_message_from_template(db, "stock_reorder_alert", variables)
  2326. await self._send_to_providers(providers, title, message, db, "stock_reorder_alert", variables=variables)
  2327. async def on_stock_break_alert(
  2328. self,
  2329. material: str,
  2330. brand: str | None,
  2331. stock_g: float,
  2332. rate_g_day: float,
  2333. days_left: int,
  2334. lead_time_days: int,
  2335. db: AsyncSession,
  2336. *,
  2337. subtype: str | None = None,
  2338. color: str | None = None,
  2339. ):
  2340. """Fire when a stock break is detected (stock runs out before lead time).
  2341. ``subtype`` and ``color`` as for :meth:`on_stock_reorder_alert`.
  2342. """
  2343. providers = await self._get_providers_for_event(db, "on_stock_break_alert", None)
  2344. if not providers:
  2345. return
  2346. variables = {
  2347. "material": material,
  2348. "subtype": subtype or "",
  2349. "brand": brand or "",
  2350. "color": color or "",
  2351. "stock_g": f"{stock_g:.0f}",
  2352. "rate_g_day": f"{rate_g_day:.1f}",
  2353. "days_left": str(days_left),
  2354. "lead_time_days": str(lead_time_days),
  2355. }
  2356. title, message = await self._build_message_from_template(db, "stock_break_alert", variables)
  2357. await self._send_to_providers(providers, title, message, db, "stock_break_alert", variables=variables)
  2358. async def _queue_for_digest(
  2359. self,
  2360. provider: NotificationProvider,
  2361. event_type: str,
  2362. title: str,
  2363. message: str,
  2364. db: AsyncSession,
  2365. printer_id: int | None = None,
  2366. printer_name: str | None = None,
  2367. ):
  2368. """Queue a notification for later delivery in the daily digest."""
  2369. try:
  2370. queue_entry = NotificationDigestQueue(
  2371. provider_id=provider.id,
  2372. event_type=event_type,
  2373. title=title,
  2374. message=message,
  2375. printer_id=printer_id,
  2376. printer_name=printer_name,
  2377. )
  2378. db.add(queue_entry)
  2379. await db.commit()
  2380. logger.info("Queued notification for digest: %s for provider %s", event_type, provider.name)
  2381. except Exception as e:
  2382. logger.warning("Failed to queue notification for digest: %s", e)
  2383. async def send_digest(self, provider_id: int):
  2384. """Send all queued notifications as a single digest for a provider."""
  2385. from backend.app.core.database import async_session
  2386. async with async_session() as db:
  2387. # Get the provider
  2388. result = await db.execute(select(NotificationProvider).where(NotificationProvider.id == provider_id))
  2389. provider = result.scalar_one_or_none()
  2390. if not provider or not provider.enabled:
  2391. return
  2392. # Get all queued notifications for this provider
  2393. result = await db.execute(
  2394. select(NotificationDigestQueue)
  2395. .where(NotificationDigestQueue.provider_id == provider_id)
  2396. .order_by(NotificationDigestQueue.created_at)
  2397. )
  2398. queue_entries = list(result.scalars().all())
  2399. if not queue_entries:
  2400. logger.debug("No queued notifications for provider %s", provider.name)
  2401. return
  2402. # Build digest message
  2403. title = f"Daily Digest - {len(queue_entries)} Events"
  2404. # Group by event type
  2405. events_by_type: dict[str, list] = {}
  2406. for entry in queue_entries:
  2407. if entry.event_type not in events_by_type:
  2408. events_by_type[entry.event_type] = []
  2409. events_by_type[entry.event_type].append(entry)
  2410. # Format the digest body
  2411. body_parts = []
  2412. for event_type, entries in events_by_type.items():
  2413. event_label = event_type.replace("_", " ").title()
  2414. body_parts.append(f"== {event_label} ({len(entries)}) ==")
  2415. for entry in entries:
  2416. time_str = entry.created_at.strftime("%H:%M")
  2417. printer_info = f"[{entry.printer_name}] " if entry.printer_name else ""
  2418. body_parts.append(f" {time_str} {printer_info}{entry.title}")
  2419. body_parts.append("")
  2420. body = "\n".join(body_parts)
  2421. # Send the digest
  2422. success, error = await self._send_to_provider(provider, title, body, db)
  2423. # Log the digest
  2424. await self._log_notification(
  2425. db=db,
  2426. provider_id=provider.id,
  2427. event_type="daily_digest",
  2428. title=title,
  2429. message=body,
  2430. success=success,
  2431. error_message=error if not success else None,
  2432. )
  2433. # Clear the queue
  2434. for entry in queue_entries:
  2435. await db.delete(entry)
  2436. await db.commit()
  2437. if success:
  2438. logger.info("Sent daily digest with %s events to %s", len(queue_entries), provider.name)
  2439. else:
  2440. logger.warning("Failed to send daily digest to %s: %s", provider.name, error)
  2441. async def check_and_send_digests(self):
  2442. """Check all providers and send digests if it's their scheduled time."""
  2443. from backend.app.core.database import async_session
  2444. current_time = datetime.now().strftime("%H:%M")
  2445. # Avoid duplicate checks within the same minute
  2446. if current_time == self._last_digest_check:
  2447. return
  2448. self._last_digest_check = current_time
  2449. async with async_session() as db:
  2450. # Find all providers with digest enabled at this time
  2451. result = await db.execute(
  2452. select(NotificationProvider).where(
  2453. NotificationProvider.enabled.is_(True),
  2454. NotificationProvider.daily_digest_enabled.is_(True),
  2455. NotificationProvider.daily_digest_time == current_time,
  2456. )
  2457. )
  2458. providers = result.scalars().all()
  2459. for provider in providers:
  2460. try:
  2461. await self.send_digest(provider.id)
  2462. except Exception as e:
  2463. logger.error("Error sending digest for provider %s: %s", provider.id, e)
  2464. def start_digest_scheduler(self):
  2465. """Start the background scheduler for daily digest notifications."""
  2466. if self._digest_scheduler_task is None:
  2467. self._digest_scheduler_task = asyncio.create_task(self._digest_scheduler_loop())
  2468. logger.info("Notification digest scheduler started")
  2469. def stop_digest_scheduler(self):
  2470. """Stop the background scheduler for daily digests."""
  2471. if self._digest_scheduler_task:
  2472. self._digest_scheduler_task.cancel()
  2473. self._digest_scheduler_task = None
  2474. logger.info("Notification digest scheduler stopped")
  2475. async def _digest_scheduler_loop(self):
  2476. """Background loop that checks for scheduled digests every minute."""
  2477. while True:
  2478. try:
  2479. await self.check_and_send_digests()
  2480. except Exception as e:
  2481. logger.error("Error in digest scheduler: %s", e)
  2482. # Wait until the next minute
  2483. await asyncio.sleep(60)
  2484. # Global instance
  2485. notification_service = NotificationService()