Coverage for clients/__init__.py: 82%

896 statements  

« prev     ^ index     » next       coverage.py v7.14.1, created at 2026-06-11 15:40 +0000

1#!/usr/bin/env python3 

2# -*- coding: utf-8 -*- 

3 

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

21 

22 

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) 

42 

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 

52 

53 

54class HermesAlreadyNotifiedException(Exception): 

55 """Raised when an exception has already been notified, to avoid a second 

56 notification""" 

57 

58 

59class HermesDatamodelUpdatePurgeError(Exception): 

60 """Raised when an error is met when purging data on datamodel update""" 

61 

62 

63class HermesClientHandlerError(Exception): 

64 """Raised when an exception is met during client handler call""" 

65 

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) 

74 

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 ) 

83 

84 if purgeCurrentFileFromTrace: 

85 # Purging current file infos from traceback 

86 lines = [line for line in lines if __file__ not in line] 

87 

88 return "".join(lines).strip() 

89 

90 

91class HermesClientCache(LocalCache): 

92 """Hermes client data to cache""" 

93 

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 ) 

105 

106 self.queueErrors: dict[str, str] = from_json_dict.get("queueErrors", {}) 

107 """Dictionary containing current objects in error queue, for notifications""" 

108 

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

113 

114 self.exception: str | None = from_json_dict.get("exception") 

115 """String containing latest exception trace""" 

116 

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

125 

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) 

129 

130 

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

136 

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

146 

147 def __init__(self, config: HermesConfig): 

148 """Instantiate a new client""" 

149 

150 __hermes__.logger.info(f"Starting {config['appname']} v{HERMES_VERSION}") 

151 

152 # Setup the signals handler 

153 config.setSignalsHandler(self.__signalHandler) 

154 

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 = {} 

162 

163 self.__previousconfig: HermesConfig = HermesConfig.loadcachefile( 

164 "_hermesconfig" 

165 ) 

166 """Previous config (from cache)""" 

167 

168 self.__msgbus: AbstractMessageBusConsumerPlugin = self.__config["hermes"][ 

169 "plugins" 

170 ]["messagebus"]["plugininstance"] 

171 self.__msgbus.setTimeout( 

172 self.__config["hermes-client"]["updateInterval"] * 1000 

173 ) 

174 

175 self.__cache: HermesClientCache = HermesClientCache.loadcachefile( 

176 f"_{self.__config['appname']}" 

177 ) 

178 """Cached attributes""" 

179 self.__cache.setCacheFilename(f"_{self.__config['appname']}") 

180 

181 self.__startTime: datetime | None = None 

182 """Datetime when mainloop was started""" 

183 

184 self.__isPaused: datetime | None = None 

185 """Contains pause datetime if standard processing is paused, None otherwise""" 

186 

187 self.__isStopped: bool = False 

188 """mainloop() will run until this var is set to True""" 

189 

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

193 

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

211 

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

217 

218 self.__newdatamodel: Datamodel = Datamodel(config=self.__config) 

219 """New datamodel (from current config)""" 

220 

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 

229 

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 ) 

237 

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

242 

243 self.__trashbin_lastpurge: datetime = datetime(year=1, month=1, day=1) 

244 """Datetime when latest trashbin purge was ran""" 

245 

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

250 

251 self.__errorQueue_lastretry: datetime = datetime(year=1, month=1, day=1) 

252 """Datetime when latest error queue retry was ran""" 

253 

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

257 

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

261 

262 self.__isAnErrorRetry: bool = False 

263 """Indicate to handler whether the current event is being processed as part of 

264 an error retry""" 

265 

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

271 

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

278 

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 

284 

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 ) 

290 

291 return deepcopy(obj) 

292 

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}"] 

301 

302 # Create an empty DataObjectList of same type as cache 

303 res = type(cache)(objlist=[]) 

304 

305 res.extend(cache) 

306 res.extend(trashbin) 

307 return res 

308 

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 

315 

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 ) 

323 

324 subparsers = self.__parser.add_subparsers(help="Sub-commands") 

325 

326 # Quit 

