package org.opendaylight.controller.clustering.it.provider;
import akka.actor.ActorRef;
-import akka.actor.ActorSystem;
+import akka.dispatch.Futures;
import akka.dispatch.OnComplete;
import akka.pattern.Patterns;
import com.google.common.base.Strings;
-import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.SettableFuture;
import java.util.HashMap;
import org.opendaylight.yangtools.concepts.ListenerRegistration;
import org.opendaylight.yangtools.concepts.ObjectRegistration;
import org.opendaylight.yangtools.yang.binding.InstanceIdentifier;
-import org.opendaylight.yangtools.yang.common.RpcError.ErrorType;
+import org.opendaylight.yangtools.yang.common.ErrorTag;
+import org.opendaylight.yangtools.yang.common.ErrorType;
import org.opendaylight.yangtools.yang.common.RpcResult;
import org.opendaylight.yangtools.yang.common.RpcResultBuilder;
import org.opendaylight.yangtools.yang.data.api.schema.NormalizedNode;
import org.slf4j.LoggerFactory;
import scala.concurrent.duration.FiniteDuration;
-public class MdsalLowLevelTestProvider implements OdlMdsalLowlevelControlService {
+public final class MdsalLowLevelTestProvider implements OdlMdsalLowlevelControlService {
private static final Logger LOG = LoggerFactory.getLogger(MdsalLowLevelTestProvider.class);
- private final RpcProviderService rpcRegistry;
private final ObjectRegistration<OdlMdsalLowlevelControlService> registration;
private final DistributedDataStoreInterface configDataStore;
private final BindingNormalizedNodeSerializer bindingNormalizedNodeSerializer;
private final DOMDataBroker domDataBroker;
private final NotificationPublishService notificationPublishService;
private final NotificationService notificationService;
- private final DOMSchemaService schemaService;
private final ClusterSingletonServiceProvider singletonService;
private final DOMRpcProviderService domRpcService;
private final DOMDataTreeChangeService domDataTreeChangeService;
- private final ActorSystem actorSystem;
private final Map<InstanceIdentifier<?>, DOMRpcImplementationRegistration<RoutedGetConstantService>>
routedRegistrations = new HashMap<>();
private IdIntsListener idIntsListener;
private final Map<String, PublishNotificationsTask> publishNotificationsTasks = new HashMap<>();
- public MdsalLowLevelTestProvider(final RpcProviderService rpcRegistry,
+ public MdsalLowLevelTestProvider(
+ // FIXME: do not depend on this service
+ final RpcProviderService rpcRegistry,
final DOMRpcProviderService domRpcService,
final ClusterSingletonServiceProvider singletonService,
+ // FIXME: do not depend on this service
final DOMSchemaService schemaService,
final BindingNormalizedNodeSerializer bindingNormalizedNodeSerializer,
final NotificationPublishService notificationPublishService,
final NotificationService notificationService,
final DOMDataBroker domDataBroker,
final DistributedDataStoreInterface configDataStore,
+ // FIXME: do not depend on this service
final ActorSystemProvider actorSystemProvider) {
- this.rpcRegistry = rpcRegistry;
this.domRpcService = domRpcService;
this.singletonService = singletonService;
- this.schemaService = schemaService;
this.bindingNormalizedNodeSerializer = bindingNormalizedNodeSerializer;
this.notificationPublishService = notificationPublishService;
this.notificationService = notificationService;
this.domDataBroker = domDataBroker;
this.configDataStore = configDataStore;
- this.actorSystem = actorSystemProvider.getActorSystem();
domDataTreeChangeService = domDataBroker.getExtensions().getInstance(DOMDataTreeChangeService.class);
LOG.info("In unregisterSingletonConstant");
if (getSingletonConstantRegistration == null) {
- return RpcResultBuilder.<UnregisterSingletonConstantOutput>failed().withError(ErrorType.RPC, "data-missing",
- "No prior RPC was registered").buildFuture();
+ return RpcResultBuilder.<UnregisterSingletonConstantOutput>failed()
+ .withError(ErrorType.RPC, ErrorTag.DATA_MISSING, "No prior RPC was registered")
+ .buildFuture();
}
try {
LOG.info("In subscribeDtcl - input: {}", input);
if (dtclReg != null) {
- return RpcResultBuilder.<SubscribeDtclOutput>failed().withError(ErrorType.RPC,
- "data-exists", "There is already a DataTreeChangeListener registered for id-ints").buildFuture();
+ return RpcResultBuilder.<SubscribeDtclOutput>failed().withError(ErrorType.RPC, ErrorTag.DATA_EXISTS,
+ "There is already a DataTreeChangeListener registered for id-ints")
+ .buildFuture();
}
idIntsListener = new IdIntsListener();
LOG.info("In subscribeYnl - input: {}", input);
if (ynlRegistrations.containsKey(input.getId())) {
- return RpcResultBuilder.<SubscribeYnlOutput>failed().withError(ErrorType.RPC,
- "data-exists", "There is already a listener registered for id: " + input.getId()).buildFuture();
+ return RpcResultBuilder.<SubscribeYnlOutput>failed()
+ .withError(ErrorType.RPC, ErrorTag.DATA_EXISTS,
+ "There is already a listener registered for id: " + input.getId())
+ .buildFuture();
}
ynlRegistrations.put(input.getId(),
routedRegistrations.remove(input.getContext());
if (rpcRegistration == null) {
- return RpcResultBuilder.<UnregisterBoundConstantOutput>failed().withError(
- ErrorType.RPC, "data-missing", "No prior RPC was registered for " + input.getContext()).buildFuture();
+ return RpcResultBuilder.<UnregisterBoundConstantOutput>failed()
+ .withError(ErrorType.RPC, ErrorTag.DATA_MISSING,
+ "No prior RPC was registered for " + input.getContext())
+ .buildFuture();
}
rpcRegistration.close();
LOG.info("In registerSingletonConstant - input: {}", input);
if (input.getConstant() == null) {
- return RpcResultBuilder.<RegisterSingletonConstantOutput>failed().withError(
- ErrorType.RPC, "invalid-value", "Constant value is null").buildFuture();
+ return RpcResultBuilder.<RegisterSingletonConstantOutput>failed()
+ .withError(ErrorType.RPC, ErrorTag.INVALID_VALUE, "Constant value is null")
+ .buildFuture();
}
getSingletonConstantRegistration =
LOG.info("In unregisterConstant");
if (globalGetConstantRegistration == null) {
- return RpcResultBuilder.<UnregisterConstantOutput>failed().withError(
- ErrorType.RPC, "data-missing", "No prior RPC was registered").buildFuture();
+ return RpcResultBuilder.<UnregisterConstantOutput>failed()
+ .withError(ErrorType.RPC, ErrorTag.DATA_MISSING, "No prior RPC was registered")
+ .buildFuture();
}
globalGetConstantRegistration.close();
globalGetConstantRegistration = null;
- return Futures.immediateFuture(RpcResultBuilder.success(new UnregisterConstantOutputBuilder().build()).build());
+ return RpcResultBuilder.success(new UnregisterConstantOutputBuilder().build()).buildFuture();
}
@Override
LOG.info("In unregisterFlappingSingleton");
if (flappingSingletonService == null) {
- return RpcResultBuilder.<UnregisterFlappingSingletonOutput>failed().withError(
- ErrorType.RPC, "data-missing", "No prior RPC was registered").buildFuture();
+ return RpcResultBuilder.<UnregisterFlappingSingletonOutput>failed()
+ .withError(ErrorType.RPC, ErrorTag.DATA_MISSING, "No prior RPC was registered")
+ .buildFuture();
}
final long flapCount = flappingSingletonService.setInactive();
if (input.getContext() == null) {
return RpcResultBuilder.<RegisterBoundConstantOutput>failed().withError(
- ErrorType.RPC, "invalid-value", "Context value is null").buildFuture();
+ ErrorType.RPC, ErrorTag.INVALID_VALUE, "Context value is null").buildFuture();
}
if (input.getConstant() == null) {
return RpcResultBuilder.<RegisterBoundConstantOutput>failed().withError(
- ErrorType.RPC, "invalid-value", "Constant value is null").buildFuture();
+ ErrorType.RPC, ErrorTag.INVALID_VALUE, "Constant value is null").buildFuture();
}
if (routedRegistrations.containsKey(input.getContext())) {
return RpcResultBuilder.<RegisterBoundConstantOutput>failed().withError(ErrorType.RPC,
- "data-exists", "There is already an rpc registered for context: " + input.getContext()).buildFuture();
+ ErrorTag.DATA_EXISTS, "There is already an rpc registered for context: " + input.getContext())
+ .buildFuture();
}
final DOMRpcImplementationRegistration<RoutedGetConstantService> rpcRegistration =
LOG.info("In registerFlappingSingleton");
if (flappingSingletonService != null) {
- return RpcResultBuilder.<RegisterFlappingSingletonOutput>failed().withError(ErrorType.RPC,
- "data-exists", "There is already an rpc registered").buildFuture();
+ return RpcResultBuilder.<RegisterFlappingSingletonOutput>failed()
+ .withError(ErrorType.RPC, ErrorTag.DATA_EXISTS, "There is already an rpc registered")
+ .buildFuture();
}
flappingSingletonService = new FlappingSingletonService(singletonService);
LOG.info("In unsubscribeDtcl");
if (idIntsListener == null || dtclReg == null) {
- return RpcResultBuilder.<UnsubscribeDtclOutput>failed().withError(
- ErrorType.RPC, "data-missing", "No prior listener was registered").buildFuture();
+ return RpcResultBuilder.<UnsubscribeDtclOutput>failed()
+ .withError(ErrorType.RPC, ErrorTag.DATA_MISSING, "No prior listener was registered")
+ .buildFuture();
}
long timeout = 120L;
dtclReg = null;
if (!idIntsListener.hasTriggered()) {
- return RpcResultBuilder.<UnsubscribeDtclOutput>failed().withError(ErrorType.APPLICATION, "operation-failed",
- "id-ints listener has not received any notifications.").buildFuture();
+ return RpcResultBuilder.<UnsubscribeDtclOutput>failed()
+ .withError(ErrorType.APPLICATION, ErrorTag.OPERATION_FAILED,
+ "id-ints listener has not received any notifications.")
+ .buildFuture();
}
try (DOMDataTreeReadTransaction rTx = domDataBroker.newReadOnlyTransaction()) {
WriteTransactionsHandler.ID_INT_YID).get();
if (!readResult.isPresent()) {
- return RpcResultBuilder.<UnsubscribeDtclOutput>failed().withError(ErrorType.APPLICATION, "data-missing",
- "No data read from id-ints list").buildFuture();
+ return RpcResultBuilder.<UnsubscribeDtclOutput>failed()
+ .withError(ErrorType.APPLICATION, ErrorTag.DATA_MISSING, "No data read from id-ints list")
+ .buildFuture();
}
final boolean nodesEqual = idIntsListener.checkEqual(readResult.get());
idIntsListener.diffWithLocalCopy(readResult.get()));
}
- return RpcResultBuilder.success(new UnsubscribeDtclOutputBuilder().setCopyMatches(nodesEqual))
+ return RpcResultBuilder.success(new UnsubscribeDtclOutputBuilder().setCopyMatches(nodesEqual).build())
.buildFuture();
} catch (final InterruptedException | ExecutionException e) {
LOG.info("In unsubscribeYnl - input: {}", input);
if (!ynlRegistrations.containsKey(input.getId())) {
- return RpcResultBuilder.<UnsubscribeYnlOutput>failed().withError(
- ErrorType.RPC, "data-missing", "No prior listener was registered for " + input.getId()).buildFuture();
+ return RpcResultBuilder.<UnsubscribeYnlOutput>failed()
+ .withError(ErrorType.RPC, ErrorTag.DATA_MISSING,
+ "No prior listener was registered for " + input.getId())
+ .buildFuture();
}
final ListenerRegistration<YnlListener> reg = ynlRegistrations.remove(input.getId());
final PublishNotificationsTask task = publishNotificationsTasks.get(input.getId());
if (task == null) {
- return Futures.immediateFuture(RpcResultBuilder.success(
- new CheckPublishNotificationsOutputBuilder().setActive(false)).build());
+ return RpcResultBuilder.success(new CheckPublishNotificationsOutputBuilder().setActive(false).build())
+ .buildFuture();
}
final CheckPublishNotificationsOutputBuilder checkPublishNotificationsOutputBuilder =
final String shardName = input.getShardName();
if (Strings.isNullOrEmpty(shardName)) {
- return RpcResultBuilder.<ShutdownShardReplicaOutput>failed().withError(ErrorType.RPC, "bad-element",
- shardName + "is not a valid shard name").buildFuture();
+ return RpcResultBuilder.<ShutdownShardReplicaOutput>failed()
+ .withError(ErrorType.RPC, ErrorTag.BAD_ELEMENT, shardName + "is not a valid shard name")
+ .buildFuture();
}
return shutdownShardGracefully(shardName, new ShutdownShardReplicaOutputBuilder().build());
long timeoutInMS = Math.max(context.getDatastoreContext().getShardRaftConfig()
.getElectionTimeOutInterval().$times(3).toMillis(), 10000);
final FiniteDuration duration = FiniteDuration.apply(timeoutInMS, TimeUnit.MILLISECONDS);
- final scala.concurrent.Promise<Boolean> shutdownShardAsk = akka.dispatch.Futures.promise();
+ final scala.concurrent.Promise<Boolean> shutdownShardAsk = Futures.promise();
context.findLocalShardAsync(shardName).onComplete(new OnComplete<ActorRef>() {
@Override
LOG.info("In registerConstant - input: {}", input);
if (input.getConstant() == null) {
- return RpcResultBuilder.<RegisterConstantOutput>failed().withError(
- ErrorType.RPC, "invalid-value", "Constant value is null").buildFuture();
+ return RpcResultBuilder.<RegisterConstantOutput>failed()
+ .withError(ErrorType.RPC, ErrorTag.INVALID_VALUE, "Constant value is null")
+ .buildFuture();
}
if (globalGetConstantRegistration != null) {
- return RpcResultBuilder.<RegisterConstantOutput>failed().withError(ErrorType.RPC,
- "data-exists", "There is already an rpc registered").buildFuture();
+ return RpcResultBuilder.<RegisterConstantOutput>failed()
+ .withError(ErrorType.RPC, ErrorTag.DATA_EXISTS, "There is already an rpc registered")
+ .buildFuture();
}
globalGetConstantRegistration = GetConstantService.registerNew(domRpcService, input.getConstant());