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.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 %> +
+ +| <%= 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 %> + | +
<%= @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.