Developer SDK

Go examples

Install the SDK · View source on GitHub ↗

messages/main.go

package main

import (
	"context"
	"crypto/rand"
	"encoding/binary"
	"encoding/json"
	"fmt"
	integration "github.com/urnetwork/examples/go/integration"
	sdk "github.com/urnetwork/sdk/v2026"
	"os"
	"os/signal"
	"strings"
	"time"
)

type event struct {
	kind      string
	source    string
	bytes     []byte
	peers     *sdk.NetworkPeers
	supported bool
}
type listener struct{ events chan event }

func (l *listener) put(e event) {
	select {
	case l.events <- e:
	default:
		fmt.Fprintln(os.Stderr, "event queue full; dropped")
	}
}
func (l *listener) SubprotocolMessage(protocol int32, source *sdk.Id, bytes []byte) {
	if protocol == subprotocol && source != nil && len(bytes) <= 4112 {
		l.put(event{kind: "message", source: source.String(), bytes: append([]byte(nil), bytes...)})
	}
}
func (l *listener) NetworkPeersChanged(peers *sdk.NetworkPeers) {
	l.put(event{kind: "peers", peers: peers})
}
func (l *listener) Result(ids *sdk.IntList, ok bool) {
	supported := false
	if ok && ids != nil {
		for i := 0; i < ids.Len(); i++ {
			if ids.Get(i) == subprotocol {
				supported = true
			}
		}
	}
	l.put(event{kind: "query", supported: supported})
}
func showPeers(peers *sdk.NetworkPeers) {
	if peers == nil {
		fmt.Println("peers unavailable (no snapshot)")
		return
	}
	fmt.Println("disconnected:", peers.DisconnectedCount)
	if peers.Connected == nil {
		fmt.Println("connected peers unavailable")
	}
	if peers.Connected != nil {
		for i := 0; i < peers.Connected.Len(); i++ {
			p := peers.Connected.Get(i)
			if p == nil {
				continue
			}
			b, _ := json.Marshal(p)
			fmt.Printf("%s color=%s\n", b, p.ColorHex())
		}
	}
}
func run() error {
	args := os.Args[1:]
	if len(args) == 1 && args[0] == "--self-test" {
		if e := codecSelfTest(); e != nil {
			return e
		}
		fmt.Println("URMS codec self-test passed")
		return nil
	}
	if len(args) == 1 && args[0] == "--version" {
		fmt.Println(sdk.Version)
		return nil
	}
	if len(args) == 0 || (args[0] != "self" && args[0] != "peers" && args[0] != "watch" && args[0] != "send") || (args[0] == "send" && len(args) < 3) {
		return fmt.Errorf("usage: --self-test | --version | self | peers | watch | send CLIENT_ID TEXT")
	}
	var destination *sdk.Id
	var err error
	if args[0] == "send" {
		destination, err = sdk.ParseId(args[1])
		if err != nil {
			return err
		}
	}
	session, err := integration.Open(false)
	if err != nil {
		return err
	}
	defer session.Close()
	device := session.Device
	l := &listener{make(chan event, 256)}
	sub, err := device.EnableSubprotocol(subprotocol, l)
	if err != nil {
		return err
	}
	defer sub.Close()
	peersSub := device.AddNetworkPeersChangeListener(l)
	defer peersSub.Close()
	device.SetProvideMode(sdk.ProvideModeNetwork)
	fmt.Println("self:", device.GetClientId())
	if args[0] == "self" {
		return nil
	}
	snapshot := device.GetNetworkPeers()
	showPeers(snapshot)
	if args[0] == "peers" && snapshot != nil {
		return nil
	}
	ctx, cancel := signal.NotifyContext(context.Background(), os.Interrupt)
	defer cancel()
	tick := time.NewTicker(100 * time.Millisecond)
	defer tick.Stop()
	deadline := time.Now().Add(30 * time.Second)
	queried := false
	var random [8]byte
	if _, err = rand.Read(random[:]); err != nil {
		return err
	}
	pending := binary.BigEndian.Uint64(random[:])
	if pending == 0 {
		pending = 1
	}
	for {
		select {
		case <-ctx.Done():
			return nil
		case <-tick.C:
			if destination != nil && !queried && device.GetProviderConnected() {
				queried = true
				device.QuerySubprotocols(destination, 10000, l)
			}
			if args[0] != "watch" && time.Now().After(deadline) {
				return fmt.Errorf("connection, peer snapshot, query, or ACK timed out")
			}
		case e := <-l.events:
			switch e.kind {
			case "peers":
				showPeers(e.peers)
				if args[0] == "peers" && e.peers != nil {
					return nil
				}
			case "query":
				if !e.supported {
					return fmt.Errorf("peer query failed or peer does not advertise 4096")
				}
				b, err := encode(frame{1, pending, strings.Join(args[2:], " ")})
				if err != nil {
					return err
				}
				if !device.SendSubprotocolBytes(subprotocol, destination, b) {
					return fmt.Errorf("SDK did not enqueue message")
				}
				fmt.Println("sent:", pending, "waiting for ACK")
				deadline = time.Now().Add(10 * time.Second)
			case "message":
				f, err := decode(e.bytes)
				if err != nil {
					fmt.Fprintln(os.Stderr, "rejected frame:", err)
					continue
				}
				fmt.Printf("kind=%d source=%s id=%d text=%q\n", f.kind, e.source, f.id, f.text)
				if f.kind == 1 {
					source, err := sdk.ParseId(e.source)
					if err != nil {
						continue
					}
					ack, _ := encode(frame{2, f.id, ""})
					if !device.SendSubprotocolBytes(subprotocol, source, ack) {
						fmt.Fprintln(os.Stderr, "ACK was not enqueued")
					}
				} else if destination != nil && e.source == destination.String() && f.id == pending {
					return nil
				}
			}
		}
	}
}
func main() {
	if err := run(); err != nil {
		fmt.Fprintln(os.Stderr, err)
		os.Exit(1)
	}
}