Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@

## New features
- [PagerDuty] Add NTLM authentication support for HTTPS proxies - [#1772](https://github.com/jertel/elastalert2/pull/1772) - @sauravnz
- [Elasticsearch] Add opt-in flag to allow queries on degraded Elasticsearch indices - [#1777](https://github.com/jertel/elastalert2/pull/1777) - @fhaldi

## Other changes
- [Helm] Omit enabled flag in probe definitions - [#1775](https://github.com/jertel/elastalert2/pull/1775) - @jim-barber-he
Expand Down
3 changes: 3 additions & 0 deletions docs/source/configuration.rst
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,9 @@ rule will no longer be run until either ElastAlert 2 restarts or the rule file h

``show_disabled_rules``: If true, ElastAlert 2 show the disable rules' list when finishes the execution. This defaults to True.

``allow_queries_on_degraded_indices``: If true, Elastalert 2 will still query Elasticsearch when some shards of the
data stream are unavailable. The default is ``false``.

``notify_alert``: List of alerters to execute upon encountering a system error. System errors occur when an unexpected exception is thrown during rule processing. For additional notifications, such as when ElastAlert 2 background tests encounter problems, or when connectivity to the data storage system is lost, enable ``notify_all_errors``.

See the :ref:`Alerts` section for the list of available alerters and their parameters.
Expand Down
7 changes: 6 additions & 1 deletion docs/source/ruletypes.rst
Original file line number Diff line number Diff line change
Expand Up @@ -379,14 +379,19 @@ ca_certs

``ca_certs``: Path to a CA cert bundle to use to verify SSL connections (Optional, string, no default)


disable_rules_on_error
^^^^^^^^^^^^^^^^^^^^^^

``disable_rules_on_error``: If true, ElastAlert 2 will disable rules which throw uncaught (not EAException) exceptions. It
will upload a traceback message to ``elastalert_metadata`` and if ``notify_email`` is set, send an email notification. The
rule will no longer be run until either ElastAlert 2 restarts or the rule file has been modified. This defaults to ``True``.

allow_queries_on_degraded_indices
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^

``allow_queries_on_degraded_indices``: If true, Elastalert 2 will still query Elasticsearch when some shards of the
data stream are unavailable. The default is ``false``.

es_conn_timeout
^^^^^^^^^^^^^^^

Expand Down
1 change: 1 addition & 0 deletions elastalert/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,7 @@ def load_conf(args, defaults=None, overrides=None):
conf.setdefault('scroll_keepalive', '30s')
conf.setdefault('max_scrolling_count', 990) # Avoid stack overflow in run_query, note that 1000 is Python's stack limit
conf.setdefault('disable_rules_on_error', True)
conf.setdefault('allow_queries_on_degraded_indices', False)
conf.setdefault('scan_subdirectories', True)
conf.setdefault('rules_loader', 'file')
conf.setdefault('custom_pretty_ts_format', None)
Expand Down
36 changes: 28 additions & 8 deletions elastalert/elastalert.py
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,7 @@ def __init__(self, args):
self.alert_time_limit = self.conf['alert_time_limit']
self.old_query_limit = self.conf['old_query_limit']
self.disable_rules_on_error = self.conf['disable_rules_on_error']
self.allow_queries_on_degraded_indices = self.conf['allow_queries_on_degraded_indices']
self.notify_email = self.conf.get('notify_email', [])
self.notify_all_errors = self.conf.get('notify_all_errors', False)
self.notify_alert = self.conf.get('notify_alert', [])
Expand Down Expand Up @@ -404,14 +405,33 @@ def get_hits(self, rule, starttime, endtime, index, scroll=False):

self.thread_data.total_hits = int(res['hits']['total']['value'])

if len(res.get('_shards', {}).get('failures', [])) > 0:
try:
errs = [e['reason']['reason'] for e in res['_shards']['failures'] if 'Failed to parse' in e['reason']['reason']]
if len(errs):
raise ElasticsearchException(errs)
except (TypeError, KeyError):
# Different versions of ES have this formatted in different ways. Fallback to str-ing the whole thing
raise ElasticsearchException(str(res['_shards']['failures']))
# Handles shard failures
if rule.get('allow_queries_on_degraded_indices', self.allow_queries_on_degraded_indices):
# Allow queries on degraded data streams
shard_failures = res.get('_shards', {}).get('failures', [])
if len(shard_failures) > 0:
reasons = [str(f) for f in shard_failures]
# Trigger Exception only if Elasticsearch query failed to parse
parse_errs = [r for r in reasons if 'Failed to parse' in r]
if parse_errs:
raise ElasticsearchException(parse_errs)
# If query still worked but returns some errors, continue...
failed_indices = ', '.join(sorted({str(f.get('index')) for f in shard_failures if isinstance(f, dict)}))
if len(failed_indices) > 1024:
failed_indices = failed_indices[:1024] + '... (%d characters removed)' % (len(failed_indices) - 1024)
elastalert_logger.warning(
'Rule %s: %d shard(s) failed, continuing with partial results. Affected indices: %s',
rule['name'], res['_shards'].get('failed', len(shard_failures)), failed_indices)
else:
# Stop rule if any error with at least one shard of a data stream
if len(res.get('_shards', {}).get('failures', [])) > 0:
try:
errs = [e['reason']['reason'] for e in res['_shards']['failures'] if 'Failed to parse' in e['reason']['reason']]
if len(errs):
raise ElasticsearchException(errs)
except (TypeError, KeyError):
# Different versions of ES have this formatted in different ways. Fallback to str-ing the whole thing
raise ElasticsearchException(str(res['_shards']['failures']))

elastalert_logger.debug(str(res))
except ElasticsearchException as e:
Expand Down
1 change: 1 addition & 0 deletions elastalert/schema.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -302,6 +302,7 @@ properties:
query_key: *arrayOfString
replace_dots_in_field_names: {type: boolean}
scan_entire_timeframe: {type: boolean}
allow_queries_on_degraded_indices: {type: boolean}

## System Error Notifications
notify_alert: *stringOrArrayOfStringsOrObjects
Expand Down
1 change: 1 addition & 0 deletions elastalert/test_rule.py
Original file line number Diff line number Diff line change
Expand Up @@ -449,6 +449,7 @@ def run_rule_test(self):
'old_query_limit': {'weeks': 1},
'run_every': {'minutes': 5},
'disable_rules_on_error': False,
'allow_queries_on_degraded_indices': False,
'buffer_time': {'minutes': 45},
'scroll_keepalive': '30s'
}
Expand Down
25 changes: 25 additions & 0 deletions tests/base_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,31 @@ def test_no_hits(ea):
assert ea.rules[0]['type'].add_data.call_count == 0


@pytest.mark.parametrize('set_in_rule', [False, True])
@pytest.mark.parametrize('allow_degraded, failures, expected_success', [
# All shards working: the option changes nothing
(False, [], True),
(True, [], True),
# Default: an unavailable shard stops the run
(False, [{'index': '.ds-index-test-00001', 'reason': {'type': 'no_shard_available_action_exception', 'reason': None}}], False),
# Option on: the run continues with the hits from the working shards
(True, [{'index': '.ds-index-test-00001', 'reason': {'type': 'no_shard_available_action_exception', 'reason': None}}], True),
(True, [{'index': '.ds-index-test-00001', 'reason': 'No shard available'}], True),
# Option on: a query that fails to parse stops the run (same behavior as default)
(True, [{'index': '.ds-index-test-00001', 'reason': {'type': 'query_shard_exception', 'reason': 'Failed to parse query'}}], False)])
def test_query_with_shard_failures(ea, set_in_rule, allow_degraded, failures, expected_success):
if set_in_rule:
# The rule's value must win over the global one, so give the global the opposite value
ea.allow_queries_on_degraded_indices = not allow_degraded
ea.rules[0]['allow_queries_on_degraded_indices'] = allow_degraded
else:
ea.allow_queries_on_degraded_indices = allow_degraded
response = generate_hits([START_TIMESTAMP, END_TIMESTAMP])
response['_shards'] = {'failures': failures}
ea.thread_data.current_es.search.return_value = response
assert ea.run_query(ea.rules[0], START, END) == expected_success


def test_no_terms_hits(ea):
ea.rules[0]['use_terms_query'] = True
ea.rules[0]['query_key'] = 'QWERTY'
Expand Down
1 change: 1 addition & 0 deletions tests/config_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ def test_config_loads():
assert conf['writeback_index'] == 'elastalert_status'

assert conf['alert_time_limit'] == datetime.timedelta(days=2)
assert conf['allow_queries_on_degraded_indices'] is False


def test_config_defaults():
Expand Down
1 change: 1 addition & 0 deletions tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -134,6 +134,7 @@ def ea():
'max_query_size': 10000,
'old_query_limit': datetime.timedelta(weeks=1),
'disable_rules_on_error': False,
'allow_queries_on_degraded_indices': False,
'scroll_keepalive': '30s',
'custom_pretty_ts_format': '%Y-%m-%d %H:%M'}
elastalert.util.elasticsearch_client = mock_es_client
Expand Down
Loading