Skip to content

Commit 83d594f

Browse files
committed
fix: genesis sync
1 parent dc73af9 commit 83d594f

15 files changed

Lines changed: 381 additions & 165 deletions

File tree

server/controller/cloud/kubernetes_gather/kubernetes_gather.go

Lines changed: 43 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
package kubernetes_gather
1818

1919
import (
20+
"errors"
2021
"regexp"
2122
"strings"
2223
"time"
@@ -33,6 +34,7 @@ import (
3334
"github.com/deepflowio/deepflow/server/controller/db/metadb"
3435
metadbmodel "github.com/deepflowio/deepflow/server/controller/db/metadb/model"
3536
"github.com/deepflowio/deepflow/server/controller/genesis"
37+
cmodel "github.com/deepflowio/deepflow/server/controller/model"
3638
"github.com/deepflowio/deepflow/server/controller/statsd"
3739
"github.com/deepflowio/deepflow/server/libs/logger"
3840
)
@@ -63,6 +65,7 @@ type KubernetesGather struct {
6365
podGroupLcuuids mapset.Set
6466
podNetworkLcuuidCIDRs networkLcuuidCIDRs
6567
nodeNetworkLcuuidCIDRs networkLcuuidCIDRs
68+
vinterfaceData []cmodel.GenesisVinterface
6669
podIPToLcuuid map[string]string
6770
nodeIPToLcuuid map[string]string
6871
namespaceToLcuuid map[string]string
@@ -96,6 +99,10 @@ func NewKubernetesGather(db *metadb.DB, domain *metadbmodel.Domain, subDomain *m
9699
var err error
97100

98101
domainConfigJson, err = simplejson.NewJson([]byte(domain.Config))
102+
if err != nil {
103+
log.Error(err, logger.NewORGPrefix(db.ORGID))
104+
return nil
105+
}
99106
portNameRegex := domainConfigJson.Get("node_port_name_regex").MustString()
100107
if portNameRegex == "" {
101108
portNameRegex = common.DEFAULT_PORT_NAME_REGEX
@@ -202,6 +209,7 @@ func NewKubernetesGather(db *metadb.DB, domain *metadbmodel.Domain, subDomain *m
202209
customTagLenMax: cfg.CustomTagLenMax,
203210
isSubDomain: isSubDomain,
204211
podGroupLcuuids: mapset.NewSet(),
212+
vinterfaceData: []cmodel.GenesisVinterface{},
205213
nodeNetworkLcuuidCIDRs: networkLcuuidCIDRs{},
206214
podNetworkLcuuidCIDRs: networkLcuuidCIDRs{},
207215
podIPToLcuuid: map[string]string{},
@@ -277,7 +285,7 @@ func (k *KubernetesGather) pgSpecGenerateConnections(nsName, pgName, pgLcuuid st
277285
if !ok {
278286
continue
279287
}
280-
cmName := ref.Get("Name").MustString()
288+
cmName := ref.Get("name").MustString()
281289
cmLcuuid, ok := k.configMapToLcuuid[[2]string{nsName, cmName}]
282290
if !ok {
283291
log.Infof("pod group (%s) imported env config map (%s) not found", pgName, cmName, logger.NewORGPrefix(k.orgID))
@@ -328,6 +336,7 @@ func (k *KubernetesGather) GetKubernetesGatherData() (model.KubernetesGatherReso
328336
// 任务循环的是同一个实例,所以这里要对关联关系进行初始化
329337
k.azLcuuid = ""
330338
k.k8sEntries = nil
339+
k.vinterfaceData = []cmodel.GenesisVinterface{}
331340
k.podNetworkLcuuidCIDRs = networkLcuuidCIDRs{}
332341
k.nodeNetworkLcuuidCIDRs = networkLcuuidCIDRs{}
333342
k.podGroupLcuuids = mapset.NewSet()
@@ -345,6 +354,39 @@ func (k *KubernetesGather) GetKubernetesGatherData() (model.KubernetesGatherReso
345354
k.nsServiceNameToService = map[string]map[string]map[string]int{}
346355
k.cloudStatsd = statsd.NewCloudStatsd()
347356

357+
if genesis.GenesisService == nil {
358+
errMsg := "genesis service is nil"
359+
log.Warning(errMsg, logger.NewORGPrefix(k.orgID))
360+
return model.KubernetesGatherResource{
361+
ErrorState: common.RESOURCE_STATE_CODE_EXIT,
362+
ErrorMessage: errMsg,
363+
}, errors.New(errMsg)
364+
}
365+
366+
k8sEntries, err := k.getKubernetesEntries()
367+
if err != nil {
368+
log.Warning(err.Error(), logger.NewORGPrefix(k.orgID))
369+
return model.KubernetesGatherResource{
370+
ErrorState: common.RESOURCE_STATE_CODE_WARNING,
371+
ErrorMessage: err.Error(),
372+
}, err
373+
}
374+
k.k8sEntries = k8sEntries
375+
376+
gsData, err := genesis.GenesisService.GetGenesisSyncResponse(k.orgID)
377+
if err != nil {
378+
return model.KubernetesGatherResource{
379+
ErrorState: common.RESOURCE_STATE_CODE_EXIT,
380+
ErrorMessage: err.Error(),
381+
}, err
382+
}
383+
for _, v := range gsData.Vinterfaces {
384+
if v.KubernetesClusterID != k.ClusterID {
385+
continue
386+
}
387+
k.vinterfaceData = append(k.vinterfaceData, v)
388+
}
389+
348390
region, err := k.getRegion()
349391
if err != nil {
350392
return model.KubernetesGatherResource{}, err
@@ -368,15 +410,6 @@ func (k *KubernetesGather) GetKubernetesGatherData() (model.KubernetesGatherReso
368410
}, err
369411
}
370412

371-
k.k8sEntries, err = k.getKubernetesEntries()
372-
if err != nil {
373-
log.Warning(err.Error(), logger.NewORGPrefix(k.orgID))
374-
return model.KubernetesGatherResource{
375-
ErrorState: common.RESOURCE_STATE_CODE_WARNING,
376-
ErrorMessage: err.Error(),
377-
}, err
378-
}
379-
380413
podCluster, err := k.getPodCluster()
381414
if err != nil {
382415
return model.KubernetesGatherResource{}, err

server/controller/cloud/kubernetes_gather/vinterface_and_ip.go

Lines changed: 4 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@
1717
package kubernetes_gather
1818

1919
import (
20-
"errors"
2120
"regexp"
2221
"sort"
2322
"strings"
@@ -28,7 +27,6 @@ import (
2827
cloudcommon "github.com/deepflowio/deepflow/server/controller/cloud/common"
2928
"github.com/deepflowio/deepflow/server/controller/cloud/model"
3029
"github.com/deepflowio/deepflow/server/controller/common"
31-
"github.com/deepflowio/deepflow/server/controller/genesis"
3230
"github.com/deepflowio/deepflow/server/libs/logger"
3331
"github.com/mikioh/ipaddr"
3432
)
@@ -58,21 +56,7 @@ func (k *KubernetesGather) getVInterfacesAndIPs() (nodeSubnets, podSubnets []mod
5856
return
5957
}
6058

61-
// 获取vinterface API返回中host ip与其上所有node ip的对应关系
62-
if genesis.GenesisService == nil {
63-
err = errors.New("genesis service is nil")
64-
return
65-
}
66-
genesisData, err := genesis.GenesisService.GetGenesisSyncResponse(k.orgID)
67-
if err != nil {
68-
log.Error(err.Error(), logger.NewORGPrefix(k.orgID))
69-
return
70-
}
71-
vData := genesisData.Vinterfaces
72-
for _, vItem := range vData {
73-
if vItem.KubernetesClusterID != k.ClusterID {
74-
continue
75-
}
59+
for _, vItem := range k.vinterfaceData {
7660
deviceType := vItem.DeviceType
7761
if deviceType == "docker-host" || deviceType == "kvm-host" {
7862
hostIP := vItem.HostIP
@@ -96,10 +80,7 @@ func (k *KubernetesGather) getVInterfacesAndIPs() (nodeSubnets, podSubnets []mod
9680
}
9781
}
9882
// 生成device_uuid或uuid和pod lcuuid的对应关系
99-
for _, vItem := range vData {
100-
if vItem.KubernetesClusterID != k.ClusterID {
101-
continue
102-
}
83+
for _, vItem := range k.vinterfaceData {
10384
if vItem.DeviceType != "docker-container" {
10485
continue
10586
}
@@ -135,10 +116,7 @@ func (k *KubernetesGather) getVInterfacesAndIPs() (nodeSubnets, podSubnets []mod
135116
}
136117

137118
// 处理POD IP,生成port,ip,cidrs信息
138-
for _, vItem := range vData {
139-
if vItem.KubernetesClusterID != k.ClusterID {
140-
continue
141-
}
119+
for _, vItem := range k.vinterfaceData {
142120
if vItem.DeviceType != "docker-container" {
143121
continue
144122
}
@@ -402,10 +380,7 @@ func (k *KubernetesGather) getVInterfacesAndIPs() (nodeSubnets, podSubnets []mod
402380

403381
// 处理nodeIP,生成port,ip,cidrs信息
404382
nodeVinterfaceLcuuids := mapset.NewSet()
405-
for _, vItem := range vData {
406-
if vItem.KubernetesClusterID != k.ClusterID {
407-
continue
408-
}
383+
for _, vItem := range k.vinterfaceData {
409384
deviceType := vItem.DeviceType
410385
if deviceType != "docker-host" && deviceType != "kvm-host" {
411386
continue

server/controller/cloud/kubernetes_gather_task.go

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -122,10 +122,14 @@ func (k *KubernetesGatherTask) run(rSignal *queue.OverwriteQueue) {
122122
kResource, err := k.kubernetesGather.GetKubernetesGatherData()
123123
// 这里因为任务内部没有对成功的状态赋值状态码,在这里统一处理了
124124
if err != nil {
125-
kResource.ErrorMessage = fmt.Sprintf("%s %s", time.Now().Format(common.GO_BIRTHDAY), err.Error())
125+
if kResource.ErrorState == common.RESOURCE_STATE_CODE_EXIT {
126+
log.Infof("kubernetes gather (%s) assemble failed: %s", k.kubernetesGather.Name, err.Error(), logger.NewORGPrefix(k.orgID))
127+
return
128+
}
126129
if kResource.ErrorState == 0 {
127130
kResource.ErrorState = common.RESOURCE_STATE_CODE_EXCEPTION
128131
}
132+
kResource.ErrorMessage = fmt.Sprintf("%s %s", time.Now().Format(common.GO_BIRTHDAY), err.Error())
129133
} else {
130134
kResource.ErrorState = common.RESOURCE_STATE_CODE_SUCCESS
131135
}

server/controller/common/const.go

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -622,11 +622,12 @@ const (
622622
)
623623

624624
const (
625-
RESOURCE_STATE_CODE_SUCCESS = 1
626-
RESOURCE_STATE_CODE_DELETING = 2
627-
RESOURCE_STATE_CODE_EXCEPTION = 3
628-
RESOURCE_STATE_CODE_WARNING = 4
629-
RESOURCE_STATE_CODE_NO_LICENSE = 5
625+
RESOURCE_STATE_CODE_SUCCESS = 1 + iota
626+
RESOURCE_STATE_CODE_DELETING
627+
RESOURCE_STATE_CODE_EXCEPTION
628+
RESOURCE_STATE_CODE_WARNING
629+
RESOURCE_STATE_CODE_NO_LICENSE
630+
RESOURCE_STATE_CODE_EXIT
630631
)
631632

632633
const (

server/controller/genesis/common/type.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,11 +25,13 @@ import (
2525

2626
type GenesisSync interface {
2727
Start()
28+
GetVtapUpdatedVersion(key string) (uint64, bool)
2829
GetGenesisSyncData(orgID int) GenesisSyncDataResponse
2930
GetGenesisSyncResponse(orgID int) (GenesisSyncDataResponse, error)
3031
}
3132

3233
type GenesisSyncType interface {
34+
GetInfo() string
3335
GetLcuuid() string
3436
GetVtapID() uint32
3537
}
@@ -82,6 +84,7 @@ type VIFRPCMessage struct {
8284
MessageType int
8385
TeamID uint32
8486
VtapID uint32
87+
Version uint64
8588
Peer string
8689
K8SClusterID string
8790
Key string

server/controller/genesis/config/config.go

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,14 +19,15 @@ package config
1919
type GenesisConfig struct {
2020
AgingTime float64 `default:"86400" yaml:"aging_time"`
2121
VinterfaceAgingTime float64 `default:"300" yaml:"vinterface_aging_time"`
22-
AgentHeartBeat float64 `default:"60" yaml:"agent_heart_beat"`
22+
AgentHeartBeat float64 `default:"10" yaml:"agent_heart_beat"`
2323
HostIPs []string `yaml:"host_ips"`
2424
LocalIPRanges []string `yaml:"local_ip_ranges"`
2525
ExcludeIPRanges []string `yaml:"exclude_ip_ranges"`
2626
QueueLengths int `default:"60" yaml:"queue_length"`
2727
DataPersistenceInterval int `default:"60" yaml:"data_persistence_interval"`
2828
MultiNSMode bool `default:"false" yaml:"multi_ns_mode"`
2929
SingleVPCMode bool `default:"false" yaml:"single_vpc_mode"`
30+
LogDetailEnabled bool `default:"false" yaml:"log_detail_enabled"`
3031
Database string `default:"mysql" yaml:"database"`
3132
DefaultVPCName string `default:"default-public-vpc" yaml:"default_vpc_name"`
3233
IgnoreNICRegex string `default:"^(kube-ipvs)" yaml:"ignore_nic_regex"`

server/controller/genesis/grpc/server.go

Lines changed: 25 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -260,7 +260,15 @@ func (g *SynchronizerServer) GenesisSync(ctx context.Context, request *agent.Gen
260260
_, enabled := g.workloadResourceEnabledCache.Get(fmt.Sprintf("%d-%s", orgID, groupShortLcuuid))
261261

262262
platformData := request.GetPlatformData()
263-
if version == localVersion || platformData == nil {
263+
if version == localVersion {
264+
// It is necessary to verify whether there are available data for use,
265+
// If there is no data, needs to be re-reported.
266+
if v, ok := g.gsync.GetVtapUpdatedVersion(vtap); !ok || v != version {
267+
g.vtapToVersion.Store(vtap, uint64(0))
268+
log.Infof("genesis sync re-reporting, not data for keep alive, reset local version, from ip %s vtap_id %v", remote, vtapID, logger.NewORGPrefix(orgID))
269+
return &agent.GenesisSyncResponse{}, nil
270+
}
271+
264272
// If the worload-v is modified to be enabled during the period of continuous heartbeat,
265273
// it will trigger the re-reporting of the full data.
266274
if _, ok := g.workloadResourceChangeEnabledCache.Get(vtap); ok {
@@ -278,6 +286,7 @@ func (g *SynchronizerServer) GenesisSync(ctx context.Context, request *agent.Gen
278286
VtapID: vtapID,
279287
ORGID: orgID,
280288
TeamID: uint32(teamID),
289+
Version: version,
281290
MessageType: common.TYPE_RENEW,
282291
Message: request,
283292
StorageRefresh: refresh,
@@ -287,6 +296,20 @@ func (g *SynchronizerServer) GenesisSync(ctx context.Context, request *agent.Gen
287296
return &agent.GenesisSyncResponse{Version: &localVersion}, nil
288297
}
289298

299+
if platformData == nil {
300+
log.Debugf("genesis sync received version %v platform is nil from ip %s vtap_id %v, local version %v", version, remote, vtapID, localVersion, logger.NewORGPrefix(orgID))
301+
return &agent.GenesisSyncResponse{}, nil
302+
}
303+
304+
// 当采集器为容器类型时(cluster id 非空)
305+
// - 采集器未注册(vtapID==0),即使没有 Interfaces 也需要处理 vinterface 来让采集器能够注册
306+
// - 采集器已经注册(vtapID!=0),采集器重启会出现 Interfaces 为空的情况
307+
// 为了避免 vinterface 异常增删,丢弃当前消息
308+
if k8sClusterID != "" && len(platformData.Interfaces) == 0 && vtapID != 0 {
309+
log.Infof("genesis sync received version %v message with empty interfaces from ip %s vtap_id %v, local version %v", version, remote, vtapID, localVersion, logger.NewORGPrefix(orgID))
310+
return &agent.GenesisSyncResponse{Version: &version}, nil
311+
}
312+
290313
log.Infof("genesis sync received version %v -> %v from ip %s vtap_id %v", localVersion, version, remote, vtapID, logger.NewORGPrefix(orgID))
291314
g.genesisSyncQueue.Put(
292315
common.VIFRPCMessage{
@@ -295,6 +318,7 @@ func (g *SynchronizerServer) GenesisSync(ctx context.Context, request *agent.Gen
295318
VtapID: vtapID,
296319
ORGID: orgID,
297320
TeamID: uint32(teamID),
321+
Version: version,
298322
K8SClusterID: k8sClusterID,
299323
MessageType: common.TYPE_UPDATE,
300324
Message: request,

server/controller/genesis/store/kubernetes/store.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -244,7 +244,7 @@ func (k *KubernetesStorage) generateCache() {
244244
} else {
245245
endpoint = net.JoinHostPort(domain.ControllerIP, strconv.Itoa(k.listenNodePort))
246246
}
247-
cacheMap[k.formatKey(db.ORGID, domain.ClusterID)] = common.ClusterDest{
247+
cacheMap[k.formatKey(db.ORGID, subDomain.ClusterID)] = common.ClusterDest{
248248
Endpoint: endpoint,
249249
DomainLcuuid: domain.Lcuuid,
250250
SubDomainLcuuid: subDomain.Lcuuid,

0 commit comments

Comments
 (0)