2022-09-30 17:07:13 +00:00
|
|
|
package handlers
|
|
|
|
|
|
|
|
import (
|
|
|
|
"reflect"
|
|
|
|
"strings"
|
|
|
|
|
2022-10-06 12:34:39 +00:00
|
|
|
"github.com/ledgerwatch/erigon/cmd/lightclient/sentinel/communication"
|
2022-09-30 17:07:13 +00:00
|
|
|
"github.com/ledgerwatch/log/v3"
|
|
|
|
"github.com/libp2p/go-libp2p/core/network"
|
|
|
|
)
|
|
|
|
|
2022-10-06 12:34:39 +00:00
|
|
|
// curryStreamHandler converts a func(ctx *communication.StreamContext, dat communication.Packet) error to func(network.Stream)
|
2022-09-30 17:07:13 +00:00
|
|
|
// this allows us to write encoding non specific type safe handler without performance overhead
|
2022-10-06 12:34:39 +00:00
|
|
|
func curryStreamHandler[T communication.Packet](newcodec func(network.Stream) communication.StreamCodec, fn func(ctx *communication.StreamContext, v T) error) func(network.Stream) {
|
2022-09-30 17:07:13 +00:00
|
|
|
return func(s network.Stream) {
|
2022-09-30 22:25:17 +00:00
|
|
|
defer s.Close()
|
2022-09-30 17:07:13 +00:00
|
|
|
sd := newcodec(s)
|
|
|
|
var t T
|
|
|
|
val := t.Clone().(T)
|
|
|
|
ctx, err := sd.Decode(val)
|
|
|
|
if err != nil {
|
|
|
|
// the stream reset error is ignored, because
|
|
|
|
if !strings.Contains(err.Error(), "stream reset") {
|
|
|
|
log.Debug("fail to decode packet", "err", err, "path", ctx.Protocol, "pkt", reflect.TypeOf(val))
|
|
|
|
}
|
|
|
|
return
|
|
|
|
}
|
|
|
|
err = fn(ctx, val)
|
|
|
|
if err != nil {
|
|
|
|
log.Debug("failed handling packet", "err", err, "path", ctx.Protocol, "pkt", reflect.TypeOf(val))
|
|
|
|
return
|
|
|
|
}
|
|
|
|
log.Trace("[ReqResp] Req->Host", "from", ctx.Stream.ID(), "endpoint", ctx.Protocol, "msg", val)
|
|
|
|
}
|
|
|
|
}
|