@@ -35,8 +35,9 @@ func (pp *PathProcessor) assemblePacketIBCMessage(
3535
3636func (pp * PathProcessor ) getUnrelayedPacketsAndAcksAndToDelete (ctx context.Context , pathEndPacketFlowMessages pathEndPacketFlowMessages ) pathEndPacketFlowResponse {
3737 res := pathEndPacketFlowResponse {
38- ToDeleteSrc : make (map [string ][]uint64 ),
39- ToDeleteDst : make (map [string ][]uint64 ),
38+ ToDeleteSrc : make (map [string ][]uint64 ),
39+ ToDeleteDst : make (map [string ][]uint64 ),
40+ ToDeleteDstChannel : make (map [string ][]ChannelKey ),
4041 }
4142
4243MsgTransferLoop:
@@ -51,12 +52,42 @@ MsgTransferLoop:
5152 continue MsgTransferLoop
5253 }
5354 }
54- for timeoutSeq := range pathEndPacketFlowMessages .SrcMsgTimeout {
55+
56+ for timeoutSeq , msgTimeout := range pathEndPacketFlowMessages .SrcMsgTimeout {
5557 if transferSeq == timeoutSeq {
56- // we have a timeout for this packet, so packet flow is complete
57- // remove all retention of this sequence number
58- res .ToDeleteSrc [chantypes .EventTypeSendPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeSendPacket ], transferSeq )
59- res .ToDeleteSrc [chantypes .EventTypeTimeoutPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeTimeoutPacket ], transferSeq )
58+ if msgTimeout .ChannelOrder == chantypes .ORDERED .String () {
59+ // For ordered channel packets, flow is not done until channel-close-confirm is observed.
60+ if pathEndPacketFlowMessages .DstMsgChannelCloseConfirm == nil {
61+ // have not observed a channel-close-confirm yet for this channel, send it if ready.
62+ // will come back through here next block if not yet ready.
63+ closeChan := channelIBCMessage {
64+ eventType : chantypes .EventTypeChannelCloseConfirm ,
65+ info : provider.ChannelInfo {
66+ Height : msgTimeout .Height ,
67+ PortID : msgTimeout .SourcePort ,
68+ ChannelID : msgTimeout .SourceChannel ,
69+ CounterpartyPortID : msgTimeout .DestPort ,
70+ CounterpartyChannelID : msgTimeout .DestChannel ,
71+ Order : orderFromString (msgTimeout .ChannelOrder ),
72+ },
73+ }
74+
75+ if pathEndPacketFlowMessages .Src .shouldSendChannelMessage (closeChan , pathEndPacketFlowMessages .Dst ) {
76+ res .DstChannelMessage = append (res .DstChannelMessage , closeChan )
77+ }
78+ } else {
79+ // ordered channel, and we have a channel close confirm, so packet-flow and channel-close-flow is complete.
80+ // remove all retention of this sequence number and this channel-close-confirm.
81+ res .ToDeleteDstChannel [chantypes .EventTypeChannelCloseConfirm ] = append (res .ToDeleteDstChannel [chantypes .EventTypeChannelCloseConfirm ], pathEndPacketFlowMessages .ChannelKey .Counterparty ())
82+ res .ToDeleteSrc [chantypes .EventTypeSendPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeSendPacket ], transferSeq )
83+ res .ToDeleteSrc [chantypes .EventTypeTimeoutPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeTimeoutPacket ], transferSeq )
84+ }
85+ } else {
86+ // unordered channel, and we have a timeout for this packet, so packet flow is complete
87+ // remove all retention of this sequence number
88+ res .ToDeleteSrc [chantypes .EventTypeSendPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeSendPacket ], transferSeq )
89+ res .ToDeleteSrc [chantypes .EventTypeTimeoutPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeTimeoutPacket ], transferSeq )
90+ }
6091 continue MsgTransferLoop
6192 }
6293 }
@@ -128,9 +159,18 @@ MsgTransferLoop:
128159 res .ToDeleteDst [chantypes .EventTypeRecvPacket ] = append (res .ToDeleteDst [chantypes .EventTypeRecvPacket ], ackSeq )
129160 res .ToDeleteSrc [chantypes .EventTypeAcknowledgePacket ] = append (res .ToDeleteSrc [chantypes .EventTypeAcknowledgePacket ], ackSeq )
130161 }
131- for timeoutSeq := range pathEndPacketFlowMessages .SrcMsgTimeout {
132- res .ToDeleteSrc [chantypes .EventTypeSendPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeSendPacket ], timeoutSeq )
133- res .ToDeleteSrc [chantypes .EventTypeTimeoutPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeTimeoutPacket ], timeoutSeq )
162+ for timeoutSeq , msgTimeout := range pathEndPacketFlowMessages .SrcMsgTimeout {
163+ if msgTimeout .ChannelOrder == chantypes .ORDERED .String () {
164+ if pathEndPacketFlowMessages .DstMsgChannelCloseConfirm != nil {
165+ // For ordered channel packets, flow is not done until channel-close-confirm is observed.
166+ res .ToDeleteDstChannel [chantypes .EventTypeChannelCloseConfirm ] = append (res .ToDeleteDstChannel [chantypes .EventTypeChannelCloseConfirm ], pathEndPacketFlowMessages .ChannelKey .Counterparty ())
167+ res .ToDeleteSrc [chantypes .EventTypeSendPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeSendPacket ], timeoutSeq )
168+ res .ToDeleteSrc [chantypes .EventTypeTimeoutPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeTimeoutPacket ], timeoutSeq )
169+ }
170+ } else {
171+ res .ToDeleteSrc [chantypes .EventTypeSendPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeSendPacket ], timeoutSeq )
172+ res .ToDeleteSrc [chantypes .EventTypeTimeoutPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeTimeoutPacket ], timeoutSeq )
173+ }
134174 }
135175 for timeoutOnCloseSeq := range pathEndPacketFlowMessages .SrcMsgTimeoutOnClose {
136176 res .ToDeleteSrc [chantypes .EventTypeSendPacket ] = append (res .ToDeleteSrc [chantypes .EventTypeSendPacket ], timeoutOnCloseSeq )
@@ -512,25 +552,41 @@ func (pp *PathProcessor) processLatestMessages(ctx context.Context, messageLifec
512552 pathEnd2ProcessRes := make ([]pathEndPacketFlowResponse , len (channelPairs ))
513553
514554 for i , pair := range channelPairs {
555+ var pathEnd1ChannelCloseConfirm , pathEnd2ChannelCloseConfirm * provider.ChannelInfo
556+
557+ if pathEnd1ChanCloseConfirmMsgs , ok := pp .pathEnd1 .messageCache .ChannelHandshake [chantypes .EventTypeChannelCloseConfirm ]; ok {
558+ if pathEnd1ChannelCloseConfirmMsg , ok := pathEnd1ChanCloseConfirmMsgs [pair .pathEnd1ChannelKey ]; ok {
559+ pathEnd1ChannelCloseConfirm = & pathEnd1ChannelCloseConfirmMsg
560+ }
561+ }
562+
563+ if pathEnd2ChanCloseConfirmMsgs , ok := pp .pathEnd2 .messageCache .ChannelHandshake [chantypes .EventTypeChannelCloseConfirm ]; ok {
564+ if pathEnd2ChannelCloseConfirmMsg , ok := pathEnd2ChanCloseConfirmMsgs [pair .pathEnd2ChannelKey ]; ok {
565+ pathEnd2ChannelCloseConfirm = & pathEnd2ChannelCloseConfirmMsg
566+ }
567+ }
568+
515569 pathEnd1PacketFlowMessages := pathEndPacketFlowMessages {
516- Src : pp .pathEnd1 ,
517- Dst : pp .pathEnd2 ,
518- ChannelKey : pair .pathEnd1ChannelKey ,
519- SrcMsgTransfer : pp .pathEnd1 .messageCache .PacketFlow [pair .pathEnd1ChannelKey ][chantypes .EventTypeSendPacket ],
520- DstMsgRecvPacket : pp .pathEnd2 .messageCache .PacketFlow [pair .pathEnd2ChannelKey ][chantypes .EventTypeRecvPacket ],
521- SrcMsgAcknowledgement : pp .pathEnd1 .messageCache .PacketFlow [pair .pathEnd1ChannelKey ][chantypes .EventTypeAcknowledgePacket ],
522- SrcMsgTimeout : pp .pathEnd1 .messageCache .PacketFlow [pair .pathEnd1ChannelKey ][chantypes .EventTypeTimeoutPacket ],
523- SrcMsgTimeoutOnClose : pp .pathEnd1 .messageCache .PacketFlow [pair .pathEnd1ChannelKey ][chantypes .EventTypeTimeoutPacketOnClose ],
570+ Src : pp .pathEnd1 ,
571+ Dst : pp .pathEnd2 ,
572+ ChannelKey : pair .pathEnd1ChannelKey ,
573+ SrcMsgTransfer : pp .pathEnd1 .messageCache .PacketFlow [pair .pathEnd1ChannelKey ][chantypes .EventTypeSendPacket ],
574+ DstMsgRecvPacket : pp .pathEnd2 .messageCache .PacketFlow [pair .pathEnd2ChannelKey ][chantypes .EventTypeRecvPacket ],
575+ SrcMsgAcknowledgement : pp .pathEnd1 .messageCache .PacketFlow [pair .pathEnd1ChannelKey ][chantypes .EventTypeAcknowledgePacket ],
576+ SrcMsgTimeout : pp .pathEnd1 .messageCache .PacketFlow [pair .pathEnd1ChannelKey ][chantypes .EventTypeTimeoutPacket ],
577+ SrcMsgTimeoutOnClose : pp .pathEnd1 .messageCache .PacketFlow [pair .pathEnd1ChannelKey ][chantypes .EventTypeTimeoutPacketOnClose ],
578+ DstMsgChannelCloseConfirm : pathEnd2ChannelCloseConfirm ,
524579 }
525580 pathEnd2PacketFlowMessages := pathEndPacketFlowMessages {
526- Src : pp .pathEnd2 ,
527- Dst : pp .pathEnd1 ,
528- ChannelKey : pair .pathEnd2ChannelKey ,
529- SrcMsgTransfer : pp .pathEnd2 .messageCache .PacketFlow [pair .pathEnd2ChannelKey ][chantypes .EventTypeSendPacket ],
530- DstMsgRecvPacket : pp .pathEnd1 .messageCache .PacketFlow [pair .pathEnd1ChannelKey ][chantypes .EventTypeRecvPacket ],
531- SrcMsgAcknowledgement : pp .pathEnd2 .messageCache .PacketFlow [pair .pathEnd2ChannelKey ][chantypes .EventTypeAcknowledgePacket ],
532- SrcMsgTimeout : pp .pathEnd2 .messageCache .PacketFlow [pair .pathEnd2ChannelKey ][chantypes .EventTypeTimeoutPacket ],
533- SrcMsgTimeoutOnClose : pp .pathEnd2 .messageCache .PacketFlow [pair .pathEnd2ChannelKey ][chantypes .EventTypeTimeoutPacketOnClose ],
581+ Src : pp .pathEnd2 ,
582+ Dst : pp .pathEnd1 ,
583+ ChannelKey : pair .pathEnd2ChannelKey ,
584+ SrcMsgTransfer : pp .pathEnd2 .messageCache .PacketFlow [pair .pathEnd2ChannelKey ][chantypes .EventTypeSendPacket ],
585+ DstMsgRecvPacket : pp .pathEnd1 .messageCache .PacketFlow [pair .pathEnd1ChannelKey ][chantypes .EventTypeRecvPacket ],
586+ SrcMsgAcknowledgement : pp .pathEnd2 .messageCache .PacketFlow [pair .pathEnd2ChannelKey ][chantypes .EventTypeAcknowledgePacket ],
587+ SrcMsgTimeout : pp .pathEnd2 .messageCache .PacketFlow [pair .pathEnd2ChannelKey ][chantypes .EventTypeTimeoutPacket ],
588+ SrcMsgTimeoutOnClose : pp .pathEnd2 .messageCache .PacketFlow [pair .pathEnd2ChannelKey ][chantypes .EventTypeTimeoutPacketOnClose ],
589+ DstMsgChannelCloseConfirm : pathEnd1ChannelCloseConfirm ,
534590 }
535591
536592 pathEnd1ProcessRes [i ] = pp .getUnrelayedPacketsAndAcksAndToDelete (ctx , pathEnd1PacketFlowMessages )
@@ -540,7 +596,10 @@ func (pp *PathProcessor) processLatestMessages(ctx context.Context, messageLifec
540596 // concatenate applicable messages for pathend
541597 pathEnd1ConnectionMessages , pathEnd2ConnectionMessages := pp .connectionMessagesToSend (pathEnd1ConnectionHandshakeRes , pathEnd2ConnectionHandshakeRes )
542598 pathEnd1ChannelMessages , pathEnd2ChannelMessages := pp .channelMessagesToSend (pathEnd1ChannelHandshakeRes , pathEnd2ChannelHandshakeRes )
543- pathEnd1PacketMessages , pathEnd2PacketMessages := pp .packetMessagesToSend (channelPairs , pathEnd1ProcessRes , pathEnd2ProcessRes )
599+
600+ pathEnd1PacketMessages , pathEnd2PacketMessages , pathEnd1ChanCloseMessages , pathEnd2ChanCloseMessages := pp .packetMessagesToSend (channelPairs , pathEnd1ProcessRes , pathEnd2ProcessRes )
601+ pathEnd1ChannelMessages = append (pathEnd1ChannelMessages , pathEnd1ChanCloseMessages ... )
602+ pathEnd2ChannelMessages = append (pathEnd2ChannelMessages , pathEnd2ChanCloseMessages ... )
544603
545604 pathEnd1Messages := pathEndMessages {
546605 connectionMessages : pathEnd1ConnectionMessages ,
@@ -626,6 +685,7 @@ func (pp *PathProcessor) assembleMessage(
626685 return
627686 }
628687 om .Append (message )
688+
629689}
630690
631691func (pp * PathProcessor ) assembleAndSendMessages (
@@ -869,29 +929,47 @@ func (pp *PathProcessor) connectionMessagesToSend(pathEnd1ConnectionHandshakeRes
869929 return pathEnd1ConnectionMessages , pathEnd2ConnectionMessages
870930}
871931
872- func (pp * PathProcessor ) packetMessagesToSend (channelPairs []channelPair , pathEnd1ProcessRes []pathEndPacketFlowResponse , pathEnd2ProcessRes []pathEndPacketFlowResponse ) ([]packetIBCMessage , []packetIBCMessage ) {
932+ func (pp * PathProcessor ) packetMessagesToSend (
933+ channelPairs []channelPair ,
934+ pathEnd1ProcessRes []pathEndPacketFlowResponse ,
935+ pathEnd2ProcessRes []pathEndPacketFlowResponse ,
936+ ) ([]packetIBCMessage , []packetIBCMessage , []channelIBCMessage , []channelIBCMessage ) {
873937 pathEnd1PacketLen := 0
874938 pathEnd2PacketLen := 0
939+ pathEnd1ChannelLen := 0
940+ pathEnd2ChannelLen := 0
941+
875942 for i := 0 ; i < len (channelPairs ); i ++ {
876943 pathEnd1PacketLen += len (pathEnd2ProcessRes [i ].DstMessages ) + len (pathEnd1ProcessRes [i ].SrcMessages )
877944 pathEnd2PacketLen += len (pathEnd1ProcessRes [i ].DstMessages ) + len (pathEnd2ProcessRes [i ].SrcMessages )
945+ pathEnd1ChannelLen += len (pathEnd2ProcessRes [i ].DstChannelMessage )
946+ pathEnd2ChannelLen += len (pathEnd1ProcessRes [i ].DstChannelMessage )
878947 }
879948
880949 pathEnd1PacketMessages := make ([]packetIBCMessage , 0 , pathEnd1PacketLen )
881950 pathEnd2PacketMessages := make ([]packetIBCMessage , 0 , pathEnd2PacketLen )
882951
952+ pathEnd1ChannelMessage := make ([]channelIBCMessage , 0 , pathEnd1ChannelLen )
953+ pathEnd2ChannelMessage := make ([]channelIBCMessage , 0 , pathEnd2ChannelLen )
954+
883955 for i , channelPair := range channelPairs {
884956 pathEnd1PacketMessages = append (pathEnd1PacketMessages , pathEnd2ProcessRes [i ].DstMessages ... )
885957 pathEnd1PacketMessages = append (pathEnd1PacketMessages , pathEnd1ProcessRes [i ].SrcMessages ... )
886958
887959 pathEnd2PacketMessages = append (pathEnd2PacketMessages , pathEnd1ProcessRes [i ].DstMessages ... )
888960 pathEnd2PacketMessages = append (pathEnd2PacketMessages , pathEnd2ProcessRes [i ].SrcMessages ... )
889961
962+ pathEnd1ChannelMessage = append (pathEnd1ChannelMessage , pathEnd2ProcessRes [i ].DstChannelMessage ... )
963+ pathEnd2ChannelMessage = append (pathEnd2ChannelMessage , pathEnd1ProcessRes [i ].DstChannelMessage ... )
964+
965+ pp .pathEnd1 .messageCache .ChannelHandshake .DeleteMessages (pathEnd2ProcessRes [i ].ToDeleteDstChannel )
966+ pp .pathEnd2 .messageCache .ChannelHandshake .DeleteMessages (pathEnd1ProcessRes [i ].ToDeleteDstChannel )
967+
890968 pp .pathEnd1 .messageCache .PacketFlow [channelPair .pathEnd1ChannelKey ].DeleteMessages (pathEnd1ProcessRes [i ].ToDeleteSrc , pathEnd2ProcessRes [i ].ToDeleteDst )
891969 pp .pathEnd2 .messageCache .PacketFlow [channelPair .pathEnd2ChannelKey ].DeleteMessages (pathEnd2ProcessRes [i ].ToDeleteSrc , pathEnd1ProcessRes [i ].ToDeleteDst )
892970 pp .pathEnd1 .packetProcessing [channelPair .pathEnd1ChannelKey ].deleteMessages (pathEnd1ProcessRes [i ].ToDeleteSrc , pathEnd2ProcessRes [i ].ToDeleteDst )
893971 pp .pathEnd2 .packetProcessing [channelPair .pathEnd2ChannelKey ].deleteMessages (pathEnd2ProcessRes [i ].ToDeleteSrc , pathEnd1ProcessRes [i ].ToDeleteDst )
894972 }
895973
896- return pathEnd1PacketMessages , pathEnd2PacketMessages
974+ return pathEnd1PacketMessages , pathEnd2PacketMessages , pathEnd1ChannelMessage , pathEnd2ChannelMessage
897975}
0 commit comments