Skip to content
Merged
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
15 changes: 12 additions & 3 deletions libs/chainconsensus/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -50,18 +50,20 @@ type handler struct {
unknownRequestsResultByID map[string]*unknownRequest
unknownRequestsOrderedByTimeout *list.List[*unknownRequest]
unknownRequestTTL time.Duration
maxUnknownRequestsCacheSize int
}

func NewHandler(lggr logger.Logger, poller Poller, metrics metrics.ConsensusMetrics, unknownRequestTTL time.Duration) Handler {
return newHandler(lggr, poller, metrics, unknownRequestTTL)
func NewHandler(lggr logger.Logger, poller Poller, metrics metrics.ConsensusMetrics, unknownRequestTTL time.Duration, maxUnknownRequestsCacheSize int) Handler {
return newHandler(lggr, poller, metrics, unknownRequestTTL, maxUnknownRequestsCacheSize)
}

func newHandler(lggr logger.Logger, poller Poller, metrics metrics.ConsensusMetrics, unknownRequestTTL time.Duration) *handler {
func newHandler(lggr logger.Logger, poller Poller, metrics metrics.ConsensusMetrics, unknownRequestTTL time.Duration, maxUnknownRequestsCacheSize int) *handler {
r := &handler{
requests: requests.NewStoreWithStatsCollector[*requestCtx](metrics),
unknownRequestsResultByID: make(map[string]*unknownRequest),
unknownRequestsOrderedByTimeout: list.New[*unknownRequest](),
unknownRequestTTL: unknownRequestTTL,
maxUnknownRequestsCacheSize: maxUnknownRequestsCacheSize,
poller: poller,
metrics: metrics,
}
Expand Down Expand Up @@ -210,6 +212,13 @@ func (s *handler) completeRequest(id string, reply types.Reply) error {
defer s.lock.Unlock()
request := s.requests.Get(id)
if request == nil {
if s.maxUnknownRequestsCacheSize > 0 && len(s.unknownRequestsResultByID) >= s.maxUnknownRequestsCacheSize {
s.lggr.Warnf("unknown requests cache is full, evicting oldest request")
if oldest := s.unknownRequestsOrderedByTimeout.Front(); oldest != nil {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It would be nice to have visibility into the size of the unknown requests cache.
I'm ok with logging on evictiond due to overflow or a metric

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

added log line

delete(s.unknownRequestsResultByID, oldest.Value.ID)
s.unknownRequestsOrderedByTimeout.Remove(oldest)
}
}
uRequest := &unknownRequest{
ID: id,
ExpiresAt: time.Now().Add(s.unknownRequestTTL),
Expand Down
39 changes: 35 additions & 4 deletions libs/chainconsensus/handler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ import (
func TestGetRequestIDs(t *testing.T) {
poller := mocks.NewPoller(t)
poller.EXPECT().Enqueue(mock.Anything, mock.Anything)
handler := NewHandler(logger.Test(t), poller, test.GetConsensusMetrics(t), time.Second)
handler := NewHandler(logger.Test(t), poller, test.GetConsensusMetrics(t), time.Second, 1000)
addRequestToHandler := func(t *testing.T, ctx context.Context, id string) {
request := types.NewEventuallyConsistentRequest(id, nil)
_, err := handler.Handle(ctx, request)
Expand Down Expand Up @@ -72,7 +72,7 @@ func TestGetRequestIDs(t *testing.T) {
func TestGetRequest(t *testing.T) {
poller := mocks.NewPoller(t)
poller.EXPECT().Enqueue(mock.Anything, mock.Anything).Maybe()
handler := NewHandler(logger.Test(t), poller, test.GetConsensusMetrics(t), time.Second)
handler := NewHandler(logger.Test(t), poller, test.GetConsensusMetrics(t), time.Second, 1000)
addRequestToHandler := func(t *testing.T, ctx context.Context, id string) {
request := types.NewAggregatableRequest(id, nil)
_, err := handler.Handle(ctx, request)
Expand Down Expand Up @@ -102,7 +102,7 @@ func TestGetRequest(t *testing.T) {

func TestCompleteRequest(t *testing.T) {
newHandler := func(t *testing.T, lggr logger.Logger, poller Poller) *handler {
h := newHandler(lggr, poller, test.GetConsensusMetrics(t), time.Second)
h := newHandler(lggr, poller, test.GetConsensusMetrics(t), time.Second, 1000)
require.NoError(t, h.Start(t.Context()))
t.Cleanup(func() {
require.NoError(t, h.Close())
Expand Down Expand Up @@ -276,7 +276,7 @@ func TestCompleteRequest(t *testing.T) {

func TestHandle(t *testing.T) {
poller := mocks.NewPoller(t)
handler := NewHandler(logger.Test(t), poller, test.GetConsensusMetrics(t), time.Second)
handler := NewHandler(logger.Test(t), poller, test.GetConsensusMetrics(t), time.Second, 1000)
require.NoError(t, handler.Start(t.Context()))
t.Cleanup(func() {
require.NoError(t, handler.Close())
Expand Down Expand Up @@ -309,3 +309,34 @@ func mustMarshalProto(t *testing.T, msg proto.Message) []byte {
require.NoError(t, err)
return data
}

func TestCompleteRequest_UnknownRequestsCacheEviction(t *testing.T) {
const maxCacheSize = 3
h := newHandler(logger.Test(t), nil, test.GetConsensusMetrics(t), time.Minute, maxCacheSize)

completeUnknown := func(t *testing.T, id string) {
require.NoError(t, h.CompleteProtoRequest(id, &types.RequestReport{
Report: &types.RequestReport_EventuallyConsistent{EventuallyConsistent: []byte(id)},
}))
}

completeUnknown(t, "req-1")
completeUnknown(t, "req-2")
completeUnknown(t, "req-3")

h.lock.RLock()
require.Len(t, h.unknownRequestsResultByID, maxCacheSize)
require.Contains(t, h.unknownRequestsResultByID, "req-1")
h.lock.RUnlock()

// cache is full; completing one more unknown request must evict the oldest ("req-1")
completeUnknown(t, "req-4")

h.lock.RLock()
require.Len(t, h.unknownRequestsResultByID, maxCacheSize)
require.NotContains(t, h.unknownRequestsResultByID, "req-1")
require.Contains(t, h.unknownRequestsResultByID, "req-2")
require.Contains(t, h.unknownRequestsResultByID, "req-3")
require.Contains(t, h.unknownRequestsResultByID, "req-4")
h.lock.RUnlock()
}
Loading