Loading transfer_service/job_handler.py +2 −1 Original line number Diff line number Diff line Loading @@ -6,6 +6,7 @@ from get_job_amqp_server import GetJobAMQPServer from abort_job_amqp_server import AbortJobAMQPServer from store_amqp_server import StoreAMQPServer class JobHandler(object): def __init__(self, host): Loading @@ -24,6 +25,6 @@ class JobHandler(object): else: sys.exit(f"FATAL: unknown server type {srvType}.") def run(self): def start(self): for srv in self.amqpServerList: srv.start() transfer_service/preprocessor.py +8 −18 Original line number Diff line number Diff line import threading import time from multiprocessing import Process from config import Config from job_queue import JobQueue from store_preprocessor import StorePreprocessor class Preprocessor(threading.Thread): class Preprocessor(Process): def __init__(self): threading.Thread.__init__(self) config = Config("vos_ts.conf") self.params = config.loadSection("scheduling") self.maxPendingJobs = self.params.getint("max_pending_jobs") self.maxReadyJobs = self.params.getint("max_ready_jobs") self.pendingQueue = JobQueue("pending") self.readyQueue = JobQueue("ready") self.storePreprocessor = StorePreprocessor() super(Preprocessor, self).__init__() def run(self): while True: if(self.readyQueue.len() <= self.maxReadyJobs and self.pendingQueue.len() > 0): jobObj = self.pendingQueue.getJob() self.storePreprocessor.prepare(jobObj["jobInfo"]["userName"]) self.storePreprocessor.start() self.pendingQueue.moveJobTo("ready") print("Job MOVED:") time.sleep(1) # Test #p = Preprocessor() #p.start() """ This method must be implemented by inherited classes """ pass No newline at end of file transfer_service/store_preprocessor.py +17 −4 Original line number Diff line number Diff line Loading @@ -6,16 +6,18 @@ import os import shutil import sys import time from datetime import datetime as dt from checksum import Checksum from file_grouper import FileGrouper from db_connector import DbConnector from node import Node from preprocessor import Preprocessor from config import Config class StorePreprocessor(object): class StorePreprocessor(Preprocessor): def __init__(self): self.md5calc = Checksum() Loading @@ -30,6 +32,7 @@ class StorePreprocessor(object): self.params.getint("port"), self.params["db"]) self.dbConn.connect() super(StorePreprocessor, self).__init__() # Scan is performed only on the first level! def scan(self): Loading Loading @@ -73,7 +76,7 @@ class StorePreprocessor(object): os.chown(os.path.join(folder, f), 0, 0) os.chmod(os.path.join(folder, f), 0o555) def start(self): def execute(self): # First scan to find crowded dirs [ dirs, files ] = self.scan() Loading Loading @@ -144,7 +147,17 @@ class StorePreprocessor(object): self.dbConn.disconnect() def run(self): while True: if(self.readyQueue.len() <= self.maxReadyJobs and self.pendingQueue.len() > 0): jobObj = self.pendingQueue.getJob() self.prepare(jobObj["jobInfo"]["userName"]) self.execute() self.pendingQueue.moveJobTo("ready") print("Job MOVED:") time.sleep(1) # Test #sp = StorePreprocessor() #sp.prepare("curban") #sp.start() No newline at end of file #sp.execute() No newline at end of file transfer_service/transfer_service.py +7 −10 Original line number Diff line number Diff line Loading @@ -2,7 +2,7 @@ import time import os from job_handler import JobHandler from preprocessor import Preprocessor from store_preprocessor import StorePreprocessor from config import Config Loading @@ -12,7 +12,7 @@ class TransferService(object): config = Config("vos_ts.conf") self.params = config.loadSection("amqp") self.jobHandler = JobHandler(self.params["host"]) self.preprocessor = Preprocessor() self.storePreprocessor = StorePreprocessor() # PullFromVOSpace self.jobHandler.addAMQPServer('start', 'start_job_queue') Loading @@ -22,13 +22,10 @@ class TransferService(object): # Push self.jobHandler.addAMQPServer('store', 'store_job_queue') # Preprocessor self.preprocessor.start() def run(self): self.jobHandler.run() def start(self): self.storePreprocessor.start() self.jobHandler.start() ts = TransferService() ts.run() time.sleep(3) ts.start() print("Transfer service is RUNNING...") No newline at end of file Loading
transfer_service/job_handler.py +2 −1 Original line number Diff line number Diff line Loading @@ -6,6 +6,7 @@ from get_job_amqp_server import GetJobAMQPServer from abort_job_amqp_server import AbortJobAMQPServer from store_amqp_server import StoreAMQPServer class JobHandler(object): def __init__(self, host): Loading @@ -24,6 +25,6 @@ class JobHandler(object): else: sys.exit(f"FATAL: unknown server type {srvType}.") def run(self): def start(self): for srv in self.amqpServerList: srv.start()
transfer_service/preprocessor.py +8 −18 Original line number Diff line number Diff line import threading import time from multiprocessing import Process from config import Config from job_queue import JobQueue from store_preprocessor import StorePreprocessor class Preprocessor(threading.Thread): class Preprocessor(Process): def __init__(self): threading.Thread.__init__(self) config = Config("vos_ts.conf") self.params = config.loadSection("scheduling") self.maxPendingJobs = self.params.getint("max_pending_jobs") self.maxReadyJobs = self.params.getint("max_ready_jobs") self.pendingQueue = JobQueue("pending") self.readyQueue = JobQueue("ready") self.storePreprocessor = StorePreprocessor() super(Preprocessor, self).__init__() def run(self): while True: if(self.readyQueue.len() <= self.maxReadyJobs and self.pendingQueue.len() > 0): jobObj = self.pendingQueue.getJob() self.storePreprocessor.prepare(jobObj["jobInfo"]["userName"]) self.storePreprocessor.start() self.pendingQueue.moveJobTo("ready") print("Job MOVED:") time.sleep(1) # Test #p = Preprocessor() #p.start() """ This method must be implemented by inherited classes """ pass No newline at end of file
transfer_service/store_preprocessor.py +17 −4 Original line number Diff line number Diff line Loading @@ -6,16 +6,18 @@ import os import shutil import sys import time from datetime import datetime as dt from checksum import Checksum from file_grouper import FileGrouper from db_connector import DbConnector from node import Node from preprocessor import Preprocessor from config import Config class StorePreprocessor(object): class StorePreprocessor(Preprocessor): def __init__(self): self.md5calc = Checksum() Loading @@ -30,6 +32,7 @@ class StorePreprocessor(object): self.params.getint("port"), self.params["db"]) self.dbConn.connect() super(StorePreprocessor, self).__init__() # Scan is performed only on the first level! def scan(self): Loading Loading @@ -73,7 +76,7 @@ class StorePreprocessor(object): os.chown(os.path.join(folder, f), 0, 0) os.chmod(os.path.join(folder, f), 0o555) def start(self): def execute(self): # First scan to find crowded dirs [ dirs, files ] = self.scan() Loading Loading @@ -144,7 +147,17 @@ class StorePreprocessor(object): self.dbConn.disconnect() def run(self): while True: if(self.readyQueue.len() <= self.maxReadyJobs and self.pendingQueue.len() > 0): jobObj = self.pendingQueue.getJob() self.prepare(jobObj["jobInfo"]["userName"]) self.execute() self.pendingQueue.moveJobTo("ready") print("Job MOVED:") time.sleep(1) # Test #sp = StorePreprocessor() #sp.prepare("curban") #sp.start() No newline at end of file #sp.execute() No newline at end of file
transfer_service/transfer_service.py +7 −10 Original line number Diff line number Diff line Loading @@ -2,7 +2,7 @@ import time import os from job_handler import JobHandler from preprocessor import Preprocessor from store_preprocessor import StorePreprocessor from config import Config Loading @@ -12,7 +12,7 @@ class TransferService(object): config = Config("vos_ts.conf") self.params = config.loadSection("amqp") self.jobHandler = JobHandler(self.params["host"]) self.preprocessor = Preprocessor() self.storePreprocessor = StorePreprocessor() # PullFromVOSpace self.jobHandler.addAMQPServer('start', 'start_job_queue') Loading @@ -22,13 +22,10 @@ class TransferService(object): # Push self.jobHandler.addAMQPServer('store', 'store_job_queue') # Preprocessor self.preprocessor.start() def run(self): self.jobHandler.run() def start(self): self.storePreprocessor.start() self.jobHandler.start() ts = TransferService() ts.run() time.sleep(3) ts.start() print("Transfer service is RUNNING...") No newline at end of file