from logging import Formatter, StreamHandler, getLogger from ingestor.managers.ingestor_manager import IngestorManager from os import getenv from time import sleep def main(): # Get os parameters kafka_servers = getenv("KAFKA_SERVERS", "localhost:9092") redis_host = getenv("REDIS_HOST", "localhost") redis_port = int(getenv("REDIS_PORT", 6379)) lease_ttl = int(getenv("LEASE_TTL", 10)) heartbeat_ttl = int(getenv("HEARTBEAT_TTL", 20)) pod_id = getenv("HOSTNAME", "localhost") poll_interval = int(getenv("POLL_INTERVAL", 5)) kafka_servers = kafka_servers.split(",") logger = getLogger(__name__) logger.setLevel(getenv("LOG_LEVEL", "INFO")) handler = StreamHandler() formatter = Formatter( '%(asctime)s - %(name)s - %(levelname)s - %(message)s') handler.setFormatter(formatter) logger.addHandler(handler) ingestor_manager = IngestorManager( kafka_servers, redis_host, redis_port, lease_ttl, heartbeat_ttl, pod_id, poll_interval, logger ) # Declare ingestor ative ingestor_manager.declare_active() # Get slot lease acquired = ingestor_manager.get_slot_leases() logger.info(f"Acquired slots: {acquired}") if not acquired: logger.warning("No slots available") else: # Subscribe to acquired slots ingestor_manager.update_opc_servers() ingestor_manager.subscribe_to_tags(acquired) while True: # Declare ingestor as active ingestor_manager.declare_active() logger.info("Polling for slot updates...") # Get active ingestors ingestors = ingestor_manager.get_active_ingestors() number_of_slots = ingestor_manager.get_number_of_slots() if not ingestor_manager.managed_tags and number_of_slots > 0: # This ingestor is active and has no slots, so we need to try to # acquire a slot lease logger.info("No slots acquired, trying to acquire a slot lease") acquired = ingestor_manager.get_slot_leases(1) if not acquired: logger.info("No slots acquired") else: # Subscribe to acquired slots ingestor_manager.update_opc_servers() ingestor_manager.subscribe_to_tags(acquired) ingestor_diff = number_of_slots - len(ingestors) slot_diff = len(ingestor_manager.managed_tags) - 1 if ingestor_diff > 0: # Some ingestors are innactive, so theres "ingestor_diff" slots available logger.info(f"Slots available: {ingestor_diff}") # Get slot lease acquired = ingestor_manager.get_slot_leases(ingestor_diff) if not acquired: logger.info("No slots acquired") else: # Subscribe to acquired slots ingestor_manager.update_opc_servers() ingestor_manager.subscribe_to_tags(acquired) elif slot_diff > 0: # Some ingestors are active and without slots, so we need to drop overleases = list(ingestor_manager.managed_tags.keys())[1:] ingestor_manager.drop_slot_leases(overleases) logger.info( f"Active ingestors: {ingestors}, " f"Number of slots: {number_of_slots}, " f"Managed tags: {ingestor_manager.managed_tags}" f"Managed servers: {ingestor_manager.opc_managers}" ) if not ingestor_manager.managed_tags: # No slots acquired logger.info("No slots acquired in this loop") # Update opc servers ingestor_manager.update_slot_config() # Sleep for poll interval sleep(poll_interval) if __name__ == "__main__": main()