Outgoing message throttle for ActionScript. This uses a timer rather than a
frame listener so that it has a chance of working out of the box on Thane. git-svn-id: svn+ssh://src.earth.threerings.net/narya/trunk@5432 542714f4-19e9-0310-aa3c-eee0fc999fb1
This commit is contained in:
@@ -22,14 +22,14 @@
|
|||||||
package com.threerings.presents.client {
|
package com.threerings.presents.client {
|
||||||
|
|
||||||
import flash.errors.IOError;
|
import flash.errors.IOError;
|
||||||
|
|
||||||
import flash.events.Event;
|
import flash.events.Event;
|
||||||
import flash.events.IOErrorEvent;
|
import flash.events.IOErrorEvent;
|
||||||
|
import flash.events.TimerEvent;
|
||||||
|
|
||||||
import flash.net.Socket;
|
import flash.net.Socket;
|
||||||
|
|
||||||
import flash.utils.ByteArray;
|
import flash.utils.ByteArray;
|
||||||
import flash.utils.Endian;
|
import flash.utils.Endian;
|
||||||
|
import flash.utils.Timer;
|
||||||
|
|
||||||
import com.threerings.util.Log;
|
import com.threerings.util.Log;
|
||||||
|
|
||||||
@@ -59,13 +59,30 @@ public class Communicator
|
|||||||
_outBuffer = new ByteArray();
|
_outBuffer = new ByteArray();
|
||||||
_outBuffer.endian = Endian.BIG_ENDIAN;
|
_outBuffer.endian = Endian.BIG_ENDIAN;
|
||||||
_outStream = new ObjectOutputStream(_outBuffer);
|
_outStream = new ObjectOutputStream(_outBuffer);
|
||||||
|
|
||||||
_inStream = new ObjectInputStream();
|
_inStream = new ObjectInputStream();
|
||||||
|
|
||||||
|
// start up our message writer
|
||||||
|
_writer = new Timer(1);
|
||||||
|
_writer.addEventListener(TimerEvent.TIMER, sendPendingMessages);
|
||||||
|
_writer.start();
|
||||||
|
|
||||||
_portIdx = 0;
|
_portIdx = 0;
|
||||||
logonToPort();
|
logonToPort();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public function logoff () :void
|
||||||
|
{
|
||||||
|
if (_socket == null) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
postMessage(new LogoffRequest());
|
||||||
|
}
|
||||||
|
|
||||||
|
public function postMessage (msg :UpstreamMessage) :void
|
||||||
|
{
|
||||||
|
_outq.push(msg);
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* This method is strangely named, and it does two things which is
|
* This method is strangely named, and it does two things which is
|
||||||
* bad style. Either log on to the next port, or save that the port
|
* bad style. Either log on to the next port, or save that the port
|
||||||
@@ -80,8 +97,7 @@ public class Communicator
|
|||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
if (_portIdx != 0) {
|
if (_portIdx != 0) {
|
||||||
_client.reportLogonTribulations(
|
_client.reportLogonTribulations(new LogonError(AuthCodes.TRYING_NEXT_PORT, true));
|
||||||
new LogonError(AuthCodes.TRYING_NEXT_PORT, true));
|
|
||||||
|
|
||||||
removeListeners();
|
removeListeners();
|
||||||
}
|
}
|
||||||
@@ -94,8 +110,7 @@ public class Communicator
|
|||||||
_socket.addEventListener(Event.CLOSE, socketClosed);
|
_socket.addEventListener(Event.CLOSE, socketClosed);
|
||||||
|
|
||||||
_frameReader = new FrameReader(_socket);
|
_frameReader = new FrameReader(_socket);
|
||||||
_frameReader.addEventListener(FrameAvailableEvent.FRAME_AVAILABLE,
|
_frameReader.addEventListener(FrameAvailableEvent.FRAME_AVAILABLE, inputFrameReceived);
|
||||||
inputFrameReceived);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
var host :String = _client.getHostname();
|
var host :String = _client.getHostname();
|
||||||
@@ -108,8 +123,7 @@ public class Communicator
|
|||||||
_portIdx = -1; // indicate that we're no longer trying new ports
|
_portIdx = -1; // indicate that we're no longer trying new ports
|
||||||
|
|
||||||
} else {
|
} else {
|
||||||
Log.getLog(this).info(
|
log.info("Connecting [host=" + host + ", port=" + port + "].");
|
||||||
"Connecting [host=" + host + ", port=" + port + "].");
|
|
||||||
_socket.connect(host, port);
|
_socket.connect(host, port);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -122,24 +136,7 @@ public class Communicator
|
|||||||
_socket.removeEventListener(IOErrorEvent.IO_ERROR, socketError);
|
_socket.removeEventListener(IOErrorEvent.IO_ERROR, socketError);
|
||||||
_socket.removeEventListener(Event.CLOSE, socketClosed);
|
_socket.removeEventListener(Event.CLOSE, socketClosed);
|
||||||
|
|
||||||
_frameReader.removeEventListener(FrameAvailableEvent.FRAME_AVAILABLE,
|
_frameReader.removeEventListener(FrameAvailableEvent.FRAME_AVAILABLE, inputFrameReceived);
|
||||||
inputFrameReceived);
|
|
||||||
}
|
|
||||||
|
|
||||||
public function logoff () :void
|
|
||||||
{
|
|
||||||
if (_socket == null) {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
sendMessage(new LogoffRequest());
|
|
||||||
|
|
||||||
shutdown(null);
|
|
||||||
}
|
|
||||||
|
|
||||||
public function postMessage (msg :UpstreamMessage) :void
|
|
||||||
{
|
|
||||||
sendMessage(msg); // send it now: we have no out queue
|
|
||||||
}
|
}
|
||||||
|
|
||||||
protected function shutdown (logonError :Error) :void
|
protected function shutdown (logonError :Error) :void
|
||||||
@@ -149,7 +146,7 @@ public class Communicator
|
|||||||
try {
|
try {
|
||||||
_socket.close();
|
_socket.close();
|
||||||
} catch (err :Error) {
|
} catch (err :Error) {
|
||||||
Log.getLog(this).warning("Error closing failed socket [error=" + err + "].");
|
log.warning("Error closing failed socket [error=" + err + "].");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
removeListeners();
|
removeListeners();
|
||||||
@@ -160,21 +157,59 @@ public class Communicator
|
|||||||
_outBuffer = null;
|
_outBuffer = null;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (_writer != null) {
|
||||||
|
_writer.stop();
|
||||||
|
_writer = null;
|
||||||
|
}
|
||||||
|
|
||||||
_client.notifyObservers(ClientEvent.CLIENT_DID_LOGOFF, null);
|
_client.notifyObservers(ClientEvent.CLIENT_DID_LOGOFF, null);
|
||||||
_client.cleanup(logonError);
|
_client.cleanup(logonError);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Sends all pending messages from our outgoing message queue. If we hit our throttle while
|
||||||
|
* sending, we stop and wait for the next time around when we'll try sending them again.
|
||||||
|
*/
|
||||||
|
protected function sendPendingMessages (event :TimerEvent) :void
|
||||||
|
{
|
||||||
|
while (_outq.length > 0) {
|
||||||
|
// if we've exceeded our throttle, stop for now
|
||||||
|
if (_client.getOutgoingMessageThrottle().throttleOp()) {
|
||||||
|
if (_tqsize != _outq.length) {
|
||||||
|
// only log when our outq size changes
|
||||||
|
_tqsize = _outq.length;
|
||||||
|
log.info("Throttling outgoing messages", "queue", _outq.length,
|
||||||
|
"throttle", _client.getOutgoingMessageThrottle());
|
||||||
|
}
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
_tqsize = 0;
|
||||||
|
|
||||||
|
// grab the next message from the queue and send it
|
||||||
|
var msg :UpstreamMessage = (_outq.shift() as UpstreamMessage);
|
||||||
|
sendMessage(msg);
|
||||||
|
|
||||||
|
// if this was a logoff request, shutdown
|
||||||
|
if (msg is LogoffRequest) {
|
||||||
|
shutdown(null);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Writes a single message to our outgoing socket.
|
||||||
|
*/
|
||||||
protected function sendMessage (msg :UpstreamMessage) :void
|
protected function sendMessage (msg :UpstreamMessage) :void
|
||||||
{
|
{
|
||||||
if (_outStream == null) {
|
if (_outStream == null) {
|
||||||
Log.getLog(this).warning("No socket, dropping msg [msg=" + msg + "].");
|
log.warning("No socket, dropping msg [msg=" + msg + "].");
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
// write the message (ends up in _outBuffer)
|
// write the message (ends up in _outBuffer)
|
||||||
_outStream.writeObject(msg);
|
_outStream.writeObject(msg);
|
||||||
|
|
||||||
// Log.debug("outBuffer: " + StringUtil.unhexlate(_outBuffer));
|
// Log.debug("outBuffer: " + StringUtil.unhexlate(_outBuffer));
|
||||||
|
|
||||||
// Frame it by writing the length, then the bytes.
|
// Frame it by writing the length, then the bytes.
|
||||||
// We add 4 to the length, because the length is of the entire frame
|
// We add 4 to the length, because the length is of the entire frame
|
||||||
@@ -192,47 +227,25 @@ public class Communicator
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Returns the time at which we last sent a packet to the server.
|
* Called when a frame of data from the server is ready to be decoded into a DownstreamMessage.
|
||||||
*/
|
|
||||||
internal function getLastWrite () :uint
|
|
||||||
{
|
|
||||||
return _lastWrite;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Makes a note of the time at which we last communicated with the server.
|
|
||||||
*/
|
|
||||||
internal function updateWriteStamp () :void
|
|
||||||
{
|
|
||||||
_lastWrite = flash.utils.getTimer();
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Called when a frame of data from the server is ready to be
|
|
||||||
* decoded into a DownstreamMessage.
|
|
||||||
*/
|
*/
|
||||||
protected function inputFrameReceived (event :FrameAvailableEvent) :void
|
protected function inputFrameReceived (event :FrameAvailableEvent) :void
|
||||||
{
|
{
|
||||||
// convert the frame data into a message from the server
|
// convert the frame data into a message from the server
|
||||||
var frameData :ByteArray = event.getFrameData();
|
var frameData :ByteArray = event.getFrameData();
|
||||||
//Log.debug("length of in frame: " + frameData.length);
|
|
||||||
//Log.debug("inBuffer: " + StringUtil.unhexlate(frameData));
|
|
||||||
_inStream.setSource(frameData);
|
_inStream.setSource(frameData);
|
||||||
var msg :DownstreamMessage;
|
var msg :DownstreamMessage;
|
||||||
try {
|
try {
|
||||||
msg = (_inStream.readObject() as DownstreamMessage);
|
msg = (_inStream.readObject() as DownstreamMessage);
|
||||||
} catch (e :Error) {
|
} catch (e :Error) {
|
||||||
var log :Log = Log.getLog(this);
|
|
||||||
log.warning("Error processing downstream message: " + e);
|
log.warning("Error processing downstream message: " + e);
|
||||||
log.logStackTrace(e);
|
log.logStackTrace(e);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (frameData.bytesAvailable > 0) {
|
if (frameData.bytesAvailable > 0) {
|
||||||
Log.getLog(this).warning(
|
log.warning("Beans! We didn't fully read a frame, is there a bug in some streaming " +
|
||||||
"Beans! We didn't fully read a frame, surely there's " +
|
"code? [bytesLeftOver=" + frameData.bytesAvailable + ", msg=" + msg + "].");
|
||||||
"a bug in some streaming code. " +
|
|
||||||
"[bytesLeftOver=" + frameData.bytesAvailable + ", msg=" + msg + "].");
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if (_omgr != null) {
|
if (_omgr != null) {
|
||||||
@@ -241,7 +254,7 @@ public class Communicator
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
// Otherwise, this would be the AuthResponse to our logon attempt.
|
// otherwise, this would be the AuthResponse to our logon attempt
|
||||||
var rsp :AuthResponse = (msg as AuthResponse);
|
var rsp :AuthResponse = (msg as AuthResponse);
|
||||||
var data :AuthResponseData = rsp.getData();
|
var data :AuthResponseData = rsp.getData();
|
||||||
if (data.code !== AuthResponseData.SUCCESS) {
|
if (data.code !== AuthResponseData.SUCCESS) {
|
||||||
@@ -261,9 +274,8 @@ public class Communicator
|
|||||||
{
|
{
|
||||||
logonToPort(true);
|
logonToPort(true);
|
||||||
// well that's great! let's logon
|
// well that's great! let's logon
|
||||||
var req :AuthRequest = new AuthRequest(
|
postMessage(new AuthRequest(_client.getCredentials(), _client.getVersion(),
|
||||||
_client.getCredentials(), _client.getVersion(), _client.getBootGroups());
|
_client.getBootGroups()));
|
||||||
sendMessage(req);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -280,9 +292,9 @@ public class Communicator
|
|||||||
}
|
}
|
||||||
|
|
||||||
// total failure
|
// total failure
|
||||||
Log.getLog(this).warning("socket error: " + event + ", target=" + event.target);
|
log.warning("Socket error: " + event, "target", event.target);
|
||||||
Log.dumpStack();
|
Log.dumpStack();
|
||||||
shutdown(new Error("socket closed unexpectedly."));
|
shutdown(new Error("Socket closed unexpectedly."));
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -290,9 +302,25 @@ public class Communicator
|
|||||||
*/
|
*/
|
||||||
protected function socketClosed (event :Event) :void
|
protected function socketClosed (event :Event) :void
|
||||||
{
|
{
|
||||||
Log.getLog(this).info("socket was closed: " + event);
|
log.info("Socket was closed: " + event);
|
||||||
_client.notifyObservers(ClientEvent.CLIENT_CONNECTION_FAILED);
|
_client.notifyObservers(ClientEvent.CLIENT_CONNECTION_FAILED);
|
||||||
logoff();
|
shutdown(null);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Returns the time at which we last sent a packet to the server.
|
||||||
|
*/
|
||||||
|
internal function getLastWrite () :uint
|
||||||
|
{
|
||||||
|
return _lastWrite;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Makes a note of the time at which we last communicated with the server.
|
||||||
|
*/
|
||||||
|
internal function updateWriteStamp () :void
|
||||||
|
{
|
||||||
|
_lastWrite = flash.utils.getTimer();
|
||||||
}
|
}
|
||||||
|
|
||||||
protected var _client :Client;
|
protected var _client :Client;
|
||||||
@@ -305,9 +333,14 @@ public class Communicator
|
|||||||
protected var _frameReader :FrameReader;
|
protected var _frameReader :FrameReader;
|
||||||
|
|
||||||
protected var _socket :Socket;
|
protected var _socket :Socket;
|
||||||
|
|
||||||
protected var _lastWrite :uint;
|
protected var _lastWrite :uint;
|
||||||
|
|
||||||
|
protected var _outq :Array = new Array();
|
||||||
|
protected var _writer :Timer;
|
||||||
|
protected var _tqsize :int = 0;
|
||||||
|
|
||||||
|
protected const log :Log = Log.getLog(this);
|
||||||
|
|
||||||
/** The current port we'll try to connect to. */
|
/** The current port we'll try to connect to. */
|
||||||
protected var _portIdx :int = -1;
|
protected var _portIdx :int = -1;
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user