Loading transfer_service/abort_job_amqp_server.py +15 −5 Original line number Diff line number Diff line from amqp_server import AMQPServer from job_cache import JobCache from db_connector import DbConnector #from job_cache import JobCache from config import Config Loading @@ -8,13 +9,22 @@ class AbortJobAMQPServer(AMQPServer): 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 = config.loadSection("job_cache") #self.jobCache = JobCache(self.params["host"], # self.params.getint("port"), # self.params.getint("db_read")) self.params = config.loadSection("file_catalog") self.dbConn = DbConnector(self.params["user"], self.params["password"], self.params["host"], self.params.getint("port"), self.params.getint("db_read")) self.params["db"]) super(AbortJobAMQPServer, self).__init__(host, port, queue) def execute_callback(self, requestBody): self.dbConn.connect() # do something... self.dbConn.disconnect() return 42 def run(self): Loading transfer_service/get_job_amqp_server.py +20 −8 Original line number Diff line number Diff line from amqp_server import AMQPServer from job_cache import JobCache from db_connector import DbConnector #from job_cache import JobCache from config import Config Loading @@ -8,17 +9,28 @@ class GetJobAMQPServer(AMQPServer): 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 = config.loadSection("job_cache") #self.jobCache = JobCache(self.params["host"], # self.params.getint("port"), # self.params.getint("db_read")) self.params = config.loadSection("file_catalog") self.dbConn = DbConnector(self.params["user"], self.params["password"], self.params["host"], self.params.getint("port"), self.params.getint("db_read")) self.params["db"]) super(GetJobAMQPServer, self).__init__(host, port, queue) def execute_callback(self, requestBody): if "jobId" in requestBody: redis_res = self.jobCache.get(requestBody["jobId"]) print(f"Redis response: {redis_res}") return redis_res self.dbConn.connect() #redis_res = self.jobCache.get(requestBody["jobId"]) dbResponse = self.dbConn.getJob(requestBody["jobId"]) self.dbConn.disconnect() #print(f"Redis response: {redis_res}") print(f"Db response: {dbResponse}") #return redis_res return dbResponse else: return 42 Loading transfer_service/start_job_amqp_server.py +28 −10 Original line number Diff line number Diff line Loading @@ -2,8 +2,9 @@ import threading # only for testing purposes import time # only for testing purposes from amqp_server import AMQPServer from db_connector import DbConnector from job import Job from job_cache import JobCache #from job_cache import JobCache from config import Config Loading @@ -12,22 +13,36 @@ class StartJobAMQPServer(AMQPServer): 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 = config.loadSection("job_cache") #self.jobCache = JobCache(self.params["host"], # self.params.getint("port"), # self.params.getint("db_read")) self.params = config.loadSection("file_catalog") self.dbConn = DbConnector(self.params["user"], self.params["password"], self.params["host"], self.params.getint("port"), self.params.getint("db_read")) self.params["db"]) super(StartJobAMQPServer, self).__init__(host, port, queue) def execute_callback(self, requestBody): self.job = Job() self.job.setType("pullToVoSpace") self.job.setInfo(requestBody) self.job.setPhase("PENDING") self.jobCache.set(self.job) redis_res = self.jobCache.get(self.job.jobId) print(f"Redis response: {redis_res}") self.dbConn.connect() self.job.setOwnerId("0000") #self.jobCache.set(self.job) #redis_res = self.jobCache.get(self.job.jobId) self.dbConn.insertJob(self.job) dbResponse = self.dbConn.getJob(self.job.jobId) self.dbConn.disconnect() #print(f"Redis response: {redis_res}") print(f"Db response: {dbResponse}") t = threading.Thread(target = self.fake_job) # only for testing purposes t.start() # only for testing purposes return redis_res #return redis_res return dbResponse def run(self): print(f"Starting AMQP server of type {self.type}...") Loading @@ -38,5 +53,8 @@ class StartJobAMQPServer(AMQPServer): time.sleep(10) print("fake_job: changing job state...") self.job.setPhase("EXECUTING") self.jobCache.set(self.job) self.dbConn.connect() #self.jobCache.set(self.job) self.dbConn.insertJob(self.job) self.dbConn.disconnect() print("fake_job: state changed!") No newline at end of file Loading
transfer_service/abort_job_amqp_server.py +15 −5 Original line number Diff line number Diff line from amqp_server import AMQPServer from job_cache import JobCache from db_connector import DbConnector #from job_cache import JobCache from config import Config Loading @@ -8,13 +9,22 @@ class AbortJobAMQPServer(AMQPServer): 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 = config.loadSection("job_cache") #self.jobCache = JobCache(self.params["host"], # self.params.getint("port"), # self.params.getint("db_read")) self.params = config.loadSection("file_catalog") self.dbConn = DbConnector(self.params["user"], self.params["password"], self.params["host"], self.params.getint("port"), self.params.getint("db_read")) self.params["db"]) super(AbortJobAMQPServer, self).__init__(host, port, queue) def execute_callback(self, requestBody): self.dbConn.connect() # do something... self.dbConn.disconnect() return 42 def run(self): Loading
transfer_service/get_job_amqp_server.py +20 −8 Original line number Diff line number Diff line from amqp_server import AMQPServer from job_cache import JobCache from db_connector import DbConnector #from job_cache import JobCache from config import Config Loading @@ -8,17 +9,28 @@ class GetJobAMQPServer(AMQPServer): 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 = config.loadSection("job_cache") #self.jobCache = JobCache(self.params["host"], # self.params.getint("port"), # self.params.getint("db_read")) self.params = config.loadSection("file_catalog") self.dbConn = DbConnector(self.params["user"], self.params["password"], self.params["host"], self.params.getint("port"), self.params.getint("db_read")) self.params["db"]) super(GetJobAMQPServer, self).__init__(host, port, queue) def execute_callback(self, requestBody): if "jobId" in requestBody: redis_res = self.jobCache.get(requestBody["jobId"]) print(f"Redis response: {redis_res}") return redis_res self.dbConn.connect() #redis_res = self.jobCache.get(requestBody["jobId"]) dbResponse = self.dbConn.getJob(requestBody["jobId"]) self.dbConn.disconnect() #print(f"Redis response: {redis_res}") print(f"Db response: {dbResponse}") #return redis_res return dbResponse else: return 42 Loading
transfer_service/start_job_amqp_server.py +28 −10 Original line number Diff line number Diff line Loading @@ -2,8 +2,9 @@ import threading # only for testing purposes import time # only for testing purposes from amqp_server import AMQPServer from db_connector import DbConnector from job import Job from job_cache import JobCache #from job_cache import JobCache from config import Config Loading @@ -12,22 +13,36 @@ class StartJobAMQPServer(AMQPServer): 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 = config.loadSection("job_cache") #self.jobCache = JobCache(self.params["host"], # self.params.getint("port"), # self.params.getint("db_read")) self.params = config.loadSection("file_catalog") self.dbConn = DbConnector(self.params["user"], self.params["password"], self.params["host"], self.params.getint("port"), self.params.getint("db_read")) self.params["db"]) super(StartJobAMQPServer, self).__init__(host, port, queue) def execute_callback(self, requestBody): self.job = Job() self.job.setType("pullToVoSpace") self.job.setInfo(requestBody) self.job.setPhase("PENDING") self.jobCache.set(self.job) redis_res = self.jobCache.get(self.job.jobId) print(f"Redis response: {redis_res}") self.dbConn.connect() self.job.setOwnerId("0000") #self.jobCache.set(self.job) #redis_res = self.jobCache.get(self.job.jobId) self.dbConn.insertJob(self.job) dbResponse = self.dbConn.getJob(self.job.jobId) self.dbConn.disconnect() #print(f"Redis response: {redis_res}") print(f"Db response: {dbResponse}") t = threading.Thread(target = self.fake_job) # only for testing purposes t.start() # only for testing purposes return redis_res #return redis_res return dbResponse def run(self): print(f"Starting AMQP server of type {self.type}...") Loading @@ -38,5 +53,8 @@ class StartJobAMQPServer(AMQPServer): time.sleep(10) print("fake_job: changing job state...") self.job.setPhase("EXECUTING") self.jobCache.set(self.job) self.dbConn.connect() #self.jobCache.set(self.job) self.dbConn.insertJob(self.job) self.dbConn.disconnect() print("fake_job: state changed!") No newline at end of file