收集 MuleSoft Anypoint 平台日志

支持:

本文档介绍了如何使用 Google Cloud Storage 将 MuleSoft Anypoint 平台日志中的审核跟踪事件注入到 Google Security Operations。

准备工作

请确保满足以下前提条件:

  • Google SecOps 实例
  • 启用了 Cloud Storage API 的 GCP 项目
  • 创建和管理 GCS 存储分区的权限
  • 管理 GCS 存储分区的 IAM 政策的权限
  • 创建 Cloud Run 函数、Pub/Sub 主题和 Cloud Scheduler 作业的权限
  • 创建服务账号的权限
  • 对 MuleSoft Anypoint 平台的特权访问权限

获取 MuleSoft 组织 ID

  1. 登录 Anypoint 平台
  2. 前往访问管理 > 组织
  3. 业务组 表中,点击您组织的名称。
  4. 复制组织 ID (例如 0a12b3c4-d5e6-789f-1021-1a2b34cd5e6f)。

或者,前往 MuleSoft 业务组,然后从网址中复制 ID。

创建 MuleSoft 关联的应用

  1. 登录 Anypoint 平台
  2. 前往访问管理 > 关联的应用 > 创建应用
  3. 提供以下配置详细信息:
    • 应用名称:输入唯一名称(例如 Google SecOps export)。
    • 选择应用以自己的名义进行操作(客户端凭据)
  4. 依次点击添加范围 > 审核日志查看器 > 下一步
  5. 选择您需要其日志的每个业务组。
  6. 依次点击下一步 > 添加范围
  7. 点击保存 ,然后复制客户端 ID客户端密钥

