Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 10 additions & 1 deletion config.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ type NodeAddress = conn.NodeAddress
// WithNodes returns a NodeSource containing the provided mapping from
// application-specific IDs to types implementing NodeAddress.
// Node IDs must be greater than 0.
func WithNodes[T NodeAddress](nodes map[uint32]T) NodeSource {
func WithNodes[T NodeAddress](nodes map[ID]T) NodeSource {
return conn.WithNodes(nodes)
}

Expand All @@ -25,6 +25,15 @@ func WithNodeList(addrsList []string) NodeSource {
return conn.WithNodeList(addrsList)
}

// ID identifies a node. The application chooses the IDs of configured nodes
// with [WithNodes], or [WithNodeList] assigns them in list order. Servers
// configured with [WithPeers] must agree on the ID of each peer, because a
// server announces its own ID when it connects. ID 0 is reserved:
// it marks a handler-only server and a client that announces no ID. A server
// assigns IDs from 2^20 upward to back-channel clients, skipping configured
// IDs. Under [WithStreamDedup], the peer with the lower ID dials the other.
type ID = conn.ID

// Node encapsulates the state of a node on which a remote procedure call can be
// performed. Nodes are created as part of a [Config] built with [NewConfig].
type Node = conn.Node
Expand Down
14 changes: 7 additions & 7 deletions doc/dev-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -186,7 +186,7 @@ classDiagram
}
class RequestHandler {
<<interface>>
HandleRequest(ctx, msg, release, send)
HandleRequest(ctx, senderID, msg, release, send)
}
class BidiStream {
<<interface>>
Expand All @@ -198,7 +198,7 @@ classDiagram
AcceptPeer(ctx, stream) InboundChannel
}
class Transport {
id uint32
id ID
shared bool
Enqueue(Request)
StoreChannel(Channel)
Expand Down Expand Up @@ -258,12 +258,12 @@ classDiagram
direction LR
class Config {
<<slice>>
NodeIDs() []uint32
NodeIDs() []ID
Size() int
Context(parent) ConfigContext
}
class Node {
id uint32
id ID
addr string
Context(parent) NodeContext
IsShared() bool
Expand Down Expand Up @@ -445,7 +445,7 @@ sequenceDiagram

The root `gorums.Server` ties the layers together on the server side.
It registers a `stream.Server` with its gRPC server, uses its `conn.InboundManager` as the server's `stream.PeerAcceptor`, and serves as the `stream.RequestHandler` for inbound streams, for its local node, and for back-channel requests on its outbound peer configuration.
`Server.HandleRequest` creates a `ServerContext` per request and runs the registered `Handler`, wrapped by any `ServerInterceptor`s.
`Server.HandleRequest` creates a `ServerContext` per request, records the sender ID that the channel passes in, and runs the registered `Handler`, wrapped by any `ServerInterceptor`s.
A `gorums.Message` wraps a `stream.Message` together with its decoded proto payload.

```mermaid
Expand All @@ -454,7 +454,7 @@ classDiagram
class Server {
handlers map[string]Handler
RegisterHandler(method, Handler)
HandleRequest(ctx, msg, release, send)
HandleRequest(ctx, senderID, msg, release, send)
ConnectedPeers() Config
ConnectedClients() Config
}
Expand Down Expand Up @@ -519,7 +519,7 @@ sequenceDiagram
SS->>IC: Serve(): receive loop
P->>IC: request Message
IC->>D: push handler job
D->>H: HandleRequest(ctx, msg, release, send)
D->>H: HandleRequest(ctx, senderID, msg, release, send)
H->>H: Handler(ServerContext, Message)
H->>IC: send(reply): push on sendQueue
IC->>P: Send(reply) on send loop
Expand Down
121 changes: 116 additions & 5 deletions doc/user-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -650,31 +650,31 @@ func majorityWrite(responses *gorums.Responses[*WriteResponse]) (*WriteResponse,
return nil, gorums.ErrIncomplete
}

// Process the map[uint32]*WriteResponse
// Process the map[gorums.ID]*WriteResponse
return aggregateWrites(replies), nil
}

// allSuccessful collects all successful responses; errored nodes are skipped.
func allSuccessful(responses *gorums.Responses[*Response]) map[uint32]*Response {
func allSuccessful(responses *gorums.Responses[*Response]) map[gorums.ID]*Response {
return responses.Results().IgnoreErrors().CollectAll()
}

// allResponses collects all responses; errored nodes are included with a zero value.
// Only use this when you know all nodes will succeed or you handle zero values.
func allResponses(responses *gorums.Responses[*Response]) map[uint32]*Response {
func allResponses(responses *gorums.Responses[*Response]) map[gorums.ID]*Response {
return responses.Results().CollectAll()
}
```

#### Per-Node Error Inspection

`CollectN` and `CollectAll` return `map[uint32]Resp` and do not preserve errors.
`CollectN` and `CollectAll` return `map[gorums.ID]Resp` and do not preserve errors.
When you need the actual error from a specific node, range over the results sequence directly:

```go
func inspectPerNodeErrors(responses *gorums.Responses[*Response]) {
var errs []error
replies := make(map[uint32]*Response)
replies := make(map[gorums.ID]*Response)

for result := range responses.Results() {
if result.Err != nil {
Expand Down Expand Up @@ -1543,6 +1543,117 @@ It returns no error, because closing a connection has no failure that a caller c

A `Server` configured with `WithPeers` owns its peer configuration, and `Stop` and `GracefulStop` close it.

## Node IDs

Gorums identifies each node with a `gorums.ID`.
The type is an alias for `uint32`, so a value from a protobuf `uint32` field needs no conversion.

### Rules for Node IDs

* ID 0 is reserved.
It marks a handler-only server and a client that announces no ID.
* With `WithNodes`, the application chooses the ID of each node.
With `WithNodeList`, Gorums assigns IDs in list order, starting after the highest ID in the configuration's connection pool, or from 1 for a new pool.
* Servers configured with `WithPeers` must agree on the ID of each peer, because a server announces its own ID when it connects.
* A server assigns IDs from 2^20 upward to back-channel clients, and it skips configured IDs.
* Under `WithStreamDedup`, the peer with the lower ID dials the other.

### Choosing Node IDs

Many protocols number their replicas from 1 to n, and use the number for leader rotation, slice indexes, or signature bitfields.
Use these numbers directly as Gorums IDs:

```go
type replica struct {
addr string
}

func (r replica) Addr() string { return r.addr }

nodes := map[gorums.ID]replica{
1: {addr: "10.0.0.1:9000"},
2: {addr: "10.0.0.2:9000"},
3: {addr: "10.0.0.3:9000"},
}
srv := gorums.NewServer(
gorums.WithPeers(myID, gorums.WithNodes(nodes), dialOpts),
)
```

Give each node an explicit ID in the cluster configuration, and let each node read its own ID from that configuration.
An explicit ID stays the same when the address list changes.
If you derive IDs from an address list instead, sort the list first, as the storage example does, so that all nodes compute the same IDs.

### Mapping Application IDs to Node IDs

An application may identify nodes by strings, such as host names or UUIDs.
Then give each node a Gorums ID in the shared cluster configuration, and keep a map in each direction:

```go
type member struct {
Name string // application ID
ID gorums.ID // Gorums ID
Address string
}

func (m member) Addr() string { return m.Address }

byID := make(map[gorums.ID]member)
idOf := make(map[string]gorums.ID)
for _, m := range members {
byID[m.ID] = m
idOf[m.Name] = m.ID
}
cfg, closeFn, err := gorums.NewConfig(gorums.WithNodes(byID), dialOpts)
```

`cfg.Node(idOf[name])` returns the node for an application ID, or nil if the configuration has no such node.
`byID[id].Name` returns the application ID for a Gorums ID.

Nodes that join without a shared configuration can derive their IDs from their names with a hash, so that all nodes compute the same ID without coordination:

```go
func idOf(name string) gorums.ID {
h := fnv.New32a()
h.Write([]byte(name))
return max(h.Sum32(), 1) // ID 0 is reserved
}
```

Two names can hash to the same 32-bit ID; the probability is about one in a million for 100 nodes.
Check for duplicate IDs when you build the node map.
Do not use hashed IDs in a protocol that needs IDs from 1 to n.

### Finding the Sender of a Request

A handler gets the ID of the node that sent a request from `ServerContext.SenderID`.
With `Config.Node`, it can then find the sender's node:

```go
func (srv *storageSrv) WriteQC(ctx gorums.ServerContext, req *WriteRequest) (*WriteResponse, error) {
if sender := ctx.PeerConfig().Node(ctx.SenderID()); sender != nil {
log.Printf("write from peer %d at %s", sender.ID(), sender.Address())
}
// ...
}
```

`SenderID` identifies the sender as this server knows it:

* a peer configured with `WithPeers` has its configured ID;
* a back-channel client has the ID that this server assigned to it, as in `ConnectedClients()`; the ID changes when the client reconnects;
* a request on a connection that this server dialed has the dialed node's ID;
* a request that this server sends to itself has this server's own ID.

`SenderID` returns 0 for a client that announces no ID, or an ID that this server does not know.

The sender is the last hop.
When a message travels through several nodes, for example down a tree, `SenderID` identifies the node that forwarded it, not the node where the message started.
Put the origin in the message if the handler needs it.

The sender asserts its ID when it connects, and the server does not authenticate it.
A Byzantine fault-tolerant protocol must still sign its messages and check the signer.

## Latency-Based Node Selection

Gorums tracks the round-trip latency to each node as an exponentially weighted
Expand Down
8 changes: 4 additions & 4 deletions examples/storage/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@ func runLocalCluster(srvOpts gorums.ServerOption) ([]string, func(), error) {
addrs := make([]string, len(servers))
for i, srv := range servers {
addrs[i] = srv.Addr()
registerAndServe(srv, uint32(i+1))
registerAndServe(srv, gorums.ID(i+1))
}

// Wait for all servers to see each other before opening the client REPL.
Expand All @@ -99,19 +99,19 @@ func runLocalCluster(srvOpts gorums.ServerOption) ([]string, func(), error) {
// assignment regardless of the order of addresses in the input list.
// The server's own address must be included in the peer list.
// It returns an error if the server's address is not found in the peer list.
func peerConfig(address string, peers []string) (uint32, gorums.NodeSource, error) {
func peerConfig(address string, peers []string) (gorums.ID, gorums.NodeSource, error) {
sorted := slices.Clone(peers)
slices.Sort(sorted)
idx := slices.Index(sorted, address)
if idx < 0 {
return 0, nil, fmt.Errorf("server address %q not found in -addrs list", address)
}
return uint32(idx + 1), gorums.WithNodeList(sorted), nil
return gorums.ID(idx + 1), gorums.WithNodeList(sorted), nil
}

// registerAndServe registers the storage service on srv and starts serving in
// a background goroutine. The server log output is labelled with the node ID.
func registerAndServe(srv *gorums.Server, id uint32) {
func registerAndServe(srv *gorums.Server, id gorums.ID) {
storage := newStorageServer(os.Stderr, fmt.Sprintf("node %d", id))
pb.RegisterStorageServer(srv, storage)
go func() {
Expand Down
2 changes: 1 addition & 1 deletion gorumstest/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ import (
// QuorumCallError returns a [gorums.QuorumCallError] with cause
// [gorums.ErrIncomplete] and one node error for each entry in nodeErrors.
// The node errors are in ascending node ID order.
func QuorumCallError(nodeErrors map[uint32]error) gorums.QuorumCallError {
func QuorumCallError(nodeErrors map[gorums.ID]error) gorums.QuorumCallError {
errs := make([]conn.NodeError, 0, len(nodeErrors))
for _, nodeID := range slices.Sorted(maps.Keys(nodeErrors)) {
errs = append(errs, conn.NewNodeError(nodeID, nodeErrors[nodeID]))
Expand Down
8 changes: 3 additions & 5 deletions gorumstest/gorumstest.go
Original file line number Diff line number Diff line change
Expand Up @@ -186,12 +186,10 @@ func Node(t testing.TB, srvFn func(i int) ServerIface, opts ...Option) *gorums.N
}

// PeerNode returns the node with the given id in cfg, or fails the test.
func PeerNode(t testing.TB, cfg gorums.Config, id uint32) *gorums.Node {
func PeerNode(t testing.TB, cfg gorums.Config, id gorums.ID) *gorums.Node {
t.Helper()
for _, node := range cfg {
if node.ID() == id {
return node
}
if node := cfg.Node(id); node != nil {
return node
}
t.Fatalf("node %d not in config %v", id, cfg.NodeIDs())
return nil
Expand Down
22 changes: 15 additions & 7 deletions internal/conn/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -115,8 +115,8 @@ func (c Config) Extend(nodes NodeSource) (Config, error) {
}

// NodeIDs returns a slice of this configuration's Node IDs.
func (c Config) NodeIDs() []uint32 {
ids := make([]uint32, len(c))
func (c Config) NodeIDs() []ID {
ids := make([]ID, len(c))
for i, node := range c {
ids[i] = node.ID()
}
Expand Down Expand Up @@ -146,14 +146,22 @@ func (c Config) Equal(b Config) bool {
return true
}

// Node returns the node in c with the given ID, or nil if c has no such node.
func (c Config) Node(id ID) *Node {
if i := slices.IndexFunc(c, func(n *Node) bool { return n.id == id }); i >= 0 {
return c[i]
}
return nil
}

// Contains reports whether c contains a node with the given ID.
func (c Config) Contains(id uint32) bool {
return slices.ContainsFunc(c, func(n *Node) bool { return n.id == id })
func (c Config) Contains(id ID) bool {
return c.Node(id) != nil
}

// Add returns a new Config containing nodes from c and nodes with the specified IDs.
// Duplicate IDs and IDs not found in the manager are ignored.
func (c Config) Add(ids ...uint32) Config {
func (c Config) Add(ids ...ID) Config {
if len(c) == 0 {
return nil
}
Expand Down Expand Up @@ -186,7 +194,7 @@ func (c Config) Union(other Config) Config {
}

// Remove returns a new Config excluding nodes with the specified IDs.
func (c Config) Remove(ids ...uint32) Config {
func (c Config) Remove(ids ...ID) Config {
if len(c) == 0 {
return nil
}
Expand Down Expand Up @@ -325,7 +333,7 @@ func (c Config) WithoutErrors(err QuorumCallError, errorTypes ...error) Config {
return false
}
// Build a set of node IDs to exclude.
excludeSet := newSet[uint32]()
excludeSet := newSet[ID]()
for _, ne := range err.errors {
if exclude(ne.cause) {
excludeSet.add(ne.nodeID)
Expand Down
26 changes: 26 additions & 0 deletions internal/conn/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,32 @@ func TestConfigWatch(t *testing.T) {
})
}

func TestConfigNode(t *testing.T) {
n1 := newTestNode(1, stream.NewChannelWithState(nil))
n3 := newTestNode(3, stream.NewChannelWithState(nil))
tests := []struct {
name string
cfg Config
id ID
want *Node
}{
{name: "First", cfg: Config{n1, n3}, id: 1, want: n1},
{name: "Last", cfg: Config{n1, n3}, id: 3, want: n3},
{name: "NotSortedByID", cfg: Config{n3, n1}, id: 1, want: n1},
{name: "Absent", cfg: Config{n1, n3}, id: 2, want: nil},
{name: "ReservedZero", cfg: Config{n1, n3}, id: 0, want: nil},
{name: "Empty", cfg: Config{}, id: 1, want: nil},
{name: "Nil", cfg: nil, id: 1, want: nil},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := tt.cfg.Node(tt.id); got != tt.want {
t.Errorf("Node(%d) = node %d, want node %d", tt.id, got.ID(), tt.want.ID())
}
})
}
}

func TestConfigSort(t *testing.T) {
const unmeasured = -1 * time.Second
// makeNode returns a node with the given latency and last error.
Expand Down
2 changes: 1 addition & 1 deletion internal/conn/dial_options.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ type DialOptions struct {
SendBufferSize uint
Metadata metadata.MD
Handler stream.RequestHandler
LocalNodeID uint32 // if non-zero, skip setting handler on this node ID
LocalNodeID ID // if non-zero, skip setting handler on this node ID
StreamDedup bool // reuse a lower-ID peer's dialed stream instead of dialing back
InboundManager *InboundManager // set when the configuration carries a server (peer or back-channel client); enables eager reconnect and, with StreamDedup, borrowing
Err error // records misuse of a dial option; surfaced by NewConfig
Expand Down
Loading
Loading