20
21
// NewClusterMetadataLoader creates a new [ClusterMetadataLoader] that loads cluster metadata from the database.
22
>
func NewClusterMetadataLoader(manager persistence.ClusterMetadataManager, logger log.Logger) *ClusterMetadataLoader {
cluster_metadata_loader.go
23
>
return &ClusterMetadataLoader{
24
>
manager: manager,
25
>
logger: logger,
26
>
}
27
>
}
28
29
// LoadAndMergeWithStaticConfig loads cluster metadata from the database and merges it with the static config.
30
>
func (c *ClusterMetadataLoader) LoadAndMergeWithStaticConfig(ctx context.Context, svc *config.Config) error {
cluster_metadata_loader.go
31
>
iter := cluster.GetAllClustersIter(ctx, c.manager)
32
>
33
>
for iter.HasNext() {
34
>
item, err := iter.Next()
35
>
if err != nil {
36
return err
37
}
39
>
c.mergeMetadataFromDBWithStaticConfig(svc, item.ClusterName, newMetadata)
40
}
42
}
43
44
>
func (c *ClusterMetadataLoader) mergeMetadataFromDBWithStaticConfig(svc *config.Config, clusterName string, newMetadata *cluster.ClusterInformation) {
cluster_metadata_loader.go
45
>
c.backfillShardCount(svc, newMetadata)
46
>
if currentMetadata, ok := svc.ClusterMetadata.ClusterInformation[clusterName]; ok {
47
>
c.reconcileMetadata(svc, clusterName, currentMetadata, newMetadata)
48
>
}
49
>
svc.ClusterMetadata.ClusterInformation[clusterName] = *newMetadata
50
}
51