-
-
Notifications
You must be signed in to change notification settings - Fork 44
/
ThreadManager.java
125 lines (112 loc) · 3 KB
/
ThreadManager.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
package me.nallar.tickthreading.minecraft;
import java.util.Collection;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.atomic.AtomicInteger;
import me.nallar.tickthreading.Log;
public class ThreadManager {
private final BlockingQueue<Runnable> taskQueue = new LinkedBlockingQueue<Runnable>();
private final String namePrefix;
private final Set<Thread> workThreads = new HashSet<Thread>();
private final Object readyLock = new Object();
private final AtomicInteger waiting = new AtomicInteger(0);
private final Runnable killTask = new Runnable() {
@Override
public void run() {
}
};
private void newThread(String name) {
Thread workThread = new WorkThread();
workThread.setName(name);
workThread.setDaemon(true);
workThread.start();
workThreads.add(workThread);
}
public ThreadManager(int threads, String name) {
namePrefix = name;
addThreads(threads);
}
public void waitForCompletion() {
while (waiting.get() > 0) {
try {
synchronized (readyLock) {
readyLock.wait(0, 100);
}
} catch (InterruptedException ignored) {
}
}
}
public void run(final List<? extends Runnable> tasks) {
Runnable arrayRunnable = new Runnable() {
private final AtomicInteger index = new AtomicInteger(0);
private final int size = tasks.size();
@Override
public void run() {
int c;
while ((c = index.getAndIncrement()) < size) {
tasks.get(c).run();
}
}
};
for (int i = 0, len = workThreads.size(); i < len; i++) {
runBackground(arrayRunnable);
}
waitForCompletion();
}
public void run(Collection<? extends Runnable> tasks) {
for (Runnable runnable : tasks) {
runBackground(runnable);
}
waitForCompletion();
}
public void runBackground(Runnable runnable) {
if (taskQueue.add(runnable)) {
waiting.incrementAndGet();
} else {
Log.severe("Failed to add " + runnable);
}
}
private void addThreads(int number) {
number += workThreads.size();
for (int i = workThreads.size() + 1; i <= number; i++) {
newThread(namePrefix + " - " + i);
}
}
private void killThreads(int number) {
for (int i = 0; i < number; i++) {
taskQueue.add(killTask);
}
}
public void stop() {
killThreads(workThreads.size());
}
private class WorkThread extends Thread {
@Override
public void run() {
while (true) {
try {
Runnable runnable;
synchronized (taskQueue) {
runnable = taskQueue.take();
}
if (runnable == killTask) {
workThreads.remove(this);
return;
}
runnable.run();
} catch (InterruptedException ignored) {
} catch (Exception e) {
Log.severe("Unhandled exception in worker thread", e);
}
if (waiting.decrementAndGet() == 0) {
synchronized (readyLock) {
readyLock.notify();
}
}
}
}
}
}