diff --git a/README.md b/README.md index 3028ac7..4bc3056 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. `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). +> **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, and `om scan azure` attaches recent Activity Log 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, CloudTrail context, rule packs | +| **3** *(now)* | Graph accuracy, crown-jewel datastores, attack-path findings, CloudTrail and Activity Log context, rule packs | Details: [docs/ROADMAP.md](./docs/ROADMAP.md) diff --git a/collectors/README.md b/collectors/README.md index 4d1cf4f..6001f43 100644 --- a/collectors/README.md +++ b/collectors/README.md @@ -22,7 +22,7 @@ The **demo** collector loads a fixed environment (internet-exposed web tier, pro 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. +The **Azure** collector records logical SQL servers. The same scan reads administrative Activity Log 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. Entra ID sign-in logs are not part of this scan. 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. The **Kubernetes** collector records Ingress objects and whether a NetworkPolicy selects each pod. A pod is internet-reachable from a LoadBalancer or an Ingress backend unless that policy governs ingress and does not allow the world. NodePort alone is not exposure. diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index c9a6110..7dadf8c 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 and recent CloudTrail management events | +| **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; Azure emits recent Activity Log 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. `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. +- Cloud audit logs as graph context. `om scan aws` stores CloudTrail management events from the last 24 hours, and `om scan azure` stores administrative Activity Log events from the same window, 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. GCP Cloud Audit Logs, Entra ID sign-in 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 1ca8167..d159581 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. `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. +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, and `om scan azure` attaches recent Activity Log events, to the identity and resource they name. Path queries return an event when that resource is on the path. Correctness and operability come first: @@ -52,7 +52,7 @@ The path, in order: - [x] **Attack path as the finding** — [#57](https://github.com/OpenSourceOM/core/issues/57) - [x] **Cloud audit logs as graph context** — [#33](https://github.com/OpenSourceOM/core/issues/33) -`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. +`om scan aws` reads CloudTrail management events from the last 24 hours. `om scan azure` reads administrative Activity Log events over the same window. 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. GCP Cloud Audit Logs, Entra ID sign-in 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. diff --git a/docs/adr/001-graph-schema-v0.md b/docs/adr/001-graph-schema-v0.md index afa5b19..6ef5da6 100644 --- a/docs/adr/001-graph-schema-v0.md +++ b/docs/adr/001-graph-schema-v0.md @@ -68,7 +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. +- **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`. `om scan azure` calls the Activity Log list API for the subscription over the same 24 hours and stops after four pages. A fixed list of administrative operations is kept when the event's object id is an identity in the batch and the resource is on an exposed path. A child id, such as a security rule or firewall rule, walks up to the inventoried resource. A role-assignment event uses the assignment scope when the assignment id is not itself a node. A failed lookup omits the property. Entra ID sign-in logs and storage data-plane reads are not in the Activity Log. 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 968031d..85581a8 100644 --- a/go.mod +++ b/go.mod @@ -9,6 +9,7 @@ require ( github.com/Azure/azure-sdk-for-go/sdk/azidentity v1.14.1 github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/authorization/armauthorization/v2 v2.2.0 github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/compute/armcompute/v6 v6.4.0 + github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/monitor/armmonitor v0.13.0 github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/network/armnetwork/v8 v8.0.0 github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/resources/armresources/v3 v3.0.1 github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/storage/armstorage/v3 v3.0.0 diff --git a/go.sum b/go.sum index 355729f..1e8c5ab 100644 --- a/go.sum +++ b/go.sum @@ -32,12 +32,14 @@ github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/authorization/armauthoriza github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/authorization/armauthorization/v2 v2.2.0/go.mod h1:/pz8dyNQe+Ey3yBp/XuYz7oqX8YDNWVpPB0hH3XWfbc= github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/compute/armcompute/v6 v6.4.0 h1:z7Mqz6l0EFH549GvHEqfjKvi+cRScxLWbaoeLm9wxVQ= github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/compute/armcompute/v6 v6.4.0/go.mod h1:v6gbfH+7DG7xH2kUNs+ZJ9tF6O3iNnR85wMtmr+F54o= -github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/internal/v3 v3.1.1 h1:1kpY4qe+BGAH2ykv4baVSqyx+AY5VjXeJ15SldlU6hs= -github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/internal/v3 v3.1.1/go.mod h1:nT6cWpWdUt+g81yuKmjeYPUtI73Ak3yQIT4PVVsCEEQ= +github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/internal/v3 v3.2.0 h1:+lnLQhKh3cgSOIOVH61UZ3s/l9d+bAZp5d/spt1+7UI= +github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/internal/v3 v3.2.0/go.mod h1:tStOHrivWUrcBolspvKV70Us1ckESYGYSHdG4LX8zyY= +github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/monitor/armmonitor v0.13.0 h1:c7r8eBbYWf2JbQFinuEbHsqq+ukY1tVIgAxt0uND2Fo= +github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/monitor/armmonitor v0.13.0/go.mod h1:HCaM3KUBkHyt9NJLP/gFdMa16WWzygEQE5oUw9NjiD4= github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/network/armnetwork/v8 v8.0.0 h1:7QO7GhGat25QEYL4h607O9zNNTUlAv8PbSesW6Ol5Gg= github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/network/armnetwork/v8 v8.0.0/go.mod h1:mCqeYzwyjn/pw0JVqHJMIzfUQJrlcV0YjTg5b0NK+F0= -github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/resources/armdeployments v0.2.0 h1:bYq3jfB2x36hslKMHyge3+esWzROtJNk/4dCjsKlrl4= -github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/resources/armdeployments v0.2.0/go.mod h1:fewgRjNVE84QVVh798sIMFb7gPXPp7NmnekGnboSnXk= +github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/resources/armdeployments v1.0.0 h1:67nFqWXpo0x5Nz0XEb1yI7s8D+EHy8NsTinYw9sZnLk= +github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/resources/armdeployments v1.0.0/go.mod h1:fewgRjNVE84QVVh798sIMFb7gPXPp7NmnekGnboSnXk= github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/resources/armresources v1.2.0 h1:Dd+RhdJn0OTtVGaeDLZpcumkIVCtA/3/Fo42+eoYvVM= github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/resources/armresources v1.2.0/go.mod h1:5kakwfW5CjC9KK+Q4wjXAg+ShuIm2mBMua0ZFj2C8PE= github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/resources/armresources/v3 v3.0.1 h1:guyQA4b8XB2sbJZXzUnOF9mn0WDBv/ZT7me9wTipKtE= diff --git a/internal/cmd/scan.go b/internal/cmd/scan.go index 27b1772..694fad4 100644 --- a/internal/cmd/scan.go +++ b/internal/cmd/scan.go @@ -47,7 +47,7 @@ var scanAWSCmd = &cobra.Command{ var scanAzureCmd = &cobra.Command{ Use: "azure", - Short: "Scan the current Azure subscription (VMs, storage, RBAC)", + Short: "Scan the current Azure subscription (VMs, storage, RBAC) and recent Activity Log events", RunE: func(cmd *cobra.Command, args []string) error { cfg := loadConfig() collector := azure.NewCollector(cfg.AzureSubscriptionID, cfg.AzureLocation) diff --git a/internal/collectors/azure/audit.go b/internal/collectors/azure/audit.go new file mode 100644 index 0000000..c6f83f6 --- /dev/null +++ b/internal/collectors/azure/audit.go @@ -0,0 +1,484 @@ +// Copyright 2026 OpenSourceOM +// SPDX-License-Identifier: Apache-2.0 + +package azure + +import ( + "context" + "fmt" + "strings" + "time" + + "github.com/Azure/azure-sdk-for-go/sdk/azcore" + "github.com/Azure/azure-sdk-for-go/sdk/azcore/runtime" + "github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/monitor/armmonitor" + "github.com/OpenSourceOM/core/internal/graph" +) + +const ( + auditLookupWindow = 24 * time.Hour + auditMaxPages = 4 + auditMaxPerNode = 5 + activityTimeLayout = "2006-01-02T15:04:05.0000000Z" + objectIDClaim = "http://schemas.microsoft.com/identity/claims/objectidentifier" +) + +// auditOperationNames are administrative operations that show a principal +// changing an identity, datastore, workload, or network control. +// The Activity Log API does not return storage data-plane reads, and Entra +// ID sign-in logs are a different source. +var auditOperationNames = map[string]struct{}{ + "microsoft.compute/virtualmachines/write": {}, + "microsoft.network/networksecuritygroups/write": {}, + "microsoft.network/networksecuritygroups/securityrules/write": {}, + "microsoft.network/networksecuritygroups/securityrules/delete": {}, + "microsoft.network/publicipaddresses/write": {}, + "microsoft.network/publicipaddresses/delete": {}, + "microsoft.network/networkinterfaces/write": {}, + "microsoft.storage/storageaccounts/write": {}, + "microsoft.storage/storageaccounts/blobservices/containers/write": {}, + "microsoft.authorization/roleassignments/write": {}, + "microsoft.authorization/roleassignments/delete": {}, + "microsoft.sql/servers/write": {}, + "microsoft.sql/servers/firewallrules/write": {}, + "microsoft.sql/servers/firewallrules/delete": {}, + "microsoft.sql/servers/virtualnetworkrules/write": {}, + "microsoft.sql/servers/virtualnetworkrules/delete": {}, +} + +// activityPager is one page of Activity Log events. +type activityPager interface { + More() bool + NextPage(context.Context) ([]*armmonitor.EventData, error) +} + +// activityLogAPI opens a paged Activity Log query. A nil client omits events. +type activityLogAPI interface { + ListPager(since, until time.Time) (activityPager, error) +} + +type armActivityLogs struct { + client *armmonitor.ActivityLogsClient +} + +func (c *Collector) activityLogs(cred azcore.TokenCredential) activityLogAPI { + if cred == nil || c.SubscriptionID == "" { + return nil + } + client, err := armmonitor.NewActivityLogsClient(c.SubscriptionID, cred, nil) + if err != nil { + return nil + } + return &armActivityLogs{client: client} +} + +func (a *armActivityLogs) ListPager(since, until time.Time) (activityPager, error) { + if a == nil || a.client == nil { + return nil, nil + } + pager := a.client.NewListPager(activityFilter(since, until), nil) + return &sdkActivityPager{inner: pager}, nil +} + +type sdkActivityPager struct { + inner *runtime.Pager[armmonitor.ActivityLogsClientListResponse] +} + +func (p *sdkActivityPager) More() bool { + return p.inner.More() +} + +func (p *sdkActivityPager) NextPage(ctx context.Context) ([]*armmonitor.EventData, error) { + page, err := p.inner.NextPage(ctx) + if err != nil { + return nil, err + } + return page.Value, nil +} + +func activityFilter(since, until time.Time) string { + return fmt.Sprintf("eventTimestamp ge '%s' and eventTimestamp le '%s'", + since.UTC().Format(activityTimeLayout), + until.UTC().Format(activityTimeLayout)) +} + +type matchedAudit struct { + event graph.AuditEvent + when time.Time +} + +// attachAuditEvents stores recent Activity Log 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 activityLogAPI, batch *graph.Batch) { + if client == nil || batch == nil || c.SubscriptionID == "" { + return + } + exposed := exposedNodeIDs(batch) + if len(exposed) == 0 { + return + } + until := time.Now().UTC() + pager, err := client.ListPager(until.Add(-auditLookupWindow), until) + if err != nil || pager == nil { + return + } + events, _ := listActivityEvents(ctx, pager) + index := buildAzureAuditIndex(batch) + var matched []matchedAudit + seen := map[string]bool{} + 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 { + addAuditEvent(batch, item.event.PrincipalNodeID, item.event, counts) + if item.event.ResourceNodeID != item.event.PrincipalNodeID { + addAuditEvent(batch, item.event.ResourceNodeID, item.event, counts) + } + } +} + +func listActivityEvents(ctx context.Context, pager activityPager) ([]*armmonitor.EventData, error) { + if pager == nil { + return nil, nil + } + var events []*armmonitor.EventData + for page := 0; page < auditMaxPages && pager.More(); page++ { + batch, err := pager.NextPage(ctx) + if err != nil { + return events, err + } + events = append(events, batch...) + } + return events, nil +} + +func (c *Collector) matchAuditEvent(event *armmonitor.EventData, index *azureAuditIndex, exposed map[string]bool) (matchedAudit, bool) { + if event == nil || event.EventTimestamp == nil { + return matchedAudit{}, false + } + name := operationName(event) + if _, ok := auditOperationNames[strings.ToLower(name)]; !ok { + return matchedAudit{}, false + } + id := strings.TrimSpace(safeString(event.EventDataID)) + if id == "" || !c.eventInSubscription(event) { + return matchedAudit{}, false + } + principalID, ok := index.principal(event) + if !ok { + return matchedAudit{}, false + } + resourceRaw, resourceID, ok := index.resource(event, principalID, exposed) + if !ok { + return matchedAudit{}, false + } + principal := strings.TrimSpace(safeString(event.Caller)) + if principal == "" { + principal = objectID(event) + } + return matchedAudit{ + when: event.EventTimestamp.UTC(), + event: graph.AuditEvent{ + ID: id, + Name: name, + Time: event.EventTimestamp.UTC().Format(time.RFC3339), + Principal: principal, + Resource: resourceRaw, + PrincipalNodeID: principalID, + ResourceNodeID: resourceID, + SourceIP: clientIP(event), + ReadOnly: requestReadOnly(event), + }, + }, true +} + +func (c *Collector) eventInSubscription(event *armmonitor.EventData) bool { + sub := strings.ToLower(strings.TrimSpace(c.SubscriptionID)) + if sub == "" { + return false + } + if got := strings.ToLower(strings.TrimSpace(safeString(event.SubscriptionID))); got != "" { + return got == sub + } + rid := strings.ToLower(safeString(event.ResourceID)) + return strings.Contains(rid, "/subscriptions/"+sub+"/") +} + +func operationName(event *armmonitor.EventData) string { + if event.OperationName == nil || event.OperationName.Value == nil { + return "" + } + return strings.TrimSpace(*event.OperationName.Value) +} + +func objectID(event *armmonitor.EventData) string { + if id := claim(event.Claims, objectIDClaim); id != "" { + return id + } + if id := claim(event.Claims, "oid"); id != "" { + return id + } + caller := strings.TrimSpace(safeString(event.Caller)) + if isObjectID(caller) { + return caller + } + return "" +} + +func claim(claims map[string]*string, key string) string { + if claims == nil { + return "" + } + return strings.TrimSpace(safeString(claims[key])) +} + +func isObjectID(raw string) bool { + if len(raw) != 36 { + return false + } + for i, r := range raw { + switch i { + case 8, 13, 18, 23: + if r != '-' { + return false + } + default: + if (r < '0' || r > '9') && (r < 'a' || r > 'f') && (r < 'A' || r > 'F') { + return false + } + } + } + return true +} + +func clientIP(event *armmonitor.EventData) string { + if event.HTTPRequest == nil { + return "" + } + return strings.TrimSpace(safeString(event.HTTPRequest.ClientIPAddress)) +} + +func requestReadOnly(event *armmonitor.EventData) bool { + if event.HTTPRequest == nil || event.HTTPRequest.Method == nil { + return false + } + return strings.EqualFold(strings.TrimSpace(*event.HTTPRequest.Method), "GET") +} + +func 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 azureAuditIndex struct { + byKey map[string]string + isIdentity map[string]bool +} + +func buildAzureAuditIndex(batch *graph.Batch) *azureAuditIndex { + index := &azureAuditIndex{ + byKey: map[string]string{}, + isIdentity: map[string]bool{}, + } + for _, node := range batch.Nodes { + if node.Type == graph.NodeIdentity && node.Name != "" { + index.isIdentity[node.ID] = true + index.add("principal:"+strings.ToLower(node.Name), node.ID) + } + if rid, ok := stringProp(node.Properties, "resource_id"); ok { + index.add("id:"+strings.ToLower(rid), node.ID) + } + } + return index +} + +func (index *azureAuditIndex) 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 *azureAuditIndex) get(key string) (string, bool) { + id, ok := index.byKey[key] + return id, ok && id != "" +} + +func (index *azureAuditIndex) principal(event *armmonitor.EventData) (string, bool) { + raw := objectID(event) + if raw == "" { + return "", false + } + id, ok := index.get("principal:" + strings.ToLower(raw)) + if !ok || !index.isIdentity[id] { + return "", false + } + return id, true +} + +func (index *azureAuditIndex) resource(event *armmonitor.EventData, principalID string, exposed map[string]bool) (string, string, bool) { + var candidates []string + if rid := strings.TrimSpace(safeString(event.ResourceID)); rid != "" { + candidates = append(candidates, rid) + } + if event.Authorization != nil { + if scope := strings.TrimSpace(safeString(event.Authorization.Scope)); scope != "" { + candidates = append(candidates, scope) + } + } + var fallbackRaw, fallbackID string + for _, raw := range candidates { + id, ok := index.resolve(raw) + 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] { + raw := strings.TrimSpace(safeString(event.Caller)) + if raw == "" { + raw = objectID(event) + } + if raw == "" { + return "", "", false + } + return raw, principalID, true + } + return "", "", false +} + +func (index *azureAuditIndex) resolve(raw string) (string, bool) { + current := strings.TrimRight(strings.TrimSpace(raw), "/") + if current == "" || current == "*" { + return "", false + } + for { + if id, ok := index.get("id:" + strings.ToLower(current)); ok { + return id, true + } + if subscriptionRoot(current) { + return "", false + } + i := strings.LastIndex(current, "/") + if i <= 0 { + return "", false + } + current = current[:i] + } +} + +func subscriptionRoot(raw string) bool { + rest, ok := strings.CutPrefix(strings.ToLower(strings.TrimRight(raw, "/")), "/subscriptions/") + return ok && rest != "" && !strings.Contains(rest, "/") +} + +func stringProp(props map[string]any, key string) (string, bool) { + if props == nil { + return "", false + } + value, ok := props[key].(string) + return value, ok && value != "" +} diff --git a/internal/collectors/azure/audit_test.go b/internal/collectors/azure/audit_test.go new file mode 100644 index 0000000..b4c7eeb --- /dev/null +++ b/internal/collectors/azure/audit_test.go @@ -0,0 +1,320 @@ +// Copyright 2026 OpenSourceOM +// SPDX-License-Identifier: Apache-2.0 + +package azure + +import ( + "context" + "errors" + "fmt" + "strings" + "testing" + "time" + + "github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/monitor/armmonitor" + "github.com/OpenSourceOM/core/internal/graph" +) + +const ( + auditSub = "00000000-0000-0000-0000-000000000000" + auditOtherSub = "99999999-9999-9999-9999-999999999999" + auditPrincipal = "11111111-1111-1111-1111-111111111111" + auditVMID = "/subscriptions/" + auditSub + "/resourceGroups/rg/providers/Microsoft.Compute/virtualMachines/web" + auditStorageID = "/subscriptions/" + auditSub + "/resourceGroups/rg/providers/Microsoft.Storage/storageAccounts/logs" + auditNSGID = "/subscriptions/" + auditSub + "/resourceGroups/rg/providers/Microsoft.Network/networkSecurityGroups/web" +) + +func TestAttachAuditEventsLinksVMWriteOnExposedPath(t *testing.T) { + c := NewCollector(auditSub, "eastus") + batch := exposedAuditBatch(c) + web := c.nodeID("workload", "web") + identity := c.nodeID("identity", auditPrincipal) + when := time.Date(2026, 9, 30, 15, 0, 0, 0, time.UTC) + event := activityEvent("evt-vm", "Microsoft.Compute/virtualMachines/write", auditVMID, when) + client := &fakeActivity{pager: &scriptedPager{pages: [][]*armmonitor.EventData{{event, event}}}} + + c.attachAuditEvents(context.Background(), client, &batch) + + webEvents := eventsOn(&batch, web) + identityEvents := eventsOn(&batch, identity) + if len(webEvents) != 1 || len(identityEvents) != 1 { + t.Fatalf("web events = %d, identity events = %d, want one each", len(webEvents), len(identityEvents)) + } + if webEvents[0].ID != "evt-vm" || webEvents[0].ResourceNodeID != web || webEvents[0].PrincipalNodeID != identity { + t.Fatalf("event = %+v", webEvents[0]) + } + if webEvents[0].ReadOnly || webEvents[0].SourceIP != "203.0.113.9" || webEvents[0].Principal != "alice@contoso.com" { + t.Fatalf("event = %+v", webEvents[0]) + } + if client.calls != 1 { + t.Fatalf("lookups = %d, want 1", client.calls) + } + if !client.since.Equal(client.until.Add(-auditLookupWindow)) { + t.Fatalf("window = %s to %s", client.since, client.until) + } +} + +func TestAttachAuditEventsSkipsUnexposedVM(t *testing.T) { + c := NewCollector(auditSub, "eastus") + 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 := &fakeActivity{pager: &scriptedPager{pages: [][]*armmonitor.EventData{ + {activityEvent("evt-vm", "Microsoft.Compute/virtualMachines/write", auditVMID, when)}, + }}} + + 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) + } + if client.calls != 0 { + t.Fatalf("lookups = %d, want none when nothing is exposed", client.calls) + } +} + +func TestAttachAuditEventsOmitsOnLookupError(t *testing.T) { + c := NewCollector(auditSub, "eastus") + batch := exposedAuditBatch(c) + client := &fakeActivity{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 := NewCollector(auditSub, "eastus") + batch := exposedAuditBatch(c) + web := c.nodeID("workload", "web") + base := time.Date(2026, 9, 30, 15, 0, 0, 0, time.UTC) + var first, second []*armmonitor.EventData + for i := 0; i < 3; i++ { + first = append(first, activityEvent(fmt.Sprintf("evt-%d", i), "Microsoft.Compute/virtualMachines/write", auditVMID, base.Add(time.Duration(i)*time.Minute))) + } + for i := 3; i < 6; i++ { + second = append(second, activityEvent(fmt.Sprintf("evt-%d", i), "Microsoft.Compute/virtualMachines/write", auditVMID, base.Add(time.Duration(i)*time.Minute))) + } + client := &fakeActivity{pager: &scriptedPager{pages: [][]*armmonitor.EventData{first, second}}} + + c.attachAuditEvents(context.Background(), client, &batch) + + got := eventsOn(&batch, web) + 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 := NewCollector(auditSub, "eastus") + batch := exposedAuditBatch(c) + pager := &scriptedPager{endless: true} + client := &fakeActivity{pager: pager} + + c.attachAuditEvents(context.Background(), client, &batch) + + if pager.n != auditMaxPages { + t.Fatalf("pages = %d, want the page cap %d", pager.n, auditMaxPages) + } +} + +func TestAttachAuditEventsKeepsPublicStorageChange(t *testing.T) { + c := NewCollector(auditSub, "eastus") + batch := exposedAuditBatch(c) + logs := c.nodeID("datastore", "logs") + when := time.Date(2026, 9, 30, 16, 0, 0, 0, time.UTC) + client := &fakeActivity{pager: &scriptedPager{pages: [][]*armmonitor.EventData{ + {activityEvent("evt-storage", "microsoft.storage/storageAccounts/write", strings.ToLower(auditStorageID), when)}, + }}} + + c.attachAuditEvents(context.Background(), client, &batch) + + got := eventsOn(&batch, logs) + if len(got) != 1 || got[0].Name != "microsoft.storage/storageAccounts/write" || got[0].ResourceNodeID != logs { + t.Fatalf("storage events = %+v", got) + } +} + +func TestAttachAuditEventsWalksChildResources(t *testing.T) { + c := NewCollector(auditSub, "eastus") + batch := exposedAuditBatch(c) + nsg := c.nodeID("network", "web") + logs := c.nodeID("datastore", "logs") + when := time.Date(2026, 9, 30, 16, 0, 0, 0, time.UTC) + rule := activityEvent("evt-rule", "Microsoft.Network/networkSecurityGroups/securityRules/write", auditNSGID+"/securityRules/allow-ssh", when) + assignmentID := "/subscriptions/" + auditSub + "/providers/Microsoft.Authorization/roleAssignments/aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee" + assignment := activityEvent("evt-role", "Microsoft.Authorization/roleAssignments/write", assignmentID, when.Add(time.Minute)) + scope := auditStorageID + action := "Microsoft.Authorization/roleAssignments/write" + assignment.Authorization = &armmonitor.SenderAuthorization{Scope: &scope, Action: &action} + + client := &fakeActivity{pager: &scriptedPager{pages: [][]*armmonitor.EventData{{rule, assignment}}}} + c.attachAuditEvents(context.Background(), client, &batch) + + rules := eventsOn(&batch, nsg) + if len(rules) != 1 || rules[0].ID != "evt-rule" || rules[0].Resource != auditNSGID+"/securityRules/allow-ssh" { + t.Fatalf("nsg events = %+v", rules) + } + roles := eventsOn(&batch, logs) + if len(roles) != 1 || roles[0].ID != "evt-role" || roles[0].Resource != auditStorageID { + t.Fatalf("storage events = %+v", roles) + } +} + +func TestAttachAuditEventsDropsDataPlaneAndOtherSubscriptions(t *testing.T) { + c := NewCollector(auditSub, "eastus") + batch := exposedAuditBatch(c) + when := time.Date(2026, 9, 30, 15, 0, 0, 0, time.UTC) + data := activityEvent("evt-data", "Microsoft.Storage/storageAccounts/listKeys/action", auditStorageID, when) + other := activityEvent("evt-other", "Microsoft.Storage/storageAccounts/write", auditStorageID, when) + otherSub := auditOtherSub + other.SubscriptionID = &otherSub + emailOnly := activityEvent("evt-email", "Microsoft.Compute/virtualMachines/write", auditVMID, when) + emailOnly.Claims = nil + client := &fakeActivity{pager: &scriptedPager{pages: [][]*armmonitor.EventData{{data, other, emailOnly}}}} + + c.attachAuditEvents(context.Background(), client, &batch) + + if n := countAuditEvents(&batch); n != 0 { + t.Fatalf("stored %d events that should not match", n) + } +} + +func TestAttachAuditEventsUsesCallerObjectID(t *testing.T) { + c := NewCollector(auditSub, "eastus") + batch := exposedAuditBatch(c) + when := time.Date(2026, 9, 30, 15, 0, 0, 0, time.UTC) + event := activityEvent("evt-mi", "Microsoft.Compute/virtualMachines/write", auditVMID, when) + event.Claims = nil + event.Caller = ptr(auditPrincipal) + client := &fakeActivity{pager: &scriptedPager{pages: [][]*armmonitor.EventData{{event}}}} + + c.attachAuditEvents(context.Background(), client, &batch) + + if len(eventsOn(&batch, c.nodeID("workload", "web"))) != 1 { + t.Fatal("object id caller was not attached") + } +} + +func TestActivityFilterIsSubscriptionWindow(t *testing.T) { + since := time.Date(2026, 9, 29, 15, 4, 5, 0, time.FixedZone("EDT", -4*60*60)) + until := time.Date(2026, 9, 30, 15, 4, 5, 0, time.UTC) + got := activityFilter(since, until) + want := "eventTimestamp ge '2026-09-29T19:04:05.0000000Z' and eventTimestamp le '2026-09-30T15:04:05.0000000Z'" + if got != want { + t.Fatalf("filter = %q", got) + } +} + +func exposedAuditBatch(c *Collector) graph.Batch { + return graph.Batch{ + Nodes: []graph.Node{ + { + ID: c.nodeID("workload", "web"), Type: graph.NodeWorkload, Name: "web", + Properties: map[string]any{"resource_id": auditVMID}, + }, + { + ID: c.nodeID("datastore", "logs"), Type: graph.NodeDatastore, Name: "logs", + Properties: map[string]any{"resource_id": auditStorageID, "public_access": true}, + }, + { + ID: c.nodeID("identity", auditPrincipal), Type: graph.NodeIdentity, Name: auditPrincipal, + }, + { + ID: c.nodeID("network", "web"), Type: graph.NodeNetwork, Name: "web", + Properties: map[string]any{"resource_id": auditNSGID}, + }, + }, + Edges: []graph.Edge{ + {SourceID: graph.InternetNodeID, TargetID: c.nodeID("workload", "web"), Type: graph.EdgeReachable}, + {SourceID: c.nodeID("workload", "web"), TargetID: c.nodeID("identity", auditPrincipal), Type: graph.EdgeAssumes}, + {SourceID: c.nodeID("identity", auditPrincipal), TargetID: c.nodeID("datastore", "logs"), Type: graph.EdgeCanAccess}, + {SourceID: c.nodeID("workload", "web"), TargetID: c.nodeID("network", "web"), Type: graph.EdgeAffects}, + }, + } +} + +func activityEvent(id, operation, resource string, when time.Time) *armmonitor.EventData { + sub := auditSub + caller := "alice@contoso.com" + oid := auditPrincipal + method := "PUT" + ip := "203.0.113.9" + return &armmonitor.EventData{ + EventDataID: &id, + OperationName: &armmonitor.LocalizableString{Value: &operation}, + EventTimestamp: &when, + Caller: &caller, + SubscriptionID: &sub, + ResourceID: &resource, + Claims: map[string]*string{ + objectIDClaim: &oid, + }, + HTTPRequest: &armmonitor.HTTPRequestInfo{ClientIPAddress: &ip, Method: &method}, + } +} + +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 fakeActivity struct { + pager activityPager + err error + calls int + since, until time.Time +} + +func (f *fakeActivity) ListPager(since, until time.Time) (activityPager, error) { + f.calls++ + f.since = since + f.until = until + if f.err != nil { + return nil, f.err + } + return f.pager, nil +} + +type scriptedPager struct { + pages [][]*armmonitor.EventData + n int + endless bool +} + +func (p *scriptedPager) More() bool { + if p.endless { + return true + } + return p.n < len(p.pages) +} + +func (p *scriptedPager) NextPage(context.Context) ([]*armmonitor.EventData, error) { + p.n++ + if p.endless { + return nil, nil + } + return p.pages[p.n-1], nil +} diff --git a/internal/collectors/azure/collector.go b/internal/collectors/azure/collector.go index c982eb5..2e152de 100644 --- a/internal/collectors/azure/collector.go +++ b/internal/collectors/azure/collector.go @@ -69,6 +69,8 @@ func (c *Collector) Collect(ctx context.Context) (graph.Batch, error) { return graph.Batch{}, err } c.linkAzureAccess(&batch, uses, assignments, roles) + // A failed lookup leaves audit_events unset. Inventory still ingests. + c.attachAuditEvents(ctx, c.activityLogs(cred), &batch) return batch, nil } diff --git a/internal/graph/audit.go b/internal/graph/audit.go index c89dc54..ec4ddbe 100644 --- a/internal/graph/audit.go +++ b/internal/graph/audit.go @@ -6,8 +6,8 @@ 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. +// events attached to an identity or resource. The AWS and Azure collectors +// write it. A plugin may set the same shape. The GCP collector does not. const AuditEventsProperty = "audit_events" // AuditEvent is one management event attached to existing graph nodes.