Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .github/workflows/docker-publish.yml
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,9 @@ jobs:
- name: keymanagementagent
dockerfile: docker/Dockerfile.keymanagementagent.dockerfile
context: .
- name: sigletagent
dockerfile: docker/Dockerfile.sigletagent.dockerfile
context: .
- name: pmanager
dockerfile: docker/Dockerfile.pmanager.dockerfile
context: .
Expand Down
19 changes: 17 additions & 2 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ IH_DIR=agent/orchestration/ih
REG_DIR=agent/orchestration/registration
ONBOARDING_DIR=agent/orchestration/onboarding
JWTLET_AGENT_DIR=agent/orchestration/jwtletagent
SIGLET_AGENT_DIR=agent/orchestration/siglet
KEY_MANAGEMENT_AGENT_DIR=agent/lifecycle/keymanagementagent
AGENT_COMMON=agent/common
KIND_CLUSTER_NAME=edcv
Expand Down Expand Up @@ -72,6 +73,7 @@ build:
$(MAKE) -C $(REG_DIR) build
$(MAKE) -C $(ONBOARDING_DIR) build
$(MAKE) -C $(JWTLET_AGENT_DIR) build
$(MAKE) -C $(SIGLET_AGENT_DIR) build
$(MAKE) -C $(KEY_MANAGEMENT_AGENT_DIR) build

build-pmanager:
Expand All @@ -91,6 +93,7 @@ build-all:
$(MAKE) -C $(REG_DIR) build-all
$(MAKE) -C $(ONBOARDING_DIR) build-all
$(MAKE) -C $(JWTLET_AGENT_DIR) build-all
$(MAKE) -C $(SIGLET_AGENT_DIR) build-all
$(MAKE) -C $(KEY_MANAGEMENT_AGENT_DIR) build-all

#==============================================================================
Expand All @@ -109,6 +112,7 @@ test: install-gotestsum
$(MAKE) -C $(REG_DIR) test
$(MAKE) -C $(ONBOARDING_DIR) test
$(MAKE) -C $(JWTLET_AGENT_DIR) test
$(MAKE) -C $(SIGLET_AGENT_DIR) test
$(MAKE) -C $(KEY_MANAGEMENT_AGENT_DIR) test
$(MAKE) -C $(ASSEMBLY_DIR) test
$(MAKE) -C $(AGENT_COMMON) test
Expand Down Expand Up @@ -163,6 +167,7 @@ clean:
$(MAKE) -C $(REG_DIR) clean
$(MAKE) -C $(ONBOARDING_DIR) clean
$(MAKE) -C $(JWTLET_AGENT_DIR) clean
$(MAKE) -C $(SIGLET_AGENT_DIR) clean
$(MAKE) -C $(KEY_MANAGEMENT_AGENT_DIR) clean

#==============================================================================
Expand Down Expand Up @@ -194,7 +199,7 @@ generate-docs:
# Docker Commands - Handled at Top Level
#==============================================================================

docker-build: docker-build-pmanager docker-build-tmanager docker-build-jwtletagent docker-build-edcvagent docker-build-ihagent docker-build-regagent docker-build-obagent docker-build-keymanagementagent
docker-build: docker-build-pmanager docker-build-tmanager docker-build-jwtletagent docker-build-edcvagent docker-build-ihagent docker-build-regagent docker-build-obagent docker-build-keymanagementagent docker-build-sigletagent

docker-build-pmanager:
@echo "Building pmanager Docker image..."
Expand Down Expand Up @@ -228,7 +233,11 @@ docker-build-keymanagementagent:
@echo "Building Key Management agent Docker image..."
docker buildx build -f docker/Dockerfile.keymanagementagent.dockerfile -t $(DOCKER_REGISTRY)keymanagementagent:$(DOCKER_TAG) .

docker-clean: docker-clean-pmanager docker-clean-tmanager docker-clean-jwtletagent docker-clean-edcvagent docker-clean-ihagent docker-clean-regagent docker-clean-obagent docker-clean-keymanagementagent
docker-build-sigletagent:
@echo "Building Siglet agent Docker image..."
docker buildx build -f docker/Dockerfile.sigletagent.dockerfile -t $(DOCKER_REGISTRY)sigletagent:$(DOCKER_TAG) .

docker-clean: docker-clean-pmanager docker-clean-tmanager docker-clean-jwtletagent docker-clean-edcvagent docker-clean-ihagent docker-clean-regagent docker-clean-obagent docker-clean-keymanagementagent docker-clean-sigletagent

docker-clean-pmanager:
docker rmi $(DOCKER_REGISTRY)pmanager:$(DOCKER_TAG) || true
Expand All @@ -254,6 +263,9 @@ docker-clean-obagent:
docker-clean-keymanagementagent:
docker rmi $(DOCKER_REGISTRY)keymanagementagent:$(DOCKER_TAG) || true

docker-clean-sigletagent:
docker rmi $(DOCKER_REGISTRY)sigletagent:$(DOCKER_TAG) || true

