gRPC
Every node can serve a gRPC API to manage itself, the hub it is part of and the pipelines of the nodes on that hub. This guide walks through the most common tasks in Go, with the Go client library. Every method, with its request and response messages, is described in the API Reference.
The code in this guide is the example program in
examples of the client library, and it
is run against real nodes by the test suite, so it works as shown.
Start the nodes
The guide uses two nodes: a client node that joins a hub, and the node that hosts the hub server. Both serve the gRPC API and run pipelines.
Start the client node with the --grpc flag:
np up --grpc 127.0.0.1:10432
Start the hub server node with a config file that enables the gRPC API with the
grpcBridge section:
node:
name: My Server
http:
address: 127.0.0.1:8769
hubServer:
authenticationByRequest:
timeout: 5m
grpcBridge:
address: 127.0.0.1:10564
pipelines: {}
np -c config.yaml up
Connect
The Go client library is github.com/nanoping-labs/nanoping-api-go/v11, with the version of the NanoPing release it was made from:
go get github.com/nanoping-labs/nanoping-api-go/v11
client.Connect returns a client with a field for every service, such as c.HubClient or c.Pipelines.
Each field is the generated gRPC client of that service, so every method in the reference is available on it.
c, err := client.Connect("127.0.0.1:10432")
if err != nil {
return err
}
defer c.Close()
The connection is made on the first call, so a node that cannot be reached shows up as an Unavailable
error there.
Name the node
A node is named after the machine it runs on. Give it a name, and metadata that other nodes on the hub can
filter on. A node started with a config file keeps the name and metadata of that file, and the call fails
with PermissionDenied.
// SetNodeInformation names the client node and gives it metadata that other
// nodes on the hub can filter on.
func SetNodeInformation(ctx context.Context, c *client.Client, out io.Writer) error {
local, err := c.HubClient.GetLocalNodeId(ctx, &hub_client.GetLocalNodeIdRequest{})
if err != nil {
return err
}
response, err := c.HubClient.SetNode(ctx, &hub_client.SetNodeRequest{
NodeId: local.NodeId,
Name: "My NanoPing Client",
Metadata: &hub_node.NodeMetadata{
Items: map[string]*hub_node.NodeMetadataItem{
"type": {Value: &hub_node.NodeMetadataItem_StringValue{StringValue: "my-type"}},
"location": {Value: &hub_node.NodeMetadataItem_StringValue{StringValue: "Denmark, Aalborg"}},
},
},
})
if err != nil {
return fmt.Errorf("set node information: %w", err)
}
fmt.Fprintf(out, "Named node %s %q\n", response.Node.Id, response.Node.Name)
return nil
}
Watch the connection to the hub
The connection state is sent when the stream opens and every time it changes. client.Receive turns the
stream into a loop that ends when the stream does.
// WatchConnectionState prints the connection state of the client node every
// time it changes, until ctx is canceled.
func WatchConnectionState(ctx context.Context, c *client.Client, out io.Writer) error {
stream, err := c.HubClient.StreamConnectionState(ctx, &hub_client.StreamConnectionStateRequest{})
if err != nil {
return err
}
go func() {
for state, err := range client.Receive(stream) {
if err != nil {
return
}
fmt.Fprintf(out, "Connection state: %s\n", state.State)
}
}()
return nil
}
Join the hub
Joining a hub takes two sides. The client node asks to join, and waits while the request is shown on the hub server:
// JoinHub joins the client node to the hub server. The call waits until the
// hub server accepted or rejected the request. A node that already is part of
// the hub stays so.
func JoinHub(ctx context.Context, c *client.Client, serverHttpAddress string, out io.Writer) error {
err := c.JoinHub(ctx, serverHttpAddress)
if status.Code(err) == codes.AlreadyExists {
fmt.Fprintln(out, "Already part of the hub")
return nil
}
if errors.Is(err, client.ErrJoinRejected) {
return errors.New("the hub server rejected the node")
}
if err != nil {
return err
}
fmt.Fprintln(out, "Joined the hub")
return nil
}
On the hub server node, the request is approved or rejected in the dashboard, with
np authentication-requests, or through the API:
// AcceptJoinRequests accepts every node that asks to join the hub server,
// until ctx is canceled.
func AcceptJoinRequests(ctx context.Context, server *client.Client) error {
return server.AnswerJoinRequests(ctx, func(request *hub_server.AuthenticationByRequestRequest) bool {
return true
})
}
Once joined, the node stays part of the hub after a restart, and asking again fails with AlreadyExists.
Find nodes on the hub
The hub server lists the nodes on the hub, optionally filtered by their metadata or by the services they run:
// PrintNodes prints the nodes on the hub: all of them, the ones with metadata
// "type" set to "my-type", and the ones running pipelines.
func PrintNodes(ctx context.Context, c *client.Client, out io.Writer) error {
all, err := c.HubServer.GetNodes(ctx, &hub_server.GetNodesRequest{})
if err != nil {
return err
}
fmt.Fprintln(out, "All nodes:")
for _, node := range all.Nodes {
fmt.Fprintf(out, " %s\n", node.Name)
}
myType, err := c.HubServer.GetNodes(ctx, &hub_server.GetNodesRequest{
MetadataFilters: map[string]*hub_node.NodeMetadataItem{
"type": {Value: &hub_node.NodeMetadataItem_StringValue{StringValue: "my-type"}},
},
})
if err != nil {
return err
}
fmt.Fprintln(out, "Nodes with type my-type:")
for _, node := range myType.Nodes {
fmt.Fprintf(out, " %s\n", node.Name)
}
withPipelines, err := PipelineNodes(ctx, c)
if err != nil {
return err
}
fmt.Fprintln(out, "Nodes running pipelines:")
for _, node := range withPipelines {
fmt.Fprintf(out, " %s\n", node.Name)
}
return nil
}
// PipelineNodes returns the nodes on the hub that run pipelines.
func PipelineNodes(ctx context.Context, c *client.Client) ([]*hub_node.Node, error) {
response, err := c.HubServer.GetNodes(ctx, &hub_server.GetNodesRequest{
ServiceFilters: []hub_node.NodeService{hub_node.NodeService_PIPELINES},
})
if err != nil {
return nil, err
}
return response.Nodes, nil
}
Manage pipelines
Every pipeline call names the node it acts on, so one connection manages the pipelines of every node on the hub. Watch the node's pipelines first: the call returns once the stream is ready, so the changes made after it are all reported.
// WatchPipelines prints an event every time a pipeline on the node is
// created, updated or deleted, until ctx is canceled. It returns once the
// stream is ready, so no change made after it returns is missed.
func WatchPipelines(ctx context.Context, c *client.Client, node *hub_node.Node, out io.Writer) error {
stream, err := c.Pipelines.StreamPipelines(ctx, &pipelines.StreamPipelinesRequest{NodeId: node.Id})
if err != nil {
return err
}
if _, err := stream.Header(); err != nil {
return err
}
go func() {
for event, err := range client.Receive(stream) {
if err != nil {
return
}
switch e := event.Event.(type) {
case *pipelines.StreamPipelinesEvent_Create:
fmt.Fprintf(out, " Event on %s: created %q\n", node.Name, e.Create.Pipeline.Name)
case *pipelines.StreamPipelinesEvent_Update:
fmt.Fprintf(out, " Event on %s: updated %q\n", node.Name, e.Update.Pipeline.Name)
case *pipelines.StreamPipelinesEvent_Delete:
fmt.Fprintf(out, " Event on %s: deleted %q\n", node.Name, e.Delete.Pipeline.Name)
}
}
}()
return nil
}
Then create a pipeline, rename it, read it back and delete it. An update replaces every setting, so it sends the current ones along with the change.
// ManagePipeline creates a pipeline on the node, renames it, reads it back and
// deletes it again.
func ManagePipeline(ctx context.Context, c *client.Client, nodeId string, out io.Writer) error {
created, err := c.Pipelines.CreatePipeline(ctx, &pipelines.CreatePipelineRequest{
NodeId: nodeId,
Name: "My First Pipeline",
JsonConfig: PipelineConfig,
RestartPolicy: &pipelines.RestartPolicy{
RestartPolicy: &pipelines.RestartPolicy_OnFailure{
OnFailure: &pipelines.RestartPolicyOnFailure{MaxRestarts: 3},
},
},
OptionalDefaultLoggingLevel: &pipelines.CreatePipelineRequest_DefaultLoggingLevel{
DefaultLoggingLevel: logging.Level_DEBUG,
},
})
if err != nil {
return fmt.Errorf("create pipeline: %w", err)
}
pipeline := created.Pipeline
fmt.Fprintf(out, " Created %q\n", pipeline.Name)
// An update replaces every setting, so send the current ones along with
// the new name.
updated, err := c.Pipelines.UpdatePipeline(ctx, &pipelines.UpdatePipelineRequest{
NodeId: nodeId,
Id: pipeline.Id,
Name: "My Renamed Pipeline",
JsonConfig: pipeline.JsonConfig,
RestartPolicy: pipeline.RestartPolicy,
StartupPolicy: pipeline.StartupPolicy,
DefaultLoggingLevel: pipeline.DefaultLoggingLevel,
InstructionsTimeout: pipeline.InstructionsTimeout,
})
if err != nil {
return fmt.Errorf("update pipeline: %w", err)
}
fmt.Fprintf(out, " Renamed to %q\n", updated.Pipeline.Name)
fetched, err := c.Pipelines.GetPipelineById(ctx, &pipelines.GetPipelineByIdRequest{NodeId: nodeId, Id: pipeline.Id})
if err != nil {
return fmt.Errorf("get pipeline: %w", err)
}
fmt.Fprintf(out, " Read back %q\n", fetched.Pipeline.Name)
_, err = c.Pipelines.DeletePipeline(ctx, &pipelines.DeletePipelineRequest{NodeId: nodeId, Id: pipeline.Id})
if err != nil {
return fmt.Errorf("delete pipeline: %w", err)
}
fmt.Fprintf(out, " Deleted %q\n", fetched.Pipeline.Name)
return nil
}
The pipeline configuration is a JSON string:
// PipelineConfig is a pipeline that sends traffic from a traffic source to a
// traffic sink on the same node.
const PipelineConfig = `{
"version": 13,
"plumr_config": {
"pipeline": {
"traffic_sink-1": {
"traffic_sink": {
"input": "[traffic_sink-1|in:0]-[uniform_traffic_source-1|out:0]",
"mtu": 1500
}
},
"uniform_traffic_source-1": {
"uniform_traffic_source": {
"output": "[traffic_sink-1|in:0]-[uniform_traffic_source-1|out:0]",
"total_packets": 5000,
"interval": 25
}
}
}
}
}`
Putting it all together
Run connects to both nodes and runs every step:
// Run connects to both nodes and runs every step of the guide.
func Run(ctx context.Context, addresses Addresses, out io.Writer) error {
c, err := client.Connect(addresses.Client)
if err != nil {
return err
}
defer c.Close()
server, err := client.Connect(addresses.Server)
if err != nil {
return err
}
defer server.Close()
ctx, cancel := context.WithCancel(ctx)
defer cancel()
out = &lockedWriter{out: out}
if err := WatchConnectionState(ctx, c, out); err != nil {
return err
}
if err := SetNodeInformation(ctx, c, out); err != nil {
return err
}
go func() { _ = AcceptJoinRequests(ctx, server) }()
if err := JoinHub(ctx, c, addresses.ServerHttp, out); err != nil {
return err
}
if err := PrintNodes(ctx, c, out); err != nil {
return err
}
nodes, err := PipelineNodes(ctx, c)
if err != nil {
return err
}
for _, node := range nodes {
fmt.Fprintf(out, "Managing a pipeline on %s\n", node.Name)
if err := WatchPipelines(ctx, c, node, out); err != nil {
return err
}
if err := ManagePipeline(ctx, c, node.Id, out); err != nil {
return err
}
}
return nil
}
main.go runs it with the addresses of the two nodes:
// Runs the gRPC guide of the NanoPing docs against two nodes.
//
// Start the client node with
//
// np up --grpc 127.0.0.1:10432
//
// and the hub server node with `np -c config.yaml up`, where config.yaml is
//
// node:
// name: My Server
// http:
// address: 127.0.0.1:8769
// hubServer:
// authenticationByRequest:
// timeout: 5m
// grpcBridge:
// address: 127.0.0.1:10564
// pipelines: {}
package main
import (
"context"
"flag"
"fmt"
"os"
"os/signal"
"nanoping/examples/guide"
)
func main() {
var addresses guide.Addresses
flag.StringVar(&addresses.Client, "client", "127.0.0.1:10432", "gRPC API of the client node")
flag.StringVar(&addresses.Server, "server", "127.0.0.1:10564", "gRPC API of the hub server node")
flag.StringVar(&addresses.ServerHttp, "server-http", "http://127.0.0.1:8769", "HTTP address of the hub server node")
flag.Parse()
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt)
defer stop()
if err := guide.Run(ctx, addresses, os.Stdout); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
}
Running it prints something like:
Connection state: UNAUTHENTICATED
Named node 1c127c74-feb5-4fef-aa1b-d5e69b8fce3c "My NanoPing Client"
Connection state: CONNECTED
Joined the hub
All nodes:
My Server
My NanoPing Client
Nodes with type my-type:
My NanoPing Client
Nodes running pipelines:
My Server
My NanoPing Client
Managing a pipeline on My Server
Created "My First Pipeline"
Event on My Server: created "My First Pipeline"
Event on My Server: updated "My Renamed Pipeline"
Renamed to "My Renamed Pipeline"
Read back "My Renamed Pipeline"
Deleted "My Renamed Pipeline"
Managing a pipeline on My NanoPing Client
...
Other languages
The gRPC API is described by the protos in grpc_bridge/protobuf. Generate a client from them for any
language gRPC supports. Without gRPC, use the HTTP REST API instead.