This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12011-6a0f36278d31563d9a9c9243f899b1c522cf3a2d in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 1356ec14b9bd7ffd5fbfe2f4d20013816148f3ca Author: Goutam Adwant <[email protected]> AuthorDate: Tue Sep 8 16:39:52 2026 +0000 [Feature][Connector-V2] Add SNMPv2c SET sink connector (#12011) Signed-off-by: goutamadwant <[email protected]> --- .github/workflows/labeler/label-scope-conf.yml | 7 + docs/en/connectors/changelog/connector-snmp.md | 1 + docs/en/connectors/sink/SNMP.md | 137 +++++++++ docs/zh/connectors/changelog/connector-snmp.md | 1 + docs/zh/connectors/sink/SNMP.md | 135 +++++++++ plugin-mapping.properties | 1 + .../seatunnel/snmp/client/SnmpTargetFactory.java | 51 ++++ .../{SnmpSourceOptions.java => SnmpOptions.java} | 19 +- .../seatunnel/snmp/config/SnmpSinkConfig.java | 131 +++++++++ .../seatunnel/snmp/config/SnmpSinkOptions.java | 53 ++++ .../seatunnel/snmp/config/SnmpSourceConfig.java | 39 +-- .../seatunnel/snmp/config/SnmpSourceOptions.java | 29 +- .../SnmpTargetConfig.java} | 29 +- .../snmp/exception/SnmpConnectorErrorCode.java | 6 +- .../Snmp4jSetClient.java} | 67 ++--- .../SnmpSetClient.java} | 27 +- .../SnmpSetRequest.java} | 29 +- .../connectors/seatunnel/snmp/sink/SnmpSink.java | 59 ++++ .../seatunnel/snmp/sink/SnmpSinkFactory.java | 70 +++++ .../seatunnel/snmp/sink/SnmpSinkRowConverter.java | 316 +++++++++++++++++++++ .../seatunnel/snmp/sink/SnmpSinkWriter.java | 93 ++++++ .../seatunnel/snmp/source/Snmp4jClient.java | 13 +- .../snmp/client/SnmpTargetFactoryTest.java | 47 +++ .../seatunnel/snmp/config/SnmpSinkConfigTest.java | 119 ++++++++ .../snmp/config/SnmpSourceConfigTest.java | 34 +++ .../seatunnel/snmp/sink/Snmp4jSetClientTest.java | 160 +++++++++++ .../seatunnel/snmp/sink/SnmpSinkFactoryTest.java | 189 ++++++++++++ .../snmp/sink/SnmpSinkRowConverterTest.java | 246 ++++++++++++++++ .../seatunnel/snmp/sink/SnmpSinkWriterTest.java | 133 +++++++++ .../core/starter/utils/ConfigShadeUtils.java | 2 +- .../core/starter/utils/ConfigBuilderTest.java | 19 ++ .../connector-snmp-e2e/pom.xml | 42 +++ .../seatunnel/e2e/connector/snmp/SnmpAgent.java | 95 +++++++ .../seatunnel/e2e/connector/snmp/SnmpSinkIT.java | 109 +++++++ .../src/test/resources/fake_to_snmp.conf | 50 ++++ seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml | 1 + 36 files changed, 2384 insertions(+), 175 deletions(-) diff --git a/.github/workflows/labeler/label-scope-conf.yml b/.github/workflows/labeler/label-scope-conf.yml index 00c7cf4609..141c3f8469 100644 --- a/.github/workflows/labeler/label-scope-conf.yml +++ b/.github/workflows/labeler/label-scope-conf.yml @@ -179,6 +179,13 @@ google-pubsub: - changed-files: - any-glob-to-any-file: seatunnel-connectors-v2/connector-google-pubsub/** - all-globs-to-all-files: '!seatunnel-connectors-v2/connector-!(google-pubsub)/**' +snmp: + - all: + - changed-files: + - any-glob-to-any-file: + - seatunnel-connectors-v2/connector-snmp/** + - seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/** + - all-globs-to-all-files: '!seatunnel-connectors-v2/connector-!(snmp)/**' azure-queue-storage: - all: - changed-files: diff --git a/docs/en/connectors/changelog/connector-snmp.md b/docs/en/connectors/changelog/connector-snmp.md index f70bbecaab..832e97161d 100644 --- a/docs/en/connectors/changelog/connector-snmp.md +++ b/docs/en/connectors/changelog/connector-snmp.md @@ -3,5 +3,6 @@ | Change | Commit | Version | | --- | --- | --- | | [Feature][Connector-V2] Add SNMPv2c polling source connector | - | Next | +| [Feature][Connector-V2] Add SNMPv2c SET sink connector | - | Next | </details> diff --git a/docs/en/connectors/sink/SNMP.md b/docs/en/connectors/sink/SNMP.md new file mode 100644 index 0000000000..09069855ab --- /dev/null +++ b/docs/en/connectors/sink/SNMP.md @@ -0,0 +1,137 @@ +import ChangeLog from '../changelog/connector-snmp.md'; + +# SNMP + +> SNMPv2c SET sink connector + +## Description + +The SNMP sink writes each input row to one SNMP agent by sending one synchronous SNMPv2c SET request. +The V1 scope is deliberately limited to SET operations. It does not send traps or informs, and it does not support SNMPv1 or SNMPv3. + +Every row supplies a numeric OID, a string value, and an SMI value type. The corresponding field names are configurable. +The default mapping consumes the `oid`, `value`, and `value_type` fields emitted by the SNMP source; additional fields such as +`agent` and `poll_time` are ignored by the sink. + +## Support Those Engines + +> Spark<br/> +> Flink<br/> +> SeaTunnel Zeta<br/> + +## Key Features + +- [x] [batch](../../introduction/concepts/connector-v2-features.md) +- [x] [stream](../../introduction/concepts/connector-v2-features.md) +- [x] [parallelism](../../introduction/concepts/connector-v2-features.md) + +## Supported DataSource Info + +The connector uses SNMP4J and supports SNMPv2c agents reachable over UDP. + +| Datasource | Supported Versions | Dependency | +|------------|--------------------|------------| +| SNMP agent | SNMPv2c | [Download](https://mvnrepository.com/artifact/org.apache.seatunnel/connector-snmp) | + +## Sink Options + +| Name | Type | Required | Default | Description | +|------------------|--------|----------|--------------|-------------| +| host | String | Yes | - | SNMP agent host name or IP address. Do not include a protocol or port. | +| port | Int | No | 161 | SNMP agent UDP port. | +| community | String | Yes | - | SNMPv2c community credential. The connector does not write this value to its logs or errors. | +| timeout_millis | Long | No | 5000 | Timeout in milliseconds for each SET request attempt. | +| retries | Int | No | 1 | Number of retries after the initial SET request attempt. A value of `0` sends one attempt. | +| oid_field | String | No | oid | Input `STRING` field containing the numeric OID to set. | +| value_field | String | No | value | Input `STRING` field containing the value to set. | +| value_type_field | String | No | value_type | Input `STRING` field containing the SMI value type. | +| common-options | | No | - | [Common Sink Options](../common-options/sink-common-options.md), including `plugin_input`. | + +Factory option validation checks nonblank strings, the port range, a positive timeout, and nonnegative retries before sink construction. +The three mapped fields must exist in the input schema, must use `STRING`, and must refer to distinct fields. Schema errors are +rejected while the job is created. Null values and blank OID or value-type fields are rejected before a network request is sent. The value field is validated according to its SMI type; an empty `OctetString` or `OctetStringHex` is valid, and text `OctetString` whitespace is preserved. + +## Supported SMI Value Types + +The `value_type` comparison is case-insensitive and ignores whitespace, `_`, and `-` characters. +The sink accepts both the documented names and SNMP4J syntax strings emitted by the SNMP source, +including `Counter`, `Gauge`, `OCTET STRING`, and `OBJECT IDENTIFIER`. + +| Value type | Accepted value | +|------------|----------------| +| `Integer32` or `Integer` | Signed 32-bit decimal integer. | +| `UnsignedInteger32` or `UnsignedInteger` | Decimal integer from 0 through 4294967295. | +| `Counter32` or `Counter` | Decimal integer from 0 through 4294967295. | +| `Gauge32` or `Gauge` | Decimal integer from 0 through 4294967295. | +| `TimeTicks` | Decimal count of hundredths of a second from 0 through 4294967295, or the SNMP4J source format `[days, ]hours:mm:ss.hh`. | +| `Counter64` | Decimal integer from 0 through 18446744073709551615. | +| `OctetString` or `OCTET STRING` | UTF-8 text represented by the input string. | +| `OctetStringHex` | An even number of hexadecimal characters, such as `00ff10`. | +| `OID` or `OBJECT IDENTIFIER` | Numeric object identifier. Leading dots are accepted. | +| `IpAddress` | Dotted IPv4 address. | + +`OctetString` is a textual mapping. Use `OctetStringHex` when byte-for-byte binary content is required. + +## Example + +```hocon +env { + parallelism = 1 + job.mode = "BATCH" +} + +source { + FakeSource { + plugin_output = "snmp_updates" + schema = { + fields { + oid = string + value = string + value_type = string + } + } + rows = [ + { + kind = INSERT + fields = { + oid = "1.3.6.1.2.1.1.5.0" + value = "router-1" + value_type = "OctetString" + } + } + ] + } +} + +sink { + SNMP { + plugin_input = "snmp_updates" + host = "192.0.2.10" + port = 161 + community = "replace-with-your-community" + timeout_millis = 3000 + retries = 1 + } +} +``` + +Replace the placeholder community before running the example. Supply the real credential outside checked-in job files. +`community` is automatically masked when the parsed job configuration is logged. There is no need to add it to +`shade.options` for log masking; that option also participates in configuration shading/encryption. + +## Delivery, Failure, and Security Behavior + +- One successful `write` call means the agent returned a successful SNMP response for that row. +- A timeout after all configured attempts or a non-zero SNMP response error status fails the sink task. +- A row can block for approximately `timeout_millis * (retries + 1)` before it fails. Keep this below the job's checkpoint timeout. +- SNMP4J retransmits a timed-out request. A late response can therefore make a non-idempotent OID observe the same SET more than once. +- The sink has no transactional commit protocol or recoverable writer state. Engine recovery can repeat a SET request, so delivery is at-least-once. +- Parallel writers can update the same OID out of order. Use parallelism 1 when update order matters. +- Row kinds are not interpreted as CDC operations. Every input row, including update or delete row kinds, is treated as a SET request. +- Treat `community` as a credential. Supply it through configuration substitution or another secret-management path, and do not place a real value in job files committed to source control. +- SNMPv2c provides no wire encryption or integrity protection. The community and SET payload are sent in cleartext; use only a trusted private network or a protected tunnel such as a VPN. +- Traps, informs, SNMPv1, and SNMPv3 are outside this V1 contract. + +See [Common Sink Options](../common-options/sink-common-options.md) for options such as `plugin_input`. + +<ChangeLog /> diff --git a/docs/zh/connectors/changelog/connector-snmp.md b/docs/zh/connectors/changelog/connector-snmp.md index f70bbecaab..832e97161d 100644 --- a/docs/zh/connectors/changelog/connector-snmp.md +++ b/docs/zh/connectors/changelog/connector-snmp.md @@ -3,5 +3,6 @@ | Change | Commit | Version | | --- | --- | --- | | [Feature][Connector-V2] Add SNMPv2c polling source connector | - | Next | +| [Feature][Connector-V2] Add SNMPv2c SET sink connector | - | Next | </details> diff --git a/docs/zh/connectors/sink/SNMP.md b/docs/zh/connectors/sink/SNMP.md new file mode 100644 index 0000000000..f23b0b9702 --- /dev/null +++ b/docs/zh/connectors/sink/SNMP.md @@ -0,0 +1,135 @@ +import ChangeLog from '../changelog/connector-snmp.md'; + +# SNMP + +> SNMPv2c SET Sink 连接器 + +## 描述 + +SNMP Sink 为每一行输入向一个 SNMP Agent 发送一次同步 SNMPv2c SET 请求。 +V1 范围仅包括 SET 操作,不发送 Trap 或 Inform,也不支持 SNMPv1 或 SNMPv3。 + +每一行需要提供数字 OID、字符串值和 SMI 值类型,对应的字段名可以配置。 +默认映射会读取 SNMP Source 输出的 `oid`、`value` 和 `value_type` 字段;Sink 会忽略 `agent`、`poll_time` 等额外字段。 + +## 支持的引擎 + +> Spark<br/> +> Flink<br/> +> SeaTunnel Zeta<br/> + +## 主要特性 + +- [x] [批处理](../../introduction/concepts/connector-v2-features.md) +- [x] [流处理](../../introduction/concepts/connector-v2-features.md) +- [x] [并行度](../../introduction/concepts/connector-v2-features.md) + +## 支持的数据源信息 + +连接器使用 SNMP4J,支持通过 UDP 访问的 SNMPv2c Agent。 + +| 数据源 | 支持版本 | 依赖 | +|--------|----------|------| +| SNMP Agent | SNMPv2c | [下载](https://mvnrepository.com/artifact/org.apache.seatunnel/connector-snmp) | + +## Sink 配置项 + +| 名称 | 类型 | 是否必填 | 默认值 | 描述 | +|------------------|--------|----------|------------|------| +| host | String | 是 | - | SNMP Agent 主机名或 IP 地址,不要包含协议或端口。 | +| port | Int | 否 | 161 | SNMP Agent 的 UDP 端口。 | +| community | String | 是 | - | SNMPv2c community 凭证。连接器不会把该值写入日志或错误信息。 | +| timeout_millis | Long | 否 | 5000 | 每次 SET 请求尝试的超时时间,单位为毫秒。 | +| retries | Int | 否 | 1 | 首次 SET 请求失败后的重试次数。`0` 表示只发送一次。 | +| oid_field | String | 否 | oid | 包含待设置数字 OID 的输入 `STRING` 字段。 | +| value_field | String | 否 | value | 包含待设置值的输入 `STRING` 字段。 | +| value_type_field | String | 否 | value_type | 包含 SMI 值类型的输入 `STRING` 字段。 | +| common-options | | 否 | - | [通用 Sink 配置项](../common-options/sink-common-options.md),包括 `plugin_input`。 | + +工厂配置校验会在构造 Sink 之前检查字符串非空、端口范围、超时时间为正数以及重试次数非负。 +三个映射字段必须存在于输入 Schema 中、类型必须为 `STRING`,并且不能指向同一个字段。Schema 错误会在创建任务时被拒绝。 +空值以及空白的 OID 或值类型字段会在发送网络请求前被拒绝。值字段会根据其 SMI 类型进行校验;空的 `OctetString` 或 `OctetStringHex` 是有效值,文本 `OctetString` 的前后空白会被保留。 + +## 支持的 SMI 值类型 + +`value_type` 不区分大小写,并忽略空白、`_` 和 `-` 字符。 +Sink 同时接受文档中的类型名和 SNMP Source 输出的 SNMP4J 语法字符串,包括 `Counter`、`Gauge`、`OCTET STRING` 和 `OBJECT IDENTIFIER`。 + +| 值类型 | 可接受的值 | +|--------|------------| +| `Integer32` 或 `Integer` | 有符号 32 位十进制整数。 | +| `UnsignedInteger32` 或 `UnsignedInteger` | 0 到 4294967295 的十进制整数。 | +| `Counter32` 或 `Counter` | 0 到 4294967295 的十进制整数。 | +| `Gauge32` 或 `Gauge` | 0 到 4294967295 的十进制整数。 | +| `TimeTicks` | 0 到 4294967295 的十进制百分之一秒计数,或 SNMP Source 使用的 SNMP4J 格式 `[days, ]hours:mm:ss.hh`。 | +| `Counter64` | 0 到 18446744073709551615 的十进制整数。 | +| `OctetString` 或 `OCTET STRING` | 输入字符串表示的 UTF-8 文本。 | +| `OctetStringHex` | 偶数个十六进制字符,例如 `00ff10`。 | +| `OID` 或 `OBJECT IDENTIFIER` | 数字对象标识符,可以带前导点。 | +| `IpAddress` | 点分 IPv4 地址。 | + +`OctetString` 用于文本映射。如果需要逐字节保存二进制内容,请使用 `OctetStringHex`。 + +## 示例 + +```hocon +env { + parallelism = 1 + job.mode = "BATCH" +} + +source { + FakeSource { + plugin_output = "snmp_updates" + schema = { + fields { + oid = string + value = string + value_type = string + } + } + rows = [ + { + kind = INSERT + fields = { + oid = "1.3.6.1.2.1.1.5.0" + value = "router-1" + value_type = "OctetString" + } + } + ] + } +} + +sink { + SNMP { + plugin_input = "snmp_updates" + host = "192.0.2.10" + port = 161 + community = "replace-with-your-community" + timeout_millis = 3000 + retries = 1 + } +} +``` + +运行示例前请替换 community 占位值,并在已提交到源码的任务文件之外提供真实凭据。 +记录解析后的任务配置时,`community` 会自动脱敏。无需为日志脱敏将其加入 `shade.options`; +该选项还参与配置的遮蔽和加密流程。 + +## 投递、失败和安全语义 + +- 一次 `write` 成功表示 Agent 已对该行返回成功的 SNMP 响应。 +- 所有配置尝试完成后仍超时,或 SNMP 响应包含非零错误状态时,Sink Task 会失败。 +- 一行在失败前可能阻塞约 `timeout_millis * (retries + 1)`。请确保该时间小于任务的 Checkpoint 超时时间。 +- SNMP4J 会重发超时请求。迟到的响应可能导致非幂等 OID 多次观察到同一次 SET。 +- Sink 没有事务提交协议或可恢复的 Writer 状态。引擎恢复后可能重复发送 SET,因此投递语义为至少一次。 +- 多个并行 Writer 可能乱序更新同一个 OID。如果更新顺序很重要,请使用并行度 1。 +- Sink 不会把 RowKind 解释为 CDC 操作。所有输入行(包括更新或删除类型)都会作为 SET 请求处理。 +- 请把 `community` 视为凭证,通过配置替换或其他密钥管理方式提供,不要把真实值提交到源码中的任务文件。 +- SNMPv2c 不提供传输加密或完整性保护,community 和 SET 负载会以明文发送。请仅在可信私有网络中使用,或通过 VPN 等受保护隧道传输。 +- Trap、Inform、SNMPv1 和 SNMPv3 不属于 V1 范围。 + +`plugin_input` 等配置请参阅[通用 Sink 配置项](../common-options/sink-common-options.md)。 + +<ChangeLog /> diff --git a/plugin-mapping.properties b/plugin-mapping.properties index 6baf75df38..1ac16a0fc5 100644 --- a/plugin-mapping.properties +++ b/plugin-mapping.properties @@ -164,6 +164,7 @@ seatunnel.sink.ActiveMQ = connector-activemq seatunnel.source.MQTT = connector-mqtt seatunnel.sink.MQTT = connector-mqtt seatunnel.source.SNMP = connector-snmp +seatunnel.sink.SNMP = connector-snmp seatunnel.source.Python = connector-python seatunnel.source.Prometheus = connector-prometheus seatunnel.sink.Prometheus = connector-prometheus diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/client/SnmpTargetFactory.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/client/SnmpTargetFactory.java new file mode 100644 index 0000000000..4d53a4249d --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/client/SnmpTargetFactory.java @@ -0,0 +1,51 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.client; + +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpTargetConfig; +import org.apache.seatunnel.connectors.seatunnel.snmp.exception.SnmpConnectorErrorCode; +import org.apache.seatunnel.connectors.seatunnel.snmp.exception.SnmpConnectorException; + +import org.snmp4j.CommunityTarget; +import org.snmp4j.Target; +import org.snmp4j.mp.SnmpConstants; +import org.snmp4j.smi.OctetString; +import org.snmp4j.smi.UdpAddress; + +/** Builds SNMP4J targets without exposing credentials through logs or error messages. */ +public final class SnmpTargetFactory { + + private SnmpTargetFactory() {} + + public static Target create(SnmpTargetConfig config) { + CommunityTarget target = new CommunityTarget(); + try { + target.setAddress(new UdpAddress(config.getHost() + "/" + config.getPort())); + } catch (IllegalArgumentException e) { + throw new SnmpConnectorException( + SnmpConnectorErrorCode.INVALID_CONFIG, + "Invalid SNMP agent address " + config.getHost() + ":" + config.getPort(), + e); + } + target.setCommunity(new OctetString(config.getCommunity())); + target.setVersion(SnmpConstants.version2c); + target.setTimeout(config.getTimeoutMillis()); + target.setRetries(config.getRetries()); + return target; + } +} diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceOptions.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpOptions.java similarity index 78% copy from seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceOptions.java copy to seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpOptions.java index c309801f32..dc17c22c4f 100644 --- a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceOptions.java +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpOptions.java @@ -20,9 +20,8 @@ package org.apache.seatunnel.connectors.seatunnel.snmp.config; import org.apache.seatunnel.api.configuration.Option; import org.apache.seatunnel.api.configuration.Options; -import java.util.List; - -public final class SnmpSourceOptions { +/** Options shared by SNMP source and sink connections. */ +public final class SnmpOptions { public static final String CONNECTOR_IDENTITY = "SNMP"; @@ -41,12 +40,6 @@ public final class SnmpSourceOptions { .noDefaultValue() .withDescription("SNMPv2c community credential"); - public static final Option<List<String>> OIDS = - Options.key("oids") - .listType() - .noDefaultValue() - .withDescription("Numeric OIDs to retrieve with SNMP GET"); - public static final Option<Long> TIMEOUT_MILLIS = Options.key("timeout_millis") .longType() @@ -59,11 +52,5 @@ public final class SnmpSourceOptions { .defaultValue(1) .withDescription("Number of retries after the initial SNMP request attempt"); - public static final Option<Long> POLL_INTERVAL_MILLIS = - Options.key("poll_interval_millis") - .longType() - .defaultValue(60000L) - .withDescription("Interval in milliseconds between streaming polls"); - - private SnmpSourceOptions() {} + private SnmpOptions() {} } diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSinkConfig.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSinkConfig.java new file mode 100644 index 0000000000..3efa98f417 --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSinkConfig.java @@ -0,0 +1,131 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.config; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; + +import java.io.Serializable; +import java.util.HashSet; +import java.util.Set; + +/** Validated runtime configuration for the SNMPv2c SET sink. */ +public final class SnmpSinkConfig implements Serializable, SnmpTargetConfig { + + private static final long serialVersionUID = 1L; + + private final String host; + private final int port; + private final String community; + private final long timeoutMillis; + private final int retries; + private final String oidField; + private final String valueField; + private final String valueTypeField; + + public SnmpSinkConfig(ReadonlyConfig config) { + String configuredHost = config.get(SnmpSinkOptions.HOST); + if (isBlank(configuredHost)) { + throw new IllegalArgumentException("SNMP sink host must not be blank"); + } + this.host = configuredHost.trim(); + this.port = config.get(SnmpSinkOptions.PORT); + this.community = config.get(SnmpSinkOptions.COMMUNITY); + this.timeoutMillis = config.get(SnmpSinkOptions.TIMEOUT_MILLIS); + this.retries = config.get(SnmpSinkOptions.RETRIES); + this.oidField = requireField(config.get(SnmpSinkOptions.OID_FIELD), "oid_field"); + this.valueField = requireField(config.get(SnmpSinkOptions.VALUE_FIELD), "value_field"); + this.valueTypeField = + requireField(config.get(SnmpSinkOptions.VALUE_TYPE_FIELD), "value_type_field"); + validateTarget(); + validateDistinctFields(); + } + + private void validateTarget() { + if (port < 1 || port > 65535) { + throw new IllegalArgumentException("SNMP sink port must be between 1 and 65535"); + } + if (isBlank(community)) { + throw new IllegalArgumentException("SNMP sink community must not be blank"); + } + if (timeoutMillis <= 0) { + throw new IllegalArgumentException("SNMP sink timeout_millis must be greater than 0"); + } + if (retries < 0) { + throw new IllegalArgumentException("SNMP sink retries must not be negative"); + } + } + + private static String requireField(String configuredField, String optionName) { + if (isBlank(configuredField)) { + throw new IllegalArgumentException("SNMP sink " + optionName + " must not be blank"); + } + return configuredField.trim(); + } + + private static boolean isBlank(String value) { + return value == null || value.trim().isEmpty(); + } + + private void validateDistinctFields() { + Set<String> fields = new HashSet<>(); + fields.add(oidField); + fields.add(valueField); + fields.add(valueTypeField); + if (fields.size() != 3) { + throw new IllegalArgumentException( + "SNMP sink oid_field, value_field, and value_type_field must be distinct"); + } + } + + @Override + public String getHost() { + return host; + } + + @Override + public int getPort() { + return port; + } + + @Override + public String getCommunity() { + return community; + } + + @Override + public long getTimeoutMillis() { + return timeoutMillis; + } + + @Override + public int getRetries() { + return retries; + } + + public String getOidField() { + return oidField; + } + + public String getValueField() { + return valueField; + } + + public String getValueTypeField() { + return valueTypeField; + } +} diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSinkOptions.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSinkOptions.java new file mode 100644 index 0000000000..7034675e26 --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSinkOptions.java @@ -0,0 +1,53 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.config; + +import org.apache.seatunnel.api.configuration.Option; +import org.apache.seatunnel.api.configuration.Options; + +/** Configuration options for the SNMPv2c SET sink. */ +public final class SnmpSinkOptions { + + public static final String CONNECTOR_IDENTITY = SnmpOptions.CONNECTOR_IDENTITY; + + public static final Option<String> HOST = SnmpOptions.HOST; + public static final Option<Integer> PORT = SnmpOptions.PORT; + public static final Option<String> COMMUNITY = SnmpOptions.COMMUNITY; + public static final Option<Long> TIMEOUT_MILLIS = SnmpOptions.TIMEOUT_MILLIS; + public static final Option<Integer> RETRIES = SnmpOptions.RETRIES; + + public static final Option<String> OID_FIELD = + Options.key("oid_field") + .stringType() + .defaultValue("oid") + .withDescription("Input STRING field containing the numeric OID to set"); + + public static final Option<String> VALUE_FIELD = + Options.key("value_field") + .stringType() + .defaultValue("value") + .withDescription("Input STRING field containing the value to set"); + + public static final Option<String> VALUE_TYPE_FIELD = + Options.key("value_type_field") + .stringType() + .defaultValue("value_type") + .withDescription("Input STRING field containing the SMI value type"); + + private SnmpSinkOptions() {} +} diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceConfig.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceConfig.java index 09065192e5..affb3da361 100644 --- a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceConfig.java +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceConfig.java @@ -30,7 +30,7 @@ import java.util.Set; import java.util.regex.Pattern; /** Validated runtime configuration for the SNMP source connector. */ -public final class SnmpSourceConfig implements Serializable { +public final class SnmpSourceConfig implements Serializable, SnmpTargetConfig { private static final long serialVersionUID = 1L; @@ -79,6 +79,25 @@ public final class SnmpSourceConfig implements Serializable { } } + private static boolean isBlank(String value) { + return value == null || value.trim().isEmpty(); + } + + @Override + public String getHost() { + return host; + } + + @Override + public int getPort() { + return port; + } + + @Override + public String getCommunity() { + return community; + } + private static List<OID> parseOids(List<String> configuredOids) { if (configuredOids == null || configuredOids.isEmpty()) { throw new IllegalArgumentException("SNMP source oids must not be empty"); @@ -110,30 +129,16 @@ public final class SnmpSourceConfig implements Serializable { return Collections.unmodifiableList(parsed); } - private static boolean isBlank(String value) { - return value == null || value.trim().isEmpty(); - } - - public String getHost() { - return host; - } - - public int getPort() { - return port; - } - - public String getCommunity() { - return community; - } - public List<OID> getOids() { return oids; } + @Override public long getTimeoutMillis() { return timeoutMillis; } + @Override public int getRetries() { return retries; } diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceOptions.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceOptions.java index c309801f32..501e5fc6b7 100644 --- a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceOptions.java +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceOptions.java @@ -24,22 +24,13 @@ import java.util.List; public final class SnmpSourceOptions { - public static final String CONNECTOR_IDENTITY = "SNMP"; + public static final String CONNECTOR_IDENTITY = SnmpOptions.CONNECTOR_IDENTITY; - public static final Option<String> HOST = - Options.key("host") - .stringType() - .noDefaultValue() - .withDescription("SNMP agent host name or IP address"); + public static final Option<String> HOST = SnmpOptions.HOST; - public static final Option<Integer> PORT = - Options.key("port").intType().defaultValue(161).withDescription("SNMP agent UDP port"); + public static final Option<Integer> PORT = SnmpOptions.PORT; - public static final Option<String> COMMUNITY = - Options.key("community") - .stringType() - .noDefaultValue() - .withDescription("SNMPv2c community credential"); + public static final Option<String> COMMUNITY = SnmpOptions.COMMUNITY; public static final Option<List<String>> OIDS = Options.key("oids") @@ -47,17 +38,9 @@ public final class SnmpSourceOptions { .noDefaultValue() .withDescription("Numeric OIDs to retrieve with SNMP GET"); - public static final Option<Long> TIMEOUT_MILLIS = - Options.key("timeout_millis") - .longType() - .defaultValue(5000L) - .withDescription("Timeout in milliseconds for each SNMP request attempt"); + public static final Option<Long> TIMEOUT_MILLIS = SnmpOptions.TIMEOUT_MILLIS; - public static final Option<Integer> RETRIES = - Options.key("retries") - .intType() - .defaultValue(1) - .withDescription("Number of retries after the initial SNMP request attempt"); + public static final Option<Integer> RETRIES = SnmpOptions.RETRIES; public static final Option<Long> POLL_INTERVAL_MILLIS = Options.key("poll_interval_millis") diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpTargetConfig.java similarity index 54% copy from seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java copy to seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpTargetConfig.java index bb87272872..e785652053 100644 --- a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpTargetConfig.java @@ -15,29 +15,20 @@ * limitations under the License. */ -package org.apache.seatunnel.connectors.seatunnel.snmp.exception; +package org.apache.seatunnel.connectors.seatunnel.snmp.config; -import org.apache.seatunnel.common.exception.SeaTunnelErrorCode; +import java.io.Serializable; -public enum SnmpConnectorErrorCode implements SeaTunnelErrorCode { - CONNECTION_FAILED("SNMP-01", "SNMP client initialization failed"), - POLL_FAILED("SNMP-02", "SNMP poll failed"); +/** Serializable target settings shared by SNMP source and sink clients. */ +public interface SnmpTargetConfig extends Serializable { - private final String code; - private final String description; + String getHost(); - SnmpConnectorErrorCode(String code, String description) { - this.code = code; - this.description = description; - } + int getPort(); - @Override - public String getCode() { - return code; - } + String getCommunity(); - @Override - public String getDescription() { - return description; - } + long getTimeoutMillis(); + + int getRetries(); } diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java index bb87272872..d660e34246 100644 --- a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java @@ -21,7 +21,11 @@ import org.apache.seatunnel.common.exception.SeaTunnelErrorCode; public enum SnmpConnectorErrorCode implements SeaTunnelErrorCode { CONNECTION_FAILED("SNMP-01", "SNMP client initialization failed"), - POLL_FAILED("SNMP-02", "SNMP poll failed"); + POLL_FAILED("SNMP-02", "SNMP poll failed"), + WRITE_FAILED("SNMP-03", "SNMP SET request failed"), + INVALID_CONFIG("SNMP-04", "Invalid SNMP connector configuration"), + INVALID_ROW("SNMP-05", "Invalid SNMP sink row"), + CLOSE_FAILED("SNMP-06", "SNMP client close failed"); private final String code; private final String description; diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/source/Snmp4jClient.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/Snmp4jSetClient.java similarity index 53% copy from seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/source/Snmp4jClient.java copy to seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/Snmp4jSetClient.java index f4095e9fda..f297194425 100644 --- a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/source/Snmp4jClient.java +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/Snmp4jSetClient.java @@ -15,44 +15,38 @@ * limitations under the License. */ -package org.apache.seatunnel.connectors.seatunnel.snmp.source; +package org.apache.seatunnel.connectors.seatunnel.snmp.sink; -import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSourceConfig; +import org.apache.seatunnel.connectors.seatunnel.snmp.client.SnmpTargetFactory; +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSinkConfig; -import org.snmp4j.CommunityTarget; import org.snmp4j.PDU; import org.snmp4j.Snmp; import org.snmp4j.Target; import org.snmp4j.event.ResponseEvent; -import org.snmp4j.mp.SnmpConstants; -import org.snmp4j.smi.OID; -import org.snmp4j.smi.OctetString; -import org.snmp4j.smi.UdpAddress; -import org.snmp4j.smi.Variable; import org.snmp4j.smi.VariableBinding; import org.snmp4j.transport.DefaultUdpTransportMapping; import java.io.IOException; -import java.util.ArrayList; -import java.util.List; -/** SNMPv2c client backed by SNMP4J. */ -final class Snmp4jClient implements SnmpClient { +/** SNMPv2c SET client backed by SNMP4J. */ +final class Snmp4jSetClient implements SnmpSetClient { - private final SnmpSourceConfig config; + private final SnmpSinkConfig config; private final Snmp snmp; private final Target target; - Snmp4jClient(SnmpSourceConfig config) throws IOException { + Snmp4jSetClient(SnmpSinkConfig config) throws IOException { this(config, new Snmp(new DefaultUdpTransportMapping())); } - Snmp4jClient(SnmpSourceConfig config, Snmp snmp) throws IOException { + Snmp4jSetClient(SnmpSinkConfig config, Snmp snmp) throws IOException { this.config = config; - this.target = buildTarget(config); + Target createdTarget; try { + createdTarget = SnmpTargetFactory.create(config); snmp.listen(); - } catch (IOException e) { + } catch (IOException | RuntimeException e) { try { snmp.close(); } catch (IOException closeException) { @@ -60,15 +54,16 @@ final class Snmp4jClient implements SnmpClient { } throw e; } + this.target = createdTarget; this.snmp = snmp; } @Override - public List<SnmpRecord> get(List<OID> oids) throws IOException { - ResponseEvent event = snmp.send(buildGetRequest(oids), target); + public void set(SnmpSetRequest request) throws IOException { + ResponseEvent event = snmp.send(buildSetRequest(request), target); if (event == null || event.getResponse() == null) { throw new IOException( - "SNMP request timed out for agent " + "SNMP SET request timed out for agent " + config.getHost() + ":" + config.getPort()); @@ -84,7 +79,6 @@ final class Snmp4jClient implements SnmpClient { + ") at index " + response.getErrorIndex()); } - return extractRecords(response); } @Override @@ -92,35 +86,10 @@ final class Snmp4jClient implements SnmpClient { snmp.close(); } - static PDU buildGetRequest(List<OID> oids) { + static PDU buildSetRequest(SnmpSetRequest request) { PDU pdu = new PDU(); - pdu.setType(PDU.GET); - for (OID oid : oids) { - pdu.add(new VariableBinding(oid)); - } + pdu.setType(PDU.SET); + pdu.add(new VariableBinding(request.getOid(), request.getValue())); return pdu; } - - static Target buildTarget(SnmpSourceConfig config) { - CommunityTarget target = new CommunityTarget(); - target.setAddress(new UdpAddress(config.getHost() + "/" + config.getPort())); - target.setCommunity(new OctetString(config.getCommunity())); - target.setVersion(SnmpConstants.version2c); - target.setTimeout(config.getTimeoutMillis()); - target.setRetries(config.getRetries()); - return target; - } - - static List<SnmpRecord> extractRecords(PDU response) { - List<SnmpRecord> records = new ArrayList<>(response.size()); - for (VariableBinding binding : response.getVariableBindings()) { - Variable variable = binding.getVariable(); - records.add( - new SnmpRecord( - binding.getOid().toString(), - variable.toString(), - variable.getSyntaxString())); - } - return records; - } } diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSetClient.java similarity index 54% copy from seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java copy to seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSetClient.java index bb87272872..5d0e40c4e4 100644 --- a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSetClient.java @@ -15,29 +15,12 @@ * limitations under the License. */ -package org.apache.seatunnel.connectors.seatunnel.snmp.exception; +package org.apache.seatunnel.connectors.seatunnel.snmp.sink; -import org.apache.seatunnel.common.exception.SeaTunnelErrorCode; +import java.io.Closeable; +import java.io.IOException; -public enum SnmpConnectorErrorCode implements SeaTunnelErrorCode { - CONNECTION_FAILED("SNMP-01", "SNMP client initialization failed"), - POLL_FAILED("SNMP-02", "SNMP poll failed"); +interface SnmpSetClient extends Closeable { - private final String code; - private final String description; - - SnmpConnectorErrorCode(String code, String description) { - this.code = code; - this.description = description; - } - - @Override - public String getCode() { - return code; - } - - @Override - public String getDescription() { - return description; - } + void set(SnmpSetRequest request) throws IOException; } diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSetRequest.java similarity index 55% copy from seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java copy to seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSetRequest.java index bb87272872..ce8070b365 100644 --- a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/exception/SnmpConnectorErrorCode.java +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSetRequest.java @@ -15,29 +15,26 @@ * limitations under the License. */ -package org.apache.seatunnel.connectors.seatunnel.snmp.exception; +package org.apache.seatunnel.connectors.seatunnel.snmp.sink; -import org.apache.seatunnel.common.exception.SeaTunnelErrorCode; +import org.snmp4j.smi.OID; +import org.snmp4j.smi.Variable; -public enum SnmpConnectorErrorCode implements SeaTunnelErrorCode { - CONNECTION_FAILED("SNMP-01", "SNMP client initialization failed"), - POLL_FAILED("SNMP-02", "SNMP poll failed"); +final class SnmpSetRequest { - private final String code; - private final String description; + private final OID oid; + private final Variable value; - SnmpConnectorErrorCode(String code, String description) { - this.code = code; - this.description = description; + SnmpSetRequest(OID oid, Variable value) { + this.oid = oid; + this.value = value; } - @Override - public String getCode() { - return code; + OID getOid() { + return oid; } - @Override - public String getDescription() { - return description; + Variable getValue() { + return value; } } diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSink.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSink.java new file mode 100644 index 0000000000..65f4f5a474 --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSink.java @@ -0,0 +1,59 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.sink; + +import org.apache.seatunnel.api.sink.SinkWriter; +import org.apache.seatunnel.api.table.catalog.CatalogTable; +import org.apache.seatunnel.api.table.type.SeaTunnelRow; +import org.apache.seatunnel.api.table.type.SeaTunnelRowType; +import org.apache.seatunnel.connectors.seatunnel.common.sink.AbstractSimpleSink; +import org.apache.seatunnel.connectors.seatunnel.common.sink.AbstractSinkWriter; +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSinkConfig; +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSinkOptions; + +import java.util.Optional; + +/** SeaTunnel sink that applies input rows through SNMPv2c SET requests. */ +public final class SnmpSink extends AbstractSimpleSink<SeaTunnelRow, Void> { + + private final SnmpSinkConfig config; + private final CatalogTable catalogTable; + private final SeaTunnelRowType rowType; + + public SnmpSink(SnmpSinkConfig config, CatalogTable catalogTable) { + this.config = config; + this.catalogTable = catalogTable; + this.rowType = catalogTable.getSeaTunnelRowType(); + new SnmpSinkRowConverter(config, rowType); + } + + @Override + public String getPluginName() { + return SnmpSinkOptions.CONNECTOR_IDENTITY; + } + + @Override + public AbstractSinkWriter<SeaTunnelRow, Void> createWriter(SinkWriter.Context context) { + return new SnmpSinkWriter(config, rowType); + } + + @Override + public Optional<CatalogTable> getWriteCatalogTable() { + return Optional.of(catalogTable); + } +} diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkFactory.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkFactory.java new file mode 100644 index 0000000000..f90fc44db5 --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkFactory.java @@ -0,0 +1,70 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.sink; + +import org.apache.seatunnel.api.configuration.util.Conditions; +import org.apache.seatunnel.api.configuration.util.OptionRule; +import org.apache.seatunnel.api.table.connector.TableSink; +import org.apache.seatunnel.api.table.factory.Factory; +import org.apache.seatunnel.api.table.factory.TableSinkFactory; +import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext; +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSinkConfig; +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSinkOptions; + +import com.google.auto.service.AutoService; + +/** Creates SNMPv2c SET sinks from table factory configuration. */ +@AutoService(Factory.class) +public final class SnmpSinkFactory implements TableSinkFactory { + + @Override + public String factoryIdentifier() { + return SnmpSinkOptions.CONNECTOR_IDENTITY; + } + + @Override + public OptionRule optionRule() { + return OptionRule.builder() + .required(SnmpSinkOptions.HOST, Conditions.notBlank(SnmpSinkOptions.HOST)) + .required(SnmpSinkOptions.COMMUNITY, Conditions.notBlank(SnmpSinkOptions.COMMUNITY)) + .optional( + SnmpSinkOptions.PORT, + Conditions.greaterOrEqual(SnmpSinkOptions.PORT, 1) + .and(Conditions.lessOrEqual(SnmpSinkOptions.PORT, 65535))) + .optional( + SnmpSinkOptions.TIMEOUT_MILLIS, + Conditions.greaterThan(SnmpSinkOptions.TIMEOUT_MILLIS, 0L)) + .optional( + SnmpSinkOptions.RETRIES, + Conditions.greaterOrEqual(SnmpSinkOptions.RETRIES, 0)) + .optional(SnmpSinkOptions.OID_FIELD, Conditions.notBlank(SnmpSinkOptions.OID_FIELD)) + .optional( + SnmpSinkOptions.VALUE_FIELD, + Conditions.notBlank(SnmpSinkOptions.VALUE_FIELD)) + .optional( + SnmpSinkOptions.VALUE_TYPE_FIELD, + Conditions.notBlank(SnmpSinkOptions.VALUE_TYPE_FIELD)) + .build(); + } + + @Override + public TableSink createSink(TableSinkFactoryContext context) { + SnmpSinkConfig config = new SnmpSinkConfig(context.getOptions()); + return () -> new SnmpSink(config, context.getCatalogTable()); + } +} diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkRowConverter.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkRowConverter.java new file mode 100644 index 0000000000..3e424ff072 --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkRowConverter.java @@ -0,0 +1,316 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.sink; + +import org.apache.seatunnel.api.table.type.SeaTunnelDataType; +import org.apache.seatunnel.api.table.type.SeaTunnelRow; +import org.apache.seatunnel.api.table.type.SeaTunnelRowType; +import org.apache.seatunnel.api.table.type.SqlType; +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSinkConfig; +import org.apache.seatunnel.connectors.seatunnel.snmp.exception.SnmpConnectorErrorCode; +import org.apache.seatunnel.connectors.seatunnel.snmp.exception.SnmpConnectorException; + +import org.snmp4j.smi.Counter32; +import org.snmp4j.smi.Counter64; +import org.snmp4j.smi.Gauge32; +import org.snmp4j.smi.Integer32; +import org.snmp4j.smi.IpAddress; +import org.snmp4j.smi.OID; +import org.snmp4j.smi.OctetString; +import org.snmp4j.smi.TimeTicks; +import org.snmp4j.smi.UnsignedInteger32; +import org.snmp4j.smi.Variable; + +import java.io.Serializable; +import java.math.BigInteger; +import java.nio.charset.StandardCharsets; +import java.util.Arrays; +import java.util.Locale; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +/** Converts the configured input fields into one validated SNMP SET binding. */ +final class SnmpSinkRowConverter implements Serializable { + + private static final long serialVersionUID = 1L; + private static final long MAX_UNSIGNED_32 = 0xFFFFFFFFL; + private static final BigInteger MAX_UNSIGNED_64 = new BigInteger("18446744073709551615"); + private static final Pattern NUMERIC_OID = Pattern.compile("^\\.?[0-9]+(\\.[0-9]+)+$"); + private static final Pattern HEX_VALUE = Pattern.compile("^[0-9a-fA-F]*$"); + private static final Pattern IPV4_VALUE = Pattern.compile("^(?:[0-9]{1,3}\\.){3}[0-9]{1,3}$"); + private static final Pattern DECIMAL_INTEGER = Pattern.compile("^[+-]?[0-9]+$"); + private static final Pattern FORMATTED_TIME_TICKS = + Pattern.compile( + "^(?:([0-9]+) (day|days), )?([0-9]{1,2}):([0-9]{2}):([0-9]{2})\\.([0-9]{2})$"); + private static final BigInteger TICKS_PER_DAY = BigInteger.valueOf(8_640_000L); + private static final BigInteger TICKS_PER_HOUR = BigInteger.valueOf(360_000L); + private static final BigInteger TICKS_PER_MINUTE = BigInteger.valueOf(6_000L); + private static final BigInteger TICKS_PER_SECOND = BigInteger.valueOf(100L); + + private final int rowArity; + private final int oidIndex; + private final int valueIndex; + private final int valueTypeIndex; + + SnmpSinkRowConverter(SnmpSinkConfig config, SeaTunnelRowType rowType) { + this.rowArity = rowType.getTotalFields(); + this.oidIndex = requireStringField(rowType, config.getOidField(), "oid_field"); + this.valueIndex = requireStringField(rowType, config.getValueField(), "value_field"); + this.valueTypeIndex = + requireStringField(rowType, config.getValueTypeField(), "value_type_field"); + } + + SnmpSetRequest convert(SeaTunnelRow row) { + if (row.getArity() != rowArity) { + throw invalidRow( + "Input row arity " + + row.getArity() + + " does not match the configured schema arity " + + rowArity); + } + + String oid = requireNonBlankRowValue(row, oidIndex, "OID"); + String value = requireNonNullRowValue(row, valueIndex, "value"); + String valueType = requireNonBlankRowValue(row, valueTypeIndex, "value type"); + return new SnmpSetRequest(parseOid(oid), parseVariable(valueType, value)); + } + + private static int requireStringField( + SeaTunnelRowType rowType, String fieldName, String optionName) { + int index = rowType.indexOf(fieldName, false); + if (index < 0) { + throw invalidConfig( + "Option `" + + optionName + + "` references unknown field `" + + fieldName + + "`. Available fields are " + + Arrays.toString(rowType.getFieldNames())); + } + SeaTunnelDataType<?> dataType = rowType.getFieldType(index); + if (dataType.getSqlType() != SqlType.STRING) { + throw invalidConfig( + "Field `" + + fieldName + + "` configured by `" + + optionName + + "` must use STRING type, but was " + + dataType.getSqlType()); + } + return index; + } + + private static String requireNonBlankRowValue( + SeaTunnelRow row, int index, String fieldDescription) { + String stringValue = requireNonNullRowValue(row, index, fieldDescription); + if (stringValue.trim().isEmpty()) { + throw invalidRow("SNMP sink " + fieldDescription + " field must not be blank"); + } + return stringValue.trim(); + } + + private static String requireNonNullRowValue( + SeaTunnelRow row, int index, String fieldDescription) { + Object value = row.getField(index); + if (value == null) { + throw invalidRow("SNMP sink " + fieldDescription + " field must not be null"); + } + if (!(value instanceof String)) { + throw invalidRow("SNMP sink " + fieldDescription + " field must contain a STRING"); + } + return (String) value; + } + + static OID parseOid(String configuredOid) { + String value = configuredOid.trim(); + if (!NUMERIC_OID.matcher(value).matches()) { + throw invalidRow("SNMP sink OID must be numeric: " + configuredOid); + } + if (value.charAt(0) == '.') { + value = value.substring(1); + } + try { + OID oid = new OID(value); + if (!oid.isValid()) { + throw invalidRow("SNMP sink OID is invalid: " + configuredOid); + } + return oid; + } catch (SnmpConnectorException e) { + throw e; + } catch (RuntimeException e) { + throw invalidRow("SNMP sink OID is invalid: " + configuredOid, e); + } + } + + static Variable parseVariable(String configuredType, String value) { + String normalizedType = normalizeType(configuredType); + try { + switch (normalizedType) { + case "INTEGER": + case "INTEGER32": + return new Integer32(Integer.parseInt(value)); + case "UNSIGNEDINTEGER": + case "UNSIGNEDINTEGER32": + return new UnsignedInteger32(parseUnsigned32(value, configuredType)); + case "COUNTER": + case "COUNTER32": + return new Counter32(parseUnsigned32(value, configuredType)); + case "GAUGE": + case "GAUGE32": + return new Gauge32(parseUnsigned32(value, configuredType)); + case "TIMETICKS": + return new TimeTicks(parseTimeTicks(value, configuredType)); + case "COUNTER64": + return new Counter64(parseUnsigned64(value, configuredType)); + case "OCTETSTRING": + return new OctetString(value.getBytes(StandardCharsets.UTF_8)); + case "OCTETSTRINGHEX": + return new OctetString(parseHex(value)); + case "OBJECTIDENTIFIER": + case "OID": + return parseOid(value); + case "IPADDRESS": + validateIpv4(value); + return new IpAddress(value); + default: + throw invalidRow( + "Unsupported SNMP sink value type `" + + configuredType + + "`. Supported types are Integer32, UnsignedInteger32, " + + "Counter/Counter32, Gauge/Gauge32, TimeTicks, Counter64, " + + "OctetString/OCTET STRING, OctetStringHex, " + + "OID/OBJECT IDENTIFIER, and IpAddress"); + } + } catch (SnmpConnectorException e) { + throw e; + } catch (RuntimeException e) { + throw invalidRow("SNMP sink value is invalid for type `" + configuredType + "`", e); + } + } + + private static String normalizeType(String configuredType) { + StringBuilder normalized = new StringBuilder(configuredType.length()); + for (int index = 0; index < configuredType.length(); index++) { + char character = configuredType.charAt(index); + if (character != '_' && character != '-' && !Character.isWhitespace(character)) { + normalized.append(character); + } + } + return normalized.toString().toUpperCase(Locale.ROOT); + } + + private static long parseUnsigned32(String value, String configuredType) { + long parsed = Long.parseLong(value); + if (parsed < 0 || parsed > MAX_UNSIGNED_32) { + throw invalidRow( + "SNMP sink value for type `" + + configuredType + + "` must be between 0 and " + + MAX_UNSIGNED_32); + } + return parsed; + } + + private static long parseTimeTicks(String value, String configuredType) { + if (DECIMAL_INTEGER.matcher(value).matches()) { + return parseUnsigned32(value, configuredType); + } + + Matcher matcher = FORMATTED_TIME_TICKS.matcher(value); + if (!matcher.matches()) { + throw invalidRow( + "SNMP sink TimeTicks value must be an unsigned decimal count or use the " + + "SNMP4J format `[days, ]hours:mm:ss.hh`"); + } + + BigInteger days = + matcher.group(1) == null ? BigInteger.ZERO : new BigInteger(matcher.group(1)); + String dayUnit = matcher.group(2); + if ((matcher.group(1) != null && BigInteger.ZERO.equals(days)) + || (BigInteger.ONE.equals(days) && !"day".equals(dayUnit)) + || (!BigInteger.ONE.equals(days) && "day".equals(dayUnit))) { + throw invalidRow("SNMP sink TimeTicks day unit does not match its value"); + } + + int hours = Integer.parseInt(matcher.group(3)); + int minutes = Integer.parseInt(matcher.group(4)); + int seconds = Integer.parseInt(matcher.group(5)); + int hundredths = Integer.parseInt(matcher.group(6)); + if (hours > 23 || minutes > 59 || seconds > 59) { + throw invalidRow("SNMP sink TimeTicks formatted value is outside clock bounds"); + } + + BigInteger ticks = + days.multiply(TICKS_PER_DAY) + .add(BigInteger.valueOf(hours).multiply(TICKS_PER_HOUR)) + .add(BigInteger.valueOf(minutes).multiply(TICKS_PER_MINUTE)) + .add(BigInteger.valueOf(seconds).multiply(TICKS_PER_SECOND)) + .add(BigInteger.valueOf(hundredths)); + if (ticks.compareTo(BigInteger.valueOf(MAX_UNSIGNED_32)) > 0) { + throw invalidRow("SNMP sink TimeTicks value must be between 0 and " + MAX_UNSIGNED_32); + } + return ticks.longValue(); + } + + private static long parseUnsigned64(String value, String configuredType) { + BigInteger parsed = new BigInteger(value); + if (parsed.signum() < 0 || parsed.compareTo(MAX_UNSIGNED_64) > 0) { + throw invalidRow( + "SNMP sink value for type `" + + configuredType + + "` must be between 0 and " + + MAX_UNSIGNED_64); + } + return parsed.longValue(); + } + + private static byte[] parseHex(String value) { + if ((value.length() & 1) != 0 || !HEX_VALUE.matcher(value).matches()) { + throw invalidRow( + "SNMP sink OctetStringHex value must contain an even number of hexadecimal characters"); + } + byte[] bytes = new byte[value.length() / 2]; + for (int index = 0; index < value.length(); index += 2) { + bytes[index / 2] = (byte) Integer.parseInt(value.substring(index, index + 2), 16); + } + return bytes; + } + + private static void validateIpv4(String value) { + if (!IPV4_VALUE.matcher(value).matches()) { + throw invalidRow("SNMP sink IpAddress value must be a dotted IPv4 address"); + } + for (String octet : value.split("\\.")) { + if (Integer.parseInt(octet) > 255) { + throw invalidRow("SNMP sink IpAddress value must be a dotted IPv4 address"); + } + } + } + + private static SnmpConnectorException invalidConfig(String message) { + return new SnmpConnectorException(SnmpConnectorErrorCode.INVALID_CONFIG, message); + } + + private static SnmpConnectorException invalidRow(String message) { + return new SnmpConnectorException(SnmpConnectorErrorCode.INVALID_ROW, message); + } + + private static SnmpConnectorException invalidRow(String message, Throwable cause) { + return new SnmpConnectorException(SnmpConnectorErrorCode.INVALID_ROW, message, cause); + } +} diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkWriter.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkWriter.java new file mode 100644 index 0000000000..a5d72ccdc0 --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkWriter.java @@ -0,0 +1,93 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.sink; + +import org.apache.seatunnel.api.table.type.SeaTunnelRow; +import org.apache.seatunnel.api.table.type.SeaTunnelRowType; +import org.apache.seatunnel.connectors.seatunnel.common.sink.AbstractSinkWriter; +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSinkConfig; +import org.apache.seatunnel.connectors.seatunnel.snmp.exception.SnmpConnectorErrorCode; +import org.apache.seatunnel.connectors.seatunnel.snmp.exception.SnmpConnectorException; + +import java.io.IOException; + +/** Writes each row as one synchronous SNMPv2c SET request. */ +public final class SnmpSinkWriter extends AbstractSinkWriter<SeaTunnelRow, Void> { + + private final SnmpSinkConfig config; + private final SnmpSinkRowConverter converter; + private final SnmpSetClient client; + + public SnmpSinkWriter(SnmpSinkConfig config, SeaTunnelRowType rowType) { + this(config, rowType, Snmp4jSetClient::new); + } + + SnmpSinkWriter( + SnmpSinkConfig config, SeaTunnelRowType rowType, SnmpSetClientFactory clientFactory) { + this.config = config; + this.converter = new SnmpSinkRowConverter(config, rowType); + try { + this.client = clientFactory.create(config); + } catch (IOException e) { + throw new SnmpConnectorException( + SnmpConnectorErrorCode.CONNECTION_FAILED, + "Failed to initialize SNMP SET client for agent " + + config.getHost() + + ":" + + config.getPort(), + e); + } + } + + @Override + public void write(SeaTunnelRow row) { + SnmpSetRequest request = converter.convert(row); + try { + client.set(request); + } catch (IOException e) { + throw new SnmpConnectorException( + SnmpConnectorErrorCode.WRITE_FAILED, + "Failed to set OID " + + request.getOid() + + " on SNMP agent " + + config.getHost() + + ":" + + config.getPort(), + e); + } + } + + @Override + public void close() { + try { + client.close(); + } catch (IOException e) { + throw new SnmpConnectorException( + SnmpConnectorErrorCode.CLOSE_FAILED, + "Failed to close SNMP SET client for agent " + + config.getHost() + + ":" + + config.getPort(), + e); + } + } + + interface SnmpSetClientFactory { + SnmpSetClient create(SnmpSinkConfig config) throws IOException; + } +} diff --git a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/source/Snmp4jClient.java b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/source/Snmp4jClient.java index f4095e9fda..18430c46d7 100644 --- a/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/source/Snmp4jClient.java +++ b/seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/source/Snmp4jClient.java @@ -17,17 +17,14 @@ package org.apache.seatunnel.connectors.seatunnel.snmp.source; +import org.apache.seatunnel.connectors.seatunnel.snmp.client.SnmpTargetFactory; import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSourceConfig; -import org.snmp4j.CommunityTarget; import org.snmp4j.PDU; import org.snmp4j.Snmp; import org.snmp4j.Target; import org.snmp4j.event.ResponseEvent; -import org.snmp4j.mp.SnmpConstants; import org.snmp4j.smi.OID; -import org.snmp4j.smi.OctetString; -import org.snmp4j.smi.UdpAddress; import org.snmp4j.smi.Variable; import org.snmp4j.smi.VariableBinding; import org.snmp4j.transport.DefaultUdpTransportMapping; @@ -102,13 +99,7 @@ final class Snmp4jClient implements SnmpClient { } static Target buildTarget(SnmpSourceConfig config) { - CommunityTarget target = new CommunityTarget(); - target.setAddress(new UdpAddress(config.getHost() + "/" + config.getPort())); - target.setCommunity(new OctetString(config.getCommunity())); - target.setVersion(SnmpConstants.version2c); - target.setTimeout(config.getTimeoutMillis()); - target.setRetries(config.getRetries()); - return target; + return SnmpTargetFactory.create(config); } static List<SnmpRecord> extractRecords(PDU response) { diff --git a/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/client/SnmpTargetFactoryTest.java b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/client/SnmpTargetFactoryTest.java new file mode 100644 index 0000000000..c360bc1acc --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/client/SnmpTargetFactoryTest.java @@ -0,0 +1,47 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.client; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSinkConfig; +import org.apache.seatunnel.connectors.seatunnel.snmp.exception.SnmpConnectorException; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.Map; + +class SnmpTargetFactoryTest { + + @Test + void testInvalidAddressUsesConnectorExceptionWithoutCommunity() { + Map<String, Object> values = new HashMap<>(); + values.put("host", "invalid/host"); + values.put("community", "private-community"); + SnmpSinkConfig config = new SnmpSinkConfig(ReadonlyConfig.fromMap(values)); + + SnmpConnectorException exception = + Assertions.assertThrows( + SnmpConnectorException.class, () -> SnmpTargetFactory.create(config)); + + Assertions.assertTrue(exception.getMessage().contains("SNMP agent address")); + Assertions.assertFalse(exception.getMessage().contains("private-community")); + Assertions.assertInstanceOf(IllegalArgumentException.class, exception.getCause()); + } +} diff --git a/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSinkConfigTest.java b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSinkConfigTest.java new file mode 100644 index 0000000000..86a8d30e24 --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSinkConfigTest.java @@ -0,0 +1,119 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.config; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.Map; + +class SnmpSinkConfigTest { + + @Test + void testDefaultsAndSharedTargetOptions() { + SnmpSinkConfig config = new SnmpSinkConfig(ReadonlyConfig.fromMap(baseConfig())); + + Assertions.assertEquals("127.0.0.1", config.getHost()); + Assertions.assertEquals(161, config.getPort()); + Assertions.assertEquals("unit-test-community", config.getCommunity()); + Assertions.assertEquals(5000L, config.getTimeoutMillis()); + Assertions.assertEquals(1, config.getRetries()); + Assertions.assertEquals("oid", config.getOidField()); + Assertions.assertEquals("value", config.getValueField()); + Assertions.assertEquals("value_type", config.getValueTypeField()); + } + + @Test + void testCustomFieldMappingIsTrimmed() { + Map<String, Object> values = baseConfig(); + values.put("host", " 127.0.0.1 "); + values.put("community", " unit-test-community "); + values.put("oid_field", " target_oid "); + values.put("value_field", " target_value "); + values.put("value_type_field", " target_type "); + + SnmpSinkConfig config = new SnmpSinkConfig(ReadonlyConfig.fromMap(values)); + + Assertions.assertEquals("127.0.0.1", config.getHost()); + Assertions.assertEquals(" unit-test-community ", config.getCommunity()); + Assertions.assertEquals("target_oid", config.getOidField()); + Assertions.assertEquals("target_value", config.getValueField()); + Assertions.assertEquals("target_type", config.getValueTypeField()); + } + + @Test + void testInvalidTargetOptionsDoNotDiscloseCommunity() { + Map<String, Object> invalidPort = baseConfig(); + invalidPort.put("port", 65536); + + IllegalArgumentException exception = + Assertions.assertThrows( + IllegalArgumentException.class, + () -> new SnmpSinkConfig(ReadonlyConfig.fromMap(invalidPort))); + + Assertions.assertFalse(exception.getMessage().contains("unit-test-community")); + + Map<String, Object> invalidTimeout = baseConfig(); + invalidTimeout.put("timeout_millis", 0L); + Assertions.assertThrows( + IllegalArgumentException.class, + () -> new SnmpSinkConfig(ReadonlyConfig.fromMap(invalidTimeout))); + + Map<String, Object> invalidRetries = baseConfig(); + invalidRetries.put("retries", -1); + Assertions.assertThrows( + IllegalArgumentException.class, + () -> new SnmpSinkConfig(ReadonlyConfig.fromMap(invalidRetries))); + } + + @Test + void testMappedFieldsMustBeNonBlankAndDistinct() { + Map<String, Object> blank = baseConfig(); + blank.put("oid_field", " "); + Assertions.assertThrows( + IllegalArgumentException.class, + () -> new SnmpSinkConfig(ReadonlyConfig.fromMap(blank))); + + Map<String, Object> duplicate = baseConfig(); + duplicate.put("value_field", " oid "); + Assertions.assertThrows( + IllegalArgumentException.class, + () -> new SnmpSinkConfig(ReadonlyConfig.fromMap(duplicate))); + } + + @Test + void testDirectConstructionRejectsBlankHostAndCommunity() { + for (String key : new String[] {"host", "community"}) { + Map<String, Object> values = baseConfig(); + values.put(key, " \t\n"); + Assertions.assertThrows( + IllegalArgumentException.class, + () -> new SnmpSinkConfig(ReadonlyConfig.fromMap(values))); + } + } + + public static Map<String, Object> baseConfig() { + Map<String, Object> values = new HashMap<>(); + values.put("host", "127.0.0.1"); + values.put("community", "unit-test-community"); + return values; + } +} diff --git a/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceConfigTest.java b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceConfigTest.java index 9435fc1028..7cd3aa9396 100644 --- a/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceConfigTest.java +++ b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/config/SnmpSourceConfigTest.java @@ -22,12 +22,26 @@ import org.apache.seatunnel.api.configuration.ReadonlyConfig; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import java.io.ByteArrayInputStream; +import java.io.ObjectInputStream; import java.util.Arrays; +import java.util.Base64; import java.util.HashMap; import java.util.Map; class SnmpSourceConfigTest { + private static final String LEGACY_SERIALIZED_CONFIG = + "rO0ABXNyAEZvcmcuYXBhY2hlLnNlYXR1bm5lbC5jb25uZWN0b3JzLnNlYXR1bm5lbC5zbm1wLmNv" + + "bmZpZy5Tbm1wU291cmNlQ29uZmlnAAAAAAAAAAECAAdKABJwb2xsSW50ZXJ2YWxNaWxsaXNJAARw" + + "b3J0SQAHcmV0cmllc0oADXRpbWVvdXRNaWxsaXNMAAljb21tdW5pdHl0ABJMamF2YS9sYW5nL1N0" + + "cmluZztMAARob3N0cQB+AAFMAARvaWRzdAAQTGphdmEvdXRpbC9MaXN0O3hwAAAAAAAA6mAAAACh" + + "AAAAAQAAAAAAABOIdAAGcHVibGljdAAJMTI3LjAuMC4xc3IAI2phdmEudXRpbC5Db2xsZWN0aW9u" + + "cyRTaW5nbGV0b25MaXN0Ku8pEDynm5cCAAFMAAdlbGVtZW50dAASTGphdmEvbGFuZy9PYmplY3Q7" + + "eHBzcgASb3JnLnNubXA0ai5zbWkuT0lEaGJUgLBTOnQCAAFbAAV2YWx1ZXQAAltJeHIAH29yZy5z" + + "bm1wNGouc21pLkFic3RyYWN0VmFyaWFibGUTXwXE8DKuiAIAAHhwdXIAAltJTbpgJnbqsqUCAAB4" + + "cAAAAAQAAAABAAAAAwAAAAYAAAAB"; + @Test void testDefaultsAndOidNormalization() { SnmpSourceConfig config = new SnmpSourceConfig(ReadonlyConfig.fromMap(baseConfig())); @@ -117,6 +131,26 @@ class SnmpSourceConfigTest { () -> new SnmpSourceConfig(ReadonlyConfig.fromMap(interval))); } + @Test + void testDeserializesLegacySerializedLayout() throws Exception { + byte[] serialized = Base64.getDecoder().decode(LEGACY_SERIALIZED_CONFIG); + + SnmpSourceConfig config; + try (ObjectInputStream input = + new ObjectInputStream(new ByteArrayInputStream(serialized))) { + config = (SnmpSourceConfig) input.readObject(); + } + + Assertions.assertEquals("127.0.0.1", config.getHost()); + Assertions.assertEquals(161, config.getPort()); + Assertions.assertEquals("public", config.getCommunity()); + Assertions.assertEquals(5000L, config.getTimeoutMillis()); + Assertions.assertEquals(1, config.getRetries()); + Assertions.assertEquals(60000L, config.getPollIntervalMillis()); + Assertions.assertEquals(1, config.getOids().size()); + Assertions.assertEquals("1.3.6.1", config.getOids().get(0).toString()); + } + static Map<String, Object> baseConfig() { Map<String, Object> values = new HashMap<>(); values.put("host", "127.0.0.1"); diff --git a/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/Snmp4jSetClientTest.java b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/Snmp4jSetClientTest.java new file mode 100644 index 0000000000..7f516bdeb9 --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/Snmp4jSetClientTest.java @@ -0,0 +1,160 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.sink; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSinkConfig; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.snmp4j.CommandResponderEvent; +import org.snmp4j.MessageException; +import org.snmp4j.PDU; +import org.snmp4j.Snmp; +import org.snmp4j.mp.StatusInformation; +import org.snmp4j.smi.Integer32; +import org.snmp4j.smi.OID; +import org.snmp4j.smi.UdpAddress; +import org.snmp4j.smi.VariableBinding; +import org.snmp4j.transport.DefaultUdpTransportMapping; + +import java.io.IOException; +import java.net.InetAddress; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.atomic.AtomicReference; + +class Snmp4jSetClientTest { + + @Test + void testBuildsSingleBindingSetRequest() { + SnmpSetRequest request = + new SnmpSetRequest(new OID("1.3.6.1.2.1.1.5.0"), new Integer32(42)); + + PDU pdu = Snmp4jSetClient.buildSetRequest(request); + + Assertions.assertEquals(PDU.SET, pdu.getType()); + Assertions.assertEquals(1, pdu.size()); + Assertions.assertEquals("1.3.6.1.2.1.1.5.0", pdu.get(0).getOid().toString()); + Assertions.assertEquals("42", pdu.get(0).getVariable().toString()); + } + + @Test + void testSendsSetToLoopbackSnmpAgent() throws Exception { + AtomicReference<VariableBinding> received = new AtomicReference<>(); + try (LoopbackAgent agent = new LoopbackAgent(PDU.noError, received); + Snmp4jSetClient client = new Snmp4jSetClient(config(agent.getPort()))) { + client.set( + new SnmpSetRequest( + new OID("1.3.6.1.2.1.1.5.0"), + SnmpSinkRowConverter.parseVariable("OctetString", "router-1"))); + } + + Assertions.assertNotNull(received.get()); + Assertions.assertEquals("1.3.6.1.2.1.1.5.0", received.get().getOid().toString()); + Assertions.assertEquals("router-1", received.get().getVariable().toString()); + } + + @Test + void testRemoteErrorStatusIsReported() throws Exception { + AtomicReference<VariableBinding> received = new AtomicReference<>(); + try (LoopbackAgent agent = new LoopbackAgent(PDU.notWritable, received); + Snmp4jSetClient client = new Snmp4jSetClient(config(agent.getPort()))) { + IOException exception = + Assertions.assertThrows( + IOException.class, + () -> + client.set( + new SnmpSetRequest( + new OID("1.3.6.1.2.1.1.5.0"), + new Integer32(42)))); + + Assertions.assertTrue( + exception.getMessage().contains("error status " + PDU.notWritable)); + Assertions.assertTrue(exception.getMessage().contains("index 1")); + Assertions.assertFalse(exception.getMessage().contains("unit-test-community")); + } + } + + private static SnmpSinkConfig config(int port) { + Map<String, Object> values = new HashMap<>(); + values.put("host", "127.0.0.1"); + values.put("community", "unit-test-community"); + values.put("port", port); + values.put("timeout_millis", 1000L); + values.put("retries", 0); + return new SnmpSinkConfig(ReadonlyConfig.fromMap(values)); + } + + private static final class LoopbackAgent implements AutoCloseable { + private final int errorStatus; + private final AtomicReference<VariableBinding> received; + private final DefaultUdpTransportMapping transport; + private final Snmp agent; + private final AtomicReference<MessageException> responseFailure = new AtomicReference<>(); + + private LoopbackAgent(int errorStatus, AtomicReference<VariableBinding> received) + throws IOException { + this.errorStatus = errorStatus; + this.received = received; + this.transport = + new DefaultUdpTransportMapping( + new UdpAddress(InetAddress.getLoopbackAddress(), 0)); + this.agent = new Snmp(transport); + agent.addCommandResponder(this::respond); + agent.listen(); + } + + private int getPort() { + return transport.getListenAddress().getPort(); + } + + private void respond(CommandResponderEvent event) { + PDU request = event.getPDU(); + received.set(request.get(0)); + PDU response = new PDU(request); + response.setType(PDU.RESPONSE); + response.setErrorStatus(errorStatus); + response.setErrorIndex(errorStatus == PDU.noError ? 0 : 1); + try { + event.getMessageDispatcher() + .returnResponsePdu( + event.getMessageProcessingModel(), + event.getSecurityModel(), + event.getSecurityName(), + event.getSecurityLevel(), + response, + event.getMaxSizeResponsePDU(), + event.getStateReference(), + new StatusInformation()); + event.setProcessed(true); + } catch (MessageException e) { + responseFailure.set(e); + } + } + + @Override + public void close() throws IOException { + agent.close(); + if (responseFailure.get() != null) { + throw new IOException( + "Loopback SNMP agent failed to send response", responseFailure.get()); + } + } + } +} diff --git a/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkFactoryTest.java b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkFactoryTest.java new file mode 100644 index 0000000000..215678cc8a --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkFactoryTest.java @@ -0,0 +1,189 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.sink; + +import org.apache.seatunnel.api.configuration.Option; +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.api.configuration.util.ConfigValidator; +import org.apache.seatunnel.api.configuration.util.OptionRule; +import org.apache.seatunnel.api.configuration.util.OptionValidationException; +import org.apache.seatunnel.api.table.catalog.CatalogTable; +import org.apache.seatunnel.api.table.catalog.PhysicalColumn; +import org.apache.seatunnel.api.table.catalog.TableIdentifier; +import org.apache.seatunnel.api.table.catalog.TableSchema; +import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext; +import org.apache.seatunnel.api.table.type.BasicType; +import org.apache.seatunnel.common.utils.SerializationUtils; +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSinkOptions; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +class SnmpSinkFactoryTest { + + @Test + void testOptionRuleAcceptsDefaultsAndNumericBoundaries() { + validate(baseConfig()); + Map<String, Object> values = baseConfig(); + values.put("port", 1); + values.put("timeout_millis", 1L); + values.put("retries", 0); + validate(values); + values.put("port", 65535); + values.put("timeout_millis", Long.MAX_VALUE); + values.put("retries", Integer.MAX_VALUE); + validate(values); + } + + @Test + void testOptionRuleRejectsMissingRequiredOptions() { + for (String key : new String[] {"host", "community"}) { + Map<String, Object> values = baseConfig(); + values.remove(key); + Assertions.assertThrows(OptionValidationException.class, () -> validate(values), key); + } + } + + @Test + void testOptionRuleRejectsBlankStrings() { + for (String key : + new String[] { + "host", "community", "oid_field", "value_field", "value_type_field" + }) { + assertInvalidOption(key, ""); + assertInvalidOption(key, " \t\n"); + } + } + + @Test + void testOptionRuleRejectsInvalidNumericValuesWithoutDisclosingCommunity() { + assertInvalidOption("port", 0); + assertInvalidOption("port", 65536); + assertInvalidOption("timeout_millis", 0L); + assertInvalidOption("timeout_millis", -1L); + assertInvalidOption("retries", -1); + } + + @Test + void testOptionRulePreservesWhitespaceForRuntimeNormalization() { + Map<String, Object> values = baseConfig(); + values.put("host", " 127.0.0.1 "); + values.put("community", " unit-test-community "); + values.put("oid_field", " oid "); + values.put("value_field", " value "); + values.put("value_type_field", " value_type "); + ReadonlyConfig config = ReadonlyConfig.fromMap(values); + ConfigValidator.of(config).validate(new SnmpSinkFactory().optionRule()); + Assertions.assertEquals(" unit-test-community ", config.get(SnmpSinkOptions.COMMUNITY)); + } + + private static void assertInvalidOption(String key, Object value) { + Map<String, Object> values = baseConfig(); + values.put(key, value); + OptionValidationException exception = + Assertions.assertThrows(OptionValidationException.class, () -> validate(values)); + Assertions.assertTrue(exception.getMessage().contains(key)); + Assertions.assertFalse(exception.getMessage().contains("unit-test-community")); + } + + private static void validate(Map<String, Object> values) { + ConfigValidator.of(ReadonlyConfig.fromMap(values)) + .validate(new SnmpSinkFactory().optionRule()); + } + + @Test + void testFactoryIdentityAndOptions() { + SnmpSinkFactory factory = new SnmpSinkFactory(); + + Assertions.assertEquals(SnmpSinkOptions.CONNECTOR_IDENTITY, factory.factoryIdentifier()); + OptionRule rule = factory.optionRule(); + List<Option<?>> required = + rule.getRequiredOptions().stream() + .flatMap(group -> group.getOptions().stream()) + .collect(Collectors.toList()); + Assertions.assertTrue(required.contains(SnmpSinkOptions.HOST)); + Assertions.assertTrue(required.contains(SnmpSinkOptions.COMMUNITY)); + Assertions.assertTrue(rule.getOptionalOptions().contains(SnmpSinkOptions.PORT)); + Assertions.assertTrue(rule.getOptionalOptions().contains(SnmpSinkOptions.TIMEOUT_MILLIS)); + Assertions.assertTrue(rule.getOptionalOptions().contains(SnmpSinkOptions.RETRIES)); + Assertions.assertTrue(rule.getOptionalOptions().contains(SnmpSinkOptions.OID_FIELD)); + Assertions.assertTrue(rule.getOptionalOptions().contains(SnmpSinkOptions.VALUE_FIELD)); + Assertions.assertTrue(rule.getOptionalOptions().contains(SnmpSinkOptions.VALUE_TYPE_FIELD)); + } + + @Test + void testFactoryCreatesSerializableSinkWithSourceCompatibleSchema() { + CatalogTable catalogTable = catalogTable(); + TableSinkFactoryContext context = + new TableSinkFactoryContext( + catalogTable, + ReadonlyConfig.fromMap(baseConfig()), + getClass().getClassLoader()); + + SnmpSink sink = (SnmpSink) new SnmpSinkFactory().createSink(context).createSink(); + SnmpSink restored = SerializationUtils.deserialize(SerializationUtils.serialize(sink)); + + Assertions.assertEquals("SNMP", restored.getPluginName()); + Assertions.assertArrayEquals( + catalogTable.getSeaTunnelRowType().getFieldNames(), + restored.getWriteCatalogTable() + .orElseThrow(() -> new AssertionError("Sink catalog table is missing")) + .getSeaTunnelRowType() + .getFieldNames()); + } + + private static CatalogTable catalogTable() { + TableSchema schema = + TableSchema.builder() + .column( + PhysicalColumn.of( + "agent", BasicType.STRING_TYPE, 0, false, null, null)) + .column( + PhysicalColumn.of( + "oid", BasicType.STRING_TYPE, 0, false, null, null)) + .column( + PhysicalColumn.of( + "value", BasicType.STRING_TYPE, 0, false, null, null)) + .column( + PhysicalColumn.of( + "value_type", BasicType.STRING_TYPE, 0, false, null, null)) + .column( + PhysicalColumn.of( + "poll_time", BasicType.LONG_TYPE, 0, false, null, null)) + .build(); + return CatalogTable.of( + TableIdentifier.of("default", "default", "snmp_sink_test"), + schema, + Collections.emptyMap(), + Collections.emptyList(), + "SNMP sink test table"); + } + + private static Map<String, Object> baseConfig() { + Map<String, Object> values = new HashMap<>(); + values.put("host", "127.0.0.1"); + values.put("community", "unit-test-community"); + return values; + } +} diff --git a/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkRowConverterTest.java b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkRowConverterTest.java new file mode 100644 index 0000000000..60bc34ed01 --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkRowConverterTest.java @@ -0,0 +1,246 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.sink; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.api.table.type.BasicType; +import org.apache.seatunnel.api.table.type.SeaTunnelDataType; +import org.apache.seatunnel.api.table.type.SeaTunnelRow; +import org.apache.seatunnel.api.table.type.SeaTunnelRowType; +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSinkConfig; +import org.apache.seatunnel.connectors.seatunnel.snmp.exception.SnmpConnectorException; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.snmp4j.smi.Counter32; +import org.snmp4j.smi.Counter64; +import org.snmp4j.smi.Gauge32; +import org.snmp4j.smi.Integer32; +import org.snmp4j.smi.IpAddress; +import org.snmp4j.smi.OID; +import org.snmp4j.smi.OctetString; +import org.snmp4j.smi.TimeTicks; +import org.snmp4j.smi.UnsignedInteger32; +import org.snmp4j.smi.Variable; + +import java.nio.charset.StandardCharsets; +import java.util.HashMap; +import java.util.Map; + +class SnmpSinkRowConverterTest { + + @Test + void testConvertsSourceCompatibleSchema() { + SnmpSinkRowConverter converter = new SnmpSinkRowConverter(config(), sourceRowType()); + + SnmpSetRequest request = + converter.convert( + new SeaTunnelRow( + new Object[] { + "127.0.0.1:161", + ".1.3.6.1.2.1.1.5.0", + "router-1", + "OctetString", + 1234L + })); + + Assertions.assertEquals("1.3.6.1.2.1.1.5.0", request.getOid().toString()); + Assertions.assertEquals("router-1", request.getValue().toString()); + Assertions.assertInstanceOf(OctetString.class, request.getValue()); + } + + @Test + void testSupportsDocumentedSmiTypes() { + assertVariable("Integer32", "-42", Integer32.class, "-42"); + assertVariable("UnsignedInteger32", "4294967295", UnsignedInteger32.class, "4294967295"); + assertVariable("Counter32", "12", Counter32.class, "12"); + assertVariable("Gauge32", "13", Gauge32.class, "13"); + assertVariable("TimeTicks", "14", TimeTicks.class, "0:00:00.14"); + assertVariable( + "Counter64", "18446744073709551615", Counter64.class, "18446744073709551615"); + assertVariable("OctetString", "router-1", OctetString.class, "router-1"); + + Variable hex = SnmpSinkRowConverter.parseVariable("OctetStringHex", "00ff10"); + Assertions.assertInstanceOf(OctetString.class, hex); + Assertions.assertArrayEquals( + new byte[] {0x00, (byte) 0xff, 0x10}, ((OctetString) hex).getValue()); + + assertVariable("OID", "1.3.6.1.2.1", OID.class, "1.3.6.1.2.1"); + assertVariable("IpAddress", "192.0.2.10", IpAddress.class, "192.0.2.10"); + } + + @Test + void testSupportsSnmp4jSourceSyntaxStringsAndValues() { + assertSourceVariable(new Integer32(-42), Integer32.class); + assertSourceVariable(new UnsignedInteger32(11), Gauge32.class); + assertSourceVariable(new Counter32(12), Counter32.class); + assertSourceVariable(new Gauge32(13), Gauge32.class); + assertSourceVariable(new TimeTicks(14), TimeTicks.class); + assertSourceVariable(new TimeTicks(172_800_014L), TimeTicks.class); + assertSourceVariable(new Counter64(15), Counter64.class); + assertSourceVariable(new OctetString("router-1"), OctetString.class); + assertSourceVariable(new OID("1.3.6.1.2.1"), OID.class); + assertSourceVariable(new IpAddress("192.0.2.10"), IpAddress.class); + } + + @Test + void testPreservesOctetStringPayload() { + SnmpSinkRowConverter converter = new SnmpSinkRowConverter(config(), sinkRowType()); + + SnmpSetRequest whitespace = + converter.convert(row("1.3.6.1.2.1.1.5.0", " router-1 ", " OctetString ")); + Assertions.assertEquals(" router-1 ", whitespace.getValue().toString()); + + SnmpSetRequest empty = converter.convert(row("1.3.6.1.2.1.1.5.0", "", "OctetString")); + Assertions.assertEquals(0, ((OctetString) empty.getValue()).length()); + + SnmpSetRequest utf8 = + converter.convert(row("1.3.6.1.2.1.1.5.0", "router-\u03b1", "OctetString")); + Assertions.assertArrayEquals( + "router-\u03b1".getBytes(StandardCharsets.UTF_8), + ((OctetString) utf8.getValue()).getValue()); + } + + @Test + void testRejectsInvalidSchemaBeforeClientCreation() { + SeaTunnelRowType missingField = + new SeaTunnelRowType( + new String[] {"oid", "value"}, + new SeaTunnelDataType[] {BasicType.STRING_TYPE, BasicType.STRING_TYPE}); + SnmpConnectorException missing = + Assertions.assertThrows( + SnmpConnectorException.class, + () -> new SnmpSinkRowConverter(config(), missingField)); + Assertions.assertTrue(missing.getMessage().contains("value_type")); + + SeaTunnelRowType wrongType = + new SeaTunnelRowType( + new String[] {"oid", "value", "value_type"}, + new SeaTunnelDataType[] { + BasicType.STRING_TYPE, BasicType.INT_TYPE, BasicType.STRING_TYPE + }); + SnmpConnectorException invalidType = + Assertions.assertThrows( + SnmpConnectorException.class, + () -> new SnmpSinkRowConverter(config(), wrongType)); + Assertions.assertTrue(invalidType.getMessage().contains("must use STRING")); + } + + @Test + void testRejectsInvalidRowsBeforeNetworkIo() { + SnmpSinkRowConverter converter = new SnmpSinkRowConverter(config(), sinkRowType()); + + Assertions.assertThrows( + SnmpConnectorException.class, + () -> converter.convert(row("not-an-oid", "1", "Integer32"))); + Assertions.assertThrows( + SnmpConnectorException.class, + () -> converter.convert(row("1.3.6.1.2.1.1.5.0", "1", "UnknownType"))); + Assertions.assertThrows( + SnmpConnectorException.class, + () -> converter.convert(row("1.3.6.1.2.1.1.5.0", "4294967296", "Counter32"))); + Assertions.assertThrows( + SnmpConnectorException.class, + () -> converter.convert(row("1.3.6.1.2.1.1.5.0", "0fg1", "OctetStringHex"))); + Assertions.assertThrows( + SnmpConnectorException.class, + () -> converter.convert(row("1.3.6.1.2.1.1.5.0", "999.0.2.1", "IpAddress"))); + Assertions.assertThrows( + SnmpConnectorException.class, + () -> converter.convert(row("1.3.6.1.2.1.1.5.0", null, "OctetString"))); + Assertions.assertThrows( + SnmpConnectorException.class, + () -> + converter.convert( + new SeaTunnelRow( + new Object[] {"1.3.6.1.2.1.1.5.0", 1, "Integer32"}))); + Assertions.assertThrows( + SnmpConnectorException.class, + () -> + converter.convert( + row("1.3.6.1.2.1.1.5.0", "18446744073709551616", "Counter64"))); + Assertions.assertThrows( + SnmpConnectorException.class, + () -> converter.convert(row("1.3.6.1.2.1.1.5.0", "0:60:00.00", "TimeTicks"))); + Assertions.assertThrows( + SnmpConnectorException.class, + () -> + converter.convert( + row("1.3.6.1.2.1.1.5.0", "2 day, 0:00:00.00", "TimeTicks"))); + Assertions.assertThrows( + SnmpConnectorException.class, + () -> + converter.convert( + row("1.3.6.1.2.1.1.5.0", "0 days, 0:00:00.00", "TimeTicks"))); + Assertions.assertThrows( + SnmpConnectorException.class, + () -> + converter.convert( + row("1.3.6.1.2.1.1.5.0", "498 days, 0:00:00.00", "TimeTicks"))); + } + + private static void assertVariable( + String type, + String value, + Class<? extends Variable> expectedClass, + String expectedText) { + Variable variable = SnmpSinkRowConverter.parseVariable(type, value); + Assertions.assertInstanceOf(expectedClass, variable); + Assertions.assertEquals(expectedText, variable.toString()); + } + + private static void assertSourceVariable( + Variable sourceVariable, Class<? extends Variable> expectedClass) { + Variable converted = + SnmpSinkRowConverter.parseVariable( + sourceVariable.getSyntaxString(), sourceVariable.toString()); + Assertions.assertInstanceOf(expectedClass, converted); + Assertions.assertEquals(sourceVariable.toString(), converted.toString()); + } + + private static SeaTunnelRow row(String oid, String value, String valueType) { + return new SeaTunnelRow(new Object[] {oid, value, valueType}); + } + + static SeaTunnelRowType sinkRowType() { + return new SeaTunnelRowType( + new String[] {"oid", "value", "value_type"}, + new SeaTunnelDataType[] { + BasicType.STRING_TYPE, BasicType.STRING_TYPE, BasicType.STRING_TYPE + }); + } + + private static SeaTunnelRowType sourceRowType() { + return new SeaTunnelRowType( + new String[] {"agent", "oid", "value", "value_type", "poll_time"}, + new SeaTunnelDataType[] { + BasicType.STRING_TYPE, + BasicType.STRING_TYPE, + BasicType.STRING_TYPE, + BasicType.STRING_TYPE, + BasicType.LONG_TYPE + }); + } + + static SnmpSinkConfig config() { + Map<String, Object> values = new HashMap<>(); + values.put("host", "127.0.0.1"); + values.put("community", "unit-test-community"); + return new SnmpSinkConfig(ReadonlyConfig.fromMap(values)); + } +} diff --git a/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkWriterTest.java b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkWriterTest.java new file mode 100644 index 0000000000..a268a01931 --- /dev/null +++ b/seatunnel-connectors-v2/connector-snmp/src/test/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkWriterTest.java @@ -0,0 +1,133 @@ +/* + * 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.seatunnel.connectors.seatunnel.snmp.sink; + +import org.apache.seatunnel.api.table.type.SeaTunnelRow; +import org.apache.seatunnel.connectors.seatunnel.snmp.config.SnmpSinkConfig; +import org.apache.seatunnel.connectors.seatunnel.snmp.exception.SnmpConnectorException; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.io.IOException; + +class SnmpSinkWriterTest { + + @Test + void testWritesOneSynchronousSetPerRowAndClosesClient() { + FakeSnmpSetClient client = new FakeSnmpSetClient(); + SnmpSinkWriter writer = writer(client); + + writer.write(row()); + writer.close(); + + Assertions.assertEquals(1, client.writeCount); + Assertions.assertEquals("1.3.6.1.2.1.1.5.0", client.request.getOid().toString()); + Assertions.assertEquals("router-1", client.request.getValue().toString()); + Assertions.assertTrue(client.closed); + } + + @Test + void testInvalidRowDoesNotReachClient() { + FakeSnmpSetClient client = new FakeSnmpSetClient(); + SnmpSinkWriter writer = writer(client); + + Assertions.assertThrows( + SnmpConnectorException.class, + () -> + writer.write( + new SeaTunnelRow( + new Object[] {"invalid", "router-1", "OctetString"}))); + + Assertions.assertEquals(0, client.writeCount); + } + + @Test + void testWriteFailureUsesConnectorErrorWithoutCommunity() { + FakeSnmpSetClient client = new FakeSnmpSetClient(); + client.writeFailure = new IOException("remote timeout"); + SnmpSinkWriter writer = writer(client); + + SnmpConnectorException exception = + Assertions.assertThrows(SnmpConnectorException.class, () -> writer.write(row())); + + Assertions.assertTrue(exception.getMessage().contains("SNMP-03")); + Assertions.assertTrue(exception.getMessage().contains("1.3.6.1.2.1.1.5.0")); + Assertions.assertFalse(exception.getMessage().contains("unit-test-community")); + } + + @Test + void testConnectionAndCloseFailuresUseConnectorErrorsWithoutCommunity() { + SnmpSinkConfig config = SnmpSinkRowConverterTest.config(); + SnmpConnectorException connectionFailure = + Assertions.assertThrows( + SnmpConnectorException.class, + () -> + new SnmpSinkWriter( + config, + SnmpSinkRowConverterTest.sinkRowType(), + ignored -> { + throw new IOException("bind failed"); + })); + Assertions.assertTrue(connectionFailure.getMessage().contains("SNMP-01")); + Assertions.assertFalse(connectionFailure.getMessage().contains("unit-test-community")); + + FakeSnmpSetClient client = new FakeSnmpSetClient(); + client.closeFailure = new IOException("close failed"); + SnmpConnectorException closeFailure = + Assertions.assertThrows(SnmpConnectorException.class, () -> writer(client).close()); + Assertions.assertTrue(closeFailure.getMessage().contains("SNMP-06")); + Assertions.assertFalse(closeFailure.getMessage().contains("unit-test-community")); + } + + private static SnmpSinkWriter writer(FakeSnmpSetClient client) { + return new SnmpSinkWriter( + SnmpSinkRowConverterTest.config(), + SnmpSinkRowConverterTest.sinkRowType(), + ignored -> client); + } + + private static SeaTunnelRow row() { + return new SeaTunnelRow(new Object[] {"1.3.6.1.2.1.1.5.0", "router-1", "OctetString"}); + } + + private static final class FakeSnmpSetClient implements SnmpSetClient { + private int writeCount; + private boolean closed; + private SnmpSetRequest request; + private IOException writeFailure; + private IOException closeFailure; + + @Override + public void set(SnmpSetRequest request) throws IOException { + writeCount++; + this.request = request; + if (writeFailure != null) { + throw writeFailure; + } + } + + @Override + public void close() throws IOException { + closed = true; + if (closeFailure != null) { + throw closeFailure; + } + } + } +} diff --git a/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigShadeUtils.java b/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigShadeUtils.java index 52fa94c6a1..706ffc8237 100644 --- a/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigShadeUtils.java +++ b/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigShadeUtils.java @@ -55,7 +55,7 @@ public final class ConfigShadeUtils { new String[] {"password", "username", "auth", "token", "access_key", "secret_key"}; private static final String[] DEFAULT_LOG_MASK_ONLY_KEYWORDS = - new String[] {"sasl.jaas.config"}; + new String[] {"sasl.jaas.config", "community"}; private static final Map<String, ConfigShade> CONFIG_SHADES = new HashMap<>(); diff --git a/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/utils/ConfigBuilderTest.java b/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/utils/ConfigBuilderTest.java index 43f711de3c..fe1b71f7a9 100644 --- a/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/utils/ConfigBuilderTest.java +++ b/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/utils/ConfigBuilderTest.java @@ -101,6 +101,25 @@ public class ConfigBuilderTest { Assertions.assertEquals("localhost:9092", desensitizedKafkaConfig.get("bootstrap.servers")); } + @Test + public void testConfigDesensitizationMasksSnmpCommunity() { + Map<String, Object> sink = new LinkedHashMap<>(); + sink.put("host", "127.0.0.1"); + sink.put("community", "private-community"); + + Map<String, Object> config = new LinkedHashMap<>(); + config.put("sink", Arrays.asList(sink)); + + Map<String, Object> desensitized = + ConfigBuilder.configDesensitization( + config, ConfigShadeUtils.getLogDesensitizationOptions(null)); + List<?> sinks = (List<?>) desensitized.get("sink"); + Map<?, ?> desensitizedSink = (Map<?, ?>) sinks.get(0); + + Assertions.assertEquals("******", desensitizedSink.get("community")); + Assertions.assertEquals("127.0.0.1", desensitizedSink.get("host")); + } + @Test public void testConfigDesensitizationMasksS3CredentialOptions() { Map<String, Object> accessKeyConfig = new LinkedHashMap<>(); diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/pom.xml b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/pom.xml new file mode 100644 index 0000000000..b5321d9de2 --- /dev/null +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/pom.xml @@ -0,0 +1,42 @@ +<?xml version="1.0" encoding="UTF-8"?> +<!-- + 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. +--> +<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/xsd/maven-4.0.0.xsd"> + <modelVersion>4.0.0</modelVersion> + <parent> + <groupId>org.apache.seatunnel</groupId> + <artifactId>seatunnel-connector-v2-e2e</artifactId> + <version>${revision}</version> + </parent> + + <artifactId>connector-snmp-e2e</artifactId> + <name>SeaTunnel : E2E : Connector V2 : SNMP</name> + + <dependencies> + <dependency> + <groupId>org.apache.seatunnel</groupId> + <artifactId>connector-fake</artifactId> + <version>${project.version}</version> + <scope>test</scope> + </dependency> + <dependency> + <groupId>org.apache.seatunnel</groupId> + <artifactId>connector-snmp</artifactId> + <version>${project.version}</version> + <scope>test</scope> + </dependency> + </dependencies> +</project> diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/src/test/java/org/apache/seatunnel/e2e/connector/snmp/SnmpAgent.java b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/src/test/java/org/apache/seatunnel/e2e/connector/snmp/SnmpAgent.java new file mode 100644 index 0000000000..325e904035 --- /dev/null +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/src/test/java/org/apache/seatunnel/e2e/connector/snmp/SnmpAgent.java @@ -0,0 +1,95 @@ +/* + * 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.seatunnel.e2e.connector.snmp; + +import org.snmp4j.CommandResponderEvent; +import org.snmp4j.PDU; +import org.snmp4j.Snmp; +import org.snmp4j.mp.StatusInformation; +import org.snmp4j.smi.UdpAddress; +import org.snmp4j.smi.VariableBinding; +import org.snmp4j.transport.DefaultUdpTransportMapping; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.Paths; +import java.nio.file.StandardOpenOption; +import java.util.concurrent.CountDownLatch; + +/** Minimal SNMP responder used to verify the sink over a real UDP transport. */ +public final class SnmpAgent { + + private static final int PORT = 1161; + private static final Path OUTPUT = Paths.get("/tmp/snmp-set.txt"); + + private SnmpAgent() {} + + public static void main(String[] args) throws Exception { + DefaultUdpTransportMapping transport = + new DefaultUdpTransportMapping(new UdpAddress("0.0.0.0/" + PORT)); + Snmp agent = new Snmp(transport); + agent.addCommandResponder(SnmpAgent::respond); + Runtime.getRuntime().addShutdownHook(new Thread(() -> close(agent))); + agent.listen(); + System.out.println("snmp-agent-ready"); + new CountDownLatch(1).await(); + } + + private static void respond(CommandResponderEvent event) { + PDU request = event.getPDU(); + if (request == null || request.getType() != PDU.SET || request.size() != 1) { + return; + } + + VariableBinding binding = request.get(0); + try { + Files.write( + OUTPUT, + (binding.getOid() + "=" + binding.getVariable() + System.lineSeparator()) + .getBytes(StandardCharsets.UTF_8), + StandardOpenOption.CREATE, + StandardOpenOption.TRUNCATE_EXISTING); + PDU response = new PDU(request); + response.setType(PDU.RESPONSE); + response.setErrorStatus(PDU.noError); + response.setErrorIndex(0); + event.getMessageDispatcher() + .returnResponsePdu( + event.getMessageProcessingModel(), + event.getSecurityModel(), + event.getSecurityName(), + event.getSecurityLevel(), + response, + event.getMaxSizeResponsePDU(), + event.getStateReference(), + new StatusInformation()); + event.setProcessed(true); + } catch (IOException e) { + throw new IllegalStateException("Failed to handle SNMP SET request", e); + } + } + + private static void close(Snmp agent) { + try { + agent.close(); + } catch (IOException e) { + System.err.println("Failed to close SNMP agent: " + e.getMessage()); + } + } +} diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/src/test/java/org/apache/seatunnel/e2e/connector/snmp/SnmpSinkIT.java b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/src/test/java/org/apache/seatunnel/e2e/connector/snmp/SnmpSinkIT.java new file mode 100644 index 0000000000..0198713c15 --- /dev/null +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/src/test/java/org/apache/seatunnel/e2e/connector/snmp/SnmpSinkIT.java @@ -0,0 +1,109 @@ +/* + * 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.seatunnel.e2e.connector.snmp; + +import org.apache.seatunnel.e2e.common.TestResource; +import org.apache.seatunnel.e2e.common.TestSuiteBase; +import org.apache.seatunnel.e2e.common.container.EngineType; +import org.apache.seatunnel.e2e.common.container.TestContainer; +import org.apache.seatunnel.e2e.common.junit.DisabledOnContainer; + +import org.awaitility.Awaitility; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.TestTemplate; +import org.snmp4j.Snmp; +import org.testcontainers.containers.Container; +import org.testcontainers.containers.GenericContainer; +import org.testcontainers.containers.wait.strategy.Wait; +import org.testcontainers.lifecycle.Startables; +import org.testcontainers.utility.DockerImageName; +import org.testcontainers.utility.MountableFile; + +import java.net.URISyntaxException; +import java.nio.file.Path; +import java.nio.file.Paths; +import java.time.Duration; +import java.util.stream.Stream; + +@DisabledOnContainer( + value = {}, + type = {EngineType.SPARK, EngineType.FLINK}, + disabledReason = "The first SNMP sink E2E slice targets the Zeta engine") +public class SnmpSinkIT extends TestSuiteBase implements TestResource { + + private static final DockerImageName JAVA_IMAGE = + DockerImageName.parse("eclipse-temurin:11-jre-jammy"); + private static final String EXPECTED_SET = "1.3.6.1.2.1.1.4.0=seatunnel-e2e"; + + private GenericContainer<?> snmpAgent; + + @Override + @BeforeAll + public void startUp() throws URISyntaxException { + Path testClasses = + Paths.get( + SnmpAgent.class + .getProtectionDomain() + .getCodeSource() + .getLocation() + .toURI()); + Path snmp4jJar = + Paths.get(Snmp.class.getProtectionDomain().getCodeSource().getLocation().toURI()); + + snmpAgent = + new GenericContainer<>(JAVA_IMAGE) + .withNetwork(NETWORK) + .withNetworkAliases("snmp-agent") + .withCopyFileToContainer( + MountableFile.forHostPath(testClasses), "/opt/snmp-agent/classes") + .withCopyFileToContainer( + MountableFile.forHostPath(snmp4jJar), "/opt/snmp-agent/snmp4j.jar") + .withCommand( + "java", + "-cp", + "/opt/snmp-agent/classes:/opt/snmp-agent/snmp4j.jar", + SnmpAgent.class.getName()) + .waitingFor( + Wait.forLogMessage(".*snmp-agent-ready.*\\n", 1) + .withStartupTimeout(Duration.ofMinutes(1))); + Startables.deepStart(Stream.of(snmpAgent)).join(); + } + + @Override + public void tearDown() { + if (snmpAgent != null) { + snmpAgent.close(); + } + } + + @TestTemplate + public void testFakeSourceWritesSnmpSet(TestContainer container) throws Exception { + Container.ExecResult execResult = container.executeJob("/fake_to_snmp.conf"); + Assertions.assertEquals(0, execResult.getExitCode(), execResult.getStderr()); + + Awaitility.await() + .atMost(Duration.ofSeconds(10)) + .untilAsserted( + () -> { + Container.ExecResult result = + snmpAgent.execInContainer("cat", "/tmp/snmp-set.txt"); + Assertions.assertEquals(0, result.getExitCode(), result.getStderr()); + Assertions.assertEquals(EXPECTED_SET, result.getStdout().trim()); + }); + } +} diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/src/test/resources/fake_to_snmp.conf b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/src/test/resources/fake_to_snmp.conf new file mode 100644 index 0000000000..71d2beba28 --- /dev/null +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-snmp-e2e/src/test/resources/fake_to_snmp.conf @@ -0,0 +1,50 @@ +# +# 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. +# + +env { + parallelism = 1 + job.mode = "BATCH" +} + +source { + FakeSource { + row.num = 1 + schema = { + fields { + oid = string + value = string + value_type = string + } + } + rows = [ + { + kind = INSERT + fields = ["1.3.6.1.2.1.1.4.0", "seatunnel-e2e", "OctetString"] + } + ] + } +} + +sink { + SNMP { + host = "snmp-agent" + port = 1161 + community = "e2e-community" + timeout_millis = 5000 + retries = 0 + } +} diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml b/seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml index eebf1fce58..fa8c1d20bb 100644 --- a/seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml @@ -73,6 +73,7 @@ <module>connector-druid-e2e</module> <module>connector-google-firestore-e2e</module> <module>connector-google-pubsub-e2e</module> + <module>connector-snmp-e2e</module> <module>connector-azure-queue-storage-e2e</module> <module>connector-rocketmq-e2e</module> <!-- <module>connector-file-obs-e2e</module>-->
