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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 20 additions & 0 deletions backend/pkg/api/connect/service/kafkaconnect/v1/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,15 @@ import (
"log/slog"
"net/http"
"sort"
"strconv"

"connectrpc.com/connect"
"github.com/cloudhut/common/rest"
con "github.com/cloudhut/connect-client"
"github.com/redpanda-data/common-go/api/pagination"
"google.golang.org/genproto/googleapis/rpc/errdetails"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/emptypb"

apierrors "github.com/redpanda-data/console/backend/pkg/api/connect/errors"
Expand Down Expand Up @@ -362,6 +366,22 @@ func (*Service) matchError(err *rest.Error) *connect.Error {
v1.Reason_REASON_KAFKA_CONNECT_API_ERROR.String(),
),
)
case http.StatusServiceUnavailable:
// Kafka Connect behind connect-gate answers 503 while its workers are
// scaled to zero and booting; that is retryable, and the gate's retry
// hint travels as RetryInfo plus ErrorInfo metadata.
errInfo := apierrors.NewErrorInfo(v1.Reason_REASON_KAFKA_CONNECT_API_ERROR.String())
details := []proto.Message{}
if se, ok := kafkaconnect.AsStartingError(err); ok {
errInfo = apierrors.NewErrorInfo(
v1.Reason_REASON_KAFKA_CONNECT_API_ERROR.String(),
apierrors.KeyVal{Key: "reason", Value: kafkaconnect.StartingReason},
apierrors.KeyVal{Key: "phase", Value: se.Phase},
apierrors.KeyVal{Key: "retry_after_seconds", Value: strconv.Itoa(se.RetryAfterSeconds())},
)
details = append(details, &errdetails.RetryInfo{RetryDelay: durationpb.New(se.RetryAfter)})
}
return apierrors.NewConnectError(connect.CodeUnavailable, err.Err, errInfo, details...)
default:
return apierrors.NewConnectError(
connect.CodeInternal,
Expand Down
20 changes: 20 additions & 0 deletions backend/pkg/api/connect/service/kafkaconnect/v1alpha2/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,15 @@ import (
"log/slog"
"net/http"
"sort"
"strconv"

"connectrpc.com/connect"
"github.com/cloudhut/common/rest"
con "github.com/cloudhut/connect-client"
"github.com/redpanda-data/common-go/api/pagination"
"google.golang.org/genproto/googleapis/rpc/errdetails"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/emptypb"

apierrors "github.com/redpanda-data/console/backend/pkg/api/connect/errors"
Expand Down Expand Up @@ -362,6 +366,22 @@ func (*Service) matchError(err *rest.Error) *connect.Error {
v1alpha2.Reason_REASON_KAFKA_CONNECT_API_ERROR.String(),
),
)
case http.StatusServiceUnavailable:
// Kafka Connect behind connect-gate answers 503 while its workers are
// scaled to zero and booting; that is retryable, and the gate's retry
// hint travels as RetryInfo plus ErrorInfo metadata.
errInfo := apierrors.NewErrorInfo(v1alpha2.Reason_REASON_KAFKA_CONNECT_API_ERROR.String())
details := []proto.Message{}
if se, ok := kafkaconnect.AsStartingError(err); ok {
errInfo = apierrors.NewErrorInfo(
v1alpha2.Reason_REASON_KAFKA_CONNECT_API_ERROR.String(),
apierrors.KeyVal{Key: "reason", Value: kafkaconnect.StartingReason},
apierrors.KeyVal{Key: "phase", Value: se.Phase},
apierrors.KeyVal{Key: "retry_after_seconds", Value: strconv.Itoa(se.RetryAfterSeconds())},
)
details = append(details, &errdetails.RetryInfo{RetryDelay: durationpb.New(se.RetryAfter)})
}
return apierrors.NewConnectError(connect.CodeUnavailable, err.Err, errInfo, details...)
default:
return apierrors.NewConnectError(
connect.CodeInternal,
Expand Down
24 changes: 12 additions & 12 deletions backend/pkg/api/handle_kafka_connect.go
Original file line number Diff line number Diff line change
Expand Up @@ -133,20 +133,20 @@ func (api *API) handlePutConnectorConfig() http.HandlerFunc {
var req putConnectorConfigRequest
restErr := rest.Decode(w, r, &req)
if restErr != nil {
rest.SendRESTError(w, r, api.Logger, restErr)
api.sendKafkaConnectError(w, r, restErr)
return
}

cInfo, restErr := api.ConnectSvc.PutConnectorConfig(r.Context(), clusterName, connectorName, req.ToClientRequest())
if restErr != nil {
rest.SendRESTError(w, r, api.Logger, restErr)
api.sendKafkaConnectError(w, r, restErr)
return
}

// restart the instance and all the tasks
restErr = api.ConnectSvc.RestartConnector(r.Context(), clusterName, connectorName, true, false)
if restErr != nil {
rest.SendRESTError(w, r, api.Logger, restErr)
api.sendKafkaConnectError(w, r, restErr)
return
}

Expand All @@ -162,13 +162,13 @@ func (api *API) handlePutValidateConnectorConfig() http.HandlerFunc {
var req map[string]any
restErr := rest.Decode(w, r, &req)
if restErr != nil {
rest.SendRESTError(w, r, api.Logger, restErr)
api.sendKafkaConnectError(w, r, restErr)
return
}

cInfo, restErr := api.ConnectSvc.ValidateConnectorConfig(r.Context(), clusterName, pluginClassName, req)
if restErr != nil {
rest.SendRESTError(w, r, api.Logger, restErr)
api.sendKafkaConnectError(w, r, restErr)
return
}
rest.SendResponse(w, r, api.Logger, http.StatusOK, cInfo)
Expand Down Expand Up @@ -199,13 +199,13 @@ func (api *API) handleCreateConnector() http.HandlerFunc {
var req createConnectorRequest
restErr := rest.Decode(w, r, &req)
if restErr != nil {
rest.SendRESTError(w, r, api.Logger, restErr)
api.sendKafkaConnectError(w, r, restErr)
return
}

cInfo, restErr := api.ConnectSvc.CreateConnector(r.Context(), clusterName, req.ToClientRequest())
if restErr != nil {
rest.SendRESTError(w, r, api.Logger, restErr)
api.sendKafkaConnectError(w, r, restErr)
return
}

Expand All @@ -223,7 +223,7 @@ func (api *API) handleDeleteConnector() http.HandlerFunc {

restErr := api.ConnectSvc.DeleteConnector(ctx, clusterName, connector)
if restErr != nil {
rest.SendRESTError(w, r, api.Logger, restErr)
api.sendKafkaConnectError(w, r, restErr)
return
}

Expand All @@ -241,7 +241,7 @@ func (api *API) handlePauseConnector() http.HandlerFunc {

restErr := api.ConnectSvc.PauseConnector(ctx, clusterName, connector)
if restErr != nil {
rest.SendRESTError(w, r, api.Logger, restErr)
api.sendKafkaConnectError(w, r, restErr)
return
}

Expand All @@ -259,7 +259,7 @@ func (api *API) handleResumeConnector() http.HandlerFunc {

restErr := api.ConnectSvc.ResumeConnector(ctx, clusterName, connector)
if restErr != nil {
rest.SendRESTError(w, r, api.Logger, restErr)
api.sendKafkaConnectError(w, r, restErr)
return
}

Expand All @@ -277,7 +277,7 @@ func (api *API) handleRestartConnector() http.HandlerFunc {

restErr := api.ConnectSvc.RestartConnector(ctx, clusterName, connector, true, false)
if restErr != nil {
rest.SendRESTError(w, r, api.Logger, restErr)
api.sendKafkaConnectError(w, r, restErr)
return
}

Expand Down Expand Up @@ -306,7 +306,7 @@ func (api *API) handleRestartConnectorTask() http.HandlerFunc {

restErr := api.ConnectSvc.RestartConnectorTask(ctx, clusterName, connector, taskID)
if restErr != nil {
rest.SendRESTError(w, r, api.Logger, restErr)
api.sendKafkaConnectError(w, r, restErr)
return
}

Expand Down
51 changes: 51 additions & 0 deletions backend/pkg/api/handle_kafka_connect_starting.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
// Copyright 2026 Redpanda Data, Inc.
//
// Use of this software is governed by the Business Source License
// included in the file https://github.com/redpanda-data/redpanda/blob/dev/licenses/bsl.md
//
// As of the Change Date specified in that file, in accordance with
// the Business Source License, use of this software will be governed
// by the Apache License, Version 2.0

package api

import (
"net/http"
"strconv"

"github.com/cloudhut/common/rest"

"github.com/redpanda-data/console/backend/pkg/connect"
)

// kafkaConnectStartingResponse is the body sent while the Kafka Connect workers
// are scaled to zero and starting. It extends the usual rest.Error body so
// existing clients keep working and newer ones can retry on their own.
type kafkaConnectStartingResponse struct {
StatusCode int `json:"statusCode"`
Message string `json:"message"`
Reason string `json:"reason"`
Phase string `json:"phase"`
RetryAfterSeconds int `json:"retryAfterSeconds"`
EstimatedWaitSeconds int `json:"estimatedWaitSeconds"`
}

// sendKafkaConnectError writes a Kafka Connect service error. A "starting" error
// gets a Retry-After header and the retry hints in the body; anything else is
// sent as a plain rest.Error.
func (api *API) sendKafkaConnectError(w http.ResponseWriter, r *http.Request, restErr *rest.Error) {
se, ok := connect.AsStartingError(restErr)
if !ok {
rest.SendRESTError(w, r, api.Logger, restErr)
return
}
w.Header().Set("Retry-After", strconv.Itoa(se.RetryAfterSeconds()))
rest.SendResponse(w, r, api.Logger, http.StatusServiceUnavailable, kafkaConnectStartingResponse{
StatusCode: http.StatusServiceUnavailable,
Message: se.Message,
Reason: connect.StartingReason,
Phase: se.Phase,
RetryAfterSeconds: se.RetryAfterSeconds(),
EstimatedWaitSeconds: int(se.EstimatedWait.Round(1e9).Seconds()),
})
}
4 changes: 4 additions & 0 deletions backend/pkg/connect/create_connector.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ func (s *Service) CreateConnector(ctx context.Context, clusterName string, req c
}
req.Config = s.Interceptor.ConsoleToKafkaConnect(className, req.Config)

ctx, starting := withStartingCapture(ctx)
cInfo, err := c.Client.CreateConnector(ctx, req)
connectorClass := getMapValueOrString(cInfo.Config, "connector.class", "unknown")
cInfo = con.ConnectorInfo{
Expand All @@ -48,6 +49,9 @@ func (s *Service) CreateConnector(ctx context.Context, clusterName string, req c
}

if err != nil {
if starting.err != nil {
return con.ConnectorInfo{}, startingRestError(starting.err, "create connector", slog.String("cluster_name", clusterName))
}
return con.ConnectorInfo{}, &rest.Error{
Err: fmt.Errorf("failed to create connector: %w", err),
Status: GetStatusCodeFromAPIError(err, http.StatusInternalServerError),
Expand Down
4 changes: 4 additions & 0 deletions backend/pkg/connect/delete_connector.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,12 @@ func (s *Service) DeleteConnector(ctx context.Context, clusterName string, conne
return restErr
}

ctx, starting := withStartingCapture(ctx)
err := c.Client.DeleteConnector(ctx, connector)
if err != nil {
if starting.err != nil {
return startingRestError(starting.err, "delete connector", slog.String("cluster_name", clusterName), slog.String("connector", connector))
}
return &rest.Error{
Err: err,
Status: GetStatusCodeFromAPIError(err, http.StatusServiceUnavailable),
Expand Down
4 changes: 4 additions & 0 deletions backend/pkg/connect/pause_connector.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,12 @@ func (s *Service) PauseConnector(ctx context.Context, clusterName string, connec
return restErr
}

ctx, starting := withStartingCapture(ctx)
err := c.Client.PauseConnector(ctx, connector)
if err != nil {
if starting.err != nil {
return startingRestError(starting.err, "pause connector", slog.String("cluster_name", clusterName), slog.String("connector", connector))
}
return &rest.Error{
Err: err,
Status: GetStatusCodeFromAPIError(err, http.StatusInternalServerError),
Expand Down
4 changes: 4 additions & 0 deletions backend/pkg/connect/put_connector_config.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ func (s *Service) PutConnectorConfig(ctx context.Context, clusterName string, co
}
req.Config = s.Interceptor.ConsoleToKafkaConnect(className, req.Config)

ctx, starting := withStartingCapture(ctx)
cInfo, err := c.Client.PutConnectorConfig(ctx, connectorName, req)
connectorClass := getMapValueOrString(cInfo.Config, "connector.class", "unknown")
cInfo = con.ConnectorInfo{
Expand All @@ -48,6 +49,9 @@ func (s *Service) PutConnectorConfig(ctx context.Context, clusterName string, co
}

if err != nil {
if starting.err != nil {
return con.ConnectorInfo{}, startingRestError(starting.err, "update connector config", slog.String("cluster_name", clusterName), slog.String("connector_name", connectorName))
}
return con.ConnectorInfo{}, &rest.Error{
Err: fmt.Errorf("failed to patch connector config: %w", err),
Status: GetStatusCodeFromAPIError(err, http.StatusInternalServerError),
Expand Down
4 changes: 4 additions & 0 deletions backend/pkg/connect/restart_connector.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,12 @@ func (s *Service) RestartConnector(ctx context.Context, clusterName string, conn
return restErr
}

ctx, starting := withStartingCapture(ctx)
err := c.Client.RestartConnector(ctx, connector, connect.RestartConnectorOptions{IncludeTasks: restartTasks, OnlyFailed: restartOnlyFailed})
if err != nil {
if starting.err != nil {
return startingRestError(starting.err, "restart connector", slog.String("cluster_name", clusterName), slog.String("connector", connector))
}
return &rest.Error{
Err: err,
Status: GetStatusCodeFromAPIError(err, http.StatusServiceUnavailable),
Expand Down
4 changes: 4 additions & 0 deletions backend/pkg/connect/restart_task.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,12 @@ func (s *Service) RestartConnectorTask(ctx context.Context, clusterName string,
return restErr
}

ctx, starting := withStartingCapture(ctx)
err := c.Client.RestartConnectorTask(ctx, connector, taskID)
if err != nil {
if starting.err != nil {
return startingRestError(starting.err, "restart connector task", slog.String("cluster_name", clusterName), slog.String("connector", connector), slog.Int("task_id", taskID))
}
return &rest.Error{
Err: err,
Status: http.StatusServiceUnavailable,
Expand Down
6 changes: 5 additions & 1 deletion backend/pkg/connect/resume_connector.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,12 +26,16 @@ func (s *Service) ResumeConnector(ctx context.Context, clusterName string, conne
return restErr
}

ctx, starting := withStartingCapture(ctx)
err := c.Client.ResumeConnector(ctx, connector)
if err != nil {
if starting.err != nil {
return startingRestError(starting.err, "resume connector", slog.String("cluster_name", clusterName), slog.String("connector", connector))
}
return &rest.Error{
Err: err,
Status: GetStatusCodeFromAPIError(err, http.StatusServiceUnavailable),
Message: fmt.Sprintf("Failed to pause connector: %v", err.Error()),
Message: fmt.Sprintf("Failed to resume connector: %v", err.Error()),
InternalLogs: []slog.Attr{slog.String("cluster_name", clusterName), slog.String("connector", connector)},
IsSilent: false,
}
Expand Down
4 changes: 4 additions & 0 deletions backend/pkg/connect/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,10 @@ func NewService(cfg config.KafkaConnect, logger *slog.Logger) (*Service, error)

// Create client
client := con.NewClient(opts...)
// Capture the details of a connect-gate "starting" 503, which the
// client itself reduces to a code and a message.
hc := client.GetClient()
hc.Transport = newGateTransport(hc.Transport)
clientsByCluster[clusterCfg.Name] = &ClientWithConfig{
Client: client,
Cfg: clusterCfg,
Expand Down
Loading
Loading