diff --git a/.github/workflows/docker-publish.yml b/.github/workflows/docker-publish.yml index 9a5cdb4..112e105 100644 --- a/.github/workflows/docker-publish.yml +++ b/.github/workflows/docker-publish.yml @@ -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: . diff --git a/Makefile b/Makefile index 9e9a102..7d3950f 100644 --- a/Makefile +++ b/Makefile @@ -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 @@ -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: @@ -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 #============================================================================== @@ -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 @@ -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 #============================================================================== @@ -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..." @@ -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 @@ -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 #============================================================================== @@ -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)') diff --git a/agent/common/Makefile b/agent/common/Makefile index 8697405..92f46ad 100644 --- a/agent/common/Makefile +++ b/agent/common/Makefile @@ -1,4 +1,4 @@ -cd.PHONY: test clean +.PHONY: test clean # Build settings BUILD_DIR=bin diff --git a/agent/common/controlplane/controlplane.go b/agent/common/controlplane/controlplane.go index 36731b8..9d778e0 100644 --- a/agent/common/controlplane/controlplane.go +++ b/agent/common/controlplane/controlplane.go @@ -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. "-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 @@ -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 diff --git a/agent/common/controlplane/controlplane_test.go b/agent/common/controlplane/controlplane_test.go index 8b47d45..2036a11 100644 --- a/agent/common/controlplane/controlplane_test.go +++ b/agent/common/controlplane/controlplane_test.go @@ -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") +} diff --git a/agent/lifecycle/keymanagementagent/siglet/siglet_api_client.go b/agent/common/siglet/siglet_api_client.go similarity index 100% rename from agent/lifecycle/keymanagementagent/siglet/siglet_api_client.go rename to agent/common/siglet/siglet_api_client.go diff --git a/agent/lifecycle/keymanagementagent/siglet/siglet_api_client_test.go b/agent/common/siglet/siglet_api_client_test.go similarity index 100% rename from agent/lifecycle/keymanagementagent/siglet/siglet_api_client_test.go rename to agent/common/siglet/siglet_api_client_test.go diff --git a/agent/common/siglet/transfer_type_client.go b/agent/common/siglet/transfer_type_client.go new file mode 100644 index 0000000..b53708c --- /dev/null +++ b/agent/common/siglet/transfer_type_client.go @@ -0,0 +1,213 @@ +/* + * Copyright (c) 2026 Metaform Systems, Inc. + * + * This program and the accompanying materials are made available under the + * terms of the Apache License, Version 2.0 which is available at + * https://www.apache.org/licenses/LICENSE-2.0 + * + * SPDX-License-Identifier: Apache-2.0 + * + * Contributors: + * Metaform Systems, Inc. - initial API and implementation + * + */ + +package siglet + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + + "github.com/eclipse-cfm/cfm/common/token" +) + +// EndpointMapping resolves a data endpoint dynamically from flow metadata (see Siglet's +// endpoint_mappings). Used as an alternative to a static Endpoint on a TransferType. +type EndpointMapping struct { + Key string `json:"key"` + Value string `json:"value"` + Endpoint string `json:"endpoint"` +} + +// TransferType describes how Siglet maps a single DPS transfer type (the flow profile) to a data +// endpoint and a token source. It mirrors a Siglet `[[transfer_types]]` block. +type TransferType struct { + // TransferType is the flow profile, e.g. "HttpData-PULL". + TransferType string `json:"transferType"` + // EndpointType is the backend type, e.g. "HTTP". + EndpointType string `json:"endpointType"` + // TokenSource is either "provider" or "client". Optional on input; defaults to "provider". + TokenSource string `json:"tokenSource"` + // Endpoint is the static data endpoint. Use this OR EndpointMappings, not both. + Endpoint string `json:"endpoint,omitempty"` + // EndpointMappings resolves the endpoint dynamically from flow metadata. + EndpointMappings []EndpointMapping `json:"endpointMappings,omitempty"` + // TxRenewalSupport indicates whether the transfer type supports transfer renewal. Defaults to false. + TxRenewalSupport bool `json:"txRenewalSupport"` +} + +// TransferTypeMapping is the complete set of transfer-type mappings for a single participant context. +// Configuring a context replaces its whole map (there is no per-entry patch). +type TransferTypeMapping struct { + ParticipantContextID string `json:"participantContextId"` + Mappings map[string]TransferType `json:"mappings"` +} + +// TransferTypeMappingClient interacts with Siglet's Management API to CRUD per-participant-context +// transfer-type mappings (the `/transfer-type-mappings` resource). See the Siglet runtime docs. +type TransferTypeMappingClient interface { + CreateTransferTypeMapping(ctx context.Context, mapping TransferTypeMapping) error + GetTransferTypeMapping(ctx context.Context, participantContextID string) (*TransferTypeMapping, error) + ReplaceTransferTypeMapping(ctx context.Context, mapping TransferTypeMapping) error + DeleteTransferTypeMapping(ctx context.Context, participantContextID string) error +} + +// NewTransferTypeMappingClient constructs a client for Siglet's transfer-type-mapping Management API. +func NewTransferTypeMappingClient(httpClient *http.Client, tokenProvider token.TokenProvider, url string) TransferTypeMappingClient { + return &httpManagementApiClient{ + httpClient: httpClient, + tokenProvider: tokenProvider, + BaseURL: url, + } +} + +func (s httpManagementApiClient) CreateTransferTypeMapping(ctx context.Context, mapping TransferTypeMapping) error { + accessToken, err := s.tokenProvider.GetToken(ctx, ScopeApiWrite, mapping.ParticipantContextID) + if err != nil { + return fmt.Errorf("failed to get API access token: %w", err) + } + + payload, err := json.Marshal(mapping) + if err != nil { + return fmt.Errorf("failed to marshal transfer type mapping: %w", err) + } + + url := fmt.Sprintf("%s/transfer-type-mappings", s.BaseURL) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewBuffer(payload)) + if err != nil { + return err + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer "+accessToken) + + resp, err := s.httpClient.Do(req) + if err != nil { + return fmt.Errorf("failed to create Siglet transfer type mapping: %w", err) + } + defer drainAndClose(resp) + + if resp.StatusCode != http.StatusCreated { + return fmt.Errorf("failed to create Siglet transfer type mapping: received status code %d", resp.StatusCode) + } + return nil +} + +func (s httpManagementApiClient) GetTransferTypeMapping(ctx context.Context, participantContextID string) (*TransferTypeMapping, error) { + accessToken, err := s.tokenProvider.GetToken(ctx, ScopeApiRead, participantContextID) + if err != nil { + return nil, fmt.Errorf("failed to get API access token: %w", err) + } + + url := fmt.Sprintf("%s/transfer-type-mappings/%s", s.BaseURL, participantContextID) + req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) + if err != nil { + return nil, fmt.Errorf("failed to create Siglet transfer type mapping request: %w", err) + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer "+accessToken) + + resp, err := s.httpClient.Do(req) + if err != nil { + return nil, fmt.Errorf("failed to get Siglet transfer type mapping: %w", err) + } + defer drainAndClose(resp) + + if resp.StatusCode == http.StatusNotFound { + return nil, nil + } + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("failed to get Siglet transfer type mapping: received status code %d", resp.StatusCode) + } + + body, err := io.ReadAll(resp.Body) + if err != nil { + return nil, fmt.Errorf("failed to read response body: %w", err) + } + + var mapping TransferTypeMapping + if err := json.Unmarshal(body, &mapping); err != nil { + return nil, fmt.Errorf("failed to unmarshal transfer type mapping: %w", err) + } + return &mapping, nil +} + +func (s httpManagementApiClient) ReplaceTransferTypeMapping(ctx context.Context, mapping TransferTypeMapping) error { + accessToken, err := s.tokenProvider.GetToken(ctx, ScopeApiWrite, mapping.ParticipantContextID) + if err != nil { + return fmt.Errorf("failed to get API access token: %w", err) + } + + payload, err := json.Marshal(mapping) + if err != nil { + return fmt.Errorf("failed to marshal transfer type mapping: %w", err) + } + + url := fmt.Sprintf("%s/transfer-type-mappings/%s", s.BaseURL, mapping.ParticipantContextID) + req, err := http.NewRequestWithContext(ctx, http.MethodPut, url, bytes.NewBuffer(payload)) + if err != nil { + return err + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer "+accessToken) + + resp, err := s.httpClient.Do(req) + if err != nil { + return fmt.Errorf("failed to replace Siglet transfer type mapping: %w", err) + } + defer drainAndClose(resp) + + if resp.StatusCode != http.StatusNoContent { + return fmt.Errorf("failed to replace Siglet transfer type mapping: received status code %d", resp.StatusCode) + } + return nil +} + +func (s httpManagementApiClient) DeleteTransferTypeMapping(ctx context.Context, participantContextID string) error { + accessToken, err := s.tokenProvider.GetToken(ctx, ScopeApiWrite, participantContextID) + if err != nil { + return fmt.Errorf("failed to get API access token: %w", err) + } + + url := fmt.Sprintf("%s/transfer-type-mappings/%s", s.BaseURL, participantContextID) + req, err := http.NewRequestWithContext(ctx, http.MethodDelete, url, nil) + if err != nil { + return fmt.Errorf("failed to create Siglet transfer type mapping request: %w", err) + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer "+accessToken) + + resp, err := s.httpClient.Do(req) + if err != nil { + return fmt.Errorf("failed to delete Siglet transfer type mapping: %w", err) + } + defer drainAndClose(resp) + + // treat an already-absent mapping as success so dispose is idempotent + if resp.StatusCode == http.StatusNotFound { + return nil + } + if resp.StatusCode != http.StatusNoContent { + return fmt.Errorf("error deleting Siglet transfer type mapping: received status code %d", resp.StatusCode) + } + return nil +} + +// drainAndClose drains and closes the response body to avoid connection/resource leaks. +func drainAndClose(resp *http.Response) { + _, _ = io.Copy(io.Discard, resp.Body) + _ = resp.Body.Close() +} diff --git a/agent/common/siglet/transfer_type_client_test.go b/agent/common/siglet/transfer_type_client_test.go new file mode 100644 index 0000000..89fc380 --- /dev/null +++ b/agent/common/siglet/transfer_type_client_test.go @@ -0,0 +1,205 @@ +/* + * Copyright (c) 2026 Metaform Systems, Inc. + * + * This program and the accompanying materials are made available under the + * terms of the Apache License, Version 2.0 which is available at + * https://www.apache.org/licenses/LICENSE-2.0 + * + * SPDX-License-Identifier: Apache-2.0 + * + * Contributors: + * Metaform Systems, Inc. - initial API and implementation + * + */ + +package siglet + +import ( + "encoding/json" + "fmt" + "io" + "net/http" + "net/http/httptest" + "testing" + + "github.com/eclipse-cfm/cfm/common/mocks" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" +) + +func sampleMapping() TransferTypeMapping { + return TransferTypeMapping{ + ParticipantContextID: "participant-1", + Mappings: map[string]TransferType{ + "HttpData-PULL": { + TransferType: "HttpData-PULL", + EndpointType: "HTTP", + TokenSource: "provider", + Endpoint: "https://data.provider.example.com/assets", + TxRenewalSupport: true, + }, + }, + } +} + +func TestHttpApiClient_CreateTransferTypeMapping(t *testing.T) { + var received TransferTypeMapping + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodPost && r.URL.Path == "/transfer-type-mappings" { + assert.Equal(t, "Bearer test token", r.Header.Get("Authorization")) + assert.Equal(t, "application/json", r.Header.Get("Content-Type")) + body, err := io.ReadAll(r.Body) + require.NoError(t, err) + require.NoError(t, json.Unmarshal(body, &received)) + w.WriteHeader(http.StatusCreated) + } else { + w.WriteHeader(http.StatusNotFound) + } + })) + defer server.Close() + + client := NewTransferTypeMappingClient(&http.Client{}, newTokenProvider(t), server.URL) + err := client.CreateTransferTypeMapping(t.Context(), sampleMapping()) + require.NoError(t, err) + assert.Equal(t, "participant-1", received.ParticipantContextID) + assert.Equal(t, "HTTP", received.Mappings["HttpData-PULL"].EndpointType) + assert.True(t, received.Mappings["HttpData-PULL"].TxRenewalSupport) +} + +func TestHttpApiClient_CreateTransferTypeMapping_AuthError(t *testing.T) { + tp := mocks.NewMockTokenProvider(t) + tp.On("GetToken", mock.Anything, mock.Anything, mock.Anything).Return("", fmt.Errorf("auth boom")) + client := NewTransferTypeMappingClient(&http.Client{}, tp, "http://foo.bar") + + err := client.CreateTransferTypeMapping(t.Context(), sampleMapping()) + require.ErrorContains(t, err, "auth boom") +} + +func TestHttpApiClient_CreateTransferTypeMapping_ApiReturnsError(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusConflict) + })) + defer server.Close() + + client := NewTransferTypeMappingClient(&http.Client{}, newTokenProvider(t), server.URL) + err := client.CreateTransferTypeMapping(t.Context(), sampleMapping()) + require.ErrorContains(t, err, "received status code 409") +} + +func TestHttpApiClient_GetTransferTypeMapping(t *testing.T) { + expected := sampleMapping() + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodGet && r.URL.Path == "/transfer-type-mappings/participant-1" { + assert.Equal(t, "Bearer test token", r.Header.Get("Authorization")) + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + _ = json.NewEncoder(w).Encode(expected) + } else { + w.WriteHeader(http.StatusInternalServerError) + } + })) + defer server.Close() + + client := NewTransferTypeMappingClient(&http.Client{}, newTokenProvider(t), server.URL) + result, err := client.GetTransferTypeMapping(t.Context(), "participant-1") + require.NoError(t, err) + require.NotNil(t, result) + assert.Equal(t, expected, *result) +} + +func TestHttpApiClient_GetTransferTypeMapping_NotFound(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusNotFound) + })) + defer server.Close() + + client := NewTransferTypeMappingClient(&http.Client{}, newTokenProvider(t), server.URL) + result, err := client.GetTransferTypeMapping(t.Context(), "participant-1") + require.NoError(t, err) + assert.Nil(t, result) +} + +func TestHttpApiClient_GetTransferTypeMapping_InvalidJson(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte("not valid json")) + })) + defer server.Close() + + client := NewTransferTypeMappingClient(&http.Client{}, newTokenProvider(t), server.URL) + result, err := client.GetTransferTypeMapping(t.Context(), "participant-1") + require.ErrorContains(t, err, "failed to unmarshal transfer type mapping") + assert.Nil(t, result) +} + +func TestHttpApiClient_ReplaceTransferTypeMapping(t *testing.T) { + var received TransferTypeMapping + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodPut && r.URL.Path == "/transfer-type-mappings/participant-1" { + assert.Equal(t, "Bearer test 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.StatusNoContent) + } else { + w.WriteHeader(http.StatusNotFound) + } + })) + defer server.Close() + + client := NewTransferTypeMappingClient(&http.Client{}, newTokenProvider(t), server.URL) + err := client.ReplaceTransferTypeMapping(t.Context(), sampleMapping()) + require.NoError(t, err) + assert.Equal(t, "provider", received.Mappings["HttpData-PULL"].TokenSource) +} + +func TestHttpApiClient_ReplaceTransferTypeMapping_ApiReturnsError(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + defer server.Close() + + client := NewTransferTypeMappingClient(&http.Client{}, newTokenProvider(t), server.URL) + err := client.ReplaceTransferTypeMapping(t.Context(), sampleMapping()) + require.ErrorContains(t, err, "received status code 500") +} + +func TestHttpApiClient_DeleteTransferTypeMapping(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodDelete && r.URL.Path == "/transfer-type-mappings/participant-1" { + assert.Equal(t, "Bearer test token", r.Header.Get("Authorization")) + w.WriteHeader(http.StatusNoContent) + } else { + w.WriteHeader(http.StatusInternalServerError) + } + })) + defer server.Close() + + client := NewTransferTypeMappingClient(&http.Client{}, newTokenProvider(t), server.URL) + err := client.DeleteTransferTypeMapping(t.Context(), "participant-1") + require.NoError(t, err) +} + +func TestHttpApiClient_DeleteTransferTypeMapping_NotFoundIsSuccess(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusNotFound) + })) + defer server.Close() + + client := NewTransferTypeMappingClient(&http.Client{}, newTokenProvider(t), server.URL) + err := client.DeleteTransferTypeMapping(t.Context(), "participant-1") + require.NoError(t, err) +} + +func TestHttpApiClient_DeleteTransferTypeMapping_ApiReturnsError(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusBadRequest) + })) + defer server.Close() + + client := NewTransferTypeMappingClient(&http.Client{}, newTokenProvider(t), server.URL) + err := client.DeleteTransferTypeMapping(t.Context(), "participant-1") + require.ErrorContains(t, err, "received status code 400") +} diff --git a/agent/lifecycle/keymanagementagent/handler/handler.go b/agent/lifecycle/keymanagementagent/handler/handler.go index f2a2138..12511ef 100644 --- a/agent/lifecycle/keymanagementagent/handler/handler.go +++ b/agent/lifecycle/keymanagementagent/handler/handler.go @@ -19,7 +19,7 @@ import ( "context" "github.com/eclipse-cfm/cfm/agent/common/controlplane" - "github.com/eclipse-cfm/cfm/agent/lifecycle/keymanagementagent/siglet" + "github.com/eclipse-cfm/cfm/agent/common/siglet" "github.com/eclipse-cfm/cfm/common/lifecycleagent" "github.com/eclipse-cfm/cfm/common/system" ) diff --git a/agent/lifecycle/keymanagementagent/launcher/launcher.go b/agent/lifecycle/keymanagementagent/launcher/launcher.go index 236aa9f..28c8c2c 100644 --- a/agent/lifecycle/keymanagementagent/launcher/launcher.go +++ b/agent/lifecycle/keymanagementagent/launcher/launcher.go @@ -18,8 +18,8 @@ import ( "net/http" "github.com/eclipse-cfm/cfm/agent/common/controlplane" + "github.com/eclipse-cfm/cfm/agent/common/siglet" "github.com/eclipse-cfm/cfm/agent/lifecycle/keymanagementagent/handler" - "github.com/eclipse-cfm/cfm/agent/lifecycle/keymanagementagent/siglet" "github.com/eclipse-cfm/cfm/assembly/httpclient" "github.com/eclipse-cfm/cfm/assembly/serviceapi" "github.com/eclipse-cfm/cfm/common/lifecycleagent" diff --git a/agent/orchestration/edcv/Makefile b/agent/orchestration/edcv/Makefile index fa5cd76..2427052 100644 --- a/agent/orchestration/edcv/Makefile +++ b/agent/orchestration/edcv/Makefile @@ -1,4 +1,4 @@ -cd.PHONY: build test clean run server +.PHONY: build test clean run server # Binary name SERVER_BINARY=edcvagent diff --git a/agent/orchestration/jwtletagent/activity/activity.go b/agent/orchestration/jwtletagent/activity/activity.go index a5cd90b..1126497 100644 --- a/agent/orchestration/jwtletagent/activity/activity.go +++ b/agent/orchestration/jwtletagent/activity/activity.go @@ -27,7 +27,7 @@ import ( "github.com/eclipse-cfm/cfm/agent/common/controlplane" "github.com/eclipse-cfm/cfm/agent/common/identityhub" "github.com/eclipse-cfm/cfm/agent/common/issuerservice" - "github.com/eclipse-cfm/cfm/agent/lifecycle/keymanagementagent/siglet" + "github.com/eclipse-cfm/cfm/agent/common/siglet" "github.com/eclipse-cfm/cfm/common/system" "github.com/eclipse-cfm/cfm/common/token" "github.com/eclipse-cfm/cfm/pmanager/api" diff --git a/agent/orchestration/onboarding/Makefile b/agent/orchestration/onboarding/Makefile index c4bfaa1..a861280 100644 --- a/agent/orchestration/onboarding/Makefile +++ b/agent/orchestration/onboarding/Makefile @@ -1,4 +1,4 @@ -cd.PHONY: build test clean run server +.PHONY: build test clean run server # Binary name SERVER_BINARY=obagent diff --git a/agent/orchestration/registration/Makefile b/agent/orchestration/registration/Makefile index ed265ce..4dfc222 100644 --- a/agent/orchestration/registration/Makefile +++ b/agent/orchestration/registration/Makefile @@ -1,4 +1,4 @@ -cd.PHONY: build test clean run server +.PHONY: build test clean run server # Binary name SERVER_BINARY=regagent diff --git a/agent/orchestration/siglet/Makefile b/agent/orchestration/siglet/Makefile new file mode 100644 index 0000000..4af8c08 --- /dev/null +++ b/agent/orchestration/siglet/Makefile @@ -0,0 +1,48 @@ +.PHONY: build test clean run server + +# Binary name +SERVER_BINARY=sigletagent + +# Build settings +BUILD_DIR=bin +SERVER_PATH=./cmd/server/main.go + +DOCKER_TAG=latest + +# Environment variables +export CGO_ENABLED=0 + +# Install development tools +install-tools: + +# Build the application +build: build-server + +# Build the server +build-server: + go build -o $(BUILD_DIR)/$(SERVER_BINARY) $(SERVER_PATH) + +# Run tests +test: + $(TEST_CMD) + +# Clean build artifacts +clean: + rm -rf $(BUILD_DIR) + go clean + +# Run the server in development mode +dev-server: build-server + ./$(BUILD_DIR)/$(SERVER_BINARY) + +# Run the server in production mode +server: build-server + ./$(BUILD_DIR)/$(SERVER_BINARY) + +# Build for multiple platforms +build-all: + # Server binaries + GOOS=linux GOARCH=amd64 go build -o $(BUILD_DIR)/$(SERVER_BINARY)-linux-amd64 $(SERVER_PATH) + GOOS=darwin GOARCH=amd64 go build -o $(BUILD_DIR)/$(SERVER_BINARY)-darwin-amd64 $(SERVER_PATH) + GOOS=darwin GOARCH=arm64 go build -o $(BUILD_DIR)/$(SERVER_BINARY)-darwin-arm64 $(SERVER_PATH) + GOOS=windows GOARCH=amd64 go build -o $(BUILD_DIR)/$(SERVER_BINARY)-windows-amd64.exe $(SERVER_PATH) diff --git a/agent/orchestration/siglet/activity/activity.go b/agent/orchestration/siglet/activity/activity.go new file mode 100644 index 0000000..7e7a207 --- /dev/null +++ b/agent/orchestration/siglet/activity/activity.go @@ -0,0 +1,225 @@ +// Copyright (c) 2026 Metaform Systems, Inc +// +// This program and the accompanying materials are made available under the +// terms of the Apache License, Version 2.0 which is available at +// https://www.apache.org/licenses/LICENSE-2.0 +// +// SPDX-License-Identifier: Apache-2.0 +// +// Contributors: +// Metaform Systems, Inc. - initial API and implementation +// + +// Package activity implements the deploy/dispose activity for the Siglet data-plane agent. On deploy +// it reads the transfer-type mappings from the cfm.dataplane VPA properties, configures them in +// Siglet for the participant context, and registers the Siglet data-plane instance with the control +// plane. On dispose it reverses both operations. +package activity + +import ( + "context" + "encoding/json" + "fmt" + "strings" + + "github.com/eclipse-cfm/cfm/agent/common/controlplane" + "github.com/eclipse-cfm/cfm/agent/common/siglet" + . "github.com/eclipse-cfm/cfm/common/collection" + "github.com/eclipse-cfm/cfm/common/model" + "github.com/eclipse-cfm/cfm/common/system" + "github.com/eclipse-cfm/cfm/pmanager/api" + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/trace" +) + +// TransferTypeMappingsKey is the key in the cfm.dataplane VPA properties that carries the map of +// transfer-type mappings (transferType -> mapping) to configure in Siglet. +const TransferTypeMappingsKey = "transferTypeMappings" + +// dataPlaneIDSuffix is appended to the participant context id to derive the data-plane instance id. +const dataPlaneIDSuffix = "-siglet" + +// defaultTokenSource is applied to a transfer-type mapping whose tokenSource is left unset. +const defaultTokenSource = "provider" + +// dataflowsPathTemplate is the Siglet DPS signaling path the control plane sends flow events to. The +// single verb is the participant context id. It is appended to the configured Siglet signaling URL. +const dataflowsPathTemplate = "/api/v1/%s/dataflows" + +type Config struct { + system.LogMonitor + // SigletSignalingURL is the base URL of the Siglet signaling API (scheme://host:port). The + // per-participant DPS endpoint registered with the control plane is derived from it. + SigletSignalingURL string + TransferTypeMappingClient siglet.TransferTypeMappingClient + DataPlaneClient controlplane.DataPlaneRegistrationClient +} + +type SigletActivityProcessor struct { + api.BaseActivityProcessor + monitor system.LogMonitor + sigletSignalingURL string + transferTypeMappingClient siglet.TransferTypeMappingClient + dataPlaneClient controlplane.DataPlaneRegistrationClient + tracer trace.Tracer +} + +func NewProcessor(config *Config) *SigletActivityProcessor { + return &SigletActivityProcessor{ + monitor: config.LogMonitor, + sigletSignalingURL: config.SigletSignalingURL, + transferTypeMappingClient: config.TransferTypeMappingClient, + dataPlaneClient: config.DataPlaneClient, + tracer: otel.GetTracerProvider().Tracer("cfm.agent.siglet"), + } +} + +type sigletData struct { + ParticipantContextId string `json:"participantContextId" validate:"required"` +} + +func (p SigletActivityProcessor) ProcessDeploy(ctx api.ActivityContext) api.ActivityResult { + spanCtx, span := p.tracer.Start(ctx.Context(), "cfm.agent.siglet.deploy") + defer span.End() + + props, err := ctx.VpaProperties(model.DataPlaneType) + if err != nil { + span.RecordError(err) + return api.ActivityResult{Result: api.ActivityResultFatalError, Error: fmt.Errorf("error reading data plane VPA properties for orchestration %s: %w", ctx.OID(), err)} + } + + // If there is no data plane VPA, or it carries no transfer-type mappings, there is nothing to do. + mappings, err := extractMappings(props) + if err != nil { + span.RecordError(err) + return api.ActivityResult{Result: api.ActivityResultFatalError, Error: fmt.Errorf("error parsing transfer type mappings for orchestration %s: %w", ctx.OID(), err)} + } + if len(mappings) == 0 { + p.monitor.Infof("No data plane transfer type mappings found for orchestration %s; nothing to configure", ctx.OID()) + return api.ActivityResult{Result: api.ActivityResultComplete} + } + + var data sigletData + if err := ctx.ReadValues(&data); err != nil { + span.RecordError(err) + return api.ActivityResult{Result: api.ActivityResultFatalError, Error: fmt.Errorf("error processing Siglet activity for orchestration %s: %w", ctx.OID(), err)} + } + participantContextId := data.ParticipantContextId + span.SetAttributes(attribute.String("cfm.participantContextId", participantContextId)) + + return p.handleDeployAction(spanCtx, participantContextId, mappings) +} + +func (p SigletActivityProcessor) ProcessDispose(ctx api.ActivityContext) api.ActivityResult { + var data sigletData + if err := ctx.ReadValues(&data); err != nil { + return api.ActivityResult{Result: api.ActivityResultFatalError, Error: fmt.Errorf("error processing Siglet activity for orchestration %s: %w", ctx.OID(), err)} + } + return p.handleDisposeAction(ctx.Context(), data.ParticipantContextId) +} + +// handleDeployAction configures the transfer-type mappings in Siglet (upsert) and registers the +// Siglet data-plane instance with the control plane. +func (p SigletActivityProcessor) handleDeployAction(ctx context.Context, participantContextId string, mappings map[string]siglet.TransferType) api.ActivityResult { + mapping := siglet.TransferTypeMapping{ + ParticipantContextID: participantContextId, + Mappings: mappings, + } + + // upsert: replace when a mapping already exists, otherwise create + existing, err := p.transferTypeMappingClient.GetTransferTypeMapping(ctx, participantContextId) + if err != nil { + return api.ActivityResult{Result: api.ActivityResultFatalError, Error: fmt.Errorf("cannot read transfer type mapping from Siglet: %w", err)} + } + if existing != nil { + if err := p.transferTypeMappingClient.ReplaceTransferTypeMapping(ctx, mapping); err != nil { + return api.ActivityResult{Result: api.ActivityResultFatalError, Error: fmt.Errorf("cannot replace transfer type mapping in Siglet: %w", err)} + } + } else { + if err := p.transferTypeMappingClient.CreateTransferTypeMapping(ctx, mapping); err != nil { + return api.ActivityResult{Result: api.ActivityResultFatalError, Error: fmt.Errorf("cannot create transfer type mapping in Siglet: %w", err)} + } + } + + if err := p.dataPlaneClient.RegisterDataPlane(ctx, participantContextId, p.dataPlaneRegistration(participantContextId, mappings)); err != nil { + return api.ActivityResult{Result: api.ActivityResultFatalError, Error: fmt.Errorf("cannot register data plane in control plane: %w", err)} + } + + p.monitor.Infof("Siglet activity for participant '%s' completed successfully", participantContextId) + return api.ActivityResult{Result: api.ActivityResultComplete} +} + +// handleDisposeAction removes the transfer-type mappings from Siglet and unregisters the data-plane +// instance from the control plane. Errors are logged but not propagated, so a failure does not block +// rollback of sibling agents. +func (p SigletActivityProcessor) handleDisposeAction(ctx context.Context, participantContextId string) api.ActivityResult { + var errors []error + + if err := p.transferTypeMappingClient.DeleteTransferTypeMapping(ctx, participantContextId); err != nil { + errors = append(errors, err) + } + if err := p.dataPlaneClient.UnregisterDataPlane(ctx, participantContextId, dataPlaneID(participantContextId)); err != nil { + errors = append(errors, err) + } + + if len(errors) > 0 { + errStrings := Collect(Map(From(errors), func(err error) string { return err.Error() })) + p.monitor.Warnf("one or more errors occurred while disposing Siglet data plane for '%s': [%s]", participantContextId, strings.Join(errStrings, ", ")) + } + return api.ActivityResult{Result: api.ActivityResultComplete} +} + +// extractMappings reads and decodes the transfer-type mappings from the data plane VPA properties. +// It returns an empty map (not an error) when the properties or the mappings key are absent. +func extractMappings(props map[string]any) (map[string]siglet.TransferType, error) { + if props == nil { + return nil, nil + } + raw, ok := props[TransferTypeMappingsKey] + if !ok || raw == nil { + return nil, nil + } + // round-trip through JSON to decode the untyped property bag into typed mappings + encoded, err := json.Marshal(raw) + if err != nil { + return nil, err + } + var mappings map[string]siglet.TransferType + if err := json.Unmarshal(encoded, &mappings); err != nil { + return nil, err + } + // tokenSource is optional on input; default it to "provider" when unset + for name, mapping := range mappings { + if mapping.TokenSource == "" { + mapping.TokenSource = defaultTokenSource + mappings[name] = mapping + } + } + return mappings, nil +} + +// dataPlaneRegistration builds the control-plane registration for the Siglet data plane. The transfer +// types are derived from the configured mappings and the endpoint is the participant-scoped Siglet +// DPS signaling endpoint. +func (p SigletActivityProcessor) dataPlaneRegistration(participantContextId string, mappings map[string]siglet.TransferType) controlplane.DataPlaneRegistration { + transferTypes := make([]string, 0, len(mappings)) + for transferType := range mappings { + transferTypes = append(transferTypes, transferType) + } + return controlplane.DataPlaneRegistration{ + ID: dataPlaneID(participantContextId), + TransferTypes: transferTypes, + Endpoint: p.dataPlaneEndpoint(participantContextId), + } +} + +// dataPlaneEndpoint builds the participant-scoped Siglet DPS signaling endpoint the control plane +// sends flow events to, e.g. http://siglet...:8081/api/v1//dataflows. +func (p SigletActivityProcessor) dataPlaneEndpoint(participantContextId string) string { + return strings.TrimRight(p.sigletSignalingURL, "/") + fmt.Sprintf(dataflowsPathTemplate, participantContextId) +} + +func dataPlaneID(participantContextId string) string { + return participantContextId + dataPlaneIDSuffix +} diff --git a/agent/orchestration/siglet/activity/activity_test.go b/agent/orchestration/siglet/activity/activity_test.go new file mode 100644 index 0000000..1a6bfa9 --- /dev/null +++ b/agent/orchestration/siglet/activity/activity_test.go @@ -0,0 +1,293 @@ +// Copyright (c) 2026 Metaform Systems, Inc +// +// This program and the accompanying materials are made available under the +// terms of the Apache License, Version 2.0 which is available at +// https://www.apache.org/licenses/LICENSE-2.0 +// +// SPDX-License-Identifier: Apache-2.0 +// +// Contributors: +// Metaform Systems, Inc. - initial API and implementation +// + +package activity + +import ( + "context" + "fmt" + "testing" + + "github.com/eclipse-cfm/cfm/agent/common/controlplane" + "github.com/eclipse-cfm/cfm/agent/common/siglet" + "github.com/eclipse-cfm/cfm/common/model" + "github.com/eclipse-cfm/cfm/common/system" + "github.com/eclipse-cfm/cfm/pmanager/api" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type ConfigOptions func(*Config) + +func WithSiglet(client siglet.TransferTypeMappingClient) ConfigOptions { + return func(config *Config) { config.TransferTypeMappingClient = client } +} + +func WithDataPlane(client controlplane.DataPlaneRegistrationClient) ConfigOptions { + return func(config *Config) { config.DataPlaneClient = client } +} + +func validConfig(opts ...ConfigOptions) *Config { + c := Config{ + LogMonitor: system.NoopMonitor{}, + SigletSignalingURL: "http://siglet.edc-v.svc.cluster.local:8081", + TransferTypeMappingClient: &MockSigletClient{}, + DataPlaneClient: &MockDataPlaneClient{}, + } + for _, opt := range opts { + opt(&c) + } + return &c +} + +func mappingsProps() map[string]any { + return map[string]any{ + TransferTypeMappingsKey: map[string]any{ + "HttpData-PULL": map[string]any{ + "transferType": "HttpData-PULL", + "endpointType": "HTTP", + "tokenSource": "provider", + "endpoint": "https://data.provider.example.com/assets", + }, + }, + } +} + +// validVpaData returns a VPA data slice with a single dataplane entry. If properties is nil, the +// entry has no properties field. +func validVpaData(properties map[string]any) []any { + entry := map[string]any{ + "vpaType": model.DataPlaneType.String(), + } + if properties != nil { + entry["properties"] = properties + } + return []any{entry} +} + +func processingDataWith(vpaData []any) map[string]any { + pd := map[string]any{ + "participantContextId": "participant-1", + } + if vpaData != nil { + pd[model.VPAData] = vpaData + } + return pd +} + +func newContext(pd map[string]any, discriminator api.Discriminator) api.ActivityContext { + activity := api.Activity{ID: "test-activity", Type: "siglet", Discriminator: discriminator} + return api.NewActivityContext(context.Background(), "orch-123", activity, pd, make(map[string]any)) +} + +func TestSiglet_Deploy_HappyPath(t *testing.T) { + sigletClient := &MockSigletClient{} + dpClient := &MockDataPlaneClient{} + processor := NewProcessor(validConfig(WithSiglet(sigletClient), WithDataPlane(dpClient))) + + result := processor.ProcessDeploy(newContext(processingDataWith(validVpaData(mappingsProps())), api.DeployDiscriminator)) + + assert.Equal(t, api.ActivityResultType(api.ActivityResultComplete), result.Result) + assert.NoError(t, result.Error) + assert.True(t, sigletClient.created, "expected a transfer type mapping to be created") + require.NotNil(t, sigletClient.lastMapping) + assert.Equal(t, "participant-1", sigletClient.lastMapping.ParticipantContextID) + assert.Contains(t, sigletClient.lastMapping.Mappings, "HttpData-PULL") + assert.True(t, dpClient.registered, "expected the data plane to be registered") + assert.Equal(t, "participant-1-siglet", dpClient.lastRegistration.ID) + assert.Equal(t, []string{"HttpData-PULL"}, dpClient.lastRegistration.TransferTypes) + assert.Equal(t, "http://siglet.edc-v.svc.cluster.local:8081/api/v1/participant-1/dataflows", dpClient.lastRegistration.Endpoint) +} + +func TestSiglet_Deploy_DefaultsTokenSourceAndRenewal(t *testing.T) { + sigletClient := &MockSigletClient{} + processor := NewProcessor(validConfig(WithSiglet(sigletClient))) + + props := map[string]any{ + TransferTypeMappingsKey: map[string]any{ + // tokenSource omitted -> should default to "provider"; txRenewalSupport omitted -> false + "HttpData-PULL": map[string]any{ + "transferType": "HttpData-PULL", + "endpointType": "HTTP", + "endpoint": "https://data.provider.example.com/assets", + }, + }, + } + + result := processor.ProcessDeploy(newContext(processingDataWith(validVpaData(props)), api.DeployDiscriminator)) + + assert.Equal(t, api.ActivityResultType(api.ActivityResultComplete), result.Result) + require.NotNil(t, sigletClient.lastMapping) + tt := sigletClient.lastMapping.Mappings["HttpData-PULL"] + assert.Equal(t, "provider", tt.TokenSource) + assert.False(t, tt.TxRenewalSupport) +} + +func TestSiglet_Deploy_UpsertReplacesExisting(t *testing.T) { + sigletClient := &MockSigletClient{existing: &siglet.TransferTypeMapping{ParticipantContextID: "participant-1"}} + processor := NewProcessor(validConfig(WithSiglet(sigletClient))) + + result := processor.ProcessDeploy(newContext(processingDataWith(validVpaData(mappingsProps())), api.DeployDiscriminator)) + + assert.Equal(t, api.ActivityResultType(api.ActivityResultComplete), result.Result) + assert.True(t, sigletClient.replaced, "expected an existing mapping to be replaced") + assert.False(t, sigletClient.created, "expected create not to be called when a mapping exists") +} + +func TestSiglet_Deploy_NoDataPlaneVpa_CompletesDoingNothing(t *testing.T) { + sigletClient := &MockSigletClient{} + dpClient := &MockDataPlaneClient{} + processor := NewProcessor(validConfig(WithSiglet(sigletClient), WithDataPlane(dpClient))) + + result := processor.ProcessDeploy(newContext(processingDataWith(nil), api.DeployDiscriminator)) + + assert.Equal(t, api.ActivityResultType(api.ActivityResultComplete), result.Result) + assert.NoError(t, result.Error) + assert.False(t, sigletClient.created) + assert.False(t, sigletClient.replaced) + assert.False(t, dpClient.registered) +} + +func TestSiglet_Deploy_NoMappings_CompletesDoingNothing(t *testing.T) { + sigletClient := &MockSigletClient{} + dpClient := &MockDataPlaneClient{} + processor := NewProcessor(validConfig(WithSiglet(sigletClient), WithDataPlane(dpClient))) + + // dataplane VPA present but no transfer-type mappings property + result := processor.ProcessDeploy(newContext(processingDataWith(validVpaData(map[string]any{})), api.DeployDiscriminator)) + + assert.Equal(t, api.ActivityResultType(api.ActivityResultComplete), result.Result) + assert.False(t, sigletClient.created) + assert.False(t, dpClient.registered) +} + +func TestSiglet_Deploy_MissingParticipantContextId(t *testing.T) { + processor := NewProcessor(validConfig()) + + pd := processingDataWith(validVpaData(mappingsProps())) + delete(pd, "participantContextId") + + result := processor.ProcessDeploy(newContext(pd, api.DeployDiscriminator)) + + assert.Equal(t, api.ActivityResultType(api.ActivityResultFatalError), result.Result) + require.Error(t, result.Error) + assert.Contains(t, result.Error.Error(), "orch-123") +} + +func TestSiglet_Deploy_SigletError(t *testing.T) { + processor := NewProcessor(validConfig(WithSiglet(&MockSigletClient{createErr: fmt.Errorf("siglet boom")}))) + + result := processor.ProcessDeploy(newContext(processingDataWith(validVpaData(mappingsProps())), api.DeployDiscriminator)) + + assert.Equal(t, api.ActivityResultType(api.ActivityResultFatalError), result.Result) + assert.ErrorContains(t, result.Error, "siglet boom") +} + +func TestSiglet_Deploy_ControlPlaneError(t *testing.T) { + processor := NewProcessor(validConfig(WithDataPlane(&MockDataPlaneClient{registerErr: fmt.Errorf("cp boom")}))) + + result := processor.ProcessDeploy(newContext(processingDataWith(validVpaData(mappingsProps())), api.DeployDiscriminator)) + + assert.Equal(t, api.ActivityResultType(api.ActivityResultFatalError), result.Result) + assert.ErrorContains(t, result.Error, "cp boom") +} + +func TestSiglet_Dispose_HappyPath(t *testing.T) { + sigletClient := &MockSigletClient{} + dpClient := &MockDataPlaneClient{} + processor := NewProcessor(validConfig(WithSiglet(sigletClient), WithDataPlane(dpClient))) + + result := processor.ProcessDispose(newContext(processingDataWith(nil), api.DisposeDiscriminator)) + + assert.Equal(t, api.ActivityResultType(api.ActivityResultComplete), result.Result) + assert.True(t, sigletClient.deleted) + assert.True(t, dpClient.unregistered) + assert.Equal(t, "participant-1-siglet", dpClient.lastDataPlaneID) +} + +func TestSiglet_Dispose_ErrorsStillComplete(t *testing.T) { + processor := NewProcessor(validConfig( + WithSiglet(&MockSigletClient{deleteErr: fmt.Errorf("siglet boom")}), + WithDataPlane(&MockDataPlaneClient{unregisterErr: fmt.Errorf("cp boom")}), + )) + + result := processor.ProcessDispose(newContext(processingDataWith(nil), api.DisposeDiscriminator)) + + // dispose must complete even when downstream clients error, so sibling rollback isn't blocked + assert.Equal(t, api.ActivityResultType(api.ActivityResultComplete), result.Result) +} + +func TestSiglet_Dispose_MissingParticipantContextId(t *testing.T) { + processor := NewProcessor(validConfig()) + + result := processor.ProcessDispose(newContext(map[string]any{}, api.DisposeDiscriminator)) + + assert.Equal(t, api.ActivityResultType(api.ActivityResultFatalError), result.Result) + require.Error(t, result.Error) +} + +// --- mocks --- + +type MockSigletClient struct { + existing *siglet.TransferTypeMapping + createErr error + replaceErr error + getErr error + deleteErr error + created bool + replaced bool + deleted bool + lastMapping *siglet.TransferTypeMapping +} + +func (m *MockSigletClient) CreateTransferTypeMapping(_ context.Context, mapping siglet.TransferTypeMapping) error { + m.created = true + m.lastMapping = &mapping + return m.createErr +} + +func (m *MockSigletClient) GetTransferTypeMapping(_ context.Context, _ string) (*siglet.TransferTypeMapping, error) { + return m.existing, m.getErr +} + +func (m *MockSigletClient) ReplaceTransferTypeMapping(_ context.Context, mapping siglet.TransferTypeMapping) error { + m.replaced = true + m.lastMapping = &mapping + return m.replaceErr +} + +func (m *MockSigletClient) DeleteTransferTypeMapping(_ context.Context, _ string) error { + m.deleted = true + return m.deleteErr +} + +type MockDataPlaneClient struct { + registerErr error + unregisterErr error + registered bool + unregistered bool + lastDataPlaneID string + lastRegistration controlplane.DataPlaneRegistration +} + +func (m *MockDataPlaneClient) RegisterDataPlane(_ context.Context, _ string, registration controlplane.DataPlaneRegistration) error { + m.registered = true + m.lastRegistration = registration + m.lastDataPlaneID = registration.ID + return m.registerErr +} + +func (m *MockDataPlaneClient) UnregisterDataPlane(_ context.Context, _ string, dataPlaneID string) error { + m.unregistered = true + m.lastDataPlaneID = dataPlaneID + return m.unregisterErr +} diff --git a/agent/orchestration/siglet/cmd/server/main.go b/agent/orchestration/siglet/cmd/server/main.go new file mode 100644 index 0000000..6a4bedc --- /dev/null +++ b/agent/orchestration/siglet/cmd/server/main.go @@ -0,0 +1,23 @@ +// Copyright (c) 2026 Metaform Systems, Inc +// +// This program and the accompanying materials are made available under the +// terms of the Apache License, Version 2.0 which is available at +// https://www.apache.org/licenses/LICENSE-2.0 +// +// SPDX-License-Identifier: Apache-2.0 +// +// Contributors: +// +// Metaform Systems, Inc. - initial API and implementation + +package main + +import ( + "github.com/eclipse-cfm/cfm/agent/orchestration/siglet/launcher" + "github.com/eclipse-cfm/cfm/common/runtime" +) + +// The entry point for the Siglet agent runtime. +func main() { + launcher.LaunchAndWaitSignal(runtime.CreateSignalShutdownChan()) +} diff --git a/agent/orchestration/siglet/launcher/launcher.go b/agent/orchestration/siglet/launcher/launcher.go new file mode 100644 index 0000000..5cb3094 --- /dev/null +++ b/agent/orchestration/siglet/launcher/launcher.go @@ -0,0 +1,89 @@ +// Copyright (c) 2026 Metaform Systems, Inc +// +// This program and the accompanying materials are made available under the +// terms of the Apache License, Version 2.0 which is available at +// https://www.apache.org/licenses/LICENSE-2.0 +// +// SPDX-License-Identifier: Apache-2.0 +// +// Contributors: +// Metaform Systems, Inc. - initial API and implementation +// + +package launcher + +import ( + "net/http" + + "github.com/eclipse-cfm/cfm/agent/common/controlplane" + "github.com/eclipse-cfm/cfm/agent/common/siglet" + "github.com/eclipse-cfm/cfm/agent/orchestration/siglet/activity" + "github.com/eclipse-cfm/cfm/assembly/httpclient" + "github.com/eclipse-cfm/cfm/assembly/serviceapi" + "github.com/eclipse-cfm/cfm/common/runtime" + "github.com/eclipse-cfm/cfm/common/system" + "github.com/eclipse-cfm/cfm/common/tokenexchange" + "github.com/eclipse-cfm/cfm/pmanager/api" + "github.com/eclipse-cfm/cfm/pmanager/natsagent" +) + +const ( + ActivityType = "siglet-activity" + sigletManagementURLKey = "siglet.management.url" + sigletSignalingURLKey = "siglet.signaling.url" + controlPlaneURLKey = "controlplane.url" + tokenExchangeURLKey = "tokenexchange.url" + tokenFilePathKey = "tokenexchange.tokenFilePath" + audienceKey = "tokenexchange.audience" +) + +func LaunchAndWaitSignal(shutdown <-chan struct{}) { + config := natsagent.LauncherConfig{ + AgentName: "Siglet Agent", + ServiceName: "cfm.agent.siglet", + ConfigPrefix: "sigletagent", + ActivityType: ActivityType, + AssemblyProvider: func() []system.ServiceAssembly { + return []system.ServiceAssembly{ + &httpclient.HttpClientServiceAssembly{}, + } + }, + NewProcessor: func(ctx *natsagent.AgentContext) api.ActivityProcessor { + httpClient := ctx.Registry.Resolve(serviceapi.HttpClientKey).(http.Client) + sigletManagementURL := ctx.Config.GetString(sigletManagementURLKey) + sigletSignalingURL := ctx.Config.GetString(sigletSignalingURLKey) + cpURL := ctx.Config.GetString(controlPlaneURLKey) + tokenExchangeURL := ctx.Config.GetString(tokenExchangeURLKey) + tokenFilePath := ctx.Config.GetString(tokenFilePathKey) + audience := ctx.Config.GetString(audienceKey) + + if err := runtime.CheckRequiredParams( + sigletManagementURLKey, sigletManagementURL, + sigletSignalingURLKey, sigletSignalingURL, + controlPlaneURLKey, cpURL, + tokenExchangeURLKey, tokenExchangeURL, + tokenFilePathKey, tokenFilePath, + audienceKey, audience, + ); err != nil { + panic(err) + } + + provider := tokenexchange.NewTokenExchangeProvider(tokenFilePath, + tokenexchange.WithTokenExchangeUrl(tokenExchangeURL), + tokenexchange.WithTokenExchangeAudience(audience), + tokenexchange.WithHttpClient(&httpClient)) + + return activity.NewProcessor(&activity.Config{ + LogMonitor: ctx.Monitor, + SigletSignalingURL: sigletSignalingURL, + TransferTypeMappingClient: siglet.NewTransferTypeMappingClient(&httpClient, provider, sigletManagementURL), + DataPlaneClient: controlplane.HttpManagementAPIClient{ + BaseURL: cpURL, + TokenProvider: provider, + HttpClient: &httpClient, + }, + }) + }, + } + natsagent.LaunchAgent(shutdown, config) +} diff --git a/docker/Dockerfile.sigletagent.dockerfile b/docker/Dockerfile.sigletagent.dockerfile new file mode 100644 index 0000000..658cda5 --- /dev/null +++ b/docker/Dockerfile.sigletagent.dockerfile @@ -0,0 +1,33 @@ +# Copyright (c) 2026 Metaform Systems, Inc +# +# This program and the accompanying materials are made available under the +# terms of the Apache License, Version 2.0 which is available at +# https://www.apache.org/licenses/LICENSE-2.0 +# +# SPDX-License-Identifier: Apache-2.0 +# +# Contributors: +# Metaform Systems, Inc. - initial API and implementation +# + +FROM --platform=$BUILDPLATFORM golang:1.25-alpine AS builder +ARG TARGETOS +ARG TARGETARCH + +WORKDIR /app + +COPY go.mod go.sum ./ +RUN go mod download + +COPY . . + +# Build the agent binary +RUN CGO_ENABLED=0 GOOS=$TARGETOS GOARCH=$TARGETARCH \ + go build -ldflags="-s -w" -o bin/sigletagent ./agent/orchestration/siglet/cmd/server/main.go + +# Production stage +FROM gcr.io/distroless/static-debian12:nonroot + +COPY --from=builder /app/bin/sigletagent /sigletagent + +ENTRYPOINT ["/sigletagent"] diff --git a/pmanager/Makefile b/pmanager/Makefile index 562e27a..c74b2a9 100644 --- a/pmanager/Makefile +++ b/pmanager/Makefile @@ -1,4 +1,4 @@ -cd.PHONY: build test clean run server dev-server +.PHONY: build test clean run server dev-server # Binary name SERVER_BINARY=pmanager