Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 18 additions & 4 deletions dpdispatcher/submission.py
Original file line number Diff line number Diff line change
Expand Up @@ -326,6 +326,8 @@ def run_submission(
clean: bool | str = True,
check_interval: int = 30,
continue_on_failure: bool | None = None,
raise_on_failure: bool = True,
include_failed_results: bool = False,
) -> dict[str, Any]: # noqa: ANN401
"""Execute the submission and monitor it until completion.

Expand All @@ -347,6 +349,12 @@ def run_submission(
Continue monitoring remaining jobs after retry exhaustion. If omitted,
use the policy stored on this submission (which defaults to ``False``).
An explicit value overrides the persisted policy for this run.
raise_on_failure : bool, default=True
Raise after downloads when continued jobs contain terminal failures.
Set false for callers that need to inspect partial results themselves.
include_failed_results : bool, default=False
Download declared files from failed jobs when available, skipping
missing files so valid partial outputs can be inspected.

Returns
-------
Expand Down Expand Up @@ -443,7 +451,10 @@ def run_submission(
self.handle_unexpected_submission_state(continue_on_failure=True)
else:
self.handle_unexpected_submission_state()
results_downloaded = self.try_download_result()
if include_failed_results:
results_downloaded = self.try_download_result(include_failed=True)
else:
results_downloaded = self.try_download_result()
all_jobs_genuinely_finished = (
all_jobs_genuinely_finished and results_downloaded
)
Expand Down Expand Up @@ -478,7 +489,7 @@ def run_submission(
"preserving remote workdir for debugging at: "
f"{machine.context.remote_root}"
)
if continue_on_failure:
if continue_on_failure and raise_on_failure:
self.raise_for_failed_jobs()
return self.serialize()

Expand Down Expand Up @@ -599,14 +610,17 @@ def try_download_error_info(self) -> None:
f"Could not download error file for job {job.job_hash}: {e}"
)

def try_download_result(self) -> bool:
def try_download_result(self, include_failed: bool = False) -> bool:
"""Download results, retrying transient failures for up to 24 hours."""
start_time = time.time()
retry_interval = 60 # retry every 1 minute
success = False
while not success:
try:
self.download_jobs()
if include_failed:
self.download_jobs(include_failed=True)
else:
self.download_jobs()
success = True
except FileNotFoundError as e:
# retry will never success if the file is not found
Expand Down
37 changes: 37 additions & 0 deletions tests/test_clean_strategy.py
Original file line number Diff line number Diff line change
Expand Up @@ -100,9 +100,46 @@ def test_invalid_strategy_raises_before_upload(self) -> None:
mock_upload.assert_not_called()


class TestFailedResultPolicy(unittest.TestCase):
"""The opt-in run policy forwards failed-result downloads."""

def test_run_submission_requests_failed_results(self):
sub = Submission.__new__(Submission)
sub.belonging_jobs = [MagicMock(job_state=JobStatus.finished)]
sub.belonging_tasks = []
sub.submission_hash = "test_hash"
sub.machine = MagicMock()
sub.resources = MagicMock()
sub.resources.strategy = {"ratio_unfinished": 0.0}
sub.resources.wait_time = 0
sub.try_recover_from_json = MagicMock()
sub.update_submission_state = MagicMock()
sub.check_all_finished = MagicMock(return_value=True)
sub.handle_unexpected_submission_state = MagicMock()
sub.try_download_result = MagicMock(return_value=True)
sub.try_download_error_info = MagicMock()
sub.submission_to_json = MagicMock()
sub.serialize = MagicMock(return_value={})

sub.run_submission(
clean=False,
check_interval=0,
include_failed_results=True,
)

sub.try_download_result.assert_called_once_with(include_failed=True)


class TestDownloadResult(unittest.TestCase):
"""Result-download status must distinguish success from retry exhaustion."""

def test_include_failed_download_requests_failed_tasks(self) -> None:
sub = Submission.__new__(Submission)
sub.download_jobs = MagicMock()

self.assertTrue(sub.try_download_result(include_failed=True))
sub.download_jobs.assert_called_once_with(include_failed=True)

def test_successful_download_returns_true(self) -> None:
sub = Submission.__new__(Submission)
sub.download_jobs = MagicMock()
Expand Down
Loading