picoclaw/pkg/channels/buzz/buzz.go
PeterChrz 5acf91a136
feat(channels/buzz): implement ReactionCapable for 👀 emoji reactions
Add ReactToMessage method that publishes a NIP-25 kind:7 reaction event
with content 👀, scoped to the same channel via the h tag and pointing at
the inbound message via the e tag. The undo function is a no-op since
NIP-25 does not define a standard way to remove a reaction.

This enables the BaseChannel auto-reaction pipeline: when a message
arrives on a Buzz channel, the bot reacts with 👀 before processing,
giving the Emperor visual confirmation that the message was received.
2026-08-12 09:41:32 -04:00

328 lines
9.1 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
}
// kindReaction is the NIP-25 reaction event kind.
const kindReaction = 7
// ReactToMessage implements channels.ReactionCapable.
// It publishes a NIP-25 kind:7 reaction event with content 👀 scoped to the
// same channel as the reacted message. The undo function is a no-op because
// NIP-25 does not define a standard way to remove a reaction.
func (c *BuzzChannel) ReactToMessage(ctx context.Context, chatID, messageID string) (func(), error) {
if !c.IsRunning() {
return func() {}, channels.ErrNotRunning
}
if messageID == "" {
return func() {}, nil
}
evt := nostr.Event{
PubKey: c.publicKey,
CreatedAt: nostr.Now(),
Kind: kindReaction,
Tags: nostr.Tags{
nostr.Tag{"e", messageID},
nostr.Tag{"h", chatID},
},
Content: "👀",
}
if err := evt.Sign(c.secretKey); err != nil {
return func() {}, fmt.Errorf("buzz reaction sign failed: %w", err)
}
if err := c.relay.Publish(ctx, evt); err != nil {
return func() {}, fmt.Errorf("buzz reaction publish failed: %w", err)
}
logger.DebugCF("buzz", "Reaction sent", map[string]any{
"channel": chatID,
"event_id": evt.ID,
"target_id": messageID,
})
return func() {}, 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)
}