circuitbreaker/peer_controller.go

522 lines
12 KiB
Go
Raw Permalink Normal View History

2022-11-29 15:18:13 +01:00
package main
import (
2022-11-29 17:44:55 +01:00
"container/list"
"context"
"time"
"github.com/lightningnetwork/lnd/lnwire"
2022-12-31 09:40:33 +01:00
"github.com/lightningnetwork/lnd/routing/route"
2023-01-03 13:01:26 +01:00
"github.com/paulbellamy/ratecounter"
2022-11-29 15:18:13 +01:00
"go.uber.org/zap"
"golang.org/x/time/rate"
)
2023-01-03 13:01:26 +01:00
type eventCounter struct {
fail *ratecounter.RateCounter
success *ratecounter.RateCounter
reject *ratecounter.RateCounter
}
type eventType int
const (
eventSuccess eventType = iota
eventFail
eventReject
)
func newEventCounter(interval time.Duration) *eventCounter {
return &eventCounter{
fail: ratecounter.NewRateCounter(interval),
success: ratecounter.NewRateCounter(interval),
reject: ratecounter.NewRateCounter(interval),
}
}
func (e *eventCounter) Incr(event eventType) {
switch event {
case eventSuccess:
e.success.Incr(1)
case eventFail:
e.fail.Incr(1)
case eventReject:
e.reject.Incr(1)
default:
panic("unknown event type")
}
}
func (e *eventCounter) Rates() (int64, int64, int64) {
return e.success.Rate(), e.fail.Rate(), e.reject.Rate()
}
2022-11-29 15:18:13 +01:00
type peerController struct {
2023-01-03 13:01:26 +01:00
cfg Limit
limiter *rate.Limiter
logger *zap.SugaredLogger
interceptChan chan peerInterceptEvent
resolvedChan chan peerResolvedEvent
2023-01-03 13:01:26 +01:00
updateLimitChan chan Limit
getStateChan chan chan *peerState
rateCounters []*eventCounter
2022-11-29 17:44:55 +01:00
htlcs map[circuitKey]*inFlightHtlc
2022-12-31 09:40:33 +01:00
lastChannelSync time.Time
pubKey route.Vertex
lnd lndclient
now func() time.Time
htlcCompleted func(context.Context, *HtlcInfo) error
2022-11-29 17:44:55 +01:00
}
type inFlightHtlc struct {
addedTs time.Time
incomingMsat lnwire.MilliSatoshi
outgoingMsat lnwire.MilliSatoshi
}
2022-11-29 17:44:55 +01:00
type peerInterceptEvent struct {
interceptEvent
peerInitiated bool
2022-11-29 15:18:13 +01:00
}
type peerResolvedEvent struct {
resolvedEvent
outgoingPeer *route.Vertex
}
2023-01-03 13:01:26 +01:00
type peerState struct {
counts []rateCounts
queueLen int64
pendingHtlcCount int64
}
type rateCounts struct {
success, fail, reject int64
}
var rateCounterIntervals = []time.Duration{time.Hour, 24 * time.Hour}
2022-12-31 09:40:33 +01:00
type peerControllerCfg struct {
logger *zap.SugaredLogger
limit Limit
burstSize int
htlcs map[circuitKey]*inFlightHtlc
lnd lndclient
pubKey route.Vertex
now func() time.Time
htlcCompleted func(context.Context, *HtlcInfo) error
2022-12-31 09:40:33 +01:00
}
func newPeerController(cfg *peerControllerCfg) *peerController {
2023-01-03 13:01:26 +01:00
logger := cfg.logger.With(
"peer", cfg.pubKey.String(),
)
2022-11-29 15:18:13 +01:00
// Skip if no interval set.
2023-01-03 13:01:26 +01:00
limiter := rate.NewLimiter(getRate(cfg.limit.MaxHourlyRate), cfg.burstSize)
2022-11-29 15:18:13 +01:00
logger.Infow("Peer controller initialized",
2023-01-03 13:01:26 +01:00
"maxHourlyRate", cfg.limit.MaxHourlyRate,
"maxPendingHtlcs", cfg.limit.MaxPending,
"mode", cfg.limit.Mode)
2022-11-29 17:44:55 +01:00
// Log initial pending htlcs.
2022-12-31 09:40:33 +01:00
for h := range cfg.htlcs {
2022-12-15 16:07:53 +00:00
logger.Infow("Initial pending htlc", "channel", h.channel, "htlc", h.htlc)
2022-11-29 17:44:55 +01:00
}
2022-11-29 15:18:13 +01:00
2023-01-03 13:01:26 +01:00
rateCounters := make([]*eventCounter, len(rateCounterIntervals))
for idx, interval := range rateCounterIntervals {
rateCounters[idx] = newEventCounter(interval)
}
2022-11-29 15:18:13 +01:00
return &peerController{
2023-01-03 13:01:26 +01:00
cfg: cfg.limit,
2022-12-31 09:40:33 +01:00
limiter: limiter,
logger: logger,
interceptChan: make(chan peerInterceptEvent),
resolvedChan: make(chan peerResolvedEvent),
2023-01-03 13:01:26 +01:00
updateLimitChan: make(chan Limit),
getStateChan: make(chan chan *peerState),
2022-12-31 09:40:33 +01:00
htlcs: cfg.htlcs,
2023-01-03 13:01:26 +01:00
rateCounters: rateCounters,
2022-12-31 09:40:33 +01:00
lnd: cfg.lnd,
pubKey: cfg.pubKey,
lastChannelSync: cfg.now(),
now: cfg.now,
htlcCompleted: cfg.htlcCompleted,
2022-11-29 15:18:13 +01:00
}
}
2023-01-03 13:01:26 +01:00
func (p *peerController) state(ctx context.Context) (*peerState, error) {
respChan := make(chan *peerState)
select {
case p.getStateChan <- respChan:
case <-ctx.Done():
return nil, ctx.Err()
}
select {
case state := <-respChan:
return state, nil
case <-ctx.Done():
return nil, ctx.Err()
}
}
func (p *peerController) rateInternal() []rateCounts {
allRateCounts := make([]rateCounts, len(p.rateCounters))
for idx, counter := range p.rateCounters {
success, fail, reject := counter.Rates()
allRateCounts[idx] = rateCounts{
success: success,
fail: fail,
reject: reject,
}
}
return allRateCounts
}
func (p *peerController) updateLimit(ctx context.Context, limit Limit) error {
select {
case p.updateLimitChan <- limit:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
2022-12-31 09:40:33 +01:00
func (p *peerController) newHtlcAllowed() bool {
2023-01-03 13:01:26 +01:00
return p.cfg.MaxPending == 0 ||
len(p.htlcs) < int(p.cfg.MaxPending)
2022-12-31 09:40:33 +01:00
}
func (p *peerController) syncPendingHtlcs(ctx context.Context) (bool, error) {
p.logger.Infow("Syncing pending htlcs")
2023-01-03 13:01:26 +01:00
allHtlcs, err := p.lnd.getPendingIncomingHtlcs(ctx, &p.pubKey)
2022-12-31 09:40:33 +01:00
if err != nil {
return false, err
}
2023-01-03 13:01:26 +01:00
htlcs := allHtlcs[p.pubKey]
p.lastChannelSync = p.now()
2022-12-31 09:40:33 +01:00
deletes := false
for key := range p.htlcs {
2023-01-03 13:01:26 +01:00
if htlcs != nil {
if _, ok := htlcs[key]; ok {
continue
}
2022-12-31 09:40:33 +01:00
}
// Htlc is no longer pending on incoming side. Must have missed an htlc
// event. Clear it from our list.
p.markHtlcComplete(ctx, key, nil)
2022-12-31 09:40:33 +01:00
logger := p.keyLogger(key)
logger.Infow("Cleaning up dangling htlc")
deletes = true
}
return deletes, nil
}
2022-11-29 17:44:55 +01:00
func (p *peerController) run(ctx context.Context) error {
queue := list.New()
2022-11-29 15:18:13 +01:00
2022-11-29 17:44:55 +01:00
var reservation *rate.Reservation
2022-11-29 15:18:13 +01:00
2022-11-29 17:44:55 +01:00
for {
// New htlcs are allowed when the number of pending htlcs is below the
// limit, or no limit has been set.
2022-12-31 09:40:33 +01:00
newHtlcAllowed := p.newHtlcAllowed()
// If no new htlcs are allowed and we've not synced recently, re-sync.
// Sometimes htlc events aren't broadcast by lnd, and this keeps our
// pending htlc count accurate.
if !newHtlcAllowed && time.Since(p.lastChannelSync) > time.Minute {
deletes, err := p.syncPendingHtlcs(ctx)
if err != nil {
return err
}
// When dangling htlcs are removed, re-evaluate whether a new htlc
// is allowed.
if deletes {
newHtlcAllowed = p.newHtlcAllowed()
}
}
2022-11-29 15:18:13 +01:00
2022-11-29 17:44:55 +01:00
// If an htlc can be forwarded, make a reservation on the rate limiter
// if it does not already exist.
if queue.Len() > 0 && newHtlcAllowed && reservation == nil {
reservation = p.limiter.Reserve()
}
2022-11-29 15:18:13 +01:00
2022-11-29 17:44:55 +01:00
// Create a delay channel based on the rate limiter delay. If there is
// no htlc to forward or the pending limit has been reached, use a nil
// channel to skip the select case.
var delayChan <-chan time.Time
if reservation != nil {
delayChan = time.After(reservation.Delay())
}
2022-11-29 15:18:13 +01:00
2022-11-29 17:44:55 +01:00
select {
// A new htlc is intercepted. Depending on the mode the controller is
// running in, the htlc will either be queued or handled immediately.
case event := <-p.interceptChan:
logger := p.keyLogger(event.circuitKey)
2022-11-29 15:18:13 +01:00
2022-12-31 09:40:33 +01:00
// Replays can happen when the htlcs map is initialized with a
// pending htlc on startup, and then a forward event happens for
2022-12-31 09:45:58 +01:00
// that htlc. For those htlcs, just resume.
2022-11-29 17:44:55 +01:00
_, ok := p.htlcs[event.circuitKey]
if ok {
2022-12-31 09:45:58 +01:00
if err := event.resume(true); err != nil {
return err
}
2022-11-29 17:44:55 +01:00
logger.Infow("Replay")
continue
}
2023-02-13 14:29:36 +01:00
mode := p.cfg.Mode
2022-12-31 09:54:40 +01:00
switch {
2023-02-13 14:28:12 +01:00
// Don't check limits in block mode and move onwards to failing the
// htlc.
case mode == ModeBlock:
logger.Infow("Htlc blocked")
2023-01-03 13:01:26 +01:00
// If there is a queue, then don't jump the queue.
case queue.Len() > 0:
2022-11-29 17:44:55 +01:00
2023-01-03 13:01:26 +01:00
// Check if new htlcs are allowed.
case !newHtlcAllowed:
logger.Infow("Pending htlc limit exceeded")
2022-11-29 17:44:55 +01:00
2023-02-13 14:29:36 +01:00
// Check the rate limit.
case !p.limiter.Allow():
logger.Infow("Rate limit exceeded")
// All signs green, forward the htlc.
default:
2022-12-31 09:54:40 +01:00
if err := p.forward(event.interceptEvent); err != nil {
2022-11-29 17:44:55 +01:00
return err
}
continue
}
2022-11-29 15:18:13 +01:00
2022-12-31 09:54:40 +01:00
// Queue if in one of the queue modes.
2023-01-03 13:01:26 +01:00
if mode == ModeQueue ||
(mode == ModeQueuePeerInitiated && event.peerInitiated) {
2022-11-29 15:18:13 +01:00
2022-12-31 09:54:40 +01:00
queue.PushFront(event)
logger.Infow("Queued", "queueLen", queue.Len())
2022-11-29 17:44:55 +01:00
continue
}
2022-12-31 09:54:40 +01:00
// Otherwise fail directly.
if err := event.resume(false); err != nil {
2022-11-29 17:44:55 +01:00
return err
}
2023-01-03 13:01:26 +01:00
p.incrCounter(eventReject)
2022-11-29 17:44:55 +01:00
// There are items in the queue, max pending htlcs has not yet been
// reached, and the rate limit delay has passed. Take the oldest item
// from the queue and forward it.
case <-delayChan:
listItem := queue.Back()
if listItem == nil {
panic("list empty")
}
queue.Remove(listItem)
event := listItem.Value.(peerInterceptEvent)
if err := p.forward(event.interceptEvent); err != nil {
return err
}
// Reservation has been used. Clear it so that a new reservation can
// be requested.
reservation = nil
// An htlc has been resolved in lnd. Remove it from the pending htlcs
// map to free up the slot for another htlc.
2023-01-03 13:01:26 +01:00
case resolvedEvent := <-p.resolvedChan:
key := resolvedEvent.incomingCircuitKey
2023-01-03 13:01:26 +01:00
2022-11-29 17:44:55 +01:00
_, ok := p.htlcs[key]
if !ok {
// Do not log here, because the event is still coming even for
// htlcs that were failed. We don't want to spam the log.
continue
}
p.markHtlcComplete(ctx, key, &resolvedEvent)
2022-11-29 17:44:55 +01:00
2023-01-03 13:01:26 +01:00
// Update rate counters.
if resolvedEvent.settled {
p.incrCounter(eventSuccess)
} else {
p.incrCounter(eventFail)
}
2022-11-29 17:44:55 +01:00
logger := p.keyLogger(key)
2023-01-03 13:01:26 +01:00
logger.Infow("Resolved htlc", "settled", resolvedEvent.settled,
"pending_htlcs", len(p.htlcs))
case limit := <-p.updateLimitChan:
p.logger.Infow("Updating peer controller", "limit", limit)
p.cfg = limit
p.limiter.SetLimit(getRate(limit.MaxHourlyRate))
case respChan := <-p.getStateChan:
counts := p.rateInternal()
select {
case respChan <- &peerState{
counts: counts,
queueLen: int64(queue.Len()),
pendingHtlcCount: int64(len(p.htlcs)),
}:
case <-ctx.Done():
return ctx.Err()
}
2022-11-29 17:44:55 +01:00
case <-ctx.Done():
return ctx.Err()
}
2022-11-29 15:18:13 +01:00
}
2022-11-29 17:44:55 +01:00
}
2022-11-29 15:18:13 +01:00
// markHtlcComplete removes the resolved htlc provided from the peerController's
// inFlight set and reports the completed htlc.
func (p *peerController) markHtlcComplete(ctx context.Context, key circuitKey,
resolution *peerResolvedEvent) {
// Lookup the HLTC to get its timestamp.
inFlight, ok := p.htlcs[key]
if !ok {
return
}
// Remove from our list of active HTLCs.
delete(p.htlcs, key)
// If no resolution is provided, we don't know the outcome of this HTLC (we
// re-synced and it was no longer present), so there is no further action to
// take.
if resolution == nil {
return
}
// If we couldn't look up an outgoing peer for the HTLC, either:
// 1. The outgoing channel never existed (since this is not validated on intercept)
// 2. The outgoing channel is pending close at time of resolution (edge case)
if resolution.outgoingPeer == nil {
return
}
// Track available HTLC information and report to handler.
htlcInfo := &HtlcInfo{
addTime: inFlight.addedTs,
resolveTime: resolution.timestamp,
settled: resolution.settled,
incomingMsat: inFlight.incomingMsat,
outgoingMsat: inFlight.outgoingMsat,
incomingCircuit: key,
outgoingCircuit: resolution.outgoingCircuitKey,
incomingPeer: p.pubKey,
outgoingPeer: *resolution.outgoingPeer,
}
if err := p.htlcCompleted(ctx, htlcInfo); err != nil {
p.logger.Infof("Mark htlc complete failed: %v", err)
}
}
2023-01-03 13:01:26 +01:00
func (p *peerController) incrCounter(event eventType) {
for _, counter := range p.rateCounters {
counter.Incr(event)
}
}
func getRate(maxHourlyRate int64) rate.Limit {
if maxHourlyRate == 0 {
return rate.Inf
}
return rate.Limit(float64(maxHourlyRate) / 3600)
}
2022-11-29 17:44:55 +01:00
func (p *peerController) forward(event interceptEvent) error {
p.htlcs[event.circuitKey] = &inFlightHtlc{
addedTs: p.now(),
incomingMsat: event.incomingMsat,
outgoingMsat: event.outgoingMsat,
}
2022-11-29 15:18:13 +01:00
2022-11-29 17:44:55 +01:00
err := event.resume(true)
if err != nil {
return err
}
2022-11-29 15:18:13 +01:00
2022-11-29 17:44:55 +01:00
logger := p.keyLogger(event.circuitKey)
logger.Infow("Forwarded", "pending_htlcs", len(p.htlcs))
return nil
2022-11-29 15:18:13 +01:00
}
2022-11-29 17:44:55 +01:00
func (p *peerController) process(ctx context.Context,
event peerInterceptEvent) error {
select {
case p.interceptChan <- event:
return nil
case <-ctx.Done():
return ctx.Err()
2022-11-29 15:18:13 +01:00
}
2022-11-29 17:44:55 +01:00
}
2022-11-29 15:18:13 +01:00
2022-11-29 17:44:55 +01:00
func (p *peerController) resolved(ctx context.Context,
key peerResolvedEvent) error {
2022-11-29 15:18:13 +01:00
2022-11-29 17:44:55 +01:00
select {
case p.resolvedChan <- key:
return nil
case <-ctx.Done():
return ctx.Err()
}
2022-11-29 15:18:13 +01:00
}
func (p *peerController) keyLogger(key circuitKey) *zap.SugaredLogger {
return p.logger.With(
"htlc", key.htlc,
"channel", key.channel)
}