Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
22 changed files
with
348 additions
and
30 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,45 @@ | ||
<?xml version="1.0" encoding="UTF-8"?> | ||
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/maven-v4_0_0.xsd"> | ||
<parent> | ||
<groupId>org.opennms.features.telemetry</groupId> | ||
<artifactId>org.opennms.features.telemetry.adapters</artifactId> | ||
<version>22.0.0-SNAPSHOT</version> | ||
</parent> | ||
<modelVersion>4.0.0</modelVersion> | ||
<groupId>org.opennms.features.telemetry.adapters</groupId> | ||
<artifactId>org.opennms.features.telemetry.adapters.flow</artifactId> | ||
<name>OpenNMS :: Features :: Telemetry :: Adapters :: Flow</name> | ||
<packaging>bundle</packaging> | ||
<build> | ||
<plugins> | ||
<plugin> | ||
<groupId>org.apache.felix</groupId> | ||
<artifactId>maven-bundle-plugin</artifactId> | ||
<extensions>true</extensions> | ||
<configuration> | ||
<instructions> | ||
<Bundle-RequiredExecutionEnvironment>JavaSE-1.8</Bundle-RequiredExecutionEnvironment> | ||
<Bundle-SymbolicName>${project.artifactId}</Bundle-SymbolicName> | ||
<Bundle-Version>${project.version}</Bundle-Version> | ||
</instructions> | ||
</configuration> | ||
</plugin> | ||
</plugins> | ||
</build> | ||
<dependencies> | ||
<dependency> | ||
<groupId>org.opennms.features.telemetry.adapters</groupId> | ||
<artifactId>org.opennms.features.telemetry.adapters.collection</artifactId> | ||
<version>${project.version}</version> | ||
</dependency> | ||
<dependency> | ||
<groupId>org.osgi</groupId> | ||
<artifactId>org.osgi.core</artifactId> | ||
<scope>provided</scope> | ||
</dependency> | ||
<dependency> | ||
<groupId>org.mongodb</groupId> | ||
<artifactId>bson</artifactId> | ||
</dependency> | ||
</dependencies> | ||
</project> |
148 changes: 148 additions & 0 deletions
148
...y/adapters/flow/src/main/java/org/opennms/netmgt/telemetry/adapters/flow/FlowAdapter.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,148 @@ | ||
/******************************************************************************* | ||
* This file is part of OpenNMS(R). | ||
* | ||
* Copyright (C) 2017-2017 The OpenNMS Group, Inc. | ||
* OpenNMS(R) is Copyright (C) 1999-2017 The OpenNMS Group, Inc. | ||
* | ||
* OpenNMS(R) is a registered trademark of The OpenNMS Group, Inc. | ||
* | ||
* OpenNMS(R) is free software: you can redistribute it and/or modify | ||
* it under the terms of the GNU Affero General Public License as published | ||
* by the Free Software Foundation, either version 3 of the License, | ||
* or (at your option) any later version. | ||
* | ||
* OpenNMS(R) is distributed in the hope that it will be useful, | ||
* but WITHOUT ANY WARRANTY; without even the implied warranty of | ||
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the | ||
* GNU Affero General Public License for more details. | ||
* | ||
* You should have received a copy of the GNU Affero General Public License | ||
* along with OpenNMS(R). If not, see: | ||
* http://www.gnu.org/licenses/ | ||
* | ||
* For more information contact: | ||
* OpenNMS(R) Licensing <license@opennms.org> | ||
* http://www.opennms.org/ | ||
* http://www.opennms.com/ | ||
*******************************************************************************/ | ||
|
||
package org.opennms.netmgt.telemetry.adapters.flow; | ||
|
||
import java.io.File; | ||
import java.net.InetAddress; | ||
import java.net.UnknownHostException; | ||
import java.util.Optional; | ||
|
||
import org.bson.RawBsonDocument; | ||
import org.opennms.netmgt.collection.api.CollectionAgent; | ||
import org.opennms.netmgt.collection.api.CollectionAgentFactory; | ||
import org.opennms.netmgt.collection.api.CollectionSet; | ||
import org.opennms.netmgt.dao.api.InterfaceToNodeCache; | ||
import org.opennms.netmgt.dao.api.NodeDao; | ||
import org.opennms.netmgt.telemetry.adapters.api.TelemetryMessage; | ||
import org.opennms.netmgt.telemetry.adapters.api.TelemetryMessageLog; | ||
import org.opennms.netmgt.telemetry.adapters.collection.AbstractPersistingAdapter; | ||
import org.opennms.netmgt.telemetry.adapters.collection.CollectionSetWithAgent; | ||
import org.opennms.netmgt.telemetry.adapters.collection.ScriptedCollectionSetBuilder; | ||
import org.osgi.framework.BundleContext; | ||
import org.slf4j.Logger; | ||
import org.slf4j.LoggerFactory; | ||
import org.springframework.beans.factory.annotation.Autowired; | ||
import org.springframework.transaction.support.TransactionOperations; | ||
|
||
public class FlowAdapter extends AbstractPersistingAdapter { | ||
|
||
private static final Logger LOG = LoggerFactory.getLogger(FlowAdapter.class); | ||
|
||
@Autowired | ||
private CollectionAgentFactory collectionAgentFactory; | ||
|
||
@Autowired | ||
private InterfaceToNodeCache interfaceToNodeCache; | ||
|
||
@Autowired | ||
private NodeDao nodeDao; | ||
|
||
@Autowired | ||
private TransactionOperations transactionTemplate; | ||
|
||
private BundleContext bundleContext; | ||
|
||
private String script; | ||
|
||
private final ThreadLocal<ScriptedCollectionSetBuilder> scriptedCollectionSetBuilders = new ThreadLocal<ScriptedCollectionSetBuilder>() { | ||
@Override | ||
protected ScriptedCollectionSetBuilder initialValue() { | ||
try { | ||
if (bundleContext != null) { | ||
return new ScriptedCollectionSetBuilder(new File(script), bundleContext); | ||
} else { | ||
return new ScriptedCollectionSetBuilder(new File(script)); | ||
} | ||
} catch (Exception e) { | ||
LOG.error("Failed to create builder for script '{}'.", script, e); | ||
return null; | ||
} | ||
} | ||
}; | ||
|
||
public String getScript() { | ||
return this.script; | ||
} | ||
|
||
public void setScript(String script) { | ||
this.script = script; | ||
} | ||
|
||
@Override | ||
public Optional<CollectionSetWithAgent> handleMessage(final TelemetryMessage message, final TelemetryMessageLog messageLog) throws Exception { | ||
final RawBsonDocument flow = new RawBsonDocument(message.getByteArray()); | ||
|
||
LOG.warn("Flow: {}", flow.toJson()); | ||
|
||
CollectionAgent agent = null; | ||
try { | ||
final InetAddress inetAddress = InetAddress.getByName(messageLog.getSourceAddress()); | ||
final Optional<Integer> nodeId = this.interfaceToNodeCache.getFirstNodeId(messageLog.getLocation(), inetAddress); | ||
if (nodeId.isPresent()) { | ||
// NOTE: This will throw a IllegalArgumentException if the nodeId/inetAddress pair does not exist in the database | ||
agent = this.collectionAgentFactory.createCollectionAgent(Integer.toString(nodeId.get()), inetAddress); | ||
} | ||
} catch (UnknownHostException e) { | ||
LOG.debug("Could not convert source address: {}", messageLog.getSourceAddress()); | ||
} | ||
|
||
if (agent == null) { | ||
LOG.warn("Unable to find node for address: {}", messageLog.getSourceAddress()); | ||
return Optional.empty(); | ||
} | ||
|
||
final ScriptedCollectionSetBuilder builder = this.scriptedCollectionSetBuilders.get(); | ||
if (builder == null) { | ||
throw new Exception(String.format("Error compiling script '%s'. See logs for details.", script)); | ||
} | ||
final CollectionSet collectionSet = builder.build(agent, flow); | ||
return Optional.of(new CollectionSetWithAgent(agent, collectionSet)); | ||
} | ||
|
||
public void setCollectionAgentFactory(CollectionAgentFactory collectionAgentFactory) { | ||
this.collectionAgentFactory = collectionAgentFactory; | ||
} | ||
|
||
public void setInterfaceToNodeCache(InterfaceToNodeCache interfaceToNodeCache) { | ||
this.interfaceToNodeCache = interfaceToNodeCache; | ||
} | ||
|
||
public void setNodeDao(NodeDao nodeDao) { | ||
this.nodeDao = nodeDao; | ||
} | ||
|
||
public void setTransactionTemplate(TransactionOperations transactionTemplate) { | ||
this.transactionTemplate = transactionTemplate; | ||
} | ||
|
||
public void setBundleContext(BundleContext bundleContext) { | ||
this.bundleContext = bundleContext; | ||
} | ||
|
||
} |
68 changes: 68 additions & 0 deletions
68
...ers/flow/src/main/java/org/opennms/netmgt/telemetry/adapters/flow/FlowAdapterFactory.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,68 @@ | ||
/******************************************************************************* | ||
* This file is part of OpenNMS(R). | ||
* | ||
* Copyright (C) 2017-2017 The OpenNMS Group, Inc. | ||
* OpenNMS(R) is Copyright (C) 1999-2017 The OpenNMS Group, Inc. | ||
* | ||
* OpenNMS(R) is a registered trademark of The OpenNMS Group, Inc. | ||
* | ||
* OpenNMS(R) is free software: you can redistribute it and/or modify | ||
* it under the terms of the GNU Affero General Public License as published | ||
* by the Free Software Foundation, either version 3 of the License, | ||
* or (at your option) any later version. | ||
* | ||
* OpenNMS(R) is distributed in the hope that it will be useful, | ||
* but WITHOUT ANY WARRANTY; without even the implied warranty of | ||
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the | ||
* GNU Affero General Public License for more details. | ||
* | ||
* You should have received a copy of the GNU Affero General Public License | ||
* along with OpenNMS(R). If not, see: | ||
* http://www.gnu.org/licenses/ | ||
* | ||
* For more information contact: | ||
* OpenNMS(R) Licensing <license@opennms.org> | ||
* http://www.opennms.org/ | ||
* http://www.opennms.com/ | ||
*******************************************************************************/ | ||
|
||
package org.opennms.netmgt.telemetry.adapters.flow; | ||
|
||
import java.util.Map; | ||
|
||
import org.opennms.netmgt.telemetry.adapters.api.Adapter; | ||
import org.opennms.netmgt.telemetry.adapters.collection.AbstractCollectionAdapterFactory; | ||
import org.opennms.netmgt.telemetry.config.api.Protocol; | ||
import org.osgi.framework.BundleContext; | ||
import org.springframework.beans.BeanWrapper; | ||
import org.springframework.beans.PropertyAccessorFactory; | ||
|
||
public class FlowAdapterFactory extends AbstractCollectionAdapterFactory { | ||
|
||
public FlowAdapterFactory(BundleContext bundleContext) { | ||
super(bundleContext); | ||
} | ||
|
||
@Override | ||
public Class<? extends Adapter> getAdapterClass() { | ||
return FlowAdapter.class; | ||
} | ||
|
||
@Override | ||
public Adapter createAdapter(Protocol protocol, Map<String, String> properties) { | ||
final FlowAdapter adapter = new FlowAdapter(); | ||
adapter.setProtocol(protocol); | ||
adapter.setCollectionAgentFactory(getCollectionAgentFactory()); | ||
adapter.setInterfaceToNodeCache(getInterfaceToNodeCache()); | ||
adapter.setNodeDao(getNodeDao()); | ||
adapter.setTransactionTemplate(getTransactionTemplate()); | ||
adapter.setFilterDao(getFilterDao()); | ||
adapter.setPersisterFactory(getPersisterFactory()); | ||
adapter.setBundleContext(getBundleContext()); | ||
|
||
final BeanWrapper wrapper = PropertyAccessorFactory.forBeanPropertyAccess(adapter); | ||
wrapper.setPropertyValues(properties); | ||
return adapter; | ||
} | ||
|
||
} |
36 changes: 36 additions & 0 deletions
36
features/telemetry/adapters/flow/src/main/resources/OSGI-INF/blueprint/blueprint.xml
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,36 @@ | ||
<blueprint xmlns="http://www.osgi.org/xmlns/blueprint/v1.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" | ||
xmlns:cm="http://aries.apache.org/blueprint/xmlns/blueprint-cm/v1.1.0" xmlns:ext="http://aries.apache.org/blueprint/xmlns/blueprint-ext/v1.1.0" | ||
xsi:schemaLocation=" | ||
http://www.osgi.org/xmlns/blueprint/v1.0.0 | ||
http://www.osgi.org/xmlns/blueprint/v1.0.0/blueprint.xsd | ||
http://aries.apache.org/blueprint/xmlns/blueprint-cm/v1.1.0 | ||
http://aries.apache.org/schemas/blueprint-cm/blueprint-cm-1.1.0.xsd | ||
http://aries.apache.org/blueprint/xmlns/blueprint-ext/v1.1.0 | ||
http://aries.apache.org/schemas/blueprint-ext/blueprint-ext-1.1.xsd | ||
"> | ||
|
||
<bean id="flowAdapterFactory" class="org.opennms.netmgt.telemetry.adapters.flow.FlowAdapterFactory"> | ||
<argument ref="blueprintBundleContext" /> | ||
<property name="collectionAgentFactory" ref="collectionAgentFactory" /> | ||
<property name="interfaceToNodeCache" ref="interfaceToNodeCache" /> | ||
<property name="nodeDao" ref="nodeDao" /> | ||
<property name="transactionTemplate" ref="transactionTemplate" /> | ||
<property name="filterDao" ref="filterDao" /> | ||
<property name="persisterFactory" ref="persisterFactory" /> | ||
</bean> | ||
|
||
<service id="flowFactoryService" ref="flowAdapterFactory" interface="org.opennms.netmgt.telemetry.adapters.api.AdapterFactory"> | ||
<service-properties> | ||
<entry key="registration.export" value="true" /> | ||
<entry key="type" value="org.opennms.netmgt.telemetry.adapters.flow.FlowAdapter" /> | ||
</service-properties> | ||
</service> | ||
|
||
<reference id="collectionAgentFactory" interface="org.opennms.netmgt.collection.api.CollectionAgentFactory" /> | ||
<reference id="interfaceToNodeCache" interface="org.opennms.netmgt.dao.api.InterfaceToNodeCache" /> | ||
<reference id="nodeDao" interface="org.opennms.netmgt.dao.api.NodeDao" /> | ||
<reference id="filterDao" interface="org.opennms.netmgt.filter.api.FilterDao" /> | ||
<reference id="transactionTemplate" interface="org.springframework.transaction.support.TransactionOperations" /> | ||
<reference id="persisterFactory" interface="org.opennms.netmgt.collection.api.PersisterFactory" /> | ||
|
||
</blueprint> |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.