Skip to content

Commit ff16c40

Browse files
authored
feat: svc-redis 僵尸实例问题与分配逻辑优化(part3) (#2756)
1 parent a5a3a0e commit ff16c40

3 files changed

Lines changed: 94 additions & 6 deletions

File tree

apiserver/paasng/paasng/accessories/servicehub/remote/client.py

Lines changed: 37 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
"""Client for remote services"""
1919

2020
import logging
21+
import time
2122
from contextlib import contextmanager
2223
from dataclasses import MISSING, dataclass
2324
from typing import Dict, List
@@ -41,6 +42,8 @@ class RemoteSvcConfig:
4142
provision_params_tmpl: Dict
4243
jwt_auth_conf: Dict
4344
prefer_async_delete: bool = True
45+
# 是否使用幂等创建实例接口
46+
prefer_idempotent_provision: bool = False
4447
is_ready: bool = True
4548

4649
@classmethod
@@ -86,6 +89,7 @@ def __post_init__(self):
8689
self.retrieve_instance_by_name_url = urljoin(self.endpoint_url, "services/{service_id}/instances/?name={name}")
8790
self.update_inst_config_url = urljoin(self.endpoint_url, "instances/{instance_id}/config/")
8891
self.create_instance_url = urljoin(self.endpoint_url, "services/{service_id}/instances/{instance_id}/")
92+
self.idem_prov = urljoin(self.endpoint_url, "services/{service_id}/instances/idem_prov/")
8993
self.delete_instance_url = urljoin(self.endpoint_url, "instances/{instance_id}/")
9094
self.async_delete_instance_url = urljoin(self.endpoint_url, "instances/{instance_id}/async_delete")
9195
# 增强服务绑定
@@ -132,6 +136,25 @@ def validate_resp(resp: requests.Response):
132136
response_text=resp.text,
133137
)
134138

139+
def _post_until_ready(self, url: str, payload: Dict) -> requests.Response:
140+
"""Retry POST requests until the async endpoint returns 200/201 or timeout."""
141+
deadline = time.monotonic() + self.REQUEST_CREATE_TIMEOUT
142+
143+
poll_interval = 0.5
144+
while time.monotonic() < deadline:
145+
resp = requests.post(url, json=payload, auth=self.auth, timeout=self.REQUEST_CREATE_TIMEOUT)
146+
self.validate_resp(resp)
147+
# 200,201 表示资源已就绪,可以返回了
148+
# 202,表示请求已接受,但资源未就绪,需要继续轮询
149+
if resp.status_code in {200, 201}:
150+
return resp
151+
time.sleep(poll_interval)
152+
poll_interval = min(poll_interval * 2, 4)
153+
154+
raise RemoteClientError(
155+
f"POST to {desensitize_url(url)} not ready after retrying for {self.REQUEST_CREATE_TIMEOUT} seconds"
156+
)
157+
135158
def get_meta_info(self) -> Dict:
136159
"""Get service's meta info
137160
@@ -214,6 +237,19 @@ def provision_instance(self, service_id: str, plan_id: str, instance_id: str, pa
214237
self.validate_resp(resp)
215238
return resp.json()
216239

240+
def idempotent_provision_instance(self, service_id: str, plan_id: str, params: Dict) -> Dict:
241+
"""Idempotently provision a new instance,
242+
`params['engine_app_name']` treated as the idempotency key, so it's required
243+
244+
:raises: RemoteClientError
245+
:return: <instance dict>
246+
"""
247+
url = self.config.idem_prov.format(service_id=service_id)
248+
payload = {"plan_id": plan_id, "params": params}
249+
with wrap_request_exc(self):
250+
resp = self._post_until_ready(url, payload)
251+
return resp.json()
252+
217253
def retrieve_instance(self, instance_id: str) -> Dict:
218254
"""Retrieve a provisioned instance info
219255
@@ -291,8 +327,7 @@ def update_instance_config(self, instance_id: str, config: Dict):
291327
def create_client_side_instance(self, service_id: str, instance_id: str, params: Dict):
292328
url = self.config.create_client_side_instance_url.format(service_id=service_id, instance_id=instance_id)
293329
with wrap_request_exc(self):
294-
resp = requests.post(url, json=params, auth=self.auth, timeout=self.REQUEST_CREATE_TIMEOUT)
295-
self.validate_resp(resp)
330+
resp = self._post_until_ready(url, params)
296331
return resp.json()
297332

298333
def destroy_client_side_instance(self, instance_id: str):

apiserver/paasng/paasng/accessories/servicehub/remote/manager.py

