from logging import Logger from typing import Dict, List from ingestor.managers.data_manager import DataManager from ingestor.managers.opc_manager import OpcManager from ingestor.managers.resource_manager import ResourceManager class IngestorManager(): def __init__(self, kafka_servers: str, redis_host: str, redis_port: int, lease_ttl: int, heartbeat_ttl: int, pod_id: str, number_of_ingestors: int, poll_interval: int, logger: Logger): self.data_manager = DataManager(kafka_servers, logger) self.opc_managers = {} self.resource_manager = ResourceManager( redis_host, redis_port, lease_ttl, heartbeat_ttl, pod_id ) self.number_of_ingestors = number_of_ingestors self.poll_interval = poll_interval self.logger = logger self.managed_tags = {} self.opc_servers = {} def update_opc_servers(self): for slot, slot_config in self.managed_tags.items(): for server, server_config in slot_config.items(): server_config = server_config.copy() server_config.pop('tags', None) if server not in self.opc_servers: self.opc_managers[server] = OpcManager.initialize_from_config( server_config, self.data_manager, self.logger ) self.opc_managers[server].connect() elif self.opc_servers[server] != server_config: del self.opc_managers[server] self.opc_managers[server] = OpcManager.initialize_from_config( server_config, self.data_manager, self.logger ) self.opc_managers[server].connect() self.opc_servers[server] = server_config def declare_active(self): self.resource_manager.ingestor_heartbeat() def get_active_ingestors(self) -> List[str]: self.resource_manager.get_all_ingestors() def get_slot_leases(self, max_slots: int = 1) -> Dict: acquired = {} for i in range(1, self.number_of_ingestors + 1): if self.resource_manager.lease_tag(str(i)): self.logger.info(f"Leased slot {i}") acquired[str(i)] = self.resource_manager.get_tag_slot(str(i)) if len(acquired) >= max_slots: self.managed_tags.update(acquired) return acquired self.logger.warning( f"Unable to acquire {max_slots} slots. " f"Only {acquired} slots were leased." ) self.managed_tags.update(acquired) return acquired def update_slot_config(self) -> Dict: for slot, slot_config in self.managed_tags.items(): update = self.resource_manager.get_tag_slot(slot) if update != slot_config: self.logger.info( f"Slot {slot} configuration updated. " f"Old: {slot_config}, New: {update}" ) self.managed_tags[slot] = update self.update_opc_servers() self.subscribe_to_tags(update) return update def drop_slot_leases(self, ids: List[str]) -> None: for id in ids: self.resource_manager.drop_tag_lease(id) def subscribe_to_tags(self, tags: Dict) -> None: for slot, slot_config in tags.items(): for server, server_config in slot_config.items(): if server not in self.opc_managers: self.logger.error( f"Server {server} not found in opc_managers." ) continue if slot not in self.opc_managers[server].subscriptions: self.opc_managers[server].create_subscription( slot ) self.opc_managers[server].subscribe( slot, server_config['tags'], self.poll_interval )