001/** 002 * Licensed to the Apache Software Foundation (ASF) under one 003 * or more contributor license agreements. See the NOTICE file 004 * distributed with this work for additional information 005 * regarding copyright ownership. The ASF licenses this file 006 * to you under the Apache License, Version 2.0 (the 007 * "License"); you may not use this file except in compliance 008 * with the License. You may obtain a copy of the License at 009 * 010 * http://www.apache.org/licenses/LICENSE-2.0 011 * 012 * Unless required by applicable law or agreed to in writing, software 013 * distributed under the License is distributed on an "AS IS" BASIS, 014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. 015 * See the License for the specific language governing permissions and 016 * limitations under the License. 017 */ 018 019package org.apache.oozie.util.db; 020 021import com.google.common.base.Preconditions; 022import com.google.common.base.Predicate; 023import com.google.common.collect.Sets; 024import org.apache.directory.api.util.Strings; 025import org.apache.oozie.util.XLog; 026 027import javax.annotation.Nullable; 028import javax.persistence.PersistenceException; 029import java.sql.Array; 030import java.sql.Blob; 031import java.sql.CallableStatement; 032import java.sql.Clob; 033import java.sql.Connection; 034import java.sql.DatabaseMetaData; 035import java.sql.NClob; 036import java.sql.PreparedStatement; 037import java.sql.SQLClientInfoException; 038import java.sql.SQLException; 039import java.sql.SQLWarning; 040import java.sql.SQLXML; 041import java.sql.Savepoint; 042import java.sql.Statement; 043import java.sql.Struct; 044import java.util.Map; 045import java.util.Properties; 046import java.util.Set; 047import java.util.concurrent.Executor; 048 049public class FailingConnectionWrapper implements Connection { 050 private static final XLog LOG = XLog.getLog(FailingConnectionWrapper.class); 051 052 private final Connection delegate; 053 private RuntimeExceptionInjector<PersistenceException> injector; 054 private Predicate<String> predicate; 055 056 public FailingConnectionWrapper(final Connection delegate, final int failurePercent, 057 @Nullable final Predicate<String> predicate) { 058 this.delegate = delegate; 059 injector = new RuntimeExceptionInjector<>(PersistenceException.class, failurePercent); 060 if (predicate == null) { 061 this.predicate = new OozieDmlStatementPredicate(); 062 } else { 063 this.predicate = predicate; 064 } 065 } 066 067 @Override 068 public Statement createStatement() throws SQLException { 069 return delegate.createStatement(); 070 } 071 072 @Override 073 public PreparedStatement prepareStatement(final String sql) throws SQLException { 074 return delegate.prepareStatement(sql); 075 } 076 077 @Override 078 public CallableStatement prepareCall(final String sql) throws SQLException { 079 return delegate.prepareCall(sql); 080 } 081 082 @Override 083 public String nativeSQL(final String sql) throws SQLException { 084 return delegate.nativeSQL(sql); 085 } 086 087 @Override 088 public void setAutoCommit(final boolean autoCommit) throws SQLException { 089 delegate.setAutoCommit(autoCommit); 090 } 091 092 @Override 093 public boolean getAutoCommit() throws SQLException { 094 return delegate.getAutoCommit(); 095 } 096 097 @Override 098 public void commit() throws SQLException { 099 delegate.commit(); 100 } 101 102 @Override 103 public void rollback() throws SQLException { 104 delegate.rollback(); 105 } 106 107 @Override 108 public void close() throws SQLException { 109 delegate.close(); 110 } 111 112 @Override 113 public boolean isClosed() throws SQLException { 114 return delegate.isClosed(); 115 } 116 117 @Override 118 public DatabaseMetaData getMetaData() throws SQLException { 119 return delegate.getMetaData(); 120 } 121 122 @Override 123 public void setReadOnly(final boolean readOnly) throws SQLException { 124 delegate.setReadOnly(readOnly); 125 } 126 127 @Override 128 public boolean isReadOnly() throws SQLException { 129 return delegate.isReadOnly(); 130 } 131 132 @Override 133 public void setCatalog(final String catalog) throws SQLException { 134 delegate.setCatalog(catalog); 135 } 136 137 @Override 138 public String getCatalog() throws SQLException { 139 return delegate.getCatalog(); 140 } 141 142 @Override 143 public void setTransactionIsolation(final int level) throws SQLException { 144 delegate.setTransactionIsolation(level); 145 } 146 147 @Override 148 public int getTransactionIsolation() throws SQLException { 149 return delegate.getTransactionIsolation(); 150 } 151 152 @Override 153 public SQLWarning getWarnings() throws SQLException { 154 return delegate.getWarnings(); 155 } 156 157 @Override 158 public void clearWarnings() throws SQLException { 159 delegate.clearWarnings(); 160 } 161 162 @Override 163 public Statement createStatement(final int resultSetType, final int resultSetConcurrency) throws SQLException { 164 return delegate.createStatement(resultSetType, resultSetConcurrency); 165 } 166 167 @Override 168 public PreparedStatement prepareStatement(final String sql, final int resultSetType, final int resultSetConcurrency) 169 throws SQLException { 170 if (predicate.apply(sql)) { 171 LOG.trace("Injecting random failure. Preparing this statement might fail."); 172 injector.inject(String.format("Deliberately failing to prepare statement. [sql=%s]", sql)); 173 } 174 175 LOG.trace("Preparing statement. [sql={0}]", sql); 176 return delegate.prepareStatement(sql, resultSetType, resultSetConcurrency); 177 } 178 179 @Override 180 public CallableStatement prepareCall(final String sql, final int resultSetType, final int resultSetConcurrency) 181 throws SQLException { 182 return delegate.prepareCall(sql, resultSetType, resultSetConcurrency); 183 } 184 185 @Override 186 public Map<String, Class<?>> getTypeMap() throws SQLException { 187 return delegate.getTypeMap(); 188 } 189 190 @Override 191 public void setTypeMap(final Map<String, Class<?>> map) throws SQLException { 192 delegate.setTypeMap(map); 193 } 194 195 @Override 196 public void setHoldability(final int holdability) throws SQLException { 197 delegate.setHoldability(holdability); 198 } 199 200 @Override 201 public int getHoldability() throws SQLException { 202 return delegate.getHoldability(); 203 } 204 205 @Override 206 public Savepoint setSavepoint() throws SQLException { 207 return delegate.setSavepoint(); 208 } 209 210 @Override 211 public Savepoint setSavepoint(final String name) throws SQLException { 212 return delegate.setSavepoint(name); 213 } 214 215 @Override 216 public void rollback(final Savepoint savepoint) throws SQLException { 217 delegate.rollback(); 218 } 219 220 @Override 221 public void releaseSavepoint(final Savepoint savepoint) throws SQLException { 222 delegate.releaseSavepoint(savepoint); 223 } 224 225 @Override 226 public Statement createStatement(final int resultSetType, final int resultSetConcurrency, final int resultSetHoldability) 227 throws SQLException { 228 return delegate.createStatement(resultSetType, resultSetConcurrency, resultSetHoldability); 229 } 230 231 @Override 232 public PreparedStatement prepareStatement(final String sql, final int resultSetType, final int resultSetConcurrency, 233 final int resultSetHoldability) throws SQLException { 234 return delegate.prepareStatement(sql, resultSetType, resultSetConcurrency, resultSetHoldability); 235 } 236 237 @Override 238 public CallableStatement prepareCall(final String sql, final int resultSetType, final int resultSetConcurrency, 239 final int resultSetHoldability) throws SQLException { 240 return delegate.prepareCall(sql, resultSetType, resultSetConcurrency, resultSetHoldability); 241 } 242 243 @Override 244 public PreparedStatement prepareStatement(final String sql, final int autoGeneratedKeys) throws SQLException { 245 return delegate.prepareStatement(sql, autoGeneratedKeys); 246 } 247 248 @Override 249 public PreparedStatement prepareStatement(final String sql, final int[] columnIndexes) throws SQLException { 250 return delegate.prepareStatement(sql, columnIndexes); 251 } 252 253 @Override 254 public PreparedStatement prepareStatement(final String sql, final String[] columnNames) throws SQLException { 255 return delegate.prepareStatement(sql, columnNames); 256 } 257 258 @Override 259 public Clob createClob() throws SQLException { 260 return delegate.createClob(); 261 } 262 263 @Override 264 public Blob createBlob() throws SQLException { 265 return delegate.createBlob(); 266 } 267 268 @Override 269 public NClob createNClob() throws SQLException { 270 return delegate.createNClob(); 271 } 272 273 @Override 274 public SQLXML createSQLXML() throws SQLException { 275 return delegate.createSQLXML(); 276 } 277 278 @Override 279 public boolean isValid(final int timeout) throws SQLException { 280 return delegate.isValid(timeout); 281 } 282 283 @Override 284 public void setClientInfo(final String name, final String value) throws SQLClientInfoException { 285 delegate.setClientInfo(name, value); 286 } 287 288 @Override 289 public void setClientInfo(final Properties properties) throws SQLClientInfoException { 290 delegate.setClientInfo(properties); 291 } 292 293 @Override 294 public String getClientInfo(final String name) throws SQLException { 295 return delegate.getClientInfo(name); 296 } 297 298 @Override 299 public Properties getClientInfo() throws SQLException { 300 return delegate.getClientInfo(); 301 } 302 303 @Override 304 public Array createArrayOf(final String typeName, final Object[] elements) throws SQLException { 305 return delegate.createArrayOf(typeName, elements); 306 } 307 308 @Override 309 public Struct createStruct(final String typeName, final Object[] attributes) throws SQLException { 310 return delegate.createStruct(typeName, attributes); 311 } 312 313 @Override 314 public void setSchema(final String schema) throws SQLException { 315 delegate.setSchema(schema); 316 } 317 318 @Override 319 public String getSchema() throws SQLException { 320 return delegate.getSchema(); 321 } 322 323 @Override 324 public void abort(final Executor executor) throws SQLException { 325 delegate.abort(executor); 326 } 327 328 @Override 329 public void setNetworkTimeout(final Executor executor, final int milliseconds) throws SQLException { 330 delegate.setNetworkTimeout(executor, milliseconds); 331 } 332 333 @Override 334 public int getNetworkTimeout() throws SQLException { 335 return delegate.getNetworkTimeout(); 336 } 337 338 @Override 339 public <T> T unwrap(final Class<T> iface) throws SQLException { 340 return delegate.unwrap(iface); 341 } 342 343 @Override 344 public boolean isWrapperFor(final Class<?> iface) throws SQLException { 345 return delegate.isWrapperFor(iface); 346 } 347 348 static class OozieDmlStatementPredicate implements Predicate<String> { 349 private static final Set<String> DML_PREFIXES = Sets.newHashSet( 350 "SELECT ", "INSERT INTO ", "UPDATE ", "DELETE FROM "); 351 private static final Set<String> OOZIE_TABLE_NAMES = Sets.newHashSet( 352 "BUNDLE_ACTIONS", "BUNDLE_JOBS", "COORD_ACTIONS", "COORD_JOBS", "SLA_REGISTRATION", "SLA_SUMMARY", 353 "WF_ACTIONS", "WF_JOBS"); 354 355 @Override 356 public boolean apply(@Nullable String input) { 357 Preconditions.checkArgument(Strings.isNotEmpty(input)); 358 359 boolean isDmlStatement = false; 360 for (final String dmlPrefix : DML_PREFIXES) { 361 if (input.toUpperCase().startsWith(dmlPrefix)) { 362 isDmlStatement = true; 363 } 364 } 365 366 boolean isOozieTable = false; 367 for (final String oozieTableName : OOZIE_TABLE_NAMES) { 368 if (input.toUpperCase().contains(oozieTableName)) { 369 isOozieTable = true; 370 } 371 } 372 373 return isDmlStatement && isOozieTable; 374 } 375 } 376}