package log
import (
"bytes"
"context"
"flag"
"fmt"
"strconv"
"sync"
"time"
"github.com/google/trillian"
"github.com/google/trillian/monitoring"
"github.com/google/trillian/quota"
"github.com/google/trillian/storage"
"github.com/google/trillian/storage/tree"
"github.com/google/trillian/types"
"github.com/google/trillian/util/clock"
"github.com/transparency-dev/merkle/compact"
"github.com/transparency-dev/merkle/rfc6962"
"google.golang.org/protobuf/types/known/timestamppb"
"k8s.io/klog/v2"
)
const logIDLabel = "logid"
var (
sequencerOnce sync.Once
seqBatches monitoring.Counter
seqTreeSize monitoring.Gauge
seqLatency monitoring.Histogram
seqDequeueLatency monitoring.Histogram
seqGetRootLatency monitoring.Histogram
seqInitTreeLatency monitoring.Histogram
seqWriteTreeLatency monitoring.Histogram
seqUpdateLeavesLatency monitoring.Histogram
seqSetNodesLatency monitoring.Histogram
seqStoreRootLatency monitoring.Histogram
seqCounter monitoring.Counter
seqMergeDelay monitoring.Histogram
seqTimestamp monitoring.Gauge
QuotaIncreaseFactor = 1.1
)
var _ = flag.String("tree_ids_with_no_ephemeral_nodes", "*", "[Deprecated] Comma-separated list of tree IDs for which storing the ephemeral nodes is disabled, or * to disable it for all trees")
func quotaIncreaseFactor() float64 {
if QuotaIncreaseFactor < 1 {
QuotaIncreaseFactor = 1
return 1
}
return QuotaIncreaseFactor
}
func InitMetrics(mf monitoring.MetricFactory) {
sequencerOnce.Do(func() {
if mf == nil {
mf = monitoring.InertMetricFactory{}
}
quota.InitMetrics(mf)
seqBatches = mf.NewCounter("sequencer_batches", "Number of sequencer batch operations", logIDLabel)
seqTreeSize = mf.NewGauge("sequencer_tree_size", "Tree size of last SLR signed", logIDLabel)
seqTimestamp = mf.NewGauge("sequencer_tree_timestamp", "Time of last SLR signed in ms since epoch", logIDLabel)
seqLatency = mf.NewHistogram("sequencer_latency", "Latency of sequencer batch operation in seconds", logIDLabel)
seqDequeueLatency = mf.NewHistogram("sequencer_latency_dequeue", "Latency of dequeue-leaves part of sequencer batch operation in seconds", logIDLabel)
seqGetRootLatency = mf.NewHistogram("sequencer_latency_get_root", "Latency of get-root part of sequencer batch operation in seconds", logIDLabel)
seqInitTreeLatency = mf.NewHistogram("sequencer_latency_init_tree", "Latency of init-tree part of sequencer batch operation in seconds", logIDLabel)
seqWriteTreeLatency = mf.NewHistogram("sequencer_latency_write_tree", "Latency of write-tree part of sequencer batch operation in seconds", logIDLabel)
seqUpdateLeavesLatency = mf.NewHistogram("sequencer_latency_update_leaves", "Latency of update-leaves part of sequencer batch operation in seconds", logIDLabel)
seqSetNodesLatency = mf.NewHistogram("sequencer_latency_set_nodes", "Latency of set-nodes part of sequencer batch operation in seconds", logIDLabel)
seqStoreRootLatency = mf.NewHistogram("sequencer_latency_store_root", "Latency of store-root part of sequencer batch operation in seconds", logIDLabel)
seqCounter = mf.NewCounter("sequencer_sequenced", "Number of leaves sequenced", logIDLabel)
seqMergeDelay = mf.NewHistogram("sequencer_merge_delay", "Delay between queuing and integration of leaves", logIDLabel)
})
}
func initCompactRangeFromStorage(ctx context.Context, root *types.LogRootV1, tx storage.LogTreeTX) (*compact.Range, error) {
fact := compact.RangeFactory{Hash: rfc6962.DefaultHasher.HashChildren}
if root.TreeSize == 0 {
return fact.NewEmptyRange(0), nil
}
ids := compact.RangeNodes(0, root.TreeSize, nil)
nodes, err := tx.GetMerkleNodes(ctx, ids)
if err != nil {
return nil, fmt.Errorf("failed to read tree nodes: %v", err)
}
if got, want := len(nodes), len(ids); got != want {
return nil, fmt.Errorf("failed to get %d nodes, got %d", want, got)
}
hashes := make([][]byte, len(nodes))
for i, node := range nodes {
hashes[i] = node.Hash
}
cr, err := fact.NewRange(0, root.TreeSize, hashes)
if err != nil {
return nil, fmt.Errorf("failed to create compact.Range: %v", err)
}
hash, err := cr.GetRootHash(nil)
if err != nil {
return nil, fmt.Errorf("failed to compute the root hash: %v", err)
}
if want := root.RootHash; !bytes.Equal(hash, want) {
return nil, fmt.Errorf("root hash mismatch: got %x, want %x", hash, want)
}
return cr, nil
}
func buildNodesFromNodeMap(nodeMap map[compact.NodeID][]byte) []tree.Node {
nodes := make([]tree.Node, 0, len(nodeMap))
for id, hash := range nodeMap {
nodes = append(nodes, tree.Node{ID: id, Hash: hash})
}
return nodes
}
func prepareLeaves(leaves []*trillian.LogLeaf, begin uint64, label string, timeSource clock.TimeSource) error {
now := timeSource.Now()
integrateAt := timestamppb.New(now)
if err := integrateAt.CheckValid(); err != nil {
return fmt.Errorf("got invalid integrate timestamp: %w", err)
}
for i, leaf := range leaves {
if got, want := leaf.LeafIndex, begin+uint64(i); got < 0 || got != int64(want) {
return fmt.Errorf("got invalid leaf index: %v, want: %v", got, want)
}
leaf.IntegrateTimestamp = integrateAt
if leaf.QueueTimestamp != nil && leaf.QueueTimestamp.Seconds != 0 {
if err := leaf.QueueTimestamp.CheckValid(); err != nil {
return fmt.Errorf("got invalid queue timestamp: %w", err)
}
queueTS := leaf.QueueTimestamp.AsTime()
mergeDelay := now.Sub(queueTS)
seqMergeDelay.Observe(mergeDelay.Seconds(), label)
}
}
return nil
}
func updateCompactRange(cr *compact.Range, leaves []*trillian.LogLeaf, label string) (map[compact.NodeID][]byte, []byte, error) {
nodeMap := make(map[compact.NodeID][]byte)
store := func(id compact.NodeID, hash []byte) { nodeMap[id] = hash }
for _, leaf := range leaves {
idx := leaf.LeafIndex
if size := cr.End(); idx < 0 || idx != int64(size) {
return nil, nil, fmt.Errorf("leaf index mismatch: got %d, want %d", idx, size)
}
if err := cr.Append(leaf.MerkleLeafHash, store); err != nil {
return nil, nil, err
}
}
hash, err := cr.GetRootHash(nil)
if err != nil {
return nil, nil, err
}
return nodeMap, hash, nil
}
type sequencingTask interface {
fetch(ctx context.Context, limit int, cutoff time.Time) ([]*trillian.LogLeaf, error)
update(ctx context.Context, leaves []*trillian.LogLeaf) error
}
type sequencingTaskData struct {
label string
treeSize uint64
timeSource clock.TimeSource
tx storage.LogTreeTX
}
type logSequencingTask sequencingTaskData
func (s *logSequencingTask) fetch(ctx context.Context, limit int, cutoff time.Time) ([]*trillian.LogLeaf, error) {
start := s.timeSource.Now()
leaves, err := s.tx.DequeueLeaves(ctx, limit, cutoff)
if err != nil {
return nil, fmt.Errorf("%v: Sequencer failed to dequeue leaves: %v", s.label, err)
}
seqDequeueLatency.Observe(clock.SecondsSince(s.timeSource, start), s.label)
for i, leaf := range leaves {
leaf.LeafIndex = int64(s.treeSize + uint64(i))
if got := leaf.LeafIndex; got < 0 {
return nil, fmt.Errorf("%v: leaf index overflow: %d", s.label, got)
}
}
return leaves, nil
}
func (s *logSequencingTask) update(ctx context.Context, leaves []*trillian.LogLeaf) error {
start := s.timeSource.Now()
if err := s.tx.UpdateSequencedLeaves(ctx, leaves); err != nil {
return fmt.Errorf("%v: Sequencer failed to update sequenced leaves: %v", s.label, err)
}
seqUpdateLeavesLatency.Observe(clock.SecondsSince(s.timeSource, start), s.label)
return nil
}
type preorderedLogSequencingTask sequencingTaskData
func (s *preorderedLogSequencingTask) fetch(ctx context.Context, limit int, cutoff time.Time) ([]*trillian.LogLeaf, error) {
start := s.timeSource.Now()
leaves, err := s.tx.DequeueLeaves(ctx, limit, cutoff)
if err != nil {
return nil, fmt.Errorf("%v: Sequencer failed to load sequenced leaves: %v", s.label, err)
}
seqDequeueLatency.Observe(clock.SecondsSince(s.timeSource, start), s.label)
return leaves, nil
}
func (s *preorderedLogSequencingTask) update(ctx context.Context, leaves []*trillian.LogLeaf) error {
return nil
}
func IntegrateBatch(ctx context.Context, tree *trillian.Tree, limit int, guardWindow, maxRootDurationInterval time.Duration, ts clock.TimeSource, ls storage.LogStorage, qm quota.Manager) (int, error) {
start := ts.Now()
label := strconv.FormatInt(tree.TreeId, 10)
numLeaves := 0
var newLogRoot *types.LogRootV1
var newSLR *trillian.SignedLogRoot
err := ls.ReadWriteTransaction(ctx, tree, func(ctx context.Context, tx storage.LogTreeTX) error {
stageStart := ts.Now()
defer seqBatches.Inc(label)
defer func() { seqLatency.Observe(clock.SecondsSince(ts, start), label) }()
sth, err := tx.LatestSignedLogRoot(ctx)
if err != nil || sth == nil {
return fmt.Errorf("%v: Sequencer failed to get latest root: %v", tree.TreeId, err)
}
var currentRoot types.LogRootV1
if err := currentRoot.UnmarshalBinary(sth.LogRoot); err != nil {
return fmt.Errorf("%v: Sequencer failed to unmarshal latest root: %v", tree.TreeId, err)
}
seqGetRootLatency.Observe(clock.SecondsSince(ts, stageStart), label)
seqTreeSize.Set(float64(currentRoot.TreeSize), label)
if currentRoot.RootHash == nil {
klog.Warningf("%v: Fresh log - no previous TreeHeads exist.", tree.TreeId)
return storage.ErrTreeNeedsInit
}
taskData := &sequencingTaskData{
label: label,
treeSize: currentRoot.TreeSize,
timeSource: ts,
tx: tx,
}
var st sequencingTask
switch tree.TreeType {
case trillian.TreeType_LOG:
st = (*logSequencingTask)(taskData)
case trillian.TreeType_PREORDERED_LOG:
st = (*preorderedLogSequencingTask)(taskData)
default:
return fmt.Errorf("IntegrateBatch not supported for TreeType %v", tree.TreeType)
}
sequencedLeaves, err := st.fetch(ctx, limit, start.Add(-guardWindow))
if err != nil {
return fmt.Errorf("%v: Sequencer failed to load sequenced batch: %v", tree.TreeId, err)
}
numLeaves = len(sequencedLeaves)
if numLeaves == 0 {
nowNanos := ts.Now().UnixNano()
interval := time.Duration(nowNanos - int64(currentRoot.TimestampNanos))
if maxRootDurationInterval == 0 || interval < maxRootDurationInterval {
klog.V(1).Infof("%v: No leaves sequenced in this signing operation", tree.TreeId)
return nil
}
klog.Infof("%v: Force new root generation as %v since last root", tree.TreeId, interval)
}
stageStart = ts.Now()
cr, err := initCompactRangeFromStorage(ctx, ¤tRoot, tx)
if err != nil {
return fmt.Errorf("%v: compact range init failed: %v", tree.TreeId, err)
}
seqInitTreeLatency.Observe(clock.SecondsSince(ts, stageStart), label)
stageStart = ts.Now()
if err := prepareLeaves(sequencedLeaves, cr.End(), label, ts); err != nil {
return err
}
nodeMap, newRoot, err := updateCompactRange(cr, sequencedLeaves, label)
if err != nil {
return err
}
seqWriteTreeLatency.Observe(clock.SecondsSince(ts, stageStart), label)
if err := st.update(ctx, sequencedLeaves); err != nil {
return err
}
stageStart = ts.Now()
targetNodes := buildNodesFromNodeMap(nodeMap)
if err := tx.SetMerkleNodes(ctx, targetNodes); err != nil {
return fmt.Errorf("%v: Sequencer failed to set Merkle nodes: %v", tree.TreeId, err)
}
seqSetNodesLatency.Observe(clock.SecondsSince(ts, stageStart), label)
stageStart = ts.Now()
if cr.End() == 0 {
newRoot = rfc6962.DefaultHasher.EmptyRoot()
}
newLogRoot = &types.LogRootV1{
RootHash: newRoot,
TimestampNanos: uint64(ts.Now().UnixNano()),
TreeSize: cr.End(),
}
seqTreeSize.Set(float64(newLogRoot.TreeSize), label)
seqTimestamp.Set(float64(time.Duration(newLogRoot.TimestampNanos)*time.Nanosecond/
time.Millisecond), label)
if newLogRoot.TimestampNanos <= currentRoot.TimestampNanos {
return fmt.Errorf("%v: refusing to sign root with timestamp earlier than previous root (%d <= %d)", tree.TreeId, newLogRoot.TimestampNanos, currentRoot.TimestampNanos)
}
logRoot, err := newLogRoot.MarshalBinary()
if err != nil {
return fmt.Errorf("%v: signer failed to marshal root: %v", tree.TreeId, err)
}
newSLR := &trillian.SignedLogRoot{LogRoot: logRoot}
if err := tx.StoreSignedLogRoot(ctx, newSLR); err != nil {
return fmt.Errorf("%v: failed to write updated tree root: %v", tree.TreeId, err)
}
seqStoreRootLatency.Observe(clock.SecondsSince(ts, stageStart), label)
return nil
})
if err != nil {
return 0, err
}
replenishQuota(ctx, numLeaves, tree.TreeId, qm)
seqCounter.Add(float64(numLeaves), label)
if newSLR != nil {
klog.Infof("%v: sequenced %v leaves, size %v", tree.TreeId, numLeaves, newLogRoot.TreeSize)
}
return numLeaves, nil
}
func replenishQuota(ctx context.Context, numLeaves int, treeID int64, qm quota.Manager) {
if numLeaves > 0 {
tokens := int(float64(numLeaves) * quotaIncreaseFactor())
specs := []quota.Spec{
{Group: quota.Tree, Kind: quota.Read, TreeID: treeID},
{Group: quota.Tree, Kind: quota.Write, TreeID: treeID},
{Group: quota.Global, Kind: quota.Read},
{Group: quota.Global, Kind: quota.Write},
}
klog.V(2).Infof("%v: replenishing %d tokens (numLeaves = %d)", treeID, tokens, numLeaves)
err := qm.PutTokens(ctx, tokens, specs)
if err != nil {
klog.Warningf("%v: failed to replenish %d tokens: %v", treeID, tokens, err)
}
quota.Metrics.IncReplenished(tokens, specs, err == nil)
}
}