From ad552d4e2f7257837b1ea73ca6ca1f338cd80b34 Mon Sep 17 00:00:00 2001 From: Woonduk Kang Date: Thu, 11 Jun 2015 15:53:33 +0900 Subject: [PATCH] refactoring collector - UDPReceiver - extract L4CheckFilter, NetworkAvailabilityCheckPacketFilter --- .../receiver/udp/BaseUDPHandlerFactory.java | 39 +++------- .../udp/ChunkedUDPPacketHandlerFactory.java | 23 +++--- .../receiver/udp/L4PacketFilter.java | 44 +++++++++++ .../NetworkAvailabilityCheckPacketFilter.java | 73 +++++++++++++++++++ .../collector/receiver/udp/TBaseFilter.java | 39 ++++++++++ .../receiver/udp/TBaseFilterChain.java | 48 ++++++++++++ .../collector/receiver/udp/UDPReceiver.java | 18 +++-- 7 files changed, 234 insertions(+), 50 deletions(-) create mode 100644 collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/L4PacketFilter.java create mode 100644 collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/NetworkAvailabilityCheckPacketFilter.java create mode 100644 collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/TBaseFilter.java create mode 100644 collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/TBaseFilterChain.java diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/BaseUDPHandlerFactory.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/BaseUDPHandlerFactory.java index 002870960..02bf14f04 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/BaseUDPHandlerFactory.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/BaseUDPHandlerFactory.java @@ -27,7 +27,6 @@ import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; import org.springframework.util.Assert; -import java.io.IOException; import java.net.*; /** @@ -43,6 +42,8 @@ public class BaseUDPHandlerFactory implements PacketHa private UDPReceiver receiver; private final DispatchHandler dispatchHandler; + private TBaseFilter filter = TBaseFilter.CONTINUE_FILTER; + public BaseUDPHandlerFactory(DispatchHandler dispatchHandler) { if (dispatchHandler == null) { throw new NullPointerException("dispatchHandler must not be null"); @@ -54,6 +55,13 @@ public class BaseUDPHandlerFactory implements PacketHa this.receiver = receiver; } + public void setFilter(TBaseFilter filter) { + if (filter == null) { + throw new NullPointerException("filter must not be null"); + } + this.filter = filter; + } + @Override public PacketHandler createPacketHandler() { return new DispatchPacket(); @@ -75,19 +83,7 @@ public class BaseUDPHandlerFactory implements PacketHa TBase tBase = null; try { tBase = deserializer.deserialize(packet.getData()); - if (tBase instanceof L4Packet) { - if (logger.isDebugEnabled()) { - L4Packet l4Packet = (L4Packet) tBase; - logger.debug("udp l4 packet {}", l4Packet.getHeader()); - } - return; - } - // Network port availability check packet - if (tBase instanceof NetworkAvailabilityCheckPacket) { - if (logger.isDebugEnabled()) { - logger.debug("received udp network availability check packet."); - } - responseOK(packet); + if (filter.filter(tBase, packet) == TBaseFilter.BREAK) { return; } // dispatch signifies business logic execution @@ -109,21 +105,6 @@ public class BaseUDPHandlerFactory implements PacketHa } } } - - private void responseOK(DatagramPacket packet) { - try { - byte[] okBytes = NetworkAvailabilityCheckPacket.DATA_OK; - DatagramPacket pongPacket = new DatagramPacket(okBytes, okBytes.length, packet.getSocketAddress()); - receiver.getSocket().send(pongPacket); - } catch (IOException e) { - if (logger.isWarnEnabled()) { - logger.warn("pong error. SendSocketAddress:{} Cause:{}", packet.getSocketAddress(), e.getMessage(), e); - } - if (logger.isDebugEnabled()) { - logger.debug("packet dump hex:{}", PacketUtils.dumpDatagramPacket(packet)); - } - } - } } } diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/ChunkedUDPPacketHandlerFactory.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/ChunkedUDPPacketHandlerFactory.java index dab5dfb7f..6a526a724 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/ChunkedUDPPacketHandlerFactory.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/ChunkedUDPPacketHandlerFactory.java @@ -44,6 +44,7 @@ public class ChunkedUDPPacketHandlerFactory implements private UDPReceiver receiver; private final DispatchHandler dispatchHandler; + private TBaseFilter filter = TBaseFilter.CONTINUE_FILTER; public ChunkedUDPPacketHandlerFactory(DispatchHandler dispatchHandler) { if (dispatchHandler == null) { @@ -56,6 +57,12 @@ public class ChunkedUDPPacketHandlerFactory implements this.receiver = receiver; } + public void setFilter(TBaseFilter filter) { + if (filter == null) { + throw new NullPointerException("filter must not be null"); + } + this.filter = filter; + } @Override public void afterPropertiesSet() throws Exception { @@ -82,20 +89,8 @@ public class ChunkedUDPPacketHandlerFactory implements } for (TBase tBase : list) { - if (tBase instanceof L4Packet) { - if (logger.isDebugEnabled()) { - L4Packet l4Packet = (L4Packet) tBase; - logger.debug("udp l4 packet {}", l4Packet.getHeader()); - } - continue; - } - // Network port availability check packet - if (tBase instanceof NetworkAvailabilityCheckPacket) { - if (logger.isDebugEnabled()) { - logger.debug("received udp network availability check packet."); - } - responseOK(packet); - continue; + if (filter.filter(tBase, packet) == TBaseFilter.BREAK) { + return; } // dispatch signifies business logic execution dispatchHandler.dispatchSendMessage(tBase); diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/L4PacketFilter.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/L4PacketFilter.java new file mode 100644 index 000000000..1e5295b1e --- /dev/null +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/L4PacketFilter.java @@ -0,0 +1,44 @@ +/* + * Copyright 2014 NAVER Corp. + * + * Licensed 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 com.navercorp.pinpoint.collector.receiver.udp; + +import com.navercorp.pinpoint.thrift.io.L4Packet; +import org.apache.thrift.TBase; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.net.DatagramPacket; + +/** + * @author emeroad + */ +public class L4PacketFilter implements TBaseFilter { + + private final Logger logger = LoggerFactory.getLogger(this.getClass()); + + @Override + public boolean filter(TBase tBase, DatagramPacket packet) { + if (tBase instanceof L4Packet) { + if (logger.isDebugEnabled()) { + L4Packet l4Packet = (L4Packet) tBase; + logger.debug("udp l4 packet {}", l4Packet.getHeader()); + } + return BREAK; + } + return CONTINUE; + } +} diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/NetworkAvailabilityCheckPacketFilter.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/NetworkAvailabilityCheckPacketFilter.java new file mode 100644 index 000000000..d4df95478 --- /dev/null +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/NetworkAvailabilityCheckPacketFilter.java @@ -0,0 +1,73 @@ +/* + * Copyright 2014 NAVER Corp. + * + * Licensed 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 com.navercorp.pinpoint.collector.receiver.udp; + +import com.navercorp.pinpoint.collector.util.PacketUtils; +import com.navercorp.pinpoint.thrift.io.NetworkAvailabilityCheckPacket; +import org.apache.thrift.TBase; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.IOException; +import java.net.DatagramPacket; +import java.net.DatagramSocket; +import java.net.SocketException; + +/** + * @author emeroad + */ +public class NetworkAvailabilityCheckPacketFilter implements TBaseFilter { + private final Logger logger = LoggerFactory.getLogger(this.getClass()); + + private final DatagramSocket socket; + + public NetworkAvailabilityCheckPacketFilter() { + try { + this.socket = new DatagramSocket(); + } catch (SocketException ex) { + throw new RuntimeException("socket create fail. error:" + ex.getMessage(), ex); + } + } + + @Override + public boolean filter(TBase tBase, DatagramPacket packet) { + // Network port availability check packet + if (tBase instanceof NetworkAvailabilityCheckPacket) { + if (logger.isInfoEnabled()) { + logger.info("received udp network availability check packet."); + } + responseOK(packet); + return BREAK; + } + return CONTINUE; + } + + private void responseOK(DatagramPacket packet) { + try { + byte[] okBytes = NetworkAvailabilityCheckPacket.DATA_OK; + DatagramPacket pongPacket = new DatagramPacket(okBytes, okBytes.length, packet.getSocketAddress()); + socket.send(pongPacket); + } catch (IOException e) { + if (logger.isWarnEnabled()) { + logger.warn("pong error. SendSocketAddress:{} Cause:{}", packet.getSocketAddress(), e.getMessage(), e); + } + if (logger.isDebugEnabled()) { + logger.debug("packet dump hex:{}", PacketUtils.dumpDatagramPacket(packet)); + } + } + } +} diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/TBaseFilter.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/TBaseFilter.java new file mode 100644 index 000000000..ce3c31b2e --- /dev/null +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/TBaseFilter.java @@ -0,0 +1,39 @@ +/* + * Copyright 2014 NAVER Corp. + * + * Licensed 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 com.navercorp.pinpoint.collector.receiver.udp; + +import org.apache.thrift.TBase; + +import java.net.DatagramPacket; + +/** + * @author emeroad + */ +public interface TBaseFilter { + boolean CONTINUE = true; + boolean BREAK = false; + + boolean filter(TBase tBase, DatagramPacket packet); + + + public static final TBaseFilter CONTINUE_FILTER = new TBaseFilter() { + @Override + public boolean filter(TBase tBase, DatagramPacket packet) { + return CONTINUE; + } + }; +} diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/TBaseFilterChain.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/TBaseFilterChain.java new file mode 100644 index 000000000..f5e6f0c1d --- /dev/null +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/TBaseFilterChain.java @@ -0,0 +1,48 @@ +/* + * Copyright 2014 NAVER Corp. + * + * Licensed 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 com.navercorp.pinpoint.collector.receiver.udp; + +import org.apache.thrift.TBase; + +import java.net.DatagramPacket; +import java.util.ArrayList; +import java.util.List; + +/** + * @author emeroad + */ +public class TBaseFilterChain implements TBaseFilter { + + private final List filterChain = new ArrayList(); + + public void addTBaseFilter(TBaseFilter tBaseFilter) { + if (tBaseFilter == null) { + throw new NullPointerException("tBaseFilter must not be null"); + } + this.filterChain.add(tBaseFilter); + } + + @Override + public boolean filter(TBase tBase, DatagramPacket packet) { + for (TBaseFilter tBaseFilter : filterChain) { + if (tBaseFilter.filter(tBase, packet) == TBaseFilter.BREAK) { + return BREAK; + } + } + return TBaseFilter.CONTINUE; + } +} diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/UDPReceiver.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/UDPReceiver.java index da5bc5e0b..31abb82be 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/UDPReceiver.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/UDPReceiver.java @@ -142,7 +142,7 @@ public class UDPReceiver implements DataReceiver { if (pooledPacket == null) { continue; } - DatagramPacket packet = pooledPacket.getObject(); + final DatagramPacket packet = pooledPacket.getObject(); if (packet.getLength() == 0) { if (debugEnabled) { logger.debug("length is 0 ip:{}, port:{}", packet.getAddress(), packet.getPort()); @@ -156,12 +156,7 @@ public class UDPReceiver implements DataReceiver { Runnable dispatchTask = wrapDispatchTask(pooledPacket); worker.execute(dispatchTask); } catch (RejectedExecutionException ree) { - rejectedCounter.inc(); - final int error = rejectedExecutionCount.incrementAndGet(); - final int mod = 100; - if ((error % mod) == 0) { - logger.warn("RejectedExecutionCount={}", error); - } + handleRejectedExecutionException(ree); } } if (logger.isInfoEnabled()) { @@ -169,6 +164,15 @@ public class UDPReceiver implements DataReceiver { } } + private void handleRejectedExecutionException(RejectedExecutionException ree) { + rejectedCounter.inc(); + final int error = rejectedExecutionCount.incrementAndGet(); + final int mod = 100; + if ((error % mod) == 0) { + logger.warn("RejectedExecutionCount={}", error); + } + } + private Runnable wrapDispatchTask(PooledObject pooledPacket) { PacketHandler dispatchPacket = packetHandlerFactory.createPacketHandler(); PooledPacketWrap pooledPacketWrap = new PooledPacketWrap(dispatchPacket, pooledPacket);