Add destination workflow config
This commit is contained in:
@@ -55,11 +55,12 @@ type Destination struct {
|
||||
Transform Transform `yaml:"transform"`
|
||||
PathMap PathMapping `yaml:"path_mapping"`
|
||||
Links *Links `yaml:"links"`
|
||||
State StatePolicy `yaml:"state"`
|
||||
Reconciliation ReconciliationPolicy `yaml:"reconciliation"`
|
||||
Takeover TakeoverPolicy `yaml:"takeover"`
|
||||
Workflow string `yaml:"workflow"`
|
||||
State StatePolicy `yaml:"-"`
|
||||
Reconciliation ReconciliationPolicy `yaml:"-"`
|
||||
Takeover TakeoverPolicy `yaml:"-"`
|
||||
Retention RetentionPolicy `yaml:"retention"`
|
||||
Transfer TransferPolicy `yaml:"transfer"`
|
||||
Transfer TransferPolicy `yaml:"-"`
|
||||
}
|
||||
|
||||
type Backend struct {
|
||||
|
||||
@@ -42,6 +42,11 @@ const (
|
||||
LinkPrimarySource = "source"
|
||||
)
|
||||
|
||||
const (
|
||||
WorkflowAdditive = "additive"
|
||||
WorkflowReplacement = "replacement"
|
||||
)
|
||||
|
||||
const (
|
||||
ReconciliationModeReplace = "replace"
|
||||
ReconciliationModeMerge = "merge"
|
||||
@@ -96,6 +101,9 @@ func ApplyDefaults(cfg *Config) {
|
||||
if destination.Links != nil && destination.Links.Primary == "" {
|
||||
destination.Links.Primary = LinkPrimaryAuto
|
||||
}
|
||||
if destination.Workflow == "" {
|
||||
destination.Workflow = WorkflowAdditive
|
||||
}
|
||||
if destination.State.Mode == "" {
|
||||
destination.State.Mode = StateModeSingleOwner
|
||||
}
|
||||
|
||||
@@ -30,17 +30,8 @@ pipelines:
|
||||
if got, want := cfg.Pipelines[0].Validation.OnDigestMismatch, ValidationActionFail; got != want {
|
||||
t.Fatalf("validation default = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := destination.Transfer.OnDestinationOlder, TransferActionReplace; got != want {
|
||||
t.Fatalf("transfer default = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := destination.Reconciliation.Mode, ReconciliationModeReplace; got != want {
|
||||
t.Fatalf("reconciliation mode default = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := destination.State.Mode, StateModeSingleOwner; got != want {
|
||||
t.Fatalf("state mode default = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := destination.Takeover.Mode, TakeoverModeSamePipeline; got != want {
|
||||
t.Fatalf("takeover mode default = %q, want %q", got, want)
|
||||
if got, want := destination.Workflow, WorkflowAdditive; got != want {
|
||||
t.Fatalf("workflow default = %q, want %q", got, want)
|
||||
}
|
||||
if destination.Retention.Prune.Enabled {
|
||||
t.Fatal("retention.prune.enabled default = true, want false")
|
||||
@@ -176,7 +167,7 @@ pipelines:
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFileAcceptsExplicitReconciliationModes(t *testing.T) {
|
||||
func TestLoadFileAcceptsExplicitWorkflows(t *testing.T) {
|
||||
cfg := loadConfig(t, `
|
||||
pipelines:
|
||||
- id: reports
|
||||
@@ -187,94 +178,19 @@ pipelines:
|
||||
- id: archive
|
||||
backend: local
|
||||
path: /archive
|
||||
reconciliation:
|
||||
mode: replace
|
||||
workflow: additive
|
||||
- id: web
|
||||
backend: local
|
||||
path: /web
|
||||
reconciliation:
|
||||
mode: merge
|
||||
workflow: replacement
|
||||
`)
|
||||
|
||||
destinations := cfg.Pipelines[0].Destinations
|
||||
if got, want := destinations[0].Reconciliation.Mode, ReconciliationModeReplace; got != want {
|
||||
t.Fatalf("archive reconciliation mode = %q, want %q", got, want)
|
||||
if got, want := destinations[0].Workflow, WorkflowAdditive; got != want {
|
||||
t.Fatalf("archive workflow = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := destinations[1].Reconciliation.Mode, ReconciliationModeMerge; got != want {
|
||||
t.Fatalf("web reconciliation mode = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFileAcceptsExplicitStateModes(t *testing.T) {
|
||||
cfg := loadConfig(t, `
|
||||
pipelines:
|
||||
- id: reports
|
||||
source:
|
||||
backend: local
|
||||
path: /source
|
||||
destinations:
|
||||
- id: archive
|
||||
backend: local
|
||||
path: /archive
|
||||
state:
|
||||
mode: single_owner
|
||||
- id: web
|
||||
backend: local
|
||||
path: /web
|
||||
state:
|
||||
mode: shared_root
|
||||
`)
|
||||
|
||||
destinations := cfg.Pipelines[0].Destinations
|
||||
if got, want := destinations[0].State.Mode, StateModeSingleOwner; got != want {
|
||||
t.Fatalf("archive state mode = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := destinations[1].State.Mode, StateModeSharedRoot; got != want {
|
||||
t.Fatalf("web state mode = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFileAcceptsExplicitTakeoverModes(t *testing.T) {
|
||||
cfg := loadConfig(t, `
|
||||
pipelines:
|
||||
- id: reports
|
||||
source:
|
||||
backend: local
|
||||
path: /source
|
||||
destinations:
|
||||
- id: same-pipeline
|
||||
backend: local
|
||||
path: /same-pipeline
|
||||
takeover:
|
||||
mode: same_pipeline
|
||||
- id: same-source
|
||||
backend: local
|
||||
path: /same-source
|
||||
takeover:
|
||||
mode: same_source
|
||||
- id: any-managed
|
||||
backend: local
|
||||
path: /any-managed
|
||||
takeover:
|
||||
mode: any_managed
|
||||
- id: never
|
||||
backend: local
|
||||
path: /never
|
||||
takeover:
|
||||
mode: never
|
||||
`)
|
||||
|
||||
destinations := cfg.Pipelines[0].Destinations
|
||||
wants := []string{
|
||||
TakeoverModeSamePipeline,
|
||||
TakeoverModeSameSource,
|
||||
TakeoverModeAnyManaged,
|
||||
TakeoverModeNever,
|
||||
}
|
||||
for index, want := range wants {
|
||||
if got := destinations[index].Takeover.Mode; got != want {
|
||||
t.Fatalf("destinations[%d].takeover.mode = %q, want %q", index, got, want)
|
||||
}
|
||||
if got, want := destinations[1].Workflow, WorkflowReplacement; got != want {
|
||||
t.Fatalf("web workflow = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -899,8 +815,63 @@ pipelines:
|
||||
`, "backend ftp is unsupported")
|
||||
}
|
||||
|
||||
func TestLoadFileRejectsInvalidTransferAction(t *testing.T) {
|
||||
func TestLoadFileRejectsInvalidWorkflow(t *testing.T) {
|
||||
assertLoadError(t, `
|
||||
pipelines:
|
||||
- id: reports
|
||||
source:
|
||||
backend: local
|
||||
path: /source
|
||||
destinations:
|
||||
- id: archive
|
||||
backend: local
|
||||
path: /archive
|
||||
workflow: append
|
||||
`, "workflow must be additive or replacement")
|
||||
}
|
||||
|
||||
func TestLoadFileRejectsLegacyDestinationPolicyFields(t *testing.T) {
|
||||
tests := map[string]string{
|
||||
"state": `
|
||||
pipelines:
|
||||
- id: reports
|
||||
source:
|
||||
backend: local
|
||||
path: /source
|
||||
destinations:
|
||||
- id: archive
|
||||
backend: local
|
||||
path: /archive
|
||||
state:
|
||||
mode: single_owner
|
||||
`,
|
||||
"reconciliation": `
|
||||
pipelines:
|
||||
- id: reports
|
||||
source:
|
||||
backend: local
|
||||
path: /source
|
||||
destinations:
|
||||
- id: archive
|
||||
backend: local
|
||||
path: /archive
|
||||
reconciliation:
|
||||
mode: replace
|
||||
`,
|
||||
"takeover": `
|
||||
pipelines:
|
||||
- id: reports
|
||||
source:
|
||||
backend: local
|
||||
path: /source
|
||||
destinations:
|
||||
- id: archive
|
||||
backend: local
|
||||
path: /archive
|
||||
takeover:
|
||||
mode: same_pipeline
|
||||
`,
|
||||
"transfer": `
|
||||
pipelines:
|
||||
- id: reports
|
||||
source:
|
||||
@@ -911,40 +882,14 @@ pipelines:
|
||||
backend: local
|
||||
path: /archive
|
||||
transfer:
|
||||
on_destination_older: overwrite
|
||||
`, "on_destination_older must be replace or fail")
|
||||
}
|
||||
|
||||
func TestLoadFileRejectsInvalidTakeoverMode(t *testing.T) {
|
||||
assertLoadError(t, `
|
||||
pipelines:
|
||||
- id: reports
|
||||
source:
|
||||
backend: local
|
||||
path: /source
|
||||
destinations:
|
||||
- id: archive
|
||||
backend: local
|
||||
path: /archive
|
||||
takeover:
|
||||
mode: unmanaged
|
||||
`, "takeover.mode must be same_pipeline, same_source, any_managed, or never")
|
||||
}
|
||||
|
||||
func TestLoadFileRejectsUnknownTakeoverFields(t *testing.T) {
|
||||
assertLoadError(t, `
|
||||
pipelines:
|
||||
- id: reports
|
||||
source:
|
||||
backend: local
|
||||
path: /source
|
||||
destinations:
|
||||
- id: archive
|
||||
backend: local
|
||||
path: /archive
|
||||
takeover:
|
||||
surprise: true
|
||||
`, "field surprise not found")
|
||||
on_destination_older: replace
|
||||
`,
|
||||
}
|
||||
for name, body := range tests {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
assertLoadError(t, body, "field "+name+" not found")
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFileRejectsInvalidValidationAction(t *testing.T) {
|
||||
@@ -1020,8 +965,6 @@ func TestExampleConfigsLoad(t *testing.T) {
|
||||
"../../examples/local-index.yml",
|
||||
"../../examples/fan-out.yml",
|
||||
"../../examples/archive-and-latest.yml",
|
||||
"../../examples/merge-reconciliation.yml",
|
||||
"../../examples/shared-root.yml",
|
||||
"../../examples/http-upload-local.yml",
|
||||
"../../examples/ssh-destination.yml",
|
||||
"../../examples/s3-destination.yml",
|
||||
|
||||
@@ -74,11 +74,8 @@ func Validate(cfg Config) error {
|
||||
errs = validatePublishTransformPolicy(errs, destinationContext, destination.Publish, destination.Transform)
|
||||
errs = validatePathMapping(errs, destinationContext+".path_mapping", destination.PathMap)
|
||||
errs = validateLinks(errs, destinationContext+".links", destination.Links)
|
||||
errs = validateStatePolicy(errs, destinationContext+".state", destination.State)
|
||||
errs = validateReconciliationPolicy(errs, destinationContext+".reconciliation", destination.Reconciliation)
|
||||
errs = validateTakeoverPolicy(errs, destinationContext+".takeover", destination.Takeover)
|
||||
errs = validateWorkflow(errs, destinationContext+".workflow", destination.Workflow)
|
||||
errs = validateRetentionPolicy(errs, destinationContext+".retention", destination.Retention)
|
||||
errs = validateTransferPolicy(errs, destinationContext+".transfer", destination.Transfer)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -369,18 +366,9 @@ func validateLinks(errs ValidationErrors, context string, links *Links) Validati
|
||||
return errs
|
||||
}
|
||||
|
||||
func validateReconciliationPolicy(errs ValidationErrors, context string, policy ReconciliationPolicy) ValidationErrors {
|
||||
if policy.Mode != ReconciliationModeReplace && policy.Mode != ReconciliationModeMerge {
|
||||
errs = append(errs, context+".mode must be "+ReconciliationModeReplace+" or "+ReconciliationModeMerge)
|
||||
}
|
||||
return errs
|
||||
}
|
||||
|
||||
func validateTakeoverPolicy(errs ValidationErrors, context string, policy TakeoverPolicy) ValidationErrors {
|
||||
switch policy.Mode {
|
||||
case TakeoverModeSamePipeline, TakeoverModeSameSource, TakeoverModeAnyManaged, TakeoverModeNever:
|
||||
default:
|
||||
errs = append(errs, context+".mode must be "+TakeoverModeSamePipeline+", "+TakeoverModeSameSource+", "+TakeoverModeAnyManaged+", or "+TakeoverModeNever)
|
||||
func validateWorkflow(errs ValidationErrors, context, workflow string) ValidationErrors {
|
||||
if workflow != WorkflowAdditive && workflow != WorkflowReplacement {
|
||||
errs = append(errs, context+" must be "+WorkflowAdditive+" or "+WorkflowReplacement)
|
||||
}
|
||||
return errs
|
||||
}
|
||||
@@ -401,26 +389,3 @@ func validateRetentionPolicy(errs ValidationErrors, context string, policy Reten
|
||||
}
|
||||
return errs
|
||||
}
|
||||
|
||||
func validateStatePolicy(errs ValidationErrors, context string, policy StatePolicy) ValidationErrors {
|
||||
if policy.Mode != StateModeSingleOwner && policy.Mode != StateModeSharedRoot {
|
||||
errs = append(errs, context+".mode must be "+StateModeSingleOwner+" or "+StateModeSharedRoot)
|
||||
}
|
||||
return errs
|
||||
}
|
||||
|
||||
func validateTransferPolicy(errs ValidationErrors, context string, policy TransferPolicy) ValidationErrors {
|
||||
if policy.OnDestinationSame != TransferActionSkip && policy.OnDestinationSame != TransferActionFail {
|
||||
errs = append(errs, context+".on_destination_same must be skip or fail")
|
||||
}
|
||||
if policy.OnDestinationOlder != TransferActionReplace && policy.OnDestinationOlder != TransferActionFail {
|
||||
errs = append(errs, context+".on_destination_older must be replace or fail")
|
||||
}
|
||||
if policy.OnDestinationNewer != TransferActionSkip && policy.OnDestinationNewer != TransferActionFail && policy.OnDestinationNewer != TransferActionReplace {
|
||||
errs = append(errs, context+".on_destination_newer must be skip, replace, or fail")
|
||||
}
|
||||
if policy.OnConflict != TransferActionFail && policy.OnConflict != TransferActionReplace {
|
||||
errs = append(errs, context+".on_conflict must be fail or replace")
|
||||
}
|
||||
return errs
|
||||
}
|
||||
|
||||
@@ -51,29 +51,6 @@ func TestValidateChecksPublishTransformPolicy(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateAcceptsForceReplacementTransferActions(t *testing.T) {
|
||||
cfg := Config{Pipelines: []Pipeline{{
|
||||
ID: "reports",
|
||||
Source: Backend{
|
||||
Backend: BackendLocal,
|
||||
Path: "/source",
|
||||
},
|
||||
Destinations: []Destination{{
|
||||
ID: "archive",
|
||||
Backend: BackendLocal,
|
||||
Path: "/destination",
|
||||
Transfer: TransferPolicy{
|
||||
OnDestinationNewer: TransferActionReplace,
|
||||
OnConflict: TransferActionReplace,
|
||||
},
|
||||
}},
|
||||
}}}
|
||||
ApplyDefaults(&cfg)
|
||||
if err := Validate(cfg); err != nil {
|
||||
t.Fatalf("Validate() error = %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidatePathMapping(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
@@ -108,51 +85,15 @@ func TestValidatePathMapping(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateReconciliationPolicy(t *testing.T) {
|
||||
func TestValidateWorkflow(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
mode string
|
||||
wantErr bool
|
||||
name string
|
||||
workflow string
|
||||
wantErr bool
|
||||
}{
|
||||
{name: "replace", mode: ReconciliationModeReplace},
|
||||
{name: "merge", mode: ReconciliationModeMerge},
|
||||
{name: "invalid", mode: "append", wantErr: true},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
cfg := Config{Pipelines: []Pipeline{{
|
||||
ID: "reports",
|
||||
Source: Backend{Backend: BackendLocal, Path: "/source"},
|
||||
Destinations: []Destination{{
|
||||
ID: "archive",
|
||||
Backend: BackendLocal,
|
||||
Path: "/destination",
|
||||
Reconciliation: ReconciliationPolicy{Mode: tt.mode},
|
||||
}},
|
||||
}}}
|
||||
ApplyDefaults(&cfg)
|
||||
err := Validate(cfg)
|
||||
if tt.wantErr && err == nil {
|
||||
t.Fatal("Validate() error = nil, want error")
|
||||
}
|
||||
if !tt.wantErr && err != nil {
|
||||
t.Fatalf("Validate() error = %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateTakeoverPolicy(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
mode string
|
||||
wantErr bool
|
||||
}{
|
||||
{name: "same pipeline", mode: TakeoverModeSamePipeline},
|
||||
{name: "same source", mode: TakeoverModeSameSource},
|
||||
{name: "any managed", mode: TakeoverModeAnyManaged},
|
||||
{name: "never", mode: TakeoverModeNever},
|
||||
{name: "invalid", mode: "unmanaged", wantErr: true},
|
||||
{name: "additive", workflow: WorkflowAdditive},
|
||||
{name: "replacement", workflow: WorkflowReplacement},
|
||||
{name: "invalid", workflow: "append", wantErr: true},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
@@ -163,7 +104,7 @@ func TestValidateTakeoverPolicy(t *testing.T) {
|
||||
ID: "archive",
|
||||
Backend: BackendLocal,
|
||||
Path: "/destination",
|
||||
Takeover: TakeoverPolicy{Mode: tt.mode},
|
||||
Workflow: tt.workflow,
|
||||
}},
|
||||
}}}
|
||||
ApplyDefaults(&cfg)
|
||||
@@ -178,7 +119,7 @@ func TestValidateTakeoverPolicy(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateTakeoverPolicyReportsFieldContext(t *testing.T) {
|
||||
func TestValidateWorkflowReportsFieldContext(t *testing.T) {
|
||||
cfg := Config{Pipelines: []Pipeline{{
|
||||
ID: "reports",
|
||||
Source: Backend{Backend: BackendLocal, Path: "/source"},
|
||||
@@ -186,7 +127,7 @@ func TestValidateTakeoverPolicyReportsFieldContext(t *testing.T) {
|
||||
ID: "archive",
|
||||
Backend: BackendLocal,
|
||||
Path: "/destination",
|
||||
Takeover: TakeoverPolicy{Mode: "unmanaged"},
|
||||
Workflow: "append",
|
||||
}},
|
||||
}}}
|
||||
ApplyDefaults(&cfg)
|
||||
@@ -194,46 +135,12 @@ func TestValidateTakeoverPolicyReportsFieldContext(t *testing.T) {
|
||||
if err == nil {
|
||||
t.Fatal("Validate() error = nil, want error")
|
||||
}
|
||||
want := "pipelines[0].destinations[0].takeover.mode must be same_pipeline, same_source, any_managed, or never"
|
||||
want := "pipelines[0].destinations[0].workflow must be additive or replacement"
|
||||
if !strings.Contains(err.Error(), want) {
|
||||
t.Fatalf("Validate() error = %q, want %q", err, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateStatePolicy(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
mode string
|
||||
wantErr bool
|
||||
}{
|
||||
{name: "single owner", mode: StateModeSingleOwner},
|
||||
{name: "shared root", mode: StateModeSharedRoot},
|
||||
{name: "invalid", mode: "shared", wantErr: true},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
cfg := Config{Pipelines: []Pipeline{{
|
||||
ID: "reports",
|
||||
Source: Backend{Backend: BackendLocal, Path: "/source"},
|
||||
Destinations: []Destination{{
|
||||
ID: "archive",
|
||||
Backend: BackendLocal,
|
||||
Path: "/destination",
|
||||
State: StatePolicy{Mode: tt.mode},
|
||||
}},
|
||||
}}}
|
||||
ApplyDefaults(&cfg)
|
||||
err := Validate(cfg)
|
||||
if tt.wantErr && err == nil {
|
||||
t.Fatal("Validate() error = nil, want error")
|
||||
}
|
||||
if !tt.wantErr && err != nil {
|
||||
t.Fatalf("Validate() error = %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateRetentionPolicy(t *testing.T) {
|
||||
olderThan := Duration(24 * time.Hour)
|
||||
zeroDuration := Duration(0)
|
||||
|
||||
Reference in New Issue
Block a user