import argparse import sys from dbdeploy.base import config from dbdeploy.base.init import MainApp from dbdeploy.logic import acquire_deploy_lock from dbdeploy.logic import check_kafka_topic_exists from dbdeploy.logic import clean_up_previous_runs from dbdeploy.logic import create_kafka_topics from dbdeploy.logic import delete_kafka_topics from dbdeploy.logic import prepare_job_list from dbdeploy.logic import get_jobs_statuses from dbdeploy.logic import monitor_dlq from dbdeploy.logic import parse_db_pr_file from dbdeploy.logic import push_to_kafka from dbdeploy.logic import release_deploy_lock from dbdeploy.logic import start_kafka_connector from dbdeploy.logic import stop_kafka_connector from dbdeploy.logic import track_kafka_consumer_offsets from dbdeploy.logic import update_job_status from dbdeploy.logic import update_kafka_cluster_config from dbdeploy.util.exceptions import LockAcquisitionError from dbdeploy.util.exceptions import JobListGenerationError log = config.LOGGER def main(file_contents: str) -> None: """Main application flow. Args: file_contents: The contents of XML file from pull request """ log.debug('==== START PROCESSING ====') is_lock_acquired = acquire_deploy_lock( timeout=config.LOCK_ACQUISITION_TIMEOUT, force_lock_release=config.FORCE_LOCK_RELEASE) if not is_lock_acquired: raise LockAcquisitionError() clean_up_previous_runs() try: changeset_list = parse_db_pr_file(file_contents) changeset_list = get_jobs_statuses(changeset_list=changeset_list) job_list = prepare_job_list(changeset_list=changeset_list) except Exception as err: # release the lock in case there is an Exception release_deploy_lock() raise JobListGenerationError(str(err)) for job in job_list: dlq_monitor_thread = None log.debug(f'Running changeset id: {job.changeset_id}') try: if not job.skip_sink_connector: # Type snowflake-mysql or snowflake-neo4j: job = create_kafka_topics(job=job) job = start_kafka_connector(job=job) dlq_monitor_thread = monitor_dlq(job) else: # Type neo4j or snowflake-kafka or kafka: update_kafka_cluster_config(job=job) job = check_kafka_topic_exists(job=job) job = push_to_kafka(job=job) if not job.skip_sink_connector: # Type snowflake-mysql or snowflake-neo4j: job = track_kafka_consumer_offsets(job=job) log.debug('==== END PROCESSING ====') except Exception as err: log.error(f'Caught an exception. Error: {err}') log.info('Cleaning up...') if not job.dlq_error: # flag dlq to stop the thread in case main() got an exception job.main_error = err raise err finally: log.info('Finally...') # release the lock in case there is an Exception if not job.is_done: release_deploy_lock() if dlq_monitor_thread and dlq_monitor_thread.is_alive(): dlq_monitor_thread.join() job = update_job_status(job=job) if job.connector_started and not job.skip_sink_connector: stop_kafka_connector(job=job) # in case of exception delete topics after # connector is deleted # to make sure nothing is produced into dlq # and dlq topic is not recreated after delete if job.topics_created and not job.skip_sink_connector: delete_kafka_topics(job=job) release_deploy_lock() log.debug('==== STOPPED PROCESSING ====') if __name__ == "__main__": parser = argparse.ArgumentParser() parser.add_argument( 'inputfile', nargs='?', default=sys.stdin, type=argparse.FileType('r'), help='Please provide an input file, or pipe it via stdin') args = parser.parse_args() if not args.inputfile: sys.exit(0) app = MainApp(main, args.inputfile.read()) app.run()