Skip to content
This repository was archived by the owner on Apr 10, 2026. It is now read-only.

Commit 2c2c6a7

Browse files
authored
path processor race (#995)
* fix path processor indexing * clone maps before passing to pathprocessor * Merge should Clone * pre-size maps
1 parent a4206f0 commit 2c2c6a7

3 files changed

Lines changed: 37 additions & 25 deletions

File tree

relayer/chains/cosmos/cosmos_chain_processor.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -398,12 +398,12 @@ func (ccp *CosmosChainProcessor) queryCycle(ctx context.Context, persistence *qu
398398
pp.HandleNewData(chainID, processor.ChainProcessorCacheData{
399399
LatestBlock: ccp.latestBlock,
400400
LatestHeader: latestHeader,
401-
IBCMessagesCache: ibcMessagesCache,
401+
IBCMessagesCache: ibcMessagesCache.Clone(),
402402
InSync: ccp.inSync,
403403
ClientState: clientState,
404404
ConnectionStateCache: ccp.connectionStateCache.FilterForClient(clientID),
405405
ChannelStateCache: ccp.channelStateCache.FilterForClient(clientID, ccp.channelConnections, ccp.connectionClients),
406-
IBCHeaderCache: ibcHeaderCache,
406+
IBCHeaderCache: ibcHeaderCache.Clone(),
407407
})
408408
}
409409

relayer/processor/types.go

Lines changed: 32 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,19 @@ type IBCMessagesCache struct {
7777
ChannelHandshake ChannelMessagesCache
7878
}
7979

80+
// Clone makes a deep copy of an IBCMessagesCache.
81+
func (c IBCMessagesCache) Clone() IBCMessagesCache {
82+
x := IBCMessagesCache{
83+
PacketFlow: make(ChannelPacketMessagesCache, len(c.PacketFlow)),
84+
ConnectionHandshake: make(ConnectionMessagesCache, len(c.ConnectionHandshake)),
85+
ChannelHandshake: make(ChannelMessagesCache, len(c.ChannelHandshake)),
86+
}
87+
x.PacketFlow.Merge(c.PacketFlow)
88+
x.ConnectionHandshake.Merge(c.ConnectionHandshake)
89+
x.ChannelHandshake.Merge(c.ChannelHandshake)
90+
return x
91+
}
92+
8093
// NewIBCMessagesCache returns an empty IBCMessagesCache.
8194
func NewIBCMessagesCache() IBCMessagesCache {
8295
return IBCMessagesCache{
@@ -257,12 +270,10 @@ func (c PacketMessagesCache) DeleteMessages(toDelete ...map[string][]uint64) {
257270
// Merge merges another ChannelPacketMessagesCache into this one.
258271
func (c ChannelPacketMessagesCache) Merge(other ChannelPacketMessagesCache) {
259272
for channelKey, messageCache := range other {
260-
_, ok := c[channelKey]
261-
if !ok {
262-
c[channelKey] = messageCache
263-
} else {
264-
c[channelKey].Merge(messageCache)
273+
if _, ok := c[channelKey]; !ok {
274+
c[channelKey] = make(PacketMessagesCache)
265275
}
276+
c[channelKey].Merge(messageCache)
266277
}
267278
}
268279

@@ -304,12 +315,10 @@ func (c ChannelPacketMessagesCache) Retain(k ChannelKey, m string, pi provider.P
304315
// Merge merges another PacketMessagesCache into this one.
305316
func (c PacketMessagesCache) Merge(other PacketMessagesCache) {
306317
for ibcMessage, messageCache := range other {
307-
_, ok := c[ibcMessage]
308-
if !ok {
309-
c[ibcMessage] = messageCache
310-
} else {
311-
c[ibcMessage].Merge(messageCache)
318+
if _, ok := c[ibcMessage]; !ok {
319+
c[ibcMessage] = make(PacketSequenceCache)
312320
}
321+
c[ibcMessage].Merge(messageCache)
313322
}
314323
}
315324

@@ -323,12 +332,10 @@ func (c PacketSequenceCache) Merge(other PacketSequenceCache) {
323332
// Merge merges another ConnectionMessagesCache into this one.
324333
func (c ConnectionMessagesCache) Merge(other ConnectionMessagesCache) {
325334
for ibcMessage, messageCache := range other {
326-
_, ok := c[ibcMessage]
327-
if !ok {
328-
c[ibcMessage] = messageCache
329-
} else {
330-
c[ibcMessage].Merge(messageCache)
335+
if _, ok := c[ibcMessage]; !ok {
336+
c[ibcMessage] = make(ConnectionMessageCache)
331337
}
338+
c[ibcMessage].Merge(messageCache)
332339
}
333340
}
334341

@@ -360,12 +367,10 @@ func (c ConnectionMessageCache) Merge(other ConnectionMessageCache) {
360367
// Merge merges another ChannelMessagesCache into this one.
361368
func (c ChannelMessagesCache) Merge(other ChannelMessagesCache) {
362369
for ibcMessage, messageCache := range other {
363-
_, ok := c[ibcMessage]
364-
if !ok {
365-
c[ibcMessage] = messageCache
366-
} else {
367-
c[ibcMessage].Merge(messageCache)
370+
if _, ok := c[ibcMessage]; !ok {
371+
c[ibcMessage] = make(ChannelMessageCache)
368372
}
373+
c[ibcMessage].Merge(messageCache)
369374
}
370375
}
371376

@@ -397,6 +402,13 @@ func (c ChannelMessageCache) Merge(other ChannelMessageCache) {
397402
// IBCHeaderCache holds a mapping of IBCHeaders for their block height.
398403
type IBCHeaderCache map[uint64]provider.IBCHeader
399404

405+
// Clone makes a deep copy of an IBCHeaderCache.
406+
func (c IBCHeaderCache) Clone() IBCHeaderCache {
407+
x := make(IBCHeaderCache, len(c))
408+
x.Merge(c)
409+
return x
410+
}
411+
400412
// Merge merges another IBCHeaderCache into this one.
401413
func (c IBCHeaderCache) Merge(other IBCHeaderCache) {
402414
for k, v := range other {

relayer/strategies.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,7 @@ func StartRelayer(
4949
}
5050

5151
ePaths := make([]path, len(paths))
52-
for _, np := range paths {
52+
for i, np := range paths {
5353
pathName := np.Name
5454
p := np.Path
5555

@@ -62,10 +62,10 @@ func StartRelayer(
6262
filterSrc = append(filterSrc, ruleSrc)
6363
filterDst = append(filterDst, ruleDst)
6464
}
65-
ePaths = append(ePaths, path{
65+
ePaths[i] = path{
6666
src: processor.NewPathEnd(pathName, p.Src.ChainID, p.Src.ClientID, filter.Rule, filterSrc),
6767
dst: processor.NewPathEnd(pathName, p.Dst.ChainID, p.Dst.ClientID, filter.Rule, filterDst),
68-
})
68+
}
6969
}
7070

7171
go relayerStartEventProcessor(ctx, log, chainProcessors, ePaths, initialBlockHistory, maxTxSize, maxMsgLength, memo, errorChan, metrics)

0 commit comments

Comments
 (0)