CVS update: JGroups/src/org/jgroups/protocols/pbcast NAKACK.java
"Bela Ban" <[email protected]> Tue, 21 Sep 2010 07:00:27 +0000
| Newsgroups | gmane.comp.java.javagroups.cvs |
|---|---|
| Message-ID | <[email protected]> |
User: belaban
Date: 10/09/21 07:00:27
Modified: src/org/jgroups/protocols/pbcast NAKACK.java
Log:
- removed delivered_msgs
- OOB messages are not removed (minus 1 lock) and we don't check for OOB messages availability (minus 1 lock): OOb messages are removed after CAS acquisition
- JIRA: https://jira.jboss.org/browse/JGRP-1104
Revision Changes Path
1.259 +14 -65 JGroups/src/org/jgroups/protocols/pbcast/NAKACK.java
Index: NAKACK.java
===================================================================
RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/pbcast/NAKACK.java,v
retrieving revision 1.258
retrieving revision 1.259
diff -u -r1.258 -r1.259
--- NAKACK.java 20 Sep 2010 11:51:10 -0000 1.258
+++ NAKACK.java 21 Sep 2010 07:00:27 -0000 1.259
@@ -32,7 +32,7 @@
* instead of the requester by setting use_mcast_xmit to true.
*
* @author Bela Ban
- * @version $Id: NAKACK.java,v 1.258 2010/09/20 11:51:10 belaban Exp $
+ * @version $Id: NAKACK.java,v 1.259 2010/09/21 07:00:27 belaban Exp $
*/
@MBean(description="Reliable transmission multipoint FIFO protocol")
@DeprecatedProperty(names={"max_xmit_size", "eager_lock_release", "stats_list_size"})
@@ -237,17 +237,14 @@
protected final BoundedList<String> digest_history=new BoundedList<String>(10);
- /** <em>Regular</em> messages which have been added, but not removed */
- private final AtomicInteger undelivered_msgs=new AtomicInteger(0);
-
public NAKACK() {
}
- @ManagedAttribute
- public int getUndeliveredMessages() {
- return undelivered_msgs.get();
+ @Deprecated
+ public static int getUndeliveredMessages() {
+ return 0;
}
public long getXmitRequestsReceived() {return xmit_reqs_received;}
@@ -391,7 +388,7 @@
* @return
* @deprecated removed in 2.6
*/
- public long getMaxXmitSize() {
+ public static long getMaxXmitSize() {
return -1;
}
@@ -659,8 +656,7 @@
if(hdr == null)
break; // pass up (e.g. unicast msg)
- // discard messages while not yet server (i.e., until JOIN has returned)
- if(!is_server) {
+ if(!is_server) { // discard messages while not yet server (i.e., until JOIN has returned)
if(log.isTraceEnabled())
log.trace("message was discarded (not yet server)");
return null;
@@ -749,8 +745,7 @@
try { // incrementing seqno and adding the msg to sent_msgs needs to be atomic
msg_id=seqno +1;
msg.putHeader(this.id, NakAckHeader.createMessageHeader(msg_id));
- if(win.add(msg_id, msg) && !msg.isFlagSet(Message.OOB))
- undelivered_msgs.incrementAndGet();
+ win.add(msg_id, msg);
seqno=msg_id;
}
catch(Throwable t) {
@@ -795,20 +790,14 @@
if(leaving)
return;
if(log.isWarnEnabled() && log_discard_msgs)
- log.warn(local_addr + ": dropped message from " + sender +
- " (not in xmit_table), keys are " + xmit_table.keySet() +", view=" + view);
+ log.warn(local_addr + ": dropped message from " + sender + " (not in table " + xmit_table.keySet() +"), view=" + view);
return;
}
boolean loopback=local_addr.equals(sender);
- boolean added_to_window=false;
- boolean added=loopback || (added_to_window=win.add(hdr.seqno, msg));
+ boolean added=loopback || win.add(hdr.seqno, msg);
- if(added_to_window && !msg.isFlagSet(Message.OOB))
- undelivered_msgs.incrementAndGet();
-
- // message is passed up if OOB. Later, when remove() is called, we discard it. This affects ordering !
- // http://jira.jboss.com/jira/browse/JGRP-379
+ // OOB msg is passed up. When removed, we discard it. Affects ordering: http://jira.jboss.com/jira/browse/JGRP-379
if(msg.isFlagSet(Message.OOB)) {
if(added) {
if(loopback)
@@ -818,17 +807,6 @@
up_prot.up(new Event(Event.MSG, msg));
}
}
- List<Message> msgs;
- while(!(msgs=win.removeOOBMessages()).isEmpty()) {
- for(Message tmp_msg: msgs) {
- if(tmp_msg.setTransientFlagIfAbsent(Message.OOB_DELIVERED)) {
- up_prot.up(new Event(Event.MSG, tmp_msg));
- }
- }
- }
-
- if(!(win.hasMessagesToRemove() && undelivered_msgs.get() > 0))
- return;
}
// Efficient way of checking whether another thread is already processing messages from 'sender'.
@@ -840,17 +818,6 @@
return;
}
- // Prevents concurrent passing up of messages by different threads (http://jira.jboss.com/jira/browse/JGRP-198);
- // this is all the more important once we have a threadless stack (http://jira.jboss.com/jira/browse/JGRP-181),
- // where lots of threads can come up to this point concurrently, but only 1 is allowed to pass at a time
- // We *can* deliver messages from *different* senders concurrently, e.g. reception of P1, Q1, P2, Q2 can result in
- // delivery of P1, Q1, Q2, P2: FIFO (implemented by NAKACK) says messages need to be delivered in the
- // order in which they were sent by the sender
- int num_regular_msgs_removed=0;
-
- // 2nd line of defense: in case of an exception, remove() might not be called, therefore processing would never
- // be set back to false. If we get an exception and released_processing is not true, then we set
- // processing to false in the finally clause
boolean released_processing=false;
try {
while(true) {
@@ -862,22 +829,11 @@
}
for(final Message msg_to_deliver: msgs) {
-
// discard OOB msg if it has already been delivered (http://jira.jboss.com/jira/browse/JGRP-379)
- if(msg_to_deliver.isFlagSet(Message.OOB)) {
- if(msg_to_deliver.setTransientFlagIfAbsent(Message.OOB_DELIVERED)) {
- timer.execute(new Runnable() {
- public void run() {
- up_prot.up(new Event(Event.MSG, msg_to_deliver));
- }
- });
- }
+ if(msg_to_deliver.isFlagSet(Message.OOB) && !msg_to_deliver.setTransientFlagIfAbsent(Message.OOB_DELIVERED))
continue;
- }
- num_regular_msgs_removed++;
- // Changed by bela Jan 29 2003: not needed (see above)
- //msg_to_deliver.removeHeader(getName());
+ //msg_to_deliver.removeHeader(getName()); // Changed by bela Jan 29 2003: not needed (see above)
try {
up_prot.up(new Event(Event.MSG, msg_to_deliver));
}
@@ -888,11 +844,6 @@
}
}
finally {
- // We keep track of regular messages that we added, but couldn't remove (because of ordering).
- // When we have such messages pending, then even OOB threads will remove and process them
- // http://jira.jboss.com/jira/browse/JGRP-781
- undelivered_msgs.addAndGet(-num_regular_msgs_removed);
-
// processing is always set in win.remove(processing) above and never here ! This code is just a
// 2nd line of defense should there be an exception before win.remove(processing) sets processing
if(!released_processing)
@@ -973,7 +924,7 @@
}
continue;
}
- sendXmitRsp(xmit_requester, msg, i);
+ sendXmitRsp(xmit_requester, msg);
}
}
@@ -1011,9 +962,8 @@
* to preserve the original message's properties, such as src, headers etc.
* @param dest
* @param msg
- * @param seqno
*/
- private void sendXmitRsp(Address dest, Message msg, long seqno) {
+ private void sendXmitRsp(Address dest, Message msg) {
Buffer buf;
if(msg == null) {
if(log.isErrorEnabled())
@@ -1579,7 +1529,6 @@
win.destroy();
}
xmit_table.clear();
- undelivered_msgs.set(0);
}
------------------------------------------------------------------------------
Start uncovering the many advantages of virtual appliances
and start using them to simplify application deployment and
accelerate your shift to cloud computing.
http://p.sf.net/sfu/novell-sfdev2dev