diff --git a/event-gateway/gateway-controller/cmd/controller/main.go b/event-gateway/gateway-controller/cmd/controller/main.go index c93fec239..c196b5e91 100644 --- a/event-gateway/gateway-controller/cmd/controller/main.go +++ b/event-gateway/gateway-controller/cmd/controller/main.go @@ -73,6 +73,8 @@ import ( eventgateway "github.com/wso2/api-platform/event-gateway/gateway-controller/pkg/api/eventgateway" eventgatewayconfig "github.com/wso2/api-platform/event-gateway/gateway-controller/pkg/config" + "github.com/wso2/api-platform/event-gateway/gateway-controller/pkg/controlplanehooks" + "github.com/wso2/api-platform/event-gateway/gateway-controller/pkg/dbschema" "github.com/wso2/api-platform/event-gateway/gateway-controller/pkg/eventlistener" "github.com/wso2/api-platform/event-gateway/gateway-controller/pkg/handler" "github.com/wso2/api-platform/event-gateway/gateway-controller/pkg/hubtopic" @@ -201,6 +203,14 @@ func main() { } defer db.Close() + // websub_apis, webbroker_apis, and webhook_secrets are event-gateway-specific + // tables that core's own schema scripts do not define. Apply them against the + // same database connection core just opened. + if err := dbschema.Apply(context.Background(), db.GetDB(), cfg.Controller.Storage.Type); err != nil { + log.Error("Failed to initialize event-gateway database schema", slog.Any("error", err)) + os.Exit(1) + } + var eventHubInstance eventhub.EventHub var eventHubStorage storage.Storage ehBackendCfg := toBackendConfig(cfg) @@ -471,6 +481,7 @@ func main() { lazyResourceXDSManager, templateDefinitions, subscriptionSnapshotManager, eventHubInstance, secretsService, webhookSecretStore, webhookSecretSnapshotManager, ) + cpClient.SetControlPlaneEventGatewayHooks(controlplanehooks.Hooks{}) if err := cpClient.Start(); err != nil { log.Error("Failed to start control plane client", slog.Any("error", err)) } diff --git a/event-gateway/gateway-controller/pkg/controlplanehooks/events.go b/event-gateway/gateway-controller/pkg/controlplanehooks/events.go new file mode 100644 index 000000000..fa82b003f --- /dev/null +++ b/event-gateway/gateway-controller/pkg/controlplanehooks/events.go @@ -0,0 +1,134 @@ +/* + * Copyright (c) 2026, WSO2 LLC. (https://www.wso2.com). + * + * WSO2 LLC. licenses this file to you under the Apache License, + * Version 2.0 (the "License"); you may not use this file except + * in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +// Package controlplanehooks implements gateway-controller (core)'s +// controlplane.ControlPlaneEventGatewayHooks interface, moved out of core's +// pkg/controlplane/client.go and pkg/controlplane/events.go. +package controlplanehooks + +import "time" + +// WebSubAPIDeployedEventPayload represents the payload of a WebSub API deployment event. +type WebSubAPIDeployedEventPayload struct { + APIID string `json:"apiId"` + DeploymentID string `json:"deploymentId"` + PerformedAt time.Time `json:"performedAt"` +} + +// WebSubAPIDeployedEvent represents the complete WebSub API deployment event. +type WebSubAPIDeployedEvent struct { + Type string `json:"type"` + Payload WebSubAPIDeployedEventPayload `json:"payload"` + Timestamp string `json:"timestamp"` + CorrelationID string `json:"correlationId"` +} + +// WebSubAPIUndeployedEventPayload represents the payload of a WebSub API undeployment event. +type WebSubAPIUndeployedEventPayload struct { + APIID string `json:"apiId"` + DeploymentID string `json:"deploymentId"` + PerformedAt time.Time `json:"performedAt"` +} + +// WebSubAPIUndeployedEvent represents the complete WebSub API undeployment event. +type WebSubAPIUndeployedEvent struct { + Type string `json:"type"` + Payload WebSubAPIUndeployedEventPayload `json:"payload"` + Timestamp string `json:"timestamp"` + CorrelationID string `json:"correlationId"` +} + +// WebSubAPIDeletedEventPayload represents the payload of a WebSub API deletion event. +type WebSubAPIDeletedEventPayload struct { + APIID string `json:"apiId"` +} + +// WebSubAPIDeletedEvent represents the complete WebSub API deletion event. +type WebSubAPIDeletedEvent struct { + Type string `json:"type"` + Payload WebSubAPIDeletedEventPayload `json:"payload"` + Timestamp string `json:"timestamp"` + CorrelationID string `json:"correlationId"` +} + +// WebBrokerAPIDeployedEventPayload represents the payload of a WebBroker API deployment event. +type WebBrokerAPIDeployedEventPayload struct { + APIID string `json:"apiId"` + DeploymentID string `json:"deploymentId"` + PerformedAt time.Time `json:"performedAt"` +} + +// WebBrokerAPIDeployedEvent represents the complete WebBroker API deployment event. +type WebBrokerAPIDeployedEvent struct { + Type string `json:"type"` + Payload WebBrokerAPIDeployedEventPayload `json:"payload"` + Timestamp string `json:"timestamp"` + CorrelationID string `json:"correlationId"` +} + +// WebBrokerAPIUndeployedEventPayload represents the payload of a WebBroker API undeployment event. +type WebBrokerAPIUndeployedEventPayload struct { + APIID string `json:"apiId"` + DeploymentID string `json:"deploymentId"` + PerformedAt time.Time `json:"performedAt"` +} + +// WebBrokerAPIUndeployedEvent represents the complete WebBroker API undeployment event. +type WebBrokerAPIUndeployedEvent struct { + Type string `json:"type"` + Payload WebBrokerAPIUndeployedEventPayload `json:"payload"` + Timestamp string `json:"timestamp"` + CorrelationID string `json:"correlationId"` +} + +// WebBrokerAPIDeletedEventPayload represents the payload of a WebBroker API deletion event. +type WebBrokerAPIDeletedEventPayload struct { + APIID string `json:"apiId"` +} + +// WebBrokerAPIDeletedEvent represents the complete WebBroker API deletion event. +type WebBrokerAPIDeletedEvent struct { + Type string `json:"type"` + Payload WebBrokerAPIDeletedEventPayload `json:"payload"` + Timestamp string `json:"timestamp"` + CorrelationID string `json:"correlationId"` +} + +// platformHmacSecretEventPayload is the payload for websub.hmacsecret.* events. +type platformHmacSecretEventPayload struct { + ArtifactUUID string `json:"artifactUuid"` + SecretName string `json:"secretName"` +} + +// platformHmacSecretInfo is the per-secret DTO returned by the internal HMAC endpoint. +type platformHmacSecretInfo struct { + Name string `json:"name"` + Secret string `json:"secret"` +} + +// platformHmacSecretsResponse is the response body from GET /websub-apis/:id/secrets. +type platformHmacSecretsResponse struct { + ArtifactID string `json:"artifactId"` + Secrets []platformHmacSecretInfo `json:"secrets"` +} + +// hmacSecretInfo is the internal view of a platform-managed HMAC secret. +type hmacSecretInfo struct { + Name string + Plaintext string +} diff --git a/event-gateway/gateway-controller/pkg/controlplanehooks/hooks.go b/event-gateway/gateway-controller/pkg/controlplanehooks/hooks.go new file mode 100644 index 000000000..0f7f77505 --- /dev/null +++ b/event-gateway/gateway-controller/pkg/controlplanehooks/hooks.go @@ -0,0 +1,641 @@ +/* + * Copyright (c) 2026, WSO2 LLC. (https://www.wso2.com). + * + * WSO2 LLC. licenses this file to you under the Apache License, + * Version 2.0 (the "License"); you may not use this file except + * in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package controlplanehooks + +import ( + "encoding/json" + "log/slog" + "time" + + "github.com/wso2/api-platform/common/eventhub" + "github.com/wso2/api-platform/gateway/gateway-controller/pkg/controlplane" + "github.com/wso2/api-platform/gateway/gateway-controller/pkg/models" + "github.com/wso2/api-platform/gateway/gateway-controller/pkg/storage" +) + +// Hooks implements controlplane.ControlPlaneEventGatewayHooks, supplying +// WebSub/WebBroker control-plane WebSocket event handling (deploy/undeploy/ +// delete and HMAC secret sync). Built entirely on top of *controlplane.Client's +// exported accessors — see gateway/gateway-controller/pkg/controlplane/eventgateway_hooks.go. +type Hooks struct{} + +// HandleWebSubAPIDeployed handles websub.deployed events. +func (Hooks) HandleWebSubAPIDeployed(c *controlplane.Client, event map[string]any) { + logger := c.Logger() + logger.Debug("WebSub API Deployment Event", + slog.Any("payload", event["payload"]), + slog.Any("timestamp", event["timestamp"]), + slog.Any("correlationId", event["correlationId"]), + ) + + eventBytes, err := json.Marshal(event) + if err != nil { + logger.Error("Failed to marshal WebSub API deployment event for parsing", slog.Any("error", err)) + return + } + + var deployedEvent WebSubAPIDeployedEvent + if err := json.Unmarshal(eventBytes, &deployedEvent); err != nil { + logger.Error("Failed to parse WebSub API deployment event", slog.Any("error", err)) + return + } + + apiID := deployedEvent.Payload.APIID + if apiID == "" { + logger.Error("API ID is empty in WebSub API deployment event") + return + } + + logger.Info("Processing WebSub API deployment", + slog.String("api_id", apiID), + slog.String("deployment_id", deployedEvent.Payload.DeploymentID), + slog.String("correlation_id", deployedEvent.CorrelationID), + ) + + // Fetch WebSub API definition from control plane + zipData, err := c.APIUtilsService().FetchResourceZip("/websub-apis/"+apiID, "WebSub API definition") + if err != nil { + logger.Error("Failed to fetch WebSub API definition", + slog.String("api_id", apiID), + slog.Any("error", err), + ) + c.SendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "websub", "deploy", "failed", + deployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") + return + } + + yamlData, err := c.APIUtilsService().ExtractYAMLFromZip(zipData) + if err != nil { + logger.Error("Failed to extract YAML from WebSub API ZIP", + slog.String("api_id", apiID), + slog.Any("error", err), + ) + c.SendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "websub", "deploy", "failed", + deployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") + return + } + + // Ensure any {{ secret "handle" }} references in the YAML are in local + // storage before rendering. + c.SyncSecretRefsFromYAML(yamlData, deployedEvent.CorrelationID) + + performedAt := deployedEvent.Payload.PerformedAt.Truncate(time.Millisecond) + if performedAt.IsZero() { + performedAt = time.Now().Truncate(time.Millisecond) + } + // Reuse the existing local UUID for a bottom-up (DP->CP) synced API so the + // control-plane deploy is an in-place update. + result, err := c.APIUtilsService().CreateAPIFromYAML(yamlData, c.ResolveLocalArtifactID(apiID), deployedEvent.Payload.DeploymentID, &performedAt, deployedEvent.CorrelationID, c.DeploymentService()) + if err != nil { + logger.Error("Failed to create WebSub API from YAML", + slog.String("api_id", apiID), + slog.Any("error", err), + ) + c.SendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "websub", "deploy", "failed", + deployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") + return + } + + if result.IsStale { + logger.Debug("Skipped stale WebSub API deploy event (newer version exists in DB)", + slog.String("api_id", apiID), + slog.String("deployment_id", deployedEvent.Payload.DeploymentID), + ) + return + } + + // Load platform-managed HMAC secrets into the webhook secret store. + if result.StoredConfig != nil { + syncHmacSecretsForArtifact(c, result.StoredConfig.UUID) + } + + c.SendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "websub", "deploy", "success", + deployedEvent.Payload.PerformedAt, "") + + logger.Info("Successfully processed WebSub API deployment event", + slog.String("api_id", apiID), + slog.String("correlation_id", deployedEvent.CorrelationID), + ) +} + +// HandleWebSubAPIUndeployed handles websub.undeployed events. +func (Hooks) HandleWebSubAPIUndeployed(c *controlplane.Client, event map[string]any) { + logger := c.Logger() + logger.Debug("WebSub API Undeployment Event", + slog.Any("payload", event["payload"]), + slog.Any("timestamp", event["timestamp"]), + slog.Any("correlationId", event["correlationId"]), + ) + + eventBytes, err := json.Marshal(event) + if err != nil { + logger.Error("Failed to marshal WebSub API undeployment event for parsing", slog.Any("error", err)) + return + } + + var undeployedEvent WebSubAPIUndeployedEvent + if err := json.Unmarshal(eventBytes, &undeployedEvent); err != nil { + logger.Error("Failed to parse WebSub API undeployment event", slog.Any("error", err)) + return + } + + apiID := undeployedEvent.Payload.APIID + if apiID == "" { + logger.Error("API ID is empty in WebSub API undeployment event") + return + } + + apiConfig, err := c.FindAPIConfig(apiID) + if err != nil { + if storage.IsNotFoundError(err) { + logger.Warn("WebSub API configuration not found for undeployment", + slog.String("api_id", apiID), + ) + c.SendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "websub", "undeploy", "success", + undeployedEvent.Payload.PerformedAt, "") + return + } + logger.Error("Failed to fetch WebSub API configuration for undeployment", + slog.String("api_id", apiID), + slog.String("correlation_id", undeployedEvent.CorrelationID), + slog.Any("error", err), + ) + c.SendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "websub", "undeploy", "failed", + undeployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") + return + } + + if apiConfig.DeploymentID != "" && undeployedEvent.Payload.DeploymentID != "" && + apiConfig.DeploymentID != undeployedEvent.Payload.DeploymentID { + logger.Warn("Ignoring stale WebSub API undeploy event: deployment ID mismatch", + slog.String("api_id", apiID), + slog.String("event_deployment_id", undeployedEvent.Payload.DeploymentID), + slog.String("current_deployment_id", apiConfig.DeploymentID), + ) + c.SendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "websub", "undeploy", "failed", + undeployedEvent.Payload.PerformedAt, "DEPLOYMENT_ID_MISMATCH") + return + } + + performedAt := undeployedEvent.Payload.PerformedAt.Truncate(time.Millisecond) + if performedAt.IsZero() { + performedAt = time.Now().Truncate(time.Millisecond) + } + apiConfig.DesiredState = models.StateUndeployed + apiConfig.DeploymentID = undeployedEvent.Payload.DeploymentID + apiConfig.DeployedAt = &performedAt + apiConfig.UpdatedAt = time.Now() + + affected, err := c.DB().UpsertConfig(apiConfig) + if err != nil { + logger.Error("Failed to upsert config for WebSub API undeployment", + slog.String("api_id", apiID), + slog.Any("error", err), + ) + c.SendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "websub", "undeploy", "failed", + undeployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") + return + } + if !affected { + logger.Debug("Skipped stale WebSub API undeploy event (newer version exists in DB)", + slog.String("api_id", apiID), + slog.String("deployment_id", undeployedEvent.Payload.DeploymentID), + ) + return + } + + evt := eventhub.Event{ + EventType: eventhub.EventTypeAPI, + Action: "UPDATE", + EntityID: apiID, + EventID: undeployedEvent.CorrelationID, + } + if err := c.EventHub().PublishEvent(c.GatewayID(), evt); err != nil { + logger.Error("Failed to publish WebSub API undeployment event", slog.Any("error", err)) + } + + c.SendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "websub", "undeploy", "success", + undeployedEvent.Payload.PerformedAt, "") + + logger.Info("Successfully processed WebSub API undeployment event", + slog.String("api_id", apiID), + slog.String("correlation_id", undeployedEvent.CorrelationID), + ) +} + +// HandleWebSubAPIDeleted handles websub.deleted events. +func (Hooks) HandleWebSubAPIDeleted(c *controlplane.Client, event map[string]any) { + logger := c.Logger() + logger.Debug("WebSub API Deleted Event", + slog.Any("payload", event["payload"]), + slog.Any("timestamp", event["timestamp"]), + slog.Any("correlationId", event["correlationId"]), + ) + + eventBytes, err := json.Marshal(event) + if err != nil { + logger.Error("Failed to marshal WebSub API deleted event for parsing", slog.Any("error", err)) + return + } + + var deletedEvent WebSubAPIDeletedEvent + if err := json.Unmarshal(eventBytes, &deletedEvent); err != nil { + logger.Error("Failed to parse WebSub API deleted event", slog.Any("error", err)) + return + } + + apiID := deletedEvent.Payload.APIID + if apiID == "" { + logger.Error("API ID is empty in WebSub API deleted event") + return + } + + apiConfig, err := c.FindAPIConfig(apiID) + if err != nil { + if storage.IsNotFoundError(err) { + logger.Warn("WebSub API configuration not found for deletion", + slog.String("api_id", apiID), + ) + cleanupHmacSecretsForArtifact(c, apiID) + return + } + logger.Error("Failed to fetch WebSub API configuration for deletion", + slog.String("api_id", apiID), + slog.String("correlation_id", deletedEvent.CorrelationID), + slog.Any("error", err), + ) + return + } + + c.PerformFullAPIDeletion(apiID, apiConfig, deletedEvent.CorrelationID) + cleanupHmacSecretsForArtifact(c, apiConfig.UUID) +} + +// HandleWebBrokerAPIDeployed handles webbroker.deployed events. +func (Hooks) HandleWebBrokerAPIDeployed(c *controlplane.Client, event map[string]any) { + logger := c.Logger() + logger.Debug("WebBroker API Deployment Event", + slog.Any("payload", event["payload"]), + slog.Any("timestamp", event["timestamp"]), + slog.Any("correlationId", event["correlationId"]), + ) + + eventBytes, err := json.Marshal(event) + if err != nil { + logger.Error("Failed to marshal WebBroker API deployment event for parsing", slog.Any("error", err)) + return + } + + var deployedEvent WebBrokerAPIDeployedEvent + if err := json.Unmarshal(eventBytes, &deployedEvent); err != nil { + logger.Error("Failed to parse WebBroker API deployment event", slog.Any("error", err)) + return + } + + apiID := deployedEvent.Payload.APIID + if apiID == "" { + logger.Error("API ID is empty in WebBroker API deployment event") + return + } + + logger.Info("Processing WebBroker API deployment", + slog.String("api_id", apiID), + slog.String("deployment_id", deployedEvent.Payload.DeploymentID), + slog.String("correlation_id", deployedEvent.CorrelationID), + ) + + // Fetch WebBroker API definition from control plane + zipData, err := c.APIUtilsService().FetchResourceZip("/webbroker-apis/"+apiID, "WebBroker API definition") + if err != nil { + logger.Error("Failed to fetch WebBroker API definition", + slog.String("api_id", apiID), + slog.Any("error", err), + ) + c.SendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "webbroker", "deploy", "failed", + deployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") + return + } + + yamlData, err := c.APIUtilsService().ExtractYAMLFromZip(zipData) + if err != nil { + logger.Error("Failed to extract YAML from WebBroker API ZIP", + slog.String("api_id", apiID), + slog.Any("error", err), + ) + c.SendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "webbroker", "deploy", "failed", + deployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") + return + } + + // Ensure any {{ secret "handle" }} references in the YAML are in local + // storage before rendering. + c.SyncSecretRefsFromYAML(yamlData, deployedEvent.CorrelationID) + + performedAt := deployedEvent.Payload.PerformedAt.Truncate(time.Millisecond) + if performedAt.IsZero() { + performedAt = time.Now().Truncate(time.Millisecond) + } + // Reuse the existing local UUID for a bottom-up (DP->CP) synced API so the + // control-plane deploy is an in-place update. + result, err := c.APIUtilsService().CreateAPIFromYAML(yamlData, c.ResolveLocalArtifactID(apiID), deployedEvent.Payload.DeploymentID, &performedAt, deployedEvent.CorrelationID, c.DeploymentService()) + if err != nil { + logger.Error("Failed to create WebBroker API from YAML", + slog.String("api_id", apiID), + slog.Any("error", err), + ) + c.SendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "webbroker", "deploy", "failed", + deployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") + return + } + + if result.IsStale { + logger.Debug("Skipped stale WebBroker API deploy event (newer version exists in DB)", + slog.String("api_id", apiID), + slog.String("deployment_id", deployedEvent.Payload.DeploymentID), + ) + return + } + + c.SendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "webbroker", "deploy", "success", + deployedEvent.Payload.PerformedAt, "") + + logger.Info("Successfully processed WebBroker API deployment event", + slog.String("api_id", apiID), + slog.String("correlation_id", deployedEvent.CorrelationID), + ) +} + +// HandleWebBrokerAPIUndeployed handles webbroker.undeployed events. +func (Hooks) HandleWebBrokerAPIUndeployed(c *controlplane.Client, event map[string]any) { + logger := c.Logger() + logger.Debug("WebBroker API Undeployment Event", + slog.Any("payload", event["payload"]), + slog.Any("timestamp", event["timestamp"]), + slog.Any("correlationId", event["correlationId"]), + ) + + eventBytes, err := json.Marshal(event) + if err != nil { + logger.Error("Failed to marshal WebBroker API undeployment event for parsing", slog.Any("error", err)) + return + } + + var undeployedEvent WebBrokerAPIUndeployedEvent + if err := json.Unmarshal(eventBytes, &undeployedEvent); err != nil { + logger.Error("Failed to parse WebBroker API undeployment event", slog.Any("error", err)) + return + } + + apiID := undeployedEvent.Payload.APIID + if apiID == "" { + logger.Error("API ID is empty in WebBroker API undeployment event") + return + } + + apiConfig, err := c.FindAPIConfig(apiID) + if err != nil { + if storage.IsNotFoundError(err) { + logger.Warn("WebBroker API configuration not found for undeployment", + slog.String("api_id", apiID), + ) + c.SendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "webbroker", "undeploy", "success", + undeployedEvent.Payload.PerformedAt, "") + return + } + logger.Error("Failed to fetch WebBroker API configuration for undeployment", + slog.String("api_id", apiID), + slog.String("correlation_id", undeployedEvent.CorrelationID), + slog.Any("error", err), + ) + c.SendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "webbroker", "undeploy", "failed", + undeployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") + return + } + + if apiConfig.DeploymentID != "" && undeployedEvent.Payload.DeploymentID != "" && + apiConfig.DeploymentID != undeployedEvent.Payload.DeploymentID { + logger.Warn("Ignoring stale WebBroker API undeploy event: deployment ID mismatch", + slog.String("api_id", apiID), + slog.String("event_deployment_id", undeployedEvent.Payload.DeploymentID), + slog.String("current_deployment_id", apiConfig.DeploymentID), + ) + c.SendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "webbroker", "undeploy", "failed", + undeployedEvent.Payload.PerformedAt, "DEPLOYMENT_ID_MISMATCH") + return + } + + performedAt := undeployedEvent.Payload.PerformedAt.Truncate(time.Millisecond) + if performedAt.IsZero() { + performedAt = time.Now().Truncate(time.Millisecond) + } + apiConfig.DesiredState = models.StateUndeployed + apiConfig.DeploymentID = undeployedEvent.Payload.DeploymentID + apiConfig.DeployedAt = &performedAt + apiConfig.UpdatedAt = time.Now() + + affected, err := c.DB().UpsertConfig(apiConfig) + if err != nil { + logger.Error("Failed to upsert config for WebBroker API undeployment", + slog.String("api_id", apiID), + slog.Any("error", err), + ) + c.SendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "webbroker", "undeploy", "failed", + undeployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") + return + } + if !affected { + logger.Debug("Skipped stale WebBroker API undeploy event (newer version exists in DB)", + slog.String("api_id", apiID), + slog.String("deployment_id", undeployedEvent.Payload.DeploymentID), + ) + return + } + + evt := eventhub.Event{ + EventType: eventhub.EventTypeAPI, + Action: "UPDATE", + EntityID: apiID, + EventID: undeployedEvent.CorrelationID, + } + if err := c.EventHub().PublishEvent(c.GatewayID(), evt); err != nil { + logger.Error("Failed to publish WebBroker API undeployment event", slog.Any("error", err)) + } + + c.SendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "webbroker", "undeploy", "success", + undeployedEvent.Payload.PerformedAt, "") + + logger.Info("Successfully processed WebBroker API undeployment event", + slog.String("api_id", apiID), + slog.String("correlation_id", undeployedEvent.CorrelationID), + ) +} + +// HandleWebBrokerAPIDeleted handles webbroker.deleted events. +func (Hooks) HandleWebBrokerAPIDeleted(c *controlplane.Client, event map[string]any) { + logger := c.Logger() + logger.Debug("WebBroker API Deleted Event", + slog.Any("payload", event["payload"]), + slog.Any("timestamp", event["timestamp"]), + slog.Any("correlationId", event["correlationId"]), + ) + + eventBytes, err := json.Marshal(event) + if err != nil { + logger.Error("Failed to marshal WebBroker API deleted event for parsing", slog.Any("error", err)) + return + } + + var deletedEvent WebBrokerAPIDeletedEvent + if err := json.Unmarshal(eventBytes, &deletedEvent); err != nil { + logger.Error("Failed to parse WebBroker API deleted event", slog.Any("error", err)) + return + } + + apiID := deletedEvent.Payload.APIID + if apiID == "" { + logger.Error("API ID is empty in WebBroker API deleted event") + return + } + + apiConfig, err := c.FindAPIConfig(apiID) + if err != nil { + if storage.IsNotFoundError(err) { + logger.Warn("WebBroker API configuration not found for deletion; running orphan cleanup", + slog.String("api_id", apiID), + ) + c.CleanupOrphanedResources(apiID, deletedEvent.CorrelationID) + return + } + logger.Error("Failed to fetch WebBroker API configuration for deletion", + slog.String("api_id", apiID), + slog.String("correlation_id", deletedEvent.CorrelationID), + slog.Any("error", err), + ) + return + } + + c.PerformFullAPIDeletion(apiID, apiConfig, deletedEvent.CorrelationID) +} + +// HandleWebSubAPIHmacSecretEvent handles websub.hmacsecret.created/updated/deleted +// events from platform-API. It re-syncs all platform-managed HMAC secrets for +// the affected artifact. +func (Hooks) HandleWebSubAPIHmacSecretEvent(c *controlplane.Client, event map[string]any) { + logger := c.Logger() + payloadBytes, err := json.Marshal(event["payload"]) + if err != nil { + logger.Error("Failed to marshal HMAC secret event payload", slog.Any("error", err)) + return + } + var payload platformHmacSecretEventPayload + if err := json.Unmarshal(payloadBytes, &payload); err != nil { + logger.Error("Failed to parse HMAC secret event payload", slog.Any("error", err)) + return + } + if payload.ArtifactUUID == "" { + logger.Warn("HMAC secret event missing artifactUuid, skipping") + return + } + logger.Info("Processing platform HMAC secret event", + slog.Any("type", event["type"]), + slog.String("artifact_uuid", payload.ArtifactUUID), + slog.String("secret_name", payload.SecretName), + ) + syncHmacSecretsForArtifact(c, payload.ArtifactUUID) +} + +// syncHmacSecretsForArtifact fetches all platform-managed HMAC secrets for a WebSub API +// artifact from platform-API and loads them into the in-memory webhook secret store. +// It replaces any previously loaded secrets for this artifact atomically (clear then re-add). +func syncHmacSecretsForArtifact(c *controlplane.Client, artifactID string) { + store := c.WebhookSecretStore() + if store == nil { + return + } + logger := c.Logger() + + secrets, err := fetchWebSubAPIHmacSecrets(c, artifactID) + if err != nil { + logger.Warn("Failed to fetch platform HMAC secrets for WebSub API", + slog.String("artifact_id", artifactID), + slog.Any("error", err)) + return + } + + if err := store.RemoveAllByAPI(artifactID); err != nil { + logger.Warn("Failed to clear existing HMAC secrets for WebSub API", + slog.String("artifact_id", artifactID), + slog.Any("error", err)) + return + } + + for _, s := range secrets { + if err := store.Store(artifactID, s.Name, s.Plaintext); err != nil { + logger.Warn("Failed to store platform HMAC secret in memory", + slog.String("artifact_id", artifactID), + slog.String("secret_name", s.Name), + slog.Any("error", err)) + } + } + + if err := c.RefreshWebhookSecretSnapshot(); err != nil { + logger.Warn("Failed to refresh webhook secret xDS snapshot after platform sync", + slog.String("artifact_id", artifactID), + slog.Any("error", err)) + } + + logger.Info("Loaded platform HMAC secrets for WebSub API", + slog.String("artifact_id", artifactID), + slog.Int("count", len(secrets))) +} + +// cleanupHmacSecretsForArtifact removes all in-memory HMAC secrets for an artifact and +// refreshes the xDS snapshot. Called on WebSub API deletion (found and not-found paths). +func cleanupHmacSecretsForArtifact(c *controlplane.Client, artifactID string) { + store := c.WebhookSecretStore() + if store == nil { + return + } + logger := c.Logger() + if err := store.RemoveAllByAPI(artifactID); err != nil { + logger.Warn("Failed to remove HMAC secrets from store during WebSub API cleanup", + slog.String("artifact_id", artifactID), + slog.Any("error", err)) + } + if err := c.RefreshWebhookSecretSnapshot(); err != nil { + logger.Warn("Failed to refresh webhook secret xDS snapshot after WebSub API cleanup", + slog.String("artifact_id", artifactID), + slog.Any("error", err)) + } +} + +// fetchWebSubAPIHmacSecrets fetches the plaintext HMAC secrets for a WebSub API artifact +// from the platform-API internal endpoint. +func fetchWebSubAPIHmacSecrets(c *controlplane.Client, artifactID string) ([]hmacSecretInfo, error) { + var response platformHmacSecretsResponse + if err := c.APIUtilsService().FetchResourceJSON("/websub-apis/"+artifactID+"/secrets", "WebSub API HMAC secrets", &response); err != nil { + return nil, err + } + + secrets := make([]hmacSecretInfo, 0, len(response.Secrets)) + for _, s := range response.Secrets { + secrets = append(secrets, hmacSecretInfo{Name: s.Name, Plaintext: s.Secret}) + } + return secrets, nil +} diff --git a/event-gateway/gateway-controller/pkg/dbschema/dbschema.go b/event-gateway/gateway-controller/pkg/dbschema/dbschema.go new file mode 100644 index 000000000..b02f6a8e8 --- /dev/null +++ b/event-gateway/gateway-controller/pkg/dbschema/dbschema.go @@ -0,0 +1,76 @@ +/* + * Copyright (c) 2026, WSO2 LLC. (https://www.wso2.com). + * + * WSO2 LLC. licenses this file to you under the Apache License, + * Version 2.0 (the "License"); you may not use this file except + * in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +// Package dbschema owns the database schema for tables that are specific to +// event-gateway-controller: websub_apis, webbroker_apis, and webhook_secrets. +// These used to live in gateway-controller (core)'s own schema scripts, but +// core has no use for them on its own — only this module's kinds and features +// reference them. Apply runs this module's own idempotent DDL against the +// exact same database connection/backend that core's Storage already opened +// (via Storage.GetDB()), immediately after core's own schema has been +// applied, so both sets of tables end up in the same database without core's +// schema files ever needing to know about these tables. +// +// The embedded .sql files here mirror gateway/gateway-controller/pkg/storage's +// own gateway-controller-db*.sql naming and per-dialect split. +package dbschema + +import ( + "context" + "database/sql" + _ "embed" + "fmt" +) + +//go:embed eventgateway-db.sql +var sqliteSchema string + +//go:embed eventgateway-db.postgres.sql +var postgresSchema string + +//go:embed eventgateway-db.sqlserver.sql +var sqlserverSchema string + +// forBackend returns the idempotent DDL for the given storage backend type +// ("sqlite", "postgres", or "sqlserver" — matching storage.BackendConfig.Type). +func forBackend(backendType string) (string, error) { + switch backendType { + case "sqlite": + return sqliteSchema, nil + case "postgres": + return postgresSchema, nil + case "sqlserver": + return sqlserverSchema, nil + default: + return "", fmt.Errorf("unsupported storage backend for event-gateway schema: %s", backendType) + } +} + +// Apply creates the websub_apis, webbroker_apis, and webhook_secrets tables +// (if they don't already exist) against db, using the DDL dialect for +// backendType. Safe to call on every startup — every statement is idempotent. +func Apply(ctx context.Context, db *sql.DB, backendType string) error { + schema, err := forBackend(backendType) + if err != nil { + return err + } + if _, err := db.ExecContext(ctx, schema); err != nil { + return fmt.Errorf("failed to apply event-gateway database schema: %w", err) + } + return nil +} diff --git a/event-gateway/gateway-controller/pkg/dbschema/eventgateway-db.postgres.sql b/event-gateway/gateway-controller/pkg/dbschema/eventgateway-db.postgres.sql new file mode 100644 index 000000000..db73682a8 --- /dev/null +++ b/event-gateway/gateway-controller/pkg/dbschema/eventgateway-db.postgres.sql @@ -0,0 +1,39 @@ +-- PostgreSQL Schema for Event-Gateway-Controller-specific tables +-- Applied against the same database gateway-controller (core) opened, after +-- core's own schema (gateway-controller-db.postgres.sql) has been applied. +-- Core's schema scripts do not define these tables — only this module does. + +CREATE TABLE IF NOT EXISTS websub_apis ( + uuid TEXT NOT NULL, + gateway_id TEXT NOT NULL, + configuration TEXT NOT NULL, + PRIMARY KEY (gateway_id, uuid), + FOREIGN KEY(gateway_id, uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE +); + +CREATE TABLE IF NOT EXISTS webbroker_apis ( + uuid TEXT NOT NULL, + gateway_id TEXT NOT NULL, + configuration TEXT NOT NULL, + PRIMARY KEY (gateway_id, uuid), + FOREIGN KEY(gateway_id, uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE +); + +-- Per-API HMAC secrets for the websub-hmac-auth policy. +-- Ciphertext is AES-256-GCM encrypted; plaintext is never stored. +CREATE TABLE IF NOT EXISTS webhook_secrets ( + uuid TEXT NOT NULL, + gateway_id TEXT NOT NULL, + artifact_uuid TEXT NOT NULL, + name TEXT NOT NULL, + display_name TEXT NOT NULL, + ciphertext BYTEA NOT NULL, + status TEXT NOT NULL DEFAULT 'active' CHECK(status IN ('active', 'revoked')), + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (gateway_id, uuid), + UNIQUE (gateway_id, artifact_uuid, name), + FOREIGN KEY (gateway_id, artifact_uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE +); + +CREATE INDEX IF NOT EXISTS idx_webhook_secrets_artifact ON webhook_secrets(gateway_id, artifact_uuid); diff --git a/event-gateway/gateway-controller/pkg/dbschema/eventgateway-db.sql b/event-gateway/gateway-controller/pkg/dbschema/eventgateway-db.sql new file mode 100644 index 000000000..14557d9d6 --- /dev/null +++ b/event-gateway/gateway-controller/pkg/dbschema/eventgateway-db.sql @@ -0,0 +1,39 @@ +-- SQLite Schema for Event-Gateway-Controller-specific tables +-- Applied against the same database gateway-controller (core) opened, after +-- core's own schema (gateway-controller-db.sql) has been applied. Core's +-- schema scripts do not define these tables — only this module does. + +CREATE TABLE IF NOT EXISTS websub_apis ( + uuid TEXT NOT NULL, + gateway_id TEXT NOT NULL, + configuration TEXT NOT NULL, + PRIMARY KEY (gateway_id, uuid), + FOREIGN KEY(gateway_id, uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE +); + +CREATE TABLE IF NOT EXISTS webbroker_apis ( + uuid TEXT NOT NULL, + gateway_id TEXT NOT NULL, + configuration TEXT NOT NULL, + PRIMARY KEY (gateway_id, uuid), + FOREIGN KEY(gateway_id, uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE +); + +-- Per-API HMAC secrets for the websub-hmac-auth policy. +-- Ciphertext is AES-256-GCM encrypted; plaintext is never stored. +CREATE TABLE IF NOT EXISTS webhook_secrets ( + uuid TEXT NOT NULL, + gateway_id TEXT NOT NULL, + artifact_uuid TEXT NOT NULL, + name TEXT NOT NULL, + display_name TEXT NOT NULL, + ciphertext BLOB NOT NULL, + status TEXT NOT NULL DEFAULT 'active' CHECK(status IN ('active', 'revoked')), + created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (gateway_id, uuid), + UNIQUE (gateway_id, artifact_uuid, name), + FOREIGN KEY (gateway_id, artifact_uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE +); + +CREATE INDEX IF NOT EXISTS idx_webhook_secrets_artifact ON webhook_secrets(gateway_id, artifact_uuid); diff --git a/event-gateway/gateway-controller/pkg/dbschema/eventgateway-db.sqlserver.sql b/event-gateway/gateway-controller/pkg/dbschema/eventgateway-db.sqlserver.sql new file mode 100644 index 000000000..aa63a652f --- /dev/null +++ b/event-gateway/gateway-controller/pkg/dbschema/eventgateway-db.sqlserver.sql @@ -0,0 +1,43 @@ +-- SQL Server Schema for Event-Gateway-Controller-specific tables +-- Applied against the same database gateway-controller (core) opened, after +-- core's own schema (gateway-controller-db.sqlserver.sql) has been applied. +-- Core's schema scripts do not define these tables — only this module does. +-- Every object is guarded by IF OBJECT_ID/IF NOT EXISTS so the batch is +-- idempotent, matching core's own SQL Server schema conventions. + +IF OBJECT_ID(N'dbo.websub_apis', N'U') IS NULL +CREATE TABLE dbo.websub_apis ( + uuid NVARCHAR(64) NOT NULL, + gateway_id NVARCHAR(64) NOT NULL, + configuration NVARCHAR(MAX) NOT NULL, + PRIMARY KEY (gateway_id, uuid), + FOREIGN KEY(gateway_id, uuid) REFERENCES dbo.artifacts(gateway_id, uuid) ON DELETE CASCADE +); + +IF OBJECT_ID(N'dbo.webbroker_apis', N'U') IS NULL +CREATE TABLE dbo.webbroker_apis ( + uuid NVARCHAR(64) NOT NULL, + gateway_id NVARCHAR(64) NOT NULL, + configuration NVARCHAR(MAX) NOT NULL, + PRIMARY KEY (gateway_id, uuid), + FOREIGN KEY(gateway_id, uuid) REFERENCES dbo.artifacts(gateway_id, uuid) ON DELETE CASCADE +); + +IF OBJECT_ID(N'dbo.webhook_secrets', N'U') IS NULL +CREATE TABLE dbo.webhook_secrets ( + uuid NVARCHAR(64) NOT NULL, + gateway_id NVARCHAR(64) NOT NULL, + artifact_uuid NVARCHAR(64) NOT NULL, + name NVARCHAR(255) NOT NULL, + display_name NVARCHAR(255) NOT NULL, + ciphertext VARBINARY(MAX) NOT NULL, + status NVARCHAR(20) NOT NULL CHECK(status IN ('active', 'revoked')) DEFAULT 'active', + created_at DATETIME2(7) NOT NULL DEFAULT SYSUTCDATETIME(), + updated_at DATETIME2(7) NOT NULL DEFAULT SYSUTCDATETIME(), + PRIMARY KEY (gateway_id, uuid), + CONSTRAINT uq_webhook_secrets_artifact_name UNIQUE (gateway_id, artifact_uuid, name), + FOREIGN KEY (gateway_id, artifact_uuid) REFERENCES dbo.artifacts(gateway_id, uuid) ON DELETE CASCADE +); + +IF NOT EXISTS (SELECT 1 FROM sys.indexes WHERE name = N'idx_webhook_secrets_artifact' AND object_id = OBJECT_ID(N'dbo.webhook_secrets')) +CREATE INDEX idx_webhook_secrets_artifact ON dbo.webhook_secrets(gateway_id, artifact_uuid); diff --git a/event-gateway/gateway-controller/pkg/kindsupport/kindsupport.go b/event-gateway/gateway-controller/pkg/kindsupport/kindsupport.go index e56d8e034..9fc0a86d4 100644 --- a/event-gateway/gateway-controller/pkg/kindsupport/kindsupport.go +++ b/event-gateway/gateway-controller/pkg/kindsupport/kindsupport.go @@ -45,6 +45,9 @@ import ( // Register wires WebSubApi/WebBrokerApi kind support into every core registry // that gateway-controller exposes for kinds it doesn't know about natively. func Register() { + storage.RegisterKindResourceTable("WebSubApi", "websub_apis") + storage.RegisterKindResourceTable("WebBrokerApi", "webbroker_apis") + storage.RegisterKindUnmarshaler("WebSubApi", unmarshalWebSubAPI) storage.RegisterKindUnmarshaler("WebBrokerApi", unmarshalWebBrokerApi) diff --git a/gateway/gateway-controller/cmd/controller/main.go b/gateway/gateway-controller/cmd/controller/main.go index a5b2a8bed..5c3fd672f 100644 --- a/gateway/gateway-controller/cmd/controller/main.go +++ b/gateway/gateway-controller/cmd/controller/main.go @@ -817,17 +817,6 @@ func generateAuthConfig(config *config.Config) commonmodels.AuthConfig { "PUT /rest-apis/{id}": {"admin", "developer"}, "DELETE /rest-apis/{id}": {"admin", "developer"}, - "POST /websub-apis": {"admin", "developer"}, - "GET /websub-apis": {"admin", "developer"}, - "GET /websub-apis/{id}": {"admin", "developer"}, - "PUT /websub-apis/{id}": {"admin", "developer"}, - "DELETE /websub-apis/{id}": {"admin", "developer"}, - - "POST /webbroker-apis": {"admin", "developer"}, - "GET /webbroker-apis": {"admin", "developer"}, - "GET /webbroker-apis/{id}": {"admin", "developer"}, - "DELETE /webbroker-apis/{id}": {"admin", "developer"}, - "GET /certificates": {"admin", "developer"}, "POST /certificates": {"admin", "developer"}, "DELETE /certificates/{id}": {"admin"}, @@ -877,23 +866,6 @@ func generateAuthConfig(config *config.Config) commonmodels.AuthConfig { "POST /llm-proxies/{id}/api-keys/{apiKeyName}/regenerate": {"admin", "consumer"}, "DELETE /llm-proxies/{id}/api-keys/{apiKeyName}": {"admin", "consumer"}, - "POST /websub-apis/{id}/api-keys": {"admin", "consumer"}, - "GET /websub-apis/{id}/api-keys": {"admin", "consumer"}, - "PUT /websub-apis/{id}/api-keys/{apiKeyName}": {"admin", "consumer"}, - "POST /websub-apis/{id}/api-keys/{apiKeyName}/regenerate": {"admin", "consumer"}, - "DELETE /websub-apis/{id}/api-keys/{apiKeyName}": {"admin", "consumer"}, - - "POST /websub-apis/{id}/secrets": {"admin", "consumer"}, - "GET /websub-apis/{id}/secrets": {"admin", "consumer"}, - "DELETE /websub-apis/{id}/secrets/{secretName}": {"admin", "consumer"}, - "POST /websub-apis/{id}/secrets/{secretName}/regenerate": {"admin", "consumer"}, - - "POST /webbroker-apis/{id}/api-keys": {"admin", "consumer"}, - "GET /webbroker-apis/{id}/api-keys": {"admin", "consumer"}, - "PUT /webbroker-apis/{id}/api-keys/{apiKeyName}": {"admin", "consumer"}, - "POST /webbroker-apis/{id}/api-keys/{apiKeyName}/regenerate": {"admin", "consumer"}, - "DELETE /webbroker-apis/{id}/api-keys/{apiKeyName}": {"admin", "consumer"}, - // Root-level subscription endpoints "POST /subscriptions": {"admin", "developer"}, "GET /subscriptions": {"admin", "developer"}, diff --git a/gateway/gateway-controller/pkg/api/handlers/handlers.go b/gateway/gateway-controller/pkg/api/handlers/handlers.go index cd139e3e2..bd995ddcc 100644 --- a/gateway/gateway-controller/pkg/api/handlers/handlers.go +++ b/gateway/gateway-controller/pkg/api/handlers/handlers.go @@ -363,11 +363,6 @@ func (s *APIServer) waitForDeploymentAndPush(configID string, correlationID stri s.deploymentPusher().WaitForDeploymentAndPush(configID, correlationID, minDeployedAt, log) } -// publishWebSubEvent publishes an event for WebSub API lifecycle changes. -func (s *APIServer) publishWebSubEvent(eventType eventhub.EventType, action, entityID, correlationID string, logger *slog.Logger) { - (&handlerkit.EventPublisher{EventHub: s.eventHub, GatewayID: s.gatewayID}).PublishEvent(eventType, action, entityID, correlationID, logger) -} - // GetConfigDump implements the GET /config_dump endpoint func (s *APIServer) GetConfigDump(w http.ResponseWriter, r *http.Request) { log := middleware.GetLogger(r, s.logger) diff --git a/gateway/gateway-controller/pkg/controlplane/client.go b/gateway/gateway-controller/pkg/controlplane/client.go index 83bd1cbe8..d2bb60672 100644 --- a/gateway/gateway-controller/pkg/controlplane/client.go +++ b/gateway/gateway-controller/pkg/controlplane/client.go @@ -154,6 +154,7 @@ type Client struct { webhookSecretSnapshotManager WebhookSecretSnapshotRefresher secretSyncer secretSyncer secretHashCache sync.Map // handle → last-known Platform API hash (string) + eventGatewayHooks ControlPlaneEventGatewayHooks // DP->CP push retry tuning. pushMaxAttempts int @@ -1352,19 +1353,19 @@ func (c *Client) handleMessage(messageType int, message []byte) { case "mcpproxy.deleted": c.handleMCPProxyDeletedEvent(event) case "websub.deployed": - c.handleWebSubAPIDeployedEvent(event) + c.dispatchEventGatewayHook(event["type"], func(h ControlPlaneEventGatewayHooks) { h.HandleWebSubAPIDeployed(c, event) }) case "websub.undeployed": - c.handleWebSubAPIUndeployedEvent(event) + c.dispatchEventGatewayHook(event["type"], func(h ControlPlaneEventGatewayHooks) { h.HandleWebSubAPIUndeployed(c, event) }) case "websub.deleted": - c.handleWebSubAPIDeletedEvent(event) + c.dispatchEventGatewayHook(event["type"], func(h ControlPlaneEventGatewayHooks) { h.HandleWebSubAPIDeleted(c, event) }) case "websub.hmacsecret.created", "websub.hmacsecret.updated", "websub.hmacsecret.deleted": - c.handleWebSubAPIHmacSecretEvent(event) + c.dispatchEventGatewayHook(event["type"], func(h ControlPlaneEventGatewayHooks) { h.HandleWebSubAPIHmacSecretEvent(c, event) }) case "webbroker.deployed": - c.handleWebBrokerAPIDeployedEvent(event) + c.dispatchEventGatewayHook(event["type"], func(h ControlPlaneEventGatewayHooks) { h.HandleWebBrokerAPIDeployed(c, event) }) case "webbroker.undeployed": - c.handleWebBrokerAPIUndeployedEvent(event) + c.dispatchEventGatewayHook(event["type"], func(h ControlPlaneEventGatewayHooks) { h.HandleWebBrokerAPIUndeployed(c, event) }) case "webbroker.deleted": - c.handleWebBrokerAPIDeletedEvent(event) + c.dispatchEventGatewayHook(event["type"], func(h ControlPlaneEventGatewayHooks) { h.HandleWebBrokerAPIDeleted(c, event) }) case "application.updated": c.handleApplicationUpdatedEvent(event) default: @@ -2584,587 +2585,6 @@ func (c *Client) handleLLMProxyDeletedEvent(event map[string]interface{}) { ) } -// syncHmacSecretsForArtifact fetches all platform-managed HMAC secrets for a WebSub API -// artifact from platform-API and loads them into the in-memory webhook secret store. -// It replaces any previously loaded secrets for this artifact atomically (clear then re-add). -func (c *Client) syncHmacSecretsForArtifact(artifactID string) { - if c.webhookSecretStore == nil { - return - } - - secrets, err := c.apiUtilsService.FetchWebSubAPIHmacSecrets(artifactID) - if err != nil { - c.logger.Warn("Failed to fetch platform HMAC secrets for WebSub API", - slog.String("artifact_id", artifactID), - slog.Any("error", err)) - return - } - - if err := c.webhookSecretStore.RemoveAllByAPI(artifactID); err != nil { - c.logger.Warn("Failed to clear existing HMAC secrets for WebSub API", - slog.String("artifact_id", artifactID), - slog.Any("error", err)) - return - } - - for _, s := range secrets { - if err := c.webhookSecretStore.Store(artifactID, s.Name, s.Plaintext); err != nil { - c.logger.Warn("Failed to store platform HMAC secret in memory", - slog.String("artifact_id", artifactID), - slog.String("secret_name", s.Name), - slog.Any("error", err)) - } - } - - if c.webhookSecretSnapshotManager != nil { - if err := c.webhookSecretSnapshotManager.RefreshSnapshot(); err != nil { - c.logger.Warn("Failed to refresh webhook secret xDS snapshot after platform sync", - slog.String("artifact_id", artifactID), - slog.Any("error", err)) - } - } - - c.logger.Info("Loaded platform HMAC secrets for WebSub API", - slog.String("artifact_id", artifactID), - slog.Int("count", len(secrets))) -} - -// cleanupHmacSecretsForArtifact removes all in-memory HMAC secrets for an artifact and -// refreshes the xDS snapshot. Called on WebSub API deletion (found and not-found paths). -func (c *Client) cleanupHmacSecretsForArtifact(artifactID string) { - if c.webhookSecretStore == nil { - return - } - if err := c.webhookSecretStore.RemoveAllByAPI(artifactID); err != nil { - c.logger.Warn("Failed to remove HMAC secrets from store during WebSub API cleanup", - slog.String("artifact_id", artifactID), - slog.Any("error", err)) - } - if c.webhookSecretSnapshotManager != nil { - if err := c.webhookSecretSnapshotManager.RefreshSnapshot(); err != nil { - c.logger.Warn("Failed to refresh webhook secret xDS snapshot after WebSub API cleanup", - slog.String("artifact_id", artifactID), - slog.Any("error", err)) - } - } -} - -// platformHmacSecretEventPayload is the payload for websub.hmacsecret.* events. -type platformHmacSecretEventPayload struct { - ArtifactUUID string `json:"artifactUuid"` - SecretName string `json:"secretName"` -} - -func (c *Client) handleWebSubAPIDeployedEvent(event map[string]any) { - c.logger.Debug("WebSub API Deployment Event", - slog.Any("payload", event["payload"]), - slog.Any("timestamp", event["timestamp"]), - slog.Any("correlationId", event["correlationId"]), - ) - - eventBytes, err := json.Marshal(event) - if err != nil { - c.logger.Error("Failed to marshal WebSub API deployment event for parsing", - slog.Any("error", err), - ) - return - } - - var deployedEvent WebSubAPIDeployedEvent - if err := json.Unmarshal(eventBytes, &deployedEvent); err != nil { - c.logger.Error("Failed to parse WebSub API deployment event", - slog.Any("error", err), - ) - return - } - - apiID := deployedEvent.Payload.APIID - if apiID == "" { - c.logger.Error("API ID is empty in WebSub API deployment event") - return - } - - c.logger.Info("Processing WebSub API deployment", - slog.String("api_id", apiID), - slog.String("deployment_id", deployedEvent.Payload.DeploymentID), - slog.String("correlation_id", deployedEvent.CorrelationID), - ) - - // Fetch WebSub API definition from control plane - zipData, err := c.apiUtilsService.FetchWebSubAPIDefinition(apiID) - if err != nil { - c.logger.Error("Failed to fetch WebSub API definition", - slog.String("api_id", apiID), - slog.Any("error", err), - ) - c.sendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "websub", "deploy", "failed", - deployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") - return - } - - yamlData, err := c.apiUtilsService.ExtractYAMLFromZip(zipData) - if err != nil { - c.logger.Error("Failed to extract YAML from WebSub API ZIP", - slog.String("api_id", apiID), - slog.Any("error", err), - ) - c.sendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "websub", "deploy", "failed", - deployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") - return - } - - // Ensure any {{ secret "handle" }} references in the YAML are in local - // storage before rendering. - c.syncSecretRefsFromYAML(yamlData, deployedEvent.CorrelationID) - - performedAt := deployedEvent.Payload.PerformedAt.Truncate(time.Millisecond) - if performedAt.IsZero() { - performedAt = time.Now().Truncate(time.Millisecond) - } - // Reuse the existing local UUID for a bottom-up (DP->CP) synced API so the - // control-plane deploy is an in-place update - result, err := c.apiUtilsService.CreateAPIFromYAML(yamlData, c.resolveLocalArtifactID(apiID), deployedEvent.Payload.DeploymentID, &performedAt, deployedEvent.CorrelationID, c.deploymentService) - if err != nil { - c.logger.Error("Failed to create WebSub API from YAML", - slog.String("api_id", apiID), - slog.Any("error", err), - ) - c.sendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "websub", "deploy", "failed", - deployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") - return - } - - if result.IsStale { - c.logger.Debug("Skipped stale WebSub API deploy event (newer version exists in DB)", - slog.String("api_id", apiID), - slog.String("deployment_id", deployedEvent.Payload.DeploymentID), - ) - return - } - - // Load platform-managed HMAC secrets into the webhook secret store. - if result.StoredConfig != nil { - c.syncHmacSecretsForArtifact(result.StoredConfig.UUID) - } - - c.sendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "websub", "deploy", "success", - deployedEvent.Payload.PerformedAt, "") - - c.logger.Info("Successfully processed WebSub API deployment event", - slog.String("api_id", apiID), - slog.String("correlation_id", deployedEvent.CorrelationID), - ) -} - -func (c *Client) handleWebSubAPIUndeployedEvent(event map[string]any) { - c.logger.Debug("WebSub API Undeployment Event", - slog.Any("payload", event["payload"]), - slog.Any("timestamp", event["timestamp"]), - slog.Any("correlationId", event["correlationId"]), - ) - - eventBytes, err := json.Marshal(event) - if err != nil { - c.logger.Error("Failed to marshal WebSub API undeployment event for parsing", - slog.Any("error", err), - ) - return - } - - var undeployedEvent WebSubAPIUndeployedEvent - if err := json.Unmarshal(eventBytes, &undeployedEvent); err != nil { - c.logger.Error("Failed to parse WebSub API undeployment event", - slog.Any("error", err), - ) - return - } - - apiID := undeployedEvent.Payload.APIID - if apiID == "" { - c.logger.Error("API ID is empty in WebSub API undeployment event") - return - } - - apiConfig, err := c.findAPIConfig(apiID) - if err != nil { - if storage.IsNotFoundError(err) { - c.logger.Warn("WebSub API configuration not found for undeployment", - slog.String("api_id", apiID), - ) - c.sendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "websub", "undeploy", "success", - undeployedEvent.Payload.PerformedAt, "") - return - } - c.logger.Error("Failed to fetch WebSub API configuration for undeployment", - slog.String("api_id", apiID), - slog.String("correlation_id", undeployedEvent.CorrelationID), - slog.Any("error", err), - ) - c.sendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "websub", "undeploy", "failed", - undeployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") - return - } - - if apiConfig.DeploymentID != "" && undeployedEvent.Payload.DeploymentID != "" && - apiConfig.DeploymentID != undeployedEvent.Payload.DeploymentID { - c.logger.Warn("Ignoring stale WebSub API undeploy event: deployment ID mismatch", - slog.String("api_id", apiID), - slog.String("event_deployment_id", undeployedEvent.Payload.DeploymentID), - slog.String("current_deployment_id", apiConfig.DeploymentID), - ) - c.sendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "websub", "undeploy", "failed", - undeployedEvent.Payload.PerformedAt, "DEPLOYMENT_ID_MISMATCH") - return - } - - performedAt := undeployedEvent.Payload.PerformedAt.Truncate(time.Millisecond) - if performedAt.IsZero() { - performedAt = time.Now().Truncate(time.Millisecond) - } - apiConfig.DesiredState = models.StateUndeployed - apiConfig.DeploymentID = undeployedEvent.Payload.DeploymentID - apiConfig.DeployedAt = &performedAt - apiConfig.UpdatedAt = time.Now() - - affected, err := c.db.UpsertConfig(apiConfig) - if err != nil { - c.logger.Error("Failed to upsert config for WebSub API undeployment", - slog.String("api_id", apiID), - slog.Any("error", err), - ) - c.sendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "websub", "undeploy", "failed", - undeployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") - return - } - if !affected { - c.logger.Debug("Skipped stale WebSub API undeploy event (newer version exists in DB)", - slog.String("api_id", apiID), - slog.String("deployment_id", undeployedEvent.Payload.DeploymentID), - ) - return - } - - evt := eventhub.Event{ - EventType: eventhub.EventTypeAPI, - Action: "UPDATE", - EntityID: apiID, - EventID: undeployedEvent.CorrelationID, - } - if err := c.eventHub.PublishEvent(c.gatewayID, evt); err != nil { - c.logger.Error("Failed to publish WebSub API undeployment event", slog.Any("error", err)) - } - - c.sendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "websub", "undeploy", "success", - undeployedEvent.Payload.PerformedAt, "") - - c.logger.Info("Successfully processed WebSub API undeployment event", - slog.String("api_id", apiID), - slog.String("correlation_id", undeployedEvent.CorrelationID), - ) -} - -func (c *Client) handleWebSubAPIDeletedEvent(event map[string]any) { - c.logger.Debug("WebSub API Deleted Event", - slog.Any("payload", event["payload"]), - slog.Any("timestamp", event["timestamp"]), - slog.Any("correlationId", event["correlationId"]), - ) - - eventBytes, err := json.Marshal(event) - if err != nil { - c.logger.Error("Failed to marshal WebSub API deleted event for parsing", - slog.Any("error", err), - ) - return - } - - var deletedEvent WebSubAPIDeletedEvent - if err := json.Unmarshal(eventBytes, &deletedEvent); err != nil { - c.logger.Error("Failed to parse WebSub API deleted event", - slog.Any("error", err), - ) - return - } - - apiID := deletedEvent.Payload.APIID - if apiID == "" { - c.logger.Error("API ID is empty in WebSub API deleted event") - return - } - - apiConfig, err := c.findAPIConfig(apiID) - if err != nil { - if storage.IsNotFoundError(err) { - c.logger.Warn("WebSub API configuration not found for deletion", - slog.String("api_id", apiID), - ) - c.cleanupHmacSecretsForArtifact(apiID) - return - } - c.logger.Error("Failed to fetch WebSub API configuration for deletion", - slog.String("api_id", apiID), - slog.String("correlation_id", deletedEvent.CorrelationID), - slog.Any("error", err), - ) - return - } - - c.performFullAPIDeletion(apiID, apiConfig, deletedEvent.CorrelationID) - c.cleanupHmacSecretsForArtifact(apiConfig.UUID) -} - -func (c *Client) handleWebBrokerAPIDeployedEvent(event map[string]any) { - c.logger.Debug("WebBroker API Deployment Event", - slog.Any("payload", event["payload"]), - slog.Any("timestamp", event["timestamp"]), - slog.Any("correlationId", event["correlationId"]), - ) - - eventBytes, err := json.Marshal(event) - if err != nil { - c.logger.Error("Failed to marshal WebBroker API deployment event for parsing", - slog.Any("error", err), - ) - return - } - - var deployedEvent WebBrokerAPIDeployedEvent - if err := json.Unmarshal(eventBytes, &deployedEvent); err != nil { - c.logger.Error("Failed to parse WebBroker API deployment event", - slog.Any("error", err), - ) - return - } - - apiID := deployedEvent.Payload.APIID - if apiID == "" { - c.logger.Error("API ID is empty in WebBroker API deployment event") - return - } - - c.logger.Info("Processing WebBroker API deployment", - slog.String("api_id", apiID), - slog.String("deployment_id", deployedEvent.Payload.DeploymentID), - slog.String("correlation_id", deployedEvent.CorrelationID), - ) - - // Fetch WebBroker API definition from control plane - zipData, err := c.apiUtilsService.FetchWebBrokerAPIDefinition(apiID) - if err != nil { - c.logger.Error("Failed to fetch WebBroker API definition", - slog.String("api_id", apiID), - slog.Any("error", err), - ) - c.sendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "webbroker", "deploy", "failed", - deployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") - return - } - - yamlData, err := c.apiUtilsService.ExtractYAMLFromZip(zipData) - if err != nil { - c.logger.Error("Failed to extract YAML from WebBroker API ZIP", - slog.String("api_id", apiID), - slog.Any("error", err), - ) - c.sendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "webbroker", "deploy", "failed", - deployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") - return - } - - // Ensure any {{ secret "handle" }} references in the YAML are in local - // storage before rendering. - c.syncSecretRefsFromYAML(yamlData, deployedEvent.CorrelationID) - - performedAt := deployedEvent.Payload.PerformedAt.Truncate(time.Millisecond) - if performedAt.IsZero() { - performedAt = time.Now().Truncate(time.Millisecond) - } - // Reuse the existing local UUID for a bottom-up (DP->CP) synced API so the - // control-plane deploy is an in-place update - result, err := c.apiUtilsService.CreateAPIFromYAML(yamlData, c.resolveLocalArtifactID(apiID), deployedEvent.Payload.DeploymentID, &performedAt, deployedEvent.CorrelationID, c.deploymentService) - if err != nil { - c.logger.Error("Failed to create WebBroker API from YAML", - slog.String("api_id", apiID), - slog.Any("error", err), - ) - c.sendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "webbroker", "deploy", "failed", - deployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") - return - } - - if result.IsStale { - c.logger.Debug("Skipped stale WebBroker API deploy event (newer version exists in DB)", - slog.String("api_id", apiID), - slog.String("deployment_id", deployedEvent.Payload.DeploymentID), - ) - return - } - - c.sendDeploymentAck(deployedEvent.Payload.DeploymentID, apiID, "webbroker", "deploy", "success", - deployedEvent.Payload.PerformedAt, "") - - c.logger.Info("Successfully processed WebBroker API deployment event", - slog.String("api_id", apiID), - slog.String("correlation_id", deployedEvent.CorrelationID), - ) -} - -func (c *Client) handleWebBrokerAPIUndeployedEvent(event map[string]any) { - c.logger.Debug("WebBroker API Undeployment Event", - slog.Any("payload", event["payload"]), - slog.Any("timestamp", event["timestamp"]), - slog.Any("correlationId", event["correlationId"]), - ) - - eventBytes, err := json.Marshal(event) - if err != nil { - c.logger.Error("Failed to marshal WebBroker API undeployment event for parsing", - slog.Any("error", err), - ) - return - } - - var undeployedEvent WebBrokerAPIUndeployedEvent - if err := json.Unmarshal(eventBytes, &undeployedEvent); err != nil { - c.logger.Error("Failed to parse WebBroker API undeployment event", - slog.Any("error", err), - ) - return - } - - apiID := undeployedEvent.Payload.APIID - if apiID == "" { - c.logger.Error("API ID is empty in WebBroker API undeployment event") - return - } - - apiConfig, err := c.findAPIConfig(apiID) - if err != nil { - if storage.IsNotFoundError(err) { - c.logger.Warn("WebBroker API configuration not found for undeployment", - slog.String("api_id", apiID), - ) - c.sendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "webbroker", "undeploy", "success", - undeployedEvent.Payload.PerformedAt, "") - return - } - c.logger.Error("Failed to fetch WebBroker API configuration for undeployment", - slog.String("api_id", apiID), - slog.String("correlation_id", undeployedEvent.CorrelationID), - slog.Any("error", err), - ) - c.sendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "webbroker", "undeploy", "failed", - undeployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") - return - } - - if apiConfig.DeploymentID != "" && undeployedEvent.Payload.DeploymentID != "" && - apiConfig.DeploymentID != undeployedEvent.Payload.DeploymentID { - c.logger.Warn("Ignoring stale WebBroker API undeploy event: deployment ID mismatch", - slog.String("api_id", apiID), - slog.String("event_deployment_id", undeployedEvent.Payload.DeploymentID), - slog.String("current_deployment_id", apiConfig.DeploymentID), - ) - c.sendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "webbroker", "undeploy", "failed", - undeployedEvent.Payload.PerformedAt, "DEPLOYMENT_ID_MISMATCH") - return - } - - performedAt := undeployedEvent.Payload.PerformedAt.Truncate(time.Millisecond) - if performedAt.IsZero() { - performedAt = time.Now().Truncate(time.Millisecond) - } - apiConfig.DesiredState = models.StateUndeployed - apiConfig.DeploymentID = undeployedEvent.Payload.DeploymentID - apiConfig.DeployedAt = &performedAt - apiConfig.UpdatedAt = time.Now() - - affected, err := c.db.UpsertConfig(apiConfig) - if err != nil { - c.logger.Error("Failed to upsert config for WebBroker API undeployment", - slog.String("api_id", apiID), - slog.Any("error", err), - ) - c.sendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "webbroker", "undeploy", "failed", - undeployedEvent.Payload.PerformedAt, "GATEWAY_PROCESSING_ERROR") - return - } - if !affected { - c.logger.Debug("Skipped stale WebBroker API undeploy event (newer version exists in DB)", - slog.String("api_id", apiID), - slog.String("deployment_id", undeployedEvent.Payload.DeploymentID), - ) - return - } - - evt := eventhub.Event{ - EventType: eventhub.EventTypeAPI, - Action: "UPDATE", - EntityID: apiID, - EventID: undeployedEvent.CorrelationID, - } - if err := c.eventHub.PublishEvent(c.gatewayID, evt); err != nil { - c.logger.Error("Failed to publish WebBroker API undeployment event", slog.Any("error", err)) - } - - c.sendDeploymentAck(undeployedEvent.Payload.DeploymentID, apiID, "webbroker", "undeploy", "success", - undeployedEvent.Payload.PerformedAt, "") - - c.logger.Info("Successfully processed WebBroker API undeployment event", - slog.String("api_id", apiID), - slog.String("correlation_id", undeployedEvent.CorrelationID), - ) -} - -func (c *Client) handleWebBrokerAPIDeletedEvent(event map[string]any) { - c.logger.Debug("WebBroker API Deleted Event", - slog.Any("payload", event["payload"]), - slog.Any("timestamp", event["timestamp"]), - slog.Any("correlationId", event["correlationId"]), - ) - - eventBytes, err := json.Marshal(event) - if err != nil { - c.logger.Error("Failed to marshal WebBroker API deleted event for parsing", - slog.Any("error", err), - ) - return - } - - var deletedEvent WebBrokerAPIDeletedEvent - if err := json.Unmarshal(eventBytes, &deletedEvent); err != nil { - c.logger.Error("Failed to parse WebBroker API deleted event", - slog.Any("error", err), - ) - return - } - - apiID := deletedEvent.Payload.APIID - if apiID == "" { - c.logger.Error("API ID is empty in WebBroker API deleted event") - return - } - - apiConfig, err := c.findAPIConfig(apiID) - if err != nil { - if storage.IsNotFoundError(err) { - c.logger.Warn("WebBroker API configuration not found for deletion; running orphan cleanup", - slog.String("api_id", apiID), - ) - c.cleanupOrphanedResources(apiID, deletedEvent.CorrelationID) - return - } - c.logger.Error("Failed to fetch WebBroker API configuration for deletion", - slog.String("api_id", apiID), - slog.String("correlation_id", deletedEvent.CorrelationID), - slog.Any("error", err), - ) - return - } - - c.performFullAPIDeletion(apiID, apiConfig, deletedEvent.CorrelationID) -} - func (c *Client) handleMCPProxyDeploymentEvent(event map[string]any) { c.logger.Debug("MCP Proxy Deployment Event", slog.Any("payload", event["payload"]), @@ -4603,29 +4023,3 @@ func (c *Client) pushGatewayManifestOnConnect(gatewayID string) { slog.Int("policy_count", len(policies)), ) } - -// handleWebSubAPIHmacSecretEvent handles websub.hmacsecret.created/updated/deleted events -// from platform-API. It re-syncs all platform-managed HMAC secrets for the affected artifact. -func (c *Client) handleWebSubAPIHmacSecretEvent(event map[string]any) { - payloadRaw, _ := event["payload"] - payloadBytes, err := json.Marshal(payloadRaw) - if err != nil { - c.logger.Error("Failed to marshal HMAC secret event payload", slog.Any("error", err)) - return - } - var payload platformHmacSecretEventPayload - if err := json.Unmarshal(payloadBytes, &payload); err != nil { - c.logger.Error("Failed to parse HMAC secret event payload", slog.Any("error", err)) - return - } - if payload.ArtifactUUID == "" { - c.logger.Warn("HMAC secret event missing artifactUuid, skipping") - return - } - c.logger.Info("Processing platform HMAC secret event", - slog.Any("type", event["type"]), - slog.String("artifact_uuid", payload.ArtifactUUID), - slog.String("secret_name", payload.SecretName), - ) - c.syncHmacSecretsForArtifact(payload.ArtifactUUID) -} diff --git a/gateway/gateway-controller/pkg/controlplane/eventgateway_hooks.go b/gateway/gateway-controller/pkg/controlplane/eventgateway_hooks.go new file mode 100644 index 000000000..3606d7eca --- /dev/null +++ b/gateway/gateway-controller/pkg/controlplane/eventgateway_hooks.go @@ -0,0 +1,145 @@ +/* + * Copyright (c) 2026, WSO2 LLC. (https://www.wso2.com). + * + * WSO2 LLC. licenses this file to you under the Apache License, + * Version 2.0 (the "License"); you may not use this file except + * in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package controlplane + +import ( + "log/slog" + "time" + + "github.com/wso2/api-platform/common/eventhub" + "github.com/wso2/api-platform/common/webhooksecret" + "github.com/wso2/api-platform/gateway/gateway-controller/pkg/models" + "github.com/wso2/api-platform/gateway/gateway-controller/pkg/storage" + "github.com/wso2/api-platform/gateway/gateway-controller/pkg/utils" +) + +// ControlPlaneEventGatewayHooks is the extension point through which an +// external event-gateway-controller binary supplies WebSub/WebBroker +// control-plane WebSocket event handling (deploy/undeploy/delete and HMAC +// secret sync). Core never implements this interface itself; it is only ever +// satisfied by code living outside this module. See +// SetControlPlaneEventGatewayHooks. +type ControlPlaneEventGatewayHooks interface { + HandleWebSubAPIDeployed(c *Client, event map[string]any) + HandleWebSubAPIUndeployed(c *Client, event map[string]any) + HandleWebSubAPIDeleted(c *Client, event map[string]any) + HandleWebSubAPIHmacSecretEvent(c *Client, event map[string]any) + HandleWebBrokerAPIDeployed(c *Client, event map[string]any) + HandleWebBrokerAPIUndeployed(c *Client, event map[string]any) + HandleWebBrokerAPIDeleted(c *Client, event map[string]any) +} + +// SetControlPlaneEventGatewayHooks registers the event-gateway control-plane +// extension. Passing nil (the default) means this binary has no event-gateway +// support compiled in — incoming websub.*/webbroker.* control-plane events +// are logged and dropped. +func (c *Client) SetControlPlaneEventGatewayHooks(h ControlPlaneEventGatewayHooks) { + c.eventGatewayHooks = h +} + +// The following exported accessors/wrappers exist solely so that a +// ControlPlaneEventGatewayHooks implementation living outside this module can +// reuse the same generic control-plane sync primitives every other kind uses, +// without duplicating them. + +// Logger returns the client's logger. +func (c *Client) Logger() *slog.Logger { + return c.logger +} + +// DB returns the client's storage handle. +func (c *Client) DB() storage.Storage { + return c.db +} + +// EventHub returns the client's EventHub instance. +func (c *Client) EventHub() eventhub.EventHub { + return c.eventHub +} + +// GatewayID returns the gateway ID this client is running for. +func (c *Client) GatewayID() string { + return c.gatewayID +} + +// APIUtilsService returns the client's platform-API HTTP helper. +func (c *Client) APIUtilsService() *utils.APIUtilsService { + return c.apiUtilsService +} + +// DeploymentService returns the client's generic API deployment service. +func (c *Client) DeploymentService() *utils.APIDeploymentService { + return c.deploymentService +} + +// WebhookSecretStore returns the client's in-memory webhook-secret store, or +// nil if none was configured. +func (c *Client) WebhookSecretStore() *webhooksecret.WebhookSecretStore { + return c.webhookSecretStore +} + +// RefreshWebhookSecretSnapshot refreshes the webhook-secret xDS snapshot via +// the configured WebhookSecretSnapshotRefresher. No-op if none was configured. +func (c *Client) RefreshWebhookSecretSnapshot() error { + if c.webhookSecretSnapshotManager == nil { + return nil + } + return c.webhookSecretSnapshotManager.RefreshSnapshot() +} + +// FindAPIConfig exposes findAPIConfig for use by ControlPlaneEventGatewayHooks implementations. +func (c *Client) FindAPIConfig(apiID string) (*models.StoredConfig, error) { + return c.findAPIConfig(apiID) +} + +// ResolveLocalArtifactID exposes resolveLocalArtifactID for use by ControlPlaneEventGatewayHooks implementations. +func (c *Client) ResolveLocalArtifactID(id string) string { + return c.resolveLocalArtifactID(id) +} + +// dispatchEventGatewayHook invokes fn with the registered hooks, or logs and +// drops the event if this binary has no event-gateway support compiled in. +func (c *Client) dispatchEventGatewayHook(eventType any, fn func(ControlPlaneEventGatewayHooks)) { + if c.eventGatewayHooks == nil { + c.logger.Warn("Received event-gateway control-plane event but no event-gateway support is compiled into this binary", + slog.Any("type", eventType)) + return + } + fn(c.eventGatewayHooks) +} + +// SendDeploymentAck exposes sendDeploymentAck for use by ControlPlaneEventGatewayHooks implementations. +func (c *Client) SendDeploymentAck(deploymentID, artifactID, resourceType, action, status string, performedAt time.Time, errorCode string) { + c.sendDeploymentAck(deploymentID, artifactID, resourceType, action, status, performedAt, errorCode) +} + +// PerformFullAPIDeletion exposes performFullAPIDeletion for use by ControlPlaneEventGatewayHooks implementations. +func (c *Client) PerformFullAPIDeletion(apiID string, apiConfig *models.StoredConfig, correlationID string) { + c.performFullAPIDeletion(apiID, apiConfig, correlationID) +} + +// CleanupOrphanedResources exposes cleanupOrphanedResources for use by ControlPlaneEventGatewayHooks implementations. +func (c *Client) CleanupOrphanedResources(apiID, correlationID string) { + c.cleanupOrphanedResources(apiID, correlationID) +} + +// SyncSecretRefsFromYAML exposes syncSecretRefsFromYAML for use by ControlPlaneEventGatewayHooks implementations. +func (c *Client) SyncSecretRefsFromYAML(yamlData []byte, correlationID string) { + c.syncSecretRefsFromYAML(yamlData, correlationID) +} diff --git a/gateway/gateway-controller/pkg/controlplane/events.go b/gateway/gateway-controller/pkg/controlplane/events.go index 4e9fb6ac6..ecde9ae79 100644 --- a/gateway/gateway-controller/pkg/controlplane/events.go +++ b/gateway/gateway-controller/pkg/controlplane/events.go @@ -285,91 +285,11 @@ type MCPProxyDeletedEvent struct { CorrelationID string `json:"correlationId"` } -// WebSubAPIDeployedEventPayload represents the payload of a WebSub API deployment event -type WebSubAPIDeployedEventPayload struct { - APIID string `json:"apiId"` - DeploymentID string `json:"deploymentId"` - PerformedAt time.Time `json:"performedAt"` -} - -// WebSubAPIDeployedEvent represents the complete WebSub API deployment event -type WebSubAPIDeployedEvent struct { - Type string `json:"type"` - Payload WebSubAPIDeployedEventPayload `json:"payload"` - Timestamp string `json:"timestamp"` - CorrelationID string `json:"correlationId"` -} - -// WebSubAPIUndeployedEventPayload represents the payload of a WebSub API undeployment event -type WebSubAPIUndeployedEventPayload struct { - APIID string `json:"apiId"` - DeploymentID string `json:"deploymentId"` - PerformedAt time.Time `json:"performedAt"` -} - -// WebSubAPIUndeployedEvent represents the complete WebSub API undeployment event -type WebSubAPIUndeployedEvent struct { - Type string `json:"type"` - Payload WebSubAPIUndeployedEventPayload `json:"payload"` - Timestamp string `json:"timestamp"` - CorrelationID string `json:"correlationId"` -} - -// WebSubAPIDeletedEventPayload represents the payload of a WebSub API deletion event -type WebSubAPIDeletedEventPayload struct { - APIID string `json:"apiId"` -} - -// WebSubAPIDeletedEvent represents the complete WebSub API deletion event -type WebSubAPIDeletedEvent struct { - Type string `json:"type"` - Payload WebSubAPIDeletedEventPayload `json:"payload"` - Timestamp string `json:"timestamp"` - CorrelationID string `json:"correlationId"` -} - -// WebBrokerAPIDeployedEventPayload represents the payload of a WebBroker API deployment event -type WebBrokerAPIDeployedEventPayload struct { - APIID string `json:"apiId"` - DeploymentID string `json:"deploymentId"` - PerformedAt time.Time `json:"performedAt"` -} - -// WebBrokerAPIDeployedEvent represents the complete WebBroker API deployment event -type WebBrokerAPIDeployedEvent struct { - Type string `json:"type"` - Payload WebBrokerAPIDeployedEventPayload `json:"payload"` - Timestamp string `json:"timestamp"` - CorrelationID string `json:"correlationId"` -} - -// WebBrokerAPIUndeployedEventPayload represents the payload of a WebBroker API undeployment event -type WebBrokerAPIUndeployedEventPayload struct { - APIID string `json:"apiId"` - DeploymentID string `json:"deploymentId"` - PerformedAt time.Time `json:"performedAt"` -} - -// WebBrokerAPIUndeployedEvent represents the complete WebBroker API undeployment event -type WebBrokerAPIUndeployedEvent struct { - Type string `json:"type"` - Payload WebBrokerAPIUndeployedEventPayload `json:"payload"` - Timestamp string `json:"timestamp"` - CorrelationID string `json:"correlationId"` -} - -// WebBrokerAPIDeletedEventPayload represents the payload of a WebBroker API deletion event -type WebBrokerAPIDeletedEventPayload struct { - APIID string `json:"apiId"` -} - -// WebBrokerAPIDeletedEvent represents the complete WebBroker API deletion event -type WebBrokerAPIDeletedEvent struct { - Type string `json:"type"` - Payload WebBrokerAPIDeletedEventPayload `json:"payload"` - Timestamp string `json:"timestamp"` - CorrelationID string `json:"correlationId"` -} +// Note: WebSub/WebBroker deploy/undeploy/delete event payload types +// (WebSubAPIDeployedEvent, WebBrokerAPIDeployedEvent, etc.) are NOT defined +// here. They are event-gateway-specific and owned by the +// event-gateway-controller module (event-gateway/gateway-controller/pkg/controlplanehooks), +// which implements controlplane.ControlPlaneEventGatewayHooks. // SubscriptionCreatedEventPayload represents the payload of a subscription created event. type SubscriptionCreatedEventPayload struct { diff --git a/gateway/gateway-controller/pkg/storage/gateway-controller-db.postgres.sql b/gateway/gateway-controller/pkg/storage/gateway-controller-db.postgres.sql index f9ee5e65a..2a1e6539d 100644 --- a/gateway/gateway-controller/pkg/storage/gateway-controller-db.postgres.sql +++ b/gateway/gateway-controller/pkg/storage/gateway-controller-db.postgres.sql @@ -37,21 +37,10 @@ CREATE TABLE IF NOT EXISTS rest_apis ( FOREIGN KEY(gateway_id, uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE ); -CREATE TABLE IF NOT EXISTS websub_apis ( - uuid TEXT NOT NULL, - gateway_id TEXT NOT NULL, - configuration TEXT NOT NULL, - PRIMARY KEY (gateway_id, uuid), - FOREIGN KEY(gateway_id, uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE -); - -CREATE TABLE IF NOT EXISTS webbroker_apis ( - uuid TEXT NOT NULL, - gateway_id TEXT NOT NULL, - configuration TEXT NOT NULL, - PRIMARY KEY (gateway_id, uuid), - FOREIGN KEY(gateway_id, uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE -); +-- Note: websub_apis, webbroker_apis, and webhook_secrets tables are NOT defined +-- here. They are event-gateway-specific and owned by the event-gateway-controller +-- module (event-gateway/gateway-controller/pkg/dbschema), which applies its own +-- supplemental DDL against this same database at startup. CREATE TABLE IF NOT EXISTS llm_providers ( uuid TEXT NOT NULL, @@ -231,21 +220,5 @@ CREATE TABLE IF NOT EXISTS secrets ( PRIMARY KEY (gateway_id, handle) ); --- Per-API HMAC secrets for the websub-hmac-auth policy. --- Ciphertext is AES-256-GCM encrypted; plaintext is never stored. -CREATE TABLE IF NOT EXISTS webhook_secrets ( - uuid TEXT NOT NULL, - gateway_id TEXT NOT NULL, - artifact_uuid TEXT NOT NULL, - name TEXT NOT NULL, - display_name TEXT NOT NULL, - ciphertext BYTEA NOT NULL, - status TEXT NOT NULL DEFAULT 'active' CHECK(status IN ('active', 'revoked')), - created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, - updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, - PRIMARY KEY (gateway_id, uuid), - UNIQUE (gateway_id, artifact_uuid, name), - FOREIGN KEY (gateway_id, artifact_uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE -); - -CREATE INDEX IF NOT EXISTS idx_webhook_secrets_artifact ON webhook_secrets(gateway_id, artifact_uuid); +-- Note: webhook_secrets (per-API HMAC secrets for the websub-hmac-auth policy) +-- is also owned by event-gateway/gateway-controller/pkg/dbschema — see note above. diff --git a/gateway/gateway-controller/pkg/storage/gateway-controller-db.sql b/gateway/gateway-controller/pkg/storage/gateway-controller-db.sql index 6d4bb1efd..3f94f1439 100644 --- a/gateway/gateway-controller/pkg/storage/gateway-controller-db.sql +++ b/gateway/gateway-controller/pkg/storage/gateway-controller-db.sql @@ -37,21 +37,10 @@ CREATE TABLE IF NOT EXISTS rest_apis ( FOREIGN KEY(gateway_id, uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE ); -CREATE TABLE IF NOT EXISTS websub_apis ( - uuid TEXT NOT NULL, - gateway_id TEXT NOT NULL, - configuration TEXT NOT NULL, - PRIMARY KEY (gateway_id, uuid), - FOREIGN KEY(gateway_id, uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE -); - -CREATE TABLE IF NOT EXISTS webbroker_apis ( - uuid TEXT NOT NULL, - gateway_id TEXT NOT NULL, - configuration TEXT NOT NULL, - PRIMARY KEY (gateway_id, uuid), - FOREIGN KEY(gateway_id, uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE -); +-- Note: websub_apis, webbroker_apis, and webhook_secrets tables are NOT defined +-- here. They are event-gateway-specific and owned by the event-gateway-controller +-- module (event-gateway/gateway-controller/pkg/dbschema), which applies its own +-- supplemental DDL against this same database at startup. CREATE TABLE IF NOT EXISTS llm_providers ( uuid TEXT NOT NULL, @@ -289,23 +278,7 @@ CREATE TABLE IF NOT EXISTS secrets ( PRIMARY KEY (gateway_id, handle) ); --- Per-API HMAC secrets for the websub-hmac-auth policy. --- Ciphertext is AES-256-GCM encrypted; plaintext is never stored. -CREATE TABLE IF NOT EXISTS webhook_secrets ( - uuid TEXT NOT NULL, - gateway_id TEXT NOT NULL, - artifact_uuid TEXT NOT NULL, - name TEXT NOT NULL, - display_name TEXT NOT NULL, - ciphertext BLOB NOT NULL, - status TEXT NOT NULL DEFAULT 'active' CHECK(status IN ('active', 'revoked')), - created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, - updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, - PRIMARY KEY (gateway_id, uuid), - UNIQUE (gateway_id, artifact_uuid, name), - FOREIGN KEY (gateway_id, artifact_uuid) REFERENCES artifacts(gateway_id, uuid) ON DELETE CASCADE -); - -CREATE INDEX IF NOT EXISTS idx_webhook_secrets_artifact ON webhook_secrets(gateway_id, artifact_uuid); +-- Note: webhook_secrets (per-API HMAC secrets for the websub-hmac-auth policy) +-- is also owned by event-gateway/gateway-controller/pkg/dbschema — see note above. PRAGMA user_version = 4; diff --git a/gateway/gateway-controller/pkg/storage/gateway-controller-db.sqlserver.sql b/gateway/gateway-controller/pkg/storage/gateway-controller-db.sqlserver.sql index d024e5e00..568dd7c65 100644 --- a/gateway/gateway-controller/pkg/storage/gateway-controller-db.sqlserver.sql +++ b/gateway/gateway-controller/pkg/storage/gateway-controller-db.sqlserver.sql @@ -50,23 +50,10 @@ CREATE TABLE dbo.rest_apis ( FOREIGN KEY(gateway_id, uuid) REFERENCES dbo.artifacts(gateway_id, uuid) ON DELETE CASCADE ); -IF OBJECT_ID(N'dbo.websub_apis', N'U') IS NULL -CREATE TABLE dbo.websub_apis ( - uuid NVARCHAR(64) NOT NULL, - gateway_id NVARCHAR(64) NOT NULL, - configuration NVARCHAR(MAX) NOT NULL, - PRIMARY KEY (gateway_id, uuid), - FOREIGN KEY(gateway_id, uuid) REFERENCES dbo.artifacts(gateway_id, uuid) ON DELETE CASCADE -); - -IF OBJECT_ID(N'dbo.webbroker_apis', N'U') IS NULL -CREATE TABLE dbo.webbroker_apis ( - uuid NVARCHAR(64) NOT NULL, - gateway_id NVARCHAR(64) NOT NULL, - configuration NVARCHAR(MAX) NOT NULL, - PRIMARY KEY (gateway_id, uuid), - FOREIGN KEY(gateway_id, uuid) REFERENCES dbo.artifacts(gateway_id, uuid) ON DELETE CASCADE -); +-- Note: websub_apis, webbroker_apis, and webhook_secrets tables are NOT defined +-- here. They are event-gateway-specific and owned by the event-gateway-controller +-- module (event-gateway/gateway-controller/pkg/dbschema), which applies its own +-- supplemental DDL against this same database at startup. IF OBJECT_ID(N'dbo.llm_providers', N'U') IS NULL CREATE TABLE dbo.llm_providers ( @@ -263,22 +250,5 @@ CREATE TABLE dbo.secrets ( PRIMARY KEY (gateway_id, handle) ); --- Table for encrypted per-artifact webhook secrets (gateway-scoped) -IF OBJECT_ID(N'dbo.webhook_secrets', N'U') IS NULL -CREATE TABLE dbo.webhook_secrets ( - uuid NVARCHAR(64) NOT NULL, - gateway_id NVARCHAR(64) NOT NULL, - artifact_uuid NVARCHAR(64) NOT NULL, - name NVARCHAR(255) NOT NULL, - display_name NVARCHAR(255) NOT NULL, - ciphertext VARBINARY(MAX) NOT NULL, - status NVARCHAR(20) NOT NULL CHECK(status IN ('active', 'revoked')) DEFAULT 'active', - created_at DATETIME2(7) NOT NULL DEFAULT SYSUTCDATETIME(), - updated_at DATETIME2(7) NOT NULL DEFAULT SYSUTCDATETIME(), - PRIMARY KEY (gateway_id, uuid), - CONSTRAINT uq_webhook_secrets_artifact_name UNIQUE (gateway_id, artifact_uuid, name), - FOREIGN KEY (gateway_id, artifact_uuid) REFERENCES dbo.artifacts(gateway_id, uuid) ON DELETE CASCADE -); - -IF NOT EXISTS (SELECT 1 FROM sys.indexes WHERE name = N'idx_webhook_secrets_artifact' AND object_id = OBJECT_ID(N'dbo.webhook_secrets')) -CREATE INDEX idx_webhook_secrets_artifact ON dbo.webhook_secrets(gateway_id, artifact_uuid); +-- Note: webhook_secrets (per-API HMAC secrets for the websub-hmac-auth policy) +-- is also owned by event-gateway/gateway-controller/pkg/dbschema — see note above. diff --git a/gateway/gateway-controller/pkg/storage/sql_store.go b/gateway/gateway-controller/pkg/storage/sql_store.go index dec66da84..8167a0d70 100644 --- a/gateway/gateway-controller/pkg/storage/sql_store.go +++ b/gateway/gateway-controller/pkg/storage/sql_store.go @@ -279,10 +279,6 @@ func kindToResourceTable(kind string) (string, error) { switch kind { case "RestApi": return "rest_apis", nil - case "WebSubApi": - return "websub_apis", nil - case "WebBrokerApi": - return "webbroker_apis", nil case "LlmProvider": return "llm_providers", nil case "LlmProviderTemplate": @@ -292,10 +288,34 @@ func kindToResourceTable(kind string) (string, error) { case "Mcp": return "mcp_proxies", nil default: + if table, ok := extraResourceTables[kind]; ok { + return table, nil + } return "", fmt.Errorf("unknown kind: %s", kind) } } +// extraResourceTables holds kind->table entries for kinds not known to core +// (e.g. "WebSubApi"/"WebBrokerApi") — registered by an event-gateway-controller +// binary via RegisterKindResourceTable. This mirrors kindUnmarshalers below: +// core's own schema scripts never define these tables (see +// event-gateway/gateway-controller/pkg/dbschema), only the module that +// registers a kind here knows which table backs it. +var extraResourceTables = map[string]string{} + +// builtinResourceTables lists the per-kind tables core defines natively. +// GetAllConfigs unions these with every table in extraResourceTables so +// cross-kind listing also covers kinds registered by an external module. +var builtinResourceTables = []string{"rest_apis", "llm_providers", "llm_proxies", "mcp_proxies"} + +// RegisterKindResourceTable registers the resource table name for an artifact +// kind not known to core. Intended to be called from an init() (or equivalent +// startup wiring) in a binary that links in support for that kind, before any +// Storage method for that kind is used. +func RegisterKindResourceTable(kind, table string) { + extraResourceTables[kind] = table +} + // kindUnmarshalers holds JSON-unmarshaling functions for artifact kinds not // known to core — registered by an event-gateway-controller binary (e.g. for // "WebSubApi"/"WebBrokerApi") via RegisterKindUnmarshaler. This mirrors the @@ -993,57 +1013,38 @@ func (s *sqlStore) GetConfigByKindNameAndVersion(kind, displayName, version stri return &cfg, nil } -// GetAllConfigs retrieves all artifact configurations +// getAllConfigsColumns is the shared SELECT list every per-table block in +// GetAllConfigs uses, so scanConfigRows sees a consistent column order +// regardless of which table (built-in or externally-registered) a row came from. +const getAllConfigsColumns = `a.uuid, a.kind, a.handle, a.display_name, a.version, a.data_version, r.configuration, a.desired_state, + a.deployment_id, a.origin, a.created_at, a.updated_at, a.deployed_at, + a.cp_sync_status, a.cp_sync_info, a.cp_artifact_id` + +// GetAllConfigs retrieves all artifact configurations. // TODO: (renuka) Remove this method once the in memory cache is removed. func (s *sqlStore) GetAllConfigs() ([]*models.StoredConfig, error) { - // Use UNION ALL across all type tables joined with artifacts - query := ` - SELECT a.uuid, a.kind, a.handle, a.display_name, a.version, a.data_version, r.configuration, a.desired_state, - a.deployment_id, a.origin, a.created_at, a.updated_at, a.deployed_at, - a.cp_sync_status, a.cp_sync_info, a.cp_artifact_id + // Union every built-in resource table with every externally-registered one + // (see RegisterKindResourceTable) so cross-kind listing also covers kinds + // core doesn't know about natively (e.g. WebSubApi/WebBrokerApi). + tables := append([]string{}, builtinResourceTables...) + for _, table := range extraResourceTables { + tables = append(tables, table) + } + sort.Strings(tables) + + blocks := make([]string, len(tables)) + args := make([]interface{}, len(tables)) + for i, table := range tables { + blocks[i] = fmt.Sprintf(` + SELECT %s FROM artifacts a - JOIN rest_apis r ON a.uuid = r.uuid AND a.gateway_id = r.gateway_id - WHERE a.gateway_id = ? - - UNION ALL - - SELECT a.uuid, a.kind, a.handle, a.display_name, a.version, a.data_version, w.configuration, a.desired_state, - a.deployment_id, a.origin, a.created_at, a.updated_at, a.deployed_at, - a.cp_sync_status, a.cp_sync_info, a.cp_artifact_id - FROM artifacts a - JOIN websub_apis w ON a.uuid = w.uuid AND a.gateway_id = w.gateway_id - WHERE a.gateway_id = ? - - UNION ALL - - SELECT a.uuid, a.kind, a.handle, a.display_name, a.version, a.data_version, lp.configuration, a.desired_state, - a.deployment_id, a.origin, a.created_at, a.updated_at, a.deployed_at, - a.cp_sync_status, a.cp_sync_info, a.cp_artifact_id - FROM artifacts a - JOIN llm_providers lp ON a.uuid = lp.uuid AND a.gateway_id = lp.gateway_id - WHERE a.gateway_id = ? - - UNION ALL - - SELECT a.uuid, a.kind, a.handle, a.display_name, a.version, a.data_version, lx.configuration, a.desired_state, - a.deployment_id, a.origin, a.created_at, a.updated_at, a.deployed_at, - a.cp_sync_status, a.cp_sync_info, a.cp_artifact_id - FROM artifacts a - JOIN llm_proxies lx ON a.uuid = lx.uuid AND a.gateway_id = lx.gateway_id - WHERE a.gateway_id = ? - - UNION ALL - - SELECT a.uuid, a.kind, a.handle, a.display_name, a.version, a.data_version, m.configuration, a.desired_state, - a.deployment_id, a.origin, a.created_at, a.updated_at, a.deployed_at, - a.cp_sync_status, a.cp_sync_info, a.cp_artifact_id - FROM artifacts a - JOIN mcp_proxies m ON a.uuid = m.uuid AND a.gateway_id = m.gateway_id - WHERE a.gateway_id = ? - ORDER BY created_at DESC - ` + JOIN %s r ON a.uuid = r.uuid AND a.gateway_id = r.gateway_id + WHERE a.gateway_id = ?`, getAllConfigsColumns, table) + args[i] = s.gatewayId + } + query := strings.Join(blocks, "\n\n\t\tUNION ALL\n") + "\n\t\tORDER BY created_at DESC" - rows, err := s.query(query, s.gatewayId, s.gatewayId, s.gatewayId, s.gatewayId, s.gatewayId) + rows, err := s.query(query, args...) if err != nil { return nil, fmt.Errorf("failed to query configurations: %w", err) } diff --git a/gateway/gateway-controller/pkg/storage/sqlite_test.go b/gateway/gateway-controller/pkg/storage/sqlite_test.go index 05c21d1ad..2f786c20f 100644 --- a/gateway/gateway-controller/pkg/storage/sqlite_test.go +++ b/gateway/gateway-controller/pkg/storage/sqlite_test.go @@ -84,7 +84,6 @@ func TestSQLiteStorage_SchemaInitialization(t *testing.T) { tables := []string{ "artifacts", "rest_apis", - "websub_apis", "llm_providers", "llm_proxies", "mcp_proxies", diff --git a/gateway/gateway-controller/pkg/utils/api_utils.go b/gateway/gateway-controller/pkg/utils/api_utils.go index 01537ba23..53f0e7383 100644 --- a/gateway/gateway-controller/pkg/utils/api_utils.go +++ b/gateway/gateway-controller/pkg/utils/api_utils.go @@ -44,12 +44,17 @@ import ( "github.com/wso2/api-platform/gateway/gateway-controller/pkg/models" ) +// defaultMaxResponseBytes is the fallback ceiling applied to response bodies +// read from the platform-API when PlatformAPIConfig.MaxResponseBytes is unset. +const defaultMaxResponseBytes = 50 << 20 // 50 MiB + // PlatformAPIConfig contains configuration for fetching API definitions type PlatformAPIConfig struct { BaseURL string // Base URL for API requests Token string // Authentication token InsecureSkipVerify bool // Skip TLS verification Timeout time.Duration // Request timeout + MaxResponseBytes int64 // Maximum bytes read from a single response body } // APIUtilsService provides utilities for API operations @@ -74,6 +79,9 @@ func NewAPIUtilsService(config PlatformAPIConfig, logger *slog.Logger) *APIUtils if config.Timeout == 0 { config.Timeout = 30 * time.Second } + if config.MaxResponseBytes <= 0 { + config.MaxResponseBytes = defaultMaxResponseBytes + } if config.InsecureSkipVerify { logger.Warn("TLS certificate verification disabled for API utils (insecure_skip_verify=true)") } @@ -640,16 +648,22 @@ func (s *APIUtilsService) FetchMCPProxyDefinition(proxyID string) ([]byte, error return bodyBytes, nil } -// FetchWebSubAPIDefinition downloads the WebSub API definition as a zip file from the control plane -func (s *APIUtilsService) FetchWebSubAPIDefinition(apiID string) ([]byte, error) { - apiURL := s.getBaseURL() + "/websub-apis/" + apiID - - s.logger.Debug("Fetching WebSub API definition", - slog.String("api_id", apiID), - slog.String("url", apiURL), +// FetchResourceZip performs a generic authenticated GET against +// {baseURL}{resourcePath}, expecting a zip response, and returns the raw +// bytes. resourceLabel is used only for log/error messages (e.g. "WebSub API +// definition"). Extracted as a reusable primitive so external modules with +// their own resource kinds (e.g. event-gateway-controller's WebSub/WebBroker +// APIs) can fetch zip definitions without duplicating this HTTP/auth +// boilerplate the way FetchAPIDefinition/FetchLLMProviderDefinition/ +// FetchLLMProxyDefinition above do for kinds known to core. +func (s *APIUtilsService) FetchResourceZip(resourcePath, resourceLabel string) ([]byte, error) { + url := s.getBaseURL() + resourcePath + + s.logger.Debug("Fetching "+resourceLabel, + slog.String("url", url), ) - req, err := http.NewRequest("GET", apiURL, nil) + req, err := http.NewRequest("GET", url, nil) if err != nil { return nil, fmt.Errorf("failed to create request: %w", err) } @@ -659,67 +673,60 @@ func (s *APIUtilsService) FetchWebSubAPIDefinition(apiID string) ([]byte, error) resp, err := s.client.Do(req) if err != nil { - return nil, fmt.Errorf("failed to fetch WebSub API definition: %w", err) + return nil, fmt.Errorf("failed to fetch %s: %w", resourceLabel, err) } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { - bodyBytes, _ := io.ReadAll(resp.Body) - return nil, fmt.Errorf("WebSub API request failed with status %d: %s", resp.StatusCode, string(bodyBytes)) + bodyBytes, _ := io.ReadAll(io.LimitReader(resp.Body, s.config.MaxResponseBytes)) + return nil, fmt.Errorf("%s request failed with status %d: %s", resourceLabel, resp.StatusCode, string(bodyBytes)) } - bodyBytes, err := io.ReadAll(resp.Body) + bodyBytes, err := io.ReadAll(io.LimitReader(resp.Body, s.config.MaxResponseBytes+1)) if err != nil { return nil, fmt.Errorf("failed to read response body: %w", err) } + if int64(len(bodyBytes)) > s.config.MaxResponseBytes { + return nil, fmt.Errorf("%s response exceeds maximum allowed size", resourceLabel) + } - s.logger.Debug("Successfully fetched WebSub API definition", - slog.String("api_id", apiID), + s.logger.Debug("Successfully fetched "+resourceLabel, slog.Int("size_bytes", len(bodyBytes)), ) return bodyBytes, nil } -// FetchWebBrokerAPIDefinition downloads the WebBroker API definition as a zip file from the control plane -func (s *APIUtilsService) FetchWebBrokerAPIDefinition(apiID string) ([]byte, error) { - apiURL := s.getBaseURL() + "/webbroker-apis/" + apiID +// FetchResourceJSON performs a generic authenticated GET against +// {baseURL}{resourcePath}, expecting a JSON response, and decodes it into out. +// See FetchResourceZip for why this exists as a reusable primitive. +func (s *APIUtilsService) FetchResourceJSON(resourcePath, resourceLabel string, out any) error { + url := s.getBaseURL() + resourcePath - s.logger.Debug("Fetching WebBroker API definition", - slog.String("api_id", apiID), - slog.String("url", apiURL), - ) - - req, err := http.NewRequest("GET", apiURL, nil) + req, err := http.NewRequest("GET", url, nil) if err != nil { - return nil, fmt.Errorf("failed to create request: %w", err) + return fmt.Errorf("failed to create request: %w", err) } req.Header.Add("api-key", s.config.Token) - req.Header.Add("Accept", "application/zip") + req.Header.Add("Accept", "application/json") resp, err := s.client.Do(req) if err != nil { - return nil, fmt.Errorf("failed to fetch WebBroker API definition: %w", err) + return fmt.Errorf("failed to fetch %s: %w", resourceLabel, err) } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { - bodyBytes, _ := io.ReadAll(resp.Body) - return nil, fmt.Errorf("WebBroker API request failed with status %d: %s", resp.StatusCode, string(bodyBytes)) + bodyBytes, _ := io.ReadAll(io.LimitReader(resp.Body, s.config.MaxResponseBytes)) + return fmt.Errorf("%s request failed with status %d: %s", resourceLabel, resp.StatusCode, string(bodyBytes)) } - bodyBytes, err := io.ReadAll(resp.Body) - if err != nil { - return nil, fmt.Errorf("failed to read response body: %w", err) + if err := json.NewDecoder(io.LimitReader(resp.Body, s.config.MaxResponseBytes)).Decode(out); err != nil { + return fmt.Errorf("failed to decode %s response: %w", resourceLabel, err) } - s.logger.Debug("Successfully fetched WebBroker API definition", - slog.String("api_id", apiID), - slog.Int("size_bytes", len(bodyBytes)), - ) - - return bodyBytes, nil + return nil } // CreateMCPProxyFromYAML creates an MCP proxy configuration from YAML data using the MCP deployment service @@ -1235,69 +1242,11 @@ func MapToStruct(data map[string]interface{}, out interface{}) error { return nil } -// platformHmacSecretInfo is the per-secret DTO returned by the internal HMAC endpoint. -type platformHmacSecretInfo struct { - Name string `json:"name"` - Secret string `json:"secret"` -} - -// platformHmacSecretsResponse is the response body from GET /websub-apis/:id/secrets. -type platformHmacSecretsResponse struct { - ArtifactID string `json:"artifactId"` - Secrets []platformHmacSecretInfo `json:"secrets"` -} - -// HmacSecretInfo is the public view of a platform-managed HMAC secret. -type HmacSecretInfo struct { - Name string - Plaintext string -} - -// FetchWebSubAPIHmacSecrets fetches the plaintext HMAC secrets for a WebSub API artifact -// from the platform-API internal endpoint. -func (s *APIUtilsService) FetchWebSubAPIHmacSecrets(artifactID string) ([]HmacSecretInfo, error) { - secretsURL := s.getBaseURL() + "/websub-apis/" + artifactID + "/secrets" - - s.logger.Debug("Fetching WebSub API HMAC secrets", - slog.String("artifact_id", artifactID), - slog.String("url", secretsURL), - ) - - req, err := http.NewRequest("GET", secretsURL, nil) - if err != nil { - return nil, fmt.Errorf("failed to create HMAC secrets request: %w", err) - } - req.Header.Add("api-key", s.config.Token) - req.Header.Add("Accept", "application/json") - - resp, err := s.client.Do(req) - if err != nil { - return nil, fmt.Errorf("failed to fetch HMAC secrets: %w", err) - } - defer resp.Body.Close() - - if resp.StatusCode != http.StatusOK { - bodyBytes, _ := io.ReadAll(resp.Body) - return nil, fmt.Errorf("HMAC secrets request failed with status %d: %s", resp.StatusCode, string(bodyBytes)) - } - - var response platformHmacSecretsResponse - if err := json.NewDecoder(resp.Body).Decode(&response); err != nil { - return nil, fmt.Errorf("failed to decode HMAC secrets response: %w", err) - } - - secrets := make([]HmacSecretInfo, 0, len(response.Secrets)) - for _, s := range response.Secrets { - secrets = append(secrets, HmacSecretInfo{Name: s.Name, Plaintext: s.Secret}) - } - - s.logger.Debug("Successfully fetched WebSub API HMAC secrets", - slog.String("artifact_id", artifactID), - slog.Int("count", len(secrets)), - ) - - return secrets, nil -} +// Note: WebSub HMAC secret fetching (platformHmacSecretInfo, +// platformHmacSecretsResponse, HmacSecretInfo, FetchWebSubAPIHmacSecrets) is +// NOT defined here. It is event-gateway-specific and owned by the +// event-gateway-controller module (event-gateway/gateway-controller/pkg/controlplanehooks), +// built on top of FetchResourceJSON above. // CheckArtifactsExist checks which artifact UUIDs still exist on the platform. // Returns the subset of provided UUIDs that exist. Used during sync to avoid diff --git a/gateway/gateway-controller/tests/integration/schema_test.go b/gateway/gateway-controller/tests/integration/schema_test.go index c427a4935..7d22caef6 100644 --- a/gateway/gateway-controller/tests/integration/schema_test.go +++ b/gateway/gateway-controller/tests/integration/schema_test.go @@ -160,7 +160,7 @@ func TestSchemaInitialization(t *testing.T) { // Verify per-resource-type tables exist t.Run("ResourceTypeTablesExist", func(t *testing.T) { - tables := []string{"rest_apis", "websub_apis", "llm_providers", "llm_proxies", "mcp_proxies"} + tables := []string{"rest_apis", "llm_providers", "llm_proxies", "mcp_proxies"} for _, table := range tables { var tableName string err := rawDB.QueryRow("SELECT name FROM sqlite_master WHERE type='table' AND name=?", table).Scan(&tableName) diff --git a/go.work.sum b/go.work.sum index 745aa0f62..a118b0e2e 100644 --- a/go.work.sum +++ b/go.work.sum @@ -2262,6 +2262,7 @@ github.com/chzyer/test v0.0.0-20210722231415-061457976a23 h1:dZ0/VyGgQdVGAss6Ju0 github.com/chzyer/test v0.0.0-20210722231415-061457976a23/go.mod h1:Q3SI9o4m/ZMnBNeIyt5eFwwo7qiLfzFZmjNmxjkiQlU= github.com/chzyer/test v1.0.0 h1:p3BQDXSxOhOG0P9z6/hGnII4LGiEPOYBhs8asl/fC04= github.com/chzyer/test v1.0.0/go.mod h1:2JlltgoNkt4TW/z9V/IzDdFaMTM2JPIi26O1pF38GC8= +github.com/cilium/ebpf v0.16.0 h1:+BiEnHL6Z7lXnlGUsXQPPAE7+kenAd4ES8MQ5min0Ok= github.com/cilium/ebpf v0.16.0/go.mod h1:L7u2Blt2jMM/vLAVgjxluxtBKlz3/GWjB0dMOEngfwE= github.com/clbanning/x2j v0.0.0-20191024224557-825249438eec/go.mod h1:jMjuTZXRI4dUb/I5gc9Hdhagfvm9+RyrPryS/auMzxE= github.com/client9/misspell v0.3.4 h1:ta993UF76GwbvJcIo3Y68y/M3WxlpEHPWIGDkJYwzJI= @@ -2606,6 +2607,7 @@ github.com/godbus/dbus v0.0.0-20190726142602-4481cbc300e2 h1:ZpnhV/YsD2/4cESfV5+ github.com/godbus/dbus v0.0.0-20190726142602-4481cbc300e2/go.mod h1:bBOAhwG1umN6/6ZUMtDFBMQR8jRg9O75tm9K00oMsK4= github.com/godbus/dbus/v5 v5.0.3/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= github.com/godbus/dbus/v5 v5.0.6/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= +github.com/godbus/dbus/v5 v5.1.0 h1:4KLkAxT3aOY8Li4FRJe/KvhoNFFxo0m6fNuFUO8QJUk= github.com/godbus/dbus/v5 v5.1.0/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= github.com/godror/godror v0.40.4 h1:X1e7hUd02GDaLWKZj40Z7L0CP0W9TrGgmPQZw6+anBg= github.com/godror/godror v0.40.4/go.mod h1:i8YtVTHUJKfFT3wTat4A9UoqScUtZXiYB9Rf3SVARgc=