go-ethereum/consensus/bor/heimdall/client.go
Arpit Temani b48b89f852
all: implement milestones (#961)
* Milestone Implementation

* Merge branch 'POS-347' into reciept-e2e-test

* Changes for testing, will be removed after testing

* Debugged the error

* Just for testing purpose

* refactor debug api methods, rename whitelist -> checkpoint

* remove first iteration based vars

* fix linters

* Rewind Changes

* Error changes

* RewindBack function in bor_checkpoint_verifier

* Testcases

* Added the fetch test for milestone and checkpoint

* Debugged the lint changes

* Debugged the lint changes

* Debugged the lint changes

* Debugged the lint changes

* Improved the error in Miner test file

* Improved the error of pointing to the wrong function

* Locking the sprint after the vote has been made on it.

* Adding more logs for testing

* Adding more logs for testing

* Adding more logs for testing

* Implemented the NoAckMilestone fetching mechanism

* Testcases for milestone implementation

* Testing code  for fetchNoAckMilestone and fetchLastNoAckMilestone

* Testing changes

* refactor else-if

* Corrected the number of params in bor_ext.go

* Dummy API for testing

* Defined the GetVoteOnRootHash in interface

* Defined the GetVoteOnRootHash in interface

* Made changes in the web3ext file

* Added the GetVoteOnRootHash in PublicBlockChain API

* Added the GetVoteOnRootHash in filterBackend

* Added the log of Root and RooHash

* Removed the 0x from rootHash

* Just for testing purpose

* "GetVoteOnRootHash" mock implementation

* bor_test.go

* Added the test for milestone implementation

* Added service for fetching milestone by ID

* Improved the comments

* Removed the duplicate code

* use setter for borVerifier

* use setter for borVerifier

* refactor handleNoAckMilestone

* remove code repetition with retry function

* Converged the repetitive code

* after CR

* persistence

* persistence implementation

* feature flag

* Persistence Changes

* cr

* initial

* fix

* fix

* Whitelist Flag

* 1 Add:Included the milestone flag  2.Add:Hardlimit the rewind to maximum of 255 blocks

* Chg:Updated go.mod file

* Remove:Dubai Hardfork code

* Add:checked errors for call functions to the Db, Rmv: Remote Header variable from the IsValidPeer() function

* Fix:Linting issues'

* Add:MilestoneGRPC functions

* Fix:Lint issues

* Fix:Lint issues

* Fix: TestFetchMilestoneFromMockHeimdall

* Fix:Integrations tests

* Add:Test for sprint length and milestone changes

* Add:Functionality to fetch the finalized block

* Chg:Changed default val of TriesInmemory to 1024

* fix:Some functions of heimdallGRPC client

* Restored the GRPC functionality, was commented out for  developing purpose

* Fix:Bor_checkpoint_Verfier function

* Test:Added the chain Rewinding test

* Test:Added the Sprint Length + Milestone merge test

* Add:Implemented the future milestone

* Add:Future milestone changes

* Add:Future milestone changes

* Chg: Voting on endBlockHash rather than rootHash

* Chg: Changed the logic of future milestone from rootHash checking to endBlockHash checking

* Fix:Using endBockHash while verifying the incoming milestone

* Chg:Variable names for better readiblity

* Fix:Testing changes

* Add:metrics for milestone implementation

* Add:Metrics for milestone implementatian

* Fix:Order of statements in a function for better optimization

* Chg:Removed unrequired file

* Fix:new variable intialization

* Add:Comment to increase readiblity

* Fix:Logs

* Chg:Name of GetVoteOnRootHash to GetVoteOnHash

* Fix:Linting issues

* Fixed linting issues

* Rmv: Unnecessary logs and Add:Skip test for long tests

* Fix:Checking current chain with whitelisted milestone or checkpoint in Finalized block function

* Fix:Test

* Fix:Whitelisting of Milestone and Checkpoint process

* Fix: Milestone JSON structure

* Chg:Testcases changes

* Fix:Change from VoteOnRootHash to VoteOnHash

* Fix:Variable name fix

* Fix:Finalized API

* internal/jsre/deps: update web3.js bundle

* Fix:milestone verifier

* Chg:Handling the long future chain import issue

* Fix:Lint issues

* Fix:TestLowDiffLongChain and TestPrunedImportSide tests, used hardcoded value 128 instead of DefaultTriesInMemory value

* Chg:Testcode for producing metrics

* Chg:Milestong polling value to 32 secs

* Add:Testcases

* Add:Implemented the check to fetch the milestoneId from heimdall before locking the fork

* Added GRPC method for FetchMilestoneID

* Fix:lint issue

* Fix:lint issue

* Skiped out the tests which were mainly used to produce the supporting data

* remove vcs build when running snyk

* Add:Improved the logs and comments

* fix linters

* Skipped some test as they are panic due to timeout issue in github

* Chg:Variable name LockerSprintNumber to LockedMilestoneNumber for better readablity and clarity

* Chg:Conflicting variable names in milestone test file

* Chg:Conflicting function names in milestone test file

* fix : minor fix in TestInsertingSpanSizeBlocks

* Fix:Mocking issue in TestInsertingSpanSizeBlocks

* Fix:GRPC Polyproto Version

* eth/downloader: skip peer drop due to whitelisting err

* eth, tests/bor: bug fixes and minor refactor

* Add:Implemented the milestone related functions in the HeimdallApp

* Fix:Lint Errors & Remove:Redundant Code

* Fix:Testing Errors

* Fix:Bor integeration tests

* Fix:Test errors

* update heimdall client mock files

* remove unused arguments

* remove redundant code

* Chg:Changed the milestone polling intervals

* Add: added block finality from whitelisted checkpoint

* skip future chain validation

* Add:confirmation check of 16 blocks over the end block while voting for the milestone in GetVoteHash() function

* Chg:Included endBlockNum in UnlockMutex function

* Add:Property based test for milestone

* Fix:Opening the lock while processing future milestone

* Add:Property based test for futureMilestone

* Defined the value of TempTriesInMemory

* Fixed the finalized api

* Fixed lint issues

* eth: add logs while fetching and rewinding

* fix linters: use default returns instead of recursive calls

* Fix:Milestone intergration test

* Add:GetVoteHash fn in mock backend

* tests/bor: fix mock span

* tests/bor: remove t.Parallel()

* use bor namespace in ethclient, fix mock function

---------

Co-authored-by: Vaibhav Jindal <vaibhavjindal29@gmail.com>
Co-authored-by: VaibhavJindal <74560896+VAIBHAVJINDAL3012@users.noreply.github.com>
Co-authored-by: Manav Darji <manavdarji.india@gmail.com>
Co-authored-by: Evgeny Danienko <6655321@bk.ru>
Co-authored-by: Shivam Sharma <shivam691999@gmail.com>
Co-authored-by: Anshal Shukla <shukla.anshal85@gmail.com>
2023-08-28 18:42:21 +05:30

460 lines
11 KiB
Go

package heimdall
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"sort"
"time"
"github.com/ethereum/go-ethereum/consensus/bor/clerk"
"github.com/ethereum/go-ethereum/consensus/bor/heimdall/checkpoint"
"github.com/ethereum/go-ethereum/consensus/bor/heimdall/milestone"
"github.com/ethereum/go-ethereum/consensus/bor/heimdall/span"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/metrics"
)
var (
// ErrShutdownDetected is returned if a shutdown was detected
ErrShutdownDetected = errors.New("shutdown detected")
ErrNoResponse = errors.New("got a nil response")
ErrNotSuccessfulResponse = errors.New("error while fetching data from Heimdall")
ErrNotInRejectedList = errors.New("milestoneID doesn't exist in rejected list")
ErrNotInMilestoneList = errors.New("milestoneID doesn't exist in Heimdall")
)
const (
stateFetchLimit = 50
apiHeimdallTimeout = 5 * time.Second
retryCall = 5 * time.Second
)
type StateSyncEventsResponse struct {
Height string `json:"height"`
Result []*clerk.EventRecordWithTime `json:"result"`
}
type SpanResponse struct {
Height string `json:"height"`
Result span.HeimdallSpan `json:"result"`
}
type HeimdallClient struct {
urlString string
client http.Client
closeCh chan struct{}
}
type Request struct {
client http.Client
url *url.URL
start time.Time
}
func NewHeimdallClient(urlString string) *HeimdallClient {
return &HeimdallClient{
urlString: urlString,
client: http.Client{
Timeout: apiHeimdallTimeout,
},
closeCh: make(chan struct{}),
}
}
const (
fetchStateSyncEventsFormat = "from-id=%d&to-time=%d&limit=%d"
fetchStateSyncEventsPath = "clerk/event-record/list"
fetchCheckpoint = "/checkpoints/%s"
fetchCheckpointCount = "/checkpoints/count"
fetchMilestone = "/milestone/latest"
fetchMilestoneCount = "/milestone/count"
fetchLastNoAckMilestone = "/milestone/lastNoAck"
fetchNoAckMilestone = "/milestone/noAck/%s"
fetchMilestoneID = "/milestone/ID/%s"
fetchSpanFormat = "bor/span/%d"
)
func (h *HeimdallClient) StateSyncEvents(ctx context.Context, fromID uint64, to int64) ([]*clerk.EventRecordWithTime, error) {
eventRecords := make([]*clerk.EventRecordWithTime, 0)
for {
url, err := stateSyncURL(h.urlString, fromID, to)
if err != nil {
return nil, err
}
log.Info("Fetching state sync events", "queryParams", url.RawQuery)
ctx = withRequestType(ctx, stateSyncRequest)
response, err := FetchWithRetry[StateSyncEventsResponse](ctx, h.client, url, h.closeCh)
if err != nil {
return nil, err
}
if response == nil || response.Result == nil {
// status 204
break
}
eventRecords = append(eventRecords, response.Result...)
if len(response.Result) < stateFetchLimit {
break
}
fromID += uint64(stateFetchLimit)
}
sort.SliceStable(eventRecords, func(i, j int) bool {
return eventRecords[i].ID < eventRecords[j].ID
})
return eventRecords, nil
}
func (h *HeimdallClient) Span(ctx context.Context, spanID uint64) (*span.HeimdallSpan, error) {
url, err := spanURL(h.urlString, spanID)
if err != nil {
return nil, err
}
ctx = withRequestType(ctx, spanRequest)
response, err := FetchWithRetry[SpanResponse](ctx, h.client, url, h.closeCh)
if err != nil {
return nil, err
}
return &response.Result, nil
}
// FetchCheckpoint fetches the checkpoint from heimdall
func (h *HeimdallClient) FetchCheckpoint(ctx context.Context, number int64) (*checkpoint.Checkpoint, error) {
url, err := checkpointURL(h.urlString, number)
if err != nil {
return nil, err
}
ctx = withRequestType(ctx, checkpointRequest)
response, err := FetchWithRetry[checkpoint.CheckpointResponse](ctx, h.client, url, h.closeCh)
if err != nil {
return nil, err
}
return &response.Result, nil
}
// FetchMilestone fetches the checkpoint from heimdall
func (h *HeimdallClient) FetchMilestone(ctx context.Context) (*milestone.Milestone, error) {
url, err := milestoneURL(h.urlString)
if err != nil {
return nil, err
}
ctx = withRequestType(ctx, milestoneRequest)
response, err := FetchWithRetry[milestone.MilestoneResponse](ctx, h.client, url, h.closeCh)
if err != nil {
return nil, err
}
return &response.Result, nil
}
// FetchCheckpointCount fetches the checkpoint count from heimdall
func (h *HeimdallClient) FetchCheckpointCount(ctx context.Context) (int64, error) {
url, err := checkpointCountURL(h.urlString)
if err != nil {
return 0, err
}
ctx = withRequestType(ctx, checkpointCountRequest)
response, err := FetchWithRetry[checkpoint.CheckpointCountResponse](ctx, h.client, url, h.closeCh)
if err != nil {
return 0, err
}
return response.Result.Result, nil
}
// FetchMilestoneCount fetches the milestone count from heimdall
func (h *HeimdallClient) FetchMilestoneCount(ctx context.Context) (int64, error) {
url, err := milestoneCountURL(h.urlString)
if err != nil {
return 0, err
}
ctx = withRequestType(ctx, milestoneCountRequest)
response, err := FetchWithRetry[milestone.MilestoneCountResponse](ctx, h.client, url, h.closeCh)
if err != nil {
return 0, err
}
return response.Result.Count, nil
}
// FetchLastNoAckMilestone fetches the last no-ack-milestone from heimdall
func (h *HeimdallClient) FetchLastNoAckMilestone(ctx context.Context) (string, error) {
url, err := lastNoAckMilestoneURL(h.urlString)
if err != nil {
return "", err
}
ctx = withRequestType(ctx, milestoneLastNoAckRequest)
response, err := FetchWithRetry[milestone.MilestoneLastNoAckResponse](ctx, h.client, url, h.closeCh)
if err != nil {
return "", err
}
return response.Result.Result, nil
}
// FetchNoAckMilestone fetches the last no-ack-milestone from heimdall
func (h *HeimdallClient) FetchNoAckMilestone(ctx context.Context, milestoneID string) error {
url, err := noAckMilestoneURL(h.urlString, milestoneID)
if err != nil {
return err
}
ctx = withRequestType(ctx, milestoneNoAckRequest)
response, err := FetchWithRetry[milestone.MilestoneNoAckResponse](ctx, h.client, url, h.closeCh)
if err != nil {
return err
}
if !response.Result.Result {
return fmt.Errorf("%w: milestoneID %q", ErrNotInRejectedList, milestoneID)
}
return nil
}
// FetchMilestoneID fetches the bool result from Heimdal whether the ID corresponding
// to the given milestone is in process in Heimdall
func (h *HeimdallClient) FetchMilestoneID(ctx context.Context, milestoneID string) error {
url, err := milestoneIDURL(h.urlString, milestoneID)
if err != nil {
return err
}
ctx = withRequestType(ctx, milestoneIDRequest)
response, err := FetchWithRetry[milestone.MilestoneIDResponse](ctx, h.client, url, h.closeCh)
if err != nil {
return err
}
if !response.Result.Result {
return fmt.Errorf("%w: milestoneID %q", ErrNotInMilestoneList, milestoneID)
}
return nil
}
// FetchWithRetry returns data from heimdall with retry
func FetchWithRetry[T any](ctx context.Context, client http.Client, url *url.URL, closeCh chan struct{}) (*T, error) {
// request data once
request := &Request{client: client, url: url, start: time.Now()}
result, err := Fetch[T](ctx, request)
if err == nil {
return result, nil
}
// attempt counter
attempt := 1
log.Warn("an error while trying fetching from Heimdall", "attempt", attempt, "error", err)
// create a new ticker for retrying the request
ticker := time.NewTicker(retryCall)
defer ticker.Stop()
const logEach = 5
retryLoop:
for {
log.Info("Retrying again in 5 seconds to fetch data from Heimdall", "path", url.Path, "attempt", attempt)
attempt++
select {
case <-ctx.Done():
log.Debug("Shutdown detected, terminating request by context.Done")
return nil, ctx.Err()
case <-closeCh:
log.Debug("Shutdown detected, terminating request by closing")
return nil, ErrShutdownDetected
case <-ticker.C:
request = &Request{client: client, url: url, start: time.Now()}
result, err = Fetch[T](ctx, request)
if err != nil {
if attempt%logEach == 0 {
log.Warn("an error while trying fetching from Heimdall", "attempt", attempt, "error", err)
}
continue retryLoop
}
return result, nil
}
}
}
// Fetch returns data from heimdall
func Fetch[T any](ctx context.Context, request *Request) (*T, error) {
isSuccessful := false
defer func() {
if metrics.EnabledExpensive {
sendMetrics(ctx, request.start, isSuccessful)
}
}()
result := new(T)
body, err := internalFetchWithTimeout(ctx, request.client, request.url)
if err != nil {
return nil, err
}
if body == nil {
return nil, ErrNoResponse
}
err = json.Unmarshal(body, result)
if err != nil {
return nil, err
}
isSuccessful = true
return result, nil
}
func spanURL(urlString string, spanID uint64) (*url.URL, error) {
return makeURL(urlString, fmt.Sprintf(fetchSpanFormat, spanID), "")
}
func stateSyncURL(urlString string, fromID uint64, to int64) (*url.URL, error) {
queryParams := fmt.Sprintf(fetchStateSyncEventsFormat, fromID, to, stateFetchLimit)
return makeURL(urlString, fetchStateSyncEventsPath, queryParams)
}
func checkpointURL(urlString string, number int64) (*url.URL, error) {
url := ""
if number == -1 {
url = fmt.Sprintf(fetchCheckpoint, "latest")
} else {
url = fmt.Sprintf(fetchCheckpoint, fmt.Sprint(number))
}
return makeURL(urlString, url, "")
}
func milestoneURL(urlString string) (*url.URL, error) {
url := fetchMilestone
return makeURL(urlString, url, "")
}
func checkpointCountURL(urlString string) (*url.URL, error) {
return makeURL(urlString, fetchCheckpointCount, "")
}
func milestoneCountURL(urlString string) (*url.URL, error) {
return makeURL(urlString, fetchMilestoneCount, "")
}
func lastNoAckMilestoneURL(urlString string) (*url.URL, error) {
return makeURL(urlString, fetchLastNoAckMilestone, "")
}
func noAckMilestoneURL(urlString string, id string) (*url.URL, error) {
url := fmt.Sprintf(fetchNoAckMilestone, id)
return makeURL(urlString, url, "")
}
func milestoneIDURL(urlString string, id string) (*url.URL, error) {
url := fmt.Sprintf(fetchMilestoneID, id)
return makeURL(urlString, url, "")
}
func makeURL(urlString, rawPath, rawQuery string) (*url.URL, error) {
u, err := url.Parse(urlString)
if err != nil {
return nil, err
}
u.Path = rawPath
u.RawQuery = rawQuery
return u, err
}
// internal fetch method
func internalFetch(ctx context.Context, client http.Client, u *url.URL) ([]byte, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u.String(), nil)
if err != nil {
return nil, err
}
res, err := client.Do(req)
if err != nil {
return nil, err
}
defer res.Body.Close()
// check status code
if res.StatusCode != 200 && res.StatusCode != 204 {
return nil, fmt.Errorf("%w: response code %d", ErrNotSuccessfulResponse, res.StatusCode)
}
// unmarshall data from buffer
if res.StatusCode == 204 {
return nil, nil
}
// get response
body, err := io.ReadAll(res.Body)
if err != nil {
return nil, err
}
return body, nil
}
func internalFetchWithTimeout(ctx context.Context, client http.Client, url *url.URL) ([]byte, error) {
ctx, cancel := context.WithTimeout(ctx, apiHeimdallTimeout)
defer cancel()
// request data once
return internalFetch(ctx, client, url)
}
// Close sends a signal to stop the running process
func (h *HeimdallClient) Close() {
close(h.closeCh)
h.client.CloseIdleConnections()
}