mirror of
https://github.com/kaspanet/kaspad.git
synced 2025-06-24 23:12:31 +00:00

* [NOD-1125] Write a skeleton for starting IBD. * [NOD-1125] Add WaitForIBDStart to Peer. * [NOD-1125] Move functions around. * [NOD-1125] Fix merge errors. * [NOD-1125] Fix a comment. * [NOD-1125] Implement sendGetBlockLocator. * [NOD-1125] Begin implementing findIBDLowHash. * [NOD-1125] Finish implementing findIBDLowHash. * [NOD-1125] Rename findIBDLowHash to findHighestSharedBlockHash. * [NOD-1125] Implement downloadBlocks. * [NOD-1125] Implement msgIBDBlock. * [NOD-1125] Implement msgIBDBlock. * [NOD-1125] Fix message types for HandleIBD. * [NOD-1125] Write a skeleton for requesting selected tip hashes. * [NOD-1125] Write a skeleton for the rest of the IBD requests. * [NOD-1125] Implement HandleGetBlockLocator. * [NOD-1125] Fix wrong timeout. * [NOD-1125] Fix compilation error. * [NOD-1125] Implement HandleGetBlocks. * [NOD-1125] Fix compilation errors. * [NOD-1125] Fix merge errors. * [NOD-1125] Implement selectPeerForIBD. * [NOD-1125] Implement RequestSelectedTip. * [NOD-1125] Implement HandleGetSelectedTip. * [NOD-1125] Make go lint happy. * [NOD-1125] Add minGetSelectedTipInterval. * [NOD-1125] Call StartIBDIfRequired where needed. * [NOD-1125] Fix merge errors. * [NOD-1125] Remove a redundant line. * [NOD-1125] Rename shouldContinue to shouldStop. * [NOD-1125] Lowercasify an error message. * [NOD-1125] Shuffle statements around in findHighestSharedBlockHash. * [NOD-1125] Rename hasRecentlyReceivedBlock to isDAGTimeCurrent. * [NOD-1125] Scope minGetSelectedTipInterval. * [NOD-1125] Handle an unhandled error. * [NOD-1125] Use AddUint32 instead of LoadUint32 + StoreUint32. * [NOD-1125] Use AddUint32 instead of LoadUint32 + StoreUint32. * [NOD-1125] Use SwapUint32 instead of AddUint32. * [NOD-1125] Remove error from requestSelectedTips. * [NOD-1125] Actually stop IBD when it should stop. * [NOD-1125] Actually stop RequestSelectedTip when it should stop. * [NOD-1125] Don't ban peers that send us delayed blocks during IBD. * [NOD-1125] Make unexpected message type messages nicer. * [NOD-1125] Remove Peer.ready and make HandleHandshake return it to guarantee we never operate on a non-initialized peer. * [NOD-1125] Remove errors associated with Peer.ready. * [NOD-1125] Extract maxHashesInMsgIBDBlocks to a const. * [NOD-1125] Move the ibd package into flows. * [NOD-1125] Start IBD if required after getting an unknown block inv. * [NOD-1125] Don't request blocks during relay if we're in the middle of IBD. * [NOD-1125] Remove AddBlockLocatorHash. * [NOD-1125] Extract runIBD to a seperate function. * [NOD-1125] Extract runSelectedTipRequest to a seperate function. * [NOD-1125] Remove EnqueueWithTimeout. * [NOD-1125] Increase the capacity of the outgoingRoute. * [NOD-1125] Fix some bad names. * [NOD-1125] Fix a comment. * [NOD-1125] Simplify a comment. * [NOD-1125] Move WaitFor... functions into their respective run... functions. * [NOD-1125] Return default values in case of error. * [NOD-1125] Use CmdXXX in error messages. * [NOD-1125] Use MaxInvPerMsg in outgoingRouteMaxMessages instead of MaxBlockLocatorsPerMsg. * [NOD-1125] Fix a comment. * [NOD-1125] Disconnect a peer that sends us a delayed block during IBD. * [NOD-1125] Use StoreUint32 instead of SwapUint32. * [NOD-1125] Add a comment. * [NOD-1125] Don't ban peers that send us delayed blocks.
99 lines
2.9 KiB
Go
99 lines
2.9 KiB
Go
package router
|
|
|
|
import (
|
|
"github.com/kaspanet/kaspad/wire"
|
|
"github.com/pkg/errors"
|
|
)
|
|
|
|
const outgoingRouteMaxMessages = wire.MaxInvPerMsg + defaultMaxMessages
|
|
|
|
// OnRouteCapacityReachedHandler is a function that is to
|
|
// be called when one of the routes reaches capacity.
|
|
type OnRouteCapacityReachedHandler func()
|
|
|
|
// Router routes messages by type to their respective
|
|
// input channels
|
|
type Router struct {
|
|
incomingRoutes map[wire.MessageCommand]*Route
|
|
outgoingRoute *Route
|
|
|
|
onRouteCapacityReachedHandler OnRouteCapacityReachedHandler
|
|
}
|
|
|
|
// NewRouter creates a new empty router
|
|
func NewRouter() *Router {
|
|
router := Router{
|
|
incomingRoutes: make(map[wire.MessageCommand]*Route),
|
|
outgoingRoute: newRouteWithCapacity(outgoingRouteMaxMessages),
|
|
}
|
|
router.outgoingRoute.setOnCapacityReachedHandler(func() {
|
|
router.onRouteCapacityReachedHandler()
|
|
})
|
|
return &router
|
|
}
|
|
|
|
// SetOnRouteCapacityReachedHandler sets the onRouteCapacityReachedHandler
|
|
// function for this router
|
|
func (r *Router) SetOnRouteCapacityReachedHandler(onRouteCapacityReachedHandler OnRouteCapacityReachedHandler) {
|
|
r.onRouteCapacityReachedHandler = onRouteCapacityReachedHandler
|
|
}
|
|
|
|
// AddIncomingRoute registers the messages of types `messageTypes` to
|
|
// be routed to the given `route`
|
|
func (r *Router) AddIncomingRoute(messageTypes []wire.MessageCommand) (*Route, error) {
|
|
route := NewRoute()
|
|
for _, messageType := range messageTypes {
|
|
if _, ok := r.incomingRoutes[messageType]; ok {
|
|
return nil, errors.Errorf("a route for '%s' already exists", messageType)
|
|
}
|
|
r.incomingRoutes[messageType] = route
|
|
}
|
|
route.setOnCapacityReachedHandler(func() {
|
|
r.onRouteCapacityReachedHandler()
|
|
})
|
|
return route, nil
|
|
}
|
|
|
|
// RemoveRoute unregisters the messages of types `messageTypes` from
|
|
// the router
|
|
func (r *Router) RemoveRoute(messageTypes []wire.MessageCommand) error {
|
|
for _, messageType := range messageTypes {
|
|
if _, ok := r.incomingRoutes[messageType]; !ok {
|
|
return errors.Errorf("a route for '%s' does not exist", messageType)
|
|
}
|
|
delete(r.incomingRoutes, messageType)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// EnqueueIncomingMessage enqueues the given message to the
|
|
// appropriate route
|
|
func (r *Router) EnqueueIncomingMessage(message wire.Message) (isOpen bool, err error) {
|
|
route, ok := r.incomingRoutes[message.Command()]
|
|
if !ok {
|
|
return false, errors.Errorf("a route for '%s' does not exist", message.Command())
|
|
}
|
|
return route.Enqueue(message), nil
|
|
}
|
|
|
|
// OutgoingRoute returns the outgoing route
|
|
func (r *Router) OutgoingRoute() *Route {
|
|
return r.outgoingRoute
|
|
}
|
|
|
|
// Close shuts down the router by closing all registered
|
|
// incoming routes and the outgoing route
|
|
func (r *Router) Close() error {
|
|
incomingRoutes := make(map[*Route]struct{})
|
|
for _, route := range r.incomingRoutes {
|
|
incomingRoutes[route] = struct{}{}
|
|
}
|
|
for route := range incomingRoutes {
|
|
err := route.Close()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return r.outgoingRoute.Close()
|
|
}
|