CVS update: JGroups/src/org/jgroups/protocols UFC.java MFC.java FlowControl.java
"Bela Ban" <[email protected]>
| Newsgroups | gmane.comp.java.javagroups.cvs |
|---|---|
| Message-ID | <[email protected]> |
User: belaban
Date: 10/09/13 09:08:37
Modified: src/org/jgroups/protocols UFC.java MFC.java FlowControl.java
Log:
ns
Revision Changes Path
1.4 +11 -20 JGroups/src/org/jgroups/protocols/UFC.java
Index: UFC.java
===================================================================
RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/UFC.java,v
retrieving revision 1.3
retrieving revision 1.4
diff -u -r1.3 -r1.4
--- UFC.java 9 Sep 2010 11:34:47 -0000 1.3
+++ UFC.java 13 Sep 2010 09:08:36 -0000 1.4
@@ -31,7 +31,7 @@
* <li>Receivers don't send the full credits (max_credits), but rather the actual number of bytes received
* <ol/>
* @author Bela Ban
- * @version $Id: UFC.java,v 1.3 2010/09/09 11:34:47 belaban Exp $
+ * @version $Id: UFC.java,v 1.4 2010/09/13 09:08:36 belaban Exp $
*/
@MBean(description="Simple flow control protocol based on a credit system")
public class UFC extends FlowControl {
@@ -92,27 +92,22 @@
public void stop() {
super.stop();
- for(final Credit cred: sent.values()) {
- synchronized(cred) {
- cred.set(max_credits);
- cred.notifyAll();
- }
- }
+ for(Credit cred: sent.values())
+ cred.set(max_credits);
}
- protected Object handleDownMessage(final Event evt, final Message msg, int length) {
- Address dest=msg.getDest();
- if(dest == null || dest.isMulticastAddress()) {
+ protected Object handleDownMessage(final Event evt, final Message msg, Address dest, int length) {
+ if(dest == null || dest.isMulticastAddress()) { // 2nd line of defense, not really needed
log.error(getClass().getSimpleName() + " doesn't handle multicast messages; passing message down");
return down_prot.down(evt);
}
Credit cred=sent.get(dest);
if(cred == null) {
- log.error("destination " + dest + " not found; passing message down");
+ log.warn("destination " + dest + " not found; passing message down");
return down_prot.down(evt);
}
@@ -137,8 +132,7 @@
if(mbrs == null) return;
// add members not in membership to received and sent hashmap (with full credits)
- for(int i=0; i < mbrs.size(); i++) {
- Address addr=mbrs.elementAt(i);
+ for(Address addr: mbrs) {
if(!sent.containsKey(addr))
sent.put(addr, new Credit(max_credits));
}
@@ -153,16 +147,13 @@
protected void handleCredit(Address sender, long increase) {
- if(sender == null) return;
- StringBuilder sb=null;
-
- Credit cred=sent.get(sender);
- if(cred == null)
+ Credit cred;
+ if(sender == null || (cred=sent.get(sender)) == null || increase <= 0)
return;
- long new_credit=Math.min(max_credits, cred.get() + increase);
+ long new_credit=Math.min(max_credits, cred.get() + increase);
if(log.isTraceEnabled()) {
- sb=new StringBuilder();
+ StringBuilder sb=new StringBuilder();
sb.append("received " + increase + " credits from ").append(sender).append(", old credits: ").append(cred)
.append(", new credits: ").append(new_credit);
log.trace(sb);
1.6 +3 -4 JGroups/src/org/jgroups/protocols/MFC.java
Index: MFC.java
===================================================================
RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/MFC.java,v
retrieving revision 1.5
retrieving revision 1.6
diff -u -r1.5 -r1.6
--- MFC.java 10 Sep 2010 10:55:18 -0000 1.5
+++ MFC.java 13 Sep 2010 09:08:36 -0000 1.6
@@ -33,7 +33,7 @@
* <li>Receivers don't send the full credits (max_credits), but rather the actual number of bytes received
* <ol/>
* @author Bela Ban
- * @version $Id: MFC.java,v 1.5 2010/09/10 10:55:18 belaban Exp $
+ * @version $Id: MFC.java,v 1.6 2010/09/13 09:08:36 belaban Exp $
*/
@MBean(description="Simple flow control protocol based on a credit system")
public class MFC extends FlowControl {
@@ -95,9 +95,8 @@
credits.clear();
}
- protected Object handleDownMessage(final Event evt, final Message msg, int length) {
- Address dest=msg.getDest();
- if(dest != null && !dest.isMulticastAddress()) {
+ protected Object handleDownMessage(final Event evt, final Message msg, Address dest, int length) {
+ if(dest != null && !dest.isMulticastAddress()) { // 2nd line of defense, not really needed
log.error(getClass().getSimpleName() + " doesn't handle unicast messages; passing message down");
return down_prot.down(evt);
}
1.7 +31 -66 JGroups/src/org/jgroups/protocols/FlowControl.java
Index: FlowControl.java
===================================================================
RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/FlowControl.java,v
retrieving revision 1.6
retrieving revision 1.7
diff -u -r1.6 -r1.7
--- FlowControl.java 10 Sep 2010 10:53:42 -0000 1.6
+++ FlowControl.java 13 Sep 2010 09:08:36 -0000 1.7
@@ -16,21 +16,10 @@
* Simple flow control protocol based on a credit system. Each sender has a number of credits (bytes
* to send). When the credits have been exhausted, the sender blocks. Each receiver also keeps track of
* how many credits it has received from a sender. When credits for a sender fall below a threshold,
- * the receiver sends more credits to the sender. Works for both unicast and multicast messages.
- * <p/>
- * Note that this protocol must be located towards the top of the stack, or all down_threads from JChannel to this
- * protocol must be set to false ! This is in order to block JChannel.send()/JChannel.down().
- * <br/>This is the second simplified implementation of the same model. The algorithm is sketched out in
- * doc/FlowControl.txt
- * <br/>
- * Changes (Brian) April 2006:
- * <ol>
- * <li>Receivers now send credits to a sender when more than min_credits have been received (rather than when min_credits
- * are left)
- * <li>Receivers don't send the full credits (max_credits), but rather the actual number of bytes received
- * <ol/>
+ * the receiver sends more credits to the sender.
+ *
* @author Bela Ban
- * @version $Id: FlowControl.java,v 1.6 2010/09/10 10:53:42 belaban Exp $
+ * @version $Id: FlowControl.java,v 1.7 2010/09/13 09:08:36 belaban Exp $
*/
@MBean(description="Simple flow control protocol based on a credit system")
public abstract class FlowControl extends Protocol {
@@ -65,9 +54,6 @@
*/
protected Map<Long,Long> max_block_times=null;
- /** Keeps track of the end time after which a message should not get blocked anymore */
- // protected static final ThreadLocal<Long> end_time=new ThreadLocal<Long>();
-
/**
* If we've received (min_threshold * max_credits) bytes from P, we send more credits to P. Example: if
@@ -107,20 +93,17 @@
/**
- * Keeps track of credits / member at the receiver's side. Keys are members, values are credits left (in bytes).
- * For each receive, the credits for the sender are decremented by the size of the received message.
- * When the credits fall below the threshold, we refill and send a REPLENISH message to the sender.
- * The sender blocks until REPLENISH message is received.
+ * Keeps track of credits per member at the receiver. For each message, the credits for the sender are decremented
+ * by the size of the received message. When the credits fall below the threshold, we refill and send a REPLENISH
+ * message to the sender.
*/
protected final Map<Address,Credit> received=new ConcurrentHashMap<Address,Credit>(11);
- /**
- * Whether FlowControl is still running, this is set to false when the protocol terminates (on stop())
- */
+ /** Whether FlowControl is still running, this is set to false when the protocol terminates (on stop()) */
protected volatile boolean running=true;
-
+
protected boolean frag_size_received=false;
@@ -167,7 +150,6 @@
this.min_credits=min_credits;
}
-
public abstract int getNumberOfBlockings();
public long getMaxBlockTime() {
@@ -288,7 +270,6 @@
if(length <= entry.getKey())
break;
}
-
return retval != null? retval : 0;
}
@@ -303,9 +284,9 @@
/**
- * Allows to unblock a blocked sender from an external program, e.g. JMX
+ * Allows to unblock all blocked senders from an external program, e.g. JMX
*/
- @ManagedOperation(description="Unblock a sender")
+ @ManagedOperation(description="Unblocks all senders")
public void unblock() {
;
}
@@ -322,9 +303,7 @@
log.warn("No fragmentation protocol was found. When flow control is used, we recommend " +
"a fragmentation protocol, due to http://jira.jboss.com/jira/browse/JGRP-590");
}
-
running=true;
-
}
public void stop() {
@@ -358,11 +337,12 @@
log.trace("bypassing flow control because of synchronous response " + Thread.currentThread());
break;
}
+ return handleDownMessage(evt, msg, dest, length);
- return handleDownMessage(evt, msg, length);
case Event.CONFIG:
handleConfigEvent((Map<String,Object>)evt.getArg());
break;
+
case Event.VIEW_CHANGE:
handleViewChange(((View)evt.getArg()).getMembers());
break;
@@ -376,9 +356,6 @@
switch(evt.getType()) {
case Event.MSG:
-
- // JGRP-465. We only deal with msgs to avoid having to use a concurrent collection; ignore views,
- // suspicions, etc which can come up on unusual threads.
Message msg=(Message)evt.getArg();
if(msg.isFlagSet(Message.NO_FC))
break;
@@ -457,7 +434,7 @@
}
- protected abstract Object handleDownMessage(final Event evt, final Message msg, int length);
+ protected abstract Object handleDownMessage(final Event evt, final Message msg, Address dest, int length);
@@ -470,16 +447,11 @@
* @return long Number of credits to be sent. Greater than 0 if credits needs to be sent, 0 otherwise
*/
protected long adjustCredit(Map<Address,Credit> map, Address sender, int length) {
- if(sender == null || length == 0)
- return 0;
-
- Credit cred=map.get(sender);
- if(cred == null)
+ Credit cred;
+ if(sender == null || length == 0 || (cred=map.get(sender)) == null)
return 0;
-
if(log.isTraceEnabled())
log.trace(sender + " used " + length + " credits, " + (cred.get() - length) + " remaining");
-
return cred.decrementAndGet(length);
}
@@ -530,30 +502,23 @@
protected void handleViewChange(Vector<Address> mbrs) {
- Address addr;
if(mbrs == null) return;
if(log.isTraceEnabled()) log.trace("new membership: " + mbrs);
-
// add members not in membership to received and sent hashmap (with full credits)
- for(int i=0; i < mbrs.size(); i++) {
- addr=mbrs.elementAt(i);
+ for(Address addr: mbrs) {
if(!received.containsKey(addr))
received.put(addr, new Credit(max_credits));
}
// remove members that left
for(Iterator<Address> it=received.keySet().iterator(); it.hasNext();) {
- addr=it.next();
+ Address addr=it.next();
if(!mbrs.contains(addr))
it.remove();
}
}
-
- protected static long computeLowestCredit(Map<Address,Credit> m) {
- Collection<Credit> credits=m.values();
- return Collections.min(credits).get();
- }
+
protected static String printMap(Map<Address,Credit> m) {
StringBuilder sb=new StringBuilder();
@@ -565,7 +530,7 @@
- protected class Credit implements Comparable {
+ protected class Credit {
protected long credits_left;
protected int num_blockings=0;
protected long total_blocking_time=0;
@@ -578,10 +543,8 @@
protected synchronized boolean decrementIfEnoughCredits(long credits, long timeout) {
- if(credits <= credits_left) {
- credits_left-=credits;
+ if(decrement(credits))
return true;
- }
if(timeout <= 0)
return false;
@@ -597,6 +560,11 @@
num_blockings++;
}
+ return decrement(credits);
+ }
+
+
+ protected boolean decrement(long credits) {
if(credits <= credits_left) {
credits_left-=credits;
return true;
@@ -616,10 +584,9 @@
}
- protected synchronized long increment(long credits) {
- long retval=credits_left=Math.min(max_credits, credits_left + credits);
+ protected synchronized void increment(long credits) {
+ credits_left=Math.min(max_credits, credits_left + credits);
notifyAll();
- return retval;
}
protected synchronized boolean needToSendCreditRequest() {
@@ -637,17 +604,15 @@
protected synchronized long get() {return credits_left;}
- protected synchronized void set(long new_credits) {credits_left=Math.min(max_credits, new_credits);}
-
+ protected synchronized void set(long new_credits) {
+ credits_left=Math.min(max_credits, new_credits);
+ notifyAll();
+ }
public String toString() {
return String.valueOf(credits_left);
}
- public int compareTo(Object o) {
- Credit other=(Credit)o;
- return credits_left < other.credits_left ? -1 : credits_left > other.credits_left ? 1 : 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