forked from ReactiveX/RxJava
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathTestScheduler.java
More file actions
180 lines (151 loc) · 5.47 KB
/
Copy pathTestScheduler.java
File metadata and controls
180 lines (151 loc) · 5.47 KB
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
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
/**
* Copyright 2014 Netflix, Inc.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package rx.schedulers;
import java.util.Comparator;
import java.util.PriorityQueue;
import java.util.Queue;
import java.util.concurrent.TimeUnit;
import rx.Scheduler;
import rx.Subscription;
import rx.functions.Action0;
import rx.subscriptions.BooleanSubscription;
import rx.subscriptions.Subscriptions;
/**
* The {@code TestScheduler} is useful for debugging. It allows you to test schedules of events by manually
* advancing the clock at whatever pace you choose.
*/
public class TestScheduler extends Scheduler {
private final Queue<TimedAction> queue = new PriorityQueue<TimedAction>(11, new CompareActionsByTime());
private static long counter = 0;
private static final class TimedAction {
private final long time;
private final Action0 action;
private final Worker scheduler;
private final long count = counter++; // for differentiating tasks at same time
private TimedAction(Worker scheduler, long time, Action0 action) {
this.time = time;
this.action = action;
this.scheduler = scheduler;
}
@Override
public String toString() {
return String.format("TimedAction(time = %d, action = %s)", time, action.toString());
}
}
private static class CompareActionsByTime implements Comparator<TimedAction> {
@Override
public int compare(TimedAction action1, TimedAction action2) {
if (action1.time == action2.time) {
return Long.valueOf(action1.count).compareTo(Long.valueOf(action2.count));
} else {
return Long.valueOf(action1.time).compareTo(Long.valueOf(action2.time));
}
}
}
// Storing time in nanoseconds internally.
private long time;
@Override
public long now() {
return TimeUnit.NANOSECONDS.toMillis(time);
}
/**
* Moves the Scheduler's clock forward by a specified amount of time.
*
* @param delayTime
* the amount of time to move the Scheduler's clock forward
* @param unit
* the units of time that {@code delayTime} is expressed in
*/
public void advanceTimeBy(long delayTime, TimeUnit unit) {
advanceTimeTo(time + unit.toNanos(delayTime), TimeUnit.NANOSECONDS);
}
/**
* Moves the Scheduler's clock to a particular moment in time.
*
* @param delayTime
* the point in time to move the Scheduler's clock to
* @param unit
* the units of time that {@code delayTime} is expressed in
*/
public void advanceTimeTo(long delayTime, TimeUnit unit) {
long targetTime = unit.toNanos(delayTime);
triggerActions(targetTime);
}
/**
* Triggers any actions that have not yet been triggered and that are scheduled to be triggered at or
* before this Scheduler's present time.
*/
public void triggerActions() {
triggerActions(time);
}
private void triggerActions(long targetTimeInNanos) {
while (!queue.isEmpty()) {
TimedAction current = queue.peek();
if (current.time > targetTimeInNanos) {
break;
}
// if scheduled time is 0 (immediate) use current virtual time
time = current.time == 0 ? time : current.time;
queue.remove();
// Only execute if not unsubscribed
if (!current.scheduler.isUnsubscribed()) {
current.action.call();
}
}
time = targetTimeInNanos;
}
@Override
public Worker createWorker() {
return new InnerTestScheduler();
}
private final class InnerTestScheduler extends Worker {
private final BooleanSubscription s = new BooleanSubscription();
@Override
public void unsubscribe() {
s.unsubscribe();
}
@Override
public boolean isUnsubscribed() {
return s.isUnsubscribed();
}
@Override
public Subscription schedule(Action0 action, long delayTime, TimeUnit unit) {
final TimedAction timedAction = new TimedAction(this, time + unit.toNanos(delayTime), action);
queue.add(timedAction);
return Subscriptions.create(new Action0() {
@Override
public void call() {
queue.remove(timedAction);
}
});
}
@Override
public Subscription schedule(Action0 action) {
final TimedAction timedAction = new TimedAction(this, 0, action);
queue.add(timedAction);
return Subscriptions.create(new Action0() {
@Override
public void call() {
queue.remove(timedAction);
}
});
}
@Override
public long now() {
return TestScheduler.this.now();
}
}
}