diff --git a/.github/workflows/e2e.alarm.yaml b/.github/workflows/e2e.alarm.yaml new file mode 100644 index 000000000..87c2ea2c3 --- /dev/null +++ b/.github/workflows/e2e.alarm.yaml @@ -0,0 +1,57 @@ +# 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. + +name: E2E + +on: + pull_request: + paths: + - '**' + - '!**.md' + schedule: + - cron: '0 18 * * *' + +env: + SW_AGENT_JDK_VERSION: 8 + +jobs: + Alarm: + name: Alarm + runs-on: ubuntu-latest + strategy: + matrix: + storage: ['h2', 'mysql', 'es6', 'es7', 'influxdb', 'postgres'] + env: + SW_STORAGE: ${{ matrix.storage }} + steps: + - uses: actions/checkout@v2 + with: + submodules: true + - name: Set Skip Env Var + uses: ./.github/actions/skip + - name: Cache local Maven repository + if: env.SKIP_CI != 'true' + uses: actions/cache@v2 + with: + path: ~/.m2/repository + key: ${{ runner.os }}-maven-${{ hashFiles('**/pom.xml') }} + restore-keys: | + ${{ runner.os }}-maven- + - name: Run E2E Test + if: env.SKIP_CI != 'true' + uses: ./.github/actions/e2e-test + with: + test_class: org.apache.skywalking.e2e.alarm.AlarmE2E diff --git a/.github/workflows/e2e.yaml b/.github/workflows/e2e.yaml index 565561d9b..e1c071f5f 100644 --- a/.github/workflows/e2e.yaml +++ b/.github/workflows/e2e.yaml @@ -71,8 +71,6 @@ jobs: class: org.apache.skywalking.e2e.GatewayE2E - name: Meter class: org.apache.skywalking.e2e.meter.MeterE2E - - name: Alarm - class: org.apache.skywalking.e2e.alarm.AlarmE2E - name: Zabbix class: org.apache.skywalking.e2e.zabbix.ZabbixE2E steps: diff --git a/CHANGES.md b/CHANGES.md index e9bc8714f..da0490198 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -22,6 +22,7 @@ Release Notes. * BugFix: filter invalid Envoy access logs whose socket address is empty. * Fix K8s monitoring the incorrect metrics calculate. * Loop alarm into event system. +* Support alarm tags. #### UI * Add logo for kong plugin. diff --git a/dist-material/config-examples/alarm-settings.yml b/dist-material/config-examples/alarm-settings.yml index c5f7d551e..c8bbab35e 100644 --- a/dist-material/config-examples/alarm-settings.yml +++ b/dist-material/config-examples/alarm-settings.yml @@ -29,6 +29,8 @@ rules: # How many times of checks, the alarm keeps silence after alarm triggered, default as same as period. silence-period: 10 message: Successful rate of endpoint {name} is lower than 75% + tags: + level: WARNING service_resp_time_rule: metrics-name: service_resp_time # [Optional] Default, match all services in this metrics @@ -39,6 +41,8 @@ rules: op: ">" period: 10 count: 1 + tags: + level: CRITICAL service_instance_resp_time_rule: metrics-name: service_instance_resp_time op: ">" diff --git a/docs/en/setup/backend/backend-alarm.md b/docs/en/setup/backend/backend-alarm.md index e7ba3fba7..21535ac48 100644 --- a/docs/en/setup/backend/backend-alarm.md +++ b/docs/en/setup/backend/backend-alarm.md @@ -30,6 +30,7 @@ Alarm rule is constituted by following keys - **Exclude labels**. The following labels of the metric are excluded in this rule. - **Include labels regex**. Provide a regex to include labels. If both setting the include label list and include label regex, both rules will take effect. - **Exclude labels regex**. Provide a regex to exclude labels. If both setting the exclude label list and exclude label regex, both rules will take effect. +- **Tags**. Tags are key/value pairs that are attached to alarms. Tags are intended to be used to specify identifying attributes of alarms that are meaningful and relevant to users. *The settings of labels is required by meter-system which intends to store metrics from label-system platform, just like Prometheus, Micrometer, etc. The function supports the above four settings should implement `LabeledValueHolder`.* @@ -57,7 +58,7 @@ Composite rule is constituted by the following keys - **Rule name**. Unique name, show in alarm message. Must end with `_rule`. - **Expression**. Specify how to compose rules, support `&&`, `||`, `()`. - **Message**. Specify the notification message when rule triggered. - +- **Tags**. Tags are key/value pairs that are attached to alarms. Tags are intended to be used to specify identifying attributes of alarms that are meaningful and relevant to users. ```yaml rules: # Rule unique name, must be ended with `_rule`. @@ -74,6 +75,8 @@ rules: silence-period: 10 # Specify if the rule can send notification or just as an condition of composite rule only-as-condition: false + tags: + level: WARNING service_percent_rule: metrics-name: service_percent # [Optional] Default, match all services in this metrics @@ -115,6 +118,8 @@ composite-rules: # Must satisfied percent rule and resp time rule expression: service_percent_rule && service_resp_time_percentile_rule message: Service {name} successful rate is less than 80% and P50 of response time is over 1000ms + tags: + level: CRITICAL ``` @@ -143,7 +148,7 @@ Webhook requires the peer is a web container. The alarm message will send throug - **ruleName**. The rule name you configured in `alarm-settings.yml`. - **alarmMessage**. Alarm text message. - **startTime**. Alarm time measured in milliseconds, between the current time and midnight, January 1, 1970 UTC. - +- **tags**. The tags you configured in `alarm-settings.yml` Example as following ```json [{ @@ -154,7 +159,11 @@ Example as following "id1": "", "ruleName": "service_resp_time_rule", "alarmMessage": "alarmMessage xxxx", - "startTime": 1560524171000 + "startTime": 1560524171000, + "tags": [{ + "key": "level", + "value": "WARNING" + }] }, { "scopeId": 1, "scope": "SERVICE", @@ -163,7 +172,11 @@ Example as following "id1": "", "ruleName": "service_resp_time_rule", "alarmMessage": "alarmMessage yyy", - "startTime": 1560524171000 + "startTime": 1560524171000, + "tags": [{ + "key": "level", + "value": "CRITICAL" + }] }] ``` @@ -182,6 +195,17 @@ message AlarmMessage { string ruleName = 6; string alarmMessage = 7; int64 startTime = 8; + AlarmTags tags = 9; +} + +message AlarmTags { + // String key, String value pair. + repeated KeyStringValuePair data = 1; +} + +message KeyStringValuePair { + string key = 1; + string value = 2; } ``` diff --git a/docs/en/setup/backend/configuration-vocabulary.md b/docs/en/setup/backend/configuration-vocabulary.md index 139d4cafe..5d8746ec3 100644 --- a/docs/en/setup/backend/configuration-vocabulary.md +++ b/docs/en/setup/backend/configuration-vocabulary.md @@ -31,6 +31,7 @@ core|default|role|Option values, `Mixed/Receiver/Aggregator`. **Receiver** mode | - | - | endpointNameMaxLength| Max length limitation of endpoint name. The max length of service + endpoint names should be less than 240.|SW_ENDPOINT_NAME_MAX_LENGTH|150| | - | - | searchableTracesTags | Define the set of span tag keys, which should be searchable through the GraphQL. Multiple values should be separated through the comma. | SW_SEARCHABLE_TAG_KEYS | http.method,status_code,db.type,db.instance,mq.queue,mq.topic,mq.broker| | - | - | searchableLogsTags | Define the set of log tag keys, which should be searchable through the GraphQL. Multiple values should be separated through the comma. | SW_SEARCHABLE_LOGS_TAG_KEYS | level | +| - | - | searchableAlarmTags | Define the set of alarm tag keys, which should be searchable through the GraphQL. Multiple values should be separated through the comma. | SW_SEARCHABLE_ALARM_TAG_KEYS | level | | - | - | gRPCThreadPoolSize|Pool size of gRPC server| SW_CORE_GRPC_THREAD_POOL_SIZE | CPU core * 4| | - | - | gRPCThreadPoolQueueSize| The queue size of gRPC server| SW_CORE_GRPC_POOL_QUEUE_SIZE | 10000| | - | - | maxConcurrentCallsPerConnection | The maximum number of concurrent calls permitted for each incoming connection. Defaults to no limit. | SW_CORE_GRPC_MAX_CONCURRENT_CALL | - | diff --git a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/AlarmModuleProvider.java b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/AlarmModuleProvider.java index 88c3467c4..6a780a20d 100644 --- a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/AlarmModuleProvider.java +++ b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/AlarmModuleProvider.java @@ -67,7 +67,7 @@ public class AlarmModuleProvider extends ModuleProvider { alarmRulesWatcher = new AlarmRulesWatcher(rules, this); notifyHandler = new NotifyHandler(alarmRulesWatcher, getManager()); - notifyHandler.init(new AlarmStandardPersistence()); + notifyHandler.init(new AlarmStandardPersistence(getManager())); this.registerServiceImplementation(MetricsNotify.class, notifyHandler); } @@ -91,5 +91,4 @@ public class AlarmModuleProvider extends ModuleProvider { ConfigurationModule.NAME }; } - } diff --git a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/AlarmRule.java b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/AlarmRule.java index 31f38e716..86bced434 100644 --- a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/AlarmRule.java +++ b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/AlarmRule.java @@ -19,6 +19,7 @@ package org.apache.skywalking.oap.server.core.alarm.provider; import java.util.ArrayList; +import java.util.Map; import java.util.Objects; import lombok.AllArgsConstructor; import lombok.Builder; @@ -52,6 +53,7 @@ public class AlarmRule { private int silencePeriod; private String message; private boolean onlyAsCondition; + private Map tags; @Override public boolean equals(final Object o) { diff --git a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/CompositeAlarmRule.java b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/CompositeAlarmRule.java index 28aedc17c..116a19253 100644 --- a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/CompositeAlarmRule.java +++ b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/CompositeAlarmRule.java @@ -25,6 +25,8 @@ import lombok.NoArgsConstructor; import lombok.Setter; import lombok.ToString; +import java.util.Map; + @Builder @NoArgsConstructor @AllArgsConstructor @@ -35,4 +37,5 @@ public class CompositeAlarmRule { private String alarmRuleName; private String expression; private String message; + private Map tags; } diff --git a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/CompositeRuleEvaluator.java b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/CompositeRuleEvaluator.java index 9c2b49d8c..7fa761d7b 100644 --- a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/CompositeRuleEvaluator.java +++ b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/CompositeRuleEvaluator.java @@ -25,6 +25,7 @@ import org.apache.skywalking.oap.server.core.Const; import org.apache.skywalking.oap.server.core.alarm.AlarmMessage; import org.apache.skywalking.oap.server.core.alarm.MetaInAlarm; import org.apache.skywalking.oap.server.core.alarm.provider.expression.Expression; +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; import java.util.ArrayList; import java.util.HashMap; @@ -32,6 +33,7 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; +import java.util.stream.Collectors; /** * Evaluate composite rule using expression eval @@ -85,6 +87,7 @@ public class CompositeRuleEvaluator { String alarmMessage = formatMessage(message, compositeAlarmRule.getMessage(), compositeAlarmRule.getExpression()); message.setAlarmMessage(alarmMessage); message.setPeriod(headMsg.getPeriod()); + message.setTags(compositeAlarmRule.getTags().entrySet().stream().map(e -> new Tag(e.getKey(), e.getValue())).collect(Collectors.toList())); compositeRuleMessages.add(message); } }); diff --git a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/RulesReader.java b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/RulesReader.java index fea2e8c8a..2bf138a0a 100644 --- a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/RulesReader.java +++ b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/RulesReader.java @@ -24,6 +24,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.HashMap; import org.apache.skywalking.oap.server.core.alarm.provider.dingtalk.DingtalkSettings; import org.apache.skywalking.oap.server.core.alarm.provider.feishu.FeishuSettings; import org.apache.skywalking.oap.server.core.alarm.provider.grpc.GRPCAlarmSetting; @@ -107,7 +108,7 @@ public class RulesReader { alarmRule.setMessage( (String) settings.getOrDefault("message", "Alarm caused by Rule " + alarmRule .getAlarmRuleName())); - + alarmRule.setTags((Map) settings.getOrDefault("tags", new HashMap())); rules.getRules().add(alarmRule); } }); @@ -203,6 +204,7 @@ public class RulesReader { compositeAlarmRule.setExpression(expression); compositeAlarmRule.setMessage( (String) settings.getOrDefault("message", "Alarm caused by Rule " + ruleName)); + compositeAlarmRule.setTags((Map) settings.getOrDefault("tags", new HashMap(0))); rules.getCompositeRules().add(compositeAlarmRule); } }); diff --git a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/RunningRule.java b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/RunningRule.java index 348af0938..5ebd951e2 100644 --- a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/RunningRule.java +++ b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/RunningRule.java @@ -28,6 +28,7 @@ import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.locks.ReentrantLock; import java.util.regex.Pattern; +import java.util.stream.Collectors; import lombok.RequiredArgsConstructor; import lombok.ToString; @@ -35,6 +36,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.skywalking.apm.util.StringUtil; import org.apache.skywalking.oap.server.core.alarm.AlarmMessage; import org.apache.skywalking.oap.server.core.alarm.MetaInAlarm; +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; import org.apache.skywalking.oap.server.core.analysis.metrics.DataTable; import org.apache.skywalking.oap.server.core.analysis.metrics.DoubleValueHolder; import org.apache.skywalking.oap.server.core.analysis.metrics.IntValueHolder; @@ -74,6 +76,7 @@ public class RunningRule { private final Pattern excludeLabelsRegex; private final AlarmMessageFormatter formatter; private final boolean onlyAsCondition; + private final List tags; public RunningRule(AlarmRule alarmRule) { metricsName = alarmRule.getMetricsName(); @@ -104,6 +107,7 @@ public class RunningRule { Pattern.compile(alarmRule.getExcludeLabelsRegex()) : null; this.formatter = new AlarmMessageFormatter(alarmRule.getMessage()); this.onlyAsCondition = alarmRule.isOnlyAsCondition(); + this.tags = alarmRule.getTags().entrySet().stream().map(e -> new Tag(e.getKey(), e.getValue())).collect(Collectors.toList()); } /** @@ -237,6 +241,7 @@ public class RunningRule { alarmMessage.setOnlyAsCondition(this.onlyAsCondition); alarmMessage.setStartTime(System.currentTimeMillis()); alarmMessage.setPeriod(this.period); + alarmMessage.setTags(this.tags); alarmMessageList.add(alarmMessage); } }); diff --git a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/grpc/GRPCCallback.java b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/grpc/GRPCCallback.java index 0b99a47a9..4d109ee95 100644 --- a/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/grpc/GRPCCallback.java +++ b/oap-server/server-alarm-plugin/src/main/java/org/apache/skywalking/oap/server/core/alarm/provider/grpc/GRPCCallback.java @@ -25,6 +25,8 @@ import lombok.extern.slf4j.Slf4j; import org.apache.skywalking.oap.server.core.alarm.AlarmCallback; import org.apache.skywalking.oap.server.core.alarm.AlarmMessage; import org.apache.skywalking.oap.server.core.alarm.grpc.AlarmServiceGrpc; +import org.apache.skywalking.oap.server.core.alarm.grpc.AlarmTags; +import org.apache.skywalking.oap.server.core.alarm.grpc.KeyStringValuePair; import org.apache.skywalking.oap.server.core.alarm.grpc.Response; import org.apache.skywalking.oap.server.core.alarm.provider.AlarmRulesWatcher; import org.apache.skywalking.oap.server.library.client.grpc.GRPCClient; @@ -107,7 +109,9 @@ public class GRPCCallback implements AlarmCallback { builder.setRuleName(message.getRuleName()); builder.setAlarmMessage(message.getAlarmMessage()); builder.setStartTime(message.getStartTime()); - + AlarmTags.Builder alarmTagsBuilder = AlarmTags.newBuilder(); + message.getTags().forEach(m -> alarmTagsBuilder.addData(KeyStringValuePair.newBuilder().setKey(m.getKey()).setValue(m.getValue()).build())); + builder.setTags(alarmTagsBuilder.build()); streamObserver.onNext(builder.build()); }); diff --git a/oap-server/server-alarm-plugin/src/main/proto/alarm-hook.proto b/oap-server/server-alarm-plugin/src/main/proto/alarm-hook.proto index 1f3a34fd5..75b9c31ec 100644 --- a/oap-server/server-alarm-plugin/src/main/proto/alarm-hook.proto +++ b/oap-server/server-alarm-plugin/src/main/proto/alarm-hook.proto @@ -35,6 +35,17 @@ message AlarmMessage { string ruleName = 6; string alarmMessage = 7; int64 startTime = 8; + AlarmTags tags = 9; +} + +message AlarmTags { + // String key, String value pair. + repeated KeyStringValuePair data = 1; +} + +message KeyStringValuePair { + string key = 1; + string value = 2; } message Response { diff --git a/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/AlarmRulesWatcherTest.java b/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/AlarmRulesWatcherTest.java index e9277a6d8..af5f29beb 100644 --- a/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/AlarmRulesWatcherTest.java +++ b/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/AlarmRulesWatcherTest.java @@ -22,6 +22,8 @@ import java.io.IOException; import java.io.Reader; import java.util.ArrayList; import java.util.Arrays; +import java.util.HashMap; + import org.apache.skywalking.oap.server.configuration.api.ConfigChangeWatcher; import org.apache.skywalking.oap.server.library.util.ResourceUtils; import org.junit.Before; @@ -60,6 +62,9 @@ public class AlarmRulesWatcherTest { .op(">") .period(1) .silencePeriod(2) + .tags(new HashMap() {{ + put("key", "value"); + }}) .threshold("2"); @Before diff --git a/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/CompositeRuleEvaluatorTest.java b/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/CompositeRuleEvaluatorTest.java index 500873160..2da3e98be 100644 --- a/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/CompositeRuleEvaluatorTest.java +++ b/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/CompositeRuleEvaluatorTest.java @@ -25,6 +25,7 @@ import org.junit.Before; import org.junit.Test; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; import static org.hamcrest.CoreMatchers.is; @@ -43,7 +44,9 @@ public class CompositeRuleEvaluatorTest { @Test public void testEvaluateMessageWithAndOp() { List compositeAlarmRules = new ArrayList<>(); - CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "a_rule && b_rule", "composite rule {name},{id} triggered!"); + CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "a_rule && b_rule", "composite rule {name},{id} triggered!", new HashMap() {{ + put("key", "value"); + }}); compositeAlarmRules.add(compositeAlarmRule); List alarmMessages = getAlarmMessages(); List compositeMsgs = ruleEvaluate.evaluate(compositeAlarmRules, alarmMessages); @@ -52,13 +55,17 @@ public class CompositeRuleEvaluatorTest { assertThat(compositeMsgs.get(0).getRuleName(), is("dummy")); assertThat(compositeMsgs.get(0).getId0(), is("id0")); assertThat(compositeMsgs.get(0).getId1(), is("id1")); + assertThat(compositeMsgs.get(0).getTags().get(0).getKey(), is("key")); + assertThat(compositeMsgs.get(0).getTags().get(0).getValue(), is("value")); assertThat(compositeMsgs.get(0).isOnlyAsCondition(), is(false)); } @Test public void testEvaluateMessageWithFormatMessage() { List compositeAlarmRules = new ArrayList<>(); - CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "a_rule && b_rule", "composite rule {name} triggered!"); + CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "a_rule && b_rule", "composite rule {name} triggered!", new HashMap() {{ + put("key", "value"); + }}); compositeAlarmRules.add(compositeAlarmRule); List alarmMessages = getAlarmMessages(); List compositeMsgs = ruleEvaluate.evaluate(compositeAlarmRules, alarmMessages); @@ -67,13 +74,15 @@ public class CompositeRuleEvaluatorTest { assertThat(compositeMsgs.get(0).getRuleName(), is("dummy")); assertThat(compositeMsgs.get(0).getId0(), is("id0")); assertThat(compositeMsgs.get(0).getId1(), is("id1")); + assertThat(compositeMsgs.get(0).getTags().get(0).getKey(), is("key")); + assertThat(compositeMsgs.get(0).getTags().get(0).getValue(), is("value")); assertThat(compositeMsgs.get(0).isOnlyAsCondition(), is(false)); } @Test public void testEvaluateMessageWithNotExistsRule() { List compositeAlarmRules = new ArrayList<>(); - CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "a_rule && not_exist_rule", "composite rule triggered!"); + CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "a_rule && not_exist_rule", "composite rule triggered!", new HashMap<>()); compositeAlarmRules.add(compositeAlarmRule); List alarmMessages = getAlarmMessages(); List compositeMsgs = ruleEvaluate.evaluate(compositeAlarmRules, alarmMessages); @@ -83,7 +92,7 @@ public class CompositeRuleEvaluatorTest { @Test public void testEvaluateMessageWithException() { List compositeAlarmRules = new ArrayList<>(); - CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "a_rule + b_rule", "composite rule triggered!"); + CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "a_rule + b_rule", "composite rule triggered!", new HashMap<>()); compositeAlarmRules.add(compositeAlarmRule); List alarmMessages = getAlarmMessages(); List compositeMsgs = ruleEvaluate.evaluate(compositeAlarmRules, alarmMessages); @@ -125,7 +134,9 @@ public class CompositeRuleEvaluatorTest { @Test public void testEvaluateMessageWithOrOp() { List compositeAlarmRules = new ArrayList<>(); - CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "a_rule || b_rule", "composite rule triggered!"); + CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "a_rule || b_rule", "composite rule triggered!", new HashMap() {{ + put("key", "value"); + }}); compositeAlarmRules.add(compositeAlarmRule); List alarmMessages = getAlarmMessages(); alarmMessages.remove(0); @@ -135,13 +146,15 @@ public class CompositeRuleEvaluatorTest { assertThat(compositeMsgs.get(0).getRuleName(), is("dummy")); assertThat(compositeMsgs.get(0).getId0(), is("id0")); assertThat(compositeMsgs.get(0).getId1(), is("id1")); + assertThat(compositeMsgs.get(0).getTags().get(0).getKey(), is("key")); + assertThat(compositeMsgs.get(0).getTags().get(0).getValue(), is("value")); assertThat(compositeMsgs.get(0).isOnlyAsCondition(), is(false)); } @Test public void testEvaluateMessageWithParenthesisAndOp() { List compositeAlarmRules = new ArrayList<>(); - CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "(a_rule || b_rule) && c_rule", "composite rule triggered!"); + CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "(a_rule || b_rule) && c_rule", "composite rule triggered!", new HashMap<>()); compositeAlarmRules.add(compositeAlarmRule); List alarmMessages = getAlarmMessages(); alarmMessages.remove(alarmMessages.size() - 1); @@ -152,7 +165,7 @@ public class CompositeRuleEvaluatorTest { @Test public void testEvaluateMessageWithParenthesisAndOrOp() { List compositeAlarmRules = new ArrayList<>(); - CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "(a_rule && b_rule) || c_rule", "composite rule triggered!"); + CompositeAlarmRule compositeAlarmRule = new CompositeAlarmRule("dummy", "(a_rule && b_rule) || c_rule", "composite rule triggered!", new HashMap<>()); compositeAlarmRules.add(compositeAlarmRule); List alarmMessages = getAlarmMessages(); List compositeMsgs = ruleEvaluate.evaluate(compositeAlarmRules, alarmMessages); diff --git a/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/RunningRuleTest.java b/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/RunningRuleTest.java index fc5368d63..a7229ec9a 100644 --- a/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/RunningRuleTest.java +++ b/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/RunningRuleTest.java @@ -23,6 +23,7 @@ import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.HashMap; import lombok.Getter; import lombok.Setter; import org.apache.skywalking.oap.server.core.Const; @@ -60,7 +61,9 @@ public class RunningRuleTest { alarmRule.setThreshold("75"); alarmRule.setCount(3); alarmRule.setPeriod(15); - + alarmRule.setTags(new HashMap() {{ + put("key", "value"); + }}); RunningRule runningRule = new RunningRule(alarmRule); LocalDateTime startTime = TIME_BUCKET_FORMATTER.parseLocalDateTime("201808301434"); long timeInPeriod1 = 201808301434L; @@ -88,7 +91,9 @@ public class RunningRuleTest { alarmRule.setCount(3); alarmRule.setPeriod(15); alarmRule.setMessage("Successful rate of endpoint {name} is lower than 75%"); - + alarmRule.setTags(new HashMap() {{ + put("key", "value"); + }}); RunningRule runningRule = new RunningRule(alarmRule); LocalDateTime startTime = TIME_BUCKET_FORMATTER.parseLocalDateTime("201808301440"); @@ -125,7 +130,9 @@ public class RunningRuleTest { alarmRule.setCount(3); alarmRule.setPeriod(15); alarmRule.setMessage("response percentile of endpoint {name} is lower than expected values"); - + alarmRule.setTags(new HashMap() {{ + put("key", "value"); + }}); RunningRule runningRule = new RunningRule(alarmRule); LocalDateTime startTime = TIME_BUCKET_FORMATTER.parseLocalDateTime("201808301440"); @@ -179,7 +186,9 @@ public class RunningRuleTest { alarmRule.setCount(3); alarmRule.setPeriod(15); //alarmRule.setSilencePeriod(0); - + alarmRule.setTags(new HashMap() {{ + put("key", "value"); + }}); RunningRule runningRule = new RunningRule(alarmRule); LocalDateTime startTime = TIME_BUCKET_FORMATTER.parseLocalDateTime("201808301441"); @@ -224,7 +233,9 @@ public class RunningRuleTest { alarmRule.setCount(3); alarmRule.setPeriod(15); alarmRule.setSilencePeriod(2); - + alarmRule.setTags(new HashMap() {{ + put("key", "value"); + }}); RunningRule runningRule = new RunningRule(alarmRule); long timeInPeriod1 = 201808301434L; @@ -261,7 +272,9 @@ public class RunningRuleTest { alarmRule.setPeriod(15); alarmRule.setMessage("Successful rate of endpoint {name} is lower than 75%"); alarmRule.setExcludeNames(Lists.newArrayList("Service_123")); - + alarmRule.setTags(new HashMap() {{ + put("key", "value"); + }}); RunningRule runningRule = new RunningRule(alarmRule); long timeInPeriod1 = 201808301434L; @@ -293,7 +306,9 @@ public class RunningRuleTest { alarmRule.setPeriod(10); alarmRule.setMessage("Response time of service instance {name} is more than 1000ms in 2 minutes of last 10 minutes"); alarmRule.setIncludeNamesRegex("Service\\_1(\\d)+"); - + alarmRule.setTags(new HashMap() {{ + put("key", "value"); + }}); RunningRule runningRule = new RunningRule(alarmRule); long timeInPeriod1 = 201808301434L; @@ -325,7 +340,9 @@ public class RunningRuleTest { alarmRule.setPeriod(10); alarmRule.setMessage("Response time of service instance {name} is more than 1000ms in 2 minutes of last 10 minutes"); alarmRule.setExcludeNamesRegex("Service\\_2(\\d)+"); - + alarmRule.setTags(new HashMap() {{ + put("key", "value"); + }}); RunningRule runningRule = new RunningRule(alarmRule); long timeInPeriod1 = 201808301434L; @@ -574,7 +591,9 @@ public class RunningRuleTest { alarmRule.setCount(3); alarmRule.setPeriod(15); alarmRule.setMessage("response percentile of endpoint {name} is lower than expected value"); - + alarmRule.setTags(new HashMap() {{ + put("key", "value"); + }}); RunningRule runningRule = new RunningRule(alarmRule); LocalDateTime startTime = TIME_BUCKET_FORMATTER.parseLocalDateTime("201808301440"); diff --git a/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/grpc/GRPChookCallbackTest.java b/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/grpc/GRPChookCallbackTest.java index 5607dc416..012d0a565 100644 --- a/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/grpc/GRPChookCallbackTest.java +++ b/oap-server/server-alarm-plugin/src/test/java/org/apache/skywalking/oap/server/core/alarm/provider/grpc/GRPChookCallbackTest.java @@ -19,10 +19,13 @@ package org.apache.skywalking.oap.server.core.alarm.provider.grpc; import com.google.common.collect.Lists; + +import java.util.Arrays; import java.util.List; import org.apache.skywalking.oap.server.core.alarm.AlarmMessage; import org.apache.skywalking.oap.server.core.alarm.provider.AlarmRulesWatcher; import org.apache.skywalking.oap.server.core.alarm.provider.Rules; +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; import org.apache.skywalking.oap.server.core.query.enumeration.Scope; import org.junit.Before; import org.junit.Test; @@ -73,7 +76,7 @@ public class GRPChookCallbackTest { alarmMessage.setAlarmMessage("message"); alarmMessage.setRuleName("mock_rule"); alarmMessage.setStartTime(System.currentTimeMillis()); - + alarmMessage.setTags(Arrays.asList(new Tag("key", "value"))); alarmMessageList = Lists.newArrayList(alarmMessage); } } diff --git a/oap-server/server-bootstrap/src/main/resources/alarm-settings.yml b/oap-server/server-bootstrap/src/main/resources/alarm-settings.yml index b7fa97952..87ad2f53e 100755 --- a/oap-server/server-bootstrap/src/main/resources/alarm-settings.yml +++ b/oap-server/server-bootstrap/src/main/resources/alarm-settings.yml @@ -38,6 +38,8 @@ rules: op: ">" period: 10 count: 1 + tags: + level: WARNING webhooks: # - http://127.0.0.1/notify/ diff --git a/oap-server/server-bootstrap/src/main/resources/application.yml b/oap-server/server-bootstrap/src/main/resources/application.yml index c9bab4639..239858b00 100755 --- a/oap-server/server-bootstrap/src/main/resources/application.yml +++ b/oap-server/server-bootstrap/src/main/resources/application.yml @@ -102,6 +102,8 @@ core: searchableTracesTags: ${SW_SEARCHABLE_TAG_KEYS:http.method,status_code,db.type,db.instance,mq.queue,mq.topic,mq.broker} # Define the set of log tag keys, which should be searchable through the GraphQL. searchableLogsTags: ${SW_SEARCHABLE_LOGS_TAG_KEYS:level} + # Define the set of alarm tag keys, which should be searchable through the GraphQL. + searchableAlarmTags: ${SW_SEARCHABLE_ALARM_TAG_KEYS:level} # The number of threads used to synchronously refresh the metrics data to the storage. syncThreads: ${SW_CORE_SYNC_THREADS:2} # The maximum number of processes supported for each synchronous storage operation. When the number of the flush data is greater than this value, it will be assigned to multiple cores for execution. diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleConfig.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleConfig.java index 2a55fa458..22f75a9df 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleConfig.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleConfig.java @@ -133,7 +133,14 @@ public class CoreModuleConfig extends ModuleConfig { @Setter @Getter private String searchableLogsTags = ""; - + /** + * Define the set of Alarm tag keys, which should be searchable through the GraphQL. + * + * @since 8.6.0 + */ + @Setter + @Getter + private String searchableAlarmTags = ""; /** * The number of threads used to synchronously refresh the metrics data to the storage. * diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmMessage.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmMessage.java index fdf3c1467..c778c988c 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmMessage.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmMessage.java @@ -20,6 +20,9 @@ package org.apache.skywalking.oap.server.core.alarm; import lombok.Getter; import lombok.Setter; +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; + +import java.util.List; /** * Alarm message represents the details of each alarm. @@ -34,6 +37,7 @@ public class AlarmMessage { private String id1; private String ruleName; private String alarmMessage; + private List tags; private long startTime; private transient int period; private transient boolean onlyAsCondition; diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmRecord.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmRecord.java index 12ddfbce9..8fc26a369 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmRecord.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmRecord.java @@ -20,16 +20,21 @@ package org.apache.skywalking.oap.server.core.alarm; import lombok.Getter; import lombok.Setter; +import org.apache.skywalking.apm.util.StringUtil; import org.apache.skywalking.oap.server.core.Const; import org.apache.skywalking.oap.server.core.analysis.Stream; +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; import org.apache.skywalking.oap.server.core.analysis.record.Record; import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor; import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine; import org.apache.skywalking.oap.server.core.source.ScopeDeclaration; import org.apache.skywalking.oap.server.core.storage.StorageHashMapBuilder; import org.apache.skywalking.oap.server.core.storage.annotation.Column; +import org.apache.skywalking.oap.server.library.util.CollectionUtils; +import java.util.Base64; import java.util.HashMap; +import java.util.List; import java.util.Map; import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.ALARM; @@ -48,6 +53,8 @@ public class AlarmRecord extends Record { public static final String START_TIME = "start_time"; public static final String ALARM_MESSAGE = "alarm_message"; public static final String RULE_NAME = "rule_name"; + public static final String TAGS = "tags"; + public static final String TAGS_RAW_DATA = "tags_raw_data"; @Override public String id() { @@ -68,6 +75,13 @@ public class AlarmRecord extends Record { private String alarmMessage; @Column(columnName = RULE_NAME) private String ruleName; + @Column(columnName = TAGS) + private List tagsInString; + @Column(columnName = TAGS_RAW_DATA) + private byte[] tagsRawData; + @Setter + @Getter + private List tags; public static class Builder implements StorageHashMapBuilder { @@ -82,6 +96,12 @@ public class AlarmRecord extends Record { map.put(START_TIME, storageData.getStartTime()); map.put(TIME_BUCKET, storageData.getTimeBucket()); map.put(RULE_NAME, storageData.getRuleName()); + if (CollectionUtils.isEmpty(storageData.getTagsRawData())) { + map.put(TAGS_RAW_DATA, Const.EMPTY_STRING); + } else { + map.put(TAGS_RAW_DATA, new String(Base64.getEncoder().encode(storageData.getTagsRawData()))); + } + map.put(TAGS, storageData.getTagsInString()); return map; } @@ -96,6 +116,12 @@ public class AlarmRecord extends Record { record.setStartTime(((Number) dbMap.get(START_TIME)).longValue()); record.setTimeBucket(((Number) dbMap.get(TIME_BUCKET)).longValue()); record.setRuleName((String) dbMap.get(RULE_NAME)); + if (StringUtil.isEmpty((String) dbMap.get(TAGS_RAW_DATA))) { + record.setTagsRawData(new byte[] {}); + } else { + // Don't read the tags as they has been in the data binary already. + record.setTagsRawData(Base64.getDecoder().decode((String) dbMap.get(TAGS_RAW_DATA))); + } return record; } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmStandardPersistence.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmStandardPersistence.java index eb0c8edeb..0764a2040 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmStandardPersistence.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmStandardPersistence.java @@ -18,11 +18,20 @@ package org.apache.skywalking.oap.server.core.alarm; +import com.google.common.base.Charsets; +import com.google.gson.Gson; +import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.analysis.TimeBucket; +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor; +import org.apache.skywalking.oap.server.core.config.ConfigService; +import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.ArrayList; +import java.util.Collection; +import java.util.HashSet; import java.util.List; /** @@ -31,6 +40,12 @@ import java.util.List; public class AlarmStandardPersistence implements AlarmCallback { private static final Logger LOGGER = LoggerFactory.getLogger(AlarmStandardPersistence.class); + private final Gson gson = new Gson(); + private final ModuleManager manager; + + public AlarmStandardPersistence(ModuleManager manager) { + this.manager = manager; + } @Override public void doAlarm(List alarmMessage) { @@ -48,8 +63,25 @@ public class AlarmStandardPersistence implements AlarmCallback { record.setStartTime(message.getStartTime()); record.setTimeBucket(TimeBucket.getRecordTimeBucket(message.getStartTime())); record.setRuleName(message.getRuleName()); - + Collection tags = appendSearchableTags(message.getTags()); + record.setTags(new ArrayList<>(tags)); + record.setTagsRawData(gson.toJson(message.getTags()).getBytes(Charsets.UTF_8)); + record.setTagsInString(Tag.Util.toStringList(new ArrayList<>(tags))); RecordStreamProcessor.getInstance().in(record); }); } + + private Collection appendSearchableTags(List tags) { + final ConfigService configService = manager.find(CoreModule.NAME) + .provider() + .getService(ConfigService.class); + HashSet alarmTags = new HashSet<>(); + tags.forEach(tag -> { + if (configService.getSearchableAlarmTags().contains(tag.getKey())) { + final Tag alarmTag = new Tag(tag.getKey(), tag.getValue()); + alarmTags.add(alarmTag); + } + }); + return alarmTags; + } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/config/ConfigService.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/config/ConfigService.java index dbf5a396b..e6bb868f3 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/config/ConfigService.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/config/ConfigService.java @@ -28,11 +28,13 @@ public class ConfigService implements Service { private final int gRPCPort; private final String searchableTracesTags; private final String searchableLogsTags; + private final String searchableAlarmTags; public ConfigService(CoreModuleConfig moduleConfig) { this.gRPCHost = moduleConfig.getGRPCHost(); this.gRPCPort = moduleConfig.getGRPCPort(); this.searchableTracesTags = moduleConfig.getSearchableTracesTags(); this.searchableLogsTags = moduleConfig.getSearchableLogsTags(); + this.searchableAlarmTags = moduleConfig.getSearchableAlarmTags(); } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/AlarmQueryService.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/AlarmQueryService.java index 2b8bdb9c3..fb06c1907 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/AlarmQueryService.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/AlarmQueryService.java @@ -19,6 +19,9 @@ package org.apache.skywalking.oap.server.core.query; import java.io.IOException; +import java.util.List; + +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; import org.apache.skywalking.oap.server.core.query.type.Alarms; import org.apache.skywalking.oap.server.core.query.type.Pagination; import org.apache.skywalking.oap.server.core.storage.StorageModule; @@ -43,8 +46,8 @@ public class AlarmQueryService implements Service { } public Alarms getAlarm(final Integer scopeId, final String keyword, final Pagination paging, final long startTB, - final long endTB) throws IOException { + final long endTB, final List tags) throws IOException { PaginationUtils.Page page = PaginationUtils.INSTANCE.exchange(paging); - return getAlarmQueryDAO().getAlarm(scopeId, keyword, page.getLimit(), page.getFrom(), startTB, endTB); + return getAlarmQueryDAO().getAlarm(scopeId, keyword, page.getLimit(), page.getFrom(), startTB, endTB, tags); } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/type/AlarmMessage.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/type/AlarmMessage.java index a8f8b23d1..328a98514 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/type/AlarmMessage.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/type/AlarmMessage.java @@ -22,6 +22,9 @@ import lombok.Getter; import lombok.Setter; import org.apache.skywalking.oap.server.core.query.enumeration.Scope; +import java.util.ArrayList; +import java.util.List; + @Getter @Setter public class AlarmMessage { @@ -30,4 +33,9 @@ public class AlarmMessage { private String id; private String message; private Long startTime; + private final List tags; + + public AlarmMessage() { + tags = new ArrayList<>(); + } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/query/IAlarmQueryDAO.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/query/IAlarmQueryDAO.java index 3278e1dee..65f91b376 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/query/IAlarmQueryDAO.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/query/IAlarmQueryDAO.java @@ -19,11 +19,37 @@ package org.apache.skywalking.oap.server.core.storage.query; import java.io.IOException; +import java.util.Base64; +import java.util.List; + +import com.google.common.base.Charsets; +import com.google.gson.Gson; +import com.google.gson.reflect.TypeToken; +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; import org.apache.skywalking.oap.server.core.query.type.Alarms; +import org.apache.skywalking.oap.server.core.query.type.KeyValue; import org.apache.skywalking.oap.server.core.storage.DAO; public interface IAlarmQueryDAO extends DAO { + Gson GSON = new Gson(); + Alarms getAlarm(final Integer scopeId, final String keyword, final int limit, final int from, final long startTB, - final long endTB) throws IOException; + final long endTB, final List tags) throws IOException; + + /** + * Parser the raw tags. + */ + default void parserDataBinaryBase64(String dataBinaryBase64, List tags) { + parserDataBinary(Base64.getDecoder().decode(dataBinaryBase64), tags); + } + + /** + * Parser the raw tags. + */ + default void parserDataBinary(byte[] dataBinary, List tags) { + List tagList = GSON.fromJson(new String(dataBinary, Charsets.UTF_8), new TypeToken>() { + }.getType()); + tagList.forEach(pair -> tags.add(new KeyValue(pair.getKey(), pair.getValue()))); + } } diff --git a/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/AlarmQuery.java b/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/AlarmQuery.java index d64f3ed48..9bac9cb56 100644 --- a/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/AlarmQuery.java +++ b/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/AlarmQuery.java @@ -20,7 +20,10 @@ package org.apache.skywalking.oap.query.graphql.resolver; import com.coxautodev.graphql.tools.GraphQLQueryResolver; import java.io.IOException; +import java.util.List; + import org.apache.skywalking.oap.server.core.CoreModule; +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; import org.apache.skywalking.oap.server.core.query.AlarmQueryService; import org.apache.skywalking.oap.server.core.query.enumeration.Scope; import org.apache.skywalking.oap.server.core.query.input.Duration; @@ -29,6 +32,8 @@ import org.apache.skywalking.oap.server.core.query.type.Alarms; import org.apache.skywalking.oap.server.core.query.type.Pagination; import org.apache.skywalking.oap.server.library.module.ModuleManager; +import static java.util.Objects.nonNull; + public class AlarmQuery implements GraphQLQueryResolver { private final ModuleManager moduleManager; @@ -50,13 +55,18 @@ public class AlarmQuery implements GraphQLQueryResolver { } public Alarms getAlarm(final Duration duration, final Scope scope, final String keyword, - final Pagination paging) throws IOException { + final Pagination paging, final List tags) throws IOException { Integer scopeId = null; if (scope != null) { scopeId = scope.getScopeId(); } - + long startSecondTB = 0; + long endSecondTB = 0; + if (nonNull(duration)) { + startSecondTB = duration.getStartTimeBucketInSec(); + endSecondTB = duration.getEndTimeBucketInSec(); + } return getQueryService().getAlarm( - scopeId, keyword, paging, duration.getStartTimeBucketInSec(), duration.getEndTimeBucketInSec()); + scopeId, keyword, paging, startSecondTB, endSecondTB, tags); } } diff --git a/oap-server/server-query-plugin/query-graphql-plugin/src/main/resources/query-protocol b/oap-server/server-query-plugin/query-graphql-plugin/src/main/resources/query-protocol index a53e89490..0c2388ba1 160000 --- a/oap-server/server-query-plugin/query-graphql-plugin/src/main/resources/query-protocol +++ b/oap-server/server-query-plugin/query-graphql-plugin/src/main/resources/query-protocol @@ -1 +1 @@ -Subproject commit a53e89490cb9ddf064685c2ec2b2ae9329d26ca5 +Subproject commit 0c2388ba18cfc1b1b103ddad71f9765bd21dff6e diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/AlarmQueryEsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/AlarmQueryEsDAO.java index 1c4636561..e224ff74c 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/AlarmQueryEsDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/AlarmQueryEsDAO.java @@ -19,19 +19,25 @@ package org.apache.skywalking.oap.server.storage.plugin.elasticsearch.query; import com.google.common.base.Strings; + import java.io.IOException; +import java.util.List; import java.util.Objects; + import org.apache.skywalking.oap.server.core.alarm.AlarmRecord; +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; import org.apache.skywalking.oap.server.core.query.enumeration.Scope; import org.apache.skywalking.oap.server.core.query.type.AlarmMessage; import org.apache.skywalking.oap.server.core.query.type.Alarms; import org.apache.skywalking.oap.server.core.storage.query.IAlarmQueryDAO; import org.apache.skywalking.oap.server.library.client.elasticsearch.ElasticSearchClient; +import org.apache.skywalking.oap.server.library.util.CollectionUtils; import org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base.EsDAO; import org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base.IndexController; import org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base.MatchCNameBuilder; import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.index.query.BoolQueryBuilder; +import org.elasticsearch.index.query.QueryBuilder; import org.elasticsearch.index.query.QueryBuilders; import org.elasticsearch.search.SearchHit; import org.elasticsearch.search.builder.SearchSourceBuilder; @@ -45,19 +51,27 @@ public class AlarmQueryEsDAO extends EsDAO implements IAlarmQueryDAO { @Override public Alarms getAlarm(final Integer scopeId, final String keyword, final int limit, final int from, - final long startTB, final long endTB) throws IOException { + final long startTB, final long endTB, final List tags) throws IOException { SearchSourceBuilder sourceBuilder = SearchSourceBuilder.searchSource(); BoolQueryBuilder boolQueryBuilder = QueryBuilders.boolQuery(); - boolQueryBuilder.must().add(QueryBuilders.rangeQuery(AlarmRecord.TIME_BUCKET).gte(startTB).lte(endTB)); + List mustQueryList = boolQueryBuilder.must(); + + if (startTB != 0 && endTB != 0) { + mustQueryList.add(QueryBuilders.rangeQuery(AlarmRecord.TIME_BUCKET).gte(startTB).lte(endTB)); + } if (Objects.nonNull(scopeId)) { - boolQueryBuilder.must().add(QueryBuilders.termQuery(AlarmRecord.SCOPE, scopeId.intValue())); + mustQueryList.add(QueryBuilders.termQuery(AlarmRecord.SCOPE, scopeId.intValue())); } if (!Strings.isNullOrEmpty(keyword)) { String matchCName = MatchCNameBuilder.INSTANCE.build(AlarmRecord.ALARM_MESSAGE); - boolQueryBuilder.must().add(QueryBuilders.matchPhraseQuery(matchCName, keyword)); + mustQueryList.add(QueryBuilders.matchPhraseQuery(matchCName, keyword)); + } + + if (CollectionUtils.isNotEmpty(tags)) { + tags.forEach(tag -> mustQueryList.add(QueryBuilders.termQuery(AlarmRecord.TAGS, tag.toString()))); } sourceBuilder.query(boolQueryBuilder).sort(AlarmRecord.START_TIME, SortOrder.DESC); @@ -65,7 +79,7 @@ public class AlarmQueryEsDAO extends EsDAO implements IAlarmQueryDAO { sourceBuilder.from(from); SearchResponse response = getClient() - .search(IndexController.LogicIndicesRegister.getPhysicalTableName(AlarmRecord.INDEX_NAME), sourceBuilder); + .search(IndexController.LogicIndicesRegister.getPhysicalTableName(AlarmRecord.INDEX_NAME), sourceBuilder); Alarms alarms = new Alarms(); alarms.setTotal((int) response.getHits().totalHits); @@ -80,6 +94,9 @@ public class AlarmQueryEsDAO extends EsDAO implements IAlarmQueryDAO { message.setStartTime(alarmRecord.getStartTime()); message.setScope(Scope.Finder.valueOf(alarmRecord.getScope())); message.setScopeId(alarmRecord.getScope()); + if (!CollectionUtils.isEmpty(alarmRecord.getTagsRawData())) { + parserDataBinary(alarmRecord.getTagsRawData(), message.getTags()); + } alarms.getMsgs().add(message); } return alarms; diff --git a/oap-server/server-storage-plugin/storage-elasticsearch7-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch7/query/AlarmQueryEs7DAO.java b/oap-server/server-storage-plugin/storage-elasticsearch7-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch7/query/AlarmQueryEs7DAO.java index ebcf1398e..91dd56b1c 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch7-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch7/query/AlarmQueryEs7DAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch7-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch7/query/AlarmQueryEs7DAO.java @@ -20,18 +20,23 @@ package org.apache.skywalking.oap.server.storage.plugin.elasticsearch7.query; import com.google.common.base.Strings; import java.io.IOException; +import java.util.List; import java.util.Objects; import org.apache.skywalking.oap.server.core.alarm.AlarmRecord; +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; +import org.apache.skywalking.oap.server.core.analysis.record.Record; import org.apache.skywalking.oap.server.core.query.enumeration.Scope; import org.apache.skywalking.oap.server.core.query.type.AlarmMessage; import org.apache.skywalking.oap.server.core.query.type.Alarms; import org.apache.skywalking.oap.server.core.storage.query.IAlarmQueryDAO; import org.apache.skywalking.oap.server.library.client.elasticsearch.ElasticSearchClient; +import org.apache.skywalking.oap.server.library.util.CollectionUtils; import org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base.EsDAO; import org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base.IndexController; import org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base.MatchCNameBuilder; import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.index.query.BoolQueryBuilder; +import org.elasticsearch.index.query.QueryBuilder; import org.elasticsearch.index.query.QueryBuilders; import org.elasticsearch.search.SearchHit; import org.elasticsearch.search.builder.SearchSourceBuilder; @@ -46,10 +51,16 @@ public class AlarmQueryEs7DAO extends EsDAO implements IAlarmQueryDAO { @Override public Alarms getAlarm(final Integer scopeId, final String keyword, final int limit, final int from, - final long startTB, final long endTB) throws IOException { + final long startTB, final long endTB, final List tags) throws IOException { SearchSourceBuilder sourceBuilder = SearchSourceBuilder.searchSource(); BoolQueryBuilder boolQueryBuilder = QueryBuilders.boolQuery(); + List mustQueryList = boolQueryBuilder.must(); + + if (startTB != 0 && endTB != 0) { + mustQueryList.add(QueryBuilders.rangeQuery(Record.TIME_BUCKET).gte(startTB).lte(endTB)); + } + boolQueryBuilder.must().add(QueryBuilders.rangeQuery(AlarmRecord.TIME_BUCKET).gte(startTB).lte(endTB)); if (Objects.nonNull(scopeId)) { @@ -61,6 +72,10 @@ public class AlarmQueryEs7DAO extends EsDAO implements IAlarmQueryDAO { boolQueryBuilder.must().add(QueryBuilders.matchPhraseQuery(matchCName, keyword)); } + if (CollectionUtils.isNotEmpty(tags)) { + tags.forEach(tag -> mustQueryList.add(QueryBuilders.termQuery(AlarmRecord.TAGS, tag.toString()))); + } + sourceBuilder.query(boolQueryBuilder).sort(AlarmRecord.START_TIME, SortOrder.DESC); sourceBuilder.size(limit); sourceBuilder.from(from); @@ -81,6 +96,9 @@ public class AlarmQueryEs7DAO extends EsDAO implements IAlarmQueryDAO { message.setStartTime(alarmRecord.getStartTime()); message.setScope(Scope.Finder.valueOf(alarmRecord.getScope())); message.setScopeId(alarmRecord.getScope()); + if (!CollectionUtils.isEmpty(alarmRecord.getTagsRawData())) { + parserDataBinary(alarmRecord.getTagsRawData(), message.getTags()); + } alarms.getMsgs().add(message); } return alarms; diff --git a/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/base/InfluxInsertRequest.java b/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/base/InfluxInsertRequest.java index 45c877f62..ffc09a16a 100644 --- a/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/base/InfluxInsertRequest.java +++ b/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/base/InfluxInsertRequest.java @@ -21,6 +21,7 @@ package org.apache.skywalking.oap.server.storage.plugin.influxdb.base; import com.google.common.collect.Maps; import java.util.Map; import java.util.concurrent.TimeUnit; +import org.apache.skywalking.oap.server.core.alarm.AlarmRecord; import org.apache.skywalking.oap.server.core.analysis.manual.log.LogRecord; import org.apache.skywalking.oap.server.core.analysis.manual.segment.SegmentRecord; import org.apache.skywalking.oap.server.core.storage.StorageHashMapBuilder; @@ -42,7 +43,7 @@ public class InfluxInsertRequest implements InsertRequest, UpdateRequest { public InfluxInsertRequest(Model model, T storageData, StorageHashMapBuilder storageBuilder) { final Map objectMap = storageBuilder.entity2Storage(storageData); - if (SegmentRecord.INDEX_NAME.equals(model.getName()) || LogRecord.INDEX_NAME.equals(model.getName())) { + if (SegmentRecord.INDEX_NAME.equals(model.getName()) || LogRecord.INDEX_NAME.equals(model.getName()) || AlarmRecord.INDEX_NAME.equals(model.getName())) { objectMap.remove(SegmentRecord.TAGS); } diff --git a/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/base/RecordDAO.java b/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/base/RecordDAO.java index 28e7c5b46..5cedff65f 100644 --- a/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/base/RecordDAO.java +++ b/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/base/RecordDAO.java @@ -24,6 +24,7 @@ import java.util.Map; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import org.apache.skywalking.apm.commons.datacarrier.common.AtomicRangeInteger; +import org.apache.skywalking.oap.server.core.alarm.AlarmRecord; import org.apache.skywalking.oap.server.core.analysis.TimeBucket; import org.apache.skywalking.oap.server.core.analysis.manual.log.LogRecord; import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; @@ -62,6 +63,8 @@ public class RecordDAO implements IRecordDAO { rawTags = ((SegmentRecord) record).getTagsRawData(); } else if (LogRecord.INDEX_NAME.equals(model.getName())) { rawTags = ((LogRecord) record).getTags(); + } else if (AlarmRecord.INDEX_NAME.equals(model.getName())) { + rawTags = ((AlarmRecord) record).getTags(); } if (nonNull(rawTags)) { Map> collect = rawTags.stream() diff --git a/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/query/AlarmQuery.java b/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/query/AlarmQuery.java index a9c67e598..231ef04b9 100644 --- a/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/query/AlarmQuery.java +++ b/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/query/AlarmQuery.java @@ -21,17 +21,21 @@ package org.apache.skywalking.oap.server.storage.plugin.influxdb.query; import java.io.IOException; import java.util.List; import java.util.Objects; + import lombok.extern.slf4j.Slf4j; import org.apache.skywalking.oap.server.core.alarm.AlarmRecord; +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; import org.apache.skywalking.oap.server.core.query.type.AlarmMessage; import org.apache.skywalking.oap.server.core.query.type.Alarms; import org.apache.skywalking.oap.server.core.query.enumeration.Scope; import org.apache.skywalking.oap.server.core.storage.query.IAlarmQueryDAO; +import org.apache.skywalking.oap.server.library.util.CollectionUtils; import org.apache.skywalking.oap.server.storage.plugin.influxdb.InfluxClient; import org.elasticsearch.common.Strings; import org.influxdb.dto.Query; import org.influxdb.dto.QueryResult; import org.influxdb.querybuilder.SelectQueryImpl; +import org.influxdb.querybuilder.WhereNested; import org.influxdb.querybuilder.WhereQueryImpl; import static org.influxdb.querybuilder.BuiltQuery.QueryBuilder.contains; @@ -50,13 +54,14 @@ public class AlarmQuery implements IAlarmQueryDAO { @Override public Alarms getAlarm(Integer scopeId, String keyword, int limit, int from, long startTB, - long endTB) throws IOException { + long endTB, List tags) throws IOException { WhereQueryImpl recallQuery = select() .function("top", AlarmRecord.START_TIME, limit + from).as(AlarmRecord.START_TIME) .column(AlarmRecord.ID0) .column(AlarmRecord.ALARM_MESSAGE) .column(AlarmRecord.SCOPE) + .column(AlarmRecord.TAGS_RAW_DATA) .from(client.getDatabase(), AlarmRecord.INDEX_NAME) .where(); if (startTB > 0 && endTB > 0) { @@ -69,6 +74,13 @@ public class AlarmQuery implements IAlarmQueryDAO { if (Objects.nonNull(scopeId)) { recallQuery.and(eq(AlarmRecord.SCOPE, scopeId)); } + if (CollectionUtils.isNotEmpty(tags)) { + WhereNested> nested = recallQuery.andNested(); + for (final Tag tag : tags) { + nested.and(contains(tag.getKey(), "'" + tag.getValue() + "'")); + } + nested.close(); + } WhereQueryImpl countQuery = select().count(AlarmRecord.ID0) .from(client.getDatabase(), AlarmRecord.INDEX_NAME) @@ -106,7 +118,10 @@ public class AlarmQuery implements IAlarmQueryDAO { message.setMessage((String) values.get(3)); message.setScope(scope); message.setScopeId(sid); - + String dataBinaryBase64 = (String) values.get(5); + if (!com.google.common.base.Strings.isNullOrEmpty(dataBinaryBase64)) { + parserDataBinaryBase64(dataBinaryBase64, message.getTags()); + } alarms.getMsgs().add(message); }); return alarms; diff --git a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/H2StorageProvider.java b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/H2StorageProvider.java index 98c0612e5..b074566ff 100644 --- a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/H2StorageProvider.java +++ b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/H2StorageProvider.java @@ -140,7 +140,12 @@ public class H2StorageProvider extends ModuleProvider { this.registerServiceImplementation( IMetadataQueryDAO.class, new H2MetadataQueryDAO(h2Client, config.getMetadataQueryMaxSize())); this.registerServiceImplementation(IAggregationQueryDAO.class, new H2AggregationQueryDAO(h2Client)); - this.registerServiceImplementation(IAlarmQueryDAO.class, new H2AlarmQueryDAO(h2Client)); + this.registerServiceImplementation(IAlarmQueryDAO.class, new H2AlarmQueryDAO( + h2Client, + getManager(), + config.getMaxSizeOfArrayColumn(), + config.getNumOfSearchableValuesPerTag() + )); this.registerServiceImplementation( IHistoryDeleteDAO.class, new H2HistoryDeleteDAO(h2Client)); this.registerServiceImplementation(ITopNRecordsQueryDAO.class, new H2TopNRecordsQueryDAO(h2Client)); diff --git a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2AlarmQueryDAO.java b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2AlarmQueryDAO.java index ddba6f7bb..472b4fea3 100644 --- a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2AlarmQueryDAO.java +++ b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2AlarmQueryDAO.java @@ -18,12 +18,18 @@ package org.apache.skywalking.oap.server.storage.plugin.jdbc.h2.dao; +import org.apache.skywalking.oap.server.core.Const; +import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.alarm.AlarmRecord; +import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag; +import org.apache.skywalking.oap.server.core.config.ConfigService; import org.apache.skywalking.oap.server.core.query.type.AlarmMessage; import org.apache.skywalking.oap.server.core.query.type.Alarms; import org.apache.skywalking.oap.server.core.query.enumeration.Scope; import org.apache.skywalking.oap.server.core.storage.query.IAlarmQueryDAO; import org.apache.skywalking.oap.server.library.client.jdbc.hikaricp.JDBCHikariCPClient; +import org.apache.skywalking.oap.server.library.module.ModuleManager; +import org.apache.skywalking.oap.server.library.util.CollectionUtils; import org.elasticsearch.common.Strings; import java.io.IOException; @@ -31,6 +37,7 @@ import java.sql.Connection; import java.sql.ResultSet; import java.sql.SQLException; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; import java.util.Objects; @@ -38,13 +45,33 @@ public class H2AlarmQueryDAO implements IAlarmQueryDAO { private JDBCHikariCPClient client; - public H2AlarmQueryDAO(JDBCHikariCPClient client) { + private final ModuleManager manager; + private final int maxSizeOfArrayColumn; + private final int numOfSearchValuesPerTag; + private List searchableTagKeys; + + public H2AlarmQueryDAO(JDBCHikariCPClient client, + final ModuleManager manager, + final int maxSizeOfArrayColumn, + final int numOfSearchValuesPerTag) { this.client = client; + this.manager = manager; + this.maxSizeOfArrayColumn = maxSizeOfArrayColumn; + this.numOfSearchValuesPerTag = numOfSearchValuesPerTag; } @Override public Alarms getAlarm(Integer scopeId, String keyword, int limit, int from, long startTB, - long endTB) throws IOException { + long endTB, final List tags) throws IOException { + if (searchableTagKeys == null) { + final ConfigService configService = manager.find(CoreModule.NAME) + .provider() + .getService(ConfigService.class); + searchableTagKeys = Arrays.asList(configService.getSearchableAlarmTags().split(Const.COMMA)); + if (searchableTagKeys.size() > maxSizeOfArrayColumn) { + searchableTagKeys = searchableTagKeys.subList(0, maxSizeOfArrayColumn); + } + } StringBuilder sql = new StringBuilder(); List parameters = new ArrayList<>(10); sql.append("from ").append(AlarmRecord.INDEX_NAME).append(" where "); @@ -64,12 +91,31 @@ public class H2AlarmQueryDAO implements IAlarmQueryDAO { sql.append(" and ").append(AlarmRecord.ALARM_MESSAGE).append(" like concat('%',?,'%') "); parameters.add(keyword); } + if (CollectionUtils.isNotEmpty(tags)) { + for (final Tag tag : tags) { + final int foundIdx = searchableTagKeys.indexOf(tag.getKey()); + if (foundIdx > -1) { + sql.append(" and ("); + for (int i = 0; i < numOfSearchValuesPerTag; i++) { + final String physicalColumn = AlarmRecord.TAGS + "_" + (foundIdx * numOfSearchValuesPerTag + i); + sql.append(physicalColumn).append(" = ? "); + parameters.add(tag.toString()); + if (i != numOfSearchValuesPerTag - 1) { + sql.append(" or "); + } + } + sql.append(")"); + } else { + return new Alarms(); + } + } + } sql.append(" order by ").append(AlarmRecord.START_TIME).append(" desc "); Alarms alarms = new Alarms(); try (Connection connection = client.getConnection()) { - try (ResultSet resultSet = client.executeQuery(connection, "select count(1) total from (select 1 " + sql.toString() + " )", parameters + try (ResultSet resultSet = client.executeQuery(connection, buildCountStatement(sql.toString()), parameters .toArray(new Object[0]))) { while (resultSet.next()) { alarms.setTotal(resultSet.getInt("total")); @@ -86,7 +132,10 @@ public class H2AlarmQueryDAO implements IAlarmQueryDAO { message.setStartTime(resultSet.getLong(AlarmRecord.START_TIME)); message.setScope(Scope.Finder.valueOf(resultSet.getInt(AlarmRecord.SCOPE))); message.setScopeId(resultSet.getInt(AlarmRecord.SCOPE)); - + String dataBinaryBase64 = resultSet.getString(AlarmRecord.TAGS_RAW_DATA); + if (!com.google.common.base.Strings.isNullOrEmpty(dataBinaryBase64)) { + parserDataBinaryBase64(dataBinaryBase64, message.getTags()); + } alarms.getMsgs().add(message); } } @@ -97,6 +146,10 @@ public class H2AlarmQueryDAO implements IAlarmQueryDAO { return alarms; } + protected String buildCountStatement(String sql) { + return "select count(1) total from (select 1 " + sql + " )"; + } + protected void buildLimit(StringBuilder sql, int from, int limit) { sql.append(" LIMIT ").append(limit); sql.append(" OFFSET ").append(from); diff --git a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2AlarmRecordBuilder.java b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2AlarmRecordBuilder.java new file mode 100644 index 000000000..100e8401b --- /dev/null +++ b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2AlarmRecordBuilder.java @@ -0,0 +1,79 @@ +/* + * 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.skywalking.oap.server.storage.plugin.jdbc.h2.dao; + +import org.apache.skywalking.apm.util.StringUtil; +import org.apache.skywalking.oap.server.core.Const; +import org.apache.skywalking.oap.server.core.alarm.AlarmRecord; +import org.apache.skywalking.oap.server.core.analysis.record.Record; +import org.apache.skywalking.oap.server.library.util.CollectionUtils; + +import java.util.Base64; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +public class H2AlarmRecordBuilder extends AbstractSearchTagBuilder { + + public H2AlarmRecordBuilder(final int maxSizeOfArrayColumn, + final int numOfSearchableValuesPerTag, + final List searchTagKeys) { + super(maxSizeOfArrayColumn, numOfSearchableValuesPerTag, searchTagKeys, AlarmRecord.TAGS); + } + + @Override + public Record storage2Entity(final Map dbMap) { + AlarmRecord record = new AlarmRecord(); + record.setScope(((Number) dbMap.get(AlarmRecord.SCOPE)).intValue()); + record.setName((String) dbMap.get(AlarmRecord.NAME)); + record.setId0((String) dbMap.get(AlarmRecord.ID0)); + record.setId1((String) dbMap.get(AlarmRecord.ID1)); + record.setAlarmMessage((String) dbMap.get(AlarmRecord.ALARM_MESSAGE)); + record.setStartTime(((Number) dbMap.get(AlarmRecord.START_TIME)).longValue()); + record.setTimeBucket(((Number) dbMap.get(AlarmRecord.TIME_BUCKET)).longValue()); + record.setRuleName((String) dbMap.get(AlarmRecord.RULE_NAME)); + if (StringUtil.isEmpty((String) dbMap.get(AlarmRecord.TAGS_RAW_DATA))) { + record.setTagsRawData(new byte[] {}); + } else { + // Don't read the tags as they has been in the data binary already. + record.setTagsRawData(Base64.getDecoder().decode((String) dbMap.get(AlarmRecord.TAGS_RAW_DATA))); + } + return record; + } + + @Override + public Map entity2Storage(final Record record) { + AlarmRecord storageData = (AlarmRecord) record; + Map map = new HashMap<>(); + map.put(AlarmRecord.SCOPE, storageData.getScope()); + map.put(AlarmRecord.NAME, storageData.getName()); + map.put(AlarmRecord.ID0, storageData.getId0()); + map.put(AlarmRecord.ID1, storageData.getId1()); + map.put(AlarmRecord.ALARM_MESSAGE, storageData.getAlarmMessage()); + map.put(AlarmRecord.START_TIME, storageData.getStartTime()); + map.put(AlarmRecord.TIME_BUCKET, storageData.getTimeBucket()); + map.put(AlarmRecord.RULE_NAME, storageData.getRuleName()); + if (CollectionUtils.isEmpty(storageData.getTagsRawData())) { + map.put(AlarmRecord.TAGS_RAW_DATA, Const.EMPTY_STRING); + } else { + map.put(AlarmRecord.TAGS_RAW_DATA, new String(Base64.getEncoder().encode(storageData.getTagsRawData()))); + } + analysisSearchTag(storageData.getTags(), map); + return map; + } +} diff --git a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2RecordDAO.java b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2RecordDAO.java index 4c621cee7..c8ba888b5 100644 --- a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2RecordDAO.java +++ b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2RecordDAO.java @@ -24,6 +24,7 @@ import java.util.Map; import org.apache.skywalking.oap.server.core.Const; import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.UnexpectedException; +import org.apache.skywalking.oap.server.core.alarm.AlarmRecord; import org.apache.skywalking.oap.server.core.analysis.manual.log.LogRecord; import org.apache.skywalking.oap.server.core.analysis.manual.segment.SegmentRecord; import org.apache.skywalking.oap.server.core.analysis.record.Record; @@ -73,6 +74,18 @@ public class H2RecordDAO extends H2SQLExecutor implements IRecordDAO { Arrays.asList(configService.getSearchableLogsTags() .split(Const.COMMA)) ); + } else if (AlarmRecord.class.equals( + storageBuilder.getClass().getMethod("storage2Entity", Map.class).getReturnType())) { + this.maxSizeOfArrayColumn = maxSizeOfArrayColumn; + final ConfigService configService = manager.find(CoreModule.NAME) + .provider() + .getService(ConfigService.class); + this.storageBuilder = new H2AlarmRecordBuilder( + maxSizeOfArrayColumn, + numOfSearchableValuesPerTag, + Arrays.asList(configService.getSearchableAlarmTags() + .split(Const.COMMA)) + ); } else { this.maxSizeOfArrayColumn = 1; this.storageBuilder = storageBuilder; diff --git a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/mysql/MySQLAlarmQueryDAO.java b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/mysql/MySQLAlarmQueryDAO.java index 59e69f426..fc85f65d6 100644 --- a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/mysql/MySQLAlarmQueryDAO.java +++ b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/mysql/MySQLAlarmQueryDAO.java @@ -18,89 +18,24 @@ package org.apache.skywalking.oap.server.storage.plugin.jdbc.mysql; -import org.apache.skywalking.oap.server.core.alarm.AlarmRecord; -import org.apache.skywalking.oap.server.core.query.type.AlarmMessage; -import org.apache.skywalking.oap.server.core.query.type.Alarms; -import org.apache.skywalking.oap.server.core.query.enumeration.Scope; -import org.apache.skywalking.oap.server.core.storage.query.IAlarmQueryDAO; import org.apache.skywalking.oap.server.library.client.jdbc.hikaricp.JDBCHikariCPClient; -import org.elasticsearch.common.Strings; +import org.apache.skywalking.oap.server.library.module.ModuleManager; +import org.apache.skywalking.oap.server.storage.plugin.jdbc.h2.dao.H2AlarmQueryDAO; -import java.io.IOException; -import java.sql.Connection; -import java.sql.ResultSet; -import java.sql.SQLException; -import java.util.ArrayList; -import java.util.List; -import java.util.Objects; +public class MySQLAlarmQueryDAO extends H2AlarmQueryDAO { -public class MySQLAlarmQueryDAO implements IAlarmQueryDAO { - - private JDBCHikariCPClient client; - - public MySQLAlarmQueryDAO(JDBCHikariCPClient client) { - this.client = client; + public MySQLAlarmQueryDAO(final JDBCHikariCPClient client, + final ModuleManager manager, + final int maxSizeOfArrayColumn, final int numOfSearchValuesPerTag) { + super(client, manager, maxSizeOfArrayColumn, numOfSearchValuesPerTag); } @Override - public Alarms getAlarm(Integer scopeId, String keyword, int limit, int from, long startTB, - long endTB) throws IOException { - StringBuilder sql = new StringBuilder(); - List parameters = new ArrayList<>(10); - sql.append("from ").append(AlarmRecord.INDEX_NAME).append(" where "); - sql.append(" 1=1 "); - if (Objects.nonNull(scopeId)) { - sql.append(" and ").append(AlarmRecord.SCOPE).append(" = ?"); - parameters.add(scopeId.intValue()); - } - if (startTB != 0 && endTB != 0) { - sql.append(" and ").append(AlarmRecord.TIME_BUCKET).append(" >= ?"); - parameters.add(startTB); - sql.append(" and ").append(AlarmRecord.TIME_BUCKET).append(" <= ?"); - parameters.add(endTB); - } - - if (!Strings.isNullOrEmpty(keyword)) { - sql.append(" and ").append(AlarmRecord.ALARM_MESSAGE).append(" like concat('%',?,'%') "); - parameters.add(keyword); - } - sql.append(" order by ").append(AlarmRecord.START_TIME).append(" desc "); - - Alarms alarms = new Alarms(); - try (Connection connection = client.getConnection()) { - - try (ResultSet resultSet = client.executeQuery(connection, buildCountStatement(sql.toString()), parameters - .toArray(new Object[0]))) { - while (resultSet.next()) { - alarms.setTotal(resultSet.getInt("total")); - } - } - - this.buildLimit(sql, from, limit); - - try (ResultSet resultSet = client.executeQuery(connection, "select * " + sql.toString(), parameters.toArray(new Object[0]))) { - while (resultSet.next()) { - AlarmMessage message = new AlarmMessage(); - message.setId(resultSet.getString(AlarmRecord.ID0)); - message.setMessage(resultSet.getString(AlarmRecord.ALARM_MESSAGE)); - message.setStartTime(resultSet.getLong(AlarmRecord.START_TIME)); - message.setScope(Scope.Finder.valueOf(resultSet.getInt(AlarmRecord.SCOPE))); - message.setScopeId(resultSet.getInt(AlarmRecord.SCOPE)); - - alarms.getMsgs().add(message); - } - } - } catch (SQLException e) { - throw new IOException(e); - } - - return alarms; - } - protected void buildLimit(StringBuilder sql, int from, int limit) { sql.append(" LIMIT ").append(from).append(", ").append(limit); } + @Override protected String buildCountStatement(String sql) { return "select count(1) total " + sql; } diff --git a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/mysql/MySQLStorageProvider.java b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/mysql/MySQLStorageProvider.java index 26e1758bf..233431ad2 100644 --- a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/mysql/MySQLStorageProvider.java +++ b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/mysql/MySQLStorageProvider.java @@ -127,7 +127,11 @@ public class MySQLStorageProvider extends ModuleProvider { this.registerServiceImplementation( IMetadataQueryDAO.class, new H2MetadataQueryDAO(mysqlClient, config.getMetadataQueryMaxSize())); this.registerServiceImplementation(IAggregationQueryDAO.class, new MySQLAggregationQueryDAO(mysqlClient)); - this.registerServiceImplementation(IAlarmQueryDAO.class, new MySQLAlarmQueryDAO(mysqlClient)); + this.registerServiceImplementation(IAlarmQueryDAO.class, new MySQLAlarmQueryDAO( + mysqlClient, + getManager(), + config.getMaxSizeOfArrayColumn(), + config.getNumOfSearchableValuesPerTag())); this.registerServiceImplementation( IHistoryDeleteDAO.class, new H2HistoryDeleteDAO(mysqlClient)); this.registerServiceImplementation(ITopNRecordsQueryDAO.class, new H2TopNRecordsQueryDAO(mysqlClient)); @@ -169,6 +173,13 @@ public class MySQLStorageProvider extends ModuleProvider { + "] > maxSizeOfArrayColumn[" + config.getMaxSizeOfArrayColumn() + "]. Potential out of bound in the runtime."); } + final int numOfSearchableAlarmTags = configService.getSearchableAlarmTags().split(Const.COMMA).length; + if (numOfSearchableAlarmTags * config.getNumOfSearchableValuesPerTag() > config.getMaxSizeOfArrayColumn()) { + throw new ModuleStartException("Size of searchableAlarmTags[" + numOfSearchableAlarmTags + + "] * numOfSearchableValuesPerTag[" + config.getNumOfSearchableValuesPerTag() + + "] > maxSizeOfArrayColumn[" + config.getMaxSizeOfArrayColumn() + + "]. Potential out of bound in the runtime."); + } try { mysqlClient.connect(); diff --git a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/postgresql/PostgreSQLStorageProvider.java b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/postgresql/PostgreSQLStorageProvider.java index 75c69bc1d..52800b332 100644 --- a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/postgresql/PostgreSQLStorageProvider.java +++ b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/postgresql/PostgreSQLStorageProvider.java @@ -127,7 +127,12 @@ public class PostgreSQLStorageProvider extends ModuleProvider { this.registerServiceImplementation( IMetadataQueryDAO.class, new H2MetadataQueryDAO(postgresqlClient, config.getMetadataQueryMaxSize())); this.registerServiceImplementation(IAggregationQueryDAO.class, new PostgreSQLAggregationQueryDAO(postgresqlClient)); - this.registerServiceImplementation(IAlarmQueryDAO.class, new PostgreSQLAlarmQueryDAO(postgresqlClient)); + this.registerServiceImplementation(IAlarmQueryDAO.class, new PostgreSQLAlarmQueryDAO( + postgresqlClient, + getManager(), + config.getMaxSizeOfArrayColumn(), + config.getNumOfSearchableValuesPerTag() + )); this.registerServiceImplementation( IHistoryDeleteDAO.class, new H2HistoryDeleteDAO(postgresqlClient)); this.registerServiceImplementation(ITopNRecordsQueryDAO.class, new H2TopNRecordsQueryDAO(postgresqlClient)); @@ -168,6 +173,13 @@ public class PostgreSQLStorageProvider extends ModuleProvider { + "] > maxSizeOfArrayColumn[" + config.getMaxSizeOfArrayColumn() + "]. Potential out of bound in the runtime."); } + final int numOfSearchableAlarmTags = configService.getSearchableAlarmTags().split(Const.COMMA).length; + if (numOfSearchableAlarmTags * config.getNumOfSearchableValuesPerTag() > config.getMaxSizeOfArrayColumn()) { + throw new ModuleStartException("Size of searchableAlarmTags[" + numOfSearchableAlarmTags + + "] * numOfSearchableValuesPerTag[" + config.getNumOfSearchableValuesPerTag() + + "] > maxSizeOfArrayColumn[" + config.getMaxSizeOfArrayColumn() + + "]. Potential out of bound in the runtime."); + } try { postgresqlClient.connect(); diff --git a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/postgresql/dao/PostgreSQLAlarmQueryDAO.java b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/postgresql/dao/PostgreSQLAlarmQueryDAO.java index de49bdc68..3b08b4ee3 100644 --- a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/postgresql/dao/PostgreSQLAlarmQueryDAO.java +++ b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/postgresql/dao/PostgreSQLAlarmQueryDAO.java @@ -19,12 +19,15 @@ package org.apache.skywalking.oap.server.storage.plugin.jdbc.postgresql.dao; import org.apache.skywalking.oap.server.library.client.jdbc.hikaricp.JDBCHikariCPClient; +import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.apache.skywalking.oap.server.storage.plugin.jdbc.mysql.MySQLAlarmQueryDAO; public class PostgreSQLAlarmQueryDAO extends MySQLAlarmQueryDAO { - - public PostgreSQLAlarmQueryDAO(JDBCHikariCPClient client) { - super(client); + + public PostgreSQLAlarmQueryDAO(final JDBCHikariCPClient client, + final ModuleManager manager, + final int maxSizeOfArrayColumn, final int numOfSearchValuesPerTag) { + super(client, manager, maxSizeOfArrayColumn, numOfSearchValuesPerTag); } @Override diff --git a/oap-server/server-storage-plugin/storage-tidb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/tidb/TiDBStorageProvider.java b/oap-server/server-storage-plugin/storage-tidb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/tidb/TiDBStorageProvider.java index 0bfc9d884..05d057cab 100644 --- a/oap-server/server-storage-plugin/storage-tidb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/tidb/TiDBStorageProvider.java +++ b/oap-server/server-storage-plugin/storage-tidb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/tidb/TiDBStorageProvider.java @@ -130,7 +130,11 @@ public class TiDBStorageProvider extends ModuleProvider { this.registerServiceImplementation( IMetadataQueryDAO.class, new H2MetadataQueryDAO(mysqlClient, config.getMetadataQueryMaxSize())); this.registerServiceImplementation(IAggregationQueryDAO.class, new MySQLAggregationQueryDAO(mysqlClient)); - this.registerServiceImplementation(IAlarmQueryDAO.class, new MySQLAlarmQueryDAO(mysqlClient)); + this.registerServiceImplementation(IAlarmQueryDAO.class, new MySQLAlarmQueryDAO( + mysqlClient, + getManager(), + config.getMaxSizeOfArrayColumn(), + config.getNumOfSearchableValuesPerTag())); this.registerServiceImplementation( IHistoryDeleteDAO.class, new H2HistoryDeleteDAO(mysqlClient)); this.registerServiceImplementation(ITopNRecordsQueryDAO.class, new H2TopNRecordsQueryDAO(mysqlClient)); @@ -174,6 +178,13 @@ public class TiDBStorageProvider extends ModuleProvider { + "] > maxSizeOfArrayColumn[" + config.getMaxSizeOfArrayColumn() + "]. Potential out of bound in the runtime."); } + final int numOfSearchableAlarmTags = configService.getSearchableAlarmTags().split(Const.COMMA).length; + if (numOfSearchableAlarmTags * config.getNumOfSearchableValuesPerTag() > config.getMaxSizeOfArrayColumn()) { + throw new ModuleStartException("Size of searchableAlarmTags[" + numOfSearchableAlarmTags + + "] * numOfSearchableValuesPerTag[" + config.getNumOfSearchableValuesPerTag() + + "] > maxSizeOfArrayColumn[" + config.getMaxSizeOfArrayColumn() + + "]. Potential out of bound in the runtime."); + } try { mysqlClient.connect(); diff --git a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/SimpleQueryClient.java b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/SimpleQueryClient.java index a45244d21..b69ad6f87 100644 --- a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/SimpleQueryClient.java +++ b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/SimpleQueryClient.java @@ -387,7 +387,8 @@ public class SimpleQueryClient { .replace("{start}", query.start()) .replace("{end}", query.end()) .replace("{pageSize}", "20") - .replace("{needTotal}", "true"); + .replace("{needTotal}", "true") + .replace("{tags}", objectMapper.writeValueAsString(query.tags())); LOGGER.info("Query: {}", queryString); final ResponseEntity> responseEntity = restTemplate.exchange( new RequestEntity<>(queryString, HttpMethod.POST, URI.create(endpointUrl)), @@ -398,7 +399,7 @@ public class SimpleQueryClient { if (responseEntity.getStatusCode() != HttpStatus.OK) { throw new RuntimeException("Response status != 200, actual: " + responseEntity.getStatusCode()); } - + LOGGER.info("Result: {}", responseEntity.getBody()); return Objects.requireNonNull(responseEntity.getBody()).getData().getGetAlarm(); } diff --git a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/Alarm.java b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/Alarm.java new file mode 100644 index 000000000..6a2a2fec4 --- /dev/null +++ b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/Alarm.java @@ -0,0 +1,34 @@ +/* + * 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.skywalking.e2e.alarm; + +import lombok.Data; +import lombok.experimental.Accessors; +import org.apache.skywalking.e2e.common.KeyValue; + +import java.util.List; + +@Data +@Accessors(chain = true) +public class Alarm { + private String startTime; + private String scope; + private String id; + private String message; + private List tags; +} diff --git a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/AlarmMatcher.java b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/AlarmMatcher.java index 74ac9c215..f047b2b26 100644 --- a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/AlarmMatcher.java +++ b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/AlarmMatcher.java @@ -18,20 +18,44 @@ package org.apache.skywalking.e2e.alarm; import lombok.Data; +import org.apache.skywalking.e2e.common.KeyValue; +import org.apache.skywalking.e2e.common.KeyValueMatcher; import org.apache.skywalking.e2e.verification.AbstractMatcher; +import java.util.List; + +import static java.util.Objects.nonNull; +import static org.assertj.core.api.Assertions.fail; + @Data -public class AlarmMatcher extends AbstractMatcher { +public class AlarmMatcher extends AbstractMatcher { private String startTime; private String scope; private String id; private String message; + private List tags; @Override - public void verify(AlarmMatcher alarmMatcher) { - doVerify(this.startTime, alarmMatcher.startTime); - doVerify(this.scope, alarmMatcher.scope); - doVerify(this.id, alarmMatcher.id); - doVerify(this.message, alarmMatcher.message); + public void verify(Alarm alarm) { + doVerify(this.startTime, alarm.getStartTime()); + doVerify(this.scope, alarm.getScope()); + doVerify(this.id, alarm.getId()); + doVerify(this.message, alarm.getMessage()); + if (nonNull(getTags())) { + for (final KeyValueMatcher matcher : getTags()) { + boolean matched = false; + for (final KeyValue keyValue : alarm.getTags()) { + try { + matcher.verify(keyValue); + matched = true; + } catch (Throwable ignore) { + + } + } + if (!matched) { + fail("\nExpected: %s\n Actual: %s", getTags(), alarm.getTags()); + } + } + } } } diff --git a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/AlarmQuery.java b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/AlarmQuery.java index 37a7d1dea..bc9af0091 100644 --- a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/AlarmQuery.java +++ b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/AlarmQuery.java @@ -19,5 +19,32 @@ package org.apache.skywalking.e2e.alarm; import org.apache.skywalking.e2e.AbstractQuery; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.HashMap; + public class AlarmQuery extends AbstractQuery { + private List> tags = Collections.emptyList(); + + public List> tags() { + return tags; + } + + public AlarmQuery tags(List> tags) { + this.tags = tags; + return this; + } + + public AlarmQuery addTag(String key, String value) { + if (Collections.EMPTY_LIST.equals(tags)) { + tags = new ArrayList<>(); + } + Map tag = new HashMap<>(); + tag.put("key", key); + tag.put("value", value); + tags.add(tag); + return this; + } } diff --git a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/AlarmsMatcher.java b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/AlarmsMatcher.java index b922df868..0b188659c 100644 --- a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/AlarmsMatcher.java +++ b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/AlarmsMatcher.java @@ -21,31 +21,41 @@ import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.junit.Assert; +import java.util.LinkedList; +import java.util.List; + import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.fail; @Data @Slf4j public class AlarmsMatcher { - private GetAlarm alarm; - public void verify(final GetAlarm alarm) { - LOGGER.info("alarms:{} matchers:{}", alarm, this.alarm); - Assert.assertEquals(this.alarm.getTotal(), alarm.getTotal()); + private int total; - assertThat(this.alarm.getMsgs()).hasSameSizeAs(alarm.getMsgs()); + private List matchers; - for (int i = 0; i < this.alarm.getMsgs().size(); i++) { + public AlarmsMatcher() { + this.matchers = new LinkedList<>(); + } + + public void verify(final GetAlarm alarms) { + LOGGER.info("alarms:{} matchers:{}", alarms, this.matchers); + Assert.assertEquals(this.total, alarms.getTotal()); + + assertThat(this.matchers).hasSameSizeAs(alarms.getMsgs()); + + for (int i = 0; i < this.matchers.size(); i++) { boolean matched = false; - for (AlarmMatcher alarmMatcher : alarm.getMsgs()) { + for (Alarm alarm : alarms.getMsgs()) { try { - this.alarm.getMsgs().get(i).verify(alarmMatcher); + this.matchers.get(i).verify(alarm); matched = true; } catch (Throwable ignored) { } } if (!matched) { - fail("\nExpected: %s\nActual: %s", this.alarm, alarm); + fail("\nExpected: %s\nActual: %s", this.matchers, alarms); } } } diff --git a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/GetAlarm.java b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/GetAlarm.java index 8b8e388f9..650604aa8 100644 --- a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/GetAlarm.java +++ b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/GetAlarm.java @@ -24,5 +24,5 @@ import java.util.List; @Data public class GetAlarm { private int total; - private List msgs; + private List msgs; } diff --git a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarm.java b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarm.java new file mode 100644 index 000000000..20b2b181c --- /dev/null +++ b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarm.java @@ -0,0 +1,36 @@ +/* + * 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.skywalking.e2e.alarm; + +import lombok.Data; +import org.apache.skywalking.e2e.common.KeyValue; + +import java.util.List; + +@Data +public class HookAlarm { + private String scopeId; + private String scope; + private String name; + private String id0; + private String id1; + private String ruleName; + private String alarmMessage; + private String startTime; + private List tags; +} diff --git a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarmMatcher.java b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarmMatcher.java index aee9168fe..3401b8e03 100644 --- a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarmMatcher.java +++ b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarmMatcher.java @@ -18,10 +18,17 @@ package org.apache.skywalking.e2e.alarm; import lombok.Data; +import org.apache.skywalking.e2e.common.KeyValue; +import org.apache.skywalking.e2e.common.KeyValueMatcher; import org.apache.skywalking.e2e.verification.AbstractMatcher; +import java.util.List; + +import static java.util.Objects.nonNull; +import static org.assertj.core.api.Assertions.fail; + @Data -public class HookAlarmMatcher extends AbstractMatcher { +public class HookAlarmMatcher extends AbstractMatcher { private String scopeId; private String scope; private String name; @@ -30,16 +37,33 @@ public class HookAlarmMatcher extends AbstractMatcher { private String ruleName; private String alarmMessage; private String startTime; + private List tags; @Override - public void verify(HookAlarmMatcher hookAlarmMatcher) { - doVerify(this.scopeId, hookAlarmMatcher.getScopeId()); - doVerify(this.scope, hookAlarmMatcher.getScope()); - doVerify(this.name, hookAlarmMatcher.getName()); - doVerify(this.id0, hookAlarmMatcher.getId0()); - doVerify(this.id1, hookAlarmMatcher.getId1()); - doVerify(this.ruleName, hookAlarmMatcher.getRuleName()); - doVerify(this.alarmMessage, hookAlarmMatcher.getAlarmMessage()); - doVerify(this.startTime, hookAlarmMatcher.getStartTime()); + public void verify(HookAlarm hookAlarm) { + doVerify(this.scopeId, hookAlarm.getScopeId()); + doVerify(this.scope, hookAlarm.getScope()); + doVerify(this.name, hookAlarm.getName()); + doVerify(this.id0, hookAlarm.getId0()); + doVerify(this.id1, hookAlarm.getId1()); + doVerify(this.ruleName, hookAlarm.getRuleName()); + doVerify(this.alarmMessage, hookAlarm.getAlarmMessage()); + doVerify(this.startTime, hookAlarm.getStartTime()); + if (nonNull(getTags())) { + for (final KeyValueMatcher matcher : getTags()) { + boolean matched = false; + for (final KeyValue keyValue : hookAlarm.getTags()) { + try { + matcher.verify(keyValue); + matched = true; + } catch (Throwable ignore) { + + } + } + if (!matched) { + fail("\nExpected: %s\n Actual: %s", getTags(), hookAlarm.getTags()); + } + } + } } } diff --git a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarms.java b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarms.java index c4f92d7a8..3809bc864 100644 --- a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarms.java +++ b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarms.java @@ -20,35 +20,12 @@ package org.apache.skywalking.e2e.alarm; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; -import org.apache.skywalking.e2e.verification.AbstractMatcher; import java.util.ArrayList; -import static org.assertj.core.api.Assertions.assertThat; -import static org.assertj.core.api.Assertions.fail; - @Data @NoArgsConstructor @AllArgsConstructor -public class HookAlarms extends AbstractMatcher { - private ArrayList messages; - - @Override - public void verify(HookAlarms hookAlarms) { - assertThat(this.messages).hasSameSizeAs(hookAlarms.getMessages()); - - for (int i = 0; i < this.messages.size(); i++) { - boolean matched = false; - for (HookAlarmMatcher alarmMatcher : hookAlarms.getMessages()) { - try { - this.messages.get(i).verify(alarmMatcher); - matched = true; - } catch (Throwable ignored) { - } - } - if (!matched) { - fail("\nExpected: %s\nActual: %s", this.messages.get(i), hookAlarms.getMessages()); - } - } - } +public class HookAlarms { + private ArrayList messages; } diff --git a/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarmsMatcher.java b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarmsMatcher.java new file mode 100644 index 000000000..2b6ebc110 --- /dev/null +++ b/test/e2e/e2e-data/src/main/java/org/apache/skywalking/e2e/alarm/HookAlarmsMatcher.java @@ -0,0 +1,54 @@ +/* + * 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.skywalking.e2e.alarm; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; +import org.apache.skywalking.e2e.verification.AbstractMatcher; + +import java.util.ArrayList; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.fail; + +@Data +@NoArgsConstructor +@AllArgsConstructor +public class HookAlarmsMatcher extends AbstractMatcher { + private ArrayList messages; + + @Override + public void verify(HookAlarms hookAlarms) { + assertThat(this.messages).hasSameSizeAs(hookAlarms.getMessages()); + + for (int i = 0; i < this.messages.size(); i++) { + boolean matched = false; + for (HookAlarm hookAlarm : hookAlarms.getMessages()) { + try { + this.messages.get(i).verify(hookAlarm); + matched = true; + } catch (Throwable ignored) { + } + } + if (!matched) { + fail("\nExpected: %s\nActual: %s", this.messages.get(i), hookAlarms.getMessages()); + } + } + } +} diff --git a/test/e2e/e2e-data/src/main/resources/read-alarms.gql b/test/e2e/e2e-data/src/main/resources/read-alarms.gql index 24e85af01..fc9787c84 100644 --- a/test/e2e/e2e-data/src/main/resources/read-alarms.gql +++ b/test/e2e/e2e-data/src/main/resources/read-alarms.gql @@ -14,10 +14,11 @@ # limitations under the License. { - "query":"query queryData($duration: Duration!, $paging: Pagination!) { + "query":"query queryData($duration: Duration!, $paging: Pagination!, $tags: [AlarmTag]) { getAlarm: getAlarm( duration: $duration, - paging: $paging + paging: $paging, + tags: $tags ) { total msgs { @@ -25,6 +26,9 @@ scope id message + tags { + key, value + } } }}", "variables":{ @@ -36,6 +40,7 @@ "paging":{ "pageSize":"{pageSize}", "needTotal":"{needTotal}" - } + }, + "tags": {tags} } } \ No newline at end of file diff --git a/test/e2e/e2e-service-provider/src/main/java/org/apache/skywalking/e2e/KeyValue.java b/test/e2e/e2e-service-provider/src/main/java/org/apache/skywalking/e2e/KeyValue.java new file mode 100644 index 000000000..64bb447d1 --- /dev/null +++ b/test/e2e/e2e-service-provider/src/main/java/org/apache/skywalking/e2e/KeyValue.java @@ -0,0 +1,29 @@ +/* + * 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.skywalking.e2e; + +import lombok.Data; +import lombok.experimental.Accessors; + +@Data +@Accessors(chain = true) +public class KeyValue { + + private String key; + private String value; +} diff --git a/test/e2e/e2e-service-provider/src/main/java/org/apache/skywalking/e2e/controller/AlarmController.java b/test/e2e/e2e-service-provider/src/main/java/org/apache/skywalking/e2e/controller/AlarmController.java index dcc118451..fd4467028 100644 --- a/test/e2e/e2e-service-provider/src/main/java/org/apache/skywalking/e2e/controller/AlarmController.java +++ b/test/e2e/e2e-service-provider/src/main/java/org/apache/skywalking/e2e/controller/AlarmController.java @@ -23,6 +23,7 @@ import lombok.Data; import lombok.Getter; import lombok.NoArgsConstructor; import lombok.Setter; +import org.apache.skywalking.e2e.KeyValue; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; @@ -64,6 +65,7 @@ public class AlarmController { private String ruleName; private String alarmMessage; private long startTime; + private List tags; } /** diff --git a/test/e2e/e2e-test/docker/alarm/alarm-settings.yml b/test/e2e/e2e-test/docker/alarm/alarm-settings.yml index 76d21f311..f765d408d 100755 --- a/test/e2e/e2e-test/docker/alarm/alarm-settings.yml +++ b/test/e2e/e2e-test/docker/alarm/alarm-settings.yml @@ -24,6 +24,9 @@ rules: silence-period: 1 message: Response time of service {name} is more than 10ms in 1 minutes of last 10 minutes. only-as-condition: false + tags: + level: WARNING + receivers: lisi # service sla > 1% service_sla_rule: metrics-name: service_sla @@ -38,6 +41,9 @@ composite-rules: comp_rule: expression: service_resp_time_rule && service_sla_rule message: Service {name} response time is more than 10ms and sla is more than 1%. + tags: + level: CRITICAL + receivers: zhangsan webhooks: - http://provider:9090/alarm/receive diff --git a/test/e2e/e2e-test/docker/alarm/docker-compose.es6.yml b/test/e2e/e2e-test/docker/alarm/docker-compose.es6.yml new file mode 100644 index 000000000..ba12a4370 --- /dev/null +++ b/test/e2e/e2e-test/docker/alarm/docker-compose.es6.yml @@ -0,0 +1,65 @@ +# 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. + +version: '2.1' + +services: + es: + image: elastic/elasticsearch:6.3.2 + expose: + - 9200 + networks: + - e2e + environment: + - discovery.type=single-node + healthcheck: + test: ["CMD", "bash", "-c", "cat < /dev/null > /dev/tcp/127.0.0.1/9200"] + interval: 5s + timeout: 60s + retries: 120 + + oap: + extends: + file: ../base-compose.yml + service: oap + environment: + SW_STORAGE: elasticsearch + SW_SEARCHABLE_ALARM_TAG_KEYS: level,receivers + depends_on: + es: + condition: service_healthy + volumes: + - ./alarm-settings.yml:/skywalking/config/alarm-settings.yml + + ui: + extends: + file: ../base-compose.yml + service: ui + depends_on: + oap: + condition: service_healthy + + provider: + extends: + file: ../base-compose.yml + service: provider + environment: + SW_AGENT_NAME: e2e-service-provider + depends_on: + oap: + condition: service_healthy + +networks: + e2e: diff --git a/test/e2e/e2e-test/docker/alarm/docker-compose.es7.yml b/test/e2e/e2e-test/docker/alarm/docker-compose.es7.yml new file mode 100644 index 000000000..276dc6ec5 --- /dev/null +++ b/test/e2e/e2e-test/docker/alarm/docker-compose.es7.yml @@ -0,0 +1,65 @@ +# 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. + +version: '2.1' + +services: + es: + image: elastic/elasticsearch:7.4.2 + expose: + - 9200 + networks: + - e2e + environment: + - discovery.type=single-node + healthcheck: + test: ["CMD", "bash", "-c", "cat < /dev/null > /dev/tcp/127.0.0.1/9200"] + interval: 5s + timeout: 60s + retries: 120 + + oap: + extends: + file: ../base-compose.yml + service: oap-es7 + environment: + SW_STORAGE: elasticsearch7 + SW_SEARCHABLE_ALARM_TAG_KEYS: level,receivers + depends_on: + es: + condition: service_healthy + volumes: + - ./alarm-settings.yml:/skywalking/config/alarm-settings.yml + + ui: + extends: + file: ../base-compose.yml + service: ui + depends_on: + oap: + condition: service_healthy + + provider: + extends: + file: ../base-compose.yml + service: provider + environment: + SW_AGENT_NAME: e2e-service-provider + depends_on: + oap: + condition: service_healthy + +networks: + e2e: diff --git a/test/e2e/e2e-test/docker/alarm/docker-compose.h2.yml b/test/e2e/e2e-test/docker/alarm/docker-compose.h2.yml new file mode 100644 index 000000000..38899dc7a --- /dev/null +++ b/test/e2e/e2e-test/docker/alarm/docker-compose.h2.yml @@ -0,0 +1,66 @@ +# 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. + +version: '2.1' + +services: + h2db: + build: + context: . + dockerfile: ../Dockerfile.h2 + networks: + - e2e + expose: + - 1521 + healthcheck: + test: [ "CMD", "sh", "-c", "nc -z 127.0.0.1 1521" ] + interval: 5s + timeout: 60s + retries: 120 + + oap: + extends: + file: ../base-compose.yml + service: oap + environment: + SW_STORAGE: h2 + SW_STORAGE_H2_URL: jdbc:h2:tcp://h2db:1521/skywalking-oap-db + SW_SEARCHABLE_ALARM_TAG_KEYS: level,receivers + depends_on: + h2db: + condition: service_healthy + volumes: + - ./alarm-settings.yml:/skywalking/config/alarm-settings.yml + + ui: + extends: + file: ../base-compose.yml + service: ui + depends_on: + oap: + condition: service_healthy + + provider: + extends: + file: ../base-compose.yml + service: provider + environment: + SW_AGENT_NAME: e2e-service-provider + depends_on: + oap: + condition: service_healthy + +networks: + e2e: diff --git a/test/e2e/e2e-test/docker/alarm/docker-compose.yml b/test/e2e/e2e-test/docker/alarm/docker-compose.influxdb.yml similarity index 76% rename from test/e2e/e2e-test/docker/alarm/docker-compose.yml rename to test/e2e/e2e-test/docker/alarm/docker-compose.influxdb.yml index e38132530..370e33372 100644 --- a/test/e2e/e2e-test/docker/alarm/docker-compose.yml +++ b/test/e2e/e2e-test/docker/alarm/docker-compose.influxdb.yml @@ -16,10 +16,28 @@ version: '2.1' services: + influxdb: + image: influxdb:1.7.9 + expose: + - 8086 + networks: + - e2e + healthcheck: + test: ["CMD", "bash", "-c", "cat < /dev/null > /dev/tcp/127.0.0.1/8086"] + interval: 5s + timeout: 60s + retries: 120 + oap: extends: file: ../base-compose.yml service: oap + environment: + SW_STORAGE: influxdb + SW_SEARCHABLE_ALARM_TAG_KEYS: level,receivers + depends_on: + influxdb: + condition: service_healthy volumes: - ./alarm-settings.yml:/skywalking/config/alarm-settings.yml diff --git a/test/e2e/e2e-test/docker/alarm/docker-compose.mysql.yml b/test/e2e/e2e-test/docker/alarm/docker-compose.mysql.yml new file mode 100644 index 000000000..57e51b241 --- /dev/null +++ b/test/e2e/e2e-test/docker/alarm/docker-compose.mysql.yml @@ -0,0 +1,68 @@ +# 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. + +version: '2.1' + +services: + mysql: + image: mysql/mysql-server:8.0.13 + networks: + - e2e + expose: + - 3306 + environment: + - MYSQL_ROOT_PASSWORD=root@1234 + - MYSQL_DATABASE=swtest + - MYSQL_ROOT_HOST=% + healthcheck: + test: ["CMD", "bash", "-c", "cat < /dev/null > /dev/tcp/127.0.0.1/3306"] + interval: 5s + timeout: 60s + retries: 120 + + oap: + extends: + file: ../base-compose.yml + service: oap + environment: + SW_STORAGE: mysql + SW_SEARCHABLE_ALARM_TAG_KEYS: level,receivers + depends_on: + mysql: + condition: service_healthy + volumes: + - ./alarm-settings.yml:/skywalking/config/alarm-settings.yml + entrypoint: ['sh', '-c', '/download-mysql.sh && /skywalking/docker-entrypoint.sh'] + + ui: + extends: + file: ../base-compose.yml + service: ui + depends_on: + oap: + condition: service_healthy + + provider: + extends: + file: ../base-compose.yml + service: provider + environment: + SW_AGENT_NAME: e2e-service-provider + depends_on: + oap: + condition: service_healthy + +networks: + e2e: diff --git a/test/e2e/e2e-test/docker/alarm/docker-compose.postgres.yml b/test/e2e/e2e-test/docker/alarm/docker-compose.postgres.yml new file mode 100644 index 000000000..d47107416 --- /dev/null +++ b/test/e2e/e2e-test/docker/alarm/docker-compose.postgres.yml @@ -0,0 +1,67 @@ +# 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. + +version: '2.1' + +services: + postgres: + image: postgres:13 + networks: + - e2e + expose: + - 5432 + environment: + - POSTGRES_PASSWORD=123456 + - POSTGRES_DB=skywalking + healthcheck: + test: ["CMD", "bash", "-c", "cat < /dev/null > /dev/tcp/127.0.0.1/5432"] + interval: 5s + timeout: 60s + retries: 120 + + oap: + extends: + file: ../base-compose.yml + service: oap + environment: + SW_STORAGE: postgresql + SW_JDBC_URL: "jdbc:postgresql://postgres:5432/skywalking" + SW_SEARCHABLE_ALARM_TAG_KEYS: level,receivers + depends_on: + postgres: + condition: service_healthy + volumes: + - ./alarm-settings.yml:/skywalking/config/alarm-settings.yml + + ui: + extends: + file: ../base-compose.yml + service: ui + depends_on: + oap: + condition: service_healthy + + provider: + extends: + file: ../base-compose.yml + service: provider + environment: + SW_AGENT_NAME: e2e-service-provider + depends_on: + oap: + condition: service_healthy + +networks: + e2e: diff --git a/test/e2e/e2e-test/src/test/java/org/apache/skywalking/e2e/alarm/AlarmE2E.java b/test/e2e/e2e-test/src/test/java/org/apache/skywalking/e2e/alarm/AlarmE2E.java index 319c5794e..facde91e3 100644 --- a/test/e2e/e2e-test/src/test/java/org/apache/skywalking/e2e/alarm/AlarmE2E.java +++ b/test/e2e/e2e-test/src/test/java/org/apache/skywalking/e2e/alarm/AlarmE2E.java @@ -47,7 +47,7 @@ import static org.apache.skywalking.e2e.utils.Yamls.load; public class AlarmE2E extends SkyWalkingTestAdapter { @SuppressWarnings("unused") - @DockerCompose("docker/alarm/docker-compose.yml") + @DockerCompose("docker/alarm/docker-compose.${SW_STORAGE}.yml") protected DockerComposeContainer justForSideEffects; @SuppressWarnings("unused") @@ -81,9 +81,10 @@ public class AlarmE2E extends SkyWalkingTestAdapter { @RetryableTest @Order(2) - void basicAlarm() throws Exception { + void + basicAlarm() throws Exception { // Wait all alarm notified(single and compose) - validate("expected/alarm/silence-before-graphql.yml", "expected/alarm/silence-before-webhook.yml"); + validate("expected/alarm/silence-before-graphql-warn.yml", "expected/alarm/silence-before-graphql-critical.yml", "expected/alarm/silence-before-webhook.yml"); // Wait silence period finished TimeUnit.SECONDS.sleep(90); @@ -93,19 +94,22 @@ public class AlarmE2E extends SkyWalkingTestAdapter { @Order(3) void afterSilenceAlarm() throws Exception { // Retry to send request and check silence config - validate("expected/alarm/silence-after-graphql.yml", "expected/alarm/silence-after-webhook.yml"); + validate("expected/alarm/silence-after-graphql-warn.yml", "expected/alarm/silence-after-graphql-critical.yml", "expected/alarm/silence-after-webhook.yml"); } - private void validate(String alarmFile, String hookFile) throws Exception { + private void validate(String alarmFileWarn, String alarmFileCritical, String hookFile) throws Exception { // validate graphql - GetAlarm alarms = graphql.readAlarms(new AlarmQuery().start(startTime).end(now())); + GetAlarm alarms = graphql.readAlarms(new AlarmQuery().start(startTime).end(now()).addTag("level", "WARNING").addTag("receivers", "lisi")); LOGGER.info("alarms query: {}", alarms); - load(alarmFile).as(AlarmsMatcher.class).verify(alarms); + load(alarmFileWarn).as(AlarmsMatcher.class).verify(alarms); + alarms = graphql.readAlarms(new AlarmQuery().start(startTime).end(now()).addTag("level", "CRITICAL").addTag("receivers", "zhangsan")); + LOGGER.info("alarms query: {}", alarms); + load(alarmFileCritical).as(AlarmsMatcher.class).verify(alarms); // validate web hook receiver ResponseEntity responseEntity = restTemplate.postForEntity("http://" + serviceHostPort.host() + ":" + serviceHostPort.port() + "/alarm/read", null, HookAlarms.class); LOGGER.info("alarms hook: {}", responseEntity.getBody()); - load(hookFile).as(HookAlarms.class).verify(responseEntity.getBody()); + load(hookFile).as(HookAlarmsMatcher.class).verify(responseEntity.getBody()); } } diff --git a/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-graphql-critical.yml b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-graphql-critical.yml new file mode 100644 index 000000000..af0012b24 --- /dev/null +++ b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-graphql-critical.yml @@ -0,0 +1,36 @@ +# 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. + +total: 2 +matchers: + - startTime: gt 0 + scope: Service + id: ZTJlLXNlcnZpY2UtcHJvdmlkZXI=.1 + message: Service e2e-service-provider response time is more than 10ms and sla is more than 1%. + tags: + - key: level + value: CRITICAL + - key: receivers + value: zhangsan + - startTime: gt 0 + scope: Service + id: ZTJlLXNlcnZpY2UtcHJvdmlkZXI=.1 + message: Service e2e-service-provider response time is more than 10ms and sla is more than 1%. + tags: + - key: level + value: CRITICAL + - key: receivers + value: zhangsan + diff --git a/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-graphql-warn.yml b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-graphql-warn.yml new file mode 100644 index 000000000..3809e6620 --- /dev/null +++ b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-graphql-warn.yml @@ -0,0 +1,35 @@ +# 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. + +total: 2 +matchers: + - startTime: gt 0 + scope: Service + id: ZTJlLXNlcnZpY2UtcHJvdmlkZXI=.1 + message: Response time of service e2e-service-provider is more than 10ms in 1 minutes of last 10 minutes. + tags: + - key: level + value: WARNING + - key: receivers + value: lisi + - startTime: gt 0 + scope: Service + id: ZTJlLXNlcnZpY2UtcHJvdmlkZXI=.1 + message: Response time of service e2e-service-provider is more than 10ms in 1 minutes of last 10 minutes. + tags: + - key: level + value: WARNING + - key: receivers + value: lisi diff --git a/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-graphql.yml b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-graphql.yml deleted file mode 100644 index 329685020..000000000 --- a/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-graphql.yml +++ /dev/null @@ -1,34 +0,0 @@ -# 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. - -alarm: - total: 4 - msgs: - - startTime: gt 0 - scope: Service - id: ZTJlLXNlcnZpY2UtcHJvdmlkZXI=.1 - message: Service e2e-service-provider response time is more than 10ms and sla is more than 1%. - - startTime: gt 0 - scope: Service - id: ZTJlLXNlcnZpY2UtcHJvdmlkZXI=.1 - message: Response time of service e2e-service-provider is more than 10ms in 1 minutes of last 10 minutes. - - startTime: gt 0 - scope: Service - id: ZTJlLXNlcnZpY2UtcHJvdmlkZXI=.1 - message: Service e2e-service-provider response time is more than 10ms and sla is more than 1%. - - startTime: gt 0 - scope: Service - id: ZTJlLXNlcnZpY2UtcHJvdmlkZXI=.1 - message: Response time of service e2e-service-provider is more than 10ms in 1 minutes of last 10 minutes. diff --git a/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-webhook.yml b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-webhook.yml index 669c44f7e..634e9044a 100644 --- a/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-webhook.yml +++ b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-after-webhook.yml @@ -21,6 +21,11 @@ messages: id1: "" ruleName: service_resp_time_rule alarmMessage: Response time of service e2e-service-provider is more than 10ms in 1 minutes of last 10 minutes. + tags: + - key: level + value: WARNING + - key: receivers + value: lisi startTime: gt 0 - scopeId: 1 scope: SERVICE @@ -30,6 +35,11 @@ messages: ruleName: comp_rule alarmMessage: Service e2e-service-provider response time is more than 10ms and sla is more than 1%. startTime: gt 0 + tags: + - key: level + value: CRITICAL + - key: receivers + value: zhangsan - scopeId: 1 scope: SERVICE name: e2e-service-provider @@ -37,6 +47,11 @@ messages: id1: "" ruleName: service_resp_time_rule alarmMessage: Response time of service e2e-service-provider is more than 10ms in 1 minutes of last 10 minutes. + tags: + - key: level + value: WARNING + - key: receivers + value: lisi startTime: gt 0 - scopeId: 1 scope: SERVICE @@ -46,4 +61,9 @@ messages: ruleName: comp_rule alarmMessage: Service e2e-service-provider response time is more than 10ms and sla is more than 1%. startTime: gt 0 + tags: + - key: level + value: CRITICAL + - key: receivers + value: zhangsan diff --git a/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-graphql.yml b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-graphql-critical.yml similarity index 65% rename from test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-graphql.yml rename to test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-graphql-critical.yml index 5661e188e..4ae75a21e 100644 --- a/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-graphql.yml +++ b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-graphql-critical.yml @@ -13,14 +13,14 @@ # See the License for the specific language governing permissions and # limitations under the License. -alarm: - total: 2 - msgs: - - startTime: gt 0 - scope: Service - id: ZTJlLXNlcnZpY2UtcHJvdmlkZXI=.1 - message: Service e2e-service-provider response time is more than 10ms and sla is more than 1%. - - startTime: gt 0 - scope: Service - id: ZTJlLXNlcnZpY2UtcHJvdmlkZXI=.1 - message: Response time of service e2e-service-provider is more than 10ms in 1 minutes of last 10 minutes. +total: 1 +matchers: + - startTime: gt 0 + scope: Service + id: ZTJlLXNlcnZpY2UtcHJvdmlkZXI=.1 + message: Service e2e-service-provider response time is more than 10ms and sla is more than 1%. + tags: + - key: level + value: CRITICAL + - key: receivers + value: zhangsan diff --git a/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-graphql-warn.yml b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-graphql-warn.yml new file mode 100644 index 000000000..d95de2fd5 --- /dev/null +++ b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-graphql-warn.yml @@ -0,0 +1,26 @@ +# 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. + +total: 1 +matchers: + - startTime: gt 0 + scope: Service + id: ZTJlLXNlcnZpY2UtcHJvdmlkZXI=.1 + message: Response time of service e2e-service-provider is more than 10ms in 1 minutes of last 10 minutes. + tags: + - key: level + value: WARNING + - key: receivers + value: lisi diff --git a/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-webhook.yml b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-webhook.yml index 51bba0011..61f6c79d4 100644 --- a/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-webhook.yml +++ b/test/e2e/e2e-test/src/test/resources/expected/alarm/silence-before-webhook.yml @@ -22,6 +22,11 @@ messages: ruleName: service_resp_time_rule alarmMessage: Response time of service e2e-service-provider is more than 10ms in 1 minutes of last 10 minutes. startTime: gt 0 + tags: + - key: level + value: WARNING + - key: receivers + value: lisi - scopeId: 1 scope: SERVICE name: e2e-service-provider @@ -30,3 +35,8 @@ messages: ruleName: comp_rule alarmMessage: Service e2e-service-provider response time is more than 10ms and sla is more than 1%. startTime: gt 0 + tags: + - key: level + value: CRITICAL + - key: receivers + value: zhangsan