Skip to content

Commit df51fb0

Browse files
Feature/shared ingress uses config nats broker (#777)
* share ingress uses config map now * added unit test for ingress * refactored tests
1 parent 1c2af24 commit df51fb0

4 files changed

Lines changed: 355 additions & 14 deletions

File tree

cmd/ingress/main.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,18 +19,21 @@ package main
1919
import (
2020
"os"
2121

22+
"knative.dev/pkg/configmap"
2223
"knative.dev/pkg/injection"
2324
"knative.dev/pkg/injection/sharedmain"
2425
"knative.dev/pkg/signals"
2526

2627
"knative.dev/eventing-natss/pkg/broker/ingress"
28+
"knative.dev/eventing-natss/pkg/common/configloader/fsloader"
2729
)
2830

2931
func main() {
3032
component := "natsjs-broker-ingress"
3133

3234
ctx := signals.NewContext()
3335
ctx = sharedmain.WithHealthProbesDisabled(ctx)
36+
ctx = fsloader.WithLoader(ctx, configmap.Load)
3437

3538
ns := os.Getenv("NAMESPACE")
3639
if ns != "" {

config/broker/500-shared-ingress.yaml

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -42,10 +42,12 @@ spec:
4242
valueFrom:
4343
fieldRef:
4444
fieldPath: metadata.name
45-
- name: NATS_URL
46-
value: "nats://nats.nats-io.svc.cluster.local:4222"
4745
- name: PORT
4846
value: "8080"
47+
volumeMounts:
48+
- name: config-nats-broker
49+
mountPath: /etc/config-nats-broker
50+
readOnly: true
4951
ports:
5052
- containerPort: 8080
5153
name: http
@@ -81,6 +83,10 @@ spec:
8183
- ALL
8284
seccompProfile:
8385
type: RuntimeDefault
86+
volumes:
87+
- name: config-nats-broker
88+
configMap:
89+
name: config-nats-broker
8490
---
8591
apiVersion: v1
8692
kind: Service

pkg/broker/ingress/controller.go

Lines changed: 32 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ import (
2323
"time"
2424

2525
"github.com/kelseyhightower/envconfig"
26+
nats "github.com/nats-io/nats.go"
2627
"go.uber.org/zap"
2728
corev1 "k8s.io/api/core/v1"
2829
apierrs "k8s.io/apimachinery/pkg/api/errors"
@@ -36,7 +37,9 @@ import (
3637

3738
configmapinformer "knative.dev/pkg/client/injection/kube/informers/core/v1/configmap"
3839

40+
"knative.dev/eventing-natss/pkg/broker/constants"
3941
"knative.dev/eventing-natss/pkg/broker/contract"
42+
"knative.dev/eventing-natss/pkg/common/configloader/fsloader"
4043
commonnats "knative.dev/eventing-natss/pkg/common/nats"
4144
)
4245

@@ -49,29 +52,23 @@ const (
4952
)
5053

5154
type envConfig struct {
52-
NatsURL string `envconfig:"NATS_URL" required:"true"`
53-
Port int `envconfig:"PORT" default:"8080"`
55+
Port int `envconfig:"PORT" default:"8080"`
5456
}
5557

5658
// NewController creates a new shared ingress controller
5759
func NewController(ctx context.Context, _ configmap.Watcher) *controller.Impl {
5860
logger := logging.FromContext(ctx)
5961

60-
// Load environment configuration
61-
env := &envConfig{}
62-
if err := envconfig.Process("", env); err != nil {
62+
env, err := loadEnvConfig()
63+
if err != nil {
6364
logger.Fatalw("Failed to process environment variables", zap.Error(err))
6465
}
6566

66-
logger.Infow("Starting shared broker ingress",
67-
zap.String("nats_url", env.NatsURL),
68-
zap.Int("port", env.Port),
69-
)
67+
logger.Infow("Starting shared broker ingress", zap.Int("port", env.Port))
7068

71-
// Create NATS connection using URL from environment variable
72-
natsConn, err := commonnats.NewNatsConnFromURL(ctx, env.NatsURL)
69+
natsConn, err := buildNatsConn(ctx)
7370
if err != nil {
74-
logger.Fatalw("Failed to create NATS connection", zap.Error(err))
71+
logger.Fatalw("Failed to initialize NATS connection", zap.Error(err))
7572
}
7673

7774
// Create JetStream context
@@ -154,6 +151,29 @@ func NewController(ctx context.Context, _ configmap.Watcher) *controller.Impl {
154151
return impl
155152
}
156153

154+
// loadEnvConfig reads environment variables into envConfig.
155+
func loadEnvConfig() (*envConfig, error) {
156+
env := &envConfig{}
157+
return env, envconfig.Process("", env)
158+
}
159+
160+
// buildNatsConn loads NATS configuration from the mounted ConfigMap and creates a connection.
161+
func buildNatsConn(ctx context.Context) (*nats.Conn, error) {
162+
fsLoader, err := fsloader.Get(ctx)
163+
if err != nil {
164+
return nil, fmt.Errorf("failed to get ConfigMap loader from context: %w", err)
165+
}
166+
cmData, err := fsLoader(constants.SettingsConfigMapMountPath)
167+
if err != nil {
168+
return nil, fmt.Errorf("failed to load NATS ConfigMap: %w", err)
169+
}
170+
natsConfig, err := commonnats.LoadEventingNatsConfig(cmData)
171+
if err != nil {
172+
return nil, fmt.Errorf("failed to parse NATS configuration: %w", err)
173+
}
174+
return commonnats.NewNatsConn(ctx, natsConfig)
175+
}
176+
157177
// filterContractConfigMap returns true if the object is the contract ConfigMap
158178
func filterContractConfigMap(obj interface{}) bool {
159179
if obj == nil {

0 commit comments

Comments
 (0)