main.mx raw
1 package main
2
3 import (
4 "runtime"
5 "git.smesh.lol/moxie/pkg/mxutil"
6 "git.smesh.lol/musiquay/web/common/helpers"
7 "git.smesh.lol/musiquay/web/common/jsbridge/relayproxy"
8 "git.smesh.lol/musiquay/web/common/jsbridge/ws"
9 "git.smesh.lol/musiquay/web/common/mw"
10 "git.smesh.lol/nostr/pkg/core"
11 "git.smesh.lol/musiquay/web/common/relay"
12 )
13
14 // Relay-proxy worker.
15 //
16 // Owns: WS pool to remote nostr relays, subscription router, MLS event
17 // subscription routing (kinds 443/445/1059 -> MLS worker via supervisor).
18 //
19 // IDB is owned by the Store Worker. This worker sends S_* messages to the
20 // supervisor which routes them to the Store Worker, and receives SR_* responses.
21 //
22 // Wire format: JSON-array MW strings.
23 // page -> worker:
24 // ["REQ", subID, filter]
25 // ["CLOSE", subID]
26 // ["PROXY", subID, filter, [relayURLs]]
27 // ["EVENT", signedEventJSON] -- publish to ps.writeRelays
28 // ["PUBLISH_TO", signedEventJSON, [relayURLs]] -- explicit targets
29 // ["SET_PUBKEY", hex]
30 // ["SET_WRITE_RELAYS", [relayURLs]]
31 // ["ENC_KEY", hex]
32 // ["MLS_SUB", [relayURLs], [groupIDs]] -- open persistent MLS subs
33 // ["MLS_UPDATE_GROUPS", [groupIDs]] -- update kind 445 #h filter
34 // ["MLS_FETCH_KP", peer, [relayURLs]] -- one-shot key package fetch
35 // ["SR_QUERY", reqID, eventsJSON] -- response from Store Worker
36 // worker -> page:
37 // ["READY"]
38 // ["EVENT", subID, event]
39 // ["EOSE", subID]
40 // ["OK", eventID, ok, message]
41 // ["MLS_EVENT", evJSON] -- kind 443/445/1059 -> MLS worker
42 // ["MLS_KP_RESULT", peer, evJSON] -- one-shot KP fetch result -> MLS worker
43 // ["S_PUT_EVENT", evJSON] -- to Store Worker (fire-and-forget)
44 // ["S_ENC_KEY", hexKey] -- to Store Worker (fire-and-forget)
45 // ["S_QUERY", reqID, filterJSON] -- to Store Worker
46
47 type peerRelayInfo struct {
48 urls []string // relays the peer reads from (or both read+write)
49 ts int64 // created_at of the kind 10002 event
50 }
51
52 type clientSub struct {
53 filter *nostr.Filter
54 filterRaw string
55 }
56
57 type proxySub struct {
58 remoteIDs map[string]bool
59 relayCount int32
60 timer int32
61 done bool
62 live bool
63 // Initial load ends when BOTH sources are done: the local store query and
64 // the relays (or the fallback timer). Whichever finishes last decides, so
65 // events from the slower source are rendered as part of the initial load
66 // instead of being buffered as new posts.
67 localDone bool
68 remoteDone bool
69 }
70
71 // proxyState is the relay-proxy worker's mutable state. Package globals are
72 // immutable outside initialization, so its timers, subscription maps and
73 // in-flight fetch bookkeeping live in one self-mutating type reached through
74 // a package-level pointer.
75 type proxyState struct {
76 // eoseTimeoutMs: how long to wait for relays to EOSE before emitting EOSE
77 // to the consumer. Short value gets the consumer unblocked quickly so it
78 // can render cached data; the linger window catches stragglers.
79 eoseTimeoutMs int32
80 // lingerMs: after EOSE is emitted to the consumer, keep the remote subs
81 // open this long to forward late-arriving events. Solves the slow-relay
82 // race where events arrived just after the EOSE timeout fired.
83 lingerMs int32
84
85 clientSubs map[string]*clientSub
86 proxySubs map[string]*proxySub
87 remoteToProxy map[string]string
88 rpool *relay.Pool
89 writeRelays []string
90 myPubkey string
91
92 mlsSubIDs map[string]bool // MLS event subscriptions (kinds 443/445/1059)
93 mlsKPFetchID string // in-flight one-shot KP fetch prefix
94 mlsKPFetchPeer string // peer pubkey for in-flight KP fetch
95 mlsKPFetchSubs []string // per-relay sub IDs for in-flight fetch
96 mlsGroupIDs []string // current #h filter for kind 445
97
98 storeCBs map[int32]func(string) // reqID -> callback for SR_QUERY responses
99 nextStoreID int32
100
101 // Peer relay list cache (kind 10002).
102 peerRelayCache map[string]*peerRelayInfo // pubkey -> cached relay list
103 rlFetchID string // in-flight relay list fetch sub prefix
104 rlFetchPeer string // peer for in-flight relay list fetch
105 rlFetchSubs []string // per-relay sub IDs
106 rlFetchCB func([]string) // callback when relay list is resolved
107 rlFetchEOSE int32 // count of EOSE received
108 }
109
110 var ps *proxyState
111
112 func initState() {
113 if ps != nil {
114 return
115 }
116 ps = &proxyState{}
117 ps.eoseTimeoutMs = 3000
118 ps.lingerMs = 25000
119 ps.clientSubs = map[string]*clientSub{}
120 ps.proxySubs = map[string]*proxySub{}
121 ps.remoteToProxy = map[string]string{}
122 ps.rpool = relay.NewPool()
123 ps.mlsSubIDs = map[string]bool{}
124 ps.storeCBs = map[int32]func(string){}
125 ps.peerRelayCache = map[string]*peerRelayInfo{}
126 }
127
128 func main() {
129 initState()
130 relayproxy.WorkerOnMessage(handleMessage)
131 relayproxy.WorkerPost(`["READY"]`)
132 }
133
134 // storeQuery sends S_QUERY to the Store Worker and calls fn with the eventsJSON response.
135 func storeQuery(filterRaw string, fn func(string)) {
136 ps.nextStoreID++
137 reqID := ps.nextStoreID
138 ps.storeCBs[reqID] = fn
139 relayproxy.WorkerPost(`["S_QUERY",` | helpers.Itoa(int64(reqID)) | `,` | filterRaw | `]`)
140 rid := reqID
141 relayproxy.WorkerSetTimeout(30000, func() {
142 if _, ok2 := ps.storeCBs[rid]; ok2 {
143 delete(ps.storeCBs, rid)
144 }
145 })
146 }
147
148 func handleMessage(msg string) {
149 w := mw.New(msg)
150 msgType := w.Str()
151 switch msgType {
152 case "REQ":
153 subID := dup(w.Str())
154 filterRaw := dup(w.Raw())
155 handleReq(subID, filterRaw)
156 case "CLOSE":
157 subID := dup(w.Str())
158 handleClose(subID)
159 case "PROXY":
160 subID := dup(w.Str())
161 filterRaw := dup(w.Raw())
162 relayURLs := dupStrs(w.Strs())
163 handleProxy(subID, filterRaw, relayURLs, false)
164 case "PROXY_LIVE":
165 subID := dup(w.Str())
166 filterRaw := dup(w.Raw())
167 relayURLs := dupStrs(w.Strs())
168 handleProxy(subID, filterRaw, relayURLs, true)
169 case "EVENT":
170 eventRaw := w.Raw()
171 handleEventPublish(eventRaw)
172 case "PUBLISH_TO":
173 eventRaw := w.Raw()
174 relayURLs := w.Strs()
175 handlePublishTo(eventRaw, relayURLs)
176 case "SET_WRITE_RELAYS":
177 urls := dupStrs(w.Strs())
178 handleSetWriteRelays(urls)
179 case "SET_PUBKEY":
180 ps.myPubkey = dup(w.Str())
181 case "ENC_KEY":
182 hexKey := w.Str()
183 relayproxy.WorkerPost(`["S_ENC_KEY",` | jstr(hexKey) | `]`)
184 case "SET_PROXY_EOSE_MS":
185 ms := int32(w.Num())
186 if ms > 0 {
187 ps.eoseTimeoutMs = ms
188 }
189 case "MLS_SUB":
190 urls := dupStrs(w.Strs())
191 groupIDs := dupStrs(w.Strs())
192 handleMLSSub(urls, groupIDs)
193 case "MLS_UPDATE_GROUPS":
194 groupIDs := dupStrs(w.Strs())
195 handleMLSUpdateGroups(groupIDs)
196 case "MLS_FETCH_KP":
197 peer := dup(w.Str())
198 urls := dupStrs(w.Strs())
199 handleMLSFetchKP(peer, urls)
200 // Store Worker response for storeQuery.
201 case "SR_QUERY":
202 reqID := int32(w.Num())
203 eventsJSON := w.Raw()
204 if fn, ok2 := ps.storeCBs[reqID]; ok2 {
205 fn(eventsJSON)
206 delete(ps.storeCBs, reqID)
207 }
208 }
209 }
210
211 func handleMLSSub(relayURLs, groupIDs []string) {
212 println("[relay-proxy] handleMLSSub: urls=" | helpers.Itoa(int64(len(relayURLs))) | " groupIDs=" | helpers.Itoa(int64(len(groupIDs))) | " pubkey=" | ps.myPubkey[:16] | "...")
213 if ps.myPubkey == "" || len(relayURLs) == 0 {
214 println("[relay-proxy] handleMLSSub: no pubkey or no URLs, aborting")
215 return
216 }
217 for rSubID := range ps.mlsSubIDs {
218 for _, c := range ps.rpool.AllConns() {
219 c.CloseSubscription(rSubID)
220 }
221 }
222 ps.mlsSubIDs = map[string]bool{}
223 ps.mlsGroupIDs = groupIDs
224
225 for _, url := range relayURLs {
226 url = normalizeRelayURL(url)
227 if !isAllowedRelay(url) {
228 println("[relay-proxy] handleMLSSub: blocked url=" | url)
229 continue
230 }
231 suffix := urlSuffix(url)
232 idP := "mlsp_" | suffix
233 ps.mlsSubIDs[idP] = true
234 c := getConn(url)
235 println("[relay-proxy] handleMLSSub: sub " | idP | " kinds=[443,1059] #p=" | ps.myPubkey[:16] | "... on " | url)
236 c.Subscribe(idP, []*nostr.Filter{{
237 Kinds: []uint32{443, 1059},
238 Tags: map[string][]string{"#p": {ps.myPubkey}},
239 }})
240 if len(ps.mlsGroupIDs) > 0 {
241 idH := "mlsh_" | suffix
242 ps.mlsSubIDs[idH] = true
243 println("[relay-proxy] handleMLSSub: sub " | idH | " kind=445 #h groups=" | helpers.Itoa(int64(len(ps.mlsGroupIDs))))
244 c.Subscribe(idH, []*nostr.Filter{{
245 Kinds: []uint32{445},
246 Tags: map[string][]string{"#h": ps.mlsGroupIDs},
247 }})
248 }
249 }
250 }
251
252 func handleMLSUpdateGroups(groupIDs []string) {
253 ps.mlsGroupIDs = groupIDs
254 for _, c := range ps.rpool.AllConns() {
255 suffix := urlSuffix(c.URL)
256 idH := "mlsh_" | suffix
257 c.CloseSubscription(idH)
258 delete(ps.mlsSubIDs, idH)
259 if len(ps.mlsGroupIDs) > 0 && ps.mlsSubIDs["mlsp_"|suffix] {
260 ps.mlsSubIDs[idH] = true
261 c.Subscribe(idH, []*nostr.Filter{{
262 Kinds: []uint32{445},
263 Tags: map[string][]string{"#h": ps.mlsGroupIDs},
264 }})
265 }
266 }
267 }
268
269 func handleMLSFetchKP(peer string, relayURLs []string) {
270 println("[relay-proxy] handleMLSFetchKP: peer=" | peer[:16] | "... urls=" | helpers.Itoa(int64(len(relayURLs))))
271 if peer == "" {
272 println("[relay-proxy] handleMLSFetchKP: no peer")
273 relayproxy.WorkerPost(`["MLS_KP_RESULT",` | jstr(peer) | `,null]`)
274 return
275 }
276 // Resolve shared relays: fetch peer's kind 10002, intersect with ours.
277 // Falls back to provided URLs if no intersection found.
278 fallback := relayURLs
279 if len(fallback) == 0 {
280 fallback = ps.writeRelays
281 }
282 resolveSharedRelays(peer, fallback, func(urls []string) {
283 if len(urls) == 0 {
284 println("[relay-proxy] handleMLSFetchKP: no relays after resolution")
285 relayproxy.WorkerPost(`["MLS_KP_RESULT",` | jstr(peer) | `,null]`)
286 return
287 }
288 doKPFetch(peer, urls)
289 })
290 }
291
292 func doKPFetch(peer string, relayURLs []string) {
293 if ps.mlsKPFetchID != "" {
294 closeKPFetchSubs()
295 }
296 ps.mlsKPFetchID = "mlskp_" | peer[:8]
297 ps.mlsKPFetchPeer = peer
298 ps.mlsKPFetchSubs = nil
299 println("[relay-proxy] doKPFetch: peer=" | peer[:16] | "... urls=" | helpers.Itoa(int64(len(relayURLs))))
300 for _, rawURL := range relayURLs {
301 url := normalizeRelayURL(rawURL)
302 if !isAllowedRelay(url) {
303 println("[relay-proxy] doKPFetch: blocked url=" | url)
304 continue
305 }
306 doKPFetchDirect(peer, url)
307 return
308 }
309 println("[relay-proxy] doKPFetch: no valid URLs")
310 ps.mlsKPFetchID = ""
311 relayproxy.WorkerPost(`["MLS_KP_RESULT",` | jstr(peer) | `,null]`)
312 }
313
314 // doKPFetchDirect opens a dedicated WebSocket to fetch a single KP event.
315 // Bypasses the pool connection to avoid state interference.
316 func doKPFetchDirect(peer, url string) {
317 subID := "kpq"
318 reqMsg := `["REQ","` | subID | `",{"kinds":[443],"authors":["` | peer | `"],"limit":1}]`
319 println("[relay-proxy] doKPFetchDirect: opening dedicated ws to " | url)
320 println("[relay-proxy] doKPFetchDirect: REQ=" | reqMsg)
321
322 done := false
323 fetchPeer := peer
324
325 wsConn := ws.Dial(url,
326 func(connID int32, data string) {
327 if done {
328 return
329 }
330 println("[relay-proxy] doKPFetchDirect: recv len=" | helpers.Itoa(int64(len(data))))
331 label, rSubID, payload := nostr.ParseRelayMessage(data)
332 println("[relay-proxy] doKPFetchDirect: label=" | label | " subID=" | rSubID)
333 if label == "EVENT" && rSubID == subID {
334 ev := nostr.ParseEvent(payload)
335 if ev != nil && ev.Kind == 443 {
336 done = true
337 evJSON := ev.ToJSON()
338 println("[relay-proxy] doKPFetchDirect: FOUND KP for " | fetchPeer[:16] | "...")
339 ps.mlsKPFetchID = ""
340 ws.Close(ws.Conn(connID))
341 relayproxy.WorkerPost(`["MLS_KP_RESULT",` | jstr(fetchPeer) | `,` | evJSON | `]`)
342 }
343 }
344 if label == "EOSE" && rSubID == subID && !done {
345 done = true
346 println("[relay-proxy] doKPFetchDirect: EOSE with no KP for " | fetchPeer[:16] | "...")
347 ps.mlsKPFetchID = ""
348 ws.Close(ws.Conn(connID))
349 relayproxy.WorkerPost(`["MLS_KP_RESULT",` | jstr(fetchPeer) | `,null]`)
350 }
351 },
352 func(connID int32) {
353 println("[relay-proxy] doKPFetchDirect: connected, sending REQ")
354 ws.Send(ws.Conn(connID), reqMsg)
355 },
356 func(connID int32, code int32, reason string) {
357 if !done {
358 done = true
359 println("[relay-proxy] doKPFetchDirect: ws closed code=" | helpers.Itoa(int64(code)))
360 ps.mlsKPFetchID = ""
361 relayproxy.WorkerPost(`["MLS_KP_RESULT",` | jstr(fetchPeer) | `,null]`)
362 }
363 },
364 func(connID int32) {
365 if !done {
366 done = true
367 println("[relay-proxy] doKPFetchDirect: ws error")
368 ps.mlsKPFetchID = ""
369 relayproxy.WorkerPost(`["MLS_KP_RESULT",` | jstr(fetchPeer) | `,null]`)
370 }
371 },
372 )
373 _ = wsConn
374
375 relayproxy.WorkerSetTimeout(10000, func() {
376 if !done {
377 done = true
378 println("[relay-proxy] doKPFetchDirect: TIMEOUT for " | fetchPeer[:16] | "...")
379 ps.mlsKPFetchID = ""
380 relayproxy.WorkerPost(`["MLS_KP_RESULT",` | jstr(fetchPeer) | `,null]`)
381 }
382 })
383 }
384
385 func closeKPFetchSubs() {
386 for _, sid := range ps.mlsKPFetchSubs {
387 for _, c := range ps.rpool.AllConns() {
388 c.CloseSubscription(sid)
389 }
390 }
391 ps.mlsKPFetchSubs = nil
392 }
393
394 // parseRelayList extracts relay URLs from a kind 10002 event.
395 // Returns URLs where the peer reads (or has no marker = both).
396 func parseRelayList(ev *nostr.Event) (ss []string) {
397 var urls []string
398 for _, tag := range ev.Tags {
399 if len(tag) < 2 || tag[0] != "r" {
400 continue
401 }
402 url := tag[1]
403 if len(url) < 6 {
404 continue
405 }
406 // If marker present, include only "read" or no marker (= both).
407 if len(tag) >= 3 && tag[2] == "write" {
408 continue
409 }
410 urls = mxutil.Ensure(urls, 1)
411 urls = push(urls, normalizeRelayURL(url))
412 }
413 return urls
414 }
415
416 // resolveSharedRelays checks the cache for the peer's relay list.
417 // If cached, returns intersection immediately via cb.
418 // If not cached, fetches kind 10002 from connected relays, then calls cb.
419 func resolveSharedRelays(peer string, fallback []string, cb func([]string)) {
420 if info, ok2 := ps.peerRelayCache[peer]; ok2 {
421 shared := intersectRelays(ps.writeRelays, info.urls)
422 println("[relay-proxy] resolveSharedRelays: cached peer=" | peer[:16] | "... shared=" | helpers.Itoa(int64(len(shared))))
423 if len(shared) > 0 {
424 cb(shared)
425 return
426 }
427 cb(fallback)
428 return
429 }
430 fetchPeerRelayList(peer, func(urls []string) {
431 if len(urls) > 0 {
432 ps.peerRelayCache[peer] = &peerRelayInfo{urls: urls}
433 shared := intersectRelays(ps.writeRelays, urls)
434 println("[relay-proxy] resolveSharedRelays: fetched peer=" | peer[:16] | "... peerURLs=" | helpers.Itoa(int64(len(urls))) | " shared=" | helpers.Itoa(int64(len(shared))))
435 if len(shared) > 0 {
436 cb(shared)
437 return
438 }
439 } else {
440 println("[relay-proxy] resolveSharedRelays: no 10002 for peer=" | peer[:16] | "... using fallback")
441 }
442 cb(fallback)
443 })
444 }
445
446 func fetchPeerRelayList(peer string, cb func([]string)) {
447 if ps.rlFetchID != "" {
448 closeRLFetchSubs()
449 }
450 ps.rlFetchID = "rl10k_" | peer[:8]
451 ps.rlFetchPeer = peer
452 ps.rlFetchSubs = nil
453 ps.rlFetchCB = cb
454 ps.rlFetchEOSE = 0
455
456 connURLs := ps.rpool.URLs()
457 if len(connURLs) == 0 {
458 connURLs = ps.writeRelays
459 }
460 for _, url := range connURLs {
461 url = normalizeRelayURL(url)
462 if !isAllowedRelay(url) {
463 continue
464 }
465 suffix := urlSuffix(url)
466 subID := ps.rlFetchID | "_" | suffix
467 ps.rlFetchSubs = mxutil.Ensure(ps.rlFetchSubs, 1)
468 ps.rlFetchSubs = push(ps.rlFetchSubs, subID)
469 c := getConn(url)
470 println("[relay-proxy] fetchPeerRelayList: sub " | subID | " kind=10002 author=" | peer[:16] | "... on " | url)
471 c.Subscribe(subID, []*nostr.Filter{{
472 Kinds: []uint32{10002},
473 Authors: []string{peer},
474 Limit: 1,
475 }})
476 }
477 fetchPrefix := ps.rlFetchID
478 relayproxy.WorkerSetTimeout(5000, func() {
479 if ps.rlFetchID == fetchPrefix {
480 println("[relay-proxy] fetchPeerRelayList: TIMEOUT for " | ps.rlFetchPeer[:16] | "...")
481 fn := ps.rlFetchCB
482 closeRLFetchSubs()
483 ps.rlFetchID = ""
484 ps.rlFetchCB = nil
485 if fn != nil {
486 fn(nil)
487 }
488 }
489 })
490 }
491
492 func closeRLFetchSubs() {
493 for _, sid := range ps.rlFetchSubs {
494 for _, c := range ps.rpool.AllConns() {
495 c.CloseSubscription(sid)
496 }
497 }
498 ps.rlFetchSubs = nil
499 }
500
501 func intersectRelays(ours, theirs []string) (ss []string) {
502 var out []string
503 for _, u := range ours {
504 for _, t := range theirs {
505 if normalizeRelayURL(u) == normalizeRelayURL(t) {
506 out = mxutil.Ensure(out, 1)
507 out = push(out, u)
508 break
509 }
510 }
511 }
512 return out
513 }
514
515 func isHex(s string) (ok bool) {
516 for i := 0; i < len(s); i++ {
517 c := s[i]
518 if !((c >= '0' && c <= '9') || (c >= 'a' && c <= 'f')) {
519 return false
520 }
521 }
522 return true
523 }
524
525 func handleEventPublish(eventRaw string) {
526 ev := nostr.ParseEvent(eventRaw)
527 if ev == nil {
528 return
529 }
530 relayproxy.WorkerPost(`["S_PUT_EVENT",` | eventRaw | `]`)
531 for _, url := range ps.writeRelays {
532 if isAllowedRelay(url) {
533 getConn(url).Publish(ev)
534 }
535 }
536 relayproxy.WorkerPost(`["OK",` | jstr(ev.ID) | `,true,""]`)
537 }
538
539 func handlePublishTo(eventRaw string, relayURLs []string) {
540 ev := nostr.ParseEvent(eventRaw)
541 if ev == nil {
542 println("[relay-proxy] handlePublishTo: parse failed")
543 return
544 }
545 println("[relay-proxy] handlePublishTo: kind=" | helpers.Itoa(int64(ev.Kind)) | " id=" | ev.ID[:16] | "... to " | helpers.Itoa(int64(len(relayURLs))) | " relays")
546 relayproxy.WorkerPost(`["S_PUT_EVENT",` | eventRaw | `]`)
547 for _, url := range relayURLs {
548 url = normalizeRelayURL(url)
549 if isAllowedRelay(url) {
550 c := getConn(url)
551 oStr := "closed"
552 if c.IsOpen() {
553 oStr = "open"
554 }
555 println("[relay-proxy] handlePublishTo: publishing to " | url | " conn=" | oStr)
556 c.Publish(ev)
557 }
558 }
559 relayproxy.WorkerPost(`["OK",` | jstr(ev.ID) | `,true,""]`)
560 }
561
562 func handleSetWriteRelays(urls []string) {
563 ps.writeRelays = nil
564 for _, u := range urls {
565 nu := normalizeRelayURL(u)
566 if isAllowedRelay(nu) {
567 ps.writeRelays = mxutil.Ensure(ps.writeRelays, 1)
568 ps.writeRelays = push(ps.writeRelays, nu)
569 }
570 }
571 println("[relay-proxy] SET_WRITE_RELAYS: " | helpers.Itoa(int64(len(ps.writeRelays))) | " relays")
572 for _, u := range ps.writeRelays {
573 println("[relay-proxy] relay: " | u)
574 }
575 }
576
577 func handleReq(subID, filterRaw string) {
578 // Re-subscribing on an id restarts it, exactly as handleProxy does.
579 cleanupProxy(subID)
580 // The subscription record and the closure that answers the store query both
581 // outlive this frame; allocate them in the root arena the maps live in.
582 prev := runtime.CurrentArena()
583 runtime.SovereignSetArena(runtime.RootArena())
584 f := nostr.ParseFilter(filterRaw)
585 if f == nil {
586 runtime.SovereignRestoreArena(prev)
587 return
588 }
589 ps.clientSubs[subID] = &clientSub{filter: f, filterRaw: filterRaw}
590 // A plain REQ is the same subscription as a PROXY one with no relays, so it
591 // goes through the same registry and fan-out: the remote half is complete
592 // from the start and the store query alone ends the initial load. Keeping
593 // it in ps.proxySubs means one EOSE policy and one place CLOSE has to clean up.
594 ps.proxySubs[subID] = &proxySub{remoteIDs: map[string]bool{}, remoteDone: true}
595 sid := subID
596 storeQuery(filterRaw, func(eventsJSON string) {
597 events := nostr.ParseEventsJSON(eventsJSON)
598 for _, ev := range events {
599 relayproxy.WorkerPost(`["EVENT",` | jstr(sid) | `,` | ev.ToJSON() | `]`)
600 }
601 if info, ok := ps.proxySubs[sid]; ok {
602 info.localDone = true
603 proxyMaybeEOSE(sid)
604 }
605 })
606 runtime.SovereignRestoreArena(prev)
607 }
608
609 func handleClose(subID string) {
610 delete(ps.clientSubs, subID)
611 cleanupProxy(subID)
612 }
613
614 func handleProxy(subID, filterRaw string, relayURLs []string, live bool) {
615 cleanupProxy(subID)
616
617 // Everything this handler stores in the worker's maps - the subscription
618 // record, its filter, the ids and filter text it is keyed by, and the
619 // closures it hands to the timer service - outlives this frame and is read
620 // by later messages. Allocate the handler's state in the root arena the
621 // maps themselves live in. One PROXY message is one subscription, so the
622 // scratch that stays is bounded by the subscription count.
623 prev := runtime.CurrentArena()
624 runtime.SovereignSetArena(runtime.RootArena())
625
626 f := nostr.ParseFilter(filterRaw)
627 if f == nil {
628 runtime.SovereignRestoreArena(prev)
629 return
630 }
631 ps.clientSubs[subID] = &clientSub{filter: f, filterRaw: filterRaw}
632
633 sid := subID
634 hasSearchField := false
635 for i := 0; i+8 <= len(filterRaw); i++ {
636 if filterRaw[i:i+8] == "\"search\"" {
637 hasSearchField = true
638 break
639 }
640 }
641 if hasSearchField {
642 // No local query runs for a search filter, so that source is complete.
643 } else {
644 storeQuery(filterRaw, func(eventsJSON string) {
645 events := nostr.ParseEventsJSON(eventsJSON)
646 for _, ev := range events {
647 relayproxy.WorkerPost(`["EVENT",` | jstr(sid) | `,` | ev.ToJSON() | `]`)
648 }
649 if info, ok := ps.proxySubs[sid]; ok {
650 info.localDone = true
651 proxyMaybeEOSE(sid)
652 }
653 })
654 }
655
656 remoteIDs := map[string]bool{}
657 base := "p_" | subID | "_"
658
659 seen := map[string]bool{}
660 deduped := []string{:0:len(relayURLs)}
661 for _, url := range relayURLs {
662 url = normalizeRelayURL(url)
663 if seen[url] {
664 continue
665 }
666 seen[url] = true
667 deduped = push(deduped, url)
668 }
669
670 ps.proxySubs[subID] = &proxySub{
671 remoteIDs: remoteIDs,
672 relayCount: len(deduped),
673 live: live,
674 }
675
676 for _, url := range deduped {
677 if !isAllowedRelay(url) {
678 continue
679 }
680 suffix := urlSuffix(url)
681 rSubID := base | suffix
682 remoteIDs[rSubID] = true
683 ps.remoteToProxy[rSubID] = subID
684 c := getConn(url)
685 c.Subscribe(rSubID, []*nostr.Filter{f})
686 }
687
688 proxyID := subID
689 if hasSearchField {
690 // Nothing to wait for locally.
691 if info, ok := ps.proxySubs[subID]; ok {
692 info.localDone = true
693 }
694 }
695 ps.proxySubs[subID].timer = relayproxy.WorkerSetTimeout(ps.eoseTimeoutMs, func() {
696 info, ok := ps.proxySubs[proxyID]
697 if !ok || info.done {
698 return
699 }
700 info.remoteDone = true
701 proxyMaybeEOSE(proxyID)
702 if !info.done {
703 // The timer has fired but the local store query has not answered.
704 // End the initial load anyway after a short grace so the consumer
705 // is never left waiting on it; 30s is too long for a feed.
706 relayproxy.WorkerSetTimeout(1000, func() {
707 cur, ok3 := ps.proxySubs[proxyID]
708 if !ok3 || cur.done {
709 return
710 }
711 cur.localDone = true
712 proxyMaybeEOSE(proxyID)
713 })
714 }
715 if !info.done || info.live {
716 return
717 }
718 // Linger: keep the remote subs open for ps.lingerMs after EOSE so
719 // events that arrive late from slow relays still get forwarded. The
720 // consumer already got its EOSE so it can proceed, but stragglers
721 // arriving during the linger window still match findProxySub and
722 // reach the app. Closing immediately on EOSE timer was the source of
723 // inconsistent profile resolution.
724 info.timer = relayproxy.WorkerSetTimeout(ps.lingerMs, func() {
725 cur, ok3 := ps.proxySubs[proxyID]
726 if !ok3 {
727 return
728 }
729 for rSubID := range cur.remoteIDs {
730 delete(ps.remoteToProxy, rSubID)
731 for _, c := range ps.rpool.AllConns() {
732 c.CloseSubscription(rSubID)
733 }
734 }
735 delete(ps.proxySubs, proxyID)
736 delete(ps.clientSubs, proxyID)
737 })
738 })
739 runtime.SovereignRestoreArena(prev)
740 }
741
742 // proxyMaybeEOSE emits the consumer's EOSE once the initial load is complete
743 // from both sources. The linger block that follows keeps the remote
744 // subscriptions open for a while so late stragglers still reach the app.
745 func proxyMaybeEOSE(proxyID string) {
746 info, ok := ps.proxySubs[proxyID]
747 if !ok || info.done {
748 return
749 }
750 if !info.localDone || !info.remoteDone {
751 return
752 }
753 info.done = true
754 if _, ok2 := ps.clientSubs[proxyID]; ok2 {
755 relayproxy.WorkerPost(`["EOSE",` | jstr(proxyID) | `]`)
756 }
757 }
758
759 func cleanupProxy(proxyID string) {
760 info, ok := ps.proxySubs[proxyID]
761 if !ok {
762 return
763 }
764 relayproxy.WorkerClearTimeout(info.timer)
765
766 if !info.done {
767 if _, ok2 := ps.clientSubs[proxyID]; ok2 {
768 relayproxy.WorkerPost(`["EOSE",` | jstr(proxyID) | `]`)
769 }
770 }
771
772 for rSubID := range info.remoteIDs {
773 delete(ps.remoteToProxy, rSubID)
774 for _, c := range ps.rpool.AllConns() {
775 c.CloseSubscription(rSubID)
776 }
777 }
778 delete(ps.proxySubs, proxyID)
779 }
780
781 func getConn(url string) (p *relay.Conn) {
782 c := ps.rpool.Connect(url)
783 if c.ScheduleReconnect == nil {
784 wireConn(c)
785 }
786 return c
787 }
788
789 func wireConn(c *relay.Conn) {
790 url := c.URL
791 c.ScheduleReconnect = func(delayMs int32, fn func()) {
792 relayproxy.WorkerSetTimeout(delayMs, fn)
793 }
794 c.SetOnEvent(func(rSubID string, ev *nostr.Event) {
795 if ps.mlsSubIDs[rSubID] {
796 evJSON := ev.ToJSON()
797 println("[relay-proxy] MLS_EVENT: sub=" | rSubID | " kind=" | helpers.Itoa(int64(ev.Kind)) | " from=" | ev.PubKey[:16] | "...")
798 relayproxy.WorkerPost(`["MLS_EVENT",` | evJSON | `]`)
799 return
800 }
801 if ps.mlsKPFetchID != "" && len(rSubID) >= len(ps.mlsKPFetchID) && rSubID[:len(ps.mlsKPFetchID)] == ps.mlsKPFetchID && ev.Kind == 443 {
802 evJSON := ev.ToJSON()
803 peer := ps.mlsKPFetchPeer
804 println("[relay-proxy] MLS_KP_RESULT: found KP for " | peer[:16] | "... on sub=" | rSubID)
805 closeKPFetchSubs()
806 ps.mlsKPFetchID = ""
807 relayproxy.WorkerPost(`["MLS_KP_RESULT",` | jstr(peer) | `,` | evJSON | `]`)
808 return
809 }
810 if ps.rlFetchID != "" && len(rSubID) >= len(ps.rlFetchID) && rSubID[:len(ps.rlFetchID)] == ps.rlFetchID && ev.Kind == 10002 {
811 urls := parseRelayList(ev)
812 peer := ps.rlFetchPeer
813 println("[relay-proxy] RL_RESULT: peer=" | peer[:16] | "... urls=" | helpers.Itoa(int64(len(urls))))
814 fn := ps.rlFetchCB
815 closeRLFetchSubs()
816 ps.rlFetchID = ""
817 ps.rlFetchCB = nil
818 if fn != nil {
819 fn(urls)
820 }
821 return
822 }
823 // Proxy event path: skip CheckSig here - the Verify Supervisor in
824 // wasm-host.mjs intercepts ["EVENT",...] and verifies before forwarding
825 // to the app worker. Unverified events never reach app business logic.
826 // A relay frame that cannot be routed to a client subscription is
827 // reported rather than dropped silently: this path losing frames was
828 // invisible because the only symptom was an empty feed.
829 proxyID, info := findProxySub(rSubID)
830 if info == nil {
831 relayproxy.WorkerPost(`["PROXY_DROP",` | jstr(rSubID) | `,"nomap"]`)
832 return
833 }
834 cs, ok := ps.clientSubs[proxyID]
835 if !ok {
836 relayproxy.WorkerPost(`["PROXY_DROP",` | jstr(rSubID) | `,"nosub"]`)
837 return
838 }
839 _ = cs
840 // No client-side filter re-check: the relay delivered this event for a
841 // subscription it holds the filter for, and re-applying a cached copy
842 // here dropped stored events whenever that cached filter was stale
843 // (the second page load in a session reused the client subID).
844 evJSON := ev.ToJSON()
845 relayproxy.WorkerPost(`["S_PUT_EVENT",` | evJSON | `]`)
846 relayproxy.WorkerPost(`["EVENT",` | jstr(proxyID) | `,` | evJSON | `]`)
847 relayproxy.WorkerPost(`["SEEN_ON",` | jstr(ev.ID) | `,` | jstr(url) | `]`)
848 })
849 c.SetOnEOSE(func(rSubID string) {
850 if ps.mlsSubIDs[rSubID] || (ps.mlsKPFetchID != "" && len(rSubID) >= len(ps.mlsKPFetchID) && rSubID[:len(ps.mlsKPFetchID)] == ps.mlsKPFetchID) {
851 println("[relay-proxy] MLS EOSE: sub=" | rSubID | " relay=" | url)
852 }
853 if ps.rlFetchID != "" && len(rSubID) >= len(ps.rlFetchID) && rSubID[:len(ps.rlFetchID)] == ps.rlFetchID {
854 ps.rlFetchEOSE++
855 println("[relay-proxy] RL EOSE: sub=" | rSubID | " relay=" | url | " eose=" | helpers.Itoa(int64(ps.rlFetchEOSE)) | "/" | helpers.Itoa(int64(len(ps.rlFetchSubs))))
856 if ps.rlFetchEOSE >= len(ps.rlFetchSubs) {
857 fn := ps.rlFetchCB
858 closeRLFetchSubs()
859 ps.rlFetchID = ""
860 ps.rlFetchCB = nil
861 if fn != nil {
862 fn(nil)
863 }
864 }
865 }
866 })
867 c.SetOnOK(func(eventID string, ok bool, msg string) {
868 okStr := "true"
869 if !ok {
870 okStr = "false"
871 println("[relay-proxy] OK REJECTED: id=" | eventID[:16] | "... ok=" | okStr | " msg=" | msg | " relay=" | url)
872 }
873 relayproxy.WorkerPost(`["OK",` | jstr(eventID) | `,` | okStr | `,` | jstr(msg) | `]`)
874 if ok {
875 relayproxy.WorkerPost(`["SEEN_ON",` | jstr(eventID) | `,` | jstr(url) | `]`)
876 }
877 })
878 }
879
880 func findProxySub(rSubID string) (s string, p *proxySub) {
881 proxyID, ok := ps.remoteToProxy[rSubID]
882 if !ok {
883 return "", nil
884 }
885 info, ok := ps.proxySubs[proxyID]
886 if !ok {
887 delete(ps.remoteToProxy, rSubID)
888 return "", nil
889 }
890 return proxyID, info
891 }
892
893 func normalizeRelayURL(url string) (s string) {
894 for len(url) > 0 && url[len(url)-1] == '/' {
895 url = url[:len(url)-1]
896 }
897 return url
898 }
899
900 func isAllowedRelay(url string) (ok bool) {
901 if len(url) >= 6 && url[:6] == "wss://" {
902 return true
903 }
904 if len(url) >= 16 && url[:16] == "ws://localhost:" {
905 return true
906 }
907 if len(url) >= 15 && url[:15] == "ws://127.0.0.1:" {
908 return true
909 }
910 return false
911 }
912
913 func urlSuffix(url string) (s string) {
914 // Extract host from wss://host/ or wss://host
915 start := 0
916 for i := 0; i < len(url); i++ {
917 if url[i] == '/' && i+1 < len(url) && url[i+1] == '/' {
918 start = i + 2
919 break
920 }
921 }
922 end := len(url)
923 for i := start; i < len(url); i++ {
924 if url[i] == '/' || url[i] == ':' {
925 end = i
926 break
927 }
928 }
929 host := url[start:end]
930 out := []byte{:0:len(host)}
931 for i := 0; i < len(host); i++ {
932 c := host[i]
933 if (c >= 'a' && c <= 'z') || (c >= '0' && c <= '9') {
934 out = push(out, c)
935 }
936 }
937 if len(out) > 12 {
938 out = out[len(out)-12:]
939 }
940 return string(out)
941 }
942
943 func jstr(s string) (sv string) { return helpers.JsonString(s) }
944
945 func dup(s string) (sv string) {
946 if len(s) == 0 {
947 return ""
948 }
949 b := []byte{:len(s)}
950 copy(b, s)
951 return string(b)
952 }
953
954 func dupStrs(ss []string) (ss2 []string) {
955 if len(ss) == 0 {
956 return nil
957 }
958 out := []string{:0:len(ss)}
959 for _, s := range ss {
960 out = push(out, dup(s))
961 }
962 return out
963 }
964