mirror of
https://github.com/luxfi/node.git
synced 2026-07-27 03:39:39 +00:00
node: ZAP-native everywhere — kill every //go:build grpc path
The companion commit `7496606282` retired the gRPC fallback inside
vms/rpcchainvm/. This commit closes the loop by deleting every
remaining `//go:build grpc` file under node/, removing the dual-build
plumbing entirely. ZAP is the only wire protocol; there is no
`-tags=grpc` opt-in.
Per the "one and only one way to do everything" mandate, the
backwards-compat scaffolding (gRPC adapters, protoc stubs, connect-go
example handlers, OTLP gRPC exporter, x/sync gRPC sync engine,
keystore-over-gRPC client/server, rpcwarp gRPC signer, gRPC alias
reader) is removed forward-only — no aliases, no deprecation period.
Deletions (84 files):
- proto/pb/{aliasreader,http,io,keystore,message,messenger,net,p2p,
platformvm,rpcdb,sdk,sender,sharedmemory,signer,sync,
validatorstate,vm,warp}/ — protoc stubs (entire tree)
- proto/{p2p,platformvm,sync,vm}/*_grpc.go — gRPC type re-exports
- db/rpcdb/{grpc_server,grpc_client,grpc_test}.go — rpcdb gRPC adapter
- service/keystore/rpckeystore/ — keystore-over-gRPC (dead consumer)
- x/sync/ — entire merkledb sync engine (100% gRPC-tagged)
- internal/ids/rpcaliasreader/ — gRPC alias reader (dead consumer)
- connectproto/ — connect-go XSVM ping handler scaffolding
- vms/platformvm/warp/rpcwarp/{client,server}.go — gRPC warp signer
- vms/platformvm/network/warp.go — protobuf-based warp justification
handler (warp_zap.go retains the no-op verifier consumers expect)
- vms/components/message/message_grpc.go + message_test.go
- vms/example/xsvm/api/ping.go + vm_http_grpc.go + cmd/{run,xsvm}/
- wallet/network/primary/examples/sign-l1-validator-* (5 dead
example main packages)
- trace/{exporter_grpc,exporter_type,exporter_type_test,noop,tracer}.go
(OTLP gRPC exporter + duplicate Tracer types — trace_zap.go has
the canonical Tracer interface + no-op tracer)
- message/bft_grpc.go (Simplex BFT wrapper)
Edits (9 files): every surviving `_zap.go` drops its `//go:build !grpc`
constraint (the files are unconditional now) and drops stale
"ZAP version" / "ZAP mode" inline comments.
Doc updates:
- LLM.md ZAP Transport section: remove the `-tags=grpc` opt-in
language. Replace with "ZAP is the only wire protocol... there is
one and only one way". Update Latest Tag to v1.26.31. Update the
rpcdb topology section to reflect single-adapter state.
go.mod: connectrpc.com/connect, connectrpc.com/grpcreflect,
otlptracegrpc moved out of direct deps. google.golang.org/grpc +
protobuf demoted to indirect (still pulled transitively via luxfi/dex).
Verified:
- `go build ./...` (default, no tags) clean
- `go test ./db/rpcdb/... ./trace/... ./message/... ./vms/rpcchainvm/...
./vms/components/message/... ./vms/platformvm/network/...
./vms/example/xsvm/... ./service/keystore/... ./proto/...` clean
- `grep -rln '//go:build grpc' --include='*.go' node/` returns zero
- `grep -rln '//go:build !grpc' --include='*.go' node/` returns zero
- `grep -rln 'google.golang.org/grpc' --include='*.go' node/` returns
zero (only transitive deps remain in go.sum)
Pre-existing TestGraniteNetworkIDConfiguration failure in tests/ is
unrelated — fails on the parent commit too (verified via stash).
This commit is contained in:
@@ -1,370 +0,0 @@
|
||||
//go:build grpc
|
||||
|
||||
// Copyright (C) 2019-2026, Lux Industries Inc. All rights reserved.
|
||||
// See the file LICENSE for licensing terms.
|
||||
|
||||
package rpcdb
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"sync"
|
||||
|
||||
"google.golang.org/protobuf/types/known/emptypb"
|
||||
|
||||
"github.com/luxfi/database"
|
||||
"github.com/luxfi/math/set"
|
||||
"github.com/luxfi/utils"
|
||||
|
||||
rpcdbpb "github.com/luxfi/node/proto/pb/rpcdb"
|
||||
)
|
||||
|
||||
var (
|
||||
_ database.Database = (*GRPCClient)(nil)
|
||||
_ database.Batch = (*grpcBatch)(nil)
|
||||
_ database.Iterator = (*grpcIterator)(nil)
|
||||
)
|
||||
|
||||
// GRPCClient is a database.Database that talks to a GRPCServer over
|
||||
// gRPC. Symmetric with GRPCServer — both are pure transport adapters
|
||||
// against the Layer-B wire types in github.com/luxfi/protocol/rpcdb.
|
||||
type GRPCClient struct {
|
||||
client rpcdbpb.DatabaseClient
|
||||
|
||||
closed utils.Atomic[bool]
|
||||
}
|
||||
|
||||
// NewGRPCClient wraps a gRPC client connection to a remote rpcdb
|
||||
// service.
|
||||
func NewGRPCClient(client rpcdbpb.DatabaseClient) *GRPCClient {
|
||||
return &GRPCClient{client: client}
|
||||
}
|
||||
|
||||
// Has attempts to return if the database has a key with the provided value.
|
||||
func (c *GRPCClient) Has(key []byte) (bool, error) {
|
||||
resp, err := c.client.Has(context.Background(), &rpcdbpb.HasRequest{Key: key})
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return resp.Has, codeToErr(resp.Err)
|
||||
}
|
||||
|
||||
// Get attempts to return the value mapped to the key.
|
||||
func (c *GRPCClient) Get(key []byte) ([]byte, error) {
|
||||
resp, err := c.client.Get(context.Background(), &rpcdbpb.GetRequest{Key: key})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return resp.Value, codeToErr(resp.Err)
|
||||
}
|
||||
|
||||
// Put attempts to set the value this key maps to.
|
||||
func (c *GRPCClient) Put(key, value []byte) error {
|
||||
resp, err := c.client.Put(context.Background(), &rpcdbpb.PutRequest{
|
||||
Key: key,
|
||||
Value: value,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return codeToErr(resp.Err)
|
||||
}
|
||||
|
||||
// Delete attempts to remove any mapping from the key.
|
||||
func (c *GRPCClient) Delete(key []byte) error {
|
||||
resp, err := c.client.Delete(context.Background(), &rpcdbpb.DeleteRequest{Key: key})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return codeToErr(resp.Err)
|
||||
}
|
||||
|
||||
// NewBatch returns a new batch.
|
||||
func (c *GRPCClient) NewBatch() database.Batch {
|
||||
return &grpcBatch{db: c}
|
||||
}
|
||||
|
||||
func (c *GRPCClient) NewIterator() database.Iterator {
|
||||
return c.NewIteratorWithStartAndPrefix(nil, nil)
|
||||
}
|
||||
|
||||
func (c *GRPCClient) NewIteratorWithStart(start []byte) database.Iterator {
|
||||
return c.NewIteratorWithStartAndPrefix(start, nil)
|
||||
}
|
||||
|
||||
func (c *GRPCClient) NewIteratorWithPrefix(prefix []byte) database.Iterator {
|
||||
return c.NewIteratorWithStartAndPrefix(nil, prefix)
|
||||
}
|
||||
|
||||
// NewIteratorWithStartAndPrefix returns a new iterator.
|
||||
func (c *GRPCClient) NewIteratorWithStartAndPrefix(start, prefix []byte) database.Iterator {
|
||||
resp, err := c.client.NewIteratorWithStartAndPrefix(
|
||||
context.Background(),
|
||||
&rpcdbpb.NewIteratorWithStartAndPrefixRequest{Start: start, Prefix: prefix},
|
||||
)
|
||||
if err != nil {
|
||||
return &database.IteratorError{Err: err}
|
||||
}
|
||||
return newGRPCIterator(c, resp.Id)
|
||||
}
|
||||
|
||||
// Compact attempts to optimize the space utilization in the provided range.
|
||||
func (c *GRPCClient) Compact(start, limit []byte) error {
|
||||
resp, err := c.client.Compact(context.Background(), &rpcdbpb.CompactRequest{
|
||||
Start: start,
|
||||
Limit: limit,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return codeToErr(resp.Err)
|
||||
}
|
||||
|
||||
// Close attempts to close the database.
|
||||
func (c *GRPCClient) Close() error {
|
||||
c.closed.Set(true)
|
||||
resp, err := c.client.Close(context.Background(), &rpcdbpb.CloseRequest{})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return codeToErr(resp.Err)
|
||||
}
|
||||
|
||||
// Sync — rpc database delegates sync to the underlying database on the
|
||||
// server side. Operations are already synchronous over RPC; no-op here
|
||||
// matches the legacy internal/database/rpcdb behavior.
|
||||
func (c *GRPCClient) Sync() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *GRPCClient) HealthCheck(ctx context.Context) (interface{}, error) {
|
||||
health, err := c.client.HealthCheck(ctx, &emptypb.Empty{})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return json.RawMessage(health.Details), nil
|
||||
}
|
||||
|
||||
// Backup is not supported over the gRPC database client.
|
||||
func (c *GRPCClient) Backup(_ io.Writer, _ uint64) (uint64, error) {
|
||||
return 0, errors.New("rpcdb: backup not supported")
|
||||
}
|
||||
|
||||
// Load is not supported over the gRPC database client.
|
||||
func (c *GRPCClient) Load(_ io.Reader) error {
|
||||
return errors.New("rpcdb: load not supported")
|
||||
}
|
||||
|
||||
// codeToErr maps a generated-pb Error enum back to a database sentinel.
|
||||
// Lives next to the gRPC adapter so the adapter knows the inverse of
|
||||
// toPbErr. Mirrors rpcdb.CodeToErr in service.go but for the
|
||||
// protobuf-typed enum.
|
||||
func codeToErr(code rpcdbpb.Error) error {
|
||||
switch code {
|
||||
case rpcdbpb.Error_ERROR_CLOSED:
|
||||
return database.ErrClosed
|
||||
case rpcdbpb.Error_ERROR_NOT_FOUND:
|
||||
return database.ErrNotFound
|
||||
default:
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
type grpcBatch struct {
|
||||
database.BatchOps
|
||||
|
||||
db *GRPCClient
|
||||
}
|
||||
|
||||
func (b *grpcBatch) Write() error {
|
||||
request := &rpcdbpb.WriteBatchRequest{}
|
||||
keySet := set.NewSet[string](len(b.Ops))
|
||||
for i := len(b.Ops) - 1; i >= 0; i-- {
|
||||
op := b.Ops[i]
|
||||
key := string(op.Key)
|
||||
if keySet.Contains(key) {
|
||||
continue
|
||||
}
|
||||
keySet.Add(key)
|
||||
|
||||
if op.Delete {
|
||||
request.Deletes = append(request.Deletes, &rpcdbpb.DeleteRequest{Key: op.Key})
|
||||
} else {
|
||||
request.Puts = append(request.Puts, &rpcdbpb.PutRequest{
|
||||
Key: op.Key,
|
||||
Value: op.Value,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
resp, err := b.db.client.WriteBatch(context.Background(), request)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return codeToErr(resp.Err)
|
||||
}
|
||||
|
||||
func (b *grpcBatch) Inner() database.Batch {
|
||||
return b
|
||||
}
|
||||
|
||||
type grpcIterator struct {
|
||||
db *GRPCClient
|
||||
id uint64
|
||||
|
||||
data []*rpcdbpb.PutRequest
|
||||
fetchedData chan []*rpcdbpb.PutRequest
|
||||
|
||||
errLock sync.RWMutex
|
||||
err error
|
||||
|
||||
reqUpdateError chan chan struct{}
|
||||
|
||||
once sync.Once
|
||||
onClose chan struct{}
|
||||
onClosed chan struct{}
|
||||
}
|
||||
|
||||
func newGRPCIterator(db *GRPCClient, id uint64) *grpcIterator {
|
||||
it := &grpcIterator{
|
||||
db: db,
|
||||
id: id,
|
||||
fetchedData: make(chan []*rpcdbpb.PutRequest),
|
||||
reqUpdateError: make(chan chan struct{}),
|
||||
onClose: make(chan struct{}),
|
||||
onClosed: make(chan struct{}),
|
||||
}
|
||||
go it.fetch()
|
||||
return it
|
||||
}
|
||||
|
||||
// Invariant: fetch is the only thread with access to send requests to
|
||||
// the server's iterator. This is needed because iterators are not
|
||||
// thread safe and the server expects the client (us) to only ever
|
||||
// issue one request at a time for a given iterator id.
|
||||
func (it *grpcIterator) fetch() {
|
||||
defer func() {
|
||||
resp, err := it.db.client.IteratorRelease(context.Background(), &rpcdbpb.IteratorReleaseRequest{Id: it.id})
|
||||
if err != nil {
|
||||
it.setError(err)
|
||||
} else {
|
||||
it.setError(codeToErr(resp.Err))
|
||||
}
|
||||
|
||||
close(it.fetchedData)
|
||||
close(it.onClosed)
|
||||
}()
|
||||
|
||||
for {
|
||||
resp, err := it.db.client.IteratorNext(context.Background(), &rpcdbpb.IteratorNextRequest{Id: it.id})
|
||||
if err != nil {
|
||||
it.setError(err)
|
||||
return
|
||||
}
|
||||
|
||||
if len(resp.Data) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
for {
|
||||
select {
|
||||
case it.fetchedData <- resp.Data:
|
||||
case onUpdated := <-it.reqUpdateError:
|
||||
it.updateError()
|
||||
close(onUpdated)
|
||||
continue
|
||||
case <-it.onClose:
|
||||
return
|
||||
}
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Next attempts to move the iterator to the next element.
|
||||
func (it *grpcIterator) Next() bool {
|
||||
if it.db.closed.Get() {
|
||||
it.data = nil
|
||||
it.setError(database.ErrClosed)
|
||||
return false
|
||||
}
|
||||
if len(it.data) > 1 {
|
||||
it.data[0] = nil
|
||||
it.data = it.data[1:]
|
||||
return true
|
||||
}
|
||||
|
||||
it.data = <-it.fetchedData
|
||||
return len(it.data) > 0
|
||||
}
|
||||
|
||||
// Error returns any error that occurred while iterating.
|
||||
func (it *grpcIterator) Error() error {
|
||||
if err := it.getError(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
onUpdated := make(chan struct{})
|
||||
select {
|
||||
case it.reqUpdateError <- onUpdated:
|
||||
<-onUpdated
|
||||
case <-it.onClosed:
|
||||
}
|
||||
|
||||
return it.getError()
|
||||
}
|
||||
|
||||
// Key returns the key of the current element.
|
||||
func (it *grpcIterator) Key() []byte {
|
||||
if len(it.data) == 0 {
|
||||
return nil
|
||||
}
|
||||
return it.data[0].Key
|
||||
}
|
||||
|
||||
// Value returns the value of the current element.
|
||||
func (it *grpcIterator) Value() []byte {
|
||||
if len(it.data) == 0 {
|
||||
return nil
|
||||
}
|
||||
return it.data[0].Value
|
||||
}
|
||||
|
||||
// Release frees any resources held by the iterator.
|
||||
func (it *grpcIterator) Release() {
|
||||
it.once.Do(func() {
|
||||
close(it.onClose)
|
||||
<-it.onClosed
|
||||
})
|
||||
}
|
||||
|
||||
func (it *grpcIterator) updateError() {
|
||||
resp, err := it.db.client.IteratorError(context.Background(), &rpcdbpb.IteratorErrorRequest{Id: it.id})
|
||||
if err != nil {
|
||||
it.setError(err)
|
||||
} else {
|
||||
it.setError(codeToErr(resp.Err))
|
||||
}
|
||||
}
|
||||
|
||||
func (it *grpcIterator) setError(err error) {
|
||||
if err == nil {
|
||||
return
|
||||
}
|
||||
|
||||
it.errLock.Lock()
|
||||
defer it.errLock.Unlock()
|
||||
|
||||
if it.err == nil {
|
||||
it.err = err
|
||||
}
|
||||
}
|
||||
|
||||
func (it *grpcIterator) getError() error {
|
||||
it.errLock.RLock()
|
||||
defer it.errLock.RUnlock()
|
||||
|
||||
return it.err
|
||||
}
|
||||
@@ -1,159 +0,0 @@
|
||||
//go:build grpc
|
||||
|
||||
// Copyright (C) 2019-2026, Lux Industries Inc. All rights reserved.
|
||||
// See the file LICENSE for licensing terms.
|
||||
|
||||
package rpcdb
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"google.golang.org/protobuf/types/known/emptypb"
|
||||
|
||||
"github.com/luxfi/database"
|
||||
rpcdbpb "github.com/luxfi/node/proto/pb/rpcdb"
|
||||
rpcdb "github.com/luxfi/protocol/rpcdb"
|
||||
)
|
||||
|
||||
// GRPCServer is the gRPC transport adapter for the rpcdb Service. It
|
||||
// wraps *Service and translates gRPC's protobuf-generated request /
|
||||
// response types into the transport-neutral wire types in
|
||||
// github.com/luxfi/protocol/rpcdb. Pure adapter — no storage logic
|
||||
// lives here.
|
||||
type GRPCServer struct {
|
||||
rpcdbpb.UnimplementedDatabaseServer
|
||||
svc *Service
|
||||
}
|
||||
|
||||
// NewGRPCServer wraps a database.Database for serving over gRPC.
|
||||
func NewGRPCServer(db database.Database) *GRPCServer {
|
||||
return NewGRPCServerFromService(NewService(db))
|
||||
}
|
||||
|
||||
// NewGRPCServerFromService wraps an existing Service for serving over
|
||||
// gRPC. Symmetric with NewZAPServerFromService.
|
||||
func NewGRPCServerFromService(svc *Service) *GRPCServer {
|
||||
return &GRPCServer{svc: svc}
|
||||
}
|
||||
|
||||
// Compile-time interface check.
|
||||
var _ rpcdbpb.DatabaseServer = (*GRPCServer)(nil)
|
||||
|
||||
func toPbErr(code rpcdb.Error) rpcdbpb.Error {
|
||||
switch code {
|
||||
case rpcdb.Error_ERROR_NOT_FOUND:
|
||||
return rpcdbpb.Error_ERROR_NOT_FOUND
|
||||
case rpcdb.Error_ERROR_CLOSED:
|
||||
return rpcdbpb.Error_ERROR_CLOSED
|
||||
default:
|
||||
return rpcdbpb.Error_ERROR_UNSPECIFIED
|
||||
}
|
||||
}
|
||||
|
||||
func (g *GRPCServer) Has(ctx context.Context, req *rpcdbpb.HasRequest) (*rpcdbpb.HasResponse, error) {
|
||||
resp, err := g.svc.Has(ctx, &rpcdb.HasRequest{Key: req.Key})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &rpcdbpb.HasResponse{Has: resp.Has, Err: toPbErr(resp.Err)}, nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer) Get(ctx context.Context, req *rpcdbpb.GetRequest) (*rpcdbpb.GetResponse, error) {
|
||||
resp, err := g.svc.Get(ctx, &rpcdb.GetRequest{Key: req.Key})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &rpcdbpb.GetResponse{Value: resp.Value, Err: toPbErr(resp.Err)}, nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer) Put(ctx context.Context, req *rpcdbpb.PutRequest) (*rpcdbpb.PutResponse, error) {
|
||||
resp, err := g.svc.Put(ctx, &rpcdb.PutRequest{Key: req.Key, Value: req.Value})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &rpcdbpb.PutResponse{Err: toPbErr(resp.Err)}, nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer) Delete(ctx context.Context, req *rpcdbpb.DeleteRequest) (*rpcdbpb.DeleteResponse, error) {
|
||||
resp, err := g.svc.Delete(ctx, &rpcdb.DeleteRequest{Key: req.Key})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &rpcdbpb.DeleteResponse{Err: toPbErr(resp.Err)}, nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer) WriteBatch(ctx context.Context, req *rpcdbpb.WriteBatchRequest) (*rpcdbpb.WriteBatchResponse, error) {
|
||||
puts := make([]*rpcdb.PutRequest, len(req.Puts))
|
||||
for i, p := range req.Puts {
|
||||
puts[i] = &rpcdb.PutRequest{Key: p.Key, Value: p.Value}
|
||||
}
|
||||
dels := make([]*rpcdb.DeleteRequest, len(req.Deletes))
|
||||
for i, d := range req.Deletes {
|
||||
dels[i] = &rpcdb.DeleteRequest{Key: d.Key}
|
||||
}
|
||||
resp, err := g.svc.WriteBatch(ctx, &rpcdb.WriteBatchRequest{Puts: puts, Deletes: dels})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &rpcdbpb.WriteBatchResponse{Err: toPbErr(resp.Err)}, nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer) Compact(ctx context.Context, req *rpcdbpb.CompactRequest) (*rpcdbpb.CompactResponse, error) {
|
||||
resp, err := g.svc.Compact(ctx, &rpcdb.CompactRequest{Start: req.Start, Limit: req.Limit})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &rpcdbpb.CompactResponse{Err: toPbErr(resp.Err)}, nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer) Close(ctx context.Context, _ *rpcdbpb.CloseRequest) (*rpcdbpb.CloseResponse, error) {
|
||||
resp, err := g.svc.Close(ctx, &rpcdb.CloseRequest{})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &rpcdbpb.CloseResponse{Err: toPbErr(resp.Err)}, nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer) HealthCheck(ctx context.Context, _ *emptypb.Empty) (*rpcdbpb.HealthCheckResponse, error) {
|
||||
resp, err := g.svc.HealthCheck(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &rpcdbpb.HealthCheckResponse{Details: resp.Details}, nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer) NewIteratorWithStartAndPrefix(ctx context.Context, req *rpcdbpb.NewIteratorWithStartAndPrefixRequest) (*rpcdbpb.NewIteratorWithStartAndPrefixResponse, error) {
|
||||
resp, err := g.svc.NewIteratorWithStartAndPrefix(ctx, &rpcdb.NewIteratorWithStartAndPrefixRequest{Start: req.Start, Prefix: req.Prefix})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &rpcdbpb.NewIteratorWithStartAndPrefixResponse{Id: resp.Id}, nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer) IteratorNext(ctx context.Context, req *rpcdbpb.IteratorNextRequest) (*rpcdbpb.IteratorNextResponse, error) {
|
||||
resp, err := g.svc.IteratorNext(ctx, &rpcdb.IteratorNextRequest{Id: req.Id})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
data := make([]*rpcdbpb.PutRequest, len(resp.Data))
|
||||
for i, e := range resp.Data {
|
||||
data[i] = &rpcdbpb.PutRequest{Key: e.Key, Value: e.Value}
|
||||
}
|
||||
return &rpcdbpb.IteratorNextResponse{Data: data}, nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer) IteratorError(ctx context.Context, req *rpcdbpb.IteratorErrorRequest) (*rpcdbpb.IteratorErrorResponse, error) {
|
||||
resp, err := g.svc.IteratorError(ctx, &rpcdb.IteratorErrorRequest{Id: req.Id})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &rpcdbpb.IteratorErrorResponse{Err: toPbErr(resp.Err)}, nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer) IteratorRelease(ctx context.Context, req *rpcdbpb.IteratorReleaseRequest) (*rpcdbpb.IteratorReleaseResponse, error) {
|
||||
resp, err := g.svc.IteratorRelease(ctx, &rpcdb.IteratorReleaseRequest{Id: req.Id})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &rpcdbpb.IteratorReleaseResponse{Err: toPbErr(resp.Err)}, nil
|
||||
}
|
||||
@@ -1,139 +0,0 @@
|
||||
//go:build grpc
|
||||
|
||||
// Copyright (C) 2019-2026, Lux Industries Inc. All rights reserved.
|
||||
// See the file LICENSE for licensing terms.
|
||||
|
||||
package rpcdb
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"github.com/luxfi/database/corruptabledb"
|
||||
"github.com/luxfi/database/dbtest"
|
||||
"github.com/luxfi/database/memdb"
|
||||
"github.com/luxfi/log"
|
||||
"github.com/luxfi/vm/rpc/grpcutils"
|
||||
|
||||
rpcdbpb "github.com/luxfi/node/proto/pb/rpcdb"
|
||||
)
|
||||
|
||||
type testGRPCDatabase struct {
|
||||
client *GRPCClient
|
||||
server *memdb.Database
|
||||
}
|
||||
|
||||
func setupGRPCDB(t testing.TB) *testGRPCDatabase {
|
||||
require := require.New(t)
|
||||
|
||||
db := &testGRPCDatabase{
|
||||
server: memdb.New(),
|
||||
}
|
||||
|
||||
listener, err := grpcutils.NewListener()
|
||||
require.NoError(err)
|
||||
serverCloser := grpcutils.ServerCloser{}
|
||||
|
||||
server := grpcutils.NewServer()
|
||||
rpcdbpb.RegisterDatabaseServer(server, NewGRPCServer(db.server))
|
||||
serverCloser.Add(server)
|
||||
|
||||
go grpcutils.Serve(listener, server)
|
||||
|
||||
conn, err := grpcutils.Dial(listener.Addr().String())
|
||||
require.NoError(err)
|
||||
|
||||
db.client = NewGRPCClient(rpcdbpb.NewDatabaseClient(conn))
|
||||
|
||||
t.Cleanup(func() {
|
||||
serverCloser.Stop()
|
||||
_ = conn.Close()
|
||||
_ = listener.Close()
|
||||
})
|
||||
|
||||
return db
|
||||
}
|
||||
|
||||
func TestGRPCInterface(t *testing.T) {
|
||||
for name, test := range dbtest.Tests {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
db := setupGRPCDB(t)
|
||||
test(t, db.client)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func FuzzGRPCKeyValue(f *testing.F) {
|
||||
db := setupGRPCDB(f)
|
||||
dbtest.FuzzKeyValue(f, db.client)
|
||||
}
|
||||
|
||||
func FuzzGRPCNewIteratorWithPrefix(f *testing.F) {
|
||||
db := setupGRPCDB(f)
|
||||
dbtest.FuzzNewIteratorWithPrefix(f, db.client)
|
||||
}
|
||||
|
||||
func FuzzGRPCNewIteratorWithStartAndPrefix(f *testing.F) {
|
||||
db := setupGRPCDB(f)
|
||||
dbtest.FuzzNewIteratorWithStartAndPrefix(f, db.client)
|
||||
}
|
||||
|
||||
func BenchmarkGRPCInterface(b *testing.B) {
|
||||
for _, size := range dbtest.BenchmarkSizes {
|
||||
keys, values := dbtest.SetupBenchmark(b, size[0], size[1], size[2])
|
||||
for name, bench := range dbtest.Benchmarks {
|
||||
b.Run(fmt.Sprintf("rpcdb_%d_pairs_%d_keys_%d_values_%s", size[0], size[1], size[2], name), func(b *testing.B) {
|
||||
db := setupGRPCDB(b)
|
||||
bench(b, db.client, keys, values)
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestGRPCHealthCheck(t *testing.T) {
|
||||
scenarios := []struct {
|
||||
name string
|
||||
testFn func(db *corruptabledb.Database) error
|
||||
wantErr bool
|
||||
wantErrMsg string
|
||||
}{
|
||||
{
|
||||
name: "healthcheck success",
|
||||
testFn: func(_ *corruptabledb.Database) error {
|
||||
return nil
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "healthcheck failed db closed",
|
||||
testFn: func(db *corruptabledb.Database) error {
|
||||
return db.Close()
|
||||
},
|
||||
wantErr: true,
|
||||
wantErrMsg: "closed",
|
||||
},
|
||||
}
|
||||
for _, scenario := range scenarios {
|
||||
t.Run(scenario.name, func(t *testing.T) {
|
||||
require := require.New(t)
|
||||
|
||||
baseDB := setupGRPCDB(t)
|
||||
db := corruptabledb.New(baseDB.server, log.NoLog{})
|
||||
defer db.Close()
|
||||
require.NoError(scenario.testFn(db))
|
||||
|
||||
_, err := db.HealthCheck(context.Background())
|
||||
if scenario.wantErr {
|
||||
require.Error(err) //nolint:forbidigo
|
||||
require.Contains(err.Error(), scenario.wantErrMsg)
|
||||
return
|
||||
}
|
||||
require.NoError(err)
|
||||
|
||||
_, err = baseDB.client.HealthCheck(context.Background())
|
||||
require.NoError(err)
|
||||
})
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user