diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 615ced1..f7334a1 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -137,7 +137,11 @@ jobs: git diff --exit-code -- gen packages/connect/src/gen test -z "$(git status --porcelain -- gen packages/connect/src/gen)" - - run: go test ./... + # The bootstrap and ingestion-worker packages share this job's migrated + # Postgres database. Run their package binaries serially so a worker test + # cannot claim another package's queued fixture. + - name: Run database-backed Go tests + run: go test -p 1 ./... - name: Run bounded Go worker smokes run: npm run smoke:workers:go diff --git a/CHANGELOG.md b/CHANGELOG.md index 072cb42..fa1982b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,12 +4,24 @@ All notable changes to Aperio are recorded here. Release entries are tied to a s ## [Unreleased] +## [0.1.0] - 2026-08-21 + +- Added tenant-scoped, hashed API tokens with read, write, and admin scopes, + expiry, revocation, last-used state, and audit records. +- Added `aperioctl` commands for health checks, findings, connectors, sync, + SIEM destinations, and API-token lifecycle operations. +- Added durable connector-sync and rule-run receipts plus authenticated + operator health for connector freshness, ingestion queues, SIEM delivery, + and recent rule execution. +- Added a Prometheus endpoint that is disabled until a dedicated scrape token + is configured and does not expose tenant or resource labels. +- Added a strict CEL/YAML detection engine with versioned rules for GitHub + public repositories, Slack MFA and external shared channels, and Google + Workspace external sharing. +- Added tenant rule disablement, severity overrides, scoped auto-resolution, + an in-memory backtest API, and an explicit connector support matrix. - Added a locally-owned review preflight that checks workflow action pinning and reports required validation commands. - Removed vendor-backed review and code-writing workflows from the public repository. - Added a production Compose bundle with Postgres, NATS, API, web, ingestion, SIEM, and migration services. - Added release, upgrade, backup, security, and contributor documentation. - Updated the Node dependency lock and Go modules to patched release lines. - -## [0.1.0] - Not released - -Initial public product baseline. The first release date will be recorded when the signed image and source tag are published. diff --git a/Makefile b/Makefile index 38490ce..b90b73f 100644 --- a/Makefile +++ b/Makefile @@ -289,7 +289,7 @@ test-go: ## Run Go unit tests .PHONY: test-go-db test-go-db: require-env ## Run Go tests including DB-backed routes (needs Postgres) @$(MAKE) --no-print-directory db-up migrate - @$(LOAD_ENV) APERIO_TEST_DATABASE_URL="$$(node $(DEV_CONFIG) go-database-url)" go test ./... + @$(LOAD_ENV) APERIO_TEST_DATABASE_URL="$$(node $(DEV_CONFIG) go-database-url)" go test -p 1 ./... .PHONY: test-api test-api: require-env ## Run the TypeScript/node test suite diff --git a/docs/detection-support-matrix.md b/docs/detection-support-matrix.md new file mode 100644 index 0000000..2e59b52 --- /dev/null +++ b/docs/detection-support-matrix.md @@ -0,0 +1,52 @@ +# Detection support matrix + +This matrix describes the payload fields the current Go ingestion worker can +actually evaluate. A connector is not treated as rule-supported merely because +credentials can be stored; the provider must enqueue one of the listed event +types with the listed fields. + +| Provider | Rule | Version | Event types | Required payload fields | Auto-resolution input | Default | Notes | +| --- | --- | --- | --- | --- | --- | --- | --- | +| GitHub | `github.public_repository_created` | `1.0.0` | `PUBLIC_REPOSITORY_CREATED`, `REPOSITORY_PUBLICIZED` (aliases accepted) | `repository.full_name`; `repository.visibility` or `repository.private` | `REPOSITORY_PRIVATE`, `REPOSITORY_PRIVATEIZED`, or `REPOSITORY_VISIBILITY_CHANGED` with private visibility | On | Declarative YAML pack; the finding target is the repository and the dedupe subject is the same repository. | +| GitHub | `github.branch_protection_disabled` | `catalog` | `BRANCH_PROTECTION_DISABLED`, `BRANCH_PROTECTION_RULE_DELETED`, `BRANCH_PROTECTION_RULE_UPDATED` | repository name; branch/ref/rule pattern when available | Not currently automatic | On | Hardcoded compatibility rule. Updated rules fire only when the payload indicates weakened settings. | +| GitHub | `github.oauth_app_installed` | `catalog` | `OAUTH_APP_INSTALLED`, `GITHUB_APP_INSTALLED`, `ORG_OAUTH_APP_ACCESS_APPROVED` | app name or ID; scopes/permissions when available | Not currently automatic | On | Payloads with no scope list are retained for review; known low-risk scoped installs are skipped. | +| GitHub | `github.deploy_key_added` | `1.0.0` catalog | `DEPLOY_KEY_ADDED`, `DEPLOY_KEY_CREATED` (aliases accepted) | repository name; key title/name/ID; `key.write_enabled` when available | Not currently automatic | Off | Existing connector catalog exposes this opt-in check. Write-enabled keys escalate to HIGH; missing write metadata remains MEDIUM. | +| Slack | `slack.mfa_disabled` | `1.0.0` | `MFA_DISABLED`, `TWO_FACTOR_AUTH_DISABLED` (aliases accepted) | `user.email` or `user.id` | `MFA_ENABLED` or `TWO_FACTOR_AUTH_ENABLED` with the same user field | On | Declarative YAML pack. The clean event is never treated as a disablement. | +| Slack | `slack.external_shared_channel_created` | `1.0.0` | `EXTERNAL_SHARED_CHANNEL_CREATED`, `SHARED_CHANNEL_INVITE_ACCEPTED` | channel name/ID; external organization/team name | Not currently automatic | On | Declarative YAML pack; the finding is emitted only when both channel and external-organization identity are present. | +| Slack | `slack.workspace_invite_link_enabled` | `catalog` | `WORKSPACE_INVITE_LINK_ENABLED`, `INVITE_LINK_CREATED` | workspace/team name when available | Not currently automatic | On | Hardcoded compatibility rule. | +| Slack | `slack.app_installed` | `catalog` | `APP_INSTALLED`, `APP_APPROVED`, `APP_SCOPES_APPROVED` | app name/ID; scopes when available | Not currently automatic | Off | Existing rule escalates to HIGH when scopes include admin, file-history, or channel-history access; no claim is made when scope data is absent. | +| Google Workspace | `google_workspace.external_sharing_enabled` | `1.0.0` | `EXTERNAL_SHARING_ENABLED` | `parameters.visibility`; document title/id/type/owner when available | `EXTERNAL_SHARING_DISABLED`, `DRIVE_FILE_VISIBILITY_CHANGED`, or `DRIVE_FILE_PRIVATE` with private/domain/internal visibility | On | Declarative YAML pack supports both resource metadata and Reports API parameter shapes. | + +## Tenant overrides + +The worker already reads `integration_connections.disabled_checks`. A severity +override can be stored in the existing `disabled_check_metadata` JSON object +without a new table: + +```json +{ + "slack.mfa_disabled": { + "reason": "tenant risk policy", + "severity": "HIGH", + "expiresAt": "2026-12-31T00:00:00Z" + } +} +``` + +Only `CRITICAL`, `HIGH`, `MEDIUM`, `LOW`, and `INFO` are accepted. Expired or +malformed entries are ignored. Disabling a rule still uses the existing +`disabled_checks` array and its expiry behavior. + +## Known gaps + +- GitHub secret-scanning alerts, dependabot alerts, membership changes, and + webhook delivery are not advertised here because this checkout does not + currently normalize those provider payloads into supported ingestion event + types. +- Slack message/file export volume, guest lifecycle, and user deactivation are + not advertised because the available payload contract does not guarantee the + required fields. +- Rule efficacy rollups, persisted `rule_version` columns, community-pack + signatures, and stateful/correlation rules remain follow-up work. The + evaluator emits versioned drafts now so those consumers can be added without + changing rule semantics. diff --git a/go.mod b/go.mod index 76680ff..33bf454 100644 --- a/go.mod +++ b/go.mod @@ -3,6 +3,7 @@ module github.com/writer/aperio go 1.26.6 require ( + cel.dev/cel-go v0.32.0 connectrpc.com/connect v1.20.0 github.com/go-pdf/fpdf v0.9.0 github.com/jackc/pgx/v5 v5.10.0 @@ -10,16 +11,25 @@ require ( github.com/writer/cerebro/sdk/go/cerebroapi v0.0.0-20260617190440-784f5eee34f2 golang.org/x/crypto v0.55.0 google.golang.org/protobuf v1.36.12 + gopkg.in/yaml.v3 v3.0.1 ) require ( + cel.dev/expr v0.25.1 // indirect + github.com/antlr4-go/antlr/v4 v4.13.1 // indirect github.com/jackc/pgpassfile v1.0.0 // indirect github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect github.com/jackc/puddle/v2 v2.2.2 // indirect github.com/klauspost/compress v1.18.5 // indirect + github.com/kr/text v0.2.0 // indirect github.com/nats-io/nkeys v0.4.15 // indirect github.com/nats-io/nuid v1.0.1 // indirect + github.com/rogpeppe/go-internal v1.16.0 // indirect + go.yaml.in/yaml/v3 v3.0.4 // indirect + golang.org/x/exp v0.0.0-20240823005443-9b4947da3948 // indirect golang.org/x/sync v0.22.0 // indirect golang.org/x/sys v0.47.0 // indirect golang.org/x/text v0.41.0 // indirect + google.golang.org/genproto/googleapis/api v0.0.0-20240826202546-f6391c0de4c7 // indirect + google.golang.org/genproto/googleapis/rpc v0.0.0-20240826202546-f6391c0de4c7 // indirect ) diff --git a/go.sum b/go.sum index e570121..822b944 100644 --- a/go.sum +++ b/go.sum @@ -1,5 +1,12 @@ +cel.dev/cel-go v0.32.0 h1:irvpFKr5EuGPyxeME03ERh0rii1TX+BDAnB9eL3IvNk= +cel.dev/cel-go v0.32.0/go.mod h1:DnVip7tpJSsgZymwfT+m1tnEVy3ivAjSMXPx12YrMkU= +cel.dev/expr v0.25.1 h1:1KrZg61W6TWSxuNZ37Xy49ps13NUovb66QLprthtwi4= +cel.dev/expr v0.25.1/go.mod h1:hrXvqGP6G6gyx8UAHSHJ5RGk//1Oj5nXQ2NI02Nrsg4= connectrpc.com/connect v1.20.0 h1:6TNDAB+WeNd2uolWNlYczB5E0KNNaVMNUEx8JEUsPmQ= connectrpc.com/connect v1.20.0/go.mod h1:A2ygJrukXwWy32vkCAAHNVguZrqZ+jeZ9rGRnGR4dN4= +github.com/antlr4-go/antlr/v4 v4.13.1 h1:SqQKkuVZ+zWkMMNkjy5FZe5mr5WURWnlpmOuzYWrPrQ= +github.com/antlr4-go/antlr/v4 v4.13.1/go.mod h1:GKmUxMtwp6ZgGwZSva4eWPC5mS6vUAmOABFgjdkM7Nw= +github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -17,6 +24,10 @@ github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= github.com/klauspost/compress v1.18.5 h1:/h1gH5Ce+VWNLSWqPzOVn6XBO+vJbCNGvjoaGBFW2IE= github.com/klauspost/compress v1.18.5/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/kr/pretty v0.3.0 h1:WgNl7dwNpEZ6jJ9k1snq4pZsg7DOEN8hP9Xw0Tsjwk0= +github.com/kr/pretty v0.3.0/go.mod h1:640gp4NfQd8pI5XOwp5fnNeVWj67G7CFk/SaSQn7NBk= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/nats-io/nats.go v1.53.1 h1:Otsq3uLc/kLdjmkNHkXH0jBqwUquwdKFoe3fq6/3/Xo= github.com/nats-io/nats.go v1.53.1/go.mod h1:26HypzazeOkyO3/mqd1zZd53STJN0EjCYF9Uy2ZOBno= github.com/nats-io/nkeys v0.4.15 h1:JACV5jRVO9V856KOapQ7x+EY8Jo3qw1vJt/9Jpwzkk4= @@ -25,6 +36,8 @@ github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw= github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/rogpeppe/go-internal v1.16.0 h1:O9DK+vNMDVGLr2BeZqmpLeMjiMNkuXfcqntWbZV6S5g= +github.com/rogpeppe/go-internal v1.16.0/go.mod h1:DrUVZyrJU+txYW5/1kwtXQSMFio52ZOxX7yM1VHvnxs= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= @@ -32,17 +45,27 @@ github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/writer/cerebro/sdk/go/cerebroapi v0.0.0-20260617190440-784f5eee34f2 h1:+Jsrg6W9h3uUA8wt3N2hXtC3ZbiVVbvrPdIR01ClEXQ= github.com/writer/cerebro/sdk/go/cerebroapi v0.0.0-20260617190440-784f5eee34f2/go.mod h1:dHVu/CuhHRejAdS2yHB7oBCnfeT9ZVD4yxDHdPPzmhw= +go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= +go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M= golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis= +golang.org/x/exp v0.0.0-20240823005443-9b4947da3948 h1:kx6Ds3MlpiUHKj7syVnbp57++8WpuKPcR5yjLBjvLEA= +golang.org/x/exp v0.0.0-20240823005443-9b4947da3948/go.mod h1:akd2r19cwCdwSwWeIdzYQGa/EZZyqcOdwWiwj5L5eKQ= golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek= golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8= golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M= +google.golang.org/genproto/googleapis/api v0.0.0-20240826202546-f6391c0de4c7 h1:YcyjlL1PRr2Q17/I0dPk2JmYS5CDXfcdb2Z3YRioEbw= +google.golang.org/genproto/googleapis/api v0.0.0-20240826202546-f6391c0de4c7/go.mod h1:OCdP9MfskevB/rbYvHTsXTtKC+3bHWajPdoKgjcYkfo= +google.golang.org/genproto/googleapis/rpc v0.0.0-20240826202546-f6391c0de4c7 h1:2035KHhUv+EpyB+hWgJnaWKJOdX1E95w2S8Rr4uWKTs= +google.golang.org/genproto/googleapis/rpc v0.0.0-20240826202546-f6391c0de4c7/go.mod h1:UqMtugtsSgubUsoxbuAoiCXvqvErP7Gf0so0mK9tHxU= google.golang.org/protobuf v1.36.12 h1:pJOKDDOyeXErUroCihFAd5LQuwXBSpVnKGrj5o/fwxc= google.golang.org/protobuf v1.36.12/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/internal/detection/README.md b/internal/detection/README.md new file mode 100644 index 0000000..b352e63 --- /dev/null +++ b/internal/detection/README.md @@ -0,0 +1,36 @@ +# Declarative detection rules + +The built-in pack is embedded from `rules/*.yaml` and compiled once when the +ingestion worker starts evaluating a job. `LoadPack` can validate an external +directory for a bounded backtest or a future rule-management API. + +Each rule has a stable `id` and semantic `version`. `source` is an allow-list +for provider/event types; `when.expression` is CEL evaluated with only these +bindings: + +- `event`: immutable event metadata and provider payload +- `org.config`: non-secret tenant configuration supplied by the caller +- `org.allowlists`: non-secret tenant allowlists supplied by the caller +- `now`: RFC3339 timestamp for deterministic evaluation + +Finding strings and evidence use a logic-free template subset. Dotted paths, +`first_nonempty(path, path, ...)`, and +`external_recipient(path, path, owner_path)` are supported. Templates cannot +call CEL, access the filesystem, or execute Go code. Unknown YAML fields, +invalid semantic versions, invalid severities, oversized expressions, and +duplicate rule IDs are rejected before compilation. A pack has one active +semantic version per rule ID; replacing a rule version is an explicit pack +rollout rather than an in-place overlap. + +`auto_resolve_when` emits a resolution draft only. Persistence must scope the +state transition by `organization_id`, `integration_id`, rule ID, and rendered +dedupe target; the evaluator never mutates storage. + +The evaluator exposes version-aware dedupe material for backtests. The current +worker preserves the existing persisted ID-plus-target hash during migration +and records the rule version in finding evidence; changing that database key +requires a separate collision/backfill rollout. + +The JSON schema in `rule.schema.json` is a portable authoring contract. The Go +loader additionally compiles every expression, which catches CEL errors before +activation. diff --git a/internal/detection/engine.go b/internal/detection/engine.go new file mode 100644 index 0000000..120cce9 --- /dev/null +++ b/internal/detection/engine.go @@ -0,0 +1,549 @@ +package detection + +import ( + "crypto/sha256" + "encoding/hex" + "errors" + "fmt" + "regexp" + "sort" + "strings" + "time" + + "cel.dev/cel-go/cel" + "cel.dev/cel-go/common/types" +) + +const maxBacktestSamples = 20 + +var templatePattern = regexp.MustCompile(`\{\{\s*([^{}]+?)\s*\}\}`) + +type compiledRule struct { + rule Rule + when cel.Program + autoResolve cel.Program +} + +// Engine compiles a rule pack once and evaluates it without I/O. CEL's +// standard environment is intentionally limited to three dynamic bindings; +// rules cannot call Go functions or access process state. +type Engine struct { + rules []compiledRule +} + +// NewEngine validates and compiles rules. Compilation is the trust boundary: +// callers should refuse to activate a pack when this returns an error. +func NewEngine(rules []Rule) (*Engine, error) { + rules, err := ValidateRules(rules) + if err != nil { + return nil, err + } + env, err := cel.NewEnv( + cel.Variable("event", cel.DynType), + cel.Variable("org", cel.DynType), + cel.Variable("now", cel.StringType), + ) + if err != nil { + return nil, fmt.Errorf("create CEL environment: %w", err) + } + compiled := make([]compiledRule, 0, len(rules)) + for _, rule := range rules { + when, err := compileExpression(env, rule.ID, rule.When.Expression) + if err != nil { + return nil, err + } + var autoResolve cel.Program + if rule.AutoResolveWhen != nil { + autoResolve, err = compileExpression(env, rule.ID+" auto_resolve_when", rule.AutoResolveWhen.Expression) + if err != nil { + return nil, err + } + } + compiled = append(compiled, compiledRule{rule: rule, when: when, autoResolve: autoResolve}) + } + return &Engine{rules: compiled}, nil +} + +func compileExpression(env *cel.Env, name, expression string) (cel.Program, error) { + ast, issues := env.Compile(expression) + if issues != nil && issues.Err() != nil { + return nil, fmt.Errorf("compile rule %q expression: %w", name, issues.Err()) + } + program, err := env.Program(ast) + if err != nil { + return nil, fmt.Errorf("build rule %q program: %w", name, err) + } + return program, nil +} + +// Rules returns a stable copy of the compiled pack metadata. +func (e *Engine) Rules() []Rule { + if e == nil { + return nil + } + out := make([]Rule, 0, len(e.rules)) + for _, item := range e.rules { + out = append(out, item.rule) + } + return out +} + +// Evaluate produces one draft per matching rule. Event type filtering and +// tenant overrides happen before CEL evaluation. A runtime expression error +// is returned rather than treated as a match; callers can fail closed and +// preserve an error/metrics record for operators. +func (e *Engine) Evaluate(event Event, org OrgContext, overrides Overrides) ([]FindingDraft, error) { + if e == nil { + return nil, errors.New("detection engine is nil") + } + if event.Payload == nil { + event.Payload = map[string]any{} + } + now := event.OccurredAt.UTC() + if now.IsZero() { + now = time.Now().UTC() + } + out := make([]FindingDraft, 0, 2) + for _, item := range e.rules { + rule := item.rule + if strings.TrimSpace(event.Provider) == "" || !strings.EqualFold(strings.TrimSpace(event.Provider), strings.TrimSpace(rule.Source.Provider)) { + continue + } + if !matchesEventType(rule.Source.EventTypes, event.EventType) { + continue + } + if overrides.Disabled != nil && overrides.Disabled[rule.ID] { + continue + } + matched, err := evaluateBoolean(item.when, event, org, now) + if err != nil { + return nil, fmt.Errorf("evaluate rule %s: %w", rule.ID, err) + } + if !matched { + continue + } + target, err := renderTemplate(rule.Dedupe.TargetTemplate, event, org, now, true) + if err != nil { + return nil, fmt.Errorf("render rule %s dedupe target: %w", rule.ID, err) + } + findingTarget := target + if strings.TrimSpace(rule.Finding.TargetTemplate) != "" { + findingTarget, err = renderTemplate(rule.Finding.TargetTemplate, event, org, now, true) + if err != nil { + return nil, fmt.Errorf("render rule %s finding target: %w", rule.ID, err) + } + } + title, err := renderTemplate(rule.Finding.Title, event, org, now, false) + if err != nil { + return nil, fmt.Errorf("render rule %s title: %w", rule.ID, err) + } + description, err := renderTemplate(rule.Finding.Description, event, org, now, false) + if err != nil { + return nil, fmt.Errorf("render rule %s description: %w", rule.ID, err) + } + steps := make([]string, 0, len(rule.Finding.RemediationSteps)) + for _, step := range rule.Finding.RemediationSteps { + rendered, err := renderTemplate(step, event, org, now, false) + if err != nil { + return nil, fmt.Errorf("render rule %s remediation step: %w", rule.ID, err) + } + steps = append(steps, rendered) + } + evidence := make(map[string]any, len(rule.Finding.Evidence)) + for key, template := range rule.Finding.Evidence { + rendered, err := renderTemplate(template, event, org, now, false) + if err != nil { + return nil, fmt.Errorf("render rule %s evidence %q: %w", rule.ID, key, err) + } + evidence[key] = rendered + } + severity := strings.ToUpper(strings.TrimSpace(rule.Severity)) + if override := strings.TrimSpace(overrides.SeverityOverrides[rule.ID]); override != "" { + severity = strings.ToUpper(override) + if !validSeverity(severity) { + return nil, fmt.Errorf("severity override for %s is invalid: %q", rule.ID, override) + } + } + out = append(out, FindingDraft{ + RuleID: rule.ID, + RuleVersion: rule.Version, + Title: title, + Description: description, + Severity: severity, + RiskScore: clampRiskScore(severity, rule.RiskScore), + Tags: append([]string(nil), rule.Tags...), + RemediationSteps: steps, + Target: findingTarget, + DedupeTarget: target, + Evidence: evidence, + }) + } + return out, nil +} + +// AutoResolve evaluates clean-event predicates for open finding state. It +// does not mutate storage; the ingestion worker must perform the scoped +// transition using the returned rule ID and target. +func (e *Engine) AutoResolve(event Event, org OrgContext, overrides Overrides) ([]ResolutionDraft, error) { + if e == nil { + return nil, errors.New("detection engine is nil") + } + now := event.OccurredAt.UTC() + if now.IsZero() { + now = time.Now().UTC() + } + out := make([]ResolutionDraft, 0, 1) + for _, item := range e.rules { + rule := item.rule + if item.autoResolve == nil || strings.TrimSpace(event.Provider) == "" || !strings.EqualFold(strings.TrimSpace(event.Provider), strings.TrimSpace(rule.Source.Provider)) { + continue + } + if !matchesEventType(rule.Source.EventTypes, event.EventType) { + continue + } + if overrides.Disabled != nil && overrides.Disabled[rule.ID] { + continue + } + matched, err := evaluateBoolean(item.autoResolve, event, org, now) + if err != nil { + return nil, fmt.Errorf("evaluate auto-resolve rule %s: %w", rule.ID, err) + } + if !matched { + continue + } + target, err := renderTemplate(rule.Dedupe.TargetTemplate, event, org, now, true) + if err != nil { + // A clean event without a stable subject must never resolve an + // unrelated finding. Fail closed for this rule. + continue + } + out = append(out, ResolutionDraft{ + RuleID: rule.ID, + RuleVersion: rule.Version, + DedupeTarget: target, + Evidence: map[string]any{ + "ruleId": rule.ID, + "ruleVersion": rule.Version, + "subject": target, + "resolution": "auto_resolve_when", + "eventType": event.EventType, + }, + }) + } + return out, nil +} + +// Backtest replays events through one rule. It is intentionally in-memory so +// callers can use fixture events in CI or a bounded database query without +// granting the evaluator database access. +func (e *Engine) Backtest(ruleID string, events []Event, org OrgContext, overrides Overrides) (BacktestReport, error) { + item, ok := e.findRule(ruleID) + if !ok { + return BacktestReport{}, fmt.Errorf("rule %q is not loaded", ruleID) + } + report := BacktestReport{RuleID: item.rule.ID, RuleVersion: item.rule.Version, Events: len(events)} + for _, event := range events { + if !strings.EqualFold(strings.TrimSpace(event.Provider), strings.TrimSpace(item.rule.Source.Provider)) || !matchesEventType(item.rule.Source.EventTypes, event.EventType) { + continue + } + report.Candidates++ + findings, err := e.Evaluate(event, org, overrides) + if err != nil { + return report, err + } + for _, finding := range findings { + if finding.RuleID != item.rule.ID { + continue + } + report.Matches++ + if len(report.MatchSamples) < maxBacktestSamples { + report.MatchSamples = append(report.MatchSamples, BacktestMatch{ + EventType: event.EventType, + OccurredAt: event.OccurredAt.UTC().Format(time.RFC3339Nano), + Target: finding.Target, + DedupeTarget: finding.DedupeTarget, + Severity: finding.Severity, + }) + } + } + resolutions, err := e.AutoResolve(event, org, overrides) + if err != nil { + return report, err + } + for _, resolution := range resolutions { + if resolution.RuleID != item.rule.ID { + continue + } + report.Resolutions++ + if len(report.MatchSamples) < maxBacktestSamples { + report.MatchSamples = append(report.MatchSamples, BacktestMatch{ + EventType: event.EventType, + OccurredAt: event.OccurredAt.UTC().Format(time.RFC3339Nano), + DedupeTarget: resolution.DedupeTarget, + Resolution: true, + }) + } + } + } + return report, nil +} + +func (e *Engine) findRule(id string) (compiledRule, bool) { + for _, item := range e.rules { + if item.rule.ID == id { + return item, true + } + } + return compiledRule{}, false +} + +func evaluateBoolean(program cel.Program, event Event, org OrgContext, now time.Time) (bool, error) { + activation := map[string]any{ + "event": map[string]any{ + "organization_id": event.OrganizationID, + "integration_id": event.IntegrationID, + "provider": event.Provider, + "event_type": event.EventType, + "event_type_normalized": normalizeEventType(event.EventType), + "source": event.Source, + "actor": event.Actor, + "occurred_at": event.OccurredAt.UTC().Format(time.RFC3339Nano), + "payload": event.Payload, + }, + "org": map[string]any{ + "config": org.Config, + "allowlists": org.Allowlists, + }, + "now": now.UTC().Format(time.RFC3339Nano), + } + value, _, err := program.Eval(activation) + if err != nil { + return false, err + } + boolean, ok := value.(types.Bool) + if !ok { + return false, fmt.Errorf("expression returned %s, want bool", value.Type().TypeName()) + } + return bool(boolean), nil +} + +func renderTemplate(template string, event Event, org OrgContext, now time.Time, required bool) (string, error) { + trimmed := strings.TrimSpace(template) + if trimmed == "" { + if required { + return "", errors.New("template is empty") + } + return "", nil + } + activation := map[string]any{ + "event": map[string]any{ + "organization_id": event.OrganizationID, + "integration_id": event.IntegrationID, + "provider": event.Provider, + "event_type": event.EventType, + "event_type_normalized": normalizeEventType(event.EventType), + "source": event.Source, + "actor": event.Actor, + "occurred_at": event.OccurredAt.UTC().Format(time.RFC3339Nano), + "payload": event.Payload, + }, + "org": map[string]any{ + "config": org.Config, + "allowlists": org.Allowlists, + }, + "now": now.UTC().Format(time.RFC3339Nano), + } + missing := "" + result := templatePattern.ReplaceAllStringFunc(template, func(match string) string { + parts := templatePattern.FindStringSubmatch(match) + if len(parts) != 2 { + missing = match + return "" + } + value, ok := resolveTemplateValue(strings.TrimSpace(parts[1]), activation) + if !ok || value == nil { + missing = parts[1] + return "" + } + return fmt.Sprint(value) + }) + if missing != "" && required { + return "", fmt.Errorf("template path %q is missing", missing) + } + if required && strings.TrimSpace(result) == "" { + return "", errors.New("template rendered an empty value") + } + return result, nil +} + +// resolveTemplateValue supports a deliberately tiny, logic-free template +// vocabulary. Plain dotted paths cover most rules. first_nonempty(...) is +// useful when a provider emits the same field under two documented payload +// shapes. external_recipient(...) selects the first email outside the owner +// domain and handles either a scalar or a JSON array. No loops, conditionals, +// functions, or arbitrary CEL are accepted here. +func resolveTemplateValue(expression string, activation map[string]any) (any, bool) { + if strings.HasPrefix(expression, "first_nonempty(") && strings.HasSuffix(expression, ")") { + args := splitTemplateArgs(strings.TrimSuffix(strings.TrimPrefix(expression, "first_nonempty("), ")")) + for _, arg := range args { + value, ok := resolvePath(activation, strings.TrimSpace(arg)) + if ok && !isEmptyTemplateValue(value) { + return value, true + } + } + return nil, false + } + if strings.HasPrefix(expression, "external_recipient(") && strings.HasSuffix(expression, ")") { + args := splitTemplateArgs(strings.TrimSuffix(strings.TrimPrefix(expression, "external_recipient("), ")")) + ownerDomain := "" + if len(args) > 0 { + if value, ok := resolvePath(activation, strings.TrimSpace(args[len(args)-1])); ok { + ownerDomain = strings.ToLower(strings.TrimSpace(fmt.Sprint(value))) + } + if index := strings.LastIndex(ownerDomain, "@"); index >= 0 { + ownerDomain = ownerDomain[index+1:] + } + } + for _, arg := range args[:max(0, len(args)-1)] { + value, ok := resolvePath(activation, strings.TrimSpace(arg)) + if !ok { + continue + } + for _, candidate := range flattenTemplateValues(value) { + email := strings.TrimSpace(fmt.Sprint(candidate)) + if email == "" || !strings.Contains(email, "@") { + continue + } + at := strings.LastIndex(email, "@") + if ownerDomain == "" || !strings.EqualFold(email[at+1:], ownerDomain) { + return email, true + } + } + } + return nil, false + } + return resolvePath(activation, expression) +} + +func splitTemplateArgs(value string) []string { + parts := strings.Split(value, ",") + out := make([]string, 0, len(parts)) + for _, part := range parts { + if strings.TrimSpace(part) != "" { + out = append(out, strings.TrimSpace(part)) + } + } + return out +} + +func flattenTemplateValues(value any) []any { + switch typed := value.(type) { + case []any: + out := make([]any, 0, len(typed)) + for _, item := range typed { + out = append(out, flattenTemplateValues(item)...) + } + return out + case []string: + out := make([]any, 0, len(typed)) + for _, item := range typed { + out = append(out, item) + } + return out + default: + return []any{value} + } +} + +func isEmptyTemplateValue(value any) bool { + if value == nil { + return true + } + if text, ok := value.(string); ok { + return strings.TrimSpace(text) == "" + } + return false +} + +func max(left, right int) int { + if left > right { + return left + } + return right +} + +func resolvePath(root map[string]any, path string) (any, bool) { + var current any = root + for _, part := range strings.Split(path, ".") { + object, ok := current.(map[string]any) + if !ok { + return nil, false + } + current, ok = object[part] + if !ok { + return nil, false + } + } + return current, true +} + +func matchesEventType(allowed []string, actual string) bool { + normalized := normalizeEventType(actual) + for _, candidate := range allowed { + if normalized == normalizeEventType(candidate) { + return true + } + } + return false +} + +func clampRiskScore(severity string, score int) int { + if score == 0 { + score = map[string]int{"CRITICAL": 90, "HIGH": 75, "MEDIUM": 55, "LOW": 30, "INFO": 10}[severity] + } + floor, ceiling := 1, 100 + switch severity { + case "CRITICAL": + floor = 90 + case "HIGH": + floor, ceiling = 60, 89 + case "MEDIUM": + floor, ceiling = 40, 74 + case "LOW": + floor, ceiling = 20, 54 + case "INFO": + floor, ceiling = 1, 29 + } + if score < floor { + return floor + } + if score > ceiling { + return ceiling + } + return score +} + +// DedupeMaterial is provided for callers that need to inspect the exact +// versioned material without duplicating the hash implementation. +func DedupeMaterial(rule Rule, target string) string { + return rule.ID + "@" + rule.Version + ":" + target +} + +// VersionedDedupeHash returns a stable content hash for external backtest +// reports and future version-aware persistence. The current worker keeps the +// legacy ID-plus-target persisted key so existing findings do not fork during +// this migration, and stores RuleVersion in finding provenance. +func VersionedDedupeHash(rule Rule, target string) string { + digest := sha256.Sum256([]byte(DedupeMaterial(rule, target))) + return hex.EncodeToString(digest[:]) +} + +// SortedRuleIDs is useful for support-matrix and diagnostics output. +func SortedRuleIDs(rules []Rule) []string { + out := make([]string, 0, len(rules)) + for _, rule := range rules { + out = append(out, rule.ID) + } + sort.Strings(out) + return out +} diff --git a/internal/detection/engine_test.go b/internal/detection/engine_test.go new file mode 100644 index 0000000..403c981 --- /dev/null +++ b/internal/detection/engine_test.go @@ -0,0 +1,198 @@ +package detection + +import ( + "encoding/json" + "os" + "path/filepath" + "reflect" + "testing" + "time" +) + +func TestBuiltinPackCompilesAndEvaluatesMigratedRules(t *testing.T) { + rules, err := LoadEmbeddedPack(BuiltinFS, "rules/*.yaml") + if err != nil { + t.Fatalf("load built-in rules: %v", err) + } + if got, want := len(rules), 4; got != want { + t.Fatalf("built-in rules = %d, want %d", got, want) + } + engine, err := NewEngine(rules) + if err != nil { + t.Fatalf("compile built-in rules: %v", err) + } + + cases := []struct { + name string + event Event + wantRule string + wantTitle string + wantTarget string + wantSubject string + }{ + { + name: "github public repository", + event: Event{ + Provider: "GITHUB", EventType: "repository.publicized", OccurredAt: testTime(), + Payload: map[string]any{"repository": map[string]any{"full_name": "writer/aperio", "private": false, "visibility": "public"}}, + }, + wantRule: "github.public_repository_created", wantTitle: "Public GitHub repository created", wantTarget: "writer/aperio", wantSubject: "writer/aperio", + }, + { + name: "slack mfa", + event: Event{ + Provider: "SLACK", EventType: "mfa.disabled", Actor: "admin@example.com", OccurredAt: testTime(), + Payload: map[string]any{"user": map[string]any{"email": "user@example.com", "id": "U123"}}, + }, + wantRule: "slack.mfa_disabled", wantTitle: "Slack multi-factor authentication disabled", wantTarget: "user@example.com", wantSubject: "user@example.com", + }, + { + name: "slack external shared channel", + event: Event{ + Provider: "SLACK", EventType: "EXTERNAL_SHARED_CHANNEL_CREATED", Actor: "admin@example.com", OccurredAt: testTime(), + Payload: map[string]any{ + "channel": map[string]any{"name": "customer-data"}, + "external_organization": map[string]any{"name": "Partner Co"}, + }, + }, + wantRule: "slack.external_shared_channel_created", wantTitle: "Slack external shared channel created", wantTarget: "customer-data", wantSubject: "customer-data:Partner Co", + }, + { + name: "google external sharing", + event: Event{ + Provider: "GOOGLE_WORKSPACE", EventType: "EXTERNAL_SHARING_ENABLED", OccurredAt: testTime(), + Payload: map[string]any{ + "resource": map[string]any{"id": "drive_file_123", "name": "Board Deck"}, + "parameters": map[string]any{"doc_title": "Board Deck", "doc_id": "drive_file_123", "doc_type": "presentation", "owner": "owner@writer.com", "visibility": "public_on_the_web", "shared_with": "partner@external.example"}, + }, + }, + wantRule: "google_workspace.external_sharing_enabled", wantTitle: "Google Workspace external sharing enabled", wantTarget: "Board Deck", wantSubject: "drive_file_123", + }, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + findings, err := engine.Evaluate(tc.event, OrgContext{}, Overrides{}) + if err != nil { + t.Fatalf("evaluate: %v", err) + } + if len(findings) != 1 { + t.Fatalf("findings = %#v, want one", findings) + } + finding := findings[0] + if finding.RuleID != tc.wantRule || finding.Title != tc.wantTitle || finding.Target != tc.wantTarget || finding.DedupeTarget != tc.wantSubject { + t.Fatalf("finding = %#v", finding) + } + if finding.RuleVersion != "1.0.0" || finding.Evidence["ruleVersion"] != nil { + t.Fatalf("rule version should be a first-class field, not duplicated in evidence: %#v", finding) + } + }) + } +} + +func TestEngineOverridesAndAutoResolution(t *testing.T) { + rules, err := LoadEmbeddedPack(BuiltinFS, "rules/*.yaml") + if err != nil { + t.Fatal(err) + } + engine, err := NewEngine(rules) + if err != nil { + t.Fatal(err) + } + event := Event{ + Provider: "SLACK", EventType: "MFA_DISABLED", OccurredAt: testTime(), + Payload: map[string]any{"user": map[string]any{"email": "user@example.com"}}, + } + findings, err := engine.Evaluate(event, OrgContext{}, Overrides{SeverityOverrides: map[string]string{"slack.mfa_disabled": "LOW"}}) + if err != nil || len(findings) != 1 { + t.Fatalf("severity override evaluation = %#v, err=%v", findings, err) + } + if findings[0].Severity != "LOW" || findings[0].RiskScore >= 55 { + t.Fatalf("severity override did not lower finding: %#v", findings[0]) + } + disabled, err := engine.Evaluate(event, OrgContext{}, Overrides{Disabled: map[string]bool{"slack.mfa_disabled": true}}) + if err != nil || len(disabled) != 0 { + t.Fatalf("disabled evaluation = %#v, err=%v", disabled, err) + } + clean := event + clean.EventType = "two-factor auth enabled" + resolutions, err := engine.AutoResolve(clean, OrgContext{}, Overrides{}) + if err != nil { + t.Fatalf("auto resolve: %v", err) + } + want := []ResolutionDraft{{RuleID: "slack.mfa_disabled", RuleVersion: "1.0.0", DedupeTarget: "user@example.com", Evidence: map[string]any{ + "ruleId": "slack.mfa_disabled", "ruleVersion": "1.0.0", "subject": "user@example.com", "resolution": "auto_resolve_when", "eventType": "two-factor auth enabled", + }}} + if !reflect.DeepEqual(resolutions, want) { + t.Fatalf("resolutions = %#v, want %#v", resolutions, want) + } +} + +func TestVersionedDedupeMaterialChangesOnlyWithRuleVersion(t *testing.T) { + rule := Rule{ID: "example.rule", Version: "1.0.0"} + first := VersionedDedupeHash(rule, "subject-1") + rule.Version = "1.1.0" + second := VersionedDedupeHash(rule, "subject-1") + if first == second || DedupeMaterial(rule, "subject-1") != "example.rule@1.1.0:subject-1" { + t.Fatalf("versioned dedupe material did not change: first=%s second=%s", first, second) + } +} + +func TestBacktestCountsMatchesAndResolutions(t *testing.T) { + rules, err := LoadEmbeddedPack(BuiltinFS, "rules/*.yaml") + if err != nil { + t.Fatal(err) + } + engine, err := NewEngine(rules) + if err != nil { + t.Fatal(err) + } + events := []Event{ + {Provider: "GITHUB", EventType: "PUBLIC_REPOSITORY_CREATED", OccurredAt: testTime(), Payload: map[string]any{"repository": map[string]any{"full_name": "acme/api", "private": false, "visibility": "public"}}}, + {Provider: "GITHUB", EventType: "REPOSITORY_PRIVATE", OccurredAt: testTime().Add(time.Minute), Payload: map[string]any{"repository": map[string]any{"full_name": "acme/api", "private": true, "visibility": "private"}}}, + {Provider: "GITHUB", EventType: "repository.deleted", OccurredAt: testTime(), Payload: map[string]any{"repository": map[string]any{"full_name": "acme/api", "private": false, "visibility": "public"}}}, + } + report, err := engine.Backtest("github.public_repository_created", events, OrgContext{}, Overrides{}) + if err != nil { + t.Fatal(err) + } + if report.Candidates != 2 || report.Matches != 1 || report.Resolutions != 1 { + t.Fatalf("backtest report = %#v", report) + } +} + +func TestBacktestReplaysRepositoryFixture(t *testing.T) { + rules, err := LoadEmbeddedPack(BuiltinFS, "rules/*.yaml") + if err != nil { + t.Fatal(err) + } + engine, err := NewEngine(rules) + if err != nil { + t.Fatal(err) + } + raw, err := os.ReadFile(filepath.Join("..", "..", "tests", "fixtures", "worker-parity", "github-public-repository.json")) + if err != nil { + t.Fatal(err) + } + var fixture struct { + Positive fixtureEvent `json:"positive"` + Negative fixtureEvent `json:"negative"` + } + if err := json.Unmarshal(raw, &fixture); err != nil { + t.Fatal(err) + } + report, err := engine.Backtest("github.public_repository_created", []Event{fixture.Positive.Event, fixture.Negative.Event}, OrgContext{}, Overrides{}) + if err != nil { + t.Fatal(err) + } + if report.Candidates != 2 || report.Matches != 1 || report.Resolutions != 1 { + t.Fatalf("fixture backtest report = %#v", report) + } +} + +type fixtureEvent struct { + Event Event `json:"payload"` +} + +func testTime() time.Time { + return time.Date(2026, 6, 6, 0, 0, 0, 0, time.UTC) +} diff --git a/internal/detection/loader.go b/internal/detection/loader.go new file mode 100644 index 0000000..2ec77c5 --- /dev/null +++ b/internal/detection/loader.go @@ -0,0 +1,261 @@ +package detection + +import ( + "embed" + "errors" + "fmt" + "io/fs" + "os" + "path/filepath" + "regexp" + "sort" + "strings" + + "gopkg.in/yaml.v3" +) + +const ( + maxRuleBytes = 128 * 1024 + maxExpressionBytes = 16 * 1024 + maxTemplateBytes = 8 * 1024 +) + +var ( + semverPattern = regexp.MustCompile(`^(0|[1-9][0-9]*)\.(0|[1-9][0-9]*)\.(0|[1-9][0-9]*)$`) + ruleIDPattern = regexp.MustCompile(`^[a-z0-9][a-z0-9_.-]{1,159}$`) +) + +// BuiltinFS contains the reviewed P1 rules shipped with Aperio. Callers can +// use LoadPack on a filesystem for operator-supplied packs; production worker +// startup uses BuiltinEngine so it does not depend on the process cwd. +// +//go:embed rules/*.yaml +var BuiltinFS embed.FS + +// LoadPack reads and validates all YAML rule files directly in dir. Nested +// directories are not walked implicitly: packs are versioned directories and +// callers must choose the exact provider pack they intend to activate. +func LoadPack(dir string) ([]Rule, error) { + entries, err := os.ReadDir(dir) + if err != nil { + return nil, fmt.Errorf("read rule pack %q: %w", dir, err) + } + paths := make([]string, 0, len(entries)) + for _, entry := range entries { + if entry.IsDir() { + continue + } + ext := strings.ToLower(filepath.Ext(entry.Name())) + if ext == ".yaml" || ext == ".yml" { + paths = append(paths, filepath.Join(dir, entry.Name())) + } + } + sort.Strings(paths) + if len(paths) == 0 { + return nil, fmt.Errorf("rule pack %q contains no YAML files", dir) + } + rules := make([]Rule, 0, len(paths)) + for _, path := range paths { + rule, err := loadRuleFile(path, os.ReadFile) + if err != nil { + return nil, err + } + rules = append(rules, rule) + } + return ValidateRules(rules) +} + +// LoadEmbeddedPack loads YAML rules from an fs.FS, which makes the built-in +// rules easy to test without relying on the machine filesystem. +func LoadEmbeddedPack(source fs.FS, pattern string) ([]Rule, error) { + paths, err := fs.Glob(source, pattern) + if err != nil { + return nil, fmt.Errorf("glob embedded rule pack %q: %w", pattern, err) + } + sort.Strings(paths) + if len(paths) == 0 { + return nil, fmt.Errorf("embedded rule pack %q contains no YAML files", pattern) + } + rules := make([]Rule, 0, len(paths)) + for _, path := range paths { + rule, err := loadRuleFile(path, func(name string) ([]byte, error) { + return fs.ReadFile(source, name) + }) + if err != nil { + return nil, err + } + rules = append(rules, rule) + } + return ValidateRules(rules) +} + +func loadRuleFile(path string, readFile func(string) ([]byte, error)) (Rule, error) { + raw, err := readFile(path) + if err != nil { + return Rule{}, fmt.Errorf("read rule %q: %w", path, err) + } + if len(raw) == 0 || len(raw) > maxRuleBytes { + return Rule{}, fmt.Errorf("rule %q must be between 1 and %d bytes", path, maxRuleBytes) + } + var rule Rule + decoder := yaml.NewDecoder(strings.NewReader(string(raw))) + decoder.KnownFields(true) + if err := decoder.Decode(&rule); err != nil { + return Rule{}, fmt.Errorf("decode rule %q: %w", path, err) + } + if err := validateRule(rule); err != nil { + return Rule{}, fmt.Errorf("validate rule %q: %w", path, err) + } + return rule, nil +} + +// ValidateRules validates a complete pack and rejects duplicate rule IDs. +// There is one active version per ID in an Engine; a version change is a +// replacement that must be rolled out deliberately rather than two programs +// competing for the same persisted finding key. Returning a fresh slice also +// makes a caller's subsequent append unable to mutate the loader's registry +// accidentally. +func ValidateRules(rules []Rule) ([]Rule, error) { + if len(rules) == 0 { + return nil, errors.New("rule pack is empty") + } + seen := make(map[string]string, len(rules)) + out := make([]Rule, len(rules)) + copy(out, rules) + for index, rule := range out { + if err := validateRule(rule); err != nil { + return nil, fmt.Errorf("rule %d: %w", index, err) + } + normalized := rule + normalized.ID = strings.TrimSpace(rule.ID) + normalized.Version = strings.TrimSpace(rule.Version) + normalized.Name = strings.TrimSpace(rule.Name) + normalized.Severity = strings.ToUpper(strings.TrimSpace(rule.Severity)) + normalized.Source.Provider = strings.ToUpper(strings.TrimSpace(rule.Source.Provider)) + normalized.Source.EventTypes = append([]string(nil), rule.Source.EventTypes...) + for eventIndex, eventType := range normalized.Source.EventTypes { + normalized.Source.EventTypes[eventIndex] = strings.TrimSpace(eventType) + } + normalized.Tags = append([]string(nil), rule.Tags...) + normalized.Finding.RemediationSteps = append([]string(nil), rule.Finding.RemediationSteps...) + if rule.Finding.Evidence != nil { + normalized.Finding.Evidence = make(map[string]string, len(rule.Finding.Evidence)) + for key, value := range rule.Finding.Evidence { + normalized.Finding.Evidence[key] = value + } + } + out[index] = normalized + if previousVersion, duplicate := seen[normalized.ID]; duplicate { + return nil, fmt.Errorf("duplicate rule id %s (versions %s and %s cannot be active together)", normalized.ID, previousVersion, normalized.Version) + } + seen[normalized.ID] = normalized.Version + } + return out, nil +} + +func validateRule(rule Rule) error { + if !ruleIDPattern.MatchString(strings.TrimSpace(rule.ID)) { + return fmt.Errorf("id must match %s", ruleIDPattern.String()) + } + if !semverPattern.MatchString(strings.TrimSpace(rule.Version)) { + return errors.New("version must be semantic version X.Y.Z") + } + if strings.TrimSpace(rule.Name) == "" || len(rule.Name) > 220 { + return errors.New("name is required and must be at most 220 characters") + } + rule.Severity = strings.ToUpper(strings.TrimSpace(rule.Severity)) + if !validSeverity(rule.Severity) { + return fmt.Errorf("severity %q is not one of CRITICAL, HIGH, MEDIUM, LOW, INFO", rule.Severity) + } + if rule.RiskScore < 0 || rule.RiskScore > 100 { + return errors.New("risk_score must be between 0 and 100") + } + provider := strings.ToUpper(strings.TrimSpace(rule.Source.Provider)) + if provider == "" { + return errors.New("source.provider is required") + } + if len(rule.Source.EventTypes) == 0 { + return errors.New("source.event_types must contain at least one event type") + } + for index, eventType := range rule.Source.EventTypes { + if strings.TrimSpace(eventType) == "" { + return fmt.Errorf("source.event_types[%d] is empty", index) + } + } + if err := validateExpression("when.expression", rule.When.Expression); err != nil { + return err + } + if strings.TrimSpace(rule.Dedupe.TargetTemplate) == "" { + return errors.New("dedupe.target_template is required") + } + if len(rule.Dedupe.TargetTemplate) > maxTemplateBytes { + return fmt.Errorf("dedupe.target_template must be at most %d bytes", maxTemplateBytes) + } + if strings.TrimSpace(rule.Finding.Title) == "" || len(rule.Finding.Title) > 220 { + return errors.New("finding.title is required and must be at most 220 characters") + } + if len(rule.Finding.TargetTemplate) > maxTemplateBytes { + return fmt.Errorf("finding.target_template must be at most %d bytes", maxTemplateBytes) + } + if strings.TrimSpace(rule.Finding.Description) == "" { + return errors.New("finding.description is required") + } + for index, step := range rule.Finding.RemediationSteps { + if strings.TrimSpace(step) == "" { + return fmt.Errorf("finding.remediation_steps[%d] is empty", index) + } + if len(step) > maxTemplateBytes { + return fmt.Errorf("finding.remediation_steps[%d] is too long", index) + } + } + for key, value := range rule.Finding.Evidence { + if strings.TrimSpace(key) == "" { + return errors.New("finding.evidence contains an empty key") + } + if len(value) > maxTemplateBytes { + return fmt.Errorf("finding.evidence[%q] is too long", key) + } + } + if rule.AutoResolveWhen != nil { + if err := validateExpression("auto_resolve_when.expression", rule.AutoResolveWhen.Expression); err != nil { + return err + } + } + return nil +} + +func validateExpression(name, expression string) error { + if strings.TrimSpace(expression) == "" { + return fmt.Errorf("%s is required", name) + } + if len(expression) > maxExpressionBytes { + return fmt.Errorf("%s must be at most %d bytes", name, maxExpressionBytes) + } + return nil +} + +func validSeverity(value string) bool { + switch value { + case "CRITICAL", "HIGH", "MEDIUM", "LOW", "INFO": + return true + default: + return false + } +} + +func normalizeEventType(value string) string { + var builder strings.Builder + separator := false + for _, char := range strings.ToUpper(strings.TrimSpace(value)) { + if (char >= 'A' && char <= 'Z') || (char >= '0' && char <= '9') { + builder.WriteRune(char) + separator = false + continue + } + if !separator { + builder.WriteByte('_') + separator = true + } + } + return strings.Trim(builder.String(), "_") +} diff --git a/internal/detection/loader_test.go b/internal/detection/loader_test.go new file mode 100644 index 0000000..cb3846e --- /dev/null +++ b/internal/detection/loader_test.go @@ -0,0 +1,62 @@ +package detection + +import ( + "os" + "path/filepath" + "strings" + "testing" +) + +func TestLoadPackRejectsUnknownFieldsAndDuplicateVersions(t *testing.T) { + dir := t.TempDir() + valid := `id: example.rule +version: 1.0.0 +name: Example rule +severity: HIGH +risk_score: 75 +source: + provider: GITHUB + event_types: [EVENT] +when: + expression: event.payload.ok == true +dedupe: + target_template: "{{ event.payload.id }}" +finding: + title: Example + description: Example description +` + if err := os.WriteFile(filepath.Join(dir, "valid.yaml"), []byte(valid), 0o600); err != nil { + t.Fatal(err) + } + if _, err := LoadPack(dir); err != nil { + t.Fatalf("valid pack rejected: %v", err) + } + unknown := strings.Replace(valid, "name: Example rule", "name: Example rule\nunknown: true", 1) + if err := os.WriteFile(filepath.Join(dir, "unknown.yaml"), []byte(unknown), 0o600); err != nil { + t.Fatal(err) + } + if _, err := LoadPack(dir); err == nil || !strings.Contains(err.Error(), "field unknown not found") { + t.Fatalf("unknown field error = %v", err) + } +} + +func TestValidateRulesRejectsInvalidVersionsAndDuplicateIDs(t *testing.T) { + rules, err := LoadEmbeddedPack(BuiltinFS, "rules/*.yaml") + if err != nil { + t.Fatal(err) + } + duplicate := append(append([]Rule(nil), rules...), rules[0]) + if _, err := ValidateRules(duplicate); err == nil || !strings.Contains(err.Error(), "duplicate rule") { + t.Fatalf("duplicate error = %v", err) + } + otherVersion := rules[0] + otherVersion.Version = "2.0.0" + if _, err := ValidateRules([]Rule{rules[0], otherVersion}); err == nil || !strings.Contains(err.Error(), "cannot be active together") { + t.Fatalf("duplicate id across versions error = %v", err) + } + invalid := rules[0] + invalid.Version = "1.0" + if _, err := ValidateRules([]Rule{invalid}); err == nil || !strings.Contains(err.Error(), "semantic version") { + t.Fatalf("version error = %v", err) + } +} diff --git a/internal/detection/rule.schema.json b/internal/detection/rule.schema.json new file mode 100644 index 0000000..6ae28d6 --- /dev/null +++ b/internal/detection/rule.schema.json @@ -0,0 +1,55 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "https://aperio.example/schema/detection-rule-v1.json", + "title": "Aperio declarative detection rule", + "type": "object", + "additionalProperties": false, + "required": ["id", "version", "name", "severity", "risk_score", "source", "when", "dedupe", "finding"], + "properties": { + "id": {"type": "string", "pattern": "^[a-z0-9][a-z0-9_.-]{1,159}$"}, + "version": {"type": "string", "pattern": "^(0|[1-9][0-9]*)\\.(0|[1-9][0-9]*)\\.(0|[1-9][0-9]*)$"}, + "name": {"type": "string", "minLength": 1, "maxLength": 220}, + "severity": {"enum": ["CRITICAL", "HIGH", "MEDIUM", "LOW", "INFO"]}, + "risk_score": {"type": "integer", "minimum": 0, "maximum": 100}, + "tags": {"type": "array", "items": {"type": "string"}}, + "source": { + "type": "object", + "additionalProperties": false, + "required": ["provider", "event_types"], + "properties": { + "provider": {"type": "string", "minLength": 1}, + "event_types": {"type": "array", "minItems": 1, "items": {"type": "string", "minLength": 1}} + } + }, + "when": { + "type": "object", + "additionalProperties": false, + "required": ["expression"], + "properties": {"expression": {"type": "string", "minLength": 1, "maxLength": 16384}} + }, + "dedupe": { + "type": "object", + "additionalProperties": false, + "required": ["target_template"], + "properties": {"target_template": {"type": "string", "minLength": 1, "maxLength": 8192}} + }, + "finding": { + "type": "object", + "additionalProperties": false, + "required": ["title", "description"], + "properties": { + "target_template": {"type": "string", "maxLength": 8192}, + "title": {"type": "string", "minLength": 1, "maxLength": 220}, + "description": {"type": "string", "minLength": 1}, + "remediation_steps": {"type": "array", "items": {"type": "string", "minLength": 1}}, + "evidence": {"type": "object", "additionalProperties": {"type": "string"}} + } + }, + "auto_resolve_when": { + "type": "object", + "additionalProperties": false, + "required": ["expression"], + "properties": {"expression": {"type": "string", "minLength": 1, "maxLength": 16384}} + } + } +} diff --git a/internal/detection/rules/github_public_repository.yaml b/internal/detection/rules/github_public_repository.yaml new file mode 100644 index 0000000..b1af779 --- /dev/null +++ b/internal/detection/rules/github_public_repository.yaml @@ -0,0 +1,36 @@ +id: github.public_repository_created +version: 1.0.0 +name: Public GitHub repository created +severity: CRITICAL +risk_score: 95 +tags: + - data.public_exposure +source: + provider: GITHUB + event_types: + - PUBLIC_REPOSITORY_CREATED + - REPOSITORY_PUBLICIZED + - REPOSITORY_PRIVATE + - REPOSITORY_PRIVATEIZED + - REPOSITORY_VISIBILITY_CHANGED +when: + expression: >- + event.payload.repository.visibility == "public" || + event.payload.repository.private == false +dedupe: + target_template: "{{ event.payload.repository.full_name }}" +finding: + title: "Public GitHub repository created" + description: "A repository was created or changed to public visibility, which can expose source code, secrets, or customer data." + remediation_steps: + - "Confirm the repository is approved for public release." + - "Set repository visibility to private if public access is not explicitly authorized." + - "Run secret scanning and branch protection checks before allowing continued public access." + evidence: + repository: "{{ event.payload.repository.full_name }}" + subject: "{{ event.payload.repository.full_name }}" + visibility: "{{ event.payload.repository.visibility }}" +auto_resolve_when: + expression: >- + event.payload.repository.private == true || + event.payload.repository.visibility == "private" diff --git a/internal/detection/rules/google_external_sharing.yaml b/internal/detection/rules/google_external_sharing.yaml new file mode 100644 index 0000000..bb8c371 --- /dev/null +++ b/internal/detection/rules/google_external_sharing.yaml @@ -0,0 +1,47 @@ +id: google_workspace.external_sharing_enabled +version: 1.0.0 +name: Google Workspace external sharing enabled +severity: HIGH +risk_score: 75 +tags: + - data.external_share + - policy.weakened +source: + provider: GOOGLE_WORKSPACE + event_types: + - EXTERNAL_SHARING_ENABLED + - EXTERNAL_SHARING_DISABLED + - DRIVE_FILE_VISIBILITY_CHANGED + - DRIVE_FILE_PRIVATE +when: + expression: >- + event.payload.parameters.visibility == "public_on_the_web" || + event.payload.parameters.visibility == "shared_with_external_user" || + event.payload.parameters.visibility == "anyone_with_link" || + event.payload.parameters.visibility == "external" +dedupe: + target_template: "{{ first_nonempty(event.payload.resource.id, event.payload.parameters.doc_id) }}" +finding: + target_template: "{{ first_nonempty(event.payload.parameters.doc_title, event.payload.resource.name) }}" + title: "Google Workspace external sharing enabled" + description: "A Google Workspace resource was configured for external sharing, which may expose regulated or confidential data." + remediation_steps: + - "Confirm the external recipient and business purpose are approved." + - "Restrict the resource to the tenant domain or approved collaborators." + - "Review the resource audit trail for additional access while it was externally shared." + evidence: + fileName: "{{ first_nonempty(event.payload.parameters.doc_title, event.payload.resource.name) }}" + fileId: "{{ first_nonempty(event.payload.resource.id, event.payload.parameters.doc_id) }}" + fileType: "{{ event.payload.parameters.doc_type }}" + owner: "{{ event.payload.parameters.owner }}" + visibility: "{{ event.payload.parameters.visibility }}" + driveType: "Shared drive" + subject: "{{ first_nonempty(event.payload.resource.id, event.payload.parameters.doc_id) }}" + externalActor: "{{ external_recipient(event.payload.parameters.shared_with, event.payload.parameters.permission_change_grantee, event.payload.parameters.owner) }}" + docTitle: "{{ first_nonempty(event.payload.parameters.doc_title, event.payload.resource.name) }}" + docType: "{{ event.payload.parameters.doc_type }}" +auto_resolve_when: + expression: >- + event.payload.parameters.visibility == "private" || + event.payload.parameters.visibility == "domain" || + event.payload.parameters.visibility == "internal" diff --git a/internal/detection/rules/slack_external_shared_channel.yaml b/internal/detection/rules/slack_external_shared_channel.yaml new file mode 100644 index 0000000..2b3f48e --- /dev/null +++ b/internal/detection/rules/slack_external_shared_channel.yaml @@ -0,0 +1,31 @@ +id: slack.external_shared_channel_created +version: 1.0.0 +name: Slack external shared channel created +severity: HIGH +risk_score: 75 +tags: + - data.external_share +source: + provider: SLACK + event_types: + - EXTERNAL_SHARED_CHANNEL_CREATED + - SHARED_CHANNEL_INVITE_ACCEPTED +when: + expression: >- + (event.payload.channel.name != "" || event.payload.channel.id != "") && + (event.payload.external_organization.name != "" || event.payload.external_team.name != "" || event.payload.target_team.name != "") +dedupe: + target_template: "{{ first_nonempty(event.payload.channel.name, event.payload.channel.id, event.payload.conversation.name, event.payload.conversation.id) }}:{{ first_nonempty(event.payload.external_organization.name, event.payload.external_team.name, event.payload.target_team.name) }}" +finding: + target_template: "{{ first_nonempty(event.payload.channel.name, event.payload.channel.id, event.payload.conversation.name, event.payload.conversation.id) }}" + title: "Slack external shared channel created" + description: "A Slack channel was shared with an external organization, expanding conversation and file visibility outside the tenant." + remediation_steps: + - "Confirm the shared channel and external organization are approved." + - "Restrict channel membership and file sharing if the collaboration is still required." + - "Disconnect the shared channel if it was created unexpectedly." + evidence: + channel: "{{ first_nonempty(event.payload.channel.name, event.payload.channel.id, event.payload.conversation.name, event.payload.conversation.id) }}" + externalOrg: "{{ first_nonempty(event.payload.external_organization.name, event.payload.external_team.name, event.payload.target_team.name) }}" + actor: "{{ event.actor }}" + subject: "{{ first_nonempty(event.payload.channel.name, event.payload.channel.id, event.payload.conversation.name, event.payload.conversation.id) }}:{{ first_nonempty(event.payload.external_organization.name, event.payload.external_team.name, event.payload.target_team.name) }}" diff --git a/internal/detection/rules/slack_mfa_disabled.yaml b/internal/detection/rules/slack_mfa_disabled.yaml new file mode 100644 index 0000000..045bbfe --- /dev/null +++ b/internal/detection/rules/slack_mfa_disabled.yaml @@ -0,0 +1,35 @@ +id: slack.mfa_disabled +version: 1.0.0 +name: Slack multi-factor authentication disabled +severity: CRITICAL +risk_score: 90 +tags: + - auth.mfa_weakened +source: + provider: SLACK + event_types: + - MFA_DISABLED + - TWO_FACTOR_AUTH_DISABLED + - MFA_ENABLED + - TWO_FACTOR_AUTH_ENABLED +when: + expression: >- + event.event_type_normalized != "MFA_ENABLED" && + event.event_type_normalized != "TWO_FACTOR_AUTH_ENABLED" && + (event.payload.user.email != "" || event.payload.user.id != "") +dedupe: + target_template: "{{ first_nonempty(event.payload.user.email, event.payload.user.id, event.actor) }}" +finding: + title: "Slack multi-factor authentication disabled" + description: "A Slack user disabled MFA, increasing the likelihood of account takeover and lateral movement." + remediation_steps: + - "Re-enable MFA for the affected Slack user immediately." + - "Force a session reset for the affected account." + - "Review recent login history and connected Slack apps for suspicious activity." + evidence: + user: "{{ first_nonempty(event.payload.user.email, event.payload.user.id, event.actor) }}" + subject: "{{ first_nonempty(event.payload.user.email, event.payload.user.id, event.actor) }}" +auto_resolve_when: + expression: >- + event.event_type_normalized == "MFA_ENABLED" || + event.event_type_normalized == "TWO_FACTOR_AUTH_ENABLED" diff --git a/internal/detection/types.go b/internal/detection/types.go new file mode 100644 index 0000000..0658296 --- /dev/null +++ b/internal/detection/types.go @@ -0,0 +1,137 @@ +// Package detection evaluates versioned, declarative detection rules. +// +// The package deliberately owns no database or provider credentials. A rule +// receives an immutable event and organization context and returns finding +// drafts. Persistence, deduplication against SecurityFinding, and provider +// remediation remain in the ingestion worker. This keeps rule evaluation +// deterministic and makes the same engine usable by backtests and tests. +package detection + +import "time" + +// Rule is the YAML representation of one stateless detection rule. +// +// Rule IDs and versions are public contracts. A rule version changes when its +// predicate, finding semantics, or dedupe target changes; callers can retain +// the old version while evaluating a new one during a controlled migration. +type Rule struct { + ID string `yaml:"id" json:"id"` + Version string `yaml:"version" json:"version"` + Name string `yaml:"name" json:"name"` + Severity string `yaml:"severity" json:"severity"` + RiskScore int `yaml:"risk_score" json:"riskScore"` + Tags []string `yaml:"tags,omitempty" json:"tags,omitempty"` + Source Source `yaml:"source" json:"source"` + When Condition `yaml:"when" json:"when"` + Dedupe DedupeSpec `yaml:"dedupe" json:"dedupe"` + Finding FindingSpec `yaml:"finding" json:"finding"` + AutoResolveWhen *Condition `yaml:"auto_resolve_when,omitempty" json:"autoResolveWhen,omitempty"` +} + +// Source restricts evaluation before CEL runs. Event types are normalized in +// the same way as ingestion worker event types, so providers can use either +// audit-log names (for example repository.publicized) or canonical names. +type Source struct { + Provider string `yaml:"provider" json:"provider"` + EventTypes []string `yaml:"event_types" json:"eventTypes"` +} + +// Condition is intentionally only an expression string. CEL is compiled at +// pack-load time and executes without access to Go functions, I/O, or a +// database. +type Condition struct { + Expression string `yaml:"expression" json:"expression"` +} + +// DedupeSpec defines the stable subject used by the worker's tenant-scoped +// dedupe key. Templates are deliberately logic-free. +type DedupeSpec struct { + TargetTemplate string `yaml:"target_template" json:"targetTemplate"` +} + +// FindingSpec is the operator-facing output of a rule. +type FindingSpec struct { + TargetTemplate string `yaml:"target_template,omitempty" json:"targetTemplate,omitempty"` + Title string `yaml:"title" json:"title"` + Description string `yaml:"description" json:"description"` + RemediationSteps []string `yaml:"remediation_steps" json:"remediationSteps"` + Evidence map[string]string `yaml:"evidence,omitempty" json:"evidence,omitempty"` +} + +// Event is the provider-neutral input to the evaluator. +type Event struct { + OrganizationID string + IntegrationID string + Provider string + EventType string + Source string + Actor string + OccurredAt time.Time + Payload map[string]any +} + +// OrgContext contains non-secret tenant configuration explicitly made +// available to a rule. It is kept as data rather than a callback so a rule +// cannot perform I/O or escape the evaluator sandbox. +type OrgContext struct { + Config map[string]any + Allowlists map[string]any +} + +// Overrides are applied after a rule is selected. Disabled is the existing +// integration disabled_checks mechanism. SeverityOverrides is intentionally a +// separate map so a caller can preserve disabled-check expiry metadata while +// adding a severity policy in the same integration configuration JSON. +type Overrides struct { + Disabled map[string]bool + SeverityOverrides map[string]string +} + +// FindingDraft is pure evaluator output. DedupeTarget is the rendered, +// version-stable subject; the worker adds organization/integration identity +// before hashing the persisted dedupe key and records RuleVersion alongside +// the finding for migration-safe provenance. +type FindingDraft struct { + RuleID string + RuleVersion string + Title string + Description string + Severity string + RiskScore int + Tags []string + RemediationSteps []string + Target string + DedupeTarget string + Evidence map[string]any +} + +// ResolutionDraft identifies an open finding that a later clean event should +// resolve. The worker owns the state transition and must scope it by tenant, +// integration, rule ID, and rendered target. +type ResolutionDraft struct { + RuleID string + RuleVersion string + DedupeTarget string + Evidence map[string]any +} + +// BacktestReport summarizes deterministic replay over a fixture/event slice. +type BacktestReport struct { + RuleID string `json:"ruleId"` + RuleVersion string `json:"ruleVersion"` + Events int `json:"events"` + Candidates int `json:"candidates"` + Matches int `json:"matches"` + Resolutions int `json:"resolutions"` + MatchSamples []BacktestMatch `json:"matchSamples,omitempty"` +} + +// BacktestMatch is a small evidence sample suitable for CLI/API rendering. +type BacktestMatch struct { + EventType string `json:"eventType"` + OccurredAt string `json:"occurredAt"` + Target string `json:"target"` + DedupeTarget string `json:"dedupeTarget"` + Severity string `json:"severity,omitempty"` + Resolution bool `json:"resolution,omitempty"` +} diff --git a/internal/ingestionworker/declarative_rules.go b/internal/ingestionworker/declarative_rules.go new file mode 100644 index 0000000..b067416 --- /dev/null +++ b/internal/ingestionworker/declarative_rules.go @@ -0,0 +1,127 @@ +package ingestionworker + +import ( + "sync" + + "github.com/writer/aperio/internal/detection" +) + +type declarativeResolution struct { + RuleID string + RuleVersion string + DedupeTarget string +} + +var ( + builtinDetectionOnce sync.Once + builtinDetectionEngine *detection.Engine + builtinDetectionErr error + + // These rules are the first migration slice. The hardcoded evaluators stay + // in worker.go as a fail-closed fallback if a pack cannot compile, but are + // not run when the reviewed declarative engine is healthy. + declarativeRuleIDs = map[string]bool{ + "github.public_repository_created": true, + "slack.mfa_disabled": true, + "slack.external_shared_channel_created": true, + "google_workspace.external_sharing_enabled": true, + } +) + +func builtinDetection() (*detection.Engine, error) { + builtinDetectionOnce.Do(func() { + var rules []detection.Rule + rules, builtinDetectionErr = detection.LoadEmbeddedPack(detection.BuiltinFS, "rules/*.yaml") + if builtinDetectionErr != nil { + return + } + builtinDetectionEngine, builtinDetectionErr = detection.NewEngine(rules) + }) + return builtinDetectionEngine, builtinDetectionErr +} + +func evaluateDeclarativeRules(payload JobPayload, disabledChecks []string, severityOverrides map[string]string) ([]Finding, bool) { + engine, err := builtinDetection() + if err != nil { + return nil, false + } + disabled := make(map[string]bool, len(disabledChecks)) + for _, key := range disabledChecks { + disabled[key] = true + } + drafts, err := engine.Evaluate(toDetectionEvent(payload), detection.OrgContext{}, detection.Overrides{ + Disabled: disabled, + SeverityOverrides: severityOverrides, + }) + if err != nil { + // A malformed provider payload must not suppress the established + // evaluator. Keep the migration fail-closed and let the caller use + // the hardcoded parity implementation for this event. + return nil, false + } + out := make([]Finding, 0, len(drafts)) + for _, draft := range drafts { + out = append(out, Finding{ + RuleID: draft.RuleID, + RuleVersion: draft.RuleVersion, + Title: draft.Title, + Description: draft.Description, + Severity: draft.Severity, + RiskScore: draft.RiskScore, + RemediationSteps: draft.RemediationSteps, + Target: draft.Target, + DedupeTarget: draft.DedupeTarget, + Evidence: draft.Evidence, + Tags: draft.Tags, + }) + } + return out, true +} + +// DeclarativeAutoResolutions is intentionally pure. The worker or a future +// backtest/API caller supplies the returned rule/target pair to its own +// tenant-scoped state transition. No database access occurs here. +func DeclarativeAutoResolutions(payload JobPayload, disabledChecks []string) ([]detection.ResolutionDraft, bool) { + engine, err := builtinDetection() + if err != nil { + return nil, false + } + disabled := make(map[string]bool, len(disabledChecks)) + for _, key := range disabledChecks { + disabled[key] = true + } + resolutions, err := engine.AutoResolve(toDetectionEvent(payload), detection.OrgContext{}, detection.Overrides{Disabled: disabled}) + if err != nil { + return nil, false + } + return resolutions, true +} + +func declarativeResolutionTargets(payload JobPayload, disabledChecks []string) ([]declarativeResolution, bool) { + drafts, loaded := DeclarativeAutoResolutions(payload, disabledChecks) + if !loaded { + return nil, false + } + out := make([]declarativeResolution, 0, len(drafts)) + for _, draft := range drafts { + out = append(out, declarativeResolution{ + RuleID: draft.RuleID, + RuleVersion: draft.RuleVersion, + DedupeTarget: draft.DedupeTarget, + }) + } + return out, true +} + +func toDetectionEvent(payload JobPayload) detection.Event { + return detection.Event{ + OrganizationID: payload.OrganizationID, + IntegrationID: payload.IntegrationID, + Provider: payload.Provider, + EventType: payload.EventType, + Source: payload.Source, + Actor: payload.Actor, + OccurredAt: payload.OccurredAt, + Payload: payload.Payload, + } +} diff --git a/internal/ingestionworker/declarative_rules_test.go b/internal/ingestionworker/declarative_rules_test.go new file mode 100644 index 0000000..c8bf970 --- /dev/null +++ b/internal/ingestionworker/declarative_rules_test.go @@ -0,0 +1,92 @@ +package ingestionworker + +import ( + "reflect" + "testing" + "time" +) + +func TestEvaluateUsesDeclarativeRuleVersionAndTenantSeverityOverride(t *testing.T) { + payload := JobPayload{ + OrganizationID: "org_1", + IntegrationID: "int_slack_1", + Provider: "SLACK", + EventType: "mfa.disabled", + Source: "slack-audit-log", + OccurredAt: time.Date(2026, 6, 6, 1, 0, 0, 0, time.UTC), + Payload: map[string]any{ + "user": map[string]any{"id": "U123", "email": "user@example.com"}, + }, + } + findings := Evaluate(payload, nil) + if len(findings) != 1 || findings[0].RuleVersion != "1.0.0" { + t.Fatalf("declarative findings = %#v", findings) + } + if !reflect.DeepEqual(findings[0].Evidence, map[string]any{"user": "user@example.com", "subject": "user@example.com"}) { + t.Fatalf("declarative evidence = %#v", findings[0].Evidence) + } + overridden := EvaluateWithSeverityOverrides(payload, nil, map[string]string{"slack.mfa_disabled": "LOW"}) + if len(overridden) != 1 || overridden[0].Severity != SeverityLow || overridden[0].RiskScore != RiskScoreFor(SeverityLow) { + t.Fatalf("severity override = %#v", overridden) + } +} + +func TestGitHubDeployKeyRuleEscalatesWriteAccess(t *testing.T) { + payload := JobPayload{ + OrganizationID: "org_1", + IntegrationID: "int_github_1", + Provider: "GITHUB", + EventType: "deploy_key.added", + OccurredAt: time.Date(2026, 6, 6, 1, 0, 0, 0, time.UTC), + Payload: map[string]any{ + "repository": map[string]any{"full_name": "acme/api"}, + "key": map[string]any{"title": "build-bot", "write_enabled": true}, + }, + } + findings := Evaluate(payload, nil) + if len(findings) != 1 { + t.Fatalf("deploy key findings = %#v", findings) + } + if findings[0].RuleID != "github.deploy_key_added" || findings[0].Severity != SeverityHigh || findings[0].DedupeTarget != "acme/api:build-bot" { + t.Fatalf("deploy key finding = %#v", findings[0]) + } + missingIdentity := payload + missingIdentity.Payload = map[string]any{"repository": map[string]any{"full_name": "acme/api"}} + if findings := Evaluate(missingIdentity, nil); len(findings) != 0 { + t.Fatalf("deploy key without key identity should remain unsupported: %#v", findings) + } +} + +func TestDeclarativeAutoResolutionRequiresCleanEvent(t *testing.T) { + payload := JobPayload{ + OrganizationID: "org_1", + IntegrationID: "int_slack_1", + Provider: "SLACK", + EventType: "two-factor auth enabled", + OccurredAt: time.Date(2026, 6, 6, 1, 0, 0, 0, time.UTC), + Payload: map[string]any{"user": map[string]any{"email": "user@example.com"}}, + } + resolutions, loaded := DeclarativeAutoResolutions(payload, nil) + if !loaded || len(resolutions) != 1 || resolutions[0].RuleID != "slack.mfa_disabled" || resolutions[0].DedupeTarget != "user@example.com" { + t.Fatalf("auto resolutions = %#v loaded=%t", resolutions, loaded) + } +} + +func TestDeclarativeRuleCatalogVersionsAreStable(t *testing.T) { + for ruleID := range declarativeRuleIDs { + found := false + for _, entry := range RuleCatalog { + if entry.ID != ruleID { + continue + } + found = true + if entry.Version != "1.0.0" { + t.Errorf("catalog rule %s version = %q, want 1.0.0", ruleID, entry.Version) + } + break + } + if !found { + t.Errorf("declarative rule %s has no catalog entry", ruleID) + } + } +} diff --git a/internal/ingestionworker/detection_packs.go b/internal/ingestionworker/detection_packs.go index 707e1a7..d9234b1 100644 --- a/internal/ingestionworker/detection_packs.go +++ b/internal/ingestionworker/detection_packs.go @@ -27,8 +27,8 @@ var DetectionPacks = []DetectionPack{ ID: "aperio.github.core.v1", Provider: "GITHUB", Name: "GitHub repository hygiene", - Description: "Repository exposure, branch-protection drift, and risky GitHub OAuth app activity.", - Version: "1.1.0", + Description: "Repository exposure, branch-protection drift, risky OAuth app activity, and deploy-key additions.", + Version: "1.2.0", }, { ID: "aperio.slack.core.v1", diff --git a/internal/ingestionworker/disabled_check_metadata.go b/internal/ingestionworker/disabled_check_metadata.go index 9647959..30c0750 100644 --- a/internal/ingestionworker/disabled_check_metadata.go +++ b/internal/ingestionworker/disabled_check_metadata.go @@ -9,6 +9,7 @@ import ( type disabledCheckMetadataEntry struct { Reason string `json:"reason"` ExpiresAt string `json:"expiresAt"` + Severity string `json:"severity,omitempty"` } func decodeDisabledCheckMetadata(raw string) map[string]disabledCheckMetadataEntry { @@ -50,3 +51,25 @@ func applyDisabledCheckExpiry(disabled []string, metadata map[string]disabledChe } return effective } + +// severityOverridesFromMetadata reads the optional severity policy stored +// alongside disabled_checks. It intentionally accepts only canonical finding +// severities and ignores expired/invalid entries so a malformed tenant policy +// cannot make the worker fail or emit an invalid database enum. +func severityOverridesFromMetadata(metadata map[string]disabledCheckMetadataEntry, now time.Time) map[string]string { + out := make(map[string]string) + for key, entry := range metadata { + severity := strings.ToUpper(strings.TrimSpace(entry.Severity)) + if key == "" || !validFindingSeverity(severity) { + continue + } + if expiresAt := strings.TrimSpace(entry.ExpiresAt); expiresAt != "" { + expires, err := time.Parse(time.RFC3339, expiresAt) + if err != nil || !expires.After(now) { + continue + } + } + out[key] = severity + } + return out +} diff --git a/internal/ingestionworker/disabled_check_metadata_test.go b/internal/ingestionworker/disabled_check_metadata_test.go index d134323..05bfae7 100644 --- a/internal/ingestionworker/disabled_check_metadata_test.go +++ b/internal/ingestionworker/disabled_check_metadata_test.go @@ -22,3 +22,20 @@ func TestApplyDisabledCheckExpiry(t *testing.T) { t.Fatalf("effective disabled checks = %v, want %v", got, want) } } + +func TestSeverityOverridesFromMetadata(t *testing.T) { + now := time.Date(2026, 6, 16, 8, 0, 0, 0, time.UTC) + got := severityOverridesFromMetadata(map[string]disabledCheckMetadataEntry{ + "slack.mfa_disabled": {Severity: "high", ExpiresAt: now.Add(time.Hour).Format(time.RFC3339)}, + "github.public_repository_created": {Severity: "LOW"}, + "slack.expired": {Severity: "CRITICAL", ExpiresAt: now.Add(-time.Minute).Format(time.RFC3339)}, + "slack.invalid": {Severity: "urgent"}, + }, now) + want := map[string]string{ + "slack.mfa_disabled": SeverityHigh, + "github.public_repository_created": SeverityLow, + } + if !reflect.DeepEqual(got, want) { + t.Fatalf("severity overrides = %#v, want %#v", got, want) + } +} diff --git a/internal/ingestionworker/rules_catalog.go b/internal/ingestionworker/rules_catalog.go index f943d4f..8d94020 100644 --- a/internal/ingestionworker/rules_catalog.go +++ b/internal/ingestionworker/rules_catalog.go @@ -15,6 +15,7 @@ package ingestionworker // cross-provider categorization (see tags.go). type RuleCatalogEntry struct { ID string + Version string Provider string Title string Description string @@ -31,6 +32,7 @@ type RuleCatalogEntry struct { var RuleCatalog = []RuleCatalogEntry{ { ID: "github.public_repository_created", + Version: "1.0.0", Provider: "GITHUB", Title: "Public GitHub repository created", Description: "A repository was created or changed to public visibility, which can expose source code, secrets, or customer data.", @@ -65,8 +67,22 @@ var RuleCatalog = []RuleCatalogEntry{ Intent: "Adversary obtains a long-lived third-party authorization that can read repositories or administer organization resources.", Tags: []string{TagOAuthRiskyGrant, TagDataAccess}, }, + { + ID: "github.deploy_key_added", + Version: "1.0.0", + Provider: "GITHUB", + Title: "GitHub deploy key added", + Description: "A deploy key was added to a repository; write-enabled keys can bypass normal user and review controls.", + Severity: "MEDIUM", + EventTypes: []string{"DEPLOY_KEY_ADDED", "DEPLOY_KEY_CREATED"}, + PackID: "aperio.github.core.v1", + MitreTechniques: []string{"T1098"}, + Intent: "Adversary establishes a non-user credential that can read or write repository content without normal identity controls.", + Tags: []string{TagDataAccess, TagPolicyWeakened}, + }, { ID: "slack.mfa_disabled", + Version: "1.0.0", Provider: "SLACK", Title: "Slack multi-factor authentication disabled", Description: "A Slack user disabled MFA, increasing the likelihood of account takeover and lateral movement.", @@ -79,6 +95,7 @@ var RuleCatalog = []RuleCatalogEntry{ }, { ID: "slack.external_shared_channel_created", + Version: "1.0.0", Provider: "SLACK", Title: "Slack external shared channel created", Description: "A Slack channel was shared with an external organization, expanding conversation and file visibility outside the tenant.", @@ -163,6 +180,7 @@ var RuleCatalog = []RuleCatalogEntry{ }, { ID: "google_workspace.external_sharing_enabled", + Version: "1.0.0", Provider: "GOOGLE_WORKSPACE", Title: "Google Drive external sharing enabled", Description: "A Drive resource was shared outside the tenant domain or set to a public visibility scope.", diff --git a/internal/ingestionworker/worker.go b/internal/ingestionworker/worker.go index 142d157..9d032cf 100644 --- a/internal/ingestionworker/worker.go +++ b/internal/ingestionworker/worker.go @@ -41,6 +41,12 @@ var supportedIngestionEventTypes = map[string][]string{ "GITHUB": { "PUBLIC_REPOSITORY_CREATED", "repository.publicized", + "REPOSITORY_PRIVATE", + "REPOSITORY_PRIVATEIZED", + "REPOSITORY_VISIBILITY_CHANGED", + "repository.private", + "repository.privateized", + "repository.visibility_changed", "BRANCH_PROTECTION_DISABLED", "BRANCH_PROTECTION_RULE_DELETED", "BRANCH_PROTECTION_RULE_UPDATED", @@ -53,12 +59,20 @@ var supportedIngestionEventTypes = map[string][]string{ "oauth_app.installed", "github_app.installed", "org.oauth_app_access_approved", + "DEPLOY_KEY_ADDED", + "DEPLOY_KEY_CREATED", + "deploy_key.added", + "deploy_key.created", }, "SLACK": { "MFA_DISABLED", "TWO_FACTOR_AUTH_DISABLED", "mfa.disabled", "two-factor auth disabled", + "MFA_ENABLED", + "TWO_FACTOR_AUTH_ENABLED", + "mfa.enabled", + "two-factor auth enabled", "EXTERNAL_SHARED_CHANNEL_CREATED", "SHARED_CHANNEL_INVITE_ACCEPTED", "external_shared_channel.created", @@ -103,6 +117,12 @@ var supportedIngestionEventTypes = map[string][]string{ "GOOGLE_WORKSPACE": { "EXTERNAL_SHARING_ENABLED", "external.sharing.enabled", + "EXTERNAL_SHARING_DISABLED", + "DRIVE_FILE_VISIBILITY_CHANGED", + "DRIVE_FILE_PRIVATE", + "external.sharing.disabled", + "drive.file.visibility.changed", + "drive.file.private", "SUPER_ADMIN_GRANTED", "super.admin.granted", "ADMIN_ROLE_GRANTED", @@ -193,6 +213,7 @@ type JobPayload struct { type Finding struct { RuleID string + RuleVersion string Title string Description string Severity string @@ -291,6 +312,7 @@ type integrationConfig struct { Provider string ExternalAccountID string DisabledChecks []string + SeverityOverrides map[string]string EncryptedAccessToken string EncryptedRefreshToken sql.NullString EncryptedWebhookSecret sql.NullString @@ -316,12 +338,21 @@ func (w *Worker) WithCerebroFanout(fanout CerebroFindingFanout) *Worker { } func Evaluate(payload JobPayload, disabledChecks []string) []Finding { + return EvaluateWithSeverityOverrides(payload, disabledChecks, nil) +} + +// EvaluateWithSeverityOverrides preserves the existing disabled_checks +// contract while allowing a tenant policy layer to lower or raise a rule's +// severity. The declarative migration slice is authoritative when its pack +// compiles; hardcoded rules remain the compatibility fallback. +func EvaluateWithSeverityOverrides(payload JobPayload, disabledChecks []string, severityOverrides map[string]string) []Finding { disabled := map[string]struct{}{} for _, check := range disabledChecks { disabled[check] = struct{}{} } - findings := []Finding{} - if _, ok := disabled["github.public_repository_created"]; !ok { + declarativeFindings, declarativeLoaded := evaluateDeclarativeRules(payload, disabledChecks, severityOverrides) + findings := append([]Finding{}, declarativeFindings...) + if _, ok := disabled["github.public_repository_created"]; !ok && (!declarativeLoaded || !declarativeRuleIDs["github.public_repository_created"]) { if finding, ok := evaluateGitHubPublicRepository(payload); ok { findings = append(findings, finding) } @@ -336,7 +367,12 @@ func Evaluate(payload JobPayload, disabledChecks []string) []Finding { findings = append(findings, finding) } } - if _, ok := disabled["slack.mfa_disabled"]; !ok { + if _, ok := disabled["github.deploy_key_added"]; !ok { + if finding, ok := evaluateGitHubDeployKeyAdded(payload); ok { + findings = append(findings, finding) + } + } + if _, ok := disabled["slack.mfa_disabled"]; !ok && (!declarativeLoaded || !declarativeRuleIDs["slack.mfa_disabled"]) { if finding, ok := evaluateSlackMFADisabled(payload); ok { findings = append(findings, finding) } @@ -376,7 +412,7 @@ func Evaluate(payload JobPayload, disabledChecks []string) []Finding { findings = append(findings, finding) } } - if _, ok := disabled["google_workspace.external_sharing_enabled"]; !ok { + if _, ok := disabled["google_workspace.external_sharing_enabled"]; !ok && (!declarativeLoaded || !declarativeRuleIDs["google_workspace.external_sharing_enabled"]) { if finding, ok := evaluateGoogleExternalSharingEnabled(payload); ok { findings = append(findings, finding) } @@ -471,9 +507,24 @@ func Evaluate(payload JobPayload, disabledChecks []string) []Finding { findings = append(findings, finding) } } + for index := range findings { + if override := strings.ToUpper(strings.TrimSpace(severityOverrides[findings[index].RuleID])); override != "" && validFindingSeverity(override) { + findings[index].Severity = override + findings[index].RiskScore = RiskScoreFor(override) + } + } return findings } +func validFindingSeverity(value string) bool { + switch value { + case SeverityCritical, SeverityHigh, SeverityMedium, SeverityLow, SeverityInfo: + return true + default: + return false + } +} + func evaluateGitHubPublicRepository(payload JobPayload) (Finding, bool) { if payload.Provider != "GITHUB" { return Finding{}, false @@ -617,6 +668,69 @@ func evaluateGitHubOAuthAppInstalled(payload JobPayload) (Finding, bool) { }, true } +func evaluateGitHubDeployKeyAdded(payload JobPayload) (Finding, bool) { + if payload.Provider != "GITHUB" { + return Finding{}, false + } + switch normalizeEventType(payload.EventType) { + case "DEPLOY_KEY_ADDED", "DEPLOY_KEY_CREATED": + default: + return Finding{}, false + } + repository := firstNonEmpty( + nestedString(payload.Payload, "repository", "full_name"), + nestedString(payload.Payload, "repository", "name"), + nestedString(payload.Payload, "repo"), + ) + key := firstNonEmpty( + nestedString(payload.Payload, "key", "title"), + nestedString(payload.Payload, "key", "name"), + nestedString(payload.Payload, "deploy_key", "title"), + nestedString(payload.Payload, "deploy_key", "id"), + nestedString(payload.Payload, "key_id"), + ) + // A deploy-key event without a repository and key identity cannot produce + // a stable target or a trustworthy remediation. Leave it unsupported until + // the provider adapter supplies the fields listed in the support matrix. + if repository == "" || key == "" { + return Finding{}, false + } + writeEnabled, hasWriteEnabled := nestedBool(payload.Payload, "key", "write_enabled") + if !hasWriteEnabled { + writeEnabled, hasWriteEnabled = nestedBool(payload.Payload, "deploy_key", "write_enabled") + } + severity := SeverityMedium + riskScore := RiskScoreFor(SeverityMedium, 4) + if hasWriteEnabled && writeEnabled { + severity = SeverityHigh + riskScore = RiskScoreFor(SeverityHigh, 7) + } + subject := repository + ":" + key + return Finding{ + RuleID: "github.deploy_key_added", + RuleVersion: "1.0.0", + Title: "GitHub deploy key added", + Description: "A deploy key was added to a GitHub repository; write-enabled keys can bypass normal user and review controls.", + Severity: severity, + RiskScore: riskScore, + Tags: []string{TagDataAccess, TagPolicyWeakened}, + RemediationSteps: []string{ + "Confirm the deploy key belongs to an approved automation system.", + "Disable write access unless the automation requires repository writes.", + "Remove the key and rotate its credential if it was not approved.", + }, + Target: repository, + DedupeTarget: subject, + Evidence: compactEvidence(map[string]any{ + "repository": repository, + "key": key, + "writeEnabled": writeEnabled, + "actor": payload.Actor, + "subject": subject, + }), + }, true +} + func evaluateSlackMFADisabled(payload JobPayload) (Finding, bool) { if payload.Provider != "SLACK" { return Finding{}, false @@ -748,12 +862,19 @@ func evaluateSlackAppInstalled(payload JobPayload) (Finding, bool) { stringArray(payload.Payload["scopes"]), stringArray(nestedRecord(payload.Payload, "app")["scopes"])..., )) + severity := SeverityMedium + riskScore := RiskScoreFor(SeverityMedium, 7) + scopeBlob := strings.ToLower(strings.Join(scopes, " ")) + if strings.Contains(scopeBlob, "admin") || strings.Contains(scopeBlob, "files:read") || strings.Contains(scopeBlob, "channels:history") { + severity = SeverityHigh + riskScore = RiskScoreFor(SeverityHigh, 6) + } return Finding{ RuleID: "slack.app_installed", Title: "Third-party Slack app installed", Description: "A third-party Slack app was installed with user, channel, admin, or file scopes.", - Severity: SeverityMedium, - RiskScore: RiskScoreFor(SeverityMedium, 7), + Severity: severity, + RiskScore: riskScore, Tags: []string{TagOAuthRiskyGrant, TagDataAccess}, RemediationSteps: []string{ "Confirm the Slack app is approved for the workspace.", @@ -2349,7 +2470,7 @@ func (w *Worker) process(ctx context.Context, item job) error { if err != nil { return w.fail(ctx, item, fmt.Errorf("parse payload: %w", err).Error()) } - findings, err := w.findingsForJob(ctx, payload, item) + findings, resolutions, err := w.evaluateJob(ctx, payload, item) if err != nil { return w.fail(ctx, item, fmt.Errorf("load findings: %w", err).Error()) } @@ -2379,6 +2500,23 @@ func (w *Worker) process(ctx context.Context, item job) error { `, eventID, item.OrganizationID, item.IntegrationID, item.ID, item.Provider, item.EventType, item.Source, nullableString(item.Actor), string(item.Payload), item.OccurredAt).Scan(&eventID); err != nil { return fail(fmt.Errorf("upsert ingested event: %w", err)) } + for _, resolution := range resolutions { + resolvedID, changed, err := resolveDeclarativeFinding(ctx, tx, payload, resolution, eventID) + if err != nil { + return fail(fmt.Errorf("auto-resolve finding: %w", err)) + } + if changed { + lifecycleEvents = append(lifecycleEvents, FindingLifecycleEvent{ + FindingID: resolvedID, + OrganizationID: payload.OrganizationID, + IntegrationID: payload.IntegrationID, + PreviousStatus: "OPEN", + NextStatus: "RESOLVED", + OccurredAt: payload.OccurredAt, + ResolutionNote: "Declarative rule observed a clean provider state", + }) + } + } for _, finding := range findings { persisted, err := upsertFinding(ctx, tx, payload, finding, eventID) if err != nil { @@ -2452,18 +2590,27 @@ func (w *Worker) fail(ctx context.Context, item job, message string) error { } func (w *Worker) findingsForJob(ctx context.Context, payload JobPayload, item job) ([]Finding, error) { + findings, _, err := w.evaluateJob(ctx, payload, item) + return findings, err +} + +func (w *Worker) evaluateJob(ctx context.Context, payload JobPayload, item job) ([]Finding, []declarativeResolution, error) { builtinRun := observability.StartRuleRun(ctx, w.db, item.OrganizationID, item.IntegrationID, item.Provider, item.ID, observability.RulePackBuiltIn, "v1", builtInRuleCount(item.Provider)) config, err := w.loadIntegrationConfig(ctx, item) if err != nil { builtinRun.Finish(ctx, "FAILED", 0, err) - return nil, err + return nil, nil, err } if err := config.validateForJob(item); err != nil { builtinRun.Finish(ctx, "FAILED", 0, err) - return nil, err + return nil, nil, err } - findings := Evaluate(payload, config.DisabledChecks) + findings := EvaluateWithSeverityOverrides(payload, config.DisabledChecks, config.SeverityOverrides) builtinRun.Finish(ctx, "SUCCEEDED", len(findings), nil) + resolutions, loaded := declarativeResolutionTargets(payload, config.DisabledChecks) + if !loaded { + resolutions = nil + } customRun := observability.StartRuleRun(ctx, w.db, item.OrganizationID, item.IntegrationID, item.Provider, item.ID, observability.RulePackCustom, "v1", 0) customRules, err := w.loadCustomRules(ctx, item.IntegrationID) if err != nil { @@ -2472,7 +2619,7 @@ func (w *Worker) findingsForJob(ctx context.Context, payload JobPayload, item jo // otherwise mask real-finding ingestion. Log via the caller's // observability surface and fall through with the built-ins. customRun.Finish(ctx, "FAILED", 0, err) - return findings, nil + return findings, resolutions, nil } if len(customRules) > 0 { customRun.SetRulesEvaluated(len(customRules)) @@ -2482,7 +2629,7 @@ func (w *Worker) findingsForJob(ctx context.Context, payload JobPayload, item jo } else { customRun.Finish(ctx, "SUCCEEDED", 0, nil) } - return findings, nil + return findings, resolutions, nil } func (w *Worker) loadCustomRules(ctx context.Context, integrationID string) ([]CustomRule, error) { @@ -2548,7 +2695,10 @@ func (w *Worker) loadIntegrationConfig(ctx context.Context, item job) (integrati if err := json.Unmarshal([]byte(rawDisabledChecks), &config.DisabledChecks); err != nil { return integrationConfig{}, errIntegrationConfigurationIncomplete } - config.DisabledChecks = applyDisabledCheckExpiry(config.DisabledChecks, decodeDisabledCheckMetadata(rawDisabledMetadata), time.Now().UTC()) + metadata := decodeDisabledCheckMetadata(rawDisabledMetadata) + now := time.Now().UTC() + config.DisabledChecks = applyDisabledCheckExpiry(config.DisabledChecks, metadata, now) + config.SeverityOverrides = severityOverridesFromMetadata(metadata, now) return config, nil } @@ -2686,6 +2836,52 @@ func upsertFinding(ctx context.Context, tx *sql.Tx, payload JobPayload, finding return persisted, err } +// resolveDeclarativeFinding performs the state transition requested by a +// pure auto-resolve draft. The lookup is tenant and integration scoped, and +// only OPEN findings transition; muted, already-resolved, or another +// integration's finding is never changed by a clean provider event. +func resolveDeclarativeFinding(ctx context.Context, tx *sql.Tx, payload JobPayload, resolution declarativeResolution, eventID string) (string, bool, error) { + placeholder := Finding{ + RuleID: resolution.RuleID, + RuleVersion: resolution.RuleVersion, + Target: resolution.DedupeTarget, + DedupeTarget: resolution.DedupeTarget, + } + dedupe := DedupeKey(payload, placeholder) + evidence, err := json.Marshal(map[string]any{ + "ruleId": resolution.RuleID, + "ruleVersion": resolution.RuleVersion, + "subject": resolution.DedupeTarget, + "sourceEventId": eventID, + "eventType": payload.EventType, + "resolution": "auto_resolve_when", + }) + if err != nil { + return "", false, err + } + var findingID string + err = tx.QueryRowContext(ctx, ` + UPDATE security_findings + SET status = 'RESOLVED'::"FindingStatus", + resolved_at = $1, + resolved_by_id = NULL, + evidence = COALESCE(evidence, '{}'::jsonb) || $2::jsonb + WHERE organization_id = $3 + AND integration_id = $4 + AND dedupe_key = $5 + AND COALESCE(evidence->>'ruleId', '') = $6 + AND status = 'OPEN'::"FindingStatus" + RETURNING id + `, payload.OccurredAt, string(evidence), payload.OrganizationID, payload.IntegrationID, dedupe, resolution.RuleID).Scan(&findingID) + if errors.Is(err, sql.ErrNoRows) { + return "", false, nil + } + if err != nil { + return "", false, err + } + return findingID, true, nil +} + func buildFindingEvidence(payload JobPayload, finding Finding, eventID string) map[string]any { subject := finding.Target if strings.TrimSpace(finding.DedupeTarget) != "" { @@ -2700,6 +2896,9 @@ func buildFindingEvidence(payload JobPayload, finding Finding, eventID string) m "eventType": payload.EventType, "sourceEventId": eventID, } + if strings.TrimSpace(finding.RuleVersion) != "" { + evidence["ruleVersion"] = finding.RuleVersion + } addNonEmptyEvidence(evidence, "actor", payload.Actor) addNonEmptyEvidence(evidence, "application", nestedString(payload.Payload, "application")) addNonEmptyEvidence(evidence, "sourceIp", nestedString(payload.Payload, "ipAddress")) @@ -2872,6 +3071,9 @@ func findingPayload(payload JobPayload, finding Finding, eventID string, persist "eventType": payload.EventType, "actor": actor, } + if strings.TrimSpace(finding.RuleVersion) != "" { + record["ruleVersion"] = finding.RuleVersion + } addOAuthClaimRecordFields(record, payload, finding) return siemdispatcher.Payload{ Kind: "finding", diff --git a/tests/fixtures/migration-ownership/migration-matrix.json b/tests/fixtures/migration-ownership/migration-matrix.json index a086c2a..b5d75ce 100644 --- a/tests/fixtures/migration-ownership/migration-matrix.json +++ b/tests/fixtures/migration-ownership/migration-matrix.json @@ -617,6 +617,21 @@ "command:npm run worker:ingestion -- -once -limit 1" ] }, + { + "id": "internal-detection-go-declarative", + "state": "go-default", + "covers": [ + "repo-file:internal/detection/*" + ], + "owner": "Go declarative detection engine and reviewed rule packs", + "rationale": "internal/detection owns the bounded CEL evaluator, strict YAML loader, versioned built-in rule pack, and in-memory backtest contract used by the Go ingestion worker; it has no independent runtime or provider credentials.", + "evidence": [ + "test:internal/detection/loader_test.go", + "test:internal/detection/engine_test.go", + "source:internal/detection/rules", + "validator:go-tests" + ] + }, { "id": "internal-siemdispatcher-go-default", "state": "go-default", diff --git a/tests/fixtures/worker-parity/ingestion-rule-matrix.json b/tests/fixtures/worker-parity/ingestion-rule-matrix.json index e0258a6..a952ce8 100644 --- a/tests/fixtures/worker-parity/ingestion-rule-matrix.json +++ b/tests/fixtures/worker-parity/ingestion-rule-matrix.json @@ -123,7 +123,13 @@ ], "goClaimedEventTypes": [ "PUBLIC_REPOSITORY_CREATED", - "repository.publicized" + "repository.publicized", + "REPOSITORY_PRIVATE", + "REPOSITORY_PRIVATEIZED", + "REPOSITORY_VISIBILITY_CHANGED", + "repository.private", + "repository.privateized", + "repository.visibility_changed" ], "fixtures": [ "tests/fixtures/worker-parity/github-public-repository.json" @@ -136,6 +142,35 @@ "cutoverBlockers": [], "goDefaultBlocked": false }, + { + "ruleId": "github.deploy_key_added", + "provider": "GITHUB", + "state": "go-default", + "typescriptEventAliases": [ + "DEPLOY_KEY_ADDED", + "DEPLOY_KEY_CREATED" + ], + "typescriptPayloadPredicates": [ + "repository and deploy-key identity are present; write_enabled escalates severity when true" + ], + "goClaimedEventTypes": [ + "DEPLOY_KEY_ADDED", + "DEPLOY_KEY_CREATED", + "deploy_key.added", + "deploy_key.created" + ], + "fixtures": [ + "tests/fixtures/worker-parity/saas-pack-rules.json" + ], + "tests": [ + "tests/ingestion-parity-matrix.test.ts", + "internal/ingestionworker/worker_test.go", + "internal/ingestionworker/worker_db_test.go", + "internal/ingestionworker/declarative_rules_test.go" + ], + "cutoverBlockers": [], + "goDefaultBlocked": false + }, { "ruleId": "github.branch_protection_disabled", "provider": "GITHUB", @@ -213,7 +248,11 @@ "MFA_DISABLED", "TWO_FACTOR_AUTH_DISABLED", "mfa.disabled", - "two-factor auth disabled" + "two-factor auth disabled", + "MFA_ENABLED", + "TWO_FACTOR_AUTH_ENABLED", + "mfa.enabled", + "two-factor auth enabled" ], "fixtures": [ "tests/fixtures/worker-parity/slack-mfa-disabled.json" @@ -449,7 +488,13 @@ "typescriptPayloadPredicates": [], "goClaimedEventTypes": [ "EXTERNAL_SHARING_ENABLED", - "external.sharing.enabled" + "external.sharing.enabled", + "EXTERNAL_SHARING_DISABLED", + "DRIVE_FILE_VISIBILITY_CHANGED", + "DRIVE_FILE_PRIVATE", + "external.sharing.disabled", + "drive.file.visibility.changed", + "drive.file.private" ], "fixtures": [ "tests/fixtures/worker-parity/google-admin-oauth-rules.json" diff --git a/tests/fixtures/worker-parity/saas-pack-rules.json b/tests/fixtures/worker-parity/saas-pack-rules.json index 2d74131..38b32dd 100644 --- a/tests/fixtures/worker-parity/saas-pack-rules.json +++ b/tests/fixtures/worker-parity/saas-pack-rules.json @@ -23,6 +23,20 @@ "scopes": ["repo", "admin:org"] } }, + { + "provider": "GITHUB", + "eventType": "DEPLOY_KEY_ADDED", + "ruleId": "github.deploy_key_added", + "payload": { + "repository": { + "full_name": "acme/api" + }, + "key": { + "title": "build-bot", + "write_enabled": true + } + } + }, { "provider": "SLACK", "eventType": "EXTERNAL_SHARED_CHANNEL_CREATED", diff --git a/tests/ingestion-parity-matrix.test.ts b/tests/ingestion-parity-matrix.test.ts index b90cc04..29610e9 100644 --- a/tests/ingestion-parity-matrix.test.ts +++ b/tests/ingestion-parity-matrix.test.ts @@ -87,6 +87,7 @@ const expectedRuleIds = [ "atlassian.org_admin_granted", "atlassian.public_space_created", "github.branch_protection_disabled", + "github.deploy_key_added", "github.oauth_app_installed", "github.public_repository_created", "google_workspace.admin_external_recovery_email", diff --git a/tests/migration-ownership-guardrails.test.ts b/tests/migration-ownership-guardrails.test.ts index 30a3c0c..08446ed 100644 --- a/tests/migration-ownership-guardrails.test.ts +++ b/tests/migration-ownership-guardrails.test.ts @@ -107,6 +107,7 @@ function inventoryItems() { ...filesUnder("workers", (file) => file.endsWith(".ts")), ...filesUnder("apps/mcp", (file) => file.endsWith(".ts")), ...filesUnder("internal/bootstrap", (file) => file.endsWith(".go")), + ...filesUnder("internal/detection", (file) => /\.(?:go|ya?ml|json|md)$/.test(file)), ...filesUnder("internal/ingestionworker", (file) => file.endsWith(".go")), ...filesUnder("internal/mcpbroker", (file) => file.endsWith(".go")), ...filesUnder("internal/siemdispatcher", (file) => file.endsWith(".go")), @@ -753,7 +754,7 @@ test("validator and CI gates include contracts, audit, worker smoke, and secret assert.match(ci, /npm run smoke:e2e/); assert.match(ci, /make lint/); assert.match(ci, /needs: \[verify-shard, go-connect, e2e-smoke\]/); - assert.match(ci, /go test \.\/\.\.\./); + assert.match(ci, /go test -p 1 \.\/\.\.\./); assert.match(contracts, /buf\/cmd\/buf@v1\.59\.0 lint/); assert.match(contracts, /buf\/cmd\/buf@v1\.59\.0 breaking/); assert.match(contracts, /git diff --exit-code -- gen packages\/connect\/src\/gen/);