Skip to content
Merged
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
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand Down
11 changes: 8 additions & 3 deletions src/q8s/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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}")
Expand Down Expand Up @@ -203,5 +208,5 @@ def jupyter(
raise typer.Exit(code=1)


# if __name__ == "__main__":
app()
if __name__ == "__main__":
app()
98 changes: 65 additions & 33 deletions src/q8s/execution.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand All @@ -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

Expand Down Expand Up @@ -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.
"""
Expand All @@ -270,24 +275,24 @@ 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,
)
],
),
)

self.core_api_instance.create_namespaced_secret(
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.
"""
Expand All @@ -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(
Expand Down Expand Up @@ -392,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(
Expand All @@ -411,20 +424,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"
Expand Down Expand Up @@ -466,30 +483,45 @@ 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.
"""

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)

stream = self.__complete_and_get_job_status()
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)

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:
self.abort()
return "Task interrupted by user", "stderr"
except Exception:
return "An error occurred.", "stderr"
finally:
self.__delete_job()
# self.__delete_job()
pass

def abort(self):
Expand Down
Loading
Loading