327 sp_quit = subparsers.add_parser("quit", help=f"Stop {self.__config['appname']}") 

328 sp_quit.set_defaults(func=self.__sock_quit) 

329 

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) 

335 

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) 

341 

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 ) 

363 

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 

370 

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 

387 

388 if reply is None: # Error was met 

389 reply = SocketMessageToClient(retcode=1, retmsg=retmsg) 

390 

391 return reply 

392 

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

398 

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 ) 

406 

407 if self.__isPaused: 

408 return SocketMessageToClient( 

409 retcode=1, 

410 retmsg=f"Error: {self.__config['appname']} is already paused", 

411 ) 

412 

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

418 

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 ) 

426 

427 if not self.__isPaused: 

428 return SocketMessageToClient( 

429 retcode=1, retmsg=f"Error: {self.__config['appname']} is not paused" 

430 ) 

431 

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

437 

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 

458 

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

467 

468 return SocketMessageToClient(retcode=0, retmsg=msg) 

469 

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 

475 

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 ) 

483 

484 if value < 0: 

485 raise ValueError(f"Specified step {value=} must be greater or equal to 0") 

486 

487 self.__currentStep = value 

488 

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 

494 

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 ) 

502 

503 self.__isPartiallyProcessed = value 

504 

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 

510 

511 def mainLoop(self): 

512 """Client main loop""" 

513 self.__startTime = datetime.now() 

514 

515 self.__checkDatamodelWarnings() 

516 # TODO: implement a check to ensure subclasses required data types and 

517 # attributes exist in datamodel 

518 

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

527 

528 if self.__sock is not None: 

529 self.__sock.startProcessMessagesDaemon(appname=__hermes__.appname) 

530 

531 # Reduce sleep duration during functional tests to speed them up 

532 sleepDuration = 1 if self.__numberOfLoopToProcess is None else 0.05 

533 

534 isFirstLoopIteration: bool = True 

535 while not self.__isStopped: 

536 self.__saveRequired = False 

537 currentConfigCanBeSaved = True 

538 

539 try: 

540 with self.__msgbus: 

541 if self.__isPaused or self.__numberOfLoopToProcess == 0: 

542 sleep(sleepDuration) 

543 continue 

544 

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 ) 

579 

580 self.__notifyException(None) 

581 

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

628 

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

642 

643 # Only used in functionnal tests 

644 if self.__numberOfLoopToProcess: 

645 self.__numberOfLoopToProcess -= 1 

646 

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 

652 

653 # All events processed in __retryErrorQueue() require this attribute to be True 

654 self.__isAnErrorRetry = True 

655 

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

664 

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 

675 

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 

686 

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 

747 

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

761 

762 self.__isAnErrorRetry = False # End of __retryErrorQueue() 

763 

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 

774 

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 

782 

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 ) 

808 

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 

813 

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 

821 

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 

831 

832 def __hasAtLeastBeganInitialization(self) -> bool: 

833 return ( 

834 self.__cache.initstartoffset is not None 

835 and self.__datamodel.hasRemoteSchema() 

836 ) 

837 

838 def __canBeInitialized(self) -> bool: 

839 self.__msgbus.seekToBeginning() 

840 

841 # List of (start, stop) offsets of initsync sequences found 

842 initSyncFound: list[tuple[Any, Any]] = [] 

843 

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 

863 

864 if not initSyncFound: 

865 return False 

866 

867 if self.__useFirstInitsyncSequence: 

868 start, stop = initSyncFound[0] 

869 else: 

870 start, stop = initSyncFound[-1] 

871 

872 if self.__cache.nextoffset is None or self.__cache.nextoffset < start: 

873 self.__cache.nextoffset = start 

874 

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 

882 

883 def __updateSchema(self, newSchema: Dataschema): 

884 if self.__datamodel.forcePurgeOfTrashedObjectsWithoutNewPkeys( 

885 self.__datamodel.remote_schema, newSchema 

886 ): 

887 self.__emptyTrashBin(force=True) 

888 

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

895 

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 ) 

903 

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

910 

911 def __processEvents(self, isInitSync=False): 

