2009-06-11 02:48:49 +02:00
|
|
|
/*
|
2011-12-09 06:18:34 +01:00
|
|
|
* Copyright 2011 The Netty Project
|
2009-06-11 02:48:49 +02:00
|
|
|
*
|
2011-12-09 06:18:34 +01:00
|
|
|
* The Netty Project 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:
|
2009-06-11 02:48:49 +02:00
|
|
|
*
|
2011-12-09 06:18:34 +01:00
|
|
|
* http://www.apache.org/licenses/LICENSE-2.0
|
2009-06-11 02:48:49 +02:00
|
|
|
*
|
2009-08-28 09:15:49 +02:00
|
|
|
* Unless required by applicable law or agreed to in writing, software
|
|
|
|
* distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
|
2011-12-09 06:18:34 +01:00
|
|
|
* WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
|
2009-08-28 09:15:49 +02:00
|
|
|
* License for the specific language governing permissions and limitations
|
|
|
|
* under the License.
|
2009-06-11 02:48:49 +02:00
|
|
|
*/
|
2011-12-09 04:38:59 +01:00
|
|
|
package io.netty.channel.socket.nio;
|
2009-06-11 02:48:49 +02:00
|
|
|
|
2012-02-18 23:02:56 +01:00
|
|
|
import static io.netty.channel.Channels.fireChannelOpen;
|
|
|
|
import io.netty.channel.ChannelException;
|
|
|
|
import io.netty.channel.ChannelFactory;
|
2012-03-20 09:43:00 +01:00
|
|
|
import io.netty.channel.ChannelFuture;
|
2012-02-18 23:02:56 +01:00
|
|
|
import io.netty.channel.ChannelPipeline;
|
|
|
|
import io.netty.channel.ChannelSink;
|
2012-04-02 11:57:32 +02:00
|
|
|
import io.netty.channel.Channels;
|
2012-02-18 23:02:56 +01:00
|
|
|
import io.netty.channel.socket.DatagramChannelConfig;
|
2012-04-02 11:07:11 +02:00
|
|
|
import io.netty.util.internal.DetectionUtil;
|
2009-06-11 02:48:49 +02:00
|
|
|
|
|
|
|
import java.io.IOException;
|
2009-06-11 08:10:46 +02:00
|
|
|
import java.net.InetAddress;
|
2009-06-11 02:48:49 +02:00
|
|
|
import java.net.InetSocketAddress;
|
2009-06-11 08:10:46 +02:00
|
|
|
import java.net.NetworkInterface;
|
2012-03-20 09:43:00 +01:00
|
|
|
import java.net.SocketAddress;
|
2012-04-02 11:07:11 +02:00
|
|
|
import java.net.SocketException;
|
2009-06-11 02:48:49 +02:00
|
|
|
import java.nio.channels.DatagramChannel;
|
2012-04-02 11:07:11 +02:00
|
|
|
import java.nio.channels.MembershipKey;
|
|
|
|
import java.util.ArrayList;
|
|
|
|
import java.util.HashMap;
|
|
|
|
import java.util.Iterator;
|
|
|
|
import java.util.List;
|
|
|
|
import java.util.Map;
|
2009-06-11 02:48:49 +02:00
|
|
|
|
|
|
|
/**
|
2011-12-09 04:38:59 +01:00
|
|
|
* Provides an NIO based {@link io.netty.channel.socket.DatagramChannel}.
|
2009-06-11 02:48:49 +02:00
|
|
|
*/
|
2012-03-29 17:07:19 +02:00
|
|
|
public final class NioDatagramChannel extends AbstractNioChannel implements io.netty.channel.socket.DatagramChannel {
|
2009-06-11 02:48:49 +02:00
|
|
|
|
|
|
|
/**
|
|
|
|
* The {@link DatagramChannelConfig}.
|
|
|
|
*/
|
|
|
|
private final NioDatagramChannelConfig config;
|
2012-04-02 11:07:11 +02:00
|
|
|
private Map<InetAddress, List<MembershipKey>> memberships;
|
2012-02-18 23:02:56 +01:00
|
|
|
|
2011-11-02 19:17:10 +01:00
|
|
|
static NioDatagramChannel create(ChannelFactory factory,
|
|
|
|
ChannelPipeline pipeline, ChannelSink sink, NioDatagramWorker worker) {
|
|
|
|
NioDatagramChannel instance =
|
|
|
|
new NioDatagramChannel(factory, pipeline, sink, worker);
|
|
|
|
fireChannelOpen(instance);
|
|
|
|
return instance;
|
|
|
|
}
|
|
|
|
|
|
|
|
private NioDatagramChannel(final ChannelFactory factory,
|
2009-06-11 02:48:49 +02:00
|
|
|
final ChannelPipeline pipeline, final ChannelSink sink,
|
2009-06-12 04:47:57 +02:00
|
|
|
final NioDatagramWorker worker) {
|
2012-03-29 17:07:19 +02:00
|
|
|
super(null, factory, pipeline, sink, worker, new NioDatagramJdkChannel(openNonBlockingChannel()));
|
2012-04-02 11:57:32 +02:00
|
|
|
config = new DefaultNioDatagramChannelConfig(getJdkChannel().getChannel());
|
2009-06-11 02:48:49 +02:00
|
|
|
}
|
|
|
|
|
2012-02-18 23:02:56 +01:00
|
|
|
private static DatagramChannel openNonBlockingChannel() {
|
2009-06-11 02:48:49 +02:00
|
|
|
try {
|
|
|
|
final DatagramChannel channel = DatagramChannel.open();
|
|
|
|
channel.configureBlocking(false);
|
|
|
|
return channel;
|
|
|
|
} catch (final IOException e) {
|
|
|
|
throw new ChannelException("Failed to open a DatagramChannel.", e);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
2012-03-29 17:07:19 +02:00
|
|
|
@Override
|
|
|
|
protected NioDatagramJdkChannel getJdkChannel() {
|
|
|
|
return (NioDatagramJdkChannel) super.getJdkChannel();
|
|
|
|
}
|
|
|
|
|
2012-03-06 19:26:32 +01:00
|
|
|
@Override
|
|
|
|
public NioDatagramWorker getWorker() {
|
|
|
|
return (NioDatagramWorker) super.getWorker();
|
|
|
|
}
|
|
|
|
|
2010-11-12 01:45:39 +01:00
|
|
|
@Override
|
2009-06-11 02:48:49 +02:00
|
|
|
public boolean isBound() {
|
2012-03-29 17:07:19 +02:00
|
|
|
return isOpen() && getJdkChannel().isSocketBound();
|
2009-06-11 02:48:49 +02:00
|
|
|
}
|
|
|
|
|
2010-11-12 01:45:39 +01:00
|
|
|
@Override
|
2009-06-11 02:48:49 +02:00
|
|
|
public boolean isConnected() {
|
2012-03-29 17:07:19 +02:00
|
|
|
return getJdkChannel().isConnected();
|
2009-06-11 02:48:49 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
|
|
|
protected boolean setClosed() {
|
|
|
|
return super.setClosed();
|
|
|
|
}
|
|
|
|
|
2010-11-12 01:45:39 +01:00
|
|
|
@Override
|
2009-06-11 02:48:49 +02:00
|
|
|
public NioDatagramChannelConfig getConfig() {
|
|
|
|
return config;
|
|
|
|
}
|
|
|
|
|
2010-11-12 01:45:39 +01:00
|
|
|
@Override
|
2012-04-02 11:57:32 +02:00
|
|
|
public ChannelFuture joinGroup(InetAddress multicastAddress) {
|
2012-04-02 11:07:11 +02:00
|
|
|
try {
|
2012-04-02 11:57:32 +02:00
|
|
|
return joinGroup(multicastAddress, NetworkInterface.getByInetAddress(getLocalAddress().getAddress()), null);
|
2012-04-02 11:07:11 +02:00
|
|
|
} catch (SocketException e) {
|
2012-04-02 11:57:32 +02:00
|
|
|
return Channels.failedFuture(this, e);
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
2009-06-11 08:10:46 +02:00
|
|
|
}
|
|
|
|
|
2010-11-12 01:45:39 +01:00
|
|
|
@Override
|
2012-04-02 11:57:32 +02:00
|
|
|
public ChannelFuture joinGroup(InetSocketAddress multicastAddress, NetworkInterface networkInterface) {
|
|
|
|
return joinGroup(multicastAddress.getAddress(), networkInterface, null);
|
2009-06-11 08:10:46 +02:00
|
|
|
}
|
|
|
|
|
2012-04-02 11:07:11 +02:00
|
|
|
/**
|
|
|
|
* Joins the specified multicast group at the specified interface using the specified source.
|
|
|
|
*/
|
2012-04-02 11:57:32 +02:00
|
|
|
public ChannelFuture joinGroup(InetAddress multicastAddress, NetworkInterface networkInterface, InetAddress source) {
|
2012-04-02 11:07:11 +02:00
|
|
|
if (DetectionUtil.javaVersion() < 7) {
|
|
|
|
throw new UnsupportedOperationException();
|
|
|
|
} else {
|
|
|
|
if (multicastAddress == null) {
|
|
|
|
throw new NullPointerException("multicastAddress");
|
|
|
|
}
|
|
|
|
|
|
|
|
if (networkInterface == null) {
|
|
|
|
throw new NullPointerException("networkInterface");
|
|
|
|
}
|
|
|
|
|
|
|
|
try {
|
|
|
|
MembershipKey key = getJdkChannel().getChannel().join(multicastAddress, networkInterface);
|
|
|
|
synchronized (this) {
|
|
|
|
if (memberships == null) {
|
|
|
|
memberships = new HashMap<InetAddress, List<MembershipKey>>();
|
|
|
|
|
|
|
|
}
|
|
|
|
List<MembershipKey> keys = memberships.get(multicastAddress);
|
|
|
|
if (keys == null) {
|
|
|
|
keys = new ArrayList<MembershipKey>();
|
|
|
|
memberships.put(multicastAddress, keys);
|
|
|
|
}
|
|
|
|
|
|
|
|
keys.add(key);
|
|
|
|
}
|
2012-04-02 11:57:32 +02:00
|
|
|
} catch (Throwable e) {
|
|
|
|
return Channels.failedFuture(this, e);
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
|
|
|
}
|
2012-04-02 11:57:32 +02:00
|
|
|
return Channels.succeededFuture(this);
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
|
|
|
|
2010-11-12 01:45:39 +01:00
|
|
|
@Override
|
2012-04-02 11:57:32 +02:00
|
|
|
public ChannelFuture leaveGroup(InetAddress multicastAddress) {
|
2012-04-02 11:07:11 +02:00
|
|
|
try {
|
2012-04-02 11:57:32 +02:00
|
|
|
return leaveGroup(multicastAddress, NetworkInterface.getByInetAddress(getLocalAddress().getAddress()), null);
|
2012-04-02 11:07:11 +02:00
|
|
|
} catch (SocketException e) {
|
2012-04-02 11:57:32 +02:00
|
|
|
return Channels.failedFuture(this, e);
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
|
|
|
|
2009-06-11 08:10:46 +02:00
|
|
|
}
|
|
|
|
|
2010-11-12 01:45:39 +01:00
|
|
|
@Override
|
2012-04-02 11:57:32 +02:00
|
|
|
public ChannelFuture leaveGroup(InetSocketAddress multicastAddress,
|
2009-06-11 08:10:46 +02:00
|
|
|
NetworkInterface networkInterface) {
|
2012-04-02 11:57:32 +02:00
|
|
|
return leaveGroup(multicastAddress.getAddress(), networkInterface, null);
|
2009-06-11 08:10:46 +02:00
|
|
|
}
|
2012-02-18 23:02:56 +01:00
|
|
|
|
2012-04-02 11:07:11 +02:00
|
|
|
/**
|
|
|
|
* Leave the specified multicast group at the specified interface using the specified source.
|
|
|
|
*/
|
2012-04-02 11:57:32 +02:00
|
|
|
public ChannelFuture leaveGroup(InetAddress multicastAddress,
|
2012-04-02 11:07:11 +02:00
|
|
|
NetworkInterface networkInterface, InetAddress source) {
|
|
|
|
if (DetectionUtil.javaVersion() < 7) {
|
|
|
|
throw new UnsupportedOperationException();
|
|
|
|
} else {
|
|
|
|
if (multicastAddress == null) {
|
|
|
|
throw new NullPointerException("multicastAddress");
|
|
|
|
}
|
|
|
|
|
|
|
|
if (networkInterface == null) {
|
|
|
|
throw new NullPointerException("networkInterface");
|
|
|
|
}
|
|
|
|
|
|
|
|
synchronized (this) {
|
|
|
|
if (memberships != null) {
|
|
|
|
List<MembershipKey> keys = memberships.get(multicastAddress);
|
|
|
|
if (keys != null) {
|
|
|
|
Iterator<MembershipKey> keyIt = keys.iterator();
|
|
|
|
|
2012-04-02 11:57:32 +02:00
|
|
|
while (keyIt.hasNext()) {
|
2012-04-02 11:07:11 +02:00
|
|
|
MembershipKey key = keyIt.next();
|
|
|
|
if (networkInterface.equals(key.networkInterface())) {
|
|
|
|
if (source == null && key.sourceAddress() == null || (source != null && source.equals(key.sourceAddress()))) {
|
|
|
|
key.drop();
|
|
|
|
keyIt.remove();
|
|
|
|
}
|
|
|
|
|
|
|
|
}
|
|
|
|
}
|
|
|
|
if (keys.isEmpty()) {
|
|
|
|
memberships.remove(multicastAddress);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2012-04-02 11:57:32 +02:00
|
|
|
return Channels.succeededFuture(this);
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
/**
|
|
|
|
* Block the given sourceToBlock address for the given multicastAddress on the given networkInterface
|
|
|
|
*
|
|
|
|
*/
|
2012-04-02 11:57:32 +02:00
|
|
|
public ChannelFuture block(InetAddress multicastAddress,
|
2012-04-02 11:07:11 +02:00
|
|
|
NetworkInterface networkInterface, InetAddress sourceToBlock) {
|
|
|
|
if (DetectionUtil.javaVersion() < 7) {
|
|
|
|
throw new UnsupportedOperationException();
|
|
|
|
} else {
|
|
|
|
if (multicastAddress == null) {
|
|
|
|
throw new NullPointerException("multicastAddress");
|
|
|
|
}
|
|
|
|
if (sourceToBlock == null) {
|
|
|
|
throw new NullPointerException("sourceToBlock");
|
|
|
|
}
|
|
|
|
|
|
|
|
if (networkInterface == null) {
|
|
|
|
throw new NullPointerException("networkInterface");
|
|
|
|
}
|
|
|
|
synchronized (this) {
|
|
|
|
if (memberships != null) {
|
|
|
|
List<MembershipKey> keys = memberships.get(multicastAddress);
|
|
|
|
for (MembershipKey key: keys) {
|
|
|
|
if (networkInterface.equals(key.networkInterface())) {
|
|
|
|
try {
|
|
|
|
key.block(sourceToBlock);
|
|
|
|
} catch (IOException e) {
|
2012-04-02 11:57:32 +02:00
|
|
|
return Channels.failedFuture(this, e);
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2012-04-02 11:57:32 +02:00
|
|
|
return Channels.succeededFuture(this);
|
2012-04-02 11:07:11 +02:00
|
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
}
|
2012-04-02 11:57:32 +02:00
|
|
|
|
2012-04-02 11:07:11 +02:00
|
|
|
/**
|
|
|
|
* Block the given sourceToBlock address for the given multicastAddress
|
|
|
|
*
|
|
|
|
*/
|
2012-04-02 11:57:32 +02:00
|
|
|
public ChannelFuture block(InetAddress multicastAddress, InetAddress sourceToBlock) {
|
2012-04-02 11:07:11 +02:00
|
|
|
try {
|
|
|
|
block(multicastAddress, NetworkInterface.getByInetAddress(getLocalAddress().getAddress()), sourceToBlock);
|
|
|
|
} catch (SocketException e) {
|
2012-04-02 11:57:32 +02:00
|
|
|
return Channels.failedFuture(this, e);
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
2012-04-02 11:57:32 +02:00
|
|
|
return Channels.succeededFuture(this);
|
|
|
|
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
|
|
|
|
2012-03-20 09:43:00 +01:00
|
|
|
@Override
|
|
|
|
public ChannelFuture write(Object message, SocketAddress remoteAddress) {
|
|
|
|
if (remoteAddress == null || remoteAddress.equals(getRemoteAddress())) {
|
|
|
|
return super.write(message, null);
|
|
|
|
} else {
|
2012-04-02 14:20:40 +02:00
|
|
|
return Channels.write(this, message, remoteAddress);
|
2012-03-20 09:43:00 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
}
|
2009-06-11 02:48:49 +02:00
|
|
|
}
|