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