Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -214,14 +214,12 @@ At least one of `name`, `kind`, or `labels` must be set.

## EventWatcherReplacement

Only one of `yamlField`, `jsonField`, `HCLField`, or `regex` may be set alongside `file`.
Only one of `yamlField` or `regex` may be set alongside `file`.

| Field | Type | Description | Required |
| --- | --- | --- | --- |
| `file` | string | Path to the file to update. | Yes |
| `yamlField` | string | YAML path to the field to update. Must start with `$`. e.g. `$.foo.bar[0].baz`. | No |
| `jsonField` | string | JSON path to the field to update. | No |
| `HCLField` | string | HCL path to the field to update. | No |
| `regex` | string | Regular expression specifying what to replace. Only the first capturing group `()` is replaced. e.g. `host.xz/foo/bar:(v[0-9].[0-9].[0-9])`. | No |

## DriftDetection
Expand Down
20 changes: 16 additions & 4 deletions pkg/app/piped/eventwatcher/eventwatcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -350,6 +350,7 @@ func (w *watcher) execute(ctx context.Context, repo git.Repo, repoID string, eve
gitUpdateEvent = false
branchHandledEvents = make(map[string][]*pipedservice.ReportEventStatusesRequest_Event, len(eventCfgs))
gitNoChangeEvents = make([]*pipedservice.ReportEventStatusesRequest_Event, 0)
failedEvents = make([]*pipedservice.ReportEventStatusesRequest_Event, 0)
)
for _, e := range eventCfgs {
for _, cfg := range e.Configs {
Expand Down Expand Up @@ -411,7 +412,7 @@ func (w *watcher) execute(ctx context.Context, repo git.Repo, repoID string, eve
Status: model.EventStatus_EVENT_FAILURE,
StatusDescription: fmt.Sprintf("Failed to change files: %v", err),
}
branchHandledEvents[branchName] = append(branchHandledEvents[branchName], handledEvent)
failedEvents = append(failedEvents, handledEvent)
continue
}

Expand Down Expand Up @@ -446,6 +447,11 @@ func (w *watcher) execute(ctx context.Context, repo git.Repo, repoID string, eve
}
w.logger.Info(fmt.Sprintf("successfully made %d events OUTDATED", len(outDatedEvents)))
}
if len(failedEvents) > 0 {
if _, err := w.apiClient.ReportEventStatuses(ctx, &pipedservice.ReportEventStatusesRequest{Events: failedEvents}); err != nil {
return fmt.Errorf("failed to report event statuses: %w", err)
}
}

