Skip to content

Commit 814ec4d

Browse files
harshachclaude
authored andcommitted
fix(governance): validate Flowable pool connections to survive DB failover (#28835)
The runtime Flowable engine builds its own MyBatis PooledDataSource from raw JDBC settings and does not validate pooled connections, so when the database drops a connection (Aurora/RDS failover, maintenance restart, or idle reaper) the async-executor poll threads (ResetExpiredJobsRunnable) keep borrowing the dead connection and fail with "PSQLException: FATAL: terminating connection due to administrator command". Enable MyBatis pool ping (SELECT 1) on the runtime engine so any connection idle past 30s is validated and transparently replaced before reuse. The migration path uses setDataSource() and is unaffected. Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent 7141ea6 commit 814ec4d

2 files changed

Lines changed: 82 additions & 0 deletions

File tree

openmetadata-service/src/main/java/org/openmetadata/service/governance/workflows/WorkflowHandler.java

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,13 @@ public class WorkflowHandler {
7676
@Getter private static volatile boolean initialized = false;
7777
private final boolean isMigrationContext;
7878

79+
private static final String CONNECTION_VALIDATION_QUERY = "SELECT 1";
80+
81+
// Validate any pooled connection idle longer than this before reuse. Kept below Flowable's
82+
// 60s reset-expired-jobs interval so the periodic async-executor threads always re-validate,
83+
// while connections in active sub-second use skip the check and pay no overhead.
84+
private static final int CONNECTION_PING_NOT_USED_FOR_MILLIS = 30000;
85+
7986
private WorkflowHandler(OpenMetadataApplicationConfig config, boolean isMigrationContext) {
8087
this.isMigrationContext = isMigrationContext;
8188
StandaloneProcessEngineConfiguration processEngineConfiguration =
@@ -179,6 +186,7 @@ public void initializeNewProcessEngine(
179186
processEngineConfiguration.setJdbcPassword(
180187
currentProcessEngineConfiguration.getJdbcPassword());
181188
processEngineConfiguration.setJdbcDriver(currentProcessEngineConfiguration.getJdbcDriver());
189+
configureConnectionPoolHealthChecks(processEngineConfiguration);
182190
}
183191
processEngineConfiguration.setDatabaseType(currentProcessEngineConfiguration.getDatabaseType());
184192
processEngineConfiguration.setDatabaseSchemaUpdate(
@@ -239,6 +247,25 @@ public void initializeNewProcessEngine(
239247
.addMapper(SqlMapper.class);
240248
}
241249

250+
/**
251+
* Enable connection-pool health checks on the runtime engine's MyBatis pool.
252+
*
253+
* <p>Unlike the application's HikariCP pool, the Flowable runtime engine builds its own MyBatis
254+
* {@code PooledDataSource} from the raw JDBC settings and does not validate pooled connections.
255+
* When the database drops a connection out from under the pool — an Aurora/RDS failover, a
256+
* maintenance restart, or an idle-connection reaper, all of which surface as {@code
257+
* PSQLException: terminating connection due to administrator command} — the async executor's
258+
* polling threads (e.g. {@code ResetExpiredJobsRunnable}) keep borrowing the dead connection and
259+
* failing until the pool happens to recycle it. Pool ping runs a lightweight validation query on
260+
* any connection idle past the threshold and transparently replaces it before handing it out.
261+
*/
262+
private static void configureConnectionPoolHealthChecks(
263+
StandaloneProcessEngineConfiguration processEngineConfiguration) {
264+
processEngineConfiguration.setJdbcPingEnabled(true);
265+
processEngineConfiguration.setJdbcPingQuery(CONNECTION_VALIDATION_QUERY);
266+
processEngineConfiguration.setJdbcPingConnectionNotUsedFor(CONNECTION_PING_NOT_USED_FOR_MILLIS);
267+
}
268+
242269
public static void initialize(OpenMetadataApplicationConfig config) {
243270
initialize(config, false);
244271
}

openmetadata-service/src/test/java/org/openmetadata/service/governance/workflows/WorkflowHandlerSchemaUpdateTest.java

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,11 +18,13 @@
1818
import static org.junit.jupiter.api.Assertions.assertThrows;
1919
import static org.junit.jupiter.api.Assertions.assertTrue;
2020
import static org.mockito.ArgumentMatchers.any;
21+
import static org.mockito.ArgumentMatchers.anyBoolean;
2122
import static org.mockito.Mockito.RETURNS_DEEP_STUBS;
2223
import static org.mockito.Mockito.lenient;
2324
import static org.mockito.Mockito.mock;
2425
import static org.mockito.Mockito.mockConstruction;
2526
import static org.mockito.Mockito.mockStatic;
27+
import static org.mockito.Mockito.never;
2628
import static org.mockito.Mockito.verify;
2729
import static org.mockito.Mockito.when;
2830

@@ -165,6 +167,59 @@ void runtimeModeSetsDbSchemaUpdateFalse() {
165167
}
166168
}
167169

170+
@Test
171+
void runtimeModeEnablesConnectionPoolHealthChecks() {
172+
try (MockedConstruction<StandaloneProcessEngineConfiguration> engineMock =
173+
mockConstruction(
174+
StandaloneProcessEngineConfiguration.class,
175+
(mock, ctx) ->
176+
when(mock.buildProcessEngine())
177+
.thenThrow(new FlowableWrongDbException("7.2.0.2", "7.1.0.0")));
178+
MockedStatic<ProcessEngines> ignored = mockStatic(ProcessEngines.class);
179+
MockedStatic<Entity> entityMock = mockStatic(Entity.class);
180+
MockedStatic<PipelineServiceClientFactory> pscMock =
181+
mockStatic(PipelineServiceClientFactory.class)) {
182+
183+
setupEntityMock(entityMock);
184+
pscMock
185+
.when(() -> PipelineServiceClientFactory.createPipelineServiceClient(any()))
186+
.thenReturn(null);
187+
188+
assertThrows(
189+
IllegalStateException.class, () -> WorkflowHandler.initialize(buildMockConfig(), false));
190+
191+
StandaloneProcessEngineConfiguration engineConfig = engineMock.constructed().getLast();
192+
verify(engineConfig).setJdbcPingEnabled(true);
193+
verify(engineConfig).setJdbcPingQuery("SELECT 1");
194+
verify(engineConfig).setJdbcPingConnectionNotUsedFor(30000);
195+
}
196+
}
197+
198+
@Test
199+
void migrationModeDoesNotEnableConnectionPoolPing() {
200+
ProcessEngine mockEngine = mock(ProcessEngine.class, RETURNS_DEEP_STUBS);
201+
202+
try (MockedConstruction<StandaloneProcessEngineConfiguration> engineMock =
203+
mockConstruction(
204+
StandaloneProcessEngineConfiguration.class,
205+
(mock, ctx) -> when(mock.buildProcessEngine()).thenReturn(mockEngine));
206+
MockedStatic<ProcessEngines> ignored = mockStatic(ProcessEngines.class);
207+
MockedStatic<Entity> entityMock = mockStatic(Entity.class);
208+
MockedStatic<PipelineServiceClientFactory> pscMock =
209+
mockStatic(PipelineServiceClientFactory.class)) {
210+
211+
setupEntityMock(entityMock);
212+
pscMock
213+
.when(() -> PipelineServiceClientFactory.createPipelineServiceClient(any()))
214+
.thenReturn(null);
215+
216+
WorkflowHandler.initialize(buildMockConfig(), true);
217+
218+
StandaloneProcessEngineConfiguration engineConfig = engineMock.constructed().getLast();
219+
verify(engineConfig, never()).setJdbcPingEnabled(anyBoolean());
220+
}
221+
}
222+
168223
// ── Helpers ──────────────────────────────────────────────────────────────────
169224

170225
private void setupEntityMock(MockedStatic<Entity> entityMock) {

0 commit comments

Comments
 (0)