mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-20 10:52:25 +00:00
cmd/swarm/swarm-smoke-pss: add symmetric send msg
This commit is contained in:
parent
86f5697fca
commit
e03a8a75dd
2 changed files with 106 additions and 55 deletions
|
|
@ -104,10 +104,16 @@ func main() {
|
|||
|
||||
app.Commands = []cli.Command{
|
||||
{
|
||||
Name: "pss-asym",
|
||||
Aliases: []string{"pa"},
|
||||
Name: "asym",
|
||||
Aliases: []string{"a"},
|
||||
Usage: "PSS: send and receive multiple messages across random nodes using asymmetric encryption",
|
||||
Action: wrapCliCommand("pss-asym", pssAsymCheck),
|
||||
Action: wrapCliCommand("asym", pssAsymCheck),
|
||||
},
|
||||
{
|
||||
Name: "sym",
|
||||
Aliases: []string{"s"},
|
||||
Usage: "PSS: send and receive multiple messages across random nodes using symmetric encryption",
|
||||
Action: wrapCliCommand("sym", pssSymCheck),
|
||||
},
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -19,15 +19,10 @@ import (
|
|||
cli "gopkg.in/urfave/cli.v1"
|
||||
)
|
||||
|
||||
const (
|
||||
pssModeRaw = iota
|
||||
pssModeAsym
|
||||
pssModeSym
|
||||
)
|
||||
|
||||
type pssJob struct {
|
||||
msg []byte
|
||||
mode int
|
||||
sender *pssNode
|
||||
receiver *pssNode
|
||||
msg []byte
|
||||
}
|
||||
|
||||
type pssNode struct {
|
||||
|
|
@ -40,15 +35,23 @@ type pssNode struct {
|
|||
}
|
||||
|
||||
type pssSession struct {
|
||||
topic pss.Topic
|
||||
msgC chan pss.APIMsg
|
||||
nodes []*pssNode
|
||||
requiredJobs map[string]*pssJob
|
||||
topic pss.Topic
|
||||
msgC chan pss.APIMsg
|
||||
nodes []*pssNode
|
||||
jobs map[string]*pssJob
|
||||
}
|
||||
|
||||
// pssAsymCheck sends one or more messages (depending on pssMessageCount)
|
||||
// using asymmetric encryption across random nodes.
|
||||
type pssTestFn func(ctx *cli.Context, session *pssSession, tuid string) error
|
||||
|
||||
func pssAsymCheck(ctx *cli.Context, tuid string) error {
|
||||
return pssCheck(ctx, tuid, "asym", pssAsymDo)
|
||||
}
|
||||
|
||||
func pssSymCheck(ctx *cli.Context, tuid string) error {
|
||||
return pssCheck(ctx, tuid, "sym", pssSymDo)
|
||||
}
|
||||
|
||||
func pssCheck(ctx *cli.Context, tuid string, tag string, fn pssTestFn) error {
|
||||
|
||||
// use input seed if it has been set
|
||||
if inputSeed != 0 {
|
||||
|
|
@ -59,9 +62,8 @@ func pssAsymCheck(ctx *cli.Context, tuid string) error {
|
|||
pssMessageCount = 1
|
||||
log.Warn(fmt.Sprintf("message count should be a positive number. Defaulting to %d", pssMessageCount))
|
||||
}
|
||||
log.Info("pss-asym test started", "msgCount", pssMessageCount)
|
||||
log.Info(fmt.Sprintf("pss.%s test started", tag), "msgCount", pssMessageCount)
|
||||
|
||||
errc := make(chan error)
|
||||
session := pssSetup()
|
||||
if len(session.nodes) <= 1 {
|
||||
return errors.New("at least 2 nodes are required to be working")
|
||||
|
|
@ -72,10 +74,11 @@ func pssAsymCheck(ctx *cli.Context, tuid string) error {
|
|||
}
|
||||
}()
|
||||
|
||||
errc := make(chan error)
|
||||
go func() {
|
||||
var failCount, successCount int64
|
||||
for i := 0; i < pssMessageCount; i++ {
|
||||
err := pssAsymDo(ctx, session, tuid)
|
||||
err := fn(ctx, session, tuid)
|
||||
if err != nil {
|
||||
failCount++
|
||||
log.Error("error sending pss msg", "err", err)
|
||||
|
|
@ -84,10 +87,10 @@ func pssAsymCheck(ctx *cli.Context, tuid string) error {
|
|||
}
|
||||
}
|
||||
|
||||
metrics.GetOrRegisterCounter("pss.asym.failMsg", nil).Inc(failCount)
|
||||
metrics.GetOrRegisterCounter("pss.asym.successMsg", nil).Inc(successCount)
|
||||
metrics.GetOrRegisterCounter(fmt.Sprintf("pss.%s.failMsg", tag), nil).Inc(failCount)
|
||||
metrics.GetOrRegisterCounter(fmt.Sprintf("pss.%s.successMsg", tag), nil).Inc(successCount)
|
||||
|
||||
log.Info("pss-asym test ended", "success", successCount, "failures", failCount)
|
||||
log.Info(fmt.Sprintf("pss.%s test ended", tag), "success", successCount, "failures", failCount)
|
||||
|
||||
if failCount > 0 {
|
||||
errc <- errors.New("some messages were not delivered")
|
||||
|
|
@ -99,11 +102,11 @@ func pssAsymCheck(ctx *cli.Context, tuid string) error {
|
|||
select {
|
||||
case err := <-errc:
|
||||
if err != nil {
|
||||
metrics.GetOrRegisterCounter("pss.asym.fail", nil).Inc(1)
|
||||
metrics.GetOrRegisterCounter(fmt.Sprintf("pss.%s.fail", tag), nil).Inc(1)
|
||||
}
|
||||
return err
|
||||
case <-time.After(time.Duration(timeout) * time.Second):
|
||||
metrics.GetOrRegisterCounter("pss.asym.timeout", nil).Inc(1)
|
||||
metrics.GetOrRegisterCounter(fmt.Sprintf("pss.%s.timeout", tag), nil).Inc(1)
|
||||
return fmt.Errorf("timeout after %v sec", timeout)
|
||||
}
|
||||
|
||||
|
|
@ -116,9 +119,9 @@ func pssSetup() *pssSession {
|
|||
log.Trace("pss random topic", "topic", topic.String())
|
||||
|
||||
session := &pssSession{
|
||||
msgC: make(chan pss.APIMsg),
|
||||
requiredJobs: make(map[string]*pssJob),
|
||||
topic: topic,
|
||||
msgC: make(chan pss.APIMsg),
|
||||
jobs: make(map[string]*pssJob),
|
||||
topic: topic,
|
||||
}
|
||||
|
||||
// set up the necessary info for each pss node
|
||||
|
|
@ -171,51 +174,46 @@ func pssSetup() *pssSession {
|
|||
return session
|
||||
}
|
||||
|
||||
// pssAsymDo sends a single PSS message between two random nodes using asymetric encryption
|
||||
func pssAsymDo(ctx *cli.Context, session *pssSession, tuid string) error {
|
||||
|
||||
senderNodeIdx := rand.Intn(len(session.nodes))
|
||||
senderNode := session.nodes[senderNodeIdx]
|
||||
// genJob generates a random message that will be sent
|
||||
// from a random sending node to a random receiving node
|
||||
func (s *pssSession) genJob() pssJob {
|
||||
senderNodeIdx := rand.Intn(len(s.nodes))
|
||||
senderNode := s.nodes[senderNodeIdx]
|
||||
log.Trace("sender node", "pss_baseAddr", hexutil.Encode(senderNode.addr), "host", hosts[senderNodeIdx])
|
||||
|
||||
// receiving node has to be different than sender
|
||||
recvNodeIdx := senderNodeIdx
|
||||
for recvNodeIdx == senderNodeIdx {
|
||||
recvNodeIdx = rand.Intn(len(session.nodes))
|
||||
recvNodeIdx = rand.Intn(len(s.nodes))
|
||||
}
|
||||
recvNode := session.nodes[recvNodeIdx]
|
||||
recvNode := s.nodes[recvNodeIdx]
|
||||
log.Trace("recv node", "pss_baseAddr", hexutil.Encode(recvNode.addr), "host", hosts[recvNodeIdx])
|
||||
|
||||
// set recipient in pivot node
|
||||
err := senderNode.client.Call(nil, "pss_setPeerPublicKey", recvNode.pubkey, session.topic, "0x")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// create new message and add it to job index to check for receives
|
||||
randomMsg := testutil.RandomBytes(seed, 128)
|
||||
// change seed so that the next random message is different
|
||||
seed = seed + 1
|
||||
|
||||
j := pssJob{
|
||||
sender: senderNode,
|
||||
receiver: recvNode,
|
||||
msg: randomMsg,
|
||||
}
|
||||
|
||||
msgIdx := toMsgIdx(randomMsg)
|
||||
session.requiredJobs[msgIdx] = &pssJob{
|
||||
msg: randomMsg,
|
||||
mode: pssModeAsym,
|
||||
}
|
||||
s.jobs[msgIdx] = &j
|
||||
|
||||
// send the message
|
||||
hostIdx := session.nodes[recvNodeIdx].hostIdx
|
||||
log.Debug("sending msg", "job", hexutil.Encode([]byte(msgIdx)), "sender", hosts[senderNodeIdx], "recv", hosts[hostIdx])
|
||||
err = senderNode.client.Call(nil, "pss_sendAsym", recvNode.pubkey, session.topic, hexutil.Encode(randomMsg))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
log.Debug("generated job", "job", hexutil.Encode([]byte(msgIdx)), "sender", hosts[j.sender.hostIdx], "recv", hosts[j.receiver.hostIdx])
|
||||
|
||||
// receive the message and check for its content
|
||||
return j
|
||||
}
|
||||
|
||||
// waitForJob blocks until a msg is received or a timeout is reached
|
||||
func (s *pssSession) waitForMsg() error {
|
||||
select {
|
||||
case res := <-session.msgC:
|
||||
case res := <-s.msgC:
|
||||
resIdx := toMsgIdx(res.Msg)
|
||||
resJob, ok := session.requiredJobs[resIdx]
|
||||
resJob, ok := s.jobs[resIdx]
|
||||
if !ok {
|
||||
return fmt.Errorf("corrupt message for job %s", resIdx)
|
||||
}
|
||||
|
|
@ -229,6 +227,53 @@ func pssAsymDo(ctx *cli.Context, session *pssSession, tuid string) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// pssSymDo sends a single PSS message between two random nodes using symmetric encryption
|
||||
func pssSymDo(ctx *cli.Context, session *pssSession, tuid string) error {
|
||||
j := session.genJob()
|
||||
|
||||
symkey := make([]byte, 32)
|
||||
c, err := rand.Read(symkey)
|
||||
if err != nil {
|
||||
return err
|
||||
} else if c < 32 {
|
||||
return fmt.Errorf("symkey size mismatch, expected 32 got %d", c)
|
||||
}
|
||||
|
||||
var senderSymKeyID string
|
||||
err = j.sender.client.Call(&senderSymKeyID, "pss_setSymmetricKey", symkey, session.topic, hexutil.Encode(j.receiver.addr), true)
|
||||
if err != nil {
|
||||
log.Error("erro setting sym key on the sender", "err", err)
|
||||
return err
|
||||
}
|
||||
|
||||
var recvSymKeyID string
|
||||
err = j.receiver.client.Call(&recvSymKeyID, "pss_setSymmetricKey", symkey, session.topic, hexutil.Encode(j.sender.addr), true)
|
||||
if err != nil {
|
||||
log.Error("error setting sym key on the receiver", "err", err)
|
||||
return err
|
||||
}
|
||||
|
||||
err = j.sender.client.Call(nil, "pss_sendSym", senderSymKeyID, session.topic, hexutil.Encode(j.msg))
|
||||
if err != nil {
|
||||
log.Error("error sending message using sym encryption", "err", err)
|
||||
return err
|
||||
}
|
||||
|
||||
return session.waitForMsg()
|
||||
}
|
||||
|
||||
// pssAsymDo sends a single PSS message between two random nodes using asymmetric encryption
|
||||
func pssAsymDo(ctx *cli.Context, session *pssSession, tuid string) error {
|
||||
j := session.genJob()
|
||||
|
||||
err := j.sender.client.Call(nil, "pss_sendAsym", j.receiver.pubkey, session.topic, hexutil.Encode(j.msg))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return session.waitForMsg()
|
||||
}
|
||||
|
||||
func toMsgIdx(msg []byte) string {
|
||||
h := sha1.New()
|
||||
h.Write(msg)
|
||||
|
|
|
|||
Loading…
Reference in a new issue