Commit c9fbce81 authored by Vasilis Katopodis's avatar Vasilis Katopodis
Browse files

feat: fetch kubeconfig from registry

parent 6308173d
Loading
Loading
Loading
Loading
+20 −0
Original line number Diff line number Diff line
package messages

// RegistryData mirrors org.etsi.osl.hypo.registry.model.registry.RegistryData.
type RegistryData struct {
	Id         string                 `json:"id"`
	SecretData map[string]interface{} `json:"secretData"`
}

// KubeConfigResponse mirrors org.etsi.osl.hypo.registry.model.kafka.outgoing.VaultSecretResponseMessage,
// received on the registry-response-kubernetes-config-outgoing-channel topic.
type KubeConfigResponse struct {
	Id      string       `json:"id"`
	Data    RegistryData `json:"data"`
	Status  string       `json:"status"`
	Message string       `json:"message"`
}

// KubeConfigSecretKey is the Vault secret map key holding the base64-encoded kubeconfig,
// matching RegistryService.KUBE_CONFIG_KEY on the registry side.
const KubeConfigSecretKey = "kubeConfig"
+53 −0
Original line number Diff line number Diff line
package messages

import (
	"encoding/json"
	"testing"
)

func TestKubeConfigResponse_UnmarshalSuccess(t *testing.T) {
	raw := `{
		"id": "svc-1",
		"data": {
			"id": "hypo/kubernetes-config/telenor/svc-1",
			"secretData": {"kubeConfig": "base64content"}
		},
		"status": "SUCCESS",
		"message": ""
	}`

	var response KubeConfigResponse
	if err := json.Unmarshal([]byte(raw), &response); err != nil {
		t.Fatalf("unexpected error: %v", err)
	}

	if response.Status != "SUCCESS" {
		t.Errorf("expected status SUCCESS, got %q", response.Status)
	}
	if response.Data.Id != "hypo/kubernetes-config/telenor/svc-1" {
		t.Errorf("unexpected data.id: %q", response.Data.Id)
	}
	kubeConfig, ok := response.Data.SecretData[KubeConfigSecretKey].(string)
	if !ok || kubeConfig != "base64content" {
		t.Errorf("expected secretData[%q] = 'base64content', got %v", KubeConfigSecretKey, response.Data.SecretData[KubeConfigSecretKey])
	}
}

func TestKubeConfigResponse_UnmarshalFailure(t *testing.T) {
	raw := `{"id": "svc-1", "data": null, "status": "FAILURE", "message": "vault unreachable"}`

	var response KubeConfigResponse
	if err := json.Unmarshal([]byte(raw), &response); err != nil {
		t.Fatalf("unexpected error: %v", err)
	}

	if response.Status != "FAILURE" {
		t.Errorf("expected status FAILURE, got %q", response.Status)
	}
	if response.Message != "vault unreachable" {
		t.Errorf("unexpected message: %q", response.Message)
	}
	if response.Data.SecretData != nil {
		t.Errorf("expected nil secretData, got %v", response.Data.SecretData)
	}
}
+133 −21
Original line number Diff line number Diff line
@@ -2,6 +2,7 @@ package handlers

