I simulated a cluster where the leader succeeded and a worker failed. wait_for_job() still returned success and never checked worker pods:
from unittest import mock
from unittest.mock import MagicMock
from kinetic.backend import pathways_client
leader = MagicMock()
leader.status.phase = "Succeeded"
leader.status.container_statuses = None
worker = MagicMock()
worker.metadata.name = "keras-pathways-job-abc-1"
worker.status.phase = "Failed"
core = MagicMock()
core.read_namespaced_pod.return_value = leader
core.list_namespaced_pod.return_value.items = [leader, worker]
streamer = MagicMock()
streamer.__enter__ = MagicMock(return_value=streamer)
streamer.__exit__ = MagicMock(return_value=False)
with mock.patch("kinetic.backend.k8s_utils.core_v1", return_value=core), \
mock.patch("kinetic.backend.pathways_client.LogStreamer", return_value=streamer):
result = pathways_client.wait_for_job("job-abc")
print(result)
print(core.read_namespaced_pod.call_args_list)
print(core.list_namespaced_pod.called)
Output:
Found pod: keras-pathways-job-abc-0
[REMOTE] Job keras-pathways-job-abc completed successfully
success
[call('keras-pathways-job-abc-0', 'default')]
False
Expected output should have been an error or failed status.
I simulated a cluster where the leader succeeded and a worker failed. wait_for_job() still returned success and never checked worker pods:
Output:
Expected output should have been an error or failed status.