From 7af3f93cc45f0271f5aa20ba55d6ee60f1979e39 Mon Sep 17 00:00:00 2001 From: Michael Hunger Date: Tue, 15 Feb 2011 11:21:50 +0100 Subject: [PATCH] Chained Transaction Manager --- .../support/ChainedTransactionManager.java | 142 ++++++++++++++++ .../DefaultSynchronizationManager.java | 24 +++ .../neo4j/support/MultiTransactionStatus.java | 151 ++++++++++++++++++ .../neo4j/support/SynchronizationManager.java | 13 ++ .../ChainedTransactionManagerTest.java | 71 ++++++++ 5 files changed, 401 insertions(+) create mode 100644 spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/ChainedTransactionManager.java create mode 100644 spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/DefaultSynchronizationManager.java create mode 100644 spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/MultiTransactionStatus.java create mode 100644 spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/SynchronizationManager.java create mode 100644 spring-data-neo4j/src/test/java/org/springframework/data/graph/neo4j/support/ChainedTransactionManagerTest.java diff --git a/spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/ChainedTransactionManager.java b/spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/ChainedTransactionManager.java new file mode 100644 index 000000000..10f30c5a6 --- /dev/null +++ b/spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/ChainedTransactionManager.java @@ -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 transactionManagers = new ArrayList(); + private SynchronizationManager sychronizationManager ; + + ChainedTransactionManager(SynchronizationManager sychronizationManager) { + this.sychronizationManager = sychronizationManager; + } + + public ChainedTransactionManager() { + this(new DefaultSynchronizationManager()); + } + + public void setTransactionManagers(List 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 Iterable reverse(Collection collection) { + List list = new ArrayList(collection); + Collections.reverse(list); + return list; + } + + + private PlatformTransactionManager getLastTransactionManager() { + return transactionManagers.get(lastTransactionManagerIndex()); + } + + private int lastTransactionManagerIndex() { + return transactionManagers.size() - 1; + } + +} \ No newline at end of file diff --git a/spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/DefaultSynchronizationManager.java b/spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/DefaultSynchronizationManager.java new file mode 100644 index 000000000..4a7eddd0e --- /dev/null +++ b/spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/DefaultSynchronizationManager.java @@ -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(); + } +} diff --git a/spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/MultiTransactionStatus.java b/spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/MultiTransactionStatus.java new file mode 100644 index 000000000..4ecac2018 --- /dev/null +++ b/spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/MultiTransactionStatus.java @@ -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 transactionStatuses = + Collections.synchronizedMap(new HashMap()); + + private boolean newSynchonization; + + public MultiTransactionStatus(PlatformTransactionManager mainTransactionManager) { + this.mainTransactionManager = mainTransactionManager; + } + + + private Map 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 savepoints=new HashMap(); + + 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(); + } + } +} \ No newline at end of file diff --git a/spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/SynchronizationManager.java b/spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/SynchronizationManager.java new file mode 100644 index 000000000..1bbc2f9d9 --- /dev/null +++ b/spring-data-neo4j/src/main/java/org/springframework/data/graph/neo4j/support/SynchronizationManager.java @@ -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(); +} diff --git a/spring-data-neo4j/src/test/java/org/springframework/data/graph/neo4j/support/ChainedTransactionManagerTest.java b/spring-data-neo4j/src/test/java/org/springframework/data/graph/neo4j/support/ChainedTransactionManagerTest.java new file mode 100644 index 000000000..98eafe0ac --- /dev/null +++ b/spring-data-neo4j/src/test/java/org/springframework/data/graph/neo4j/support/ChainedTransactionManagerTest.java @@ -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.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; + } + } +}