Merged 0.4

This commit is contained in:
Roy Hyunjin Han 2013-04-26 09:42:16 -07:00
commit 4f7185b074
11 changed files with 673 additions and 388 deletions

9
.gitignore vendored
View file

@ -1,8 +1,9 @@
*~ *~
*.sw[op]
*.py[cod]
*.egg
*.egg-info *.egg-info
*.pyc
*.swo
*.swp
.coverage
build build
dist dist
sdist
.coverage

View file

@ -1,8 +1,14 @@
0.4
---
- Added support for server-side callbacks thanks to Zac Lee
- Added low-level _SocketIO to remove cyclic references
- Merged Channel functionality into BaseNamespace thanks to Alexandre Bourget
0.3 0.3
--- ---
- Added support for secure connections - Added support for secure connections
- Added socketIO.wait() - Added socketIO.wait()
- Improved exception handling in heartbeatThread and namespaceThread - Improved exception handling in _RhythmicThread and _ListenerThread
0.2 0.2
--- ---

View file

@ -1,4 +1,4 @@
Copyright (c) 2012 Roy Hyunjin Han and contributors Copyright (c) 2013 Roy Hyunjin Han and contributors
Permission is hereby granted, free of charge, to any person obtaining a copy of this software and associated documentation files (the "Software"), to deal in the Software without restriction, including without limitation the rights to use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the Software, and to permit persons to whom the Software is furnished to do so, subject to the following conditions: Permission is hereby granted, free of charge, to any person obtaining a copy of this software and associated documentation files (the "Software"), to deal in the Software without restriction, including without limitation the rights to use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the Software, and to permit persons to whom the Software is furnished to do so, subject to the following conditions:

View file

