package config import ( "context" "encoding/json" "fmt" "sync" "dataengine/common/report/model" "gitea.redpowerfuture.com/red-future/common/db/gfdb" "github.com/gogf/gf/v2/frame/g" ) // ConfigLoader 配置加载器 type ConfigLoader struct { mu sync.RWMutex // 缓存 businessCache map[string]*model.BusinessConfig reportCache map[string]*model.ReportConfig fieldCache map[string][]model.FieldConfig extractCache map[string][]model.ExtractConfig } var ( defaultLoader *ConfigLoader once sync.Once ) // GetLoader 获取配置加载器单例 func GetLoader() *ConfigLoader { once.Do(func() { defaultLoader = &ConfigLoader{ businessCache: make(map[string]*model.BusinessConfig), reportCache: make(map[string]*model.ReportConfig), fieldCache: make(map[string][]model.FieldConfig), extractCache: make(map[string][]model.ExtractConfig), } }) return defaultLoader } // GetBusiness 获取业务配置 func (l *ConfigLoader) GetBusiness(ctx context.Context, businessCode string) (*model.BusinessConfig, error) { l.mu.RLock() if biz, ok := l.businessCache[businessCode]; ok { l.mu.RUnlock() return biz, nil } l.mu.RUnlock() var biz model.BusinessConfig r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_business_config WHERE business_code = $1 AND deleted_at IS NULL LIMIT 1", businessCode) if err != nil { return nil, fmt.Errorf("查询业务配置失败: %w", err) } if r.IsEmpty() { return nil, fmt.Errorf("业务配置不存在: %s", businessCode) } if err = r[0].Struct(&biz); err != nil { return nil, err } if biz.Config == nil { biz.Config = make(map[string]interface{}) } l.mu.Lock() l.businessCache[businessCode] = &biz l.mu.Unlock() return &biz, nil } // GetReport 获取报表配置 func (l *ConfigLoader) GetReport(ctx context.Context, businessCode, reportCode string) (*model.ReportConfig, error) { key := businessCode + ":" + reportCode l.mu.RLock() if rpt, ok := l.reportCache[key]; ok { l.mu.RUnlock() return rpt, nil } l.mu.RUnlock() var rpt model.ReportConfig r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_report_config WHERE business_code = $1 AND report_code = $2 AND deleted_at IS NULL LIMIT 1", businessCode, reportCode) if err != nil { return nil, fmt.Errorf("查询报表配置失败: %w", err) } if r.IsEmpty() { return nil, fmt.Errorf("报表配置不存在: %s/%s", businessCode, reportCode) } if err = r[0].Struct(&rpt); err != nil { return nil, err } if rpt.PrimaryKeys == nil { rpt.PrimaryKeys = []string{"id"} } if rpt.ConflictKeys == nil { rpt.ConflictKeys = []string{"stat_date"} } if rpt.Config == nil { rpt.Config = make(map[string]interface{}) } l.mu.Lock() l.reportCache[key] = &rpt l.mu.Unlock() return &rpt, nil } // GetFields 获取报表字段配置 func (l *ConfigLoader) GetFields(ctx context.Context, businessCode, reportCode string) ([]model.FieldConfig, error) { key := businessCode + ":" + reportCode l.mu.RLock() if fields, ok := l.fieldCache[key]; ok { l.mu.RUnlock() return fields, nil } l.mu.RUnlock() r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_field_config WHERE business_code = $1 AND report_code = $2 AND deleted_at IS NULL ORDER BY sort_order ASC", businessCode, reportCode) if err != nil { return nil, err } var fields []model.FieldConfig for _, record := range r { var f model.FieldConfig if err := record.Struct(&f); err != nil { return nil, err } if f.ValidAggregates == nil { f.ValidAggregates = []string{} } if f.FilterOperators == nil { f.FilterOperators = []string{"=", "!=", ">", "<", ">=", "<=", "IN", "LIKE", "BETWEEN"} } fields = append(fields, f) } l.mu.Lock() l.fieldCache[key] = fields l.mu.Unlock() return fields, nil } // GetAllFields 获取报表全部字段(含 INACTIVE,不含已删除,用于编辑回显) func (l *ConfigLoader) GetAllFields(ctx context.Context, businessCode, reportCode string) ([]model.FieldConfig, error) { r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_field_config WHERE business_code = $1 AND report_code = $2 AND deleted_at IS NULL ORDER BY sort_order ASC", businessCode, reportCode) if err != nil { return nil, err } var fields []model.FieldConfig for _, record := range r { var f model.FieldConfig if err := record.Struct(&f); err != nil { return nil, err } if f.ValidAggregates == nil { f.ValidAggregates = []string{} } if f.FilterOperators == nil { f.FilterOperators = []string{"=", "!=", ">", "<", ">=", "<=", "IN", "LIKE", "BETWEEN"} } fields = append(fields, f) } return fields, nil } // GetFieldMap 获取字段配置Map func (l *ConfigLoader) GetFieldMap(ctx context.Context, businessCode, reportCode string) (map[string]*model.FieldConfig, error) { fields, err := l.GetFields(ctx, businessCode, reportCode) if err != nil { return nil, err } fieldMap := make(map[string]*model.FieldConfig) for i := range fields { fieldMap[fields[i].FieldCode] = &fields[i] } return fieldMap, nil } // GetExtractConfigs 获取抽取配置 func (l *ConfigLoader) GetExtractConfigs(ctx context.Context, businessCode, reportCode string) ([]model.ExtractConfig, error) { key := businessCode + ":" + reportCode l.mu.RLock() if configs, ok := l.extractCache[key]; ok { l.mu.RUnlock() return configs, nil } l.mu.RUnlock() r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_extract_config WHERE business_code = $1 AND report_code = $2 AND deleted_at IS NULL ORDER BY id ASC", businessCode, reportCode) if err != nil { return nil, err } var configs []model.ExtractConfig for _, record := range r { var ec model.ExtractConfig if err := record.Struct(&ec); err != nil { return nil, err } if ec.JoinConfigs == nil { ec.JoinConfigs = []model.JoinConfig{} } if ec.FieldMappings == nil { ec.FieldMappings = []model.FieldMapping{} } if ec.TransformRules == nil { ec.TransformRules = []model.TransformRule{} } if ec.GroupByFields == nil { ec.GroupByFields = []string{} } if ec.ExtractMode == "" { ec.ExtractMode = model.ExtractModeDirect } configs = append(configs, ec) } l.mu.Lock() l.extractCache[key] = configs l.mu.Unlock() return configs, nil } // GetExtractLog 获取抽取记录 func (l *ConfigLoader) GetExtractLog(ctx context.Context, businessCode, reportCode, extractCode, statDate string) (*model.ExtractLog, error) { var log model.ExtractLog r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_extract_log WHERE business_code = $1 AND report_code = $2 AND extract_code = $3 AND stat_date = $4 LIMIT 1", businessCode, reportCode, extractCode, statDate) if err != nil { return nil, err } if r.IsEmpty() { return nil, nil } if err = r[0].Struct(&log); err != nil { return nil, err } return &log, nil } // CreateExtractLog 创建抽取记录 func (l *ConfigLoader) CreateExtractLog(ctx context.Context, log *model.ExtractLog) error { data, _ := json.Marshal(log) var m map[string]interface{} json.Unmarshal(data, &m) delete(m, "id") _, err := gfdb.DB(ctx).Model(ctx, "report_extract_log").Data(m).Save() return err } // UpdateExtractLog 更新抽取记录 func (l *ConfigLoader) UpdateExtractLog(ctx context.Context, log *model.ExtractLog) error { data, _ := json.Marshal(log) var m map[string]interface{} json.Unmarshal(data, &m) delete(m, "id") delete(m, "created_at") _, err := gfdb.DB(ctx).Model(ctx, "report_extract_log"). Where("business_code", log.BusinessCode). Where("report_code", log.ReportCode). Where("extract_code", log.ExtractCode). Where("stat_date", log.StatDate). Data(m). Update() return err } // InvalidateCache 失效缓存 func (l *ConfigLoader) InvalidateCache(businessCode, reportCode string) { l.mu.Lock() delete(l.businessCache, businessCode) delete(l.reportCache, businessCode+":"+reportCode) delete(l.fieldCache, businessCode+":"+reportCode) delete(l.extractCache, businessCode+":"+reportCode) l.mu.Unlock() } // InvalidateBusinessCache 只失效业务缓存(不影响报表/字段) func (l *ConfigLoader) InvalidateBusinessCache(businessCode string) { l.mu.Lock() delete(l.businessCache, businessCode) l.mu.Unlock() } // ============================================================ // CRUD: BusinessConfig // ============================================================ // CreateBusiness 创建业务配置 func (l *ConfigLoader) CreateBusiness(ctx context.Context, biz *model.BusinessConfig) (int64, error) { data, _ := json.Marshal(biz) var m map[string]interface{} json.Unmarshal(data, &m) delete(m, "id") delete(m, "created_at") delete(m, "updated_at") delete(m, "deleted_at") result, err := gfdb.DB(ctx).Model(ctx, "report_business_config").Data(m).Insert() if err != nil { return 0, fmt.Errorf("创建业务配置失败: %w", err) } id, _ := result.LastInsertId() l.InvalidateBusinessCache(biz.BusinessCode) return id, nil } // UpdateBusiness 更新业务配置 func (l *ConfigLoader) UpdateBusiness(ctx context.Context, biz *model.BusinessConfig) error { data, _ := json.Marshal(biz) var m map[string]interface{} json.Unmarshal(data, &m) delete(m, "id") delete(m, "created_at") delete(m, "deleted_at") _, err := gfdb.DB(ctx).Model(ctx, "report_business_config"). Where("id", biz.Id). Data(m). Update() if err != nil { return fmt.Errorf("更新业务配置失败: %w", err) } l.InvalidateBusinessCache(biz.BusinessCode) return nil } // DeleteBusiness 删除业务配置(软删除) func (l *ConfigLoader) DeleteBusiness(ctx context.Context, id int64, businessCode string) error { _, err := gfdb.DB(ctx).Model(ctx, "report_business_config"). Where("id", id). Data(map[string]interface{}{ "status": model.StatusInactive, "deleted_at": "NOW()", }). Update() if err != nil { return fmt.Errorf("删除业务配置失败: %w", err) } l.InvalidateBusinessCache(businessCode) return nil } // GetBusinessByID 根据ID获取业务配置 func (l *ConfigLoader) GetBusinessByID(ctx context.Context, id int64) (*model.BusinessConfig, error) { var biz model.BusinessConfig r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_business_config WHERE id = $1 LIMIT 1", id) if err != nil { g.Log().Infof(ctx, "[GetBusinessByID] id=%d, err=%v", id, err) return nil, fmt.Errorf("查询业务配置失败: %w", err) } if r.IsEmpty() { return nil, fmt.Errorf("业务配置不存在: id=%d", id) } if err = r[0].Struct(&biz); err != nil { return nil, err } g.Log().Infof(ctx, "[GetBusinessByID] id=%d, biz.Id=%d, biz.BusinessCode=%s", id, biz.Id, biz.BusinessCode) return &biz, nil } // ============================================================ // CRUD: ReportConfig // ============================================================ // CreateReport 创建报表配置 func (l *ConfigLoader) CreateReport(ctx context.Context, rpt *model.ReportConfig) (int64, error) { data, _ := json.Marshal(rpt) var m map[string]interface{} json.Unmarshal(data, &m) delete(m, "id") delete(m, "created_at") delete(m, "updated_at") delete(m, "deleted_at") result, err := gfdb.DB(ctx).Model(ctx, "report_report_config").Data(m).Insert() if err != nil { return 0, fmt.Errorf("创建报表配置失败: %w", err) } id, _ := result.LastInsertId() l.InvalidateCache(rpt.BusinessCode, rpt.ReportCode) return id, nil } // UpdateReport 更新报表配置 func (l *ConfigLoader) UpdateReport(ctx context.Context, rpt *model.ReportConfig) error { data, _ := json.Marshal(rpt) var m map[string]interface{} json.Unmarshal(data, &m) delete(m, "id") delete(m, "created_at") delete(m, "deleted_at") _, err := gfdb.DB(ctx).Model(ctx, "report_report_config"). Where("id", rpt.Id). Data(m). Update() if err != nil { return fmt.Errorf("更新报表配置失败: %w", err) } l.InvalidateCache(rpt.BusinessCode, rpt.ReportCode) return nil } // DeleteReport 删除报表配置(软删除) func (l *ConfigLoader) DeleteReport(ctx context.Context, id int64, businessCode, reportCode string) error { _, err := gfdb.DB(ctx).Model(ctx, "report_report_config"). Where("id", id). Data(map[string]interface{}{ "status": model.StatusInactive, "deleted_at": "NOW()", }). Update() if err != nil { return fmt.Errorf("删除报表配置失败: %w", err) } l.InvalidateCache(businessCode, reportCode) return nil } // GetReportByID 根据ID获取报表配置 func (l *ConfigLoader) GetReportByID(ctx context.Context, id int64) (*model.ReportConfig, error) { var rpt model.ReportConfig r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_report_config WHERE id = $1 LIMIT 1", id) if err != nil { return nil, fmt.Errorf("查询报表配置失败: %w", err) } if r.IsEmpty() { return nil, fmt.Errorf("报表配置不存在: id=%d", id) } if err = r[0].Struct(&rpt); err != nil { return nil, err } return &rpt, nil } // ============================================================ // CRUD: FieldConfig // ============================================================ // CreateField 创建字段配置 func (l *ConfigLoader) CreateField(ctx context.Context, field *model.FieldConfig) (int64, error) { data, _ := json.Marshal(field) var m map[string]interface{} json.Unmarshal(data, &m) delete(m, "id") delete(m, "created_at") delete(m, "updated_at") delete(m, "deleted_at") result, err := gfdb.DB(ctx).Model(ctx, "report_field_config").Data(m).Insert() if err != nil { return 0, fmt.Errorf("创建字段配置失败: %w", err) } id, _ := result.LastInsertId() l.InvalidateCache(field.BusinessCode, field.ReportCode) return id, nil } // UpdateField 更新字段配置 func (l *ConfigLoader) UpdateField(ctx context.Context, field *model.FieldConfig) error { data, _ := json.Marshal(field) var m map[string]interface{} json.Unmarshal(data, &m) delete(m, "id") delete(m, "created_at") delete(m, "deleted_at") _, err := gfdb.DB(ctx).Model(ctx, "report_field_config"). Where("id", field.Id). Data(m). Update() if err != nil { return fmt.Errorf("更新字段配置失败: %w", err) } l.InvalidateCache(field.BusinessCode, field.ReportCode) return nil } // DeleteField 删除字段配置(软删除) func (l *ConfigLoader) DeleteField(ctx context.Context, id int64, businessCode, reportCode string) error { _, err := gfdb.DB(ctx).Model(ctx, "report_field_config"). Where("id", id). Data(map[string]interface{}{ "status": model.StatusInactive, "deleted_at": "NOW()", }). Update() if err != nil { return fmt.Errorf("删除字段配置失败: %w", err) } l.InvalidateCache(businessCode, reportCode) return nil } // GetFieldByID 根据ID获取字段配置 func (l *ConfigLoader) GetFieldByID(ctx context.Context, id int64) (*model.FieldConfig, error) { var field model.FieldConfig r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_field_config WHERE id = $1 LIMIT 1", id) if err != nil { return nil, fmt.Errorf("查询字段配置失败: %w", err) } if r.IsEmpty() { return nil, fmt.Errorf("字段配置不存在: id=%d", id) } if err = r[0].Struct(&field); err != nil { return nil, err } return &field, nil } // ============================================================ // CRUD: ExtractConfig // ============================================================ // CreateExtractConfig 创建抽取配置 func (l *ConfigLoader) CreateExtractConfig(ctx context.Context, ec *model.ExtractConfig) (int64, error) { data, _ := json.Marshal(ec) var m map[string]interface{} json.Unmarshal(data, &m) delete(m, "id") delete(m, "created_at") delete(m, "updated_at") delete(m, "deleted_at") result, err := gfdb.DB(ctx).Model(ctx, "report_extract_config").Data(m).Insert() if err != nil { return 0, fmt.Errorf("创建抽取配置失败: %w", err) } id, _ := result.LastInsertId() l.InvalidateCache(ec.BusinessCode, ec.ReportCode) return id, nil } // UpdateExtractConfig 更新抽取配置 func (l *ConfigLoader) UpdateExtractConfig(ctx context.Context, ec *model.ExtractConfig) error { data, _ := json.Marshal(ec) var m map[string]interface{} json.Unmarshal(data, &m) delete(m, "id") delete(m, "created_at") delete(m, "deleted_at") _, err := gfdb.DB(ctx).Model(ctx, "report_extract_config"). Where("id", ec.Id). Data(m). Update() if err != nil { return fmt.Errorf("更新抽取配置失败: %w", err) } l.InvalidateCache(ec.BusinessCode, ec.ReportCode) return nil } // DeleteExtractConfig 删除抽取配置(软删除) func (l *ConfigLoader) DeleteExtractConfig(ctx context.Context, id int64, businessCode, reportCode string) error { _, err := gfdb.DB(ctx).Model(ctx, "report_extract_config"). Where("id", id). Data(map[string]interface{}{ "status": model.StatusInactive, "deleted_at": "NOW()", }). Update() if err != nil { return fmt.Errorf("删除抽取配置失败: %w", err) } l.InvalidateCache(businessCode, reportCode) return nil } // GetExtractConfigByID 根据ID获取抽取配置 func (l *ConfigLoader) GetExtractConfigByID(ctx context.Context, id int64) (*model.ExtractConfig, error) { var ec model.ExtractConfig r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_extract_config WHERE id = $1 LIMIT 1", id) if err != nil { return nil, fmt.Errorf("查询抽取配置失败: %w", err) } if r.IsEmpty() { return nil, fmt.Errorf("抽取配置不存在: id=%d", id) } if err = r[0].Struct(&ec); err != nil { return nil, err } return &ec, nil } // GetAllBusinesses 获取所有业务配置 func (l *ConfigLoader) GetAllBusinesses(ctx context.Context) ([]model.BusinessConfig, error) { r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_business_config WHERE deleted_at IS NULL ORDER BY id ASC", ) // (removed status filter) if err != nil { return nil, err } var businesses []model.BusinessConfig for _, record := range r { var biz model.BusinessConfig if err := record.Struct(&biz); err != nil { return nil, err } businesses = append(businesses, biz) } return businesses, nil } // GetAllReports 获取所有报表配置 func (l *ConfigLoader) GetAllReports(ctx context.Context, businessCode string) ([]model.ReportConfig, error) { r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_report_config WHERE business_code = $1 AND deleted_at IS NULL ORDER BY id ASC", businessCode) if err != nil { return nil, err } var reports []model.ReportConfig for _, record := range r { var rpt model.ReportConfig if err := record.Struct(&rpt); err != nil { return nil, err } reports = append(reports, rpt) } return reports, nil } // ListBusinesses 分页获取业务列表 func (l *ConfigLoader) ListBusinesses(ctx context.Context, pageNum, pageSize int) ([]model.BusinessConfig, int, error) { total, err := gfdb.DB(ctx).GetAll(ctx, "SELECT COUNT(*) AS cnt FROM report_business_config WHERE deleted_at IS NULL") if err != nil { return nil, 0, err } count := 0 if !total.IsEmpty() { count = total[0]["cnt"].Int() } offset := (pageNum - 1) * pageSize r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_business_config WHERE deleted_at IS NULL ORDER BY id ASC LIMIT $1 OFFSET $2", pageSize, offset) if err != nil { return nil, 0, err } var businesses []model.BusinessConfig for _, record := range r { var biz model.BusinessConfig if err := record.Struct(&biz); err != nil { return nil, 0, err } businesses = append(businesses, biz) } return businesses, count, nil } // ListReports 分页获取报表列表 func (l *ConfigLoader) ListReports(ctx context.Context, businessCode, reportName string, pageNum, pageSize int) ([]model.ReportConfig, int, error) { namePattern := "%" + reportName + "%" total, err := gfdb.DB(ctx).GetAll(ctx, "SELECT COUNT(*) AS cnt FROM report_report_config WHERE business_code = $1 AND report_name ILIKE $2 AND deleted_at IS NULL", businessCode, namePattern) if err != nil { return nil, 0, err } count := 0 if !total.IsEmpty() { count = total[0]["cnt"].Int() } offset := (pageNum - 1) * pageSize r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_report_config WHERE business_code = $1 AND report_name ILIKE $2 AND deleted_at IS NULL ORDER BY id ASC LIMIT $3 OFFSET $4", businessCode, namePattern, pageSize, offset) if err != nil { return nil, 0, err } var reports []model.ReportConfig for _, record := range r { var rpt model.ReportConfig if err := record.Struct(&rpt); err != nil { return nil, 0, err } reports = append(reports, rpt) } return reports, count, nil } // ListExtractConfigs 分页获取抽取配置列表 func (l *ConfigLoader) ListExtractConfigs(ctx context.Context, businessCode, reportCode string, pageNum, pageSize int) ([]model.ExtractConfig, int, error) { total, err := gfdb.DB(ctx).GetAll(ctx, "SELECT COUNT(*) AS cnt FROM report_extract_config WHERE business_code = $1 AND report_code = $2 AND deleted_at IS NULL", businessCode, reportCode) if err != nil { return nil, 0, err } count := 0 if !total.IsEmpty() { count = total[0]["cnt"].Int() } offset := (pageNum - 1) * pageSize r, err := gfdb.DB(ctx).GetAll(ctx, "SELECT * FROM report_extract_config WHERE business_code = $1 AND report_code = $2 AND deleted_at IS NULL ORDER BY id ASC LIMIT $3 OFFSET $4", businessCode, reportCode, pageSize, offset) if err != nil { return nil, 0, err } var configs []model.ExtractConfig for _, record := range r { var ec model.ExtractConfig if err := record.Struct(&ec); err != nil { return nil, 0, err } configs = append(configs, ec) } l.mu.Lock() l.extractCache[businessCode+":"+reportCode] = configs l.mu.Unlock() return configs, count, nil } // GetReportFields 获取报表可用字段(按角色分类) func (l *ConfigLoader) GetReportFields(ctx context.Context, businessCode, reportCode string) (*model.GetReportFieldsResp, error) { fields, err := l.GetFields(ctx, businessCode, reportCode) if err != nil { return nil, err } resp := &model.GetReportFieldsResp{ BusinessCode: businessCode, ReportCode: reportCode, Dimensions: []model.FieldConfig{}, Indicators: []model.FieldConfig{}, Filters: []model.FieldConfig{}, } for _, f := range fields { switch f.FieldRole { case model.RoleDimension: resp.Dimensions = append(resp.Dimensions, f) case model.RoleIndicator: resp.Indicators = append(resp.Indicators, f) case model.RoleFilter, model.RoleFilterOnly: resp.Filters = append(resp.Filters, f) } } return resp, nil }