Loading transfer_service/abort_job_amqp_server.py +2 −2 Original line number Diff line number Diff line Loading @@ -5,14 +5,14 @@ from config import Config class AbortJobAMQPServer(AMQPServer): def __init__(self, host, queue): def __init__(self, host, port, queue): self.type = "abort" config = Config("vos_ts.conf") self.params = config.loadSection("job_cache") self.jobCache = JobCache(self.params["host"], self.params.getint("port"), self.params.getint("db_read")) super(AbortJobAMQPServer, self).__init__(host, queue) super(AbortJobAMQPServer, self).__init__(host, port, queue) def execute_callback(self, requestBody): return 42 Loading transfer_service/amqp_server.py +3 −2 Original line number Diff line number Diff line Loading @@ -5,11 +5,12 @@ import json class AMQPServer(threading.Thread): def __init__(self, host, queue): def __init__(self, host, port, queue): threading.Thread.__init__(self) self.host = host self.port = port self.queue = queue self.connection = pika.BlockingConnection(pika.ConnectionParameters(host = self.host)) self.connection = pika.BlockingConnection(pika.ConnectionParameters(host = self.host, port = self.port)) self.channel = self.connection.channel(); self.channel.queue_declare(queue = self.queue) self.channel.basic_qos(prefetch_count = 1) Loading transfer_service/get_job_amqp_server.py +2 −2 Original line number Diff line number Diff line Loading @@ -5,14 +5,14 @@ from config import Config class GetJobAMQPServer(AMQPServer): def __init__(self, host, queue): def __init__(self, host, port, queue): self.type = "poll" config = Config("vos_ts.conf") self.params = config.loadSection("job_cache") self.jobCache = JobCache(self.params["host"], self.params.getint("port"), self.params.getint("db_read")) super(GetJobAMQPServer, self).__init__(host, queue) super(GetJobAMQPServer, self).__init__(host, port, queue) def execute_callback(self, requestBody): if "jobId" in requestBody: Loading transfer_service/job_handler.py +6 −5 Original line number Diff line number Diff line Loading @@ -9,19 +9,20 @@ from store_amqp_server import StoreAMQPServer class JobHandler(object): def __init__(self, host): def __init__(self, host, port): self.host = host self.port = port self.amqpServerList = [] def addAMQPServer(self, srvType, rpcQueue): if srvType == 'start': self.amqpServerList.append(StartJobAMQPServer(self.host, rpcQueue)) self.amqpServerList.append(StartJobAMQPServer(self.host, self.port, rpcQueue)) elif srvType == 'poll': self.amqpServerList.append(GetJobAMQPServer(self.host, rpcQueue)) self.amqpServerList.append(GetJobAMQPServer(self.host, self.port, rpcQueue)) elif srvType == 'abort': self.amqpServerList.append(AbortJobAMQPServer(self.host, rpcQueue)) self.amqpServerList.append(AbortJobAMQPServer(self.host, self.port, rpcQueue)) elif srvType == 'store': self.amqpServerList.append(StoreAMQPServer(self.host, rpcQueue)) self.amqpServerList.append(StoreAMQPServer(self.host, self.port, rpcQueue)) else: sys.exit(f"FATAL: unknown server type {srvType}.") Loading transfer_service/start_job_amqp_server.py +2 −2 Original line number Diff line number Diff line Loading @@ -9,14 +9,14 @@ from config import Config class StartJobAMQPServer(AMQPServer): def __init__(self, host, queue): def __init__(self, host, port, queue): self.type = "start" config = Config("vos_ts.conf") self.params = config.loadSection("job_cache") self.jobCache = JobCache(self.params["host"], self.params.getint("port"), self.params.getint("db_read")) super(StartJobAMQPServer, self).__init__(host, queue) super(StartJobAMQPServer, self).__init__(host, port, queue) def execute_callback(self, requestBody): self.job = Job() Loading Loading
transfer_service/abort_job_amqp_server.py +2 −2 Original line number Diff line number Diff line Loading @@ -5,14 +5,14 @@ from config import Config class AbortJobAMQPServer(AMQPServer): def __init__(self, host, queue): def __init__(self, host, port, queue): self.type = "abort" config = Config("vos_ts.conf") self.params = config.loadSection("job_cache") self.jobCache = JobCache(self.params["host"], self.params.getint("port"), self.params.getint("db_read")) super(AbortJobAMQPServer, self).__init__(host, queue) super(AbortJobAMQPServer, self).__init__(host, port, queue) def execute_callback(self, requestBody): return 42 Loading
transfer_service/amqp_server.py +3 −2 Original line number Diff line number Diff line Loading @@ -5,11 +5,12 @@ import json class AMQPServer(threading.Thread): def __init__(self, host, queue): def __init__(self, host, port, queue): threading.Thread.__init__(self) self.host = host self.port = port self.queue = queue self.connection = pika.BlockingConnection(pika.ConnectionParameters(host = self.host)) self.connection = pika.BlockingConnection(pika.ConnectionParameters(host = self.host, port = self.port)) self.channel = self.connection.channel(); self.channel.queue_declare(queue = self.queue) self.channel.basic_qos(prefetch_count = 1) Loading
transfer_service/get_job_amqp_server.py +2 −2 Original line number Diff line number Diff line Loading @@ -5,14 +5,14 @@ from config import Config class GetJobAMQPServer(AMQPServer): def __init__(self, host, queue): def __init__(self, host, port, queue): self.type = "poll" config = Config("vos_ts.conf") self.params = config.loadSection("job_cache") self.jobCache = JobCache(self.params["host"], self.params.getint("port"), self.params.getint("db_read")) super(GetJobAMQPServer, self).__init__(host, queue) super(GetJobAMQPServer, self).__init__(host, port, queue) def execute_callback(self, requestBody): if "jobId" in requestBody: Loading
transfer_service/job_handler.py +6 −5 Original line number Diff line number Diff line Loading @@ -9,19 +9,20 @@ from store_amqp_server import StoreAMQPServer class JobHandler(object): def __init__(self, host): def __init__(self, host, port): self.host = host self.port = port self.amqpServerList = [] def addAMQPServer(self, srvType, rpcQueue): if srvType == 'start': self.amqpServerList.append(StartJobAMQPServer(self.host, rpcQueue)) self.amqpServerList.append(StartJobAMQPServer(self.host, self.port, rpcQueue)) elif srvType == 'poll': self.amqpServerList.append(GetJobAMQPServer(self.host, rpcQueue)) self.amqpServerList.append(GetJobAMQPServer(self.host, self.port, rpcQueue)) elif srvType == 'abort': self.amqpServerList.append(AbortJobAMQPServer(self.host, rpcQueue)) self.amqpServerList.append(AbortJobAMQPServer(self.host, self.port, rpcQueue)) elif srvType == 'store': self.amqpServerList.append(StoreAMQPServer(self.host, rpcQueue)) self.amqpServerList.append(StoreAMQPServer(self.host, self.port, rpcQueue)) else: sys.exit(f"FATAL: unknown server type {srvType}.") Loading
transfer_service/start_job_amqp_server.py +2 −2 Original line number Diff line number Diff line Loading @@ -9,14 +9,14 @@ from config import Config class StartJobAMQPServer(AMQPServer): def __init__(self, host, queue): def __init__(self, host, port, queue): self.type = "start" config = Config("vos_ts.conf") self.params = config.loadSection("job_cache") self.jobCache = JobCache(self.params["host"], self.params.getint("port"), self.params.getint("db_read")) super(StartJobAMQPServer, self).__init__(host, queue) super(StartJobAMQPServer, self).__init__(host, port, queue) def execute_callback(self, requestBody): self.job = Job() Loading