client.go raw

   1  package nostr
   2  
   3  import (
   4  	"context"
   5  	"encoding/json"
   6  	"fmt"
   7  	"sync"
   8  
   9  	"github.com/coder/websocket"
  10  )
  11  
  12  // OKResponse is the relay's response to an EVENT submission.
  13  type OKResponse struct {
  14  	EventID  string
  15  	Accepted bool
  16  	Message  string
  17  }
  18  
  19  // Client connects to a Nostr relay and handles subscriptions.
  20  type Client struct {
  21  	URL  string
  22  	conn *websocket.Conn
  23  	mu   sync.Mutex
  24  
  25  	// Events receives validated events from active subscriptions.
  26  	Events chan *Event
  27  
  28  	// Notices receives NOTICE messages from the relay.
  29  	Notices chan string
  30  
  31  	// OKs receives OK responses from the relay after EVENT submissions.
  32  	OKs chan OKResponse
  33  }
  34  
  35  // Connect establishes a WebSocket connection to the relay.
  36  func Connect(ctx context.Context, url string) (*Client, error) {
  37  	conn, _, err := websocket.Dial(ctx, url, nil)
  38  	if err != nil {
  39  		return nil, fmt.Errorf("dial %s: %w", url, err)
  40  	}
  41  	conn.SetReadLimit(1 << 20) // 1MB
  42  
  43  	c := &Client{
  44  		URL:     url,
  45  		conn:    conn,
  46  		Events:  make(chan *Event, 256),
  47  		Notices: make(chan string, 16),
  48  		OKs:     make(chan OKResponse, 16),
  49  	}
  50  
  51  	return c, nil
  52  }
  53  
  54  // Subscribe sends a REQ message to the relay.
  55  func (c *Client) Subscribe(ctx context.Context, subID string, filters ...Filter) error {
  56  	msg := make([]any, 0, 2+len(filters))
  57  	msg = append(msg, "REQ", subID)
  58  	for _, f := range filters {
  59  		msg = append(msg, f)
  60  	}
  61  
  62  	data, err := json.Marshal(msg)
  63  	if err != nil {
  64  		return err
  65  	}
  66  
  67  	c.mu.Lock()
  68  	defer c.mu.Unlock()
  69  	return c.conn.Write(ctx, websocket.MessageText, data)
  70  }
  71  
  72  // Close sends a CLOSE message for a subscription.
  73  func (c *Client) CloseSubscription(ctx context.Context, subID string) error {
  74  	msg, _ := json.Marshal([]string{"CLOSE", subID})
  75  	c.mu.Lock()
  76  	defer c.mu.Unlock()
  77  	return c.conn.Write(ctx, websocket.MessageText, msg)
  78  }
  79  
  80  // Publish sends an EVENT message to the relay.
  81  func (c *Client) Publish(ctx context.Context, event *Event) error {
  82  	msg, _ := json.Marshal([]any{"EVENT", event})
  83  	c.mu.Lock()
  84  	defer c.mu.Unlock()
  85  	return c.conn.Write(ctx, websocket.MessageText, msg)
  86  }
  87  
  88  // Listen reads messages from the relay and dispatches them.
  89  // It blocks until the context is cancelled or the connection closes.
  90  // Events are validated before being sent to the Events channel.
  91  func (c *Client) Listen(ctx context.Context) error {
  92  	for {
  93  		_, data, err := c.conn.Read(ctx)
  94  		if err != nil {
  95  			return err
  96  		}
  97  
  98  		// Parse the envelope: ["TYPE", ...]
  99  		var envelope []json.RawMessage
 100  		if json.Unmarshal(data, &envelope) != nil || len(envelope) < 2 {
 101  			continue
 102  		}
 103  
 104  		var msgType string
 105  		if json.Unmarshal(envelope[0], &msgType) != nil {
 106  			continue
 107  		}
 108  
 109  		switch msgType {
 110  		case "EVENT":
 111  			// ["EVENT", "<sub_id>", <event>]
 112  			if len(envelope) < 3 {
 113  				continue
 114  			}
 115  			var event Event
 116  			if json.Unmarshal(envelope[2], &event) != nil {
 117  				continue
 118  			}
 119  			// Validate before delivering.
 120  			if event.Valid() {
 121  				select {
 122  				case c.Events <- &event:
 123  				default:
 124  					// Events channel full — drop.
 125  				}
 126  			}
 127  
 128  		case "EOSE":
 129  			// End of stored events — we don't need to act on this yet.
 130  
 131  		case "NOTICE":
 132  			if len(envelope) >= 2 {
 133  				var notice string
 134  				if json.Unmarshal(envelope[1], &notice) == nil {
 135  					select {
 136  					case c.Notices <- notice:
 137  					default:
 138  					}
 139  				}
 140  			}
 141  
 142  		case "OK":
 143  			// ["OK", "<event_id>", <accepted>, "<message>"]
 144  			if len(envelope) >= 4 {
 145  				var eventID string
 146  				var accepted bool
 147  				var message string
 148  				json.Unmarshal(envelope[1], &eventID)
 149  				json.Unmarshal(envelope[2], &accepted)
 150  				json.Unmarshal(envelope[3], &message)
 151  				select {
 152  				case c.OKs <- OKResponse{EventID: eventID, Accepted: accepted, Message: message}:
 153  				default:
 154  				}
 155  			}
 156  
 157  		case "CLOSED":
 158  			// Subscription closed by relay.
 159  		}
 160  	}
 161  }
 162  
 163  // Disconnect closes the WebSocket connection.
 164  func (c *Client) Disconnect() {
 165  	c.conn.Close(websocket.StatusNormalClosure, "")
 166  }
 167