diff --git a/pkg/registry/apis/provisioning/jobs/export/worker.go b/pkg/registry/apis/provisioning/jobs/export/worker.go index c8347746584..ca4df587476 100644 --- a/pkg/registry/apis/provisioning/jobs/export/worker.go +++ b/pkg/registry/apis/provisioning/jobs/export/worker.go @@ -4,12 +4,12 @@ import ( "context" "errors" "fmt" - "os" "time" 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/repository" + gogit "github.com/grafana/grafana/pkg/registry/apis/provisioning/repository/go-git" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" ) @@ -57,9 +57,14 @@ func (r *ExportWorker) Process(ctx context.Context, repo repository.Repository, return err } + writer := gogit.Progress(func(line string) { + progress.SetMessage(ctx, line) + }, "finished") + cloneOptions := repository.CloneOptions{ Timeout: 10 * time.Minute, PushOnWrites: false, + Progress: writer, BeforeFn: func() error { progress.SetMessage(ctx, "clone target") // :( the branch is now baked into the repo @@ -73,7 +78,7 @@ func (r *ExportWorker) Process(ctx context.Context, repo repository.Repository, pushOptions := repository.PushOptions{ Timeout: 10 * time.Minute, - Progress: os.Stdout, + Progress: writer, BeforeFn: func() error { progress.SetMessage(ctx, "push changes") return nil diff --git a/pkg/registry/apis/provisioning/jobs/export/worker_test.go b/pkg/registry/apis/provisioning/jobs/export/worker_test.go index 7a2bfaf0c4f..77dbadedb70 100644 --- a/pkg/registry/apis/provisioning/jobs/export/worker_test.go +++ b/pkg/registry/apis/provisioning/jobs/export/worker_test.go @@ -4,18 +4,18 @@ import ( "context" "errors" "fmt" - "os" "testing" "time" + "github.com/stretchr/testify/assert" + mock "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + v0alpha1 "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/repository" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" - "github.com/stretchr/testify/assert" - mock "github.com/stretchr/testify/mock" - "github.com/stretchr/testify/require" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) func TestExportWorker_IsSupported(t *testing.T) { @@ -265,7 +265,7 @@ func TestExportWorker_ProcessCloneAndPushOptions(t *testing.T) { mockCloneFn.On("Execute", mock.Anything, mockRepo, mock.MatchedBy(func(opts repository.CloneOptions) bool { return opts.Timeout == 10*time.Minute && !opts.PushOnWrites && opts.BeforeFn != nil }), mock.MatchedBy(func(opts repository.PushOptions) bool { - return opts.Timeout == 10*time.Minute && opts.Progress == os.Stdout && opts.BeforeFn != nil + return opts.Timeout == 10*time.Minute && opts.Progress != nil && opts.BeforeFn != nil }), mock.Anything).Return(func(ctx context.Context, repo repository.Repository, cloneOpts repository.CloneOptions, pushOpts repository.PushOptions, fn func(repository.Repository, bool) error) error { // Execute both BeforeFn functions to verify progress messages assert.NoError(t, cloneOpts.BeforeFn()) diff --git a/pkg/registry/apis/provisioning/jobs/migrate/legacy.go b/pkg/registry/apis/provisioning/jobs/migrate/legacy.go index e4841589585..ad2754bc444 100644 --- a/pkg/registry/apis/provisioning/jobs/migrate/legacy.go +++ b/pkg/registry/apis/provisioning/jobs/migrate/legacy.go @@ -1,17 +1,16 @@ package migrate import ( - "bufio" "context" "errors" "fmt" - "io" "time" "github.com/grafana/grafana-app-sdk/logging" 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/repository" + gogit "github.com/grafana/grafana/pkg/registry/apis/provisioning/repository/go-git" ) type LegacyMigrator struct { @@ -38,14 +37,9 @@ func NewLegacyMigrator( func (m *LegacyMigrator) Migrate(ctx context.Context, rw repository.ReaderWriter, options provisioning.MigrateJobOptions, progress jobs.JobProgressRecorder) error { namespace := rw.Config().Namespace - reader, writer := io.Pipe() - go func() { - scanner := bufio.NewScanner(reader) - for scanner.Scan() { - progress.SetMessage(ctx, scanner.Text()) - } - }() - + writer := gogit.Progress(func(line string) { + progress.SetMessage(ctx, line) + }, "finished") cloneOptions := repository.CloneOptions{ PushOnWrites: options.History, // TODO: make this configurable diff --git a/pkg/registry/apis/provisioning/repository/go-git/progress.go b/pkg/registry/apis/provisioning/repository/go-git/progress.go new file mode 100644 index 00000000000..29e840499ad --- /dev/null +++ b/pkg/registry/apis/provisioning/repository/go-git/progress.go @@ -0,0 +1,47 @@ +package gogit + +import ( + "bufio" + "bytes" + "io" +) + +func Progress(lines func(line string), final string) io.WriteCloser { + reader, writer := io.Pipe() + scanner := bufio.NewScanner(reader) + scanner.Split(scanLines) + go func() { + for scanner.Scan() { + line := scanner.Text() + if line != "" { + lines(line) + } + } + lines(final) + }() + return writer +} + +// Copied from bufio.ScanLines and modifed to accept standalone \r as input +func scanLines(data []byte, atEOF bool) (advance int, token []byte, err error) { + if atEOF && len(data) == 0 { + return 0, nil, nil + } + if i := bytes.IndexByte(data, '\r'); i >= 0 { + // We have a full newline-terminated line. + return i + 1, data[0:i], nil + } + + // Support standalone newlines also + if i := bytes.IndexByte(data, '\n'); i >= 0 { + // We have a full newline-terminated line. + return i + 1, data[0:i], nil + } + + // If we're at EOF, we have a final, non-terminated line. Return it. + if atEOF { + return len(data), data, nil + } + // Request more data. + return 0, nil, nil +} diff --git a/pkg/registry/apis/provisioning/repository/go-git/progress_test.go b/pkg/registry/apis/provisioning/repository/go-git/progress_test.go new file mode 100644 index 00000000000..d52168c8cd5 --- /dev/null +++ b/pkg/registry/apis/provisioning/repository/go-git/progress_test.go @@ -0,0 +1,58 @@ +package gogit + +import ( + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestProgressParsing(t *testing.T) { + tests := []struct { + name string + input string + expect []string + }{ + { + name: "no breaks", + input: "some text", + expect: []string{"some text"}, + }, + { + name: "with cr", + input: "hello\rworld", + expect: []string{"hello", "world"}, + }, + { + name: "with nl", + input: "hello\nworld", + expect: []string{"hello", "world"}, + }, + { + name: "with cr+nl", + input: "hello\r\nworld", + expect: []string{"hello", "world"}, + }, + } + for _, tt := range tests { + lastLine := "***LAST*LINE***" + t.Run(tt.name, func(t *testing.T) { + lines := []string{} + writer := Progress(func(line string) { + lines = append(lines, line) + }, lastLine) + _, _ = writer.Write([]byte(tt.input)) + err := writer.Close() + require.NoError(t, err) + + assert.EventuallyWithT(t, func(c *assert.CollectT) { + assert.NotEmpty(c, lines) + assert.Equal(c, lastLine, lines[len(lines)-1]) + + // Compare the results + require.Equal(c, tt.expect, lines[0:len(lines)-1]) + }, time.Millisecond*100, time.Microsecond*50) + }) + } +}