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
26 changes: 16 additions & 10 deletions common/mqueue.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ class Partitioners:
fallback_partitioner = DefaultPartitioner()

@classmethod
def org_id_partitioner(cls, key: str, all_partitions: [], available_partitions: []):
def org_id_partitioner(cls, key: bytes, all_partitions: [], available_partitions: []):
# pylint: disable=broad-except
"""Kafka producer partitioner, expects org_id as a key and selects partition by modulo."""
org_id_hash = None
Expand Down Expand Up @@ -143,9 +143,11 @@ async def send_one(self, msg, key=None, headers=None):
try:
data = bytes(json.dumps(msg).encode("utf-8"))
res = await self.client.send_and_wait(self.topic, value=data, key=self._serialize_key(key), headers=headers)
LOGGER.debug(res)
LOGGER.debug("Sent message to Kafka topic %s: %s", self.topic, res)
except KafkaError:
self.connected = False
LOGGER.exception("Failed to send message to Kafka topic %s", self.topic)
raise

async def send_many(self, msg_list, key=None, headers=None):
"""Send list of messages"""
Expand All @@ -154,27 +156,31 @@ async def send_many(self, msg_list, key=None, headers=None):
for msg in msg_list:
data = bytes(json.dumps(msg).encode("utf-8"))
res = await self.client.send_and_wait(self.topic, value=data, key=self._serialize_key(key), headers=headers)
LOGGER.debug(res)
LOGGER.debug("Sent message to Kafka topic %s: %s", self.topic, res)
except KafkaError:
self.connected = False
LOGGER.exception("Failed to send message batch to Kafka topic %s (batch size %d)", self.topic, len(msg_list))
raise

async def send_raw(self, msg: bytes, key=None, headers=None):
"""Logic around sending raw message"""
await self.start()
try:
res = await self.client.send_and_wait(self.topic, value=msg, key=self._serialize_key(key), headers=headers)
LOGGER.debug(res)
LOGGER.debug("Sent raw message to Kafka topic %s: %s", self.topic, res)
except KafkaError:
self.connected = False
LOGGER.exception("Failed to send raw message to Kafka topic %s", self.topic)
raise

def send(self, msg, key=None, loop=None, headers=None):
async def send(self, msg, key=None, headers=None):
"""Sends a message"""
return asyncio.ensure_future(self.send_one(msg, key=key, headers=headers), loop=loop)
await self.send_one(msg, key=key, headers=headers)

def send_list(self, msgs, key=None, loop=None, headers=None):
async def send_list(self, msgs, key=None, headers=None):
"""Sends list of messages"""
return asyncio.ensure_future(self.send_many(msgs, key=key, headers=headers), loop=loop)
await self.send_many(msgs, key=key, headers=headers)

def send_bytes(self, msg: bytes, key=None, loop=None, headers=None):
async def send_bytes(self, msg: bytes, key=None, headers=None):
"""Sends a message, where message is already encoded to bytes"""
return asyncio.ensure_future(self.send_raw(msg, key=key, headers=headers), loop=loop)
await self.send_raw(msg, key=key, headers=headers)
25 changes: 16 additions & 9 deletions common/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@

import pytz
import requests
from aiokafka.errors import KafkaError
from dateutil.parser import isoparse
from prometheus_client import Counter
from psycopg2.extras import Json
Expand Down Expand Up @@ -107,7 +108,7 @@ def format_datetime(datetime_obj):
return str(datetime_obj) if datetime_obj else None


def send_msg_to_payload_tracker(producer, msg_dict, status, status_msg=None, loop=None, service="vulnerability"):
async def send_msg_to_payload_tracker(producer, msg_dict, status, status_msg=None, service="vulnerability"):
"""prepare and send message to payload-tracker"""
request_id = (msg_dict.get("platform_metadata", {}) or {}).get("request_id")
if not request_id:
Expand All @@ -123,11 +124,15 @@ 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)
LOGGER.debug("Sent message to topic %s: %s", producer.topic, str(tracking_payload))
LOGGER.debug("Sending message to topic %s: %s", producer.topic, str(tracking_payload))
try:
await producer.send(tracking_payload)
except KafkaError:
LOGGER.error("Sending Kafka producer.topic %s message failed", producer.topic)
pass # Work should compate and the caller should not treat the flow as failed, best-effort


def send_remediations_update(producer, inventory_id: str, cves: list, loop=None) -> None:
async def send_remediations_update(producer, inventory_id: str, cves: list) -> None:
"""
Send a message in format of remediations application updates using given Kafka producer

Expand All @@ -138,7 +143,7 @@ def send_remediations_update(producer, inventory_id: str, cves: list, loop=None)
loop (asyncio.event_loop, optional): asyncio event loop to be used by producer
"""
msg = {"host_id": inventory_id, "issues": ["vulnerabilities:{}".format(cve) for cve in cves]}
producer.send(msg, loop=loop)
await producer.send(msg)