@ -2,12 +2,6 @@ socketIO-client
=============== ===============
Here is a socket.io_ client library for Python. You can use it to write test code for your socket.io_ server. Here is a socket.io_ client library for Python. You can use it to write test code for your socket.io_ server.
Thanks to rod_ for the `StackOverflow question and answer`__ on which this code is based.
Thanks to liris_ for websocket-client_ and to guille_ for the `socket.io specification`_.
Thanks to `Paul Kienzle`_, `Josh VanderLinden`_, `Ian Fitzpatrick`_ for submitting code to expand support of the socket.io protocol.
Installation Installation
------------ ------------
@ -22,7 +16,7 @@ Installation
source $VIRTUAL_ENV/bin/activate source $VIRTUAL_ENV/bin/activate
# Install package # Install package
easy_install -U socketIO-client pip install -U socketIO-client
Usage Usage
@ -36,31 +30,32 @@ Emit. ::
from socketIO_client import SocketIO from socketIO_client import SocketIO
socketIO = SocketIO('localhost', 8000) with SocketIO('localhost', 8000) as socketIO:
socketIO.emit('aaa', {'bbb': 'ccc'}) socketIO.emit('aaa')
socketIO.wait(seconds=1) socketIO.wait(seconds=1)
Emit with callback. :: Emit with callback. ::
from socketIO_client import SocketIO from socketIO_client import SocketIO
def on_response(*args): def on_bbb_response(*args):
print args print 'on_bbb_response', args
socketIO = SocketIO('localhost', 8000) with SocketIO('localhost', 8000) as socketIO:
socketIO.emit('aaa', {'bbb': 'ccc'}, on_response) socketIO.emit('bbb', {'xxx': 'yyy'}, on_bbb_response)
socketIO.wait(forCallbacks=True) socketIO.wait_for_callbacks(seconds=1)
Define events. :: Define events. ::
from socketIO_client import SocketIO from socketIO_client import SocketIO
def on_ddd(*args): def on_aaa_response(*args):
print args print 'on_aaa_response', args
socketIO = SocketIO('localhost', 8000) socketIO = SocketIO('localhost', 8000)
socketIO.on('ddd', on_ddd) socketIO.on('aaa_response', on_aaa_response)
socketIO.wait() socketIO.emit('aaa')
socketIO.wait(seconds=1)
Define events in a namespace. :: Define events in a namespace. ::
@ -68,11 +63,14 @@ Define events in a namespace. ::
class Namespace(BaseNamespace): class Namespace(BaseNamespace):
def on_ddd(self, *args): def on_aaa_response(self, *args):
self.socketIO.emit('eee', {'fff': 'ggg'}) print 'on_aaa_response', args
self.emit('bbb')
socketIO = SocketIO('localhost', 8000, Namespace) socketIO = SocketIO('localhost', 8000)
socketIO.wait() socketIO.define(Namespace)
socketIO.emit('aaa')
socketIO.wait(seconds=1)
Define standard events. :: Define standard events. ::
@ -80,44 +78,40 @@ Define standard events. ::
class Namespace(BaseNamespace): class Namespace(BaseNamespace):
def on_connect(self, socketIO): def on_connect(self):
print '[Connected]' print '[Connected]'
def on_disconnect(self): socketIO = SocketIO('localhost', 8000)
print '[Disconnected]' socketIO.define(Namespace)
socketIO.wait(seconds=1)
def on_error(self, name, message): Define different namespaces on a single socket. ::
print '[Error] %s: %s' % (name, message)
def on_message(self, id, message):
print '[Message] %s: %s' % (id, message)
socketIO = SocketIO('localhost', 8000, Namespace)
socketIO.wait()
Define different behavior for different channels on a single socket. ::
from socketIO_client import SocketIO, BaseNamespace from socketIO_client import SocketIO, BaseNamespace
class MainNamespace(BaseNamespace):
def on_aaa(self, *args):
print 'aaa', args
class ChatNamespace(BaseNamespace): class ChatNamespace(BaseNamespace):
def on_bbb(self, *args): def on_aaa_response(self, *args):
print 'bbb', args print 'on_aaa_response', args
class NewsNamespace(BaseNamespace): class NewsNamespace(BaseNamespace):
def on_ccc(self, *args): def on_aaa_response(self, *args):
print 'ccc', args print 'on_aaa_response', args
mainSocket = SocketIO('localhost', 8000, MainNamespace) socketIO = SocketIO('localhost', 8000)
chatSocket = mainSocket.connect('/chat', ChatNamespace) chatNamespace = socketIO.define(ChatNamespace, '/chat')
newsSocket = mainSocket.connect('/news', NewsNamespace) newsNamespace = socketIO.define(NewsNamespace, '/news')
mainSocket.wait()
chatNamespace.emit('aaa')
newsNamespace.emit('aaa')
socketIO.wait(seconds=1)
Open secure websockets (HTTPS / WSS) behind a proxy. ::
SocketIO('localhost', 8000,
secure=True,
proxies={'https': 'https://proxy.example.com:8080'})
License License
@ -125,14 +119,31 @@ License
This software is available under the MIT License. This software is available under the MIT License.
Credits
-------
- `Guillermo Rauch`_ wrote the `socket.io specification`_.
- `Hiroki Ohtani`_ wrote websocket-client_.
- rod_ wrote a `prototype for a Python client to a socket.io server`_ on StackOverflow.
- `Alexandre Bourget`_ wrote gevent-socketio_, which is a socket.io server written in Python.
- `Paul Kienzle`_, `Zac Lee`_, `Josh VanderLinden`_, `Ian Fitzpatrick`_, `Lucas Klein`_ submitted code to expand support of the socket.io protocol.
.. _socket.io: http://socket.io .. _socket.io: http://socket.io
.. _rod: http://stackoverflow.com/users/370115/rod
.. _StackOverflowQA: http://stackoverflow.com/questions/6692908/formatting-messages-to-send-to-socket-io-node-js-server-from-python-client .. _Guillermo Rauch: https://github.com/guille
__ StackOverflowQA_
.. _liris: https://github.com/liris
.. _websocket-client: https://github.com/liris/websocket-client
.. _guille: https://github.com/guille
.. _socket.io specification: https://github.com/LearnBoost/socket.io-spec .. _socket.io specification: https://github.com/LearnBoost/socket.io-spec
.. _Hiroki Ohtani: https://github.com/liris
.. _websocket-client: https://github.com/liris/websocket-client
.. _rod: http://stackoverflow.com/users/370115/rod
.. _prototype for a Python client to a socket.io server: http://stackoverflow.com/questions/6692908/formatting-messages-to-send-to-socket-io-node-js-server-from-python-client
.. _Alexandre Bourget: https://github.com/abourget
.. _gevent-socketio: https://github.com/abourget/gevent-socketio
.. _Paul Kienzle: https://github.com/pkienzle .. _Paul Kienzle: https://github.com/pkienzle
.. _Zac Lee: https://github.com/zratic
.. _Josh VanderLinden: https://github.com/codekoala .. _Josh VanderLinden: https://github.com/codekoala
.. _Ian Fitzpatrick: https://github.com/GraphEffect .. _Ian Fitzpatrick: https://github.com/GraphEffect
.. _Lucas Klein: https://github.com/lukashed

7
TODO.goals Normal file
View file

@ -0,0 +1,7 @@
= Resolve pull requests
Resolve issues
Investigate issue #8
Examine forks
Integrate Sajal's fork #7
Integrate Francis's fork #10
Integrate Paul's fork

View file

71
serve_tests.js Normal file
View file

