Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

[App] Improve cluster creation / deletion experience #15458

Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
32 commits
Select commit Hold shift + click to select a range
d3215ad
Wait by default and notify on state changes
luca3rd Nov 1, 2022
14670e2
tt
nmiculinic Nov 4, 2022
c93525c
Further update
nmiculinic Nov 4, 2022
9b8db79
updated docs
nmiculinic Nov 4, 2022
fe90dd6
Improve user messages and test output
luca3rd Nov 5, 2022
2e1e272
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Nov 5, 2022
2969668
fixup! Improve user messages and test output
luca3rd Nov 5, 2022
5d8bf34
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Nov 5, 2022
a855134
fixup! fixup! Improve user messages and test output
luca3rd Nov 5, 2022
f469b8d
fix mypy and flake8
luca3rd Nov 5, 2022
9404c32
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Nov 5, 2022
7bd79b7
PR feedback
luca3rd Nov 7, 2022
5d9b2f0
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Nov 7, 2022
3331cc1
format doctest
luca3rd Nov 7, 2022
36a966c
Fix tests in other files
luca3rd Nov 7, 2022
d71da41
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Nov 7, 2022
d6fabcf
Merge branch 'master' into ENG-1513-for-cluster-creation-creation-def…
luca3rd Nov 7, 2022
7a2ebba
Fixup param name
luca3rd Nov 7, 2022
cdab084
Address some feedback
luca3rd Nov 9, 2022
d15c2ba
Merge branch 'master' into ENG-1513-for-cluster-creation-creation-def…
nicolai86 Nov 9, 2022
5c12f39
Remove frivolous tests
luca3rd Nov 9, 2022
d1ab5cc
Merge branch 'master' into ENG-1513-for-cluster-creation-creation-def…
nicolai86 Nov 10, 2022
b36de6b
Merge branch 'master' into ENG-1513-for-cluster-creation-creation-def…
luca3rd Nov 10, 2022
09a9920
Merge branch 'master' into ENG-1513-for-cluster-creation-creation-def…
luca3rd Nov 17, 2022
51a9732
Fixup changelog
luca3rd Nov 17, 2022
01bcba3
Merge branch 'master' into ENG-1513-for-cluster-creation-creation-def…
awaelchli Nov 21, 2022
e6cec27
Use StrEnum from lightning_utilities
luca3rd Nov 21, 2022
355b6a1
[pre-commit.ci] auto fixes from pre-commit.com hooks
pre-commit-ci[bot] Nov 21, 2022
c390d05
Merge branch 'master' into ENG-1513-for-cluster-creation-creation-def…
Borda Nov 22, 2022
138bf8a
Merge branch 'master' into ENG-1513-for-cluster-creation-creation-def…
luca3rd Nov 28, 2022
42a6b4a
Rename failed to error
luca3rd Nov 28, 2022
1314e74
Merge branch 'master' into ENG-1513-for-cluster-creation-creation-def…
luca3rd Nov 28, 2022
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: 2 additions & 13 deletions docs/source-app/workflows/byoc/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ Parameters
^^^^^^^^^^

+------------------------+----------------------------------------------------------------------------------------------------+
|Parameter | Descritption |
|Parameter | Description |
+========================+====================================================================================================+
| provider | The cloud provider where your cluster is located. |
| | |
Expand All @@ -78,18 +78,7 @@ Parameters
+------------------------+----------------------------------------------------------------------------------------------------+
| region | AWS region containing compute resources |
+------------------------+----------------------------------------------------------------------------------------------------+
| instance-types | Instance types that you want to support, for computer jobs within the cluster. |
| | |
| | For now, this is the AWS instance types supported by the cluster. |
+------------------------+----------------------------------------------------------------------------------------------------+
| enable-performance | Specifies if the cluster uses cost savings mode. |
| | |
| | In cost saving mode the number of compute nodes is reduced to one, reducing the cost for clusters |
| | with low utilization. |
+------------------------+----------------------------------------------------------------------------------------------------+
| edit-before-creation | Enables interactive editing of requests before submitting it to Lightning AI. |
+------------------------+----------------------------------------------------------------------------------------------------+
| wait | Waits for the cluster to be in a RUNNING state. Only use this for debugging. |
luca3rd marked this conversation as resolved.
Show resolved Hide resolved
| async | Cluster creation will happen in the background. |
+------------------------+----------------------------------------------------------------------------------------------------+

