[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",