Provisioning: Add Standalone Job Controller Without Job Processing (#109610)

* Add standalone job controller
* Add makefile
* Add limit on the current implementation
* Move job controllers to app package
* Add TLS flags
This commit is contained in:
Roberto Jiménez Sánchez
2025-08-25 08:48:40 +00:00
committed by GitHub
parent 9204a08207
commit e7ccefcf92
24 changed files with 665 additions and 25 deletions
@@ -0,0 +1,45 @@
package auth
import (
"context"
"fmt"
"net/http"
"github.com/grafana/authlib/authn"
provisioning "github.com/grafana/grafana/apps/provisioning/pkg/apis/provisioning/v0alpha1"
utilnet "k8s.io/apimachinery/pkg/util/net"
)
// tokenExchanger abstracts the token exchange client for testability.
type tokenExchanger interface {
Exchange(ctx context.Context, req authn.TokenExchangeRequest) (*authn.TokenExchangeResponse, error)
}
// RoundTripper injects an exchanged access token for the provisioning API into outgoing requests.
type RoundTripper struct {
client tokenExchanger
transport http.RoundTripper
}
// NewRoundTripper constructs a RoundTripper that exchanges the provided token per request
// and forwards the request to the provided base transport.
func NewRoundTripper(tokenExchangeClient tokenExchanger, base http.RoundTripper) *RoundTripper {
return &RoundTripper{
client: tokenExchangeClient,
transport: base,
}
}
func (t *RoundTripper) RoundTrip(req *http.Request) (*http.Response, error) {
tokenResponse, err := t.client.Exchange(req.Context(), authn.TokenExchangeRequest{
Audiences: []string{provisioning.GROUP},
Namespace: "*",
})
if err != nil {
return nil, fmt.Errorf("failed to exchange token: %w", err)
}
req = utilnet.CloneRequest(req)
req.Header.Set("X-Access-Token", "Bearer "+tokenResponse.Token)
return t.transport.RoundTrip(req)
}
@@ -0,0 +1,63 @@
package auth
import (
"context"
"io"
"net/http"
"net/http/httptest"
"testing"
"github.com/grafana/authlib/authn"
)
type fakeExchanger struct {
resp *authn.TokenExchangeResponse
err error
}
func (f *fakeExchanger) Exchange(_ context.Context, req authn.TokenExchangeRequest) (*authn.TokenExchangeResponse, error) {
return f.resp, f.err
}
// roundTripperFunc allows building a stub transport inline
type roundTripperFunc func(*http.Request) (*http.Response, error)
func (f roundTripperFunc) RoundTrip(r *http.Request) (*http.Response, error) { return f(r) }
func TestRoundTripper_SetsAccessTokenHeader(t *testing.T) {
tr := NewRoundTripper(&fakeExchanger{resp: &authn.TokenExchangeResponse{Token: "abc123"}}, roundTripperFunc(func(r *http.Request) (*http.Response, error) {
got := r.Header.Get("X-Access-Token")
if got != "Bearer abc123" {
t.Fatalf("expected X-Access-Token header 'Bearer abc123', got %q", got)
}
// Return a minimal response; body must be non-nil per http.RoundTripper contract
rr := httptest.NewRecorder()
rr.WriteHeader(http.StatusOK)
return rr.Result(), nil
}))
req, _ := http.NewRequestWithContext(context.Background(), http.MethodGet, "http://example", nil)
resp, err := tr.RoundTrip(req)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
// drain and close body
_, _ = io.Copy(io.Discard, resp.Body)
_ = resp.Body.Close()
}
func TestRoundTripper_PropagatesExchangeError(t *testing.T) {
tr := NewRoundTripper(&fakeExchanger{err: io.EOF}, roundTripperFunc(func(_ *http.Request) (*http.Response, error) {
t.Fatal("transport should not be called on exchange error")
return nil, nil
}))
req, _ := http.NewRequestWithContext(context.Background(), http.MethodGet, "http://example", nil)
resp, err := tr.RoundTrip(req)
if err == nil {
if resp != nil && resp.Body != nil {
_ = resp.Body.Close()
}
t.Fatalf("expected error, got nil")
}
}
@@ -0,0 +1,89 @@
package controller
import (
"context"
"time"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apiserver/pkg/endpoints/request"
"k8s.io/client-go/tools/cache"
"github.com/grafana/grafana-app-sdk/logging"
provisioning "github.com/grafana/grafana/apps/provisioning/pkg/apis/provisioning/v0alpha1"
client "github.com/grafana/grafana/apps/provisioning/pkg/generated/clientset/versioned/typed/provisioning/v0alpha1"
informer "github.com/grafana/grafana/apps/provisioning/pkg/generated/informers/externalversions/provisioning/v0alpha1"
"github.com/grafana/grafana/pkg/apimachinery/identity"
)
const (
historyJobControllerLoggerName = "provisioning-historyjob-controller"
)
// HistoryJobController manages the cleanup of old HistoryJob entries.
type HistoryJobController struct {
client client.ProvisioningV0alpha1Interface
logger logging.Logger
expirationTime time.Duration
}
// NewHistoryJobController creates a new HistoryJobController.
func NewHistoryJobController(
provisioningClient client.ProvisioningV0alpha1Interface,
historyJobInformer informer.HistoricJobInformer,
expirationTime time.Duration,
) (*HistoryJobController, error) {
c := &HistoryJobController{
client: provisioningClient,
logger: logging.DefaultLogger.With("logger", historyJobControllerLoggerName),
expirationTime: expirationTime,
}
// Use the resync events from the shared informer to trigger cleanup for each job
_, err := historyJobInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
c.cleanupJob(obj)
},
UpdateFunc: func(oldObj, newObj interface{}) {
c.cleanupJob(newObj)
},
})
if err != nil {
return nil, err
}
return c, nil
}
func (c *HistoryJobController) cleanupJob(obj interface{}) {
job, ok := obj.(*provisioning.HistoricJob)
if !ok {
c.logger.Error("Expected HistoricJob but got", "type", obj)
return
}
age := time.Since(job.CreationTimestamp.Time)
if age > c.expirationTime {
namespace := job.Namespace
ctx, _, err := identity.WithProvisioningIdentity(context.Background(), namespace)
if err != nil {
c.logger.Error("Failed to set provisioning identity for cleanup", "error", err)
return
}
ctx = request.WithNamespace(ctx, namespace)
err = c.client.HistoricJobs(job.Namespace).Delete(ctx, job.Name, metav1.DeleteOptions{})
if err != nil && !apierrors.IsNotFound(err) {
c.logger.Error("Failed to delete expired HistoryJob",
"namespace", job.Namespace,
"name", job.Name,
"age", age,
"error", err)
} else {
c.logger.Info("Deleted expired HistoryJob",
"namespace", job.Namespace,
"name", job.Name,
"age", age)
}
}
}
+58
View File
@@ -0,0 +1,58 @@
package controller
import (
"k8s.io/client-go/tools/cache"
"github.com/grafana/grafana-app-sdk/logging"
informer "github.com/grafana/grafana/apps/provisioning/pkg/generated/informers/externalversions/provisioning/v0alpha1"
)
const (
jobControllerLoggerName = "provisioning-job-controller"
)
// JobController manages job create notifications.
type JobController struct {
jobSynced cache.InformerSynced
logger logging.Logger
// notification channel for job create events (replaces InsertNotifications)
notifications chan struct{}
}
// NewJobController creates a new JobController.
func NewJobController(
jobInformer informer.JobInformer,
) (*JobController, error) {
jc := &JobController{
jobSynced: jobInformer.Informer().HasSynced,
logger: logging.DefaultLogger.With("logger", jobControllerLoggerName),
notifications: make(chan struct{}, 1),
}
_, err := jobInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
// Send notification for job create events (replaces InsertNotifications)
jc.sendNotification()
},
})
if err != nil {
return nil, err
}
return jc, nil
}
// InsertNotifications returns a channel that receives notifications when jobs are created.
// This replaces the InsertNotifications method from persistentstore.go.
func (jc *JobController) InsertNotifications() chan struct{} {
return jc.notifications
}
func (jc *JobController) sendNotification() {
select {
case jc.notifications <- struct{}{}:
default:
// Don't block if there's already a notification waiting
}
}
@@ -0,0 +1,85 @@
package controller
import (
"context"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
provisioning "github.com/grafana/grafana/apps/provisioning/pkg/apis/provisioning/v0alpha1"
provisioningfake "github.com/grafana/grafana/apps/provisioning/pkg/generated/clientset/versioned/fake"
provisioninginformers "github.com/grafana/grafana/apps/provisioning/pkg/generated/informers/externalversions"
)
func TestJobController_New(t *testing.T) {
client := provisioningfake.NewSimpleClientset()
informerFactory := provisioninginformers.NewSharedInformerFactory(client, time.Second)
jobInformer := informerFactory.Provisioning().V0alpha1().Jobs()
controller, err := NewJobController(jobInformer)
require.NoError(t, err)
assert.NotNil(t, controller)
assert.NotNil(t, controller.notifications)
}
func TestJobController_InsertNotifications(t *testing.T) {
client := provisioningfake.NewSimpleClientset()
informerFactory := provisioninginformers.NewSharedInformerFactory(client, time.Second)
jobInformer := informerFactory.Provisioning().V0alpha1().Jobs()
controller, err := NewJobController(jobInformer)
require.NoError(t, err)
notifications := controller.InsertNotifications()
assert.NotNil(t, notifications)
// Test that notification is sent
controller.sendNotification()
select {
case <-notifications:
// Success - notification received
case <-time.After(time.Second):
t.Fatal("Expected notification but didn't receive one")
}
}
func TestJobController_NotificationOnJobCreate(t *testing.T) {
client := provisioningfake.NewSimpleClientset()
informerFactory := provisioninginformers.NewSharedInformerFactory(client, time.Second)
jobInformer := informerFactory.Provisioning().V0alpha1().Jobs()
controller, err := NewJobController(jobInformer)
require.NoError(t, err)
// Start informer and wait for cache sync
ctx, cancel := context.WithTimeout(context.Background(), time.Second*5)
defer cancel()
informerFactory.Start(ctx.Done())
informerFactory.WaitForCacheSync(ctx.Done())
// Get notifications channel
notifications := controller.InsertNotifications()
// Create a job - this should trigger a notification
_, err = client.ProvisioningV0alpha1().Jobs("default").Create(ctx, &provisioning.Job{
ObjectMeta: metav1.ObjectMeta{
Name: "test-job",
Namespace: "default",
},
}, metav1.CreateOptions{})
require.NoError(t, err)
// Wait for notification
select {
case <-notifications:
// Success - notification received
case <-time.After(time.Second * 2):
t.Fatal("Expected notification but didn't receive one")
}
}