From 4b7d8bd2c480efc81b6ec58ecda561a9d843de4d Mon Sep 17 00:00:00 2001 From: Olivier Notteghem Date: Wed, 16 Sep 2026 19:00:29 +0000 Subject: [PATCH 1/3] Prototype live compact execution log streaming --- .../com/google/devtools/build/lib/bazel/BUILD | 5 + .../bazel/ExecutionLogGrpcOutputStream.java | 190 ++++++++++++++++++ .../build/lib/bazel/SpawnLogModule.java | 16 +- .../lib/exec/CompactSpawnLogContext.java | 35 +++- .../build/lib/exec/ExecutionOptions.java | 11 + .../io/AsynchronousMessageOutputStream.java | 14 ++ src/main/protobuf/BUILD | 16 ++ src/main/protobuf/execution_log_stream.proto | 48 +++++ .../AsynchronousMessageOutputStreamTest.java | 22 ++ 9 files changed, 352 insertions(+), 5 deletions(-) create mode 100644 src/main/java/com/google/devtools/build/lib/bazel/ExecutionLogGrpcOutputStream.java create mode 100644 src/main/protobuf/execution_log_stream.proto diff --git a/src/main/java/com/google/devtools/build/lib/bazel/BUILD b/src/main/java/com/google/devtools/build/lib/bazel/BUILD index d710a9fd23c8cd..fb28a757956c46 100644 --- a/src/main/java/com/google/devtools/build/lib/bazel/BUILD +++ b/src/main/java/com/google/devtools/build/lib/bazel/BUILD @@ -121,6 +121,7 @@ java_library( java_library( name = "spawn_log_module", srcs = [ + "ExecutionLogGrpcOutputStream.java", "SpawnLogModule.java", ], deps = [ @@ -139,8 +140,12 @@ java_library( "//src/main/java/com/google/devtools/build/lib/vfs:output_service", "//src/main/java/com/google/devtools/build/lib/vfs:pathfragment", "//src/main/protobuf:failure_details_java_proto", + "//src/main/protobuf:execution_log_stream_java_grpc", + "//src/main/protobuf:execution_log_stream_java_proto", "//third_party:guava", "//third_party:jsr305", + "//third_party/grpc-java:grpc-jar", + "@com_google_protobuf//:protobuf_java", ], ) diff --git a/src/main/java/com/google/devtools/build/lib/bazel/ExecutionLogGrpcOutputStream.java b/src/main/java/com/google/devtools/build/lib/bazel/ExecutionLogGrpcOutputStream.java new file mode 100644 index 00000000000000..cac62fd2e1bfd9 --- /dev/null +++ b/src/main/java/com/google/devtools/build/lib/bazel/ExecutionLogGrpcOutputStream.java @@ -0,0 +1,190 @@ +// Copyright 2026 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 + +package com.google.devtools.build.lib.bazel; + +import com.google.common.hash.Hasher; +import com.google.common.hash.Hashing; +import com.google.devtools.build.lib.exec.ExecutionLogServiceGrpc; +import com.google.devtools.build.lib.exec.ExecutionLogStreamRequest; +import com.google.devtools.build.lib.exec.ExecutionLogStreamResponse; +import com.google.devtools.build.lib.exec.Finish; +import com.google.devtools.build.lib.exec.Header; +import com.google.protobuf.ByteString; +import io.grpc.ManagedChannel; +import io.grpc.ManagedChannelBuilder; +import io.grpc.stub.StreamObserver; +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.io.OutputStream; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +/** Streams compact execution-log bytes to a dedicated gRPC service as they are produced. */ +final class ExecutionLogGrpcOutputStream extends OutputStream { + private static final long CLOSE_TIMEOUT_SECONDS = 30; + private static final int MAX_CHUNK_BYTES = 64 * 1024; + private static final long MAX_CHUNK_DELAY_MILLIS = 100; + + private final ManagedChannel channel; + private final StreamObserver requests; + private final CountDownLatch responseDone = new CountDownLatch(1); + private final AtomicReference failure = new AtomicReference<>(); + private final AtomicReference response = new AtomicReference<>(); + private final Hasher hasher = Hashing.sha256().newHasher(); + private final ByteArrayOutputStream pending = new ByteArrayOutputStream(MAX_CHUNK_BYTES); + private final ScheduledExecutorService flusher = + Executors.newSingleThreadScheduledExecutor( + runnable -> { + Thread thread = new Thread(runnable, "execution-log-grpc-flusher"); + thread.setDaemon(true); + return thread; + }); + private long sizeBytes; + private boolean closed; + + ExecutionLogGrpcOutputStream(String endpoint, String invocationId, String logName) { + String target = endpoint.startsWith("grpc://") ? endpoint.substring("grpc://".length()) : endpoint; + channel = ManagedChannelBuilder.forTarget(target).usePlaintext().build(); + requests = + ExecutionLogServiceGrpc.newStub(channel) + .stream( + new StreamObserver<>() { + @Override + public void onNext(ExecutionLogStreamResponse value) { + response.set(value); + } + + @Override + public void onError(Throwable error) { + failure.compareAndSet(null, error); + responseDone.countDown(); + } + + @Override + public void onCompleted() { + responseDone.countDown(); + } + }); + requests.onNext( + ExecutionLogStreamRequest.newBuilder() + .setHeader( + Header.newBuilder() + .setInvocationId(invocationId) + .setLogName(logName) + .setFormat("compact") + .setCompression("zstd")) + .build()); + flusher.scheduleAtFixedRate( + this::flushFromTimer, + MAX_CHUNK_DELAY_MILLIS, + MAX_CHUNK_DELAY_MILLIS, + TimeUnit.MILLISECONDS); + } + + @Override + public synchronized void write(int value) throws IOException { + write(new byte[] {(byte) value}, 0, 1); + } + + @Override + public synchronized void write(byte[] data, int offset, int length) throws IOException { + if (closed) { + throw new IOException("execution-log stream is closed"); + } + throwIfFailed(); + if (length == 0) { + return; + } + hasher.putBytes(data, offset, length); + sizeBytes += length; + pending.write(data, offset, length); + if (pending.size() >= MAX_CHUNK_BYTES) { + flushPending(); + } + } + + /** + * The compact-log writer flushes after each record so zstd emits bytes promptly. Keep those + * bytes locally until either a small time bound or chunk size is reached. + */ + @Override + public synchronized void flush() throws IOException { + throwIfFailed(); + } + + @Override + public synchronized void close() throws IOException { + if (closed) { + return; + } + closed = true; + try { + throwIfFailed(); + flushPending(); + String digest = hasher.hash().toString(); + requests.onNext( + ExecutionLogStreamRequest.newBuilder() + .setFinish(Finish.newBuilder().setSizeBytes(sizeBytes).setSha256(digest)) + .build()); + requests.onCompleted(); + if (!responseDone.await(CLOSE_TIMEOUT_SECONDS, TimeUnit.SECONDS)) { + throw new IOException("timed out waiting for execution-log stream response"); + } + throwIfFailed(); + ExecutionLogStreamResponse result = response.get(); + if (result == null + || result.getCommittedSize() != sizeBytes + || !result.getSha256().equals(digest)) { + throw new IOException("execution-log stream response failed integrity validation"); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IOException("interrupted while closing execution-log stream", e); + } finally { + flusher.shutdownNow(); + channel.shutdownNow(); + } + } + + private synchronized void flushFromTimer() { + if (closed || pending.size() == 0) { + return; + } + try { + flushPending(); + } catch (IOException e) { + failure.compareAndSet(null, e); + } + } + + private void flushPending() throws IOException { + if (pending.size() == 0) { + return; + } + byte[] data = pending.toByteArray(); + pending.reset(); + try { + requests.onNext( + ExecutionLogStreamRequest.newBuilder().setData(ByteString.copyFrom(data)).build()); + } catch (RuntimeException e) { + failure.compareAndSet(null, e); + throw new IOException("failed to stream execution log", e); + } + } + + private void throwIfFailed() throws IOException { + Throwable error = failure.get(); + if (error != null) { + throw new IOException("execution-log stream failed", error); + } + } +} diff --git a/src/main/java/com/google/devtools/build/lib/bazel/SpawnLogModule.java b/src/main/java/com/google/devtools/build/lib/bazel/SpawnLogModule.java index 2dfce7f7124e50..368b4343a695e7 100644 --- a/src/main/java/com/google/devtools/build/lib/bazel/SpawnLogModule.java +++ b/src/main/java/com/google/devtools/build/lib/bazel/SpawnLogModule.java @@ -136,6 +136,16 @@ private void initOutputs(CommandEnvironment env) throws IOException { outputPath = getAbsolutePath(logPath, env); outputStream = new BufferedOutputStream(outputPath.getOutputStream(), OUTPUT_BUFFER_SIZE); displayName = outputPath.toString(); + } else if (executionOptions.executionLogCompactFile != null + && !executionOptions.executionLogStreamEndpoint.isEmpty()) { + outputStream = + new BufferedOutputStream( + new ExecutionLogGrpcOutputStream( + executionOptions.executionLogStreamEndpoint, + env.getCommandId().toString(), + logName), + OUTPUT_BUFFER_SIZE); + displayName = logName + "-live-stream"; } else if (bepOptions.streamingLogFileUploads) { // Path is empty but streaming is enabled. BuildEventArtifactUploader uploader = @@ -157,7 +167,8 @@ private void initOutputs(CommandEnvironment env) throws IOException { FailureDetail.newBuilder() .setMessage( "--execution_log_{compact,binary,json}_file is empty, but" - + " --experimental_stream_log_file_uploads is not enabled." + + " neither --experimental_execution_log_stream_endpoint nor" + + " --experimental_stream_log_file_uploads is enabled." + " Execution log will not be uploaded to the BEP.") .setExecutionOptions( FailureDetails.ExecutionOptions.newBuilder() @@ -194,7 +205,8 @@ private void initOutputs(CommandEnvironment env) throws IOException { env.getRuntime().getFileSystem().getDigestFunction(), xattrProvider, env.getCommandId(), - env.getReporter()); + env.getReporter(), + !executionOptions.executionLogStreamEndpoint.isEmpty() && logPath.isEmpty()); } else { boolean binaryElseJson = executionOptions.executionLogBinaryFile != null; // Use a well-known temporary path to avoid accumulation of potentially large files in /tmp diff --git a/src/main/java/com/google/devtools/build/lib/exec/CompactSpawnLogContext.java b/src/main/java/com/google/devtools/build/lib/exec/CompactSpawnLogContext.java index ad052764758793..14f79329849823 100644 --- a/src/main/java/com/google/devtools/build/lib/exec/CompactSpawnLogContext.java +++ b/src/main/java/com/google/devtools/build/lib/exec/CompactSpawnLogContext.java @@ -187,6 +187,33 @@ public CompactSpawnLogContext( UUID invocationId, ExtendedEventHandler reporter) throws IOException, InterruptedException { + this( + out, + displayName, + execRoot, + workspaceName, + siblingRepositoryLayout, + remoteOptions, + digestHashFunction, + xattrProvider, + invocationId, + reporter, + /* flushAfterWrite= */ false); + } + + public CompactSpawnLogContext( + BufferedOutputStream out, + String displayName, + PathFragment execRoot, + String workspaceName, + boolean siblingRepositoryLayout, + @Nullable RemoteOptions remoteOptions, + DigestHashFunction digestHashFunction, + XattrProvider xattrProvider, + UUID invocationId, + ExtendedEventHandler reporter, + boolean flushAfterWrite) + throws IOException, InterruptedException { this.execRoot = execRoot; this.workspaceName = workspaceName; this.siblingRepositoryLayout = siblingRepositoryLayout; @@ -195,16 +222,18 @@ public CompactSpawnLogContext( this.xattrProvider = xattrProvider; this.invocationId = invocationId; this.reporter = reporter; - this.outputStream = getOutputStream(out, displayName); + this.outputStream = getOutputStream(out, displayName, flushAfterWrite); logInvocation(); } - private static MessageOutputStream getOutputStream(OutputStream out, String name) + private static MessageOutputStream getOutputStream( + OutputStream out, String name, boolean flushAfterWrite) throws IOException { // Use an AsynchronousMessageOutputStream so that compression and I/O occur in a separate // thread. This ensures concurrent writes don't tear and avoids blocking execution. - return new AsynchronousMessageOutputStream<>(name, new ZstdOutputStream(out)); + return new AsynchronousMessageOutputStream<>( + name, new ZstdOutputStream(out), flushAfterWrite); } private void logInvocation() throws IOException, InterruptedException { diff --git a/src/main/java/com/google/devtools/build/lib/exec/ExecutionOptions.java b/src/main/java/com/google/devtools/build/lib/exec/ExecutionOptions.java index 54201c74b21e7a..e688794ad17d48 100644 --- a/src/main/java/com/google/devtools/build/lib/exec/ExecutionOptions.java +++ b/src/main/java/com/google/devtools/build/lib/exec/ExecutionOptions.java @@ -512,6 +512,17 @@ public boolean usingLocalTestJobs() { + " terminal output).") public PathFragment executionLogCompactFile; + @Option( + name = "experimental_execution_log_stream_endpoint", + defaultValue = "", + documentationCategory = OptionDocumentationCategory.UNDOCUMENTED, + effectTags = {OptionEffectTag.UNKNOWN}, + help = + "A gRPC endpoint that receives compact execution-log bytes while actions are still" + + " executing. Use with --execution_log_compact_file=true. This is independent of" + + " BEP artifact uploads.") + public String executionLogStreamEndpoint; + @Option( name = "execution_log_sort", defaultValue = "true", diff --git a/src/main/java/com/google/devtools/build/lib/util/io/AsynchronousMessageOutputStream.java b/src/main/java/com/google/devtools/build/lib/util/io/AsynchronousMessageOutputStream.java index 7620ce64ef1120..2da86d48bad651 100644 --- a/src/main/java/com/google/devtools/build/lib/util/io/AsynchronousMessageOutputStream.java +++ b/src/main/java/com/google/devtools/build/lib/util/io/AsynchronousMessageOutputStream.java @@ -57,6 +57,17 @@ public AsynchronousMessageOutputStream(Path path) throws IOException { } public AsynchronousMessageOutputStream(String name, OutputStream out) { + this(name, out, /* flushAfterWrite= */ false); + } + + /** + * Creates an asynchronous message stream. + * + *

{@code flushAfterWrite} is useful for live transports where complete messages should + * become visible to the receiver before the stream is closed. + */ + public AsynchronousMessageOutputStream( + String name, OutputStream out, boolean flushAfterWrite) { writerThread = new Thread( () -> { @@ -64,6 +75,9 @@ public AsynchronousMessageOutputStream(String name, OutputStream out) { byte[] data; while ((data = queue.take()) != POISON_PILL) { out.write(data); + if (flushAfterWrite) { + out.flush(); + } } } catch (InterruptedException e) { // Exit quietly. diff --git a/src/main/protobuf/BUILD b/src/main/protobuf/BUILD index f8fce154cb489e..971191854bb412 100644 --- a/src/main/protobuf/BUILD +++ b/src/main/protobuf/BUILD @@ -202,6 +202,22 @@ java_grpc_library( deps = [":command_server_java_proto"], ) +proto_library( + name = "execution_log_stream_proto", + srcs = ["execution_log_stream.proto"], +) + +java_proto_library( + name = "execution_log_stream_java_proto", + deps = [":execution_log_stream_proto"], +) + +java_grpc_library( + name = "execution_log_stream_java_grpc", + srcs = [":execution_log_stream_proto"], + deps = [":execution_log_stream_java_proto"], +) + cc_proto_library( name = "command_server_cc_proto", deps = [":command_server_proto"], diff --git a/src/main/protobuf/execution_log_stream.proto b/src/main/protobuf/execution_log_stream.proto new file mode 100644 index 00000000000000..f8a61cc71e7c51 --- /dev/null +++ b/src/main/protobuf/execution_log_stream.proto @@ -0,0 +1,48 @@ +// Copyright 2026 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 + +syntax = "proto3"; + +package build.bazel.executionlog.v1; + +option java_multiple_files = true; +option java_package = "com.google.devtools.build.lib.exec"; + +// Receives an execution log while Bazel is still executing actions. +// +// The data messages contain consecutive bytes of the zstd-compressed compact +// execution log. The server can decode complete ExecLogEntry records from the +// prefix received so far. The final digest covers exactly the data bytes. +service ExecutionLogService { + rpc Stream(stream ExecutionLogStreamRequest) returns (ExecutionLogStreamResponse); +} + +message ExecutionLogStreamRequest { + oneof payload { + Header header = 1; + bytes data = 2; + Finish finish = 3; + } +} + +message Header { + string invocation_id = 1; + string log_name = 2; + string format = 3; + string compression = 4; +} + +message Finish { + int64 size_bytes = 1; + string sha256 = 2; +} + +message ExecutionLogStreamResponse { + int64 committed_size = 1; + string sha256 = 2; +} diff --git a/src/test/java/com/google/devtools/build/lib/util/io/AsynchronousMessageOutputStreamTest.java b/src/test/java/com/google/devtools/build/lib/util/io/AsynchronousMessageOutputStreamTest.java index 7cb4e90e4159cd..6eec7b7dd5ada7 100644 --- a/src/test/java/com/google/devtools/build/lib/util/io/AsynchronousMessageOutputStreamTest.java +++ b/src/test/java/com/google/devtools/build/lib/util/io/AsynchronousMessageOutputStreamTest.java @@ -30,6 +30,7 @@ import java.util.Random; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.atomic.AtomicInteger; import org.junit.After; import org.junit.Before; import org.junit.Test; @@ -203,4 +204,25 @@ public void testWriteAfterCloseThrowsException() throws Exception { assertThrows(IllegalStateException.class, () -> out.write(generateRandomMessage())); } + + @Test + public void testFlushAfterWriteMakesMessagesVisibleBeforeClose() throws Exception { + AtomicInteger flushes = new AtomicInteger(); + OutputStream recordingOutputStream = + new ByteArrayOutputStream() { + @Override + public void flush() { + flushes.incrementAndGet(); + } + }; + AsynchronousMessageOutputStream out = + new AsynchronousMessageOutputStream<>( + "", recordingOutputStream, /* flushAfterWrite= */ true); + + out.write(generateRandomMessage()); + out.write(generateRandomMessage()); + out.close(); + + assertThat(flushes.get()).isAtLeast(2); + } } From 2b0c976e7bbb9076fdc01f5e47f8f829eea8fe42 Mon Sep 17 00:00:00 2001 From: Olivier Notteghem Date: Wed, 16 Sep 2026 04:56:00 +0000 Subject: [PATCH 2/3] Allow fork release builds without BuildBuddy secret --- .github/workflows/build-and-publish-bazel.yml | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) diff --git a/.github/workflows/build-and-publish-bazel.yml b/.github/workflows/build-and-publish-bazel.yml index a9d6dbc5af22f6..2f65b06b14a22e 100644 --- a/.github/workflows/build-and-publish-bazel.yml +++ b/.github/workflows/build-and-publish-bazel.yml @@ -132,13 +132,17 @@ jobs: - name: Build Bazel run: | set -euo pipefail - if [[ -z "${BUILDBUDDY_API_KEY}" ]]; then - echo 'Missing the BUILDBUDDY_API_KEY Actions secret.' >&2 - exit 1 + remote_args=() + if [[ -n "${BUILDBUDDY_API_KEY}" ]]; then + remote_args+=( + --config=buildbuddy-cache + "--remote_header=x-buildbuddy-api-key=${BUILDBUDDY_API_KEY}" + ) + else + echo 'BUILDBUDDY_API_KEY is unavailable; building without the remote cache.' fi ./bazelisk --batch --nohome_rc --nosystem_rc build \ - --config=buildbuddy-cache \ - "--remote_header=x-buildbuddy-api-key=${BUILDBUDDY_API_KEY}" \ + "${remote_args[@]}" \ --verbose_failures \ --compilation_mode=opt \ --jobs=2 \ From e91368ba0acedc4e5906af8ce9af351ad4a5436d Mon Sep 17 00:00:00 2001 From: Olivier Notteghem Date: Wed, 16 Sep 2026 04:57:39 +0000 Subject: [PATCH 3/3] Make no-cache release builds portable --- .github/workflows/build-and-publish-bazel.yml | 21 ++++++++++--------- 1 file changed, 11 insertions(+), 10 deletions(-) diff --git a/.github/workflows/build-and-publish-bazel.yml b/.github/workflows/build-and-publish-bazel.yml index 2f65b06b14a22e..c4e38a9014c03b 100644 --- a/.github/workflows/build-and-publish-bazel.yml +++ b/.github/workflows/build-and-publish-bazel.yml @@ -132,23 +132,24 @@ jobs: - name: Build Bazel run: | set -euo pipefail - remote_args=() + build_args=( + --verbose_failures + --compilation_mode=opt + --jobs=2 + --stamp + "--embed_label=${RELEASE_TAG}" + //src:bazel + ) if [[ -n "${BUILDBUDDY_API_KEY}" ]]; then - remote_args+=( + build_args=( --config=buildbuddy-cache "--remote_header=x-buildbuddy-api-key=${BUILDBUDDY_API_KEY}" + "${build_args[@]}" ) else echo 'BUILDBUDDY_API_KEY is unavailable; building without the remote cache.' fi - ./bazelisk --batch --nohome_rc --nosystem_rc build \ - "${remote_args[@]}" \ - --verbose_failures \ - --compilation_mode=opt \ - --jobs=2 \ - --stamp \ - "--embed_label=${RELEASE_TAG}" \ - //src:bazel + ./bazelisk --batch --nohome_rc --nosystem_rc build "${build_args[@]}" - name: Verify and package Bazel env: