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 permanently ('verification failed'
then 'authentication already failed' on retries).
Add a 2s sleep between RelayConnect and Auth() to give the background
read loop time to receive and store the challenge. The relay typically
delivers it within a few hundred milliseconds.
237 lines
6.4 KiB
Go
237 lines
6.4 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 the
|
|
// WebSocket handshake completes. go-nostr stores the challenge in an
|
|
// unexported field populated by the background read loop. If we call
|
|
// Auth() before the challenge arrives, the auth event carries an empty
|
|
// challenge tag and the relay rejects it permanently ("verification
|
|
// failed" → "authentication already failed" on all subsequent attempts).
|
|
//
|
|
// We cannot inspect the challenge field directly (unexported), so we wait
|
|
// briefly to give the read loop time to receive and store it. The relay
|
|
// typically delivers the challenge within a few hundred milliseconds.
|
|
time.Sleep(2 * time.Second)
|
|
|
|
if err := relay.Auth(c.ctx, func(evt *nostr.Event) error {
|
|
return evt.Sign(c.secretKey)
|
|
}); err != nil {
|
|
_ = relay.Close()
|
|
c.cancel()
|
|
return fmt.Errorf("buzz NIP-42 auth failed: %w", err)
|
|
}
|
|
|
|
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)
|
|
}
|