Calling relay.Auth() immediately after RelayConnect signed an auth event with an empty challenge tag. go-nostr keeps the relay's challenge on an unexported field populated by its reader goroutine when the AUTH envelope arrives, so the value is not yet set at connect time, and relays reject the resulting event. Subscribe first and run the handshake only once the relay answers "auth-required". That guarantees the challenge has been read: the CLOSED envelope is processed after AUTH on the relay's single reader goroutine, and receiving it over a channel establishes the happens-before edge that makes the read safe. A second rejection after authenticating is a permissions failure rather than a mistimed handshake, so it fails instead of retrying. Both waits are bounded by a 15s timeout, and relays that accept without auth still work via EndOfStoredEvents. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
287 lines
8 KiB
Go
287 lines
8 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
|
|
|
|
// subscribeTimeout bounds how long Start waits for the relay to accept or
|
|
// reject a subscription before giving up.
|
|
const subscribeTimeout = 15 * time.Second
|
|
|
|
// 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
|
|
|
|
channelIDs := []string(c.config.Channels)
|
|
filters := nostr.Filters{{
|
|
Kinds: []int{kindStreamMessage},
|
|
Tags: nostr.TagMap{"h": channelIDs},
|
|
}}
|
|
|
|
sub, err := c.subscribeAuthed(c.ctx, filters)
|
|
if err != nil {
|
|
_ = relay.Close()
|
|
c.cancel()
|
|
return 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
|
|
}
|
|
|
|
// subscribeAuthed subscribes, performing the NIP-42 handshake if the relay
|
|
// demands it.
|
|
//
|
|
// The handshake MUST be driven by the relay's rejection rather than attempted
|
|
// eagerly after connect. go-nostr stores the relay's challenge on an unexported
|
|
// field populated by its reader goroutine when the AUTH envelope arrives;
|
|
// calling Auth() straight after RelayConnect signs an auth event with an empty
|
|
// challenge tag, which every relay rejects. Waiting for "auth-required"
|
|
// guarantees the challenge has been read — the CLOSED envelope is processed
|
|
// after the AUTH envelope on that same goroutine, and receiving it over a
|
|
// channel establishes the happens-before edge that makes the read safe.
|
|
func (c *BuzzChannel) subscribeAuthed(
|
|
ctx context.Context,
|
|
filters nostr.Filters,
|
|
) (*nostr.Subscription, error) {
|
|
sub, err := c.relay.Subscribe(ctx, filters)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("buzz subscribe failed: %w", err)
|
|
}
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil, ctx.Err()
|
|
|
|
case <-sub.EndOfStoredEvents:
|
|
// Relay accepted the subscription without requiring authentication.
|
|
return sub, nil
|
|
|
|
case reason := <-sub.ClosedReason:
|
|
if !strings.HasPrefix(reason, "auth-required") {
|
|
return nil, fmt.Errorf("buzz relay closed subscription: %s", reason)
|
|
}
|
|
sub.Unsub()
|
|
|
|
if err := c.relay.Auth(ctx, func(evt *nostr.Event) error {
|
|
return evt.Sign(c.secretKey)
|
|
}); err != nil {
|
|
return nil, fmt.Errorf("buzz NIP-42 auth failed: %w", err)
|
|
}
|
|
logger.InfoCF("buzz", "Authenticated to relay", map[string]any{
|
|
"relay": c.config.RelayURL,
|
|
"pubkey": c.publicKey,
|
|
})
|
|
|
|
authed, err := c.relay.Subscribe(ctx, filters)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("buzz subscribe after auth failed: %w", err)
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil, ctx.Err()
|
|
case <-authed.EndOfStoredEvents:
|
|
return authed, nil
|
|
case reason := <-authed.ClosedReason:
|
|
// A second rejection means the identity is not permitted, not that
|
|
// the handshake was mistimed — retrying would loop forever.
|
|
return nil, fmt.Errorf("buzz relay rejected subscription after auth: %s", reason)
|
|
case <-time.After(subscribeTimeout):
|
|
return nil, fmt.Errorf("buzz timed out waiting for subscription after auth")
|
|
}
|
|
|
|
case <-time.After(subscribeTimeout):
|
|
return nil, fmt.Errorf("buzz timed out waiting for relay to accept subscription")
|
|
}
|
|
}
|
|
|
|
// 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)
|
|
}
|