diff --git a/packages/syft-job/src/syft_job/job_runner.py b/packages/syft-job/src/syft_job/job_runner.py index ec4bade4294..427e744823a 100644 --- a/packages/syft-job/src/syft_job/job_runner.py +++ b/packages/syft-job/src/syft_job/job_runner.py @@ -34,6 +34,15 @@ def get_job_timeout_seconds() -> int: ) +def format_timeout(timeout: float) -> str: + """Render a timeout in seconds for humans: whole minutes when exact, else seconds.""" + if timeout >= 60 and timeout % 60 == 0: + minutes = int(timeout // 60) + return f"{minutes} minute" if minutes == 1 else f"{minutes} minutes" + seconds = f"{timeout:g}" + return f"{seconds} second" if seconds == "1" else f"{seconds} seconds" + + IS_IN_JOB_ENV_VAR = "SYFT_IS_IN_JOB" # A job runs from a fresh copy of its approved submission in the system temp @@ -296,7 +305,7 @@ def _execute_job_streaming(self, ref: JobRef, timeout: int, run_dir: Path) -> in _kill_process_tree(process.pid) process.wait() timed_out = True - print(f" Job {job_name} timed out after {timeout // 60} minutes") + print(f" Job {job_name} timed out after {format_timeout(timeout)}") stdout_f.write("\n--- PROCESS TIMED OUT ---\n") stderr_f.write("\n--- PROCESS TIMED OUT ---\n") break @@ -367,7 +376,7 @@ def _execute_job_captured(self, ref: JobRef, timeout: int, run_dir: Path) -> int returncode = -1 stdout = (stdout or "") + "\n--- PROCESS TIMED OUT ---\n" stderr = (stderr or "") + "\n--- PROCESS TIMED OUT ---\n" - print(f" Job {job_name} timed out after {timeout // 60} minutes") + print(f" Job {job_name} timed out after {format_timeout(timeout)}") staging_dir = self.manager.staging_dir(ref) staging_dir.mkdir(parents=True, exist_ok=True) @@ -503,7 +512,7 @@ def _run_approved_job( self._report_result(ref, returncode) except subprocess.TimeoutExpired: - print(f" Job {job_name} timed out after {timeout // 60} minutes") + print(f" Job {job_name} timed out after {format_timeout(timeout)}") self._finalize(ref, -1) except Exception as e: print(f" Error executing job {job_name}: {e}") diff --git a/packages/syft-job/tests/test_format_timeout.py b/packages/syft-job/tests/test_format_timeout.py new file mode 100644 index 00000000000..5e13b6b9612 --- /dev/null +++ b/packages/syft-job/tests/test_format_timeout.py @@ -0,0 +1,20 @@ +import pytest + +from syft_job.job_runner import format_timeout + + +@pytest.mark.parametrize( + "timeout,expected", + [ + (30, "30 seconds"), + (1, "1 second"), + (59, "59 seconds"), + (60, "1 minute"), + (90, "90 seconds"), + (120, "2 minutes"), + (1800, "30 minutes"), + (2.5, "2.5 seconds"), + ], +) +def test_format_timeout(timeout, expected): + assert format_timeout(timeout) == expected diff --git a/packages/syft-job/tests/test_job_flow.py b/packages/syft-job/tests/test_job_flow.py index 4ca5b38320f..259ad7f5a88 100644 --- a/packages/syft-job/tests/test_job_flow.py +++ b/packages/syft-job/tests/test_job_flow.py @@ -237,7 +237,7 @@ def test_submission_validation(tmp_path: Path): """ -def test_timeout_does_not_hang_runner(tmp_path: Path): +def test_timeout_does_not_hang_runner(tmp_path: Path, capsys): syftbox = tmp_path / "SyftBox" syftbox.mkdir() @@ -263,3 +263,4 @@ def test_timeout_does_not_hang_runner(tmp_path: Path): # 3s job timeout + venv setup + tree-kill cleanup should fit well under 60s. assert elapsed < 60, f"process_approved_jobs took {elapsed:.1f}s — likely hung" assert do_client.jobs[0].status == "failed" + assert "timed out after 3 seconds" in capsys.readouterr().out