-
Notifications
You must be signed in to change notification settings - Fork 504
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
[Memory fix] Remove memory leak in event subscriptions (#741)
* Remove dangling pointers in the subscription struct * Remove leftover log
- Loading branch information
1 parent
68df30d
commit 52b6261
Showing
2 changed files
with
50 additions
and
105 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,76 +1,53 @@ | ||
package blockchain | ||
|
||
import ( | ||
"sync" | ||
"testing" | ||
"time" | ||
|
||
"github.com/0xPolygon/polygon-edge/types" | ||
"github.com/stretchr/testify/assert" | ||
) | ||
|
||
func TestSubscriptionLinear(t *testing.T) { | ||
e := &eventStream{} | ||
func TestSubscription(t *testing.T) { | ||
t.Parallel() | ||
|
||
// add a genesis block to eventstream | ||
e.push(&Event{ | ||
NewChain: []*types.Header{ | ||
{Number: 0}, | ||
}, | ||
}) | ||
var ( | ||
e = &eventStream{} | ||
sub = e.subscribe() | ||
caughtEventNum = uint64(0) | ||
event = &Event{ | ||
NewChain: []*types.Header{ | ||
{ | ||
Number: 100, | ||
}, | ||
}, | ||
} | ||
|
||
sub := e.subscribe() | ||
wg sync.WaitGroup | ||
) | ||
|
||
eventCh := make(chan *Event) | ||
defer sub.Close() | ||
|
||
go func() { | ||
for { | ||
task := sub.GetEvent() | ||
eventCh <- task | ||
} | ||
}() | ||
updateCh := sub.GetEventCh() | ||
|
||
for i := 1; i < 10; i++ { | ||
evnt := &Event{} | ||
wg.Add(1) | ||
|
||
evnt.AddNewHeader(&types.Header{Number: uint64(i)}) | ||
e.push(evnt) | ||
go func() { | ||
defer wg.Done() | ||
|
||
// it should fire updateCh | ||
select { | ||
case evnt := <-eventCh: | ||
if evnt.NewChain[0].Number != uint64(i) { | ||
t.Fatal("bad") | ||
} | ||
case <-time.After(1 * time.Second): | ||
t.Fatal("timeout") | ||
case ev := <-updateCh: | ||
caughtEventNum = ev.NewChain[0].Number | ||
case <-time.After(5 * time.Second): | ||
} | ||
} | ||
} | ||
|
||
func TestSubscriptionSlowConsumer(t *testing.T) { | ||
e := &eventStream{} | ||
}() | ||
|
||
e.push(&Event{ | ||
NewChain: []*types.Header{ | ||
{Number: 0}, | ||
}, | ||
}) | ||
// Send the event to the channel | ||
e.push(event) | ||
|
||
sub := e.subscribe() | ||
// Wait for the event to be parsed | ||
wg.Wait() | ||
|
||
// send multiple events | ||
for i := 1; i < 10; i++ { | ||
e.push(&Event{ | ||
NewChain: []*types.Header{ | ||
{Number: uint64(i)}, | ||
}, | ||
}) | ||
} | ||
|
||
// consume events now | ||
for i := 1; i < 10; i++ { | ||
evnt := sub.GetEvent() | ||
if evnt.NewChain[0].Number != uint64(i) { | ||
t.Fatal("bad") | ||
} | ||
} | ||
assert.Equal(t, event.NewChain[0].Number, caughtEventNum) | ||
} |