rpcdb: decomplect into Layer A/B/C topology

Layer A — wire framing: github.com/luxfi/api/zap (unchanged)
Layer B — service spec: github.com/luxfi/proto/rpcdb (transport-agnostic
   data carriers, was node/proto/zap/rpcdb)
Layer C — service impl + transports: node/db/rpcdb/
   - service.go      transport-neutral Service wrapping database.Database
   - zap_server.go   ZAP transport adapter (default; used by cevm)
   - grpc_server.go  gRPC transport adapter (-tags=grpc)

One Service. Many transport adapters. Adding a transport is a new
adapter file wrapping *Service; storage logic stays in service.go.

Removed (-1390 LOC of duplicate/dead code):
 - node/proto/zap/rpcdb/rpcdb.go (moved to luxfi/proto/rpcdb)
 - node/proto/rpcdb/rpcdb_zap.go (dead HandlerRegistry path, no callers)
 - node/proto/rpcdb/rpcdb_grpc.go (duplicated gRPC server, replaced by
   node/db/rpcdb/grpc_server.go)
 - node/db/rpcdb/db.go (re-export shim, replaced by Service)

Consumers updated to import the canonical adapter at node/db/rpcdb:
 - vms/rpcchainvm/zap/client.go
 - vms/rpcchainvm/zap/client_dbserver_test.go
 - vms/rpcchainvm/zap/cevm_e2e_test.go

go.mod: add `replace github.com/luxfi/proto => ../proto` for in-tree
  development of the central wire-types module.

