Skip to content
Open
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
10 changes: 6 additions & 4 deletions pkg/app/piped/cmd/piped/piped.go
Original file line number Diff line number Diff line change
Expand Up @@ -834,19 +834,21 @@ func (p *piped) sendPipedMeta(ctx context.Context, client pipedservice.Client, c
}
}

for retry.WaitNext(ctx) {
if res, err := client.ReportPipedMeta(ctx, req); err == nil {
_, err = retry.Do(ctx, func() (interface{}, error) {
res, err := client.ReportPipedMeta(ctx, req)
if err == nil {
cfg.Name = res.Name
if cfg.WebAddress == "" {
cfg.WebAddress = res.WebBaseUrl
}
return nil
return nil, nil
}
logger.Warn("failed to report piped meta to control-plane, wait to the next retry",
zap.Int("calls", retry.Calls()),
zap.Error(err),
)
}
return nil, err
})

return err
}
Expand Down
42 changes: 21 additions & 21 deletions pkg/app/piped/controller/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -674,24 +674,24 @@ func (c *controller) startNewScheduler(ctx context.Context, d *model.Deployment)

func (c *controller) getMostRecentlySuccessfulDeployment(ctx context.Context, applicationID string) (*model.ApplicationDeploymentReference, error) {
var (
err error
resp *pipedservice.GetApplicationMostRecentDeploymentResponse
retry = pipedservice.NewRetry(3)
req = &pipedservice.GetApplicationMostRecentDeploymentRequest{
ApplicationId: applicationID,
Status: model.DeploymentStatus_DEPLOYMENT_SUCCESS,
}
)

for retry.WaitNext(ctx) {
if resp, err = c.apiClient.GetApplicationMostRecentDeployment(ctx, req); err == nil {
return resp.Deployment, nil
}
if !pipedservice.Retriable(err) {
return nil, err
d, err := retry.Do(ctx, func() (interface{}, error) {
resp, err := c.apiClient.GetApplicationMostRecentDeployment(ctx, req)
if err != nil {
return nil, pipedservice.NewRetriableErr(err)
}
return resp.Deployment, nil
})
if err != nil {
return nil, err
}
return nil, err
return d.(*model.ApplicationDeploymentReference), nil
}

func (c *controller) shouldStartPlanningDeployment(ctx context.Context, d *model.Deployment) (plannable, cancel bool, cancelReason string, err error) {
Expand All @@ -715,7 +715,6 @@ func (c *controller) shouldStartPlanningDeployment(ctx context.Context, d *model

func (c *controller) cancelDeployment(ctx context.Context, d *model.Deployment, reason string) error {
var (
err error
req = &pipedservice.ReportDeploymentCompletedRequest{
DeploymentId: d.Id,
Status: model.DeploymentStatus_DEPLOYMENT_CANCELLED,
Expand All @@ -728,12 +727,13 @@ func (c *controller) cancelDeployment(ctx context.Context, d *model.Deployment,
retry = pipedservice.NewRetry(10)
)

for retry.WaitNext(ctx) {
if _, err = c.apiClient.ReportDeploymentCompleted(ctx, req); err == nil {
return nil
_, err := retry.Do(ctx, func() (interface{}, error) {
_, err := c.apiClient.ReportDeploymentCompleted(ctx, req)
if err != nil {
return nil, fmt.Errorf("failed to report deployment status to control-plane: %v", err)
}
err = fmt.Errorf("failed to report deployment status to control-plane: %v", err)
}
return nil, nil
})
return err
}

Expand All @@ -749,19 +749,19 @@ func (l appLiveResourceLister) ListKubernetesResources() ([]provider.Manifest, b

func reportApplicationDeployingStatus(ctx context.Context, c apiClient, appID string, deploying bool) error {
var (
err error
retry = pipedservice.NewRetry(10)
req = &pipedservice.ReportApplicationDeployingStatusRequest{
ApplicationId: appID,
Deploying: deploying,
}
)

for retry.WaitNext(ctx) {
if _, err = c.ReportApplicationDeployingStatus(ctx, req); err == nil {
return nil
_, err := retry.Do(ctx, func() (interface{}, error) {
_, err := c.ReportApplicationDeployingStatus(ctx, req)
if err != nil {
return nil, fmt.Errorf("failed to report application deploying status to control-plane: %w", err)
}
err = fmt.Errorf("failed to report application deploying status to control-plane: %w", err)
}
return nil, nil
})
return err
}
33 changes: 18 additions & 15 deletions pkg/app/piped/controller/planner.go
Original file line number Diff line number Diff line change
Expand Up @@ -273,12 +273,13 @@ func (p *planner) reportDeploymentPlanned(ctx context.Context, out pln.Output) e
})
}()

for retry.WaitNext(ctx) {
if _, err = p.apiClient.ReportDeploymentPlanned(ctx, req); err == nil {
return nil
_, err = retry.Do(ctx, func() (interface{}, error) {
_, err := p.apiClient.ReportDeploymentPlanned(ctx, req)
if err != nil {
return nil, fmt.Errorf("failed to report deployment status to control-plane: %v", err)
}
err = fmt.Errorf("failed to report deployment status to control-plane: %v", err)
}
return nil, nil
})

if err != nil {
p.logger.Error("failed to mark deployment to be planned", zap.Error(err))
Expand Down Expand Up @@ -319,12 +320,13 @@ func (p *planner) reportDeploymentFailed(ctx context.Context, reason string) err
})
}()

for retry.WaitNext(ctx) {
if _, err = p.apiClient.ReportDeploymentCompleted(ctx, req); err == nil {
return nil
_, err = retry.Do(ctx, func() (interface{}, error) {
_, err := p.apiClient.ReportDeploymentCompleted(ctx, req)
if err != nil {
return nil, fmt.Errorf("failed to report deployment status to control-plane: %v", err)
}
err = fmt.Errorf("failed to report deployment status to control-plane: %v", err)
}
return nil, nil
})

if err != nil {
p.logger.Error("failed to mark deployment to be failed", zap.Error(err))
Expand Down Expand Up @@ -365,12 +367,13 @@ func (p *planner) reportDeploymentCancelled(ctx context.Context, commander, reas
})
}()

for retry.WaitNext(ctx) {
if _, err = p.apiClient.ReportDeploymentCompleted(ctx, req); err == nil {
return nil
_, err = retry.Do(ctx, func() (interface{}, error) {
_, err := p.apiClient.ReportDeploymentCompleted(ctx, req)
if err != nil {
return nil, fmt.Errorf("failed to report deployment status to control-plane: %v", err)
}
err = fmt.Errorf("failed to report deployment status to control-plane: %v", err)
}
return nil, nil
})

if err != nil {
p.logger.Error("failed to mark deployment to be cancelled", zap.Error(err))
Expand Down
46 changes: 24 additions & 22 deletions pkg/app/piped/controller/scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -678,14 +678,13 @@ func (s *scheduler) reportStageStatus(ctx context.Context, stageID string, statu
// Update stage status at local.
s.stageStatuses[stageID] = status

// Update stage status on the remote.
for retry.WaitNext(ctx) {
_, err = s.apiClient.ReportStageStatusChanged(ctx, req)
if err == nil {
break
_, err = retry.Do(ctx, func() (interface{}, error) {
_, err := s.apiClient.ReportStageStatusChanged(ctx, req)
if err != nil {
return nil, fmt.Errorf("failed to report stage status to control-plane: %v", err)
}
err = fmt.Errorf("failed to report stage status to control-plane: %v", err)
}
return nil, nil
})

return err
}
Expand All @@ -704,12 +703,13 @@ func (s *scheduler) reportDeploymentStatusChanged(ctx context.Context, status mo
)

// Update deployment status on remote.
for retry.WaitNext(ctx) {
if _, err = s.apiClient.ReportDeploymentStatusChanged(ctx, req); err == nil {
return nil
_, err = retry.Do(ctx, func() (interface{}, error) {
_, err := s.apiClient.ReportDeploymentStatusChanged(ctx, req)
if err != nil {
return nil, fmt.Errorf("failed to report deployment status to control-plane: %v", err)
}
err = fmt.Errorf("failed to report deployment status to control-plane: %v", err)
}
return nil, nil
})

return err
}
Expand Down Expand Up @@ -782,12 +782,13 @@ func (s *scheduler) reportDeploymentCompleted(ctx context.Context, status model.
}()

// Update deployment status on remote.
for retry.WaitNext(ctx) {
if _, err = s.apiClient.ReportDeploymentCompleted(ctx, req); err == nil {
return nil
_, err = retry.Do(ctx, func() (interface{}, error) {
_, err := s.apiClient.ReportDeploymentCompleted(ctx, req)
if err != nil {
return nil, fmt.Errorf("failed to report deployment status to control-plane: %w", err)
}
err = fmt.Errorf("failed to report deployment status to control-plane: %w", err)
}
return nil, nil
})

return err
}
Expand Down Expand Up @@ -825,12 +826,13 @@ func (s *scheduler) reportMostRecentlySuccessfulDeployment(ctx context.Context)
retry = pipedservice.NewRetry(10)
)

for retry.WaitNext(ctx) {
if _, err = s.apiClient.ReportApplicationMostRecentDeployment(ctx, req); err == nil {
return nil
_, err = retry.Do(ctx, func() (interface{}, error) {
_, err := s.apiClient.ReportApplicationMostRecentDeployment(ctx, req)
if err != nil {
return nil, fmt.Errorf("failed to report most recent successful deployment: %w", err)
}
err = fmt.Errorf("failed to report most recent successful deployment: %w", err)
}
return nil, nil
})

return err
}
Expand Down
12 changes: 5 additions & 7 deletions pkg/app/piped/executor/lambda/lambda.go
Original file line number Diff line number Diff line change
Expand Up @@ -308,22 +308,20 @@ func build(ctx context.Context, in *executor.Input, client provider.Client, fm p

in.LogPersister.Info("Waiting to update lambda function in progress...")
retry := backoff.NewRetry(provider.RequestRetryTime, backoff.NewConstant(provider.RetryIntervalDuration))
publishFunctionSucceed := false
startWaitingStamp := time.Now()
for retry.WaitNext(ctx) {
_, err = retry.Do(ctx, func() (interface{}, error) {
// Commit version for applied Lambda function.
// Note: via the current docs of [Lambda.PublishVersion](https://docs.aws.amazon.com/sdk-for-go/api/service/lambda/#Lambda.PublishVersion)
// AWS Lambda doesn't publish a version if the function's configuration and code haven't changed since the last version.
// But currently, unchanged revision is able to make publish (versionId++) as usual.
version, err = client.PublishFunction(ctx, fm)
if err != nil {
in.Logger.Error("Failed publish new version for Lambda function")
} else {
publishFunctionSucceed = true
break
return nil, err
}
}
if !publishFunctionSucceed {
return nil, nil
})
if err != nil {
in.LogPersister.Errorf("Failed to commit new version for Lambda function %s: %v", fm.Spec.Name, err)
return
}
Expand Down
17 changes: 7 additions & 10 deletions pkg/app/piped/platformprovider/lambda/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -284,9 +284,7 @@ func (c *client) UpdateFunctionFromSource(ctx context.Context, fm FunctionManife

func (c *client) updateFunctionConfiguration(ctx context.Context, fm FunctionManifest) error {
retry := backoff.NewRetry(RequestRetryTime, backoff.NewConstant(RetryIntervalDuration))
updateFunctionConfigurationSucceed := false
var err error
for retry.WaitNext(ctx) {
_, err := retry.Do(ctx, func() (interface{}, error) {
configInput := &lambda.UpdateFunctionConfigurationInput{
FunctionName: aws.String(fm.Spec.Name),
Role: aws.String(fm.Spec.Role),
Expand Down Expand Up @@ -314,16 +312,15 @@ func (c *client) updateFunctionConfiguration(ctx context.Context, fm FunctionMan
SubnetIds: fm.Spec.VPCConfig.SubnetIDs,
}
}
_, err = c.client.UpdateFunctionConfiguration(ctx, configInput)
_, err := c.client.UpdateFunctionConfiguration(ctx, configInput)
if err != nil {
c.logger.Error("Failed to update function configuration")
} else {
updateFunctionConfigurationSucceed = true
break
return nil, fmt.Errorf("failed to update configuration for Lambda function %s: %w", fm.Spec.Name, err)
}
}
if !updateFunctionConfigurationSucceed {
return fmt.Errorf("failed to update configuration for Lambda function %s: %w", fm.Spec.Name, err)
return nil, nil
})
if err != nil {
return err
}

// Wait until function updated successfully.
Expand Down
18 changes: 9 additions & 9 deletions pkg/app/piped/trigger/cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,22 +58,22 @@ func (s *lastTriggeredCommitStore) Put(applicationID, commit string) error {

func (s *lastTriggeredCommitStore) getLastTriggeredDeployment(ctx context.Context, applicationID string) (*model.ApplicationDeploymentReference, error) {
var (
err error
resp *pipedservice.GetApplicationMostRecentDeploymentResponse
retry = pipedservice.NewRetry(3)
req = &pipedservice.GetApplicationMostRecentDeploymentRequest{
ApplicationId: applicationID,
Status: model.DeploymentStatus_DEPLOYMENT_PENDING,
}
)

for retry.WaitNext(ctx) {
if resp, err = s.apiClient.GetApplicationMostRecentDeployment(ctx, req); err == nil {
return resp.Deployment, nil
}
if !pipedservice.Retriable(err) {
return nil, err
d, err := retry.Do(ctx, func() (interface{}, error) {
resp, err := s.apiClient.GetApplicationMostRecentDeployment(ctx, req)
if err != nil {
return nil, pipedservice.NewRetriableErr(err)
}
return resp.Deployment, nil
})
if err != nil {
return nil, err
}
return nil, err
return d.(*model.ApplicationDeploymentReference), nil
}
12 changes: 6 additions & 6 deletions pkg/app/piped/trigger/deployment.go
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,6 @@ func buildDeployment(

func reportMostRecentlyTriggeredDeployment(ctx context.Context, client apiClient, d *model.Deployment) error {
var (
err error
req = &pipedservice.ReportApplicationMostRecentDeploymentRequest{
ApplicationId: d.ApplicationId,
Status: model.DeploymentStatus_DEPLOYMENT_PENDING,
Expand All @@ -129,11 +128,12 @@ func reportMostRecentlyTriggeredDeployment(ctx context.Context, client apiClient
retry = pipedservice.NewRetry(10)
)

for retry.WaitNext(ctx) {
if _, err = client.ReportApplicationMostRecentDeployment(ctx, req); err == nil {
return nil
_, err := retry.Do(ctx, func() (interface{}, error) {
_, err := client.ReportApplicationMostRecentDeployment(ctx, req)
if err != nil {
return nil, fmt.Errorf("failed to report most recent successful deployment: %w", err)
}
err = fmt.Errorf("failed to report most recent successful deployment: %w", err)
}
return nil, nil
})
return err
}
6 changes: 2 additions & 4 deletions pkg/backoff/backoff.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,6 @@ type Backoff interface {

type Retry interface {
Do(ctx context.Context, operation func() (interface{}, error)) (interface{}, error)
WaitNext(ctx context.Context) bool
Calls() int
}

Expand All @@ -62,8 +61,7 @@ type retry struct {
backoff Backoff
}

// TODO: Find all using of WaitNext and replace by Do to avoid panic.
func (r *retry) WaitNext(ctx context.Context) bool {
func (r *retry) waitNext(ctx context.Context) bool {
defer func() {
r.calls++
}()
Expand Down Expand Up @@ -98,7 +96,7 @@ func (r *retry) WaitNext(ctx context.Context) bool {
func (r *retry) Do(ctx context.Context, operation func() (interface{}, error)) (interface{}, error) {
var err error

for r.WaitNext(ctx) {
for r.waitNext(ctx) {
var data interface{}
data, err = operation()
if err == nil {
Expand Down
Loading
Loading