科研技能库/分析作业
数据分析
未发现用户侧风险

分析作业

在Cartography模块中添加一个摄取后分析作业(JSON Cypher文件),以在同步后丰富图形。适用于计算互联网暴露、传播继承权限、链接Human/标准本体节点、风险评分或数据加载后的跨资源分析。

文件预览

2 个文件
references
SKILL.md
7.1 KB · 可预览
---
name: analysis-jobs
description: Add a post-ingestion analysis job (JSON Cypher file) to a Cartography module to enrich the graph after sync. Use when the user asks to compute internet exposure, propagate inherited permissions, link Human / canonical ontology nodes, score risk, or add cross-resource analysis after data is loaded.
---

# analysis-jobs

Analysis jobs are post-ingestion Cypher queries (JSON files) that enrich the graph with computed relationships and properties. They run **after** data is loaded and perform cross-node work that cannot be done during the initial load.

## When to use analysis jobs

Use them when you need to:

1. Compute properties that depend on multiple nodes / relationships.
2. Create relationships that span across resource types.
3. Perform transitive closure (e.g. inherited permissions).
4. Enrich data after all resources of a type are loaded.

**Do NOT** use analysis jobs for:

1. Simple node-to-node relationships (use the data model — see `add-relationship`).
2. Properties that can be computed during `transform()`.
3. Relationships already present in the source data.

## Critical rules