创建 Google Cloud Storage 存储桶

  1. 转到 Google Cloud 控制台
  2. 选择您的项目,或创建新项目。
  3. 在导航菜单中,前往 Cloud Storage > 存储分区
  4. 点击创建存储桶
  5. 提供以下配置详细信息:

    设置
    指定存储桶的名称 输入全局唯一的名称(例如 mulesoft-audit-logs
    位置类型 根据您的需求进行选择(区域、双区域、多区域)
    位置 选择位置(例如 us-central1
    存储类别 标准(建议用于经常访问的日志)
    访问权限控制 统一(推荐)
    保护工具 可选:启用对象版本控制或保留政策
  6. 点击创建

为 Cloud Run 函数创建服务帐号

Cloud Run 函数需要一个具有写入 GCS 存储桶权限的服务帐号。

创建服务帐号

  1. GCP 控制台 中,前往 IAM 和管理 > 服务账号
  2. 点击创建服务账号
  3. 提供以下配置详细信息:
    • 服务账号名称:输入 mulesoft-logs-collector-sa
    • 服务账号说明:输入 Service account for Cloud Run function to collect MuleSoft Anypoint logs
  4. 点击创建并继续
  5. 向此服务帐号授予对项目的访问权限 部分:
    1. 点击选择角色
    2. 搜索并选择 Storage Object Admin
    3. 点击 + 添加其他角色
    4. 搜索并选择 Cloud Run Invoker
    5. 点击 + 添加其他角色
    6. 搜索并选择 Cloud Functions Invoker
  6. 点击继续
  7. 点击完成

这些角色是以下操作所必需的:

  • Storage Object Admin:将日志写入 GCS 存储桶并管理状态文件
  • Cloud Run Invoker:允许 Pub/Sub 调用该函数
  • Cloud Functions Invoker:允许函数调用

授予对 GCS 存储桶的 IAM 权限

向服务帐号授予对 GCS 存储桶的写入权限:

  1. 前往 Cloud Storage > 存储分区
  2. 点击您的存储桶名称。
  3. 前往权限 标签页。
  4. 点击授予访问权限
  5. 提供以下配置详细信息:
    • 添加主账号:输入服务帐号电子邮件地址(例如 mulesoft-logs-collector-sa@PROJECT_ID.iam.gserviceaccount.com)。
    • 分配角色:选择 Storage Object Admin
  6. 点击保存

创建 Pub/Sub 主题

创建一个 Pub/Sub 主题,Cloud Scheduler 将向该主题发布消息,而 Cloud Run 函数将订阅该主题。

  1. GCP 控制台 中,前往 Pub/Sub > 主题
  2. 点击创建主题
  3. 提供以下配置详细信息:
    • 主题 ID:输入 mulesoft-audit-trigger
    • 其他设置保留默认值。
  4. 点击创建

创建 Cloud Run 函数以收集日志

Cloud Run 函数由 Cloud Scheduler 中的 Pub/Sub 消息触发,以从 MuleSoft Anypoint API 中提取日志并将其写入 GCS。

  1. GCP 控制台 中,前往 Cloud Run
  2. 点击创建服务
  3. 选择函数 (使用内嵌编辑器创建函数)。
  4. 配置 部分中,提供以下配置详细信息:

    设置
    服务名称 mulesoft-audit-collector
    区域 选择与您的 GCS 存储桶匹配的区域(例如 us-central1
    运行时 选择 Python 3.12 或更高版本
  5. 触发器(可选) 部分中:

    1. 点击 + 添加触发器
    2. 选择 Cloud Pub/Sub
    3. 选择 Cloud Pub/Sub 主题 中,选择主题 mulesoft-audit-trigger
    4. 点击保存
  6. 身份验证 部分中:

    1. 选择需要进行身份验证
    2. 选中 Identity and Access Management (IAM)
  7. 向下滚动并展开容器、网络、安全性

  8. 前往安全性 标签页:

    • 服务账号:选择服务账号 mulesoft-logs-collector-sa
  9. 前往容器 标签页:

    1. 点击变量和密钥
    2. 为每个环境变量点击 + 添加变量
    变量名称 示例值
    MULE_ORG_ID your_org_id
    CLIENT_ID your_client_id
    CLIENT_SECRET your_client_secret
    GCS_BUCKET mulesoft-audit-logs
  10. 变量和密钥 标签页中向下滚动到请求

    • 请求超时:输入 600 秒(10 分钟)。
  11. 前往容器 中的设置 标签页:

    • 资源 部分中:
      • 内存:选择 512 MiB 或更高。
      • CPU:选择 1
    • 点击完成
  12. 滚动到执行环境

    • 选择默认 (推荐)。
  13. 修订版本伸缩 部分中:

    • 最小实例数:输入 0
    • 最大实例数:输入 100(或根据预期负载进行调整)。
  14. 点击创建

  15. 等待服务创建完成(1-2 分钟)。

  16. 服务创建完成后,系统会自动打开内嵌代码编辑器

添加函数代码

  1. 函数入口点 中输入 main
  2. 在内嵌代码编辑器中,创建两个文件:

    • 第一个文件:main.py:
    import functions_framework
    from google.cloud import storage
    import json
    import os
    import urllib3
    from datetime import datetime, timedelta, timezone
    import uuid
    import time
    
    # Initialize HTTP client with timeouts
    http = urllib3.PoolManager(
        timeout=urllib3.Timeout(connect=5.0, read=30.0),
        retries=False,
    )
    
    # Initialize Storage client
    storage_client = storage.Client()
    
    # MuleSoft API endpoints
    TOKEN_URL = "https://anypoint.mulesoft.com/accounts/api/v2/oauth2/token"
    
    @functions_framework.cloud_event
    def main(cloud_event):
        """
        Cloud Run function triggered by Pub/Sub to fetch MuleSoft audit logs and write to GCS.
    
        Args:
            cloud_event: CloudEvent object containing Pub/Sub message
        """
    
        # Get environment variables
        org_id = os.environ.get('MULE_ORG_ID')
        client_id = os.environ.get('CLIENT_ID')
        client_secret = os.environ.get('CLIENT_SECRET')
        bucket_name = os.environ.get('GCS_BUCKET')
    
        if not all([org_id, client_id, client_secret, bucket_name]):
            print('Error: Missing required environment variables')
            return
    
        query_url = f"https://anypoint.mulesoft.com/audit/v2/organizations/{org_id}/query"
    
        try:
            # Get GCS bucket
            bucket = storage_client.bucket(bucket_name)
    
            # Get OAuth token
            token = get_token(client_id, client_secret)
    
            # Calculate time range (last 24 hours)
            now = datetime.now(timezone.utc).replace(microsecond=0)
            start = now - timedelta(days=1)
    
            print(f'Fetching audit logs from {start.isoformat()} to {now.isoformat()}')
    
            # Fetch audit logs
            events = list(fetch_audit(query_url, token, start, now))
    
            # Upload to GCS
            if events:
                upload_to_gcs(bucket, events, start)
                print(f'Uploaded {len(events)} events')
            else:
                print('No events in the last 24 hours')
    
        except Exception as e:
            print(f'Error processing logs: {str(e)}')
            raise
    
    def get_token(client_id, client_secret):
        """Get OAuth 2.0 access token from MuleSoft."""
        data = {
            'grant_type': 'client_credentials',
            'client_id': client_id,
            'client_secret': client_secret
        }
    
        encoded_data = urllib3.request.urlencode(data).encode('utf-8')
    
        backoff = 1.0
        max_retries = 3
    
        for attempt in range(max_retries):
            try:
                response = http.request(
                    'POST',
                    TOKEN_URL,
                    body=encoded_data,
                    headers={'Content-Type': 'application/x-www-form-urlencoded'}
                )
    
                if response.status == 429:
                    retry_after = int(response.headers.get('Retry-After', str(int(backoff))))
                    print(f'Rate limited (429) on token request. Retrying after {retry_after}s...')
                    time.sleep(retry_after)
                    backoff = min(backoff * 2, 30.0)
                    continue
    
                if response.status != 200:
                    raise Exception(f'Failed to get token: {response.status} - {response.data.decode()}')
    
                token_data = json.loads(response.data.decode('utf-8'))
                return token_data['access_token']
    
            except Exception as e:
                if attempt == max_retries - 1:
                    raise
                print(f'Token request failed (attempt {attempt + 1}/{max_retries}): {e}')
                time.sleep(backoff)
                backoff = min(backoff * 2, 30.0)
    
        raise Exception('Failed to get token after maximum retries')
    
    def fetch_audit(query_url, token, start, end):
        """Fetch audit logs from MuleSoft API with pagination."""
        headers = {
            'Authorization': f'Bearer {token}',
            'Content-Type': 'application/json'
        }
    
        body = {
            'startDate': f"{start.isoformat(timespec='milliseconds')}Z",
            'endDate': f"{end.isoformat(timespec='milliseconds')}Z",
            'limit': 200,
            'offset': 0,
            'ascending': False
        }
    
        backoff = 1.0
    
        while True:
            try:
                response = http.request(
                    'POST',
                    query_url,
                    body=json.dumps(body).encode('utf-8'),
                    headers=headers
                )
    
                # Handle rate limiting with exponential backoff
                if response.status == 429:
                    retry_after = int(response.headers.get('Retry-After', str(int(backoff))))
                    print(f'Rate limited (429). Retrying after {retry_after}s...')
                    time.sleep(retry_after)
                    backoff = min(backoff * 2, 30.0)
                    continue
    
                backoff = 1.0
    
                if response.status != 200:
                    print(f'HTTP Error: {response.status}')
                    response_text = response.data.decode('utf-8')
                    print(f'Response body: {response_text}')
                    break
    
                data = json.loads(response.data.decode('utf-8'))
    
                if not data.get('data'):
                    break
    
                yield from data['data']
                body['offset'] += body['limit']
    
            except Exception as e:
                print(f'Error fetching audit logs: {e}')
                break
    
    def upload_to_gcs(bucket, events, timestamp):
        """Upload events to GCS as compressed JSON."""
        import gzip
        import io
    
        # Create blob name with timestamp and UUID
        blob_name = f"{timestamp.strftime('%Y/%m/%d')}/mulesoft-audit-{uuid.uuid4()}.json.gz"
    
        # Compress events
        buf = io.BytesIO()
        with gzip.GzipFile(fileobj=buf, mode='w') as gz:
            for event in events:
                gz.write((json.dumps(event) + '\n').encode('utf-8'))
    
        buf.seek(0)
    
        # Upload to GCS
        blob = bucket.blob(blob_name)
        blob.upload_from_file(buf, content_type='application/gzip')
    
        print(f'Uploaded to gs://{bucket.name}/{blob_name}')
    
    • 第二个文件:requirements.txt:
    functions-framework==3.*
    google-cloud-storage==2.*
    urllib3>=2.0.0
    
  3. 点击部署 以保存和部署函数。

  4. 等待部署完成(2-3 分钟)。

重要注意事项

速率限制 :审核日志查询端点在三个控制平面中按 IP 应用速率限制。美国控制平面允许每个 IP 每分钟 700 个请求,而欧盟和政府控制平面允许每个 IP 每分钟 40 个请求。该函数实现了指数退避算法,以自动处理速率限制。

令牌到期 :访问令牌通常会在颁发后约 30 到 60 分钟到期。该函数会为每次执行请求新令牌。对于频繁执行的生产环境部署,请考虑使用刷新逻辑实现令牌缓存。

审核日志保留 :审核日志的默认保留期限为一年。如果您的组织是在 2023 年 7 月 10 日之前创建的,并且您没有手动更改保留期限,则保留期限为六年。如果您需要在配置的保留期限过后继续保留日志,请定期下载日志。

创建 Cloud Scheduler 作业

Cloud Scheduler 会定期向 Pub/Sub 主题发布消息,从而触发 Cloud Run 函数。

  1. GCP 控制台 中,前往 Cloud Scheduler
  2. 点击创建作业
  3. 提供以下配置详细信息:

    设置
    名称 daily-mulesoft-audit-export
    区域 选择与 Cloud Run 函数相同的区域
    频率 0 2 * * *(每天 02:00 UTC 运行)
    时区 选择时区(建议使用 UTC)
    目标类型 Pub/Sub
    主题 选择主题 mulesoft-audit-trigger
    邮件正文 {}(空 JSON 对象)
  4. 点击创建

测试调度器作业

  1. Cloud Scheduler 控制台中,找到您的作业。
  2. 点击强制运行 以手动触发。
  3. 等待几秒钟,然后前往 Cloud Run > 服务 > mulesoft-audit-collector > 日志
  4. 验证函数是否已成功执行。
  5. 检查 GCS 存储桶以确认日志已写入。

检索 Google SecOps 服务帐号

Google SecOps 使用唯一的服务帐号从您的 GCS 存储桶中读取数据。您必须向此服务帐号授予对您的存储桶的访问权限。

获取服务帐号电子邮件地址

  1. 前往 SIEM 设置 > Feed
  2. 点击 Add New Feed (添加新 Feed)。
  3. 点击配置单个 Feed
  4. Feed 名称 字段中,输入 Feed 的名称(例如 MuleSoft Logs)。
  5. 选择 Google Cloud Storage V2 作为来源类型
  6. 选择 Mulesoft 作为日志类型
  7. 点击 Get Service Account (获取服务账号)。系统会显示唯一的服务帐号电子邮件地址,例如:

    chronicle-12345678@chronicle-gcp-prod.iam.gserviceaccount.com
    
  8. 复制此电子邮件地址,以便在下一步中使用。

向 Google SecOps 服务帐号授予 IAM 权限

Google SecOps 服务帐号需要对您的 GCS 存储桶具有 Storage Object Viewer 角色。

  1. 前往 Cloud Storage > 存储分区
  2. 点击您的存储桶名称。
  3. 前往权限 标签页。
  4. 点击授予访问权限
  5. 提供以下配置详细信息:
    • 添加主账号:粘贴 Google SecOps 服务帐号电子邮件地址。
    • 分配角色:选择 Storage Object Viewer
  6. 点击保存

在 Google SecOps 中配置 Feed 以注入 MuleSoft 日志

  1. 前往 SIEM 设置 > Feed
  2. 点击 Add New Feed (添加新 Feed)。
  3. 点击配置单个 Feed
  4. Feed 名称 字段中,输入 Feed 的名称(例如 MuleSoft Logs)。
  5. 选择 Google Cloud Storage V2 作为来源类型
  6. 选择 Mulesoft 作为日志类型
  7. 点击下一步
  8. 为以下输入参数指定值:

    • 存储桶网址:输入 GCS 存储桶 URI:

      gs://mulesoft-audit-logs/
      
      • mulesoft-audit-logs 替换为存储桶的实际名称。
    • 来源删除选项:根据您的偏好选择删除选项:

      • 永不:转移后永不删除任何文件(建议用于测试)。
      • 删除已转移的文件:成功转移后删除文件。
      • 删除已转移的文件和空目录:成功转移后删除文件和空目录。

    • 最长文件存在时间:包含在过去指定天数内修改的文件。默认值为 180 天。

    • 资产命名空间资产命名空间

    • 提取标签:要应用于此 Feed 中的事件的标签。

  9. 点击下一步

  10. 完成 屏幕中检查新 Feed 配置,然后点击提交

需要更多帮助?获得社区成员和 Google SecOps 专业人士的解答。