From 08ced14629c7fb6cbe67ef36969a5842fea01c62 Mon Sep 17 00:00:00 2001 From: Lily Zhang Date: Fri, 18 Jul 2025 08:14:06 -0700 Subject: [PATCH 1/2] handle agbot inbound queue blocking Signed-off-by: Lily Zhang --- agreementbot/consumer_protocol_handler.go | 36 +++++++++++++++++------ config/constants.go | 2 +- 2 files changed, 28 insertions(+), 10 deletions(-) diff --git a/agreementbot/consumer_protocol_handler.go b/agreementbot/consumer_protocol_handler.go index 560634ad8..580d83f36 100644 --- a/agreementbot/consumer_protocol_handler.go +++ b/agreementbot/consumer_protocol_handler.go @@ -239,16 +239,34 @@ func (b *BaseConsumerProtocolHandler) DispatchProtocolMessage(cmd *NewProtocolMe // Figure out what kind of message this is if reply, rerr := cph.AgreementProtocolHandler("", "", "").ValidateReply(string(cmd.Message)); rerr == nil { - agreementWork := NewHandleReply(reply, cmd.From, cmd.PubKey, cmd.MessageId) - cph.WorkQueue().InboundHigh() <- &agreementWork - if glog.V(5) { - glog.Infof(BCPHlogstring(b.Name(), fmt.Sprintf("queued reply message"))) + if ag, err := b.db.FindSingleAgreementByAgreementId(reply.AgreementId(), reply.Protocol(), []persistence.AFilter{}); err != nil { + glog.Errorf(BCPHlogstring(b.Name(), fmt.Sprintf("error finding agreement %v in the db", reply.AgreementId()))) + } else if ag == nil { + glog.Warningf(BCPHlogstring(b.Name(), fmt.Sprintf("reply ignored, cannot find agreement %v in the db", reply.AgreementId()))) + } else if ag.DeviceId != cmd.From { + glog.Warningf(BCPHlogstring(b.Name(), fmt.Sprintf("reply ignored, reply message for %v came from id %v but agreement is with %v", reply.AgreementId(), cmd.From, ag.DeviceId))) + } else { + agreementWork := NewHandleReply(reply, cmd.From, cmd.PubKey, cmd.MessageId) + cph.WorkQueue().InboundHigh() <- &agreementWork + if glog.V(5) { + glog.Infof(BCPHlogstring(b.Name(), fmt.Sprintf("queued reply message"))) + } } - } else if _, aerr := cph.AgreementProtocolHandler("", "", "").ValidateDataReceivedAck(string(cmd.Message)); aerr == nil { - agreementWork := NewHandleDataReceivedAck(string(cmd.Message), cmd.From, cmd.PubKey, cmd.MessageId) - cph.WorkQueue().InboundHigh() <- &agreementWork - if glog.V(5) { - glog.Infof(BCPHlogstring(b.Name(), fmt.Sprintf("queued data received ack message"))) + } else if dra, aerr := cph.AgreementProtocolHandler("", "", "").ValidateDataReceivedAck(string(cmd.Message)); aerr == nil { + if drAck, ok := dra.(*abstractprotocol.BaseDataReceivedAck); !ok { + glog.Errorf(BCPHlogstring(b.Name(), fmt.Sprintf("unable to cast Data Received Ack %v to %v Proposal Reply, is %T", dra, cph.Name(), dra))) + } else if ag, err := b.db.FindSingleAgreementByAgreementId(drAck.AgreementId(), reply.Protocol(), []persistence.AFilter{}); err != nil { + glog.Errorf(BCPHlogstring(b.Name(), fmt.Sprintf("error finding agreement %v in the db", drAck.AgreementId()))) + } else if ag == nil { + glog.Warningf(BCPHlogstring(b.Name(), fmt.Sprintf("data received ack ignored, cannot find agreement %v in the db", drAck.AgreementId()))) + } else if ag.DeviceId != cmd.From { + glog.Warningf(BCPHlogstring(b.Name(), fmt.Sprintf("data received ack ignored, data received ack message for %v came from id %v but agreement is with %v", drAck.AgreementId(), cmd.From, ag.DeviceId))) + } else { + agreementWork := NewHandleDataReceivedAck(string(cmd.Message), cmd.From, cmd.PubKey, cmd.MessageId) + cph.WorkQueue().InboundHigh() <- &agreementWork + if glog.V(5) { + glog.Infof(BCPHlogstring(b.Name(), fmt.Sprintf("queued data received ack message"))) + } } } else if can, cerr := cph.AgreementProtocolHandler("", "", "").ValidateCancel(string(cmd.Message)); cerr == nil { // Before dispatching the cancel to a worker thread, make sure it's a valid cancel diff --git a/config/constants.go b/config/constants.go index 0bd8abb89..eada75401 100644 --- a/config/constants.go +++ b/config/constants.go @@ -103,7 +103,7 @@ const AnaxAPIPortDefault = "8510" const AgbotAgreementBatchSize_DEFAULT = 300 // The default max agreement bot work queue size. This is essentially the maximum queue depth for a given agbot protocol worker pool. -const AgbotAgreementQueueSize_DEFAULT = 300 +const AgbotAgreementQueueSize_DEFAULT = 600 // The default scaling factor applied to Agreement Queue size inorder to keep the message queue full. const AgbotMessageQueueScale_DEFAULT = 33.0 From f6c76523df99de85e50e40f98e1b570d6c3d0142 Mon Sep 17 00:00:00 2001 From: Lily Zhang Date: Fri, 18 Jul 2025 08:56:32 -0700 Subject: [PATCH 2/2] update Signed-off-by: Lily Zhang --- agreementbot/basic_protocol_handler.go | 59 +++++++++++++++++++------- 1 file changed, 44 insertions(+), 15 deletions(-) diff --git a/agreementbot/basic_protocol_handler.go b/agreementbot/basic_protocol_handler.go index d48216f38..4c8fde921 100644 --- a/agreementbot/basic_protocol_handler.go +++ b/agreementbot/basic_protocol_handler.go @@ -220,24 +220,53 @@ func (b *BasicProtocolHandler) HandleExtensionMessage(cmd *NewProtocolMessageCom // Figure out what kind of message this is if verify, perr := b.agreementPH.ValidateAgreementVerify(string(cmd.Message)); perr == nil { - agreementWork := NewBAgreementVerification(verify, cmd.From, cmd.PubKey, cmd.MessageId) - b.WorkQueue().InboundHigh() <- &agreementWork - glog.V(5).Infof(BsCPHlogString(fmt.Sprintf("queued agreement verify message"))) - + if ag, err := b.db.FindSingleAgreementByAgreementId(verify.AgreementId(), verify.Protocol(), []persistence.AFilter{}); err != nil { + glog.Errorf(BCPHlogstring(b.Name(), fmt.Sprintf("error finding agreement %v in the db", verify.AgreementId()))) + } else if ag == nil { + glog.Warningf(BCPHlogstring(b.Name(), fmt.Sprintf("agreement verify message ignored, cannot find agreement %v in the db", verify.AgreementId()))) + } else if ag.DeviceId != cmd.From { + glog.Warningf(BCPHlogstring(b.Name(), fmt.Sprintf("agreement verify message ignored, reply message for %v came from id %v but agreement is with %v", verify.AgreementId(), cmd.From, ag.DeviceId))) + } else { + agreementWork := NewBAgreementVerification(verify, cmd.From, cmd.PubKey, cmd.MessageId) + b.WorkQueue().InboundHigh() <- &agreementWork + glog.V(5).Infof(BsCPHlogString(fmt.Sprintf("queued agreement verify message"))) + } } else if verifyr, perr := b.agreementPH.ValidateAgreementVerifyReply(string(cmd.Message)); perr == nil { - agreementWork := NewBAgreementVerificationReply(verifyr, cmd.From, cmd.PubKey, cmd.MessageId) - b.WorkQueue().InboundHigh() <- &agreementWork - glog.V(5).Infof(BsCPHlogString(fmt.Sprintf("queued agreement verify reply message"))) - + if ag, err := b.db.FindSingleAgreementByAgreementId(verifyr.AgreementId(), verifyr.Protocol(), []persistence.AFilter{}); err != nil { + glog.Errorf(BCPHlogstring(b.Name(), fmt.Sprintf("error finding agreement %v in the db", verifyr.AgreementId()))) + } else if ag == nil { + glog.Warningf(BCPHlogstring(b.Name(), fmt.Sprintf("agreement verify reply message ignored, cannot find agreement %v in the db", verifyr.AgreementId()))) + } else if ag.DeviceId != cmd.From { + glog.Warningf(BCPHlogstring(b.Name(), fmt.Sprintf("agreement verify reply message ignored, reply message for %v came from id %v but agreement is with %v", verifyr.AgreementId(), cmd.From, ag.DeviceId))) + } else { + agreementWork := NewBAgreementVerificationReply(verifyr, cmd.From, cmd.PubKey, cmd.MessageId) + b.WorkQueue().InboundHigh() <- &agreementWork + glog.V(5).Infof(BsCPHlogString(fmt.Sprintf("queued agreement verify reply message"))) + } } else if update, perr := b.agreementPH.ValidateUpdate(string(cmd.Message)); perr == nil { - agreementWork := NewBAgreementUpdate(update, cmd.From, cmd.PubKey, cmd.MessageId) - b.WorkQueue().InboundHigh() <- &agreementWork - glog.V(5).Infof(BsCPHlogString(fmt.Sprintf("queued agreement update message"))) + if ag, err := b.db.FindSingleAgreementByAgreementId(update.AgreementId(), update.Protocol(), []persistence.AFilter{}); err != nil { + glog.Errorf(BCPHlogstring(b.Name(), fmt.Sprintf("error finding agreement %v in the db", update.AgreementId()))) + } else if ag == nil { + glog.Warningf(BCPHlogstring(b.Name(), fmt.Sprintf("agreement update message ignored, cannot find agreement %v in the db", update.AgreementId()))) + } else if ag.DeviceId != cmd.From { + glog.Warningf(BCPHlogstring(b.Name(), fmt.Sprintf("agreement update message ignored, reply message for %v came from id %v but agreement is with %v", update.AgreementId(), cmd.From, ag.DeviceId))) + } else { + agreementWork := NewBAgreementUpdate(update, cmd.From, cmd.PubKey, cmd.MessageId) + b.WorkQueue().InboundHigh() <- &agreementWork + glog.V(5).Infof(BsCPHlogString(fmt.Sprintf("queued agreement update message"))) + } } else if updater, perr := b.agreementPH.ValidateUpdateReply(string(cmd.Message)); perr == nil { - agreementWork := NewBAgreementUpdateReply(updater, cmd.From, cmd.PubKey, cmd.MessageId) - b.WorkQueue().InboundHigh() <- &agreementWork - glog.V(5).Infof(BsCPHlogString(fmt.Sprintf("queued areement update reply message"))) - + if ag, err := b.db.FindSingleAgreementByAgreementId(updater.AgreementId(), updater.Protocol(), []persistence.AFilter{}); err != nil { + glog.Errorf(BCPHlogstring(b.Name(), fmt.Sprintf("error finding agreement %v in the db", update.AgreementId()))) + } else if ag == nil { + glog.Warningf(BCPHlogstring(b.Name(), fmt.Sprintf("agreement update reply message ignored, cannot find agreement %v in the db", updater.AgreementId()))) + } else if ag.DeviceId != cmd.From { + glog.Warningf(BCPHlogstring(b.Name(), fmt.Sprintf("agreement update reply message ignored, reply message for %v came from id %v but agreement is with %v", updater.AgreementId(), cmd.From, ag.DeviceId))) + } else { + agreementWork := NewBAgreementUpdateReply(updater, cmd.From, cmd.PubKey, cmd.MessageId) + b.WorkQueue().InboundHigh() <- &agreementWork + glog.V(5).Infof(BsCPHlogString(fmt.Sprintf("queued agreement update reply message"))) + } } else { glog.V(5).Infof(BsCPHlogString(fmt.Sprintf("ignoring message: %v because it is an unknown type", string(cmd.Message)))) return errors.New(BsCPHlogString(fmt.Sprintf("unknown protocol msg %s", cmd.Message)))