Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,8 @@ If you were at `firehose-core` version `1.0.0` and are bumping to `1.1.0`, you s

### Fixed

- `index-builder` now serves the gRPC health check on `--index-builder-grpc-listen-addr`, which it never listened on.

- Merger no longer moves canonical one-block files to the forked-blocks store when readers write the same block with a different LIB, as happens on Polygon PoS / Amoy. Only blocks with another ID at the same height are moved.

- The reader now always clamps a decoded block's `lib_num` down to its own block number when the node/plugin reports `lib_num` greater than the block number (invalid; `lib_num` equal to the block number is still valid), logs an error and increments `reader_node_invalid_libnum_clamped_count`, instead of letting the bad value reach the relayer where it moved LIB past head and silently stalled it forever.
Expand All @@ -30,6 +32,8 @@ If you were at `firehose-core` version `1.0.0` and are bumping to `1.1.0`, you s

### Added

- `index-builder` now exposes HTTP `/healthz` on `--index-builder-http-healthz-addr` (default `:10019`), like the merger and relayer.

- Store URLs, such as `--common-merged-blocks-store-url`, accept `compression_config` to tune how files are written, matching the store's compression: a zstd level with an optional window in MiB (`best`, `better/32`), or a gzip level from `1` to `9`. For example `gs://bucket/merged-blocks?compression_config=best/32`. Files written with any setting are read back without configuration. An invalid value makes opening the store fail.

