62 lines
1.6 KiB
C#
62 lines
1.6 KiB
C#
|
using System;
|
|||
|
using System.Threading;
|
|||
|
namespace ln.threading
|
|||
|
{
|
|||
|
public class DynamicPool : Pool
|
|||
|
{
|
|||
|
public int Timeout { get; set; }
|
|||
|
|
|||
|
public DynamicPool() :this(Environment.ProcessorCount){}
|
|||
|
public DynamicPool(int maxPoolSize)
|
|||
|
:base(maxPoolSize)
|
|||
|
{
|
|||
|
Timeout = 15000;
|
|||
|
}
|
|||
|
|
|||
|
public override void Start()
|
|||
|
{
|
|||
|
if ((State == PoolState.RUN) || (State == PoolState.SHUTDOWN))
|
|||
|
throw new NotSupportedException("Pool can only be started if not running");
|
|||
|
|
|||
|
State = PoolState.RUN;
|
|||
|
}
|
|||
|
|
|||
|
protected override PoolJob WaitForJob(PoolThread poolThread)
|
|||
|
{
|
|||
|
if (State == PoolState.RUN)
|
|||
|
{
|
|||
|
lock (waitingThreads)
|
|||
|
{
|
|||
|
PoolJob poolJob = Dequeue(poolThread);
|
|||
|
if (poolJob != null)
|
|||
|
return poolJob;
|
|||
|
|
|||
|
waitingThreads.Add(poolThread);
|
|||
|
Monitor.Wait(waitingThreads, Timeout);
|
|||
|
waitingThreads.Remove(poolThread);
|
|||
|
|
|||
|
poolJob = Dequeue(poolThread);
|
|||
|
return poolJob;
|
|||
|
}
|
|||
|
}
|
|||
|
return null;
|
|||
|
}
|
|||
|
|
|||
|
protected override void PulseWaitingThread()
|
|||
|
{
|
|||
|
lock (waitingThreads)
|
|||
|
{
|
|||
|
if (waitingThreads.Count > 0)
|
|||
|
{
|
|||
|
Monitor.Pulse(waitingThreads);
|
|||
|
}
|
|||
|
else if (CurrentPoolSize < PoolSize)
|
|||
|
{
|
|||
|
CreatePoolThread();
|
|||
|
}
|
|||
|
}
|
|||
|
}
|
|||
|
|
|||
|
}
|
|||
|
}
|