diff --git a/.betterer.results b/.betterer.results index df376f544fa..711fa68a381 100644 --- a/.betterer.results +++ b/.betterer.results @@ -5752,10 +5752,8 @@ exports[`better eslint`] = { [0, 0, 0, "No untranslated strings in text props. Wrap text with or use t()", "6"], [0, 0, 0, "No untranslated strings in text props. Wrap text with or use t()", "7"], [0, 0, 0, "No untranslated strings in text props. Wrap text with or use t()", "8"], - [0, 0, 0, "No untranslated strings in text props. Wrap text with or use t()", "9"], - [0, 0, 0, "No untranslated strings in text props. Wrap text with or use t()", "10"], - [0, 0, 0, "No untranslated strings. Wrap text with ", "11"], - [0, 0, 0, "No untranslated strings. Wrap text with ", "12"] + [0, 0, 0, "No untranslated strings. Wrap text with ", "9"], + [0, 0, 0, "No untranslated strings. Wrap text with ", "10"] ], "public/app/features/provisioning/FileHistoryPage.tsx:5381": [ [0, 0, 0, "No untranslated strings. Wrap text with ", "0"], @@ -5782,8 +5780,7 @@ exports[`better eslint`] = { [0, 0, 0, "No untranslated strings in text props. Wrap text with or use t()", "1"], [0, 0, 0, "No untranslated strings. Wrap text with ", "2"], [0, 0, 0, "No untranslated strings. Wrap text with ", "3"], - [0, 0, 0, "No untranslated strings. Wrap text with ", "4"], - [0, 0, 0, "No untranslated strings. Wrap text with ", "5"] + [0, 0, 0, "No untranslated strings. Wrap text with ", "4"] ], "public/app/features/provisioning/RepositoryHealth.tsx:5381": [ [0, 0, 0, "No untranslated strings in text props. Wrap text with or use t()", "0"], @@ -5793,12 +5790,14 @@ exports[`better eslint`] = { ], "public/app/features/provisioning/RepositoryListPage.tsx:5381": [ [0, 0, 0, "No untranslated strings in text props. Wrap text with or use t()", "0"], - [0, 0, 0, "No untranslated strings. Wrap text with ", "1"], + [0, 0, 0, "No untranslated strings in text props. Wrap text with or use t()", "1"], [0, 0, 0, "No untranslated strings. Wrap text with ", "2"], [0, 0, 0, "No untranslated strings. Wrap text with ", "3"], [0, 0, 0, "No untranslated strings. Wrap text with ", "4"], [0, 0, 0, "No untranslated strings. Wrap text with ", "5"], - [0, 0, 0, "No untranslated strings. Wrap text with ", "6"] + [0, 0, 0, "No untranslated strings. Wrap text with ", "6"], + [0, 0, 0, "No untranslated strings. Wrap text with ", "7"], + [0, 0, 0, "No untranslated strings. Wrap text with ", "8"] ], "public/app/features/provisioning/RepositoryOverview.tsx:5381": [ [0, 0, 0, "No untranslated strings. Wrap text with ", "0"], diff --git a/pkg/apimachinery/identity/context.go b/pkg/apimachinery/identity/context.go index 598abbeaae2..b5a4484936a 100644 --- a/pkg/apimachinery/identity/context.go +++ b/pkg/apimachinery/identity/context.go @@ -111,6 +111,9 @@ var serviceIdentityPermissions = getWildcardPermissions( "datasources:delete", "alert.provisioning:write", "alert.provisioning.secrets:read", + "users:read", // accesscontrol.ActionUsersRead, + "org.users:read", // accesscontrol.ActionOrgUsersRead, + "teams:read", // accesscontrol.ActionTeamsRead, ) func IsServiceIdentity(ctx context.Context) bool { diff --git a/pkg/apis/provisioning/v0alpha1/jobs.go b/pkg/apis/provisioning/v0alpha1/jobs.go index a4e01d98743..86643858a9f 100644 --- a/pkg/apis/provisioning/v0alpha1/jobs.go +++ b/pkg/apis/provisioning/v0alpha1/jobs.go @@ -136,6 +136,7 @@ func (in *JobStatus) ToSyncStatus(jobId string) SyncStatus { type JobResourceSummary struct { Group string `json:"group,omitempty"` Resource string `json:"resource,omitempty"` + Total int64 `json:"total,omitempty"` // the count (if known) Create int64 `json:"create,omitempty"` Update int64 `json:"update,omitempty"` diff --git a/pkg/apis/provisioning/v0alpha1/settings.go b/pkg/apis/provisioning/v0alpha1/settings.go index 65ad68fd7dc..63d202f7ee1 100644 --- a/pkg/apis/provisioning/v0alpha1/settings.go +++ b/pkg/apis/provisioning/v0alpha1/settings.go @@ -9,6 +9,11 @@ import ( type RepositoryViewList struct { metav1.TypeMeta `json:",inline"` + // The backend is using legacy storage + // FIXME: Not sure where this should be exposed... but we need it somewhere + // The UI should force the onboarding workflow when this is true + LegacyStorage bool `json:"legacyStorage,omitempty"` + // +mapType=atomic Items []RepositoryView `json:"items"` } diff --git a/pkg/apis/provisioning/v0alpha1/zz_generated.openapi.go b/pkg/apis/provisioning/v0alpha1/zz_generated.openapi.go index 93e3e4bcf76..fd73980698f 100644 --- a/pkg/apis/provisioning/v0alpha1/zz_generated.openapi.go +++ b/pkg/apis/provisioning/v0alpha1/zz_generated.openapi.go @@ -569,6 +569,12 @@ func schema_pkg_apis_provisioning_v0alpha1_JobResourceSummary(ref common.Referen Format: "", }, }, + "total": { + SchemaProps: spec.SchemaProps{ + Type: []string{"integer"}, + Format: "int64", + }, + }, "create": { SchemaProps: spec.SchemaProps{ Type: []string{"integer"}, @@ -1156,6 +1162,13 @@ func schema_pkg_apis_provisioning_v0alpha1_RepositoryViewList(ref common.Referen Format: "", }, }, + "legacyStorage": { + SchemaProps: spec.SchemaProps{ + Description: "The backend is using legacy storage FIXME: Not sure where this should be exposed... but we need it somewhere The UI should force the onboarding workflow when this is true", + Type: []string{"boolean"}, + Format: "", + }, + }, "items": { VendorExtensible: spec.VendorExtensible{ Extensions: spec.Extensions{ diff --git a/pkg/registry/apis/dashboard/legacy/migrate.go b/pkg/registry/apis/dashboard/legacy/migrate.go index c1236ef4f2d..5a5cb951aa5 100644 --- a/pkg/registry/apis/dashboard/legacy/migrate.go +++ b/pkg/registry/apis/dashboard/legacy/migrate.go @@ -14,7 +14,10 @@ import ( "github.com/grafana/grafana/pkg/apimachinery/utils" dashboard "github.com/grafana/grafana/pkg/apis/dashboard" folders "github.com/grafana/grafana/pkg/apis/folder/v0alpha1" + "github.com/grafana/grafana/pkg/infra/db" + "github.com/grafana/grafana/pkg/services/provisioning" "github.com/grafana/grafana/pkg/services/sqlstore" + "github.com/grafana/grafana/pkg/storage/legacysql" "github.com/grafana/grafana/pkg/storage/unified/apistore" "github.com/grafana/grafana/pkg/storage/unified/resource" ) @@ -22,7 +25,6 @@ import ( type MigrateOptions struct { Namespace string Store resource.BatchStoreClient - Writer resource.BatchResourceWriter LargeObjects apistore.LargeObjectSupport BlobStore resource.BlobStoreClient Resources []schema.GroupResource @@ -36,6 +38,15 @@ type LegacyMigrator interface { Migrate(ctx context.Context, opts MigrateOptions) (*resource.BatchResponse, error) } +// This can migrate Folders, Dashboards and LibraryPanels +func ProvideLegacyMigrator( + sql db.DB, // direct access to tables + provisioning provisioning.ProvisioningService, // only needed for dashboard settings +) LegacyMigrator { + dbp := legacysql.NewDatabaseProvider(sql) + return NewDashboardAccess(dbp, authlib.OrgNamespaceFormatter, nil, provisioning, false) +} + type BlobStoreInfo struct { Count int64 Size int64 @@ -49,6 +60,9 @@ func (a *dashboardSqlAccess) Migrate(ctx context.Context, opts MigrateOptions) ( if err != nil { return nil, err } + if opts.Progress == nil { + opts.Progress = func(count int, msg string) {} // noop + } // Migrate everything if len(opts.Resources) < 1 { diff --git a/pkg/registry/apis/provisioning/controller/repository.go b/pkg/registry/apis/provisioning/controller/repository.go index 2cbbd48e5af..08a5f7a6920 100644 --- a/pkg/registry/apis/provisioning/controller/repository.go +++ b/pkg/registry/apis/provisioning/controller/repository.go @@ -346,8 +346,10 @@ func (rc *RepositoryController) process(item *queueItem) error { } sync = &provisioning.SyncJobOptions{} case shouldResync: - logger.Info("handle repository resync") - sync = &provisioning.SyncJobOptions{Incremental: true} + if obj.Spec.Sync.Enabled { + logger.Info("handle repository resync") + sync = &provisioning.SyncJobOptions{Incremental: true} + } default: logger.Info("handle unknown repository situation") } diff --git a/pkg/registry/apis/provisioning/jobs/export/folders.go b/pkg/registry/apis/provisioning/jobs/export/folders.go index 0d2b640178f..c47e667e130 100644 --- a/pkg/registry/apis/provisioning/jobs/export/folders.go +++ b/pkg/registry/apis/provisioning/jobs/export/folders.go @@ -1,39 +1,141 @@ package export import ( + "context" + "errors" "fmt" - "golang.org/x/net/context" + apierrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" - "k8s.io/client-go/dynamic" + "k8s.io/apimachinery/pkg/runtime/schema" - apiutils "github.com/grafana/grafana/pkg/apimachinery/utils" - "github.com/grafana/grafana/pkg/cmd/grafana-cli/logger" + folders "github.com/grafana/grafana/pkg/apis/folder/v0alpha1" + provisioning "github.com/grafana/grafana/pkg/apis/provisioning/v0alpha1" + "github.com/grafana/grafana/pkg/registry/apis/dashboard/legacy" + "github.com/grafana/grafana/pkg/registry/apis/provisioning/repository" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" + "github.com/grafana/grafana/pkg/storage/unified/parquet" + "github.com/grafana/grafana/pkg/storage/unified/resource" ) -func readFolders(ctx context.Context, client dynamic.ResourceInterface, skip string) (*resources.FolderTree, error) { - // TODO: handle pagination - rawList, err := client.List(ctx, metav1.ListOptions{Limit: 10000}) +var ( + _ resource.BatchResourceWriter = (*folderReader)(nil) +) + +type folderReader struct { + tree *resources.FolderTree + targetRepoName string + summary *provisioning.JobResourceSummary +} + +// Close implements resource.BatchResourceWriter. +func (f *folderReader) Close() error { + return nil +} + +// CloseWithResults implements resource.BatchResourceWriter. +func (f *folderReader) CloseWithResults() (*resource.BatchResponse, error) { + return &resource.BatchResponse{}, nil +} + +// Write implements resource.BatchResourceWriter. +func (f *folderReader) Write(ctx context.Context, key *resource.ResourceKey, value []byte) error { + item := &unstructured.Unstructured{} + err := item.UnmarshalJSON(value) if err != nil { - return nil, fmt.Errorf("failed to list folders: %w", err) + return err } - if rawList.GetContinue() != "" { - return nil, fmt.Errorf("unable to list all folders in one request: %s", rawList.GetContinue()) + err = f.tree.AddUnstructured(item, f.targetRepoName) + if err != nil { + f.summary.Errors = append(f.summary.Errors, err.Error()) + } + return nil +} + +func (r *exportJob) loadFolders(ctx context.Context) error { + logger := r.logger + status := r.jobStatus + status.Message = "reading folder tree" + r.maybeNotify(ctx) + + summary := r.getSummary(schema.GroupResource{ + Group: folders.GROUP, + Resource: folders.RESOURCE, + }) + + reader := &folderReader{ + tree: resources.NewEmptyFolderTree(), + targetRepoName: r.target.Config().Name, + summary: summary, } - // filter out the folders we already own - rawFolders := make([]unstructured.Unstructured, 0, len(rawList.Items)) - for _, f := range rawList.Items { - repoName := f.GetAnnotations()[apiutils.AnnoKeyRepoName] - if repoName == skip { - logger.Info("skip as folder is already in repository", "folder", f.GetName()) - continue + if r.legacy != nil { + _, err := r.legacy.Migrate(ctx, legacy.MigrateOptions{ + Namespace: r.namespace, + Resources: []schema.GroupResource{{ + Group: folders.GROUP, + Resource: folders.RESOURCE, + }}, + Store: parquet.NewBatchResourceWriterClient(reader), + }) + if err != nil { + return fmt.Errorf("unable to read folders from legacy storage %w", err) + } + } else { + client := r.client.Resource(schema.GroupVersionResource{ + Group: folders.GROUP, + Version: folders.VERSION, + Resource: folders.RESOURCE, + }) + + rawList, err := client.List(ctx, metav1.ListOptions{Limit: 10000}) + if err != nil { + return fmt.Errorf("failed to list folders: %w", err) + } + if rawList.GetContinue() != "" { + return fmt.Errorf("unable to list all folders in one request: %s", rawList.GetContinue()) + } + for _, item := range rawList.Items { + err = reader.tree.AddUnstructured(&item, reader.targetRepoName) + if err != nil { + summary.Errors = append(summary.Errors, err.Error()) + } + } + } + + // first create folders + // NOTE: this is required so that empty folders exist when finished + status.Message = "writing folders" + err := reader.tree.Walk(ctx, func(ctx context.Context, folder resources.Folder) error { + p := folder.Path + "/" + if r.prefix != "" { + p = r.prefix + "/" + p + } + logger := logger.With("path", p) + + _, err := r.target.Read(ctx, p, r.ref) + if err != nil && !(errors.Is(err, repository.ErrFileNotFound) || apierrors.IsNotFound(err)) { + logger.Error("failed to check if folder exists before writing", "error", err) + return fmt.Errorf("failed to check if folder exists before writing: %w", err) + } else if err == nil { + logger.Info("folder already exists") + summary.Noop++ + return nil } - rawFolders = append(rawFolders, f) + // Create with an empty body will make a folder (or .keep file if unsupported) + if err := r.target.Create(ctx, p, r.ref, nil, "export folder `"+p+"`"); err != nil { + logger.Error("failed to write a folder in repository", "error", err) + return fmt.Errorf("failed to write folder in repo: %w", err) + } + summary.Create++ + logger.Debug("successfully exported folder") + return nil + }) + if err != nil { + return fmt.Errorf("failed to write folders: %w", err) } - - return resources.NewFolderTreeFromUnstructure(ctx, rawFolders), nil + r.foldersTree = reader.tree + return nil } diff --git a/pkg/registry/apis/provisioning/jobs/export/job.go b/pkg/registry/apis/provisioning/jobs/export/job.go index 168406578f4..abe1e2cad54 100644 --- a/pkg/registry/apis/provisioning/jobs/export/job.go +++ b/pkg/registry/apis/provisioning/jobs/export/job.go @@ -3,21 +3,17 @@ package export import ( "context" "encoding/json" - "errors" "fmt" "time" - apierrors "k8s.io/apimachinery/pkg/api/errors" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime/schema" "github.com/grafana/grafana-app-sdk/logging" "github.com/grafana/grafana/pkg/apimachinery/utils" - folders "github.com/grafana/grafana/pkg/apis/folder/v0alpha1" provisioning "github.com/grafana/grafana/pkg/apis/provisioning/v0alpha1" - "github.com/grafana/grafana/pkg/cmd/grafana-cli/logger" "github.com/grafana/grafana/pkg/infra/slugify" + "github.com/grafana/grafana/pkg/registry/apis/dashboard/legacy" "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" "github.com/grafana/grafana/pkg/registry/apis/provisioning/repository" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" @@ -26,19 +22,22 @@ import ( // ExportJob holds all context for a running job type exportJob struct { - logger logging.Logger - client *resources.DynamicClient // Read from - target repository.Repository // Write to + logger logging.Logger + client *resources.DynamicClient // Read from + target repository.Repository // Write to + legacy legacy.LegacyMigrator + namespace string progress jobs.ProgressFn progressInterval time.Duration progressLast time.Time foldersTree *resources.FolderTree + userInfo map[string]repository.CommitSignature prefix string // from options (now clean+safe) ref string // from options (only git) keepIdentifier bool - addAuthorInfo bool + withHistory bool jobStatus *provisioning.JobStatus summary map[string]*provisioning.JobResourceSummary @@ -55,17 +54,18 @@ func newExportJob(ctx context.Context, prefix = safepath.Clean(prefix) } return &exportJob{ + namespace: target.Config().Namespace, target: target, client: client, logger: logging.FromContext(ctx), progress: progress, progressLast: time.Now(), - progressInterval: time.Second * 10, + progressInterval: time.Second * 5, prefix: prefix, ref: options.Branch, keepIdentifier: options.Identifier, - addAuthorInfo: options.History, + withHistory: options.History, jobStatus: &provisioning.JobStatus{ State: provisioning.JobStateWorking, @@ -77,6 +77,7 @@ func newExportJob(ctx context.Context, // Send progress messages to any listeners func (r *exportJob) maybeNotify(ctx context.Context) { if time.Since(r.progressLast) > r.progressInterval { + r.progressLast = time.Now() err := r.progress(ctx, *r.jobStatus) if err != nil { r.logger.Warn("unable to send progress", "err", err) @@ -98,86 +99,6 @@ func (r *exportJob) getSummary(gr schema.GroupResource) *provisioning.JobResourc return summary } -func (r *exportJob) loadFolders(ctx context.Context) error { - logger := r.logger - targetRepoName := r.target.Config().Name - status := r.jobStatus - status.Message = "reading folder tree" - foldersTree, err := readFolders(ctx, r.client.Resource(schema.GroupVersionResource{ - Group: folders.GROUP, - Version: folders.VERSION, - Resource: folders.RESOURCE, - }), targetRepoName) - - summary := r.getSummary(schema.GroupResource{ - Group: folders.GROUP, - Resource: folders.RESOURCE, - }) - - // first create folders - // TODO! this should not be necessary if writing to a path also makes the parents - status.Message = "writing folders" - err = foldersTree.Walk(ctx, func(ctx context.Context, folder resources.Folder) error { - p := folder.Path + "/" - if r.prefix != "" { - p = r.prefix + "/" + p - } - logger := logger.With("path", p) - - _, err = r.target.Read(ctx, p, r.ref) - if err != nil && !(errors.Is(err, repository.ErrFileNotFound) || apierrors.IsNotFound(err)) { - logger.Error("failed to check if folder exists before writing", "error", err) - return fmt.Errorf("failed to check if folder exists before writing: %w", err) - } else if err == nil { - logger.Info("folder already exists") - summary.Noop++ - return nil - } - - // Create with an empty body will make a folder (or .keep file if unsupported) - if err := r.target.Create(ctx, p, r.ref, nil, "export folder `"+p+"`"); err != nil { - logger.Error("failed to write a folder in repository", "error", err) - return fmt.Errorf("failed to write folder in repo: %w", err) - } - summary.Create++ - logger.Debug("successfully exported folder") - return nil - }) - if err != nil { - return fmt.Errorf("failed to write folders: %w", err) - } - r.foldersTree = foldersTree - return nil -} - -func (r *exportJob) export(ctx context.Context, kind schema.GroupVersionResource) error { - r.jobStatus.Message = "Exporting " + kind.Resource + "..." - r.maybeNotify(ctx) - client := r.client.Resource(kind) - summary := r.getSummary(kind.GroupResource()) - - continueToken := "" - for { - list, err := client.List(ctx, metav1.ListOptions{Limit: 100, Continue: continueToken}) - if err != nil { - return fmt.Errorf("error executing list: %w", err) - } - - for _, item := range list.Items { - if err = r.add(ctx, summary, &item); err != nil { - return fmt.Errorf("error adding value: %w", err) - } - } - - continueToken = list.GetContinue() - if continueToken == "" { - break - } - } - - return nil -} - func (r *exportJob) add(ctx context.Context, summary *provisioning.JobResourceSummary, obj *unstructured.Unstructured) error { if err := ctx.Err(); err != nil { return err @@ -194,7 +115,7 @@ func (r *exportJob) add(ctx context.Context, summary *provisioning.JobResourceSu if commitMessage == "" { g := item.GetGeneration() if g > 0 { - commitMessage = fmt.Sprintf("Generation: %d, ResourceVersion: %s", g, item.GetResourceVersion()) + commitMessage = fmt.Sprintf("Generation: %d", g) } else { commitMessage = "exported from grafana" } @@ -213,14 +134,21 @@ func (r *exportJob) add(ctx context.Context, summary *provisioning.JobResourceSu } folder := item.GetFolder() + // Add the author in context (if available) + ctx = r.withAuthorSignature(ctx, item) + // Get the absolute path of the folder fid, ok := r.foldersTree.DirPath(folder, "") if !ok { - logger.Error("folder of item was not in tree of repository") - return fmt.Errorf("folder of item was not in tree of repository") + fid = resources.Folder{ + Path: "__folder_not_found/" + slugify.Slugify(folder), + } + r.logger.Error("folder of item was not in tree of repository") } + // Clear the metadata delete(obj.Object, "metadata") + if r.keepIdentifier { item.SetName(name) // keep the identifier in the metadata } @@ -244,14 +172,11 @@ func (r *exportJob) add(ctx context.Context, summary *provisioning.JobResourceSu } } - // Add the author in context (if available) - ctx = r.withAuthorSignature(ctx, item) - // Write the file err = r.target.Write(ctx, fileName, r.ref, body, commitMessage) if err != nil { summary.Error++ - logger.Error("failed to write a file in repository", "error", err) + r.logger.Error("failed to write a file in repository", "error", err) if len(summary.Errors) < 20 { summary.Errors = append(summary.Errors, fmt.Sprintf("error writing: %s", fileName)) } @@ -263,24 +188,26 @@ func (r *exportJob) add(ctx context.Context, summary *provisioning.JobResourceSu } func (r *exportJob) withAuthorSignature(ctx context.Context, item utils.GrafanaMetaAccessor) context.Context { - if !r.addAuthorInfo { + if r.userInfo == nil { return ctx } + id := item.GetUpdatedBy() + if id == "" { + id = item.GetCreatedBy() + } + if id == "" { + id = "grafana" + } - sig := repository.CommitSignature{ - Name: item.GetUpdatedBy(), - When: item.GetCreationTimestamp().Time, + sig := r.userInfo[id] // lookup + if sig.Name == "" && sig.Email == "" { + sig.Name = id } - if sig.Name == "" { - sig.Name = item.GetCreatedBy() - } - if sig.Name == "" { - return ctx // no user info - } - // TODO: convert internal id to name+email - updated, _ := item.GetUpdatedTimestamp() - if updated != nil { - sig.When = *updated + t, err := item.GetUpdatedTimestamp() + if err == nil && t != nil { + sig.When = *t + } else { + sig.When = item.GetCreationTimestamp().Time } return repository.WithAuthorSignature(ctx, sig) } diff --git a/pkg/registry/apis/provisioning/jobs/export/resources.go b/pkg/registry/apis/provisioning/jobs/export/resources.go new file mode 100644 index 00000000000..a332b6ece41 --- /dev/null +++ b/pkg/registry/apis/provisioning/jobs/export/resources.go @@ -0,0 +1,132 @@ +package export + +import ( + "context" + "fmt" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + + "github.com/grafana/grafana-app-sdk/logging" + dashboards "github.com/grafana/grafana/pkg/apis/dashboard" + provisioning "github.com/grafana/grafana/pkg/apis/provisioning/v0alpha1" + "github.com/grafana/grafana/pkg/registry/apis/dashboard/legacy" + "github.com/grafana/grafana/pkg/storage/unified/parquet" + "github.com/grafana/grafana/pkg/storage/unified/resource" +) + +var ( + _ resource.BatchResourceWriter = (*resourceReader)(nil) +) + +type resourceReader struct { + job *exportJob + summary *provisioning.JobResourceSummary + logger logging.Logger +} + +// Close implements resource.BatchResourceWriter. +func (f *resourceReader) Close() error { + return nil +} + +// CloseWithResults implements resource.BatchResourceWriter. +func (f *resourceReader) CloseWithResults() (*resource.BatchResponse, error) { + return &resource.BatchResponse{}, nil +} + +// Write implements resource.BatchResourceWriter. +func (f *resourceReader) Write(ctx context.Context, key *resource.ResourceKey, value []byte) error { + item := &unstructured.Unstructured{} + err := item.UnmarshalJSON(value) + if err != nil { + return err + } + err = f.job.add(ctx, f.summary, item) + if err != nil { + f.logger.Warn("error adding from legacy", "name", key.Name, "err", err) + f.summary.Errors = append(f.summary.Errors, fmt.Sprintf("%s: %s", key.Name, err.Error())) + if len(f.summary.Errors) > 50 { + return err + } + } + return nil +} + +func (r *exportJob) loadResources(ctx context.Context) error { + kinds := []schema.GroupVersionResource{{ + Group: dashboards.GROUP, + Resource: dashboards.DASHBOARD_RESOURCE, + Version: "v1alpha1", + }} + + for _, kind := range kinds { + r.jobStatus.Message = "Exporting " + kind.Resource + "..." + if r.legacy != nil { + gr := kind.GroupResource() + reader := &resourceReader{ + summary: r.getSummary(gr), + job: r, + logger: r.logger, + } + opts := legacy.MigrateOptions{ + Namespace: r.namespace, + WithHistory: r.withHistory, + Resources: []schema.GroupResource{gr}, + Store: parquet.NewBatchResourceWriterClient(reader), + OnlyCount: true, // first get the count + } + stats, err := r.legacy.Migrate(ctx, opts) + if err != nil { + return fmt.Errorf("unable to count legacy items %w", err) + } + if len(stats.Summary) > 0 { + count := stats.Summary[0].Count + history := stats.Summary[0].History + if history > count { + count = history // the number of items we will process + } + reader.summary.Total = count + } + + opts.OnlyCount = false // this time actually write + _, err = r.legacy.Migrate(ctx, opts) + if err != nil { + return fmt.Errorf("error running legacy migrate %s %w", kind.Resource, err) + } + } + + if err := r.loadResourcesFromAPIServer(ctx, kind); err != nil { + return fmt.Errorf("error loading %s %w", kind.Resource, err) + } + } + return nil +} + +func (r *exportJob) loadResourcesFromAPIServer(ctx context.Context, kind schema.GroupVersionResource) error { + r.maybeNotify(ctx) + client := r.client.Resource(kind) + summary := r.getSummary(kind.GroupResource()) + + continueToken := "" + for { + list, err := client.List(ctx, metav1.ListOptions{Limit: 100, Continue: continueToken}) + if err != nil { + return fmt.Errorf("error executing list: %w", err) + } + + for _, item := range list.Items { + if err = r.add(ctx, summary, &item); err != nil { + return fmt.Errorf("error adding value: %w", err) + } + } + + continueToken = list.GetContinue() + if continueToken == "" { + break + } + } + + return nil +} diff --git a/pkg/registry/apis/provisioning/jobs/export/users.go b/pkg/registry/apis/provisioning/jobs/export/users.go new file mode 100644 index 00000000000..32548a172c4 --- /dev/null +++ b/pkg/registry/apis/provisioning/jobs/export/users.go @@ -0,0 +1,60 @@ +package export + +import ( + "context" + "fmt" + "strings" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + + iam "github.com/grafana/grafana/pkg/apis/iam/v0alpha1" + "github.com/grafana/grafana/pkg/registry/apis/provisioning/repository" +) + +func (r *exportJob) loadUsers(ctx context.Context) error { + status := r.jobStatus + status.Message = "reading user info" + r.maybeNotify(ctx) + + client := r.client.Resource(schema.GroupVersionResource{ + Group: iam.GROUP, + Version: iam.VERSION, + Resource: iam.UserResourceInfo.GroupResource().Resource, + }) + + rawList, err := client.List(ctx, metav1.ListOptions{Limit: 10000}) + if err != nil { + return fmt.Errorf("failed to list users: %w", err) + } + if rawList.GetContinue() != "" { + return fmt.Errorf("unable to list all users in one request: %s", rawList.GetContinue()) + } + + var ok bool + r.userInfo = make(map[string]repository.CommitSignature) + for _, item := range rawList.Items { + sig := repository.CommitSignature{} + sig.Name, ok, err = unstructured.NestedString(item.Object, "spec", "login") + if !ok || err != nil { + continue + } + sig.Email, ok, err = unstructured.NestedString(item.Object, "spec", "email") + if !ok || err != nil { + continue + } + + if sig.Name == sig.Email { + if sig.Name == "" { + sig.Name = item.GetName() + } else if strings.Contains(sig.Email, "@") { + sig.Email = "" // don't use the same value for name+email + } + } + + r.userInfo["user:"+item.GetName()] = sig + } + + return nil +} diff --git a/pkg/registry/apis/provisioning/jobs/export/worker.go b/pkg/registry/apis/provisioning/jobs/export/worker.go index 4478ac0ab0c..bc27197b907 100644 --- a/pkg/registry/apis/provisioning/jobs/export/worker.go +++ b/pkg/registry/apis/provisioning/jobs/export/worker.go @@ -5,23 +5,45 @@ import ( "fmt" "os" - "k8s.io/apimachinery/pkg/runtime/schema" - - dashboards "github.com/grafana/grafana/pkg/apis/dashboard" provisioning "github.com/grafana/grafana/pkg/apis/provisioning/v0alpha1" + "github.com/grafana/grafana/pkg/registry/apis/dashboard/legacy" "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" "github.com/grafana/grafana/pkg/registry/apis/provisioning/repository" gogit "github.com/grafana/grafana/pkg/registry/apis/provisioning/repository/go-git" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" + "github.com/grafana/grafana/pkg/registry/apis/provisioning/secrets" + "github.com/grafana/grafana/pkg/storage/legacysql/dualwrite" ) type ExportWorker struct { - clients *resources.ClientFactory + // Tempdir for repo clones clonedir string + + // When exporting from apiservers + clients *resources.ClientFactory + + // Check where values are currently saved + storageStatus dualwrite.Service + + // Support reading from history + legacyMigrator legacy.LegacyMigrator + + secrets secrets.Service } -func NewExportWorker(clients *resources.ClientFactory, clonedir string) *ExportWorker { - return &ExportWorker{clients, clonedir} +func NewExportWorker(clients *resources.ClientFactory, + legacyMigrator legacy.LegacyMigrator, + storageStatus dualwrite.Service, + secrets secrets.Service, + clonedir string, +) *ExportWorker { + return &ExportWorker{ + clonedir, + clients, + storageStatus, + legacyMigrator, + secrets, + } } func (r *ExportWorker) IsSupported(ctx context.Context, job provisioning.Job) bool { @@ -51,7 +73,7 @@ func (r *ExportWorker) Process(ctx context.Context, repo repository.Repository, buffered, err = gogit.Clone(ctx, repo.Config(), gogit.GoGitCloneOptions{ Root: r.clonedir, SingleCommitBeforePush: !options.History, - }, os.Stdout) + }, r.secrets, os.Stdout) if err != nil { return &provisioning.JobStatus{ State: provisioning.JobStateError, @@ -59,7 +81,7 @@ func (r *ExportWorker) Process(ctx context.Context, repo repository.Repository, }, nil } - // New empty branch + // New empty branch (same on main???) if options.Branch != "" { _, err := buffered.NewEmptyBranch(ctx, options.Branch) if err != nil { @@ -75,27 +97,32 @@ func (r *ExportWorker) Process(ctx context.Context, repo repository.Repository, dynamicClient, _, err := r.clients.New(repo.Config().Namespace) if err != nil { - return nil, fmt.Errorf("namespace mismatch") + return nil, fmt.Errorf("error getting client %w", err) } worker := newExportJob(ctx, repo, *options, dynamicClient, progress) + if options.History { + err = worker.loadUsers(ctx) + if err != nil { + return nil, fmt.Errorf("error loading users %w", err) + } + } + + // Read from legacy if not yet using unified storage + if dualwrite.IsReadingLegacyDashboardsAndFolders(ctx, r.storageStatus) { + worker.legacy = r.legacyMigrator + } + // Load and write all folders err = worker.loadFolders(ctx) if err != nil { return worker.jobStatus, err } - kinds := []schema.GroupVersionResource{{ - Group: dashboards.GROUP, - Version: "v1alpha1", - Resource: dashboards.DASHBOARD_RESOURCE, - }} - for _, kind := range kinds { - err = worker.export(ctx, kind) - if err != nil { - return worker.jobStatus, err - } + err = worker.loadResources(ctx) + if err != nil { + return worker.jobStatus, err } status := worker.jobStatus diff --git a/pkg/registry/apis/provisioning/jobs/progress.go b/pkg/registry/apis/provisioning/jobs/progress.go index 51b4c1d468c..3ba89f5335d 100644 --- a/pkg/registry/apis/provisioning/jobs/progress.go +++ b/pkg/registry/apis/provisioning/jobs/progress.go @@ -56,8 +56,8 @@ func (r *JobProgressRecorder) Record(ctx context.Context, result JobResourceResu } r.results = append(r.results, result) - logger := logging.FromContext(ctx) if result.Error != nil { + logger := logging.FromContext(ctx) logger.Error("job resource operation failed", "err", result.Error, "path", result.Path, "resource", result.Resource, "group", result.Group, "action", result.Action, "name", result.Name) r.errors = append(r.errors, result.Error.Error()) } diff --git a/pkg/registry/apis/provisioning/jobs/sync/legacy.go b/pkg/registry/apis/provisioning/jobs/sync/legacy.go new file mode 100644 index 00000000000..12ad0e59586 --- /dev/null +++ b/pkg/registry/apis/provisioning/jobs/sync/legacy.go @@ -0,0 +1,66 @@ +package sync + +import ( + "context" + "fmt" + "time" + + "google.golang.org/grpc/metadata" + "k8s.io/apimachinery/pkg/runtime/schema" + + "github.com/grafana/grafana-app-sdk/logging" + dashboard "github.com/grafana/grafana/pkg/apis/dashboard" + folders "github.com/grafana/grafana/pkg/apis/folder/v0alpha1" + "github.com/grafana/grafana/pkg/storage/unified/resource" +) + +func (r *SyncWorker) wipeUnifiedAndSetMigratedFlag(ctx context.Context, ns string) error { + kinds := []schema.GroupResource{{ + Group: folders.GROUP, + Resource: folders.RESOURCE, + }, { + Group: dashboard.GROUP, + Resource: dashboard.DASHBOARD_RESOURCE, + }} + + for _, gr := range kinds { + status, _ := r.storageStatus.Status(ctx, gr) + if status.ReadUnified { + return fmt.Errorf("unexpected state - already using unified storage for: %s", gr) + } + if status.Migrating > 0 { + if time.Since(time.UnixMilli(status.Migrating)) < time.Second*30 { + return fmt.Errorf("another migration job is running for: %s", gr) + } + } + settings := resource.BatchSettings{ + RebuildCollection: true, // wipes everything in the collection + Collection: []*resource.ResourceKey{{ + Namespace: ns, + Group: gr.Group, + Resource: gr.Resource, + }}, + } + ctx = metadata.NewOutgoingContext(ctx, settings.ToMD()) + stream, err := r.batch.BatchProcess(ctx) + if err != nil { + return fmt.Errorf("error clearing unified %s / %w", gr, err) + } + stats, err := stream.CloseAndRecv() + if err != nil { + return fmt.Errorf("error clearing unified %s / %w", gr, err) + } + logger := logging.FromContext(ctx) + logger.Error("cleared unified stoage", "stats", stats) + + status.Migrated = time.Now().UnixMilli() // but not really... since the sync is starting + status.ReadUnified = true + status.WriteLegacy = false // keep legacy "clean" + _, err = r.storageStatus.Update(ctx, status) + if err != nil { + return err + } + } + + return nil +} diff --git a/pkg/registry/apis/provisioning/jobs/sync/worker.go b/pkg/registry/apis/provisioning/jobs/sync/worker.go index 45839839106..63bf4725548 100644 --- a/pkg/registry/apis/provisioning/jobs/sync/worker.go +++ b/pkg/registry/apis/provisioning/jobs/sync/worker.go @@ -28,25 +28,42 @@ import ( "github.com/grafana/grafana/pkg/registry/apis/provisioning/repository" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" "github.com/grafana/grafana/pkg/registry/apis/provisioning/safepath" + "github.com/grafana/grafana/pkg/storage/legacysql/dualwrite" + "github.com/grafana/grafana/pkg/storage/unified/resource" ) // SyncWorker synchronizes the external repo with grafana database // this function updates the status for both the job and the referenced repository type SyncWorker struct { - client client.ProvisioningV0alpha1Interface + // Used to update the repository status with sync info + client client.ProvisioningV0alpha1Interface + + // Lists the values saved in grafana database + lister resources.ResourceLister + + // Parses fields saved in remore repository parsers *resources.ParserFactory - lister resources.ResourceLister + + // Check if the system is using unified storage + storageStatus dualwrite.Service + + // Direct access to unified storage... to wipe any existing values! + batch resource.BatchStoreClient } func NewSyncWorker( client client.ProvisioningV0alpha1Interface, parsers *resources.ParserFactory, lister resources.ResourceLister, + storageStatus dualwrite.Service, + batch resource.BatchStoreClient, ) *SyncWorker { return &SyncWorker{ - client: client, - parsers: parsers, - lister: lister, + client: client, + parsers: parsers, + lister: lister, + storageStatus: storageStatus, + batch: batch, } } @@ -80,6 +97,16 @@ func (r *SyncWorker) Process(ctx context.Context, return nil, fmt.Errorf("failed to create sync job: %w", err) } + // Check if we are onboarding from legacy storage + if dualwrite.IsReadingLegacyDashboardsAndFolders(ctx, r.storageStatus) { + if job.Spec.Sync.Incremental { + return nil, fmt.Errorf("incremental sync not suppored from legacy state") + } + if err = r.wipeUnifiedAndSetMigratedFlag(ctx, job.Namespace); err != nil { + return nil, err + } + } + // Execute the job syncError := syncJob.run(ctx, *job.Spec.Sync) jobStatus := progress.Complete(ctx, syncError) diff --git a/pkg/registry/apis/provisioning/register.go b/pkg/registry/apis/provisioning/register.go index 28003b9cb3b..07b8cafbe79 100644 --- a/pkg/registry/apis/provisioning/register.go +++ b/pkg/registry/apis/provisioning/register.go @@ -30,6 +30,7 @@ import ( clientset "github.com/grafana/grafana/pkg/generated/clientset/versioned" informers "github.com/grafana/grafana/pkg/generated/informers/externalversions" listers "github.com/grafana/grafana/pkg/generated/listers/provisioning/v0alpha1" + "github.com/grafana/grafana/pkg/registry/apis/dashboard/legacy" "github.com/grafana/grafana/pkg/registry/apis/provisioning/controller" "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs/export" @@ -46,6 +47,7 @@ import ( "github.com/grafana/grafana/pkg/services/rendering" grafanasecrets "github.com/grafana/grafana/pkg/services/secrets" "github.com/grafana/grafana/pkg/setting" + "github.com/grafana/grafana/pkg/storage/legacysql/dualwrite" "github.com/grafana/grafana/pkg/storage/unified/blob" "github.com/grafana/grafana/pkg/storage/unified/resource" ) @@ -78,6 +80,9 @@ type APIBuilder struct { tester *RepositoryTester resourceLister resources.ResourceLister repositoryLister listers.RepositoryLister + legacyMigrator legacy.LegacyMigrator + storageStatus dualwrite.Service + unified resource.ResourceClient secrets secrets.Service } @@ -90,11 +95,13 @@ func NewAPIBuilder( webhookSecretKey string, features featuremgmt.FeatureToggles, render rendering.Service, - index resource.RepositoryIndexClient, + unified resource.ResourceClient, blobstore blob.PublicBlobStore, clonedir string, // where repo clones are managed configProvider apiserver.RestConfigProvider, ghFactory github.ClientFactory, + legacyMigrator legacy.LegacyMigrator, + storageStatus dualwrite.Service, secrets secrets.Service, ) *APIBuilder { clientFactory := resources.NewFactory(configProvider) @@ -110,8 +117,11 @@ func NewAPIBuilder( }, render: render, clonedir: clonedir, - resourceLister: resources.NewResourceLister(index), + resourceLister: resources.NewResourceLister(unified), blobstore: blobstore, + legacyMigrator: legacyMigrator, + storageStatus: storageStatus, + unified: unified, secrets: secrets, } } @@ -128,6 +138,8 @@ func RegisterAPIService( client resource.ResourceClient, // implements resource.RepositoryClient configProvider apiserver.RestConfigProvider, ghFactory github.ClientFactory, + legacyMigrator legacy.LegacyMigrator, + storageStatus dualwrite.Service, // FIXME: use multi-tenant service when one exists. In this state, we can't make this a multi-tenant service! secretssvc grafanasecrets.Service, ) (*APIBuilder, error) { @@ -153,7 +165,9 @@ func RegisterAPIService( builder := NewAPIBuilder(folderResolver, urlProvider, cfg.SecretKey, features, render, client, store, filepath.Join(cfg.DataPath, "clone"), // where repositories are cloned (temporarialy for now) - configProvider, ghFactory, secrets.NewSingleTenant(secretssvc)) + configProvider, ghFactory, + legacyMigrator, storageStatus, + secrets.NewSingleTenant(secretssvc)) apiregistration.RegisterAPI(builder) return builder, nil } @@ -462,11 +476,19 @@ func (b *APIBuilder) GetPostStartHooks() (map[string]genericapiserver.PostStartH } b.repositoryLister = repoInformer.Lister() - b.jobs.Register(export.NewExportWorker(b.client, b.clonedir)) + b.jobs.Register(export.NewExportWorker( + b.client, + b.legacyMigrator, + b.storageStatus, + b.secrets, + b.clonedir, + )) b.jobs.Register(sync.NewSyncWorker( c.ProvisioningV0alpha1(), b.parsers, b.resourceLister, + b.storageStatus, + b.unified, )) renderer := pullrequest.NewRenderer(b.render, b.blobstore) diff --git a/pkg/registry/apis/provisioning/repository/go-git/wrapper.go b/pkg/registry/apis/provisioning/repository/go-git/wrapper.go index d569fe39ffc..e15df09038e 100644 --- a/pkg/registry/apis/provisioning/repository/go-git/wrapper.go +++ b/pkg/registry/apis/provisioning/repository/go-git/wrapper.go @@ -22,11 +22,10 @@ import ( provisioning "github.com/grafana/grafana/pkg/apis/provisioning/v0alpha1" "github.com/grafana/grafana/pkg/registry/apis/provisioning/repository" + "github.com/grafana/grafana/pkg/registry/apis/provisioning/secrets" ) -var ( - _ repository.Repository = (*GoGitRepo)(nil) -) +var _ repository.Repository = (*GoGitRepo)(nil) type GoGitCloneOptions struct { Root string // tempdir (when empty, memory??) @@ -36,8 +35,9 @@ type GoGitCloneOptions struct { } type GoGitRepo struct { - config *provisioning.Repository - opts GoGitCloneOptions + config *provisioning.Repository + opts GoGitCloneOptions + decryptedPassword string repo *git.Repository tree *git.Worktree @@ -50,6 +50,7 @@ func Clone( ctx context.Context, config *provisioning.Repository, opts GoGitCloneOptions, + secrets secrets.Service, progress io.Writer, // os.Stdout ) (*GoGitRepo, error) { gitcfg := config.Spec.GitHub @@ -59,7 +60,13 @@ func Clone( if opts.Root == "" { return nil, fmt.Errorf("missing root config") } - err := os.MkdirAll(opts.Root, 0700) + + decrypted, err := secrets.Decrypt(ctx, []byte(gitcfg.Token)) + if err != nil { + return nil, fmt.Errorf("error decrypting token %w", err) + } + + err = os.MkdirAll(opts.Root, 0700) if err != nil { return nil, err } @@ -68,8 +75,7 @@ func Clone( return nil, err } - url := fmt.Sprintf("/%s.git", gitcfg.URL) - + url := fmt.Sprintf("%s.git", gitcfg.URL) repo, err := git.PlainOpen(dir) if err != nil { if !errors.Is(err, git.ErrRepositoryNotExists) { @@ -78,8 +84,8 @@ func Clone( repo, err = git.PlainCloneContext(ctx, dir, false, &git.CloneOptions{ Auth: &githttp.BasicAuth{ - Username: "grafana", // this can be anything except an empty string for PAT - Password: gitcfg.Token, // TODO... will need to get from a service! + Username: "grafana", // this can be anything except an empty string for PAT + Password: string(decrypted), // TODO... will need to get from a service! }, URL: url, ReferenceName: plumbing.ReferenceName(gitcfg.Branch), @@ -117,11 +123,12 @@ func Clone( } return &GoGitRepo{ - config: config, - opts: opts, - tree: worktree, - repo: repo, - dir: dir, + config: config, + opts: opts, + tree: worktree, + decryptedPassword: string(decrypted), + repo: repo, + dir: dir, }, nil } @@ -174,7 +181,7 @@ func (g *GoGitRepo) Push(ctx context.Context, progress io.Writer) error { Progress: progress, Auth: &githttp.BasicAuth{ // reuse logic from clone? Username: "grafana", - Password: g.config.Spec.GitHub.Token, + Password: g.decryptedPassword, }, }) } @@ -204,7 +211,6 @@ func (g *GoGitRepo) ReadTree(ctx context.Context, ref string) ([]repository.File entries = append(entries, entry) return err }) - if err != nil { return nil, err } diff --git a/pkg/registry/apis/provisioning/repository/go-git/wrapper_test.go b/pkg/registry/apis/provisioning/repository/go-git/wrapper_test.go index 82ee7bbae50..e6d6fc05e66 100644 --- a/pkg/registry/apis/provisioning/repository/go-git/wrapper_test.go +++ b/pkg/registry/apis/provisioning/repository/go-git/wrapper_test.go @@ -45,6 +45,7 @@ func TestGoGitWrapper(t *testing.T) { // one commit (not 11) SingleCommitBeforePush: true, }, + nil, // TODO: add a mock os.Stdout) require.NoError(t, err) diff --git a/pkg/registry/apis/provisioning/resources/tree.go b/pkg/registry/apis/provisioning/resources/tree.go index 3525b37b02d..11e48fe4b15 100644 --- a/pkg/registry/apis/provisioning/resources/tree.go +++ b/pkg/registry/apis/provisioning/resources/tree.go @@ -6,10 +6,11 @@ import ( "sort" "strings" - apiutils "github.com/grafana/grafana/pkg/apimachinery/utils" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + + "github.com/grafana/grafana/pkg/apimachinery/utils" folders "github.com/grafana/grafana/pkg/apis/folder/v0alpha1" provisioning "github.com/grafana/grafana/pkg/apis/provisioning/v0alpha1" - "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" ) // FolderTree contains the entire set of folders (at a given snapshot in time) of the Grafana instance. @@ -97,33 +98,21 @@ func NewEmptyFolderTree() *FolderTree { } } -func NewFolderTreeFromUnstructure(ctx context.Context, rawFolders []unstructured.Unstructured) *FolderTree { - tree := make(map[string]string, len(rawFolders)) - folders := make(map[string]Folder, len(rawFolders)) - - for _, rf := range rawFolders { - name := rf.GetName() - // TODO: Can I use MetaAccessor here? - parent := rf.GetAnnotations()[apiutils.AnnoKeyFolder] - tree[name] = parent - - id := Folder{ - Title: name, - ID: name, - // TODO: should not this be be the annotation itself? - Path: "", // We'll set this later in the DirPath function :) - } - if title, ok, _ := unstructured.NestedString(rf.Object, "spec", "title"); ok { - // If the title doesn't exist (it should), we'll just use the K8s name. - id.Title = title - } - folders[name] = id +func (t *FolderTree) AddUnstructured(item *unstructured.Unstructured, skipRepo string) error { + meta, err := utils.MetaAccessor(item) + if err != nil { + return err } - - return &FolderTree{ - tree: tree, - folders: folders, + if meta.GetRepositoryName() == skipRepo { + return nil // skip it... already in tree? } + folder := Folder{ + Title: meta.FindTitle(item.GetName()), + ID: item.GetName(), + } + t.tree[folder.ID] = meta.GetFolder() + t.folders[folder.ID] = folder + return nil } func NewFolderTreeFromResourceList(resources *provisioning.ResourceList) *FolderTree { diff --git a/pkg/registry/apis/provisioning/routes.go b/pkg/registry/apis/provisioning/routes.go index 326d68558e3..b9a0ddc77c6 100644 --- a/pkg/registry/apis/provisioning/routes.go +++ b/pkg/registry/apis/provisioning/routes.go @@ -12,6 +12,7 @@ import ( authlib "github.com/grafana/authlib/types" provisioning "github.com/grafana/grafana/pkg/apis/provisioning/v0alpha1" "github.com/grafana/grafana/pkg/services/apiserver/builder" + "github.com/grafana/grafana/pkg/storage/legacysql/dualwrite" "github.com/grafana/grafana/pkg/util/errhttp" ) @@ -146,7 +147,8 @@ func (b *APIBuilder) handleSettings(w http.ResponseWriter, r *http.Request) { } settings := provisioning.RepositoryViewList{ - Items: make([]provisioning.RepositoryView, len(all)), + Items: make([]provisioning.RepositoryView, len(all)), + LegacyStorage: dualwrite.IsReadingLegacyDashboardsAndFolders(ctx, b.storageStatus), } for i, val := range all { settings.Items[i] = provisioning.RepositoryView{ diff --git a/pkg/server/wire.go b/pkg/server/wire.go index d7ea2b2b577..c3c94f816fb 100644 --- a/pkg/server/wire.go +++ b/pkg/server/wire.go @@ -36,6 +36,7 @@ import ( "github.com/grafana/grafana/pkg/middleware/csrf" "github.com/grafana/grafana/pkg/middleware/loggermw" apiregistry "github.com/grafana/grafana/pkg/registry/apis" + "github.com/grafana/grafana/pkg/registry/apis/dashboard/legacy" "github.com/grafana/grafana/pkg/registry/apis/provisioning/repository/github" appregistry "github.com/grafana/grafana/pkg/registry/apps" "github.com/grafana/grafana/pkg/services/accesscontrol" @@ -205,6 +206,7 @@ var wireBasicSet = wire.NewSet( uss.ProvideService, wire.Bind(new(usagestats.Service), new(*uss.UsageStats)), validator.ProvideService, + legacy.ProvideLegacyMigrator, pluginsintegration.WireSet, pluginDashboards.ProvideFileStoreManager, wire.Bind(new(pluginDashboards.FileStore), new(*pluginDashboards.FileStoreManager)), diff --git a/pkg/storage/legacysql/dualwrite/managed_mode3.go b/pkg/storage/legacysql/dualwrite/managed_mode3.go index af1cee40613..48dd2d65fff 100644 --- a/pkg/storage/legacysql/dualwrite/managed_mode3.go +++ b/pkg/storage/legacysql/dualwrite/managed_mode3.go @@ -37,20 +37,20 @@ func (m *service) NewStorage(gr schema.GroupResource, } return &mangedMode3{ - service: m, - legacy: legacy, - unified: storage, - target: grafanarest.NewDualWriter(grafanarest.Mode3, legacy, storage, m.reg, gr.String()), - gr: gr, + service: m, + legacy: legacy, + unified: storage, + dualwrite: grafanarest.NewDualWriter(grafanarest.Mode3, legacy, storage, m.reg, gr.String()), + gr: gr, }, nil } type mangedMode3 struct { - service Service - legacy grafanarest.LegacyStorage - unified grafanarest.Storage - target grafanarest.Storage - gr schema.GroupResource + service Service + legacy grafanarest.LegacyStorage + unified grafanarest.Storage + dualwrite grafanarest.Storage + gr schema.GroupResource } func (d *mangedMode3) Get(ctx context.Context, name string, options *metav1.GetOptions) (runtime.Object, error) { @@ -67,68 +67,78 @@ func (d *mangedMode3) List(ctx context.Context, options *metainternalversion.Lis return d.legacy.List(ctx, options) } -func (d *mangedMode3) isMigrating(ctx context.Context) error { +func (d *mangedMode3) getWriter(ctx context.Context) (grafanarest.Storage, error) { status, ok := d.service.Status(ctx, d.gr) if ok && status.Migrating > 0 { - return &apierrors.StatusError{ + return nil, &apierrors.StatusError{ ErrStatus: metav1.Status{ Code: http.StatusServiceUnavailable, Message: "the system is migrating", }, } } - return nil + if status.WriteLegacy { + if status.WriteUnified { + return d.dualwrite, nil + } + return d.legacy, nil // only write legacy (mode0) + } + return d.unified, nil // only write unified (mode4) } func (d *mangedMode3) Create(ctx context.Context, in runtime.Object, createValidation rest.ValidateObjectFunc, options *metav1.CreateOptions) (runtime.Object, error) { - if err := d.isMigrating(ctx); err != nil { + store, err := d.getWriter(ctx) + if err != nil { return nil, err } - return d.target.Create(ctx, in, createValidation, options) + return store.Create(ctx, in, createValidation, options) } func (d *mangedMode3) Update(ctx context.Context, name string, objInfo rest.UpdatedObjectInfo, createValidation rest.ValidateObjectFunc, updateValidation rest.ValidateObjectUpdateFunc, forceAllowCreate bool, options *metav1.UpdateOptions) (runtime.Object, bool, error) { - if err := d.isMigrating(ctx); err != nil { + store, err := d.getWriter(ctx) + if err != nil { return nil, false, err } - return d.target.Update(ctx, name, objInfo, createValidation, updateValidation, forceAllowCreate, options) + return store.Update(ctx, name, objInfo, createValidation, updateValidation, forceAllowCreate, options) } func (d *mangedMode3) Delete(ctx context.Context, name string, deleteValidation rest.ValidateObjectFunc, options *metav1.DeleteOptions) (runtime.Object, bool, error) { - if err := d.isMigrating(ctx); err != nil { + store, err := d.getWriter(ctx) + if err != nil { return nil, false, err } - return d.target.Delete(ctx, name, deleteValidation, options) + return store.Delete(ctx, name, deleteValidation, options) } // DeleteCollection overrides the behavior of the generic DualWriter and deletes from both LegacyStorage and Storage. func (d *mangedMode3) DeleteCollection(ctx context.Context, deleteValidation rest.ValidateObjectFunc, options *metav1.DeleteOptions, listOptions *metainternalversion.ListOptions) (runtime.Object, error) { - if err := d.isMigrating(ctx); err != nil { + store, err := d.getWriter(ctx) + if err != nil { return nil, err } - return d.target.DeleteCollection(ctx, deleteValidation, options, listOptions) + return store.DeleteCollection(ctx, deleteValidation, options, listOptions) } func (d *mangedMode3) Destroy() { - d.target.Destroy() + d.dualwrite.Destroy() } func (d *mangedMode3) GetSingularName() string { - return d.target.GetSingularName() + return d.unified.GetSingularName() } func (d *mangedMode3) NamespaceScoped() bool { - return d.target.NamespaceScoped() + return d.unified.NamespaceScoped() } func (d *mangedMode3) New() runtime.Object { - return d.target.New() + return d.unified.New() } func (d *mangedMode3) NewList() runtime.Object { - return d.target.NewList() + return d.unified.NewList() } func (d *mangedMode3) ConvertToTable(ctx context.Context, object runtime.Object, tableOptions runtime.Object) (*metav1.Table, error) { - return d.target.ConvertToTable(ctx, object, tableOptions) + return d.unified.ConvertToTable(ctx, object, tableOptions) } diff --git a/pkg/storage/legacysql/dualwrite/service.go b/pkg/storage/legacysql/dualwrite/service.go index 1f231414bf5..563a14aa6b6 100644 --- a/pkg/storage/legacysql/dualwrite/service.go +++ b/pkg/storage/legacysql/dualwrite/service.go @@ -20,11 +20,10 @@ func ProvideService(features featuremgmt.FeatureToggles, reg prometheus.Register } return &service{ - db: newFileDB(path), - reg: reg, - enabled: features.IsEnabledGlobally(featuremgmt.FlagManagedDualWriter), - // TODO: when we can "export" from legacy, this can enabled along with provisioning - // || features.IsEnabledGlobally(featuremgmt.FlagProvisioning), // required for git provisioning + db: newFileDB(path), + reg: reg, + enabled: features.IsEnabledGlobally(featuremgmt.FlagManagedDualWriter) || + features.IsEnabledGlobally(featuremgmt.FlagProvisioning), // required for git provisioning } } diff --git a/pkg/storage/legacysql/dualwrite/storage_sql_mig.go b/pkg/storage/legacysql/dualwrite/storage_sql_mig.go index 3bb4650ff78..ce7a6f3f876 100644 --- a/pkg/storage/legacysql/dualwrite/storage_sql_mig.go +++ b/pkg/storage/legacysql/dualwrite/storage_sql_mig.go @@ -2,18 +2,23 @@ package dualwrite import "github.com/grafana/grafana/pkg/services/sqlstore/migrator" +// Not yet used... but you get the idea func AddUnifiedStatusMigrations(mg *migrator.Migrator) { resourceStorageStatus := migrator.Table{ Name: "resource_storage_status", Columns: []*migrator.Column{ {Name: "group", Type: migrator.DB_NVarchar, Length: 190, Nullable: false}, {Name: "resource", Type: migrator.DB_NVarchar, Length: 190, Nullable: false}, - {Name: "migrated", Type: migrator.DB_BigInt, Nullable: false}, // Timestamp when we can start trusting unified storage - {Name: "migrating", Type: migrator.DB_BigInt, Nullable: false}, // Actively running a migration (start timestamp) + {Name: "write_legacy", Type: migrator.DB_Bool, Nullable: false, Default: "TRUE"}, + {Name: "write_unified", Type: migrator.DB_Bool, Nullable: false, Default: "TRUE"}, + {Name: "read_unified", Type: migrator.DB_Bool, Nullable: false}, + {Name: "migrating", Type: migrator.DB_BigInt, Nullable: false}, // Timestamp Actively running a migration (start timestamp) + {Name: "migrated", Type: migrator.DB_BigInt, Nullable: false}, // Timestamp job finished + {Name: "runtime", Type: migrator.DB_Bool, Nullable: false, Default: "TRUE"}, {Name: "update_key", Type: migrator.DB_BigInt, Nullable: false}, // optimistic lock key -- required for update }, Indices: []*migrator.Index{ - {Cols: []string{"group", "resource", "namespace"}, Type: migrator.UniqueIndex}, + {Cols: []string{"group", "resource"}, Type: migrator.UniqueIndex}, }, } mg.AddMigration("create resource_storage_status table", migrator.NewAddTableMigration(resourceStorageStatus)) diff --git a/pkg/storage/legacysql/dualwrite/types.go b/pkg/storage/legacysql/dualwrite/types.go index 5dc4401db6b..3000ba91947 100644 --- a/pkg/storage/legacysql/dualwrite/types.go +++ b/pkg/storage/legacysql/dualwrite/types.go @@ -22,7 +22,7 @@ type StorageStatus struct { Migrated int64 `json:"migrated" xorm:"migrated"` // Timestamp when a migration *started* this should be cleared when finished - // While migrating all write commands will be unavaliable + // While migrating all write commands will be unavailable Migrating int64 `json:"migrating" xorm:"migrating"` // When false, the behavior will not change at runtime diff --git a/pkg/storage/legacysql/dualwrite/utils.go b/pkg/storage/legacysql/dualwrite/utils.go new file mode 100644 index 00000000000..27832aecc0d --- /dev/null +++ b/pkg/storage/legacysql/dualwrite/utils.go @@ -0,0 +1,18 @@ +package dualwrite + +import ( + "golang.org/x/net/context" + "k8s.io/apimachinery/pkg/runtime/schema" + + dashboard "github.com/grafana/grafana/pkg/apis/dashboard" + folders "github.com/grafana/grafana/pkg/apis/folder/v0alpha1" +) + +func IsReadingLegacyDashboardsAndFolders(ctx context.Context, svc Service) bool { + f := svc.ReadFromUnified(ctx, folders.FolderResourceInfo.GroupResource()) + d := svc.ReadFromUnified(ctx, schema.GroupResource{ + Group: dashboard.GROUP, + Resource: dashboard.DASHBOARD_RESOURCE, + }) + return !(f && d) +} diff --git a/pkg/tests/apis/openapi_snapshots/provisioning.grafana.app-v0alpha1.json b/pkg/tests/apis/openapi_snapshots/provisioning.grafana.app-v0alpha1.json index 68f0273d1bf..e42bcfce665 100644 --- a/pkg/tests/apis/openapi_snapshots/provisioning.grafana.app-v0alpha1.json +++ b/pkg/tests/apis/openapi_snapshots/provisioning.grafana.app-v0alpha1.json @@ -3032,6 +3032,10 @@ "resource": { "type": "string" }, + "total": { + "type": "integer", + "format": "int64" + }, "update": { "type": "integer", "format": "int64" @@ -3486,6 +3490,10 @@ "kind": { "description": "Kind is a string value representing the REST resource this object represents. Servers may infer this from the endpoint the client submits requests to. Cannot be updated. In CamelCase. More info: https://git.k8s.io/community/contributors/devel/sig-architecture/api-conventions.md#types-kinds", "type": "string" + }, + "legacyStorage": { + "description": "The backend is using legacy storage FIXME: Not sure where this should be exposed... but we need it somewhere The UI should force the onboarding workflow when this is true", + "type": "boolean" } } }, diff --git a/public/app/features/provisioning/ConfigForm.tsx b/public/app/features/provisioning/ConfigForm.tsx index 0808dc5c090..1a47060bab5 100644 --- a/public/app/features/provisioning/ConfigForm.tsx +++ b/public/app/features/provisioning/ConfigForm.tsx @@ -233,7 +233,7 @@ export function ConfigForm({ data }: ConfigFormProps) { /> - +
diff --git a/public/app/features/provisioning/ExportToRepository.tsx b/public/app/features/provisioning/ExportToRepository.tsx index 81a549c1f2b..562d3594eb8 100644 --- a/public/app/features/provisioning/ExportToRepository.tsx +++ b/public/app/features/provisioning/ExportToRepository.tsx @@ -1,7 +1,6 @@ -import { Controller, useForm } from 'react-hook-form'; +import { useForm } from 'react-hook-form'; import { Box, Button, Field, FieldSet, Input, Stack, Switch, Text } from '@grafana/ui'; -import { FolderPicker } from 'app/core/components/Select/FolderPicker'; import ProgressBar from './ProgressBar'; import { Repository, useCreateRepositoryExportMutation, useListJobQuery, ExportJobOptions } from './api'; @@ -14,7 +13,7 @@ export function ExportToRepository({ repo }: Props) { const [exportRepo, exportQuery] = useCreateRepositoryExportMutation(); const exportName = exportQuery.data?.metadata?.name; - const { register, control, formState, handleSubmit } = useForm({ + const { register, formState, handleSubmit } = useForm({ defaultValues: { history: true, prefix: '', @@ -37,14 +36,6 @@ export function ExportToRepository({ repo }: Props) {
- - } - /> - - {isGit && ( diff --git a/public/app/features/provisioning/RecentJobs.tsx b/public/app/features/provisioning/RecentJobs.tsx index ceed8ebe607..15266445f75 100644 --- a/public/app/features/provisioning/RecentJobs.tsx +++ b/public/app/features/provisioning/RecentJobs.tsx @@ -193,7 +193,7 @@ export function RecentJobs({ repo }: Props) { key={items?.length} data={items.slice(0, 20)} columns={jobColumns} - getRowId={(item) => item.metadata?.resourceVersion || ''} + getRowId={(item) => `${item.metadata?.name}`} renderExpandedRow={(row) => } /> )} diff --git a/public/app/features/provisioning/RepositoryListPage.tsx b/public/app/features/provisioning/RepositoryListPage.tsx index aa00bea056c..14c9f716a1e 100644 --- a/public/app/features/provisioning/RepositoryListPage.tsx +++ b/public/app/features/provisioning/RepositoryListPage.tsx @@ -11,6 +11,7 @@ import { Stack, TextLink, Text, + Alert, } from '@grafana/ui'; import { Page } from 'app/core/components/Page/Page'; @@ -18,17 +19,23 @@ import { DeleteRepositoryButton } from './DeleteRepositoryButton'; import { SetupWarnings } from './SetupWarnings'; import { StatusBadge } from './StatusBadge'; import { SyncRepository } from './SyncRepository'; -import { Repository, ResourceCount } from './api'; +import { Repository, ResourceCount, useGetFrontendSettingsQuery } from './api'; import { NEW_URL, PROVISIONING_URL } from './constants'; import { useRepositoryList } from './hooks'; export default function RepositoryListPage() { const [items, isLoading] = useRepositoryList({ watch: true }); + const settings = useGetFrontendSettingsQuery(); return ( + {settings.data?.legacyStorage && ( + + Require running the onboarding wizard to convert from legacy to unified + + )} diff --git a/public/app/features/provisioning/SetupWarnings.tsx b/public/app/features/provisioning/SetupWarnings.tsx index 86db36e1a48..18989d2b92a 100644 --- a/public/app/features/provisioning/SetupWarnings.tsx +++ b/public/app/features/provisioning/SetupWarnings.tsx @@ -25,12 +25,6 @@ kubernetesFoldersServiceV2 = true # If you want easy kubectl setup development mode grafanaAPIServerEnsureKubectlAccess = true -[unified_storage.dashboards.dashboard.grafana.app] -dualWriterMode = 5 - -[unified_storage.folders.folder.grafana.app] -dualWriterMode = 5 - # For Github webhook support, you will need something like: [server] root_url = https://supreme-exact-beetle.ngrok-free.app`; diff --git a/public/app/features/provisioning/api/endpoints.gen.ts b/public/app/features/provisioning/api/endpoints.gen.ts index 9a16ac019ab..db71cc177eb 100644 --- a/public/app/features/provisioning/api/endpoints.gen.ts +++ b/public/app/features/provisioning/api/endpoints.gen.ts @@ -687,6 +687,7 @@ export type JobResourceSummary = { /** No action required (useful for sync) */ noop?: number; resource?: string; + total?: number; update?: number; write?: number; }; @@ -1067,6 +1068,8 @@ export type RepositoryViewList = { items: RepositoryView[]; /** Kind is a string value representing the REST resource this object represents. Servers may infer this from the endpoint the client submits requests to. Cannot be updated. In CamelCase. More info: https://git.k8s.io/community/contributors/devel/sig-architecture/api-conventions.md#types-kinds */ kind?: string; + /** The backend is using legacy storage FIXME: Not sure where this should be exposed... but we need it somewhere The UI should force the onboarding workflow when this is true */ + legacyStorage?: boolean; }; export type ResourceStats = { /** APIVersion defines the versioned schema of this representation of an object. Servers should convert recognized schemas to the latest internal value, and may reject unrecognized values. More info: https://git.k8s.io/community/contributors/devel/sig-architecture/api-conventions.md#resources */