912 remote_event: Event 

913 schema: Dataschema | None = None 

914 evcategory: str = "initsync" if isInitSync else "base" 

915 

916 self.__msgbus.seek(self.__cache.nextoffset) 

917 

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 

924 

925 # TODO: implement data consistency check if event.evcategory==initsync 

926 # and evcategory==base 

927 

928 if remote_event.evcategory != evcategory: 

929 self.__cache.nextoffset = remote_event.offset + 1 

930 continue 

931 

932 if isInitSync: 

933 if remote_event.eventtype == "init-start": 

934 schema = Dataschema.from_json(remote_event.objattrs) 

935 self.__updateSchema(schema) 

936 continue 

937 

938 if remote_event.eventtype == "init-stop": 

939 self.__cache.nextoffset = remote_event.offset + 1 

940 break 

941 

942 if schema is None and not self.__hasAtLeastBeganInitialization(): 

943 msg = "Invalid initsync sequence met, ignoring" 

944 __hermes__.logger.critical(msg) 

945 return 

946 

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 ) 

961 

962 self.__cache.nextoffset = remote_event.offset + 1 

963 

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 

968 

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 

983 

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 ) 

1002 

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] 

1012 

1013 trashbin = self.__datamodel.remotedata[f"trashbin_{remote_event.objtype}"] 

1014 was_in_trashbin = remote_event.objpkey in trashbin 

1015 

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 ) 

1023 

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 

1067 

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) 

1079 

1080 case "modified": 

1081 self.__remoteModified(remote_event, local_event, simulateOnly) 

1082 

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 

1109 

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 

1126 

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 

1138 

1139 secretAttrs = self.__datamodel.local_schema.secretsAttributesOf( 

1140 local_event.objtype 

1141 ) 

1142 __hermes__.logger.debug( 

1143 f"__processLocalEvent({local_event.toString(secretAttrs)})" 

1144 ) 

1145 

1146 self.__saveRequired = True 

1147 

1148 if not simulateOnly: 

1149 # Reset current step 

1150 self.currentStep = local_event.step 

1151 self.isPartiallyProcessed = local_event.isPartiallyProcessed 

1152 

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 

1183 

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 

1243 

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

1251 

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 ) 

1261 

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 ) 

1272 

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

1278 

1279 l_obj = self.__datamodel.createLocalDataobject( 

1280 local_ev.objtype, local_ev.objattrs 

1281 ) 

1282 

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 ) 

1292 

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 ) 

1302 

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 ] 

1318 

1319 r_obj = self.__datamodel.createRemoteDataobject( 

1320 remote_event.objtype, remote_event.objattrs 

1321 ) 

1322 

1323 self.__processLocalEvent( 

1324 remote_event, 

1325 local_event, 

1326 enqueueEventWithError=False, 

1327 simulateOnly=simulateOnly, 

1328 ) 

1329 

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) 

1340 

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 ] 

1352 

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 

1358 

1359 l_obj_trash_complete: DataObject = deepcopy( 

1360 trashbin_complete.get(local_ev.objpkey) 

1361 ) 

1362 

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 

1368 

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 ) 

1378 

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) 

1384 

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) 

1389 

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) 

1407 

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] 

1418 

1419 if not simulateOnly: 

1420 r_cachedobj: DataObject = maincache.get(remote_event.objpkey) 

1421 r_obj = Datamodel.getUpdatedObject(r_cachedobj, remote_event.objattrs) 

1422 

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 ) 

1431 

1432 self.__processLocalEvent( 

1433 remote_event, 

1434 local_event, 

1435 enqueueEventWithError=False, 

1436 simulateOnly=simulateOnly, 

1437 ) 

1438 

1439 # Update remote object in cache 

1440 if not simulateOnly: 

1441 maincache.replace(r_obj) 

1442 

1443 # May not exist 

1444 if cache_complete is not None: 

1445 cache_complete.replace(r_obj_complete) 

1446 

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] 

1453 

1454 if not simulateOnly: 

1455 l_cachedobj: DataObject = maincache.get(local_ev.objpkey) 

