Implement upstream message overflow for presents sessions
* Subclasses use setOverflowLimit to enable * handleMessage appends to overflow queue and starts interval if overflow is currently active or if the throttle fails a message * handleMessage ends session if overflow limit is reached * New handleOverflowMessages dispatches messages from the overflow queue until none are left or throttle fails git-svn-id: svn+ssh://src.earth.threerings.net/narya/trunk@5417 542714f4-19e9-0310-aa3c-eee0fc999fb1
This commit is contained in:
@@ -23,6 +23,7 @@ package com.threerings.presents.server;
|
|||||||
|
|
||||||
import java.net.InetAddress;
|
import java.net.InetAddress;
|
||||||
|
|
||||||
|
import java.util.LinkedList;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.TimeZone;
|
import java.util.TimeZone;
|
||||||
|
|
||||||
@@ -34,7 +35,9 @@ import com.google.inject.Inject;
|
|||||||
|
|
||||||
import com.samskivert.util.IntMap;
|
import com.samskivert.util.IntMap;
|
||||||
import com.samskivert.util.IntMaps;
|
import com.samskivert.util.IntMaps;
|
||||||
|
import com.samskivert.util.Interval;
|
||||||
import com.samskivert.util.ResultListener;
|
import com.samskivert.util.ResultListener;
|
||||||
|
import com.samskivert.util.RunQueue;
|
||||||
import com.samskivert.util.Throttle;
|
import com.samskivert.util.Throttle;
|
||||||
|
|
||||||
import com.threerings.util.Name;
|
import com.threerings.util.Name;
|
||||||
@@ -323,6 +326,11 @@ public class PresentsClient
|
|||||||
|
|
||||||
// clear out the client object so that we know the session is over
|
// clear out the client object so that we know the session is over
|
||||||
_clobj = null;
|
_clobj = null;
|
||||||
|
|
||||||
|
// Cancel overflow processing
|
||||||
|
if (_overflow != null) {
|
||||||
|
_overflow.stop();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -386,10 +394,62 @@ public class PresentsClient
|
|||||||
// ", msg=" + message + "].");
|
// ", msg=" + message + "].");
|
||||||
return;
|
return;
|
||||||
|
|
||||||
|
} else if (_overflow != null && _overflow.size() > 0) {
|
||||||
|
// We've got some overflow, the new message has to dispatch after those
|
||||||
|
if (_overflow.size() >= _overflow.limit) {
|
||||||
|
// Can't take any more overflow, bail out
|
||||||
|
handleThrottleExceeded();
|
||||||
|
|
||||||
|
} else {
|
||||||
|
_overflow.enqueue(message);
|
||||||
|
_overflow.start();
|
||||||
|
}
|
||||||
|
return;
|
||||||
|
|
||||||
} else if (_throttle.throttleOp(message.received)) {
|
} else if (_throttle.throttleOp(message.received)) {
|
||||||
|
// We're throttled, check if overflow is allowed
|
||||||
|
if (_overflow != null && _overflow.limit > 0) {
|
||||||
|
// Enter overflow mode and save message for later
|
||||||
|
_overflow.enqueue(message);
|
||||||
|
_overflow.start();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
handleThrottleExceeded();
|
handleThrottleExceeded();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
dispatch(message);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Process the overflowed messages one at a time as long as the throttle will let us.
|
||||||
|
*/
|
||||||
|
protected void handleOverflowMessages ()
|
||||||
|
{
|
||||||
|
// Special case, we've already reached our limit since overflow was kicked off
|
||||||
|
if (_throttle == null) {
|
||||||
|
_overflow.stop();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
long now = System.currentTimeMillis();
|
||||||
|
while (true) {
|
||||||
|
// We're done, the next message will be business as usual
|
||||||
|
if (_overflow.size() == 0) {
|
||||||
|
_overflow.stop();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Dispatch the next message if it isn't throttled
|
||||||
|
if (_throttle.throttleOp(now)) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
dispatch(_overflow.dequeue());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
protected void dispatch (UpstreamMessage message)
|
||||||
|
{
|
||||||
// we dispatch to a message dispatcher that is specialized for the particular class of
|
// we dispatch to a message dispatcher that is specialized for the particular class of
|
||||||
// message that we received
|
// message that we received
|
||||||
MessageDispatcher disp = _disps.get(message.getClass());
|
MessageDispatcher disp = _disps.get(message.getClass());
|
||||||
@@ -528,12 +588,13 @@ public class PresentsClient
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Creates our incoming message throttle.
|
* Creates our incoming message throttle. Subclasses can override this to customize the message
|
||||||
|
* rate limit.
|
||||||
*/
|
*/
|
||||||
protected Throttle createIncomingMessageThrottle ()
|
protected Throttle createIncomingMessageThrottle ()
|
||||||
{
|
{
|
||||||
// more than 100 messages in 10 seconds and you're audi like 5000
|
// more than 100 messages in 10 seconds and you're audi like 5000
|
||||||
return new Throttle(100, 10 * 1000L);
|
return new Throttle(DEFAULT_THROTTLE_MESSAGE_LIMIT, DEFAULT_THROTTLE_TIME_LIMIT);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -541,11 +602,26 @@ public class PresentsClient
|
|||||||
*/
|
*/
|
||||||
protected void handleThrottleExceeded ()
|
protected void handleThrottleExceeded ()
|
||||||
{
|
{
|
||||||
log.warning("Client sent more than 100 msgs in 10 seconds, disconnecting " + this + ".");
|
log.warning(
|
||||||
|
"Client exceeded message limit, disconnecting", "client", this, "throttle", _throttle);
|
||||||
safeEndSession();
|
safeEndSession();
|
||||||
_throttle = null;
|
_throttle = null;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Allow throttled messages to hang around until permitted to send.
|
||||||
|
*/
|
||||||
|
protected void setOverflowLimit (int limit)
|
||||||
|
{
|
||||||
|
if (_overflow == null) {
|
||||||
|
if (limit == 0) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
_overflow = new Overflow();
|
||||||
|
}
|
||||||
|
_overflow.limit = limit;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Makes a note that this client is no longer subscribed to this object. The subscription map
|
* Makes a note that this client is no longer subscribed to this object. The subscription map
|
||||||
* is used to clean up after the client when it goes away.
|
* is used to clean up after the client when it goes away.
|
||||||
@@ -1008,6 +1084,71 @@ public class PresentsClient
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Contains information about message overflow.
|
||||||
|
*/
|
||||||
|
protected class Overflow
|
||||||
|
{
|
||||||
|
/** Maxmimum size of queue before disconnection. */
|
||||||
|
public int limit;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Starts processing the overflowed messages every so often.
|
||||||
|
*/
|
||||||
|
public void start ()
|
||||||
|
{
|
||||||
|
if (_interval == null) {
|
||||||
|
_interval = new Interval(_omgr) {
|
||||||
|
public void expired () {
|
||||||
|
handleOverflowMessages();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
_interval.schedule(OVERFLOW_CHECK_INTERVAL, true);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Stops processing overflowed messages.
|
||||||
|
*/
|
||||||
|
public void stop ()
|
||||||
|
{
|
||||||
|
if (_interval != null) {
|
||||||
|
_interval.cancel();
|
||||||
|
_interval = null;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Gets the number of overflowed messages.
|
||||||
|
*/
|
||||||
|
public int size ()
|
||||||
|
{
|
||||||
|
return _queue.size();
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Adds a message to the end of the queue.
|
||||||
|
*/
|
||||||
|
public void enqueue (UpstreamMessage message)
|
||||||
|
{
|
||||||
|
_queue.add(message);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Removes and returns the least recently added message.
|
||||||
|
*/
|
||||||
|
public UpstreamMessage dequeue ()
|
||||||
|
{
|
||||||
|
return _queue.remove();
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Interval for checking the queue. */
|
||||||
|
protected Interval _interval;
|
||||||
|
|
||||||
|
/** Previously throttled messages waiting to be sent. */
|
||||||
|
protected LinkedList<UpstreamMessage> _queue = new LinkedList<UpstreamMessage>();
|
||||||
|
}
|
||||||
|
|
||||||
@Inject protected ClientManager _clmgr;
|
@Inject protected ClientManager _clmgr;
|
||||||
@Inject protected ConnectionManager _conmgr;
|
@Inject protected ConnectionManager _conmgr;
|
||||||
@Inject protected PresentsDObjectMgr _omgr;
|
@Inject protected PresentsDObjectMgr _omgr;
|
||||||
@@ -1034,6 +1175,9 @@ public class PresentsClient
|
|||||||
/** Prevent the client from sending too many messages too frequently. */
|
/** Prevent the client from sending too many messages too frequently. */
|
||||||
protected Throttle _throttle = createIncomingMessageThrottle();
|
protected Throttle _throttle = createIncomingMessageThrottle();
|
||||||
|
|
||||||
|
/** Overflow data, null unless active. */
|
||||||
|
protected Overflow _overflow;
|
||||||
|
|
||||||
// keep these for kicks and giggles
|
// keep these for kicks and giggles
|
||||||
protected int _messagesIn;
|
protected int _messagesIn;
|
||||||
protected int _messagesOut;
|
protected int _messagesOut;
|
||||||
@@ -1046,6 +1190,15 @@ public class PresentsClient
|
|||||||
* ended. */
|
* ended. */
|
||||||
protected static final long FLUSH_TIME = 7 * 60 * 1000L;
|
protected static final long FLUSH_TIME = 7 * 60 * 1000L;
|
||||||
|
|
||||||
|
/** Maximum number of messages in allowed in a time period. */
|
||||||
|
protected static final int DEFAULT_THROTTLE_MESSAGE_LIMIT = 100;
|
||||||
|
|
||||||
|
/** Time period over which the message limit applies. */
|
||||||
|
protected static final long DEFAULT_THROTTLE_TIME_LIMIT = 10*1000;
|
||||||
|
|
||||||
|
/** Time between checks of the overflow queue. */
|
||||||
|
protected static final int OVERFLOW_CHECK_INTERVAL = 250;
|
||||||
|
|
||||||
// register our message dispatchers
|
// register our message dispatchers
|
||||||
static {
|
static {
|
||||||
_disps.put(SubscribeRequest.class, new SubscribeDispatcher());
|
_disps.put(SubscribeRequest.class, new SubscribeDispatcher());
|
||||||
|
|||||||
@@ -61,4 +61,9 @@ public class VariableThrottle extends Throttle
|
|||||||
}
|
}
|
||||||
_ops = ops;
|
_ops = ops;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public String toString ()
|
||||||
|
{
|
||||||
|
return "VariableThrottle [" + _ops.length + " per " + _period + "ms]";
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user