k8s: Add a dev only feature flag and simple service to get a client (#60204)
This commit is contained in:
@@ -22,6 +22,7 @@ import (
|
||||
"github.com/grafana/grafana/pkg/services/searchV2"
|
||||
"github.com/grafana/grafana/pkg/services/stats"
|
||||
"github.com/grafana/grafana/pkg/services/store/entity/httpentitystore"
|
||||
"github.com/grafana/grafana/pkg/services/store/k8saccess"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"github.com/prometheus/client_golang/prometheus/promhttp"
|
||||
@@ -253,6 +254,7 @@ func ProvideHTTPServer(opts ServerOptions, cfg *setting.Cfg, routeRegister routi
|
||||
annotationRepo annotations.Repository, tagService tag.Service, searchv2HTTPService searchV2.SearchHTTPService,
|
||||
queryLibraryHTTPService querylibrary.HTTPService, queryLibraryService querylibrary.Service, oauthTokenService oauthtoken.OAuthTokenService,
|
||||
statsService stats.Service,
|
||||
k8saccess k8saccess.K8SAccess, // required so that the router is registered
|
||||
) (*HTTPServer, error) {
|
||||
web.Env = cfg.Env
|
||||
m := web.New()
|
||||
|
||||
@@ -121,6 +121,7 @@ import (
|
||||
"github.com/grafana/grafana/pkg/services/store"
|
||||
"github.com/grafana/grafana/pkg/services/store/entity/httpentitystore"
|
||||
"github.com/grafana/grafana/pkg/services/store/entity/sqlstash"
|
||||
"github.com/grafana/grafana/pkg/services/store/k8saccess"
|
||||
"github.com/grafana/grafana/pkg/services/store/kind"
|
||||
"github.com/grafana/grafana/pkg/services/store/resolver"
|
||||
"github.com/grafana/grafana/pkg/services/store/sanitizer"
|
||||
@@ -352,6 +353,7 @@ var wireBasicSet = wire.NewSet(
|
||||
wire.Bind(new(tag.Service), new(*tagimpl.Service)),
|
||||
authnimpl.ProvideService,
|
||||
wire.Bind(new(authn.Service), new(*authnimpl.Service)),
|
||||
k8saccess.ProvideK8SAccess,
|
||||
)
|
||||
|
||||
var wireSet = wire.NewSet(
|
||||
|
||||
@@ -154,6 +154,12 @@ var (
|
||||
Description: "Configurable storage for dashboards, datasources, and resources",
|
||||
State: FeatureStateAlpha,
|
||||
},
|
||||
{
|
||||
Name: "k8s",
|
||||
Description: "Explore native k8s integrations",
|
||||
State: FeatureStateAlpha,
|
||||
RequiresDevMode: true,
|
||||
},
|
||||
{
|
||||
Name: "dashboardsFromStorage",
|
||||
Description: "Load dashboards from the generic storage interface",
|
||||
|
||||
@@ -115,6 +115,10 @@ const (
|
||||
// Configurable storage for dashboards, datasources, and resources
|
||||
FlagStorage = "storage"
|
||||
|
||||
// FlagK8s
|
||||
// Explore native k8s integrations
|
||||
FlagK8s = "k8s"
|
||||
|
||||
// FlagDashboardsFromStorage
|
||||
// Load dashboards from the generic storage interface
|
||||
FlagDashboardsFromStorage = "dashboardsFromStorage"
|
||||
|
||||
@@ -25,6 +25,7 @@ func TestFeatureToggleFiles(t *testing.T) {
|
||||
"live-config": true,
|
||||
"live-pipeline": true,
|
||||
"live-service-web-worker": true,
|
||||
"k8s": true, // Camle case does not like this one
|
||||
}
|
||||
|
||||
t.Run("check registry constraints", func(t *testing.T) {
|
||||
|
||||
@@ -152,13 +152,24 @@ func (s *ServiceImpl) getServerAdminNode(c *models.ReqContext) *navtree.NavLink
|
||||
}
|
||||
|
||||
if hasAccess(ac.ReqGrafanaAdmin, ac.EvalPermission(ac.ActionSettingsRead)) && s.features.IsEnabled(featuremgmt.FlagStorage) {
|
||||
adminNavLinks = append(adminNavLinks, &navtree.NavLink{
|
||||
storage := &navtree.NavLink{
|
||||
Text: "Storage",
|
||||
Id: "storage",
|
||||
SubTitle: "Manage file storage",
|
||||
Icon: "cube",
|
||||
Url: s.cfg.AppSubURL + "/admin/storage",
|
||||
})
|
||||
}
|
||||
adminNavLinks = append(adminNavLinks, storage)
|
||||
|
||||
if s.features.IsEnabled(featuremgmt.FlagK8s) {
|
||||
storage.Children = append(storage.Children, &navtree.NavLink{
|
||||
Text: "Kubernetes",
|
||||
Id: "k8s",
|
||||
SubTitle: "Manage k8s storage",
|
||||
Icon: "cube",
|
||||
Url: s.cfg.AppSubURL + "/admin/storage/k8s",
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
if s.cfg.LDAPEnabled && hasAccess(ac.ReqGrafanaAdmin, ac.EvalPermission(ac.ActionLDAPStatusRead)) {
|
||||
|
||||
@@ -0,0 +1,99 @@
|
||||
package k8saccess
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"net/url"
|
||||
|
||||
"github.com/grafana/grafana/pkg/models"
|
||||
"github.com/grafana/grafana/pkg/web"
|
||||
"k8s.io/apimachinery/pkg/runtime/schema"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/client-go/rest"
|
||||
)
|
||||
|
||||
type clientWrapper struct {
|
||||
err error
|
||||
baseURL *url.URL
|
||||
client *kubernetes.Clientset
|
||||
config *rest.Config
|
||||
httpClient *http.Client
|
||||
}
|
||||
|
||||
func newClientWrapper(config *rest.Config) *clientWrapper {
|
||||
if config.UserAgent == "" {
|
||||
config.UserAgent = rest.DefaultKubernetesUserAgent()
|
||||
}
|
||||
|
||||
url, _, err := defaultServerUrlFor(config)
|
||||
wrapper := &clientWrapper{
|
||||
config: config,
|
||||
baseURL: url,
|
||||
err: err,
|
||||
}
|
||||
|
||||
if err == nil && config != nil {
|
||||
// share the transport between all clients
|
||||
wrapper.httpClient, wrapper.err = rest.HTTPClientFor(config)
|
||||
if wrapper.err == nil {
|
||||
wrapper.client, wrapper.err = kubernetes.NewForConfigAndClient(config, wrapper.httpClient)
|
||||
}
|
||||
}
|
||||
|
||||
return wrapper
|
||||
}
|
||||
|
||||
func (s *clientWrapper) getInfo() map[string]interface{} {
|
||||
info := make(map[string]interface{}, 0)
|
||||
|
||||
if s.err != nil {
|
||||
info["error"] = s.err.Error()
|
||||
}
|
||||
|
||||
if s.baseURL != nil {
|
||||
info["baseURL"] = s.baseURL.String()
|
||||
}
|
||||
|
||||
if s.client != nil {
|
||||
v, err := s.client.ServerVersion()
|
||||
if err != nil {
|
||||
info["version_error"] = err.Error()
|
||||
}
|
||||
if v != nil {
|
||||
info["k8s.version"] = v
|
||||
}
|
||||
}
|
||||
return info
|
||||
}
|
||||
|
||||
// defaultServerUrlFor is shared between IsConfigTransportTLS and RESTClientFor. It
|
||||
// requires Host and Version to be set prior to being called.
|
||||
func defaultServerUrlFor(config *rest.Config) (*url.URL, string, error) {
|
||||
// TODO: move the default to secure when the apiserver supports TLS by default
|
||||
// config.Insecure is taken to mean "I want HTTPS but don't bother checking the certs against a CA."
|
||||
hasCA := len(config.CAFile) != 0 || len(config.CAData) != 0
|
||||
hasCert := len(config.CertFile) != 0 || len(config.CertData) != 0
|
||||
defaultTLS := hasCA || hasCert || config.Insecure
|
||||
host := config.Host
|
||||
if host == "" {
|
||||
host = "localhost"
|
||||
}
|
||||
|
||||
if config.GroupVersion != nil {
|
||||
return rest.DefaultServerURL(host, config.APIPath, *config.GroupVersion, defaultTLS)
|
||||
}
|
||||
return rest.DefaultServerURL(host, config.APIPath, schema.GroupVersion{}, defaultTLS)
|
||||
}
|
||||
|
||||
func (s *clientWrapper) doProxy(c *models.ReqContext) {
|
||||
if s.baseURL == nil {
|
||||
c.Resp.WriteHeader(500)
|
||||
return
|
||||
}
|
||||
|
||||
params := web.Params(c.Req)
|
||||
path := params["*"]
|
||||
|
||||
url := s.baseURL.JoinPath(path)
|
||||
|
||||
_, _ = c.Resp.Write([]byte("TODO, proxy: " + url.String()))
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
package k8saccess
|
||||
|
||||
import (
|
||||
"github.com/grafana/grafana/pkg/api/response"
|
||||
"github.com/grafana/grafana/pkg/api/routing"
|
||||
"github.com/grafana/grafana/pkg/middleware"
|
||||
"github.com/grafana/grafana/pkg/models"
|
||||
)
|
||||
|
||||
type httpHelper struct {
|
||||
access *k8sAccess
|
||||
}
|
||||
|
||||
func newHTTPHelper(access *k8sAccess, router routing.RouteRegister) *httpHelper {
|
||||
s := &httpHelper{
|
||||
access: access,
|
||||
}
|
||||
|
||||
// Must be admin for everything
|
||||
router.Group("/api/k8s", func(k8sRoute routing.RouteRegister) {
|
||||
k8sRoute.Get("/info", middleware.ReqOrgAdmin, routing.Wrap(s.showClientInfo))
|
||||
k8sRoute.Any("/proxy/*", middleware.ReqOrgAdmin, s.doProxy)
|
||||
})
|
||||
|
||||
return s
|
||||
}
|
||||
|
||||
func (s *httpHelper) showClientInfo(c *models.ReqContext) response.Response {
|
||||
if s.access.sys != nil {
|
||||
info := s.access.sys.getInfo()
|
||||
if s.access.sys.err != nil {
|
||||
return response.JSON(500, info)
|
||||
}
|
||||
return response.JSON(200, info)
|
||||
}
|
||||
return response.JSON(500, map[string]interface{}{
|
||||
"error": "no client initialized",
|
||||
})
|
||||
}
|
||||
|
||||
func (s *httpHelper) doProxy(c *models.ReqContext) {
|
||||
// TODO... this does not yet do a real proxy
|
||||
if s.access.sys != nil {
|
||||
if s.access.sys.err == nil {
|
||||
s.access.sys.doProxy(c)
|
||||
} else {
|
||||
c.Resp.WriteHeader(500)
|
||||
}
|
||||
return
|
||||
}
|
||||
_, _ = c.Resp.Write([]byte("??"))
|
||||
}
|
||||
@@ -0,0 +1,80 @@
|
||||
package k8saccess
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
|
||||
"github.com/grafana/grafana/pkg/api/routing"
|
||||
"github.com/grafana/grafana/pkg/registry"
|
||||
"github.com/grafana/grafana/pkg/services/featuremgmt"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/client-go/rest"
|
||||
"k8s.io/client-go/tools/clientcmd"
|
||||
)
|
||||
|
||||
type K8SAccess interface {
|
||||
registry.CanBeDisabled
|
||||
|
||||
// Get the system client
|
||||
GetSystemClient() *kubernetes.Clientset
|
||||
}
|
||||
|
||||
var _ K8SAccess = &k8sAccess{}
|
||||
|
||||
type k8sAccess struct {
|
||||
enabled bool
|
||||
apihelper *httpHelper
|
||||
sys *clientWrapper
|
||||
}
|
||||
|
||||
func ProvideK8SAccess(toggles featuremgmt.FeatureToggles, router routing.RouteRegister) K8SAccess {
|
||||
access := &k8sAccess{
|
||||
enabled: toggles.IsEnabled(featuremgmt.FlagK8s),
|
||||
}
|
||||
|
||||
// Skips setting up any HTTP routing
|
||||
if !access.enabled {
|
||||
return access // dummy
|
||||
}
|
||||
|
||||
// If we are in a cluster, this is the
|
||||
config, err := rest.InClusterConfig()
|
||||
|
||||
// Look for kube config setup
|
||||
if err != nil {
|
||||
var home string
|
||||
var configBytes []byte
|
||||
home, err = os.UserHomeDir()
|
||||
if err == nil {
|
||||
fpath := filepath.Join(home, ".kube", "config")
|
||||
//nolint:gosec
|
||||
configBytes, err = os.ReadFile(fpath)
|
||||
if err == nil {
|
||||
config, err = clientcmd.RESTConfigFromKubeConfig(configBytes)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if err == nil && config != nil {
|
||||
access.sys = newClientWrapper(config)
|
||||
} else {
|
||||
access.sys = &clientWrapper{
|
||||
err: err,
|
||||
}
|
||||
}
|
||||
|
||||
access.apihelper = newHTTPHelper(access, router)
|
||||
return access
|
||||
}
|
||||
|
||||
func (s *k8sAccess) IsDisabled() bool {
|
||||
return !s.enabled
|
||||
}
|
||||
|
||||
// Return access to the system k8s client
|
||||
func (s *k8sAccess) GetSystemClient() *kubernetes.Clientset {
|
||||
if s.sys != nil {
|
||||
return s.sys.client
|
||||
}
|
||||
return nil
|
||||
}
|
||||
Reference in New Issue
Block a user