Commit e2ae28e2 authored by Cristiano Urban's avatar Cristiano Urban
Browse files

Removed job_cache + added code to store jobs on Postgres db.

parent a6779680
Loading
Loading
Loading
Loading
+19 −11
Original line number Diff line number Diff line
@@ -11,8 +11,8 @@ import json
#from enum import Enum

from amqp_server import AMQPServer
from db_connector import DbConnector
from job import Job
from job_cache import JobCache
from job_queue import JobQueue
from config import Config

@@ -23,10 +23,12 @@ class StoreAMQPServer(AMQPServer):
        self.type = "store"
        self.storeAck = False
        config = Config("vos_ts.conf")       
        self.params = config.loadSection("job_cache")
        self.jobCache = JobCache(self.params["host"], 
        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_write"))
                                  self.params["db"])
        self.pendingQueueWrite = JobQueue("pending")        
        self.job = None
        self.username = None
@@ -38,10 +40,11 @@ class StoreAMQPServer(AMQPServer):
        if "requestType" not in requestBody or "userName" not in requestBody:
            response = { "errorCode": 1, "errorMsg": "Malformed request, missing parameters." }
        elif requestBody["requestType"] == "CSTORE" or requestBody["requestType"] == "HSTORE":
            user = requestBody["userName"]
            self.job = Job()
            self.job.setType("other")            
            self.job.setInfo(requestBody)
            self.job.setPhase("PENDING")            
            user = requestBody["userName"]
            folderPath = "/home/" + user + "/store"
            userInfo = self.userInfo(user)
            # Check if the user exists on the transfer node
@@ -74,10 +77,13 @@ class StoreAMQPServer(AMQPServer):
                self.storeAck = False
                user = requestBody["userName"]
                self.prepare(user)
                self.jobCache.set(self.job)
                redisResponse = self.jobCache.get(self.job.jobId)
                self.dbConn.connect()
                self.job.setOwnerId(self.dbConn.getRapId(self.username))
                self.dbConn.insertJob(self.job)
                dbResponse = self.dbConn.getJob(self.job.jobId)
                self.dbConn.disconnect()
                self.pendingQueueWrite.insertJob(self.job)
                if "error" in redisResponse:
                if "error" in dbResponse:
                    response = { "responseType": "ERROR",
                                 "errorCode": 5,
                                 "errorMsg": "Job creation failed." }
@@ -110,6 +116,8 @@ class StoreAMQPServer(AMQPServer):
                os.chown(os.path.join(folder, f), 0, 0)
                os.chmod(os.path.join(folder, f), 0o555)

    # parse /etc/passwd on the transfer node and get
    # username, uid and gid
    def userInfo(self, username):
        fp = open("/etc/passwd", 'r')
        for line in fp: