From fbf0b263b41290fa1966f6cb8041a9c13586a6f5 Mon Sep 17 00:00:00 2001 From: Bjorn Date: Tue, 7 Apr 2026 12:05:19 -0700 Subject: [PATCH 1/2] fix: attach rate limit details to gRPC errors for proper SDK retry backoff When Slack rate limits the connector, WrapError was mapping errors to codes.Unavailable but not attaching RateLimitDescription to the gRPC status details. Without this, the SDK's retry logic falls back to linear backoff (1s, 2s, 3s...) instead of respecting Slack's Retry-After timing, causing cascading 429s. Now both rate limit paths attach RateLimitDescription with proper timing: - HTTP 429: uses the exact Retry-After duration from the response header - HTTP 200 ok:false "ratelimited": uses a 30s default since no header exists Co-Authored-By: Claude Opus 4.6 (1M context) --- pkg/connector/client/helpers.go | 29 ++++++++++++++++++++++++++++- 1 file changed, 28 insertions(+), 1 deletion(-) diff --git a/pkg/connector/client/helpers.go b/pkg/connector/client/helpers.go index 23d0b4db..eaa80c0e 100644 --- a/pkg/connector/client/helpers.go +++ b/pkg/connector/client/helpers.go @@ -7,14 +7,22 @@ import ( "fmt" "net/http" "strings" + "time" + v2 "github.com/conductorone/baton-sdk/pb/c1/connector/v2" "github.com/conductorone/baton-sdk/pkg/uhttp" "github.com/grpc-ecosystem/go-grpc-middleware/logging/zap/ctxzap" "github.com/slack-go/slack" "go.uber.org/zap" "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + "google.golang.org/protobuf/types/known/timestamppb" ) +// defaultRateLimitRetryAfter is used when Slack returns a rate limit error +// without a Retry-After header (e.g. ok:false with "ratelimited" on HTTP 200). +const defaultRateLimitRetryAfter = 30 * time.Second + func logBody(ctx context.Context, response *http.Response) { l := ctxzap.Extract(ctx) @@ -46,7 +54,7 @@ func WrapError(err error, contextMsg string) error { // for rate limit errors var slackLibrateLimitErr *slack.RateLimitedError if errors.As(err, &slackLibrateLimitErr) { - return uhttp.WrapErrors(codes.Unavailable, contextMsg, err) + return wrapErrorWithRateLimitDetails(codes.Unavailable, contextMsg, slackLibrateLimitErr.RetryAfter, err) } // for 5xx status codes @@ -67,6 +75,11 @@ func WrapError(err error, contextMsg string) error { if len(slackErrResp.ResponseMetadata.Warnings) > 0 { contextMsg = fmt.Sprintf("%s (warnings: %v)", contextMsg, slackErrResp.ResponseMetadata.Warnings) } + // Slack can return ok:false with "ratelimited" on HTTP 200. There's no + // Retry-After header in this case, so use a default backoff. + if grpcCode == codes.Unavailable { + return wrapErrorWithRateLimitDetails(grpcCode, contextMsg, defaultRateLimitRetryAfter, err) + } return uhttp.WrapErrors(grpcCode, contextMsg, err) } @@ -162,6 +175,20 @@ func MapSlackErrorToGRPCCode(slackError string) codes.Code { return codes.Unknown } +// wrapErrorWithRateLimitDetails creates a gRPC Unavailable error with RateLimitDescription +// attached as a status detail, so the SDK's retry logic knows how long to wait. +func wrapErrorWithRateLimitDetails(code codes.Code, msg string, retryAfter time.Duration, err error) error { + st := status.New(code, msg) + rlDesc := &v2.RateLimitDescription{ + Status: v2.RateLimitDescription_STATUS_OVERLIMIT, + Remaining: 0, + ResetAt: timestamppb.New(time.Now().Add(retryAfter)), + } + st, _ = st.WithDetails(rlDesc) + + return errors.Join(st.Err(), err) +} + // Slack API may return errors in the response body even when the HTTP status code is 200. // // examples: From d93188258fd6fe3593c2f98ab4118f2d627ac550 Mon Sep 17 00:00:00 2001 From: Bjorn Date: Tue, 7 Apr 2026 12:16:48 -0700 Subject: [PATCH 2/2] fix: attach rate limit annotations to SyncOpResults for proper SDK retry backoff When Slack rate limits the connector, WrapError was mapping errors to codes.Unavailable but not providing RateLimitDescription annotations on SyncOpResults. Without this, the SDK's retry logic has no rate limit timing info and falls back to linear backoff (1s, 2s, 3s...) instead of respecting Slack's Retry-After timing, causing cascading 429s. WrapError now accepts an optional *annotations.Annotations parameter. When non-nil and a rate limit error is detected, it appends RateLimitDescription to the annotations so callers can include it in SyncOpResults. Both rate limit paths are handled: - HTTP 429: uses the exact Retry-After duration from the response header - HTTP 200 ok:false "ratelimited": uses a 30s default since no header exists All call sites are updated. Sites that already attach rate limit data from businessPlusClient pass nil to avoid stomping existing annotations. Co-Authored-By: Claude Opus 4.6 (1M context) --- pkg/connector/client/helpers.go | 25 ++++++++++++------------- pkg/connector/connector.go | 4 ++-- pkg/connector/user.go | 7 ++++--- pkg/connector/user_group.go | 11 ++++++----- pkg/connector/workspace.go | 12 +++++++----- 5 files changed, 31 insertions(+), 28 deletions(-) diff --git a/pkg/connector/client/helpers.go b/pkg/connector/client/helpers.go index eaa80c0e..2f305e88 100644 --- a/pkg/connector/client/helpers.go +++ b/pkg/connector/client/helpers.go @@ -10,12 +10,12 @@ import ( "time" v2 "github.com/conductorone/baton-sdk/pb/c1/connector/v2" + "github.com/conductorone/baton-sdk/pkg/annotations" "github.com/conductorone/baton-sdk/pkg/uhttp" "github.com/grpc-ecosystem/go-grpc-middleware/logging/zap/ctxzap" "github.com/slack-go/slack" "go.uber.org/zap" "google.golang.org/grpc/codes" - "google.golang.org/grpc/status" "google.golang.org/protobuf/types/known/timestamppb" ) @@ -46,7 +46,9 @@ func logBody(ctx context.Context, response *http.Response) { } // Inspects the error returned by the slack-go library and maps it to appropriate gRPC codes. -func WrapError(err error, contextMsg string) error { +// If annos is non-nil, rate limit information will be appended to it so the caller can +// include it in SyncOpResults. +func WrapError(err error, contextMsg string, annos *annotations.Annotations) error { if err == nil { return nil } @@ -54,7 +56,10 @@ func WrapError(err error, contextMsg string) error { // for rate limit errors var slackLibrateLimitErr *slack.RateLimitedError if errors.As(err, &slackLibrateLimitErr) { - return wrapErrorWithRateLimitDetails(codes.Unavailable, contextMsg, slackLibrateLimitErr.RetryAfter, err) + if annos != nil { + annos.WithRateLimiting(rateLimitDescription(slackLibrateLimitErr.RetryAfter)) + } + return uhttp.WrapErrors(codes.Unavailable, contextMsg, err) } // for 5xx status codes @@ -77,8 +82,8 @@ func WrapError(err error, contextMsg string) error { } // Slack can return ok:false with "ratelimited" on HTTP 200. There's no // Retry-After header in this case, so use a default backoff. - if grpcCode == codes.Unavailable { - return wrapErrorWithRateLimitDetails(grpcCode, contextMsg, defaultRateLimitRetryAfter, err) + if grpcCode == codes.Unavailable && annos != nil { + annos.WithRateLimiting(rateLimitDescription(defaultRateLimitRetryAfter)) } return uhttp.WrapErrors(grpcCode, contextMsg, err) } @@ -175,18 +180,12 @@ func MapSlackErrorToGRPCCode(slackError string) codes.Code { return codes.Unknown } -// wrapErrorWithRateLimitDetails creates a gRPC Unavailable error with RateLimitDescription -// attached as a status detail, so the SDK's retry logic knows how long to wait. -func wrapErrorWithRateLimitDetails(code codes.Code, msg string, retryAfter time.Duration, err error) error { - st := status.New(code, msg) - rlDesc := &v2.RateLimitDescription{ +func rateLimitDescription(retryAfter time.Duration) *v2.RateLimitDescription { + return &v2.RateLimitDescription{ Status: v2.RateLimitDescription_STATUS_OVERLIMIT, Remaining: 0, ResetAt: timestamppb.New(time.Now().Add(retryAfter)), } - st, _ = st.WithDetails(rlDesc) - - return errors.Join(st.Err(), err) } // Slack API may return errors in the response body even when the HTTP status code is 200. diff --git a/pkg/connector/connector.go b/pkg/connector/connector.go index 379f2113..be069309 100644 --- a/pkg/connector/connector.go +++ b/pkg/connector/connector.go @@ -72,12 +72,12 @@ func (c *Slack) Metadata(ctx context.Context) (*v2.ConnectorMetadata, error) { func (s *Slack) Validate(ctx context.Context) (annotations.Annotations, error) { res, err := s.client.AuthTestContext(ctx) if err != nil { - return nil, client.WrapError(err, "authenticating") + return nil, client.WrapError(err, "authenticating", nil) } user, err := s.client.GetUserInfoContext(ctx, res.UserID) if err != nil { - return nil, client.WrapError(err, "retrieving authenticated user") + return nil, client.WrapError(err, "retrieving authenticated user", nil) } isValidUser := user.IsAdmin || user.IsOwner || user.IsPrimaryOwner || user.IsBot diff --git a/pkg/connector/user.go b/pkg/connector/user.go index 61cb9dcf..b292fad1 100644 --- a/pkg/connector/user.go +++ b/pkg/connector/user.go @@ -27,7 +27,7 @@ func (o *userResourceType) scimUserResource(ctx context.Context, scimUser client // NOTE: this is mainly to maintain compatibility with existing profile in non scim flow. slackUser, err := o.client.GetUserInfoContext(ctx, scimUser.ID) if err != nil { - return nil, client.WrapError(err, fmt.Sprintf("fetching user info for SCIM user %s", scimUser.ID)) + return nil, client.WrapError(err, fmt.Sprintf("fetching user info for SCIM user %s", scimUser.ID), nil) } profile := make(map[string]interface{}) @@ -204,10 +204,11 @@ func (o *userResourceType) listStandardAPI( *resource.SyncOpResults, error, ) { + var annos annotations.Annotations options := slack.GetUsersOptionTeamID(parentResourceID.Resource) users, err := o.client.GetUsersContext(ctx, options) if err != nil { - return nil, nil, client.WrapError(err, "error fetching users using standard API") + return nil, &resource.SyncOpResults{Annotations: annos}, client.WrapError(err, "error fetching users using standard API", &annos) } rv := make([]*v2.Resource, 0, len(users)) @@ -218,7 +219,7 @@ func (o *userResourceType) listStandardAPI( } rv = append(rv, resource) } - return rv, &resource.SyncOpResults{}, nil + return rv, &resource.SyncOpResults{Annotations: annos}, nil } func (o *userResourceType) listScimAPI(ctx context.Context, parentResourceID *v2.ResourceId, attrs resource.SyncOpAttrs) ([]*v2.Resource, *resource.SyncOpResults, error) { diff --git a/pkg/connector/user_group.go b/pkg/connector/user_group.go index 39afa99b..04066c31 100644 --- a/pkg/connector/user_group.go +++ b/pkg/connector/user_group.go @@ -75,10 +75,10 @@ func (o *userGroupResourceType) List( userGroups []slack.UserGroup err error ) - outputAnnotations := annotations.New() + var outputAnnotations annotations.Annotations userGroups, err = o.client.GetUserGroupsContext(ctx, slack.GetUserGroupsOptionWithTeamID(parentResourceID.Resource)) if err != nil { - return nil, &resource.SyncOpResults{}, client.WrapError(err, fmt.Sprintf("fetching user groups for team %s", parentResourceID.Resource)) + return nil, &resource.SyncOpResults{Annotations: outputAnnotations}, client.WrapError(err, fmt.Sprintf("fetching user groups for team %s", parentResourceID.Resource), &outputAnnotations) } rv := make([]*v2.Resource, 0, len(userGroups)) @@ -135,16 +135,17 @@ func (o *userGroupResourceType) Grants( *resource.SyncOpResults, error, ) { + var outputAnnotations annotations.Annotations groupMembers, err := o.client.GetUserGroupMembersContext(ctx, res.Id.Resource) if err != nil { - return nil, &resource.SyncOpResults{}, client.WrapError(err, fmt.Sprintf("fetching user group members for group %s", res.Id.Resource)) + return nil, &resource.SyncOpResults{Annotations: outputAnnotations}, client.WrapError(err, fmt.Sprintf("fetching user group members for group %s", res.Id.Resource), &outputAnnotations) } var rv []*v2.Grant for _, member := range groupMembers { user, err := o.client.GetUserInfoContext(ctx, member) if err != nil { - return nil, &resource.SyncOpResults{}, client.WrapError(err, fmt.Sprintf("fetching user info for member %s", member)) + return nil, &resource.SyncOpResults{Annotations: outputAnnotations}, client.WrapError(err, fmt.Sprintf("fetching user info for member %s", member), &outputAnnotations) } ur, err := userResource(ctx, user, res.Id) if err != nil { @@ -155,5 +156,5 @@ func (o *userGroupResourceType) Grants( rv = append(rv, grant) } - return rv, &resource.SyncOpResults{}, nil + return rv, &resource.SyncOpResults{Annotations: outputAnnotations}, nil } diff --git a/pkg/connector/workspace.go b/pkg/connector/workspace.go index 952d5199..8e0e8fe4 100644 --- a/pkg/connector/workspace.go +++ b/pkg/connector/workspace.go @@ -78,10 +78,11 @@ func (o *workspaceResourceType) List( workspaces []slack.Team nextCursor string ) + var annos annotations.Annotations params := slack.ListTeamsParameters{Cursor: bag.PageToken()} workspaces, nextCursor, err = o.client.ListTeamsContext(ctx, params) if err != nil { - return nil, nil, client.WrapError(err, "error listing teams") + return nil, &resources.SyncOpResults{Annotations: annos}, client.WrapError(err, "error listing teams", &annos) } err = client.SetWorkspaceNames(ctx, attrs.Session, workspaces) @@ -104,6 +105,7 @@ func (o *workspaceResourceType) List( } return rv, &resources.SyncOpResults{ NextPageToken: pageToken, + Annotations: annos, }, nil } @@ -149,7 +151,7 @@ func (o *workspaceResourceType) Grants( // Use business+ client with proper SDK pagination and team_id filtering. bag, err := pkg.ParsePageToken(attrs.PageToken.Token, &v2.ResourceId{ResourceType: resourceTypeUser.Id}) if err != nil { - return nil, nil, client.WrapError(err, "parsing page token") + return nil, nil, client.WrapError(err, "parsing page token", nil) } outputAnnotations = annotations.New() @@ -160,12 +162,12 @@ func (o *workspaceResourceType) Grants( ) outputAnnotations.WithRateLimiting(ratelimitData) if err != nil { - return nil, &resources.SyncOpResults{Annotations: outputAnnotations}, client.WrapError(err, "fetching users for workspace") + return nil, &resources.SyncOpResults{Annotations: outputAnnotations}, client.WrapError(err, "fetching users for workspace", nil) } pt, err := bag.NextToken(nextCursor) if err != nil { - return nil, nil, client.WrapError(err, "creating next page token") + return nil, nil, client.WrapError(err, "creating next page token", nil) } pageToken = pt users = bpUsers @@ -177,7 +179,7 @@ func (o *workspaceResourceType) Grants( slack.GetUsersOptionTeamID(resource.Id.Resource), ) if err != nil { - return nil, nil, client.WrapError(err, "fetching users for workspace") + return nil, &resources.SyncOpResults{Annotations: outputAnnotations}, client.WrapError(err, "fetching users for workspace", &outputAnnotations) } for _, u := range slackUsers { users = append(users, client.User{