1. **Pick the right scope.** Global jobs run after all accounts/projects/tenants (`run_analysis_job`). Scoped jobs run once per account (`run_scoped_analysis_job`). Use dependency checking (`run_analysis_and_ensure_deps`) when a job needs specific upstream modules.
2. **Use iterative queries for large datasets.** They must return `COUNT(*) AS TotalCompleted`.
3. **Document each query** with `__comment__`.
4. **Clean up stale data** that the analysis job creates (don't leave orphan edges between syncs).
5. **Order statements correctly to avoid read windows.**
    - **Properties:** clean up first (`REMOVE n.attr`), then SET. Cleanup of attributes can usually run in a single transaction.
    - **Relationships:** MERGE first, then DELETE stale (`WHERE r.lastupdated <> $UPDATE_TAG`). Iterative DELETE commits per batch, so a leading DELETE of relationships exposes a graph with those edges missing to concurrent readers until the MERGE finishes. MERGE is idempotent and bumps `r.lastupdated`, so the trailing DELETE only targets edges that genuinely no longer have a current basis. Canonical example: `cartography/data/jobs/analysis/aws_lambda_ecr.json`.

## Instructions

### Step 1 — Pick global vs scoped

| Type    | Runs                                  | Location                                    | Helper                          |
| ------- | ------------------------------------- | ------------------------------------------- | ------------------------------- |
| Global  | Once after all accounts / projects    | `cartography/data/jobs/analysis/`           | `run_analysis_job()`            |
| Scoped  | Once per account / project / tenant   | `cartography/data/jobs/scoped_analysis/`    | `run_scoped_analysis_job()`     |

Examples:

- Internet exposure that needs to see all security groups across all accounts -> **global**.
- IAM instance profile analysis that runs per AWS account -> **scoped**.

### Step 2 — Author the JSON file

```json
{
  "name": "Human-readable name for logging",
  "statements": [
    {
      "__comment__": "Optional comment explaining this query",
      "query": "MATCH (n:NodeType) WHERE ... SET n.property = value",
      "iterative": false
    },
    {
      "__comment__": "Iterative queries for large datasets",
      "query": "MATCH (n:NodeType) WHERE n.property IS NULL WITH n LIMIT $LIMIT_SIZE SET n.property = value RETURN COUNT(*) AS TotalCompleted",
      "iterative": true,
      "iterationsize": 1000
    }
  ]
}
```

### Step 3 — Write the queries

**Non-iterative** — single execution, OK for queries touching a manageable number of nodes:

```json
{
  "query": "MATCH (instance:GCPInstance) WHERE ... SET instance.exposed_internet = true",
  "iterative": false
}
```

**Iterative** — required for large datasets. Must return `TotalCompleted`:

```json
{
  "query": "MATCH (n:Node) WHERE n.stale = true WITH n LIMIT $LIMIT_SIZE DELETE n RETURN COUNT(*) AS TotalCompleted",
  "iterative": true,
  "iterationsize": 1000
}
```

### Step 4 — Available parameters

`common_job_parameters` is forwarded into the query. Typical params:

- `$UPDATE_TAG` — current sync timestamp.
- `$LIMIT_SIZE` — set automatically by the iterative runner.
- Module-specific (`$AWS_ID`, `$PROJECT_ID`, ...).

### Step 5 — Wire the call into your module

#### Pattern A — global analysis at end of ingestion

```python
from cartography.util import run_analysis_job

@timeit
def start_your_module_ingestion(neo4j_session: neo4j.Session, config: Config) -> None:
    common_job_parameters = {"UPDATE_TAG": config.update_tag}

    for account in accounts:
        _sync_one_account(neo4j_session, account, config.update_tag, common_job_parameters)

    run_analysis_job(
        "your_module_exposure_analysis.json",
        neo4j_session,
        common_job_parameters,
    )
```

#### Pattern B — scoped per account/project

```python
from cartography.util import run_scoped_analysis_job

def _sync_one_account(neo4j_session, account_id, update_tag, common_job_parameters):
    common_job_parameters["ACCOUNT_ID"] = account_id

    sync_resources(neo4j_session, account_id, update_tag, common_job_parameters)

    run_scoped_analysis_job(
        "your_module_account_analysis.json",
        neo4j_session,
        common_job_parameters,
    )
```

#### Pattern C — conditional with dependency checking

```python
from cartography.util import run_analysis_and_ensure_deps

def _perform_analysis(requested_syncs, neo4j_session, common_job_parameters):
    run_analysis_and_ensure_deps(
        "your_module_combined_analysis.json",
        {"ec2:instance", "ec2:security_group"},  # required upstream syncs
        set(requested_syncs),
        common_job_parameters,
        neo4j_session,
    )
```

### Step 6 — Test it

Add an integration test that:

1. Calls `sync()` with mocked external boundaries.
2. Asserts the analysis-produced edges / properties using `check_nodes` / `check_rels`.

See the `create-module` skill for testing conventions.

## Best practices

1. **Right scope.** Global runs after all accounts; scoped runs per-account.
2. **Use dep-checking** (`run_analysis_and_ensure_deps`) when a job requires upstream modules.
3. **Document queries** with `__comment__`.
4. **Test analysis jobs** with integration tests.
5. **Use iterative queries** for large datasets.
6. **Clean up stale data** the job creates.

## Common issues

- Job runs before the upstream module — switch to `run_analysis_and_ensure_deps` with the right deps.
- Iterative query never terminates — make sure it returns `COUNT(*) AS TotalCompleted` and the matched set shrinks each iteration.
- Wrong scope — global query reading per-account state can be empty if it runs in the wrong place.

For broader troubleshooting, see the `troubleshooting` skill.

## References (load on demand)

- `references/examples.md` — GCP, AWS, Semgrep wiring examples plus the audit table of modules with proper analysis-job integration.

SKILL.md

元数据
nameanalysis-jobs
description在Cartography模块中添加一个摄取后分析作业(JSON Cypher文件),以在同步后丰富图形。适用于计算互联网暴露、传播继承权限、链接Human/标准本体节点、风险评分或数据加载后的跨资源分析。

analysis-jobs

分析作业是摄取后的Cypher查询(JSON文件),它们通过计算关系和属性来丰富图。它们在数据加载之后运行,执行初始加载期间无法完成的跨节点工作。

何时使用分析作业

当需要实现以下功能时使用:

  1. 计算依赖于多个节点/关系的属性。
  2. 创建跨资源类型的关系。
  3. 执行传递闭包(例如继承权限)。
  4. 在某一类型的所有资源加载完成后丰富数据。

不要将分析作业用于:

  1. 简单的节点间关系(使用数据模型 — 参见 add-relationship)。
  2. 可在 transform() 期间计算的属性。
  3. 源数据中已经存在的关系。

关键规则

  1. 选择正确的范围。 全局作业在所有账户/项目/租户之后运行(run_analysis_job)。按范围作业每个账户运行一次(run_scoped_analysis_job)。当作业依赖特定上游模块时,使用依赖检查(run_analysis_and_ensure_deps)。
  2. 对大数据集使用迭代查询。 它们必须返回 COUNT(*) AS TotalCompleted。
  3. 为每个查询添加文档,使用 __comment__。
  4. 清理过时数据:分析作业创建的数据要清理(不要在两次同步之间留下孤立边)。
  5. 正确排序语句以避免读窗口问题。
    • 属性: 先清理(REMOVE n.attr),然后 SET。属性清理通常可以在单个事务中运行。
    • 关系: 先 MERGE,再 DELETE 过时项(WHERE r.lastupdated <> $UPDATE_TAG)。迭代 DELETE 按批次提交,因此如果在 MERGE 之前执行 DELETE,会导致在 MERGE 完成之前图中缺失这些边,并发读会看到不完整图。MERGE 是幂等的,并会更新 r.lastupdated,因此尾部的 DELETE 仅针对真正失去当前依据的边。标准示例:cartography/data/jobs/analysis/aws_lambda_ecr.json。

操作说明

步骤1 — 选择全局或按范围

类型运行时机位置辅助函数
全局所有账户/项目之后运行一次cartography/data/jobs/analysis/run_analysis_job()
按范围每个账户/项目/租户运行一次cartography/data/jobs/scoped_analysis/run_scoped_analysis_job()

示例:

  • 需要查看所有账户中所有安全组的互联网暴露情况 -> 全局。
  • 按 AWS 账户运行的 IAM 实例配置文件分析 -> 按范围。

步骤2 — 编写 JSON 文件

json
{
  "name": "日志中可读的名称",
  "statements": [
    {
      "__comment__": "解释此查询的可选注释",
      "query": "MATCH (n:NodeType) WHERE ... SET n.property = value",
      "iterative": false
    },
    {
      "__comment__": "大数据集的迭代查询",
      "query": "MATCH (n:NodeType) WHERE n.property IS NULL WITH n LIMIT $LIMIT_SIZE SET n.property = value RETURN COUNT(*) AS TotalCompleted",
      "iterative": true,
      "iterationsize": 1000
    }
  ]
}

步骤3 — 编写查询

非迭代 — 单次执行,适用于涉及节点数可控的查询:

json
{
  "query": "MATCH (instance:GCPInstance) WHERE ... SET instance.exposed_internet = true",
  "iterative": false
}

迭代 — 大数据集必须使用。必须返回 TotalCompleted:

json
{
  "query": "MATCH (n:Node) WHERE n.stale = true WITH n LIMIT $LIMIT_SIZE DELETE n RETURN COUNT(*) AS TotalCompleted",
  "iterative": true,
  "iterationsize": 1000
}

步骤4 — 可用参数

common_job_parameters 会转发到查询中。典型参数:

  • $UPDATE_TAG — 当前同步时间戳。
  • $LIMIT_SIZE — 由迭代运行器自动设置。
  • 模块特定的($AWS_ID、$PROJECT_ID 等)。

步骤5 — 将调用接入你的模块

模式 A — 在摄取末尾的全局分析

python
from cartography.util import run_analysis_job

@timeit
def start_your_module_ingestion(neo4j_session: neo4j.Session, config: Config) -> None:
    common_job_parameters = {"UPDATE_TAG": config.update_tag}

    for account in accounts:
        _sync_one_account(neo4j_session, account, config.update_tag, common_job_parameters)

    run_analysis_job(
        "your_module_exposure_analysis.json",
        neo4j_session,
        common_job_parameters,
    )

模式 B — 每个账户/项目的按范围分析

python
from cartography.util import run_scoped_analysis_job

def _sync_one_account(neo4j_session, account_id, update_tag, common_job_parameters):
    common_job_parameters["ACCOUNT_ID"] = account_id

    sync_resources(neo4j_session, account_id, update_tag, common_job_parameters)

    run_scoped_analysis_job(
        "your_module_account_analysis.json",
        neo4j_session,
        common_job_parameters,
    )

模式 C — 带依赖检查的条件执行

python
from cartography.util import run_analysis_and_ensure_deps

def _perform_analysis(requested_syncs, neo4j_session, common_job_parameters):
    run_analysis_and_ensure_deps(
        "your_module_combined_analysis.json",
        {"ec2:instance", "ec2:security_group"},  # 所需的上游同步
        set(requested_syncs),
        common_job_parameters,
        neo4j_session,
    )

步骤6 — 测试

添加一个集成测试,该测试:

  1. 调用 sync() 并使用模拟的外部边界。
  2. 使用 check_nodes / check_rels 断言分析产生的边/属性。

有关测试约定,参见 create-module 技能。

最佳实践

  1. 正确的范围。 全局在所有账户后运行;按范围每个账户运行。
  2. 使用依赖检查(run_analysis_and_ensure_deps)当作业依赖上游模块时。
  3. 使用 __comment__ 为查询添加文档。
  4. 用集成测试测试分析作业。
  5. 对大数据集使用迭代查询。
  6. 清理作业创建的过时数据。

常见问题

  • 作业在上游模块之前运行 — 切换为 run_analysis_and_ensure_deps 并指定正确的依赖项。
  • 迭代查询永远不会终止 — 确保它返回 COUNT(*) AS TotalCompleted 并且每次迭代匹配集都会缩小。
  • 范围错误 — 如果在错误的位置运行,读取按账户状态的全局查询可能为空。

更广泛的故障排除请参见 troubleshooting 技能。

参考资料(按需加载)

  • references/examples.md — GCP、AWS、Semgrep 集成示例,以及正确集成分析作业的模块审计表。