Commit 00e4ecb4 authored by Vasilis Katopodis's avatar Vasilis Katopodis
Browse files

fix: resolve kafka topics from service.yml

parent e98a7783
Loading
Loading
Loading
Loading
Loading
+1 −1
Original line number Diff line number Diff line
@@ -7,4 +7,4 @@ vendor/
CLAUDE.md
bin/
coverage.out
todo.md
+1 −1
Original line number Diff line number Diff line
@@ -18,7 +18,7 @@ The service is driven entirely by Kafka messages:
   across all its namespaces.
4. On every change (debounced), the session assesses each workload's health,
   aggregates it to one release-level verdict, and emits a `Diagnosis` report on
   `monitor-result`.
   `service-monitor-result`.
5. The session runs until a `monitor-stop` message — or a new request with the
   same `serviceId` — ends it. A final report is emitted on the way out.

+17 −16
Original line number Diff line number Diff line
@@ -14,17 +14,12 @@ import (
	"k8s.io/client-go/tools/clientcmd"

	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/base.go/pkg/adapters/kogito"
	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/base.go/pkg/communication"
	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/base.go/pkg/config"
	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/base.go/pkg/util"
	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/service.monitor/cmd/datamodels/messages"
	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/service.monitor/cmd/sessions"
)

const monitorResultTopic string = "monitor-result"
const kubeConfigRequestTopic string = "registry-retrieve-kubernetes-config"
const KubeConfigResponseTopic string = "registry-response-kubernetes-config-outgoing-channel"

