diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 22b211ea..53a793e4 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -24,6 +24,7 @@ jobs: go_sdk: ${{ steps.filter.outputs.go_sdk }} rust_sdk: ${{ steps.filter.outputs.rust_sdk }} deezer: ${{ steps.filter.outputs.deezer }} + remote_ws: ${{ steps.filter.outputs.remote_ws }} steps: - uses: actions/checkout@v4 - uses: dorny/paths-filter@v3 @@ -45,6 +46,9 @@ jobs: deezer: - 'deezer/**' - '.github/workflows/tests.yml' + remote_ws: + - 'remote-ws/**' + - '.github/workflows/tests.yml' cli: needs: changes @@ -128,3 +132,31 @@ jobs: # CI so the suite never depends on the external API or its rate limits. - name: Run Deezer service tests run: go test ./... + + remote-ws: + needs: changes + if: needs.changes.outputs.remote_ws == 'true' + runs-on: ubuntu-latest + defaults: + run: + working-directory: remote-ws + steps: + - uses: actions/checkout@v4 + - uses: erlef/setup-beam@v1 + with: + otp-version: "27" + elixir-version: "1.18" + # The relay's unit tests use in-memory doubles (config/test.exs + # start_externals: false), so no Postgres / Redis / NATS services are needed. + - name: Cache deps and build + uses: actions/cache@v4 + with: + path: | + remote-ws/deps + remote-ws/_build + key: ${{ runner.os }}-mix-${{ hashFiles('remote-ws/mix.lock') }} + restore-keys: ${{ runner.os }}-mix- + - name: Install dependencies + run: mix deps.get + - name: Run tests + run: mix test diff --git a/remote-ws/.formatter.exs b/remote-ws/.formatter.exs new file mode 100644 index 00000000..580c5079 --- /dev/null +++ b/remote-ws/.formatter.exs @@ -0,0 +1,4 @@ +[ + import_deps: [:ecto, :ecto_sql, :phoenix], + inputs: ["*.{ex,exs}", "{config,lib,test}/**/*.{ex,exs}"] +] diff --git a/remote-ws/.gitignore b/remote-ws/.gitignore new file mode 100644 index 00000000..235a7554 --- /dev/null +++ b/remote-ws/.gitignore @@ -0,0 +1,14 @@ +# Elixir build artifacts and fetched dependencies +/_build/ +/cover/ +/deps/ +/doc/ +/.fetch +erl_crash.dump +*.ez +*.beam +/config/*.secret.exs +.elixir_ls/ + +# Env files +.env diff --git a/remote-ws/README.md b/remote-ws/README.md new file mode 100644 index 00000000..219ba45b --- /dev/null +++ b/remote-ws/README.md @@ -0,0 +1,67 @@ +# remote-ws + +An **Elixir / Phoenix / Ecto** port of the player remote-control WebSocket +server that currently lives in `apps/api/src/websocket/handler.ts`. + +Phase 1: this is a **1:1 port** — the existing Node `/ws` endpoint stays +authoritative. This service speaks the identical wire protocol and publishes the +same NATS events so it can eventually replace the Node relay with no client or +consumer changes. + +## What it does + +A raw-JSON WebSocket relay at `GET /ws` that lets the web/mobile miniplayers and +player devices (the Rocksky CLI, the Rockbox companion) control each other: + +- `register` a device, scoped by the JWT's DID +- broadcast a device's now-playing (`track`) and transport (`status`) to the + user's other devices, enriched from Redis/Postgres +- relay `command`s (play/pause/next/previous/seek) to devices +- publish `rocksky.song.changed` / `rocksky.song.stopped` to NATS (with the same + 15s stop debounce and `ws_lastsong` gating as the Node server) + +Because it's a raw-JSON protocol (not Phoenix Channel envelopes), the socket is a +`WebSock` handler served by Bandit — **not** a Phoenix Channel — so every existing +client works unchanged. + +## Architecture (map from the Node handler) + +| Node (`handler.ts`) | Elixir | +| --- | --- | +| `@hono/node-ws` upgrade at `/ws` | `WebSockAdapter.upgrade` → `RemoteWs.Ws.Connection` (`WebSock`) | +| `devices` / `deviceNames` / `userDevices` maps | `RemoteWs.Devices` over a duplicate `Registry` keyed by DID (auto-cleans on disconnect) | +| `pendingStop` + `setTimeout` | `RemoteWs.StopDebouncer` GenServer | +| `verifyToken` (jsonwebtoken) | `RemoteWs.Auth` (Joken HS256, ignore-expiration + jti revoke check) | +| `ctx.redis` | `RemoteWs.Redis` behaviour → `RemoteWs.Redis.Redix` | +| `ctx.db` (drizzle) | `RemoteWs.Store` behaviour → `RemoteWs.Store.Ecto` + `RemoteWs.Repo` | +| `ctx.nc` (NATS) | `RemoteWs.Nats` behaviour → `RemoteWs.Nats.Gnat` | +| enrichment + gating (lines 64-273) | `RemoteWs.NowPlaying` | +| `onMessage` dispatch | `RemoteWs.Ws.Handler` | + +Redis, NATS, and the read store sit behind behaviours so the intricate gating and +debounce logic is unit-tested against in-memory doubles — no live Redis / NATS / +Postgres required. + +## Environment + +Reuses the **same variable names** as `apps/api` (share one environment): + +| Var | Purpose | +| --- | --- | +| `JWT_SECRET` | HS256 secret for verifying bearer tokens | +| `XATA_POSTGRES_URL` | Postgres connection URL (Ecto) | +| `REDIS_URL` | Redis URL (default `redis://localhost:6379`) | +| `NATS_URL` | NATS URL (default `nats://localhost:4222`) | +| `REMOTE_WS_PORT` | HTTP listen port (service-specific; default `4000`) | +| `SECRET_KEY_BASE` | Phoenix endpoint secret (prod only) | + +## Develop + +```bash +mix deps.get +mix test # 29 tests, no external services needed +REMOTE_WS_PORT=4000 mix phx.server +``` + +Deployed via `systemd/rocksky-remote-ws.service`. Tests run in CI through the +`remote-ws` job in `.github/workflows/tests.yml`. diff --git a/remote-ws/config/config.exs b/remote-ws/config/config.exs new file mode 100644 index 00000000..d1f09888 --- /dev/null +++ b/remote-ws/config/config.exs @@ -0,0 +1,31 @@ +import Config + +# Compile-time configuration shared by all environments. Environment-specific +# values (secrets, URLs) are read at runtime — see config/runtime.exs. + +config :remote_ws, + ecto_repos: [RemoteWs.Repo], + # Pluggable side-effect adapters — swapped for in-memory doubles in test so the + # relay logic (enrichment, gating, debounce) is unit-testable without live + # Redis / NATS / Postgres. See config/test.exs. + redis: RemoteWs.Redis.Redix, + nats: RemoteWs.Nats.Gnat, + store: RemoteWs.Store.Ecto, + # Whether to boot the external clients (Repo, Redix, NATS) in the supervision + # tree. Disabled in test. + start_externals: true + +# The raw-JSON WebSocket relay is served by a Bandit-backed Phoenix endpoint. +config :remote_ws, RemoteWsWeb.Endpoint, + adapter: Bandit.PhoenixAdapter, + url: [host: "localhost"], + render_errors: [formats: [json: RemoteWsWeb.ErrorJSON], layout: false], + pubsub_server: RemoteWs.PubSub + +config :phoenix, :json_library, Jason + +config :logger, :console, + format: "$time $metadata[$level] $message\n", + metadata: [:request_id, :did] + +import_config "#{config_env()}.exs" diff --git a/remote-ws/config/dev.exs b/remote-ws/config/dev.exs new file mode 100644 index 00000000..4f322fd7 --- /dev/null +++ b/remote-ws/config/dev.exs @@ -0,0 +1,11 @@ +import Config + +config :remote_ws, RemoteWsWeb.Endpoint, + # Bind on all interfaces in dev; the port comes from runtime.exs (PORT). + http: [ip: {0, 0, 0, 0}, port: String.to_integer(System.get_env("REMOTE_WS_PORT") || "4000")], + server: true, + debug_errors: true, + check_origin: false, + secret_key_base: String.duplicate("dev", 30) + +config :logger, level: :debug diff --git a/remote-ws/config/runtime.exs b/remote-ws/config/runtime.exs new file mode 100644 index 00000000..e2752885 --- /dev/null +++ b/remote-ws/config/runtime.exs @@ -0,0 +1,41 @@ +import Config + +# Runtime configuration — reads the SAME environment variable names as the +# existing Node API (apps/api) so the two services share one environment: +# JWT_SECRET, XATA_POSTGRES_URL, REDIS_URL, NATS_URL. +# The listen port is service-specific: REMOTE_WS_PORT (so it never clashes with +# the Node API's PORT). + +# JWT signing secret (mirrors apps/api env.JWT_SECRET). Only override when set so +# the test config's fixed secret stays intact. +if secret = System.get_env("JWT_SECRET") do + config :remote_ws, :jwt_secret, secret +end + +# Redis connection URL (mirrors apps/api env.REDIS_URL). +config :remote_ws, :redis_url, System.get_env("REDIS_URL") || "redis://localhost:6379" + +# NATS connection URL (mirrors apps/api env.NATS_URL). +config :remote_ws, :nats_url, System.get_env("NATS_URL") || "nats://localhost:4222" + +# Postgres (mirrors apps/api env.XATA_POSTGRES_URL). +if database_url = System.get_env("XATA_POSTGRES_URL") do + config :remote_ws, RemoteWs.Repo, + url: database_url, + pool_size: String.to_integer(System.get_env("POOL_SIZE") || "10"), + ssl: true +end + +if config_env() == :prod do + secret_key_base = + System.get_env("SECRET_KEY_BASE") || + raise "SECRET_KEY_BASE is required in production" + + config :remote_ws, RemoteWsWeb.Endpoint, + http: [ + ip: {0, 0, 0, 0}, + port: String.to_integer(System.get_env("REMOTE_WS_PORT") || "4000") + ], + secret_key_base: secret_key_base, + server: true +end diff --git a/remote-ws/config/test.exs b/remote-ws/config/test.exs new file mode 100644 index 00000000..ecd50f32 --- /dev/null +++ b/remote-ws/config/test.exs @@ -0,0 +1,23 @@ +import Config + +# Tests exercise the ported relay logic against in-memory doubles — no live +# Redis, NATS, or Postgres. The externals are not booted and the endpoint does +# not listen. +config :remote_ws, + redis: RemoteWs.Test.RedisMemory, + nats: RemoteWs.Test.NatsCollector, + store: RemoteWs.Test.StoreStub, + start_externals: false + +config :remote_ws, RemoteWsWeb.Endpoint, + http: [ip: {127, 0, 0, 1}, port: 4002], + server: false, + secret_key_base: String.duplicate("test", 30) + +# The JWT signing secret used by the token verifier (mirrors env.JWT_SECRET). +config :remote_ws, :jwt_secret, "test-secret" + +# Short debounce so song.stopped tests don't wait 15s. +config :remote_ws, :stop_debounce_ms, 30 + +config :logger, level: :warning diff --git a/remote-ws/lib/remote_ws/application.ex b/remote-ws/lib/remote_ws/application.ex new file mode 100644 index 00000000..5a3f6678 --- /dev/null +++ b/remote-ws/lib/remote_ws/application.ex @@ -0,0 +1,62 @@ +defmodule RemoteWs.Application do + @moduledoc false + use Application + + @impl true + def start(_type, _args) do + children = + [ + # PubSub is referenced by the Phoenix endpoint config. + {Phoenix.PubSub, name: RemoteWs.PubSub}, + # Tracks connected devices per user (did) — replaces the Node handler's + # module-level `devices` / `deviceNames` / `userDevices` maps. Entries + # are owned by each connection process, so they auto-clean on disconnect. + {Registry, keys: :duplicate, name: RemoteWs.Devices.Registry}, + # Debounced song.stopped timers — replaces the `pendingStop` map. + RemoteWs.StopDebouncer, + RemoteWsWeb.Endpoint + ] ++ external_children() + + opts = [strategy: :one_for_one, name: RemoteWs.Supervisor] + Supervisor.start_link(children, opts) + end + + @impl true + def config_change(changed, _new, removed) do + RemoteWsWeb.Endpoint.config_change(changed, removed) + :ok + end + + # Repo, Redis and NATS clients are only booted when configured (never in test, + # where in-memory doubles are used instead). + defp external_children do + if Application.get_env(:remote_ws, :start_externals, true) do + [RemoteWs.Repo, redis_child(), nats_child()] + |> Enum.reject(&is_nil/1) + else + [] + end + end + + defp redis_child do + url = Application.get_env(:remote_ws, :redis_url, "redis://localhost:6379") + {Redix, {url, [name: RemoteWs.Redix]}} + end + + defp nats_child do + %{host: host, port: port} = + RemoteWs.Nats.Gnat.parse_url(Application.get_env(:remote_ws, :nats_url)) + + %{ + id: :gnat_conn, + start: + {Gnat.ConnectionSupervisor, :start_link, + [ + %{ + name: RemoteWs.Gnat, + connection_settings: [%{host: host, port: port}] + } + ]} + } + end +end diff --git a/remote-ws/lib/remote_ws/auth.ex b/remote-ws/lib/remote_ws/auth.ex new file mode 100644 index 00000000..1fe30dfe --- /dev/null +++ b/remote-ws/lib/remote_ws/auth.ex @@ -0,0 +1,43 @@ +defmodule RemoteWs.Auth do + @moduledoc """ + Bearer-token verification — a port of apps/api lib/verifyToken.ts. + + Verifies the HS256 signature against `:remote_ws, :jwt_secret` (mirrors + env.JWT_SECRET) WITHOUT enforcing expiration (matches jwt.verify's + `ignoreExpiration: true`). If the token carries `type == "access_token"` with a + `jti`, the jti must still exist in the access_tokens table or the token is + treated as revoked. + """ + + @access_token_type "access_token" + + @type verified :: %{did: String.t() | nil, jti: String.t() | nil, type: String.t() | nil} + + @spec verify_token(String.t()) :: {:ok, verified} | {:error, term()} + def verify_token(bearer) when is_binary(bearer) and bearer != "" do + secret = Application.fetch_env!(:remote_ws, :jwt_secret) + signer = Joken.Signer.create("HS256", secret) + + case Joken.verify(bearer, signer) do + {:ok, claims} -> + if claims["type"] == @access_token_type and is_binary(claims["jti"]) do + if RemoteWs.Store.access_token_exists?(claims["jti"]) do + {:ok, extract(claims)} + else + {:error, :revoked} + end + else + {:ok, extract(claims)} + end + + {:error, reason} -> + {:error, reason} + end + end + + def verify_token(_), do: {:error, :invalid_token} + + defp extract(claims) do + %{did: claims["did"], jti: claims["jti"], type: claims["type"]} + end +end diff --git a/remote-ws/lib/remote_ws/devices.ex b/remote-ws/lib/remote_ws/devices.ex new file mode 100644 index 00000000..ded63d5b --- /dev/null +++ b/remote-ws/lib/remote_ws/devices.ex @@ -0,0 +1,61 @@ +defmodule RemoteWs.Devices do + @moduledoc """ + Connected-device bookkeeping, scoped by user DID. Replaces the Node handler's + module-level `devices`, `deviceNames`, and `userDevices` maps with a duplicate + `Registry` keyed by DID: each connection process registers one entry, so + disconnects clean up automatically (replacing the explicit `onClose` handler). + + Frames are delivered to a connection by sending `{:push, frame}` to its process + (see RemoteWs.Ws.Connection.handle_info/2). + """ + + @registry RemoteWs.Devices.Registry + + @doc "Register the CALLING process as a device for `did`. Call from the connection process." + def register(did, device_id, name) do + {:ok, _} = Registry.register(@registry, did, %{device_id: device_id, name: name}) + :ok + end + + @doc "All {pid, %{device_id, name}} entries for a user." + def list(did), do: Registry.lookup(@registry, did) + + @doc "The clientName registered for a given device_id under `did`, or nil." + def name_of(_did, nil), do: nil + + def name_of(did, device_id) do + Enum.find_value(list(did), fn {_pid, %{device_id: id, name: name}} -> + if id == device_id, do: name + end) + end + + @doc "Send a frame to every one of the user's connected devices (including the sender)." + def broadcast(did, frame) do + for {pid, _} <- list(did), do: send(pid, {:push, frame}) + :ok + end + + @doc "Send a frame to every device EXCEPT the one with `except_device_id`." + def broadcast_except(did, except_device_id, frame) do + for {pid, %{device_id: id}} <- list(did), id != except_device_id do + send(pid, {:push, frame}) + end + + :ok + end + + @doc """ + Send a frame to the single device `target_device_id`. Returns :ok if a matching + device was found, :not_found otherwise. + """ + def send_to(did, target_device_id, frame) do + case Enum.find(list(did), fn {_pid, %{device_id: id}} -> id == target_device_id end) do + {pid, _} -> + send(pid, {:push, frame}) + :ok + + nil -> + :not_found + end + end +end diff --git a/remote-ws/lib/remote_ws/nats.ex b/remote-ws/lib/remote_ws/nats.ex new file mode 100644 index 00000000..28661bd6 --- /dev/null +++ b/remote-ws/lib/remote_ws/nats.ex @@ -0,0 +1,12 @@ +defmodule RemoteWs.Nats do + @moduledoc """ + NATS publish surface (mirrors ctx.nc.publish in apps/api). Backed by Gnat in + production; a collector in test. Dispatches to `:remote_ws, :nats`. + """ + + @callback publish(String.t(), iodata()) :: :ok + + defp impl, do: Application.get_env(:remote_ws, :nats, RemoteWs.Nats.Gnat) + + def publish(subject, payload), do: impl().publish(subject, payload) +end diff --git a/remote-ws/lib/remote_ws/nats/gnat.ex b/remote-ws/lib/remote_ws/nats/gnat.ex new file mode 100644 index 00000000..2c6f0be7 --- /dev/null +++ b/remote-ws/lib/remote_ws/nats/gnat.ex @@ -0,0 +1,29 @@ +defmodule RemoteWs.Nats.Gnat do + @moduledoc "Gnat-backed NATS adapter (production)." + @behaviour RemoteWs.Nats + require Logger + + @conn RemoteWs.Gnat + + @impl true + def publish(subject, payload) do + try do + Gnat.pub(@conn, subject, IO.iodata_to_binary(payload)) + catch + kind, reason -> + Logger.error("NATS publish failed for #{subject}: #{inspect({kind, reason})}") + end + + :ok + end + + @doc """ + Parse a `nats://host:port` URL into the shape Gnat's connection settings want. + """ + def parse_url(url) when is_binary(url) do + uri = URI.parse(url) + %{host: uri.host || "localhost", port: uri.port || 4222} + end + + def parse_url(_), do: %{host: "localhost", port: 4222} +end diff --git a/remote-ws/lib/remote_ws/now_playing.ex b/remote-ws/lib/remote_ws/now_playing.ex new file mode 100644 index 00000000..b47475e2 --- /dev/null +++ b/remote-ws/lib/remote_ws/now_playing.ex @@ -0,0 +1,202 @@ +defmodule RemoteWs.NowPlaying do + @moduledoc """ + Now-playing enrichment and the song.changed / song.stopped gating — a faithful + port of apps/api/src/websocket/handler.ts (the `data.type === "track"` and the + status branches). + + All Redis keys, TTLs, and the `ws_lastsong` gate semantics match the Node + implementation 1:1 so both servers behave identically against shared Redis. + """ + + alias RemoteWs.{Nats, Redis, StopDebouncer, Store} + + @day_seconds 86_400 + + @doc """ + Enrich a "track" now-playing payload and emit song.changed when appropriate. + Returns the enriched `data` map (string keys) to broadcast. `source` is the + device name used in the song.changed event. + """ + def handle_track(did, data, source) do + title = data["title"] + artist = data["artist"] + album = data["album"] + sha256 = sha256_hex(String.downcase("#{title} - #{artist} - #{album}")) + + cached_track = Redis.get("track:#{sha256}") + cached_likes = Redis.get("likes:#{did}:#{sha256}") + + {data, liked} = resolve_liked(data, did, sha256, cached_likes) + + duration_ms = data["duration_ms"] || data["duration"] || 0 + {data, duration_ms} = resolve_metadata(data, did, sha256, liked, duration_ms, cached_track) + + maybe_emit_song_changed(did, sha256, %{ + title: title, + artist: artist, + album: album, + album_art: data["album_art"], + duration_ms: duration_ms, + source: source + }) + + data + end + + @doc """ + Handle a status payload (`data.status`): persist the status, and drive the + song.stopped debounce / ws_lastsong reactivation exactly like the Node handler. + """ + def handle_status(did, data) do + status = data["status"] + Redis.set_ex("nowplaying:#{did}:status", 3, "#{status}") + + ws_was_playing = Redis.exists("ws_lastsong:#{did}") + + cond do + status == 0 and ws_was_playing -> + # Do NOT delete ws_lastsong here — only the debounce timer does, so a + # status=1 within the window keeps the gate intact. + Redis.set_ex("stopped:#{did}", @day_seconds, "1") + StopDebouncer.schedule(did) + + status == 1 -> + handle_resume(did) + + true -> + :ok + end + + :ok + end + + # ---- liked resolution (handler.ts lines 82-98) ---- + + defp resolve_liked(data, _did, _sha256, cached_likes) when is_binary(cached_likes) do + liked = Jason.decode!(cached_likes)["liked"] + {Map.put(data, "liked", liked), liked} + end + + defp resolve_liked(data, did, sha256, _cached_likes) do + liked = Store.liked?(did, sha256) + Redis.set_ex("likes:#{did}:#{sha256}", 2, Jason.encode!(%{liked: liked})) + {Map.put(data, "liked", liked), liked} + end + + # ---- track metadata resolution (handler.ts lines 100-146) ---- + + defp resolve_metadata(data, did, sha256, liked, duration_ms, cached_track) + when is_binary(cached_track) do + cached = Jason.decode!(cached_track) + + data = + data + |> Map.put("album_art", cached["albumArt"]) + |> Map.put("song_uri", cached["uri"]) + |> Map.put("album_uri", cached["albumUri"]) + |> Map.put("artist_uri", cached["artistUri"]) + + duration_ms = cached["duration"] || duration_ms + write_nowplaying(did, data, sha256, liked) + {data, duration_ms} + end + + defp resolve_metadata(data, did, sha256, liked, duration_ms, _cached_track_nil) do + case Store.get_track_by_sha256(sha256) do + nil -> + {data, duration_ms} + + track -> + data = + data + |> Map.put("album_art", track.album_art) + |> Map.put("song_uri", track.uri) + |> Map.put("album_uri", track.album_uri) + |> Map.put("artist_uri", track.artist_uri) + + duration_ms = track.duration || duration_ms + + Redis.set_ex( + "track:#{sha256}", + 10, + Jason.encode!(%{ + albumArt: track.album_art, + uri: track.uri, + albumUri: track.album_uri, + artistUri: track.artist_uri, + duration: track.duration, + liked: liked + }) + ) + + write_nowplaying(did, data, sha256, liked) + {data, duration_ms} + end + end + + defp write_nowplaying(did, data, sha256, liked) do + payload = data |> Map.put("sha256", sha256) |> Map.put("liked", liked) + Redis.set_ex("nowplaying:#{did}", 3, Jason.encode!(payload)) + end + + # ---- song.changed gate (handler.ts lines 159-208) ---- + + defp maybe_emit_song_changed(did, sha256, track) do + last_song_sha = Redis.get("lastsong:#{did}") + + if last_song_sha != sha256 do + if Redis.exists("ws_lastsong:#{did}") do + StopDebouncer.cancel(did) + Redis.set_ex("lastsong:#{did}", @day_seconds, sha256) + Redis.set_ex("ws_lastsong:#{did}", @day_seconds, sha256) + Redis.del("stopped:#{did}") + + Nats.publish( + "rocksky.song.changed", + Jason.encode!(%{ + did: did, + track: %{ + name: track.title, + artist: track.artist, + album: track.album, + albumCoverUrl: track.album_art, + duration_ms: track.duration_ms, + source: track.source + } + }) + ) + end + end + + :ok + end + + # ---- status=1 resume (handler.ts lines 243-272) ---- + + defp handle_resume(did) do + if StopDebouncer.has_pending?(did) do + # Cancelled before firing — song.stopped was never published, PDS record + # still exists, ws_lastsong still set. + StopDebouncer.cancel(did) + Redis.del("stopped:#{did}") + else + # No pending timer — the 15s stop already fired: ws_lastsong was deleted and + # song.stopped published. Restore ws_lastsong from the saved value, then + # delete lastsong so the next heartbeat re-publishes song.changed. + case Redis.get("stopped:#{did}") do + nil -> + :ok + + _ -> + saved_sha = Redis.get("lastsong:#{did}") + Redis.del("stopped:#{did}") + Redis.del("lastsong:#{did}") + if saved_sha, do: Redis.set_ex("ws_lastsong:#{did}", @day_seconds, saved_sha) + end + end + end + + defp sha256_hex(str) do + :crypto.hash(:sha256, str) |> Base.encode16(case: :lower) + end +end diff --git a/remote-ws/lib/remote_ws/redis.ex b/remote-ws/lib/remote_ws/redis.ex new file mode 100644 index 00000000..8903b77e --- /dev/null +++ b/remote-ws/lib/remote_ws/redis.ex @@ -0,0 +1,19 @@ +defmodule RemoteWs.Redis do + @moduledoc """ + The subset of Redis the relay uses (mirrors ctx.redis in apps/api). Backed by + Redix in production; an in-memory map in test. Dispatches to the module in + `:remote_ws, :redis`. + """ + + @callback get(String.t()) :: String.t() | nil + @callback set_ex(String.t(), non_neg_integer(), String.t()) :: :ok + @callback del(String.t()) :: :ok + @callback exists(String.t()) :: boolean() + + defp impl, do: Application.get_env(:remote_ws, :redis, RemoteWs.Redis.Redix) + + def get(key), do: impl().get(key) + def set_ex(key, seconds, value), do: impl().set_ex(key, seconds, value) + def del(key), do: impl().del(key) + def exists(key), do: impl().exists(key) +end diff --git a/remote-ws/lib/remote_ws/redis/redix.ex b/remote-ws/lib/remote_ws/redis/redix.ex new file mode 100644 index 00000000..0690949d --- /dev/null +++ b/remote-ws/lib/remote_ws/redis/redix.ex @@ -0,0 +1,34 @@ +defmodule RemoteWs.Redis.Redix do + @moduledoc "Redix-backed Redis adapter (production)." + @behaviour RemoteWs.Redis + + @conn RemoteWs.Redix + + @impl true + def get(key) do + case Redix.command(@conn, ["GET", key]) do + {:ok, value} -> value + _ -> nil + end + end + + @impl true + def set_ex(key, seconds, value) do + Redix.command(@conn, ["SET", key, value, "EX", to_string(seconds)]) + :ok + end + + @impl true + def del(key) do + Redix.command(@conn, ["DEL", key]) + :ok + end + + @impl true + def exists(key) do + case Redix.command(@conn, ["EXISTS", key]) do + {:ok, n} when is_integer(n) -> n > 0 + _ -> false + end + end +end diff --git a/remote-ws/lib/remote_ws/repo.ex b/remote-ws/lib/remote_ws/repo.ex new file mode 100644 index 00000000..c8d5a0bb --- /dev/null +++ b/remote-ws/lib/remote_ws/repo.ex @@ -0,0 +1,5 @@ +defmodule RemoteWs.Repo do + use Ecto.Repo, + otp_app: :remote_ws, + adapter: Ecto.Adapters.Postgres +end diff --git a/remote-ws/lib/remote_ws/schema/access_token.ex b/remote-ws/lib/remote_ws/schema/access_token.ex new file mode 100644 index 00000000..0515dfd9 --- /dev/null +++ b/remote-ws/lib/remote_ws/schema/access_token.ex @@ -0,0 +1,9 @@ +defmodule RemoteWs.Schema.AccessToken do + @moduledoc "Mirrors apps/api schema/access-tokens (the `access_tokens` table)." + use Ecto.Schema + + @primary_key {:id, :string, source: :xata_id, autogenerate: false} + schema "access_tokens" do + field :jti, :string + end +end diff --git a/remote-ws/lib/remote_ws/schema/loved_track.ex b/remote-ws/lib/remote_ws/schema/loved_track.ex new file mode 100644 index 00000000..c2efe2d7 --- /dev/null +++ b/remote-ws/lib/remote_ws/schema/loved_track.ex @@ -0,0 +1,10 @@ +defmodule RemoteWs.Schema.LovedTrack do + @moduledoc "Mirrors apps/api schema/loved-tracks (the `loved_tracks` table)." + use Ecto.Schema + + @primary_key {:id, :string, source: :xata_id, autogenerate: false} + schema "loved_tracks" do + field :user_id, :string + field :track_id, :string + end +end diff --git a/remote-ws/lib/remote_ws/schema/track.ex b/remote-ws/lib/remote_ws/schema/track.ex new file mode 100644 index 00000000..00c6c413 --- /dev/null +++ b/remote-ws/lib/remote_ws/schema/track.ex @@ -0,0 +1,18 @@ +defmodule RemoteWs.Schema.Track do + @moduledoc "Mirrors apps/api schema/tracks (the `tracks` table)." + use Ecto.Schema + + @primary_key {:id, :string, source: :xata_id, autogenerate: false} + schema "tracks" do + field :title, :string + field :artist, :string + field :album, :string + field :album_artist, :string + field :album_art, :string + field :duration, :integer + field :sha256, :string + field :uri, :string + field :album_uri, :string + field :artist_uri, :string + end +end diff --git a/remote-ws/lib/remote_ws/schema/user.ex b/remote-ws/lib/remote_ws/schema/user.ex new file mode 100644 index 00000000..cd5e9747 --- /dev/null +++ b/remote-ws/lib/remote_ws/schema/user.ex @@ -0,0 +1,10 @@ +defmodule RemoteWs.Schema.User do + @moduledoc "Mirrors apps/api schema/users (the `users` table)." + use Ecto.Schema + + @primary_key {:id, :string, source: :xata_id, autogenerate: false} + schema "users" do + field :did, :string + field :handle, :string + end +end diff --git a/remote-ws/lib/remote_ws/stop_debouncer.ex b/remote-ws/lib/remote_ws/stop_debouncer.ex new file mode 100644 index 00000000..54ea553b --- /dev/null +++ b/remote-ws/lib/remote_ws/stop_debouncer.ex @@ -0,0 +1,69 @@ +defmodule RemoteWs.StopDebouncer do + @moduledoc """ + Debounces `rocksky.song.stopped` by DID — a port of the Node handler's + `pendingStop` map. A status=0 schedules a fire after `stop_debounce_ms` (15s by + default). A status=1 (or a track change) within the window cancels it, so a + paused player oscillating status 0/1 doesn't produce a PDS delete→create loop. + + When a timer fires it deletes `ws_lastsong:` from Redis and publishes + `rocksky.song.stopped` — exactly the Node timer callback. + """ + use GenServer + + @default_ms 15_000 + + # ---- API ---- + + def start_link(opts), do: GenServer.start_link(__MODULE__, opts, name: __MODULE__) + + @doc "Schedule (or reschedule) the debounced song.stopped for `did`." + def schedule(did), do: GenServer.cast(__MODULE__, {:schedule, did}) + + @doc "Cancel a pending song.stopped for `did` (no-op if none pending)." + def cancel(did), do: GenServer.cast(__MODULE__, {:cancel, did}) + + @doc "Whether a song.stopped is currently pending for `did`." + def has_pending?(did), do: GenServer.call(__MODULE__, {:has_pending?, did}) + + # ---- Server ---- + + @impl true + def init(_opts), do: {:ok, %{timers: %{}}} + + @impl true + def handle_cast({:schedule, did}, state) do + state = cancel_timer(state, did) + ref = Process.send_after(self(), {:fire, did}, debounce_ms()) + {:noreply, put_in(state.timers[did], ref)} + end + + @impl true + def handle_cast({:cancel, did}, state) do + {:noreply, cancel_timer(state, did)} + end + + @impl true + def handle_call({:has_pending?, did}, _from, state) do + {:reply, Map.has_key?(state.timers, did), state} + end + + @impl true + def handle_info({:fire, did}, state) do + RemoteWs.Redis.del("ws_lastsong:#{did}") + RemoteWs.Nats.publish("rocksky.song.stopped", Jason.encode!(%{did: did})) + {:noreply, %{state | timers: Map.delete(state.timers, did)}} + end + + defp cancel_timer(state, did) do + case Map.pop(state.timers, did) do + {nil, _} -> + state + + {ref, timers} -> + Process.cancel_timer(ref) + %{state | timers: timers} + end + end + + defp debounce_ms, do: Application.get_env(:remote_ws, :stop_debounce_ms, @default_ms) +end diff --git a/remote-ws/lib/remote_ws/store.ex b/remote-ws/lib/remote_ws/store.ex new file mode 100644 index 00000000..fad5a6d8 --- /dev/null +++ b/remote-ws/lib/remote_ws/store.ex @@ -0,0 +1,17 @@ +defmodule RemoteWs.Store do + @moduledoc """ + Read model for the relay: track metadata, like status, and access-token + existence. Backed by Ecto in production (RemoteWs.Store.Ecto); swapped for an + in-memory stub in test. Dispatches to the module in `:remote_ws, :store`. + """ + + @callback get_track_by_sha256(String.t()) :: map() | nil + @callback liked?(String.t(), String.t()) :: boolean() + @callback access_token_exists?(String.t()) :: boolean() + + defp impl, do: Application.get_env(:remote_ws, :store, RemoteWs.Store.Ecto) + + def get_track_by_sha256(sha256), do: impl().get_track_by_sha256(sha256) + def liked?(did, sha256), do: impl().liked?(did, sha256) + def access_token_exists?(jti), do: impl().access_token_exists?(jti) +end diff --git a/remote-ws/lib/remote_ws/store/ecto.ex b/remote-ws/lib/remote_ws/store/ecto.ex new file mode 100644 index 00000000..8b89edf0 --- /dev/null +++ b/remote-ws/lib/remote_ws/store/ecto.ex @@ -0,0 +1,43 @@ +defmodule RemoteWs.Store.Ecto do + @moduledoc "Ecto-backed Store — the production read model (Postgres/Xata)." + @behaviour RemoteWs.Store + + import Ecto.Query + + alias RemoteWs.Repo + alias RemoteWs.Schema.{AccessToken, LovedTrack, Track, User} + + @impl true + def get_track_by_sha256(sha256) do + from(t in Track, + where: t.sha256 == ^sha256, + select: %{ + album_art: t.album_art, + uri: t.uri, + album_uri: t.album_uri, + artist_uri: t.artist_uri, + duration: t.duration + }, + limit: 1 + ) + |> Repo.one() + end + + @impl true + def liked?(did, sha256) do + from(lt in LovedTrack, + join: t in Track, + on: lt.track_id == t.id, + join: u in User, + on: lt.user_id == u.id, + where: u.did == ^did and t.sha256 == ^sha256 + ) + |> Repo.exists?() + end + + @impl true + def access_token_exists?(jti) do + from(a in AccessToken, where: a.jti == ^jti) + |> Repo.exists?() + end +end diff --git a/remote-ws/lib/remote_ws/ws/connection.ex b/remote-ws/lib/remote_ws/ws/connection.ex new file mode 100644 index 00000000..0b66ebbc --- /dev/null +++ b/remote-ws/lib/remote_ws/ws/connection.ex @@ -0,0 +1,53 @@ +defmodule RemoteWs.Ws.Connection do + @moduledoc """ + The per-connection WebSocket process (WebSock behaviour, served by Bandit). + Mirrors one client socket in the Node handler: + + * raw `"ping"` → `"pong"` (handler.ts lines 55-57) + * every other text frame is decoded as JSON and dispatched to + RemoteWs.Ws.Handler + * `{:push, frame}` messages from other connection processes (broadcasts / + targeted commands) are forwarded to this socket + + Device registry entries are owned by this process, so a disconnect cleans them + up automatically — no explicit onClose needed. + """ + @behaviour WebSock + + alias RemoteWs.Ws.Handler + + @impl true + def init(_opts), do: {:ok, %{device_id: nil, did: nil}} + + @impl true + def handle_in({"ping", [opcode: :text]}, state) do + {:push, {:text, "pong"}, state} + end + + def handle_in({text, [opcode: :text]}, state) do + case Jason.decode(text) do + {:ok, msg} when is_map(msg) -> + {frames, new_state} = Handler.handle(msg, state) + push(frames, new_state) + + _ -> + {:ok, state} + end + end + + # Ignore binary frames. + def handle_in({_data, _opts}, state), do: {:ok, state} + + @impl true + def handle_info({:push, frame}, state) do + {:push, {:text, frame}, state} + end + + def handle_info(_other, state), do: {:ok, state} + + @impl true + def terminate(_reason, _state), do: :ok + + defp push([], state), do: {:ok, state} + defp push(frames, state), do: {:push, Enum.map(frames, &{:text, &1}), state} +end diff --git a/remote-ws/lib/remote_ws/ws/handler.ex b/remote-ws/lib/remote_ws/ws/handler.ex new file mode 100644 index 00000000..1276db10 --- /dev/null +++ b/remote-ws/lib/remote_ws/ws/handler.ex @@ -0,0 +1,120 @@ +defmodule RemoteWs.Ws.Handler do + @moduledoc """ + Dispatches a decoded inbound message to the register / command / message logic + — a port of the `onMessage` body in apps/api/src/websocket/handler.ts. + + Pure with respect to the socket: it returns `{frames, new_state}` where `frames` + is a list of JSON strings to push back to THIS connection, and performs + broadcasts to other connections via RemoteWs.Devices. `conn_state` is the + per-connection `%{device_id, did}`. + + Auth failures and unknown messages are swallowed (return `{[], state}`), + mirroring the Node handler's try/catch which logs and drops. + """ + + alias RemoteWs.{Auth, Devices, NowPlaying} + + @type state :: %{device_id: String.t() | nil, did: String.t() | nil} + + @spec handle(map(), state()) :: {[String.t()], state()} + def handle(%{"type" => "register"} = msg, state), do: register(msg, state) + def handle(%{"type" => "command"} = msg, state), do: command(msg, state) + def handle(%{"type" => "message"} = msg, state), do: device_message(msg, state) + def handle(_msg, state), do: {[], state} + + # ---- register (handler.ts lines 320-354) ---- + + defp register(%{"clientName" => client_name, "token" => token}, state) + when is_binary(client_name) do + case Auth.verify_token(token) do + {:ok, %{did: did}} when is_binary(did) -> + device_id = Ecto.UUID.generate() + Devices.register(did, device_id, client_name) + + # Announce to the user's OTHER devices. + Devices.broadcast_except( + did, + device_id, + Jason.encode!(%{ + type: "device_registered", + deviceId: device_id, + clientName: client_name + }) + ) + + reply = Jason.encode!(%{status: "registered", deviceId: device_id}) + {[reply], %{state | device_id: device_id, did: did}} + + _ -> + {[], state} + end + end + + defp register(_msg, state), do: {[], state} + + # ---- command (handler.ts lines 286-317) ---- + + defp command(%{"action" => action, "token" => token} = msg, state) do + case Auth.verify_token(token) do + {:ok, %{did: did}} when is_binary(did) -> + out = Jason.encode!(command_out(msg["type"], action, msg["args"])) + target = msg["target"] + + if is_binary(target) and Devices.send_to(did, target, out) == :ok do + :ok + else + Devices.broadcast(did, out) + end + + {[], state} + + _ -> + {[], state} + end + end + + defp command(_msg, state), do: {[], state} + + # `args` is omitted when nil, matching JSON.stringify dropping `undefined`. + defp command_out(type, action, nil), do: %{type: type, action: action} + defp command_out(type, action, args), do: %{type: type, action: action, args: args} + + # ---- device message: track / status (handler.ts lines 64-283) ---- + + defp device_message(%{"data" => data, "device_id" => device_id, "token" => token}, state) + when is_map(data) do + case Auth.verify_token(token) do + {:ok, %{did: did}} when is_binary(did) -> + data = + if data["type"] == "track" do + source = source_name(did, device_id, state.device_id) + NowPlaying.handle_track(did, data, source) + else + NowPlaying.handle_status(did, data) + data + end + + device_name = Devices.name_of(did, device_id) || Devices.name_of(did, state.device_id) + + Devices.broadcast(did, Jason.encode!(broadcast_envelope(data, device_id, device_name))) + {[], state} + + _ -> + {[], state} + end + end + + defp device_message(_msg, state), do: {[], state} + + # source for song.changed: name of the message's device, else this connection's. + defp source_name(did, device_id, own_device_id) do + Devices.name_of(did, device_id) || Devices.name_of(did, own_device_id) || "websocket" + end + + # device_name omitted when nil (JSON.stringify drops undefined). + defp broadcast_envelope(data, device_id, nil), + do: %{type: "message", data: data, device_id: device_id} + + defp broadcast_envelope(data, device_id, device_name), + do: %{type: "message", data: data, device_id: device_id, device_name: device_name} +end diff --git a/remote-ws/lib/remote_ws_web/endpoint.ex b/remote-ws/lib/remote_ws_web/endpoint.ex new file mode 100644 index 00000000..c49c00c6 --- /dev/null +++ b/remote-ws/lib/remote_ws_web/endpoint.ex @@ -0,0 +1,13 @@ +defmodule RemoteWsWeb.Endpoint do + use Phoenix.Endpoint, otp_app: :remote_ws + + plug Plug.RequestId + plug Plug.Telemetry, event_prefix: [:phoenix, :endpoint] + + plug Plug.Parsers, + parsers: [:json], + pass: ["*/*"], + json_decoder: Phoenix.json_library() + + plug RemoteWsWeb.Router +end diff --git a/remote-ws/lib/remote_ws_web/error_json.ex b/remote-ws/lib/remote_ws_web/error_json.ex new file mode 100644 index 00000000..e664b0ff --- /dev/null +++ b/remote-ws/lib/remote_ws_web/error_json.ex @@ -0,0 +1,6 @@ +defmodule RemoteWsWeb.ErrorJSON do + # Minimal error renderer referenced by the endpoint's :render_errors config. + def render(template, _assigns) do + %{errors: %{detail: Phoenix.Controller.status_message_from_template(template)}} + end +end diff --git a/remote-ws/lib/remote_ws_web/router.ex b/remote-ws/lib/remote_ws_web/router.ex new file mode 100644 index 00000000..355e0ab7 --- /dev/null +++ b/remote-ws/lib/remote_ws_web/router.ex @@ -0,0 +1,18 @@ +defmodule RemoteWsWeb.Router do + use Phoenix.Router + + pipeline :api do + plug :accepts, ["json"] + end + + # The player remote-control relay. A raw-JSON WebSocket (NOT a Phoenix + # Channel) to stay 1:1 wire-compatible with the existing Node /ws server. + scope "/", RemoteWsWeb do + get "/ws", WsController, :upgrade + end + + scope "/", RemoteWsWeb do + pipe_through :api + get "/health", WsController, :health + end +end diff --git a/remote-ws/lib/remote_ws_web/ws_controller.ex b/remote-ws/lib/remote_ws_web/ws_controller.ex new file mode 100644 index 00000000..e18fa9cd --- /dev/null +++ b/remote-ws/lib/remote_ws_web/ws_controller.ex @@ -0,0 +1,15 @@ +defmodule RemoteWsWeb.WsController do + use Phoenix.Controller, formats: [:json] + + # Upgrade the HTTP request to a WebSocket handled by RemoteWs.Ws.Connection. + # Mirrors Hono's `upgradeWebSocket(handleWebsocket)` mounted at GET /ws. + def upgrade(conn, _params) do + conn + |> WebSockAdapter.upgrade(RemoteWs.Ws.Connection, %{}, timeout: 60_000) + |> halt() + end + + def health(conn, _params) do + json(conn, %{status: "ok"}) + end +end diff --git a/remote-ws/mix.exs b/remote-ws/mix.exs new file mode 100644 index 00000000..77d8d291 --- /dev/null +++ b/remote-ws/mix.exs @@ -0,0 +1,48 @@ +defmodule RemoteWs.MixProject do + use Mix.Project + + def project do + [ + app: :remote_ws, + version: "0.1.0", + elixir: "~> 1.16", + elixirc_paths: elixirc_paths(Mix.env()), + start_permanent: Mix.env() == :prod, + aliases: aliases(), + deps: deps() + ] + end + + # Run "mix help compile.app" to learn about applications. + def application do + [ + mod: {RemoteWs.Application, []}, + extra_applications: [:logger, :crypto] + ] + end + + # Specifies which paths to compile per environment. + defp elixirc_paths(:test), do: ["lib", "test/support"] + defp elixirc_paths(_), do: ["lib"] + + defp deps do + [ + {:phoenix, "~> 1.7.14"}, + {:bandit, "~> 1.5"}, + {:websock_adapter, "~> 0.5"}, + {:ecto_sql, "~> 3.12"}, + {:postgrex, ">= 0.0.0"}, + {:redix, "~> 1.5"}, + {:gnat, "~> 1.9"}, + {:joken, "~> 2.6"}, + {:jason, "~> 1.4"} + ] + end + + defp aliases do + [ + setup: ["deps.get"], + test: ["test"] + ] + end +end diff --git a/remote-ws/mix.lock b/remote-ws/mix.lock new file mode 100644 index 00000000..224fa480 --- /dev/null +++ b/remote-ws/mix.lock @@ -0,0 +1,36 @@ +%{ + "bandit": {:hex, :bandit, "1.12.3", "23f49ae03d86365b0caff3b3a1eba5d981066b1ecafe2b19c0bef5ad36c211ce", [:mix], [{:hpax, "~> 1.0", [hex: :hpax, repo: "hexpm", optional: false]}, {:plug, "~> 1.18", [hex: :plug, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}, {:thousand_island, "~> 1.5", [hex: :thousand_island, repo: "hexpm", optional: false]}, {:websock, "~> 0.5", [hex: :websock, repo: "hexpm", optional: false]}], "hexpm", "a253ec03f391755b2126e4181ee2fee05c75b712b407aa399de1831b5088c58e"}, + "castore": {:hex, :castore, "1.0.20", "455e48f7115eca98c9f2b0e7a152b5a2e8f2a8a4f964c96e95bd31645ee5fa59", [:mix], [], "hexpm", "940eafbfd8b14bee649f083bc11b3b54ec555b54c3e4ea8213351ff6fee39c10"}, + "chacha20": {:hex, :chacha20, "1.0.4", "0359d8f9a32269271044c1b471d5cf69660c362a7c61a98f73a05ef0b5d9eb9e", [:mix], [], "hexpm", "2027f5d321ae9903f1f0da7f51b0635ad6b8819bc7fe397837930a2011bc2349"}, + "connection": {:hex, :connection, "1.1.0", "ff2a49c4b75b6fb3e674bfc5536451607270aac754ffd1bdfe175abe4a6d7a68", [:mix], [], "hexpm", "722c1eb0a418fbe91ba7bd59a47e28008a189d47e37e0e7bb85585a016b2869c"}, + "curve25519": {:hex, :curve25519, "1.0.6", "e1381598b4d16691cd569709fda6f54d7f87918ed2c19a051714a338e7e56279", [:mix], [], "hexpm", "dd714c2da20cdddca81f0b659236c8fb2119e58697326da50785b5c1bc64af9d"}, + "db_connection": {:hex, :db_connection, "2.10.2", "ae391e803a5adff104da913c2fc1c0c14a37f8b10001dcef568796e1fb7bf95c", [:mix], [{:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "510b14482330f1af6490a2fa0efd8d4f1435d1529b165647df22ac0f2df0fa93"}, + "decimal": {:hex, :decimal, "3.1.1", "430d87b04011ce6cbd4fd205be758311a81f87d552d40904abd00f015935b1d0", [:mix], [], "hexpm", "c5f25f2ced74a0587d03e6023f595db8e924c9d3922c8c8ffd9edfc4498cf1f6"}, + "ecto": {:hex, :ecto, "3.14.1", "7b740d87bdf45996aa0c2c2e081640906f10caa7ce5ba328fd294c7d49d0cc6f", [:mix], [{:decimal, "~> 3.0", [hex: :decimal, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "24b991956796700f467d0a3ef3d303138a3ef9ddddf8b98f43758ee067b20a30"}, + "ecto_sql": {:hex, :ecto_sql, "3.14.0", "06446ab8410d2f85bfbb80857ee224ab3b693700cbb38f6535d507449a627b2e", [:mix], [{:db_connection, "~> 2.9", [hex: :db_connection, repo: "hexpm", optional: false]}, {:decimal, "~> 3.0", [hex: :decimal, repo: "hexpm", optional: false]}, {:ecto, "~> 3.14.0", [hex: :ecto, repo: "hexpm", optional: false]}, {:myxql, "~> 0.8", [hex: :myxql, repo: "hexpm", optional: true]}, {:postgrex, "~> 0.19 or ~> 1.0", [hex: :postgrex, repo: "hexpm", optional: true]}, {:tds, "~> 2.1.1 or ~> 2.2", [hex: :tds, repo: "hexpm", optional: true]}, {:telemetry, "~> 0.4.0 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "f4d8d36faf294c9417b5a37ec7ac8217ee2abdef5fcf197ba690f361548d3949"}, + "ed25519": {:hex, :ed25519, "1.5.1", "c1e7d96188ff241fe9924ed06640d3e60ce5c19f9c3f1079db850f4d0ea1f7dc", [:mix], [], "hexpm", "f83f7b346edfb91f0672682dee8672d16a26bfa8a7b85ee9795cb7fb2825b8bb"}, + "equivalex": {:hex, :equivalex, "1.0.3", "170d9a82ae066e0020dfe1cf7811381669565922eb3359f6c91d7e9a1124ff74", [:mix], [], "hexpm", "46fa311adb855117d36e461b9c0ad2598f72110ad17ad73d7533c78020e045fc"}, + "gnat": {:hex, :gnat, "1.16.0", "6ad99dc29beef871c7b89f8fef9e62d4e2c1a40e7b68deddc7574abfccee2a5d", [:mix], [{:connection, "~> 1.1", [hex: :connection, repo: "hexpm", optional: false]}, {:jason, "~> 1.1", [hex: :jason, repo: "hexpm", optional: false]}, {:nimble_parsec, "~> 0.5 or ~> 1.0", [hex: :nimble_parsec, repo: "hexpm", optional: false]}, {:nkeys, "~> 0.2", [hex: :nkeys, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "812e7d4bb2b9c93787a428fb1ebda1df538e993cf2c6ca9d855b1338c348e219"}, + "hpax": {:hex, :hpax, "1.0.4", "777de5d433b0fbdc7c418159c8055910faa8047ffdb3d6b31098d2a46cd7685c", [:mix], [], "hexpm", "afc7cb142ebcc2d01ce7816190b98ce5dd49e799111b24249f3443d730f377ca"}, + "jason": {:hex, :jason, "1.4.5", "2e3a008590b0b8d7388c20293e9dcc9cf3e5d642fd2a114e4cbbb52e595d940a", [:mix], [{:decimal, "~> 1.0 or ~> 2.0 or ~> 3.0", [hex: :decimal, repo: "hexpm", optional: true]}], "hexpm", "b0c823996102bcd0239b3c2444eb00409b72f6a140c1950bc8b457d836b30684"}, + "joken": {:hex, :joken, "2.6.2", "5daaf82259ca603af4f0b065475099ada1b2b849ff140ccd37f4b6828ca6892a", [:mix], [{:jose, "~> 1.11.10", [hex: :jose, repo: "hexpm", optional: false]}], "hexpm", "5134b5b0a6e37494e46dbf9e4dad53808e5e787904b7c73972651b51cce3d72b"}, + "jose": {:hex, :jose, "1.11.12", "06e62b467b61d3726cbc19e9b5489f7549c37993de846dfb3ee8259f9ed208b3", [:mix, :rebar3], [], "hexpm", "31e92b653e9210b696765cdd885437457de1add2a9011d92f8cf63e4641bab7b"}, + "kcl": {:hex, :kcl, "1.5.1", "7284fa9445ec887a2fa77c8d0c7205a5b0156959882748e65056089882b6d388", [:mix], [{:curve25519, ">= 1.0.4", [hex: :curve25519, repo: "hexpm", optional: false]}, {:ed25519, ">= 1.5.0", [hex: :ed25519, repo: "hexpm", optional: false]}, {:poly1305, "~> 1.0", [hex: :poly1305, repo: "hexpm", optional: false]}, {:salsa20, "~> 1.0", [hex: :salsa20, repo: "hexpm", optional: false]}], "hexpm", "254cf4eaa0f103b4e687ecf49158c6e350fc32134f446ef1cf028565509953c4"}, + "mime": {:hex, :mime, "2.0.7", "b8d739037be7cd402aee1ba0306edfdef982687ee7e9859bee6198c1e7e2f128", [:mix], [], "hexpm", "6171188e399ee16023ffc5b76ce445eb6d9672e2e241d2df6050f3c771e80ccd"}, + "nimble_options": {:hex, :nimble_options, "1.1.1", "e3a492d54d85fc3fd7c5baf411d9d2852922f66e69476317787a7b2bb000a61b", [:mix], [], "hexpm", "821b2470ca9442c4b6984882fe9bb0389371b8ddec4d45a9504f00a66f650b44"}, + "nimble_parsec": {:hex, :nimble_parsec, "1.4.2", "8efba0122db06df95bfaa78f791344a89352ba04baedd3849593bfce4d0dc1c6", [:mix], [], "hexpm", "4b21398942dda052b403bbe1da991ccd03a053668d147d53fb8c4e0efe09c973"}, + "nkeys": {:hex, :nkeys, "0.3.1", "5129d5df23f6762b26b1c18506942958015f4d81174ab15bf42e9509e4a45b2c", [:mix], [{:ed25519, "~> 1.3", [hex: :ed25519, repo: "hexpm", optional: false]}, {:kcl, "~> 1.4", [hex: :kcl, repo: "hexpm", optional: false]}], "hexpm", "80d8d1d62ac9c5127ad776d8f435e5e1cc732985f6235b22e6c157808b44c108"}, + "phoenix": {:hex, :phoenix, "1.7.24", "4cb76aed6d3f03878893769020e97c4394ee95b62b2b2d6313c20f66d7d37baa", [:mix], [{:castore, ">= 0.0.0", [hex: :castore, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:phoenix_pubsub, "~> 2.1", [hex: :phoenix_pubsub, repo: "hexpm", optional: false]}, {:phoenix_template, "~> 1.0", [hex: :phoenix_template, repo: "hexpm", optional: false]}, {:phoenix_view, "~> 2.0", [hex: :phoenix_view, repo: "hexpm", optional: true]}, {:plug, "~> 1.14", [hex: :plug, repo: "hexpm", optional: false]}, {:plug_cowboy, "~> 2.7", [hex: :plug_cowboy, repo: "hexpm", optional: true]}, {:plug_crypto, "~> 1.2 or ~> 2.0", [hex: :plug_crypto, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}, {:websock_adapter, "~> 0.5.3", [hex: :websock_adapter, repo: "hexpm", optional: false]}], "hexpm", "a283a9d91517116166244fdd8ed8b405c142774dc39b4a4047d179cecca9c09f"}, + "phoenix_pubsub": {:hex, :phoenix_pubsub, "2.2.0", "ff3a5616e1bed6804de7773b92cbccfc0b0f473faf1f63d7daf1206c7aeaaa6f", [:mix], [], "hexpm", "adc313a5bf7136039f63cfd9668fde73bba0765e0614cba80c06ac9460ff3e96"}, + "phoenix_template": {:hex, :phoenix_template, "1.0.4", "e2092c132f3b5e5b2d49c96695342eb36d0ed514c5b252a77048d5969330d639", [:mix], [{:phoenix_html, "~> 2.14.2 or ~> 3.0 or ~> 4.0", [hex: :phoenix_html, repo: "hexpm", optional: true]}], "hexpm", "2c0c81f0e5c6753faf5cca2f229c9709919aba34fab866d3bc05060c9c444206"}, + "plug": {:hex, :plug, "1.20.3", "56c480c633ec2ce10140e236e15233bf576e1d323887d7c96711bd02ab5160db", [:mix], [{:mime, "~> 1.0 or ~> 2.0", [hex: :mime, repo: "hexpm", optional: false]}, {:plug_crypto, "~> 1.1.1 or ~> 1.2 or ~> 2.0", [hex: :plug_crypto, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4.3 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "be266aee1b8536ef6409d58cf39a3121319f0ec47cfa1b24024485aa0e76ad76"}, + "plug_crypto": {:hex, :plug_crypto, "2.2.0", "144014737daaf485407f5ed77daeaad74d651b216a28c87543f8cc7043f8efc8", [:mix], [], "hexpm", "83a95744ab1c75876542b6fab135fcc176280e0f301a111c1f757fddcec95d2c"}, + "poly1305": {:hex, :poly1305, "1.0.4", "7cdc8961a0a6e00a764835918cdb8ade868044026df8ef5d718708ea6cc06611", [:mix], [{:chacha20, "~> 1.0", [hex: :chacha20, repo: "hexpm", optional: false]}, {:equivalex, "~> 1.0", [hex: :equivalex, repo: "hexpm", optional: false]}], "hexpm", "e14e684661a5195e149b3139db4a1693579d4659d65bba115a307529c47dbc3b"}, + "postgrex": {:hex, :postgrex, "0.22.3", "bf65941737ee7a9adbe4a64c91080310d11703da343e8ac9188aacb9eb9f6f02", [:mix], [{:db_connection, "~> 2.9", [hex: :db_connection, repo: "hexpm", optional: false]}, {:decimal, "~> 1.5 or ~> 2.0 or ~> 3.0", [hex: :decimal, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:table, "~> 0.1.0", [hex: :table, repo: "hexpm", optional: true]}], "hexpm", "f018c13752b2b46e8d35d7e2d84c3276557cbfd880769109021a1d0ee36c1cfe"}, + "redix": {:hex, :redix, "1.6.0", "694179c7a3c71bffac8848fcbfe16f6f6e2a1e020f00a2ad01bc6484815470d1", [:mix], [{:castore, "~> 0.1.0 or ~> 1.0", [hex: :castore, repo: "hexpm", optional: true]}, {:nimble_options, "~> 0.5.0 or ~> 1.0", [hex: :nimble_options, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4.0 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "b2eccb05e02f21c0c3ca57513e6bacb4dd48e6406dadbd7ff9fbe07bd6745999"}, + "salsa20": {:hex, :salsa20, "1.0.4", "404cbea1fa8e68a41bcc834c0a2571ac175580fec01cc38cc70c0fb9ffc87e9b", [:mix], [], "hexpm", "745ddcd8cfa563ddb0fd61e7ce48d5146279a2cf7834e1da8441b369fdc58ac6"}, + "telemetry": {:hex, :telemetry, "1.4.2", "a0cb522801dffb1c49fe6e30561badffc7b6d0e180db1300df759faa22062855", [:rebar3], [], "hexpm", "928f6495066506077862c0d1646609eed891a4326bee3126ba54b60af61febb1"}, + "thousand_island": {:hex, :thousand_island, "1.5.0", "f50a213cac97262b6d5ebb85745aa2c00fec1413191e6e66834788d45425cecb", [:mix], [{:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "708923d40523e43cf99041ab37a0d4b0ec426ac6438fa3716ab23d919eaeb412"}, + "websock": {:hex, :websock, "0.5.3", "2f69a6ebe810328555b6fe5c831a851f485e303a7c8ce6c5f675abeb20ebdadc", [:mix], [], "hexpm", "6105453d7fac22c712ad66fab1d45abdf049868f253cf719b625151460b8b453"}, + "websock_adapter": {:hex, :websock_adapter, "0.5.9", "43dc3ba6d89ef5dec5b1d0a39698436a1e856d000d84bf31a3149862b01a287f", [:mix], [{:bandit, ">= 0.6.0", [hex: :bandit, repo: "hexpm", optional: true]}, {:plug, "~> 1.14", [hex: :plug, repo: "hexpm", optional: false]}, {:plug_cowboy, "~> 2.6", [hex: :plug_cowboy, repo: "hexpm", optional: true]}, {:websock, "~> 0.5", [hex: :websock, repo: "hexpm", optional: false]}], "hexpm", "5534d5c9adad3c18a0f58a9371220d75a803bf0b9a3d87e6fe072faaeed76a08"}, +} diff --git a/remote-ws/test/remote_ws/auth_test.exs b/remote-ws/test/remote_ws/auth_test.exs new file mode 100644 index 00000000..e95d094d --- /dev/null +++ b/remote-ws/test/remote_ws/auth_test.exs @@ -0,0 +1,36 @@ +defmodule RemoteWs.AuthTest do + use RemoteWs.Test.Case + alias RemoteWs.Auth + + test "verifies a valid token and extracts did" do + token = Token.sign(%{"did" => "did:plc:alice"}) + assert {:ok, %{did: "did:plc:alice"}} = Auth.verify_token(token) + end + + test "rejects a token signed with a different secret" do + signer = Joken.Signer.create("HS256", "wrong-secret") + {:ok, token} = Joken.Signer.sign(%{"did" => "x"}, signer) + assert {:error, _} = Auth.verify_token(token) + end + + test "access_token type is revoked when jti not in store" do + token = Token.sign(%{"did" => "d", "type" => "access_token", "jti" => "j1"}) + assert {:error, :revoked} = Auth.verify_token(token) + end + + test "access_token type is accepted when jti still exists" do + StoreStub.put_token("j2") + token = Token.sign(%{"did" => "d", "type" => "access_token", "jti" => "j2"}) + assert {:ok, %{did: "d"}} = Auth.verify_token(token) + end + + test "ignores expiration (mirrors jwt.verify ignoreExpiration)" do + token = Token.sign(%{"did" => "d", "exp" => 1}) + assert {:ok, %{did: "d"}} = Auth.verify_token(token) + end + + test "empty / non-string bearer is invalid" do + assert {:error, _} = Auth.verify_token("") + assert {:error, _} = Auth.verify_token(nil) + end +end diff --git a/remote-ws/test/remote_ws/devices_test.exs b/remote-ws/test/remote_ws/devices_test.exs new file mode 100644 index 00000000..f1cc2df8 --- /dev/null +++ b/remote-ws/test/remote_ws/devices_test.exs @@ -0,0 +1,49 @@ +defmodule RemoteWs.DevicesTest do + use RemoteWs.Test.Case + alias RemoteWs.Devices + + defp new_did, do: "did:plc:#{System.unique_integer([:positive])}" + + test "broadcast reaches all registered devices" do + did = new_did() + a = FakeDevice.start(did, "a", "A", self()) + b = FakeDevice.start(did, "b", "B", self()) + + Devices.broadcast(did, "hello") + + assert_receive {:pushed, ^a, "hello"} + assert_receive {:pushed, ^b, "hello"} + end + + test "send_to targets a single device" do + did = new_did() + a = FakeDevice.start(did, "a", "A", self()) + _b = FakeDevice.start(did, "b", "B", self()) + + assert :ok = Devices.send_to(did, "a", "x") + assert_receive {:pushed, ^a, "x"} + refute_receive {:pushed, _, "x"}, 100 + end + + test "send_to returns :not_found for an unknown device" do + assert :not_found = Devices.send_to(new_did(), "nope", "x") + end + + test "name_of resolves a device's client name" do + did = new_did() + FakeDevice.start(did, "cli", "Rocksky CLI", self()) + assert Devices.name_of(did, "cli") == "Rocksky CLI" + assert Devices.name_of(did, "missing") == nil + end + + test "broadcast_except skips the given device" do + did = new_did() + a = FakeDevice.start(did, "a", "A", self()) + b = FakeDevice.start(did, "b", "B", self()) + + Devices.broadcast_except(did, "a", "y") + + assert_receive {:pushed, ^b, "y"} + refute_receive {:pushed, ^a, "y"}, 100 + end +end diff --git a/remote-ws/test/remote_ws/now_playing_test.exs b/remote-ws/test/remote_ws/now_playing_test.exs new file mode 100644 index 00000000..8254e074 --- /dev/null +++ b/remote-ws/test/remote_ws/now_playing_test.exs @@ -0,0 +1,79 @@ +defmodule RemoteWs.NowPlayingTest do + use RemoteWs.Test.Case + alias RemoteWs.NowPlaying + + @did "did:plc:np" + + defp track_data, do: %{"type" => "track", "title" => "t", "artist" => "a", "album" => "al"} + + test "enriches liked + metadata and writes redis caches" do + sha = sha_of("t", "a", "al") + + StoreStub.put_track(sha, %{ + album_art: "art", + uri: "song://x", + album_uri: "al://x", + artist_uri: "ar://x", + duration: 1000 + }) + + StoreStub.set_liked(@did, sha) + + data = NowPlaying.handle_track(@did, track_data(), "CLI") + + assert data["liked"] == true + assert data["album_art"] == "art" + assert data["song_uri"] == "song://x" + assert data["album_uri"] == "al://x" + assert data["artist_uri"] == "ar://x" + assert RedisMemory.get("nowplaying:#{@did}") + assert RedisMemory.get("track:#{sha}") + assert RedisMemory.get("likes:#{@did}:#{sha}") + end + + test "uses cached like status without touching the store" do + sha = sha_of("t", "a", "al") + RedisMemory.put("likes:#{@did}:#{sha}", Jason.encode!(%{liked: true})) + # No StoreStub.set_liked → if it consulted the store it would be false. + data = NowPlaying.handle_track(@did, track_data(), "CLI") + assert data["liked"] == true + end + + test "does NOT emit song.changed when ws source is inactive" do + NowPlaying.handle_track(@did, track_data(), "CLI") + refute_receive {:nats, "rocksky.song.changed", _}, 100 + end + + test "emits song.changed when ws_lastsong active and track changed" do + sha = sha_of("t", "a", "al") + + StoreStub.put_track(sha, %{ + album_art: "art", + uri: "s", + album_uri: "al", + artist_uri: "ar", + duration: 4200 + }) + + RedisMemory.put("ws_lastsong:#{@did}", "old-sha") + + NowPlaying.handle_track(@did, track_data(), "Rocksky CLI") + + assert_receive {:nats, "rocksky.song.changed", payload} + assert payload["did"] == @did + assert payload["track"]["name"] == "t" + assert payload["track"]["source"] == "Rocksky CLI" + assert payload["track"]["duration_ms"] == 4200 + assert payload["track"]["albumCoverUrl"] == "art" + assert RedisMemory.get("lastsong:#{@did}") == sha + assert RedisMemory.get("ws_lastsong:#{@did}") == sha + end + + test "does NOT emit song.changed for the same track (sha match)" do + sha = sha_of("t", "a", "al") + RedisMemory.put("ws_lastsong:#{@did}", sha) + RedisMemory.put("lastsong:#{@did}", sha) + NowPlaying.handle_track(@did, track_data(), "CLI") + refute_receive {:nats, "rocksky.song.changed", _}, 100 + end +end diff --git a/remote-ws/test/remote_ws/status_debounce_test.exs b/remote-ws/test/remote_ws/status_debounce_test.exs new file mode 100644 index 00000000..d53d76ab --- /dev/null +++ b/remote-ws/test/remote_ws/status_debounce_test.exs @@ -0,0 +1,38 @@ +defmodule RemoteWs.StatusDebounceTest do + use RemoteWs.Test.Case + alias RemoteWs.NowPlaying + + test "status=0 with active ws source schedules song.stopped, then fires" do + did = "did:plc:stop1" + RedisMemory.put("ws_lastsong:#{did}", "sha") + + NowPlaying.handle_status(did, %{"status" => 0}) + + # Debounce is 30ms in test config; give it ample room. + assert_receive {:nats, "rocksky.song.stopped", %{"did" => ^did}}, 500 + refute RedisMemory.exists("ws_lastsong:#{did}") + end + + test "status=1 within the window cancels a pending song.stopped" do + did = "did:plc:stop2" + RedisMemory.put("ws_lastsong:#{did}", "sha") + + NowPlaying.handle_status(did, %{"status" => 0}) + NowPlaying.handle_status(did, %{"status" => 1}) + + refute_receive {:nats, "rocksky.song.stopped", _}, 200 + assert RedisMemory.exists("ws_lastsong:#{did}") + end + + test "status=0 without an active ws source does nothing" do + did = "did:plc:stop3" + NowPlaying.handle_status(did, %{"status" => 0}) + refute_receive {:nats, "rocksky.song.stopped", _}, 200 + end + + test "status persists to redis" do + did = "did:plc:stop4" + NowPlaying.handle_status(did, %{"status" => 2}) + assert RedisMemory.get("nowplaying:#{did}:status") == "2" + end +end diff --git a/remote-ws/test/remote_ws/ws/connection_test.exs b/remote-ws/test/remote_ws/ws/connection_test.exs new file mode 100644 index 00000000..e8152b64 --- /dev/null +++ b/remote-ws/test/remote_ws/ws/connection_test.exs @@ -0,0 +1,23 @@ +defmodule RemoteWs.Ws.ConnectionTest do + use RemoteWs.Test.Case + alias RemoteWs.Ws.Connection + + @state %{device_id: nil, did: nil} + + test "a ping frame yields pong" do + assert {:push, {:text, "pong"}, _} = + Connection.handle_in({"ping", [opcode: :text]}, @state) + end + + test "invalid JSON is ignored" do + assert {:ok, _} = Connection.handle_in({"{not json", [opcode: :text]}, @state) + end + + test "binary frames are ignored" do + assert {:ok, _} = Connection.handle_in({<<1, 2, 3>>, [opcode: :binary]}, @state) + end + + test "a {:push, frame} message is forwarded to the socket" do + assert {:push, {:text, "hi"}, _} = Connection.handle_info({:push, "hi"}, @state) + end +end diff --git a/remote-ws/test/remote_ws/ws/handler_test.exs b/remote-ws/test/remote_ws/ws/handler_test.exs new file mode 100644 index 00000000..f55ad005 --- /dev/null +++ b/remote-ws/test/remote_ws/ws/handler_test.exs @@ -0,0 +1,129 @@ +defmodule RemoteWs.Ws.HandlerTest do + use RemoteWs.Test.Case + alias RemoteWs.Ws.Handler + + defp new_did, do: "did:plc:#{System.unique_integer([:positive])}" + defp state, do: %{device_id: nil, did: nil} + + test "register replies 'registered' and announces to other devices" do + did = new_did() + token = Token.sign(%{"did" => did}) + other = FakeDevice.start(did, "other", "Other", self()) + + {frames, new_state} = + Handler.handle( + %{"type" => "register", "clientName" => "Rocksky CLI", "token" => token}, + state() + ) + + assert [reply] = frames + decoded = Jason.decode!(reply) + assert decoded["status"] == "registered" + assert is_binary(decoded["deviceId"]) + assert new_state.did == did + assert new_state.device_id == decoded["deviceId"] + + assert_receive {:pushed, ^other, frame} + ann = Jason.decode!(frame) + assert ann["type"] == "device_registered" + assert ann["clientName"] == "Rocksky CLI" + end + + test "command with target routes to one device; no target broadcasts to all" do + did = new_did() + token = Token.sign(%{"did" => did}) + a = FakeDevice.start(did, "a", "A", self()) + b = FakeDevice.start(did, "b", "B", self()) + st = %{device_id: "sender", did: did} + + Handler.handle( + %{"type" => "command", "action" => "pause", "target" => "a", "token" => token}, + st + ) + + assert_receive {:pushed, ^a, frame} + assert Jason.decode!(frame) == %{"type" => "command", "action" => "pause"} + refute_receive {:pushed, ^b, _}, 100 + + Handler.handle(%{"type" => "command", "action" => "next", "token" => token}, st) + assert_receive {:pushed, ^a, f1} + assert_receive {:pushed, ^b, f2} + assert Jason.decode!(f1)["action"] == "next" + assert Jason.decode!(f2)["action"] == "next" + end + + test "seek command carries its position in args" do + did = new_did() + token = Token.sign(%{"did" => did}) + a = FakeDevice.start(did, "a", "A", self()) + st = %{device_id: "sender", did: did} + + Handler.handle( + %{ + "type" => "command", + "action" => "seek", + "args" => %{"position" => 12_345}, + "token" => token + }, + st + ) + + assert_receive {:pushed, ^a, frame} + + assert Jason.decode!(frame) == %{ + "type" => "command", + "action" => "seek", + "args" => %{"position" => 12_345} + } + end + + test "device track message broadcasts enriched data with device_name" do + did = new_did() + token = Token.sign(%{"did" => did}) + sha = sha_of("t", "a", "al") + + StoreStub.put_track(sha, %{ + album_art: "art", + uri: "s", + album_uri: "al", + artist_uri: "ar", + duration: 1000 + }) + + cli = FakeDevice.start(did, "cli", "Rocksky CLI", self()) + st = %{device_id: "cli", did: did} + + data = %{ + "type" => "track", + "title" => "t", + "artist" => "a", + "album" => "al", + "is_playing" => true + } + + Handler.handle( + %{"type" => "message", "data" => data, "device_id" => "cli", "token" => token}, + st + ) + + assert_receive {:pushed, ^cli, frame} + env = Jason.decode!(frame) + assert env["type"] == "message" + assert env["device_id"] == "cli" + assert env["device_name"] == "Rocksky CLI" + assert env["data"]["album_art"] == "art" + assert env["data"]["is_playing"] == true + assert env["data"]["liked"] == false + end + + test "an invalid token is ignored (no reply, no broadcast)" do + did = new_did() + _dev = FakeDevice.start(did, "a", "A", self()) + + {frames, _st} = + Handler.handle(%{"type" => "register", "clientName" => "x", "token" => "garbage"}, state()) + + assert frames == [] + refute_receive {:pushed, _, _}, 100 + end +end diff --git a/remote-ws/test/support/case.ex b/remote-ws/test/support/case.ex new file mode 100644 index 00000000..b6593037 --- /dev/null +++ b/remote-ws/test/support/case.ex @@ -0,0 +1,29 @@ +defmodule RemoteWs.Test.Case do + @moduledoc """ + Base case for relay tests: resets the in-memory Redis/Store doubles and points + the NATS double at the current test process (so publishes arrive as + `{:nats, subject, decoded}` messages). Not async — the doubles are shared. + """ + use ExUnit.CaseTemplate + + using do + quote do + import RemoteWs.Test.Case, only: [sha_of: 3] + alias RemoteWs.Test.{FakeDevice, RedisMemory, StoreStub, Token} + end + end + + setup do + RemoteWs.Test.RedisMemory.reset() + RemoteWs.Test.StoreStub.reset() + Application.put_env(:remote_ws, :nats_test_pid, self()) + on_exit(fn -> Application.delete_env(:remote_ws, :nats_test_pid) end) + :ok + end + + @doc "Compute the now-playing sha256 the same way NowPlaying does." + def sha_of(title, artist, album) do + :crypto.hash(:sha256, String.downcase("#{title} - #{artist} - #{album}")) + |> Base.encode16(case: :lower) + end +end diff --git a/remote-ws/test/support/doubles.ex b/remote-ws/test/support/doubles.ex new file mode 100644 index 00000000..a1b73444 --- /dev/null +++ b/remote-ws/test/support/doubles.ex @@ -0,0 +1,115 @@ +defmodule RemoteWs.Test.RedisMemory do + @moduledoc "In-memory Redis double (ignores TTLs) for tests." + @behaviour RemoteWs.Redis + use Agent + + def start_link(_opts \\ []), do: Agent.start_link(fn -> %{} end, name: __MODULE__) + def reset, do: Agent.update(__MODULE__, fn _ -> %{} end) + + # Seed a key directly (test helper). + def put(key, value), do: Agent.update(__MODULE__, &Map.put(&1, key, value)) + + @impl true + def get(key), do: Agent.get(__MODULE__, &Map.get(&1, key)) + + @impl true + def set_ex(key, _seconds, value), do: Agent.update(__MODULE__, &Map.put(&1, key, value)) + + @impl true + def del(key), do: Agent.update(__MODULE__, &Map.delete(&1, key)) + + @impl true + def exists(key), do: Agent.get(__MODULE__, &Map.has_key?(&1, key)) +end + +defmodule RemoteWs.Test.StoreStub do + @moduledoc "Configurable in-memory Store double for tests." + @behaviour RemoteWs.Store + use Agent + + def start_link(_opts \\ []), + do: Agent.start_link(fn -> empty() end, name: __MODULE__) + + def reset, do: Agent.update(__MODULE__, fn _ -> empty() end) + + defp empty, do: %{tracks: %{}, likes: MapSet.new(), tokens: MapSet.new()} + + # Test seeders. + def put_track(sha256, map), do: Agent.update(__MODULE__, &put_in(&1.tracks[sha256], map)) + + def set_liked(did, sha256), + do: Agent.update(__MODULE__, &%{&1 | likes: MapSet.put(&1.likes, {did, sha256})}) + + def put_token(jti), do: Agent.update(__MODULE__, &%{&1 | tokens: MapSet.put(&1.tokens, jti)}) + + @impl true + def get_track_by_sha256(sha256), do: Agent.get(__MODULE__, &Map.get(&1.tracks, sha256)) + + @impl true + def liked?(did, sha256), do: Agent.get(__MODULE__, &MapSet.member?(&1.likes, {did, sha256})) + + @impl true + def access_token_exists?(jti), do: Agent.get(__MODULE__, &MapSet.member?(&1.tokens, jti)) +end + +defmodule RemoteWs.Test.NatsCollector do + @moduledoc """ + NATS double that forwards each publish to the test process registered in + `:remote_ws, :nats_test_pid` as `{:nats, subject, decoded_payload}`. + """ + @behaviour RemoteWs.Nats + + @impl true + def publish(subject, payload) do + case Application.get_env(:remote_ws, :nats_test_pid) do + pid when is_pid(pid) -> send(pid, {:nats, subject, Jason.decode!(payload)}) + _ -> :ok + end + + :ok + end +end + +defmodule RemoteWs.Test.Token do + @moduledoc "Mint HS256 JWTs signed with the test secret." + def sign(claims) do + secret = Application.fetch_env!(:remote_ws, :jwt_secret) + signer = Joken.Signer.create("HS256", secret) + {:ok, token} = Joken.Signer.sign(claims, signer) + token + end +end + +defmodule RemoteWs.Test.FakeDevice do + @moduledoc """ + A stand-in connection process: registers itself as a device and relays every + `{:push, frame}` it receives to `report_to` as `{:pushed, self(), frame}`. + """ + alias RemoteWs.Devices + + def start(did, device_id, name, report_to) do + pid = + spawn(fn -> + Devices.register(did, device_id, name) + send(report_to, {:registered, device_id, self()}) + loop(report_to) + end) + + receive do + {:registered, ^device_id, ^pid} -> pid + after + 1000 -> raise "FakeDevice #{device_id} failed to register" + end + end + + defp loop(report_to) do + receive do + {:push, frame} -> + send(report_to, {:pushed, self(), frame}) + loop(report_to) + + :stop -> + :ok + end + end +end diff --git a/remote-ws/test/test_helper.exs b/remote-ws/test/test_helper.exs new file mode 100644 index 00000000..15133f89 --- /dev/null +++ b/remote-ws/test/test_helper.exs @@ -0,0 +1,6 @@ +ExUnit.start() + +# Shared in-memory doubles for Redis and the read store (the app itself does not +# start Redis/Postgres/NATS in test — see config/test.exs start_externals: false). +{:ok, _} = RemoteWs.Test.RedisMemory.start_link() +{:ok, _} = RemoteWs.Test.StoreStub.start_link() diff --git a/systemd/rocksky-remote-ws.service b/systemd/rocksky-remote-ws.service new file mode 100644 index 00000000..de7e6a8b --- /dev/null +++ b/systemd/rocksky-remote-ws.service @@ -0,0 +1,27 @@ +[Unit] +Description=Rocksky Remote Control WebSocket (Elixir/Phoenix) +After=network.target +Wants=network-online.target + +[Service] +Type=simple +User=root +WorkingDirectory=/root/github/rocksky/remote-ws +# Fetch/compile prod deps on (re)deploy, then start the Phoenix endpoint that +# serves the /ws relay. Secrets (JWT_SECRET, XATA_POSTGRES_URL, REDIS_URL, +# NATS_URL, SECRET_KEY_BASE, REMOTE_WS_PORT) are injected by doppler. +ExecStartPre=/bin/bash -ic 'mix deps.get --only prod && MIX_ENV=prod mix compile' +ExecStart=/bin/bash -ic 'exec doppler run -- env MIX_ENV=prod mix phx.server' +Restart=on-failure +RestartSec=5 +StandardOutput=journal +StandardError=journal +SyslogIdentifier=rocksky-remote-ws + +# Environment +Environment=HOME=/root +Environment=USER=root +Environment=MIX_ENV=prod + +[Install] +WantedBy=multi-user.target