From bb41e906d68ea5ead7d678d6e266c8f074e6f41f Mon Sep 17 00:00:00 2001 From: Lukas Holecek Date: Oct 06 2020 15:48:09 +0000 Subject: Refactor making decision Moves parsing request data and processing decision to another submodule. Consumers call the function directly instead making a full decision request. Signed-off-by: Lukas Holecek --- diff --git a/greenwave/api_v1.py b/greenwave/api_v1.py index 6a7f8f4..2f94a56 100644 --- a/greenwave/api_v1.py +++ b/greenwave/api_v1.py @@ -1,76 +1,29 @@ # SPDX-License-Identifier: GPL-2.0+ import logging -import datetime import random from flask import Blueprint, request, current_app, jsonify, url_for, redirect, Response -from werkzeug.exceptions import BadRequest, NotFound, UnsupportedMediaType +from werkzeug.exceptions import BadRequest from prometheus_client import generate_latest from greenwave import __version__ -from greenwave.policies import (summarize_answers, - RemotePolicy, - OnDemandPolicy, - _missing_decision_contexts_in_parent_policies) -from greenwave.resources import ResultsRetriever, WaiversRetriever +from greenwave.policies import ( + RemotePolicy, + _missing_decision_contexts_in_parent_policies, +) from greenwave.safe_yaml import SafeYAMLError from greenwave.utils import insert_headers, jsonp -from greenwave.waivers import waive_answers -from greenwave.subjects.factory import ( - create_subject, - create_subject_from_data, - UnknownSubjectDataError, -) from greenwave.monitor import ( registry, decision_exception_counter, decision_request_duration_seconds, ) +import greenwave.decision api = (Blueprint('api_v1', __name__)) log = logging.getLogger(__name__) -def _decision_subject(data): - try: - subject = create_subject_from_data(data) - except UnknownSubjectDataError: - log.info('Could not detect subject_identifier.') - raise BadRequest('Could not detect subject_identifier.') - - return subject - - -def _decision_subjects_for_request(data): - """ - Greenwave < 0.8 accepted a list of arbitrary dicts for the 'subject'. - Now we expect a specific type and identifier. - This maps from the old style to the new, for backwards compatibility. - - Note that WaiverDB has a very similar helper function, for compatibility - with WaiverDB < 0.11, but it accepts a single subject dict. Here we accept - a list. - """ - if 'subject' in data: - subjects = data['subject'] - if (not isinstance(subjects, list) or not subjects or - not all(isinstance(entry, dict) for entry in subjects)): - log.error('Invalid subject, must be a list of dicts') - raise BadRequest('Invalid subject, must be a list of dicts') - - for subject in subjects: - yield _decision_subject(subject) - else: - if 'subject_type' not in data: - log.error('Missing required "subject_type" parameter') - raise BadRequest('Missing required "subject_type" parameter') - if 'subject_identifier' not in data: - log.error('Missing required "subject_identifier" parameter') - raise BadRequest('Missing required "subject_identifier" parameter') - - yield create_subject(data['subject_type'], data['subject_identifier']) - - @api.route('/version', methods=['GET']) def version(): """ @@ -403,137 +356,7 @@ def make_decision(): :statuscode 504: Timeout while querying an upstream """ # noqa: E501 data = request.get_json() - if data: - if not data.get('product_version'): - log.error('Missing required product version') - raise BadRequest('Missing required product version') - if not data.get('decision_context') and not data.get('rules'): - log.error('Either decision_context or rules is required.') - raise BadRequest('Either decision_context or rules is required.') - else: - log.error('No JSON payload in request') - raise UnsupportedMediaType('No JSON payload in request') - - log.debug('New decision request for data: %s', data) - product_version = data['product_version'] - - decision_context = data.get('decision_context', None) - rules = data.get('rules', []) - if decision_context and rules: - log.error('Cannot have both decision_context and rules') - raise BadRequest('Cannot have both decision_context and rules') - - on_demand_policies = [] - if rules: - request_data = {key: data[key] for key in data if key not in ('subject', 'subject_type')} - for subject in _decision_subjects_for_request(data): - request_data['subject_type'] = subject.type - request_data['subject_identifier'] = subject.identifier - on_demand_policy = OnDemandPolicy.create_from_json(request_data) - on_demand_policies.append(on_demand_policy) - - verbose = data.get('verbose', False) - if not isinstance(verbose, bool): - log.error('Invalid verbose flag, must be a bool') - raise BadRequest('Invalid verbose flag, must be a bool') - ignore_results = data.get('ignore_result', []) - ignore_waivers = data.get('ignore_waiver', []) - when = data.get('when') - - if when: - try: - datetime.datetime.strptime(when, '%Y-%m-%dT%H:%M:%S.%f') - except ValueError: - raise BadRequest('Invalid "when" parameter, must be in ISO8601 format') - - answers = [] - verbose_results = [] - applicable_policies = [] - retriever_args = {'when': when} - results_retriever = ResultsRetriever( - ignore_ids=ignore_results, - url=current_app.config['RESULTSDB_API_URL'], - **retriever_args) - waivers_retriever = WaiversRetriever( - ignore_ids=ignore_waivers, - url=current_app.config['WAIVERDB_API_URL'], - **retriever_args) - waiver_filters = [] - - policies = on_demand_policies or current_app.config['policies'] - for subject in _decision_subjects_for_request(data): - subject_policies = [ - policy for policy in policies - if policy.matches( - decision_context=decision_context, - product_version=product_version, - subject=subject) - ] - - if not subject_policies: - if subject.ignore_missing_policy: - continue - - log.error( - 'Cannot find any applicable policies for %s subjects at gating point %s in %s', - subject.type, decision_context, product_version) - raise NotFound( - 'Cannot find any applicable policies for %s subjects at gating point %s in %s' % ( - subject.type, decision_context, product_version)) - - if verbose: - # Retrieve test results and waivers for all items when verbose output is requested. - verbose_results.extend(results_retriever.retrieve(subject)) - waiver_filters.append(dict( - subject_type=subject.type, - subject_identifier=subject.identifier, - product_version=product_version, - )) - - for policy in subject_policies: - answers.extend( - policy.check( - product_version, - subject, - results_retriever)) - - applicable_policies.extend(subject_policies) - - if not verbose: - for answer in answers: - if not answer.is_satisfied: - waiver_filters.append(dict( - subject_type=answer.subject.type, - subject_identifier=answer.subject.identifier, - product_version=product_version, - testcase=answer.test_case_name, - )) - - if waiver_filters: - waivers = waivers_retriever.retrieve(waiver_filters) - else: - waivers = [] - answers = waive_answers(answers, waivers) - - response = { - 'policies_satisfied': all(answer.is_satisfied for answer in answers), - 'summary': summarize_answers(answers), - 'satisfied_requirements': - [answer.to_json() for answer in answers if answer.is_satisfied], - 'unsatisfied_requirements': - [answer.to_json() for answer in answers if not answer.is_satisfied] - } - - # Check if on-demand policy was specified - if not rules: - response.update({'applicable_policies': [policy.id for policy in applicable_policies]}) - - if verbose: - response.update({ - 'results': list({result['id']: result for result in verbose_results}.values()), - 'waivers': list({waiver['id']: waiver for waiver in waivers}.values()), - }) - + response = greenwave.decision.make_decision(data, current_app.config) log.debug('Response: %s', response) resp = jsonify(response) resp = insert_headers(resp) diff --git a/greenwave/consumers/consumer.py b/greenwave/consumers/consumer.py index 165d3ab..3003898 100644 --- a/greenwave/consumers/consumer.py +++ b/greenwave/consumers/consumer.py @@ -6,6 +6,8 @@ import requests import fedmsg.consumers import greenwave.app_factory +import greenwave.decision + from greenwave.monitor import ( publish_decision_exceptions_result_counter, messaging_tx_to_send_counter, messaging_tx_stopped_counter, @@ -13,8 +15,6 @@ from greenwave.monitor import ( from greenwave.policies import applicable_decision_context_product_version_pairs from greenwave.utils import right_before_this_time -import greenwave.resources - try: import fedora_messaging.api import fedora_messaging.exceptions @@ -161,10 +161,10 @@ class Consumer(fedmsg.consumers.FedmsgConsumer): log.debug('querying greenwave at: %s', greenwave_url) try: - decision = greenwave.resources.retrieve_decision(greenwave_url, request_data) + decision = greenwave.decision.make_decision(request_data, self.flask_app.config) request_data['when'] = right_before_this_time(submit_time) - old_decision = greenwave.resources.retrieve_decision(greenwave_url, request_data) + old_decision = greenwave.decision.make_decision(request_data, self.flask_app.config) log.debug('old decision: %s', old_decision) except requests.exceptions.HTTPError as e: log.exception('Failed to retrieve decision for data=%s, error: %s', request_data, e) diff --git a/greenwave/decision.py b/greenwave/decision.py new file mode 100644 index 0000000..52296a4 --- /dev/null +++ b/greenwave/decision.py @@ -0,0 +1,199 @@ +# SPDX-License-Identifier: GPL-2.0+ +import logging +import datetime + +from werkzeug.exceptions import ( + BadRequest, + NotFound, + UnsupportedMediaType, +) + +from greenwave.policies import ( + summarize_answers, + OnDemandPolicy, +) +from greenwave.resources import ResultsRetriever, WaiversRetriever +from greenwave.subjects.factory import ( + create_subject, + create_subject_from_data, + UnknownSubjectDataError, +) +from greenwave.waivers import waive_answers + +log = logging.getLogger(__name__) + + +def _decision_subject(data): + try: + subject = create_subject_from_data(data) + except UnknownSubjectDataError: + log.info('Could not detect subject_identifier.') + raise BadRequest('Could not detect subject_identifier.') + + return subject + + +def _decision_subjects_for_request(data): + """ + Greenwave < 0.8 accepted a list of arbitrary dicts for the 'subject'. + Now we expect a specific type and identifier. + This maps from the old style to the new, for backwards compatibility. + + Note that WaiverDB has a very similar helper function, for compatibility + with WaiverDB < 0.11, but it accepts a single subject dict. Here we accept + a list. + """ + if 'subject' in data: + subjects = data['subject'] + if (not isinstance(subjects, list) or not subjects or + not all(isinstance(entry, dict) for entry in subjects)): + log.error('Invalid subject, must be a list of dicts') + raise BadRequest('Invalid subject, must be a list of dicts') + + for subject in subjects: + yield _decision_subject(subject) + else: + if 'subject_type' not in data: + log.error('Missing required "subject_type" parameter') + raise BadRequest('Missing required "subject_type" parameter') + if 'subject_identifier' not in data: + log.error('Missing required "subject_identifier" parameter') + raise BadRequest('Missing required "subject_identifier" parameter') + + yield create_subject(data['subject_type'], data['subject_identifier']) + + +def make_decision(data, config): + if not data: + log.error('No JSON payload in request') + raise UnsupportedMediaType('No JSON payload in request') + + if not data.get('product_version'): + log.error('Missing required product version') + raise BadRequest('Missing required product version') + + if not data.get('decision_context') and not data.get('rules'): + log.error('Either decision_context or rules is required.') + raise BadRequest('Either decision_context or rules is required.') + + log.debug('New decision request for data: %s', data) + product_version = data['product_version'] + + decision_context = data.get('decision_context', None) + rules = data.get('rules', []) + if decision_context and rules: + log.error('Cannot have both decision_context and rules') + raise BadRequest('Cannot have both decision_context and rules') + + on_demand_policies = [] + if rules: + request_data = {key: data[key] for key in data if key not in ('subject', 'subject_type')} + for subject in _decision_subjects_for_request(data): + request_data['subject_type'] = subject.type + request_data['subject_identifier'] = subject.identifier + on_demand_policy = OnDemandPolicy.create_from_json(request_data) + on_demand_policies.append(on_demand_policy) + + verbose = data.get('verbose', False) + if not isinstance(verbose, bool): + log.error('Invalid verbose flag, must be a bool') + raise BadRequest('Invalid verbose flag, must be a bool') + ignore_results = data.get('ignore_result', []) + ignore_waivers = data.get('ignore_waiver', []) + when = data.get('when') + + if when: + try: + datetime.datetime.strptime(when, '%Y-%m-%dT%H:%M:%S.%f') + except ValueError: + raise BadRequest('Invalid "when" parameter, must be in ISO8601 format') + + answers = [] + verbose_results = [] + applicable_policies = [] + retriever_args = {'when': when} + results_retriever = ResultsRetriever( + ignore_ids=ignore_results, + url=config['RESULTSDB_API_URL'], + **retriever_args) + waivers_retriever = WaiversRetriever( + ignore_ids=ignore_waivers, + url=config['WAIVERDB_API_URL'], + **retriever_args) + waiver_filters = [] + + policies = on_demand_policies or config['policies'] + for subject in _decision_subjects_for_request(data): + subject_policies = [ + policy for policy in policies + if policy.matches( + decision_context=decision_context, + product_version=product_version, + subject=subject) + ] + + if not subject_policies: + if subject.ignore_missing_policy: + continue + + log.error( + 'Cannot find any applicable policies for %s subjects at gating point %s in %s', + subject.type, decision_context, product_version) + raise NotFound( + 'Cannot find any applicable policies for %s subjects at gating point %s in %s' % ( + subject.type, decision_context, product_version)) + + if verbose: + # Retrieve test results and waivers for all items when verbose output is requested. + verbose_results.extend(results_retriever.retrieve(subject)) + waiver_filters.append(dict( + subject_type=subject.type, + subject_identifier=subject.identifier, + product_version=product_version, + )) + + for policy in subject_policies: + answers.extend( + policy.check( + product_version, + subject, + results_retriever)) + + applicable_policies.extend(subject_policies) + + if not verbose: + for answer in answers: + if not answer.is_satisfied: + waiver_filters.append(dict( + subject_type=answer.subject.type, + subject_identifier=answer.subject.identifier, + product_version=product_version, + testcase=answer.test_case_name, + )) + + if waiver_filters: + waivers = waivers_retriever.retrieve(waiver_filters) + else: + waivers = [] + answers = waive_answers(answers, waivers) + + response = { + 'policies_satisfied': all(answer.is_satisfied for answer in answers), + 'summary': summarize_answers(answers), + 'satisfied_requirements': + [answer.to_json() for answer in answers if answer.is_satisfied], + 'unsatisfied_requirements': + [answer.to_json() for answer in answers if not answer.is_satisfied] + } + + # Check if on-demand policy was specified + if not rules: + response.update({'applicable_policies': [policy.id for policy in applicable_policies]}) + + if verbose: + response.update({ + 'results': list({result['id']: result for result in verbose_results}.values()), + 'waivers': list({waiver['id']: waiver for waiver in waivers}.values()), + }) + + return response diff --git a/greenwave/resources.py b/greenwave/resources.py index a55cb96..caf7290 100644 --- a/greenwave/resources.py +++ b/greenwave/resources.py @@ -183,10 +183,3 @@ def retrieve_yaml_remote_rule(url): response = requests_session.request('GET', url) response.raise_for_status() return response.content - - -# NOTE - not cached. -def retrieve_decision(greenwave_url, data): - response = requests_session.post(greenwave_url, json=data) - response.raise_for_status() - return response.json() diff --git a/greenwave/tests/test_resultsdb_consumer.py b/greenwave/tests/test_resultsdb_consumer.py index fa8a9a9..339b749 100644 --- a/greenwave/tests/test_resultsdb_consumer.py +++ b/greenwave/tests/test_resultsdb_consumer.py @@ -12,6 +12,36 @@ from greenwave.policies import Policy from greenwave.subjects.factory import create_subject +@pytest.fixture(autouse=True) +def mock_retrieve_decision(): + with mock.patch('greenwave.decision.make_decision') as mocked: + def retrieve_decision(data, config): + #pylint: disable=unused-argument + if 'when' in data: + return None + return {} + mocked.side_effect = retrieve_decision + yield mocked + + +@pytest.fixture +def mock_retrieve_results(): + with mock.patch('greenwave.resources.ResultsRetriever.retrieve') as mocked: + yield mocked + + +@pytest.fixture +def mock_retrieve_scm_from_koji(): + with mock.patch('greenwave.resources.retrieve_scm_from_koji') as mocked: + yield mocked + + +@pytest.fixture +def mock_retrieve_yaml_remote_rule(): + with mock.patch('greenwave.resources.retrieve_yaml_remote_rule') as mocked: + yield mocked + + def announcement_subject(message): cls = greenwave.consumers.resultsdb.ResultsDBHandler @@ -119,14 +149,9 @@ parameters = [ @pytest.mark.parametrize("config,publish", parameters) -@mock.patch('greenwave.resources.ResultsRetriever.retrieve') -@mock.patch('greenwave.resources.retrieve_decision') -@mock.patch('greenwave.resources.retrieve_scm_from_koji') -@mock.patch('greenwave.resources.retrieve_yaml_remote_rule') def test_remote_rule_decision_change( mock_retrieve_yaml_remote_rule, mock_retrieve_scm_from_koji, - mock_retrieve_decision, mock_retrieve_results, config, publish): @@ -166,12 +191,6 @@ def test_remote_rule_decision_change( } mock_retrieve_results.return_value = [result] - def retrieve_decision(url, data): - #pylint: disable=unused-argument - if 'when' in data: - return None - return {} - mock_retrieve_decision.side_effect = retrieve_decision mock_retrieve_scm_from_koji.return_value = ('rpms', nvr, 'c3c47a08a66451cb9686c49f040776ed35a0d1bb') @@ -227,14 +246,9 @@ def test_remote_rule_decision_change( @pytest.mark.parametrize("config,publish", parameters) -@mock.patch('greenwave.resources.ResultsRetriever.retrieve') -@mock.patch('greenwave.resources.retrieve_decision') -@mock.patch('greenwave.resources.retrieve_scm_from_koji') -@mock.patch('greenwave.resources.retrieve_yaml_remote_rule') def test_remote_rule_decision_change_not_matching( mock_retrieve_yaml_remote_rule, mock_retrieve_scm_from_koji, - mock_retrieve_decision, mock_retrieve_results, config, publish): @@ -274,12 +288,6 @@ def test_remote_rule_decision_change_not_matching( } mock_retrieve_results.return_value = [result] - def retrieve_decision(url, data): - #pylint: disable=unused-argument - if 'when' in data: - return None - return {} - mock_retrieve_decision.side_effect = retrieve_decision mock_retrieve_scm_from_koji.return_value = ('rpms', nvr, 'c3c47a08a66451cb9686c49f040776ed35a0d1bb') @@ -374,14 +382,9 @@ def test_guess_product_version_failure(nvr): @pytest.mark.parametrize("config,publish", parameters) -@mock.patch('greenwave.resources.ResultsRetriever.retrieve') -@mock.patch('greenwave.resources.retrieve_decision') -@mock.patch('greenwave.resources.retrieve_scm_from_koji') -@mock.patch('greenwave.resources.retrieve_yaml_remote_rule') def test_decision_change_for_modules( mock_retrieve_yaml_remote_rule, mock_retrieve_scm_from_koji, - mock_retrieve_decision, mock_retrieve_results, config, publish): @@ -425,12 +428,6 @@ def test_decision_change_for_modules( } mock_retrieve_results.return_value = [result] - def retrieve_decision(url, data): - #pylint: disable=unused-argument - if 'when' in data: - return None - return {} - mock_retrieve_decision.side_effect = retrieve_decision mock_retrieve_scm_from_koji.return_value = ('modules', nsvc, '97273b80dd568bd15f9636b695f6001ecadb65e0') @@ -485,11 +482,7 @@ def test_decision_change_for_modules( } -@mock.patch('greenwave.resources.ResultsRetriever.retrieve') -@mock.patch('greenwave.resources.retrieve_decision') -def test_real_fedora_messaging_msg( - mock_retrieve_decision, - mock_retrieve_results): +def test_real_fedora_messaging_msg(mock_retrieve_results): message = { 'msg': { 'task': { @@ -573,13 +566,6 @@ def test_real_fedora_messaging_msg( } mock_retrieve_results.return_value = [result] - def retrieve_decision(url, data): - #pylint: disable=unused-argument - if 'when' in data: - return None - return {} - mock_retrieve_decision.side_effect = retrieve_decision - hub = mock.MagicMock() hub.config = { 'environment': 'environment', @@ -610,11 +596,7 @@ def test_real_fedora_messaging_msg( } -@mock.patch('greenwave.resources.ResultsRetriever.retrieve') -@mock.patch('greenwave.resources.retrieve_decision') -def test_container_brew_build( - mock_retrieve_decision, - mock_retrieve_results): +def test_container_brew_build(mock_retrieve_results): message = { 'msg': { 'submit_time': '2019-08-27T13:57:53.490376', @@ -653,13 +635,6 @@ def test_container_brew_build( } mock_retrieve_results.return_value = [result] - def retrieve_decision(url, data): - #pylint: disable=unused-argument - if 'when' in data: - return None - return {} - mock_retrieve_decision.side_effect = retrieve_decision - hub = mock.MagicMock() hub.config = { 'environment': 'environment',