package relay import ( "runtime" "git.smesh.lol/musiquay/web/common/jsbridge/ws" "git.smesh.lol/nostr/pkg/core" ) // State constants for connection readiness. const ( StateConnecting = 0 StateOpen = 1 StateClosed = 2 ) // Conn is a single relay WebSocket connection. type Conn struct { URL string wsConn ws.Conn state int32 subs map[string]*Sub // Callbacks set before Dial returns. onReady func(bool) onEvent func(string, *nostr.Event) onEOSE func(string) onOK func(string, bool, string) onAuth func(string) closing bool // true when Close() was called intentionally everOpened bool // false until the first successful open pendingSend []string // REQ/CLOSE queued when the socket was not writable failCount int32 // consecutive connection failures for backoff // ScheduleReconnect is set by the consumer to provide a delayed callback. // Receives the desired delay in milliseconds so the consumer can apply it. ScheduleReconnect func(delayMs int32, fn func()) } // Dial opens a connection to a relay. // Call OnReady to receive the open/fail notification, then Start to begin processing. func Dial(url string) (c *Conn) { // The connection outlives this frame: the caller keeps it for the life of // the connection, and the ws callbacks registered in dial() fire long // after Dial returns. A frame-arena Conn is copied into the caller's // arena on return, which leaves both the callbacks and the address-keyed // sovereign arena that holds them bound to the old address - the caller's // Subscribe() writes then land in a different object than the one the open // callback flushes. prev := runtime.CurrentArena() runtime.SovereignSetArena(runtime.RootArena()) c = &Conn{ URL: url, state: StateConnecting, subs: map[string]*Sub{}, } c.dial() runtime.SovereignRestoreArena(prev) return c } func (c *Conn) dial() { c.wsConn = ws.Dial( c.URL, func(connID int32, data string) { c.handleMessage(data) }, func(connID int32) { c.state = StateOpen c.failCount = 0 // reset backoff on successful connect first := !c.everOpened c.everOpened = true if c.onReady != nil { c.onReady(true) c.onReady = nil } // Only a RE-open has to re-establish subscriptions. On the first // open every REQ is already in hand: Subscribe sent it, or queued it // in pendingSend when the socket was not writable. Flushing here as // well sent each REQ a second time - and that second send went out // as a zero-filled buffer of the same length, which the relay could // not parse and which stopped it answering that connection, so the // later subscriptions on it (the feed) never got their stored event. c.flushPendingSend() if !first { c.flushSubs() } }, func(connID int32, code int32, reason string) { c.state = StateClosed if c.onReady != nil { c.onReady(false) c.onReady = nil } c.maybeReconnect() }, func(connID int32) { c.state = StateClosed if c.onReady != nil { c.onReady(false) c.onReady = nil } c.maybeReconnect() }, ) } func (c *Conn) maybeReconnect() { if c.closing || len(c.subs) == 0 || c.ScheduleReconnect == nil { return } c.failCount++ // Exponential backoff: 5s → 30s → 5min → 30min cap. // Dead relays stop spamming logs after a few attempts. var delayMs int32 switch { case c.failCount <= 2: delayMs = 5000 case c.failCount <= 5: delayMs = 30000 case c.failCount <= 8: delayMs = 300000 // 5 min default: delayMs = 1800000 // 30 min } c.state = StateConnecting // prevent pool from creating a duplicate c.ScheduleReconnect(delayMs, func() { if c.closing { return } c.dial() }) } // OnReady sets a callback that fires once when the connection opens (true) or fails (false). func (c *Conn) OnReady(fn func(bool)) { if c.state == StateOpen { fn(true) return } if c.state == StateClosed { fn(false) return } c.onReady = fn } // IsOpen returns whether the connection is open. func (c *Conn) IsOpen() (ok bool) { return c.state == StateOpen } func (c *Conn) handleMessage(msg string) { label, subID, payload := nostr.ParseRelayMessage(msg) switch label { case "EVENT": ev := nostr.ParseEvent(payload) if ev == nil { return } if sub, ok2 := c.subs[subID]; ok2 { if sub.OnEvent != nil { sub.OnEvent(ev) } } if c.onEvent != nil { c.onEvent(subID, ev) } case "EOSE": if sub, ok2 := c.subs[subID]; ok2 { sub.gotEOSE = true if sub.OnEOSE != nil { sub.OnEOSE() } } if c.onEOSE != nil { c.onEOSE(subID) } case "OK": ok := len(payload) > 0 && payload[0] == 't' reason := "" idx := indexOf(payload, ':') if idx >= 0 && idx+1 < len(payload) { reason = payload[idx+1:] } if c.onOK != nil { c.onOK(subID, ok, reason) } case "AUTH": if c.onAuth != nil { c.onAuth(payload) } case "NOTICE": _ = payload } } // Subscribe sends a REQ and tracks the subscription. // If the connection is still opening, the REQ is deferred to flushSubs(). func (c *Conn) Subscribe(id string, filters []*nostr.Filter) (s *Sub) { // The subscription outlives this frame: the Conn keeps it in c.subs and the // ws callbacks read it long after Subscribe returns. Build it, and the id it // is keyed by, in the root arena the Conn itself lives in. filters are the // caller's: they must already live somewhere that outlives this call. prev := runtime.CurrentArena() runtime.SovereignSetArena(runtime.RootArena()) idCopy := []byte{:len(id)} copy(idCopy, id) sub := &Sub{ ID: string(idCopy), Filters: filters, conn: c, } c.subs[sub.ID] = sub runtime.SovereignRestoreArena(prev) msg := "[\"REQ\",\"" | sub.ID | "\"" for _, f := range filters { msg |= "," | f.Serialize() } msg |= "]" c.sendOrQueue(msg) return sub } // Publish sends an EVENT message. Queues if the connection is still opening. func (c *Conn) Publish(ev *nostr.Event) { msg := "[\"EVENT\"," | eventJSON(ev) | "]" if c.state != StateOpen { c.queueSend(msg) return } c.sendOrQueue(msg) } // CloseSubscription sends a CLOSE message. func (c *Conn) CloseSubscription(id string) { delete(c.subs, id) msg := "[\"CLOSE\",\"" | id | "\"]" c.sendOrQueue(msg) } // Send sends a raw JSON message string. func (c *Conn) Send(msg string) { ws.Send(c.wsConn, msg) } // Close closes the connection intentionally (no reconnect). func (c *Conn) Close() { c.closing = true c.state = StateClosed for id, sub := range c.subs { sub.Filters = nil sub.OnEvent = nil sub.OnEOSE = nil sub.conn = nil delete(c.subs, id) } c.subs = nil c.onEvent = nil c.onEOSE = nil c.onOK = nil c.onAuth = nil c.onReady = nil c.ScheduleReconnect = nil ws.Close(c.wsConn) } // SetOnEvent sets a global event handler (all subscriptions). func (c *Conn) SetOnEvent(fn func(string, *nostr.Event)) { c.onEvent = fn } // SetOnEOSE sets a global EOSE handler. func (c *Conn) SetOnEOSE(fn func(string)) { c.onEOSE = fn } // SetOnOK sets a handler for OK responses. func (c *Conn) SetOnOK(fn func(string, bool, string)) { c.onOK = fn } // SetOnAuth sets a handler for AUTH challenges. func (c *Conn) SetOnAuth(fn func(string)) { c.onAuth = fn } // sendOrQueue writes msg when the worker's socket is actually writable and // queues it otherwise. // // The Moxie state flag and the WebSocket's readyState are set by different // callbacks, so state can read "open" while the socket is still connecting; // ws.Send reports that as false. Ignoring that result silently dropped the // REQ, leaving the subscription absent on that relay with no error anywhere. // A client-visible symptom: the feed subscription - the one with an event to // deliver - never received its stored event, while subscriptions answered // from the local store looked healthy. func (c *Conn) sendOrQueue(msg string) { if ws.ReadyState(c.wsConn) == 1 && ws.Send(c.wsConn, msg) { return } c.queueSend(msg) } // queueSend stores msg for a later flush. The copy into the root arena is the // point: msg was built in the caller's frame, and a string written into the // receiver has to outlive that frame. Pushing the frame-backed string instead // left pendingSend holding freed memory, and flushing it sent a zero-filled // buffer of the right length - which the relay cannot parse, after which it // stops answering that connection. func (c *Conn) queueSend(msg string) { prev := runtime.CurrentArena() runtime.SovereignSetArena(runtime.RootArena()) cp := []byte{:len(msg)} copy(cp, msg) kept := string(cp) runtime.SovereignRestoreArena(prev) if len(c.pendingSend) >= 256 { c.pendingSend = push(c.pendingSend[1:], kept) return } c.pendingSend = push(c.pendingSend, kept) } // flushPendingSend re-sends what could not be written earlier. func (c *Conn) flushPendingSend() { if len(c.pendingSend) == 0 { return } for _, m := range c.pendingSend { ws.Send(c.wsConn, m) } c.pendingSend = c.pendingSend[:0] } // flushSubs re-sends REQ for all stored subscriptions (used after WS opens). func (c *Conn) flushSubs() { for _, sub := range c.subs { msg := "[\"REQ\",\"" | sub.ID | "\"" for _, f := range sub.Filters { msg |= "," | f.Serialize() } msg |= "]" ws.Send(c.wsConn, msg) } } func eventJSON(ev *nostr.Event) (s string) { // Everything is appended to this buffer within its capacity: `|` does not // grow, so the size has to be right before the first write. The helpers // return their own complete piece, never a buffer they were handed: a // returned slice comes back with len == cap (the spare capacity is not part // of the value), and appending to that hands back a full slice, so the next // append fails loud with "push past end of slice capacity". buf := []byte{:0:jsonSize(ev)} buf = buf | "{" buf = buf | "\"id\":\"" buf = buf | escapeJSON(ev.ID) buf = buf | "\",\"pubkey\":\"" buf = buf | escapeJSON(ev.PubKey) buf = buf | "\",\"created_at\":" buf = buf | itoa(ev.CreatedAt) buf = buf | ",\"kind\":" buf = buf | itoa(int64(ev.Kind)) buf = buf | ",\"tags\":" buf = buf | tagsJSON(ev.Tags) buf = buf | ",\"content\":\"" buf = buf | escapeJSON(ev.Content) buf = buf | "\",\"sig\":\"" buf = buf | escapeJSON(ev.Sig) buf = buf | "\"}" return string(buf) } // jsonSize upper-bounds the serialized event: escaping can double a byte and // each tag element costs its quotes plus a separator. func jsonSize(ev *nostr.Event) (n int32) { n = 256 + len(ev.ID)*2 + len(ev.PubKey)*2 + len(ev.Sig)*2 + len(ev.Content)*2 for _, tag := range ev.Tags { n += 4 for _, el := range tag { n += len(el)*2 + 4 } } return n } // tagsJSON renders the tag array whole. Its own buffer is sized from the tags, // so the caller appends a finished piece rather than a half-built buffer. func tagsJSON(tags [][]string) (out string) { n := int32(2) for _, tag := range tags { n += 4 for _, el := range tag { n += len(el)*2 + 4 } } buf := []byte{:0:n} buf = buf | "[" for i, tag := range tags { if i > 0 { buf = buf | "," } buf = buf | "[" for j, s := range tag { if j > 0 { buf = buf | "," } buf = buf | "\"" buf = buf | escapeJSON(s) buf = buf | "\"" } buf = buf | "]" } buf = buf | "]" return string(buf) } // escapeJSON returns the JSON string body for s, quotes excluded. func escapeJSON(s string) (out string) { buf := []byte{:0:len(s)*2 + 8} for i := 0; i < len(s); i++ { c := s[i] switch c { case '"': buf = buf | "\\\"" case '\\': buf = buf | "\\\\" case '\n': buf = buf | "\\n" case '\r': buf = buf | "\\r" case '\t': buf = buf | "\\t" default: buf = push(buf, c) } } return string(buf) } func itoa(n int64) (s string) { if n == 0 { return "0" } neg := false if n < 0 { neg = true n = -n } var b [20]byte i := len(b) for n > 0 { i-- b[i] = byte('0' + n%10) n /= 10 } if neg { i-- b[i] = '-' } return string(b[i:]) } func indexOf(s string, c byte) (n int32) { for i := 0; i < len(s); i++ { if s[i] == c { return i } } return -1 }