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
10 changes: 4 additions & 6 deletions block/internal/syncing/da_retriever.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,8 @@ import (
pb "github.com/evstack/ev-node/types/pb/evnode/v1"
)

const dAFetcherTimeout = 10 * time.Second
// defaultDATimeout is the default timeout for DA retrieval operations
const defaultDATimeout = 10 * time.Second

// DARetriever handles DA retrieval operations for syncing
type DARetriever struct {
Expand Down Expand Up @@ -64,9 +65,6 @@ func NewDARetriever(
// RetrieveFromDA retrieves blocks from the specified DA height and returns height events
func (r *DARetriever) RetrieveFromDA(ctx context.Context, daHeight uint64) ([]common.DAHeightEvent, error) {
r.logger.Debug().Uint64("da_height", daHeight).Msg("retrieving from DA")
ctx, cancel := context.WithTimeout(ctx, dAFetcherTimeout)
defer cancel()

blobsResp, err := r.fetchBlobs(ctx, daHeight)
if err != nil {
return nil, err
Expand All @@ -84,14 +82,14 @@ func (r *DARetriever) RetrieveFromDA(ctx context.Context, daHeight uint64) ([]co
// fetchBlobs retrieves blobs from the DA layer
func (r *DARetriever) fetchBlobs(ctx context.Context, daHeight uint64) (coreda.ResultRetrieve, error) {
// Retrieve from both namespaces
headerRes := types.RetrieveWithHelpers(ctx, r.da, r.logger, daHeight, r.namespaceBz)
headerRes := types.RetrieveWithHelpers(ctx, r.da, r.logger, daHeight, r.namespaceBz, defaultDATimeout)

// If namespaces are the same, return header result
if bytes.Equal(r.namespaceBz, r.namespaceDataBz) {
return headerRes, r.validateBlobResponse(headerRes, daHeight)
}

dataRes := types.RetrieveWithHelpers(ctx, r.da, r.logger, daHeight, r.namespaceDataBz)
dataRes := types.RetrieveWithHelpers(ctx, r.da, r.logger, daHeight, r.namespaceDataBz, defaultDATimeout)

// Validate responses
headerErr := r.validateBlobResponse(headerRes, daHeight)
Expand Down
10 changes: 8 additions & 2 deletions types/da.go
Original file line number Diff line number Diff line change
Expand Up @@ -114,15 +114,19 @@ func SubmitWithHelpers(
// RetrieveWithHelpers performs blob retrieval using the underlying DA layer,
// handling error mapping to produce a ResultRetrieve.
// It mimics the logic previously found in da.DAClient.Retrieve.
// requestTimeout defines the timeout for the each retrieval request.
func RetrieveWithHelpers(
ctx context.Context,
da coreda.DA,
logger zerolog.Logger,
dataLayerHeight uint64,
namespace []byte,
requestTimeout time.Duration,
) coreda.ResultRetrieve {
// 1. Get IDs
idsResult, err := da.GetIDs(ctx, dataLayerHeight, namespace)
getIDsCtx, cancel := context.WithTimeout(ctx, requestTimeout)
defer cancel()
idsResult, err := da.GetIDs(getIDsCtx, dataLayerHeight, namespace)
if err != nil {
// Handle specific "not found" error
if strings.Contains(err.Error(), coreda.ErrBlobNotFound.Error()) {
Expand Down Expand Up @@ -177,7 +181,9 @@ func RetrieveWithHelpers(
for i := 0; i < len(idsResult.IDs); i += batchSize {
end := min(i+batchSize, len(idsResult.IDs))

batchBlobs, err := da.Get(ctx, idsResult.IDs[i:end], namespace)
getBlobsCtx, cancel := context.WithTimeout(ctx, requestTimeout)
batchBlobs, err := da.Get(getBlobsCtx, idsResult.IDs[i:end], namespace)
cancel()
if err != nil {
// Handle errors during Get
logger.Error().Uint64("height", dataLayerHeight).Int("num_ids", len(idsResult.IDs)).Err(err).Msg("Retrieve helper: Failed to get blobs")
Expand Down
52 changes: 51 additions & 1 deletion types/da_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -229,7 +229,7 @@ func TestRetrieveWithHelpers(t *testing.T) {
mockDA.On("Get", mock.Anything, tc.getIDsResult.IDs, mock.Anything).Return(mockBlobs, tc.getBlobsErr)
}

result := types.RetrieveWithHelpers(context.Background(), mockDA, logger, dataLayerHeight, encodedNamespace.Bytes())
result := types.RetrieveWithHelpers(context.Background(), mockDA, logger, dataLayerHeight, encodedNamespace.Bytes(), 5*time.Second)

assert.Equal(t, tc.expectedCode, result.Code)
assert.Equal(t, tc.expectedHeight, result.Height)
Expand All @@ -246,3 +246,53 @@ func TestRetrieveWithHelpers(t *testing.T) {
})
}
}

func TestRetrieveWithHelpers_Timeout(t *testing.T) {
logger := zerolog.Nop()
dataLayerHeight := uint64(100)
encodedNamespace := coreda.NamespaceFromString("test-namespace")

t.Run("timeout during GetIDs", func(t *testing.T) {
mockDA := mocks.NewMockDA(t)

// Mock GetIDs to block until context is cancelled
mockDA.On("GetIDs", mock.Anything, dataLayerHeight, mock.Anything).Run(func(args mock.Arguments) {
ctx := args.Get(0).(context.Context)
<-ctx.Done() // Wait for context cancellation
}).Return(nil, context.DeadlineExceeded)

// Use a very short timeout to ensure it triggers
result := types.RetrieveWithHelpers(context.Background(), mockDA, logger, dataLayerHeight, encodedNamespace.Bytes(), 1*time.Millisecond)

assert.Equal(t, coreda.StatusError, result.Code)
assert.Contains(t, result.Message, "failed to get IDs")
assert.Contains(t, result.Message, "context deadline exceeded")
mockDA.AssertExpectations(t)
})

t.Run("timeout during Get", func(t *testing.T) {
mockDA := mocks.NewMockDA(t)
mockIDs := [][]byte{[]byte("id1")}
mockTimestamp := time.Now()

// Mock GetIDs to succeed
mockDA.On("GetIDs", mock.Anything, dataLayerHeight, mock.Anything).Return(&coreda.GetIDsResult{
IDs: mockIDs,
Timestamp: mockTimestamp,
}, nil)

// Mock Get to block until context is cancelled
mockDA.On("Get", mock.Anything, mockIDs, mock.Anything).Run(func(args mock.Arguments) {
ctx := args.Get(0).(context.Context)
<-ctx.Done() // Wait for context cancellation
}).Return(nil, context.DeadlineExceeded)

// Use a very short timeout to ensure it triggers
result := types.RetrieveWithHelpers(context.Background(), mockDA, logger, dataLayerHeight, encodedNamespace.Bytes(), 1*time.Millisecond)

assert.Equal(t, coreda.StatusError, result.Code)
assert.Contains(t, result.Message, "failed to get blobs for batch")
assert.Contains(t, result.Message, "context deadline exceeded")
mockDA.AssertExpectations(t)
})
}
Loading