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..77367a0b5 100644
--- a/api/src/main/java/com/google/appengine/setup/AppLogsWriter.java
+++ b/api/src/main/java/com/google/appengine/setup/AppLogsWriter.java
@@ -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;
@@ -62,8 +63,10 @@
* 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
+ * 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 {
@@ -87,6 +90,7 @@ class AppLogsWriter {
private int flushCount = 0;
private Future currentFlush;
private Stopwatch stopwatch;
+ private final ReentrantLock lock = new ReentrantLock();
/**
* Construct an AppLogsWriter instance.
@@ -147,99 +151,33 @@ 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()
- .setLevel(record.getLevel().ordinal())
- .setTimestampUsec(record.getTimestamp())
- .setMessage(record.getMessage())
- .build();
- int maxEncodingSize = 1000; // logLine.maxEncodingSize();
- Future 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 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();
}
}
@@ -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 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.
- *
- * 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 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;
- }
- }
- }
- }
-
- private synchronized void flushAndWaitLegacy() {
- waitForCurrentFlushLegacy();
- if (buffer.size() > 0) {
- currentFlush = doFlush();
- waitForCurrentFlushLegacy();
- }
- }
-
- private void waitForFlush(Future 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.
- *
- * 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) {
diff --git a/google3/third_party/java_src/appengine_standard/e2etest.md b/google3/third_party/java_src/appengine_standard/e2etest.md
deleted file mode 100644
index 8c5c5db27..000000000
--- a/google3/third_party/java_src/appengine_standard/e2etest.md
+++ /dev/null
@@ -1,301 +0,0 @@
-
-
-
-
-
-# App Engine Standard Java: Virtual Threads Concurrency & Multi-CL Architecture Plan
-
-This document details the end-to-end (E2E) verification plan for the Java 21+ Virtual Threads concurrency bug (b/514813839), explains how virtual thread scheduling behaves across App Engine connector modes, and provides the architectural mapping for all changelists in the current development workspace.
-
-> **Important Configuration Note:** App Engine Standard Java applications exclusively use `WEB-INF/appengine-web.xml` for all application configuration, instance class sizing, system properties, and environment variables. `app.yaml` is not used by App Engine Standard Java runtimes.
-
----
-
-## 1. Executive Summary & The Concurrency Bug (b/514813839)
-
-Under Java 21 (prior to JEP 491 / Java 24 where synchronized monitor pinning was eliminated), web applications running on App Engine Standard with virtual threads enabled via the system property `-Dappengine.use.virtualthreads=true` in `appengine-web.xml` suffer from severe latency degradation, carrier thread pool starvation, and container deadlocks on fractional or low-core instance classes (such as F1 and F2 instances with $\le 1024\text{ MB}$ RAM).
-
-### Root Causes
-
-1. **Carrier Pool Thrashing:** By default, `VirtualThreads.getDefaultVirtualThreadsExecutor()` delegates to `ForkJoinPool.commonPool()`, which sizes itself to physical host CPU core counts (often 64+ cores on container host machines). When container CPU share is throttled (e.g., 0.5 CPU on an F1 instance), 64+ active carrier threads fight for CPU slices, causing severe OS context switching thrashing.
-2. **Monitor Lock Carrier Pinning:** When application logs are written, `AppLogsWriter` flushes log batches via an asynchronous gRPC call (`ApiProxyDelegate.makeAsyncCall("logservice", "Flush")`). Previously, when calling `flushAndWait()`, the virtual thread invoked `slowFlush.get()` inside a synchronized monitor block (`synchronized (this)` in `api/setup`, and `synchronized (lock)` in `runtime/impl`). Because Java 21 virtual threads cannot unmount while holding a synchronized monitor, the virtual thread pinned its OS carrier thread while waiting for network I/O.
-
-### The Starvation Deadlock Under Load
-
-On an F1 instance where carrier threads are bounded or throttled:
-1. Request A calls `flushAndWait()` inside `synchronized`, pinning its carrier thread waiting for the gRPC network call.
-2. Request B arrives concurrently and tries to write a log (`addLogRecordAndMaybeFlush`), attempting to enter `synchronized (lock)` and blocking.
-3. Because Request B is blocked on the monitor and Request A has pinned the carrier thread waiting for network I/O, the entire container deadlocks, request queues overflow, and `500 Internal Server Error` / `502 Bad Gateway` timeouts occur.
-
----
-
-## 2. Global Workspace Plan: Multi-CL Architectural Relationship
-
-Our development workspace contains three active pending changelists and one submitted prerequisite changelist that work in concert to overhaul concurrency safety, thread reuse cleanliness, container wiring, and dependency hygiene across the App Engine Java runtime stack:
-
-```
-+---------------------------------------------------------------------------------------------------+
-| APP ENGINE STANDARD JAVA RUNTIME |
-+---------------------------------------------------------------------------------------------------+
-| [CL 946974028] (Active Pending) Virtual Threads Carrier Capping & Lock Decoupling (b/514813839) |
-| -> Bounds ForkJoinPool parallelism & moves gRPC I/O (waitForFlush) outside synchronized blocks. |
-+---------------------------------------------------------------------------------------------------+
-| [CL 946869615] (Active Pending) Datastore ThreadLocal Stack Eviction & Rollback Trap (b/494621464)|
-| -> Purges orphaned transaction state on pooled/reused threads (Jetty / Virtual Threads). |
-+---------------------------------------------------------------------------------------------------+
-| [CL 947510879] (Active Pending) Jetty 12 / 12.1 Jakarta EE 11 (ee11) & Servlet Wiring |
-| -> Modernizes the web server container hosting our virtual threads executor. |
-+---------------------------------------------------------------------------------------------------+
-| [CL 950649645] (Submitted / OCL 948873605) Third-Party Dependency Upgrades |
-| -> Upgrades unpinned libraries & removes deprecated bridges while preserving spec exclusions. |
-+---------------------------------------------------------------------------------------------------+
-```
-
-### Detailed Breakdown of Related CLs
-
-* **CL 946974028 (b/514813839) — Virtual Threads Carrier Capping & Lock Decoupling (Active Pending):**
- * **Action:** Implements dynamic carrier parallelism capping in `JavaRuntimeMain.configureVirtualThreadParallelism()` and `JettyServletEngineAdapter.start()`:
- * $\le 512\text{ MB}$ (F1 Class / 0.5 CPU) $\rightarrow$ 1 carrier core
- * $\le 1024\text{ MB}$ (F2 Class / 1 CPU) $\rightarrow$ 2 carrier cores
- * $> 1024\text{ MB}$ (F4 / Backend Classes) $\rightarrow$ 4 carrier cores
- * In `AppLogsWriter`, extracts `pendingFlush` inside the monitor and awaits `waitForFlush(pendingFlush)` strictly outside `synchronized`, preventing carrier pinning during gRPC network I/O. Also strips circular `logger.info()` recursion and prevents premature flush nulling.
-
-* **CL 946869615 (b/494621464) — Datastore Rollback Trap & ThreadLocal Stack Eviction (Active Pending):**
- * **Relationship:** Directly addresses thread reuse under concurrent request execution. In container environments where threads are pooled and reused across incoming requests (such as Jetty 12 `QueuedThreadPool` or virtual threads), if an application catches a transaction exception without rolling back, or if `BeginTransaction` fails, orphaned transaction state on `ThreadLocal` stacks poisons subsequent requests on that thread.
- * **Action:** Updates `doRollbackAsync()` in `InternalTransactionV3` to return a completed future when a transaction is inactive/failed, and adds lazy eviction in `TransactionStackImpl.peek()` to guarantee thread-local cleanliness under concurrent thread reuse.
-
-* **CL 947510879 — Jetty 12.1 / EE 11 Jakarta Servlet Wiring (Active Pending):**
- * **Relationship:** Modernizes the web server engine (`JettyServletEngineAdapter`) that hosts our virtual threads executor.
- * **Action:** Wires up EE 11 and Jakarta Servlet API support across `runtime-impl-jetty121.jar`, `runtime-shared-jetty121-ee11.jar`, and `runtime_servlets.jar`, resolving strict dependency errors and missing `ee11` symbol exports so modern Java 21+ applications run cleanly on Jetty 12.1.
-
-* **CL 950649645 (Submitted / OCL 948873605) — Dependency Upgrades & google-http-client-jackson Removal:**
- * **Relationship:** Ensures build reactor and runtime hygiene across all modules.
- * **Action:** Upgrades unpinned dependencies (Guava, Flogger, Protobuf, JUnit, Mockito, etc.) and removes the deprecated `google-http-client-jackson:1.29.2` bridge (replacing it with `google-http-client-jackson2:1.47.1`) while strictly protecting runtime spec boundaries (`` exclusions for Servlet API, Lucene 2.9.4, Jasper, etc.).
-
----
-
-## 3. Connector Architecture: RPC Mode vs. HTTP Connector Mode
-
-**Question:** Does `appengine.use.virtualthreads` work effectively in non-HTTP connector (RPC / CGI stubby) mode?
-**Answer:** Yes. Virtual threads operate identically across both RPC connector mode and HTTP connector mode.
-
-### How Request Dispatches Work in Jetty 12 / 12.1 (`JettyServletEngineAdapter`)
-
-1. **Connector Ingestion:**
- * **RPC Connector Mode (`appengine.use.HttpConnector=false` / default):** Incoming requests arrive from App Engine frontends via `DelegateConnector` (`rpcConnector`), which extends Jetty's `AbstractConnector`.
- * **HTTP Connector Mode (`appengine.use.HttpConnector=true`):** Incoming requests arrive directly over HTTP/HTTP2 via HTTP connectors configured by `AppVersionHandlerFactory`.
-2. **Shared Thread Pool Execution:**
- * Both connectors are attached to the same Jetty `Server` instance and share the server's `QueuedThreadPool` (`server.getThreadPool()`).
- * In `DelegateConnector.java`, incoming RPC requests are dispatched via `getExecutor().execute(runnable)`, which resolves directly to `server.getThreadPool()`.
- * When `appengine.use.virtualthreads=true` is set, `JettyServletEngineAdapter.start()` explicitly invokes `threadPool.setVirtualThreadsExecutor(virtualThreadsExecutor)`. Consequently, whether running over default RPC connectors or HTTP connectors, Jetty dispatches 100% of incoming request tasks onto virtual threads.
-
----
-
-## 4. Production E2E Load Test Runbook
-
-To definitively prove in production that the bug causes starvation/timeouts under load and that CL 946974028 eliminates it, follow this E2E testing protocol using the custom runtime bundling method documented in `TRYLATESTBITSINPROD.md`.
-
-### Step 1: Build Local Runtime Deployment Jars
-
-From your local checkout containing CL 946974028, build the runtime deployment artifacts:
-
-```bash
-./mvnw clean install -DskipTests
-```
-
-This produces the 7 core runtime jars under `runtime/deployment/target/runtime-deployment-*/`:
-* `runtime-impl-jetty12.jar`
-* `runtime-impl-jetty121.jar`
-* `runtime-main.jar`
-* `runtime-shared-jetty12.jar`
-* `runtime-shared-jetty12-ee10.jar`
-* `runtime-shared-jetty121-ee8.jar`
-* `runtime-shared-jetty121-ee11.jar`
-
-### Step 2: Configure the Test Application Harness
-
-Create a Java 21 App Engine Standard test application (e.g., Servlet or Spring Boot on Jetty 12/12.1) configured with an `F1` instance class in `WEB-INF/appengine-web.xml` (`F1`, giving $\le 512\text{ MB}$ RAM, which bounds the carrier pool to 1 thread under our fix).
-
-In your app's `pom.xml`, use `copy-rename-maven-plugin` (as detailed in `TRYLATESTBITSINPROD.md`) to copy the locally built runtime jars into `WEB-INF/lib/` and rename them to root names (e.g., `runtime-main.jar`).
-
-In `WEB-INF/appengine-web.xml`, configure the custom entrypoint and system properties:
-
-```xml
-
-
- java21
- F1
- true
-
-
-
-
-
-
- java
- --add-opens java.base/java.lang=ALL-UNNAMED
- --add-opens java.base/java.nio.charset=ALL-UNNAMED
- --add-opens java.base/java.util.concurrent=ALL-UNNAMED
- --add-opens java.logging/java.util.logging=ALL-UNNAMED
- -showversion -XX:+PrintCommandLineFlags
- -Djava.class.path=runtime-main.jar
- -Dclasspath.runtimebase=.:
- com/google/apphosting/runtime/JavaRuntimeMainWithDefaults
- --fixed_application_path=.
- .
-
-
-```
-
-Add a test servlet endpoint `/test-log-starvation` designed to force concurrent log buffering and async gRPC flushing.
-
-Depending on whether your App Engine application is running with modern **Jakarta EE** (default for Java 21 on Jetty 12/12.1, using `jakarta.servlet.*`) or legacy **Java EE 8** (enabled via `` in `appengine-web.xml`, using `javax.servlet.*`), implement the test servlet as follows:
-
-#### Jakarta EE (EE 10 / EE 11 — `jakarta.servlet.*`, Default for Java 21)
-
-```java
-package com.google.appengine.test;
-
-import jakarta.servlet.annotation.WebServlet;
-import jakarta.servlet.http.HttpServlet;
-import jakarta.servlet.http.HttpServletRequest;
-import jakarta.servlet.http.HttpServletResponse;
-import java.io.IOException;
-import java.util.Arrays;
-import java.util.logging.Logger;
-
-@WebServlet("/test-log-starvation")
-public class LogStarvationServlet extends HttpServlet {
- private static final Logger logger = Logger.getLogger(LogStarvationServlet.class.getName());
-
- @Override
- protected void doGet(HttpServletRequest req, HttpServletResponse resp) throws IOException {
- // 1. Generate a log payload larger than SMALL_FLUSH (~20KB) to trigger an async gRPC flush
- char[] chars = new char[25000];
- Arrays.fill(chars, 'x');
- logger.info("Load test log batch: " + new String(chars));
-
- // 2. Simulate brief application processing while the async gRPC log flush is in flight
- try {
- Thread.sleep(25);
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- }
-
- // 3. Request completion invokes RequestManager.clearEnvironmentForCurrentThread() -> ApiProxy.flushLogs(),
- // which calls AppLogsWriter.flushAndWait(). Under legacy locking, this blocks inside synchronized (lock).
- resp.setStatus(HttpServletResponse.SC_OK);
- resp.getWriter().write("OK");
- }
-}
-```
-
-#### Java EE 8 (EE 8 — `javax.servlet.*`, for apps with ``)
-
-```java
-package com.google.appengine.test;
-
-import java.io.IOException;
-import java.util.Arrays;
-import java.util.logging.Logger;
-import javax.servlet.annotation.WebServlet;
-import javax.servlet.http.HttpServlet;
-import javax.servlet.http.HttpServletRequest;
-import javax.servlet.http.HttpServletResponse;
-
-@WebServlet("/test-log-starvation")
-public class LogStarvationServletEE8 extends HttpServlet {
- private static final Logger logger = Logger.getLogger(LogStarvationServletEE8.class.getName());
-
- @Override
- protected void doGet(HttpServletRequest req, HttpServletResponse resp) throws IOException {
- char[] chars = new char[25000];
- Arrays.fill(chars, 'x');
- logger.info("Load test log batch: " + new String(chars));
-
- try {
- Thread.sleep(25);
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- }
-
- resp.setStatus(HttpServletResponse.SC_OK);
- resp.getWriter().write("OK");
- }
-}
-```
-
-### Step 3: Phase 1 — Prove the Bug (Baseline Load Test on Unpatched Runtime)
-
-Deploy the test application to production without the CL 946974028 fixes (or by deploying against standard production base runtime jars) with virtual threads enabled in `WEB-INF/appengine-web.xml`:
-
-```xml
-
-
-
-```
-
-*(Alternatively, configure via `` in `appengine-web.xml`).*
-
-Execute a concurrent load test using ApacheBench (`ab`), `hey`, or Locust:
-
-```bash
-# Send 50 concurrent requests continuously for 45 seconds
-hey -c 50 -z 45s https://.appspot.com/test-log-starvation
-```
-
-#### Observed Metrics & Failure Proof (The Bug)
-* **Container Deadlocks & Timeouts:** Throughput collapses. You will observe a high percentage of `500 Internal Server Error` and `502 Bad Gateway` responses as request queues overflow.
-* **Carrier Pinning:** In Cloud Logging / Sherlog trace analysis, request threads show long blocking times waiting on monitor acquisition inside `AppLogsWriter.flushAndWait()` while the single F1 carrier thread is pinned in `slowFlush.get()`.
-* **Instance Thrashing:** Cloud Monitoring shows extreme CPU throttling and container instance restarts due to health-check starvation.
-
----
-
-### Step 4: Phase 2 — Prove the Fix (Verification Load Test with CL 946974028)
-
-Deploy the test application booted with your custom bundled runtime jars containing CL 946974028. Execute the load test across both connector modes:
-
-#### Test A: Default RPC Connector Mode
-In `WEB-INF/appengine-web.xml`, verify default RPC mode:
-
-```xml
-
-
-
-
-```
-
-#### Test B: HTTP Connector Mode
-In `WEB-INF/appengine-web.xml`, switch to HTTP connector mode:
-
-```xml
-
-
-
-
-```
-
-Execute the exact same load test against both deployments:
-
-```bash
-hey -c 50 -z 45s https://.appspot.com/test-log-starvation
-```
-
-#### Observed Metrics & Verification Proof (The Fix)
-* **Zero Deadlocks & Zero Timeouts:** 100% of requests return `200 OK` without a single 500/502 error across both RPC and HTTP connector modes.
-* **Smooth, Stable Throughput:** Because `waitForFlush(pendingFlush)` executes outside `synchronized`, virtual threads unmount cleanly during gRPC network waits without pinning the 1-carrier F1 pool. Concurrent requests acquire the monitor immediately, buffer their logs, and maintain stable p95/p99 latency under load.
-* **Deterministic Parallelism:** Inspection of thread dumps or runtime telemetry confirms `ForkJoinPool` carrier threads remain bounded to 1 core on F1 instances, eliminating OS scheduling thrashing.
\ No newline at end of file
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..c379d9e14 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,60 @@ private void addLogLinesAndMaybeFlushLegacy(Iterable appLogLines) {
if (maxSecondsBetweenFlush > 0
&& stopwatch.elapsed().compareTo(Duration.ofSeconds(maxSecondsBetweenFlush)) >= 0) {
- waitForCurrentFlushAndStartNewFlushLegacy();
- }
- }
-
- @GuardedBy("lock")
- private void waitForCurrentFlushAndStartNewFlushLegacy() {
- waitForCurrentFlushLegacy();
- if (genericResponse.getAppLogCount() > 0) {
- currentFlush = doFlush();
+ waitForCurrentFlushAndStartNewFlush();
}
}
- @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.
+ * Starts an asynchronous flush. This method may block if flushes
+ * are backed up.
*/
@GuardedBy("lock")
- private @Nullable Future getPendingFlushLocked() {
- if (currentFlush != null && !currentFlush.isDone() && !currentFlush.isCancelled()) {
- return currentFlush;
+ private void waitForCurrentFlushAndStartNewFlush() {
+ waitForCurrentFlush();
+ if (genericResponse.getAppLogCount() > 0) {
+ currentFlush = doFlush();
}
- 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.
+ * Initiates a synchronous flush. This method will always block
+ * until any pending flushes and its own flush completes.
*/
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 +360,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..e5e4432d9 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,67 +463,7 @@ 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()))
.thenReturn(slowFlush);
@@ -578,69 +518,60 @@ 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()))
@@ -654,10 +585,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,9 +597,7 @@ 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(
() -> {
@@ -679,7 +607,6 @@ public void testCustomerIssue_carrierPoolStarvationUnderLegacyLocking() throws E
});
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 +614,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/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:
- *
- *
- * - F1 Class (<= 512MB, 0.5 CPU) -> 1 carrier thread
- *
- F2 Class (<= 1024MB, 1 CPU) -> 2 carrier threads
- *
- Backend/F4+ Classes (> 1024MB, 2+ CPU) -> 4 carrier threads
- *
- *
- * 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_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..3939e2416 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
@@ -17,7 +17,6 @@
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;
@@ -144,24 +143,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);
- }
-}