Skip to content

Commit 30a13bb

Browse files
committed
Skip logging error on duplicate key gen
1 parent 57314fb commit 30a13bb

4 files changed

Lines changed: 37 additions & 18 deletions

File tree

pkg/eventconsumer/event_consumer.go

Lines changed: 17 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -137,13 +137,11 @@ func (ec *eventConsumer) handleKeyGenEvent(natMsg *nats.Msg) {
137137
walletID := msg.WalletID
138138
ecdsaSession, err := ec.node.CreateKeyGenSession(mpc.SessionTypeECDSA, walletID, ec.mpcThreshold, ec.genKeyResultQueue)
139139
if err != nil {
140-
logger.Error("Failed to create ECDSA key generation session", err, "walletID", walletID)
141140
ec.handleKeygenSessionError(walletID, err, "Failed to create ECDSA key generation session", natMsg)
142141
return
143142
}
144143
eddsaSession, err := ec.node.CreateKeyGenSession(mpc.SessionTypeEDDSA, walletID, ec.mpcThreshold, ec.genKeyResultQueue)
145144
if err != nil {
146-
logger.Error("Failed to create EdDSA key generation session", err, "walletID", walletID)
147145
ec.handleKeygenSessionError(walletID, err, "Failed to create EdDSA key generation session", natMsg)
148146
return
149147
}
@@ -225,7 +223,11 @@ func (ec *eventConsumer) handleKeyGenEvent(natMsg *nats.Msg) {
225223
}
226224

227225
key := fmt.Sprintf(mpc.TypeGenerateWalletResultFmt, walletID)
228-
if err := ec.genKeyResultQueue.Enqueue(key, payload, &messaging.EnqueueOptions{IdempotententKey: key}); err != nil {
226+
if err := ec.genKeyResultQueue.Enqueue(
227+
key,
228+
payload,
229+
&messaging.EnqueueOptions{IdempotententKey: composeKeygenIdempotentKey(walletID, natMsg)},
230+
); err != nil {
229231
logger.Error("Failed to publish key generation success message", err)
230232
ec.handleKeygenSessionError(walletID, err, "Failed to publish key generation success message", natMsg)
231233
return
@@ -238,14 +240,6 @@ func (ec *eventConsumer) handleKeyGenEvent(natMsg *nats.Msg) {
238240
func (ec *eventConsumer) handleKeygenSessionError(walletID string, err error, contextMsg string, natMsg *nats.Msg) {
239241
fullErrMsg := fmt.Sprintf("%s: %v", contextMsg, err)
240242
errorCode := event.GetErrorCodeFromError(err)
241-
242-
logger.Warn("Keygen session error",
243-
"walletID", walletID,
244-
"error", err.Error(),
245-
"errorCode", errorCode,
246-
"context", contextMsg,
247-
)
248-
249243
keygenResult := event.KeygenResultEvent{
250244
ResultType: event.ResultTypeError,
251245
ErrorCode: string(errorCode),
@@ -263,7 +257,7 @@ func (ec *eventConsumer) handleKeygenSessionError(walletID string, err error, co
263257

264258
key := fmt.Sprintf(mpc.TypeGenerateWalletResultFmt, walletID)
265259
err = ec.genKeyResultQueue.Enqueue(key, keygenResultBytes, &messaging.EnqueueOptions{
266-
IdempotententKey: key,
260+
IdempotententKey: composeKeygenIdempotentKey(walletID, natMsg),
267261
})
268262
if err != nil {
269263
logger.Error("Failed to enqueue keygen result event", err,
@@ -795,3 +789,14 @@ func sessionTypeFromKeyType(keyType types.KeyType) (mpc.SessionType, error) {
795789
return "", fmt.Errorf("unsupported key type: %v", keyType)
796790
}
797791
}
792+
793+
func composeKeygenIdempotentKey(walletID string, natMsg *nats.Msg) string {
794+
var uniqueKey string
795+
sid := natMsg.Header.Get("SessionID")
796+
if sid != "" {
797+
uniqueKey = fmt.Sprintf("%s:%s", walletID, sid)
798+
} else {
799+
uniqueKey = walletID
800+
}
801+
return fmt.Sprintf(mpc.TypeGenerateWalletResultFmt, uniqueKey)
802+
}

pkg/eventconsumer/keygen_consumer.go

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import (
99
"github.com/fystack/mpcium/pkg/logger"
1010
"github.com/fystack/mpcium/pkg/messaging"
1111
"github.com/fystack/mpcium/pkg/mpc"
12+
"github.com/google/uuid"
1213
"github.com/nats-io/nats.go"
1314
"github.com/nats-io/nats.go/jetstream"
1415
)
@@ -142,7 +143,10 @@ func (sc *keygenConsumer) handleKeygenEvent(msg jetstream.Msg) {
142143
}()
143144

144145
// Publish the signing event with the reply inbox.
145-
if err := sc.pubsub.PublishWithReply(MPCGenerateEvent, replyInbox, msg.Data()); err != nil {
146+
headers := map[string]string{
147+
"SessionID": uuid.New().String(),
148+
}
149+
if err := sc.pubsub.PublishWithReply(MPCGenerateEvent, replyInbox, msg.Data(), headers); err != nil {
146150
logger.Error("KeygenConsumer: Failed to publish signing event with reply", err)
147151
_ = msg.Nak()
148152
return

pkg/eventconsumer/sign_consumer.go

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import (
99
"github.com/fystack/mpcium/pkg/logger"
1010
"github.com/fystack/mpcium/pkg/messaging"
1111
"github.com/fystack/mpcium/pkg/mpc"
12+
"github.com/google/uuid"
1213
"github.com/nats-io/nats.go"
1314
"github.com/nats-io/nats.go/jetstream"
1415
"github.com/spf13/viper"
@@ -158,7 +159,10 @@ func (sc *signingConsumer) handleSigningEvent(msg jetstream.Msg) {
158159
}()
159160

160161
// Publish the signing event with the reply inbox.
161-
if err := sc.pubsub.PublishWithReply(MPCSignEvent, replyInbox, msg.Data()); err != nil {
162+
headers := map[string]string{
163+
"SessionID": uuid.New().String(),
164+
}
165+
if err := sc.pubsub.PublishWithReply(MPCSignEvent, replyInbox, msg.Data(), headers); err != nil {
162166
logger.Error("SigningConsumer: Failed to publish signing event with reply", err)
163167
_ = msg.Nak()
164168
return

pkg/messaging/pubsub.go

Lines changed: 10 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@ type Subscription interface {
1111

1212
type PubSub interface {
1313
Publish(topic string, message []byte) error
14-
PublishWithReply(topic, reply string, data []byte) error
14+
PublishWithReply(topic, reply string, data []byte, headers map[string]string) error
1515
Subscribe(topic string, handler func(msg *nats.Msg)) (Subscription, error)
1616
}
1717

@@ -36,12 +36,18 @@ func (n *natsPubSub) Publish(topic string, message []byte) error {
3636
return n.natsConn.Publish(topic, message)
3737
}
3838

39-
func (n *natsPubSub) PublishWithReply(topic, reply string, data []byte) error {
40-
return n.natsConn.PublishMsg(&nats.Msg{
39+
func (n *natsPubSub) PublishWithReply(topic, reply string, data []byte, headers map[string]string) error {
40+
msg := &nats.Msg{
4141
Subject: topic,
4242
Reply: reply,
4343
Data: data,
44-
})
44+
Header: nats.Header{},
45+
}
46+
for k, v := range headers {
47+
msg.Header.Set(k, v)
48+
}
49+
err := n.natsConn.PublishMsg(msg)
50+
return err
4551
}
4652

4753
func (n *natsPubSub) Subscribe(topic string, handler func(msg *nats.Msg)) (Subscription, error) {

0 commit comments

Comments
 (0)