Fix recreate client after destroy (#11780)
* Fix recreate client after destroy * Fix uts
This commit is contained in:
parent
16c61c53c6
commit
d0a1bd0143
|
|
@ -42,6 +42,7 @@ import org.apache.dubbo.rpc.model.ApplicationModel;
|
|||
import org.apache.dubbo.rpc.model.FrameworkModel;
|
||||
import org.apache.dubbo.rpc.model.ModuleModel;
|
||||
import org.apache.dubbo.rpc.model.ServiceMetadata;
|
||||
import org.apache.dubbo.rpc.protocol.ReferenceCountInvokerWrapper;
|
||||
import org.apache.dubbo.rpc.protocol.injvm.InjvmInvoker;
|
||||
import org.apache.dubbo.rpc.protocol.injvm.InjvmProtocol;
|
||||
import org.apache.dubbo.rpc.service.GenericService;
|
||||
|
|
@ -479,10 +480,11 @@ class ReferenceConfigTest {
|
|||
|
||||
referenceConfig.init();
|
||||
Assertions.assertTrue(referenceConfig.getInvoker() instanceof MockClusterInvoker);
|
||||
Invoker<?> withFilter = ((MockClusterInvoker<?>) referenceConfig.getInvoker()).getDirectory().getAllInvokers().get(0);
|
||||
Assertions.assertTrue(withFilter instanceof ListenerInvokerWrapper);
|
||||
Assertions.assertTrue(((ListenerInvokerWrapper<?>) withFilter).getInvoker() instanceof InjvmInvoker);
|
||||
URL url = withFilter.getUrl();
|
||||
Invoker<?> withCount = ((MockClusterInvoker<?>) referenceConfig.getInvoker()).getDirectory().getAllInvokers().get(0);
|
||||
Assertions.assertTrue(withCount instanceof ReferenceCountInvokerWrapper);
|
||||
Assertions.assertTrue(((ReferenceCountInvokerWrapper<?>) withCount).getInvoker() instanceof ListenerInvokerWrapper);
|
||||
Assertions.assertTrue(((ListenerInvokerWrapper<?>)(((ReferenceCountInvokerWrapper<?>) withCount).getInvoker())).getInvoker() instanceof InjvmInvoker);
|
||||
URL url = withCount.getUrl();
|
||||
Assertions.assertEquals("application1", url.getParameter("application"));
|
||||
Assertions.assertEquals("value1", url.getParameter("key1"));
|
||||
Assertions.assertEquals("value2", url.getParameter("key2"));
|
||||
|
|
|
|||
|
|
@ -0,0 +1,57 @@
|
|||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one or more
|
||||
* contributor license agreements. See the NOTICE file distributed with
|
||||
* this work for additional information regarding copyright ownership.
|
||||
* The ASF licenses this file to You 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 org.apache.dubbo.rpc.protocol;
|
||||
|
||||
import org.apache.dubbo.common.URL;
|
||||
import org.apache.dubbo.common.extension.Activate;
|
||||
import org.apache.dubbo.common.utils.UrlUtils;
|
||||
import org.apache.dubbo.rpc.Exporter;
|
||||
import org.apache.dubbo.rpc.Invoker;
|
||||
import org.apache.dubbo.rpc.Protocol;
|
||||
import org.apache.dubbo.rpc.RpcException;
|
||||
|
||||
@Activate(order = Integer.MIN_VALUE + 1000)
|
||||
public class InvokerCountWrapper implements Protocol {
|
||||
private final Protocol protocol;
|
||||
|
||||
public InvokerCountWrapper(Protocol protocol) {
|
||||
this.protocol = protocol;
|
||||
}
|
||||
|
||||
@Override
|
||||
public int getDefaultPort() {
|
||||
return protocol.getDefaultPort();
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> Exporter<T> export(Invoker<T> invoker) throws RpcException {
|
||||
return protocol.export(invoker);
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> Invoker<T> refer(Class<T> type, URL url) throws RpcException {
|
||||
if (UrlUtils.isRegistry(url)) {
|
||||
return protocol.refer(type, url);
|
||||
}
|
||||
return new ReferenceCountInvokerWrapper<>(protocol.refer(type, url));
|
||||
}
|
||||
|
||||
@Override
|
||||
public void destroy() {
|
||||
protocol.destroy();
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,87 @@
|
|||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one or more
|
||||
* contributor license agreements. See the NOTICE file distributed with
|
||||
* this work for additional information regarding copyright ownership.
|
||||
* The ASF licenses this file to You 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 org.apache.dubbo.rpc.protocol;
|
||||
|
||||
import org.apache.dubbo.common.URL;
|
||||
import org.apache.dubbo.common.constants.LoggerCodeConstants;
|
||||
import org.apache.dubbo.common.logger.ErrorTypeAwareLogger;
|
||||
import org.apache.dubbo.common.logger.LoggerFactory;
|
||||
import org.apache.dubbo.rpc.Invocation;
|
||||
import org.apache.dubbo.rpc.Invoker;
|
||||
import org.apache.dubbo.rpc.Result;
|
||||
import org.apache.dubbo.rpc.RpcException;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.locks.ReadWriteLock;
|
||||
import java.util.concurrent.locks.ReentrantReadWriteLock;
|
||||
|
||||
public class ReferenceCountInvokerWrapper<T> implements Invoker<T> {
|
||||
private final ErrorTypeAwareLogger logger = LoggerFactory.getErrorTypeAwareLogger(ReferenceCountInvokerWrapper.class);
|
||||
private final Invoker<T> invoker;
|
||||
|
||||
private final ReadWriteLock lock = new ReentrantReadWriteLock();
|
||||
private final AtomicBoolean destroyed = new AtomicBoolean(false);
|
||||
|
||||
public ReferenceCountInvokerWrapper(Invoker<T> invoker) {
|
||||
this.invoker = invoker;
|
||||
}
|
||||
|
||||
@Override
|
||||
public URL getUrl() {
|
||||
return invoker.getUrl();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isAvailable() {
|
||||
return !destroyed.get() && invoker.isAvailable();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void destroy() {
|
||||
try {
|
||||
lock.writeLock().lock();
|
||||
destroyed.set(true);
|
||||
} finally {
|
||||
lock.writeLock().unlock();
|
||||
}
|
||||
invoker.destroy();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Class<T> getInterface() {
|
||||
return invoker.getInterface();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Result invoke(Invocation invocation) throws RpcException {
|
||||
try {
|
||||
lock.readLock().lock();
|
||||
if (destroyed.get()) {
|
||||
logger.warn(LoggerCodeConstants.PROTOCOL_CLOSED_SERVER, "", "",
|
||||
"Remote invoker has been destroyed, and unable to invoke anymore.");
|
||||
throw new RpcException("This invoker has been destroyed!");
|
||||
}
|
||||
return invoker.invoke(invocation);
|
||||
} finally {
|
||||
lock.readLock().unlock();
|
||||
}
|
||||
}
|
||||
|
||||
public Invoker<T> getInvoker() {
|
||||
return invoker;
|
||||
}
|
||||
}
|
||||
|
|
@ -2,3 +2,4 @@ listener=org.apache.dubbo.rpc.protocol.ProtocolListenerWrapper
|
|||
mock=org.apache.dubbo.rpc.support.MockProtocol
|
||||
serializationwrapper=org.apache.dubbo.rpc.protocol.ProtocolSerializationWrapper
|
||||
securitywrapper=org.apache.dubbo.rpc.protocol.ProtocolSecurityWrapper
|
||||
invokercount=org.apache.dubbo.rpc.protocol.InvokerCountWrapper
|
||||
|
|
|
|||
Loading…
Reference in New Issue