From c3878cdc6d0b83ca5cae8f93e615f0506f75df68 Mon Sep 17 00:00:00 2001 From: Emiliano Frascaro Date: Sun, 9 Aug 2026 19:52:47 +0200 Subject: [PATCH] Aggiunge registry nodi stream e provisioning Hetzner/lab per lo scale-out. Prepara l'architettura multi-nodo (MediaMTX+ffmpeg) con assignment URL per sessione, admin di provision/drain e provider Cloud/DNS astratti verso mltv-stream.net. Co-authored-by: Cursor --- .../admin/stream_nodes_controller.rb | 56 +++ .../api/v1/stream_sessions_controller.rb | 1 + .../app/jobs/cleanup_expired_sessions_job.rb | 2 +- backend/app/models/stream_node.rb | 38 ++ backend/app/models/stream_session.rb | 28 +- backend/app/services/mediamtx/client.rb | 7 + .../app/services/recordings/cleanup_local.rb | 2 +- .../recordings/upload_from_session.rb | 2 +- backend/app/services/sessions/create.rb | 16 +- backend/app/services/sessions/pause.rb | 4 +- backend/app/services/sessions/stop.rb | 2 +- .../app/services/streams/cloud_providers.rb | 15 + .../services/streams/cloud_providers/base.rb | 35 ++ .../streams/cloud_providers/hetzner.rb | 183 ++++++++++ .../streams/cloud_providers/local_lab.rb | 45 +++ .../streams/cloud_providers/proxmox_lab.rb | 162 +++++++++ backend/app/services/streams/dns_providers.rb | 14 + .../services/streams/dns_providers/base.rb | 21 ++ .../services/streams/dns_providers/hetzner.rb | 105 ++++++ .../app/services/streams/dns_providers/lab.rb | 47 +++ .../app/services/streams/node_provisioner.rb | 152 ++++++++ backend/app/services/streams/node_registry.rb | 57 +++ backend/app/services/streams/youtube_relay.rb | 6 +- .../app/services/webhooks/mediamtx_handler.rb | 4 +- backend/app/services/youtube/live_pipeline.rb | 2 +- .../views/admin/stream_nodes/index.html.erb | 54 +++ backend/app/views/layouts/admin.html.erb | 1 + backend/config/locales/admin.de.yml | 19 + backend/config/locales/admin.en.yml | 22 ++ backend/config/locales/admin.es.yml | 19 + backend/config/locales/admin.fr.yml | 19 + backend/config/locales/admin.it.yml | 22 ++ backend/config/routes.rb | 5 + .../20260809120000_create_stream_nodes.rb | 30 ++ backend/db/schema.rb | 28 +- backend/db/seeds.rb | 3 + backend/lib/tasks/streams_nodes.rake | 32 ++ .../spec/models/stream_session_node_spec.rb | 37 ++ .../streams/cloud_providers/hetzner_spec.rb | 82 +++++ .../streams/dns_providers/hetzner_spec.rb | 35 ++ .../streams/node_provisioner_cloud_spec.rb | 41 +++ .../services/streams/node_provisioner_spec.rb | 64 ++++ .../services/streams/node_registry_spec.rb | 79 +++++ docs/ARCHITECTURE.md | 8 +- .../HETZNER_CLOUD_RELAY_OVERFLOW.md | 8 +- docs/infrastructure/STREAMING_AUTOSCALE.md | 332 ++++++++++++++++++ infra/.env.production.example | 18 + infra/stream-node/cloud-init.yaml | 41 +++ 48 files changed, 1980 insertions(+), 25 deletions(-) create mode 100644 backend/app/controllers/admin/stream_nodes_controller.rb create mode 100644 backend/app/models/stream_node.rb create mode 100644 backend/app/services/streams/cloud_providers.rb create mode 100644 backend/app/services/streams/cloud_providers/base.rb create mode 100644 backend/app/services/streams/cloud_providers/hetzner.rb create mode 100644 backend/app/services/streams/cloud_providers/local_lab.rb create mode 100644 backend/app/services/streams/cloud_providers/proxmox_lab.rb create mode 100644 backend/app/services/streams/dns_providers.rb create mode 100644 backend/app/services/streams/dns_providers/base.rb create mode 100644 backend/app/services/streams/dns_providers/hetzner.rb create mode 100644 backend/app/services/streams/dns_providers/lab.rb create mode 100644 backend/app/services/streams/node_provisioner.rb create mode 100644 backend/app/services/streams/node_registry.rb create mode 100644 backend/app/views/admin/stream_nodes/index.html.erb create mode 100644 backend/db/migrate/20260809120000_create_stream_nodes.rb create mode 100644 backend/lib/tasks/streams_nodes.rake create mode 100644 backend/spec/models/stream_session_node_spec.rb create mode 100644 backend/spec/services/streams/cloud_providers/hetzner_spec.rb create mode 100644 backend/spec/services/streams/dns_providers/hetzner_spec.rb create mode 100644 backend/spec/services/streams/node_provisioner_cloud_spec.rb create mode 100644 backend/spec/services/streams/node_provisioner_spec.rb create mode 100644 backend/spec/services/streams/node_registry_spec.rb create mode 100644 docs/infrastructure/STREAMING_AUTOSCALE.md create mode 100644 infra/stream-node/cloud-init.yaml diff --git a/backend/app/controllers/admin/stream_nodes_controller.rb b/backend/app/controllers/admin/stream_nodes_controller.rb new file mode 100644 index 0000000..3a492ba --- /dev/null +++ b/backend/app/controllers/admin/stream_nodes_controller.rb @@ -0,0 +1,56 @@ +# frozen_string_literal: true + +module Admin + class StreamNodesController < Admin::BaseController + def index + Streams::NodeRegistry.ensure_home_from_env! + @nodes = StreamNode.order(:role, :slug) + @dns_provider = ENV.fetch("STREAM_DNS_PROVIDER", "lab") + @cloud_provider = ENV.fetch("STREAM_CLOUD_PROVIDER", "local_lab") + @lab_hosts = lab_dns_snippet + @hetzner_configured = ENV["HCLOUD_TOKEN"].present? + end + + def create + kind = params[:kind].to_s + node = + if kind == "cloud" + Streams::NodeProvisioner.new.provision_cloud! + else + Streams::NodeProvisioner.new.provision_lab! + end + redirect_to admin_stream_nodes_path, + notice: t("admin.flash.stream_node_created", slug: node.slug) + rescue Streams::NodeProvisioner::Error, Streams::CloudProviders::Error, Streams::DnsProviders::Error, + KeyError => e + redirect_to admin_stream_nodes_path, alert: e.message + end + + def drain + node = StreamNode.find(params[:id]) + Streams::NodeProvisioner.new.drain!(node) + redirect_to admin_stream_nodes_path, notice: t("admin.flash.stream_node_draining", slug: node.slug) + rescue Streams::NodeProvisioner::Error => e + redirect_to admin_stream_nodes_path, alert: e.message + end + + def destroy + node = StreamNode.find(params[:id]) + Streams::NodeProvisioner.new.decommission!(node) + redirect_to admin_stream_nodes_path, notice: t("admin.flash.stream_node_destroyed", slug: node.slug) + rescue Streams::NodeProvisioner::BusyError, Streams::NodeProvisioner::Error, + Streams::CloudProviders::Error, Streams::DnsProviders::Error => e + redirect_to admin_stream_nodes_path, alert: e.message + end + + private + + def lab_dns_snippet + return "" unless @dns_provider == "lab" + + Streams::DnsProviders::Lab.new.hosts_file_snippet + rescue Redis::BaseError + "" + end + end +end diff --git a/backend/app/controllers/api/v1/stream_sessions_controller.rb b/backend/app/controllers/api/v1/stream_sessions_controller.rb index 4ca1845..e1a7185 100644 --- a/backend/app/controllers/api/v1/stream_sessions_controller.rb +++ b/backend/app/controllers/api/v1/stream_sessions_controller.rb @@ -180,6 +180,7 @@ module Api platform: session.platform, rtmp_ingest_url: session.rtmp_ingest_url, hls_playback_url: session.hls_playback_url, + stream_node: session.stream_node&.slug, watch_page_url: session.matchlivetv_platform? ? session.watch_page_url : nil, share_url: session.share_url, youtube_watch_url: session.youtube_watch_url, diff --git a/backend/app/jobs/cleanup_expired_sessions_job.rb b/backend/app/jobs/cleanup_expired_sessions_job.rb index b7d521f..91addb8 100644 --- a/backend/app/jobs/cleanup_expired_sessions_job.rb +++ b/backend/app/jobs/cleanup_expired_sessions_job.rb @@ -6,7 +6,7 @@ class CleanupExpiredSessionsJob .where("updated_at < ?", 6.hours.ago) .find_each do |session| session.fail! if session.may_fail? - Mediamtx::Client.new.delete_path(session) + Mediamtx::Client.for_session(session).delete_path(session) end end end diff --git a/backend/app/models/stream_node.rb b/backend/app/models/stream_node.rb new file mode 100644 index 0000000..02954d3 --- /dev/null +++ b/backend/app/models/stream_node.rb @@ -0,0 +1,38 @@ +# frozen_string_literal: true + +# Nodo streaming (MediaMTX [+ relay ffmpeg]). Fase 0: registry + assignment URL. +class StreamNode < ApplicationRecord + ROLES = %w[home cloud lab].freeze + STATUSES = %w[provisioning ready draining offline error].freeze + PROVIDERS = %w[local proxmox_lab hetzner].freeze + + has_many :stream_sessions, dependent: :nullify + + validates :slug, presence: true, uniqueness: true + validates :hostname, presence: true + validates :role, inclusion: { in: ROLES } + validates :status, inclusion: { in: STATUSES } + validates :provider, inclusion: { in: PROVIDERS } + validates :rtmp_base_url, :hls_base_url, :api_base_url, presence: true + validates :max_publishers, :max_relays, numericality: { greater_than: 0 } + + scope :ready, -> { where(status: "ready") } + scope :allocatable, -> { ready } + + # Sessioni che occupano uno slot path MediaMTX su questo nodo. + def occupying_sessions + stream_sessions.where(status: %w[idle connecting live reconnecting paused]) + end + + def active_publishers + occupying_sessions.count + end + + def free_slots + [max_publishers - active_publishers, 0].max + end + + def allocatable? + status == "ready" && free_slots.positive? + end +end diff --git a/backend/app/models/stream_session.rb b/backend/app/models/stream_session.rb index e8ea493..861166b 100644 --- a/backend/app/models/stream_session.rb +++ b/backend/app/models/stream_session.rb @@ -7,6 +7,7 @@ class StreamSession < ApplicationRecord belongs_to :match belongs_to :user + belongs_to :stream_node, optional: true has_many :stream_events, dependent: :destroy has_one :score_state, dependent: :destroy has_many :device_states, dependent: :destroy @@ -74,7 +75,7 @@ class StreamSession < ApplicationRecord def rtmp_ingest_url # RootEncoder richiede rtmp://host:port/app/stream (due segmenti). # MediaMTX path = live/match_{uuid} (no ?token= nel path). - "#{MatchLiveTv.mediamtx_rtmp_url}/#{mediamtx_path_name}" + "#{rtmp_base_url.chomp('/')}/#{mediamtx_path_name}" end def mediamtx_path_name @@ -90,8 +91,29 @@ class StreamSession < ApplicationRecord end def hls_playback_url - base = MatchLiveTv.hls_public_url.chomp("/") - "#{base}/#{effective_hls_path_name}/index.m3u8" + "#{hls_base_url.chomp('/')}/#{effective_hls_path_name}/index.m3u8" + end + + def rtmp_base_url + stream_node&.rtmp_base_url.presence || MatchLiveTv.mediamtx_rtmp_url + end + + def hls_base_url + stream_node&.hls_base_url.presence || MatchLiveTv.hls_public_url + end + + def mediamtx_api_base_url + stream_node&.api_base_url.presence || MatchLiveTv.mediamtx_api_url + end + + def mediamtx_internal_rtmp_url + stream_node&.internal_rtmp_url.presence || + ENV.fetch("MEDIAMTX_INTERNAL_RTMP_URL", "rtmp://mediamtx:1935") + end + + def mediamtx_internal_hls_url + stream_node&.internal_hls_url.presence || + ENV.fetch("MEDIAMTX_HLS_URL", "http://mediamtx:8888") end def effective_hls_path_name diff --git a/backend/app/services/mediamtx/client.rb b/backend/app/services/mediamtx/client.rb index 60ec0ee..cc0a5d6 100644 --- a/backend/app/services/mediamtx/client.rb +++ b/backend/app/services/mediamtx/client.rb @@ -4,7 +4,12 @@ module Mediamtx class Client class Error < StandardError; end + def self.for_session(session) + new(base_url: session.mediamtx_api_base_url) + end + def initialize(base_url: MatchLiveTv.mediamtx_api_url) + @base_url = base_url @conn = Faraday.new(url: base_url) do |f| f.request :json f.response :json @@ -12,6 +17,8 @@ module Mediamtx end end + attr_reader :base_url + def create_path(session) path = session.mediamtx_path_name # record: false finché non c'è publisher — con alwaysAvailable MediaMTX registrerebbe diff --git a/backend/app/services/recordings/cleanup_local.rb b/backend/app/services/recordings/cleanup_local.rb index 993bdb7..9a0401d 100644 --- a/backend/app/services/recordings/cleanup_local.rb +++ b/backend/app/services/recordings/cleanup_local.rb @@ -95,7 +95,7 @@ module Recordings def cleanup_mediamtx_path(session) return unless session.status.in?(%w[ended error]) - Mediamtx::Client.new.delete_path(session) + Mediamtx::Client.for_session(session).delete_path(session) rescue Mediamtx::Client::Error => e @logger.warn("[Recordings::CleanupLocal] delete_path #{session.id}: #{e.message}") end diff --git a/backend/app/services/recordings/upload_from_session.rb b/backend/app/services/recordings/upload_from_session.rb index cccb92a..34c6eb0 100644 --- a/backend/app/services/recordings/upload_from_session.rb +++ b/backend/app/services/recordings/upload_from_session.rb @@ -95,7 +95,7 @@ module Recordings def cleanup_mediamtx_path Mediamtx::PublisherSync.forget_recording_state!(@session.id) - Mediamtx::Client.new.delete_path(@session) + Mediamtx::Client.for_session(@session).delete_path(@session) rescue Mediamtx::Client::Error => e Rails.logger.warn("[Recordings::UploadFromSession] delete_path: #{e.message}") end diff --git a/backend/app/services/sessions/create.rb b/backend/app/services/sessions/create.rb index f5d2768..421dae4 100644 --- a/backend/app/services/sessions/create.rb +++ b/backend/app/services/sessions/create.rb @@ -34,11 +34,19 @@ module Sessions end StreamSession.transaction do + session.stream_node = Streams::NodeRegistry.allocate! session.save! Scoring::Engine.ensure_score_for(session) - mtx = Mediamtx::Client.new - mtx.create_path(session) - log_event(session, "pairing", { created: true, platform: session.platform }) + Mediamtx::Client.for_session(session).create_path(session) + log_event( + session, + "pairing", + { + created: true, + platform: session.platform, + stream_node: session.stream_node&.slug + } + ) end if session.platform == "youtube" @@ -46,6 +54,8 @@ module Sessions end session + rescue Streams::NodeRegistry::NoCapacityError => e + raise Teams::EntitlementError.new(e.message, code: "stream_capacity_exhausted") end private diff --git a/backend/app/services/sessions/pause.rb b/backend/app/services/sessions/pause.rb index b21fe6d..b54049f 100644 --- a/backend/app/services/sessions/pause.rb +++ b/backend/app/services/sessions/pause.rb @@ -9,11 +9,11 @@ module Sessions @session.pause! if @session.may_pause? Mediamtx::PublisherSync.forget_recording_state!(@session.id) begin - Mediamtx::Client.new.set_path_recording(@session, enabled: false) + Mediamtx::Client.for_session(@session).set_path_recording(@session, enabled: false) rescue Mediamtx::Client::Error => e Rails.logger.warn("[Sessions::Pause] disable recording: #{e.message}") end - Mediamtx::Client.new.set_always_available(@session, enabled: true) + Mediamtx::Client.for_session(@session).set_always_available(@session, enabled: true) # Slate su path per HLS in pausa; RTMP telefono si ferma via comando app. log_event("paused") SessionChannel.broadcast_message(@session, { type: "command", action: "pause_stream" }) diff --git a/backend/app/services/sessions/stop.rb b/backend/app/services/sessions/stop.rb index 2217024..96c9a74 100644 --- a/backend/app/services/sessions/stop.rb +++ b/backend/app/services/sessions/stop.rb @@ -33,7 +33,7 @@ module Sessions end def remove_mediamtx_paths! - Mediamtx::Client.new.delete_path(@session) + Mediamtx::Client.for_session(@session).delete_path(@session) rescue Mediamtx::Client::Error => e Rails.logger.warn("[Sessions::Stop] delete_path #{@session.id}: #{e.message}") end diff --git a/backend/app/services/streams/cloud_providers.rb b/backend/app/services/streams/cloud_providers.rb new file mode 100644 index 0000000..751228d --- /dev/null +++ b/backend/app/services/streams/cloud_providers.rb @@ -0,0 +1,15 @@ +# frozen_string_literal: true + +module Streams + module CloudProviders + def self.build(name = ENV.fetch("STREAM_CLOUD_PROVIDER", "local_lab")) + case name.to_s + when "local_lab" then LocalLab.new + when "proxmox_lab" then ProxmoxLab.new + when "hetzner" then Hetzner.new + else + raise Error, "STREAM_CLOUD_PROVIDER sconosciuto: #{name}" + end + end + end +end diff --git a/backend/app/services/streams/cloud_providers/base.rb b/backend/app/services/streams/cloud_providers/base.rb new file mode 100644 index 0000000..82e5c66 --- /dev/null +++ b/backend/app/services/streams/cloud_providers/base.rb @@ -0,0 +1,35 @@ +# frozen_string_literal: true + +module Streams + module CloudProviders + class Error < StandardError; end + + # Descrittore restituito da create_node / list. + Instance = Struct.new( + :id, :name, :public_ip, :private_ip, :status, :raw, + keyword_init: true + ) + + class Base + def create_node(name:, labels: {}) + raise NotImplementedError + end + + def destroy_node(instance_id) + raise NotImplementedError + end + + def list_nodes(labels: {}) + raise NotImplementedError + end + + def wait_until_running(instance_id, timeout: 120) + raise NotImplementedError + end + + def public_ip(instance_id) + raise NotImplementedError + end + end + end +end diff --git a/backend/app/services/streams/cloud_providers/hetzner.rb b/backend/app/services/streams/cloud_providers/hetzner.rb new file mode 100644 index 0000000..bde3b53 --- /dev/null +++ b/backend/app/services/streams/cloud_providers/hetzner.rb @@ -0,0 +1,183 @@ +# frozen_string_literal: true + +require "faraday" + +module Streams + module CloudProviders + # Hetzner Cloud — create/destroy server per nodi stream. + # + # ENV: + # HCLOUD_TOKEN (obbligatorio) + # HCLOUD_LOCATION (default fsn1) + # HCLOUD_SERVER_TYPE (default cpx21) + # HCLOUD_IMAGE (default debian-12) + # HCLOUD_SSH_KEY (nome chiave, default matchlivetv-stream) + # HCLOUD_NETWORK_ID (opzionale, private network / WireGuard prep) + # HCLOUD_USER_DATA_FILE (opzionale, cloud-init path) + class Hetzner < Base + API = "https://api.hetzner.cloud/v1" + + def initialize(token: ENV.fetch("HCLOUD_TOKEN"), conn: nil) + @token = token + @conn = conn + end + + def create_node(name:, labels: {}) + body = { + name: name, + server_type: ENV.fetch("HCLOUD_SERVER_TYPE", "cpx21"), + image: ENV.fetch("HCLOUD_IMAGE", "debian-12"), + location: ENV.fetch("HCLOUD_LOCATION", "fsn1"), + start_after_create: true, + labels: default_labels.merge(stringify_labels(labels)), + ssh_keys: [ENV.fetch("HCLOUD_SSH_KEY", "matchlivetv-stream")], + public_net: { + enable_ipv4: true, + enable_ipv6: false + } + } + network_id = ENV["HCLOUD_NETWORK_ID"].presence + body[:networks] = [network_id.to_i] if network_id + user_data = cloud_init_user_data + body[:user_data] = user_data if user_data.present? + + data = post("/servers", body) + server = data["server"] || {} + action = data["action"] + wait_action!(action) if action + instance = wait_until_running(server["id"].to_s) + instance.name = name + instance + end + + def destroy_node(instance_id) + delete("/servers/#{instance_id}") + true + end + + def list_nodes(labels: {}) + params = {} + label_selector = labels.map { |k, v| "#{k}=#{v}" }.join(",") + params[:label_selector] = label_selector if label_selector.present? + params[:label_selector] ||= "matchlivetv=true,role=stream-node" + + data = get("/servers", params) + Array(data["servers"]).map { |s| instance_from_server(s) } + end + + def wait_until_running(instance_id, timeout: 180) + deadline = Time.now + timeout + loop do + data = get("/servers/#{instance_id}") + server = data["server"] + status = server["status"] + if status == "running" + return instance_from_server(server) + end + raise Error, "Timeout attesa server Hetzner #{instance_id} (status=#{status})" if Time.now >= deadline + + sleep 3 + end + end + + def public_ip(instance_id) + wait_until_running(instance_id).public_ip + end + + private + + def default_labels + { + "matchlivetv" => "true", + "role" => "stream-node", + "env" => ENV.fetch("STREAM_NODE_ENV", "prod") + } + end + + def stringify_labels(labels) + labels.to_h.transform_keys(&:to_s).transform_values(&:to_s) + end + + def instance_from_server(server) + public_ip = server.dig("public_net", "ipv4", "ip") + private_ip = Array(server["private_net"]).first&.dig("ip") + Instance.new( + id: server["id"].to_s, + name: server["name"], + public_ip: public_ip, + private_ip: private_ip.presence || public_ip, + status: server["status"], + raw: server + ) + end + + def cloud_init_user_data + path = ENV["HCLOUD_USER_DATA_FILE"].presence + return File.read(path) if path && File.file?(path) + + ENV["HCLOUD_USER_DATA"].presence + end + + def conn + @conn ||= Faraday.new(url: API) do |f| + f.request :json + f.response :json, content_type: /\bjson$/ + f.adapter Faraday.default_adapter + end + end + + def auth_headers + { "Authorization" => "Bearer #{@token}" } + end + + def get(path, params = {}) + response = conn.get(path) do |req| + req.headers.update(auth_headers) + req.params.update(params) + end + unwrap!(response) + end + + def post(path, body) + response = conn.post(path) do |req| + req.headers.update(auth_headers) + req.body = body + end + unwrap!(response) + end + + def delete(path) + response = conn.delete(path) do |req| + req.headers.update(auth_headers) + end + return {} if response.status == 204 + + unwrap!(response) + end + + def unwrap!(response) + unless response.success? + raise Error, "Hetzner Cloud API #{response.status}: #{response.body.inspect}" + end + + response.body.is_a?(Hash) ? response.body : {} + end + + def wait_action!(action, timeout: 120) + return unless action.is_a?(Hash) && action["id"] + + deadline = Time.now + timeout + id = action["id"] + loop do + data = get("/actions/#{id}") + status = data.dig("action", "status") + return if status == "success" + raise Error, "Hetzner action #{id} failed: #{data.inspect}" if status == "error" + raise Error, "Timeout action Hetzner #{id}" if Time.now >= deadline + + sleep 2 + end + end + end + end +end diff --git a/backend/app/services/streams/cloud_providers/local_lab.rb b/backend/app/services/streams/cloud_providers/local_lab.rb new file mode 100644 index 0000000..437c6ed --- /dev/null +++ b/backend/app/services/streams/cloud_providers/local_lab.rb @@ -0,0 +1,45 @@ +# frozen_string_literal: true + +module Streams + module CloudProviders + # Lab senza API Proxmox: simula create/destroy e riusa MediaMTX home per i path. + # Utile per testare registry, assignment e admin senza secondi host. + class LocalLab < Base + def create_node(name:, labels: {}) + Instance.new( + id: "sim-#{name}", + name: name, + public_ip: labels[:public_ip].presence || "127.0.0.1", + private_ip: labels[:private_ip].presence || "127.0.0.1", + status: "running", + raw: { simulated: true, labels: labels } + ) + end + + def destroy_node(instance_id) + true + end + + def list_nodes(labels: {}) + StreamNode.where(provider: "local", role: "lab").map do |node| + Instance.new( + id: node.provider_instance_id, + name: node.slug, + public_ip: node.metadata["public_ip"], + private_ip: node.metadata["private_ip"], + status: node.status == "ready" ? "running" : node.status, + raw: node.metadata + ) + end + end + + def wait_until_running(instance_id, timeout: 120) + Instance.new(id: instance_id, name: instance_id, status: "running") + end + + def public_ip(instance_id) + "127.0.0.1" + end + end + end +end diff --git a/backend/app/services/streams/cloud_providers/proxmox_lab.rb b/backend/app/services/streams/cloud_providers/proxmox_lab.rb new file mode 100644 index 0000000..852425b --- /dev/null +++ b/backend/app/services/streams/cloud_providers/proxmox_lab.rb @@ -0,0 +1,162 @@ +# frozen_string_literal: true + +require "faraday" + +module Streams + module CloudProviders + # Clone/start/stop di VM template su Proxmox VE (API token). + # + # ENV richiesti: + # PROXMOX_API_URL, PROXMOX_TOKEN_ID, PROXMOX_TOKEN_SECRET, + # PROXMOX_NODE, PROXMOX_TEMPLATE_VMID + class ProxmoxLab < Base + def initialize( + api_url: ENV.fetch("PROXMOX_API_URL"), + token_id: ENV.fetch("PROXMOX_TOKEN_ID"), + token_secret: ENV.fetch("PROXMOX_TOKEN_SECRET"), + node: ENV.fetch("PROXMOX_NODE"), + template_vmid: ENV.fetch("PROXMOX_TEMPLATE_VMID"), + verify_ssl: ENV.fetch("PROXMOX_VERIFY_SSL", "false") == "true" + ) + @api_url = api_url.to_s.chomp("/") + @token_id = token_id + @token_secret = token_secret + @node = node + @template_vmid = template_vmid.to_i + @verify_ssl = verify_ssl + end + + def create_node(name:, labels: {}) + newid = next_vmid + post("/nodes/#{@node}/qemu/#{@template_vmid}/clone", { + newid: newid, + name: name, + full: 1, + target: @node + }) + post("/nodes/#{@node}/qemu/#{newid}/status/start", {}) + wait_until_running(newid.to_s) + ip = public_ip(newid.to_s) + Instance.new( + id: newid.to_s, + name: name, + public_ip: ip, + private_ip: ip, + status: "running", + raw: { node: @node, vmid: newid, labels: labels } + ) + end + + def destroy_node(instance_id) + vmid = instance_id.to_i + begin + post("/nodes/#{@node}/qemu/#{vmid}/status/stop", { timeout: 30 }) + rescue Error + # già spenta + end + sleep 2 + delete("/nodes/#{@node}/qemu/#{vmid}", { purge: 1 }) + true + end + + def list_nodes(labels: {}) + items = get("/nodes/#{@node}/qemu") + Array(items).filter_map do |row| + name = row["name"].to_s + next unless name.start_with?("mltv-stream-") || name.start_with?("ingest-") + + Instance.new( + id: row["vmid"].to_s, + name: name, + public_ip: nil, + private_ip: nil, + status: row["status"], + raw: row + ) + end + end + + def wait_until_running(instance_id, timeout: 180) + deadline = Time.now + timeout + loop do + status = get("/nodes/#{@node}/qemu/#{instance_id}/status/current") + return Instance.new(id: instance_id.to_s, status: "running", raw: status) if status["status"] == "running" + raise Error, "Timeout attesa VM #{instance_id}" if Time.now >= deadline + + sleep 3 + end + end + + def public_ip(instance_id) + agent = get("/nodes/#{@node}/qemu/#{instance_id}/agent/network-get-interfaces") + interfaces = agent.is_a?(Hash) ? agent["result"] : nil + Array(interfaces).each do |iface| + Array(iface["ip-addresses"]).each do |addr| + ip = addr["ip-address"].to_s + next if ip.blank? || ip.start_with?("127.") || ip.include?(":") + + return ip + end + end + ENV["STREAM_LAB_FALLBACK_IP"].presence || "127.0.0.1" + rescue Error + ENV["STREAM_LAB_FALLBACK_IP"].presence || "127.0.0.1" + end + + private + + def next_vmid + used = Array(get("/cluster/resources", type: "vm")).map { |r| r["vmid"].to_i } + candidate = ENV.fetch("PROXMOX_VMID_START", "9100").to_i + candidate += 1 while used.include?(candidate) + candidate + end + + def conn + @conn ||= Faraday.new(url: "#{@api_url}/api2/json") do |f| + f.request :url_encoded + f.response :json, content_type: /\bjson$/ + f.adapter Faraday.default_adapter + f.ssl[:verify] = @verify_ssl + end + end + + def auth_headers + { "Authorization" => "PVEAPIToken=#{@token_id}=#{@token_secret}" } + end + + def get(path, params = {}) + response = conn.get(path) do |req| + req.headers.update(auth_headers) + req.params.update(params) + end + unwrap!(response) + end + + def post(path, body = {}) + response = conn.post(path) do |req| + req.headers.update(auth_headers) + req.body = body + end + unwrap!(response) + end + + def delete(path, params = {}) + response = conn.delete(path) do |req| + req.headers.update(auth_headers) + req.params.update(params) + end + unwrap!(response) + end + + def unwrap!(response) + unless response.success? + raise Error, "Proxmox API #{response.status}: #{response.body.inspect}" + end + + body = response.body + body.is_a?(Hash) && body.key?("data") ? body["data"] : body + end + end + end +end diff --git a/backend/app/services/streams/dns_providers.rb b/backend/app/services/streams/dns_providers.rb new file mode 100644 index 0000000..8c0cc75 --- /dev/null +++ b/backend/app/services/streams/dns_providers.rb @@ -0,0 +1,14 @@ +# frozen_string_literal: true + +module Streams + module DnsProviders + def self.build(name = ENV.fetch("STREAM_DNS_PROVIDER", "lab")) + case name.to_s + when "lab" then Lab.new + when "hetzner" then Hetzner.new + else + raise Error, "STREAM_DNS_PROVIDER sconosciuto: #{name}" + end + end + end +end diff --git a/backend/app/services/streams/dns_providers/base.rb b/backend/app/services/streams/dns_providers/base.rb new file mode 100644 index 0000000..0d84f54 --- /dev/null +++ b/backend/app/services/streams/dns_providers/base.rb @@ -0,0 +1,21 @@ +# frozen_string_literal: true + +module Streams + module DnsProviders + class Error < StandardError; end + + class Base + def upsert_a(name, ip) + raise NotImplementedError + end + + def delete_a(name) + raise NotImplementedError + end + + def resolve(name) + raise NotImplementedError + end + end + end +end diff --git a/backend/app/services/streams/dns_providers/hetzner.rb b/backend/app/services/streams/dns_providers/hetzner.rb new file mode 100644 index 0000000..9a5a673 --- /dev/null +++ b/backend/app/services/streams/dns_providers/hetzner.rb @@ -0,0 +1,105 @@ +# frozen_string_literal: true + +require "faraday" +require "cgi" + +module Streams + module DnsProviders + # Hetzner Cloud DNS (Console) — zone mltv-stream.net. + # + # ENV: + # HCLOUD_TOKEN (stesso del Cloud) + # STREAM_DNS_ZONE (default mltv-stream.net) + # STREAM_DNS_TTL (default 60) + class Hetzner < Base + API = "https://api.hetzner.cloud/v1" + + def initialize(token: ENV.fetch("HCLOUD_TOKEN"), zone: nil, conn: nil) + @token = token + @zone = zone || ENV.fetch("STREAM_DNS_ZONE", "mltv-stream.net") + @ttl = ENV.fetch("STREAM_DNS_TTL", "60").to_i + @conn = conn + end + + def upsert_a(name, ip) + rr_name = relative_name(name) + delete_a(name) + post("/zones/#{CGI.escape(@zone)}/rrsets", { + name: rr_name, + type: "A", + ttl: @ttl, + records: [{ value: ip.to_s, comment: "matchlivetv stream-node" }], + labels: { "matchlivetv" => "true", "role" => "stream-node" } + }) + true + end + + def delete_a(name) + rr_name = relative_name(name) + encoded = CGI.escape(rr_name) + response = conn.delete("/zones/#{CGI.escape(@zone)}/rrsets/#{encoded}/A") do |req| + req.headers.update(auth_headers) + end + return true if response.status == 404 || response.status == 204 || response.success? + + raise Error, "Hetzner DNS API #{response.status}: #{response.body.inspect}" + end + + def resolve(name) + rr_name = relative_name(name) + data = get("/zones/#{CGI.escape(@zone)}/rrsets", name: rr_name, type: "A") + rrset = Array(data["rrsets"]).first + Array(rrset&.dig("records")).first&.dig("value") + rescue Error + nil + end + + private + + def relative_name(name) + host = name.to_s.strip.downcase.delete_suffix(".") + suffix = ".#{@zone}" + return "@" if host == @zone + return host.delete_suffix(suffix) if host.end_with?(suffix) + + host + end + + def conn + @conn ||= Faraday.new(url: API) do |f| + f.request :json + f.response :json, content_type: /\bjson$/ + f.adapter Faraday.default_adapter + end + end + + def auth_headers + { "Authorization" => "Bearer #{@token}" } + end + + def get(path, params = {}) + response = conn.get(path) do |req| + req.headers.update(auth_headers) + req.params.update(params) + end + unwrap!(response) + end + + def post(path, body) + response = conn.post(path) do |req| + req.headers.update(auth_headers) + req.body = body + end + unwrap!(response) + end + + def unwrap!(response) + unless response.success? + raise Error, "Hetzner DNS API #{response.status}: #{response.body.inspect}" + end + + response.body.is_a?(Hash) ? response.body : {} + end + end + end +end diff --git a/backend/app/services/streams/dns_providers/lab.rb b/backend/app/services/streams/dns_providers/lab.rb new file mode 100644 index 0000000..1fe8d5d --- /dev/null +++ b/backend/app/services/streams/dns_providers/lab.rb @@ -0,0 +1,47 @@ +# frozen_string_literal: true + +module Streams + module DnsProviders + # DNS lab in Redis (e dump hosts). Nessuna chiamata al registrar. + class Lab < Base + REDIS_KEY = "stream_dns:a_records" + + def initialize(redis: nil) + @redis = redis + end + + def upsert_a(name, ip) + host = normalize(name) + redis.hset(REDIS_KEY, host, ip.to_s) + true + end + + def delete_a(name) + redis.hdel(REDIS_KEY, normalize(name)) + true + end + + def resolve(name) + redis.hget(REDIS_KEY, normalize(name)) + end + + def all_records + redis.hgetall(REDIS_KEY) + end + + def hosts_file_snippet + all_records.sort.map { |host, ip| "#{ip}\t#{host}" }.join("\n") + end + + private + + def normalize(name) + name.to_s.strip.downcase.delete_suffix(".") + end + + def redis + @redis ||= Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0")) + end + end + end +end diff --git a/backend/app/services/streams/node_provisioner.rb b/backend/app/services/streams/node_provisioner.rb new file mode 100644 index 0000000..92c9060 --- /dev/null +++ b/backend/app/services/streams/node_provisioner.rb @@ -0,0 +1,152 @@ +# frozen_string_literal: true + +module Streams + # Provisiona / decommissiona nodi stream (lab o cloud) e aggiorna DNS + registry. + class NodeProvisioner + class Error < StandardError; end + class BusyError < Error; end + + LAB_DNS_SUFFIX = -> { ENV.fetch("STREAM_LAB_DNS_SUFFIX", "lab.mltv-stream.net") } + CLOUD_DNS_SUFFIX = -> { ENV.fetch("STREAM_CLOUD_DNS_SUFFIX", ENV.fetch("STREAM_DNS_ZONE", "mltv-stream.net")) } + + def initialize(cloud: nil, dns: nil) + @cloud = cloud + @dns = dns + end + + def provision_lab!(prefix: "ingest-lab") + provision!( + prefix: prefix, + role: "lab", + dns_suffix: LAB_DNS_SUFFIX.call, + cloud: cloud_provider(ENV.fetch("STREAM_CLOUD_PROVIDER", "local_lab")), + dns: dns_provider(ENV.fetch("STREAM_DNS_PROVIDER", "lab")), + max: ENV.fetch("STREAM_LAB_MAX_PUBLISHERS", "2").to_i, + use_node_hostname: ENV["STREAM_LAB_USE_NODE_HOSTNAME"] == "1" + ) + end + + def provision_cloud!(prefix: "ingest") + provision!( + prefix: prefix, + role: "cloud", + dns_suffix: CLOUD_DNS_SUFFIX.call, + cloud: cloud_provider("hetzner"), + dns: dns_provider("hetzner"), + max: ENV.fetch("STREAM_CLOUD_MAX_PUBLISHERS", "4").to_i, + use_node_hostname: true + ) + end + + def decommission!(node) + raise Error, "Non si può decommissionare il nodo home" if node.slug == Streams::NodeRegistry::HOME_SLUG + if node.occupying_sessions.exists? + raise BusyError, "Nodo #{node.slug} ha ancora sessioni attive" + end + + node.update!(status: "draining") + cloud = cloud_for_node(node) + dns = dns_for_node(node) + cloud.destroy_node(node.provider_instance_id) if node.provider_instance_id.present? + dns.delete_a(node.hostname) if node.hostname.present? + node.destroy! + true + end + + def drain!(node) + raise Error, "Non si può mettere in drain il nodo home" if node.slug == Streams::NodeRegistry::HOME_SLUG + + node.update!(status: "draining") + node + end + + private + + def provision!(prefix:, role:, dns_suffix:, cloud:, dns:, max:, use_node_hostname:) + Streams::NodeRegistry.ensure_home_from_env! + home = StreamNode.find_by!(slug: Streams::NodeRegistry::HOME_SLUG) + slug = next_slug(prefix) + hostname = "#{slug}.#{dns_suffix}" + + instance = cloud.create_node( + name: "mltv-stream-#{slug}", + labels: { role: "stream-node", env: role == "cloud" ? "prod" : "lab" } + ) + ip = instance.public_ip.presence || "127.0.0.1" + private_ip = instance.private_ip.presence || ip + dns.upsert_a(hostname, ip) + + simulated = instance.raw.is_a?(Hash) && (instance.raw[:simulated] || instance.raw["simulated"]) + api_base = simulated ? home.api_base_url : "http://#{private_ip}:9997" + internal_rtmp = simulated ? home.internal_rtmp_url : "rtmp://#{private_ip}:1935" + internal_hls = simulated ? home.internal_hls_url : "http://#{private_ip}:8888" + + StreamNode.create!( + slug: slug, + hostname: hostname, + role: role, + status: "ready", + provider: provider_name_for(cloud, role: role), + provider_instance_id: instance.id, + rtmp_base_url: use_node_hostname ? "rtmp://#{hostname}:1935" : home.rtmp_base_url, + hls_base_url: use_node_hostname ? "https://#{hostname}/hls" : home.hls_base_url, + api_base_url: api_base, + internal_rtmp_url: internal_rtmp, + internal_hls_url: internal_hls, + max_publishers: max, + max_relays: max, + last_health_at: Time.current, + metadata: { + "public_ip" => ip, + "private_ip" => private_ip, + "simulated" => simulated, + "cloud_raw" => instance.raw + } + ) + end + + def next_slug(prefix) + used = StreamNode.where("slug LIKE ?", "#{prefix}-%").pluck(:slug) + n = 1 + loop do + candidate = format("%s-%02d", prefix, n) + return candidate unless used.include?(candidate) + + n += 1 + end + end + + def cloud_provider(name) + @cloud || Streams::CloudProviders.build(name) + end + + def dns_provider(name) + @dns || Streams::DnsProviders.build(name) + end + + def provider_name_for(cloud, role: nil) + return "hetzner" if role.to_s == "cloud" + + case cloud + when Streams::CloudProviders::ProxmoxLab then "proxmox_lab" + when Streams::CloudProviders::Hetzner then "hetzner" + else "local" + end + end + + def cloud_for_node(node) + case node.provider + when "hetzner" then Streams::CloudProviders::Hetzner.new + when "proxmox_lab" then Streams::CloudProviders::ProxmoxLab.new + else Streams::CloudProviders::LocalLab.new + end + end + + def dns_for_node(node) + case node.provider + when "hetzner" then Streams::DnsProviders::Hetzner.new + else Streams::DnsProviders::Lab.new + end + end + end +end diff --git a/backend/app/services/streams/node_registry.rb b/backend/app/services/streams/node_registry.rb new file mode 100644 index 0000000..92bc9b2 --- /dev/null +++ b/backend/app/services/streams/node_registry.rb @@ -0,0 +1,57 @@ +# frozen_string_literal: true + +module Streams + # Assegna un StreamNode a una nuova sessione (least-loaded tra i ready). + # Garantisce il nodo "home" derivato dagli ENV MediaMTX attuali. + class NodeRegistry + class NoCapacityError < StandardError; end + + HOME_SLUG = "home" + + class << self + def ensure_home_from_env! + StreamNode.find_or_initialize_by(slug: HOME_SLUG).tap do |node| + node.assign_attributes( + hostname: home_hostname, + role: "home", + status: "ready", + provider: "local", + rtmp_base_url: MatchLiveTv.mediamtx_rtmp_url, + hls_base_url: MatchLiveTv.hls_public_url, + api_base_url: MatchLiveTv.mediamtx_api_url, + internal_rtmp_url: ENV.fetch("MEDIAMTX_INTERNAL_RTMP_URL", "rtmp://mediamtx:1935"), + internal_hls_url: ENV.fetch("MEDIAMTX_HLS_URL", "http://mediamtx:8888"), + max_publishers: ENV.fetch("STREAM_NODE_HOME_MAX_PUBLISHERS", "6").to_i, + max_relays: ENV.fetch("STREAM_NODE_HOME_MAX_RELAYS", "6").to_i + ) + node.save! + end + end + + def allocate! + ensure_home_from_env! + + node = StreamNode.ready + .to_a + .select(&:allocatable?) + .min_by { |n| [n.active_publishers, n.role == "home" ? 1 : 0, n.slug] } + + raise NoCapacityError, "Nessun nodo streaming con slot liberi" if node.nil? + + node + end + + private + + def home_hostname + ENV["STREAM_NODE_HOME_HOSTNAME"].presence || + begin + uri = URI.parse(MatchLiveTv.mediamtx_rtmp_url.sub(/\Artmps?:\/\//, "http://")) + uri.host.presence + rescue URI::InvalidURIError + nil + end || "home" + end + end + end +end diff --git a/backend/app/services/streams/youtube_relay.rb b/backend/app/services/streams/youtube_relay.rb index 3381474..a1cfec4 100644 --- a/backend/app/services/streams/youtube_relay.rb +++ b/backend/app/services/streams/youtube_relay.rb @@ -146,20 +146,20 @@ module Streams end def mediamtx_intake_source(session) - base = ENV.fetch("MEDIAMTX_INTERNAL_RTMP_URL", "rtmp://mediamtx:1935") + base = session.mediamtx_internal_rtmp_url if Mediamtx::PublisherOnline.active?(session) return [:rtmp, "#{base.chomp('/')}/#{session.mediamtx_path_name}"] end path = session.mediamtx_path_name - hls = ENV.fetch("MEDIAMTX_HLS_URL", "http://mediamtx:8888").chomp("/") + hls = session.mediamtx_internal_hls_url.chomp("/") [:hls, "#{hls}/#{path}/index.m3u8"] end def intake_available?(session) return true if Mediamtx::PublisherOnline.active?(session) - info = Mediamtx::Client.new.list_paths.find { |i| i["name"] == session.mediamtx_path_name } + info = Mediamtx::Client.for_session(session).list_paths.find { |i| i["name"] == session.mediamtx_path_name } info && (info["ready"] || info["online"] || info["available"]) rescue StandardError false diff --git a/backend/app/services/webhooks/mediamtx_handler.rb b/backend/app/services/webhooks/mediamtx_handler.rb index 705aa95..caf7290 100644 --- a/backend/app/services/webhooks/mediamtx_handler.rb +++ b/backend/app/services/webhooks/mediamtx_handler.rb @@ -68,13 +68,13 @@ module Webhooks def enable_recording(session) return unless session.match.team.entitlements.recording_enabled_for_mediamtx? - Mediamtx::Client.new.set_path_recording(session, enabled: true) + Mediamtx::Client.for_session(session).set_path_recording(session, enabled: true) rescue Mediamtx::Client::Error => e Rails.logger.warn("[MediamtxHandler] enable recording: #{e.message}") end def disable_recording(session) - Mediamtx::Client.new.set_path_recording(session, enabled: false) + Mediamtx::Client.for_session(session).set_path_recording(session, enabled: false) rescue Mediamtx::Client::Error => e Rails.logger.warn("[MediamtxHandler] disable recording: #{e.message}") end diff --git a/backend/app/services/youtube/live_pipeline.rb b/backend/app/services/youtube/live_pipeline.rb index 570f463..d323eb7 100644 --- a/backend/app/services/youtube/live_pipeline.rb +++ b/backend/app/services/youtube/live_pipeline.rb @@ -63,7 +63,7 @@ module Youtube return end - Mediamtx::Client.new.set_always_available(session, enabled: false) + Mediamtx::Client.for_session(session).set_always_available(session, enabled: false) session.go_live! if session.may_go_live? session.reconnect! if session.reconnecting? && session.may_reconnect? diff --git a/backend/app/views/admin/stream_nodes/index.html.erb b/backend/app/views/admin/stream_nodes/index.html.erb new file mode 100644 index 0000000..dd8f2a2 --- /dev/null +++ b/backend/app/views/admin/stream_nodes/index.html.erb @@ -0,0 +1,54 @@ +<% content_for :body_class, "admin-body" %> +