@ -0,0 +1,71 @@
var io = require('socket.io').listen(8000);
var main = io.of('').on('connection', function(socket) {
socket.on('message', function(data, fn) {
if (fn) { // Client expects a callback
if (data) {
fn(data);
} else {
fn();
}
} else if (typeof data === 'object') {
socket.json.send(data ? data : 'message_response'); // object or null
} else {
socket.send(data ? data : 'message_response'); // string or ''
}
});
socket.on('emit', function() {
socket.emit('emit_response');
});
socket.on('emit_with_payload', function(payload) {
socket.emit('emit_with_payload_response', payload);
});
socket.on('emit_with_multiple_payloads', function(payload, payload) {
socket.emit('emit_with_multiple_payloads_response', payload, payload);
});
socket.on('emit_with_callback', function(fn) {
fn();
});
socket.on('emit_with_callback_with_payload', function(fn) {
fn(PAYLOAD);
});
socket.on('emit_with_callback_with_multiple_payloads', function(fn) {
fn(PAYLOAD, PAYLOAD);
});
socket.on('emit_with_event', function(payload) {
socket.emit('emit_with_event_response', payload);
});
socket.on('ack', function(payload) {
socket.emit('ack_response', payload, function(payload) {
socket.emit('ack_callback_response', payload);
});
});
socket.on('aaa', function() {
socket.emit('aaa_response', PAYLOAD);
});
socket.on('bbb', function(payload, fn) {
if (fn) {
fn(payload);
}
});
});
var chat = io.of('/chat').on('connection', function (socket) {
socket.on('emit_with_payload', function(payload) {
socket.emit('emit_with_payload_response', payload);
});
socket.on('aaa', function() {
socket.emit('aaa_response', 'in chat');
});
});
var news = io.of('/news').on('connection', function (socket) {
socket.on('emit_with_payload', function(payload) {
socket.emit('emit_with_payload_response', payload);
});
socket.on('aaa', function() {
socket.emit('aaa_response', 'in news');
});
});
var PAYLOAD = {'xxx': 'yyy'};

View file

@ -1,29 +0,0 @@
'Launch this server in another terminal window before running tests'
from socketio import socketio_manage
from socketio.namespace import BaseNamespace
from socketio.server import SocketIOServer
class Namespace(BaseNamespace):
def on_aaa(self, *args):
self.socket.send_packet(dict(
type='event',
name='ddd',
args=args,
endpoint=self.ns_name))
class Application(object):
def __call__(self, environ, start_response):
socketio_manage(environ, {
'': Namespace,
'/chat': Namespace,
'/news': Namespace,
})
if __name__ == '__main__':
socketIOServer = SocketIOServer(('0.0.0.0', 8000), Application())
socketIOServer.serve_forever()

2
setup.py Executable file → Normal file
View file

