diff --git a/backend/app/controllers/api/v1/base_controller.rb b/backend/app/controllers/api/v1/base_controller.rb index 440d9fb..8d114a0 100644 --- a/backend/app/controllers/api/v1/base_controller.rb +++ b/backend/app/controllers/api/v1/base_controller.rb @@ -2,6 +2,7 @@ module Api module V1 class BaseController < ApplicationController rescue_from Teams::EntitlementError, with: :render_entitlement_error + rescue_from Streams::IngestUnavailableError, with: :render_ingest_unavailable rescue_from BrandingAttachments::CoverUploadError, with: :render_cover_upload_error rescue_from Youtube::BroadcastService::Error, with: :render_youtube_error @@ -22,6 +23,13 @@ module Api }, status: :forbidden end + def render_ingest_unavailable(error) + render json: { + error: error.message, + error_code: error.code + }, status: :service_unavailable + end + def render_cover_upload_error(error) render json: { error: error.message, error_code: "cover_upload_invalid" }, status: :unprocessable_entity end diff --git a/backend/app/services/mediamtx/client.rb b/backend/app/services/mediamtx/client.rb index e89bd5a..d6b9526 100644 --- a/backend/app/services/mediamtx/client.rb +++ b/backend/app/services/mediamtx/client.rb @@ -4,6 +4,9 @@ module Mediamtx class Client class Error < StandardError; end + CREATE_PATH_RETRIES = -> { ENV.fetch("MEDIAMTX_CREATE_RETRIES", "5").to_i } + CREATE_PATH_RETRY_BASE_SECS = -> { ENV.fetch("MEDIAMTX_CREATE_RETRY_BASE_SECS", "0.4").to_f } + def self.for_session(session) new(base_url: session.mediamtx_api_base_url) end @@ -19,6 +22,19 @@ module Mediamtx attr_reader :base_url + # Health probe for CPX readiness (GET /v3/paths/list). + def reachable?(timeout: 2) + conn = Faraday.new(url: @base_url) do |f| + f.adapter Faraday.default_adapter + f.options.open_timeout = timeout + f.options.timeout = timeout + end + response = conn.get("/v3/paths/list") + response.success? + rescue Faraday::Error + false + end + def create_path(session) path = session.mediamtx_path_name # record: false finché non c'è publisher — con alwaysAvailable MediaMTX registrerebbe @@ -30,7 +46,9 @@ module Mediamtx body[:alwaysAvailable] = true body[:alwaysAvailableFile] = slate_file_path(session) # YouTube: telefono → MediaMTX; relay copy verso RTMPS in sidekiq. - response = @conn.post("/v3/config/paths/add/#{CGI.escape(path)}", body) + response = with_connection_retries("create_path #{path}") do + @conn.post("/v3/config/paths/add/#{CGI.escape(path)}", body) + end unless response.success? err = response.body.is_a?(Hash) ? response.body["error"] : response.body raise Error, "MediaMTX path create failed: #{response.status} #{err}" @@ -161,6 +179,23 @@ module Mediamtx private + def with_connection_retries(label) + attempts = [CREATE_PATH_RETRIES.call, 1].max + base = CREATE_PATH_RETRY_BASE_SECS.call + try = 0 + begin + try += 1 + yield + rescue Faraday::ConnectionFailed, Faraday::TimeoutError => e + raise if try >= attempts + + sleep_secs = base * (2**(try - 1)) + Rails.logger.warn("[Mediamtx::Client] #{label} retry #{try}/#{attempts} after #{e.class}: #{e.message} (sleep #{sleep_secs}s)") + sleep(sleep_secs) + retry + end + end + def recording_body(session, enabled:) ent = session.match.team.entitlements can_record = ent.recording_enabled_for_mediamtx? diff --git a/backend/app/services/sessions/create.rb b/backend/app/services/sessions/create.rb index 0e3c743..c517032 100644 --- a/backend/app/services/sessions/create.rb +++ b/backend/app/services/sessions/create.rb @@ -58,6 +58,13 @@ module Sessions session rescue Streams::NodeRegistry::NoCapacityError => e raise Teams::EntitlementError.new(e.message, code: "stream_capacity_exhausted") + rescue Faraday::ConnectionFailed, Faraday::TimeoutError => e + raise Streams::IngestUnavailableError, "Ingest temporaneamente non disponibile (#{e.class})" + rescue Mediamtx::Client::Error => e + if e.message.to_s.match?(/failed to open|Connection refused|Timeout|timed out/i) + raise Streams::IngestUnavailableError, e.message + end + raise end private diff --git a/backend/app/services/streams/autoscaler.rb b/backend/app/services/streams/autoscaler.rb index 902eb52..95b5e60 100644 --- a/backend/app/services/streams/autoscaler.rb +++ b/backend/app/services/streams/autoscaler.rb @@ -136,6 +136,8 @@ module Streams def reconcile! actions = [] + actions.concat(promote_provisioning_nodes!) + actions.concat(reclaim_stuck_provisioning!) m = self.class.metrics if need_capacity?(m) && can_provision?(m) @@ -170,6 +172,30 @@ module Streams private + def promote_provisioning_nodes! + actions = [] + StreamNode.where(status: "provisioning").find_each do |node| + next unless Streams::NodeHealth.promote_if_healthy!(node) + + actions << :"ready_#{node.slug}" + Rails.logger.info("[Streams::Autoscaler] promoted #{node.slug} to ready") + end + actions + end + + def reclaim_stuck_provisioning! + actions = [] + stuck_after = ENV.fetch("STREAM_NODE_PROVISIONING_STUCK_MINUTES", "15").to_i.minutes.ago + StreamNode.where(status: "provisioning").where("created_at < ?", stuck_after).find_each do |node| + @provisioner.decommission!(node) + actions << :"reclaim_#{node.slug}" + Rails.logger.warn("[Streams::Autoscaler] decommissioned stuck provisioning #{node.slug}") + rescue NodeProvisioner::BusyError, NodeProvisioner::Error => e + Rails.logger.warn("[Streams::Autoscaler] reclaim #{node.slug}: #{e.message}") + end + actions + end + def need_capacity?(m) m[:free_slots] <= self.class.soft_free_slots end diff --git a/backend/app/services/streams/ingest_unavailable_error.rb b/backend/app/services/streams/ingest_unavailable_error.rb new file mode 100644 index 0000000..2261579 --- /dev/null +++ b/backend/app/services/streams/ingest_unavailable_error.rb @@ -0,0 +1,13 @@ +# frozen_string_literal: true + +module Streams + # Ingest MediaMTX non raggiungibile (es. CPX ancora in bootstrap). API → 503 retryable. + class IngestUnavailableError < StandardError + attr_reader :code + + def initialize(message = "Ingest temporaneamente non disponibile", code: "stream_ingest_unavailable") + super(message) + @code = code + end + end +end diff --git a/backend/app/services/streams/node_health.rb b/backend/app/services/streams/node_health.rb new file mode 100644 index 0000000..e1580ff --- /dev/null +++ b/backend/app/services/streams/node_health.rb @@ -0,0 +1,27 @@ +# frozen_string_literal: true + +module Streams + # Probe reachability of MediaMTX API on a stream node (:9997). + class NodeHealth + def self.mediamtx_up?(node, timeout: 2) + new(node, timeout: timeout).mediamtx_up? + end + + def self.promote_if_healthy!(node, timeout: 2) + return false unless node.status == "provisioning" + return false unless mediamtx_up?(node, timeout: timeout) + + node.update!(status: "ready", last_health_at: Time.current) + true + end + + def initialize(node, timeout: 2) + @node = node + @timeout = timeout + end + + def mediamtx_up? + Mediamtx::Client.new(base_url: @node.api_base_url).reachable?(timeout: @timeout) + end + end +end diff --git a/backend/app/services/streams/node_provisioner.rb b/backend/app/services/streams/node_provisioner.rb index 2062d4f..c61f2f2 100644 --- a/backend/app/services/streams/node_provisioner.rb +++ b/backend/app/services/streams/node_provisioner.rb @@ -80,11 +80,15 @@ module Streams urls = urls_for_node(role: role, home: home, hostname: hostname, simulated: simulated, private_ip: private_ip, public_ip: ip, use_node_hostname: use_node_hostname) - StreamNode.create!( + # Simulated/lab che riusa MediaMTX home: subito ready. Cloud reale: provisioning + # finché :9997 risponde (evita allocate → Faraday Connection refused → 500). + initial_status = simulated ? "ready" : "provisioning" + + node = StreamNode.create!( slug: slug, hostname: hostname, role: role, - status: "ready", + status: initial_status, provider: provider_name_for(cloud, role: role), provider_instance_id: instance.id, rtmp_base_url: urls.fetch(:rtmp_base_url), @@ -102,6 +106,34 @@ module Streams "cloud_raw" => instance.raw } ) + + wait_until_mediamtx_ready!(node) unless simulated + node.reload + end + + def wait_until_mediamtx_ready!(node, timeout: nil, interval: nil) + timeout ||= ENV.fetch("STREAM_NODE_READY_TIMEOUT_SECS", "180").to_i + interval ||= ENV.fetch("STREAM_NODE_READY_POLL_SECS", "3").to_f + deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout + + loop do + if Streams::NodeHealth.promote_if_healthy!(node) + Rails.logger.info("[Streams::NodeProvisioner] #{node.slug} MediaMTX ready") + return node + end + + remaining = deadline - Process.clock_gettime(Process::CLOCK_MONOTONIC) + if remaining <= 0 + Rails.logger.warn( + "[Streams::NodeProvisioner] #{node.slug} ancora provisioning dopo #{timeout}s " \ + "(api=#{node.api_base_url}); autoscaler continuerà a promuovere" + ) + return node + end + + sleep([interval, remaining].min) + node.reload + end end def urls_for_node(role:, home:, hostname:, simulated:, private_ip:, public_ip: nil, use_node_hostname:) diff --git a/backend/spec/requests/api/v1/sessions_ingest_unavailable_spec.rb b/backend/spec/requests/api/v1/sessions_ingest_unavailable_spec.rb new file mode 100644 index 0000000..d561433 --- /dev/null +++ b/backend/spec/requests/api/v1/sessions_ingest_unavailable_spec.rb @@ -0,0 +1,52 @@ +# frozen_string_literal: true + +require "rails_helper" + +RSpec.describe "API session create ingest unavailable", type: :request do + let!(:user) do + User.create!(email: "ingest-#{SecureRandom.hex(4)}@example.com", name: "Coach", password: "Password123", role: "coach") + end + let!(:club) { Club.create!(name: "IngestClub", sport: "volleyball") } + let!(:membership) { club.club_memberships.create!(user: user, role: "owner") } + let!(:team) { club.teams.create!(name: "Tigers", sport: "volleyball", slug: "ingest-tigers-#{SecureRandom.hex(3)}") } + let!(:match) { team.matches.create!(opponent_name: "Opp", scheduled_at: 1.hour.from_now) } + + def auth_headers + post "/api/v1/auth/login", params: { email: user.email, password: "Password123" } + token = response.parsed_body["access_token"] + { "Authorization" => "Bearer #{token}", "Content-Type" => "application/json" } + end + + before do + plan = Plan.find_or_initialize_by(slug: "premium_full") + plan.name ||= "Premium Full" + plan.features = (plan.features || {}).merge( + "platforms" => %w[matchlivetv youtube], + "youtube_enabled" => true, + "concurrent_streams_limit" => 10, + "recordings_enabled" => true, + "recording_days" => 90 + ) + plan.save! + club.create_subscription!(plan: plan, status: "active") if club.subscription.blank? + + allow_any_instance_of(Teams::Entitlements).to receive(:assert_can_stream_on!) + allow_any_instance_of(Teams::Entitlements).to receive(:assert_concurrent_stream!) + allow(Streams::CoverSlateEnsurer).to receive(:ensure_for!) + allow(Streams::SlateDistributor).to receive(:ensure_for!) + Streams::NodeRegistry.ensure_home_from_env! + end + + it "returns 503 when MediaMTX connection fails after retries" do + allow_any_instance_of(Mediamtx::Client).to receive(:create_path) + .and_raise(Faraday::ConnectionFailed.new("Connection refused")) + + post "/api/v1/matches/#{match.id}/sessions", + params: { platform: "matchlivetv", privacy_status: "private" }.to_json, + headers: auth_headers + + expect(response).to have_http_status(:service_unavailable) + body = response.parsed_body + expect(body["error_code"]).to eq("stream_ingest_unavailable") + end +end diff --git a/backend/spec/services/mediamtx/client_create_path_retry_spec.rb b/backend/spec/services/mediamtx/client_create_path_retry_spec.rb new file mode 100644 index 0000000..19f4cbd --- /dev/null +++ b/backend/spec/services/mediamtx/client_create_path_retry_spec.rb @@ -0,0 +1,50 @@ +# frozen_string_literal: true + +require "rails_helper" + +RSpec.describe Mediamtx::Client do + describe "#create_path" do + let(:session) do + user = User.create!(email: "mtx-#{SecureRandom.hex(4)}@example.com", name: "M", password: "Password123", role: "coach") + club = Club.create!(name: "MtxClub-#{SecureRandom.hex(3)}", sport: "volleyball") + team = club.teams.create!(name: "T", sport: "volleyball", slug: "mtx-#{SecureRandom.hex(4)}") + match = team.matches.create!(opponent_name: "X", scheduled_at: 1.hour.from_now) + StreamSession.create!(match: match, user: user, platform: "matchlivetv", status: "idle") + end + + it "retries Faraday connection failures then succeeds" do + ENV["MEDIAMTX_CREATE_RETRIES"] = "3" + ENV["MEDIAMTX_CREATE_RETRY_BASE_SECS"] = "0" + + client = described_class.new(base_url: "http://mtx.test:9997") + conn = instance_double(Faraday::Connection) + client.instance_variable_set(:@conn, conn) + + fail_once = Faraday::ConnectionFailed.new("Connection refused") + ok = instance_double(Faraday::Response, success?: true, status: 200, body: {}) + + expect(conn).to receive(:post).once.and_raise(fail_once) + expect(conn).to receive(:post).once.and_return(ok) + + expect(client.create_path(session)).to eq(true) + ensure + ENV.delete("MEDIAMTX_CREATE_RETRIES") + ENV.delete("MEDIAMTX_CREATE_RETRY_BASE_SECS") + end + + it "raises after exhausting connection retries" do + ENV["MEDIAMTX_CREATE_RETRIES"] = "2" + ENV["MEDIAMTX_CREATE_RETRY_BASE_SECS"] = "0" + + client = described_class.new(base_url: "http://mtx.test:9997") + conn = instance_double(Faraday::Connection) + client.instance_variable_set(:@conn, conn) + allow(conn).to receive(:post).and_raise(Faraday::ConnectionFailed.new("Connection refused")) + + expect { client.create_path(session) }.to raise_error(Faraday::ConnectionFailed) + ensure + ENV.delete("MEDIAMTX_CREATE_RETRIES") + ENV.delete("MEDIAMTX_CREATE_RETRY_BASE_SECS") + end + end +end diff --git a/backend/spec/services/streams/autoscaler_spec.rb b/backend/spec/services/streams/autoscaler_spec.rb index 1ab061e..f86665e 100644 --- a/backend/spec/services/streams/autoscaler_spec.rb +++ b/backend/spec/services/streams/autoscaler_spec.rb @@ -224,4 +224,41 @@ RSpec.describe Streams::Autoscaler do expect(described_class.within_budget?(1)).to eq(false) end end + + it "promotes provisioning nodes when MediaMTX becomes reachable" do + with_env( + "STREAM_AUTOSCALE_ENABLED" => "1", + "STREAM_AUTOSCALE_WARM_SPARE" => "0", + "STREAM_AUTOSCALE_SOFT_FREE_SLOTS" => "0", + "STREAM_AUTOSCALE_KIND" => "lab", + "MEDIAMTX_API_URL" => "http://mtx-home:9997", + "MEDIAMTX_RTMP_URL" => "rtmp://home.example:1935", + "HLS_PUBLIC_URL" => "https://home.example/hls" + ) do + Streams::NodeRegistry.ensure_home_from_env! + node = StreamNode.create!( + slug: "ingest-lab-prov", + hostname: "prov.lab", + role: "lab", + status: "provisioning", + provider: "local", + rtmp_base_url: "rtmp://h:1935", + hls_base_url: "https://h/hls", + api_base_url: "http://h:9997", + max_publishers: 2, + max_relays: 2 + ) + allow(Streams::NodeHealth).to receive(:promote_if_healthy!) do |n| + n.update!(status: "ready", last_health_at: Time.current) + true + end + + provisioner = instance_double(Streams::NodeProvisioner) + expect(provisioner).not_to receive(:provision_lab!) + + result = described_class.reconcile!(provisioner: provisioner) + expect(result.actions).to include(:"ready_ingest-lab-prov") + expect(node.reload.status).to eq("ready") + end + end end diff --git a/backend/spec/services/streams/node_health_spec.rb b/backend/spec/services/streams/node_health_spec.rb new file mode 100644 index 0000000..97b9566 --- /dev/null +++ b/backend/spec/services/streams/node_health_spec.rb @@ -0,0 +1,39 @@ +# frozen_string_literal: true + +require "rails_helper" + +RSpec.describe Streams::NodeHealth do + let!(:node) do + StreamNode.create!( + slug: "ingest-health-01", + hostname: "ingest-health-01.mltv-stream.net", + role: "cloud", + status: "provisioning", + provider: "hetzner", + rtmp_base_url: "rtmp://h:1935", + hls_base_url: "https://h/hls", + api_base_url: "http://203.0.113.10:9997", + max_publishers: 4, + max_relays: 4 + ) + end + + after { node.destroy } + + it "promotes provisioning node when MediaMTX is reachable" do + client = instance_double(Mediamtx::Client, reachable?: true) + allow(Mediamtx::Client).to receive(:new).with(base_url: node.api_base_url).and_return(client) + + expect(described_class.promote_if_healthy!(node)).to eq(true) + expect(node.reload.status).to eq("ready") + expect(node.last_health_at).to be_present + end + + it "does not promote when MediaMTX is down" do + client = instance_double(Mediamtx::Client, reachable?: false) + allow(Mediamtx::Client).to receive(:new).with(base_url: node.api_base_url).and_return(client) + + expect(described_class.promote_if_healthy!(node)).to eq(false) + expect(node.reload.status).to eq("provisioning") + end +end diff --git a/backend/spec/services/streams/node_provisioner_cloud_spec.rb b/backend/spec/services/streams/node_provisioner_cloud_spec.rb index e167bee..32fad7c 100644 --- a/backend/spec/services/streams/node_provisioner_cloud_spec.rb +++ b/backend/spec/services/streams/node_provisioner_cloud_spec.rb @@ -24,11 +24,15 @@ RSpec.describe "Streams::NodeProvisioner cloud" do ENV["STREAM_CLOUD_DNS_SUFFIX"] = "mltv-stream.net" ENV["STREAM_CLOUD_MAX_PUBLISHERS"] = "4" ENV["STREAM_CLOUD_PUBLIC_CONTROL"] = "0" + ENV["STREAM_NODE_READY_TIMEOUT_SECS"] = "0" + + allow(Streams::NodeHealth).to receive(:promote_if_healthy!).and_return(false) 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.status).to eq("provisioning") 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") @@ -39,6 +43,45 @@ RSpec.describe "Streams::NodeProvisioner cloud" do %w[ MEDIAMTX_API_URL MEDIAMTX_RTMP_URL HLS_PUBLIC_URL STREAM_CLOUD_DNS_SUFFIX STREAM_CLOUD_MAX_PUBLISHERS STREAM_CLOUD_PUBLIC_CONTROL + STREAM_NODE_READY_TIMEOUT_SECS + ].each { |k| ENV.delete(k) } + end + + it "marks cloud node ready when MediaMTX answers during wait" do + cloud = instance_double( + Streams::CloudProviders::Hetzner, + create_node: Streams::CloudProviders::Instance.new( + id: "100", + 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_PUBLIC_CONTROL"] = "0" + ENV["STREAM_NODE_READY_TIMEOUT_SECS"] = "5" + ENV["STREAM_NODE_READY_POLL_SECS"] = "0.01" + + allow(Streams::NodeHealth).to receive(:promote_if_healthy!) do |node| + node.update!(status: "ready", last_health_at: Time.current) + true + end + + node = Streams::NodeProvisioner.new(cloud: cloud, dns: dns).provision_cloud! + expect(node.status).to eq("ready") + ensure + %w[ + MEDIAMTX_API_URL MEDIAMTX_RTMP_URL HLS_PUBLIC_URL + STREAM_CLOUD_DNS_SUFFIX STREAM_CLOUD_PUBLIC_CONTROL + STREAM_NODE_READY_TIMEOUT_SECS STREAM_NODE_READY_POLL_SECS ].each { |k| ENV.delete(k) } end end