From 50f97b86d0a544521614627dfcabc483c788adda Mon Sep 17 00:00:00 2001 From: Abderrahim Kitouni Date: Mon, 3 Aug 2026 16:00:03 +0100 Subject: [PATCH] _cas: avoid blocking calls to long grpc methods Long-running synchronous grpc calls block the thread where they are called, which can cause the scheduler to be unable to terminate the job. This commit replaces them with usage of the "asynchronous" future API. The result() method of a grpc future will not block the thread while it's waiting, and thus allows terminating the job without waiting for the blocking call to return. This covers the LocalCAS methods {Fetch,Upload}{Tree,MissingBlobs}, and the CAS method FindMissingBlobs. As an added benefit, we try to cancel the grpc calls when terminated. This depends on the server side (buildbox-casd) promptly cancelling the remote calls, which it doesn't consistently do at the moment. But that can be fixed independently. Fixes https://github.com/apache/buildstream/issues/2157 --- src/buildstream/_cas/cascache.py | 46 +++++++++++++++++++++++++++---- src/buildstream/_cas/casremote.py | 16 +++++++++-- 2 files changed, 54 insertions(+), 8 deletions(-) diff --git a/src/buildstream/_cas/cascache.py b/src/buildstream/_cas/cascache.py index 68fd4b610..61a84b694 100644 --- a/src/buildstream/_cas/cascache.py +++ b/src/buildstream/_cas/cascache.py @@ -127,7 +127,13 @@ def contains_directory(self, digest): request.fetch_file_blobs = not self._remote_cache try: - local_cas.FetchTree(request) + response_future = local_cas.FetchTree.future(request) + try: + response_future.result() + except: + response_future.cancel() + raise + if not self._remote_cache: return True except grpc.RpcError as e: @@ -141,7 +147,13 @@ def contains_directory(self, digest): request = local_cas_pb2.UploadTreeRequest() request.root_digest.CopyFrom(digest) try: - local_cas.UploadTree(request) + response_future = local_cas.UploadTree.future(request) + try: + response_future.result() + except: + response_future.cancel() + raise + return True except grpc.RpcError as e: if e.code() == grpc.StatusCode.NOT_FOUND: @@ -231,7 +243,12 @@ def ensure_tree(self, tree): request.root_digest.CopyFrom(tree) request.fetch_file_blobs = True - local_cas.FetchTree(request) + response_future = local_cas.FetchTree.future(request) + try: + response_future.result() + except: + response_future.cancel() + raise # fetch_directory(): # @@ -252,7 +269,13 @@ def fetch_directory(self, remote, dir_digest): request.fetch_file_blobs = False try: - local_cas.FetchTree(request) + response_future = local_cas.FetchTree.future(request) + try: + response_future.result() + except: + response_future.cancel() + raise + except grpc.RpcError as e: if e.code() == grpc.StatusCode.NOT_FOUND: raise BlobNotFound( @@ -536,7 +559,13 @@ def missing_blobs(self, blobs, *, remote=None): d.CopyFrom(required_digest) try: - response = cas.FindMissingBlobs(request) + response_future = cas.FindMissingBlobs.future(request) + try: + response = response_future.result() + except: + response_future.cancel() + raise + except grpc.RpcError as e: if e.code() == grpc.StatusCode.INVALID_ARGUMENT and e.details().startswith("Invalid instance name"): raise CASCacheError("Unsupported buildbox-casd version: FindMissingBlobs failed") from e @@ -566,7 +595,12 @@ def required_blobs_for_directory(self, directory_digest, *, excluded_subdirs=Non request.root_digest.CopyFrom(directory_digest) request.fetch_file_blobs = False - local_cas.FetchTree(request) + response_future = local_cas.FetchTree.future(request) + try: + response_future.result() + except: + response_future.cancel() + raise # parse directory, and recursively add blobs diff --git a/src/buildstream/_cas/casremote.py b/src/buildstream/_cas/casremote.py index 3fefb9d47..f04d8d7be 100644 --- a/src/buildstream/_cas/casremote.py +++ b/src/buildstream/_cas/casremote.py @@ -92,7 +92,13 @@ def send(self, *, missing_blobs=None): local_cas = self._remote.casd.get_local_cas() for request in self._requests: - batch_response = local_cas.FetchMissingBlobs(request) + batch_response_future = local_cas.FetchMissingBlobs.future(request) + + try: + batch_response = batch_response_future.result() + except: + batch_response_future.cancel() + raise for response in batch_response.responses: if response.status.code == code_pb2.NOT_FOUND: @@ -146,7 +152,13 @@ def send(self): local_cas = self._remote.casd.get_local_cas() for request in self._requests: - batch_response = local_cas.UploadMissingBlobs(request) + batch_response_future = local_cas.UploadMissingBlobs.future(request) + + try: + batch_response = batch_response_future.result() + except: + batch_response_future.cancel() + raise for response in batch_response.responses: if response.status.code != code_pb2.OK: