Rev pulsar client code.

Updates Pulsar client through https://github.com/galaxyproject/pulsar/commit/f1a9e83674dd98affbed476ba36b3ad07fa5e88b. The most important and substantial change here is a fix for galaxyproject/pulsar#63 provided by @dctrud.
This commit is contained in:
John Chilton
2015-04-10 12:21:44 -04:00
parent 800e83fef6
commit aad36d0776
3 changed files with 8 additions and 4 deletions
+1 -1
View File
@@ -83,7 +83,7 @@ class PulsarExchange(object):
connection.drain_events(timeout=self.__timeout)
except socket.timeout:
pass
except (IOError, socket.error), exc:
except (IOError, socket.error) as exc:
self.__handle_io_error(exc, heartbeat_thread)
except BaseException:
log.exception("Problem consuming queue, consumer quitting in problematic fashion!")
+6 -2
View File
@@ -1,5 +1,6 @@
import os
from json import dumps
from json import loads
from .destination import submit_params
from .setup_handler import build as build_setup_handler
@@ -170,10 +171,13 @@ class JobClient(BaseJobClient):
input_path = path
if contents:
input_path = None
if action_type == 'transfer':
# action type == 'message' should either copy or transfer
# depending on default not just fallback to transfer.
if action_type in ['transfer', 'message']:
return self._upload_file(args, contents, input_path)
elif action_type == 'copy':
pulsar_path = self._raw_execute('path', args)
path_response = self._raw_execute('path', args)
pulsar_path = loads(path_response)['path']
copy(path, pulsar_path)
return {'path': pulsar_path}
+1 -1
View File
@@ -348,7 +348,7 @@ class TransferTracker(object):
if action.staging_needed:
local_action = action.staging_action_local
if local_action:
response = self.client.put_file(path, type, name=name, contents=contents)
response = self.client.put_file(path, type, name=name, contents=contents, action_type=action.action_type)
def get_path():
return response['path']