def ensure_minimal_schema_version():
Expand Down Expand Up @@ -292,7 +297,9 @@ async def wrapper(*args, **kwargs):
return decorator


def send_notifications(notif_topic, new_sys_vulns, mit_sys_vulns, unmit_sys_vulns, rh_account_id, org_id, inventory_id, display_name=None):
async def send_notifications(
notif_topic, new_sys_vulns, mit_sys_vulns, unmit_sys_vulns, rh_account_id, org_id, inventory_id, display_name=None
):
"""Sends kafka message to notificator with system_vulnerabilities"""
if not new_sys_vulns and not mit_sys_vulns and not unmit_sys_vulns:
return
Expand All @@ -312,10 +319,10 @@ def send_notifications(notif_topic, new_sys_vulns, mit_sys_vulns, unmit_sys_vuln
],
}
LOGGER.debug("Sending evaluation result to notificator: %s", msg)
notif_topic.send(msg)
await notif_topic.send(msg)


def send_inventory_views(
async def send_inventory_views(
inventory_views_topic,
request_id,
inventory_id,
Expand Down Expand Up @@ -357,7 +364,7 @@ def send_inventory_views(
("request_id", bytes(request_id, "ascii")),
]
LOGGER.debug("Sending evaluation result to inventory views: %s", msg)
inventory_views_topic.send(msg, headers=headers)
await inventory_views_topic.send(msg, headers=headers)


def create_task_and_log(coro, logger, loop):
Expand Down
12 changes: 6 additions & 6 deletions evaluator/processor.py
Original file line number Diff line number Diff line change
Expand Up @@ -378,8 +378,8 @@ async def _evaluate_system(
cves_with_known_exploits += 1
await self._mark_system_evaluated(total_cves, system_platform, conn)

send_remediations_update(self.remediations_results, inventory_id, fixable_sys_vuln_rows)
send_notifications(
await send_remediations_update(self.remediations_results, inventory_id, fixable_sys_vuln_rows)
await send_notifications(
self.evaluator_results,
new_system_vulns,
[],
Expand All @@ -389,7 +389,7 @@ async def _evaluate_system(
inventory_id,
system_platform.display_name,
)
send_inventory_views(
await send_inventory_views(
self.inventory_views_results,
request_id,
inventory_id,
Expand All @@ -415,12 +415,12 @@ async def evaluate_system(
await self._evaluate_system(inventory_id, org_id, request_id, request_timestamp, recalc_event_id=recalc_event_id)
except EvaluatorException as ex:
LOGGER.error(str(ex))
send_msg_to_payload_tracker(self.payload_tracker, msg, "error", status_msg="evaluation failed", loop=self.loop)
await send_msg_to_payload_tracker(self.payload_tracker, msg, "error", status_msg="evaluation failed")
return
except VmaasErrorException as ex:
LOGGER.error(str(ex))
VMAAS_ERRORS_SKIP.inc()
send_msg_to_payload_tracker(self.payload_tracker, msg, "error", status_msg="evaluation failed", loop=self.loop)
await send_msg_to_payload_tracker(self.payload_tracker, msg, "error", status_msg="evaluation failed")
return

send_msg_to_payload_tracker(self.payload_tracker, msg, "success", status_msg="evaluation succeeded", loop=self.loop)
await send_msg_to_payload_tracker(self.payload_tracker, msg, "success", status_msg="evaluation succeeded")
10 changes: 3 additions & 7 deletions grouper/queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -161,14 +161,10 @@ 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(
self.payload_tracker, msg, "success", status_msg="unchanged system, not sending to evaluator", loop=self.loop
)
await send_msg_to_payload_tracker(self.payload_tracker, msg, "success", status_msg="unchanged system, not sending to evaluator")
return

CHANGED_SYSTEM.inc()
LOGGER.info("sending upload message to evaluator: %s, org_id: %s", inventory_id, org_id)
send_msg_to_payload_tracker(
self.payload_tracker, msg, "processing", status_msg="changed system, sending to evaluator", loop=self.loop
)
self.evaluator.send(msg)
await send_msg_to_payload_tracker(self.payload_tracker, msg, "processing", status_msg="changed system, sending to evaluator")
await self.evaluator.send(msg)
12 changes: 6 additions & 6 deletions listener/advisor_processor.py
Original file line number Diff line number Diff line change
Expand Up @@ -168,7 +168,7 @@ def _parse_input_metadata(self, msg: AdvisorMsg) -> (str, str, str):
msg.msg["input"]["timestamp"],
)

def _send_for_evaluation(
async def _send_for_evaluation(
self, org_id: str, inventory_id: str, request_id: str, reporter: str, timestamp: str, import_status: ImportStatus
):
"""Send message to evaluator to evaluate"""
Expand All @@ -185,11 +185,11 @@ def _send_for_evaluation(
},
"timestamp": timestamp,
}
self.grouper.send(msg, loop=self.loop, key=org_id)
await self.grouper.send(msg, key=org_id)

def _send_to_payload_tracker(self, status: str, msg: AdvisorMsg, message=None):
async 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)
await send_msg_to_payload_tracker(self.payload_tracker, msg.msg["input"], status, status_msg=message)

async def _process_upload(self, msg: AdvisorMsg):
"""Process message from advisor"""
Expand Down Expand Up @@ -221,8 +221,8 @@ 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")
await self._send_for_evaluation(org_id, inventory_id, request_id, reporter, timestamp, import_status)
await self._send_to_payload_tracker("received", msg, message="system received from advisor, sending to grouper")

async def process_msg(self, msg: AdvisorMsg):
"""Process single advisor msg"""
Expand Down
12 changes: 6 additions & 6 deletions listener/inventory_processor.py
Original file line number Diff line number Diff line change
Expand Up @@ -449,7 +449,7 @@ def _parse_input_metadata(self, msg: InventoryMsg) -> (str, str, str):
"""Extract reporter, request_id and timestamp from inventory message"""
return msg.msg["host"]["reporter"], (msg.msg.get("platform_metadata") or {}).get("request_id", ""), msg.msg["timestamp"]

def _send_for_evaluation(
async def _send_for_evaluation(
self, org_id: str, inventory_id: str, request_id: str, reporter: str, timestamp: str, import_status: ImportStatus
):
"""Send message to evaluator to evaluate"""
Expand All @@ -466,11 +466,11 @@ def _send_for_evaluation(
},
"timestamp": timestamp,
}
self.grouper.send(msg, loop=self.loop, key=org_id)
await self.grouper.send(msg, key=org_id)

def _send_to_payload_tracker(self, status: str, msg: InventoryMsg, message=None):
async 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)
await send_msg_to_payload_tracker(self.payload_tracker, msg.msg, status, status_msg=message)

async def _process_upload(self, msg: InventoryMsg):
"""Process upload message defined by QueueItem"""
Expand Down Expand Up @@ -544,8 +544,8 @@ 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")
await self._send_for_evaluation(org_id, inventory_id, request_id, reporter, timestamp, import_status)
await self._send_to_payload_tracker("received", msg, message="system received from inventory, sending to grouper")

async def _process_delete(self, msg: InventoryMsg):
"""Process inventory delete message"""
Expand Down
8 changes: 4 additions & 4 deletions listener/listener.py
Original file line number Diff line number Diff line change
Expand Up @@ -152,9 +152,9 @@ async def init(self):
"""Async constructor"""
await self.inventory_msg_processor.init()

def _send_payload_tracker_error(self, msg: dict, reason: str):
async def _send_payload_tracker_error(self, msg: dict, reason: str):
# since messages come separetely, inform separetly user about error in message
send_msg_to_payload_tracker(self.payload_tracker, msg, "error", status_msg=reason)
await send_msg_to_payload_tracker(self.payload_tracker, msg, "error", status_msg=reason)

async def _consume_inventory_msg(self, msg: dict) -> InventoryMsgType:
"""Consumes inventory message"""
Expand All @@ -163,7 +163,7 @@ async def _consume_inventory_msg(self, msg: dict) -> InventoryMsgType:
except InvalidInventoryMsg as exc:
LOGGER.error("obtained invalid inventory msg: %s", exc)
try:
self._send_payload_tracker_error(msg, "obtained invalid inventory msg")
await self._send_payload_tracker_error(msg, "obtained invalid inventory msg")
except KeyError:
pass
return InventoryMsgType.UNKNOWN
Expand All @@ -185,7 +185,7 @@ async def _consume_advisor_msg(self, msg: dict):
except InvalidAdvisorMsg as exc:
LOGGER.error("obtained invalid advisor msg: %s", exc)
try:
self._send_payload_tracker_error(msg["input"], "obtained invalid advisor msg")
await self._send_payload_tracker_error(msg["input"], "obtained invalid advisor msg")
except KeyError:
pass
return
Expand Down
26 changes: 16 additions & 10 deletions manager/admin_handler.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,14 @@
BATCH_SEMAPHORE = asyncio.BoundedSemaphore(CFG.re_evaluation_kafka_batches)


async def _send_recalc_batch_and_release(msgs):
"""Send a recalc batch to Kafka and release the batch semaphore"""
try:
await EVALUATOR_QUEUE.send_list(msgs)
finally:
BATCH_SEMAPHORE.release()


class TaskomaticRun(PutRequest):
"""PUT to /v1/taskomatic/run"""

Expand Down Expand Up @@ -433,17 +441,15 @@ def handle_delete(cls, **kwargs):
class RecalcBase:

@classmethod
def _create_kafka_msg_task(cls, rows, loop):
msgs = [
def _create_kafka_msg(cls, rows):
return [
{
"type": EvaluatorMessageType.RE_EVALUATE_SYSTEM,
"host": {"id": str(inventory_id), "org_id": org_id},
"timestamp": str(datetime.now(timezone.utc)),
}
for inventory_id, org_id in rows
]
task = EVALUATOR_QUEUE.send_list(msgs, loop=loop)
return task, len(msgs)


class RecalcAccounts(PutRequest, RecalcBase):
Expand Down Expand Up @@ -472,10 +478,9 @@ def handle_put(cls, **kwargs):
if not rows:
BATCH_SEMAPHORE.release()
break
task, msg_count = cls._create_kafka_msg_task(rows, loop)
total_scheduled += msg_count
task.add_done_callback(lambda x: BATCH_SEMAPHORE.release())
loop.run_until_complete(task)
msgs = cls._create_kafka_msg(rows)
loop.run_until_complete(_send_recalc_batch_and_release(msgs))
total_scheduled += len(msgs)

return f"{total_scheduled} systems scheduled for re-evaluation", 200

Expand Down Expand Up @@ -503,8 +508,9 @@ def handle_put(cls, **kwargs):
if not rows:
return f"{total_scheduled} systems scheduled for re-evaluation", 200

task, total_scheduled = cls._create_kafka_msg_task(rows, loop)
loop.run_until_complete(task)
msgs = cls._create_kafka_msg(rows)
loop.run_until_complete(EVALUATOR_QUEUE.send_list(msgs))
total_scheduled = len(msgs)

return f"{total_scheduled} systems scheduled for re-evaluation", 200

Expand Down
9 changes: 4 additions & 5 deletions notificator/notificator_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -118,7 +118,7 @@ def _create_notif_events(self, cve_id):
}
]

def _send_kafka_notif(
async def _send_kafka_notif(
self,
org_id: str,
inventory_id: str,
Comment thread
sourcery-ai[bot] marked this conversation as resolved.
Expand All @@ -127,7 +127,6 @@ def _send_kafka_notif(
events: list,
only_admins=False,
ignore_user_preferences=False,
loop=None,
):
"""
Sends kafka msg to notification kafka.
Expand Down Expand Up @@ -157,7 +156,7 @@ def _send_kafka_notif(
msg["org_id"] = org_id
msg["context"]["org_id"] = org_id
LOGGER.debug("Sending notification: %s", msg)
self.notifications_topic.send(msg, loop=loop)
await self.notifications_topic.send(msg)

async def _register_notified_acc(self, cve_id: int, notif_events: set(NotificationType), rh_account_id):
"""Registers new notified accounts into notified_accounts table"""
Expand Down Expand Up @@ -198,7 +197,7 @@ async def _process_normal_queue(self):
else:
# account-level: existing dedup logic
if not await self._is_already_notified(item.rh_account_id, item.cve_id, notif_event):
self._send_kafka_notif(item.org_id, inventory_id, display_name, notif_event.value, cve_events, loop=self.loop)
await self._send_kafka_notif(item.org_id, inventory_id, display_name, notif_event.value, cve_events)
LOGGER.info("Sent %s, cve=%s, org_id=%s", notif_event.value, item.cve, item.org_id)
SENT_NOTIFICATIONS.inc()
new_notified.add(notif_event)
Expand All @@ -213,7 +212,7 @@ async def _process_normal_queue(self):

# send batched system-level notifications (one message per system per notification type)
for (inventory_id, org_id, display_name, notif_event), events in system_notif_batch.items():
self._send_kafka_notif(org_id, inventory_id, display_name, notif_event.value, events, loop=self.loop)
await self._send_kafka_notif(org_id, inventory_id, display_name, notif_event.value, events)
LOGGER.info("Sent %s, inventory_id=%s, cve_count=%s", notif_event.value, inventory_id, len(events))
SENT_NOTIFICATIONS.inc()

Expand Down
Loading
Loading