diff --git a/block/manager_test.go b/block/manager_test.go index bfe291ec40..a85c58f2a5 100644 --- a/block/manager_test.go +++ b/block/manager_test.go @@ -30,7 +30,7 @@ import ( func getManager(t *testing.T, backend goDA.DA) *Manager { logger := test.NewFileLoggerCustom(t, test.TempLogFileName(t, t.Name())) return &Manager{ - dalc: &da.DAClient{DA: backend, GasPrice: -1, GasMultiplier: -1, Logger: logger}, + dalc: da.NewDAClient(backend, -1, -1, nil, logger), blockCache: NewBlockCache(), logger: logger, } diff --git a/cmd/rollkit/commands/run_node.go b/cmd/rollkit/commands/run_node.go index 8ff5931a65..f59d1adf54 100644 --- a/cmd/rollkit/commands/run_node.go +++ b/cmd/rollkit/commands/run_node.go @@ -4,9 +4,8 @@ import ( "context" "fmt" "math/rand" - "net" + "net/url" "os" - "strconv" cmtcmd "github.com/cometbft/cometbft/cmd/cometbft/commands" cometconf "github.com/cometbft/cometbft/config" @@ -20,10 +19,8 @@ import ( cometproxy "github.com/cometbft/cometbft/proxy" comettypes "github.com/cometbft/cometbft/types" comettime "github.com/cometbft/cometbft/types/time" - "google.golang.org/grpc" - "google.golang.org/grpc/credentials/insecure" - "github.com/rollkit/go-da/proxy" + proxy "github.com/rollkit/go-da/proxy/jsonrpc" goDATest "github.com/rollkit/go-da/test" "github.com/spf13/cobra" @@ -115,6 +112,16 @@ func NewRunNodeCmd() *cobra.Command { // initialize the metrics metrics := rollnode.DefaultMetricsProvider(cometconf.DefaultInstrumentationConfig()) + // use mock jsonrpc da server by default + if !cmd.Flags().Lookup("rollkit.da_address").Changed { + srv, err := startMockDAServJSONRPC(cmd.Context()) + if err != nil { + return fmt.Errorf("failed to launch mock da server: %w", err) + } + // nolint:errcheck,gosec + defer func() { srv.Stop(cmd.Context()) }() + } + // create the rollkit node rollnode, err := rollnode.NewNode( context.Background(), @@ -130,9 +137,6 @@ func NewRunNodeCmd() *cobra.Command { return fmt.Errorf("failed to create new rollkit node: %w", err) } - // start mock da server - startMockGRPCServ() - // Launch the RPC server server := rollrpc.NewServer(rollnode, config.RPC, logger) err = server.Start() @@ -173,11 +177,6 @@ func NewRunNodeCmd() *cobra.Command { if !cmd.Flags().Lookup("rollkit.aggregator").Changed { rollkitConfig.Aggregator = true } - - // use mock da server by default - if !cmd.Flags().Lookup("rollkit.da_address").Changed { - rollkitConfig.DAAddress = ":7980" - } return cmd } @@ -193,18 +192,15 @@ func addNodeFlags(cmd *cobra.Command) { rollconf.AddFlags(cmd) } -// startMockGRPCServ starts a mock gRPC server for the dummy DA -func startMockGRPCServ() *grpc.Server { - srv := proxy.NewServer(goDATest.NewDummyDA(), grpc.Creds(insecure.NewCredentials())) - lis, err := net.Listen("tcp", "127.0.0.1"+":"+strconv.Itoa(7980)) +// startMockDAServJSONRPC starts a mock JSONRPC server +func startMockDAServJSONRPC(ctx context.Context) (*proxy.Server, error) { + addr, _ := url.Parse(rollkitConfig.DAAddress) + srv := proxy.NewServer(addr.Hostname(), addr.Port(), goDATest.NewDummyDA()) + err := srv.Start(ctx) if err != nil { - fmt.Println(err) - return nil + return nil, err } - go func() { - _ = srv.Serve(lis) - }() - return srv + return srv, nil } // TODO (Ferret-san): modify so that it initiates files with rollkit configurations by default diff --git a/cmd/rollkit/docs/rollkit_start.md b/cmd/rollkit/docs/rollkit_start.md index 2f7473e62f..e3922113c9 100644 --- a/cmd/rollkit/docs/rollkit_start.md +++ b/cmd/rollkit/docs/rollkit_start.md @@ -30,7 +30,8 @@ rollkit start [flags] --proxy_app string proxy app address, or one of: 'kvstore', 'persistent_kvstore' or 'noop' for local testing. (default "tcp://127.0.0.1:26658") --rollkit.aggregator run node in aggregator mode --rollkit.block_time duration block time (for aggregator mode) (default 1s) - --rollkit.da_address string DA address (host:port) (default ":26650") + --rollkit.da_address string DA address (host:port) (default "http://localhost:26658") + --rollkit.da_auth_token string DA auth token --rollkit.da_block_time duration DA chain block time (for syncing) (default 15s) --rollkit.da_gas_multiplier float DA gas price multiplier for retrying blob transactions (default -1) --rollkit.da_gas_price float DA gas price for blob transactions (default -1) diff --git a/config/config.go b/config/config.go index e14da5faf3..bb560c7da3 100644 --- a/config/config.go +++ b/config/config.go @@ -14,6 +14,8 @@ const ( FlagAggregator = "rollkit.aggregator" // FlagDAAddress is a flag for specifying the data availability layer address FlagDAAddress = "rollkit.da_address" + // FlagDAAuthToken is a flag for specifying the data availability layer auth token + FlagDAAuthToken = "rollkit.da_auth_token" // #nosec G101 // FlagBlockTime is a flag for specifying the block time FlagBlockTime = "rollkit.block_time" // FlagDABlockTime is a flag for specifying the data availability layer block time @@ -45,6 +47,7 @@ type NodeConfig struct { Aggregator bool `mapstructure:"aggregator"` BlockManagerConfig `mapstructure:",squash"` DAAddress string `mapstructure:"da_address"` + DAAuthToken string `mapstructure:"da_auth_token"` Light bool `mapstructure:"light"` HeaderConfig `mapstructure:",squash"` LazyAggregator bool `mapstructure:"lazy_aggregator"` @@ -106,6 +109,7 @@ func GetNodeConfig(nodeConf *NodeConfig, cmConf *cmcfg.Config) { func (nc *NodeConfig) GetViperConfig(v *viper.Viper) error { nc.Aggregator = v.GetBool(FlagAggregator) nc.DAAddress = v.GetString(FlagDAAddress) + nc.DAAuthToken = v.GetString(FlagDAAuthToken) nc.DAGasPrice = v.GetFloat64(FlagDAGasPrice) nc.DAGasMultiplier = v.GetFloat64(FlagDAGasMultiplier) nc.DANamespace = v.GetString(FlagDANamespace) @@ -127,6 +131,7 @@ func AddFlags(cmd *cobra.Command) { cmd.Flags().Bool(FlagAggregator, def.Aggregator, "run node in aggregator mode") cmd.Flags().Bool(FlagLazyAggregator, def.LazyAggregator, "wait for transactions, don't build empty blocks") cmd.Flags().String(FlagDAAddress, def.DAAddress, "DA address (host:port)") + cmd.Flags().String(FlagDAAuthToken, def.DAAuthToken, "DA auth token") cmd.Flags().Duration(FlagBlockTime, def.BlockTime, "block time (for aggregator mode)") cmd.Flags().Duration(FlagDABlockTime, def.DABlockTime, "DA chain block time (for syncing)") cmd.Flags().Float64(FlagDAGasPrice, def.DAGasPrice, "DA gas price for blob transactions") diff --git a/config/defaults.go b/config/defaults.go index 9bed356d36..ab88395025 100644 --- a/config/defaults.go +++ b/config/defaults.go @@ -26,7 +26,7 @@ var DefaultNodeConfig = NodeConfig{ BlockTime: 1 * time.Second, DABlockTime: 15 * time.Second, }, - DAAddress: ":26650", + DAAddress: "http://localhost:26658", DAGasPrice: -1, DAGasMultiplier: -1, Light: false, diff --git a/da/da.go b/da/da.go index 96a50ef778..a325830ae9 100644 --- a/da/da.go +++ b/da/da.go @@ -16,12 +16,12 @@ import ( pb "github.com/rollkit/rollkit/types/pb/rollkit" ) -var ( - // submitTimeout is the timeout for block submission - submitTimeout = 60 * time.Second +const ( + // defaultSubmitTimeout is the timeout for block submission + defaultSubmitTimeout = 60 * time.Second - // retrieveTimeout is the timeout for block retrieval - retrieveTimeout = 60 * time.Second + // defaultRetrieveTimeout is the timeout for block retrieval + defaultRetrieveTimeout = 60 * time.Second ) var ( @@ -95,11 +95,26 @@ type ResultRetrieveBlocks struct { // DAClient is a new DA implementation. type DAClient struct { - DA goDA.DA - GasPrice float64 - GasMultiplier float64 - Namespace goDA.Namespace - Logger log.Logger + DA goDA.DA + GasPrice float64 + GasMultiplier float64 + Namespace goDA.Namespace + SubmitTimeout time.Duration + RetrieveTimeout time.Duration + Logger log.Logger +} + +// NewDAClient returns a new DA client. +func NewDAClient(da goDA.DA, gasPrice, gasMultiplier float64, ns goDA.Namespace, logger log.Logger) *DAClient { + return &DAClient{ + DA: da, + GasPrice: gasPrice, + GasMultiplier: gasMultiplier, + Namespace: ns, + SubmitTimeout: defaultSubmitTimeout, + RetrieveTimeout: defaultRetrieveTimeout, + Logger: logger, + } } // SubmitBlocks submits blocks to DA. @@ -133,7 +148,7 @@ func (dac *DAClient) SubmitBlocks(ctx context.Context, blocks []*types.Block, ma }, } } - ctx, cancel := context.WithTimeout(ctx, submitTimeout) + ctx, cancel := context.WithTimeout(ctx, dac.SubmitTimeout) defer cancel() ids, err := dac.DA.Submit(ctx, blobs, gasPrice, dac.Namespace) if err != nil { @@ -200,7 +215,7 @@ func (dac *DAClient) RetrieveBlocks(ctx context.Context, dataLayerHeight uint64) } } - ctx, cancel := context.WithTimeout(ctx, retrieveTimeout) + ctx, cancel := context.WithTimeout(ctx, dac.RetrieveTimeout) defer cancel() blobs, err := dac.DA.Get(ctx, ids, dac.Namespace) if err != nil { diff --git a/da/da.md b/da/da.md index c986b867f4..9c87cd63f7 100644 --- a/da/da.md +++ b/da/da.md @@ -4,7 +4,11 @@ Rollkit provides a wrapper for [go-da][go-da], a generic data availability inter ## Details -`DAClient` under the hood uses a GRPC implementation of the [go-da][go-da] DA interface. Using the `DAAddress` specified in the node's config, node creates a GRPC connection to it using go-da's gprc implementation [grpc-proxy][grpc-proxy] which is then used under the hood of `DAClient` to communicate with the underlying DA. +`DAClient` can connect via either gRPC or JSON-RPC transports using the [go-da][go-da] [proxy/grpc][proxy/grpc] or [proxy/jsonrpc][proxy/jsonrpc] implementations. The connection can be configured using the following cli flags: + +* `--rollkit.da_address`: url address of the DA service (default: "grpc://localhost:26650") +* `--rollkit.da_auth_token`: authentication token of the DA service +* `--rollkit.da_namespace`: namespace to use when submitting blobs to the DA service Given a set of blocks to be submitted to DA by the block manager, the `SubmitBlocks` first encodes the blocks using protobuf (the encoded data are called blobs) and invokes the `Submit` method on the underlying DA implementation. On successful submission (`StatusSuccess`), the DA block height which included in the rollup blocks is returned. @@ -29,9 +33,12 @@ See [da implementation] [2] [celestia-da][celestia-da] -[3] [grpc-proxy][grpc-proxy] +[3] [proxy/grpc][proxy/grpc] + +[4] [proxy/jsonrpc][proxy/jsonrpc] [da implementation]: https://github.com/rollkit/rollkit/blob/main/da/da.go [go-da]: https://github.com/rollkit/go-da [celestia-da]: https://github.com/rollkit/celestia-da -[grpc-proxy]: https://github.com/rollkit/go-da/tree/main/proxy +[proxy/grpc]: https://github.com/rollkit/go-da/tree/main/proxy/grpc +[proxy/jsonrpc]: https://github.com/rollkit/go-da/tree/main/proxy/jsonrpc diff --git a/da/da_test.go b/da/da_test.go index 7387fed57c..6e556e0d93 100644 --- a/da/da_test.go +++ b/da/da_test.go @@ -4,11 +4,10 @@ import ( "bytes" "context" "errors" - "fmt" "math/rand" "net" + "net/url" "os" - "strconv" "testing" "time" @@ -16,27 +15,52 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "google.golang.org/grpc" - "google.golang.org/grpc/credentials/insecure" "github.com/rollkit/go-da" - "github.com/rollkit/go-da/proxy" + proxygrpc "github.com/rollkit/go-da/proxy/grpc" + proxyjsonrpc "github.com/rollkit/go-da/proxy/jsonrpc" goDATest "github.com/rollkit/go-da/test" "github.com/rollkit/rollkit/da/mock" "github.com/rollkit/rollkit/types" ) -const mockDaBlockTime = 100 * time.Millisecond +const ( + // MockDABlockTime is the mock da block time + MockDABlockTime = 100 * time.Millisecond + + // MockDAAddress is the mock address for the gRPC server + MockDAAddress = "grpc://localhost:7980" + + // MockDAAddressHTTP is mock address for the JSONRPC server + MockDAAddressHTTP = "http://localhost:7988" + + // MockDANamespace is the mock namespace + MockDANamespace = "00000000000000000000000000000000000000000000000000deadbeef" +) +// TestMain starts the mock gRPC and JSONRPC DA services +// gRPC service listens on MockDAAddress +// JSONRPC service listens on MockDAAddressHTTP +// Ports were chosen to be sufficiently different from defaults (26650, 26658) +// Static ports are used to keep client configuration simple +// NOTE: this should be unique per test package to avoid +// "bind: listen address already in use" because multiple packages +// are tested in parallel func TestMain(m *testing.M) { - srv := startMockGRPCServ() - if srv == nil { + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + jsonrpcSrv := startMockDAServJSONRPC(ctx) + if jsonrpcSrv == nil { os.Exit(1) } + grpcSrv := startMockDAServGRPC() exitCode := m.Run() // teardown servers - srv.GracefulStop() + // nolint:errcheck,gosec + jsonrpcSrv.Stop(context.Background()) + grpcSrv.Stop() os.Exit(exitCode) } @@ -44,7 +68,7 @@ func TestMain(m *testing.M) { func TestMockDAErrors(t *testing.T) { t.Run("submit_timeout", func(t *testing.T) { mockDA := &mock.MockDA{} - dalc := &DAClient{DA: mockDA, GasPrice: -1, GasMultiplier: -1, Logger: log.TestingLogger()} + dalc := NewDAClient(mockDA, -1, -1, nil, log.TestingLogger()) blocks := []*types.Block{types.GetRandomBlock(1, 0)} var blobs []da.Blob for _, block := range blocks { @@ -62,7 +86,7 @@ func TestMockDAErrors(t *testing.T) { }) t.Run("max_blob_size_error", func(t *testing.T) { mockDA := &mock.MockDA{} - dalc := &DAClient{DA: mockDA, GasPrice: -1, GasMultiplier: -1, Logger: log.TestingLogger()} + dalc := NewDAClient(mockDA, -1, -1, nil, log.TestingLogger()) // Set up the mock to return an error for MaxBlobSize mockDA.On("MaxBlobSize").Return(uint64(0), errors.New("unable to get DA max blob size")) doTestMaxBlockSizeError(t, dalc) @@ -70,12 +94,17 @@ func TestMockDAErrors(t *testing.T) { } func TestSubmitRetrieve(t *testing.T) { - dummyClient := &DAClient{DA: goDATest.NewDummyDA(), GasPrice: -1, Logger: log.TestingLogger()} - grpcClient, err := startMockGRPCClient() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + dummyClient := NewDAClient(goDATest.NewDummyDA(), -1, -1, nil, log.TestingLogger()) + jsonrpcClient, err := startMockDAClientJSONRPC(ctx) + require.NoError(t, err) + grpcClient := startMockDAClientGRPC() require.NoError(t, err) clients := map[string]*DAClient{ - "dummy": dummyClient, - "grpc": grpcClient, + "dummy": dummyClient, + "jsonrpc": jsonrpcClient, + "grpc": grpcClient, } tests := []struct { name string @@ -97,26 +126,44 @@ func TestSubmitRetrieve(t *testing.T) { } } -func startMockGRPCServ() *grpc.Server { - srv := proxy.NewServer(goDATest.NewDummyDA(), grpc.Creds(insecure.NewCredentials())) - lis, err := net.Listen("tcp", "127.0.0.1"+":"+strconv.Itoa(7980)) +func startMockDAServGRPC() *grpc.Server { + server := proxygrpc.NewServer(goDATest.NewDummyDA(), grpc.Creds(insecure.NewCredentials())) + addr, _ := url.Parse(MockDAAddress) + lis, err := net.Listen("tcp", addr.Host) if err != nil { - fmt.Println(err) - return nil + panic(err) } go func() { - _ = srv.Serve(lis) + _ = server.Serve(lis) }() + return server +} + +func startMockDAClientGRPC() *DAClient { + client := proxygrpc.NewClient() + addr, _ := url.Parse(MockDAAddress) + if err := client.Start(addr.Host, grpc.WithTransportCredentials(insecure.NewCredentials())); err != nil { + panic(err) + } + return NewDAClient(client, -1, -1, nil, log.TestingLogger()) +} + +func startMockDAServJSONRPC(ctx context.Context) *proxyjsonrpc.Server { + addr, _ := url.Parse(MockDAAddressHTTP) + srv := proxyjsonrpc.NewServer(addr.Hostname(), addr.Port(), goDATest.NewDummyDA()) + err := srv.Start(ctx) + if err != nil { + panic(err) + } return srv } -func startMockGRPCClient() (*DAClient, error) { - client := proxy.NewClient() - err := client.Start("127.0.0.1:7980", grpc.WithTransportCredentials(insecure.NewCredentials())) +func startMockDAClientJSONRPC(ctx context.Context) (*DAClient, error) { + client, err := proxyjsonrpc.NewClient(ctx, MockDAAddressHTTP, "") if err != nil { return nil, err } - return &DAClient{DA: client, GasPrice: -1, GasMultiplier: -1, Logger: log.TestingLogger()}, nil + return NewDAClient(&client.DA, -1, -1, nil, log.TestingLogger()), nil } func doTestSubmitTimeout(t *testing.T, dalc *DAClient, blocks []*types.Block) { @@ -127,7 +174,7 @@ func doTestSubmitTimeout(t *testing.T, dalc *DAClient, blocks []*types.Block) { require.NoError(t, err) assert := assert.New(t) - submitTimeout = 50 * time.Millisecond + dalc.SubmitTimeout = 50 * time.Millisecond resp := dalc.SubmitBlocks(ctx, blocks, maxBlobSize, -1) assert.Contains(resp.Message, "context deadline exceeded", "should return context timeout error") } @@ -176,7 +223,7 @@ func doTestSubmitRetrieve(t *testing.T, dalc *DAClient) { blocks[i] = types.GetRandomBlock(batch*numBatches+uint64(i), rand.Int()%20) //nolint:gosec } submitAndRecordBlocks(blocks) - time.Sleep(time.Duration(rand.Int63() % mockDaBlockTime.Milliseconds())) //nolint:gosec + time.Sleep(time.Duration(rand.Int63() % MockDABlockTime.Milliseconds())) //nolint:gosec } validateBlockRetrieval := func(height uint64, expectedCount int) { diff --git a/da/mock/cmd/main.go b/da/mock/cmd/main.go index 12fb217077..3517fa9130 100644 --- a/da/mock/cmd/main.go +++ b/da/mock/cmd/main.go @@ -1,34 +1,43 @@ package main import ( + "context" "flag" + "fmt" "log" - "net" - "strconv" + "net/url" + "os" + "os/signal" + "syscall" - "google.golang.org/grpc" - "google.golang.org/grpc/credentials/insecure" - - "github.com/rollkit/go-da/proxy" + proxy "github.com/rollkit/go-da/proxy/jsonrpc" goDATest "github.com/rollkit/go-da/test" ) +const ( + // MockDAAddress is the mock address for the gRPC server + MockDAAddress = "grpc://localhost:7980" +) + func main() { var ( host string - port int + port string ) - flag.IntVar(&port, "port", 7980, "listening port") - flag.StringVar(&host, "host", "0.0.0.0", "listening address") + addr, _ := url.Parse(MockDAAddress) + flag.StringVar(&port, "port", addr.Port(), "listening port") + flag.StringVar(&host, "host", addr.Hostname(), "listening address") flag.Parse() - lis, err := net.Listen("tcp", host+":"+strconv.Itoa(port)) - if err != nil { - log.Panic(err) - } - log.Println("Listening on:", lis.Addr()) - srv := proxy.NewServer(goDATest.NewDummyDA(), grpc.Creds(insecure.NewCredentials())) - if err := srv.Serve(lis); err != nil { + srv := proxy.NewServer(host, port, goDATest.NewDummyDA()) + log.Printf("Listening on: %s:%s", host, port) + if err := srv.Start(context.Background()); err != nil { log.Fatal("error while serving:", err) } + + interrupt := make(chan os.Signal, 1) + signal.Notify(interrupt, os.Interrupt, syscall.SIGINT) + <-interrupt + fmt.Println("\nCtrl+C pressed. Exiting...") + os.Exit(0) } diff --git a/go.mod b/go.mod index 0edfd43ead..a6f758c636 100644 --- a/go.mod +++ b/go.mod @@ -22,7 +22,7 @@ require ( github.com/multiformats/go-multiaddr v0.12.2 github.com/pkg/errors v0.9.1 github.com/prometheus/client_golang v1.19.0 - github.com/rollkit/go-da v0.4.0 + github.com/rollkit/go-da v0.5.0 github.com/rs/cors v1.10.1 github.com/spf13/cobra v1.8.0 github.com/spf13/viper v1.18.2 @@ -61,6 +61,7 @@ require ( github.com/docker/go-units v0.5.0 // indirect github.com/dustin/go-humanize v1.0.1 // indirect github.com/elastic/gosigar v0.14.2 // indirect + github.com/filecoin-project/go-jsonrpc v0.3.1 // indirect github.com/flynn/noise v1.0.0 // indirect github.com/francoispqt/gojay v1.2.13 // indirect github.com/fsnotify/fsnotify v1.7.0 // indirect @@ -181,6 +182,7 @@ require ( golang.org/x/sys v0.18.0 // indirect golang.org/x/text v0.14.0 // indirect golang.org/x/tools v0.16.0 // indirect + golang.org/x/xerrors v0.0.0-20220907171357-04be3eba64a2 // indirect gonum.org/v1/gonum v0.12.0 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20240123012728-ef4313101c80 // indirect gopkg.in/ini.v1 v1.67.0 // indirect diff --git a/go.sum b/go.sum index 21b90894c0..1795796baf 100644 --- a/go.sum +++ b/go.sum @@ -356,6 +356,8 @@ github.com/fatih/color v1.10.0/go.mod h1:ELkj/draVOlAH/xkhN6mQ50Qd0MPOk5AAr3maGE github.com/fatih/color v1.12.0/go.mod h1:ELkj/draVOlAH/xkhN6mQ50Qd0MPOk5AAr3maGEBuJM= github.com/fatih/color v1.13.0/go.mod h1:kLAiJbzzSOZDVNGyDpeOxJ47H46qBXwg5ILebYFFOfk= github.com/fatih/structtag v1.2.0/go.mod h1:mBJUNpUnHmRKrKlQQlmCrh5PuhftFbNv8Ys4/aAZl94= +github.com/filecoin-project/go-jsonrpc v0.3.1 h1:qwvAUc5VwAkooquKJmfz9R2+F8znhiqcNHYjEp/NM10= +github.com/filecoin-project/go-jsonrpc v0.3.1/go.mod h1:jBSvPTl8V1N7gSTuCR4bis8wnQnIjHbRPpROol6iQKM= github.com/firefart/nonamedreturns v1.0.4/go.mod h1:TDhe/tjI1BXo48CmYbUduTV7BdIga8MAO/xbKdcVsGI= github.com/flynn/go-shlex v0.0.0-20150515145356-3f9db97f8568/go.mod h1:xEzjJPgXI435gkrCt3MPfRiAkVrwSbHsst4LCFVfpJc= github.com/flynn/noise v1.0.0 h1:DlTHqmzmvcEiKj+4RYo/imoswx/4r6iBlCMfVtrMXpQ= @@ -1362,8 +1364,8 @@ github.com/rogpeppe/go-internal v1.6.1/go.mod h1:xXDCJY+GAPziupqXw64V24skbSoqbTE github.com/rogpeppe/go-internal v1.8.1/go.mod h1:JeRgkft04UBgHMgCIwADu4Pn6Mtm5d4nPKWu0nJ5d+o= github.com/rogpeppe/go-internal v1.11.0 h1:cWPaGQEPrBb5/AsnsZesgZZ9yb1OQ+GOISoDNXVBh4M= github.com/rogpeppe/go-internal v1.11.0/go.mod h1:ddIwULY96R17DhadqLgMfk9H9tvdUzkipdSkR5nkCZA= -github.com/rollkit/go-da v0.4.0 h1:/s7ZrVq7DC2aK8UXIvB7rsXrZ2mVGRw7zrexcxRvhlw= -github.com/rollkit/go-da v0.4.0/go.mod h1:Kef0XI5ecEKd3TXzI8S+9knAUJnZg0svh2DuXoCsPlM= +github.com/rollkit/go-da v0.5.0 h1:sQpZricNS+2TLx3HMjNWhtRfqtvVC/U4pWHpfUz3eN4= +github.com/rollkit/go-da v0.5.0/go.mod h1:VsUeAoPvKl4Y8wWguu/VibscYiFFePkkrvZWyTjZHww= github.com/rs/cors v1.7.0/go.mod h1:gFx+x8UowdsKA9AchylcLynDq+nNFfI8FkUZdN/jGCU= github.com/rs/cors v1.8.2/go.mod h1:XyqrcTp5zjWr1wsJ8PIRZssZ8b/WMcMf71DJnit4EMU= github.com/rs/cors v1.10.1 h1:L0uuZVXIKlI1SShY2nhFfo44TYvDPQ1w4oFkUJNfhyo= @@ -2153,6 +2155,8 @@ golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8T golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20220411194840-2f41105eb62f/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20220517211312-f3a8303e98df/go.mod h1:K8+ghG5WaK9qNqU5K3HdILfMLy1f3aNYFI/wnl100a8= +golang.org/x/xerrors v0.0.0-20220907171357-04be3eba64a2 h1:H2TDz8ibqkAF6YGhCdN3jS9O0/s90v0rJh3X/OLHEUk= +golang.org/x/xerrors v0.0.0-20220907171357-04be3eba64a2/go.mod h1:K8+ghG5WaK9qNqU5K3HdILfMLy1f3aNYFI/wnl100a8= gonum.org/v1/gonum v0.0.0-20180816165407-929014505bf4/go.mod h1:Y+Yx5eoAFn32cQvJDxZx5Dpnq+c3wtXuadVZAcxbbBo= gonum.org/v1/gonum v0.8.2/go.mod h1:oe/vMfY3deqTw+1EZJhuvEW2iwGF1bW9wwu7XCu0+v0= gonum.org/v1/gonum v0.12.0 h1:xKuo6hzt+gMav00meVPUlXwSdoEJP46BR+wdxQEFK2o= diff --git a/node/full.go b/node/full.go index c320c8e3df..497e3d9edc 100644 --- a/node/full.go +++ b/node/full.go @@ -14,8 +14,6 @@ import ( "github.com/libp2p/go-libp2p/core/crypto" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/promhttp" - "google.golang.org/grpc" - "google.golang.org/grpc/credentials/insecure" abci "github.com/cometbft/cometbft/abci/types" llcfg "github.com/cometbft/cometbft/config" @@ -26,7 +24,8 @@ import ( rpcclient "github.com/cometbft/cometbft/rpc/client" cmtypes "github.com/cometbft/cometbft/types" - goDAProxy "github.com/rollkit/go-da/proxy" + proxyda "github.com/rollkit/go-da/proxy" + "github.com/rollkit/rollkit/block" "github.com/rollkit/rollkit/config" "github.com/rollkit/rollkit/da" @@ -228,12 +227,14 @@ func initDALC(nodeConfig config.NodeConfig, dalcKV ds.TxnDatastore, logger log.L if err != nil { return nil, fmt.Errorf("error decoding namespace: %w", err) } - daClient := goDAProxy.NewClient() - err = daClient.Start(nodeConfig.DAAddress, grpc.WithTransportCredentials(insecure.NewCredentials())) + + client, err := proxyda.NewClient(nodeConfig.DAAddress, nodeConfig.DAAuthToken) if err != nil { - return nil, fmt.Errorf("error while establishing GRPC connection to DA layer: %w", err) + return nil, fmt.Errorf("error while establishing connection to DA layer: %w", err) } - return &da.DAClient{DA: daClient, Namespace: namespace, GasPrice: nodeConfig.DAGasPrice, GasMultiplier: nodeConfig.DAGasMultiplier, Logger: logger.With("module", "da_client")}, nil + + return da.NewDAClient(client, nodeConfig.DAGasPrice, nodeConfig.DAGasMultiplier, + namespace, logger.With("module", "da_client")), nil } func initMempool(logger log.Logger, proxyApp proxy.AppConns, memplMetrics *mempool.Metrics) *mempool.CListMempool { diff --git a/node/full_client_test.go b/node/full_client_test.go index d05badbe67..edf96a64e7 100644 --- a/node/full_client_test.go +++ b/node/full_client_test.go @@ -69,8 +69,8 @@ func getRPC(t *testing.T) (*mocks.Application, *FullClient) { node, err := newFullNode( ctx, config.NodeConfig{ - DAAddress: MockServerAddr, - DANamespace: MockNamespace, + DAAddress: MockDAAddress, + DANamespace: MockDANamespace, }, key, signingKey, @@ -172,7 +172,7 @@ func TestGenesisChunked(t *testing.T) { signingKey, _, _ := crypto.GenerateEd25519Key(crand.Reader) ctx, cancel := context.WithCancel(context.Background()) defer cancel() - n, _ := newFullNode(ctx, config.NodeConfig{DAAddress: MockServerAddr, DANamespace: MockNamespace}, privKey, signingKey, proxy.NewLocalClientCreator(mockApp), genDoc, DefaultMetricsProvider(cmconfig.DefaultInstrumentationConfig()), test.NewFileLogger(t)) + n, _ := newFullNode(ctx, config.NodeConfig{DAAddress: MockDAAddress, DANamespace: MockDANamespace}, privKey, signingKey, proxy.NewLocalClientCreator(mockApp), genDoc, DefaultMetricsProvider(cmconfig.DefaultInstrumentationConfig()), test.NewFileLogger(t)) rpc := NewFullClient(n) @@ -544,8 +544,8 @@ func TestTx(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() node, err := newFullNode(ctx, config.NodeConfig{ - DAAddress: MockServerAddr, - DANamespace: MockNamespace, + DAAddress: MockDAAddress, + DANamespace: MockDANamespace, Aggregator: true, BlockManagerConfig: config.BlockManagerConfig{ BlockTime: 1 * time.Second, // blocks must be at least 1 sec apart for adjacent headers to get verified correctly @@ -802,8 +802,8 @@ func TestMempool2Nodes(t *testing.T) { defer cancel() // make node1 an aggregator, so that node2 can start gracefully node1, err := newFullNode(ctx, config.NodeConfig{ - DAAddress: MockServerAddr, - DANamespace: MockNamespace, + DAAddress: MockDAAddress, + DANamespace: MockDANamespace, Aggregator: true, P2P: config.P2PConfig{ ListenAddress: "/ip4/127.0.0.1/tcp/9001", @@ -814,8 +814,8 @@ func TestMempool2Nodes(t *testing.T) { require.NotNil(node1) node2, err := newFullNode(ctx, config.NodeConfig{ - DAAddress: MockServerAddr, - DANamespace: MockNamespace, + DAAddress: MockDAAddress, + DANamespace: MockDANamespace, P2P: config.P2PConfig{ ListenAddress: "/ip4/127.0.0.1/tcp/9002", Seeds: "/ip4/127.0.0.1/tcp/9001/p2p/" + id1.Loggable()["peerID"].(string), @@ -877,8 +877,8 @@ func TestStatus(t *testing.T) { node, err := newFullNode( context.Background(), config.NodeConfig{ - DAAddress: MockServerAddr, - DANamespace: MockNamespace, + DAAddress: MockDAAddress, + DANamespace: MockDANamespace, P2P: config.P2PConfig{ ListenAddress: "/ip4/0.0.0.0/tcp/26656", }, @@ -1013,8 +1013,8 @@ func TestFutureGenesisTime(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() node, err := newFullNode(ctx, config.NodeConfig{ - DAAddress: MockServerAddr, - DANamespace: MockNamespace, + DAAddress: MockDAAddress, + DANamespace: MockDANamespace, Aggregator: true, BlockManagerConfig: config.BlockManagerConfig{ BlockTime: 200 * time.Millisecond, diff --git a/node/full_node_integration_test.go b/node/full_node_integration_test.go index 649a6f3658..e62f8148a8 100644 --- a/node/full_node_integration_test.go +++ b/node/full_node_integration_test.go @@ -60,7 +60,7 @@ func TestAggregatorMode(t *testing.T) { } ctx, cancel := context.WithCancel(context.Background()) defer cancel() - node, err := newFullNode(ctx, config.NodeConfig{DAAddress: MockServerAddr, DANamespace: MockNamespace, Aggregator: true, BlockManagerConfig: blockManagerConfig}, key, signingKey, proxy.NewLocalClientCreator(app), genesisDoc, DefaultMetricsProvider(cmconfig.DefaultInstrumentationConfig()), log.TestingLogger()) + node, err := newFullNode(ctx, config.NodeConfig{DAAddress: MockDAAddress, DANamespace: MockDANamespace, Aggregator: true, BlockManagerConfig: blockManagerConfig}, key, signingKey, proxy.NewLocalClientCreator(app), genesisDoc, DefaultMetricsProvider(cmconfig.DefaultInstrumentationConfig()), log.TestingLogger()) require.NoError(err) require.NotNil(node) @@ -183,8 +183,8 @@ func TestLazyAggregator(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() node, err := NewNode(ctx, config.NodeConfig{ - DAAddress: MockServerAddr, - DANamespace: MockNamespace, + DAAddress: MockDAAddress, + DANamespace: MockDANamespace, Aggregator: true, BlockManagerConfig: blockManagerConfig, LazyAggregator: true, @@ -648,8 +648,8 @@ func createNode(ctx context.Context, n int, aggregator bool, isLight bool, keys node, err := NewNode( ctx, config.NodeConfig{ - DAAddress: MockServerAddr, - DANamespace: MockNamespace, + DAAddress: MockDAAddress, + DANamespace: MockDANamespace, P2P: p2pConfig, Aggregator: aggregator, BlockManagerConfig: bmConfig, diff --git a/node/full_node_test.go b/node/full_node_test.go index 53a8fba4b2..52556cb568 100644 --- a/node/full_node_test.go +++ b/node/full_node_test.go @@ -192,11 +192,7 @@ func TestPendingBlocks(t *testing.T) { mockDA.On("MaxBlobSize", mock.Anything).Return(uint64(10240), nil) mockDA.On("Submit", mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return(nil, errors.New("DA not available")) - dac := &da.DAClient{ - DA: mockDA, - Namespace: goDA.Namespace(MockNamespace), - GasPrice: 1234, - } + dac := da.NewDAClient(mockDA, 1234, -1, goDA.Namespace(MockDANamespace), nil) dbPath, err := os.MkdirTemp("", "testdb") require.NoError(t, err) defer func() { @@ -281,8 +277,8 @@ func createAggregatorWithPersistence(ctx context.Context, dbPath string, dalc *d ctx, config.NodeConfig{ DBPath: dbPath, - DAAddress: MockServerAddr, - DANamespace: MockNamespace, + DAAddress: MockDAAddress, + DANamespace: MockDANamespace, Aggregator: true, BlockManagerConfig: config.BlockManagerConfig{ BlockTime: 100 * time.Millisecond, diff --git a/node/helpers_test.go b/node/helpers_test.go index ef583d13c0..938307f94d 100644 --- a/node/helpers_test.go +++ b/node/helpers_test.go @@ -18,10 +18,10 @@ import ( ) func getMockDA(t *testing.T) *da.DAClient { - namespace := make([]byte, len(MockNamespace)/2) - _, err := hex.Decode(namespace, []byte(MockNamespace)) + namespace := make([]byte, len(MockDANamespace)/2) + _, err := hex.Decode(namespace, []byte(MockDANamespace)) require.NoError(t, err) - return &da.DAClient{DA: goDATest.NewDummyDA(), Namespace: namespace, GasPrice: -1, GasMultiplier: -1, Logger: log.TestingLogger()} + return da.NewDAClient(goDATest.NewDummyDA(), -1, -1, namespace, log.TestingLogger()) } func TestMockTester(t *testing.T) { diff --git a/node/node_test.go b/node/node_test.go index 758594c7e1..443d6a74cc 100644 --- a/node/node_test.go +++ b/node/node_test.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "net" + "net/url" "os" "testing" @@ -21,18 +22,20 @@ import ( "google.golang.org/grpc/credentials/insecure" - goDAproxy "github.com/rollkit/go-da/proxy" + goDAproxy "github.com/rollkit/go-da/proxy/grpc" goDATest "github.com/rollkit/go-da/test" ) -// MockServerAddr is the address used by the mock gRPC service -// NOTE: this should be unique per test package to avoid -// "bind: listen address already in use" because multiple packages -// are tested in parallel -var MockServerAddr = "127.0.0.1:7990" - -// MockNamespace is a sample namespace used by the mock DA client -var MockNamespace = "00000000000000000000000000000000000000000000000000deadbeef" +const ( + // MockDAAddress is the address used by the mock gRPC service + // NOTE: this should be unique per test package to avoid + // "bind: listen address already in use" because multiple packages + // are tested in parallel + MockDAAddress = "grpc://localhost:7990" + + // MockDANamespace is a sample namespace used by the mock DA client + MockDANamespace = "00000000000000000000000000000000000000000000000000deadbeef" +) // TestMain does setup and teardown on the test package // to make the mock gRPC service available to the nodes @@ -51,7 +54,8 @@ func TestMain(m *testing.M) { func startMockGRPCServ() *grpc.Server { srv := goDAproxy.NewServer(goDATest.NewDummyDA(), grpc.Creds(insecure.NewCredentials())) - lis, err := net.Listen("tcp", MockServerAddr) + addr, _ := url.Parse(MockDAAddress) + lis, err := net.Listen("tcp", addr.Host) if err != nil { panic(err) } @@ -103,7 +107,7 @@ func setupTestNode(ctx context.Context, t *testing.T, nodeType NodeType) (Node, // newTestNode creates a new test node based on the NodeType. func newTestNode(ctx context.Context, t *testing.T, nodeType NodeType) (Node, ed25519.PrivKey, error) { - config := config.NodeConfig{DAAddress: MockServerAddr, DANamespace: MockNamespace} + config := config.NodeConfig{DAAddress: MockDAAddress, DANamespace: MockDANamespace} switch nodeType { case Light: config.Light = true diff --git a/rpc/json/helpers_test.go b/rpc/json/helpers_test.go index d21bc3fa28..18c66476f4 100644 --- a/rpc/json/helpers_test.go +++ b/rpc/json/helpers_test.go @@ -24,16 +24,20 @@ import ( "github.com/rollkit/rollkit/types" ) +const ( + // MockDAAddress is the mock address for the gRPC server + MockDAAddress = "grpc://localhost:7980" + + // MockDANamespace is the mock namespace + MockDANamespace = "00000000000000000000000000000000000000000000000000deadbeef" +) + func prepareProposalResponse(_ context.Context, req *abci.RequestPrepareProposal) (*abci.ResponsePrepareProposal, error) { return &abci.ResponsePrepareProposal{ Txs: req.Txs, }, nil } -var MockServerAddr = ":7980" - -var MockNamespace = "deadbeef" - // copied from rpc func getRPC(t *testing.T) (*mocks.Application, rpcclient.Client) { t.Helper() @@ -79,7 +83,7 @@ func getRPC(t *testing.T) (*mocks.Application, rpcclient.Client) { genesisValidators := []cmtypes.GenesisValidator{ {Address: pubKey.Address(), PubKey: pubKey, Power: int64(100), Name: "gen #1"}, } - n, err := node.NewNode(context.Background(), config.NodeConfig{DAAddress: MockServerAddr, DANamespace: MockNamespace, Aggregator: true, BlockManagerConfig: config.BlockManagerConfig{BlockTime: 1 * time.Second}, Light: false}, key, signingKey, proxy.NewLocalClientCreator(app), &cmtypes.GenesisDoc{ChainID: "test", Validators: genesisValidators}, node.DefaultMetricsProvider(cmconfig.DefaultInstrumentationConfig()), log.TestingLogger()) + n, err := node.NewNode(context.Background(), config.NodeConfig{DAAddress: MockDAAddress, DANamespace: MockDANamespace, Aggregator: true, BlockManagerConfig: config.BlockManagerConfig{BlockTime: 1 * time.Second}, Light: false}, key, signingKey, proxy.NewLocalClientCreator(app), &cmtypes.GenesisDoc{ChainID: "test", Validators: genesisValidators}, node.DefaultMetricsProvider(cmconfig.DefaultInstrumentationConfig()), log.TestingLogger()) require.NoError(err) require.NotNil(n)