Coverage for custom_components/remote_logger/exporter.py: 87%
170 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-10-06 00:18 +0000
« prev ^ index » next coverage.py v7.15.4, created at 2026-10-06 00:18 +0000
1import asyncio
2import json
3import logging
4import threading
5from abc import abstractmethod
6from dataclasses import dataclass
7from typing import TYPE_CHECKING, Any
9from homeassistant.auth import EVENT_USER_ADDED, EVENT_USER_REMOVED, EVENT_USER_UPDATED, HomeAssistant
10from homeassistant.components.automation import EVENT_AUTOMATION_TRIGGERED
11from homeassistant.components.script import EVENT_SCRIPT_STARTED
12from homeassistant.const import EVENT_COMPONENT_LOADED, EVENT_STATE_CHANGED
13from homeassistant.core import EVENT_CALL_SERVICE, EVENT_SERVICE_REGISTERED, EVENT_SERVICE_REMOVED, Event, callback
14from homeassistant.helpers.area_registry import EVENT_AREA_REGISTRY_UPDATED
15from homeassistant.helpers.category_registry import EVENT_CATEGORY_REGISTRY_UPDATED
16from homeassistant.helpers.device_registry import EVENT_DEVICE_REGISTRY_UPDATED
17from homeassistant.helpers.entity_registry import EVENT_ENTITY_REGISTRY_UPDATED
18from homeassistant.helpers.floor_registry import EVENT_FLOOR_REGISTRY_UPDATED
19from homeassistant.helpers.label_registry import EVENT_LABEL_REGISTRY_UPDATED
20from homeassistant.util import dt as dt_util
22from custom_components.remote_logger.const import BATCH_FLUSH_INTERVAL_SECONDS
24if TYPE_CHECKING:
25 import datetime as dt
26 from collections.abc import Mapping
28_LOGGER = logging.getLogger(__name__)
30_HA_EVENT_BODY_ATTRIBUTE_KEYS: frozenset[str] = frozenset({"entity_id", "domain", "service"})
33def _event_data_serializer(obj: Any) -> Any:
34 """JSON serializer for HA event data objects."""
35 if hasattr(obj, "as_dict"):
36 return obj.as_dict()
37 if hasattr(obj, "value"):
38 return obj.value
39 return str(obj)
42@dataclass
43class LogMessage:
44 payload: Any
45 sent: bool = False
48class LogSubmission:
49 @abstractmethod
50 def for_display(self) -> dict[str, Any]:
51 pass
54class LogExporter:
55 """Base class for log exporters"""
57 logger_type: str
59 def __init__(self, hass: HomeAssistant) -> None:
60 self._hass: HomeAssistant = hass
61 self.name: str = self.logger_type
62 self.destination: tuple[str, ...]
63 self.tz = dt_util.get_default_time_zone()
65 self._batch_max_size: int
66 self.event_count: int = 0
67 self.last_event: dt.datetime | None = None
68 self.posting_count: int = 0
69 self.last_posting: dt.datetime | None = None
70 self.format_error_count: int = 0
71 self.last_format_error_message: str | None = None
72 self.last_format_error: dt.datetime | None = None
73 self.posting_error_count: int = 0
74 self.last_posting_error_message: str | None = None
75 self.last_posting_error: dt.datetime | None = None
77 self._buffer: list[LogMessage] = []
78 self.self_source: str = f"custom_components/remote_logger/{self.logger_type}"
79 self.last_sent_payload: LogSubmission | None = None
80 self.flushing: threading.Event = threading.Event()
82 async def disable_buffer(self) -> None:
83 """Flush logs and prevent future buffering, use for shutdowns"""
84 self._batch_max_size = 0
85 self.flushing.clear()
86 await self.flush()
88 @callback
89 def handle_event(self, event: Event) -> None:
90 self.on_event()
91 if (
92 event.data
93 and event.data.get("source")
94 and len(event.data["source"]) == 2
95 and self.self_source in event.data["source"][0]
96 ):
97 # prevent log loops
98 return
99 try:
100 record: LogMessage = self.create_log_record(event.data, event.event_type, event.time_fired)
101 self._buffer.append(record)
103 if len(self._buffer) >= self._batch_max_size:
104 self._hass.async_create_task(self.flush())
105 except Exception as e: # ruff: ignore[blind-except]
106 _LOGGER.error("remote_logger: %s event handler failure %s on %s", self.logger_type, e, event.data)
107 self.on_format_error(str(e))
109 def handle_entry(self, entry: Mapping[str, Any], time_fired: dt.datetime) -> None:
110 self.on_event()
111 if entry and entry.get("source") and len(entry["source"]) == 2 and self.self_source in entry["source"][0]:
112 # prevent log loops
113 return
114 try:
115 record: LogMessage = self.create_log_record(entry, None, time_fired)
116 self._buffer.append(record)
118 if len(self._buffer) >= self._batch_max_size:
119 try:
120 asyncio.get_running_loop().create_task(self.flush())
121 except RuntimeError:
122 self._hass.create_task(self.flush())
123 except Exception as e: # ruff: ignore[blind-except]
124 _LOGGER.error("remote_logger: %s entry handler failure %s on %s", self.logger_type, e, entry)
125 self.on_format_error(str(e))
127 @abstractmethod
128 def create_log_record(
129 self,
130 event_data: Mapping[str, Any],
131 event_type: str | None = None,
132 time_fired: dt.datetime | None = None,
133 message_override: list[str] | None = None,
134 level_override: str | None = None,
135 state_only: bool = False,
136 ) -> LogMessage:
137 pass
139 @callback
140 def handle_ha_event(self, event_type: str, event: Event, state_only: bool = False, event_body: bool = False) -> None:
141 """Handle a non-system-log HA event (lifecycle, core change, or custom)."""
142 self.on_event()
143 title_fields: list[str] = ["message", "id", "entity_id", "name", "component", "device_id"]
144 try:
145 if (
146 event_type == EVENT_CALL_SERVICE
147 and event.data.get("domain") == "system_log"
148 and event.data.get("service") == "write"
149 ):
150 # don't double count log events
151 return
152 if event_type == EVENT_STATE_CHANGED:
153 old_state: str = (event.data["old_state"] and event.data["old_state"].state) or "N/A"
154 new_state: str = (event.data["new_state"] and event.data["new_state"].state) or "N/A"
155 message: list[str] = [event_type, ":", event.data["entity_id"], old_state, "->", new_state]
156 elif event_type in (EVENT_CALL_SERVICE, EVENT_SERVICE_REGISTERED, EVENT_SERVICE_REMOVED):
157 message = [event_type, ":", event.data["domain"], event.data["service"]]
158 elif event_type == EVENT_COMPONENT_LOADED:
159 message = [event_type, ":", event.data["component"]]
160 elif event_type in (EVENT_SCRIPT_STARTED, EVENT_AUTOMATION_TRIGGERED):
161 message = [event_type, ":", event.data["entity_id"]]
162 elif event_type == EVENT_DEVICE_REGISTRY_UPDATED:
163 message = [event_type, ":", event.data["device_id"], event.data.get("action", "???")]
164 elif event_type == EVENT_ENTITY_REGISTRY_UPDATED:
165 message = [event_type, ":", event.data["entity_id"], event.data.get("action", "???")]
166 elif event_type == EVENT_LABEL_REGISTRY_UPDATED:
167 message = [event_type, ":", event.data["label_id"], event.data.get("action", "???")]
168 elif event_type == EVENT_AREA_REGISTRY_UPDATED:
169 message = [event_type, ":", event.data["area_id"], event.data.get("action", "???")]
170 elif event_type == EVENT_CATEGORY_REGISTRY_UPDATED:
171 message = [event_type, ":", event.data["category_id"], event.data.get("action", "???")]
172 elif event_type == EVENT_FLOOR_REGISTRY_UPDATED:
173 message = [event_type, ":", event.data["floor_id"], event.data.get("action", "???")]
174 elif event_type in (EVENT_USER_ADDED, EVENT_USER_REMOVED, EVENT_USER_UPDATED):
175 message = [event_type, ":", event.data["user_id"]]
176 elif event_type == "autoarm_change":
177 message = [
178 event_type,
179 ":",
180 event.data["original_state"],
181 "->",
182 event.data["new_state"],
183 " from ",
184 event.data["change_source"],
185 ]
186 elif any(v in event.data for v in title_fields):
187 message = [event_type, ":"] + [event.data[v] for v in title_fields if v in event.data]
188 else:
189 message = [event_type]
191 if event_body:
192 flat_event: dict[str, Any] = {k: v for k, v in event.data.items() if k not in _HA_EVENT_BODY_ATTRIBUTE_KEYS}
193 event_data: dict[str, Any] = {k: v for k, v in event.data.items() if k in _HA_EVENT_BODY_ATTRIBUTE_KEYS}
194 if flat_event:
195 event_data["event.data"] = json.dumps(flat_event, default=_event_data_serializer, indent=2)
196 else:
197 event_data = dict(event.data)
199 record: LogMessage = self.create_log_record(
200 event_data,
201 event.event_type,
202 event.time_fired,
203 message_override=message,
204 level_override="INFO",
205 state_only=state_only,
206 )
207 self._buffer.append(record)
208 if len(self._buffer) >= self._batch_max_size:
209 self._hass.async_create_task(self.flush())
210 except Exception as e: # ruff: ignore[blind-except]
211 _LOGGER.error("remote_logger: %s ha_event handler failure %s on %s", self.logger_type, e, event_type)
212 self.on_format_error(str(e))
214 @callback
215 @abstractmethod
216 async def flush(self) -> None:
217 pass
219 async def flush_loop(self) -> None:
220 """Periodically flush buffered log records."""
221 self.flushing.set()
222 try:
223 while self.flushing.is_set():
224 await asyncio.sleep(BATCH_FLUSH_INTERVAL_SECONDS)
225 await self.flush()
226 except asyncio.CancelledError:
227 _LOGGER.debug("AUTOARM log flush cancelled")
228 raise
230 async def close(self) -> None:
231 """Clean up resources (no-op for HTTP-based exporter)."""
233 def on_format_error(self, message: str) -> None:
234 self.format_error_count += 1
235 self.last_format_error_message = message
236 self.last_format_error = dt_util.now()
238 def on_posting_error(self, message: str) -> None:
239 self.posting_error_count += 1
240 self.last_posting_error_message = message
241 self.last_posting_error = dt_util.now()
243 def on_success(self) -> None:
244 self.posting_count += 1
245 self.last_posting = dt_util.now()
247 def on_event(self) -> None:
248 self.event_count += 1
249 self.last_event = dt_util.now()
251 @abstractmethod
252 def log_direct(self, event_name: str, message: str, level: str, attributes: dict[str, Any] | None = None) -> None:
253 """Buffer a custom syslog record without requiring a HA Event."""