Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions cmd/livepeer/starter/flags.go
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,8 @@ func NewLivepeerConfig(fs *flag.FlagSet) LivepeerConfig {
cfg.MaxTotalEV = fs.String("maxTotalEV", *cfg.MaxTotalEV, "The maximum acceptable expected value for one PM payment")
// Broadcaster deposit multiplier to determine max acceptable ticket faceValue
cfg.DepositMultiplier = fs.Int("depositMultiplier", *cfg.DepositMultiplier, "The deposit multiplier used to determine max acceptable faceValue for PM tickets")
// Orchestrator automatic pruning of dead winning tickets
cfg.TicketPrune = fs.Bool("ticketPrune", *cfg.TicketPrune, "Enable automatic pruning of expired and permanently failed winning tickets from the ticket queue database")
// Orchestrator base pricing info
cfg.PricePerUnit = fs.String("pricePerUnit", "0", "The price per 'pixelsPerUnit' amount pixels. Can be specified in wei or a custom currency in the format <price><currency> (e.g. 0.50USD). When using a custom currency, a corresponding price feed must be configured with -priceFeedAddr")
// Unit of pixels for both O's pricePerUnit and B's maxPricePerUnit
Expand Down
4 changes: 4 additions & 0 deletions cmd/livepeer/starter/starter.go
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,7 @@ type LivepeerConfig struct {
MaxTicketEV *string
MaxTotalEV *string
DepositMultiplier *int
TicketPrune *bool
PricePerUnit *string
PixelsPerUnit *string
PriceFeedAddr *string
Expand Down Expand Up @@ -268,6 +269,7 @@ func DefaultLivepeerConfig() LivepeerConfig {
defaultMaxTicketEV := "3000000000000"
defaultMaxTotalEV := "20000000000000"
defaultDepositMultiplier := 1
defaultTicketPrune := true
defaultMaxPricePerUnit := "0"
defaultMaxPricePerCapability := ""
defaultIgnoreMaxPriceIfNeeded := false
Expand Down Expand Up @@ -393,6 +395,7 @@ func DefaultLivepeerConfig() LivepeerConfig {
MaxTicketEV: &defaultMaxTicketEV,
MaxTotalEV: &defaultMaxTotalEV,
DepositMultiplier: &defaultDepositMultiplier,
TicketPrune: &defaultTicketPrune,
MaxPricePerUnit: &defaultMaxPricePerUnit,
MaxPricePerCapability: &defaultMaxPricePerCapability,
IgnoreMaxPriceIfNeeded: &defaultIgnoreMaxPriceIfNeeded,
Expand Down Expand Up @@ -987,6 +990,7 @@ func StartLivepeer(ctx context.Context, cfg LivepeerConfig) {
RedeemGas: redeemGas,
SuggestGasPrice: client.Backend().SuggestGasPrice,
RPCTimeout: ethRPCTimeout,
TicketPrune: *cfg.TicketPrune,
}

if *cfg.Orchestrator {
Expand Down
26 changes: 26 additions & 0 deletions common/db.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ type DB struct {
winningTicketCount *sql.Stmt
markWinningTicketRedeemed *sql.Stmt
removeWinningTicket *sql.Stmt
clearDeadWinningTickets *sql.Stmt
insertMiniHeader *sql.Stmt
findLatestMiniHeader *sql.Stmt
findAllMiniHeadersSortedByNumber *sql.Stmt
Expand Down Expand Up @@ -324,6 +325,15 @@ func InitDB(dbPath string) (*DB, error) {
}
d.markWinningTicketRedeemed = stmt

// Clear dead tickets
stmt, err = db.Prepare("DELETE FROM ticketQueue WHERE (redeemedAt IS NULL AND txHash IS NULL AND creationRound < ?) OR (redeemedAt IS NOT NULL AND txHash=?)")
if err != nil {
glog.Error("Unable to prepare clearDeadWinningTickets ", err)
d.Close()
return nil, err
}
d.clearDeadWinningTickets = stmt

// Insert block header
stmt, err = db.Prepare("INSERT INTO blockheaders(number, parent, hash, logs) VALUES(?, ?, ?, ?)")
if err != nil {
Expand Down Expand Up @@ -405,6 +415,9 @@ func (db *DB) Close() {
if db.removeWinningTicket != nil {
db.removeWinningTicket.Close()
}
if db.clearDeadWinningTickets != nil {
db.clearDeadWinningTickets.Close()
}
if db.insertMiniHeader != nil {
db.insertMiniHeader.Close()
}
Expand Down Expand Up @@ -749,6 +762,19 @@ func (db *DB) RemoveWinningTicket(ticket *pm.SignedTicket) error {
return nil
}

// ClearDeadWinningTickets removes winning tickets that can never be redeemed: unredeemed
// tickets with a creationRound older than minCreationRound (outside the redemption validity
// window) and tickets marked redeemed with a zero txHash (permanently failed redemptions).
// Live tickets within the validity window and successfully redeemed tickets are kept.
// It returns the number of removed tickets.
func (db *DB) ClearDeadWinningTickets(minCreationRound int64) (int64, error) {
res, err := db.clearDeadWinningTickets.Exec(minCreationRound, ethcommon.Hash{}.Hex())
if err != nil {
return 0, errors.Wrapf(err, "failed clearing dead winning tickets minCreationRound=%v", minCreationRound)
}
return res.RowsAffected()
}

// SelectEarliestWinningTicket selects the earliest stored winning ticket for a 'sender' that is not expired and not yet redeemed
func (db *DB) SelectEarliestWinningTicket(sender ethcommon.Address, minCreationRound int64) (*pm.SignedTicket, error) {

Expand Down
77 changes: 77 additions & 0 deletions common/db_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1188,6 +1188,83 @@ func TestRemoveWinningTicket(t *testing.T) {
require.Equal(count, 0)
}

func TestClearDeadWinningTickets(t *testing.T) {
assert := assert.New(t)
dbh, dbraw, err := TempDB(t)
defer dbh.Close()
defer dbraw.Close()
require := require.New(t)
require.Nil(err)

minCreationRound := int64(100)

// live unredeemed ticket within the redemption validity window
_, ticket, sig, recipientRand := defaultWinningTicket(t)
ticket.CreationRound = minCreationRound
liveTicket := &pm.SignedTicket{
Ticket: ticket,
Sig: sig,
RecipientRand: recipientRand,
}
err = dbh.StoreWinningTicket(liveTicket)
require.Nil(err)

// expired unredeemed ticket outside the redemption validity window
_, ticket, sig, recipientRand = defaultWinningTicket(t)
ticket.CreationRound = minCreationRound - 1
expiredTicket := &pm.SignedTicket{
Ticket: ticket,
Sig: sig,
RecipientRand: recipientRand,
}
err = dbh.StoreWinningTicket(expiredTicket)
require.Nil(err)

// ticket whose redemption permanently failed (marked redeemed with a zero txHash)
_, ticket, sig, recipientRand = defaultWinningTicket(t)
ticket.CreationRound = minCreationRound
failedTicket := &pm.SignedTicket{
Ticket: ticket,
Sig: sig,
RecipientRand: recipientRand,
}
err = dbh.StoreWinningTicket(failedTicket)
require.Nil(err)
err = dbh.MarkWinningTicketRedeemed(failedTicket, ethcommon.Hash{})
require.Nil(err)

// successfully redeemed ticket (real txHash)
_, ticket, sig, recipientRand = defaultWinningTicket(t)
ticket.CreationRound = minCreationRound - 1
redeemedTicket := &pm.SignedTicket{
Ticket: ticket,
Sig: sig,
RecipientRand: recipientRand,
}
err = dbh.StoreWinningTicket(redeemedTicket)
require.Nil(err)
err = dbh.MarkWinningTicketRedeemed(redeemedTicket, pm.RandHash())
require.Nil(err)

require.Equal(4, getRowCountOrFatal("SELECT count(sig) FROM ticketQueue", dbraw, t))

// removes only the expired unredeemed and permanently failed tickets
count, err := dbh.ClearDeadWinningTickets(minCreationRound)
assert.Nil(err)
assert.Equal(int64(2), count)

assert.Equal(2, getRowCountOrFatal("SELECT count(sig) FROM ticketQueue", dbraw, t))
// the live unredeemed ticket is kept
assert.Equal(1, getRowCountOrFatal("SELECT count(sig) FROM ticketQueue WHERE redeemedAt IS NULL AND txHash IS NULL", dbraw, t))
// the successfully redeemed ticket is kept
assert.Equal(1, getRowCountOrFatal("SELECT count(sig) FROM ticketQueue WHERE redeemedAt IS NOT NULL", dbraw, t))

// clearing again removes nothing
count, err = dbh.ClearDeadWinningTickets(minCreationRound)
assert.Nil(err)
assert.Equal(int64(0), count)
}

func TestInsertMiniHeader_ReturnsFindLatestMiniHeader(t *testing.T) {
dbh, dbraw, err := TempDB(t)
defer dbh.Close()
Expand Down
4 changes: 3 additions & 1 deletion pm/queue.go
Original file line number Diff line number Diff line change
Expand Up @@ -168,7 +168,9 @@ func isNonRetryableTicketErr(err error) bool {
// Depends on logic in eth.client.CheckTx()
strings.Contains(err.Error(), "transaction failed") ||
// Arbitrum L2 happens to return zero as the L1 block hash which results in this non-retryable error
strings.Contains(err.Error(), "ticket creationRound does not have a block hash")
strings.Contains(err.Error(), "ticket creationRound does not have a block hash") ||
// Ticket is expired on-chain (past its redemption window); retrying can never succeed
strings.Contains(err.Error(), "ticket is expired")
}

func (q *ticketQueue) isRecipientActive(addr ethcommon.Address) bool {
Expand Down
12 changes: 11 additions & 1 deletion pm/queue_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -189,10 +189,20 @@ func TestTicketQueueLoop_IsNonRetryableTicketErr_MarkAsRedeemed(t *testing.T) {
consumeQueue(qc)
assert.True(ts.submitted[fmt.Sprintf("%x", ticket.Sig)])

// Test that ticket is not marked as redeemed if there is an error checking the tx, but the tx did not fail
// Test that ticket is marked as redeemed if it is expired on-chain
ticket = defaultSignedTicket(sender, 2)
addTicket(ticket)

qc = &queueConsumer{
redemptionErr: errors.New("execution reverted: ticket is expired"),
}
consumeQueue(qc)
assert.True(ts.submitted[fmt.Sprintf("%x", ticket.Sig)])

// Test that ticket is not marked as redeemed if there is an error checking the tx, but the tx did not fail
ticket = defaultSignedTicket(sender, 3)
addTicket(ticket)

qc = &queueConsumer{
redemptionErr: errors.New("some other error"),
}
Expand Down
39 changes: 39 additions & 0 deletions pm/sendermonitor.go
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,16 @@ type LocalSenderMonitorConfig struct {
RedeemGas int
SuggestGasPrice func(context.Context) (*big.Int, error)
RPCTimeout time.Duration

// TicketPrune enables automatic pruning of dead winning tickets
// (expired unredeemed or permanently failed) from the ticket store
TicketPrune bool
}

// ticketStoreCleaner describes a TicketStore implementation that is also
// capable of clearing dead winning tickets from persistent storage
type ticketStoreCleaner interface {
ClearDeadWinningTickets(minCreationRound int64) (int64, error)
}

type LocalSenderMonitor struct {
Expand Down Expand Up @@ -328,12 +338,41 @@ func (sm *LocalSenderMonitor) startCleanupLoop() {
select {
case <-ticker.C:
sm.cleanup()
if sm.cfg.TicketPrune {
sm.pruneDeadTickets()
}
case <-sm.quit:
return
}
}
}

// pruneDeadTickets removes dead winning tickets from the ticket store.
// A ticket is dead if it is unredeemed and outside of the redemption validity
// window (the redemption loop can never select it again) or if its redemption
// permanently failed. The cutoff matches the redemption window exactly
// (ticketValidityPeriod rounds behind the last initialized round): a ticket is
// pruned as soon as it falls out of the window, on the next cleanup tick,
// rather than a round later. The prune condition is the exact complement of the
// selection condition (creationRound >= LastInit-ticketValidityPeriod), so a
// still-redeemable ticket is never pruned.
func (sm *LocalSenderMonitor) pruneDeadTickets() {
cleaner, ok := sm.ticketStore.(ticketStoreCleaner)
if !ok {
return
}

minCreationRound := new(big.Int).Sub(sm.tm.LastInitializedRound(), big.NewInt(ticketValidityPeriod)).Int64()
count, err := cleaner.ClearDeadWinningTickets(minCreationRound)
if err != nil {
glog.Errorf("Unable to prune dead winning tickets err=%q", err)
return
}
if count > 0 {
glog.Infof("Pruned dead winning tickets count=%d", count)
}
}

// cleanup removes tracked remote senders that have exceeded
// their ttl
func (sm *LocalSenderMonitor) cleanup() {
Expand Down
Loading