Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
59 changes: 44 additions & 15 deletions agreementbot/basic_protocol_handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)))
Expand Down
36 changes: 27 additions & 9 deletions agreementbot/consumer_protocol_handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion config/constants.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading