diff --git a/GEOSldas_App/util/inputs/obs_preproc/SMOS_preproc/SMOSproc_main.py b/GEOSldas_App/util/inputs/obs_preproc/SMOS_preproc/SMOSproc_main.py index 19bb729..76c99c9 100644 --- a/GEOSldas_App/util/inputs/obs_preproc/SMOS_preproc/SMOSproc_main.py +++ b/GEOSldas_App/util/inputs/obs_preproc/SMOS_preproc/SMOSproc_main.py @@ -28,12 +28,24 @@ TMP_NC_PATH = config['tmp_nc_path'] OUT_REG_PATH = config['out_reg_path'] + +try: + env_fail_fast = os.environ.get('FAIL_FAST').strip().lower() + fail_fast = env_fail_fast not in ['false'] + logging.info(f"fail_fast is set to {fail_fast}.") + +except: + logging.info(f"fail_fast is not set. Executing code with default behavior.") + fail_fast = False + def run_in_isolated_process(func, *args): """ Runs a function in an isolated child process and waits for it to finish. This prevents memory leaks or state changes in external libraries from affecting the main script. """ + logging.info(f"Initiation of {func.__name__}.") + p = multiprocessing.Process(target=func, args=args) p.start() p.join() @@ -115,9 +127,12 @@ def main(): try: ncflist = process_ee_to_nc(current_date) except RuntimeError as e: - logging.error(f"Halting processing for {date_str} due to error: {e}") - current_date += timedelta(days=1) - continue + logging.error(f"Halting processing for {date_str} due to error: {e} related to building nc file list") + if fail_fast: + sys.exit(1) + else: + current_date += timedelta(days=1) + continue # Step 2: Preprocess NetCDF files into REG for fnc in ncflist: @@ -125,8 +140,12 @@ def main(): try: run_in_isolated_process(preprocess_nc, fnc, config) except RuntimeError as e: - logging.error(f"Skipping {fnc} due to error: {e}") - continue + if fail_fast: + logging.error(f"Failure in preprocessing of {fnc} due to: {e}. Halting processing") + sys.exit(1) + else: + logging.error(f"Skipping {fnc} due to error: {e} in preprocessing of nc files") + continue # Move processed file to 'to_delete' dest = os.path.join(TMP_NC_PATH, 'to_delete', os.path.basename(fnc)) @@ -140,18 +159,82 @@ def main(): try: run_in_isolated_process(SCLF1C_reg2fit, OUT_REG_PATH, current_date, next_date, '_A') except RuntimeError as e: - logging.error(f"SCLF1C_reg2fit failed for Ascending (_A) on {date_str}: {e}") + if fail_fast: + logging.error(f"SCLF1C_reg2fit failed for Ascending (_A) on {date_str}: {e}. Exiting program.") + sys.exit(1) + else: + logging.error(f"SCLF1C_reg2fit failed for Ascending (_A) on {date_str}: {e}") logging.info("Running SCLF1C_reg2fit for Descending (_D)") try: run_in_isolated_process(SCLF1C_reg2fit, OUT_REG_PATH, current_date, next_date, '_D') except RuntimeError as e: - logging.error(f"SCLF1C_reg2fit failed for Descending (_D) on {date_str}: {e}") + if fail_fast: + logging.error(f"SCLF1C_reg2fit failed for Descending (_D) on {date_str}: {e}. Exiting program.") + sys.exit(1) + else: + logging.error(f"SCLF1C_reg2fit failed for Descending (_D) on {date_str}: {e}") logging.info(f"=== Completed processing to Tb40 for {date_str} ===") # Advance to the next day current_date += timedelta(days=1) + while current_date < end_time: + date_str = current_date.strftime('%Y%m%d') + logging.info(f"=== Starting processing for {date_str} ===") + + # Step 1: Convert files to NetCDF + try: + ncflist = process_ee_to_nc(current_date) + except RuntimeError as e: + logging.error(f"Halting processing for {date_str} due to error: {e} related to building nc file list") + + if fail_fast: + sys.exit(1) + else: + current_date += timedelta(days=1) + continue + + # Step 2: Preprocess NetCDF files into REG + for fnc in ncflist: + logging.info(f"Preprocessing NC file: {fnc}") + try: + run_in_isolated_process(preprocess_nc, fnc, config) + except RuntimeError as e: + if fail_fast: + logging.error(f"Failure in preprocessing of {fnc} due to: {e}. Halting processing") + sys.exit(1) + else: + logging.error(f"Skipping {fnc} due to error: {e} in preprocessing of nc files") + continue + + # Move processed file to 'to_delete' + dest = os.path.join(TMP_NC_PATH, 'to_delete', os.path.basename(fnc)) + shutil.move(fnc, dest) + logging.info(f"Moved {fnc} to cleanup directory.") + + # Step 3: Run REG to FIT processing for Ascending and Descending + next_date = current_date + timedelta(days=1) + + logging.info("Running SCLF1C_reg2fit for Ascending (_A)") + try: + run_in_isolated_process(SCLF1C_reg2fit, OUT_REG_PATH, current_date, next_date, '_A') + except RuntimeError as e: + if fail_fast: + logging.error(f"SCLF1C_reg2fit failed for Ascending (_A) on {date_str}: {e}. Exiting program.") + sys.exit(1) + else: + logging.error(f"SCLF1C_reg2fit failed for Ascending (_A) on {date_str}: {e}") + + logging.info("Running SCLF1C_reg2fit for Descending (_D)") + try: + run_in_isolated_process(SCLF1C_reg2fit, OUT_REG_PATH, current_date, next_date, '_D') + except RuntimeError as e: + if fail_fast: + logging.error(f"SCLF1C_reg2fit failed for Descending (_D) on {date_str}: {e}. Exiting program.") + sys.exit(1) + else: + logging.error(f"SCLF1C_reg2fit failed for Descending (_D) on {date_str}: {e}") if __name__ == '__main__': diff --git a/GEOSldas_App/util/inputs/obs_preproc/SMOS_preproc/modules.rc b/GEOSldas_App/util/inputs/obs_preproc/SMOS_preproc/modules.rc new file mode 100755 index 0000000..3ec1c0f --- /dev/null +++ b/GEOSldas_App/util/inputs/obs_preproc/SMOS_preproc/modules.rc @@ -0,0 +1,2 @@ +#!/bin/bash +module load python/GEOSpyD