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
8 changes: 6 additions & 2 deletions sagemaker-core/src/sagemaker/core/processing.py
Original file line number Diff line number Diff line change
Expand Up @@ -361,7 +361,9 @@ def run(

if not isinstance(self.sagemaker_session, PipelineSession):
self.jobs.append(self.latest_job)
if wait:
# Local processing jobs run to completion inside _start_new and only exist
# locally, so there is nothing to poll or stream from CloudWatch.
if wait and not isinstance(self.sagemaker_session, LocalSession):
if logs:
logs_for_processing_job(
sagemaker_session=self.sagemaker_session,
Expand Down Expand Up @@ -940,7 +942,9 @@ def run(

if not isinstance(self.sagemaker_session, PipelineSession):
self.jobs.append(self.latest_job)
if wait:
# Local processing jobs run to completion inside _start_new and only exist
# locally, so there is nothing to poll or stream from CloudWatch.
if wait and not isinstance(self.sagemaker_session, LocalSession):
if logs:
logs_for_processing_job(
sagemaker_session=self.sagemaker_session,
Expand Down
67 changes: 67 additions & 0 deletions sagemaker-core/tests/unit/test_processing.py
Original file line number Diff line number Diff line change
Expand Up @@ -1165,6 +1165,73 @@ def test_run_with_wait_and_logs_uses_session(self, mock_session):
if os.path.exists(temp_file):
os.unlink(temp_file)

@pytest.mark.parametrize("logs", [True, False])
def test_run_in_local_mode_does_not_poll_or_stream_logs(self, mock_session, logs):
"""Local processing jobs finish inside _start_new; run(wait=True) must not reach AWS."""
from sagemaker.core.local.local_session import LocalSession

local_session = Mock(spec=LocalSession)
local_session.sagemaker_config = {}
local_session.boto_region_name = "us-west-2"
processor = Processor(
role="arn:aws:iam::123456789012:role/SageMakerRole",
image_uri="test-image:latest",
instance_count=1,
instance_type="local",
sagemaker_session=local_session,
)
mock_job = Mock()
mock_job.processing_job_name = "test-processing-job"

with (
patch.object(processor, "_start_new", return_value=mock_job),
patch("sagemaker.core.processing.logs_for_processing_job") as mock_logs,
patch("sagemaker.core.processing._wait_for_processing_job") as mock_wait,
):
processor.run(wait=True, logs=logs)

mock_logs.assert_not_called()
mock_wait.assert_not_called()

@pytest.mark.parametrize("logs", [True, False])
def test_script_processor_run_in_local_mode_does_not_poll_or_stream_logs(
self, mock_session, logs
):
from sagemaker.core.local.local_session import LocalSession

local_session = Mock(spec=LocalSession)
local_session.sagemaker_config = {}
local_session.boto_region_name = "us-west-2"
local_session.default_bucket_prefix = None
local_session.default_bucket.return_value = "test-bucket"
processor = ScriptProcessor(
role="arn:aws:iam::123456789012:role/SageMakerRole",
image_uri="test-image:latest",
command=["python3"],
instance_count=1,
instance_type="local",
sagemaker_session=local_session,
)
mock_job = Mock()
mock_job.processing_job_name = "test-processing-job"

with tempfile.NamedTemporaryFile(mode="w", delete=False, suffix=".py") as f:
f.write("print('test')")
temp_file = f.name
try:
with (
patch.object(processor, "_start_new", return_value=mock_job),
patch("sagemaker.core.s3.S3Uploader.upload", return_value="s3://bucket/code.py"),
patch("sagemaker.core.processing.logs_for_processing_job") as mock_logs,
patch("sagemaker.core.processing._wait_for_processing_job") as mock_wait,
):
processor.run(code=temp_file, wait=True, logs=logs)
finally:
os.unlink(temp_file)

mock_logs.assert_not_called()
mock_wait.assert_not_called()

def test_wait_for_processing_job_uses_session_region(self, mock_session):
"""_wait_for_processing_job describes via the session's region-aware client."""
mock_session.sagemaker_client.describe_processing_job = Mock(
Expand Down
Loading