"""Lambda trigger_dag function module.""" import os import requests from lambdacommon.common_config import logger def trigger_airflow_dag(dag_id): """Triggers an Airflow DAG and returns the response.""" airflow_url = os.getenv("AIRFLOW_URL") # api_key = os.getenv("AIRFLOW_API_KEY") # headers = { # "Content-Type": "application/json", # "Authorization": f"Bearer {api_key}" # } headers = { "Content-Type": "application/json" } url = f"{airflow_url}/dags/{dag_id}/dagRuns" payload = { "conf": {} } response = requests.post(url, headers=headers, json=payload) return response def handler(event, context): """Lambda entry point.""" try: dag_name = event['dag_name'] response = trigger_airflow_dag(dag_name) return { 'status': 'OK', 'status_code': response.status_code, 'response': response.json() } except Exception as e: logger.exception(str(e)) raise e