feat: implement global settings service, and decouple from the repository layer in business services

This commit is contained in:
Fu Diwei
2026-06-29 15:30:28 +08:00
committed by RHQYZ
parent 0fbb3395f4
commit c2c04e3274
17 changed files with 212 additions and 106 deletions
+4 -10
View File
@@ -5,7 +5,7 @@ import (
"fmt"
"github.com/certimate-go/certimate/internal/domain"
"github.com/certimate-go/certimate/internal/repository"
"github.com/certimate-go/certimate/internal/settings"
xmaps "github.com/certimate-go/certimate/pkg/utils/maps"
)
@@ -32,15 +32,9 @@ func CreateACMEConfig(ctx context.Context, options *ACMEConfigOptions) (*ACMECon
providerAccessCfg := options.CAProviderAccessConfig
if provider.String() == "" {
// follow global settings
// TODO: decoupling from the repository layer
settingsRepo := repository.NewSettingsRepository()
settings, _ := settingsRepo.GetByName(ctx, domain.SettingsNameSSLProvider)
if settings != nil {
sslProviderSettings := settings.Content.AsSSLProvider()
provider = sslProviderSettings.Provider
providerAccessCfg = sslProviderSettings.Configs[sslProviderSettings.Provider]
}
globalSettingsForSSLProvider := settings.GetGlobalSettingsForSSLProvider()
provider = globalSettingsForSSLProvider.Provider
providerAccessCfg = globalSettingsForSSLProvider.Configs[globalSettingsForSSLProvider.Provider]
}
if provider.String() == "" {
+6 -18
View File
@@ -4,7 +4,6 @@ import (
"archive/zip"
"bytes"
"context"
"errors"
"fmt"
"log/slog"
"strings"
@@ -15,6 +14,7 @@ import (
"github.com/certimate-go/certimate/internal/certacme"
"github.com/certimate-go/certimate/internal/domain"
"github.com/certimate-go/certimate/internal/domain/dtos"
"github.com/certimate-go/certimate/internal/settings"
xcert "github.com/certimate-go/certimate/pkg/utils/cert"
xcertpfx "github.com/certimate-go/certimate/pkg/utils/cert/pfx"
)
@@ -22,14 +22,12 @@ import (
type CertificateService struct {
acmeAccountRepo acmeAccountRepository
certificateRepo certificateRepository
settingsRepo settingsRepository
}
func NewCertificateService(acmeAccountRepo acmeAccountRepository, certificateRepo certificateRepository, settingsRepo settingsRepository) *CertificateService {
func NewCertificateService(acmeAccountRepo acmeAccountRepository, certificateRepo certificateRepository) *CertificateService {
return &CertificateService{
acmeAccountRepo: acmeAccountRepo,
certificateRepo: certificateRepo,
settingsRepo: settingsRepo,
}
}
@@ -259,20 +257,10 @@ func (s *CertificateService) RevokeCertificate(ctx context.Context, req *dtos.Ce
}
func (s *CertificateService) cleanupExpiredCertificates(ctx context.Context) error {
settings, err := s.settingsRepo.GetByName(ctx, domain.SettingsNamePersistence)
if err != nil {
if errors.Is(err, domain.ErrRecordNotFound) {
return nil
}
app.GetLogger().Error("failed to get persistence settings", slog.Any("error", err))
return err
}
persistenceSettings := settings.Content.AsPersistence()
if persistenceSettings.CertificatesRetentionMaxDays != 0 {
ret, err := s.certificateRepo.DeleteWithExprs(context.Background(),
dbx.NewExp(fmt.Sprintf("validityNotAfter<DATETIME('now', '-%d days')", persistenceSettings.CertificatesRetentionMaxDays)),
globalSettingsForPersistence := settings.GetGlobalSettingsForPersistence()
if globalSettingsForPersistence.CertificatesRetentionMaxDays != 0 {
ret, err := s.certificateRepo.DeleteWithExprs(ctx,
dbx.NewExp(fmt.Sprintf("validityNotAfter<DATETIME('now', '-%d days')", globalSettingsForPersistence.CertificatesRetentionMaxDays)),
)
if err != nil {
app.GetLogger().Error("failed to delete expired certificates", slog.Any("error", err))
-4
View File
@@ -17,7 +17,3 @@ type certificateRepository interface {
Save(ctx context.Context, certificate *domain.Certificate) (*domain.Certificate, error)
DeleteWithExprs(ctx context.Context, exprs ...dbx.Expression) (int, error)
}
type settingsRepository interface {
GetByName(ctx context.Context, name string) (*domain.Settings, error)
}
+2 -3
View File
@@ -26,11 +26,10 @@ func BindRouter(router *router.Router[*core.RequestEvent]) {
workflowRunRepo := repository.NewWorkflowRunRepository()
acmeAccountRepo := repository.NewACMEAccountRepository()
certificateRepo := repository.NewCertificateRepository()
settingsRepo := repository.NewSettingsRepository()
statisticsRepo := repository.NewStatisticsRepository()
certificateSvc = certificate.NewCertificateService(acmeAccountRepo, certificateRepo, settingsRepo)
workflowSvc = workflow.NewWorkflowService(workflowRepo, workflowRunRepo, settingsRepo)
certificateSvc = certificate.NewCertificateService(acmeAccountRepo, certificateRepo)
workflowSvc = workflow.NewWorkflowService(workflowRepo, workflowRunRepo)
statisticsSvc = statistics.NewStatisticsService(statisticsRepo)
notifySvc = notify.NewNotifyService(accessRepo)
+2 -3
View File
@@ -14,10 +14,9 @@ func Setup() {
workflowRunRepo := repository.NewWorkflowRunRepository()
acmeAccountRepo := repository.NewACMEAccountRepository()
certificateRepo := repository.NewCertificateRepository()
settingsRepo := repository.NewSettingsRepository()
workflowSvc := workflow.NewWorkflowService(workflowRepo, workflowRunRepo, settingsRepo)
certificateSvc := certificate.NewCertificateService(acmeAccountRepo, certificateRepo, settingsRepo)
workflowSvc := workflow.NewWorkflowService(workflowRepo, workflowRunRepo)
certificateSvc := certificate.NewCertificateService(acmeAccountRepo, certificateRepo)
if err := initWorkflowScheduler(workflowSvc); err != nil {
app.GetLogger().Error("failed to init workflow scheduler", slog.Any("error", err))
+71
View File
@@ -0,0 +1,71 @@
package settings
import (
"context"
"github.com/pocketbase/pocketbase/core"
"github.com/certimate-go/certimate/internal/app"
"github.com/certimate-go/certimate/internal/domain"
)
func registerSettingsRecordEvents() {
pb := app.GetApp()
pb.OnRecordCreateRequest(domain.CollectionNameSettings).BindFunc(func(e *core.RecordRequestEvent) error {
if err := e.Next(); err != nil {
return err
}
if err := onSettingsRecordCreateOrUpdate(e.Request.Context(), e.App, e.Record); err != nil {
app.GetLogger().Error(err.Error())
return err
}
return nil
})
pb.OnRecordUpdateRequest(domain.CollectionNameSettings).BindFunc(func(e *core.RecordRequestEvent) error {
if err := e.Next(); err != nil {
return err
}
if err := onSettingsRecordCreateOrUpdate(e.Request.Context(), e.App, e.Record); err != nil {
app.GetLogger().Error(err.Error())
return err
}
return nil
})
pb.OnRecordDeleteRequest(domain.CollectionNameSettings).BindFunc(func(e *core.RecordRequestEvent) error {
if err := e.Next(); err != nil {
return err
}
if err := onSettingsRecordDelete(e.Request.Context(), e.App, e.Record); err != nil {
app.GetLogger().Error(err.Error())
return err
}
return nil
})
}
func onSettingsRecordCreateOrUpdate(_ context.Context, pb core.App, record *core.Record) error {
sn := record.GetString("name")
if sn != "" {
content := make(domain.SettingsContent)
record.UnmarshalJSONField("content", &content)
pb.Store().Set(buildPbStoreKey(sn), content)
}
return nil
}
func onSettingsRecordDelete(_ context.Context, pb core.App, record *core.Record) error {
sn := record.GetString("name")
if sn != "" {
pb.Store().Remove(buildPbStoreKey(sn))
}
return nil
}
+30
View File
@@ -0,0 +1,30 @@
package settings
import (
"github.com/certimate-go/certimate/internal/app"
)
func initPbSettings() {
pb := app.GetApp()
settings := pb.Settings()
changed := false
if settings.Meta.AppName != app.AppName {
settings.Meta.AppName = app.AppName
changed = true
}
if settings.Batch.Enabled != true {
settings.Batch.Enabled = true
settings.Batch.MaxRequests = 1000
settings.Batch.Timeout = 30
changed = true
}
if changed {
if err := pb.Save(settings); err != nil {
panic(err)
}
}
}
+47
View File
@@ -0,0 +1,47 @@
package settings
import (
"context"
"fmt"
"strings"
"github.com/certimate-go/certimate/internal/app"
"github.com/certimate-go/certimate/internal/domain"
"github.com/certimate-go/certimate/internal/repository"
)
func GetGlobalSettingsForSSLProvider() domain.SettingsContentForSSLProvider {
pb := app.GetApp()
name := domain.SettingsNameSSLProvider
content := pb.Store().Get(buildPbStoreKey(name))
if content == nil {
content = domain.SettingsContent{}
}
return *(content.(domain.SettingsContent)).AsSSLProvider()
}
func GetGlobalSettingsForPersistence() domain.SettingsContentForPersistence {
pb := app.GetApp()
name := domain.SettingsNamePersistence
content := pb.Store().Get(buildPbStoreKey(name))
if content == nil {
content = domain.SettingsContent{}
}
return *(content.(domain.SettingsContent)).AsPersistence()
}
func registerSettingsStoreByName(settingsName string) error {
settingsRepo := repository.NewSettingsRepository()
settings, err := settingsRepo.GetByName(context.Background(), settingsName)
if err != nil {
return err
}
pb := app.GetApp()
pb.Store().Set(buildPbStoreKey(settingsName), settings.Content)
return nil
}
func buildPbStoreKey(settingsName string) string {
return fmt.Sprintf("%s|settings|%s", strings.ToLower(app.AppName), strings.ToLower(settingsName))
}
+13
View File
@@ -0,0 +1,13 @@
package settings
import (
"github.com/certimate-go/certimate/internal/domain"
)
func Setup() {
initPbSettings()
registerSettingsStoreByName(domain.SettingsNameSSLProvider)
registerSettingsStoreByName(domain.SettingsNamePersistence)
registerSettingsRecordEvents()
}
-4
View File
@@ -20,7 +20,3 @@ type workflowOutputRepository interface {
GetByWorkflowIdAndNodeId(ctx context.Context, workflowId string, workflowNodeId string) (*domain.WorkflowOutput, error)
Save(ctx context.Context, workflowOutput *domain.WorkflowOutput) (*domain.WorkflowOutput, error)
}
type settingsRepository interface {
GetByName(ctx context.Context, name string) (*domain.Settings, error)
}
@@ -18,6 +18,7 @@ import (
"github.com/certimate-go/certimate/internal/certacme"
"github.com/certimate-go/certimate/internal/domain"
"github.com/certimate-go/certimate/internal/repository"
"github.com/certimate-go/certimate/internal/settings"
"github.com/certimate-go/certimate/internal/tools/mproc"
xcert "github.com/certimate-go/certimate/pkg/utils/cert"
xcertkey "github.com/certimate-go/certimate/pkg/utils/cert/key"
@@ -53,7 +54,6 @@ const (
type bizApplyNodeExecutor struct {
nodeExecutor
settingsRepo settingsRepository
accessRepo accessRepository
certificateRepo certificateRepository
wfoutputRepo workflowOutputRepository
@@ -380,12 +380,9 @@ func (ne *bizApplyNodeExecutor) execObtainCertificate(execCtx *NodeExecutionCont
// 构造证书申请时所需的 lego 配置项
legoCertifierCfg := &lego.NewConfig(nil).Certificate
settings, _ := ne.settingsRepo.GetByName(execCtx.Context(), domain.SettingsNameSSLProvider)
if settings != nil {
sslProviderSettings := settings.Content.AsSSLProvider()
if sslProviderSettings.Timeout > 0 {
legoCertifierCfg.Timeout = time.Duration(sslProviderSettings.Timeout) * time.Second
}
globalSettingsForPersistence := settings.GetGlobalSettingsForSSLProvider()
if globalSettingsForPersistence.Timeout > 0 {
legoCertifierCfg.Timeout = time.Duration(globalSettingsForPersistence.Timeout) * time.Second
}
// 如果启用多进程模式,发送指令
@@ -494,7 +491,6 @@ func (ne *bizApplyNodeExecutor) setVariablesOfResult(execCtx *NodeExecutionConte
func newBizApplyNodeExecutor() NodeExecutor {
return &bizApplyNodeExecutor{
nodeExecutor: nodeExecutor{logger: slog.Default()},
settingsRepo: repository.NewSettingsRepository(),
accessRepo: repository.NewAccessRepository(),
certificateRepo: repository.NewCertificateRepository(),
wfoutputRepo: repository.NewWorkflowOutputRepository(),
+7 -9
View File
@@ -2,7 +2,6 @@ package workflow
import (
"context"
"fmt"
"github.com/pocketbase/pocketbase/core"
@@ -17,7 +16,7 @@ func registerWorkflowRecordEvents() {
return err
}
if err := onWorkflowRecordCreateOrUpdate(e.Request.Context(), e.Record); err != nil {
if err := onWorkflowRecordCreateOrUpdate(e.Request.Context(), e.App, e.Record); err != nil {
app.GetLogger().Error(err.Error())
return err
}
@@ -29,7 +28,7 @@ func registerWorkflowRecordEvents() {
return err
}
if err := onWorkflowRecordCreateOrUpdate(e.Request.Context(), e.Record); err != nil {
if err := onWorkflowRecordCreateOrUpdate(e.Request.Context(), e.App, e.Record); err != nil {
app.GetLogger().Error(err.Error())
return err
}
@@ -41,7 +40,7 @@ func registerWorkflowRecordEvents() {
return err
}
if err := onWorkflowRecordDelete(e.Request.Context(), e.Record); err != nil {
if err := onWorkflowRecordDelete(e.Request.Context(), e.App, e.Record); err != nil {
app.GetLogger().Error(err.Error())
return err
}
@@ -50,7 +49,7 @@ func registerWorkflowRecordEvents() {
})
}
func onWorkflowRecordCreateOrUpdate(_ context.Context, record *core.Record) error {
func onWorkflowRecordCreateOrUpdate(_ context.Context, _ core.App, record *core.Record) error {
scheduler := app.GetScheduler()
// 向数据库插入/更新时,同时更新定时任务
@@ -60,7 +59,7 @@ func onWorkflowRecordCreateOrUpdate(_ context.Context, record *core.Record) erro
// 如果非定时触发或未启用,移除定时任务
if !enabled || trigger != domain.WorkflowTriggerTypeScheduled.String() {
scheduler.Remove(fmt.Sprintf("workflow#%s", record.Id))
scheduler.Remove(buildPbJobKey(record.Id))
return nil
}
@@ -72,12 +71,11 @@ func onWorkflowRecordCreateOrUpdate(_ context.Context, record *core.Record) erro
return nil
}
func onWorkflowRecordDelete(_ context.Context, record *core.Record) error {
func onWorkflowRecordDelete(_ context.Context, _ core.App, record *core.Record) error {
scheduler := app.GetScheduler()
// 从数据库删除时,同时移除定时任务
jobId := fmt.Sprintf("workflow#%s", record.Id)
scheduler.Remove(jobId)
scheduler.Remove(buildPbJobKey(record.Id))
return nil
}
+7 -3
View File
@@ -16,13 +16,13 @@ import (
func registerWorkflowJob(workflowSrv *WorkflowService, workflowId string, triggerCron string) error {
scheduler := app.GetScheduler()
jobId := fmt.Sprintf("workflow#%s", workflowId)
job, _ := lo.Find(scheduler.Jobs(), func(j *cron.Job) bool { return j.Id() == jobId })
jobKey := buildPbJobKey(workflowId)
job, _ := lo.Find(scheduler.Jobs(), func(j *cron.Job) bool { return j.Id() == jobKey })
if job != nil && job.Expression() == triggerCron {
return nil
}
err := scheduler.Add(jobId, triggerCron, func() {
err := scheduler.Add(jobKey, triggerCron, func() {
app.GetLogger().Info(fmt.Sprintf("workflow #%s is triggered ...", workflowId))
_, err := workflowSrv.StartRun(context.Background(), &dtos.WorkflowStartRunReq{
@@ -41,3 +41,7 @@ func registerWorkflowJob(workflowSrv *WorkflowService, workflowId string, trigge
app.GetLogger().Info(fmt.Sprintf("registered cron job for workflow #%s", workflowId), slog.String("cron", triggerCron))
return nil
}
func buildPbJobKey(workflowId string) string {
return fmt.Sprintf("workflow#%s", workflowId)
}
+6 -18
View File
@@ -12,6 +12,7 @@ import (
"github.com/certimate-go/certimate/internal/app"
"github.com/certimate-go/certimate/internal/domain"
"github.com/certimate-go/certimate/internal/domain/dtos"
"github.com/certimate-go/certimate/internal/settings"
"github.com/certimate-go/certimate/internal/workflow/dispatcher"
)
@@ -20,16 +21,14 @@ type WorkflowService struct {
workflowRepo workflowRepository
workflowRunRepo workflowRunRepository
settingsRepo settingsRepository
}
func NewWorkflowService(workflowRepo workflowRepository, workflowRunRepo workflowRunRepository, settingsRepo settingsRepository) *WorkflowService {
func NewWorkflowService(workflowRepo workflowRepository, workflowRunRepo workflowRunRepository) *WorkflowService {
srv := &WorkflowService{
dispatcher: dispatcher.GetSingletonDispatcher(),
workflowRepo: workflowRepo,
workflowRunRepo: workflowRunRepo,
settingsRepo: settingsRepo,
}
return srv
}
@@ -136,23 +135,12 @@ func (s *WorkflowService) Shutdown(ctx context.Context) {
}
func (s *WorkflowService) cleanupHistoryRuns(ctx context.Context) error {
settings, err := s.settingsRepo.GetByName(ctx, domain.SettingsNamePersistence)
if err != nil {
if errors.Is(err, domain.ErrRecordNotFound) {
return nil
}
app.GetLogger().Error("failed to get persistence settings", slog.Any("error", err))
return err
}
persistenceSettings := settings.Content.AsPersistence()
if persistenceSettings.WorkflowRunsRetentionMaxDays != 0 {
ret, err := s.workflowRunRepo.DeleteWithExprs(
ctx,
globalSettingsForPersistence := settings.GetGlobalSettingsForPersistence()
if globalSettingsForPersistence.WorkflowRunsRetentionMaxDays != 0 {
ret, err := s.workflowRunRepo.DeleteWithExprs(ctx,
dbx.NewExp(fmt.Sprintf("status!='%s'", domain.WorkflowRunStatusTypePending)),
dbx.NewExp(fmt.Sprintf("status!='%s'", domain.WorkflowRunStatusTypeProcessing)),
dbx.NewExp(fmt.Sprintf("endedAt<DATETIME('now', '-%d days')", persistenceSettings.WorkflowRunsRetentionMaxDays)),
dbx.NewExp(fmt.Sprintf("endedAt<DATETIME('now', '-%d days')", globalSettingsForPersistence.WorkflowRunsRetentionMaxDays)),
)
if err != nil {
app.GetLogger().Error("failed to delete workflow history runs", slog.Any("error", err))
-4
View File
@@ -20,7 +20,3 @@ type workflowRunRepository interface {
SaveWithCascading(ctx context.Context, workflowRun *domain.WorkflowRun) (*domain.WorkflowRun, error)
DeleteWithExprs(ctx context.Context, exprs ...dbx.Expression) (int, error)
}
type settingsRepository interface {
GetByName(ctx context.Context, name string) (*domain.Settings, error)
}
-1
View File
@@ -16,7 +16,6 @@ func thisSvcInst() *WorkflowService {
thisSvc = NewWorkflowService(
repository.NewWorkflowRepository(),
repository.NewWorkflowRunRepository(),
repository.NewSettingsRepository(),
)
})
return thisSvc
+13 -21
View File
@@ -17,6 +17,7 @@ import (
"github.com/certimate-go/certimate/internal/app"
"github.com/certimate-go/certimate/internal/rest/routes"
"github.com/certimate-go/certimate/internal/scheduler"
"github.com/certimate-go/certimate/internal/settings"
"github.com/certimate-go/certimate/internal/workflow"
"github.com/certimate-go/certimate/ui"
@@ -48,11 +49,13 @@ func main() {
pflag.StringVar(&flagHttp, "http", "127.0.0.1:8090", "HTTP server address")
pflag.Parse()
pb.OnServe().BindFunc(func(e *core.ServeEvent) error {
scheduler.Setup()
workflow.Setup()
routes.BindRouter(e.Router)
return e.Next()
pb.OnBootstrap().BindFunc(func(e *core.BootstrapEvent) error {
if err := e.Next(); err != nil {
return err
}
settings.Setup()
return nil
})
pb.OnServe().Bind(&hook.Handler[*core.ServeEvent]{
@@ -66,26 +69,15 @@ func main() {
})
pb.OnServe().BindFunc(func(e *core.ServeEvent) error {
slog.Info("[CERTIMATE] Visit the website: http://" + flagHttp)
return e.Next()
})
scheduler.Setup()
workflow.Setup()
routes.BindRouter(e.Router)
pb.OnBootstrap().BindFunc(func(e *core.BootstrapEvent) error {
err := e.Next()
if err != nil {
if err := e.Next(); err != nil {
return err
}
settings := pb.Settings()
if !settings.Batch.Enabled {
settings.Batch.Enabled = true
settings.Batch.MaxRequests = 1000
settings.Batch.Timeout = 30
if err := pb.Save(settings); err != nil {
return err
}
}
slog.Info("[CERTIMATE] Visit the website: http://" + flagHttp)
return nil
})