kubemodel.go 3.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109
  1. package kubemodel
  2. import (
  3. "errors"
  4. "fmt"
  5. "time"
  6. "github.com/opencost/opencost/core/pkg/log"
  7. "github.com/opencost/opencost/core/pkg/model/kubemodel"
  8. "github.com/opencost/opencost/core/pkg/source"
  9. )
  10. type KubeModel struct {
  11. ds source.OpenCostDataSource
  12. forceV1 bool
  13. clusterUID string
  14. }
  15. func NewKubeModel(clusterUID string, forceV1 bool, dataSource source.OpenCostDataSource) (*KubeModel, error) {
  16. if dataSource == nil {
  17. return nil, errors.New("OpenCostDataSource cannot be nil")
  18. }
  19. km := &KubeModel{
  20. ds: dataSource,
  21. forceV1: forceV1,
  22. clusterUID: clusterUID,
  23. }
  24. km.clusterUID = clusterUID
  25. log.Debugf("NewKubeModel(%s)", km.clusterUID)
  26. return km, nil
  27. }
  28. type computeFunc func(*kubemodel.KubeModelSet, time.Time, time.Time) error
  29. // ComputeKubeModel uses the CostModel instance to compute an KubeModelSet
  30. // for the window defined by the given start and end times. The KubeModels
  31. // returned are unaggregated (i.e. down to the container level).
  32. func (km *KubeModel) ComputeKubeModelSet(start, end time.Time) (*kubemodel.KubeModelSet, error) {
  33. kms := kubemodel.NewKubeModelSet(start, end)
  34. computeFuncs := km.computeFuncs(start, end)
  35. for _, f := range computeFuncs {
  36. if err := f(kms, start, end); err != nil {
  37. kms.Error(err)
  38. return kms, fmt.Errorf("error computing kubemodel for (%s, %s): %w", start.Format(time.DateTime), end.Format(time.DateTime), err)
  39. }
  40. }
  41. kms.Metadata.CompletedAt = time.Now().UTC()
  42. return kms, nil
  43. }
  44. // computeFuncs returns the set of compute functions to run for the window,
  45. // in struct field order. If cluster_info is not yet reporting the
  46. // complete_kubemodel label for any cluster in the result, the source has not
  47. // been upgraded to emit a full kubemodel, so only the minimal set of
  48. // resources (cluster, namespaces, resource quotas) is computed. The
  49. // FORCE_KUBEMODEL_V1 env var skips the check entirely and always returns the
  50. // minimal set.
  51. func (km *KubeModel) computeFuncs(start, end time.Time) []computeFunc {
  52. KubeModelV1ComputeFucs := []computeFunc{
  53. km.computeCluster,
  54. km.computeNamespaces,
  55. km.computeResourceQuotas,
  56. }
  57. if km.forceV1 {
  58. return KubeModelV1ComputeFucs
  59. }
  60. results, err := km.ds.Metrics().QueryClusterKubeModelVersion(start, end).Await()
  61. if err != nil {
  62. log.Errorf("computeFuncs: querying cluster complete kubemodel: %s", err)
  63. return KubeModelV1ComputeFucs
  64. }
  65. // If the window contains a result which is missing the version number then the window will produce incomplete KubeModel data
  66. // and should not export
  67. for _, res := range results {
  68. if res.Version == "" {
  69. return KubeModelV1ComputeFucs
  70. }
  71. }
  72. return []computeFunc{
  73. km.computeCluster,
  74. km.computeNamespaces,
  75. km.computeResourceQuotas,
  76. km.computeServices,
  77. km.computeDeployments,
  78. km.computeStatefulSets,
  79. km.computeDaemonSets,
  80. km.computeJobs,
  81. km.computeCronJobs,
  82. km.computeReplicaSets,
  83. km.computeNodes,
  84. km.computePersistentVolumes,
  85. km.computePersistentVolumeClaims,
  86. km.computePods,
  87. km.computeContainers,
  88. //km.computeDCGMDevices,
  89. }
  90. }