Skip to content
Open
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
3 changes: 2 additions & 1 deletion common/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -123,8 +123,9 @@ def send_msg_to_payload_tracker(producer, msg_dict, status, status_msg=None, loo
}
if status_msg:
tracking_payload["status_msg"] = status_msg
producer.send(tracking_payload, loop=loop)
fut = producer.send(tracking_payload, loop=loop)
LOGGER.debug("Sent message to topic %s: %s", producer.topic, str(tracking_payload))
return fut


def send_remediations_update(producer, inventory_id: str, cves: list, loop=None) -> None:
Expand Down
1 change: 1 addition & 0 deletions grouper/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ class QueueItem:
request_id: str

otel_context: otel_context_api.Context | None = None
topic: str = ""
second_upload_event: asyncio.Event = None

def __post_init__(self):
Expand Down
50 changes: 20 additions & 30 deletions grouper/grouper.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,6 @@
from common.telemetry import instrument_outbound_http
from common.telemetry import instrument_psycopg
from common.telemetry import instrument_psycopg2
from common.telemetry import use_otel_context
from common.utils import create_task_and_log

from .common import CFG
Expand Down Expand Up @@ -85,35 +84,26 @@ async def consume_message(self, msg: ConsumerRecord, unlock: asyncio.BoundedSema

parent_ctx = extract_context_from_msg_headers(msg.headers)

with use_otel_context(parent_ctx):
with TRACER.start_as_current_span(
f"process {msg.topic}",
attributes={
"messaging.system": "kafka",
"messaging.operation.name": "process",
"messaging.destination.name": msg.topic,
"vulnerability.inventory_id": inventory_id,
"vulnerability.msg_type": msg_type.value,
},
):
if msg_type is GrouperMessageType.INVENTORY_UPLOAD:
self.queue.push_inventory_msg(
org_id,
inventory_id,
reporter,
changed,
(msg_dict.get("platform_metadata") or {}).get("request_id", ""),
otel_context=parent_ctx,
)
elif msg_type is GrouperMessageType.ADVISOR_UPLOAD:
self.queue.push_advisor_msg(
org_id,
inventory_id,
reporter,
changed,
(msg_dict.get("platform_metadata") or {}).get("request_id", ""),
otel_context=parent_ctx,
)
if msg_type is GrouperMessageType.INVENTORY_UPLOAD:
self.queue.push_inventory_msg(
org_id,
inventory_id,
reporter,
changed,
(msg_dict.get("platform_metadata") or {}).get("request_id", ""),
otel_context=parent_ctx,
topic=msg.topic,
)
elif msg_type is GrouperMessageType.ADVISOR_UPLOAD:
self.queue.push_advisor_msg(
org_id,
inventory_id,
reporter,
changed,
(msg_dict.get("platform_metadata") or {}).get("request_id", ""),
otel_context=parent_ctx,
topic=msg.topic,
)

async def _start_grouping_inventory(self) -> None:
"""Start of the grouping inventory uploads"""
Expand Down
115 changes: 76 additions & 39 deletions grouper/queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,10 +9,13 @@
from typing import Dict

from opentelemetry import context as otel_context_api
from opentelemetry import trace

from common.constants import EvaluatorMessageType
from common.logging import get_logger
from common.mqueue import MQWriter
from common.telemetry import get_tracer
from common.telemetry import threadctx
from common.telemetry import use_otel_context
from common.utils import create_task_and_log
from common.utils import send_msg_to_payload_tracker
Expand All @@ -26,9 +29,11 @@
from .common import QUEUE_SIZE
from .common import UNCHANGED_SYSTEM
from .common import BoundedSemaphorePrometheus
from .common import GrouperMessageType
from .common import QueueItem

LOGGER = get_logger(__name__)
TRACER = get_tracer(__name__)


class GrouperQueue:
Expand Down Expand Up @@ -64,14 +69,15 @@ def push_inventory_msg(
inventory_changed: bool,
request_id: str,
otel_context: otel_context_api.Context | None = None,
topic: str = "",
) -> None:
"""Push inventory upload message to queue"""
is_updated = False

item = self._queue.get(inventory_id)
if not item:
LOGGER.info("pushing listener upload to queue for system: %s, org_id: %s", inventory_id, org_id)
item = QueueItem(True, inventory_changed, False, False, request_id, otel_context=otel_context)
item = QueueItem(True, inventory_changed, False, False, request_id, otel_context=otel_context, topic=topic)
self._queue[inventory_id] = item
else:
is_updated = True
Expand All @@ -86,6 +92,8 @@ def push_inventory_msg(
item.request_id = request_id
if otel_context is not None:
item.otel_context = otel_context
if topic and not item.topic:
item.topic = topic

if item.inventory_upload and item.advisor_upload:
LOGGER.info("obtained both uploads for system: %s, account: %s, releasing lock", inventory_id, org_id)
Expand All @@ -102,14 +110,15 @@ def push_advisor_msg(
advisor_changed: bool,
request_id: str,
otel_context: otel_context_api.Context | None = None,
topic: str = "",
) -> None:
"""Push advisor message to queue"""
is_updated = False

item = self._queue.get(inventory_id)
if not item:
LOGGER.info("pushing advisor upload to queue for system: %s, org_id: %s", inventory_id, org_id)
item = QueueItem(False, False, True, advisor_changed, request_id, otel_context=otel_context)
item = QueueItem(False, False, True, advisor_changed, request_id, otel_context=otel_context, topic=topic)
self._queue[inventory_id] = item
else:
is_updated = True
Expand All @@ -124,6 +133,8 @@ def push_advisor_msg(
item.request_id = request_id
if otel_context is not None and item.otel_context is None:
item.otel_context = otel_context
if topic and not item.topic:
item.topic = topic

if item.inventory_upload and item.advisor_upload:
LOGGER.info("obtained both messages for system: %s, account: %s, releasing lock", inventory_id, org_id)
Expand All @@ -135,40 +146,61 @@ def push_advisor_msg(
async def _start_item_processing(self, org_id: str, inventory_id: str, reporter: str) -> None:
"""Single queue item waiting coroutine"""
item = self._queue[inventory_id]
topic = item.topic or (CFG.grouper_inventory_topic if item.inventory_upload else CFG.grouper_advisor_topic)
msg_type = GrouperMessageType.INVENTORY_UPLOAD.value if item.inventory_upload else GrouperMessageType.ADVISOR_UPLOAD.value

# RHSM systems do not need to wait for advisor message
if reporter == "rhsm-system-profile-bridge":
LOGGER.debug(
"reporter %s skipped waiting for %s message for system: %s, account: %s",
reporter,
"inventory" if item.advisor_upload else "advisor",
inventory_id,
org_id,
)
else:
LOGGER.debug(
"starting waiting for %s msg for system: %s, account: %s",
"inventory" if item.advisor_upload else "advisor",
inventory_id,
org_id,
)

try:
await asyncio.wait_for(item.second_upload_event.wait(), CFG.grouper_messages_timeout_sec)
PAIR_HIT.inc()
except asyncio.TimeoutError:
LOGGER.info("timing out while waiting for message for system: %s, account: %s", inventory_id, org_id)
PAIR_MISS.inc()
self._queue.pop(inventory_id, None)

if item.inventory_upload:
self.max_inventory_msgs.release()
if item.advisor_upload:
self.max_advisor_msgs.release()

await self._send_for_evaluation(item, org_id, inventory_id)

QUEUE_SIZE.dec()
with use_otel_context(item.otel_context):
with TRACER.start_as_current_span(
f"process {topic}",
kind=trace.SpanKind.CONSUMER,
attributes={
"messaging.system": "kafka",
"messaging.operation.name": "process",
"messaging.destination.name": topic,
"vulnerability.inventory_id": inventory_id,
"vulnerability.msg_type": msg_type,
},
) as span:
if org_id:
span.set_attribute("rh.org_id", org_id)
threadctx.org_id = org_id
if item.request_id:
span.set_attribute("rh.request_id", item.request_id)
threadctx.request_id = item.request_id

# RHSM systems do not need to wait for advisor message
if reporter == "rhsm-system-profile-bridge":
LOGGER.debug(
"reporter %s skipped waiting for %s message for system: %s, account: %s",
reporter,
"inventory" if item.advisor_upload else "advisor",
inventory_id,
org_id,
)
else:
LOGGER.debug(
"starting waiting for %s msg for system: %s, account: %s",
"inventory" if item.advisor_upload else "advisor",
inventory_id,
org_id,
)

try:
await asyncio.wait_for(item.second_upload_event.wait(), CFG.grouper_messages_timeout_sec)
PAIR_HIT.inc()
except asyncio.TimeoutError:
LOGGER.info("timing out while waiting for message for system: %s, account: %s", inventory_id, org_id)
PAIR_MISS.inc()
self._queue.pop(inventory_id, None)

if item.inventory_upload:
self.max_inventory_msgs.release()
if item.advisor_upload:
self.max_advisor_msgs.release()

await self._send_for_evaluation(item, org_id, inventory_id)

QUEUE_SIZE.dec()

async def _send_for_evaluation(self, item: QueueItem, org_id: str, inventory_id: str) -> None:
"""Sends message to evaluate a system"""
Expand All @@ -184,15 +216,20 @@ async def _send_for_evaluation(self, item: QueueItem, org_id: str, inventory_id:
if (not item.inventory_changed and not item.advisor_changed) and not CFG.disable_optimisation:
UNCHANGED_SYSTEM.inc()
LOGGER.info("skipping evaluation, system not changed: %s, org_id: %s", inventory_id, org_id)
send_msg_to_payload_tracker(
tracker_fut = send_msg_to_payload_tracker(
self.payload_tracker, msg, "success", status_msg="unchanged system, not sending to evaluator", loop=self.loop
)
if tracker_fut:
await tracker_fut
return

CHANGED_SYSTEM.inc()
LOGGER.info("sending upload message to evaluator: %s, org_id: %s", inventory_id, org_id)
send_msg_to_payload_tracker(
tracker_fut = send_msg_to_payload_tracker(
self.payload_tracker, msg, "processing", status_msg="changed system, sending to evaluator", loop=self.loop
)
with use_otel_context(item.otel_context):
self.evaluator.send(msg)
evaluator_fut = self.evaluator.send(msg)
if tracker_fut:
await tracker_fut
if evaluator_fut:
await evaluator_fut
12 changes: 8 additions & 4 deletions listener/advisor_processor.py
Original file line number Diff line number Diff line change
Expand Up @@ -185,11 +185,11 @@ def _send_for_evaluation(
},
"timestamp": timestamp,
}
self.grouper.send(msg, loop=self.loop, key=org_id)
return self.grouper.send(msg, loop=self.loop, key=org_id)

def _send_to_payload_tracker(self, status: str, msg: AdvisorMsg, message=None):
"""Send payload tracker message"""
send_msg_to_payload_tracker(self.payload_tracker, msg.msg["input"], status, status_msg=message, loop=self.loop)
return send_msg_to_payload_tracker(self.payload_tracker, msg.msg["input"], status, status_msg=message, loop=self.loop)

async def _process_upload(self, msg: AdvisorMsg):
"""Process message from advisor"""
Expand Down Expand Up @@ -221,8 +221,12 @@ async def _process_upload(self, msg: AdvisorMsg):
LOGGER.info(
"advisor data inserted, system: %s, org_id: %s, reporter: %s, request_id: %s", inventory_id, org_id, reporter, request_id
)
self._send_for_evaluation(org_id, inventory_id, request_id, reporter, timestamp, import_status)
self._send_to_payload_tracker("received", msg, message="system received from advisor, sending to grouper")
eval_fut = self._send_for_evaluation(org_id, inventory_id, request_id, reporter, timestamp, import_status)
tracker_fut = self._send_to_payload_tracker("received", msg, message="system received from advisor, sending to grouper")
if eval_fut:
await eval_fut
if tracker_fut:
await tracker_fut

async def process_msg(self, msg: AdvisorMsg):
"""Process single advisor msg"""
Expand Down
12 changes: 8 additions & 4 deletions listener/inventory_processor.py
Original file line number Diff line number Diff line change
Expand Up @@ -467,11 +467,11 @@ def _send_for_evaluation(
},
"timestamp": timestamp,
}
self.grouper.send(msg, loop=self.loop, key=org_id)
return self.grouper.send(msg, loop=self.loop, key=org_id)

def _send_to_payload_tracker(self, status: str, msg: InventoryMsg, message=None):
"""Send payload tracker message"""
send_msg_to_payload_tracker(self.payload_tracker, msg.msg, status, status_msg=message, loop=self.loop)
return send_msg_to_payload_tracker(self.payload_tracker, msg.msg, status, status_msg=message, loop=self.loop)

async def _process_upload(self, msg: InventoryMsg):
"""Process upload message defined by QueueItem"""
Expand Down Expand Up @@ -547,8 +547,12 @@ async def _process_upload(self, msg: InventoryMsg):
LOGGER.info(
"Inventory data processed, system: %s, org_id: %s, reporter: %s, request_id: %s", inventory_id, org_id, reporter, request_id
)
self._send_for_evaluation(org_id, inventory_id, request_id, reporter, timestamp, import_status)
self._send_to_payload_tracker("received", msg, message="system received from inventory, sending to grouper")
eval_fut = self._send_for_evaluation(org_id, inventory_id, request_id, reporter, timestamp, import_status)
tracker_fut = self._send_to_payload_tracker("received", msg, message="system received from inventory, sending to grouper")
if eval_fut:
await eval_fut
if tracker_fut:
await tracker_fut

async def _process_delete(self, msg: InventoryMsg):
"""Process inventory delete message"""
Expand Down
Loading
Loading