diff --git a/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py b/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py index 9c64ea05c9..ccba453624 100644 --- a/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py +++ b/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py @@ -180,36 +180,45 @@ def _process_s3_input( paths_to_compress_buffer.add_file(object_metadata) -def _write_failed_path_log( - invalid_path_messages: List[str], logs_directory: Path, job_id: Any +def _write_user_failure_log( + title: str, + content: List[str], + logs_directory: Path, + job_id: Any, + filename_suffix: str, ) -> Optional[Path]: """ - Writes the error messages in `invalid_path_messages` to a log file, - `{logs_directory}/user/failed_paths_{job_id}.txt`. The directory will be created if it doesn't - already exist. - :param invalid_path_messages: + Writes a user-oriented failure log to + `{logs_directory}/user/job_{job_id}_{filename_suffix}.txt`. The `{logs_directory}/user` + directory will be created if it does not already exist. + + :param title: + :param content: :param logs_directory: :param job_id: - :return: Path to the written log file or `None` if error is encountered. + :param filename_suffix: + :return: Path to the written log file relative to `logs_directory`, or `None` on error. """ - - user_logs_dir = Path(logs_directory) / "user" + relative_log_path = Path("user") / f"job_{job_id}_{filename_suffix}.txt" + user_logs_dir = logs_directory / relative_log_path.parent try: user_logs_dir.mkdir(parents=True, exist_ok=True) - except Exception: + except Exception as e: + logger.error("Failed to create user logs directory: '%s' - %s", user_logs_dir, e) return None - log_path = user_logs_dir / f"failed_paths_{job_id}.txt" + log_path = logs_directory / relative_log_path try: with log_path.open("w", encoding="utf-8") as f: timestamp = datetime.datetime.now().isoformat(timespec="seconds") - f.write(f"Failed input paths log.\nGenerated at {timestamp}.\n\n") - for msg in invalid_path_messages: - f.write(f"{msg.rstrip()}\n") - except Exception: + f.write(f"{title}\nGenerated at {timestamp}.\n\n") + for item in content: + f.write(f"{item.rstrip()}\n") + except Exception as e: + logger.error("Failed to write compression failure user log: '%s' - %s", log_path, e) return None - return log_path + return relative_log_path def search_and_schedule_new_tasks( @@ -274,18 +283,22 @@ def search_and_schedule_new_tasks( if input_type == InputType.FS.value: invalid_path_messages = _process_fs_input_paths(input_config, paths_to_compress_buffer) if len(invalid_path_messages) > 0: - base_msg = "At least one of your input paths could not be processed." - - user_log_path = _write_failed_path_log( - invalid_path_messages, clp_config.logs_directory, job_id + user_log_relative_path = _write_user_failure_log( + title="Failed input paths log.", + content=invalid_path_messages, + logs_directory=clp_config.logs_directory, + job_id=job_id, + filename_suffix="failed_paths", + ) + if user_log_relative_path is None: + err_msg = "Failed to write user log for invalid input paths." + raise RuntimeError(err_msg) + + error_msg = ( + "At least one of your input paths could not be processed." + f" See the error log at '{user_log_relative_path}' inside your configured logs" + " directory (`logs_directory`) for more details." ) - if user_log_path is None: - error_msg = base_msg + ( - f" Check the compression scheduler logs in {clp_config.logs_directory} for" - " more details." - ) - else: - error_msg = base_msg + f" Check {user_log_path} for more details." update_compression_job_metadata( db_cursor, @@ -384,7 +397,7 @@ def search_and_schedule_new_tasks( scheduled_jobs[job_id] = job -def poll_running_jobs(db_conn, db_cursor): +def poll_running_jobs(logs_directory: Path, db_conn, db_cursor): """ Poll for running jobs and update their status. """ @@ -395,7 +408,7 @@ def poll_running_jobs(db_conn, db_cursor): for job_id, job in scheduled_jobs.items(): job_success = True duration = 0.0 - error_message = "" + error_messages: List[str] = [] try: returned_results = job.result_handle.get_result() @@ -412,7 +425,9 @@ def poll_running_jobs(db_conn, db_cursor): ) else: job_success = False - error_message += f"task {task_result.task_id}: {task_result.error_message}\n" + error_messages.append( + f"task {task_result.task_id}: {task_result.error_message}" + ) logger.error( f"Compression task job-{job_id}-task-{task_result.task_id} failed with" f" error: {task_result.error_message}." @@ -434,12 +449,30 @@ def poll_running_jobs(db_conn, db_cursor): ) else: logger.error(f"Job {job_id} failed. See worker logs or status_msg for details.") + + error_log_relative_path = _write_user_failure_log( + title="Compression task errors.", + content=error_messages, + logs_directory=logs_directory, + job_id=job_id, + filename_suffix="task_errors", + ) + if error_log_relative_path is None: + err_msg = "Failed to write user log for failed compression job." + raise RuntimeError(err_msg) + + error_msg = ( + "One or more compression tasks failed." + f" See the error log at '{error_log_relative_path}' inside your configured logs" + " directory (`logs_directory`) for more details." + ) + update_compression_job_metadata( db_cursor, job_id, dict( status=CompressionJobStatus.FAILED, - status_msg=error_message, + status_msg=error_msg, ), ) db_conn.commit() @@ -515,7 +548,7 @@ def main(argv): clp_metadata_db_connection_config, task_manager, ) - poll_running_jobs(db_conn, db_cursor) + poll_running_jobs(clp_config.logs_directory, db_conn, db_cursor) time.sleep(clp_config.compression_scheduler.jobs_poll_delay) except KeyboardInterrupt: logger.info("Forcefully shutting down")