import (
	"context"
	"encoding/base64"
	"encoding/json"
	"fmt"

@@ -20,24 +21,55 @@ import (
)

const monitorResultTopic string = "monitor-result"
const kubeConfigRequestTopic string = "registry-retrieve-kubernetes-config"

// KubeConfigResponseTopic is bound to RunKubeConfigResponse in cmd/lcmStart.go.
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: it's already unique per in-flight diagnosis session
// (sessions.Manager enforces that), and since resume state travels fully in each message's own
// headers — no shared pending-request map — there's nothing for a separate correlation id to
// disambiguate.
const (
	headerNamespace       = "namespace"
	headerHelmReleaseName = "helmReleaseName"
	headerServiceID       = "serviceId"
	headerProcessID       = "processId"
)

type DiagnosisHandler struct {
	logger                   *util.Logger
	manager                  *sessions.Manager
	emitter                  *goka.Emitter
	kubeConfigRequestEmitter *goka.Emitter
}

// NewDiagnosisHandler creates the handler along with a long-lived emitter for monitor-result.
// The emitter is independent of any single goka.Context: a watch session keeps producing
// health reports long after the goka callback that started it (Run) has returned, and a
// goka.Context is only valid for the duration of that one callback.
// NewDiagnosisHandler creates the handler along with long-lived emitters for monitor-result and
// registry-retrieve-kubernetes-config. Both are independent of any single goka.Context: a watch
// session keeps producing health reports long after the goka callback that started it (Run) has
// returned, and a goka.Context is only valid for the duration of that one callback.
func NewDiagnosisHandler(serviceLogger *util.Logger, _ *communication.TopicHandler, serviceConfig *config.BaseConfig, manager *sessions.Manager) *DiagnosisHandler {
	emitter, err := goka.NewEmitter(serviceConfig.Kafka.Brokers, goka.Stream(monitorResultTopic), new(codec.String))
	serviceLogger.Fatal(err, "could not create emitter for "+monitorResultTopic)

	return &DiagnosisHandler{logger: serviceLogger, manager: manager, emitter: emitter}
	kubeConfigRequestEmitter, err := goka.NewEmitter(serviceConfig.Kafka.Brokers, goka.Stream(kubeConfigRequestTopic), new(codec.String))
	serviceLogger.Fatal(err, "could not create emitter for "+kubeConfigRequestTopic)

	return &DiagnosisHandler{
		logger:                   serviceLogger,
		manager:                  manager,
		emitter:                  emitter,
		kubeConfigRequestEmitter: kubeConfigRequestEmitter,
	}
}

// Run starts the kubeconfig retrieval for a diagnosis request. It does not build the Kubernetes
// client or start the watch session itself — goka.Context is only valid for this callback's
// lifetime, and the registry round trip is asynchronous, so blocking here would stall the whole
// partition processor. Instead it publishes the retrieve request with enough state riding on the
// Kafka headers (see resumeStateFromHeaders) for RunKubeConfigResponse to resume statelessly once
// the registry answers, from any pod.
func (d *DiagnosisHandler) Run(ctx goka.Context, request kogito.RequestData[messages.DiagnosisRequest]) {
	d.logger.Info(" ----- New Diagnosis request for release '%s' in namespace '%s' -----",
		request.Data.HelmReleaseName, request.Data.NameSpace)
@@ -45,37 +77,117 @@ func (d *DiagnosisHandler) Run(ctx goka.Context, request kogito.RequestData[mess
	headers := ctx.Headers()
	processID := request.ProccessID

	// TODO: fetch the kubeconfig from registry-api using request.Data.KubernetesId
	// instead of treating it as the kubeconfig itself.
	k8sCfg, err := clientcmd.RESTConfigFromKubeConfig([]byte(request.Data.KubernetesId))
	requestHeaders := buildKubeConfigRequestHeaders(headers, request.Data, processID)

	promise, err := d.kubeConfigRequestEmitter.EmitWithHeaders("", request.Data.KubernetesId, requestHeaders)
	if err != nil {
		d.logger.Error(err, "could not build config from kubeconfig")
		d.logger.Error(err, "failed to emit kubeconfig retrieve request")
		d.emit(messages.ResponseToDiagnosisRequest{ServiceId: request.Data.ServiceId, Status: "FAILURE", Message: err.Error()}, processID, headers)
		return
	}
	promise.Then(func(err error) {
		if err != nil {
			d.logger.Warning("failed to emit to %s: %v", kubeConfigRequestTopic, err)
		}
	})
}

// buildKubeConfigRequestHeaders copies the original request headers (notably Authorization,
// required by the registry's KafkaUtils.extractTokenFromKafkaMessage) and adds the resume state
// RunKubeConfigResponse needs to pick this diagnosis request back up statelessly. This is the
// send-side counterpart of resumeStateFromHeaders.
func buildKubeConfigRequestHeaders(original goka.Headers, data messages.DiagnosisRequest, processID string) goka.Headers {
	headers := make(goka.Headers, len(original)+4)
	for k, v := range original {
		headers[k] = v
	}
	headers[headerNamespace] = []byte(data.NameSpace)
	headers[headerHelmReleaseName] = []byte(data.HelmReleaseName)
	headers[headerServiceID] = []byte(data.ServiceId)
	headers[headerProcessID] = []byte(processID)
	return headers
}

// kubeconfigResumeState is extracted from the Kafka headers of a registry response, letting any
// pod resume a diagnosis request statelessly — no in-memory pending-request map is kept between
// Run and RunKubeConfigResponse.
type kubeconfigResumeState struct {
	Namespace       string
	HelmReleaseName string
	ServiceID       string
	ProcessID       string
}

func resumeStateFromHeaders(headers goka.Headers) kubeconfigResumeState {
	return kubeconfigResumeState{
		Namespace:       string(headers[headerNamespace]),
		HelmReleaseName: string(headers[headerHelmReleaseName]),
		ServiceID:       string(headers[headerServiceID]),
		ProcessID:       string(headers[headerProcessID]),
	}
}

// RunKubeConfigResponse resumes a diagnosis request once the registry answers on
// KubeConfigResponseTopic. On success it continues exactly where Run left off before the
// registry integration: build the REST config, start the watch session, and drain its health
// reports for as long as the session runs.
func (d *DiagnosisHandler) RunKubeConfigResponse(ctx goka.Context, msg interface{}) {
	headers := ctx.Headers()
	resumeState := resumeStateFromHeaders(headers)

	var response messages.KubeConfigResponse
	if err := json.Unmarshal([]byte(msg.(string)), &response); err != nil {
		d.logger.Error(err, "(ServiceId: %s) could not decode kubeconfig response", resumeState.ServiceID)
		d.emit(messages.ResponseToDiagnosisRequest{ServiceId: resumeState.ServiceID, Status: "FAILURE", Message: err.Error()}, resumeState.ProcessID, headers)
		return
	}

	if response.Status != "SUCCESS" || response.Data.SecretData == nil {
		failureMessage := response.Message
		if failureMessage == "" {
			failureMessage = "kubeconfig not found"
		}
		d.emit(messages.ResponseToDiagnosisRequest{ServiceId: resumeState.ServiceID, Status: "FAILURE", Message: failureMessage}, resumeState.ProcessID, headers)
		return
	}

	kubeConfigBase64, _ := response.Data.SecretData[messages.KubeConfigSecretKey].(string)
	kubeConfigBytes, err := base64.StdEncoding.DecodeString(kubeConfigBase64)
	if err != nil {
		d.logger.Error(err, "(ServiceId: %s) could not decode kubeconfig from base64", resumeState.ServiceID)
		d.emit(messages.ResponseToDiagnosisRequest{ServiceId: resumeState.ServiceID, Status: "FAILURE", Message: err.Error()}, resumeState.ProcessID, headers)
		return
	}

	k8sCfg, err := clientcmd.RESTConfigFromKubeConfig(kubeConfigBytes)
	if err != nil {
		d.logger.Error(err, "could not build config from kubeconfig")
		d.emit(messages.ResponseToDiagnosisRequest{ServiceId: resumeState.ServiceID, Status: "FAILURE", Message: err.Error()}, resumeState.ProcessID, headers)
		return
	}

	dynamicClient, err := dynamic.NewForConfig(k8sCfg)
	if err != nil {
		d.logger.Error(err, "could not create dynamic Kubernetes client")
		d.emit(messages.ResponseToDiagnosisRequest{ServiceId: request.Data.ServiceId, Status: "FAILURE", Message: err.Error()}, processID, headers)
		d.emit(messages.ResponseToDiagnosisRequest{ServiceId: resumeState.ServiceID, Status: "FAILURE", Message: err.Error()}, resumeState.ProcessID, headers)
		return
	}

	labelSelector := fmt.Sprintf("app.kubernetes.io/managed-by=Helm,app.kubernetes.io/instance=%s",
		request.Data.HelmReleaseName)
		resumeState.HelmReleaseName)

	session := sessions.New(request.Data.NameSpace, request.Data.ServiceId, labelSelector, dynamicClient, d.logger)
	session := sessions.New(resumeState.Namespace, resumeState.ServiceID, labelSelector, dynamicClient, d.logger)
	cancel := session.Start(context.Background())
	token := d.manager.Register(request.Data.ServiceId, cancel)
	token := d.manager.Register(resumeState.ServiceID, cancel)

	// The session keeps emitting health reports for as long as the watch runs, well past
	// this callback's lifetime. Draining it in a goroutine instead of blocking Run() on it
	// lets the shared partition-processor goroutine move on to the next message — notably
	// the StopRequest, which is the only thing that can end this session.
	// this callback's lifetime. Draining it in a goroutine instead of blocking on it lets the
	// shared partition-processor goroutine move on to the next message — notably the
	// StopRequest, which is the only thing that can end this session.
	go func() {
		defer d.manager.Deregister(request.Data.ServiceId, token)
		defer d.manager.Deregister(resumeState.ServiceID, token)
		for resp := range session.EmitCh() {
			d.emit(resp, processID, headers)
			d.emit(resp, resumeState.ProcessID, headers)
		}
	}()
}
+77 −0
Original line number Diff line number Diff line
package handlers

import (
	"testing"

	"github.com/lovoo/goka"

	"labs.etsi.org/rep/osl/hypo/code/org.etsi.osl.hypo.core/service.monitor/cmd/datamodels/messages"
)

func TestBuildKubeConfigRequestHeaders_SetsResumeState(t *testing.T) {
	original := goka.Headers{"Authorization": []byte("Bearer token")}
	data := messages.DiagnosisRequest{
		ServiceId:       "svc-1",
		KubernetesId:    "kubeconfig-id",
		NameSpace:       "ns-1",
		HelmReleaseName: "release-1",
	}

	headers := buildKubeConfigRequestHeaders(original, data, "proc-1")

	if string(headers["Authorization"]) != "Bearer token" {
		t.Errorf("expected Authorization header to be preserved, got %q", headers["Authorization"])
	}
	if string(headers[headerNamespace]) != "ns-1" {
		t.Errorf("expected namespace header 'ns-1', got %q", headers[headerNamespace])
	}
	if string(headers[headerHelmReleaseName]) != "release-1" {
		t.Errorf("expected helmReleaseName header 'release-1', got %q", headers[headerHelmReleaseName])
	}
	if string(headers[headerServiceID]) != "svc-1" {
		t.Errorf("expected serviceId header 'svc-1', got %q", headers[headerServiceID])
	}
	if string(headers[headerProcessID]) != "proc-1" {
		t.Errorf("expected processId header 'proc-1', got %q", headers[headerProcessID])
	}
}

func TestBuildKubeConfigRequestHeaders_DoesNotMutateOriginal(t *testing.T) {
	original := goka.Headers{"Authorization": []byte("Bearer token")}
	data := messages.DiagnosisRequest{ServiceId: "svc-1", NameSpace: "ns-1", HelmReleaseName: "release-1"}

	buildKubeConfigRequestHeaders(original, data, "proc-1")

	if len(original) != 1 {
		t.Errorf("expected original headers map to still have 1 entry, got %d", len(original))
	}
}

func TestResumeStateFromHeaders(t *testing.T) {
	headers := goka.Headers{
		headerNamespace:       []byte("ns-1"),
		headerHelmReleaseName: []byte("release-1"),
		headerServiceID:       []byte("svc-1"),
		headerProcessID:       []byte("proc-1"),
	}

	state := resumeStateFromHeaders(headers)

	want := kubeconfigResumeState{
		Namespace:       "ns-1",
		HelmReleaseName: "release-1",
		ServiceID:       "svc-1",
		ProcessID:       "proc-1",
	}
	if state != want {
		t.Errorf("resumeStateFromHeaders() = %+v, want %+v", state, want)
	}
}

func TestResumeStateFromHeaders_MissingKeysYieldEmptyStrings(t *testing.T) {
	state := resumeStateFromHeaders(goka.Headers{})

	if state != (kubeconfigResumeState{}) {
		t.Errorf("expected zero-value state for empty headers, got %+v", state)
	}
}
+3 −0
Original line number Diff line number Diff line
@@ -38,6 +38,9 @@ func (lcm *Lcm) Init() {
		kogito.NewProcess[messages.DiagnosisRequest],
	)

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

	stopHandler := handlers.NewStopHandler(lcm.logger, lcm.topicHandler, lcm.serviceConfig, manager)
	communication.BindTopicsToHandler(
		lcm.topicHandler,
Loading