package main import ( "git.smesh.lol/moxie/pkg/mxutil" "fmt" "net/url" "os" "strconv" "time" "git.smesh.lol/morly/pkg/blossom" "git.smesh.lol/morly/pkg/broadcast" "git.smesh.lol/morly/pkg/mediaproxy" "git.smesh.lol/nostr/pkg/envelope" "git.smesh.lol/nostr/pkg/ws" "git.smesh.lol/morly/pkg/relay/config" "git.smesh.lol/morly/pkg/relay/server" ) var version string func initMainGlobals() { version = "0.6.11" crawlSeeds = []string{ "wss://relay.damus.io", "wss://nos.lol", } bootstrapSeeds = []string{ "wss://purplepag.es", "wss://relay.primal.net", "wss://relay.damus.io", "wss://nos.lol", "wss://nostr.wine", } } func main() { initMainGlobals() if len(os.Args) < 2 { runRelay(os.Args[1:]) return } switch os.Args[1] { case "relay": runRelay(os.Args[2:]) case "sync": runSync(os.Args[2:]) case "outbox": runOutbox(os.Args[2:]) case "crawl": runCrawl(os.Args[2:]) case "version", "-v", "--version": os.Stdout.Write([]byte("musiquay " | version | "\n")) case "help", "-h", "--help": config.PrintHelp() default: if len(os.Args[1]) > 0 && os.Args[1][0] == '-' { runRelay(os.Args[1:]) } else { fmt.Fprintf(os.Stderr, "unknown command: %s\n", os.Args[1]) config.PrintHelp() os.Exit(1) } } } func runRelay(_ [][]byte) { // Register signal handlers before any domain is forked. A spawn child // inherits the disposition: with the handler installed, a SIGTERM to the // process group (which the test fixture sends) lets the child finish its // loop and dump its coverage instead of taking the default action and // dying with its private counter table. server.InitSignals() // PreSpawn forks the broadcast domain before config is loaded and before // the first request touches the store, so relay mode does not pay a fork // in the middle of serving. broadcast.PreSpawn() cfg := config.Load() listenAddr := cfg.Addr() // The database-engine domain owns the store: root sends it requests and // never touches storage itself. db := server.SpawnDB(cfg) srv := server.New(db, cfg) srv.Version = version bsrv, err5 := blossom.New(cfg.BlossomDir) if err5 != nil { fmt.Fprintf(os.Stderr, "blossom: %v\n", err5) os.Exit(1) } // Derive self-host prefix from ORLY_RELAY_URL for proxy deadlock avoidance. selfHost := "git.smesh.lol/morly/" if u, err4 := url.Parse(cfg.RelayURL); err4 == nil && len(u.Host) > 0 { selfHost = u.Host | "/" } srv.Fallback = func(method, path string, headers map[string]string, body []byte) (int32, map[string]string, []byte) { switch { case hasPrefix(path, "/proxy/"): rest := path[len("/proxy/"):] if hasPrefix(rest, selfHost) { direct := rest[len(selfHost)-1:] return 302, map[string]string{"Location": direct, "Access-Control-Allow-Origin": "*"}, nil } return handleProxy(rest, selfHost) case hasPrefix(path, "/blossom/"): return bsrv.HandleRawWithUpstream(method, path[len("/blossom"):], headers, body, cfg.BlossomUpstream) case path == "/__version": stamp := int64(0) if info1, err2 := os.Stat(cfg.StaticDir | "/app.wasm"); err2 == nil { stamp = info1.ModTime().Unix() } else if info2, err1 := os.Stat(cfg.StaticDir | "/_.mjs"); err1 == nil { stamp = info2.ModTime().Unix() } return 200, map[string]string{ "Content-Type": "application/json", "Access-Control-Allow-Origin": "*", "Cache-Control": "no-store, no-cache, must-revalidate", "Pragma": "no-cache", "Expires": "0", }, []byte(`{"v":"` | version | "+" | string(strconv.Itoa(int32(stamp))) | `"}`) case path == "/.well-known/nostr.json": return 200, map[string]string{ "Content-Type": "application/json", "Access-Control-Allow-Origin": "*", }, []byte(`{"names":{"mleku":"4c800257a588a82849d049817c2bdaad984b25a45ad9f6dad66e47d3b47e3b2f","bridge":"cf1ae33ad5f229dabd7d733ce37b0165126aebf581e4094df9373f77e00cb696"},"relays":{"4c800257a588a82849d049817c2bdaad984b25a45ad9f6dad66e47d3b47e3b2f":["wss://smesh.lol"],"cf1ae33ad5f229dabd7d733ce37b0165126aebf581e4094df9373f77e00cb696":["wss://smesh.lol","wss://relay.orly.dev"]}}`) case isBlossomPath(path): return bsrv.HandleRawWithUpstream(method, path, headers, body, cfg.BlossomUpstream) } return serveStatic(cfg.StaticDir, path) } srv.OnReady = func() { fmt.Fprintln(os.Stderr, cfg.AppName | " " | version | " listening on " | listenAddr) if cfg.CrawlerEnabled { spawnCrawler(listenAddr) } if cfg.SyncPubkey != "" { spawnOutbox(listenAddr, cfg.SyncPubkey) } if len(cfg.RelayPeers) > 0 { for _, peer := range cfg.RelayPeers { spawnSync(listenAddr, peer) } } } if err3 := srv.ListenAndServe(listenAddr); err3 != nil { srv.Close() fmt.Fprintf(os.Stderr, "listen: %v\n", err3) os.Exit(1) } srv.Close() } func serveStatic(dir, path string) (n2 int32, m map[string]string, out []byte) { if path == "" || path == "/" { path = "/index.html" } if decoded, err2 := url.PathUnescape(path); err2 == nil { path = decoded } data, err1 := os.ReadFile(dir | path) if err1 != nil { if hasFileExtension(path) { return 404, map[string]string{"Content-Type": "text/plain"}, []byte("404 not found\n") } data, err1 = os.ReadFile(dir | "/index.html") if err1 != nil { return 404, map[string]string{"Content-Type": "text/plain"}, []byte("404 not found\n") } return 200, map[string]string{ "Content-Type": "text/html; charset=utf-8", "Cache-Control": "no-cache", "Cross-Origin-Opener-Policy": "same-origin", "Cross-Origin-Embedder-Policy": "require-corp", }, data } ct := "application/octet-stream" switch { case hasSuffix(path, ".html"): ct = "text/html; charset=utf-8" case hasSuffix(path, ".js"), hasSuffix(path, ".mjs"): ct = "application/javascript" case hasSuffix(path, ".css"): ct = "text/css" case hasSuffix(path, ".json"): ct = "application/json" case hasSuffix(path, ".svg"): ct = "image/svg+xml" case hasSuffix(path, ".png"): ct = "image/png" case hasSuffix(path, ".ico"): ct = "image/x-icon" case hasSuffix(path, ".wasm"): ct = "application/wasm" case hasSuffix(path, ".webp"): ct = "image/webp" case hasSuffix(path, ".woff2"): ct = "font/woff2" case hasSuffix(path, ".xpi"): ct = "application/x-xpinstall" } h := map[string]string{ "Content-Type": ct, "Cross-Origin-Opener-Policy": "same-origin", "Cross-Origin-Embedder-Policy": "require-corp", "Cross-Origin-Resource-Policy": "same-origin", "Cache-Control": "no-cache", } if path == "/$sw/wasm-host-sw.mjs" { h["Service-Worker-Allowed"] = "/" } if hasSuffix(path, ".xpi") { h["Content-Disposition"] = "attachment; filename=\"musiquay-signer.xpi\"" } return 200, h, data } // isBlossomPath matches /<64hex> or /<64hex>. with no further slashes. func isBlossomPath(path string) (ok bool) { if len(path) < 65 || path[0] != '/' { return false } for i := 1; i < 65; i++ { c := path[i] if !((c >= '0' && c <= '9') || (c >= 'a' && c <= 'f') || (c >= 'A' && c <= 'F')) { return false } } if len(path) == 65 { return true } if path[65] != '.' { return false } for i := 66; i < len(path); i++ { if path[i] == '/' { return false } } return true } // hasFileExtension returns true only for paths ending in a known static file // extension. Generic TLD-like suffixes (.wine, .lol, .dev etc.) that appear // in relay-URL-based SPA routes must NOT be treated as file paths. func hasFileExtension(path string) (ok bool) { known := []string{ ".html", ".htm", ".js", ".mjs", ".css", ".wasm", ".svg", ".png", ".jpg", ".jpeg", ".webp", ".gif", ".ico", ".woff", ".woff2", ".ttf", ".otf", ".json", ".xpi", ".map", ".txt", } for _, ext := range known { el := len(ext) pl := len(path) if pl > el && path[pl-el:] == ext { return true } } return false } // handleProxy fetches path (host/path format) via https and re-serves with // COEP-compatible CORP headers. path arrives without the /proxy/ prefix and // without a scheme; https:// is prepended before fetching. // // Self-referential paths (smesh.lol/...) are redirected directly so the // relay does not deadlock trying to serve itself while blocked in Fetch. func handleProxy(path, selfHost string) (n2 int32, m map[string]string, out []byte) { if path == "" { return 400, map[string]string{"Content-Type": "text/plain"}, []byte("missing path\n") } if hasPrefix(path, selfHost) { direct := path[len(selfHost)-1:] return 302, map[string]string{ "Location": direct, "Access-Control-Allow-Origin": "*", }, nil } target := "https://" | path status, upstream, body, err := mediaproxy.Fetch(target, 32*1024*1024) if err != nil { return 502, map[string]string{"Content-Type": "text/plain"}, []byte("proxy: " | err.Error() | "\n") } if status < 200 || status >= 300 { return status, map[string]string{"Content-Type": "text/plain"}, []byte(fmt.Sprintf("upstream %d\n", status)) } ct := upstream["content-type"] if !proxyAllowedCT(ct) { return 415, map[string]string{"Content-Type": "text/plain"}, []byte("content-type not allowed: " | ct | "\n") } return 200, map[string]string{ "Content-Type": ct, "Cross-Origin-Resource-Policy": "cross-origin", "Cache-Control": "public, max-age=86400", "Access-Control-Allow-Origin": "*", }, body } func proxyAllowedCT(ct string) (ok bool) { return hasPrefix(ct, "image/") || hasPrefix(ct, "video/") || ct == "application/octet-stream" } // --- sync command --- // syncKinds is the set of event kinds synced from follows' write relays. // Matches the feed view filter plus reactions and zaps for notifications. const syncKindsJSON = `[1,6,7,1111,9735]` // runSync handles "musiquay sync [--authors ] [--since ] [local-url]". // With no --authors it falls back to the old empty-filter behaviour. func runSync(args []string) { if len(args) < 1 { fmt.Fprintln(os.Stderr, "usage: musiquay sync [--authors hex1,hex2,...] [--since unix-ts] [local-url]") os.Exit(1) } // Parse flags. var authors string var sinceStr string rest := args for len(rest) >= 2 { switch rest[0] { case "--authors": authors = rest[1] rest = rest[2:] case "--since": sinceStr = rest[1] rest = rest[2:] default: goto doneFlags } } doneFlags: if len(rest) < 1 { fmt.Fprintln(os.Stderr, "sync: missing remote-url") os.Exit(1) } remoteURL := rest[0] localURL := "ws://127.0.0.1:3335" if len(rest) >= 2 { localURL = rest[1] } sinceTs := int64(0) if sinceStr != "" { for _, c := range sinceStr { if c >= '0' && c <= '9' { sinceTs = sinceTs*10 + int64(c-'0') } } } for { latest := syncOnce(remoteURL, localURL, authors, sinceTs) if latest > sinceTs { sinceTs = latest } fmt.Fprintln(os.Stderr, "sync: disconnected, reconnecting in 30s...") time.Sleep(30 * time.Second) } } func syncOnce(remoteURL, localURL, authors string, sinceTs int64) (n int64) { fmt.Fprintf(os.Stderr, "sync: connecting to remote %s\n", remoteURL) remote, err7 := ws.Dial(remoteURL) if err7 != nil { fmt.Fprintf(os.Stderr, "sync: remote connect error: %v\n", err7) return sinceTs } defer remote.Close() local, err6 := ws.Dial(localURL) if err6 != nil { fmt.Fprintf(os.Stderr, "sync: local connect error: %v\n", err6) return sinceTs } defer local.Close() // Build filter JSON directly - simpler than constructing filter.F structs. var filterJSON []byte filterJSON = filterJSON | `{"kinds":` filterJSON = filterJSON | syncKindsJSON if authors != "" { // authors is comma-separated hex pubkeys. filterJSON = filterJSON | `,"authors":[` first := true start := 0 for i := 0; i <= len(authors); i++ { if i == len(authors) || authors[i] == ',' { pk := authors[start:i] if len(pk) == 64 { if !first { filterJSON = filterJSON | "," } filterJSON = filterJSON | "\"" filterJSON = filterJSON | pk filterJSON = filterJSON | "\"" first = false } start = i + 1 } } filterJSON = filterJSON | "]" } if sinceTs > 0 { // since = checkpoint minus 1 hour to catch any late events. since := sinceTs - 3600 filterJSON = filterJSON | `,"since":` filterJSON = filterJSON | itoa64(since) } filterJSON = filterJSON | "}" reqJSON := []byte(`["REQ","sync",`) reqJSON = reqJSON | filterJSON reqJSON = reqJSON | "]" if authors != "" { nAuthors := 1 for i := 0; i < len(authors); i++ { if authors[i] == ',' { nAuthors++ } } fmt.Fprintf(os.Stderr, "sync: subscribing kinds=%s authors=%d since=%d\n", syncKindsJSON, nAuthors, sinceTs) } else { fmt.Fprintf(os.Stderr, "sync: subscribing kinds=%s (all authors)\n", syncKindsJSON) } if err5 := remote.WriteText(reqJSON); err5 != nil { fmt.Fprintf(os.Stderr, "sync: subscribe error: %v\n", err5) return sinceTs } var forwarded int64 var latestTs int64 eosed := false for { op, payload, err4 := remote.ReadMessage() if err4 != nil { fmt.Fprintf(os.Stderr, "sync: read error (%d forwarded): %v\n", forwarded, err4) return latestTs } if op == ws.OpClose { fmt.Fprintf(os.Stderr, "sync: remote closed (%d forwarded)\n", forwarded) return latestTs } if op != ws.OpText { continue } label, rem, _ := envelope.Identify(payload) switch label { case envelope.EventLabel: var es envelope.EventSubmission if _, err3 := es.Unmarshal(rem); err3 != nil { continue } if es.E == nil { continue } fwd := &envelope.EventSubmission{E: es.E} if err2 := local.WriteText(fwd.Marshal(nil)); err2 != nil { fmt.Fprintf(os.Stderr, "sync: local publish error: %v\n", err2) return latestTs } // Drain the OK response to keep the local socket buffer clear. if _, _, err1 := local.ReadMessage(); err1 != nil { fmt.Fprintf(os.Stderr, "sync: local read error: %v\n", err1) return latestTs } forwarded++ if ts := int64(es.E.CreatedAt); ts > latestTs { latestTs = ts } if forwarded%1000 == 0 { fmt.Fprintf(os.Stderr, "sync: %d events forwarded (latest=%d)\n", forwarded, latestTs) } case envelope.EOSELabel: if !eosed { eosed = true fmt.Fprintf(os.Stderr, "sync: EOSE - historical sync complete (%d forwarded). streaming live...\n", forwarded) } } } } // itoa64 converts an int64 to its decimal string representation. func itoa64(n int64) (s string) { if n == 0 { return "0" } neg := n < 0 if neg { n = -n } buf := [20]byte{} pos := 20 for n > 0 { pos-- buf[pos] = byte('0' + n%10) n /= 10 } if neg { pos-- buf[pos] = '-' } return string(buf[pos:]) } // --- outbox sync --- // runOutbox implements the outbox model: query the local store for a user's // follows and their kind 10002 relay lists, then spawn one sync process per // distinct write relay carrying only the authors who write there. func runOutbox(args []string) { if len(args) < 2 { fmt.Fprintln(os.Stderr, "usage: musiquay outbox ") os.Exit(1) } pubkey := args[0] localURL := args[1] fmt.Fprintf(os.Stderr, "outbox: starting for pubkey %s...\n", pubkey[:8]) // Track running sync children: remoteURL → pid children := map[string]int32{} for { relayAuthors := outboxDiscover(pubkey, localURL) if len(relayAuthors) == 0 { fmt.Fprintln(os.Stderr, "outbox: no write relays found yet, retrying in 2m") time.Sleep(2 * time.Minute) continue } // Kill syncs for relays no longer needed. for url2, pid := range children { if _, still := relayAuthors[url2]; !still { fmt.Fprintf(os.Stderr, "outbox: relay removed, killing sync pid=%d url=%s\n", pid, url2) proc, _ := os.FindProcess(pid) if proc != nil { proc.Kill() } delete(children, url2) } } // Spawn syncs for new relays. for url1, authors := range relayAuthors { if _, running := children[url1]; running { continue } authorsArg := "" for i, pk := range authors { if i > 0 { authorsArg = authorsArg | "," } authorsArg = authorsArg | pk } cmd := os.Args[0] | " sync --authors " | authorsArg | " " | url1 | " " | localURL argv := []string{"/bin/sh", "-c", cmd} attr := &os.ProcAttr{} proc, err := os.StartProcess("/bin/sh", argv, attr) if err != nil { fmt.Fprintf(os.Stderr, "outbox: spawn failed for %s: %v\n", url1, err) continue } fmt.Fprintf(os.Stderr, "outbox: spawned sync pid=%d url=%s authors=%d\n", proc.Pid, url1, len(authors)) children[url1] = proc.Pid } // Re-discover every 30 minutes to pick up follow list changes. time.Sleep(30 * time.Minute) } } // outboxDiscover queries the local relay for the user's follows (kind 3) and // their relay lists (kind 10002), returning a map of writeRelayURL → []pubkeyHex. func outboxDiscover(pubkey, localURL string) (m map[string][]string) { local, err7 := ws.Dial(localURL) if err7 != nil { fmt.Fprintf(os.Stderr, "outbox: local connect error: %v\n", err7) return nil } defer local.Close() // Step 1: fetch the user's kind 3 (follows list). req3 := []byte(`["REQ","ob-k3",{"kinds":[3],"authors":["` | pubkey | `"],"limit":1}]`) if err6 := local.WriteText(req3); err6 != nil { fmt.Fprintf(os.Stderr, "outbox: REQ k3 error: %v\n", err6) return nil } var follows []string for { op, payload, err4 := local.ReadMessage() if err4 != nil || op == ws.OpClose { break } if op != ws.OpText { continue } label, rem, _ := envelope.Identify(payload) if label == envelope.EOSELabel { break } if label != envelope.EventLabel { continue } var er envelope.EventResult if _, err3 := er.Unmarshal(rem); err3 != nil || er.Event == nil { continue } ev := er.Event if ev.Kind != 3 || ev.Tags == nil { continue } for _, t := range ev.Tags.GetAll([]byte("p")) { if t.Len() >= 2 { pk := string(t.ValueHex()) if len(pk) == 64 { follows = mxutil.Ensure(follows, 1) follows = push(follows, pk) } } } } if len(follows) == 0 { fmt.Fprintln(os.Stderr, "outbox: no follows in local store, bootstrapping from seed relays") follows = outboxBootstrapFollows(pubkey, localURL) if len(follows) == 0 { return nil } } // Always include the user's own pubkey. follows = mxutil.Ensure(follows, 1) follows = push(follows, pubkey) fmt.Fprintf(os.Stderr, "outbox: found %d follows\n", len(follows)) // Step 2: fetch kind 10002 for all follows (batched). // Build authors array. authorsJSON := `["` | follows[0] | `"` for _, pk := range follows[1:] { authorsJSON = authorsJSON | `,"` | pk | `"` } authorsJSON = authorsJSON | `]` req10002 := []byte(`["REQ","ob-rl",{"kinds":[10002],"authors":` | authorsJSON | `,"limit":` | itoa64(int64(len(follows)*2)) | `}]`) if err5 := local.WriteText(req10002); err5 != nil { fmt.Fprintf(os.Stderr, "outbox: REQ k10002 error: %v\n", err5) return nil } // Map pubkey → write relay URLs. writeRelays := map[string][]string{} // pubkey → []relayURL for { op, payload, err2 := local.ReadMessage() if err2 != nil || op == ws.OpClose { break } if op != ws.OpText { continue } label, rem, _ := envelope.Identify(payload) if label == envelope.EOSELabel { break } if label != envelope.EventLabel { continue } var er envelope.EventResult if _, err1 := er.Unmarshal(rem); err1 != nil || er.Event == nil { continue } ev := er.Event if ev.Kind != 10002 || ev.Tags == nil { continue } // Key by the hex spelling the follow list uses: a p-tag's // ValueHex is hex, so a raw 32-byte event pubkey never matched a // follow and every author fell back to the default relay. pk := hexEnc(ev.Pubkey) for _, t := range ev.Tags.GetAll([]byte("r")) { if t.Len() < 2 { continue } relayURL := string(t.Value()) if !hasPrefix(relayURL, "wss://") && !hasPrefix(relayURL, "ws://") { continue } // marker is t.T[2] if present: "read", "write", or absent (both). marker := "" if t.Len() >= 3 { marker = string(t.T[2]) } isWrite := marker == "write" || marker == "" if isWrite { writeRelays[pk] = push(writeRelays[pk], relayURL) } } } // Invert: relayURL → []pubkeys that write there. relayAuthors := map[string][]string{} for _, pk := range follows { urls, ok := writeRelays[pk] if !ok { // No kind 10002 found: fall back to the default outbox relay. urls = []string{"wss://relay.damus.io"} } for _, u := range urls { relayAuthors[u] = push(relayAuthors[u], pk) } } fmt.Fprintf(os.Stderr, "outbox: %d write relays discovered\n", len(relayAuthors)) return relayAuthors } func spawnOutbox(localAddr, pubkey string) { host := localAddr if hasPrefix(host, "0.0.0.0:") { host = "127.0.0.1:" | host[len("0.0.0.0:"):] } // Small delay so the relay's epoll loop is accepting before we connect. cmd := "sleep 3 && " | os.Args[0] | " outbox " | pubkey | " ws://" | host argv := []string{"/bin/sh", "-c", cmd} attr := &os.ProcAttr{} _, err := os.StartProcess("/bin/sh", argv, attr) if err != nil { fmt.Fprintf(os.Stderr, "outbox: spawn failed: %v\n", err) } } // outboxBootstrapFollows fetches the user's kind 3 from seed relays and // publishes it into the local store so outboxDiscover can find it. // Returns the extracted follow pubkeys. func outboxBootstrapFollows(pubkey, localURL string) (ss []string) { for _, seed := range bootstrapSeeds { fmt.Fprintf(os.Stderr, "outbox: bootstrap: trying %s\n", seed) remote, err4 := ws.Dial(seed) if err4 != nil { continue } req := []byte(`["REQ","ob-boot",{"kinds":[3,10002],"authors":["` | pubkey | `"],"limit":2}]`) if err3 := remote.WriteText(req); err3 != nil { remote.Close() continue } local, lerr := ws.Dial(localURL) if lerr != nil { remote.Close() continue } var follows []string evCount := 0 for { op, payload, err2 := remote.ReadMessage() if err2 != nil || op == ws.OpClose { fmt.Fprintf(os.Stderr, "outbox: bootstrap: %s closed after %d events (err=%v)\n", seed, evCount, err2) break } if op != ws.OpText { continue } label, rem, _ := envelope.Identify(payload) if label == envelope.EOSELabel { fmt.Fprintf(os.Stderr, "outbox: bootstrap: %s EOSE after %d events\n", seed, evCount) break } if label != envelope.EventLabel { continue } var er envelope.EventResult if _, err1 := er.Unmarshal(rem); err1 != nil || er.Event == nil { continue } ev := er.Event evCount++ fmt.Fprintf(os.Stderr, "outbox: bootstrap: %s event kind=%d\n", seed, ev.Kind) // Publish to local store. fwd := &envelope.EventSubmission{E: ev} local.WriteText(fwd.Marshal(nil)) local.ReadMessage() // drain OK // Extract follows from kind 3. Use ValueHex() because "p" tag values // may be binary-encoded (32 bytes) rather than hex strings (64 chars). if ev.Kind == 3 && ev.Tags != nil { for _, t := range ev.Tags.GetAll([]byte("p")) { if t.Len() >= 2 { pk := string(t.ValueHex()) if len(pk) == 64 { follows = mxutil.Ensure(follows, 1) follows = push(follows, pk) } } } } } remote.Close() local.Close() if len(follows) > 0 { fmt.Fprintf(os.Stderr, "outbox: bootstrap: found %d follows from %s\n", len(follows), seed) // Second pass: fetch kind 10002 for all follows so outboxDiscover // can map them to write relays without another bootstrap cycle. outboxBootstrapRelayLists(follows, localURL, seed) return follows } } fmt.Fprintln(os.Stderr, "outbox: bootstrap: no follows found in seed relays") return nil } // outboxBootstrapRelayLists fetches kind 10002 for all follows from one relay // and publishes them to local so outboxDiscover can build the write relay map. func outboxBootstrapRelayLists(follows []string, localURL, seedURL string) { if len(follows) == 0 { return } remote, err4 := ws.Dial(seedURL) if err4 != nil { return } defer remote.Close() local, lerr := ws.Dial(localURL) if lerr != nil { return } defer local.Close() // Build authors JSON for all follows. authorsJSON := `["` | follows[0] | `"` for _, pk := range follows[1:] { authorsJSON = authorsJSON | `,"` | pk | `"` } authorsJSON = authorsJSON | `]` lim := itoa64(int64(len(follows) * 2)) req := []byte(`["REQ","ob-rl2",{"kinds":[10002],"authors":` | authorsJSON | `,"limit":` | lim | `}]`) if err3 := remote.WriteText(req); err3 != nil { return } count := 0 for { op, payload, err2 := remote.ReadMessage() if err2 != nil || op == ws.OpClose { break } if op != ws.OpText { continue } label, rem, _ := envelope.Identify(payload) if label == envelope.EOSELabel { break } if label != envelope.EventLabel { continue } var er envelope.EventResult if _, err1 := er.Unmarshal(rem); err1 != nil || er.Event == nil { continue } fwd := &envelope.EventSubmission{E: er.Event} local.WriteText(fwd.Marshal(nil)) local.ReadMessage() // drain OK count++ } fmt.Fprintf(os.Stderr, "outbox: bootstrap: stored %d kind 10002 events from %s\n", count, seedURL) } func hasPrefix(s, prefix string) (ok bool) { return len(s) >= len(prefix) && s[:len(prefix)] == prefix } func hasSuffix(s, suffix string) (ok bool) { return len(s) >= len(suffix) && s[len(s)-len(suffix):] == suffix } func spawnCrawler(listenAddr string) { host := listenAddr if hasPrefix(host, "0.0.0.0:") { host = "127.0.0.1:" | host[len("0.0.0.0:"):] } cmd := os.Args[0] | " crawl ws://" | host argv := []string{"/bin/sh", "-c", cmd} attr := &os.ProcAttr{} _, err := os.StartProcess("/bin/sh", argv, attr) if err != nil { fmt.Fprintf(os.Stderr, "crawl: spawn failed: %v\n", err) } } func spawnSync(listenAddr, remoteURL string) { host := listenAddr if hasPrefix(host, "0.0.0.0:") { host = "127.0.0.1:" | host[len("0.0.0.0:"):] } cmd := os.Args[0] | " sync " | remoteURL | " ws://" | host argv := []string{"/bin/sh", "-c", cmd} attr := &os.ProcAttr{} _, err := os.StartProcess("/bin/sh", argv, attr) if err != nil { fmt.Fprintf(os.Stderr, "sync: spawn failed: %v\n", err) } } // --- crawl command --- var crawlSeeds []string // bootstrapSeeds are tried specifically for kind 3/10002 lookups. // purplepag.es is designed for profile and contact list events. var bootstrapSeeds []string // Directory event kinds to fetch. const crawlKindsFilter = `[0,3,5,1984,10000,10002,10050]` func clog(out *os.File, format string, args ...fmt.Stringer) { ts := time.Now().Format("15:04:05") fmt.Fprintf(out, ts|" "|format|"\n", args) } // relayDB tracks known relays with frequency scores. // Higher score = seen more often in kind 10002/10050 events = higher priority. type relayDB struct { score map[string]int32 order []string // URLs sorted by descending score } func newRelayDB() (p *relayDB) { return &relayDB{score: map[string]int32{}} } func (db *relayDB) add(relayURL string, weight int32) { db.score[relayURL] += weight } // sorted returns relay URLs ordered by descending frequency. func (db *relayDB) sorted() (ss []string) { urls := []string{:0:len(db.score)} for u := range db.score { urls = mxutil.Ensure(urls, 1) urls = push(urls, u) } // Simple insertion sort by score descending. for i := 1; i < len(urls); i++ { for j := i; j > 0 && db.score[urls[j]] > db.score[urls[j-1]]; j-- { urls[j], urls[j-1] = urls[j-1], urls[j] } } return urls } func runCrawl(args []string) { var err error var out *os.File out, err = os.OpenFile("/tmp/musiquay-crawl.log", os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0644) if err != nil { out = os.Stderr } localURL := "ws://127.0.0.1:3334" if len(args) >= 1 { localURL = args[0] } clog(out, "started pid=%d local=%s", os.Getpid(), localURL) db := newRelayDB() // Seed relays get a high initial score. for _, s := range crawlSeeds { db.add(s, 100) } pass := 0 for { pass++ clog(out, "=== pass %d, %d relays known ===", pass, len(db.score)) ok := crawlPass(localURL, db, out) if ok { clog(out, "pass complete, sleeping 5m") time.Sleep(5 * time.Minute) } else { clog(out, "pass failed, retrying in 30s") time.Sleep(30 * time.Second) } } } func crawlPass(localURL string, db *relayDB, out *os.File) (ok bool) { relays := db.sorted() if len(relays) == 0 { clog(out, "no relays known") return false } totalEvents := 0 for i, relayURL := range relays { clog(out, "[%d/%d] crawling %s (score %d)", i+1, len(relays), relayURL, db.score[relayURL]) events := crawlRelay(relayURL, out) if len(events) == 0 { clog(out, " %s → 0 events", relayURL) time.Sleep(1 * time.Second) continue } clog(out, " %s → %d events", relayURL, len(events)) // Extract new relay URLs from the events before publishing. for _, raw := range events { crawlExtractRelays(raw, db) } // Publish batch to local relay. published := crawlPublishBatch(localURL, events, out) clog(out, " published %d/%d to local", published, len(events)) totalEvents += published time.Sleep(1 * time.Second) } clog(out, "total %d events from %d relays", totalEvents, len(relays)) return true } // crawlRelay connects to one relay and subscribes to directory events. // Returns raw EVENT messages suitable for republishing. func crawlRelay(relayURL string, out *os.File) (ss [][]byte) { remote, err4 := ws.Dial(relayURL) if err4 != nil { clog(out, " dial %s FAILED: %v", relayURL, err4) return nil } defer remote.Close() reqJSON := []byte(`["REQ","cr",{"kinds":` | crawlKindsFilter | `,"limit":200}]`) if err3 := remote.WriteText(reqJSON); err3 != nil { clog(out, " write REQ to %s failed: %v", relayURL, err3) return nil } var events [][]byte for { op, payload, err2 := remote.ReadMessage() if err2 != nil { break } if op != ws.OpText { continue } label, rem, _ := envelope.Identify(payload) if label == envelope.EOSELabel { break } if label == envelope.EventLabel { var er envelope.EventResult if _, err1 := er.Unmarshal(rem); err1 == nil && er.Event != nil { es := &envelope.EventSubmission{E: er.Event} events = mxutil.Ensure(events, 1) events = push(events, es.Marshal(nil)) } } _ = rem } return events } // crawlExtractRelays parses an EVENT submission and adds discovered relay // URLs to the database with a frequency bump. func crawlExtractRelays(raw []byte, db *relayDB) { _, rem, err2 := envelope.Identify(raw) if err2 != nil { return } var es envelope.EventSubmission if _, err1 := es.Unmarshal(rem); err1 != nil || es.E == nil { return } ev := es.E if (ev.Kind != 10002 && ev.Kind != 10050) || ev.Tags == nil { return } for _, t := range ev.Tags.GetAll([]byte("r")) { if t.Len() >= 2 { relayURL := string(t.Value()) if len(relayURL) > 5 && (hasPrefix(relayURL, "wss://") || hasPrefix(relayURL, "ws://")) { db.add(relayURL, 1) } } } } // crawlPublishBatch publishes a batch of EVENT messages to the local relay. func crawlPublishBatch(localURL string, events [][]byte, out *os.File) (n int32) { local, err2 := ws.Dial(localURL) if err2 != nil { clog(out, " local connect failed: %v", err2) return 0 } defer local.Close() for _, evBytes := range events { local.WriteText(evBytes) } // Drain OKs - one per event sent. count := 0 for count < len(events) { _, _, err1 := local.ReadMessage() if err1 != nil { break } count++ } return count } func hexEnc(b []byte) (s string) { const hx = "0123456789abcdef" out := []byte{:len(b)*2} for i, v := range b { out[i*2] = hx[v>>4] out[i*2+1] = hx[v&0x0f] } return string(out) } func i64str(n int64) (s string) { if n == 0 { return "0" } neg := false if n < 0 { neg = true n = -n } var buf [20]byte i := 19 for n > 0 { buf[i] = byte('0' + n%10) i-- n /= 10 } if neg { buf[i] = '-' i-- } return string(buf[i+1:]) } func crawlAppendUniq(ss []string, s string) (ss2 []string) { for _, x := range ss { if x == s { return ss } } return push(ss, s) }