From cd56310cbed2237c82950b0c486626a03c452b08 Mon Sep 17 00:00:00 2001 From: Josema Camacho Date: Fri, 10 Oct 2025 13:24:45 +0200 Subject: [PATCH] feat(attack-paths): prowler to cartography data model --- api/Dockerfile | 2 + api/src/backend/api/v1/views.py | 26 +++ .../tasks/jobs/cartography/__init__.py | 29 +++ .../tasks/jobs/cartography/aws/__init__.py | 31 +++ .../backend/tasks/jobs/cartography/aws/ecs.py | 211 +++++++++++++++++ .../backend/tasks/jobs/cartography/aws/iam.py | 135 +++++++++++ .../backend/tasks/jobs/cartography/aws/s3.py | 221 ++++++++++++++++++ docker-compose-dev.yml | 31 +++ 8 files changed, 686 insertions(+) create mode 100644 api/src/backend/tasks/jobs/cartography/__init__.py create mode 100644 api/src/backend/tasks/jobs/cartography/aws/__init__.py create mode 100644 api/src/backend/tasks/jobs/cartography/aws/ecs.py create mode 100644 api/src/backend/tasks/jobs/cartography/aws/iam.py create mode 100644 api/src/backend/tasks/jobs/cartography/aws/s3.py diff --git a/api/Dockerfile b/api/Dockerfile index a3cfb21782..9904d6b530 100644 --- a/api/Dockerfile +++ b/api/Dockerfile @@ -58,6 +58,8 @@ ENV PATH="/home/prowler/.local/bin:$PATH" RUN poetry install --no-root && \ rm -rf ~/.cache/pip +RUN poetry run python -m pip install cartography==0.117.0 + RUN poetry run python "$(poetry env info --path)/src/prowler/prowler/providers/m365/lib/powershell/m365_powershell.py" COPY src/backend/ ./backend/ diff --git a/api/src/backend/api/v1/views.py b/api/src/backend/api/v1/views.py index 8be9f3aed3..1bc69800a8 100644 --- a/api/src/backend/api/v1/views.py +++ b/api/src/backend/api/v1/views.py @@ -1,10 +1,12 @@ import glob import logging import os + from datetime import datetime, timedelta, timezone from urllib.parse import urljoin import sentry_sdk + from allauth.socialaccount.models import SocialAccount, SocialApp from allauth.socialaccount.providers.github.views import GitHubOAuth2Adapter from allauth.socialaccount.providers.google.views import GoogleOAuth2Adapter @@ -205,6 +207,7 @@ from api.v1.serializers import ( UserSerializer, UserUpdateSerializer, ) +from tasks.jobs.cartography import sync_scan_to_cartography logger = logging.getLogger(BackendLogger.API) @@ -1528,6 +1531,29 @@ class ProviderViewSet(BaseRLSViewSet): }, ) + @extend_schema( + tags=["Provider"], + summary="Sync resources to Cartography", + description=( + "Synchronous endpoint to trigger Cartography sync without full validation. " + "Intended for development/testing; a Celery-driven integration may replace this later." + ), + request=None, + responses={200: OpenApiResponse(description="Sync completed successfully")}, + ) + @action(detail=True, methods=["post"], url_name="cartography_sync") + def cartography_sync(self, request, pk=None): + get_object_or_404(Provider, pk=pk) + # Always run synchronously + payload = request.data if isinstance(request.data, dict) else {} + + result = sync_scan_to_cartography( + tenant_id=self.request.tenant_id, + provider_id=pk, + scan_id=payload.get("scan_id"), + ) + return Response(data={"result": result}, status=status.HTTP_201_CREATED) + @extend_schema_view( list=extend_schema( diff --git a/api/src/backend/tasks/jobs/cartography/__init__.py b/api/src/backend/tasks/jobs/cartography/__init__.py new file mode 100644 index 0000000000..b7e9d58716 --- /dev/null +++ b/api/src/backend/tasks/jobs/cartography/__init__.py @@ -0,0 +1,29 @@ +from celery.utils.log import get_task_logger +from neo4j import GraphDatabase + +from tasks.jobs.cartography.aws import sync_aws + +logger = get_task_logger(__name__) + + +def sync_scan_to_cartography( + tenant_id: str, + provider_id: str, + scan_id: str, +): + """ + Sync scan data to Cartography. + """ + + logger.info(f"Sync Cartography - Tenant {tenant_id} - Provider {provider_id} - Scan {scan_id}") + + # TODO: Get Neo4j parameters from settings + with GraphDatabase.driver("bolt://neo4j:7687", auth=("neo4j", "neo4j")) as driver: + with driver.session() as neo4j_session: + # TODO: Depending on the provider type use the appropriate sync function + return sync_aws( + tenant_id=tenant_id, + provider_id=provider_id, + scan_id=scan_id, + neo4j_session=neo4j_session, + ) diff --git a/api/src/backend/tasks/jobs/cartography/aws/__init__.py b/api/src/backend/tasks/jobs/cartography/aws/__init__.py new file mode 100644 index 0000000000..3eb0fdb78a --- /dev/null +++ b/api/src/backend/tasks/jobs/cartography/aws/__init__.py @@ -0,0 +1,31 @@ +from datetime import datetime, timezone +from typing import Any + +import neo4j + +from tasks.jobs.cartography.aws.s3 import sync_aws_s3 +from tasks.jobs.cartography.aws.ecs import sync_aws_ecs +from tasks.jobs.cartography.aws.iam import sync_aws_iam + + +def sync_aws( + tenant_id: str, + provider_id: str, + scan_id: str, + regions: list[str], + neo4j_session: neo4j.Session, +) -> dict[str, Any]: + """ + Sync AWS resources for a specific tenant and provider. + """ + + regions = [] # TODO: Get `regions` from scan + + update_tag = int(datetime.now(tz=timezone.utc).timestamp() * 1000) # TODO: Check if is calculated right + common_job_parameters = {"UPDATE_TAG": update_tag} # TODO: Add other stuff to `common_job_parameters` + + return { + "s3": sync_aws_s3(tenant_id, provider_id, scan_id, regions, neo4j_session, update_tag, common_job_parameters), + # "ecs": sync_aws_ecs(tenant_id, provider_id, scan_id, regions, neo4j_session, update_tag, common_job_parameters), + # "iam": sync_aws_iam(tenant_id, provider_id, scan_id, neo4j_session, update_tag, common_job_parameters), + } diff --git a/api/src/backend/tasks/jobs/cartography/aws/ecs.py b/api/src/backend/tasks/jobs/cartography/aws/ecs.py new file mode 100644 index 0000000000..34e0b62cde --- /dev/null +++ b/api/src/backend/tasks/jobs/cartography/aws/ecs.py @@ -0,0 +1,211 @@ +from __future__ import annotations + +import json +from datetime import datetime, timezone +from typing import Any, Dict, List, Optional + +from cartography.intel.aws import ecs as carto_ecs +from celery.utils.log import get_task_logger +from neo4j import GraphDatabase + +from api.db_utils import rls_transaction +from api.models import Provider, Resource, ResourceScanSummary + +logger = get_task_logger(__name__) + + +def sync_aws_ecs( + tenant_id: str, + provider_id: str, + scan_id: Optional[str], + regions: List[str], + neo4j_conf: Dict[str, Any], +) -> Dict[str, Any]: + try: + from neo4j import GraphDatabase as _ # ensure import present + except Exception as e: + logger.error(f"Neo4j not available: {e}") + return {"error": str(e)} + + with rls_transaction(tenant_id): + provider = Provider.objects.get(pk=provider_id) + account_id = provider.uid + + base_qs = Resource.objects.filter(provider_id=provider_id, service="ecs") + if scan_id: + rss_ids = ResourceScanSummary.objects.filter( + tenant_id=tenant_id, scan_id=scan_id, service="ecs" + ).values_list("resource_id", flat=True) + base_qs = base_qs.filter(id__in=list(rss_ids)) + if regions: + base_qs = base_qs.filter(region__in=regions) + + ecs_resources = list( + base_qs.only("uid", "name", "type", "region", "metadata", "details") + ) + + # Build in-memory maps for patched extractors + clusters_by_region: Dict[str, List[Dict[str, Any]]] = {} + cluster_arns_by_region: Dict[str, List[str]] = {} + clusters_by_arn: Dict[str, Dict[str, Any]] = {} + services_by_cluster: Dict[str, List[Dict[str, Any]]] = {} + tasks_by_cluster: Dict[str, List[Dict[str, Any]]] = {} + tds_by_arn: Dict[str, Dict[str, Any]] = {} + container_instances_by_cluster: Dict[str, List[Dict[str, Any]]] = {} + + for r in ecs_resources: + region = r.region or "" + obj = _ecs_collect_items([r])[0] + if r.type in ("cluster", "ecs_cluster"): + arn = obj.get("clusterArn") or r.uid + clusters_by_arn[arn] = obj + cluster_arns_by_region.setdefault(region, []).append(arn) + clusters_by_region.setdefault(region, []).append(obj) + elif r.type in ("service", "ecs_service"): + cluster_arn = obj.get("clusterArn") or _ecs_cluster_arn_from_uid(obj.get("_uid"), region, account_id) + if cluster_arn: + services_by_cluster.setdefault(cluster_arn, []).append(obj) + elif r.type in ("task", "ecs_task"): + cluster_arn = obj.get("clusterArn") or _ecs_cluster_arn_from_uid(obj.get("_uid"), region, account_id) + if cluster_arn: + tasks_by_cluster.setdefault(cluster_arn, []).append(obj) + td_arn = obj.get("taskDefinitionArn") + if td_arn: + tds_by_arn.setdefault(td_arn, {"taskDefinitionArn": td_arn}) + elif r.type in ("task_definition", "ecs_task_definition"): + arn = obj.get("taskDefinitionArn") or r.uid + tds_by_arn[arn] = obj + + uri = neo4j_conf.get("uri") + user = neo4j_conf.get("user") or neo4j_conf.get("username") + password = neo4j_conf.get("password") + database = neo4j_conf.get("database") + if not all([uri, user, password]): + logger.error("Neo4j configuration incomplete: require uri, user, password") + return {"error": "missing_neo4j_config"} + + update_tag = int(datetime.now(tz=timezone.utc).timestamp() * 1000) + common_job_parameters = {"UPDATE_TAG": update_tag, "AWS_ID": account_id} + driver = GraphDatabase.driver(uri, auth=(user, password)) + + # Save originals references (not restored per request) + def _patched_get_ecs_cluster_arns(_boto3_session, region): + return cluster_arns_by_region.get(region, []) + + def _patched_get_ecs_clusters(_boto3_session, cluster_arns, region): + return [clusters_by_arn.get(arn) for arn in cluster_arns if clusters_by_arn.get(arn)] + + def _patched_get_ecs_container_instances(_boto3_session, cluster_arn, region): + return container_instances_by_cluster.get(cluster_arn, []) + + def _patched_get_ecs_services(_boto3_session, cluster_arn, region): + return services_by_cluster.get(cluster_arn, []) + + def _patched_get_ecs_tasks(_boto3_session, cluster_arn, region): + return tasks_by_cluster.get(cluster_arn, []) + + def _patched_get_ecs_task_definitions(_boto3_session, task_definition_arns, region): + out = [] + for arn in task_definition_arns or []: + td = tds_by_arn.get(arn) + if td: + out.append(td) + return out + + # Apply patches + setattr(carto_ecs, "get_ecs_cluster_arns", _patched_get_ecs_cluster_arns) + setattr(carto_ecs, "get_ecs_clusters", _patched_get_ecs_clusters) + setattr(carto_ecs, "get_ecs_container_instances", _patched_get_ecs_container_instances) + setattr(carto_ecs, "get_ecs_services", _patched_get_ecs_services) + setattr(carto_ecs, "get_ecs_tasks", _patched_get_ecs_tasks) + setattr(carto_ecs, "get_ecs_task_definitions", _patched_get_ecs_task_definitions) + + try: + with driver.session(database=database) if database else driver.session() as neo4j_session: + class _Boto3SessionStub: + pass + + boto3_session = _Boto3SessionStub() + try: + carto_ecs.sync( + neo4j_session, + boto3_session, + account_id, + regions, + update_tag, + common_job_parameters, + ) + except TypeError: + try: + carto_ecs.sync( + neo4j_session, + boto3_session, + account_id, + regions, + update_tag, + ) + except TypeError: + for region in regions or list(cluster_arns_by_region.keys()): + carto_ecs.sync( + neo4j_session, + boto3_session, + account_id, + region, + update_tag, + common_job_parameters, + ) + finally: + try: + driver.close() + except Exception: + pass + + return {"regions": len(regions or cluster_arns_by_region.keys())} + + +def _ecs_collect_items(resources: List[Resource]) -> List[Dict[str, Any]]: + items: List[Dict[str, Any]] = [] + for r in resources or []: + payload = None + for raw in (getattr(r, "metadata", None), getattr(r, "details", None)): + if not raw: + continue + try: + data = json.loads(raw) if isinstance(raw, str) else raw + if isinstance(data, dict): + payload = data + break + except Exception: + continue + obj: Dict[str, Any] = payload.copy() if isinstance(payload, dict) else {} + if r.uid: + obj.setdefault("arn", r.uid) + if r.name: + obj.setdefault("clusterName", r.name) + obj.setdefault("serviceName", r.name) + if r.type in ("cluster", "ecs_cluster"): + obj.setdefault("clusterArn", r.uid) + if not obj.get("clusterName") and r.uid: + obj["clusterName"] = r.uid.split("/")[-1] + elif r.type in ("service", "ecs_service"): + obj.setdefault("serviceArn", r.uid) + if not obj.get("serviceName") and r.uid: + obj["serviceName"] = r.uid.split("/")[-1] + elif r.type in ("task", "ecs_task"): + obj.setdefault("taskArn", r.uid) + elif r.type in ("task_definition", "ecs_task_definition"): + if "taskDefinition" in obj and isinstance(obj["taskDefinition"], dict): + obj.update(obj["taskDefinition"]) + obj.setdefault("taskDefinitionArn", r.uid) + obj["_uid"] = r.uid + items.append(obj) + return items + + +def _ecs_cluster_arn_from_uid(uid: Optional[str], region: str, account_id: str) -> Optional[str]: + if not uid: + return None + if ":ecs:" in uid and ":cluster/" in uid: + return uid + name = uid.split("/")[-1] + return f"arn:aws:ecs:{region}:{account_id}:cluster/{name}" diff --git a/api/src/backend/tasks/jobs/cartography/aws/iam.py b/api/src/backend/tasks/jobs/cartography/aws/iam.py new file mode 100644 index 0000000000..6501c05d86 --- /dev/null +++ b/api/src/backend/tasks/jobs/cartography/aws/iam.py @@ -0,0 +1,135 @@ +from __future__ import annotations + +import json +from datetime import datetime, timezone +from typing import Any, Dict, List, Optional + +from cartography.intel.aws import iam as carto_iam +from celery.utils.log import get_task_logger +from neo4j import GraphDatabase + +from api.db_utils import rls_transaction +from api.models import Provider, Resource, ResourceScanSummary + +logger = get_task_logger(__name__) + + +def sync_aws_iam( + tenant_id: str, + provider_id: str, + scan_id: Optional[str], + neo4j_conf: Dict[str, Any], +) -> Dict[str, Any]: + try: + from neo4j import GraphDatabase as _ # ensure import present + except Exception as e: + logger.error(f"Neo4j not available: {e}") + return {"error": str(e)} + + with rls_transaction(tenant_id): + provider = Provider.objects.get(pk=provider_id) + account_id = provider.uid + + base_qs = Resource.objects.filter( + provider_id=provider_id, service="iam", type__in=["role", "iam_role"] + ) + if scan_id: + rss_ids = ResourceScanSummary.objects.filter( + tenant_id=tenant_id, scan_id=scan_id, service="iam" + ).values_list("resource_id", flat=True) + base_qs = base_qs.filter(id__in=list(rss_ids)) + role_resources = list( + base_qs.only("uid", "name", "metadata", "details", "inserted_at") + ) + + roles: List[Dict[str, Any]] = [] + for r in role_resources: + role_obj: Dict[str, Any] = { + "Arn": r.uid, + "RoleName": r.name or r.uid.split("/")[-1], + "RoleId": r.uid.split("/")[-1], + "Path": "/", + "CreateDate": (r.inserted_at or datetime.now(tz=timezone.utc)).isoformat(), + "AssumeRolePolicyDocument": {"Statement": []}, + } + for raw in (getattr(r, "metadata", None), getattr(r, "details", None)): + if not raw: + continue + try: + data = json.loads(raw) if isinstance(raw, str) else raw + except Exception: + continue + if not isinstance(data, dict): + continue + pol = ( + data.get("AssumeRolePolicyDocument") + or data.get("AssumeRolePolicy") + or data.get("assume_role_policy_document") + ) + if pol and isinstance(pol, dict) and pol.get("Statement"): + role_obj["AssumeRolePolicyDocument"] = pol + break + roles.append(role_obj) + + uri = neo4j_conf.get("uri") + user = neo4j_conf.get("user") or neo4j_conf.get("username") + password = neo4j_conf.get("password") + database = neo4j_conf.get("database") + if not all([uri, user, password]): + logger.error("Neo4j configuration incomplete: require uri, user, password") + return {"error": "missing_neo4j_config"} + + update_tag = int(datetime.now(tz=timezone.utc).timestamp() * 1000) + common_job_parameters = {"UPDATE_TAG": update_tag, "AWS_ID": account_id} + driver = GraphDatabase.driver(uri, auth=(user, password)) + + # Provide minimal patches for IAM getters + def _patched_get_iam_roles(_boto3_session): + return roles + + def _empty_list(*args, **kwargs): + return [] + + setattr(carto_iam, "get_iam_roles", _patched_get_iam_roles) + # Some versions may use list_roles instead + if hasattr(carto_iam, "list_roles"): + setattr(carto_iam, "list_roles", _patched_get_iam_roles) + for fname in [ + "get_iam_users", + "get_iam_groups", + "get_iam_policies", + "get_iam_role_inline_policies", + "get_iam_role_attached_policies", + "get_iam_instance_profiles", + ]: + if hasattr(carto_iam, fname): + setattr(carto_iam, fname, _empty_list) + + try: + with driver.session(database=database) if database else driver.session() as neo4j_session: + class _Boto3SessionStub: + pass + + boto3_session = _Boto3SessionStub() + try: + carto_iam.sync( + neo4j_session, + boto3_session, + account_id, + update_tag, + common_job_parameters, + ) + except TypeError: + carto_iam.sync( + neo4j_session, + boto3_session, + account_id, + update_tag, + ) + finally: + try: + driver.close() + except Exception: + pass + + return {"roles": len(roles)} diff --git a/api/src/backend/tasks/jobs/cartography/aws/s3.py b/api/src/backend/tasks/jobs/cartography/aws/s3.py new file mode 100644 index 0000000000..ee1e0fd0b2 --- /dev/null +++ b/api/src/backend/tasks/jobs/cartography/aws/s3.py @@ -0,0 +1,221 @@ +import json + +from datetime import datetime, timezone +from typing import Any, Generator + +import neo4j + +from cartography.intel.aws import s3 as cartography_s3 +from celery.utils.log import get_task_logger + +from api.db_utils import rls_transaction +from api.models import Provider, Resource, ResourceScanSummary + +logger = get_task_logger(__name__) + + +def sync_aws_s3( + tenant_id: str, + provider_id: str, + scan_id: str, + regions: list[str], + neo4j_session: neo4j.Session, + update_tag: int, + common_job_parameters: dict[str, Any], +) -> dict[str, Any]: + """ + Monkey-patch Cartography S3 extractors and call cartography.intel.aws.s3.sync. + + Patched functions: + - get_s3_bucket_list + - get_s3_bucket_details + - _sync_s3_notifications (uses Prowler data and Cartography's original _load_s3_notifications) + """ + + with rls_transaction(tenant_id): + provider = Provider.objects.get(pk=provider_id) + account_id = provider.uid + + bucket_list = _build_s3_bucket_list(tenant_id, provider_id, scan_id, regions) + bucket_data = _build_s3_bucket_details(bucket_list) + bucket_notifications = _build_s3_notifications(tenant_id, provider_id, scan_id, regions) + + def _patched_get_s3_bucket_list(_boto3_session): + return bucket_list + + def _patched_get_s3_bucket_details(_boto3_session, _bucket_data): + return bucket_data + + def _patched_sync_s3_notifications(_neo4j_session, _boto3_session, _bucket_data, _update_tag): + cartography_s3._load_s3_notifications(_neo4j_session, bucket_notifications, _update_tag) + + # Apply patches without restoring originals by request + setattr(cartography_s3, "get_s3_bucket_list", _patched_get_s3_bucket_list) + setattr(cartography_s3, "get_s3_bucket_details", _patched_get_s3_bucket_details) + setattr(cartography_s3, "_sync_s3_notifications", _patched_sync_s3_notifications) + + cartography_s3.sync( + neo4j_session, + None, + account_id, + update_tag, + common_job_parameters, + ) + + # Stats + bucket_notifications_count = 0 + for bn in bucket_notifications: + bucket_notifications_count += len(bn.get("TopicConfigurations", []) or []) + bucket_notifications_count += len(bn.get("QueueConfigurations", []) or []) + bucket_notifications_count += len(bn.get("LambdaFunctionConfigurations", []) or []) + + return {"buckets": len(bucket_data.get("Buckets", [])), "notifications": bucket_notifications_count} + + +def _build_s3_bucket_list( + tenant_id: str, + provider_id: str, + scan_id: str | None, + regions: list[str], +) -> dict[str, Any]: + bucket_items: list[dict[str, Any]] = [] + + with rls_transaction(tenant_id): + base_qs = Resource.objects.filter( + provider_id=provider_id, + service="s3", + type__in=["bucket", "s3_bucket"], + ) + if scan_id: + rss_ids = ResourceScanSummary.objects.filter( + tenant_id=tenant_id, scan_id=scan_id, service="s3" + ).values_list("resource_id", flat=True) + base_qs = base_qs.filter(id__in=list(rss_ids)) + if regions: + base_qs = base_qs.filter(region__in=regions) + + resources = list(base_qs.only("name", "region", "metadata", "uid", "inserted_at")) + + owner: dict[str, Any] = {"DisplayName": None, "ID": None} + + for r in resources: + name = r.name or _s3_derive_bucket_name_from_uid(r.uid) + creation_date = ( + datetime.now(tz=timezone.utc).isoformat() + if not getattr(r, "inserted_at", None) + else r.inserted_at.replace(tzinfo=timezone.utc).isoformat() + ) + bucket_items.append({"Name": name, "CreationDate": creation_date, "Region": r.region}) + + return {"Owner": owner, "Buckets": bucket_items} + + +def _s3_derive_bucket_name_from_uid(uid: str) -> str: + if uid and ":s3:::" in uid: + try: + return uid.split(":s3:::", 1)[1] + except Exception: + return uid + return uid or "unknown-bucket" + + +def _build_s3_bucket_details( + bucket_data: dict[str, Any], +) -> Generator[ + tuple[ + str, + dict[str, Any] | None, + dict[str, Any] | None, + dict[str, Any] | None, + dict[str, Any] | None, + dict[str, Any] | None, + dict[str, Any] | None, + dict[str, Any] | None, + ], + None, + None, +]: + for b in bucket_data.get("Buckets", []): + yield (b["Name"], None, None, None, None, None, None, None) + + +def _build_s3_notifications( + tenant_id: str, + provider_id: str, + scan_id: str | None, + regions: list[str], +) -> list[dict[str, Any]]: + """ + Return a list of per-bucket notifications. + + Shape (best-effort, Cartography-compatible): + { + "bucket": , + "TopicConfigurations": [...], + "QueueConfigurations": [...], + "LambdaFunctionConfigurations": [...], + "EventBridgeConfiguration": {...} | None, + } + """ + + notifications: list[dict[str, Any]] = [] + + with rls_transaction(tenant_id): + base_qs = Resource.objects.filter( + provider_id=provider_id, + service="s3", + type__in=["bucket", "s3_bucket"], + ) + if scan_id: + rss_ids = ResourceScanSummary.objects.filter( + tenant_id=tenant_id, scan_id=scan_id, service="s3" + ).values_list("resource_id", flat=True) + base_qs = base_qs.filter(id__in=list(rss_ids)) + if regions: + base_qs = base_qs.filter(region__in=regions) + + resources = list(base_qs.only("name", "metadata", "details", "uid")) + + for r in resources: + name = r.name or _s3_derive_bucket_name_from_uid(r.uid) + conf_obj: dict[str, Any] | None= None + + for raw in (getattr(r, "metadata", None), getattr(r, "details", None)): + if not raw: + continue + try: + data = json.loads(raw) if isinstance(raw, str) else raw + except Exception: + continue + + candidate = data.get("notification_config") if isinstance(data, dict) else None + if not candidate and isinstance(data, dict): + has_keys = any( + k in data + for k in ( + "TopicConfigurations", + "QueueConfigurations", + "LambdaFunctionConfigurations", + "EventBridgeConfiguration", + ) + ) + if has_keys: + candidate = data + + if candidate and isinstance(candidate, dict): + conf_obj = candidate + break + + if not conf_obj: + continue + + notification = { + "bucket": name, + "TopicConfigurations": conf_obj.get("TopicConfigurations") or [], + "QueueConfigurations": conf_obj.get("QueueConfigurations") or [], + "LambdaFunctionConfigurations": conf_obj.get("LambdaFunctionConfigurations") or [], + "EventBridgeConfiguration": conf_obj.get("EventBridgeConfiguration"), + } + notifications.append(notification) + + return notifications diff --git a/docker-compose-dev.yml b/docker-compose-dev.yml index 05ed89c397..95c3cf091a 100644 --- a/docker-compose-dev.yml +++ b/docker-compose-dev.yml @@ -78,6 +78,37 @@ services: timeout: 5s retries: 3 + # TODO: Configure Neo4j properly, and also add it to the _production_ compose file + neo4j: + image: neo4j:2025.09.0-community-bullseye + hostname: "neo4j" + volumes: + - ./_data/neo4j:/data + environment: + # Raise memory limits: + - NEO4J_server_memory_pagecache_size=1G + - NEO4J_server_memory_heap_initial__size=1G + - NEO4J_server_memory_heap_max__size=1G + # Auth: + - NEO4J_AUTH=neo4j/neo4j + # Add APOC and GDS: + - apoc.export.file.enabled=true + - apoc.import.file.enabled=true + - apoc.import.file.use_neo4j_config=true + - NEO4J_PLUGINS=["graph-data-science", "apoc"] + - NEO4J_dbms_security_procedures_allowlist=gds.*, apoc.* + - NEO4J_dbms_security_procedures_unrestricted=gds.*, apoc.* + # Networking: + - dbms.connector.bolt.listen_address=0.0.0.0:7687 + ports: + - 7474:7474 + - 7687:7687 + healthcheck: + test: ["CMD", "wget", "--no-verbose", "http://localhost:7474"] + interval: 10s + timeout: 10s + retries: 10 + worker-dev: build: context: ./api