Chained Transaction Manager
This commit is contained in:
@@ -0,0 +1,142 @@
|
||||
package org.springframework.data.graph.neo4j.support;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.springframework.transaction.*;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collection;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* @author mh
|
||||
* @since 14.02.11
|
||||
*/
|
||||
public class ChainedTransactionManager implements PlatformTransactionManager {
|
||||
|
||||
protected Log logger = LogFactory.getLog(getClass());
|
||||
|
||||
private final List<PlatformTransactionManager> transactionManagers = new ArrayList<PlatformTransactionManager>();
|
||||
private SynchronizationManager sychronizationManager ;
|
||||
|
||||
ChainedTransactionManager(SynchronizationManager sychronizationManager) {
|
||||
this.sychronizationManager = sychronizationManager;
|
||||
}
|
||||
|
||||
public ChainedTransactionManager() {
|
||||
this(new DefaultSynchronizationManager());
|
||||
}
|
||||
|
||||
public void setTransactionManagers(List<PlatformTransactionManager> transactionManagers) {
|
||||
this.transactionManagers.addAll(transactionManagers);
|
||||
}
|
||||
|
||||
|
||||
|
||||
|
||||
@Override
|
||||
public MultiTransactionStatus getTransaction(TransactionDefinition definition) throws TransactionException {
|
||||
|
||||
MultiTransactionStatus mts = new MultiTransactionStatus(transactionManagers.get(0)/*First TM is main TM*/);
|
||||
|
||||
if (!sychronizationManager.isSynchronizationActive()) {
|
||||
sychronizationManager.initSynchronization();
|
||||
mts.setNewSynchonization();
|
||||
}
|
||||
|
||||
for (PlatformTransactionManager transactionManager : transactionManagers) {
|
||||
mts.registerTransactionManager(definition, transactionManager);
|
||||
}
|
||||
|
||||
return mts;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void commit(TransactionStatus status) throws TransactionException {
|
||||
|
||||
MultiTransactionStatus multiTransactionStatus = (MultiTransactionStatus) status;
|
||||
|
||||
boolean commit = true;
|
||||
Exception commitException = null;
|
||||
PlatformTransactionManager commitExceptionTransactionManager = null;
|
||||
|
||||
for (PlatformTransactionManager transactionManager : reverse(transactionManagers)) {
|
||||
if (commit) {
|
||||
try {
|
||||
multiTransactionStatus.commit(transactionManager);
|
||||
} catch (Exception ex) {
|
||||
commit = false;
|
||||
commitException = ex;
|
||||
commitExceptionTransactionManager = transactionManager;
|
||||
}
|
||||
} else {
|
||||
//after unsucessfull commit we must try to rollback remaining transaction managers
|
||||
try {
|
||||
multiTransactionStatus.rollback(transactionManager);
|
||||
} catch (Exception ex) {
|
||||
logger.warn("Rollback exception (after commit) (" + transactionManager + ") " + ex.getMessage(), ex);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (multiTransactionStatus.isNewSynchonization()){
|
||||
sychronizationManager.clearSynchronization();
|
||||
}
|
||||
|
||||
if (commitException != null) {
|
||||
boolean firstTransactionManagerFailed = commitExceptionTransactionManager == getLastTransactionManager();
|
||||
int transactionState = firstTransactionManagerFailed ? HeuristicCompletionException.STATE_ROLLED_BACK : HeuristicCompletionException.STATE_MIXED;
|
||||
throw new HeuristicCompletionException(transactionState, commitException);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void rollback(TransactionStatus status) throws TransactionException {
|
||||
|
||||
Exception rollbackException = null;
|
||||
PlatformTransactionManager rollbackExceptionTransactionManager = null;
|
||||
|
||||
|
||||
MultiTransactionStatus multiTransactionStatus = (MultiTransactionStatus) status;
|
||||
|
||||
for (PlatformTransactionManager transactionManager : reverse(transactionManagers)) {
|
||||
try {
|
||||
multiTransactionStatus.rollback(transactionManager);
|
||||
} catch (Exception ex) {
|
||||
if (rollbackException == null) {
|
||||
rollbackException = ex;
|
||||
rollbackExceptionTransactionManager = transactionManager;
|
||||
} else {
|
||||
logger.warn("Rollback exception (" + transactionManager + ") " + ex.getMessage(), ex);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (multiTransactionStatus.isNewSynchonization()){
|
||||
sychronizationManager.clearSynchronization();
|
||||
}
|
||||
|
||||
if (rollbackException != null) {
|
||||
throw new UnexpectedRollbackException("Rollback exception, originated at ("+rollbackExceptionTransactionManager+") "+
|
||||
rollbackException.getMessage(), rollbackException);
|
||||
}
|
||||
}
|
||||
|
||||
private <T> Iterable<T> reverse(Collection<T> collection) {
|
||||
List<T> list = new ArrayList<T>(collection);
|
||||
Collections.reverse(list);
|
||||
return list;
|
||||
}
|
||||
|
||||
|
||||
private PlatformTransactionManager getLastTransactionManager() {
|
||||
return transactionManagers.get(lastTransactionManagerIndex());
|
||||
}
|
||||
|
||||
private int lastTransactionManagerIndex() {
|
||||
return transactionManagers.size() - 1;
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,24 @@
|
||||
package org.springframework.data.graph.neo4j.support;
|
||||
|
||||
import org.springframework.transaction.support.TransactionSynchronizationManager;
|
||||
|
||||
/**
|
||||
* @author mh
|
||||
* @since 15.02.11
|
||||
*/
|
||||
public class DefaultSynchronizationManager implements SynchronizationManager {
|
||||
@Override
|
||||
public void initSynchronization() {
|
||||
TransactionSynchronizationManager.initSynchronization();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isSynchronizationActive() {
|
||||
return TransactionSynchronizationManager.isSynchronizationActive();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void clearSynchronization() {
|
||||
TransactionSynchronizationManager.clear();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,151 @@
|
||||
package org.springframework.data.graph.neo4j.support;
|
||||
|
||||
import org.springframework.transaction.PlatformTransactionManager;
|
||||
import org.springframework.transaction.TransactionDefinition;
|
||||
import org.springframework.transaction.TransactionException;
|
||||
import org.springframework.transaction.TransactionStatus;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* @author mh
|
||||
* @since 14.02.11
|
||||
*/
|
||||
public class MultiTransactionStatus implements TransactionStatus {
|
||||
|
||||
|
||||
private PlatformTransactionManager mainTransactionManager;
|
||||
|
||||
private Map<PlatformTransactionManager, TransactionStatus> transactionStatuses =
|
||||
Collections.synchronizedMap(new HashMap<PlatformTransactionManager, TransactionStatus>());
|
||||
|
||||
private boolean newSynchonization;
|
||||
|
||||
public MultiTransactionStatus(PlatformTransactionManager mainTransactionManager) {
|
||||
this.mainTransactionManager = mainTransactionManager;
|
||||
}
|
||||
|
||||
|
||||
private Map<PlatformTransactionManager, TransactionStatus> getTransactionStatuses() {
|
||||
return transactionStatuses;
|
||||
}
|
||||
|
||||
private TransactionStatus getMainTransactionStatus() {
|
||||
return transactionStatuses.get(mainTransactionManager);
|
||||
}
|
||||
|
||||
|
||||
public void setNewSynchonization() {
|
||||
this.newSynchonization = true;
|
||||
}
|
||||
|
||||
public boolean isNewSynchonization() {
|
||||
return newSynchonization;
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public boolean isNewTransaction() {
|
||||
return getMainTransactionStatus().isNewTransaction();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasSavepoint() {
|
||||
return getMainTransactionStatus().hasSavepoint();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setRollbackOnly() {
|
||||
for(TransactionStatus ts : transactionStatuses.values() ){
|
||||
ts.setRollbackOnly();
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isRollbackOnly() {
|
||||
return getMainTransactionStatus().isRollbackOnly();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isCompleted() {
|
||||
return getMainTransactionStatus().isCompleted();
|
||||
}
|
||||
|
||||
|
||||
private static class SavePoints {
|
||||
Map<TransactionStatus,Object> savepoints=new HashMap<TransactionStatus, Object>();
|
||||
|
||||
private void addSavePoint(TransactionStatus status, Object savepoint) {
|
||||
this.savepoints.put(status, savepoint);
|
||||
}
|
||||
|
||||
private void save(TransactionStatus transactionStatus) {
|
||||
Object savepoint = transactionStatus.createSavepoint();
|
||||
addSavePoint(transactionStatus, savepoint);
|
||||
}
|
||||
|
||||
|
||||
public void rollback() {
|
||||
for (TransactionStatus transactionStatus : savepoints.keySet()) {
|
||||
transactionStatus.rollbackToSavepoint(savepointFor(transactionStatus));
|
||||
}
|
||||
}
|
||||
|
||||
private Object savepointFor(TransactionStatus transactionStatus) {
|
||||
return savepoints.get(transactionStatus);
|
||||
}
|
||||
|
||||
public void release() {
|
||||
for (TransactionStatus transactionStatus : savepoints.keySet()) {
|
||||
transactionStatus.releaseSavepoint(savepointFor(transactionStatus));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Object createSavepoint() throws TransactionException {
|
||||
SavePoints savePoints = new SavePoints();
|
||||
|
||||
for (TransactionStatus transactionStatus : transactionStatuses.values()) {
|
||||
savePoints.save(transactionStatus);
|
||||
}
|
||||
return savePoints;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void rollbackToSavepoint(Object savepoint) throws TransactionException {
|
||||
SavePoints savePoints= (SavePoints) savepoint;
|
||||
savePoints.rollback();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void releaseSavepoint(Object savepoint) throws TransactionException {
|
||||
((SavePoints)savepoint).release();
|
||||
}
|
||||
|
||||
public void registerTransactionManager(TransactionDefinition definition, PlatformTransactionManager transactionManager) {
|
||||
getTransactionStatuses().put(transactionManager, transactionManager.getTransaction(definition));
|
||||
}
|
||||
|
||||
void commit(PlatformTransactionManager transactionManager) {
|
||||
TransactionStatus transactionStatus = getTransactionStatus(transactionManager);
|
||||
transactionManager.commit(transactionStatus);
|
||||
}
|
||||
|
||||
private TransactionStatus getTransactionStatus(PlatformTransactionManager transactionManager) {
|
||||
return this.getTransactionStatuses().get(transactionManager);
|
||||
}
|
||||
|
||||
void rollback(PlatformTransactionManager transactionManager) {
|
||||
transactionManager.rollback(getTransactionStatus(transactionManager));
|
||||
}
|
||||
|
||||
@Override
|
||||
public void flush() {
|
||||
for (TransactionStatus transactionStatus : transactionStatuses.values()) {
|
||||
transactionStatus.flush();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
package org.springframework.data.graph.neo4j.support;
|
||||
|
||||
/**
|
||||
* @author mh
|
||||
* @since 15.02.11
|
||||
*/
|
||||
public interface SynchronizationManager {
|
||||
void initSynchronization();
|
||||
|
||||
boolean isSynchronizationActive();
|
||||
|
||||
void clearSynchronization();
|
||||
}
|
||||
@@ -0,0 +1,71 @@
|
||||
package org.springframework.data.graph.neo4j.support;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.springframework.transaction.PlatformTransactionManager;
|
||||
import org.springframework.transaction.TransactionDefinition;
|
||||
import org.springframework.transaction.TransactionException;
|
||||
import org.springframework.transaction.TransactionStatus;
|
||||
import org.springframework.transaction.support.DefaultTransactionDefinition;
|
||||
|
||||
import java.util.Arrays;
|
||||
|
||||
import static junit.framework.Assert.assertTrue;
|
||||
|
||||
/**
|
||||
* @author mh
|
||||
* @since 15.02.11
|
||||
*/
|
||||
public class ChainedTransactionManagerTest {
|
||||
|
||||
@Test
|
||||
public void shouldCompleteSuccessfully() throws Exception {
|
||||
ChainedTransactionManager tm = new ChainedTransactionManager(new NullSynchronizationManager());
|
||||
TestPlatformTransactionManager transactionManager = new TestPlatformTransactionManager();
|
||||
tm.setTransactionManagers(Arrays.<PlatformTransactionManager>asList(transactionManager));
|
||||
MultiTransactionStatus transaction = tm.getTransaction(new DefaultTransactionDefinition());
|
||||
tm.commit(transaction);
|
||||
|
||||
assertTrue("TM didn't commit", transactionManager.isCommited());
|
||||
}
|
||||
|
||||
private static class NullSynchronizationManager implements SynchronizationManager {
|
||||
@Override
|
||||
public void initSynchronization() {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isSynchronizationActive() {
|
||||
return true;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void clearSynchronization() {
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
private static class TestPlatformTransactionManager implements PlatformTransactionManager {
|
||||
|
||||
private boolean commited;
|
||||
|
||||
@Override
|
||||
public TransactionStatus getTransaction(TransactionDefinition definition) throws TransactionException {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void commit(TransactionStatus status) throws TransactionException {
|
||||
commited = true;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void rollback(TransactionStatus status) throws TransactionException {
|
||||
|
||||
}
|
||||
|
||||
public boolean isCommited() {
|
||||
return commited;
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user