zfs: introduce pkg zfs/zfscmd for command logging, status, prometheus metrics
refs #196
This commit is contained in:
@@ -0,0 +1,122 @@
|
||||
// Package zfscmd provides a wrapper around packate os/exec.
|
||||
// Functionality provided by the wrapper:
|
||||
// - logging start and end of command execution
|
||||
// - status report of active commands
|
||||
// - prometheus metrics of runtimes
|
||||
package zfscmd
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/zrepl/zrepl/util/circlog"
|
||||
)
|
||||
|
||||
type Cmd struct {
|
||||
cmd *exec.Cmd
|
||||
ctx context.Context
|
||||
mtx sync.RWMutex
|
||||
startedAt, waitReturnedAt time.Time
|
||||
}
|
||||
|
||||
func CommandContext(ctx context.Context, name string, arg ...string) *Cmd {
|
||||
cmd := exec.CommandContext(ctx, name, arg...)
|
||||
return &Cmd{cmd: cmd, ctx: ctx}
|
||||
}
|
||||
|
||||
// err.(*exec.ExitError).Stderr will NOT be set
|
||||
func (c *Cmd) CombinedOutput() (o []byte, err error) {
|
||||
o, err = c.cmd.CombinedOutput()
|
||||
return
|
||||
}
|
||||
|
||||
// err.(*exec.ExitError).Stderr will be set
|
||||
func (c *Cmd) Output() (o []byte, err error) {
|
||||
o, err = c.cmd.Output()
|
||||
return
|
||||
}
|
||||
|
||||
// Careful: err.(*exec.ExitError).Stderr will not be set, even if you don't open an StderrPipe
|
||||
func (c *Cmd) StdoutPipeWithErrorBuf() (p io.ReadCloser, errBuf *circlog.CircularLog, err error) {
|
||||
p, err = c.cmd.StdoutPipe()
|
||||
errBuf = circlog.MustNewCircularLog(1 << 15)
|
||||
c.cmd.Stderr = errBuf
|
||||
return p, errBuf, err
|
||||
}
|
||||
|
||||
type Stdio struct {
|
||||
Stdin io.ReadCloser
|
||||
Stdout io.Writer
|
||||
Stderr io.Writer
|
||||
}
|
||||
|
||||
func (c *Cmd) SetStdio(stdio Stdio) {
|
||||
c.cmd.Stdin = stdio.Stdin
|
||||
c.cmd.Stderr = stdio.Stderr
|
||||
c.cmd.Stdout = stdio.Stdout
|
||||
}
|
||||
|
||||
func (c *Cmd) String() string {
|
||||
return strings.Join(c.cmd.Args, " ") // includes argv[0] if initialized with CommandContext, that's the only way we o it
|
||||
}
|
||||
|
||||
func (c *Cmd) log() Logger {
|
||||
return getLogger(c.ctx).WithField("cmd", c.String())
|
||||
}
|
||||
|
||||
func (c *Cmd) Start() (err error) {
|
||||
startPreLogging(c, time.Now())
|
||||
|
||||
err = c.cmd.Start()
|
||||
now := time.Now()
|
||||
|
||||
c.mtx.Lock()
|
||||
c.startedAt = now
|
||||
c.mtx.Unlock()
|
||||
|
||||
startPostReport(c, err, now)
|
||||
startPostLogging(c, err, now)
|
||||
return err
|
||||
}
|
||||
|
||||
// only call this after a successful call to .Start()
|
||||
func (c *Cmd) Process() *os.Process {
|
||||
if c.startedAt.IsZero() {
|
||||
panic("calling Process() only allowed after successful call to Start()")
|
||||
}
|
||||
return c.cmd.Process
|
||||
}
|
||||
|
||||
func (c *Cmd) Wait() (err error) {
|
||||
waitPreLogging(c, time.Now())
|
||||
|
||||
err = c.cmd.Wait()
|
||||
now := time.Now()
|
||||
|
||||
if !c.waitReturnedAt.IsZero() {
|
||||
// ignore duplicate waits
|
||||
return
|
||||
}
|
||||
|
||||
c.mtx.Lock()
|
||||
c.waitReturnedAt = now
|
||||
c.mtx.Unlock()
|
||||
|
||||
waitPostReport(c, now)
|
||||
waitPostLogging(c, err, now)
|
||||
waitPostPrometheus(c, err, now)
|
||||
return err
|
||||
}
|
||||
|
||||
// returns 0 if the command did not yet finish
|
||||
func (c *Cmd) Runtime() time.Duration {
|
||||
if c.waitReturnedAt.IsZero() {
|
||||
return 0
|
||||
}
|
||||
return c.waitReturnedAt.Sub(c.startedAt)
|
||||
}
|
||||
@@ -0,0 +1,39 @@
|
||||
package zfscmd
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/zrepl/zrepl/logger"
|
||||
)
|
||||
|
||||
type contextKey int
|
||||
|
||||
const (
|
||||
contextKeyLogger contextKey = iota
|
||||
contextKeyJobID
|
||||
)
|
||||
|
||||
type Logger = logger.Logger
|
||||
|
||||
func WithJobID(ctx context.Context, jobID string) context.Context {
|
||||
return context.WithValue(ctx, contextKeyJobID, jobID)
|
||||
}
|
||||
|
||||
func getJobIDOrDefault(ctx context.Context, def string) string {
|
||||
ret, ok := ctx.Value(contextKeyJobID).(string)
|
||||
if !ok {
|
||||
return def
|
||||
}
|
||||
return ret
|
||||
}
|
||||
|
||||
func WithLogger(ctx context.Context, log Logger) context.Context {
|
||||
return context.WithValue(ctx, contextKeyLogger, log)
|
||||
}
|
||||
|
||||
func getLogger(ctx context.Context) Logger {
|
||||
if l, ok := ctx.Value(contextKeyLogger).(Logger); ok {
|
||||
return l
|
||||
}
|
||||
return logger.NewNullLogger()
|
||||
}
|
||||
@@ -0,0 +1,40 @@
|
||||
package zfscmd
|
||||
|
||||
import (
|
||||
"time"
|
||||
)
|
||||
|
||||
// Implementation Note:
|
||||
//
|
||||
// Pre-events logged with debug
|
||||
// Post-event without error logged with info
|
||||
// Post-events with error logged at error level
|
||||
|
||||
func startPreLogging(c *Cmd, now time.Time) {
|
||||
c.log().Debug("starting command")
|
||||
}
|
||||
|
||||
func startPostLogging(c *Cmd, err error, now time.Time) {
|
||||
if err == nil {
|
||||
c.log().Info("started command")
|
||||
} else {
|
||||
c.log().WithError(err).Error("cannot start command")
|
||||
}
|
||||
}
|
||||
|
||||
func waitPreLogging(c *Cmd, now time.Time) {
|
||||
c.log().Debug("start waiting")
|
||||
}
|
||||
|
||||
func waitPostLogging(c *Cmd, err error, now time.Time) {
|
||||
log := c.log().
|
||||
WithField("total_time_s", c.Runtime().Seconds()).
|
||||
WithField("systemtime_s", c.cmd.ProcessState.SystemTime().Seconds()).
|
||||
WithField("usertime_s", c.cmd.ProcessState.UserTime().Seconds())
|
||||
|
||||
if err == nil {
|
||||
log.Info("command exited without error")
|
||||
} else {
|
||||
log.WithError(err).Error("command exited with error")
|
||||
}
|
||||
}
|
||||
Executable
+7
@@ -0,0 +1,7 @@
|
||||
#!/usr/bin/env bash
|
||||
set -euo pipefail
|
||||
|
||||
echo "to stderr" 1>&2
|
||||
echo "to stdout"
|
||||
|
||||
exit "$1"
|
||||
@@ -0,0 +1,60 @@
|
||||
package zfscmd
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"io"
|
||||
"os/exec"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
const testBin = "./zfscmd_platform_test.bash"
|
||||
|
||||
func TestCmdStderrBehaviorOutput(t *testing.T) {
|
||||
|
||||
stdout, err := exec.Command(testBin, "0").Output()
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, []byte("to stdout\n"), stdout)
|
||||
|
||||
stdout, err = exec.Command(testBin, "1").Output()
|
||||
assert.Equal(t, []byte("to stdout\n"), stdout)
|
||||
require.Error(t, err)
|
||||
ee, ok := err.(*exec.ExitError)
|
||||
require.True(t, ok)
|
||||
require.Equal(t, ee.Stderr, []byte("to stderr\n"))
|
||||
}
|
||||
|
||||
func TestCmdStderrBehaviorCombinedOutput(t *testing.T) {
|
||||
|
||||
stdio, err := exec.Command(testBin, "0").CombinedOutput()
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, "to stderr\nto stdout\n", string(stdio))
|
||||
|
||||
stdio, err = exec.Command(testBin, "1").CombinedOutput()
|
||||
require.Equal(t, "to stderr\nto stdout\n", string(stdio))
|
||||
require.Error(t, err)
|
||||
ee, ok := err.(*exec.ExitError)
|
||||
require.True(t, ok)
|
||||
require.Empty(t, ee.Stderr) // !!!! maybe not what one would expect
|
||||
}
|
||||
|
||||
func TestCmdStderrBehaviorStdoutPipe(t *testing.T) {
|
||||
cmd := exec.Command(testBin, "1")
|
||||
stdoutPipe, err := cmd.StdoutPipe()
|
||||
require.NoError(t, err)
|
||||
err = cmd.Start()
|
||||
require.NoError(t, err)
|
||||
defer cmd.Wait()
|
||||
var stdout bytes.Buffer
|
||||
_, err = io.Copy(&stdout, stdoutPipe)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, "to stdout\n", stdout.String())
|
||||
|
||||
err = cmd.Wait()
|
||||
require.Error(t, err)
|
||||
ee, ok := err.(*exec.ExitError)
|
||||
require.True(t, ok)
|
||||
require.Empty(t, ee.Stderr) // !!!!! probably not what one would expect if we only redirect stdout
|
||||
}
|
||||
@@ -0,0 +1,73 @@
|
||||
package zfscmd
|
||||
|
||||
import (
|
||||
"time"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
)
|
||||
|
||||
var metrics struct {
|
||||
totaltime *prometheus.HistogramVec
|
||||
systemtime *prometheus.HistogramVec
|
||||
usertime *prometheus.HistogramVec
|
||||
}
|
||||
|
||||
var timeLabels = []string{"jobid", "zfsbinary", "zfsverb"}
|
||||
var timeBuckets = []float64{0.01, 0.1, 0.2, 0.5, 0.75, 1, 2, 5, 10, 60}
|
||||
|
||||
func init() {
|
||||
metrics.totaltime = prometheus.NewHistogramVec(prometheus.HistogramOpts{
|
||||
Namespace: "zrepl",
|
||||
Subsystem: "zfscmd",
|
||||
Name: "runtime",
|
||||
Help: "number of seconds that the command took from start until wait returned",
|
||||
Buckets: timeBuckets,
|
||||
}, timeLabels)
|
||||
metrics.systemtime = prometheus.NewHistogramVec(prometheus.HistogramOpts{
|
||||
Namespace: "zrepl",
|
||||
Subsystem: "zfscmd",
|
||||
Name: "systemtime",
|
||||
Help: "https://golang.org/pkg/os/#ProcessState.SystemTime",
|
||||
Buckets: timeBuckets,
|
||||
}, timeLabels)
|
||||
metrics.usertime = prometheus.NewHistogramVec(prometheus.HistogramOpts{
|
||||
Namespace: "zrepl",
|
||||
Subsystem: "zfscmd",
|
||||
Name: "usertime",
|
||||
Help: "https://golang.org/pkg/os/#ProcessState.UserTime",
|
||||
Buckets: timeBuckets,
|
||||
}, timeLabels)
|
||||
|
||||
}
|
||||
|
||||
func RegisterMetrics(r prometheus.Registerer) {
|
||||
r.MustRegister(metrics.totaltime)
|
||||
r.MustRegister(metrics.systemtime)
|
||||
r.MustRegister(metrics.usertime)
|
||||
}
|
||||
|
||||
func waitPostPrometheus(c *Cmd, err error, now time.Time) {
|
||||
|
||||
if len(c.cmd.Args) < 2 {
|
||||
getLogger(c.ctx).WithField("args", c.cmd.Args).
|
||||
Warn("prometheus: cannot turn zfs command into metric")
|
||||
return
|
||||
}
|
||||
|
||||
// Note: do not start parsing other aspects
|
||||
// of the ZFS command line. This is not the suitable layer
|
||||
// for such a task.
|
||||
|
||||
jobid := getJobIDOrDefault(c.ctx, "_nojobid")
|
||||
|
||||
labelValues := []string{jobid, c.cmd.Args[0], c.cmd.Args[1]}
|
||||
|
||||
metrics.totaltime.
|
||||
WithLabelValues(labelValues...).
|
||||
Observe(c.Runtime().Seconds())
|
||||
metrics.systemtime.WithLabelValues(labelValues...).
|
||||
Observe(c.cmd.ProcessState.SystemTime().Seconds())
|
||||
metrics.usertime.WithLabelValues(labelValues...).
|
||||
Observe(c.cmd.ProcessState.UserTime().Seconds())
|
||||
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
package zfscmd
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
type Report struct {
|
||||
Active []ActiveCommand
|
||||
}
|
||||
|
||||
type ActiveCommand struct {
|
||||
Path string
|
||||
Args []string
|
||||
StartedAt time.Time
|
||||
}
|
||||
|
||||
func GetReport() *Report {
|
||||
active.mtx.RLock()
|
||||
defer active.mtx.RUnlock()
|
||||
var activeCommands []ActiveCommand
|
||||
for c := range active.cmds {
|
||||
c.mtx.RLock()
|
||||
activeCommands = append(activeCommands, ActiveCommand{
|
||||
Path: c.cmd.Path,
|
||||
Args: c.cmd.Args,
|
||||
StartedAt: c.startedAt,
|
||||
})
|
||||
c.mtx.RUnlock()
|
||||
}
|
||||
return &Report{
|
||||
Active: activeCommands,
|
||||
}
|
||||
}
|
||||
|
||||
var active struct {
|
||||
mtx sync.RWMutex
|
||||
cmds map[*Cmd]bool
|
||||
}
|
||||
|
||||
func init() {
|
||||
active.cmds = make(map[*Cmd]bool)
|
||||
}
|
||||
|
||||
func startPostReport(c *Cmd, err error, now time.Time) {
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
active.mtx.Lock()
|
||||
prev := active.cmds[c]
|
||||
if prev {
|
||||
panic("impl error: duplicate active command")
|
||||
}
|
||||
active.cmds[c] = true
|
||||
active.mtx.Unlock()
|
||||
}
|
||||
|
||||
func waitPostReport(c *Cmd, now time.Time) {
|
||||
active.mtx.Lock()
|
||||
defer active.mtx.Unlock()
|
||||
prev := active.cmds[c]
|
||||
if !prev {
|
||||
panic(fmt.Sprintf("impl error: onWaitDone must only be called on an active command: %s", c))
|
||||
}
|
||||
delete(active.cmds, c)
|
||||
}
|
||||
Reference in New Issue
Block a user