[jboss-cvs] JBoss Messaging SVN: r4921 - in trunk: src/main/org/jboss/messaging/core/transaction/impl and 3 other directories.
jboss-cvs-commits at lists.jboss.org
jboss-cvs-commits at lists.jboss.org
Mon Sep 8 18:21:13 EDT 2008
Author: clebert.suconic at jboss.com
Date: 2008-09-08 18:21:13 -0400 (Mon, 08 Sep 2008)
New Revision: 4921
Modified:
trunk/src/main/org/jboss/messaging/core/journal/impl/JournalImpl.java
trunk/src/main/org/jboss/messaging/core/transaction/impl/TransactionImpl.java
trunk/tests/src/org/jboss/messaging/tests/integration/xa/BasicXaRecoveryTest.java
trunk/tests/src/org/jboss/messaging/tests/unit/core/journal/impl/AlignedJournalImplTest.java
trunk/tests/src/org/jboss/messaging/tests/unit/core/transaction/impl/TransactionImplTest.java
Log:
https://jira.jboss.org/jira/browse/JBMESSAGING-1418 - Persisting transaction even when using NonPersistentMessages
Modified: trunk/src/main/org/jboss/messaging/core/journal/impl/JournalImpl.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/journal/impl/JournalImpl.java 2008-09-08 16:06:47 UTC (rev 4920)
+++ trunk/src/main/org/jboss/messaging/core/journal/impl/JournalImpl.java 2008-09-08 22:21:13 UTC (rev 4921)
@@ -143,7 +143,7 @@
private final Map<Long, PosFiles> posFilesMap = new ConcurrentHashMap<Long, PosFiles>();
- private final Map<Long, JournalTransaction> transactionInfos = new ConcurrentHashMap<Long, JournalTransaction>();
+ private final ConcurrentMap<Long, JournalTransaction> transactionInfos = new ConcurrentHashMap<Long, JournalTransaction>();
private final ConcurrentMap<Long, TransactionCallback> transactionCallbacks = new ConcurrentHashMap<Long, TransactionCallback>();
@@ -631,13 +631,8 @@
throw new IllegalStateException("Journal must be loaded first");
}
- JournalTransaction tx = transactionInfos.get(txID);
+ JournalTransaction tx = getTransactionInfo(txID);
- if (tx == null)
- {
- throw new IllegalStateException("Cannot find tx with id " + txID);
- }
-
ByteBuffer bb = writePrepareTransaction(PREPARE_RECORD, txID, tx, xid);
lock.acquire();
@@ -1015,6 +1010,14 @@
{
TransactionHolder tx = transactions.get(transactionID);
+
+ if (tx == null)
+ {
+ tx = new TransactionHolder(transactionID);
+ transactions.put(transactionID, tx);
+ }
+
+
// We need to read it even if transaction was not found, or the reading checks would fail
byte xidData[] = new byte[preparedTransactionDataSize];
@@ -1029,11 +1032,14 @@
tx.xidData = xidData;
JournalTransaction journalTransaction = transactionInfos.get(transactionID);
+
if (journalTransaction == null)
{
- throw new IllegalStateException("Cannot find tx " + transactionID);
+ journalTransaction = new JournalTransaction();
+
+ transactionInfos.put(transactionID, journalTransaction);
}
-
+
boolean healthy = checkTransactionHealth(journalTransaction, orderedFiles, values);
if (healthy)
@@ -1942,7 +1948,11 @@
{
tx = new JournalTransaction();
- transactionInfos.put(txID, tx);
+ JournalTransaction trans = transactionInfos.putIfAbsent(txID, tx);
+ if (trans != null)
+ {
+ tx = trans;
+ }
}
return tx;
@@ -1957,7 +1967,11 @@
if (callback == null)
{
callback = new TransactionCallback();
- transactionCallbacks.put(transactionId, callback);
+ TransactionCallback callbackCheck = transactionCallbacks.putIfAbsent(transactionId, callback);
+ if (callbackCheck != null)
+ {
+ callback = callbackCheck;
+ }
}
if (callback.errorMessage != null)
Modified: trunk/src/main/org/jboss/messaging/core/transaction/impl/TransactionImpl.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/transaction/impl/TransactionImpl.java 2008-09-08 16:06:47 UTC (rev 4920)
+++ trunk/src/main/org/jboss/messaging/core/transaction/impl/TransactionImpl.java 2008-09-08 22:21:13 UTC (rev 4921)
@@ -230,10 +230,7 @@
pageMessages();
- if (containsPersistent)
- {
- storageManager.prepare(id, xid);
- }
+ storageManager.prepare(id, xid);
state = State.PREPARED;
}
@@ -273,7 +270,7 @@
pageMessages();
}
- if (containsPersistent)
+ if (containsPersistent || xid != null)
{
storageManager.commit(id);
}
@@ -317,7 +314,7 @@
}
}
- if (containsPersistent)
+ if (containsPersistent || xid != null)
{
storageManager.rollback(id);
}
Modified: trunk/tests/src/org/jboss/messaging/tests/integration/xa/BasicXaRecoveryTest.java
===================================================================
--- trunk/tests/src/org/jboss/messaging/tests/integration/xa/BasicXaRecoveryTest.java 2008-09-08 16:06:47 UTC (rev 4920)
+++ trunk/tests/src/org/jboss/messaging/tests/integration/xa/BasicXaRecoveryTest.java 2008-09-08 22:21:13 UTC (rev 4921)
@@ -219,7 +219,44 @@
{
testMultipleTxReceiveWithRollback(true);
}
+
+ public void testNonPersistent() throws Exception
+ {
+ testNonPersistent(true);
+ testNonPersistent(false);
+ }
+
+ public void testNonPersistent(boolean commit) throws Exception
+ {
+ Xid xid = new XidImpl("xa1".getBytes(), 1, new GUID().toString().getBytes());
+
+ ClientMessage m1 = createTextMessage("m1", false);
+ ClientMessage m2 = createTextMessage("m2", false);
+ ClientMessage m3 = createTextMessage("m3", false);
+ ClientMessage m4 = createTextMessage("m4", false);
+
+ clientSession.start(xid, XAResource.TMNOFLAGS);
+ clientProducer.send(m1);
+ clientProducer.send(m2);
+ clientProducer.send(m3);
+ clientProducer.send(m4);
+ clientSession.end(xid, XAResource.TMSUCCESS);
+ clientSession.prepare(xid);
+
+ stopAndRestartServer();
+
+ Xid[] xids = clientSession.recover(XAResource.TMSTARTRSCAN);
+
+ assertEquals(xids.length, 1);
+ assertEquals(xids[0].getFormatId(), xid.getFormatId());
+ assertEqualsByteArrays(xids[0].getBranchQualifier(), xid.getBranchQualifier());
+ assertEqualsByteArrays(xids[0].getGlobalTransactionId(), xid.getGlobalTransactionId());
+ xids = clientSession.recover(XAResource.TMENDRSCAN);
+ assertEquals(xids.length, 0);
+ clientSession.commit(xid, true);
+ }
+
public void testBasicSendWithCommit(boolean stopServer) throws Exception
{
Xid xid = new XidImpl("xa1".getBytes(), 1, new GUID().toString().getBytes());
@@ -973,7 +1010,12 @@
private ClientMessage createTextMessage(String s)
{
- ClientMessage message = clientSession.createClientMessage(JBossTextMessage.TYPE, true, 0, System.currentTimeMillis(), (byte) 1);
+ return createTextMessage(s, true);
+ }
+
+ private ClientMessage createTextMessage(String s, boolean durable)
+ {
+ ClientMessage message = clientSession.createClientMessage(JBossTextMessage.TYPE, durable, 0, System.currentTimeMillis(), (byte) 1);
message.getBody().putString(s);
message.getBody().flip();
return message;
Modified: trunk/tests/src/org/jboss/messaging/tests/unit/core/journal/impl/AlignedJournalImplTest.java
===================================================================
--- trunk/tests/src/org/jboss/messaging/tests/unit/core/journal/impl/AlignedJournalImplTest.java 2008-09-08 16:06:47 UTC (rev 4920)
+++ trunk/tests/src/org/jboss/messaging/tests/unit/core/journal/impl/AlignedJournalImplTest.java 2008-09-08 22:21:13 UTC (rev 4921)
@@ -1125,7 +1125,52 @@
assertEquals(10, records.size());
}
+ // It should be ok to write records on NIO, and later read then on AIO
+ public void testEmptyPrepare() throws Exception
+ {
+ final int JOURNAL_SIZE = 512 * 4;
+
+ setupJournal(JOURNAL_SIZE, 1);
+
+ journalImpl.appendPrepareRecord(2l, new SimpleEncoding(10, (byte)'j'));
+
+ journalImpl.forceMoveNextFile();
+
+ journalImpl.appendAddRecord(1l, (byte)0, new SimpleEncoding(10, (byte)'k'));
+
+ setupJournal(JOURNAL_SIZE, 1);
+
+ assertEquals(1, journalImpl.getDataFilesCount());
+
+ assertEquals(1, transactions.size());
+
+ journalImpl.forceMoveNextFile();
+
+ setupJournal(JOURNAL_SIZE, 1);
+ assertEquals(1, journalImpl.getDataFilesCount());
+
+ assertEquals(1, transactions.size());
+
+ journalImpl.appendCommitRecord(2l);
+
+ journalImpl.appendDeleteRecord(1l);
+
+ journalImpl.forceMoveNextFile();
+
+ setupJournal(JOURNAL_SIZE, 0);
+
+ journalImpl.forceMoveNextFile();
+ journalImpl.debugWait();
+ journalImpl.checkAndReclaimFiles();
+
+ assertEquals(0, transactions.size());
+ assertEquals(0, journalImpl.getDataFilesCount());
+
+ }
+
+
+
public void testReclaimingAfterConcurrentAddsAndDeletes() throws Exception
{
final int JOURNAL_SIZE = 10 * 1024;
Modified: trunk/tests/src/org/jboss/messaging/tests/unit/core/transaction/impl/TransactionImplTest.java
===================================================================
--- trunk/tests/src/org/jboss/messaging/tests/unit/core/transaction/impl/TransactionImplTest.java 2008-09-08 16:06:47 UTC (rev 4920)
+++ trunk/tests/src/org/jboss/messaging/tests/unit/core/transaction/impl/TransactionImplTest.java 2008-09-08 22:21:13 UTC (rev 4921)
@@ -269,10 +269,17 @@
}
- tx = createTransactionXA();
-
+ CreatedTrans resultTrans = createTransactionXA();
+ tx = resultTrans.tx;
assertEquals(Transaction.State.ACTIVE, tx.getState());
+ EasyMock.reset(resultTrans.sm);
+
+ resultTrans.sm.prepare(EasyMock.eq(resultTrans.txId), EasyMock.eq(resultTrans.xid));
+ resultTrans.sm.commit(resultTrans.txId);
+
+ EasyMock.replay(resultTrans.sm);
+
tx.prepare();
tx.commit();
@@ -332,8 +339,20 @@
//OK
}
- tx = createTransactionXA();
+ EasyMock.verify(resultTrans.sm);
+
+ resultTrans = createTransactionXA();
+ tx = resultTrans.tx;
+
+
+ EasyMock.reset(resultTrans.sm);
+
+ resultTrans.sm.prepare(resultTrans.txId, resultTrans.xid);
+ resultTrans.sm.rollback(resultTrans.txId);
+
+ EasyMock.replay(resultTrans.sm);
+
assertEquals(Transaction.State.ACTIVE, tx.getState());
tx.prepare();
@@ -393,7 +412,11 @@
catch (IllegalStateException e)
{
//OK
- }
+ }
+
+
+ EasyMock.verify(resultTrans.sm);
+
}
// public void testSendCommit() throws Exception
@@ -636,28 +659,44 @@
return tx;
}
- private Transaction createTransactionXA()
+ private CreatedTrans createTransactionXA() throws Exception
{
- StorageManager sm = EasyMock.createStrictMock(StorageManager.class);
+ CreatedTrans trans = new CreatedTrans();
- PostOffice po = EasyMock.createStrictMock(PostOffice.class);
+ trans.sm = EasyMock.createStrictMock(StorageManager.class);
- final long txID = 123L;
+ trans.po = EasyMock.createMock(PostOffice.class);
- EasyMock.expect(sm.generateTransactionID()).andReturn(txID);
+ EasyMock.expect(trans.po.getPagingManager()).andStubReturn(null);
+
+ trans.txId = 123L;
- EasyMock.replay(sm);
+ trans.xid = randomXid();
- Xid xid = randomXid();
+ EasyMock.expect(trans.sm.generateTransactionID()).andReturn(trans.txId);
+
+ EasyMock.replay(trans.sm, trans.po);
+
+ trans.tx = new TransactionImpl(trans.xid, trans.sm, trans.po);
- Transaction tx = new TransactionImpl(xid, sm, po);
+ EasyMock.verify(trans.sm, trans.po);
- EasyMock.verify(sm);
+ EasyMock.reset(trans.sm, trans.po);
- return tx;
+ return trans;
}
// Inner classes -----------------------------------------------------------------------
+
+
+ class CreatedTrans
+ {
+ TransactionImpl tx;
+ PostOffice po;
+ StorageManager sm;
+ Xid xid;
+ long txId;
+ }
}
More information about the jboss-cvs-commits
mailing list