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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 17 additions & 12 deletions .github/workflows/build-and-publish-bazel.yml
Original file line number Diff line number Diff line change
Expand Up @@ -132,19 +132,24 @@ 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
fi
./bazelisk --batch --nohome_rc --nosystem_rc build \
--config=buildbuddy-cache \
"--remote_header=x-buildbuddy-api-key=${BUILDBUDDY_API_KEY}" \
--verbose_failures \
--compilation_mode=opt \
--jobs=2 \
--stamp \
"--embed_label=${RELEASE_TAG}" \
build_args=(
--verbose_failures
--compilation_mode=opt
--jobs=2
--stamp
"--embed_label=${RELEASE_TAG}"
//src:bazel
)
if [[ -n "${BUILDBUDDY_API_KEY}" ]]; then
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 "${build_args[@]}"

- name: Verify and package Bazel
env:
Expand Down
5 changes: 5 additions & 0 deletions src/main/java/com/google/devtools/build/lib/bazel/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,7 @@ java_library(
java_library(
name = "spawn_log_module",
srcs = [
"ExecutionLogGrpcOutputStream.java",
"SpawnLogModule.java",
],
deps = [
Expand All @@ -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",
],
)

Expand Down
Original file line number Diff line number Diff line change
@@ -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<ExecutionLogStreamRequest> requests;
private final CountDownLatch responseDone = new CountDownLatch(1);
private final AtomicReference<Throwable> failure = new AtomicReference<>();
private final AtomicReference<ExecutionLogStreamResponse> 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);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand All @@ -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()
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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<ExecLogEntry> getOutputStream(OutputStream out, String name)
private static MessageOutputStream<ExecLogEntry> 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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,13 +57,27 @@ public AsynchronousMessageOutputStream(Path path) throws IOException {
}

public AsynchronousMessageOutputStream(String name, OutputStream out) {
this(name, out, /* flushAfterWrite= */ false);
}

/**
* Creates an asynchronous message stream.
*
* <p>{@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(
() -> {
try {
byte[] data;
while ((data = queue.take()) != POISON_PILL) {
out.write(data);
if (flushAfterWrite) {
out.flush();
}
}
} catch (InterruptedException e) {
// Exit quietly.
Expand Down
Loading
Loading