http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-db/jena-dboe-transaction/src/main/java/org/apache/jena/dboe/transaction/txn/TransactionCoordinator.java ---------------------------------------------------------------------- diff --git a/jena-db/jena-dboe-transaction/src/main/java/org/apache/jena/dboe/transaction/txn/TransactionCoordinator.java b/jena-db/jena-dboe-transaction/src/main/java/org/apache/jena/dboe/transaction/txn/TransactionCoordinator.java index 92510c5..7c63bbf 100644 --- a/jena-db/jena-dboe-transaction/src/main/java/org/apache/jena/dboe/transaction/txn/TransactionCoordinator.java +++ b/jena-db/jena-dboe-transaction/src/main/java/org/apache/jena/dboe/transaction/txn/TransactionCoordinator.java @@ -21,10 +21,7 @@ package org.apache.jena.dboe.transaction.txn; import static org.apache.jena.dboe.transaction.txn.journal.JournalEntryType.UNDO; import java.nio.ByteBuffer ; -import java.util.ArrayList ; -import java.util.Iterator ; -import java.util.List ; -import java.util.Objects ; +import java.util.*; import java.util.concurrent.ConcurrentHashMap ; import java.util.concurrent.Semaphore ; import java.util.concurrent.atomic.AtomicLong ; @@ -501,8 +498,9 @@ public class TransactionCoordinator { // A read transaction can be promoted if writer does not start // This TransactionCoordinator provides Serializable, Read-lock-free - // execution. With no item locking, a read can only be promoted - // if no writer started since the reader started. + // execution. With no item locking, a read can only be promoted + // if no writer started since the reader started or if it is "read committed", + // seeing changes made since it started and comitted at the poiont of promotion. /* The version of the data - incremented when transaction commits. * This is the version with repest to the last commited transaction. @@ -542,7 +540,7 @@ public class TransactionCoordinator { // Detemine ReadWrite for the transaction start from initial TxnType. private static ReadWrite initialMode(TxnType txnType) { - return (txnType == TxnType.WRITE) ? ReadWrite.WRITE : ReadWrite.READ; + return TxnType.initial(txnType); } private ComponentGroup chooseComponents(ComponentGroup components, TxnType txnType) { @@ -620,6 +618,7 @@ public class TransactionCoordinator { releaseWriterLock(); return false ; } + promoteActiveTransaction(transaction); } return true; } @@ -659,6 +658,7 @@ public class TransactionCoordinator { releaseWriterLock(); return false ; } + promoteActiveTransaction(transaction); } return true ; } @@ -736,9 +736,8 @@ public class TransactionCoordinator { notifyAbortFinish(transaction) ; } - // Active transactions: this is (the missing) ConcurrentHashSet - private final static Object dummy = new Object() ; - private ConcurrentHashMap<Transaction, Object> activeTransactions = new ConcurrentHashMap<>() ; + // Active transactions. + private Set<Transaction> activeTransactions = ConcurrentHashMap.newKeySet(); private AtomicLong activeTransactionCount = new AtomicLong(0) ; private AtomicLong activeReadersCount = new AtomicLong(0) ; private AtomicLong activeWritersCount = new AtomicLong(0) ; @@ -753,10 +752,17 @@ public class TransactionCoordinator { case WRITE: countBeginWrite.incrementAndGet() ; activeWritersCount.incrementAndGet() ; break ; } activeTransactionCount.incrementAndGet() ; - activeTransactions.put(transaction, dummy) ; + activeTransactions.add(transaction) ; } } + + private void promoteActiveTransaction(Transaction transaction) { + // Called for a real promote as READ-> WRITE + activeReadersCount.decrementAndGet(); + activeWritersCount.incrementAndGet(); + } + private void finishActiveTransaction(Transaction transaction) { synchronized(coordinatorLock) { // Idempotent.
http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-db/jena-dboe-transaction/src/main/java/org/apache/jena/dboe/transaction/txn/TransactionalBase.java ---------------------------------------------------------------------- diff --git a/jena-db/jena-dboe-transaction/src/main/java/org/apache/jena/dboe/transaction/txn/TransactionalBase.java b/jena-db/jena-dboe-transaction/src/main/java/org/apache/jena/dboe/transaction/txn/TransactionalBase.java index d4fd8f5..4c424a8 100644 --- a/jena-db/jena-dboe-transaction/src/main/java/org/apache/jena/dboe/transaction/txn/TransactionalBase.java +++ b/jena-db/jena-dboe-transaction/src/main/java/org/apache/jena/dboe/transaction/txn/TransactionalBase.java @@ -20,6 +20,7 @@ package org.apache.jena.dboe.transaction.txn; import java.util.Objects ; +import org.apache.jena.atlas.lib.Lib; import org.apache.jena.atlas.logging.Log ; import org.apache.jena.query.ReadWrite ; import org.apache.jena.query.TxnType; @@ -156,7 +157,7 @@ public class TransactionalBase implements TransactionalSystem { @Override public ReadWrite transactionMode() { checkRunning() ; - Transaction txn = getTxn() ; + Transaction txn = Lib.readThreadLocal(theTxn) ; if ( txn != null ) return txn.getMode() ; return null ; @@ -165,7 +166,7 @@ public class TransactionalBase implements TransactionalSystem { @Override public TxnType transactionType() { checkRunning() ; - Transaction txn = getTxn() ; + Transaction txn = Lib.readThreadLocal(theTxn) ; if ( txn != null ) return txn.getTxnType() ; return null ; @@ -173,18 +174,9 @@ public class TransactionalBase implements TransactionalSystem { @Override public boolean isInTransaction() { - return getTxn() != null; + return Lib.readThreadLocal(theTxn) != null; } - private Transaction getTxn() { - // tricky - touching theTxn causes it to initialize. - Transaction txn = theTxn.get() ; - if ( txn != null ) - return txn; - theTxn.remove() ; - return null ; - } - @Override final public TransactionInfo getTransactionInfo() { @@ -194,12 +186,7 @@ public class TransactionalBase implements TransactionalSystem { @Override final public Transaction getThreadTransaction() { - Transaction txn = theTxn.get() ; - // XXX Use getTxn() ?? - // Touched the thread local so it is defined now. -// if ( txn == null ) -// theTxn.remove() ; - return txn ; + return Lib.readThreadLocal(theTxn); } /** Get the transaction, checking there is one */ http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-db/jena-dboe-transaction/src/test/java/org/apache/jena/dboe/transaction/TestTransactionLifecycle.java ---------------------------------------------------------------------- diff --git a/jena-db/jena-dboe-transaction/src/test/java/org/apache/jena/dboe/transaction/TestTransactionLifecycle.java b/jena-db/jena-dboe-transaction/src/test/java/org/apache/jena/dboe/transaction/TestTransactionLifecycle.java index ba73fa8..72230a0 100644 --- a/jena-db/jena-dboe-transaction/src/test/java/org/apache/jena/dboe/transaction/TestTransactionLifecycle.java +++ b/jena-db/jena-dboe-transaction/src/test/java/org/apache/jena/dboe/transaction/TestTransactionLifecycle.java @@ -18,58 +18,66 @@ package org.apache.jena.dboe.transaction; -import static org.junit.Assert.fail ; +import static org.junit.Assert.* ; import org.junit.Test ; import org.apache.jena.dboe.transaction.txn.TransactionException; import org.apache.jena.query.ReadWrite ; +import org.apache.jena.query.TxnType; /** * Tests of transaction lifecycle in one JVM. - * Journal independent. - * Not testing recovery or writing to the journal. + * Not testing recovery or writing to the journal. + * Tests on a "unit", not directly on the TransactionCoordinator. */ public class TestTransactionLifecycle extends AbstractTestTxn { - @Test public void txn_read_end() { + + @Test public void txn_read_end_RW() { unit.begin(ReadWrite.READ); unit.end() ; checkClear() ; } + @Test public void txn_read_end() { + unit.begin(TxnType.READ); + unit.end() ; + checkClear() ; + } + @Test public void txn_read_end_end() { - unit.begin(ReadWrite.READ); + unit.begin(TxnType.READ); unit.end() ; unit.end() ; checkClear() ; } @Test public void txn_read_abort() { - unit.begin(ReadWrite.READ); + unit.begin(TxnType.READ); unit.abort() ; checkClear() ; } @Test public void txn_read_commit() { - unit.begin(ReadWrite.READ); + unit.begin(TxnType.READ); unit.commit() ; checkClear() ; } @Test public void txn_read_abort_end() { - unit.begin(ReadWrite.READ); + unit.begin(TxnType.READ); unit.abort() ; unit.end() ; checkClear() ; } @Test public void txn_read_commit_end() { - unit.begin(ReadWrite.READ); + unit.begin(TxnType.READ); unit.commit() ; unit.end() ; checkClear() ; } @Test public void txn_read_commit_abort() { - unit.begin(ReadWrite.READ); + unit.begin(TxnType.READ); unit.commit() ; try { unit.abort() ; fail() ; } catch (TransactionException ex) { /* Expected : can continue */ } @@ -78,7 +86,7 @@ public class TestTransactionLifecycle extends AbstractTestTxn { } @Test public void txn_read_commit_commit() { - unit.begin(ReadWrite.READ); + unit.begin(TxnType.READ); unit.commit() ; try { unit.commit() ; fail() ; } catch (TransactionException ex) { /* Expected : can continue */ } @@ -87,7 +95,7 @@ public class TestTransactionLifecycle extends AbstractTestTxn { } @Test public void txn_read_abort_commit() { - unit.begin(ReadWrite.READ); + unit.begin(TxnType.READ); unit.abort() ; try { unit.commit() ; fail() ; } catch (TransactionException ex) { /* Expected : can continue */ } @@ -96,7 +104,7 @@ public class TestTransactionLifecycle extends AbstractTestTxn { } @Test public void txn_read_abort_abort() { - unit.begin(ReadWrite.READ); + unit.begin(TxnType.READ); unit.abort() ; try { unit.abort() ; fail() ; } catch (TransactionException ex) { /* Expected : can continue */ } @@ -106,71 +114,78 @@ public class TestTransactionLifecycle extends AbstractTestTxn { @Test(expected=TransactionException.class) public void txn_begin_read_begin_read() { - unit.begin(ReadWrite.READ); - unit.begin(ReadWrite.READ); + unit.begin(TxnType.READ); + unit.begin(TxnType.READ); } @Test(expected=TransactionException.class) public void txn_begin_read_begin_write() { - unit.begin(ReadWrite.READ); - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.READ); + unit.begin(TxnType.WRITE); } @Test(expected=TransactionException.class) public void txn_begin_write_begin_read() { - unit.begin(ReadWrite.WRITE); - unit.begin(ReadWrite.READ); + unit.begin(TxnType.WRITE); + unit.begin(TxnType.READ); } @Test(expected=TransactionException.class) public void txn_begin_write_begin_write() { - unit.begin(ReadWrite.WRITE); - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.WRITE); + unit.begin(TxnType.WRITE); } @Test(expected=TransactionException.class) public void txn_write_begin_end() { - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.WRITE); unit.end() ; checkClear() ; } @Test public void txn_write_abort() { - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.WRITE); unit.abort() ; checkClear() ; } @Test public void txn_write_commit() { - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.WRITE); unit.commit() ; checkClear() ; } @Test public void txn_write_abort_end() { - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.WRITE); unit.abort() ; unit.end() ; checkClear() ; } @Test public void txn_write_abort_end_end() { - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.WRITE); unit.abort() ; unit.end() ; unit.end() ; checkClear() ; } - @Test public void txn_write_commit_end() { + @Test public void txn_write_commit_end_RW() { unit.begin(ReadWrite.WRITE); unit.commit() ; unit.end() ; checkClear() ; } + @Test public void txn_write_commit_end() { + unit.begin(TxnType.WRITE); + unit.commit() ; + unit.end() ; + checkClear() ; + } + @Test public void txn_write_commit_end_end() { - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.WRITE); unit.commit() ; unit.end() ; unit.end() ; @@ -179,7 +194,7 @@ public class TestTransactionLifecycle extends AbstractTestTxn { @Test public void txn_write_commit_abort() { // commit-abort - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.WRITE); unit.commit() ; try { unit.abort() ; fail() ; } catch (TransactionException ex) { /* Expected : can continue */ } @@ -189,7 +204,7 @@ public class TestTransactionLifecycle extends AbstractTestTxn { @Test public void txn_write_commit_commit() { // commit-commit - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.WRITE); unit.commit() ; try { unit.commit() ; fail() ; } catch (TransactionException ex) { /* Expected : can continue */ } @@ -199,7 +214,7 @@ public class TestTransactionLifecycle extends AbstractTestTxn { @Test public void txn_write_abort_commit() { // abort-commit - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.WRITE); unit.abort() ; try { unit.commit() ; fail() ; } catch (TransactionException ex) { /* Expected : can continue */ } @@ -209,7 +224,7 @@ public class TestTransactionLifecycle extends AbstractTestTxn { @Test public void txn_write_abort_abort() { // abort-abort - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.WRITE); unit.abort() ; try { unit.abort() ; fail() ; } catch (TransactionException ex) { /* Expected : can continue */ } @@ -217,14 +232,110 @@ public class TestTransactionLifecycle extends AbstractTestTxn { checkClear() ; } + @Test public void txn_read_promote_commit() { + unit.begin(TxnType.READ); + try { unit.promote(); fail(); } + // Exception is correct - it is illegal to call promote in a TxnType.READ + catch (TransactionException ex) { /* Expected : can continue */ } + unit.end() ; + checkClear() ; + } + + @Test public void txn_readpromote_promote_commit() { + unit.begin(TxnType.READ_PROMOTE); + boolean b = unit.promote(); + assertTrue(b); + unit.commit(); + unit.end() ; + checkClear() ; + } + + @Test public void txn_readpromote_promote_promote_commit() { + unit.begin(TxnType.READ_PROMOTE); + boolean b = unit.promote(); + assertTrue(b); + unit.commit(); + unit.end() ; + checkClear() ; + } + + @Test public void txn_readpromote_promote_abort() { + unit.begin(TxnType.READ_PROMOTE); + boolean b = unit.promote(); + assertTrue(b); + unit.abort(); + unit.end() ; + checkClear() ; + } + + @Test public void txn_readcomittedpromote_promote_commit() { + unit.begin(TxnType.READ_COMMITTED_PROMOTE); + boolean b = unit.promote(); + assertTrue(b); + unit.commit(); + unit.end() ; + checkClear() ; + } + + @Test public void txn_readcomittedpromote_promote_abort() { + unit.begin(TxnType.READ_COMMITTED_PROMOTE); + boolean b = unit.promote(); + assertTrue(b); + unit.abort(); + unit.end() ; + checkClear() ; + } + + @Test public void txn_readpromote_promote_end() { + unit.begin(TxnType.READ_PROMOTE); + boolean b = unit.promote(); + assertTrue(b); + try { unit.end() ; } + catch (TransactionException ex) { /* Expected : check clearup */ } + checkClear() ; + } + + @Test public void txn_readcomittedpromote_promote_end() { + unit.begin(TxnType.READ_COMMITTED_PROMOTE); + boolean b = unit.promote(); + assertTrue(b); + try { unit.end() ; } + catch (TransactionException ex) { /* Expected : check clearup */ } + checkClear() ; + } + + @Test public void txn_readpromote_commit_promote() { + unit.begin(TxnType.READ_PROMOTE); + unit.commit(); // READ commit. + try { unit.promote() ; } + catch (TransactionException ex) { /* Expected : check clearup */ } + checkClear() ; + } + + @Test public void txn_readpromote_abort_promote() { + unit.begin(TxnType.READ_PROMOTE); + unit.abort(); // READ abort. + try { unit.promote() ; } + catch (TransactionException ex) { /* Expected : check clearup */ } + checkClear() ; + } + + @Test public void txn_readpromote_end_promote() { + unit.begin(TxnType.READ_PROMOTE); + unit.end(); // READ end + try { unit.promote() ; } + catch (TransactionException ex) { /* Expected : check clearup */ } + checkClear() ; + } + private void read() { - unit.begin(ReadWrite.READ); + unit.begin(TxnType.READ); unit.end() ; checkClear() ; } private void write() { - unit.begin(ReadWrite.WRITE); + unit.begin(TxnType.WRITE); unit.commit() ; unit.end() ; checkClear() ; http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-db/jena-dboe-transaction/src/test/java/org/apache/jena/dboe/transaction/TestTransactionLifecycle2.java ---------------------------------------------------------------------- diff --git a/jena-db/jena-dboe-transaction/src/test/java/org/apache/jena/dboe/transaction/TestTransactionLifecycle2.java b/jena-db/jena-dboe-transaction/src/test/java/org/apache/jena/dboe/transaction/TestTransactionLifecycle2.java index 0e6c29d..4fe50a2 100644 --- a/jena-db/jena-dboe-transaction/src/test/java/org/apache/jena/dboe/transaction/TestTransactionLifecycle2.java +++ b/jena-db/jena-dboe-transaction/src/test/java/org/apache/jena/dboe/transaction/TestTransactionLifecycle2.java @@ -36,16 +36,12 @@ import org.junit.Test ; /** * Details tests of the transaction lifecycle in one JVM - * including tests beyond the TransactionalComponentLifecycle - * - * Journal independent. + * including tests beyond the TransactionalComponentLifecycle + * Tests directly on the TransactionCoordinator. */ public class TestTransactionLifecycle2 { // org.junit.rules.ExternalResource ? protected TransactionCoordinator txnMgr ; -// protected TransInteger counter1 = new TransInteger(0) ; -// protected TransInteger counter2 = new TransInteger(0) ; -// protected TransMonitor monitor = new TransMonitor() ; @Before public void setup() { Journal jrnl = Journal.create(Location.mem()) ; http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-fuseki1/src/main/java/org/apache/jena/fuseki/servlets/HttpAction.java ---------------------------------------------------------------------- diff --git a/jena-fuseki1/src/main/java/org/apache/jena/fuseki/servlets/HttpAction.java b/jena-fuseki1/src/main/java/org/apache/jena/fuseki/servlets/HttpAction.java index a418972..081d167 100644 --- a/jena-fuseki1/src/main/java/org/apache/jena/fuseki/servlets/HttpAction.java +++ b/jena-fuseki1/src/main/java/org/apache/jena/fuseki/servlets/HttpAction.java @@ -121,7 +121,7 @@ public class HttpAction isTransactional = false ; } else { // Nothing to build on. Be safe. - transactional = new TransactionalMutex(dsg.getLock()) ; + transactional = TransactionalLock.create(dsg.getLock()) ; isTransactional = false ; } } http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionLocal.java ---------------------------------------------------------------------- diff --git a/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionLocal.java b/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionLocal.java index 8be2704..f2d64ae 100644 --- a/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionLocal.java +++ b/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionLocal.java @@ -292,21 +292,15 @@ public class RDFConnectionLocal implements RDFConnection { throw new ARQException("closed"); } - @Override - public void begin(ReadWrite readWrite) { dataset.begin(readWrite); } - - @Override - public void commit() { dataset.commit(); } - - @Override - public void abort() { dataset.abort(); } - - @Override - public boolean isInTransaction() { return dataset.isInTransaction(); } - - @Override - public void end() { dataset.end(); } - - + @Override public void begin() { dataset.begin(); } + @Override public void begin(TxnType txnType) { dataset.begin(txnType); } + @Override public void begin(ReadWrite mode) { dataset.begin(mode); } + @Override public boolean promote() { return dataset.promote(); } + @Override public void commit() { dataset.commit(); } + @Override public void abort() { dataset.abort(); } + @Override public boolean isInTransaction() { return dataset.isInTransaction(); } + @Override public void end() { dataset.end(); } + @Override public ReadWrite transactionMode() { return dataset.transactionMode(); } + @Override public TxnType transactionType() { return dataset.transactionType(); } } http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionModular.java ---------------------------------------------------------------------- diff --git a/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionModular.java b/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionModular.java index 77e473c..97bd3b5 100644 --- a/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionModular.java +++ b/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionModular.java @@ -18,10 +18,7 @@ package org.apache.jena.rdfconnection; -import org.apache.jena.query.Dataset; -import org.apache.jena.query.Query; -import org.apache.jena.query.QueryExecution; -import org.apache.jena.query.ReadWrite; +import org.apache.jena.query.*; import org.apache.jena.rdf.model.Model; import org.apache.jena.sparql.core.Transactional; import org.apache.jena.update.UpdateRequest; @@ -37,18 +34,16 @@ public class RDFConnectionModular implements RDFConnection { private final RDFDatasetConnection datasetConnection; private final Transactional transactional; - @Override - public void begin(ReadWrite readWrite) { transactional.begin(readWrite); } - @Override - public void commit() { transactional.commit(); } - @Override - public void abort() { transactional.abort(); } - @Override - public void end() { transactional.end(); } - @Override - public boolean isInTransaction() { - return transactional.isInTransaction(); - } + @Override public void begin() { transactional.begin(); } + @Override public void begin(TxnType txnType) { transactional.begin(txnType); } + @Override public void begin(ReadWrite mode) { transactional.begin(mode); } + @Override public boolean promote() { return transactional.promote(); } + @Override public void commit() { transactional.commit(); } + @Override public void abort() { transactional.abort(); } + @Override public boolean isInTransaction() { return transactional.isInTransaction(); } + @Override public void end() { transactional.end(); } + @Override public ReadWrite transactionMode() { return transactional.transactionMode(); } + @Override public TxnType transactionType() { return transactional.transactionType(); } public RDFConnectionModular(SparqlQueryConnection queryConnection , SparqlUpdateConnection updateConnection , http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionRemote.java ---------------------------------------------------------------------- diff --git a/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionRemote.java b/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionRemote.java index 8f5b359..3666bc4 100644 --- a/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionRemote.java +++ b/jena-rdfconnection/src/main/java/org/apache/jena/rdfconnection/RDFConnectionRemote.java @@ -22,7 +22,6 @@ import static java.util.Objects.requireNonNull; import java.io.File; import java.io.InputStream; -import java.util.concurrent.locks.ReentrantLock; import java.util.function.Supplier; import org.apache.http.HttpEntity; @@ -42,6 +41,7 @@ import org.apache.jena.riot.web.HttpResponseLib; import org.apache.jena.sparql.ARQException; import org.apache.jena.sparql.core.DatasetGraph; import org.apache.jena.sparql.core.Transactional; +import org.apache.jena.sparql.core.TransactionalLock; import org.apache.jena.update.UpdateExecutionFactory; import org.apache.jena.update.UpdateProcessor; import org.apache.jena.update.UpdateRequest; @@ -416,64 +416,17 @@ public class RDFConnectionRemote implements RDFConnection { throw ex; } - /** Engine for the transaction lifecycle. MR+SW */ - static class TxnLifecycle implements Transactional { - private ReentrantLock lock = new ReentrantLock(); - private ThreadLocal<ReadWrite> mode = ThreadLocal.withInitial(()->null); - @Override - public void begin(ReadWrite readWrite) { - if ( readWrite == ReadWrite.WRITE ) - lock.lock(); - mode.set(readWrite); - } - - @Override - public void commit() { - if ( mode.get() == ReadWrite.WRITE) - lock.unlock(); - mode.set(null); - } - - @Override - public void abort() { - if ( mode.get() == ReadWrite.WRITE ) - lock.unlock(); - mode.set(null); - } - - @Override - public boolean isInTransaction() { - return mode.get() != null; - } - - @Override - public void end() { - ReadWrite rw = mode.get(); - if ( rw == null ) - return; - if ( rw == ReadWrite.WRITE ) { - abort(); - return; - } - mode.set(null); - } - } - - private TxnLifecycle inner = new TxnLifecycle(); + private final Transactional txn = TransactionalLock.createMRPlusSW(); - @Override - public void begin(ReadWrite readWrite) { inner.begin(readWrite); } - - @Override - public void commit() { inner.commit(); } - - @Override - public void abort() { inner.abort(); } - - @Override - public boolean isInTransaction() { return inner.isInTransaction(); } - - @Override - public void end() { inner.end(); } + @Override public void begin() { txn.begin(); } + @Override public void begin(TxnType txnType) { txn.begin(txnType); } + @Override public void begin(ReadWrite mode) { txn.begin(mode); } + @Override public boolean promote() { return txn.promote(); } + @Override public void commit() { txn.commit(); } + @Override public void abort() { txn.abort(); } + @Override public boolean isInTransaction() { return txn.isInTransaction(); } + @Override public void end() { txn.end(); } + @Override public ReadWrite transactionMode() { return txn.transactionMode(); } + @Override public TxnType transactionType() { return txn.transactionType(); } } http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-sdb/src/main/java/org/apache/jena/sdb/store/DatasetGraphSDB.java ---------------------------------------------------------------------- diff --git a/jena-sdb/src/main/java/org/apache/jena/sdb/store/DatasetGraphSDB.java b/jena-sdb/src/main/java/org/apache/jena/sdb/store/DatasetGraphSDB.java index 0c0ccc0..7c02665 100644 --- a/jena-sdb/src/main/java/org/apache/jena/sdb/store/DatasetGraphSDB.java +++ b/jena-sdb/src/main/java/org/apache/jena/sdb/store/DatasetGraphSDB.java @@ -25,6 +25,7 @@ import org.apache.jena.graph.Graph ; import org.apache.jena.graph.Node ; import org.apache.jena.graph.Triple ; import org.apache.jena.query.ReadWrite ; +import org.apache.jena.query.TxnType; import org.apache.jena.sdb.Store ; import org.apache.jena.sdb.graph.GraphSDB ; import org.apache.jena.sdb.util.StoreUtils ; @@ -110,17 +111,19 @@ public class DatasetGraphSDB extends DatasetGraphTriplesQuads @Override public void close() { store.close() ; } - + + // Transactions for SDB are an aspect of the JDBC connection not the dataset. private final Transactional txn = new TransactionalNotSupported() ; - @Override public void begin(ReadWrite mode) { txn.begin(mode) ; } - @Override public void commit() { txn.commit() ; } - @Override public void abort() { txn.abort() ; } - @Override public boolean isInTransaction() { return txn.isInTransaction() ; } + @Override public void begin() { txn.begin(); } + @Override public void begin(TxnType txnType) { txn.begin(txnType); } + @Override public void begin(ReadWrite mode) { txn.begin(mode); } + @Override public boolean promote() { return txn.promote(); } + @Override public void commit() { txn.commit(); } + @Override public void abort() { txn.abort(); } + @Override public boolean isInTransaction() { return txn.isInTransaction(); } @Override public void end() { txn.end(); } - @Override public boolean supportsTransactions() { return false ; } - @Override public boolean supportsTransactionAbort() { return false ; } - - // Helper implementations of operations. - // Not necessarily efficient. - + @Override public ReadWrite transactionMode() { return txn.transactionMode(); } + @Override public TxnType transactionType() { return txn.transactionType(); } + @Override public boolean supportsTransactions() { return false; } + @Override public boolean supportsTransactionAbort() { return false; } } http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-tdb/src/main/java/org/apache/jena/tdb/StoreConnection.java ---------------------------------------------------------------------- diff --git a/jena-tdb/src/main/java/org/apache/jena/tdb/StoreConnection.java b/jena-tdb/src/main/java/org/apache/jena/tdb/StoreConnection.java index f24065c..f645442 100644 --- a/jena-tdb/src/main/java/org/apache/jena/tdb/StoreConnection.java +++ b/jena-tdb/src/main/java/org/apache/jena/tdb/StoreConnection.java @@ -88,10 +88,10 @@ public class StoreConnection return transactionManager.state(); } - /* + /** * @deprecated Use {@link #begin(TxnType)} */ - //@Deprecated + @Deprecated public DatasetGraphTxn begin(ReadWrite mode) { return begin(TxnType.convert(mode)); } @@ -197,8 +197,8 @@ public class StoreConnection /** Stop managing a location. Use with great care (testing only). */ public static synchronized void expel(Location location, boolean force) { StoreConnection sConn = cache.get(location) ; - if (sConn == null) return ; - + if (sConn == null) + return ; if (!force && sConn.transactionManager.activeTransactions()) throw new TDBTransactionException("Can't expel: Active transactions for location: " + location) ; http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-tdb/src/main/java/org/apache/jena/tdb/store/DatasetGraphTDB.java ---------------------------------------------------------------------- diff --git a/jena-tdb/src/main/java/org/apache/jena/tdb/store/DatasetGraphTDB.java b/jena-tdb/src/main/java/org/apache/jena/tdb/store/DatasetGraphTDB.java index 647f46e..d4390dc 100644 --- a/jena-tdb/src/main/java/org/apache/jena/tdb/store/DatasetGraphTDB.java +++ b/jena-tdb/src/main/java/org/apache/jena/tdb/store/DatasetGraphTDB.java @@ -257,12 +257,16 @@ public class DatasetGraphTDB extends DatasetGraphTriplesQuads } private final Transactional txn = new TransactionalNotSupported() ; - @Override public void begin(TxnType type) { txn.begin(type) ; } - @Override public void begin(ReadWrite mode) { txn.begin(mode) ; } - @Override public void commit() { txn.commit() ; } - @Override public void abort() { txn.abort() ; } - @Override public boolean isInTransaction() { return txn.isInTransaction() ; } + @Override public void begin() { txn.begin(); } + @Override public void begin(TxnType txnType) { txn.begin(txnType); } + @Override public void begin(ReadWrite mode) { txn.begin(mode); } + @Override public boolean promote() { return txn.promote(); } + @Override public void commit() { txn.commit(); } + @Override public void abort() { txn.abort(); } + @Override public boolean isInTransaction() { return txn.isInTransaction(); } @Override public void end() { txn.end(); } - @Override public boolean supportsTransactions() { return true ; } - @Override public boolean supportsTransactionAbort() { return false ; } + @Override public ReadWrite transactionMode() { return txn.transactionMode(); } + @Override public TxnType transactionType() { return txn.transactionType(); } + @Override public boolean supportsTransactions() { return true; } + @Override public boolean supportsTransactionAbort() { return false; } } http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/BlockMgrJournal.java ---------------------------------------------------------------------- diff --git a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/BlockMgrJournal.java b/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/BlockMgrJournal.java index 025d716..5b02e30 100644 --- a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/BlockMgrJournal.java +++ b/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/BlockMgrJournal.java @@ -75,7 +75,7 @@ public class BlockMgrJournal implements BlockMgr, TransactionLifecycle writeBlockBufferAllocator = new BufferAllocatorMem() ; reset(txn, fileRef, underlyingBlockMgr) ; - if ( txn.getMode() == ReadWrite.READ && underlyingBlockMgr instanceof BlockMgrJournal ) + if ( txn.getTxnMode() == ReadWrite.READ && underlyingBlockMgr instanceof BlockMgrJournal ) System.err.println("Two level BlockMgrJournal") ; } http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/DatasetGraphTransaction.java ---------------------------------------------------------------------- diff --git a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/DatasetGraphTransaction.java b/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/DatasetGraphTransaction.java index e37b323..5b77607 100644 --- a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/DatasetGraphTransaction.java +++ b/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/DatasetGraphTransaction.java @@ -23,6 +23,7 @@ import static java.lang.ThreadLocal.withInitial ; import org.apache.jena.atlas.lib.Sync ; import org.apache.jena.graph.Graph ; import org.apache.jena.graph.Node ; +import org.apache.jena.query.ReadWrite; import org.apache.jena.query.TxnType; import org.apache.jena.sparql.JenaTransactionException ; import org.apache.jena.sparql.core.DatasetGraph ; @@ -61,7 +62,7 @@ import org.apache.jena.tdb.store.GraphTxnTDB ; */ // Transaction per thread per DatasetGraphTransaction object. - private ThreadLocal<DatasetGraphTxn> txn = withInitial(() -> null); + private ThreadLocal<DatasetGraphTxn> dsgtxn = withInitial(() -> null); private ThreadLocal<Boolean> inTransaction = withInitial(() -> false); private final StoreConnection sConn; @@ -89,7 +90,7 @@ import org.apache.jena.tdb.store.GraphTxnTDB ; @Override public DatasetGraph getW() { if ( isInTransaction() ) { - DatasetGraphTxn dsgTxn = txn.get() ; + DatasetGraphTxn dsgTxn = dsgtxn.get() ; if ( dsgTxn.getTransaction().isRead() ) { TxnType txnType = dsgTxn.getTransaction().getTxnType(); switch(txnType) { @@ -104,7 +105,9 @@ import org.apache.jena.tdb.store.GraphTxnTDB ; // Promotion. TransactionManager txnMgr = dsgTxn.getTransaction().getTxnMgr() ; DatasetGraphTxn dsgTxn2 = txnMgr.promote(dsgTxn, txnType) ; - txn.set(dsgTxn2); + if ( dsgTxn2 == null ) + throw new JenaTransactionException("Can't promote "+txnType+"- dataset has been written to"); + dsgtxn.set(dsgTxn2); } } return super.getW() ; @@ -114,7 +117,7 @@ import org.apache.jena.tdb.store.GraphTxnTDB ; @Override public DatasetGraphTDB get() { if ( isInTransaction() ) { - DatasetGraphTxn dsgTxn = txn.get() ; + DatasetGraphTxn dsgTxn = dsgtxn.get() ; if ( dsgTxn == null ) throw new TDBTransactionException("In a transaction but no transactional DatasetGraph") ; return dsgTxn.getView() ; @@ -152,6 +155,22 @@ import org.apache.jena.tdb.store.GraphTxnTDB ; checkNotClosed() ; return inTransaction.get() ; } + + @Override + public ReadWrite transactionMode() { + checkNotClosed() ; + if ( ! isInTransaction() ) + return null; + return dsgtxn.get().getTransaction().getTxnMode(); + } + + @Override + public TxnType transactionType() { + checkNotClosed() ; + if ( ! isInTransaction() ) + return null; + return dsgtxn.get().getTransaction().getTxnType(); + } public boolean isClosed() { return isClosed ; @@ -182,34 +201,41 @@ import org.apache.jena.tdb.store.GraphTxnTDB ; protected void _begin(TxnType txnType) { checkNotClosed() ; DatasetGraphTxn dsgTxn = sConn.begin(txnType) ; - txn.set(dsgTxn) ; + dsgtxn.set(dsgTxn) ; inTransaction.set(true) ; } @Override protected boolean _promote() { + // Promotion (TDB1) is a reset of the DatasetGraphTxn. checkNotClosed() ; - return txn.get().promote(); + DatasetGraphTxn dsgTxn = dsgtxn.get(); + Transaction transaction = dsgTxn.getTransaction(); + DatasetGraphTxn dsgTxn2 = transaction.getTxnMgr().promote(dsgTxn, transaction.getTxnType()); + if ( dsgTxn2 == null ) + return false; + dsgtxn.set(dsgTxn2) ; + return true; } - + @Override protected void _commit() { checkNotClosed() ; - txn.get().commit() ; + dsgtxn.get().commit() ; inTransaction.set(false) ; } @Override protected void _abort() { checkNotClosed() ; - txn.get().abort() ; + dsgtxn.get().abort() ; inTransaction.set(false) ; } @Override protected void _end() { checkNotClosed() ; - DatasetGraphTxn dsg = txn.get() ; + DatasetGraphTxn dsg = dsgtxn.get() ; // It's null if end() already called. if ( dsg == null ) { TDB.logInfo.warn("Transaction already ended") ; @@ -217,11 +243,11 @@ import org.apache.jena.tdb.store.GraphTxnTDB ; } try { // begin(W)..end() throws an exception. - txn.get().end() ; + dsgtxn.get().end() ; } finally { // May already be false due to .commit/.abort. inTransaction.set(false) ; - txn.set(null) ; + dsgtxn.set(null) ; } } @@ -272,7 +298,7 @@ import org.apache.jena.tdb.store.GraphTxnTDB ; TDB.logInfo.warn("Attempt to close a DatasetGraphTransaction while a transaction is active - ignored close (" + getLocation() + ")") ; return ; } - txn.remove() ; + dsgtxn.remove() ; inTransaction.remove() ; isClosed = true ; } http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/DatasetGraphTxn.java ---------------------------------------------------------------------- diff --git a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/DatasetGraphTxn.java b/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/DatasetGraphTxn.java index 98d9ec2..cfa36e7 100644 --- a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/DatasetGraphTxn.java +++ b/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/DatasetGraphTxn.java @@ -18,6 +18,7 @@ package org.apache.jena.tdb.transaction; +import org.apache.jena.atlas.lib.NotImplemented; import org.apache.jena.query.ReadWrite ; import org.apache.jena.sparql.core.DatasetGraphWrapper; import org.apache.jena.tdb.store.DatasetGraphTDB; @@ -50,6 +51,12 @@ public class DatasetGraphTxn extends DatasetGraphWrapper { } @Override + public boolean promote() { + //transaction.getTxnMgr().promote(this, ??) + throw new NotImplemented("DatasetGraphTxn.promote"); + } + + @Override public void commit() { transaction.commit(); } http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/Transaction.java ---------------------------------------------------------------------- diff --git a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/Transaction.java b/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/Transaction.java index 9323435..9e79a00 100644 --- a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/Transaction.java +++ b/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/Transaction.java @@ -38,6 +38,7 @@ public class Transaction private final TransactionManager txnMgr ; private final Journal journal ; private final TxnType txnType ; + private final TxnType originalTxnType ; private final ReadWrite mode ; private final List<ObjectFileTrans> objectFileTrans = new ArrayList<>() ; @@ -56,7 +57,7 @@ public class Transaction private boolean changesPending ; - public Transaction(DatasetGraphTDB dsg, long version, TxnType txnType, ReadWrite mode, long id, String label, TransactionManager txnMgr) { + public Transaction(DatasetGraphTDB dsg, long version, TxnType txnType, ReadWrite mode, long id, TxnType originalTxnType, String label, TransactionManager txnMgr) { this.id = id ; if (label == null ) label = "Txn" ; @@ -71,6 +72,7 @@ public class Transaction this.basedsg = dsg ; this.version = version ; this.txnType = txnType ; + this.originalTxnType = originalTxnType ; this.mode = mode ; this.journal = ( txnMgr == null ) ? null : txnMgr.getJournal() ; activedsg = null ; // Don't know yet. @@ -273,8 +275,9 @@ public class Transaction } } - public TxnType getTxnType() { return txnType ; } - public ReadWrite getMode() { return mode ; } + public TxnType getTxnType() { return originalTxnType ; } + public TxnType getCurrentTxnType() { return txnType ; } + public ReadWrite getTxnMode() { return mode ; } public boolean isRead() { return mode == ReadWrite.READ ; } public boolean isWrite() { return mode == ReadWrite.WRITE ; } http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/TransactionManager.java ---------------------------------------------------------------------- diff --git a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/TransactionManager.java b/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/TransactionManager.java index ad01342..10e8b44 100644 --- a/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/TransactionManager.java +++ b/jena-tdb/src/main/java/org/apache/jena/tdb/transaction/TransactionManager.java @@ -26,6 +26,7 @@ import static org.apache.jena.tdb.transaction.TransactionManager.TxnPoint.CLOSE import java.io.File ; import java.util.*; import java.util.concurrent.BlockingQueue ; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.LinkedBlockingDeque ; import java.util.concurrent.Semaphore ; import java.util.concurrent.atomic.AtomicLong ; @@ -48,7 +49,7 @@ public class TransactionManager private static boolean checking = true ; private static Logger log = LoggerFactory.getLogger(TransactionManager.class) ; - private Set<Transaction> activeTransactions = new HashSet<>() ; + private Set<Transaction> activeTransactions = ConcurrentHashMap.newKeySet(); synchronized public boolean activeTransactions() { return !activeTransactions.isEmpty() ; } // Setting this true cause the TransactionManager to keep lists of transactions @@ -162,6 +163,7 @@ public class TransactionManager void transactionCloses(Transaction txn) ; void readerStarts(Transaction txn) ; void readerFinishes(Transaction txn) ; + void transactionPromotes(Transaction txn) ; void writerStarts(Transaction txn) ; void writerCommits(Transaction txn) ; void writerAborts(Transaction txn) ; @@ -173,6 +175,7 @@ public class TransactionManager @Override public void transactionCloses(Transaction txn) {} @Override public void readerStarts(Transaction txn) {} @Override public void readerFinishes(Transaction txn) {} + @Override public void transactionPromotes(Transaction txn) {} @Override public void writerStarts(Transaction txn) {} @Override public void writerCommits(Transaction txn) {} @Override public void writerAborts(Transaction txn) {} @@ -182,6 +185,7 @@ public class TransactionManager TSM_Logger() {} @Override public void readerStarts(Transaction txn) { log("start", txn) ; } @Override public void readerFinishes(Transaction txn) { log("finish", txn) ; } + @Override public void transactionPromotes(Transaction txn) { log("promote", txn) ; } @Override public void writerStarts(Transaction txn) { log("begin", txn) ; } @Override public void writerCommits(Transaction txn) { log("commit", txn) ; } @Override public void writerAborts(Transaction txn) { log("abort", txn) ; } @@ -192,12 +196,12 @@ public class TransactionManager TSM_LoggerDebug() {} @Override public void readerStarts(Transaction txn) { logInternal("start", txn) ; } @Override public void readerFinishes(Transaction txn) { logInternal("finish", txn) ; } + @Override public void transactionPromotes(Transaction txn) { logInternal("promote", txn) ; } @Override public void writerStarts(Transaction txn) { logInternal("begin", txn) ; } @Override public void writerCommits(Transaction txn) { logInternal("commit", txn) ; } @Override public void writerAborts(Transaction txn) { logInternal("abort", txn) ; } } - // Mixes stats and state variables :-( class TSM_Counters implements TSM { TSM_Counters() {} @Override public void transactionStarts(Transaction txn) { activeTransactions.add(txn) ; } @@ -205,6 +209,7 @@ public class TransactionManager @Override public void transactionCloses(Transaction txn) { } @Override public void readerStarts(Transaction txn) { inc(activeReaders) ; } @Override public void readerFinishes(Transaction txn) { dec(activeReaders) ; inc(finishedReaders); } + @Override public void transactionPromotes(Transaction txn) { dec(activeReaders) ; inc(finishedReaders); inc(activeWriters); } @Override public void writerStarts(Transaction txn) { inc(activeWriters) ; } @Override public void writerCommits(Transaction txn) { dec(activeWriters) ; inc(committedWriters) ; } @Override public void writerAborts(Transaction txn) { dec(activeWriters) ; inc(abortedWriters) ; } @@ -318,7 +323,11 @@ public class TransactionManager // _begin(readWrite) ; // } - public DatasetGraphTxn begin(TxnType mode, String label) { + public DatasetGraphTxn begin(TxnType txnType, String label) { + return beginInternal(txnType, txnType, label); + } + + private DatasetGraphTxn beginInternal(TxnType txnType, TxnType originalTxnType, String label) { // The exclusivitylock surrounds the entire transaction cycle. // Paired with notifyCommit, notifyAbort. startNonExclusive(); @@ -326,14 +335,14 @@ public class TransactionManager // Not synchronized (else blocking on semaphore will never wake up // because Semaphore.release is inside synchronized). // Allow only one active writer. - if ( mode == TxnType.WRITE ) { + if ( txnType == TxnType.WRITE ) { // Writers take a WRITE permit from the semaphore to ensure there // is at most one active writer, else the attempt to start the // transaction blocks. acquireWriterLock(true) ; } // entry synchronized part - return begin$(mode, label) ; + return begin$(txnType, originalTxnType, label) ; } /** Ensure a DatasetGraphTxn is for a write transaction. @@ -349,23 +358,29 @@ public class TransactionManager * There is no point retrying - later committed changes have been made and will remain. * <p> * "read committed" will always succeed but the app needs to be aware that data access before the promotion - * is no longer valid. It may need to check it. + * is no longer valid. It may need to check it. + * <p> + * Return null for "no promote" due to intermediate commits. */ /*package*/ DatasetGraphTxn promote(DatasetGraphTxn dsgtxn, TxnType txnType) throws TDBTransactionException { Transaction txn = dsgtxn.getTransaction() ; if ( txn.getState() != TxnState.ACTIVE ) throw new TDBTransactionException("promote: transaction is not active") ; - if ( txn.getMode() == ReadWrite.WRITE ) + if ( txn.getTxnMode() == ReadWrite.WRITE ) return dsgtxn ; - if ( txn.getTxnType() == TxnType.READ ) - throw new TDBTransactionException("promote: transaction is a read transaction") ; + if ( txn.getTxnType() == TxnType.READ ) { + txn.abort(); + throw new TDBTransactionException("promote: transaction is a READ transaction") ; + } // Read commit - pick up whatever is current at the point setup. // Can also promote - may need to wait for active writers. // Go through begin for the writers lock. if ( txnType == TxnType.READ_COMMITTED_PROMOTE ) { - DatasetGraphTxn dsgtxn2 = begin(TxnType.WRITE, txn.getLabel()) ; + // Full begin cycle. + DatasetGraphTxn dsgtxn2 = beginInternal(TxnType.WRITE, txnType, txn.getLabel()) ; // Junk the old one. + noteTxnPromote(txn, dsgtxn2.getTransaction()); return dsgtxn2 ; } @@ -377,9 +392,8 @@ public class TransactionManager // acquireWriterLock returning. It catches many cases without needing // to acquire the writer lock. - if ( txn.getVersion() != version.get() ) { - throw new TDBTransactionException("Dataset changed - can't promote") ; - } + if ( txn.getVersion() != version.get() ) + return null; // Put ourselves in the serialization timeline of the dataset - that, is grab a writer // lock as a step toward promotion. We can then test properly because no other writer @@ -390,19 +404,20 @@ public class TransactionManager acquireWriterLock(true) ; // Do the synchronized stuff. - return promote2$(dsgtxn) ; + return promote2$(dsgtxn, txnType) ; } synchronized - private DatasetGraphTxn promote2$(DatasetGraphTxn dsgtxn) { + private DatasetGraphTxn promote2$(DatasetGraphTxn dsgtxn, TxnType originalTxnType) { Transaction txn = dsgtxn.getTransaction() ; // Writers may have happened between the first check of the active writers may have committed. if ( txn.getVersion() != version.get() ) { releaseWriterLock(); - throw new TDBTransactionException("Active writer changed the dataset - can't promote") ; + return null; } - // Use begin$ - we have the writers lock. - DatasetGraphTxn dsgtxn2 = begin$(TxnType.WRITE, txn.getLabel()) ; + // Use begin$ (not beginInternal) - we have the writers lock. + DatasetGraphTxn dsgtxn2 = begin$(TxnType.WRITE, originalTxnType, txn.getLabel()) ; + noteTxnPromote(txn, dsgtxn2.getTransaction()); return dsgtxn2 ; } @@ -411,7 +426,7 @@ public class TransactionManager // of the low level objects directly so we'll play safe. synchronized - private DatasetGraphTxn begin$(TxnType txnType, String label) { + private DatasetGraphTxn begin$(TxnType txnType, TxnType originalTxnType, String label) { Objects.requireNonNull(txnType); if ( txnType == TxnType.WRITE && activeWriters.get() > 0 ) // Guard throw new TDBTransactionException("Existing active write transaction") ; @@ -426,7 +441,7 @@ public class TransactionManager } DatasetGraphTDB dsg = determineBaseDataset() ; - Transaction txn = createTransaction(dsg, txnType, label) ; + Transaction txn = createTransaction(dsg, txnType, originalTxnType, label) ; log("begin$", txn) ; @@ -446,7 +461,7 @@ public class TransactionManager for ( TransactionLifecycle component : components ) component.begin(dsgTxn.getTransaction()) ; - noteStartTxn(txn) ; + noteTxnStart(txn) ; return dsgTxn ; } @@ -463,14 +478,16 @@ public class TransactionManager dsg = commitedAwaitingFlush.get(commitedAwaitingFlush.size() - 1).getActiveDataset().getView() ; return dsg ; } - private Transaction createTransaction(DatasetGraphTDB dsg, TxnType txnType, String label) { - Transaction txn = new Transaction(dsg, version.get(), txnType, initialMode(txnType), transactionId.getAndIncrement(), label, this) ; + private Transaction createTransaction(DatasetGraphTDB dsg, TxnType txnType, TxnType originalTxnType, String label) { + if ( originalTxnType == null ) + originalTxnType = txnType; + Transaction txn = new Transaction(dsg, version.get(), txnType, initialMode(txnType), transactionId.getAndIncrement(), originalTxnType, label, this) ; return txn ; } // State. private static ReadWrite initialMode(TxnType txnType) { - return (txnType == TxnType.WRITE) ? ReadWrite.WRITE : ReadWrite.READ; + return TxnType.initial(txnType); } private DatasetGraphTxn createDSGTxn(DatasetGraphTDB dsg, Transaction txn, ReadWrite mode) { @@ -510,7 +527,7 @@ public class TransactionManager noteTxnCommit(transaction) ; - switch ( transaction.getMode() ) { + switch ( transaction.getTxnMode() ) { case READ: break ; case WRITE: version.incrementAndGet() ; @@ -548,7 +565,7 @@ public class TransactionManager noteTxnAbort(transaction) ; - switch ( transaction.getMode() ) + switch ( transaction.getTxnMode() ) { case READ: break ; case WRITE: releaseWriterLock(); @@ -801,7 +818,7 @@ public class TransactionManager try { Transaction txn2 = queue.take() ; - if ( txn2.getMode() == ReadWrite.READ ) + if ( txn2.getTxnMode() == ReadWrite.READ ) continue ; if ( log() ) log(" Flush delayed commit of "+txn2.getLabel(), txn) ; @@ -841,8 +858,8 @@ public class TransactionManager log.error("There are now active transactions") ; } - private void noteStartTxn(Transaction transaction) { - switch (transaction.getMode()) + private void noteTxnStart(Transaction transaction) { + switch (transaction.getTxnMode()) { case READ : readerStarts(transaction) ; break ; case WRITE : writerStarts(transaction) ; break ; @@ -850,8 +867,15 @@ public class TransactionManager transactionStarts(transaction) ; } + private void noteTxnPromote(Transaction transaction, Transaction transaction2) { + activeTransactions.remove(transaction); + activeTransactions.add(transaction2); + transactionPromotes(transaction) ; + } + + private void noteTxnCommit(Transaction transaction) { - switch (transaction.getMode()) + switch (transaction.getTxnMode()) { case READ : readerFinishes(transaction) ; break ; case WRITE : writerCommits(transaction) ; break ; @@ -860,7 +884,7 @@ public class TransactionManager } private void noteTxnAbort(Transaction transaction) { - switch (transaction.getMode()) + switch (transaction.getTxnMode()) { case READ : readerFinishes(transaction) ; break ; case WRITE : writerAborts(transaction) ; break ; @@ -986,6 +1010,12 @@ public class TransactionManager if ( tsm != null ) tsm.readerFinishes(txn) ; } + + private void transactionPromotes(Transaction txn) { + for ( TSM tsm : actions ) + if ( tsm != null ) + tsm.transactionPromotes(txn); + } private void writerStarts(Transaction txn) { for ( TSM tsm : actions ) http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-tdb/src/test/java/org/apache/jena/tdb/extra/T_TDBWriteTransaction.java ---------------------------------------------------------------------- diff --git a/jena-tdb/src/test/java/org/apache/jena/tdb/extra/T_TDBWriteTransaction.java b/jena-tdb/src/test/java/org/apache/jena/tdb/extra/T_TDBWriteTransaction.java deleted file mode 100644 index c01e190..0000000 --- a/jena-tdb/src/test/java/org/apache/jena/tdb/extra/T_TDBWriteTransaction.java +++ /dev/null @@ -1,151 +0,0 @@ -/* - * 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 - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -/* ****************************************************************************** - * Licensed Materials - Property of IBM - * (c) Copyright IBM Corporation 2011. All Rights Reserved. - * - * Note to U.S. Government Users Restricted Rights: Use, - * duplication or disclosure restricted by GSA ADP Schedule - * Contract with IBM Corp. - * *****************************************************************************/ - -// Submitted JENA-256 - -package org.apache.jena.tdb.extra ; - -import java.util.ArrayList ; -import java.util.List ; - -import org.apache.jena.atlas.lib.FileOps ; -import org.apache.jena.atlas.logging.LogCtl ; -import org.apache.jena.query.Dataset ; -import org.apache.jena.query.ReadWrite ; -import org.apache.jena.rdf.model.Model ; -import org.apache.jena.rdf.model.Property ; -import org.apache.jena.rdf.model.Resource ; -import org.apache.jena.tdb.TDBFactory ; -import org.apache.jena.tdb.base.block.FileMode ; -import org.apache.jena.tdb.base.file.Location ; -import org.apache.jena.tdb.sys.SystemTDB ; -import org.apache.jena.tdb.transaction.Journal ; -import org.apache.jena.tdb.transaction.JournalControl ; -import org.apache.jena.tdb.transaction.TransactionManager ; - -public class T_TDBWriteTransaction { - - private final static int TOTAL = 100 ; - static boolean bracketWithReader = true ; - - final static String INDEX_INFO_SUBJECT = "http://test.net/xmlns/test/1.0/Triple-Indexer"; - final static String TIMESTAMP_PREDICATE = "http://test.net/xmlns/test/1.0/lastProcessedTimestamp"; - final static String URI_PREDICATE = "http://test.net/xmlns/test/1.0/lastProcessedUri"; - final static String VERSION_PREDICATE = "http://test.net/xmlns/test/1.0/indexVersion"; - final static String INDEX_SIZE_PREDICATE = "http://test.net/xmlns/test/1.0/indexSize"; - - public static void main(String[] args) { - if ( true ) - SystemTDB.setFileMode(FileMode.direct) ; - - LogCtl.setLog4j() ; - TransactionManager.QueueBatchSize = 10; - -// if (args.length == 0) { -// System.out.println("Provide index location"); -// return; -// } - - String location = "DBX" ; - FileOps.ensureDir(location) ; - //FileOps.clearDirectory(location) ; - bracketWithReader = false ; - - run(location) ; -// StoreConnection.make(location).forceRecoverFromJournal() ; -// run(location) ; - } - - static public void run(String location) - { - if ( false ) - { - Journal journal = Journal.create(Location.create(location)) ; - JournalControl.print(journal) ; - journal.close() ; - } - //String location = args[0]; // + "/" + UUID.randomUUID().toString(); - - //String baseGraphName = "com.ibm.test.graphNamePrefix."; - - long totalExecTime = 0L; - long size = 0; - Dataset dataset = TDBFactory.createDataset(location); - - Dataset dataset1 = TDBFactory.createDataset(location); - - if ( bracketWithReader ) - dataset1.begin(ReadWrite.READ) ; - - for (int i = 0; i < TOTAL; i++) { - List<String> lastProcessedUris = new ArrayList<>(); - for (int j = 0; j < 10*i; j++) { - String lastProcessedUri = "http://test.net/xmlns/test/1.0/someUri" + j; - lastProcessedUris.add(lastProcessedUri); - } - //Dataset dataset = TDBFactory.createDataset(location); - //String graphName = baseGraphName + i; - long t = System.currentTimeMillis(); - - try { - dataset.begin(ReadWrite.WRITE); - Model m = dataset.getDefaultModel(); - - m.removeAll(); - Resource subject = m.createResource(INDEX_INFO_SUBJECT); - Property predicate = m.createProperty(TIMESTAMP_PREDICATE); - m.addLiteral(subject, predicate, System.currentTimeMillis()); - predicate = m.createProperty(URI_PREDICATE); - for (String uri : lastProcessedUris) { - m.add(subject, predicate, m.createResource(uri)); - } - predicate = m.createProperty(VERSION_PREDICATE); - m.addLiteral(subject, predicate, 1.0); - - size += m.size() + 1; - - predicate = m.createProperty(INDEX_SIZE_PREDICATE); - m.addLiteral(subject, predicate, size); - - dataset.commit(); - } catch (Throwable e) { - dataset.abort(); - throw new RuntimeException(e); - } finally { - dataset.end(); - long writeOperationDuration = System.currentTimeMillis() - t; - totalExecTime += writeOperationDuration; - System.out.println("Write operation " + i + " took " + writeOperationDuration + "ms"); - } - } - if ( bracketWithReader ) - dataset1.end() ; - - System.out.println("All " + TOTAL + " write operations wrote " + size + " triples and took " + totalExecTime + "ms"); - } - -} http://git-wip-us.apache.org/repos/asf/jena/blob/edab900a/jena-tdb/src/test/java/org/apache/jena/tdb/extra/T_TimeoutTDBPattern.java ---------------------------------------------------------------------- diff --git a/jena-tdb/src/test/java/org/apache/jena/tdb/extra/T_TimeoutTDBPattern.java b/jena-tdb/src/test/java/org/apache/jena/tdb/extra/T_TimeoutTDBPattern.java deleted file mode 100644 index 107bf2c..0000000 --- a/jena-tdb/src/test/java/org/apache/jena/tdb/extra/T_TimeoutTDBPattern.java +++ /dev/null @@ -1,112 +0,0 @@ -/** - * 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 - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.jena.tdb.extra; - -import java.text.MessageFormat ; -import java.util.Date ; -import java.util.concurrent.TimeUnit ; - -import org.apache.jena.query.* ; -import org.apache.jena.rdf.model.* ; -import org.apache.jena.tdb.TDBFactory ; - -// From Jena-289. - -public class T_TimeoutTDBPattern -{ - private static final int timeout1_sec = 3; - private static final int timeout2_sec = 5; - - private static final int RESOURCES = 100000; - private static final int COMMIT_EVERY = 1000; - private static final int TRIPLES_PER_RESOURCE = 100; - private static final String RES_NS = "http://example.com/"; - private static final String PROP_NS = "http://example.org/ns/1.0/"; - - public static void main(String[] args) { - String location = "DB_Jena289" ; - Dataset ds = TDBFactory.createDataset(location); - - if (ds.asDatasetGraph().isEmpty()) - create(ds) ; - - // 10M triples. - // No match to { ?a ?b ?c . ?c ?d ?e } - - final String sparql = "SELECT * WHERE { ?a ?b ?c . ?c ?d ?e }"; - - Query query = QueryFactory.create(sparql); - - ds.begin(ReadWrite.READ); - System.out.println(MessageFormat.format("{0,date} {0,time} Executing query [timeout1={1}s timeout2={2}s]: {3}", - new Date(System.currentTimeMillis()), timeout1_sec, timeout2_sec, sparql)); - try(QueryExecution qexec = QueryExecutionFactory.create(query, ds)) { - if ( true ) - qexec.setTimeout(timeout1_sec, TimeUnit.SECONDS, timeout2_sec, TimeUnit.SECONDS); - long start = System.nanoTime() ; - long finish = start ; - ResultSet rs = qexec.execSelect(); - - try { - long x = ResultSetFormatter.consume(rs) ; - finish = System.nanoTime() ; - System.out.println("Results: "+x) ; - } catch (QueryCancelledException ex) - { - finish = System.nanoTime() ; - System.out.println("Cancelled") ; - } - System.out.printf("%.2fs\n",(finish-start)/(1000.0*1000.0*1000.0)) ; - } catch (Throwable t) { - t.printStackTrace(); // OOME - } finally { - ds.end(); - ds.close(); - System.out.println(MessageFormat.format("{0,date} {0,time} Finished", - new Date(System.currentTimeMillis()))); - } - } - - private static void create(Dataset ds) - { - for (int iR = 0; iR < RESOURCES; iR++) { // 100,000 - if (iR % COMMIT_EVERY == 0) { - if (ds.isInTransaction()) { - ds.commit(); - ds.end(); - } - ds.begin(ReadWrite.WRITE); - } - - Model model = ModelFactory.createDefaultModel(); - Resource res = model.createResource(RES_NS + "resource" + iR); - for (int iP = 0; iP < TRIPLES_PER_RESOURCE; iP++) { // 100 - Property prop = ResourceFactory.createProperty(PROP_NS, "property" + iP); - model.add(res, prop, model.createTypedLiteral("Property value " + iP)); - } - //ds.addNamedModel(res.getURI(), model); - ds.getDefaultModel().add(model); - System.out.println("Created " + res.getURI()); - } - ds.commit(); - ds.end(); - } -} - -
