From 4d5f30860815f7141881db6ce1995fb091cbddfe Mon Sep 17 00:00:00 2001 From: Vlad Stirbu Date: Wed, 21 Jan 2026 17:05:00 +0200 Subject: [PATCH 1/7] feat: enhance job specification and cleanup procedure --- src/q8s/execution.py | 61 +++++++++++++++++++++++++++++--------------- 1 file changed, 40 insertions(+), 21 deletions(-) diff --git a/src/q8s/execution.py b/src/q8s/execution.py index eff29cb..0603270 100644 --- a/src/q8s/execution.py +++ b/src/q8s/execution.py @@ -191,7 +191,13 @@ def __create_job_object_from_workload(self, workload: Workload) -> client.V1Job: ) # Create the specification of deployment - spec = client.V1JobSpec(template=template) # , ttl_seconds_after_finished=10 + spec = client.V1JobSpec( + template=template, + # cleanup job and associate resources + ttl_seconds_after_finished=10, + # do not retry failed jobs + backoff_limit=0, + ) # Instantiate the job object job_spec = client.V1Job( @@ -212,12 +218,11 @@ def __create_job_object_from_workload(self, workload: Workload) -> client.V1Job: self.__create_config_map_object_from_workload(job, workload=workload) self.__progress.console.print("Application code attached") - self.__create_environment_secret() + self.__create_environment_secret(job) self.__progress.console.print("Environment variables created") if self.registry_pat: - self.__create_registry_credentials_secret() - + self.__create_registry_credentials_secret(job) self.__progress.advance(prepare_task, 1) return job @@ -252,7 +257,7 @@ def __create_config_map_object_from_workload( namespace=self.namespace, body=configmap ) - def __create_environment_secret(self): + def __create_environment_secret(self, job: client.V1Job): """ Create a Secret object with the environment variables. """ @@ -270,16 +275,16 @@ def __create_environment_secret(self): metadata=client.V1ObjectMeta( name=self.name, namespace=self.namespace, - # owner_references=[ - # client.V1OwnerReference( - # api_version="v1", - # kind="Job", - # name=job.metadata.name, - # uid=job.metadata.uid, - # # block_owner_deletion=True, - # # controller=True, - # ) - # ], + owner_references=[ + client.V1OwnerReference( + api_version="batch/v1", + kind="Job", + name=job.metadata.name, + uid=job.metadata.uid, + # block_owner_deletion=True, + # controller=True, + ) + ], ), ) @@ -287,7 +292,7 @@ def __create_environment_secret(self): namespace=self.namespace, body=secret ) - def __create_registry_credentials_secret(self): + def __create_registry_credentials_secret(self, job: client.V1Job): """ Create a Secret object with the registry credentials. """ @@ -314,6 +319,16 @@ def __create_registry_credentials_secret(self): metadata=client.V1ObjectMeta( name=self.__registry_credentials_secret_name(), namespace=self.namespace, + owner_references=[ + client.V1OwnerReference( + api_version="batch/v1", + kind="Job", + name=job.metadata.name, + uid=job.metadata.uid, + # block_owner_deletion=True, + # controller=True, + ) + ], ), data={ ".dockerconfigjson": base64.b64encode( @@ -411,20 +426,24 @@ def __complete_and_get_job_status(self): # Job execution completed if event["object"].status.active is None: + print(event["object"].status.conditions) + # Failed - if event["object"].status.conditions is None: + if event["object"].status.conditions[-1].type == "Failed": message = "Failed" color = "red" w.stop() result = "stderr" # Succeeded - else: + elif event["object"].status.conditions[-1].type == "Complete": message = event["object"].status.conditions[-1].type color = "green" + w.stop() + + else: + pass - if event["object"].status.conditions[-1].type == "Complete": - w.stop() # Job schedukled elif event["type"] == "ADDED": message = "Scheduled" @@ -489,7 +508,7 @@ def execute_workload(self, workload: Workload) -> tuple[str, str]: except Exception: return "An error occurred.", "stderr" finally: - self.__delete_job() + # self.__delete_job() pass def abort(self): From ecbcaf29df0245593c2bc236c430b25c24837218 Mon Sep 17 00:00:00 2001 From: Vlad Stirbu Date: Wed, 21 Jan 2026 17:44:35 +0200 Subject: [PATCH 2/7] feat: add job submission option to execute command --- src/q8s/cli.py | 7 ++++++- src/q8s/execution.py | 36 ++++++++++++++++++++++++------------ 2 files changed, 30 insertions(+), 13 deletions(-) diff --git a/src/q8s/cli.py b/src/q8s/cli.py index 0f06f43..283e599 100644 --- a/src/q8s/cli.py +++ b/src/q8s/cli.py @@ -115,6 +115,9 @@ def execute( envvar="REGISTRY_PAT", ), ] = None, + submit: Annotated[ + bool, typer.Option(help="Submit job and exit without waiting for completion") + ] = False, args: Annotated[list[str], typer.Argument(help="Additional arguments")] = None, ): project = Project() @@ -147,7 +150,9 @@ def execute( workload = Workload.from_entry_script(entry_script=file) workload.set_args(args or []) - output, stream_name = k8s_context.execute_workload(workload=workload) + output, stream_name = k8s_context.execute_workload( + workload=workload, submit=submit + ) print(f"output:\n{output}") print(f"output stream: {stream_name}") diff --git a/src/q8s/execution.py b/src/q8s/execution.py index 0603270..f08d367 100644 --- a/src/q8s/execution.py +++ b/src/q8s/execution.py @@ -10,7 +10,7 @@ from dxf import DXF from dxf.exceptions import DXFUnauthorizedError from kubernetes import client, config, watch -from rich.progress import Progress +from rich.progress import Progress, TaskID from q8s.enums import Target from q8s.plugins.cpu_job import CPUJobTemplatePlugin @@ -407,14 +407,12 @@ def __get_pods_in_job(self): return pod_name - def __complete_and_get_job_status(self): + def __complete_and_get_job_status(self, execute_task: TaskID): """ Wait for the job to complete and get its status. """ result = "stdout" - execute_task = self.__progress.add_task("[cyan]Executing job...", total=1) - w = watch.Watch() for event in w.stream( @@ -485,7 +483,9 @@ def __prepare_environment(self): return env - def execute_workload(self, workload: Workload) -> tuple[str, str]: + def execute_workload( + self, workload: Workload, submit: bool = False + ) -> tuple[str, str]: """ Execute the given workload. """ @@ -493,16 +493,28 @@ def execute_workload(self, workload: Workload) -> tuple[str, str]: try: self.__create_job_object_from_workload(workload=workload) - if self.jupyter_logger is not None: - self.jupyter_logger(f"Job {self.name} created") + execute_task = self.__progress.add_task("[cyan]Executing job...", total=1) + + if submit is False: + if self.jupyter_logger is not None: + self.jupyter_logger(f"Job {self.name} created") + + stream = self.__complete_and_get_job_status(execute_task=execute_task) - stream = self.__complete_and_get_job_status() + job = self.__get_pods_in_job() + logs = self.__get_job_logs(job) + self.__progress.console.print("Fetched job logs") + + return logs, stream + else: + self.__progress.update( + execute_task, + description=f"[cyan]Executing job... [orange3]Queued {self.name}", + ) - job = self.__get_pods_in_job() - logs = self.__get_job_logs(job) - self.__progress.console.print("Fetched job logs") + self.__progress.advance(execute_task, 1) - return logs, stream + return f"Job {self.name} submitted successfully.", "stdout" except KeyboardInterrupt: return "Task interrupted by user", "stderr" except Exception: From d455e80e639ece82942e24de981e306b02ff0cd2 Mon Sep 17 00:00:00 2001 From: Vlad Stirbu Date: Thu, 22 Jan 2026 19:32:58 +0200 Subject: [PATCH 3/7] fix: handle job deletion on user interruption --- src/q8s/execution.py | 1 + 1 file changed, 1 insertion(+) diff --git a/src/q8s/execution.py b/src/q8s/execution.py index f08d367..f8b58cc 100644 --- a/src/q8s/execution.py +++ b/src/q8s/execution.py @@ -516,6 +516,7 @@ def execute_workload( return f"Job {self.name} submitted successfully.", "stdout" except KeyboardInterrupt: + self.__delete_job() return "Task interrupted by user", "stderr" except Exception: return "An error occurred.", "stderr" From 36c579fbae08199a9248be6ea4a24dc8f0a22ede Mon Sep 17 00:00:00 2001 From: Vlad Stirbu Date: Mon, 26 Jan 2026 10:43:10 +0200 Subject: [PATCH 4/7] fix: correct main app invocation in cli.py and enhance test coverage for K8sContext --- src/q8s/cli.py | 4 +-- tests/test_execution.py | 78 +++++++++++++++++++++++++++++++++++++++-- 2 files changed, 78 insertions(+), 4 deletions(-) diff --git a/src/q8s/cli.py b/src/q8s/cli.py index 283e599..164563b 100644 --- a/src/q8s/cli.py +++ b/src/q8s/cli.py @@ -208,5 +208,5 @@ def jupyter( raise typer.Exit(code=1) -# if __name__ == "__main__": -app() +if __name__ == "__main__": + app() diff --git a/tests/test_execution.py b/tests/test_execution.py index 9f8cfa4..3783fc1 100644 --- a/tests/test_execution.py +++ b/tests/test_execution.py @@ -1,10 +1,11 @@ import unittest -from unittest.mock import patch +from unittest.mock import Mock, patch import requests from dxf.exceptions import DXFUnauthorizedError -from q8s.execution import ContainerImageValidator +from q8s.execution import ContainerImageValidator, K8sContext +from q8s.workload import Workload class MockResponse: @@ -99,5 +100,78 @@ def test_validate_invalid_reference_raises(self, MockDXF): self.assertIn("Invalid container image reference", str(ctx.exception)) +class TestK8sContextExecuteWorkload(unittest.TestCase): + def _make_context(self): + ctx = K8sContext.__new__(K8sContext) + progress = Mock() + progress.console = Mock() + progress.add_task.return_value = "task-id" + ctx._K8sContext__progress = progress + ctx.jupyter_logger = Mock() + ctx.name = "qubernetes-job-test" + return ctx, progress + + def test_execute_workload_submit_false_returns_logs(self): + ctx, progress = self._make_context() + ctx._K8sContext__create_job_object_from_workload = Mock() + ctx._K8sContext__complete_and_get_job_status = Mock(return_value="stdout") + ctx._K8sContext__get_pods_in_job = Mock(return_value="pod-1") + ctx._K8sContext__get_job_logs = Mock(return_value="log output") + + workload = Workload.from_code("print('hi')") + result = ctx.execute_workload(workload, submit=False) + + self.assertEqual(result, ("log output", "stdout")) + ctx._K8sContext__create_job_object_from_workload.assert_called_once_with( + workload=workload + ) + ctx._K8sContext__complete_and_get_job_status.assert_called_once_with( + execute_task="task-id" + ) + ctx._K8sContext__get_job_logs.assert_called_once_with("pod-1") + progress.console.print.assert_called_once_with("Fetched job logs") + + def test_execute_workload_submit_true_skips_wait(self): + ctx, progress = self._make_context() + ctx._K8sContext__create_job_object_from_workload = Mock() + ctx._K8sContext__complete_and_get_job_status = Mock() + + workload = Workload.from_code("print('hi')") + result = ctx.execute_workload(workload, submit=True) + + self.assertEqual( + result, ("Job qubernetes-job-test submitted successfully.", "stdout") + ) + ctx._K8sContext__complete_and_get_job_status.assert_not_called() + progress.update.assert_called_once() + progress.advance.assert_called_once_with("task-id", 1) + + def test_execute_workload_keyboard_interrupt_aborts(self): + ctx, _ = self._make_context() + ctx._K8sContext__create_job_object_from_workload = Mock( + side_effect=KeyboardInterrupt + ) + ctx.abort = Mock() + + workload = Workload.from_code("print('hi')") + result = ctx.execute_workload(workload, submit=False) + + self.assertEqual(result, ("Task interrupted by user", "stderr")) + ctx.abort.assert_called_once() + + def test_execute_workload_exception_returns_error(self): + ctx, _ = self._make_context() + ctx._K8sContext__create_job_object_from_workload = Mock( + side_effect=Exception("boom") + ) + ctx.abort = Mock() + + workload = Workload.from_code("print('hi')") + result = ctx.execute_workload(workload, submit=False) + + self.assertEqual(result, ("An error occurred.", "stderr")) + ctx.abort.assert_not_called() + + if __name__ == "__main__": unittest.main() From ad0345256a0d99ca3504693029565d4bfab1bab7 Mon Sep 17 00:00:00 2001 From: Vlad Stirbu Date: Mon, 26 Jan 2026 10:52:01 +0200 Subject: [PATCH 5/7] fix: replace job deletion with abort handling on user interruption --- src/q8s/execution.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/q8s/execution.py b/src/q8s/execution.py index f8b58cc..108c22e 100644 --- a/src/q8s/execution.py +++ b/src/q8s/execution.py @@ -516,7 +516,7 @@ def execute_workload( return f"Job {self.name} submitted successfully.", "stdout" except KeyboardInterrupt: - self.__delete_job() + self.abort() return "Task interrupted by user", "stderr" except Exception: return "An error occurred.", "stderr" From 0a79e04b23a948327e507c9d7e4094de4cdd6c48 Mon Sep 17 00:00:00 2001 From: Vlad Stirbu Date: Mon, 26 Jan 2026 10:59:56 +0200 Subject: [PATCH 6/7] feat: add tests for job status completion handling in K8sContext --- tests/test_execution.py | 89 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 89 insertions(+) diff --git a/tests/test_execution.py b/tests/test_execution.py index 3783fc1..85a6e0d 100644 --- a/tests/test_execution.py +++ b/tests/test_execution.py @@ -1,4 +1,5 @@ import unittest +from types import SimpleNamespace from unittest.mock import Mock, patch import requests @@ -13,6 +14,21 @@ def __init__(self, status_code): self.status_code = status_code +class DummyWatch: + def __init__(self, events): + self.events = events + self.stopped = False + + def stream(self, *args, **kwargs): + for event in self.events: + if self.stopped: + break + yield event + + def stop(self): + self.stopped = True + + class TestContainerImageValidator(unittest.TestCase): @patch("q8s.execution.DXF") @@ -173,5 +189,78 @@ def test_execute_workload_exception_returns_error(self): ctx.abort.assert_not_called() +class TestK8sContextCompleteJobStatus(unittest.TestCase): + def _make_context(self): + ctx = K8sContext.__new__(K8sContext) + progress = Mock() + progress.console = Mock() + ctx._K8sContext__progress = progress + ctx.jupyter_logger = None + ctx.name = "qubernetes-job-test" + ctx.namespace = "default" + ctx.batch_api_instance = Mock() + return ctx, progress + + def _make_event(self, name, active, condition_type, event_type="MODIFIED"): + conditions = [] + if condition_type is not None: + conditions = [SimpleNamespace(type=condition_type)] + obj = SimpleNamespace( + metadata=SimpleNamespace(name=name), + status=SimpleNamespace(active=active, conditions=conditions), + ) + return {"object": obj, "type": event_type} + + def test_complete_and_get_job_status_failed_returns_stderr(self): + ctx, progress = self._make_context() + event = self._make_event(ctx.name, None, "Failed") + dummy_watch = DummyWatch([event]) + + with patch("q8s.execution.watch.Watch", return_value=dummy_watch): + result = ctx._K8sContext__complete_and_get_job_status( + execute_task="task-id" + ) + + self.assertEqual(result, "stderr") + self.assertTrue(dummy_watch.stopped) + progress.update.assert_called_once() + self.assertIn("Failed", progress.update.call_args.kwargs["description"]) + progress.advance.assert_called_once_with("task-id", 1) + + def test_complete_and_get_job_status_complete_returns_stdout(self): + ctx, progress = self._make_context() + event = self._make_event(ctx.name, None, "Complete") + dummy_watch = DummyWatch([event]) + + with patch("q8s.execution.watch.Watch", return_value=dummy_watch): + result = ctx._K8sContext__complete_and_get_job_status( + execute_task="task-id" + ) + + self.assertEqual(result, "stdout") + self.assertTrue(dummy_watch.stopped) + progress.update.assert_called_once() + self.assertIn("Complete", progress.update.call_args.kwargs["description"]) + progress.advance.assert_called_once_with("task-id", 1) + + def test_complete_and_get_job_status_running_keeps_stdout(self): + ctx, progress = self._make_context() + events = [ + self._make_event(ctx.name, 1, None, event_type="ADDED"), + self._make_event(ctx.name, 1, None, event_type="MODIFIED"), + ] + dummy_watch = DummyWatch(events) + + with patch("q8s.execution.watch.Watch", return_value=dummy_watch): + result = ctx._K8sContext__complete_and_get_job_status( + execute_task="task-id" + ) + + self.assertEqual(result, "stdout") + self.assertFalse(dummy_watch.stopped) + self.assertEqual(progress.update.call_count, 2) + progress.advance.assert_called_once_with("task-id", 1) + + if __name__ == "__main__": unittest.main() From f900ffd3345ebeaebab7214489601e5e15b82f54 Mon Sep 17 00:00:00 2001 From: Vlad Stirbu Date: Mon, 26 Jan 2026 11:33:11 +0200 Subject: [PATCH 7/7] fix: update project version to 0.13.0 in pyproject.toml --- pyproject.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pyproject.toml b/pyproject.toml index 3f02a80..a0d4867 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "q8s" -version = "0.12.0" +version = "0.13.0" description = "Kernel extension for executing quantum programs in simulators on q8s clusters" authors = [{ name = "Vlad Stirbu", email = "vstirbu@gmail.com" }] readme = "README.md"