@@ -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
0 commit comments