package mdm
import (
"context"
crand "crypto/rand"
"fmt"
"math/rand"
"strconv"
"sync"
"time"
"github.com/google/trillian"
"github.com/google/trillian/client"
"github.com/google/trillian/monitoring"
"github.com/google/trillian/types"
"k8s.io/klog/v2"
)
type MergeDelayMonitor struct {
client []*client.LogClient
logID int64
opts MergeDelayOptions
}
type MergeDelayOptions struct {
ParallelAdds int
LeafSize int
NewLeafChance int
EmitInterval time.Duration
Deadline time.Duration
MinMergeDelay time.Duration
MetricFactory monitoring.MetricFactory
}
func NewMonitor(ctx context.Context, logID int64, cl trillian.TrillianLogClient, adminCl trillian.TrillianAdminClient, opts MergeDelayOptions) (*MergeDelayMonitor, error) {
if opts.MetricFactory == nil {
opts.MetricFactory = monitoring.InertMetricFactory{}
}
if opts.EmitInterval <= 0 {
opts.EmitInterval = 10 * time.Second
}
if opts.Deadline <= 0 {
opts.Deadline = 60 * time.Second
}
metricsOnce.Do(func() { initMetrics(opts.MetricFactory) })
tree, err := adminCl.GetTree(ctx, &trillian.GetTreeRequest{TreeId: logID})
if err != nil {
return nil, fmt.Errorf("failed to get tree %d: %v", logID, err)
}
verifier, err := client.NewLogVerifierFromTree(tree)
if err != nil {
return nil, fmt.Errorf("failed to build verifier: %v", err)
}
clients := make([]*client.LogClient, opts.ParallelAdds)
for i := 0; i < opts.ParallelAdds; i++ {
clients[i] = client.New(logID, cl, verifier, types.LogRootV1{})
clients[i].MinMergeDelay = opts.MinMergeDelay
}
mdm := MergeDelayMonitor{logID: logID, opts: opts, client: clients}
return &mdm, nil
}
func (m *MergeDelayMonitor) Monitor(ctx context.Context) error {
klog.Infof("starting %d parallel monitor instances", m.opts.ParallelAdds)
errs := make(chan error, m.opts.ParallelAdds)
var wg sync.WaitGroup
wg.Add(m.opts.ParallelAdds)
for i := 0; i < m.opts.ParallelAdds; i++ {
go func(i int) {
defer wg.Done()
if err := m.monitor(ctx, i); err != nil {
errs <- err
}
}(i)
}
ticker := time.NewTicker(m.opts.EmitInterval)
defer ticker.Stop()
klog.V(1).Infof("start stats ticker every %v", m.opts.EmitInterval)
go func(c <-chan time.Time) {
for range c {
countT, totalT := m.Stats(true)
if countT > 0 {
klog.Infof("new leaves: %d in %f secs, average %f secs", countT, totalT, totalT/float64(countT))
}
countF, totalF := m.Stats(false)
if countF > 0 {
klog.Infof("dup leaves: %d in %f secs, average %f secs", countF, totalF, totalF/float64(countF))
}
}
}(ticker.C)
wg.Wait()
close(errs)
var lastErr error
for err := range errs {
klog.Errorf("monitor failure: %v", err)
lastErr = err
}
return lastErr
}
func (m *MergeDelayMonitor) monitor(ctx context.Context, idx int) error {
logIDLabel := strconv.FormatInt(m.logID, 10)
data := make([]byte, m.opts.LeafSize)
createNew := true
for {
if rand.Intn(100) < m.opts.NewLeafChance {
createNew = true
}
if createNew {
if _, err := crand.Read(data); err != nil {
return fmt.Errorf("Read(): %v", err)
}
}
klog.V(1).Infof("[%d] submit new=%t leaf and wait for inclusion (within %v)", idx, createNew, m.opts.Deadline)
cctx, cancel := context.WithTimeout(ctx, m.opts.Deadline)
start := time.Now()
err := m.client[idx].AddLeaf(cctx, data)
cancel()
if err != nil {
return fmt.Errorf("failed to QueueLeaf: %v", err)
}
mergeDelay := time.Since(start)
mergeDelayDist.Observe(mergeDelay.Seconds(), logIDLabel, newLeafLabel[createNew])
klog.V(1).Infof("[%d] merge delay for new=%t leaf = %v", idx, createNew, mergeDelay)
select {
case <-ctx.Done():
return nil
default:
}
createNew = false
}
}
func (m *MergeDelayMonitor) Stats(newLeaf bool) (uint64, float64) {
logIDLabel := strconv.FormatInt(m.logID, 10)
return mergeDelayDist.Info(logIDLabel, newLeafLabel[newLeaf])
}
var (
metricsOnce sync.Once
mergeDelayDist monitoring.Histogram
newLeafLabel = map[bool]string{true: "true", false: "false"}
)
func initMetrics(mf monitoring.MetricFactory) {
mergeDelayDist = mf.NewHistogram("merge_delay", "Merge delay for submitted leaves", "log_id", "new_leaf")
}