From 4b438773240f4ecf275aa5f6b7c6af3fb7a9e8a5 Mon Sep 17 00:00:00 2001 From: Ezequiel Victorero Date: Thu, 28 Aug 2025 17:40:45 -0300 Subject: [PATCH] ShortURL: Use the k8s API in the cleanup process (#109938) --- pkg/api/short_url_test.go | 4 + pkg/registry/apps/shorturl/legacy_storage.go | 81 ++++++++---- pkg/server/wire_gen.go | 4 +- pkg/services/cleanup/cleanup.go | 120 ++++++++++++++++-- pkg/services/shorturls/shorturl.go | 1 + .../shorturls/shorturlimpl/shorturl.go | 4 + pkg/services/shorturls/shorturlimpl/store.go | 17 +++ 7 files changed, 192 insertions(+), 39 deletions(-) diff --git a/pkg/api/short_url_test.go b/pkg/api/short_url_test.go index 09059dc91c9..6d3da5fbe0c 100644 --- a/pkg/api/short_url_test.go +++ b/pkg/api/short_url_test.go @@ -85,6 +85,10 @@ type fakeShortURLService struct { createConvertShortURLToDTO func(shortURL *shorturls.ShortUrl, appURL string) *dtos.ShortURL } +func (s *fakeShortURLService) List(ctx context.Context, orgID int64) ([]*shorturls.ShortUrl, error) { + return nil, nil +} + func (s *fakeShortURLService) GetShortURLByUID(ctx context.Context, user *user.SignedInUser, uid string) (*shorturls.ShortUrl, error) { return nil, nil } diff --git a/pkg/registry/apps/shorturl/legacy_storage.go b/pkg/registry/apps/shorturl/legacy_storage.go index 35f538b5806..8ee7a59e899 100644 --- a/pkg/registry/apps/shorturl/legacy_storage.go +++ b/pkg/registry/apps/shorturl/legacy_storage.go @@ -61,7 +61,25 @@ func (s *legacyStorage) ConvertToTable(ctx context.Context, object runtime.Objec } func (s *legacyStorage) List(ctx context.Context, options *internalversion.ListOptions) (runtime.Object, error) { - return nil, fmt.Errorf("List for shorturl not implemented") + orgID, err := request.OrgIDForList(ctx) + if err != nil { + return nil, err + } + + shortURLs, err := s.service.List(ctx, orgID) + if err != nil { + if errors.Is(err, shorturls.ErrShortURLNotFound) { + return shorturl.ShortURLKind().ZeroListValue(), nil // return empty list if no short URLs found + } + return nil, err + } + + list := &shorturl.ShortURLList{} + for idx := range shortURLs { + list.Items = append(list.Items, *convertToK8sResource(shortURLs[idx], s.namespacer)) + } + + return list, nil } func (s *legacyStorage) Get(ctx context.Context, name string, options *metav1.GetOptions) (runtime.Object, error) { @@ -70,14 +88,10 @@ func (s *legacyStorage) Get(ctx context.Context, name string, options *metav1.Ge return nil, err } - // Convert identity.Requester to *user.SignedInUser - var signedInUser *user.SignedInUser - if authnIdentity, ok := requester.(*authn.Identity); ok { - signedInUser = authnIdentity.SignedInUser() - } else if userIdentity, ok := requester.(*user.SignedInUser); ok { - signedInUser = userIdentity - } else { - return nil, fmt.Errorf("unsupported identity type") + // Convert any identity.Requester to *user.SignedInUser + signedInUser, err := convertRequesterToSignedInUser(requester) + if err != nil { + return nil, fmt.Errorf("failed to convert requester: %w", err) } dto, err := s.service.GetShortURLByUID(ctx, signedInUser, name) @@ -103,14 +117,11 @@ func (s *legacyStorage) Create(ctx context.Context, if err != nil { return nil, err } - // Convert identity.Requester to *user.SignedInUser - var signedInUser *user.SignedInUser - if authnIdentity, ok := requester.(*authn.Identity); ok { - signedInUser = authnIdentity.SignedInUser() - } else if userIdentity, ok := requester.(*user.SignedInUser); ok { - signedInUser = userIdentity - } else { - return nil, fmt.Errorf("unsupported identity type") + + // Convert any identity.Requester to *user.SignedInUser + signedInUser, err := convertRequesterToSignedInUser(requester) + if err != nil { + return nil, fmt.Errorf("failed to convert requester: %w", err) } if createValidation != nil { @@ -147,14 +158,10 @@ func (s *legacyStorage) Update(ctx context.Context, return nil, false, err } - // Convert identity.Requester to *user.SignedInUser - var signedInUser *user.SignedInUser - if authnIdentity, ok := requester.(*authn.Identity); ok { - signedInUser = authnIdentity.SignedInUser() - } else if userIdentity, ok := requester.(*user.SignedInUser); ok { - signedInUser = userIdentity - } else { - return nil, false, fmt.Errorf("unsupported identity type") + // Convert any identity.Requester to *user.SignedInUser + signedInUser, err := convertRequesterToSignedInUser(requester) + if err != nil { + return nil, false, fmt.Errorf("failed to convert requester: %w", err) } shortURL, err := s.service.GetShortURLByUID(ctx, signedInUser, name) @@ -199,3 +206,27 @@ func (s *legacyStorage) Delete(ctx context.Context, name string, deleteValidatio func (s *legacyStorage) DeleteCollection(ctx context.Context, deleteValidation rest.ValidateObjectFunc, options *metav1.DeleteOptions, listOptions *internalversion.ListOptions) (runtime.Object, error) { return nil, fmt.Errorf("DeleteCollection for shorturl not implemented") } + +// convertRequesterToSignedInUser converts any identity.Requester to *user.SignedInUser +// This is needed because some legacy shorturls service methods still expect SignedInUser +func convertRequesterToSignedInUser(requester identity.Requester) (*user.SignedInUser, error) { + // If it's already a SignedInUser, return it directly + if signedInUser, ok := requester.(*user.SignedInUser); ok { + return signedInUser, nil + } + + // If it's a StaticRequester (service identity), convert it + if staticRequester, ok := requester.(*identity.StaticRequester); ok { + return &user.SignedInUser{ + UserID: staticRequester.UserID, // Used for CreatedBy field + OrgID: staticRequester.OrgID, // Used in SQL queries + }, nil + } + + // If it's an authn.Identity, use its SignedInUser method + if authnIdentity, ok := requester.(*authn.Identity); ok { + return authnIdentity.SignedInUser(), nil + } + + return nil, fmt.Errorf("unsupported identity type") +} diff --git a/pkg/server/wire_gen.go b/pkg/server/wire_gen.go index 65a13359a15..c760b9d8e57 100644 --- a/pkg/server/wire_gen.go +++ b/pkg/server/wire_gen.go @@ -590,7 +590,7 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api deleteExpiredService := image.ProvideDeleteExpiredService(dBstore) tempuserService := tempuserimpl.ProvideService(sqlStore, cfg) cleanupServiceImpl := annotationsimpl.ProvideCleanupService(sqlStore, cfg) - cleanUpService := cleanup.ProvideService(cfg, serverLockService, shortURLService, sqlStore, queryHistoryService, dashverService, serviceImpl, deleteExpiredService, tempuserService, tracingService, cleanupServiceImpl, dashboardService, dBstore) + cleanUpService := cleanup.ProvideService(cfg, featureToggles, serverLockService, shortURLService, sqlStore, queryHistoryService, dashverService, serviceImpl, deleteExpiredService, tempuserService, tracingService, cleanupServiceImpl, dBstore, eventualRestConfigProvider, orgService) secretsKVStore, err := kvstore2.ProvideService(sqlStore, secretsService) if err != nil { return nil, err @@ -1165,7 +1165,7 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac deleteExpiredService := image.ProvideDeleteExpiredService(dBstore) tempuserService := tempuserimpl.ProvideService(sqlStore, cfg) cleanupServiceImpl := annotationsimpl.ProvideCleanupService(sqlStore, cfg) - cleanUpService := cleanup.ProvideService(cfg, serverLockService, shortURLService, sqlStore, queryHistoryService, dashverService, serviceImpl, deleteExpiredService, tempuserService, tracingService, cleanupServiceImpl, dashboardService, dBstore) + cleanUpService := cleanup.ProvideService(cfg, featureToggles, serverLockService, shortURLService, sqlStore, queryHistoryService, dashverService, serviceImpl, deleteExpiredService, tempuserService, tracingService, cleanupServiceImpl, dBstore, eventualRestConfigProvider, orgService) secretsKVStore, err := kvstore2.ProvideService(sqlStore, secretsService) if err != nil { return nil, err diff --git a/pkg/services/cleanup/cleanup.go b/pkg/services/cleanup/cleanup.go index ac045743b75..3f5890a29f6 100644 --- a/pkg/services/cleanup/cleanup.go +++ b/pkg/services/cleanup/cleanup.go @@ -11,16 +11,26 @@ import ( "time" "go.opentelemetry.io/otel/attribute" + k8serrors "k8s.io/apimachinery/pkg/api/errors" + v1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/client-go/dynamic" + "github.com/grafana/grafana/apps/shorturl/pkg/apis/shorturl/v1alpha1" + "github.com/grafana/grafana/pkg/apimachinery/identity" "github.com/grafana/grafana/pkg/infra/db" "github.com/grafana/grafana/pkg/infra/log" "github.com/grafana/grafana/pkg/infra/serverlock" "github.com/grafana/grafana/pkg/infra/tracing" "github.com/grafana/grafana/pkg/services/annotations" - "github.com/grafana/grafana/pkg/services/dashboards" + grafanaapiserver "github.com/grafana/grafana/pkg/services/apiserver" + "github.com/grafana/grafana/pkg/services/apiserver/endpoints/request" "github.com/grafana/grafana/pkg/services/dashboardsnapshots" dashver "github.com/grafana/grafana/pkg/services/dashboardversion" + "github.com/grafana/grafana/pkg/services/featuremgmt" "github.com/grafana/grafana/pkg/services/ngalert/image" + "github.com/grafana/grafana/pkg/services/org" "github.com/grafana/grafana/pkg/services/queryhistory" "github.com/grafana/grafana/pkg/services/shorturls" tempuser "github.com/grafana/grafana/pkg/services/temp_user" @@ -36,6 +46,7 @@ type CleanUpService struct { tracer tracing.Tracer store db.DB Cfg *setting.Cfg + Features featuremgmt.FeatureToggles ServerLockService *serverlock.ServerLockService ShortURLService shorturls.Service QueryHistoryService queryhistory.Service @@ -44,16 +55,18 @@ type CleanUpService struct { deleteExpiredImageService *image.DeleteExpiredService tempUserService tempuser.Service annotationCleaner annotations.Cleaner - dashboardService dashboards.DashboardService alertRuleService AlertRuleService + clientConfigProvider grafanaapiserver.RestConfigProvider + orgService org.Service } -func ProvideService(cfg *setting.Cfg, serverLockService *serverlock.ServerLockService, +func ProvideService(cfg *setting.Cfg, Features featuremgmt.FeatureToggles, serverLockService *serverlock.ServerLockService, shortURLService shorturls.Service, sqlstore db.DB, queryHistoryService queryhistory.Service, dashboardVersionService dashver.Service, dashSnapSvc dashboardsnapshots.Service, deleteExpiredImageService *image.DeleteExpiredService, - tempUserService tempuser.Service, tracer tracing.Tracer, annotationCleaner annotations.Cleaner, dashboardService dashboards.DashboardService, service AlertRuleService) *CleanUpService { + tempUserService tempuser.Service, tracer tracing.Tracer, annotationCleaner annotations.Cleaner, service AlertRuleService, clientConfigProvider grafanaapiserver.RestConfigProvider, orgService org.Service) *CleanUpService { s := &CleanUpService{ Cfg: cfg, + Features: Features, ServerLockService: serverLockService, ShortURLService: shortURLService, QueryHistoryService: queryHistoryService, @@ -65,8 +78,9 @@ func ProvideService(cfg *setting.Cfg, serverLockService *serverlock.ServerLockSe tempUserService: tempUserService, tracer: tracer, annotationCleaner: annotationCleaner, - dashboardService: dashboardService, alertRuleService: service, + clientConfigProvider: clientConfigProvider, + orgService: orgService, } return s } @@ -83,7 +97,7 @@ func (j cleanUpJob) String() string { func (srv *CleanUpService) Run(ctx context.Context) error { srv.cleanUpTmpFiles(ctx) - ticker := time.NewTicker(time.Minute * 10) + ticker := time.NewTicker(time.Minute * 1) for { select { case <-ticker.C: @@ -272,16 +286,98 @@ func (srv *CleanUpService) expireOldVerifications(ctx context.Context) { func (srv *CleanUpService) deleteStaleShortURLs(ctx context.Context) { logger := srv.log.FromContext(ctx) - cmd := shorturls.DeleteShortUrlCommand{ - OlderThan: time.Now().Add(-time.Duration(srv.Cfg.ShortLinkExpiration*24) * time.Hour), - } - if err := srv.ShortURLService.DeleteStaleShortURLs(ctx, &cmd); err != nil { - logger.Error("Problem deleting stale short urls", "error", err.Error()) + if srv.Features.IsEnabledGlobally(featuremgmt.FlagKubernetesShortURLs) { + srv.deleteStaleKubernetesShortURLs(ctx) } else { - logger.Debug("Deleted short urls", "rows affected", cmd.NumDeleted) + cmd := shorturls.DeleteShortUrlCommand{ + OlderThan: time.Now().Add(-time.Duration(srv.Cfg.ShortLinkExpiration*24) * time.Hour), + } + if err := srv.ShortURLService.DeleteStaleShortURLs(ctx, &cmd); err != nil { + logger.Error("Problem deleting stale short urls", "error", err.Error()) + } else { + logger.Debug("Deleted short urls", "rows affected", cmd.NumDeleted) + } } } +func (srv *CleanUpService) deleteStaleKubernetesShortURLs(ctx context.Context) { + logger := srv.log.FromContext(ctx) + logger.Debug("Starting deleting expired Kubernetes shortURLs") + + // Create the dynamic client for Kubernetes API + restConfig, err := srv.clientConfigProvider.GetRestConfig(ctx) + if err != nil { + logger.Error("Failed to get REST config for Kubernetes client", "error", err.Error()) + return + } + + client, err := dynamic.NewForConfig(restConfig) + if err != nil { + logger.Error("Failed to create Kubernetes client", "error", err.Error()) + return + } + + // Set up the GroupVersionResource for shortURLs + gvr := schema.GroupVersionResource{ + Group: v1alpha1.ShortURLKind().Group(), + Version: v1alpha1.ShortURLKind().Version(), + Resource: v1alpha1.ShortURLKind().Plural(), + } + + // Calculate the expiration time + expirationTime := time.Now().Add(-time.Duration(srv.Cfg.ShortLinkExpiration*24) * time.Hour) + expirationTimestamp := expirationTime.Unix() + deletedCount := 0 + + // List and delete expired shortURLs across all namespaces + orgs, err := srv.orgService.Search(ctx, &org.SearchOrgsQuery{}) + if err != nil { + logger.Error("Failed to list organizations", "error", err.Error()) + return + } + + for _, o := range orgs { + ctx, _ := identity.WithServiceIdentity(ctx, o.ID) + namespaceMapper := request.GetNamespaceMapper(srv.Cfg) + shortURLs, err := client.Resource(gvr).Namespace(namespaceMapper(o.ID)).List(ctx, v1.ListOptions{}) + if err != nil { + logger.Error("Failed to list shortURLs", "error", err.Error()) + return + } + // Check each shortURL for expiration + for _, item := range shortURLs.Items { + // Convert unstructured object to ShortURL struct + var shortURL v1alpha1.ShortURL + err := runtime.DefaultUnstructuredConverter.FromUnstructured(item.Object, &shortURL) + if err != nil { + logger.Error("Failed to convert unstructured object to ShortURL", "name", item.GetName(), "namespace", item.GetNamespace(), "error", err.Error()) + continue + } + + // Only delete if lastSeenAt is 0 (meaning it has not been accessed) and the creation time is older than the expiration time + if shortURL.Status.LastSeenAt == 0 && shortURL.CreationTimestamp.Unix() < expirationTimestamp { + namespace := shortURL.Namespace + name := shortURL.Name + + err := client.Resource(gvr).Namespace(namespace).Delete(ctx, name, v1.DeleteOptions{}) + if err != nil { + // Check if it's a "not found" error, which is expected if the resource was already deleted + if k8serrors.IsNotFound(err) { + logger.Debug("ShortURL already deleted", "name", name, "namespace", namespace) + } else { + logger.Error("Failed to delete expired shortURL", "name", name, "namespace", namespace, "error", err.Error()) + } + } else { + deletedCount++ + logger.Debug("Successfully deleted expired shortURL", "name", name, "namespace", namespace, "creationTime", shortURL.CreationTimestamp.Unix(), "expirationTime", expirationTimestamp) + } + } + } + } + + logger.Debug("Deleted expired Kubernetes shortURLs", "count", deletedCount) +} + func (srv *CleanUpService) deleteStaleQueryHistory(ctx context.Context) { logger := srv.log.FromContext(ctx) // Delete query history from 14+ days ago with exception of starred queries diff --git a/pkg/services/shorturls/shorturl.go b/pkg/services/shorturls/shorturl.go index 14852e9432a..d62f469f85c 100644 --- a/pkg/services/shorturls/shorturl.go +++ b/pkg/services/shorturls/shorturl.go @@ -13,4 +13,5 @@ type Service interface { UpdateLastSeenAt(ctx context.Context, shortURL *ShortUrl) error DeleteStaleShortURLs(ctx context.Context, cmd *DeleteShortUrlCommand) error ConvertShortURLToDTO(shortURL *ShortUrl, appURL string) *dtos.ShortURL + List(ctx context.Context, orgID int64) ([]*ShortUrl, error) } diff --git a/pkg/services/shorturls/shorturlimpl/shorturl.go b/pkg/services/shorturls/shorturlimpl/shorturl.go index 9f83ef70d66..60d4496e979 100644 --- a/pkg/services/shorturls/shorturlimpl/shorturl.go +++ b/pkg/services/shorturls/shorturlimpl/shorturl.go @@ -37,6 +37,10 @@ func (s ShortURLService) UpdateLastSeenAt(ctx context.Context, shortURL *shortur return s.SQLStore.Update(ctx, shortURL) } +func (s ShortURLService) List(ctx context.Context, orgID int64) ([]*shorturls.ShortUrl, error) { + return s.SQLStore.List(ctx, orgID) +} + func (s ShortURLService) CreateShortURL(ctx context.Context, user *user.SignedInUser, cmd *dtos.CreateShortURLCmd) (*shorturls.ShortUrl, error) { relPath := strings.TrimSpace(cmd.Path) diff --git a/pkg/services/shorturls/shorturlimpl/store.go b/pkg/services/shorturls/shorturlimpl/store.go index 0fe54fc677d..6e0a072d23c 100644 --- a/pkg/services/shorturls/shorturlimpl/store.go +++ b/pkg/services/shorturls/shorturlimpl/store.go @@ -13,6 +13,7 @@ type store interface { Update(ctx context.Context, shortURL *shorturls.ShortUrl) error Insert(ctx context.Context, shortURL *shorturls.ShortUrl) error Delete(ctx context.Context, cmd *shorturls.DeleteShortUrlCommand) error + List(ctx context.Context, orgID int64) ([]*shorturls.ShortUrl, error) } type sqlStore struct { @@ -84,3 +85,19 @@ func (s sqlStore) Delete(ctx context.Context, cmd *shorturls.DeleteShortUrlComma return nil }) } + +func (s sqlStore) List(ctx context.Context, orgID int64) ([]*shorturls.ShortUrl, error) { + var shortURLs []*shorturls.ShortUrl + err := s.db.WithDbSession(ctx, func(dbSession *db.Session) error { + err := dbSession.Where("org_id = ?", orgID).Find(&shortURLs) + if err != nil { + return err + } + return nil + }) + if err != nil { + return nil, err + } + + return shortURLs, nil +}