Fix tasks

This commit is contained in:
mvdbeek
2021-11-29 11:13:10 +01:00
parent 7ec0963a7c
commit 64442a90d5
3 changed files with 15 additions and 14 deletions
+5 -5
View File
@@ -122,8 +122,8 @@ class TaskPathRewriter:
return os.path.join(self.working_directory, os.path.basename(job_file_name))
def get_path_rewriter(outputs_to_working_directory, working_directory=None, outputs_directory=None):
if outputs_to_working_directory:
return OutputsToWorkingDirectoryPathRewriter(working_directory, outputs_directory)
else:
return NullDatasetPathRewriter()
def get_path_rewriter(outputs_to_working_directory, working_directory, outputs_directory, is_task):
job_dataset_path_rewriter = OutputsToWorkingDirectoryPathRewriter(working_directory, outputs_directory) if outputs_to_working_directory else NullDatasetPathRewriter()
if is_task:
return TaskPathRewriter(working_directory, job_dataset_path_rewriter=job_dataset_path_rewriter)
return job_dataset_path_rewriter
+6 -2
View File
@@ -33,6 +33,7 @@ class JobIO(Dictifiable):
'check_job_script_integrity',
'check_job_script_integrity_count',
'check_job_script_integrity_sleep',
'is_task',
)
def __init__(
@@ -54,7 +55,8 @@ class JobIO(Dictifiable):
_file_sources,
check_job_script_integrity,
check_job_script_integrity_count,
check_job_script_integrity_sleep):
check_job_script_integrity_sleep,
is_task=False):
self.sa_session = sa_session
self.job = job
self.working_directory = working_directory
@@ -73,6 +75,7 @@ class JobIO(Dictifiable):
self.check_job_script_integrity = check_job_script_integrity
self.check_job_script_integrity_count = check_job_script_integrity_count
self.check_job_script_integrity_sleep = check_job_script_integrity_sleep
self.is_task = is_task
self.output_paths = None
self.output_hdas_and_paths = None
self._dataset_path_rewriter = None
@@ -98,7 +101,8 @@ class JobIO(Dictifiable):
self._dataset_path_rewriter = get_path_rewriter(
outputs_to_working_directory=self.outputs_to_working_directory,
working_directory=self.working_directory,
outputs_directory=self.outputs_directory
outputs_directory=self.outputs_directory,
is_task=self.is_task,
)
return self._dataset_path_rewriter
+4 -7
View File
@@ -922,6 +922,7 @@ class JobWrapper(HasResourceParameters):
Wraps a 'model.Job' with convenience methods for running processes and
state management.
"""
is_task = False
def __init__(self, job, queue: 'JobHandlerQueue', use_persisted_destination=False):
self.job_id = job.id
@@ -944,7 +945,6 @@ class JobWrapper(HasResourceParameters):
self._setup_working_directory(job=job)
# the path rewriter needs destination params, so it cannot be set up until after the destination has been
# resolved
self._dataset_path_rewriter = None
self._job_io = None
self.tool_provided_job_metadata = None
self.job_runner_mapper = JobRunnerMapper(self, queue.dispatcher.url_to_destination, self.app.job_config)
@@ -1006,6 +1006,7 @@ class JobWrapper(HasResourceParameters):
check_job_script_integrity=self.app.config.check_job_script_integrity,
check_job_script_integrity_count=self.app.config.check_job_script_integrity_count,
check_job_script_integrity_sleep=self.app.config.check_job_script_integrity_sleep,
is_task=self.is_task,
)
return self._job_io
@@ -2251,6 +2252,8 @@ class TaskWrapper(JobWrapper):
"""
# Abstract this to be more useful for running tasks that *don't* necessarily compose a job.
is_task = True
def __init__(self, task, queue):
self.task_id = task.id
super().__init__(task.job, queue)
@@ -2260,12 +2263,6 @@ class TaskWrapper(JobWrapper):
self.prepare_input_files_cmds = None
self.status = task.states.NEW
@property
def dataset_path_rewriter(self):
if self._dataset_path_rewriter is None:
self._dataset_path_rewriter = TaskPathRewriter(self.working_directory, super().dataset_path_rewriter)
return self._dataset_path_rewriter
def can_split(self):
# Should the job handler split this job up? TaskWrapper should
# always return False as the job has already been split.