Provisioning: legacy → git → unified (#100481)

This commit is contained in:
Ryan McKinley
2025-02-13 17:13:32 +03:00
committed by GitHub
parent 45ebb0b1f9
commit 06de4c0041
33 changed files with 708 additions and 273 deletions
+7 -8
View File
@@ -5752,10 +5752,8 @@ exports[`better eslint`] = {
[0, 0, 0, "No untranslated strings in text props. Wrap text with <Trans /> or use t()", "6"],
[0, 0, 0, "No untranslated strings in text props. Wrap text with <Trans /> or use t()", "7"],
[0, 0, 0, "No untranslated strings in text props. Wrap text with <Trans /> or use t()", "8"],
[0, 0, 0, "No untranslated strings in text props. Wrap text with <Trans /> or use t()", "9"],
[0, 0, 0, "No untranslated strings in text props. Wrap text with <Trans /> or use t()", "10"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "11"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "12"]
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "9"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "10"]
],
"public/app/features/provisioning/FileHistoryPage.tsx:5381": [
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "0"],
@@ -5782,8 +5780,7 @@ exports[`better eslint`] = {
[0, 0, 0, "No untranslated strings in text props. Wrap text with <Trans /> or use t()", "1"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "2"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "3"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "4"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "5"]
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "4"]
],
"public/app/features/provisioning/RepositoryHealth.tsx:5381": [
[0, 0, 0, "No untranslated strings in text props. Wrap text with <Trans /> 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 <Trans /> or use t()", "0"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "1"],
[0, 0, 0, "No untranslated strings in text props. Wrap text with <Trans /> or use t()", "1"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "2"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "3"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "4"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "5"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "6"]
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "6"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "7"],
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "8"]
],
"public/app/features/provisioning/RepositoryOverview.tsx:5381": [
[0, 0, 0, "No untranslated strings. Wrap text with <Trans />", "0"],
+3
View File
@@ -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 {
+1
View File
@@ -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"`
@@ -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"`
}
@@ -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{
+15 -1
View File
@@ -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 {
@@ -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")
}
@@ -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
}
+39 -112
View File
@@ -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)
}
@@ -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
}
@@ -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
}
@@ -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
@@ -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())
}
@@ -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
}
@@ -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)
+26 -4
View File
@@ -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)
@@ -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
}
@@ -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)
@@ -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 {
+3 -1
View File
@@ -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{
+2
View File
@@ -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)),
@@ -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)
}
+4 -5
View File
@@ -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
}
}
@@ -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))
+1 -1
View File
@@ -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
+18
View File
@@ -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)
}
@@ -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"
}
}
},
@@ -233,7 +233,7 @@ export function ConfigForm({ data }: ConfigFormProps) {
/>
</Field>
<Field label={'Interval (seconds)'}>
<Input {...register('sync.intervalSeconds', {valueAsNumber: true})} type={'number'} placeholder={'60'} />
<Input {...register('sync.intervalSeconds', { valueAsNumber: true })} type={'number'} placeholder={'60'} />
</Field>
</FieldSet>
<FieldSet label="Advanced Settings">
@@ -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<ExportJobOptions>({
const { register, formState, handleSubmit } = useForm<ExportJobOptions>({
defaultValues: {
history: true,
prefix: '',
@@ -37,14 +36,6 @@ export function ExportToRepository({ repo }: Props) {
<Box paddingTop={2}>
<form onSubmit={handleSubmit(onSubmit)}>
<FieldSet label="Export from grafana into repository">
<Field label={'Source folder'} description="Select where we should read data (or empty for everything)">
<Controller
control={control}
name={'folder'}
render={({ field: { ref, ...field } }) => <FolderPicker {...field} />}
/>
</Field>
{isGit && (
<Field label="Target Branch" description={'The target branch. This will be created and emptied first'}>
<Input placeholder={repo.spec?.github?.branch} {...register('branch')} />
@@ -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) => <ExpandedRow row={row} />}
/>
)}
@@ -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 (
<Page navId="provisioning" subTitle="View and manage your configured repositories">
<Page.Contents isLoading={isLoading}>
<SetupWarnings />
{settings.data?.legacyStorage && (
<Alert title="Legacy Storage" severity="error">
Require running the onboarding wizard to convert from legacy to unified
</Alert>
)}
<RepositoryListPageContent items={items} />
</Page.Contents>
</Page>
@@ -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`;
@@ -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 */