import argparse import datetime import os import subprocess import sys from pathlib import Path import dotenv from jinja2 import Environment, FileSystemLoader from tadas.platform.metrics import DurationMetrics from tadas.platform import context as contexts from tadas.platform import config from tadas.platform import locking as lock_utils from tadas.platform import logging as logging_utils from tadas.platform import sentry logger = logging_utils.get_file_logger(__file__) sentry.init_sentry() PROJECT_ROOT = Path(__file__).resolve().parents[2] FILES_DIR = PROJECT_ROOT / 'backend' / 'files' LOCK_FILE_PATH = FILES_DIR / 'launcher.lock' LOCK_TTL = datetime.timedelta(hours=2) def render_templates(source_dir, target_dir, context): env = Environment(loader=FileSystemLoader(source_dir), keep_trailing_newline=True) for root, dirs, files in os.walk(source_dir): rel_root = os.path.relpath(root, source_dir) output_root = os.path.join(target_dir, rel_root) os.makedirs(output_root, exist_ok=True) for filename in files: source_path = os.path.join(root, filename) if filename.endswith('.j2'): template_path = os.path.relpath(source_path, source_dir) output_filename = filename[:-3] # remove .j2 output_path = os.path.join(output_root, output_filename) template = env.get_template(template_path) rendered_content = template.render(**context) with open(output_path, 'w') as f: f.write(rendered_content) logger.info(f"Rendered: {template_path} -> {output_path}") def prepare(): render_templates( source_dir=f'{PROJECT_ROOT}/templates', target_dir=PROJECT_ROOT, context=dict( env=os.environ, ) ) logger.info('Listing content of backend/files ...') for file in FILES_DIR.glob('*'): modified_at = file.stat().st_mtime modified_at_str = datetime.datetime.fromtimestamp(modified_at).strftime('%Y-%m-%d %H:%M:%S') logger.info(f' {file} ({modified_at_str}) size: {file.stat().st_size} bytes') logger.info('DONE (Listing content of backend/files)') def run_cli(module, args): assert module args = args or [] command = ["python", "-m", module, *args] logger.info(f"Running command: {' '.join(command)}") return subprocess.run( command, cwd=PROJECT_ROOT, env=os.environ, ) def run_makefile(tadas_model_version, args: list[str], env: dict[str, str] = None): assert tadas_model_version env = env or {} assert all(isinstance(k, str) for k in env.keys()) assert all(isinstance(v, str) for v in env.values()) args = args or [] assert all(isinstance(v, str) for v in args) MODELS_PATH = PROJECT_ROOT / 'tadas' / 'models' model_path = MODELS_PATH / tadas_model_version assert model_path.exists(), f"Model {model_path} does not exist" command = ["make", *args] make_env = {**os.environ, **env, 'MODEL_VERSION': tadas_model_version} logger.info(f"Running command: {' '.join(command)} in {MODELS_PATH} with env {env} MODEL_VERSION={tadas_model_version}") return subprocess.run( command, cwd=MODELS_PATH, env=make_env, check=True ) def print_config(): logger.info("Configuration:") for option in config.DEFAULTS.keys(): config.get(option) # it logs the value def create_argsparser(): parser = argparse.ArgumentParser( description='Launcher for tadas-etl' ) parser.add_argument('--env', default='', help='Override ENV values in format of "NAME=value;NAME=value" (semicolon separated)') return parser def parse_env_values(env_str): """ Parse env values from string in format of NAME=value, colon separated :param env_str: :return: """ result = {} if not env_str: return result for key_value_str in env_str.split('&'): key, value = key_value_str.split('=', maxsplit=1) result[key] = value return result def launch(): logger.info("Launching the application...") legacy_versions_mapping = { 'initial': 't2503_reportingdb', 'v2025sep': 't2509_snowflake', } tadas_model_version = config.get('MODEL_VERSION') tadas_model_version = legacy_versions_mapping.get(tadas_model_version, tadas_model_version) make_tasks = config.get('MAKE_TASKS') report_date = config.get('REPORT_DATE') context = { 'report_date': report_date, 'model_version': tadas_model_version, 'environment': config.get('ENVIRONMENT'), 'make_tasks': make_tasks, } metrics = DurationMetrics( category='model_run_timings', context=contexts.load_context(), ) with contexts.add_context(contex=context), metrics: completed_process = run_makefile( tadas_model_version=tadas_model_version, args=make_tasks, env={ 'REPORT_DATE': report_date }, ) if completed_process.returncode != 0: logger.error(f"Failed to run makefile. Return code: {completed_process.returncode}") def main(args=None): parser = create_argsparser() args = parser.parse_args(args=args) env_str = args.env configure(env_str) if config.get('LAUNCHER_RELEASE_LOCK'): logger.info(f"Force releasing (possibly existing) lock: {LOCK_FILE_PATH}") LOCK_FILE_PATH.unlink(missing_ok=True) try: with lock_utils.try_lock(lock_path=LOCK_FILE_PATH, ttl=LOCK_TTL): prepare() launch() except lock_utils.LockAcquireException: logger.warning(f"Could not acquire lock. Another instance is running. Exiting...") def configure(env_str): env_overrides = parse_env_values(env_str) for key, value in env_overrides.items(): logger.info(f"Overriding ENV {key}={value}") os.environ[key] = value print_config() if __name__ == '__main__': logging_utils.init_logging(app_context='launcher') dotenv.load_dotenv() main(sys.argv[1:])