import importlib import os from pathlib import Path from typing import List, Tuple from dagster import ( AssetChecksDefinition, AssetKey, AssetSpec, AssetsDefinition, DagsterInvalidDefinitionError, Definitions, JobDefinition, ScheduleDefinition, ResourceDefinition, SensorDefinition, ) from dagster_dbt import DbtCliResource from dagster_snowflake import SnowflakeResource from src.assets.snowflake.resources import snowflake as snowflake_resource def auto_load_from_directory( directory_path: str, ) -> Tuple[List[JobDefinition], List[AssetsDefinition], List[ScheduleDefinition], List[SensorDefinition], List[AssetChecksDefinition], dict]: jobs = [] schedules = [] assets = [] sensors = [] asset_checks = [] resources = {} py_files = Path(directory_path).rglob("*.py") base_dir = Path(os.path.dirname(__file__)) for file_path in py_files: module_path = str(file_path).replace("/", ".").replace("\\", ".") relative_path = file_path.relative_to(base_dir.parent) module_path = str(relative_path).replace("/", ".").replace("\\", ".") if module_path.endswith(".py"): module_path = module_path[:-3] try: module = importlib.import_module(module_path) for attr_name in dir(module): attr = getattr(module, attr_name) # Skip specific dbt_assets decorated functions if they follow the pattern # But allow our test assets through if (hasattr(attr, '__wrapped__') and hasattr(attr.__wrapped__, '__name__') and attr.__wrapped__.__name__ == 'dbt_models' and 'test' not in attr_name.lower()): continue # Handle both individual items and lists if isinstance(attr, list): # If it's a list, check each item's type for item in attr: if isinstance(item, JobDefinition): jobs.append(item) elif isinstance(item, (AssetsDefinition, AssetSpec)): assets.append(item) elif isinstance(item, ScheduleDefinition): schedules.append(item) elif isinstance(item, SensorDefinition): sensors.append(item) elif isinstance(item, AssetChecksDefinition): asset_checks.append(item) elif isinstance(item, (ResourceDefinition, DbtCliResource, SnowflakeResource)): resources[attr_name] = item else: # Handle individual items as before if isinstance(attr, JobDefinition): jobs.append(attr) elif isinstance(attr, (AssetsDefinition, AssetSpec)): assets.append(attr) elif isinstance(attr, ScheduleDefinition): schedules.append(attr) elif isinstance(attr, SensorDefinition): sensors.append(attr) elif isinstance(attr, AssetChecksDefinition): asset_checks.append(attr) elif isinstance(attr, (ResourceDefinition, DbtCliResource, SnowflakeResource)): resources[attr_name] = attr except Exception as e: error_msg = f"CRITICAL: Failed to load module {file_path}: {e}" print(f"{'='*80}") print(f"ERROR: {error_msg}") print(f"{'='*80}") # Re-raise the exception to fail loudly if "dbt_analytics" in str(file_path): #TODO[COR-518]: it should fail for other modules too. TBD raise RuntimeError(error_msg) from e return jobs, assets, schedules, sensors, asset_checks, resources base_dir = os.path.dirname(__file__) jobs_dir = os.path.join(base_dir, "jobs") bots_dir = os.path.join(base_dir, "assets/bots") rds_glue_dir = os.path.join(base_dir, "assets/rds_glue") postgres_dir = os.path.join(base_dir, "assets/postgres") snowflake_dir = os.path.join(base_dir, "assets/snowflake") trending_dir = os.path.join(base_dir, "assets/trending") recommendation_dir = os.path.join(base_dir, "assets/recommendation") schedules_dir = os.path.join(base_dir, "schedules") dbt_dir = os.path.join(base_dir, "assets/dbt") utils_dir = os.path.join(base_dir, "utils") # Load from each directory jobs, job_assets, job_schedules, job_sensors, job_asset_checks, job_resources = auto_load_from_directory(jobs_dir) schedule_jobs, schedule_assets, schedules, schedule_sensors, schedule_asset_checks, schedule_resources = auto_load_from_directory(schedules_dir) bot_jobs, bot_assets, bot_schedules, bot_sensors, bot_asset_checks, bot_resources = auto_load_from_directory(bots_dir) rds_glue_jobs, rds_glue_assets, rds_glue_schedules, rds_glue_sensors, rds_glue_asset_checks, rds_glue_resources = auto_load_from_directory(rds_glue_dir) postgres_jobs, postgres_assets, postgres_schedules, postgres_sensors, postgres_asset_checks, postgres_resources = auto_load_from_directory(postgres_dir) snowflake_jobs, snowflake_assets, snowflake_schedules, snowflake_sensors, snowflake_asset_checks, snowflake_resources = auto_load_from_directory(snowflake_dir) trending_jobs, trending_assets, trending_schedules, trending_sensors, trending_asset_checks, trending_resources = auto_load_from_directory(trending_dir) recommendation_jobs, recommendation_assets, recommendation_schedules, recommendation_sensors, recommendation_asset_checks, recommendation_resources = auto_load_from_directory( recommendation_dir ) dbt_jobs, dbt_assets, dbt_schedules, dbt_sensors, dbt_asset_checks, dbt_resources = auto_load_from_directory(dbt_dir) utils_jobs, utils_assets, utils_schedules, utils_sensors, utils_asset_checks, utils_resources = auto_load_from_directory(utils_dir) # Combine everything all_jobs = ( jobs + schedule_jobs + bot_jobs + rds_glue_jobs + postgres_jobs + snowflake_jobs + trending_jobs + recommendation_jobs + dbt_jobs + utils_jobs ) all_assets = ( job_assets + schedule_assets + bot_assets + rds_glue_assets + postgres_assets + snowflake_assets + trending_assets + recommendation_assets + dbt_assets + utils_assets ) all_schedules = ( schedules + job_schedules + bot_schedules + rds_glue_schedules + postgres_schedules + snowflake_schedules + trending_schedules + recommendation_schedules + dbt_schedules + utils_schedules ) all_sensors = ( job_sensors + schedule_sensors + bot_sensors + rds_glue_sensors + postgres_sensors + snowflake_sensors + trending_sensors + recommendation_sensors + dbt_sensors + utils_sensors ) all_asset_checks = ( job_asset_checks + schedule_asset_checks + bot_asset_checks + rds_glue_asset_checks + postgres_asset_checks + snowflake_asset_checks + trending_asset_checks + recommendation_asset_checks + dbt_asset_checks + utils_asset_checks ) all_resources = { **job_resources, **schedule_resources, **bot_resources, **rds_glue_resources, **postgres_resources, # Added this - it was missing! **snowflake_resources, **trending_resources, **recommendation_resources, **dbt_resources, **utils_resources, } # CRITICAL: Explicitly ensure snowflake resource is available with the correct key # This guarantees asset checks and assets can find it all_resources["snowflake"] = snowflake_resource # Optional: Debug logging to verify resources are loaded correctly # Uncomment these lines if you want to verify what's being loaded # print("🔍 Available resource keys:", list(all_resources.keys())) # print("🔍 Snowflake resource exists:", "snowflake" in all_resources) defs = Definitions( assets=all_assets, jobs=all_jobs, schedules=all_schedules, sensors=all_sensors, asset_checks=all_asset_checks, resources=all_resources, ) Definitions.validate_loadable(defs)