[Provisioning] Record resource changes by group and resource in sync (#100535)
This commit is contained in:
@@ -0,0 +1,206 @@
|
||||
package jobs
|
||||
|
||||
import (
|
||||
"context"
|
||||
"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/repository"
|
||||
)
|
||||
|
||||
// MaybeNotifyProgress will only notify if a certain amount of time has passed
|
||||
// or if the job completed
|
||||
func MaybeNotifyProgress(threshold time.Duration, fn ProgressFn) ProgressFn {
|
||||
var last time.Time
|
||||
|
||||
return func(ctx context.Context, status provisioning.JobStatus) error {
|
||||
if status.Finished != 0 || last.IsZero() || time.Since(last) > threshold {
|
||||
last = time.Now()
|
||||
return fn(ctx, status)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// FIXME: ProgressRecorder should be moved to jobs package and initialized in the queue
|
||||
|
||||
type JobResourceResult struct {
|
||||
Name string
|
||||
Resource string
|
||||
Group string
|
||||
Path string
|
||||
Action repository.FileAction
|
||||
Error error
|
||||
}
|
||||
|
||||
type JobProgressRecorder struct {
|
||||
total int
|
||||
ref string
|
||||
message string
|
||||
results []JobResourceResult
|
||||
errors []string
|
||||
progressFn ProgressFn
|
||||
}
|
||||
|
||||
func NewJobProgressRecorder(progressFn ProgressFn) *JobProgressRecorder {
|
||||
return &JobProgressRecorder{
|
||||
progressFn: MaybeNotifyProgress(15*time.Second, progressFn),
|
||||
}
|
||||
}
|
||||
|
||||
func (r *JobProgressRecorder) Record(ctx context.Context, result JobResourceResult) {
|
||||
if r.results == nil {
|
||||
r.results = make([]JobResourceResult, 0)
|
||||
}
|
||||
r.results = append(r.results, result)
|
||||
|
||||
logger := logging.FromContext(ctx)
|
||||
if result.Error != nil {
|
||||
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())
|
||||
}
|
||||
|
||||
r.notify(ctx)
|
||||
}
|
||||
|
||||
func (r *JobProgressRecorder) SetMessage(msg string) {
|
||||
r.message = msg
|
||||
}
|
||||
|
||||
func (r *JobProgressRecorder) GetMessage() string {
|
||||
return r.message
|
||||
}
|
||||
|
||||
func (r *JobProgressRecorder) SetRef(ref string) {
|
||||
r.ref = ref
|
||||
}
|
||||
|
||||
func (r *JobProgressRecorder) GetRef() string {
|
||||
return r.ref
|
||||
}
|
||||
|
||||
func (r *JobProgressRecorder) SetTotal(total int) {
|
||||
r.total = total
|
||||
}
|
||||
|
||||
func (r *JobProgressRecorder) Errors() []string {
|
||||
return r.errors
|
||||
}
|
||||
|
||||
func (r *JobProgressRecorder) summary() []*provisioning.JobResourceSummary {
|
||||
if len(r.results) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Group results by resource+group
|
||||
groupedResults := make(map[string][]JobResourceResult)
|
||||
for _, result := range r.results {
|
||||
key := result.Resource + ":" + result.Group
|
||||
groupedResults[key] = append(groupedResults[key], result)
|
||||
}
|
||||
|
||||
summaries := make([]*provisioning.JobResourceSummary, 0)
|
||||
for _, results := range groupedResults {
|
||||
if len(results) == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
// Count actions
|
||||
actions := make(map[repository.FileAction]int64)
|
||||
var errors []string
|
||||
for _, result := range results {
|
||||
if result.Error != nil {
|
||||
errors = append(errors, result.Error.Error())
|
||||
} else {
|
||||
actions[result.Action]++
|
||||
}
|
||||
}
|
||||
|
||||
// Create summary for this group
|
||||
|
||||
// Default to unknown if resource or group is empty
|
||||
resource := results[0].Resource
|
||||
if resource == "" {
|
||||
resource = "unknown"
|
||||
}
|
||||
|
||||
group := results[0].Group
|
||||
if group == "" {
|
||||
group = "unknown"
|
||||
}
|
||||
|
||||
summary := &provisioning.JobResourceSummary{
|
||||
Resource: resource,
|
||||
Group: group,
|
||||
Delete: actions[repository.FileActionDeleted],
|
||||
Update: actions[repository.FileActionUpdated],
|
||||
Create: actions[repository.FileActionCreated],
|
||||
Write: actions[repository.FileActionCreated] + actions[repository.FileActionUpdated],
|
||||
Error: int64(len(errors)),
|
||||
Noop: actions[repository.FileActionIgnored],
|
||||
Errors: errors,
|
||||
}
|
||||
|
||||
summaries = append(summaries, summary)
|
||||
}
|
||||
|
||||
return summaries
|
||||
}
|
||||
|
||||
func (r *JobProgressRecorder) progress() float64 {
|
||||
if r.total == 0 {
|
||||
return 0
|
||||
}
|
||||
|
||||
return float64(r.total - len(r.results)/r.total*100)
|
||||
}
|
||||
|
||||
func (r *JobProgressRecorder) notify(ctx context.Context) {
|
||||
jobStatus := provisioning.JobStatus{
|
||||
State: provisioning.JobStateWorking,
|
||||
Message: r.message,
|
||||
Errors: r.Errors(),
|
||||
Progress: r.progress(),
|
||||
Summary: r.summary(),
|
||||
}
|
||||
|
||||
logger := logging.FromContext(ctx)
|
||||
if err := r.progressFn(ctx, jobStatus); err != nil {
|
||||
logger.Warn("error notifying progress", "err", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (r *JobProgressRecorder) Complete(ctx context.Context, err error) *provisioning.JobStatus {
|
||||
// Initialize base job status
|
||||
jobStatus := provisioning.JobStatus{
|
||||
// TODO: do we really need to set this one here?
|
||||
// Started: job.Status.Started,
|
||||
// TODO: do we really need to set this one here?
|
||||
Finished: time.Now().UnixMilli(),
|
||||
State: provisioning.JobStateSuccess,
|
||||
Message: "completed successfully",
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
jobStatus.State = provisioning.JobStateError
|
||||
jobStatus.Message = err.Error()
|
||||
}
|
||||
|
||||
jobStatus.Summary = r.summary()
|
||||
jobStatus.Errors = r.Errors()
|
||||
|
||||
// Check for errors during execution
|
||||
if len(jobStatus.Errors) > 0 && jobStatus.State != provisioning.JobStateError {
|
||||
jobStatus.State = provisioning.JobStateError
|
||||
jobStatus.Message = "completed with errors"
|
||||
}
|
||||
|
||||
// Override message if progress have a more explicit message
|
||||
if r.message != "" && jobStatus.State != provisioning.JobStateError {
|
||||
jobStatus.Message = r.message
|
||||
}
|
||||
|
||||
return &jobStatus
|
||||
}
|
||||
@@ -4,9 +4,9 @@ import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"path"
|
||||
"reflect"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
@@ -57,42 +57,37 @@ func (r *SyncWorker) IsSupported(ctx context.Context, job provisioning.Job) bool
|
||||
func (r *SyncWorker) Process(ctx context.Context,
|
||||
repo repository.Repository,
|
||||
job provisioning.Job,
|
||||
progress jobs.ProgressFn,
|
||||
progressFn jobs.ProgressFn,
|
||||
) (*provisioning.JobStatus, error) {
|
||||
cfg := repo.Config()
|
||||
|
||||
logger := logging.FromContext(ctx).With("job", job.GetName(), "namespace", job.GetNamespace())
|
||||
patch, err := json.Marshal(map[string]any{
|
||||
|
||||
// Update sync status at the start
|
||||
data := map[string]any{
|
||||
"status": map[string]any{
|
||||
"sync": job.Status.ToSyncStatus(job.Name),
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
cfg, err = r.client.Repositories(cfg.Namespace).
|
||||
Patch(ctx, cfg.Name, types.MergePatchType, patch, metav1.PatchOptions{}, "status")
|
||||
if err != nil {
|
||||
logger.Warn("unable to update repo with job status", "err", err)
|
||||
if err := r.patchStatus(ctx, cfg, data); err != nil {
|
||||
return nil, fmt.Errorf("update repo with job status at start: %w", err)
|
||||
}
|
||||
|
||||
// Execute the sync task
|
||||
jobStatus, syncStatus, syncError := r.sync(ctx, repo, *job.Spec.Sync, progress)
|
||||
if syncStatus == nil {
|
||||
syncStatus = &provisioning.SyncStatus{}
|
||||
// Create job
|
||||
progress := jobs.NewJobProgressRecorder(progressFn)
|
||||
syncJob, err := r.createJob(ctx, repo, progress)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to create sync job: %w", err)
|
||||
}
|
||||
syncStatus.JobID = job.Name
|
||||
syncStatus.Started = job.Status.Started
|
||||
syncStatus.Finished = time.Now().UnixMilli()
|
||||
if syncError != nil {
|
||||
syncStatus.State = provisioning.JobStateError
|
||||
syncStatus.Message = []string{
|
||||
"error running sync",
|
||||
syncError.Error(),
|
||||
}
|
||||
} else if syncStatus.State == "" {
|
||||
syncStatus.State = provisioning.JobStateSuccess
|
||||
|
||||
// Execute the job
|
||||
syncError := syncJob.run(ctx, *job.Spec.Sync)
|
||||
jobStatus := progress.Complete(ctx, syncError)
|
||||
|
||||
// Create sync status and set hash if successful
|
||||
syncStatus := jobStatus.ToSyncStatus(job.Name)
|
||||
if syncStatus.State == provisioning.JobStateSuccess {
|
||||
syncStatus.Hash = progress.GetRef()
|
||||
}
|
||||
|
||||
// Update the resource stats -- give the index some time to catch up
|
||||
@@ -105,67 +100,46 @@ func (r *SyncWorker) Process(ctx context.Context,
|
||||
stats = &provisioning.ResourceStats{}
|
||||
}
|
||||
|
||||
patch, err = json.Marshal(map[string]any{
|
||||
data = map[string]any{
|
||||
"status": map[string]any{
|
||||
"sync": syncStatus,
|
||||
"stats": stats.Items,
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
_, err = r.client.Repositories(cfg.Namespace).
|
||||
Patch(ctx, cfg.Name, types.MergePatchType, patch, metav1.PatchOptions{}, "status")
|
||||
if err != nil {
|
||||
logger.Warn("unable to update repo with job status", "err", err)
|
||||
if err := r.patchStatus(ctx, cfg, data); err != nil {
|
||||
return nil, fmt.Errorf("update repo with job final status: %w", err)
|
||||
}
|
||||
|
||||
if syncError != nil {
|
||||
return nil, syncError
|
||||
}
|
||||
|
||||
return jobStatus, nil
|
||||
}
|
||||
|
||||
// start a job and run it
|
||||
func (r *SyncWorker) sync(ctx context.Context,
|
||||
repo repository.Repository,
|
||||
options provisioning.SyncJobOptions,
|
||||
progress jobs.ProgressFn,
|
||||
) (*provisioning.JobStatus, *provisioning.SyncStatus, error) {
|
||||
func (r *SyncWorker) createJob(ctx context.Context, repo repository.Repository, progress *jobs.JobProgressRecorder) (*syncJob, error) {
|
||||
cfg := repo.Config()
|
||||
if !cfg.Spec.Sync.Enabled {
|
||||
return &provisioning.JobStatus{
|
||||
State: provisioning.JobStateError,
|
||||
Message: "sync is not enabled",
|
||||
}, nil, nil
|
||||
return nil, errors.New("sync is not enabled")
|
||||
}
|
||||
|
||||
parser, err := r.parsers.GetParser(ctx, repo)
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf("failed to get parser for %s: %w", cfg.Name, err)
|
||||
return nil, fmt.Errorf("failed to get parser for %s: %w", cfg.Name, err)
|
||||
}
|
||||
|
||||
dynamicClient := parser.Client()
|
||||
if repo.Config().Namespace != dynamicClient.GetNamespace() {
|
||||
return nil, nil, fmt.Errorf("namespace mismatch")
|
||||
return nil, fmt.Errorf("namespace mismatch")
|
||||
}
|
||||
|
||||
job := &syncJob{
|
||||
summary: provisioning.JobResourceSummary{
|
||||
Resource: "all",
|
||||
Group: "all",
|
||||
},
|
||||
repository: repo,
|
||||
options: options,
|
||||
progress: progress,
|
||||
progressInterval: time.Second * 15, // how often we update the status
|
||||
progressLast: time.Now(),
|
||||
parser: parser,
|
||||
lister: r.lister,
|
||||
logger: logging.FromContext(ctx),
|
||||
jobStatus: &provisioning.JobStatus{
|
||||
State: provisioning.JobStateWorking,
|
||||
},
|
||||
syncStatus: &provisioning.SyncStatus{},
|
||||
repository: repo,
|
||||
progress: progress,
|
||||
parser: parser,
|
||||
lister: r.lister,
|
||||
folders: dynamicClient.Resource(schema.GroupVersionResource{
|
||||
Group: folders.GROUP,
|
||||
Version: folders.VERSION,
|
||||
@@ -178,49 +152,36 @@ func (r *SyncWorker) sync(ctx context.Context,
|
||||
}),
|
||||
}
|
||||
|
||||
// Execute the job
|
||||
err = job.run(ctx)
|
||||
return job, nil
|
||||
}
|
||||
|
||||
func (r *SyncWorker) patchStatus(ctx context.Context, repo *provisioning.Repository, data interface{}) error {
|
||||
patch, err := json.Marshal(data)
|
||||
if err != nil {
|
||||
job.logger.Warn("error running job", "err", err)
|
||||
job.jobStatus.State = provisioning.JobStateError
|
||||
job.jobStatus.Message = "Sync error: " + err.Error()
|
||||
job.syncStatus.State = job.jobStatus.State
|
||||
job.syncStatus.Message = append(job.syncStatus.Message, job.jobStatus.Message)
|
||||
} else if len(job.jobStatus.Errors) > 0 {
|
||||
job.jobStatus.State = provisioning.JobStateError
|
||||
} else if !job.jobStatus.State.Finished() {
|
||||
job.jobStatus.State = provisioning.JobStateSuccess
|
||||
return fmt.Errorf("unable to marshal patch data: %w", err)
|
||||
}
|
||||
|
||||
if !reflect.DeepEqual(job.summary, provisioning.JobResourceSummary{}) {
|
||||
job.jobStatus.Summary = []*provisioning.JobResourceSummary{&job.summary}
|
||||
_, err = r.client.Repositories(repo.Namespace).
|
||||
Patch(ctx, repo.Name, types.MergePatchType, patch, metav1.PatchOptions{}, "status")
|
||||
if err != nil {
|
||||
return fmt.Errorf("unable to update repo with job status: %w", err)
|
||||
}
|
||||
return job.jobStatus, job.syncStatus, nil
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// created once for each sync execution
|
||||
type syncJob struct {
|
||||
repository repository.Repository
|
||||
options provisioning.SyncJobOptions
|
||||
progress jobs.ProgressFn
|
||||
progressInterval time.Duration
|
||||
progressLast time.Time
|
||||
logger logging.Logger
|
||||
|
||||
repository repository.Repository
|
||||
progress *jobs.JobProgressRecorder
|
||||
parser *resources.Parser
|
||||
lister resources.ResourceLister
|
||||
folders dynamic.ResourceInterface
|
||||
dashboards dynamic.ResourceInterface
|
||||
folderLookup *resources.FolderTree
|
||||
|
||||
jobStatus *provisioning.JobStatus
|
||||
syncStatus *provisioning.SyncStatus
|
||||
|
||||
// FIXME: generic summary for now (not typed)
|
||||
summary provisioning.JobResourceSummary
|
||||
}
|
||||
|
||||
func (r *syncJob) run(ctx context.Context) error {
|
||||
func (r *syncJob) run(ctx context.Context, options provisioning.SyncJobOptions) error {
|
||||
// Ensure the configured folder exists and is managed by the repository
|
||||
cfg := r.repository.Config()
|
||||
rootFolder := resources.RootFolder(cfg)
|
||||
@@ -242,19 +203,19 @@ func (r *syncJob) run(ctx context.Context) error {
|
||||
if err != nil {
|
||||
return fmt.Errorf("getting latest ref: %w", err)
|
||||
}
|
||||
r.progress.SetRef(currentRef)
|
||||
|
||||
if cfg.Status.Sync.Hash != "" && r.options.Incremental {
|
||||
r.syncStatus.Hash = currentRef
|
||||
if cfg.Status.Sync.Hash != "" && options.Incremental {
|
||||
if currentRef == cfg.Status.Sync.Hash {
|
||||
message := "same commit as last sync"
|
||||
r.syncStatus.Hash = currentRef
|
||||
r.syncStatus.State = provisioning.JobStateSuccess
|
||||
r.syncStatus.Message = append(r.syncStatus.Message, message)
|
||||
r.syncStatus.Incremental = true
|
||||
r.jobStatus.Message = message
|
||||
r.progress.SetMessage("same commit as last sync")
|
||||
return nil
|
||||
}
|
||||
return r.applyVersionedChanges(ctx, versionedRepo, cfg.Status.Sync.Hash, currentRef)
|
||||
|
||||
if err := r.applyVersionedChanges(ctx, versionedRepo, cfg.Status.Sync.Hash, currentRef); err != nil {
|
||||
return fmt.Errorf("apply versioned changes: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
@@ -272,12 +233,8 @@ func (r *syncJob) run(ctx context.Context) error {
|
||||
return fmt.Errorf("error calculating changes: %w", err)
|
||||
}
|
||||
|
||||
r.syncStatus.Hash = currentRef
|
||||
if len(changes) == 0 {
|
||||
message := "no changes to sync"
|
||||
r.syncStatus.State = provisioning.JobStateSuccess
|
||||
r.syncStatus.Message = append(r.syncStatus.Message, message)
|
||||
r.jobStatus.Message = message
|
||||
r.progress.SetMessage("no changes to sync")
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -285,69 +242,67 @@ func (r *syncJob) run(ctx context.Context) error {
|
||||
r.folderLookup = resources.NewFolderTreeFromResourceList(target)
|
||||
|
||||
// Now apply the changes
|
||||
return r.applyChanges(ctx, changes)
|
||||
r.applyChanges(ctx, changes)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *syncJob) applyChanges(ctx context.Context, changes []ResourceFileChange) error {
|
||||
func (r *syncJob) applyChanges(ctx context.Context, changes []ResourceFileChange) {
|
||||
// Do the longest paths first (important for delete)
|
||||
sort.Slice(changes, func(i, j int) bool {
|
||||
return len(changes[i].Path) > len(changes[j].Path)
|
||||
})
|
||||
|
||||
r.progress.SetTotal(len(changes))
|
||||
r.progress.SetMessage("replicating changes")
|
||||
|
||||
// Create folder structure first
|
||||
for _, change := range changes {
|
||||
if len(r.jobStatus.Errors) > 20 {
|
||||
r.jobStatus.Errors = append(r.jobStatus.Errors, "too many errors, stopping")
|
||||
r.jobStatus.State = provisioning.JobStateError
|
||||
return nil
|
||||
if len(r.progress.Errors()) > 20 {
|
||||
r.progress.Record(ctx, jobs.JobResourceResult{
|
||||
Name: change.Existing.Name,
|
||||
Resource: change.Existing.Resource,
|
||||
Group: change.Existing.Group,
|
||||
Path: change.Path,
|
||||
// FIXME: should we use a skipped action instead? or a different action type?
|
||||
Action: repository.FileActionIgnored,
|
||||
Error: errors.New("too many errors"),
|
||||
})
|
||||
continue
|
||||
}
|
||||
r.maybeNotify(ctx)
|
||||
|
||||
if change.Action == repository.FileActionDeleted {
|
||||
result := jobs.JobResourceResult{
|
||||
Name: change.Existing.Name,
|
||||
Resource: change.Existing.Resource,
|
||||
Group: change.Existing.Group,
|
||||
Path: change.Path,
|
||||
Action: change.Action,
|
||||
}
|
||||
|
||||
if change.Existing == nil || change.Existing.Name == "" {
|
||||
r.logger.Error("deleted file is missing existing reference", "file", change.Path)
|
||||
r.jobStatus.Errors = append(r.jobStatus.Errors,
|
||||
fmt.Sprintf("unable to delete resource from: %s", change.Path))
|
||||
result.Error = errors.New("missing existing reference")
|
||||
r.progress.Record(ctx, result)
|
||||
continue
|
||||
}
|
||||
|
||||
client, err := r.client(change.Existing.Resource)
|
||||
if err != nil {
|
||||
r.logger.Warn("unable to get client for deleted object", "file", change.Path, "err", err, "obj", change.Existing)
|
||||
r.jobStatus.Errors = append(r.jobStatus.Errors,
|
||||
fmt.Sprintf("unsupported object: %s / %s", change.Path, change.Existing.Resource))
|
||||
result.Error = fmt.Errorf("unable to get client for deleted object: %w", err)
|
||||
r.progress.Record(ctx, result)
|
||||
continue
|
||||
}
|
||||
err = client.Delete(ctx, change.Existing.Name, metav1.DeleteOptions{})
|
||||
if err != nil {
|
||||
r.logger.Warn("deleting error", "file", change.Path, "err", err)
|
||||
r.jobStatus.Errors = append(r.jobStatus.Errors,
|
||||
fmt.Sprintf("error deleting %s: %s // %s", change.Existing.Resource, change.Existing.Name, change.Path))
|
||||
} else {
|
||||
r.summary.Delete++
|
||||
}
|
||||
|
||||
result.Error = client.Delete(ctx, change.Existing.Name, metav1.DeleteOptions{})
|
||||
r.progress.Record(ctx, result)
|
||||
continue
|
||||
}
|
||||
|
||||
// Write the resource file
|
||||
err := r.writeResourceFromFile(ctx, change.Path, "", change.Action)
|
||||
if err != nil {
|
||||
r.logger.Warn("write resource error", "file", change.Path, "err", err)
|
||||
r.jobStatus.Errors = append(r.jobStatus.Errors,
|
||||
fmt.Sprintf("error writing: %s // %s", change.Path, err.Error()))
|
||||
}
|
||||
r.progress.Record(ctx, r.writeResourceFromFile(ctx, change.Path, "", change.Action))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Convert git changes into resource file changes
|
||||
func (r *syncJob) maybeNotify(ctx context.Context) {
|
||||
if time.Since(r.progressLast) > r.progressInterval {
|
||||
err := r.progress(ctx, *r.jobStatus)
|
||||
if err != nil {
|
||||
r.logger.Warn("unable to send progress", "err", err)
|
||||
}
|
||||
}
|
||||
r.progress.SetMessage("changes replicated")
|
||||
}
|
||||
|
||||
// Convert git changes into resource file changes
|
||||
@@ -358,130 +313,141 @@ func (r *syncJob) applyVersionedChanges(ctx context.Context, repo repository.Ver
|
||||
}
|
||||
|
||||
if len(diff) < 1 {
|
||||
message := "no changes detected between commits"
|
||||
r.syncStatus.State = provisioning.JobStateSuccess
|
||||
r.syncStatus.Message = append(r.syncStatus.Message, message)
|
||||
r.syncStatus.Incremental = true
|
||||
r.jobStatus.Message = message
|
||||
r.progress.SetMessage("no changes detected between commits")
|
||||
return nil
|
||||
}
|
||||
|
||||
r.progress.SetTotal(len(diff))
|
||||
r.progress.SetMessage("replicating versioned changes")
|
||||
|
||||
for _, change := range diff {
|
||||
if len(r.jobStatus.Errors) > 20 {
|
||||
r.jobStatus.Errors = append(r.jobStatus.Errors, "too many errors to continue")
|
||||
return nil
|
||||
if len(r.progress.Errors()) > 20 {
|
||||
r.progress.Record(ctx, jobs.JobResourceResult{
|
||||
Path: change.Path,
|
||||
// FIXME: should we use a skipped action instead? or a different action type?
|
||||
Action: repository.FileActionIgnored,
|
||||
Error: errors.New("too many errors"),
|
||||
})
|
||||
continue
|
||||
}
|
||||
r.maybeNotify(ctx)
|
||||
|
||||
switch change.Action {
|
||||
case repository.FileActionCreated, repository.FileActionUpdated:
|
||||
err = r.writeResourceFromFile(ctx, change.Path, change.Ref, change.Action)
|
||||
if err != nil {
|
||||
r.logger.Warn("error writing", "change", change, "err", err)
|
||||
r.jobStatus.Errors = append(r.jobStatus.Errors,
|
||||
fmt.Sprintf("error loading: %s / %s", change.Path, err.Error()))
|
||||
}
|
||||
|
||||
r.progress.Record(ctx, r.writeResourceFromFile(ctx, change.Path, change.Ref, change.Action))
|
||||
case repository.FileActionDeleted:
|
||||
err = r.deleteObject(ctx, change.Path, change.PreviousRef)
|
||||
if err != nil {
|
||||
r.jobStatus.Errors = append(r.jobStatus.Errors,
|
||||
fmt.Sprintf("error deleting file: %s / %s", change.Path, err.Error()))
|
||||
continue
|
||||
}
|
||||
|
||||
r.progress.Record(ctx, r.deleteObject(ctx, change.Path, change.PreviousRef))
|
||||
case repository.FileActionRenamed:
|
||||
// 1. Delete
|
||||
err = r.deleteObject(ctx, change.Path, change.PreviousRef)
|
||||
if err != nil {
|
||||
r.jobStatus.Errors = append(r.jobStatus.Errors,
|
||||
fmt.Sprintf("error deleting renamed file: %s / %s", change.Path, err.Error()))
|
||||
result := r.deleteObject(ctx, change.Path, change.PreviousRef)
|
||||
if result.Error != nil {
|
||||
r.progress.Record(ctx, result)
|
||||
continue
|
||||
}
|
||||
|
||||
// 2. Create
|
||||
err = r.writeResourceFromFile(ctx, change.Path, change.Ref, repository.FileActionCreated)
|
||||
if err != nil {
|
||||
r.logger.Warn("error writing", "change", change, "err", err)
|
||||
r.jobStatus.Errors = append(r.jobStatus.Errors,
|
||||
fmt.Sprintf("error loading: %s / %s", change.Path, err.Error()))
|
||||
}
|
||||
r.progress.Record(ctx, r.writeResourceFromFile(ctx, change.Path, change.Ref, repository.FileActionCreated))
|
||||
case repository.FileActionIgnored:
|
||||
r.progress.Record(ctx, jobs.JobResourceResult{
|
||||
Path: change.Path,
|
||||
Action: repository.FileActionIgnored,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
r.progress.SetMessage("versioned changes replicated")
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *syncJob) deleteObject(ctx context.Context, path string, ref string) error {
|
||||
func (r *syncJob) deleteObject(ctx context.Context, path string, ref string) jobs.JobResourceResult {
|
||||
info, err := r.repository.Read(ctx, path, ref)
|
||||
result := jobs.JobResourceResult{
|
||||
Path: path,
|
||||
Action: repository.FileActionDeleted,
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
result.Error = fmt.Errorf("failed to read file: %w", err)
|
||||
return result
|
||||
}
|
||||
|
||||
obj, gvk, _ := resources.DecodeYAMLObject(bytes.NewBuffer(info.Data))
|
||||
if obj == nil {
|
||||
return fmt.Errorf("no object found in: %s", path)
|
||||
result.Error = errors.New("no object found")
|
||||
return result
|
||||
}
|
||||
|
||||
// Find the referenced file
|
||||
objName, _ := resources.NamesFromHashedRepoPath(r.repository.Config().Name, path)
|
||||
result.Name = objName
|
||||
result.Resource = gvk.Kind
|
||||
result.Group = gvk.Group
|
||||
|
||||
client, err := r.client(gvk.Kind)
|
||||
if err != nil {
|
||||
return err
|
||||
result.Error = fmt.Errorf("unable to get client for deleted object: %w", err)
|
||||
return result
|
||||
}
|
||||
|
||||
err = client.Delete(ctx, objName, metav1.DeleteOptions{})
|
||||
if err != nil {
|
||||
return fmt.Errorf("deleting error: %s, %w", path, err)
|
||||
result.Error = fmt.Errorf("failed to delete: %w", err)
|
||||
return result
|
||||
}
|
||||
r.summary.Delete++
|
||||
return nil
|
||||
|
||||
return result
|
||||
}
|
||||
|
||||
func (r *syncJob) writeResourceFromFile(ctx context.Context, path string, ref string, action repository.FileAction) error {
|
||||
if resources.ShouldIgnorePath(path) {
|
||||
return nil // skip
|
||||
func (r *syncJob) writeResourceFromFile(ctx context.Context, path string, ref string, action repository.FileAction) jobs.JobResourceResult {
|
||||
result := jobs.JobResourceResult{
|
||||
Path: path,
|
||||
Action: action,
|
||||
}
|
||||
|
||||
// Make sure the parent folders exist
|
||||
folder, err := r.ensureFolderPathExists(ctx, path)
|
||||
if err != nil {
|
||||
return err // fail when we can not make folders
|
||||
if resources.ShouldIgnorePath(path) {
|
||||
result.Action = repository.FileActionIgnored
|
||||
return result
|
||||
}
|
||||
|
||||
// Read the referenced file
|
||||
fileInfo, err := r.repository.Read(ctx, path, ref)
|
||||
if err != nil {
|
||||
return err
|
||||
result.Error = fmt.Errorf("failed to read file: %w", err)
|
||||
return result
|
||||
}
|
||||
|
||||
parsed, err := r.parser.Parse(ctx, fileInfo, false) // no validation
|
||||
if err != nil {
|
||||
return err
|
||||
result.Error = fmt.Errorf("failed to parse file: %w", err)
|
||||
return result
|
||||
}
|
||||
|
||||
// Make sure the parent folders exist
|
||||
folder, err := r.ensureFolderPathExists(ctx, path)
|
||||
if err != nil {
|
||||
result.Error = fmt.Errorf("failed to ensure folder path exists: %w", err)
|
||||
return result
|
||||
}
|
||||
|
||||
parsed.Meta.SetFolder(folder)
|
||||
parsed.Meta.SetUID("") // clear identifiers
|
||||
parsed.Meta.SetResourceVersion("") // clear identifiers
|
||||
|
||||
result.Name = parsed.Obj.GetName()
|
||||
result.Resource = parsed.GVR.Resource
|
||||
result.Group = parsed.GVK.Group
|
||||
|
||||
switch action {
|
||||
case repository.FileActionCreated:
|
||||
_, err = parsed.Client.Create(ctx, parsed.Obj, metav1.CreateOptions{})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
r.summary.Create++
|
||||
|
||||
result.Error = err
|
||||
case repository.FileActionUpdated:
|
||||
_, err = parsed.Client.Update(ctx, parsed.Obj, metav1.UpdateOptions{})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
r.summary.Update++
|
||||
|
||||
result.Error = err
|
||||
default:
|
||||
return fmt.Errorf("unexpected action: %s", action)
|
||||
result.Error = fmt.Errorf("unsupported action: %s", action)
|
||||
}
|
||||
return nil
|
||||
return result
|
||||
}
|
||||
|
||||
// ensureFolderPathExists creates the folder structure in the cluster.
|
||||
@@ -503,8 +469,7 @@ func (r *syncJob) ensureFolderPathExists(ctx context.Context, filePath string) (
|
||||
return f.ID, nil
|
||||
}
|
||||
|
||||
traverse := ""
|
||||
|
||||
var traverse string
|
||||
for i, part := range strings.Split(f.Path, "/") {
|
||||
if i == 0 {
|
||||
traverse = part
|
||||
@@ -527,6 +492,7 @@ func (r *syncJob) ensureFolderPathExists(ctx context.Context, filePath string) (
|
||||
r.folderLookup.Add(f, parent)
|
||||
parent = f.ID
|
||||
}
|
||||
|
||||
return f.ID, err
|
||||
}
|
||||
|
||||
@@ -574,10 +540,23 @@ func (r *syncJob) ensureFolderExists(ctx context.Context, folder resources.Folde
|
||||
Timestamp: nil, // ???&info.Modified.Time,
|
||||
})
|
||||
|
||||
result := jobs.JobResourceResult{
|
||||
Name: folder.ID,
|
||||
Resource: folders.RESOURCE,
|
||||
Group: folders.GROUP,
|
||||
Path: folder.Path,
|
||||
Action: repository.FileActionCreated,
|
||||
Error: err,
|
||||
}
|
||||
|
||||
if _, err := r.folders.Create(ctx, obj, metav1.CreateOptions{}); err != nil {
|
||||
result.Error = fmt.Errorf("failed to create folder: %w", err)
|
||||
r.progress.Record(ctx, result)
|
||||
return fmt.Errorf("failed to create folder: %w", err)
|
||||
}
|
||||
|
||||
r.progress.Record(ctx, result)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -105,6 +105,7 @@ const (
|
||||
FileActionCreated FileAction = "created"
|
||||
FileActionUpdated FileAction = "updated"
|
||||
FileActionDeleted FileAction = "deleted"
|
||||
FileActionIgnored FileAction = "ignored"
|
||||
|
||||
// Renamed actions may be reconstructed as delete then create
|
||||
FileActionRenamed FileAction = "renamed"
|
||||
|
||||
Reference in New Issue
Block a user