#==============================================================================
# Load images into KinD Cluster
#==============================================================================
Expand Down Expand Up @@ -282,6 +294,9 @@ load-into-kind-jwtletagent: docker-build-jwtletagent
load-into-kind-keymanagementagent: docker-build-keymanagementagent
kind load docker-image -n $(KIND_CLUSTER_NAME) $(DOCKER_REGISTRY)keymanagementagent:$(DOCKER_TAG)

load-into-kind-sigletagent: docker-build-sigletagent
kind load docker-image -n $(KIND_CLUSTER_NAME) $(DOCKER_REGISTRY)sigletagent:$(DOCKER_TAG)

# builds and loads all images into KinD cluster. Will require kind to be installed and a kind cluster named KIND_CLUSTER_NAME running.
load-into-kind: docker-build
kind load docker-image -n $(KIND_CLUSTER_NAME) $$(docker images --format "{{.Repository}}:{{.Tag}}" | grep '^$(DOCKER_REGISTRY).*:$(DOCKER_TAG)')
Expand Down
2 changes: 1 addition & 1 deletion agent/common/Makefile
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
cd.PHONY: test clean
.PHONY: test clean

# Build settings
BUILD_DIR=bin
Expand Down
86 changes: 86 additions & 0 deletions agent/common/controlplane/controlplane.go
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,25 @@ type ManagementAPIClient interface {
DeleteParticipantContext(ctx context.Context, participantContextID string) error
}

// DataPlaneRegistration describes a data-plane instance to register with the control plane for a
// participant context. For a Siglet data plane, Endpoint is the DPS signaling endpoint and the
// transfer types are the ones configured as transfer-type mappings in Siglet.
type DataPlaneRegistration struct {
// ID is the unique identifier of the data-plane instance, e.g. "<participant>-siglet".
ID string `json:"dataplaneId"`
// TransferTypes are the transfer types the data plane supports, e.g. "HttpData-PULL".
TransferTypes []string `json:"transferTypes"`
// Endpoint is the data plane's DPS signaling endpoint the control plane sends flow events to.
Endpoint string `json:"endpoint"`
}

// DataPlaneRegistrationClient registers and unregisters data-plane instances with the EDC control
// plane, scoped to a participant context. HttpManagementAPIClient implements this interface.
type DataPlaneRegistrationClient interface {
RegisterDataPlane(ctx context.Context, participantContextID string, registration DataPlaneRegistration) error
UnregisterDataPlane(ctx context.Context, participantContextID string, dataPlaneID string) error
}

type HttpManagementAPIClient struct {
BaseURL string
TokenProvider token.TokenProvider
Expand Down Expand Up @@ -205,6 +224,73 @@ func (h HttpManagementAPIClient) PatchConfig(ctx context.Context, participantCon
return nil
}

// RegisterDataPlane registers a data-plane instance with the control plane for the given participant
// context via PUT /v5beta/participants/{participantContextID}/dataplanes.
func (h HttpManagementAPIClient) RegisterDataPlane(ctx context.Context, participantContextID string, registration DataPlaneRegistration) error {
accessToken, err := h.TokenProvider.GetToken(ctx, ScopeApiAdmin, participantContextID)
if err != nil {
return fmt.Errorf("failed to get API access token: %w", err)
}

payload, err := json.Marshal(registration)
if err != nil {
return err
}

url := fmt.Sprintf("%s%s/%s/dataplanes", h.BaseURL, CreateParticipantURL, participantContextID)
req, err := http.NewRequestWithContext(ctx, http.MethodPut, url, bytes.NewBuffer(payload))
if err != nil {
return err
}
req.Header.Set("Content-Type", applicationJSON)
req.Header.Set("Authorization", "Bearer "+accessToken)
resp, err := h.HttpClient.Do(req)
if err != nil {
return fmt.Errorf("failed to register data plane on control plane: %w", err)
}

defer h.closeResponse(resp)

if resp.StatusCode < http.StatusOK || resp.StatusCode >= http.StatusBadRequest {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("failed to register data plane on control plane: received status code %d, body: %s", resp.StatusCode, string(body))
}
return nil
}

// UnregisterDataPlane removes a previously registered data-plane instance from the control plane via
// DELETE /v5beta/participants/{participantContextID}/dataplanes/{dataPlaneID}.
func (h HttpManagementAPIClient) UnregisterDataPlane(ctx context.Context, participantContextID string, dataPlaneID string) error {
accessToken, err := h.TokenProvider.GetToken(ctx, ScopeApiAdmin, participantContextID)
if err != nil {
return fmt.Errorf("failed to get API access token: %w", err)
}

url := fmt.Sprintf("%s%s/%s/dataplanes/%s", h.BaseURL, CreateParticipantURL, participantContextID, dataPlaneID)
req, err := http.NewRequestWithContext(ctx, http.MethodDelete, url, nil)
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+accessToken)
resp, err := h.HttpClient.Do(req)
if err != nil {
return fmt.Errorf("failed to unregister data plane on control plane: %w", err)
}

defer h.closeResponse(resp)

switch {
case resp.StatusCode == http.StatusNotFound:
// treat an already-absent data plane as success so dispose is idempotent
return nil
case resp.StatusCode >= http.StatusOK && resp.StatusCode < http.StatusBadRequest:
return nil
default:
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("failed to unregister data plane on control plane: received status code %d, body: %s", resp.StatusCode, string(body))
}
}

