Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 40 additions & 6 deletions src/buildstream/_cas/cascache.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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:
Expand Down Expand Up @@ -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():
#
Expand All @@ -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(
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand Down
16 changes: 14 additions & 2 deletions src/buildstream/_cas/casremote.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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:
Expand Down
Loading