Tests: db/rpcdb (3/3) + rpcchainvm/zap (4/4 incl. cevm cross-process
e2e — KP_META written via ZAP db channel, no fallback to local zapdb).
This commit is contained in:
Hanzo AI
2026-05-15 15:47:44 -07:00
parent 697b9ccfc3
commit 8efcdff444
13 changed files with 1123 additions and 1390 deletions
+12 -1
View File
@@ -284,11 +284,22 @@ go build -tags=grpc # gRPC support (for testing/compatibility)
```
**Key Packages:**
- `github.com/luxfi/api/zap` - Core wire protocol and message types
- `github.com/luxfi/api/zap` - Core wire protocol and message types (Layer A)
- `github.com/luxfi/proto/rpcdb` - rpcdb service spec / data carriers (Layer B)
- `github.com/luxfi/node/db/rpcdb` - rpcdb Service + ZAP/gRPC transport adapters (Layer C)
- `github.com/luxfi/vm/rpc/sender` - p2p.Sender over ZAP/gRPC
- `vms/rpcchainvm/sender/` - Node-side sender implementation
- `vms/platformvm/warp/zwarp/` - ZAP-based warp signing client/server
**rpcdb Layered Topology (post-2026-05 reorg):**
- Layer A — wire framing: `github.com/luxfi/api/zap` (independent module)
- Layer B — rpcdb service spec: `github.com/luxfi/proto/rpcdb` (transport-agnostic data carriers)
- Layer C — rpcdb impl: `node/db/rpcdb/{service.go, grpc_server.go, zap_server.go}`
- `service.go` — transport-neutral `Service` wrapping `database.Database`
- `zap_server.go` (default) — ZAP transport adapter (used by cevm)
- `grpc_server.go` (`-tags=grpc`) — gRPC transport adapter
- One Service, many transport adapters. Adding a transport = new file wrapping `*Service`.
**Wire Protocol Format:**
```
[4 bytes: length][1 byte: message type][payload...]
-21
View File
@@ -1,21 +0,0 @@
//go:build grpc
// Copyright (C) 2019-2025, Lux Industries Inc. All rights reserved.
// See the file LICENSE for licensing terms.
// Package rpcdb re-exports the proto/rpcdb package for backwards compatibility
package rpcdb
import "github.com/luxfi/node/proto/rpcdb"
// Type aliases for backwards compatibility
type (
DatabaseClient = rpcdb.DatabaseClient
DatabaseServer = rpcdb.DatabaseServer
)
// Function aliases for backwards compatibility
var (
NewClient = rpcdb.NewClient
NewServer = rpcdb.NewServer
)
+159
View File
@@ -0,0 +1,159 @@
//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/proto/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/proto/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
}
+264
View File
@@ -0,0 +1,264 @@
// Copyright (C) 2019-2026, Lux Industries Inc. All rights reserved.
// See the file LICENSE for licensing terms.
// Package rpcdb is luxd's canonical Layer C home for the rpcdb service:
// one Service implementation backed by a database.Database, with a
// transport adapter per supported wire protocol.
//
// Layered topology:
//
// Layer A — wire framing (github.com/luxfi/api/zap)
// Layer B — service spec (data carriers) (github.com/luxfi/proto/rpcdb)
// Layer C — service impl + transports (this package)
//
// One Service. Many transport adapters. Adding a new transport is a
// new file that wraps `*Service`; the storage logic stays here and
// stays orthogonal to framing concerns.
//
// Files:
// - service.go (this) — transport-neutral Service struct + methods
// - grpc_server.go (build tag `grpc`) — gRPC adapter
// - zap_server.go (default) — ZAP adapter
package rpcdb
import (
"context"
"errors"
"sync"
"github.com/luxfi/database"
rpcdb "github.com/luxfi/proto/rpcdb"
)
// errUnknownIterator is the sentinel returned by IteratorNext / Error
// / Release when the caller hands in an iterator id that was never
// allocated (or has already been released). Transport adapters should
// surface this however their wire spec demands.
var errUnknownIterator = errors.New("rpcdb: unknown iterator")
// iterationBatchBytes is the upper bound on bytes batched into a
// single IteratorNext response. Picked to match the cevm-side
// expectation (RemoteZapDB::iterator_next reads up to ~128 KiB per
// roundtrip).
const iterationBatchBytes = 128 * 1024
// Service is the rpcdb service implementation. It holds one
// `database.Database` and a server-managed iterator pool, and exposes
// every rpcdb operation as a (request → response) method on Go data
// carriers from the proto/rpcdb wire-types package.
//
// Service is transport-agnostic: it knows nothing about gRPC, ZAP, or
// any other framing. Each transport adapter wraps `*Service` and
// translates its own wire format into these method calls.
type Service struct {
db database.Database
// iteratorMu serializes mutations of nextIteratorID and iterators.
// It does NOT serialize calls to a particular iterator — the
// underlying database.Iterator is documented as not safe for
// concurrent use, and the caller is expected to respect that.
iteratorMu sync.RWMutex
nextIteratorID uint64
iterators map[uint64]database.Iterator
}
// NewService wraps `db` so it can be served over any transport.
func NewService(db database.Database) *Service {
return &Service{
db: db,
iterators: make(map[uint64]database.Iterator),
}
}
// DB returns the underlying database. Adapters that need to surface
// transport-specific behavior (eg. health-check serialization) can
// reach in for it; everyone else uses the methods.
func (s *Service) DB() database.Database {
return s.db
}
// CloseIterators releases every server-side iterator. Adapters call
// this on shutdown so iterators don't leak past the listener.
func (s *Service) CloseIterators() {
s.iteratorMu.Lock()
defer s.iteratorMu.Unlock()
for id, it := range s.iterators {
it.Release()
delete(s.iterators, id)
}
}
// errToCode maps a database error to the wire enum.
func errToCode(err error) rpcdb.Error {
switch {
case err == nil:
return rpcdb.Error_ERROR_UNSPECIFIED
case errors.Is(err, database.ErrNotFound):
return rpcdb.Error_ERROR_NOT_FOUND
case errors.Is(err, database.ErrClosed):
return rpcdb.Error_ERROR_CLOSED
default:
// Any other error is reported as CLOSED on the wire — the
// proto enum doesn't have a richer error space. Adapters that
// can carry a transport-level error message should also pass
// the original err out-of-band.
return rpcdb.Error_ERROR_CLOSED
}
}
// CodeToErr maps a wire-typed Error back to a database sentinel. Used
// by remote clients on the receive side; lives here so server and
// client agree on the inverse of errToCode.
func CodeToErr(code rpcdb.Error) error {
switch code {
case rpcdb.Error_ERROR_UNSPECIFIED:
return nil
case rpcdb.Error_ERROR_NOT_FOUND:
return database.ErrNotFound
case rpcdb.Error_ERROR_CLOSED:
return database.ErrClosed
default:
return database.ErrClosed
}
}
// Has services rpcdb.Has.
func (s *Service) Has(_ context.Context, req *rpcdb.HasRequest) (*rpcdb.HasResponse, error) {
has, err := s.db.Has(req.Key)
if err != nil {
return &rpcdb.HasResponse{Has: false, Err: errToCode(err)}, nil
}
return &rpcdb.HasResponse{Has: has, Err: rpcdb.Error_ERROR_UNSPECIFIED}, nil
}
// Get services rpcdb.Get.
func (s *Service) Get(_ context.Context, req *rpcdb.GetRequest) (*rpcdb.GetResponse, error) {
value, err := s.db.Get(req.Key)
if err != nil {
return &rpcdb.GetResponse{Value: nil, Err: errToCode(err)}, nil
}
return &rpcdb.GetResponse{Value: value, Err: rpcdb.Error_ERROR_UNSPECIFIED}, nil
}
// Put services rpcdb.Put.
func (s *Service) Put(_ context.Context, req *rpcdb.PutRequest) (*rpcdb.PutResponse, error) {
return &rpcdb.PutResponse{Err: errToCode(s.db.Put(req.Key, req.Value))}, nil
}
// Delete services rpcdb.Delete.
func (s *Service) Delete(_ context.Context, req *rpcdb.DeleteRequest) (*rpcdb.DeleteResponse, error) {
return &rpcdb.DeleteResponse{Err: errToCode(s.db.Delete(req.Key))}, nil
}
// WriteBatch services rpcdb.WriteBatch.
func (s *Service) WriteBatch(_ context.Context, req *rpcdb.WriteBatchRequest) (*rpcdb.WriteBatchResponse, error) {
batch := s.db.NewBatch()
for _, p := range req.Puts {
if err := batch.Put(p.Key, p.Value); err != nil {
return &rpcdb.WriteBatchResponse{Err: errToCode(err)}, nil
}
}
for _, d := range req.Deletes {
if err := batch.Delete(d.Key); err != nil {
return &rpcdb.WriteBatchResponse{Err: errToCode(err)}, nil
}
}
return &rpcdb.WriteBatchResponse{Err: errToCode(batch.Write())}, nil
}
// Compact services rpcdb.Compact.
func (s *Service) Compact(_ context.Context, req *rpcdb.CompactRequest) (*rpcdb.CompactResponse, error) {
return &rpcdb.CompactResponse{Err: errToCode(s.db.Compact(req.Start, req.Limit))}, nil
}
// Close services rpcdb.Close.
func (s *Service) Close(_ context.Context, _ *rpcdb.CloseRequest) (*rpcdb.CloseResponse, error) {
return &rpcdb.CloseResponse{Err: errToCode(s.db.Close())}, nil
}
// HealthCheck services rpcdb.HealthCheck.
func (s *Service) HealthCheck(ctx context.Context) (*rpcdb.HealthCheckResponse, error) {
health, err := s.db.HealthCheck(ctx)
if err != nil {
return &rpcdb.HealthCheckResponse{Details: []byte(err.Error())}, nil
}
details := []byte("healthy")
if health != nil {
if str, ok := health.(string); ok {
details = []byte(str)
}
}
return &rpcdb.HealthCheckResponse{Details: details}, nil
}
// NewIteratorWithStartAndPrefix services rpcdb.NewIteratorWithStartAndPrefix.
func (s *Service) NewIteratorWithStartAndPrefix(_ context.Context, req *rpcdb.NewIteratorWithStartAndPrefixRequest) (*rpcdb.NewIteratorWithStartAndPrefixResponse, error) {
it := s.db.NewIteratorWithStartAndPrefix(req.Start, req.Prefix)
s.iteratorMu.Lock()
id := s.nextIteratorID
s.nextIteratorID++
s.iterators[id] = it
s.iteratorMu.Unlock()
return &rpcdb.NewIteratorWithStartAndPrefixResponse{Id: id}, nil
}
// IteratorNext services rpcdb.IteratorNext, returning up to
// iterationBatchBytes worth of (key, value) pairs per call. An empty
// Data slice signals exhaustion.
//
// Returned bytes are deep-copied so the caller can hold them past the
// next iterator advance; the underlying database.Iterator is allowed
// to recycle its internal buffers.
func (s *Service) IteratorNext(_ context.Context, req *rpcdb.IteratorNextRequest) (*rpcdb.IteratorNextResponse, error) {
s.iteratorMu.RLock()
it, ok := s.iterators[req.Id]
s.iteratorMu.RUnlock()
if !ok {
return nil, errUnknownIterator
}
var (
size int
data []*rpcdb.PutRequest
)
for size < iterationBatchBytes && it.Next() {
k := it.Key()
v := it.Value()
size += len(k) + len(v)
kc := make([]byte, len(k))
copy(kc, k)
vc := make([]byte, len(v))
copy(vc, v)
data = append(data, &rpcdb.PutRequest{Key: kc, Value: vc})
}
return &rpcdb.IteratorNextResponse{Data: data}, nil
}
// IteratorError services rpcdb.IteratorError.
func (s *Service) IteratorError(_ context.Context, req *rpcdb.IteratorErrorRequest) (*rpcdb.IteratorErrorResponse, error) {
s.iteratorMu.RLock()
it, ok := s.iterators[req.Id]
s.iteratorMu.RUnlock()
if !ok {
return &rpcdb.IteratorErrorResponse{Err: rpcdb.Error_ERROR_NOT_FOUND}, nil
}
return &rpcdb.IteratorErrorResponse{Err: errToCode(it.Error())}, nil
}
// IteratorRelease services rpcdb.IteratorRelease. Idempotent: calling
// it twice for the same id returns UNSPECIFIED on the second call.
func (s *Service) IteratorRelease(_ context.Context, req *rpcdb.IteratorReleaseRequest) (*rpcdb.IteratorReleaseResponse, error) {
s.iteratorMu.Lock()
it, ok := s.iterators[req.Id]
if !ok {
s.iteratorMu.Unlock()
return &rpcdb.IteratorReleaseResponse{Err: rpcdb.Error_ERROR_UNSPECIFIED}, nil
}
delete(s.iterators, req.Id)
s.iteratorMu.Unlock()
dbErr := it.Error()
it.Release()
return &rpcdb.IteratorReleaseResponse{Err: errToCode(dbErr)}, nil
}
+407
View File
@@ -0,0 +1,407 @@
// Copyright (C) 2019-2026, Lux Industries Inc. All rights reserved.
// See the file LICENSE for licensing terms.
package rpcdb
import (
"context"
"errors"
"fmt"
"net"
"sync"
zapwire "github.com/luxfi/api/zap"
"github.com/luxfi/database"
rpcdb "github.com/luxfi/proto/rpcdb"
)
// ZAP db channel MsgType IDs. These are the wire-level Layer-A
// dispatch tags for the rpcdb service over ZAP. They live in their
// own listener (one per VM plugin); they do NOT collide with the VM
// lifecycle MsgTypes (1..31) or sender (40..49) or warp (50..59)
// because each listener has its own dispatch table.
//
// Order is the canonical Lane A assignment — the cevm side hard-codes
// these numbers in `RemoteZapDB::call`. Do not reorder.
const (
MsgDBHas zapwire.MessageType = 1
MsgDBGet zapwire.MessageType = 2
MsgDBPut zapwire.MessageType = 3
MsgDBDelete zapwire.MessageType = 4
MsgDBWriteBatch zapwire.MessageType = 5
MsgDBCompact zapwire.MessageType = 6
MsgDBClose zapwire.MessageType = 7
MsgDBHealthCheck zapwire.MessageType = 8
MsgDBIteratorNew zapwire.MessageType = 9
MsgDBIteratorNext zapwire.MessageType = 10
MsgDBIteratorError zapwire.MessageType = 11
MsgDBIteratorRelease zapwire.MessageType = 12
)
// Wire-byte values for the Error enum. These match
// rpcdb.Error_ERROR_* but are encoded as a single byte on the ZAP
// wire so the cevm side can decode them without pulling in protobuf.
const (
dbErrUnspecified uint8 = 0
dbErrClosed uint8 = 1
dbErrNotFound uint8 = 2
)
func errCodeToByte(code rpcdb.Error) uint8 {
switch code {
case rpcdb.Error_ERROR_CLOSED:
return dbErrClosed
case rpcdb.Error_ERROR_NOT_FOUND:
return dbErrNotFound
default:
return dbErrUnspecified
}
}
// ZAPServer is the ZAP transport adapter for the rpcdb Service. It
// owns the listener and the dispatch loop; the actual storage logic
// lives in *Service. To swap the wire format (eg. to add framing
// over a different reliable byte stream) write a new file with a new
// adapter — never edit Service.
type ZAPServer struct {
svc *Service
listener *zapwire.Listener
server *zapwire.Server
// closeOnce guards Close so callers can call it from both the
// client lifecycle and a deferred test cleanup without panicking
// on a double-close of the underlying listener.
closeOnce sync.Once
}
// NewZAPServer wraps a database.Database for serving over ZAP.
//
// Equivalent to NewZAPServerFromService(NewService(db)) — kept as a
// one-liner because cevm-side test fixtures and the production
// rpcchainvm/zap path both build it this way.
func NewZAPServer(db database.Database) *ZAPServer {
return NewZAPServerFromService(NewService(db))
}
// NewZAPServerFromService wraps an existing Service for serving over
// ZAP. Useful when a single Service needs multiple transport adapters
// at once (eg. tests that want both ZAP and direct in-process calls).
func NewZAPServerFromService(svc *Service) *ZAPServer {
return &ZAPServer{svc: svc}
}
// Listen binds the ZAP listener to addr.
func (s *ZAPServer) Listen(addr string) error {
listener, err := zapwire.Listen(addr, nil)
if err != nil {
return fmt.Errorf("rpcdb zap: listen %s: %w", addr, err)
}
s.listener = listener
s.server = zapwire.NewServer(listener, zapwire.HandlerFunc(s.handle))
return nil
}
// ListenOn wraps an existing net.Listener (for ephemeral ports etc.).
func (s *ZAPServer) ListenOn(raw net.Listener) {
s.listener = zapwire.NewListener(raw, nil)
s.server = zapwire.NewServer(s.listener, zapwire.HandlerFunc(s.handle))
}
// Addr returns the bound address (or nil if not yet listening).
func (s *ZAPServer) Addr() net.Addr {
if s.listener == nil {
return nil
}
return s.listener.Addr()
}
// Serve blocks until ctx is cancelled or Close is called.
func (s *ZAPServer) Serve(ctx context.Context) error {
if s.server == nil {
return errors.New("rpcdb zap: not initialized — call Listen first")
}
return s.server.Serve(ctx)
}
// Close releases all iterators and closes the listener. Caller is
// expected to also cancel the context passed to Serve so the accept
// loop exits cleanly. Safe to call multiple times.
//
// We deliberately do NOT call s.server.Close() because the upstream
// zapwire.Server.Close races with in-flight accept (it nils its conns
// map mid-Serve, causing "assignment to entry in nil map"). Closing
// the listener is enough to make Accept return; ctx cancellation does
// the rest.
func (s *ZAPServer) Close() error {
var err error
s.closeOnce.Do(func() {
s.svc.CloseIterators()
if s.listener != nil {
err = s.listener.Close()
}
})
return err
}
func (s *ZAPServer) handle(ctx context.Context, msgType zapwire.MessageType, payload []byte) (zapwire.MessageType, []byte, error) {
switch msgType {
case MsgDBHas:
return s.handleHas(ctx, payload)
case MsgDBGet:
return s.handleGet(ctx, payload)
case MsgDBPut:
return s.handlePut(ctx, payload)
case MsgDBDelete:
return s.handleDelete(ctx, payload)
case MsgDBWriteBatch:
return s.handleWriteBatch(ctx, payload)
case MsgDBCompact:
return s.handleCompact(ctx, payload)
case MsgDBClose:
return s.handleClose(ctx, payload)
case MsgDBHealthCheck:
return s.handleHealthCheck(ctx)
case MsgDBIteratorNew:
return s.handleIteratorNew(ctx, payload)
case MsgDBIteratorNext:
return s.handleIteratorNext(ctx, payload)
case MsgDBIteratorError:
return s.handleIteratorError(ctx, payload)
case MsgDBIteratorRelease:
return s.handleIteratorRelease(ctx, payload)
default:
return 0, nil, fmt.Errorf("rpcdb zap: unknown msg type %d", msgType)
}
}
func (s *ZAPServer) handleHas(ctx context.Context, payload []byte) (zapwire.MessageType, []byte, error) {
r := zapwire.NewReader(payload)
key, err := r.ReadBytes()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb Has decode: %w", err)
}
resp, err := s.svc.Has(ctx, &rpcdb.HasRequest{Key: key})
if err != nil {
return 0, nil, err
}
buf := zapwire.GetBuffer()
defer zapwire.PutBuffer(buf)
if resp.Has {
buf.WriteUint8(1)
} else {
buf.WriteUint8(0)
}
buf.WriteUint8(errCodeToByte(resp.Err))
return MsgDBHas, append([]byte(nil), buf.Bytes()...), nil
}
func (s *ZAPServer) handleGet(ctx context.Context, payload []byte) (zapwire.MessageType, []byte, error) {
r := zapwire.NewReader(payload)
key, err := r.ReadBytes()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb Get decode: %w", err)
}
resp, err := s.svc.Get(ctx, &rpcdb.GetRequest{Key: key})
if err != nil {
return 0, nil, err
}
buf := zapwire.GetBuffer()
defer zapwire.PutBuffer(buf)
buf.WriteBytes(resp.Value)
buf.WriteUint8(errCodeToByte(resp.Err))
return MsgDBGet, append([]byte(nil), buf.Bytes()...), nil
}
func (s *ZAPServer) handlePut(ctx context.Context, payload []byte) (zapwire.MessageType, []byte, error) {
r := zapwire.NewReader(payload)
key, err := r.ReadBytes()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb Put decode key: %w", err)
}
value, err := r.ReadBytes()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb Put decode value: %w", err)
}
resp, err := s.svc.Put(ctx, &rpcdb.PutRequest{Key: key, Value: value})
if err != nil {
return 0, nil, err
}
buf := zapwire.GetBuffer()
defer zapwire.PutBuffer(buf)
buf.WriteUint8(errCodeToByte(resp.Err))
return MsgDBPut, append([]byte(nil), buf.Bytes()...), nil
}
func (s *ZAPServer) handleDelete(ctx context.Context, payload []byte) (zapwire.MessageType, []byte, error) {
r := zapwire.NewReader(payload)
key, err := r.ReadBytes()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb Delete decode: %w", err)
}
resp, err := s.svc.Delete(ctx, &rpcdb.DeleteRequest{Key: key})
if err != nil {
return 0, nil, err
}
buf := zapwire.GetBuffer()
defer zapwire.PutBuffer(buf)
buf.WriteUint8(errCodeToByte(resp.Err))
return MsgDBDelete, append([]byte(nil), buf.Bytes()...), nil
}
func (s *ZAPServer) handleWriteBatch(ctx context.Context, payload []byte) (zapwire.MessageType, []byte, error) {
r := zapwire.NewReader(payload)
nputs, err := r.ReadUint32()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb WriteBatch decode nputs: %w", err)
}
puts := make([]*rpcdb.PutRequest, 0, nputs)
for i := uint32(0); i < nputs; i++ {
key, err := r.ReadBytes()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb WriteBatch put[%d] key: %w", i, err)
}
value, err := r.ReadBytes()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb WriteBatch put[%d] value: %w", i, err)
}
puts = append(puts, &rpcdb.PutRequest{Key: key, Value: value})
}
ndels, err := r.ReadUint32()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb WriteBatch decode ndels: %w", err)
}
dels := make([]*rpcdb.DeleteRequest, 0, ndels)
for i := uint32(0); i < ndels; i++ {
key, err := r.ReadBytes()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb WriteBatch del[%d]: %w", i, err)
}
dels = append(dels, &rpcdb.DeleteRequest{Key: key})
}
resp, err := s.svc.WriteBatch(ctx, &rpcdb.WriteBatchRequest{Puts: puts, Deletes: dels})
if err != nil {
return 0, nil, err
}
buf := zapwire.GetBuffer()
defer zapwire.PutBuffer(buf)
buf.WriteUint8(errCodeToByte(resp.Err))
return MsgDBWriteBatch, append([]byte(nil), buf.Bytes()...), nil
}
func (s *ZAPServer) handleCompact(ctx context.Context, payload []byte) (zapwire.MessageType, []byte, error) {
r := zapwire.NewReader(payload)
start, err := r.ReadBytes()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb Compact decode start: %w", err)
}
limit, err := r.ReadBytes()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb Compact decode limit: %w", err)
}
resp, err := s.svc.Compact(ctx, &rpcdb.CompactRequest{Start: start, Limit: limit})
if err != nil {
return 0, nil, err
}
buf := zapwire.GetBuffer()
defer zapwire.PutBuffer(buf)
buf.WriteUint8(errCodeToByte(resp.Err))
return MsgDBCompact, append([]byte(nil), buf.Bytes()...), nil
}
func (s *ZAPServer) handleClose(ctx context.Context, _ []byte) (zapwire.MessageType, []byte, error) {
resp, err := s.svc.Close(ctx, &rpcdb.CloseRequest{})
if err != nil {
return 0, nil, err
}
buf := zapwire.GetBuffer()
defer zapwire.PutBuffer(buf)
buf.WriteUint8(errCodeToByte(resp.Err))
return MsgDBClose, append([]byte(nil), buf.Bytes()...), nil
}
func (s *ZAPServer) handleHealthCheck(ctx context.Context) (zapwire.MessageType, []byte, error) {
resp, err := s.svc.HealthCheck(ctx)
if err != nil {
return 0, nil, err
}
buf := zapwire.GetBuffer()
defer zapwire.PutBuffer(buf)
buf.WriteBytes(resp.Details)
return MsgDBHealthCheck, append([]byte(nil), buf.Bytes()...), nil
}
func (s *ZAPServer) handleIteratorNew(ctx context.Context, payload []byte) (zapwire.MessageType, []byte, error) {
r := zapwire.NewReader(payload)
start, err := r.ReadBytes()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb IteratorNew decode start: %w", err)
}
prefix, err := r.ReadBytes()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb IteratorNew decode prefix: %w", err)
}
resp, err := s.svc.NewIteratorWithStartAndPrefix(ctx, &rpcdb.NewIteratorWithStartAndPrefixRequest{Start: start, Prefix: prefix})
if err != nil {
return 0, nil, err
}
buf := zapwire.GetBuffer()
defer zapwire.PutBuffer(buf)
buf.WriteUint64(resp.Id)
return MsgDBIteratorNew, append([]byte(nil), buf.Bytes()...), nil
}
func (s *ZAPServer) handleIteratorNext(ctx context.Context, payload []byte) (zapwire.MessageType, []byte, error) {
r := zapwire.NewReader(payload)
id, err := r.ReadUint64()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb IteratorNext decode id: %w", err)
}
resp, err := s.svc.IteratorNext(ctx, &rpcdb.IteratorNextRequest{Id: id})
if err != nil {
return 0, nil, err
}
buf := zapwire.GetBuffer()
defer zapwire.PutBuffer(buf)
buf.WriteUint32(uint32(len(resp.Data)))
for _, e := range resp.Data {
buf.WriteBytes(e.Key)
buf.WriteBytes(e.Value)
}
return MsgDBIteratorNext, append([]byte(nil), buf.Bytes()...), nil
}
func (s *ZAPServer) handleIteratorError(ctx context.Context, payload []byte) (zapwire.MessageType, []byte, error) {
r := zapwire.NewReader(payload)
id, err := r.ReadUint64()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb IteratorError decode id: %w", err)
}
resp, err := s.svc.IteratorError(ctx, &rpcdb.IteratorErrorRequest{Id: id})
if err != nil {
return 0, nil, err
}
buf := zapwire.GetBuffer()
defer zapwire.PutBuffer(buf)
buf.WriteUint8(errCodeToByte(resp.Err))
return MsgDBIteratorError, append([]byte(nil), buf.Bytes()...), nil
}
func (s *ZAPServer) handleIteratorRelease(ctx context.Context, payload []byte) (zapwire.MessageType, []byte, error) {
r := zapwire.NewReader(payload)
id, err := r.ReadUint64()
if err != nil {
return 0, nil, fmt.Errorf("rpcdb IteratorRelease decode id: %w", err)
}
resp, err := s.svc.IteratorRelease(ctx, &rpcdb.IteratorReleaseRequest{Id: id})
if err != nil {
return 0, nil, err
}
buf := zapwire.GetBuffer()
defer zapwire.PutBuffer(buf)
buf.WriteUint8(errCodeToByte(resp.Err))
return MsgDBIteratorRelease, append([]byte(nil), buf.Bytes()...), nil
}
+275
View File
@@ -0,0 +1,275 @@
// Copyright (C) 2019-2026, Lux Industries Inc. All rights reserved.
// See the file LICENSE for licensing terms.
package rpcdb
import (
"context"
"net"
"testing"
"time"
"github.com/stretchr/testify/require"
zapwire "github.com/luxfi/api/zap"
"github.com/luxfi/database"
"github.com/luxfi/database/memdb"
)
// startZAPServer spawns a ZAPServer over an in-memory listener bound to
// 127.0.0.1:0, returns the addr + the server (for shutdown).
func startZAPServer(t *testing.T, db database.Database) (string, *ZAPServer, context.CancelFunc) {
t.Helper()
listener, err := net.Listen("tcp", "127.0.0.1:0")
require.NoError(t, err)
s := NewZAPServer(db)
s.ListenOn(listener)
ctx, cancel := context.WithCancel(context.Background())
go func() {
_ = s.Serve(ctx)
}()
// Wait for Serve to fully install accept loop.
time.Sleep(50 * time.Millisecond)
return listener.Addr().String(), s, cancel
}
// TestZAPServer_HasGetPutDelete exercises the four primitive ops.
func TestZAPServer_HasGetPutDelete(t *testing.T) {
require := require.New(t)
mem := memdb.New()
addr, srv, cancel := startZAPServer(t, mem)
defer cancel()
defer srv.Close()
conn, err := zapwire.Dial(context.Background(), addr, nil)
require.NoError(err)
defer conn.Close()
ctx := context.Background()
// Has (key absent) → has=false, err=NotFound
{
buf := zapwire.GetBuffer()
buf.WriteBytes([]byte("foo"))
_, resp, err := conn.Call(ctx, MsgDBHas, buf.Bytes())
zapwire.PutBuffer(buf)
require.NoError(err)
r := zapwire.NewReader(resp)
hasB, _ := r.ReadUint8()
errB, _ := r.ReadUint8()
require.Equal(uint8(0), hasB)
require.Equal(dbErrUnspecified, errB) // memdb returns no err on Has-miss
}
// Put
{
buf := zapwire.GetBuffer()
buf.WriteBytes([]byte("foo"))
buf.WriteBytes([]byte("bar"))
_, resp, err := conn.Call(ctx, MsgDBPut, buf.Bytes())
zapwire.PutBuffer(buf)
require.NoError(err)
r := zapwire.NewReader(resp)
errB, _ := r.ReadUint8()
require.Equal(dbErrUnspecified, errB)
}
// Has (key present) → has=true
{
buf := zapwire.GetBuffer()
buf.WriteBytes([]byte("foo"))
_, resp, err := conn.Call(ctx, MsgDBHas, buf.Bytes())
zapwire.PutBuffer(buf)
require.NoError(err)
r := zapwire.NewReader(resp)
hasB, _ := r.ReadUint8()
errB, _ := r.ReadUint8()
require.Equal(uint8(1), hasB)
require.Equal(dbErrUnspecified, errB)
}
// Get
{
buf := zapwire.GetBuffer()
buf.WriteBytes([]byte("foo"))
_, resp, err := conn.Call(ctx, MsgDBGet, buf.Bytes())
zapwire.PutBuffer(buf)
require.NoError(err)
r := zapwire.NewReader(resp)
val, _ := r.ReadBytes()
errB, _ := r.ReadUint8()
require.Equal([]byte("bar"), val)
require.Equal(dbErrUnspecified, errB)
}
// Get (missing) → empty value, err=NotFound
{
buf := zapwire.GetBuffer()
buf.WriteBytes([]byte("absent"))
_, resp, err := conn.Call(ctx, MsgDBGet, buf.Bytes())
zapwire.PutBuffer(buf)
require.NoError(err)
r := zapwire.NewReader(resp)
val, _ := r.ReadBytes()
errB, _ := r.ReadUint8()
require.Empty(val)
require.Equal(dbErrNotFound, errB)
}
// Delete
{
buf := zapwire.GetBuffer()
buf.WriteBytes([]byte("foo"))
_, resp, err := conn.Call(ctx, MsgDBDelete, buf.Bytes())
zapwire.PutBuffer(buf)
require.NoError(err)
r := zapwire.NewReader(resp)
errB, _ := r.ReadUint8()
require.Equal(dbErrUnspecified, errB)
}
// Get post-delete → NotFound
{
buf := zapwire.GetBuffer()
buf.WriteBytes([]byte("foo"))
_, resp, err := conn.Call(ctx, MsgDBGet, buf.Bytes())
zapwire.PutBuffer(buf)
require.NoError(err)
r := zapwire.NewReader(resp)
_, _ = r.ReadBytes()
errB, _ := r.ReadUint8()
require.Equal(dbErrNotFound, errB)
}
}
// TestZAPServer_WriteBatch_AndIterate covers batched writes plus the
// iterator path used by cevm's load_persisted snapshot_with_prefix.
func TestZAPServer_WriteBatch_AndIterate(t *testing.T) {
require := require.New(t)
mem := memdb.New()
addr, srv, cancel := startZAPServer(t, mem)
defer cancel()
defer srv.Close()
conn, err := zapwire.Dial(context.Background(), addr, nil)
require.NoError(err)
defer conn.Close()
ctx := context.Background()
// WriteBatch: 3 puts (prefix p:), 1 delete (no-op since absent).
{
buf := zapwire.GetBuffer()
buf.WriteUint32(3)
buf.WriteBytes([]byte("p:1"))
buf.WriteBytes([]byte("v1"))
buf.WriteBytes([]byte("p:2"))
buf.WriteBytes([]byte("v2"))
buf.WriteBytes([]byte("q:3"))
buf.WriteBytes([]byte("v3"))
buf.WriteUint32(1)
buf.WriteBytes([]byte("absent"))
_, resp, err := conn.Call(ctx, MsgDBWriteBatch, buf.Bytes())
zapwire.PutBuffer(buf)
require.NoError(err)
r := zapwire.NewReader(resp)
errB, _ := r.ReadUint8()
require.Equal(dbErrUnspecified, errB)
}
// IteratorNew with prefix "p:"
var iterID uint64
{
buf := zapwire.GetBuffer()
buf.WriteBytes(nil) // start
buf.WriteBytes([]byte("p:")) // prefix
_, resp, err := conn.Call(ctx, MsgDBIteratorNew, buf.Bytes())
zapwire.PutBuffer(buf)
require.NoError(err)
r := zapwire.NewReader(resp)
iterID, _ = r.ReadUint64()
}
// IteratorNext — collect.
collected := map[string]string{}
for {
buf := zapwire.GetBuffer()
buf.WriteUint64(iterID)
_, resp, err := conn.Call(ctx, MsgDBIteratorNext, buf.Bytes())
zapwire.PutBuffer(buf)
require.NoError(err)
r := zapwire.NewReader(resp)
n, _ := r.ReadUint32()
if n == 0 {
break
}
for i := uint32(0); i < n; i++ {
k, _ := r.ReadBytes()
v, _ := r.ReadBytes()
collected[string(k)] = string(v)
}
}
// IteratorRelease
{
buf := zapwire.GetBuffer()
buf.WriteUint64(iterID)
_, resp, err := conn.Call(ctx, MsgDBIteratorRelease, buf.Bytes())
zapwire.PutBuffer(buf)
require.NoError(err)
r := zapwire.NewReader(resp)
errB, _ := r.ReadUint8()
require.Equal(dbErrUnspecified, errB)
}
require.Equal(map[string]string{"p:1": "v1", "p:2": "v2"}, collected)
}
// TestZAPServer_ParityWithMemDB writes/reads via wire and via direct
// memdb Get to confirm both paths see the same state — proves the ZAP
// round-trip preserves bytes exactly.
func TestZAPServer_ParityWithMemDB(t *testing.T) {
require := require.New(t)
mem := memdb.New()
addr, srv, cancel := startZAPServer(t, mem)
defer cancel()
defer srv.Close()
conn, err := zapwire.Dial(context.Background(), addr, nil)
require.NoError(err)
defer conn.Close()
ctx := context.Background()
key := []byte{0x42, 0x00, 0x01, 0xff}
value := make([]byte, 256)
for i := range value {
value[i] = byte(i)
}
// Put via wire.
buf := zapwire.GetBuffer()
buf.WriteBytes(key)
buf.WriteBytes(value)
_, _, err = conn.Call(ctx, MsgDBPut, buf.Bytes())
zapwire.PutBuffer(buf)
require.NoError(err)
// Read directly.
got, dErr := mem.Get(key)
require.NoError(dErr)
require.Equal(value, got)
// Read via wire.
buf2 := zapwire.GetBuffer()
buf2.WriteBytes(key)
_, resp, err := conn.Call(ctx, MsgDBGet, buf2.Bytes())
zapwire.PutBuffer(buf2)
require.NoError(err)
r := zapwire.NewReader(resp)
wireVal, _ := r.ReadBytes()
require.Equal(value, wireVal)
}
+3
View File
@@ -226,6 +226,7 @@ require (
github.com/luxfi/formatting v1.0.1
github.com/luxfi/go-bip32 v1.0.2
github.com/luxfi/math/big v0.1.0 // indirect
github.com/luxfi/proto v0.0.0-00010101000000-000000000000
github.com/luxfi/sampler v1.0.0 // indirect
github.com/luxfi/tls v1.0.3 // indirect
github.com/mattn/go-colorable v0.1.14 // indirect
@@ -241,3 +242,5 @@ require (
)
exclude github.com/ethereum/go-ethereum v1.10.26
replace github.com/luxfi/proto => ../proto
-548
View File
@@ -1,548 +0,0 @@
//go:build grpc
// Copyright (C) 2019-2025, Lux Industries Inc. All rights reserved.
// See the file LICENSE for licensing terms.
// Package rpcdb provides database RPC functionality
package rpcdb
import (
"context"
"errors"
"io"
"sync"
"github.com/luxfi/database"
rpcdbpb "github.com/luxfi/node/proto/pb/rpcdb"
"google.golang.org/protobuf/types/known/emptypb"
)
var (
_ database.Database = (*DatabaseClient)(nil)
_ rpcdbpb.DatabaseServer = (*DatabaseServer)(nil)
)
// DatabaseServer is a database server that listens over RPC.
type DatabaseServer struct {
rpcdbpb.UnimplementedDatabaseServer
db database.Database
iteratorLock sync.RWMutex
nextIteratorID uint64
iterators map[uint64]database.Iterator
}
// NewServer returns a new database server
func NewServer(db database.Database) *DatabaseServer {
return &DatabaseServer{
db: db,
iterators: make(map[uint64]database.Iterator),
}
}
func (db *DatabaseServer) Has(ctx context.Context, req *rpcdbpb.HasRequest) (*rpcdbpb.HasResponse, error) {
has, err := db.db.Has(req.Key)
if err != nil {
return &rpcdbpb.HasResponse{
Has: false,
Err: errorToRPCError(err),
}, nil
}
return &rpcdbpb.HasResponse{
Has: has,
Err: rpcdbpb.Error_ERROR_UNSPECIFIED,
}, nil
}
func (db *DatabaseServer) Get(ctx context.Context, req *rpcdbpb.GetRequest) (*rpcdbpb.GetResponse, error) {
value, err := db.db.Get(req.Key)
if err != nil {
return &rpcdbpb.GetResponse{
Value: nil,
Err: errorToRPCError(err),
}, nil
}
return &rpcdbpb.GetResponse{
Value: value,
Err: rpcdbpb.Error_ERROR_UNSPECIFIED,
}, nil
}
func (db *DatabaseServer) Put(ctx context.Context, req *rpcdbpb.PutRequest) (*rpcdbpb.PutResponse, error) {
err := db.db.Put(req.Key, req.Value)
return &rpcdbpb.PutResponse{
Err: errorToRPCError(err),
}, nil
}
func (db *DatabaseServer) Delete(ctx context.Context, req *rpcdbpb.DeleteRequest) (*rpcdbpb.DeleteResponse, error) {
err := db.db.Delete(req.Key)
return &rpcdbpb.DeleteResponse{
Err: errorToRPCError(err),
}, nil
}
func (db *DatabaseServer) WriteBatch(ctx context.Context, req *rpcdbpb.WriteBatchRequest) (*rpcdbpb.WriteBatchResponse, error) {
batch := db.db.NewBatch()
for _, put := range req.Puts {
if err := batch.Put(put.Key, put.Value); err != nil {
return &rpcdbpb.WriteBatchResponse{
Err: errorToRPCError(err),
}, nil
}
}
for _, del := range req.Deletes {
if err := batch.Delete(del.Key); err != nil {
return &rpcdbpb.WriteBatchResponse{
Err: errorToRPCError(err),
}, nil
}
}
err := batch.Write()
return &rpcdbpb.WriteBatchResponse{
Err: errorToRPCError(err),
}, nil
}
func (db *DatabaseServer) Compact(ctx context.Context, req *rpcdbpb.CompactRequest) (*rpcdbpb.CompactResponse, error) {
err := db.db.Compact(req.Start, req.Limit)
return &rpcdbpb.CompactResponse{
Err: errorToRPCError(err),
}, nil
}
func (db *DatabaseServer) Close(ctx context.Context, req *rpcdbpb.CloseRequest) (*rpcdbpb.CloseResponse, error) {
err := db.db.Close()
return &rpcdbpb.CloseResponse{
Err: errorToRPCError(err),
}, nil
}
func (db *DatabaseServer) HealthCheck(ctx context.Context, req *emptypb.Empty) (*rpcdbpb.HealthCheckResponse, error) {
health, err := db.db.HealthCheck(ctx)
if err != nil {
return &rpcdbpb.HealthCheckResponse{
Details: []byte(err.Error()),
}, nil
}
// Convert health data to bytes
healthBytes := []byte("healthy")
if health != nil {
if str, ok := health.(string); ok {
healthBytes = []byte(str)
}
}
return &rpcdbpb.HealthCheckResponse{
Details: healthBytes,
}, nil
}
func (db *DatabaseServer) NewIteratorWithStartAndPrefix(ctx context.Context, req *rpcdbpb.NewIteratorWithStartAndPrefixRequest) (*rpcdbpb.NewIteratorWithStartAndPrefixResponse, error) {
it := db.db.NewIteratorWithStartAndPrefix(req.Start, req.Prefix)
db.iteratorLock.Lock()
id := db.nextIteratorID
db.nextIteratorID++
db.iterators[id] = it
db.iteratorLock.Unlock()
return &rpcdbpb.NewIteratorWithStartAndPrefixResponse{
Id: id,
}, nil
}
func (db *DatabaseServer) IteratorNext(ctx context.Context, req *rpcdbpb.IteratorNextRequest) (*rpcdbpb.IteratorNextResponse, error) {
db.iteratorLock.RLock()
it, ok := db.iterators[req.Id]
db.iteratorLock.RUnlock()
if !ok {
return &rpcdbpb.IteratorNextResponse{
Data: nil,
}, nil
}
// Collect data until we have a reasonable batch or iterator is exhausted
var data []*rpcdbpb.PutRequest
const maxBatchSize = 100
for i := 0; i < maxBatchSize && it.Next(); i++ {
data = append(data, &rpcdbpb.PutRequest{
Key: it.Key(),
Value: it.Value(),
})
}
return &rpcdbpb.IteratorNextResponse{
Data: data,
}, nil
}
func (db *DatabaseServer) IteratorError(ctx context.Context, req *rpcdbpb.IteratorErrorRequest) (*rpcdbpb.IteratorErrorResponse, error) {
db.iteratorLock.RLock()
it, ok := db.iterators[req.Id]
db.iteratorLock.RUnlock()
if !ok {
return &rpcdbpb.IteratorErrorResponse{
Err: rpcdbpb.Error_ERROR_NOT_FOUND,
}, nil
}
err := it.Error()
return &rpcdbpb.IteratorErrorResponse{
Err: errorToRPCError(err),
}, nil
}
func (db *DatabaseServer) IteratorRelease(ctx context.Context, req *rpcdbpb.IteratorReleaseRequest) (*rpcdbpb.IteratorReleaseResponse, error) {
db.iteratorLock.Lock()
it, ok := db.iterators[req.Id]
if ok {
it.Release()
delete(db.iterators, req.Id)
}
db.iteratorLock.Unlock()
var err rpcdbpb.Error
if !ok {
err = rpcdbpb.Error_ERROR_NOT_FOUND
} else {
err = rpcdbpb.Error_ERROR_UNSPECIFIED
}
return &rpcdbpb.IteratorReleaseResponse{
Err: err,
}, nil
}
// DatabaseClient is a database client
type DatabaseClient struct {
client rpcdbpb.DatabaseClient
}
// NewClient returns a new database client
func NewClient(client rpcdbpb.DatabaseClient) *DatabaseClient {
return &DatabaseClient{
client: client,
}
}
func (db *DatabaseClient) Has(key []byte) (bool, error) {
resp, err := db.client.Has(context.Background(), &rpcdbpb.HasRequest{
Key: key,
})
if err != nil {
return false, err
}
if resp.Err != rpcdbpb.Error_ERROR_UNSPECIFIED {
return false, rpcErrorToError(resp.Err)
}
return resp.Has, nil
}
func (db *DatabaseClient) Get(key []byte) ([]byte, error) {
resp, err := db.client.Get(context.Background(), &rpcdbpb.GetRequest{
Key: key,
})
if err != nil {
return nil, err
}
if resp.Err != rpcdbpb.Error_ERROR_UNSPECIFIED {
return nil, rpcErrorToError(resp.Err)
}
return resp.Value, nil
}
func (db *DatabaseClient) Put(key []byte, value []byte) error {
resp, err := db.client.Put(context.Background(), &rpcdbpb.PutRequest{
Key: key,
Value: value,
})
if err != nil {
return err
}
if resp.Err != rpcdbpb.Error_ERROR_UNSPECIFIED {
return rpcErrorToError(resp.Err)
}
return nil
}
func (db *DatabaseClient) Delete(key []byte) error {
resp, err := db.client.Delete(context.Background(), &rpcdbpb.DeleteRequest{
Key: key,
})
if err != nil {
return err
}
if resp.Err != rpcdbpb.Error_ERROR_UNSPECIFIED {
return rpcErrorToError(resp.Err)
}
return nil
}
func (db *DatabaseClient) NewBatch() database.Batch {
return &batch{
db: db,
puts: []*rpcdbpb.PutRequest{},
deletes: []*rpcdbpb.DeleteRequest{},
}
}
func (db *DatabaseClient) NewIterator() database.Iterator {
return db.NewIteratorWithStartAndPrefix(nil, nil)
}
func (db *DatabaseClient) NewIteratorWithStart(start []byte) database.Iterator {
return db.NewIteratorWithStartAndPrefix(start, nil)
}
func (db *DatabaseClient) NewIteratorWithPrefix(prefix []byte) database.Iterator {
return db.NewIteratorWithStartAndPrefix(nil, prefix)
}
func (db *DatabaseClient) NewIteratorWithStartAndPrefix(start []byte, prefix []byte) database.Iterator {
resp, err := db.client.NewIteratorWithStartAndPrefix(context.Background(), &rpcdbpb.NewIteratorWithStartAndPrefixRequest{
Start: start,
Prefix: prefix,
})
if err != nil {
return &iterator{
db: db,
err: err,
}
}
return &iterator{
db: db,
id: resp.Id,
}
}
// Backup is not supported over the RPC database client.
func (db *DatabaseClient) Backup(_ io.Writer, _ uint64) (uint64, error) {
return 0, errors.New("rpcdb: backup not supported")
}
// Load is not supported over the RPC database client.
func (db *DatabaseClient) Load(_ io.Reader) error {
return errors.New("rpcdb: load not supported")
}
func (db *DatabaseClient) Compact(start []byte, limit []byte) error {
resp, err := db.client.Compact(context.Background(), &rpcdbpb.CompactRequest{
Start: start,
Limit: limit,
})
if err != nil {
return err
}
if resp.Err != rpcdbpb.Error_ERROR_UNSPECIFIED {
return rpcErrorToError(resp.Err)
}
return nil
}
func (db *DatabaseClient) Sync() error {
// RPC database sync is a no-op on client side
// The server handles persistence
return nil
}
func (db *DatabaseClient) Close() error {
resp, err := db.client.Close(context.Background(), &rpcdbpb.CloseRequest{})
if err != nil {
return err
}
if resp.Err != rpcdbpb.Error_ERROR_UNSPECIFIED {
return rpcErrorToError(resp.Err)
}
return nil
}
func (db *DatabaseClient) HealthCheck(ctx context.Context) (interface{}, error) {
resp, err := db.client.HealthCheck(ctx, &emptypb.Empty{})
if err != nil {
return nil, err
}
return resp.Details, nil
}
type batch struct {
db *DatabaseClient
puts []*rpcdbpb.PutRequest
deletes []*rpcdbpb.DeleteRequest
size int
}
func (b *batch) Put(key, value []byte) error {
b.puts = append(b.puts, &rpcdbpb.PutRequest{
Key: key,
Value: value,
})
b.size += len(key) + len(value)
return nil
}
func (b *batch) Delete(key []byte) error {
b.deletes = append(b.deletes, &rpcdbpb.DeleteRequest{
Key: key,
})
b.size += len(key)
return nil
}
func (b *batch) Size() int {
return b.size
}
func (b *batch) Write() error {
resp, err := b.db.client.WriteBatch(context.Background(), &rpcdbpb.WriteBatchRequest{
Puts: b.puts,
Deletes: b.deletes,
})
if err != nil {
return err
}
if resp.Err != rpcdbpb.Error_ERROR_UNSPECIFIED {
return rpcErrorToError(resp.Err)
}
return nil
}
func (b *batch) Reset() {
b.puts = b.puts[:0]
b.deletes = b.deletes[:0]
b.size = 0
}
func (b *batch) Replay(w database.KeyValueWriterDeleter) error {
for _, put := range b.puts {
if err := w.Put(put.Key, put.Value); err != nil {
return err
}
}
for _, del := range b.deletes {
if err := w.Delete(del.Key); err != nil {
return err
}
}
return nil
}
func (b *batch) Inner() database.Batch {
return b
}
type iterator struct {
db *DatabaseClient
id uint64
data []*rpcdbpb.PutRequest
dataIndex int
err error
closed bool
}
func (it *iterator) Next() bool {
if it.err != nil || it.closed {
return false
}
// If we have data left from previous fetch, use it
if it.dataIndex < len(it.data)-1 {
it.dataIndex++
return true
}
// Fetch next batch
resp, err := it.db.client.IteratorNext(context.Background(), &rpcdbpb.IteratorNextRequest{
Id: it.id,
})
if err != nil {
it.err = err
return false
}
if len(resp.Data) == 0 {
return false
}
it.data = resp.Data
it.dataIndex = 0
return true
}
func (it *iterator) Error() error {
if it.err != nil {
return it.err
}
resp, err := it.db.client.IteratorError(context.Background(), &rpcdbpb.IteratorErrorRequest{
Id: it.id,
})
if err != nil {
return err
}
if resp.Err != rpcdbpb.Error_ERROR_UNSPECIFIED {
return rpcErrorToError(resp.Err)
}
return nil
}
func (it *iterator) Key() []byte {
if it.dataIndex >= 0 && it.dataIndex < len(it.data) {
return it.data[it.dataIndex].Key
}
return nil
}
func (it *iterator) Value() []byte {
if it.dataIndex >= 0 && it.dataIndex < len(it.data) {
return it.data[it.dataIndex].Value
}
return nil
}
func (it *iterator) Release() {
if it.closed {
return
}
it.closed = true
it.db.client.IteratorRelease(context.Background(), &rpcdbpb.IteratorReleaseRequest{
Id: it.id,
})
}
// Helper functions for error conversion
func errorToRPCError(err error) rpcdbpb.Error {
if err == nil {
return rpcdbpb.Error_ERROR_UNSPECIFIED
}
switch err {
case database.ErrNotFound:
return rpcdbpb.Error_ERROR_NOT_FOUND
case database.ErrClosed:
return rpcdbpb.Error_ERROR_CLOSED
default:
// For any other error, we'll use CLOSED to indicate an error occurred
// since we only have limited error types in the proto
return rpcdbpb.Error_ERROR_CLOSED
}
}
func rpcErrorToError(err rpcdbpb.Error) error {
switch err {
case rpcdbpb.Error_ERROR_UNSPECIFIED:
return nil
case rpcdbpb.Error_ERROR_NOT_FOUND:
return database.ErrNotFound
case rpcdbpb.Error_ERROR_CLOSED:
return database.ErrClosed
default:
return database.ErrClosed
}
}
-653
View File
@@ -1,653 +0,0 @@
//go:build !grpc
// Copyright (C) 2019-2025, Lux Industries Inc. All rights reserved.
// See the file LICENSE for licensing terms.
// Package rpcdb provides the default ZAP-backed database RPC.
//
// The host (luxd) holds a `database.Database` (currently zapdb under
// the hood) and exposes it to a VM plugin process over the ZAP
// transport. The plugin sees a `database.Database` interface that
// forwards every call across the wire.
//
// This is the default build — it uses ZAP for transport and pure-Go
// types from proto/zap/rpcdb on the wire (no protobuf). The grpc
// build tag selects rpcdb_grpc.go instead, which uses gRPC + the
// protobuf-generated types under proto/pb/rpcdb. The two are mutually
// exclusive: build with `-tags=grpc` or with no tag, never both.
package rpcdb
import (
"context"
"errors"
"io"
"sync"
"github.com/luxfi/database"
zaprpcdb "github.com/luxfi/node/proto/zap/rpcdb"
)
// Re-export the wire types so callers that import "proto/rpcdb"
// see a single namespace regardless of which transport is active.
type (
Error = zaprpcdb.Error
HasRequest = zaprpcdb.HasRequest
HasResponse = zaprpcdb.HasResponse
GetRequest = zaprpcdb.GetRequest
GetResponse = zaprpcdb.GetResponse
PutRequest = zaprpcdb.PutRequest
PutResponse = zaprpcdb.PutResponse
DeleteRequest = zaprpcdb.DeleteRequest
DeleteResponse = zaprpcdb.DeleteResponse
WriteBatchRequest = zaprpcdb.WriteBatchRequest
WriteBatchResponse = zaprpcdb.WriteBatchResponse
CompactRequest = zaprpcdb.CompactRequest
CompactResponse = zaprpcdb.CompactResponse
CloseRequest = zaprpcdb.CloseRequest
CloseResponse = zaprpcdb.CloseResponse
HealthCheckResponse = zaprpcdb.HealthCheckResponse
NewIteratorWithStartAndPrefixRequest = zaprpcdb.NewIteratorWithStartAndPrefixRequest
NewIteratorWithStartAndPrefixResponse = zaprpcdb.NewIteratorWithStartAndPrefixResponse
IteratorNextRequest = zaprpcdb.IteratorNextRequest
IteratorNextResponse = zaprpcdb.IteratorNextResponse
IteratorErrorRequest = zaprpcdb.IteratorErrorRequest
IteratorErrorResponse = zaprpcdb.IteratorErrorResponse
IteratorReleaseRequest = zaprpcdb.IteratorReleaseRequest
IteratorReleaseResponse = zaprpcdb.IteratorReleaseResponse
)
const (
Error_ERROR_UNSPECIFIED = zaprpcdb.Error_ERROR_UNSPECIFIED
Error_ERROR_CLOSED = zaprpcdb.Error_ERROR_CLOSED
Error_ERROR_NOT_FOUND = zaprpcdb.Error_ERROR_NOT_FOUND
)
// Transport is the minimal RPC dispatcher abstraction the rpcdb
// client uses to talk to the server. The default ZAP transport
// implementation in vms/rpcchainvm/zap satisfies this interface; any
// other in-process or socket-based transport that can dispatch
// (method name, request) → response also works.
//
// Keeping this interface in the rpcdb package — rather than importing
// luxfi/api/zap directly — avoids a circular dependency between the
// proto layer and the application layer.
type Transport interface {
// Call dispatches a request to the named method and writes the
// response into out. The on-the-wire encoding of req and out is
// the transport's responsibility; this package only commits to
// the named-method shape.
Call(ctx context.Context, method string, req any, out any) error
}
// Method names — kept in one place so the client and server stay in
// sync.
const (
methodHas = "rpcdb.Has"
methodGet = "rpcdb.Get"
methodPut = "rpcdb.Put"
methodDelete = "rpcdb.Delete"
methodWriteBatch = "rpcdb.WriteBatch"
methodCompact = "rpcdb.Compact"
methodClose = "rpcdb.Close"
methodHealthCheck = "rpcdb.HealthCheck"
methodNewIteratorWithStartAndPrefix = "rpcdb.NewIteratorWithStartAndPrefix"
methodIteratorNext = "rpcdb.IteratorNext"
methodIteratorError = "rpcdb.IteratorError"
methodIteratorRelease = "rpcdb.IteratorRelease"
)
// DatabaseServer hosts a `database.Database` over a ZAP transport.
// The server registers its methods with the transport's handler
// surface so a remote DatabaseClient can call them. Every database
// call (Has/Get/Put/…) becomes one round-trip on the underlying
// connection, which today is a Z-Wing-encrypted ZAP channel.
type DatabaseServer struct {
db database.Database
iteratorLock sync.RWMutex
nextIteratorID uint64
iterators map[uint64]database.Iterator
}
// NewServer wraps `db` for service over a ZAP transport.
func NewServer(db database.Database) *DatabaseServer {
return &DatabaseServer{
db: db,
iterators: make(map[uint64]database.Iterator),
}
}
// Has services rpcdb.Has.
func (s *DatabaseServer) Has(_ context.Context, req *HasRequest) (*HasResponse, error) {
has, err := s.db.Has(req.Key)
if err != nil {
return &HasResponse{Has: false, Err: errorToRPCError(err)}, nil
}
return &HasResponse{Has: has, Err: Error_ERROR_UNSPECIFIED}, nil
}
// Get services rpcdb.Get.
func (s *DatabaseServer) Get(_ context.Context, req *GetRequest) (*GetResponse, error) {
value, err := s.db.Get(req.Key)
if err != nil {
return &GetResponse{Value: nil, Err: errorToRPCError(err)}, nil
}
return &GetResponse{Value: value, Err: Error_ERROR_UNSPECIFIED}, nil
}
// Put services rpcdb.Put.
func (s *DatabaseServer) Put(_ context.Context, req *PutRequest) (*PutResponse, error) {
return &PutResponse{Err: errorToRPCError(s.db.Put(req.Key, req.Value))}, nil
}
// Delete services rpcdb.Delete.
func (s *DatabaseServer) Delete(_ context.Context, req *DeleteRequest) (*DeleteResponse, error) {
return &DeleteResponse{Err: errorToRPCError(s.db.Delete(req.Key))}, nil
}
// WriteBatch services rpcdb.WriteBatch.
func (s *DatabaseServer) WriteBatch(_ context.Context, req *WriteBatchRequest) (*WriteBatchResponse, error) {
batch := s.db.NewBatch()
for _, p := range req.Puts {
if err := batch.Put(p.Key, p.Value); err != nil {
return &WriteBatchResponse{Err: errorToRPCError(err)}, nil
}
}
for _, d := range req.Deletes {
if err := batch.Delete(d.Key); err != nil {
return &WriteBatchResponse{Err: errorToRPCError(err)}, nil
}
}
return &WriteBatchResponse{Err: errorToRPCError(batch.Write())}, nil
}
// Compact services rpcdb.Compact.
func (s *DatabaseServer) Compact(_ context.Context, req *CompactRequest) (*CompactResponse, error) {
return &CompactResponse{Err: errorToRPCError(s.db.Compact(req.Start, req.Limit))}, nil
}
// Close services rpcdb.Close.
func (s *DatabaseServer) Close(_ context.Context, _ *CloseRequest) (*CloseResponse, error) {
return &CloseResponse{Err: errorToRPCError(s.db.Close())}, nil
}
// HealthCheck services rpcdb.HealthCheck.
func (s *DatabaseServer) HealthCheck(ctx context.Context) (*HealthCheckResponse, error) {
health, err := s.db.HealthCheck(ctx)
if err != nil {
return &HealthCheckResponse{Details: []byte(err.Error())}, nil
}
details := []byte("healthy")
if health != nil {
if str, ok := health.(string); ok {
details = []byte(str)
}
}
return &HealthCheckResponse{Details: details}, nil
}
// NewIteratorWithStartAndPrefix services rpcdb.NewIteratorWithStartAndPrefix.
func (s *DatabaseServer) NewIteratorWithStartAndPrefix(_ context.Context, req *NewIteratorWithStartAndPrefixRequest) (*NewIteratorWithStartAndPrefixResponse, error) {
it := s.db.NewIteratorWithStartAndPrefix(req.Start, req.Prefix)
s.iteratorLock.Lock()
id := s.nextIteratorID
s.nextIteratorID++
s.iterators[id] = it
s.iteratorLock.Unlock()
return &NewIteratorWithStartAndPrefixResponse{Id: id}, nil
}
// IteratorNext services rpcdb.IteratorNext.
func (s *DatabaseServer) IteratorNext(_ context.Context, req *IteratorNextRequest) (*IteratorNextResponse, error) {
s.iteratorLock.RLock()
it, ok := s.iterators[req.Id]
s.iteratorLock.RUnlock()
if !ok {
return &IteratorNextResponse{Data: nil}, nil
}
const maxBatchSize = 100
var data []*PutRequest
for i := 0; i < maxBatchSize && it.Next(); i++ {
data = append(data, &PutRequest{Key: it.Key(), Value: it.Value()})
}
return &IteratorNextResponse{Data: data}, nil
}
// IteratorError services rpcdb.IteratorError.
func (s *DatabaseServer) IteratorError(_ context.Context, req *IteratorErrorRequest) (*IteratorErrorResponse, error) {
s.iteratorLock.RLock()
it, ok := s.iterators[req.Id]
s.iteratorLock.RUnlock()
if !ok {
return &IteratorErrorResponse{Err: Error_ERROR_NOT_FOUND}, nil
}
return &IteratorErrorResponse{Err: errorToRPCError(it.Error())}, nil
}
// IteratorRelease services rpcdb.IteratorRelease.
func (s *DatabaseServer) IteratorRelease(_ context.Context, req *IteratorReleaseRequest) (*IteratorReleaseResponse, error) {
s.iteratorLock.Lock()
it, ok := s.iterators[req.Id]
if ok {
it.Release()
delete(s.iterators, req.Id)
}
s.iteratorLock.Unlock()
if !ok {
return &IteratorReleaseResponse{Err: Error_ERROR_NOT_FOUND}, nil
}
return &IteratorReleaseResponse{Err: Error_ERROR_UNSPECIFIED}, nil
}
// Register installs `s` on `r` so a remote DatabaseClient can dial in.
// `r` is a method-name → handler registry; the ZAP transport in
// vms/rpcchainvm/zap satisfies this.
type HandlerRegistry interface {
HandleFunc(method string, handler func(context.Context, any, any) error)
}
// Register binds every database method on `r`. The handler signatures
// reflect the `(ctx, request) -> response` shape of every rpcdb call;
// the transport is responsible for unmarshaling the request and
// marshaling the response.
func (s *DatabaseServer) Register(r HandlerRegistry) {
r.HandleFunc(methodHas, func(ctx context.Context, in any, out any) error {
req := in.(*HasRequest)
resp, err := s.Has(ctx, req)
if err != nil {
return err
}
*(out.(*HasResponse)) = *resp
return nil
})
r.HandleFunc(methodGet, func(ctx context.Context, in any, out any) error {
req := in.(*GetRequest)
resp, err := s.Get(ctx, req)
if err != nil {
return err
}
*(out.(*GetResponse)) = *resp
return nil
})
r.HandleFunc(methodPut, func(ctx context.Context, in any, out any) error {
req := in.(*PutRequest)
resp, err := s.Put(ctx, req)
if err != nil {
return err
}
*(out.(*PutResponse)) = *resp
return nil
})
r.HandleFunc(methodDelete, func(ctx context.Context, in any, out any) error {
req := in.(*DeleteRequest)
resp, err := s.Delete(ctx, req)
if err != nil {
return err
}
*(out.(*DeleteResponse)) = *resp
return nil
})
r.HandleFunc(methodWriteBatch, func(ctx context.Context, in any, out any) error {
req := in.(*WriteBatchRequest)
resp, err := s.WriteBatch(ctx, req)
if err != nil {
return err
}
*(out.(*WriteBatchResponse)) = *resp
return nil
})
r.HandleFunc(methodCompact, func(ctx context.Context, in any, out any) error {
req := in.(*CompactRequest)
resp, err := s.Compact(ctx, req)
if err != nil {
return err
}
*(out.(*CompactResponse)) = *resp
return nil
})
r.HandleFunc(methodClose, func(ctx context.Context, in any, out any) error {
req := in.(*CloseRequest)
resp, err := s.Close(ctx, req)
if err != nil {
return err
}
*(out.(*CloseResponse)) = *resp
return nil
})
r.HandleFunc(methodHealthCheck, func(ctx context.Context, _ any, out any) error {
resp, err := s.HealthCheck(ctx)
if err != nil {
return err
}
*(out.(*HealthCheckResponse)) = *resp
return nil
})
r.HandleFunc(methodNewIteratorWithStartAndPrefix, func(ctx context.Context, in any, out any) error {
req := in.(*NewIteratorWithStartAndPrefixRequest)
resp, err := s.NewIteratorWithStartAndPrefix(ctx, req)
if err != nil {
return err
}
*(out.(*NewIteratorWithStartAndPrefixResponse)) = *resp
return nil
})
r.HandleFunc(methodIteratorNext, func(ctx context.Context, in any, out any) error {
req := in.(*IteratorNextRequest)
resp, err := s.IteratorNext(ctx, req)
if err != nil {
return err
}
*(out.(*IteratorNextResponse)) = *resp
return nil
})
r.HandleFunc(methodIteratorError, func(ctx context.Context, in any, out any) error {
req := in.(*IteratorErrorRequest)
resp, err := s.IteratorError(ctx, req)
if err != nil {
return err
}
*(out.(*IteratorErrorResponse)) = *resp
return nil
})
r.HandleFunc(methodIteratorRelease, func(ctx context.Context, in any, out any) error {
req := in.(*IteratorReleaseRequest)
resp, err := s.IteratorRelease(ctx, req)
if err != nil {
return err
}
*(out.(*IteratorReleaseResponse)) = *resp
return nil
})
}
// DatabaseClient is the plugin-side database that forwards every call
// to a remote DatabaseServer over a ZAP Transport. Implements
// database.Database, so the plugin can drop it in unmodified.
type DatabaseClient struct {
t Transport
}
// NewClient builds a DatabaseClient that talks to a DatabaseServer
// across `t`.
func NewClient(t Transport) *DatabaseClient {
return &DatabaseClient{t: t}
}
// Compile-time interface check.
var _ database.Database = (*DatabaseClient)(nil)
// Has forwards rpcdb.Has.
func (c *DatabaseClient) Has(key []byte) (bool, error) {
resp := &HasResponse{}
if err := c.t.Call(context.Background(), methodHas, &HasRequest{Key: key}, resp); err != nil {
return false, err
}
if resp.Err != Error_ERROR_UNSPECIFIED {
return false, rpcErrorToError(resp.Err)
}
return resp.Has, nil
}
// Get forwards rpcdb.Get.
func (c *DatabaseClient) Get(key []byte) ([]byte, error) {
resp := &GetResponse{}
if err := c.t.Call(context.Background(), methodGet, &GetRequest{Key: key}, resp); err != nil {
return nil, err
}
if resp.Err != Error_ERROR_UNSPECIFIED {
return nil, rpcErrorToError(resp.Err)
}
return resp.Value, nil
}
// Put forwards rpcdb.Put.
func (c *DatabaseClient) Put(key, value []byte) error {
resp := &PutResponse{}
if err := c.t.Call(context.Background(), methodPut, &PutRequest{Key: key, Value: value}, resp); err != nil {
return err
}
if resp.Err != Error_ERROR_UNSPECIFIED {
return rpcErrorToError(resp.Err)
}
return nil
}
// Delete forwards rpcdb.Delete.
func (c *DatabaseClient) Delete(key []byte) error {
resp := &DeleteResponse{}
if err := c.t.Call(context.Background(), methodDelete, &DeleteRequest{Key: key}, resp); err != nil {
return err
}
if resp.Err != Error_ERROR_UNSPECIFIED {
return rpcErrorToError(resp.Err)
}
return nil
}
// NewBatch returns an in-memory batch that flushes via WriteBatch.
func (c *DatabaseClient) NewBatch() database.Batch {
return &batch{c: c}
}
// NewIterator returns a forward-only iterator over the entire DB.
func (c *DatabaseClient) NewIterator() database.Iterator {
return c.NewIteratorWithStartAndPrefix(nil, nil)
}
// NewIteratorWithStart returns a forward-only iterator from start.
func (c *DatabaseClient) NewIteratorWithStart(start []byte) database.Iterator {
return c.NewIteratorWithStartAndPrefix(start, nil)
}
// NewIteratorWithPrefix returns a forward-only iterator over keys
// matching prefix.
func (c *DatabaseClient) NewIteratorWithPrefix(prefix []byte) database.Iterator {
return c.NewIteratorWithStartAndPrefix(nil, prefix)
}
// NewIteratorWithStartAndPrefix forwards rpcdb.NewIteratorWithStartAndPrefix.
func (c *DatabaseClient) NewIteratorWithStartAndPrefix(start, prefix []byte) database.Iterator {
resp := &NewIteratorWithStartAndPrefixResponse{}
err := c.t.Call(context.Background(), methodNewIteratorWithStartAndPrefix,
&NewIteratorWithStartAndPrefixRequest{Start: start, Prefix: prefix}, resp)
if err != nil {
return &iterator{c: c, err: err}
}
return &iterator{c: c, id: resp.Id}
}
// Backup is not supported across an RPC boundary.
func (*DatabaseClient) Backup(io.Writer, uint64) (uint64, error) {
return 0, errors.New("rpcdb: backup not supported")
}
// Load is not supported across an RPC boundary.
func (*DatabaseClient) Load(io.Reader) error {
return errors.New("rpcdb: load not supported")
}
// Compact forwards rpcdb.Compact.
func (c *DatabaseClient) Compact(start, limit []byte) error {
resp := &CompactResponse{}
if err := c.t.Call(context.Background(), methodCompact, &CompactRequest{Start: start, Limit: limit}, resp); err != nil {
return err
}
if resp.Err != Error_ERROR_UNSPECIFIED {
return rpcErrorToError(resp.Err)
}
return nil
}
// Sync is a no-op on the client; the server owns persistence.
func (*DatabaseClient) Sync() error { return nil }
// Close forwards rpcdb.Close.
func (c *DatabaseClient) Close() error {
resp := &CloseResponse{}
if err := c.t.Call(context.Background(), methodClose, &CloseRequest{}, resp); err != nil {
return err
}
if resp.Err != Error_ERROR_UNSPECIFIED {
return rpcErrorToError(resp.Err)
}
return nil
}
// HealthCheck forwards rpcdb.HealthCheck.
func (c *DatabaseClient) HealthCheck(ctx context.Context) (interface{}, error) {
resp := &HealthCheckResponse{}
if err := c.t.Call(ctx, methodHealthCheck, struct{}{}, resp); err != nil {
return nil, err
}
return resp.Details, nil
}
// batch buffers Put/Delete locally and writes once on flush.
type batch struct {
c *DatabaseClient
puts []*PutRequest
deletes []*DeleteRequest
size int
}
func (b *batch) Put(key, value []byte) error {
b.puts = append(b.puts, &PutRequest{Key: key, Value: value})
b.size += len(key) + len(value)
return nil
}
func (b *batch) Delete(key []byte) error {
b.deletes = append(b.deletes, &DeleteRequest{Key: key})
b.size += len(key)
return nil
}
func (b *batch) Size() int { return b.size }
func (b *batch) Write() error {
resp := &WriteBatchResponse{}
if err := b.c.t.Call(context.Background(), methodWriteBatch, &WriteBatchRequest{Puts: b.puts, Deletes: b.deletes}, resp); err != nil {
return err
}
if resp.Err != Error_ERROR_UNSPECIFIED {
return rpcErrorToError(resp.Err)
}
return nil
}
func (b *batch) Reset() {
b.puts = b.puts[:0]
b.deletes = b.deletes[:0]
b.size = 0
}
func (b *batch) Replay(w database.KeyValueWriterDeleter) error {
for _, p := range b.puts {
if err := w.Put(p.Key, p.Value); err != nil {
return err
}
}
for _, d := range b.deletes {
if err := w.Delete(d.Key); err != nil {
return err
}
}
return nil
}
func (b *batch) Inner() database.Batch { return b }
// iterator paginates IteratorNext calls.
type iterator struct {
c *DatabaseClient
id uint64
data []*PutRequest
dataIndex int
err error
closed bool
}
func (it *iterator) Next() bool {
if it.err != nil || it.closed {
return false
}
if it.dataIndex < len(it.data)-1 {
it.dataIndex++
return true
}
resp := &IteratorNextResponse{}
if err := it.c.t.Call(context.Background(), methodIteratorNext, &IteratorNextRequest{Id: it.id}, resp); err != nil {
it.err = err
return false
}
if len(resp.Data) == 0 {
return false
}
it.data = resp.Data
it.dataIndex = 0
return true
}
func (it *iterator) Error() error {
if it.err != nil {
return it.err
}
resp := &IteratorErrorResponse{}
if err := it.c.t.Call(context.Background(), methodIteratorError, &IteratorErrorRequest{Id: it.id}, resp); err != nil {
return err
}
if resp.Err != Error_ERROR_UNSPECIFIED {
return rpcErrorToError(resp.Err)
}
return nil
}
func (it *iterator) Key() []byte {
if it.dataIndex >= 0 && it.dataIndex < len(it.data) {
return it.data[it.dataIndex].Key
}
return nil
}
func (it *iterator) Value() []byte {
if it.dataIndex >= 0 && it.dataIndex < len(it.data) {
return it.data[it.dataIndex].Value
}
return nil
}
func (it *iterator) Release() {
if it.closed {
return
}
it.closed = true
resp := &IteratorReleaseResponse{}
_ = it.c.t.Call(context.Background(), methodIteratorRelease, &IteratorReleaseRequest{Id: it.id}, resp)
}
func errorToRPCError(err error) Error {
if err == nil {
return Error_ERROR_UNSPECIFIED
}
switch err {
case database.ErrNotFound:
return Error_ERROR_NOT_FOUND
case database.ErrClosed:
return Error_ERROR_CLOSED
default:
return Error_ERROR_CLOSED
}
}
func rpcErrorToError(err Error) error {
switch err {
case Error_ERROR_UNSPECIFIED:
return nil
case Error_ERROR_NOT_FOUND:
return database.ErrNotFound
case Error_ERROR_CLOSED:
return database.ErrClosed
default:
return database.ErrClosed
}
}
-164
View File
@@ -1,164 +0,0 @@
// Copyright (C) 2019-2025, Lux Industries Inc. All rights reserved.
// See the file LICENSE for licensing terms.
// Package rpcdb provides the ZAP-native database RPC types — the
// shapes that flow over the wire when a VM plugin (separate process)
// reads/writes the host's database via the ZAP transport.
//
// This package is the single source of truth for the on-the-wire
// types; the grpc-tagged variant under proto/pb/rpcdb is generated
// from these shapes and is mutually exclusive with the ZAP default
// (you build with either `-tags=grpc` or no tag, never both).
//
// The actual storage backend on the server side is luxfi/database —
// which is currently zapdb under the hood. The same interface is
// presented to plugin clients regardless of which DB engine the host
// runs.
package rpcdb
// Error is the on-the-wire error code returned by every database RPC.
type Error int32
const (
// Error_ERROR_UNSPECIFIED is the success / no-error sentinel.
Error_ERROR_UNSPECIFIED Error = 0
// Error_ERROR_CLOSED maps to database.ErrClosed.
Error_ERROR_CLOSED Error = 1
// Error_ERROR_NOT_FOUND maps to database.ErrNotFound.
Error_ERROR_NOT_FOUND Error = 2
)
func (e Error) String() string {
switch e {
case Error_ERROR_CLOSED:
return "CLOSED"
case Error_ERROR_NOT_FOUND:
return "NOT_FOUND"
default:
return "UNSPECIFIED"
}
}
// HasRequest is the wire payload for Database.Has.
type HasRequest struct {
Key []byte
}
// HasResponse is the wire reply for Database.Has.
type HasResponse struct {
Has bool
Err Error
}
// GetRequest is the wire payload for Database.Get.
type GetRequest struct {
Key []byte
}
// GetResponse is the wire reply for Database.Get.
type GetResponse struct {
Value []byte
Err Error
}
// PutRequest is the wire payload for Database.Put (and the PUT entries
// in WriteBatch / IteratorNext).
type PutRequest struct {
Key []byte
Value []byte
}
// PutResponse is the wire reply for Database.Put.
type PutResponse struct {
Err Error
}
// DeleteRequest is the wire payload for Database.Delete.
type DeleteRequest struct {
Key []byte
}
// DeleteResponse is the wire reply for Database.Delete.
type DeleteResponse struct {
Err Error
}
// WriteBatchRequest is the wire payload for Database.WriteBatch.
type WriteBatchRequest struct {
Puts []*PutRequest
Deletes []*DeleteRequest
}
// WriteBatchResponse is the wire reply for Database.WriteBatch.
type WriteBatchResponse struct {
Err Error
}
// CompactRequest is the wire payload for Database.Compact.
type CompactRequest struct {
Start []byte
Limit []byte
}
// CompactResponse is the wire reply for Database.Compact.
type CompactResponse struct {
Err Error
}
// CloseRequest is the wire payload for Database.Close.
type CloseRequest struct{}
// CloseResponse is the wire reply for Database.Close.
type CloseResponse struct {
Err Error
}
// HealthCheckResponse is the wire reply for Database.HealthCheck.
type HealthCheckResponse struct {
Details []byte
}
// NewIteratorWithStartAndPrefixRequest is the wire payload for
// Database.NewIteratorWithStartAndPrefix.
type NewIteratorWithStartAndPrefixRequest struct {
Start []byte
Prefix []byte
}
// NewIteratorWithStartAndPrefixResponse is the wire reply for
// Database.NewIteratorWithStartAndPrefix.
type NewIteratorWithStartAndPrefixResponse struct {
Id uint64
}
// IteratorNextRequest is the wire payload for Database.IteratorNext.
type IteratorNextRequest struct {
Id uint64
}
// IteratorNextResponse is the wire reply for Database.IteratorNext.
// Data is the next batch of (key, value) pairs to read; an empty
// batch indicates the iterator has been exhausted.
type IteratorNextResponse struct {
Data []*PutRequest
}
// IteratorErrorRequest is the wire payload for Database.IteratorError.
type IteratorErrorRequest struct {
Id uint64
}
// IteratorErrorResponse is the wire reply for Database.IteratorError.
type IteratorErrorResponse struct {
Err Error
}
// IteratorReleaseRequest is the wire payload for Database.IteratorRelease.
type IteratorReleaseRequest struct {
Id uint64
}
// IteratorReleaseResponse is the wire reply for Database.IteratorRelease.
type IteratorReleaseResponse struct {
Err Error
}
+1 -1
View File
@@ -20,7 +20,7 @@ import (
zapwire "github.com/luxfi/api/zap"
"github.com/luxfi/database/memdb"
rpcdbzap "github.com/luxfi/vm/proto/rpcdb"
rpcdbzap "github.com/luxfi/node/db/rpcdb"
)
// TestCEvm_DialsZAPdbServer is a true cross-process end-to-end test:
+1 -1
View File
@@ -23,7 +23,7 @@ import (
"github.com/luxfi/log"
"github.com/luxfi/version"
"github.com/luxfi/vm/chain"
rpcdbzap "github.com/luxfi/vm/proto/rpcdb"
rpcdbzap "github.com/luxfi/node/db/rpcdb"
)
var (
+1 -1
View File
@@ -17,7 +17,7 @@ import (
"github.com/luxfi/database/memdb"
"github.com/luxfi/log"
luxruntime "github.com/luxfi/runtime"
rpcdbzap "github.com/luxfi/vm/proto/rpcdb"
rpcdbzap "github.com/luxfi/node/db/rpcdb"
)
// TestInitialize_DBServerAddrPopulated confirms that Client.Initialize