diff --git a/scanners/findings-processor/.env.example b/scanners/findings-processor/.env.example new file mode 100644 index 000000000..556c894d5 --- /dev/null +++ b/scanners/findings-processor/.env.example @@ -0,0 +1,16 @@ +NATS_URL= +NATS_STREAM= +NATS_SUBJECT= +NATS_CONSUMER_DURABLE= +NATS_QUEUE_GROUP= +NATS_ACK_WAIT= +NATS_MAX_DELIVER= +NATS_MAX_ACK_PENDING= + +DB_URL= +DB_USER= +DB_NAME= +DB_PASS= + +LOG_LEVEL= +LOG_PRETTY= diff --git a/scanners/findings-processor/Dockerfile b/scanners/findings-processor/Dockerfile new file mode 100644 index 000000000..68ee1ab68 --- /dev/null +++ b/scanners/findings-processor/Dockerfile @@ -0,0 +1,19 @@ +FROM golang:1.25.0-alpine3.22 AS build + +WORKDIR /src + +COPY go.mod go.sum ./ +RUN go mod download + +COPY . . +RUN CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -o /out/findings-processor ./cmd/service + +FROM alpine:3.22 + +RUN addgroup -S app && adduser -S app -G app +USER app + +WORKDIR /app +COPY --from=build /out/findings-processor /app/findings-processor + +ENTRYPOINT ["/app/findings-processor"] diff --git a/scanners/findings-processor/Makefile b/scanners/findings-processor/Makefile new file mode 100644 index 000000000..18976bc46 --- /dev/null +++ b/scanners/findings-processor/Makefile @@ -0,0 +1,51 @@ +.PHONY: help run test test-race test-integration build fmt fmt-check vet lint tidy ci + +GO ?= go +SERVICE_BIN ?= findings-processor +BUILD_DIR ?= bin + +help: + @printf "Targets:\n" + @printf " make run - Run the service\n" + @printf " make test - Run all tests\n" + @printf " make test-race - Run tests with race detector\n" + @printf " make test-integration - Run integration tests (requires Docker)\n" + @printf " make build - Build service binary\n" + @printf " make fmt - Format Go files\n" + @printf " make fmt-check - Check formatting (no changes)\n" + @printf " make vet - Run go vet\n" + @printf " make lint - Run fmt-check + vet\n" + @printf " make tidy - Tidy modules\n" + @printf " make ci - Lint, test, build\n" + +run: + $(GO) run ./cmd/service + +test: + $(GO) test ./... + +test-integration: + $(GO) test -tags integration ./... + +test-race: + $(GO) test -race ./... + +build: + mkdir -p $(BUILD_DIR) + CGO_ENABLED=0 $(GO) build -o $(BUILD_DIR)/$(SERVICE_BIN) ./cmd/service + +fmt: + $(GO) fmt ./... + +fmt-check: + @test -z "$$($(GO)fmt -l .)" || (printf "Unformatted files found. Run 'make fmt'.\n" && exit 1) + +vet: + $(GO) vet ./... + +lint: fmt-check vet + +tidy: + $(GO) mod tidy + +ci: lint test build diff --git a/scanners/findings-processor/README.md b/scanners/findings-processor/README.md new file mode 100644 index 000000000..cbc459230 --- /dev/null +++ b/scanners/findings-processor/README.md @@ -0,0 +1,67 @@ +# Findings Processor + +Consumes finding events from NATS JetStream and upserts normalized finding documents into ArangoDB `additionalFindings`. + +## What it does + +- Subscribes to `scans.findings.*` (configurable) +- Validates incoming payloads +- Writes findings to ArangoDB +- Acknowledges messages with explicit `Ack` / `Nak` / `Term` behavior + +## Event contract (current) + +Expected JSON fields: + +- Required: `source`, `findingType`, `domainKey`, `subject`, `confidence`, `observedAt` +- Optional: `severity`, `reasonCode`, `evidence`, `attributes` + +`observedAt` must be RFC3339. + +## Local development + +1. Copy env template: + +```bash +cp .env.example .env +``` + +2. Fill required vars in `.env` (at minimum DB and NATS settings). + +3. Run service: + +```bash +go run ./cmd/service +``` + +## Environment variables + +| Variable | Default | Notes | +| ----------------------- | ----------------------- | ---------------------------- | +| `NATS_URL` | `nats://localhost:4222` | NATS server URL | +| `NATS_STREAM` | `SCANS` | JetStream stream name | +| `NATS_SUBJECT` | `scans.findings.*` | Subscription subject | +| `NATS_CONSUMER_DURABLE` | `findings-processor` | Durable consumer name | +| `NATS_ACK_WAIT` | `30s` | Ack timeout | +| `NATS_MAX_DELIVER` | `10` | Max redeliveries | +| `NATS_MAX_ACK_PENDING` | `256` | Max pending unacked messages | +| `DB_URL` | `http://localhost:8529` | ArangoDB URL | +| `DB_USER` | _(none)_ | ArangoDB user | +| `DB_NAME` | _(none)_ | ArangoDB database | +| `DB_PASS` | _(empty)_ | ArangoDB password | +| `LOG_LEVEL` | `info` | Zerolog global level | +| `LOG_PRETTY` | `true` | Human-readable logs | + +## Docker + +Build: + +```bash +docker build -t findings-processor . +``` + +Run: + +```bash +docker run --rm --env-file .env findings-processor +``` diff --git a/scanners/findings-processor/cloudbuild.yaml b/scanners/findings-processor/cloudbuild.yaml new file mode 100644 index 000000000..604dddf37 --- /dev/null +++ b/scanners/findings-processor/cloudbuild.yaml @@ -0,0 +1,52 @@ +steps: + - name: "golang:1.25" + id: ci-checks + dir: scanners/findings-processor + entrypoint: "bash" + args: + - "-c" + - | + make ci + + - name: "golang:1.25" + id: integration-tests + dir: scanners/findings-processor + entrypoint: "bash" + args: + - "-c" + - | + make test-integration + + - name: "gcr.io/cloud-builders/docker" + id: generate-image-name + entrypoint: "bash" + dir: scanners/findings-processor + args: + - "-c" + - | + echo "northamerica-northeast1-docker.pkg.dev/track-compliance/tracker/findings-processor:$(echo $BRANCH_NAME | sed 's/[^a-zA-Z0-9]/-/g')-$SHORT_SHA-$(date +%s)" > /workspace/imagename + + - name: "gcr.io/cloud-builders/docker" + id: build + entrypoint: "bash" + dir: scanners/findings-processor + args: + - "-c" + - | + image=$(cat /workspace/imagename) + docker build -t $image . + + - name: "gcr.io/cloud-builders/docker" + id: push-if-master + entrypoint: "bash" + dir: scanners/findings-processor + args: + - "-c" + - | + if [[ "$BRANCH_NAME" == "master" ]] + then + image=$(cat /workspace/imagename) + docker push $image + else + exit 0 + fi diff --git a/scanners/findings-processor/cmd/service/main.go b/scanners/findings-processor/cmd/service/main.go new file mode 100644 index 000000000..a2d909c32 --- /dev/null +++ b/scanners/findings-processor/cmd/service/main.go @@ -0,0 +1,19 @@ +package main + +import ( + "github.com/canada-ca/tracker/scanners/findings-processor/internal/config" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/runner" + "github.com/rs/zerolog/log" +) + +func main() { + cfg, err := config.Load() + if err != nil { + log.Fatal().Err(err).Msg("failed to load config") + } + config.SetupLogger(cfg) + + if err := runner.Run(cfg); err != nil { + log.Fatal().Err(err).Msg("findings processor failed") + } +} diff --git a/scanners/findings-processor/go.mod b/scanners/findings-processor/go.mod new file mode 100644 index 000000000..ea962654a --- /dev/null +++ b/scanners/findings-processor/go.mod @@ -0,0 +1,79 @@ +module github.com/canada-ca/tracker/scanners/findings-processor + +go 1.25.0 + +require ( + github.com/joho/godotenv v1.5.1 + github.com/kelseyhightower/envconfig v1.4.0 + github.com/nats-io/nats.go v1.52.0 + github.com/rs/zerolog v1.34.0 + github.com/testcontainers/testcontainers-go v0.44.0 + github.com/testcontainers/testcontainers-go/modules/nats v0.44.0 +) + +require ( + dario.cat/mergo v1.0.2 // indirect + github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c // indirect + github.com/Microsoft/go-winio v0.6.2 // indirect + github.com/arangodb/go-velocypack v0.0.0-20200318135517-5af53c29c67e // indirect + github.com/cenkalti/backoff/v4 v4.3.0 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/containerd/errdefs v1.0.0 // indirect + github.com/containerd/errdefs/pkg v0.3.0 // indirect + github.com/containerd/log v0.1.0 // indirect + github.com/containerd/platforms v0.2.1 // indirect + github.com/cpuguy83/dockercfg v0.3.2 // indirect + github.com/davecgh/go-spew v1.1.1 // indirect + github.com/dchest/siphash v1.2.3 // indirect + github.com/distribution/reference v0.6.0 // indirect + github.com/docker/go-connections v0.7.0 // indirect + github.com/docker/go-units v0.5.0 // indirect + github.com/ebitengine/purego v0.10.1 // indirect + github.com/felixge/httpsnoop v1.1.0 // indirect + github.com/go-logr/logr v1.4.3 // indirect + github.com/go-logr/stdr v1.2.2 // indirect + github.com/go-ole/go-ole v1.3.0 // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/kkdai/maglev v0.2.0 // indirect + github.com/lufia/plan9stats v0.0.0-20260330125221-c963978e514e // indirect + github.com/magiconair/properties v1.8.10 // indirect + github.com/mattn/go-colorable v0.1.13 // indirect + github.com/mattn/go-isatty v0.0.20 // indirect + github.com/moby/docker-image-spec v1.3.1 // indirect + github.com/moby/go-archive v0.2.0 // indirect + github.com/moby/moby/api v1.55.0 // indirect + github.com/moby/moby/client v0.5.0 // indirect + github.com/moby/patternmatcher v0.6.1 // indirect + github.com/moby/sys/sequential v0.7.0 // indirect + github.com/moby/sys/user v0.4.0 // indirect + github.com/moby/sys/userns v0.1.0 // indirect + github.com/moby/term v0.5.2 // indirect + github.com/opencontainers/go-digest v1.0.0 // indirect + github.com/opencontainers/image-spec v1.1.1 // indirect + github.com/pkg/errors v0.9.1 // indirect + github.com/pmezard/go-difflib v1.0.0 // indirect + github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect + github.com/shirou/gopsutil/v4 v4.26.6 // indirect + github.com/sirupsen/logrus v1.9.4 // indirect + github.com/stretchr/testify v1.11.1 // indirect + github.com/tklauser/go-sysconf v0.4.0 // indirect + github.com/tklauser/numcpus v0.12.0 // indirect + github.com/yusufpapurcu/wmi v1.2.4 // indirect + go.opentelemetry.io/auto/sdk v1.2.1 // indirect + go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.69.0 // indirect + go.opentelemetry.io/otel v1.44.0 // indirect + go.opentelemetry.io/otel/metric v1.44.0 // indirect + go.opentelemetry.io/otel/trace v1.44.0 // indirect + golang.org/x/net v0.56.0 // indirect + golang.org/x/text v0.40.0 // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect +) + +require ( + github.com/arangodb/go-driver/v2 v2.3.1 + github.com/klauspost/compress v1.18.6 // indirect + github.com/nats-io/nkeys v0.4.15 // indirect + github.com/nats-io/nuid v1.0.1 // indirect + golang.org/x/crypto v0.54.0 // indirect + golang.org/x/sys v0.47.0 // indirect +) diff --git a/scanners/findings-processor/go.sum b/scanners/findings-processor/go.sum new file mode 100644 index 000000000..0bb0805b8 --- /dev/null +++ b/scanners/findings-processor/go.sum @@ -0,0 +1,180 @@ +dario.cat/mergo v1.0.2 h1:85+piFYR1tMbRrLcDwR18y4UKJ3aH1Tbzi24VRW1TK8= +dario.cat/mergo v1.0.2/go.mod h1:E/hbnu0NxMFBjpMIE34DRGLWqDy0g5FuKDhCb31ngxA= +github.com/AdaLogics/go-fuzz-headers v0.0.0-20240806141605-e8a1dd7889d6 h1:He8afgbRMd7mFxO99hRNu+6tazq8nFF9lIwo9JFroBk= +github.com/AdaLogics/go-fuzz-headers v0.0.0-20240806141605-e8a1dd7889d6/go.mod h1:8o94RPi1/7XTJvwPpRSzSUedZrtlirdB3r9Z20bi2f8= +github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c h1:udKWzYgxTojEKWjV8V+WSxDXJ4NFATAsZjh8iIbsQIg= +github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c/go.mod h1:xomTg63KZ2rFqZQzSB4Vz2SUXa1BpHTVz9L5PTmPC4E= +github.com/Microsoft/go-winio v0.6.2 h1:F2VQgta7ecxGYO8k3ZZz3RS8fVIXVxONVUPlNERoyfY= +github.com/Microsoft/go-winio v0.6.2/go.mod h1:yd8OoFMLzJbo9gZq8j5qaps8bJ9aShtEA8Ipt1oGCvU= +github.com/arangodb/go-driver/v2 v2.3.1 h1:km44FjBl6Uh+rD+0mPLLCkOv7QNOTl460Eid7jwQFMQ= +github.com/arangodb/go-driver/v2 v2.3.1/go.mod h1:Mi++s/SLvrrXKmLpAy84ivZ3fuyUBPu7ydmyJo3PMlk= +github.com/arangodb/go-velocypack v0.0.0-20200318135517-5af53c29c67e h1:Xg+hGrY2LcQBbxd0ZFdbGSyRKTYMZCfBbw/pMJFOk1g= +github.com/arangodb/go-velocypack v0.0.0-20200318135517-5af53c29c67e/go.mod h1:mq7Shfa/CaixoDxiyAAc5jZ6CVBAyPaNQCGS7mkj4Ho= +github.com/cenkalti/backoff/v4 v4.3.0 h1:MyRJ/UdXutAwSAT+s3wNd7MfTIcy71VQueUuFK343L8= +github.com/cenkalti/backoff/v4 v4.3.0/go.mod h1:Y3VNntkOUPxTVeUxJ/G5vcM//AlwfmyYozVcomhLiZE= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/containerd/errdefs v1.0.0 h1:tg5yIfIlQIrxYtu9ajqY42W3lpS19XqdxRQeEwYG8PI= +github.com/containerd/errdefs v1.0.0/go.mod h1:+YBYIdtsnF4Iw6nWZhJcqGSg/dwvV7tyJ/kCkyJ2k+M= +github.com/containerd/errdefs/pkg v0.3.0 h1:9IKJ06FvyNlexW690DXuQNx2KA2cUJXx151Xdx3ZPPE= +github.com/containerd/errdefs/pkg v0.3.0/go.mod h1:NJw6s9HwNuRhnjJhM7pylWwMyAkmCQvQ4GpJHEqRLVk= +github.com/containerd/log v0.1.0 h1:TCJt7ioM2cr/tfR8GPbGf9/VRAX8D2B4PjzCpfX540I= +github.com/containerd/log v0.1.0/go.mod h1:VRRf09a7mHDIRezVKTRCrOq78v577GXq3bSa3EhrzVo= +github.com/containerd/platforms v0.2.1 h1:zvwtM3rz2YHPQsF2CHYM8+KtB5dvhISiXh5ZpSBQv6A= +github.com/containerd/platforms v0.2.1/go.mod h1:XHCb+2/hzowdiut9rkudds9bE5yJ7npe7dG/wG+uFPw= +github.com/coreos/go-systemd/v22 v22.5.0/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc= +github.com/cpuguy83/dockercfg v0.3.2 h1:DlJTyZGBDlXqUZ2Dk2Q3xHs/FtnooJJVaad2S9GKorA= +github.com/cpuguy83/dockercfg v0.3.2/go.mod h1:sugsbF4//dDlL/i+S+rtpIWp+5h0BHJHfjj5/jFyUJc= +github.com/creack/pty v1.1.24 h1:bJrF4RRfyJnbTJqzRLHzcGaZK1NeM5kTC9jGgovnR1s= +github.com/creack/pty v1.1.24/go.mod h1:08sCNb52WyoAwi2QDyzUCTgcvVFhUzewun7wtTfvcwE= +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= +github.com/dchest/siphash v1.2.2/go.mod h1:q+IRvb2gOSrUnYoPqHiyHXS0FOBBOdl6tONBlVnOnt4= +github.com/dchest/siphash v1.2.3 h1:QXwFc8cFOR2dSa/gE6o/HokBMWtLUaNDVd+22aKHeEA= +github.com/dchest/siphash v1.2.3/go.mod h1:0NvQU092bT0ipiFN++/rXm69QG9tVxLAlQHIXMPAkHc= +github.com/distribution/reference v0.6.0 h1:0IXCQ5g4/QMHHkarYzh5l+u8T3t73zM5QvfrDyIgxBk= +github.com/distribution/reference v0.6.0/go.mod h1:BbU0aIcezP1/5jX/8MP0YiH4SdvB5Y4f/wlDRiLyi3E= +github.com/docker/go-connections v0.7.0 h1:6SsRfJddP22WMrCkj19x9WKjEDTB+ahsdiGYf0mN39c= +github.com/docker/go-connections v0.7.0/go.mod h1:no1qkHdjq7kLMGUXYAduOhYPSJxxvgWBh7ogVvptn3Q= +github.com/docker/go-units v0.5.0 h1:69rxXcBk27SvSaaxTtLh/8llcHD8vYHT7WSdRZ/jvr4= +github.com/docker/go-units v0.5.0/go.mod h1:fgPhTUdO+D/Jk86RDLlptpiXQzgHJF7gydDDbaIK4Dk= +github.com/ebitengine/purego v0.10.1 h1:dewVBCBT2GaMu1SrNTYxQhgQBethzfhiwvZiLGP/qyY= +github.com/ebitengine/purego v0.10.1/go.mod h1:iIjxzd6CiRiOG0UyXP+V1+jWqUXVjPKLAI0mRfJZTmQ= +github.com/felixge/httpsnoop v1.1.0 h1:3YtUj32ZZkqZtt3sZZsClsymw/QDuVfpNhoA31zeORc= +github.com/felixge/httpsnoop v1.1.0/go.mod h1:Zqxgdd+1Rkcz8euOqdr7lqgCRJztwr5hp9vDSi5UZCE= +github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= +github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= +github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= +github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= +github.com/go-ole/go-ole v1.2.6/go.mod h1:pprOEPIfldk/42T2oK7lQ4v4JSDwmV0As9GaiUsvbm0= +github.com/go-ole/go-ole v1.3.0 h1:Dt6ye7+vXGIKZ7Xtk4s6/xVdGDQynvom7xCFEdWr6uE= +github.com/go-ole/go-ole v1.3.0/go.mod h1:5LS6F96DhAwUc7C+1HLexzMXY1xGRSryjyPPKW6zv78= +github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0= +github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4= +github.com/kelseyhightower/envconfig v1.4.0 h1:Im6hONhd3pLkfDFsbRgu68RDNkGF1r3dvMUtDTo2cv8= +github.com/kelseyhightower/envconfig v1.4.0/go.mod h1:cccZRl6mQpaq41TPp5QxidR+Sa3axMbJDNb//FQX6Gg= +github.com/kkdai/maglev v0.2.0 h1:w6DCW0kAA6fstZqXkrBrlgIC3jeIRXkjOYea/m6EK/Y= +github.com/kkdai/maglev v0.2.0/go.mod h1:d+mt8Lmt3uqi9aRb/BnPjzD0fy+ETs1vVXiGRnqHVZ4= +github.com/klauspost/compress v1.18.6 h1:2jupLlAwFm95+YDR+NwD2MEfFO9d4z4Prjl1XXDjuao= +github.com/klauspost/compress v1.18.6/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/lufia/plan9stats v0.0.0-20260330125221-c963978e514e h1:Q6MvJtQK/iRcRtzAscm/zF23XxJlbECiGPyRicsX+Ak= +github.com/lufia/plan9stats v0.0.0-20260330125221-c963978e514e/go.mod h1:autxFIvghDt3jPTLoqZ9OZ7s9qTGNAWmYCjVFWPX/zg= +github.com/magiconair/properties v1.8.10 h1:s31yESBquKXCV9a/ScB3ESkOjUYYv+X0rg8SYxI99mE= +github.com/magiconair/properties v1.8.10/go.mod h1:Dhd985XPs7jluiymwWYZ0G4Z61jb3vdS329zhj2hYo0= +github.com/mattn/go-colorable v0.1.13 h1:fFA4WZxdEF4tXPZVKMLwD8oUnCTTo08duU7wxecdEvA= +github.com/mattn/go-colorable v0.1.13/go.mod h1:7S9/ev0klgBDR4GtXTXX8a3vIGJpMovkB8vQcUbaXHg= +github.com/mattn/go-isatty v0.0.16/go.mod h1:kYGgaQfpe5nmfYZH+SKPsOc2e4SrIfOl2e/yFXSvRLM= +github.com/mattn/go-isatty v0.0.19/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= +github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +github.com/moby/docker-image-spec v1.3.1 h1:jMKff3w6PgbfSa69GfNg+zN/XLhfXJGnEx3Nl2EsFP0= +github.com/moby/docker-image-spec v1.3.1/go.mod h1:eKmb5VW8vQEh/BAr2yvVNvuiJuY6UIocYsFu/DxxRpo= +github.com/moby/go-archive v0.2.0 h1:zg5QDUM2mi0JIM9fdQZWC7U8+2ZfixfTYoHL7rWUcP8= +github.com/moby/go-archive v0.2.0/go.mod h1:mNeivT14o8xU+5q1YnNrkQVpK+dnNe/K6fHqnTg4qPU= +github.com/moby/moby/api v1.55.0 h1:2/sexvQyqIWS8pRSCFddBfpW2qE7vR7FCL+vN8pxwMc= +github.com/moby/moby/api v1.55.0/go.mod h1:+RQ6wluLwtYaTd1WnPLykIDPekkuyD/ROWQClE83pzs= +github.com/moby/moby/client v0.5.0 h1:5XhyPk2fuOWf6RlSFa3MkIIgDZkF25xToXW8Q/BH7cc= +github.com/moby/moby/client v0.5.0/go.mod h1:rcVpF8ncl9vo5gaIBdol6CnbEtSj1uxMvEV/UrykF/s= +github.com/moby/patternmatcher v0.6.1 h1:qlhtafmr6kgMIJjKJMDmMWq7WLkKIo23hsrpR3x084U= +github.com/moby/patternmatcher v0.6.1/go.mod h1:hDPoyOpDY7OrrMDLaYoY3hf52gNCR/YOUYxkhApJIxc= +github.com/moby/sys/sequential v0.7.0 h1:ASQNGNROJSuOO6LL6bPHbKvuZu6NU8P4ldPWk31zj/8= +github.com/moby/sys/sequential v0.7.0/go.mod h1:NfSTAp6V3fw4tmkD62PEcOKeZKquXT8VKCkf7aVR79o= +github.com/moby/sys/user v0.4.0 h1:jhcMKit7SA80hivmFJcbB1vqmw//wU61Zdui2eQXuMs= +github.com/moby/sys/user v0.4.0/go.mod h1:bG+tYYYJgaMtRKgEmuueC0hJEAZWwtIbZTB+85uoHjs= +github.com/moby/sys/userns v0.1.0 h1:tVLXkFOxVu9A64/yh59slHVv9ahO9UIev4JZusOLG/g= +github.com/moby/sys/userns v0.1.0/go.mod h1:IHUYgu/kao6N8YZlp9Cf444ySSvCmDlmzUcYfDHOl28= +github.com/moby/term v0.5.2 h1:6qk3FJAFDs6i/q3W/pQ97SX192qKfZgGjCQqfCJkgzQ= +github.com/moby/term v0.5.2/go.mod h1:d3djjFCrjnB+fl8NJux+EJzu0msscUP+f8it8hPkFLc= +github.com/nats-io/nats.go v1.52.0 h1:n3avV4VBsCgsdwh71TppsTwtv+QdPs7ntSKM8qJLGsc= +github.com/nats-io/nats.go v1.52.0/go.mod h1:26HypzazeOkyO3/mqd1zZd53STJN0EjCYF9Uy2ZOBno= +github.com/nats-io/nkeys v0.4.15 h1:JACV5jRVO9V856KOapQ7x+EY8Jo3qw1vJt/9Jpwzkk4= +github.com/nats-io/nkeys v0.4.15/go.mod h1:CpMchTXC9fxA5zrMo4KpySxNjiDVvr8ANOSZdiNfUrs= +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/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8Oi/yOhh5U= +github.com/opencontainers/go-digest v1.0.0/go.mod h1:0JzlMkj0TRzQZfJkVvzbP0HBR3IKzErnv2BNG4W4MAM= +github.com/opencontainers/image-spec v1.1.1 h1:y0fUlFfIZhPF1W537XOLg0/fcx6zcHCJwooC2xJA040= +github.com/opencontainers/image-spec v1.1.1/go.mod h1:qpqAh3Dmcf36wStyyWU+kCeDgrGnAve2nCC8+7h8Q0M= +github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= +github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= +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/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 h1:o4JXh1EVt9k/+g42oCprj/FisM4qX9L3sZB3upGN2ZU= +github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55/go.mod h1:OmDBASR4679mdNQnz2pUhc2G8CO2JrUAVFDRBDP/hJE= +github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= +github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= +github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0= +github.com/rs/zerolog v1.34.0 h1:k43nTLIwcTVQAncfCw4KZ2VY6ukYoZaBPNOE8txlOeY= +github.com/rs/zerolog v1.34.0/go.mod h1:bJsvje4Z08ROH4Nhs5iH600c3IkWhwp44iRc54W6wYQ= +github.com/shirou/gopsutil/v4 v4.26.6 h1:Mzr/npDtQC/xpeEuQKHZt8Zo9CmPvhTj8nkR8w5TLDs= +github.com/shirou/gopsutil/v4 v4.26.6/go.mod h1:LZ6ewCSkBqUpvSOf+LsTGnRinC6iaNUNMGBtDkJBaLQ= +github.com/sirupsen/logrus v1.9.4 h1:TsZE7l11zFCLZnZ+teH4Umoq5BhEIfIzfRDZ1Uzql2w= +github.com/sirupsen/logrus v1.9.4/go.mod h1:ftWc9WdOfJ0a92nsE2jF5u5ZwH8Bv2zdeOC42RjbV2g= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/objx v0.5.3 h1:jmXUvGomnU1o3W/V5h2VEradbpJDwGrzugQQvL0POH4= +github.com/stretchr/objx v0.5.3/go.mod h1:rDQraq+vQZU7Fde9LOZLr8Tax6zZvy4kuNKF+QYS+U0= +github.com/stretchr/testify v1.5.1/go.mod h1:5W2xD1RspED5o8YsWQXVCued0rvSQ+mT+I5cxcmMvtA= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/testcontainers/testcontainers-go v0.44.0 h1:/Fwh6HY1mIikhnm9e7HwoxGycx0lzRAE0f5VQpjFxzI= +github.com/testcontainers/testcontainers-go v0.44.0/go.mod h1:IcnwQrYTO86xHXu5bvMaBH7ATlbS3Qn1M1QWW3c66rE= +github.com/testcontainers/testcontainers-go/modules/nats v0.44.0 h1:xGgxnCy6BnmIUUQXQmlYVl7hLx/gwXjJ2S6ccOz+JbA= +github.com/testcontainers/testcontainers-go/modules/nats v0.44.0/go.mod h1:UfIi/50Rj5pl3ixym03CO6kLQL5MIogZnGZj4OTJbh0= +github.com/tklauser/go-sysconf v0.4.0 h1:7H0uAN+7RkwWRaxhYXDLqa5V3LPrJeV8wmD9dRUgPQU= +github.com/tklauser/go-sysconf v0.4.0/go.mod h1:8mTNWyog7H+MpKijp4VmKJAd2bbYQ2zuUwkYRbUArPI= +github.com/tklauser/numcpus v0.12.0 h1:NR85qdvHA9pFse3x3weVZ0r0ST8R6l5RHbZrlRaqob4= +github.com/tklauser/numcpus v0.12.0/go.mod h1:ABHeXzJnr/qqwguhClkZKT1/8VABcYrsyUiUGobwWJg= +github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0= +github.com/yusufpapurcu/wmi v1.2.4/go.mod h1:SBZ9tNy3G9/m5Oi98Zks0QjeHVDvuK0qfxQmPyzfmi0= +go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= +go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= +go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.69.0 h1:8tvICD4vSTOOsNrsI4Ljf6C+6UKvpTEH5XY3JMoyPoo= +go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.69.0/go.mod h1:z9+yiacE0IHRqM4qFfkbt/JYlmYXgss8GY/jXoNuPJI= +go.opentelemetry.io/otel v1.44.0 h1:JjwHmHpA4iZ3wBxluu2fbbE7j4kqlE8jXyAyPXH7HqU= +go.opentelemetry.io/otel v1.44.0/go.mod h1:BMgjTHL9WPRlRjL2oZCBTL4whCGtXch2H4BhOPIAyYc= +go.opentelemetry.io/otel/metric v1.44.0 h1:1w0gILTcHdr3YI+ixLyjemwrVnsMURbTZFrSYCdDdmc= +go.opentelemetry.io/otel/metric v1.44.0/go.mod h1:8O7hanEPBNgEMmybD3s2VBKcgWOCsA6tzHBPODAiquo= +go.opentelemetry.io/otel/sdk v1.44.0 h1:nHYwb9lK+fJPU/dnT6s7W7Z8itMWyqrnVfbheVYrZ58= +go.opentelemetry.io/otel/sdk v1.44.0/go.mod h1:Osuydd3Se74nqjAKxid74N5eC+jfEqfTegHRnq58oK0= +go.opentelemetry.io/otel/sdk/metric v1.44.0 h1:3LlKgI+VjbVsjNRFZJZAJ30WjXC5VkNRks6si09iEfI= +go.opentelemetry.io/otel/sdk/metric v1.44.0/go.mod h1:5B5pMARnXxKhltooO4xUuCBorl65a4EpnTalObqOigA= +go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/6gtIk= +go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE= +golang.org/x/crypto v0.54.0 h1:YLIA59K4fiNzHzjnZt2tUJQjQtUWfWbeHBqKtk3eScw= +golang.org/x/crypto v0.54.0/go.mod h1:KWL8ny2AZdGR2cWmzeHrp2azQPGogOv+HeQaVEXC2dk= +golang.org/x/net v0.56.0 h1:Rw8j/hFzGvJUZwNBXnAtf5sVDVt+65SK2C7IxCxZt5o= +golang.org/x/net v0.56.0/go.mod h1:D3Ku6r+V6JROoZK144D2XfMHFcMq/0zSfLelVTCFKec= +golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20201204225414-ed752295db88/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20210616094352-59db8d763f22/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.1.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.12.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= +golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/term v0.45.0 h1:NwWyBmoJCbfTHpxrWoZ9C6/VxOf7ic219I8xZZFdrf0= +golang.org/x/term v0.45.0/go.mod h1:9aqxs0blBcrm/n0L9QW0aRVD+ktan8ssZromtqJC43w= +golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs= +golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY= +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.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gotest.tools/v3 v3.5.2 h1:7koQfIKdy+I8UTetycgUqXWSDwpgv193Ka+qRsmBY8Q= +gotest.tools/v3 v3.5.2/go.mod h1:LtdLGcnqToBH83WByAAi/wiwSFCArdFIUV/xxN4pcjA= +pgregory.net/rapid v1.2.0 h1:keKAYRcjm+e1F0oAuU5F5+YPAWcyxNNRK2wud503Gnk= +pgregory.net/rapid v1.2.0/go.mod h1:PY5XlDGj0+V1FCq0o192FdRhpKHGTRIWBgqjDBTrq04= diff --git a/scanners/findings-processor/internal/config/config.go b/scanners/findings-processor/internal/config/config.go new file mode 100644 index 000000000..f97de1e5b --- /dev/null +++ b/scanners/findings-processor/internal/config/config.go @@ -0,0 +1,58 @@ +package config + +import ( + "errors" + "os" + "strings" + "time" + + "github.com/joho/godotenv" + "github.com/kelseyhightower/envconfig" + "github.com/rs/zerolog" + "github.com/rs/zerolog/log" +) + +type Config struct { + NATSURL string `envconfig:"NATS_URL" default:"nats://localhost:4222"` + NATSStream string `envconfig:"NATS_STREAM" default:"SCANS"` + NATSSubject string `envconfig:"NATS_SUBJECT" default:"scans.findings.*"` + NATSDurable string `envconfig:"NATS_CONSUMER_DURABLE" default:"findings-processor"` + NATSAckWait time.Duration `envconfig:"NATS_ACK_WAIT" default:"30s"` + NATSMaxDeliver int `envconfig:"NATS_MAX_DELIVER" default:"10"` + NATSMaxPending int `envconfig:"NATS_MAX_ACK_PENDING" default:"256"` + + DBURL string `envconfig:"DB_URL" default:"http://localhost:8529"` + DBUser string `envconfig:"DB_USER"` + DBName string `envconfig:"DB_NAME"` + DBPassword string `envconfig:"DB_PASS"` + + LogLevel string `envconfig:"LOG_LEVEL" default:"info"` + PrettyLogOutput bool `envconfig:"LOG_PRETTY" default:"true"` +} + +func Load() (Config, error) { + if err := godotenv.Load(); err != nil && !errors.Is(err, os.ErrNotExist) { + return Config{}, err + } + + var cfg Config + if err := envconfig.Process("", &cfg); err != nil { + return Config{}, err + } + + return cfg, nil +} + +func SetupLogger(cfg Config) { + level, err := zerolog.ParseLevel(strings.ToLower(cfg.LogLevel)) + if err != nil { + level = zerolog.InfoLevel + } + + zerolog.SetGlobalLevel(level) + if cfg.PrettyLogOutput { + log.Logger = log.Output(zerolog.ConsoleWriter{Out: os.Stdout, TimeFormat: time.RFC3339}) + } + + log.Info().Str("logLevel", zerolog.GlobalLevel().String()).Msg("logger initialized") +} diff --git a/scanners/findings-processor/internal/config/config_test.go b/scanners/findings-processor/internal/config/config_test.go new file mode 100644 index 000000000..8b936d1ea --- /dev/null +++ b/scanners/findings-processor/internal/config/config_test.go @@ -0,0 +1,82 @@ +package config + +import ( + "testing" + "time" +) + +func TestLoad(t *testing.T) { + // Load() calls godotenv.Load(), which reads a .env file from the working + // directory if one exists. None is present in this package during tests, + // so envconfig defaults/overrides are what's under test here. + + t.Run("defaults when no env vars set", func(t *testing.T) { + cfg, err := Load() + if err != nil { + t.Fatalf("Load() error = %v", err) + } + if cfg.NATSURL != "nats://localhost:4222" { + t.Errorf("NATSURL = %q, want default", cfg.NATSURL) + } + if cfg.NATSAckWait != 30*time.Second { + t.Errorf("NATSAckWait = %v, want 30s", cfg.NATSAckWait) + } + if cfg.NATSMaxDeliver != 10 { + t.Errorf("NATSMaxDeliver = %d, want 10", cfg.NATSMaxDeliver) + } + if !cfg.PrettyLogOutput { + t.Error("PrettyLogOutput = false, want true (default)") + } + }) + + t.Run("env vars override defaults", func(t *testing.T) { + t.Setenv("NATS_URL", "nats://example.com:4222") + t.Setenv("NATS_MAX_DELIVER", "3") + t.Setenv("DB_NAME", "findings") + t.Setenv("LOG_PRETTY", "false") + + cfg, err := Load() + if err != nil { + t.Fatalf("Load() error = %v", err) + } + if cfg.NATSURL != "nats://example.com:4222" { + t.Errorf("NATSURL = %q, want override", cfg.NATSURL) + } + if cfg.NATSMaxDeliver != 3 { + t.Errorf("NATSMaxDeliver = %d, want 3", cfg.NATSMaxDeliver) + } + if cfg.DBName != "findings" { + t.Errorf("DBName = %q, want %q", cfg.DBName, "findings") + } + if cfg.PrettyLogOutput { + t.Error("PrettyLogOutput = true, want false (override)") + } + }) + + t.Run("invalid duration returns error", func(t *testing.T) { + t.Setenv("NATS_ACK_WAIT", "not-a-duration") + + if _, err := Load(); err == nil { + t.Fatal("Load() error = nil, want error for invalid NATS_ACK_WAIT") + } + }) +} + +func TestSetupLogger(t *testing.T) { + tests := []struct { + name string + logLevel string + }{ + {name: "valid level", logLevel: "warn"}, + {name: "invalid level falls back to info", logLevel: "not-a-level"}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + // SetupLogger mutates the global zerolog level/writer as a side + // effect; assert it doesn't panic and completes for both a valid + // and an invalid level. + SetupLogger(Config{LogLevel: tt.logLevel, PrettyLogOutput: true}) + }) + } +} diff --git a/scanners/findings-processor/internal/database/client.go b/scanners/findings-processor/internal/database/client.go new file mode 100644 index 000000000..658459b53 --- /dev/null +++ b/scanners/findings-processor/internal/database/client.go @@ -0,0 +1,19 @@ +package database + +import ( + "github.com/arangodb/go-driver/v2/arangodb" + "github.com/arangodb/go-driver/v2/connection" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/config" +) + +func CreateDBClient(cfg config.Config) (arangodb.Client, error) { + endpoint := connection.NewRoundRobinEndpoints([]string{cfg.DBURL}) + conn := connection.NewHttp2Connection(connection.DefaultHTTP2ConfigurationWrapper(endpoint, false)) + + auth := connection.NewBasicAuth(cfg.DBUser, cfg.DBPassword) + if err := conn.SetAuthentication(auth); err != nil { + return nil, err + } + + return arangodb.NewClient(conn), nil +} diff --git a/scanners/findings-processor/internal/database/client_test.go b/scanners/findings-processor/internal/database/client_test.go new file mode 100644 index 000000000..d138c7a7e --- /dev/null +++ b/scanners/findings-processor/internal/database/client_test.go @@ -0,0 +1,26 @@ +package database + +import ( + "testing" + + "github.com/canada-ca/tracker/scanners/findings-processor/internal/config" +) + +func TestCreateDBClient(t *testing.T) { + t.Parallel() + + // CreateDBClient only builds the client/auth wiring; it doesn't dial the + // server, so this exercises the real construction path without a live + // ArangoDB instance. + client, err := CreateDBClient(config.Config{ + DBURL: "http://localhost:8529", + DBUser: "root", + DBPassword: "secret", + }) + if err != nil { + t.Fatalf("CreateDBClient() error = %v", err) + } + if client == nil { + t.Fatal("CreateDBClient() returned nil client") + } +} diff --git a/scanners/findings-processor/internal/database/upsert.go b/scanners/findings-processor/internal/database/upsert.go new file mode 100644 index 000000000..82e8ba5ce --- /dev/null +++ b/scanners/findings-processor/internal/database/upsert.go @@ -0,0 +1,81 @@ +package database + +import ( + "context" + + "github.com/arangodb/go-driver/v2/arangodb" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/model" + "github.com/rs/zerolog/log" +) + +type findingUpdatePatch struct { + LastSeen string `json:"lastSeen"` + OccurrenceCount int `json:"occurrenceCount"` +} + +func readFinding(ctx context.Context, db arangodb.Database, key string) (*model.FindingDocument, error) { + var finding model.FindingDocument + options := arangodb.QueryOptions{ + Count: true, + BindVars: map[string]interface{}{ + "key": key, + }, + } + query := "FOR f IN additionalFindings FILTER f.findingKey == @key LIMIT 1 RETURN f" + + cursor, err := db.Query(ctx, query, &options) + if err != nil { + return nil, err + } + defer cursor.Close() + + if cursor.Count() == 0 { + return nil, nil + } + + if _, err := cursor.ReadDocument(ctx, &finding); err != nil { + return nil, err + } + return &finding, nil +} + +func UpsertFinding(ctx context.Context, db arangodb.Database, evt model.FindingEvent) error { + findingsCol, err := db.GetCollection(ctx, "additionalFindings", nil) + if err != nil { + log.Warn().Err(err).Msg("failed to find collection") + return err + } + + key := evt.DeriveFindingKey() + + finding, err := readFinding(ctx, db, key) + if err != nil { + log.Warn().Err(err).Msg("failed to check finding existence") + return err + } + + if finding != nil { + patch := findingUpdatePatch{ + LastSeen: evt.ObservedAt, + OccurrenceCount: finding.OccurrenceCount + 1, + } + + if _, err := findingsCol.UpdateDocument(ctx, finding.Key, patch); err != nil { + log.Warn().Err(err).Msg("failed to update finding doc") + return err + } + } else { + newDoc, err := model.NewFindingDocumentFromEvent(evt) + if err != nil { + log.Warn().Err(err).Msg("failed to create doc") + return err + } + + if _, err := findingsCol.CreateDocument(ctx, newDoc); err != nil { + log.Warn().Err(err).Msg("failed to create doc") + return err + } + } + + return nil +} diff --git a/scanners/findings-processor/internal/database/upsert_integration_test.go b/scanners/findings-processor/internal/database/upsert_integration_test.go new file mode 100644 index 000000000..846a40118 --- /dev/null +++ b/scanners/findings-processor/internal/database/upsert_integration_test.go @@ -0,0 +1,205 @@ +//go:build integration + +package database_test + +import ( + "context" + "fmt" + "testing" + "time" + + "github.com/arangodb/go-driver/v2/arangodb" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/config" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/database" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/model" + "github.com/testcontainers/testcontainers-go" + "github.com/testcontainers/testcontainers-go/wait" +) + +const rootPassword = "test-password" + +func startArangoDB(t *testing.T) arangodb.Database { + t.Helper() + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute) + defer cancel() + + req := testcontainers.ContainerRequest{ + Image: "arangodb:3.11", + ExposedPorts: []string{"8529/tcp"}, + Env: map[string]string{"ARANGO_ROOT_PASSWORD": rootPassword}, + WaitingFor: wait.ForListeningPort("8529/tcp").WithStartupTimeout(90 * time.Second), + } + container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{ + ContainerRequest: req, + Started: true, + }) + if err != nil { + t.Fatalf("failed to start arangodb container: %v", err) + } + t.Cleanup(func() { + if err := container.Terminate(context.Background()); err != nil { + t.Logf("failed to terminate arangodb container: %v", err) + } + }) + + host, err := container.Host(ctx) + if err != nil { + t.Fatalf("failed to get container host: %v", err) + } + port, err := container.MappedPort(ctx, "8529") + if err != nil { + t.Fatalf("failed to get mapped port: %v", err) + } + + dbName := "findings_test" + cfg := config.Config{ + DBURL: fmt.Sprintf("http://%s:%s", host, port.Port()), + DBUser: "root", + DBPassword: rootPassword, + DBName: dbName, + } + + client, err := database.CreateDBClient(cfg) + if err != nil { + t.Fatalf("CreateDBClient() error = %v", err) + } + + // The server takes a moment after the port opens before it accepts requests. + var db arangodb.Database + deadline := time.Now().Add(30 * time.Second) + for { + db, err = client.CreateDatabase(ctx, dbName, nil) + if err == nil { + break + } + if time.Now().After(deadline) { + t.Fatalf("CreateDatabase() error = %v", err) + } + time.Sleep(500 * time.Millisecond) + } + + if _, err := db.CreateCollectionV2(ctx, "additionalFindings", &arangodb.CreateCollectionPropertiesV2{}); err != nil { + t.Fatalf("CreateCollection() error = %v", err) + } + + return db +} + +func TestIntegration_UpsertFinding_CreatesThenUpdates(t *testing.T) { + db := startArangoDB(t) + ctx := context.Background() + + evt := model.FindingEvent{ + Source: "scanner", + FindingType: "tls-weak", + DomainKey: "example-domain", + Subject: "example.com", + Confidence: "high", + ReasonCode: "weak-cipher", + ObservedAt: "2024-01-01T00:00:00Z", + } + + if err := database.UpsertFinding(ctx, db, evt); err != nil { + t.Fatalf("UpsertFinding() first call error = %v", err) + } + + col, err := db.GetCollection(ctx, "additionalFindings", nil) + if err != nil { + t.Fatalf("GetCollection() error = %v", err) + } + + query := "FOR f IN additionalFindings FILTER f.findingKey == @key LIMIT 1 RETURN f" + readOne := func() model.FindingDocument { + cursor, err := db.Query(ctx, query, &arangodb.QueryOptions{ + BindVars: map[string]interface{}{"key": evt.DeriveFindingKey()}, + }) + if err != nil { + t.Fatalf("Query() error = %v", err) + } + defer cursor.Close() + + var doc model.FindingDocument + if _, err := cursor.ReadDocument(ctx, &doc); err != nil { + t.Fatalf("ReadDocument() error = %v", err) + } + return doc + } + + first := readOne() + if first.OccurrenceCount != 1 { + t.Errorf("OccurrenceCount after create = %d, want 1", first.OccurrenceCount) + } + if first.Status != "active" { + t.Errorf("Status after create = %q, want %q", first.Status, "active") + } + + // Second event for the same domain/source/type/subject/reason should + // update the existing document rather than create a new one. + evt.ObservedAt = "2024-01-02T00:00:00Z" + if err := database.UpsertFinding(ctx, db, evt); err != nil { + t.Fatalf("UpsertFinding() second call error = %v", err) + } + + second := readOne() + if second.OccurrenceCount != 2 { + t.Errorf("OccurrenceCount after update = %d, want 2", second.OccurrenceCount) + } + if second.Key != first.Key { + t.Errorf("update created a new document: first key %q, second key %q", first.Key, second.Key) + } + if second.LastSeen.Format(time.RFC3339) != "2024-01-02T00:00:00Z" { + t.Errorf("LastSeen = %v, want 2024-01-02T00:00:00Z", second.LastSeen) + } + + count, err := col.Count(ctx) + if err != nil { + t.Fatalf("Count() error = %v", err) + } + if count != 1 { + t.Errorf("collection document count = %d, want 1 (update, not duplicate insert)", count) + } +} + +func TestIntegration_UpsertFinding_MissingCollectionReturnsError(t *testing.T) { + db := startArangoDB(t) + ctx := context.Background() + + // Drop the collection UpsertFinding expects, to exercise the error path. + col, err := db.GetCollection(ctx, "additionalFindings", nil) + if err != nil { + t.Fatalf("GetCollection() error = %v", err) + } + if err := col.Remove(ctx); err != nil { + t.Fatalf("Remove() error = %v", err) + } + + err = database.UpsertFinding(ctx, db, model.FindingEvent{ + DomainKey: "d", + Source: "s", + FindingType: "t", + Subject: "sub", + ObservedAt: "2024-01-01T00:00:00Z", + }) + if err == nil { + t.Fatal("UpsertFinding() error = nil, want error for missing collection") + } +} + +func TestIntegration_UpsertFinding_InvalidObservedAtReturnsError(t *testing.T) { + db := startArangoDB(t) + ctx := context.Background() + + // No existing document for this key, so UpsertFinding takes the create + // path, where NewFindingDocumentFromEvent fails to parse ObservedAt. + err := database.UpsertFinding(ctx, db, model.FindingEvent{ + DomainKey: "d", + Source: "s", + FindingType: "t", + Subject: "sub", + ObservedAt: "not-a-timestamp", + }) + if err == nil { + t.Fatal("UpsertFinding() error = nil, want error for an unparsable ObservedAt") + } +} diff --git a/scanners/findings-processor/internal/model/finding.go b/scanners/findings-processor/internal/model/finding.go new file mode 100644 index 000000000..2f0156be1 --- /dev/null +++ b/scanners/findings-processor/internal/model/finding.go @@ -0,0 +1,118 @@ +package model + +import ( + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "strings" + "time" +) + +func ParseEvent(payload []byte) (FindingEvent, error) { + var event FindingEvent + if err := json.Unmarshal(payload, &event); err != nil { + return FindingEvent{}, err + } + + return event, nil +} + +type FindingEvent struct { + Source string `json:"source"` + FindingType string `json:"findingType"` + DomainKey string `json:"domainKey"` + Subject string `json:"subject"` + Confidence string `json:"confidence"` + Severity string `json:"severity,omitempty"` + ReasonCode string `json:"reasonCode,omitempty"` + ObservedAt string `json:"observedAt"` + Evidence map[string]any `json:"evidence,omitempty"` + Attributes map[string]any `json:"attributes,omitempty"` +} + +type FindingDocument struct { + Key string `json:"_key,omitempty"` + FindingKey string `json:"findingKey"` + Domain string `json:"domain"` + Source string `json:"source"` + FindingType string `json:"findingType"` + DomainKey string `json:"domainKey"` + Subject string `json:"subject"` + Confidence string `json:"confidence"` + Severity string `json:"severity,omitempty"` + ReasonCode string `json:"reasonCode,omitempty"` + FirstSeen time.Time `json:"firstSeen"` + LastSeen time.Time `json:"lastSeen"` + Evidence map[string]any `json:"evidence,omitempty"` + Attributes map[string]any `json:"attributes,omitempty"` + OccurrenceCount int `json:"occurrenceCount"` + Raw map[string]any `json:"raw"` + Status string `json:"status"` +} + +func (e FindingEvent) DeriveFindingKey() string { + keyArgs := []string{e.DomainKey, e.Source, e.FindingType, e.Subject, e.ReasonCode} + data := strings.Join(keyArgs, "\x00") + h := sha256.Sum256([]byte(data)) + return hex.EncodeToString(h[:]) +} + +func NewFindingDocumentFromEvent(e FindingEvent) (FindingDocument, error) { + observedAt, err := time.Parse(time.RFC3339, e.ObservedAt) + if err != nil { + return FindingDocument{}, err + } + + domain := fmt.Sprintf("domains/%s", e.DomainKey) + + return FindingDocument{ + FindingKey: e.DeriveFindingKey(), + Domain: domain, + DomainKey: e.DomainKey, + Source: e.Source, + FindingType: e.FindingType, + Subject: e.Subject, + Status: "active", + Confidence: e.Confidence, + Severity: e.Severity, + ReasonCode: e.ReasonCode, + FirstSeen: observedAt, + LastSeen: observedAt, + OccurrenceCount: 1, + Evidence: nonNilMap(e.Evidence), + Attributes: nonNilMap(e.Attributes), + Raw: eventToMap(e), + }, nil +} + +func nonNilMap(m map[string]any) map[string]any { + if m == nil { + return map[string]any{} + } + return m +} + +func eventToMap(e FindingEvent) map[string]any { + b, _ := json.Marshal(e) + out := map[string]any{} + _ = json.Unmarshal(b, &out) + return out +} + +func Validate(e FindingEvent) error { + required := []string{ + e.Source, e.FindingType, + e.DomainKey, e.Subject, e.Confidence, e.ObservedAt, + } + for _, v := range required { + if strings.TrimSpace(v) == "" { + return errors.New("missing required fields") + } + } + if _, err := time.Parse(time.RFC3339, e.ObservedAt); err != nil { + return errors.New("observedAt must be RFC3339") + } + return nil +} diff --git a/scanners/findings-processor/internal/model/finding_test.go b/scanners/findings-processor/internal/model/finding_test.go new file mode 100644 index 000000000..27fea5bcb --- /dev/null +++ b/scanners/findings-processor/internal/model/finding_test.go @@ -0,0 +1,223 @@ +package model + +import ( + "reflect" + "strings" + "testing" + "time" +) + +func TestParseEvent(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + payload string + want FindingEvent + wantErr bool + }{ + { + name: "valid payload", + payload: `{"source":"scanner","findingType":"tls-weak","domainKey":"abc","subject":"example.com","confidence":"high","observedAt":"2024-01-01T00:00:00Z"}`, + want: FindingEvent{ + Source: "scanner", + FindingType: "tls-weak", + DomainKey: "abc", + Subject: "example.com", + Confidence: "high", + ObservedAt: "2024-01-01T00:00:00Z", + }, + }, + { + name: "malformed json", + payload: `{"source":`, + wantErr: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + got, err := ParseEvent([]byte(tt.payload)) + if (err != nil) != tt.wantErr { + t.Fatalf("ParseEvent() error = %v, wantErr %v", err, tt.wantErr) + } + if tt.wantErr { + return + } + if !reflect.DeepEqual(got, tt.want) { + t.Errorf("ParseEvent() = %+v, want %+v", got, tt.want) + } + }) + } +} + +func TestFindingEvent_DeriveFindingKey(t *testing.T) { + t.Parallel() + + base := FindingEvent{ + DomainKey: "domain", + Source: "scanner", + FindingType: "tls-weak", + Subject: "example.com", + ReasonCode: "weak-cipher", + } + + key := base.DeriveFindingKey() + if key == "" { + t.Fatal("DeriveFindingKey() returned empty string") + } + if len(key) != 64 { + t.Errorf("DeriveFindingKey() len = %d, want 64 (sha256 hex)", len(key)) + } + + t.Run("stable across calls", func(t *testing.T) { + t.Parallel() + if got := base.DeriveFindingKey(); got != key { + t.Errorf("DeriveFindingKey() = %q, want %q", got, key) + } + }) + + t.Run("changes when a field changes", func(t *testing.T) { + t.Parallel() + other := base + other.Subject = "other.example.com" + if got := other.DeriveFindingKey(); got == key { + t.Error("DeriveFindingKey() did not change when Subject changed") + } + }) +} + +func TestNewFindingDocumentFromEvent(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + event FindingEvent + wantErr bool + }{ + { + name: "valid event populates document", + event: FindingEvent{ + Source: "scanner", + FindingType: "tls-weak", + DomainKey: "abc", + Subject: "example.com", + Confidence: "high", + Severity: "medium", + ReasonCode: "weak-cipher", + ObservedAt: "2024-01-01T00:00:00Z", + }, + }, + { + name: "nil evidence and attributes become empty maps", + event: FindingEvent{ + Source: "scanner", + FindingType: "tls-weak", + DomainKey: "abc", + Subject: "example.com", + Confidence: "high", + ObservedAt: "2024-01-01T00:00:00Z", + }, + }, + { + name: "invalid observedAt", + event: FindingEvent{ + DomainKey: "abc", + ObservedAt: "not-a-timestamp", + }, + wantErr: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + doc, err := NewFindingDocumentFromEvent(tt.event) + if (err != nil) != tt.wantErr { + t.Fatalf("NewFindingDocumentFromEvent() error = %v, wantErr %v", err, tt.wantErr) + } + if tt.wantErr { + return + } + + wantObservedAt, _ := time.Parse(time.RFC3339, tt.event.ObservedAt) + if !doc.FirstSeen.Equal(wantObservedAt) || !doc.LastSeen.Equal(wantObservedAt) { + t.Errorf("FirstSeen/LastSeen = %v/%v, want %v", doc.FirstSeen, doc.LastSeen, wantObservedAt) + } + if doc.FindingKey != tt.event.DeriveFindingKey() { + t.Errorf("FindingKey = %q, want %q", doc.FindingKey, tt.event.DeriveFindingKey()) + } + if doc.Domain != "domains/"+tt.event.DomainKey { + t.Errorf("Domain = %q, want %q", doc.Domain, "domains/"+tt.event.DomainKey) + } + if doc.Status != "active" { + t.Errorf("Status = %q, want %q", doc.Status, "active") + } + if doc.OccurrenceCount != 1 { + t.Errorf("OccurrenceCount = %d, want 1", doc.OccurrenceCount) + } + if doc.Evidence == nil || doc.Attributes == nil { + t.Error("Evidence/Attributes should never be nil") + } + if doc.Raw == nil { + t.Error("Raw should be populated from the source event") + } + }) + } +} + +func TestValidate(t *testing.T) { + t.Parallel() + + valid := FindingEvent{ + Source: "scanner", + FindingType: "tls-weak", + DomainKey: "abc", + Subject: "example.com", + Confidence: "high", + ObservedAt: "2024-01-01T00:00:00Z", + } + + tests := []struct { + name string + mutate func(e FindingEvent) FindingEvent + wantErr bool + }{ + { + name: "valid event", + mutate: func(e FindingEvent) FindingEvent { return e }, + }, + { + name: "missing source", + mutate: func(e FindingEvent) FindingEvent { e.Source = ""; return e }, + wantErr: true, + }, + { + name: "blank subject (whitespace only)", + mutate: func(e FindingEvent) FindingEvent { e.Subject = " "; return e }, + wantErr: true, + }, + { + name: "observedAt not RFC3339", + mutate: func(e FindingEvent) FindingEvent { e.ObservedAt = "2024-01-01"; return e }, + wantErr: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + err := Validate(tt.mutate(valid)) + if (err != nil) != tt.wantErr { + t.Fatalf("Validate() error = %v, wantErr %v", err, tt.wantErr) + } + if err != nil && !strings.Contains(err.Error(), "required") && !strings.Contains(err.Error(), "RFC3339") { + t.Errorf("Validate() error message %q not descriptive", err.Error()) + } + }) + } +} diff --git a/scanners/findings-processor/internal/runner/processor.go b/scanners/findings-processor/internal/runner/processor.go new file mode 100644 index 000000000..7d76272ba --- /dev/null +++ b/scanners/findings-processor/internal/runner/processor.go @@ -0,0 +1,142 @@ +package runner + +import ( + "context" + "fmt" + "os/signal" + "syscall" + "time" + + "github.com/arangodb/go-driver/v2/arangodb" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/config" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/database" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/model" + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" + "github.com/rs/zerolog/log" +) + +func Run(cfg config.Config) error { + nc, err := nats.Connect( + cfg.NATSURL, + nats.MaxReconnects(-1), + nats.ReconnectHandler(func(c *nats.Conn) { + log.Info().Str("url", c.ConnectedUrl()).Msg("nats reconnected") + }), + nats.DisconnectErrHandler(func(c *nats.Conn, err error) { + log.Warn().Err(err).Msg("nats disconnected") + }), + nats.ClosedHandler(func(c *nats.Conn) { + log.Info().Msg("nats connection closed") + }), + ) + if err != nil { + return fmt.Errorf("failed to connect to NATS: %w", err) + } + defer nc.Close() + + js, err := jetstream.New(nc) + if err != nil { + return fmt.Errorf("failed to create JetStream context: %w", err) + } + + ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer stop() + + client, err := database.CreateDBClient(cfg) + if err != nil { + return fmt.Errorf("failed to create ArangoDB client: %w", err) + } + + dbCtx, cancelDB := context.WithTimeout(ctx, 10*time.Second) + defer cancelDB() + + db, err := client.GetDatabase(dbCtx, cfg.DBName, nil) + if err != nil { + return fmt.Errorf("get database failed: %w", err) + } + + cons, err := js.CreateOrUpdateConsumer(ctx, cfg.NATSStream, jetstream.ConsumerConfig{ + Durable: cfg.NATSDurable, + AckPolicy: jetstream.AckExplicitPolicy, + AckWait: cfg.NATSAckWait, + MaxDeliver: cfg.NATSMaxDeliver, + MaxAckPending: cfg.NATSMaxPending, + FilterSubject: cfg.NATSSubject, + }) + if err != nil { + return fmt.Errorf("create/update consumer failed: %w", err) + } + + handler := func(msg jetstream.Msg) { + upsertCtx, cancelUpsert := context.WithTimeout(context.Background(), 5*time.Second) + defer cancelUpsert() + + var ackErr error + switch HandleEvent(upsertCtx, db, msg.Data()) { + case "ack": + ackErr = msg.Ack() + case "nak": + ackErr = msg.Nak() + case "term": + ackErr = msg.Term() + default: + ackErr = msg.Nak() + } + if ackErr != nil { + log.Warn().Err(ackErr).Msg("failed to ack/nak/term message") + } + } + + consumeCtx, err := cons.Consume(handler, jetstream.ConsumeErrHandler(func(_ jetstream.ConsumeContext, err error) { + log.Warn().Err(err).Msg("consume error") + })) + if err != nil { + return fmt.Errorf("failed to create consumer context: %w", err) + } + + log.Info(). + Str("stream", cfg.NATSStream). + Str("subject", cfg.NATSSubject). + Str("durable", cfg.NATSDurable). + Msg("findings processor started") + + <-ctx.Done() + log.Info().Msg("shutdown signal received, draining consumer") + + consumeCtx.Drain() + select { + case <-consumeCtx.Closed(): + log.Info().Msg("consumer drained") + case <-time.After(15 * time.Second): + log.Warn().Msg("timed out waiting for consumer to drain") + } + + return nil +} + +func HandleEvent(ctx context.Context, db arangodb.Database, payload []byte) string { + evt, err := model.ParseEvent(payload) + if err != nil { + log.Warn().Err(err).Msg("invalid json payload") + return "term" + } + + if err := model.Validate(evt); err != nil { + log.Warn().Err(err).Msg("invalid event payload") + return "term" + } + + if err := database.UpsertFinding(ctx, db, evt); err != nil { + log.Warn().Err(err).Msg("failed to upsert finding") + return "nak" + } + + log.Info(). + Str("source", evt.Source). + Str("findingType", evt.FindingType). + Str("domainKey", evt.DomainKey). + Str("subject", evt.Subject). + Msg("received finding event") + return "ack" +} diff --git a/scanners/findings-processor/internal/runner/processor_integration_test.go b/scanners/findings-processor/internal/runner/processor_integration_test.go new file mode 100644 index 000000000..a02be6d93 --- /dev/null +++ b/scanners/findings-processor/internal/runner/processor_integration_test.go @@ -0,0 +1,124 @@ +//go:build integration + +package runner_test + +import ( + "context" + "fmt" + "testing" + "time" + + "github.com/arangodb/go-driver/v2/arangodb" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/config" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/database" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/runner" + "github.com/testcontainers/testcontainers-go" + "github.com/testcontainers/testcontainers-go/wait" +) + +const rootPassword = "test-password" + +func startArangoDB(t *testing.T) (arangodb.Database, config.Config) { + t.Helper() + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute) + defer cancel() + + req := testcontainers.ContainerRequest{ + Image: "arangodb:3.11", + ExposedPorts: []string{"8529/tcp"}, + Env: map[string]string{"ARANGO_ROOT_PASSWORD": rootPassword}, + WaitingFor: wait.ForListeningPort("8529/tcp").WithStartupTimeout(90 * time.Second), + } + container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{ + ContainerRequest: req, + Started: true, + }) + if err != nil { + t.Fatalf("failed to start arangodb container: %v", err) + } + t.Cleanup(func() { + if err := container.Terminate(context.Background()); err != nil { + t.Logf("failed to terminate arangodb container: %v", err) + } + }) + + host, err := container.Host(ctx) + if err != nil { + t.Fatalf("failed to get container host: %v", err) + } + port, err := container.MappedPort(ctx, "8529") + if err != nil { + t.Fatalf("failed to get mapped port: %v", err) + } + + dbName := "findings_test" + cfg := config.Config{ + DBURL: fmt.Sprintf("http://%s:%s", host, port.Port()), + DBUser: "root", + DBPassword: rootPassword, + DBName: dbName, + } + + client, err := database.CreateDBClient(cfg) + if err != nil { + t.Fatalf("CreateDBClient() error = %v", err) + } + + var db arangodb.Database + deadline := time.Now().Add(30 * time.Second) + for { + db, err = client.CreateDatabase(ctx, dbName, nil) + if err == nil { + break + } + if time.Now().After(deadline) { + t.Fatalf("CreateDatabase() error = %v", err) + } + time.Sleep(500 * time.Millisecond) + } + + if _, err := db.CreateCollectionV2(ctx, "additionalFindings", &arangodb.CreateCollectionPropertiesV2{}); err != nil { + t.Fatalf("CreateCollection() error = %v", err) + } + + return db, cfg +} + +func TestIntegration_HandleEvent(t *testing.T) { + db, _ := startArangoDB(t) + + tests := []struct { + name string + payload string + want string + }{ + { + name: "valid event is acked and persisted", + payload: `{"source":"scanner","findingType":"tls-weak","domainKey":"d1","subject":"example.com","confidence":"high","observedAt":"2024-01-01T00:00:00Z"}`, + want: "ack", + }, + { + name: "malformed json is termed", + payload: `not json`, + want: "term", + }, + { + name: "missing required field is termed", + payload: `{"source":"scanner","findingType":"tls-weak","domainKey":"d1","subject":"example.com","observedAt":"2024-01-01T00:00:00Z"}`, + want: "term", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + + got := runner.HandleEvent(ctx, db, []byte(tt.payload)) + if got != tt.want { + t.Errorf("HandleEvent() = %q, want %q", got, tt.want) + } + }) + } +} diff --git a/scanners/findings-processor/internal/runner/processor_test.go b/scanners/findings-processor/internal/runner/processor_test.go new file mode 100644 index 000000000..b91d6c5ea --- /dev/null +++ b/scanners/findings-processor/internal/runner/processor_test.go @@ -0,0 +1,20 @@ +package runner + +import ( + "strings" + "testing" + + "github.com/canada-ca/tracker/scanners/findings-processor/internal/config" +) + +func TestRun_ConnectNATSError(t *testing.T) { + t.Parallel() + + err := Run(config.Config{NATSURL: "nats://127.0.0.1:1"}) + if err == nil { + t.Fatal("Run() error = nil, want error connecting to an unreachable NATS URL") + } + if !strings.Contains(err.Error(), "failed to connect to NATS") { + t.Errorf("Run() error = %q, want it to mention the NATS connection failure", err.Error()) + } +} diff --git a/scanners/findings-processor/internal/runner/run_integration_test.go b/scanners/findings-processor/internal/runner/run_integration_test.go new file mode 100644 index 000000000..90fc408b8 --- /dev/null +++ b/scanners/findings-processor/internal/runner/run_integration_test.go @@ -0,0 +1,190 @@ +//go:build integration + +package runner_test + +import ( + "context" + "fmt" + "os" + "strings" + "syscall" + "testing" + "time" + + "github.com/arangodb/go-driver/v2/arangodb" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/config" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/model" + "github.com/canada-ca/tracker/scanners/findings-processor/internal/runner" + "github.com/nats-io/nats.go" + natscontainer "github.com/testcontainers/testcontainers-go/modules/nats" +) + +func startNATS(t *testing.T) string { + t.Helper() + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute) + defer cancel() + + container, err := natscontainer.Run(ctx, "nats:2.11.7") + if err != nil { + t.Fatalf("failed to start nats container: %v", err) + } + t.Cleanup(func() { + if err := container.Terminate(context.Background()); err != nil { + t.Logf("failed to terminate nats container: %v", err) + } + }) + + url, err := container.ConnectionString(ctx) + if err != nil { + t.Fatalf("failed to get nats connection string: %v", err) + } + return url +} + +// runUntilPersistedThenStop starts runner.Run(cfg) in the background, waits +// for findingKey to show up in Arango (proving the message was consumed and +// upserted), sends SIGINT to trigger the same shutdown path Run() uses in +// production, and returns the error Run() exited with. +func runUntilPersistedThenStop(t *testing.T, cfg config.Config, db arangodb.Database, findingKey string) error { + t.Helper() + + errCh := make(chan error, 1) + go func() { + errCh <- runner.Run(cfg) + }() + + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) + defer cancel() + + query := "FOR f IN additionalFindings FILTER f.findingKey == @key LIMIT 1 RETURN f" + for { + cursor, err := db.Query(ctx, query, &arangodb.QueryOptions{ + Count: true, + BindVars: map[string]interface{}{"key": findingKey}, + }) + if err == nil { + count := cursor.Count() + cursor.Close() + if count > 0 { + break + } + } + + select { + case <-ctx.Done(): + t.Fatalf("timed out waiting for finding %q to be persisted", findingKey) + case <-time.After(200 * time.Millisecond): + } + } + + if err := syscall.Kill(os.Getpid(), syscall.SIGINT); err != nil { + t.Fatalf("failed to send SIGINT: %v", err) + } + + select { + case err := <-errCh: + return err + case <-time.After(15 * time.Second): + t.Fatal("runner.Run() did not return after SIGINT") + return nil + } +} + +func TestIntegration_Run_ProcessesEventAndShutsDownCleanly(t *testing.T) { + db, dbCfg := startArangoDB(t) + natsURL := startNATS(t) + + nc, err := nats.Connect(natsURL) + if err != nil { + t.Fatalf("nats.Connect() error = %v", err) + } + defer nc.Close() + + js, err := nc.JetStream() + if err != nil { + t.Fatalf("JetStream() error = %v", err) + } + + t.Run("direct subscribe", func(t *testing.T) { + stream := "SCANS_direct" + subject := "scans.findings.direct" + domainKey := "domain-direct" + + if _, err := js.AddStream(&nats.StreamConfig{ + Name: stream, + Subjects: []string{subject}, + }); err != nil { + t.Fatalf("AddStream() error = %v", err) + } + + evt := model.FindingEvent{ + Source: "scanner", + FindingType: "tls-weak", + DomainKey: domainKey, + Subject: "example.com", + Confidence: "high", + ObservedAt: "2024-01-01T00:00:00Z", + } + payload := fmt.Sprintf( + `{"source":%q,"findingType":%q,"domainKey":%q,"subject":%q,"confidence":%q,"observedAt":%q}`, + evt.Source, evt.FindingType, evt.DomainKey, evt.Subject, evt.Confidence, evt.ObservedAt, + ) + if _, err := js.Publish(subject, []byte(payload)); err != nil { + t.Fatalf("Publish() error = %v", err) + } + + cfg := dbCfg + cfg.NATSURL = natsURL + cfg.NATSStream = stream + cfg.NATSSubject = subject + cfg.NATSDurable = "findings-processor" + cfg.NATSAckWait = 30 * time.Second + cfg.NATSMaxDeliver = 5 + cfg.NATSMaxPending = 64 + + if err := runUntilPersistedThenStop(t, cfg, db, evt.DeriveFindingKey()); err != nil { + t.Errorf("runner.Run() returned error = %v, want nil (clean shutdown)", err) + } + }) + + t.Run("get database error", func(t *testing.T) { + cfg := dbCfg + cfg.NATSURL = natsURL + cfg.NATSStream = "SCANS_BAD_DB" + cfg.NATSSubject = "scans.findings.bad-db" + cfg.NATSDurable = "findings-processor" + cfg.DBName = "database-that-does-not-exist" + + if _, err := js.AddStream(&nats.StreamConfig{ + Name: cfg.NATSStream, + Subjects: []string{cfg.NATSSubject}, + }); err != nil { + t.Fatalf("AddStream() error = %v", err) + } + + err := runner.Run(cfg) + if err == nil { + t.Fatal("runner.Run() error = nil, want error for a non-existent database") + } + if !strings.Contains(err.Error(), "get database failed") { + t.Errorf("runner.Run() error = %q, want it to mention the database lookup failure", err.Error()) + } + }) + + t.Run("consumer error on stream mismatch", func(t *testing.T) { + cfg := dbCfg + cfg.NATSURL = natsURL + cfg.NATSStream = "SCANS_DOES_NOT_EXIST" + cfg.NATSSubject = "scans.findings.no-such-stream" + cfg.NATSDurable = "findings-processor" + + err := runner.Run(cfg) + if err == nil { + t.Fatal("runner.Run() error = nil, want error for a subscribe against a non-existent stream") + } + if !strings.Contains(err.Error(), "create/update consumer failed") { + t.Errorf("runner.Run() error = %q, want it to mention the subscribe failure", err.Error()) + } + }) +}