Handling partial metadata requests

This commit is contained in:
James Barnsley
2018-03-19 16:09:27 +13:00
parent f615c6c8d0
commit 2cde9fa786
6 changed files with 79 additions and 94 deletions

View File

@ -42,6 +42,30 @@ class IrisCore(object):
snapcast_listener = False
##
# Create a new snapcast TCP connection
#
# @return socket
##
def new_snapcast_socket(self):
if not self.config['iris'].get('snapcast_enabled'):
logger.error("Iris Snapcast not enabled")
raise Exception("Snapcast not enabled")
try:
snapcast = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
snapcast.connect((self.config['iris']['snapcast_host'], self.config['iris']['snapcast_port']))
except socket.gaierror, e:
logger.error("Iris could not connect to Snapcast: %s" % e)
raise Exception(e);
except socket.error, e:
logger.error("Iris could not connect to Snapcast: %s" % e)
raise Exception(e);
return snapcast
##
# Create our ongoing notification listener
#
@ -56,30 +80,13 @@ class IrisCore(object):
logger.error("Iris could not connect to Snapcast: %s" % e)
##
# Create a new snapcast TCP connection
#
# @return socket
##
def new_snapcast_socket(self):
try:
snapcast = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
snapcast.connect((self.config['iris']['snapcast_host'], self.config['iris']['snapcast_port']))
except socket.gaierror, e:
logger.error("Iris could not connect to Snapcast: %s" % e)
raise Exception(e);
except socket.error, e:
logger.error("Iris could not connect to Snapcast: %s" % e)
raise Exception(e);
return snapcast
##
# Disconnect our Snapcast listener
##
def snapcast_disconnect(self):
if (self.snapcast_listener):
def snapcast_disconnect_listener(self):
if self.snapcast_listener:
self.snapcast_listener.close()
self.snapcast_listener = None
##
@ -90,7 +97,7 @@ class IrisCore(object):
# @param data = string
##
def snapcast_handle_message(self, data):
logger.info("Iris received Snapcast message: "+data)
logger.debug("Iris received Snapcast message: "+data)
try:
data = json.loads(data)
@ -126,11 +133,10 @@ class IrisCore(object):
# Create our connection
try:
snapcast_socket = self.new_snapcast_socket()
except e:
callback(error={
'status': 0,
except Exception, e:
callback(response=None, error={
'message': "Could not connect to Snapcast",
'data': e
'data': str(e)
})
return
@ -139,9 +145,9 @@ class IrisCore(object):
snapcast_socket.send(data.encode('ascii')+b"\n")
except socket.error, e:
logger.error("Iris could not send request to Snapcast: %s" % e)
callback(error={
'status': 0,
'message': "Failed to send request to Snapcast"
callback(response=None, error={
'message': "Failed to send request to Snapcast",
'data': str(e)
})
return
@ -151,9 +157,9 @@ class IrisCore(object):
response = snapcast_socket.recv(8192)
except socket.error, e:
logger.error("Iris failed to receive Snapcast response: %s" % e)
callback(error={
'status': 0,
'message': "Failed to receive Snapcast response"
callback(response=None, error={
'message': "Failed to receive Snapcast response",
'data': str(e)
})
snapcast_socket.close()
@ -166,8 +172,7 @@ class IrisCore(object):
callback(response=response)
except:
logger.error("Iris received malformed Snapcast response: "+response)
callback(error={
'status': 0,
callback(response=None, error={
'message': "Malformed Snapcast response",
'data': response
})
@ -240,30 +245,31 @@ class IrisCore(object):
def broadcast(self, *args, **kwargs):
callback = kwargs.get('callback', None)
message = kwargs.get('message', None)
data = kwargs.get('data', None)
if 'jsonrpc' not in message:
message['jsonrpc'] = '2.0'
if 'jsonrpc' not in data:
data['jsonrpc'] = '2.0'
for connection in self.connections.itervalues():
send_to_this_connection = True
# Don't send the broadcast to the origin, naturally
if 'connection_id' in message:
if connection['connection_id'] == message["connection_id"]:
if 'connection_id' in data:
if connection['connection_id'] == data["connection_id"]:
send_to_this_connection = False
if send_to_this_connection:
connection['connection'].write_message(json_encode(message))
connection['connection'].write_message(json_encode(data))
response = {
'result': 'Broadcast to '+str(len(self.connections))+' connections'
'message': 'Broadcast to '+str(len(self.connections))+' connections'
}
if (callback):
callback(response)
else:
return response
return response
##
# Connections
@ -315,7 +321,7 @@ class IrisCore(object):
}
)
self.broadcast(message={
self.broadcast(data={
'method': 'pusher_connection_added',
'params': {
'connection': client
@ -327,7 +333,7 @@ class IrisCore(object):
try:
client = self.connections[connection_id]['client']
del self.connections[connection_id]
self.broadcast(message={
self.broadcast(data={
'method': "pusher_connection_removed",
'params': {
'connection': client
@ -343,7 +349,7 @@ class IrisCore(object):
if connection_id in self.connections:
self.connections[connection_id]['client']['username'] = data['username']
self.broadcast(message={
self.broadcast(data={
'method': "pusher_connection_updated",
'params': {
'connection': self.connections[connection_id]['client']
@ -555,14 +561,14 @@ class IrisCore(object):
if starting:
self.core.playback.play()
self.broadcast(message={
self.broadcast(data={
'method': "radio_started",
'params': {
'radio': self.radio
}
})
else:
self.broadcast(message={
self.broadcast(data={
'method': "radio_changed",
'params': {
'radio': self.radio
@ -604,7 +610,7 @@ class IrisCore(object):
self.core.tracklist.set_consume(self.initial_consume)
self.core.playback.stop()
self.broadcast(message={
self.broadcast(data={
'method': "radio_stopped",
'params': {
'radio': self.radio
@ -701,12 +707,12 @@ class IrisCore(object):
for tlid in data['tlids']:
item = {
'tlid': tlid,
'added_from': data['added_from'],
'added_by': data['added_by']
'added_from': data['added_from'] if 'added_from' in data else None,
'added_by': data['added_by'] if 'added_by' in data else None
}
self.queue_metadata['tlid_'+str(tlid)] = item
self.broadcast(message={
self.broadcast(data={
'method': 'queue_metadata_changed',
'params': {
'queue_metadata': self.queue_metadata
@ -779,7 +785,7 @@ class IrisCore(object):
token['expires_at'] = time.time() + token['expires_in']
self.spotify_token = token
self.broadcast(message={
self.broadcast(data={
'method': 'spotify_token_changed',
'params': {
'spotify_token': self.spotify_token

View File

@ -19,12 +19,12 @@ class IrisFrontend(pykka.ThreadingActor, CoreListener):
def on_start(self):
logger.info('Starting Iris '+mem.iris.version)
#if mem.iris.config['iris']['snapcast_enabed']:
# Create our listening socket for Snapcast notifications
mem.iris.create_snapcast_listener()
# Create our listening socket for Snapcast event notifications (only if enabled)
if mem.iris.config['iris'].get('snapcast_enabled'):
mem.iris.create_snapcast_listener()
def on_stop(self):
mem.iris.snapcast_disconnect()
mem.iris.snapcast_disconnect_listener()
def track_playback_ended(self, tl_track, time_position):
mem.iris.check_for_radio_update()

View File

@ -62,8 +62,7 @@ class WebsocketHandler(tornado.websocket.WebSocketHandler):
def on_message(self, message):
print message
logger.debug("Iris websocket message received: "+message)
message = json_decode(message)
@ -73,7 +72,7 @@ class WebsocketHandler(tornado.websocket.WebSocketHandler):
id = None
if 'jsonrpc' not in message:
self.handle_response(id=id, error={'id': id, 'code': 32602, 'message': 'Invalid JSON-RPC request (missing property "jsonrpc")'})
self.handle_result(id=id, error={'id': id, 'code': 32602, 'message': 'Invalid JSON-RPC request (missing property "jsonrpc")'})
if 'params' in message:
params = message['params']
@ -90,13 +89,13 @@ class WebsocketHandler(tornado.websocket.WebSocketHandler):
# make sure the method exists
if hasattr(mem.iris, message['method']):
getattr(mem.iris, message['method'])(data=params, callback=lambda response, error=False: self.handle_response(id=id, method=message['method'], response=response, error=error))
getattr(mem.iris, message['method'])(data=params, callback=lambda response, error=False: self.handle_result(id=id, method=message['method'], response=response, error=error))
else:
self.handle_response(error={'id': id, 'code': 32601, 'message': 'Method "'+message['method']+'" does not exist'}, id=id)
self.handle_result(error={'id': id, 'code': 32601, 'message': 'Method "'+message['method']+'" does not exist'}, id=id)
return
else:
self.handle_response(error={'id': id, 'code': 32602, 'message': 'Method key missing from request'}, id=id)
self.handle_result(error={'id': id, 'code': 32602, 'message': 'Method key missing from request'}, id=id)
return
@ -107,7 +106,7 @@ class WebsocketHandler(tornado.websocket.WebSocketHandler):
# Handle a response from our core
# This is just our callback from an Async request
##
def handle_response(self, *args, **kwargs):
def handle_result(self, *args, **kwargs):
id = kwargs.get('id', False)
method = kwargs.get('method', None)
response = kwargs.get('response', None)
@ -162,45 +161,45 @@ class HttpHandler(tornado.web.RequestHandler):
@tornado.web.asynchronous
def get(self, slug=None):
id = time.time()
id = int(time.time())
# make sure the method exists
if hasattr(mem.iris, slug):
getattr(mem.iris, slug)(request=self.request, callback=lambda response, error=False: self.handle_response(id=id, method=slug, response=response, error=error))
getattr(mem.iris, slug)(request=self.request, callback=lambda response, error=False: self.handle_result(id=id, method=slug, response=response, error=error))
else:
self.handle_response(id=self.request_id, error={'code': 32601, 'message': "Method "+slug+" does not exist"})
self.handle_result(id=self.request_id, error={'code': 32601, 'message': "Method "+slug+" does not exist"})
return
@tornado.web.asynchronous
def post(self, slug=None):
id = time.time()
id = int(time.time())
try:
params = json.loads(self.request.body.decode('utf-8'))
except:
self.handle_response(id=id, error={'code': 32700, 'message': "Missing or invalid payload"})
self.handle_result(id=id, error={'code': 32700, 'message': "Missing or invalid payload"})
return
# make sure the method exists
if hasattr(mem.iris, slug):
try:
getattr(mem.iris, slug)(data=params, request=self.request, callback=lambda response=False, error=False: self.handle_response(id=id, method=slug, response=response, error=error))
getattr(mem.iris, slug)(data=params, request=self.request, callback=lambda response=False, error=False: self.handle_result(id=id, method=slug, response=response, error=error))
except urllib2.HTTPError as e:
self.handle_response(id=id, error={'code': 32601, 'message': "Invalid JSON payload"})
self.handle_result(id=id, error={'code': 32601, 'message': "Invalid JSON payload"})
return
else:
self.handle_response(id=id, error={'code': 32601, 'message': "Method "+slug+" does not exist"})
self.handle_result(id=id, error={'code': 32601, 'message': "Method "+slug+" does not exist"})
return
##
# Handle a response from our core
# This is just our callback from an Async request
##
def handle_response(self, *args, **kwargs):
def handle_result(self, *args, **kwargs):
id = kwargs.get('id', None)
method = kwargs.get('method', None)
response = kwargs.get('response', None)
@ -213,6 +212,7 @@ class HttpHandler(tornado.web.RequestHandler):
if error:
request_response['error'] = error
self.set_status(400)
# Log error with Sentry
#mem.iris.raven_client.captureMessage(data.message)
@ -240,6 +240,7 @@ class HttpHandler(tornado.web.RequestHandler):
# Write our response
self.write(request_response)
self.finish()

View File

@ -53209,11 +53209,6 @@ var PusherMiddleware = function () {
case 'PUSHER_START_UPGRADE':
_reactGa2.default.event({ category: 'Pusher', action: 'Upgrade', label: '' });
request(store, 'upgrade').then(function (response) {
if (response.error) {
console.error(response.error);
return false;
}
if (response.upgrade_successful) {
store.dispatch(uiActions.createNotification({ content: 'Upgrade complete' }));
} else {
@ -53232,10 +53227,6 @@ var PusherMiddleware = function () {
request(store, 'set_username', {
username: action.username
}).then(function (response) {
if (response.error) {
console.error(response.error);
return false;
}
response.type = 'PUSHER_USERNAME_CHANGED';
store.dispatch(response);
}, function (error) {
@ -53256,9 +53247,7 @@ var PusherMiddleware = function () {
break;
case 'PUSHER_GET_CONFIG':
console.log(action);
request(store, 'get_config').then(function (response) {
console.log(response);
store.dispatch({
type: 'PUSHER_CONFIG',
config: response.config

File diff suppressed because one or more lines are too long

View File

@ -233,11 +233,6 @@ const PusherMiddleware = (function(){
request(store, 'upgrade')
.then(
response => {
if (response.error){
console.error(response.error)
return false
}
if (response.upgrade_successful){
store.dispatch(uiActions.createNotification({content: 'Upgrade complete'}));
} else {
@ -263,10 +258,6 @@ const PusherMiddleware = (function(){
})
.then(
response => {
if (response.error){
console.error(response.error)
return false
}
response.type = 'PUSHER_USERNAME_CHANGED'
store.dispatch(response)
},
@ -299,11 +290,9 @@ const PusherMiddleware = (function(){
break;
case 'PUSHER_GET_CONFIG':
console.log(action)
request(store, 'get_config')
.then(
response => {
console.log(response)
store.dispatch({
type: 'PUSHER_CONFIG',
config: response.config