Provisioning: introduce interface for git clones (#103175)
* Delegate clone to export in migrate from API server * Clonable interface * Root from register.go * Call option push on write * Fix linting
This commit is contained in:
@@ -11,41 +11,29 @@ import (
|
|||||||
provisioning "github.com/grafana/grafana/pkg/apis/provisioning/v0alpha1"
|
provisioning "github.com/grafana/grafana/pkg/apis/provisioning/v0alpha1"
|
||||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs"
|
"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/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/resources"
|
||||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/secrets"
|
|
||||||
"github.com/grafana/grafana/pkg/storage/legacysql/dualwrite"
|
"github.com/grafana/grafana/pkg/storage/legacysql/dualwrite"
|
||||||
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
|
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
|
||||||
"k8s.io/client-go/dynamic"
|
"k8s.io/client-go/dynamic"
|
||||||
)
|
)
|
||||||
|
|
||||||
type ExportWorker struct {
|
type ExportWorker struct {
|
||||||
// Tempdir for repo clones
|
|
||||||
clonedir string
|
|
||||||
|
|
||||||
// required to create clients
|
// required to create clients
|
||||||
clientFactory *resources.ClientFactory
|
clientFactory *resources.ClientFactory
|
||||||
|
|
||||||
// Check where values are currently saved
|
// Check where values are currently saved
|
||||||
storageStatus dualwrite.Service
|
storageStatus dualwrite.Service
|
||||||
|
|
||||||
// Decrypt secrets in config
|
|
||||||
secrets secrets.Service
|
|
||||||
|
|
||||||
parsers *resources.ParserFactory
|
parsers *resources.ParserFactory
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewExportWorker(clientFactory *resources.ClientFactory,
|
func NewExportWorker(clientFactory *resources.ClientFactory,
|
||||||
storageStatus dualwrite.Service,
|
storageStatus dualwrite.Service,
|
||||||
secrets secrets.Service,
|
|
||||||
clonedir string,
|
|
||||||
parsers *resources.ParserFactory,
|
parsers *resources.ParserFactory,
|
||||||
) *ExportWorker {
|
) *ExportWorker {
|
||||||
return &ExportWorker{
|
return &ExportWorker{
|
||||||
clonedir,
|
|
||||||
clientFactory,
|
clientFactory,
|
||||||
storageStatus,
|
storageStatus,
|
||||||
secrets,
|
|
||||||
parsers,
|
parsers,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -67,27 +55,25 @@ func (r *ExportWorker) Process(ctx context.Context, repo repository.Repository,
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
// Use the existing clone if already checked out
|
var clone repository.ClonedRepository
|
||||||
buffered, ok := repo.(*gogit.GoGitRepo)
|
if clonable, ok := repo.(repository.ClonableRepository); ok {
|
||||||
if !ok && repo.Config().Spec.GitHub != nil {
|
|
||||||
progress.SetMessage(ctx, "clone target")
|
progress.SetMessage(ctx, "clone target")
|
||||||
buffered, err = gogit.Clone(ctx, repo.Config(), gogit.GoGitCloneOptions{
|
clone, err = clonable.Clone(ctx, repository.CloneOptions{
|
||||||
Root: r.clonedir,
|
PushOnWrites: false,
|
||||||
SingleCommitBeforePush: true,
|
|
||||||
// TODO: make this configurable
|
// TODO: make this configurable
|
||||||
Timeout: 10 * time.Minute,
|
Timeout: 10 * time.Minute,
|
||||||
}, r.secrets, os.Stdout)
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("unable to clone target: %w", err)
|
return fmt.Errorf("unable to clone target: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
repo = buffered // send all writes to the buffered repo
|
|
||||||
defer func() {
|
defer func() {
|
||||||
if err := buffered.Remove(ctx); err != nil {
|
if err := clone.Remove(ctx); err != nil {
|
||||||
logging.FromContext(ctx).Error("failed to remove cloned repository after export", "err", err)
|
logging.FromContext(ctx).Error("failed to remove cloned repository after export", "err", err)
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
|
// Use the cloned repo for all operations
|
||||||
|
repo = clone
|
||||||
options.Branch = "" // :( the branch is now baked into the repo
|
options.Branch = "" // :( the branch is now baked into the repo
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -180,12 +166,13 @@ func (r *ExportWorker) Process(ctx context.Context, repo repository.Repository,
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if buffered != nil {
|
if clone != nil {
|
||||||
progress.SetMessage(ctx, "push changes")
|
progress.SetMessage(ctx, "push changes")
|
||||||
if err := buffered.Push(ctx, gogit.GoGitPushOptions{
|
if err := clone.Push(ctx, repository.PushOptions{
|
||||||
// TODO: make this configurable
|
// TODO: make this configurable
|
||||||
Timeout: 10 * time.Minute,
|
Timeout: 10 * time.Minute,
|
||||||
}, os.Stdout); err != nil {
|
Progress: os.Stdout,
|
||||||
|
}); err != nil {
|
||||||
return fmt.Errorf("error pushing changes: %w", err)
|
return fmt.Errorf("error pushing changes: %w", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -15,9 +15,7 @@ import (
|
|||||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs/export"
|
"github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs/export"
|
||||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs/sync"
|
"github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs/sync"
|
||||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/repository"
|
"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/resources"
|
||||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/secrets"
|
|
||||||
"github.com/grafana/grafana/pkg/storage/legacysql/dualwrite"
|
"github.com/grafana/grafana/pkg/storage/legacysql/dualwrite"
|
||||||
"github.com/grafana/grafana/pkg/storage/unified/resource"
|
"github.com/grafana/grafana/pkg/storage/unified/resource"
|
||||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||||
@@ -26,9 +24,6 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type MigrationWorker struct {
|
type MigrationWorker struct {
|
||||||
// Tempdir for repo clones
|
|
||||||
clonedir string
|
|
||||||
|
|
||||||
// temporary... while we still do an import
|
// temporary... while we still do an import
|
||||||
parsers *resources.ParserFactory
|
parsers *resources.ParserFactory
|
||||||
|
|
||||||
@@ -41,9 +36,6 @@ type MigrationWorker struct {
|
|||||||
// Direct access to unified storage... use carefully!
|
// Direct access to unified storage... use carefully!
|
||||||
bulk resource.BulkStoreClient
|
bulk resource.BulkStoreClient
|
||||||
|
|
||||||
// Decrypt secret from config object
|
|
||||||
secrets secrets.Service
|
|
||||||
|
|
||||||
// Delegate the export to the export worker
|
// Delegate the export to the export worker
|
||||||
exportWorker *export.ExportWorker
|
exportWorker *export.ExportWorker
|
||||||
|
|
||||||
@@ -56,18 +48,14 @@ func NewMigrationWorker(
|
|||||||
parsers *resources.ParserFactory, // should not be necessary!
|
parsers *resources.ParserFactory, // should not be necessary!
|
||||||
storageStatus dualwrite.Service,
|
storageStatus dualwrite.Service,
|
||||||
batch resource.BulkStoreClient,
|
batch resource.BulkStoreClient,
|
||||||
secrets secrets.Service,
|
|
||||||
exportWorker *export.ExportWorker,
|
exportWorker *export.ExportWorker,
|
||||||
syncWorker *sync.SyncWorker,
|
syncWorker *sync.SyncWorker,
|
||||||
clonedir string,
|
|
||||||
) *MigrationWorker {
|
) *MigrationWorker {
|
||||||
return &MigrationWorker{
|
return &MigrationWorker{
|
||||||
clonedir,
|
|
||||||
parsers,
|
parsers,
|
||||||
storageStatus,
|
storageStatus,
|
||||||
legacyMigrator,
|
legacyMigrator,
|
||||||
batch,
|
batch,
|
||||||
secrets,
|
|
||||||
exportWorker,
|
exportWorker,
|
||||||
syncWorker,
|
syncWorker,
|
||||||
}
|
}
|
||||||
@@ -84,17 +72,33 @@ func (w *MigrationWorker) Process(ctx context.Context, repo repository.Repositor
|
|||||||
return errors.New("missing migrate settings")
|
return errors.New("missing migrate settings")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
progress.SetTotal(ctx, 10) // will show a progress bar
|
||||||
|
rw, ok := repo.(repository.ReaderWriter)
|
||||||
|
if !ok {
|
||||||
|
return errors.New("migration job submitted targeting repository that is not a ReaderWriter")
|
||||||
|
}
|
||||||
|
parser, err := w.parsers.GetParser(ctx, rw)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("error getting parser: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if dualwrite.IsReadingLegacyDashboardsAndFolders(ctx, w.storageStatus) {
|
||||||
|
return w.migrateFromLegacy(ctx, rw, parser, *options, progress)
|
||||||
|
}
|
||||||
|
|
||||||
|
return w.migrateFromAPIServer(ctx, rw, parser, *options, progress)
|
||||||
|
}
|
||||||
|
|
||||||
|
// migrateFromLegacy will export the resources from legacy storage and import them into the target repository
|
||||||
|
func (w *MigrationWorker) migrateFromLegacy(ctx context.Context, rw repository.ReaderWriter, parser *resources.Parser, options provisioning.MigrateJobOptions, progress jobs.JobProgressRecorder) error {
|
||||||
var (
|
var (
|
||||||
err error
|
err error
|
||||||
buffered *gogit.GoGitRepo
|
clone repository.ClonedRepository
|
||||||
)
|
)
|
||||||
|
|
||||||
isFromLegacy := dualwrite.IsReadingLegacyDashboardsAndFolders(ctx, w.storageStatus)
|
clonable, ok := rw.(repository.ClonableRepository)
|
||||||
progress.SetTotal(ctx, 10) // will show a progress bar
|
if ok {
|
||||||
|
progress.SetMessage(ctx, "clone "+rw.Config().Spec.GitHub.URL)
|
||||||
// TODO: we should fail fast if migration is not possible and not always clone the repository.
|
|
||||||
if repo.Config().Spec.GitHub != nil {
|
|
||||||
progress.SetMessage(ctx, "clone "+repo.Config().Spec.GitHub.URL)
|
|
||||||
reader, writer := io.Pipe()
|
reader, writer := io.Pipe()
|
||||||
go func() {
|
go func() {
|
||||||
scanner := bufio.NewScanner(reader)
|
scanner := bufio.NewScanner(reader)
|
||||||
@@ -103,43 +107,24 @@ func (w *MigrationWorker) Process(ctx context.Context, repo repository.Repositor
|
|||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
buffered, err = gogit.Clone(ctx, repo.Config(), gogit.GoGitCloneOptions{
|
clone, err = clonable.Clone(ctx, repository.CloneOptions{
|
||||||
Root: w.clonedir,
|
PushOnWrites: options.History,
|
||||||
SingleCommitBeforePush: !(options.History && isFromLegacy),
|
|
||||||
// TODO: make this configurable
|
// TODO: make this configurable
|
||||||
Timeout: 10 * time.Minute,
|
Timeout: 10 * time.Minute,
|
||||||
}, w.secrets, writer)
|
Progress: writer,
|
||||||
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("unable to clone target: %w", err)
|
return fmt.Errorf("unable to clone target: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
repo = buffered // send all writes to the buffered repo
|
rw = clone // send all writes to the buffered repo
|
||||||
defer func() {
|
defer func() {
|
||||||
if err := buffered.Remove(ctx); err != nil {
|
if err := clone.Remove(ctx); err != nil {
|
||||||
logging.FromContext(ctx).Error("failed to remove cloned repository after migrate", "err", err)
|
logging.FromContext(ctx).Error("failed to remove cloned repository after migrate", "err", err)
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
||||||
rw, ok := repo.(repository.ReaderWriter)
|
|
||||||
if !ok {
|
|
||||||
return errors.New("migration job submitted targeting repository that is not a ReaderWriter")
|
|
||||||
}
|
|
||||||
|
|
||||||
if isFromLegacy {
|
|
||||||
return w.migrateFromLegacy(ctx, rw, buffered, *options, progress)
|
|
||||||
}
|
|
||||||
|
|
||||||
return w.migrateFromAPIServer(ctx, rw, *options, progress)
|
|
||||||
}
|
|
||||||
|
|
||||||
// migrateFromLegacy will export the resources from legacy storage and import them into the target repository
|
|
||||||
func (w *MigrationWorker) migrateFromLegacy(ctx context.Context, rw repository.ReaderWriter, buffered *gogit.GoGitRepo, options provisioning.MigrateJobOptions, progress jobs.JobProgressRecorder) error {
|
|
||||||
parser, err := w.parsers.GetParser(ctx, rw)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("error getting parser: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
var userInfo map[string]repository.CommitSignature
|
var userInfo map[string]repository.CommitSignature
|
||||||
if options.History {
|
if options.History {
|
||||||
progress.SetMessage(ctx, "loading users")
|
progress.SetMessage(ctx, "loading users")
|
||||||
@@ -197,7 +182,7 @@ func (w *MigrationWorker) migrateFromLegacy(ctx context.Context, rw repository.R
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if buffered != nil {
|
if clone != nil {
|
||||||
progress.SetMessage(ctx, "pushing changes")
|
progress.SetMessage(ctx, "pushing changes")
|
||||||
reader, writer := io.Pipe()
|
reader, writer := io.Pipe()
|
||||||
go func() {
|
go func() {
|
||||||
@@ -207,10 +192,11 @@ func (w *MigrationWorker) migrateFromLegacy(ctx context.Context, rw repository.R
|
|||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
if err := buffered.Push(ctx, gogit.GoGitPushOptions{
|
if err := clone.Push(ctx, repository.PushOptions{
|
||||||
// TODO: make this configurable
|
// TODO: make this configurable
|
||||||
Timeout: 10 * time.Minute,
|
Timeout: 10 * time.Minute,
|
||||||
}, writer); err != nil {
|
Progress: writer,
|
||||||
|
}); err != nil {
|
||||||
return fmt.Errorf("error pushing changes: %w", err)
|
return fmt.Errorf("error pushing changes: %w", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -244,7 +230,7 @@ func (w *MigrationWorker) migrateFromLegacy(ctx context.Context, rw repository.R
|
|||||||
}
|
}
|
||||||
|
|
||||||
// migrateFromAPIServer will export the resources from unified storage and import them into the target repository
|
// migrateFromAPIServer will export the resources from unified storage and import them into the target repository
|
||||||
func (w *MigrationWorker) migrateFromAPIServer(ctx context.Context, repo repository.ReaderWriter, options provisioning.MigrateJobOptions, progress jobs.JobProgressRecorder) error {
|
func (w *MigrationWorker) migrateFromAPIServer(ctx context.Context, repo repository.ReaderWriter, parser *resources.Parser, options provisioning.MigrateJobOptions, progress jobs.JobProgressRecorder) error {
|
||||||
progress.SetMessage(ctx, "exporting unified storage resources")
|
progress.SetMessage(ctx, "exporting unified storage resources")
|
||||||
exportJob := provisioning.Job{
|
exportJob := provisioning.Job{
|
||||||
Spec: provisioning.JobSpec{
|
Spec: provisioning.JobSpec{
|
||||||
@@ -259,6 +245,7 @@ func (w *MigrationWorker) migrateFromAPIServer(ctx context.Context, repo reposit
|
|||||||
|
|
||||||
// Reset the results after the export as pull will operate on the same resources
|
// Reset the results after the export as pull will operate on the same resources
|
||||||
progress.ResetResults()
|
progress.ResetResults()
|
||||||
|
|
||||||
progress.SetMessage(ctx, "pulling resources")
|
progress.SetMessage(ctx, "pulling resources")
|
||||||
syncJob := provisioning.Job{
|
syncJob := provisioning.Job{
|
||||||
Spec: provisioning.JobSpec{
|
Spec: provisioning.JobSpec{
|
||||||
@@ -273,11 +260,6 @@ func (w *MigrationWorker) migrateFromAPIServer(ctx context.Context, repo reposit
|
|||||||
}
|
}
|
||||||
|
|
||||||
progress.SetMessage(ctx, "removing unprovisioned resources")
|
progress.SetMessage(ctx, "removing unprovisioned resources")
|
||||||
parser, err := w.parsers.GetParser(ctx, repo)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("error getting parser: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
return parser.Clients().ForEachUnmanagedResource(ctx, func(client dynamic.ResourceInterface, item *unstructured.Unstructured) error {
|
return parser.Clients().ForEachUnmanagedResource(ctx, func(client dynamic.ResourceInterface, item *unstructured.Unstructured) error {
|
||||||
result := jobs.JobResourceResult{
|
result := jobs.JobResourceResult{
|
||||||
Name: item.GetName(),
|
Name: item.GetName(),
|
||||||
@@ -286,7 +268,7 @@ func (w *MigrationWorker) migrateFromAPIServer(ctx context.Context, repo reposit
|
|||||||
Action: repository.FileActionDeleted,
|
Action: repository.FileActionDeleted,
|
||||||
}
|
}
|
||||||
|
|
||||||
if err = client.Delete(ctx, item.GetName(), metav1.DeleteOptions{}); err != nil {
|
if err := client.Delete(ctx, item.GetName(), metav1.DeleteOptions{}); err != nil {
|
||||||
result.Error = fmt.Errorf("failed to delete folder: %w", err)
|
result.Error = fmt.Errorf("failed to delete folder: %w", err)
|
||||||
progress.Record(ctx, result)
|
progress.Record(ctx, result)
|
||||||
return result.Error
|
return result.Error
|
||||||
|
|||||||
@@ -45,6 +45,7 @@ import (
|
|||||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs/sync"
|
"github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs/sync"
|
||||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/repository"
|
"github.com/grafana/grafana/pkg/registry/apis/provisioning/repository"
|
||||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/repository/github"
|
"github.com/grafana/grafana/pkg/registry/apis/provisioning/repository/github"
|
||||||
|
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/resources"
|
||||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/safepath"
|
"github.com/grafana/grafana/pkg/registry/apis/provisioning/safepath"
|
||||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/secrets"
|
"github.com/grafana/grafana/pkg/registry/apis/provisioning/secrets"
|
||||||
@@ -557,8 +558,6 @@ func (b *APIBuilder) GetPostStartHooks() (map[string]genericapiserver.PostStartH
|
|||||||
exportWorker := export.NewExportWorker(
|
exportWorker := export.NewExportWorker(
|
||||||
b.parsers.ClientFactory,
|
b.parsers.ClientFactory,
|
||||||
b.storageStatus,
|
b.storageStatus,
|
||||||
b.secrets,
|
|
||||||
b.clonedir,
|
|
||||||
b.parsers,
|
b.parsers,
|
||||||
)
|
)
|
||||||
syncWorker := sync.NewSyncWorker(
|
syncWorker := sync.NewSyncWorker(
|
||||||
@@ -572,10 +571,8 @@ func (b *APIBuilder) GetPostStartHooks() (map[string]genericapiserver.PostStartH
|
|||||||
b.parsers,
|
b.parsers,
|
||||||
b.storageStatus,
|
b.storageStatus,
|
||||||
b.unified,
|
b.unified,
|
||||||
b.secrets,
|
|
||||||
exportWorker,
|
exportWorker,
|
||||||
syncWorker,
|
syncWorker,
|
||||||
b.clonedir,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// Pull request worker
|
// Pull request worker
|
||||||
@@ -1134,7 +1131,11 @@ func (b *APIBuilder) AsRepository(ctx context.Context, r *provisioning.Repositor
|
|||||||
r.GetName(),
|
r.GetName(),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
return repository.NewGitHub(ctx, r, b.ghFactory, b.secrets, webhookURL)
|
cloneFn := func(ctx context.Context, opts repository.CloneOptions) (repository.ClonedRepository, error) {
|
||||||
|
return gogit.Clone(ctx, b.clonedir, r, opts, b.secrets)
|
||||||
|
}
|
||||||
|
|
||||||
|
return repository.NewGitHub(ctx, r, b.ghFactory, b.secrets, webhookURL, cloneFn)
|
||||||
default:
|
default:
|
||||||
return nil, fmt.Errorf("unknown repository type (%s)", r.Spec.Type)
|
return nil, fmt.Errorf("unknown repository type (%s)", r.Spec.Type)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -36,6 +36,8 @@ type githubRepository struct {
|
|||||||
|
|
||||||
owner string
|
owner string
|
||||||
repo string
|
repo string
|
||||||
|
|
||||||
|
cloneFn CloneFn
|
||||||
}
|
}
|
||||||
|
|
||||||
var (
|
var (
|
||||||
@@ -45,6 +47,7 @@ var (
|
|||||||
_ Writer = (*githubRepository)(nil)
|
_ Writer = (*githubRepository)(nil)
|
||||||
_ Reader = (*githubRepository)(nil)
|
_ Reader = (*githubRepository)(nil)
|
||||||
_ RepositoryWithURLs = (*githubRepository)(nil)
|
_ RepositoryWithURLs = (*githubRepository)(nil)
|
||||||
|
_ ClonableRepository = (*githubRepository)(nil)
|
||||||
)
|
)
|
||||||
|
|
||||||
func NewGitHub(
|
func NewGitHub(
|
||||||
@@ -53,6 +56,7 @@ func NewGitHub(
|
|||||||
factory *pgh.Factory,
|
factory *pgh.Factory,
|
||||||
secrets secrets.Service,
|
secrets secrets.Service,
|
||||||
webhookURL string,
|
webhookURL string,
|
||||||
|
cloneFn CloneFn,
|
||||||
) (*githubRepository, error) {
|
) (*githubRepository, error) {
|
||||||
owner, repo, err := parseOwnerRepo(config.Spec.GitHub.URL)
|
owner, repo, err := parseOwnerRepo(config.Spec.GitHub.URL)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -73,6 +77,7 @@ func NewGitHub(
|
|||||||
webhookURL: webhookURL,
|
webhookURL: webhookURL,
|
||||||
owner: owner,
|
owner: owner,
|
||||||
repo: repo,
|
repo: repo,
|
||||||
|
cloneFn: cloneFn,
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -896,6 +901,10 @@ func (r *githubRepository) OnDelete(ctx context.Context) error {
|
|||||||
return r.deleteWebhook(ctx)
|
return r.deleteWebhook(ctx)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (r *githubRepository) Clone(ctx context.Context, opts CloneOptions) (ClonedRepository, error) {
|
||||||
|
return r.cloneFn(ctx, opts)
|
||||||
|
}
|
||||||
|
|
||||||
func (r *githubRepository) logger(ctx context.Context, ref string) (context.Context, logging.Logger) {
|
func (r *githubRepository) logger(ctx context.Context, ref string) (context.Context, logging.Logger) {
|
||||||
logger := logging.FromContext(ctx)
|
logger := logging.FromContext(ctx)
|
||||||
|
|
||||||
|
|||||||
@@ -45,30 +45,10 @@ func init() {
|
|||||||
|
|
||||||
var _ repository.Repository = (*GoGitRepo)(nil)
|
var _ repository.Repository = (*GoGitRepo)(nil)
|
||||||
|
|
||||||
type GoGitCloneOptions struct {
|
|
||||||
Root string // tempdir (when empty, memory??)
|
|
||||||
|
|
||||||
// If the branch does not exist, create it
|
|
||||||
CreateIfNotExists bool
|
|
||||||
|
|
||||||
// Skip intermediate commits and commit all before push
|
|
||||||
SingleCommitBeforePush bool
|
|
||||||
|
|
||||||
// Maximum allowed size for repository clone in bytes (0 means no limit)
|
|
||||||
MaxSize int64
|
|
||||||
|
|
||||||
// Maximum time allowed for clone operation in seconds (0 means no limit)
|
|
||||||
Timeout time.Duration
|
|
||||||
}
|
|
||||||
|
|
||||||
type GoGitPushOptions struct {
|
|
||||||
Timeout time.Duration
|
|
||||||
}
|
|
||||||
|
|
||||||
type GoGitRepo struct {
|
type GoGitRepo struct {
|
||||||
config *provisioning.Repository
|
config *provisioning.Repository
|
||||||
opts GoGitCloneOptions
|
|
||||||
decryptedPassword string
|
decryptedPassword string
|
||||||
|
opts repository.CloneOptions
|
||||||
|
|
||||||
repo *git.Repository
|
repo *git.Repository
|
||||||
tree *git.Worktree
|
tree *git.Worktree
|
||||||
@@ -79,12 +59,12 @@ type GoGitRepo struct {
|
|||||||
// As structured, it is valid for one context and should not be shared across multiple requests
|
// As structured, it is valid for one context and should not be shared across multiple requests
|
||||||
func Clone(
|
func Clone(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
|
root string,
|
||||||
config *provisioning.Repository,
|
config *provisioning.Repository,
|
||||||
opts GoGitCloneOptions,
|
opts repository.CloneOptions,
|
||||||
secrets secrets.Service,
|
secrets secrets.Service,
|
||||||
progress io.Writer, // os.Stdout
|
) (repository.ClonedRepository, error) {
|
||||||
) (*GoGitRepo, error) {
|
if root == "" {
|
||||||
if opts.Root == "" {
|
|
||||||
return nil, fmt.Errorf("missing root config")
|
return nil, fmt.Errorf("missing root config")
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -101,15 +81,20 @@ func Clone(
|
|||||||
return nil, fmt.Errorf("error decrypting token: %w", err)
|
return nil, fmt.Errorf("error decrypting token: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := os.MkdirAll(opts.Root, 0700); err != nil {
|
if err := os.MkdirAll(root, 0700); err != nil {
|
||||||
return nil, fmt.Errorf("create root dir: %w", err)
|
return nil, fmt.Errorf("create root dir: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
dir, err := mkdirTempClone(opts.Root, config)
|
dir, err := mkdirTempClone(root, config)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("create temp clone dir: %w", err)
|
return nil, fmt.Errorf("create temp clone dir: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
progress := opts.Progress
|
||||||
|
if progress == nil {
|
||||||
|
progress = io.Discard
|
||||||
|
}
|
||||||
|
|
||||||
repo, worktree, err := clone(ctx, config, opts, decrypted, dir, progress)
|
repo, worktree, err := clone(ctx, config, opts, decrypted, dir, progress)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if err := os.RemoveAll(dir); err != nil {
|
if err := os.RemoveAll(dir); err != nil {
|
||||||
@@ -121,15 +106,15 @@ func Clone(
|
|||||||
|
|
||||||
return &GoGitRepo{
|
return &GoGitRepo{
|
||||||
config: config,
|
config: config,
|
||||||
opts: opts,
|
|
||||||
tree: worktree,
|
tree: worktree,
|
||||||
|
opts: opts,
|
||||||
decryptedPassword: string(decrypted),
|
decryptedPassword: string(decrypted),
|
||||||
repo: repo,
|
repo: repo,
|
||||||
dir: dir,
|
dir: dir,
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func clone(ctx context.Context, config *provisioning.Repository, opts GoGitCloneOptions, decrypted []byte, dir string, progress io.Writer) (*git.Repository, *git.Worktree, error) {
|
func clone(ctx context.Context, config *provisioning.Repository, opts repository.CloneOptions, decrypted []byte, dir string, progress io.Writer) (*git.Repository, *git.Worktree, error) {
|
||||||
gitcfg := config.Spec.GitHub
|
gitcfg := config.Spec.GitHub
|
||||||
url := fmt.Sprintf("%s.git", gitcfg.URL)
|
url := fmt.Sprintf("%s.git", gitcfg.URL)
|
||||||
|
|
||||||
@@ -198,17 +183,22 @@ func mkdirTempClone(root string, config *provisioning.Repository) (string, error
|
|||||||
return os.MkdirTemp(root, fmt.Sprintf("clone-%s-%s-", config.Namespace, config.Name))
|
return os.MkdirTemp(root, fmt.Sprintf("clone-%s-%s-", config.Namespace, config.Name))
|
||||||
}
|
}
|
||||||
|
|
||||||
// Affer making changes to the worktree, push changes
|
// After making changes to the worktree, push changes
|
||||||
func (g *GoGitRepo) Push(ctx context.Context, opts GoGitPushOptions, progress io.Writer) error {
|
func (g *GoGitRepo) Push(ctx context.Context, opts repository.PushOptions) error {
|
||||||
timeout := maxOperationTimeout
|
timeout := maxOperationTimeout
|
||||||
if opts.Timeout > 0 {
|
if opts.Timeout > 0 {
|
||||||
timeout = opts.Timeout
|
timeout = opts.Timeout
|
||||||
}
|
}
|
||||||
|
|
||||||
|
progress := opts.Progress
|
||||||
|
if progress == nil {
|
||||||
|
progress = io.Discard
|
||||||
|
}
|
||||||
|
|
||||||
ctx, cancel := context.WithTimeout(ctx, timeout)
|
ctx, cancel := context.WithTimeout(ctx, timeout)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
if g.opts.SingleCommitBeforePush {
|
if !g.opts.PushOnWrites {
|
||||||
_, err := g.tree.Commit("exported from grafana", &git.CommitOptions{
|
_, err := g.tree.Commit("exported from grafana", &git.CommitOptions{
|
||||||
All: true, // Add everything that changed
|
All: true, // Add everything that changed
|
||||||
})
|
})
|
||||||
@@ -336,7 +326,7 @@ func (g *GoGitRepo) Write(ctx context.Context, fpath string, ref string, data []
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Skip commit for each file
|
// Skip commit for each file
|
||||||
if g.opts.SingleCommitBeforePush {
|
if !g.opts.PushOnWrites {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -45,7 +45,7 @@ func TestGoGitWrapper(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
wrap, err := Clone(ctx, &v0alpha1.Repository{
|
wrap, err := Clone(ctx, "testdata/clone", &v0alpha1.Repository{
|
||||||
ObjectMeta: v1.ObjectMeta{
|
ObjectMeta: v1.ObjectMeta{
|
||||||
Namespace: "ns",
|
Namespace: "ns",
|
||||||
Name: "unit-tester",
|
Name: "unit-tester",
|
||||||
@@ -57,14 +57,13 @@ func TestGoGitWrapper(t *testing.T) {
|
|||||||
},
|
},
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
GoGitCloneOptions{
|
repository.CloneOptions{
|
||||||
Root: "testdata/clone", // where things are cloned,
|
PushOnWrites: false,
|
||||||
// one commit (not 11)
|
CreateIfNotExists: true,
|
||||||
SingleCommitBeforePush: true,
|
Progress: os.Stdout,
|
||||||
CreateIfNotExists: true,
|
|
||||||
},
|
},
|
||||||
&dummySecret{},
|
&dummySecret{},
|
||||||
os.Stdout)
|
)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
tree, err := wrap.ReadTree(ctx, "")
|
tree, err := wrap.ReadTree(ctx, "")
|
||||||
@@ -89,9 +88,10 @@ func TestGoGitWrapper(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fmt.Printf("push...\n")
|
fmt.Printf("push...\n")
|
||||||
err = wrap.Push(ctx, GoGitPushOptions{
|
err = wrap.Push(ctx, repository.PushOptions{
|
||||||
Timeout: 10,
|
Timeout: 10,
|
||||||
}, os.Stdout)
|
Progress: os.Stdout,
|
||||||
|
})
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -4,8 +4,10 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"io"
|
||||||
"io/fs"
|
"io/fs"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"time"
|
||||||
|
|
||||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||||
"k8s.io/apimachinery/pkg/util/validation/field"
|
"k8s.io/apimachinery/pkg/util/validation/field"
|
||||||
@@ -44,6 +46,40 @@ type FileInfo struct {
|
|||||||
Modified *metav1.Time
|
Modified *metav1.Time
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type CloneFn func(ctx context.Context, opts CloneOptions) (ClonedRepository, error)
|
||||||
|
|
||||||
|
type CloneOptions struct {
|
||||||
|
// If the branch does not exist, create it
|
||||||
|
CreateIfNotExists bool
|
||||||
|
|
||||||
|
// Push on every write
|
||||||
|
PushOnWrites bool
|
||||||
|
|
||||||
|
// Maximum allowed size for repository clone in bytes (0 means no limit)
|
||||||
|
MaxSize int64
|
||||||
|
|
||||||
|
// Maximum time allowed for clone operation in seconds (0 means no limit)
|
||||||
|
Timeout time.Duration
|
||||||
|
|
||||||
|
// Progress is the writer to report progress to
|
||||||
|
Progress io.Writer
|
||||||
|
}
|
||||||
|
|
||||||
|
type ClonableRepository interface {
|
||||||
|
Clone(ctx context.Context, opts CloneOptions) (ClonedRepository, error)
|
||||||
|
}
|
||||||
|
|
||||||
|
type PushOptions struct {
|
||||||
|
Timeout time.Duration
|
||||||
|
Progress io.Writer
|
||||||
|
}
|
||||||
|
|
||||||
|
type ClonedRepository interface {
|
||||||
|
ReaderWriter
|
||||||
|
Push(ctx context.Context, opts PushOptions) error
|
||||||
|
Remove(ctx context.Context) error
|
||||||
|
}
|
||||||
|
|
||||||
// An entry in the file tree, as returned by 'ReadFileTree'. Like FileInfo, but contains less information.
|
// An entry in the file tree, as returned by 'ReadFileTree'. Like FileInfo, but contains less information.
|
||||||
type FileTreeEntry struct {
|
type FileTreeEntry struct {
|
||||||
// The path to the file from the base path given (if any).
|
// The path to the file from the base path given (if any).
|
||||||
|
|||||||
Reference in New Issue
Block a user