// Header keys carrying resume state across the registry round trip (see RunKubeConfigResponse).
// serviceId doubles as the correlation key.
const (
@@ -39,21 +34,27 @@ type DiagnosisHandler struct {
	logger                   *util.Logger
	manager                  *sessions.Manager
	emitter                  *goka.Emitter
	monitorResultTopic       string
	kubeConfigRequestEmitter *goka.Emitter
	kubeConfigRequestTopic   string
	KubeConfigResponseTopic  string
}

func NewDiagnosisHandler(serviceLogger *util.Logger, _ *communication.TopicHandler, serviceConfig *config.BaseConfig, manager *sessions.Manager) *DiagnosisHandler {
func NewDiagnosisHandler(serviceLogger *util.Logger, serviceConfig *config.BaseConfig, manager *sessions.Manager, monitorResultTopic, kubeConfigRequestTopic, kubeConfigResponseTopic string) *DiagnosisHandler {
	emitter, err := goka.NewEmitter(serviceConfig.Kafka.Brokers, goka.Stream(monitorResultTopic), new(codec.String))
	serviceLogger.Fatal(err, "could not create emitter for "+monitorResultTopic)
	serviceLogger.Fatal(err, "could not create emitter for %s", monitorResultTopic)

	kubeConfigRequestEmitter, err := goka.NewEmitter(serviceConfig.Kafka.Brokers, goka.Stream(kubeConfigRequestTopic), new(codec.String))
	serviceLogger.Fatal(err, "could not create emitter for "+kubeConfigRequestTopic)
	serviceLogger.Fatal(err, "could not create emitter for %s", kubeConfigRequestTopic)

	return &DiagnosisHandler{
		logger:                   serviceLogger,
		manager:                  manager,
		emitter:                  emitter,
		monitorResultTopic:       monitorResultTopic,
		kubeConfigRequestEmitter: kubeConfigRequestEmitter,
		kubeConfigRequestTopic:   kubeConfigRequestTopic,
		KubeConfigResponseTopic:  kubeConfigResponseTopic,
	}
}

@@ -83,10 +84,10 @@ func (d *DiagnosisHandler) Run(ctx goka.Context, request kogito.RequestData[mess
	}
	promise.Then(func(err error) {
		if err != nil {
			d.logger.Warning("failed to emit to %s: %v", kubeConfigRequestTopic, err)
			d.logger.Warning("failed to emit to %s: %v", d.kubeConfigRequestTopic, err)
			return
		}
		d.logger.Debug("(ServiceId: %s) kubeconfig retrieve request delivered to %s", request.Data.ServiceId, kubeConfigRequestTopic)
		d.logger.Debug("(ServiceId: %s) kubeconfig retrieve request delivered to %s", request.Data.ServiceId, d.kubeConfigRequestTopic)
	})
}

@@ -145,7 +146,7 @@ func (d *DiagnosisHandler) RunKubeConfigResponse(ctx goka.Context, msg interface
	}

	d.logger.Info("(ServiceId: %s) received kubeconfig response on %s (id: %s, status: %s)",
		resumeState.ServiceID, KubeConfigResponseTopic, response.Id, response.Status)
		resumeState.ServiceID, d.KubeConfigResponseTopic, response.Id, response.Status)

	if response.Status != "SUCCESS" || response.Data.SecretData == nil {
		failureMessage := response.Message
@@ -200,26 +201,26 @@ func (d *DiagnosisHandler) emit(response messages.ResponseToDiagnosisRequest, pr
		SpecificationVersion: kogito.SpecificationVersion,
		ID:                   uuid.New().String(),
		Source:               source,
		Type:                 monitorResultTopic,
		Type:                 d.monitorResultTopic,
		ProccessID:           processID,
		Data:                 response,
	}

	jsonData, err := json.Marshal(kogitoResponse)
	if err != nil {
		d.logger.Warning("failed to encode response for %s: %v", monitorResultTopic, err)
		d.logger.Warning("failed to encode response for %s: %v", d.monitorResultTopic, err)
		return
	}
	d.logger.Debug("(ServiceId: %s) publishing to %s: %s", response.ServiceId, monitorResultTopic, string(jsonData))
	d.logger.Debug("(ServiceId: %s) publishing to %s: %s", response.ServiceId, d.monitorResultTopic, string(jsonData))

	promise, err := d.emitter.EmitWithHeaders("", string(jsonData), headers)
	if err != nil {
		d.logger.Warning("failed to emit to %s: %v", monitorResultTopic, err)
		d.logger.Warning("failed to emit to %s: %v", d.monitorResultTopic, err)
		return
	}
	promise.Then(func(err error) {
		if err != nil {
			d.logger.Warning("failed to emit to %s: %v", monitorResultTopic, err)
			d.logger.Warning("failed to emit to %s: %v", d.monitorResultTopic, err)
		}
	})
}
+1 −3
Original line number Diff line number Diff line
@@ -6,8 +6,6 @@ import (
	"github.com/lovoo/goka"

	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/base.go/pkg/adapters/kogito"
	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/base.go/pkg/communication"
	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/base.go/pkg/config"
	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/base.go/pkg/util"
	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/service.monitor/cmd/datamodels/messages"
	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/service.monitor/cmd/sessions"
@@ -18,7 +16,7 @@ type StopHandler struct {
	manager *sessions.Manager
}

func NewStopHandler(serviceLogger *util.Logger, _ *communication.TopicHandler, _ *config.BaseConfig, manager *sessions.Manager) *StopHandler {
func NewStopHandler(serviceLogger *util.Logger, manager *sessions.Manager) *StopHandler {
	return &StopHandler{logger: serviceLogger, manager: manager}
}

+35 −7
Original line number Diff line number Diff line
@@ -2,6 +2,7 @@ package cmd

import (
	"context"
	"fmt"

	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/base.go/pkg/adapters/kogito"
	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/base.go/pkg/communication"
@@ -19,6 +20,26 @@ type Lcm struct {
	done          chan bool
}

type streamConfig struct {
	In  string
	Out string
}

func (lcm *Lcm) resolveKafkaTopics() (map[string]streamConfig, error) {
	topicNames := []string{"monitor", "kubeconfig"}
	topics := make(map[string]streamConfig, len(topicNames))

	for _, name := range topicNames {
		in, out, err := lcm.topicHandler.GetTopics(name)
		if err != nil {
			return nil, fmt.Errorf("failed to get %s topics: %w", name, err)
		}
		topics[name] = streamConfig{In: in, Out: out}
	}

	return topics, nil
}

func NewLcm(currentLogger *util.Logger, currentTopicHandler *communication.TopicHandler, currentServiceConfig *config.BaseConfig) Lcm {
	return Lcm{
		logger:        currentLogger,
@@ -30,24 +51,31 @@ func NewLcm(currentLogger *util.Logger, currentTopicHandler *communication.Topic
func (lcm *Lcm) Init() {
	manager := sessions.NewManager()

	diagnosisHandler := handlers.NewDiagnosisHandler(lcm.logger, lcm.topicHandler, lcm.serviceConfig, manager)
	communication.BindTopicsToHandler(
	topics, err := lcm.resolveKafkaTopics()
	lcm.logger.Fatal(err, "failed to resolve kafka topics")

	diagnosisHandler := handlers.NewDiagnosisHandler(lcm.logger, lcm.serviceConfig, manager,
		topics["monitor"].Out, topics["kubeconfig"].In, topics["kubeconfig"].Out)
	err = communication.BindTopicsToHandler(
		lcm.topicHandler,
		"monitor",
		diagnosisHandler.Run,
		kogito.NewProcess[messages.DiagnosisRequest],
	)
	lcm.logger.Fatal(err, "failed to bind topics for key 'monitor'")

	err := lcm.topicHandler.AddInputStreamTopic(handlers.KubeConfigResponseTopic, diagnosisHandler.RunKubeConfigResponse)
	lcm.logger.Fatal(err, "could not add input stream topic "+handlers.KubeConfigResponseTopic)
	err = lcm.topicHandler.AddInputStreamTopic(diagnosisHandler.KubeConfigResponseTopic, diagnosisHandler.RunKubeConfigResponse)
	lcm.logger.Fatal(err, "could not add input stream topic %s", diagnosisHandler.KubeConfigResponseTopic)

	stopHandler := handlers.NewStopHandler(lcm.logger, lcm.topicHandler, lcm.serviceConfig, manager)
	communication.BindTopicsToHandler(
	stopHandler := handlers.NewStopHandler(lcm.logger, manager)
	if err := communication.BindTopicsToHandler(
		lcm.topicHandler,
		"monitor-stop",
		stopHandler.Run,
		kogito.NewProcess[messages.StopRequest],
	)
	); err != nil {
		lcm.logger.Error(err, "failed to bind topics for key 'monitor-stop'")
	}
}

func (lcm *Lcm) StartReceivers() {
Loading