diff --git a/vastai/serverless/server/lib/metrics.py b/vastai/serverless/server/lib/metrics.py index 2909210..647b711 100644 --- a/vastai/serverless/server/lib/metrics.py +++ b/vastai/serverless/server/lib/metrics.py @@ -151,16 +151,13 @@ def _set_mtoken(self, mtoken: str) -> None: #######################################Private####################################### async def __send_delete_requests_and_reset(self): - async def post(report_addr: str, idxs: list[int], success_flag: bool) -> bool: + async def post(report_addr: str, requests: list[dict]) -> bool: data = { "worker_id": self.id, "mtoken": self.mtoken, - "request_idxs": idxs, - "success": success_flag, + "requests": requests, } - log.debug( - f"Deleting requests that {'succeeded' if success_flag else 'failed'}: {data['request_idxs']}" - ) + log.debug(f"Deleting requests: {[r['request_idx'] for r in requests]}") full_path = report_addr.rstrip("/") + "/delete_requests/" for attempt in range(1, 4): try: @@ -177,13 +174,12 @@ async def post(report_addr: str, idxs: list[int], success_flag: bool) -> bool: log.debug(f"Retrying delete_request, attempt: {attempt}") return False - # Take a snapshot of what we plan to send this tick. - # New arrivals after this snapshot will remain in the queue for the next tick. - snapshot = list(self.model_metrics.requests_deleting) - success_idxs = [r.request_idx for r in snapshot if r.success is True] - failed_idxs = [r.request_idx for r in snapshot if r.success is False] + requests_payload = [ + {"request_idx": r.request_idx, "success": r.success, "status": r.status} + for r in self.model_metrics.requests_deleting + ] - if not success_idxs and not failed_idxs: + if not requests_payload: return # nothing to do for report_addr in self.report_addr: @@ -191,17 +187,10 @@ async def post(report_addr: str, idxs: list[int], success_flag: bool) -> bool: if report_addr == "https://cloud.vast.ai/api/v0": # Patch: ignore the Redis API report_addr continue - sent_success = True - sent_failed = True - - if success_idxs: - sent_success = await post(report_addr, success_idxs, True) - if failed_idxs: - sent_failed = await post(report_addr, failed_idxs, False) - if sent_success and sent_failed: + if await post(report_addr, requests_payload): # Remove only the items we actually sent from the live queue. - sent_set = set(success_idxs) | set(failed_idxs) + sent_set = {r["request_idx"] for r in requests_payload} self.model_metrics.requests_deleting[:] = [ r for r in self.model_metrics.requests_deleting if r.request_idx not in sent_set