[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