|
13 | 13 | from django_rq import job |
14 | 14 |
|
15 | 15 | from dje.models import get_unsecured_manager |
| 16 | +from notification.models import fire_webhooks |
16 | 17 | from policy.engine import evaluate_rules |
17 | 18 |
|
18 | 19 | logger = logging.getLogger(__name__) |
19 | 20 |
|
20 | 21 |
|
| 22 | +def fire_policy_webhooks(product, new_violations, resolved_count): |
| 23 | + """Fire policy webhooks for newly detected or resolved violations.""" |
| 24 | + if new_violations: |
| 25 | + lines = [ |
| 26 | + f"- {violation.rule_label}: {violation.violation_count} violation(s)" |
| 27 | + for violation in new_violations |
| 28 | + ] |
| 29 | + payload = { |
| 30 | + "text": (f"[DejaCode] Policy violations detected for {product}\n" + "\n".join(lines)) |
| 31 | + } |
| 32 | + fire_webhooks("policy.violation_detected", instance=product, payload_override=payload) |
| 33 | + |
| 34 | + if resolved_count: |
| 35 | + payload = { |
| 36 | + "text": (f"[DejaCode] {resolved_count} policy violation(s) resolved for {product}") |
| 37 | + } |
| 38 | + fire_webhooks("policy.violation_resolved", instance=product, payload_override=payload) |
| 39 | + |
| 40 | + |
21 | 41 | @job |
22 | 42 | def evaluate_product_rules_task(product_uuid): |
23 | | - """Evaluate all active PolicyRules for the given product.""" |
| 43 | + """Evaluate all active policy rules for the given product and fire webhooks on changes.""" |
24 | 44 | Product = apps.get_model("product_portfolio", "product") |
25 | 45 |
|
26 | 46 | try: |
27 | | - product = Product.unsecured_objects.select_related("dataspace").get(uuid=product_uuid) |
| 47 | + product = get_unsecured_manager(Product).get(uuid=product_uuid) |
28 | 48 | except Product.DoesNotExist: |
29 | 49 | logger.error(f"evaluate_product_rules_task: product {product_uuid} not found, skipping.") |
30 | 50 | return |
31 | 51 |
|
32 | 52 | logger.info(f"Evaluating policy rules for product {product}") |
33 | | - violations = evaluate_rules(product) |
34 | | - logger.info(f"Policy rules evaluated for {product}: {len(violations)} active violation(s).") |
| 53 | + new_violations, resolved_count = evaluate_rules(product) |
| 54 | + logger.info( |
| 55 | + f"Policy rules evaluated for {product}: " |
| 56 | + f"{len(new_violations)} new violation(s), {resolved_count} resolved." |
| 57 | + ) |
| 58 | + fire_policy_webhooks(product, new_violations, resolved_count) |
35 | 59 |
|
36 | 60 |
|
37 | 61 | @job |
38 | 62 | def evaluate_all_products_rules_task(include_locked=False): |
39 | | - """Enqueue evaluate_product_rules_task for every product, skipping locked ones by default.""" |
| 63 | + """Evaluate policy rules for every product directly, skipping locked ones by default.""" |
40 | 64 | Product = apps.get_model("product_portfolio", "product") |
41 | 65 |
|
42 | | - products = Product.unsecured_objects.select_related("dataspace") |
| 66 | + products = get_unsecured_manager(Product).select_related("dataspace") |
43 | 67 | if not include_locked: |
44 | 68 | products = products.exclude_locked() |
45 | 69 |
|
46 | 70 | count = products.count() |
47 | | - logger.info(f"Queuing policy rule evaluation for {count} product(s).") |
| 71 | + logger.info(f"Starting policy rule evaluation for {count} product(s).") |
| 72 | + |
48 | 73 | for product in products: |
49 | | - evaluate_product_rules_task.delay(product_uuid=product.uuid) |
| 74 | + logger.info(f"Evaluating policy rules for product {product}") |
| 75 | + new_violations, resolved_count = evaluate_rules(product) |
| 76 | + logger.info( |
| 77 | + f"Policy rules evaluated for {product}: " |
| 78 | + f"{len(new_violations)} new violation(s), {resolved_count} resolved." |
| 79 | + ) |
| 80 | + fire_policy_webhooks(product, new_violations, resolved_count) |
| 81 | + |
| 82 | + logger.info(f"Policy rule evaluation complete for {count} product(s).") |
0 commit comments