kubemodel.go 3.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116
  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. // Compute every resource from one consistent state of the data source by running the compute
  35. // functions against a copy of the KubeModel whose data source is pinned.
  36. ds, release := source.PinDataSource(km.ds)
  37. defer release()
  38. pinned := *km
  39. pinned.ds = ds
  40. computeFuncs := pinned.computeFuncs(start, end)
  41. for _, f := range computeFuncs {
  42. if err := f(kms, start, end); err != nil {
  43. kms.Error(err)
  44. return kms, fmt.Errorf("error computing kubemodel for (%s, %s): %w", start.Format(time.DateTime), end.Format(time.DateTime), err)
  45. }
  46. }
  47. kms.Metadata.CompletedAt = time.Now().UTC()
  48. return kms, nil
  49. }
  50. // computeFuncs returns the set of compute functions to run for the window,
  51. // in struct field order. If cluster_info is not yet reporting the
  52. // complete_kubemodel label for any cluster in the result, the source has not
  53. // been upgraded to emit a full kubemodel, so only the minimal set of
  54. // resources (cluster, namespaces, resource quotas) is computed. The
  55. // FORCE_KUBEMODEL_V1 env var skips the check entirely and always returns the
  56. // minimal set.
  57. func (km *KubeModel) computeFuncs(start, end time.Time) []computeFunc {
  58. KubeModelV1ComputeFucs := []computeFunc{
  59. km.computeCluster,
  60. km.computeNamespaces,
  61. km.computeResourceQuotas,
  62. }
  63. if km.forceV1 {
  64. return KubeModelV1ComputeFucs
  65. }
  66. results, err := km.ds.Metrics().QueryClusterKubeModelVersion(start, end).Await()
  67. if err != nil {
  68. log.Errorf("computeFuncs: querying cluster complete kubemodel: %s", err)
  69. return KubeModelV1ComputeFucs
  70. }
  71. // If the window contains a result which is missing the version number then the window will produce incomplete KubeModel data
  72. // and should not export
  73. for _, res := range results {
  74. if res.Version == "" {
  75. return KubeModelV1ComputeFucs
  76. }
  77. }
  78. return []computeFunc{
  79. km.computeCluster,
  80. km.computeNamespaces,
  81. km.computeResourceQuotas,
  82. km.computeServices,
  83. km.computeDeployments,
  84. km.computeStatefulSets,
  85. km.computeDaemonSets,
  86. km.computeJobs,
  87. km.computeCronJobs,
  88. km.computeReplicaSets,
  89. km.computeNodes,
  90. km.computePersistentVolumes,
  91. km.computePersistentVolumeClaims,
  92. km.computePods,
  93. km.computeContainers,
  94. km.computeDevices,
  95. }
  96. }