- On zstd stores, `compression_config` also takes decoder settings, comma separated after the level or alone: `lowmem=false` gives each decoder a history buffer of twice the window, and `pool=<name>` reuses decoders across reads and across every store using that name, in place of the `default` pool. `pool=none` gives each file a decoder of its own. For example `--common-merged-blocks-store-url=gs://bucket/merged-blocks?compression_config=best/32,lowmem=false,pool=blocks` and `--substreams-state-store-url=gs://bucket/states?compression_config=better/16,pool=cache`. Substreams tier2 gets them through the URLs tier1 sends. Reading 100 MiB `best/32` merged blocks files 10 at a time, adding `lowmem=false` to the `default` pool used 70% less CPU and 57% more peak heap. `lowmem=false` only helps files larger than their window plus about 1 MiB.
Expand Down
13 changes: 8 additions & 5 deletions cmd/apps/index_builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ func RegisterIndexBuilderApp[B firecore.Block](chain *firecore.Chain[B], rootLog
Description: "App the builds indexes out of Firehose blocks",
RegisterFlags: func(cmd *cobra.Command) error {
cmd.Flags().String("index-builder-grpc-listen-addr", firecore.IndexBuilderServiceAddr, "Address to listen for grpc-based healthz check")
cmd.Flags().String("index-builder-http-healthz-addr", firecore.IndexBuilderHTTPHealthzAddr, "Address to listen on for the HTTP /healthz endpoint. Set to an empty string to disable. Returns 200 when ready, 503 otherwise.")
cmd.Flags().Uint64("index-builder-index-size", 10000, "Size of index bundles that will be created")
cmd.Flags().Uint64("index-builder-start-block", 0, "Block number to start indexing")
cmd.Flags().Uint64("index-builder-stop-block", 0, "Block number to stop indexing")
Expand Down Expand Up @@ -81,11 +82,13 @@ func RegisterIndexBuilderApp[B firecore.Block](chain *firecore.Chain[B], rootLog
})

app := index_builder.New(&index_builder.Config{
BlockHandler: handler,
StartBlockResolver: startBlockResolver,
EndBlock: stopBlockNum,
MergedBlocksStoreURL: mergedBlocksStoreURL,
GRPCListenAddr: viper.GetString("index-builder-grpc-listen-addr"),
BlockHandler: handler,
StartBlockResolver: startBlockResolver,
EndBlock: stopBlockNum,
MergedBlocksStoreURL: mergedBlocksStoreURL,
GRPCListenAddr: viper.GetString("index-builder-grpc-listen-addr"),
HTTPHealthzListenAddr: viper.GetString("index-builder-http-healthz-addr"),
IsPendingShutdown: runtime.IsPendingShutdown,
})

return app, nil
Expand Down
1 change: 1 addition & 0 deletions cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,7 @@ func Main[B firecore.Block](chain *firecore.Chain[B]) {
--substreams-tier1-grpc-listen-addr :10016
--substreams-tier2-grpc-listen-addr :10017
--index-builder-grpc-listen-addr :10009
--index-builder-http-healthz-addr :10019

Client connection addresses:
--common-live-blocks-addr :10014
Expand Down
1 change: 1 addition & 0 deletions cmd/shift_ports.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ var portFlags = []string{
"substreams-tier1-grpc-listen-addr",
"substreams-tier2-grpc-listen-addr",
"index-builder-grpc-listen-addr",
"index-builder-http-healthz-addr",

// Client-side connection addresses (defaults point to services above)
"common-live-blocks-addr",
Expand Down
1 change: 1 addition & 0 deletions constants.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ var (
SubstreamsTier1GRPCServingAddr string = ":10016"
SubstreamsTier2GRPCServingAddr string = ":10017"
RelayerHTTPHealthzAddr string = ":10018"
IndexBuilderHTTPHealthzAddr string = ":10019"

// Data storage default locations
BlocksCacheDirectory string = "file://{data-dir}/storage/blocks-cache"
Expand Down
42 changes: 42 additions & 0 deletions index-builder/app/index-builder/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@ package index_builder
import (
"context"
"fmt"
"net/http"
"time"

"github.com/streamingfast/bstream"
"github.com/streamingfast/dgrpc"
Expand All @@ -22,6 +24,10 @@ type Config struct {
MergedBlocksStoreURL string
ForkedBlocksStoreURL string
GRPCListenAddr string

HTTPHealthzListenAddr string

IsPendingShutdown func() bool `json:"-"`
}

type App struct {
Expand Down Expand Up @@ -59,6 +65,7 @@ func (a *App) Run() error {
startBlock,
a.config.EndBlock,
blockStore,
a.config.GRPCListenAddr,
)

gs, err := dgrpc.NewInternalClient(a.config.GRPCListenAddr)
Expand All @@ -72,12 +79,47 @@ func (a *App) Run() error {
a.OnTerminating(indexBuilder.Shutdown)
indexBuilder.OnTerminated(a.Shutdown)

if a.config.HTTPHealthzListenAddr != "" {
a.startHTTPHealthzServer(indexBuilder)
}

go indexBuilder.Launch()

zlog.Info("index builder running")
return nil
}

func (a *App) startHTTPHealthzServer(indexBuilder *index_builder.IndexBuilder) {
mux := http.NewServeMux()
mux.HandleFunc("/healthz", func(w http.ResponseWriter, _ *http.Request) {
if a.config.IsPendingShutdown != nil && a.config.IsPendingShutdown() {
http.Error(w, "not ready: shutting down", http.StatusServiceUnavailable)
return
}
resp, err := indexBuilder.Check(context.Background(), &pbhealth.HealthCheckRequest{})
if err != nil || resp.Status != pbhealth.HealthCheckResponse_SERVING {
http.Error(w, "not ready", http.StatusServiceUnavailable)
return
}
w.Write([]byte("ready\n"))
})

srv := &http.Server{Addr: a.config.HTTPHealthzListenAddr, Handler: mux}
a.OnTerminating(func(_ error) {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
_ = srv.Shutdown(ctx)
})

zlog.Info("starting index builder http healthz server", zap.String("addr", a.config.HTTPHealthzListenAddr))
go func() {
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
zlog.Error("index builder http healthz server failed", zap.Error(err))
a.Shutdown(err)
}
}()
}

func (a *App) IsReady() bool {
if a.readinessProbe == nil {
return false
Expand Down
47 changes: 47 additions & 0 deletions index-builder/app/index-builder/app_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
package index_builder

import (
"context"
"io"
"net"
"net/http"
"testing"
"time"

"github.com/streamingfast/bstream"
pbbstream "github.com/streamingfast/bstream/pb/sf/bstream/v1"
"github.com/stretchr/testify/require"
)

func TestApp_ServesHTTPHealthz(t *testing.T) {
app := New(&Config{
BlockHandler: bstream.HandlerFunc(func(*pbbstream.Block, interface{}) error { return nil }),
StartBlockResolver: func(context.Context) (uint64, error) { return 0, nil },
MergedBlocksStoreURL: "file://" + t.TempDir(),
GRPCListenAddr: freeAddr(t),
HTTPHealthzListenAddr: freeAddr(t),
})
require.NoError(t, app.Run())
defer app.Shutdown(nil)

url := "http://" + app.config.HTTPHealthzListenAddr + "/healthz"
require.Eventually(t, func() bool {
resp, err := http.Get(url)
if err != nil {
return false
}
defer resp.Body.Close()
body, _ := io.ReadAll(resp.Body)
return resp.StatusCode == http.StatusOK && string(body) == "ready\n"
}, 10*time.Second, 100*time.Millisecond)

require.Eventually(t, app.IsReady, 10*time.Second, 100*time.Millisecond)
}

func freeAddr(t *testing.T) string {
t.Helper()
listener, err := net.Listen("tcp", "127.0.0.1:0")
require.NoError(t, err)
defer listener.Close()
return listener.Addr().String()
}
11 changes: 11 additions & 0 deletions index-builder/healthz.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,17 @@ func (app *IndexBuilder) Check(ctx context.Context, in *pbhealth.HealthCheckRequ
}, nil
}

func (app *IndexBuilder) List(ctx context.Context, in *pbhealth.HealthListRequest) (*pbhealth.HealthListResponse, error) {
status := pbhealth.HealthCheckResponse_SERVING
return &pbhealth.HealthListResponse{
Statuses: map[string]*pbhealth.HealthCheckResponse{
"index-builder": &pbhealth.HealthCheckResponse{
Status: status,
},
},
}, nil
}

// Watch is basic GRPC Healthcheck as a stream
func (app *IndexBuilder) Watch(req *pbhealth.HealthCheckRequest, stream pbhealth.Health_WatchServer) error {
err := stream.Send(&pbhealth.HealthCheckResponse{
Expand Down
17 changes: 11 additions & 6 deletions index-builder/index-builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,21 +27,26 @@ type IndexBuilder struct {
handler bstream.Handler

blocksStore dstore.Store

grpcListenAddr string
}

func NewIndexBuilder(logger *zap.Logger, handler bstream.Handler, startBlockNum, stopBlockNum uint64, blockStore dstore.Store) *IndexBuilder {
func NewIndexBuilder(logger *zap.Logger, handler bstream.Handler, startBlockNum, stopBlockNum uint64, blockStore dstore.Store, grpcListenAddr string) *IndexBuilder {
return &IndexBuilder{
Shutter: shutter.New(),
startBlockNum: startBlockNum,
stopBlockNum: stopBlockNum,
handler: handler,
blocksStore: blockStore,
Shutter: shutter.New(),
startBlockNum: startBlockNum,
stopBlockNum: stopBlockNum,
handler: handler,
blocksStore: blockStore,
grpcListenAddr: grpcListenAddr,

logger: logger,
}
}

func (app *IndexBuilder) Launch() {
app.startGRPCServer()

err := app.launch()
if errors.Is(err, stream.ErrStopBlockReached) {
app.logger.Info("index builder reached stop block", zap.Uint64("stop_block_num", app.stopBlockNum))
Expand Down
20 changes: 20 additions & 0 deletions index-builder/server.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
package index_builder

import (
dgrpcfactory "github.com/streamingfast/dgrpc/server/factory"
pbhealth "google.golang.org/grpc/health/grpc_health_v1"
)

func (app *IndexBuilder) startGRPCServer() {
gs := dgrpcfactory.ServerFromOptions()
gs.OnTerminated(app.Shutdown)
app.logger.Info("grpc server created")

app.OnTerminated(func(_ error) {
gs.Shutdown(0)
})
pbhealth.RegisterHealthServer(gs.ServiceRegistrar(), app)
app.logger.Info("server registered")

go gs.Launch(app.grpcListenAddr)
}
45 changes: 45 additions & 0 deletions index-builder/server_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
package index_builder

import (
"context"
"net"
"testing"
"time"

"github.com/streamingfast/bstream"
pbbstream "github.com/streamingfast/bstream/pb/sf/bstream/v1"
"github.com/streamingfast/dstore"
"github.com/stretchr/testify/require"
"go.uber.org/zap"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
pbhealth "google.golang.org/grpc/health/grpc_health_v1"
)

func TestIndexBuilder_ServesGRPCHealth(t *testing.T) {
listener, err := net.Listen("tcp", "127.0.0.1:0")
require.NoError(t, err)
addr := listener.Addr().String()
require.NoError(t, listener.Close())

store, err := dstore.NewDBinStore("file://" + t.TempDir())
require.NoError(t, err)

handler := bstream.HandlerFunc(func(*pbbstream.Block, interface{}) error { return nil })
app := NewIndexBuilder(zap.NewNop(), handler, 0, 0, store, addr)
go app.Launch()
defer app.Shutdown(nil)

conn, err := grpc.NewClient(addr, grpc.WithTransportCredentials(insecure.NewCredentials()))
require.NoError(t, err)
defer conn.Close()

client := pbhealth.NewHealthClient(conn)
require.Eventually(t, func() bool {
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()

resp, err := client.Check(ctx, &pbhealth.HealthCheckRequest{})
return err == nil && resp.Status == pbhealth.HealthCheckResponse_SERVING
}, 10*time.Second, 100*time.Millisecond)
}
6 changes: 6 additions & 0 deletions shift_ports_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,12 @@ func TestShiftAddressPort(t *testing.T) {
offset: 100,
expected: ":10109",
},
{
name: "index builder http healthz default port",
addr: ":10019",
offset: 100,
expected: ":10119",
},
{
name: "reader node grpc default port",
addr: ":10010",
Expand Down
Loading