The go-nostr RelayConnect returns immediately after the WebSocket handshake, but the relay has not yet sent the AUTH challenge. Calling relay.Auth() right away signs the auth event with an empty challenge tag, which the Buzz relay rejects with 'auth-required: verification failed'. Retry the Auth() call up to 3 times with 1s backoff between attempts, giving the relay time to deliver the challenge.
249 lines
6.5 KiB
Go
249 lines
6.5 KiB
Go
// Package buzz implements a Channel for Buzz, a Nostr-based relay chat.
|
|
//
|
|
// Wire format: chat messages are kind:9 events scoped to a channel by an "h"
|
|
// tag. Mentions are "p" tags carrying the mentioned pubkey. The relay requires
|
|
// NIP-42 authentication (kind:22242) before it accepts a subscription.
|
|
package buzz
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/nbd-wtf/go-nostr"
|
|
"github.com/nbd-wtf/go-nostr/nip19"
|
|
|
|
"github.com/sipeed/picoclaw/pkg/bus"
|
|
"github.com/sipeed/picoclaw/pkg/channels"
|
|
"github.com/sipeed/picoclaw/pkg/config"
|
|
"github.com/sipeed/picoclaw/pkg/logger"
|
|
)
|
|
|
|
// kindStreamMessage is the Buzz chat message kind (NIP-29 style).
|
|
const kindStreamMessage = 9
|
|
|
|
// BuzzChannel implements the Channel interface for a Buzz relay.
|
|
type BuzzChannel struct {
|
|
*channels.BaseChannel
|
|
bc *config.Channel
|
|
config *config.BuzzSettings
|
|
|
|
secretKey string
|
|
publicKey string
|
|
|
|
relay *nostr.Relay
|
|
sub *nostr.Subscription
|
|
ctx context.Context
|
|
cancel context.CancelFunc
|
|
wg sync.WaitGroup
|
|
}
|
|
|
|
// NewBuzzChannel creates a new Buzz channel.
|
|
func NewBuzzChannel(
|
|
bc *config.Channel,
|
|
cfg *config.BuzzSettings,
|
|
messageBus *bus.MessageBus,
|
|
) (*BuzzChannel, error) {
|
|
if cfg.RelayURL == "" {
|
|
return nil, fmt.Errorf("buzz relay_url is required")
|
|
}
|
|
if len(cfg.Channels) == 0 {
|
|
return nil, fmt.Errorf("buzz channels is required: at least one channel ID to join")
|
|
}
|
|
|
|
sk, err := normalizeSecretKey(cfg.PrivateKey.String())
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
pk, err := nostr.GetPublicKey(sk)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("buzz private_key is not a valid secret key: %w", err)
|
|
}
|
|
|
|
base := channels.NewBaseChannel("buzz", cfg, messageBus, bc.AllowFrom,
|
|
channels.WithGroupTrigger(bc.GroupTrigger),
|
|
channels.WithReasoningChannelID(bc.ReasoningChannelID),
|
|
)
|
|
|
|
return &BuzzChannel{
|
|
BaseChannel: base,
|
|
bc: bc,
|
|
config: cfg,
|
|
secretKey: sk,
|
|
publicKey: pk,
|
|
}, nil
|
|
}
|
|
|
|
// normalizeSecretKey accepts either a 64-char hex secret key or an nsec1
|
|
// bech32 string and returns the hex form.
|
|
func normalizeSecretKey(key string) (string, error) {
|
|
key = strings.TrimSpace(key)
|
|
if key == "" {
|
|
return "", fmt.Errorf("buzz private_key is required")
|
|
}
|
|
|
|
if strings.HasPrefix(key, "nsec1") {
|
|
prefix, value, err := nip19.Decode(key)
|
|
if err != nil {
|
|
return "", fmt.Errorf("buzz private_key: invalid nsec: %w", err)
|
|
}
|
|
if prefix != "nsec" {
|
|
return "", fmt.Errorf("buzz private_key: expected nsec, got %s", prefix)
|
|
}
|
|
sk, ok := value.(string)
|
|
if !ok {
|
|
return "", fmt.Errorf("buzz private_key: unexpected nsec payload")
|
|
}
|
|
return sk, nil
|
|
}
|
|
|
|
if len(key) != 64 {
|
|
return "", fmt.Errorf("buzz private_key must be 64-char hex or an nsec1 string")
|
|
}
|
|
return key, nil
|
|
}
|
|
|
|
// Start connects to the relay, authenticates via NIP-42, and subscribes to the
|
|
// configured channels.
|
|
func (c *BuzzChannel) Start(ctx context.Context) error {
|
|
logger.InfoC("buzz", "Starting Buzz channel")
|
|
c.ctx, c.cancel = context.WithCancel(ctx)
|
|
|
|
relay, err := nostr.RelayConnect(c.ctx, c.config.RelayURL)
|
|
if err != nil {
|
|
c.cancel()
|
|
return fmt.Errorf("buzz relay connect failed: %w", err)
|
|
}
|
|
c.relay = relay
|
|
|
|
// NIP-42: the relay sends an AUTH challenge asynchronously after connect.
|
|
// The go-nostr library stores the challenge in an unexported field that is
|
|
// populated by the read loop. If we call Auth() too quickly the challenge
|
|
// is empty and the relay rejects the event with "verification failed".
|
|
// Retry with backoff to give the relay time to deliver the challenge.
|
|
const maxAuthAttempts = 3
|
|
var authErr error
|
|
for attempt := 0; attempt < maxAuthAttempts; attempt++ {
|
|
if attempt > 0 {
|
|
logger.WarnCF("buzz", "Retrying NIP-42 auth", map[string]any{
|
|
"attempt": attempt + 1,
|
|
"delay": "1s",
|
|
})
|
|
time.Sleep(1 * time.Second)
|
|
}
|
|
authErr = relay.Auth(c.ctx, func(evt *nostr.Event) error {
|
|
return evt.Sign(c.secretKey)
|
|
})
|
|
if authErr == nil {
|
|
break
|
|
}
|
|
logger.WarnCF("buzz", "NIP-42 auth attempt failed", map[string]any{
|
|
"attempt": attempt + 1,
|
|
"error": authErr.Error(),
|
|
})
|
|
}
|
|
if authErr != nil {
|
|
_ = relay.Close()
|
|
c.cancel()
|
|
return fmt.Errorf("buzz NIP-42 auth failed: %w", authErr)
|
|
}
|
|
|
|
channelIDs := []string(c.config.Channels)
|
|
filters := nostr.Filters{{
|
|
Kinds: []int{kindStreamMessage},
|
|
Tags: nostr.TagMap{"h": channelIDs},
|
|
}}
|
|
|
|
sub, err := relay.Subscribe(c.ctx, filters)
|
|
if err != nil {
|
|
_ = relay.Close()
|
|
c.cancel()
|
|
return fmt.Errorf("buzz subscribe failed: %w", err)
|
|
}
|
|
c.sub = sub
|
|
|
|
c.wg.Add(1)
|
|
go func() {
|
|
defer c.wg.Done()
|
|
c.consume(sub)
|
|
}()
|
|
|
|
c.SetRunning(true)
|
|
logger.InfoCF("buzz", "Buzz channel started", map[string]any{
|
|
"relay": c.config.RelayURL,
|
|
"pubkey": c.publicKey,
|
|
"channels": len(channelIDs),
|
|
})
|
|
return nil
|
|
}
|
|
|
|
// Stop closes the subscription and disconnects from the relay.
|
|
func (c *BuzzChannel) Stop(ctx context.Context) error {
|
|
logger.InfoC("buzz", "Stopping Buzz channel")
|
|
c.SetRunning(false)
|
|
|
|
if c.cancel != nil {
|
|
c.cancel()
|
|
}
|
|
if c.sub != nil {
|
|
c.sub.Unsub()
|
|
}
|
|
if c.relay != nil {
|
|
if err := c.relay.Close(); err != nil {
|
|
logger.WarnCF("buzz", "Relay close failed", map[string]any{"error": err.Error()})
|
|
}
|
|
}
|
|
c.wg.Wait()
|
|
|
|
logger.InfoC("buzz", "Buzz channel stopped")
|
|
return nil
|
|
}
|
|
|
|
// Send publishes a kind:9 message scoped to the target channel.
|
|
func (c *BuzzChannel) Send(ctx context.Context, msg bus.OutboundMessage) ([]string, error) {
|
|
if !c.IsRunning() {
|
|
return nil, channels.ErrNotRunning
|
|
}
|
|
|
|
target := msg.ChatID
|
|
if target == "" {
|
|
return nil, fmt.Errorf("chat ID is empty: %w", channels.ErrSendFailed)
|
|
}
|
|
if strings.TrimSpace(msg.Content) == "" {
|
|
return nil, nil
|
|
}
|
|
|
|
tags := nostr.Tags{nostr.Tag{"h", target}}
|
|
if c.config.ReplyInThread && msg.ReplyToMessageID != "" {
|
|
tags = append(tags, nostr.Tag{"e", msg.ReplyToMessageID, "", "reply"})
|
|
}
|
|
|
|
evt := nostr.Event{
|
|
PubKey: c.publicKey,
|
|
CreatedAt: nostr.Now(),
|
|
Kind: kindStreamMessage,
|
|
Tags: tags,
|
|
Content: msg.Content,
|
|
}
|
|
if err := evt.Sign(c.secretKey); err != nil {
|
|
return nil, fmt.Errorf("buzz sign failed: %w", errJoin(err, channels.ErrSendFailed))
|
|
}
|
|
|
|
if err := c.relay.Publish(ctx, evt); err != nil {
|
|
return nil, fmt.Errorf("buzz publish failed: %w", errJoin(err, channels.ErrSendFailed))
|
|
}
|
|
|
|
logger.DebugCF("buzz", "Message sent", map[string]any{
|
|
"channel": target,
|
|
"event_id": evt.ID,
|
|
})
|
|
return []string{evt.ID}, nil
|
|
}
|
|
|
|
// errJoin wraps err so that errors.Is(result, sentinel) holds for the sentinel
|
|
// while preserving the underlying cause in the message.
|
|
func errJoin(err, sentinel error) error {
|
|
return fmt.Errorf("%v: %w", err, sentinel)
|
|
}
|