func (h HttpManagementAPIClient) closeResponse(resp *http.Response) {
func() {
// drain and close response body to avoid connection/resource leak
Expand Down
100 changes: 100 additions & 0 deletions agent/common/controlplane/controlplane_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -370,3 +370,103 @@ func TestPatchConfig_ServerError(t *testing.T) {
err := client.PatchConfig(t.Context(), "test-participant", ParticipantContextConfig{ParticipantContextID: "test-participant"})
require.ErrorContains(t, err, "received status code 500")
}

func TestRegisterDataPlane(t *testing.T) {
var received map[string]any
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodPut && r.URL.Path == CreateParticipantURL+"/test-participant/dataplanes" {
require.Equal(t, "Bearer token", r.Header.Get("Authorization"))
body, err := io.ReadAll(r.Body)
require.NoError(t, err)
require.NoError(t, json.Unmarshal(body, &received))
w.WriteHeader(http.StatusOK)
} else {
w.WriteHeader(http.StatusNotFound)
}
}))
defer server.Close()

tp := mocks.NewMockTokenProvider(t)
tp.On("GetToken", mock.Anything, mock.Anything, mock.Anything).Return("token", nil)
client := HttpManagementAPIClient{BaseURL: server.URL, TokenProvider: tp, HttpClient: &http.Client{}}

err := client.RegisterDataPlane(t.Context(), "test-participant", DataPlaneRegistration{
ID: "test-participant-siglet",
TransferTypes: []string{"HttpData-PULL"},
Endpoint: "http://siglet.edc-v.svc.cluster.local:8081/api/v1/test-participant/dataflows",
})
require.NoError(t, err)
require.Equal(t, "test-participant-siglet", received["dataplaneId"])
require.Equal(t, "http://siglet.edc-v.svc.cluster.local:8081/api/v1/test-participant/dataflows", received["endpoint"])
}

func TestRegisterDataPlane_AuthError(t *testing.T) {
tp := mocks.NewMockTokenProvider(t)
tp.On("GetToken", mock.Anything, mock.Anything, mock.Anything).Return("", fmt.Errorf("test error"))
client := HttpManagementAPIClient{BaseURL: "http://foo.bar", TokenProvider: tp, HttpClient: &http.Client{}}

err := client.RegisterDataPlane(t.Context(), "test-participant", DataPlaneRegistration{ID: "test-participant-siglet"})
require.ErrorContains(t, err, "test error")
}

func TestRegisterDataPlane_ServerError(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
_, _ = w.Write([]byte("boom"))
}))
defer server.Close()

tp := mocks.NewMockTokenProvider(t)
tp.On("GetToken", mock.Anything, mock.Anything, mock.Anything).Return("token", nil)
client := HttpManagementAPIClient{BaseURL: server.URL, TokenProvider: tp, HttpClient: &http.Client{}}

err := client.RegisterDataPlane(t.Context(), "test-participant", DataPlaneRegistration{ID: "test-participant-siglet"})
require.ErrorContains(t, err, "received status code 500")
}

func TestUnregisterDataPlane(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodDelete && r.URL.Path == CreateParticipantURL+"/test-participant/dataplanes/siglet-1" {
require.Equal(t, "Bearer token", r.Header.Get("Authorization"))
w.WriteHeader(http.StatusOK)
} else {
w.WriteHeader(http.StatusInternalServerError)
}
}))
defer server.Close()

tp := mocks.NewMockTokenProvider(t)
tp.On("GetToken", mock.Anything, mock.Anything, mock.Anything).Return("token", nil)
client := HttpManagementAPIClient{BaseURL: server.URL, TokenProvider: tp, HttpClient: &http.Client{}}

err := client.UnregisterDataPlane(t.Context(), "test-participant", "siglet-1")
require.NoError(t, err)
}

func TestUnregisterDataPlane_NotFoundIsSuccess(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusNotFound)
}))
defer server.Close()

tp := mocks.NewMockTokenProvider(t)
tp.On("GetToken", mock.Anything, mock.Anything, mock.Anything).Return("token", nil)
client := HttpManagementAPIClient{BaseURL: server.URL, TokenProvider: tp, HttpClient: &http.Client{}}

err := client.UnregisterDataPlane(t.Context(), "test-participant", "siglet-1")
require.NoError(t, err)
}

func TestUnregisterDataPlane_ServerError(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
}))
defer server.Close()

tp := mocks.NewMockTokenProvider(t)
tp.On("GetToken", mock.Anything, mock.Anything, mock.Anything).Return("token", nil)
client := HttpManagementAPIClient{BaseURL: server.URL, TokenProvider: tp, HttpClient: &http.Client{}}

err := client.UnregisterDataPlane(t.Context(), "test-participant", "siglet-1")
require.ErrorContains(t, err, "received status code 500")
}
Loading