blob: 0b681f72641360ecf344170cf8310d0c6c7c42df [file] [log] [blame]
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you 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
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* See the License for the specific language governing permissions and
* limitations under the License.
package org.apache.oozie.util.db;
import org.apache.oozie.util.XLog;
import javax.annotation.Nullable;
import javax.persistence.PersistenceException;
import java.sql.Array;
import java.sql.Blob;
import java.sql.CallableStatement;
import java.sql.Clob;
import java.sql.Connection;
import java.sql.DatabaseMetaData;
import java.sql.NClob;
import java.sql.PreparedStatement;
import java.sql.SQLClientInfoException;
import java.sql.SQLException;
import java.sql.SQLWarning;
import java.sql.SQLXML;
import java.sql.Savepoint;
import java.sql.Statement;
import java.sql.Struct;
import java.util.Map;
import java.util.Properties;
import java.util.Set;
import java.util.concurrent.Executor;
import java.util.function.Predicate;
public class FailingConnectionWrapper implements Connection {
private static final XLog LOG = XLog.getLog(FailingConnectionWrapper.class);
private final Connection delegate;
private RuntimeExceptionInjector<PersistenceException> injector;
private Predicate<String> predicate;
public FailingConnectionWrapper(final Connection delegate, final int failurePercent,
@Nullable final Predicate<String> predicate) {
this.delegate = delegate;
injector = new RuntimeExceptionInjector<>(PersistenceException.class, failurePercent);
if (predicate == null) {
this.predicate = new OozieDmlStatementPredicate();
} else {
this.predicate = predicate;
public Statement createStatement() throws SQLException {
return delegate.createStatement();
public PreparedStatement prepareStatement(final String sql) throws SQLException {
return delegate.prepareStatement(sql);
public CallableStatement prepareCall(final String sql) throws SQLException {
return delegate.prepareCall(sql);
public String nativeSQL(final String sql) throws SQLException {
return delegate.nativeSQL(sql);
public void setAutoCommit(final boolean autoCommit) throws SQLException {
public boolean getAutoCommit() throws SQLException {
return delegate.getAutoCommit();
public void commit() throws SQLException {
public void rollback() throws SQLException {
public void close() throws SQLException {
public boolean isClosed() throws SQLException {
return delegate.isClosed();
public DatabaseMetaData getMetaData() throws SQLException {
return delegate.getMetaData();
public void setReadOnly(final boolean readOnly) throws SQLException {
public boolean isReadOnly() throws SQLException {
return delegate.isReadOnly();
public void setCatalog(final String catalog) throws SQLException {
public String getCatalog() throws SQLException {
return delegate.getCatalog();
public void setTransactionIsolation(final int level) throws SQLException {
public int getTransactionIsolation() throws SQLException {
return delegate.getTransactionIsolation();
public SQLWarning getWarnings() throws SQLException {
return delegate.getWarnings();
public void clearWarnings() throws SQLException {
public Statement createStatement(final int resultSetType, final int resultSetConcurrency) throws SQLException {
return delegate.createStatement(resultSetType, resultSetConcurrency);
public PreparedStatement prepareStatement(final String sql, final int resultSetType, final int resultSetConcurrency)
throws SQLException {
if (predicate.test(sql)) {
LOG.trace("Injecting random failure. Preparing this statement might fail.");
injector.inject(String.format("Deliberately failing to prepare statement. [sql=%s]", sql));
LOG.trace("Preparing statement. [sql={0}]", sql);
return delegate.prepareStatement(sql, resultSetType, resultSetConcurrency);
public CallableStatement prepareCall(final String sql, final int resultSetType, final int resultSetConcurrency)
throws SQLException {
return delegate.prepareCall(sql, resultSetType, resultSetConcurrency);
public Map<String, Class<?>> getTypeMap() throws SQLException {
return delegate.getTypeMap();
public void setTypeMap(final Map<String, Class<?>> map) throws SQLException {
public void setHoldability(final int holdability) throws SQLException {
public int getHoldability() throws SQLException {
return delegate.getHoldability();
public Savepoint setSavepoint() throws SQLException {
return delegate.setSavepoint();
public Savepoint setSavepoint(final String name) throws SQLException {
return delegate.setSavepoint(name);
public void rollback(final Savepoint savepoint) throws SQLException {
public void releaseSavepoint(final Savepoint savepoint) throws SQLException {
public Statement createStatement(final int resultSetType, final int resultSetConcurrency, final int resultSetHoldability)
throws SQLException {
return delegate.createStatement(resultSetType, resultSetConcurrency, resultSetHoldability);
public PreparedStatement prepareStatement(final String sql, final int resultSetType, final int resultSetConcurrency,
final int resultSetHoldability) throws SQLException {
return delegate.prepareStatement(sql, resultSetType, resultSetConcurrency, resultSetHoldability);
public CallableStatement prepareCall(final String sql, final int resultSetType, final int resultSetConcurrency,
final int resultSetHoldability) throws SQLException {
return delegate.prepareCall(sql, resultSetType, resultSetConcurrency, resultSetHoldability);
public PreparedStatement prepareStatement(final String sql, final int autoGeneratedKeys) throws SQLException {
return delegate.prepareStatement(sql, autoGeneratedKeys);
public PreparedStatement prepareStatement(final String sql, final int[] columnIndexes) throws SQLException {
return delegate.prepareStatement(sql, columnIndexes);
public PreparedStatement prepareStatement(final String sql, final String[] columnNames) throws SQLException {
return delegate.prepareStatement(sql, columnNames);
public Clob createClob() throws SQLException {
return delegate.createClob();
public Blob createBlob() throws SQLException {
return delegate.createBlob();
public NClob createNClob() throws SQLException {
return delegate.createNClob();
public SQLXML createSQLXML() throws SQLException {
return delegate.createSQLXML();
public boolean isValid(final int timeout) throws SQLException {
return delegate.isValid(timeout);
public void setClientInfo(final String name, final String value) throws SQLClientInfoException {
delegate.setClientInfo(name, value);
public void setClientInfo(final Properties properties) throws SQLClientInfoException {
public String getClientInfo(final String name) throws SQLException {
return delegate.getClientInfo(name);
public Properties getClientInfo() throws SQLException {
return delegate.getClientInfo();
public Array createArrayOf(final String typeName, final Object[] elements) throws SQLException {
return delegate.createArrayOf(typeName, elements);
public Struct createStruct(final String typeName, final Object[] attributes) throws SQLException {
return delegate.createStruct(typeName, attributes);
public void setSchema(final String schema) throws SQLException {
public String getSchema() throws SQLException {
return delegate.getSchema();
public void abort(final Executor executor) throws SQLException {
public void setNetworkTimeout(final Executor executor, final int milliseconds) throws SQLException {
delegate.setNetworkTimeout(executor, milliseconds);
public int getNetworkTimeout() throws SQLException {
return delegate.getNetworkTimeout();
public <T> T unwrap(final Class<T> iface) throws SQLException {
return delegate.unwrap(iface);
public boolean isWrapperFor(final Class<?> iface) throws SQLException {
return delegate.isWrapperFor(iface);
static class OozieDmlStatementPredicate implements Predicate<String> {
private static final Set<String> DML_PREFIXES = Sets.newHashSet(
private static final Set<String> OOZIE_TABLE_NAMES = Sets.newHashSet(
public boolean test(@Nullable String input) {
boolean isDmlStatement = false;
for (final String dmlPrefix : DML_PREFIXES) {
if (input.toUpperCase().startsWith(dmlPrefix)) {
isDmlStatement = true;
boolean isOozieTable = false;
for (final String oozieTableName : OOZIE_TABLE_NAMES) {
if (input.toUpperCase().contains(oozieTableName)) {
isOozieTable = true;
return isDmlStatement && isOozieTable;