收集 Asset Panda 日志

解析器版本:1.0

支持的平台:

本文档介绍了如何使用 Google Cloud Storage V2 将 Asset Panda 日志注入到 Google Security Operations。

Asset Panda 是一款基于云的资产管理平台,可让组织在整个生命周期内跟踪、管理和报告实体资产和数字资产。Asset Panda REST API 可用于访问变更日志,这些日志会记录所有资产组的字段级修改、资产转移、操作历史记录和用户活动。这些更改日志可以通过 API 收集,并写入 GCS 存储桶以供 Google SecOps 提取。

准备工作

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

  • Google SecOps 实例
  • 已启用 Cloud Storage API 的 GCP 项目
  • 创建和管理 GCS 存储分区的权限
  • 创建 Cloud Run 服务、Pub/Sub 主题和 Cloud Scheduler 作业的权限
  • 具有管理员访问权限的 Asset Panda 账号(API 配置必需)
  • 具有读取权限的 Asset Panda API 密钥和 Secret

生成 Asset Panda API 凭据

  1. 以管理员身份登录您的 Asset Panda 账号。
  2. 点击右上角的设置(齿轮)图标。
  3. 选择 API 配置

创建 API 密钥

  1. API 配置页面上,点击 Create New API Key
  2. 提供以下配置详细信息:

    • 名称:为密钥输入一个描述性名称(例如 SIEM Change Log Collector
    • 权限:选择读取
  3. 点击保存

  4. 记录页面底部显示的以下凭据:

    • API 密钥(也称为客户端 ID)
    • API 密钥(也称为客户端密钥)

识别群组 ID

Cloud Run 函数需要群组 ID 才能查询对象及其更改日志。如需查找群组 ID,请执行以下操作:

  1. Asset Panda 中,依次前往设置 > 群组设置
  2. 点击群组名称旁边的修改链接。
  3. 记下浏览器网址末尾附近的数值。这是群组 ID(例如 12345)。
  4. 针对您要收集更改日志的每个群组重复执行上述步骤。

测试 API 访问权限

  • 在继续进行集成之前,请先测试您的凭据:

    API_KEY="your-api-key"
    API_SECRET="your-api-secret"
    
    # Test authentication by retrieving account settings
    curl -s "https://api.assetpanda.com/v3/settings" \
      -H "Accept: application/json" \
      -H "Access-Key-Id: ${API_KEY}" \
      -H "Access-Key-Secret: ${API_SECRET}"
    

创建 Google Cloud Storage 存储桶

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

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

为 Cloud Run 函数创建服务账号

Cloud Run 函数需要一个服务账号,该账号具有写入 GCS 存储桶的权限,并且可以由 Pub/Sub 调用。

创建服务账号

  1. GCP 控制台中,依次前往 IAM 和管理 > 服务账号
  2. 点击创建服务账号
  3. 提供以下配置详细信息:
    • 服务账号名称:输入 asset-panda-logs-collector-sa
    • 服务账号说明:输入 Service account for Cloud Run function to collect Asset Panda change 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. 提供以下配置详细信息:
    • 添加主账号:输入服务账号电子邮件地址(例如 asset-panda-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:输入 asset-panda-logs-trigger
    • 将其他设置保留为默认值
  4. 点击创建

创建 Cloud Run 函数以收集日志

Cloud Run 函数将由来自 Cloud Scheduler 的 Pub/Sub 消息触发,以从 Asset Panda REST API 中提取更改日志并将其写入 GCS。

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

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

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

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

  8. 前往安全性标签页:

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

    1. 点击变量和密钥
    2. 为每个环境变量点击+ 添加变量
    变量名称 示例值 说明
    GCS_BUCKET asset-panda-change-logs GCS 存储桶名称
    GCS_PREFIX asset-panda 日志文件的前缀
    STATE_KEY asset-panda/state.json 状态文件路径
    AP_API_KEY your-api-key Asset Panda API 密钥
    AP_API_SECRET your-api-secret Asset Panda API 密钥
    AP_GROUP_IDS 12345,67890 要监控的以逗号分隔的群组 ID
    PAGE_SIZE 50 每个 API 页面的记录数(最多 50 条)
    MAX_RECORDS 10000 每次运行的记录数上限
    LOOKBACK_HOURS 2 初始回溯期
  10. 变量和 Secret 部分中,向下滚动到请求

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

    • 资源部分中:
      • 内存:选择 512 MiB 或更高值
      • CPU:选择 1
  12. 修订版本伸缩部分中:

    • 实例数下限:输入 0
    • 实例数上限:输入 100(或根据预期负载进行调整)
  13. 点击创建

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

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

添加函数代码

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

    • 第一个文件 - main.py:

      import functions_framework
      from google.cloud import storage
      import json
      import os
      import urllib3
      from datetime import datetime, timezone, timedelta
      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()
      
      # Environment variables
      GCS_BUCKET = os.environ.get('GCS_BUCKET')
      GCS_PREFIX = os.environ.get('GCS_PREFIX', 'asset-panda').strip('/')
      STATE_KEY = os.environ.get('STATE_KEY') or f"{GCS_PREFIX}/state.json"
      AP_API_KEY = os.environ.get('AP_API_KEY')
      AP_API_SECRET = os.environ.get('AP_API_SECRET')
      AP_GROUP_IDS = os.environ.get('AP_GROUP_IDS', '')
      PAGE_SIZE = int(os.environ.get('PAGE_SIZE', '50'))
      MAX_RECORDS = int(os.environ.get('MAX_RECORDS', '10000'))
      LOOKBACK_HOURS = int(os.environ.get('LOOKBACK_HOURS', '2'))
      
      BASE_URL = 'https://api.assetpanda.com'
      
      # Rate limiting: 400 calls per 3 minutes
      RATE_LIMIT_CALLS = 400
      RATE_LIMIT_WINDOW = 180  # seconds
      call_timestamps = []
      
      def rate_limit_wait():
        """Enforce rate limiting of 400 calls per 3-minute window."""
        global call_timestamps
        now = time.time()
        call_timestamps = [t for t in call_timestamps if now - t < RATE_LIMIT_WINDOW]
        if len(call_timestamps) >= RATE_LIMIT_CALLS:
          sleep_time = RATE_LIMIT_WINDOW - (now - call_timestamps[0]) + 1
          print(f"Rate limit approaching. Sleeping {sleep_time:.1f}s...")
          time.sleep(sleep_time)
          call_timestamps = [t for t in call_timestamps if time.time() - t < RATE_LIMIT_WINDOW]
        call_timestamps.append(time.time())
      
      def api_request(method, path, params=None):
        """Make an authenticated request to the Asset Panda API."""
        rate_limit_wait()
      
        url = f"{BASE_URL}{path}"
        headers = {
          'Accept': 'application/json',
          'Content-Type': 'application/json',
          'Access-Key-Id': AP_API_KEY,
          'Access-Key-Secret': AP_API_SECRET,
        }
      
        backoff = 1.0
        for attempt in range(3):
          try:
            if method == 'GET':
              response = http.request(method, url, headers=headers, fields=params)
            else:
              body = json.dumps(params or {}).encode('utf-8')
              response = http.request(method, url, headers=headers, body=body)
      
            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, 60.0)
              continue
      
            if response.status == 200:
              return json.loads(response.data.decode('utf-8'))
      
            print(f"HTTP {response.status}: {response.data.decode('utf-8')[:500]}")
            return None
      
          except Exception as e:
            print(f"Request error (attempt {attempt + 1}): {e}")
            if attempt < 2:
              time.sleep(backoff)
              backoff *= 2
      
        return None
      
      def parse_datetime(value):
        """Parse ISO datetime string to datetime object."""
        if not value:
          return None
        if value.endswith("Z"):
          value = value[:-1] + "+00:00"
        try:
          return datetime.fromisoformat(value)
        except Exception:
          return None
      
      @functions_framework.cloud_event
      def main(cloud_event):
        """
        Cloud Run function triggered by Pub/Sub to fetch
        Asset Panda change logs and write to GCS.
      
        Args:
          cloud_event: CloudEvent object containing Pub/Sub message
        """
        if not all([GCS_BUCKET, AP_API_KEY, AP_API_SECRET, AP_GROUP_IDS]):
          print('Error: Missing required environment variables')
          return
      
        group_ids = [gid.strip() for gid in AP_GROUP_IDS.split(',') if gid.strip()]
        if not group_ids:
          print('Error: No group IDs configured')
          return
      
        try:
          bucket = storage_client.bucket(GCS_BUCKET)
      
          # Load state
          state = load_state(bucket, STATE_KEY)
      
          # Determine time window
          now = datetime.now(timezone.utc)
          last_time = None
      
          if isinstance(state, dict) and state.get("last_run_time"):
            try:
              last_time = parse_datetime(state["last_run_time"])
              if last_time:
                last_time = last_time - timedelta(minutes=2)
            except Exception as e:
              print(f"Warning: Could not parse last_run_time: {e}")
      
          if last_time is None:
            last_time = now - timedelta(hours=LOOKBACK_HOURS)
      
          print(f"Collecting change logs from {last_time.isoformat()} to {now.isoformat()}")
      
          all_records = []
      
          for group_id in group_ids:
            if len(all_records) >= MAX_RECORDS:
              print(f"Reached max_records limit ({MAX_RECORDS})")
              break
      
            records = fetch_group_change_logs(group_id, last_time, now)
            all_records.extend(records)
            print(f"Group {group_id}: collected {len(records)} change log entries")
      
          if not all_records:
            print("No new change log records found.")
            save_state(bucket, STATE_KEY, now.isoformat())
            return
      
          # Write to GCS as NDJSON
          timestamp = now.strftime('%Y%m%d_%H%M%S')
          object_key = f"{GCS_PREFIX}/logs_{timestamp}.ndjson"
          blob = bucket.blob(object_key)
      
          ndjson = '\n'.join(
            [json.dumps(record, ensure_ascii=False) for record in all_records]
          ) + '\n'
          blob.upload_from_string(ndjson, content_type='application/x-ndjson')
      
          print(f"Wrote {len(all_records)} records to gs://{GCS_BUCKET}/{object_key}")
      
          # Update state
          save_state(bucket, STATE_KEY, now.isoformat())
          print(f"Successfully processed {len(all_records)} records")
      
        except Exception as e:
          print(f'Error processing change logs: {str(e)}')
          raise
      
      def fetch_group_change_logs(group_id, start_time, end_time):
        """
        Fetch change logs for all objects in a group.
      
        Args:
          group_id: Asset Panda group ID
          start_time: Start time for log query
          end_time: End time for log query
      
        Returns:
          List of change log records
        """
        records = []
        offset = 0
      
        # Iterate through objects in the group
        while True:
          if len(records) >= MAX_RECORDS:
            break
      
          search_data = api_request(
            'POST',
            f'/v3/groups/{group_id}/search_objects',
            {'offset': offset, 'limit': PAGE_SIZE}
          )
      
          if not search_data:
            break
      
          objects = search_data.get('objects', [])
          if not objects:
            break
      
          for obj in objects:
            object_id = obj.get('id')
            if not object_id:
              continue
      
            change_logs = fetch_object_change_logs(
              object_id, group_id, start_time, end_time
            )
            records.extend(change_logs)
      
            if len(records) >= MAX_RECORDS:
              break
      
          total = search_data.get('totals', {}).get('objects', 0)
          offset += len(objects)
          if offset >= total:
            break
      
        return records
      
      def fetch_object_change_logs(object_id, group_id, start_time, end_time):
        """
        Fetch change logs for a specific object.
      
        Args:
          object_id: Asset Panda object ID
          group_id: Asset Panda group ID
          start_time: Start time filter
          end_time: End time filter
      
        Returns:
          List of change log records with metadata
        """
        records = []
        offset = 0
      
        while True:
          data = api_request(
            'GET',
            f'/v3/entity_objects/{object_id}/change_logs',
            {'limit': PAGE_SIZE, 'offset': offset}
          )
      
          if not data:
            break
      
          logs = data.get('change_logs', data.get('data', []))
          if isinstance(logs, dict):
            logs = [logs]
      
          if not logs:
            break
      
          for log_entry in logs:
            log_time = parse_datetime(
              log_entry.get('created_at') or log_entry.get('updated_at', '')
            )
      
            if log_time and log_time < start_time:
              return records
            if log_time and log_time > end_time:
              continue
      
            log_entry['_group_id'] = group_id
            log_entry['_object_id'] = object_id
            records.append(log_entry)
      
          if len(logs) < PAGE_SIZE:
            break
      
          offset += len(logs)
      
        return records
      
      def load_state(bucket, key):
        """Load state from GCS."""
        try:
          blob = bucket.blob(key)
          if blob.exists():
            state_data = blob.download_as_text()
            return json.loads(state_data)
        except Exception as e:
          print(f"Warning: Could not load state: {e}")
        return {}
      
      def save_state(bucket, key, last_run_time_iso):
        """Save the last run timestamp to GCS state file."""
        try:
          state = {'last_run_time': last_run_time_iso}
          blob = bucket.blob(key)
          blob.upload_from_string(
            json.dumps(state, indent=2),
            content_type='application/json'
          )
          print(f"Saved state: last_run_time={last_run_time_iso}")
        except Exception as e:
          print(f"Warning: Could not save state: {e}")
      
    • 第二个文件 - requirements.txt:

      functions-framework==3.*
      google-cloud-storage==2.*
      urllib3>=2.0.0
      
  3. 点击部署以保存并部署该函数。

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

创建 Cloud Scheduler 作业

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

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

    设置
    名称 asset-panda-logs-collector-hourly
    区域 选择与 Cloud Run 函数相同的区域
    频率 0 * * * *(每小时一次,整点时)
    时区 选择时区(建议选择 UTC)
    目标类型 Pub/Sub
    主题 选择主题 asset-panda-logs-trigger
    消息正文 {}(空 JSON 对象)
  4. 点击创建

时间表频率选项

根据日志量和延迟时间要求选择频次:

频率 Cron 表达式 使用场景
每小时 0 * * * * 标准(推荐)
每 2 小时 0 */2 * * * 调低音量
每 6 小时 0 */6 * * * 低成交量、批处理

测试集成

  1. Cloud Scheduler 控制台中,找到您的作业。
  2. 点击强制运行以手动触发作业。
  3. 等待几秒钟。
  4. 前往 Cloud Run > 服务
  5. 点击函数名称 asset-panda-logs-collector
  6. 点击日志标签页。
  7. 验证函数是否已成功执行。查找:

    Collecting change logs from YYYY-MM-DDTHH:MM:SS+00:00 to YYYY-MM-DDTHH:MM:SS+00:00
    Group 12345: collected X change log entries
    Wrote X records to gs://asset-panda-change-logs/asset-panda/logs_YYYYMMDD_HHMMSS.ndjson
    Successfully processed X records
    
  8. 前往 Cloud Storage > 存储分区

  9. 点击您的存储桶名称。

  10. 前往前缀文件夹 asset-panda/

  11. 验证是否已创建具有当前时间戳的新 .ndjson 文件。

如果您在日志中看到错误,请执行以下操作:

  • HTTP 401:检查环境变量中的 API 凭据。验证 API 密钥和 Secret 是否正确。
  • HTTP 403:验证 API 密钥是否已启用读取权限。
  • HTTP 429:速率限制 - 函数将自动重试并进行退避。
  • 结果为空:验证配置的群组 ID 是否存在,以及是否包含具有更改历史记录的对象。
  • 缺少环境变量:检查是否已设置所有必需的变量。

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

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

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

  9. 点击下一步

  10. 为以下输入参数指定值:

    • 存储桶网址:输入带有前缀路径的 GCS 存储桶 URI:

      gs://asset-panda-change-logs/asset-panda/
      
      • 替换:
        • asset-panda-change-logs:您的 GCS 存储桶名称。
        • asset-panda:存储日志的可选前缀/文件夹路径(留空表示根目录)。
    • 来源删除选项:根据您的偏好选择删除选项:

      • 永不:转移后永不删除任何文件(建议用于测试)。
      • 删除已转移的文件:在成功转移后删除文件。
      • 删除已转移的文件和空目录:成功转移后删除文件和空目录。
    • 文件存在时间上限:包含在过去指定天数内修改的文件(默认值为 180 天)

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

    • 注入标签:要应用于此 Feed 中事件的标签

  11. 点击下一步

  12. 最终确定界面中查看新的 Feed 配置,然后点击提交

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

Google SecOps 服务账号需要您的 GCS 存储桶的 Storage Object Viewer 角色。

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

UDM 映射表

日志字段 UDM 映射 逻辑
array_values_label additional.fields 已合并
data_label additional.fields 已合并
field_values_label additional.fields 已合并
low_values_label additional.fields 已合并
value_ids_label additional.fields 已合并
values_label additional.fields 已合并
display_with_secondary intermediary.user.user_display_name 直接映射
secondary_name intermediary.user.userid 直接映射
created_at metadata.event_timestamp 解析为 RFC 3339
has_user metadata.event_type 已映射:trueUSER_UNCATEGORIZED
object_version_ids metadata.product_version 直接映射
share_url metadata.url_back_to_product 直接映射
display_name principal.user.user_display_name 直接映射
id principal.user.userid 直接映射
account_id_label security_result.detection_fields 已合并
change_source_label security_result.detection_fields 已合并
change_trigger_label security_result.detection_fields 已合并
created_at_label security_result.detection_fields 已合并
date_format_label security_result.detection_fields 已合并
docusign_envelope_id_label security_result.detection_fields 已合并
embedded_into_object_id_label security_result.detection_fields 已合并
entity_action_id_label security_result.detection_fields 已合并
entity_id_label security_result.detection_fields 已合并
entity_id_label1 security_result.detection_fields 已合并
entity_key_label security_result.detection_fields 已合并
google_calendar_sync_label security_result.detection_fields 已合并
gps_coordinates_label security_result.detection_fields 已合并
has_audit_history_label security_result.detection_fields 已合并
id_label security_result.detection_fields 已合并
is_archived_label security_result.detection_fields 已合并
is_deletable_label security_result.detection_fields 已合并
is_editable_label security_result.detection_fields 已合并
is_locked_label security_result.detection_fields 已合并
linked_action_object_id_label security_result.detection_fields 已合并
modifier_id_label security_result.detection_fields 已合并
next_step_reservation_uid_label security_result.detection_fields 已合并
object_appreciation_label security_result.detection_fields 已合并
object_depreciation_label security_result.detection_fields 已合并
oid_label security_result.detection_fields 已合并
old_id_label security_result.detection_fields 已合并
parent_action_object_id_label security_result.detection_fields 已合并
predefined_forms_label security_result.detection_fields 已合并
reservation_notification_label security_result.detection_fields 已合并
reservation_uid_label security_result.detection_fields 已合并
returned_label security_result.detection_fields 已合并
state_label security_result.detection_fields 已合并
status_label security_result.detection_fields 已合并
updated_at_label security_result.detection_fields 已合并
user_id_label security_result.detection_fields 已合并
version_label security_result.detection_fields 已合并
不适用 metadata.event_type 常量:USER_UNCATEGORIZED

更新日志

查看相应解析器的更改日志

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