annotate Implab/Parallels/WorkerPool.cs @ 14:e943453e5039 promises

Implemented interllocked queue fixed promise syncronization
author cin
date Wed, 06 Nov 2013 17:49:12 +0400
parents b0feb5b9ad1c
children 0f982f9b7d4d
Ignore whitespace changes - Everywhere: Within whitespace: At end of lines:
rev   line source
12
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
1 using System;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
2 using System.Collections.Generic;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
3 using System.Linq;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
4 using System.Text;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
5 using System.Threading;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
6 using System.Diagnostics;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
7
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
8 namespace Implab.Parallels {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
9 public class WorkerPool : IDisposable {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
10 readonly int m_minThreads;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
11 readonly int m_maxThreads;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
12 int m_runningThreads;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
13 object m_lock = new object();
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
14
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
15 bool m_disposed = false;
13
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
16
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
17 // this event will signal that workers can try to fetch a task from queue or the pool has been disposed
12
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
18 ManualResetEvent m_hasTasks = new ManualResetEvent(false);
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
19 Queue<Action> m_queue = new Queue<Action>();
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
20
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
21 public WorkerPool(int min, int max) {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
22 if (min < 0)
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
23 throw new ArgumentOutOfRangeException("min");
13
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
24 if (max <= 0)
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
25 throw new ArgumentOutOfRangeException("max");
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
26
12
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
27 if (min > max)
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
28 min = max;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
29 m_minThreads = min;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
30 m_maxThreads = max;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
31
13
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
32 InitPool();
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
33 }
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
34
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
35 public WorkerPool(int max)
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
36 : this(0, max) {
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
37 }
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
38
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
39 public WorkerPool() {
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
40 int maxThreads, maxCP;
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
41 ThreadPool.GetMaxThreads(out maxThreads, out maxCP);
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
42
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
43 m_minThreads = 0;
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
44 m_maxThreads = maxThreads;
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
45
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
46 InitPool();
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
47 }
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
48
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
49 void InitPool() {
12
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
50 for (int i = 0; i < m_minThreads; i++)
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
51 StartWorker();
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
52 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
53
13
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
54 public int ThreadCount {
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
55 get {
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
56 return m_runningThreads;
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
57 }
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
58 }
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
59
12
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
60 public Promise<T> Invoke<T>(Func<T> task) {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
61 if (m_disposed)
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
62 throw new ObjectDisposedException(ToString());
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
63 if (task == null)
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
64 throw new ArgumentNullException("task");
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
65
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
66 var promise = new Promise<T>();
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
67
13
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
68 var queueLen = EnqueueTask(delegate() {
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
69 try {
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
70 promise.Resolve(task());
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
71 } catch (Exception e) {
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
72 promise.Reject(e);
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
73 }
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
74 });
12
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
75
13
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
76 if (queueLen > 1)
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
77 StartWorker();
12
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
78
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
79 return promise;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
80 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
81
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
82 bool StartWorker() {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
83 var current = m_runningThreads;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
84 // use spins to allocate slot for the new thread
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
85 do {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
86 if (current >= m_maxThreads)
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
87 // no more slots left
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
88 return false;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
89 } while (current != Interlocked.CompareExchange(ref m_runningThreads, current + 1, current));
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
90
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
91 // slot successfully allocated
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
92
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
93 var worker = new Thread(this.Worker);
13
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
94 worker.IsBackground = true;
12
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
95 worker.Start();
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
96
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
97 return true;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
98 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
99
13
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
100 int EnqueueTask(Action task) {
12
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
101 Debug.Assert(task != null);
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
102 lock (m_queue) {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
103 m_queue.Enqueue(task);
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
104 m_hasTasks.Set();
13
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
105 return m_queue.Count;
12
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
106 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
107 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
108
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
109 bool FetchTask(out Action task) {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
110 task = null;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
111
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
112 while (true) {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
113
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
114 m_hasTasks.WaitOne();
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
115
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
116 if (m_disposed)
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
117 return false;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
118
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
119 lock (m_queue) {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
120 if (m_queue.Count > 0) {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
121 task = m_queue.Dequeue();
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
122 return true;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
123 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
124
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
125 // no tasks left
13
b0feb5b9ad1c small fixes, WorkerPool still incomplete
cin
parents: 12
diff changeset
126 // signal that no more tasks left, current lock ensures that this event won't suppress newly added task
12
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
127 m_hasTasks.Reset();
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
128 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
129
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
130 bool exit = true;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
131
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
132 var current = m_runningThreads;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
133 do {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
134 if (current <= m_minThreads) {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
135 exit = false; // this thread should return and wait for the new events
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
136 break;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
137 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
138 } while (current != Interlocked.CompareExchange(ref m_runningThreads, current - 1, current));
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
139
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
140 if (exit)
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
141 return false;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
142 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
143 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
144
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
145 void Worker() {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
146 Action task;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
147 while (FetchTask(out task))
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
148 task();
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
149 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
150
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
151 protected virtual void Dispose(bool disposing) {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
152 if (disposing) {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
153 lock (m_lock) {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
154 if (m_disposed)
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
155 return;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
156 m_disposed = true;
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
157 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
158 m_hasTasks.Set();
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
159 GC.SuppressFinalize(this);
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
160 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
161 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
162
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
163 public void Dispose() {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
164 Dispose(true);
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
165 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
166
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
167 ~WorkerPool() {
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
168 Dispose(false);
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
169 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
170 }
eb418ba8275b refactoring, added WorkerPool
cin
parents:
diff changeset
171 }