From c64455ccfa13bff91f313329dbef8cfe92a57782 Mon Sep 17 00:00:00 2001 From: Marcus Pasell <3690498+rickyrombo@users.noreply.github.com> Date: Thu, 27 Aug 2026 21:53:53 -0700 Subject: [PATCH] feat(api): route stream resolution through named hosts during the chain migration Plays are recorded by whichever node serves the audio -- logTrackListen runs at the top of mediorum's serveBlob, before it 307s to storage -- and they never travel through the relay, so the queue that carries ManageEntity writes to the new chain does not carry them. That matters for days during the genesis migration. Nodes move to the new chain in batches while the indexer still reads the old one, so every play recorded by an already-migrated node lands on a chain nobody is indexing. The fleet migration is measured in days, and roughly a quarter of plays would be lost across that window. playRoutingHosts, when set, puts those hosts first when resolving a stream URL. Naming nodes that stay on the old chain keeps plays on the chain the indexer is reading. Unset, this is inert. The original url and mirrors are kept behind the routing hosts rather than replaced: store-all nodes hold nearly everything, but a fresh upload has not necessarily replicated, and such a track must still stream -- with the play following whichever node serves it. tryFindWorkingUrl probes in order, so a routing host that cannot serve costs one request and falls through. Bandwidth is not a concern: those nodes 307 to presigned storage rather than streaming bytes, so the added cost is one request and a URL signature per play. Tests cover ordering, fallback retention, dedupe against existing mirrors, that the path and the signature mediorum parses for attribution survive the host rewrite, bare-host and full-URL forms, and that an unparseable link degrades to current behaviour. Confirmed failing with the routing disabled. --- api/play_routing.go | 89 ++++++++++++++++++++++++++++++++++++++++ api/play_routing_test.go | 89 ++++++++++++++++++++++++++++++++++++++++ api/v1_track_stream.go | 5 +++ config/config.go | 21 ++++++++++ 4 files changed, 204 insertions(+) create mode 100644 api/play_routing.go create mode 100644 api/play_routing_test.go diff --git a/api/play_routing.go b/api/play_routing.go new file mode 100644 index 00000000..d6752d0a --- /dev/null +++ b/api/play_routing.go @@ -0,0 +1,89 @@ +package api + +import ( + "net/url" + "strings" + + "api.audius.co/api/dbv1" +) + +// withPlayRoutingHosts returns a copy of link whose candidate hosts are the +// configured routing hosts first, then the link's own url and mirrors. +// +// A play is recorded by whichever node serves the audio -- logTrackListen runs +// at the top of mediorum's serveBlob, before it 307s to storage -- and plays +// never travel through the relay. So the host in the URL the API hands back is +// what decides which chain the play lands on. +// +// During the genesis migration the fleet is split across two chains for days. +// An already-migrated node writes its plays to the new chain while the indexer +// is still reading the old one, and those plays are indexed by nobody. Naming +// hosts that stay on the old chain keeps every play on the chain the indexer is +// actually reading. Cleared at the cutover, after which plays follow the node +// serving them again. See cmd/genesis-writer/ROLLOUT.md, Runbook steps 5 and 13. +// +// The original url and mirrors are kept as fallbacks rather than replaced. The +// routing hosts are store-all nodes and hold essentially everything, but +// replication of a fresh upload is not instant, and a track they do not have +// yet must still be streamable. tryFindWorkingUrl probes in order, so a routing +// host that cannot serve costs one request and falls through. +func withPlayRoutingHosts(link *dbv1.MediaLink, hosts []string) *dbv1.MediaLink { + if link == nil || len(hosts) == 0 { + return link + } + + primary, err := url.Parse(link.Url) + if err != nil { + return link + } + + routed := &dbv1.MediaLink{} + seen := make(map[string]bool, len(hosts)+len(link.Mirrors)+1) + + // Host-only comparison: mirrors are recorded as hosts, and the routing hosts + // are configured as hosts, so normalising to that avoids listing the same + // node twice under different spellings. + add := func(raw string) { + h := hostOf(raw) + if h == "" || seen[h] { + return + } + seen[h] = true + if routed.Url == "" { + u := *primary + u.Host = h + routed.Url = u.String() + return + } + routed.Mirrors = append(routed.Mirrors, h) + } + + for _, h := range hosts { + add(h) + } + add(link.Url) + for _, m := range link.Mirrors { + add(m) + } + + if routed.Url == "" { + return link + } + return routed +} + +// hostOf accepts either a bare host or a full URL and returns the host. +func hostOf(s string) string { + s = strings.TrimSpace(s) + if s == "" { + return "" + } + if !strings.Contains(s, "//") { + s = "https://" + s + } + u, err := url.Parse(s) + if err != nil { + return "" + } + return u.Host +} diff --git a/api/play_routing_test.go b/api/play_routing_test.go new file mode 100644 index 00000000..3c59ab17 --- /dev/null +++ b/api/play_routing_test.go @@ -0,0 +1,89 @@ +package api + +import ( + "net/url" + "testing" + + "api.audius.co/api/dbv1" + "github.com/stretchr/testify/require" +) + +const streamPath = "/tracks/cidstream/abc?signature=sig" + +func hostsOf(t *testing.T, link *dbv1.MediaLink) []string { + t.Helper() + u, err := url.Parse(link.Url) + require.NoError(t, err) + return append([]string{u.Host}, link.Mirrors...) +} + +// Unconfigured, this must be exactly inert -- it ships ahead of the migration +// and sits dormant in production until someone sets the env var. +func TestPlayRoutingIsInertWhenUnconfigured(t *testing.T) { + link := &dbv1.MediaLink{Url: "https://node-a.example" + streamPath, Mirrors: []string{"node-b.example"}} + require.Same(t, link, withPlayRoutingHosts(link, nil)) + require.Same(t, link, withPlayRoutingHosts(link, []string{})) + require.Nil(t, withPlayRoutingHosts(nil, []string{"x.example"})) +} + +// Routing hosts go first, because tryFindWorkingUrl probes in order and the +// first host that can serve is the one that records the play. +func TestPlayRoutingHostsAreTriedFirst(t *testing.T) { + link := &dbv1.MediaLink{Url: "https://node-a.example" + streamPath, Mirrors: []string{"node-b.example"}} + routed := withPlayRoutingHosts(link, []string{"creatornode.audius.co", "v.monophonic.digital"}) + + require.Equal(t, + []string{"creatornode.audius.co", "v.monophonic.digital", "node-a.example", "node-b.example"}, + hostsOf(t, routed)) +} + +// The original hosts stay as fallbacks. Store-all nodes hold nearly everything, +// but a freshly uploaded track may not have replicated yet, and it still has to +// be streamable -- just with the play landing on whichever node serves it. +func TestPlayRoutingKeepsOriginalHostsAsFallback(t *testing.T) { + link := &dbv1.MediaLink{Url: "https://node-a.example" + streamPath, Mirrors: []string{"node-b.example"}} + routed := withPlayRoutingHosts(link, []string{"creatornode.audius.co"}) + + require.Contains(t, hostsOf(t, routed), "node-a.example") + require.Contains(t, hostsOf(t, routed), "node-b.example") +} + +// The path and query -- including the signature mediorum parses to attribute the +// listen -- must survive the host rewrite, or the play is recorded against the +// wrong user or not at all. +func TestPlayRoutingPreservesPathAndSignature(t *testing.T) { + link := &dbv1.MediaLink{Url: "https://node-a.example" + streamPath} + routed := withPlayRoutingHosts(link, []string{"creatornode.audius.co"}) + + u, err := url.Parse(routed.Url) + require.NoError(t, err) + require.Equal(t, "creatornode.audius.co", u.Host) + require.Equal(t, "/tracks/cidstream/abc", u.Path) + require.Equal(t, "sig", u.Query().Get("signature")) +} + +// A routing host that is already the primary or a mirror must not be probed +// twice; duplicates would waste a request and could double-count if one of them +// ever lost its skip_play_count. +func TestPlayRoutingDeduplicatesHosts(t *testing.T) { + link := &dbv1.MediaLink{Url: "https://creatornode.audius.co" + streamPath, Mirrors: []string{"node-b.example"}} + routed := withPlayRoutingHosts(link, []string{"creatornode.audius.co", "node-b.example"}) + + require.Equal(t, []string{"creatornode.audius.co", "node-b.example"}, hostsOf(t, routed)) +} + +// Hosts may be configured bare or as full URLs; both must normalise to the same +// thing so a scheme in the env var does not silently create a duplicate. +func TestPlayRoutingAcceptsBareHostsAndUrls(t *testing.T) { + link := &dbv1.MediaLink{Url: "https://node-a.example" + streamPath} + bare := withPlayRoutingHosts(link, []string{"creatornode.audius.co"}) + full := withPlayRoutingHosts(link, []string{"https://creatornode.audius.co"}) + require.Equal(t, hostsOf(t, bare), hostsOf(t, full)) +} + +// An unparseable link is returned untouched rather than dropped: streaming +// should degrade to current behaviour, never fail, because of this feature. +func TestPlayRoutingLeavesUnparseableLinkAlone(t *testing.T) { + link := &dbv1.MediaLink{Url: "://not a url"} + require.Same(t, link, withPlayRoutingHosts(link, []string{"creatornode.audius.co"})) +} diff --git a/api/v1_track_stream.go b/api/v1_track_stream.go index 98393d22..1645317e 100644 --- a/api/v1_track_stream.go +++ b/api/v1_track_stream.go @@ -52,6 +52,11 @@ func (app *ApiServer) v1TrackStream(c *fiber.Ctx) error { } func (app *ApiServer) redirectToStream(c *fiber.Ctx, stream *dbv1.MediaLink) error { + // Temporary, for the genesis migration: prefer hosts that stay on the old + // chain so plays are not split across two chains while the fleet migrates. + // No-op when unconfigured. See withPlayRoutingHosts. + stream = withPlayRoutingHosts(stream, app.config.PlayRoutingHosts) + streamURL := tryFindWorkingUrl(stream) if skipPlayCount := c.Query("skip_play_count"); skipPlayCount != "" { diff --git a/config/config.go b/config/config.go index ed1a8439..15aabc55 100644 --- a/config/config.go +++ b/config/config.go @@ -89,6 +89,18 @@ type Config struct { // confirmed_block < NewChainFlushFromBlock before sending — trimming rows already // covered by the backfill. // NewChainInsecureSkipVerify disables TLS verification for the new chain endpoint (e.g. localstack). + // PlayRoutingHosts, when non-empty, are tried first when resolving a track + // stream URL. Plays are recorded by whichever node serves the audio, never + // through the relay, so this is what decides which chain a play lands on. + // + // It exists for the genesis migration: while the fleet is split across two + // chains, an already-migrated node writes its plays to the new chain even + // though the indexer is still reading the old one, and every play in that + // window goes to a chain nobody reads. Pointing this at nodes that stay on + // the old chain keeps plays where the indexer is. Cleared at the cutover. + // See cmd/genesis-writer/ROLLOUT.md, Runbook steps 5 and 13. + PlayRoutingHosts []string + NewChainURL string NewChainQueueEnabled bool NewChainFlushEnabled bool @@ -362,6 +374,15 @@ func init() { Cfg.FeaturedAudienceUserID = int32(parsed) } + // Genesis migration: temporary play routing (see the struct field). + if v := strings.TrimSpace(os.Getenv("playRoutingHosts")); v != "" { + for _, h := range strings.Split(v, ",") { + if h = strings.TrimSpace(h); h != "" { + Cfg.PlayRoutingHosts = append(Cfg.PlayRoutingHosts, h) + } + } + } + // Genesis migration dual-write queue Cfg.NewChainURL = os.Getenv("newChainUrl") Cfg.NewChainQueueEnabled = os.Getenv("newChainQueueEnabled") == "true"