package log
import (
"context"
"fmt"
"reflect"
"sort"
"strconv"
"sync"
"time"
"github.com/google/trillian/extension"
"github.com/google/trillian/monitoring"
"github.com/google/trillian/storage"
"github.com/google/trillian/util/clock"
"github.com/google/trillian/util/election"
"golang.org/x/sync/semaphore"
"k8s.io/klog/v2"
)
var (
DefaultTimeout = 60 * time.Second
once sync.Once
knownLogs monitoring.Gauge
resignations monitoring.Counter
isMaster monitoring.Gauge
signingRuns monitoring.Counter
failedSigningRuns monitoring.Counter
entriesAdded monitoring.Counter
batchesAdded monitoring.Counter
)
func createMetrics(mf monitoring.MetricFactory) {
if mf == nil {
mf = monitoring.InertMetricFactory{}
}
knownLogs = mf.NewGauge("known_logs", "Set to 1 for known logs (whether this instance is master or not)", logIDLabel)
resignations = mf.NewCounter("master_resignations", "Number of mastership resignations", logIDLabel)
isMaster = mf.NewGauge("is_master", "Whether this instance is master (0/1)", logIDLabel)
signingRuns = mf.NewCounter("signing_runs", "Number of times a signing run has succeeded", logIDLabel)
failedSigningRuns = mf.NewCounter("failed_signing_runs", "Number of times a signing run has failed", logIDLabel)
entriesAdded = mf.NewCounter("entries_added", "Number of entries added to the log", logIDLabel)
batchesAdded = mf.NewCounter("batches_added", "Number of times a non zero number of entries was added", logIDLabel)
}
type Operation interface {
ExecutePass(ctx context.Context, logID int64, info *OperationInfo) (int, error)
}
type OperationInfo struct {
Registry extension.Registry
BatchSize int
TimeSource clock.TimeSource
ElectionConfig election.RunnerConfig
RunInterval time.Duration
NumWorkers int
Timeout time.Duration
}
type OperationManager struct {
info OperationInfo
logOperation Operation
runnerWG sync.WaitGroup
runnerCancels map[string]context.CancelFunc
pendingResignations chan election.Resignation
tracker *election.MasterTracker
logNames map[int64]string
lastHeld []int64
idsMutex sync.Mutex
}
func NewOperationManager(info OperationInfo, logOperation Operation) *OperationManager {
once.Do(func() {
createMetrics(info.Registry.MetricFactory)
})
if info.Timeout == 0 {
info.Timeout = DefaultTimeout
}
tracker := election.NewMasterTracker(nil, func(id string, v bool) {
val := 0.0
if v {
val = 1.0
}
isMaster.Set(val, id)
})
return &OperationManager{
info: info,
logOperation: logOperation,
runnerCancels: make(map[string]context.CancelFunc),
pendingResignations: make(chan election.Resignation, 100),
tracker: tracker,
logNames: make(map[int64]string),
}
}
func (o *OperationManager) logName(ctx context.Context, logID int64) string {
o.idsMutex.Lock()
defer o.idsMutex.Unlock()
if name, ok := o.logNames[logID]; ok {
return name
}
tree, err := storage.GetTree(ctx, o.info.Registry.AdminStorage, logID)
if err != nil {
klog.Errorf("%v: failed to get log info: %v", logID, err)
return "<err>"
}
name := tree.DisplayName
if name == "" {
name = fmt.Sprintf("<log-%d>", logID)
}
o.logNames[logID] = name
return o.logNames[logID]
}
func (o *OperationManager) heldInfo(ctx context.Context, logIDs []int64) string {
names := make([]string, 0, len(logIDs))
for _, logID := range logIDs {
names = append(names, o.logName(ctx, logID))
}
sort.Strings(names)
result := "master for:"
for _, name := range names {
result += " " + name
}
return result
}
func (o *OperationManager) masterFor(ctx context.Context, allIDs []int64) ([]int64, error) {
if o.info.Registry.ElectionFactory == nil {
return allIDs, nil
}
allStringIDs := make([]string, 0, len(allIDs))
for _, id := range allIDs {
s := strconv.FormatInt(id, 10)
allStringIDs = append(allStringIDs, s)
}
for _, logID := range allStringIDs {
knownLogs.Set(1, logID)
if o.runnerCancels[logID] == nil {
o.tracker.Set(logID, false)
o.runnerCancels[logID] = o.runElectionWithRestarts(ctx, logID)
}
}
held := o.tracker.Held()
heldIDs := make([]int64, 0, len(allIDs))
sort.Strings(allStringIDs)
for _, s := range held {
if i := sort.SearchStrings(allStringIDs, s); i >= len(allStringIDs) || allStringIDs[i] != s {
continue
}
id, err := strconv.ParseInt(s, 10, 64)
if err != nil {
return nil, fmt.Errorf("failed to parse logID %v as int64", s)
}
heldIDs = append(heldIDs, id)
}
return heldIDs, nil
}
func (o *OperationManager) runElectionWithRestarts(ctx context.Context, logID string) context.CancelFunc {
klog.Infof("create master election goroutine for %v", logID)
cctx, cancel := context.WithCancel(ctx)
run := func(ctx context.Context) {
e, err := o.info.Registry.ElectionFactory.NewElection(ctx, logID)
if err != nil {
klog.Errorf("failed to create election for %v: %v", logID, err)
return
}
config := o.info.ElectionConfig
r := election.NewRunner(logID, &config, o.tracker, cancel, e)
r.Run(ctx, o.pendingResignations)
}
o.runnerWG.Add(1)
go func(ctx context.Context) {
defer o.runnerWG.Done()
for ctx.Err() == nil {
run(ctx)
const pause = time.Duration(5 * time.Second)
if err := clock.SleepSource(ctx, pause, o.info.TimeSource); err != nil {
break
}
}
}(cctx)
return cancel
}
func (o *OperationManager) updateHeldIDs(ctx context.Context, logIDs, activeIDs []int64) {
heldInfo := o.heldInfo(ctx, logIDs)
msg := fmt.Sprintf("Acting as master for %d / %d active logs: %s", len(logIDs), len(activeIDs), heldInfo)
o.idsMutex.Lock()
defer o.idsMutex.Unlock()
if !reflect.DeepEqual(logIDs, o.lastHeld) {
o.lastHeld = make([]int64, len(logIDs))
copy(o.lastHeld, logIDs)
klog.Info(msg)
if o.info.Registry.SetProcessStatus != nil {
o.info.Registry.SetProcessStatus(heldInfo)
}
} else {
klog.V(1).Info(msg)
}
}
func (o *OperationManager) getLogsAndExecutePass(ctx context.Context) error {
runCtx, cancel := context.WithTimeout(ctx, o.info.Timeout)
defer cancel()
activeIDs, err := o.info.Registry.LogStorage.GetActiveLogIDs(ctx)
if err != nil {
return fmt.Errorf("failed to list active log IDs: %v", err)
}
logIDs, err := o.masterFor(ctx, activeIDs)
if err != nil {
return fmt.Errorf("failed to determine log IDs we're master for: %v", err)
}
o.updateHeldIDs(ctx, logIDs, activeIDs)
executePassForAll(runCtx, &o.info, o.logOperation, logIDs)
return nil
}
func (o *OperationManager) OperationSingle(ctx context.Context) {
if err := o.getLogsAndExecutePass(ctx); err != nil {
klog.Errorf("failed to perform operation: %v", err)
}
}
func (o *OperationManager) OperationLoop(ctx context.Context) {
klog.Infof("Log operation manager starting")
for {
if err := o.operateOnce(ctx); err != nil {
klog.Infof("Log operation manager shutting down")
break
}
}
for logID, cancel := range o.runnerCancels {
if cancel != nil {
klog.V(1).Infof("cancel election runner for %s", logID)
cancel()
}
}
close(o.pendingResignations)
for r := range o.pendingResignations {
resignations.Inc(r.ID)
r.Execute(ctx)
}
klog.Infof("wait for termination of election runners...")
o.runnerWG.Wait()
klog.Infof("wait for termination of election runners...done")
}
func (o *OperationManager) operateOnce(ctx context.Context) error {
start := o.info.TimeSource.Now()
if err := o.getLogsAndExecutePass(ctx); err != nil {
if ctx.Err() != nil {
klog.Errorf("failed to execute operation on logs: %v", err)
}
}
klog.V(1).Infof("Log operation manager pass complete")
doneResigning := false
for !doneResigning {
select {
case r := <-o.pendingResignations:
resignations.Inc(r.ID)
r.Execute(ctx)
default:
doneResigning = true
}
}
select {
case <-ctx.Done():
return ctx.Err()
default:
}
duration := o.info.TimeSource.Now().Sub(start)
wait := o.info.RunInterval - duration
if wait > 0 {
klog.V(1).Infof("Processing started at %v for %v; wait %v before next run", start, duration, wait)
if err := clock.SleepContext(ctx, wait); err != nil {
return err
}
} else {
klog.V(1).Infof("Processing started at %v for %v; start next run immediately", start, duration)
}
return nil
}
func executePassForAll(ctx context.Context, info *OperationInfo, op Operation, logIDs []int64) {
startBatch := info.TimeSource.Now()
numWorkers := info.NumWorkers
if numWorkers <= 0 {
klog.Warning("Running executor with NumWorkers <= 0, assuming 1")
numWorkers = 1
}
klog.V(1).Infof("Running executor with %d worker(s)", numWorkers)
sem := semaphore.NewWeighted(int64(numWorkers))
var wg sync.WaitGroup
for _, logID := range logIDs {
if err := sem.Acquire(ctx, 1); err != nil {
break
}
wg.Add(1)
go func(logID int64) {
defer wg.Done()
defer sem.Release(1)
if err := executePass(ctx, info, op, logID); err != nil {
klog.Errorf("ExecutePass(%v) failed: %v", logID, err)
}
}(logID)
}
wg.Wait()
d := clock.SecondsSince(info.TimeSource, startBatch)
klog.V(1).Infof("Group run completed in %.2f seconds", d)
}
func executePass(ctx context.Context, info *OperationInfo, op Operation, logID int64) error {
label := strconv.FormatInt(logID, 10)
start := info.TimeSource.Now()
count, err := op.ExecutePass(ctx, logID, info)
if err != nil {
failedSigningRuns.Inc(label)
return err
}
signingRuns.Inc(label)
if count > 0 {
d := clock.SecondsSince(info.TimeSource, start)
klog.Infof("%v: processed %d items in %.2f seconds (%.2f qps)", logID, count, d, float64(count)/d)
entriesAdded.Add(float64(count), label)
batchesAdded.Inc(label)
} else {
klog.V(1).Infof("%v: no items to process", logID)
}
return nil
}