if !gitUpdateEvent {
return nil
Expand Down Expand Up @@ -671,6 +677,12 @@ func (w *watcher) updateValues(ctx context.Context, repo git.Repo, repoID string
// commitFiles commits changes if the data in Git is different from the latest event.
// If there are no changes to commit, it returns errNoChanges.
func (w *watcher) commitFiles(ctx context.Context, latestEvent *model.Event, eventName, commitMsg, gitPath string, replacements []config.EventWatcherReplacement, repo git.Repo, newBranch bool) (string, error) {
for _, r := range replacements {
if err := r.Validate(); err != nil {
return "", err
}
}

// Determine files to be changed by comparing with the latest event.
changes := make(map[string][]byte, len(replacements))
for _, r := range replacements {
Expand All @@ -689,9 +701,9 @@ func (w *watcher) commitFiles(ctx context.Context, latestEvent *model.Event, eve
case r.YAMLField != "":
newContent, upToDate, err = modifyYAML(path, r.YAMLField, latestEvent.Data)
case r.JSONField != "":
// TODO: Empower Event watcher to parse JSON format
return "", fmt.Errorf("jsonField replacements are not supported")
case r.HCLField != "":
// TODO: Empower Event watcher to parse HCL format
return "", fmt.Errorf("HCLField replacements are not supported")
case r.Regex != "":
newContent, upToDate, err = modifyText(path, r.Regex, latestEvent.Data)
}
Expand All @@ -702,11 +714,11 @@ func (w *watcher) commitFiles(ctx context.Context, latestEvent *model.Event, eve
if upToDate {
continue
}

if err := os.WriteFile(path, newContent, os.ModePerm); err != nil {
w.logger.Error("failed to write file", zap.Error(err))
return "", err
}

changes[filePath] = newContent
}
if len(changes) == 0 {
Expand Down
63 changes: 63 additions & 0 deletions pkg/app/piped/eventwatcher/eventwatcher_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,17 @@
package eventwatcher

import (
"context"
"os"
"path/filepath"
"testing"

"github.com/stretchr/testify/assert"
"go.uber.org/mock/gomock"

config "github.com/pipe-cd/pipecd/pkg/config"
"github.com/pipe-cd/pipecd/pkg/git/gittest"
"github.com/pipe-cd/pipecd/pkg/model"
)

func TestConvertStr(t *testing.T) {
Expand Down Expand Up @@ -81,6 +89,61 @@ func TestConvertStr(t *testing.T) {
}
}

func TestCommitFilesDoesNotTruncateUnsupportedReplacementFiles(t *testing.T) {
t.Parallel()

for _, tc := range []struct {
name string
replacement config.EventWatcherReplacement
wantError string
}{
{
name: "JSON field",
replacement: config.EventWatcherReplacement{
File: "version.json",
JSONField: "$.image",
},
wantError: "replacement has an unsupported jsonField",
},
{
name: "HCL field",
replacement: config.EventWatcherReplacement{
File: "version.hcl",
HCLField: "image",
},
wantError: "replacement has an unsupported HCLField",
},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()

dir := t.TempDir()
path := filepath.Join(dir, tc.replacement.File)
original := []byte("must not be truncated\n")
assert.NoError(t, os.WriteFile(path, original, 0o600))

ctrl := gomock.NewController(t)
repo := gittest.NewMockRepo(ctrl)
w := &watcher{}
_, err := w.commitFiles(
context.Background(),
&model.Event{Data: "new-value"},
"image-update",
"",
"",
[]config.EventWatcherReplacement{tc.replacement},
repo,
false,
)
assert.EqualError(t, err, tc.wantError)

actual, readErr := os.ReadFile(path)
assert.NoError(t, readErr)
assert.Equal(t, original, actual)
})
}
}

func TestModifyYAML(t *testing.T) {
t.Parallel()

Expand Down
20 changes: 16 additions & 4 deletions pkg/app/pipedv1/eventwatcher/eventwatcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -339,6 +339,7 @@ func (w *watcher) execute(ctx context.Context, repo git.Repo, repoID string, eve
outDatedDuration = time.Hour
gitUpdateEvent = false
branchHandledEvents = make(map[string][]*pipedservice.ReportEventStatusesRequest_Event, len(eventCfgs))
failedEvents = make([]*pipedservice.ReportEventStatusesRequest_Event, 0)
)
for _, e := range eventCfgs {
for _, cfg := range e.Configs {
Expand Down Expand Up @@ -398,7 +399,7 @@ func (w *watcher) execute(ctx context.Context, repo git.Repo, repoID string, eve
Status: model.EventStatus_EVENT_FAILURE,
StatusDescription: fmt.Sprintf("Failed to change files: %v", err),
}
branchHandledEvents[branchName] = append(branchHandledEvents[branchName], handledEvent)
failedEvents = append(failedEvents, handledEvent)
continue
}
handledEvent := &pipedservice.ReportEventStatusesRequest_Event{
Expand Down Expand Up @@ -426,6 +427,11 @@ func (w *watcher) execute(ctx context.Context, repo git.Repo, repoID string, eve
}
w.logger.Info(fmt.Sprintf("successfully made %d events OUTDATED", len(outDatedEvents)))
}
if len(failedEvents) > 0 {
if _, err := w.apiClient.ReportEventStatuses(ctx, &pipedservice.ReportEventStatusesRequest{Events: failedEvents}); err != nil {
return fmt.Errorf("failed to report event statuses: %w", err)
}
}

if !gitUpdateEvent {
return nil
Expand Down Expand Up @@ -607,6 +613,12 @@ func (w *watcher) updateValues(ctx context.Context, repo git.Repo, repoID string

// commitFiles commits changes if the data in Git is different from the latest event.
func (w *watcher) commitFiles(ctx context.Context, latestEvent *model.Event, eventName, commitMsg, gitPath string, replacements []config.EventWatcherReplacement, repo git.Repo, newBranch bool) (string, error) {
for _, r := range replacements {
if err := r.Validate(); err != nil {
return "", err
}
}

// Determine files to be changed by comparing with the latest event.
changes := make(map[string][]byte, len(replacements))
for _, r := range replacements {
Expand All @@ -625,9 +637,9 @@ func (w *watcher) commitFiles(ctx context.Context, latestEvent *model.Event, eve
case r.YAMLField != "":
newContent, upToDate, err = modifyYAML(path, r.YAMLField, latestEvent.Data)
case r.JSONField != "":
// TODO: Empower Event watcher to parse JSON format
return "", fmt.Errorf("jsonField replacements are not supported")
case r.HCLField != "":
// TODO: Empower Event watcher to parse HCL format
return "", fmt.Errorf("HCLField replacements are not supported")
case r.Regex != "":
newContent, upToDate, err = modifyText(path, r.Regex, latestEvent.Data)
}
Expand All @@ -637,10 +649,10 @@ func (w *watcher) commitFiles(ctx context.Context, latestEvent *model.Event, eve
if upToDate {
continue
}

if err := os.WriteFile(path, newContent, os.ModePerm); err != nil {
return "", fmt.Errorf("failed to write file: %w", err)
}

changes[filePath] = newContent
}
if len(changes) == 0 {
Expand Down
63 changes: 63 additions & 0 deletions pkg/app/pipedv1/eventwatcher/eventwatcher_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,17 @@
package eventwatcher

import (
"context"
"os"
"path/filepath"
"testing"

"github.com/stretchr/testify/assert"
"go.uber.org/mock/gomock"

config "github.com/pipe-cd/pipecd/pkg/configv1"
"github.com/pipe-cd/pipecd/pkg/git/gittest"
"github.com/pipe-cd/pipecd/pkg/model"
)

func TestConvertStr(t *testing.T) {
Expand Down Expand Up @@ -81,6 +89,61 @@ func TestConvertStr(t *testing.T) {
}
}

func TestCommitFilesDoesNotTruncateUnsupportedReplacementFiles(t *testing.T) {
t.Parallel()

for _, tc := range []struct {
name string
replacement config.EventWatcherReplacement
wantError string
}{
{
name: "JSON field",
replacement: config.EventWatcherReplacement{
File: "version.json",
JSONField: "$.image",
},
wantError: "replacement has an unsupported jsonField",
},
{
name: "HCL field",
replacement: config.EventWatcherReplacement{
File: "version.hcl",
HCLField: "image",
},
wantError: "replacement has an unsupported HCLField",
},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()

dir := t.TempDir()
path := filepath.Join(dir, tc.replacement.File)
original := []byte("must not be truncated\n")
assert.NoError(t, os.WriteFile(path, original, 0o600))

ctrl := gomock.NewController(t)
repo := gittest.NewMockRepo(ctrl)
w := &watcher{}
_, err := w.commitFiles(
context.Background(),
&model.Event{Data: "new-value"},
"image-update",
"",
"",
[]config.EventWatcherReplacement{tc.replacement},
repo,
false,
)
assert.EqualError(t, err, tc.wantError)

actual, readErr := os.ReadFile(path)
assert.NoError(t, readErr)
assert.Equal(t, original, actual)
})
}
}

func TestModifyYAML(t *testing.T) {
t.Parallel()

Expand Down
36 changes: 36 additions & 0 deletions pkg/config/event_watcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,15 @@ type EventWatcherHandlerConfig struct {
Replacements []EventWatcherReplacement `json:"replacements"`
}

func (c EventWatcherHandlerConfig) Validate() error {
for _, r := range c.Replacements {
if err := r.Validate(); err != nil {
return err
}
}
return nil
}

type EventWatcherReplacement struct {
// The path to the file to be updated.
File string `json:"file"`
Expand All @@ -88,6 +97,33 @@ type EventWatcherReplacement struct {
Regex string `json:"regex"`
}

func (r EventWatcherReplacement) Validate() error {
if r.File == "" {
return fmt.Errorf("replacement has no file name")
}
if r.JSONField != "" {
return fmt.Errorf("replacement has an unsupported jsonField")
}
if r.HCLField != "" {
return fmt.Errorf("replacement has an unsupported HCLField")
}

count := 0
if r.YAMLField != "" {
count++
}
if r.Regex != "" {
count++
}
if count == 0 {
return fmt.Errorf("replacement has no field")
}
if count > 1 {
return fmt.Errorf("replacement has multiple fields")
}
return nil
}

// EventWatcherHandlerType represents the type of an event watcher handler.
type EventWatcherHandlerType string

Expand Down
Loading