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

1import asyncio 

2import json 

3import logging 

4import threading 

5from abc import abstractmethod 

6from dataclasses import dataclass 

7from typing import TYPE_CHECKING, Any 

8 

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 

21 

22from custom_components.remote_logger.const import BATCH_FLUSH_INTERVAL_SECONDS 

23 

24if TYPE_CHECKING: 

25 import datetime as dt 

26 from collections.abc import Mapping 

27 

28_LOGGER = logging.getLogger(__name__) 

29 

30_HA_EVENT_BODY_ATTRIBUTE_KEYS: frozenset[str] = frozenset({"entity_id", "domain", "service"}) 

31 

32 

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) 

40 

41 

42@dataclass 

43class LogMessage: 

44 payload: Any 

45 sent: bool = False 

46 

47 

48class LogSubmission: 

49 @abstractmethod 

50 def for_display(self) -> dict[str, Any]: 

51 pass 

52 

53 

54class LogExporter: 

55 """Base class for log exporters""" 

56 

57 logger_type: str 

58 

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() 

64 

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 

76 

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() 

81 

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() 

87 

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) 

102 

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)) 

108 

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) 

117 

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)) 

126 

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 

138 

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] 

190 

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) 

198 

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)) 

213 

214 @callback 

215 @abstractmethod 

216 async def flush(self) -> None: 

217 pass 

218 

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 

229 

230 async def close(self) -> None: 

231 """Clean up resources (no-op for HTTP-based exporter).""" 

232 

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() 

237 

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() 

242 

243 def on_success(self) -> None: 

244 self.posting_count += 1 

245 self.last_posting = dt_util.now() 

246 

247 def on_event(self) -> None: 

248 self.event_count += 1 

249 self.last_event = dt_util.now() 

250 

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."""