@ -9,7 +9,7 @@ CHANGES = open(os.path.join(here, 'CHANGES.rst')).read()
setup( setup(
name='socketIO-client', name='socketIO-client',
version='0.3', version='0.4',
description='A socket.io client library', description='A socket.io client library',
long_description=README + '\n\n' + CHANGES, long_description=README + '\n\n' + CHANGES,
license='MIT', license='MIT',

View file

@ -1,35 +1,55 @@
import websocket import socket
from anyjson import dumps, loads from json import dumps, loads
from threading import Thread, Event from threading import Thread, Event
from time import sleep from time import sleep
from urllib import urlopen from urllib import urlopen
from websocket import WebSocketConnectionClosedException, create_connection
__version__ = '0.3' PROTOCOL = 1 # socket.io protocol version
PROTOCOL = 1 # SocketIO protocol version
class BaseNamespace(object): # pragma: no cover class BaseNamespace(object): # pragma: no cover
'Define socket.io behavior'
def __init__(self, socketIO): def __init__(self, _socketIO, path):
self.socketIO = socketIO self._socketIO = _socketIO
self._path = path
self._callbackByEvent = {}
self.initialize()
def on_connect(self, socketIO): def initialize(self):
'Initialize custom variables here; you can override this method'
pass
def on_connect(self):
'Called when socket is connecting; you can override this method'
pass pass
def on_disconnect(self): def on_disconnect(self):
'Called when socket is disconnecting; you can override this method'
pass pass
def on_error(self, reason, advice): def on_error(self, reason, advice):
'Called when server sends an error; you can override this method'
print '[Error] %s' % advice print '[Error] %s' % advice
def on_message(self, messageData): def on_message(self, data):
print '[Message] %s' % messageData 'Called when server sends a message; you can override this method'
print '[Message] %s' % data
def on_(self, eventName, *eventArguments): def on_event(self, event, *args):
print '[Event] %s%s' % (eventName, eventArguments) """
Called when server emits an event; you can override this method.
Called only if the program cannot find a more specific event handler,
such as one defined by namespace.on('my_event', my_function).
"""
callback, args = find_callback(args)
arguments = [repr(_) for _ in args]
if callback:
arguments.append('callback(*args)')
callback(*args)
print '[Event] %s(%s)' % (event, ', '.join(arguments))
def on_open(self, *args): def on_open(self, *args):
print '[Open]', args print '[Open]', args
@ -43,143 +63,96 @@ class BaseNamespace(object): # pragma: no cover
def on_reconnect(self, *args): def on_reconnect(self, *args):
print '[Reconnect]', args print '[Reconnect]', args
def message(self, data='', callback=None):
self._socketIO.message(data, callback, path=self._path)
def emit(self, event, *args, **kw):
kw['path'] = self._path
self._socketIO.emit(event, *args, **kw)
def on(self, event, callback):
'Define a callback to handle a custom event emitted by the server'
self._callbackByEvent[event] = callback
def _get_eventCallback(self, event):
# Check callbacks defined by on()
try:
return self._callbackByEvent[event]
except KeyError:
pass
# Check callbacks defined explicitly or use on_event()
callback = lambda *args: self.on_event(event, *args)
return getattr(self, 'on_' + event.replace(' ', '_'), callback)
class SocketIO(object): class SocketIO(object):
messageID = 0 def __init__(self, host, port, secure=False, proxies=None):
"""
Create a socket.io client that connects to a socket.io server
at the specified host and port. Set secure=True to use HTTPS / WSS.
def __init__(self, host, port, Namespace=BaseNamespace, secure=False): SocketIO('localhost', 8000, secure=True,
self.host = host proxies={'https': 'https://proxy.example.com:8080'})
self.port = int(port) """
self.namespace = Namespace(self) self._socketIO = _SocketIO(host, port, secure, proxies)
self.secure = secure self._namespaceByPath = {}
self.__connect() self.define(BaseNamespace) # Define default namespace
heartbeatInterval = self.heartbeatTimeout - 2 self._rhythmicThread = _RhythmicThread(
self.heartbeatThread = RhythmicThread(heartbeatInterval, self._socketIO.heartbeatInterval,
self._send_heartbeat) self._socketIO.send_heartbeat)
self.heartbeatThread.start() self._rhythmicThread.start()
self.channelByName = {} self._listenerThread = _ListenerThread(
self.callbackByEvent = {} self._socketIO,
self.namespaceThread = ListenerThread(self) self._namespaceByPath)
self.namespaceThread.start() self._listenerThread.start()
def __del__(self): # pragma: no cover def __enter__(self):
self.heartbeatThread.cancel() return self
self.namespaceThread.cancel()
self.connection.close()
def __connect(self): def __exit__(self, exc_type, exc_value, traceback):
baseURL = '%s:%d/socket.io/%s' % (self.host, self.port, PROTOCOL) self.disconnect()
try:
response = urlopen('%s://%s/' % (
'https' if self.secure else 'http', baseURL))
except IOError: # pragma: no cover
raise SocketIOError('Could not start connection')
if 200 != response.getcode(): # pragma: no cover
raise SocketIOError('Could not establish connection')
responseParts = response.readline().split(':')
self.sessionID = responseParts[0]
self.heartbeatTimeout = int(responseParts[1])
self.connectionTimeout = int(responseParts[2])
self.supportedTransports = responseParts[3].split(',')
if 'websocket' not in self.supportedTransports:
raise SocketIOError('Could not parse handshake') # pragma: no cover
socketURL = '%s://%s/websocket/%s' % (
'wss' if self.secure else 'ws', baseURL, self.sessionID)
self.connection = websocket.create_connection(socketURL)
def _recv_packet(self): def __del__(self):
code, packetID, channelName, data = -1, None, None, None self.disconnect(close=False)
packet = self.connection.recv()
packetParts = packet.split(':', 3)
packetCount = len(packetParts)
if 4 == packetCount:
code, packetID, channelName, data = packetParts
elif 3 == packetCount:
code, packetID, channelName = packetParts
elif 1 == packetCount: # pragma: no cover
code = packetParts[0]
return int(code), packetID, channelName, data
def _send_packet(self, code, channelName='', data='', callback=None):
self.connection.send(':'.join([
str(code),
self.set_callback(callback) if callback else '',
channelName,
data]))
def disconnect(self, channelName=''):
self._send_packet(0, channelName)
if channelName:
del self.channelByName[channelName]
else:
self.__del__()
@property @property
def connected(self): def connected(self):
return self.connection.connected return self._socketIO.connected
def connect(self, channelName, Namespace=BaseNamespace): def disconnect(self, path='', close=True):
channel = Channel(self, channelName, Namespace) if self.connected:
self.channelByName[channelName] = channel self._socketIO.disconnect(path, close)
self._send_packet(1, channelName) if path:
return channel del self._namespaceByPath[path]
def _send_heartbeat(self):
try:
self._send_packet(2)
except:
self.__del__()
def message(self, messageData, callback=None, channelName=''):
if isinstance(messageData, basestring):
code = 3
data = messageData
else: else:
code = 4 self._rhythmicThread.cancel()
data = dumps(messageData) self._listenerThread.cancel()
self._send_packet(code, channelName, data, callback)
def emit(self, eventName, *eventArguments, **eventKeywords): def define(self, Namespace, path=''):
code = 5 if path:
if callable(eventArguments[-1]): self._socketIO.connect(path)
callback = eventArguments[-1] namespace = Namespace(self._socketIO, path)
eventArguments = eventArguments[:-1] self._namespaceByPath[path] = namespace
else: return namespace
callback = None
channelName = eventKeywords.get('channelName', '')
data = dumps(dict(name=eventName, args=eventArguments))
self._send_packet(code, channelName, data, callback)
def get_callback(self, channelName, eventName): def get_namespace(self, path=''):
'Get callback associated with channelName and eventName' return self._namespaceByPath[path]
socketIO = self.channelByName[channelName] if channelName else self
try:
return socketIO.callbackByEvent[eventName]
except KeyError:
pass
namespace = socketIO.namespace
def callback_(*eventArguments): def on(self, event, callback, path=''):
return namespace.on_(eventName, *eventArguments) return self.get_namespace(path).on(event, callback)
return getattr(namespace, name_callback(eventName), callback_)
def set_callback(self, callback): def message(self, data='', callback=None, path=''):
'Set callback that will be called after receiving an acknowledgment' self._socketIO.message(data, callback, path)
self.messageID += 1
self.namespaceThread.set_callback(self.messageID, callback)
return '%s+' % self.messageID
def on(self, eventName, callback): def emit(self, event, *args, **kw):
self.callbackByEvent[eventName] = callback self._socketIO.emit(event, *args, **kw)
def wait(self, seconds=None, forCallbacks=False): def wait(self, seconds=None):
if forCallbacks: if seconds:
self.namespaceThread.wait_for_callbacks(seconds) self._listenerThread.wait(seconds)
elif seconds:
sleep(seconds)
else: else:
try: try:
while self.connected: while self.connected:
@ -187,149 +160,281 @@ class SocketIO(object):
except KeyboardInterrupt: except KeyboardInterrupt:
pass pass
def wait_for_callbacks(self, seconds=None):
class Channel(object): self._listenerThread.wait_for_callbacks(seconds)
def __init__(self, socketIO, channelName, Namespace):
self.socketIO = socketIO
self.channelName = channelName
self.namespace = Namespace(self)
self.callbackByEvent = {}
def disconnect(self):
self.socketIO.disconnect(self.channelName)
def emit(self, eventName, *eventArguments):
self.socketIO.emit(eventName, *eventArguments,
channelName=self.channelName)
def message(self, messageData, callback=None):
self.socketIO.message(messageData, callback,
channelName=self.channelName)
def on(self, eventName, eventCallback):
self.callbackByEvent[eventName] = eventCallback
class ListenerThread(Thread): class _RhythmicThread(Thread):
'Process messages from SocketIO server' 'Execute call every few seconds'
daemon = True daemon = True
def __init__(self, socketIO): def __init__(self, intervalInSeconds, call, *args, **kw):
super(ListenerThread, self).__init__() super(_RhythmicThread, self).__init__()
self.socketIO = socketIO
self.done = Event()
self.waitingForCallbacks = Event()
self.callbackByMessageID = {}
self.get_callback = self.socketIO.get_callback
def run(self):
while not self.done.is_set():
try:
code, packetID, channelName, data = self.socketIO._recv_packet()
except:
continue
try:
delegate = {
0: self.on_disconnect,
1: self.on_connect,
2: self.on_heartbeat,
3: self.on_message,
4: self.on_json,
5: self.on_event,
6: self.on_acknowledgment,
7: self.on_error,
}[code]
except KeyError:
continue
delegate(packetID, channelName, data)
def cancel(self):
self.done.set()
def wait_for_callbacks(self, seconds):
self.waitingForCallbacks.set()
self.join(seconds)
def set_callback(self, messageID, callback):
self.callbackByMessageID[messageID] = callback
def on_disconnect(self, packetID, channelName, data):
callback = self.get_callback(channelName, 'disconnect')
callback()
def on_connect(self, packetID, channelName, data):
callback = self.get_callback(channelName, 'connect')
callback(self.socketIO)
def on_heartbeat(self, packetID, channelName, data):
pass
def on_message(self, packetID, channelName, data):
callback = self.get_callback(channelName, 'message')
callback(data)
def on_json(self, packetID, channelName, data):
callback = self.get_callback(channelName, 'message')
callback(loads(data))
def on_event(self, packetID, channelName, data):
valueByName = loads(data)
eventName = valueByName['name']
eventArguments = valueByName.get('args', [])
callback = self.get_callback(channelName, eventName)
callback(*eventArguments)
def on_acknowledgment(self, packetID, channelName, data):
dataParts = data.split('+', 1)
messageID = int(dataParts[0])
arguments = loads(dataParts[1]) or []
try:
callback = self.callbackByMessageID[messageID]
except KeyError:
pass
else:
del self.callbackByMessageID[messageID]
callback(*arguments)
callbackCount = len(self.callbackByMessageID)
if self.waitingForCallbacks.is_set() and not callbackCount:
self.cancel()
def on_error(self, packetID, channelName, data):
reason, advice = data.split('+', 1)
callback = self.get_callback(channelName, 'error')
callback(reason, advice)
class RhythmicThread(Thread):
'Execute rhythmicFunction every few seconds'
daemon = True
def __init__(self, intervalInSeconds, rhythmicFunction, *args, **kw):
super(RhythmicThread, self).__init__()
self.intervalInSeconds = intervalInSeconds self.intervalInSeconds = intervalInSeconds
self.rhythmicFunction = rhythmicFunction self.call = call
self.args = args self.args = args
self.kw = kw self.kw = kw
self.done = Event() self.done = Event()
def run(self): def run(self):
try: while not self.done.is_set():
while not self.done.is_set(): self.call(*self.args, **self.kw)
self.rhythmicFunction(*self.args, **self.kw) self.done.wait(self.intervalInSeconds)
self.done.wait(self.intervalInSeconds)
except:
pass
def cancel(self): def cancel(self):
self.done.set() self.done.set()
class _ListenerThread(Thread):
'Process messages from socket.io server'
daemon = True
def __init__(self, _socketIO, _namespaceByPath):
super(_ListenerThread, self).__init__()
self._socketIO = _socketIO
self._namespaceByPath = _namespaceByPath
self.done = Event()
self.ready = Event()
self.ready.set()
def cancel(self):
self.done.set()
def wait(self, seconds):
self.done.wait(seconds)
def wait_for_callbacks(self, seconds):
self.ready.clear()
self.ready.wait(seconds)
def get_ackCallback(self, packetID):
return lambda *args: self._socketIO.ack(packetID, *args)
def run(self):
while not self.done.is_set():
try:
code, packetID, path, data = self._socketIO.recv_packet()
except SocketIOConnectionError, error:
print error
return
except SocketIOPacketError, error:
print error
continue
try:
namespace = self._namespaceByPath[path]
except KeyError:
print 'Received unexpected path (%s)' % path
continue
try:
delegate = {
'0': self.on_disconnect,
'1': self.on_connect,
'2': self.on_heartbeat,
'3': self.on_message,
'4': self.on_json,
'5': self.on_event,
'6': self.on_ack,
'7': self.on_error,
}[code]
except KeyError:
print 'Received unexpected code (%s)' % code
continue
delegate(packetID, namespace._get_eventCallback, data)
def on_disconnect(self, packetID, get_eventCallback, data):
get_eventCallback('disconnect')()
def on_connect(self, packetID, get_eventCallback, data):
get_eventCallback('connect')()
def on_heartbeat(self, packetID, get_eventCallback, data):
pass
def on_message(self, packetID, get_eventCallback, data):
args = [data]
if packetID:
args.append(self.get_ackCallback(packetID))
get_eventCallback('message')(*args)
def on_json(self, packetID, get_eventCallback, data):
args = [loads(data)]
if packetID:
args.append(self.get_ackCallback(packetID))
get_eventCallback('message')(*args)
def on_event(self, packetID, get_eventCallback, data):
valueByName = loads(data)
event = valueByName['name']
args = valueByName.get('args', [])
if packetID:
args.append(self.get_ackCallback(packetID))
get_eventCallback(event)(*args)
def on_ack(self, packetID, get_eventCallback, data):
dataParts = data.split('+', 1)
messageID = int(dataParts[0])
args = loads(dataParts[1]) if len(dataParts) > 1 else []
callback = self._socketIO.get_messageCallback(messageID)
if not callback:
return
callback(*args)
if not self._socketIO.has_messageCallback:
self.ready.set()
def on_error(self, packetID, get_eventCallback, data):
reason, advice = data.split('+', 1)
get_eventCallback('error')(reason, advice)
class _SocketIO(object):
'Low-level interface to remove cyclic references in child threads'
messageID = 0
def __init__(self, host, port, secure, proxies):
baseURL = '%s:%d/socket.io/%s' % (host, port, PROTOCOL)
targetScheme = 'https' if secure else 'http'
targetURL = '%s://%s/' % (targetScheme, baseURL)
try:
response = urlopen(targetURL, proxies=proxies)
except IOError: # pragma: no cover
raise SocketIOError('Could not start connection')
if 200 != response.getcode(): # pragma: no cover
raise SocketIOError('Could not establish connection')
responseParts = response.readline().split(':')
sessionID = responseParts[0]
heartbeatTimeout = int(responseParts[1])
# connectionTimeout = int(responseParts[2])
supportedTransports = responseParts[3].split(',')
if 'websocket' not in supportedTransports:
raise SocketIOError('Could not parse handshake')
socketScheme = 'wss' if secure else 'ws'
socketURL = '%s://%s/websocket/%s' % (socketScheme, baseURL, sessionID)
self.connection = create_connection(socketURL)
self.heartbeatInterval = heartbeatTimeout - 2
self.callbackByMessageID = {}
def __del__(self):
self.disconnect(close=False)
def disconnect(self, path='', close=True):
if not self.connected:
return
if path:
self.send_packet(0, path)
elif close:
self.connection.close()
def connect(self, path):
self.send_packet(1, path)
def send_heartbeat(self):
try:
self.send_packet(2)
except SocketIOPacketError:
print 'Could not send heartbeat'
pass
def message(self, data, callback, path):
if isinstance(data, basestring):
code = 3
packetData = data
else:
code = 4
packetData = dumps(data, ensure_ascii=False)
self.send_packet(code, path, packetData, callback)
def emit(self, event, *args, **kw):
callback, args = find_callback(args, kw)
packetData = dumps(dict(name=event, args=args), ensure_ascii=False)
path = kw.get('path', '')
self.send_packet(5, path, packetData, callback)
def ack(self, packetID, *args):
packetID = packetID.rstrip('+')
packetData = '%s+%s' % (
packetID,
dumps(args, ensure_ascii=False),
) if args else packetID
self.send_packet(6, data=packetData)
def set_messageCallback(self, callback):
'Set callback that will be called after receiving an acknowledgment'
self.messageID += 1
self.callbackByMessageID[self.messageID] = callback
return '%s+' % self.messageID
def get_messageCallback(self, messageID):
try:
callback = self.callbackByMessageID[messageID]
del self.callbackByMessageID[messageID]
return callback
except KeyError:
return
@property
def has_messageCallback(self):
return True if self.callbackByMessageID else False
def recv_packet(self):
try:
packet = self.connection.recv()
except WebSocketConnectionClosedException:
text = 'Lost connection (Connection closed)'
raise SocketIOConnectionError(text)
except socket.timeout:
text = 'Lost connection (Connection timed out)'
raise SocketIOConnectionError(text)
except socket.error:
text = 'Lost connection'
raise SocketIOConnectionError(text)
try:
packetParts = packet.split(':', 3)
except AttributeError:
raise SocketIOPacketError('Received invalid packet (%s)' % packet)
packetCount = len(packetParts)
code, packetID, path, data = None, None, None, None
if 4 == packetCount:
code, packetID, path, data = packetParts
elif 3 == packetCount:
code, packetID, path = packetParts
elif 1 == packetCount:
code = packetParts[0]
return code, packetID, path, data
def send_packet(self, code, path='', data='', callback=None):
packetID = self.set_messageCallback(callback) if callback else ''
packetParts = [str(code), packetID, path, data]
try:
packet = ':'.join(packetParts)
self.connection.send(packet)
except socket.error:
raise SocketIOPacketError('Could not send packet')
@property
def connected(self):
return self.connection.connected
class SocketIOError(Exception): class SocketIOError(Exception):
pass pass
def name_callback(eventName): class SocketIOConnectionError(SocketIOError):
return 'on_' + eventName.replace(' ', '_') pass
class SocketIOPacketError(SocketIOError):
pass
def find_callback(args, kw=None):
'Return callback whether passed as a last argument or as a keyword'
if args and callable(args[-1]):
return args[-1], args[:-1]
try:
return kw['callback'], args
except (KeyError, TypeError):
return None, args

View file

@ -1,61 +1,174 @@
from socketIO_client import SocketIO, BaseNamespace from socketIO_client import SocketIO, BaseNamespace, find_callback
from time import sleep
from unittest import TestCase from unittest import TestCase
PAYLOAD = {'bbb': 'ccc'} HOST = 'localhost'
ON_RESPONSE_CALLED = False PORT = 8000
DATA = 'xxx'
PAYLOAD = {'xxx': 'yyy'}
class TestSocketIO(TestCase): class TestSocketIO(TestCase):
def setUp(self):
self.socketIO = SocketIO(HOST, PORT)
self.called_on_response = False
def tearDown(self):
del self.socketIO
def on_response(self, *args):
self.called_on_response = True
for arg in args:
if isinstance(arg, dict):
self.assertEqual(arg, PAYLOAD)
else:
self.assertEqual(arg, DATA)
def is_connected(self, socketIO, connected):
childThreads = [
socketIO._rhythmicThread,
socketIO._listenerThread,
]
for childThread in childThreads:
self.assertEqual(not connected, childThread.done.is_set())
self.assertEqual(connected, socketIO.connected)
def test_disconnect(self): def test_disconnect(self):
socketIO = SocketIO('localhost', 8000) 'Terminate child threads after disconnect'
socketIO.disconnect() self.is_connected(self.socketIO, True)
self.assertEqual(socketIO.connected, False) self.socketIO.disconnect()
self.is_connected(self.socketIO, False)
# Use context manager
with SocketIO(HOST, PORT) as self.socketIO:
self.is_connected(self.socketIO, True)
self.is_connected(self.socketIO, False)
def test_message(self):
'Message'
self.socketIO.define(Namespace)
self.socketIO.message()
self.socketIO.wait(0.1)
namespace = self.socketIO.get_namespace()
self.assertEqual(namespace.response, 'message_response')
def test_message_with_data(self):
'Message with data'
self.socketIO.define(Namespace)
self.socketIO.message(DATA)
self.socketIO.wait(0.1)
namespace = self.socketIO.get_namespace()
self.assertEqual(namespace.response, DATA)
def test_message_with_payload(self):
'Message with payload'
self.socketIO.define(Namespace)
self.socketIO.message(PAYLOAD)
self.socketIO.wait(0.1)
namespace = self.socketIO.get_namespace()
self.assertEqual(namespace.response, PAYLOAD)
def test_message_with_callback(self):
'Message with callback'
self.socketIO.message(callback=self.on_response)
self.socketIO.wait_for_callbacks(seconds=0.1)
self.assertEqual(self.called_on_response, True)
def test_message_with_callback_with_data(self):
'Message with callback with data'
self.socketIO.message(DATA, self.on_response)
self.socketIO.wait_for_callbacks(seconds=0.1)
self.assertEqual(self.called_on_response, True)
def test_emit(self): def test_emit(self):
socketIO = SocketIO('localhost', 8000, Namespace) 'Emit'
socketIO.emit('aaa', PAYLOAD) self.socketIO.define(Namespace)
sleep(0.5) self.socketIO.emit('emit')
self.assertEqual(socketIO.namespace.payload, PAYLOAD) self.socketIO.wait(0.1)
self.assertEqual(self.socketIO.get_namespace().argsByEvent, {
'emit_response': (),
})
def test_emit_with_payload(self):
'Emit with payload'
self.socketIO.define(Namespace)
self.socketIO.emit('emit_with_payload', PAYLOAD)
self.socketIO.wait(0.1)
self.assertEqual(self.socketIO.get_namespace().argsByEvent, {
'emit_with_payload_response': (PAYLOAD,),
})
def test_emit_with_multiple_payloads(self):
'Emit with multiple payloads'
self.socketIO.define(Namespace)
self.socketIO.emit('emit_with_multiple_payloads', PAYLOAD, PAYLOAD)
self.socketIO.wait(0.1)
self.assertEqual(self.socketIO.get_namespace().argsByEvent, {
'emit_with_multiple_payloads_response': (PAYLOAD, PAYLOAD),
})
def test_emit_with_callback(self): def test_emit_with_callback(self):
global ON_RESPONSE_CALLED 'Emit with callback'
ON_RESPONSE_CALLED = False self.socketIO.emit('emit_with_callback', self.on_response)
socketIO = SocketIO('localhost', 8000) self.socketIO.wait_for_callbacks(seconds=0.1)
socketIO.emit('aaa', PAYLOAD, on_response) self.assertEqual(self.called_on_response, True)
socketIO.wait(forCallbacks=True)
self.assertEqual(ON_RESPONSE_CALLED, True)
def test_events(self): def test_emit_with_callback_with_payload(self):
global ON_RESPONSE_CALLED 'Emit with callback with payload'
ON_RESPONSE_CALLED = False self.socketIO.emit('emit_with_callback_with_payload',
socketIO = SocketIO('localhost', 8000) self.on_response)
socketIO.on('ddd', on_response) self.socketIO.wait_for_callbacks(seconds=0.1)
socketIO.emit('aaa', PAYLOAD) self.assertEqual(self.called_on_response, True)
sleep(0.5)
self.assertEqual(ON_RESPONSE_CALLED, True)
def test_channels(self): def test_emit_with_callback_with_multiple_payloads(self):
mainSocket = SocketIO('localhost', 8000, Namespace) 'Emit with callback with multiple payloads'
chatSocket = mainSocket.connect('/chat', Namespace) self.socketIO.emit('emit_with_callback_with_multiple_payloads',
newsSocket = mainSocket.connect('/news', Namespace) self.on_response)
newsSocket.emit('aaa', PAYLOAD) self.socketIO.wait_for_callbacks(seconds=0.1)
sleep(0.5) self.assertEqual(self.called_on_response, True)
self.assertNotEqual(mainSocket.namespace.payload, PAYLOAD)
self.assertNotEqual(chatSocket.namespace.payload, PAYLOAD) def test_emit_with_event(self):
self.assertEqual(newsSocket.namespace.payload, PAYLOAD) 'Emit to trigger an event'
self.socketIO.on('emit_with_event_response', self.on_response)
self.socketIO.emit('emit_with_event', PAYLOAD)
self.socketIO.wait_for_callbacks(0.1)
self.assertEqual(self.called_on_response, True)
def test_ack(self):
'Trigger server callback'
self.socketIO.define(Namespace)
self.socketIO.emit('ack', PAYLOAD)
self.socketIO.wait(0.1)
self.assertEqual(self.socketIO.get_namespace().argsByEvent, {
'ack_response': (PAYLOAD,),
'ack_callback_response': (PAYLOAD,),
})
def test_namespaces(self):
'Behave differently in different namespaces'
mainNamespace = self.socketIO.define(Namespace)
chatNamespace = self.socketIO.define(Namespace, '/chat')
newsNamespace = self.socketIO.define(Namespace, '/news')
newsNamespace.emit('emit_with_payload', PAYLOAD)
self.socketIO.wait(0.1)
self.assertEqual(mainNamespace.argsByEvent, {})
self.assertEqual(chatNamespace.argsByEvent, {})
self.assertEqual(newsNamespace.argsByEvent, {
'emit_with_payload_response': (PAYLOAD,),
})
class Namespace(BaseNamespace): class Namespace(BaseNamespace):
payload = None def initialize(self):
self.response = None
self.argsByEvent = {}
def on_ddd(self, data): def on_message(self, data):
self.payload = data self.response = data
def on_event(self, event, *args):
def on_response(*args): callback, args = find_callback(args)
global ON_RESPONSE_CALLED if callback:
ON_RESPONSE_CALLED = True callback(*args)
self.argsByEvent[event] = args