[9.1.1] Cache CAS interactions at lower layer to fix lost input handling (#29677)
The error handling for uploads in `UploadTask` depends on the exec path
of the uploaded file, which wasn't part of the `casUploadCache` cache
key. Fix this by moving deduplication logic for CAS uploads and
`FindMissing` calls into `RemoteCacheClient`.
Avoids spurious Bazel failures caused by lost inputs that can't be
recovered from via build or action rewinding due to
`CacheNotFoundException` being marked with exec paths of inputs to
concurrent actions (see the new test case).
No
- [x] I have added tests for the new use cases (if any).
- [ ] I have updated the documentation (if applicable).
RELNOTES: Fixed an issue that caused Bazel to fail on a lost input even
with build or action rewinding enabled.
Closes https://github.com/bazelbuild/bazel/pull/29551.
PiperOrigin-RevId: 922775351
Change-Id: Ia22580325c3bc92b56d202e5c24d3630fc8c27a0
(cherry picked from commit
https://github.com/bazelbuild/bazel/commit/100272284dde7a7504e7581604ddcf826788b6c8,
includes
https://github.com/bazelbuild/bazel/commit/5572634408f7ead56ebb544ab037ffa125a0154b
and
https://github.com/bazelbuild/bazel/commit/8f5af363182885e08179d47313d417e40fa922d5)
Fixes #29575
---------
Co-authored-by: tjgq <tjgq@google.com>
diff --git a/src/main/java/com/google/devtools/build/lib/bazel/repository/downloader/BUILD b/src/main/java/com/google/devtools/build/lib/bazel/repository/downloader/BUILD
index 82527ea..5eaef77c 100644
--- a/src/main/java/com/google/devtools/build/lib/bazel/repository/downloader/BUILD
+++ b/src/main/java/com/google/devtools/build/lib/bazel/repository/downloader/BUILD
@@ -25,9 +25,9 @@
"//src/main/java/com/google/devtools/build/lib/concurrent:thread_safety",
"//src/main/java/com/google/devtools/build/lib/events",
"//src/main/java/com/google/devtools/build/lib/profiler",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
"//src/main/java/com/google/devtools/build/lib/util",
"//src/main/java/com/google/devtools/build/lib/util:os",
+ "//src/main/java/com/google/devtools/build/lib/util:string",
"//src/main/java/com/google/devtools/build/lib/util:string_encoding",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//src/main/java/com/google/devtools/build/lib/vfs:pathfragment",
diff --git a/src/main/java/com/google/devtools/build/lib/bazel/repository/downloader/DownloadProgressEvent.java b/src/main/java/com/google/devtools/build/lib/bazel/repository/downloader/DownloadProgressEvent.java
index e006d0f..08f1045 100644
--- a/src/main/java/com/google/devtools/build/lib/bazel/repository/downloader/DownloadProgressEvent.java
+++ b/src/main/java/com/google/devtools/build/lib/bazel/repository/downloader/DownloadProgressEvent.java
@@ -14,8 +14,9 @@
package com.google.devtools.build.lib.bazel.repository.downloader;
+import static com.google.devtools.build.lib.util.StringUtilities.bytesCountToDisplayString;
+
import com.google.devtools.build.lib.events.ExtendedEventHandler;
-import com.google.devtools.build.lib.remote.util.Utils;
import java.net.URI;
import java.text.DecimalFormat;
import java.text.DecimalFormatSymbols;
@@ -87,10 +88,10 @@
double ratio = totalBytesDouble != 0 ? bytesRead / totalBytesDouble : 1;
// 10.1 MiB (20.2%)
return String.format(
- "%s (%s)", Utils.bytesCountToDisplayString(bytesRead), PERCENTAGE_FORMAT.format(ratio));
+ "%s (%s)", bytesCountToDisplayString(bytesRead), PERCENTAGE_FORMAT.format(ratio));
} else {
// 10.1 MiB (10,590,000B)
- return String.format("%s (%,dB)", Utils.bytesCountToDisplayString(bytesRead), bytesRead);
+ return String.format("%s (%,dB)", bytesCountToDisplayString(bytesRead), bytesRead);
}
} else {
return "";
diff --git a/src/main/java/com/google/devtools/build/lib/buildeventservice/client/BUILD b/src/main/java/com/google/devtools/build/lib/buildeventservice/client/BUILD
index 34a8ba4..3372ffe 100644
--- a/src/main/java/com/google/devtools/build/lib/buildeventservice/client/BUILD
+++ b/src/main/java/com/google/devtools/build/lib/buildeventservice/client/BUILD
@@ -20,7 +20,7 @@
],
deps = [
"//src/main/java/com/google/devtools/build/lib/buildeventstream/proto:build_event_stream_java_proto",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:tracing_metadata_utils",
"//third_party:auto_value",
"//third_party:error_prone_annotations",
"//third_party:guava",
diff --git a/src/main/java/com/google/devtools/build/lib/exec/BUILD b/src/main/java/com/google/devtools/build/lib/exec/BUILD
index 4051469..e8563a5 100644
--- a/src/main/java/com/google/devtools/build/lib/exec/BUILD
+++ b/src/main/java/com/google/devtools/build/lib/exec/BUILD
@@ -282,7 +282,7 @@
"//src/main/java/com/google/devtools/build/lib/events",
"//src/main/java/com/google/devtools/build/lib/profiler",
"//src/main/java/com/google/devtools/build/lib/remote/options",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
"//src/main/java/com/google/devtools/build/lib/util:string_encoding",
"//src/main/java/com/google/devtools/build/lib/util/io",
"//src/main/java/com/google/devtools/build/lib/util/io:io-proto",
diff --git a/src/main/java/com/google/devtools/build/lib/remote/BUILD b/src/main/java/com/google/devtools/build/lib/remote/BUILD
index a14b00f..5a3cdfb 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/BUILD
+++ b/src/main/java/com/google/devtools/build/lib/remote/BUILD
@@ -112,8 +112,13 @@
"//src/main/java/com/google/devtools/build/lib/remote/logging",
"//src/main/java/com/google/devtools/build/lib/remote/merkletree",
"//src/main/java/com/google/devtools/build/lib/remote/options",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:async_task_cache",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_output_stream",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:rx_futures",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:rx_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:tracing_metadata_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:utils",
"//src/main/java/com/google/devtools/build/lib/remote/zstd",
"//src/main/java/com/google/devtools/build/lib/rules:repository/repo_recorded_input",
"//src/main/java/com/google/devtools/build/lib/skyframe:action_execution_value",
@@ -126,6 +131,7 @@
"//src/main/java/com/google/devtools/build/lib/util:detailed_exit_code",
"//src/main/java/com/google/devtools/build/lib/util:exit_code",
"//src/main/java/com/google/devtools/build/lib/util:os",
+ "//src/main/java/com/google/devtools/build/lib/util:string",
"//src/main/java/com/google/devtools/build/lib/util:string_encoding",
"//src/main/java/com/google/devtools/build/lib/util:temp_path_generator",
"//src/main/java/com/google/devtools/build/lib/util/io",
@@ -187,7 +193,7 @@
deps = [
"//src/main/java/com/google/devtools/build/lib/profiler",
"//src/main/java/com/google/devtools/build/lib/remote/grpc",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:rx_futures",
"//third_party:guava",
"//third_party:netty",
"//third_party:rxjava3",
@@ -227,7 +233,6 @@
"//src/main/java/com/google/devtools/build/lib/clock",
"//src/main/java/com/google/devtools/build/lib/packages",
"//src/main/java/com/google/devtools/build/lib/remote/options",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
"//src/main/java/com/google/devtools/build/lib/skyframe:coverage_report_value",
"//src/main/java/com/google/devtools/build/lib/skyframe:sky_functions",
"//src/main/java/com/google/devtools/build/lib/vfs:pathfragment",
@@ -248,7 +253,9 @@
"//src/main/java/com/google/devtools/build/lib/actions:file_metadata",
"//src/main/java/com/google/devtools/build/lib/events",
"//src/main/java/com/google/devtools/build/lib/profiler",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:async_task_cache",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:rx_futures",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:utils",
"//src/main/java/com/google/devtools/build/lib/util:temp_path_generator",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//src/main/java/com/google/devtools/build/lib/vfs:pathfragment",
diff --git a/src/main/java/com/google/devtools/build/lib/remote/CombinedCache.java b/src/main/java/com/google/devtools/build/lib/remote/CombinedCache.java
index ba7aeb3..1dd7931 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/CombinedCache.java
+++ b/src/main/java/com/google/devtools/build/lib/remote/CombinedCache.java
@@ -20,8 +20,8 @@
import static com.google.common.util.concurrent.Futures.immediateFuture;
import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
import static com.google.devtools.build.lib.remote.common.ProgressStatusListener.NO_ACTION;
-import static com.google.devtools.build.lib.remote.util.Utils.bytesCountToDisplayString;
import static com.google.devtools.build.lib.remote.util.Utils.getFromFuture;
+import static com.google.devtools.build.lib.util.StringUtilities.bytesCountToDisplayString;
import build.bazel.remote.execution.v2.ActionResult;
import build.bazel.remote.execution.v2.CacheCapabilities;
@@ -48,9 +48,7 @@
import com.google.devtools.build.lib.remote.common.RemoteCacheClient.ActionKey;
import com.google.devtools.build.lib.remote.common.RemoteCacheClient.Blob;
import com.google.devtools.build.lib.remote.disk.DiskCacheClient;
-import com.google.devtools.build.lib.remote.util.AsyncTaskCache;
import com.google.devtools.build.lib.remote.util.DigestUtil;
-import com.google.devtools.build.lib.remote.util.RxFutures;
import com.google.devtools.build.lib.server.FailureDetails.FailureDetail;
import com.google.devtools.build.lib.server.FailureDetails.RemoteExecution;
import com.google.devtools.build.lib.server.FailureDetails.RemoteExecution.Code;
@@ -61,7 +59,6 @@
import com.google.errorprone.annotations.CanIgnoreReturnValue;
import com.google.protobuf.ByteString;
import io.netty.util.AbstractReferenceCounted;
-import io.reactivex.rxjava3.core.Completable;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.io.OutputStream;
@@ -98,7 +95,6 @@
SpawnCheckingCacheEvent.create("remote-cache");
private final CountDownLatch closeCountDownLatch = new CountDownLatch(1);
- protected final AsyncTaskCache.NoResult<Digest> casUploadCache = AsyncTaskCache.NoResult.create();
@SuppressWarnings("AllowVirtualThreads")
private final ListeningExecutorService virtualThreadExecutor =
@@ -386,21 +382,11 @@
ListenableFuture<Void> remoteCacheFuture = Futures.immediateVoidFuture();
if (remoteCacheClient != null && context.getWriteCachePolicy().allowRemoteCache()) {
if (chunkingSupported && digest.getSizeBytes() > chunking.config().chunkingThreshold()) {
- Completable upload =
- casUploadCache.execute(
- digest,
- RxFutures.toCompletable(
- () -> uploadChunked(context, digest, file), directExecutor()),
- /* force= */ false);
- remoteCacheFuture = RxFutures.toListenableFuture(upload);
+ remoteCacheFuture =
+ remoteCacheClient.dedupUpload(
+ digest, () -> uploadChunked(context, digest, file), /* force= */ false);
} else {
- Completable upload =
- casUploadCache.execute(
- digest,
- RxFutures.toCompletable(
- () -> remoteCacheClient.uploadFile(context, digest, file), directExecutor()),
- /* force= */ false);
- remoteCacheFuture = RxFutures.toListenableFuture(upload);
+ remoteCacheFuture = remoteCacheClient.uploadFile(context, digest, file, /* force= */ false);
}
}
@@ -451,14 +437,7 @@
ListenableFuture<Void> remoteCacheFuture = Futures.immediateVoidFuture();
if (remoteCacheClient != null && context.getWriteCachePolicy().allowRemoteCache()) {
- Completable upload =
- casUploadCache.execute(
- digest,
- RxFutures.toCompletable(
- () -> remoteCacheClient.uploadBlob(context, digest, blob), directExecutor()),
- /* force= */ false);
-
- remoteCacheFuture = RxFutures.toListenableFuture(upload);
+ remoteCacheFuture = remoteCacheClient.uploadBlob(context, digest, blob, /* force= */ false);
}
return Futures.whenAllSucceed(diskCacheFuture, remoteCacheFuture)
@@ -819,9 +798,9 @@
if (diskCacheClient != null) {
diskCacheClient.close();
}
- casUploadCache.shutdown();
virtualThreadExecutor.shutdown();
if (remoteCacheClient != null) {
+ remoteCacheClient.shutdownUploads();
remoteCacheClient.close();
}
@@ -842,14 +821,18 @@
/** Waits for active network I/Os to finish. */
public void awaitTermination() throws InterruptedException {
- casUploadCache.awaitTermination();
+ if (remoteCacheClient != null) {
+ remoteCacheClient.awaitUploadTermination();
+ }
closeCountDownLatch.await();
virtualThreadExecutor.awaitTermination(Long.MAX_VALUE, TimeUnit.NANOSECONDS);
}
/** Shuts the cache down and cancels active network I/Os. */
public void shutdownNow() {
- casUploadCache.shutdownNow();
+ if (remoteCacheClient != null) {
+ remoteCacheClient.shutdownUploadsNow();
+ }
virtualThreadExecutor.shutdownNow();
}
diff --git a/src/main/java/com/google/devtools/build/lib/remote/GrpcCacheClient.java b/src/main/java/com/google/devtools/build/lib/remote/GrpcCacheClient.java
index b9bc6ed..3e4eac4 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/GrpcCacheClient.java
+++ b/src/main/java/com/google/devtools/build/lib/remote/GrpcCacheClient.java
@@ -83,7 +83,7 @@
/** A RemoteActionCache implementation that uses gRPC calls to a remote cache server. */
@ThreadSafe
-public class GrpcCacheClient implements RemoteCacheClient, MissingDigestsFinder {
+public class GrpcCacheClient extends RemoteCacheClient implements MissingDigestsFinder {
private static final GoogleLogger logger = GoogleLogger.forEnclosingClass();
private final CallCredentialsProvider callCredentialsProvider;
@@ -572,7 +572,7 @@
}
@Override
- public ListenableFuture<Void> uploadBlob(
+ public ListenableFuture<Void> uploadBlobImpl(
RemoteActionExecutionContext context, Digest digest, Blob blob) {
return Futures.catchingAsync(
uploadChunker(
diff --git a/src/main/java/com/google/devtools/build/lib/remote/RemoteExecutionCache.java b/src/main/java/com/google/devtools/build/lib/remote/RemoteExecutionCache.java
index 2102f62..7264eda 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/RemoteExecutionCache.java
+++ b/src/main/java/com/google/devtools/build/lib/remote/RemoteExecutionCache.java
@@ -46,6 +46,7 @@
import com.google.devtools.build.lib.remote.disk.DiskCacheClient;
import com.google.devtools.build.lib.remote.merkletree.MerkleTree;
import com.google.devtools.build.lib.remote.merkletree.MerkleTreeUploader;
+import com.google.devtools.build.lib.remote.util.AsyncTaskCache;
import com.google.devtools.build.lib.remote.util.DigestUtil;
import com.google.devtools.build.lib.remote.util.RxUtils.TransferResult;
import com.google.devtools.build.lib.vfs.Path;
@@ -86,6 +87,14 @@
ListenableFuture<Boolean> isAvailableLocally(RemoteActionExecutionContext context, Path path);
}
+ /**
+ * Deduplicates concurrent {@code findMissingDigests} queries for the same digest across
+ * overlapping {@link #ensureInputsPresent} invocations. Results reported as "present" stay cached
+ * until explicitly invalidated, but a "missing" result is invalidated as soon as the triggered
+ * upload attempt terminates.
+ */
+ private final AsyncTaskCache<Digest, Boolean> findMissingCache = AsyncTaskCache.create();
+
private RemotePathChecker remotePathChecker =
new RemotePathChecker() {
@Override
@@ -199,7 +208,8 @@
RemoteActionExecutionContext context,
RemotePathResolver remotePathResolver,
Digest digest,
- Path path) {
+ Path path,
+ boolean force) {
return Futures.transformAsync(
remotePathChecker.isAvailableLocally(context, path),
isAvailableLocally -> {
@@ -217,7 +227,7 @@
throw new CacheNotFoundException(digest, path.getPathString());
}
}
- return remoteCacheClient.uploadFile(context, digest, path);
+ return remoteCacheClient.uploadFile(context, digest, path, force);
},
directExecutor());
}
@@ -226,7 +236,7 @@
public ListenableFuture<Void> uploadVirtualActionInput(
RemoteActionExecutionContext context, Digest digest, VirtualActionInput virtualActionInput) {
return remoteCacheClient.uploadBlob(
- context, digest, new VirtualActionInputBlob(virtualActionInput));
+ context, digest, new VirtualActionInputBlob(virtualActionInput), /* force= */ false);
}
private record VirtualActionInputBlob(VirtualActionInput virtualActionInput) implements Blob {
@@ -269,7 +279,8 @@
@Override
public ListenableFuture<Void> uploadBlob(
RemoteActionExecutionContext context, Digest digest, byte[] data) {
- return remoteCacheClient.uploadBlob(context, digest, () -> new ByteArrayInputStream(data));
+ return remoteCacheClient.uploadBlob(
+ context, digest, () -> new ByteArrayInputStream(data), /* force= */ false);
}
private ListenableFuture<Void> uploadBlob(
@@ -277,15 +288,16 @@
Digest digest,
MerkleTree.Uploadable merkleTree,
Map<Digest, Message> additionalInputs,
- @Nullable RemotePathResolver remotePathResolver) {
- var upload = merkleTree.upload(this, context, remotePathResolver, digest);
+ @Nullable RemotePathResolver remotePathResolver,
+ boolean force) {
+ var upload = merkleTree.upload(this, context, remotePathResolver, digest, force);
if (upload.isPresent()) {
return upload.get();
}
Message message = additionalInputs.get(digest);
if (message != null) {
- return remoteCacheClient.uploadBlob(context, digest, message.toByteString());
+ return remoteCacheClient.uploadBlob(context, digest, message.toByteString(), force);
}
return immediateFailedFuture(
@@ -344,32 +356,40 @@
uploadTask.disposable = new AtomicReference<>();
uploadTask.completion = Completable.fromObservable(completion);
Completable upload =
- casUploadCache.execute(
- digest,
- Single.<Boolean>create(
+ findMissingCache
+ .execute(
+ digest,
+ Single.<Boolean>create(
continuation -> {
uploadTask.continuation = continuation;
emitter.onSuccess(uploadTask);
- })
- .flatMapCompletable(
- shouldUpload -> {
- if (!shouldUpload) {
- return Completable.complete();
- }
-
- return toCompletable(
+ }),
+ /* onAlreadyRunning= */ () -> emitter.onSuccess(uploadTask),
+ /* onAlreadyFinished= */ () -> emitter.onSuccess(uploadTask),
+ force)
+ .flatMapCompletable(
+ shouldUpload -> {
+ if (!shouldUpload) {
+ return Completable.complete();
+ }
+ return toCompletable(
() ->
uploadBlob(
context,
uploadTask.digest,
merkleTree,
additionalInputs,
- remotePathResolver),
- directExecutor());
- }),
- /* onAlreadyRunning= */ () -> emitter.onSuccess(uploadTask),
- /* onAlreadyFinished= */ emitter::onComplete,
- force);
+ remotePathResolver,
+ force),
+ directExecutor())
+ // On success, the digest is now present remotely: replace the cached
+ // "missing" answer with "present" so late callers (or the next
+ // ensureInputsPresent invocation) skip both findMissingDigests and
+ // the upload path. On failure, invalidate so a subsequent caller
+ // (e.g., after action rewinding) re-queries the remote.
+ .doOnComplete(() -> findMissingCache.put(digest, false))
+ .doOnError(t -> findMissingCache.invalidate(digest));
+ });
upload.subscribe(
new CompletableObserver() {
@Override
diff --git a/src/main/java/com/google/devtools/build/lib/remote/RemoteExternalOverlayFileSystem.java b/src/main/java/com/google/devtools/build/lib/remote/RemoteExternalOverlayFileSystem.java
index 030d235..466d72a 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/RemoteExternalOverlayFileSystem.java
+++ b/src/main/java/com/google/devtools/build/lib/remote/RemoteExternalOverlayFileSystem.java
@@ -20,6 +20,7 @@
import static com.google.devtools.build.lib.remote.util.Utils.getFromFuture;
import static com.google.devtools.build.lib.remote.util.Utils.waitForBulkTransfer;
import static com.google.devtools.build.lib.util.StringEncoding.unicodeToInternal;
+import static com.google.devtools.build.lib.util.StringUtilities.bytesCountToDisplayString;
import build.bazel.remote.execution.v2.Digest;
import build.bazel.remote.execution.v2.Directory;
@@ -38,7 +39,6 @@
import com.google.devtools.build.lib.remote.common.RemoteActionExecutionContext;
import com.google.devtools.build.lib.remote.util.DigestUtil;
import com.google.devtools.build.lib.remote.util.TracingMetadataUtils;
-import com.google.devtools.build.lib.remote.util.Utils;
import com.google.devtools.build.lib.server.FailureDetails;
import com.google.devtools.build.lib.skyframe.SkyFunctions;
import com.google.devtools.build.lib.vfs.DetailedIOException;
@@ -707,7 +707,7 @@
@Override
public String getProgress() {
- return "(%s)".formatted(Utils.bytesCountToDisplayString(info.getSize()));
+ return "(%s)".formatted(bytesCountToDisplayString(info.getSize()));
}
@Override
diff --git a/src/main/java/com/google/devtools/build/lib/remote/chunking/BUILD b/src/main/java/com/google/devtools/build/lib/remote/chunking/BUILD
index 4f9dd29..773fb3a 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/chunking/BUILD
+++ b/src/main/java/com/google/devtools/build/lib/remote/chunking/BUILD
@@ -18,7 +18,7 @@
"FastCdcChunker.java",
],
deps = [
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
"//third_party:guava",
"//third_party:jsr305",
"@remoteapis//:build_bazel_remote_execution_v2_remote_execution_java_proto",
diff --git a/src/main/java/com/google/devtools/build/lib/remote/common/BUILD b/src/main/java/com/google/devtools/build/lib/remote/common/BUILD
index dd9fda1..e4ada2b 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/common/BUILD
+++ b/src/main/java/com/google/devtools/build/lib/remote/common/BUILD
@@ -18,7 +18,7 @@
":cache_not_found_exception",
"//src/main/java/com/google/devtools/build/lib/actions:artifacts",
"//src/main/java/com/google/devtools/build/lib/actions:important_output_handler",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
"//src/main/java/com/google/devtools/build/lib/vfs:pathfragment",
"//third_party:guava",
],
@@ -59,10 +59,13 @@
"//src/main/java/com/google/devtools/build/lib/concurrent:thread_safety",
"//src/main/java/com/google/devtools/build/lib/exec:spawn_input_expander",
"//src/main/java/com/google/devtools/build/lib/exec:spawn_runner",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:async_task_cache",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:rx_futures",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//src/main/java/com/google/devtools/build/lib/vfs:pathfragment",
"//third_party:guava",
"//third_party:jsr305",
+ "//third_party:rxjava3",
"@com_google_protobuf//:protobuf_java",
"@googleapis//google/longrunning:longrunning_java_proto",
"@remoteapis//:build_bazel_remote_execution_v2_remote_execution_java_proto",
diff --git a/src/main/java/com/google/devtools/build/lib/remote/common/RemoteCacheClient.java b/src/main/java/com/google/devtools/build/lib/remote/common/RemoteCacheClient.java
index fb351ee..6271fba 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/common/RemoteCacheClient.java
+++ b/src/main/java/com/google/devtools/build/lib/remote/common/RemoteCacheClient.java
@@ -14,14 +14,21 @@
package com.google.devtools.build.lib.remote.common;
+import static com.google.common.util.concurrent.MoreExecutors.directExecutor;
+
import build.bazel.remote.execution.v2.Action;
import build.bazel.remote.execution.v2.ActionResult;
import build.bazel.remote.execution.v2.Digest;
import build.bazel.remote.execution.v2.ServerCapabilities;
+import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Preconditions;
+import com.google.common.collect.ImmutableSet;
import com.google.common.util.concurrent.ListenableFuture;
+import com.google.devtools.build.lib.remote.util.AsyncTaskCache;
+import com.google.devtools.build.lib.remote.util.RxFutures;
import com.google.devtools.build.lib.vfs.Path;
import com.google.protobuf.ByteString;
+import io.reactivex.rxjava3.functions.Supplier;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
@@ -30,14 +37,21 @@
import javax.annotation.Nullable;
/**
- * An interface for a remote caching protocol.
+ * Base class for a remote caching protocol.
+ *
+ * <p>Concurrent uploads of the same digest are deduplicated: only one network upload is performed
+ * per digest at a time, and subsequent callers attach as observers to the in-flight upload.
+ * Implementations provide the raw network calls via the {@code *Impl} methods.
*
* <p>Implementations must be thread-safe.
*/
-public interface RemoteCacheClient extends MissingDigestsFinder {
- ServerCapabilities getServerCapabilities() throws IOException;
+public abstract class RemoteCacheClient implements MissingDigestsFinder {
- ListenableFuture<String> getAuthority();
+ private final AsyncTaskCache.NoResult<Digest> casUploadCache = AsyncTaskCache.NoResult.create();
+
+ public abstract ServerCapabilities getServerCapabilities() throws IOException;
+
+ public abstract ListenableFuture<String> getAuthority();
/**
* A key in the remote action cache. The type wraps around a {@link Digest} of an {@link Action}.
@@ -46,7 +60,7 @@
* <p>Terminology note: "action" is used here in the remote execution protocol sense, which is
* equivalent to a Bazel "spawn" (a Bazel "action" being a higher-level concept).
*/
- record ActionKey(Digest digest) {
+ public record ActionKey(Digest digest) {
public ActionKey {
Preconditions.checkNotNull(digest, "digest");
}
@@ -64,7 +78,7 @@
* @return A Future representing pending download of an action result. If an action result for
* {@code actionKey} cannot be found the result of the Future is {@code null}.
*/
- ListenableFuture<ActionResult> downloadActionResult(
+ public abstract ListenableFuture<ActionResult> downloadActionResult(
RemoteActionExecutionContext context,
ActionKey actionKey,
boolean inlineOutErr,
@@ -78,7 +92,7 @@
* @param actionResult The action result to associate with the {@code actionKey}.
* @return A Future representing pending completion of the upload.
*/
- ListenableFuture<Void> uploadActionResult(
+ public abstract ListenableFuture<Void> uploadActionResult(
RemoteActionExecutionContext context, ActionKey actionKey, ActionResult actionResult);
/**
@@ -90,7 +104,7 @@
* @return A Future representing pending completion of the download. If a BLOB for {@code digest}
* does not exist in the cache the Future fails with a {@link CacheNotFoundException}.
*/
- ListenableFuture<Void> downloadBlob(
+ public abstract ListenableFuture<Void> downloadBlob(
RemoteActionExecutionContext context, Digest digest, OutputStream out);
/**
@@ -100,7 +114,7 @@
* as late as possible.
*/
@FunctionalInterface
- interface Blob {
+ public interface Blob {
/** Get an input stream for the blob's data. Can be called multiple times. */
InputStream get() throws IOException;
@@ -114,13 +128,15 @@
/**
* Uploads a {@code file} BLOB to the CAS.
*
+ * <p>Concurrent uploads of the same digest are deduplicated. If {@code force} is true an upload
+ * that has already finished is re-executed.
+ *
* @param context the context for the action.
* @param digest The digest of the file.
* @param file The file to upload.
- * @return A future representing pending completion of the upload.
*/
- default ListenableFuture<Void> uploadFile(
- RemoteActionExecutionContext context, Digest digest, Path file) {
+ public final ListenableFuture<Void> uploadFile(
+ RemoteActionExecutionContext context, Digest digest, Path file, boolean force) {
return uploadBlob(
context,
digest,
@@ -134,34 +150,55 @@
public String description() {
return "file " + file;
}
- });
+ },
+ force);
}
/**
* Uploads a blob to the CAS.
*
+ * <p>Concurrent uploads of the same digest are deduplicated. If {@code force} is true an upload
+ * that has already finished is re-executed.
+ *
* @param context the context for the action.
* @param digest The digest of the blob.
* @param blob A supplier for the blob to upload. May be called multiple times, but is closed by
* the implementation after the upload is complete.
- * @return A future representing pending completion of the upload.
*/
- ListenableFuture<Void> uploadBlob(RemoteActionExecutionContext context, Digest digest, Blob blob);
+ public final ListenableFuture<Void> uploadBlob(
+ RemoteActionExecutionContext context, Digest digest, Blob blob, boolean force) {
+ return RxFutures.toListenableFuture(
+ casUploadCache.execute(
+ digest,
+ RxFutures.toCompletable(() -> uploadBlobImpl(context, digest, blob), directExecutor()),
+ force));
+ }
/**
* Uploads an in-memory BLOB to the CAS.
*
+ * <p>Concurrent uploads of the same digest are deduplicated. If {@code force} is true an upload
+ * that has already finished is re-executed.
+ *
* @param context the context for the action.
* @param digest The digest of the blob.
* @param data The BLOB to upload.
- * @return A future representing pending completion of the upload.
*/
- default ListenableFuture<Void> uploadBlob(
- RemoteActionExecutionContext context, Digest digest, ByteString data) {
- return uploadBlob(context, digest, data::newInput);
+ public final ListenableFuture<Void> uploadBlob(
+ RemoteActionExecutionContext context, Digest digest, ByteString data, boolean force) {
+ return uploadBlob(context, digest, (Blob) data::newInput, force);
}
/**
+ * Performs the actual network upload. Called by the deduplicating {@link #uploadBlob} wrappers.
+ *
+ * <p>Callers should use {@link #uploadBlob} instead.
+ */
+ @VisibleForTesting
+ public abstract ListenableFuture<Void> uploadBlobImpl(
+ RemoteActionExecutionContext context, Digest digest, Blob blob);
+
+ /**
* Registers a blob as the concatenation of the given chunks via SpliceBlob RPC.
*
* <p>This is used for CDC (Content-Defined Chunking) uploads. After uploading all chunks,
@@ -175,11 +212,52 @@
* is not supported by this cache client.
*/
@Nullable
- default ListenableFuture<Void> spliceBlob(
+ public ListenableFuture<Void> spliceBlob(
RemoteActionExecutionContext context, Digest blobDigest, List<Digest> chunkDigests) {
return null;
}
+ /**
+ * Deduplicates an upload by digest using the same cache as {@link #uploadFile} and {@link
+ * #uploadBlob}. For use by callers that perform their own upload logic but want to share the
+ * dedup state with the regular upload paths (e.g. chunked uploads).
+ */
+ public final ListenableFuture<Void> dedupUpload(
+ Digest digest, Supplier<ListenableFuture<Void>> upload, boolean force) {
+ return RxFutures.toListenableFuture(
+ casUploadCache.execute(digest, RxFutures.toCompletable(upload, directExecutor()), force));
+ }
+
+ /** Returns the digests currently being uploaded. */
+ public final ImmutableSet<Digest> getInProgressUploads() {
+ return casUploadCache.getInProgressTasks();
+ }
+
+ /** Returns the digests for which an upload has finished successfully. */
+ public final ImmutableSet<Digest> getFinishedUploads() {
+ return casUploadCache.getFinishedTasks();
+ }
+
+ /** Returns the number of subscribers waiting for an in-progress upload of {@code digest}. */
+ public final int getUploadSubscriberCount(Digest digest) {
+ return casUploadCache.getSubscriberCount(digest);
+ }
+
+ /** Stops accepting new uploads. */
+ public final void shutdownUploads() {
+ casUploadCache.shutdown();
+ }
+
+ /** Waits for in-progress uploads to finish. */
+ public final void awaitUploadTermination() throws InterruptedException {
+ casUploadCache.awaitTermination();
+ }
+
+ /** Cancels in-progress uploads. */
+ public final void shutdownUploadsNow() {
+ casUploadCache.shutdownNow();
+ }
+
/** Close resources associated with the remote cache. */
- void close();
+ public abstract void close();
}
diff --git a/src/main/java/com/google/devtools/build/lib/remote/disk/BUILD b/src/main/java/com/google/devtools/build/lib/remote/disk/BUILD
index 962a538..7e408a2 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/disk/BUILD
+++ b/src/main/java/com/google/devtools/build/lib/remote/disk/BUILD
@@ -20,10 +20,12 @@
"//src/main/java/com/google/devtools/build/lib/remote/common",
"//src/main/java/com/google/devtools/build/lib/remote/common:cache_not_found_exception",
"//src/main/java/com/google/devtools/build/lib/remote/options",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_output_stream",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:utils",
"//src/main/java/com/google/devtools/build/lib/server:idle_task",
"//src/main/java/com/google/devtools/build/lib/util:file_system_lock",
+ "//src/main/java/com/google/devtools/build/lib/util:string",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//third_party:flogger",
"//third_party:guava",
diff --git a/src/main/java/com/google/devtools/build/lib/remote/disk/DiskCacheGarbageCollector.java b/src/main/java/com/google/devtools/build/lib/remote/disk/DiskCacheGarbageCollector.java
index 7903db2..c4da78c 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/disk/DiskCacheGarbageCollector.java
+++ b/src/main/java/com/google/devtools/build/lib/remote/disk/DiskCacheGarbageCollector.java
@@ -14,7 +14,7 @@
package com.google.devtools.build.lib.remote.disk;
import static com.google.common.collect.ImmutableSet.toImmutableSet;
-import static com.google.devtools.build.lib.remote.util.Utils.bytesCountToDisplayString;
+import static com.google.devtools.build.lib.util.StringUtilities.bytesCountToDisplayString;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.collect.ComparisonChain;
diff --git a/src/main/java/com/google/devtools/build/lib/remote/downloader/BUILD b/src/main/java/com/google/devtools/build/lib/remote/downloader/BUILD
index 8d12624..f1ca8ee 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/downloader/BUILD
+++ b/src/main/java/com/google/devtools/build/lib/remote/downloader/BUILD
@@ -24,7 +24,8 @@
"//src/main/java/com/google/devtools/build/lib/remote:Retrier",
"//src/main/java/com/google/devtools/build/lib/remote/common",
"//src/main/java/com/google/devtools/build/lib/remote/options",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:tracing_metadata_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:utils",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//third_party:auth",
"//third_party:guava",
diff --git a/src/main/java/com/google/devtools/build/lib/remote/http/BUILD b/src/main/java/com/google/devtools/build/lib/remote/http/BUILD
index b25b40b..17de1f8 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/http/BUILD
+++ b/src/main/java/com/google/devtools/build/lib/remote/http/BUILD
@@ -24,8 +24,9 @@
"//src/main/java/com/google/devtools/build/lib/remote:Retrier",
"//src/main/java/com/google/devtools/build/lib/remote/common",
"//src/main/java/com/google/devtools/build/lib/remote/common:cache_not_found_exception",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_output_stream",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:utils",
"//third_party:auth",
"//third_party:flogger",
"//third_party:guava",
diff --git a/src/main/java/com/google/devtools/build/lib/remote/http/HttpCacheClient.java b/src/main/java/com/google/devtools/build/lib/remote/http/HttpCacheClient.java
index efab3fb..4581be3 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/http/HttpCacheClient.java
+++ b/src/main/java/com/google/devtools/build/lib/remote/http/HttpCacheClient.java
@@ -120,7 +120,7 @@
*
* <p>The implementation currently does not support transfer encoding chunked.
*/
-public final class HttpCacheClient implements RemoteCacheClient {
+public final class HttpCacheClient extends RemoteCacheClient {
private static final GoogleLogger logger = GoogleLogger.forEnclosingClass();
public static final String AC_PREFIX = "ac/";
@@ -717,7 +717,7 @@
}
@Override
- public ListenableFuture<Void> uploadBlob(
+ public ListenableFuture<Void> uploadBlobImpl(
RemoteActionExecutionContext context, Digest digest, Blob blob) {
return retrier.executeAsync(
() ->
diff --git a/src/main/java/com/google/devtools/build/lib/remote/logging/BUILD b/src/main/java/com/google/devtools/build/lib/remote/logging/BUILD
index a1d9d20..bbde938 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/logging/BUILD
+++ b/src/main/java/com/google/devtools/build/lib/remote/logging/BUILD
@@ -16,7 +16,7 @@
srcs = glob(["*.java"]),
deps = [
"//src/main/java/com/google/devtools/build/lib/clock",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:tracing_metadata_utils",
"//src/main/java/com/google/devtools/build/lib/util/io",
"//src/main/protobuf:remote_execution_log_java_proto",
"//third_party:flogger",
diff --git a/src/main/java/com/google/devtools/build/lib/remote/merkletree/BUILD b/src/main/java/com/google/devtools/build/lib/remote/merkletree/BUILD
index 5aff069..51b2904 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/merkletree/BUILD
+++ b/src/main/java/com/google/devtools/build/lib/remote/merkletree/BUILD
@@ -27,8 +27,8 @@
"//src/main/java/com/google/devtools/build/lib/remote:scrubber",
"//src/main/java/com/google/devtools/build/lib/remote/common",
"//src/main/java/com/google/devtools/build/lib/remote/common:bulk_transfer_exception",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:tracing_metadata_utils",
"//src/main/java/com/google/devtools/build/lib/skyframe:tree_artifact_value",
"//src/main/java/com/google/devtools/build/lib/util:TestType",
"//src/main/java/com/google/devtools/build/lib/util:string_encoding",
diff --git a/src/main/java/com/google/devtools/build/lib/remote/merkletree/MerkleTree.java b/src/main/java/com/google/devtools/build/lib/remote/merkletree/MerkleTree.java
index 76c652e..3896ae7 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/merkletree/MerkleTree.java
+++ b/src/main/java/com/google/devtools/build/lib/remote/merkletree/MerkleTree.java
@@ -172,7 +172,8 @@
MerkleTreeUploader uploader,
RemoteActionExecutionContext context,
RemotePathResolver remotePathResolver,
- Digest digest) {
+ Digest digest,
+ boolean force) {
return switch (blobs.get(digest)) {
case byte[] data -> Optional.of(uploader.uploadBlob(context, digest, data));
case VirtualActionInput virtualActionInput ->
@@ -189,7 +190,7 @@
: MerkleTreeComputer.PATH_ACTION_INPUT_RESOLVER;
yield Optional.of(
uploader.uploadFile(
- context, remotePathResolver, digest, pathResolver.toPath(actionInput)));
+ context, remotePathResolver, digest, pathResolver.toPath(actionInput), force));
}
case null -> Optional.empty();
default -> throw new IllegalStateException("Unexpected blob type: " + blobs.get(digest));
diff --git a/src/main/java/com/google/devtools/build/lib/remote/merkletree/MerkleTreeUploader.java b/src/main/java/com/google/devtools/build/lib/remote/merkletree/MerkleTreeUploader.java
index da34fa1..ea05f75 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/merkletree/MerkleTreeUploader.java
+++ b/src/main/java/com/google/devtools/build/lib/remote/merkletree/MerkleTreeUploader.java
@@ -32,7 +32,8 @@
RemoteActionExecutionContext context,
RemotePathResolver remotePathResolver,
Digest digest,
- Path path);
+ Path path,
+ boolean force);
/** Uploads a virtual action input to the remote cache. */
ListenableFuture<Void> uploadVirtualActionInput(
diff --git a/src/main/java/com/google/devtools/build/lib/remote/util/AsyncTaskCache.java b/src/main/java/com/google/devtools/build/lib/remote/util/AsyncTaskCache.java
index b61603d..7da0022 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/util/AsyncTaskCache.java
+++ b/src/main/java/com/google/devtools/build/lib/remote/util/AsyncTaskCache.java
@@ -33,6 +33,7 @@
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.CancellationException;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicBoolean;
import javax.annotation.concurrent.GuardedBy;
import javax.annotation.concurrent.ThreadSafe;
@@ -73,8 +74,10 @@
@GuardedBy("lock")
private final ArrayList<CompletableEmitter> terminationSubscriber = new ArrayList<>();
- @GuardedBy("lock")
- private Map<KeyT, ValueT> finished = new HashMap<>();
+ // Concurrent so that {@link #invalidate} can run without acquiring {@code lock}, which prevents
+ // lock-ordering deadlocks when invalidation is triggered from within another cache's observer
+ // notification (e.g., a doFinally on an upload completion).
+ private final ConcurrentHashMap<KeyT, ValueT> finished = new ConcurrentHashMap<>();
@GuardedBy("lock")
private Map<KeyT, Execution> inProgress = new HashMap<>();
@@ -85,9 +88,25 @@
/** Returns a set of keys for tasks which is finished. */
public ImmutableSet<KeyT> getFinishedTasks() {
- synchronized (lock) {
- return ImmutableSet.copyOf(finished.keySet());
- }
+ return ImmutableSet.copyOf(finished.keySet());
+ }
+
+ /**
+ * Removes any cached result for the given {@code key}, so that the next call to {@link #execute}
+ * for that key re-runs the task. Does not affect in-progress tasks. Safe to call concurrently
+ * with {@link #execute}.
+ */
+ public void invalidate(KeyT key) {
+ finished.remove(key);
+ }
+
+ /**
+ * Atomically replaces the cached result for {@code key} with {@code value}. The new value is
+ * visible to subsequent {@link #execute} callers. Safe to call concurrently with {@link
+ * #execute}.
+ */
+ public void put(KeyT key, ValueT value) {
+ finished.put(key, value);
}
/** Returns a set of keys for tasks which is still executing. */
@@ -299,14 +318,17 @@
return;
}
- if (!force && finished.containsKey(key)) {
- onAlreadyFinished.run();
- emitter.onSuccess(finished.get(key));
- return;
+ if (!force) {
+ ValueT cached = finished.get(key);
+ if (cached != null) {
+ onAlreadyFinished.run();
+ emitter.onSuccess(cached);
+ return;
+ }
+ } else {
+ finished.remove(key);
}
- finished.remove(key);
-
Execution execution = inProgress.get(key);
if (execution != null) {
onAlreadyRunning.run();
@@ -400,7 +422,7 @@
// Reduce retained size in case references to the cache are held after shutdown.
terminationSubscriber.trimToSize();
inProgress = new HashMap<>();
- finished = new HashMap<>();
+ finished.clear();
emitter.onComplete();
} else {
terminationSubscriber.add(emitter);
diff --git a/src/main/java/com/google/devtools/build/lib/remote/util/BUILD b/src/main/java/com/google/devtools/build/lib/remote/util/BUILD
index 41607cf..0bc9af1 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/util/BUILD
+++ b/src/main/java/com/google/devtools/build/lib/remote/util/BUILD
@@ -12,23 +12,46 @@
)
java_library(
- name = "util",
- srcs = glob(
- ["*.java"],
- exclude = [
- "DigestOutputStream.java",
- "DigestUtil.java",
- ],
- ),
+ name = "async_task_cache",
+ srcs = ["AsyncTaskCache.java"],
deps = [
- ":digest_utils",
+ "//third_party:guava",
+ "//third_party:jsr305",
+ "//third_party:rxjava3",
+ ],
+)
+
+java_library(
+ name = "rx_futures",
+ srcs = ["RxFutures.java"],
+ deps = [
+ "//third_party:guava",
+ "//third_party:jsr305",
+ "//third_party:rxjava3",
+ ],
+)
+
+java_library(
+ name = "rx_utils",
+ srcs = ["RxUtils.java"],
+ deps = [
+ "//src/main/java/com/google/devtools/build/lib/remote/common:bulk_transfer_exception",
+ "//third_party:guava",
+ "//third_party:jsr305",
+ "//third_party:rxjava3",
+ ],
+)
+
+java_library(
+ name = "utils",
+ srcs = ["Utils.java"],
+ deps = [
+ ":digest_util",
"//src/main/java/com/google/devtools/build/lib/actions",
"//src/main/java/com/google/devtools/build/lib/actions:artifacts",
"//src/main/java/com/google/devtools/build/lib/actions:execution_requirements",
- "//src/main/java/com/google/devtools/build/lib/analysis:blaze_version_info",
"//src/main/java/com/google/devtools/build/lib/authandtls",
"//src/main/java/com/google/devtools/build/lib/authandtls/credentialhelper",
- "//src/main/java/com/google/devtools/build/lib/cmdline",
"//src/main/java/com/google/devtools/build/lib/remote:ExecutionStatusException",
"//src/main/java/com/google/devtools/build/lib/remote/common",
"//src/main/java/com/google/devtools/build/lib/remote/common:bulk_transfer_exception",
@@ -38,30 +61,49 @@
"//src/main/protobuf:failure_details_java_proto",
"//third_party:guava",
"//third_party:jsr305",
- "//third_party:rxjava3",
"//third_party/grpc-java:grpc-jar",
"@com_google_protobuf//:protobuf_java",
"@com_google_protobuf//:protobuf_java_util",
"@googleapis//google/rpc:rpc_java_proto",
- "@remoteapis//:build_bazel_remote_execution_v2_remote_execution_java_grpc",
"@remoteapis//:build_bazel_remote_execution_v2_remote_execution_java_proto",
],
)
java_library(
- name = "digest_utils",
- srcs = [
- "DigestOutputStream.java",
- "DigestUtil.java",
- ],
+ name = "tracing_metadata_utils",
+ srcs = ["TracingMetadataUtils.java"],
deps = [
+ "//src/main/java/com/google/devtools/build/lib/actions",
+ "//src/main/java/com/google/devtools/build/lib/analysis:blaze_version_info",
+ "//src/main/java/com/google/devtools/build/lib/remote/options",
+ "//third_party:guava",
+ "//third_party:jsr305",
+ "//third_party/grpc-java:grpc-jar",
+ "@remoteapis//:build_bazel_remote_execution_v2_remote_execution_java_proto",
+ ],
+)
+
+java_library(
+ name = "digest_util",
+ srcs = ["DigestUtil.java"],
+ deps = [
+ ":digest_output_stream",
"//src/main/java/com/google/devtools/build/lib/remote/common",
"//src/main/java/com/google/devtools/build/lib/util:deterministic_writer",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//src/main/protobuf:spawn_java_proto",
"//third_party:guava",
- "//third_party:jsr305",
"@com_google_protobuf//:protobuf_java",
"@remoteapis//:build_bazel_remote_execution_v2_remote_execution_java_proto",
],
)
+
+java_library(
+ name = "digest_output_stream",
+ srcs = ["DigestOutputStream.java"],
+ deps = [
+ "//third_party:guava",
+ "//third_party:jsr305",
+ "@remoteapis//:build_bazel_remote_execution_v2_remote_execution_java_proto",
+ ],
+)
diff --git a/src/main/java/com/google/devtools/build/lib/remote/util/Utils.java b/src/main/java/com/google/devtools/build/lib/remote/util/Utils.java
index da3e6bb..06bcc38 100644
--- a/src/main/java/com/google/devtools/build/lib/remote/util/Utils.java
+++ b/src/main/java/com/google/devtools/build/lib/remote/util/Utils.java
@@ -28,9 +28,7 @@
import build.bazel.remote.execution.v2.Digest;
import build.bazel.remote.execution.v2.Platform;
import com.google.common.base.Ascii;
-import com.google.common.base.Preconditions;
import com.google.common.base.Throwables;
-import com.google.common.collect.ImmutableList;
import com.google.common.util.concurrent.AsyncCallable;
import com.google.common.util.concurrent.FluentFuture;
import com.google.common.util.concurrent.Futures;
@@ -76,8 +74,6 @@
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.io.OutputStream;
-import java.text.DecimalFormat;
-import java.text.DecimalFormatSymbols;
import java.time.Instant;
import java.util.Arrays;
import java.util.Locale;
@@ -545,33 +541,6 @@
}
}
- private static final ImmutableList<String> UNITS = ImmutableList.of("KiB", "MiB", "GiB", "TiB");
- // Format as single digit decimal number.
- private static final DecimalFormat BYTE_COUNT_FORMAT =
- new DecimalFormat("0.0", new DecimalFormatSymbols(Locale.US));
-
- /**
- * Converts the number of bytes to a human readable string, e.g. 1024 -> 1 KiB.
- *
- * <p>Negative numbers are not allowed.
- */
- public static String bytesCountToDisplayString(long bytes) {
- Preconditions.checkArgument(bytes >= 0);
-
- if (bytes < 1024) {
- return bytes + " B";
- }
-
- int unitIndex = 0;
- long value = bytes;
- while ((unitIndex + 1) < UNITS.size() && value >= (1 << 20)) {
- value >>= 10;
- unitIndex++;
- }
-
- return String.format("%s %s", BYTE_COUNT_FORMAT.format(value / 1024.0), UNITS.get(unitIndex));
- }
-
public static boolean shouldUploadLocalResultsToRemoteCache(
RemoteOptions remoteOptions, Map<String, String> executionInfo) {
return remoteOptions.remoteUploadLocalResults
diff --git a/src/main/java/com/google/devtools/build/lib/util/StringUtilities.java b/src/main/java/com/google/devtools/build/lib/util/StringUtilities.java
index 49d4410..35af251 100644
--- a/src/main/java/com/google/devtools/build/lib/util/StringUtilities.java
+++ b/src/main/java/com/google/devtools/build/lib/util/StringUtilities.java
@@ -13,11 +13,17 @@
// limitations under the License.
package com.google.devtools.build.lib.util;
+import static com.google.common.base.Preconditions.checkArgument;
+
import com.google.common.base.Ascii;
import com.google.common.base.Joiner;
+import com.google.common.collect.ImmutableList;
import com.google.common.escape.CharEscaperBuilder;
import com.google.common.escape.Escaper;
+import java.text.DecimalFormat;
+import java.text.DecimalFormatSymbols;
import java.util.Collection;
+import java.util.Locale;
/**
* Various utility methods operating on strings.
@@ -81,9 +87,11 @@
return result.toString();
}
+ // TODO(tjgq): Unify prettyPrintBytes and bytesCountToDisplayString.
+
/**
- * Returns an easy-to-read string approximation of a number of bytes,
- * e.g. "21MB". Note, these are IEEE units, i.e. decimal not binary powers.
+ * Returns an easy-to-read string approximation of a number of bytes, e.g. "21MB". Note, these are
+ * IEEE units, i.e. decimal not binary powers.
*/
public static String prettyPrintBytes(long bytes) {
if (bytes < 1E4) { // up to 10KB
@@ -97,6 +105,33 @@
}
}
+ private static final ImmutableList<String> UNITS = ImmutableList.of("KiB", "MiB", "GiB", "TiB");
+ // Format as single digit decimal number.
+ private static final DecimalFormat BYTE_COUNT_FORMAT =
+ new DecimalFormat("0.0", new DecimalFormatSymbols(Locale.US));
+
+ /**
+ * Converts the number of bytes to a human readable string, e.g. 1024 -> 1 KiB.
+ *
+ * <p>Negative numbers are not allowed.
+ */
+ public static String bytesCountToDisplayString(long bytes) {
+ checkArgument(bytes >= 0);
+
+ if (bytes < 1024) {
+ return bytes + " B";
+ }
+
+ int unitIndex = 0;
+ long value = bytes;
+ while ((unitIndex + 1) < UNITS.size() && value >= (1 << 20)) {
+ value >>= 10;
+ unitIndex++;
+ }
+
+ return String.format("%s %s", BYTE_COUNT_FORMAT.format(value / 1024.0), UNITS.get(unitIndex));
+ }
+
/**
* Replace control characters with visible strings.
* @return the sanitized string.
diff --git a/src/test/java/com/google/devtools/build/lib/buildeventservice/BUILD b/src/test/java/com/google/devtools/build/lib/buildeventservice/BUILD
index 0547d54..911f482 100644
--- a/src/test/java/com/google/devtools/build/lib/buildeventservice/BUILD
+++ b/src/test/java/com/google/devtools/build/lib/buildeventservice/BUILD
@@ -112,7 +112,7 @@
"//src/main/java/com/google/devtools/build/lib/buildeventstream/transports",
"//src/main/java/com/google/devtools/build/lib/network:connectivity_status",
"//src/main/java/com/google/devtools/build/lib/network:noop_connectivity",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:tracing_metadata_utils",
"//src/main/java/com/google/devtools/build/lib/util:abrupt_exit_exception",
"//src/main/java/com/google/devtools/build/lib/util:exit_code",
"//src/main/protobuf:command_line_java_proto",
@@ -134,7 +134,7 @@
"DelayingPublishBuildEventService.java",
],
deps = [
- "//src/main/java/com/google/devtools/build/lib/remote/util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:tracing_metadata_utils",
"//third_party:guava",
"//third_party:jsr305",
"//third_party:truth",
diff --git a/src/test/java/com/google/devtools/build/lib/remote/BUILD b/src/test/java/com/google/devtools/build/lib/remote/BUILD
index 3ded09a..7fe7642 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/BUILD
+++ b/src/test/java/com/google/devtools/build/lib/remote/BUILD
@@ -36,8 +36,8 @@
"//src/main/java/com/google/devtools/build/lib/actions:file_metadata",
"//src/main/java/com/google/devtools/build/lib/remote:abstract_action_input_prefetcher",
"//src/main/java/com/google/devtools/build/lib/remote/common:bulk_transfer_exception",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:utils",
"//src/main/java/com/google/devtools/build/lib/skyframe:tree_artifact_value",
"//src/main/java/com/google/devtools/build/lib/testing/vfs:spied_filesystem",
"//src/main/java/com/google/devtools/build/lib/util",
@@ -127,8 +127,10 @@
"//src/main/java/com/google/devtools/build/lib/remote/http",
"//src/main/java/com/google/devtools/build/lib/remote/merkletree",
"//src/main/java/com/google/devtools/build/lib/remote/options",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_output_stream",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:tracing_metadata_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:utils",
"//src/main/java/com/google/devtools/build/lib/runtime/commands",
"//src/main/java/com/google/devtools/build/lib/skyframe:tree_artifact_value",
"//src/main/java/com/google/devtools/build/lib/testing/vfs:spied_filesystem",
@@ -192,7 +194,6 @@
"//src/main/java/com/google/devtools/build/lib/actions",
"//src/main/java/com/google/devtools/build/lib/actions:artifacts",
"//src/main/java/com/google/devtools/build/lib/actions:file_metadata",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//src/main/java/com/google/devtools/build/lib/vfs:pathfragment",
"//third_party:guava",
@@ -281,7 +282,7 @@
"//src/main/java/com/google/devtools/build/lib/authandtls/credentialhelper:credential_module",
"//src/main/java/com/google/devtools/build/lib/dynamic",
"//src/main/java/com/google/devtools/build/lib/remote",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
"//src/main/java/com/google/devtools/build/lib/standalone",
"//src/main/java/com/google/devtools/build/lib/util:os",
"//src/main/java/com/google/devtools/build/lib/vfs",
@@ -307,7 +308,7 @@
"//src/main/java/com/google/devtools/build/lib:runtime",
"//src/main/java/com/google/devtools/build/lib/authandtls/credentialhelper:credential_module",
"//src/main/java/com/google/devtools/build/lib/remote",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:tracing_metadata_utils",
"//src/main/java/com/google/devtools/build/lib/standalone",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//src/test/java/com/google/devtools/build/lib/buildtool/util",
@@ -362,7 +363,7 @@
"//src/main/java/com/google/devtools/build/lib/events",
"//src/main/java/com/google/devtools/build/lib/remote",
"//src/main/java/com/google/devtools/build/lib/remote:store",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
"//src/main/java/com/google/devtools/build/lib/standalone",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//src/main/java/com/google/devtools/build/lib/vfs:pathfragment",
diff --git a/src/test/java/com/google/devtools/build/lib/remote/CombinedCacheTest.java b/src/test/java/com/google/devtools/build/lib/remote/CombinedCacheTest.java
index b76fe30..1017f8f 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/CombinedCacheTest.java
+++ b/src/test/java/com/google/devtools/build/lib/remote/CombinedCacheTest.java
@@ -24,7 +24,6 @@
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.verify;
-import static org.mockito.Mockito.when;
import build.bazel.remote.execution.v2.ActionResult;
import build.bazel.remote.execution.v2.CacheCapabilities;
@@ -32,6 +31,7 @@
import build.bazel.remote.execution.v2.FastCdc2020Params;
import build.bazel.remote.execution.v2.RequestMetadata;
import build.bazel.remote.execution.v2.ServerCapabilities;
+import com.google.common.collect.ImmutableClassToInstanceMap;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
import com.google.common.collect.ImmutableSet;
@@ -43,24 +43,31 @@
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import com.google.devtools.build.lib.actions.ActionInputHelper;
+import com.google.devtools.build.lib.actions.Artifact;
import com.google.devtools.build.lib.actions.ArtifactRoot;
import com.google.devtools.build.lib.actions.ArtifactRoot.RootType;
import com.google.devtools.build.lib.actions.ResourceSet;
import com.google.devtools.build.lib.actions.SimpleSpawn;
import com.google.devtools.build.lib.actions.Spawn;
+import com.google.devtools.build.lib.actions.util.ActionsTestUtil;
+import com.google.devtools.build.lib.authandtls.CallCredentialsProvider;
import com.google.devtools.build.lib.clock.JavaClock;
import com.google.devtools.build.lib.collect.nestedset.NestedSetBuilder;
import com.google.devtools.build.lib.collect.nestedset.Order;
import com.google.devtools.build.lib.exec.SpawnCheckingCacheEvent;
import com.google.devtools.build.lib.exec.SpawnRunner.SpawnExecutionContext;
import com.google.devtools.build.lib.exec.util.FakeOwner;
+import com.google.devtools.build.lib.exec.util.SpawnBuilder;
import com.google.devtools.build.lib.remote.common.BulkTransferException;
import com.google.devtools.build.lib.remote.common.RemoteActionExecutionContext;
import com.google.devtools.build.lib.remote.common.RemoteCacheClient;
import com.google.devtools.build.lib.remote.common.RemoteCacheClient.Blob;
import com.google.devtools.build.lib.remote.common.RemotePathResolver;
+import com.google.devtools.build.lib.remote.merkletree.MerkleTree;
import com.google.devtools.build.lib.remote.merkletree.MerkleTreeComputer;
+import com.google.devtools.build.lib.remote.options.RemoteOptions;
import com.google.devtools.build.lib.remote.util.DigestUtil;
+import com.google.devtools.build.lib.remote.util.FakeSpawnExecutionContext;
import com.google.devtools.build.lib.remote.util.InMemoryCacheClient;
import com.google.devtools.build.lib.remote.util.RxNoGlobalErrorsRule;
import com.google.devtools.build.lib.remote.util.TracingMetadataUtils;
@@ -74,6 +81,7 @@
import com.google.devtools.build.lib.vfs.PathFragment;
import com.google.devtools.build.lib.vfs.SyscallCache;
import com.google.devtools.build.lib.vfs.inmemoryfs.InMemoryFileSystem;
+import com.google.devtools.common.options.Options;
import com.google.protobuf.ByteString;
import java.io.IOException;
import java.io.OutputStream;
@@ -89,6 +97,7 @@
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
import javax.annotation.Nullable;
import org.junit.After;
import org.junit.Before;
@@ -282,7 +291,7 @@
@Test
@SuppressWarnings("FutureReturnValueIgnored")
public void upload_deduplicationWorks() throws IOException {
- RemoteCacheClient remoteCacheClient = mock(RemoteCacheClient.class);
+ RemoteCacheClient remoteCacheClient = spy(new InMemoryCacheClient());
AtomicInteger times = new AtomicInteger(0);
doAnswer(
invocationOnMock -> {
@@ -290,7 +299,7 @@
return SettableFuture.create();
})
.when(remoteCacheClient)
- .uploadFile(any(), any(), any());
+ .uploadBlobImpl(any(), any(), any());
CombinedCache combinedCache = newCombinedCache(remoteCacheClient);
Digest digest = fakeFileCache.createScratchInput(ActionInputHelper.fromPath("file"), "content");
Path file = execRoot.getRelative("file");
@@ -304,26 +313,16 @@
@Test
public void upload_failedUploads_doNotDeduplicate() throws Exception {
AtomicBoolean failRequest = new AtomicBoolean(true);
- InMemoryCacheClient inMemoryCacheClient = new InMemoryCacheClient();
- RemoteCacheClient remoteCacheClient = mock(RemoteCacheClient.class);
+ RemoteCacheClient remoteCacheClient = spy(new InMemoryCacheClient());
doAnswer(
invocationOnMock -> {
if (failRequest.getAndSet(false)) {
return Futures.immediateFailedFuture(new IOException("Failed"));
}
- return inMemoryCacheClient.uploadFile(
- invocationOnMock.getArgument(0),
- invocationOnMock.getArgument(1),
- invocationOnMock.getArgument(2));
+ return invocationOnMock.callRealMethod();
})
.when(remoteCacheClient)
- .uploadFile(any(), any(), any());
- doAnswer(
- invocationOnMock ->
- inMemoryCacheClient.findMissingDigests(
- invocationOnMock.getArgument(0), invocationOnMock.getArgument(1)))
- .when(remoteCacheClient)
- .findMissingDigests(any(), any());
+ .uploadBlobImpl(any(), any(), any());
CombinedCache combinedCache = newCombinedCache(remoteCacheClient);
Digest digest = fakeFileCache.createScratchInput(ActionInputHelper.fromPath("file"), "content");
Path file = execRoot.getRelative("file");
@@ -383,6 +382,141 @@
}
@Test
+ public void ensureInputsPresent_sharedMissingDigest_exceptionsHaveOwnLostInputs()
+ throws Exception {
+ RemoteCacheClient cacheProtocol = spy(new InMemoryCacheClient());
+ RemoteExecutionCache remoteCache = spy(newRemoteExecutionCache(cacheProtocol));
+
+ CountDownLatch findMissingDigestsCalls = new CountDownLatch(2);
+ doAnswer(
+ invocationOnMock -> {
+ findMissingDigestsCalls.countDown();
+ return invocationOnMock.callRealMethod();
+ })
+ .when(cacheProtocol)
+ .findMissingDigests(any(), any());
+
+ SettableFuture<Boolean> missingInputAvailable = SettableFuture.create();
+ CountDownLatch remotePathChecked = new CountDownLatch(1);
+ remoteCache.setRemotePathChecker(
+ (context, path) -> {
+ PathFragment execPath = path.relativeTo(execRoot);
+ if (execPath.equals(PathFragment.create("outputs/foo"))
+ || execPath.equals(PathFragment.create("outputs/bar"))) {
+ remotePathChecked.countDown();
+ return missingInputAvailable;
+ }
+ return immediateFuture(true);
+ });
+
+ Artifact foo = ActionsTestUtil.createArtifact(artifactRoot, "foo");
+ Artifact bar = ActionsTestUtil.createArtifact(artifactRoot, "bar");
+ Digest digest = fakeFileCache.createScratchInput(foo, "same");
+ assertThat(fakeFileCache.createScratchInput(bar, "same")).isEqualTo(digest);
+
+ Spawn fooSpawn = new SpawnBuilder().withInput(foo).build();
+ var fooContext =
+ new FakeSpawnExecutionContext(
+ fooSpawn,
+ fakeFileCache,
+ execRoot,
+ new FileOutErr(execRoot.getRelative("stdout"), execRoot.getRelative("stderr")),
+ ImmutableClassToInstanceMap.of(),
+ /* actionFileSystem= */ null);
+ var fooRemoteContext = RemoteActionExecutionContext.create(fooSpawn, fooContext, metadata);
+ var fooTree =
+ (MerkleTree.Uploadable)
+ merkleTreeComputer.buildForSpawn(
+ fooSpawn,
+ ImmutableSet.of(),
+ /* scrubber= */ null,
+ fooContext,
+ RemotePathResolver.createDefault(execRoot),
+ MerkleTreeComputer.BlobPolicy.KEEP_AND_REUPLOAD);
+
+ Spawn barSpawn = new SpawnBuilder().withInput(bar).build();
+ var barContext =
+ new FakeSpawnExecutionContext(
+ barSpawn,
+ fakeFileCache,
+ execRoot,
+ new FileOutErr(execRoot.getRelative("stdout"), execRoot.getRelative("stderr")),
+ ImmutableClassToInstanceMap.of(),
+ /* actionFileSystem= */ null);
+ var barRemoteContext = RemoteActionExecutionContext.create(barSpawn, barContext, metadata);
+ var barTree =
+ (MerkleTree.Uploadable)
+ merkleTreeComputer.buildForSpawn(
+ barSpawn,
+ ImmutableSet.of(),
+ /* scrubber= */ null,
+ barContext,
+ RemotePathResolver.createDefault(execRoot),
+ MerkleTreeComputer.BlobPolicy.KEEP_AND_REUPLOAD);
+
+ var fooFailure = new AtomicReference<Throwable>();
+ Thread fooThread =
+ new Thread(
+ () -> {
+ try {
+ remoteCache.ensureInputsPresent(
+ fooRemoteContext,
+ fooTree,
+ ImmutableMap.of(),
+ /* force= */ false,
+ RemotePathResolver.createDefault(execRoot));
+ } catch (Throwable t) {
+ if (t instanceof InterruptedException) {
+ Thread.currentThread().interrupt();
+ }
+ fooFailure.set(t);
+ }
+ });
+ fooThread.start();
+ assertThat(remotePathChecked.await(TestUtils.WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)).isTrue();
+
+ var barFailure = new AtomicReference<Throwable>();
+ Thread barThread =
+ new Thread(
+ () -> {
+ try {
+ remoteCache.ensureInputsPresent(
+ barRemoteContext,
+ barTree,
+ ImmutableMap.of(),
+ /* force= */ false,
+ RemotePathResolver.createDefault(execRoot));
+ } catch (Throwable t) {
+ if (t instanceof InterruptedException) {
+ Thread.currentThread().interrupt();
+ }
+ barFailure.set(t);
+ }
+ });
+ barThread.start();
+ assertThat(findMissingDigestsCalls.await(TestUtils.WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS))
+ .isTrue();
+
+ missingInputAvailable.set(false);
+ fooThread.join();
+ barThread.join();
+
+ assertThat(fooFailure.get()).isInstanceOf(BulkTransferException.class);
+ assertThat(
+ ((BulkTransferException) fooFailure.get())
+ .getLostArtifacts(execPath -> execPath.equals(foo.getExecPath()) ? foo : null)
+ .byDigest())
+ .containsExactly(DigestUtil.toString(digest), foo);
+
+ assertThat(barFailure.get()).isInstanceOf(BulkTransferException.class);
+ assertThat(
+ ((BulkTransferException) barFailure.get())
+ .getLostArtifacts(execPath -> execPath.equals(bar.getExecPath()) ? bar : null)
+ .byDigest())
+ .containsExactly(DigestUtil.toString(digest), bar);
+ }
+
+ @Test
public void ensureInputsPresent_interruptedDuringUploadBlobs_cancelInProgressUploadTasks()
throws Exception {
// arrange
@@ -400,7 +534,7 @@
return future;
})
.when(cacheProtocol)
- .uploadBlob(any(), any(), (Blob) any());
+ .uploadBlobImpl(any(), any(), (Blob) any());
doAnswer(
invocationOnMock -> {
SettableFuture<Void> future = SettableFuture.create();
@@ -409,7 +543,7 @@
return future;
})
.when(cacheProtocol)
- .uploadFile(any(), any(), any());
+ .uploadBlobImpl(any(), any(), any());
Path path = execRoot.getRelative("foo");
FileSystemUtils.writeContentAsLatin1(path, "bar");
@@ -439,14 +573,14 @@
thread.start();
uploadBlobCalls.await();
assertThat(futures).hasSize(2);
- assertThat(remoteCache.casUploadCache.getInProgressTasks()).isNotEmpty();
+ assertThat(cacheProtocol.getInProgressUploads()).isNotEmpty();
thread.interrupt();
ensureInputsPresentReturned.await();
// assert
- assertThat(remoteCache.casUploadCache.getInProgressTasks()).isEmpty();
- assertThat(remoteCache.casUploadCache.getFinishedTasks()).isEmpty();
+ assertThat(cacheProtocol.getInProgressUploads()).isEmpty();
+ assertThat(cacheProtocol.getFinishedUploads()).isEmpty();
for (SettableFuture<Void> future : futures) {
assertThat(future.isCancelled()).isTrue();
}
@@ -480,7 +614,7 @@
return future;
})
.when(cacheProtocol)
- .uploadBlob(any(), any(), (Blob) any());
+ .uploadBlobImpl(any(), any(), (Blob) any());
doAnswer(
invocationOnMock -> {
SettableFuture<Void> future = SettableFuture.create();
@@ -489,7 +623,7 @@
return future;
})
.when(cacheProtocol)
- .uploadFile(any(), any(), any());
+ .uploadBlobImpl(any(), any(), any());
Path path = execRoot.getRelative("foo");
FileSystemUtils.writeContentAsLatin1(path, "bar");
@@ -531,8 +665,8 @@
assertThat(futures).hasSize(2);
// assert
- assertThat(remoteCache.casUploadCache.getInProgressTasks()).hasSize(2);
- assertThat(remoteCache.casUploadCache.getFinishedTasks()).isEmpty();
+ assertThat(cacheProtocol.getInProgressUploads()).hasSize(2);
+ assertThat(cacheProtocol.getFinishedUploads()).isEmpty();
for (SettableFuture<Void> future : futures) {
assertThat(future.isCancelled()).isFalse();
}
@@ -541,8 +675,8 @@
future.set(null);
}
ensureInputsPresentReturned.await();
- assertThat(remoteCache.casUploadCache.getInProgressTasks()).isEmpty();
- assertThat(remoteCache.casUploadCache.getFinishedTasks()).hasSize(2);
+ assertThat(cacheProtocol.getInProgressUploads()).isEmpty();
+ assertThat(cacheProtocol.getFinishedUploads()).hasSize(2);
}
@Test
@@ -554,29 +688,19 @@
RemoteExecutionCache remoteCache = spy(newRemoteExecutionCache(cacheProtocol));
remoteActionExecutionContext = RemoteActionExecutionContext.create(metadata);
- ConcurrentLinkedDeque<SettableFuture<Void>> uploadBlobFutures = new ConcurrentLinkedDeque<>();
- Map<Path, SettableFuture<Void>> uploadFileFutures = Maps.newConcurrentMap();
- CountDownLatch uploadBlobCalls = new CountDownLatch(2);
- CountDownLatch uploadFileCalls = new CountDownLatch(3);
+ Map<Digest, SettableFuture<Void>> uploadFutures = Maps.newConcurrentMap();
+ // 3 unique file digests + 2 unique directory blob digests = 5 uploads total.
+ CountDownLatch uploadCalls = new CountDownLatch(5);
doAnswer(
invocationOnMock -> {
+ Digest digest = invocationOnMock.getArgument(1, Digest.class);
SettableFuture<Void> future = SettableFuture.create();
- uploadBlobFutures.add(future);
- uploadBlobCalls.countDown();
+ uploadFutures.put(digest, future);
+ uploadCalls.countDown();
return future;
})
.when(cacheProtocol)
- .uploadBlob(any(), any(), (Blob) any());
- doAnswer(
- invocationOnMock -> {
- Path file = invocationOnMock.getArgument(2, Path.class);
- SettableFuture<Void> future = SettableFuture.create();
- uploadFileFutures.put(file, future);
- uploadFileCalls.countDown();
- return future;
- })
- .when(cacheProtocol)
- .uploadFile(any(), any(), any());
+ .uploadBlobImpl(any(), any(), (Blob) any());
Path foo = execRoot.getRelative("foo");
FileSystemUtils.writeContentAsLatin1(foo, "foo");
@@ -584,6 +708,9 @@
FileSystemUtils.writeContentAsLatin1(bar, "bar");
Path qux = execRoot.getRelative("qux");
FileSystemUtils.writeContentAsLatin1(qux, "qux");
+ Digest fooDigest = digestUtil.computeAsUtf8("foo");
+ Digest barDigest = digestUtil.computeAsUtf8("bar");
+ Digest quxDigest = digestUtil.computeAsUtf8("qux");
SortedMap<PathFragment, Path> input1 = new TreeMap<>();
input1.put(PathFragment.create("foo"), foo);
@@ -635,37 +762,28 @@
// act
thread1.start();
thread2.start();
- uploadBlobCalls.await();
- uploadFileCalls.await();
- assertThat(uploadBlobFutures).hasSize(2);
- assertThat(uploadFileFutures).hasSize(3);
- assertThat(remoteCache.casUploadCache.getInProgressTasks()).hasSize(5);
+ uploadCalls.await();
+ assertThat(uploadFutures).hasSize(5);
+ assertThat(cacheProtocol.getInProgressUploads()).hasSize(5);
thread1.interrupt();
ensureInterrupted.await();
// assert
- assertThat(remoteCache.casUploadCache.getInProgressTasks()).hasSize(3);
- assertThat(remoteCache.casUploadCache.getFinishedTasks()).isEmpty();
- for (Map.Entry<Path, SettableFuture<Void>> entry : uploadFileFutures.entrySet()) {
- Path file = entry.getKey();
- SettableFuture<Void> future = entry.getValue();
- if (file.equals(foo)) {
- assertThat(future.isCancelled()).isTrue();
- } else {
- assertThat(future.isCancelled()).isFalse();
- }
- }
+ assertThat(cacheProtocol.getInProgressUploads()).hasSize(3);
+ assertThat(cacheProtocol.getFinishedUploads()).isEmpty();
+ // foo is only in tree1, so interrupting thread1 cancels it; bar is shared and qux is only in
+ // tree2, so both are kept.
+ assertThat(uploadFutures.get(fooDigest).isCancelled()).isTrue();
+ assertThat(uploadFutures.get(barDigest).isCancelled()).isFalse();
+ assertThat(uploadFutures.get(quxDigest).isCancelled()).isFalse();
- for (SettableFuture<Void> future : uploadBlobFutures) {
- future.set(null);
- }
- for (SettableFuture<Void> future : uploadFileFutures.values()) {
+ for (SettableFuture<Void> future : uploadFutures.values()) {
future.set(null);
}
ensureInputsPresentReturned.await();
- assertThat(remoteCache.casUploadCache.getInProgressTasks()).isEmpty();
- assertThat(remoteCache.casUploadCache.getFinishedTasks()).hasSize(3);
+ assertThat(cacheProtocol.getInProgressUploads()).isEmpty();
+ assertThat(cacheProtocol.getFinishedUploads()).hasSize(3);
}
@Test
@@ -674,10 +792,10 @@
remoteActionExecutionContext = RemoteActionExecutionContext.create(metadata);
doAnswer(invocationOnMock -> Futures.immediateFailedFuture(new IOException("upload failed")))
.when(cacheProtocol)
- .uploadBlob(any(), any(), (Blob) any());
+ .uploadBlobImpl(any(), any(), (Blob) any());
doAnswer(invocationOnMock -> Futures.immediateFailedFuture(new IOException("upload failed")))
.when(cacheProtocol)
- .uploadFile(any(), any(), any());
+ .uploadBlobImpl(any(), any(), any());
RemoteExecutionCache remoteCache = spy(newRemoteExecutionCache(cacheProtocol));
Path path = execRoot.getRelative("foo");
FileSystemUtils.writeContentAsLatin1(path, "bar");
@@ -700,18 +818,18 @@
@Test
public void shutdownNow_cancelInProgressUploads() throws Exception {
- RemoteCacheClient remoteCacheClient = mock(RemoteCacheClient.class);
+ RemoteCacheClient remoteCacheClient = spy(new InMemoryCacheClient());
// Return a future that never completes
doAnswer(invocationOnMock -> SettableFuture.create())
.when(remoteCacheClient)
- .uploadFile(any(), any(), any());
+ .uploadBlobImpl(any(), any(), any());
CombinedCache combinedCache = newCombinedCache(remoteCacheClient);
Digest digest = fakeFileCache.createScratchInput(ActionInputHelper.fromPath("file"), "content");
Path file = execRoot.getRelative("file");
ListenableFuture<Void> upload =
combinedCache.uploadFile(remoteActionExecutionContext, digest, file);
- assertThat(combinedCache.casUploadCache.getInProgressTasks()).contains(digest);
+ assertThat(remoteCacheClient.getInProgressUploads()).contains(digest);
combinedCache.shutdownNow();
assertThat(upload.isCancelled()).isTrue();
@@ -719,19 +837,31 @@
@Test
public void uploadFile_chunkedUpload_deduplicatesRemoteUpload() throws Exception {
- GrpcCacheClient grpcCacheClient = mock(GrpcCacheClient.class);
- when(grpcCacheClient.getServerCapabilities()).thenAnswer(unused -> chunkingCapabilities());
- when(grpcCacheClient.findMissingDigests(any(), any()))
- .thenAnswer(unused -> immediateFuture(ImmutableSet.of()));
+ // Spy on a real GrpcCacheClient so that final methods on the RemoteCacheClient base class
+ // (e.g. dedupUpload, uploadFile) execute their real implementations against a properly
+ // initialized casUploadCache.
+ GrpcCacheClient grpcCacheClient =
+ spy(
+ new GrpcCacheClient(
+ mock(ReferenceCountedChannel.class),
+ mock(CallCredentialsProvider.class),
+ Options.getDefaults(RemoteOptions.class),
+ mock(RemoteRetrier.class),
+ digestUtil));
+ doAnswer(unused -> chunkingCapabilities()).when(grpcCacheClient).getServerCapabilities();
+ doAnswer(unused -> immediateFuture(ImmutableSet.of()))
+ .when(grpcCacheClient)
+ .findMissingDigests(any(), any());
CountDownLatch spliceStarted = new CountDownLatch(1);
SettableFuture<Void> spliceFuture = SettableFuture.create();
- when(grpcCacheClient.spliceBlob(any(), any(), any()))
- .thenAnswer(
+ doAnswer(
unused -> {
spliceStarted.countDown();
return spliceFuture;
- });
+ })
+ .when(grpcCacheClient)
+ .spliceBlob(any(), any(), any());
CombinedCache combinedCache =
new CombinedCache(
@@ -755,7 +885,7 @@
ListenableFuture<Void> secondUpload =
combinedCache.uploadFile(remoteActionExecutionContext, digest, file);
- assertThat(combinedCache.casUploadCache.getSubscriberCount(digest)).isEqualTo(2);
+ assertThat(grpcCacheClient.getUploadSubscriberCount(digest)).isEqualTo(2);
verify(grpcCacheClient).findMissingDigests(any(), any());
verify(grpcCacheClient).spliceBlob(any(), any(), any());
diff --git a/src/test/java/com/google/devtools/build/lib/remote/InMemoryCombinedCache.java b/src/test/java/com/google/devtools/build/lib/remote/InMemoryCombinedCache.java
index 466e5f8..fe779a8 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/InMemoryCombinedCache.java
+++ b/src/test/java/com/google/devtools/build/lib/remote/InMemoryCombinedCache.java
@@ -74,7 +74,9 @@
Digest addContents(RemoteActionExecutionContext context, byte[] bytes)
throws IOException, InterruptedException {
Digest digest = digestUtil.compute(bytes);
- Utils.getFromFuture(remoteCacheClient.uploadBlob(context, digest, ByteString.copyFrom(bytes)));
+ Utils.getFromFuture(
+ remoteCacheClient.uploadBlob(
+ context, digest, ByteString.copyFrom(bytes), /* force= */ false));
return digest;
}
diff --git a/src/test/java/com/google/devtools/build/lib/remote/RemoteExecutionServiceTest.java b/src/test/java/com/google/devtools/build/lib/remote/RemoteExecutionServiceTest.java
index 358cc25..72e1c00 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/RemoteExecutionServiceTest.java
+++ b/src/test/java/com/google/devtools/build/lib/remote/RemoteExecutionServiceTest.java
@@ -105,7 +105,6 @@
import com.google.devtools.build.lib.remote.RemoteScrubbing.Config;
import com.google.devtools.build.lib.remote.common.BulkTransferException;
import com.google.devtools.build.lib.remote.common.RemoteActionExecutionContext;
-import com.google.devtools.build.lib.remote.common.RemoteCacheClient.Blob;
import com.google.devtools.build.lib.remote.common.RemoteExecutionClient;
import com.google.devtools.build.lib.remote.common.RemotePathResolver;
import com.google.devtools.build.lib.remote.common.RemotePathResolver.DefaultRemotePathResolver;
@@ -2679,7 +2678,7 @@
return future;
})
.when(cache.remoteCacheClient)
- .uploadBlob(any(), any(), (Blob) any());
+ .uploadBlobImpl(any(), any(), any());
ActionInput input = ActionInputHelper.fromPath("inputs/foo");
fakeFileCache.createScratchInput(input, "input-foo");
RemoteExecutionService service = newRemoteExecutionService();
diff --git a/src/test/java/com/google/devtools/build/lib/remote/chunking/BUILD b/src/test/java/com/google/devtools/build/lib/remote/chunking/BUILD
index 598fbad..90b0e38 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/chunking/BUILD
+++ b/src/test/java/com/google/devtools/build/lib/remote/chunking/BUILD
@@ -27,7 +27,7 @@
],
deps = [
"//src/main/java/com/google/devtools/build/lib/remote/chunking",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//third_party:guava",
"//third_party:junit4",
@@ -52,7 +52,7 @@
main_class = "org.openjdk.jmh.Main",
deps = [
"//src/main/java/com/google/devtools/build/lib/remote/chunking",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//src/main/java/com/google/devtools/build/lib/vfs/bazel",
"//third_party:jmh",
diff --git a/src/test/java/com/google/devtools/build/lib/remote/disk/BUILD b/src/test/java/com/google/devtools/build/lib/remote/disk/BUILD
index 599a099..f030d82 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/disk/BUILD
+++ b/src/test/java/com/google/devtools/build/lib/remote/disk/BUILD
@@ -25,8 +25,8 @@
"//src/main/java/com/google/devtools/build/lib/remote/common",
"//src/main/java/com/google/devtools/build/lib/remote/common:cache_not_found_exception",
"//src/main/java/com/google/devtools/build/lib/remote/disk",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:utils",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//src/main/java/com/google/devtools/build/lib/vfs/bazel",
"//src/main/java/com/google/devtools/build/lib/vfs/inmemoryfs",
diff --git a/src/test/java/com/google/devtools/build/lib/remote/downloader/BUILD b/src/test/java/com/google/devtools/build/lib/remote/downloader/BUILD
index 8aee210..e7ae8bf 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/downloader/BUILD
+++ b/src/test/java/com/google/devtools/build/lib/remote/downloader/BUILD
@@ -32,8 +32,9 @@
"//src/main/java/com/google/devtools/build/lib/remote/common",
"//src/main/java/com/google/devtools/build/lib/remote/downloader",
"//src/main/java/com/google/devtools/build/lib/remote/options",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:tracing_metadata_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:utils",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//src/main/java/com/google/devtools/common/options",
"//src/test/java/com/google/devtools/build/lib/remote/util",
@@ -41,7 +42,6 @@
"//src/test/java/com/google/devtools/build/lib/testutil:TestUtils",
"//third_party:auth",
"//third_party:guava",
- "//third_party:jsr305",
"//third_party:junit4",
"//third_party:mockito",
"//third_party:rxjava3",
diff --git a/src/test/java/com/google/devtools/build/lib/remote/downloader/GrpcRemoteDownloaderTest.java b/src/test/java/com/google/devtools/build/lib/remote/downloader/GrpcRemoteDownloaderTest.java
index 5542a64..6fd99ee 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/downloader/GrpcRemoteDownloaderTest.java
+++ b/src/test/java/com/google/devtools/build/lib/remote/downloader/GrpcRemoteDownloaderTest.java
@@ -237,7 +237,9 @@
final RemoteCacheClient cacheClient = new InMemoryCacheClient();
final GrpcRemoteDownloader downloader = newDownloader(cacheClient);
- getFromFuture(cacheClient.uploadBlob(context, contentDigest, ByteString.copyFrom(content)));
+ getFromFuture(
+ cacheClient.uploadBlob(
+ context, contentDigest, ByteString.copyFrom(content), /* force= */ false));
final byte[] downloaded =
downloadBlob(
downloader, URI.create("http://example.com/content.txt"), Optional.<Checksum>empty());
@@ -351,7 +353,9 @@
final GrpcRemoteDownloader downloader = newDownloader(cacheClient, /* httpDownloader= */ null);
// Add a cache entry for the empty Digest to verify that the implementation checks the status
// before fetching the digest.
- getFromFuture(cacheClient.uploadBlob(context, Digest.getDefaultInstance(), ByteString.EMPTY));
+ getFromFuture(
+ cacheClient.uploadBlob(
+ context, Digest.getDefaultInstance(), ByteString.EMPTY, /* force= */ false));
var exception =
assertThrows(
@@ -395,7 +399,9 @@
final RemoteCacheClient cacheClient = new InMemoryCacheClient();
final GrpcRemoteDownloader downloader = newDownloader(cacheClient);
- getFromFuture(cacheClient.uploadBlob(context, contentDigest, ByteString.copyFrom(content)));
+ getFromFuture(
+ cacheClient.uploadBlob(
+ context, contentDigest, ByteString.copyFrom(content), /* force= */ false));
final byte[] downloaded =
downloadBlob(
downloader,
@@ -435,7 +441,8 @@
final GrpcRemoteDownloader downloader = newDownloader(cacheClient);
getFromFuture(
- cacheClient.uploadBlob(context, contentDigest, ByteString.copyFromUtf8("wrong content")));
+ cacheClient.uploadBlob(
+ context, contentDigest, ByteString.copyFromUtf8("wrong content"), /* force= */ false));
IOException e =
assertThrows(
diff --git a/src/test/java/com/google/devtools/build/lib/remote/http/BUILD b/src/test/java/com/google/devtools/build/lib/remote/http/BUILD
index 91b6c87..ec1b56a 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/http/BUILD
+++ b/src/test/java/com/google/devtools/build/lib/remote/http/BUILD
@@ -27,8 +27,9 @@
"//src/main/java/com/google/devtools/build/lib/remote:Retrier",
"//src/main/java/com/google/devtools/build/lib/remote/common",
"//src/main/java/com/google/devtools/build/lib/remote/http",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:tracing_metadata_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:utils",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//src/main/java/com/google/devtools/common/options",
"//src/test/java/com/google/devtools/build/lib:test_runner",
diff --git a/src/test/java/com/google/devtools/build/lib/remote/http/HttpCacheClientTest.java b/src/test/java/com/google/devtools/build/lib/remote/http/HttpCacheClientTest.java
index cc1dc9d..2c140d4 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/http/HttpCacheClientTest.java
+++ b/src/test/java/com/google/devtools/build/lib/remote/http/HttpCacheClientTest.java
@@ -366,7 +366,7 @@
ByteString data = ByteString.copyFrom("foo bar", StandardCharsets.UTF_8);
Digest digest = DIGEST_UTIL.compute(data.toByteArray());
- blobStore.uploadBlob(remoteActionExecutionContext, digest, data).get();
+ blobStore.uploadBlob(remoteActionExecutionContext, digest, data, /* force= */ false).get();
assertThat(cacheContents).hasSize(1);
String cacheKey = "/cas/" + digest.getHash();
@@ -420,7 +420,8 @@
blobStore.uploadBlob(
remoteActionExecutionContext,
DIGEST_UTIL.compute(data),
- ByteString.copyFrom(data))));
+ ByteString.copyFrom(data),
+ /* force= */ false)));
} finally {
testServer.stop(server);
}
@@ -491,7 +492,8 @@
blobStore.uploadBlob(
remoteActionExecutionContext,
DIGEST_UTIL.compute(data.toByteArray()),
- data)));
+ data,
+ /* force= */ false)));
assertThat(e.getCause()).isInstanceOf(TooLongFrameException.class);
} finally {
testServer.stop(server);
@@ -774,7 +776,10 @@
byte[] data = "File Contents".getBytes(StandardCharsets.US_ASCII);
blobStore
.uploadBlob(
- remoteActionExecutionContext, DIGEST_UTIL.compute(data), ByteString.copyFrom(data))
+ remoteActionExecutionContext,
+ DIGEST_UTIL.compute(data),
+ ByteString.copyFrom(data),
+ /* force= */ false)
.get();
verify(credentials, times(1)).refresh();
verify(credentials, times(2)).getRequestMetadata(any(URI.class));
@@ -839,7 +844,8 @@
blobStore.uploadBlob(
remoteActionExecutionContext,
DIGEST_UTIL.compute(oneByte),
- ByteString.copyFrom(oneByte)));
+ ByteString.copyFrom(oneByte),
+ /* force= */ false));
fail("Exception expected.");
} catch (Exception e) {
assertThat(e).isInstanceOf(HttpException.class);
diff --git a/src/test/java/com/google/devtools/build/lib/remote/logging/BUILD b/src/test/java/com/google/devtools/build/lib/remote/logging/BUILD
index c457ccc..326f60c 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/logging/BUILD
+++ b/src/test/java/com/google/devtools/build/lib/remote/logging/BUILD
@@ -19,7 +19,7 @@
test_class = "com.google.devtools.build.lib.AllTests",
deps = [
"//src/main/java/com/google/devtools/build/lib/remote/logging",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
"//src/main/java/com/google/devtools/build/lib/util/io",
"//src/main/protobuf:remote_execution_log_java_proto",
"//src/test/java/com/google/devtools/build/lib:test_runner",
diff --git a/src/test/java/com/google/devtools/build/lib/remote/merkletree/BUILD b/src/test/java/com/google/devtools/build/lib/remote/merkletree/BUILD
index 89b4f8b..e6a10c1 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/merkletree/BUILD
+++ b/src/test/java/com/google/devtools/build/lib/remote/merkletree/BUILD
@@ -28,7 +28,7 @@
"//src/main/java/com/google/devtools/build/lib/clock",
"//src/main/java/com/google/devtools/build/lib/remote/common",
"//src/main/java/com/google/devtools/build/lib/remote/merkletree",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
"//src/main/java/com/google/devtools/build/lib/skyframe:tree_artifact_value",
"//src/main/java/com/google/devtools/build/lib/util/io",
"//src/main/java/com/google/devtools/build/lib/vfs",
diff --git a/src/test/java/com/google/devtools/build/lib/remote/merkletree/MerkleTreeComputerTest.java b/src/test/java/com/google/devtools/build/lib/remote/merkletree/MerkleTreeComputerTest.java
index 05af964..e849db0a 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/merkletree/MerkleTreeComputerTest.java
+++ b/src/test/java/com/google/devtools/build/lib/remote/merkletree/MerkleTreeComputerTest.java
@@ -217,7 +217,8 @@
RemoteActionExecutionContext context,
RemotePathResolver remotePathResolver,
Digest digest,
- Path path) {
+ Path path,
+ boolean force) {
return immediateVoidFuture();
}
diff --git a/src/test/java/com/google/devtools/build/lib/remote/util/BUILD b/src/test/java/com/google/devtools/build/lib/remote/util/BUILD
index d496c5c..fccacf9 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/util/BUILD
+++ b/src/test/java/com/google/devtools/build/lib/remote/util/BUILD
@@ -32,7 +32,9 @@
"//src/main/java/com/google/devtools/build/lib/remote/common",
"//src/main/java/com/google/devtools/build/lib/remote/common:bulk_transfer_exception",
"//src/main/java/com/google/devtools/build/lib/remote/common:cache_not_found_exception",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:async_task_cache",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:rx_futures",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:rx_utils",
"//src/main/java/com/google/devtools/build/lib/util/io",
"//src/main/java/com/google/devtools/build/lib/vfs",
"//src/main/java/com/google/devtools/build/lib/vfs:pathfragment",
diff --git a/src/test/java/com/google/devtools/build/lib/remote/util/InMemoryCacheClient.java b/src/test/java/com/google/devtools/build/lib/remote/util/InMemoryCacheClient.java
index d8d5c1e..af9b9f2 100644
--- a/src/test/java/com/google/devtools/build/lib/remote/util/InMemoryCacheClient.java
+++ b/src/test/java/com/google/devtools/build/lib/remote/util/InMemoryCacheClient.java
@@ -39,7 +39,7 @@
import java.util.stream.Collectors;
/** A {@link RemoteCacheClient} that stores its contents in memory. */
-public class InMemoryCacheClient implements RemoteCacheClient {
+public class InMemoryCacheClient extends RemoteCacheClient {
private final ListeningExecutorService executorService =
MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(100));
@@ -140,7 +140,7 @@
}
@Override
- public ListenableFuture<Void> uploadBlob(
+ public ListenableFuture<Void> uploadBlobImpl(
RemoteActionExecutionContext context, Digest digest, Blob blob) {
try {
cas.put(digest, blob.get().readAllBytes());
diff --git a/src/test/java/com/google/devtools/build/lib/remote/util/UtilsTest.java b/src/test/java/com/google/devtools/build/lib/remote/util/UtilsTest.java
deleted file mode 100644
index 6ec5583..0000000
--- a/src/test/java/com/google/devtools/build/lib/remote/util/UtilsTest.java
+++ /dev/null
@@ -1,39 +0,0 @@
-// Copyright 2021 The Bazel Authors. All rights reserved.
-//
-// Licensed under the Apache License, Version 2.0 (the "License");
-// you may not use this file except in compliance with the License.
-// You may obtain a copy of the License at
-//
-// http://www.apache.org/licenses/LICENSE-2.0
-//
-// Unless required by applicable law or agreed to in writing, software
-// distributed under the License is distributed on an "AS IS" BASIS,
-// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-// See the License for the specific language governing permissions and
-// limitations under the License.
-package com.google.devtools.build.lib.remote.util;
-
-import static com.google.common.truth.Truth.assertThat;
-import static com.google.devtools.build.lib.remote.util.Utils.bytesCountToDisplayString;
-
-import org.junit.Test;
-import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
-
-/** Tests for {@link Utils}. */
-@RunWith(JUnit4.class)
-public class UtilsTest {
- @Test
- public void bytesCountToDisplayString_works() {
- assertThat(bytesCountToDisplayString(1000)).isEqualTo("1000 B");
- assertThat(bytesCountToDisplayString(1 << 10)).isEqualTo("1.0 KiB");
- assertThat(bytesCountToDisplayString((1 << 10) + (1 << 10) / 10)).isEqualTo("1.1 KiB");
- assertThat(bytesCountToDisplayString(1 << 20)).isEqualTo("1.0 MiB");
- assertThat(bytesCountToDisplayString((1 << 20) + (1 << 20) / 10)).isEqualTo("1.1 MiB");
- assertThat(bytesCountToDisplayString(1 << 30)).isEqualTo("1.0 GiB");
- assertThat(bytesCountToDisplayString((1 << 30) + (1 << 30) / 10)).isEqualTo("1.1 GiB");
- assertThat(bytesCountToDisplayString(1L << 40)).isEqualTo("1.0 TiB");
- assertThat(bytesCountToDisplayString((1L << 40) + (1L << 40) / 10)).isEqualTo("1.1 TiB");
- assertThat(bytesCountToDisplayString(1L << 50)).isEqualTo("1024.0 TiB");
- }
-}
diff --git a/src/test/java/com/google/devtools/build/lib/util/StringUtilitiesTest.java b/src/test/java/com/google/devtools/build/lib/util/StringUtilitiesTest.java
index e7a65dd..4716d47 100644
--- a/src/test/java/com/google/devtools/build/lib/util/StringUtilitiesTest.java
+++ b/src/test/java/com/google/devtools/build/lib/util/StringUtilitiesTest.java
@@ -14,6 +14,7 @@
package com.google.devtools.build.lib.util;
import static com.google.common.truth.Truth.assertThat;
+import static com.google.devtools.build.lib.util.StringUtilities.bytesCountToDisplayString;
import static com.google.devtools.build.lib.util.StringUtilities.joinLines;
import static com.google.devtools.build.lib.util.StringUtilities.prettyPrintBytes;
@@ -74,6 +75,20 @@
}
@Test
+ public void testBytesCountToDisplayString() {
+ assertThat(bytesCountToDisplayString(1000)).isEqualTo("1000 B");
+ assertThat(bytesCountToDisplayString(1 << 10)).isEqualTo("1.0 KiB");
+ assertThat(bytesCountToDisplayString((1 << 10) + (1 << 10) / 10)).isEqualTo("1.1 KiB");
+ assertThat(bytesCountToDisplayString(1 << 20)).isEqualTo("1.0 MiB");
+ assertThat(bytesCountToDisplayString((1 << 20) + (1 << 20) / 10)).isEqualTo("1.1 MiB");
+ assertThat(bytesCountToDisplayString(1 << 30)).isEqualTo("1.0 GiB");
+ assertThat(bytesCountToDisplayString((1 << 30) + (1 << 30) / 10)).isEqualTo("1.1 GiB");
+ assertThat(bytesCountToDisplayString(1L << 40)).isEqualTo("1.0 TiB");
+ assertThat(bytesCountToDisplayString((1L << 40) + (1L << 40) / 10)).isEqualTo("1.1 TiB");
+ assertThat(bytesCountToDisplayString(1L << 50)).isEqualTo("1024.0 TiB");
+ }
+
+ @Test
public void sanitizeControlChars() {
assertThat(StringUtilities.sanitizeControlChars("\000")).isEqualTo("<?>");
assertThat(StringUtilities.sanitizeControlChars("\001")).isEqualTo("<?>");
diff --git a/src/tools/diskcache/BUILD b/src/tools/diskcache/BUILD
index 90bcf16..1154eb5 100644
--- a/src/tools/diskcache/BUILD
+++ b/src/tools/diskcache/BUILD
@@ -18,7 +18,6 @@
visibility = ["//visibility:public"],
deps = [
"//src/main/java/com/google/devtools/build/lib/remote/disk",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
"//src/main/java/com/google/devtools/build/lib/unix",
"//src/main/java/com/google/devtools/build/lib/util:os",
"//src/main/java/com/google/devtools/build/lib/vfs",
diff --git a/src/tools/remote/src/main/java/com/google/devtools/build/remote/worker/BUILD b/src/tools/remote/src/main/java/com/google/devtools/build/remote/worker/BUILD
index 0f5717b..1e3dcf4 100644
--- a/src/tools/remote/src/main/java/com/google/devtools/build/remote/worker/BUILD
+++ b/src/tools/remote/src/main/java/com/google/devtools/build/remote/worker/BUILD
@@ -39,8 +39,10 @@
"//src/main/java/com/google/devtools/build/lib/remote/common",
"//src/main/java/com/google/devtools/build/lib/remote/common:cache_not_found_exception",
"//src/main/java/com/google/devtools/build/lib/remote/disk",
- "//src/main/java/com/google/devtools/build/lib/remote/util",
- "//src/main/java/com/google/devtools/build/lib/remote/util:digest_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_output_stream",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:digest_util",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:tracing_metadata_utils",
+ "//src/main/java/com/google/devtools/build/lib/remote/util:utils",
"//src/main/java/com/google/devtools/build/lib/sandbox:linux_sandbox_command_line_builder",
"//src/main/java/com/google/devtools/build/lib/shell",
"//src/main/java/com/google/devtools/build/lib/unix",