diff --git a/OpenSim/Framework/Monitoring/JobEngine.cs b/OpenSim/Framework/Monitoring/JobEngine.cs index ca4428908f..850e99b1a6 100644 --- a/OpenSim/Framework/Monitoring/JobEngine.cs +++ b/OpenSim/Framework/Monitoring/JobEngine.cs @@ -51,25 +51,11 @@ namespace OpenSim.Framework.Monitoring /// public bool IsRunning { get; private set; } - /// - /// The current job that the engine is running. - /// - /// - /// Will be null if no job is currently running. - /// - private Job m_currentJob; - public Job CurrentJob { get { return m_currentJob;} } - /// /// Number of jobs waiting to be processed. /// public int JobsWaiting { get { return m_jobQueue.Count; } } - /// - /// The timeout in milliseconds to wait for at least one event to be written when the recorder is stopping. - /// - public int RequestProcessTimeoutOnStop { get; set; } - /// /// Controls whether we need to warn in the log about exceeding the max queue size. /// @@ -84,15 +70,16 @@ namespace OpenSim.Framework.Monitoring private CancellationTokenSource m_cancelSource; private int m_timeout = -1; + private int m_concurrency = 1; - private bool m_threadRunnig = false; + private int m_numberThreads = 0; - public JobEngine(string name, string loggingName, int timeout = -1) + public JobEngine(string name, string loggingName, int timeout = -1, int concurrency = 1) { Name = name; LoggingName = loggingName; m_timeout = timeout; - RequestProcessTimeoutOnStop = 5000; + m_concurrency = concurrency; } public void Start() @@ -101,12 +88,10 @@ namespace OpenSim.Framework.Monitoring { if (IsRunning) return; - + if(m_concurrency < 1) + m_concurrency = 1; IsRunning = true; - m_cancelSource = new CancellationTokenSource(); - //WorkManager.RunInThreadPool(ProcessRequests, null, Name, false); - //m_threadRunnig = true; } } @@ -122,10 +107,10 @@ namespace OpenSim.Framework.Monitoring m_log.DebugFormat("[JobEngine] Stopping {0}", Name); IsRunning = false; - if(m_threadRunnig) + if(m_numberThreads > 0) { m_cancelSource.Cancel(); - m_threadRunnig = false; + Thread.Yield(); } } finally @@ -198,22 +183,21 @@ namespace OpenSim.Framework.Monitoring /// public bool QueueJob(Job job) { - lock(JobLock) - { - if(!IsRunning) - return false; - - if(!m_threadRunnig) - { - WorkManager.RunInThreadPool(ProcessRequests, null, Name, false); - m_threadRunnig = true; - } - } + if (!IsRunning) + return false; if (m_jobQueue.Count < m_jobQueue.BoundedCapacity) { m_jobQueue.Add(job); + lock (JobLock) + { + if (m_numberThreads < m_concurrency && m_numberThreads < m_jobQueue.Count) + { + Util.FireAndForget(ProcessRequests, null, Name, false); + ++m_numberThreads; + } + } if (!m_warnOverMaxQueue) m_warnOverMaxQueue = true; @@ -233,24 +217,15 @@ namespace OpenSim.Framework.Monitoring } } - private void ProcessRequests(Object o) + private void ProcessRequests(object o) { - while(IsRunning) + Job currentJob; + while (IsRunning) { try { - if(!m_jobQueue.TryTake(out m_currentJob, m_timeout, m_cancelSource.Token)) - { - lock(JobLock) - m_threadRunnig = false; + if(!m_jobQueue.TryTake(out currentJob, m_timeout, m_cancelSource.Token)) break; - } - } - catch (OperationCanceledException) - { - m_log.DebugFormat("[JobEngine] {0} Canceled ignoring {1} jobs in queue", - Name, m_jobQueue.Count); - break; } catch { @@ -258,25 +233,26 @@ namespace OpenSim.Framework.Monitoring } if(LogLevel >= 1) - m_log.DebugFormat("[{0}]: Processing job {1}",LoggingName,m_currentJob.Name); + m_log.DebugFormat("[{0}]: Processing job {1}",LoggingName,currentJob.Name); try { - m_currentJob.Action(); + currentJob.Action(); } catch(Exception e) { - m_log.Error( - string.Format( - "[{0}]: Job {1} failed, continuing. Exception ",LoggingName,m_currentJob.Name),e); + m_log.ErrorFormat( + "[{0}]: Job {1} failed, continuing. Exception {2}",LoggingName, currentJob.Name, e); } if(LogLevel >= 1) - m_log.DebugFormat("[{0}]: Processed job {1}",LoggingName,m_currentJob.Name); + m_log.DebugFormat("[{0}]: Processed job {1}",LoggingName,currentJob.Name); - m_currentJob.Action = null; - m_currentJob = null; + currentJob.Action = null; + currentJob = null; } + lock (JobLock) + --m_numberThreads; } public class Job