blob: 180fbdc4618b88f5c15b868706b5628a1b06687f (
plain)
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
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
|
/*
* Copyright (C) 2010, Google Inc. and others
*
* This program and the accompanying materials are made available under the
* terms of the Eclipse Distribution License v. 1.0 which is available at
* https://www.eclipse.org/org/documents/edl-v10.php.
*
* SPDX-License-Identifier: BSD-3-Clause
*/
package org.eclipse.jgit.lib;
import java.util.concurrent.Semaphore;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.locks.ReentrantLock;
/**
* Wrapper around the general {@link org.eclipse.jgit.lib.ProgressMonitor} to
* make it thread safe.
*
* Updates to the underlying ProgressMonitor are made only from the thread that
* allocated this wrapper. Callers are responsible for ensuring the allocating
* thread uses {@link #pollForUpdates()} or {@link #waitForCompletion()} to
* update the underlying ProgressMonitor.
*
* Only {@link #update(int)}, {@link #isCancelled()}, and {@link #endWorker()}
* may be invoked from a worker thread. All other methods of the ProgressMonitor
* interface can only be called from the thread that allocates this wrapper.
*/
public class ThreadSafeProgressMonitor implements ProgressMonitor {
private final ProgressMonitor pm;
private final ReentrantLock lock;
private final Thread mainThread;
private final AtomicInteger workers;
private final AtomicInteger pendingUpdates;
private final Semaphore process;
/**
* Wrap a ProgressMonitor to be thread safe.
*
* @param pm
* the underlying monitor to receive events.
*/
public ThreadSafeProgressMonitor(ProgressMonitor pm) {
this.pm = pm;
this.lock = new ReentrantLock();
this.mainThread = Thread.currentThread();
this.workers = new AtomicInteger(0);
this.pendingUpdates = new AtomicInteger(0);
this.process = new Semaphore(0);
}
/** {@inheritDoc} */
@Override
public void start(int totalTasks) {
if (!isMainThread())
throw new IllegalStateException();
pm.start(totalTasks);
}
/** {@inheritDoc} */
@Override
public void beginTask(String title, int totalWork) {
if (!isMainThread())
throw new IllegalStateException();
pm.beginTask(title, totalWork);
}
/**
* Notify the monitor a worker is starting.
*/
public void startWorker() {
startWorkers(1);
}
/**
* Notify the monitor of workers starting.
*
* @param count
* the number of worker threads that are starting.
*/
public void startWorkers(int count) {
workers.addAndGet(count);
}
/**
* Notify the monitor a worker is finished.
*/
public void endWorker() {
if (workers.decrementAndGet() == 0)
process.release();
}
/**
* Non-blocking poll for pending updates.
*
* This method can only be invoked by the same thread that allocated this
* ThreadSafeProgressMonior.
*/
public void pollForUpdates() {
assert isMainThread();
doUpdates();
}
/**
* Process pending updates and wait for workers to finish.
*
* This method can only be invoked by the same thread that allocated this
* ThreadSafeProgressMonior.
*
* @throws java.lang.InterruptedException
* if the main thread is interrupted while waiting for
* completion of workers.
*/
public void waitForCompletion() throws InterruptedException {
assert isMainThread();
while (0 < workers.get()) {
doUpdates();
process.acquire();
}
doUpdates();
}
private void doUpdates() {
int cnt = pendingUpdates.getAndSet(0);
if (0 < cnt)
pm.update(cnt);
}
/** {@inheritDoc} */
@Override
public void update(int completed) {
if (0 == pendingUpdates.getAndAdd(completed))
process.release();
}
/** {@inheritDoc} */
@Override
public boolean isCancelled() {
lock.lock();
try {
return pm.isCancelled();
} finally {
lock.unlock();
}
}
/** {@inheritDoc} */
@Override
public void endTask() {
if (!isMainThread())
throw new IllegalStateException();
pm.endTask();
}
private boolean isMainThread() {
return Thread.currentThread() == mainThread;
}
}
|