1456 l_obj = Datamodel.getUpdatedObject(l_cachedobj, local_ev.objattrs) 

1457 

1458 cache_complete, l_cachedobj_complete = Datamodel.getObjectFromCacheOrTrashbin( 

1459 self.__datamodel.localdata_complete, local_ev.objtype, local_ev.objpkey 

1460 ) 

1461 

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 ) 

1467 

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 ) 

1478 

1479 # Update local object in cache 

1480 if not simulateOnly: 

1481 maincache.replace(l_obj) 

1482 

1483 # May not exist 

1484 if cache_complete is not None: 

1485 cache_complete.replace(l_obj_complete) 

1486 

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 ] 

1502 

1503 r_cachedobj: DataObject = maincache.get(remote_event.objpkey) 

1504 r_cachedobj_complete: DataObject = maincache_complete.get(remote_event.objpkey) 

1505 

1506 self.__processLocalEvent( 

1507 remote_event, 

1508 local_event, 

1509 enqueueEventWithError=False, 

1510 simulateOnly=simulateOnly, 

1511 ) 

1512 

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) 

1524 

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) 

1530 

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 ] 

1542 

1543 l_cachedobj: DataObject = maincache.get(local_ev.objpkey) 

1544 l_cachedobj_complete: DataObject = maincache_complete.get(local_ev.objpkey) 

1545 

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 ) 

1555 

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) 

1562 

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) 

1567 

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 ) 

1577 

1578 cache, r_cachedobj = Datamodel.getObjectFromCacheOrTrashbin( 

1579 self.__datamodel.remotedata, 

1580 remote_event.objtype, 

1581 remote_event.objpkey, 

1582 ) 

1583 

1584 cache_complete, r_cachedobj_complete = Datamodel.getObjectFromCacheOrTrashbin( 

1585 self.__datamodel.remotedata_complete, 

1586 remote_event.objtype, 

1587 remote_event.objpkey, 

1588 ) 

1589 

1590 self.__processLocalEvent( 

1591 remote_event, 

1592 local_event, 

1593 enqueueEventWithError=False, 

1594 simulateOnly=simulateOnly, 

1595 ) 

1596 

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) 

1602 

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) 

1609 

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

1615 

1616 cache, l_cachedobj = Datamodel.getObjectFromCacheOrTrashbin( 

1617 self.__datamodel.localdata, 

1618 local_ev.objtype, 

1619 local_ev.objpkey, 

1620 ) 

1621 

1622 cache_complete, l_cachedobj_complete = Datamodel.getObjectFromCacheOrTrashbin( 

1623 self.__datamodel.localdata_complete, 

1624 local_ev.objtype, 

1625 local_ev.objpkey, 

1626 ) 

1627 

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 ) 

1637 

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) 

1644 

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

1651 

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 ) 

1658 

1659 kwargsstr = ", ".join([f"{k}={repr(v)}" for k, v in kwargs_filtered.items()]) 

1660 

1661 hdlr = getattr(self, handlerName, None) 

1662 

1663 if not callable(hdlr): 

1664 __hermes__.logger.debug( 

1665 f"Calling '{handlerName}({kwargsstr})': handler '{handlerName}()'" 

1666 " doesn't exists" 

1667 ) 

1668 return 

1669 

1670 __hermes__.logger.info( 

1671 f"Calling '{handlerName}({kwargsstr})' - currentStep={self.currentStep}," 

1672 f" isPartiallyProcessed={self.isPartiallyProcessed}," 

1673 f" isAnErrorRetry={self.isAnErrorRetry}" 

1674 ) 

1675 

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) 

1684 

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

1688 

1689 diff = self.__newdatamodel.diffFrom(self.__datamodel) 

1690 

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 

1697 

1698 __hermes__.logger.info(f"Datamodel has changed: {diff.dict}") 

1699 self.__saveRequired = True 

1700 

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 

1717 

1718 pkeys = ( 

1719 self.__datamodel.localdata[l_objtype].getPKeys() 

1720 | self.__datamodel.localdata[f"trashbin_{l_objtype}"].getPKeys() 

1721 ) 

1722 

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

