Skip to content
Closed
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
243 changes: 58 additions & 185 deletions api/src/main/java/com/google/appengine/setup/AppLogsWriter.java
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.locks.ReentrantLock;
import java.util.logging.Level;
import java.util.logging.Logger;

Expand Down Expand Up @@ -62,8 +63,10 @@
* from truncating individual log entries.
* <p/>
* <p>This class is thread safe and all methods accessing local state are
* synchronized. Since each request have their own instance of this class the
* only contention possible is between the original request thread and and any
* guarded by a {@link ReentrantLock} (rather than {@code synchronized}, so that a
* virtual thread blocked on a flush does not pin its carrier thread on Java 21).
* Since each request have their own instance of this class the
* only contention possible is between the original request thread and any
* child RequestThreads created by the request through the threading API.
*/
class AppLogsWriter {
Expand All @@ -87,6 +90,7 @@ class AppLogsWriter {
private int flushCount = 0;
private Future<byte[]> currentFlush;
private Stopwatch stopwatch;
private final ReentrantLock lock = new ReentrantLock();

/**
* Construct an AppLogsWriter instance.
Expand Down Expand Up @@ -147,99 +151,33 @@ public AppLogsWriter(List<UserAppLogLine> buffer, long maxBytesToFlush, int maxL
* this method may block.
*/
void addLogRecordAndMaybeFlush(LogRecord fullRecord) {
if (Boolean.getBoolean("appengine.use.virtualthreads")) {
addLogRecordAndMaybeFlushVirtualThreads(fullRecord);
} else {
addLogRecordAndMaybeFlushLegacy(fullRecord);
}
}

private void addLogRecordAndMaybeFlushVirtualThreads(LogRecord fullRecord) {
for (LogRecord record : split(fullRecord)) {
UserAppLogLine logLine = UserAppLogLine.newBuilder()
.setLevel(record.getLevel().ordinal())
.setTimestampUsec(record.getTimestamp())
.setMessage(record.getMessage())
.build();
int maxEncodingSize = 1000; // logLine.maxEncodingSize();
Future<byte[]> pendingFlush = null;
synchronized (this) {
lock.lock();
try {
for (LogRecord record : split(fullRecord)) {
UserAppLogLine logLine = UserAppLogLine.newBuilder()
.setLevel(record.getLevel().ordinal())
.setTimestampUsec(record.getTimestamp())
.setMessage(record.getMessage())
.build();
int maxEncodingSize = 1000; // logLine.maxEncodingSize();
if (maxBytesToFlush > 0 &&
(currentByteCount + maxEncodingSize) > maxBytesToFlush) {
pendingFlush = getPendingFlushLocked();
if (pendingFlush == null && buffer.size() > 0) {
currentFlush = doFlush();
}
}
}
if (pendingFlush != null) {
waitForFlush(pendingFlush);
synchronized (this) {
if (currentFlush == null || currentFlush.isDone()) {
if (buffer.size() > 0) {
currentFlush = doFlush();
} else if (currentFlush != null && currentFlush.isDone()) {
currentFlush = null;
}
}
logger.info(currentByteCount + " bytes of app logs pending, starting flush...");
waitForCurrentFlushAndStartNewFlush();
}
}
synchronized (this) {
if (buffer.size() == 0) {
stopwatch.start();
}
buffer.add(logLine);
currentByteCount += maxEncodingSize;
}
}

Future<byte[]> pendingTimeFlush = null;
synchronized (this) {
if (maxSecondsBetweenFlush > 0 &&
stopwatch.elapsed(TimeUnit.SECONDS) >= maxSecondsBetweenFlush) {
pendingTimeFlush = getPendingFlushLocked();
if (pendingTimeFlush == null && buffer.size() > 0) {
currentFlush = doFlush();
}
}
}
if (pendingTimeFlush != null) {
waitForFlush(pendingTimeFlush);
synchronized (this) {
if (currentFlush == null || currentFlush.isDone()) {
if (buffer.size() > 0) {
currentFlush = doFlush();
} else if (currentFlush != null && currentFlush.isDone()) {
currentFlush = null;
}
}
waitForCurrentFlushAndStartNewFlush();
}
}
}

private synchronized void addLogRecordAndMaybeFlushLegacy(LogRecord fullRecord) {
for (LogRecord record : split(fullRecord)) {
UserAppLogLine logLine = UserAppLogLine.newBuilder()
.setLevel(record.getLevel().ordinal())
.setTimestampUsec(record.getTimestamp())
.setMessage(record.getMessage())
.build();
int maxEncodingSize = 1000; // logLine.maxEncodingSize();
if (maxBytesToFlush > 0 &&
(currentByteCount + maxEncodingSize) > maxBytesToFlush) {
logger.info(currentByteCount + " bytes of app logs pending, starting flush...");
waitForCurrentFlushAndStartNewFlushLegacy();
}
if (buffer.size() == 0) {
stopwatch.start();
}
buffer.add(logLine);
currentByteCount += maxEncodingSize;
}

if (maxSecondsBetweenFlush > 0 &&
stopwatch.elapsed(TimeUnit.SECONDS) >= maxSecondsBetweenFlush) {
waitForCurrentFlushAndStartNewFlushLegacy();
} finally {
lock.unlock();
}
}

Expand All @@ -249,129 +187,64 @@ private synchronized void addLogRecordAndMaybeFlushLegacy(LogRecord fullRecord)
*
* @return The number of times this AppLogsWriter has initiated a flush.
*/
synchronized int waitForCurrentFlushAndStartNewFlush() {
if (Boolean.getBoolean("appengine.use.virtualthreads")) {
Future<byte[]> pending = getPendingFlushLocked();
if (pending != null) {
waitForFlush(pending);
}
int waitForCurrentFlushAndStartNewFlush() {
lock.lock();
try {
waitForCurrentFlush();
if (buffer.size() > 0) {
currentFlush = doFlush();
}
return flushCount;
} else {
return waitForCurrentFlushAndStartNewFlushLegacy();
}
}

private synchronized int waitForCurrentFlushAndStartNewFlushLegacy() {
waitForCurrentFlushLegacy();
if (buffer.size() > 0) {
currentFlush = doFlush();
} finally {
lock.unlock();
}
return flushCount;
}

/**
* Initiates a synchronous flush. This method will always block until any pending flushes and
* its own flush completes.
*
* <p>When {@code appengine.use.virtualthreads} is enabled, the actual I/O wait on {@link
* Future#get()} is performed outside of the {@code synchronized} monitor lock to allow virtual
* threads to unmount without pinning carrier threads. Otherwise, it follows legacy synchronized locking.
* Initiates a synchronous flush. This method will always block
* until any pending flushes and its own flush completes.
*/
void flushAndWait() {
if (Boolean.getBoolean("appengine.use.virtualthreads")) {
flushAndWaitVirtualThreads();
} else {
flushAndWaitLegacy();
}
}

private void flushAndWaitVirtualThreads() {
Future<byte[]> previousFlush;
synchronized (this) {
previousFlush = getPendingFlushLocked();
}
if (previousFlush != null) {
waitForFlush(previousFlush);
}

Future<byte[]> flush = null;
synchronized (this) {
if (currentFlush == null || currentFlush.isDone()) {
if (buffer.size() > 0) {
flush = currentFlush = doFlush();
} else if (currentFlush != null && currentFlush.isDone()) {
currentFlush = null;
}
} else {
flush = currentFlush;
}
}
if (flush != null) {
waitForFlush(flush);
synchronized (this) {
if (currentFlush != null && currentFlush.isDone()) {
currentFlush = null;
}
}
}
}

private synchronized void flushAndWaitLegacy() {
waitForCurrentFlushLegacy();
if (buffer.size() > 0) {
currentFlush = doFlush();
waitForCurrentFlushLegacy();
}
}

private void waitForFlush(Future<byte[]> flush) {
lock.lock();
try {
flush.get(
ApiProxyDelegate.ADDITIONAL_HTTP_TIMEOUT_BUFFER_MS + LOG_FLUSH_TIMEOUT_MS,
TimeUnit.MILLISECONDS);
} catch (InterruptedException ex) {
logger.warning("Interrupted while blocking on a log flush, setting interrupt bit and " +
"continuing. Some logs may be lost or occur out of order!");
Thread.currentThread().interrupt();
} catch (TimeoutException e) {
logger.log(Level.WARNING, "Timeout waiting for log flush to complete. "
+ "Log messages may have been lost/reordered!", e);
} catch (ExecutionException ex) {
logger.log(
Level.WARNING,
"A log flush request failed. Log messages may have been lost!", ex);
waitForCurrentFlush();
if (buffer.size() > 0) {
currentFlush = doFlush();
waitForCurrentFlush();
}
} finally {
lock.unlock();
}
}

private void waitForCurrentFlushLegacy() {
/**
* This method blocks until any outstanding flush is completed. This method
* should be called prior to {@link #doFlush()} so that it is impossible for
* the appserver to process logs out of order.
*/
private void waitForCurrentFlush() {
if (currentFlush != null) {
logger.info("Previous flush has not yet completed, blocking.");
waitForFlush(currentFlush);
try {
currentFlush.get(
ApiProxyDelegate.ADDITIONAL_HTTP_TIMEOUT_BUFFER_MS + LOG_FLUSH_TIMEOUT_MS,
TimeUnit.MILLISECONDS);
} catch (InterruptedException ex) {
logger.warning("Interrupted while blocking on a log flush, setting interrupt bit and " +
"continuing. Some logs may be lost or occur out of order!");
Thread.currentThread().interrupt();
} catch (TimeoutException e) {
logger.log(Level.WARNING, "Timeout waiting for log flush to complete. "
+ "Log messages may have been lost/reordered!", e);
} catch (ExecutionException ex) {
logger.log(
Level.WARNING,
"A log flush request failed. Log messages may have been lost!", ex);
}
currentFlush = null;
}
}

/**
* Returns the currently pending flush {@link Future} if it has not yet completed.
*
* <p>By retrieving the pending flush under {@code synchronized (this)} and returning it to the
* caller without nullifying it right away, we allow {@link #waitForFlush(Future)} (which invokes {@link
* Future#get()}) to be executed strictly outside the synchronized monitor block while ensuring
* other virtual threads see that a flush is still pending. Under Java 21 (pre-JEP 491),
* blocking inside a synchronized scope prevents virtual threads from unmounting and pins their
* carrier threads, leading to pool starvation across the web container.
*/
private synchronized Future<byte[]> getPendingFlushLocked() {
if (currentFlush != null && !currentFlush.isDone() && !currentFlush.isCancelled()) {
return currentFlush;
}
currentFlush = null;
return null;
}

private Future<byte[]> doFlush() {
UserAppLogGroup.Builder group = UserAppLogGroup.newBuilder();
for (UserAppLogLine logLine : buffer) {
Expand Down
Loading
Loading