140def timelimited(timeout=None, retry=1, timefactor=1.5, sleeptime=10, defaultrc=None):
141
142 import traceback
143 import queue
144 import multiprocessing as mp
145
146 from sys import exc_info
147 from PyJobTransforms.trfExceptions import TransformTimeoutException, TransformInternalException
148
149 msg = logging.getLogger(__name__)
150
151 def internal(func):
152
153
158 def funcWithQueue(queue, *args, **kwargs):
159 try:
160 result = func(*args, **kwargs)
161 queue.put((True, result))
162 except Exception:
163 exc0=exc_info()[0]
164 exc1=exc_info()[1]
165 exc2=traceback.format_exc()
166 msg.warning('In time limited function %s an exception occurred', func.__name__)
167 msg.warning('Original traceback:')
168 msg.warning(exc2)
169 queue.put((False,(exc0, exc1, exc2)))
170
171 def funcWithTimeout(*args, **kwargs):
172 ltimeout=timeout
173 lretry=retry
174 ltimefactor=timefactor
175 lsleeptime=sleeptime
176 ldefaultrc=defaultrc
177
178 if 'timeout' in kwargs:
179 ltimeout=kwargs.pop('timeout')
180 if 'retry' in kwargs:
181 lretry=kwargs.pop('retry')
182 if 'timefactor' in kwargs:
183 ltimefactor=kwargs.pop('timefactor')
184 if 'sleeptime' in kwargs:
185 lsleeptime=kwargs.pop('sleeptime')
186 if 'defaultrc' in kwargs:
187 ldefaultrc=kwargs.pop('defaultrc')
188
189 if ltimeout is None:
190
191 msg.debug('Running {0}: {1} {2} without timeout'.format(func, args, kwargs))
192 return func(*args, **kwargs)
193
194 n=0
195 while n<=lretry:
196 msg.info('Try %i out of %i (time limit %s s) to call %s.', n+1, retry+1, ltimeout, func.__name__)
197 starttime = time.time()
198 q=mp.Queue(maxsize=1)
199 nargs = (q,) + args
200 proc=mp.Process(target=funcWithQueue, args=nargs, kwargs=kwargs)
201 proc.start()
202 try:
203
204 flag,result = q.get(block=True, timeout=ltimeout)
205 proc.join(60)
206 msg.info('Executed call within %d s.', time.time()-starttime)
207 if flag:
208 return result
209 else:
210 msg.warning('But an exception occurred in function %s.', func.__name__)
211 msg.warning('Returning default return code %s.', ldefaultrc)
212 return ldefaultrc
213 except queue.Empty:
214
215 msg.warning('Timeout limit of %d s reached. Kill subprocess and its children.', ltimeout)
216 parent=proc.pid
217 pids=[parent]
218 pids.extend(trfUtils.listChildren(parent=parent, listOrphans = False))
219 trfUtils.infanticide(pids)
220 proc.join(60)
221 if n!=lretry:
222 msg.info('Going to sleep for %d s.', lsleeptime)
223 time.sleep(lsleeptime)
224 n+=1
225 ltimeout*=ltimefactor
226 lsleeptime*=ltimefactor
227 except IOError:
228 errMsg = "IOError while communicating with subprocess"
229 msg.error(errMsg)
230 raise TransformInternalException(trfExit.nameToCode("TRF_EXTERNAL"), errMsg)
231
232 msg.warning('All %i tries failed!', n)
233 raise TransformTimeoutException(trfExit.nameToCode('TRF_EXEC_TIMEOUT'), 'Timeout in function %s' % (func.__name__))
234
235 return funcWithTimeout
236
237 functools.update_wrapper(funcWithTimeout, func)
238
239 return internal