diff --git a/api/src/main/java/com/google/appengine/api/datastore/DatastoreServiceImpl.java b/api/src/main/java/com/google/appengine/api/datastore/DatastoreServiceImpl.java index d0444ae6b..f13a4b486 100644 --- a/api/src/main/java/com/google/appengine/api/datastore/DatastoreServiceImpl.java +++ b/api/src/main/java/com/google/appengine/api/datastore/DatastoreServiceImpl.java @@ -18,18 +18,24 @@ import static com.google.appengine.api.datastore.FutureHelper.quietGet; +import com.google.apphosting.api.ApiProxy; import java.util.Collection; import java.util.List; import java.util.Map; +import java.util.concurrent.Future; /** * An implementation of {@link DatastoreService} that farms out all calls to a provided {@link * AsyncDatastoreService}. - * */ final class DatastoreServiceImpl implements DatastoreService { private final AsyncDatastoreServiceInternal async; + static final long BEGIN_TXN_RETRY_DELAY_MS = 100; + + private static int getMaxRetries() { + return Math.max(0, Integer.getInteger("appengine.datastore.retries", 1)); + } public DatastoreServiceImpl(AsyncDatastoreServiceInternal async) { this.async = async; @@ -137,12 +143,42 @@ public KeyRangeState allocateIdRange(KeyRange range) { @Override public Transaction beginTransaction() { - return quietGet(async.beginTransaction()); + return beginTransaction(TransactionOptions.Builder.withDefaults()); } @Override public Transaction beginTransaction(TransactionOptions options) { - return quietGet(async.beginTransaction(options)); + int retries = 0; + int maxRetries = getMaxRetries(); + long delay = BEGIN_TXN_RETRY_DELAY_MS; + while (true) { + Transaction tx = null; + try { + tx = quietGet(async.beginTransaction(options)); + tx.getId(); // Force handle resolution + return tx; + } catch (DatastoreFailureException + | DatastoreTimeoutException + | ApiProxy.RPCFailedException e) { + if (tx != null) { + try { + Future unused = tx.rollbackAsync(); + } catch (RuntimeException ignored) { + // Best-effort rollback; original exception is retried or rethrown. + } + } + if (++retries > maxRetries) { + throw e; + } + try { + Thread.sleep(delay); + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + throw e; + } + delay *= 2; + } + } } @Override diff --git a/api/src/main/java/com/google/appengine/setup/AppLogsWriter.java b/api/src/main/java/com/google/appengine/setup/AppLogsWriter.java index ecc4b49aa..46ea0cc3e 100644 --- a/api/src/main/java/com/google/appengine/setup/AppLogsWriter.java +++ b/api/src/main/java/com/google/appengine/setup/AppLogsWriter.java @@ -16,55 +16,60 @@ package com.google.appengine.setup; -import com.google.common.base.Stopwatch; -import com.google.protobuf.ByteString; import com.google.apphosting.api.ApiProxy; import com.google.apphosting.api.ApiProxy.ApiConfig; import com.google.apphosting.api.ApiProxy.LogRecord; import com.google.apphosting.api.logservice.LogServicePb.FlushRequest; import com.google.apphosting.api.logservice.LogServicePb.UserAppLogGroup; import com.google.apphosting.api.logservice.LogServicePb.UserAppLogLine; +import com.google.common.base.Stopwatch; +import com.google.protobuf.ByteString; import java.util.LinkedList; import java.util.List; import java.util.concurrent.ExecutionException; 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; /** - * {@code AppsLogWriter} is responsible for batching application logs - * for a single request and sending them back to the AppServer via the - * LogService.Flush API call. - *

+ * {@code AppsLogWriter} is responsible for batching application logs for a single request and + * sending them back to the AppServer via the LogService.Flush API call. + * + *

+ * *

The current algorithm used to send logs is as follows: + * *

- *

- *

This class is also responsible for splitting large log entries - * into smaller fragments, which is unrelated to the batching - * mechanism described above but is necessary to prevent the AppServer + * + *

+ * + *

This class is also responsible for splitting large log entries into smaller fragments, which + * is unrelated to the batching mechanism described above but is necessary to prevent the AppServer * from truncating individual log entries. - *

- *

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 - * child RequestThreads created by the request through the threading API. + * + *

+ * + *

This class is thread safe and all methods accessing local state are 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 { private static final Logger logger = @@ -87,6 +92,7 @@ class AppLogsWriter { private int flushCount = 0; private Future currentFlush; private Stopwatch stopwatch; + private final ReentrantLock lock = new ReentrantLock(); /** * Construct an AppLogsWriter instance. @@ -147,231 +153,102 @@ public AppLogsWriter(List 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() + 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(); - Future pendingFlush = null; - synchronized (this) { + 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 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(); } + } finally { + lock.unlock(); } } - 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(); - } - } - - /** - * Starts an asynchronous flush. This method may block if flushes - * are backed up. - * - * @return The number of times this AppLogsWriter has initiated a flush. - */ - synchronized int waitForCurrentFlushAndStartNewFlush() { - if (Boolean.getBoolean("appengine.use.virtualthreads")) { - Future pending = getPendingFlushLocked(); - if (pending != null) { - waitForFlush(pending); - } + /** + * Starts an asynchronous flush. This method may block if flushes are backed up. + * + * @return The number of times this AppLogsWriter has initiated a flush. + */ + int waitForCurrentFlushAndStartNewFlush() { + lock.lock(); + try { + waitForCurrentFlush(); if (buffer.size() > 0) { currentFlush = doFlush(); } return flushCount; - } else { - return waitForCurrentFlushAndStartNewFlushLegacy(); + } finally { + lock.unlock(); } } - private synchronized int waitForCurrentFlushAndStartNewFlushLegacy() { - waitForCurrentFlushLegacy(); - if (buffer.size() > 0) { - currentFlush = doFlush(); - } - return flushCount; - } - - /** - * Initiates a synchronous flush. This method will always block until any pending flushes and - * its own flush completes. - * - *

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. - */ - void flushAndWait() { - if (Boolean.getBoolean("appengine.use.virtualthreads")) { - flushAndWaitVirtualThreads(); - } else { - flushAndWaitLegacy(); - } - } - - private void flushAndWaitVirtualThreads() { - Future previousFlush; - synchronized (this) { - previousFlush = getPendingFlushLocked(); - } - if (previousFlush != null) { - waitForFlush(previousFlush); - } - - Future 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; - } + /** + * Initiates a synchronous flush. This method will always block until any pending flushes and its + * own flush completes. + */ + void flushAndWait() { + lock.lock(); + try { + waitForCurrentFlush(); + if (buffer.size() > 0) { + currentFlush = doFlush(); + waitForCurrentFlush(); } + } finally { + lock.unlock(); } } - private synchronized void flushAndWaitLegacy() { - waitForCurrentFlushLegacy(); - if (buffer.size() > 0) { - currentFlush = doFlush(); - waitForCurrentFlushLegacy(); - } - } - - private void waitForFlush(Future flush) { - 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); - } - } - - 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. - * - *

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 getPendingFlushLocked() { - if (currentFlush != null && !currentFlush.isDone() && !currentFlush.isCancelled()) { - return currentFlush; - } - currentFlush = null; - return null; - } - private Future doFlush() { UserAppLogGroup.Builder group = UserAppLogGroup.newBuilder(); for (UserAppLogLine logLine : buffer) { @@ -385,7 +262,7 @@ private Future doFlush() { request.setLogs(ByteString.copyFrom(group.build().toByteArray())); ApiConfig apiConfig = new ApiConfig(); apiConfig.setDeadlineInSeconds(LOG_FLUSH_TIMEOUT_MS / 1000.0); - return ApiProxy.makeAsyncCall("logservice", "Flush", request.build().toByteArray(), apiConfig); + return ApiProxy.makeAsyncCall("logservice", "Flush", request.build().toByteArray(), apiConfig); } /** diff --git a/runtime/impl/src/main/java/com/google/apphosting/runtime/AppLogsWriter.java b/runtime/impl/src/main/java/com/google/apphosting/runtime/AppLogsWriter.java index d9297db5d..c92f726a2 100644 --- a/runtime/impl/src/main/java/com/google/apphosting/runtime/AppLogsWriter.java +++ b/runtime/impl/src/main/java/com/google/apphosting/runtime/AppLogsWriter.java @@ -30,9 +30,9 @@ import java.util.List; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; +import java.util.concurrent.locks.ReentrantLock; import java.util.logging.Level; import javax.annotation.concurrent.GuardedBy; -import org.jspecify.annotations.Nullable; /** * {@code AppsLogWriter} is responsible for batching application logs for a single request and @@ -76,7 +76,9 @@ public class AppLogsWriter { static final String LOG_TRUNCATED_SUFFIX = "\n"; static final int LOG_TRUNCATED_SUFFIX_LENGTH = LOG_TRUNCATED_SUFFIX.length(); - private final Object lock = new Object(); + // A ReentrantLock rather than synchronized, so that a virtual thread blocked on a flush while + // holding the lock does not pin its carrier thread on Java 21. + private final ReentrantLock lock = new ReentrantLock(); private final int maxLogMessageLength; private final int logCutLength; @@ -181,88 +183,28 @@ public void addLogRecordAndMaybeFlush(ApiProxy.LogRecord fullRecord) { appLogLines.add(logLineBuilder.build()); } - if (Boolean.getBoolean("appengine.use.virtualthreads")) { - addLogLinesAndMaybeFlushVirtualThreads(appLogLines); - } else { - synchronized (lock) { - addLogLinesAndMaybeFlushLegacy(appLogLines); - } - } - } - - private void addLogLinesAndMaybeFlushVirtualThreads(Iterable appLogLines) { - for (AppLogLine logLine : appLogLines) { - int serializedSize = logLine.getSerializedSize(); - - Future pendingFlush = null; - synchronized (lock) { - if (maxBytesToFlush > 0 && (currentByteCount + serializedSize) > maxBytesToFlush) { - pendingFlush = getPendingFlushLocked(); - if (pendingFlush == null && genericResponse.getAppLogCount() > 0) { - currentFlush = doFlush(); - } - } - } - if (pendingFlush != null) { - waitForFlush(pendingFlush); - synchronized (lock) { - if (currentFlush == null || currentFlush.isDone()) { - if (genericResponse.getAppLogCount() > 0) { - currentFlush = doFlush(); - } else if (currentFlush != null && currentFlush.isDone()) { - currentFlush = null; - } - } - } - } - - synchronized (lock) { - if (!stopwatch.isRunning()) { - // We only want to flush once a log message has been around for - // longer than maxSecondsBetweenFlush. So, we only start the timer - // when we add the first message so we don't include time when - // the queue is empty. - stopwatch.start(); - } - genericResponse.addAppLog(logLine); - currentByteCount += serializedSize; - } - } - - Future pendingTimeFlush = null; - synchronized (lock) { - if (maxSecondsBetweenFlush > 0 - && stopwatch.elapsed().compareTo(Duration.ofSeconds(maxSecondsBetweenFlush)) >= 0) { - pendingTimeFlush = getPendingFlushLocked(); - if (pendingTimeFlush == null && genericResponse.getAppLogCount() > 0) { - currentFlush = doFlush(); - } - } - } - if (pendingTimeFlush != null) { - waitForFlush(pendingTimeFlush); - synchronized (lock) { - if (currentFlush == null || currentFlush.isDone()) { - if (genericResponse.getAppLogCount() > 0) { - currentFlush = doFlush(); - } else if (currentFlush != null && currentFlush.isDone()) { - currentFlush = null; - } - } - } + lock.lock(); + try { + addLogLinesAndMaybeFlush(appLogLines); + } finally { + lock.unlock(); } } @GuardedBy("lock") - private void addLogLinesAndMaybeFlushLegacy(Iterable appLogLines) { + private void addLogLinesAndMaybeFlush(Iterable appLogLines) { for (AppLogLine logLine : appLogLines) { int serializedSize = logLine.getSerializedSize(); if (maxBytesToFlush > 0 && (currentByteCount + serializedSize) > maxBytesToFlush) { logger.atInfo().log("%d bytes of app logs pending, starting flush...", currentByteCount); - waitForCurrentFlushAndStartNewFlushLegacy(); + waitForCurrentFlushAndStartNewFlush(); } if (!stopwatch.isRunning()) { + // We only want to flush once a log message has been around for + // longer than maxSecondsBetweenFlush. So, we only start the timer + // when we add the first message so we don't include time when + // the queue is empty. stopwatch.start(); } genericResponse.addAppLog(logLine); @@ -271,111 +213,56 @@ private void addLogLinesAndMaybeFlushLegacy(Iterable appLogLines) { if (maxSecondsBetweenFlush > 0 && stopwatch.elapsed().compareTo(Duration.ofSeconds(maxSecondsBetweenFlush)) >= 0) { - waitForCurrentFlushAndStartNewFlushLegacy(); + waitForCurrentFlushAndStartNewFlush(); } } + /** Starts an asynchronous flush. This method may block if flushes are backed up. */ @GuardedBy("lock") - private void waitForCurrentFlushAndStartNewFlushLegacy() { - waitForCurrentFlushLegacy(); + private void waitForCurrentFlushAndStartNewFlush() { + waitForCurrentFlush(); if (genericResponse.getAppLogCount() > 0) { currentFlush = doFlush(); } } - @GuardedBy("lock") - private void waitForCurrentFlushLegacy() { - if (currentFlush != null && !currentFlush.isDone() && !currentFlush.isCancelled()) { - logger.atInfo().log("Previous flush has not yet completed, blocking."); - waitForFlush(currentFlush); - } - currentFlush = null; - } - - /** - * Returns the currently pending flush {@link Future} if it has not yet completed. - * - *

By retrieving the pending flush under {@code lock} 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 {@code synchronized (lock)} 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 in {@link java.util.concurrent.ForkJoinPool#commonPool()}, leading to - * thread pool starvation across the web container. - */ - @GuardedBy("lock") - private @Nullable Future getPendingFlushLocked() { - if (currentFlush != null && !currentFlush.isDone() && !currentFlush.isCancelled()) { - return currentFlush; - } - currentFlush = null; - return null; - } - /** * Initiates a synchronous flush. This method will always block until any pending flushes and its * own flush completes. - * - *

When {@code appengine.use.virtualthreads} is enabled, the actual I/O wait on {@link - * Future#get()} is performed outside of {@code synchronized (lock)} to allow virtual threads to - * unmount without pinning carrier threads. Otherwise, it follows legacy synchronized locking. */ public void flushAndWait() { - if (Boolean.getBoolean("appengine.use.virtualthreads")) { - flushAndWaitVirtualThreads(); - } else { - flushAndWaitLegacy(); - } - } - - private void flushAndWaitVirtualThreads() { - Future previousFlush; - synchronized (lock) { - previousFlush = getPendingFlushLocked(); - } - if (previousFlush != null) { - waitForFlush(previousFlush); - } - - Future flush = null; - synchronized (lock) { - if (currentFlush == null || currentFlush.isDone()) { - if (genericResponse.getAppLogCount() > 0) { - flush = currentFlush = doFlush(); - } else if (currentFlush != null && currentFlush.isDone()) { - currentFlush = null; - } - } else { - flush = currentFlush; - } - } - - // Wait for this flush outside the synchronized block to allow virtual threads to unmount without pinning. - if (flush != null) { - waitForFlush(flush); - synchronized (lock) { - if (currentFlush != null && currentFlush.isDone()) { - currentFlush = null; - } - } - } - } - - private void flushAndWaitLegacy() { Future flush = null; - synchronized (lock) { - waitForCurrentFlushLegacy(); + lock.lock(); + try { + waitForCurrentFlush(); if (genericResponse.getAppLogCount() > 0) { flush = currentFlush = doFlush(); } + } finally { + lock.unlock(); } + // Wait for this flush outside the lock to avoid unnecessarily blocking + // addLogRecordAndMaybeFlush() calls when flushes are not backed up. if (flush != null) { waitForFlush(flush); } } + /** + * 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. + */ + @GuardedBy("lock") + private void waitForCurrentFlush() { + if (currentFlush != null && !currentFlush.isDone() && !currentFlush.isCancelled()) { + logger.atInfo().log("Previous flush has not yet completed, blocking."); + waitForFlush(currentFlush); + } + currentFlush = null; + } + private void waitForFlush(Future flush) { try { flush.get(); @@ -469,8 +356,11 @@ List split(ApiProxy.LogRecord aRecord) { */ @VisibleForTesting void setStopwatch(Stopwatch stopwatch) { - synchronized (lock) { + lock.lock(); + try { this.stopwatch = stopwatch; + } finally { + lock.unlock(); } } diff --git a/runtime/impl/src/test/java/com/google/apphosting/runtime/AppLogsWriterTest.java b/runtime/impl/src/test/java/com/google/apphosting/runtime/AppLogsWriterTest.java index 8036dc446..056e799d0 100644 --- a/runtime/impl/src/test/java/com/google/apphosting/runtime/AppLogsWriterTest.java +++ b/runtime/impl/src/test/java/com/google/apphosting/runtime/AppLogsWriterTest.java @@ -463,69 +463,10 @@ public void testPreserveLeadingSpaceAtSplit() { } @Test - public void testFlushAndWait_doesNotHoldLockWhileWaitingOnFuture() throws Exception { - System.setProperty("appengine.use.virtualthreads", "true"); - try { - SettableFuture slowFlush = SettableFuture.create(); - when(delegate.makeAsyncCall(eq(environment), eq("logservice"), eq("Flush"), notNull(), notNull())) - .thenReturn(slowFlush); - - AppLogsWriter writer = new AppLogsWriter(response, SMALL_FLUSH, DEFAULT_MAX_LOG_LINE, 0); - writer.addLogRecordAndMaybeFlush(new LogRecord(LogRecord.Level.info, 0, "log message")); - - CountDownLatch firstFlushEntered = new CountDownLatch(1); - CountDownLatch secondFlushWaiting = new CountDownLatch(1); - CountDownLatch thirdThreadAcquiredLock = new CountDownLatch(1); - - Thread firstFlushThread = - new Thread( - () -> { - ApiProxy.setEnvironmentForCurrentThread(environment); - firstFlushEntered.countDown(); - writer.flushAndWait(); - }); - firstFlushThread.start(); - - firstFlushEntered.await(); - Thread.sleep(100); - - Thread secondFlushThread = - new Thread( - () -> { - ApiProxy.setEnvironmentForCurrentThread(environment); - secondFlushWaiting.countDown(); - writer.flushAndWait(); - }); - secondFlushThread.start(); - - secondFlushWaiting.await(); - Thread.sleep(100); - - Thread thirdThread = - new Thread( - () -> { - ApiProxy.setEnvironmentForCurrentThread(environment); - writer.addLogRecordAndMaybeFlush(new LogRecord(LogRecord.Level.info, 0, "concurrent")); - thirdThreadAcquiredLock.countDown(); - }); - thirdThread.start(); - - assertThat(thirdThreadAcquiredLock.await(3, SECONDS)).isTrue(); - - slowFlush.set(new byte[0]); - firstFlushThread.join(3000); - secondFlushThread.join(3000); - thirdThread.join(3000); - } finally { - System.clearProperty("appengine.use.virtualthreads"); - } - } - - @Test - public void testFlushAndWait_holdsLockWhenVirtualThreadsDisabled() throws Exception { - System.clearProperty("appengine.use.virtualthreads"); + public void testSecondFlushAndWaitBlocksAddUntilPriorFlushCompletes() throws Exception { SettableFuture slowFlush = SettableFuture.create(); - when(delegate.makeAsyncCall(eq(environment), eq("logservice"), eq("Flush"), notNull(), notNull())) + when(delegate.makeAsyncCall( + eq(environment), eq("logservice"), eq("Flush"), notNull(), notNull())) .thenReturn(slowFlush); AppLogsWriter writer = new AppLogsWriter(response, SMALL_FLUSH, DEFAULT_MAX_LOG_LINE, 0); @@ -563,7 +504,8 @@ public void testFlushAndWait_holdsLockWhenVirtualThreadsDisabled() throws Except new Thread( () -> { ApiProxy.setEnvironmentForCurrentThread(environment); - writer.addLogRecordAndMaybeFlush(new LogRecord(LogRecord.Level.info, 0, "concurrent")); + writer.addLogRecordAndMaybeFlush( + new LogRecord(LogRecord.Level.info, 0, "concurrent")); thirdThreadAcquiredLock.countDown(); }); thirdThread.start(); @@ -578,72 +520,66 @@ public void testFlushAndWait_holdsLockWhenVirtualThreadsDisabled() throws Except } @Test - public void testVirtualThreadsFlush_noRecursionOrPrematureNulling() throws Exception { - System.setProperty("appengine.use.virtualthreads", "true"); - try { - SettableFuture slowFlush = SettableFuture.create(); - when(delegate.makeAsyncCall(eq(environment), eq("logservice"), eq("Flush"), notNull(), notNull())) - .thenReturn(slowFlush); + public void testAddOverThresholdWaitsForPendingFlushInsteadOfStartingAnother() throws Exception { + SettableFuture slowFlush = SettableFuture.create(); + when(delegate.makeAsyncCall( + eq(environment), eq("logservice"), eq("Flush"), notNull(), notNull())) + .thenReturn(slowFlush); - AppLogsWriter writer = new AppLogsWriter(response, SMALL_FLUSH, DEFAULT_MAX_LOG_LINE, 0); - writer.addLogRecordAndMaybeFlush(new LogRecord(LogRecord.Level.info, 0, "log message")); + AppLogsWriter writer = new AppLogsWriter(response, SMALL_FLUSH, DEFAULT_MAX_LOG_LINE, 0); + writer.addLogRecordAndMaybeFlush(new LogRecord(LogRecord.Level.info, 0, "log message")); - CountDownLatch firstFlushEntered = new CountDownLatch(1); - CountDownLatch secondThreadCompleted = new CountDownLatch(1); + CountDownLatch firstFlushEntered = new CountDownLatch(1); + CountDownLatch secondThreadCompleted = new CountDownLatch(1); - Thread firstThread = - new Thread( - () -> { - ApiProxy.setEnvironmentForCurrentThread(environment); - firstFlushEntered.countDown(); - writer.flushAndWait(); - }); - firstThread.start(); + Thread firstThread = + new Thread( + () -> { + ApiProxy.setEnvironmentForCurrentThread(environment); + firstFlushEntered.countDown(); + writer.flushAndWait(); + }); + firstThread.start(); - firstFlushEntered.await(); - Thread.sleep(100); + firstFlushEntered.await(); + Thread.sleep(100); - // Now second thread adds a log record exceeding SMALL_FLUSH while firstThread is waiting on slowFlush. - // Because pendingFlush is preserved and not prematurely set to null, secondThread waits on slowFlush - // instead of starting a new doFlush() immediately or looping recursively. - String largeMessage = new String(new char[(int) SMALL_FLUSH + 10]).replace('\0', 'a'); - Thread secondThread = - new Thread( - () -> { - ApiProxy.setEnvironmentForCurrentThread(environment); - writer.addLogRecordAndMaybeFlush(new LogRecord(LogRecord.Level.info, 0, largeMessage)); - secondThreadCompleted.countDown(); - }); - secondThread.start(); + // The second thread adds a log record exceeding SMALL_FLUSH while firstThread's flush is still + // pending. It must wait for that flush rather than start a second one. + String largeMessage = new String(new char[(int) SMALL_FLUSH + 10]).replace('\0', 'a'); + Thread secondThread = + new Thread( + () -> { + ApiProxy.setEnvironmentForCurrentThread(environment); + writer.addLogRecordAndMaybeFlush( + new LogRecord(LogRecord.Level.info, 0, largeMessage)); + secondThreadCompleted.countDown(); + }); + secondThread.start(); - // verify secondThread is waiting on slowFlush rather than completing immediately or looping - assertThat(secondThreadCompleted.await(500, MILLISECONDS)).isFalse(); - verify(delegate) - .makeAsyncCall(eq(environment), eq("logservice"), eq("Flush"), notNull(), notNull()); + // secondThread is waiting on slowFlush; only one flush has been issued. + assertThat(secondThreadCompleted.await(500, MILLISECONDS)).isFalse(); + verify(delegate) + .makeAsyncCall(eq(environment), eq("logservice"), eq("Flush"), notNull(), notNull()); - slowFlush.set(new byte[0]); - firstThread.join(3000); - secondThread.join(3000); - assertThat(secondThreadCompleted.await(3, SECONDS)).isTrue(); - } finally { - System.clearProperty("appengine.use.virtualthreads"); - } + slowFlush.set(new byte[0]); + firstThread.join(3000); + secondThread.join(3000); + assertThat(secondThreadCompleted.await(3, SECONDS)).isTrue(); } /** - * Simulates the original customer issue (b/514813839) on low-core / F1 instances where the virtual - * thread carrier parallelism is 1 (`GAE_MEMORY_MB <= 512`). Under legacy synchronized locking - * (`appengine.use.virtualthreads: false`), a request calling `flushAndWait()` holds the monitor - * lock while blocking on `slowFlush.get()`. This pins the carrier worker and starves concurrent - * request threads trying to access the logging pipeline, causing container deadlocks and timeouts. + * A thread that calls {@code flushAndWait()} while a flush is still pending waits for that flush + * holding the lock, so another thread logging to the same writer blocks until the pending flush + * completes, and then proceeds. */ @Test - public void testCustomerIssue_carrierPoolStarvationUnderLegacyLocking() throws Exception { - System.clearProperty("appengine.use.virtualthreads"); - ExecutorService carrierPool = Executors.newFixedThreadPool(1); + public void testFlushAndWaitBlocksConcurrentAddWhilePriorFlushPending() throws Exception { + ExecutorService executor = Executors.newFixedThreadPool(1); try { SettableFuture slowFlush = SettableFuture.create(); - when(delegate.makeAsyncCall(eq(environment), eq("logservice"), eq("Flush"), notNull(), notNull())) + when(delegate.makeAsyncCall( + eq(environment), eq("logservice"), eq("Flush"), notNull(), notNull())) .thenReturn(slowFlush); AppLogsWriter writer = new AppLogsWriter(response, SMALL_FLUSH, DEFAULT_MAX_LOG_LINE, 0); @@ -654,10 +590,9 @@ public void testCustomerIssue_carrierPoolStarvationUnderLegacyLocking() throws E CountDownLatch flushStarted = new CountDownLatch(1); CountDownLatch concurrentLogCompleted = new CountDownLatch(1); - // Task 1 simulates Request A calling flushAndWait() while a flush is in-flight. - // Under legacy mode, Request A blocks on slowFlush.get() INSIDE the synchronized(lock) monitor. + // Task 1 calls flushAndWait() while a flush is in flight and waits for it holding the lock. Future task1 = - carrierPool.submit( + executor.submit( () -> { ApiProxy.setEnvironmentForCurrentThread(environment); flushStarted.countDown(); @@ -667,19 +602,17 @@ public void testCustomerIssue_carrierPoolStarvationUnderLegacyLocking() throws E flushStarted.await(); Thread.sleep(100); - // Task 2 simulates Request B arriving concurrently from another thread trying to log. - // Because Task 1 holds the monitor lock while blocking inside slowFlush.get(), Task 2 cannot - // acquire the lock and is completely locked out / starved. + // Task 2 logs concurrently and blocks on the lock until the pending flush completes. Thread concurrentRequestThread = new Thread( () -> { ApiProxy.setEnvironmentForCurrentThread(environment); - writer.addLogRecordAndMaybeFlush(new LogRecord(LogRecord.Level.info, 0, "request B log")); + writer.addLogRecordAndMaybeFlush( + new LogRecord(LogRecord.Level.info, 0, "request B log")); concurrentLogCompleted.countDown(); }); concurrentRequestThread.start(); - // Verify that Request B is starved and cannot complete while slowFlush is pending. assertThat(concurrentLogCompleted.await(500, MILLISECONDS)).isFalse(); slowFlush.set(new byte[0]); @@ -687,68 +620,7 @@ public void testCustomerIssue_carrierPoolStarvationUnderLegacyLocking() throws E task1.get(3, SECONDS); assertThat(concurrentLogCompleted.await(3, SECONDS)).isTrue(); } finally { - carrierPool.shutdownNow(); - } - } - - /** - * Proves that under Virtual Threads mode (`appengine.use.virtualthreads: true`) on the same - * resource-constrained 1-carrier pool, our monitor lock decoupling allows Request A to wait on - * `slowFlush.get()` strictly outside the synchronized block. This enables Request B to immediately - * acquire the monitor lock, buffer its log record, and proceed without carrier starvation. - */ - @Test - public void testCustomerIssue_carrierPoolDecoupledProceedsWithoutStarvation() throws Exception { - System.setProperty("appengine.use.virtualthreads", "true"); - ExecutorService carrierPool = Executors.newFixedThreadPool(1); - try { - SettableFuture slowFlush = SettableFuture.create(); - when(delegate.makeAsyncCall(eq(environment), eq("logservice"), eq("Flush"), notNull(), notNull())) - .thenReturn(slowFlush); - - AppLogsWriter writer = new AppLogsWriter(response, SMALL_FLUSH, DEFAULT_MAX_LOG_LINE, 0); - String largeMessage = new String(new char[(int) SMALL_FLUSH + 10]).replace('\0', 'a'); - // Initiate the first flush on the main thread so slowFlush is pending as currentFlush. - writer.addLogRecordAndMaybeFlush(new LogRecord(LogRecord.Level.info, 0, largeMessage)); - - CountDownLatch flushStarted = new CountDownLatch(1); - CountDownLatch concurrentLogCompleted = new CountDownLatch(1); - - // Task 1 simulates Request A calling flushAndWait() while a flush is in-flight. - // Under decoupled virtual threads mode, Request A retrieves slowFlush via getPendingFlushLocked() - // and waits on slowFlush.get() OUTSIDE the synchronized(lock) monitor. - Future task1 = - carrierPool.submit( - () -> { - ApiProxy.setEnvironmentForCurrentThread(environment); - flushStarted.countDown(); - writer.flushAndWait(); - }); - - flushStarted.await(); - Thread.sleep(100); - - // Task 2 simulates Request B arriving concurrently trying to log. - // Because Request A released the monitor lock before blocking on slowFlush.get(), Request B can - // immediately acquire the monitor lock and complete! - Thread concurrentRequestThread = - new Thread( - () -> { - ApiProxy.setEnvironmentForCurrentThread(environment); - writer.addLogRecordAndMaybeFlush(new LogRecord(LogRecord.Level.info, 0, "request B log")); - concurrentLogCompleted.countDown(); - }); - concurrentRequestThread.start(); - - // Verify that Request B completes immediately WITHOUT carrier starvation or deadlock! - assertThat(concurrentLogCompleted.await(3, SECONDS)).isTrue(); - - slowFlush.set(new byte[0]); - concurrentRequestThread.join(3000); - task1.get(3, SECONDS); - } finally { - carrierPool.shutdownNow(); - System.clearProperty("appengine.use.virtualthreads"); + executor.shutdownNow(); } } diff --git a/runtime/local_jetty121/src/main/java/com/google/appengine/tools/development/jetty/JettyContainerService.java b/runtime/local_jetty121/src/main/java/com/google/appengine/tools/development/jetty/JettyContainerService.java index 05b82db1b..d9807e26c 100644 --- a/runtime/local_jetty121/src/main/java/com/google/appengine/tools/development/jetty/JettyContainerService.java +++ b/runtime/local_jetty121/src/main/java/com/google/appengine/tools/development/jetty/JettyContainerService.java @@ -183,9 +183,9 @@ public void exitScope(ContextHandler.APIContext context, Request request) { // location (WEB-INF/web.xml). // Only set the descriptor if the web.xml file actually exists. // Jetty 12 throws an IllegalArgumentException if the descriptor path is invalid. - if (webXmlLocation != null && webXmlLocation.exists()) { - context.setDescriptor(webXmlLocation.getAbsolutePath()); - } + if (webXmlLocation != null && webXmlLocation.exists()) { + context.setDescriptor(webXmlLocation.getAbsolutePath()); + } // Override the web.xml that Jetty automatically prepends to other // web.xml files. This is where the DefaultServlet is registered, @@ -339,8 +339,8 @@ protected void connectContainer() throws Exception { configuration.setSendDateHeader(false); configuration.setSendServerVersion(false); configuration.setSendXPoweredBy(false); - // Try to enable virtual threads if requested on java21: - if (Boolean.getBoolean("appengine.use.virtualthreads")) { + // Try to enable virtual threads if requested on Java 21+: + if (Boolean.getBoolean("appengine.use.virtualthreads") && Runtime.version().feature() >= 21) { QueuedThreadPool threadPool = new QueuedThreadPool(); threadPool.setVirtualThreadsExecutor(VirtualThreads.getDefaultVirtualThreadsExecutor()); server = new Server(threadPool); diff --git a/runtime/local_jetty121_ee11/src/main/java/com/google/appengine/tools/development/jetty/ee11/JettyContainerService.java b/runtime/local_jetty121_ee11/src/main/java/com/google/appengine/tools/development/jetty/ee11/JettyContainerService.java index 5f040c6e9..d02b47375 100644 --- a/runtime/local_jetty121_ee11/src/main/java/com/google/appengine/tools/development/jetty/ee11/JettyContainerService.java +++ b/runtime/local_jetty121_ee11/src/main/java/com/google/appengine/tools/development/jetty/ee11/JettyContainerService.java @@ -187,10 +187,10 @@ public void exitScope( // which is fine, it just means Jetty will look for it in the default // location (WEB-INF/web.xml). // Jetty 12 throws an IllegalArgumentException if the descriptor path is invalid. - if (webXmlLocation != null && webXmlLocation.exists()) { - context.setDescriptor(webXmlLocation.getAbsolutePath()); - } - + if (webXmlLocation != null && webXmlLocation.exists()) { + context.setDescriptor(webXmlLocation.getAbsolutePath()); + } + // Override the web.xml that Jetty automatically prepends to other // web.xml files. This is where the DefaultServlet is registered, // which serves static files. We override it to disable some @@ -268,42 +268,42 @@ public void exitScope( return appRoot; } - private void enterScope(ServletContextRequest request) { + private void enterScope(ServletContextRequest request) { - // We should have a request that use its associated environment, if there is no request - // we cannot select a local environment as picking the wrong one could result in - // waiting on the LocalEnvironment API call semaphore forever. - if (request == null) { - return; - } + // We should have a request that use its associated environment, if there is no request + // we cannot select a local environment as picking the wrong one could result in + // waiting on the LocalEnvironment API call semaphore forever. + if (request == null) { + return; + } - LocalEnvironment env = - (LocalEnvironment) request.getAttribute(LocalEnvironment.class.getName()); - if (env == null) { - env = - new LocalHttpRequestEnvironment( - appEngineWebXml.getAppId(), - WebModule.getModuleName(appEngineWebXml), - appEngineWebXml.getMajorVersionId(), - instance, - getPort(), - request.getServletApiRequest(), - SOFT_DEADLINE_DELAY_MS, - modulesFilterHelper); - env.getAttributes() - .put(LocalEnvironment.API_CALL_SEMAPHORE, new Semaphore(MAX_SIMULTANEOUS_API_CALLS)); + LocalEnvironment env = + (LocalEnvironment) request.getAttribute(LocalEnvironment.class.getName()); + if (env == null) { + env = + new LocalHttpRequestEnvironment( + appEngineWebXml.getAppId(), + WebModule.getModuleName(appEngineWebXml), + appEngineWebXml.getMajorVersionId(), + instance, + getPort(), + request.getServletApiRequest(), + SOFT_DEADLINE_DELAY_MS, + modulesFilterHelper); + env.getAttributes() + .put(LocalEnvironment.API_CALL_SEMAPHORE, new Semaphore(MAX_SIMULTANEOUS_API_CALLS)); env.getAttributes().put(DEFAULT_VERSION_HOSTNAME, "localhost:" + devAppServer.getPort()); - request.setAttribute(LocalEnvironment.class.getName(), env); - environments.add(env); - addCompletionListener(request); - } - - ApiProxy.setEnvironmentForCurrentThread(env); - DevAppServerModulesFilter.injectBackendServiceCurrentApiInfo( - backendName, backendInstance, portMappingProvider.getPortMapping()); + request.setAttribute(LocalEnvironment.class.getName(), env); + environments.add(env); + addCompletionListener(request); } + ApiProxy.setEnvironmentForCurrentThread(env); + DevAppServerModulesFilter.injectBackendServiceCurrentApiInfo( + backendName, backendInstance, portMappingProvider.getPortMapping()); + } + /** Check if the application contains a JSP file. */ private static boolean applicationContainsJSP(File dir, Pattern jspPattern) { for (File file : @@ -344,8 +344,8 @@ protected void connectContainer() throws Exception { configuration.setSendDateHeader(false); configuration.setSendServerVersion(false); configuration.setSendXPoweredBy(false); - // Try to enable virtual threads if requested on java21: - if (Boolean.getBoolean("appengine.use.virtualthreads")) { + // Try to enable virtual threads if requested on Java 21+: + if (Boolean.getBoolean("appengine.use.virtualthreads") && Runtime.version().feature() >= 21) { QueuedThreadPool threadPool = new QueuedThreadPool(); threadPool.setVirtualThreadsExecutor(VirtualThreads.getDefaultVirtualThreadsExecutor()); server = new Server(threadPool); diff --git a/runtime/main/src/main/java/com/google/apphosting/runtime/JavaRuntimeMain.java b/runtime/main/src/main/java/com/google/apphosting/runtime/JavaRuntimeMain.java index 938eb3b93..ab0b1dd66 100644 --- a/runtime/main/src/main/java/com/google/apphosting/runtime/JavaRuntimeMain.java +++ b/runtime/main/src/main/java/com/google/apphosting/runtime/JavaRuntimeMain.java @@ -23,7 +23,6 @@ import java.io.InputStream; import java.lang.reflect.Method; import java.util.Properties; -import java.util.function.UnaryOperator; import java.util.logging.Level; import java.util.logging.Logger; @@ -63,9 +62,6 @@ public class JavaRuntimeMain { private static final String ALLOW_NON_RESIDENT_SESSION_ACCESS = "gae.allow_non_resident_session_access"; - /* @VisibleForTesting */ - UnaryOperator envProvider = System::getenv; - public static void main(String[] args) { new JavaRuntimeMain().load(args); } @@ -80,8 +76,6 @@ public void load(String[] args) { // Process user defined properties as soon as possible, in the simple main Classpath. processOptionalProperties(args); - configureVirtualThreadParallelism(); - String appsRoot = getApplicationRoot(args); NullSandboxPlugin plugin = new NullSandboxPlugin(); ClassPathUtils classPathUtils = new ClassPathUtils(); @@ -100,50 +94,6 @@ public void load(String[] args) { } } - /** - * Configures the global default virtual thread scheduler parallelism (carrier pool size) based on - * the GAE Standard sandbox resource limits. - * - *

This configuration is critical under sandboxed, resource-constrained container environments - * (exposing fractional or single core quotas such as 0.5 CPU or 1.0 CPU). In these environments, - * the JVM default scheduler parallelism (which defaults to the underlying physical host core - * count, often 64+) triggers heavy thread context thrashing and CPU starvation. - * - *

This method maps the memory limit (GAE_MEMORY_MB) to a safe maximum carrier thread cap: - * - *

- * - *

This method is run early during the primary JVM bootstrap (JavaRuntimeMain.main) to - * guarantee the parallelism system property is set before the virtual thread scheduler is - * initialized, as late-bound properties are ignored. Explicit user-defined overrides (e.g., via - * JAVA_OPTS) are preserved. It only takes effect when {@code appengine.use.virtualthreads} is - * enabled. - */ - /* @VisibleForTesting */ - void configureVirtualThreadParallelism() { - String memoryMbStr = envProvider.apply("GAE_MEMORY_MB"); - if (Boolean.getBoolean("appengine.use.virtualthreads") - && memoryMbStr != null - && System.getProperty("jdk.virtualThreadScheduler.parallelism") == null) { - try { - int memoryMb = Integer.parseInt(memoryMbStr); - int parallelism = memoryMb <= 512 ? 1 : memoryMb <= 1024 ? 2 : 4; - System.setProperty("jdk.virtualThreadScheduler.parallelism", String.valueOf(parallelism)); - logger.info( - "Configured virtual thread parallelism to " - + parallelism - + " based on GAE_MEMORY_MB=" - + memoryMb); - } catch (NumberFormatException e) { - logger.log(Level.WARNING, "Failed to parse GAE_MEMORY_MB: " + memoryMbStr, e); - } - } - } - /** Parse the value of the --application_root flag. */ private String getApplicationRoot(String[] args) { return getFlag(args, "application_root", null); diff --git a/runtime/main/src/test/java/com/google/apphosting/runtime/JavaRuntimeMainTest.java b/runtime/main/src/test/java/com/google/apphosting/runtime/JavaRuntimeMainTest.java index f238be089..0019f5d46 100644 --- a/runtime/main/src/test/java/com/google/apphosting/runtime/JavaRuntimeMainTest.java +++ b/runtime/main/src/test/java/com/google/apphosting/runtime/JavaRuntimeMainTest.java @@ -22,9 +22,6 @@ import java.io.File; import java.io.IOException; import java.io.PrintWriter; -import java.util.HashMap; -import java.util.Map; -import org.junit.After; import org.junit.Before; import org.junit.Rule; import org.junit.Test; @@ -40,23 +37,11 @@ public class JavaRuntimeMainTest { @Rule public TemporaryFolder temporaryFolder = new TemporaryFolder(); private JavaRuntimeMain main; - private Map mockEnv; @Before public void setUp() { System.clearProperty("disable_api_call_logging_in_apiproxy"); - System.clearProperty("jdk.virtualThreadScheduler.parallelism"); - System.setProperty("appengine.use.virtualthreads", "true"); main = new JavaRuntimeMain(); - mockEnv = new HashMap<>(); - main.envProvider = (name) -> mockEnv.get(name); - } - - @After - public void tearDown() { - System.clearProperty("disable_api_call_logging_in_apiproxy"); - System.clearProperty("jdk.virtualThreadScheduler.parallelism"); - System.clearProperty("appengine.use.virtualthreads"); } @Test @@ -110,54 +95,4 @@ public void testWithOptionalProperties() throws IOException { assertThat(main.getApplicationPath(optionalProperties)).isEqualTo(appRoot); assertThat(System.getProperty("disable_api_call_logging_in_apiproxy")).isEqualTo("true"); } - - @Test - public void testConfigureVirtualThreadParallelism_noEnv() { - main.configureVirtualThreadParallelism(); - assertThat(System.getProperty("jdk.virtualThreadScheduler.parallelism")).isNull(); - } - - @Test - public void testConfigureVirtualThreadParallelism_f1() { - mockEnv.put("GAE_MEMORY_MB", "512"); - main.configureVirtualThreadParallelism(); - assertThat(System.getProperty("jdk.virtualThreadScheduler.parallelism")).isEqualTo("1"); - } - - @Test - public void testConfigureVirtualThreadParallelism_f2() { - mockEnv.put("GAE_MEMORY_MB", "1024"); - main.configureVirtualThreadParallelism(); - assertThat(System.getProperty("jdk.virtualThreadScheduler.parallelism")).isEqualTo("2"); - } - - @Test - public void testConfigureVirtualThreadParallelism_f4() { - mockEnv.put("GAE_MEMORY_MB", "2048"); - main.configureVirtualThreadParallelism(); - assertThat(System.getProperty("jdk.virtualThreadScheduler.parallelism")).isEqualTo("4"); - } - - @Test - public void testConfigureVirtualThreadParallelism_alreadySet() { - mockEnv.put("GAE_MEMORY_MB", "512"); - System.setProperty("jdk.virtualThreadScheduler.parallelism", "10"); - main.configureVirtualThreadParallelism(); - assertThat(System.getProperty("jdk.virtualThreadScheduler.parallelism")).isEqualTo("10"); - } - - @Test - public void testConfigureVirtualThreadParallelism_flagDisabled() { - System.clearProperty("appengine.use.virtualthreads"); - mockEnv.put("GAE_MEMORY_MB", "512"); - main.configureVirtualThreadParallelism(); - assertThat(System.getProperty("jdk.virtualThreadScheduler.parallelism")).isNull(); - } - - @Test - public void testConfigureVirtualThreadParallelism_invalidEnv() { - mockEnv.put("GAE_MEMORY_MB", "invalid"); - main.configureVirtualThreadParallelism(); - assertThat(System.getProperty("jdk.virtualThreadScheduler.parallelism")).isNull(); - } } diff --git a/runtime/runtime_impl_jetty12/src/main/java/com/google/apphosting/runtime/jetty/JettyServletEngineAdapter.java b/runtime/runtime_impl_jetty12/src/main/java/com/google/apphosting/runtime/jetty/JettyServletEngineAdapter.java index 8ff1ac3e7..c4bfd803a 100644 --- a/runtime/runtime_impl_jetty12/src/main/java/com/google/apphosting/runtime/jetty/JettyServletEngineAdapter.java +++ b/runtime/runtime_impl_jetty12/src/main/java/com/google/apphosting/runtime/jetty/JettyServletEngineAdapter.java @@ -15,7 +15,6 @@ */ package com.google.apphosting.runtime.jetty; -import static com.google.apphosting.runtime.AppEngineConstants.GAE_RUNTIME; import static com.google.apphosting.runtime.AppEngineConstants.IGNORE_RESPONSE_SIZE_LIMIT; import com.google.apphosting.base.AppVersionKey; @@ -32,9 +31,8 @@ import com.google.apphosting.runtime.jetty.proxy.JettyHttpProxy; import com.google.common.flogger.GoogleLogger; import java.util.Objects; -import java.util.concurrent.Executor; -import java.util.concurrent.ForkJoinPool; import org.eclipse.jetty.server.Server; +import org.eclipse.jetty.util.VirtualThreads; import org.eclipse.jetty.util.thread.QueuedThreadPool; /** @@ -61,17 +59,10 @@ public JettyServletEngineAdapter() {} public void start(String serverInfo, ServletEngineAdapter.Config runtimeOptions) { QueuedThreadPool threadPool = new QueuedThreadPool(MAX_THREAD_POOL_THREADS, MIN_THREAD_POOL_THREADS); - // Try to enable virtual threads if requested and on java21: - if (Boolean.getBoolean("appengine.use.virtualthreads") - && ("java21".equals(GAE_RUNTIME) || "java25".equals(GAE_RUNTIME))) { - int maxParallelism = getMaxSafeCarrierParallelism(); - Executor virtualThreadsExecutor = - new ForkJoinPool( - maxParallelism, ForkJoinPool.defaultForkJoinWorkerThreadFactory, null, true); - threadPool.setVirtualThreadsExecutor(virtualThreadsExecutor); - logger.atInfo().log( - "Configuring Appengine web server virtual threads with capped carrier parallelism: %d", - maxParallelism); + // Try to enable virtual threads if requested and on Java 21+: + if (Boolean.getBoolean("appengine.use.virtualthreads") && Runtime.version().feature() >= 21) { + threadPool.setVirtualThreadsExecutor(VirtualThreads.getDefaultVirtualThreadsExecutor()); + logger.atInfo().log("Configuring Appengine web server virtual threads."); } server = @@ -147,24 +138,4 @@ public void serviceRequest(UPRequest upRequest, MutableUpResponse upResponse) th throw new UnsupportedOperationException( "serviceRequest is not supported in HTTP connector mode"); } - - /** - * Calculates a safe maximum carrier thread count based on GAE sandbox memory boundaries to - * prevent OS scheduling thrashing on fractional/low-core instances. - */ - static int getMaxSafeCarrierParallelism() { - return getMaxSafeCarrierParallelism(System.getenv("GAE_MEMORY_MB")); - } - - static int getMaxSafeCarrierParallelism(String memoryMbStr) { - if (memoryMbStr == null || memoryMbStr.isEmpty()) { - return 4; // Conservative default cap for standard runtimes - } - try { - int memoryMb = Integer.parseInt(memoryMbStr); - return memoryMb <= 512 ? 1 : memoryMb <= 1024 ? 2 : 4; - } catch (NumberFormatException e) { - return 4; // Safety Fallback - } - } } diff --git a/runtime/runtime_impl_jetty121/API_CLIENTS.md b/runtime/runtime_impl_jetty121/API_CLIENTS.md new file mode 100644 index 000000000..30af8551d --- /dev/null +++ b/runtime/runtime_impl_jetty121/API_CLIENTS.md @@ -0,0 +1,161 @@ + + +# App Engine API Client Configuration + +The App Engine Java runtime communicates with Google Cloud APIs (such as +Datastore, Task Queue, and Memcache) using an HTTP-based RPC mechanism. The +runtime includes two HTTP client implementations for this purpose: a default +client based on Jetty and an alternative client using the JDK's built-in HTTP +facilities. + +This document describes both clients and how to configure them using environment +variables and Java system properties. + +## Jetty HTTP Client (Default) + +The Jetty HTTP client is the default client used by the runtime. It is based on +the [Eclipse Jetty](https://eclipse.dev/jetty/) HTTP client and is optimized for +high performance and efficient connection management. + +By default, the client is configured to allow a maximum of 100 concurrent +threads and 100 concurrent connections for API calls. These limits help prevent +memory exhaustion on smaller App Engine instance types (like F1 or F2) and avoid +overwhelming backend services during sudden traffic spikes or high rates of +failing requests with aggressive retry logic. If the thread limit is reached, +subsequent requests are queued until a thread becomes available. + +### Configuration + +You can configure the Jetty client using the following environment variables and +system properties: + +* **`APPENGINE_API_MAX_CONNECTIONS`** (Environment Variable): Sets the maximum + number of concurrent connections in the HTTP client pool. + * Default: `100` +* **`APPENGINE_API_MAX_THREADS`** (Environment Variable): Sets the maximum + number of concurrent threads for executing API calls. If unset, this also + defaults to 100. This is the most direct way to control API call throughput + and prevent backend overload. If set to lower values (e.g. 5 or 8), + `minThreads` automatically scales down (`min(10, maxThreads)`) to prevent + startup failures. + * Default: `100` +* **`APPENGINE_API_CALLS_IDLE_TIMEOUT_MS`** (Environment Variable): Sets the + idle timeout in milliseconds for connections in the connection pool. + Connections that are idle for longer than this duration may be closed. + * Default: `58000` (58 seconds) +* **`appengine.api.use.virtualthreads`** (Java System Property): If set to + `true` on Java 21+, the client will use Java Virtual Threads to execute API + requests, avoiding platform thread stack memory overhead while honoring pool + and connection limits. + * Default: `false` + +## JDK HTTP Client + +The JDK HTTP client uses Java's built-in `HttpURLConnection` for API calls. It +is provided as an alternative to the Jetty client and can be useful for +troubleshooting network- or connection-related issues that might be specific to +one client implementation. In general, it may be less performant than the +default Jetty client. + +To use the JDK client instead of the Jetty client, set the +`APPENGINE_API_CALLS_USING_JDK_CLIENT` environment variable to any non-null +value (e.g., `true`). + +### Configuration + +When using the JDK client, you can configure its behavior with the following +settings: + +* **`appengine.api.use.virtualthreads`** (Java System Property): If set to + `true`, the JDK client will use Java Virtual Threads (when available on the + JVM) to handle API requests. This avoids platform thread stack size overhead + while keeping concurrent in-flight calls safely throttled via a semaphore to + prevent backend overload and retry storms. Note that this is different than + `appengine.use.virtualthreads`. + * Default: `false` + +The JDK client controls concurrency using the following environment variables +(both when using virtual threads and when using the platform thread pool): + +* **`APPENGINE_API_MAX_CONNECTIONS`** (Environment Variable): Sets the maximum + number of concurrent connections in the HTTP client pool. + * Default: `100` +* **`APPENGINE_API_MAX_THREADS`** (Environment Variable): Sets the maximum + number of concurrent API calls (via thread pool size or virtual thread + concurrency semaphore). This limits the number of concurrent in-flight API + calls to prevent overwhelming the backend. + * Default: `100` + +## Datastore-Specific Configuration + +### `beginTransaction` Retries + +In addition to the HTTP client settings, the Datastore client library includes +specific retry logic for `beginTransaction()` calls. These calls can fail with +transient errors such as `DatastoreFailureException`, +`DatastoreTimeoutException`, or `ApiProxy.RPCFailedException`, especially under +high contention when multiple transactions attempt to access the same entity +group simultaneously. + +To handle this, `DatastoreService.beginTransaction()` automatically retries +failed attempts with exponential backoff, starting at 100ms. You can configure +the number of retry attempts using a system property: + +* **`appengine.datastore.retries`** (Java System Property): The maximum number + of times to retry a `beginTransaction` call if it fails with + `DatastoreFailureException`, `DatastoreTimeoutException`, or + `ApiProxy.RPCFailedException`. This retry logic applies only to + `beginTransaction` calls; other Datastore operations are not retried by this + mechanism. + * Default: `1` + +```xml + + + +``` + +## Recommended value for Jetty 12.1 / Java 25 + +The Java 25 runtime environment is highly performant, and in rare cases of very +high throughput, this can lead to backend services (like Datastore) being +temporarily overloaded, which may result in exceptions like +`DatastoreFailureException: Internal Datastore Error`. If you encounter such +issues, you can throttle the rate of API calls by reducing the maximum number of +concurrent threads used by the API client. We recommend starting with a lower +value and adjusting as needed: `APPENGINE_API_MAX_THREADS=50` + +## Configuring via `appengine-web.xml` + +You can set these options by adding `` and `` +sections to your `appengine-web.xml` file. + +For example, to switch to the JDK client, enable virtual threads, increase the +thread limit to 200, and set Datastore `beginTransaction` retries to 3, you +would add: + +```xml + + + + + + + + + +``` diff --git a/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/HttpApiHostClient.java b/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/HttpApiHostClient.java index 8d44a665d..8b2e1c5cf 100644 --- a/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/HttpApiHostClient.java +++ b/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/HttpApiHostClient.java @@ -76,6 +76,12 @@ abstract class HttpApiHostClient implements APIHostClientInterface { abstract static class Config { abstract double extraTimeoutSeconds(); + /** + * The maximum number of concurrent connections to the API host. + * + *

This value is used to configure both the Jetty/JDK HTTP client's connection pool limit + * and, when applicable, the maximum number of threads in the client's thread pool. + */ abstract OptionalInt maxConnectionsPerDestination(); /** For testing that we handle missing Content-Length correctly. */ @@ -123,8 +129,36 @@ Config config() { return config; } + static int getMaxThreads(Config config) { + int maxThreads = config.maxConnectionsPerDestination().orElse(100); + if (maxThreads <= 0) { + maxThreads = 100; + } + String maxThreadsEnv = System.getenv("APPENGINE_API_MAX_THREADS"); + if (maxThreadsEnv != null) { + try { + int envMaxThreads = Integer.parseInt(maxThreadsEnv); + if (envMaxThreads > 0) { + logger.atInfo().log( + "Overriding API max threads to %d from environment variable.", envMaxThreads); + return envMaxThreads; + } else { + logger.atWarning().log( + "APPENGINE_API_MAX_THREADS must be positive: %d, using default %d", + envMaxThreads, maxThreads); + } + } catch (NumberFormatException e) { + logger.atWarning().withCause(e).log( + "Invalid value for APPENGINE_API_MAX_THREADS: %s, using default %d", + maxThreadsEnv, maxThreads); + } + } + return maxThreads; + } + static HttpApiHostClient create(String url, Config config) { - if (System.getenv("APPENGINE_API_CALLS_USING_JDK_CLIENT") != null) { + if (System.getenv("APPENGINE_API_CALLS_USING_JDK_CLIENT") != null + || Boolean.getBoolean("com.google.appengine.api.calls.using.jdk.client")) { logger.atInfo().log("Using JDK HTTP client for API calls"); return JdkHttpApiHostClient.create(url, config); } else { diff --git a/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/HttpApiHostClientFactory.java b/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/HttpApiHostClientFactory.java index 9bcab0d92..ecf13ab37 100644 --- a/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/HttpApiHostClientFactory.java +++ b/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/HttpApiHostClientFactory.java @@ -20,20 +20,44 @@ import com.google.apphosting.runtime.anyrpc.APIHostClientInterface; import com.google.apphosting.runtime.http.HttpApiHostClient.Config; +import com.google.common.flogger.GoogleLogger; import com.google.common.net.HostAndPort; import java.util.OptionalInt; /** Makes instances of {@link HttpApiHostClient}. */ public class HttpApiHostClientFactory { + private static final GoogleLogger logger = GoogleLogger.forEnclosingClass(); + private HttpApiHostClientFactory() {} /** * Creates a new HttpApiHostClient instance to talk to the HTTP-based API server on the given host * and port. This method is called reflectively from ApiHostClientFactory. + * + *

The maximum number of concurrent connections can be configured by setting the {@code + * APPENGINE_API_MAX_CONNECTIONS} environment variable to a positive integer. If set, this value + * overrides the {@code maxConcurrentRpcs} parameter. + * + * @param hostAndPort The host and port of the API server. + * @param maxConcurrentRpcs The default maximum number of concurrent RPCs, used if the environment + * variable is not set. + * @return A new {@link APIHostClientInterface} instance. */ public static APIHostClientInterface create( HostAndPort hostAndPort, OptionalInt maxConcurrentRpcs) { String url = "http://" + hostAndPort + REQUEST_ENDPOINT; + String maxConnectionsEnv = System.getenv("APPENGINE_API_MAX_CONNECTIONS"); + if (maxConnectionsEnv != null) { + try { + int maxConnections = Integer.parseInt(maxConnectionsEnv); + if (maxConnections > 0) { + maxConcurrentRpcs = OptionalInt.of(maxConnections); + } + } catch (NumberFormatException e) { + logger.atWarning().withCause(e).log( + "Failed to parse APPENGINE_API_MAX_CONNECTIONS: %s", maxConnectionsEnv); + } + } Config config = Config.builder().setMaxConnectionsPerDestination(maxConcurrentRpcs).build(); return HttpApiHostClient.create(url, config); } diff --git a/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/JdkHttpApiHostClient.java b/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/JdkHttpApiHostClient.java index cb84007e5..cc256924b 100644 --- a/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/JdkHttpApiHostClient.java +++ b/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/JdkHttpApiHostClient.java @@ -17,29 +17,40 @@ package com.google.apphosting.runtime.http; import static java.lang.Math.max; +import static java.util.concurrent.TimeUnit.SECONDS; import com.google.apphosting.base.protos.RuntimePb.APIResponse; import com.google.apphosting.runtime.anyrpc.AnyRpcCallback; import com.google.common.flogger.GoogleLogger; import com.google.common.io.ByteStreams; import com.google.common.primitives.Ints; -import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; import java.io.UncheckedIOException; +import java.lang.reflect.Method; import java.net.HttpURLConnection; import java.net.MalformedURLException; import java.net.SocketTimeoutException; +import java.net.URI; import java.net.URL; import java.util.concurrent.Executor; import java.util.concurrent.Executors; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.Semaphore; import java.util.concurrent.ThreadFactory; +import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.atomic.AtomicInteger; /** * An alternative API client that uses the JDK's built-in HTTP client. This is likely to be much * less performant than {@link JettyHttpApiHostClient} but should allow us to determine whether * communications problems we are seeing are due to the Jetty client. + * + *

By default, this client uses a bounded thread pool to execute API calls, with the maximum + * number of threads determined by configuration. If the system property {@code + * appengine.api.use.virtualthreads} is set to {@code true}, it will instead use virtual threads via + * {@link Executors#newVirtualThreadPerTaskExecutor()}, with in-flight concurrency throttled by a + * semaphore to prevent backend overload. */ class JdkHttpApiHostClient extends HttpApiHostClient { private static final GoogleLogger logger = GoogleLogger.forEnclosingClass(); @@ -50,26 +61,81 @@ class JdkHttpApiHostClient extends HttpApiHostClient { private final URL url; private final Executor executor; + private final Semaphore concurrencySemaphore; - private JdkHttpApiHostClient(Config config, URL url, Executor executor) { + private JdkHttpApiHostClient( + Config config, URL url, Executor executor, Semaphore concurrencySemaphore) { super(config); this.url = url; this.executor = executor; + this.concurrencySemaphore = concurrencySemaphore; } + /** + * Creates a {@link JdkHttpApiHostClient}. + * + *

If the system property {@code appengine.api.use.virtualthreads} is set to {@code true}, a + * virtual thread executor is used to run requests with concurrency throttled to {@code + * maxThreads}. Otherwise, a bounded {@link ThreadPoolExecutor} is created, with {@code + * maxThreads} derived from {@code config.maxConnectionsPerDestination()}. + * + * @param url The URL of the API host. + * @param config Configuration for the client, including connection limits. + * @return A new {@link JdkHttpApiHostClient}. + */ + @SuppressWarnings("AllowVirtualThreads") static JdkHttpApiHostClient create(String url, Config config) { try { - ThreadFactory factory = - runnable -> { - Thread t = new Thread(rootThreadGroup(), runnable); - t.setName("JdkHttp-" + threadCount.incrementAndGet()); - t.setDaemon(true); - return t; - }; - Executor executor = Executors.newCachedThreadPool(factory); - return new JdkHttpApiHostClient(config, new URL(url), executor); - } catch (MalformedURLException e) { - throw new UncheckedIOException(e); + Executor executor = null; + Semaphore concurrencySemaphore = null; + int maxThreads = getMaxThreads(config); + if (Boolean.getBoolean("appengine.api.use.virtualthreads")) { + try { + Method newVirtualThreadPerTaskExecutor = + Executors.class.getMethod("newVirtualThreadPerTaskExecutor"); + executor = (Executor) newVirtualThreadPerTaskExecutor.invoke(null); + concurrencySemaphore = new Semaphore(maxThreads); + logger.atInfo().log( + "Using virtual threads for JdkHttpApiHostClient with concurrency capped at %d.", + maxThreads); + } catch (ReflectiveOperationException e) { + logger.atInfo().log( + "appengine.api.use.virtualthreads is true, but virtual threads are not available on" + + " this JDK. Falling back to thread pool for JdkHttpApiHostClient."); + } + } + if (executor == null) { + ThreadFactory factory = + runnable -> { + Thread t = new Thread(rootThreadGroup(), runnable); + t.setName("JdkHttp-" + threadCount.incrementAndGet()); + t.setDaemon(true); + return t; + }; + /* + * Thread Pool Configuration & Bug Analysis: + * + * Similar to the JettyHttpApiHostClient, we explicitly bound the thread pool. + * We cap the threads at `maxConnectionsPerDestination` (which defaults to 100) + * instead of a hardcoded 200 to prevent severe memory pressure (Thread Stack sizes) + * on smaller AppEngine instance classes like F1 (256MB) or F2 (512MB). + * An unbounded thread pool allows a failing RPC to rapidly spin up thousands + * of threads under retry, which overwhelms the JVM and the internal Datastore + * Appserver connection, forcing it to respond with masking INTERNAL_ERROR fallbacks. + */ + ThreadPoolExecutor tpe = + new ThreadPoolExecutor( + maxThreads, maxThreads, 60L, SECONDS, new LinkedBlockingQueue<>(), factory); + tpe.allowCoreThreadTimeOut(true); + executor = tpe; + } + return new JdkHttpApiHostClient( + config, URI.create(url).toURL(), executor, concurrencySemaphore); + } catch (MalformedURLException | IllegalArgumentException e) { + throw new UncheckedIOException( + e instanceof MalformedURLException malformedUrlException + ? malformedUrlException + : new MalformedURLException(e.getMessage())); } } @@ -82,6 +148,13 @@ private static ThreadGroup rootThreadGroup() { return group; } + /** + * Asynchronously sends an API request to the API host using a thread pool. + * + * @param requestBytes The serialized API request. + * @param context The context for the request, including deadline information. + * @param callback Callback to be invoked with the API response or failure. + */ @Override void send( byte[] requestBytes, @@ -94,6 +167,15 @@ private void doSend( byte[] requestBytes, HttpApiHostClient.Context context, AnyRpcCallback callback) { + if (concurrencySemaphore != null) { + try { + concurrencySemaphore.acquire(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + cancelled(callback); + return; + } + } try { HttpURLConnection connection = (HttpURLConnection) url.openConnection(); connection.setDoOutput(true); @@ -102,8 +184,9 @@ private void doSend( if (context.getDeadlineNanos().isPresent()) { double deadlineSeconds = context.getDeadlineNanos().get() / 1e9; connection.addRequestProperty(DEADLINE_HEADER, Double.toString(deadlineSeconds)); - int deadlineMillis = - Ints.saturatedCast(max(1, context.getDeadlineNanos().get() / 1_000_000)); + double fallbackDeadlineSeconds = deadlineSeconds + config().extraTimeoutSeconds(); + int deadlineMillis = Ints.saturatedCast(max(1, (long) (fallbackDeadlineSeconds * 1000))); + connection.setConnectTimeout(deadlineMillis); connection.setReadTimeout(deadlineMillis); } connection.setFixedLengthStreamingMode(requestBytes.length); @@ -111,33 +194,66 @@ private void doSend( try (OutputStream out = connection.getOutputStream()) { out.write(requestBytes); } - if (connection.getResponseCode() == HttpURLConnection.HTTP_OK) { - int length = connection.getContentLength(); + int responseCode = connection.getResponseCode(); + if (responseCode == HttpURLConnection.HTTP_OK) { + int length = config().ignoreContentLength() ? -1 : connection.getContentLength(); if (length > MAX_LENGTH) { connection.getInputStream().close(); + connection.disconnect(); responseTooBig(callback); - } else { + } else if (length >= 0) { byte[] buffer = new byte[length]; try (InputStream in = connection.getInputStream()) { ByteStreams.readFully(in, buffer); // EOFException (an IOException) if too few bytes receivedResponse(buffer, length, context, callback); } + } else { + // Chunked transfer encoding or unspecified content length + byte[] buffer; + try (InputStream in = connection.getInputStream()) { + buffer = ByteStreams.limit(in, MAX_LENGTH + 1).readAllBytes(); + } + if (buffer.length > MAX_LENGTH) { + connection.disconnect(); + responseTooBig(callback); + } else { + receivedResponse(buffer, buffer.length, context, callback); + } } + } else { + String httpError = responseCode + " " + connection.getResponseMessage(); + connection.disconnect(); + logger.atWarning().log("HTTP communication got error: %s", httpError); + communicationFailure(context, httpError, callback, null); } } catch (SocketTimeoutException e) { logger.atWarning().withCause(e).log("SocketTimeoutException"); timeout(callback); - } catch (IOException e) { - logger.atWarning().withCause(e).log("IOException"); - communicationFailure(context, e.toString(), callback, e); + } catch (Throwable t) { + logger.atWarning().withCause(t).log("HTTP communication failure"); + communicationFailure(context, t.toString(), callback, t); + } finally { + if (concurrencySemaphore != null) { + concurrencySemaphore.release(); + } } } + /** + * This operation is not supported by JdkHttpApiHostClient. + * + * @throws UnsupportedOperationException always. + */ @Override public void enable() { throw new UnsupportedOperationException(); } + /** + * This operation is not supported by JdkHttpApiHostClient. + * + * @throws UnsupportedOperationException always. + */ @Override public void disable() { throw new UnsupportedOperationException(); diff --git a/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/JettyHttpApiHostClient.java b/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/JettyHttpApiHostClient.java index f55286bc6..f36151ee9 100644 --- a/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/JettyHttpApiHostClient.java +++ b/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/http/JettyHttpApiHostClient.java @@ -32,11 +32,8 @@ import java.nio.channels.ClosedSelectorException; import java.util.Arrays; import java.util.Map; -import java.util.concurrent.Executors; import java.util.concurrent.RejectedExecutionException; -import java.util.concurrent.ThreadFactory; import java.util.concurrent.TimeoutException; -import java.util.concurrent.atomic.AtomicInteger; import org.eclipse.jetty.client.BytesRequestContent; import org.eclipse.jetty.client.HttpClient; import org.eclipse.jetty.client.HttpResponseException; @@ -48,6 +45,8 @@ import org.eclipse.jetty.http.HttpHeader; import org.eclipse.jetty.http.HttpMethod; import org.eclipse.jetty.io.EofException; +import org.eclipse.jetty.util.VirtualThreads; +import org.eclipse.jetty.util.thread.QueuedThreadPool; import org.eclipse.jetty.util.thread.ScheduledExecutorScheduler; import org.eclipse.jetty.util.thread.Scheduler; @@ -55,8 +54,6 @@ class JettyHttpApiHostClient extends HttpApiHostClient { private static final GoogleLogger logger = GoogleLogger.forEnclosingClass(); - private static final AtomicInteger threadCount = new AtomicInteger(); - private final String url; private final HttpClient httpClient; @@ -66,6 +63,17 @@ private JettyHttpApiHostClient(String url, HttpClient httpClient, Config config) this.httpClient = httpClient; } + /** + * Creates and starts a {@link JettyHttpApiHostClient}. + * + *

The {@code config.maxConnectionsPerDestination()} parameter is used to configure both {@link + * HttpClient#setMaxConnectionsPerDestination(int)} and the maximum number of threads in the + * {@link QueuedThreadPool}. + * + * @param url The URL of the API host. + * @param config Configuration for the client. + * @return A new, started {@link JettyHttpApiHostClient}. + */ static JettyHttpApiHostClient create(String url, Config config) { Preconditions.checkNotNull(url); HttpClient httpClient = new HttpClient(); @@ -86,20 +94,45 @@ static JettyHttpApiHostClient create(String url, Config config) { boolean daemon = false; Scheduler scheduler = new ScheduledExecutorScheduler(schedulerName, daemon, myLoader, myThreadGroup); - ThreadFactory factory = - runnable -> { - Thread t = new Thread(myThreadGroup, runnable); - t.setName("JettyHttpApiHostClient-" + threadCount.incrementAndGet()); - t.setDaemon(true); - return t; - }; - // By default HttpClient will use a QueuedThreadPool with minThreads=8 and maxThreads=200. - // 8 threads is probably too much for most apps, especially since asynchronous I/O means that - // 8 concurrent API requests probably don't need that many threads. It's also not clear - // what advantage we'd get from using a QueuedThreadPool with a smaller minThreads value, versus - // just one of the standard java.util.concurrent pools. Here we have minThreads=1, maxThreads=∞, - // and idleTime=60 seconds. maxThreads=200 and maxThreads=∞ are probably equivalent in practice. - httpClient.setExecutor(Executors.newCachedThreadPool(factory)); + /* + * Thread Pool Configuration & Bug Analysis: + * + * In previous versions of the runtime, an unbounded CachedThreadPool was used here: + * `httpClient.setExecutor(Executors.newCachedThreadPool(factory));` + * + * Under high load (e.g., when a customer's custom retry logic aggressively retries failing + * RPCs like `BeginTransaction`), an unbounded thread pool creates thousands of threads instantly. + * This leads to a system collapse: + * 1. JVM Overload: The Java container becomes severely memory and CPU constrained. + * 2. Appserver Flooded: The avalanche of concurrent requests from the Java container floods the + * C++ Appserver proxy. + * 3. Triggering the C++ Bug Mask: Under massive load, the C++ Appserver's gRPC calls to the + * Datastore fail with UNAVAILABLE or RESOURCE_EXHAUSTED errors. + * 4. The Response: Because these aren't standard application errors, the C++ code + * (DatastoreClientHelper::DoneImpl) masks them as `Error::INTERNAL_ERROR` and returns the + * message "Internal Datastore Error" to the Java client to prevent leaking internal + * infrastructure details. + * 5. The Java client throws DatastoreFailureException, triggering the customer's loop again. + * + * To prevent this "retry storm", we explicitly use a bounded QueuedThreadPool. + * We cap the threads at `maxConnectionsPerDestination` (which defaults to 100) + * instead of a hardcoded 200 to prevent severe memory pressure (Thread Stack sizes) + * on smaller AppEngine instance classes like F1 (256MB) or F2 (512MB). + * If the system experiences a spike, Jetty will safely queue the outgoing RPCs, preventing the + * JVM and the Appserver from being overwhelmed and eliminating the INTERNAL_ERROR fallback loop. + */ + int maxThreads = getMaxThreads(config); + int minThreads = Math.max(1, Math.min(10, maxThreads)); + QueuedThreadPool threadPool = + new QueuedThreadPool(maxThreads, minThreads, 60000, null, myThreadGroup); + threadPool.setName("JettyHttpApiHostClient"); + threadPool.setDaemon(true); + if (Boolean.getBoolean("appengine.api.use.virtualthreads") + && Runtime.version().feature() >= 21) { + threadPool.setVirtualThreadsExecutor(VirtualThreads.getDefaultVirtualThreadsExecutor()); + logger.atInfo().log("Using virtual threads for JettyHttpApiHostClient."); + } + httpClient.setExecutor(threadPool); httpClient.setScheduler(scheduler); config.maxConnectionsPerDestination().ifPresent(httpClient::setMaxConnectionsPerDestination); try { @@ -220,6 +253,13 @@ && config().treatClosedChannelAsCancellation()) { } } + /** + * Asynchronously sends an API request to the API host. + * + * @param requestBytes The serialized API request. + * @param context The context for the request, including deadline information. + * @param callback Callback to be invoked with the API response or failure. + */ @Override void send( byte[] requestBytes, @@ -260,6 +300,10 @@ void send( request.send(completeListener); } + /** + * Disables the client by stopping the underlying {@link HttpClient}. Subsequent calls to {@link + * #send} may fail until {@link #enable()} is called. + */ @Override public synchronized void disable() { try { @@ -271,6 +315,7 @@ public synchronized void disable() { } } + /** Enables the client by starting the underlying {@link HttpClient}. */ @Override public synchronized void enable() { try { diff --git a/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/jetty/JettyServletEngineAdapter.java b/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/jetty/JettyServletEngineAdapter.java index fbbcde3ce..1c5512e91 100644 --- a/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/jetty/JettyServletEngineAdapter.java +++ b/runtime/runtime_impl_jetty121/src/main/java/com/google/apphosting/runtime/jetty/JettyServletEngineAdapter.java @@ -15,9 +15,7 @@ */ package com.google.apphosting.runtime.jetty; -import static com.google.apphosting.runtime.AppEngineConstants.GAE_RUNTIME; import static com.google.apphosting.runtime.AppEngineConstants.IGNORE_RESPONSE_SIZE_LIMIT; -import static com.google.common.base.Strings.isNullOrEmpty; import com.google.apphosting.base.AppVersionKey; import com.google.apphosting.base.protos.AppinfoPb; @@ -61,9 +59,8 @@ public JettyServletEngineAdapter() {} public void start(String serverInfo, ServletEngineAdapter.Config runtimeOptions) { QueuedThreadPool threadPool = new QueuedThreadPool(MAX_THREAD_POOL_THREADS, MIN_THREAD_POOL_THREADS); - // Try to enable virtual threads if requested and on java21: - if (Boolean.getBoolean("appengine.use.virtualthreads") - && ("java21".equals(GAE_RUNTIME) || "java25".equals(GAE_RUNTIME))) { + // Try to enable virtual threads if requested and on Java 21+: + if (Boolean.getBoolean("appengine.use.virtualthreads") && Runtime.version().feature() >= 21) { threadPool.setVirtualThreadsExecutor(VirtualThreads.getDefaultVirtualThreadsExecutor()); logger.atInfo().log("Configuring Appengine web server virtual threads."); } @@ -144,24 +141,4 @@ public void serviceRequest(UPRequest upRequest, MutableUpResponse upResponse) th throw new UnsupportedOperationException( "serviceRequest is not supported in HTTP connector mode"); } - - /** - * Calculates a safe maximum carrier thread count based on GAE sandbox memory boundaries to - * prevent OS scheduling thrashing on fractional/low-core instances. - */ - static int getMaxSafeCarrierParallelism() { - return getMaxSafeCarrierParallelism(System.getenv("GAE_MEMORY_MB")); - } - - static int getMaxSafeCarrierParallelism(String memoryMbStr) { - if (isNullOrEmpty(memoryMbStr)) { - return 4; // Conservative default cap for standard runtimes - } - try { - int memoryMb = Integer.parseInt(memoryMbStr); - return memoryMb <= 512 ? 1 : memoryMb <= 1024 ? 2 : 4; - } catch (NumberFormatException e) { - return 4; // Safety Fallback - } - } } diff --git a/runtime/runtime_impl_jetty121/src/test/java/com/google/apphosting/runtime/jetty/JettyServletEngineAdapterTest.java b/runtime/runtime_impl_jetty121/src/test/java/com/google/apphosting/runtime/jetty/JettyServletEngineAdapterTest.java deleted file mode 100644 index 693dfff54..000000000 --- a/runtime/runtime_impl_jetty121/src/test/java/com/google/apphosting/runtime/jetty/JettyServletEngineAdapterTest.java +++ /dev/null @@ -1,39 +0,0 @@ -/* - * Copyright 2026 Google LLC - * - * 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 - * - * https://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.apphosting.runtime.jetty; - -import static com.google.common.truth.Truth.assertThat; - -import org.junit.Test; -import org.junit.runner.RunWith; -import org.junit.runners.JUnit4; - -@RunWith(JUnit4.class) -public class JettyServletEngineAdapterTest { - - @Test - public void testGetMaxSafeCarrierParallelism_boundaries() { - assertThat(JettyServletEngineAdapter.getMaxSafeCarrierParallelism(null)).isEqualTo(4); - assertThat(JettyServletEngineAdapter.getMaxSafeCarrierParallelism("")).isEqualTo(4); - assertThat(JettyServletEngineAdapter.getMaxSafeCarrierParallelism("invalid")).isEqualTo(4); - assertThat(JettyServletEngineAdapter.getMaxSafeCarrierParallelism("256")).isEqualTo(1); - assertThat(JettyServletEngineAdapter.getMaxSafeCarrierParallelism("512")).isEqualTo(1); - assertThat(JettyServletEngineAdapter.getMaxSafeCarrierParallelism("600")).isEqualTo(2); - assertThat(JettyServletEngineAdapter.getMaxSafeCarrierParallelism("1024")).isEqualTo(2); - assertThat(JettyServletEngineAdapter.getMaxSafeCarrierParallelism("2048")).isEqualTo(4); - } -}