diff --git a/README.md b/README.md index c3df16a..3028ac7 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. 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). +> **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. `om scan aws` attaches recent CloudTrail management events to the identity and resource on that path. Further rule packs are open for contributors. See the [roadmap](./docs/ROADMAP.md). ## Why this exists @@ -160,7 +160,7 @@ Full documentation: [opensourceom.org](https://opensourceom.org) (docs at [opens | **0** | Graph schema v0, AWS collector, ingest API, `om` CLI | | **1** | Attack path queries, CVE enrichment, web UI, Azure/GCP collectors | | **2** | CSPM rules, blast radius, K8s connector, exports | -| **3** *(now)* | Graph accuracy, crown-jewel datastores, attack-path findings, rule packs, cloud audit ingest | +| **3** *(now)* | Graph accuracy, crown-jewel datastores, attack-path findings, CloudTrail context, rule packs | Details: [docs/ROADMAP.md](./docs/ROADMAP.md) diff --git a/collectors/README.md b/collectors/README.md index 4f7f667..4d1cf4f 100644 --- a/collectors/README.md +++ b/collectors/README.md @@ -20,7 +20,7 @@ Cloud and platform ingestion plugins. Each collector normalizes provider APIs in The **demo** collector loads a fixed environment (internet-exposed web tier, production database, public/private buckets, admin vs app identities, public Kubernetes service). Use it to exercise CSPM packs without cloud credentials. -The **AWS** collector emits properties the CIS pack matches on: `open_ingress`, `imdsv2`, `public_ip`, S3 `encryption` / `versioning` / `public_access_block`, and IAM user `mfa` / `unused_access_keys`. It also records RDS and Aurora instances as datastores. S3 bucket tags and the RDS `TagList` copy `sensitivity` or `data-class` onto the datastore. `sensitivity` wins when both are set. +The **AWS** collector emits properties the CIS pack matches on: `open_ingress`, `imdsv2`, `public_ip`, S3 `encryption` / `versioning` / `public_access_block`, and IAM user `mfa` / `unused_access_keys`. It also records RDS and Aurora instances as datastores. S3 bucket tags and the RDS `TagList` copy `sensitivity` or `data-class` onto the datastore. `sensitivity` wins when both are set. The same scan reads CloudTrail management events from the last 24 hours and stores the ones that name an identity and a resource on an exposed path. A failed lookup omits those events. The **Azure** collector records logical SQL servers, and the **GCP** collector records Cloud SQL instances. A workload gets `CAN_ACCESS` only when a security group, firewall, or VPC path allows it. Azure copies the crown-jewel mark from storage-account and SQL-server tags. GCP copies it from a bucket label or a Cloud SQL user label. The keys are `sensitivity` and `data-class`. A plugin may set `sensitivity` on a datastore directly. diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 84bf31a..c9a6110 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -31,7 +31,7 @@ SPDX-License-Identifier: Apache-2.0 | Component | Location | Notes | |-----------|----------|-------| -| **Collectors** | `internal/collectors/` | AWS, Azure, GCP, Kubernetes, demo; RDS, Azure SQL, and Cloud SQL are datastores; AWS emits CIS pack properties | +| **Collectors** | `internal/collectors/` | AWS, Azure, GCP, Kubernetes, demo; RDS, Azure SQL, and Cloud SQL are datastores; AWS emits CIS pack properties and recent CloudTrail management events | | **Plugin SDK** | `sdk/collector`, `internal/plugins/` | External executables; `om scan plugin` | | **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` | @@ -78,7 +78,7 @@ Making the skeleton true, in order: - 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. 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 +- Cloud audit logs as graph context. `om scan aws` stores CloudTrail management events from the last 24 hours on the identity and resource they name, when that resource is on an exposed path. `GET /v1/graph/query` returns those events in `audits` for each path that contains the resource. Azure Activity Log, GCP Cloud Audit Logs, and S3 data events such as `GetObject` are not collected. Further community rule packs (PCI and additional CIS mappings) stay open for contributors. diff --git a/docs/ROADMAP.md b/docs/ROADMAP.md index 91dec86..1ca8167 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. 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. +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. `om scan aws` attaches recent CloudTrail management events to the identity and resource they name, and path queries return an event when that resource is on the path. Correctness and operability come first: @@ -50,9 +50,11 @@ 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) - [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) +- [x] **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. +`om scan aws` reads CloudTrail management events from the last 24 hours. An event is stored on the identity and the resource it names when that resource sits on an exposed path. Named path queries list those events when the resource is on the returned path. Azure Activity Log, GCP Cloud Audit Logs, and S3 object data events such as `GetObject` are not in this slice. + +[#10](https://github.com/OpenSourceOM/core/issues/10) stays open for a contributor who wants another rule pack. ## Open source vs. commercial diff --git a/docs/adr/001-graph-schema-v0.md b/docs/adr/001-graph-schema-v0.md index 1caccfd..afa5b19 100644 --- a/docs/adr/001-graph-schema-v0.md +++ b/docs/adr/001-graph-schema-v0.md @@ -68,6 +68,7 @@ Synthetic nodes use stable global IDs (e.g. `internet:global`). - **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. +- **Cloud audit context:** `om scan aws` calls CloudTrail `LookupEvents` for the scan region, and also `us-east-1` when the scan region is different, over the last 24 hours. A fixed list of management events is kept when the event names an identity already in the batch and a resource on an exposed path. Exposed means an internet-reachable workload, a node that workload assumes or can access, a network that workload affects, a public datastore, or an identity that can access a public datastore. The event is stored as `audit_events` on both nodes, at most five per node, newest first. A failed lookup omits the property. `GET /v1/graph/query` and `om paths run` include an `audits` entry when the event's resource node is on a returned path. S3 object data events such as `GetObject` are not returned by `LookupEvents`. Azure and GCP audit logs are not collected. These events are cloud-provider evidence, not a log of who used OpenSourceOM. - **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/go.mod b/go.mod index f9fba13..968031d 100644 --- a/go.mod +++ b/go.mod @@ -14,6 +14,7 @@ require ( github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/storage/armstorage/v3 v3.0.0 github.com/aws/aws-sdk-go-v2 v1.47.1 github.com/aws/aws-sdk-go-v2/config v1.33.6 + github.com/aws/aws-sdk-go-v2/service/cloudtrail v1.53.0 github.com/aws/aws-sdk-go-v2/service/ec2 v1.336.1 github.com/aws/aws-sdk-go-v2/service/iam v1.64.1 github.com/aws/aws-sdk-go-v2/service/rds v1.129.0 diff --git a/go.sum b/go.sum index dc8d320..355729f 100644 --- a/go.sum +++ b/go.sum @@ -72,6 +72,8 @@ github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.4 h1:dD4MR81I7YkpEBRk6UP github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.4/go.mod h1:EcXV1kAFd5XwSkDHlj94gnF3q5CkJyYiIJfH8N0VmrE= github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.4 h1:7Wo47d/xn/7KttCSBd8EGYeZ7ULRFRkUHr6vkZPBzVQ= github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.4/go.mod h1:tDB2IVC1xC3vX8o+6uRlzhTxP3g1b77CZXFX/oD2FnQ= +github.com/aws/aws-sdk-go-v2/service/cloudtrail v1.53.0 h1:WMgigsEPtSgsVe+jBMqCuAF2u0j/CnSjCm3I6Ar7nFo= +github.com/aws/aws-sdk-go-v2/service/cloudtrail v1.53.0/go.mod h1:1QQJFpFapuZD93JdP+VNezwfQt88oyxqW6bdCC5xmbo= github.com/aws/aws-sdk-go-v2/service/ec2 v1.336.1 h1:qiuU5+MtLJV2CAxLZYA/GPuvrsScBIk2am+QNAoHmMM= github.com/aws/aws-sdk-go-v2/service/ec2 v1.336.1/go.mod h1:d0e0acsyS3WnFCFJiByGwnUgPpn2wAk97PTIksHN2NI= github.com/aws/aws-sdk-go-v2/service/iam v1.64.1 h1:Uwitin0mXJ7iG5rFuuja3aG9/c84LpyyZUhaTiwZj7w= diff --git a/internal/api/web/app.js b/internal/api/web/app.js index c7a0f80..36fec07 100644 --- a/internal/api/web/app.js +++ b/internal/api/web/app.js @@ -196,7 +196,12 @@ function escapeHTML(value) { .replaceAll(">", ">"); } -function renderPathDetail(paths, edges, truncation) { +function auditsForPath(audits, index) { + const row = (audits || []).find((item) => item.index === index); + return row?.events || []; +} + +function renderPathDetail(paths, edges, truncation, audits) { const panel = document.getElementById("path-detail"); if (!paths.length && !truncation) { panel.classList.add("hidden"); @@ -218,8 +223,13 @@ function renderPathDetail(paths, edges, truncation) { ``, ); } + const events = auditsForPath(audits, index).map((event) => { + const bits = [event.time, event.name, event.principal, event.resource].filter(Boolean); + return `
  • ${escapeHTML(bits.join(" "))}
  • `; + }).join(""); + const eventList = events ? `

    CloudTrail

    ` : ""; const names = path.map((node) => escapeHTML(node.name)).join(" → "); - return `

    Path ${index + 1}. ${names}

      ${hops.join("")}
    `; + return `

    Path ${index + 1}. ${names}

      ${hops.join("")}
    ${eventList}`; }).join(""); } @@ -290,7 +300,7 @@ async function refresh() { if (query) { const result = await fetchJSON(`/v1/graph/query?name=${encodeURIComponent(query)}`); const paths = result.paths || []; - renderPathDetail(paths, snapshot.edges, result.truncated ? result.truncation : ""); + renderPathDetail(paths, snapshot.edges, result.truncated ? result.truncation : "", result.audits); const built = pathGraph(paths, snapshot.edges); renderGraph(built.nodes, built.edges); return; diff --git a/internal/cmd/paths.go b/internal/cmd/paths.go index 15c3efb..a0c8ef8 100644 --- a/internal/cmd/paths.go +++ b/internal/cmd/paths.go @@ -42,8 +42,22 @@ var pathsRunCmd = &cobra.Command{ fmt.Println("No paths found.") return nil } + audits := map[int][]graph.AuditEvent{} + for _, audit := range result.Audits { + audits[audit.Index] = audit.Events + } for i, path := range result.Paths { fmt.Printf("%d. %s\n", i+1, graph.FormatPath(path)) + for _, event := range audits[i] { + fmt.Printf(" %s %s", event.Time, event.Name) + if event.Principal != "" { + fmt.Printf(" %s", event.Principal) + } + if event.Resource != "" { + fmt.Printf(" → %s", event.Resource) + } + fmt.Println() + } } return nil }, diff --git a/internal/cmd/scan.go b/internal/cmd/scan.go index e5b0801..27b1772 100644 --- a/internal/cmd/scan.go +++ b/internal/cmd/scan.go @@ -29,7 +29,7 @@ var scanCmd = &cobra.Command{ var scanAWSCmd = &cobra.Command{ Use: "aws", - Short: "Scan the current AWS account (EC2, IAM, S3, security groups)", + Short: "Scan the current AWS account (EC2, IAM, S3, security groups) and recent CloudTrail management events", RunE: func(cmd *cobra.Command, args []string) error { cfg := loadConfig() collector, err := aws.NewCollector(cmd.Context(), cfg.AWSRegion) diff --git a/internal/collectors/aws/audit.go b/internal/collectors/aws/audit.go new file mode 100644 index 0000000..1ea2053 --- /dev/null +++ b/internal/collectors/aws/audit.go @@ -0,0 +1,539 @@ +// Copyright 2026 OpenSourceOM +// SPDX-License-Identifier: Apache-2.0 + +package aws + +import ( + "context" + "encoding/json" + "strings" + "time" + + "github.com/OpenSourceOM/core/internal/graph" + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/service/cloudtrail" + trailtypes "github.com/aws/aws-sdk-go-v2/service/cloudtrail/types" +) + +const ( + auditLookupWindow = 24 * time.Hour + auditPageSize = 50 + auditMaxPages = 4 + auditMaxPerNode = 5 +) + +// auditEventNames are management events that show a principal using or +// changing an identity, datastore, workload, or network control. +// LookupEvents does not return S3 object data events such as GetObject. +var auditEventNames = map[string]struct{}{ + "AssumeRole": {}, + "AssumeRoleWithSAML": {}, + "AssumeRoleWithWebIdentity": {}, + "ConsoleLogin": {}, + "PutBucketPolicy": {}, + "PutBucketAcl": {}, + "PutPublicAccessBlock": {}, + "DeletePublicAccessBlock": {}, + "AuthorizeSecurityGroupIngress": {}, + "AuthorizeSecurityGroupEgress": {}, + "RevokeSecurityGroupIngress": {}, + "RevokeSecurityGroupEgress": {}, + "CreateAccessKey": {}, + "AttachRolePolicy": {}, + "AttachUserPolicy": {}, + "PutRolePolicy": {}, + "PutUserPolicy": {}, + "RunInstances": {}, + "ModifyDBInstance": {}, +} + +// cloudTrailAPI is the LookupEvents call the collector pages. +type cloudTrailAPI interface { + LookupEvents(context.Context, *cloudtrail.LookupEventsInput, ...func(*cloudtrail.Options)) (*cloudtrail.LookupEventsOutput, error) +} + +var _ cloudTrailAPI = (*cloudtrail.Client)(nil) + +type trailPayload struct { + UserIdentity struct { + ARN string `json:"arn"` + AccountID string `json:"accountId"` + SessionContext struct { + SessionIssuer struct { + ARN string `json:"arn"` + } `json:"sessionIssuer"` + } `json:"sessionContext"` + } `json:"userIdentity"` + SourceIPAddress string `json:"sourceIPAddress"` + RequestParameters map[string]any `json:"requestParameters"` + Resources []struct { + ARN string `json:"ARN"` + } `json:"resources"` +} + +type matchedAudit struct { + event graph.AuditEvent + when time.Time +} + +// attachAuditEvents stores recent CloudTrail management events on the +// identity and resource nodes they name. A lookup error omits the property. +// Events are kept only when the resource sits on an exposed path in this batch: +// internet-reachable, assumed by such a workload, able to reach that workload's +// data, or a public datastore and the identities that can access it. +func (c *Collector) attachAuditEvents(ctx context.Context, client cloudTrailAPI, batch *graph.Batch) { + if client == nil || batch == nil || c.AccountID == "" { + return + } + exposed := exposedNodeIDs(batch) + if len(exposed) == 0 { + return + } + index := buildAuditIndex(batch) + var matched []matchedAudit + seen := map[string]bool{} + for _, region := range c.auditRegions() { + events, _ := c.lookupRegion(ctx, client, region) + for _, event := range events { + item, ok := c.matchAuditEvent(event, index, exposed) + if !ok || seen[item.event.ID] { + continue + } + seen[item.event.ID] = true + matched = append(matched, item) + } + } + if len(matched) == 0 { + return + } + sortAuditsNewestFirst(matched) + counts := map[string]int{} + for _, item := range matched { + c.addAuditEvent(batch, item.event.PrincipalNodeID, item.event, counts) + if item.event.ResourceNodeID != item.event.PrincipalNodeID { + c.addAuditEvent(batch, item.event.ResourceNodeID, item.event, counts) + } + } +} + +func (c *Collector) auditRegions() []string { + region := c.Region + if region == "" { + region = "us-east-1" + } + regions := []string{region} + if region != "us-east-1" { + regions = append(regions, "us-east-1") + } + return regions +} + +func (c *Collector) lookupRegion(ctx context.Context, client cloudTrailAPI, region string) ([]trailtypes.Event, error) { + start := time.Now().Add(-auditLookupWindow) + var token *string + var events []trailtypes.Event + for page := 0; page < auditMaxPages; page++ { + out, err := client.LookupEvents(ctx, &cloudtrail.LookupEventsInput{ + StartTime: &start, + MaxResults: aws.Int32(auditPageSize), + NextToken: token, + }, func(o *cloudtrail.Options) { + o.Region = region + }) + if err != nil { + return events, err + } + if out == nil { + return events, nil + } + events = append(events, out.Events...) + if out.NextToken == nil || *out.NextToken == "" { + return events, nil + } + token = out.NextToken + } + return events, nil +} + +func (c *Collector) matchAuditEvent(event trailtypes.Event, index *auditIndex, exposed map[string]bool) (matchedAudit, bool) { + name := aws.ToString(event.EventName) + if _, ok := auditEventNames[name]; !ok { + return matchedAudit{}, false + } + id := aws.ToString(event.EventId) + if id == "" || event.EventTime == nil { + return matchedAudit{}, false + } + payload := parseTrailPayload(aws.ToString(event.CloudTrailEvent)) + if payload.UserIdentity.AccountID != "" && payload.UserIdentity.AccountID != c.AccountID { + return matchedAudit{}, false + } + principalARN, principalID, ok := index.principal(payload, c.AccountID) + if !ok { + return matchedAudit{}, false + } + resourceRaw, resourceID, ok := index.resource(payload, event.Resources, principalID, exposed, c.AccountID) + if !ok { + return matchedAudit{}, false + } + return matchedAudit{ + when: event.EventTime.UTC(), + event: graph.AuditEvent{ + ID: id, + Name: name, + Time: event.EventTime.UTC().Format(time.RFC3339), + Principal: principalARN, + Resource: resourceRaw, + PrincipalNodeID: principalID, + ResourceNodeID: resourceID, + SourceIP: payload.SourceIPAddress, + ReadOnly: strings.EqualFold(aws.ToString(event.ReadOnly), "true"), + }, + }, true +} + +func parseTrailPayload(raw string) trailPayload { + var payload trailPayload + if raw == "" { + return payload + } + _ = json.Unmarshal([]byte(raw), &payload) + return payload +} + +func (c *Collector) addAuditEvent(batch *graph.Batch, nodeID string, event graph.AuditEvent, counts map[string]int) { + if nodeID == "" || counts[nodeID] >= auditMaxPerNode { + return + } + for i := range batch.Nodes { + if batch.Nodes[i].ID != nodeID { + continue + } + if batch.Nodes[i].Properties == nil { + batch.Nodes[i].Properties = map[string]any{} + } + events := append(graph.AuditEventsFrom(batch.Nodes[i].Properties), event) + graph.SetAuditEvents(batch.Nodes[i].Properties, events) + counts[nodeID]++ + return + } +} + +func sortAuditsNewestFirst(events []matchedAudit) { + for i := 1; i < len(events); i++ { + item := events[i] + j := i + for j > 0 && events[j-1].when.Before(item.when) { + events[j] = events[j-1] + j-- + } + events[j] = item + } +} + +func exposedNodeIDs(batch *graph.Batch) map[string]bool { + exposed := map[string]bool{} + reachable := map[string]bool{} + for _, edge := range batch.Edges { + if edge.Type == graph.EdgeReachable && edge.SourceID == graph.InternetNodeID { + reachable[edge.TargetID] = true + exposed[edge.TargetID] = true + } + } + assumed := map[string]bool{} + for _, edge := range batch.Edges { + if edge.Type == graph.EdgeAssumes && reachable[edge.SourceID] { + assumed[edge.TargetID] = true + exposed[edge.TargetID] = true + } + } + for _, edge := range batch.Edges { + if edge.Type != graph.EdgeCanAccess { + continue + } + if reachable[edge.SourceID] || assumed[edge.SourceID] { + exposed[edge.TargetID] = true + } + } + for _, edge := range batch.Edges { + if edge.Type == graph.EdgeAffects && reachable[edge.SourceID] { + exposed[edge.TargetID] = true + } + } + public := map[string]bool{} + for _, node := range batch.Nodes { + if node.Type != graph.NodeDatastore || !boolProp(node.Properties, "public_access") { + continue + } + public[node.ID] = true + exposed[node.ID] = true + } + for _, edge := range batch.Edges { + if edge.Type == graph.EdgeCanAccess && public[edge.TargetID] { + exposed[edge.SourceID] = true + } + } + return exposed +} + +func boolProp(props map[string]any, key string) bool { + if props == nil { + return false + } + value, ok := props[key].(bool) + return ok && value +} + +type auditIndex struct { + byKey map[string]string + isIdentity map[string]bool +} + +func buildAuditIndex(batch *graph.Batch) *auditIndex { + index := &auditIndex{ + byKey: map[string]string{}, + isIdentity: map[string]bool{}, + } + for _, node := range batch.Nodes { + if node.Type == graph.NodeIdentity { + index.isIdentity[node.ID] = true + if kind, ok := stringProp(node.Properties, "principal_type"); ok && node.Name != "" { + index.add(kind+":"+node.Name, node.ID) + } + } + if arn, _ := stringProp(node.Properties, "arn"); arn != "" { + index.add("arn:"+arn, node.ID) + } + if rid, _ := stringProp(node.Properties, "resource_id"); rid != "" { + index.add("id:"+rid, node.ID) + if strings.HasPrefix(rid, "arn:") { + index.add("arn:"+rid, node.ID) + } + } + if service, _ := stringProp(node.Properties, "service"); service == "s3" && node.Name != "" { + index.add("s3:"+node.Name, node.ID) + index.add("id:"+node.Name, node.ID) + } + } + return index +} + +func (index *auditIndex) add(key, id string) { + if key == "" || id == "" { + return + } + if prev, ok := index.byKey[key]; ok && prev != id { + index.byKey[key] = "" + return + } + index.byKey[key] = id +} + +func (index *auditIndex) get(key string) (string, bool) { + id, ok := index.byKey[key] + return id, ok && id != "" +} + +func (index *auditIndex) principal(payload trailPayload, accountID string) (string, string, bool) { + candidates := []string{ + payload.UserIdentity.ARN, + payload.UserIdentity.SessionContext.SessionIssuer.ARN, + } + if role, ok := iamRoleARN(payload.UserIdentity.ARN); ok { + candidates = append(candidates, role) + } + for _, raw := range candidates { + if raw == "" || !arnInAccount(raw, accountID) { + continue + } + id, ok := index.get("arn:" + raw) + if ok && index.isIdentity[id] { + return raw, id, true + } + if role, ok := iamRoleARN(raw); ok && role != raw && arnInAccount(role, accountID) { + id, ok = index.get("arn:" + role) + if ok && index.isIdentity[id] { + return role, id, true + } + } + } + return "", "", false +} + +func (index *auditIndex) resource(payload trailPayload, listed []trailtypes.Resource, principalID string, exposed map[string]bool, accountID string) (string, string, bool) { + var candidates []string + for _, resource := range listed { + if name := aws.ToString(resource.ResourceName); name != "" { + candidates = append(candidates, name) + } + } + for _, resource := range payload.Resources { + if resource.ARN != "" { + candidates = append(candidates, resource.ARN) + } + } + for _, key := range []string{"roleArn", "bucketName", "userName", "groupId", "instanceId"} { + if value, ok := payload.RequestParameters[key].(string); ok && value != "" { + candidates = append(candidates, value) + } + } + var fallbackRaw, fallbackID string + for _, raw := range candidates { + id, ok := index.resolve(raw, accountID) + if !ok || !exposed[id] { + continue + } + if id != principalID { + return raw, id, true + } + if fallbackID == "" { + fallbackRaw, fallbackID = raw, id + } + } + if fallbackID != "" && exposed[principalID] { + return fallbackRaw, fallbackID, true + } + if exposed[principalID] { + for _, raw := range []string{payload.UserIdentity.ARN, payload.UserIdentity.SessionContext.SessionIssuer.ARN} { + if raw == "" { + continue + } + return raw, principalID, true + } + } + return "", "", false +} + +func (index *auditIndex) resolve(raw, accountID string) (string, bool) { + raw = strings.TrimSpace(raw) + if raw == "" || raw == "*" { + return "", false + } + if !arnInAccount(raw, accountID) { + return "", false + } + if id, ok := index.get("arn:" + raw); ok { + return id, true + } + if id, ok := index.get("id:" + raw); ok { + return id, true + } + if bucket, ok := s3Bucket(raw); ok { + if id, ok := index.get("s3:" + bucket); ok { + return id, true + } + } + if id, ok := ec2ResourceID(raw); ok { + if node, ok := index.get("id:" + id); ok { + return node, true + } + } + if role, ok := iamRoleARN(raw); ok { + if node, ok := index.get("arn:" + role); ok { + return node, true + } + } + if !strings.Contains(raw, ":") { + if id, ok := index.bareIdentity(raw); ok { + return id, true + } + } + return "", false +} + +func (index *auditIndex) bareIdentity(raw string) (string, bool) { + user, userOK := index.get("user:" + raw) + role, roleOK := index.get("role:" + raw) + switch { + case userOK && roleOK: + return "", false + case userOK: + return user, true + case roleOK: + return role, true + default: + return "", false + } +} + +func stringProp(props map[string]any, key string) (string, bool) { + if props == nil { + return "", false + } + value, ok := props[key].(string) + return value, ok && value != "" +} + +func arnInAccount(raw, accountID string) bool { + if !strings.HasPrefix(raw, "arn:") { + return true + } + account := arnField(raw, 4) + return account == "" || account == accountID +} + +func arnField(raw string, index int) string { + parts := strings.Split(raw, ":") + if len(parts) <= index || parts[0] != "arn" { + return "" + } + return parts[index] +} + +func s3Bucket(raw string) (string, bool) { + parts := strings.SplitN(raw, ":", 6) + if len(parts) != 6 || parts[0] != "arn" || parts[2] != "s3" { + return "", false + } + bucket := parts[5] + if i := strings.Index(bucket, "/"); i >= 0 { + bucket = bucket[:i] + } + if bucket == "" { + return "", false + } + return bucket, true +} + +func ec2ResourceID(raw string) (string, bool) { + if strings.HasPrefix(raw, "i-") || strings.HasPrefix(raw, "sg-") { + return raw, true + } + parts := strings.SplitN(raw, ":", 6) + if len(parts) != 6 || parts[0] != "arn" || parts[2] != "ec2" { + return "", false + } + switch { + case strings.HasPrefix(parts[5], "instance/"): + return strings.TrimPrefix(parts[5], "instance/"), true + case strings.HasPrefix(parts[5], "security-group/"): + return strings.TrimPrefix(parts[5], "security-group/"), true + default: + return "", false + } +} + +func iamRoleARN(raw string) (string, bool) { + if strings.Contains(raw, ":role/") && strings.HasPrefix(raw, "arn:") { + return raw, true + } + const marker = ":assumed-role/" + i := strings.Index(raw, marker) + if i < 0 || !strings.HasPrefix(raw, "arn:") { + return "", false + } + rest := raw[i+len(marker):] + name := rest + if slash := strings.Index(rest, "/"); slash >= 0 { + name = rest[:slash] + } + if name == "" { + return "", false + } + partition := arnField(raw, 1) + account := arnField(raw, 4) + if partition == "" || account == "" { + return "", false + } + return "arn:" + partition + ":iam::" + account + ":role/" + name, true +} diff --git a/internal/collectors/aws/audit_test.go b/internal/collectors/aws/audit_test.go new file mode 100644 index 0000000..eb537af --- /dev/null +++ b/internal/collectors/aws/audit_test.go @@ -0,0 +1,308 @@ +// Copyright 2026 OpenSourceOM +// SPDX-License-Identifier: Apache-2.0 + +package aws + +import ( + "context" + "errors" + "fmt" + "strings" + "testing" + "time" + + "github.com/OpenSourceOM/core/internal/graph" + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/service/cloudtrail" + trailtypes "github.com/aws/aws-sdk-go-v2/service/cloudtrail/types" +) + +func TestAttachAuditEventsLinksAssumeRoleOnExposedPath(t *testing.T) { + c := &Collector{Region: "us-east-1", AccountID: "111122223333"} + batch := exposedAuditBatch(c) + role := c.globalNodeID("identity", "Admin") + user := c.globalNodeID("identity", "user/alice") + when := time.Date(2026, 9, 30, 15, 0, 0, 0, time.UTC) + client := &fakeTrail{ + pages: map[string][]trailtypes.Event{ + "": {assumeRoleEvent("evt-assume", when, "arn:aws:iam::111122223333:user/alice", "arn:aws:iam::111122223333:role/Admin")}, + }, + } + + c.attachAuditEvents(context.Background(), client, &batch) + + roleEvents := eventsOn(&batch, role) + userEvents := eventsOn(&batch, user) + if len(roleEvents) != 1 || len(userEvents) != 1 { + t.Fatalf("role events = %d, user events = %d, want both", len(roleEvents), len(userEvents)) + } + if roleEvents[0].ID != "evt-assume" || roleEvents[0].ResourceNodeID != role || roleEvents[0].PrincipalNodeID != user { + t.Fatalf("event = %+v", roleEvents[0]) + } + if roleEvents[0].ReadOnly { + t.Fatal("AssumeRole stored as read-only") + } + if client.regions[0] != "us-east-1" || len(client.regions) != 1 { + t.Fatalf("regions = %v, want one us-east-1 lookup", client.regions) + } +} + +func TestAttachAuditEventsSkipsUnexposedRole(t *testing.T) { + c := &Collector{Region: "us-east-1", AccountID: "111122223333"} + batch := exposedAuditBatch(c) + batch.Edges = nil + for i := range batch.Nodes { + if batch.Nodes[i].Properties != nil { + delete(batch.Nodes[i].Properties, "public_access") + } + } + when := time.Date(2026, 9, 30, 15, 0, 0, 0, time.UTC) + client := &fakeTrail{pages: map[string][]trailtypes.Event{ + "": {assumeRoleEvent("evt-assume", when, "arn:aws:iam::111122223333:user/alice", "arn:aws:iam::111122223333:role/Admin")}, + }} + + c.attachAuditEvents(context.Background(), client, &batch) + + if n := countAuditEvents(&batch); n != 0 { + t.Fatalf("stored %d events on a graph with no exposed path", n) + } +} + +func TestAttachAuditEventsOmitsOnLookupError(t *testing.T) { + c := &Collector{Region: "us-east-1", AccountID: "111122223333"} + batch := exposedAuditBatch(c) + client := &fakeTrail{err: errors.New("access denied")} + + c.attachAuditEvents(context.Background(), client, &batch) + + if n := countAuditEvents(&batch); n != 0 { + t.Fatalf("stored %d events after a failed lookup", n) + } +} + +func TestAttachAuditEventsFollowsPagesAndCaps(t *testing.T) { + c := &Collector{Region: "us-east-1", AccountID: "111122223333"} + batch := exposedAuditBatch(c) + role := c.globalNodeID("identity", "Admin") + user := "arn:aws:iam::111122223333:user/alice" + roleARN := "arn:aws:iam::111122223333:role/Admin" + base := time.Date(2026, 9, 30, 15, 0, 0, 0, time.UTC) + var first, second []trailtypes.Event + for i := 0; i < 3; i++ { + first = append(first, assumeRoleEvent(fmt.Sprintf("evt-%d", i), base.Add(time.Duration(i)*time.Minute), user, roleARN)) + } + for i := 3; i < 6; i++ { + second = append(second, assumeRoleEvent(fmt.Sprintf("evt-%d", i), base.Add(time.Duration(i)*time.Minute), user, roleARN)) + } + client := &fakeTrail{ + pages: map[string][]trailtypes.Event{"": first, "page-2": second}, + next: map[string]string{"": "page-2"}, + } + + c.attachAuditEvents(context.Background(), client, &batch) + + got := eventsOn(&batch, role) + if len(got) != auditMaxPerNode { + t.Fatalf("events = %d, want the cap %d", len(got), auditMaxPerNode) + } + if got[0].ID != "evt-5" || got[len(got)-1].ID != "evt-1" { + t.Fatalf("order = %s then %s, want newest evt-5 down to evt-1", got[0].ID, got[len(got)-1].ID) + } +} + +func TestAttachAuditEventsStopsPaging(t *testing.T) { + c := &Collector{Region: "us-east-1", AccountID: "111122223333"} + batch := exposedAuditBatch(c) + client := &fakeTrail{endless: true} + + c.attachAuditEvents(context.Background(), client, &batch) + + if client.calls != auditMaxPages { + t.Fatalf("calls = %d, want the page cap %d", client.calls, auditMaxPages) + } +} + +func TestAttachAuditEventsLooksUpGlobalEvents(t *testing.T) { + c := &Collector{Region: "eu-west-1", AccountID: "111122223333"} + batch := exposedAuditBatch(c) + when := time.Date(2026, 9, 30, 15, 0, 0, 0, time.UTC) + client := &fakeTrail{ + pages: map[string][]trailtypes.Event{ + "us-east-1": {assumeRoleEvent("evt-global", when, "arn:aws:iam::111122223333:user/alice", "arn:aws:iam::111122223333:role/Admin")}, + }, + } + + c.attachAuditEvents(context.Background(), client, &batch) + + if strings.Join(client.regions, ",") != "eu-west-1,us-east-1" { + t.Fatalf("regions = %v", client.regions) + } + if len(eventsOn(&batch, c.globalNodeID("identity", "Admin"))) != 1 { + t.Fatal("us-east-1 AssumeRole was not attached") + } +} + +func TestAttachAuditEventsKeepsPublicBucketChange(t *testing.T) { + c := &Collector{Region: "us-east-1", AccountID: "111122223333"} + batch := exposedAuditBatch(c) + bucket := c.globalNodeID("datastore", "logs") + when := time.Date(2026, 9, 30, 16, 0, 0, 0, time.UTC) + payload := `{"userIdentity":{"arn":"arn:aws:iam::111122223333:user/alice","accountId":"111122223333"},"sourceIPAddress":"203.0.113.9","requestParameters":{"bucketName":"logs"}}` + client := &fakeTrail{pages: map[string][]trailtypes.Event{ + "": {{ + EventId: aws.String("evt-policy"), + EventName: aws.String("PutBucketPolicy"), + EventTime: &when, + ReadOnly: aws.String("false"), + CloudTrailEvent: aws.String(payload), + }}, + }} + + c.attachAuditEvents(context.Background(), client, &batch) + + got := eventsOn(&batch, bucket) + if len(got) != 1 || got[0].Name != "PutBucketPolicy" || got[0].ResourceNodeID != bucket { + t.Fatalf("bucket events = %+v", got) + } +} + +func TestAttachAuditEventsDropsDataEventsAndOtherAccounts(t *testing.T) { + c := &Collector{Region: "us-east-1", AccountID: "111122223333"} + batch := exposedAuditBatch(c) + when := time.Date(2026, 9, 30, 15, 0, 0, 0, time.UTC) + client := &fakeTrail{pages: map[string][]trailtypes.Event{ + "": { + assumeRoleEvent("evt-data", when, "arn:aws:iam::111122223333:user/alice", "arn:aws:iam::111122223333:role/Admin"), + assumeRoleEvent("evt-other", when, "arn:aws:iam::999999999999:user/alice", "arn:aws:iam::999999999999:role/Admin"), + }, + }} + client.pages[""][0].EventName = aws.String("GetObject") + + c.attachAuditEvents(context.Background(), client, &batch) + + if n := countAuditEvents(&batch); n != 0 { + t.Fatalf("stored %d events, want data events and other accounts dropped", n) + } +} + +func TestAttachAuditEventsMapsAssumedRoleSession(t *testing.T) { + c := &Collector{Region: "us-east-1", AccountID: "111122223333"} + batch := exposedAuditBatch(c) + role := c.globalNodeID("identity", "Admin") + bucket := c.globalNodeID("datastore", "logs") + when := time.Date(2026, 9, 30, 15, 0, 0, 0, time.UTC) + payload := `{"userIdentity":{"arn":"arn:aws:sts::111122223333:assumed-role/Admin/i-1","accountId":"111122223333","sessionContext":{"sessionIssuer":{"arn":"arn:aws:iam::111122223333:role/Admin"}}},"requestParameters":{"bucketName":"logs"}}` + client := &fakeTrail{pages: map[string][]trailtypes.Event{ + "": {{ + EventId: aws.String("evt-session"), + EventName: aws.String("PutBucketAcl"), + EventTime: &when, + CloudTrailEvent: aws.String(payload), + }}, + }} + + c.attachAuditEvents(context.Background(), client, &batch) + + got := eventsOn(&batch, bucket) + if len(got) != 1 || got[0].PrincipalNodeID != role || got[0].ResourceNodeID != bucket { + t.Fatalf("event = %+v", got) + } +} + +func exposedAuditBatch(c *Collector) graph.Batch { + workload := c.nodeID("workload", "i-1") + role := c.globalNodeID("identity", "Admin") + user := c.globalNodeID("identity", "user/alice") + bucket := c.globalNodeID("datastore", "logs") + return graph.Batch{ + Nodes: []graph.Node{ + {ID: graph.InternetNodeID, Type: graph.NodeInternet, Name: "Internet"}, + {ID: workload, Type: graph.NodeWorkload, Name: "web", Provider: "aws", Properties: map[string]any{"resource_id": "i-1"}}, + {ID: role, Type: graph.NodeIdentity, Name: "Admin", Provider: "aws", Properties: map[string]any{ + "arn": "arn:aws:iam::111122223333:role/Admin", "principal_type": "role", + }}, + {ID: user, Type: graph.NodeIdentity, Name: "alice", Provider: "aws", Properties: map[string]any{ + "arn": "arn:aws:iam::111122223333:user/alice", "principal_type": "user", + }}, + {ID: bucket, Type: graph.NodeDatastore, Name: "logs", Provider: "aws", Properties: map[string]any{ + "service": "s3", "resource_id": "logs", "public_access": true, + }}, + }, + Edges: []graph.Edge{ + {SourceID: graph.InternetNodeID, TargetID: workload, Type: graph.EdgeReachable}, + {SourceID: workload, TargetID: role, Type: graph.EdgeAssumes}, + {SourceID: role, TargetID: bucket, Type: graph.EdgeCanAccess}, + }, + } +} + +func assumeRoleEvent(id string, when time.Time, principal, role string) trailtypes.Event { + payload := fmt.Sprintf(`{"userIdentity":{"arn":%q,"accountId":%q},"sourceIPAddress":"203.0.113.4","requestParameters":{"roleArn":%q}}`, principal, arnField(principal, 4), role) + return trailtypes.Event{ + EventId: aws.String(id), + EventName: aws.String("AssumeRole"), + EventTime: &when, + ReadOnly: aws.String("false"), + CloudTrailEvent: aws.String(payload), + Resources: []trailtypes.Resource{{ + ResourceName: aws.String(role), + ResourceType: aws.String("AWS::IAM::Role"), + }}, + } +} + +func eventsOn(batch *graph.Batch, id string) []graph.AuditEvent { + for _, node := range batch.Nodes { + if node.ID == id { + return graph.AuditEventsFrom(node.Properties) + } + } + return nil +} + +func countAuditEvents(batch *graph.Batch) int { + n := 0 + for _, node := range batch.Nodes { + n += len(graph.AuditEventsFrom(node.Properties)) + } + return n +} + +type fakeTrail struct { + pages map[string][]trailtypes.Event + next map[string]string + err error + endless bool + regions []string + calls int +} + +func (f *fakeTrail) LookupEvents(_ context.Context, params *cloudtrail.LookupEventsInput, optFns ...func(*cloudtrail.Options)) (*cloudtrail.LookupEventsOutput, error) { + f.calls++ + opts := cloudtrail.Options{} + for _, fn := range optFns { + fn(&opts) + } + f.regions = append(f.regions, opts.Region) + if f.err != nil { + return nil, f.err + } + if f.endless { + token := "more" + return &cloudtrail.LookupEventsOutput{NextToken: &token}, nil + } + key := "" + if params != nil && params.NextToken != nil { + key = *params.NextToken + } + if opts.Region == "us-east-1" { + if events, ok := f.pages["us-east-1"]; ok && key == "" { + return &cloudtrail.LookupEventsOutput{Events: events}, nil + } + } + out := &cloudtrail.LookupEventsOutput{Events: f.pages[key]} + if next := f.next[key]; next != "" { + out.NextToken = aws.String(next) + } + return out, nil +} diff --git a/internal/collectors/aws/collector.go b/internal/collectors/aws/collector.go index 6343720..b72784c 100644 --- a/internal/collectors/aws/collector.go +++ b/internal/collectors/aws/collector.go @@ -11,6 +11,7 @@ import ( "github.com/OpenSourceOM/core/internal/graph" "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/service/cloudtrail" "github.com/aws/aws-sdk-go-v2/service/ec2" ec2types "github.com/aws/aws-sdk-go-v2/service/ec2/types" "github.com/aws/aws-sdk-go-v2/service/iam" @@ -86,6 +87,8 @@ func (c *Collector) Collect(ctx context.Context) (graph.Batch, error) { if err := c.linkInstanceProfileAccess(ctx, iamClient, &batch, profiles); err != nil { return graph.Batch{}, err } + // A failed lookup leaves audit_events unset. Inventory still ingests. + c.attachAuditEvents(ctx, cloudtrail.NewFromConfig(cfg), &batch) return batch, nil } diff --git a/internal/graph/audit.go b/internal/graph/audit.go new file mode 100644 index 0000000..c89dc54 --- /dev/null +++ b/internal/graph/audit.go @@ -0,0 +1,104 @@ +// Copyright 2026 OpenSourceOM +// SPDX-License-Identifier: Apache-2.0 + +package graph + +import "encoding/json" + +// AuditEventsProperty is the node property that holds recent cloud audit +// events attached to an identity or resource. The AWS collector writes it. +// A plugin may set the same shape. Azure and GCP collectors do not. +const AuditEventsProperty = "audit_events" + +// AuditEvent is one management event attached to existing graph nodes. +// It is evidence that a principal acted on a resource, not a new node type. +type AuditEvent struct { + ID string `json:"id"` + Name string `json:"name"` + Time string `json:"time"` + Principal string `json:"principal,omitempty"` + Resource string `json:"resource,omitempty"` + PrincipalNodeID string `json:"principal_node_id,omitempty"` + ResourceNodeID string `json:"resource_node_id,omitempty"` + SourceIP string `json:"source_ip,omitempty"` + ReadOnly bool `json:"read_only"` +} + +// PathAudit is the set of audit events whose resource node sits on Paths[Index]. +type PathAudit struct { + Index int `json:"index"` + Events []AuditEvent `json:"events"` +} + +// SetAuditEvents stores events on props. An empty list removes the key. +func SetAuditEvents(props map[string]any, events []AuditEvent) { + if props == nil { + return + } + if len(events) == 0 { + delete(props, AuditEventsProperty) + return + } + raw, err := json.Marshal(events) + if err != nil { + return + } + var decoded any + if err := json.Unmarshal(raw, &decoded); err != nil { + return + } + props[AuditEventsProperty] = decoded +} + +// AuditEventsFrom reads events stored by SetAuditEvents. +// A missing or unreadable value returns nil. +func AuditEventsFrom(props map[string]any) []AuditEvent { + if props == nil { + return nil + } + raw, ok := props[AuditEventsProperty] + if !ok || raw == nil { + return nil + } + encoded, err := json.Marshal(raw) + if err != nil { + return nil + } + var events []AuditEvent + if err := json.Unmarshal(encoded, &events); err != nil { + return nil + } + return events +} + +// AuditsForPaths returns events that show use of a returned path. +// An event is included when its resource node is on that path. +// The same event stored on two nodes of one path is returned once. +func AuditsForPaths(paths [][]Node) []PathAudit { + var out []PathAudit + for i, path := range paths { + onPath := make(map[string]bool, len(path)) + for _, node := range path { + onPath[node.ID] = true + } + seen := map[string]bool{} + var events []AuditEvent + for _, node := range path { + for _, event := range AuditEventsFrom(node.Properties) { + if event.ID == "" || seen[event.ID] { + continue + } + if event.ResourceNodeID == "" || !onPath[event.ResourceNodeID] { + continue + } + seen[event.ID] = true + events = append(events, event) + } + } + if len(events) == 0 { + continue + } + out = append(out, PathAudit{Index: i, Events: events}) + } + return out +} diff --git a/internal/graph/audit_finish_test.go b/internal/graph/audit_finish_test.go new file mode 100644 index 0000000..946dd68 --- /dev/null +++ b/internal/graph/audit_finish_test.go @@ -0,0 +1,26 @@ +// Copyright 2026 OpenSourceOM +// SPDX-License-Identifier: Apache-2.0 + +package graph + +import "testing" + +func TestFinishPathsIncludesAudits(t *testing.T) { + props := map[string]any{} + SetAuditEvents(props, []AuditEvent{{ + ID: "evt-1", + Name: "AssumeRole", + ResourceNodeID: "role", + }}) + paths := [][]Node{{ + {ID: InternetNodeID, Type: NodeInternet, Name: "Internet"}, + {ID: "role", Type: NodeIdentity, Name: "Admin", Properties: props}, + }} + result := finishPaths("internet-to-workload", "summary", paths, 50, false) + if len(result.Audits) != 1 || len(result.Audits[0].Events) != 1 || result.Audits[0].Events[0].ID != "evt-1" { + t.Fatalf("audits = %+v", result.Audits) + } + if result.Audits[0].Index != 0 { + t.Fatalf("index = %d", result.Audits[0].Index) + } +} diff --git a/internal/graph/audit_test.go b/internal/graph/audit_test.go new file mode 100644 index 0000000..e90d13a --- /dev/null +++ b/internal/graph/audit_test.go @@ -0,0 +1,79 @@ +// Copyright 2026 OpenSourceOM +// SPDX-License-Identifier: Apache-2.0 + +package graph_test + +import ( + "testing" + + "github.com/OpenSourceOM/core/internal/graph" +) + +func TestAuditsForPathsKeepsEventsOnTheResource(t *testing.T) { + roleID := "aws:111:global:identity:Admin" + bucketID := "aws:111:global:datastore:logs" + otherID := "aws:111:global:datastore:private" + event := graph.AuditEvent{ + ID: "evt-assume", + Name: "AssumeRole", + Time: "2026-09-30T15:00:00Z", + Principal: "arn:aws:iam::111:user/alice", + Resource: "arn:aws:iam::111:role/Admin", + PrincipalNodeID: "aws:111:global:identity:user/alice", + ResourceNodeID: roleID, + } + bucketEvent := graph.AuditEvent{ + ID: "evt-policy", + Name: "PutBucketPolicy", + Time: "2026-09-30T15:04:00Z", + Resource: "arn:aws:s3:::logs", + ResourceNodeID: bucketID, + } + + roleProps := map[string]any{} + graph.SetAuditEvents(roleProps, []graph.AuditEvent{event}) + bucketProps := map[string]any{} + graph.SetAuditEvents(bucketProps, []graph.AuditEvent{event, bucketEvent}) + + paths := [][]graph.Node{ + { + {ID: graph.InternetNodeID, Type: graph.NodeInternet, Name: "Internet"}, + {ID: roleID, Type: graph.NodeIdentity, Name: "Admin", Properties: roleProps}, + {ID: bucketID, Type: graph.NodeDatastore, Name: "logs", Properties: bucketProps}, + }, + { + {ID: otherID, Type: graph.NodeDatastore, Name: "private"}, + }, + } + + got := graph.AuditsForPaths(paths) + if len(got) != 1 { + t.Fatalf("audits = %d, want the exposed path only", len(got)) + } + if got[0].Index != 0 { + t.Fatalf("index = %d, want 0", got[0].Index) + } + if len(got[0].Events) != 2 { + t.Fatalf("events = %d, want AssumeRole once and PutBucketPolicy", len(got[0].Events)) + } + if got[0].Events[0].ID != "evt-assume" || got[0].Events[1].ID != "evt-policy" { + t.Fatalf("events = %+v", got[0].Events) + } +} + +func TestSetAuditEventsRoundTrip(t *testing.T) { + props := map[string]any{} + graph.SetAuditEvents(props, nil) + if _, ok := props[graph.AuditEventsProperty]; ok { + t.Fatal("empty events left the property set") + } + graph.SetAuditEvents(props, []graph.AuditEvent{{ + ID: "evt-1", + Name: "AssumeRole", + ReadOnly: false, + }}) + got := graph.AuditEventsFrom(props) + if len(got) != 1 || got[0].ID != "evt-1" || got[0].Name != "AssumeRole" || got[0].ReadOnly { + t.Fatalf("events = %+v", got) + } +} diff --git a/internal/graph/query.go b/internal/graph/query.go index 725ff3d..0106105 100644 --- a/internal/graph/query.go +++ b/internal/graph/query.go @@ -299,6 +299,7 @@ func finishPaths(queryName, summary string, paths [][]Node, cap int, depthCut bo Summary: summary, Truncated: truncated, Truncation: note, + Audits: AuditsForPaths(paths), } } diff --git a/internal/graph/types.go b/internal/graph/types.go index 8abf201..c0b91f6 100644 --- a/internal/graph/types.go +++ b/internal/graph/types.go @@ -69,6 +69,9 @@ type PathResult struct { // Truncation names the cap that cut the result: "path cap", "depth cap", // or "path cap and depth cap". Empty when Truncated is false. Truncation string `json:"truncation,omitempty"` + // Audits lists management events whose resource node is on a returned path. + // Index is the position of that path in Paths. + Audits []PathAudit `json:"audits,omitempty"` } type FindingView struct {