Lines changed: 15 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -110,6 +110,7 @@ def semantic_version_gte(self, version: str) -> bool:
110110
DEFAULT_META_INFO = MetaInfo(version=None)
111111
VERSION_WITH_INST_CONFIG = "0.1.0"
112112
VERSION_WITH_REST_UPSERT = "0.2.0"
113+
VERSION_WITH_IDEMPOTENT_PROVISION = "2.0.4"
113114

114115

115116
@dataclass
@@ -179,6 +180,10 @@ def supports_rest_upsert(self) -> bool:
179180
"""Check if current service supports Feature: RestFul upsert Service/Plan"""
180181
return self.meta_info.semantic_version_gte(VERSION_WITH_REST_UPSERT)
181182

183+
def supports_idempotent_provision(self) -> bool:
184+
"""Check if current service supports idempotent provision, which means provisioning an already provisioned instance will not cause error"""
185+
return self.meta_info.semantic_version_gte(VERSION_WITH_IDEMPOTENT_PROVISION)
186+
182187

183188
@dataclass
184189
class EnvClusterInfo:
@@ -237,12 +242,18 @@ def provision(self):
237242
logger.warning(f"remote service {self.get_service().name} is not ready, skip")
238243
return
239244

240-
instance_id = str(uuid.uuid4())
241245
try:
242246
params = self.render_params(self.remote_config.provision_params_tmpl)
243-
self.remote_client.provision_instance(
244-
str(self.db_obj.service_id), str(self.db_obj.plan_id), instance_id, params=params
245-
)
247+
if self.get_service().supports_idempotent_provision():
248+
resp = self.remote_client.idempotent_provision_instance(
249+
str(self.db_obj.service_id), str(self.db_obj.plan_id), params=params
250+
)
251+
instance_id = resp["uuid"]
252+
else:
253+
instance_id = str(uuid.uuid4())
254+
self.remote_client.provision_instance(
255+
str(self.db_obj.service_id), str(self.db_obj.plan_id), instance_id, params=params
256+
)
246257
except Exception as e:
247258
logger.exception(f"Error provisioning new instance for {self.db_application.name}")
248259
raise exceptions.ProvisionInstanceError(

apiserver/paasng/tests/paasng/accessories/servicehub/remote/test_manager.py

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -150,6 +150,48 @@ def test_provision(
150150
assert bool(all(mocked_provision.call_args[0]))
151151
assert mocked_provision.call_args[1]["params"]["username"] == rel.db_engine_app.name
152152

153+
@mock.patch("paas_wl.workloads.networking.egress.shim.get_cluster_egress_info")
154+
@mock.patch("paasng.accessories.servicehub.remote.client.RemoteServiceClient.provision_instance")
155+
@mock.patch("paasng.accessories.servicehub.remote.client.RemoteServiceClient.idempotent_provision_instance")
156+
def test_provision_with_idempotent_api(
157+
self,
158+
mocked_idempotent_provision,
159+
mocked_provision,
160+
get_cluster_egress_info,
161+
config,
162+
store,
163+
bk_module,
164+
bk_service,
165+
bk_plan_1,
166+
):
167+
"""Test service instance provision through idempotent API"""
168+
get_cluster_egress_info.return_value = {"egress_ips": ["1.1.1.1"], "digest_version": "foo"}
169+
mocked_idempotent_provision.return_value = {"uuid": str(uuid.uuid4())}
170+
config.prefer_idempotent_provision = True
171+
172+
plans = [bk_plan_1]
173+
mgr = RemoteServiceMgr(store=store)
174+
bk_service.plans = plans
175+
bk_service.meta_info = MetaInfo(version="2.0.4")
176+
177+
SvcBindingPolicyManager(bk_service, DEFAULT_TENANT_ID).set_uniform(plans=[plans[0].uuid])
178+
mgr.bind_service(bk_service, bk_module)
179+
180+
with mock.patch.object(mgr, "get") as get_service:
181+
get_service.return_value = bk_service
182+
env = bk_module.get_envs("stag")
183+
rel = next(mgr.list_unprovisioned_rels(env.engine_app))
184+
rel.provision()
185+
186+
assert rel.is_provisioned() is True
187+
assert str(rel.db_obj.service_instance_id) == mocked_idempotent_provision.return_value["uuid"]
188+
mocked_idempotent_provision.assert_called_once_with(
189+
str(rel.db_obj.service_id),
190+
str(rel.db_obj.plan_id),
191+
params={"username": rel.db_engine_app.name},
192+
)
193+
mocked_provision.assert_not_called()
194+
153195
@mock.patch("paasng.accessories.servicehub.remote.manager.EnvClusterInfo.get_egress_info")
154196
def test_render_params(
155197
self,

0 commit comments

Comments
 (0)