Pruner: refactor + use Task API

refs #10
This commit is contained in:
Christian Schwarz
2017-12-26 20:38:27 +01:00
parent 13562b48ed
commit 91c4a97f72
+55 -22
View File
@@ -23,46 +23,56 @@ type PruneResult struct {
Remove []zfs.FilesystemVersion Remove []zfs.FilesystemVersion
} }
func (p *Pruner) Run(ctx context.Context) (r []PruneResult, err error) { func (p *Pruner) filterFilesystems() (filesystems []*zfs.DatasetPath, stop bool) {
p.task.Enter("filter_fs")
log := ctx.Value(contextKeyLog).(Logger) defer p.task.Finish()
if p.DryRun {
log.Info("doing dry run")
}
filesystems, err := zfs.ZFSListMapping(p.DatasetFilter) filesystems, err := zfs.ZFSListMapping(p.DatasetFilter)
if err != nil { if err != nil {
log.WithError(err).Error("error applying filesystem filter") p.task.Log().WithError(err).Error("error applying filesystem filter")
return nil, err return nil, true
} }
if len(filesystems) <= 0 { if len(filesystems) <= 0 {
log.Info("no filesystems matching filter") p.task.Log().Info("no filesystems matching filter")
return nil, err return nil, true
}
return filesystems, false
} }
r = make([]PruneResult, 0, len(filesystems)) func (p *Pruner) filterVersions(fs *zfs.DatasetPath) (fsversions []zfs.FilesystemVersion, stop bool) {
p.task.Enter("filter_versions")
for _, fs := range filesystems { defer p.task.Finish()
log := p.task.Log().WithField(logFSField, fs.ToString())
log := log.WithField(logFSField, fs.ToString())
// only prune snapshots, bookmarks are kept forever // only prune snapshots, bookmarks are kept forever
snapshotFilter := NewTypedPrefixFilter(p.SnapshotPrefix, zfs.Snapshot) snapshotFilter := NewTypedPrefixFilter(p.SnapshotPrefix, zfs.Snapshot)
fsversions, err := zfs.ZFSListFilesystemVersions(fs, snapshotFilter) fsversions, err := zfs.ZFSListFilesystemVersions(fs, snapshotFilter)
if err != nil { if err != nil {
log.WithError(err).Error("error listing filesytem versions") log.WithError(err).Error("error listing filesytem versions")
continue return nil, true
} }
if len(fsversions) == 0 { if len(fsversions) == 0 {
log.WithField("prefix", p.SnapshotPrefix).Info("no filesystem versions matching prefix") log.WithField("prefix", p.SnapshotPrefix).Info("no filesystem versions matching prefix")
continue return nil, true
}
return fsversions, false
} }
func (p *Pruner) pruneFilesystem(fs *zfs.DatasetPath) (r PruneResult, valid bool) {
p.task.Enter("prune_fs")
defer p.task.Finish()
log := p.task.Log().WithField(logFSField, fs.ToString())
fsversions, stop := p.filterVersions(fs)
if stop {
return
}
p.task.Enter("prune_policy")
keep, remove, err := p.PrunePolicy.Prune(fs, fsversions) keep, remove, err := p.PrunePolicy.Prune(fs, fsversions)
p.task.Finish()
if err != nil { if err != nil {
log.WithError(err).Error("error evaluating prune policy") log.WithError(err).Error("error evaluating prune policy")
continue return
} }
log.WithField("fsversions", fsversions). log.WithField("fsversions", fsversions).
@@ -70,7 +80,7 @@ func (p *Pruner) Run(ctx context.Context) (r []PruneResult, err error) {
WithField("remove", remove). WithField("remove", remove).
Debug("prune policy debug dump") Debug("prune policy debug dump")
r = append(r, PruneResult{fs, fsversions, keep, remove}) r = PruneResult{fs, fsversions, keep, remove}
makeFields := func(v zfs.FilesystemVersion) (fields map[string]interface{}) { makeFields := func(v zfs.FilesystemVersion) (fields map[string]interface{}) {
fields = make(map[string]interface{}) fields = make(map[string]interface{})
@@ -91,14 +101,37 @@ func (p *Pruner) Run(ctx context.Context) (r []PruneResult, err error) {
// TODO special handling for EBUSY (zfs hold) // TODO special handling for EBUSY (zfs hold)
// TODO error handling for clones? just echo to cli, skip over, and exit with non-zero status code (we're idempotent) // TODO error handling for clones? just echo to cli, skip over, and exit with non-zero status code (we're idempotent)
if !p.DryRun { if !p.DryRun {
p.task.Enter("destroy")
err := zfs.ZFSDestroyFilesystemVersion(fs, v) err := zfs.ZFSDestroyFilesystemVersion(fs, v)
p.task.Finish()
if err != nil { if err != nil {
// handle
log.WithFields(fields).WithError(err).Error("error destroying version") log.WithFields(fields).WithError(err).Error("error destroying version")
} }
} }
} }
return r, true
}
func (p *Pruner) Run(ctx context.Context) (r []PruneResult, err error) {
p.task.Enter("run")
defer p.task.Finish()
if p.DryRun {
p.task.Log().Info("doing dry run")
}
filesystems, stop := p.filterFilesystems()
if stop {
return
}
r = make([]PruneResult, 0, len(filesystems))
for _, fs := range filesystems {
res, ok := p.pruneFilesystem(fs)
if ok {
r = append(r, res)
}
} }
return return