<%= t("admin.stream_nodes.title") %>

+

+ <%= t("admin.stream_nodes.providers", cloud: @cloud_provider, dns: @dns_provider) %> +

+ +

+ <%= button_to t("admin.stream_nodes.provision_lab"), admin_stream_nodes_path, method: :post, params: { kind: "lab" }, class: "admin-btn" %> + <% if @hetzner_configured %> + <%= button_to t("admin.stream_nodes.provision_cloud"), admin_stream_nodes_path, method: :post, params: { kind: "cloud" }, class: "admin-btn", + form: { data: { confirm: t("admin.stream_nodes.provision_cloud_confirm") } } %> + <% else %> + <%= t("admin.stream_nodes.hetzner_token_missing") %> + <% end %> +

+ + + + + + + + + + + + + + + <% @nodes.each do |node| %> + + + + + + + + + + <% end %> + +
<%= t("admin.stream_nodes.col.slug") %><%= t("admin.stream_nodes.col.role") %><%= t("admin.stream_nodes.col.status") %><%= t("admin.stream_nodes.col.slots") %><%= t("admin.stream_nodes.col.hostname") %><%= t("admin.stream_nodes.col.provider") %>
<%= node.slug %><%= node.role %><%= node.status %><%= node.active_publishers %> / <%= node.max_publishers %> (free <%= node.free_slots %>)<%= node.hostname %><%= node.provider %><% if node.provider_instance_id.present? %> · <%= node.provider_instance_id %><% end %> + <% if node.slug != "home" %> + <% if node.status == "ready" %> + <%= button_to t("admin.stream_nodes.drain"), drain_admin_stream_node_path(node), method: :post, class: "admin-btn admin-btn-secondary" %> + <% end %> + <%= button_to t("admin.stream_nodes.destroy"), admin_stream_node_path(node), method: :delete, class: "admin-btn admin-btn-danger", form: { data: { confirm: t("admin.stream_nodes.destroy_confirm", slug: node.slug) } } %> + <% end %> +
+ +<% if @lab_hosts.present? %> +

