mirror of
https://github.com/kubevela/kubevela.git
synced 2026-08-21 21:46:56 +00:00
* Feat: ref component Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Feat: support topology and override Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Feat: add support for external policy and workflow Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Feat: add admission control Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: disable cross namespace ref object Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Chore: refactor Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Feat: support labelSelector in ref-objects Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Feat: add pre approve for deploy step Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Chore: refactor Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: test Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Feat: support comp/trait type in override policy even not used by prototype Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Feat: support regex match for patch component name Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: labelSelector not work for cluster Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: ref workflow contains external policy Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: revision test Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Feat: parallel apply components Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Feat: add test for oam provider Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: service ref-comp & indirect trait ns Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: align namespace setting for chart Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: add strict unmarshal and reformat Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: merge with cluster rework Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Feat: patch trait-def Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: apply components + load dynamic component Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: add test for loadPoliciesInOrder Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Feat: add test for open merge Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: reformat & add test for step generator Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: add test for parse override policy related defs Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: add test for multicluster provider (expandTopology and overrideConfiguration) Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: add admission test Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: revert trait status pass in component status Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: add test for dependency in workflowstep & standalone multicluster test Signed-off-by: Somefive <yd219913@alibaba-inc.com> * Fix: add check for ref and steps in WorkflowStep & enhance ref-objects scheme check Signed-off-by: Somefive <yd219913@alibaba-inc.com>
147 lines
4.2 KiB
Go
147 lines
4.2 KiB
Go
/*
|
|
Copyright 2021 The KubeVela Authors.
|
|
|
|
Licensed under the Apache License, Version 2.0 (the "License");
|
|
you may not use this file except in compliance with the License.
|
|
You may obtain a copy of the License at
|
|
|
|
http://www.apache.org/licenses/LICENSE-2.0
|
|
|
|
Unless required by applicable law or agreed to in writing, software
|
|
distributed under the License is distributed on an "AS IS" BASIS,
|
|
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
See the License for the specific language governing permissions and
|
|
limitations under the License.
|
|
*/
|
|
|
|
package resourcekeeper
|
|
|
|
import (
|
|
"context"
|
|
"sync"
|
|
|
|
"github.com/pkg/errors"
|
|
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
|
|
kerrors "k8s.io/apimachinery/pkg/util/errors"
|
|
|
|
"github.com/oam-dev/kubevela/pkg/multicluster"
|
|
"github.com/oam-dev/kubevela/pkg/oam"
|
|
"github.com/oam-dev/kubevela/pkg/resourcetracker"
|
|
"github.com/oam-dev/kubevela/pkg/utils/apply"
|
|
)
|
|
|
|
// MaxDispatchConcurrent is the max dispatch concurrent number
|
|
var MaxDispatchConcurrent = 10
|
|
|
|
// DispatchOption option for dispatch
|
|
type DispatchOption interface {
|
|
ApplyToDispatchConfig(*dispatchConfig)
|
|
}
|
|
|
|
type dispatchConfig struct {
|
|
rtConfig
|
|
metaOnly bool
|
|
}
|
|
|
|
func newDispatchConfig(options ...DispatchOption) *dispatchConfig {
|
|
cfg := &dispatchConfig{}
|
|
for _, option := range options {
|
|
option.ApplyToDispatchConfig(cfg)
|
|
}
|
|
return cfg
|
|
}
|
|
|
|
// Dispatch dispatch resources
|
|
func (h *resourceKeeper) Dispatch(ctx context.Context, manifests []*unstructured.Unstructured, options ...DispatchOption) (err error) {
|
|
if h.applyOncePolicy != nil && h.applyOncePolicy.Enable {
|
|
options = append(options, MetaOnlyOption{})
|
|
}
|
|
// 0. check admission
|
|
if err = h.AdmissionCheck(ctx, manifests); err != nil {
|
|
return err
|
|
}
|
|
// 1. record manifests in resourcetracker
|
|
if err = h.record(ctx, manifests, options...); err != nil {
|
|
return err
|
|
}
|
|
// 2. apply manifests
|
|
if err = h.dispatch(ctx, manifests); err != nil {
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (h *resourceKeeper) record(ctx context.Context, manifests []*unstructured.Unstructured, options ...DispatchOption) error {
|
|
h.mu.Lock()
|
|
defer h.mu.Unlock()
|
|
var rootManifests []*unstructured.Unstructured
|
|
var versionManifests []*unstructured.Unstructured
|
|
|
|
for _, manifest := range manifests {
|
|
if manifest != nil {
|
|
_options := options
|
|
if h.garbageCollectPolicy != nil {
|
|
if strategy := h.garbageCollectPolicy.FindStrategy(manifest); strategy != nil {
|
|
_options = append(_options, GarbageCollectStrategyOption(*strategy))
|
|
}
|
|
}
|
|
cfg := newDispatchConfig(_options...)
|
|
if !cfg.skipRT {
|
|
if cfg.useRoot {
|
|
rootManifests = append(rootManifests, manifest)
|
|
} else {
|
|
versionManifests = append(versionManifests, manifest)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
cfg := newDispatchConfig(options...)
|
|
if len(rootManifests) != 0 {
|
|
rt, err := h.getRootRT(ctx)
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed to get resourcetracker")
|
|
}
|
|
if err = resourcetracker.RecordManifestsInResourceTracker(multicluster.ContextInLocalCluster(ctx), h.Client, rt, rootManifests, cfg.metaOnly); err != nil {
|
|
return errors.Wrapf(err, "failed to record resources in resourcetracker %s", rt.Name)
|
|
}
|
|
}
|
|
|
|
rt, err := h.getCurrentRT(ctx)
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed to get resourcetracker")
|
|
}
|
|
if err = resourcetracker.RecordManifestsInResourceTracker(multicluster.ContextInLocalCluster(ctx), h.Client, rt, versionManifests, cfg.metaOnly); err != nil {
|
|
return errors.Wrapf(err, "failed to record resources in resourcetracker %s", rt.Name)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (h *resourceKeeper) dispatch(ctx context.Context, manifests []*unstructured.Unstructured) error {
|
|
var errs []error
|
|
var l sync.Mutex
|
|
var wg sync.WaitGroup
|
|
|
|
ch := make(chan struct{}, MaxDispatchConcurrent)
|
|
applyOpts := []apply.ApplyOption{apply.MustBeControlledByApp(h.app), apply.NotUpdateRenderHashEqual()}
|
|
|
|
for i := 0; i < len(manifests); i++ {
|
|
ch <- struct{}{}
|
|
wg.Add(1)
|
|
go func(index int) {
|
|
defer wg.Done()
|
|
manifest := manifests[index]
|
|
applyCtx := multicluster.ContextWithClusterName(ctx, oam.GetCluster(manifest))
|
|
err := h.applicator.Apply(applyCtx, manifest, applyOpts...)
|
|
if err != nil {
|
|
l.Lock()
|
|
errs = append(errs, err)
|
|
l.Unlock()
|
|
}
|
|
<-ch
|
|
}(i)
|
|
}
|
|
wg.Wait()
|
|
return kerrors.NewAggregate(errs)
|
|
}
|