diff --git a/.ai/test-plans/2026-07-17-16-28-integration-tests.md b/.ai/test-plans/2026-07-17-16-28-integration-tests.md index 2565ee4..ea6bdc4 100644 --- a/.ai/test-plans/2026-07-17-16-28-integration-tests.md +++ b/.ai/test-plans/2026-07-17-16-28-integration-tests.md @@ -1134,6 +1134,37 @@ Unless overridden, tests use: --- +### IT-MSG-071: Nested CE payload — resourceId extracted correctly + +- **Validates AC:** AC-MSG-030 +- **Test Infrastructure:** NATS JetStream; nested JSON payload with `spec` object +- **Given** a CloudEvent with nested payload `{"resourceId":"res-nested","spec":{"replicas":3}}` +- **When** the message is consumed on the main topic +- **Then** `resourceId` MUST be extracted correctly (struct ignores nested fields) +- **And** the response CE MUST include the extracted `resourceId` +- **And** a nested cancel payload MUST correctly populate the deny list + +### IT-MSG-072: Delete request produces deletion-acknowledged response + +- **Validates AC:** AC-MSG-060 +- **Test Infrastructure:** NATS JetStream; responses subscriber +- **Given** a `dcm.request.delete` CloudEvent is published to the main topic +- **When** the handler returns nil (success) +- **Then** the response CE type MUST be `dcm.agent.deletion-acknowledged` +- **And** the response CE data MUST include `status: "DELETING"` + +### IT-MSG-073: publishResponseCE failure causes nak and redelivery + +- **Validates AC:** AC-MSG-055 +- **Test Infrastructure:** NATS JetStream; client with empty AgentName +- **Given** a client created with empty `AgentName` (FormatSource will fail) +- **When** a valid message is consumed and handler returns nil +- **Then** `publishResponseCE` MUST fail (FormatSource error) +- **And** the message MUST be nak'd +- **And** JetStream MUST redeliver the message (delivery count >= 2) + +--- + ## Topic 8: Resource Operation Routing ### IT-RTE-010: Route creation to Ready embedded SP @@ -1647,12 +1678,12 @@ Unless overridden, tests use: | AC-MSG-018 | IT-MSG-040 | | AC-MSG-020 | IT-MSG-050 | | AC-MSG-025 | IT-MSG-060 | -| AC-MSG-030 | IT-MSG-070 | +| AC-MSG-030 | IT-MSG-070, IT-MSG-071 | | AC-MSG-035 | IT-MSG-080 | | AC-MSG-040 | IT-MSG-090 | | AC-MSG-050 | IT-MSG-100 | -| AC-MSG-055 | IT-MSG-110 | -| AC-MSG-060 | IT-MSG-120 | +| AC-MSG-055 | IT-MSG-110, IT-MSG-073 | +| AC-MSG-060 | IT-MSG-120, IT-MSG-072 | | AC-RTE-010 | IT-RTE-010, IT-RTE-015 | | AC-RTE-020 | IT-RTE-020 | | AC-RTE-030 | IT-RTE-030 | diff --git a/.ai/test-plans/2026-07-17-16-28-unit-tests.md b/.ai/test-plans/2026-07-17-16-28-unit-tests.md index bf6f95b..ed653b6 100644 --- a/.ai/test-plans/2026-07-17-16-28-unit-tests.md +++ b/.ai/test-plans/2026-07-17-16-28-unit-tests.md @@ -499,6 +499,11 @@ This plan identifies pure logic functions suitable for isolated unit testing in - UT-MSG-034: `"agent-prod.1"` (dot separator) — valid (dots allowed in NATS subjects) - UT-MSG-035: `"agent-prod-1"` (alphanumeric + hyphens) — valid - UT-MSG-036: Empty string — invalid +- UT-MSG-037: Contains `"!"` — invalid (whitelist regex rejects) +- UT-MSG-038: Contains `"@"` — invalid (whitelist regex rejects) +- UT-MSG-039: Contains `"/"` — invalid (whitelist regex rejects) +- UT-MSG-040: Contains underscore `"test_topic"` — valid (underscores allowed) +- UT-MSG-041: Exactly 255 characters — valid (boundary acceptance) --- @@ -644,7 +649,7 @@ This plan identifies pure logic functions suitable for isolated unit testing in | AC-HMN-005 | UT-HMN-070 | | AC-DCM-050 | UT-DCM-010, UT-DCM-011–014, UT-DCM-020–023 | | AC-DCM-061 | UT-DCM-030, UT-DCM-031–037 | -| AC-MSG-010 | UT-MSG-010, UT-MSG-020, UT-MSG-030–036 | +| AC-MSG-010 | UT-MSG-010, UT-MSG-020, UT-MSG-030–041 | | AC-RTE-055 | UT-RTE-050, UT-RTE-051–053 | | AC-RTE-070 | UT-RTE-010 | | AC-RTE-075 | UT-RTE-030, UT-RTE-040 | diff --git a/Makefile b/Makefile index 32fbb59..3be7675 100644 --- a/Makefile +++ b/Makefile @@ -38,10 +38,10 @@ test: go run github.com/onsi/ginkgo/v2/ginkgo -r --randomize-all --fail-on-pending --skip-package=test/e2e test-unit: - go run github.com/onsi/ginkgo/v2/ginkgo -r --randomize-all --fail-on-pending --label-filter=unit ./internal/config ./internal/httperror ./internal/provider ./internal/health/monitor ./internal/backoff ./internal/dcm ./cmd/environment-agent + go run github.com/onsi/ginkgo/v2/ginkgo -r --randomize-all --fail-on-pending --label-filter=unit ./internal/config ./internal/httperror ./internal/provider ./internal/health/monitor ./internal/backoff ./internal/dcm ./internal/messaging ./internal/cloudevent ./cmd/environment-agent test-integration: - go run github.com/onsi/ginkgo/v2/ginkgo -r --randomize-all --fail-on-pending --label-filter=integration ./internal/apiserver ./internal/health ./internal/health/monitor ./internal/provider ./internal/dcm + go run github.com/onsi/ginkgo/v2/ginkgo -r --randomize-all --fail-on-pending --label-filter=integration ./internal/apiserver ./internal/health ./internal/health/monitor ./internal/provider ./internal/dcm ./internal/messaging test-race: go run github.com/onsi/ginkgo/v2/ginkgo -r --race --randomize-all --fail-on-pending --skip-package=test/e2e diff --git a/cmd/environment-agent/main.go b/cmd/environment-agent/main.go index 1031e25..cf4315d 100644 --- a/cmd/environment-agent/main.go +++ b/cmd/environment-agent/main.go @@ -3,6 +3,7 @@ package main import ( "context" "errors" + "fmt" "log/slog" "net" "net/http" @@ -19,21 +20,12 @@ import ( "github.com/dcm-project/environment-agent/internal/health" "github.com/dcm-project/environment-agent/internal/health/monitor" "github.com/dcm-project/environment-agent/internal/httperror" + "github.com/dcm-project/environment-agent/internal/messaging" "github.com/dcm-project/environment-agent/internal/provider" "github.com/dcm-project/environment-agent/internal/provider/service" "github.com/dcm-project/environment-agent/internal/provider/store" ) -// TODO: replace with real MessagingStatus from the NATS/messaging subsystem. -type messagingStatus struct{} - -func (messagingStatus) IsConnected() bool { return true } - -// stubConsumerLagProvider returns 0 until Topic 7 provides real NATS consumer lag. -type stubConsumerLagProvider struct{} - -func (stubConsumerLagProvider) ConsumerLag() int64 { return 0 } - // serviceTypeLister adapts ProviderService to dcm.ServiceTypeLister. type serviceTypeLister struct { providerSvc *service.ProviderService @@ -107,19 +99,27 @@ func run(ctx context.Context) int { } providerSvc.RegisterEmbedded(cfg.Provider.EmbeddedSPs) + // Messaging client — must start before registrar (provides ConsumerLagProvider) + msgClient, topicMain, err := setupMessaging(cfg, logger) + if err != nil { + logger.Error("invalid topic name", "error", err) + return 1 + } + if err := msgClient.Start(ctx); err != nil { + logger.Error("failed to start messaging client", "error", err) + return 1 + } + defer msgClient.Stop() + // DCM Registrar — created before monitor starts so callbacks can be wired // before any health transitions fire. Deferred after monitor so LIFO shuts // registrar down first. - topicName := cfg.Messaging.TopicName - if topicName == "" { - topicName = cfg.Agent.Name - } registrar, err := dcm.NewRegistrar( dcm.RegistrarConfig{ AgentName: cfg.Agent.Name, Environment: cfg.Agent.Environment, Cost: cfg.Agent.Cost, - TopicName: topicName, + TopicName: topicMain, RegistrationURL: cfg.DCM.RegistrationURL, InitialBackoff: cfg.DCM.InitialBackoff, MaxBackoff: cfg.DCM.MaxBackoff, @@ -127,7 +127,7 @@ func run(ctx context.Context) int { PrerequisiteRetryInterval: cfg.DCM.PrerequisiteRetryInterval, }, &serviceTypeLister{providerSvc: providerSvc, logger: logger}, - stubConsumerLagProvider{}, + msgClient, nil, logger, ) @@ -159,7 +159,7 @@ func run(ctx context.Context) int { <-registrar.Done() }() - healthSvc := health.NewService(messagingStatus{}) + healthSvc := health.NewService(msgClient) strictHandler := handler.New(healthSvc, providerSvc) h := oapigen.NewStrictHandlerWithOptions(strictHandler, nil, oapigen.StrictHTTPServerOptions{ ResponseErrorHandlerFunc: func(w http.ResponseWriter, r *http.Request, err error) { @@ -177,3 +177,16 @@ func run(ctx context.Context) int { logger.Info("Environment Agent stopped") return 0 } + +func setupMessaging(cfg *config.Config, logger *slog.Logger) (*messaging.Client, string, error) { + topics := messaging.DeriveTopicNames(cfg.Agent.Name, cfg.Messaging.TopicName) + if err := messaging.ValidateTopicName(topics.Main); err != nil { + return nil, "", fmt.Errorf("invalid topic name: %w", err) + } + client := messaging.NewClient(messaging.ClientConfig{ + URL: cfg.Messaging.URL, + TopicName: topics.Main, + AgentName: cfg.Agent.Name, + }, logger) + return client, topics.Main, nil +} diff --git a/cmd/environment-agent/main_test.go b/cmd/environment-agent/main_test.go index 8fdb374..f6c8851 100644 --- a/cmd/environment-agent/main_test.go +++ b/cmd/environment-agent/main_test.go @@ -24,6 +24,7 @@ var _ = Describe("run", Label("unit"), func() { GinkgoT().Setenv("AGENT_ENVIRONMENT", "test") GinkgoT().Setenv("AGENT_COST", "medium") GinkgoT().Setenv("DCM_REGISTRATION_URL", "http://localhost:8080") + GinkgoT().Setenv("AGENT_MESSAGING_URL", "nats://localhost:4222") ctx, cancel := context.WithCancel(context.Background()) cancel() diff --git a/go.mod b/go.mod index c44da24..32b9cc7 100644 --- a/go.mod +++ b/go.mod @@ -4,9 +4,12 @@ go 1.25.5 require ( github.com/caarlos0/env/v11 v11.4.1 + github.com/cloudevents/sdk-go/v2 v2.16.2 github.com/getkin/kin-openapi v0.144.0 github.com/go-chi/chi/v5 v5.3.1 github.com/google/uuid v1.6.0 + github.com/nats-io/nats-server/v2 v2.14.4 + github.com/nats-io/nats.go v1.52.0 github.com/oapi-codegen/nethttp-middleware v1.1.2 github.com/oapi-codegen/oapi-codegen/v2 v2.8.0 github.com/oapi-codegen/runtime v1.6.0 @@ -16,6 +19,7 @@ require ( require ( github.com/Masterminds/semver/v3 v3.4.0 // indirect + github.com/antithesishq/antithesis-sdk-go v0.7.2-default-no-op // indirect github.com/apapsch/go-jsonmerge/v2 v2.0.0 // indirect github.com/dprotaso/go-yit v0.0.0-20220510233725-9ba8df137936 // indirect github.com/go-logr/logr v1.4.3 // indirect @@ -23,20 +27,33 @@ require ( github.com/go-openapi/swag/jsonname v0.26.0 // indirect github.com/go-task/slim-sprig/v3 v3.0.0 // indirect github.com/google/go-cmp v0.7.0 // indirect + github.com/google/go-tpm v0.9.8 // indirect github.com/google/pprof v0.0.0-20260402051712-545e8a4df936 // indirect github.com/gorilla/mux v1.8.1 // indirect + github.com/json-iterator/go v1.1.12 // indirect + github.com/klauspost/compress v1.19.0 // indirect + github.com/minio/highwayhash v1.0.4 // indirect + github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect + github.com/modern-go/reflect2 v1.0.2 // indirect + github.com/nats-io/jwt/v2 v2.8.2 // indirect + github.com/nats-io/nkeys v0.4.16 // indirect + github.com/nats-io/nuid v1.0.1 // indirect github.com/oasdiff/yaml v0.1.1 // indirect github.com/oasdiff/yaml3 v0.0.14 // indirect github.com/santhosh-tekuri/jsonschema/v6 v6.0.2 // indirect github.com/speakeasy-api/jsonpath v0.6.3 // indirect github.com/speakeasy-api/openapi v1.24.0 // indirect github.com/vmware-labs/yaml-jsonpath v0.3.2 // indirect + go.uber.org/multierr v1.11.0 // indirect + go.uber.org/zap v1.27.0 // indirect go.yaml.in/yaml/v3 v3.0.4 // indirect + golang.org/x/crypto v0.54.0 // indirect golang.org/x/mod v0.38.0 // indirect golang.org/x/net v0.57.0 // indirect golang.org/x/sync v0.22.0 // indirect golang.org/x/sys v0.47.0 // indirect golang.org/x/text v0.40.0 // indirect + golang.org/x/time v0.15.0 // indirect golang.org/x/tools v0.48.0 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/go.sum b/go.sum index 4ee5a38..e3a6c65 100644 --- a/go.sum +++ b/go.sum @@ -1,6 +1,8 @@ github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1Xbatp0= github.com/Masterminds/semver/v3 v3.4.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM= github.com/RaveNoX/go-jsoncommentstrip v1.0.0/go.mod h1:78ihd09MekBnJnxpICcwzCMzGrKSKYe4AqU6PDYYpjk= +github.com/antithesishq/antithesis-sdk-go v0.7.2-default-no-op h1:p2zFsAzvhIpFya8AIOHIbWf7NGvO34QpLGclyf7nXj8= +github.com/antithesishq/antithesis-sdk-go v0.7.2-default-no-op/go.mod h1:FQyySiasQQM8735Ddel3MRojmy4dA1IqCeyJ5jmPMbI= github.com/apapsch/go-jsonmerge/v2 v2.0.0 h1:axGnT1gRIfimI7gJifB699GoE/oq+F2MU7Dml6nw9rQ= github.com/apapsch/go-jsonmerge/v2 v2.0.0/go.mod h1:lvDnEdqiQrp0O42VQGgmlKpxL1AP2+08jFMw88y4klk= github.com/bmatcuk/doublestar v1.1.1/go.mod h1:UD6OnuiIn0yFxxA2le/rnRU1G4RaI4UvFv1sNto9p6w= @@ -9,6 +11,8 @@ github.com/caarlos0/env/v11 v11.4.1/go.mod h1:qupehSf/Y0TUTsxKywqRt/vJjN5nz6vaui github.com/chzyer/logex v1.1.10/go.mod h1:+Ywpsq7O8HXn0nuIou7OrIPyXbp3wmkHB+jjWRnGsAI= github.com/chzyer/readline v0.0.0-20180603132655-2972be24d48e/go.mod h1:nSuG5e5PlCu98SY8svDHJxuZscDgtXS6KTTbou5AhLI= github.com/chzyer/test v0.0.0-20180213035817-a1ea475d72b1/go.mod h1:Q3SI9o4m/ZMnBNeIyt5eFwwo7qiLfzFZmjNmxjkiQlU= +github.com/cloudevents/sdk-go/v2 v2.16.2 h1:ZYDFrYke4FD+jM8TZTJJO6JhKHzOQl2oqpFK1D+NnQM= +github.com/cloudevents/sdk-go/v2 v2.16.2/go.mod h1:laOcGImm4nVJEU+PHnUrKL56CKmRL65RlQF0kRmW/kg= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM= @@ -59,6 +63,9 @@ github.com/google/go-cmp v0.4.0/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/ github.com/google/go-cmp v0.5.5/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= 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/go-tpm v0.9.8 h1:slArAR9Ft+1ybZu0lBwpSmpwhRXaa85hWtMinMyRAWo= +github.com/google/go-tpm v0.9.8/go.mod h1:h9jEsEECg7gtLis0upRBQU+GhYVH6jMjrFxI8u6bVUY= +github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= github.com/google/pprof v0.0.0-20210407192527-94a9f03dee38/go.mod h1:kpwsk12EmLew5upagYY7GY0pfYCcupk39gWOCRROcvE= github.com/google/pprof v0.0.0-20260402051712-545e8a4df936 h1:EwtI+Al+DeppwYX2oXJCETMO23COyaKGP6fHVpkpWpg= github.com/google/pprof v0.0.0-20260402051712-545e8a4df936/go.mod h1:MxpfABSjhmINe3F1It9d+8exIHFvUqtLIRCdOGNXqiI= @@ -70,7 +77,11 @@ github.com/hpcloud/tail v1.0.0/go.mod h1:ab1qPbhIpdTxEkNHXyeSf5vhxWSCs/tWer42PpO github.com/ianlancetaylor/demangle v0.0.0-20200824232613-28f6c0f3b639/go.mod h1:aSSvb/t6k1mPoxDqO4vJh6VOCGPwU4O0C2/Eqndh1Sc= github.com/joshdk/go-junit v1.0.0 h1:S86cUKIdwBHWwA6xCmFlf3RTLfVXYQfvanM5Uh+K6GE= github.com/joshdk/go-junit v1.0.0/go.mod h1:TiiV0PqkaNfFXjEiyjWM3XXrhVyCa1K4Zfga6W52ung= +github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= +github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= github.com/juju/gnuflag v0.0.0-20171113085948-2ce1bb71843d/go.mod h1:2PavIy+JPciBPrBUjwbNvtwB6RQlve+hkpll6QSNmOE= +github.com/klauspost/compress v1.19.0 h1:sXLILfc9jV2QYWkzFOPWStmcUVH2RHEB1JCdY2oVvCQ= +github.com/klauspost/compress v1.19.0/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= @@ -82,6 +93,23 @@ github.com/maruel/natural v1.1.1 h1:Hja7XhhmvEFhcByqDoHz9QZbkWey+COd9xWfCfn1ioo= github.com/maruel/natural v1.1.1/go.mod h1:v+Rfd79xlw1AgVBjbO0BEQmptqb5HvL/k9GRHB7ZKEg= github.com/mfridman/tparse v0.18.0 h1:wh6dzOKaIwkUGyKgOntDW4liXSo37qg5AXbIhkMV3vE= github.com/mfridman/tparse v0.18.0/go.mod h1:gEvqZTuCgEhPbYk/2lS3Kcxg1GmTxxU7kTC8DvP0i/A= +github.com/minio/highwayhash v1.0.4 h1:asJizugGgchQod2ja9NJlGOWq4s7KsAWr5XUc9Clgl4= +github.com/minio/highwayhash v1.0.4/go.mod h1:GGYsuwP/fPD6Y9hMiXuapVvlIUEhFhMTh0rxU3ik1LQ= +github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M= +github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk= +github.com/nats-io/jwt/v2 v2.8.2 h1:XXRgB60MSTnqsRwejQurVDs/hcv2dkt+86GjI+I/bMc= +github.com/nats-io/jwt/v2 v2.8.2/go.mod h1:Ag/56sq9OblL4JgdYufDd16Egb17Kr/8WwwuO/forVc= +github.com/nats-io/nats-server/v2 v2.14.4 h1:efgjZ8cdExAKRuqSg8UPJFprb+l7NlBtSDPhDlw3rO4= +github.com/nats-io/nats-server/v2 v2.14.4/go.mod h1:BltdpOYestjbtQSnVO2zGHdg5SGBZjt+GYTgB9LZq/I= +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.16 h1:rd5oAuLOb8mnAycB0xleuEBNS1pVVnN0fv/FF34Eypg= +github.com/nats-io/nkeys v0.4.16/go.mod h1:llLgWoI0o4z/Q57q2R1kHfmocyhGV6VG/U18Glg1Afs= +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/nxadm/tail v1.4.4/go.mod h1:kenIhsEOeOJmVchQTgglprH7qJGnHDVpk1VPCcaMI8A= github.com/nxadm/tail v1.4.8 h1:nPr65rt6Y5JFSKQO7qToXr7pePgD6Gwiw05lkbyAQTE= github.com/nxadm/tail v1.4.8/go.mod h1:+ncqLTQzXmGhMZNUePPaPqPvBxHAIsmXswZKocGu+AU= @@ -140,14 +168,24 @@ github.com/tidwall/pretty v1.2.1 h1:qjsOFOWWQl+N3RsoF5/ssm1pHmJJwhjlSbZ51I6wMl4= github.com/tidwall/pretty v1.2.1/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= github.com/tidwall/sjson v1.2.5 h1:kLy8mja+1c9jlljvWTlSazM7cKDRfJuR/bOJhcY5NcY= github.com/tidwall/sjson v1.2.5/go.mod h1:Fvgq9kS/6ociJEDnK0Fk1cpYF4FIW6ZF7LAe+6jwd28= +github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw= +github.com/valyala/bytebufferpool v1.0.0/go.mod h1:6bBcMArwyJ5K/AmCkWv1jt77kVWyCJ6HpOuEn7z0Csc= github.com/vmware-labs/yaml-jsonpath v0.3.2 h1:/5QKeCBGdsInyDCyVNLbXyilb61MXGi9NP674f9Hobk= github.com/vmware-labs/yaml-jsonpath v0.3.2/go.mod h1:U6whw1z03QyqgWdgXxvVnQ90zN1BWz5V+51Ewf8k+rQ= github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= +go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y= +go.uber.org/zap v1.27.0 h1:aJMhYGrd5QSmlpLMr2MftRKl7t8J8PTZPA732ud/XR8= +go.uber.org/zap v1.27.0/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E= 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.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= +golang.org/x/crypto v0.54.0 h1:YLIA59K4fiNzHzjnZt2tUJQjQtUWfWbeHBqKtk3eScw= +golang.org/x/crypto v0.54.0/go.mod h1:KWL8ny2AZdGR2cWmzeHrp2azQPGogOv+HeQaVEXC2dk= golang.org/x/mod v0.3.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk= golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40= @@ -179,6 +217,7 @@ golang.org/x/sys v0.0.0-20210112080510-489259a85091/go.mod h1:h1NjWce9XRLGQEsW7w golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20211216021012-1d35b9e2eb4e/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.21.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= 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.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= @@ -189,6 +228,8 @@ golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= 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= +golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U= +golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= golang.org/x/tools v0.0.0-20201224043029-2b0845dc783e/go.mod h1:emZCQorbCU4vsT4fOWvOPXz4eW1wZW4PmDk9uLelYpA= diff --git a/internal/cloudevent/builder.go b/internal/cloudevent/builder.go new file mode 100644 index 0000000..3117b65 --- /dev/null +++ b/internal/cloudevent/builder.go @@ -0,0 +1,33 @@ +// Package cloudevent provides CloudEvent v1.0 construction utilities. +package cloudevent + +import ( + "errors" + "fmt" + "time" + + cloudevents "github.com/cloudevents/sdk-go/v2" + "github.com/google/uuid" +) + +// NewCloudEvent creates a CloudEvents v1.0 event with standard agent attributes. +func NewCloudEvent(agentID, eventType string) (cloudevents.Event, error) { + source, err := FormatSource(agentID) + if err != nil { + return cloudevents.Event{}, fmt.Errorf("failed to format source: %w", err) + } + event := cloudevents.NewEvent() + event.SetID(uuid.New().String()) + event.SetSource(source) + event.SetType(eventType) + event.SetTime(time.Now().UTC()) + return event, nil +} + +// FormatSource formats the CloudEvent source attribute for an agent. +func FormatSource(agentID string) (string, error) { + if agentID == "" { + return "", errors.New("agentID must not be empty") + } + return fmt.Sprintf("dcm/agents/%s", agentID), nil +} diff --git a/internal/cloudevent/builder_test.go b/internal/cloudevent/builder_test.go new file mode 100644 index 0000000..528ca35 --- /dev/null +++ b/internal/cloudevent/builder_test.go @@ -0,0 +1,47 @@ +package cloudevent_test + +import ( + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + + "github.com/dcm-project/environment-agent/internal/cloudevent" +) + +var _ = Describe("NewCloudEvent", Label("unit"), func() { + It("includes all required CE fields with correct values (UT-XC-CE-010)", func() { + event, err := cloudevent.NewCloudEvent("my-agent-456", "dcm.status.create") + Expect(err).NotTo(HaveOccurred()) + Expect(event.SpecVersion()).To(Equal("1.0")) + Expect(event.ID()).NotTo(BeEmpty()) + Expect(event.Source()).To(Equal("dcm/agents/my-agent-456")) + Expect(event.Type()).To(Equal("dcm.status.create")) + Expect(event.Time().IsZero()).To(BeFalse()) + }) + + It("produces distinct IDs on successive calls (UT-XC-CE-030)", func() { + e1, err1 := cloudevent.NewCloudEvent("agent-a", "dcm.test") + Expect(err1).NotTo(HaveOccurred()) + e2, err2 := cloudevent.NewCloudEvent("agent-a", "dcm.test") + Expect(err2).NotTo(HaveOccurred()) + Expect(e1.ID()).NotTo(Equal(e2.ID())) + }) +}) + +var _ = Describe("FormatSource", Label("unit"), func() { + It("formats source as dcm/agents/{agentId} (UT-XC-CE-020)", func() { + result, err := cloudevent.FormatSource("my-agent-456") + Expect(err).NotTo(HaveOccurred()) + Expect(result).To(Equal("dcm/agents/my-agent-456")) + }) + + It("rejects empty agentId (UT-XC-CE-021)", func() { + _, err := cloudevent.FormatSource("") + Expect(err).To(MatchError(ContainSubstring("empty"))) + }) + + It("preserves special characters in agentId (UT-XC-CE-022)", func() { + result, err := cloudevent.FormatSource("special!chars") + Expect(err).NotTo(HaveOccurred()) + Expect(result).To(Equal("dcm/agents/special!chars")) + }) +}) diff --git a/internal/cloudevent/suite_test.go b/internal/cloudevent/suite_test.go new file mode 100644 index 0000000..dbe2f70 --- /dev/null +++ b/internal/cloudevent/suite_test.go @@ -0,0 +1,13 @@ +package cloudevent_test + +import ( + "testing" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +func TestCloudevent(t *testing.T) { + RegisterFailHandler(Fail) + RunSpecs(t, "CloudEvent Suite") +} diff --git a/internal/config/config.go b/internal/config/config.go index 34f6b27..9575eb6 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -123,6 +123,8 @@ func (c *Config) Validate() error { } // Topic 7: Messaging Integration — append-only below this line - // AGENT_MESSAGING_URL validated when messaging subsystem is wired. + if err := validateRequired("AGENT_MESSAGING_URL", c.Messaging.URL); err != nil { + return err + } return nil } diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 74fc0be..4655af1 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -85,18 +85,19 @@ var _ = Describe("Server Configuration", Label("unit"), func() { }) }) -// setValidTopicSixEnv sets all required Topic 6 env vars to valid defaults. -func setValidTopicSixEnv() { +// setValidEnv sets all required env vars to valid defaults. +func setValidEnv() { GinkgoT().Setenv("AGENT_NAME", "test-agent") GinkgoT().Setenv("AGENT_ENVIRONMENT", "test") GinkgoT().Setenv("AGENT_COST", "medium") GinkgoT().Setenv("DCM_REGISTRATION_URL", "http://localhost:8080") + GinkgoT().Setenv("AGENT_MESSAGING_URL", "nats://localhost:4222") } var _ = Describe("Topic 6 Config", Label("unit"), func() { Describe("Load", func() { It("parses Topic 6 config fields from env (UT-XC-CFG-040)", func() { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("DCM_REGISTRATION_INITIAL_BACKOFF", "2s") GinkgoT().Setenv("DCM_REGISTRATION_MAX_BACKOFF", "10m") GinkgoT().Setenv("AGENT_HEARTBEAT_INTERVAL", "45s") @@ -125,7 +126,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { }) It("rejects malformed duration string at parse time (UT-XC-CFG-035)", func() { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("AGENT_HEARTBEAT_INTERVAL", "abc") _, err := config.Load() Expect(err).To(HaveOccurred()) @@ -135,7 +136,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { Describe("Validate", func() { DescribeTable("rejects absent required field (UT-XC-CFG-010, UT-XC-CFG-011)", func(envVar string) { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv(envVar, "") cfg, err := config.Load() Expect(err).NotTo(HaveOccurred()) @@ -147,10 +148,11 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { Entry("AGENT_ENVIRONMENT", "AGENT_ENVIRONMENT"), Entry("AGENT_COST", "AGENT_COST"), Entry("DCM_REGISTRATION_URL", "DCM_REGISTRATION_URL"), + Entry("AGENT_MESSAGING_URL (UT-XC-CFG-011)", "AGENT_MESSAGING_URL"), ) It("accepts all required fields present (UT-XC-CFG-012)", func() { - setValidTopicSixEnv() + setValidEnv() cfg, err := config.Load() Expect(err).NotTo(HaveOccurred()) Expect(cfg.Validate()).To(Succeed()) @@ -158,7 +160,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { DescribeTable("rejects whitespace-only required field (UT-XC-CFG-013)", func(envVar string) { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv(envVar, " ") cfg, err := config.Load() Expect(err).NotTo(HaveOccurred()) @@ -171,7 +173,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { ) It("rejects invalid AGENT_COST value (UT-XC-CFG-020)", func() { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("AGENT_COST", "expensive") cfg, err := config.Load() Expect(err).NotTo(HaveOccurred()) @@ -183,7 +185,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { DescribeTable("accepts valid cost values (UT-XC-CFG-021, UT-XC-CFG-022, UT-XC-CFG-023)", func(cost string) { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("AGENT_COST", cost) cfg, err := config.Load() Expect(err).NotTo(HaveOccurred()) @@ -197,7 +199,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { ) It("rejects case-sensitive cost (UT-XC-CFG-024)", func() { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("AGENT_COST", "Medium") cfg, err := config.Load() Expect(err).NotTo(HaveOccurred()) @@ -207,7 +209,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { }) It("rejects empty cost (UT-XC-CFG-025)", func() { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("AGENT_COST", "") cfg, err := config.Load() Expect(err).NotTo(HaveOccurred()) @@ -217,7 +219,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { }) It("accepts heartbeat interval at minimum 5s (UT-XC-CFG-031)", func() { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("AGENT_HEARTBEAT_INTERVAL", "5s") cfg, err := config.Load() Expect(err).NotTo(HaveOccurred()) @@ -225,7 +227,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { }) It("accepts heartbeat interval at maximum 10m (UT-XC-CFG-032)", func() { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("AGENT_HEARTBEAT_INTERVAL", "10m") cfg, err := config.Load() Expect(err).NotTo(HaveOccurred()) @@ -233,7 +235,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { }) It("rejects heartbeat interval below minimum 5s (UT-XC-CFG-033)", func() { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("AGENT_HEARTBEAT_INTERVAL", "4s") cfg, err := config.Load() Expect(err).NotTo(HaveOccurred()) @@ -243,7 +245,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { }) It("rejects heartbeat interval above maximum 10m (UT-XC-CFG-034)", func() { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("AGENT_HEARTBEAT_INTERVAL", "11m") cfg, err := config.Load() Expect(err).NotTo(HaveOccurred()) @@ -254,7 +256,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { DescribeTable("integer config range (UT-XC-CFG-036)", func(value int, shouldPass bool) { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("AGENT_HEALTH_FAILURE_THRESHOLD", fmt.Sprintf("%d", value)) cfg, err := config.Load() Expect(err).NotTo(HaveOccurred()) @@ -273,7 +275,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { ) It("accepts timeout equal to interval (UT-XC-CFG-041)", func() { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("AGENT_HEALTH_CHECK_TIMEOUT", "10s") GinkgoT().Setenv("AGENT_HEALTH_CHECK_INTERVAL", "10s") cfg, err := config.Load() @@ -282,7 +284,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { }) It("accepts timeout below interval (UT-XC-CFG-042)", func() { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("AGENT_HEALTH_CHECK_TIMEOUT", "9s") GinkgoT().Setenv("AGENT_HEALTH_CHECK_INTERVAL", "10s") cfg, err := config.Load() @@ -291,7 +293,7 @@ var _ = Describe("Topic 6 Config", Label("unit"), func() { }) It("rejects initial backoff exceeding max backoff (UT-XC-CFG-032 cross-field)", func() { - setValidTopicSixEnv() + setValidEnv() GinkgoT().Setenv("DCM_REGISTRATION_INITIAL_BACKOFF", "10m") GinkgoT().Setenv("DCM_REGISTRATION_MAX_BACKOFF", "1m") cfg, err := config.Load() diff --git a/internal/messaging/client.go b/internal/messaging/client.go new file mode 100644 index 0000000..596fbdf --- /dev/null +++ b/internal/messaging/client.go @@ -0,0 +1,311 @@ +// Package messaging provides NATS/JetStream messaging client and topic management. +package messaging + +import ( + "context" + "errors" + "fmt" + "log/slog" + "math" + "sync" + "sync/atomic" + "time" + + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" + + "github.com/dcm-project/environment-agent/internal/dcm" + "github.com/dcm-project/environment-agent/internal/health" +) + +var ( + _ health.MessagingStatus = (*Client)(nil) + _ dcm.ConsumerLagProvider = (*Client)(nil) +) + +const ( + // SentinelConsumerLag is a clearly impossible value used by the stub. + SentinelConsumerLag = math.MinInt64 + + SubjectResponses = "dcm.agents.responses" + TypeCreationAcked = "dcm.agent.creation-acknowledged" + TypeDeletionAcked = "dcm.agent.deletion-acknowledged" + TypeAcked = "dcm.agent.acknowledged" + + drainTimeout = 5 * time.Second + drainBatchWait = 200 * time.Millisecond + drainBatchSize = 100 +) + +// MessageHandler is the callback type for processing consumed messages. +type MessageHandler func(ctx context.Context, msg []byte) error + +// ClientConfig holds the NATS client configuration. +// TopicName, if set, must be a pre-validated topic name (via ValidateTopicName). +// If empty, the agent name is used as the base for topic derivation. +type ClientConfig struct { + URL string + TopicName string + AgentName string +} + +// Client manages NATS/JetStream connectivity and topic consumption. +type Client struct { + cfg ClientConfig + logger *slog.Logger + topics TopicNames + + conn *nats.Conn + js jetstream.JetStream + mainCons jetstream.Consumer + consumers []jetstream.ConsumeContext + mu sync.Mutex + + connected atomic.Bool + stopped atomic.Bool + + // ponytail: set-before-Start contract — not safe for concurrent mutation + mainHandler MessageHandler + cancelHandler MessageHandler + + denyList sync.Map + setupOnce sync.Once + stopOnce sync.Once +} + +// NewClient creates a new messaging client. Does NOT connect to NATS. +func NewClient(cfg ClientConfig, logger *slog.Logger) *Client { + return &Client{ + cfg: cfg, + logger: logger, + topics: DeriveTopicNames(cfg.AgentName, cfg.TopicName), + } +} + +// Start connects to NATS, creates streams/consumers, and begins consuming. +// Non-blocking: returns nil immediately even if NATS is unreachable. +func (c *Client) Start(_ context.Context) error { + setupCtx := context.Background() + + conn, err := nats.Connect(c.cfg.URL, + nats.RetryOnFailedConnect(true), + nats.MaxReconnects(-1), + nats.ReconnectWait(2*time.Second), + nats.ConnectHandler(func(nc *nats.Conn) { + c.connected.Store(true) + c.logger.Info("NATS connected", "url", nc.ConnectedUrl()) + go c.doSetup(setupCtx, nc) + }), + nats.DisconnectErrHandler(func(_ *nats.Conn, err error) { + c.connected.Store(false) + c.logger.Warn("NATS disconnected", "error", err) + }), + nats.ReconnectHandler(func(nc *nats.Conn) { + c.connected.Store(true) + c.logger.Info("NATS reconnected", "url", nc.ConnectedUrl()) + }), + ) + if err != nil { + return fmt.Errorf("failed to connect to NATS: %w", err) + } + + c.mu.Lock() + c.conn = conn + c.mu.Unlock() + + if conn.IsConnected() { + c.connected.Store(true) + c.doSetup(setupCtx, conn) + } + + return nil +} + +func (c *Client) doSetup(ctx context.Context, conn *nats.Conn) { + if c.stopped.Load() { + return + } + c.setupOnce.Do(func() { + c.setupStreamsAndConsume(ctx, conn) + }) +} + +func (c *Client) setupStreamsAndConsume(ctx context.Context, conn *nats.Conn) { + js, err := jetstream.New(conn) + if err != nil { + c.logger.Error("failed to create JetStream context", "error", err) + return + } + + c.mu.Lock() + c.conn = conn + c.js = js + c.mu.Unlock() + + mainS, retryS, cancelS, err := c.initStreams(ctx, js) + if err != nil { + c.logger.Error("failed to initialize streams", "error", err) + return + } + + mainCons, cancelCons, err := c.initConsumers(ctx, mainS, retryS, cancelS) + if err != nil { + c.logger.Error("failed to initialize consumers", "error", err) + return + } + + c.drainCancelTopic(ctx, cancelCons) + + cc, err := cancelCons.Consume(func(msg jetstream.Msg) { + c.handleCancelMessage(msg) + _ = msg.Ack() + }) + if err != nil { + c.logger.Error("failed to start cancel consumer", "error", err) + return + } + c.mu.Lock() + c.consumers = append(c.consumers, cc) + c.mu.Unlock() + + cc, err = mainCons.Consume(func(msg jetstream.Msg) { + c.handleMainMessage(msg) + }) + if err != nil { + c.logger.Error("failed to start main consumer", "error", err) + return + } + c.mu.Lock() + c.consumers = append(c.consumers, cc) + c.mu.Unlock() +} + +func (c *Client) initStreams(ctx context.Context, js jetstream.JetStream) (jetstream.Stream, jetstream.Stream, jetstream.Stream, error) { + mainS, err := js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{ + Name: c.topics.Main, Subjects: []string{c.topics.Main}, + }) + if err != nil { + return nil, nil, nil, fmt.Errorf("failed to create main stream: %w", err) + } + retryS, err := js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{ + Name: c.topics.Main + "-retry", Subjects: []string{c.topics.Retry}, + }) + if err != nil { + return nil, nil, nil, fmt.Errorf("failed to create retry stream: %w", err) + } + cancelS, err := js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{ + Name: c.topics.Main + "-cancel", Subjects: []string{c.topics.Cancel}, + }) + if err != nil { + return nil, nil, nil, fmt.Errorf("failed to create cancel stream: %w", err) + } + return mainS, retryS, cancelS, nil +} + +func (c *Client) initConsumers(ctx context.Context, mainS, retryS, cancelS jetstream.Stream) (jetstream.Consumer, jetstream.Consumer, error) { + mainCons, err := mainS.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{ + Durable: c.topics.Main + "-consumer", AckPolicy: jetstream.AckExplicitPolicy, + }) + if err != nil { + return nil, nil, fmt.Errorf("failed to create main consumer: %w", err) + } + // ponytail: retry consumer created for Topic 9; not consumed yet + if _, err := retryS.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{ + Durable: c.topics.Main + "-retry-consumer", AckPolicy: jetstream.AckExplicitPolicy, + }); err != nil { + return nil, nil, fmt.Errorf("failed to create retry consumer: %w", err) + } + cancelCons, err := cancelS.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{ + Durable: c.topics.Main + "-cancel-consumer", AckPolicy: jetstream.AckExplicitPolicy, + }) + if err != nil { + return nil, nil, fmt.Errorf("failed to create cancel consumer: %w", err) + } + + c.mu.Lock() + c.mainCons = mainCons + c.mu.Unlock() + + return mainCons, cancelCons, nil +} + +func (c *Client) drainCancelTopic(_ context.Context, cancelCons jetstream.Consumer) { + drainCtx, drainCancel := context.WithTimeout(context.Background(), drainTimeout) + defer drainCancel() + + for drainCtx.Err() == nil { + batch, err := cancelCons.Fetch(drainBatchSize, jetstream.FetchMaxWait(drainBatchWait)) + if err != nil { + return + } + count := 0 + for msg := range batch.Messages() { + c.handleCancelMessage(msg) + _ = msg.Ack() + count++ + } + if count == 0 { + return + } + } +} + +// Stop gracefully shuts down the client. +func (c *Client) Stop() { + c.stopOnce.Do(func() { + c.stopped.Store(true) + + c.mu.Lock() + consumers := c.consumers + c.consumers = nil + conn := c.conn + c.mu.Unlock() + + for _, cc := range consumers { + cc.Stop() + } + if conn != nil { + conn.Close() + } + c.connected.Store(false) + }) +} + +// IsConnected returns the cached connectivity state (no I/O). +func (c *Client) IsConnected() bool { return c.connected.Load() } + +// ConsumerLag returns the current consumer lag from JetStream. +func (c *Client) ConsumerLag() int64 { + c.mu.Lock() + cons := c.mainCons + c.mu.Unlock() + + if cons == nil { + return 0 + } + return int64(cons.CachedInfo().NumPending) +} + +// TopicName returns the main topic name advertised to DCM. +func (c *Client) TopicName() string { return c.topics.Main } + +// Publish publishes raw bytes to the given NATS subject via JetStream. +func (c *Client) Publish(ctx context.Context, subject string, data []byte) error { + c.mu.Lock() + js := c.js + c.mu.Unlock() + + if js == nil { + return errors.New("jetstream not initialized") + } + _, err := js.Publish(ctx, subject, data) + return err +} + +// SetMainHandler sets the handler for messages on the main topic. +// Must be called before Start. +func (c *Client) SetMainHandler(h MessageHandler) { c.mainHandler = h } + +// SetCancelHandler sets the handler for messages on the cancel topic. +// Must be called before Start. +func (c *Client) SetCancelHandler(h MessageHandler) { c.cancelHandler = h } diff --git a/internal/messaging/client_integration_test.go b/internal/messaging/client_integration_test.go new file mode 100644 index 0000000..02414e9 --- /dev/null +++ b/internal/messaging/client_integration_test.go @@ -0,0 +1,780 @@ +package messaging_test + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + "os" + "sync/atomic" + "time" + + cloudevents "github.com/cloudevents/sdk-go/v2" + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + + "github.com/google/uuid" + natstest "github.com/nats-io/nats-server/v2/test" + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" + + "github.com/dcm-project/environment-agent/internal/messaging" +) + +func deleteStreams(js jetstream.JetStream, topicName string) { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + for _, suffix := range []string{"", "-retry", "-cancel"} { + _ = js.DeleteStream(ctx, topicName+suffix) + } +} + +func publishCE(ctx context.Context, js jetstream.JetStream, subject, ceType, source string, payload any) { + event := cloudevents.NewEvent() + event.SetID(uuid.New().String()) + event.SetSource(source) + event.SetType(ceType) + event.SetTime(time.Now()) + _ = event.SetData(cloudevents.ApplicationJSON, payload) + data, err := json.Marshal(event) + Expect(err).NotTo(HaveOccurred()) + _, err = js.Publish(ctx, subject, data) + Expect(err).NotTo(HaveOccurred()) +} + +var _ = Describe("Topic Management", Label("integration"), func() { + var ( + ctx context.Context + cancel context.CancelFunc + testConn *nats.Conn + testJS jetstream.JetStream + topicName string + logger *slog.Logger + ) + + BeforeEach(func() { + var err error + topicName = fmt.Sprintf("test-%s", uuid.New().String()[:8]) + logger = slog.Default() + ctx, cancel = context.WithTimeout(context.Background(), 30*time.Second) //nolint:fatcontext // Ginkgo BeforeEach pattern + + testConn, err = nats.Connect(testNATSServer.ClientURL()) + Expect(err).NotTo(HaveOccurred()) + testJS, err = jetstream.New(testConn) + Expect(err).NotTo(HaveOccurred()) + }) + + AfterEach(func() { + cancel() + deleteStreams(testJS, topicName) + testConn.Close() + }) + + It("creates three JetStream subjects at startup (IT-MSG-010)", func() { + client := messaging.NewClient(messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + }, logger) + Expect(client.Start(ctx)).To(Succeed()) + defer client.Stop() + + mainStream, err := testJS.Stream(ctx, topicName) + Expect(err).NotTo(HaveOccurred()) + Expect(mainStream.CachedInfo().Config.Subjects).To(ContainElement(topicName)) + + retryStream, err := testJS.Stream(ctx, topicName+"-retry") + Expect(err).NotTo(HaveOccurred()) + Expect(retryStream.CachedInfo().Config.Subjects).To(ContainElement(topicName + ".retry")) + + cancelStream, err := testJS.Stream(ctx, topicName+"-cancel") + Expect(err).NotTo(HaveOccurred()) + Expect(cancelStream.CachedInfo().Config.Subjects).To(ContainElement(topicName + ".cancel")) + }) + + It("creates deterministic durable consumer names derived from topic (IT-MSG-020)", func() { + client := messaging.NewClient(messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + }, logger) + Expect(client.Start(ctx)).To(Succeed()) + defer client.Stop() + + stream, err := testJS.Stream(ctx, topicName) + Expect(err).NotTo(HaveOccurred()) + + // Verify consumer names are deterministic (derived from topic name) + consNames := []string{} + lister := stream.ListConsumers(ctx) + for info := range lister.Info() { + consNames = append(consNames, info.Name) + } + Expect(consNames).NotTo(BeEmpty()) + // Consumer names must contain the topic name for determinism + for _, name := range consNames { + Expect(name).To(ContainSubstring(topicName)) + } + }) + + It("reuses existing topics on restart without error (IT-MSG-050)", func() { + cfg := messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + } + + client1 := messaging.NewClient(cfg, logger) + Expect(client1.Start(ctx)).To(Succeed()) + client1.Stop() + + client2 := messaging.NewClient(cfg, logger) + Expect(client2.Start(ctx)).To(Succeed()) + client2.Stop() + }) + + It("reuses existing consumers on restart — no duplicates (IT-MSG-030)", func() { + cfg := messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + } + + client1 := messaging.NewClient(cfg, logger) + Expect(client1.Start(ctx)).To(Succeed()) + + stream, err := testJS.Stream(ctx, topicName) + Expect(err).NotTo(HaveOccurred()) + initialConsumers := stream.CachedInfo().State.Consumers + client1.Stop() + + client2 := messaging.NewClient(cfg, logger) + Expect(client2.Start(ctx)).To(Succeed()) + defer client2.Stop() + + stream, err = testJS.Stream(ctx, topicName) + Expect(err).NotTo(HaveOccurred()) + Expect(stream.CachedInfo().State.Consumers).To(Equal(initialConsumers)) + }) +}) + +var _ = Describe("Message Durability", Label("integration"), func() { + var ( + ctx context.Context + cancel context.CancelFunc + testConn *nats.Conn + testJS jetstream.JetStream + topicName string + logger *slog.Logger + ) + + BeforeEach(func() { + var err error + topicName = fmt.Sprintf("test-%s", uuid.New().String()[:8]) + logger = slog.Default() + ctx, cancel = context.WithTimeout(context.Background(), 30*time.Second) //nolint:fatcontext // Ginkgo BeforeEach pattern + + testConn, err = nats.Connect(testNATSServer.ClientURL()) + Expect(err).NotTo(HaveOccurred()) + testJS, err = jetstream.New(testConn) + Expect(err).NotTo(HaveOccurred()) + }) + + AfterEach(func() { + cancel() + deleteStreams(testJS, topicName) + testConn.Close() + }) + + It("redelivers unacknowledged message after restart (IT-MSG-040)", func() { + cfg := messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + } + + received := make(chan []byte, 1) + client1 := messaging.NewClient(cfg, logger) + client1.SetMainHandler(func(_ context.Context, msg []byte) error { + received <- msg + return fmt.Errorf("simulated failure") + }) + Expect(client1.Start(ctx)).To(Succeed()) + + publishCE(ctx, testJS, topicName, "dcm.command.create", "dcm/test", map[string]string{"key": "value"}) + + Eventually(received, 5*time.Second).Should(Receive()) + client1.Stop() + + redelivered := make(chan []byte, 1) + client2 := messaging.NewClient(cfg, logger) + client2.SetMainHandler(func(_ context.Context, msg []byte) error { + redelivered <- msg + return nil + }) + Expect(client2.Start(ctx)).To(Succeed()) + defer client2.Stop() + + Eventually(redelivered, 5*time.Second).Should(Receive()) + }) +}) + +var _ = Describe("Topic Advertising", Label("integration"), func() { + var ( + ctx context.Context + cancel context.CancelFunc + testConn *nats.Conn + testJS jetstream.JetStream + topicName string + logger *slog.Logger + ) + + BeforeEach(func() { + var err error + topicName = fmt.Sprintf("test-%s", uuid.New().String()[:8]) + logger = slog.Default() + ctx, cancel = context.WithTimeout(context.Background(), 10*time.Second) //nolint:fatcontext // Ginkgo BeforeEach pattern + + testConn, err = nats.Connect(testNATSServer.ClientURL()) + Expect(err).NotTo(HaveOccurred()) + testJS, err = jetstream.New(testConn) + Expect(err).NotTo(HaveOccurred()) + }) + + AfterEach(func() { + cancel() + deleteStreams(testJS, topicName) + testConn.Close() + }) + + It("returns only main topic name — not retry or cancel (IT-MSG-060)", func() { + client := messaging.NewClient(messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + }, logger) + Expect(client.Start(ctx)).To(Succeed()) + defer client.Stop() + + Expect(client.TopicName()).To(Equal(topicName)) + }) +}) + +var _ = Describe("Message Consumption", Label("integration"), func() { + var ( + ctx context.Context + cancel context.CancelFunc + testConn *nats.Conn + testJS jetstream.JetStream + topicName string + logger *slog.Logger + ) + + BeforeEach(func() { + var err error + topicName = fmt.Sprintf("test-%s", uuid.New().String()[:8]) + logger = slog.Default() + ctx, cancel = context.WithTimeout(context.Background(), 30*time.Second) //nolint:fatcontext // Ginkgo BeforeEach pattern + + testConn, err = nats.Connect(testNATSServer.ClientURL()) + Expect(err).NotTo(HaveOccurred()) + testJS, err = jetstream.New(testConn) + Expect(err).NotTo(HaveOccurred()) + }) + + AfterEach(func() { + cancel() + deleteStreams(testJS, topicName) + testConn.Close() + }) + + It("invokes handler on main topic message and publishes response CE (IT-MSG-070)", func() { + cfg := messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + } + + handlerCalled := make(chan []byte, 1) + client := messaging.NewClient(cfg, logger) + client.SetMainHandler(func(_ context.Context, msg []byte) error { + handlerCalled <- msg + return nil + }) + Expect(client.Start(ctx)).To(Succeed()) + defer client.Stop() + + responseSub, err := testConn.SubscribeSync("dcm.agents.responses") + Expect(err).NotTo(HaveOccurred()) + + publishCE(ctx, testJS, topicName, "dcm.command.create", "dcm/control-plane", + map[string]string{"resourceId": "res-001"}) + + Eventually(handlerCalled, 5*time.Second).Should(Receive()) + + msg, err := responseSub.NextMsg(5 * time.Second) + Expect(err).NotTo(HaveOccurred()) + + var respEvent cloudevents.Event + Expect(json.Unmarshal(msg.Data, &respEvent)).To(Succeed()) + Expect(respEvent.Type()).To(Equal("dcm.agent.creation-acknowledged")) + }) + + It("cancel message updates deny list — blocks subsequent create for same resourceId (IT-MSG-080)", func() { + cfg := messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + } + + mainReceived := make(chan string, 5) + cancelReceived := make(chan string, 5) + + client := messaging.NewClient(cfg, logger) + client.SetMainHandler(func(_ context.Context, msg []byte) error { + var event cloudevents.Event + _ = json.Unmarshal(msg, &event) + var payload map[string]string + _ = json.Unmarshal(event.Data(), &payload) + mainReceived <- payload["resourceId"] + return nil + }) + client.SetCancelHandler(func(_ context.Context, msg []byte) error { + var event cloudevents.Event + _ = json.Unmarshal(msg, &event) + var payload map[string]string + _ = json.Unmarshal(event.Data(), &payload) + cancelReceived <- payload["resourceId"] + return nil + }) + Expect(client.Start(ctx)).To(Succeed()) + defer client.Stop() + + // Cancel res-001 + publishCE(ctx, testJS, topicName+".cancel", "dcm.command.cancel", "dcm/control-plane", + map[string]string{"resourceId": "res-001"}) + + Eventually(cancelReceived, 5*time.Second).Should(Receive(Equal("res-001"))) + + // Create for same resourceId — should be filtered by deny list + publishCE(ctx, testJS, topicName, "dcm.command.create", "dcm/control-plane", + map[string]string{"resourceId": "res-001"}) + + // Create for different resourceId — should go through (positive control) + publishCE(ctx, testJS, topicName, "dcm.command.create", "dcm/control-plane", + map[string]string{"resourceId": "res-002"}) + + // res-002 should arrive but res-001 should be filtered + Eventually(mainReceived, 5*time.Second).Should(Receive(Equal("res-002"))) + Consistently(mainReceived, 2*time.Second).ShouldNot(Receive(Equal("res-001"))) + }) + + It("drains cancel topic before processing main topic (IT-MSG-090)", func() { + cfg := messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + } + + // Pre-populate streams manually + _, err := testJS.CreateOrUpdateStream(ctx, jetstream.StreamConfig{ + Name: topicName + "-cancel", + Subjects: []string{topicName + ".cancel"}, + }) + Expect(err).NotTo(HaveOccurred()) + _, err = testJS.CreateOrUpdateStream(ctx, jetstream.StreamConfig{ + Name: topicName, + Subjects: []string{topicName}, + }) + Expect(err).NotTo(HaveOccurred()) + + // Publish cancel for "res-cancel-1" and main for "res-main-1" (different IDs to avoid deny-list) + publishCE(ctx, testJS, topicName+".cancel", "dcm.command.cancel", "dcm/control-plane", + map[string]string{"resourceId": "res-cancel-1"}) + publishCE(ctx, testJS, topicName, "dcm.command.create", "dcm/control-plane", + map[string]string{"resourceId": "res-main-1"}) + + // Start client — cancel must be processed before main + order := make(chan string, 10) + client := messaging.NewClient(cfg, logger) + client.SetCancelHandler(func(_ context.Context, _ []byte) error { + order <- "cancel" + return nil + }) + client.SetMainHandler(func(_ context.Context, _ []byte) error { + order <- "main" + return nil + }) + Expect(client.Start(ctx)).To(Succeed()) + defer client.Stop() + + // Verify ordering: cancel processed before main + var first, second string + Eventually(order, 10*time.Second).Should(Receive(&first)) + Eventually(order, 10*time.Second).Should(Receive(&second)) + Expect(first).To(Equal("cancel")) + Expect(second).To(Equal("main")) + }) + + It("drain completes within 5s timeout — main processing begins after (IT-MSG-095)", func() { + cfg := messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + } + + _, err := testJS.CreateOrUpdateStream(ctx, jetstream.StreamConfig{ + Name: topicName + "-cancel", + Subjects: []string{topicName + ".cancel"}, + }) + Expect(err).NotTo(HaveOccurred()) + _, err = testJS.CreateOrUpdateStream(ctx, jetstream.StreamConfig{ + Name: topicName, + Subjects: []string{topicName}, + }) + Expect(err).NotTo(HaveOccurred()) + + // Pre-populate cancel messages + for i := 0; i < 5; i++ { + publishCE(ctx, testJS, topicName+".cancel", "dcm.command.cancel", "dcm/control-plane", + map[string]string{"resourceId": fmt.Sprintf("res-drain-%d", i)}) + } + + // Publish a main message with distinct resourceId (won't be in deny list) + publishCE(ctx, testJS, topicName, "dcm.command.create", "dcm/control-plane", + map[string]string{"resourceId": "res-main-not-cancelled"}) + + mainProcessed := make(chan struct{}, 1) + client := messaging.NewClient(cfg, logger) + client.SetCancelHandler(func(_ context.Context, _ []byte) error { + return nil + }) + client.SetMainHandler(func(_ context.Context, _ []byte) error { + mainProcessed <- struct{}{} + return nil + }) + + startTime := time.Now() + Expect(client.Start(ctx)).To(Succeed()) + defer client.Stop() + + // Main processing must begin within drain timeout (5s) + some margin + Eventually(mainProcessed, 7*time.Second).Should(Receive()) + Expect(time.Since(startTime)).To(BeNumerically("<=", 7*time.Second)) + }) + + It("extracts resourceId from nested CE payload — struct ignores extra fields (IT-MSG-071)", func() { + cfg := messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + } + + handlerCalled := make(chan []byte, 1) + cancelProcessed := make(chan struct{}, 1) + client := messaging.NewClient(cfg, logger) + client.SetMainHandler(func(_ context.Context, msg []byte) error { + handlerCalled <- msg + return nil + }) + client.SetCancelHandler(func(_ context.Context, _ []byte) error { + cancelProcessed <- struct{}{} + return nil + }) + Expect(client.Start(ctx)).To(Succeed()) + defer client.Stop() + + responseSub, err := testConn.SubscribeSync("dcm.agents.responses") + Expect(err).NotTo(HaveOccurred()) + + // Nested payload — resourceId at top level, extra nested object + nestedPayload := map[string]any{ + "resourceId": "res-nested", + "spec": map[string]any{"replicas": 3, "image": "nginx:latest"}, + } + publishCE(ctx, testJS, topicName, "dcm.command.create", "dcm/control-plane", nestedPayload) + + Eventually(handlerCalled, 5*time.Second).Should(Receive()) + + // Verify response CE contains the extracted resourceId + msg, err := responseSub.NextMsg(5 * time.Second) + Expect(err).NotTo(HaveOccurred()) + var respEvent cloudevents.Event + Expect(json.Unmarshal(msg.Data, &respEvent)).To(Succeed()) + var respPayload map[string]any + Expect(json.Unmarshal(respEvent.Data(), &respPayload)).To(Succeed()) + Expect(respPayload["resourceId"]).To(Equal("res-nested")) + + // Also verify nested cancel populates deny list + cancelPayload := map[string]any{ + "resourceId": "res-nested-cancel", + "metadata": map[string]any{"reason": "user-requested"}, + } + publishCE(ctx, testJS, topicName+".cancel", "dcm.command.cancel", "dcm/control-plane", cancelPayload) + + // Wait for cancel handler to confirm processing (no sleep) + Eventually(cancelProcessed, 5*time.Second).Should(Receive()) + + // Main message for cancelled resourceId should be filtered + publishCE(ctx, testJS, topicName, "dcm.command.create", "dcm/control-plane", + map[string]string{"resourceId": "res-nested-cancel"}) + + // Positive control — different resourceId goes through + publishCE(ctx, testJS, topicName, "dcm.command.create", "dcm/control-plane", + map[string]string{"resourceId": "res-not-cancelled"}) + + Eventually(handlerCalled, 5*time.Second).Should(Receive()) + Consistently(func() int { return len(handlerCalled) }, 1*time.Second).Should(Equal(0)) + }) +}) + +var _ = Describe("Connection Resilience", Label("integration"), func() { + var ( + ctx context.Context + cancel context.CancelFunc + topicName string + logger *slog.Logger + ) + + BeforeEach(func() { + topicName = fmt.Sprintf("test-%s", uuid.New().String()[:8]) + logger = slog.Default() + ctx, cancel = context.WithTimeout(context.Background(), 30*time.Second) //nolint:fatcontext // Ginkgo BeforeEach pattern + }) + + AfterEach(func() { + cancel() + }) + + It("HTTP health returns 200 without NATS — IsConnected false→true on reconnect (IT-MSG-100)", func() { + // Use a dedicated port for this test's NATS server lifecycle + const reconnectPort = 14222 + reconnectURL := fmt.Sprintf("nats://127.0.0.1:%d", reconnectPort) + + // Create client pointing to not-yet-started NATS server + client := messaging.NewClient(messaging.ClientConfig{ + URL: reconnectURL, + TopicName: topicName, + AgentName: "test-agent", + }, logger) + + // Start must NOT block — non-blocking per REQ-MSG-110 + startDone := make(chan error, 1) + go func() { startDone <- client.Start(ctx) }() + Eventually(startDone, 5*time.Second).Should(Receive(Not(HaveOccurred()))) + + // IsConnected must be false when NATS is unreachable + Expect(client.IsConnected()).To(BeFalse()) + + // Start a NATS server on the expected port — client should auto-reconnect + opts := natstest.DefaultTestOptions + opts.Port = reconnectPort + opts.JetStream = true + tmpDir, err := os.MkdirTemp("", "nats-reconnect-*") + Expect(err).NotTo(HaveOccurred()) + defer func() { _ = os.RemoveAll(tmpDir) }() + opts.StoreDir = tmpDir + reconnectServer := natstest.RunServer(&opts) + defer reconnectServer.Shutdown() + + // IsConnected should transition to true after reconnection + Eventually(func() bool { return client.IsConnected() }, 10*time.Second, 100*time.Millisecond).Should(BeTrue()) + + client.Stop() + }) +}) + +var _ = Describe("Acknowledgment", Label("integration"), func() { + var ( + ctx context.Context + cancel context.CancelFunc + testConn *nats.Conn + testJS jetstream.JetStream + topicName string + logger *slog.Logger + ) + + BeforeEach(func() { + var err error + topicName = fmt.Sprintf("test-%s", uuid.New().String()[:8]) + logger = slog.Default() + ctx, cancel = context.WithTimeout(context.Background(), 30*time.Second) //nolint:fatcontext // Ginkgo BeforeEach pattern + + testConn, err = nats.Connect(testNATSServer.ClientURL()) + Expect(err).NotTo(HaveOccurred()) + testJS, err = jetstream.New(testConn) + Expect(err).NotTo(HaveOccurred()) + }) + + AfterEach(func() { + cancel() + deleteStreams(testJS, topicName) + testConn.Close() + }) + + It("acks on nil handler return, naks on error (IT-MSG-110)", func() { + cfg := messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + } + + callCount := make(chan int, 10) + handlerBlock := make(chan struct{}) + var invocation atomic.Int32 + client := messaging.NewClient(cfg, logger) + client.SetMainHandler(func(_ context.Context, _ []byte) error { + current := int(invocation.Add(1)) + callCount <- current + if current == 1 { + <-handlerBlock + return fmt.Errorf("simulated error") + } + return nil // Second invocation succeeds (ack) + }) + Expect(client.Start(ctx)).To(Succeed()) + defer client.Stop() + + publishCE(ctx, testJS, topicName, "dcm.test.ack", "dcm/test", + map[string]string{"key": "ack-test"}) + + // First delivery — handler is blocked (message in-flight, unacknowledged) + Eventually(callCount, 5*time.Second).Should(Receive(Equal(1))) + + // Unblock handler — returns error → nak → JetStream redelivers + close(handlerBlock) + + // Second delivery after nak — handler returns nil → ack + Eventually(callCount, 10*time.Second).Should(Receive(Equal(2))) + }) +}) + +var _ = Describe("CloudEvent Correlation", Label("integration"), func() { + var ( + ctx context.Context + cancel context.CancelFunc + testConn *nats.Conn + testJS jetstream.JetStream + topicName string + logger *slog.Logger + ) + + BeforeEach(func() { + var err error + topicName = fmt.Sprintf("test-%s", uuid.New().String()[:8]) + logger = slog.Default() + ctx, cancel = context.WithTimeout(context.Background(), 30*time.Second) //nolint:fatcontext // Ginkgo BeforeEach pattern + + testConn, err = nats.Connect(testNATSServer.ClientURL()) + Expect(err).NotTo(HaveOccurred()) + testJS, err = jetstream.New(testConn) + Expect(err).NotTo(HaveOccurred()) + }) + + AfterEach(func() { + cancel() + deleteStreams(testJS, topicName) + testConn.Close() + }) + + It("response CE conforms to CloudEvents v1.0 with agentName and topicName in data (IT-MSG-120)", func() { + cfg := messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + } + + client := messaging.NewClient(cfg, logger) + client.SetMainHandler(func(_ context.Context, _ []byte) error { + return nil + }) + Expect(client.Start(ctx)).To(Succeed()) + defer client.Stop() + + responseSub, err := testConn.SubscribeSync("dcm.agents.responses") + Expect(err).NotTo(HaveOccurred()) + + publishCE(ctx, testJS, topicName, "dcm.command.create", "dcm/control-plane", + map[string]string{"resourceId": "res-corr"}) + + msg, err := responseSub.NextMsg(5 * time.Second) + Expect(err).NotTo(HaveOccurred()) + + var respEvent cloudevents.Event + Expect(json.Unmarshal(msg.Data, &respEvent)).To(Succeed()) + + // CloudEvents v1.0 compliance + Expect(respEvent.SpecVersion()).To(Equal("1.0")) + Expect(respEvent.ID()).NotTo(BeEmpty()) + Expect(respEvent.Source()).NotTo(BeEmpty()) + Expect(respEvent.Type()).NotTo(BeEmpty()) + Expect(respEvent.Time().IsZero()).To(BeFalse()) + + // Correlation fields in data + var payload map[string]interface{} + Expect(json.Unmarshal(respEvent.Data(), &payload)).To(Succeed()) + Expect(payload).To(HaveKey("agentName")) + Expect(payload).To(HaveKey("topicName")) + Expect(payload["agentName"]).To(Equal("test-agent")) + Expect(payload["topicName"]).To(Equal(topicName)) + }) + + It("delete request produces deletion-acknowledged response with DELETING status (IT-MSG-072)", func() { + cfg := messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "test-agent", + } + + client := messaging.NewClient(cfg, logger) + client.SetMainHandler(func(_ context.Context, _ []byte) error { + return nil + }) + Expect(client.Start(ctx)).To(Succeed()) + defer client.Stop() + + responseSub, err := testConn.SubscribeSync("dcm.agents.responses") + Expect(err).NotTo(HaveOccurred()) + + publishCE(ctx, testJS, topicName, "dcm.request.delete", "dcm/control-plane", + map[string]string{"resourceId": "res-del-001"}) + + msg, err := responseSub.NextMsg(5 * time.Second) + Expect(err).NotTo(HaveOccurred()) + + var respEvent cloudevents.Event + Expect(json.Unmarshal(msg.Data, &respEvent)).To(Succeed()) + Expect(respEvent.Type()).To(Equal("dcm.agent.deletion-acknowledged")) + + var respPayload map[string]interface{} + Expect(json.Unmarshal(respEvent.Data(), &respPayload)).To(Succeed()) + Expect(respPayload["status"]).To(Equal("DELETING")) + Expect(respPayload["resourceId"]).To(Equal("res-del-001")) + Expect(respPayload["agentName"]).To(Equal("test-agent")) + Expect(respPayload["topicName"]).To(Equal(topicName)) + }) + + It("publishResponseCE failure causes nak and redelivery (IT-MSG-073)", func() { + cfg := messaging.ClientConfig{ + URL: testNATSServer.ClientURL(), + TopicName: topicName, + AgentName: "", // Empty AgentName → FormatSource error → publishResponseCE fails + } + + var deliveryCount atomic.Int32 + client := messaging.NewClient(cfg, logger) + client.SetMainHandler(func(_ context.Context, _ []byte) error { + deliveryCount.Add(1) + return nil + }) + Expect(client.Start(ctx)).To(Succeed()) + defer client.Stop() + + publishCE(ctx, testJS, topicName, "dcm.command.create", "dcm/control-plane", + map[string]string{"resourceId": "res-nak-001"}) + + // Message should be redelivered because publishResponseCE fails (empty AgentName) + Eventually(deliveryCount.Load, 10*time.Second, 100*time.Millisecond). + Should(BeNumerically(">=", int32(2))) + }) +}) diff --git a/internal/messaging/handlers.go b/internal/messaging/handlers.go new file mode 100644 index 0000000..0a9a891 --- /dev/null +++ b/internal/messaging/handlers.go @@ -0,0 +1,128 @@ +package messaging + +import ( + "context" + "encoding/json" + "fmt" + "strings" + "time" + + cloudevents "github.com/cloudevents/sdk-go/v2" + "github.com/nats-io/nats.go/jetstream" + + "github.com/dcm-project/environment-agent/internal/cloudevent" +) + +const nakDelay = 500 * time.Millisecond + +type ceResourcePayload struct { + ResourceID string `json:"resourceId"` +} + +func (c *Client) handleCancelMessage(msg jetstream.Msg) { + // TODO(topic8): MISSING-1 — acking malformed CE violates REQ-MSG-115; should nak instead + var event cloudevents.Event + if err := json.Unmarshal(msg.Data(), &event); err != nil { + c.logger.Warn("failed to parse cancel CE", "error", err) + if c.cancelHandler != nil { + _ = c.cancelHandler(context.Background(), msg.Data()) + } + return + } + + var payload ceResourcePayload + _ = json.Unmarshal(event.Data(), &payload) + if payload.ResourceID != "" { + c.denyList.Store(payload.ResourceID, struct{}{}) + } + + if c.cancelHandler != nil { + _ = c.cancelHandler(context.Background(), msg.Data()) + } +} + +func (c *Client) handleMainMessage(msg jetstream.Msg) { + // TODO(topic8): MISSING-1 — acking malformed CE violates REQ-MSG-115; should nak instead + var event cloudevents.Event + if err := json.Unmarshal(msg.Data(), &event); err != nil { + c.logger.Warn("failed to parse main CE", "error", err) + if c.mainHandler != nil { + if herr := c.mainHandler(context.Background(), msg.Data()); herr != nil { + _ = msg.NakWithDelay(nakDelay) + return + } + } + _ = msg.Ack() + return + } + + var payload ceResourcePayload + _ = json.Unmarshal(event.Data(), &payload) + resourceID := payload.ResourceID + + if _, loaded := c.denyList.LoadAndDelete(resourceID); loaded { + _ = msg.Ack() + return + } + + if c.mainHandler != nil { + if err := c.mainHandler(context.Background(), msg.Data()); err != nil { + _ = msg.NakWithDelay(nakDelay) + return + } + } + + if err := c.publishResponseCE(event.Type(), resourceID); err != nil { + c.logger.Error("failed to publish response CE", "error", err) + _ = msg.NakWithDelay(nakDelay) + return + } + + _ = msg.Ack() +} + +func (c *Client) publishResponseCE(incomingType, resourceID string) error { + var responseType string + switch { + case strings.Contains(incomingType, "create"): + responseType = TypeCreationAcked + case strings.Contains(incomingType, "delete"): + responseType = TypeDeletionAcked + default: + responseType = TypeAcked + } + + var status string + switch responseType { + case TypeCreationAcked: + status = "PROVISIONING" + case TypeDeletionAcked: + status = "DELETING" + default: + status = "ACKNOWLEDGED" + } + + event, err := cloudevent.NewCloudEvent(c.cfg.AgentName, responseType) + if err != nil { + return fmt.Errorf("failed to create response CE: %w", err) + } + if err := event.SetData(cloudevents.ApplicationJSON, map[string]interface{}{ + "agentName": c.cfg.AgentName, + "topicName": c.topics.Main, + "resourceId": resourceID, + "status": status, + }); err != nil { + return fmt.Errorf("failed to set response CE data: %w", err) + } + + data, err := json.Marshal(event) + if err != nil { + return fmt.Errorf("failed to marshal response CE: %w", err) + } + + c.mu.Lock() + conn := c.conn + c.mu.Unlock() + + return conn.Publish(SubjectResponses, data) +} diff --git a/internal/messaging/suite_test.go b/internal/messaging/suite_test.go new file mode 100644 index 0000000..ee3f64d --- /dev/null +++ b/internal/messaging/suite_test.go @@ -0,0 +1,41 @@ +package messaging_test + +import ( + "os" + "testing" + + natsserver "github.com/nats-io/nats-server/v2/server" + natstest "github.com/nats-io/nats-server/v2/test" + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +var testNATSServer *natsserver.Server + +var testStoreDir string + +func TestMessaging(t *testing.T) { + RegisterFailHandler(Fail) + RunSpecs(t, "Messaging Suite") +} + +var _ = BeforeSuite(func() { + var err error + testStoreDir, err = os.MkdirTemp("", "nats-test-*") + Expect(err).NotTo(HaveOccurred()) + + opts := natstest.DefaultTestOptions + opts.Port = -1 + opts.JetStream = true + opts.StoreDir = testStoreDir + testNATSServer = natstest.RunServer(&opts) +}) + +var _ = AfterSuite(func() { + if testNATSServer != nil { + testNATSServer.Shutdown() + } + if testStoreDir != "" { + _ = os.RemoveAll(testStoreDir) + } +}) diff --git a/internal/messaging/topics.go b/internal/messaging/topics.go new file mode 100644 index 0000000..333095d --- /dev/null +++ b/internal/messaging/topics.go @@ -0,0 +1,44 @@ +// Package messaging provides NATS/JetStream messaging client and topic management. +package messaging + +import ( + "errors" + "regexp" +) + +var validTopicNameRe = regexp.MustCompile(`^[A-Za-z0-9._-]+$`) + +// TopicNames holds the derived main, retry, and cancel topic names. +type TopicNames struct { + Main string + Retry string + Cancel string +} + +// DeriveTopicNames derives main/retry/cancel topic names from the agent name +// and an optional override. If topicNameOverride is non-empty, it takes precedence. +func DeriveTopicNames(agentName, topicNameOverride string) TopicNames { + base := agentName + if topicNameOverride != "" { + base = topicNameOverride + } + return TopicNames{ + Main: base, + Retry: base + ".retry", + Cancel: base + ".cancel", + } +} + +// ValidateTopicName validates that the given name conforms to NATS subject token rules. +func ValidateTopicName(name string) error { + if name == "" { + return errors.New("topic name must not be empty") + } + if len(name) > 255 { + return errors.New("topic name exceeds 255 characters") + } + if !validTopicNameRe.MatchString(name) { + return errors.New("topic name contains invalid characters (allowed: alphanumeric, hyphens, dots, underscores)") + } + return nil +} diff --git a/internal/messaging/topics_test.go b/internal/messaging/topics_test.go new file mode 100644 index 0000000..c3f2501 --- /dev/null +++ b/internal/messaging/topics_test.go @@ -0,0 +1,50 @@ +package messaging_test + +import ( + "strings" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + + "github.com/dcm-project/environment-agent/internal/messaging" +) + +var _ = Describe("DeriveTopicNames", Label("unit"), func() { + DescribeTable("derives main, retry, and cancel topic names", + func(agentName, override, expectMain, expectRetry, expectCancel string) { + result := messaging.DeriveTopicNames(agentName, override) + Expect(result.Main).To(Equal(expectMain)) + Expect(result.Retry).To(Equal(expectRetry)) + Expect(result.Cancel).To(Equal(expectCancel)) + }, + Entry("from agent name when no override (UT-MSG-010)", + "agent-prod-1", "", "agent-prod-1", "agent-prod-1.retry", "agent-prod-1.cancel"), + Entry("explicit topic name overrides agent name (UT-MSG-020)", + "agent-prod-1", "custom-topic", "custom-topic", "custom-topic.retry", "custom-topic.cancel"), + ) +}) + +var _ = Describe("ValidateTopicName", Label("unit"), func() { + DescribeTable("enforces NATS subject token rules", + func(name, errSubstring string) { + err := messaging.ValidateTopicName(name) + if errSubstring == "" { + Expect(err).NotTo(HaveOccurred()) + } else { + Expect(err).To(MatchError(ContainSubstring(errSubstring))) + } + }, + Entry("rejects space (UT-MSG-030)", "agent prod", "invalid characters"), + Entry("rejects wildcard * (UT-MSG-031)", "agent.*", "invalid characters"), + Entry("rejects full wildcard > (UT-MSG-032)", "agent>", "invalid characters"), + Entry("rejects exceeding 255 chars (UT-MSG-033)", strings.Repeat("a", 256), "255"), + Entry("accepts dot separator (UT-MSG-034)", "agent-prod.1", ""), + Entry("accepts alphanum + hyphens (UT-MSG-035)", "agent-prod-1", ""), + Entry("rejects empty string (UT-MSG-036)", "", "empty"), + Entry("rejects ! (UT-MSG-037)", "test!", "invalid"), + Entry("rejects @ (UT-MSG-038)", "test@", "invalid"), + Entry("rejects / (UT-MSG-039)", "test/a", "invalid"), + Entry("accepts underscore (UT-MSG-040)", "test_topic", ""), + Entry("accepts 255-char boundary (UT-MSG-041)", strings.Repeat("a", 255), ""), + ) +})