diff --git a/pkg/connector/client/helpers.go b/pkg/connector/client/helpers.go index 23d0b4db..2f305e88 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/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/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) @@ -38,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 } @@ -46,6 +56,9 @@ func WrapError(err error, contextMsg string) error { // for rate limit errors var slackLibrateLimitErr *slack.RateLimitedError if errors.As(err, &slackLibrateLimitErr) { + if annos != nil { + annos.WithRateLimiting(rateLimitDescription(slackLibrateLimitErr.RetryAfter)) + } return uhttp.WrapErrors(codes.Unavailable, contextMsg, err) } @@ -67,6 +80,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 && annos != nil { + annos.WithRateLimiting(rateLimitDescription(defaultRateLimitRetryAfter)) + } return uhttp.WrapErrors(grpcCode, contextMsg, err) } @@ -162,6 +180,14 @@ func MapSlackErrorToGRPCCode(slackError string) codes.Code { return codes.Unknown } +func rateLimitDescription(retryAfter time.Duration) *v2.RateLimitDescription { + return &v2.RateLimitDescription{ + Status: v2.RateLimitDescription_STATUS_OVERLIMIT, + Remaining: 0, + ResetAt: timestamppb.New(time.Now().Add(retryAfter)), + } +} + // Slack API may return errors in the response body even when the HTTP status code is 200. // // examples: 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{