diff --git a/server/controller/common/metadata/platform.go b/server/controller/common/metadata/platform.go index 3e90a655c4d..6431e74a653 100644 --- a/server/controller/common/metadata/platform.go +++ b/server/controller/common/metadata/platform.go @@ -117,7 +117,7 @@ func (m *Platform) SetDomain(domain metadbmodel.Domain) { m.teamID = domain.TeamID m.LogPrefixes = append(m.LogPrefixes, logger.NewTeamPrefix(domain.TeamID)) } - m.LogPrefixes = append(m.LogPrefixes, NewDomainPrefix(domain.Name)) + m.LogPrefixes = append(m.LogPrefixes, NewDomainPrefix(domain.Name, domain.Lcuuid)) } func (m *Platform) SetSubDomain(subDomain metadbmodel.SubDomain) { @@ -126,7 +126,7 @@ func (m *Platform) SetSubDomain(subDomain metadbmodel.SubDomain) { m.teamID = subDomain.TeamID m.LogPrefixes = append(m.LogPrefixes, logger.NewTeamPrefix(subDomain.TeamID)) } - m.LogPrefixes = append(m.LogPrefixes, NewSubDomainPrefix(subDomain.Name)) + m.LogPrefixes = append(m.LogPrefixes, NewSubDomainPrefix(subDomain.Name, subDomain.Lcuuid)) } func MetadataDomain(domain metadbmodel.Domain) func(*Platform) { @@ -149,11 +149,11 @@ type SubDomainInfo struct { metadbmodel.SubDomain } -func NewDomainPrefix(name string) logger.Prefix { +func NewDomainPrefix(name, lcuuid string) logger.Prefix { if name == "" { return &DomainIDPrefix{0} } - return &DomainNameLogPrefix{name} + return &DomainNameLogPrefix{name, lcuuid} } type DomainIDPrefix struct { @@ -165,21 +165,23 @@ func (p *DomainIDPrefix) Prefix() string { } type DomainNameLogPrefix struct { - Name string + Name string + Lcuuid string } func (p *DomainNameLogPrefix) Prefix() string { - return fmt.Sprintf("[DomainName-%s]", p.Name) + return fmt.Sprintf("[DomainName-%s-%s]", p.Name, p.Lcuuid) } -func NewSubDomainPrefix(name string) logger.Prefix { - return &SubDomainNameLogPrefix{name} +func NewSubDomainPrefix(name, lcuuid string) logger.Prefix { + return &SubDomainNameLogPrefix{name, lcuuid} } type SubDomainNameLogPrefix struct { - Name string + Name string + Lcuuid string } func (p *SubDomainNameLogPrefix) Prefix() string { - return fmt.Sprintf("[SubDomainName-%s]", p.Name) + return fmt.Sprintf("[SubDomainName-%s-%s]", p.Name, p.Lcuuid) } diff --git a/server/controller/recorder/config/config.go b/server/controller/recorder/config/config.go index a6bae2f62c1..b23a35a43d2 100644 --- a/server/controller/recorder/config/config.go +++ b/server/controller/recorder/config/config.go @@ -35,6 +35,7 @@ type RecorderConfig struct { EventCfg eventConfig.Config SelfHealCfg SelfHealConfig `yaml:"self_heal"` TagRecorderSelfHealCfg TagRecorderSelfHealConfig `yaml:"tagrecorder_self_heal"` + SkipSyncIfEmptyCfg SkipSyncIfEmptyConfig `yaml:"skip_sync_if_empty"` } func Get() *RecorderConfig { @@ -51,6 +52,11 @@ type LogDebugConfig struct { ResourceTypes []string `default:"" yaml:"resource_type"` } +type SkipSyncIfEmptyConfig struct { + Enabled bool `default:"false" yaml:"enabled"` + Resources []string `default:"" yaml:"resources"` +} + type SelfHealConfig struct { Enabled bool `default:"true" yaml:"enabled"` Resources []string `default:"" yaml:"resources"` diff --git a/server/controller/recorder/debugger.go b/server/controller/recorder/debugger.go index 332ca4e95a6..479feda823b 100644 --- a/server/controller/recorder/debugger.go +++ b/server/controller/recorder/debugger.go @@ -18,7 +18,9 @@ package recorder import ( + "fmt" "reflect" + "strings" "github.com/deepflowio/deepflow/server/controller/recorder/cache" ) @@ -66,3 +68,43 @@ func (r *Recorder) GetToolMap(domainLcuuid, subDomainLcuuid, field string) map[i } return dataSet.(map[interface{}]interface{}) } + +// GetResourceFieldCountsString returns a string representation of slice and map field counts in the resource object +// This is a more descriptive name than CountString +func GetResourceFieldCountsString(obj interface{}) string { + var parts []string + v := reflect.ValueOf(obj) + t := v.Type() + + for i := 0; i < v.NumField(); i++ { + f := v.Field(i) + name := t.Field(i).Name + + switch f.Kind() { + case reflect.Slice: + parts = append(parts, fmt.Sprintf("%s=%d", name, f.Len())) + case reflect.Map: + parts = append(parts, fmt.Sprintf("%s=%d", name, f.Len())) + for _, key := range f.MapKeys() { + sub := joinSliceFields(f.MapIndex(key)) + parts = append(parts, fmt.Sprintf("%s[%v]=%s", name, key.Interface(), sub)) + } + } + } + return strings.Join(parts, ", ") +} + +// 提取结构体中所有 slice 字段的 Name=Len 拼接 +func joinSliceFields(v reflect.Value) string { + if v.Kind() == reflect.Ptr { + v = v.Elem() + } + var parts []string + t := v.Type() + for i := 0; i < v.NumField(); i++ { + if f := v.Field(i); f.Kind() == reflect.Slice { + parts = append(parts, fmt.Sprintf("%s=%d", t.Field(i).Name, f.Len())) + } + } + return strings.Join(parts, ", ") +} diff --git a/server/controller/recorder/debugger_test.go b/server/controller/recorder/debugger_test.go new file mode 100644 index 00000000000..4607ca22d4c --- /dev/null +++ b/server/controller/recorder/debugger_test.go @@ -0,0 +1,165 @@ +package recorder + +import ( + "strings" + "testing" + "time" + + "github.com/deepflowio/deepflow/server/controller/cloud/model" +) + +// TestGetResourceFieldCountsString tests the GetResourceFieldCountsString function with various scenarios +func TestGetResourceFieldCountsString(t *testing.T) { + // Test case 1: Empty lists and maps (0 values) + t.Run("EmptyValues", func(t *testing.T) { + resource := model.Resource{ + SubDomains: []model.SubDomain{}, + VMs: []model.VM{}, + VPCs: []model.VPC{}, + SubDomainResources: map[string]model.SubDomainResource{}, + } + + result := GetResourceFieldCountsString(resource) + + // Expected result should contain all fields with 0 counts + expectedParts := []string{ + "SubDomains=0", + "VMs=0", + "VPCs=0", + "SubDomainResources=0", + } + + for _, part := range expectedParts { + if !strings.Contains(result, part) { + t.Errorf("Expected to find '%s' in result: %s", part, result) + } + } + + // Ensure no map entries are shown since all maps are empty + if strings.Contains(result, "[") && strings.Contains(result, "]") { + t.Errorf("No map entries should be present for empty maps, got: %s", result) + } + }) + + // Test case 2: Non-empty lists and maps + t.Run("WithValues", func(t *testing.T) { + resource := model.Resource{ + SubDomains: []model.SubDomain{{Lcuuid: "sd1", Name: "subdomain1"}}, + VMs: []model.VM{{Name: "test-vm-1", Lcuuid: "vm1"}, {Name: "test-vm-2", Lcuuid: "vm2"}}, + VPCs: []model.VPC{{Name: "test-vpc-1", Lcuuid: "vpc1"}}, + SubDomainResources: map[string]model.SubDomainResource{"sub1": {}}, + } + + result := GetResourceFieldCountsString(resource) + + // Check that we have the expected counts + expectedContains := []string{ + "SubDomains=1", + "VMs=2", + "VPCs=1", + "SubDomainResources=1", + } + + for _, exp := range expectedContains { + if !strings.Contains(result, exp) { + t.Errorf("Expected to find '%s' in result: %s", exp, result) + } + } + + // Check that the map entry is properly formatted + if !strings.Contains(result, "SubDomainResources[sub1]=") { + t.Errorf("Expected map entry SubDomainResources[sub1]= to be present, got: %s", result) + } + }) + + // Test case 3: Map with flat structure having 0 values and non-zero values + t.Run("MapWithFlatStructure", func(t *testing.T) { + // Test empty SubDomainResource (all slice fields are zero values) + resourceEmpty := model.Resource{ + SubDomainResources: map[string]model.SubDomainResource{ + "empty": {}, // All slice fields in SubDomainResource are zero values (empty slices) + }, + } + + resultEmpty := GetResourceFieldCountsString(resourceEmpty) + + // Should contain the map count + if !strings.Contains(resultEmpty, "SubDomainResources=1") { + t.Errorf("Expected SubDomainResources=1 in result: %s", resultEmpty) + } + + // The 'empty' entry should show all slice counts as 0 + if !strings.Contains(resultEmpty, "SubDomainResources[empty]=") { + t.Errorf("Expected SubDomainResources[empty]= entry in result: %s", resultEmpty) + } + + // Even though SubDomainResource is empty, it should still show all slice fields as 0 + // Since SubDomainResource has many slice fields, the result should contain multiple "=0" entries + zeroCountFound := strings.Contains(resultEmpty, "=0, ") || strings.HasSuffix(resultEmpty, "=0") + if !zeroCountFound { + t.Errorf("Expected to find zero counts in SubDomainResources[empty] entry, got: %s", resultEmpty) + } + + // Test SubDomainResource with non-zero values + resourceWithData := model.Resource{ + SubDomainResources: map[string]model.SubDomainResource{ + "with-data": { // SubDomainResource with some data in slices + Networks: []model.Network{{Lcuuid: "net1", Name: "network1"}, {Lcuuid: "net2", Name: "network2"}}, + Subnets: []model.Subnet{{Lcuuid: "subnet1", Name: "subnet1"}}, + Pods: []model.Pod{{Lcuuid: "pod1", Name: "pod1"}, {Lcuuid: "pod2", Name: "pod2"}, {Lcuuid: "pod3", Name: "pod3"}}, + }, + }, + } + + resultWithData := GetResourceFieldCountsString(resourceWithData) + + // Should contain the map count + if !strings.Contains(resultWithData, "SubDomainResources=1") { + t.Errorf("Expected SubDomainResources=1 in result: %s", resultWithData) + } + + // The 'with-data' entry should show the count of slices inside the SubDomainResource + if !strings.Contains(resultWithData, "SubDomainResources[with-data]=") { + t.Errorf("Expected SubDomainResources[with-data]= entry in result: %s", resultWithData) + } + + // Verify specific counts for the slices + if !strings.Contains(resultWithData, "Networks=2") { + t.Errorf("Expected Networks=2 in the with-data entry, got: %s", resultWithData) + } + if !strings.Contains(resultWithData, "Subnets=1") { + t.Errorf("Expected Subnets=1 in the with-data entry, got: %s", resultWithData) + } + if !strings.Contains(resultWithData, "Pods=3") { + t.Errorf("Expected Pods=3 in the with-data entry, got: %s", resultWithData) + } + }) + + // Test case 4: Resource with only basic fields (no slices or maps) + t.Run("BasicFieldsOnly", func(t *testing.T) { + // Create a resource with basic fields that don't contribute to the count + resource := model.Resource{ + Verified: true, + ErrorState: 0, + ErrorMessage: "", + SyncAt: time.Now(), + SubDomains: []model.SubDomain{}, // Empty slice + SubDomainResources: map[string]model.SubDomainResource{}, // Empty map + } + + result := GetResourceFieldCountsString(resource) + + // Should only show slice and map fields, not basic fields like Verified, ErrorState, etc. + if !strings.Contains(result, "SubDomains=0") { + t.Errorf("Expected SubDomains=0 in result: %s", result) + } + if !strings.Contains(result, "SubDomainResources=0") { + t.Errorf("Expected SubDomainResources=0 in result: %s", result) + } + + // Basic fields shouldn't appear in the result + if strings.Contains(result, "Verified=") { + t.Errorf("Basic fields like Verified should not appear in result: %s", result) + } + }) +} diff --git a/server/controller/recorder/domain.go b/server/controller/recorder/domain.go index ddea3df94b2..bab6be65ef7 100644 --- a/server/controller/recorder/domain.go +++ b/server/controller/recorder/domain.go @@ -28,7 +28,7 @@ import ( cloudmodel "github.com/deepflowio/deepflow/server/controller/cloud/model" "github.com/deepflowio/deepflow/server/controller/common" - mysqlmodel "github.com/deepflowio/deepflow/server/controller/db/metadb/model" + metadbmodel "github.com/deepflowio/deepflow/server/controller/db/metadb/model" "github.com/deepflowio/deepflow/server/controller/recorder/cache" rcommon "github.com/deepflowio/deepflow/server/controller/recorder/common" "github.com/deepflowio/deepflow/server/controller/recorder/config" @@ -74,14 +74,14 @@ func (d *domain) CloseStatsd() { } func (d *domain) Refresh(target string, cloudData cloudmodel.Resource) error { - log.Infof("refresh target: %s", target, d.metadata.LogPrefixes) + log.Infof("refresh target: %s, cloudData count: %s", target, GetResourceFieldCountsString(cloudData), d.metadata.LogPrefixes) switch target { case RefreshTargetDomain: log.Info("refresher started, triggered by ticker/hand", d.metadata.LogPrefixes) if err := d.refreshDomainExcludeSubDomain(cloudData); err != nil { return err } - return d.subDomains.RefreshAll(cloudData.SubDomainResources) + return d.subDomains.RefreshAll(cloudData.SubDomains, cloudData.SubDomainResources) case RefreshTargetSubDomain: log.Info("refresher started, triggered by hand", d.metadata.LogPrefixes) return d.subDomains.RefreshOne(cloudData.SubDomainResources) @@ -132,6 +132,13 @@ func (d *domain) shouldRefresh(cloudData cloudmodel.Resource) error { log.Info("domain has no vms and pods, does nothing", d.metadata.LogPrefixes) return DataMissingError } + // 检查当 SubDomains 为空时,是否需要跳过同步 + if d.metadata.Config.SkipSyncIfEmptyCfg.Enabled && + slices.Contains(d.metadata.Config.SkipSyncIfEmptyCfg.Resources, common.RESOURCE_TYPE_SUB_DOMAIN_EN) && + len(cloudData.SubDomains) == 0 { + log.Info("domain has no SubDomains, does nothing", d.metadata.LogPrefixes) + return DataMissingError + } } else { log.Info("domain is not verified, does nothing", d.metadata.LogPrefixes) return DataNotVerifiedError @@ -285,7 +292,7 @@ func (d *domain) updateSyncedAt(syncAt time.Time) { log.Infof("update domain synced_at: %s", syncAt.Format(common.GO_BIRTHDAY), d.metadata.LogPrefixes) d.fillStatsd(syncAt) - var domain mysqlmodel.Domain + var domain metadbmodel.Domain err := d.metadata.DB.Where("lcuuid = ?", d.metadata.GetDomainLcuuid()).First(&domain).Error if err != nil { log.Errorf("get domain from db failed: %s", err, d.metadata.LogPrefixes) @@ -303,7 +310,7 @@ func (d *domain) fillStatsd(syncAt time.Time) { } func (d *domain) updateStateInfo(cloudData cloudmodel.Resource) { - var domain mysqlmodel.Domain + var domain metadbmodel.Domain err := d.metadata.DB.Where("lcuuid = ?", d.metadata.GetDomainLcuuid()).First(&domain).Error if err != nil { log.Errorf("get domain from db failed: %s", err, d.metadata.LogPrefixes) @@ -314,7 +321,7 @@ func (d *domain) updateStateInfo(cloudData cloudmodel.Resource) { log.Debugf("update domain (%+v)", domain, d.metadata.LogPrefixes) for subDomainLcuuid, subDomainResource := range cloudData.SubDomainResources { - var subDomain mysqlmodel.SubDomain + var subDomain metadbmodel.SubDomain err := d.metadata.DB.Where("lcuuid = ?", subDomainLcuuid).First(&subDomain).Error if err != nil { log.Errorf("get sub_domain (lcuuid: %s) from db failed: %s", subDomainLcuuid, err, d.metadata.LogPrefixes) diff --git a/server/controller/recorder/sub_domain.go b/server/controller/recorder/sub_domain.go index c04de73b7d0..d7aa45f47a8 100644 --- a/server/controller/recorder/sub_domain.go +++ b/server/controller/recorder/sub_domain.go @@ -62,10 +62,11 @@ func (s *subDomains) CloseStatsd() { } } -func (s *subDomains) RefreshAll(cloudSubDomainResources map[string]cloudmodel.SubDomainResource) error { +func (s *subDomains) RefreshAll(cloudSubDomains []cloudmodel.SubDomain, cloudSubDomainResources map[string]cloudmodel.SubDomainResource) error { // 遍历 cloud 中的 subdomain 资源,与缓存中的 subdomain 资源对比,根据对比结果增删改 var err error for lcuuid, resource := range cloudSubDomainResources { + log.Infof("sub_domain(lcuuid=%s) will be refreshed", lcuuid, s.metadata.LogPrefixes) sd, ok := s.refreshers[lcuuid] if !ok { sd, err = s.newRefresher(lcuuid) @@ -79,10 +80,21 @@ func (s *subDomains) RefreshAll(cloudSubDomainResources map[string]cloudmodel.Su } } + lcuuidToCloudSubDomain := make(map[string]cloudmodel.SubDomain) + for _, sd := range cloudSubDomains { + lcuuidToCloudSubDomain[sd.Lcuuid] = sd + } + // 遍历 subdomain 字典,删除 cloud 未返回的 subdomain 资源 for lcuuid, sd := range s.refreshers { + log.Infof("sub_domain(lcuuid=%s) resources will be checked", lcuuid, sd.metadata.LogPrefixes) if _, ok := cloudSubDomainResources[lcuuid]; !ok { - log.Info("sub_domain will be deleted", sd.metadata.LogPrefixes) + // 当 cloud 中存在 subDomain 时,表明数据不一致,不删除 + if item, ok := lcuuidToCloudSubDomain[lcuuid]; ok { + log.Infof("sub_domain(lcuuid=%s, name=%s) is still in cloud, skip", lcuuid, item.Name, sd.metadata.LogPrefixes) + continue + } + log.Infof("sub_domain(lcuuid=%s) resources will be deleted", lcuuid, sd.metadata.LogPrefixes) sd.clear() delete(s.refreshers, lcuuid) delete(s.cacheMng.SubDomainCacheMap, lcuuid) @@ -95,6 +107,7 @@ func (s *subDomains) RefreshOne(cloudSubDomainResources map[string]cloudmodel.Su // 遍历 cloud 中的 subdomain 资源,与缓存中的 subdomain 资源对比,根据对比结果增删改 var err error for lcuuid, resource := range cloudSubDomainResources { + log.Infof("sub_domain(lcuuid=%s) will be refreshed", lcuuid, s.metadata.LogPrefixes) sd, ok := s.refreshers[lcuuid] if !ok { sd, err = s.newRefresher(lcuuid) diff --git a/server/server.yaml b/server/server.yaml index 487492c61b3..d28784aee7b 100644 --- a/server/server.yaml +++ b/server/server.yaml @@ -266,6 +266,13 @@ controller: # 资源ID限制:所有设备ID(除宿主机外)、容器节点、Ingress、工作负载、ReplicaSet、POD resource_max_id_1: 499999 # local debug log + # 是否跳过同步空资源 + skip_sync_if_empty: + # 默认总开关:false + enabled: false + # 当前版本仅支持配置 sub_domain + resources: + # - sub_domain log_debug: enabled: false detail_enabled: false