From 82e5f840fb407ec38a5aa2f385c93b0c586d3a59 Mon Sep 17 00:00:00 2001 From: O M Date: Tue, 29 Sep 2026 16:25:20 -0400 Subject: [PATCH] Write attack-path findings for workload findings that can reach a datastore. A rules run should record the path an operator has to fix, not only leave it as a named query. Co-authored-by: Cursor --- README.md | 2 +- docs/ARCHITECTURE.md | 4 +- docs/ROADMAP.md | 6 +- docs/adr/001-graph-schema-v0.md | 1 + internal/api/web/app.js | 4 + internal/export/findings.go | 2 + internal/graph/attackpath.go | 160 ++++++++++++++++++++ internal/graph/attackpath_test.go | 205 ++++++++++++++++++++++++++ internal/graph/page.go | 1 + internal/graph/pathids_test.go | 19 +++ internal/graph/store.go | 11 +- internal/graph/types.go | 2 + internal/rules/attackpath.go | 145 ++++++++++++++++++ internal/rules/attackpath_test.go | 234 ++++++++++++++++++++++++++++++ internal/rules/engine.go | 19 ++- 15 files changed, 798 insertions(+), 17 deletions(-) create mode 100644 internal/graph/attackpath.go create mode 100644 internal/graph/attackpath_test.go create mode 100644 internal/graph/pathids_test.go create mode 100644 internal/rules/attackpath.go create mode 100644 internal/rules/attackpath_test.go diff --git a/README.md b/README.md index f2b0790..714dd28 100644 --- a/README.md +++ b/README.md @@ -21,7 +21,7 @@ OpenSourceOM Core is the platform behind [opensourceom.org](https://opensourceom Traditional scanners flood you with CVEs and misconfigurations. OpenSourceOM connects the dots — showing which findings sit on paths from the internet to your sensitive data and privileged identities. -> **Status:** Early development (Phase 3). CSPM rules, identity blast radius, Kubernetes ingest, exports, a collector plugin SDK, and a Helm chart are available. CVE findings follow package and image inventory. Datastores can carry a sensitivity mark from a tag or label. Current work is graph accuracy: attack-path findings, rule packs, and cloud audit ingest. See the [roadmap](./docs/ROADMAP.md). +> **Status:** Early development (Phase 3). CSPM rules, identity blast radius, Kubernetes ingest, exports, a collector plugin SDK, and a Helm chart are available. CVE findings follow package and image inventory. Datastores can carry a sensitivity mark from a tag or label. A rules run writes an attack-path finding when a workload finding sits on a path to a datastore. Current work is cloud audit ingest, with further rule packs open for contributors. See the [roadmap](./docs/ROADMAP.md). ## Why this exists diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 563fa58..84bf31a 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -36,7 +36,7 @@ SPDX-License-Identifier: Apache-2.0 | **Graph store** | `internal/graph/`, `migrations/` | PostgreSQL `nodes` + `edges` | | **Path queries** | `internal/graph/query.go` | Named queries including `internet-to-datastore` and `internet-to-sensitive-datastore` | | **Blast radius** | `internal/graph/blastradius.go` | Reachability from identities over `CAN_ACCESS` / `ASSUMES` | -| **CSPM rules** | `internal/rules/` | Built-in policies plus embedded YAML packs | +| **CSPM rules** | `internal/rules/` | Built-in policies, embedded YAML packs, and attack-path findings | | **CVE enrichment** | `internal/enrichment/` | Match workload packages and images, then NVD or a catalog for CVSS | | **Exports** | `internal/export/` | SIEM JSONL, Slack webhooks, Jira issues | | **API** | `internal/api/` | REST + embedded console at `/`; shared API secret on `/v1` except health | @@ -77,7 +77,7 @@ Making the skeleton true, in order: - Self-hosted operability (console). Read routes honor the API secret, and the Helm chart schedules collectors when credentials are set. - CVE enrichment tied to workload inventory. `om enrich cve` writes a finding only when a workload package or image matches. - Crown-jewel mark on datastores. A tag or label named `sensitivity` or `data-class` is stored on the datastore, and `internet-to-sensitive-datastore` keeps only those paths. -- Attack path as the finding +- Attack path as the finding. A rules run writes one `attack_path` finding per existing finding on an internet-reachable workload that can reach a datastore, and stores the path as ordered node ids. - Cloud audit logs as graph context Further community rule packs (PCI and additional CIS mappings) stay open for contributors. diff --git a/docs/ROADMAP.md b/docs/ROADMAP.md index 3d6a7c9..91dec86 100644 --- a/docs/ROADMAP.md +++ b/docs/ROADMAP.md @@ -39,7 +39,7 @@ High-level plan for OpenSourceOM core. Timelines are approximate and community-d - [x] Sample environment (`om scan demo`) - [ ] Broader community rule packs (PCI and additional CIS mappings) — [#10](https://github.com/OpenSourceOM/core/issues/10) -Phases 0–2 shipped the walking skeleton. Exposure and identity edges now follow the cloud and Kubernetes. CVE findings follow package and image inventory on the workload. Datastores carry a sensitivity mark when a tag or label names one. Current work is the attack path as the thing an operator fixes. +Phases 0–2 shipped the walking skeleton. Exposure and identity edges now follow the cloud and Kubernetes. CVE findings follow package and image inventory on the workload. Datastores carry a sensitivity mark when a tag or label names one. A rules run writes an attack-path finding for each finding already on an internet-reachable workload that can reach a datastore. Current work is cloud audit logs as graph context. Correctness and operability come first: @@ -49,8 +49,8 @@ The path, in order: - [x] **CVE enrichment tied to workload inventory** — [#12](https://github.com/OpenSourceOM/core/issues/12) - [x] **Crown-jewel mark on datastores** — [#58](https://github.com/OpenSourceOM/core/issues/58) -- **Attack path as the finding** — [#57](https://github.com/OpenSourceOM/core/issues/57) -- **Cloud audit logs as graph context** — [#33](https://github.com/OpenSourceOM/core/issues/33), after those edges are trustworthy +- [x] **Attack path as the finding** — [#57](https://github.com/OpenSourceOM/core/issues/57) +- **Cloud audit logs as graph context** — [#33](https://github.com/OpenSourceOM/core/issues/33) [#10](https://github.com/OpenSourceOM/core/issues/10) stays open for a contributor who wants another rule pack. The priority above is graph accuracy. diff --git a/docs/adr/001-graph-schema-v0.md b/docs/adr/001-graph-schema-v0.md index d65fcc8..1caccfd 100644 --- a/docs/adr/001-graph-schema-v0.md +++ b/docs/adr/001-graph-schema-v0.md @@ -67,6 +67,7 @@ Synthetic nodes use stable global IDs (e.g. `internet:global`). - **Limitations:** Recursive path queries are SQL-based and capped (depth 6); not suitable for very large graphs without indexing and query optimization in later phases. - **Heuristics:** AWS `admin_access` follows attached and inline policies: an allow of `*` on `*`, or the AWS-managed `AdministratorAccess` policy. Deny statements, conditions, permission boundaries, and group policies are not evaluated. Azure `admin_access` is Owner, Contributor, User Access Administrator, or a custom role whose actions are `*` or include `Microsoft.Authorization/roleAssignments/write`. GCP `admin_access` is `roles/owner`, `roles/editor`, `roles/resourcemanager.projectIamAdmin`, or a custom role that includes `resourcemanager.projects.setIamPolicy`. Azure `CAN_ACCESS` for storage follows role-assignment scope. GCP `CAN_ACCESS` for buckets follows project and bucket IAM. Conditional assignments and bindings are not treated as access. Those identity edges do not apply to managed databases. Azure and GCP `REACHABLE` requires an internet path and an allow. A path is a public IP, or a private address behind a load balancer with a public frontend. GCP allow is an ingress firewall rule from `0.0.0.0/0` or `::/0` that selects the instance and is not covered by a higher-priority deny; no matching allow stays closed. Azure allow is an inbound NSG rule from `Internet`, `*`, or `0.0.0.0/0` that is not covered by a higher-priority deny. When both a NIC and a subnet NSG are attached, both must allow. A VM with a public IP and no NSG stays reachable, which is Azure's platform default. Network endpoint groups are not expanded. Kubernetes `REACHABLE` follows a LoadBalancer service or an Ingress backend that selects the pod. NodePort alone does not. A NetworkPolicy that selects the pod and governs ingress removes that edge unless a rule allows every source or `0.0.0.0/0` / `::/0`. Gateway API and internal-only ingress classes are not modeled. - **Crown jewels:** A datastore may carry `sensitivity`. Collectors copy it from a resource tag or label named `sensitivity` or `data-class` when the provider returns one. `sensitivity` wins when both are present. A blank value is omitted. Plugins may set the property on the node. Object contents are not read. `internet-to-datastore` stays unfiltered. `internet-to-sensitive-datastore` is that walk restricted to datastores whose `sensitivity` is a non-empty string. +- **Attack-path findings:** A rules run writes one `Finding` per existing finding on an internet-reachable workload paired with a datastore that workload can reach. Reach is a `CAN_ACCESS` edge from the workload, or `ASSUMES` to an identity that has `CAN_ACCESS`. The finding's `finding_type` is `attack_path`. Its `path` property is ordered node ids: the shortest `REACHABLE` walk from the internet to the workload, the identity when the hop uses one, then the datastore. `GET /v1/findings` returns that list on the row. A reachable workload with no datastore hop does not get this finding. A hop whose workload has no other finding does not either. `sensitivity` does not filter these rows. Per-control CSPM findings still run. - **Managed databases:** RDS (including Aurora instances), Azure SQL servers, and Cloud SQL instances are `Datastore` nodes. `public_access` is true only when the endpoint is open to the internet: RDS is publicly accessible and a security group allows the instance port from `0.0.0.0/0` or `::/0`; an Azure SQL server has public network access and a firewall rule from `0.0.0.0` to `255.255.255.255`; Cloud SQL has a public IPv4 address and an authorized network of `0.0.0.0/0` or `::/0`. Workload `CAN_ACCESS` follows that network path. RDS requires the same VPC and a security group that references the workload or contains its private IPv4 address on the instance port. Azure SQL requires a virtual-network rule for the VM's subnet, or a firewall range that contains the VM's public IP while the public endpoint is enabled. Cloud SQL requires the instance's private-IP network, or an authorized network that contains the instance's public IP. Same-account membership is not enough. VPC peering, prefix lists, Azure private endpoints, default outbound SNAT, and Cloud SQL Private Service Connect are not expanded. ## References diff --git a/internal/api/web/app.js b/internal/api/web/app.js index 57db3cb..c7a0f80 100644 --- a/internal/api/web/app.js +++ b/internal/api/web/app.js @@ -103,11 +103,15 @@ function renderFindings(findings, partial) { const props = f.properties || {}; const severity = props.severity || "info"; const score = props.normalized_score ?? ""; + const pathLine = Array.isArray(item.path) && item.path.length + ? `
${item.path.map((id) => escapeHTML(id)).join(" → ")}
` + : ""; return `
${severity}${score === "" ? "" : " " + score}
${f.name}
${props.description || props.title || ""}
+ ${pathLine}
${item.affected_resource_name || "Unknown resource"}
`; diff --git a/internal/export/findings.go b/internal/export/findings.go index f5b0fdf..fa61df9 100644 --- a/internal/export/findings.go +++ b/internal/export/findings.go @@ -22,6 +22,7 @@ type FindingRecord struct { AffectedID string `json:"affected_resource_id,omitempty"` AffectedName string `json:"affected_resource_name,omitempty"` AffectedType string `json:"affected_resource_type,omitempty"` + Path []string `json:"path,omitempty"` } func LoadFindingRecords(ctx context.Context, store *graph.Store) ([]FindingRecord, error) { @@ -34,6 +35,7 @@ func LoadFindingRecords(ctx context.Context, store *graph.Store) ([]FindingRecor AffectedID: view.AffectedResourceID, AffectedName: view.AffectedResourceName, AffectedType: view.AffectedResourceType, + Path: view.Path, }) return nil }) diff --git a/internal/graph/attackpath.go b/internal/graph/attackpath.go new file mode 100644 index 0000000..39e884b --- /dev/null +++ b/internal/graph/attackpath.go @@ -0,0 +1,160 @@ +// Copyright 2026 OpenSourceOM +// SPDX-License-Identifier: Apache-2.0 + +package graph + +import "context" + +// FindingTypeAttackPath is the finding_type of a row that records an +// internet-to-datastore path. Those rows are not themselves sources for +// another attack-path finding. +const FindingTypeAttackPath = "attack_path" + +// AttackPathCombination is one existing finding on an internet-reachable +// workload, paired with a datastore that workload can reach. +type AttackPathCombination struct { + Workload Node + Datastore Node + Finding Node + Path []string +} + +// ListAttackPathCombinations returns one row per finding on an +// internet-reachable workload and a datastore that workload can reach. +// +// Reach is a CAN_ACCESS edge from the workload, or ASSUMES to an identity +// that has CAN_ACCESS. Path is the shortest REACHABLE walk from the internet +// to the workload, then that access hop. When both hops reach the same +// datastore, the shorter path is kept. Findings whose finding_type is +// attack_path are not sources. Sensitivity is not a filter. +func (s *Store) ListAttackPathCombinations(ctx context.Context) ([]AttackPathCombination, error) { + rows, err := s.pool.Query(ctx, ` + WITH RECURSIVE reach AS ( + SELECT + e.target_id AS node_id, + ARRAY[e.source_id, e.target_id] AS node_ids, + 1 AS depth + FROM edges e + WHERE e.source_id = $1 AND e.type = $2 + + UNION ALL + + SELECT + e.target_id, + r.node_ids || e.target_id, + r.depth + 1 + FROM edges e + INNER JOIN reach r ON e.source_id = r.node_id + WHERE e.type = $2 + AND r.depth < $3 + AND NOT e.target_id = ANY (r.node_ids) + ), + workload_path AS ( + SELECT DISTINCT ON (r.node_id) + r.node_id, + r.node_ids + FROM reach r + INNER JOIN nodes w ON w.id = r.node_id AND w.type = $4 + ORDER BY r.node_id, r.depth + ), + hops AS ( + SELECT + wp.node_id AS workload_id, + d.id AS datastore_id, + f.id AS finding_id, + wp.node_ids || d.id AS path + FROM workload_path wp + INNER JOIN edges v ON v.target_id = wp.node_id AND v.type = $5 + INNER JOIN nodes f ON f.id = v.source_id AND f.type = $6 + AND COALESCE(f.properties->>'finding_type', '') <> $11 + INNER JOIN edges a ON a.source_id = wp.node_id AND a.type = $7 + INNER JOIN nodes d ON d.id = a.target_id AND d.type = $8 + WHERE NOT (d.id = ANY (wp.node_ids)) + + UNION ALL + + SELECT + wp.node_id, + d.id, + f.id, + wp.node_ids || i.id || d.id + FROM workload_path wp + INNER JOIN edges v ON v.target_id = wp.node_id AND v.type = $5 + INNER JOIN nodes f ON f.id = v.source_id AND f.type = $6 + AND COALESCE(f.properties->>'finding_type', '') <> $11 + INNER JOIN edges assumes ON assumes.source_id = wp.node_id AND assumes.type = $9 + INNER JOIN nodes i ON i.id = assumes.target_id AND i.type = $10 + INNER JOIN edges a ON a.source_id = i.id AND a.type = $7 + INNER JOIN nodes d ON d.id = a.target_id AND d.type = $8 + WHERE NOT (i.id = ANY (wp.node_ids)) + AND NOT (d.id = ANY (wp.node_ids || i.id)) + ) + SELECT DISTINCT ON (finding_id, datastore_id) + workload_id, datastore_id, finding_id, path + FROM hops + ORDER BY finding_id, datastore_id, cardinality(path), path + `, InternetNodeID, EdgeReachable, internetDatastoreDepth, NodeWorkload, + EdgeViolates, NodeFinding, EdgeCanAccess, NodeDatastore, EdgeAssumes, NodeIdentity, + FindingTypeAttackPath) + if err != nil { + return nil, err + } + defer rows.Close() + + var combos []AttackPathCombination + for rows.Next() { + var workloadID, datastoreID, findingID string + var path []string + if err := rows.Scan(&workloadID, &datastoreID, &findingID, &path); err != nil { + return nil, err + } + workload, err := s.GetNode(ctx, workloadID) + if err != nil { + return nil, err + } + datastore, err := s.GetNode(ctx, datastoreID) + if err != nil { + return nil, err + } + finding, err := s.GetNode(ctx, findingID) + if err != nil { + return nil, err + } + combos = append(combos, AttackPathCombination{ + Workload: workload, + Datastore: datastore, + Finding: finding, + Path: path, + }) + } + return combos, rows.Err() +} + +func pathIDs(props map[string]any) []string { + raw, ok := props["path"] + if !ok || raw == nil { + return nil + } + switch ids := raw.(type) { + case []string: + if len(ids) == 0 { + return nil + } + return ids + case []any: + out := make([]string, 0, len(ids)) + for _, id := range ids { + s, ok := id.(string) + if !ok { + return nil + } + out = append(out, s) + } + if len(out) == 0 { + return nil + } + return out + default: + return nil + } +} diff --git a/internal/graph/attackpath_test.go b/internal/graph/attackpath_test.go new file mode 100644 index 0000000..4f94501 --- /dev/null +++ b/internal/graph/attackpath_test.go @@ -0,0 +1,205 @@ +// Copyright 2026 OpenSourceOM +// SPDX-License-Identifier: Apache-2.0 + +package graph_test + +import ( + "context" + "strings" + "testing" + + "github.com/OpenSourceOM/core/internal/graph" +) + +func TestAttackPathCombinationUsesShortestHop(t *testing.T) { + ctx := context.Background() + store := openQueryFixture(t) + + const ( + sg = "aws:111122223333:us-east-1:network:sg-web" + web = "aws:111122223333:us-east-1:workload:web" + role = "aws:111122223333:us-east-1:identity:admin" + db = "aws:111122223333:us-east-1:datastore:prod" + cve = "finding:cve-2021-44228:" + web + ) + edge := func(src, dst, typ string) graph.Edge { + return graph.Edge{ID: src + "|" + dst + "|" + typ, SourceID: src, TargetID: dst, Type: typ} + } + batch := graph.Batch{ + Nodes: []graph.Node{ + {ID: graph.InternetNodeID, Type: graph.NodeInternet, Name: "Internet"}, + {ID: sg, Type: graph.NodeNetwork, Name: "sg-web", Provider: "aws"}, + {ID: web, Type: graph.NodeWorkload, Name: "web-1", Provider: "aws"}, + {ID: role, Type: graph.NodeIdentity, Name: "AdminRole", Provider: "aws"}, + {ID: db, Type: graph.NodeDatastore, Name: "prod-db", Provider: "aws"}, + { + ID: cve, Type: graph.NodeFinding, Name: "CVE-2021-44228", Provider: "aws", + Properties: map[string]any{ + "finding_type": "cve", + "title": "CVE-2021-44228", + "normalized_score": 100, + }, + }, + }, + Edges: []graph.Edge{ + edge(graph.InternetNodeID, sg, graph.EdgeReachable), + edge(sg, web, graph.EdgeReachable), + edge(web, role, graph.EdgeAssumes), + edge(role, db, graph.EdgeCanAccess), + edge(web, db, graph.EdgeCanAccess), + edge(cve, web, graph.EdgeViolates), + }, + } + if err := store.UpsertBatch(ctx, batch); err != nil { + t.Fatalf("seed: %v", err) + } + + combos, err := store.ListAttackPathCombinations(ctx) + if err != nil { + t.Fatalf("list: %v", err) + } + if len(combos) != 1 { + t.Fatalf("combinations = %d, want the direct hop", len(combos)) + } + want := []string{graph.InternetNodeID, sg, web, db} + if strings.Join(combos[0].Path, ",") != strings.Join(want, ",") { + t.Fatalf("path = %v, want %v", combos[0].Path, want) + } + if combos[0].Finding.ID != cve || combos[0].Datastore.ID != db { + t.Fatalf("combo = finding %s datastore %s", combos[0].Finding.ID, combos[0].Datastore.ID) + } +} + +func TestAttackPathCombinationViaAssumedIdentity(t *testing.T) { + ctx := context.Background() + store := openQueryFixture(t) + if err := store.UpsertBatch(ctx, attackPathIdentityBatch()); err != nil { + t.Fatalf("seed: %v", err) + } + + combos, err := store.ListAttackPathCombinations(ctx) + if err != nil { + t.Fatalf("list: %v", err) + } + if len(combos) != 1 { + t.Fatalf("combinations = %d, want 1", len(combos)) + } + want := []string{graph.InternetNodeID, attackPathSG, attackPathWeb, attackPathRole, attackPathDB} + if strings.Join(combos[0].Path, ",") != strings.Join(want, ",") { + t.Fatalf("path = %v, want %v", combos[0].Path, want) + } +} + +func TestAttackPathCombinationRequiresFindingOnReachableWorkload(t *testing.T) { + ctx := context.Background() + store := openQueryFixture(t) + + batch := attackPathIdentityBatch() + batch.Nodes = append(batch.Nodes, graph.Node{ + ID: attackPathPrivate, Type: graph.NodeWorkload, Name: "worker-1", Provider: "aws", + }, graph.Node{ + ID: attackPathAppRole, Type: graph.NodeIdentity, Name: "AppRole", Provider: "aws", + }, graph.Node{ + ID: attackPathPrivateDB, Type: graph.NodeDatastore, Name: "assets", Provider: "aws", + }, graph.Node{ + ID: attackPathPrivateFinding, Type: graph.NodeFinding, Name: "CVE-2021-44228", Provider: "aws", + Properties: map[string]any{"finding_type": "cve", "normalized_score": 90}, + }, graph.Node{ + ID: attackPathNetworkFinding, Type: graph.NodeFinding, Name: "Open security group", Provider: "aws", + Properties: map[string]any{"finding_type": "cspm", "normalized_score": 70}, + }) + edge := func(src, dst, typ string) graph.Edge { + return graph.Edge{ID: src + "|" + dst + "|" + typ, SourceID: src, TargetID: dst, Type: typ} + } + batch.Edges = append(batch.Edges, + edge(attackPathPrivate, attackPathAppRole, graph.EdgeAssumes), + edge(attackPathAppRole, attackPathPrivateDB, graph.EdgeCanAccess), + edge(attackPathPrivateFinding, attackPathPrivate, graph.EdgeViolates), + edge(attackPathNetworkFinding, attackPathSG, graph.EdgeViolates), + ) + if err := store.UpsertBatch(ctx, batch); err != nil { + t.Fatalf("seed: %v", err) + } + + combos, err := store.ListAttackPathCombinations(ctx) + if err != nil { + t.Fatalf("list: %v", err) + } + if len(combos) != 1 || combos[0].Workload.ID != attackPathWeb { + t.Fatalf("combinations = %+v, want only the reachable workload", combos) + } + + if _, err := store.Pool().Exec(ctx, `DELETE FROM nodes WHERE id = $1`, attackPathCVE); err != nil { + t.Fatal(err) + } + combos, err = store.ListAttackPathCombinations(ctx) + if err != nil { + t.Fatalf("list without finding: %v", err) + } + if len(combos) != 0 { + t.Fatalf("combinations = %d, want none without a finding on the workload", len(combos)) + } +} + +func TestAttackPathCombinationIgnoresAttackPathFindings(t *testing.T) { + ctx := context.Background() + store := openQueryFixture(t) + batch := attackPathIdentityBatch() + batch.Nodes[len(batch.Nodes)-1].Properties = map[string]any{ + "finding_type": graph.FindingTypeAttackPath, + "rule_id": "attack-path", + } + if err := store.UpsertBatch(ctx, batch); err != nil { + t.Fatalf("seed: %v", err) + } + combos, err := store.ListAttackPathCombinations(ctx) + if err != nil { + t.Fatalf("list: %v", err) + } + if len(combos) != 0 { + t.Fatalf("combinations = %d, want none from an attack-path finding", len(combos)) + } +} + +const ( + attackPathSG = "aws:111122223333:us-east-1:network:sg-web" + attackPathWeb = "aws:111122223333:us-east-1:workload:web" + attackPathRole = "aws:111122223333:us-east-1:identity:admin" + attackPathDB = "aws:111122223333:us-east-1:datastore:prod" + attackPathCVE = "finding:cve-2021-44228:" + attackPathWeb + attackPathPrivate = "aws:111122223333:us-east-1:workload:worker" + attackPathAppRole = "aws:111122223333:us-east-1:identity:app" + attackPathPrivateDB = "aws:111122223333:us-east-1:datastore:assets" + attackPathPrivateFinding = "finding:cve-2021-44228:" + attackPathPrivate + attackPathNetworkFinding = "finding:cspm-sg:" + attackPathSG +) + +func attackPathIdentityBatch() graph.Batch { + edge := func(src, dst, typ string) graph.Edge { + return graph.Edge{ID: src + "|" + dst + "|" + typ, SourceID: src, TargetID: dst, Type: typ} + } + return graph.Batch{ + Nodes: []graph.Node{ + {ID: graph.InternetNodeID, Type: graph.NodeInternet, Name: "Internet"}, + {ID: attackPathSG, Type: graph.NodeNetwork, Name: "sg-web", Provider: "aws"}, + {ID: attackPathWeb, Type: graph.NodeWorkload, Name: "web-1", Provider: "aws"}, + {ID: attackPathRole, Type: graph.NodeIdentity, Name: "AdminRole", Provider: "aws"}, + {ID: attackPathDB, Type: graph.NodeDatastore, Name: "prod-db", Provider: "aws"}, + { + ID: attackPathCVE, Type: graph.NodeFinding, Name: "CVE-2021-44228", Provider: "aws", + Properties: map[string]any{ + "finding_type": "cve", + "title": "CVE-2021-44228", + "normalized_score": 100, + }, + }, + }, + Edges: []graph.Edge{ + edge(graph.InternetNodeID, attackPathSG, graph.EdgeReachable), + edge(attackPathSG, attackPathWeb, graph.EdgeReachable), + edge(attackPathWeb, attackPathRole, graph.EdgeAssumes), + edge(attackPathRole, attackPathDB, graph.EdgeCanAccess), + edge(attackPathCVE, attackPathWeb, graph.EdgeViolates), + }, + } +} diff --git a/internal/graph/page.go b/internal/graph/page.go index 167134b..d9b1fe4 100644 --- a/internal/graph/page.go +++ b/internal/graph/page.go @@ -274,6 +274,7 @@ func (s *Store) ListFindings(ctx context.Context, limit int, cursor string) (Fin if targetType != nil { view.AffectedResourceType = *targetType } + view.Path = pathIDs(view.Finding.Properties) findings = append(findings, view) scores = append(scores, score) } diff --git a/internal/graph/pathids_test.go b/internal/graph/pathids_test.go new file mode 100644 index 0000000..1da29d6 --- /dev/null +++ b/internal/graph/pathids_test.go @@ -0,0 +1,19 @@ +// Copyright 2026 OpenSourceOM +// SPDX-License-Identifier: Apache-2.0 + +package graph + +import "testing" + +func TestPathIDs(t *testing.T) { + got := pathIDs(map[string]any{"path": []any{"internet:global", "web"}}) + if len(got) != 2 || got[0] != "internet:global" || got[1] != "web" { + t.Fatalf("pathIDs = %#v", got) + } + if pathIDs(map[string]any{"path": []any{"web", 1}}) != nil { + t.Fatal("mixed path should be dropped") + } + if pathIDs(map[string]any{}) != nil { + t.Fatal("missing path should be empty") + } +} diff --git a/internal/graph/store.go b/internal/graph/store.go index f0c402b..2af2851 100644 --- a/internal/graph/store.go +++ b/internal/graph/store.go @@ -64,7 +64,7 @@ type Scope struct { // ReplaceInventory upserts batch, then deletes inventory in scopes that the // batch no longer contains. Finding nodes stay unless their affected resource -// was removed. Rule evaluation removes CSPM findings that no longer match. +// was removed. Rule evaluation removes findings that no longer match. func (s *Store) ReplaceInventory(ctx context.Context, scopes []Scope, batch Batch) error { tx, err := s.pool.Begin(ctx) if err != nil { @@ -262,8 +262,9 @@ func deleteAbsentInventory(ctx context.Context, tx pgx.Tx, scopes []Scope, batch return nil } -// DeleteStaleRuleFindings removes CSPM findings for ruleID whose ids are not in -// keepIDs. Findings from other rules and from CVE enrichment are left alone. +// DeleteStaleRuleFindings removes CSPM and attack-path findings for ruleID +// whose ids are not in keepIDs. Findings from other rules and from CVE +// enrichment are left alone. func (s *Store) DeleteStaleRuleFindings(ctx context.Context, ruleID string, keepIDs []string) error { if ruleID == "" { return nil @@ -274,10 +275,10 @@ func (s *Store) DeleteStaleRuleFindings(ctx context.Context, ruleID string, keep _, err := s.pool.Exec(ctx, ` DELETE FROM nodes WHERE type = $1 - AND properties->>'finding_type' = 'cspm' + AND properties->>'finding_type' IN ('cspm', $4) AND properties->>'rule_id' = $2 AND NOT (id = ANY($3::text[])) - `, NodeFinding, ruleID, keepIDs) + `, NodeFinding, ruleID, keepIDs, FindingTypeAttackPath) return err } diff --git a/internal/graph/types.go b/internal/graph/types.go index 0de3303..8abf201 100644 --- a/internal/graph/types.go +++ b/internal/graph/types.go @@ -76,6 +76,8 @@ type FindingView struct { AffectedResourceID string `json:"affected_resource_id,omitempty"` AffectedResourceName string `json:"affected_resource_name,omitempty"` AffectedResourceType string `json:"affected_resource_type,omitempty"` + // Path is the ordered node ids recorded on an attack-path finding. + Path []string `json:"path,omitempty"` } type GraphSnapshot struct { diff --git a/internal/rules/attackpath.go b/internal/rules/attackpath.go new file mode 100644 index 0000000..2ea03ad --- /dev/null +++ b/internal/rules/attackpath.go @@ -0,0 +1,145 @@ +// Copyright 2026 OpenSourceOM +// SPDX-License-Identifier: Apache-2.0 + +package rules + +import ( + "context" + "encoding/json" + "fmt" + "strings" + + "github.com/OpenSourceOM/core/internal/graph" +) + +const ( + attackPathRuleID = "attack-path" + attackPathDefaultScore = 75 +) + +// attackPathRule runs last. RunAll persists earlier rules first, so this +// pass sees findings those rules just wrote as well as CVE and exposure +// findings already on the graph. +var attackPathRule = Rule{ + ID: attackPathRuleID, + Name: "Attack path", + Description: "Internet-reachable workload with a finding on a path to a datastore", + BaseScore: attackPathDefaultScore, + Run: ruleAttackPath, +} + +func ruleAttackPath(ctx context.Context, store *graph.Store) ([]Match, error) { + combos, err := store.ListAttackPathCombinations(ctx) + if err != nil { + return nil, err + } + matches := make([]Match, 0, len(combos)) + for _, combo := range combos { + score, ok := findingScore(combo.Finding.Properties) + if !ok { + score = attackPathDefaultScore + } + source := strings.TrimSpace(combo.Finding.Name) + if title, ok := combo.Finding.Properties["title"].(string); ok && strings.TrimSpace(title) != "" { + source = strings.TrimSpace(title) + } + matches = append(matches, Match{ + RuleID: attackPathRuleID, + Resource: combo.Workload, + Title: fmt.Sprintf("Attack path: %s to %s", combo.Workload.Name, combo.Datastore.Name), + Description: fmt.Sprintf("%s is on a path to %s.", source, combo.Datastore.Name), + BaseScore: score, + Path: combo.Path, + SourceFindingID: combo.Finding.ID, + DatastoreID: combo.Datastore.ID, + }) + } + return matches, nil +} + +func (e *Engine) persistAttackPathMatches(ctx context.Context, rule Rule, matches []Match) (int, error) { + keep := make([]string, 0, len(matches)) + for _, match := range matches { + if match.SourceFindingID == "" || match.DatastoreID == "" { + return 0, fmt.Errorf("attack path match missing source finding or datastore") + } + findingID := attackPathFindingID(match.SourceFindingID, match.DatastoreID) + if err := e.persistAttackPathFinding(ctx, rule, match, findingID); err != nil { + return 0, err + } + keep = append(keep, findingID) + } + if err := e.store.DeleteStaleRuleFindings(ctx, rule.ID, keep); err != nil { + return 0, err + } + return len(matches), nil +} + +func attackPathFindingID(sourceFindingID, datastoreID string) string { + return fmt.Sprintf("finding:%s:%s:%s", attackPathRuleID, sourceFindingID, datastoreID) +} + +func (e *Engine) persistAttackPathFinding(ctx context.Context, rule Rule, match Match, findingID string) error { + score := match.BaseScore + batch := graph.Batch{ + Nodes: []graph.Node{ + { + ID: findingID, + Type: graph.NodeFinding, + Name: rule.Name, + Provider: match.Resource.Provider, + Region: match.Resource.Region, + AccountID: match.Resource.AccountID, + Properties: graph.MustProperties(map[string]any{ + "finding_type": graph.FindingTypeAttackPath, + "rule_id": rule.ID, + "title": match.Title, + "description": match.Description, + "severity": SeverityFromScore(score), + "normalized_score": score, + "affected_resource": match.Resource.ID, + "source_finding_id": match.SourceFindingID, + "datastore_id": match.DatastoreID, + "path": match.Path, + }), + }, + }, + Edges: []graph.Edge{ + { + ID: fmt.Sprintf("%s|%s|%s", findingID, match.Resource.ID, graph.EdgeViolates), + SourceID: findingID, + TargetID: match.Resource.ID, + Type: graph.EdgeViolates, + }, + }, + } + return e.store.UpsertBatch(ctx, batch) +} + +func findingScore(props map[string]any) (int, bool) { + if props == nil { + return 0, false + } + switch v := props["normalized_score"].(type) { + case float64: + return int(v), true + case float32: + return int(v), true + case int: + return v, true + case int64: + return int(v), true + case json.Number: + n, err := v.Int64() + if err != nil { + f, ferr := v.Float64() + if ferr != nil { + return 0, false + } + return int(f), true + } + return int(n), true + default: + return 0, false + } +} diff --git a/internal/rules/attackpath_test.go b/internal/rules/attackpath_test.go new file mode 100644 index 0000000..6ee0045 --- /dev/null +++ b/internal/rules/attackpath_test.go @@ -0,0 +1,234 @@ +// Copyright 2026 OpenSourceOM +// SPDX-License-Identifier: Apache-2.0 + +package rules_test + +import ( + "context" + "strings" + "testing" + + "github.com/OpenSourceOM/core/internal/collectors/demo" + "github.com/OpenSourceOM/core/internal/graph" + "github.com/OpenSourceOM/core/internal/rules" +) + +func TestAttackPathFindingRecordsPathAndSkipsIncompleteGraphs(t *testing.T) { + ctx := context.Background() + store := openRulesFixture(t) + + const ( + web = "aws:111122223333:us-east-1:workload:web" + db = "aws:111122223333:us-east-1:datastore:prod" + sg = "aws:111122223333:us-east-1:network:sg-web" + cve = "finding:cve-2021-44228:" + web + ) + if err := store.UpsertBatch(ctx, attackPathBatch(web, db, sg, cve)); err != nil { + t.Fatalf("seed: %v", err) + } + + none, err := rules.NewEngine(store).Run(ctx, "attack-path") + if err != nil { + t.Fatalf("run without hop: %v", err) + } + if none.FindingsCreated != 0 { + t.Fatalf("FindingsCreated = %d, want 0 before the datastore hop exists", none.FindingsCreated) + } + + hop := graph.Edge{ + ID: web + "|" + db + "|" + graph.EdgeCanAccess, SourceID: web, TargetID: db, Type: graph.EdgeCanAccess, + } + if err := store.UpsertBatch(ctx, graph.Batch{Edges: []graph.Edge{hop}}); err != nil { + t.Fatalf("add hop: %v", err) + } + + result, err := rules.NewEngine(store).Run(ctx, "attack-path") + if err != nil { + t.Fatalf("run: %v", err) + } + if result.FindingsCreated != 1 { + t.Fatalf("FindingsCreated = %d, want 1", result.FindingsCreated) + } + if len(result.Matches) != 1 || len(result.Matches[0].Path) == 0 { + t.Fatalf("match path = %#v", result.Matches) + } + + findingID := "finding:attack-path:" + cve + ":" + db + node, err := store.GetNode(ctx, findingID) + if err != nil { + t.Fatalf("get finding: %v", err) + } + if node.Properties["finding_type"] != graph.FindingTypeAttackPath { + t.Fatalf("finding_type = %#v", node.Properties["finding_type"]) + } + if node.Properties["source_finding_id"] != cve || node.Properties["datastore_id"] != db { + t.Fatalf("properties = %#v", node.Properties) + } + wantPath := []string{graph.InternetNodeID, sg, web, db} + page, err := store.ListFindings(ctx, 20, "") + if err != nil { + t.Fatalf("list: %v", err) + } + var viewed graph.FindingView + for _, view := range page.Findings { + if view.Finding.ID == findingID { + viewed = view + break + } + } + if viewed.Finding.ID == "" { + t.Fatal("findings list omitted the attack-path row") + } + if strings.Join(viewed.Path, ",") != strings.Join(wantPath, ",") { + t.Fatalf("list path = %v, want %v", viewed.Path, wantPath) + } + if viewed.AffectedResourceID != web { + t.Fatalf("affected resource = %s, want the workload", viewed.AffectedResourceID) + } + if viewed.Finding.Properties["severity"] != "critical" { + t.Fatalf("severity = %#v, want the source score band", viewed.Finding.Properties["severity"]) + } + + if _, err := store.Pool().Exec(ctx, `DELETE FROM edges WHERE id = $1`, hop.ID); err != nil { + t.Fatal(err) + } + again, err := rules.NewEngine(store).Run(ctx, "attack-path") + if err != nil { + t.Fatalf("rerun: %v", err) + } + if again.FindingsCreated != 0 { + t.Fatalf("FindingsCreated = %d, want 0 after the hop is removed", again.FindingsCreated) + } + requireRuleNode(t, ctx, store, findingID, false) + requireRuleNode(t, ctx, store, cve, true) +} + +func TestRunAllWritesAttackPathWithoutDroppingControlFindings(t *testing.T) { + ctx := context.Background() + store := openRulesFixture(t) + + const ( + web = "aws:111122223333:us-east-1:workload:web" + db = "aws:111122223333:us-east-1:datastore:prod" + public = "aws:111122223333:global:datastore:public-logs" + sg = "aws:111122223333:us-east-1:network:sg-web" + ) + batch := attackPathBatch(web, db, sg, "") + batch.Nodes = append(batch.Nodes, graph.Node{ + ID: public, Type: graph.NodeDatastore, Name: "public-logs", Provider: "aws", AccountID: "111122223333", + Properties: map[string]any{"public_access": true}, + }) + batch.Edges = append(batch.Edges, graph.Edge{ + ID: web + "|" + db + "|" + graph.EdgeCanAccess, SourceID: web, TargetID: db, Type: graph.EdgeCanAccess, + }) + if err := store.UpsertBatch(ctx, batch); err != nil { + t.Fatalf("seed: %v", err) + } + + result, err := rules.NewEngine(store).RunAll(ctx) + if err != nil { + t.Fatalf("run all: %v", err) + } + if result.FindingsCreated < 2 { + t.Fatalf("FindingsCreated = %d, want the control finding and the attack path", result.FindingsCreated) + } + requireRuleNode(t, ctx, store, "finding:cspm-public-datastore:"+public, true) + + internetFinding := "finding:cspm-internet-workload:" + web + requireRuleNode(t, ctx, store, internetFinding, true) + attackID := "finding:attack-path:" + internetFinding + ":" + db + requireRuleNode(t, ctx, store, attackID, true) + + again, err := rules.NewEngine(store).RunAll(ctx) + if err != nil { + t.Fatalf("second run: %v", err) + } + if again.FindingsCreated < 2 { + t.Fatalf("second FindingsCreated = %d", again.FindingsCreated) + } + var attackPaths int + err = store.Pool().QueryRow(ctx, ` + SELECT COUNT(*) FROM nodes + WHERE type = $1 AND properties->>'finding_type' = $2 + `, graph.NodeFinding, graph.FindingTypeAttackPath).Scan(&attackPaths) + if err != nil { + t.Fatal(err) + } + if attackPaths != 1 { + t.Fatalf("attack-path findings = %d, want 1 after a second run", attackPaths) + } +} + +func TestDemoGraphWritesAttackPathThroughAdminRole(t *testing.T) { + ctx := context.Background() + store := openRulesFixture(t) + if err := store.UpsertBatch(ctx, demo.Collect()); err != nil { + t.Fatalf("seed: %v", err) + } + + if _, err := rules.NewEngine(store).RunAll(ctx); err != nil { + t.Fatalf("run: %v", err) + } + + page, err := store.ListFindings(ctx, 200, "") + if err != nil { + t.Fatalf("list: %v", err) + } + want := []string{ + graph.InternetNodeID, + demo.SecurityGroupWebID, + demo.WebInstanceID, + demo.AdminRoleID, + demo.ProdDBID, + } + var paths int + for _, view := range page.Findings { + if view.Finding.Properties["finding_type"] != graph.FindingTypeAttackPath { + continue + } + paths++ + if strings.Join(view.Path, ",") != strings.Join(want, ",") { + t.Fatalf("path = %v, want %v", view.Path, want) + } + if view.AffectedResourceID != demo.WebInstanceID { + t.Fatalf("affected = %s", view.AffectedResourceID) + } + if view.Finding.Properties["datastore_id"] != demo.ProdDBID { + t.Fatalf("datastore = %#v", view.Finding.Properties["datastore_id"]) + } + } + if paths != 3 { + t.Fatalf("attack-path findings = %d, want one per finding on web-1", paths) + } +} + +func attackPathBatch(web, db, sg, findingID string) graph.Batch { + edge := func(src, dst, typ string) graph.Edge { + return graph.Edge{ID: src + "|" + dst + "|" + typ, SourceID: src, TargetID: dst, Type: typ} + } + batch := graph.Batch{ + Nodes: []graph.Node{ + {ID: graph.InternetNodeID, Type: graph.NodeInternet, Name: "Internet"}, + {ID: sg, Type: graph.NodeNetwork, Name: "sg-web", Provider: "aws", AccountID: "111122223333"}, + {ID: web, Type: graph.NodeWorkload, Name: "web-1", Provider: "aws", Region: "us-east-1", AccountID: "111122223333"}, + {ID: db, Type: graph.NodeDatastore, Name: "prod-db", Provider: "aws", AccountID: "111122223333"}, + }, + Edges: []graph.Edge{ + edge(graph.InternetNodeID, sg, graph.EdgeReachable), + edge(sg, web, graph.EdgeReachable), + }, + } + if findingID != "" { + batch.Nodes = append(batch.Nodes, graph.Node{ + ID: findingID, Type: graph.NodeFinding, Name: "CVE-2021-44228", Provider: "aws", AccountID: "111122223333", + Properties: map[string]any{ + "finding_type": "cve", + "title": "CVE-2021-44228", + "normalized_score": 100, + "affected_resource": web, + }, + }) + batch.Edges = append(batch.Edges, edge(findingID, web, graph.EdgeViolates)) + } + return batch +} diff --git a/internal/rules/engine.go b/internal/rules/engine.go index 4d48c71..a55acdb 100644 --- a/internal/rules/engine.go +++ b/internal/rules/engine.go @@ -11,12 +11,15 @@ import ( ) type Match struct { - RuleID string `json:"rule_id"` - Resource graph.Node `json:"resource"` - Title string `json:"title"` - Description string `json:"description"` - BaseScore int `json:"base_score"` - Context GraphContext `json:"graph_context"` + RuleID string `json:"rule_id"` + Resource graph.Node `json:"resource"` + Title string `json:"title"` + Description string `json:"description"` + BaseScore int `json:"base_score"` + Context GraphContext `json:"graph_context"` + Path []string `json:"path,omitempty"` + SourceFindingID string `json:"source_finding_id,omitempty"` + DatastoreID string `json:"datastore_id,omitempty"` } type Rule struct { @@ -61,6 +64,7 @@ func init() { panic(err) } Catalog = append(Catalog, extra...) + Catalog = append(Catalog, attackPathRule) } func CatalogMap() map[string]string { @@ -123,6 +127,9 @@ func (e *Engine) Run(ctx context.Context, ruleID string) (RunResult, error) { } func (e *Engine) persistMatches(ctx context.Context, rule Rule, matches []Match) (int, error) { + if rule.ID == attackPathRuleID { + return e.persistAttackPathMatches(ctx, rule, matches) + } keep := make([]string, 0, len(matches)) for _, match := range matches { findingID := findingNodeID(rule.ID, match.Resource.ID)