Skip to content
Closed
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
15 changes: 12 additions & 3 deletions packages/syft-job/src/syft_job/job_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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}")
Expand Down
20 changes: 20 additions & 0 deletions packages/syft-job/tests/test_format_timeout.py
Original file line number Diff line number Diff line change
@@ -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
3 changes: 2 additions & 1 deletion packages/syft-job/tests/test_job_flow.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand All @@ -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
Loading