Skip to content

Commit ecfe547

Browse files
committed
change: remove targetStatus from ItemTicket
1 parent 745cc56 commit ecfe547

5 files changed

Lines changed: 38 additions & 70 deletions

File tree

src/main/java/com/ishland/flowsched/scheduler/ItemHolder.java

Lines changed: 18 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -133,60 +133,60 @@ public ItemStatus<K, V, Ctx> upgradingStatusTo() {
133133
return pair != null ? pair.right() : null;
134134
}
135135

136-
public void addTicket(ItemTicket<K, V, Ctx> ticket) {
136+
public void addTicket(ItemStatus<K, V, Ctx> targetStatus, ItemTicket ticket) {
137137
assertOpen();
138138
boolean needConsumption;
139139
synchronized (this) {
140-
final boolean add = this.tickets.checkAdd(ticket);
140+
final boolean add = this.tickets.checkAdd(targetStatus, ticket);
141141
if (!add) {
142142
throw new IllegalStateException("Ticket already exists");
143143
}
144-
this.tickets.addUnchecked(ticket);
144+
this.tickets.addUnchecked(targetStatus);
145145
createFutures();
146-
needConsumption = ticket.getTargetStatus().ordinal() <= this.getStatus().ordinal();
146+
needConsumption = targetStatus.ordinal() <= this.getStatus().ordinal();
147147
}
148148

149149
if (needConsumption) {
150150
ticket.consumeCallback();
151151
}
152-
this.validateRequestedFutures(ticket.getTargetStatus());
152+
this.validateRequestedFutures(targetStatus);
153153
}
154154

155-
public void removeTicket(ItemTicket<K, V, Ctx> ticket) {
155+
public void removeTicket(ItemStatus<K, V, Ctx> targetStatus, ItemTicket ticket) {
156156
assertOpen();
157157
synchronized (this) {
158-
final boolean remove = this.tickets.checkRemove(ticket);
158+
final boolean remove = this.tickets.checkRemove(targetStatus, ticket);
159159
if (!remove) {
160160
throw new IllegalStateException("Ticket does not exist");
161161
}
162-
this.tickets.removeUnchecked(ticket);
162+
this.tickets.removeUnchecked(targetStatus);
163163
}
164164
}
165165

166-
public void swapTicket(ItemTicket<K, V, Ctx> orig, ItemTicket<K, V, Ctx> ticket) {
166+
public void swapTicket(ItemStatus<K, V, Ctx> origStatus, ItemTicket orig, ItemStatus<K, V, Ctx> targetStatus, ItemTicket ticket) {
167167
assertOpen();
168168
boolean needConsumption;
169169
synchronized (this) {
170-
final boolean add = this.tickets.checkAdd(ticket);
170+
final boolean add = this.tickets.checkAdd(targetStatus, ticket);
171171
if (!add) {
172172
throw new IllegalStateException("Ticket already exists");
173173
}
174-
final boolean remove = this.tickets.checkRemove(orig);
174+
final boolean remove = this.tickets.checkRemove(origStatus, orig);
175175
if (!remove) {
176-
boolean value = this.tickets.checkRemove(ticket); // revert side effect
176+
boolean value = this.tickets.checkRemove(targetStatus, ticket); // revert side effect
177177
Assertions.assertTrue(value);
178178
throw new IllegalStateException("Ticket does not exist");
179179
}
180-
this.tickets.addUnchecked(ticket);
181-
this.tickets.removeUnchecked(orig);
180+
this.tickets.addUnchecked(targetStatus);
181+
this.tickets.removeUnchecked(origStatus);
182182
createFutures();
183-
needConsumption = ticket.getTargetStatus().ordinal() <= this.getStatus().ordinal();
183+
needConsumption = targetStatus.ordinal() <= this.getStatus().ordinal();
184184
}
185185

186186
if (needConsumption) {
187187
ticket.consumeCallback();
188188
}
189-
this.validateRequestedFutures(ticket.getTargetStatus());
189+
this.validateRequestedFutures(targetStatus);
190190
}
191191

192192
public void submitOp(CompletionStage<Void> op) {
@@ -286,7 +286,7 @@ private void markDirty0(StatusAdvancingScheduler<K, V, Ctx, UserData> scheduler)
286286

287287
public boolean setStatus(ItemStatus<K, V, Ctx> status, boolean isCancellation) {
288288
assertOpen();
289-
ItemTicket<K, V, Ctx>[] ticketsToFire = null;
289+
ItemTicket[] ticketsToFire = null;
290290
CompletableFuture<Void> futureToFire = null;
291291
synchronized (this) {
292292
final ItemStatus<K, V, Ctx> prevStatus = this.getStatus();
@@ -330,7 +330,7 @@ public boolean setStatus(ItemStatus<K, V, Ctx> status, boolean isCancellation) {
330330
}
331331
}
332332
if (ticketsToFire != null) {
333-
for (ItemTicket<K, V, Ctx> ticket : ticketsToFire) {
333+
for (ItemTicket ticket : ticketsToFire) {
334334
ticket.consumeCallback();
335335
}
336336
}

src/main/java/com/ishland/flowsched/scheduler/ItemTicket.java

Lines changed: 4 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -3,21 +3,19 @@
33
import java.util.Objects;
44
import java.util.concurrent.atomic.AtomicReferenceFieldUpdater;
55

6-
public class ItemTicket<K, V, Ctx> {
6+
public class ItemTicket {
77

88
private static final AtomicReferenceFieldUpdater<ItemTicket, Runnable> CALLBACK_UPDATER = AtomicReferenceFieldUpdater.newUpdater(ItemTicket.class, Runnable.class, "callback");
99

1010
private final int hashCode;
1111
private final TicketType type;
1212
private final Object source;
13-
private final ItemStatus<K, V, Ctx> targetStatus;
1413
private volatile Runnable callback = null;
1514
// private int hash = 0;
1615

17-
public ItemTicket(TicketType type, Object source, ItemStatus<K, V, Ctx> targetStatus, Runnable callback) {
16+
public ItemTicket(TicketType type, Object source, Runnable callback) {
1817
this.type = Objects.requireNonNull(type);
1918
this.source = Objects.requireNonNull(source);
20-
this.targetStatus = Objects.requireNonNull(targetStatus);
2119
this.callback = callback;
2220
this.hashCode = this.hashCode0();
2321
}
@@ -26,10 +24,6 @@ public Object getSource() {
2624
return this.source;
2725
}
2826

29-
public ItemStatus<K, V, Ctx> getTargetStatus() {
30-
return this.targetStatus;
31-
}
32-
3327
public TicketType getType() {
3428
return this.type;
3529
}
@@ -49,23 +43,16 @@ public void consumeCallback() {
4943
public boolean equals(Object o) {
5044
if (this == o) return true;
5145
if (o == null || getClass() != o.getClass()) return false;
52-
ItemTicket<?, ?, ?> that = (ItemTicket<?, ?, ?>) o;
53-
return type == that.type && Objects.equals(source, that.source) && Objects.equals(targetStatus, that.targetStatus);
46+
ItemTicket that = (ItemTicket) o;
47+
return type == that.type && Objects.equals(source, that.source);
5448
}
5549

56-
// public boolean equalsAlternative(ItemTicket<K, V, Ctx> that) {
57-
// if (this == that) return true;
58-
// if (that == null) return false;
59-
// return type == that.type && Objects.equals(source, that.source);
60-
// }
61-
6250
private int hashCode0() {
6351
// inlined version of Objects.hash(type, source, targetStatus)
6452
int result = 1;
6553

6654
result = 31 * result + type.hashCode();
6755
result = 31 * result + source.hashCode();
68-
result = 31 * result + targetStatus.hashCode();
6956
return result;
7057
}
7158

@@ -74,19 +61,6 @@ public int hashCode() {
7461
return this.hashCode;
7562
}
7663

77-
// public int hashCodeAlternative() {
78-
// int hc = hash;
79-
// if (hc == 0) {
80-
// // inlined version of Objects.hash(type, source, targetStatus)
81-
// int result = 1;
82-
//
83-
// result = 31 * result + type.hashCode();
84-
// result = 31 * result + source.hashCode();
85-
// hc = hash = result;
86-
// }
87-
// return hc;
88-
// }
89-
9064
public static class TicketType {
9165
public static TicketType DEPENDENCY = new TicketType("flowsched:dependency");
9266
public static TicketType EXTERNAL = new TicketType("flowsched:external");

src/main/java/com/ishland/flowsched/scheduler/StatusAdvancingScheduler.java

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -534,11 +534,11 @@ public ItemHolder<K, V, Ctx, UserData> addTicket(K key, Object source, ItemStatu
534534
}
535535

536536
public ItemHolder<K, V, Ctx, UserData> addTicket(K key, ItemTicket.TicketType type, Object source, ItemStatus<K, V, Ctx> targetStatus, Runnable callback) {
537-
return this.addTicket(key, new ItemTicket<>(type, source, targetStatus, callback));
537+
return this.addTicket(key, targetStatus, new ItemTicket(type, source, callback));
538538
}
539539

540-
public ItemHolder<K, V, Ctx, UserData> addTicket(K key, ItemTicket<K, V, Ctx> ticket) {
541-
if (this.getUnloadedStatus().equals(ticket.getTargetStatus())) {
540+
public ItemHolder<K, V, Ctx, UserData> addTicket(K key, ItemStatus<K, V, Ctx> targetStatus, ItemTicket ticket) {
541+
if (this.getUnloadedStatus().equals(targetStatus)) {
542542
throw new IllegalArgumentException("Cannot add ticket to unloaded status");
543543
}
544544
try {
@@ -554,7 +554,7 @@ public ItemHolder<K, V, Ctx, UserData> addTicket(K key, ItemTicket<K, V, Ctx> ti
554554
holder.busyRefCounter().incrementRefCount();
555555
}
556556
try {
557-
holder.addTicket(ticket);
557+
holder.addTicket(targetStatus, ticket);
558558
holder.consolidateMarkDirty(this);
559559
} finally {
560560
holder.busyRefCounter().decrementRefCount();
@@ -580,17 +580,17 @@ public void removeTicket(K key, ItemTicket.TicketType type, Object source, ItemS
580580
if (holder == null) {
581581
throw new IllegalStateException("No such item");
582582
}
583-
holder.removeTicket(new ItemTicket<>(type, source, targetStatus, null));
583+
holder.removeTicket(targetStatus, new ItemTicket(type, source, null));
584584
// holder may have been removed at this point, only mark it dirty if it still exists
585585
holder.tryMarkDirty(this);
586586
}
587587

588-
public void swapTicket(K key, ItemTicket<K, V, Ctx> orig, ItemTicket<K, V, Ctx> ticket) {
588+
public void swapTicket(K key, ItemStatus<K, V, Ctx> origStatus, ItemTicket orig, ItemStatus<K, V, Ctx> targetStatus, ItemTicket ticket) {
589589
ItemHolder<K, V, Ctx, UserData> holder = this.getHolder(key);
590590
if (holder == null) {
591591
throw new IllegalStateException("No such item");
592592
}
593-
holder.swapTicket(orig, ticket);
593+
holder.swapTicket(origStatus, orig, targetStatus, ticket);
594594
holder.markDirty(this);
595595
}
596596

src/main/java/com/ishland/flowsched/scheduler/TicketSet.java

Lines changed: 8 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -3,18 +3,16 @@
33
import com.ishland.flowsched.util.Assertions;
44
import it.unimi.dsi.fastutil.objects.ObjectOpenHashSet;
55

6-
import java.lang.invoke.MethodHandles;
76
import java.lang.invoke.VarHandle;
87
import java.util.Set;
9-
import java.util.concurrent.atomic.AtomicIntegerFieldUpdater;
108

119
/**
1210
* Not thread-safe
1311
*/
1412
public class TicketSet<K, V, Ctx> {
1513

1614
private final ItemStatus<K, V, Ctx> initialStatus;
17-
private final Set<ItemTicket<K, V, Ctx>>[] status2Tickets;
15+
private final Set<ItemTicket>[] status2Tickets;
1816
private final int[] status2TicketsSize;
1917
private volatile int targetStatus = 0;
2018

@@ -30,28 +28,24 @@ public TicketSet(ItemStatus<K, V, Ctx> initialStatus, ObjectFactory objectFactor
3028
VarHandle.fullFence();
3129
}
3230

33-
public boolean checkAdd(ItemTicket<K, V, Ctx> ticket) {
34-
ItemStatus<K, V, Ctx> targetStatus = ticket.getTargetStatus();
31+
public boolean checkAdd(ItemStatus<K, V, Ctx> targetStatus, ItemTicket ticket) {
3532
final boolean added = this.status2Tickets[targetStatus.ordinal()].add(ticket);
3633
return added;
3734
}
3835

39-
public void addUnchecked(ItemTicket<K, V, Ctx> ticket) {
40-
ItemStatus<K, V, Ctx> targetStatus = ticket.getTargetStatus();
36+
public void addUnchecked(ItemStatus<K, V, Ctx> targetStatus) {
4137
this.status2TicketsSize[targetStatus.ordinal()] ++;
4238
if (targetStatus.ordinal() > this.targetStatus) {
4339
this.targetStatus = targetStatus.ordinal();
4440
}
4541
}
4642

47-
public boolean checkRemove(ItemTicket<K, V, Ctx> ticket) {
48-
ItemStatus<K, V, Ctx> targetStatus = ticket.getTargetStatus();
43+
public boolean checkRemove(ItemStatus<K, V, Ctx> targetStatus, ItemTicket ticket) {
4944
final boolean removed = this.status2Tickets[targetStatus.ordinal()].remove(ticket);
5045
return removed;
5146
}
5247

53-
public void removeUnchecked(ItemTicket<K, V, Ctx> ticket) {
54-
ItemStatus<K, V, Ctx> targetStatus = ticket.getTargetStatus();
48+
public void removeUnchecked(ItemStatus<K, V, Ctx> targetStatus) {
5549
int updated = --this.status2TicketsSize[targetStatus.ordinal()];
5650
if (updated == 0) {
5751
this.updateTargetStatus();
@@ -66,20 +60,20 @@ public ItemStatus<K, V, Ctx> getTargetStatus() {
6660
return this.initialStatus.getAllStatuses()[this.targetStatus];
6761
}
6862

69-
public Set<ItemTicket<K, V, Ctx>> getTicketsForStatus(ItemStatus<K, V, Ctx> status) {
63+
public Set<ItemTicket> getTicketsForStatus(ItemStatus<K, V, Ctx> status) {
7064
return this.status2Tickets[status.ordinal()];
7165
}
7266

7367
void clear() {
74-
for (Set<ItemTicket<K, V, Ctx>> tickets : status2Tickets) {
68+
for (Set<ItemTicket> tickets : status2Tickets) {
7569
tickets.clear();
7670
}
7771

7872
VarHandle.fullFence();
7973
}
8074

8175
void assertEmpty() {
82-
for (Set<ItemTicket<K, V, Ctx>> tickets : status2Tickets) {
76+
for (Set<ItemTicket> tickets : status2Tickets) {
8377
Assertions.assertTrue(tickets.isEmpty());
8478
}
8579
}

src/test/java/com/ishland/flowsched/scheduler/SchedulerTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,7 @@ public void testSimple() {
3434
Random random = new Random();
3535
while (spamLoaderRunning.get()) {
3636
long victim = random.nextLong(key - 1);
37-
ItemHolder<Long, TestItem, TestContext, Void> holder = scheduler.addTicket(victim, TestStatus.STATE_8, null);
37+
ItemHolder<Long, TestItem, TestContext, Void> holder = scheduler.addTicket(victim, TestStatus.STATE_8, (Runnable) null);
3838
CompletableFuture<Void> future = holder.getFutureForStatus0(TestStatus.STATE_8);
3939
if (future.isCompletedExceptionally()) {
4040
Assertions.fail();

0 commit comments

Comments
 (0)