main.mx raw
1 package main
2
3 import (
4 "git.smesh.lol/moxie/pkg/mxutil"
5 "fmt"
6 "net/url"
7 "os"
8 "strconv"
9 "time"
10
11 "git.smesh.lol/morly/pkg/blossom"
12 "git.smesh.lol/morly/pkg/broadcast"
13 "git.smesh.lol/morly/pkg/mediaproxy"
14 "git.smesh.lol/nostr/pkg/envelope"
15 "git.smesh.lol/nostr/pkg/ws"
16 "git.smesh.lol/morly/pkg/relay/config"
17 "git.smesh.lol/morly/pkg/relay/server"
18 )
19
20 var version string
21
22 func initMainGlobals() {
23 version = "0.6.11"
24 crawlSeeds = []string{
25 "wss://relay.damus.io",
26 "wss://nos.lol",
27 }
28 bootstrapSeeds = []string{
29 "wss://purplepag.es",
30 "wss://relay.primal.net",
31 "wss://relay.damus.io",
32 "wss://nos.lol",
33 "wss://nostr.wine",
34 }
35 }
36
37 func main() {
38 initMainGlobals()
39
40 if len(os.Args) < 2 {
41 runRelay(os.Args[1:])
42 return
43 }
44 switch os.Args[1] {
45 case "relay":
46 runRelay(os.Args[2:])
47 case "sync":
48 runSync(os.Args[2:])
49 case "outbox":
50 runOutbox(os.Args[2:])
51 case "crawl":
52 runCrawl(os.Args[2:])
53 case "version", "-v", "--version":
54 os.Stdout.Write([]byte("musiquay " | version | "\n"))
55 case "help", "-h", "--help":
56 config.PrintHelp()
57 default:
58 if len(os.Args[1]) > 0 && os.Args[1][0] == '-' {
59 runRelay(os.Args[1:])
60 } else {
61 fmt.Fprintf(os.Stderr, "unknown command: %s\n", os.Args[1])
62 config.PrintHelp()
63 os.Exit(1)
64 }
65 }
66 }
67
68 func runRelay(_ [][]byte) {
69 // Register signal handlers before any domain is forked. A spawn child
70 // inherits the disposition: with the handler installed, a SIGTERM to the
71 // process group (which the test fixture sends) lets the child finish its
72 // loop and dump its coverage instead of taking the default action and
73 // dying with its private counter table.
74 server.InitSignals()
75
76 // PreSpawn forks the broadcast domain before config is loaded and before
77 // the first request touches the store, so relay mode does not pay a fork
78 // in the middle of serving.
79 broadcast.PreSpawn()
80 cfg := config.Load()
81 listenAddr := cfg.Addr()
82
83 // The database-engine domain owns the store: root sends it requests and
84 // never touches storage itself.
85 db := server.SpawnDB(cfg)
86 srv := server.New(db, cfg)
87 srv.Version = version
88
89 bsrv, err5 := blossom.New(cfg.BlossomDir)
90 if err5 != nil {
91 fmt.Fprintf(os.Stderr, "blossom: %v\n", err5)
92 os.Exit(1)
93 }
94
95 // Derive self-host prefix from ORLY_RELAY_URL for proxy deadlock avoidance.
96 selfHost := "git.smesh.lol/morly/"
97 if u, err4 := url.Parse(cfg.RelayURL); err4 == nil && len(u.Host) > 0 {
98 selfHost = u.Host | "/"
99 }
100
101 srv.Fallback = func(method, path string, headers map[string]string, body []byte) (int32, map[string]string, []byte) {
102 switch {
103 case hasPrefix(path, "/proxy/"):
104 rest := path[len("/proxy/"):]
105 if hasPrefix(rest, selfHost) {
106 direct := rest[len(selfHost)-1:]
107 return 302, map[string]string{"Location": direct, "Access-Control-Allow-Origin": "*"}, nil
108 }
109 return handleProxy(rest, selfHost)
110 case hasPrefix(path, "/blossom/"):
111 return bsrv.HandleRawWithUpstream(method, path[len("/blossom"):], headers, body, cfg.BlossomUpstream)
112 case path == "/__version":
113 stamp := int64(0)
114 if info1, err2 := os.Stat(cfg.StaticDir | "/app.wasm"); err2 == nil {
115 stamp = info1.ModTime().Unix()
116 } else if info2, err1 := os.Stat(cfg.StaticDir | "/_.mjs"); err1 == nil {
117 stamp = info2.ModTime().Unix()
118 }
119 return 200, map[string]string{
120 "Content-Type": "application/json",
121 "Access-Control-Allow-Origin": "*",
122 "Cache-Control": "no-store, no-cache, must-revalidate",
123 "Pragma": "no-cache",
124 "Expires": "0",
125 }, []byte(`{"v":"` | version | "+" | string(strconv.Itoa(int32(stamp))) | `"}`)
126 case path == "/.well-known/nostr.json":
127 return 200, map[string]string{
128 "Content-Type": "application/json",
129 "Access-Control-Allow-Origin": "*",
130 }, []byte(`{"names":{"mleku":"4c800257a588a82849d049817c2bdaad984b25a45ad9f6dad66e47d3b47e3b2f","bridge":"cf1ae33ad5f229dabd7d733ce37b0165126aebf581e4094df9373f77e00cb696"},"relays":{"4c800257a588a82849d049817c2bdaad984b25a45ad9f6dad66e47d3b47e3b2f":["wss://smesh.lol"],"cf1ae33ad5f229dabd7d733ce37b0165126aebf581e4094df9373f77e00cb696":["wss://smesh.lol","wss://relay.orly.dev"]}}`)
131 case isBlossomPath(path):
132 return bsrv.HandleRawWithUpstream(method, path, headers, body, cfg.BlossomUpstream)
133 }
134 return serveStatic(cfg.StaticDir, path)
135 }
136
137 srv.OnReady = func() {
138 fmt.Fprintln(os.Stderr, cfg.AppName | " " | version | " listening on " | listenAddr)
139 if cfg.CrawlerEnabled {
140 spawnCrawler(listenAddr)
141 }
142 if cfg.SyncPubkey != "" {
143 spawnOutbox(listenAddr, cfg.SyncPubkey)
144 }
145 if len(cfg.RelayPeers) > 0 {
146 for _, peer := range cfg.RelayPeers {
147 spawnSync(listenAddr, peer)
148 }
149 }
150 }
151 if err3 := srv.ListenAndServe(listenAddr); err3 != nil {
152 srv.Close()
153 fmt.Fprintf(os.Stderr, "listen: %v\n", err3)
154 os.Exit(1)
155 }
156 srv.Close()
157 }
158
159 func serveStatic(dir, path string) (n2 int32, m map[string]string, out []byte) {
160 if path == "" || path == "/" {
161 path = "/index.html"
162 }
163 if decoded, err2 := url.PathUnescape(path); err2 == nil {
164 path = decoded
165 }
166 data, err1 := os.ReadFile(dir | path)
167 if err1 != nil {
168 if hasFileExtension(path) {
169 return 404, map[string]string{"Content-Type": "text/plain"}, []byte("404 not found\n")
170 }
171 data, err1 = os.ReadFile(dir | "/index.html")
172 if err1 != nil {
173 return 404, map[string]string{"Content-Type": "text/plain"}, []byte("404 not found\n")
174 }
175 return 200, map[string]string{
176 "Content-Type": "text/html; charset=utf-8",
177 "Cache-Control": "no-cache",
178 "Cross-Origin-Opener-Policy": "same-origin",
179 "Cross-Origin-Embedder-Policy": "require-corp",
180 }, data
181 }
182 ct := "application/octet-stream"
183 switch {
184 case hasSuffix(path, ".html"):
185 ct = "text/html; charset=utf-8"
186 case hasSuffix(path, ".js"), hasSuffix(path, ".mjs"):
187 ct = "application/javascript"
188 case hasSuffix(path, ".css"):
189 ct = "text/css"
190 case hasSuffix(path, ".json"):
191 ct = "application/json"
192 case hasSuffix(path, ".svg"):
193 ct = "image/svg+xml"
194 case hasSuffix(path, ".png"):
195 ct = "image/png"
196 case hasSuffix(path, ".ico"):
197 ct = "image/x-icon"
198 case hasSuffix(path, ".wasm"):
199 ct = "application/wasm"
200 case hasSuffix(path, ".webp"):
201 ct = "image/webp"
202 case hasSuffix(path, ".woff2"):
203 ct = "font/woff2"
204 case hasSuffix(path, ".xpi"):
205 ct = "application/x-xpinstall"
206 }
207 h := map[string]string{
208 "Content-Type": ct,
209 "Cross-Origin-Opener-Policy": "same-origin",
210 "Cross-Origin-Embedder-Policy": "require-corp",
211 "Cross-Origin-Resource-Policy": "same-origin",
212 "Cache-Control": "no-cache",
213 }
214 if path == "/$sw/wasm-host-sw.mjs" {
215 h["Service-Worker-Allowed"] = "/"
216 }
217 if hasSuffix(path, ".xpi") {
218 h["Content-Disposition"] = "attachment; filename=\"musiquay-signer.xpi\""
219 }
220 return 200, h, data
221 }
222
223 // isBlossomPath matches /<64hex> or /<64hex>.<ext> with no further slashes.
224 func isBlossomPath(path string) (ok bool) {
225 if len(path) < 65 || path[0] != '/' {
226 return false
227 }
228 for i := 1; i < 65; i++ {
229 c := path[i]
230 if !((c >= '0' && c <= '9') || (c >= 'a' && c <= 'f') || (c >= 'A' && c <= 'F')) {
231 return false
232 }
233 }
234 if len(path) == 65 {
235 return true
236 }
237 if path[65] != '.' {
238 return false
239 }
240 for i := 66; i < len(path); i++ {
241 if path[i] == '/' {
242 return false
243 }
244 }
245 return true
246 }
247
248 // hasFileExtension returns true only for paths ending in a known static file
249 // extension. Generic TLD-like suffixes (.wine, .lol, .dev etc.) that appear
250 // in relay-URL-based SPA routes must NOT be treated as file paths.
251 func hasFileExtension(path string) (ok bool) {
252 known := []string{
253 ".html", ".htm",
254 ".js", ".mjs",
255 ".css",
256 ".wasm",
257 ".svg",
258 ".png", ".jpg", ".jpeg", ".webp", ".gif",
259 ".ico",
260 ".woff", ".woff2", ".ttf", ".otf",
261 ".json",
262 ".xpi",
263 ".map",
264 ".txt",
265 }
266 for _, ext := range known {
267 el := len(ext)
268 pl := len(path)
269 if pl > el && path[pl-el:] == ext {
270 return true
271 }
272 }
273 return false
274 }
275
276 // handleProxy fetches path (host/path format) via https and re-serves with
277 // COEP-compatible CORP headers. path arrives without the /proxy/ prefix and
278 // without a scheme; https:// is prepended before fetching.
279 //
280 // Self-referential paths (smesh.lol/...) are redirected directly so the
281 // relay does not deadlock trying to serve itself while blocked in Fetch.
282 func handleProxy(path, selfHost string) (n2 int32, m map[string]string, out []byte) {
283 if path == "" {
284 return 400, map[string]string{"Content-Type": "text/plain"}, []byte("missing path\n")
285 }
286 if hasPrefix(path, selfHost) {
287 direct := path[len(selfHost)-1:]
288 return 302, map[string]string{
289 "Location": direct,
290 "Access-Control-Allow-Origin": "*",
291 }, nil
292 }
293 target := "https://" | path
294 status, upstream, body, err := mediaproxy.Fetch(target, 32*1024*1024)
295 if err != nil {
296 return 502, map[string]string{"Content-Type": "text/plain"}, []byte("proxy: " | err.Error() | "\n")
297 }
298 if status < 200 || status >= 300 {
299 return status, map[string]string{"Content-Type": "text/plain"}, []byte(fmt.Sprintf("upstream %d\n", status))
300 }
301 ct := upstream["content-type"]
302 if !proxyAllowedCT(ct) {
303 return 415, map[string]string{"Content-Type": "text/plain"}, []byte("content-type not allowed: " | ct | "\n")
304 }
305 return 200, map[string]string{
306 "Content-Type": ct,
307 "Cross-Origin-Resource-Policy": "cross-origin",
308 "Cache-Control": "public, max-age=86400",
309 "Access-Control-Allow-Origin": "*",
310 }, body
311 }
312
313 func proxyAllowedCT(ct string) (ok bool) {
314 return hasPrefix(ct, "image/") || hasPrefix(ct, "video/") || ct == "application/octet-stream"
315 }
316
317 // --- sync command ---
318
319 // syncKinds is the set of event kinds synced from follows' write relays.
320 // Matches the feed view filter plus reactions and zaps for notifications.
321 const syncKindsJSON = `[1,6,7,1111,9735]`
322
323 // runSync handles "musiquay sync [--authors <hex,...>] [--since <unix-ts>] <remote-url> [local-url]".
324 // With no --authors it falls back to the old empty-filter behaviour.
325 func runSync(args []string) {
326 if len(args) < 1 {
327 fmt.Fprintln(os.Stderr, "usage: musiquay sync [--authors hex1,hex2,...] [--since unix-ts] <remote-url> [local-url]")
328 os.Exit(1)
329 }
330
331 // Parse flags.
332 var authors string
333 var sinceStr string
334 rest := args
335 for len(rest) >= 2 {
336 switch rest[0] {
337 case "--authors":
338 authors = rest[1]
339 rest = rest[2:]
340 case "--since":
341 sinceStr = rest[1]
342 rest = rest[2:]
343 default:
344 goto doneFlags
345 }
346 }
347 doneFlags:
348 if len(rest) < 1 {
349 fmt.Fprintln(os.Stderr, "sync: missing remote-url")
350 os.Exit(1)
351 }
352 remoteURL := rest[0]
353 localURL := "ws://127.0.0.1:3335"
354 if len(rest) >= 2 {
355 localURL = rest[1]
356 }
357
358 sinceTs := int64(0)
359 if sinceStr != "" {
360 for _, c := range sinceStr {
361 if c >= '0' && c <= '9' {
362 sinceTs = sinceTs*10 + int64(c-'0')
363 }
364 }
365 }
366
367 for {
368 latest := syncOnce(remoteURL, localURL, authors, sinceTs)
369 if latest > sinceTs {
370 sinceTs = latest
371 }
372 fmt.Fprintln(os.Stderr, "sync: disconnected, reconnecting in 30s...")
373 time.Sleep(30 * time.Second)
374 }
375 }
376
377 func syncOnce(remoteURL, localURL, authors string, sinceTs int64) (n int64) {
378 fmt.Fprintf(os.Stderr, "sync: connecting to remote %s\n", remoteURL)
379 remote, err7 := ws.Dial(remoteURL)
380 if err7 != nil {
381 fmt.Fprintf(os.Stderr, "sync: remote connect error: %v\n", err7)
382 return sinceTs
383 }
384 defer remote.Close()
385
386 local, err6 := ws.Dial(localURL)
387 if err6 != nil {
388 fmt.Fprintf(os.Stderr, "sync: local connect error: %v\n", err6)
389 return sinceTs
390 }
391 defer local.Close()
392
393 // Build filter JSON directly - simpler than constructing filter.F structs.
394 var filterJSON []byte
395 filterJSON = filterJSON | `{"kinds":`
396 filterJSON = filterJSON | syncKindsJSON
397 if authors != "" {
398 // authors is comma-separated hex pubkeys.
399 filterJSON = filterJSON | `,"authors":[`
400 first := true
401 start := 0
402 for i := 0; i <= len(authors); i++ {
403 if i == len(authors) || authors[i] == ',' {
404 pk := authors[start:i]
405 if len(pk) == 64 {
406 if !first {
407 filterJSON = filterJSON | ","
408 }
409 filterJSON = filterJSON | "\""
410 filterJSON = filterJSON | pk
411 filterJSON = filterJSON | "\""
412 first = false
413 }
414 start = i + 1
415 }
416 }
417 filterJSON = filterJSON | "]"
418 }
419 if sinceTs > 0 {
420 // since = checkpoint minus 1 hour to catch any late events.
421 since := sinceTs - 3600
422 filterJSON = filterJSON | `,"since":`
423 filterJSON = filterJSON | itoa64(since)
424 }
425 filterJSON = filterJSON | "}"
426
427 reqJSON := []byte(`["REQ","sync",`)
428 reqJSON = reqJSON | filterJSON
429 reqJSON = reqJSON | "]"
430
431 if authors != "" {
432 nAuthors := 1
433 for i := 0; i < len(authors); i++ {
434 if authors[i] == ',' {
435 nAuthors++
436 }
437 }
438 fmt.Fprintf(os.Stderr, "sync: subscribing kinds=%s authors=%d since=%d\n", syncKindsJSON, nAuthors, sinceTs)
439 } else {
440 fmt.Fprintf(os.Stderr, "sync: subscribing kinds=%s (all authors)\n", syncKindsJSON)
441 }
442
443 if err5 := remote.WriteText(reqJSON); err5 != nil {
444 fmt.Fprintf(os.Stderr, "sync: subscribe error: %v\n", err5)
445 return sinceTs
446 }
447
448 var forwarded int64
449 var latestTs int64
450 eosed := false
451
452 for {
453 op, payload, err4 := remote.ReadMessage()
454 if err4 != nil {
455 fmt.Fprintf(os.Stderr, "sync: read error (%d forwarded): %v\n", forwarded, err4)
456 return latestTs
457 }
458 if op == ws.OpClose {
459 fmt.Fprintf(os.Stderr, "sync: remote closed (%d forwarded)\n", forwarded)
460 return latestTs
461 }
462 if op != ws.OpText {
463 continue
464 }
465
466 label, rem, _ := envelope.Identify(payload)
467 switch label {
468 case envelope.EventLabel:
469 var es envelope.EventSubmission
470 if _, err3 := es.Unmarshal(rem); err3 != nil {
471 continue
472 }
473 if es.E == nil {
474 continue
475 }
476 fwd := &envelope.EventSubmission{E: es.E}
477 if err2 := local.WriteText(fwd.Marshal(nil)); err2 != nil {
478 fmt.Fprintf(os.Stderr, "sync: local publish error: %v\n", err2)
479 return latestTs
480 }
481 // Drain the OK response to keep the local socket buffer clear.
482 if _, _, err1 := local.ReadMessage(); err1 != nil {
483 fmt.Fprintf(os.Stderr, "sync: local read error: %v\n", err1)
484 return latestTs
485 }
486 forwarded++
487 if ts := int64(es.E.CreatedAt); ts > latestTs {
488 latestTs = ts
489 }
490 if forwarded%1000 == 0 {
491 fmt.Fprintf(os.Stderr, "sync: %d events forwarded (latest=%d)\n", forwarded, latestTs)
492 }
493 case envelope.EOSELabel:
494 if !eosed {
495 eosed = true
496 fmt.Fprintf(os.Stderr, "sync: EOSE - historical sync complete (%d forwarded). streaming live...\n", forwarded)
497 }
498 }
499 }
500 }
501
502 // itoa64 converts an int64 to its decimal string representation.
503 func itoa64(n int64) (s string) {
504 if n == 0 {
505 return "0"
506 }
507 neg := n < 0
508 if neg {
509 n = -n
510 }
511 buf := [20]byte{}
512 pos := 20
513 for n > 0 {
514 pos--
515 buf[pos] = byte('0' + n%10)
516 n /= 10
517 }
518 if neg {
519 pos--
520 buf[pos] = '-'
521 }
522 return string(buf[pos:])
523 }
524
525 // --- outbox sync ---
526
527 // runOutbox implements the outbox model: query the local store for a user's
528 // follows and their kind 10002 relay lists, then spawn one sync process per
529 // distinct write relay carrying only the authors who write there.
530 func runOutbox(args []string) {
531 if len(args) < 2 {
532 fmt.Fprintln(os.Stderr, "usage: musiquay outbox <pubkey-hex> <local-url>")
533 os.Exit(1)
534 }
535 pubkey := args[0]
536 localURL := args[1]
537
538 fmt.Fprintf(os.Stderr, "outbox: starting for pubkey %s...\n", pubkey[:8])
539
540 // Track running sync children: remoteURL → pid
541 children := map[string]int32{}
542
543 for {
544 relayAuthors := outboxDiscover(pubkey, localURL)
545 if len(relayAuthors) == 0 {
546 fmt.Fprintln(os.Stderr, "outbox: no write relays found yet, retrying in 2m")
547 time.Sleep(2 * time.Minute)
548 continue
549 }
550
551 // Kill syncs for relays no longer needed.
552 for url2, pid := range children {
553 if _, still := relayAuthors[url2]; !still {
554 fmt.Fprintf(os.Stderr, "outbox: relay removed, killing sync pid=%d url=%s\n", pid, url2)
555 proc, _ := os.FindProcess(pid)
556 if proc != nil {
557 proc.Kill()
558 }
559 delete(children, url2)
560 }
561 }
562
563 // Spawn syncs for new relays.
564 for url1, authors := range relayAuthors {
565 if _, running := children[url1]; running {
566 continue
567 }
568 authorsArg := ""
569 for i, pk := range authors {
570 if i > 0 {
571 authorsArg = authorsArg | ","
572 }
573 authorsArg = authorsArg | pk
574 }
575 cmd := os.Args[0] | " sync --authors " | authorsArg | " " | url1 | " " | localURL
576 argv := []string{"/bin/sh", "-c", cmd}
577 attr := &os.ProcAttr{}
578 proc, err := os.StartProcess("/bin/sh", argv, attr)
579 if err != nil {
580 fmt.Fprintf(os.Stderr, "outbox: spawn failed for %s: %v\n", url1, err)
581 continue
582 }
583 fmt.Fprintf(os.Stderr, "outbox: spawned sync pid=%d url=%s authors=%d\n", proc.Pid, url1, len(authors))
584 children[url1] = proc.Pid
585 }
586
587 // Re-discover every 30 minutes to pick up follow list changes.
588 time.Sleep(30 * time.Minute)
589 }
590 }
591
592 // outboxDiscover queries the local relay for the user's follows (kind 3) and
593 // their relay lists (kind 10002), returning a map of writeRelayURL → []pubkeyHex.
594 func outboxDiscover(pubkey, localURL string) (m map[string][]string) {
595 local, err7 := ws.Dial(localURL)
596 if err7 != nil {
597 fmt.Fprintf(os.Stderr, "outbox: local connect error: %v\n", err7)
598 return nil
599 }
600 defer local.Close()
601
602 // Step 1: fetch the user's kind 3 (follows list).
603 req3 := []byte(`["REQ","ob-k3",{"kinds":[3],"authors":["` | pubkey | `"],"limit":1}]`)
604 if err6 := local.WriteText(req3); err6 != nil {
605 fmt.Fprintf(os.Stderr, "outbox: REQ k3 error: %v\n", err6)
606 return nil
607 }
608
609 var follows []string
610 for {
611 op, payload, err4 := local.ReadMessage()
612 if err4 != nil || op == ws.OpClose {
613 break
614 }
615 if op != ws.OpText {
616 continue
617 }
618 label, rem, _ := envelope.Identify(payload)
619 if label == envelope.EOSELabel {
620 break
621 }
622 if label != envelope.EventLabel {
623 continue
624 }
625 var er envelope.EventResult
626 if _, err3 := er.Unmarshal(rem); err3 != nil || er.Event == nil {
627 continue
628 }
629 ev := er.Event
630 if ev.Kind != 3 || ev.Tags == nil {
631 continue
632 }
633 for _, t := range ev.Tags.GetAll([]byte("p")) {
634 if t.Len() >= 2 {
635 pk := string(t.ValueHex())
636 if len(pk) == 64 {
637 follows = mxutil.Ensure(follows, 1)
638 follows = push(follows, pk)
639 }
640 }
641 }
642 }
643 if len(follows) == 0 {
644 fmt.Fprintln(os.Stderr, "outbox: no follows in local store, bootstrapping from seed relays")
645 follows = outboxBootstrapFollows(pubkey, localURL)
646 if len(follows) == 0 {
647 return nil
648 }
649 }
650 // Always include the user's own pubkey.
651 follows = mxutil.Ensure(follows, 1)
652 follows = push(follows, pubkey)
653 fmt.Fprintf(os.Stderr, "outbox: found %d follows\n", len(follows))
654
655 // Step 2: fetch kind 10002 for all follows (batched).
656 // Build authors array.
657 authorsJSON := `["` | follows[0] | `"`
658 for _, pk := range follows[1:] {
659 authorsJSON = authorsJSON | `,"` | pk | `"`
660 }
661 authorsJSON = authorsJSON | `]`
662 req10002 := []byte(`["REQ","ob-rl",{"kinds":[10002],"authors":` | authorsJSON | `,"limit":` | itoa64(int64(len(follows)*2)) | `}]`)
663 if err5 := local.WriteText(req10002); err5 != nil {
664 fmt.Fprintf(os.Stderr, "outbox: REQ k10002 error: %v\n", err5)
665 return nil
666 }
667
668 // Map pubkey → write relay URLs.
669 writeRelays := map[string][]string{} // pubkey → []relayURL
670 for {
671 op, payload, err2 := local.ReadMessage()
672 if err2 != nil || op == ws.OpClose {
673 break
674 }
675 if op != ws.OpText {
676 continue
677 }
678 label, rem, _ := envelope.Identify(payload)
679 if label == envelope.EOSELabel {
680 break
681 }
682 if label != envelope.EventLabel {
683 continue
684 }
685 var er envelope.EventResult
686 if _, err1 := er.Unmarshal(rem); err1 != nil || er.Event == nil {
687 continue
688 }
689 ev := er.Event
690 if ev.Kind != 10002 || ev.Tags == nil {
691 continue
692 }
693 // Key by the hex spelling the follow list uses: a p-tag's
694 // ValueHex is hex, so a raw 32-byte event pubkey never matched a
695 // follow and every author fell back to the default relay.
696 pk := hexEnc(ev.Pubkey)
697 for _, t := range ev.Tags.GetAll([]byte("r")) {
698 if t.Len() < 2 {
699 continue
700 }
701 relayURL := string(t.Value())
702 if !hasPrefix(relayURL, "wss://") && !hasPrefix(relayURL, "ws://") {
703 continue
704 }
705 // marker is t.T[2] if present: "read", "write", or absent (both).
706 marker := ""
707 if t.Len() >= 3 {
708 marker = string(t.T[2])
709 }
710 isWrite := marker == "write" || marker == ""
711 if isWrite {
712 writeRelays[pk] = push(writeRelays[pk], relayURL)
713 }
714 }
715 }
716
717 // Invert: relayURL → []pubkeys that write there.
718 relayAuthors := map[string][]string{}
719 for _, pk := range follows {
720 urls, ok := writeRelays[pk]
721 if !ok {
722 // No kind 10002 found: fall back to the default outbox relay.
723 urls = []string{"wss://relay.damus.io"}
724 }
725 for _, u := range urls {
726 relayAuthors[u] = push(relayAuthors[u], pk)
727 }
728 }
729 fmt.Fprintf(os.Stderr, "outbox: %d write relays discovered\n", len(relayAuthors))
730 return relayAuthors
731 }
732
733 func spawnOutbox(localAddr, pubkey string) {
734 host := localAddr
735 if hasPrefix(host, "0.0.0.0:") {
736 host = "127.0.0.1:" | host[len("0.0.0.0:"):]
737 }
738 // Small delay so the relay's epoll loop is accepting before we connect.
739 cmd := "sleep 3 && " | os.Args[0] | " outbox " | pubkey | " ws://" | host
740 argv := []string{"/bin/sh", "-c", cmd}
741 attr := &os.ProcAttr{}
742 _, err := os.StartProcess("/bin/sh", argv, attr)
743 if err != nil {
744 fmt.Fprintf(os.Stderr, "outbox: spawn failed: %v\n", err)
745 }
746 }
747
748 // outboxBootstrapFollows fetches the user's kind 3 from seed relays and
749 // publishes it into the local store so outboxDiscover can find it.
750 // Returns the extracted follow pubkeys.
751 func outboxBootstrapFollows(pubkey, localURL string) (ss []string) {
752 for _, seed := range bootstrapSeeds {
753 fmt.Fprintf(os.Stderr, "outbox: bootstrap: trying %s\n", seed)
754 remote, err4 := ws.Dial(seed)
755 if err4 != nil {
756 continue
757 }
758 req := []byte(`["REQ","ob-boot",{"kinds":[3,10002],"authors":["` | pubkey | `"],"limit":2}]`)
759 if err3 := remote.WriteText(req); err3 != nil {
760 remote.Close()
761 continue
762 }
763 local, lerr := ws.Dial(localURL)
764 if lerr != nil {
765 remote.Close()
766 continue
767 }
768 var follows []string
769 evCount := 0
770 for {
771 op, payload, err2 := remote.ReadMessage()
772 if err2 != nil || op == ws.OpClose {
773 fmt.Fprintf(os.Stderr, "outbox: bootstrap: %s closed after %d events (err=%v)\n", seed, evCount, err2)
774 break
775 }
776 if op != ws.OpText {
777 continue
778 }
779 label, rem, _ := envelope.Identify(payload)
780 if label == envelope.EOSELabel {
781 fmt.Fprintf(os.Stderr, "outbox: bootstrap: %s EOSE after %d events\n", seed, evCount)
782 break
783 }
784 if label != envelope.EventLabel {
785 continue
786 }
787 var er envelope.EventResult
788 if _, err1 := er.Unmarshal(rem); err1 != nil || er.Event == nil {
789 continue
790 }
791 ev := er.Event
792 evCount++
793 fmt.Fprintf(os.Stderr, "outbox: bootstrap: %s event kind=%d\n", seed, ev.Kind)
794 // Publish to local store.
795 fwd := &envelope.EventSubmission{E: ev}
796 local.WriteText(fwd.Marshal(nil))
797 local.ReadMessage() // drain OK
798 // Extract follows from kind 3. Use ValueHex() because "p" tag values
799 // may be binary-encoded (32 bytes) rather than hex strings (64 chars).
800 if ev.Kind == 3 && ev.Tags != nil {
801 for _, t := range ev.Tags.GetAll([]byte("p")) {
802 if t.Len() >= 2 {
803 pk := string(t.ValueHex())
804 if len(pk) == 64 {
805 follows = mxutil.Ensure(follows, 1)
806 follows = push(follows, pk)
807 }
808 }
809 }
810 }
811 }
812 remote.Close()
813 local.Close()
814 if len(follows) > 0 {
815 fmt.Fprintf(os.Stderr, "outbox: bootstrap: found %d follows from %s\n", len(follows), seed)
816 // Second pass: fetch kind 10002 for all follows so outboxDiscover
817 // can map them to write relays without another bootstrap cycle.
818 outboxBootstrapRelayLists(follows, localURL, seed)
819 return follows
820 }
821 }
822 fmt.Fprintln(os.Stderr, "outbox: bootstrap: no follows found in seed relays")
823 return nil
824 }
825
826 // outboxBootstrapRelayLists fetches kind 10002 for all follows from one relay
827 // and publishes them to local so outboxDiscover can build the write relay map.
828 func outboxBootstrapRelayLists(follows []string, localURL, seedURL string) {
829 if len(follows) == 0 {
830 return
831 }
832 remote, err4 := ws.Dial(seedURL)
833 if err4 != nil {
834 return
835 }
836 defer remote.Close()
837 local, lerr := ws.Dial(localURL)
838 if lerr != nil {
839 return
840 }
841 defer local.Close()
842
843 // Build authors JSON for all follows.
844 authorsJSON := `["` | follows[0] | `"`
845 for _, pk := range follows[1:] {
846 authorsJSON = authorsJSON | `,"` | pk | `"`
847 }
848 authorsJSON = authorsJSON | `]`
849 lim := itoa64(int64(len(follows) * 2))
850 req := []byte(`["REQ","ob-rl2",{"kinds":[10002],"authors":` | authorsJSON | `,"limit":` | lim | `}]`)
851 if err3 := remote.WriteText(req); err3 != nil {
852 return
853 }
854 count := 0
855 for {
856 op, payload, err2 := remote.ReadMessage()
857 if err2 != nil || op == ws.OpClose {
858 break
859 }
860 if op != ws.OpText {
861 continue
862 }
863 label, rem, _ := envelope.Identify(payload)
864 if label == envelope.EOSELabel {
865 break
866 }
867 if label != envelope.EventLabel {
868 continue
869 }
870 var er envelope.EventResult
871 if _, err1 := er.Unmarshal(rem); err1 != nil || er.Event == nil {
872 continue
873 }
874 fwd := &envelope.EventSubmission{E: er.Event}
875 local.WriteText(fwd.Marshal(nil))
876 local.ReadMessage() // drain OK
877 count++
878 }
879 fmt.Fprintf(os.Stderr, "outbox: bootstrap: stored %d kind 10002 events from %s\n", count, seedURL)
880 }
881
882 func hasPrefix(s, prefix string) (ok bool) {
883 return len(s) >= len(prefix) && s[:len(prefix)] == prefix
884 }
885
886 func hasSuffix(s, suffix string) (ok bool) {
887 return len(s) >= len(suffix) && s[len(s)-len(suffix):] == suffix
888 }
889
890 func spawnCrawler(listenAddr string) {
891 host := listenAddr
892 if hasPrefix(host, "0.0.0.0:") {
893 host = "127.0.0.1:" | host[len("0.0.0.0:"):]
894 }
895 cmd := os.Args[0] | " crawl ws://" | host
896 argv := []string{"/bin/sh", "-c", cmd}
897 attr := &os.ProcAttr{}
898 _, err := os.StartProcess("/bin/sh", argv, attr)
899 if err != nil {
900 fmt.Fprintf(os.Stderr, "crawl: spawn failed: %v\n", err)
901 }
902 }
903
904 func spawnSync(listenAddr, remoteURL string) {
905 host := listenAddr
906 if hasPrefix(host, "0.0.0.0:") {
907 host = "127.0.0.1:" | host[len("0.0.0.0:"):]
908 }
909 cmd := os.Args[0] | " sync " | remoteURL | " ws://" | host
910 argv := []string{"/bin/sh", "-c", cmd}
911 attr := &os.ProcAttr{}
912 _, err := os.StartProcess("/bin/sh", argv, attr)
913 if err != nil {
914 fmt.Fprintf(os.Stderr, "sync: spawn failed: %v\n", err)
915 }
916 }
917
918 // --- crawl command ---
919
920 var crawlSeeds []string
921
922 // bootstrapSeeds are tried specifically for kind 3/10002 lookups.
923 // purplepag.es is designed for profile and contact list events.
924 var bootstrapSeeds []string
925
926 // Directory event kinds to fetch.
927 const crawlKindsFilter = `[0,3,5,1984,10000,10002,10050]`
928
929 func clog(out *os.File, format string, args ...fmt.Stringer) {
930 ts := time.Now().Format("15:04:05")
931 fmt.Fprintf(out, ts|" "|format|"\n", args)
932 }
933
934 // relayDB tracks known relays with frequency scores.
935 // Higher score = seen more often in kind 10002/10050 events = higher priority.
936 type relayDB struct {
937 score map[string]int32
938 order []string // URLs sorted by descending score
939 }
940
941 func newRelayDB() (p *relayDB) {
942 return &relayDB{score: map[string]int32{}}
943 }
944
945 func (db *relayDB) add(relayURL string, weight int32) {
946 db.score[relayURL] += weight
947 }
948
949 // sorted returns relay URLs ordered by descending frequency.
950 func (db *relayDB) sorted() (ss []string) {
951 urls := []string{:0:len(db.score)}
952 for u := range db.score {
953 urls = mxutil.Ensure(urls, 1)
954 urls = push(urls, u)
955 }
956 // Simple insertion sort by score descending.
957 for i := 1; i < len(urls); i++ {
958 for j := i; j > 0 && db.score[urls[j]] > db.score[urls[j-1]]; j-- {
959 urls[j], urls[j-1] = urls[j-1], urls[j]
960 }
961 }
962 return urls
963 }
964
965 func runCrawl(args []string) {
966 var err error
967 var out *os.File
968 out, err = os.OpenFile("/tmp/musiquay-crawl.log",
969 os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0644)
970 if err != nil {
971 out = os.Stderr
972 }
973
974 localURL := "ws://127.0.0.1:3334"
975 if len(args) >= 1 {
976 localURL = args[0]
977 }
978 clog(out, "started pid=%d local=%s", os.Getpid(), localURL)
979
980 db := newRelayDB()
981 // Seed relays get a high initial score.
982 for _, s := range crawlSeeds {
983 db.add(s, 100)
984 }
985
986 pass := 0
987 for {
988 pass++
989 clog(out, "=== pass %d, %d relays known ===", pass, len(db.score))
990 ok := crawlPass(localURL, db, out)
991 if ok {
992 clog(out, "pass complete, sleeping 5m")
993 time.Sleep(5 * time.Minute)
994 } else {
995 clog(out, "pass failed, retrying in 30s")
996 time.Sleep(30 * time.Second)
997 }
998 }
999 }
1000
1001 func crawlPass(localURL string, db *relayDB, out *os.File) (ok bool) {
1002 relays := db.sorted()
1003 if len(relays) == 0 {
1004 clog(out, "no relays known")
1005 return false
1006 }
1007
1008 totalEvents := 0
1009 for i, relayURL := range relays {
1010 clog(out, "[%d/%d] crawling %s (score %d)", i+1, len(relays), relayURL, db.score[relayURL])
1011
1012 events := crawlRelay(relayURL, out)
1013 if len(events) == 0 {
1014 clog(out, " %s → 0 events", relayURL)
1015 time.Sleep(1 * time.Second)
1016 continue
1017 }
1018 clog(out, " %s → %d events", relayURL, len(events))
1019
1020 // Extract new relay URLs from the events before publishing.
1021 for _, raw := range events {
1022 crawlExtractRelays(raw, db)
1023 }
1024
1025 // Publish batch to local relay.
1026 published := crawlPublishBatch(localURL, events, out)
1027 clog(out, " published %d/%d to local", published, len(events))
1028 totalEvents += published
1029
1030 time.Sleep(1 * time.Second)
1031 }
1032
1033 clog(out, "total %d events from %d relays", totalEvents, len(relays))
1034 return true
1035 }
1036
1037 // crawlRelay connects to one relay and subscribes to directory events.
1038 // Returns raw EVENT messages suitable for republishing.
1039 func crawlRelay(relayURL string, out *os.File) (ss [][]byte) {
1040 remote, err4 := ws.Dial(relayURL)
1041 if err4 != nil {
1042 clog(out, " dial %s FAILED: %v", relayURL, err4)
1043 return nil
1044 }
1045 defer remote.Close()
1046
1047 reqJSON := []byte(`["REQ","cr",{"kinds":` | crawlKindsFilter | `,"limit":200}]`)
1048 if err3 := remote.WriteText(reqJSON); err3 != nil {
1049 clog(out, " write REQ to %s failed: %v", relayURL, err3)
1050 return nil
1051 }
1052
1053 var events [][]byte
1054 for {
1055 op, payload, err2 := remote.ReadMessage()
1056 if err2 != nil {
1057 break
1058 }
1059 if op != ws.OpText {
1060 continue
1061 }
1062 label, rem, _ := envelope.Identify(payload)
1063 if label == envelope.EOSELabel {
1064 break
1065 }
1066 if label == envelope.EventLabel {
1067 var er envelope.EventResult
1068 if _, err1 := er.Unmarshal(rem); err1 == nil && er.Event != nil {
1069 es := &envelope.EventSubmission{E: er.Event}
1070 events = mxutil.Ensure(events, 1)
1071 events = push(events, es.Marshal(nil))
1072 }
1073 }
1074 _ = rem
1075 }
1076 return events
1077 }
1078
1079 // crawlExtractRelays parses an EVENT submission and adds discovered relay
1080 // URLs to the database with a frequency bump.
1081 func crawlExtractRelays(raw []byte, db *relayDB) {
1082 _, rem, err2 := envelope.Identify(raw)
1083 if err2 != nil {
1084 return
1085 }
1086 var es envelope.EventSubmission
1087 if _, err1 := es.Unmarshal(rem); err1 != nil || es.E == nil {
1088 return
1089 }
1090 ev := es.E
1091 if (ev.Kind != 10002 && ev.Kind != 10050) || ev.Tags == nil {
1092 return
1093 }
1094 for _, t := range ev.Tags.GetAll([]byte("r")) {
1095 if t.Len() >= 2 {
1096 relayURL := string(t.Value())
1097 if len(relayURL) > 5 && (hasPrefix(relayURL, "wss://") || hasPrefix(relayURL, "ws://")) {
1098 db.add(relayURL, 1)
1099 }
1100 }
1101 }
1102 }
1103
1104 // crawlPublishBatch publishes a batch of EVENT messages to the local relay.
1105 func crawlPublishBatch(localURL string, events [][]byte, out *os.File) (n int32) {
1106 local, err2 := ws.Dial(localURL)
1107 if err2 != nil {
1108 clog(out, " local connect failed: %v", err2)
1109 return 0
1110 }
1111 defer local.Close()
1112
1113 for _, evBytes := range events {
1114 local.WriteText(evBytes)
1115 }
1116
1117 // Drain OKs - one per event sent.
1118 count := 0
1119 for count < len(events) {
1120 _, _, err1 := local.ReadMessage()
1121 if err1 != nil {
1122 break
1123 }
1124 count++
1125 }
1126 return count
1127 }
1128
1129 func hexEnc(b []byte) (s string) {
1130 const hx = "0123456789abcdef"
1131 out := []byte{:len(b)*2}
1132 for i, v := range b {
1133 out[i*2] = hx[v>>4]
1134 out[i*2+1] = hx[v&0x0f]
1135 }
1136 return string(out)
1137 }
1138
1139 func i64str(n int64) (s string) {
1140 if n == 0 {
1141 return "0"
1142 }
1143 neg := false
1144 if n < 0 {
1145 neg = true
1146 n = -n
1147 }
1148 var buf [20]byte
1149 i := 19
1150 for n > 0 {
1151 buf[i] = byte('0' + n%10)
1152 i--
1153 n /= 10
1154 }
1155 if neg {
1156 buf[i] = '-'
1157 i--
1158 }
1159 return string(buf[i+1:])
1160 }
1161
1162 func crawlAppendUniq(ss []string, s string) (ss2 []string) {
1163 for _, x := range ss {
1164 if x == s {
1165 return ss
1166 }
1167 }
1168 return push(ss, s)
1169 }
1170
1171