Skip to content
Open
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
507 changes: 196 additions & 311 deletions .golangci-strict.yml

Large diffs are not rendered by default.

19 changes: 15 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,8 @@ Reply example:
"version": "0.0.0",
"state": 0,
"next_attempt": "0001-01-01T00:00:00Z",
"score": null
"score": null,
"last_seen": "2024-11-06T12:26:11Z"
}
]
```
Expand All @@ -72,6 +73,11 @@ Reply fields:
* `2` `Hostile` - successful connection was established, but the network byte or version is not acceptable.
* `next_attempt` - time when the next attempt to connect to the node is planned. This field is used during version selection (see the section [Version Selection](#version-selection)).
* `score` - last known value of the Score; `null` value if there was no connection.
* `last_seen` - time of the last activity of the node: the last successful handshake, Score update or disconnection. For a node that was never connected, it is the time when its address was discovered.

Legacy records without `last_seen`, including configured seeds, receive the current time during pruning after an upgrade.

Pruning runs at startup and every hour. Nodes that were not seen for 30 days are removed from storage, except configured seed peers and peers with active or pending connections. When other nodes request known peers, Fork Detector replies only with nodes that have connected successfully at least once, have a known listening port, and are connected at the moment or were seen within the last 96 hours (at most 1000 addresses).

### `GET` `/api/peers/friendly`

Expand All @@ -87,7 +93,8 @@ Reply example:
"version": "1.5.8",
"state": 1,
"next_attempt": "2024-11-06T12:36:11Z",
"score": 967128354984173501468049
"score": 967128354984173501468049,
"last_seen": "2024-11-06T12:26:11Z"
}
]
```
Expand All @@ -108,7 +115,8 @@ Reply example:
"version": "1.5.8",
"state": 1,
"next_attempt": "2024-11-06T12:36:11Z",
"score": 967129143986158086525130
"score": 967129143986158086525130,
"last_seen": "2024-11-06T12:26:11Z"
}
]
```
Expand Down Expand Up @@ -363,7 +371,10 @@ This command will start Fork Detector to work with the `MainNet` network in outg

Available command-line parameters:
* `-db` - path to the directory where the database will be stored. **This parameter is required.**
* `-log-level` - logging level, default is `INFO`. Other possible values: `DEBUG`, `INFO`, `WARN`, `ERROR`, and `FATAL`.
* `-log-level` - logging level, default is `info`. Other possible values: `debug`, `info`, `warn`, and `error`.
* `-log-type` - logger output format, default is `pretty`. Other possible values: `text` and `json`.
* `-log-network` - log the operation of network stack, turned off by default.
* `-log-network-data` - log network messages as Base64 strings, turned off by default.
* `-blockchain-type` - network type, available values are `mainnet`, `testnet`, and `stagenet`. Default is `mainnet`.
* `-peers` - a comma-separated list of node addresses to initially connect to in the format `ip:port,...,ip:port`. **This parameter is required.**
* `-api` - the address at which the Fork Detector API will be available. The default is `localhost:8080`, meaning the API will be available locally on port `8080`.
Expand Down
42 changes: 25 additions & 17 deletions api/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"errors"
"fmt"
"io/fs"
"log/slog"
"net/http"
"net/netip"
"runtime"
Expand All @@ -18,8 +19,8 @@ import (

"github.com/go-chi/chi/v5"
"github.com/go-chi/chi/v5/middleware"
"github.com/wavesplatform/gowaves/pkg/logging"
"github.com/wavesplatform/gowaves/pkg/proto"
"go.uber.org/zap"
"golang.org/x/sync/errgroup"

"github.com/alexeykiselev/waves-fork-detector/chains"
Expand All @@ -37,21 +38,21 @@ var (
// Logger is a middleware that logs the start and end of each request, along
// with some useful data about what was requested, what the response status was,
// and how long it took to return.
func Logger(l *zap.Logger) func(next http.Handler) http.Handler {
func Logger(l *slog.Logger) func(next http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
fn := func(w http.ResponseWriter, r *http.Request) {
ww := middleware.NewWrapResponseWriter(w, r.ProtoMajor)

t1 := time.Now()
defer func() {
l.Debug("Served",
zap.String("proto", r.Proto),
zap.String("path", r.URL.Path),
zap.String("remote", r.RemoteAddr),
zap.Duration("lat", time.Since(t1)),
zap.Int("status", ww.Status()),
zap.Int("size", ww.BytesWritten()),
zap.String("reqId", middleware.GetReqID(r.Context())))
slog.String("proto", r.Proto),
slog.String("path", r.URL.Path),
slog.String("remote", r.RemoteAddr),
slog.Duration("lat", time.Since(t1)),
slog.Int("status", ww.Status()),
slog.Int("size", ww.BytesWritten()),
slog.String("reqId", middleware.GetReqID(r.Context())))
}()

next.ServeHTTP(ww, r)
Expand All @@ -66,9 +67,12 @@ type API struct {
registry *peers.Registry
linkage *chains.Linkage
srv *http.Server
logger *slog.Logger
}

func NewAPI(registry *peers.Registry, linkage *chains.Linkage, bind string) (*API, error) {
func NewAPI(
registry *peers.Registry, linkage *chains.Linkage, bind string, logger *slog.Logger,
) (*API, error) {
if bind == "" {
return nil, errors.New("empty address to bin")
}
Expand All @@ -77,11 +81,10 @@ func NewAPI(registry *peers.Registry, linkage *chains.Linkage, bind string) (*AP
return nil, fmt.Errorf("failed to get swagger FS: %w", err)
}

a := API{registry: registry, linkage: linkage}
a := API{registry: registry, linkage: linkage, logger: logger}
r := chi.NewRouter()
r.Use(middleware.RequestID)
r.Use(middleware.RealIP)
r.Use(Logger(zap.L()))
r.Use(Logger(logger))
r.Use(middleware.Recoverer)
r.Use(middleware.Compress(flate.DefaultCompression))
const apiRoot = "/api"
Expand All @@ -98,20 +101,25 @@ func (a *API) Run(ctx context.Context) {
g.Go(a.runServer)
}

// Wait waits for the API server started by Run to stop and returns its error.
func (a *API) Wait() error {
return a.wait()
}

func (a *API) Shutdown() {
if err := a.srv.Shutdown(a.ctx); err != nil && !errors.Is(err, context.Canceled) {
zap.S().Errorf("Failed to shutdown API: %v", err)
a.logger.Error("Failed to shutdown API", logging.Error(err))
}
if err := a.wait(); err != nil {
zap.S().Warnf("Failed to shutdown API: %v", err)
a.logger.Warn("Failed to shutdown API", logging.Error(err))
}
zap.S().Info("API shutdown successfully")
a.logger.Info("API shutdown successfully")
}

func (a *API) runServer() error {
err := a.srv.ListenAndServe()
if err != nil && !errors.Is(err, http.ErrServerClosed) {
zap.S().Fatalf("Failed to start API: %v", err)
a.logger.Error("Failed to start API", logging.Error(err))
Comment thread
alexeykiselev marked this conversation as resolved.
return err
}
return nil
Expand Down
25 changes: 25 additions & 0 deletions api/api_internal_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
package api

import (
"log/slog"
"testing"
"time"

"github.com/stretchr/testify/require"
)

func TestWaitAfterServerClose(t *testing.T) {
a, err := NewAPI(nil, nil, "127.0.0.1:0", slog.New(slog.DiscardHandler))
require.NoError(t, err)
a.Run(t.Context())
require.NoError(t, a.srv.Close())
done := make(chan error, 1)
go func() { done <- a.Wait() }()
select {
case waitErr := <-done:
require.NoError(t, waitErr)
case <-time.After(5 * time.Second):
t.Fatal("API did not stop after closing the server")
}
a.Shutdown()
}
10 changes: 10 additions & 0 deletions api/swagger/openapi.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -252,6 +252,16 @@ components:
nullable: true
example: null
description: Last known Score value; `null` if there was no connection.
last_seen:
type: string
format: date-time
example: "2024-11-06T12:26:11Z"
description: |-
Time of the last activity of the node: the last successful handshake, Score update or disconnection.
For a node that was never connected, it is the time when its address was discovered.
Legacy records without this field, including configured seeds, receive the current time during pruning after an upgrade.
Only nodes connected at the moment or seen within the last 96 hours are advertised to other nodes.
Nodes not seen for 30 days are removed from storage, except the configured seed peers.

Head:
type: object
Expand Down
24 changes: 14 additions & 10 deletions chains/linkage.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package chains
import (
"errors"
"fmt"
"log/slog"
"math"
"math/big"
"net/netip"
Expand All @@ -12,8 +13,8 @@ import (
"time"

"github.com/syndtr/goleveldb/leveldb"
"github.com/wavesplatform/gowaves/pkg/logging"
"github.com/wavesplatform/gowaves/pkg/proto"
"go.uber.org/zap"
)

const (
Expand Down Expand Up @@ -44,9 +45,11 @@ type Linkage struct {

mu *sync.RWMutex
st *storage

logger *slog.Logger
}

func NewLinkage(path string, scheme proto.Scheme, genesis proto.Block) (*Linkage, error) {
func NewLinkage(path string, scheme proto.Scheme, genesis proto.Block, logger *slog.Logger) (*Linkage, error) {
st, err := newStorage(path, scheme)
if err != nil {
return nil, fmt.Errorf("failed to initialize Linkage: %w", err)
Expand All @@ -59,13 +62,14 @@ func NewLinkage(path string, scheme proto.Scheme, genesis proto.Block) (*Linkage
genesis: genesis.BlockID(),
mu: &sync.RWMutex{},
st: st,
logger: logger,
}, nil
}

func (l *Linkage) Close() {
err := l.st.close()
if err != nil {
zap.S().Errorf("Failed to close Linkage: %v", err)
l.logger.Error("Failed to close Linkage", logging.Error(err))
}
}

Expand Down Expand Up @@ -234,26 +238,26 @@ func (l *Linkage) LogInitialStats() {

heads, err := l.activeHeads()
if err != nil {
zap.S().Errorf("Failed to log statistics: %v", err)
l.logger.Error("Failed to log statistics", logging.Error(err))
return
}
zap.S().Infof("Heads count in storage: %d", len(heads))
l.logger.Info("Heads in storage", slog.Int("count", len(heads)))
for _, head := range heads {
b, blErr := l.Block(head.BlockID)
if blErr != nil {
zap.S().Errorf("Failed to get block: %v", blErr)
l.logger.Error("Failed to get block", logging.Error(blErr))
return
}
zap.S().Infof("\tHead '%s' at height %d", head.BlockID.String(), b.Height)
l.logger.Info("Head", slog.String("block", head.BlockID.String()), slog.Uint64("height", uint64(b.Height)))
}
leashes, err := l.st.leashes()
if err != nil {
zap.S().Errorf("Failed to log statistics: %v", err)
l.logger.Error("Failed to log statistics", logging.Error(err))
return
}
zap.S().Infof("Leashes count in storage: %d", len(leashes))
l.logger.Info("Leashes in storage", slog.Int("count", len(leashes)))
for _, lsh := range leashes {
zap.S().Infof("\tPeer '%s' on block '%s'", lsh.Addr.String(), lsh.BlockID.String())
l.logger.Info("Leash", slog.String("peer", lsh.Addr.String()), slog.String("block", lsh.BlockID.String()))
}
}

Expand Down
19 changes: 9 additions & 10 deletions chains/linkage_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package chains

import (
"fmt"
"log/slog"
"net/netip"
"testing"

Expand Down Expand Up @@ -259,8 +260,8 @@ func TestLastBlockIDs(t *testing.T) {
}

func createTestLinkageAndGenesisID(t testing.TB) (*Linkage, proto.BlockID) {
genesis := settings.TestNetSettings.Genesis
dr, err := NewLinkage(t.TempDir(), proto.TestNetScheme, genesis)
genesis := settings.MustTestNetSettings().Genesis
dr, err := NewLinkage(t.TempDir(), proto.TestNetScheme, genesis, slog.New(slog.DiscardHandler))
require.NoError(t, err)
return dr, genesis.BlockID()
}
Expand All @@ -275,13 +276,11 @@ func createNthBlock(t testing.TB, prev proto.BlockID, n int) *proto.Block {
const tsStep uint64 = 10000
id := testID(t, fmt.Sprintf("BLOCK%d", n))
return &proto.Block{
BlockHeader: proto.BlockHeader{
Version: proto.ProtobufBlockVersion,
Timestamp: tsStep * uint64(n),
Parent: prev,
NxtConsensus: proto.NxtConsensus{BaseTarget: 12345},
GeneratorPublicKey: crypto.PublicKey{},
ID: id,
},
Version: proto.ProtobufBlockVersion,
Timestamp: tsStep * uint64(n),
Parent: prev,
BaseTarget: 12345,
GeneratorPublicKey: crypto.PublicKey{},
ID: id,
}
}
4 changes: 2 additions & 2 deletions chains/storage_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ import (

func TestDoubleInitialization(t *testing.T) {
st, _ := createTestStorage(t)
genesis := settings.TestNetSettings.Genesis
genesis := settings.MustTestNetSettings().Genesis
err := st.initialize(genesis)
require.NoError(t, err)
bl, err := st.block(genesis.BlockID())
Expand Down Expand Up @@ -367,7 +367,7 @@ func BenchmarkLCAOnTwo1MBlockForks(b *testing.B) {
func createTestStorage(t testing.TB) (*storage, proto.BlockID) {
st, err := newStorage(t.TempDir(), proto.TestNetScheme)
require.NoError(t, err)
genesis := settings.TestNetSettings.Genesis
genesis := settings.MustTestNetSettings().Genesis
err = st.initialize(genesis)
require.NoError(t, err)
return st, genesis.BlockID()
Expand Down
Loading
Loading