From 426b22d2e5933ac513ee9fe9686201ceddc4d859 Mon Sep 17 00:00:00 2001 From: Tianxin Dong Date: Fri, 22 Apr 2022 13:14:51 +0800 Subject: [PATCH] Feat: add log provider (#3711) * Feat: add log provider Signed-off-by: FogDong * fix lift Signed-off-by: FogDong * fix vet Signed-off-by: FogDong * fix rebase vet Signed-off-by: FogDong --- .../v1alpha2/application/generator.go | 5 ++-- .../v1alpha2/application/generator_test.go | 16 +++++++---- pkg/stdlib/op.cue | 2 ++ pkg/stdlib/pkgs/util.cue | 7 +++++ pkg/workflow/providers/util/util.go | 27 ++++++++++++++++--- pkg/workflow/providers/util/util_test.go | 17 +++++++++++- pkg/workflow/tasks/discover.go | 15 ++++------- 7 files changed, 68 insertions(+), 21 deletions(-) diff --git a/pkg/controller/core.oam.dev/v1alpha2/application/generator.go b/pkg/controller/core.oam.dev/v1alpha2/application/generator.go index 45677169c..384d1d639 100644 --- a/pkg/controller/core.oam.dev/v1alpha2/application/generator.go +++ b/pkg/controller/core.oam.dev/v1alpha2/application/generator.go @@ -33,6 +33,7 @@ import ( "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/application/assemble" "github.com/oam-dev/kubevela/pkg/cue/model/value" "github.com/oam-dev/kubevela/pkg/cue/process" + monitorContext "github.com/oam-dev/kubevela/pkg/monitor/context" "github.com/oam-dev/kubevela/pkg/monitor/metrics" "github.com/oam-dev/kubevela/pkg/multicluster" "github.com/oam-dev/kubevela/pkg/oam" @@ -56,7 +57,7 @@ var ( // GenerateApplicationSteps generate application steps. // nolint:gocyclo -func (h *AppHandler) GenerateApplicationSteps(ctx context.Context, +func (h *AppHandler) GenerateApplicationSteps(ctx monitorContext.Context, app *v1beta1.Application, appParser *appfile.Parser, af *appfile.Appfile, @@ -68,7 +69,7 @@ func (h *AppHandler) GenerateApplicationSteps(ctx context.Context, appParser, appRev, af), h.renderComponentFunc(appParser, appRev, af)) http.Install(handlerProviders, h.r.Client, app.Namespace) pCtx := process.NewContext(generateContextDataFromApp(app, appRev.Name)) - taskDiscover := tasks.NewTaskDiscoverFromRevision(handlerProviders, h.r.pd, appRev, h.r.dm, pCtx) + taskDiscover := tasks.NewTaskDiscoverFromRevision(ctx, handlerProviders, h.r.pd, appRev, h.r.dm, pCtx) multiclusterProvider.Install(handlerProviders, h.r.Client, app) terraformProvider.Install(handlerProviders, app, func(comp common.ApplicationComponent) (*appfile.Workload, error) { return appParser.ParseWorkloadFromRevision(comp, appRev) diff --git a/pkg/controller/core.oam.dev/v1alpha2/application/generator_test.go b/pkg/controller/core.oam.dev/v1alpha2/application/generator_test.go index b5a6bc82a..bddebae3e 100644 --- a/pkg/controller/core.oam.dev/v1alpha2/application/generator_test.go +++ b/pkg/controller/core.oam.dev/v1alpha2/application/generator_test.go @@ -30,6 +30,7 @@ import ( "github.com/oam-dev/kubevela/apis/core.oam.dev/common" oamcore "github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1" + monitorContext "github.com/oam-dev/kubevela/pkg/monitor/context" "github.com/oam-dev/kubevela/pkg/oam/util" ) @@ -114,7 +115,8 @@ var _ = Describe("Test Application workflow generator", func() { handler, err := NewAppHandler(ctx, reconciler, app, appParser) Expect(err).Should(Succeed()) - taskRunner, err := handler.GenerateApplicationSteps(ctx, app, appParser, af, appRev) + logCtx := monitorContext.NewTraceContext(ctx, "") + taskRunner, err := handler.GenerateApplicationSteps(logCtx, app, appParser, af, appRev) Expect(err).To(BeNil()) Expect(len(taskRunner)).Should(BeEquivalentTo(2)) Expect(taskRunner[0].Name()).Should(BeEquivalentTo("myweb1")) @@ -155,7 +157,8 @@ var _ = Describe("Test Application workflow generator", func() { handler, err := NewAppHandler(ctx, reconciler, app, appParser) Expect(err).Should(Succeed()) - taskRunner, err := handler.GenerateApplicationSteps(ctx, app, appParser, af, appRev) + logCtx := monitorContext.NewTraceContext(ctx, "") + taskRunner, err := handler.GenerateApplicationSteps(logCtx, app, appParser, af, appRev) Expect(err).To(BeNil()) Expect(len(taskRunner)).Should(BeEquivalentTo(2)) Expect(taskRunner[0].Name()).Should(BeEquivalentTo("myweb1")) @@ -274,7 +277,8 @@ var _ = Describe("Test Application workflow generator", func() { handler, err := NewAppHandler(ctx, reconciler, app, appParser) Expect(err).Should(Succeed()) - taskRunner, err := handler.GenerateApplicationSteps(ctx, app, appParser, af, appRev) + logCtx := monitorContext.NewTraceContext(ctx, "") + taskRunner, err := handler.GenerateApplicationSteps(logCtx, app, appParser, af, appRev) Expect(err).To(BeNil()) Expect(len(taskRunner)).Should(BeEquivalentTo(2)) Expect(taskRunner[0].Name()).Should(BeEquivalentTo("myweb1")) @@ -314,7 +318,8 @@ var _ = Describe("Test Application workflow generator", func() { handler, err := NewAppHandler(ctx, reconciler, app, appParser) Expect(err).Should(Succeed()) - _, err = handler.GenerateApplicationSteps(ctx, app, appParser, af, appRev) + logCtx := monitorContext.NewTraceContext(ctx, "") + _, err = handler.GenerateApplicationSteps(logCtx, app, appParser, af, appRev) Expect(err).NotTo(BeNil()) }) @@ -351,7 +356,8 @@ var _ = Describe("Test Application workflow generator", func() { handler, err := NewAppHandler(ctx, reconciler, app, appParser) Expect(err).Should(Succeed()) - _, err = handler.GenerateApplicationSteps(ctx, app, appParser, af, appRev) + logCtx := monitorContext.NewTraceContext(ctx, "") + _, err = handler.GenerateApplicationSteps(logCtx, app, appParser, af, appRev) Expect(err).NotTo(BeNil()) }) }) diff --git a/pkg/stdlib/op.cue b/pkg/stdlib/op.cue index f6280c8f6..7b474d8b8 100644 --- a/pkg/stdlib/op.cue +++ b/pkg/stdlib/op.cue @@ -166,6 +166,8 @@ import ( #ConvertString: util.#String +#Log: util.#Log + #DateToTimestamp: time.#DateToTimestamp #TimestampToDate: time.#TimestampToDate diff --git a/pkg/stdlib/pkgs/util.cue b/pkg/stdlib/pkgs/util.cue index e137d40ba..16df0c573 100644 --- a/pkg/stdlib/pkgs/util.cue +++ b/pkg/stdlib/pkgs/util.cue @@ -15,3 +15,10 @@ str?: string ... } + +#Log: { + #do: "log" + #provider: "util" + + data: {...} +} diff --git a/pkg/workflow/providers/util/util.go b/pkg/workflow/providers/util/util.go index 07cb768ed..40911389b 100644 --- a/pkg/workflow/providers/util/util.go +++ b/pkg/workflow/providers/util/util.go @@ -19,6 +19,7 @@ package util import ( "github.com/oam-dev/kubevela/pkg/cue/model" "github.com/oam-dev/kubevela/pkg/cue/model/value" + monitorContext "github.com/oam-dev/kubevela/pkg/monitor/context" wfContext "github.com/oam-dev/kubevela/pkg/workflow/context" "github.com/oam-dev/kubevela/pkg/workflow/providers" "github.com/oam-dev/kubevela/pkg/workflow/types" @@ -29,7 +30,9 @@ const ( ProviderName = "util" ) -type provider struct{} +type provider struct { + logCtx monitorContext.Context +} func (p *provider) PatchK8sObject(ctx wfContext.Context, v *value.Value, act types.Action) error { val, err := v.LookupValue("value") @@ -72,11 +75,29 @@ func (p *provider) String(ctx wfContext.Context, v *value.Value, act types.Actio return v.FillObject(string(s), "str") } +// Log print cue value in log +func (p *provider) Log(ctx wfContext.Context, v *value.Value, act types.Action) error { + data, err := v.LookupValue("data") + if err != nil { + return err + } + s, err := data.String() + if err != nil { + return err + } + logCtx := p.logCtx.Fork("cue logs") + logCtx.Info(s) + return nil +} + // Install register handlers to provider discover. -func Install(p providers.Providers) { - prd := &provider{} +func Install(ctx monitorContext.Context, p providers.Providers) { + prd := &provider{ + logCtx: ctx, + } p.Register(ProviderName, map[string]providers.Handler{ "patch-k8s-object": prd.PatchK8sObject, "string": prd.String, + "log": prd.Log, }) } diff --git a/pkg/workflow/providers/util/util_test.go b/pkg/workflow/providers/util/util_test.go index d41a50edd..ce146bc31 100644 --- a/pkg/workflow/providers/util/util_test.go +++ b/pkg/workflow/providers/util/util_test.go @@ -17,6 +17,7 @@ package util import ( + "context" "errors" "testing" @@ -24,6 +25,7 @@ import ( "github.com/stretchr/testify/require" "github.com/oam-dev/kubevela/pkg/cue/model/value" + monitorContext "github.com/oam-dev/kubevela/pkg/monitor/context" "github.com/oam-dev/kubevela/pkg/workflow/providers" ) @@ -192,9 +194,22 @@ func TestConvertString(t *testing.T) { } } +func TestLog(t *testing.T) { + r := require.New(t) + v, err := value.NewValue(` +data: "test" +`, nil, "") + r.NoError(err) + logCtx := monitorContext.NewTraceContext(context.Background(), "") + prd := &provider{logCtx: logCtx} + err = prd.Log(nil, v, nil) + r.NoError(err) +} + func TestInstall(t *testing.T) { + logCtx := monitorContext.NewTraceContext(context.Background(), "") p := providers.NewProviders() - Install(p) + Install(logCtx, p) h, ok := p.GetHandler("util", "string") r := require.New(t) r.Equal(ok, true) diff --git a/pkg/workflow/tasks/discover.go b/pkg/workflow/tasks/discover.go index bbf6fd83d..98cb4dfe9 100644 --- a/pkg/workflow/tasks/discover.go +++ b/pkg/workflow/tasks/discover.go @@ -29,6 +29,7 @@ import ( "github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1" "github.com/oam-dev/kubevela/pkg/cue/packages" "github.com/oam-dev/kubevela/pkg/cue/process" + monitorContext "github.com/oam-dev/kubevela/pkg/monitor/context" "github.com/oam-dev/kubevela/pkg/oam/discoverymapper" "github.com/oam-dev/kubevela/pkg/velaql/providers/query" wfContext "github.com/oam-dev/kubevela/pkg/workflow/context" @@ -86,11 +87,11 @@ func suspend(step v1beta1.WorkflowStep, opt *types.GeneratorOptions) (types.Task return tr, nil } -func newTaskDiscover(providerHandlers providers.Providers, pd *packages.PackageDiscover, pCtx process.Context, templateLoader template.Loader) types.TaskDiscover { +func newTaskDiscover(ctx monitorContext.Context, providerHandlers providers.Providers, pd *packages.PackageDiscover, pCtx process.Context, templateLoader template.Loader) types.TaskDiscover { // install builtin provider workspace.Install(providerHandlers) email.Install(providerHandlers) - util.Install(providerHandlers) + util.Install(ctx, providerHandlers) return &taskDiscover{ builtins: map[string]types.TaskGenerator{ @@ -101,16 +102,10 @@ func newTaskDiscover(providerHandlers providers.Providers, pd *packages.PackageD } } -// NewTaskDiscover will create a client for load task generator. -func NewTaskDiscover(providerHandlers providers.Providers, pd *packages.PackageDiscover, cli client.Client, dm discoverymapper.DiscoveryMapper, pCtx process.Context) types.TaskDiscover { - templateLoader := template.NewWorkflowStepTemplateLoader(cli, dm) - return newTaskDiscover(providerHandlers, pd, pCtx, templateLoader) -} - // NewTaskDiscoverFromRevision will create a client for load task generator from ApplicationRevision. -func NewTaskDiscoverFromRevision(providerHandlers providers.Providers, pd *packages.PackageDiscover, rev *v1beta1.ApplicationRevision, dm discoverymapper.DiscoveryMapper, pCtx process.Context) types.TaskDiscover { +func NewTaskDiscoverFromRevision(ctx monitorContext.Context, providerHandlers providers.Providers, pd *packages.PackageDiscover, rev *v1beta1.ApplicationRevision, dm discoverymapper.DiscoveryMapper, pCtx process.Context) types.TaskDiscover { templateLoader := template.NewWorkflowStepTemplateRevisionLoader(rev, dm) - return newTaskDiscover(providerHandlers, pd, pCtx, templateLoader) + return newTaskDiscover(ctx, providerHandlers, pd, pCtx, templateLoader) } type suspendTaskRunner struct {