diff --git a/block/internal/syncing/da_retriever.go b/block/internal/syncing/da_retriever.go index 81de02b997..d0953901ab 100644 --- a/block/internal/syncing/da_retriever.go +++ b/block/internal/syncing/da_retriever.go @@ -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 { @@ -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 @@ -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) diff --git a/types/da.go b/types/da.go index 9335188a4a..e0d58710d9 100644 --- a/types/da.go +++ b/types/da.go @@ -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()) { @@ -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") diff --git a/types/da_test.go b/types/da_test.go index 59b842572b..4a111499dc 100644 --- a/types/da_test.go +++ b/types/da_test.go @@ -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) @@ -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) + }) +}