----
Expand Down
1 change: 1 addition & 0 deletions src/lightning_app/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ The format is based on [Keep a Changelog](http://keepachangelog.com/en/1.0.0/).
### Changed

- The `MultiNode` components now warn the user when running with `num_nodes > 1` locally ([#15806](https://github.com/Lightning-AI/lightning/pull/15806))
- Cluster creation and deletion now waits by default [#15458](https://github.com/Lightning-AI/lightning/pull/15458)


### Deprecated
Expand Down
206 changes: 173 additions & 33 deletions src/lightning_app/cli/cmd_clusters.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,9 +3,10 @@
import time
from datetime import datetime
from textwrap import dedent
from typing import Any, List
from typing import Any, List, Union

import click
import lightning_cloud
from lightning_cloud.openapi import (
Externalv1Cluster,
V1AWSClusterDriverSpec,
Expand All @@ -15,8 +16,10 @@
V1ClusterState,
V1ClusterType,
V1CreateClusterRequest,
V1GetClusterResponse,
V1KubernetesClusterDriver,
)
from lightning_utilities.core.enums import StrEnum
from rich.console import Console
from rich.table import Table
from rich.text import Text
Expand All @@ -25,10 +28,26 @@
from lightning_app.utilities.network import LightningClient
from lightning_app.utilities.openapi import create_openapi_object, string2dict

CLUSTER_STATE_CHECKING_TIMEOUT = 60
MAX_CLUSTER_WAIT_TIME = 5400


class ClusterState(StrEnum):
UNSPECIFIED = "unspecified"
QUEUED = "queued"
PENDING = "pending"
RUNNING = "running"
FAILED = "error"
DELETED = "deleted"

def __str__(self) -> str:
return str(self.value)

@classmethod
def from_api(cls, status: V1ClusterState) -> "ClusterState":
parsed = str(status).lower().split("_", maxsplit=2)[-1]
luca3rd marked this conversation as resolved.
Show resolved Hide resolved
return cls.from_str(parsed)


class ClusterList(Formatable):
def __init__(self, clusters: List[Externalv1Cluster]):
self.clusters = clusters
Expand Down Expand Up @@ -86,7 +105,7 @@ def create(
region: str = "us-east-1",
external_id: str = None,
edit_before_creation: bool = False,
wait: bool = False,
do_async: bool = False,
luca3rd marked this conversation as resolved.
Show resolved Hide resolved
) -> None:
"""request Lightning AI BYOC compute cluster creation.

Expand All @@ -97,7 +116,7 @@ def create(
region: AWS region containing compute resources
external_id: AWS IAM Role external ID
edit_before_creation: Enables interactive editing of requests before submitting it to Lightning AI.
wait: Waits for the cluster to be in a RUNNING state. Only use this for debugging.
do_async: Triggers cluster creation in the background and exits
"""
performance_profile = V1ClusterPerformanceProfile.DEFAULT
if cost_savings:
Expand Down Expand Up @@ -130,22 +149,31 @@ def create(
click.echo("cluster unchanged")

resp = self.api_client.cluster_service_create_cluster(body=new_body)
if wait:
_wait_for_cluster_state(self.api_client, resp.id, V1ClusterState.RUNNING)

click.echo(
dedent(
f"""\
{resp.id} is now being created... This can take up to an hour.
BYOC cluster creation triggered successfully!
This can take up to an hour to complete.

To view the status of your clusters use:
`lightning list clusters`
lightning list clusters

To view cluster logs use:
`lightning show cluster logs {resp.id}`
"""
lightning show cluster logs {cluster_name}

To delete the cluster run:
lightning delete cluster {cluster_name}
"""
)
)
background_message = "\nCluster will be created in the background!"
if do_async:
click.echo(background_message)
else:
try:
_wait_for_cluster_state(self.api_client, resp.id, V1ClusterState.RUNNING)
except KeyboardInterrupt:
click.echo(background_message)

def get_clusters(self) -> ClusterList:
resp = self.api_client.cluster_service_list_clusters(phase_not_in=[V1ClusterState.DELETED])
Expand All @@ -156,7 +184,7 @@ def list(self) -> None:
console = Console()
console.print(clusters.as_table())

def delete(self, cluster_id: str, force: bool = False, wait: bool = False) -> None:
def delete(self, cluster_id: str, force: bool = False, do_async: bool = False) -> None:
if force:
click.echo(
"""
Expand All @@ -167,47 +195,86 @@ def delete(self, cluster_id: str, force: bool = False, wait: bool = False) -> No
)
click.confirm("Do you want to continue?", abort=True)

resp: V1GetClusterResponse = self.api_client.cluster_service_get_cluster(id=cluster_id)
bucket_name = resp.spec.driver.kubernetes.aws.bucket_name

self.api_client.cluster_service_delete_cluster(id=cluster_id, force=force)
click.echo("Cluster deletion triggered successfully")
click.echo(
dedent(
f"""\
Cluster deletion triggered successfully

For safety purposes we will not delete anything in the S3 bucket associated with the cluster:
{bucket_name}

if wait:
_wait_for_cluster_state(self.api_client, cluster_id, V1ClusterState.DELETED)
You may want to delete it manually using the AWS CLI:
aws s3 rb --force s3://{bucket_name}
"""
)
)

background_message = "\nCluster will be deleted in the background!"
if do_async:
click.echo(background_message)
else:
try:
_wait_for_cluster_state(self.api_client, cluster_id, V1ClusterState.DELETED)
except KeyboardInterrupt:
click.echo(background_message)


def _wait_for_cluster_state(
api_client: LightningClient,
cluster_id: str,
target_state: V1ClusterState,
max_wait_time: int = MAX_CLUSTER_WAIT_TIME,
check_timeout: int = CLUSTER_STATE_CHECKING_TIMEOUT,
timeout_seconds: int = MAX_CLUSTER_WAIT_TIME,
poll_duration_seconds: int = 10,
) -> None:
"""_wait_for_cluster_state waits until the provided cluster has reached a desired state, or failed.

Messages will be displayed to the user as the cluster changes state.
We poll the API server for any changes

Args:
api_client: LightningClient used for polling
cluster_id: Specifies the cluster to wait for
target_state: Specifies the desired state the target cluster needs to meet
max_wait_time: Maximum duration to wait (in seconds)
check_timeout: duration between polling for the cluster state (in seconds)
timeout_seconds: Maximum duration to wait
poll_duration_seconds: duration between polling for the cluster state
"""
start = time.time()
elapsed = 0
while elapsed < max_wait_time:
cluster_resp = api_client.cluster_service_list_clusters()
new_cluster = None
for clust in cluster_resp.clusters:
if clust.id == cluster_id:
new_cluster = clust
break
if new_cluster is not None:
if new_cluster.status.phase == target_state:

click.echo(f"Waiting for cluster to be {ClusterState.from_api(target_state)}...")
while elapsed < timeout_seconds:
try:
resp: V1GetClusterResponse = api_client.cluster_service_get_cluster(id=cluster_id)
click.echo(_cluster_status_long(cluster=resp, desired_state=target_state, elapsed=elapsed))
if resp.status.phase == target_state:
break
elif new_cluster.status.phase == V1ClusterState.FAILED:
raise click.ClickException(f"Cluster {cluster_id} is in failed state.")
time.sleep(check_timeout)
elapsed = int(time.time() - start)
time.sleep(poll_duration_seconds)
elapsed = int(time.time() - start)
except lightning_cloud.openapi.rest.ApiException as e:
if e.status == 404 and target_state == V1ClusterState.DELETED:
return
raise
else:
raise click.ClickException("Max wait time elapsed")
state_str = ClusterState.from_api(target_state)
raise click.ClickException(
dedent(
f"""\
The cluster has not entered the {state_str} state within {_format_elapsed_seconds(timeout_seconds)}.

The cluster may eventually be {state_str} afterwards, please check its status using:
lighting list clusters

To view cluster logs use:
lightning show cluster logs {cluster_id}

Contact [email protected] for additional help
"""
)
)


def _check_cluster_name_is_valid(_ctx: Any, _param: Any, value: str) -> str:
Expand All @@ -219,3 +286,76 @@ def _check_cluster_name_is_valid(_ctx: Any, _param: Any, value: str) -> str:
Provide a cluster name using valid characters and try again."""
)
return value


def _cluster_status_long(cluster: V1GetClusterResponse, desired_state: V1ClusterState, elapsed: float) -> str:
"""Echos a long-form status message to the user about the cluster state.

Args:
cluster: The cluster object
elapsed: Seconds since we've started polling
"""

cluster_name = cluster.name
current_state = cluster.status.phase
current_reason = cluster.status.reason
bucket_name = cluster.spec.driver.kubernetes.aws.bucket_name

duration = _format_elapsed_seconds(elapsed)

if current_state == V1ClusterState.FAILED:
return dedent(
f"""\
The requested cluster operation for cluster {cluster_name} has errors:
{current_reason}

---
We are automatically retrying, and an automated alert has been created

WARNING: Any non-deleted cluster may be using resources.
To avoid incuring cost on your cloud provider, delete the cluster using the following command:
lightning delete cluster {cluster_name}

Contact [email protected] for additional help
"""
)

if desired_state == current_state == V1ClusterState.RUNNING:
return dedent(
f"""\
Cluster {cluster_name} is now running and ready to use.
To launch an app on this cluster use: lightning run app app.py --cloud --cluster-id {cluster_name}
"""
)

if desired_state == V1ClusterState.RUNNING:
return f"Cluster {cluster_name} is being created [elapsed={duration}]"

if desired_state == current_state == V1ClusterState.DELETED:
return dedent(
f"""\
Cluster {cluster_name} has been successfully deleted, and almost all AWS resources have been removed

For safety purposes we kept the S3 bucket associated with the cluster: {bucket_name}

You may want to delete it manually using the AWS CLI:
aws s3 rb --force s3://{bucket_name}
"""
)

if desired_state == V1ClusterState.DELETED:
return f"Cluster {cluster_name} is being deleted [elapsed={duration}]"

raise click.ClickException(f"Unknown cluster desired state {desired_state}")


def _format_elapsed_seconds(seconds: Union[float, int]) -> str:
"""Turns seconds into a duration string.

>>> _format_elapsed_seconds(5)
'05s'
>>> _format_elapsed_seconds(60)
'01m00s'
"""
minutes, seconds = divmod(seconds, 60)
return (f"{minutes:02}m" if minutes else "") + f"{seconds:02}s"
12 changes: 7 additions & 5 deletions src/lightning_app/cli/lightning_cli_create.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ def create() -> None:
type=bool,
required=False,
default=False,
hidden=True,
is_flag=True,
help=""""Use this flag to ensure that the cluster is created with a profile that is optimized for performance.
This makes runs more expensive but start-up times decrease.""",
Expand All @@ -45,16 +46,17 @@ def create() -> None:
"--edit-before-creation",
default=False,
is_flag=True,
hidden=True,
help="Edit the cluster specs before submitting them to the API server.",
)
@click.option(
"--wait",
"wait",
"--async",
"do_async",
type=bool,
required=False,
default=False,
is_flag=True,
help="Enabling this flag makes the CLI wait until the cluster is running.",
help="This flag makes the CLI return immediately and lets the cluster creation happen in the background.",
)
def create_cluster(
cluster_name: str,
Expand All @@ -64,7 +66,7 @@ def create_cluster(
provider: str,
edit_before_creation: bool,
enable_performance: bool,
wait: bool,
do_async: bool,
**kwargs: Any,
) -> None:
"""Create a Lightning AI BYOC compute cluster with your cloud provider credentials."""
Expand All @@ -79,7 +81,7 @@ def create_cluster(
external_id=external_id,
edit_before_creation=edit_before_creation,
cost_savings=not enable_performance,
wait=wait,
do_async=do_async,
)


Expand Down
Loading