1757 

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 ) 

1765 

1766 self.__datamodel.saveLocalAndRemoteData() # Save changes 

1767 

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="__") 

1773 

1774 # Start working with new datamodel 

1775 self.__datamodel = self.__newdatamodel 

1776 self.__datamodel.loadLocalAndRemoteData() 

1777 

1778 # Reload error queue to allow it to handle new datamodel 

1779 self.__datamodel.saveErrorQueue() 

1780 self.__datamodel.loadErrorQueue() 

1781 

1782 if diff.added or diff.modified: 

1783 # Generate diff events according to datamodel changes 

1784 new_local_data: dict[str, DataObjectList] = {} 

1785 

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 

1791 

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 

1798 

1799 # Fetch corresponding local type 

1800 for l_objt in self.__datamodel.typesmapping[r_objtype]: 

1801 l_objtype = f"{prefix}{l_objt}" 

1802 

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 ) 

1811 

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 ) 

1820 

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 ) 

1830 

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 

1839 

1840 # Add local object 

1841 self.__processLocalEvent( 

1842 None, event, enqueueEventWithError=True 

1843 ) 

1844 

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 ) 

1863 

1864 # Process Event and update cache if no error is met, 

1865 # enqueue event otherwise 

1866 self.__processLocalEvent( 

1867 None, event, enqueueEventWithError=True 

1868 ) 

1869 

1870 self.__config.savecachefile() # Save config to be able to rebuild datamodel 

1871 self.__datamodel.saveLocalAndRemoteData() # Save data 

1872 self.__checkDatamodelWarnings() 

1873 

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. 

1879 

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 ) 

1887 

1888 appname: str = self.__config["appname"] 

1889 

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 ] 

1904 

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 

1922 

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 } 

1938 

1939 # Error queue 

1940 res["errorQueue"] = { 

1941 "information": {}, 

1942 "warning": {}, 

1943 "error": {}, 

1944 } 

1945 

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 

1960 

1961 # Always try to get object from local cache in order to use configured 

1962 # toString template for obj repr() 

1963 objtype = localEvent.objtype 

1964 

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) 

1982 

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 } 

1993 

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] 

2001 

2002 if not verbose and not res[objname]: 

2003 del res[objname] 

2004 

2005 return res 

2006 

2007 def __notifyQueueErrors(self): 

2008 """Notify of any objects change in error queue""" 

2009 new_error = self.__status(level="error", ignoreUnhandledExceptions=True) 

2010 

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']}" 

2015 

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) 

2022 

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" 

2028 

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 

2037 

2038 def __notifyDatamodelWarnings(self): 

2039 """Notify of any data model warnings changes""" 

2040 new_errors = self.__status(level="warning", ignoreUnhandledExceptions=True) 

2041 

2042 if "datamodel" not in new_errors: 

2043 new_errors = {} 

2044 else: 

2045 new_errors = new_errors["datamodel"] 

2046 

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 ) 

2055 

2056 if new_errors: 

2057 __hermes__.logger.error("Datamodel has warnings:\n" + new_errstr) 

2058 

2059 if new_errstr != old_errstr: 

2060 if new_errors: 

2061 desc = "datamodel warnings have changed" 

2062 else: 

2063 desc = "no more datamodel warnings" 

2064 

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 

2073 

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

2078 

2079 if self.__cache.exception != trace: 

2080 if trace: 

2081 desc = "unhandled exception" 

2082 else: 

2083 desc = "no more unhandled exception" 

2084 

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 

2095 

2096 def __notifyFatalException(self, trace: str): 

2097 """Notify of any fatal exception met before, and terminate app""" 

2098 self.__isStopped = True 

2099 

2100 desc = "Unhandled fatal exception, APP WILL TERMINATE IMMEDIATELY" 

2101 NL = "\n" 

2102 

2103 __hermes__.logger.critical(f"{desc}: {trace}") 

2104 

2105 Email.send( 

2106 config=self.__config, 

2107 subject=f"[{self.__config['appname']}] {desc}", 

2108 content=f"{desc}:{NL}{NL}{trace}", 

2109 )