2009-06-11 02:48:49 +02:00
|
|
|
/*
|
2012-06-04 22:31:44 +02:00
|
|
|
* Copyright 2012 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
|
|
|
*
|
2012-06-04 22:31:44 +02: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-12-17 09:43:45 +01:00
|
|
|
import io.netty.buffer.BufType;
|
2012-06-10 04:08:43 +02:00
|
|
|
import io.netty.buffer.ByteBuf;
|
2012-06-12 10:02:00 +02:00
|
|
|
import io.netty.buffer.MessageBuf;
|
2012-02-18 23:02:56 +01:00
|
|
|
import io.netty.channel.ChannelException;
|
2012-03-20 09:43:00 +01:00
|
|
|
import io.netty.channel.ChannelFuture;
|
2012-06-19 03:39:30 +02:00
|
|
|
import io.netty.channel.ChannelMetadata;
|
2012-12-30 17:40:24 +01:00
|
|
|
import io.netty.channel.ChannelPromise;
|
2013-02-01 09:02:26 +01:00
|
|
|
import io.netty.channel.nio.AbstractNioMessageChannel;
|
2012-02-18 23:02:56 +01:00
|
|
|
import io.netty.channel.socket.DatagramChannelConfig;
|
2012-05-24 17:57:10 +02:00
|
|
|
import io.netty.channel.socket.DatagramPacket;
|
2012-05-31 11:49:39 +02:00
|
|
|
import io.netty.channel.socket.InternetProtocolFamily;
|
2013-01-11 06:03:27 +01:00
|
|
|
import io.netty.util.internal.PlatformDependent;
|
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;
|
2012-05-24 17:57:10 +02:00
|
|
|
import java.nio.ByteBuffer;
|
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;
|
2012-05-24 17:57:10 +02:00
|
|
|
import java.nio.channels.SelectionKey;
|
2012-04-02 11:07:11 +02:00
|
|
|
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
|
|
|
|
|
|
|
/**
|
2012-12-23 19:24:20 +01:00
|
|
|
* Provides an NIO based {@link io.netty.channel.socket.DatagramChannel} which can be used
|
|
|
|
* to send and receive {@link DatagramPacket}'s.
|
2009-06-11 02:48:49 +02:00
|
|
|
*/
|
2012-06-08 12:28:12 +02:00
|
|
|
public final class NioDatagramChannel
|
|
|
|
extends AbstractNioMessageChannel implements io.netty.channel.socket.DatagramChannel {
|
2009-06-11 02:48:49 +02:00
|
|
|
|
2012-12-17 09:43:45 +01:00
|
|
|
private static final ChannelMetadata METADATA = new ChannelMetadata(BufType.MESSAGE, true);
|
2012-06-19 03:39:30 +02:00
|
|
|
|
2012-05-24 17:57:10 +02:00
|
|
|
private final DatagramChannelConfig config;
|
|
|
|
private final Map<InetAddress, List<MembershipKey>> memberships =
|
|
|
|
new HashMap<InetAddress, List<MembershipKey>>();
|
|
|
|
|
|
|
|
private static DatagramChannel newSocket() {
|
|
|
|
try {
|
|
|
|
return DatagramChannel.open();
|
|
|
|
} catch (IOException e) {
|
|
|
|
throw new ChannelException("Failed to open a socket.", e);
|
|
|
|
}
|
2012-04-03 12:04:33 +02:00
|
|
|
}
|
2012-05-24 17:57:10 +02:00
|
|
|
|
2012-05-31 11:49:39 +02:00
|
|
|
private static DatagramChannel newSocket(InternetProtocolFamily ipFamily) {
|
|
|
|
if (ipFamily == null) {
|
|
|
|
return newSocket();
|
|
|
|
}
|
|
|
|
|
2013-01-11 06:03:27 +01:00
|
|
|
if (PlatformDependent.javaVersion() < 7) {
|
2012-05-31 11:49:39 +02:00
|
|
|
throw new UnsupportedOperationException();
|
|
|
|
}
|
|
|
|
|
|
|
|
try {
|
|
|
|
return DatagramChannel.open(ProtocolFamilyConverter.convert(ipFamily));
|
|
|
|
} catch (IOException e) {
|
|
|
|
throw new ChannelException("Failed to open a socket.", e);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2012-12-23 19:24:20 +01:00
|
|
|
/**
|
|
|
|
* Create a new instance which will use the Operation Systems default {@link InternetProtocolFamily}.
|
|
|
|
*/
|
2012-05-24 17:57:10 +02:00
|
|
|
public NioDatagramChannel() {
|
|
|
|
this(newSocket());
|
2011-11-02 19:17:10 +01:00
|
|
|
}
|
|
|
|
|
2012-12-23 19:24:20 +01:00
|
|
|
/**
|
|
|
|
* Create a new instance using the given {@link InternetProtocolFamily}. If {@code null} is used it will depend
|
|
|
|
* on the Operation Systems default which will be chosen.
|
|
|
|
*/
|
2012-05-31 11:49:39 +02:00
|
|
|
public NioDatagramChannel(InternetProtocolFamily ipFamily) {
|
|
|
|
this(newSocket(ipFamily));
|
|
|
|
}
|
|
|
|
|
2012-12-23 19:24:20 +01:00
|
|
|
/**
|
|
|
|
* Create a new instance from the given {@link DatagramChannel}.
|
|
|
|
*/
|
2012-05-24 17:57:10 +02:00
|
|
|
public NioDatagramChannel(DatagramChannel socket) {
|
|
|
|
this(null, socket);
|
2009-06-11 02:48:49 +02:00
|
|
|
}
|
|
|
|
|
2012-12-23 19:24:20 +01:00
|
|
|
/**
|
|
|
|
* Create a new instance from the given {@link DatagramChannel}.
|
|
|
|
*
|
|
|
|
* @param id the id to use for this instance or {@code null} if a new one should be generated.
|
|
|
|
* @param socket the {@link DatagramChannel} which will be used
|
|
|
|
*/
|
2012-05-24 17:57:10 +02:00
|
|
|
public NioDatagramChannel(Integer id, DatagramChannel socket) {
|
2012-06-07 14:06:56 +02:00
|
|
|
super(null, id, socket, SelectionKey.OP_READ);
|
2013-01-01 08:49:21 +01:00
|
|
|
config = new NioDatagramChannelConfig(this, socket);
|
2009-06-11 02:48:49 +02:00
|
|
|
}
|
|
|
|
|
2012-06-19 03:39:30 +02:00
|
|
|
@Override
|
|
|
|
public ChannelMetadata metadata() {
|
|
|
|
return METADATA;
|
|
|
|
}
|
|
|
|
|
2012-05-24 17:57:10 +02:00
|
|
|
@Override
|
|
|
|
public DatagramChannelConfig config() {
|
|
|
|
return config;
|
|
|
|
}
|
2009-06-11 02:48:49 +02:00
|
|
|
|
2012-03-29 17:07:19 +02:00
|
|
|
@Override
|
2012-05-24 17:57:10 +02:00
|
|
|
public boolean isActive() {
|
|
|
|
DatagramChannel ch = javaChannel();
|
|
|
|
return ch.isOpen() && ch.socket().isBound();
|
2012-03-29 17:07:19 +02:00
|
|
|
}
|
|
|
|
|
2012-06-19 03:43:38 +02:00
|
|
|
@Override
|
|
|
|
public boolean isConnected() {
|
|
|
|
return javaChannel().isConnected();
|
|
|
|
}
|
|
|
|
|
2012-03-06 19:26:32 +01:00
|
|
|
@Override
|
2012-05-24 17:57:10 +02:00
|
|
|
protected DatagramChannel javaChannel() {
|
|
|
|
return (DatagramChannel) super.javaChannel();
|
2012-03-06 19:26:32 +01:00
|
|
|
}
|
|
|
|
|
2010-11-12 01:45:39 +01:00
|
|
|
@Override
|
2012-05-24 17:57:10 +02:00
|
|
|
protected SocketAddress localAddress0() {
|
|
|
|
return javaChannel().socket().getLocalSocketAddress();
|
2009-06-11 02:48:49 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
2012-05-24 17:57:10 +02:00
|
|
|
protected SocketAddress remoteAddress0() {
|
|
|
|
return javaChannel().socket().getRemoteSocketAddress();
|
2009-06-11 02:48:49 +02:00
|
|
|
}
|
|
|
|
|
2010-11-12 01:45:39 +01:00
|
|
|
@Override
|
2012-05-24 17:57:10 +02:00
|
|
|
protected void doBind(SocketAddress localAddress) throws Exception {
|
2012-05-24 18:32:14 +02:00
|
|
|
javaChannel().socket().bind(localAddress);
|
2012-05-24 17:57:10 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
|
|
|
protected boolean doConnect(SocketAddress remoteAddress,
|
|
|
|
SocketAddress localAddress) throws Exception {
|
|
|
|
if (localAddress != null) {
|
|
|
|
javaChannel().socket().bind(localAddress);
|
|
|
|
}
|
|
|
|
|
|
|
|
boolean success = false;
|
|
|
|
try {
|
|
|
|
javaChannel().connect(remoteAddress);
|
|
|
|
success = true;
|
|
|
|
return true;
|
|
|
|
} finally {
|
|
|
|
if (!success) {
|
|
|
|
doClose();
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
|
|
|
protected void doFinishConnect() throws Exception {
|
|
|
|
throw new Error();
|
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
|
|
|
protected void doDisconnect() throws Exception {
|
|
|
|
javaChannel().disconnect();
|
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
|
|
|
protected void doClose() throws Exception {
|
|
|
|
javaChannel().close();
|
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
2012-06-12 10:02:00 +02:00
|
|
|
protected int doReadMessages(MessageBuf<Object> buf) throws Exception {
|
2012-05-24 17:57:10 +02:00
|
|
|
DatagramChannel ch = javaChannel();
|
2013-01-01 16:03:18 +01:00
|
|
|
ByteBuf buffer = alloc().directBuffer(config().getReceivePacketSize());
|
|
|
|
boolean free = true;
|
|
|
|
try {
|
|
|
|
ByteBuffer data = buffer.nioBuffer(buffer.writerIndex(), buffer.writableBytes());
|
2012-05-24 17:57:10 +02:00
|
|
|
|
2013-01-01 16:03:18 +01:00
|
|
|
InetSocketAddress remoteAddress = (InetSocketAddress) ch.receive(data);
|
|
|
|
if (remoteAddress == null) {
|
|
|
|
return 0;
|
|
|
|
}
|
2013-01-19 17:22:28 +01:00
|
|
|
buf.add(new DatagramPacket(buffer.writerIndex(buffer.writerIndex() + data.position()), remoteAddress));
|
2013-01-01 16:03:18 +01:00
|
|
|
free = false;
|
|
|
|
return 1;
|
|
|
|
} catch (Throwable cause) {
|
|
|
|
if (cause instanceof Error) {
|
|
|
|
throw (Error) cause;
|
|
|
|
}
|
|
|
|
if (cause instanceof RuntimeException) {
|
|
|
|
throw (RuntimeException) cause;
|
|
|
|
}
|
|
|
|
if (cause instanceof Exception) {
|
|
|
|
throw (Exception) cause;
|
|
|
|
}
|
|
|
|
throw new ChannelException(cause);
|
|
|
|
} finally {
|
|
|
|
if (free) {
|
2013-02-10 05:10:09 +01:00
|
|
|
buffer.release();
|
2013-01-01 16:03:18 +01:00
|
|
|
}
|
|
|
|
}
|
2012-05-24 17:57:10 +02:00
|
|
|
}
|
|
|
|
|
2012-05-25 15:16:25 +02:00
|
|
|
@Override
|
2012-06-12 10:02:00 +02:00
|
|
|
protected int doWriteMessages(MessageBuf<Object> buf, boolean lastSpin) throws Exception {
|
2012-05-24 17:57:10 +02:00
|
|
|
DatagramPacket packet = (DatagramPacket) buf.peek();
|
2012-06-10 04:08:43 +02:00
|
|
|
ByteBuf data = packet.data();
|
2012-08-30 07:04:13 +02:00
|
|
|
int dataLen = data.readableBytes();
|
2012-06-02 10:30:55 +02:00
|
|
|
ByteBuffer nioData;
|
2012-12-14 04:20:33 +01:00
|
|
|
if (data.nioBufferCount() == 1) {
|
2012-06-02 10:30:55 +02:00
|
|
|
nioData = data.nioBuffer();
|
|
|
|
} else {
|
2012-08-30 07:04:13 +02:00
|
|
|
nioData = ByteBuffer.allocate(dataLen);
|
2012-06-02 10:30:55 +02:00
|
|
|
data.getBytes(data.readerIndex(), nioData);
|
|
|
|
nioData.flip();
|
|
|
|
}
|
|
|
|
|
|
|
|
final int writtenBytes = javaChannel().send(nioData, packet.remoteAddress());
|
2012-05-24 17:57:10 +02:00
|
|
|
|
|
|
|
final SelectionKey key = selectionKey();
|
|
|
|
final int interestOps = key.interestOps();
|
2012-08-30 07:04:13 +02:00
|
|
|
if (writtenBytes <= 0 && dataLen > 0) {
|
2012-05-24 17:57:10 +02:00
|
|
|
// Did not write a packet.
|
|
|
|
// 1) If 'lastSpin' is false, the caller will call this method again real soon.
|
|
|
|
// - Do not update OP_WRITE.
|
|
|
|
// 2) If 'lastSpin' is true, the caller will not retry.
|
|
|
|
// - Set OP_WRITE so that the event loop calls flushForcibly() later.
|
|
|
|
if (lastSpin) {
|
|
|
|
if ((interestOps & SelectionKey.OP_WRITE) == 0) {
|
|
|
|
key.interestOps(interestOps | SelectionKey.OP_WRITE);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return 0;
|
|
|
|
}
|
|
|
|
|
|
|
|
// Wrote a packet.
|
|
|
|
buf.remove();
|
2013-01-01 16:03:18 +01:00
|
|
|
|
|
|
|
// packet was written free up buffer
|
2013-02-10 05:10:09 +01:00
|
|
|
packet.release();
|
2013-01-01 16:03:18 +01:00
|
|
|
|
2012-05-24 17:57:10 +02:00
|
|
|
if (buf.isEmpty()) {
|
|
|
|
// Wrote the outbound buffer completely - clear OP_WRITE.
|
|
|
|
if ((interestOps & SelectionKey.OP_WRITE) != 0) {
|
|
|
|
key.interestOps(interestOps & ~SelectionKey.OP_WRITE);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return 1;
|
2009-06-11 02:48:49 +02:00
|
|
|
}
|
|
|
|
|
2010-11-12 01:45:39 +01:00
|
|
|
@Override
|
2013-02-01 09:02:26 +01:00
|
|
|
public InetSocketAddress localAddress() {
|
|
|
|
return (InetSocketAddress) super.localAddress();
|
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
|
|
|
public InetSocketAddress remoteAddress() {
|
|
|
|
return (InetSocketAddress) super.remoteAddress();
|
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
2012-04-02 11:57:32 +02:00
|
|
|
public ChannelFuture joinGroup(InetAddress multicastAddress) {
|
2012-12-30 17:40:24 +01:00
|
|
|
return joinGroup(multicastAddress, newPromise());
|
2012-05-24 17:57:10 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
2012-12-30 17:40:24 +01:00
|
|
|
public ChannelFuture joinGroup(InetAddress multicastAddress, ChannelPromise promise) {
|
2012-05-24 17:57:10 +02:00
|
|
|
try {
|
|
|
|
return joinGroup(
|
|
|
|
multicastAddress,
|
|
|
|
NetworkInterface.getByInetAddress(localAddress().getAddress()),
|
2012-12-30 17:40:24 +01:00
|
|
|
null, promise);
|
2012-04-02 11:07:11 +02:00
|
|
|
} catch (SocketException e) {
|
2012-12-30 17:40:24 +01:00
|
|
|
promise.setFailure(e);
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
2012-12-30 17:40:24 +01:00
|
|
|
return promise;
|
2009-06-11 08:10:46 +02:00
|
|
|
}
|
|
|
|
|
2010-11-12 01:45:39 +01:00
|
|
|
@Override
|
2012-05-24 17:57:10 +02:00
|
|
|
public ChannelFuture joinGroup(
|
|
|
|
InetSocketAddress multicastAddress, NetworkInterface networkInterface) {
|
2012-12-30 17:40:24 +01:00
|
|
|
return joinGroup(multicastAddress, networkInterface, newPromise());
|
2009-06-11 08:10:46 +02:00
|
|
|
}
|
|
|
|
|
2012-05-24 17:57:10 +02:00
|
|
|
@Override
|
|
|
|
public ChannelFuture joinGroup(
|
|
|
|
InetSocketAddress multicastAddress, NetworkInterface networkInterface,
|
2012-12-30 17:40:24 +01:00
|
|
|
ChannelPromise promise) {
|
|
|
|
return joinGroup(multicastAddress.getAddress(), networkInterface, null, promise);
|
2012-05-24 17:57:10 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
|
|
|
public ChannelFuture joinGroup(
|
|
|
|
InetAddress multicastAddress, NetworkInterface networkInterface, InetAddress source) {
|
2012-12-30 17:40:24 +01:00
|
|
|
return joinGroup(multicastAddress, networkInterface, source, newPromise());
|
2012-05-24 17:57:10 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
|
|
|
public ChannelFuture joinGroup(
|
|
|
|
InetAddress multicastAddress, NetworkInterface networkInterface,
|
2012-12-30 17:40:24 +01:00
|
|
|
InetAddress source, ChannelPromise promise) {
|
2013-01-11 06:03:27 +01:00
|
|
|
if (PlatformDependent.javaVersion() >= 7) {
|
2012-04-02 11:07:11 +02:00
|
|
|
if (multicastAddress == null) {
|
|
|
|
throw new NullPointerException("multicastAddress");
|
|
|
|
}
|
2012-05-24 17:57:10 +02:00
|
|
|
|
2012-04-02 11:07:11 +02:00
|
|
|
if (networkInterface == null) {
|
|
|
|
throw new NullPointerException("networkInterface");
|
|
|
|
}
|
2012-05-24 17:57:10 +02:00
|
|
|
|
2012-04-02 11:07:11 +02:00
|
|
|
try {
|
2012-04-02 15:25:40 +02:00
|
|
|
MembershipKey key;
|
|
|
|
if (source == null) {
|
2012-05-24 17:57:10 +02:00
|
|
|
key = javaChannel().join(multicastAddress, networkInterface);
|
2012-04-02 15:25:40 +02:00
|
|
|
} else {
|
2012-05-24 17:57:10 +02:00
|
|
|
key = javaChannel().join(multicastAddress, networkInterface, source);
|
2012-04-02 15:25:40 +02:00
|
|
|
}
|
|
|
|
|
2012-04-02 11:07:11 +02:00
|
|
|
synchronized (this) {
|
|
|
|
List<MembershipKey> keys = memberships.get(multicastAddress);
|
|
|
|
if (keys == null) {
|
|
|
|
keys = new ArrayList<MembershipKey>();
|
|
|
|
memberships.put(multicastAddress, keys);
|
|
|
|
}
|
|
|
|
keys.add(key);
|
|
|
|
}
|
2012-05-24 17:57:10 +02:00
|
|
|
|
2012-12-30 17:40:24 +01:00
|
|
|
promise.setSuccess();
|
2012-04-02 11:57:32 +02:00
|
|
|
} catch (Throwable e) {
|
2012-12-30 17:40:24 +01:00
|
|
|
promise.setFailure(e);
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
2012-11-12 01:31:40 +01:00
|
|
|
} else {
|
|
|
|
throw new UnsupportedOperationException();
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
2012-12-30 17:40:24 +01:00
|
|
|
return promise;
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
2012-05-24 17:57:10 +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-12-30 17:40:24 +01:00
|
|
|
return leaveGroup(multicastAddress, newPromise());
|
2012-05-24 17:57:10 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
2012-12-30 17:40:24 +01:00
|
|
|
public ChannelFuture leaveGroup(InetAddress multicastAddress, ChannelPromise promise) {
|
2012-04-02 11:07:11 +02:00
|
|
|
try {
|
2012-06-08 12:28:12 +02:00
|
|
|
return leaveGroup(
|
2012-12-30 17:40:24 +01:00
|
|
|
multicastAddress, NetworkInterface.getByInetAddress(localAddress().getAddress()), null, promise);
|
2012-04-02 11:07:11 +02:00
|
|
|
} catch (SocketException e) {
|
2012-12-30 17:40:24 +01:00
|
|
|
promise.setFailure(e);
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
2012-12-30 17:40:24 +01:00
|
|
|
return promise;
|
2009-06-11 08:10:46 +02:00
|
|
|
}
|
|
|
|
|
2010-11-12 01:45:39 +01:00
|
|
|
@Override
|
2012-05-24 17:57:10 +02:00
|
|
|
public ChannelFuture leaveGroup(
|
|
|
|
InetSocketAddress multicastAddress, NetworkInterface networkInterface) {
|
2012-12-30 17:40:24 +01:00
|
|
|
return leaveGroup(multicastAddress, networkInterface, newPromise());
|
2009-06-11 08:10:46 +02:00
|
|
|
}
|
2012-02-18 23:02:56 +01:00
|
|
|
|
2012-05-24 17:57:10 +02:00
|
|
|
@Override
|
|
|
|
public ChannelFuture leaveGroup(
|
|
|
|
InetSocketAddress multicastAddress,
|
2012-12-30 17:40:24 +01:00
|
|
|
NetworkInterface networkInterface, ChannelPromise promise) {
|
|
|
|
return leaveGroup(multicastAddress.getAddress(), networkInterface, null, promise);
|
2012-05-24 17:57:10 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
|
|
|
public ChannelFuture leaveGroup(
|
|
|
|
InetAddress multicastAddress, NetworkInterface networkInterface, InetAddress source) {
|
2012-12-30 17:40:24 +01:00
|
|
|
return leaveGroup(multicastAddress, networkInterface, source, newPromise());
|
2012-05-24 17:57:10 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
|
|
|
public ChannelFuture leaveGroup(
|
|
|
|
InetAddress multicastAddress, NetworkInterface networkInterface, InetAddress source,
|
2012-12-30 17:40:24 +01:00
|
|
|
ChannelPromise promise) {
|
2013-01-11 06:03:27 +01:00
|
|
|
if (PlatformDependent.javaVersion() < 7) {
|
2012-04-02 11:07:11 +02:00
|
|
|
throw new UnsupportedOperationException();
|
2012-05-24 17:57:10 +02:00
|
|
|
}
|
|
|
|
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();
|
|
|
|
|
|
|
|
while (keyIt.hasNext()) {
|
|
|
|
MembershipKey key = keyIt.next();
|
|
|
|
if (networkInterface.equals(key.networkInterface())) {
|
2012-06-08 12:28:12 +02:00
|
|
|
if (source == null && key.sourceAddress() == null ||
|
|
|
|
source != null && source.equals(key.sourceAddress())) {
|
2012-05-24 17:57:10 +02:00
|
|
|
key.drop();
|
|
|
|
keyIt.remove();
|
|
|
|
}
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
|
|
|
}
|
2012-05-24 17:57:10 +02:00
|
|
|
if (keys.isEmpty()) {
|
|
|
|
memberships.remove(multicastAddress);
|
|
|
|
}
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2012-05-24 17:57:10 +02:00
|
|
|
|
2012-12-30 17:40:24 +01:00
|
|
|
promise.setSuccess();
|
|
|
|
return promise;
|
2012-05-24 17:57:10 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
/**
|
|
|
|
* Block the given sourceToBlock address for the given multicastAddress on the given networkInterface
|
|
|
|
*/
|
|
|
|
@Override
|
|
|
|
public ChannelFuture block(
|
|
|
|
InetAddress multicastAddress, NetworkInterface networkInterface,
|
|
|
|
InetAddress sourceToBlock) {
|
2012-12-30 17:40:24 +01:00
|
|
|
return block(multicastAddress, networkInterface, sourceToBlock, newPromise());
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
2012-05-24 17:57:10 +02:00
|
|
|
|
2012-04-02 11:07:11 +02:00
|
|
|
/**
|
|
|
|
* Block the given sourceToBlock address for the given multicastAddress on the given networkInterface
|
|
|
|
*/
|
2012-05-24 17:57:10 +02:00
|
|
|
@Override
|
|
|
|
public ChannelFuture block(
|
|
|
|
InetAddress multicastAddress, NetworkInterface networkInterface,
|
2012-12-30 17:40:24 +01:00
|
|
|
InetAddress sourceToBlock, ChannelPromise promise) {
|
2013-01-11 06:03:27 +01:00
|
|
|
if (PlatformDependent.javaVersion() < 7) {
|
2012-04-02 11:07:11 +02:00
|
|
|
throw new UnsupportedOperationException();
|
|
|
|
} else {
|
|
|
|
if (multicastAddress == null) {
|
|
|
|
throw new NullPointerException("multicastAddress");
|
|
|
|
}
|
|
|
|
if (sourceToBlock == null) {
|
|
|
|
throw new NullPointerException("sourceToBlock");
|
|
|
|
}
|
2012-05-24 17:57:10 +02:00
|
|
|
|
2012-04-02 11:07:11 +02:00
|
|
|
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-12-30 17:40:24 +01:00
|
|
|
promise.setFailure(e);
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2012-12-30 17:40:24 +01:00
|
|
|
promise.setSuccess();
|
|
|
|
return promise;
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
|
|
|
}
|
2012-05-24 17:57:10 +02:00
|
|
|
|
2012-04-02 11:07:11 +02:00
|
|
|
/**
|
|
|
|
* Block the given sourceToBlock address for the given multicastAddress
|
2012-05-24 17:57:10 +02:00
|
|
|
*
|
2012-04-02 11:07:11 +02:00
|
|
|
*/
|
2012-05-24 17:57:10 +02:00
|
|
|
@Override
|
2012-04-02 11:57:32 +02:00
|
|
|
public ChannelFuture block(InetAddress multicastAddress, InetAddress sourceToBlock) {
|
2012-12-30 17:40:24 +01:00
|
|
|
return block(multicastAddress, sourceToBlock, newPromise());
|
2012-04-02 11:07:11 +02:00
|
|
|
}
|
2012-05-24 17:57:10 +02:00
|
|
|
|
|
|
|
/**
|
|
|
|
* Block the given sourceToBlock address for the given multicastAddress
|
|
|
|
*
|
|
|
|
*/
|
2012-03-20 09:43:00 +01:00
|
|
|
@Override
|
2012-05-24 17:57:10 +02:00
|
|
|
public ChannelFuture block(
|
2012-12-30 17:40:24 +01:00
|
|
|
InetAddress multicastAddress, InetAddress sourceToBlock, ChannelPromise promise) {
|
2012-05-24 17:57:10 +02:00
|
|
|
try {
|
|
|
|
return block(
|
|
|
|
multicastAddress,
|
|
|
|
NetworkInterface.getByInetAddress(localAddress().getAddress()),
|
2012-12-30 17:40:24 +01:00
|
|
|
sourceToBlock, promise);
|
2012-05-24 17:57:10 +02:00
|
|
|
} catch (SocketException e) {
|
2012-12-30 17:40:24 +01:00
|
|
|
promise.setFailure(e);
|
2012-03-20 09:43:00 +01:00
|
|
|
}
|
2012-12-30 17:40:24 +01:00
|
|
|
return promise;
|
2012-03-20 09:43:00 +01:00
|
|
|
}
|
2009-06-11 02:48:49 +02:00
|
|
|
}
|