Coverage for opt/mealie/lib/python3.12/site-packages/mealie/services/event_bus_service/event_bus_service.py: 89%
51 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 03:04 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 03:04 +0000
1from fastapi import BackgroundTasks, Depends
2from pydantic import UUID4
3from sqlalchemy.orm.session import Session
5from mealie.core.config import get_app_settings
6from mealie.db.db_setup import generate_session
7from mealie.repos.all_repositories import get_repositories
8from mealie.schema.response.pagination import PaginationQuery
9from mealie.services.event_bus_service.event_bus_listeners import (
10 AppriseEventListener,
11 EventListenerBase,
12 WebhookEventListener,
13)
15from .event_types import Event, EventBusMessage, EventDocumentDataBase, EventTypes
17settings = get_app_settings()
18ALGORITHM = "HS256"
21class EventSource:
22 event_type: str
23 item_type: str
24 item_id: UUID4 | int
25 kwargs: dict
27 def __init__(self, event_type: str, item_type: str, item_id: UUID4 | int, **kwargs) -> None:
28 self.event_type = event_type
29 self.item_type = item_type
30 self.item_id = item_id
31 self.kwargs = kwargs
33 def dict(self) -> dict:
34 return {
35 "event_type": self.event_type,
36 "item_type": self.item_type,
37 "item_id": str(self.item_id),
38 **self.kwargs,
39 }
42class EventBusService:
43 bg: BackgroundTasks | None = None
44 session: Session | None = None
46 def __init__(
47 self,
48 bg: BackgroundTasks | None = None,
49 session: Session | None = None,
50 ) -> None:
51 self.bg = bg
52 self.session = session
54 def _get_listeners(self, group_id: UUID4, household_id: UUID4) -> list[EventListenerBase]:
55 return [
56 AppriseEventListener(group_id, household_id),
57 WebhookEventListener(group_id, household_id),
58 ]
60 def _publish_event(self, event: Event, group_id: UUID4, household_id: UUID4) -> None:
61 """Publishes the event to all listeners"""
62 for listener in self._get_listeners(group_id, household_id):
63 if subscribers := listener.get_subscribers(event):
64 listener.publish_to_subscribers(event, subscribers)
66 def dispatch(
67 self,
68 integration_id: str,
69 group_id: UUID4,
70 household_id: UUID4 | None,
71 event_type: EventTypes,
72 document_data: EventDocumentDataBase | None,
73 message: str = "",
74 ) -> None:
75 event = Event(
76 message=EventBusMessage.from_type(event_type, body=message),
77 event_type=event_type,
78 integration_id=integration_id,
79 document_data=document_data,
80 )
82 if not household_id:
83 if not self.session: 83 ↛ 84line 83 didn't jump to line 84 because the condition on line 83 was never true
84 raise ValueError("Session is required if household_id is not provided")
86 repos = get_repositories(self.session, group_id=group_id)
87 households = repos.households.page_all(PaginationQuery(page=1, per_page=-1)).items
88 household_ids = [household.id for household in households]
89 else:
90 household_ids = [household_id]
92 for household_id in household_ids:
93 if self.bg:
94 self.bg.add_task(self._publish_event, event=event, group_id=group_id, household_id=household_id)
95 else:
96 self._publish_event(event, group_id, household_id)
98 @classmethod
99 def as_dependency(
100 cls,
101 bg: BackgroundTasks,
102 session=Depends(generate_session),
103 ):
104 """Convenience method to use as a dependency in FastAPI routes"""
105 return cls(bg, session)