Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 207c943fc6 |
+13
-20
@@ -58,26 +58,19 @@ commands:
|
|||||||
git config --global user.email "zreplbot@cschwarz.com"
|
git config --global user.email "zreplbot@cschwarz.com"
|
||||||
git config --global user.name "zrepl-github-io-ci"
|
git config --global user.name "zrepl-github-io-ci"
|
||||||
|
|
||||||
# if we're pushing, we need to add the deploy key
|
# https://circleci.com/docs/2.0/add-ssh-key/#adding-multiple-keys-with-blank-hostnames
|
||||||
# which is stored as "Additional SSH Keys" in the CircleCI project settings.
|
- run: ssh-add -D
|
||||||
# We can't use the CircleCI-manage deploy key because we're pushing
|
# the default circleci ssh config only additional ssh keys for Host !github.com
|
||||||
# to a different repo than the one we're building.
|
- run:
|
||||||
- when:
|
command: |
|
||||||
condition: << parameters.push >>
|
cat > ~/.ssh/config \<<EOF
|
||||||
steps:
|
Host *
|
||||||
# https://circleci.com/docs/2.0/add-ssh-key/#adding-multiple-keys-with-blank-hostnames
|
IdentityFile /home/circleci/.ssh/id_rsa_458e62c517f6c480e40452126ce47421
|
||||||
- run: ssh-add -D
|
EOF
|
||||||
# the default circleci ssh config only additional ssh keys for Host !github.com
|
- add_ssh_keys:
|
||||||
- run:
|
fingerprints:
|
||||||
command: |
|
# deploy key for zrepl.github.io
|
||||||
cat > ~/.ssh/config \<<EOF
|
- "45:8e:62:c5:17:f6:c4:80:e4:04:52:12:6c:e4:74:21"
|
||||||
Host *
|
|
||||||
IdentityFile /home/circleci/.ssh/id_rsa_458e62c517f6c480e40452126ce47421
|
|
||||||
EOF
|
|
||||||
- add_ssh_keys:
|
|
||||||
fingerprints:
|
|
||||||
# deploy key for zrepl.github.io
|
|
||||||
- "45:8e:62:c5:17:f6:c4:80:e4:04:52:12:6c:e4:74:21"
|
|
||||||
|
|
||||||
# caller must install-docdep
|
# caller must install-docdep
|
||||||
- when:
|
- when:
|
||||||
|
|||||||
@@ -4,17 +4,8 @@
|
|||||||
|
|
||||||
ARTIFACTDIR := artifacts
|
ARTIFACTDIR := artifacts
|
||||||
|
|
||||||
ifdef ZREPL_VERSION
|
_ZREPL_VERSION := v0.6.1
|
||||||
_ZREPL_VERSION := $(ZREPL_VERSION)
|
ZREPL_PACKAGE_RELEASE := 2
|
||||||
endif
|
|
||||||
ifndef _ZREPL_VERSION
|
|
||||||
_ZREPL_VERSION := $(shell git describe --always --dirty 2>/dev/null || echo "ZREPL_BUILD_INVALID_VERSION" )
|
|
||||||
ifeq ($(_ZREPL_VERSION),ZREPL_BUILD_INVALID_VERSION) # can't use .SHELLSTATUS because Debian Stretch is still on gmake 4.1
|
|
||||||
$(error cannot infer variable ZREPL_VERSION using git and variable is not overriden by make invocation)
|
|
||||||
endif
|
|
||||||
endif
|
|
||||||
|
|
||||||
ZREPL_PACKAGE_RELEASE := 1
|
|
||||||
|
|
||||||
GO := go
|
GO := go
|
||||||
GOOS ?= $(shell bash -c 'source <($(GO) env) && echo "$$GOOS"')
|
GOOS ?= $(shell bash -c 'source <($(GO) env) && echo "$$GOOS"')
|
||||||
|
|||||||
+1
-27
@@ -121,7 +121,6 @@ type BandwidthLimit struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type Replication struct {
|
type Replication struct {
|
||||||
Triggers []*ReplicationTriggerEnum
|
|
||||||
Protection *ReplicationOptionsProtection `yaml:"protection,optional,fromdefaults"`
|
Protection *ReplicationOptionsProtection `yaml:"protection,optional,fromdefaults"`
|
||||||
Concurrency *ReplicationOptionsConcurrency `yaml:"concurrency,optional,fromdefaults"`
|
Concurrency *ReplicationOptionsConcurrency `yaml:"concurrency,optional,fromdefaults"`
|
||||||
}
|
}
|
||||||
@@ -136,32 +135,6 @@ type ReplicationOptionsConcurrency struct {
|
|||||||
SizeEstimates int `yaml:"size_estimates,optional,default=4"`
|
SizeEstimates int `yaml:"size_estimates,optional,default=4"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type ReplicationTriggerEnum struct {
|
|
||||||
Ret interface{}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *ReplicationTriggerEnum) UnmarshalYAML(u func(interface{}, bool) error) (err error) {
|
|
||||||
t.Ret, err = enumUnmarshal(u, map[string]interface{}{
|
|
||||||
"manual": &ReplicationTriggerManual{},
|
|
||||||
"periodic": &ReplicationTriggerPeriodic{},
|
|
||||||
})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
type ReplicationTriggerManual struct {
|
|
||||||
Type string `yaml:"type"`
|
|
||||||
}
|
|
||||||
|
|
||||||
type ReplicationTriggerPeriodic struct {
|
|
||||||
Type string `yaml:"type"`
|
|
||||||
Interval *PositiveDuration `yaml:"interval"`
|
|
||||||
}
|
|
||||||
|
|
||||||
type ReplicationTriggerCron struct {
|
|
||||||
Type string `yaml:"type"`
|
|
||||||
Cron CronSpec `yaml:"cron"`
|
|
||||||
}
|
|
||||||
|
|
||||||
type PropertyRecvOptions struct {
|
type PropertyRecvOptions struct {
|
||||||
Inherit []zfsprop.Property `yaml:"inherit,optional"`
|
Inherit []zfsprop.Property `yaml:"inherit,optional"`
|
||||||
Override map[zfsprop.Property]string `yaml:"override,optional"`
|
Override map[zfsprop.Property]string `yaml:"override,optional"`
|
||||||
@@ -184,6 +157,7 @@ func (j *PushJob) GetSendOptions() *SendOptions { return j.Send }
|
|||||||
type PullJob struct {
|
type PullJob struct {
|
||||||
ActiveJob `yaml:",inline"`
|
ActiveJob `yaml:",inline"`
|
||||||
RootFS string `yaml:"root_fs"`
|
RootFS string `yaml:"root_fs"`
|
||||||
|
Interval PositiveDurationOrManual `yaml:"interval"`
|
||||||
Recv *RecvOptions `yaml:"recv,fromdefaults,optional"`
|
Recv *RecvOptions `yaml:"recv,fromdefaults,optional"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+19
-30
@@ -15,7 +15,6 @@ import (
|
|||||||
|
|
||||||
"github.com/zrepl/zrepl/config"
|
"github.com/zrepl/zrepl/config"
|
||||||
"github.com/zrepl/zrepl/daemon/job/reset"
|
"github.com/zrepl/zrepl/daemon/job/reset"
|
||||||
"github.com/zrepl/zrepl/daemon/job/trigger"
|
|
||||||
"github.com/zrepl/zrepl/daemon/job/wakeup"
|
"github.com/zrepl/zrepl/daemon/job/wakeup"
|
||||||
"github.com/zrepl/zrepl/daemon/pruner"
|
"github.com/zrepl/zrepl/daemon/pruner"
|
||||||
"github.com/zrepl/zrepl/daemon/snapper"
|
"github.com/zrepl/zrepl/daemon/snapper"
|
||||||
@@ -45,8 +44,6 @@ type ActiveSide struct {
|
|||||||
promReplicationErrors prometheus.Gauge
|
promReplicationErrors prometheus.Gauge
|
||||||
promLastSuccessful prometheus.Gauge
|
promLastSuccessful prometheus.Gauge
|
||||||
|
|
||||||
triggers *trigger.Triggers
|
|
||||||
|
|
||||||
tasksMtx sync.Mutex
|
tasksMtx sync.Mutex
|
||||||
tasks activeSideTasks
|
tasks activeSideTasks
|
||||||
}
|
}
|
||||||
@@ -93,7 +90,7 @@ type activeMode interface {
|
|||||||
SenderReceiver() (logic.Sender, logic.Receiver)
|
SenderReceiver() (logic.Sender, logic.Receiver)
|
||||||
Type() Type
|
Type() Type
|
||||||
PlannerPolicy() logic.PlannerPolicy
|
PlannerPolicy() logic.PlannerPolicy
|
||||||
RunPeriodic(ctx context.Context, wakeReplication *trigger.Manual)
|
RunPeriodic(ctx context.Context, wakeUpCommon chan<- struct{})
|
||||||
SnapperReport() *snapper.Report
|
SnapperReport() *snapper.Report
|
||||||
ResetConnectBackoff()
|
ResetConnectBackoff()
|
||||||
}
|
}
|
||||||
@@ -135,8 +132,8 @@ func (m *modePush) Type() Type { return TypePush }
|
|||||||
|
|
||||||
func (m *modePush) PlannerPolicy() logic.PlannerPolicy { return *m.plannerPolicy }
|
func (m *modePush) PlannerPolicy() logic.PlannerPolicy { return *m.plannerPolicy }
|
||||||
|
|
||||||
func (m *modePush) RunPeriodic(ctx context.Context, trigger *trigger.Manual) {
|
func (m *modePush) RunPeriodic(ctx context.Context, wakeUpCommon chan<- struct{}) {
|
||||||
m.snapper.Run(ctx, trigger)
|
m.snapper.Run(ctx, wakeUpCommon)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *modePush) SnapperReport() *snapper.Report {
|
func (m *modePush) SnapperReport() *snapper.Report {
|
||||||
@@ -224,7 +221,7 @@ func (*modePull) Type() Type { return TypePull }
|
|||||||
|
|
||||||
func (m *modePull) PlannerPolicy() logic.PlannerPolicy { return *m.plannerPolicy }
|
func (m *modePull) PlannerPolicy() logic.PlannerPolicy { return *m.plannerPolicy }
|
||||||
|
|
||||||
func (m *modePull) RunPeriodic(ctx context.Context, wakeReplication *trigger.Manual) {
|
func (m *modePull) RunPeriodic(ctx context.Context, wakeUpCommon chan<- struct{}) {
|
||||||
if m.interval.Manual {
|
if m.interval.Manual {
|
||||||
GetLogger(ctx).Info("manual pull configured, periodic pull disabled")
|
GetLogger(ctx).Info("manual pull configured, periodic pull disabled")
|
||||||
// "waiting for wakeups" is printed in common ActiveSide.do
|
// "waiting for wakeups" is printed in common ActiveSide.do
|
||||||
@@ -235,7 +232,14 @@ func (m *modePull) RunPeriodic(ctx context.Context, wakeReplication *trigger.Man
|
|||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-t.C:
|
case <-t.C:
|
||||||
wakeReplication.Fire()
|
select {
|
||||||
|
case wakeUpCommon <- struct{}{}:
|
||||||
|
default:
|
||||||
|
GetLogger(ctx).
|
||||||
|
WithField("pull_interval", m.interval).
|
||||||
|
Warn("pull job took longer than pull interval")
|
||||||
|
wakeUpCommon <- struct{}{} // block anyways, to queue up the wakeup
|
||||||
|
}
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -366,11 +370,6 @@ func activeSide(g *config.Global, in *config.ActiveJob, configJob interface{}, p
|
|||||||
return nil, errors.Wrap(err, "cannot build replication driver config")
|
return nil, errors.Wrap(err, "cannot build replication driver config")
|
||||||
}
|
}
|
||||||
|
|
||||||
j.triggers, err = trigger.FromConfig(in.Replication.Triggers)
|
|
||||||
if err != nil {
|
|
||||||
return nil, errors.Wrap(err, "cannot build triggers")
|
|
||||||
}
|
|
||||||
|
|
||||||
return j, nil
|
return j, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -445,17 +444,12 @@ func (j *ActiveSide) Run(ctx context.Context) {
|
|||||||
|
|
||||||
defer log.Info("job exiting")
|
defer log.Info("job exiting")
|
||||||
|
|
||||||
|
periodicDone := make(chan struct{})
|
||||||
ctx, cancel := context.WithCancel(ctx)
|
ctx, cancel := context.WithCancel(ctx)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
periodicCtx, endTask := trace.WithTask(ctx, "periodic")
|
||||||
periodCtx, endTask := trace.WithTask(ctx, "periodic")
|
|
||||||
defer endTask()
|
|
||||||
go j.mode.RunPeriodic(periodCtx, periodicTrigger)
|
|
||||||
|
|
||||||
wakeupTrigger := wakeup.Trigger(ctx)
|
|
||||||
|
|
||||||
triggered, endTask := j.triggers.Spawn(ctx, []*trigger.Trigger{periodicTrigger, wakeupTrigger})
|
|
||||||
defer endTask()
|
defer endTask()
|
||||||
|
go j.mode.RunPeriodic(periodicCtx, periodicDone)
|
||||||
|
|
||||||
invocationCount := 0
|
invocationCount := 0
|
||||||
outer:
|
outer:
|
||||||
@@ -465,15 +459,10 @@ outer:
|
|||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
log.WithError(ctx.Err()).Info("context")
|
log.WithError(ctx.Err()).Info("context")
|
||||||
break outer
|
break outer
|
||||||
case trigger := <-triggered:
|
|
||||||
log :=
|
case <-wakeup.Wait(ctx):
|
||||||
log.WithField("trigger_id", trigger.ID())
|
j.mode.ResetConnectBackoff()
|
||||||
log.Info("triggered")
|
case <-periodicDone:
|
||||||
switch trigger {
|
|
||||||
case wakeupTrigger:
|
|
||||||
log.Info("trigger is wakeup command, resetting connection backoff")
|
|
||||||
j.mode.ResetConnectBackoff()
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
invocationCount++
|
invocationCount++
|
||||||
invocationCtx, endSpan := trace.WithSpan(ctx, fmt.Sprintf("invocation-%d", invocationCount))
|
invocationCtx, endSpan := trace.WithSpan(ctx, fmt.Sprintf("invocation-%d", invocationCount))
|
||||||
|
|||||||
@@ -24,6 +24,21 @@ func JobsFromConfig(c *config.Config, parseFlags config.ParseFlags) ([]Job, erro
|
|||||||
js[i] = j
|
js[i] = j
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// receiving-side root filesystems must not overlap
|
||||||
|
{
|
||||||
|
rfss := make([]string, 0, len(js))
|
||||||
|
for _, j := range js {
|
||||||
|
jrfs, ok := j.OwnedDatasetSubtreeRoot()
|
||||||
|
if !ok {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
rfss = append(rfss, jrfs.ToString())
|
||||||
|
}
|
||||||
|
if err := validateReceivingSidesDoNotOverlap(rfss); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
return js, nil
|
return js, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+5
-10
@@ -15,7 +15,6 @@ import (
|
|||||||
|
|
||||||
"github.com/zrepl/zrepl/config"
|
"github.com/zrepl/zrepl/config"
|
||||||
"github.com/zrepl/zrepl/daemon/filters"
|
"github.com/zrepl/zrepl/daemon/filters"
|
||||||
"github.com/zrepl/zrepl/daemon/job/trigger"
|
|
||||||
"github.com/zrepl/zrepl/daemon/job/wakeup"
|
"github.com/zrepl/zrepl/daemon/job/wakeup"
|
||||||
"github.com/zrepl/zrepl/daemon/pruner"
|
"github.com/zrepl/zrepl/daemon/pruner"
|
||||||
"github.com/zrepl/zrepl/daemon/snapper"
|
"github.com/zrepl/zrepl/daemon/snapper"
|
||||||
@@ -105,18 +104,12 @@ func (j *SnapJob) Run(ctx context.Context) {
|
|||||||
|
|
||||||
defer log.Info("job exiting")
|
defer log.Info("job exiting")
|
||||||
|
|
||||||
wakeupTrigger := wakeup.Trigger(ctx)
|
periodicDone := make(chan struct{})
|
||||||
|
|
||||||
snapshottingTrigger := trigger.New("periodic")
|
|
||||||
ctx, cancel := context.WithCancel(ctx)
|
ctx, cancel := context.WithCancel(ctx)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
periodicCtx, endTask := trace.WithTask(ctx, "snapshotting")
|
periodicCtx, endTask := trace.WithTask(ctx, "snapshotting")
|
||||||
defer endTask()
|
defer endTask()
|
||||||
go j.snapper.Run(periodicCtx, snapshottingTrigger)
|
go j.snapper.Run(periodicCtx, periodicDone)
|
||||||
|
|
||||||
triggers := trigger.Empty()
|
|
||||||
triggered, endTask := triggers.Spawn(ctx, []trigger.Trigger{snapshottingTrigger, wakeupTrigger})
|
|
||||||
defer endTask()
|
|
||||||
|
|
||||||
invocationCount := 0
|
invocationCount := 0
|
||||||
outer:
|
outer:
|
||||||
@@ -126,7 +119,9 @@ outer:
|
|||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
log.WithError(ctx.Err()).Info("context")
|
log.WithError(ctx.Err()).Info("context")
|
||||||
break outer
|
break outer
|
||||||
case <-triggered:
|
|
||||||
|
case <-wakeup.Wait(ctx):
|
||||||
|
case <-periodicDone:
|
||||||
}
|
}
|
||||||
invocationCount++
|
invocationCount++
|
||||||
|
|
||||||
|
|||||||
@@ -1,23 +0,0 @@
|
|||||||
package trigger
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
|
|
||||||
"github.com/robfig/cron/v3"
|
|
||||||
)
|
|
||||||
|
|
||||||
type Cron struct {
|
|
||||||
spec cron.Schedule
|
|
||||||
}
|
|
||||||
|
|
||||||
var _ Trigger = &Cron{}
|
|
||||||
|
|
||||||
func NewCron(spec cron.Schedule) *Cron {
|
|
||||||
return &Cron{spec: spec}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *Cron) ID() string { return "cron" }
|
|
||||||
|
|
||||||
func (t *Cron) run(ctx context.Context, signal chan<- struct{}) {
|
|
||||||
panic("unimpl: extract from cron snapper")
|
|
||||||
}
|
|
||||||
@@ -1,30 +0,0 @@
|
|||||||
package trigger
|
|
||||||
|
|
||||||
import (
|
|
||||||
"fmt"
|
|
||||||
|
|
||||||
"github.com/zrepl/zrepl/config"
|
|
||||||
)
|
|
||||||
|
|
||||||
func FromConfig(in []*config.ReplicationTriggerEnum) (*Triggers, error) {
|
|
||||||
triggers := make([]Trigger, len(in))
|
|
||||||
for i, e := range in {
|
|
||||||
var t Trigger = nil
|
|
||||||
switch te := e.Ret.(type) {
|
|
||||||
case *config.ReplicationTriggerManual:
|
|
||||||
// not a trigger
|
|
||||||
t = NewManual("manual")
|
|
||||||
case *config.ReplicationTriggerPeriodic:
|
|
||||||
t = NewPeriodic(te.Interval.Duration())
|
|
||||||
case *config.ReplicationTriggerCron:
|
|
||||||
t = NewCron(te.Cron.Schedule)
|
|
||||||
default:
|
|
||||||
return nil, fmt.Errorf("unknown trigger type %T", te)
|
|
||||||
}
|
|
||||||
triggers[i] = t
|
|
||||||
}
|
|
||||||
return &Triggers{
|
|
||||||
spawned: false,
|
|
||||||
triggers: triggers,
|
|
||||||
}, nil
|
|
||||||
}
|
|
||||||
@@ -1,12 +0,0 @@
|
|||||||
package trigger
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
|
|
||||||
"github.com/zrepl/zrepl/logger"
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
func getLogger(ctx context.Context) logger.Logger {
|
|
||||||
panic("unimpl")
|
|
||||||
}
|
|
||||||
@@ -1,33 +0,0 @@
|
|||||||
package trigger
|
|
||||||
|
|
||||||
import "context"
|
|
||||||
|
|
||||||
type Manual struct {
|
|
||||||
id string
|
|
||||||
signal chan<- struct{}
|
|
||||||
}
|
|
||||||
|
|
||||||
var _ Trigger = &Manual{}
|
|
||||||
|
|
||||||
func NewManual(id string) *Manual {
|
|
||||||
return &Manual{
|
|
||||||
id: id,
|
|
||||||
signal: nil,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *Manual) ID() string {
|
|
||||||
return t.id
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *Manual) run(ctx context.Context, signal chan<- struct{}) {
|
|
||||||
if t.signal != nil {
|
|
||||||
panic("run must only be called once")
|
|
||||||
}
|
|
||||||
t.signal = signal
|
|
||||||
}
|
|
||||||
|
|
||||||
// Panics if called before the trigger has been spanwed as part of a `Triggers`.
|
|
||||||
func (t *Manual) Fire() {
|
|
||||||
t.signal <- struct{}{}
|
|
||||||
}
|
|
||||||
@@ -1,33 +0,0 @@
|
|||||||
package trigger
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"time"
|
|
||||||
)
|
|
||||||
|
|
||||||
type Periodic struct {
|
|
||||||
interval time.Duration
|
|
||||||
}
|
|
||||||
|
|
||||||
var _ Trigger = &Periodic{}
|
|
||||||
|
|
||||||
func NewPeriodic(interval time.Duration) *Periodic {
|
|
||||||
return &Periodic{
|
|
||||||
interval: interval,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (p *Periodic) ID() string { return "periodic" }
|
|
||||||
|
|
||||||
func (p *Periodic) run(ctx context.Context, signal chan<- struct{}) {
|
|
||||||
t := time.NewTicker(p.interval)
|
|
||||||
defer t.Stop()
|
|
||||||
for {
|
|
||||||
select {
|
|
||||||
case <-t.C:
|
|
||||||
signal <- struct{}{}
|
|
||||||
case <-ctx.Done():
|
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,91 +0,0 @@
|
|||||||
package trigger
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"github.com/zrepl/zrepl/daemon/logging/trace"
|
|
||||||
)
|
|
||||||
|
|
||||||
type Triggers struct {
|
|
||||||
spawned bool
|
|
||||||
triggers []Trigger
|
|
||||||
}
|
|
||||||
|
|
||||||
type Trigger interface {
|
|
||||||
ID() string
|
|
||||||
run(context.Context, chan<- struct{})
|
|
||||||
}
|
|
||||||
|
|
||||||
func Empty() *Triggers {
|
|
||||||
return &Triggers{
|
|
||||||
spawned: false,
|
|
||||||
triggers: nil,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *Triggers) Spawn(ctx context.Context, additionalTriggers []Trigger) (chan Trigger, trace.DoneFunc) {
|
|
||||||
if t.spawned {
|
|
||||||
panic("must only spawn once")
|
|
||||||
}
|
|
||||||
t.spawned = true
|
|
||||||
t.triggers = append(t.triggers, additionalTriggers...)
|
|
||||||
sink := make(chan Trigger)
|
|
||||||
endTask := t.spawn(ctx, sink)
|
|
||||||
return sink, endTask
|
|
||||||
}
|
|
||||||
|
|
||||||
type triggering struct {
|
|
||||||
trigger Trigger
|
|
||||||
handled chan struct{}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *Triggers) spawn(ctx context.Context, sink chan Trigger) trace.DoneFunc {
|
|
||||||
ctx, endTask := trace.WithTask(ctx, "triggers")
|
|
||||||
ctx, add, wait := trace.WithTaskGroup(ctx, "trigger-tasks")
|
|
||||||
triggered := make(chan triggering, len(t.triggers))
|
|
||||||
for _, t := range t.triggers {
|
|
||||||
t := t
|
|
||||||
signal := make(chan struct{})
|
|
||||||
go add(func(ctx context.Context) {
|
|
||||||
t.run(ctx, signal)
|
|
||||||
})
|
|
||||||
go func() {
|
|
||||||
for {
|
|
||||||
select {
|
|
||||||
case <-ctx.Done():
|
|
||||||
return
|
|
||||||
case <-signal:
|
|
||||||
handled := make(chan struct{})
|
|
||||||
select {
|
|
||||||
case triggered <- triggering{trigger: t, handled: handled}:
|
|
||||||
default:
|
|
||||||
panic("this funtion ensures that there's always room in the channel")
|
|
||||||
}
|
|
||||||
select {
|
|
||||||
case <-handled:
|
|
||||||
case <-ctx.Done():
|
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
}
|
|
||||||
go func() {
|
|
||||||
defer wait()
|
|
||||||
for {
|
|
||||||
select {
|
|
||||||
case <-ctx.Done():
|
|
||||||
return
|
|
||||||
case triggering := <-triggered:
|
|
||||||
select {
|
|
||||||
case sink <- triggering.trigger:
|
|
||||||
default:
|
|
||||||
getLogger(ctx).
|
|
||||||
WithField("trigger_id", triggering.trigger.ID()).
|
|
||||||
Warn("dropping triggering because job is busy")
|
|
||||||
}
|
|
||||||
close(triggering.handled)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
return endTask
|
|
||||||
}
|
|
||||||
@@ -3,8 +3,6 @@ package wakeup
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
|
|
||||||
"github.com/zrepl/zrepl/daemon/job/trigger"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
type contextKey int
|
type contextKey int
|
||||||
@@ -19,10 +17,6 @@ func Wait(ctx context.Context) <-chan struct{} {
|
|||||||
return wc
|
return wc
|
||||||
}
|
}
|
||||||
|
|
||||||
func Trigger(ctx context.Context) trigger.Trigger {
|
|
||||||
panic("unimpl")
|
|
||||||
}
|
|
||||||
|
|
||||||
type Func func() error
|
type Func func() error
|
||||||
|
|
||||||
var AlreadyWokenUp = errors.New("already woken up")
|
var AlreadyWokenUp = errors.New("already woken up")
|
||||||
|
|||||||
@@ -63,7 +63,6 @@ type Subsystem string
|
|||||||
const (
|
const (
|
||||||
SubsysMeta Subsystem = "meta"
|
SubsysMeta Subsystem = "meta"
|
||||||
SubsysJob Subsystem = "job"
|
SubsysJob Subsystem = "job"
|
||||||
SubsysTrigger Subsystem = "trigger"
|
|
||||||
SubsysReplication Subsystem = "repl"
|
SubsysReplication Subsystem = "repl"
|
||||||
SubsysEndpoint Subsystem = "endpoint"
|
SubsysEndpoint Subsystem = "endpoint"
|
||||||
SubsysPruning Subsystem = "pruning"
|
SubsysPruning Subsystem = "pruning"
|
||||||
|
|||||||
@@ -10,7 +10,6 @@ import (
|
|||||||
|
|
||||||
"github.com/zrepl/zrepl/config"
|
"github.com/zrepl/zrepl/config"
|
||||||
"github.com/zrepl/zrepl/daemon/hooks"
|
"github.com/zrepl/zrepl/daemon/hooks"
|
||||||
"github.com/zrepl/zrepl/daemon/job/trigger"
|
|
||||||
"github.com/zrepl/zrepl/util/suspendresumesafetimer"
|
"github.com/zrepl/zrepl/util/suspendresumesafetimer"
|
||||||
"github.com/zrepl/zrepl/zfs"
|
"github.com/zrepl/zrepl/zfs"
|
||||||
)
|
)
|
||||||
@@ -43,7 +42,7 @@ type Cron struct {
|
|||||||
wakeupWhileRunningCount int
|
wakeupWhileRunningCount int
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Cron) Run(ctx context.Context, snapshotsTaken *trigger.Manual) {
|
func (s *Cron) Run(ctx context.Context, snapshotsTaken chan<- struct{}) {
|
||||||
|
|
||||||
for {
|
for {
|
||||||
now := time.Now()
|
now := time.Now()
|
||||||
@@ -76,7 +75,13 @@ func (s *Cron) Run(ctx context.Context, snapshotsTaken *trigger.Manual) {
|
|||||||
s.running = false
|
s.running = false
|
||||||
s.mtx.Unlock()
|
s.mtx.Unlock()
|
||||||
|
|
||||||
snapshotsTaken.Fire()
|
select {
|
||||||
|
case snapshotsTaken <- struct{}{}:
|
||||||
|
default:
|
||||||
|
if snapshotsTaken != nil {
|
||||||
|
getLogger(ctx).Warn("callback channel is full, discarding snapshot update event")
|
||||||
|
}
|
||||||
|
}
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -2,13 +2,11 @@ package snapper
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
|
||||||
"github.com/zrepl/zrepl/daemon/job/trigger"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
type manual struct{}
|
type manual struct{}
|
||||||
|
|
||||||
func (s *manual) Run(ctx context.Context, snapshotsTaken *trigger.Manual) {
|
func (s *manual) Run(ctx context.Context, wakeUpCommon chan<- struct{}) {
|
||||||
// nothing to do
|
// nothing to do
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -9,7 +9,6 @@ import (
|
|||||||
|
|
||||||
"github.com/pkg/errors"
|
"github.com/pkg/errors"
|
||||||
|
|
||||||
"github.com/zrepl/zrepl/daemon/job/trigger"
|
|
||||||
"github.com/zrepl/zrepl/daemon/logging/trace"
|
"github.com/zrepl/zrepl/daemon/logging/trace"
|
||||||
|
|
||||||
"github.com/zrepl/zrepl/config"
|
"github.com/zrepl/zrepl/config"
|
||||||
@@ -52,7 +51,7 @@ type periodicArgs struct {
|
|||||||
interval time.Duration
|
interval time.Duration
|
||||||
fsf zfs.DatasetFilter
|
fsf zfs.DatasetFilter
|
||||||
planArgs planArgs
|
planArgs planArgs
|
||||||
snapshotsTaken *trigger.Manual
|
snapshotsTaken chan<- struct{}
|
||||||
dryRun bool
|
dryRun bool
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -104,7 +103,7 @@ func (s State) sf() state {
|
|||||||
type updater func(u func(*Periodic)) State
|
type updater func(u func(*Periodic)) State
|
||||||
type state func(a periodicArgs, u updater) state
|
type state func(a periodicArgs, u updater) state
|
||||||
|
|
||||||
func (s *Periodic) Run(ctx context.Context, snapshotsTaken *trigger.Manual) {
|
func (s *Periodic) Run(ctx context.Context, snapshotsTaken chan<- struct{}) {
|
||||||
defer trace.WithSpanFromStackUpdateCtx(&ctx)()
|
defer trace.WithSpanFromStackUpdateCtx(&ctx)()
|
||||||
getLogger(ctx).Debug("start")
|
getLogger(ctx).Debug("start")
|
||||||
defer getLogger(ctx).Debug("stop")
|
defer getLogger(ctx).Debug("stop")
|
||||||
@@ -208,7 +207,13 @@ func periodicStateSnapshot(a periodicArgs, u updater) state {
|
|||||||
|
|
||||||
ok := plan.execute(a.ctx, false)
|
ok := plan.execute(a.ctx, false)
|
||||||
|
|
||||||
a.snapshotsTaken.Fire()
|
select {
|
||||||
|
case a.snapshotsTaken <- struct{}{}:
|
||||||
|
default:
|
||||||
|
if a.snapshotsTaken != nil {
|
||||||
|
getLogger(a.ctx).Warn("callback channel is full, discarding snapshot update event")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
return u(func(snapper *Periodic) {
|
return u(func(snapper *Periodic) {
|
||||||
if !ok {
|
if !ok {
|
||||||
|
|||||||
@@ -5,7 +5,6 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
|
|
||||||
"github.com/zrepl/zrepl/config"
|
"github.com/zrepl/zrepl/config"
|
||||||
"github.com/zrepl/zrepl/daemon/job/trigger"
|
|
||||||
"github.com/zrepl/zrepl/zfs"
|
"github.com/zrepl/zrepl/zfs"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -18,7 +17,7 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type Snapper interface {
|
type Snapper interface {
|
||||||
Run(ctx context.Context, snapshotsTaken *trigger.Manual)
|
Run(ctx context.Context, snapshotsTaken chan<- struct{})
|
||||||
Report() Report
|
Report() Report
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -258,9 +258,6 @@ On your setup, ensure that
|
|||||||
|
|
||||||
* all ``filesystems`` filter specifications are disjoint
|
* all ``filesystems`` filter specifications are disjoint
|
||||||
* no ``root_fs`` is a prefix or equal to another ``root_fs``
|
* no ``root_fs`` is a prefix or equal to another ``root_fs``
|
||||||
|
|
||||||
* For ``sink`` jobs, consider all possible ``root_fs/${client_identity}``.
|
|
||||||
|
|
||||||
* no ``filesystems`` filter matches any ``root_fs``
|
* no ``filesystems`` filter matches any ``root_fs``
|
||||||
|
|
||||||
**Exceptions to the rule**:
|
**Exceptions to the rule**:
|
||||||
|
|||||||
@@ -33,6 +33,11 @@ func NewKeepLastN(n int, regex string) (*KeepLastN, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (k KeepLastN) KeepRule(snaps []Snapshot) (destroyList []Snapshot) {
|
func (k KeepLastN) KeepRule(snaps []Snapshot) (destroyList []Snapshot) {
|
||||||
|
|
||||||
|
if k.n > len(snaps) {
|
||||||
|
return []Snapshot{}
|
||||||
|
}
|
||||||
|
|
||||||
matching, notMatching := partitionSnapList(snaps, func(snapshot Snapshot) bool {
|
matching, notMatching := partitionSnapList(snaps, func(snapshot Snapshot) bool {
|
||||||
return k.re.MatchString(snapshot.Name())
|
return k.re.MatchString(snapshot.Name())
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -90,7 +90,7 @@ func TestKeepLastN(t *testing.T) {
|
|||||||
stubSnap{"a2", false, o(12)},
|
stubSnap{"a2", false, o(12)},
|
||||||
},
|
},
|
||||||
rules: []KeepRule{
|
rules: []KeepRule{
|
||||||
MustKeepLastN(4, "a"),
|
MustKeepLastN(3, "a"),
|
||||||
},
|
},
|
||||||
expDestroy: map[string]bool{
|
expDestroy: map[string]bool{
|
||||||
"b1": true,
|
"b1": true,
|
||||||
|
|||||||
Reference in New Issue
Block a user