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: 20 additions & 2 deletions slicer_cli_web/girder_worker_plugin/direct_docker_run.py
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,20 @@ def resolve(arg, **kwargs):
return extra_volumes


def _cancel_latched(task):
return getattr(task.request, '_slicer_cli_web_canceled', False)


class DirectDockerTask(DockerTask):
@property
def canceled(self):
# Latch the first observed cancel: the base property re-inspects the
# broker on every read, so a later inspection that times out to False
# must not un-cancel a task the docker loop already stopped.
if not _cancel_latched(self):
self.request._slicer_cli_web_canceled = super().canceled
return self.request._slicer_cli_web_canceled

def __call__(self, *args, **kwargs):
extra_volumes = _resolve_direct_file_paths(args, kwargs)
if extra_volumes:
Expand All @@ -104,7 +117,7 @@ def __call__(self, *args, **kwargs):
for extra_volume in extra_volumes:
volumes.update(extra_volume._repr_json_())

super().__call__(*args, **kwargs)
return super().__call__(*args, **kwargs)


def _has_image(image):
Expand Down Expand Up @@ -140,4 +153,9 @@ def run(task, **kwargs):
output=CLIProgressCLIWriter(task.job_manager)
))

return _docker_run(task, **kwargs)
results = _docker_run(task, **kwargs)
# Drop a canceled run's results so the upload hooks are skipped: the stopped
# container never wrote its outputs, and uploading nothing would fail the
# job as an error instead of a cancel. Read the latch, not the broker, to
# keep successful runs free of an extra round-trip.
return () if _cancel_latched(task) else results
40 changes: 40 additions & 0 deletions tests/girder_worker_plugin/test_direct_docker_run.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
from os.path import basename
from unittest import mock

import pytest

Expand Down Expand Up @@ -54,3 +55,42 @@ def test_direct_docker_run(mocker, server, adminToken, file):
assert kwargs['container_args'] == [target_path]
# volumes
assert len(kwargs['volumes']) == 2


@pytest.mark.plugin('slicer_cli_web')
@pytest.mark.parametrize('canceled', [True, False])
def test_direct_docker_run_canceled_skips_result_hooks(mocker, canceled):
docker_run_mock = mocker.patch(
'slicer_cli_web.girder_worker_plugin.direct_docker_run._docker_run')
docker_run_mock.return_value = (None, )

hook = mock.Mock()
run.push_request(girder_result_hooks=[hook])
try:
if canceled:
# stand in for the docker loop having latched the cancel mid-run
run.request._slicer_cli_web_canceled = True
run(image='test', container_args=[])
finally:
run.pop_request()

docker_run_mock.assert_called_once()
# a canceled run must not upload the outputs its stopped container skipped
if canceled:
hook.transform.assert_not_called()
else:
hook.transform.assert_called_once()


@pytest.mark.plugin('slicer_cli_web')
def test_direct_docker_run_cancel_is_latched(mocker):
# once the broker reports the revocation, a later inspection that times out
# to False must not un-cancel the run.
mocker.patch('girder_worker.task.is_revoked', side_effect=[True, False, False])

run.push_request()
try:
assert run.canceled
assert run.canceled
finally:
run.pop_request()