@ -61,6 +61,14 @@ type statusWaiter struct {
// when they don't set a timeout.
// when they don't set a timeout.
var DefaultStatusWatcherTimeout = 30 * time . Second
var DefaultStatusWatcherTimeout = 30 * time . Second
// DefaultStatusComputeWorkers controls the number of concurrent goroutines
// used to compute object status per informer. This prevents the informer
// notification pipeline from being blocked by slow API calls (e.g., LIST
// ReplicaSets/Pods for Deployments) when many resources are updated
// simultaneously.
// See https://github.com/fluxcd/cli-utils/pull/20
var DefaultStatusComputeWorkers = 8
func alwaysReady ( _ * unstructured . Unstructured ) ( * status . Result , error ) {
func alwaysReady ( _ * unstructured . Unstructured ) ( * status . Result , error ) {
return & status . Result {
return & status . Result {
Status : status . CurrentStatus ,
Status : status . CurrentStatus ,
@ -76,6 +84,7 @@ func (w *statusWaiter) WatchUntilReady(resourceList ResourceList, timeout time.D
defer cancel ( )
defer cancel ( )
w . Logger ( ) . Debug ( "waiting for resources" , "count" , len ( resourceList ) , "timeout" , timeout )
w . Logger ( ) . Debug ( "waiting for resources" , "count" , len ( resourceList ) , "timeout" , timeout )
sw := watcher . NewDefaultStatusWatcher ( w . client , w . restMapper )
sw := watcher . NewDefaultStatusWatcher ( w . client , w . restMapper )
sw . StatusComputeWorkers = DefaultStatusComputeWorkers
jobSR := helmStatusReaders . NewCustomJobStatusReader ( w . restMapper )
jobSR := helmStatusReaders . NewCustomJobStatusReader ( w . restMapper )
podSR := helmStatusReaders . NewCustomPodStatusReader ( w . restMapper )
podSR := helmStatusReaders . NewCustomPodStatusReader ( w . restMapper )
// We don't want to wait on any other resources as watchUntilReady is only for Helm hooks.
// We don't want to wait on any other resources as watchUntilReady is only for Helm hooks.
@ -98,6 +107,7 @@ func (w *statusWaiter) Wait(resourceList ResourceList, timeout time.Duration) er
defer cancel ( )
defer cancel ( )
w . Logger ( ) . Debug ( "waiting for resources" , "count" , len ( resourceList ) , "timeout" , timeout )
w . Logger ( ) . Debug ( "waiting for resources" , "count" , len ( resourceList ) , "timeout" , timeout )
sw := watcher . NewDefaultStatusWatcher ( w . client , w . restMapper )
sw := watcher . NewDefaultStatusWatcher ( w . client , w . restMapper )
sw . StatusComputeWorkers = DefaultStatusComputeWorkers
sw . StatusReader = statusreaders . NewStatusReader ( w . restMapper , w . readers ... )
sw . StatusReader = statusreaders . NewStatusReader ( w . restMapper , w . readers ... )
return w . wait ( ctx , resourceList , sw )
return w . wait ( ctx , resourceList , sw )
}
}
@ -110,6 +120,7 @@ func (w *statusWaiter) WaitWithJobs(resourceList ResourceList, timeout time.Dura
defer cancel ( )
defer cancel ( )
w . Logger ( ) . Debug ( "waiting for resources" , "count" , len ( resourceList ) , "timeout" , timeout )
w . Logger ( ) . Debug ( "waiting for resources" , "count" , len ( resourceList ) , "timeout" , timeout )
sw := watcher . NewDefaultStatusWatcher ( w . client , w . restMapper )
sw := watcher . NewDefaultStatusWatcher ( w . client , w . restMapper )
sw . StatusComputeWorkers = DefaultStatusComputeWorkers
newCustomJobStatusReader := helmStatusReaders . NewCustomJobStatusReader ( w . restMapper )
newCustomJobStatusReader := helmStatusReaders . NewCustomJobStatusReader ( w . restMapper )
readers := append ( [ ] engine . StatusReader ( nil ) , w . readers ... )
readers := append ( [ ] engine . StatusReader ( nil ) , w . readers ... )
readers = append ( readers , newCustomJobStatusReader )
readers = append ( readers , newCustomJobStatusReader )