From 407ae24f04695885e75ad3663c1248981df7ba6a Mon Sep 17 00:00:00 2001 From: Prowler Bot Date: Fri, 17 Apr 2026 11:01:19 +0200 Subject: [PATCH] perf(attack-paths): cleanup task prioritization, restore default batch sizes to 1000, upgrade Cartography to 0.135.0 (#10768) Co-authored-by: Josema Camacho --- api/CHANGELOG.md | 4 + api/poetry.lock | 98 ++++++++++--------- api/pyproject.toml | 3 +- .../0090_attack_paths_cleanup_priority.py | 23 +++++ api/src/backend/config/celery.py | 4 +- .../backend/tasks/jobs/attack_paths/aws.py | 56 +++++++++-- .../backend/tasks/jobs/attack_paths/config.py | 4 +- .../tasks/jobs/attack_paths/findings.py | 16 +-- .../tasks/jobs/attack_paths/queries.py | 11 +-- .../backend/tasks/jobs/attack_paths/scan.py | 48 +++++++-- .../backend/tasks/jobs/attack_paths/sync.py | 8 +- .../tasks/tests/test_attack_paths_scan.py | 47 ++++----- prowler/CHANGELOG.md | 6 +- 13 files changed, 217 insertions(+), 111 deletions(-) create mode 100644 api/src/backend/api/migrations/0090_attack_paths_cleanup_priority.py diff --git a/api/CHANGELOG.md b/api/CHANGELOG.md index b7243786f1..0d798fb29f 100644 --- a/api/CHANGELOG.md +++ b/api/CHANGELOG.md @@ -4,6 +4,10 @@ All notable changes to the **Prowler API** are documented in this file. ## [1.25.1] (Prowler v5.24.1) +### 🔄 Changed + +- Attack Paths: Restore `SYNC_BATCH_SIZE` and `FINDINGS_BATCH_SIZE` defaults to 1000, upgrade Cartography to 0.135.0, enable Celery queue priority for cleanup task, rewrite Finding insertion, remove AWS graph cleanup and add timing logs [(#10729)](https://github.com/prowler-cloud/prowler/pull/10729) + ### 🐞 Fixed - Attack Paths: Missing `tenant_id` filter while getting related findings after scan completes [(#10722)](https://github.com/prowler-cloud/prowler/pull/10722) diff --git a/api/poetry.lock b/api/poetry.lock index 72ed62f1a3..e5be8ef1cc 100644 --- a/api/poetry.lock +++ b/api/poetry.lock @@ -1526,19 +1526,19 @@ typing-extensions = ">=4.6.0" [[package]] name = "azure-mgmt-resource" -version = "23.3.0" +version = "24.0.0" description = "Microsoft Azure Resource Management Client Library for Python" optional = false -python-versions = ">=3.8" +python-versions = ">=3.9" groups = ["main"] files = [ - {file = "azure_mgmt_resource-23.3.0-py3-none-any.whl", hash = "sha256:ab216ee28e29db6654b989746e0c85a1181f66653929d2cb6e48fba66d9af323"}, - {file = "azure_mgmt_resource-23.3.0.tar.gz", hash = "sha256:fc4f1fd8b6aad23f8af4ed1f913df5f5c92df117449dc354fea6802a2829fea4"}, + {file = "azure_mgmt_resource-24.0.0-py3-none-any.whl", hash = "sha256:27b32cd223e2784269f5a0db3c282042886ee4072d79cedc638438ece7cd0df4"}, + {file = "azure_mgmt_resource-24.0.0.tar.gz", hash = "sha256:cf6b8995fcdd407ac9ff1dd474087129429a1d90dbb1ac77f97c19b96237b265"}, ] [package.dependencies] azure-common = ">=1.1" -azure-mgmt-core = ">=1.3.2" +azure-mgmt-core = ">=1.5.0" isodate = ">=0.6.1" typing-extensions = ">=4.6.0" @@ -1822,19 +1822,19 @@ crt = ["awscrt (==0.27.6)"] [[package]] name = "cartography" -version = "0.132.0" +version = "0.135.0" description = "Explore assets and their relationships across your technical infrastructure." optional = false python-versions = ">=3.10" groups = ["main"] files = [ - {file = "cartography-0.132.0-py3-none-any.whl", hash = "sha256:c070aa51d0ab4479cb043cae70b35e7df49f2fb5f1fa95ccf10000bbeb952262"}, - {file = "cartography-0.132.0.tar.gz", hash = "sha256:7c6332bc57fd2629d7b83aee7bd95a7b2edb0d51ef746efa0461399e0b66625c"}, + {file = "cartography-0.135.0-py3-none-any.whl", hash = "sha256:c62c32a6917b8f23a8b98fe2b6c7c4a918b50f55918482966c4dae1cf5f538e1"}, + {file = "cartography-0.135.0.tar.gz", hash = "sha256:3f500cd22c3b392d00e8b49f62acc95fd4dcd559ce514aafe2eb8101133c7a49"}, ] [package.dependencies] adal = ">=1.2.4" -aioboto3 = ">=13.0.0" +aioboto3 = ">=15.0.0" azure-cli-core = ">=2.26.0" azure-identity = ">=1.5.0" azure-keyvault-certificates = ">=4.0.0" @@ -1852,9 +1852,9 @@ azure-mgmt-keyvault = ">=10.0.0" azure-mgmt-logic = ">=10.0.0" azure-mgmt-monitor = ">=3.0.0" azure-mgmt-network = ">=25.0.0" -azure-mgmt-resource = ">=10.2.0,<25.0.0" +azure-mgmt-resource = ">=24.0.0,<25" azure-mgmt-security = ">=5.0.0" -azure-mgmt-sql = ">=3.0.1,<4" +azure-mgmt-sql = ">=3.0.1" azure-mgmt-storage = ">=16.0.0" azure-mgmt-synapse = ">=2.0.0" azure-mgmt-web = ">=7.0.0" @@ -1862,38 +1862,39 @@ azure-synapse-artifacts = ">=0.17.0" backoff = ">=2.1.2" boto3 = ">=1.15.1" botocore = ">=1.18.1" -cloudflare = ">=4.1.0,<5.0.0" +cloudflare = ">=4.1.0" crowdstrike-falconpy = ">=0.5.1" -cryptography = "*" -dnspython = ">=1.15.0" -duo-client = "*" -google-api-python-client = ">=1.7.8" +cryptography = ">=45.0.0" +dnspython = ">=2.0.0" +duo-client = ">=5.5.0" +google-api-python-client = ">=2.0.0" google-auth = ">=2.37.0" google-cloud-asset = ">=1.0.0" google-cloud-resource-manager = ">=1.14.2" httpx = ">=0.24.0" kubernetes = ">=22.6.0" -marshmallow = ">=3.0.0rc7" -msgraph-sdk = "*" +marshmallow = ">=4.0.0" +msgraph-sdk = ">=1.53.0" msrestazure = ">=0.6.4" neo4j = ">=6.0.0" oci = ">=2.71.0" okta = "<1.0.0" -packageurl-python = "*" -packaging = "*" +packageurl-python = ">=0.17.0" +packaging = ">=26.0.0" pagerduty = ">=4.0.1" policyuniverse = ">=1.1.0.0" PyJWT = {version = ">=2.0.0", extras = ["crypto"]} -python-dateutil = "*" +python-dateutil = ">=2.9.0" python-digitalocean = ">=1.16.0" pyyaml = ">=5.3.1" requests = ">=2.22.0" scaleway = ">=2.10.0" slack-sdk = ">=3.37.0" -statsd = "*" +statsd = ">=4.0.0" typer = ">=0.9.0" -types-aiobotocore-ecr = "*" -xmltodict = "*" +types-aiobotocore-ecr = ">=3.1.0" +workos = ">=5.44.0" +xmltodict = ">=1.0.0" [[package]] name = "celery" @@ -5193,24 +5194,16 @@ files = [ [[package]] name = "marshmallow" -version = "3.26.2" +version = "4.3.0" description = "A lightweight library for converting complex datatypes to and from native Python datatypes." optional = false -python-versions = ">=3.9" +python-versions = ">=3.10" groups = ["main", "dev"] files = [ - {file = "marshmallow-3.26.2-py3-none-any.whl", hash = "sha256:013fa8a3c4c276c24d26d84ce934dc964e2aa794345a0f8c7e5a7191482c8a73"}, - {file = "marshmallow-3.26.2.tar.gz", hash = "sha256:bbe2adb5a03e6e3571b573f42527c6fe926e17467833660bebd11593ab8dfd57"}, + {file = "marshmallow-4.3.0-py3-none-any.whl", hash = "sha256:46c4fe6984707e3cbd485dfebbf0a59874f58d695aad05c1668d15e8c6e13b46"}, + {file = "marshmallow-4.3.0.tar.gz", hash = "sha256:fb43c53b3fe240b8f6af37223d6ef1636f927ad9bea8ab323afad95dff090880"}, ] -[package.dependencies] -packaging = ">=17.0" - -[package.extras] -dev = ["marshmallow[tests]", "pre-commit (>=3.5,<5.0)", "tox"] -docs = ["autodocsumm (==0.2.14)", "furo (==2024.8.6)", "sphinx (==8.1.3)", "sphinx-copybutton (==0.5.2)", "sphinx-issues (==5.0.0)", "sphinxext-opengraph (==0.9.1)"] -tests = ["pytest", "simplejson"] - [[package]] name = "matplotlib" version = "3.10.8" @@ -5504,14 +5497,14 @@ dev = ["bumpver", "isort", "mypy", "pylint", "pytest", "yapf"] [[package]] name = "msgraph-sdk" -version = "1.23.0" +version = "1.55.0" description = "The Microsoft Graph Python SDK" optional = false python-versions = ">=3.9" groups = ["main"] files = [ - {file = "msgraph_sdk-1.23.0-py3-none-any.whl", hash = "sha256:58e0047b4ca59fd82022c02cd73fec0170a3d84f3b76721e3db2a0314df9a58a"}, - {file = "msgraph_sdk-1.23.0.tar.gz", hash = "sha256:6dd1ba9a46f5f0ce8599fd9610133adbd9d1493941438b5d3632fce9e55ed607"}, + {file = "msgraph_sdk-1.55.0-py3-none-any.whl", hash = "sha256:c8e68ebc4b88af5111de312e7fa910a4e76ddf48a4534feadb1fb8a411c48cfc"}, + {file = "msgraph_sdk-1.55.0.tar.gz", hash = "sha256:6df691a31954a050d26b8a678968017e157d940fb377f2a8a4e17a9741b98756"}, ] [package.dependencies] @@ -6714,7 +6707,7 @@ azure-mgmt-postgresqlflexibleservers = "1.1.0" azure-mgmt-rdbms = "10.1.0" azure-mgmt-recoveryservices = "3.1.0" azure-mgmt-recoveryservicesbackup = "9.2.0" -azure-mgmt-resource = "23.3.0" +azure-mgmt-resource = "24.0.0" azure-mgmt-search = "9.1.0" azure-mgmt-security = "7.0.0" azure-mgmt-sql = "3.0.1" @@ -6740,7 +6733,7 @@ jsonschema = "4.23.0" kubernetes = "32.0.1" markdown = "3.10.2" microsoft-kiota-abstractions = "1.9.2" -msgraph-sdk = "1.23.0" +msgraph-sdk = "1.55.0" numpy = "2.0.2" oci = "2.169.0" openstacksdk = "4.2.0" @@ -6971,11 +6964,11 @@ description = "C parser in Python" optional = false python-versions = ">=3.10" groups = ["main", "dev"] +markers = "platform_python_implementation != \"PyPy\" and implementation_name != \"PyPy\"" files = [ {file = "pycparser-3.0-py3-none-any.whl", hash = "sha256:b727414169a36b7d524c1c3e31839a521725078d7b2ff038656844266160a992"}, {file = "pycparser-3.0.tar.gz", hash = "sha256:600f49d217304a5902ac3c37e1281c9fe94e4d0489de643a9504c5cdfdfc6b29"}, ] -markers = {main = "implementation_name != \"PyPy\" and platform_python_implementation != \"PyPy\"", dev = "platform_python_implementation != \"PyPy\" and implementation_name != \"PyPy\""} [[package]] name = "pydantic" @@ -7223,7 +7216,7 @@ description = "The MSALRuntime Python Interop Package" optional = false python-versions = ">=3.6" groups = ["main"] -markers = "(platform_system == \"Windows\" or platform_system == \"Darwin\" or platform_system == \"Linux\") and sys_platform == \"win32\"" +markers = "sys_platform == \"win32\" and (platform_system == \"Windows\" or platform_system == \"Darwin\" or platform_system == \"Linux\")" files = [ {file = "pymsalruntime-0.18.1-cp310-cp310-macosx_14_0_arm64.whl", hash = "sha256:0c22e2e83faa10de422bbfaacc1bb2887c9025ee8a53f0fc2e4f7db01c4a7b66"}, {file = "pymsalruntime-0.18.1-cp310-cp310-macosx_14_0_x86_64.whl", hash = "sha256:8ce2944a0f944833d047bb121396091e00287e2b6373716106da86ea99abf379"}, @@ -8821,6 +8814,23 @@ markupsafe = ">=2.1.1" [package.extras] watchdog = ["watchdog (>=2.3)"] +[[package]] +name = "workos" +version = "6.0.4" +description = "WorkOS Python Client" +optional = false +python-versions = ">=3.10" +groups = ["main"] +files = [ + {file = "workos-6.0.4-py3-none-any.whl", hash = "sha256:548668b3702673536f853ba72a7b5bbbc269e467aaf9ac4f477b6e0177df5e21"}, + {file = "workos-6.0.4.tar.gz", hash = "sha256:b0bfe8fd212b8567422c4ea3732eb33608794033eb3a69900c6b04db183c32d6"}, +] + +[package.dependencies] +cryptography = ">=46.0,<47.0" +httpx = ">=0.28,<1.0" +pyjwt = ">=2.12,<3.0" + [[package]] name = "wrapt" version = "1.17.3" @@ -9414,4 +9424,4 @@ files = [ [metadata] lock-version = "2.1" python-versions = ">=3.11,<3.13" -content-hash = "44caea5e040c54d4c726144d644e67c942b5acf7e316b1fcb2c22b0947614fbe" +content-hash = "5781e74b0692aed541fe445d6713d2dfd792bb226789501420aac4a8cb45aa2a" diff --git a/api/pyproject.toml b/api/pyproject.toml index 0d710d3a28..348cac3e11 100644 --- a/api/pyproject.toml +++ b/api/pyproject.toml @@ -38,7 +38,7 @@ dependencies = [ "matplotlib (==3.10.8)", "reportlab (==4.4.10)", "neo4j (==6.1.0)", - "cartography (==0.132.0)", + "cartography (==0.135.0)", "gevent (==25.9.1)", "werkzeug (==3.1.7)", "sqlparse (==0.5.5)", @@ -62,7 +62,6 @@ django-silk = "5.3.2" docker = "7.1.0" filelock = "3.20.3" freezegun = "1.5.1" -marshmallow = "==3.26.2" mypy = "1.10.1" pylint = "3.2.5" pytest = "9.0.3" diff --git a/api/src/backend/api/migrations/0090_attack_paths_cleanup_priority.py b/api/src/backend/api/migrations/0090_attack_paths_cleanup_priority.py new file mode 100644 index 0000000000..5ef8529b08 --- /dev/null +++ b/api/src/backend/api/migrations/0090_attack_paths_cleanup_priority.py @@ -0,0 +1,23 @@ +from django.db import migrations + +TASK_NAME = "attack-paths-cleanup-stale-scans" + + +def set_cleanup_priority(apps, schema_editor): + PeriodicTask = apps.get_model("django_celery_beat", "PeriodicTask") + PeriodicTask.objects.filter(name=TASK_NAME).update(priority=0) + + +def unset_cleanup_priority(apps, schema_editor): + PeriodicTask = apps.get_model("django_celery_beat", "PeriodicTask") + PeriodicTask.objects.filter(name=TASK_NAME).update(priority=None) + + +class Migration(migrations.Migration): + dependencies = [ + ("api", "0089_backfill_finding_group_status_muted"), + ] + + operations = [ + migrations.RunPython(set_cleanup_priority, unset_cleanup_priority), + ] diff --git a/api/src/backend/config/celery.py b/api/src/backend/config/celery.py index aaa1b1c386..c46c8a426c 100644 --- a/api/src/backend/config/celery.py +++ b/api/src/backend/config/celery.py @@ -17,8 +17,10 @@ celery_app.config_from_object("django.conf:settings", namespace="CELERY") celery_app.conf.update(result_extended=True, result_expires=None) celery_app.conf.broker_transport_options = { - "visibility_timeout": BROKER_VISIBILITY_TIMEOUT + "visibility_timeout": BROKER_VISIBILITY_TIMEOUT, + "queue_order_strategy": "priority", } +celery_app.conf.task_default_priority = 6 celery_app.conf.result_backend_transport_options = { "visibility_timeout": BROKER_VISIBILITY_TIMEOUT } diff --git a/api/src/backend/tasks/jobs/attack_paths/aws.py b/api/src/backend/tasks/jobs/attack_paths/aws.py index e7f26b5173..7248cc39e8 100644 --- a/api/src/backend/tasks/jobs/attack_paths/aws.py +++ b/api/src/backend/tasks/jobs/attack_paths/aws.py @@ -1,6 +1,8 @@ # Portions of this file are based on code from the Cartography project # (https://github.com/cartography-cncf/cartography), which is licensed under the Apache 2.0 License. +import time + from typing import Any import aioboto3 @@ -33,7 +35,7 @@ def start_aws_ingestion( For the scan progress updates: - The caller of this function (`tasks.jobs.attack_paths.scan.run`) has set it to 2. - - When the control returns to the caller, it will be set to 95. + - When the control returns to the caller, it will be set to 93. """ # Initialize variables common to all jobs @@ -89,34 +91,50 @@ def start_aws_ingestion( logger.info( f"Syncing function permission_relationships for AWS account {prowler_api_provider.uid}" ) + t0 = time.perf_counter() cartography_aws.RESOURCE_FUNCTIONS["permission_relationships"](**sync_args) + logger.info( + f"Synced function permission_relationships for AWS account {prowler_api_provider.uid} in {time.perf_counter() - t0:.3f}s" + ) db_utils.update_attack_paths_scan_progress(attack_paths_scan, 88) if "resourcegroupstaggingapi" in requested_syncs: logger.info( f"Syncing function resourcegroupstaggingapi for AWS account {prowler_api_provider.uid}" ) + t0 = time.perf_counter() cartography_aws.RESOURCE_FUNCTIONS["resourcegroupstaggingapi"](**sync_args) + logger.info( + f"Synced function resourcegroupstaggingapi for AWS account {prowler_api_provider.uid} in {time.perf_counter() - t0:.3f}s" + ) db_utils.update_attack_paths_scan_progress(attack_paths_scan, 89) logger.info( f"Syncing ec2_iaminstanceprofile scoped analysis for AWS account {prowler_api_provider.uid}" ) + t0 = time.perf_counter() cartography_aws.run_scoped_analysis_job( "aws_ec2_iaminstanceprofile.json", neo4j_session, common_job_parameters, ) + logger.info( + f"Synced ec2_iaminstanceprofile scoped analysis for AWS account {prowler_api_provider.uid} in {time.perf_counter() - t0:.3f}s" + ) db_utils.update_attack_paths_scan_progress(attack_paths_scan, 90) logger.info( f"Syncing lambda_ecr analysis for AWS account {prowler_api_provider.uid}" ) + t0 = time.perf_counter() cartography_aws.run_analysis_job( "aws_lambda_ecr.json", neo4j_session, common_job_parameters, ) + logger.info( + f"Synced lambda_ecr analysis for AWS account {prowler_api_provider.uid} in {time.perf_counter() - t0:.3f}s" + ) if all( s in requested_syncs @@ -125,25 +143,34 @@ def start_aws_ingestion( logger.info( f"Syncing lb_container_exposure scoped analysis for AWS account {prowler_api_provider.uid}" ) + t0 = time.perf_counter() cartography_aws.run_scoped_analysis_job( "aws_lb_container_exposure.json", neo4j_session, common_job_parameters, ) + logger.info( + f"Synced lb_container_exposure scoped analysis for AWS account {prowler_api_provider.uid} in {time.perf_counter() - t0:.3f}s" + ) if all(s in requested_syncs for s in ["ec2:network_acls", "ec2:load_balancer_v2"]): logger.info( f"Syncing lb_nacl_direct scoped analysis for AWS account {prowler_api_provider.uid}" ) + t0 = time.perf_counter() cartography_aws.run_scoped_analysis_job( "aws_lb_nacl_direct.json", neo4j_session, common_job_parameters, ) + logger.info( + f"Synced lb_nacl_direct scoped analysis for AWS account {prowler_api_provider.uid} in {time.perf_counter() - t0:.3f}s" + ) db_utils.update_attack_paths_scan_progress(attack_paths_scan, 91) logger.info(f"Syncing metadata for AWS account {prowler_api_provider.uid}") + t0 = time.perf_counter() cartography_aws.merge_module_sync_metadata( neo4j_session, group_type="AWSAccount", @@ -152,24 +179,23 @@ def start_aws_ingestion( update_tag=cartography_config.update_tag, stat_handler=cartography_aws.stat_handler, ) + logger.info( + f"Synced metadata for AWS account {prowler_api_provider.uid} in {time.perf_counter() - t0:.3f}s" + ) db_utils.update_attack_paths_scan_progress(attack_paths_scan, 92) # Removing the added extra field del common_job_parameters["AWS_ID"] - logger.info(f"Syncing cleanup_job for AWS account {prowler_api_provider.uid}") - cartography_aws.run_cleanup_job( - "aws_post_ingestion_principals_cleanup.json", - neo4j_session, - common_job_parameters, - ) - db_utils.update_attack_paths_scan_progress(attack_paths_scan, 93) - logger.info(f"Syncing analysis for AWS account {prowler_api_provider.uid}") + t0 = time.perf_counter() cartography_aws._perform_aws_analysis( requested_syncs, neo4j_session, common_job_parameters ) - db_utils.update_attack_paths_scan_progress(attack_paths_scan, 94) + logger.info( + f"Synced analysis for AWS account {prowler_api_provider.uid} in {time.perf_counter() - t0:.3f}s" + ) + db_utils.update_attack_paths_scan_progress(attack_paths_scan, 93) return failed_syncs @@ -234,6 +260,8 @@ def sync_aws_account( ) try: + func_t0 = time.perf_counter() + # `ecr:image_layers` uses `aioboto3_session` instead of `boto3_session` if func_name == "ecr:image_layers": cartography_aws.RESOURCE_FUNCTIONS[func_name]( @@ -257,7 +285,15 @@ def sync_aws_account( else: cartography_aws.RESOURCE_FUNCTIONS[func_name](**sync_args) + logger.info( + f"Synced function {func_name} for AWS account {prowler_api_provider.uid} in {time.perf_counter() - func_t0:.3f}s" + ) + except Exception as e: + logger.info( + f"Synced function {func_name} for AWS account {prowler_api_provider.uid} in {time.perf_counter() - func_t0:.3f}s (FAILED)" + ) + exception_message = utils.stringify_exception( e, f"Exception for AWS sync function: {func_name}" ) diff --git a/api/src/backend/tasks/jobs/attack_paths/config.py b/api/src/backend/tasks/jobs/attack_paths/config.py index 76dbdc1dc5..5f5c523ceb 100644 --- a/api/src/backend/tasks/jobs/attack_paths/config.py +++ b/api/src/backend/tasks/jobs/attack_paths/config.py @@ -8,9 +8,9 @@ from tasks.jobs.attack_paths import aws # Batch size for Neo4j write operations (resource labeling, cleanup) BATCH_SIZE = env.int("ATTACK_PATHS_BATCH_SIZE", 1000) # Batch size for Postgres findings fetch (keyset pagination page size) -FINDINGS_BATCH_SIZE = env.int("ATTACK_PATHS_FINDINGS_BATCH_SIZE", 500) +FINDINGS_BATCH_SIZE = env.int("ATTACK_PATHS_FINDINGS_BATCH_SIZE", 1000) # Batch size for temp-to-tenant graph sync (nodes and relationships per cursor page) -SYNC_BATCH_SIZE = env.int("ATTACK_PATHS_SYNC_BATCH_SIZE", 250) +SYNC_BATCH_SIZE = env.int("ATTACK_PATHS_SYNC_BATCH_SIZE", 1000) # Neo4j internal labels (Prowler-specific, not provider-specific) # - `Internet`: Singleton node representing external internet access for exposed-resource queries diff --git a/api/src/backend/tasks/jobs/attack_paths/findings.py b/api/src/backend/tasks/jobs/attack_paths/findings.py index 8ed5ccf3fc..0b2ecb4c45 100644 --- a/api/src/backend/tasks/jobs/attack_paths/findings.py +++ b/api/src/backend/tasks/jobs/attack_paths/findings.py @@ -12,6 +12,7 @@ from typing import Any, Generator from uuid import UUID import neo4j + from cartography.config import Config as CartographyConfig from celery.utils.log import get_task_logger from tasks.jobs.attack_paths.config import ( @@ -86,17 +87,21 @@ def analysis( prowler_api_provider: Provider, scan_id: str, config: CartographyConfig, -) -> None: +) -> tuple[int, int]: """ Main entry point for Prowler findings analysis. Adds resource labels and loads findings. + Returns (labeled_nodes, findings_loaded). """ - add_resource_label( + total_labeled = add_resource_label( neo4j_session, prowler_api_provider.provider, str(prowler_api_provider.uid) ) findings_data = stream_findings_with_resources(prowler_api_provider, scan_id) - load_findings(neo4j_session, findings_data, prowler_api_provider, config) + total_loaded = load_findings( + neo4j_session, findings_data, prowler_api_provider, config + ) + return total_labeled, total_loaded def add_resource_label( @@ -146,12 +151,11 @@ def load_findings( findings_batches: Generator[list[dict[str, Any]], None, None], prowler_api_provider: Provider, config: CartographyConfig, -) -> None: +) -> int: """Load Prowler findings into the graph, linking them to resources.""" query = render_cypher_template( INSERT_FINDING_TEMPLATE, { - "__ROOT_NODE_LABEL__": get_root_node_label(prowler_api_provider.provider), "__NODE_UID_FIELD__": get_node_uid_field(prowler_api_provider.provider), "__RESOURCE_LABEL__": get_provider_resource_label( prowler_api_provider.provider @@ -160,7 +164,6 @@ def load_findings( ) parameters = { - "provider_uid": str(prowler_api_provider.uid), "last_updated": config.update_tag, "prowler_version": ProwlerConfig.prowler_version, } @@ -178,6 +181,7 @@ def load_findings( neo4j_session.run(query, parameters) logger.info(f"Finished loading {total_records} records in {batch_num} batches") + return total_records # Findings Streaming (Generator-based) diff --git a/api/src/backend/tasks/jobs/attack_paths/queries.py b/api/src/backend/tasks/jobs/attack_paths/queries.py index ab13150b56..26ffa32f92 100644 --- a/api/src/backend/tasks/jobs/attack_paths/queries.py +++ b/api/src/backend/tasks/jobs/attack_paths/queries.py @@ -32,17 +32,14 @@ ADD_RESOURCE_LABEL_TEMPLATE = """ """ INSERT_FINDING_TEMPLATE = f""" - MATCH (account:__ROOT_NODE_LABEL__ {{id: $provider_uid}}) UNWIND $findings_data AS finding_data - OPTIONAL MATCH (account)-->(resource_by_uid:__RESOURCE_LABEL__) - WHERE resource_by_uid.__NODE_UID_FIELD__ = finding_data.resource_uid - WITH account, finding_data, resource_by_uid + OPTIONAL MATCH (resource_by_uid:__RESOURCE_LABEL__ {{__NODE_UID_FIELD__: finding_data.resource_uid}}) + WITH finding_data, resource_by_uid - OPTIONAL MATCH (account)-->(resource_by_id:__RESOURCE_LABEL__) + OPTIONAL MATCH (resource_by_id:__RESOURCE_LABEL__ {{id: finding_data.resource_uid}}) WHERE resource_by_uid IS NULL - AND resource_by_id.id = finding_data.resource_uid - WITH account, finding_data, COALESCE(resource_by_uid, resource_by_id) AS resource + WITH finding_data, COALESCE(resource_by_uid, resource_by_id) AS resource WHERE resource IS NOT NULL MERGE (finding:{PROWLER_FINDING_LABEL} {{id: finding_data.id}}) diff --git a/api/src/backend/tasks/jobs/attack_paths/scan.py b/api/src/backend/tasks/jobs/attack_paths/scan.py index a53a6a530f..382c231eeb 100644 --- a/api/src/backend/tasks/jobs/attack_paths/scan.py +++ b/api/src/backend/tasks/jobs/attack_paths/scan.py @@ -55,6 +55,7 @@ exception propagates to Celery. import logging import time + from typing import Any from cartography.config import Config as CartographyConfig @@ -144,6 +145,12 @@ def run(tenant_id: str, scan_id: str, task_id: str) -> dict[str, Any]: attack_paths_scan, task_id, tenant_cartography_config ) + scan_t0 = time.perf_counter() + logger.info( + f"Starting Attack Paths scan ({attack_paths_scan.id}) for " + f"{prowler_api_provider.provider.upper()} provider {prowler_api_provider.id}" + ) + subgraph_dropped = False sync_completed = False provider_gated = False @@ -169,6 +176,7 @@ def run(tenant_id: str, scan_id: str, task_id: str) -> dict[str, Any]: db_utils.update_attack_paths_scan_progress(attack_paths_scan, 2) # The real scan, where iterates over cloud services + t0 = time.perf_counter() ingestion_exceptions = utils.call_within_event_loop( cartography_ingestion_function, tmp_neo4j_session, @@ -177,19 +185,23 @@ def run(tenant_id: str, scan_id: str, task_id: str) -> dict[str, Any]: prowler_sdk_provider, attack_paths_scan, ) + logger.info( + f"Cartography ingestion completed in {time.perf_counter() - t0:.3f}s " + f"(failed_syncs={len(ingestion_exceptions)})" + ) # Post-processing: Just keeping it to be more Cartography compliant logger.info( f"Syncing Cartography ontology for AWS account {prowler_api_provider.uid}" ) cartography_ontology.run(tmp_neo4j_session, tmp_cartography_config) - db_utils.update_attack_paths_scan_progress(attack_paths_scan, 95) + db_utils.update_attack_paths_scan_progress(attack_paths_scan, 94) logger.info( f"Syncing Cartography analysis for AWS account {prowler_api_provider.uid}" ) cartography_analysis.run(tmp_neo4j_session, tmp_cartography_config) - db_utils.update_attack_paths_scan_progress(attack_paths_scan, 96) + db_utils.update_attack_paths_scan_progress(attack_paths_scan, 95) # Creating Internet node and CAN_ACCESS relationships logger.info( @@ -198,14 +210,20 @@ def run(tenant_id: str, scan_id: str, task_id: str) -> dict[str, Any]: internet.analysis( tmp_neo4j_session, prowler_api_provider, tmp_cartography_config ) + db_utils.update_attack_paths_scan_progress(attack_paths_scan, 96) # Adding Prowler Finding nodes and relationships logger.info( f"Syncing Prowler analysis for AWS account {prowler_api_provider.uid}" ) - findings.analysis( + t0 = time.perf_counter() + labeled_nodes, findings_loaded = findings.analysis( tmp_neo4j_session, prowler_api_provider, scan_id, tmp_cartography_config ) + logger.info( + f"Prowler analysis completed in {time.perf_counter() - t0:.3f}s " + f"(findings={findings_loaded}, labeled_nodes={labeled_nodes})" + ) db_utils.update_attack_paths_scan_progress(attack_paths_scan, 97) logger.info( @@ -227,22 +245,33 @@ def run(tenant_id: str, scan_id: str, task_id: str) -> dict[str, Any]: logger.info(f"Deleting existing provider graph in {tenant_database_name}") db_utils.set_provider_graph_data_ready(attack_paths_scan, False) provider_gated = True - graph_database.drop_subgraph( + + t0 = time.perf_counter() + deleted_nodes = graph_database.drop_subgraph( database=tenant_database_name, provider_id=str(prowler_api_provider.id), ) + logger.info( + f"Deleted existing provider graph in {time.perf_counter() - t0:.3f}s " + f"(deleted_nodes={deleted_nodes})" + ) subgraph_dropped = True db_utils.update_attack_paths_scan_progress(attack_paths_scan, 98) logger.info( f"Syncing graph from {tmp_database_name} into {tenant_database_name}" ) - sync.sync_graph( + t0 = time.perf_counter() + sync_result = sync.sync_graph( source_database=tmp_database_name, target_database=tenant_database_name, tenant_id=str(prowler_api_provider.tenant_id), provider_id=str(prowler_api_provider.id), ) + logger.info( + f"Synced graph in {time.perf_counter() - t0:.3f}s " + f"(nodes={sync_result['nodes']}, relationships={sync_result['relationships']})" + ) sync_completed = True db_utils.set_graph_data_ready(attack_paths_scan, True) db_utils.update_attack_paths_scan_progress(attack_paths_scan, 99) @@ -250,17 +279,16 @@ def run(tenant_id: str, scan_id: str, task_id: str) -> dict[str, Any]: logger.info(f"Clearing Neo4j cache for database {tenant_database_name}") graph_database.clear_cache(tenant_database_name) - logger.info( - f"Completed Cartography ({attack_paths_scan.id}) for " - f"{prowler_api_provider.provider.upper()} provider {prowler_api_provider.id}" - ) - logger.info(f"Dropping temporary Neo4j database {tmp_database_name}") graph_database.drop_database(tmp_database_name) db_utils.finish_attack_paths_scan( attack_paths_scan, StateChoices.COMPLETED, ingestion_exceptions ) + logger.info( + f"Attack Paths scan completed in {time.perf_counter() - scan_t0:.3f}s " + f"(state=completed, failed_syncs={len(ingestion_exceptions)})" + ) return ingestion_exceptions except Exception as e: diff --git a/api/src/backend/tasks/jobs/attack_paths/sync.py b/api/src/backend/tasks/jobs/attack_paths/sync.py index 24ffa6cf48..f720a12e82 100644 --- a/api/src/backend/tasks/jobs/attack_paths/sync.py +++ b/api/src/backend/tasks/jobs/attack_paths/sync.py @@ -5,6 +5,8 @@ This module handles syncing graph data from temporary scan databases to the tenant database, adding provider isolation labels and properties. """ +import time + from collections import defaultdict from typing import Any @@ -81,6 +83,7 @@ def sync_nodes( Source and target sessions are opened sequentially per batch to avoid holding two Bolt connections simultaneously for the entire sync duration. """ + t0 = time.perf_counter() last_id = -1 total_synced = 0 @@ -117,7 +120,7 @@ def sync_nodes( total_synced += batch_count logger.info( - f"Synced {total_synced} nodes from {source_database} to {target_database}" + f"Synced {total_synced} nodes from {source_database} to {target_database} in {time.perf_counter() - t0:.3f}s" ) return total_synced @@ -136,6 +139,7 @@ def sync_relationships( Source and target sessions are opened sequentially per batch to avoid holding two Bolt connections simultaneously for the entire sync duration. """ + t0 = time.perf_counter() last_id = -1 total_synced = 0 @@ -166,7 +170,7 @@ def sync_relationships( total_synced += batch_count logger.info( - f"Synced {total_synced} relationships from {source_database} to {target_database}" + f"Synced {total_synced} relationships from {source_database} to {target_database} in {time.perf_counter() - t0:.3f}s" ) return total_synced diff --git a/api/src/backend/tasks/tests/test_attack_paths_scan.py b/api/src/backend/tasks/tests/test_attack_paths_scan.py index 1b1beb11a9..283c0650e1 100644 --- a/api/src/backend/tasks/tests/test_attack_paths_scan.py +++ b/api/src/backend/tasks/tests/test_attack_paths_scan.py @@ -38,11 +38,14 @@ class TestAttackPathsRun: @patch("tasks.jobs.attack_paths.scan.db_utils.finish_attack_paths_scan") @patch("tasks.jobs.attack_paths.scan.db_utils.update_attack_paths_scan_progress") @patch("tasks.jobs.attack_paths.scan.db_utils.starting_attack_paths_scan") - @patch("tasks.jobs.attack_paths.scan.sync.sync_graph") - @patch("tasks.jobs.attack_paths.scan.graph_database.drop_subgraph") + @patch( + "tasks.jobs.attack_paths.scan.sync.sync_graph", + return_value={"nodes": 0, "relationships": 0}, + ) + @patch("tasks.jobs.attack_paths.scan.graph_database.drop_subgraph", return_value=0) @patch("tasks.jobs.attack_paths.scan.indexes.create_sync_indexes") @patch("tasks.jobs.attack_paths.scan.internet.analysis") - @patch("tasks.jobs.attack_paths.scan.findings.analysis") + @patch("tasks.jobs.attack_paths.scan.findings.analysis", return_value=(0, 0)) @patch("tasks.jobs.attack_paths.scan.indexes.create_findings_indexes") @patch("tasks.jobs.attack_paths.scan.cartography_ontology.run") @patch("tasks.jobs.attack_paths.scan.cartography_analysis.run") @@ -188,7 +191,7 @@ class TestAttackPathsRun: @patch("tasks.jobs.attack_paths.scan.db_utils.set_provider_graph_data_ready") @patch("tasks.jobs.attack_paths.scan.db_utils.update_attack_paths_scan_progress") @patch("tasks.jobs.attack_paths.scan.db_utils.starting_attack_paths_scan") - @patch("tasks.jobs.attack_paths.scan.findings.analysis") + @patch("tasks.jobs.attack_paths.scan.findings.analysis", return_value=(0, 0)) @patch("tasks.jobs.attack_paths.scan.internet.analysis") @patch("tasks.jobs.attack_paths.scan.indexes.create_findings_indexes") @patch("tasks.jobs.attack_paths.scan.cartography_analysis.run") @@ -287,7 +290,7 @@ class TestAttackPathsRun: @patch("tasks.jobs.attack_paths.scan.db_utils.set_provider_graph_data_ready") @patch("tasks.jobs.attack_paths.scan.db_utils.update_attack_paths_scan_progress") @patch("tasks.jobs.attack_paths.scan.db_utils.starting_attack_paths_scan") - @patch("tasks.jobs.attack_paths.scan.findings.analysis") + @patch("tasks.jobs.attack_paths.scan.findings.analysis", return_value=(0, 0)) @patch("tasks.jobs.attack_paths.scan.internet.analysis") @patch("tasks.jobs.attack_paths.scan.indexes.create_findings_indexes") @patch("tasks.jobs.attack_paths.scan.cartography_analysis.run") @@ -390,7 +393,7 @@ class TestAttackPathsRun: @patch("tasks.jobs.attack_paths.scan.db_utils.set_provider_graph_data_ready") @patch("tasks.jobs.attack_paths.scan.db_utils.update_attack_paths_scan_progress") @patch("tasks.jobs.attack_paths.scan.db_utils.starting_attack_paths_scan") - @patch("tasks.jobs.attack_paths.scan.findings.analysis") + @patch("tasks.jobs.attack_paths.scan.findings.analysis", return_value=(0, 0)) @patch("tasks.jobs.attack_paths.scan.internet.analysis") @patch("tasks.jobs.attack_paths.scan.indexes.create_findings_indexes") @patch("tasks.jobs.attack_paths.scan.cartography_analysis.run") @@ -489,14 +492,17 @@ class TestAttackPathsRun: @patch("tasks.jobs.attack_paths.scan.db_utils.set_provider_graph_data_ready") @patch("tasks.jobs.attack_paths.scan.db_utils.update_attack_paths_scan_progress") @patch("tasks.jobs.attack_paths.scan.db_utils.starting_attack_paths_scan") - @patch("tasks.jobs.attack_paths.scan.sync.sync_graph") + @patch( + "tasks.jobs.attack_paths.scan.sync.sync_graph", + return_value={"nodes": 0, "relationships": 0}, + ) @patch( "tasks.jobs.attack_paths.scan.graph_database.drop_subgraph", side_effect=RuntimeError("drop failed"), ) @patch("tasks.jobs.attack_paths.scan.indexes.create_sync_indexes") @patch("tasks.jobs.attack_paths.scan.internet.analysis") - @patch("tasks.jobs.attack_paths.scan.findings.analysis") + @patch("tasks.jobs.attack_paths.scan.findings.analysis", return_value=(0, 0)) @patch("tasks.jobs.attack_paths.scan.indexes.create_findings_indexes") @patch("tasks.jobs.attack_paths.scan.cartography_ontology.run") @patch("tasks.jobs.attack_paths.scan.cartography_analysis.run") @@ -609,7 +615,7 @@ class TestAttackPathsRun: @patch("tasks.jobs.attack_paths.scan.graph_database.drop_subgraph") @patch("tasks.jobs.attack_paths.scan.indexes.create_sync_indexes") @patch("tasks.jobs.attack_paths.scan.internet.analysis") - @patch("tasks.jobs.attack_paths.scan.findings.analysis") + @patch("tasks.jobs.attack_paths.scan.findings.analysis", return_value=(0, 0)) @patch("tasks.jobs.attack_paths.scan.indexes.create_findings_indexes") @patch("tasks.jobs.attack_paths.scan.cartography_ontology.run") @patch("tasks.jobs.attack_paths.scan.cartography_analysis.run") @@ -718,11 +724,14 @@ class TestAttackPathsRun: @patch("tasks.jobs.attack_paths.scan.db_utils.set_provider_graph_data_ready") @patch("tasks.jobs.attack_paths.scan.db_utils.update_attack_paths_scan_progress") @patch("tasks.jobs.attack_paths.scan.db_utils.starting_attack_paths_scan") - @patch("tasks.jobs.attack_paths.scan.sync.sync_graph") + @patch( + "tasks.jobs.attack_paths.scan.sync.sync_graph", + return_value={"nodes": 0, "relationships": 0}, + ) @patch("tasks.jobs.attack_paths.scan.graph_database.drop_subgraph") @patch("tasks.jobs.attack_paths.scan.indexes.create_sync_indexes") @patch("tasks.jobs.attack_paths.scan.internet.analysis") - @patch("tasks.jobs.attack_paths.scan.findings.analysis") + @patch("tasks.jobs.attack_paths.scan.findings.analysis", return_value=(0, 0)) @patch("tasks.jobs.attack_paths.scan.indexes.create_findings_indexes") @patch("tasks.jobs.attack_paths.scan.cartography_ontology.run") @patch("tasks.jobs.attack_paths.scan.cartography_analysis.run") @@ -833,14 +842,17 @@ class TestAttackPathsRun: @patch("tasks.jobs.attack_paths.scan.db_utils.set_provider_graph_data_ready") @patch("tasks.jobs.attack_paths.scan.db_utils.update_attack_paths_scan_progress") @patch("tasks.jobs.attack_paths.scan.db_utils.starting_attack_paths_scan") - @patch("tasks.jobs.attack_paths.scan.sync.sync_graph") + @patch( + "tasks.jobs.attack_paths.scan.sync.sync_graph", + return_value={"nodes": 0, "relationships": 0}, + ) @patch( "tasks.jobs.attack_paths.scan.graph_database.drop_subgraph", side_effect=RuntimeError("drop failed"), ) @patch("tasks.jobs.attack_paths.scan.indexes.create_sync_indexes") @patch("tasks.jobs.attack_paths.scan.internet.analysis") - @patch("tasks.jobs.attack_paths.scan.findings.analysis") + @patch("tasks.jobs.attack_paths.scan.findings.analysis", return_value=(0, 0)) @patch("tasks.jobs.attack_paths.scan.indexes.create_findings_indexes") @patch("tasks.jobs.attack_paths.scan.cartography_ontology.run") @patch("tasks.jobs.attack_paths.scan.cartography_analysis.run") @@ -1274,10 +1286,6 @@ class TestAttackPathsFindingsHelpers: mock_session = MagicMock() with ( - patch( - "tasks.jobs.attack_paths.findings.get_root_node_label", - return_value="AWSAccount", - ), patch( "tasks.jobs.attack_paths.findings.get_node_uid_field", return_value="arn", @@ -1294,7 +1302,6 @@ class TestAttackPathsFindingsHelpers: assert mock_session.run.call_count == 2 for call_args in mock_session.run.call_args_list: params = call_args.args[1] - assert params["provider_uid"] == str(provider.uid) assert params["last_updated"] == config.update_tag assert "findings_data" in params @@ -1673,10 +1680,6 @@ class TestAttackPathsFindingsHelpers: yield # Make it a generator with ( - patch( - "tasks.jobs.attack_paths.findings.get_root_node_label", - return_value="AWSAccount", - ), patch( "tasks.jobs.attack_paths.findings.get_node_uid_field", return_value="arn", diff --git a/prowler/CHANGELOG.md b/prowler/CHANGELOG.md index 9b8d1e3322..c3d1776a0e 100644 --- a/prowler/CHANGELOG.md +++ b/prowler/CHANGELOG.md @@ -2,16 +2,12 @@ All notable changes to the **Prowler SDK** are documented in this file. -## [5.25.0] (Prowler UNRELEASED) +## [5.24.1] (Prowler UNRELEASED) ### 🔄 Changed - bumped `msgraph-sdk` from 1.23.0 to 1.55.0 and `azure-mgmt-resource` from 23.3.0 to 24.0.0, removing `marshmallow` as is a transitively dev dependency [(#10733)](https://github.com/prowler-cloud/prowler/pull/10733) ---- - -## [5.24.1] (Prowler UNRELEASED) - ### 🐞 Fixed - Cloudflare account-scoped API tokens failing connection test in the App with `CloudflareUserTokenRequiredError` [(#10723)](https://github.com/prowler-cloud/prowler/pull/10723)