Committed by
Gerrit Code Review
Merge "Fixing flow rule batches"
Showing
9 changed files
with
109 additions
and
102 deletions
... | @@ -63,7 +63,7 @@ public final class DefaultTrafficSelector implements TrafficSelector { | ... | @@ -63,7 +63,7 @@ public final class DefaultTrafficSelector implements TrafficSelector { |
63 | 63 | ||
64 | @Override | 64 | @Override |
65 | public int hashCode() { | 65 | public int hashCode() { |
66 | - return Objects.hash(criteria); | 66 | + return criteria.hashCode(); |
67 | } | 67 | } |
68 | 68 | ||
69 | @Override | 69 | @Override | ... | ... |
... | @@ -18,7 +18,7 @@ package org.onlab.onos.net.flow; | ... | @@ -18,7 +18,7 @@ package org.onlab.onos.net.flow; |
18 | import org.onlab.onos.core.ApplicationId; | 18 | import org.onlab.onos.core.ApplicationId; |
19 | import org.onlab.onos.net.provider.Provider; | 19 | import org.onlab.onos.net.provider.Provider; |
20 | 20 | ||
21 | -import com.google.common.util.concurrent.ListenableFuture; | 21 | +import java.util.concurrent.Future; |
22 | 22 | ||
23 | /** | 23 | /** |
24 | * Abstraction of a flow rule provider. | 24 | * Abstraction of a flow rule provider. |
... | @@ -58,6 +58,6 @@ public interface FlowRuleProvider extends Provider { | ... | @@ -58,6 +58,6 @@ public interface FlowRuleProvider extends Provider { |
58 | * @param batch a batch of flow rules | 58 | * @param batch a batch of flow rules |
59 | * @return a future indicating the status of this execution | 59 | * @return a future indicating the status of this execution |
60 | */ | 60 | */ |
61 | - ListenableFuture<CompletedBatchOperation> executeBatch(BatchOperation<FlowRuleBatchEntry> batch); | 61 | + Future<CompletedBatchOperation> executeBatch(BatchOperation<FlowRuleBatchEntry> batch); |
62 | 62 | ||
63 | } | 63 | } | ... | ... |
... | @@ -196,7 +196,7 @@ public final class Criteria { | ... | @@ -196,7 +196,7 @@ public final class Criteria { |
196 | 196 | ||
197 | @Override | 197 | @Override |
198 | public int hashCode() { | 198 | public int hashCode() { |
199 | - return Objects.hash(port, type()); | 199 | + return Objects.hash(type(), port); |
200 | } | 200 | } |
201 | 201 | ||
202 | @Override | 202 | @Override |
... | @@ -242,7 +242,7 @@ public final class Criteria { | ... | @@ -242,7 +242,7 @@ public final class Criteria { |
242 | 242 | ||
243 | @Override | 243 | @Override |
244 | public int hashCode() { | 244 | public int hashCode() { |
245 | - return Objects.hash(mac, type); | 245 | + return Objects.hash(type, mac); |
246 | } | 246 | } |
247 | 247 | ||
248 | @Override | 248 | @Override |
... | @@ -288,7 +288,7 @@ public final class Criteria { | ... | @@ -288,7 +288,7 @@ public final class Criteria { |
288 | 288 | ||
289 | @Override | 289 | @Override |
290 | public int hashCode() { | 290 | public int hashCode() { |
291 | - return Objects.hash(ethType, type()); | 291 | + return Objects.hash(type(), ethType); |
292 | } | 292 | } |
293 | 293 | ||
294 | @Override | 294 | @Override |
... | @@ -336,7 +336,7 @@ public final class Criteria { | ... | @@ -336,7 +336,7 @@ public final class Criteria { |
336 | 336 | ||
337 | @Override | 337 | @Override |
338 | public int hashCode() { | 338 | public int hashCode() { |
339 | - return Objects.hash(ip, type); | 339 | + return Objects.hash(type, ip); |
340 | } | 340 | } |
341 | 341 | ||
342 | @Override | 342 | @Override |
... | @@ -382,7 +382,7 @@ public final class Criteria { | ... | @@ -382,7 +382,7 @@ public final class Criteria { |
382 | 382 | ||
383 | @Override | 383 | @Override |
384 | public int hashCode() { | 384 | public int hashCode() { |
385 | - return Objects.hash(proto, type()); | 385 | + return Objects.hash(type(), proto); |
386 | } | 386 | } |
387 | 387 | ||
388 | @Override | 388 | @Override |
... | @@ -427,7 +427,7 @@ public final class Criteria { | ... | @@ -427,7 +427,7 @@ public final class Criteria { |
427 | 427 | ||
428 | @Override | 428 | @Override |
429 | public int hashCode() { | 429 | public int hashCode() { |
430 | - return Objects.hash(vlanPcp); | 430 | + return Objects.hash(type(), vlanPcp); |
431 | } | 431 | } |
432 | 432 | ||
433 | @Override | 433 | @Override |
... | @@ -474,7 +474,7 @@ public final class Criteria { | ... | @@ -474,7 +474,7 @@ public final class Criteria { |
474 | 474 | ||
475 | @Override | 475 | @Override |
476 | public int hashCode() { | 476 | public int hashCode() { |
477 | - return Objects.hash(vlanId, type()); | 477 | + return Objects.hash(type(), vlanId); |
478 | } | 478 | } |
479 | 479 | ||
480 | @Override | 480 | @Override |
... | @@ -522,7 +522,7 @@ public final class Criteria { | ... | @@ -522,7 +522,7 @@ public final class Criteria { |
522 | 522 | ||
523 | @Override | 523 | @Override |
524 | public int hashCode() { | 524 | public int hashCode() { |
525 | - return Objects.hash(tcpPort, type); | 525 | + return Objects.hash(type, tcpPort); |
526 | } | 526 | } |
527 | 527 | ||
528 | @Override | 528 | @Override |
... | @@ -568,7 +568,7 @@ public final class Criteria { | ... | @@ -568,7 +568,7 @@ public final class Criteria { |
568 | 568 | ||
569 | @Override | 569 | @Override |
570 | public int hashCode() { | 570 | public int hashCode() { |
571 | - return Objects.hash(lambda, type); | 571 | + return Objects.hash(type, lambda); |
572 | } | 572 | } |
573 | 573 | ||
574 | @Override | 574 | @Override |
... | @@ -612,7 +612,7 @@ public final class Criteria { | ... | @@ -612,7 +612,7 @@ public final class Criteria { |
612 | 612 | ||
613 | @Override | 613 | @Override |
614 | public int hashCode() { | 614 | public int hashCode() { |
615 | - return Objects.hash(signalType, type); | 615 | + return Objects.hash(type, signalType); |
616 | } | 616 | } |
617 | 617 | ||
618 | @Override | 618 | @Override | ... | ... |
... | @@ -190,7 +190,7 @@ public final class Instructions { | ... | @@ -190,7 +190,7 @@ public final class Instructions { |
190 | 190 | ||
191 | @Override | 191 | @Override |
192 | public int hashCode() { | 192 | public int hashCode() { |
193 | - return Objects.hash(port, type()); | 193 | + return Objects.hash(type(), port); |
194 | } | 194 | } |
195 | 195 | ||
196 | @Override | 196 | @Override | ... | ... |
... | @@ -70,7 +70,7 @@ public abstract class L0ModificationInstruction implements Instruction { | ... | @@ -70,7 +70,7 @@ public abstract class L0ModificationInstruction implements Instruction { |
70 | 70 | ||
71 | @Override | 71 | @Override |
72 | public int hashCode() { | 72 | public int hashCode() { |
73 | - return Objects.hash(lambda, type(), subtype); | 73 | + return Objects.hash(type(), subtype, lambda); |
74 | } | 74 | } |
75 | 75 | ||
76 | @Override | 76 | @Override | ... | ... |
... | @@ -93,7 +93,7 @@ public abstract class L2ModificationInstruction implements Instruction { | ... | @@ -93,7 +93,7 @@ public abstract class L2ModificationInstruction implements Instruction { |
93 | 93 | ||
94 | @Override | 94 | @Override |
95 | public int hashCode() { | 95 | public int hashCode() { |
96 | - return Objects.hash(mac, type(), subtype); | 96 | + return Objects.hash(type(), subtype, mac); |
97 | } | 97 | } |
98 | 98 | ||
99 | @Override | 99 | @Override |
... | @@ -142,7 +142,7 @@ public abstract class L2ModificationInstruction implements Instruction { | ... | @@ -142,7 +142,7 @@ public abstract class L2ModificationInstruction implements Instruction { |
142 | 142 | ||
143 | @Override | 143 | @Override |
144 | public int hashCode() { | 144 | public int hashCode() { |
145 | - return Objects.hash(vlanId, type(), subtype()); | 145 | + return Objects.hash(type(), subtype(), vlanId); |
146 | } | 146 | } |
147 | 147 | ||
148 | @Override | 148 | @Override |
... | @@ -191,7 +191,7 @@ public abstract class L2ModificationInstruction implements Instruction { | ... | @@ -191,7 +191,7 @@ public abstract class L2ModificationInstruction implements Instruction { |
191 | 191 | ||
192 | @Override | 192 | @Override |
193 | public int hashCode() { | 193 | public int hashCode() { |
194 | - return Objects.hash(vlanPcp, type(), subtype()); | 194 | + return Objects.hash(type(), subtype(), vlanPcp); |
195 | } | 195 | } |
196 | 196 | ||
197 | @Override | 197 | @Override | ... | ... |
... | @@ -85,7 +85,7 @@ public abstract class L3ModificationInstruction implements Instruction { | ... | @@ -85,7 +85,7 @@ public abstract class L3ModificationInstruction implements Instruction { |
85 | 85 | ||
86 | @Override | 86 | @Override |
87 | public int hashCode() { | 87 | public int hashCode() { |
88 | - return Objects.hash(ip, type(), subtype()); | 88 | + return Objects.hash(type(), subtype(), ip); |
89 | } | 89 | } |
90 | 90 | ||
91 | @Override | 91 | @Override | ... | ... |
... | @@ -15,22 +15,12 @@ | ... | @@ -15,22 +15,12 @@ |
15 | */ | 15 | */ |
16 | package org.onlab.onos.net.flow.impl; | 16 | package org.onlab.onos.net.flow.impl; |
17 | 17 | ||
18 | -import static com.google.common.base.Preconditions.checkNotNull; | 18 | +import com.google.common.collect.ArrayListMultimap; |
19 | -import static org.slf4j.LoggerFactory.getLogger; | 19 | +import com.google.common.collect.Iterables; |
20 | -import static org.onlab.util.Tools.namedThreads; | 20 | +import com.google.common.collect.Lists; |
21 | - | 21 | +import com.google.common.collect.Maps; |
22 | -import java.util.List; | 22 | +import com.google.common.collect.Multimap; |
23 | -import java.util.Map; | 23 | +import com.google.common.collect.Sets; |
24 | -import java.util.Set; | ||
25 | -import java.util.concurrent.CancellationException; | ||
26 | -import java.util.concurrent.ExecutionException; | ||
27 | -import java.util.concurrent.ExecutorService; | ||
28 | -import java.util.concurrent.Executors; | ||
29 | -import java.util.concurrent.Future; | ||
30 | -import java.util.concurrent.TimeUnit; | ||
31 | -import java.util.concurrent.TimeoutException; | ||
32 | -import java.util.concurrent.atomic.AtomicReference; | ||
33 | - | ||
34 | import org.apache.felix.scr.annotations.Activate; | 24 | import org.apache.felix.scr.annotations.Activate; |
35 | import org.apache.felix.scr.annotations.Component; | 25 | import org.apache.felix.scr.annotations.Component; |
36 | import org.apache.felix.scr.annotations.Deactivate; | 26 | import org.apache.felix.scr.annotations.Deactivate; |
... | @@ -64,14 +54,21 @@ import org.onlab.onos.net.provider.AbstractProviderRegistry; | ... | @@ -64,14 +54,21 @@ import org.onlab.onos.net.provider.AbstractProviderRegistry; |
64 | import org.onlab.onos.net.provider.AbstractProviderService; | 54 | import org.onlab.onos.net.provider.AbstractProviderService; |
65 | import org.slf4j.Logger; | 55 | import org.slf4j.Logger; |
66 | 56 | ||
67 | -import com.google.common.collect.ArrayListMultimap; | 57 | +import java.util.List; |
68 | -import com.google.common.collect.Iterables; | 58 | +import java.util.Map; |
69 | -import com.google.common.collect.Lists; | 59 | +import java.util.Set; |
70 | -import com.google.common.collect.Maps; | 60 | +import java.util.concurrent.CancellationException; |
71 | -import com.google.common.collect.Multimap; | 61 | +import java.util.concurrent.ExecutionException; |
72 | -import com.google.common.collect.Sets; | 62 | +import java.util.concurrent.ExecutorService; |
73 | -import com.google.common.util.concurrent.Futures; | 63 | +import java.util.concurrent.Executors; |
74 | -import com.google.common.util.concurrent.ListenableFuture; | 64 | +import java.util.concurrent.Future; |
65 | +import java.util.concurrent.TimeUnit; | ||
66 | +import java.util.concurrent.TimeoutException; | ||
67 | +import java.util.concurrent.atomic.AtomicReference; | ||
68 | + | ||
69 | +import static com.google.common.base.Preconditions.checkNotNull; | ||
70 | +import static org.onlab.util.Tools.namedThreads; | ||
71 | +import static org.slf4j.LoggerFactory.getLogger; | ||
75 | 72 | ||
76 | /** | 73 | /** |
77 | * Provides implementation of the flow NB & SB APIs. | 74 | * Provides implementation of the flow NB & SB APIs. |
... | @@ -92,8 +89,7 @@ public class FlowRuleManager | ... | @@ -92,8 +89,7 @@ public class FlowRuleManager |
92 | 89 | ||
93 | private final FlowRuleStoreDelegate delegate = new InternalStoreDelegate(); | 90 | private final FlowRuleStoreDelegate delegate = new InternalStoreDelegate(); |
94 | 91 | ||
95 | - private final ExecutorService futureListeners = | 92 | + private ExecutorService futureService; |
96 | - Executors.newCachedThreadPool(namedThreads("provider-future-listeners")); | ||
97 | 93 | ||
98 | @Reference(cardinality = ReferenceCardinality.MANDATORY_UNARY) | 94 | @Reference(cardinality = ReferenceCardinality.MANDATORY_UNARY) |
99 | protected FlowRuleStore store; | 95 | protected FlowRuleStore store; |
... | @@ -106,6 +102,7 @@ public class FlowRuleManager | ... | @@ -106,6 +102,7 @@ public class FlowRuleManager |
106 | 102 | ||
107 | @Activate | 103 | @Activate |
108 | public void activate() { | 104 | public void activate() { |
105 | + futureService = Executors.newCachedThreadPool(namedThreads("provider-future-listeners")); | ||
109 | store.setDelegate(delegate); | 106 | store.setDelegate(delegate); |
110 | eventDispatcher.addSink(FlowRuleEvent.class, listenerRegistry); | 107 | eventDispatcher.addSink(FlowRuleEvent.class, listenerRegistry); |
111 | log.info("Started"); | 108 | log.info("Started"); |
... | @@ -113,7 +110,7 @@ public class FlowRuleManager | ... | @@ -113,7 +110,7 @@ public class FlowRuleManager |
113 | 110 | ||
114 | @Deactivate | 111 | @Deactivate |
115 | public void deactivate() { | 112 | public void deactivate() { |
116 | - futureListeners.shutdownNow(); | 113 | + futureService.shutdownNow(); |
117 | 114 | ||
118 | store.unsetDelegate(delegate); | 115 | store.unsetDelegate(delegate); |
119 | eventDispatcher.removeSink(FlowRuleEvent.class); | 116 | eventDispatcher.removeSink(FlowRuleEvent.class); |
... | @@ -364,6 +361,9 @@ public class FlowRuleManager | ... | @@ -364,6 +361,9 @@ public class FlowRuleManager |
364 | 361 | ||
365 | // Store delegate to re-post events emitted from the store. | 362 | // Store delegate to re-post events emitted from the store. |
366 | private class InternalStoreDelegate implements FlowRuleStoreDelegate { | 363 | private class InternalStoreDelegate implements FlowRuleStoreDelegate { |
364 | + | ||
365 | + private static final int TIMEOUT = 5000; // ms | ||
366 | + | ||
367 | // TODO: Right now we only dispatch events at individual flowEntry level. | 367 | // TODO: Right now we only dispatch events at individual flowEntry level. |
368 | // It may be more efficient for also dispatch events as a batch. | 368 | // It may be more efficient for also dispatch events as a batch. |
369 | @Override | 369 | @Override |
... | @@ -384,15 +384,21 @@ public class FlowRuleManager | ... | @@ -384,15 +384,21 @@ public class FlowRuleManager |
384 | 384 | ||
385 | FlowRuleProvider flowRuleProvider = | 385 | FlowRuleProvider flowRuleProvider = |
386 | getProvider(batchOperation.getOperations().get(0).getTarget().deviceId()); | 386 | getProvider(batchOperation.getOperations().get(0).getTarget().deviceId()); |
387 | - final ListenableFuture<CompletedBatchOperation> result = | 387 | + final Future<CompletedBatchOperation> result = |
388 | flowRuleProvider.executeBatch(batchOperation); | 388 | flowRuleProvider.executeBatch(batchOperation); |
389 | - result.addListener(new Runnable() { | 389 | + futureService.submit(new Runnable() { |
390 | @Override | 390 | @Override |
391 | public void run() { | 391 | public void run() { |
392 | - store.batchOperationComplete(FlowRuleBatchEvent.completed(request, | 392 | + CompletedBatchOperation res = null; |
393 | - Futures.getUnchecked(result))); | 393 | + try { |
394 | + res = result.get(TIMEOUT, TimeUnit.MILLISECONDS); | ||
395 | + } catch (TimeoutException | InterruptedException | ExecutionException e) { | ||
396 | + log.warn("Something went wrong with the batch operation {}", | ||
397 | + request.batchId()); | ||
398 | + } | ||
399 | + store.batchOperationComplete(FlowRuleBatchEvent.completed(request, res)); | ||
394 | } | 400 | } |
395 | - }, futureListeners); | 401 | + }); |
396 | break; | 402 | break; |
397 | 403 | ||
398 | case BATCH_OPERATION_COMPLETED: | 404 | case BATCH_OPERATION_COMPLETED: | ... | ... |
providers/openflow/flow/src/main/java/org/onlab/onos/provider/of/flow/impl/OpenFlowRuleProvider.java
... | @@ -15,21 +15,11 @@ | ... | @@ -15,21 +15,11 @@ |
15 | */ | 15 | */ |
16 | package org.onlab.onos.provider.of.flow.impl; | 16 | package org.onlab.onos.provider.of.flow.impl; |
17 | 17 | ||
18 | -import static org.slf4j.LoggerFactory.getLogger; | 18 | +import com.google.common.collect.ArrayListMultimap; |
19 | - | 19 | +import com.google.common.collect.Maps; |
20 | -import java.util.Collections; | 20 | +import com.google.common.collect.Multimap; |
21 | -import java.util.HashMap; | 21 | +import com.google.common.collect.Sets; |
22 | -import java.util.List; | 22 | +import com.google.common.util.concurrent.ExecutionList; |
23 | -import java.util.Map; | ||
24 | -import java.util.Set; | ||
25 | -import java.util.concurrent.ConcurrentHashMap; | ||
26 | -import java.util.concurrent.CountDownLatch; | ||
27 | -import java.util.concurrent.ExecutionException; | ||
28 | -import java.util.concurrent.Executor; | ||
29 | -import java.util.concurrent.TimeUnit; | ||
30 | -import java.util.concurrent.TimeoutException; | ||
31 | -import java.util.concurrent.atomic.AtomicBoolean; | ||
32 | - | ||
33 | import org.apache.felix.scr.annotations.Activate; | 23 | import org.apache.felix.scr.annotations.Activate; |
34 | import org.apache.felix.scr.annotations.Component; | 24 | import org.apache.felix.scr.annotations.Component; |
35 | import org.apache.felix.scr.annotations.Deactivate; | 25 | import org.apache.felix.scr.annotations.Deactivate; |
... | @@ -80,15 +70,23 @@ import org.projectfloodlight.openflow.protocol.errormsg.OFFlowModFailedErrorMsg; | ... | @@ -80,15 +70,23 @@ import org.projectfloodlight.openflow.protocol.errormsg.OFFlowModFailedErrorMsg; |
80 | import org.projectfloodlight.openflow.protocol.instruction.OFInstruction; | 70 | import org.projectfloodlight.openflow.protocol.instruction.OFInstruction; |
81 | import org.projectfloodlight.openflow.protocol.instruction.OFInstructionApplyActions; | 71 | import org.projectfloodlight.openflow.protocol.instruction.OFInstructionApplyActions; |
82 | import org.projectfloodlight.openflow.types.OFPort; | 72 | import org.projectfloodlight.openflow.types.OFPort; |
83 | -import org.projectfloodlight.openflow.types.U32; | ||
84 | import org.slf4j.Logger; | 73 | import org.slf4j.Logger; |
85 | 74 | ||
86 | -import com.google.common.collect.ArrayListMultimap; | 75 | +import java.util.Collections; |
87 | -import com.google.common.collect.Maps; | 76 | +import java.util.HashMap; |
88 | -import com.google.common.collect.Multimap; | 77 | +import java.util.List; |
89 | -import com.google.common.collect.Sets; | 78 | +import java.util.Map; |
90 | -import com.google.common.util.concurrent.ExecutionList; | 79 | +import java.util.Set; |
91 | -import com.google.common.util.concurrent.ListenableFuture; | 80 | +import java.util.concurrent.ConcurrentHashMap; |
81 | +import java.util.concurrent.CountDownLatch; | ||
82 | +import java.util.concurrent.ExecutionException; | ||
83 | +import java.util.concurrent.Future; | ||
84 | +import java.util.concurrent.TimeUnit; | ||
85 | +import java.util.concurrent.TimeoutException; | ||
86 | +import java.util.concurrent.atomic.AtomicBoolean; | ||
87 | +import java.util.concurrent.atomic.AtomicLong; | ||
88 | + | ||
89 | +import static org.slf4j.LoggerFactory.getLogger; | ||
92 | 90 | ||
93 | /** | 91 | /** |
94 | * Provider which uses an OpenFlow controller to detect network | 92 | * Provider which uses an OpenFlow controller to detect network |
... | @@ -124,6 +122,8 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -124,6 +122,8 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
124 | 122 | ||
125 | private final Map<Dpid, FlowStatsCollector> collectors = Maps.newHashMap(); | 123 | private final Map<Dpid, FlowStatsCollector> collectors = Maps.newHashMap(); |
126 | 124 | ||
125 | + private final AtomicLong xidCounter = new AtomicLong(0); | ||
126 | + | ||
127 | /** | 127 | /** |
128 | * Creates an OpenFlow host provider. | 128 | * Creates an OpenFlow host provider. |
129 | */ | 129 | */ |
... | @@ -154,6 +154,7 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -154,6 +154,7 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
154 | 154 | ||
155 | log.info("Stopped"); | 155 | log.info("Stopped"); |
156 | } | 156 | } |
157 | + | ||
157 | @Override | 158 | @Override |
158 | public void applyFlowRule(FlowRule... flowRules) { | 159 | public void applyFlowRule(FlowRule... flowRules) { |
159 | for (int i = 0; i < flowRules.length; i++) { | 160 | for (int i = 0; i < flowRules.length; i++) { |
... | @@ -167,7 +168,6 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -167,7 +168,6 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
167 | } | 168 | } |
168 | 169 | ||
169 | 170 | ||
170 | - | ||
171 | @Override | 171 | @Override |
172 | public void removeFlowRule(FlowRule... flowRules) { | 172 | public void removeFlowRule(FlowRule... flowRules) { |
173 | for (int i = 0; i < flowRules.length; i++) { | 173 | for (int i = 0; i < flowRules.length; i++) { |
... | @@ -188,11 +188,15 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -188,11 +188,15 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
188 | } | 188 | } |
189 | 189 | ||
190 | @Override | 190 | @Override |
191 | - public ListenableFuture<CompletedBatchOperation> executeBatch(BatchOperation<FlowRuleBatchEntry> batch) { | 191 | + public Future<CompletedBatchOperation> executeBatch(BatchOperation<FlowRuleBatchEntry> batch) { |
192 | final Set<Dpid> sws = | 192 | final Set<Dpid> sws = |
193 | Collections.newSetFromMap(new ConcurrentHashMap<Dpid, Boolean>()); | 193 | Collections.newSetFromMap(new ConcurrentHashMap<Dpid, Boolean>()); |
194 | final Map<Long, FlowRuleBatchEntry> fmXids = new HashMap<Long, FlowRuleBatchEntry>(); | 194 | final Map<Long, FlowRuleBatchEntry> fmXids = new HashMap<Long, FlowRuleBatchEntry>(); |
195 | - OFFlowMod mod = null; | 195 | + /* |
196 | + * Use identity hash map for reference equality as we could have equal | ||
197 | + * flow mods for different switches. | ||
198 | + */ | ||
199 | + Map<OFFlowMod, OpenFlowSwitch> mods = Maps.newIdentityHashMap(); | ||
196 | for (FlowRuleBatchEntry fbe : batch.getOperations()) { | 200 | for (FlowRuleBatchEntry fbe : batch.getOperations()) { |
197 | FlowRule flowRule = fbe.getTarget(); | 201 | FlowRule flowRule = fbe.getTarget(); |
198 | OpenFlowSwitch sw = controller.getSwitch(Dpid.dpid(flowRule.deviceId().uri())); | 202 | OpenFlowSwitch sw = controller.getSwitch(Dpid.dpid(flowRule.deviceId().uri())); |
... | @@ -208,6 +212,7 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -208,6 +212,7 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
208 | } | 212 | } |
209 | sws.add(new Dpid(sw.getId())); | 213 | sws.add(new Dpid(sw.getId())); |
210 | FlowModBuilder builder = FlowModBuilder.builder(flowRule, sw.factory()); | 214 | FlowModBuilder builder = FlowModBuilder.builder(flowRule, sw.factory()); |
215 | + OFFlowMod mod = null; | ||
211 | switch (fbe.getOperator()) { | 216 | switch (fbe.getOperator()) { |
212 | case ADD: | 217 | case ADD: |
213 | mod = builder.buildFlowAdd(); | 218 | mod = builder.buildFlowAdd(); |
... | @@ -222,19 +227,23 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -222,19 +227,23 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
222 | log.error("Unsupported batch operation {}", fbe.getOperator()); | 227 | log.error("Unsupported batch operation {}", fbe.getOperator()); |
223 | } | 228 | } |
224 | if (mod != null) { | 229 | if (mod != null) { |
225 | - sw.sendMsg(mod); | 230 | + mods.put(mod, sw); |
226 | - fmXids.put(mod.getXid(), fbe); | 231 | + fmXids.put(xidCounter.getAndIncrement(), fbe); |
227 | } else { | 232 | } else { |
228 | log.error("Conversion of flowrule {} failed.", flowRule); | 233 | log.error("Conversion of flowrule {} failed.", flowRule); |
229 | } | 234 | } |
230 | - | ||
231 | } | 235 | } |
232 | InstallationFuture installation = new InstallationFuture(sws, fmXids); | 236 | InstallationFuture installation = new InstallationFuture(sws, fmXids); |
233 | for (Long xid : fmXids.keySet()) { | 237 | for (Long xid : fmXids.keySet()) { |
234 | pendingFMs.put(xid, installation); | 238 | pendingFMs.put(xid, installation); |
235 | } | 239 | } |
236 | - pendingFutures.put(U32.f(batch.hashCode()), installation); | 240 | + pendingFutures.put(installation.xid(), installation); |
237 | - installation.verify(U32.f(batch.hashCode())); | 241 | + for (Map.Entry<OFFlowMod, OpenFlowSwitch> entry : mods.entrySet()) { |
242 | + OpenFlowSwitch sw = entry.getValue(); | ||
243 | + OFFlowMod mod = entry.getKey(); | ||
244 | + sw.sendMsg(mod); | ||
245 | + } | ||
246 | + installation.verify(); | ||
238 | return installation; | 247 | return installation; |
239 | } | 248 | } |
240 | 249 | ||
... | @@ -352,8 +361,9 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -352,8 +361,9 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
352 | 361 | ||
353 | } | 362 | } |
354 | 363 | ||
355 | - private class InstallationFuture implements ListenableFuture<CompletedBatchOperation> { | 364 | + private class InstallationFuture implements Future<CompletedBatchOperation> { |
356 | 365 | ||
366 | + private final Long xid; | ||
357 | private final Set<Dpid> sws; | 367 | private final Set<Dpid> sws; |
358 | private final AtomicBoolean ok = new AtomicBoolean(true); | 368 | private final AtomicBoolean ok = new AtomicBoolean(true); |
359 | private final Map<Long, FlowRuleBatchEntry> fms; | 369 | private final Map<Long, FlowRuleBatchEntry> fms; |
... | @@ -361,18 +371,22 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -361,18 +371,22 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
361 | private final Set<FlowEntry> offendingFlowMods = Sets.newHashSet(); | 371 | private final Set<FlowEntry> offendingFlowMods = Sets.newHashSet(); |
362 | 372 | ||
363 | private final CountDownLatch countDownLatch; | 373 | private final CountDownLatch countDownLatch; |
364 | - private Long pendingXid; | ||
365 | private BatchState state; | 374 | private BatchState state; |
366 | 375 | ||
367 | private final ExecutionList executionList = new ExecutionList(); | 376 | private final ExecutionList executionList = new ExecutionList(); |
368 | 377 | ||
369 | public InstallationFuture(Set<Dpid> sws, Map<Long, FlowRuleBatchEntry> fmXids) { | 378 | public InstallationFuture(Set<Dpid> sws, Map<Long, FlowRuleBatchEntry> fmXids) { |
379 | + this.xid = xidCounter.getAndIncrement(); | ||
370 | this.state = BatchState.STARTED; | 380 | this.state = BatchState.STARTED; |
371 | this.sws = sws; | 381 | this.sws = sws; |
372 | this.fms = fmXids; | 382 | this.fms = fmXids; |
373 | countDownLatch = new CountDownLatch(sws.size()); | 383 | countDownLatch = new CountDownLatch(sws.size()); |
374 | } | 384 | } |
375 | 385 | ||
386 | + public Long xid() { | ||
387 | + return xid; | ||
388 | + } | ||
389 | + | ||
376 | public void fail(OFErrorMsg msg, Dpid dpid) { | 390 | public void fail(OFErrorMsg msg, Dpid dpid) { |
377 | 391 | ||
378 | ok.set(false); | 392 | ok.set(false); |
... | @@ -434,13 +448,12 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -434,13 +448,12 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
434 | } | 448 | } |
435 | 449 | ||
436 | 450 | ||
437 | - public void verify(Long id) { | 451 | + public void verify() { |
438 | - pendingXid = id; | ||
439 | for (Dpid dpid : sws) { | 452 | for (Dpid dpid : sws) { |
440 | OpenFlowSwitch sw = controller.getSwitch(dpid); | 453 | OpenFlowSwitch sw = controller.getSwitch(dpid); |
441 | OFBarrierRequest.Builder builder = sw.factory() | 454 | OFBarrierRequest.Builder builder = sw.factory() |
442 | .buildBarrierRequest() | 455 | .buildBarrierRequest() |
443 | - .setXid(id); | 456 | + .setXid(xid); |
444 | sw.sendMsg(builder.build()); | 457 | sw.sendMsg(builder.build()); |
445 | } | 458 | } |
446 | } | 459 | } |
... | @@ -462,7 +475,6 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -462,7 +475,6 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
462 | } | 475 | } |
463 | 476 | ||
464 | } | 477 | } |
465 | - invokeCallbacks(); | ||
466 | return true; | 478 | return true; |
467 | } | 479 | } |
468 | 480 | ||
... | @@ -481,6 +493,7 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -481,6 +493,7 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
481 | countDownLatch.await(); | 493 | countDownLatch.await(); |
482 | this.state = BatchState.FINISHED; | 494 | this.state = BatchState.FINISHED; |
483 | CompletedBatchOperation result = new CompletedBatchOperation(ok.get(), offendingFlowMods); | 495 | CompletedBatchOperation result = new CompletedBatchOperation(ok.get(), offendingFlowMods); |
496 | + //FIXME do cleanup here | ||
484 | return result; | 497 | return result; |
485 | } | 498 | } |
486 | 499 | ||
... | @@ -491,6 +504,7 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -491,6 +504,7 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
491 | if (countDownLatch.await(timeout, unit)) { | 504 | if (countDownLatch.await(timeout, unit)) { |
492 | this.state = BatchState.FINISHED; | 505 | this.state = BatchState.FINISHED; |
493 | CompletedBatchOperation result = new CompletedBatchOperation(ok.get(), offendingFlowMods); | 506 | CompletedBatchOperation result = new CompletedBatchOperation(ok.get(), offendingFlowMods); |
507 | + // FIXME do cleanup here | ||
494 | return result; | 508 | return result; |
495 | } | 509 | } |
496 | throw new TimeoutException(); | 510 | throw new TimeoutException(); |
... | @@ -498,9 +512,7 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -498,9 +512,7 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
498 | 512 | ||
499 | private void cleanUp() { | 513 | private void cleanUp() { |
500 | if (isDone() || isCancelled()) { | 514 | if (isDone() || isCancelled()) { |
501 | - if (pendingXid != null) { | 515 | + pendingFutures.remove(xid); |
502 | - pendingFutures.remove(pendingXid); | ||
503 | - } | ||
504 | for (Long xid : fms.keySet()) { | 516 | for (Long xid : fms.keySet()) { |
505 | pendingFMs.remove(xid); | 517 | pendingFMs.remove(xid); |
506 | } | 518 | } |
... | @@ -509,21 +521,10 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr | ... | @@ -509,21 +521,10 @@ public class OpenFlowRuleProvider extends AbstractProvider implements FlowRulePr |
509 | 521 | ||
510 | private void removeRequirement(Dpid dpid) { | 522 | private void removeRequirement(Dpid dpid) { |
511 | countDownLatch.countDown(); | 523 | countDownLatch.countDown(); |
512 | - if (countDownLatch.getCount() == 0) { | ||
513 | - invokeCallbacks(); | ||
514 | - } | ||
515 | sws.remove(dpid); | 524 | sws.remove(dpid); |
525 | + //FIXME don't do cleanup here | ||
516 | cleanUp(); | 526 | cleanUp(); |
517 | } | 527 | } |
518 | - | ||
519 | - @Override | ||
520 | - public void addListener(Runnable runnable, Executor executor) { | ||
521 | - executionList.add(runnable, executor); | ||
522 | - } | ||
523 | - | ||
524 | - private void invokeCallbacks() { | ||
525 | - executionList.execute(); | ||
526 | - } | ||
527 | } | 528 | } |
528 | 529 | ||
529 | } | 530 | } | ... | ... |
-
Please register or login to post a comment