| 49 | |
| 50 | |
| 51 | class workThread(threading.Thread): |
| 52 | |
| 53 | def __init__(self, in_queue, outqueue, apiClient, db=None, lock=None): |
| 54 | threading.Thread.__init__(self) |
| 55 | self.inqueue = in_queue |
| 56 | self.output = outqueue |
| 57 | self.connection = apiClient.connection.__copy__() |
| 58 | self.db = None |
| 59 | self.lock = lock |
| 60 | |
| 61 | def queryAsynJob(self, job): |
| 62 | if job.jobId is None: |
| 63 | return job |
| 64 | |
| 65 | try: |
| 66 | self.lock.acquire() |
| 67 | result = self.connection.poll(job.jobId, job.responsecls).jobresult |
| 68 | except cloudstackException.CloudstackAPIException as e: |
| 69 | result = str(e) |
| 70 | finally: |
| 71 | self.lock.release() |
| 72 | |
| 73 | job.result = result |
| 74 | return job |
| 75 | |
| 76 | def executeCmd(self, job): |
| 77 | cmd = job.cmd |
| 78 | |
| 79 | jobstatus = jobStatus() |
| 80 | jobId = None |
| 81 | try: |
| 82 | self.lock.acquire() |
| 83 | |
| 84 | if cmd.isAsync == "false": |
| 85 | jobstatus.startTime = datetime.datetime.now() |
| 86 | |
| 87 | result = self.connection.marvin_request(cmd) |
| 88 | jobstatus.result = result |
| 89 | jobstatus.endTime = datetime.datetime.now() |
| 90 | jobstatus.duration =\ |
| 91 | time.mktime(jobstatus.endTime.timetuple()) - time.mktime( |
| 92 | jobstatus.startTime.timetuple()) |
| 93 | else: |
| 94 | result = self.connection.marvinRequest(cmd) |
| 95 | if result is None: |
| 96 | jobstatus.status = False |
| 97 | else: |
| 98 | jobId = result.jobid |
| 99 | jobstatus.jobId = jobId |
| 100 | try: |
| 101 | responseName =\ |
| 102 | cmd.__class__.__name__.replace("Cmd", "Response") |
| 103 | jobstatus.responsecls =\ |
| 104 | jsonHelper.getclassFromName(cmd, responseName) |
| 105 | except: |
| 106 | pass |
| 107 | jobstatus.status = True |
| 108 | except cloudstackException.CloudstackAPIException as e: |
no outgoing calls
no test coverage detected