Coverage for clients/__init__.py: 82%
896 statements
« prev ^ index » next coverage.py v7.14.1, created at 2026-06-11 15:34 +0000
« prev ^ index » next coverage.py v7.14.1, created at 2026-06-11 15:34 +0000
1#!/usr/bin/env python3
2# -*- coding: utf-8 -*-
4# Hermes : Change Data Capture (CDC) tool from any source(s) to any target
5# Copyright (C) 2023 INSA Strasbourg
6#
7# This file is part of Hermes.
8#
9# Hermes is free software: you can redistribute it and/or modify
10# it under the terms of the GNU General Public License as published by
11# the Free Software Foundation, either version 3 of the License, or
12# (at your option) any later version.
13#
14# Hermes is distributed in the hope that it will be useful,
15# but WITHOUT ANY WARRANTY; without even the implied warranty of
16# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
17# GNU General Public License for more details.
18#
19# You should have received a copy of the GNU General Public License
20# along with Hermes. If not, see <https://www.gnu.org/licenses/>.
23from clients.datamodel import Datamodel, InvalidDataError
24from lib.config import HermesConfig
25from lib.version import HERMES_VERSION
26from lib.datamodel.dataobject import DataObject
27from lib.datamodel.dataobjectlist import DataObjectList
28from lib.datamodel.datasource import Dataschema, Datasource
29from lib.datamodel.diffobject import DiffObject
30from lib.datamodel.event import Event
31from lib.datamodel.serialization import LocalCache, JSONEncoder
32from lib.plugins import AbstractMessageBusConsumerPlugin
33from lib.utils.mail import Email
34from lib.utils.socket import (
35 SockServer,
36 SocketMessageToServer,
37 SocketMessageToClient,
38 SocketArgumentParser,
39 SocketParsingError,
40 SocketParsingMessage,
41)
43from copy import deepcopy
44from datetime import datetime, timedelta
45from time import sleep
46from types import FrameType
47from typing import Any
48import argparse
49import json
50import signal
51import traceback
54class HermesAlreadyNotifiedException(Exception):
55 """Raised when an exception has already been notified, to avoid a second
56 notification"""
59class HermesDatamodelUpdatePurgeError(Exception):
60 """Raised when an error is met when purging data on datamodel update"""
63class HermesClientHandlerError(Exception):
64 """Raised when an exception is met during client handler call"""
66 def __init__(self, err: Exception | str | None):
67 if isinstance(err, Exception):
68 self.msg: str | None = HermesClientHandlerError.exceptionToString(err)
69 """Printable error message, None in a rare case where exception is raised
70 without error, but to postpone an event processing"""
71 else:
72 self.msg = err
73 super().__init__(self.msg)
75 @staticmethod
76 def exceptionToString(
77 exception: Exception, purgeCurrentFileFromTrace: bool = True
78 ) -> str:
79 """Convert the specified exception to a string containing its full trace"""
80 lines = traceback.format_exception(
81 type(exception), exception, exception.__traceback__
82 )
84 if purgeCurrentFileFromTrace:
85 # Purging current file infos from traceback
86 lines = [line for line in lines if __file__ not in line]
88 return "".join(lines).strip()
91class HermesClientCache(LocalCache):
92 """Hermes client data to cache"""
94 def __init__(self, from_json_dict: dict[str, Any] = {}):
95 super().__init__(
96 jsondataattr=[
97 "queueErrors",
98 "datamodelWarnings",
99 "exception",
100 "initstartoffset",
101 "initstopoffset",
102 "nextoffset",
103 ],
104 )
106 self.queueErrors: dict[str, str] = from_json_dict.get("queueErrors", {})
107 """Dictionary containing current objects in error queue, for notifications"""
109 self.datamodelWarnings: dict[str, dict[str, dict[str, Any]]] = (
110 from_json_dict.get("datamodelWarnings", {})
111 )
112 """Dictionary containing current datamodel warnings, for notifications"""
114 self.exception: str | None = from_json_dict.get("exception")
115 """String containing latest exception trace"""
117 self.initstartoffset: Any | None = from_json_dict.get("initstartoffset")
118 """Contains the offset of the first message of initSync sequence on message
119 bus"""
120 self.initstopoffset: Any | None = from_json_dict.get("initstopoffset")
121 """Contains the offset of the last message of initSync sequence on message
122 bus"""
123 self.nextoffset: Any | None = from_json_dict.get("nextoffset")
124 """Contains the offset of the next message to process on message bus"""
126 def savecachefile(self, cacheFilename: str | None = None):
127 """Override method only to disable backup files in cache"""
128 return super().savecachefile(cacheFilename, dontKeepBackup=True)
131class GenericClient:
132 """Superclass of all hermes-client implementations.
133 Manage all the internals of hermes for a client: datamodel updates, caching, error
134 management, trashbin and converting messages from message bus into events handlers
135 calls"""
137 __FOREIGNKEYS_POLICIES: dict[str, tuple[str]] = {
138 "disabled": tuple(),
139 "on_remove_event": ("removed",),
140 "on_every_event": ("added", "modified", "removed"),
141 }
142 """Different foreignkeys_policy settings : associate each foreignkeys_policy
143 (as key) with the list of event types that will be placed in the error queue if the
144 object concerning them is the parent (by foreign key) of an object already present
145 in the error queue"""
147 def __init__(self, config: HermesConfig):
148 """Instantiate a new client"""
150 __hermes__.logger.info(f"Starting {config['appname']} v{HERMES_VERSION}")
152 # Setup the signals handler
153 config.setSignalsHandler(self.__signalHandler)
155 self.__config: HermesConfig = config
156 """Current config"""
157 try:
158 self.config: dict[str, Any] = self.__config[self.__config["appname"]]
159 """Dict containing the client plugin configuration"""
160 except KeyError:
161 self.config = {}
163 self.__previousconfig: HermesConfig = HermesConfig.loadcachefile(
164 "_hermesconfig"
165 )
166 """Previous config (from cache)"""
168 self.__msgbus: AbstractMessageBusConsumerPlugin = self.__config["hermes"][
169 "plugins"
170 ]["messagebus"]["plugininstance"]
171 self.__msgbus.setTimeout(
172 self.__config["hermes-client"]["updateInterval"] * 1000
173 )
175 self.__cache: HermesClientCache = HermesClientCache.loadcachefile(
176 f"_{self.__config['appname']}"
177 )
178 """Cached attributes"""
179 self.__cache.setCacheFilename(f"_{self.__config['appname']}")
181 self.__startTime: datetime | None = None
182 """Datetime when mainloop was started"""
184 self.__isPaused: datetime | None = None
185 """Contains pause datetime if standard processing is paused, None otherwise"""
187 self.__isStopped: bool = False
188 """mainloop() will run until this var is set to True"""
190 self.__numberOfLoopToProcess: int | None = None
191 """**For functionnal tests only**, if a value is set, will process for *value*
192 iterations of mainloop and pause execution until a new positive value is set"""
194 self.__sock: SockServer | None = None
195 """Facultative socket to allow cli communication"""
196 if (
197 self.__config["hermes"]["cli_socket"]["path"] is not None
198 or self.__config["hermes"]["cli_socket"]["dont_manage_sockfile"] is not None
199 ):
200 self.__sock = SockServer(
201 path=self.__config["hermes"]["cli_socket"]["path"],
202 owner=self.__config["hermes"]["cli_socket"]["owner"],
203 group=self.__config["hermes"]["cli_socket"]["group"],
204 mode=self.__config["hermes"]["cli_socket"]["mode"],
205 processHdlr=self.__processSocketMessage,
206 dontManageSockfile=self.__config["hermes"]["cli_socket"][
207 "dont_manage_sockfile"
208 ],
209 )
210 self.__setupSocketParser()
212 self.__useFirstInitsyncSequence: bool = self.__config["hermes-client"][
213 "useFirstInitsyncSequence"
214 ]
215 """Indicate if we prefer using the first/oldest or last/most recent
216 initsync sequence available on message bus"""
218 self.__newdatamodel: Datamodel = Datamodel(config=self.__config)
219 """New datamodel (from current config)"""
221 if self.__previousconfig.hasData():
222 # Start app with previous datamodel to be able to check differences with
223 # new one through self.__processDatamodelUpdate()
224 self.__datamodel: Datamodel = Datamodel(config=self.__previousconfig)
225 """Current datamodel"""
226 else:
227 # As no previous datamodel is available, start app with new one
228 self.__datamodel = self.__newdatamodel
230 self.__trashbin_retention: timedelta | None = None
231 """Timedelta with delay to keep removed data in trashbin before permanently
232 deleting it. 'None' means no trashbin"""
233 if self.__config["hermes-client"]["trashbin_retention"] > 0:
234 self.__trashbin_retention = timedelta(
235 days=self.__config["hermes-client"]["trashbin_retention"]
236 )
238 self.__trashbin_purgeInterval: timedelta | None = timedelta(
239 minutes=self.__config["hermes-client"]["trashbin_purgeInterval"]
240 )
241 """Timedelta with delay between two trashbin purge attempts"""
243 self.__trashbin_lastpurge: datetime = datetime(year=1, month=1, day=1)
244 """Datetime when latest trashbin purge was ran"""
246 self.__errorQueue_retryInterval: timedelta | None = timedelta(
247 minutes=self.__config["hermes-client"]["errorQueue_retryInterval"]
248 )
249 """Timedelta with delay between two attempts of processing events in error"""
251 self.__errorQueue_lastretry: datetime = datetime(year=1, month=1, day=1)
252 """Datetime when latest error queue retry was ran"""
254 self.__currentStep: int = 0
255 """Store the step number of current event processing. Will be stored in events
256 in error queue to allow clients to resume an event where it has failed"""
258 self.__isPartiallyProcessed: bool = False
259 """Store if some data has been processed during current event processing. Will
260 be stored in events in error queue to handle autoremediation properly"""
262 self.__isAnErrorRetry: bool = False
263 """Indicate to handler whether the current event is being processed as part of
264 an error retry"""
266 self.__saveRequired: bool = False
267 """Reset to False at each loop start, and if any change is made during
268 processing, set to True in order to save all cache at the loop end.
269 Used to avoid expensive .save() calls when unnecessary
270 """
272 self.__foreignkeys_events: tuple[str] = self.__FOREIGNKEYS_POLICIES[
273 self.__config["hermes-client"]["foreignkeys_policy"]
274 ]
275 """List of event types that will be placed in the error queue if the object
276 concerning them is the parent (by foreign key) of an object already present in
277 the error queue"""
279 def getObjectFromCache(self, objtype: str, objpkey: Any) -> DataObject:
280 """Returns a deepcopy of an object from cache.
281 Raise IndexError if objtype is invalid, or if objpkey is not found
282 """
283 ds: Datasource = self.__datamodel.localdata
285 _, obj = Datamodel.getObjectFromCacheOrTrashbin(ds, objtype, objpkey)
286 if obj is None:
287 raise IndexError(
288 f"No object of {objtype=} with {objpkey=} was found in cache"
289 )
291 return deepcopy(obj)
293 def getDataobjectlistFromCache(self, objtype: str) -> DataObjectList:
294 """Returns cache of specified objtype, by reference.
295 WARNING: Any modification of the cache content will mess up your client !!!
296 Raise IndexError if objtype is invalid
297 """
298 ds: Datasource = self.__datamodel.localdata
299 cache = ds[objtype]
300 trashbin = ds[f"trashbin_{objtype}"]
302 # Create an empty DataObjectList of same type as cache
303 res = type(cache)(objlist=[])
305 res.extend(cache)
306 res.extend(trashbin)
307 return res
309 def __signalHandler(self, signalnumber: int, frame: FrameType | None):
310 """Signal handler that will be called on SIGINT and SIGTERM"""
311 __hermes__.logger.critical(
312 f"Signal '{signal.strsignal(signalnumber)}' received, terminating"
313 )
314 self.__isStopped = True
316 def __setupSocketParser(self):
317 """Set up the argparse context for unix socket commands"""
318 self.__parser = SocketArgumentParser(
319 prog=f"{self.__config['appname']}-cli",
320 description=f"Hermes client {self.__config['appname']} CLI",
321 exit_on_error=False,
322 )
324 subparsers = self.__parser.add_subparsers(help="Sub-commands")
326 # Quit
327 sp_quit = subparsers.add_parser("quit", help=f"Stop {self.__config['appname']}")
328 sp_quit.set_defaults(func=self.__sock_quit)
330 # Pause
331 sp_pause = subparsers.add_parser(
332 "pause", help="Pause processing until 'resume' command is sent"
333 )
334 sp_pause.set_defaults(func=self.__sock_pause)
336 # Resume
337 sp_resume = subparsers.add_parser(
338 "resume", help="Resume processing that has been paused with 'pause'"
339 )
340 sp_resume.set_defaults(func=self.__sock_resume)
342 # Status
343 sp_status = subparsers.add_parser(
344 "status", help=f"Show {self.__config['appname']} status"
345 )
346 sp_status.set_defaults(func=self.__sock_status)
347 sp_status.add_argument(
348 "-j",
349 "--json",
350 action="store_const",
351 const=True,
352 default=False,
353 help="Print status as json",
354 )
355 sp_status.add_argument(
356 "-v",
357 "--verbose",
358 action="store_const",
359 const=True,
360 default=False,
361 help="Output items without values",
362 )
364 def __processSocketMessage(
365 self, msg: SocketMessageToServer
366 ) -> SocketMessageToClient:
367 """Handler that process specified msg received on unix socket and returns the
368 answer to send"""
369 reply: SocketMessageToClient | None = None
371 try:
372 args = self.__parser.parse_args(msg.argv)
373 if "func" not in args:
374 raise SocketParsingMessage(self.__parser.format_help())
375 except (SocketParsingError, SocketParsingMessage) as e:
376 retmsg = str(e)
377 except argparse.ArgumentError as e:
378 retmsg = self.__parser.format_error(str(e))
379 else:
380 try:
381 reply = args.func(args)
382 except Exception as e:
383 lines = traceback.format_exception(type(e), e, e.__traceback__)
384 trace = "".join(lines).strip()
385 __hermes__.logger.critical(f"Unhandled exception: {trace}")
386 retmsg = trace
388 if reply is None: # Error was met
389 reply = SocketMessageToClient(retcode=1, retmsg=retmsg)
391 return reply
393 def __sock_quit(self, args: argparse.Namespace) -> SocketMessageToClient:
394 """Handler called when quit subcommand is requested on unix socket"""
395 self.__isStopped = True
396 __hermes__.logger.info(f"{self.__config['appname']} has been requested to quit")
397 return SocketMessageToClient(retcode=0, retmsg="")
399 def __sock_pause(self, args: argparse.Namespace) -> SocketMessageToClient:
400 """Handler called when pause subcommand is requested on unix socket"""
401 if self.__isStopped:
402 return SocketMessageToClient(
403 retcode=1,
404 retmsg=f"Error: {self.__config['appname']} is currently being stopped",
405 )
407 if self.__isPaused:
408 return SocketMessageToClient(
409 retcode=1,
410 retmsg=f"Error: {self.__config['appname']} is already paused",
411 )
413 __hermes__.logger.info(
414 f"{self.__config['appname']} has been requested to pause"
415 )
416 self.__isPaused = datetime.now()
417 return SocketMessageToClient(retcode=0, retmsg="")
419 def __sock_resume(self, args: argparse.Namespace) -> SocketMessageToClient:
420 """Handler called when resume subcommand is requested on unix socket"""
421 if self.__isStopped:
422 return SocketMessageToClient(
423 retcode=1,
424 retmsg=f"Error: {self.__config['appname']} is currently being stopped",
425 )
427 if not self.__isPaused:
428 return SocketMessageToClient(
429 retcode=1, retmsg=f"Error: {self.__config['appname']} is not paused"
430 )
432 __hermes__.logger.info(
433 f"{self.__config['appname']} has been requested to resume"
434 )
435 self.__isPaused = None
436 return SocketMessageToClient(retcode=0, retmsg="")
438 def __sock_status(self, args: argparse.Namespace) -> SocketMessageToClient:
439 """Handler called when status subcommand is requested on unix socket"""
440 status = self.__status(verbose=args.verbose)
441 if args.json:
442 msg = json.dumps(status, cls=JSONEncoder, indent=4)
443 else:
444 nl = "\n"
445 info2printable = {}
446 msg = ""
447 for objname in [self.__config["appname"]] + list(
448 status.keys() - (self.__config["appname"],)
449 ):
450 infos = status[objname]
451 msg += f"{objname}:{nl}"
452 for category in ("information", "warning", "error"):
453 if category not in infos:
454 continue
455 if not infos[category]:
456 msg += f" * {category.capitalize()}: []{nl}"
457 continue
459 msg += f" * {category.capitalize()}{nl}"
460 for infoname, infodata in infos[category].items():
461 indentedinfodata = str(infodata).replace("\n", "\n ")
462 msg += (
463 f" - {info2printable.get(infoname, infoname)}:"
464 f" {indentedinfodata}{nl}"
465 )
466 msg = msg.rstrip()
468 return SocketMessageToClient(retcode=0, retmsg=msg)
470 @property
471 def currentStep(self) -> int:
472 """Step number of current event processed.
473 Allow clients to resume an event where it has failed"""
474 return self.__currentStep
476 @currentStep.setter
477 def currentStep(self, value: int):
478 if type(value) is not int:
479 raise TypeError(
480 f"Specified step {value=} has invalid type '{type(value)}'"
481 " instead of int"
482 )
484 if value < 0:
485 raise ValueError(f"Specified step {value=} must be greater or equal to 0")
487 self.__currentStep = value
489 @property
490 def isPartiallyProcessed(self) -> bool:
491 """Indicate if some data has been processed during current event processing.
492 Required to handle autoremediation properly"""
493 return self.__isPartiallyProcessed
495 @isPartiallyProcessed.setter
496 def isPartiallyProcessed(self, value: bool):
497 if type(value) is not bool:
498 raise TypeError(
499 f"Specified isPartiallyProcessed {value=} has invalid type"
500 f" '{type(value)}' instead of bool"
501 )
503 self.__isPartiallyProcessed = value
505 @property
506 def isAnErrorRetry(self) -> bool:
507 """Read-only attribute that indicates to handler whether the current event is
508 being processed as part of an error retry"""
509 return self.__isAnErrorRetry
511 def mainLoop(self):
512 """Client main loop"""
513 self.__startTime = datetime.now()
515 self.__checkDatamodelWarnings()
516 # TODO: implement a check to ensure subclasses required data types and
517 # attributes exist in datamodel
519 if self.__datamodel.hasRemoteSchema():
520 __hermes__.logger.debug(
521 "Remote Dataschema in cache:"
522 f" {self.__datamodel.remote_schema.to_json()=}"
523 )
524 self.__datamodel.loadErrorQueue()
525 else:
526 __hermes__.logger.debug("No remote Dataschema in cache yet")
528 if self.__sock is not None:
529 self.__sock.startProcessMessagesDaemon(appname=__hermes__.appname)
531 # Reduce sleep duration during functional tests to speed them up
532 sleepDuration = 1 if self.__numberOfLoopToProcess is None else 0.05
534 isFirstLoopIteration: bool = True
535 while not self.__isStopped:
536 self.__saveRequired = False
537 currentConfigCanBeSaved = True
539 try:
540 with self.__msgbus:
541 if self.__isPaused or self.__numberOfLoopToProcess == 0:
542 sleep(sleepDuration)
543 continue
545 try:
546 if self.__hasAlreadyBeenInitialized():
547 if isFirstLoopIteration:
548 try:
549 self.__processDatamodelUpdate()
550 except Exception as e:
551 currentConfigCanBeSaved = False
552 self.__notifyFatalException(
553 HermesClientHandlerError.exceptionToString(
554 e, purgeCurrentFileFromTrace=False
555 )
556 )
557 raise HermesAlreadyNotifiedException
558 self.__retryErrorQueue()
559 self.__emptyTrashBin()
560 self.__processEvents(isInitSync=False)
561 else:
562 __hermes__.logger.info(
563 "Client hasn't ran its first initsync sequence yet"
564 )
565 if self.__canBeInitialized():
566 __hermes__.logger.info(
567 "First initsync sequence processing begins"
568 )
569 self.__processEvents(isInitSync=True)
570 if self.__hasAlreadyBeenInitialized():
571 __hermes__.logger.info(
572 "First initsync sequence processing completed"
573 )
574 else:
575 __hermes__.logger.info(
576 "No initsync sequence is available on message bus."
577 " Retry..."
578 )
580 self.__notifyException(None)
582 except HermesAlreadyNotifiedException:
583 pass
584 except InvalidDataError as e:
585 self.__notifyFatalException(
586 HermesClientHandlerError.exceptionToString(
587 e, purgeCurrentFileFromTrace=False
588 )
589 )
590 except Exception as e:
591 self.__notifyException(
592 HermesClientHandlerError.exceptionToString(
593 e, purgeCurrentFileFromTrace=False
594 )
595 )
596 finally:
597 isFirstLoopIteration = False
598 # Still could be True if an exception was raised in
599 # __retryErrorQueue()
600 self.__isAnErrorRetry = False
601 except Exception as e:
602 __hermes__.logger.warning(
603 "Message bus seems to be unavailable."
604 " Waiting 60 seconds before retrying"
605 )
606 self.__notifyException(
607 HermesClientHandlerError.exceptionToString(
608 e, purgeCurrentFileFromTrace=False
609 )
610 )
611 # Wait one second 60 times to avoid waiting too long before stopping
612 for i in range(60):
613 if self.__isStopped:
614 break
615 sleep(1)
616 finally:
617 if self.__saveRequired or self.__isStopped:
618 self.__datamodel.saveErrorQueue()
619 if self.__hasAtLeastBeganInitialization():
620 self.__datamodel.saveLocalAndRemoteData()
621 # Reset current step for "on_save" event
622 self.currentStep = 0
623 self.isPartiallyProcessed = False
624 # Call special event "on_save()"
625 self.__callHandler("", "save")
626 self.__notifyQueueErrors()
627 self.__cache.savecachefile()
629 if self.__isStopped:
630 # Only to ensure cache files version update is saved,
631 # to avoid version migrations at each restart,
632 # as those files aren't expected to be updated often
633 if (
634 currentConfigCanBeSaved
635 and self.__hasAtLeastBeganInitialization()
636 and self.__cache.nextoffset is not None
637 and self.__cache.nextoffset > self.__cache.initstartoffset
638 ):
639 # Do not save until a first event has been properly handled
640 self.__config.savecachefile()
641 self.__datamodel.remote_schema.savecachefile()
643 # Only used in functionnal tests
644 if self.__numberOfLoopToProcess:
645 self.__numberOfLoopToProcess -= 1
647 def __retryErrorQueue(self):
648 # Enforce retryInterval
649 now = datetime.now()
650 if now < self.__errorQueue_lastretry + self.__errorQueue_retryInterval:
651 return # Too early to process again
653 # All events processed in __retryErrorQueue() require this attribute to be True
654 self.__isAnErrorRetry = True
656 done = False
657 evNumbersToRetry: list[int] = []
658 eventNumber: int
659 localEvent: Event
660 remoteEvent: Event | None
661 while not done:
662 retryQueue: list[int] = []
663 previousKeys = self.__datamodel.errorqueue.keys()
665 for (
666 eventNumber,
667 remoteEvent,
668 localEvent,
669 errorMsg,
670 ) in self.__datamodel.errorqueue:
671 if evNumbersToRetry and eventNumber not in evNumbersToRetry:
672 # Ignore eventNumber absent from evNumbersToRetry,
673 # excepted on first iteration of while loop
674 continue
676 if remoteEvent is not None:
677 if self.__datamodel.errorqueue.isEventAParentOfAnotherError(
678 remoteEvent, False
679 ):
680 __hermes__.logger.info(
681 f"Won't retry remote event {remoteEvent} from error queue"
682 " as it is still a dependency of another error"
683 )
684 retryQueue.append(eventNumber)
685 continue
687 __hermes__.logger.info(
688 f"Retrying to process remote event {remoteEvent} from error"
689 " queue"
690 )
691 try:
692 self.__processRemoteEvent(
693 remoteEvent, localEvent, enqueueEventWithError=False
694 )
695 except HermesClientHandlerError as e:
696 __hermes__.logger.info(
697 f"... failed on step {self.currentStep}: {str(e)}"
698 )
699 remoteEvent.step = self.currentStep
700 remoteEvent.isPartiallyProcessed = self.isPartiallyProcessed
701 localEvent.step = self.currentStep
702 localEvent.isPartiallyProcessed = self.isPartiallyProcessed
703 self.__datamodel.errorqueue.updateErrorMsg(eventNumber, e.msg)
704 else:
705 # If event has suppressed object, eventNumber has already been
706 # purged from queue
707 self.__datamodel.errorqueue.remove(
708 eventNumber, ignoreMissingEventNumber=True
709 )
710 else:
711 if self.__datamodel.errorqueue.isEventAParentOfAnotherError(
712 localEvent, True
713 ):
714 __hermes__.logger.info(
715 f"Won't retry local event {localEvent} from error queue"
716 " as it is still a dependency of another error"
717 )
718 retryQueue.append(eventNumber)
719 continue
720 __hermes__.logger.info(
721 f"Retrying to process local event {localEvent} from error queue"
722 )
723 try:
724 self.__processLocalEvent(
725 remoteEvent, localEvent, enqueueEventWithError=False
726 )
727 except HermesClientHandlerError as e:
728 __hermes__.logger.info(
729 f"... failed on step {self.currentStep}: {str(e)}"
730 )
731 localEvent.step = self.currentStep
732 localEvent.isPartiallyProcessed = self.isPartiallyProcessed
733 self.__datamodel.errorqueue.updateErrorMsg(eventNumber, e.msg)
734 else:
735 # If event has suppressed object, eventNumber has already been
736 # purged from queue
737 self.__datamodel.errorqueue.remove(
738 eventNumber, ignoreMissingEventNumber=True
739 )
740 while self.__isPaused and not self.__isStopped:
741 sleep(1) # Allow loop to be paused
742 if self.__isStopped:
743 break # Allow loop to be interrupted if requested
744 else:
745 # Update __errorQueue_lastretry only if loop hasn't been interrupted
746 self.__errorQueue_lastretry = now
748 done = previousKeys == self.__datamodel.errorqueue.keys() or not retryQueue
749 if done:
750 if previousKeys:
751 __hermes__.logger.debug(
752 f"End of retryerrorqueue {previousKeys=}"
753 f" {self.__datamodel.errorqueue.keys()=} - {retryQueue=}"
754 )
755 else:
756 __hermes__.logger.debug(
757 "As some event have been processed, will retry ignored events"
758 f" {retryQueue}"
759 )
760 evNumbersToRetry = retryQueue.copy()
762 self.__isAnErrorRetry = False # End of __retryErrorQueue()
764 def __emptyTrashBin(self, force: bool = False):
765 # Enforce purgeInterval
766 now = datetime.now()
767 if (
768 not force
769 and now < self.__trashbin_lastpurge + self.__trashbin_purgeInterval
770 ):
771 return # Too early to process again
772 if self.__trashbin_retention is not None:
773 retentionLimit = datetime.now() - self.__trashbin_retention
775 objtype: str
776 objs: DataObjectList
777 # As we'll remove objects, process data in the datamodel reversed declaration
778 # order
779 for objtype, objs in reversed(self.__datamodel.remotedata.items()):
780 if not objtype.startswith("trashbin_"):
781 continue
783 for pkey in objs.getPKeys():
784 obj = objs.get(pkey)
785 if (
786 self.__trashbin_retention is None
787 or obj._trashbin_timestamp < retentionLimit
788 ):
789 event = Event(
790 evcategory="base", eventtype="removed", obj=obj, objattrs={}
791 )
792 __hermes__.logger.info(f"Trying to purge {repr(obj)} from trashbin")
793 if self.__datamodel.errorqueue.containsObjectByEvent(
794 event, isLocalEvent=False
795 ) and not self.__datamodel.errorqueue.isEventAParentOfAnotherError(
796 event, isLocalEvent=False
797 ):
798 try:
799 self.__processRemoteEvent(
800 event, local_event=None, enqueueEventWithError=False
801 )
802 except HermesClientHandlerError:
803 pass
804 else:
805 self.__processRemoteEvent(
806 event, local_event=None, enqueueEventWithError=True
807 )
809 while self.__isPaused and not self.__isStopped:
810 sleep(1) # Allow loop to be paused
811 if self.__isStopped:
812 break # Allow loop to be interrupted if requested
814 while self.__isPaused and not self.__isStopped:
815 sleep(1) # Allow loop to be paused
816 if self.__isStopped:
817 break # Allow loop to be interrupted if requested
818 else:
819 # Update __trashbin_lastpurge only if loop hasn't been interrupted
820 self.__trashbin_lastpurge = now
822 def __hasAlreadyBeenInitialized(self) -> bool:
823 if (
824 self.__cache.initstartoffset is None
825 or self.__cache.initstopoffset is None
826 or self.__cache.nextoffset is None
827 or self.__cache.nextoffset < self.__cache.initstopoffset
828 ):
829 return False
830 return True
832 def __hasAtLeastBeganInitialization(self) -> bool:
833 return (
834 self.__cache.initstartoffset is not None
835 and self.__datamodel.hasRemoteSchema()
836 )
838 def __canBeInitialized(self) -> bool:
839 self.__msgbus.seekToBeginning()
841 # List of (start, stop) offsets of initsync sequences found
842 initSyncFound: list[tuple[Any, Any]] = []
844 start = None
845 stop = None
846 event: Event
847 for event in self.__msgbus:
848 if event.evcategory != "initsync":
849 continue
850 if event.eventtype == "init-start":
851 start = event.offset
852 elif event.eventtype == "init-stop" and start is not None:
853 stop = event.offset
854 initSyncFound.append(
855 (start, stop),
856 )
857 if self.__useFirstInitsyncSequence:
858 break # We found the first sequence
859 else:
860 # Continue to find a new complete sequence
861 start = None
862 stop = None
864 if not initSyncFound:
865 return False
867 if self.__useFirstInitsyncSequence:
868 start, stop = initSyncFound[0]
869 else:
870 start, stop = initSyncFound[-1]
872 if self.__cache.nextoffset is None or self.__cache.nextoffset < start:
873 self.__cache.nextoffset = start
875 self.__cache.initstartoffset = start
876 self.__cache.initstopoffset = stop
877 __hermes__.logger.debug(
878 "Init sequence was found in Kafka at offsets"
879 f" [{self.__cache.initstartoffset} ; {self.__cache.initstopoffset}]"
880 )
881 return True
883 def __updateSchema(self, newSchema: Dataschema):
884 if self.__datamodel.forcePurgeOfTrashedObjectsWithoutNewPkeys(
885 self.__datamodel.remote_schema, newSchema
886 ):
887 self.__emptyTrashBin(force=True)
889 self.__datamodel.updateSchema(newSchema)
890 self.__config.savecachefile() # Save config to be able to rebuild datamodel
891 # Save and reload error queue to purge it from events of any suppressed types
892 self.__datamodel.saveErrorQueue()
893 self.__datamodel.loadErrorQueue()
894 self.__checkDatamodelWarnings()
896 def __checkDatamodelWarnings(self):
897 if self.__datamodel.unknownRemoteTypes:
898 __hermes__.logger.warning(
899 "Datamodel errors: remote types"
900 f" '{self.__datamodel.unknownRemoteTypes}'"
901 " don't exist in current Dataschema"
902 )
904 if self.__datamodel.unknownRemoteAttributes:
905 __hermes__.logger.warning(
906 "Datamodel errors: remote attributes don't exist in current"
907 f" Dataschema: {self.__datamodel.unknownRemoteAttributes}"
908 )
909 self.__notifyDatamodelWarnings()
911 def __processEvents(self, isInitSync=False):
912 remote_event: Event
913 schema: Dataschema | None = None
914 evcategory: str = "initsync" if isInitSync else "base"
916 self.__msgbus.seek(self.__cache.nextoffset)
918 for remote_event in self.__msgbus:
919 self.__saveRequired = True
920 if isInitSync and remote_event.offset > self.__cache.initstopoffset:
921 # Should never be called
922 self.__cache.nextoffset = remote_event.offset + 1
923 break
925 # TODO: implement data consistency check if event.evcategory==initsync
926 # and evcategory==base
928 if remote_event.evcategory != evcategory:
929 self.__cache.nextoffset = remote_event.offset + 1
930 continue
932 if isInitSync:
933 if remote_event.eventtype == "init-start":
934 schema = Dataschema.from_json(remote_event.objattrs)
935 self.__updateSchema(schema)
936 continue
938 if remote_event.eventtype == "init-stop":
939 self.__cache.nextoffset = remote_event.offset + 1
940 break
942 if schema is None and not self.__hasAtLeastBeganInitialization():
943 msg = "Invalid initsync sequence met, ignoring"
944 __hermes__.logger.critical(msg)
945 return
947 # Process "standard" message
948 match remote_event.eventtype:
949 case "added" | "modified" | "removed":
950 self.__processRemoteEvent(
951 remote_event, local_event=None, enqueueEventWithError=True
952 )
953 case "dataschema":
954 schema = Dataschema.from_json(remote_event.objattrs)
955 self.__updateSchema(schema)
956 case _:
957 __hermes__.logger.error(
958 "Received an event with unknown type"
959 f" '{remote_event.eventtype}': ignored"
960 )
962 self.__cache.nextoffset = remote_event.offset + 1
964 while self.__isPaused and not self.__isStopped:
965 sleep(1) # Allow loop to be paused
966 if self.__isStopped:
967 break # Allow loop to be interrupted if requested
969 def __processRemoteEvent(
970 self,
971 remote_event: Event | None,
972 local_event: Event | None,
973 enqueueEventWithError: bool,
974 simulateOnly: bool = False,
975 ):
976 secretAttrs = self.__datamodel.remote_schema.secretsAttributesOf(
977 remote_event.objtype
978 )
979 __hermes__.logger.debug(
980 f"__processRemoteEvent({remote_event.toString(secretAttrs)})"
981 )
982 self.__saveRequired = True
984 # In case of modification, try to compose full object in order to provide all
985 # attributes values, in order to render Jinja Template with several vars.
986 # In this specific case, one of the template var may have been modified,
987 # but not the other, so the event attrs are not enough to process the
988 # template rendering
989 r_obj_complete: DataObject | None = None
990 if remote_event.eventtype == "modified":
991 cache_complete, r_cachedobj_complete = (
992 Datamodel.getObjectFromCacheOrTrashbin(
993 self.__datamodel.remotedata_complete,
994 remote_event.objtype,
995 remote_event.objpkey,
996 )
997 )
998 if r_cachedobj_complete is not None:
999 r_obj_complete = Datamodel.getUpdatedObject(
1000 r_cachedobj_complete, remote_event.objattrs
1001 )
1003 # Should be always None, except when called from __retryErrorQueue()
1004 # In this case, we have to use the provided local_event, as it may contains
1005 # some extra changes stacked by autoremediation
1006 if local_event is None:
1007 local_events = self.__datamodel.convertEventToLocal(
1008 remote_event, r_obj_complete
1009 )
1010 else:
1011 local_events = [local_event]
1013 trashbin = self.__datamodel.remotedata[f"trashbin_{remote_event.objtype}"]
1014 was_in_trashbin = remote_event.objpkey in trashbin
1016 if not simulateOnly and enqueueEventWithError:
1017 hadErrors = self.__datamodel.errorqueue.containsObjectByEvent(
1018 remote_event, isLocalEvent=False
1019 )
1020 isParent = self.__datamodel.errorqueue.isEventAParentOfAnotherError(
1021 remote_event, isLocalEvent=False
1022 )
1024 for local_event in local_events:
1025 if not simulateOnly and enqueueEventWithError:
1026 if hadErrors or (
1027 isParent and remote_event.eventtype in self.__foreignkeys_events
1028 ):
1029 secretAttrs = self.__datamodel.remote_schema.secretsAttributesOf(
1030 remote_event.objtype
1031 )
1032 if hadErrors:
1033 errorMsg = (
1034 "Object in remote event"
1035 f" {remote_event.toString(secretAttrs)} already had"
1036 " unresolved errors: appending event to error queue"
1037 )
1038 else:
1039 errorMsg = (
1040 "Object in remote event"
1041 f" {remote_event.toString(secretAttrs)} is a dependency of"
1042 " an object that already had unresolved errors: appending"
1043 " event to error queue"
1044 )
1045 __hermes__.logger.warning(errorMsg)
1046 self.__processRemoteEvent(
1047 remote_event,
1048 local_event=None,
1049 enqueueEventWithError=False,
1050 simulateOnly=True,
1051 )
1052 if local_event is None:
1053 # Force empty event generation when local_event doesn't change
1054 # anything
1055 l_events = self.__datamodel.convertEventToLocal(
1056 remote_event, r_obj_complete, allowEmptyEvent=True
1057 )
1058 for local_event in l_events:
1059 self.__datamodel.errorqueue.append(
1060 remote_event, local_event, errorMsg
1061 )
1062 else:
1063 self.__datamodel.errorqueue.append(
1064 remote_event, local_event, errorMsg
1065 )
1066 continue
1068 try:
1069 match remote_event.eventtype:
1070 case "added":
1071 if self.__trashbin_retention is not None and was_in_trashbin:
1072 # Object is in trashbin, recycle it
1073 self.__remoteRecycled(
1074 remote_event, local_event, simulateOnly
1075 )
1076 else:
1077 # Add new object
1078 self.__remoteAdded(remote_event, local_event, simulateOnly)
1080 case "modified":
1081 self.__remoteModified(remote_event, local_event, simulateOnly)
1083 case "removed":
1084 # Remove object on any of these conditions:
1085 # - trashbin retention is disabled
1086 # - object is already in trashbin
1087 if self.__trashbin_retention is None or was_in_trashbin:
1088 # Remove object
1089 self.__remoteRemoved(
1090 remote_event, local_event, simulateOnly
1091 )
1092 else:
1093 # Store object in trashbin
1094 self.__remoteTrashed(
1095 remote_event, local_event, simulateOnly
1096 )
1097 except HermesClientHandlerError as e:
1098 if not simulateOnly and enqueueEventWithError:
1099 self.__processRemoteEvent(
1100 remote_event,
1101 local_event=None,
1102 enqueueEventWithError=False,
1103 simulateOnly=True,
1104 )
1105 remote_event.step = self.currentStep
1106 remote_event.isPartiallyProcessed = self.isPartiallyProcessed
1107 local_event.step = self.currentStep
1108 local_event.isPartiallyProcessed = self.isPartiallyProcessed
1110 if local_event is None:
1111 # Force empty event generation when local_event doesn't change
1112 # anything
1113 l_events = self.__datamodel.convertEventToLocal(
1114 remote_event, r_obj_complete, allowEmptyEvent=True
1115 )
1116 for local_event in l_events:
1117 self.__datamodel.errorqueue.append(
1118 remote_event, local_event, e.msg
1119 )
1120 else:
1121 self.__datamodel.errorqueue.append(
1122 remote_event, local_event, e.msg
1123 )
1124 else:
1125 raise
1127 def __processLocalEvent(
1128 self,
1129 remote_event: Event | None,
1130 local_event: Event | None,
1131 enqueueEventWithError: bool,
1132 simulateOnly: bool = False,
1133 dontStoreInTrashbin: bool = False,
1134 ):
1135 if local_event is None:
1136 __hermes__.logger.debug("__processLocalEvent(None)")
1137 return
1139 secretAttrs = self.__datamodel.local_schema.secretsAttributesOf(
1140 local_event.objtype
1141 )
1142 __hermes__.logger.debug(
1143 f"__processLocalEvent({local_event.toString(secretAttrs)})"
1144 )
1146 self.__saveRequired = True
1148 if not simulateOnly:
1149 # Reset current step
1150 self.currentStep = local_event.step
1151 self.isPartiallyProcessed = local_event.isPartiallyProcessed
1153 if not simulateOnly and enqueueEventWithError:
1154 hadErrors = self.__datamodel.errorqueue.containsObjectByEvent(
1155 local_event, isLocalEvent=True
1156 )
1157 isParent = self.__datamodel.errorqueue.isEventAParentOfAnotherError(
1158 local_event, isLocalEvent=True
1159 )
1160 if hadErrors or (
1161 isParent and local_event.eventtype in self.__foreignkeys_events
1162 ):
1163 secretAttrs = self.__datamodel.local_schema.secretsAttributesOf(
1164 local_event.objtype
1165 )
1166 if hadErrors:
1167 errorMsg = (
1168 f"Object in local event {local_event.toString(secretAttrs)}"
1169 " already had unresolved errors: appending event to error queue"
1170 )
1171 else:
1172 errorMsg = (
1173 f"Object in local event {local_event.toString(secretAttrs)}"
1174 " is a dependency of an object that already had unresolved"
1175 " errors: appending event to error queue"
1176 )
1177 __hermes__.logger.warning(errorMsg)
1178 self.__processLocalEvent(
1179 None, local_event, enqueueEventWithError=False, simulateOnly=True
1180 )
1181 self.__datamodel.errorqueue.append(remote_event, local_event, errorMsg)
1182 return
1184 trashbin = self.__datamodel.localdata[f"trashbin_{local_event.objtype}"]
1185 try:
1186 match local_event.eventtype:
1187 case "added":
1188 if (
1189 self.__trashbin_retention is not None
1190 and local_event.objpkey in trashbin
1191 ):
1192 # Object is in trashbin, recycle it
1193 self.__localRecycled(local_event, simulateOnly)
1194 else:
1195 self.__localAdded(local_event, simulateOnly)
1196 case "modified":
1197 if not simulateOnly and local_event.objpkey in trashbin:
1198 # Object is in trashbin, and cannot be modified until it is
1199 # restored
1200 if enqueueEventWithError:
1201 # As the object changes will be processed at restore,
1202 # ignore the change
1203 self.__processLocalEvent(
1204 None,
1205 local_event,
1206 enqueueEventWithError=False,
1207 simulateOnly=True,
1208 )
1209 else:
1210 # Propagate error as requested
1211 raise HermesClientHandlerError(
1212 f"Object of event {repr(local_event)} is in trashbin,"
1213 " and cannot be modified until it is restored"
1214 )
1215 else:
1216 self.__localModified(local_event, simulateOnly)
1217 case "removed":
1218 # Remove object on any of these conditions:
1219 # - trashbin retention is disabled
1220 # - object is already in trashbin
1221 # - dontStoreInTrashbin arg is True (used on datamodel update)
1222 if (
1223 self.__trashbin_retention is None
1224 or local_event.objpkey in trashbin
1225 or dontStoreInTrashbin
1226 ):
1227 self.__localRemoved(local_event, simulateOnly)
1228 else:
1229 self.__localTrashed(local_event, simulateOnly)
1230 except HermesClientHandlerError as e:
1231 if not simulateOnly and enqueueEventWithError:
1232 self.__processLocalEvent(
1233 None, local_event, enqueueEventWithError=False, simulateOnly=True
1234 )
1235 if remote_event is not None:
1236 remote_event.step = self.currentStep
1237 remote_event.isPartiallyProcessed = self.isPartiallyProcessed
1238 local_event.step = self.currentStep
1239 local_event.isPartiallyProcessed = self.isPartiallyProcessed
1240 self.__datamodel.errorqueue.append(remote_event, local_event, e.msg)
1241 else:
1242 raise
1244 def __remoteAdded(
1245 self, remote_event: Event, local_event: Event, simulateOnly: bool = False
1246 ):
1247 secretAttrs = self.__datamodel.remote_schema.secretsAttributesOf(
1248 remote_event.objtype
1249 )
1250 __hermes__.logger.debug(f"__remoteAdded({remote_event.toString(secretAttrs)})")
1252 r_obj = self.__datamodel.createRemoteDataobject(
1253 remote_event.objtype, remote_event.objattrs
1254 )
1255 self.__processLocalEvent(
1256 remote_event,
1257 local_event,
1258 enqueueEventWithError=False,
1259 simulateOnly=simulateOnly,
1260 )
1262 # Add remote object to cache
1263 if not simulateOnly:
1264 # May already been added if remote type is used in more than one local type
1265 self.__datamodel.remotedata[remote_event.objtype].append(
1266 r_obj, ignoreIfAlreadyPresent=True
1267 )
1268 # May already been added if current event is from errorqueue
1269 self.__datamodel.remotedata_complete[remote_event.objtype].append(
1270 r_obj, ignoreIfAlreadyPresent=True
1271 )
1273 def __localAdded(self, local_ev: Event, simulateOnly: bool = False):
1274 secretAttrs = self.__datamodel.local_schema.secretsAttributesOf(
1275 local_ev.objtype
1276 )
1277 __hermes__.logger.debug(f"__localAdded({local_ev.toString(secretAttrs)})")
1279 l_obj = self.__datamodel.createLocalDataobject(
1280 local_ev.objtype, local_ev.objattrs
1281 )
1283 if not simulateOnly:
1284 # Call added handler
1285 self.__callHandler(
1286 objtype=local_ev.objtype,
1287 eventtype="added",
1288 objkey=local_ev.objpkey,
1289 eventattrs=local_ev.objattrs,
1290 newobj=deepcopy(l_obj),
1291 )
1293 # Add local object to cache
1294 if not simulateOnly:
1295 self.__datamodel.localdata[local_ev.objtype].append(
1296 l_obj, ignoreIfAlreadyPresent=True
1297 )
1298 # May already been added if current event is from errorqueue
1299 self.__datamodel.localdata_complete[local_ev.objtype].append(
1300 l_obj, ignoreIfAlreadyPresent=True
1301 )
1303 def __remoteRecycled(
1304 self, remote_event: Event, local_event: Event, simulateOnly: bool = False
1305 ):
1306 secretAttrs = self.__datamodel.remote_schema.secretsAttributesOf(
1307 remote_event.objtype
1308 )
1309 __hermes__.logger.debug(
1310 f"__remoteRecycled({remote_event.toString(secretAttrs)})"
1311 )
1312 maincache = self.__datamodel.remotedata[remote_event.objtype]
1313 maincache_complete = self.__datamodel.remotedata_complete[remote_event.objtype]
1314 trashbin = self.__datamodel.remotedata[f"trashbin_{remote_event.objtype}"]
1315 trashbin_complete = self.__datamodel.remotedata_complete[
1316 f"trashbin_{remote_event.objtype}"
1317 ]
1319 r_obj = self.__datamodel.createRemoteDataobject(
1320 remote_event.objtype, remote_event.objattrs
1321 )
1323 self.__processLocalEvent(
1324 remote_event,
1325 local_event,
1326 enqueueEventWithError=False,
1327 simulateOnly=simulateOnly,
1328 )
1330 # Remove remote object from trashbin
1331 if not simulateOnly:
1332 trashbin.removeByPkey(remote_event.objpkey)
1333 trashbin_complete.removeByPkey(remote_event.objpkey)
1334 # Restore remote object, with its potential changes, in main cache
1335 if not simulateOnly:
1336 # May already been added if remote type is used in more than one local type
1337 maincache.append(r_obj, ignoreIfAlreadyPresent=True)
1338 # May already been recycled if current event is from errorqueue
1339 maincache_complete.append(r_obj, ignoreIfAlreadyPresent=True)
1341 def __localRecycled(self, local_ev: Event, simulateOnly: bool = False):
1342 secretAttrs = self.__datamodel.local_schema.secretsAttributesOf(
1343 local_ev.objtype
1344 )
1345 __hermes__.logger.debug(f"__localRecycled({local_ev.toString(secretAttrs)})")
1346 maincache = self.__datamodel.localdata[local_ev.objtype]
1347 maincache_complete = self.__datamodel.localdata_complete[local_ev.objtype]
1348 trashbin = self.__datamodel.localdata[f"trashbin_{local_ev.objtype}"]
1349 trashbin_complete = self.__datamodel.localdata_complete[
1350 f"trashbin_{local_ev.objtype}"
1351 ]
1353 l_obj = self.__datamodel.createLocalDataobject(
1354 local_ev.objtype, local_ev.objattrs
1355 )
1356 l_obj_trash: DataObject = deepcopy(trashbin.get(local_ev.objpkey))
1357 del l_obj_trash._trashbin_timestamp # Remove trashbin timestamp from object
1359 l_obj_trash_complete: DataObject = deepcopy(
1360 trashbin_complete.get(local_ev.objpkey)
1361 )
1363 # May already been recycled if current event is from errorqueue
1364 if l_obj_trash_complete is not None:
1365 del (
1366 l_obj_trash_complete._trashbin_timestamp
1367 ) # Remove trashbin timestamp from object
1369 if not simulateOnly:
1370 # Call recycled handler
1371 self.__callHandler(
1372 objtype=local_ev.objtype,
1373 eventtype="recycled",
1374 objkey=l_obj_trash.getPKey(),
1375 eventattrs=l_obj_trash.toNative(),
1376 newobj=deepcopy(l_obj_trash),
1377 )
1379 if not simulateOnly:
1380 # Remove local object from trashbin
1381 trashbin.remove(l_obj_trash)
1382 # Restore local object in main cache
1383 maincache.append(l_obj_trash, ignoreIfAlreadyPresent=True)
1385 # May already been recycled if current event is from errorqueue
1386 if l_obj_trash_complete is not None:
1387 trashbin_complete.remove(l_obj_trash_complete)
1388 maincache_complete.append(l_obj_trash_complete, ignoreIfAlreadyPresent=True)
1390 diff = l_obj.diffFrom(l_obj_trash) # Handle local object changes if any
1391 if diff and not simulateOnly:
1392 event, obj = Event.fromDiffItem(
1393 diffitem=diff,
1394 eventCategory=local_ev.evcategory,
1395 changeType="modified",
1396 )
1397 # Hack: we pass this second event (modified) to error queue in order to
1398 # postpone its processing once all caches of previous one are up to date.
1399 # Otherwise, if an error is met on this second event (modified), we'll try
1400 # to reprocess the first one (recycled)
1401 self.__datamodel.errorqueue.append(
1402 remoteEvent=None, localEvent=event, errorMsg=None
1403 )
1404 # ... and force error queue to be retried asap in order to process the
1405 # pending event
1406 self.__errorQueue_lastretry = datetime(year=1, month=1, day=1)
1408 def __remoteModified(
1409 self, remote_event: Event, local_event: Event, simulateOnly: bool = False
1410 ):
1411 secretAttrs = self.__datamodel.remote_schema.secretsAttributesOf(
1412 remote_event.objtype
1413 )
1414 __hermes__.logger.debug(
1415 f"__remoteModified({remote_event.toString(secretAttrs)})"
1416 )
1417 maincache = self.__datamodel.remotedata[remote_event.objtype]
1419 if not simulateOnly:
1420 r_cachedobj: DataObject = maincache.get(remote_event.objpkey)
1421 r_obj = Datamodel.getUpdatedObject(r_cachedobj, remote_event.objattrs)
1423 cache_complete, r_cachedobj_complete = Datamodel.getObjectFromCacheOrTrashbin(
1424 self.__datamodel.remotedata_complete,
1425 remote_event.objtype,
1426 remote_event.objpkey,
1427 )
1428 r_obj_complete = Datamodel.getUpdatedObject(
1429 r_cachedobj_complete, remote_event.objattrs
1430 )
1432 self.__processLocalEvent(
1433 remote_event,
1434 local_event,
1435 enqueueEventWithError=False,
1436 simulateOnly=simulateOnly,
1437 )
1439 # Update remote object in cache
1440 if not simulateOnly:
1441 maincache.replace(r_obj)
1443 # May not exist
1444 if cache_complete is not None:
1445 cache_complete.replace(r_obj_complete)
1447 def __localModified(self, local_ev: Event, simulateOnly: bool = False):
1448 secretAttrs = self.__datamodel.local_schema.secretsAttributesOf(
1449 local_ev.objtype
1450 )
1451 __hermes__.logger.debug(f"__localModified({local_ev.toString(secretAttrs)})")
1452 maincache = self.__datamodel.localdata[local_ev.objtype]
1454 if not simulateOnly:
1455 l_cachedobj: DataObject = maincache.get(local_ev.objpkey)
1456 l_obj = Datamodel.getUpdatedObject(l_cachedobj, local_ev.objattrs)
1458 cache_complete, l_cachedobj_complete = Datamodel.getObjectFromCacheOrTrashbin(
1459 self.__datamodel.localdata_complete, local_ev.objtype, local_ev.objpkey
1460 )
1462 # May not exist
1463 if cache_complete is not None:
1464 l_obj_complete = Datamodel.getUpdatedObject(
1465 l_cachedobj_complete, local_ev.objattrs
1466 )
1468 if not simulateOnly:
1469 # Call modified handler
1470 self.__callHandler(
1471 objtype=local_ev.objtype,
1472 eventtype="modified",
1473 objkey=local_ev.objpkey,
1474 eventattrs=local_ev.objattrs,
1475 newobj=deepcopy(l_obj),
1476 cachedobj=deepcopy(l_cachedobj),
1477 )
1479 # Update local object in cache
1480 if not simulateOnly:
1481 maincache.replace(l_obj)
1483 # May not exist
1484 if cache_complete is not None:
1485 cache_complete.replace(l_obj_complete)
1487 def __remoteTrashed(
1488 self, remote_event: Event, local_event: Event, simulateOnly: bool = False
1489 ):
1490 secretAttrs = self.__datamodel.remote_schema.secretsAttributesOf(
1491 remote_event.objtype
1492 )
1493 __hermes__.logger.debug(
1494 f"__remoteTrashed({remote_event.toString(secretAttrs)})"
1495 )
1496 maincache = self.__datamodel.remotedata[remote_event.objtype]
1497 maincache_complete = self.__datamodel.remotedata_complete[remote_event.objtype]
1498 trashbin = self.__datamodel.remotedata[f"trashbin_{remote_event.objtype}"]
1499 trashbin_complete = self.__datamodel.remotedata_complete[
1500 f"trashbin_{remote_event.objtype}"
1501 ]
1503 r_cachedobj: DataObject = maincache.get(remote_event.objpkey)
1504 r_cachedobj_complete: DataObject = maincache_complete.get(remote_event.objpkey)
1506 self.__processLocalEvent(
1507 remote_event,
1508 local_event,
1509 enqueueEventWithError=False,
1510 simulateOnly=simulateOnly,
1511 )
1513 if not simulateOnly:
1514 # Remove remote object from cache
1515 # May already been removed if remote type is used in more than one local
1516 # type
1517 if r_cachedobj is not None:
1518 maincache.remove(r_cachedobj)
1519 r_cachedobj._trashbin_timestamp = remote_event.timestamp
1520 # Add remote object to trashbin
1521 # May already been added if remote type is used in more than one local
1522 # type
1523 trashbin.append(r_cachedobj, ignoreIfAlreadyPresent=True)
1525 if r_cachedobj_complete is not None:
1526 maincache_complete.remove(r_cachedobj_complete)
1527 r_cachedobj_complete._trashbin_timestamp = remote_event.timestamp
1528 # May already been added if remote type is used in more than one local type
1529 trashbin_complete.append(r_cachedobj_complete, ignoreIfAlreadyPresent=True)
1531 def __localTrashed(self, local_ev: Event, simulateOnly: bool = False):
1532 secretAttrs = self.__datamodel.local_schema.secretsAttributesOf(
1533 local_ev.objtype
1534 )
1535 __hermes__.logger.debug(f"__localTrashed({local_ev.toString(secretAttrs)})")
1536 maincache = self.__datamodel.localdata[local_ev.objtype]
1537 maincache_complete = self.__datamodel.localdata_complete[local_ev.objtype]
1538 trashbin = self.__datamodel.localdata[f"trashbin_{local_ev.objtype}"]
1539 trashbin_complete = self.__datamodel.localdata_complete[
1540 f"trashbin_{local_ev.objtype}"
1541 ]
1543 l_cachedobj: DataObject = maincache.get(local_ev.objpkey)
1544 l_cachedobj_complete: DataObject = maincache_complete.get(local_ev.objpkey)
1546 if not simulateOnly:
1547 # Call trashed handler
1548 self.__callHandler(
1549 objtype=local_ev.objtype,
1550 eventtype="trashed",
1551 objkey=local_ev.objpkey,
1552 eventattrs=local_ev.objattrs,
1553 cachedobj=deepcopy(l_cachedobj),
1554 )
1556 if not simulateOnly:
1557 # Remove local object from cache
1558 maincache.remove(l_cachedobj)
1559 l_cachedobj._trashbin_timestamp = local_ev.timestamp
1560 # Add local object to trashbin
1561 trashbin.append(l_cachedobj, ignoreIfAlreadyPresent=True)
1563 if l_cachedobj_complete is not None:
1564 maincache_complete.remove(l_cachedobj_complete)
1565 l_cachedobj_complete._trashbin_timestamp = local_ev.timestamp
1566 trashbin_complete.append(l_cachedobj_complete, ignoreIfAlreadyPresent=True)
1568 def __remoteRemoved(
1569 self, remote_event: Event, local_event: Event, simulateOnly: bool = False
1570 ):
1571 secretAttrs = self.__datamodel.remote_schema.secretsAttributesOf(
1572 remote_event.objtype
1573 )
1574 __hermes__.logger.debug(
1575 f"__remoteRemoved({remote_event.toString(secretAttrs)})"
1576 )
1578 cache, r_cachedobj = Datamodel.getObjectFromCacheOrTrashbin(
1579 self.__datamodel.remotedata,
1580 remote_event.objtype,
1581 remote_event.objpkey,
1582 )
1584 cache_complete, r_cachedobj_complete = Datamodel.getObjectFromCacheOrTrashbin(
1585 self.__datamodel.remotedata_complete,
1586 remote_event.objtype,
1587 remote_event.objpkey,
1588 )
1590 self.__processLocalEvent(
1591 remote_event,
1592 local_event,
1593 enqueueEventWithError=False,
1594 simulateOnly=simulateOnly,
1595 )
1597 # Remove remote object from cache or trashbin
1598 # May already been removed if remote type is used in more than one local type
1599 if not simulateOnly:
1600 if r_cachedobj is not None:
1601 cache.remove(r_cachedobj)
1603 # May already been removed if current event is from errorqueue
1604 if cache_complete is not None:
1605 # May already been removed if remote type is used in more than one local
1606 # type
1607 if r_cachedobj_complete is not None:
1608 cache_complete.remove(r_cachedobj_complete)
1610 def __localRemoved(self, local_ev: Event, simulateOnly: bool = False):
1611 secretAttrs = self.__datamodel.local_schema.secretsAttributesOf(
1612 local_ev.objtype
1613 )
1614 __hermes__.logger.debug(f"__localRemoved({local_ev.toString(secretAttrs)})")
1616 cache, l_cachedobj = Datamodel.getObjectFromCacheOrTrashbin(
1617 self.__datamodel.localdata,
1618 local_ev.objtype,
1619 local_ev.objpkey,
1620 )
1622 cache_complete, l_cachedobj_complete = Datamodel.getObjectFromCacheOrTrashbin(
1623 self.__datamodel.localdata_complete,
1624 local_ev.objtype,
1625 local_ev.objpkey,
1626 )
1628 if not simulateOnly:
1629 # Call removed handler
1630 self.__callHandler(
1631 objtype=local_ev.objtype,
1632 eventtype="removed",
1633 objkey=local_ev.objpkey,
1634 eventattrs=local_ev.objattrs,
1635 cachedobj=deepcopy(l_cachedobj),
1636 )
1638 # Remove local object from cache or trashbin
1639 if not simulateOnly:
1640 cache.remove(l_cachedobj)
1641 # May already been removed if current event is from errorqueue
1642 if cache_complete is not None:
1643 cache_complete.remove(l_cachedobj_complete)
1645 def __callHandler(self, objtype: str, eventtype: str, **kwargs):
1646 if not objtype:
1647 handlerName = f"on_{eventtype}"
1648 kwargs_filtered = kwargs
1649 else:
1650 handlerName = f"on_{objtype}_{eventtype}"
1652 # Filter secrets values
1653 secretAttrs = self.__datamodel.local_schema.secretsAttributesOf(objtype)
1654 kwargs_filtered = kwargs.copy()
1655 kwargs_filtered["eventattrs"] = Event.objattrsToString(
1656 kwargs["eventattrs"], secretAttrs
1657 )
1659 kwargsstr = ", ".join([f"{k}={repr(v)}" for k, v in kwargs_filtered.items()])
1661 hdlr = getattr(self, handlerName, None)
1663 if not callable(hdlr):
1664 __hermes__.logger.debug(
1665 f"Calling '{handlerName}({kwargsstr})': handler '{handlerName}()'"
1666 " doesn't exists"
1667 )
1668 return
1670 __hermes__.logger.info(
1671 f"Calling '{handlerName}({kwargsstr})' - currentStep={self.currentStep},"
1672 f" isPartiallyProcessed={self.isPartiallyProcessed},"
1673 f" isAnErrorRetry={self.isAnErrorRetry}"
1674 )
1676 try:
1677 hdlr(**kwargs)
1678 except Exception as e:
1679 __hermes__.logger.error(
1680 f"Calling '{handlerName}({kwargsstr})': error met on step"
1681 f" {self.currentStep} '{str(e)}'"
1682 )
1683 raise HermesClientHandlerError(e)
1685 def __processDatamodelUpdate(self):
1686 """Check difference between current datamodel and previous one. If datamodel has
1687 changed, generate local events according to datamodel changes"""
1689 diff = self.__newdatamodel.diffFrom(self.__datamodel)
1691 if not diff:
1692 __hermes__.logger.info("No change in datamodel")
1693 # Start working with new datamodel
1694 self.__datamodel = self.__newdatamodel
1695 self.__datamodel.loadErrorQueue()
1696 return
1698 __hermes__.logger.info(f"Datamodel has changed: {diff.dict}")
1699 self.__saveRequired = True
1701 # Start by removed types, as it requires the previous datamodel to process data
1702 # removal
1703 if diff.removed:
1704 for l_objtype in diff.removed:
1705 __hermes__.logger.info(f"About to purge data from type '{l_objtype}'")
1706 for ltypes in self.__datamodel.typesmapping.values():
1707 if l_objtype in ltypes:
1708 break
1709 else:
1710 # l_objtype was not found in each list from
1711 # self.__datamodel.typesmapping.values()
1712 __hermes__.logger.warning(
1713 f"Requested to purge data from type '{l_objtype}', but it"
1714 " doesn't exist in previous datamodel: ignoring"
1715 )
1716 continue
1718 pkeys = (
1719 self.__datamodel.localdata[l_objtype].getPKeys()
1720 | self.__datamodel.localdata[f"trashbin_{l_objtype}"].getPKeys()
1721 )
1723 # Call remove on each object
1724 for pkey in sorted(pkeys):
1725 _, l_obj = Datamodel.getObjectFromCacheOrTrashbin(
1726 self.__datamodel.localdata, l_objtype, pkey
1727 )
1728 if l_obj:
1729 l_ev = Event(
1730 evcategory="base",
1731 eventtype="removed",
1732 obj=l_obj,
1733 objattrs={},
1734 )
1735 secretAttrs = self.__datamodel.local_schema.secretsAttributesOf(
1736 l_objtype
1737 )
1738 __hermes__.logger.debug(
1739 f"Removing local object of {pkey=}:"
1740 f" {l_ev.toString(secretAttrs)=}"
1741 )
1742 try:
1743 self.__processLocalEvent(
1744 None,
1745 l_ev,
1746 enqueueEventWithError=False,
1747 dontStoreInTrashbin=True,
1748 )
1749 except Exception:
1750 raise HermesDatamodelUpdatePurgeError(
1751 "An error was met when purging data from type"
1752 f" '{l_objtype}'. App will stop, but purging will"
1753 " resume upon restart"
1754 )
1755 else:
1756 __hermes__.logger.error(f"Local object of {pkey=} not found")
1758 # All objects have been removed, remove remaining events from
1759 # errorqueue, if any
1760 for l_obj in self.__datamodel.localdata_complete[l_objtype]:
1761 # Remove eventual events relative to current object from error queue
1762 self.__datamodel.errorqueue.purgeAllEventsOfDataObject(
1763 l_obj, isLocalObjtype=True
1764 )
1766 self.__datamodel.saveLocalAndRemoteData() # Save changes
1768 # Purge old remote and local cache files, if any
1769 __hermes__.logger.info(
1770 f"Types removed from Datamodel: {diff.removed}, purging cache files"
1771 )
1772 Datamodel.purgeOldCacheFiles(diff.removed, cacheFilePrefix="__")
1774 # Start working with new datamodel
1775 self.__datamodel = self.__newdatamodel
1776 self.__datamodel.loadLocalAndRemoteData()
1778 # Reload error queue to allow it to handle new datamodel
1779 self.__datamodel.saveErrorQueue()
1780 self.__datamodel.loadErrorQueue()
1782 if diff.added or diff.modified:
1783 # Generate diff events according to datamodel changes
1784 new_local_data: dict[str, DataObjectList] = {}
1786 # Work on "complete" copy of data cache, representing the cache that should
1787 # be without any event in error queue in order to compute diff on complete
1788 # cache representation
1789 completeRemoteData = self.__datamodel.remotedata_complete
1790 completeLocalData = self.__datamodel.localdata_complete
1792 # Loop over each remote type
1793 for r_objtype in self.__datamodel.remote_schema.objectTypes:
1794 # For each type, we'll work on classic data, and on trashbin data
1795 for prefix in ("", "trashbin_"):
1796 if r_objtype not in self.__datamodel.typesmapping:
1797 continue # Remote objtype isn't set in current Datamodel
1799 # Fetch corresponding local type
1800 for l_objt in self.__datamodel.typesmapping[r_objtype]:
1801 l_objtype = f"{prefix}{l_objt}"
1803 # Convert remote data cache to local data
1804 new_local_data[l_objtype] = (
1805 self.__datamodel.convertDataObjectListToLocal(
1806 r_objtype,
1807 completeRemoteData[f"{prefix}{r_objtype}"],
1808 l_objt,
1809 )
1810 )
1812 # Compute differences between new local data and local data
1813 # cache
1814 completeLocalDataObjtype: DataObjectList = (
1815 completeLocalData.get(l_objtype, DataObjectList([]))
1816 )
1817 datadiff = new_local_data[l_objtype].diffFrom(
1818 completeLocalDataObjtype
1819 )
1821 for changeType, difflist in datadiff.dict.items():
1822 diffitem: DiffObject | DataObject
1823 for diffitem in difflist:
1824 # Convert diffitem to local Event
1825 event, obj = Event.fromDiffItem(
1826 diffitem=diffitem,
1827 eventCategory="base",
1828 changeType=changeType,
1829 )
1831 if prefix == "trashbin_":
1832 if obj not in completeLocalDataObjtype:
1833 # Object exists in remote trashbin, but not in
1834 # local one as it has been removed before its
1835 # type was added to client's Datamodel.
1836 # Process a local "added" event, then a local
1837 # "removed" event to store local object in
1838 # trashbin
1840 # Add local object
1841 self.__processLocalEvent(
1842 None, event, enqueueEventWithError=True
1843 )
1845 # Prepare "removed" event
1846 event = Event(
1847 evcategory="base",
1848 eventtype="removed",
1849 obj=obj,
1850 objattrs={},
1851 )
1852 # Preserve object _trashbin_timestamp
1853 event.timestamp = completeRemoteData[
1854 f"{prefix}{r_objtype}"
1855 ][obj]._trashbin_timestamp
1856 else:
1857 # Preserve object _trashbin_timestamp
1858 obj._trashbin_timestamp = (
1859 completeLocalDataObjtype[
1860 obj
1861 ]._trashbin_timestamp
1862 )
1864 # Process Event and update cache if no error is met,
1865 # enqueue event otherwise
1866 self.__processLocalEvent(
1867 None, event, enqueueEventWithError=True
1868 )
1870 self.__config.savecachefile() # Save config to be able to rebuild datamodel
1871 self.__datamodel.saveLocalAndRemoteData() # Save data
1872 self.__checkDatamodelWarnings()
1874 def __status(
1875 self, verbose=False, level="information", ignoreUnhandledExceptions=False
1876 ) -> dict[str, dict[str, dict[str, Any]]]:
1877 """Returns a dict containing status for current client Datamodel and error
1878 queue.
1880 Each status contains 3 categories/levels: "information", "warning" and "error"
1881 """
1882 if level not in ("information", "warning", "error"):
1883 raise AttributeError(
1884 f"Specified level '{level}' is invalid. Possible values are"
1885 """("information", "warning", "error"):"""
1886 )
1888 appname: str = self.__config["appname"]
1890 match level:
1891 case "error":
1892 levels = ["error"]
1893 case "warning":
1894 levels = [
1895 "warning",
1896 "error",
1897 ]
1898 case "information":
1899 levels = [
1900 "information",
1901 "warning",
1902 "error",
1903 ]
1905 res = {
1906 appname: {
1907 "information": {
1908 "startTime": self.__startTime.strftime("%Y-%m-%d %H:%M:%S"),
1909 "status": "paused" if self.__isPaused else "running",
1910 "pausedSince": (
1911 self.__isPaused.strftime("%Y-%m-%d %H:%M:%S")
1912 if self.__isPaused
1913 else "None"
1914 ),
1915 },
1916 "warning": {},
1917 "error": {},
1918 },
1919 }
1920 if not ignoreUnhandledExceptions and self.__cache.exception:
1921 res[appname]["error"]["unhandledException"] = self.__cache.exception
1923 # Datamodel
1924 res["datamodel"] = {
1925 "information": {},
1926 "warning": {},
1927 "error": {},
1928 }
1929 if self.__datamodel.unknownRemoteTypes:
1930 res["datamodel"]["warning"]["unknownRemoteTypes"] = sorted(
1931 self.__datamodel.unknownRemoteTypes
1932 )
1933 if self.__datamodel.unknownRemoteAttributes:
1934 res["datamodel"]["warning"]["unknownRemoteAttributes"] = {
1935 k: sorted(v)
1936 for k, v in self.__datamodel.unknownRemoteAttributes.items()
1937 }
1939 # Error queue
1940 res["errorQueue"] = {
1941 "information": {},
1942 "warning": {},
1943 "error": {},
1944 }
1946 if self.__datamodel.errorqueue is not None:
1947 eventNumber: int
1948 remoteEvent: Event | None
1949 localEvent: Event
1950 errorMsg: str
1951 for (
1952 eventNumber,
1953 remoteEvent,
1954 localEvent,
1955 errorMsg,
1956 ) in self.__datamodel.errorqueue:
1957 if errorMsg is None:
1958 # Ignore the events in queue that are not errors
1959 continue
1961 # Always try to get object from local cache in order to use configured
1962 # toString template for obj repr()
1963 objtype = localEvent.objtype
1965 _, obj = Datamodel.getObjectFromCacheOrTrashbin(
1966 self.__datamodel.localdata_complete, objtype, localEvent.objpkey
1967 )
1968 if obj is None:
1969 _, obj = Datamodel.getObjectFromCacheOrTrashbin(
1970 self.__datamodel.localdata, objtype, localEvent.objpkey
1971 )
1972 if obj is None and remoteEvent is not None:
1973 _, obj = Datamodel.getObjectFromCacheOrTrashbin(
1974 self.__datamodel.remotedata_complete,
1975 remoteEvent.objtype,
1976 remoteEvent.objpkey,
1977 )
1978 if obj is None:
1979 obj = f"<{localEvent.objtype}[{localEvent.objpkey}]>"
1980 else:
1981 obj = repr(obj)
1983 res["errorQueue"]["error"][eventNumber] = {
1984 "objrepr": obj,
1985 "errorMsg": errorMsg,
1986 }
1987 if verbose:
1988 res["errorQueue"]["error"][eventNumber] |= {
1989 "objtype": localEvent.objtype,
1990 "objpkey": localEvent.objpkey,
1991 "objattrs": localEvent.objattrs,
1992 }
1994 # Clean empty categories
1995 for objname in list(res.keys()):
1996 for category in ("information", "warning", "error"):
1997 if category not in levels or (
1998 not verbose and not res[objname][category]
1999 ):
2000 del res[objname][category]
2002 if not verbose and not res[objname]:
2003 del res[objname]
2005 return res
2007 def __notifyQueueErrors(self):
2008 """Notify of any objects change in error queue"""
2009 new_error = self.__status(level="error", ignoreUnhandledExceptions=True)
2011 new_errors = {}
2012 if "errorQueue" in new_error:
2013 for errNumber, err in new_error["errorQueue"]["error"].items():
2014 new_errors[errNumber] = f"{err['objrepr']}: {err['errorMsg']}"
2016 new_errstr = json.dumps(
2017 new_errors,
2018 cls=JSONEncoder,
2019 indent=4,
2020 )
2021 old_errstr = json.dumps(self.__cache.queueErrors, cls=JSONEncoder, indent=4)
2023 if new_errstr != old_errstr:
2024 if new_errors:
2025 desc = "objects in error queue have changed"
2026 else:
2027 desc = "no more objects in error queue"
2029 __hermes__.logger.info(desc)
2030 Email.sendDiff(
2031 config=self.__config,
2032 contentdesc=desc,
2033 previous=old_errstr,
2034 current=new_errstr,
2035 )
2036 self.__cache.queueErrors = new_errors
2038 def __notifyDatamodelWarnings(self):
2039 """Notify of any data model warnings changes"""
2040 new_errors = self.__status(level="warning", ignoreUnhandledExceptions=True)
2042 if "datamodel" not in new_errors:
2043 new_errors = {}
2044 else:
2045 new_errors = new_errors["datamodel"]
2047 new_errstr = json.dumps(
2048 new_errors,
2049 cls=JSONEncoder,
2050 indent=4,
2051 )
2052 old_errstr = json.dumps(
2053 self.__cache.datamodelWarnings, cls=JSONEncoder, indent=4
2054 )
2056 if new_errors:
2057 __hermes__.logger.error("Datamodel has warnings:\n" + new_errstr)
2059 if new_errstr != old_errstr:
2060 if new_errors:
2061 desc = "datamodel warnings have changed"
2062 else:
2063 desc = "no more datamodel warnings"
2065 __hermes__.logger.info(desc)
2066 Email.sendDiff(
2067 config=self.__config,
2068 contentdesc=desc,
2069 previous=old_errstr,
2070 current=new_errstr,
2071 )
2072 self.__cache.datamodelWarnings = new_errors
2074 def __notifyException(self, trace: str | None):
2075 """Notify of any unhandled exception met/solved"""
2076 if trace:
2077 __hermes__.logger.critical(f"Unhandled exception: {trace}")
2079 if self.__cache.exception != trace:
2080 if trace:
2081 desc = "unhandled exception"
2082 else:
2083 desc = "no more unhandled exception"
2085 __hermes__.logger.info(desc)
2086 previous = "" if self.__cache.exception is None else self.__cache.exception
2087 current = "" if trace is None else trace
2088 Email.sendDiff(
2089 config=self.__config,
2090 contentdesc=desc,
2091 previous=previous,
2092 current=current,
2093 )
2094 self.__cache.exception = trace
2096 def __notifyFatalException(self, trace: str):
2097 """Notify of any fatal exception met before, and terminate app"""
2098 self.__isStopped = True
2100 desc = "Unhandled fatal exception, APP WILL TERMINATE IMMEDIATELY"
2101 NL = "\n"
2103 __hermes__.logger.critical(f"{desc}: {trace}")
2105 Email.send(
2106 config=self.__config,
2107 subject=f"[{self.__config['appname']}] {desc}",
2108 content=f"{desc}:{NL}{NL}{trace}",
2109 )