Files
Igor Sirotin c3768ff3fc refactor: address review on the Node facades
Move the Messaging API back out of pkg/kernel: ffi.Handle becomes a
defined type in the internal package, so kernel can hand the context to
pkg/messaging through kernel.Handle without the kernel layer knowing the
tier exists, and without the type being nameable outside this module.

Split Discovery into DiscV5, PeerExchange and DNSDiscovery, group the
node's identity and health under Debug(), and give every facade a pointer
receiver. Drop the node name, the per-operation logging that duplicates
what the library already writes, and GetFreePortIfNeeded — port 0 already
means "let the OS pick".

Heavy kernel tests now mark themselves with requiresNode and skip under
-short, so the gate runs `go test -short ./...` instead of naming tests
in a regexp.
2026-08-25 00:46:13 +01:00

183 lines
5.1 KiB
Go

package kernel
import (
"context"
"crypto/ecdsa"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"strconv"
"github.com/ethereum/go-ethereum/crypto"
"github.com/libp2p/go-libp2p/core/peer"
"github.com/logos-messaging/logos-delivery-go-bindings/internal/ffi"
"github.com/logos-messaging/logos-delivery-go-bindings/pkg/kernel/common"
"github.com/logos-messaging/logos-delivery-go-bindings/pkg/kernel/pb"
)
// Relay is a Node's relay protocol surface: gossipsub subscriptions, publishing
// and the state of the relay mesh. Take one with Node.Relay.
type Relay struct{ n *Node }
// Subscribe joins the relay mesh for a pubsub topic. Messages received on it
// arrive on Node.Messages.
func (r *Relay) Subscribe(pubsubTopic string) error {
if err := r.n.check(); err != nil {
return err
}
if pubsubTopic == "" {
return errors.New("kernel: relay subscribe: pubsub topic is empty")
}
if err := ffi.RelaySubscribe(r.n.h, pubsubTopic); err != nil {
return fmt.Errorf("kernel: relay subscribe %q: %w", pubsubTopic, err)
}
return nil
}
// Unsubscribe leaves the relay mesh for a pubsub topic.
func (r *Relay) Unsubscribe(pubsubTopic string) error {
if err := r.n.check(); err != nil {
return err
}
if pubsubTopic == "" {
return errors.New("kernel: relay unsubscribe: pubsub topic is empty")
}
if err := ffi.RelayUnsubscribe(r.n.h, pubsubTopic); err != nil {
return fmt.Errorf("kernel: relay unsubscribe %q: %w", pubsubTopic, err)
}
return nil
}
// Publish publishes a message on a pubsub topic and returns its hash. A ctx
// without a deadline gets the package default of 30s.
func (r *Relay) Publish(
ctx context.Context, pubsubTopic string, message *pb.WakuMessage,
) (common.MessageHash, error) {
if err := r.n.check(); err != nil {
return "", err
}
if err := ctx.Err(); err != nil {
return "", err
}
if message == nil {
return "", errors.New("kernel: relay publish: message is nil")
}
// The library expects the WakuMessage wire format (camelCase keys); the
// generated protobuf struct marshals content_topic, so marshal explicitly.
jsonMsg, err := json.Marshal(struct {
Payload []byte `json:"payload,omitempty"`
ContentTopic string `json:"contentTopic"`
Version *uint32 `json:"version,omitempty"`
Timestamp *int64 `json:"timestamp,omitempty"`
Meta []byte `json:"meta,omitempty"`
Ephemeral *bool `json:"ephemeral,omitempty"`
}{
Payload: message.Payload,
ContentTopic: message.ContentTopic,
Version: message.Version,
Timestamp: message.Timestamp,
Meta: message.Meta,
Ephemeral: message.Ephemeral,
})
if err != nil {
return "", err
}
hash, err := ffi.RelayPublish(r.n.h, pubsubTopic, string(jsonMsg), timeoutMillis(ctx, requestTimeout))
if err != nil {
return "", fmt.Errorf("kernel: relay publish: %w", err)
}
parsed, err := common.ToMessageHash(hash)
if err != nil {
return "", err
}
return parsed, nil
}
// AddProtectedShard registers the public key allowed to sign messages on a
// protected shard.
func (r *Relay) AddProtectedShard(clusterID, shardID uint16, pubkey *ecdsa.PublicKey) error {
if err := r.n.check(); err != nil {
return err
}
if pubkey == nil {
return errors.New("kernel: add protected shard: pubkey is nil")
}
keyHex := hex.EncodeToString(crypto.FromECDSAPub(pubkey))
if err := ffi.RelayAddProtectedShard(r.n.h, int(clusterID), int(shardID), keyHex); err != nil {
return fmt.Errorf("kernel: add protected shard: %w", err)
}
return nil
}
// PeersInMesh returns the relay mesh peers for a pubsub topic.
func (r *Relay) PeersInMesh(pubsubTopic string) (peer.IDSlice, error) {
if err := r.n.check(); err != nil {
return nil, err
}
list, err := ffi.GetPeersInMesh(r.n.h, pubsubTopic)
if err != nil {
return nil, fmt.Errorf("kernel: peers in mesh: %w", err)
}
return parsePeerIDs(list)
}
// NumPeersInMesh returns the relay mesh peer count for a pubsub topic.
func (r *Relay) NumPeersInMesh(pubsubTopic string) (int, error) {
if err := r.n.check(); err != nil {
return 0, err
}
countStr, err := ffi.GetNumPeersInMesh(r.n.h, pubsubTopic)
if err != nil {
return 0, fmt.Errorf("kernel: num peers in mesh: %w", err)
}
return strconv.Atoi(countStr)
}
// ConnectedPeers returns the connected relay peers, optionally narrowed to one
// pubsub topic.
func (r *Relay) ConnectedPeers(optPubsubTopic ...string) (peer.IDSlice, error) {
if err := r.n.check(); err != nil {
return nil, err
}
list, err := ffi.GetConnectedRelayPeers(r.n.h, optionalTopic(optPubsubTopic))
if err != nil {
return nil, fmt.Errorf("kernel: connected relay peers: %w", err)
}
return parsePeerIDs(list)
}
// NumConnectedPeers returns the connected relay peer count, optionally narrowed
// to one pubsub topic.
func (r *Relay) NumConnectedPeers(optPubsubTopic ...string) (int, error) {
if err := r.n.check(); err != nil {
return 0, err
}
countStr, err := ffi.GetNumConnectedRelayPeers(r.n.h, optionalTopic(optPubsubTopic))
if err != nil {
return 0, fmt.Errorf("kernel: num connected relay peers: %w", err)
}
return strconv.Atoi(countStr)
}
func optionalTopic(topics []string) string {
if len(topics) > 0 {
return topics[0]
}
return ""
}