From a1ed763aa0b3e646d52b60be74d693d92d2de5a1 Mon Sep 17 00:00:00 2001 From: Fu Diwei Date: Wed, 10 Sep 2025 16:38:43 +0800 Subject: [PATCH] feat: render template string with workflow variables --- internal/certificate/service.go | 2 +- internal/workflow/dispatcher/dispatcher.go | 27 ++++-- internal/workflow/engine/engine.go | 44 ++++++++-- internal/workflow/engine/executor_bizapply.go | 88 +++++++++++++++---- .../workflow/engine/executor_bizdeploy.go | 18 ++-- .../workflow/engine/executor_bizmonitor.go | 61 ++++++++++--- .../workflow/engine/executor_biznotify.go | 20 +++-- .../workflow/engine/executor_bizupload.go | 78 +++++++++++++--- internal/workflow/engine/state.go | 53 ++++++++++- 9 files changed, 315 insertions(+), 76 deletions(-) diff --git a/internal/certificate/service.go b/internal/certificate/service.go index 1facb4bac..f6936c99d 100644 --- a/internal/certificate/service.go +++ b/internal/certificate/service.go @@ -170,7 +170,7 @@ func (s *CertificateService) ValidateCertificate(ctx context.Context, req *dtos. certX509, err := xcert.ParseCertificateFromPEM(req.Certificate) if err != nil { return nil, err - } else if time.Now().After(certX509.NotAfter) { + } else if certX509.NotAfter.Before(time.Now()) { return nil, fmt.Errorf("certificate has expired at %s", certX509.NotAfter.UTC().Format(time.RFC3339)) } diff --git a/internal/workflow/dispatcher/dispatcher.go b/internal/workflow/dispatcher/dispatcher.go index 7054f0406..a9238d5f4 100644 --- a/internal/workflow/dispatcher/dispatcher.go +++ b/internal/workflow/dispatcher/dispatcher.go @@ -212,7 +212,9 @@ func (wd *workflowDispatcher) Cancel(ctx context.Context, runId string) error { } func (wd *workflowDispatcher) tryExecuteAsync(task *taskInfo) { + var workflow *domain.Workflow var workflowRun *domain.WorkflowRun + var err error // 捕获 panic defer func() { @@ -238,23 +240,28 @@ func (wd *workflowDispatcher) tryExecuteAsync(task *taskInfo) { }() // 查询运行实体,并级联更新状态 - if run, err := wd.workflowRunRepo.GetById(task.ctx, task.RunId); err != nil { + if workflowRun, err = wd.workflowRunRepo.GetById(task.ctx, task.RunId); err != nil { if !(errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded)) { wd.syslog.Error(fmt.Sprintf("failed to get workflow run #%s record", task.RunId), slog.Any("error", err)) } return } else { - workflowRun = run - - if run.Status == domain.WorkflowRunStatusTypePending { - run.Status = domain.WorkflowRunStatusTypeProcessing - wd.workflowRunRepo.SaveWithCascading(task.ctx, run) + if workflowRun.Status == domain.WorkflowRunStatusTypePending { + workflowRun.Status = domain.WorkflowRunStatusTypeProcessing + wd.workflowRunRepo.SaveWithCascading(task.ctx, workflowRun) } else { // WTF? That should be impossible! return } } + // 查询工作流实体 + workflow, err = wd.workflowRepo.GetById(task.ctx, workflowRun.WorkflowId) + if err != nil { + wd.syslog.Error(fmt.Sprintf("failed to get workflow #%s record", workflowRun.WorkflowId), slog.Any("error", err)) + return + } + // 初始化工作流引擎 logsBuf := make(domain.WorkflowLogs, 0) we := engine.NewWorkflowEngine() @@ -326,7 +333,13 @@ func (wd *workflowDispatcher) tryExecuteAsync(task *taskInfo) { // 执行工作流 wd.syslog.Info(fmt.Sprintf("workflow run #%s (work#%s) started", task.RunId, task.WorkflowId)) - we.Invoke(task.ctx, workflowRun.WorkflowId, workflowRun.Id, workflowRun.Graph) + we.Invoke(task.ctx, engine.WorkflowExecution{ + WorkflowId: workflowRun.WorkflowId, + WorkflowName: workflow.Name, + RunId: workflowRun.Id, + RunTrigger: workflowRun.Trigger, + Graph: workflowRun.Graph, + }) wd.syslog.Info(fmt.Sprintf("workflow run #%s (work#%s) stopped", task.RunId, task.WorkflowId)) } diff --git a/internal/workflow/engine/engine.go b/internal/workflow/engine/engine.go index 39a13143c..3f92b6c31 100644 --- a/internal/workflow/engine/engine.go +++ b/internal/workflow/engine/engine.go @@ -15,8 +15,16 @@ import ( "github.com/certimate-go/certimate/pkg/logging" ) +type WorkflowExecution struct { + WorkflowId string + WorkflowName string + RunId string + RunTrigger domain.WorkflowTriggerType + Graph *Graph +} + type WorkflowEngine interface { - Invoke(ctx context.Context, workflowId string, runId string, graph *Graph) error + Invoke(ctx context.Context, execution WorkflowExecution) error OnStart(callback func(ctx context.Context) error) OnEnd(callback func(ctx context.Context) error) @@ -46,22 +54,33 @@ type workflowEngine struct { var _ WorkflowEngine = (*workflowEngine)(nil) -func (we *workflowEngine) Invoke(ctx context.Context, workflowId string, runId string, runGraph *Graph) error { - we.fireOnStartHooks(ctx) - +func (we *workflowEngine) Invoke(ctx context.Context, execution WorkflowExecution) error { defer func() { if r := recover(); r != nil { we.fireOnErrorHooks(ctx, fmt.Errorf("workflow engine panic: %v", r)) } }() + we.fireOnStartHooks(ctx) + + wfIOs := newInOutManager() + + wfVars := newVariableManager() + wfVars.Set("workflow.id", execution.WorkflowId, "string") + wfVars.Set("workflow.name", execution.WorkflowName, "string") + wfVars.Set("run.id", execution.RunId, "string") + wfVars.Set("run.trigger", execution.RunTrigger, "string") + wfVars.Set("error.nodeId", "", "string") + wfVars.Set("error.nodeName", "", "string") + wfVars.Set("error.message", "", "string") + wfCtx := (&WorkflowContext{}). - SetExecutingWorkflow(workflowId, runId, runGraph). + SetExecutingWorkflow(execution.WorkflowId, execution.RunId, execution.Graph). SetEngine(we). - SetVariablesManager(newVariableManager()). - SetInputsManager(newInOutManager()). + SetInputsManager(wfIOs). + SetVariablesManager(wfVars). SetContext(ctx) - if err := we.executeBlocks(wfCtx, runGraph.Nodes); err != nil { + if err := we.executeBlocks(wfCtx, execution.Graph.Nodes); err != nil { if !errors.Is(err, ErrTerminated) { we.fireOnErrorHooks(ctx, err) return err @@ -131,6 +150,9 @@ func (we *workflowEngine) executeNode(wfCtx *WorkflowContext, node *Node) error executor.SetLogger(logger) } + wfCtx.variables.SetScoped(node.Id, stateVarKeyNodeId, node.Id, "string") + wfCtx.variables.SetScoped(node.Id, stateVarKeyNodeName, node.Data.Name, "string") + // 节点已禁用,直接跳过执行 if node.Data.Disabled { return nil @@ -141,6 +163,12 @@ func (we *workflowEngine) executeNode(wfCtx *WorkflowContext, node *Node) error execCtx := newNodeExecutionContext(wfCtx, node) execRes, err := executor.Execute(execCtx) if err != nil && !errors.Is(err, ErrTerminated) { + if !errors.Is(err, ErrBlocksException) { + wfCtx.variables.Set("error.nodeId", node.Id, "string") + wfCtx.variables.Set("error.nodeName", node.Data.Name, "string") + wfCtx.variables.Set("error.message", err.Error(), "string") + } + we.fireOnNodeErrorHooks(wfCtx.ctx, node, err) return err } diff --git a/internal/workflow/engine/executor_bizapply.go b/internal/workflow/engine/executor_bizapply.go index 7f8e16f42..e1a232a64 100644 --- a/internal/workflow/engine/executor_bizapply.go +++ b/internal/workflow/engine/executor_bizapply.go @@ -35,13 +35,18 @@ func init() { } /** - * Result Variables: - * - node.skipped: boolean - * - certificate.validity: boolean - * - certificate.daysLeft: number + * Outputs: + * - ref: "certificate": string * - * Result Outputs: - * - ref: certificate: string + * Variables: + * - "node.skipped": boolean + * - "certificate.domain": string + * - "certificate.domains": string + * - "certificate.notBefore": datetime + * - "certificate.notAfter": datetime + * - "certificate.hoursLeft": number + * - "certificate.daysLeft": number + * - "certificate.validity": boolean */ type bizApplyNodeExecutor struct { nodeExecutor @@ -62,9 +67,8 @@ func (ne *bizApplyNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeExe if err != nil { return execRes, err } else if lastCertificate != nil { - execRes.AddOutput(stateIOTypeRef, "certificate", fmt.Sprintf("%s#%s", domain.CollectionNameCertificate, lastCertificate.Id), "string") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateValidity, time.Now().After(lastCertificate.ValidityNotAfter), "boolean") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDaysLeft, int32(time.Until(lastCertificate.ValidityNotAfter).Hours()/24), "number") + ne.setOuputsOfResult(execCtx, execRes, lastCertificate, false) + ne.setVariablesOfResult(execCtx, execRes, lastCertificate) } // 检测是否可以跳过本次执行 @@ -73,12 +77,17 @@ func (ne *bizApplyNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeExe execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyNodeSkipped, true, "boolean") return execRes, nil - } else if reason != "" { - ne.logger.Info(fmt.Sprintf("re-apply, because %s", reason)) } else { - ne.logger.Info("no found last issued certificate, begin to apply") + if reason != "" { + ne.logger.Info(fmt.Sprintf("re-apply, because %s", reason)) + } else { + ne.logger.Info("no found last issued certificate, begin to apply") + } + + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyNodeSkipped, false, "boolean") } + // 申请证书 obtainResp, err := ne.executeObtain(execCtx, &nodeCfg, lastCertificate) if err != nil { return execRes, err @@ -112,17 +121,15 @@ func (ne *bizApplyNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeExe ne.logger.Info("certificate saved", slog.String("recordId", certificate.Id)) } - // 保存 ARI 记录 + // 保存 ARI 替换状态 if lastCertificate != nil && obtainResp.ARIReplaced { lastCertificate.ACMERenewed = true ne.certificateRepo.Save(execCtx.ctx, lastCertificate) } // 节点输出 - execRes.AddOutputWithPersistent(stateIOTypeRef, "certificate", fmt.Sprintf("%s#%s", domain.CollectionNameCertificate, certificate.Id), "string") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyNodeSkipped, false, "boolean") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateValidity, true, "boolean") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDaysLeft, int32(time.Until(certificate.ValidityNotAfter).Hours()/24), "number") + ne.setOuputsOfResult(execCtx, execRes, certificate, true) + ne.setVariablesOfResult(execCtx, execRes, certificate) ne.logger.Info("application completed") return execRes, nil @@ -353,6 +360,53 @@ func (ne *bizApplyNodeExecutor) executeObtain(execCtx *NodeExecutionContext, nod return obtainResp, nil } +func (ne *bizApplyNodeExecutor) setOuputsOfResult(execCtx *NodeExecutionContext, execRes *NodeExecutionResult, certificate *domain.Certificate, persistent bool) { + if certificate != nil { + key := "certificate" + value := fmt.Sprintf("%s#%s", domain.CollectionNameCertificate, certificate.Id) + if persistent { + execRes.AddOutputWithPersistent(stateIOTypeRef, key, value, "string") + } else { + execRes.AddOutput(stateIOTypeRef, key, value, "string") + } + } +} + +func (ne *bizApplyNodeExecutor) setVariablesOfResult(execCtx *NodeExecutionContext, execRes *NodeExecutionResult, certificate *domain.Certificate) { + var vDomain string + var vDomains string + var vNotBefore time.Time + var vNotAfter time.Time + var vHoursLeft int32 + var vDaysLeft int32 + var vValidity bool + + if certificate != nil { + vDomain = strings.Split(certificate.SubjectAltNames, ";")[0] + vDomains = certificate.SubjectAltNames + vNotBefore = certificate.ValidityNotBefore + vNotAfter = certificate.ValidityNotAfter + vHoursLeft = int32(math.Floor(time.Until(certificate.ValidityNotAfter).Hours())) + vDaysLeft = int32(math.Floor(time.Until(certificate.ValidityNotAfter).Hours() / 24)) + vValidity = certificate.ValidityNotAfter.After(time.Now()) + } + + execRes.AddVariable(stateVarKeyCertificateDomain, vDomain, "string") + execRes.AddVariable(stateVarKeyCertificateDomains, vDomains, "string") + execRes.AddVariable(stateVarKeyCertificateNotBefore, vNotBefore, "datetime") + execRes.AddVariable(stateVarKeyCertificateNotAfter, vNotAfter, "datetime") + execRes.AddVariable(stateVarKeyCertificateHoursLeft, vHoursLeft, "number") + execRes.AddVariable(stateVarKeyCertificateDaysLeft, vDaysLeft, "number") + execRes.AddVariable(stateVarKeyCertificateValidity, vValidity, "boolean") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDomain, vDomain, "string") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDomains, vDomains, "string") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateNotBefore, vNotBefore, "datetime") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateNotAfter, vNotAfter, "datetime") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateHoursLeft, vHoursLeft, "number") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDaysLeft, vDaysLeft, "number") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateValidity, vValidity, "boolean") +} + func newBizApplyNodeExecutor() NodeExecutor { return &bizApplyNodeExecutor{ nodeExecutor: nodeExecutor{logger: slog.Default()}, diff --git a/internal/workflow/engine/executor_bizdeploy.go b/internal/workflow/engine/executor_bizdeploy.go index 1f05e5f52..fed162d06 100644 --- a/internal/workflow/engine/executor_bizdeploy.go +++ b/internal/workflow/engine/executor_bizdeploy.go @@ -12,8 +12,11 @@ import ( ) /** - * Result Variables: - * - node.skipped: boolean + * Inputs: + * - ref: "certificate": string + * + * Variables: + * - "node.skipped": boolean */ type bizDeployNodeExecutor struct { nodeExecutor @@ -64,7 +67,11 @@ func (ne *bizDeployNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeEx return execRes, nil } else if reason != "" { ne.logger.Info(fmt.Sprintf("re-deploy, because %s", reason)) + + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyNodeSkipped, false, "boolean") } + } else { + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyNodeSkipped, false, "boolean") } // 读取部署提供商授权 @@ -77,10 +84,8 @@ func (ne *bizDeployNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeEx } } - // 初始化部署器 - deployClient := certdeploy.NewClient(certdeploy.WithLogger(ne.logger)) - // 部署证书 + deployer := certdeploy.NewClient(certdeploy.WithLogger(ne.logger)) deployReq := &certdeploy.DeployCertificateRequest{ Provider: nodeCfg.Provider, ProviderAccessConfig: providerAccessConfig, @@ -88,14 +93,13 @@ func (ne *bizDeployNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeEx Certificate: inputCertificate.Certificate, PrivateKey: inputCertificate.PrivateKey, } - if _, err := deployClient.DeployCertificate(execCtx.ctx, deployReq); err != nil { + if _, err := deployer.DeployCertificate(execCtx.ctx, deployReq); err != nil { ne.logger.Warn("could not deploy certificate") return execRes, err } // 节点输出 execRes.outputForced = true - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyNodeSkipped, false, "boolean") ne.logger.Info("deployment completed") return execRes, nil diff --git a/internal/workflow/engine/executor_bizmonitor.go b/internal/workflow/engine/executor_bizmonitor.go index 1d8f40515..01e17c051 100644 --- a/internal/workflow/engine/executor_bizmonitor.go +++ b/internal/workflow/engine/executor_bizmonitor.go @@ -17,9 +17,14 @@ import ( ) /** - * Result Variables: - * - certificate.validity: boolean - * - certificate.daysLeft: number + * Variables: + * - "certificate.domain": string + * - "certificate.domains": string + * - "certificate.notBefore": datetime + * - "certificate.notAfter": datetime + * - "certificate.hoursLeft": number + * - "certificate.daysLeft": number + * - "certificate.validity": boolean */ type bizMonitorNodeExecutor struct { nodeExecutor @@ -51,7 +56,7 @@ func (ne *bizMonitorNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeE var certs []*x509.Certificate for attempt := 0; attempt < MAX_ATTEMPTS; attempt++ { if attempt > 0 { - ne.logger.Info(fmt.Sprintf("retry %d time(s) ...", attempt, targetAddr)) + ne.logger.Info(fmt.Sprintf("retry %d time(s) ...", attempt)) select { case <-execCtx.ctx.Done(): @@ -73,8 +78,7 @@ func (ne *bizMonitorNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeE if len(certs) == 0 { ne.logger.Warn("no ssl certificates retrieved in http response") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateValidity, false, "boolean") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDaysLeft, 0, "number") + ne.setVariablesOfResult(execCtx, execRes, nil) } else { cert := certs[0] // 只取证书链中的第一个证书,即服务器证书 ne.logger.Info(fmt.Sprintf("ssl certificate retrieved (serial='%s', subject='%s', issuer='%s', not_before='%s', not_after='%s', sans='%s')", @@ -82,15 +86,13 @@ func (ne *bizMonitorNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeE cert.NotBefore.Format(time.RFC3339), cert.NotAfter.Format(time.RFC3339), strings.Join(cert.DNSNames, ";")), ) + ne.setVariablesOfResult(execCtx, execRes, cert) now := time.Now() isCertPeriodValid := now.Before(cert.NotAfter) && now.After(cert.NotBefore) isCertHostMatched := cert.VerifyHostname(targetDomain) == nil - + daysLeft := int32(math.Floor(time.Until(cert.NotAfter).Hours() / 24)) validated := isCertPeriodValid && isCertHostMatched - daysLeft := int(math.Floor(time.Until(cert.NotAfter).Hours() / 24)) - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateValidity, validated, "boolean") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDaysLeft, daysLeft, "number") if validated { ne.logger.Info(fmt.Sprintf("the certificate is valid, and will expire in %d day(s)", daysLeft)) @@ -102,6 +104,10 @@ func (ne *bizMonitorNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeE } else { ne.logger.Warn("the certificate is invalid") } + + // 除了验证证书有效期,还要确保证书与域名匹配 + execRes.AddVariable(stateVarKeyCertificateValidity, false, "boolean") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateValidity, false, "boolean") } } } @@ -146,6 +152,41 @@ func (ne *bizMonitorNodeExecutor) tryRetrievePeerCertificates(execCtx *NodeExecu return resp.TLS.PeerCertificates, nil } +func (ne *bizMonitorNodeExecutor) setVariablesOfResult(execCtx *NodeExecutionContext, execRes *NodeExecutionResult, certX509 *x509.Certificate) { + var vDomain string + var vDomains string + var vNotBefore time.Time + var vNotAfter time.Time + var vHoursLeft int32 + var vDaysLeft int32 + var vValidity bool + + if certX509 != nil { + vDomain = certX509.Subject.CommonName + vDomains = strings.Join(certX509.DNSNames, ";") + vNotBefore = certX509.NotBefore + vNotAfter = certX509.NotAfter + vHoursLeft = int32(math.Floor(time.Until(certX509.NotAfter).Hours())) + vDaysLeft = int32(math.Floor(time.Until(certX509.NotAfter).Hours() / 24)) + vValidity = certX509.NotAfter.After(time.Now()) + } + + execRes.AddVariable(stateVarKeyCertificateDomain, vDomain, "string") + execRes.AddVariable(stateVarKeyCertificateDomains, vDomains, "string") + execRes.AddVariable(stateVarKeyCertificateNotBefore, vNotBefore, "datetime") + execRes.AddVariable(stateVarKeyCertificateNotAfter, vNotAfter, "datetime") + execRes.AddVariable(stateVarKeyCertificateHoursLeft, vHoursLeft, "number") + execRes.AddVariable(stateVarKeyCertificateDaysLeft, vDaysLeft, "number") + execRes.AddVariable(stateVarKeyCertificateValidity, vValidity, "boolean") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDomain, vDomain, "string") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDomains, vDomains, "string") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateNotBefore, vNotBefore, "datetime") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateNotAfter, vNotAfter, "datetime") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateHoursLeft, vHoursLeft, "number") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDaysLeft, vDaysLeft, "number") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateValidity, vValidity, "boolean") +} + func newBizMonitorNodeExecutor() NodeExecutor { return &bizMonitorNodeExecutor{ nodeExecutor: nodeExecutor{logger: slog.Default()}, diff --git a/internal/workflow/engine/executor_biznotify.go b/internal/workflow/engine/executor_biznotify.go index 5ba43d72d..7795fcac1 100644 --- a/internal/workflow/engine/executor_biznotify.go +++ b/internal/workflow/engine/executor_biznotify.go @@ -22,8 +22,8 @@ func (ne *bizNotifyNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeEx ne.logger.Info("ready to send notification ...", slog.Any("config", nodeCfg)) // 检测是否可以跳过本次执行 - if skippable := ne.checkCanSkip(execCtx); skippable { - ne.logger.Info(fmt.Sprintf("skip this notification, because all the previous nodes have been skipped")) + if skippable, reason := ne.checkCanSkip(execCtx); skippable { + ne.logger.Info(fmt.Sprintf("skip this application, because %s", reason)) return execRes, nil } @@ -37,10 +37,8 @@ func (ne *bizNotifyNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeEx } } - // 初始化通知器 - notifyClient := notify.NewClient(notify.WithLogger(ne.logger)) - // 推送通知 + notifier := notify.NewClient(notify.WithLogger(ne.logger)) notifyReq := ¬ify.SendNotificationRequest{ Provider: nodeCfg.Provider, ProviderAccessConfig: providerAccessConfig, @@ -48,7 +46,7 @@ func (ne *bizNotifyNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeEx Subject: nodeCfg.Subject, Message: nodeCfg.Message, } - if _, err := notifyClient.SendNotification(execCtx.ctx, notifyReq); err != nil { + if _, err := notifier.SendNotification(execCtx.ctx, notifyReq); err != nil { ne.logger.Warn("could not send notification") return execRes, err } @@ -57,10 +55,10 @@ func (ne *bizNotifyNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeEx return execRes, nil } -func (ne *bizNotifyNodeExecutor) checkCanSkip(execCtx *NodeExecutionContext) (_skip bool) { +func (ne *bizNotifyNodeExecutor) checkCanSkip(execCtx *NodeExecutionContext) (_skip bool, _reason string) { thisNodeCfg := execCtx.Node.Data.Config.AsBizNotify() if !thisNodeCfg.SkipOnAllPrevSkipped { - return false + return false, "" } var total, skipped int32 @@ -72,7 +70,11 @@ func (ne *bizNotifyNodeExecutor) checkCanSkip(execCtx *NodeExecutionContext) (_s } } } - return total > 0 && skipped == total + if total == 0 || skipped != total { + return false, "" + } + + return true, "all the previous nodes have been skipped" } func newBizNotifyNodeExecutor() NodeExecutor { diff --git a/internal/workflow/engine/executor_bizupload.go b/internal/workflow/engine/executor_bizupload.go index f63510d5b..10d680963 100644 --- a/internal/workflow/engine/executor_bizupload.go +++ b/internal/workflow/engine/executor_bizupload.go @@ -7,6 +7,7 @@ import ( "crypto/tls" "fmt" "log/slog" + "math" "os" "strings" "time" @@ -20,13 +21,18 @@ import ( ) /** - * Result Variables: - * - node.skipped: boolean - * - certificate.validity: boolean - * - certificate.daysLeft: number + * Outputs: + * - ref: "certificate": string * - * Result Outputs: - * - ref: certificate: string + * Variables: + * - "node.skipped": boolean + * - "certificate.domain": string + * - "certificate.domains": string + * - "certificate.notBefore": datetime + * - "certificate.notAfter": datetime + * - "certificate.hoursLeft": number + * - "certificate.daysLeft": number + * - "certificate.validity": boolean */ type bizUploadNodeExecutor struct { nodeExecutor @@ -52,9 +58,8 @@ func (ne *bizUploadNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeEx if err != nil { return execRes, err } else if lastCertificate != nil { - execRes.AddOutput(stateIOTypeRef, "certificate", fmt.Sprintf("%s#%s", domain.CollectionNameCertificate, lastCertificate.Id), "string") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateValidity, time.Now().After(lastCertificate.ValidityNotAfter), "boolean") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDaysLeft, int32(time.Until(lastCertificate.ValidityNotAfter).Hours()/24), "number") + ne.setOuputsOfResult(execCtx, execRes, lastCertificate, false) + ne.setVariablesOfResult(execCtx, execRes, lastCertificate) } // 检测是否可以跳过本次执行 @@ -125,7 +130,7 @@ func (ne *bizUploadNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeEx certX509, err := certcrypto.ParsePEMCertificate([]byte(certPEM)) if err != nil { return execRes, err - } else if time.Now().After(certX509.NotAfter) { + } else if certX509.NotAfter.Before(time.Now()) { ne.logger.Warn(fmt.Sprintf("the uploaded certificate has expired at %s", certX509.NotAfter.UTC().Format(time.RFC3339))) } @@ -179,10 +184,8 @@ func (ne *bizUploadNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeEx } // 节点输出 - execRes.AddOutputWithPersistent(stateIOTypeRef, "certificate", fmt.Sprintf("%s#%s", domain.CollectionNameCertificate, certificate.Id), "string") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyNodeSkipped, false, "boolean") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateValidity, time.Now().After(certificate.ValidityNotAfter), "boolean") - execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDaysLeft, int32(time.Until(certificate.ValidityNotAfter).Hours()/24), "number") + ne.setOuputsOfResult(execCtx, execRes, certificate, true) + ne.setVariablesOfResult(execCtx, execRes, certificate) ne.logger.Info("uploading completed") return execRes, nil @@ -239,6 +242,53 @@ func (ne *bizUploadNodeExecutor) checkCanSkip(execCtx *NodeExecutionContext, las return false, "" } +func (ne *bizUploadNodeExecutor) setOuputsOfResult(execCtx *NodeExecutionContext, execRes *NodeExecutionResult, certificate *domain.Certificate, persistent bool) { + if certificate != nil { + key := "certificate" + value := fmt.Sprintf("%s#%s", domain.CollectionNameCertificate, certificate.Id) + if persistent { + execRes.AddOutputWithPersistent(stateIOTypeRef, key, value, "string") + } else { + execRes.AddOutput(stateIOTypeRef, key, value, "string") + } + } +} + +func (ne *bizUploadNodeExecutor) setVariablesOfResult(execCtx *NodeExecutionContext, execRes *NodeExecutionResult, certificate *domain.Certificate) { + var vDomain string + var vDomains string + var vNotBefore time.Time + var vNotAfter time.Time + var vHoursLeft int32 + var vDaysLeft int32 + var vValidity bool + + if certificate != nil { + vDomain = strings.Split(certificate.SubjectAltNames, ";")[0] + vDomains = certificate.SubjectAltNames + vNotBefore = certificate.ValidityNotBefore + vNotAfter = certificate.ValidityNotAfter + vHoursLeft = int32(math.Floor(time.Until(certificate.ValidityNotAfter).Hours())) + vDaysLeft = int32(math.Floor(time.Until(certificate.ValidityNotAfter).Hours() / 24)) + vValidity = certificate.ValidityNotAfter.After(time.Now()) + } + + execRes.AddVariable(stateVarKeyCertificateDomain, vDomain, "string") + execRes.AddVariable(stateVarKeyCertificateDomains, vDomains, "string") + execRes.AddVariable(stateVarKeyCertificateNotBefore, vNotBefore, "datetime") + execRes.AddVariable(stateVarKeyCertificateNotAfter, vNotAfter, "datetime") + execRes.AddVariable(stateVarKeyCertificateHoursLeft, vHoursLeft, "number") + execRes.AddVariable(stateVarKeyCertificateDaysLeft, vDaysLeft, "number") + execRes.AddVariable(stateVarKeyCertificateValidity, vValidity, "boolean") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDomain, vDomain, "string") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDomains, vDomains, "string") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateNotBefore, vNotBefore, "datetime") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateNotAfter, vNotAfter, "datetime") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateHoursLeft, vHoursLeft, "number") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateDaysLeft, vDaysLeft, "number") + execRes.AddVariableWithScope(execCtx.Node.Id, stateVarKeyCertificateValidity, vValidity, "boolean") +} + func newBizUploadNodeExecutor() NodeExecutor { return &bizUploadNodeExecutor{ nodeExecutor: nodeExecutor{logger: slog.Default()}, diff --git a/internal/workflow/engine/state.go b/internal/workflow/engine/state.go index 5b3b282e7..68577d56c 100644 --- a/internal/workflow/engine/state.go +++ b/internal/workflow/engine/state.go @@ -2,9 +2,12 @@ import ( "fmt" + "regexp" "slices" "strconv" + "strings" "sync" + "time" ) type VariableState struct { @@ -22,6 +25,12 @@ func (s VariableState) ValueString() string { return fmt.Sprintf("%d", s.Value) case "boolean": return strconv.FormatBool(s.Value.(bool)) + case "datetime": + valueAsTime := s.Value.(time.Time) + if valueAsTime.IsZero() { + return "-" + } + return valueAsTime.Format(time.RFC3339) default: return fmt.Sprintf("%v", s.Value) } @@ -40,6 +49,8 @@ type VariableManager interface { TakeScoped(scope string, key string) (*VariableState, bool) Remove(key string) bool RemoveScoped(scope string, key string) bool + + RenderTemplateString(template string) string } type variableManager struct { @@ -142,6 +153,35 @@ func (m *variableManager) RemoveScoped(scope string, key string) bool { return ok } +func (m *variableManager) RenderTemplateString(template string) string { + m.statesMtx.RLock() + defer m.statesMtx.RUnlock() + + replaceFunc := func(match string) string { + mustache := strings.TrimSpace(match[2 : len(match)-2]) + if mustache == "" { + return match + } + + key := mustache[1:] + if key == "" { + return match + } else if key == "$now" { + return time.Now().Format(time.RFC3339) + } + + // TODO: 支持作用域变量 + if state, ok := m.Get(key); ok { + return state.ValueString() + } + + return match + } + + re := regexp.MustCompile(`\{\{\s*(\$[^\s]+)\s*\}\}`) + return re.ReplaceAllStringFunc(template, replaceFunc) +} + func newVariableManager() VariableManager { return &variableManager{ states: make([]VariableState, 0), @@ -276,7 +316,14 @@ const ( ) const ( - stateVarKeyNodeSkipped = "node.skipped" // ValueType: "boolean" - stateVarKeyCertificateValidity = "certificate.validity" // ValueType: "boolean" - stateVarKeyCertificateDaysLeft = "certificate.daysLeft" // ValueType: "number" + stateVarKeyNodeId = "node.id" // ValueType: "string" + stateVarKeyNodeName = "node.name" // ValueType: "string" + stateVarKeyNodeSkipped = "node.skipped" // ValueType: "boolean" + stateVarKeyCertificateDomain = "certificate.domain" // ValueType: "string" + stateVarKeyCertificateDomains = "certificate.domains" // ValueType: "string" + stateVarKeyCertificateNotBefore = "certificate.notBefore" // ValueType: "datetime" + stateVarKeyCertificateNotAfter = "certificate.notAfter" // ValueType: "datetime" + stateVarKeyCertificateHoursLeft = "certificate.hoursLeft" // ValueType: "number" + stateVarKeyCertificateDaysLeft = "certificate.daysLeft" // ValueType: "number" + stateVarKeyCertificateValidity = "certificate.validity" // ValueType: "boolean" )