Skip to content

Commit 506614c

Browse files
authored
Merge pull request GoogleCloudPlatform#4491 from abbas1902/rmig
Transition DWS Flex-Start to Regional MIGs
2 parents 1e5390c + 3aa5e48 commit 506614c

4 files changed

Lines changed: 65 additions & 51 deletions

File tree

community/modules/scheduler/schedmd-slurm-gcp-v6-controller/modules/cleanup_compute/scripts/cleanup_compute.sh

Lines changed: 11 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ nodeset_name="$3"
2121
universe_domain="$4"
2222
compute_endpoint_version="$5"
2323
gcloud_dir="$6"
24+
MAX_ATTEMPTS=3
2425

2526
if [[ $# -ne 5 ]] && [[ $# -ne 6 ]]; then
2627
echo "Usage: $0 <project> <cluster_name> <nodeset_name> <universe_domain> <compute_endpoint_version> [<gcloud_dir>]"
@@ -44,14 +45,21 @@ tmpfile=$(mktemp) # have to use a temp file, since `< <(gcloud ...)` doesn't wor
4445
trap 'rm -f "$tmpfile"' EXIT
4546

4647
echo "Deleting managed instance groups"
47-
mig_filter="name:${cluster_name}-${nodeset_name}-* AND status!=STOPPING"
48+
mig_filter="name:${cluster_name}-${nodeset_name}-*"
4849
gcloud compute instance-groups managed list --format="value(self_link)" --filter="${mig_filter}" >"$tmpfile"
49-
while batch="$(head -n 2)" && [[ ${#batch} -gt 0 ]]; do
50+
while batch="$(head -n 5)" && [[ ${#batch} -gt 0 ]]; do
5051
groups=$(echo "$batch" | paste -sd " " -) # concat into a single space-separated line
5152
# The lack of quotes around ${groups} is intentional and causes each new space-separated "word" to
5253
# be treated as independent arguments. See PR#2523
5354
# shellcheck disable=SC2086
54-
gcloud compute instance-groups managed delete --quiet ${groups} || echo "Failed to delete some instance groups"
55+
for _ in $( #occasionally MIGs will fail to delete due to some active transformation happening, so let's retry
56+
seq 1 $MAX_ATTEMPTS
57+
); do
58+
if gcloud compute instance-groups managed delete --quiet ${groups}; then
59+
break
60+
fi
61+
echo "MIG deletion failed, retrying"
62+
done
5563
done <"$tmpfile"
5664
true >"$tmpfile" # Wipe contents of tmp file
5765

community/modules/scheduler/schedmd-slurm-gcp-v6-controller/modules/slurm_files/scripts/mig_flex.py

Lines changed: 28 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -45,9 +45,8 @@ def resume_flex_chunk(nodes: List[str], job_id: Optional[int], lkp: util.Lookup)
4545
assert nodes
4646
model = nodes[0]
4747
nodeset = lkp.node_nodeset(model)
48-
zones = nodeset.zone_policy_allow
49-
assert len(zones) == 1
50-
zone = zones[0]
48+
assert len(nodeset.zone_policy_allow) > 0
49+
region = lkp.node_region(model)
5150

5251
assert nodeset.dws_flex.enabled
5352

@@ -58,24 +57,20 @@ def resume_flex_chunk(nodes: List[str], job_id: Optional[int], lkp: util.Lookup)
5857
mig_name = f"{lkp.cfg.slurm_cluster_name}-{nodeset.nodeset_name}-{uid}"
5958

6059
# Create MIG
61-
req = lkp.compute.instanceGroupManagers().insert(
60+
req = lkp.compute.regionInstanceGroupManagers().insert(
6261
project=lkp.project,
63-
zone=zone,
62+
region=region,
6463
body=dict(
6564
name=mig_name,
66-
versions=[dict(
67-
instanceTemplate=nodeset.instance_template)],
65+
versions=[dict(instanceTemplate=nodeset.instance_template)],
6866
targetSize=0,
69-
# TODO(FLEX): uncomment once moved to RMIG
70-
# distributionPolicy=dict(
71-
# zones=[
72-
# dict(zone=f"zones/{z}") for z in nodeset.zone_policy_allow
73-
# ],
74-
# targetShape="ANY_SINGLE_ZONE" ),
75-
#updatePolicy = dict(
76-
# instanceRedistributionType = "NONE" ),
77-
instanceLifecyclePolicy=dict(
78-
defaultActionOnFailure= "DO_NOTHING" ), # TODO(FLEX): Not supported yet, migrate once supported
67+
distributionPolicy=dict(
68+
zones=[
69+
dict(zone=f"zones/{z}") for z in nodeset.zone_policy_allow
70+
],
71+
targetShape="ANY_SINGLE_ZONE" ),
72+
updatePolicy = dict(instanceRedistributionType = "NONE" ),
73+
instanceLifecyclePolicy=dict(defaultActionOnFailure= "DO_NOTHING" ), # TODO(FLEX): Not supported yet, migrate once supported
7974
)
8075
)
8176
util.log_api_request(req)
@@ -84,9 +79,9 @@ def resume_flex_chunk(nodes: List[str], job_id: Optional[int], lkp: util.Lookup)
8479
assert "error" not in res, f"{res}"
8580

8681
# Create resize request
87-
req = lkp.compute.instanceGroupManagerResizeRequests().insert(
82+
req = lkp.compute.regionInstanceGroupManagerResizeRequests().insert(
8883
project=lkp.project,
89-
zone=zone,
84+
region=region,
9085
instanceGroupManager=mig_name,
9186
body=dict(
9287
name="initial-resize",
@@ -105,9 +100,8 @@ def _suspend_flex_mig(mig_self_link: str, nodes: List[str], lkp: util.Lookup) ->
105100
assert nodes
106101
model = nodes[0]
107102
nodeset = lkp.node_nodeset(model)
108-
zones = nodeset.zone_policy_allow
109-
assert len(zones) == 1
110-
zone = zones[0]
103+
assert len(nodeset.zone_policy_allow) > 0
104+
region = lkp.node_region(model)
111105
project=lkp.project
112106
instanceGroupManager=util.trim_self_link(mig_self_link)
113107

@@ -118,7 +112,7 @@ def _suspend_flex_mig(mig_self_link: str, nodes: List[str], lkp: util.Lookup) ->
118112
] if inst
119113
]
120114

121-
target_mig=lkp.get_mig(lkp.project, zone, instanceGroupManager)
115+
target_mig=lkp.get_mig(lkp.project, region, instanceGroupManager)
122116
assert target_mig
123117

124118
# TODO(FLEX): This will not work if MIG didn't obtain capacity yet.
@@ -130,15 +124,15 @@ def _suspend_flex_mig(mig_self_link: str, nodes: List[str], lkp: util.Lookup) ->
130124
# - Need to `down_nodes_notify_jobs` for all nodes in MIG, make sure that it doesn't interfere with Slurm suspend-flow.
131125

132126
if target_mig["targetSize"] == len(nodes): #We can just delete the whole MIG in this case
133-
req = lkp.compute.instanceGroupManagers().delete(
127+
req = lkp.compute.regionInstanceGroupManagers().delete(
134128
project=project,
135-
zone=zone,
129+
region=region,
136130
instanceGroupManager=instanceGroupManager,
137131
)
138132
else:
139-
req = lkp.compute.instanceGroupManagers().deleteInstances(
133+
req = lkp.compute.regionInstanceGroupManagers().deleteInstances(
140134
project=project,
141-
zone=zone,
135+
region=region,
142136
instanceGroupManager=instanceGroupManager,
143137
body=dict(
144138
instances=links,
@@ -156,11 +150,10 @@ def _suspend_provisioning_inst(nodes:List[str], node_template:str, lkp: util.Loo
156150
assert nodes
157151
model = nodes[0]
158152
nodeset = lkp.node_nodeset(model)
159-
zones = nodeset.zone_policy_allow
160-
assert len(zones) == 1
161-
zone = zones[0]
153+
assert len(nodeset.zone_policy_allow) > 0
154+
region = lkp.node_region(model)
162155

163-
mig_list=lkp.get_mig_list(lkp.project, zone) #Validated via terraform that this is one
156+
mig_list=lkp.get_mig_list(lkp.project, region)
164157

165158
# FLEX (#TODO): If we enter this conditional it's likely this was called so early that MIG creation hasn't started
166159
# Consider potentially retrying? No natural mechanism for retry currently but we could
@@ -169,18 +162,18 @@ def _suspend_provisioning_inst(nodes:List[str], node_template:str, lkp: util.Loo
169162
# so until we do this is slurmsync this is a temporary workaround.
170163

171164
if not mig_list or not mig_list.get("items"):
172-
log.info("No matching MIG found to delete!")
165+
log.info("No matching MIG found to delete! Retrying...")
173166
sleep(5)
174-
mig_list=lkp.get_mig_list(lkp.project, zone)
167+
mig_list=lkp.get_mig_list(lkp.project, region)
175168
if not mig_list or not mig_list.get("items"):
176169
return
177170

178171
for mig in mig_list["items"]:
179172
if mig["instanceTemplate"] == node_template:
180173
if mig["currentActions"]["creating"] > 0 and mig["targetSize"] == mig["currentActions"]["creating"]:
181-
req = lkp.compute.instanceGroupManagers().delete(
174+
req = lkp.compute.regionInstanceGroupManagers().delete(
182175
project=lkp.project,
183-
zone=zone,
176+
region=region,
184177
instanceGroupManager=util.trim_self_link(mig["selfLink"]),
185178
)
186179

community/modules/scheduler/schedmd-slurm-gcp-v6-controller/modules/slurm_files/scripts/slurmsync.py

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,15 @@ def apply(self, nodes:List[str]) -> None:
8181
log.info(f"{len(nodes)} instances to power down ({hostlist})")
8282
run(f"{lookup().scontrol} update nodename={hostlist} state=power_down")
8383

84+
85+
@dataclass(frozen=True)
86+
class NodeActionPowerDownForce():
87+
def apply(self, nodes:List[str]) -> None:
88+
hostlist = util.to_hostlist(nodes)
89+
log.info(f"{len(nodes)} instances to power down ({hostlist})")
90+
run(f"{lookup().scontrol} update nodename={hostlist} state=power_down_force")
91+
92+
8493
@dataclass(frozen=True)
8594
class NodeActionDelete():
8695
def apply(self, nodes:List[str]) -> None:
@@ -292,7 +301,11 @@ def get_node_action(nodename: str) -> NodeAction:
292301
elif state is None:
293302
# if state is None here, the instance exists but it's not in Slurm
294303
return NodeActionUnknown(slurm_state=state, instance_state=inst.status)
295-
304+
elif lkp.is_flex_node(nodename) and "POWERING_UP" in state.flags:
305+
threshold = timedelta(seconds=int(lkp.cfg.compute_startup_scripts_timeout) * 2) #extra buffer for unexpectedly long startup scripts
306+
if util.now() - inst.creation_timestamp > threshold:
307+
log.info(f"{nodename} was unable to join the cluster after {threshold.seconds}s, potential failure on VM startup. Powering down...")
308+
return NodeActionPowerDownForce()
296309
return NodeActionUnchanged()
297310

298311

community/modules/scheduler/schedmd-slurm-gcp-v6-controller/modules/slurm_files/scripts/util.py

Lines changed: 12 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1812,18 +1812,18 @@ def _get_reservation(self, project: str, zone: str, name: str) -> Any:
18121812
project=project, zone=zone, reservation=name).execute()
18131813

18141814
@lru_cache()
1815-
def get_mig(self, project: str, zone: str, self_link:str) -> Any:
1816-
"""https://cloud.google.com/compute/docs/reference/rest/v1/instanceGroupManagers"""
1817-
return self.compute.instanceGroupManagers().get(project=project, zone=zone, instanceGroupManager=self_link).execute()
1815+
def get_mig(self, project: str, region: str, self_link:str) -> Any:
1816+
"""https://cloud.google.com/compute/docs/reference/rest/v1/regionInstanceGroupManagers"""
1817+
return self.compute.regionInstanceGroupManagers().get(project=project, region=region, instanceGroupManager=self_link).execute()
18181818

18191819
@lru_cache
1820-
def get_mig_instances(self, project: str, zone: str, self_link:str) -> Any:
1821-
return self.compute.instanceGroupManagers().listManagedInstances(project=project, zone=zone, instanceGroupManager=self_link).execute()
1820+
def get_mig_instances(self, project: str, region: str, self_link:str) -> Any:
1821+
return self.compute.regionInstanceGroupManagers().listManagedInstances(project=project, region=region, instanceGroupManager=self_link).execute()
18221822

18231823
@lru_cache()
1824-
def get_mig_list(self, project: str, zone: str) -> Any:
1825-
"""https://cloud.google.com/compute/docs/reference/rest/v1/instanceGroupManagers"""
1826-
return self.compute.instanceGroupManagers().list(project=project, zone=zone).execute()
1824+
def get_mig_list(self, project: str, region: str) -> Any:
1825+
"""https://cloud.google.com/compute/docs/reference/rest/v1/regionInstanceGroupManagers"""
1826+
return self.compute.regionInstanceGroupManagers().list(project=project, region=region).execute()
18271827

18281828
@lru_cache()
18291829
def _get_future_reservation(self, project:str, zone:str, name: str) -> Any:
@@ -2117,11 +2117,11 @@ def is_provisioning_flex_node(self, node:str) -> bool:
21172117

21182118
nodeset = self.node_nodeset(node)
21192119
zones = nodeset.zone_policy_allow
2120-
assert len(zones) == 1
2121-
zone = zones[0]
2120+
assert len(zones) > 0
2121+
region = self.node_region(node)
21222122

21232123
potential_migs=[]
2124-
mig_list=self.get_mig_list(self.project, zone)
2124+
mig_list=self.get_mig_list(self.project, region)
21252125

21262126
if not mig_list or not mig_list.get("items"):
21272127
return False
@@ -2130,7 +2130,7 @@ def is_provisioning_flex_node(self, node:str) -> bool:
21302130
if not mig.get("instanceTemplate"): #possibly an old MIG
21312131
return False
21322132
if mig["instanceTemplate"] == self.node_template(node) and mig["currentActions"]["creating"] > 0:
2133-
potential_migs.append(self.get_mig_instances(self.project, zone, trim_self_link(mig["selfLink"])))
2133+
potential_migs.append(self.get_mig_instances(self.project, region, trim_self_link(mig["selfLink"])))
21342134

21352135
if not potential_migs:
21362136
return False

0 commit comments

Comments
 (0)