<%= t("admin.stream_nodes.hosts_title") %>

+
<%= @lab_hosts %>
+<% end %> diff --git a/backend/app/views/layouts/admin.html.erb b/backend/app/views/layouts/admin.html.erb index f4b4690..5c3bd99 100644 --- a/backend/app/views/layouts/admin.html.erb +++ b/backend/app/views/layouts/admin.html.erb @@ -29,6 +29,7 @@ <%= link_to t("admin.layout.nav.billing"), admin_billing_path, class: ("active" if controller_name.in?(%w[billing billing_invoices])) %> <%= link_to t("admin.layout.nav.youtube"), admin_youtube_platform_path, class: ("active" if controller_name == "youtube") %> <%= link_to t("admin.layout.nav.sessions"), admin_sessions_path, class: ("active" if controller_name == "sessions") %> + <%= link_to t("admin.layout.nav.stream_nodes"), admin_stream_nodes_path, class: ("active" if controller_name == "stream_nodes") %> <%= link_to t("admin.layout.nav.password"), edit_admin_password_path %> <%= button_to t("admin.layout.nav.logout"), admin_logout_path, method: :delete %> <% end %> diff --git a/backend/config/locales/admin.de.yml b/backend/config/locales/admin.de.yml index d65159c..8346fcb 100644 --- a/backend/config/locales/admin.de.yml +++ b/backend/config/locales/admin.de.yml @@ -9,6 +9,7 @@ de: billing: Abrechnung youtube: YouTube sessions: Sitzungen + stream_nodes: Stream-Knoten password: Passwort logout: Abmelden flash: @@ -26,6 +27,9 @@ de: ops_acknowledged: Vorfall übernommen ops_resolved: Vorfall gelöst ops_muted: Benachrichtigungen für 24 Stunden stummgeschaltet + stream_node_created: "Lab-Knoten %{slug} bereitgestellt." + stream_node_destroyed: "Knoten %{slug} entfernt." + stream_node_draining: "Knoten %{slug} im Drain-Modus." comped_granted: "Kostenloses Abonnement %{plan} für %{club} aktiviert." comped_revoked: "Kostenloses Abonnement für %{club} widerrufen." session_already_terminated: "Sitzung bereits beendet (%{status})." @@ -111,6 +115,21 @@ de: title: Teams matches_count: "%{count} Spiele" view_all: "Vereine ansehen (%{count} Teams)" + stream_nodes: + title: Stream-Knoten + providers: "Cloud-Provider: %{cloud} · DNS-Provider: %{dns}" + provision_lab: Lab-Knoten bereitstellen + drain: Drain + destroy: Löschen + destroy_confirm: "Knoten %{slug} löschen?" + hosts_title: Lab-DNS (/etc/hosts) + col: + slug: Slug + role: Rolle + status: Status + slots: Slots + hostname: Hostname + provider: Provider ops: kpi: critical: Kritisch offen diff --git a/backend/config/locales/admin.en.yml b/backend/config/locales/admin.en.yml index 1f8d295..4c1da92 100644 --- a/backend/config/locales/admin.en.yml +++ b/backend/config/locales/admin.en.yml @@ -9,6 +9,7 @@ en: billing: Billing youtube: YouTube sessions: Sessions + stream_nodes: Stream nodes password: Password logout: Log out flash: @@ -26,6 +27,9 @@ en: ops_acknowledged: Incident acknowledged ops_resolved: Incident resolved ops_muted: Notifications muted for 24 hours + stream_node_created: "Lab node %{slug} provisioned." + stream_node_destroyed: "Node %{slug} removed." + stream_node_draining: "Node %{slug} is draining (no new sessions)." comped_granted: "%{plan} complimentary subscription activated for %{club}." comped_revoked: "Complimentary subscription revoked for %{club}." session_already_terminated: "Session already ended (%{status})." @@ -111,6 +115,24 @@ en: title: Teams matches_count: "%{count} matches" view_all: "View clubs (%{count} teams)" + stream_nodes: + title: Streaming nodes + providers: "Cloud provider: %{cloud} · DNS provider: %{dns}" + provision_lab: Provision lab node + provision_cloud: Provision Hetzner node + provision_cloud_confirm: "Create a Hetzner Cloud server + DNS record on mltv-stream.net? Billing applies until destroyed." + hetzner_token_missing: "Set HCLOUD_TOKEN to enable Cloud provisioning." + drain: Drain + destroy: Delete + destroy_confirm: "Delete node %{slug}?" + hosts_title: Lab DNS records (/etc/hosts snippet) + col: + slug: Slug + role: Role + status: Status + slots: Slots + hostname: Hostname + provider: Provider ops: kpi: critical: Critical open diff --git a/backend/config/locales/admin.es.yml b/backend/config/locales/admin.es.yml index e1f1c9c..667d3ba 100644 --- a/backend/config/locales/admin.es.yml +++ b/backend/config/locales/admin.es.yml @@ -9,6 +9,7 @@ es: billing: Facturación youtube: YouTube sessions: Sesiones + stream_nodes: Nodos stream password: Contraseña logout: Salir flash: @@ -26,6 +27,9 @@ es: ops_acknowledged: Incidencia asumida ops_resolved: Incidencia resuelta ops_muted: Notificaciones silenciadas durante 24 horas + stream_node_created: "Nodo lab %{slug} provisionado." + stream_node_destroyed: "Nodo %{slug} eliminado." + stream_node_draining: "Nodo %{slug} en drain." comped_granted: "Suscripción de cortesía %{plan} activada para %{club}." comped_revoked: "Suscripción de cortesía revocada para %{club}." session_already_terminated: "La sesión ya ha finalizado (%{status})." @@ -111,6 +115,21 @@ es: title: Equipos matches_count: "%{count} partidos" view_all: "Ver clubes (%{count} equipos)" + stream_nodes: + title: Nodos streaming + providers: "Cloud provider: %{cloud} · DNS provider: %{dns}" + provision_lab: Provisionar nodo lab + drain: Drain + destroy: Eliminar + destroy_confirm: "¿Eliminar el nodo %{slug}?" + hosts_title: DNS lab (/etc/hosts) + col: + slug: Slug + role: Rol + status: Estado + slots: Slots + hostname: Hostname + provider: Provider ops: kpi: critical: Críticas abiertas diff --git a/backend/config/locales/admin.fr.yml b/backend/config/locales/admin.fr.yml index a9a0915..0c8ab75 100644 --- a/backend/config/locales/admin.fr.yml +++ b/backend/config/locales/admin.fr.yml @@ -9,6 +9,7 @@ fr: billing: Facturation youtube: YouTube sessions: Sessions + stream_nodes: Nœuds stream password: Mot de passe logout: Déconnexion flash: @@ -26,6 +27,9 @@ fr: ops_acknowledged: Incident pris en charge ops_resolved: Incident résolu ops_muted: Notifications suspendues pendant 24 heures + stream_node_created: "Nœud lab %{slug} provisionné." + stream_node_destroyed: "Nœud %{slug} supprimé." + stream_node_draining: "Nœud %{slug} en drain." comped_granted: "Abonnement offert %{plan} activé pour %{club}." comped_revoked: "Abonnement offert révoqué pour %{club}." session_already_terminated: "Session déjà terminée (%{status})." @@ -111,6 +115,21 @@ fr: title: Équipes matches_count: "%{count} matchs" view_all: "Voir les clubs (%{count} équipes)" + stream_nodes: + title: Nœuds streaming + providers: "Cloud provider: %{cloud} · DNS provider: %{dns}" + provision_lab: Provisionner un nœud lab + drain: Drain + destroy: Supprimer + destroy_confirm: "Supprimer le nœud %{slug} ?" + hosts_title: DNS lab (/etc/hosts) + col: + slug: Slug + role: Rôle + status: Statut + slots: Slots + hostname: Hostname + provider: Provider ops: kpi: critical: Critiques ouverts diff --git a/backend/config/locales/admin.it.yml b/backend/config/locales/admin.it.yml index ffdf74e..37fc6ed 100644 --- a/backend/config/locales/admin.it.yml +++ b/backend/config/locales/admin.it.yml @@ -9,6 +9,7 @@ it: billing: Fatturazione youtube: YouTube sessions: Sessioni + stream_nodes: Nodi stream password: Password logout: Esci flash: @@ -26,6 +27,9 @@ it: ops_acknowledged: Incidente preso in carico ops_resolved: Incidente risolto ops_muted: Notifiche sospese per 24 ore + stream_node_created: "Nodo lab %{slug} provisionato." + stream_node_destroyed: "Nodo %{slug} rimosso." + stream_node_draining: "Nodo %{slug} in drain (niente nuove sessioni)." comped_granted: "Abbonamento omaggio %{plan} attivato per %{club}." comped_revoked: "Abbonamento omaggio revocato per %{club}." session_already_terminated: "Sessione già terminata (%{status})." @@ -111,6 +115,24 @@ it: title: Squadre matches_count: "%{count} partite" view_all: "Vedi società (%{count} squadre)" + stream_nodes: + title: Nodi streaming + providers: "Cloud provider: %{cloud} · DNS provider: %{dns}" + provision_lab: Provisiona nodo lab + provision_cloud: Provisiona nodo Hetzner + provision_cloud_confirm: "Creare un server Hetzner Cloud + record DNS su mltv-stream.net? Verrà addebitato fino allo spegnimento." + hetzner_token_missing: "Imposta HCLOUD_TOKEN per abilitare il provisioning Cloud." + drain: Drain + destroy: Elimina + destroy_confirm: "Eliminare il nodo %{slug}?" + hosts_title: Record DNS lab (snippet /etc/hosts) + col: + slug: Slug + role: Ruolo + status: Stato + slots: Slot + hostname: Hostname + provider: Provider ops: kpi: critical: Critici aperti diff --git a/backend/config/routes.rb b/backend/config/routes.rb index 5b37e3a..2d09cc3 100644 --- a/backend/config/routes.rb +++ b/backend/config/routes.rb @@ -110,6 +110,11 @@ Rails.application.routes.draw do post :regia_link end end + resources :stream_nodes, only: %i[index create destroy] do + member do + post :drain + end + end get "youtube/platform", to: "youtube#platform", as: :youtube_platform end diff --git a/backend/db/migrate/20260809120000_create_stream_nodes.rb b/backend/db/migrate/20260809120000_create_stream_nodes.rb new file mode 100644 index 0000000..d0a016a --- /dev/null +++ b/backend/db/migrate/20260809120000_create_stream_nodes.rb @@ -0,0 +1,30 @@ +# frozen_string_literal: true + +class CreateStreamNodes < ActiveRecord::Migration[7.2] + def change + create_table :stream_nodes, id: :uuid, default: -> { "gen_random_uuid()" } do |t| + t.string :slug, null: false + t.string :hostname, null: false + t.string :role, null: false, default: "home" + t.string :status, null: false, default: "ready" + t.string :rtmp_base_url, null: false + t.string :hls_base_url, null: false + t.string :api_base_url, null: false + t.string :internal_rtmp_url + t.string :internal_hls_url + t.integer :max_publishers, null: false, default: 6 + t.integer :max_relays, null: false, default: 6 + t.string :provider, null: false, default: "local" + t.string :provider_instance_id + t.datetime :last_health_at + t.jsonb :metadata, null: false, default: {} + t.timestamps + end + + add_index :stream_nodes, :slug, unique: true + add_index :stream_nodes, :status + add_index :stream_nodes, :role + + add_reference :stream_sessions, :stream_node, type: :uuid, foreign_key: true, null: true, index: true + end +end diff --git a/backend/db/schema.rb b/backend/db/schema.rb index 1e8e2e6..57571e4 100644 --- a/backend/db/schema.rb +++ b/backend/db/schema.rb @@ -10,7 +10,7 @@ # # It's strongly recommended that you check this file into your version control system. -ActiveRecord::Schema[7.2].define(version: 2026_06_12_120000) do +ActiveRecord::Schema[7.2].define(version: 2026_08_09_120000) do # These are extensions that must be enabled in order to support this database enable_extension "pgcrypto" enable_extension "plpgsql" @@ -255,6 +255,29 @@ ActiveRecord::Schema[7.2].define(version: 2026_06_12_120000) do t.index ["stream_session_id"], name: "index_stream_events_on_stream_session_id" end + create_table "stream_nodes", id: :uuid, default: -> { "gen_random_uuid()" }, force: :cascade do |t| + t.string "slug", null: false + t.string "hostname", null: false + t.string "role", default: "home", null: false + t.string "status", default: "ready", null: false + t.string "rtmp_base_url", null: false + t.string "hls_base_url", null: false + t.string "api_base_url", null: false + t.string "internal_rtmp_url" + t.string "internal_hls_url" + t.integer "max_publishers", default: 6, null: false + t.integer "max_relays", default: 6, null: false + t.string "provider", default: "local", null: false + t.string "provider_instance_id" + t.datetime "last_health_at" + t.jsonb "metadata", default: {}, null: false + t.datetime "created_at", null: false + t.datetime "updated_at", null: false + t.index ["role"], name: "index_stream_nodes_on_role" + t.index ["slug"], name: "index_stream_nodes_on_slug", unique: true + t.index ["status"], name: "index_stream_nodes_on_status" + end + create_table "stream_sessions", id: :uuid, default: -> { "gen_random_uuid()" }, force: :cascade do |t| t.uuid "match_id", null: false t.uuid "user_id", null: false @@ -280,10 +303,12 @@ ActiveRecord::Schema[7.2].define(version: 2026_06_12_120000) do t.datetime "updated_at", null: false t.string "regia_token_digest" t.datetime "regia_token_expires_at" + t.uuid "stream_node_id" t.index ["match_id"], name: "index_stream_sessions_on_match_id" t.index ["publish_token"], name: "index_stream_sessions_on_publish_token", unique: true t.index ["regia_token_digest"], name: "index_stream_sessions_on_regia_token_digest", unique: true t.index ["status"], name: "index_stream_sessions_on_status" + t.index ["stream_node_id"], name: "index_stream_sessions_on_stream_node_id" t.index ["user_id"], name: "index_stream_sessions_on_user_id" end @@ -411,6 +436,7 @@ ActiveRecord::Schema[7.2].define(version: 2026_06_12_120000) do add_foreign_key "score_states", "stream_sessions" add_foreign_key "stream_events", "stream_sessions" add_foreign_key "stream_sessions", "matches" + add_foreign_key "stream_sessions", "stream_nodes" add_foreign_key "stream_sessions", "users" add_foreign_key "subscriptions", "admin_accounts", column: "admin_comped_by_id" add_foreign_key "subscriptions", "clubs" diff --git a/backend/db/seeds.rb b/backend/db/seeds.rb index 9ffd5bc..403ae6c 100644 --- a/backend/db/seeds.rb +++ b/backend/db/seeds.rb @@ -43,3 +43,6 @@ end puts "Seed OK: coach@matchlivetv.test / Password123" puts "Club: #{club.name}, Team: #{team.name}, Match: #{match.opponent_name}" + +home = Streams::NodeRegistry.ensure_home_from_env! +puts "Stream node home: #{home.slug} rtmp=#{home.rtmp_base_url}" diff --git a/backend/lib/tasks/streams_nodes.rake b/backend/lib/tasks/streams_nodes.rake new file mode 100644 index 0000000..fe347f4 --- /dev/null +++ b/backend/lib/tasks/streams_nodes.rake @@ -0,0 +1,32 @@ +# frozen_string_literal: true + +namespace :streams do + namespace :nodes do + desc "Assicura il nodo home dagli ENV MediaMTX" + task ensure_home: :environment do + node = Streams::NodeRegistry.ensure_home_from_env! + puts "home ready slug=#{node.slug} rtmp=#{node.rtmp_base_url} max=#{node.max_publishers}" + end + + desc "Provisiona un nodo lab (STREAM_CLOUD_PROVIDER=local_lab|proxmox_lab)" + task provision_lab: :environment do + node = Streams::NodeProvisioner.new.provision_lab! + puts "lab node ready slug=#{node.slug} host=#{node.hostname} id=#{node.provider_instance_id}" + if ENV.fetch("STREAM_DNS_PROVIDER", "lab") == "lab" + puts "DNS lab snippet:" + puts Streams::DnsProviders::Lab.new.hosts_file_snippet + end + end + + desc "Provisiona un nodo Hetzner Cloud + DNS mltv-stream.net (richiede HCLOUD_TOKEN)" + task provision_cloud: :environment do + node = Streams::NodeProvisioner.new.provision_cloud! + puts "cloud node ready slug=#{node.slug} host=#{node.hostname} id=#{node.provider_instance_id} ip=#{node.metadata['public_ip']}" + end + + desc "Dump record DNS lab (Redis)" + task dns_lab_dump: :environment do + puts Streams::DnsProviders::Lab.new.hosts_file_snippet + end + end +end diff --git a/backend/spec/models/stream_session_node_spec.rb b/backend/spec/models/stream_session_node_spec.rb new file mode 100644 index 0000000..34fa83b --- /dev/null +++ b/backend/spec/models/stream_session_node_spec.rb @@ -0,0 +1,37 @@ +# frozen_string_literal: true + +require "rails_helper" + +RSpec.describe StreamSession do + it "builds ingest/hls URLs from the assigned stream node" do + node = StreamNode.create!( + slug: "ingest-01", + hostname: "ingest-01.mltv-stream.net", + role: "cloud", + status: "ready", + provider: "hetzner", + rtmp_base_url: "rtmp://ingest-01.mltv-stream.net:1935", + hls_base_url: "https://ingest-01.mltv-stream.net/hls", + api_base_url: "http://10.0.0.2:9997", + internal_rtmp_url: "rtmp://10.0.0.2:1935", + max_publishers: 4, + max_relays: 4 + ) + user = User.create!(email: "s@example.com", name: "S", password: "Password123", role: "coach") + club = Club.create!(name: "Club", sport: "volleyball") + team = club.teams.create!(name: "Team", sport: "volleyball", slug: "team-url") + match = team.matches.create!(opponent_name: "Opp", scheduled_at: 1.hour.from_now) + session = StreamSession.create!( + match: match, user: user, platform: "matchlivetv", status: "idle", stream_node: node + ) + + expect(session.rtmp_ingest_url).to eq( + "rtmp://ingest-01.mltv-stream.net:1935/live/match_#{session.id}" + ) + expect(session.hls_playback_url).to eq( + "https://ingest-01.mltv-stream.net/hls/live/match_#{session.id}/index.m3u8" + ) + expect(session.mediamtx_api_base_url).to eq("http://10.0.0.2:9997") + expect(session.mediamtx_internal_rtmp_url).to eq("rtmp://10.0.0.2:1935") + end +end diff --git a/backend/spec/services/streams/cloud_providers/hetzner_spec.rb b/backend/spec/services/streams/cloud_providers/hetzner_spec.rb new file mode 100644 index 0000000..53c7d52 --- /dev/null +++ b/backend/spec/services/streams/cloud_providers/hetzner_spec.rb @@ -0,0 +1,82 @@ +# frozen_string_literal: true + +require "rails_helper" + +RSpec.describe Streams::CloudProviders::Hetzner do + let(:stubs) { Faraday::Adapter::Test::Stubs.new } + let(:conn) do + Faraday.new(url: "https://api.hetzner.cloud/v1") do |f| + f.request :json + f.response :json + f.adapter :test, stubs + end + end + let(:provider) { described_class.new(token: "test-token", conn: conn) } + + def with_env(vars) + previous = vars.keys.index_with { |k| ENV[k] } + vars.each { |k, v| ENV[k] = v } + yield + ensure + previous.each { |k, v| v.nil? ? ENV.delete(k) : ENV[k] = v } + end + + after { stubs.verify_stubbed_calls } + + it "creates a server and waits until running" do + with_env( + "HCLOUD_LOCATION" => "fsn1", + "HCLOUD_SERVER_TYPE" => "cpx21", + "HCLOUD_IMAGE" => "debian-12", + "HCLOUD_SSH_KEY" => "matchlivetv-stream" + ) do + stubs.post("/servers") do |env| + body = JSON.parse(env.body) + expect(body["name"]).to eq("mltv-stream-ingest-01") + expect(body["location"]).to eq("fsn1") + [ + 201, + { "Content-Type" => "application/json" }, + { + "server" => { + "id" => 42, + "name" => "mltv-stream-ingest-01", + "status" => "initializing", + "public_net" => { "ipv4" => { "ip" => "49.13.1.2" } }, + "private_net" => [] + }, + "action" => { "id" => 7, "status" => "running" } + } + ] + end + stubs.get("/actions/7") do + [200, { "Content-Type" => "application/json" }, { "action" => { "id" => 7, "status" => "success" } }] + end + stubs.get("/servers/42") do + [ + 200, + { "Content-Type" => "application/json" }, + { + "server" => { + "id" => 42, + "name" => "mltv-stream-ingest-01", + "status" => "running", + "public_net" => { "ipv4" => { "ip" => "49.13.1.2" } }, + "private_net" => [] + } + } + ] + end + + instance = provider.create_node(name: "mltv-stream-ingest-01", labels: { role: "stream-node" }) + expect(instance.id).to eq("42") + expect(instance.public_ip).to eq("49.13.1.2") + expect(instance.status).to eq("running") + end + end + + it "destroys a server" do + stubs.delete("/servers/42") { [204, {}, ""] } + expect(provider.destroy_node("42")).to eq(true) + end +end diff --git a/backend/spec/services/streams/dns_providers/hetzner_spec.rb b/backend/spec/services/streams/dns_providers/hetzner_spec.rb new file mode 100644 index 0000000..7b7c710 --- /dev/null +++ b/backend/spec/services/streams/dns_providers/hetzner_spec.rb @@ -0,0 +1,35 @@ +# frozen_string_literal: true + +require "rails_helper" + +RSpec.describe Streams::DnsProviders::Hetzner do + let(:stubs) { Faraday::Adapter::Test::Stubs.new } + let(:conn) do + Faraday.new(url: "https://api.hetzner.cloud/v1") do |f| + f.request :json + f.response :json + f.adapter :test, stubs + end + end + let(:provider) { described_class.new(token: "test-token", zone: "mltv-stream.net", conn: conn) } + + after { stubs.verify_stubbed_calls } + + it "upserts an A record (delete then create)" do + stubs.delete("/zones/mltv-stream.net/rrsets/ingest-01/A") { [404, {}, ""] } + stubs.post("/zones/mltv-stream.net/rrsets") do |env| + body = JSON.parse(env.body) + expect(body["name"]).to eq("ingest-01") + expect(body["type"]).to eq("A") + expect(body["records"].first["value"]).to eq("49.13.1.2") + [201, { "Content-Type" => "application/json" }, { "rrset" => { "name" => "ingest-01" } }] + end + + expect(provider.upsert_a("ingest-01.mltv-stream.net", "49.13.1.2")).to eq(true) + end + + it "deletes an A record" do + stubs.delete("/zones/mltv-stream.net/rrsets/ingest-01/A") { [204, {}, ""] } + expect(provider.delete_a("ingest-01.mltv-stream.net")).to eq(true) + end +end diff --git a/backend/spec/services/streams/node_provisioner_cloud_spec.rb b/backend/spec/services/streams/node_provisioner_cloud_spec.rb new file mode 100644 index 0000000..705ad51 --- /dev/null +++ b/backend/spec/services/streams/node_provisioner_cloud_spec.rb @@ -0,0 +1,41 @@ +# frozen_string_literal: true + +require "rails_helper" + +RSpec.describe "Streams::NodeProvisioner cloud" do + it "registers a cloud node using injected providers" do + cloud = instance_double( + Streams::CloudProviders::Hetzner, + create_node: Streams::CloudProviders::Instance.new( + id: "99", + name: "mltv-stream-ingest-01", + public_ip: "49.13.9.9", + private_ip: "10.0.0.9", + status: "running", + raw: {} + ) + ) + dns = instance_double(Streams::DnsProviders::Hetzner) + allow(dns).to receive(:upsert_a) + + ENV["MEDIAMTX_API_URL"] = "http://mtx-home:9997" + ENV["MEDIAMTX_RTMP_URL"] = "rtmp://home.example:1935" + ENV["HLS_PUBLIC_URL"] = "https://home.example/hls" + ENV["STREAM_CLOUD_DNS_SUFFIX"] = "mltv-stream.net" + ENV["STREAM_CLOUD_MAX_PUBLISHERS"] = "4" + + node = Streams::NodeProvisioner.new(cloud: cloud, dns: dns).provision_cloud! + expect(node.slug).to eq("ingest-01") + expect(node.role).to eq("cloud") + expect(node.provider).to eq("hetzner") + expect(node.hostname).to eq("ingest-01.mltv-stream.net") + expect(node.rtmp_base_url).to eq("rtmp://ingest-01.mltv-stream.net:1935") + expect(node.api_base_url).to eq("http://10.0.0.9:9997") + expect(dns).to have_received(:upsert_a).with("ingest-01.mltv-stream.net", "49.13.9.9") + ensure + %w[ + MEDIAMTX_API_URL MEDIAMTX_RTMP_URL HLS_PUBLIC_URL + STREAM_CLOUD_DNS_SUFFIX STREAM_CLOUD_MAX_PUBLISHERS + ].each { |k| ENV.delete(k) } + end +end diff --git a/backend/spec/services/streams/node_provisioner_spec.rb b/backend/spec/services/streams/node_provisioner_spec.rb new file mode 100644 index 0000000..09df63a --- /dev/null +++ b/backend/spec/services/streams/node_provisioner_spec.rb @@ -0,0 +1,64 @@ +# frozen_string_literal: true + +require "rails_helper" + +RSpec.describe Streams::NodeProvisioner do + def with_env(vars) + previous = vars.keys.index_with { |k| ENV[k] } + vars.each { |k, v| ENV[k] = v } + yield + ensure + previous.each { |k, v| v.nil? ? ENV.delete(k) : ENV[k] = v } + end + + let(:redis) { Redis.new(url: ENV.fetch("REDIS_URL", "redis://redis:6379/0")) } + + before do + redis.del(Streams::DnsProviders::Lab::REDIS_KEY) + end + + it "provisions and decommissions a simulated lab node" do + with_env( + "STREAM_CLOUD_PROVIDER" => "local_lab", + "STREAM_DNS_PROVIDER" => "lab", + "STREAM_LAB_DNS_SUFFIX" => "lab.mltv-stream.net", + "MEDIAMTX_API_URL" => "http://mtx-home:9997", + "MEDIAMTX_RTMP_URL" => "rtmp://ingest-home.example:1935", + "HLS_PUBLIC_URL" => "https://ingest-home.example/hls", + "STREAM_LAB_MAX_PUBLISHERS" => "2" + ) do + provisioner = described_class.new + node = provisioner.provision_lab! + + expect(node.slug).to eq("ingest-lab-01") + expect(node.role).to eq("lab") + expect(node.status).to eq("ready") + expect(node.hostname).to eq("ingest-lab-01.lab.mltv-stream.net") + expect(node.api_base_url).to eq("http://mtx-home:9997") + expect(Streams::DnsProviders::Lab.new.resolve(node.hostname)).to eq("127.0.0.1") + + provisioner.decommission!(node) + expect(StreamNode.find_by(slug: "ingest-lab-01")).to be_nil + expect(Streams::DnsProviders::Lab.new.resolve("ingest-lab-01.lab.mltv-stream.net")).to be_nil + end + end + + it "refuses to decommission a busy node" do + with_env( + "STREAM_CLOUD_PROVIDER" => "local_lab", + "STREAM_DNS_PROVIDER" => "lab", + "MEDIAMTX_API_URL" => "http://mtx-home:9997", + "MEDIAMTX_RTMP_URL" => "rtmp://ingest-home.example:1935", + "HLS_PUBLIC_URL" => "https://ingest-home.example/hls" + ) do + node = described_class.new.provision_lab! + user = User.create!(email: "lab@example.com", name: "L", password: "Password123", role: "coach") + club = Club.create!(name: "LabClub", sport: "volleyball") + team = club.teams.create!(name: "T", sport: "volleyball", slug: "lab-t") + match = team.matches.create!(opponent_name: "X", scheduled_at: 1.hour.from_now) + StreamSession.create!(match: match, user: user, platform: "matchlivetv", status: "live", stream_node: node) + + expect { described_class.new.decommission!(node) }.to raise_error(Streams::NodeProvisioner::BusyError) + end + end +end diff --git a/backend/spec/services/streams/node_registry_spec.rb b/backend/spec/services/streams/node_registry_spec.rb new file mode 100644 index 0000000..17041b9 --- /dev/null +++ b/backend/spec/services/streams/node_registry_spec.rb @@ -0,0 +1,79 @@ +# frozen_string_literal: true + +require "rails_helper" + +RSpec.describe Streams::NodeRegistry do + def with_env(vars) + previous = vars.keys.index_with { |k| ENV[k] } + vars.each { |k, v| ENV[k] = v } + yield + ensure + previous.each { |k, v| v.nil? ? ENV.delete(k) : ENV[k] = v } + end + + let(:home_env) do + { + "MEDIAMTX_API_URL" => "http://mtx-home:9997", + "MEDIAMTX_RTMP_URL" => "rtmp://ingest-home.example:1935", + "HLS_PUBLIC_URL" => "https://ingest-home.example/hls", + "MEDIAMTX_INTERNAL_RTMP_URL" => "rtmp://mtx-home:1935", + "MEDIAMTX_HLS_URL" => "http://mtx-home:8888", + "STREAM_NODE_HOME_MAX_PUBLISHERS" => "2" + } + end + + it "creates home node from ENV" do + with_env(home_env) do + node = described_class.ensure_home_from_env! + expect(node.slug).to eq("home") + expect(node.rtmp_base_url).to eq("rtmp://ingest-home.example:1935") + expect(node.api_base_url).to eq("http://mtx-home:9997") + expect(node.max_publishers).to eq(2) + expect(node.status).to eq("ready") + end + end + + it "allocates the least loaded ready node" do + with_env(home_env) do + home = described_class.ensure_home_from_env! + cloud = StreamNode.create!( + slug: "ingest-01", + hostname: "ingest-01.mltv-stream.net", + role: "cloud", + status: "ready", + provider: "hetzner", + rtmp_base_url: "rtmp://ingest-01.mltv-stream.net:1935", + hls_base_url: "https://ingest-01.mltv-stream.net/hls", + api_base_url: "http://10.0.0.2:9997", + max_publishers: 2, + max_relays: 2 + ) + + user = User.create!(email: "n@example.com", name: "N", password: "Password123", role: "coach") + club = Club.create!(name: "C", sport: "volleyball") + team = club.teams.create!(name: "T", sport: "volleyball", slug: "t-node") + match = team.matches.create!(opponent_name: "X", scheduled_at: 1.hour.from_now) + + StreamSession.create!( + match: match, user: user, platform: "matchlivetv", status: "live", stream_node: home + ) + + expect(described_class.allocate!).to eq(cloud) + end + end + + it "raises when no capacity remains" do + with_env(home_env.merge("STREAM_NODE_HOME_MAX_PUBLISHERS" => "1")) do + home = described_class.ensure_home_from_env! + user = User.create!(email: "full@example.com", name: "F", password: "Password123", role: "coach") + club = Club.create!(name: "Full", sport: "volleyball") + team = club.teams.create!(name: "T", sport: "volleyball", slug: "t-full") + match = team.matches.create!(opponent_name: "X", scheduled_at: 1.hour.from_now) + StreamSession.create!( + match: match, user: user, platform: "matchlivetv", status: "live", stream_node: home + ) + + expect { described_class.allocate! }.to raise_error(Streams::NodeRegistry::NoCapacityError) + end + end +end diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 0ec450b..918e362 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -2,8 +2,8 @@ Documento di riferimento per l'intero sistema: rete, container Docker, flussi video, replay, monitoraggio ops e rilasci. -**Ultimo aggiornamento:** 2026-07-29 -**Produzione:** `eminux@192.168.1.146` → `/opt/matchlivetv` +**Ultimo aggiornamento:** 2026-08-09 +**Produzione:** `eminux@192.168.1.146` → `/opt/matchlivetv` (Proxmox; control plane attuale per autoscale) --- @@ -397,7 +397,8 @@ cd infra && cp .env.example .env && docker compose up -d --build | Documento | Contenuto | |-----------|-----------| | [`infrastructure/SERVER_DEPLOYMENT.md`](infrastructure/SERVER_DEPLOYMENT.md) | Bootstrap server, NPM, cron, backup | -| [`infrastructure/HETZNER_CLOUD_RELAY_OVERFLOW.md`](infrastructure/HETZNER_CLOUD_RELAY_OVERFLOW.md) | Piano overflow relay YouTube su Hetzner Cloud (picchi) | +| [`infrastructure/STREAMING_AUTOSCALE.md`](infrastructure/STREAMING_AUTOSCALE.md) | **Design attuale:** Proxmox prod + Hetzner Cloud warm spare (MediaMTX+ffmpeg), lab, DNS, disco; Auction in futuro | +| [`infrastructure/HETZNER_CLOUD_RELAY_OVERFLOW.md`](infrastructure/HETZNER_CLOUD_RELAY_OVERFLOW.md) | Piano storico overflow-only relay YouTube (superseded dal doc sopra) | | [`ANDROID_APP_LINKS.md`](ANDROID_APP_LINKS.md) | Digital Asset Links / deep link Play (`assetlinks.json`) | | [`LIVE_STREAMING.md`](LIVE_STREAMING.md) | MediaMTX, pausa, HLS, ruolo ffmpeg | | [`REPLAY_MODULE.md`](REPLAY_MODULE.md) | Garage, retention, YouTube VOD | @@ -420,3 +421,4 @@ cd infra && cp .env.example .env && docker compose up -d --build | 2026-06 | Cleanup path MediaMTX orfani (`Mediamtx::CleanupOrphanPaths`); delete path sempre a fine diretta | | 2026-07 | Rimosso overlay server (`OverlayRelay`); tabellone bruciato in app; HLS su path camera (non `*_air`); `YoutubeRelay` = copy video + AAC | | 2026-08 | `YoutubeRelay` = remux copy video+audio (niente ricodifica AAC); profilo app/slate AAC 48k mono | +| 2026-08 | Design autoscale: Proxmox prod + Hetzner Cloud warm spare (MediaMTX+ffmpeg), lab locale, secondo dominio DNS, Garage ora / B2 dopo; Auction rimandato (`STREAMING_AUTOSCALE.md`) | diff --git a/docs/infrastructure/HETZNER_CLOUD_RELAY_OVERFLOW.md b/docs/infrastructure/HETZNER_CLOUD_RELAY_OVERFLOW.md index 2d59937..6871346 100644 --- a/docs/infrastructure/HETZNER_CLOUD_RELAY_OVERFLOW.md +++ b/docs/infrastructure/HETZNER_CLOUD_RELAY_OVERFLOW.md @@ -1,10 +1,12 @@ # Piano: overflow relay YouTube su Hetzner Cloud -**Stato:** piano / design — non implementato -**Ultimo aggiornamento:** 2026-07-29 +> **Superseded (2026-08-09):** il design corrente è [`STREAMING_AUTOSCALE.md`](STREAMING_AUTOSCALE.md) (Auction + Proxmox, scale di MediaMTX **e** ffmpeg, warm spare, secondo dominio DNS, Garage ora / B2 dopo). Questo file resta come riferimento storico sul solo overflow dei relay. + +**Stato:** piano / design storico — non implementato come doc standalone +**Ultimo aggiornamento:** 2026-07-29 (header superseded 2026-08-09) **Contesto:** produzione attuale su host dedicato (es. Proxmox/Hetzner) con MediaMTX + Sidekiq sullo stesso box; obiettivo >10 live YouTube contemporanee senza abbandonare il middle-tier. -Documenti correlati: [`LIVE_STREAMING.md`](../LIVE_STREAMING.md), [`ARCHITECTURE.md`](../ARCHITECTURE.md), [`OPS_MONITORING.md`](../OPS_MONITORING.md). +Documenti correlati: [`STREAMING_AUTOSCALE.md`](STREAMING_AUTOSCALE.md), [`LIVE_STREAMING.md`](../LIVE_STREAMING.md), [`ARCHITECTURE.md`](../ARCHITECTURE.md), [`OPS_MONITORING.md`](../OPS_MONITORING.md). --- diff --git a/docs/infrastructure/STREAMING_AUTOSCALE.md b/docs/infrastructure/STREAMING_AUTOSCALE.md new file mode 100644 index 0000000..dabaa5d --- /dev/null +++ b/docs/infrastructure/STREAMING_AUTOSCALE.md @@ -0,0 +1,332 @@ +# Autoscale streaming — Proxmox attuale + Hetzner Cloud + +**Stato:** design approvato (decisioni chiuse) — pronto per implementazione +**Branch:** `feature/streaming-autoscale-hetzner` +**Ultimo aggiornamento:** 2026-08-09 (dominio ingest: `mltv-stream.net`) + +**Sostituisce / estende:** il piano overflow-only [`HETZNER_CLOUD_RELAY_OVERFLOW.md`](HETZNER_CLOUD_RELAY_OVERFLOW.md). Questo documento copre **MediaMTX + ffmpeg**, warm spare, DNS, storage, disco e lab sul Proxmox di produzione. + +Documenti correlati: [`ARCHITECTURE.md`](../ARCHITECTURE.md), [`LIVE_STREAMING.md`](../LIVE_STREAMING.md), [`REPLAY_MODULE.md`](../REPLAY_MODULE.md), [`OPS_MONITORING.md`](../OPS_MONITORING.md). + +--- + +## 1. Obiettivo + +Scalare molte dirette contemporanee (sito e/o YouTube) senza saturare un unico box: + +- **Control plane (ora):** Proxmox **già in produzione** (`192.168.1.146` / `/opt/matchlivetv`) — sito, API, DB, Redis, Garage, MediaMTX/ffmpeg “home”. +- **Data plane elastico:** **Hetzner Cloud** — nodi `MediaMTX + ffmpeg` con warm spare (+1). +- **DNS ingest:** secondo dominio su **Hetzner DNS** (API). +- **Lab:** stesso Proxmox, con `ProxmoxLabProvider`, prima di affidarsi al Cloud in produzione. +- **Futuro (non bloccante):** migrazione control plane su **Hetzner Auction + Proxmox** (vSwitch nativo, hardware dedicato). + +--- + +## 2. Decisioni chiuse + +| # | Decisione | Scelta | +|---|----------|--------| +| D1 | Host primario **ora** | **Proxmox già usato in produzione** (non attendere Auction) | +| D1b | Host primario **futuro** | Hetzner Server Auction + Proxmox (quando si migra) | +| D2 | Overflow / warm spare | **Hetzner Cloud** (attivato per questo progetto) | +| D3 | Unità di scala | Nodo **stream** = MediaMTX + worker `youtube_relay` (ffmpeg remux copy) | +| D4 | Warm pool | Almeno **1 spare ready** quando il carico si avvicina al cap | +| D5 | DNS ingest | Dominio **`mltv-stream.net`** (registrar Aruba, NS → Hetzner DNS). `matchlivetv.it` resta su Aruba intatto | +| D6 | Aruba | Nessuna delega NS; niente Floating IP prenotate | +| D7a | Rete **ora** (casa/Proxmox ↔ Cloud) | **WireGuard** (o Tailscale) privato: Redis, DB, API MediaMTX home. Non esporre Redis/Postgres su WAN | +| D7b | Rete **futuro** (Auction ↔ Cloud) | vSwitch Robot + Cloud Network (stessa location) | +| D8 | Object storage | **Garage resta**. **B2** = migrazione pronta quando misurato | +| D9 | Staging | Segmenti su disco VM `stream-*` in live; upload Garage a fine sessione | +| D10 | Prodotto | Consigliare **YouTube** per alleggerire storage/egress replay | +| D11 | Disco | Separazione volumi + **monitoraggio attivo** obbligatorio (§6) | +| D12 | Lab prima del Cloud prod | `ProxmoxLabProvider` + DNS lab sullo stesso Proxmox; poi `HetznerCloudProvider` | + +--- + +## 3. Topologia — fase attuale (Proxmox prod + Cloud) + +```text +Internet + │ HTTPS :443 · RTMP :1935 + ▼ +┌──────────────────────────────────────────────────────────┐ +│ Proxmox produzione (attuale) │ +│ eminux@192.168.1.146 · /opt/matchlivetv │ +│ │ +│ edge · rails · sidekiq-core · postgres · redis │ +│ garage · stream-home (MediaMTX + youtube_relay) │ +│ autoscaler (Rails/Sidekiq) │ +│ │ +│ Lab (opz.): clone VM stream-* via ProxmoxLabProvider │ +└────────────────────────┬─────────────────────────────────┘ + │ WireGuard / tunnel privato + ▼ +┌──────────────────────────────────────────────────────────┐ +│ Hetzner Cloud — pool stream nodes │ +│ stream-01 … stream-N MediaMTX + ffmpeg │ +│ + 1 spare ready (warm) │ +└────────────────────────┬─────────────────────────────────┘ + │ + Hetzner DNS (secondo dominio) + ingest-XX. → IP pubblico nodo +``` + +**Futuro Auction:** stesso schema, tunnel sostituito da vSwitch; control plane migrato sul bare metal Hetzner. + +**Principio:** Rails solo controllo. Live sui MediaMTX del nodo assegnato. Replay su Garage (poi eventualmente B2). + +--- + +## 4. DNS (secondo dominio) + +RTMP **non** tollera round-robin su un unico hostname. + +1. Dominio **`mltv-stream.net`** registrato su Aruba; zona DNS su Hetzner Console (progetto `matchlivetv-stream`). +2. NS Aruba del dominio → `hydrogen` / `oxygen` / `helium` Hetzner (delegazione OK). +3. Autoscaler a VM ready → upsert `A` `ingest-03.mltv-stream.net` → health-check → nodo `ready`. +4. API create/start sessione → URL del nodo assegnato: + +```text +rtmp://ingest-03.mltv-stream.net:1935/live/match_ +https://ingest-03.mltv-stream.net/hls/... # o via edge matchlivetv.it che proxya il nodo +``` + +5. Scale-in: drain → delete DNS → destroy VM. + +Hostname pubblici dei nodi vivono sotto `*.mltv-stream.net`. In lab: `LabDnsProvider` (`/etc/hosts`, dnsmasq). + +--- + +## 5. Autoscaler e warm spare + +Ogni nodo: `max_publishers` / `max_relays` (es. 4). + +| Regola | Comportamento | +|--------|----------------| +| Scale-out | `free_slots` sotto soglia soft → crea 1 nodo | +| Warm spare | Se `spare_ready < 1` vicino al carico → crea spare | +| Scale-in | Idle da X min e non unica spare → destroy | +| Floor | `stream-home` sul Proxmox sempre on | + +`CloudProvider`: `HetznerCloudProvider` (prod overflow) + `ProxmoxLabProvider` (test). +Code Sidekiq: coda `youtube_relay` isolata; sticky owner Redis. + +--- + +## 6. Disco Proxmox — requisito critico + +**Non** mettere OS, DB e video sullo stesso volume che può riempirsi al 100%. + +Vale **già** sul Proxmox di produzione (rinforzare ora), e di nuovo alla migrazione Auction. + +### 6.1 Layout volumi (target) + +| Volume / disco | Contenuto | Note | +|----------------|-----------|------| +| `os` | Proxmox / sistema host | Quasi statico | +| `vm-system` / root stack | rails, edge, redis, … | Snapshot | +| `vm-data-db` | Postgres | Separato; backup prioritari | +| `vm-data-video` | Garage + staging MediaMTX | **Disco che può riempirsi** | +| `backup` | Backup / PBS | **Mai** sullo stesso disco dei video | + +Regole: cap Garage/staging; se staging pieno → **rifiutare nuove sessioni**, non far cadere DB/API. + +### 6.2 Monitoraggio attivo (obbligatorio) + +| Metrica | Warn | Critical | +|---------|------|----------| +| Uso `%` volume video | ≥ 70% | ≥ 85% | +| Spazio libero GB | sotto soglia | < X GB | +| Inode liberi | bassi | esaurimento | +| Crescita GB/giorno Garage | anomalia | — | +| Purge/retention falliti | — | immediato | +| Staging per nodo `stream-*` | ≥ 70% | ≥ 85% + blocco live su quel nodo | + +Checklist: + +- [ ] Alert ntfy disco testati +- [ ] Vista ops “disco video” +- [ ] Test periodico “80% pieno” +- [ ] Restore di prova Postgres + +--- + +## 7. Storage: Garage ora, B2 dopo + +| Ora | **Garage** | +| Dopo | **Backblaze B2** quando metriche/€ lo giustificano | + +Client S3-ready; staging locale invariato. Spinta prodotto YouTube riduce pressione replay. + +--- + +## 8. Prodotto: YouTube consigliato + +Suggerire YouTube in app/dashboard quando ha senso → meno storage/egress MatchLiveTV; middle-tier (slate + relay) resta il valore. Live solo-sito restano supportate. + +--- + +## 9. Astrazioni codice + +```text +CloudProvider + create_node / destroy_node / list_nodes / wait_until_running / public_ip + → HetznerCloudProvider | ProxmoxLabProvider + +DnsProvider + upsert_a / delete_a + → HetznerDnsProvider | LabDnsProvider + +StreamNodeRegistry + register / mark_ready / allocate_for_session / drain / release + +Autoscaler + reconcile!(metrics) → ensure spare / scale-in idle +``` + +Tag VM Cloud: `matchlivetv`, `role=stream-node`, `env=prod|lab`. + +--- + +## 10. Lab sul Proxmox attuale + +Prima (o in parallelo) al Cloud “vero”: + +1. Template VM `stream-node` su Proxmox. +2. `CLOUD_PROVIDER=proxmox_lab` → clone/start/stop via API Proxmox. +3. DNS lab (hosts/dnsmasq). +4. Soft cap basso; N sessioni di test; verifica assignment, spare, scale-in, alert disco. +5. Poi smoke su Hetzner Cloud (1 CX piccolo) + DNS reale. + +Cosa il lab **non** replica al 100%: tempi boot Hetzner, vSwitch, RTMP 4G multi-nodo (smoke WAN dopo). + +--- + +## 11. Fasi di implementazione + +| Fase | Cosa | Esito | +|------|------|--------| +| **L — Lab Proxmox** | `LocalLab` / `ProxmoxLab`, DNS lab, admin nodi | **implementata sul branch** | +| **0 — Multi-nodo ready** | Registry + URL RTMP/HLS in API (`StreamNode`, `Streams::NodeRegistry`) | **implementata sul branch** | +| **1 — Hetzner Cloud + DNS** | `HetznerCloudProvider` + `HetznerDnsProvider`, `mltv-stream.net`, cloud-init, WireGuard (§16) | **implementata sul branch** (WG ops manuale) | +| **2 — Routing relay** | Coda `youtube_relay`, sticky owner, cap | Niente doppi ffmpeg | +| **3 — Autoscaler** | Soglie + warm spare (+1) | Automatico | +| **4 — Hardening** | Drain, budget, runbook, kill-switch | Picchi weekend | +| **A — Auction (futuro)** | Migrazione control plane + vSwitch al posto di WireGuard | Hardware dedicato Hetzner | +| **B2 (opz.)** | Switch object storage | Quando misurato | + +### Fase 0 — dettagli implementati + +- Tabella `stream_nodes` + `stream_sessions.stream_node_id` +- `Streams::NodeRegistry.ensure_home_from_env!` / `allocate!` (least-loaded) +- `Sessions::Create` assegna il nodo e crea il path MediaMTX sul client del nodo +- URL RTMP/HLS/API/internal per sessione derivati dal nodo (fallback ENV se nodo assente) +- API JSON espone `stream_node` (slug) +- Dominio ingest pubblico futuro: `*.mltv-stream.net` + +### Fase L — dettagli implementati + +- `Streams::CloudProviders` (`local_lab`, `proxmox_lab`, stub `hetzner`) +- `Streams::DnsProviders` (`lab` su Redis, stub `hetzner`) +- `Streams::NodeProvisioner` (provision / drain / decommission) +- Admin **Nodi stream** (`/admin/stream_nodes`): lista, provision lab, drain, delete +- Rake: `streams:nodes:ensure_home`, `streams:nodes:provision_lab`, `streams:nodes:dns_lab_dump` +- Default sicuro: `STREAM_CLOUD_PROVIDER=local_lab` (nessuna chiamata Proxmox finché non configuri token) + +### Fase 1 — dettagli implementati + +- `Streams::CloudProviders::Hetzner` (create/destroy server, wait running, labels) +- `Streams::DnsProviders::Hetzner` (upsert/delete A su zona `mltv-stream.net`) +- `NodeProvisioner#provision_cloud!` + bottone admin **Provisiona nodo Hetzner** +- Cloud-init: `infra/stream-node/cloud-init.yaml` (`HCLOUD_USER_DATA_FILE`) +- Decommission usa il provider del nodo (non solo ENV globale) +- WireGuard: vedi §16 + +Ordine di lavoro consigliato sul branch: **0 → L → 1 → 2 → 3 → 4**; **A** quando si decide di lasciare il Proxmox casa. + +--- + +## 12. Ordine operativo “ora” (senza Auction) + +1. Rinforzare dischi/alert sul Proxmox prod (§6). +2. Secondo dominio + zona Hetzner DNS. +3. Progetto Hetzner Cloud + immagine/snapshot `stream-node`. +4. Tunnel WireGuard Proxmox ↔ Cloud (Redis/DB/API private). +5. Implementare fasi 0 → L → 1… sul branch. +6. Cutover overflow in produzione solo dopo lab + smoke Cloud verdi. + +--- + +## 13. Criteri di successo + +| Criterio | Misura | +|----------|--------| +| Lab | Scale-out/in e assignment verificati su Proxmox senza Hetzner | +| Picco | N live oltre soft cap home con spare Cloud ready | +| Continuità | Drop telefono → slate YT attiva | +| Disco | Nessun outage DB/API per disco video; alert prima del critical | +| DNS | Aruba intatto; ingest sul secondo dominio | +| Futuro | Path chiaro verso Auction senza riscrivere autoscaler | + +--- + +## 14. Aperto / da calibrare + +| Voce | Note | +|------|------| +| Nome secondo dominio | **`mltv-stream.net`** (chiuso) | +| Soft cap `stream-home` | es. 4 vs 6 | +| Tipo VM Cloud | CPX vs CCX | +| `OVERFLOW_MAX` / budget | — | +| Spare di notte | 0 vs 1 | +| Dettaglio WireGuard | subnet, peer, MTU | +| Soglie GB disco | da size reale volume video | + +--- + +## 15. Riepilogo esecutivo + +**Ora:** restiamo sul Proxmox di produzione; attiviamo Hetzner Cloud (spare MediaMTX+ffmpeg) + Hetzner DNS (secondo dominio) + WireGuard. Testiamo prima in lab sullo stesso Proxmox. Garage resta; B2 dopo. YouTube consigliato in prodotto. Disco e alert sono vincolo non negoziabile. + +**Dopo:** eventuale Auction Hetzner sostituisce il control plane casa e WireGuard → vSwitch, senza cambiare il modello nodi/autoscaler. + +--- + +## 16. WireGuard (Proxmox casa ↔ Hetzner Cloud) + +Obiettivo: Rails/Sidekiq sull’host di produzione raggiungono `api_base_url` / `internal_rtmp_url` dei nodi Cloud (es. `http://10.0.0.9:9997`) **senza** esporre Redis/Postgres/MediaMTX API su Internet. + +### Setup consigliato (una tantum) + +1. Su Hetzner Console: crea **Network** privata (es. `10.0.0.0/16`) nella location `fsn1`, subnet `10.0.1.0/24` per i server Cloud. +2. Imposta `HCLOUD_NETWORK_ID=` così i nodi stream entrano nella network al create. +3. Su Proxmox (VM o LXC gateway): installa WireGuard; peer verso una VM “gateway” Cloud **sempre on** (CX22 piccola) *oppure* peer site-to-site verso un server Cloud fisso. +4. Alternativa più semplice per i primi test: **Tailscale** su Proxmox host + su ogni stream-node (cloud-init) — stesso effetto, meno ops. +5. Firewall Hetzner Cloud: + - WAN: `1935/tcp` (RTMP), `443/tcp` o `8888` se HLS diretto + - Solo rete privata / WG: `9997` (MediaMTX API), niente Postgres/Redis sui nodi stream +6. Verifica da Rails: `curl http://:9997/v3/paths/list` + +Finché WireGuard/Tailscale non è pronto, `api_base_url` punta comunque alla private IP (o pubblica se manca private_net): il path MediaMTX da Create fallirà dal control plane se non raggiungibile — provisionare nodi Cloud solo dopo connettività privata OK, oppure smoke con API su IP pubblico temporaneo + firewall allowlist IP casa. + +### ENV produzione (stream autoscale) + +```bash +HCLOUD_TOKEN=... +HCLOUD_LOCATION=fsn1 +HCLOUD_SERVER_TYPE=cpx21 +HCLOUD_IMAGE=debian-12 +HCLOUD_SSH_KEY=matchlivetv-stream +HCLOUD_NETWORK_ID= # opzionale +HCLOUD_USER_DATA_FILE=/opt/matchlivetv/infra/stream-node/cloud-init.yaml + +STREAM_DNS_ZONE=mltv-stream.net +STREAM_DNS_TTL=60 +STREAM_CLOUD_DNS_SUFFIX=mltv-stream.net +STREAM_CLOUD_MAX_PUBLISHERS=4 +STREAM_NODE_ENV=prod + +# Default lab sicuro; in prod per overflow usa hetzner esplicitamente dal bottone admin +STREAM_CLOUD_PROVIDER=local_lab +STREAM_DNS_PROVIDER=lab +``` diff --git a/infra/.env.production.example b/infra/.env.production.example index bc56c6b..95b9d6d 100644 --- a/infra/.env.production.example +++ b/infra/.env.production.example @@ -104,3 +104,21 @@ OPS_LOG_SUBSCRIBER=false # Sentry (opzionale — error tracking produzione) SENTRY_DSN= SENTRY_TRACES_SAMPLE_RATE=0.1 + +# --- Stream autoscale (Hetzner Cloud + DNS mltv-stream.net) --- +# Token progetto matchlivetv-stream (Read & Write). Non committare il valore reale. +HCLOUD_TOKEN= +HCLOUD_LOCATION=fsn1 +HCLOUD_SERVER_TYPE=cpx21 +HCLOUD_IMAGE=debian-12 +HCLOUD_SSH_KEY=matchlivetv-stream +HCLOUD_NETWORK_ID= +HCLOUD_USER_DATA_FILE=/opt/matchlivetv/infra/stream-node/cloud-init.yaml +STREAM_DNS_ZONE=mltv-stream.net +STREAM_DNS_TTL=60 +STREAM_CLOUD_DNS_SUFFIX=mltv-stream.net +STREAM_CLOUD_MAX_PUBLISHERS=4 +STREAM_NODE_ENV=prod +# Lab locale di default; il bottone admin "Provisiona nodo Hetzner" usa comunque i provider hetzner. +STREAM_CLOUD_PROVIDER=local_lab +STREAM_DNS_PROVIDER=lab diff --git a/infra/stream-node/cloud-init.yaml b/infra/stream-node/cloud-init.yaml new file mode 100644 index 0000000..e51f098 --- /dev/null +++ b/infra/stream-node/cloud-init.yaml @@ -0,0 +1,41 @@ +#cloud-config +# Cloud-init minimale per nodo stream Hetzner (MediaMTX + ffmpeg). +# Uso in .env produzione: +# HCLOUD_USER_DATA_FILE=/opt/matchlivetv/infra/stream-node/cloud-init.yaml +# +# Per boot più veloci: creare uno snapshot dopo il primo setup e usare HCLOUD_IMAGE=. + +package_update: true +packages: + - docker.io + - ffmpeg + - curl +runcmd: + - systemctl enable --now docker + - mkdir -p /opt/stream-node /recordings /slates + - | + cat >/opt/stream-node/docker-compose.yml <<'EOF' + services: + mediamtx: + image: bluenviron/mediamtx:latest + network_mode: host + restart: unless-stopped + volumes: + - ./mediamtx.yml:/mediamtx.yml:ro + - /recordings:/recordings + - /slates:/slates:ro + command: /mediamtx.yml + EOF + - | + cat >/opt/stream-node/mediamtx.yml <<'EOF' + logLevel: info + api: yes + apiAddress: :9997 + rtmp: yes + rtmpAddress: :1935 + hls: yes + hlsAddress: :8888 + paths: + all_others: + EOF + - cd /opt/stream-node && docker compose up -d